Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions migrations/019_queue_unique_constraints.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
-- Prevent duplicate pending operations in the queue
-- Partial unique indexes ensure one pending operation per (mod/addon, action)

CREATE UNIQUE INDEX IF NOT EXISTS idx_pending_ops_mod_action
ON pending_operations(forge_mod_id, action)
WHERE item_type = 'mod' AND forge_mod_id IS NOT NULL;

CREATE UNIQUE INDEX IF NOT EXISTS idx_pending_ops_addon_action
ON pending_operations(forge_addon_id, action)
WHERE item_type = 'addon' AND forge_addon_id IS NOT NULL;
25 changes: 22 additions & 3 deletions src/db/users.rs
Original file line number Diff line number Diff line change
Expand Up @@ -451,7 +451,7 @@ impl Database {
// ── Pending Operations CRUD ───────────────────────────────────────

pub fn insert_pending_op(&self, op: &InsertPendingOp<'_>) -> rusqlite::Result<i64> {
self.conn.execute(
match self.conn.execute(
"INSERT INTO pending_operations (action, forge_mod_id, forge_version_id, mod_name, metadata, queued_by, item_type, forge_addon_id, archive_path, source, source_url)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
params![
Expand All @@ -467,8 +467,18 @@ impl Database {
op.source,
op.source_url,
],
)?;
Ok(self.conn.last_insert_rowid())
) {
Ok(_) => Ok(self.conn.last_insert_rowid()),
Err(rusqlite::Error::SqliteFailure(err, _))
if err.extended_code == rusqlite::ffi::SQLITE_CONSTRAINT_UNIQUE =>
{
Err(rusqlite::Error::SqliteFailure(
err,
Some("Operation already queued".to_string()),
))
}
Err(e) => Err(e),
}
}

pub fn has_pending_op(&self, forge_mod_id: i64, action: QueueAction) -> rusqlite::Result<bool> {
Expand All @@ -493,6 +503,15 @@ impl Database {
Ok(count > 0)
}

pub fn has_pending_url_op(&self, url: &str, action: QueueAction) -> rusqlite::Result<bool> {
let count: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM pending_operations WHERE source_url = ?1 AND action = ?2",
params![url, action.as_str()],
|row| row.get(0),
)?;
Ok(count > 0)
}

pub fn list_pending_ops(&self) -> rusqlite::Result<Vec<PendingOperation>> {
let mut stmt = self.conn.prepare(
"SELECT id, action, forge_mod_id, forge_version_id, mod_name, metadata, queued_at, queued_by, item_type, forge_addon_id, archive_path, source, source_url
Expand Down
120 changes: 104 additions & 16 deletions src/web/handlers/mods.rs
Original file line number Diff line number Diff line change
Expand Up @@ -394,7 +394,7 @@ async fn try_queue_mod_op(
// Remove doesn't need downloading
let db = state.db.clone();
let mod_name_owned = mod_name.to_string();
web::block(move || {
match web::block(move || {
let db = db.lock();
db.insert_pending_op(&crate::db::users::InsertPendingOp {
action: QueueAction::Remove,
Expand All @@ -411,9 +411,22 @@ async fn try_queue_mod_op(
})
})
.await
.map_err(WebError::from)?
.map_err(WebError::from)?;
set_flash(session, "Mod queued for removal", FlashType::Success);
{
Ok(Ok(_)) => {
set_flash(session, "Mod queued for removal", FlashType::Success);
}
Ok(Err(rusqlite::Error::SqliteFailure(err, _)))
if err.extended_code == rusqlite::ffi::SQLITE_CONSTRAINT_UNIQUE =>
{
set_flash(
session,
"This mod removal is already queued",
FlashType::Info,
);
}
Ok(Err(e)) => return Err(WebError::from(e)),
Err(e) => return Err(WebError::from(e)),
}
}
QueueAction::Install => {
let version_id = version_id.ok_or(WebError::BadRequest("missing version_id".into()))?;
Expand Down Expand Up @@ -525,7 +538,7 @@ async fn try_queue_addon_op(
let db = state.db.clone();
let addon_name_owned = addon_name.to_string();
let username = user.username.clone();
web::block(move || {
match web::block(move || {
let db = db.lock();
db.insert_pending_op(&crate::db::users::InsertPendingOp {
action: QueueAction::Remove,
Expand All @@ -542,9 +555,22 @@ async fn try_queue_addon_op(
})
})
.await
.map_err(WebError::from)?
.map_err(WebError::from)?;
set_flash(session, "Addon queued for removal", FlashType::Success);
{
Ok(Ok(_)) => {
set_flash(session, "Addon queued for removal", FlashType::Success);
}
Ok(Err(rusqlite::Error::SqliteFailure(err, _)))
if err.extended_code == rusqlite::ffi::SQLITE_CONSTRAINT_UNIQUE =>
{
set_flash(
session,
"This addon removal is already queued",
FlashType::Info,
);
}
Ok(Err(e)) => return Err(WebError::from(e)),
Err(e) => return Err(WebError::from(e)),
}
}
QueueAction::Install | QueueAction::Update => {
let version_id = version_id.ok_or(WebError::BadRequest("missing version_id".into()))?;
Expand Down Expand Up @@ -1178,6 +1204,28 @@ async fn install_mod_from_url(

// Queue if server running
if should_queue_operation(state).await {
// Check for duplicate URL operation
let db_check = state.db.clone();
let url_check = url.to_string();
let already_queued = web::block(move || {
let db = db_check.lock();
db.has_pending_url_op(&url_check, crate::db::users::QueueAction::Install)
})
.await
.map_err(WebError::from)?
.map_err(WebError::from)?;

if already_queued {
set_flash(
session,
"This URL is already queued for installation",
FlashType::Info,
);
return Ok(HttpResponse::SeeOther()
.insert_header(("Location", "/quma/mods"))
.finish());
}

let queue_dir = state.dirs.queue_dir();
let _ = std::fs::create_dir_all(&queue_dir);

Expand All @@ -1201,7 +1249,7 @@ async fn install_mod_from_url(
let mod_name_q = mod_name.clone();
let dest_str = dest.to_string_lossy().to_string();
let url_owned = url.to_string();
let _ = web::block(move || {
match web::block(move || {
let db = db.lock();
db.insert_pending_op(&crate::db::users::InsertPendingOp {
action: crate::db::users::QueueAction::Install,
Expand All @@ -1218,14 +1266,27 @@ async fn install_mod_from_url(
})
})
.await
.map_err(WebError::from)?
.map_err(WebError::from)?;
{
Ok(Ok(_)) => {
set_flash(
session,
"Mod queued for install from URL",
FlashType::Success,
);
}
Ok(Err(rusqlite::Error::SqliteFailure(err, msg)))
if err.extended_code == rusqlite::ffi::SQLITE_CONSTRAINT_UNIQUE =>
{
set_flash(
session,
msg.as_deref().unwrap_or("Operation already queued"),
FlashType::Info,
);
}
Ok(Err(e)) => return Err(WebError::from(e).into()),
Err(e) => return Err(WebError::from(e).into()),
}

set_flash(
session,
"Mod queued for install from URL",
FlashType::Success,
);
return Ok(HttpResponse::SeeOther()
.insert_header(("Location", "/quma/mods"))
.finish());
Expand Down Expand Up @@ -2933,6 +2994,33 @@ pub async fn remove_addon(
};

let parent_mod_id = addon.parent_mod_id;

// Check if the operation should be queued
let parent_forge_mod_id_opt = {
let db = state.db.lock();
db.get_mod(parent_mod_id)
.ok()
.flatten()
.and_then(|m| m.forge_mod_id)
};
if let Some(parent_forge_mod_id) = parent_forge_mod_id_opt {
if let Some(resp) = try_queue_addon_op(
&state,
&session,
&user,
QueueAction::Remove,
addon.forge_addon_id,
None,
&addon.name,
parent_forge_mod_id,
&format!("/quma/mods/{}#queue", parent_mod_id),
)
.await?
{
return Ok(resp);
}
}

let dirs = Arc::clone(&state.dirs);
let config = state.config_cloned();

Expand Down
Loading