From 27e5df106f86b40356e95ef6db76f6d9f0fe71fc Mon Sep 17 00:00:00 2001 From: Hare Date: Fri, 4 Sep 2026 19:16:01 +0900 Subject: [PATCH] fix: complete memory lifecycle feature boundaries --- crates/worker/src/controller.rs | 10 +- .../src/feature/builtin/memory_lifecycle.rs | 304 +++++++++++++++--- crates/worker/src/feature/session.rs | 8 + crates/worker/src/worker.rs | 40 ++- 4 files changed, 312 insertions(+), 50 deletions(-) diff --git a/crates/worker/src/controller.rs b/crates/worker/src/controller.rs index 43f28539..ddef4480 100644 --- a/crates/worker/src/controller.rs +++ b/crates/worker/src/controller.rs @@ -994,13 +994,7 @@ where let worker_enabled = feature_config.worker.enabled; let sub_worker_enabled = feature_config.sub_worker.enabled; let mut feature_registry = FeatureRegistryBuilder::new(); - if feature_config.memory.enabled { - let config = memory_config.clone().ok_or_else(|| { - std::io::Error::new( - std::io::ErrorKind::InvalidInput, - "[feature.memory].enabled = true requires a [memory] configuration section", - ) - })?; + if let Some(config) = memory_config.clone() { let workspace_client = worker.workspace_client_handle(); if !workspace_client.is_available() || workspace_client.workspace_id().is_none() { return Err(std::io::Error::new( @@ -1009,7 +1003,7 @@ where )); } feature_registry.add_module( - crate::feature::builtin::memory_lifecycle::MemoryExtractionLifecycleFeature::new( + crate::feature::builtin::memory_lifecycle::MemoryLifecycleFeature::new( config, worker.committed_session_capture_handle(), worker.session_extension_handle(), diff --git a/crates/worker/src/feature/builtin/memory_lifecycle.rs b/crates/worker/src/feature/builtin/memory_lifecycle.rs index 30c14c91..4d8db2f4 100644 --- a/crates/worker/src/feature/builtin/memory_lifecycle.rs +++ b/crates/worker/src/feature/builtin/memory_lifecycle.rs @@ -20,7 +20,8 @@ use crate::feature::builtin::memory_staging_output::{ }; use crate::feature::builtin::session_explore::{SessionExploreFeature, SessionExploreState}; use crate::feature::session::{ - CommittedSessionCapture, CommittedSessionCaptureHandle, SessionExtensionHandle, + CommittedRunExit, CommittedSessionCapture, CommittedSessionCaptureHandle, + SessionExtensionHandle, }; use crate::feature::{ BackgroundTaskDeclaration, FeatureDescriptor, FeatureInstallContext, FeatureInstallError, @@ -37,7 +38,7 @@ use agen::token_counter::total_tokens_at; use manifest::WorkerManifest; use protocol::Event; -const TASK_NAME: &str = "memory-extraction"; +const TASK_NAME: &str = "memory-lifecycle"; const TASK_TIMEOUT: Duration = Duration::from_secs(300); /// Parent-Worker lifecycle Feature that observes committed runs and schedules @@ -45,12 +46,12 @@ const TASK_TIMEOUT: Duration = Duration::from_secs(300); /// Internal Worker, and staging disposition; Worker core owns only generic /// hook/task/session plumbing. #[derive(Clone)] -pub(crate) struct MemoryExtractionLifecycleFeature { - task: MemoryExtractionTask, +pub(crate) struct MemoryLifecycleFeature { + task: MemoryLifecycleTask, } #[derive(Clone)] -struct MemoryExtractionTask { +struct MemoryLifecycleTask { config: manifest::MemoryConfig, capture: CommittedSessionCaptureHandle, extensions: SessionExtensionHandle, @@ -62,7 +63,7 @@ struct MemoryExtractionTask { event_tx: Option>, } -impl MemoryExtractionLifecycleFeature { +impl MemoryLifecycleFeature { #[allow(clippy::too_many_arguments)] pub(crate) fn new( config: manifest::MemoryConfig, @@ -76,7 +77,7 @@ impl MemoryExtractionLifecycleFeature { event_tx: Option>, ) -> Self { Self { - task: MemoryExtractionTask { + task: MemoryLifecycleTask { config, capture, extensions, @@ -91,46 +92,133 @@ impl MemoryExtractionLifecycleFeature { } } -impl FeatureModule for MemoryExtractionLifecycleFeature { +impl FeatureModule for MemoryLifecycleFeature { fn descriptor(&self) -> FeatureDescriptor { - FeatureDescriptor::builtin("memory-extraction-lifecycle", "Memory Extraction Lifecycle") + FeatureDescriptor::builtin("memory-lifecycle", "Memory Lifecycle") .with_description( - "Observes terminal committed runs and schedules bounded Memory extraction.", + "Observes terminal committed runs and schedules bounded Memory extraction and Backend consolidation requests.", ) .with_background_task(BackgroundTaskDeclaration::worker_managed( TASK_NAME, - "Extract provenance-preserving Memory candidates after committed runs.", + "Extract provenance-preserving Memory candidates and request Backend consolidation after committed runs.", )) } fn install(&self, context: &mut FeatureInstallContext<'_>) -> Result<(), FeatureInstallError> { context .background_tasks() - .register(memory_extraction_task_spec(), self.task.clone()) + .register(memory_lifecycle_task_spec(), self.task.clone()) } } -fn memory_extraction_task_spec() -> BackgroundTaskSpec { +fn memory_lifecycle_task_spec() -> BackgroundTaskSpec { let declaration = BackgroundTaskDeclaration::worker_managed( TASK_NAME, - "Extract provenance-preserving Memory candidates after committed runs.", + "Extract provenance-preserving Memory candidates and request Backend consolidation after committed runs.", ); let mut spec = BackgroundTaskSpec::single_flight(declaration, TASK_TIMEOUT); spec.trigger = BackgroundTaskTrigger::RunCommitted; spec } -#[async_trait] -impl FeatureBackgroundTask for MemoryExtractionTask { - async fn run( +impl MemoryLifecycleTask { + async fn run_extraction( &self, context: BackgroundTaskContext, cancellation: BackgroundTaskCancellation, ) -> Result<(), HookError> { context.generation_fence.ensure_current()?; - let capture = self.capture.capture().map_err(hook_internal)?; - let pointer = extract_pointer(&capture)?; - if !extraction_threshold_reached(&capture, pointer.as_ref(), &self.config) { + let audit = WorkerAuditBase::new( + memory::audit::AuditWorker::MemoryExtract, + memory::audit::AuditTrigger::TokenThreshold, + self.config + .extract_model + .as_ref() + .or(Some(&self.manifest.model)) + .map(model_audit_from_manifest), + ) + .with_memory_settings(&self.config); + let capture = match self.capture.capture() { + Ok(capture) => capture, + Err(error) => { + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Failed, + format!("committed_session_capture_failed: {error}"), + None, + None, + None, + ) + .await; + return Ok(()); + } + }; + if !extraction_run_eligible(capture.run_exit) { + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Skipped, + format!("parent_run_not_finished: {:?}", capture.run_exit), + None, + None, + None, + ) + .await; + return Ok(()); + } + let pointer = match extract_pointer(&capture) { + Ok(pointer) => pointer, + Err(error) => { + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Failed, + format!("extract_pointer_invalid: {error}"), + None, + None, + None, + ) + .await; + return Ok(()); + } + }; + let Some(threshold) = self + .config + .extract_threshold + .filter(|threshold| *threshold > 0) + else { + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Skipped, + "token_threshold_disabled", + None, + None, + None, + ) + .await; + return Ok(()); + }; + let tokens_since = tokens_since_pointer(&capture, pointer.as_ref()); + if tokens_since < threshold { + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Skipped, + format!( + "token_threshold_not_reached tokens_since={tokens_since} threshold={threshold}" + ), + None, + None, + None, + ) + .await; return Ok(()); } @@ -141,6 +229,22 @@ impl FeatureBackgroundTask for MemoryExtractionTask { .min(capture.history.len()); let history_end = capture.history.len(); if history_start >= history_end || capture.entry_count == 0 { + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Skipped, + "no_new_committed_session_entries", + None, + Some(memory::audit::ExtractAudit { + session_id: Some(capture.session_id.clone()), + segment_id: Some(capture.segment_id.clone()), + history_range: Some([history_start as u64, history_end as u64]), + ..Default::default() + }), + None, + ) + .await; return Ok(()); } let view = SessionCapture::from_history_entries( @@ -155,16 +259,6 @@ impl FeatureBackgroundTask for MemoryExtractionTask { segment_id: capture.segment_id.clone(), range: [start_entry as u64, (capture.entry_count - 1) as u64], }; - let audit = WorkerAuditBase::new( - memory::audit::AuditWorker::MemoryExtract, - memory::audit::AuditTrigger::TokenThreshold, - self.config - .extract_model - .as_ref() - .or(Some(&self.manifest.model)) - .map(model_audit_from_manifest), - ) - .with_memory_settings(&self.config); let extract_audit_base = memory::audit::ExtractAudit { session_id: Some(capture.session_id.clone()), segment_id: Some(capture.segment_id.clone()), @@ -376,7 +470,86 @@ impl FeatureBackgroundTask for MemoryExtractionTask { } } -impl MemoryExtractionTask { +#[async_trait] +impl FeatureBackgroundTask for MemoryLifecycleTask { + async fn run( + &self, + context: BackgroundTaskContext, + cancellation: BackgroundTaskCancellation, + ) -> Result<(), HookError> { + let extraction = self + .run_extraction(context.clone(), cancellation.clone()) + .await; + if !cancellation.is_cancelled() { + context.generation_fence.ensure_current()?; + self.request_consolidation().await; + } + extraction + } +} + +impl MemoryLifecycleTask { + async fn request_consolidation(&self) { + let audit = WorkerAuditBase::new( + memory::audit::AuditWorker::MemoryConsolidation, + memory::audit::AuditTrigger::StagingBacklog, + self.config + .consolidation_model + .as_ref() + .or(Some(&self.manifest.model)) + .map(model_audit_from_manifest), + ) + .with_memory_settings(&self.config); + let Some((threshold_files, threshold_bytes)) = consolidation_thresholds(&self.config) + else { + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Skipped, + "consolidation_threshold_disabled", + None, + None, + None, + ) + .await; + return; + }; + match self + .workspace_client + .request_memory_staging_consolidation( + memory::backend::MemoryConsolidateStagingOperation { + force: false, + threshold_files, + threshold_bytes, + }, + ) + .await + { + Ok(output) => { + tracing::debug!( + status = output.status.as_str(), + summary = output.summary.as_str(), + "requested Backend Memory staging consolidation" + ); + } + Err(error) => { + tracing::warn!(%error, "request Backend Memory staging consolidation failed"); + audit + .emit( + self.workspace_client.as_ref(), + self.event_tx.as_ref(), + memory::audit::WorkerLifecycleStatus::Skipped, + "consolidation_backend_operation_failed", + None, + None, + None, + ) + .await; + } + } + } + async fn record_preparation_failure( &self, audit: &WorkerAuditBase, @@ -473,14 +646,30 @@ fn extract_pointer( Ok(pointer) } -fn extraction_threshold_reached( +fn consolidation_thresholds( + config: &manifest::MemoryConfig, +) -> Option<(Option, Option)> { + let threshold_files = config + .consolidation_threshold_files + .filter(|threshold| *threshold > 0); + let threshold_bytes = config + .consolidation_threshold_bytes + .filter(|threshold| *threshold > 0); + if threshold_files.is_none() && threshold_bytes.is_none() { + None + } else { + Some((threshold_files, threshold_bytes)) + } +} + +fn extraction_run_eligible(exit: CommittedRunExit) -> bool { + exit == CommittedRunExit::Finished +} + +fn tokens_since_pointer( capture: &CommittedSessionCapture, pointer: Option<&memory::ExtractPointerPayload>, - config: &manifest::MemoryConfig, -) -> bool { - if capture.history.is_empty() { - return false; - } +) -> u64 { let history_pointer = pointer .map(|pointer| pointer.processed_through_history_len) .unwrap_or(0) @@ -492,10 +681,22 @@ fn extraction_threshold_reached( .collect::>(); let current = total_tokens_at(&items, &capture.usage_history, capture.history.len()).tokens; let baseline = total_tokens_at(&items, &capture.usage_history, history_pointer).tokens; + current.saturating_sub(baseline) +} + +#[cfg(test)] +fn extraction_threshold_reached( + capture: &CommittedSessionCapture, + pointer: Option<&memory::ExtractPointerPayload>, + config: &manifest::MemoryConfig, +) -> bool { + if capture.history.is_empty() { + return false; + } let Some(threshold) = config.extract_threshold.filter(|threshold| *threshold > 0) else { return false; }; - current.saturating_sub(baseline) >= threshold + tokens_since_pointer(capture, pointer) >= threshold } #[derive(Clone)] @@ -620,6 +821,7 @@ mod tests { segment_id: "segment-1".to_string(), session_revision: history_len.try_into().unwrap(), entry_count: history_len, + run_exit: CommittedRunExit::Finished, history: (0..history_len) .map(|index| HistoryEntry { item: Item::user_message(format!("message-{index}")), @@ -637,6 +839,24 @@ mod tests { } } + #[test] + fn consolidation_thresholds_enable_backend_request_on_either_limit() { + let mut config = manifest::MemoryConfig::default(); + assert_eq!(consolidation_thresholds(&config), None); + config.consolidation_threshold_files = Some(3); + assert_eq!(consolidation_thresholds(&config), Some((Some(3), None))); + config.consolidation_threshold_files = None; + config.consolidation_threshold_bytes = Some(4096); + assert_eq!(consolidation_thresholds(&config), Some((None, Some(4096)))); + } + + #[test] + fn interrupted_parent_run_is_not_extraction_eligible() { + assert!(extraction_run_eligible(CommittedRunExit::Finished)); + assert!(!extraction_run_eligible(CommittedRunExit::NonFinal)); + assert!(!extraction_run_eligible(CommittedRunExit::Interrupted)); + } + #[test] fn normal_and_empty_extraction_require_explicit_finish() { let result = internal_result(WorkerRunResult::Finished); @@ -672,7 +892,7 @@ mod tests { #[test] fn task_scope_cancels_and_joins_before_rewrite_and_shutdown() { - let spec = memory_extraction_task_spec(); + let spec = memory_lifecycle_task_spec(); assert_eq!(spec.trigger, BackgroundTaskTrigger::RunCommitted); assert_eq!(spec.max_concurrency, 1); assert_eq!(spec.rewrite, BackgroundTaskRewritePolicy::CancelAndWait); @@ -759,8 +979,10 @@ mod tests { ); } let controller_source = include_str!("../../controller.rs"); - assert!(controller_source.contains("if feature_config.memory.enabled")); - assert!(controller_source.contains("MemoryExtractionLifecycleFeature::new")); + assert!(controller_source.contains("if let Some(config) = memory_config.clone()")); + assert!(controller_source.contains("MemoryLifecycleFeature::new")); + let lifecycle_source = include_str!("memory_lifecycle.rs"); + assert!(lifecycle_source.contains("request_memory_staging_consolidation")); let internal_worker_source = include_str!("../../internal_worker.rs"); assert!(!internal_worker_source.contains("manifest.memory = None")); } diff --git a/crates/worker/src/feature/session.rs b/crates/worker/src/feature/session.rs index 97345b10..0729f0af 100644 --- a/crates/worker/src/feature/session.rs +++ b/crates/worker/src/feature/session.rs @@ -5,6 +5,13 @@ use serde_json::Value; use crate::session_history::SessionHistoryMetadata; +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) enum CommittedRunExit { + Finished, + NonFinal, + Interrupted, +} + /// Immutable projection of one durably committed session-log location. /// /// Feature code receives this value only after the host has committed the @@ -18,6 +25,7 @@ pub(crate) struct CommittedSessionCapture { /// Monotonic committed-log revision for the captured Segment. pub(crate) session_revision: u64, pub(crate) entry_count: usize, + pub(crate) run_exit: CommittedRunExit, pub(crate) history: Vec>, pub(crate) usage_history: Vec, pub(crate) extensions: Vec<(String, Value)>, diff --git a/crates/worker/src/worker.rs b/crates/worker/src/worker.rs index c5fcb23b..4b74f647 100644 --- a/crates/worker/src/worker.rs +++ b/crates/worker/src/worker.rs @@ -44,7 +44,7 @@ use crate::feature::background::{BackgroundTaskRewriteGuard, FeatureBackgroundTa use crate::feature::builtin::memory::WorkspaceMemoryBackendError; use crate::feature::builtin::{TaskFeature, WorkerObservationProvider}; use crate::feature::session::{ - CommittedSessionCapture, CommittedSessionCaptureHandle, FeatureSessionError, + CommittedRunExit, CommittedSessionCapture, CommittedSessionCaptureHandle, FeatureSessionError, SessionExtensionHandle, }; use crate::feature::{ @@ -2108,6 +2108,19 @@ impl Worker { let entries = store .read_all(location.session_id, location.segment_id) .map_err(|error| FeatureSessionError::Capture(error.to_string()))?; + let run_exit = entries + .iter() + .rev() + .find_map(|entry| match entry { + LogEntry::RunCompleted { + result: EngineResult::Finished, + .. + } => Some(CommittedRunExit::Finished), + LogEntry::RunCompleted { .. } => Some(CommittedRunExit::NonFinal), + LogEntry::RunErrored { .. } => Some(CommittedRunExit::Interrupted), + _ => None, + }) + .unwrap_or(CommittedRunExit::NonFinal); let restored = segment_log::collect_state(&entries); let history = restore_history_entries(location.session_id, location.segment_id, &entries) @@ -2117,6 +2130,7 @@ impl Worker { segment_id: location.segment_id.to_string(), session_revision: entries.len().try_into().unwrap_or(u64::MAX), entry_count: entries.len(), + run_exit, history, usage_history: restored.usage_history, extensions: restored.extensions.into_iter().collect(), @@ -8546,6 +8560,30 @@ mod build_summary_prompt_tests { .count(), 1 ); + + worker + .commit_entry(LogEntry::RunErrored { + ts: segment_log::now_millis(), + interrupted: true, + message: "cancelled".to_string(), + }) + .unwrap(); + assert_eq!( + capture_handle.capture().unwrap().run_exit, + CommittedRunExit::Interrupted + ); + worker + .commit_entry(LogEntry::RunCompleted { + ts: segment_log::now_millis(), + interrupted: false, + result: EngineResult::Finished, + active_run_turn_count: None, + }) + .unwrap(); + assert_eq!( + capture_handle.capture().unwrap().run_exit, + CommittedRunExit::Finished + ); } fn minimal_manifest() -> WorkerManifest {