diff --git a/crates/protocol/src/lib.rs b/crates/protocol/src/lib.rs index c28ebae9..fd5190c7 100644 --- a/crates/protocol/src/lib.rs +++ b/crates/protocol/src/lib.rs @@ -154,8 +154,8 @@ pub enum Method { /// Stop the in-flight turn and transition to `Paused`. /// /// Unlike `Cancel` (which discards and returns to `Idle`), a paused - /// Worker can resume the interrupted work via `Resume`, or start a - /// fresh turn via `Run` (orphan `tool_use` items are closed with a + /// Worker can resume the interrupted work via `Resume`, or accept a + /// fresh `Submit` (orphan `tool_use` items are closed with a /// synthetic tool result before the new user message is appended). Pause, /// Request an explicit compaction while the Worker is otherwise idle. diff --git a/crates/session-store/src/fs_store.rs b/crates/session-store/src/fs_store.rs index 0ebb1498..19a2c6d0 100644 --- a/crates/session-store/src/fs_store.rs +++ b/crates/session-store/src/fs_store.rs @@ -23,7 +23,8 @@ use crate::uploaded_file::{ bind_uploaded_file, clear_uploaded_file_binding, copy_committed_uploaded_files, delete_uncommitted_uploaded_files, delete_uploaded_file, finalize_uploaded_file_binding, list_uploaded_file_refs, pin_uploaded_file, read_uploaded_file, read_uploaded_file_by_id, - release_uploaded_file_pin, uploaded_file_has_pending_owner, write_uploaded_file, + reconcile_uploaded_file_pins, release_uploaded_file_pin, uploaded_file_has_pending_owner, + write_uploaded_file, }; use crate::{ PasteArtifactLimits, SegmentId, SessionId, UploadedFileLimits, UploadedFileUploadContext, @@ -562,6 +563,18 @@ impl Store for FsStore { ) } + fn reconcile_uploaded_file_pins( + &self, + session_id: SessionId, + live_owner_ids: &[String], + ) -> Result { + let _guard = self + .append_lock + .lock() + .map_err(|_| std::io::Error::other("session store append lock was poisoned"))?; + reconcile_uploaded_file_pins(&self.paste_artifact_dir(session_id), live_owner_ids) + } + fn delete_uploaded_file( &self, session_id: SessionId, diff --git a/crates/session-store/src/store.rs b/crates/session-store/src/store.rs index b14839fa..33daae39 100644 --- a/crates/session-store/src/store.rs +++ b/crates/session-store/src/store.rs @@ -256,6 +256,17 @@ pub trait Store: Send + Sync { Err(StoreError::PasteArtifactUnsupported) } + /// Clear pending-operation pins that have no owner in restored durable + /// Worker Session state. This repairs an interrupted pin-before-checkpoint + /// acceptance without disturbing live queue owners or committed history. + fn reconcile_uploaded_file_pins( + &self, + _session_id: SessionId, + _live_owner_ids: &[String], + ) -> Result { + Ok(0) + } + /// Delete an uncommitted uploaded file owned by `session_id`. fn delete_uploaded_file( &self, diff --git a/crates/session-store/src/uploaded_file.rs b/crates/session-store/src/uploaded_file.rs index c8474b0f..6583dd25 100644 --- a/crates/session-store/src/uploaded_file.rs +++ b/crates/session-store/src/uploaded_file.rs @@ -586,6 +586,40 @@ pub(crate) fn copy_committed_uploaded_files(source_dir: &Path, target_dir: &Path Ok(copied) } +pub(crate) fn reconcile_uploaded_file_pins(dir: &Path, live_owner_ids: &[String]) -> Result { + fs::create_dir_all(dir)?; + let aggregate_lock = fs::OpenOptions::new() + .create(true) + .read(true) + .write(true) + .open(dir.join(".aggregate.lock"))?; + FileExt::lock_exclusive(&aggregate_lock)?; + let mut reconciled = 0_u64; + for entry in fs::read_dir(dir)? { + let entry = entry?; + let path = entry.path(); + let Some(file_name) = path.file_name().and_then(|name| name.to_str()) else { + continue; + }; + let Some(artifact_id) = file_name.strip_suffix(".file.json") else { + continue; + }; + let mut stored: StoredUploadedFile = serde_json::from_slice(&fs::read(&path)?)?; + let Some(owner_id) = stored.pending_owner_id.as_deref() else { + continue; + }; + if live_owner_ids.iter().any(|live| live == owner_id) { + continue; + } + stored.pending_owner_id = None; + let temp = dir.join(format!(".{artifact_id}.file.reconcile.tmp")); + fs::write(&temp, serde_json::to_vec(&stored)?)?; + fs::rename(temp, path)?; + reconciled = reconciled.saturating_add(1); + } + Ok(reconciled) +} + pub(crate) fn delete_uncommitted_uploaded_files(dir: &Path) -> Result { fs::create_dir_all(dir)?; let aggregate_lock = fs::OpenOptions::new() diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index 0b79c4f0..67c7f7ed 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -173,6 +173,14 @@ impl PendingActivationState { } } + fn live_artifact_pin_owner_ids(&self) -> Vec { + self.activating + .iter() + .chain(self.pending.iter()) + .map(|submission| submission.submission_id.clone()) + .collect() + } + fn remember_notification_receipt(&mut self, receipt: NotificationReceipt) { self.notification_receipts.push_back(receipt); while self.notification_receipts.len() > MAX_SUBMISSION_RECEIPTS { @@ -1830,6 +1838,18 @@ where .snapshot() } + fn reconcile_uploaded_file_pins(&self) -> Result { + let live_owner_ids = self + .state + .lock() + .expect("pending activation state poisoned") + .live_artifact_pin_owner_ids(); + Ok(self + .writer + .store + .reconcile_uploaded_file_pins(self.writer.state.session_id(), &live_owner_ids)?) + } + pub(crate) fn cancel( &self, submission_id: &str, @@ -6340,6 +6360,9 @@ where worker .session .restore_pending_activations(&state.extensions); + worker + .pending_submission_handle() + .reconcile_uploaded_file_pins()?; worker.apply_permissions_from_manifest(); worker.apply_prune_from_manifest(); worker.write_worker_metadata_active(SegmentLocation { @@ -9703,6 +9726,56 @@ mod build_summary_prompt_tests { assert_eq!(state.pending[1].provenance, account_b); } + #[test] + fn restore_reconciliation_clears_interrupted_acceptance_pin_for_retry() { + let temp = tempfile::tempdir().unwrap(); + let handle = PendingSubmissionHandle::for_test(temp.path()); + let session_id = handle.writer.state.session_id(); + let segment_id = handle.writer.state.location().segment_id; + let limits = session_store::UploadedFileLimits { + max_file_bytes: 1024, + max_session_bytes: 2048, + }; + let file = handle + .writer + .store + .write_uploaded_file(session_id, "retry.txt", "text/plain", b"retry", limits) + .unwrap(); + handle + .writer + .store + .pin_uploaded_file(session_id, &file, "interrupted-before-checkpoint") + .unwrap(); + drop(handle); + let handle = PendingSubmissionHandle { + state: Arc::new(Mutex::new(PendingActivationState::default())), + writer: LogWriterHandle { + store: session_store::FsStore::new(temp.path()).unwrap(), + state: SegmentState::new(session_id, segment_id, 0), + sink: SegmentLogSink::new(), + in_flight: None, + }, + }; + + assert_eq!(handle.reconcile_uploaded_file_pins().unwrap(), 1); + let accepted = handle + .accept( + "request-after-restore".into(), + vec![Segment::UploadedFile { file: file.clone() }], + false, + ) + .unwrap(); + assert!(!accepted.submission_id.is_empty()); + assert_eq!(handle.reconcile_uploaded_file_pins().unwrap(), 0); + assert!(matches!( + handle + .writer + .store + .delete_uploaded_file(session_id, &file.artifact_id), + Err(StoreError::ArtifactAlreadyCommitted) + )); + } + #[test] fn rejected_submission_rolls_back_uploaded_file_pins_acquired_before_conflict() { let temp = tempfile::tempdir().unwrap();