diff --git a/CHANGELOG.md b/CHANGELOG.md index 04cf38ea..1968dd9b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -16,7 +16,9 @@ If you are upgrading from v0.16.x, replace the binary (or run `docker pull`). If - Invalidate caches when group memberships change on an external directory. - Return OIDC errors instead of "failed to decode token". - User impersonation. -- Task queue pagination by anchor. +- Tasks: + - Delete locked tasks. + - Queue pagination by anchor. - Log viewer: All events show as `INFO`. - Registry: Allow changing object variants. - Node id renewal. diff --git a/crates/jmap/src/registry/mapping/task.rs b/crates/jmap/src/registry/mapping/task.rs index 12fad496..0d0b97c3 100644 --- a/crates/jmap/src/registry/mapping/task.rs +++ b/crates/jmap/src/registry/mapping/task.rs @@ -259,15 +259,6 @@ pub(crate) async fn task_set( continue; } - if !set.server.try_lock_task(task_id).await { - set.response.not_destroyed.append( - id, - SetError::forbidden().with_description( - "Task is currently being processed and cannot be destroyed".to_string(), - ), - ); - continue; - } locked_tasks.push(task_id); let due = task.due_timestamp(); diff --git a/crates/services/src/task_manager/lock.rs b/crates/services/src/task_manager/lock.rs index 33b9984f..cb11a2e6 100644 --- a/crates/services/src/task_manager/lock.rs +++ b/crates/services/src/task_manager/lock.rs @@ -24,7 +24,6 @@ impl TaskLockManager for Server { TaskManager(TaskManagerEvent::TaskLocked), Id = id, Details = "Task details not available", - Expires = trc::Value::Timestamp(now() + DEFAULT_LOCK_EXPIRY), ); } result diff --git a/crates/services/src/task_manager/manager.rs b/crates/services/src/task_manager/manager.rs index b8f89317..c76a61c9 100644 --- a/crates/services/src/task_manager/manager.rs +++ b/crates/services/src/task_manager/manager.rs @@ -40,7 +40,7 @@ use store::rand::seq::SliceRandom; use store::write::key::DeserializeBigEndian; use store::{ IterateParams, ValueKey, - write::{BatchBuilder, TaskQueueClass, ValueClass, now}, + write::{BatchBuilder, TaskQueueClass, ValueClass, assert::AssertValue, now}, }; use store::{SerializeInfallible, U64_LEN, rand}; use tokio::sync::{mpsc, watch}; @@ -574,6 +574,10 @@ async fn update_tasks( u64::MAX }; batch + .assert_value( + ValueClass::TaskQueue(TaskQueueClass::Task { id }), + AssertValue::Some, + ) .set( ValueClass::TaskQueue(TaskQueueClass::Due { id, due }), task.info.typ.to_id().serialize(), @@ -587,7 +591,14 @@ async fn update_tasks( } if let Err(err) = server.store().write(batch.build_all()).await { - trc::error!(err.details("Failed to remove task(s) from queue.")); + if err.matches(trc::EventType::Store(trc::StoreEvent::AssertValueFailed)) { + trc::event!( + TaskManager(TaskManagerEvent::TaskIgnored), + Reason = "Task was deleted while being processed; skipping update.", + ); + } else { + trc::error!(err.details("Failed to remove task(s) from queue.")); + } } for task in tasks {