diff --git a/crates/protocol/src/lib.rs b/crates/protocol/src/lib.rs index bfc1a49a..339d1d67 100644 --- a/crates/protocol/src/lib.rs +++ b/crates/protocol/src/lib.rs @@ -340,8 +340,7 @@ pub struct InternalWorkerRef { pub struct InternalWorkerSnapshot { pub worker: InternalWorkerRef, pub revision: u64, - #[cfg_attr(feature = "typescript", ts(type = "Array"))] - pub entries: Vec, + pub session: SessionSnapshot, #[serde(default)] pub status: WorkerStatus, #[serde(default, skip_serializing_if = "Option::is_none")] @@ -364,6 +363,106 @@ pub enum ToolResultDisposition { OutcomeUnknown, } +/// Canonical, storage-independent projection of committed session history. +/// +/// Worker protocols expose this DTO instead of append-log records. New +/// storage variants can therefore be added without teaching every client how +/// to replay the durable log format. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +pub struct SessionSnapshot { + pub entries: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(rename_all = "snake_case")] +pub enum SessionEntryProvenance { + HumanInput, + WorkerInput, + FlowInstruction, + BackendInstruction, + ModelOutput, + ToolOutput, + DerivedSummary, + LegacyUnknown, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +pub struct SessionSnapshotEntry { + /// Stable identity from durable history metadata, or a deterministic + /// identity derived from the legacy segment and log position. + pub entry_id: String, + pub provenance: SessionEntryProvenance, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub derived_from: Vec, + #[serde(flatten)] + pub data: SessionSnapshotEntryData, +} + +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum SessionSnapshotEntryData { + UserInput { + segments: Vec, + }, + Message { + role: SessionMessageRole, + content: Vec, + }, + ToolCall { + call_id: String, + name: String, + arguments: String, + }, + ToolResult { + call_id: String, + summary: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + content: Option, + is_error: bool, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + attachments: Vec, + }, + SystemItem { + item_kind: String, + content: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + #[cfg_attr(feature = "typescript", ts(type = "unknown"))] + data: Option, + }, + RunError { + message: String, + }, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(rename_all = "snake_case")] +pub enum SessionMessageRole { + User, + Assistant, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum SessionContentPart { + Text { text: String }, + Refusal { refusal: String }, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[cfg_attr(feature = "typescript", derive(ts_rs::TS))] +pub struct SessionToolAttachment { + pub media_type: String, + /// Base64-encoded durable attachment body. Public snapshots preserve the + /// committed multimodal value instead of replacing it with placeholder text. + pub data_base64: String, +} + #[derive(Debug, Clone, Serialize, Deserialize)] #[cfg_attr(feature = "typescript", derive(ts_rs::TS))] #[serde(tag = "event", content = "data", rename_all = "snake_case")] @@ -555,8 +654,7 @@ pub enum Event { /// role-specific entry events (`SegmentRotated` / `SystemItem`) — /// there is no generic "every committed entry" broadcast. Snapshot { - #[cfg_attr(feature = "typescript", ts(type = "Array"))] - entries: Vec, + session: SessionSnapshot, greeting: Greeting, #[serde(default)] status: WorkerStatus, @@ -589,14 +687,10 @@ pub enum Event { /// Server-side segment log rotated to a fresh `SegmentStart`. /// /// Fires on compaction and on auto-fork when the store head drifts - /// from the live writer's cached head. Clients drop their derived - /// view and reseed from `entry.history` exactly the way they would - /// from a connect-time `Snapshot`. - /// - /// Payload is the JSON form of `session_store::LogEntry::SegmentStart`. + /// A compaction/fork has replaced the authoritative segment. Clients drop + /// their derived view and reseed from the canonical committed snapshot. SegmentRotated { - #[cfg_attr(feature = "typescript", ts(type = "unknown"))] - entry: serde_json::Value, + session: SessionSnapshot, }, /// Current Worker controller status. Broadcast on every controller-level /// transition and included in `History` snapshots for late attach. @@ -623,11 +717,10 @@ pub enum Event { head_entries: usize, targets: Vec, }, - /// A rewind has truncated the authoritative session. `entries` is the - /// retained session-log prefix clients should use to reseed display state. + /// A rewind has truncated the authoritative session. `session` is the + /// retained canonical snapshot clients should use to reseed display state. RewindApplied { - #[cfg_attr(feature = "typescript", ts(type = "Array"))] - entries: Vec, + session: SessionSnapshot, input: Vec, summary: RewindSummary, }, @@ -1440,7 +1533,16 @@ mod tests { #[test] fn event_snapshot_format() { let event = Event::Snapshot { - entries: vec![serde_json::json!({"kind": "user_input", "ts": 1, "segments": []})], + session: SessionSnapshot { + entries: vec![SessionSnapshotEntry { + entry_id: "entry-1".into(), + provenance: SessionEntryProvenance::HumanInput, + derived_from: Vec::new(), + data: SessionSnapshotEntryData::UserInput { + segments: Vec::new(), + }, + }], + }, greeting: Greeting { worker_name: "test".into(), cwd: "/tmp".into(), @@ -1458,8 +1560,11 @@ mod tests { let json = serde_json::to_string(&event).unwrap(); let parsed: serde_json::Value = serde_json::from_str(&json).unwrap(); assert_eq!(parsed["event"], "snapshot"); - assert!(parsed["data"]["entries"].is_array()); - assert_eq!(parsed["data"]["entries"][0]["kind"], "user_input"); + assert!(parsed["data"]["session"]["entries"].is_array()); + assert_eq!( + parsed["data"]["session"]["entries"][0]["kind"], + "user_input" + ); assert_eq!(parsed["data"]["greeting"]["worker_name"], "test"); assert_eq!(parsed["data"]["greeting"]["tools"][0], "Read"); assert_eq!(parsed["data"]["greeting"]["context_window"], 200_000); @@ -1469,7 +1574,7 @@ mod tests { #[test] fn event_snapshot_in_flight_roundtrip_and_default() { - let inbound = r#"{"event":"snapshot","data":{"entries":[],"greeting":{"worker_name":"test","cwd":"/tmp","provider":"p","model":"m","scope_summary":"s","tools":[]},"status":"running"}}"#; + let inbound = r#"{"event":"snapshot","data":{"session":{"entries":[]},"greeting":{"worker_name":"test","cwd":"/tmp","provider":"p","model":"m","scope_summary":"s","tools":[]},"status":"running"}}"#; let decoded: Event = serde_json::from_str(inbound).unwrap(); match decoded { Event::Snapshot { in_flight, .. } => assert!(in_flight.is_empty()), @@ -1477,7 +1582,9 @@ mod tests { } let event = Event::Snapshot { - entries: Vec::new(), + session: SessionSnapshot { + entries: Vec::new(), + }, greeting: Greeting { worker_name: "test".into(), cwd: "/tmp".into(), @@ -1543,15 +1650,17 @@ mod tests { #[test] fn event_segment_rotated_roundtrip() { let event = Event::SegmentRotated { - entry: serde_json::json!({"kind": "segment_start", "ts": 1, "history": []}), + session: SessionSnapshot { + entries: Vec::new(), + }, }; let json = serde_json::to_string(&event).unwrap(); let parsed: serde_json::Value = serde_json::from_str(&json).unwrap(); assert_eq!(parsed["event"], "segment_rotated"); - assert_eq!(parsed["data"]["entry"]["kind"], "segment_start"); + assert!(parsed["data"]["session"]["entries"].is_array()); let decoded: Event = serde_json::from_str(&json).unwrap(); match decoded { - Event::SegmentRotated { entry } => assert_eq!(entry["kind"], "segment_start"), + Event::SegmentRotated { session } => assert!(session.entries.is_empty()), other => panic!("expected SegmentRotated, got {other:?}"), } } @@ -1627,8 +1736,8 @@ mod tests { } #[test] - fn event_snapshot_legacy_without_status_defaults_to_idle() { - let json = r#"{"event":"snapshot","data":{"entries":[],"greeting":{"worker_name":"test","cwd":"/tmp","provider":"anthropic","model":"claude","scope_summary":"","tools":[]}}}"#; + fn event_snapshot_without_status_defaults_to_idle() { + let json = r#"{"event":"snapshot","data":{"session":{"entries":[]},"greeting":{"worker_name":"test","cwd":"/tmp","provider":"anthropic","model":"claude","scope_summary":"","tools":[]}}}"#; let decoded: Event = serde_json::from_str(json).unwrap(); match decoded { Event::Snapshot { @@ -2039,11 +2148,11 @@ mod tests { } #[test] - fn legacy_snapshot_defaults_internal_workers_to_empty() { + fn snapshot_defaults_internal_workers_to_empty() { let snapshot: Event = serde_json::from_value(serde_json::json!({ "event": "snapshot", "data": { - "entries": [], + "session": { "entries": [] }, "greeting": { "worker_name": "parent", "cwd": ".", diff --git a/crates/protocol/src/typescript.rs b/crates/protocol/src/typescript.rs index 2ac5446c..bf75e324 100644 --- a/crates/protocol/src/typescript.rs +++ b/crates/protocol/src/typescript.rs @@ -8,7 +8,9 @@ use crate::{ CompletionKind, ErrorCode, Event, Greeting, InFlightBlock, InFlightSnapshot, InFlightToolCallState, InternalWorkerKind, InternalWorkerRef, InternalWorkerSnapshot, InvokeKind, MemoryWorkerEvent, Method, Permission, RewindSummary, RewindTarget, RewindTargetId, - RunResult, ScopeRule, Segment, ToolResultDisposition, TurnResult, WorkerEvent, WorkerStatus, + RunResult, ScopeRule, Segment, SessionContentPart, SessionEntryProvenance, SessionMessageRole, + SessionSnapshot, SessionSnapshotEntry, SessionSnapshotEntryData, SessionToolAttachment, + ToolResultDisposition, TurnResult, WorkerEvent, WorkerStatus, subscription::{ EventSubscriptionSelector, SubscriptionEvent, SubscriptionEventPayload, SubscriptionFrame, SubscriptionFramePayload, SubscriptionId, SubscriptionRejectionCode, SubscriptionRequest, @@ -63,6 +65,13 @@ pub fn generated_protocol_types() -> String { push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); + push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); push_decl::(&cfg, &mut output); diff --git a/crates/session-store/src/history.rs b/crates/session-store/src/history.rs index 3903213c..9558f343 100644 --- a/crates/session-store/src/history.rs +++ b/crates/session-store/src/history.rs @@ -1,5 +1,6 @@ //! Serializable history entries with restore-authoritative logical identity and origin. +use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD}; use serde::{Deserialize, Serialize}; use crate::{LoggedItem, SessionId}; @@ -175,6 +176,25 @@ pub fn legacy_segment_history( session_id: SessionId, items: impl IntoIterator, ) -> Vec { - let _ = session_id; - items.into_iter().map(legacy_logged_history).collect() + items + .into_iter() + .enumerate() + .map(|(index, item)| LoggedHistoryEntry { + item, + metadata: LoggedSessionHistoryMetadata { + // Legacy logs have no persisted entry id. Derive one solely from + // durable segment content rather than minting a new random value + // on every restore/read. The explicit LegacyUnknown origin keeps + // this compatibility identity from becoming trust authority. + entry_id: { + let mut identity = Vec::with_capacity(24); + identity.extend_from_slice(session_id.as_bytes()); + identity.extend_from_slice(&(index as u64).to_be_bytes()); + LoggedSessionHistoryEntryId(format!("l-{}", URL_SAFE_NO_PAD.encode(identity))) + }, + origin: LoggedSessionHistoryOrigin::LegacyUnknown, + derivation: None, + }, + }) + .collect() } diff --git a/crates/session-store/src/lib.rs b/crates/session-store/src/lib.rs index 9fe1fad7..07bc0bc9 100644 --- a/crates/session-store/src/lib.rs +++ b/crates/session-store/src/lib.rs @@ -34,6 +34,7 @@ pub mod event_trace; pub mod fs_store; pub mod history; pub mod logged_item; +pub mod public_snapshot; pub mod segment; pub mod segment_log; pub mod store; diff --git a/crates/session-store/src/public_snapshot.rs b/crates/session-store/src/public_snapshot.rs new file mode 100644 index 00000000..d25feb99 --- /dev/null +++ b/crates/session-store/src/public_snapshot.rs @@ -0,0 +1,402 @@ +use base64::{ + Engine as _, + engine::general_purpose::{STANDARD as BASE64, URL_SAFE_NO_PAD}, +}; +use protocol::{ + Segment, SessionContentPart, SessionEntryProvenance, SessionMessageRole, SessionSnapshot, + SessionSnapshotEntry, SessionSnapshotEntryData, SessionToolAttachment, +}; + +use crate::{ + LogEntry, LoggedContentPart, LoggedHistoryEntry, LoggedItem, LoggedRole, + LoggedSessionHistoryOrigin, SessionId, SystemItem, +}; + +/// Project a complete current-segment log. A valid segment always starts with +/// one of the two SegmentStart records; malformed partial input uses the nil +/// session only to keep the public failure projection deterministic. +pub fn project_current_session_snapshot(log: &[LogEntry]) -> SessionSnapshot { + let session_id = log.iter().find_map(|entry| match entry { + LogEntry::SegmentStart { session_id, .. } + | LogEntry::AnnotatedSegmentStart { session_id, .. } => Some(*session_id), + _ => None, + }); + project_session_snapshot(session_id.unwrap_or_else(SessionId::nil), log) +} + +/// Project the current durable segment into the only public session-history +/// representation. Append-log records remain an internal persistence format. +pub fn project_session_snapshot(session_id: SessionId, log: &[LogEntry]) -> SessionSnapshot { + let mut session_key = session_id; + let mut entries = Vec::new(); + + for (log_index, record) in log.iter().enumerate() { + match record { + LogEntry::SegmentStart { + session_id, + history, + .. + } => { + session_key = *session_id; + entries.clear(); + for (item_index, item) in history.iter().enumerate() { + if let Some(data) = project_item(item) { + entries.push(legacy_entry(&session_key, log_index, item_index, data)); + } + } + } + LogEntry::AnnotatedSegmentStart { + session_id, + history, + .. + } => { + session_key = *session_id; + entries.clear(); + extend_history(&mut entries, history, None); + } + LogEntry::UserInput { segments, .. } => entries.push(legacy_entry( + &session_key, + log_index, + 0, + SessionSnapshotEntryData::UserInput { + segments: segments.clone(), + }, + )), + LogEntry::AnnotatedUserInput { + segments, history, .. + } => extend_history(&mut entries, history, Some(segments)), + LogEntry::AssistantItem { item, .. } | LogEntry::ToolResult { item, .. } => { + if let Some(data) = project_item(item) { + entries.push(legacy_entry(&session_key, log_index, 0, data)); + } + } + LogEntry::AnnotatedAssistantItem { entry, .. } + | LogEntry::AnnotatedToolResult { entry, .. } => { + if let Some(data) = project_item(&entry.item) { + entries.push(history_entry(entry, data)); + } + } + LogEntry::SystemItem { item, .. } => entries.push(system_entry( + item, + legacy_entry_id(&session_key, log_index, 0), + SessionEntryProvenance::LegacyUnknown, + Vec::new(), + )), + LogEntry::AnnotatedSystemItem { entry, .. } => entries.push(system_entry( + &entry.item, + entry.metadata.entry_id.0.clone(), + provenance(&entry.metadata.origin), + derivation_ids(entry), + )), + LogEntry::RunErrored { message, .. } => entries.push(legacy_entry( + &session_key, + log_index, + 0, + SessionSnapshotEntryData::RunError { + message: message.clone(), + }, + )), + // Run checkpoints, configuration, usage, and extension state are + // controller/storage authority rather than committed conversation. + LogEntry::Invoke { .. } + | LogEntry::TurnEnd { .. } + | LogEntry::RunCompleted { .. } + | LogEntry::ActiveRunCheckpoint { .. } + | LogEntry::PausedTurnAbandoned { .. } + | LogEntry::ConfigChanged { .. } + | LogEntry::LlmUsage { .. } + | LogEntry::Extension { .. } => {} + } + } + + SessionSnapshot { entries } +} + +fn extend_history( + output: &mut Vec, + history: &[LoggedHistoryEntry], + input_segments: Option<&Vec>, +) { + let mut attached_segments = false; + for entry in history { + let data = if !attached_segments + && matches!( + entry.metadata.origin, + LoggedSessionHistoryOrigin::HumanInput { .. } + ) + && input_segments.is_some() + { + attached_segments = true; + SessionSnapshotEntryData::UserInput { + segments: input_segments.cloned().unwrap_or_default(), + } + } else { + let Some(data) = project_item(&entry.item) else { + continue; + }; + data + }; + output.push(history_entry(entry, data)); + } +} + +fn history_entry( + entry: &LoggedHistoryEntry, + data: SessionSnapshotEntryData, +) -> SessionSnapshotEntry { + SessionSnapshotEntry { + entry_id: entry.metadata.entry_id.0.clone(), + provenance: provenance(&entry.metadata.origin), + derived_from: entry + .metadata + .derivation + .as_ref() + .map(|derivation| { + derivation + .sources + .iter() + .map(|source| source.0.clone()) + .collect() + }) + .unwrap_or_default(), + data, + } +} + +fn derivation_ids(entry: &crate::LoggedSystemHistoryEntry) -> Vec { + entry + .metadata + .derivation + .as_ref() + .map(|derivation| { + derivation + .sources + .iter() + .map(|source| source.0.clone()) + .collect() + }) + .unwrap_or_default() +} + +fn legacy_entry( + session_key: &SessionId, + log_index: usize, + item_index: usize, + data: SessionSnapshotEntryData, +) -> SessionSnapshotEntry { + SessionSnapshotEntry { + entry_id: legacy_entry_id(session_key, log_index, item_index), + provenance: SessionEntryProvenance::LegacyUnknown, + derived_from: Vec::new(), + data, + } +} + +fn legacy_entry_id(session_key: &SessionId, log_index: usize, item_index: usize) -> String { + let mut identity = Vec::with_capacity(32); + identity.extend_from_slice(session_key.as_bytes()); + identity.extend_from_slice(&(log_index as u64).to_be_bytes()); + identity.extend_from_slice(&(item_index as u64).to_be_bytes()); + format!("l-{}", URL_SAFE_NO_PAD.encode(identity)) +} + +fn provenance(origin: &LoggedSessionHistoryOrigin) -> SessionEntryProvenance { + match origin { + LoggedSessionHistoryOrigin::HumanInput { .. } => SessionEntryProvenance::HumanInput, + LoggedSessionHistoryOrigin::WorkerInput { .. } => SessionEntryProvenance::WorkerInput, + LoggedSessionHistoryOrigin::FlowInstruction { .. } => { + SessionEntryProvenance::FlowInstruction + } + LoggedSessionHistoryOrigin::BackendInstruction { .. } => { + SessionEntryProvenance::BackendInstruction + } + LoggedSessionHistoryOrigin::ModelOutput { .. } => SessionEntryProvenance::ModelOutput, + LoggedSessionHistoryOrigin::ToolOutput { .. } => SessionEntryProvenance::ToolOutput, + LoggedSessionHistoryOrigin::DerivedSummary => SessionEntryProvenance::DerivedSummary, + LoggedSessionHistoryOrigin::LegacyUnknown => SessionEntryProvenance::LegacyUnknown, + } +} + +fn project_item(item: &LoggedItem) -> Option { + match item { + LoggedItem::Message { role, content } => { + let role = match role { + LoggedRole::User => SessionMessageRole::User, + LoggedRole::Assistant => SessionMessageRole::Assistant, + // System prompts and instruction history never cross the public + // snapshot boundary. Typed SystemItems have separate records. + LoggedRole::System => return None, + }; + Some(SessionSnapshotEntryData::Message { + role, + content: content + .iter() + .map(|part| match part { + LoggedContentPart::Text { text } => { + SessionContentPart::Text { text: text.clone() } + } + LoggedContentPart::Refusal { refusal } => SessionContentPart::Refusal { + refusal: refusal.clone(), + }, + }) + .collect(), + }) + } + LoggedItem::ToolCall { + call_id, + name, + arguments, + } => Some(SessionSnapshotEntryData::ToolCall { + call_id: call_id.clone(), + name: name.clone(), + arguments: arguments.clone(), + }), + LoggedItem::ToolResult { + call_id, + summary, + content, + is_error, + attachments, + .. + } => Some(SessionSnapshotEntryData::ToolResult { + call_id: call_id.clone(), + summary: summary.clone(), + content: content.clone(), + is_error: *is_error, + attachments: attachments + .iter() + .map(|attachment| match attachment { + crate::logged_item::LoggedAttachment::Image { mime_type, data } => { + SessionToolAttachment { + media_type: mime_type.clone(), + data_base64: BASE64.encode(data), + } + } + }) + .collect(), + }), + // Hidden model reasoning is never observable. + LoggedItem::Reasoning { .. } => None, + } +} + +fn system_entry( + item: &SystemItem, + entry_id: String, + provenance: SessionEntryProvenance, + derived_from: Vec, +) -> SessionSnapshotEntry { + let mut data = serde_json::to_value(item).ok(); + if let Some(serde_json::Value::Object(object)) = data.as_mut() { + object.remove("prompt_provenance"); + } + let item_kind = data + .as_ref() + .and_then(|value| value.get("kind")) + .and_then(serde_json::Value::as_str) + .unwrap_or("system_item") + .to_owned(); + SessionSnapshotEntry { + entry_id, + provenance, + derived_from, + data: SessionSnapshotEntryData::SystemItem { + item_kind, + content: item.history_text(), + data, + }, + } +} + +#[cfg(test)] +mod tests { + use agen::llm_client::RequestConfig; + + use super::*; + use crate::{LoggedSessionHistoryEntryId, LoggedSessionHistoryMetadata, LoggedWorkerSubject}; + + #[test] + fn legacy_projection_is_stable_and_hides_reasoning_and_system_prompts() { + let session_id = crate::new_session_id(); + let log = vec![LogEntry::SegmentStart { + ts: 1, + session_id, + system_prompt: None, + config: RequestConfig::default(), + history: vec![ + LoggedItem::Message { + role: LoggedRole::System, + content: vec![LoggedContentPart::Text { + text: "secret prompt".into(), + }], + }, + LoggedItem::Reasoning { + text: "secret reasoning".into(), + summary: Vec::new(), + encrypted_content: None, + signature: None, + }, + LoggedItem::Message { + role: LoggedRole::Assistant, + content: vec![LoggedContentPart::Text { + text: "visible".into(), + }], + }, + ], + forked_from: None, + compacted_from: None, + }]; + + let first = project_session_snapshot(session_id, &log); + let second = project_session_snapshot(session_id, &log); + assert_eq!(first, second); + assert_eq!(first.entries.len(), 1); + assert_eq!( + first.entries[0].provenance, + SessionEntryProvenance::LegacyUnknown + ); + let json = serde_json::to_string(&first).unwrap(); + assert!(!json.contains("secret prompt")); + assert!(!json.contains("secret reasoning")); + assert!(json.contains("visible")); + } + + #[test] + fn annotated_projection_preserves_identity_and_provenance() { + let session_id = crate::new_session_id(); + let metadata = LoggedSessionHistoryMetadata { + entry_id: LoggedSessionHistoryEntryId::new(), + origin: LoggedSessionHistoryOrigin::ModelOutput { + worker: LoggedWorkerSubject { + workspace_id: None, + runtime_id: None, + worker_id: "worker".into(), + }, + }, + derivation: None, + }; + let expected_id = metadata.entry_id.0.clone(); + let log = vec![LogEntry::AnnotatedSegmentStart { + ts: 1, + session_id, + system_prompt: None, + config: RequestConfig::default(), + history: vec![LoggedHistoryEntry { + item: LoggedItem::Message { + role: LoggedRole::Assistant, + content: vec![LoggedContentPart::Text { text: "ok".into() }], + }, + metadata, + }], + forked_from: None, + compacted_from: None, + }]; + + let snapshot = project_session_snapshot(session_id, &log); + assert_eq!(snapshot.entries[0].entry_id, expected_id); + assert_eq!( + snapshot.entries[0].provenance, + SessionEntryProvenance::ModelOutput + ); + } +} diff --git a/crates/session-store/src/worker_session_store.rs b/crates/session-store/src/worker_session_store.rs index 5dd063d8..ddce9e1b 100644 --- a/crates/session-store/src/worker_session_store.rs +++ b/crates/session-store/src/worker_session_store.rs @@ -12,7 +12,11 @@ use crate::event_trace::TraceEntry; use crate::segment_log::LogEntry; use crate::store::{Store, StoreError}; -use crate::{SegmentId, SessionId}; +use crate::{ + LoggedHistoryEntry, LoggedItem, LoggedSessionHistoryEntryId, LoggedSessionHistoryMetadata, + LoggedSessionHistoryOrigin, LoggedSystemHistoryEntry, SegmentId, SessionId, +}; +use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD}; use serde::{Deserialize, Serialize}; use std::fs::{self, File, OpenOptions}; use std::io::{Read, Seek, SeekFrom, Write}; @@ -20,7 +24,8 @@ use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex}; use std::time::SystemTime; -const SESSION_SCHEMA_VERSION: u32 = 2; +const SESSION_SCHEMA_VERSION: u32 = 3; +const PREVIOUS_SESSION_SCHEMA_VERSION: u32 = 2; const LEGACY_SESSION_SCHEMA_VERSION: u32 = 1; const SESSION_FILE: &str = "session.json"; const SEGMENTS_DIR: &str = "segments"; @@ -47,9 +52,11 @@ impl WorkerSessionStore { Ok(bytes) => { let mut manifest: SessionManifest = serde_json::from_slice(&bytes)?; match manifest.schema_version { - SESSION_SCHEMA_VERSION => {} - LEGACY_SESSION_SCHEMA_VERSION => { - validate_legacy_segment_logs(&root)?; + SESSION_SCHEMA_VERSION => { + validate_canonical_segment_logs(&root)?; + } + PREVIOUS_SESSION_SCHEMA_VERSION | LEGACY_SESSION_SCHEMA_VERSION => { + migrate_segment_logs_to_v3(&root, manifest.session_id)?; manifest.schema_version = SESSION_SCHEMA_VERSION; atomic_write_json(&root.join(SESSION_FILE), &manifest)?; } @@ -144,6 +151,48 @@ impl WorkerSessionStore { .join(format!("{segment_id}.trace.jsonl")) } + fn append_log_entry( + &self, + path: &Path, + session_id: SessionId, + segment_id: SegmentId, + entry: &LogEntry, + ) -> Result<(), StoreError> { + let _guard = self + .append_lock + .lock() + .map_err(|_| std::io::Error::other("Worker Session append lock was poisoned"))?; + let mut file = OpenOptions::new() + .create(true) + .read(true) + .write(true) + .append(true) + .open(path)?; + let committed_len = truncate_uncommitted_tail(&mut file)?; + file.seek(SeekFrom::Start(0))?; + let mut existing = Vec::new(); + file.read_to_end(&mut existing)?; + let line_index = parse_jsonl::(&existing)?.len(); + let entry = canonicalize_log_entry(session_id, segment_id, line_index, entry.clone()); + let line = serde_json::to_string(&entry)?; + let mut record = Vec::with_capacity(line.len() + 1); + record.extend_from_slice(line.as_bytes()); + record.push(b'\n'); + if let Err(write_error) = file.write_all(&record) { + return match file.set_len(committed_len) { + Ok(()) => Err(write_error.into()), + Err(rollback_error) => Err(std::io::Error::new( + rollback_error.kind(), + format!( + "session append failed ({write_error}) and rollback failed: {rollback_error}" + ), + ) + .into()), + }; + } + Ok(()) + } + fn append_line(&self, path: &Path, line: &str) -> Result<(), StoreError> { let _guard = self .append_lock @@ -183,7 +232,7 @@ impl Store for WorkerSessionStore { entry: &LogEntry, ) -> Result<(), StoreError> { self.ensure_session(session_id, true)?; - self.append_line(&self.log_path(segment_id), &serde_json::to_string(entry)?) + self.append_log_entry(&self.log_path(segment_id), session_id, segment_id, entry) } fn read_all( @@ -236,8 +285,9 @@ impl Store for WorkerSessionStore { ) -> Result<(), StoreError> { self.ensure_session(session_id, true)?; let mut content = Vec::new(); - for entry in entries { - serde_json::to_writer(&mut content, entry)?; + for (line_index, entry) in entries.iter().enumerate() { + let entry = canonicalize_log_entry(session_id, segment_id, line_index, entry.clone()); + serde_json::to_writer(&mut content, &entry)?; content.push(b'\n'); } atomic_write_bytes(&self.log_path(segment_id), &content)?; @@ -286,37 +336,208 @@ impl Store for WorkerSessionStore { } } -fn validate_legacy_segment_logs(root: &Path) -> Result<(), StoreError> { +fn segment_log_paths(root: &Path) -> Result, StoreError> { let segments = root.join(SEGMENTS_DIR); if !segments.exists() { - return Ok(()); + return Ok(Vec::new()); } + let mut paths = Vec::new(); for entry in fs::read_dir(&segments)? { let entry = entry?; let path = entry.path(); + let metadata = fs::symlink_metadata(&path)?; let Some(name) = path.file_name().and_then(|name| name.to_str()) else { - continue; + return Err(StoreError::Corrupt { + line: 0, + message: format!("non-UTF-8 Worker Session segment path: {}", path.display()), + }); }; - if !name.ends_with(".jsonl") || name.ends_with(".trace.jsonl") { + if name.ends_with(".trace.jsonl") || name.starts_with('.') { continue; } - let contents = fs::read_to_string(&path)?; - for (line_index, line) in contents.lines().enumerate() { - if line.trim().is_empty() { - continue; - } - serde_json::from_str::(line).map_err(|error| StoreError::Corrupt { - line: line_index + 1, + if !name.ends_with(".jsonl") { + continue; + } + if !metadata.file_type().is_file() { + return Err(StoreError::Corrupt { + line: 0, message: format!( - "cannot migrate legacy Worker Session log {}: {error}", + "Worker Session segment is not a regular file: {}", path.display() ), - })?; + }); + } + let segment_id = + name.trim_end_matches(".jsonl") + .parse() + .map_err(|_| StoreError::Corrupt { + line: 0, + message: format!("invalid Worker Session segment name: {name}"), + })?; + paths.push((segment_id, path)); + } + paths.sort_by_key(|(segment_id, _)| *segment_id); + Ok(paths) +} + +fn migrate_segment_logs_to_v3(root: &Path, session_id: SessionId) -> Result<(), StoreError> { + for (segment_id, path) in segment_log_paths(root)? { + let source = fs::read(&path)?; + let entries: Vec = parse_jsonl(&source).map_err(|error| StoreError::Corrupt { + line: 0, + message: format!( + "cannot migrate Worker Session log {}: {error}", + path.display() + ), + })?; + let canonical = entries + .into_iter() + .enumerate() + .map(|(line_index, entry)| { + canonicalize_log_entry(session_id, segment_id, line_index, entry) + }) + .collect::>(); + validate_canonical_entries(&path, &canonical)?; + let mut output = Vec::new(); + for entry in canonical { + serde_json::to_writer(&mut output, &entry)?; + output.push(b'\n'); + } + + // Opening a Session is the exclusive restore boundary, but retain an + // unchanged-source fence so a racing writer cannot be silently lost. + if fs::read(&path)? != source { + return Err(StoreError::Corrupt { + line: 0, + message: format!( + "Worker Session segment changed during migration: {}", + path.display() + ), + }); + } + atomic_write_bytes(&path, &output)?; + } + Ok(()) +} + +fn validate_canonical_segment_logs(root: &Path) -> Result<(), StoreError> { + for (_, path) in segment_log_paths(root)? { + let entries: Vec = parse_jsonl(&fs::read(&path)?)?; + validate_canonical_entries(&path, &entries)?; + } + Ok(()) +} + +fn validate_canonical_entries(path: &Path, entries: &[LogEntry]) -> Result<(), StoreError> { + for (line_index, entry) in entries.iter().enumerate() { + if matches!( + entry, + LogEntry::SegmentStart { .. } + | LogEntry::UserInput { .. } + | LogEntry::AssistantItem { .. } + | LogEntry::ToolResult { .. } + | LogEntry::SystemItem { .. } + ) { + return Err(StoreError::Corrupt { + line: line_index + 1, + message: format!( + "Worker Session schema v3 contains legacy history record in {}", + path.display() + ), + }); } } Ok(()) } +fn legacy_metadata( + _session_id: SessionId, + segment_id: SegmentId, + line_index: usize, + item_index: usize, +) -> LoggedSessionHistoryMetadata { + let mut identity = Vec::with_capacity(32); + identity.extend_from_slice(segment_id.as_bytes()); + identity.extend_from_slice(&(line_index as u64).to_be_bytes()); + identity.extend_from_slice(&(item_index as u64).to_be_bytes()); + LoggedSessionHistoryMetadata { + entry_id: LoggedSessionHistoryEntryId(format!("l-{}", URL_SAFE_NO_PAD.encode(identity))), + origin: LoggedSessionHistoryOrigin::LegacyUnknown, + derivation: None, + } +} + +fn canonicalize_log_entry( + session_id: SessionId, + segment_id: SegmentId, + line_index: usize, + entry: LogEntry, +) -> LogEntry { + match entry { + LogEntry::SegmentStart { + ts, + session_id, + system_prompt, + config, + history, + forked_from, + compacted_from, + } => LogEntry::AnnotatedSegmentStart { + ts, + session_id, + system_prompt, + config, + history: history + .into_iter() + .enumerate() + .map(|(item_index, item)| LoggedHistoryEntry { + item, + metadata: legacy_metadata(session_id, segment_id, line_index, item_index), + }) + .collect(), + forked_from, + compacted_from, + }, + LogEntry::UserInput { + ts, + segments, + extensions, + } => LogEntry::AnnotatedUserInput { + ts, + history: vec![LoggedHistoryEntry { + item: LoggedItem::from(agen::Item::user_message( + protocol::Segment::flatten_to_text(&segments), + )), + metadata: legacy_metadata(session_id, segment_id, line_index, 0), + }], + segments, + extensions, + }, + LogEntry::AssistantItem { ts, item } => LogEntry::AnnotatedAssistantItem { + ts, + entry: LoggedHistoryEntry { + item, + metadata: legacy_metadata(session_id, segment_id, line_index, 0), + }, + }, + LogEntry::ToolResult { ts, item } => LogEntry::AnnotatedToolResult { + ts, + entry: LoggedHistoryEntry { + item, + metadata: legacy_metadata(session_id, segment_id, line_index, 0), + }, + }, + LogEntry::SystemItem { ts, item } => LogEntry::AnnotatedSystemItem { + ts, + entry: LoggedSystemHistoryEntry { + item, + metadata: legacy_metadata(session_id, segment_id, line_index, 0), + }, + }, + canonical => canonical, + } +} + fn atomic_write_json(path: &Path, value: &T) -> Result<(), StoreError> { let mut bytes = serde_json::to_vec_pretty(value)?; bytes.push(b'\n'); @@ -445,7 +666,7 @@ mod tests { } #[test] - fn schema_v1_logs_are_validated_and_promoted_to_v2() { + fn schema_v1_logs_are_rewritten_and_promoted_to_v3() { let root = tempfile::tempdir().unwrap(); let session_id = new_session_id(); let segment_id = new_segment_id(); @@ -467,7 +688,7 @@ mod tests { } #[test] - fn schema_v1_migration_rejects_corrupt_log_before_manifest_update() { + fn schema_v1_migration_rejects_corrupt_log_before_v3_manifest_update() { let root = tempfile::tempdir().unwrap(); let session_id = new_session_id(); let manifest = SessionManifest { @@ -492,6 +713,135 @@ mod tests { assert_eq!(persisted.schema_version, LEGACY_SESSION_SCHEMA_VERSION); } + #[test] + fn schema_v2_migration_rewrites_legacy_records_with_stable_unknown_provenance() { + let root = tempfile::tempdir().unwrap(); + let session_id = new_session_id(); + let segment_id = new_segment_id(); + fs::create_dir_all(root.path().join(SEGMENTS_DIR)).unwrap(); + atomic_write_json( + &root.path().join(SESSION_FILE), + &SessionManifest { + schema_version: PREVIOUS_SESSION_SCHEMA_VERSION, + session_id, + }, + ) + .unwrap(); + let source = vec![ + LogEntry::SegmentStart { + ts: 1, + session_id, + system_prompt: None, + config: agen::llm_client::RequestConfig::default(), + history: vec![LoggedItem::from(agen::Item::assistant_message("prior"))], + forked_from: None, + compacted_from: None, + }, + LogEntry::UserInput { + ts: 2, + segments: vec![protocol::Segment::Text { + content: "hello".into(), + }], + extensions: Vec::new(), + }, + LogEntry::AssistantItem { + ts: 3, + item: LoggedItem::from(agen::Item::assistant_message("reply")), + }, + ]; + let path = root + .path() + .join(SEGMENTS_DIR) + .join(format!("{segment_id}.jsonl")); + let mut bytes = Vec::new(); + for entry in source { + serde_json::to_writer(&mut bytes, &entry).unwrap(); + bytes.push(b'\n'); + } + fs::write(&path, bytes).unwrap(); + + let store = WorkerSessionStore::new(root.path()).unwrap(); + let first = store.read_all(session_id, segment_id).unwrap(); + assert!(matches!(first[0], LogEntry::AnnotatedSegmentStart { .. })); + assert!(matches!(first[1], LogEntry::AnnotatedUserInput { .. })); + assert!(matches!(first[2], LogEntry::AnnotatedAssistantItem { .. })); + let first_bytes = fs::read(&path).unwrap(); + drop(store); + + let reopened = WorkerSessionStore::new(root.path()).unwrap(); + assert_eq!(fs::read(&path).unwrap(), first_bytes); + let snapshot = crate::public_snapshot::project_current_session_snapshot( + &reopened.read_all(session_id, segment_id).unwrap(), + ); + assert_eq!(snapshot.entries.len(), 3); + assert!(snapshot.entries.iter().all(|entry| { + entry.provenance == protocol::SessionEntryProvenance::LegacyUnknown + && entry.entry_id.len() <= 64 + })); + } + + #[test] + fn schema_v3_rejects_legacy_records_and_new_writes_are_canonical() { + let root = tempfile::tempdir().unwrap(); + let session_id = new_session_id(); + let segment_id = new_segment_id(); + let store = WorkerSessionStore::new(root.path()).unwrap(); + store + .create_segment( + session_id, + segment_id, + &[LogEntry::SegmentStart { + ts: 1, + session_id, + system_prompt: None, + config: agen::llm_client::RequestConfig::default(), + history: Vec::new(), + forked_from: None, + compacted_from: None, + }], + ) + .unwrap(); + store + .append( + session_id, + segment_id, + &LogEntry::UserInput { + ts: 2, + segments: vec![protocol::Segment::Text { + content: "new".into(), + }], + extensions: Vec::new(), + }, + ) + .unwrap(); + let entries = store.read_all(session_id, segment_id).unwrap(); + assert!(matches!(entries[0], LogEntry::AnnotatedSegmentStart { .. })); + assert!(matches!(entries[1], LogEntry::AnnotatedUserInput { .. })); + drop(store); + + let path = root + .path() + .join(SEGMENTS_DIR) + .join(format!("{segment_id}.jsonl")); + let mut file = OpenOptions::new().append(true).open(path).unwrap(); + serde_json::to_writer( + &mut file, + &LogEntry::SystemItem { + ts: 3, + item: crate::SystemItem::LegacyIgnored { + slug: "legacy".into(), + }, + }, + ) + .unwrap(); + file.write_all(b"\n").unwrap(); + let error = match WorkerSessionStore::new(root.path()) { + Ok(_) => panic!("schema v3 must reject a legacy history record"), + Err(error) => error, + }; + assert!(matches!(error, StoreError::Corrupt { .. })); + } + #[test] fn reopen_preserves_session_and_segment_ids() { let root = tempfile::tempdir().unwrap(); diff --git a/crates/tui/src/app.rs b/crates/tui/src/app.rs index fd8b7bc9..ba158e8b 100644 --- a/crates/tui/src/app.rs +++ b/crates/tui/src/app.rs @@ -1098,10 +1098,9 @@ impl App { self.blocks.push(Block::UserMessage { segments }); self.assistant_streaming = false; } - Event::SegmentRotated { entry } => { + Event::SegmentRotated { session } => { let retained_run_errors = self.run_error_messages.clone(); - self.reset_for_rotation(); - self.apply_log_entry_raw(&entry); + self.restore_session(&session, self.greeting.clone()); for message in retained_run_errors { self.blocks.push(Block::Alert { level: AlertLevel::Error, @@ -1408,14 +1407,14 @@ impl App { self.latest_memory_worker_event = Some(event.message); } Event::Snapshot { - entries, + session, greeting, status, in_flight, internal_workers, } => { self.rewind_refresh_fence = false; - self.restore_snapshot(&entries, greeting, in_flight); + self.restore_snapshot(&session, greeting, in_flight); self.replace_internal_worker_snapshots(internal_workers); self.set_worker_status(status); } @@ -1455,11 +1454,11 @@ impl App { } } Event::RewindApplied { - entries, + session, input, summary, } => { - self.restore_rewind_snapshot(&entries); + self.restore_rewind_snapshot(&session); self.rewind_refresh_fence = true; let restored_composer = if self.input.is_empty() { self.input.replace_with_segments(&input); @@ -2173,7 +2172,7 @@ impl App { ) -> InternalWorkerView { let mut app = App::new(snapshot.worker.name.clone()); app.mode = mode; - app.restore_entries(&snapshot.entries, None); + app.restore_session(&snapshot.session, None); app.apply_in_flight_snapshot(snapshot.in_flight); app.set_worker_status(snapshot.status); if let Some(error) = snapshot.error { @@ -2254,14 +2253,14 @@ impl App { fn restore_snapshot( &mut self, - entries: &[serde_json::Value], + session: &protocol::SessionSnapshot, greeting: protocol::Greeting, in_flight: InFlightSnapshot, ) { self.greeting = Some(greeting.clone()); self.context_window = greeting.context_window; self.session_context_tokens = greeting.context_tokens; - self.restore_entries(entries, Some(greeting)); + self.restore_session(session, Some(greeting)); self.apply_in_flight_snapshot(in_flight); } @@ -2270,7 +2269,7 @@ impl App { /// session tail; always clear/replay from it even if this TUI instance has /// somehow lost connect-time greeting metadata. Skipping the restore in /// that case would leave old post-target output visible after success. - fn restore_rewind_snapshot(&mut self, entries: &[serde_json::Value]) { + fn restore_rewind_snapshot(&mut self, session: &protocol::SessionSnapshot) { let greeting = self.greeting.clone().or_else(|| { self.blocks.iter().find_map(|b| match b { Block::Greeting(g) => Some(g.clone()), @@ -2283,7 +2282,7 @@ impl App { self.session_context_tokens = greeting.context_tokens; } let missing_greeting = greeting.is_none(); - self.restore_entries(entries, greeting); + self.restore_session(session, greeting); if missing_greeting { self.blocks.push(Block::Alert { level: AlertLevel::Warn, @@ -2293,9 +2292,9 @@ impl App { } } - fn restore_entries( + fn restore_session( &mut self, - entries: &[serde_json::Value], + session: &protocol::SessionSnapshot, greeting: Option, ) { self.run_error_messages.clear(); @@ -2309,137 +2308,90 @@ impl App { } self.assistant_streaming = false; - for entry in entries { - self.apply_log_entry_raw(entry); + for entry in &session.entries { + use protocol::{SessionContentPart, SessionMessageRole, SessionSnapshotEntryData}; + match &entry.data { + SessionSnapshotEntryData::UserInput { segments } => { + self.turn_index += 1; + self.blocks.push(Block::TurnHeader { + turn: self.turn_index, + }); + if !segments.is_empty() { + self.blocks.push(Block::UserMessage { + segments: segments.clone(), + }); + } + } + SessionSnapshotEntryData::Message { role, content } => { + let role = match role { + SessionMessageRole::User => agen::Role::User, + SessionMessageRole::Assistant => agen::Role::Assistant, + }; + let item = agen::Item::Message { + id: None, + role, + content: content + .iter() + .map(|part| match part { + SessionContentPart::Text { text } => { + agen::ContentPart::Text { text: text.clone() } + } + SessionContentPart::Refusal { refusal } => { + agen::ContentPart::Refusal { + refusal: refusal.clone(), + } + } + }) + .collect(), + status: None, + }; + let value = serde_json::to_value(item).expect("Item is Serialize"); + self.push_history_item(&value); + } + SessionSnapshotEntryData::ToolCall { + call_id, + name, + arguments, + } => { + let item = + agen::Item::tool_call(call_id.clone(), name.clone(), arguments.clone()); + let value = serde_json::to_value(item).expect("Item is Serialize"); + self.push_history_item(&value); + } + SessionSnapshotEntryData::ToolResult { + call_id, + summary, + content, + is_error, + .. + } => { + let item = agen::Item::tool_result_item( + call_id.clone(), + summary.clone(), + content.clone(), + *is_error, + ); + let value = serde_json::to_value(item).expect("Item is Serialize"); + self.push_history_item(&value); + } + SessionSnapshotEntryData::SystemItem { data, .. } => { + if let Some(data) = data { + self.apply_system_item(data); + } + } + SessionSnapshotEntryData::RunError { message } => { + self.push_run_error(message.clone()); + } + } } - self.mark_orphan_tool_calls_incomplete_pass(); } - /// Drop the derived view in preparation for replaying a new - /// `SegmentStart` (compaction / fork). Greeting is preserved - /// because the Worker identity hasn't changed. - fn reset_for_rotation(&mut self) { - let greeting = self.blocks.iter().find_map(|b| match b { - Block::Greeting(g) => Some(g.clone()), - _ => None, - }); - self.turn_index = 0; - self.blocks.clear(); - self.cache = FileCache::new(); - self.task_store = TaskStore::new(); - self.task_pane_scroll = 0; - if let Some(g) = greeting { - self.greeting = Some(g.clone()); - self.blocks.push(Block::Greeting(g)); - } - } - - /// Walk a single `LogEntry` JSON value and translate it into blocks - /// the live event path would have produced. Shared between - /// `restore_snapshot` (replay path) and `apply_log_entry` (live - /// path). - fn apply_log_entry_raw(&mut self, value: &serde_json::Value) { - let Ok(entry) = serde_json::from_value::(value.clone()) else { - return; - }; - match entry { - session_store::LogEntry::SegmentStart { history, .. } => { - for logged in history { - let item: agen::Item = logged.into(); - let item_value = serde_json::to_value(&item).expect("Item is Serialize"); - self.push_history_item(&item_value); - } - } - session_store::LogEntry::UserInput { segments, .. } => { - self.turn_index += 1; - self.blocks.push(Block::TurnHeader { - turn: self.turn_index, - }); - if !segments.is_empty() { - self.blocks.push(Block::UserMessage { segments }); - } - } - session_store::LogEntry::AssistantItem { item, .. } - | session_store::LogEntry::ToolResult { item, .. } => { - let it: agen::Item = item.into(); - let item_value = serde_json::to_value(&it).expect("Item is Serialize"); - self.push_history_item(&item_value); - } - session_store::LogEntry::SystemItem { item, .. } => { - let value = serde_json::to_value(&item).expect("SystemItem is Serialize"); - self.apply_system_item(&value); - } - session_store::LogEntry::Extension { - domain, payload, .. - } if domain == "yoi.compaction" => { - self.apply_compaction_extension(&payload); - } - session_store::LogEntry::RunErrored { message, .. } => { - self.push_run_error(message); - } - // Non-history-bearing variants don't affect the block view. - _ => {} - } - } - /// Dispatch one `SystemItem` JSON value into the appropriate block. /// /// Kind-based routing replaces the old free-text `[Notification]` / /// `[File: …]` parsing path: each kind maps directly to a typed /// block (`Block::Notify`, `Block::WorkerEvent`, …). - fn apply_compaction_extension(&mut self, payload: &serde_json::Value) { - if payload.get("kind").and_then(|value| value.as_str()) != Some("compaction_block") { - return; - } - match payload.get("state").and_then(|value| value.as_str()) { - Some("running") => { - if self.last_streaming_compact_mut().is_none() { - self.blocks.push(Block::Compact(CompactEvent::Streaming { - started_at: Instant::now(), - })); - } - } - Some("done") => { - let new_segment_id = payload - .get("new_segment_id") - .and_then(|value| value.as_str()) - .and_then(|value| value.parse::().ok()) - .unwrap_or_else(uuid::Uuid::nil); - if let Some(evt) = self.last_streaming_compact_mut() { - *evt = CompactEvent::Done { - new_segment_id, - elapsed_secs: None, - }; - } else { - self.blocks.push(Block::Compact(CompactEvent::Done { - new_segment_id, - elapsed_secs: None, - })); - } - } - Some("failed") => { - let error = payload - .get("error") - .and_then(|value| value.as_str()) - .unwrap_or("compact failed") - .to_string(); - if let Some(evt) = self.last_streaming_compact_mut() { - *evt = CompactEvent::Failed { - error, - elapsed_secs: None, - }; - } else { - self.blocks.push(Block::Compact(CompactEvent::Failed { - error, - elapsed_secs: None, - })); - } - } - _ => {} - } - } - fn apply_system_item(&mut self, value: &serde_json::Value) { let Ok(item) = serde_json::from_value::(value.clone()) else { // Unknown / forward-compat shape: fall back to rendering the @@ -2542,6 +2494,15 @@ fn fmt_millis(ms: u64) -> String { } } +#[cfg(test)] +fn public_session(values: Vec) -> protocol::SessionSnapshot { + let entries = values + .into_iter() + .map(|value| serde_json::from_value(value).expect("LogEntry deserializes")) + .collect::>(); + session_store::public_snapshot::project_current_session_snapshot(&entries) +} + fn message_text(item: &serde_json::Value) -> String { item["content"] .as_array() @@ -2685,7 +2646,7 @@ mod rewind_refresh_tests { }); app.handle_worker_event(Event::RewindApplied { - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, input: vec![Segment::text("selected rewind input")], summary: summary(3), }); @@ -2704,7 +2665,7 @@ mod rewind_refresh_tests { }); app.handle_worker_event(Event::RewindApplied { - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, input: vec![Segment::text("rewound input")], summary: summary(1), }); @@ -2747,7 +2708,7 @@ mod rewind_refresh_tests { }); app.handle_worker_event(Event::RewindApplied { - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, input: vec![Segment::text("rewound input")], summary: summary(2), }); @@ -3289,7 +3250,9 @@ mod completion_flow_tests { }; app.handle_worker_event(Event::SegmentRotated { - entry: serde_json::to_value(start).expect("LogEntry is Serialize"), + session: public_session(vec![ + serde_json::to_value(start).expect("LogEntry is Serialize"), + ]), }); app.handle_worker_event(Event::UserMessage { segments: vec![Segment::text("first persisted message")], @@ -3533,7 +3496,7 @@ mod completion_flow_tests { } #[test] - fn snapshot_renders_system_message_block_from_session_start() { + fn snapshot_excludes_system_prompt_history_from_public_blocks() { let mut app = App::new("test".into()); let session_start = session_store::LogEntry::SegmentStart { ts: 1, @@ -3549,7 +3512,7 @@ mod completion_flow_tests { let session_start_value = serde_json::to_value(&session_start).unwrap(); app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: vec![session_start_value], + session: public_session(vec![session_start_value]), status: WorkerStatus::Running, in_flight: Default::default(), internal_workers: Vec::new(), @@ -3557,10 +3520,8 @@ mod completion_flow_tests { assert!(matches!(app.worker_status, WorkerStatus::Running)); assert!(app.running); - assert!(matches!( - app.blocks.get(1), - Some(Block::SystemMessage { text }) if text == "[File: src/main.rs]\nfn main() {}" - )); + assert_eq!(app.blocks.len(), 1); + assert!(matches!(app.blocks.first(), Some(Block::Greeting(_)))); } #[test] @@ -3595,7 +3556,7 @@ mod completion_flow_tests { }; app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: vec![serde_json::to_value(run_errored).unwrap()], + session: public_session(vec![serde_json::to_value(run_errored).unwrap()]), status: WorkerStatus::Idle, in_flight: Default::default(), internal_workers: Vec::new(), @@ -3633,7 +3594,7 @@ mod completion_flow_tests { compacted_from: None, }; app.handle_worker_event(Event::SegmentRotated { - entry: serde_json::to_value(segment_start).unwrap(), + session: public_session(vec![serde_json::to_value(segment_start).unwrap()]), }); let errors = app @@ -3656,7 +3617,9 @@ mod completion_flow_tests { let mut app = App::new("test".into()); app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, status: WorkerStatus::Running, in_flight: InFlightSnapshot { blocks: vec![ @@ -3762,7 +3725,9 @@ mod completion_flow_tests { }, revision, status: WorkerStatus::Idle, - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, in_flight: protocol::InFlightSnapshot::default(), error: None, internal_workers: Vec::new(), @@ -3977,7 +3942,9 @@ mod completion_flow_tests { assert_eq!(app.selected_worker_view().worker_name, "parent"); app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, status: WorkerStatus::Idle, in_flight: Default::default(), internal_workers: Vec::new(), @@ -4026,7 +3993,9 @@ mod completion_flow_tests { }); app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, status: WorkerStatus::Idle, in_flight: Default::default(), internal_workers: vec![InternalWorkerSnapshot { @@ -4037,7 +4006,9 @@ mod completion_flow_tests { kind: protocol::InternalWorkerKind::SubWorker, }, revision: 4, - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, status: WorkerStatus::Running, error: None, in_flight: Default::default(), @@ -4193,7 +4164,9 @@ mod completion_flow_tests { greeting.context_tokens = 45_000; app.handle_worker_event(Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting, status: WorkerStatus::Idle, in_flight: Default::default(), @@ -4396,7 +4369,7 @@ mod completion_flow_tests { ]; app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: assistant_item_entries, + session: public_session(assistant_item_entries), status: WorkerStatus::Running, in_flight: Default::default(), internal_workers: Vec::new(), diff --git a/crates/tui/src/console/mod.rs b/crates/tui/src/console/mod.rs index 0a6262ee..92a3a193 100644 --- a/crates/tui/src/console/mod.rs +++ b/crates/tui/src/console/mod.rs @@ -547,7 +547,9 @@ async fn run_e2e_rewind_fixture( let mut app = App::new_with_persistent_input_history(worker_name.clone(), &workspace_root); app.connected = true; app.handle_worker_event(Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, status: WorkerStatus::Idle, greeting: Greeting { worker_name: worker_name.clone(), @@ -673,7 +675,9 @@ async fn run_e2e_rewind_fixture( if let Some(submitted_at) = pending_apply { if submitted_at.elapsed() >= apply_delay { app.handle_worker_event(Event::RewindApplied { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, input: vec![Segment::text("rewind-live-refresh")], summary: RewindSummary { truncated_to_entries: 1, @@ -2023,13 +2027,13 @@ mod tests { let mut app = App::new("agent".to_string()); app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, status: WorkerStatus::Idle, in_flight: Default::default(), internal_workers: Vec::new(), }); app.handle_worker_event(Event::RewindApplied { - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, input: vec![Segment::Text { content: "retry this".into(), }], @@ -2050,7 +2054,7 @@ mod tests { let mut app = App::new("agent".to_string()); app.handle_worker_event(Event::Snapshot { greeting: test_greeting(), - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, status: WorkerStatus::Idle, in_flight: Default::default(), internal_workers: Vec::new(), @@ -2058,7 +2062,7 @@ mod tests { type_keys(&mut app, "draft"); app.handle_worker_event(Event::RewindApplied { - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, input: vec![Segment::Text { content: "retry this".into(), }], diff --git a/crates/tui/src/dashboard/tests.rs b/crates/tui/src/dashboard/tests.rs index 32be9393..a86a3040 100644 --- a/crates/tui/src/dashboard/tests.rs +++ b/crates/tui/src/dashboard/tests.rs @@ -849,7 +849,9 @@ async fn ticket_queue_notification_sends_notify_when_socket_available() { let mut writer = JsonLineWriter::new(writer); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "test-orchestrator".to_string(), cwd: temp.path().display().to_string(), @@ -891,7 +893,9 @@ async fn send_notify_only_can_deliver_weak_notification_without_auto_run() { let mut writer = JsonLineWriter::new(writer); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "yoi".to_string(), cwd: temp.path().display().to_string(), diff --git a/crates/tui/src/worker_list.rs b/crates/tui/src/worker_list.rs index 49236489..fa2d834c 100644 --- a/crates/tui/src/worker_list.rs +++ b/crates/tui/src/worker_list.rs @@ -910,7 +910,7 @@ mod tests { timestamp_ms: 0, }), Event::Snapshot { - entries: vec![], + session: protocol::SessionSnapshot { entries: vec![] }, greeting: test_greeting(), status: WorkerStatus::Idle, in_flight: Default::default(), diff --git a/crates/worker-runtime/src/runtime.rs b/crates/worker-runtime/src/runtime.rs index 03f3b5f4..a1d2a99d 100644 --- a/crates/worker-runtime/src/runtime.rs +++ b/crates/worker-runtime/src/runtime.rs @@ -1530,7 +1530,9 @@ impl Runtime { } } Ok(protocol::Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: worker_ref.worker_id.to_string(), cwd: String::new(), @@ -3152,7 +3154,9 @@ mod tests { ), ); let snapshot = protocol::Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "parent".to_string(), cwd: "/tmp".to_string(), @@ -4581,7 +4585,16 @@ mod tests { backend.set_worker_snapshot( &detail.worker_ref, protocol::Event::Snapshot { - entries: vec![expected_entry.clone()], + session: protocol::SessionSnapshot { + entries: vec![protocol::SessionSnapshotEntry { + entry_id: "restored-log-entry".to_owned(), + provenance: protocol::SessionEntryProvenance::LegacyUnknown, + derived_from: Vec::new(), + data: protocol::SessionSnapshotEntryData::RunError { + message: expected_entry.to_string(), + }, + }], + }, greeting: protocol::Greeting { worker_name: "live-worker".to_string(), cwd: "/tmp/live".to_string(), @@ -4606,12 +4619,13 @@ mod tests { .unwrap(); match snapshot { protocol::Event::Snapshot { - entries, + session, greeting, status, .. } => { - assert_eq!(entries, vec![expected_entry]); + assert_eq!(session.entries.len(), 1); + assert_eq!(session.entries[0].entry_id, "restored-log-entry"); assert_eq!(greeting.worker_name, "live-worker"); assert_eq!(status, protocol::WorkerStatus::Running); } diff --git a/crates/worker/src/controller.rs b/crates/worker/src/controller.rs index ff943c60..4042ab68 100644 --- a/crates/worker/src/controller.rs +++ b/crates/worker/src/controller.rs @@ -84,10 +84,7 @@ impl WorkerHandle { (entries, entry_rx, in_flight) }; let event = Event::Snapshot { - entries: entries - .into_iter() - .map(|entry| serde_json::to_value(entry).expect("log entry serializes")) - .collect(), + session: session_store::public_snapshot::project_current_session_snapshot(&entries), greeting: self.shared_state.greeting.clone(), status: self.shared_state.get_status(), in_flight, @@ -1874,28 +1871,16 @@ where St: Store, { match worker.rewind_to(target, expected_head_entries) { - Ok(applied) => match applied - .entries - .into_iter() - .map(serde_json::to_value) - .collect::, _>>() - { - Ok(entries) => { - let _ = event_tx.send(Event::RewindApplied { - entries, - input: applied.input, - summary: applied.summary, - }); - true - } - Err(error) => { - let _ = event_tx.send(Event::Error { - code: ErrorCode::Internal, - message: format!("failed to encode rewind snapshot: {error}"), - }); - false - } - }, + Ok(applied) => { + let session = + session_store::public_snapshot::project_current_session_snapshot(&applied.entries); + let _ = event_tx.send(Event::RewindApplied { + session, + input: applied.input, + summary: applied.summary, + }); + true + } Err(err) => { let _ = event_tx.send(Event::Error { code: ErrorCode::InvalidRequest, @@ -2101,7 +2086,9 @@ mod tests { let mut writer = JsonLineWriter::new(w); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "parent".into(), cwd: "/tmp".into(), diff --git a/crates/worker/src/discovery.rs b/crates/worker/src/discovery.rs index 6e31005b..03a499d8 100644 --- a/crates/worker/src/discovery.rs +++ b/crates/worker/src/discovery.rs @@ -1481,7 +1481,9 @@ mod tests { let mut writer = JsonLineWriter::new(stream); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "target".into(), cwd: "/tmp".into(), @@ -1514,7 +1516,9 @@ mod tests { .unwrap(); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "target".into(), cwd: "/tmp".into(), @@ -1603,7 +1607,9 @@ mod tests { let mut writer = JsonLineWriter::new(stream); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "target".into(), cwd: "/tmp".into(), @@ -1627,7 +1633,9 @@ mod tests { let mut writer = JsonLineWriter::new(writer_half); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "target".into(), cwd: "/tmp".into(), @@ -1729,7 +1737,9 @@ mod tests { .unwrap(); writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "alerted".into(), cwd: "/tmp".into(), @@ -1779,7 +1789,9 @@ mod tests { let mut writer = JsonLineWriter::new(stream); let _ = writer .write(&Event::Snapshot { - entries: Vec::new(), + session: protocol::SessionSnapshot { + entries: Vec::new(), + }, greeting: protocol::Greeting { worker_name: "child-live".into(), cwd: "/tmp".into(), diff --git a/crates/worker/src/feature/builtin/worker_observation.rs b/crates/worker/src/feature/builtin/worker_observation.rs index 68305b2c..ad59bb3f 100644 --- a/crates/worker/src/feature/builtin/worker_observation.rs +++ b/crates/worker/src/feature/builtin/worker_observation.rs @@ -6,7 +6,7 @@ use agen::tool::{Tool, ToolDefinition, ToolError, ToolMeta, ToolOutput}; use async_trait::async_trait; use schemars::JsonSchema; use serde::{Deserialize, Serialize}; -use session_store::{LogEntry, collect_state}; +use session_store::LogEntry; use super::manage_worker::{WORKER_CONTROL_SERVICE_ID, WorkerControlService}; use crate::feature::{ @@ -61,7 +61,7 @@ pub struct WorkerObservationSubject { #[derive(Debug, Clone)] pub struct WorkerSessionCapture { pub segment_id: String, - pub entries: Vec>, + pub session: protocol::SessionSnapshot, } impl WorkerSessionCapture { @@ -69,17 +69,9 @@ impl WorkerSessionCapture { segment_id: impl Into, log_entries: &[LogEntry], ) -> Result { - let segment_id = segment_id.into(); - let state = collect_state(log_entries); - let parsed_segment_id = segment_id.parse().unwrap_or_default(); - let entries = crate::session_history::restore_history_entries( - state.session_id.unwrap_or_default(), - parsed_segment_id, - log_entries, - )?; Ok(Self { - segment_id, - entries, + segment_id: segment_id.into(), + session: session_store::public_snapshot::project_current_session_snapshot(log_entries), }) } } @@ -115,7 +107,7 @@ struct WorkspaceWorkerObservationListResponse { #[derive(Debug, Deserialize)] struct WorkspaceWorkerObservationCaptureResponse { segment_id: String, - entries: Vec, + session: protocol::SessionSnapshot, } pub struct WorkspaceClientWorkerObservationProvider { @@ -173,26 +165,9 @@ impl WorkerObservationProvider for WorkspaceClientWorkerObservationProvider { let body = workspace_response_body(response)?; let response = serde_json::from_str::(&body) .map_err(|error| WorkerObservationError::Unavailable(error.to_string()))?; - let entries = response - .entries - .into_iter() - .map(|entry| { - serde_json::from_value(entry) - .map_err(|error| WorkerObservationError::Unavailable(error.to_string())) - }) - .collect::, _>>()?; - let state = collect_state(&entries); - let segment_id = response.segment_id; - let parsed_segment_id = segment_id.parse().unwrap_or_default(); - let typed_entries = crate::session_history::restore_history_entries( - state.session_id.unwrap_or_default(), - parsed_segment_id, - &entries, - ) - .map_err(WorkerObservationError::Unavailable)?; Ok(WorkerSessionCapture { - segment_id, - entries: typed_entries, + segment_id: response.segment_id, + session: response.session, }) } } @@ -420,16 +395,9 @@ impl WorkerObservationProvider for SpawnedSubWorkerObservationProvider { .get_internal(name) .ok_or(WorkerObservationError::NotFound)?; let entries = record.session.entries(); - let state = collect_state(&entries); - let typed_entries = crate::session_history::restore_history_entries( - state.session_id.unwrap_or_default(), - Default::default(), - &entries, - ) - .map_err(WorkerObservationError::Unavailable)?; Ok(WorkerSessionCapture { segment_id: format!("subworker:{name}"), - entries: typed_entries, + session: session_store::public_snapshot::project_current_session_snapshot(&entries), }) } } @@ -699,9 +667,9 @@ async fn latest_view( .capture_worker_session(subject) .await .map_err(tool_error)?; - Ok(SessionCapture::from_history_entries( + Ok(SessionCapture::from_session_snapshot( capture.segment_id, - capture.entries, + capture.session, )) } @@ -799,16 +767,42 @@ mod tests { .clone() .into_iter() .enumerate() - .map(|(index, item)| { - let mut metadata = crate::SessionHistoryMetadata::legacy_unknown(); - metadata.entry_id = - session_store::LoggedSessionHistoryEntryId(format!("fake-{index:08}")); - agen::HistoryEntry::new(item, metadata) + .filter_map(|(index, item)| { + let data = match item { + Item::Message { role, content, .. } => { + let role = match role { + Role::User => protocol::SessionMessageRole::User, + Role::Assistant => protocol::SessionMessageRole::Assistant, + Role::System => return None, + }; + protocol::SessionSnapshotEntryData::Message { + role, + content: content + .into_iter() + .map(|part| match part { + agen::ContentPart::Text { text } => { + protocol::SessionContentPart::Text { text } + } + agen::ContentPart::Refusal { refusal } => { + protocol::SessionContentPart::Refusal { refusal } + } + }) + .collect(), + } + } + _ => return None, + }; + Some(protocol::SessionSnapshotEntry { + entry_id: format!("fake-{index:08}"), + provenance: protocol::SessionEntryProvenance::LegacyUnknown, + derived_from: Vec::new(), + data, + }) }) .collect(); Ok(WorkerSessionCapture { segment_id: "segment".to_string(), - entries, + session: protocol::SessionSnapshot { entries }, }) } } diff --git a/crates/worker/src/internal_worker.rs b/crates/worker/src/internal_worker.rs index a8d21e49..e033369e 100644 --- a/crates/worker/src/internal_worker.rs +++ b/crates/worker/src/internal_worker.rs @@ -335,7 +335,7 @@ enum InternalWorkerSessionCommand { /// task; protocol access is consumed only by the owning parent registry. #[derive(Debug, Clone)] pub(crate) struct InternalWorkerSessionSnapshot { - pub entries: Vec, + pub session: protocol::SessionSnapshot, pub status: WorkerStatus, pub error: Option, pub in_flight: InFlightSnapshot, @@ -402,7 +402,7 @@ impl InternalWorkerSessionHandle { (entries, snapshot_from_guard(&guard)) }; InternalWorkerSessionSnapshot { - entries, + session: session_store::public_snapshot::project_current_session_snapshot(&entries), status: match self.status() { InternalWorkerSessionStatus::Running => WorkerStatus::Running, InternalWorkerSessionStatus::Paused => WorkerStatus::Paused, diff --git a/crates/worker/src/ipc/protocol_session.rs b/crates/worker/src/ipc/protocol_session.rs index 9b749126..5e16069b 100644 --- a/crates/worker/src/ipc/protocol_session.rs +++ b/crates/worker/src/ipc/protocol_session.rs @@ -30,8 +30,9 @@ pub fn subscribe_worker_protocol_session(handle: &WorkerHandle) -> WorkerProtoco pub fn live_log_entry_event(entry: LogEntry) -> Option { match entry { entry @ (LogEntry::SegmentStart { .. } | LogEntry::AnnotatedSegmentStart { .. }) => { - let value = serde_json::to_value(&entry).expect("LogEntry is Serialize"); - Some(Event::SegmentRotated { entry: value }) + let session = + session_store::public_snapshot::project_current_session_snapshot(&[entry]); + Some(Event::SegmentRotated { session }) } LogEntry::UserInput { segments, .. } | LogEntry::AnnotatedUserInput { segments, .. } => { Some(Event::UserMessage { segments }) diff --git a/crates/worker/src/session_capture.rs b/crates/worker/src/session_capture.rs index 1886e4e8..c0be48d0 100644 --- a/crates/worker/src/session_capture.rs +++ b/crates/worker/src/session_capture.rs @@ -8,6 +8,10 @@ use std::sync::Arc; use crate::session_history::{SessionHistoryMetadata, WorkerHistoryProvenance}; use agen::{HistoryEntry, Item, Role}; +use protocol::{ + SessionContentPart, SessionEntryProvenance, SessionMessageRole, SessionSnapshot, + SessionSnapshotEntryData, +}; use serde::{Deserialize, Serialize}; const DEFAULT_SEARCH_LIMIT: usize = 20; @@ -225,6 +229,66 @@ pub(crate) struct SessionCapture { } impl SessionCapture { + pub(crate) fn from_session_snapshot( + segment_id: impl Into, + snapshot: SessionSnapshot, + ) -> Self { + let entries = snapshot + .entries + .into_iter() + .filter_map(|entry| { + let item = match entry.data { + SessionSnapshotEntryData::UserInput { segments } => { + Item::user_message(protocol::Segment::flatten_to_text(&segments)) + } + SessionSnapshotEntryData::Message { role, content } => { + let role = match role { + SessionMessageRole::User => Role::User, + SessionMessageRole::Assistant => Role::Assistant, + }; + Item::Message { + id: None, + role, + content: content + .into_iter() + .map(|part| match part { + SessionContentPart::Text { text } => { + agen::ContentPart::Text { text } + } + SessionContentPart::Refusal { refusal } => { + agen::ContentPart::Refusal { refusal } + } + }) + .collect(), + status: None, + } + } + SessionSnapshotEntryData::ToolCall { + call_id, + name, + arguments, + } => Item::tool_call(call_id, name, arguments), + SessionSnapshotEntryData::ToolResult { + call_id, + summary, + content, + is_error, + attachments: _, + } => Item::tool_result_item(call_id, summary, content, is_error), + // Observation deliberately excludes system items and + // controller errors from model-visible session evidence. + SessionSnapshotEntryData::SystemItem { .. } + | SessionSnapshotEntryData::RunError { .. } => return None, + }; + Some(HistoryEntry::new( + item, + public_snapshot_metadata(entry.entry_id, entry.provenance), + )) + }) + .collect(); + Self::from_history_entries(segment_id, entries) + } + pub(crate) fn new(segment_id: impl Into, items: Vec) -> Self { let entries = items .into_iter() @@ -543,6 +607,46 @@ impl SessionCapture { } } +fn public_snapshot_metadata( + entry_id: String, + provenance: SessionEntryProvenance, +) -> SessionHistoryMetadata { + let worker = session_store::LoggedWorkerSubject { + workspace_id: None, + runtime_id: None, + worker_id: "public-session-snapshot".to_owned(), + }; + let origin = match provenance { + SessionEntryProvenance::HumanInput => WorkerHistoryProvenance::HumanInput { + account_id: "public-session-snapshot".to_owned(), + }, + SessionEntryProvenance::WorkerInput => WorkerHistoryProvenance::WorkerInput { + actor: worker.clone(), + }, + SessionEntryProvenance::FlowInstruction => WorkerHistoryProvenance::FlowInstruction { + selector: "public-session-snapshot".to_owned(), + definition_id: "public-session-snapshot".to_owned(), + definition_revision: 0, + instance_id: "public-session-snapshot".to_owned(), + state_id: "public-session-snapshot".to_owned(), + }, + SessionEntryProvenance::BackendInstruction => { + WorkerHistoryProvenance::BackendInstruction { operation_id: None } + } + SessionEntryProvenance::ModelOutput => WorkerHistoryProvenance::ModelOutput { + worker: worker.clone(), + }, + SessionEntryProvenance::ToolOutput => WorkerHistoryProvenance::ToolOutput { worker }, + SessionEntryProvenance::DerivedSummary => WorkerHistoryProvenance::DerivedSummary, + SessionEntryProvenance::LegacyUnknown => WorkerHistoryProvenance::LegacyUnknown, + }; + SessionHistoryMetadata { + entry_id: session_store::LoggedSessionHistoryEntryId(entry_id), + origin, + derivation: None, + } +} + fn message_reference_kind( origin: &WorkerHistoryProvenance, provider_role: &Role, diff --git a/crates/worker/src/spawn/comm_tools.rs b/crates/worker/src/spawn/comm_tools.rs index 03b1b126..9d30d89b 100644 --- a/crates/worker/src/spawn/comm_tools.rs +++ b/crates/worker/src/spawn/comm_tools.rs @@ -278,7 +278,20 @@ mod tests { fn snapshot(entries: Vec) -> Event { Event::Snapshot { - entries, + session: protocol::SessionSnapshot { + entries: entries + .into_iter() + .enumerate() + .map(|(index, value)| protocol::SessionSnapshotEntry { + entry_id: format!("test-{index}"), + provenance: protocol::SessionEntryProvenance::LegacyUnknown, + derived_from: Vec::new(), + data: protocol::SessionSnapshotEntryData::RunError { + message: value.to_string(), + }, + }) + .collect(), + }, greeting: Greeting { worker_name: "server".into(), cwd: "/tmp".into(), diff --git a/crates/worker/src/spawn/registry.rs b/crates/worker/src/spawn/registry.rs index 86ea8bae..8e0ea0f1 100644 --- a/crates/worker/src/spawn/registry.rs +++ b/crates/worker/src/spawn/registry.rs @@ -806,11 +806,7 @@ fn internal_worker_snapshot( InternalWorkerSnapshot { worker, revision, - entries: snapshot - .entries - .into_iter() - .filter_map(|entry| serde_json::to_value(entry).ok()) - .collect(), + session: snapshot.session, status: snapshot.status, error: snapshot.error, in_flight: snapshot.in_flight, @@ -1045,7 +1041,7 @@ mod tests { let snapshots = registry.internal_worker_snapshots(); assert_eq!(snapshots.len(), 1); assert_eq!(snapshots[0].revision, 2); - assert_eq!(snapshots[0].entries.len(), 1); + assert_eq!(snapshots[0].session.entries.len(), 1); record.session.emit_test_text_delta("partial"); let streamed = tokio::time::timeout(Duration::from_secs(1), parent_rx.recv()) diff --git a/crates/worker/src/spawn/tool.rs b/crates/worker/src/spawn/tool.rs index 21a9371a..8e4843d7 100644 --- a/crates/worker/src/spawn/tool.rs +++ b/crates/worker/src/spawn/tool.rs @@ -960,9 +960,7 @@ mod tests { use crate::WorkspaceId; use agen::llm_client::event::{Event as LlmEvent, ResponseStatus, StatusEvent}; - use agen::llm_client::types::ContentPart; use agen::llm_client::{ClientError, LlmClient, Request}; - use agen::{Item, Role}; use async_trait::async_trait; use futures::Stream; use manifest::{AuthRef, ModelManifest, SchemeKind, WorkerManifest}; @@ -1252,9 +1250,11 @@ extract_threshold = 4000 ) .await .unwrap(); - assert!(first_capture.entries.iter().map(|entry| &entry.item).any(|item| { - matches!(item, Item::Message { role: Role::Assistant, content, .. } if content.iter().any(|part| matches!(part, ContentPart::Text { text } if text.contains("reviewed")))) - })); + assert!( + serde_json::to_string(&first_capture.session) + .unwrap() + .contains("reviewed") + ); let send = (crate::spawn::comm_tools::sub_worker_send_tool(registry.clone()))().1; send.execute( @@ -1274,7 +1274,7 @@ extract_threshold = 4000 ) .await .unwrap(); - assert!(latest_capture.entries.len() > first_capture.entries.len()); + assert!(latest_capture.session.entries.len() > first_capture.session.entries.len()); fail_requests.store(true, Ordering::SeqCst); send.execute( diff --git a/crates/worker/tests/controller_test.rs b/crates/worker/tests/controller_test.rs index e4ef5b31..078fdc97 100644 --- a/crates/worker/tests/controller_test.rs +++ b/crates/worker/tests/controller_test.rs @@ -839,26 +839,26 @@ async fn snapshot_includes_user_input_for_in_flight_turn() { loop { let event = reader.next::().await.unwrap().unwrap(); match event { - Event::Snapshot { entries, .. } => { - // Walk the entries, find a `LogEntry::UserInput` and - // confirm its segments flatten to our submitted text. - let mut found = false; - for value in &entries { - let entry: session_store::LogEntry = - serde_json::from_value(value.clone()).expect("LogEntry deserialise"); - if let session_store::LogEntry::UserInput { segments, .. } - | session_store::LogEntry::AnnotatedUserInput { segments, .. } = entry - { - let text = protocol::Segment::flatten_to_text(&segments); - if text == "hello in-flight" { - found = true; - break; - } + Event::Snapshot { session, .. } => { + let found = session.entries.iter().any(|entry| match &entry.data { + protocol::SessionSnapshotEntryData::UserInput { segments } => { + protocol::Segment::flatten_to_text(segments) == "hello in-flight" } - } + protocol::SessionSnapshotEntryData::Message { + role: protocol::SessionMessageRole::User, + content, + } => content.iter().any(|part| { + matches!( + part, + protocol::SessionContentPart::Text { text } + if text == "hello in-flight" + ) + }), + _ => false, + }); assert!( found, - "snapshot must carry the in-flight UserInput entry: {entries:?}" + "snapshot must carry the in-flight UserInput entry: {session:?}" ); return; } @@ -2410,17 +2410,12 @@ async fn snapshot_contains_user_input(handle: &WorkerHandle, needle: &str) -> bo loop { let event = reader.next::().await.unwrap().unwrap(); match event { - Event::Snapshot { entries, .. } => { - return entries.into_iter().any(|value| { - let entry: session_store::LogEntry = - serde_json::from_value(value).expect("LogEntry deserialise"); - match entry { - session_store::LogEntry::UserInput { segments, .. } - | session_store::LogEntry::AnnotatedUserInput { segments, .. } => { - protocol::Segment::flatten_to_text(&segments).contains(needle) - } - _ => false, + Event::Snapshot { session, .. } => { + return session.entries.into_iter().any(|entry| match entry.data { + protocol::SessionSnapshotEntryData::UserInput { segments } => { + protocol::Segment::flatten_to_text(&segments).contains(needle) } + _ => false, }); } Event::Alert(_) => continue, diff --git a/crates/worker/tests/system_prompt_template_test.rs b/crates/worker/tests/system_prompt_template_test.rs index 53ca4141..9293917d 100644 --- a/crates/worker/tests/system_prompt_template_test.rs +++ b/crates/worker/tests/system_prompt_template_test.rs @@ -203,12 +203,12 @@ async fn session_start_state_captures_rendered_prompt() { .unwrap(); let first = entries.first().expect("at least one entry"); match first { - LogEntry::SegmentStart { system_prompt, .. } => { + LogEntry::AnnotatedSegmentStart { system_prompt, .. } => { let sp = system_prompt.as_deref().expect("system prompt set"); assert!(sp.starts_with("hello")); assert!(sp.contains(&pwd.display().to_string())); } - other => panic!("expected SegmentStart as first entry, got {other:?}"), + other => panic!("expected AnnotatedSegmentStart as first entry, got {other:?}"), } } diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index cc374820..da87769a 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -8513,7 +8513,7 @@ async fn scoped_capture_worker_observation_session( message: "worker protocol closed before the session snapshot".to_string(), }) })?; - let protocol::Event::Snapshot { entries, .. } = event else { + let protocol::Event::Snapshot { session, .. } = event else { return Err(ApiError::from(Error::RuntimeOperationFailed { runtime_id: target.runtime_id.clone(), code: "worker_observation_missing_snapshot".to_string(), @@ -8522,7 +8522,7 @@ async fn scoped_capture_worker_observation_session( }; Ok(Json(serde_json::json!({ "segment_id": format!("runtime:{}:worker:{}", target.runtime_id, target.worker_id), - "entries": entries, + "session": session, }))) } @@ -16785,7 +16785,7 @@ mod tests { ) .await .unwrap(); - assert!(capture["entries"].is_array()); + assert!(capture["session"]["entries"].is_array()); let revoked = api .store diff --git a/web/workspace/src/lib/generated/protocol.ts b/web/workspace/src/lib/generated/protocol.ts index 6a357acb..39ab33c2 100644 --- a/web/workspace/src/lib/generated/protocol.ts +++ b/web/workspace/src/lib/generated/protocol.ts @@ -73,11 +73,35 @@ export type InFlightBlock = { "kind": "text", text: string, finished?: boolean, export type InFlightSnapshot = { blocks?: Array, commands?: Array, }; +export type SessionEntryProvenance = "human_input" | "worker_input" | "flow_instruction" | "backend_instruction" | "model_output" | "tool_output" | "derived_summary" | "legacy_unknown"; + +export type SessionMessageRole = "user" | "assistant"; + +export type SessionContentPart = { "kind": "text", text: string, } | { "kind": "refusal", refusal: string, }; + +export type SessionToolAttachment = { media_type: string, +/** + * Base64-encoded durable attachment body. Public snapshots preserve the + * committed multimodal value instead of replacing it with placeholder text. + */ +data_base64: string, }; + +export type SessionSnapshotEntryData = { "kind": "user_input", segments: Array, } | { "kind": "message", role: SessionMessageRole, content: Array, } | { "kind": "tool_call", call_id: string, name: string, arguments: string, } | { "kind": "tool_result", call_id: string, summary: string, content?: string | null, is_error: boolean, attachments?: Array, } | { "kind": "system_item", item_kind: string, content: string, data?: unknown, } | { "kind": "run_error", message: string, }; + +export type SessionSnapshotEntry = { +/** + * Stable identity from durable history metadata, or a deterministic + * identity derived from the legacy segment and log position. + */ +entry_id: string, provenance: SessionEntryProvenance, derived_from?: Array, } & ({ "kind": "user_input", segments: Array, } | { "kind": "message", role: SessionMessageRole, content: Array, } | { "kind": "tool_call", call_id: string, name: string, arguments: string, } | { "kind": "tool_result", call_id: string, summary: string, content?: string | null, is_error: boolean, attachments?: Array, } | { "kind": "system_item", item_kind: string, content: string, data?: unknown, } | { "kind": "run_error", message: string, }); + +export type SessionSnapshot = { entries: Array, }; + export type InternalWorkerKind = "sub_worker" | { "service": { kind: string, } }; export type InternalWorkerRef = { session_id: string, name: string, parent_session_id?: string | null, kind: InternalWorkerKind, }; -export type InternalWorkerSnapshot = { worker: InternalWorkerRef, revision: number, entries: Array, status: WorkerStatus, error?: string | null, in_flight?: InFlightSnapshot, internal_workers?: Array, }; +export type InternalWorkerSnapshot = { worker: InternalWorkerRef, revision: number, session: SessionSnapshot, status: WorkerStatus, error?: string | null, in_flight?: InFlightSnapshot, internal_workers?: Array, }; export type Greeting = { worker_name: string, cwd: string, provider: string, model: string, scope_summary: string, tools: Array, /** @@ -193,7 +217,7 @@ summary: string, * Full tool output. Absent when the tool chose to return * summary-only, or when the result was pruned. */ -output?: string | null, disposition?: ToolResultDisposition | null, is_error: boolean, } } | { "event": "usage", "data": { input_tokens: number | null, output_tokens: number | null, cache_read_input_tokens?: number | null, } } | { "event": "run_end", "data": { result: RunResult, } } | { "event": "error", "data": { code: ErrorCode, message: string, } } | { "event": "snapshot", "data": { entries: Array, greeting: Greeting, status: WorkerStatus, +output?: string | null, disposition?: ToolResultDisposition | null, is_error: boolean, } } | { "event": "usage", "data": { input_tokens: number | null, output_tokens: number | null, cache_read_input_tokens?: number | null, } } | { "event": "run_end", "data": { result: RunResult, } } | { "event": "error", "data": { code: ErrorCode, message: string, } } | { "event": "snapshot", "data": { session: SessionSnapshot, greeting: Greeting, status: WorkerStatus, /** * Unfinished model output that has already streamed in the current * run but is not yet represented by committed snapshot entries. @@ -203,4 +227,4 @@ in_flight?: InFlightSnapshot, * Parent-owned Internal Worker sessions visible to this client. * Service-private Internal Workers are deliberately excluded. */ -internal_workers?: Array, } } | { "event": "internal_worker", "data": { worker: InternalWorkerRef, revision: number, event: Event, } } | { "event": "internal_worker_removed", "data": { worker: InternalWorkerRef, revision: number, } } | { "event": "segment_rotated", "data": { entry: unknown, } } | { "event": "status", "data": { status: WorkerStatus, } } | { "event": "command", "data": { event: CommandEvent, } } | { "event": "completions", "data": { kind: CompletionKind, entries: Array, } } | { "event": "rewind_targets", "data": { head_entries: number, targets: Array, } } | { "event": "rewind_applied", "data": { entries: Array, input: Array, summary: RewindSummary, } } | { "event": "workers_listed", "data": { workers: unknown, } } | { "event": "worker_restored", "data": { result: unknown, } } | { "event": "peer_registered", "data": { result: unknown, } } | { "event": "alert", "data": Alert } | { "event": "memory_worker", "data": MemoryWorkerEvent } | { "event": "compact_start", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_done", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_failed", "data": { lifecycle: CompactionLifecycle, } } | { "event": "shutdown" }; +internal_workers?: Array, } } | { "event": "internal_worker", "data": { worker: InternalWorkerRef, revision: number, event: Event, } } | { "event": "internal_worker_removed", "data": { worker: InternalWorkerRef, revision: number, } } | { "event": "segment_rotated", "data": { session: SessionSnapshot, } } | { "event": "status", "data": { status: WorkerStatus, } } | { "event": "command", "data": { event: CommandEvent, } } | { "event": "completions", "data": { kind: CompletionKind, entries: Array, } } | { "event": "rewind_targets", "data": { head_entries: number, targets: Array, } } | { "event": "rewind_applied", "data": { session: SessionSnapshot, input: Array, summary: RewindSummary, } } | { "event": "workers_listed", "data": { workers: unknown, } } | { "event": "worker_restored", "data": { result: unknown, } } | { "event": "peer_registered", "data": { result: unknown, } } | { "event": "alert", "data": Alert } | { "event": "memory_worker", "data": MemoryWorkerEvent } | { "event": "compact_start", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_done", "data": { lifecycle: CompactionLifecycle, } } | { "event": "compact_failed", "data": { lifecycle: CompactionLifecycle, } } | { "event": "shutdown" }; diff --git a/web/workspace/src/lib/workspace/console/model.test.ts b/web/workspace/src/lib/workspace/console/model.test.ts index 39732416..418b3d58 100644 --- a/web/workspace/src/lib/workspace/console/model.test.ts +++ b/web/workspace/src/lib/workspace/console/model.test.ts @@ -43,11 +43,84 @@ function consoleLine(id: string, kind: ConsoleLine["kind"]): ConsoleLine { }; } +function canonicalSession(logEntries: unknown[]): Event extends { + event: "snapshot"; + data: infer D; +} ? D extends { session: infer S } ? S : never : never { + const entries: Record[] = []; + let sequence = 0; + const itemEntry = (item: Record) => { + const kind = item["kind"]; + if (kind === "reasoning" || item["role"] === "system") return; + entries.push({ + entry_id: `legacy-test-${sequence++}`, + provenance: "legacy_unknown", + ...item, + }); + }; + for (const raw of logEntries) { + if (!raw || typeof raw !== "object" || Array.isArray(raw)) continue; + const entry = raw as Record; + switch (entry["kind"]) { + case "segment_start": + case "annotated_segment_start": + entries.length = 0; + for (const history of Array.isArray(entry["history"]) ? entry["history"] : []) { + const value = history as Record; + itemEntry((value["item"] as Record | undefined) ?? value); + } + break; + case "user_input": + case "annotated_user_input": + entries.push({ + entry_id: `legacy-test-${sequence++}`, + provenance: "legacy_unknown", + kind: "user_input", + segments: entry["segments"] ?? [], + }); + break; + case "assistant_item": + case "tool_result": + case "annotated_assistant_item": + case "annotated_tool_result": { + const annotated = entry["entry"] as Record | undefined; + itemEntry((annotated?.["item"] as Record | undefined) ?? + (entry["item"] as Record)); + break; + } + case "system_item": + case "annotated_system_item": { + const annotated = entry["entry"] as Record | undefined; + const item = (annotated?.["item"] as Record | undefined) ?? + (entry["item"] as Record); + entries.push({ + entry_id: `legacy-test-${sequence++}`, + provenance: "legacy_unknown", + kind: "system_item", + item_kind: item?.["kind"] ?? "system_item", + content: item?.["body"] ?? item?.["message"] ?? "", + data: item, + }); + break; + } + case "run_errored": + entries.push({ + entry_id: `legacy-test-${sequence++}`, + provenance: "legacy_unknown", + kind: "run_error", + message: entry["message"] ?? "Worker run failed.", + }); + break; + } + } + return { entries } as never; +} + function snapshotEvent(cwd: string, entries: unknown[] = []): Event { return { event: "snapshot", data: { - entries, + session: canonicalSession(entries), greeting: { worker_name: "Worker", cwd, @@ -138,7 +211,7 @@ Deno.test("segment rotation retains a live error beside the real SegmentStart hi event: { event: "segment_rotated", data: { - entry: { + session: canonicalSession([{ kind: "segment_start", ts: 5, session_id: "session-1", @@ -149,7 +222,7 @@ Deno.test("segment rotation retains a live error beside the real SegmentStart hi role: "user", content: [{ kind: "text", text: "retained conversation" }], }], - }, + }]), }, } satisfies Event, }, @@ -852,7 +925,7 @@ Deno.test("compaction service activity stays nested in one lifecycle item", () = ); }); -Deno.test("snapshot normalizes orphaned running compaction to interrupted", () => { +Deno.test("snapshot excludes storage-only compaction extension records", () => { const projection = projectConsole([{ eventId: "snapshot", observedAtMs: 9_000, @@ -876,10 +949,7 @@ Deno.test("snapshot normalizes orphaned running compaction to interrupted", () = }]), }]); - assertEquals(projection.lines.length, 1); - assertEquals(projection.lines[0].compaction?.state, "interrupted"); - assertEquals(projection.lines[0].compaction?.endedAtMs, 9_000); - assertEquals(projection.lines[0].streaming, false); + assertEquals(projection.lines.length, 0); }); Deno.test("createConsoleProjector ignores stale compaction revisions", () => { @@ -1238,7 +1308,7 @@ Deno.test("projectConsole renders snapshot entries and in-flight output", () => event: { event: "snapshot", data: { - entries: [ + session: canonicalSession([ { kind: "segment_start", ts: 1, @@ -1300,7 +1370,7 @@ Deno.test("projectConsole renders snapshot entries and in-flight output", () => message: "Compacting…", }, }, - ], + ]), greeting: { worker_name: "Worker", cwd: "/repo", @@ -1332,7 +1402,6 @@ Deno.test("projectConsole renders snapshot entries and in-flight output", () => "user:new user:false", "assistant:assistant reply:false", "tool:Read(1 file)\n /tmp/a.md:false", - "status:Compacting…:true", "in_flight:partial:true", ], ); @@ -1344,7 +1413,7 @@ Deno.test("projectConsole restores system items from snapshot entries", () => { event: { event: "snapshot", data: { - entries: [{ + session: canonicalSession([{ kind: "system_item", ts: 1, item: { @@ -1352,7 +1421,7 @@ Deno.test("projectConsole restores system items from snapshot entries", () => { message: "Worker completed", body: "Child Worker coder-1 completed.", }, - }], + }]), greeting: { worker_name: "Worker", cwd: "/repo", @@ -1388,7 +1457,7 @@ Deno.test("projectConsole reseeds visible rows from segment rotation", () => { event: { event: "segment_rotated", data: { - entry: { + session: canonicalSession([{ kind: "segment_start", ts: 10, session_id: "00000000-0000-0000-0000-000000000001", @@ -1401,7 +1470,7 @@ Deno.test("projectConsole reseeds visible rows from segment rotation", () => { content: [{ kind: "text", text: "after rotation seed" }], }, ], - }, + }]), }, } satisfies Event, }, @@ -1773,7 +1842,7 @@ Deno.test("parent snapshot authoritatively replaces Internal Worker projections" kind: "sub_worker", }, revision: 4, - entries: [{ + session: canonicalSession([{ kind: "assistant_item", ts: 1, item: { @@ -1792,7 +1861,7 @@ Deno.test("parent snapshot authoritatively replaces Internal Worker projections" content: "content", is_error: false, }, - }], + }]), status: "idle", in_flight: { blocks: [{ @@ -1934,18 +2003,20 @@ Deno.test("snapshot restores TaskStore state from system history", () => { `[Session TaskStore snapshot]\n\n\`\`\`json\n{\n "tasks": [{"taskid": 3, "status": "pending", "subject": "Restored", "description": "From compaction"}]\n}\n\`\`\``; const event = snapshotEvent("/repo"); if (event.event !== "snapshot") throw new Error("snapshot fixture expected"); - event.data.entries = [{ - kind: "segment_start", - ts: 1, - session_id: "00000000-0000-0000-0000-000000000001", - system_prompt: null, - config: {}, - history: [{ - kind: "message", - role: "system", - content: [{ kind: "text", text: taskSnapshot }], + event.data.session = { + entries: [{ + entry_id: "task-reminder-1", + provenance: "backend_instruction", + kind: "system_item", + item_kind: "task_reminder", + content: taskSnapshot, + data: { + kind: "task_reminder", + body: taskSnapshot, + source: "automatic", + }, }], - }]; + }; const projection = projectConsole([{ eventId: "task-snapshot", event }]); assertEquals(projection.tasks, [{ diff --git a/web/workspace/src/lib/workspace/console/model.ts b/web/workspace/src/lib/workspace/console/model.ts index d9e56953..61f34c78 100644 --- a/web/workspace/src/lib/workspace/console/model.ts +++ b/web/workspace/src/lib/workspace/console/model.ts @@ -650,9 +650,9 @@ function projectInternalWorkerSnapshot( eventId: string, cwd: string | null, ): InternalWorkerProjection { - const console = snapshotProjectionFromEntries( + const console = snapshotProjectionFromSession( `${eventId}:internal:${snapshot.worker.session_id}:snapshot`, - snapshot.entries, + snapshot.session, cwd, ); console.status = snapshot.status; @@ -905,9 +905,9 @@ export function applyProtocolEvent( case "snapshot": { next.status = event.data.status; next.cwd = event.data.greeting.cwd; - const snapshot = snapshotProjectionFromEntries( + const snapshot = snapshotProjectionFromSession( envelope.eventId, - event.data.entries, + event.data.session, next.cwd, ); next.lines = snapshot.lines; @@ -1008,9 +1008,9 @@ export function applyProtocolEvent( break; case "segment_rotated": { const retainedErrors = next.lines.filter((line) => line.kind === "error"); - const segment = snapshotProjectionFromEntries( + const segment = snapshotProjectionFromSession( envelope.eventId, - [event.data.entry], + event.data.session, next.cwd, ); next.lines = [...segment.lines, ...retainedErrors]; @@ -1925,9 +1925,9 @@ function applyTaskSystemItem( if (typeof body === "string") applyTaskSnapshot(projection, body); } -function snapshotProjectionFromEntries( +function snapshotProjectionFromSession( eventId: string, - entries: unknown[], + snapshot: unknown, cwd: string | null, ): ConsoleProjection { const projection: ConsoleProjection = { @@ -1942,58 +1942,74 @@ function snapshotProjectionFromEntries( internalWorkers: [], removedInternalWorkers: {}, }; + const entries = isRecord(snapshot) ? arrayField(snapshot, "entries") : []; entries.forEach((entry, index) => - applyLogEntry(projection, `${eventId}-snapshot-${index}`, entry) + applySessionEntry(projection, `${eventId}-snapshot-${index}`, entry) ); return projection; } -function applyLogEntry( +function applySessionEntry( projection: ConsoleProjection, - eventId: string, - entry: unknown, + fallbackEventId: string, + value: unknown, ): void { - if (!isRecord(entry)) return; - switch (stringField(entry, "kind")) { - case "segment_start": - arrayField(entry, "history").forEach((item, index) => - applyLoggedItem(projection, `${eventId}-history-${index}`, item) - ); - break; + if (!isRecord(value)) return; + const eventId = stringField(value, "entry_id") ?? fallbackEventId; + switch (stringField(value, "kind")) { case "user_input": projection.lines.push( line( eventId, "user", "User", - segmentsToText(arrayField(entry, "segments") as Segment[]), + segmentsToText(arrayField(value, "segments") as Segment[]), ), ); break; - case "system_item": - projection.lines.push(systemItemLine(eventId, entry["item"])); - applyTaskSystemItem(projection, entry["item"]); + case "message": + applyLoggedItem(projection, eventId, { + kind: "message", + role: value["role"], + content: value["content"], + }); + break; + case "tool_call": + applyLoggedItem(projection, eventId, { + kind: "tool_call", + call_id: value["call_id"], + name: value["name"], + arguments: value["arguments"], + }); break; - case "assistant_item": case "tool_result": - applyLoggedItem(projection, eventId, entry["item"]); + applyLoggedItem(projection, eventId, { + kind: "tool_result", + call_id: value["call_id"], + summary: value["summary"], + content: value["content"], + is_error: value["is_error"], + }); break; - case "run_errored": + case "system_item": { + const item = value["data"]; + projection.lines.push(systemItemLine(eventId, item ?? value)); + applyTaskSystemItem(projection, item); + break; + } + case "run_error": projection.lines.push( line( eventId, "error", "Run error", - stringField(entry, "message") ?? "Worker run failed.", + stringField(value, "message") ?? "Worker run failed.", undefined, false, true, ), ); break; - case "extension": - applyExtensionEntry(projection, eventId, entry); - break; default: break; }