diff --git a/design-system/packages/ui/README.md b/design-system/packages/ui/README.md index 3aa189fa36..2fcb300e38 100644 --- a/design-system/packages/ui/README.md +++ b/design-system/packages/ui/README.md @@ -625,6 +625,8 @@ Compact tabs use `size="sm"` (30px, 14px icons, 4px icon gap); standard tabs ret Dialog titles use 24px bold type with their own 29px line box and normal tracking. `DialogHeader` and `DialogFooter` omit separators by default; pass `separator` for a deliberate divider. A direct `DialogBody` sibling of `DialogFooter appearance="floating"` owns the trailing scroll inset automatically. The floating footer provides the 68px centered action area and a masked blur/gradient using the current theme surface; reduced transparency and forced colors use an opaque fallback. Keep scrollable form content inside `DialogBody` instead of adding a second viewport with independent footer spacing. +`ConfirmDialog secondaryActionPlacement="start"` presents an alternative action as a text button at the start of the footer, with cancel and confirm grouped at the end. The groups wrap when space is limited, keeping cancel and confirm together. The default `inline` placement retains the standard action order and outline secondary button. + Extra-large (`xl`) dialogs have an 800px maximum width and continue shrinking within the viewport gutter. Provider editing uses the floating footer; small workspace creation retains its attached footer and existing button/input sizes. The Lab workspace pattern uses local sample paths and callbacks only. Keep `Dialog` and `Sheet` mounted and set `open={false}` to close them. They retain the last committed children during the exit animation, with interaction disabled, so clearing an owner selection does not collapse the surface. Reopening uses the latest children and cancels the pending exit. Owners that conditionally mount an editor can remove it in `onExitComplete`, which runs once after the surface unmounts. diff --git a/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.meta.ts b/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.meta.ts index 03413e61d8..3bf26da2f9 100644 --- a/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.meta.ts +++ b/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.meta.ts @@ -12,6 +12,7 @@ export const confirmDialogMeta = { { defaultValue: "warning", name: "type", type: "info | warning | error | success" }, { name: "confirmText", type: "ReactNode" }, { name: "secondaryText", type: "ReactNode" }, + { defaultValue: "inline", name: "secondaryActionPlacement", type: "inline | start" }, { name: "cancelText", type: "ReactNode" }, { name: "preview", type: "ReactNode" }, { defaultValue: "false", name: "confirmDanger", type: "boolean" }, @@ -42,5 +43,8 @@ export const confirmDialogMeta = { "layout.confirmDialog.previewRadius", "type.body.md.fontSize", "type.code.sm.fontSize", + "space.2", + "space.3", + "space.4", ], } as const satisfies ComponentMeta; diff --git a/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.module.css b/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.module.css index 16eb628a5f..ecb9190e30 100644 --- a/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.module.css +++ b/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.module.css @@ -3,7 +3,9 @@ .messageRow, .icon, .message, - .preview { + .preview, + .splitActions, + .actionGroup { box-sizing: border-box; } @@ -93,6 +95,25 @@ overflow-wrap: anywhere; } + .splitActions { + display: flex; + min-inline-size: 0; + inline-size: 100%; + flex-wrap: wrap; + align-items: center; + justify-content: space-between; + column-gap: var(--openbitfun-space-4); + row-gap: var(--openbitfun-space-3); + } + + .actionGroup { + display: flex; + flex: 0 0 auto; + align-items: center; + gap: var(--openbitfun-space-2); + margin-inline-start: auto; + } + @media (forced-colors: active) { .icon { color: CanvasText; diff --git a/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.tsx b/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.tsx index 9161e63b44..6cdcda29ae 100644 --- a/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.tsx +++ b/design-system/packages/ui/src/components/ConfirmDialog/ConfirmDialog.tsx @@ -44,6 +44,8 @@ export interface ConfirmDialogProps { open: boolean; pendingAction?: "confirm" | "secondary" | null; preview?: ReactNode; + /** Place an alternative action at the start, separate from cancel and confirm. */ + secondaryActionPlacement?: "inline" | "start"; secondaryText?: ReactNode; showCancel?: boolean; showCloseButton?: boolean; @@ -75,6 +77,7 @@ export const ConfirmDialog = forwardRef( open, pendingAction: controlledPendingAction, preview, + secondaryActionPlacement = "inline", secondaryText, showCancel = true, showCloseButton = false, @@ -91,6 +94,8 @@ export const ConfirmDialog = forwardRef( const resolvedIcon = icon === false ? null : icon ?? defaultIcons[type]; const hasMessage = message !== undefined && message !== null && message !== ""; const hasPreview = preview !== undefined && preview !== null && preview !== ""; + const hasSecondary = secondaryText !== undefined && secondaryText !== null; + const secondaryAtStart = hasSecondary && secondaryActionPlacement === "start"; const resolvedCancelText = cancelText ?? designSystem.messages.confirmCancel; const resolvedConfirmText = confirmText ?? designSystem.messages.confirmAction; @@ -125,6 +130,35 @@ export const ConfirmDialog = forwardRef( onOpenChange(false, reason); }, [busy, onOpenChange]); + const cancelButton = showCancel ? ( + + ) : null; + const secondaryButton = hasSecondary ? ( + + ) : null; + const confirmButton = ( + + ); + return ( ( ) : null} - {showCancel ? ( - - ) : null} - {secondaryText !== undefined && secondaryText !== null ? ( - - ) : null} - + {secondaryButton} +
+ {cancelButton} + {confirmButton} +
+ + ) : ( + <> + {cancelButton} + {secondaryButton} + {confirmButton} + + )}
); diff --git a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs index 21d40caa0f..cfd16dd336 100644 --- a/src/crates/assembly/core/src/agentic/coordination/coordinator.rs +++ b/src/crates/assembly/core/src/agentic/coordination/coordinator.rs @@ -11311,7 +11311,7 @@ Update the persona files and delete BOOTSTRAP.md as soon as bootstrap is complet snapshot.parent_agent_type.clone(), snapshot.build_child_session_config(None), Some(format!("session-{}", snapshot.parent_session_id)), - SessionKind::Standard, + SessionKind::EphemeralChild, ) .await?; self.session_manager @@ -21722,7 +21722,7 @@ mod tests { } #[tokio::test] - async fn btw_session_persists_relationship_and_seeds_forked_listing_baselines() { + async fn btw_session_is_ephemeral_and_seeds_forked_listing_baselines() { let (coordinator, session_manager) = test_persistent_coordinator(); // The parent lives in a registered remote workspace; the child must // inherit that record's SSH facts rather than transport hints. @@ -21811,7 +21811,7 @@ mod tests { assert_eq!( child_session.kind, - crate::agentic::core::SessionKind::Standard + crate::agentic::core::SessionKind::EphemeralChild ); assert_eq!( child_session.last_user_dialog_agent_type.as_deref(), @@ -21871,26 +21871,27 @@ mod tests { let metadata = session_manager .load_session_metadata(&session_storage_path, &child_session.session_id) .await - .expect("BTW metadata should load") - .expect("BTW metadata should exist"); - let relationship = metadata - .relationship - .expect("BTW relationship should persist"); - assert_eq!(relationship.kind, Some(SessionRelationshipKind::Btw)); - assert_eq!( - relationship.parent_session_id.as_deref(), - Some(parent_session.session_id.as_str()) - ); - assert_eq!( - relationship.parent_request_id.as_deref(), - Some("btw-request") - ); - assert_eq!( - relationship.parent_dialog_turn_id.as_deref(), - Some("parent-turn") - ); - assert_eq!(relationship.parent_turn_index, Some(2)); - assert_eq!(metadata.memory_mode, SessionMemoryMode::Disabled); + .expect("BTW metadata lookup should succeed"); + assert!(metadata.is_none(), "temporary BTW must not write history"); + assert!(!session_manager.should_persist_session_id(&child_session.session_id)); + + let reused = coordinator + .ensure_btw_session( + &parent_session.session_id, + &child_session.session_id, + None, + "btw-follow-up", + Some("parent-turn"), + Some(2), + ) + .await + .expect("temporary BTW should accept follow-up questions"); + assert_eq!(reused.kind, SessionKind::EphemeralChild); + assert!(session_manager + .load_session_metadata(&session_storage_path, &child_session.session_id) + .await + .expect("BTW metadata lookup after reuse should succeed") + .is_none()); } #[test] diff --git a/src/crates/assembly/core/src/agentic/persistence/manager.rs b/src/crates/assembly/core/src/agentic/persistence/manager.rs index 5406875a35..de4aafaa42 100644 --- a/src/crates/assembly/core/src/agentic/persistence/manager.rs +++ b/src/crates/assembly/core/src/agentic/persistence/manager.rs @@ -1124,7 +1124,7 @@ impl PersistenceManager { } } - async fn build_session_metadata( + pub(super) async fn build_session_metadata( &self, workspace_path: &Path, session: &Session, diff --git a/src/crates/assembly/core/src/agentic/persistence/session_branch.rs b/src/crates/assembly/core/src/agentic/persistence/session_branch.rs index 6f95dc6ac9..283e24113d 100644 --- a/src/crates/assembly/core/src/agentic/persistence/session_branch.rs +++ b/src/crates/assembly/core/src/agentic/persistence/session_branch.rs @@ -1,5 +1,8 @@ use super::manager::PersistenceManager; -use crate::agentic::core::{MessageContent, Session, SessionKind}; +use crate::agentic::core::{Message, MessageContent, Session, SessionKind}; +use crate::agentic::session::{EvidenceLedgerEvent, SessionPromptCache}; +use crate::agentic::skill_agent_snapshot::TurnSkillAgentSnapshot; +use crate::service::session::DialogTurnData; use crate::util::errors::{OpenBitFunError, OpenBitFunResult}; use openbitfun_services_core::session::SessionBranchBoundary; use openbitfun_services_core::session::{ @@ -10,6 +13,17 @@ pub use openbitfun_services_core::session::{SessionBranchRequest, SessionBranchR use std::path::Path; use std::time::{SystemTime, UNIX_EPOCH}; +/// Runtime-owned history; saving a fork never persists its temporary source. +pub(crate) struct TransientSessionBranchSource { + pub session: Session, + pub turns: Vec, + pub context_snapshots: Vec<(usize, Vec)>, + pub prompt_cache: Option, + pub skill_agent_snapshots: Vec<(usize, TurnSkillAgentSnapshot)>, + pub baseline_override: Option, + pub evidence_events: Vec, +} + fn clear_inherited_snapshot_capability(result: &mut serde_json::Value) { if let Some(result) = result.as_object_mut() { // File snapshots belong to the source Session and are not forked with @@ -24,6 +38,26 @@ impl PersistenceManager { &self, workspace_path: &Path, request: &SessionBranchRequest, + ) -> OpenBitFunResult { + self.branch_session_with_transient_source(workspace_path, request, None) + .await + } + + pub(crate) async fn branch_transient_session( + &self, + workspace_path: &Path, + request: &SessionBranchRequest, + source: TransientSessionBranchSource, + ) -> OpenBitFunResult { + self.branch_session_with_transient_source(workspace_path, request, Some(source)) + .await + } + + async fn branch_session_with_transient_source( + &self, + workspace_path: &Path, + request: &SessionBranchRequest, + transient_source: Option, ) -> OpenBitFunResult { openbitfun_core_types::validate_session_id(&request.source_session_id) .map_err(OpenBitFunError::Validation)?; @@ -32,18 +66,32 @@ impl PersistenceManager { .await; let _branch_allocation_guard = branch_allocation_lock.lock().await; - let source_session = self - .load_session(workspace_path, &request.source_session_id) - .await?; - let source_metadata = self - .load_session_metadata(workspace_path, &request.source_session_id) - .await? - .ok_or_else(|| { - OpenBitFunError::NotFound(format!( - "Source session metadata not found: {}", - request.source_session_id - )) - })?; + let (source_session, source_metadata) = if let Some(source) = transient_source.as_ref() { + if source.session.session_id != request.source_session_id { + return Err(OpenBitFunError::Validation( + "Transient fork source identity mismatch".to_string(), + )); + } + ( + source.session.clone(), + self.build_session_metadata(workspace_path, &source.session, None) + .await, + ) + } else { + let session = self + .load_session(workspace_path, &request.source_session_id) + .await?; + let metadata = self + .load_session_metadata(workspace_path, &request.source_session_id) + .await? + .ok_or_else(|| { + OpenBitFunError::NotFound(format!( + "Source session metadata not found: {}", + request.source_session_id + )) + })?; + (session, metadata) + }; let metadata_list = self .list_session_metadata_including_internal(workspace_path) .await?; @@ -52,12 +100,18 @@ impl PersistenceManager { &source_session.session_name, &metadata_list, ); - let source_turns = self - .load_session_turns(workspace_path, &request.source_session_id) - .await?; - let source_prompt_cache = self - .load_prompt_cache(workspace_path, &request.source_session_id) - .await?; + let source_turns = if let Some(source) = transient_source.as_ref() { + source.turns.clone() + } else { + self.load_session_turns(workspace_path, &request.source_session_id) + .await? + }; + let source_prompt_cache = if let Some(source) = transient_source.as_ref() { + source.prompt_cache.clone() + } else { + self.load_prompt_cache(workspace_path, &request.source_session_id) + .await? + }; if source_turns.is_empty() { return Err(OpenBitFunError::Validation( @@ -127,14 +181,21 @@ impl PersistenceManager { for (new_index, source_turn) in source_turns.iter().take(copied_turn_count).enumerate() { - if let Some(mut messages) = self - .load_turn_context_snapshot( + let context_snapshot = if let Some(source) = transient_source.as_ref() { + source + .context_snapshots + .iter() + .find(|(index, _)| *index == source_turn.turn_index) + .map(|(_, messages)| messages.clone()) + } else { + self.load_turn_context_snapshot( workspace_path, &request.source_session_id, source_turn.turn_index, ) .await? - { + }; + if let Some(mut messages) = context_snapshot { for message in &mut messages { if let MessageContent::ToolResult { result, .. } = &mut message.content { clear_inherited_snapshot_capability(result); @@ -148,14 +209,21 @@ impl PersistenceManager { ) .await?; } - if let Some(snapshot) = self - .load_turn_skill_agent_snapshot( + let skill_snapshot = if let Some(source) = transient_source.as_ref() { + source + .skill_agent_snapshots + .iter() + .find(|(index, _)| *index == source_turn.turn_index) + .map(|(_, snapshot)| snapshot.clone()) + } else { + self.load_turn_skill_agent_snapshot( workspace_path, &request.source_session_id, source_turn.turn_index, ) .await? - { + }; + if let Some(snapshot) = skill_snapshot { self.save_turn_skill_agent_snapshot( workspace_path, &target_session_id, @@ -170,7 +238,10 @@ impl PersistenceManager { self.save_dialog_turn(workspace_path, turn).await?; } - if let Some(last_copied_turn_index) = copied_turn_count.checked_sub(1) { + if let Some(last_copied_turn_index) = copied_turn_count + .checked_sub(1) + .filter(|_| transient_source.is_none()) + { self.copy_compression_transcripts_through( workspace_path, &request.source_session_id, @@ -184,13 +255,16 @@ impl PersistenceManager { self.save_prompt_cache(workspace_path, &target_session_id, cache) .await?; } - if let Some(snapshot) = self - .load_skill_agent_baseline_override_snapshot( + let baseline_snapshot = if let Some(source) = transient_source.as_ref() { + source.baseline_override.clone() + } else { + self.load_skill_agent_baseline_override_snapshot( workspace_path, &request.source_session_id, ) .await? - { + }; + if let Some(snapshot) = baseline_snapshot { self.save_skill_agent_baseline_override_snapshot( workspace_path, &target_session_id, @@ -202,9 +276,12 @@ impl PersistenceManager { // Copy evidence ledger events for the branched turns, rewriting // session_id to the target session so the fork inherits // checkpoints, failed commands, and partial subagent results. - let source_evidence_events = self - .load_evidence_ledger_events(workspace_path, &request.source_session_id) - .await?; + let source_evidence_events = if let Some(source) = transient_source.as_ref() { + source.evidence_events.clone() + } else { + self.load_evidence_ledger_events(workspace_path, &request.source_session_id) + .await? + }; if !source_evidence_events.is_empty() { let copied_turn_ids: std::collections::HashSet = branched_turns .iter() diff --git a/src/crates/assembly/core/src/agentic/session/session_manager.rs b/src/crates/assembly/core/src/agentic/session/session_manager.rs index 57cf05ea25..c3b446c5cb 100644 --- a/src/crates/assembly/core/src/agentic/session/session_manager.rs +++ b/src/crates/assembly/core/src/agentic/session/session_manager.rs @@ -13,6 +13,7 @@ use crate::agentic::fork_agent::normalize_incomplete_tool_calls; use crate::agentic::image_analysis::ImageContextData; use crate::agentic::keyed_lock::{KeyedAsyncLock, KeyedAsyncLockGuard}; use crate::agentic::memories::db::{MemoryDatabase, MEMORY_PHASE2_GLOBAL_JOB_KEY}; +use crate::agentic::persistence::session_branch::TransientSessionBranchSource; use crate::agentic::persistence::{MaterializedSessionReferenceTranscript, PersistenceManager}; use crate::agentic::session::revert::SessionRevertPhase; use crate::agentic::session::session_store_port::CoreSessionStorePort; @@ -68,7 +69,7 @@ use openbitfun_services_core::session::{ }; use serde::{Deserialize, Serialize}; use serde_json::json; -use std::collections::{HashSet, VecDeque}; +use std::collections::{BTreeMap, HashSet, VecDeque}; use std::path::{Path, PathBuf}; use std::sync::{Arc, Mutex}; use std::time::Instant; @@ -355,6 +356,9 @@ pub struct SessionManager { /// removed with that Session; they are never serialized into public config. transient_session_ids: Arc>, + /// Authoritative temporary history survives context compression until explicit closure. + transient_turns: Arc>, + /// Recent authoritative terminal results for live Turn settlement callers. /// The bounded cache preserves the exact execution result; persisted Turns /// remain the fallback, while transient Sessions depend on this copy. @@ -412,9 +416,16 @@ pub struct ActiveTurnPermissionMode { pub mode: PermissionMode, } +#[derive(Default)] +struct TransientSessionHistory { + turns: Vec, + context_snapshots: BTreeMap>, +} + fn clear_session_runtime_stores( session_id: &str, context_store: &SessionContextStore, + transient_turns: &DashMap, prompt_cache_store: &SessionPromptCacheStore, token_anchor_store: &TokenAnchorStore, turn_skill_agent_snapshot_store: &TurnSkillAgentSnapshotStore, @@ -423,6 +434,7 @@ fn clear_session_runtime_stores( evidence_ledger: &SessionEvidenceLedger, ) { context_store.delete_session(session_id); + transient_turns.remove(session_id); prompt_cache_store.delete_session(session_id); token_anchor_store.delete_session(session_id); turn_skill_agent_snapshot_store.delete_session(session_id); @@ -518,6 +530,7 @@ impl SessionManager { pub(crate) fn evict_loaded_session_for_test(&self, session_id: &str) { self.sessions.remove(session_id); self.transient_session_ids.remove(session_id); + self.transient_turns.remove(session_id); self.clear_turn_settlement_results(session_id); self.release_active_session_reservation(session_id); self.release_session_write_lock(session_id); @@ -1682,6 +1695,12 @@ impl SessionManager { reason: &str, ) { if !self.should_persist_session_id(session_id) { + if let Some(mut history) = self.transient_turns.get_mut(session_id) { + history.context_snapshots.insert( + turn_index, + self.context_store.get_context_messages(session_id), + ); + } return; } @@ -1991,6 +2010,7 @@ impl SessionManager { sessions: Arc::new(DashMap::new()), active_turn_permission_modes: Arc::new(DashMap::new()), transient_session_ids: Arc::new(DashMap::new()), + transient_turns: Arc::new(DashMap::new()), turn_settlement_results: Arc::new(DashMap::new()), turn_settlement_result_order: Arc::new(Mutex::new(VecDeque::new())), active_session_capacity: Arc::new(Semaphore::new(config.max_active_sessions)), @@ -2493,6 +2513,7 @@ impl SessionManager { let sessions = self.sessions.clone(); let active_turn_permission_modes = self.active_turn_permission_modes.clone(); let transient_session_ids = self.transient_session_ids.clone(); + let transient_turns = self.transient_turns.clone(); let turn_settlement_results = self.turn_settlement_results.clone(); let turn_settlement_result_order = self.turn_settlement_result_order.clone(); let active_session_capacity = self.active_session_capacity.clone(); @@ -2530,6 +2551,7 @@ impl SessionManager { sessions, active_turn_permission_modes, transient_session_ids, + transient_turns, turn_settlement_results, turn_settlement_result_order, active_session_capacity, @@ -4976,6 +4998,7 @@ impl SessionManager { clear_session_runtime_stores( session_id, self.context_store.as_ref(), + self.transient_turns.as_ref(), self.prompt_cache_store.as_ref(), self.token_anchor_store.as_ref(), self.turn_skill_agent_snapshot_store.as_ref(), @@ -5021,6 +5044,7 @@ impl SessionManager { clear_session_runtime_stores( session_id, self.context_store.as_ref(), + self.transient_turns.as_ref(), self.prompt_cache_store.as_ref(), self.token_anchor_store.as_ref(), self.turn_skill_agent_snapshot_store.as_ref(), @@ -6210,6 +6234,7 @@ impl SessionManager { clear_session_runtime_stores( session_id, self.context_store.as_ref(), + self.transient_turns.as_ref(), self.prompt_cache_store.as_ref(), self.token_anchor_store.as_ref(), self.turn_skill_agent_snapshot_store.as_ref(), @@ -7069,28 +7094,27 @@ impl SessionManager { .add_message(session_id, message.with_turn_id(turn_id.clone())); } + let turn_data = DialogTurnData::new_with_kind( + kind, + turn_id.clone(), + turn_index, + session_id.to_string(), + if kind == DialogTurnKind::UserDialog { + agent_type.clone() + } else { + None + }, + UserMessageData { + id: format!("{}-user", turn_id), + content: user_input, + timestamp: SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .as_millis() as u64, + metadata: user_message_metadata, + }, + ); if self.should_persist_session_id(session_id) { - let turn_data = DialogTurnData::new_with_kind( - kind, - turn_id.clone(), - turn_index, - session_id.to_string(), - if kind == DialogTurnKind::UserDialog { - agent_type.clone() - } else { - None - }, - UserMessageData { - id: format!("{}-user", turn_id), - content: user_input, - timestamp: SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .as_millis() as u64, - metadata: user_message_metadata, - }, - ); - // Clone the session data out of the DashMap guard before awaiting I/O. let session_snapshot = self.sessions.get(session_id).map(|s| s.clone()); // Ref guard released -- DashMap shard lock is free. @@ -7102,6 +7126,12 @@ impl SessionManager { self.persistence_manager .save_dialog_turn(&workspace_path, &turn_data) .await?; + } else { + self.transient_turns + .entry(session_id.to_string()) + .or_default() + .turns + .push(turn_data); } self.persist_context_snapshot_for_turn_best_effort(session_id, turn_index, "turn_started") @@ -7527,6 +7557,16 @@ impl SessionManager { self.persistence_manager .save_dialog_turn(&workspace_path, &turn) .await?; + } else { + let mut history = self + .transient_turns + .entry(session_id.to_string()) + .or_default(); + if let Some(current) = history.turns.get_mut(turn_index) { + *current = turn.clone(); + } else { + history.turns.push(turn.clone()); + } } let session_snapshot = if let Some(mut session) = self.sessions.get_mut(session_id) { @@ -7787,7 +7827,116 @@ impl SessionManager { } } - /// Complete dialog turn + /// Read the canonical turn from its durable or temporary owner. + async fn load_runtime_dialog_turn( + &self, + workspace_path: &Path, + session_id: &str, + turn_index: usize, + ) -> OpenBitFunResult> { + if self.should_persist_session_id(session_id) { + self.persistence_manager + .load_dialog_turn(workspace_path, session_id, turn_index) + .await + } else { + Ok(self + .transient_turns + .get(session_id) + .and_then(|history| history.turns.get(turn_index).cloned())) + } + } + + async fn save_runtime_dialog_turn( + &self, + workspace_path: &Path, + turn: &DialogTurnData, + ) -> OpenBitFunResult<()> { + if self.should_persist_session_id(&turn.session_id) { + self.persistence_manager + .save_dialog_turn(workspace_path, turn) + .await + } else { + let mut history = self + .transient_turns + .get_mut(&turn.session_id) + .ok_or_else(|| { + OpenBitFunError::NotFound(format!( + "Temporary session history not found: {}", + turn.session_id + )) + })?; + let current = history.turns.get_mut(turn.turn_index).ok_or_else(|| { + OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn.turn_id)) + })?; + *current = turn.clone(); + Ok(()) + } + } + + /// Called under the source Session mutation lock after execution has settled. + pub(crate) fn transient_branch_source_locked( + &self, + session_id: &str, + ) -> OpenBitFunResult { + let session = self + .get_session(session_id) + .ok_or_else(|| OpenBitFunError::NotFound(format!("Session not found: {session_id}")))?; + if !matches!( + session.state, + SessionState::Idle | SessionState::Error { .. } + ) { + return Err(OpenBitFunError::Validation( + "Stop the temporary conversation before saving it".to_string(), + )); + } + let history = self.transient_turns.get(session_id).ok_or_else(|| { + OpenBitFunError::Validation( + "Temporary conversation has no submitted questions".to_string(), + ) + })?; + let turns = history.turns.clone(); + let context_snapshots = history + .context_snapshots + .iter() + .map(|(index, messages)| (*index, messages.clone())) + .collect(); + drop(history); + if turns.len() != session.dialog_turn_ids.len() + || turns + .iter() + .any(|turn| turn.status == TurnStatus::InProgress) + { + return Err(OpenBitFunError::Validation( + "Temporary conversation history has not settled".to_string(), + )); + } + let skill_agent_snapshots = turns + .iter() + .filter_map(|turn| { + self.turn_skill_agent_snapshot_store + .get_snapshot(session_id, turn.turn_index) + .map(|snapshot| (turn.turn_index, snapshot)) + }) + .collect(); + let evidence_events = turns + .iter() + .flat_map(|turn| self.evidence_events_for_turn(session_id, &turn.turn_id)) + .collect(); + Ok(TransientSessionBranchSource { + session, + turns, + context_snapshots, + prompt_cache: self.prompt_cache_store.get_cache(session_id), + skill_agent_snapshots, + baseline_override: self + .skill_agent_baseline_override_snapshot_store + .get(session_id) + .map(|snapshot| snapshot.clone()), + evidence_events, + }) + } + + /// Complete a dialog turn in its history owner. pub async fn complete_dialog_turn( &self, session_id: &str, @@ -7798,7 +7947,9 @@ impl SessionManager { finish_reason: Option, has_final_response: Option, ) -> OpenBitFunResult<()> { - if !self.should_persist_session_id(session_id) { + if !self.should_persist_session_id(session_id) + && !self.transient_turns.contains_key(session_id) + { debug!( "Skipping dialog turn persistence for transient session completion: session_id={}, turn_id={}, response_len={}, rounds={}", session_id, @@ -7848,8 +7999,7 @@ impl SessionManager { let _mutation_guard = self.acquire_session_mutation(session_id).await?; let mut turn = self - .persistence_manager - .load_dialog_turn(&workspace_path, session_id, turn_index) + .load_runtime_dialog_turn(&workspace_path, session_id, turn_index) .await? .ok_or_else(|| { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) @@ -7950,11 +8100,8 @@ impl SessionManager { turn.end_time = Some(completion_timestamp); // Persist - if self.should_persist_session_id(session_id) { - self.persistence_manager - .save_dialog_turn(&workspace_path, &turn) - .await?; - } + self.save_runtime_dialog_turn(&workspace_path, &turn) + .await?; debug!( "Dialog turn completed: turn_id={}, rounds={}, tools={}", @@ -8245,7 +8392,9 @@ impl SessionManager { generation_messages: &[Message], ) -> OpenBitFunResult<()> { let _mutation_guard = self.acquire_session_mutation(session_id).await?; - if !self.should_persist_session_id(session_id) { + if !self.should_persist_session_id(session_id) + && !self.transient_turns.contains_key(session_id) + { debug!( "Skipping dialog turn persistence for transient session failure: session_id={}, turn_id={}, error={}", session_id, turn_id, error @@ -8270,8 +8419,7 @@ impl SessionManager { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) })?; let mut turn = self - .persistence_manager - .load_dialog_turn(&workspace_path, session_id, turn_index) + .load_runtime_dialog_turn(&workspace_path, session_id, turn_index) .await? .ok_or_else(|| { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) @@ -8306,11 +8454,8 @@ impl SessionManager { ) .await; } - if self.should_persist_session_id(session_id) { - self.persistence_manager - .save_dialog_turn(&workspace_path, &turn) - .await?; - } + self.save_runtime_dialog_turn(&workspace_path, &turn) + .await?; debug!( "Dialog turn marked as failed: turn_id={}, turn_index={}, error={}", @@ -8341,7 +8486,9 @@ impl SessionManager { generation_messages: &[Message], ) -> OpenBitFunResult<()> { let _mutation_guard = self.acquire_session_mutation(session_id).await?; - if !self.should_persist_session_id(session_id) { + if !self.should_persist_session_id(session_id) + && !self.transient_turns.contains_key(session_id) + { debug!( "Skipping dialog turn persistence for transient session cancellation: session_id={}, turn_id={}", session_id, turn_id @@ -8366,8 +8513,7 @@ impl SessionManager { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) })?; let mut turn = self - .persistence_manager - .load_dialog_turn(&workspace_path, session_id, turn_index) + .load_runtime_dialog_turn(&workspace_path, session_id, turn_index) .await? .ok_or_else(|| { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) @@ -8402,8 +8548,7 @@ impl SessionManager { .await; } - self.persistence_manager - .save_dialog_turn(&workspace_path, &turn) + self.save_runtime_dialog_turn(&workspace_path, &turn) .await?; debug!( @@ -9021,7 +9166,9 @@ impl SessionManager { duration_ms: u64, snapshot_reason: &str, ) -> OpenBitFunResult<()> { - if !self.should_persist_session_id(session_id) { + if !self.should_persist_session_id(session_id) + && !self.transient_turns.contains_key(session_id) + { debug!( "Skipping turn persistence for transient session completion: session_id={}, turn_id={}, rounds={}, duration_ms={}", session_id, @@ -9049,8 +9196,7 @@ impl SessionManager { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) })?; let mut turn = self - .persistence_manager - .load_dialog_turn(&workspace_path, session_id, turn_index) + .load_runtime_dialog_turn(&workspace_path, session_id, turn_index) .await? .ok_or_else(|| { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) @@ -9072,11 +9218,8 @@ impl SessionManager { ) .await; - if self.should_persist_session_id(session_id) { - self.persistence_manager - .save_dialog_turn(&workspace_path, &turn) - .await?; - } + self.save_runtime_dialog_turn(&workspace_path, &turn) + .await?; Ok(()) } @@ -9124,7 +9267,9 @@ impl SessionManager { model_rounds: Vec, snapshot_reason: &str, ) -> OpenBitFunResult<()> { - if !self.should_persist_session_id(session_id) { + if !self.should_persist_session_id(session_id) + && !self.transient_turns.contains_key(session_id) + { debug!( "Skipping turn persistence for transient session failure: session_id={}, turn_id={}, rounds={}, error={}", session_id, @@ -9152,8 +9297,7 @@ impl SessionManager { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) })?; let mut turn = self - .persistence_manager - .load_dialog_turn(&workspace_path, session_id, turn_index) + .load_runtime_dialog_turn(&workspace_path, session_id, turn_index) .await? .ok_or_else(|| { OpenBitFunError::NotFound(format!("Dialog turn not found: {}", turn_id)) @@ -9178,11 +9322,8 @@ impl SessionManager { ) .await; - if self.should_persist_session_id(session_id) { - self.persistence_manager - .save_dialog_turn(&workspace_path, &turn) - .await?; - } + self.save_runtime_dialog_turn(&workspace_path, &turn) + .await?; debug!( "Turn marked as failed: turn_id={}, turn_index={}, error={}", @@ -9724,6 +9865,7 @@ impl SessionManager { let session_mutation_locks = self.session_mutation_locks.clone(); let session_write_locks = self.session_write_locks.clone(); let context_store = self.context_store.clone(); + let transient_turns = self.transient_turns.clone(); let prompt_cache_store = self.prompt_cache_store.clone(); let token_anchor_store = self.token_anchor_store.clone(); let turn_skill_agent_snapshot_store = self.turn_skill_agent_snapshot_store.clone(); @@ -9823,6 +9965,7 @@ impl SessionManager { clear_session_runtime_stores( &candidate.session_id, context_store.as_ref(), + transient_turns.as_ref(), prompt_cache_store.as_ref(), token_anchor_store.as_ref(), turn_skill_agent_snapshot_store.as_ref(), diff --git a/src/crates/assembly/core/src/product_runtime.rs b/src/crates/assembly/core/src/product_runtime.rs index 357de3d7d3..44af140a86 100644 --- a/src/crates/assembly/core/src/product_runtime.rs +++ b/src/crates/assembly/core/src/product_runtime.rs @@ -2166,7 +2166,7 @@ impl CoreSessionOperationsPort { )); } let session_manager = self.coordinator.get_session_manager(); - let expected_workspace = { + let (expected_workspace, is_transient) = { let _guard = session_manager .acquire_session_mutation(&source_session_id) .await @@ -2174,16 +2174,23 @@ impl CoreSessionOperationsPort { session_manager .validate_session_storage_path_binding(&source_session_id, storage_path) .map_err(runtime_port_error)?; - let persisted = self - .persistence - .load_session_header(storage_path, &source_session_id) - .await - .map_err(runtime_port_error)?; - let binding = fork_workspace_binding(&persisted, workspace_id) + let loaded = session_manager.get_session(&source_session_id); + let is_transient = loaded.as_ref().is_some_and(|session| { + session.kind == crate::agentic::core::SessionKind::EphemeralChild + }); + let source = match loaded.as_ref() { + Some(session) if is_transient => session.clone(), + _ => self + .persistence + .load_session_header(storage_path, &source_session_id) + .await + .map_err(runtime_port_error)?, + }; + let binding = fork_workspace_binding(&source, workspace_id) .await .map_err(runtime_port_error)?; let execution_id = binding.workspace_id.as_deref().expect("resolved binding"); - if let Some(loaded) = session_manager.get_session(&source_session_id) { + if let Some(loaded) = loaded { fork_workspace_binding(&loaded, execution_id) .await .map_err(runtime_port_error)?; @@ -2199,8 +2206,43 @@ impl CoreSessionOperationsPort { }, ) .map_err(runtime_port_error)?; - execution_id.to_owned() + (execution_id.to_owned(), is_transient) }; + if is_transient { + let _guard = session_manager + .acquire_session_mutation(&source_session_id) + .await + .map_err(runtime_port_error)?; + session_manager + .validate_session_storage_path_binding(&source_session_id, storage_path) + .map_err(runtime_port_error)?; + let source = session_manager + .transient_branch_source_locked(&source_session_id) + .map_err(runtime_port_error)?; + fork_workspace_binding(&source.session, &expected_workspace) + .await + .map_err(runtime_port_error)?; + let source_turn_id = match source_turn_id { + Some(id) => id, + None => latest_persisted_turn_id(&source.turns).map_err(runtime_port_error)?, + }; + let result = self + .persistence + .branch_transient_session( + storage_path, + &SessionBranchRequest { + source_session_id: source_session_id.clone(), + source_turn_id, + boundary, + }, + source, + ) + .await + .map_err(runtime_port_error)?; + return self + .finish_session_fork(storage_path, &source_session_id, result) + .await; + } if !session_manager .is_session_loaded_from_storage_path(storage_path, &source_session_id) .map_err(runtime_port_error)? @@ -2293,9 +2335,19 @@ impl CoreSessionOperationsPort { ) .await .map_err(runtime_port_error)?; + self.finish_session_fork(storage_path, &source_session_id_for_coordination, result) + .await + } + + async fn finish_session_fork( + &self, + storage_path: &Path, + source_session_id: &str, + result: crate::agentic::persistence::SessionBranchResult, + ) -> PortResult { if let Err(error) = self .coordinator - .initialize_fork_coordination(&source_session_id_for_coordination, &result.session_id) + .initialize_fork_coordination(source_session_id, &result.session_id) .await { if let Err(cleanup_error) = self @@ -4197,6 +4249,189 @@ mod tests { .len(), 2 ); + + // BTW keeps its authoritative turns in memory, independently of + // model context compression, until the user chooses to save it. + let temporary = session_manager + .create_session_with_id_and_details( + Some(format!("btw-{connection_id}")), + "Temporary question".to_string(), + "Standard".to_string(), + crate::agentic::core::SessionConfig { + workspace_id: Some(record.id.clone()), + workspace_path: Some(remote_path.clone()), + remote_connection_id: Some(connection_id.to_string()), + remote_ssh_host: Some(host.to_string()), + ..Default::default() + }, + Some(format!("session-{source_session_id}")), + SessionKind::EphemeralChild, + ) + .await + .expect("temporary BTW"); + session_manager + .replace_context_messages( + &temporary.session_id, + vec![crate::agentic::core::Message::user( + "Inherited parent context".to_string(), + )], + ) + .await; + for index in 0..3 { + let turn_id = session_manager + .start_dialog_turn( + &temporary.session_id, + "Standard".to_string(), + format!("Question {index}"), + Some(format!("btw-turn-{index}")), + None, + None, + ) + .await + .expect("temporary question"); + let messages = + vec![ + crate::agentic::core::Message::assistant(format!("Answer {index}")) + .with_turn_id(turn_id.clone()) + .with_round_id(format!("btw-round-{index}")), + ]; + session_manager + .add_messages(&temporary.session_id, messages.clone()) + .await + .unwrap(); + match index { + 0 => session_manager + .complete_dialog_turn( + &temporary.session_id, + &turn_id, + format!("Answer {index}"), + &messages, + crate::agentic::core::TurnStats::default(), + Some("stop".to_string()), + Some(true), + ) + .await + .unwrap(), + 1 => session_manager + .cancel_dialog_turn_with_messages( + &temporary.session_id, + &turn_id, + &messages, + ) + .await + .unwrap(), + _ => session_manager + .fail_dialog_turn_with_messages( + &temporary.session_id, + &turn_id, + "test failure".to_string(), + &messages, + ) + .await + .unwrap(), + } + session_manager.reset_session_state_if_processing(&temporary.session_id, &turn_id); + } + let compacted = vec![ + crate::agentic::core::Message::user("Inherited parent context".to_string()), + crate::agentic::core::Message::assistant("Compacted temporary context".to_string()), + ]; + session_manager + .replace_context_messages(&temporary.session_id, compacted.clone()) + .await; + assert!(persistence + .load_session_metadata(&storage_path, &temporary.session_id) + .await + .unwrap() + .is_none()); + let old_boundary = port + .fork_session_at_turn(AgentSessionForkAtTurnRequest { + workspace_id: Some(record.id.clone()), + workspace_path: String::new(), + source_session_id: temporary.session_id.clone(), + source_turn_id: "btw-turn-0".to_string(), + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + .expect("temporary historical forks preserve the matching context snapshot"); + let older_context = persistence + .load_turn_context_snapshot(&storage_path, &old_boundary.session_id, 0) + .await + .unwrap() + .unwrap(); + let older_context = serde_json::to_string(&older_context).unwrap(); + assert!(older_context.contains("Inherited parent context")); + assert!(older_context.contains("Answer 0")); + assert!(!older_context.contains("Compacted temporary context")); + let saved = port + .fork_session_at_turn(AgentSessionForkAtTurnRequest { + workspace_id: Some(record.id.clone()), + workspace_path: String::new(), + source_session_id: temporary.session_id.clone(), + source_turn_id: "btw-turn-2".to_string(), + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + .expect("save temporary BTW through the existing fork port"); + let saved_session = persistence + .load_session(&storage_path, &saved.session_id) + .await + .unwrap(); + assert_eq!(saved_session.kind, SessionKind::Standard); + assert_eq!(saved_session.created_by, None); + assert_eq!( + saved_session.config.remote_connection_id.as_deref(), + Some(connection_id) + ); + assert_eq!(saved_session.config.remote_ssh_host.as_deref(), Some(host)); + let saved_turns = persistence + .load_session_turns(&storage_path, &saved.session_id) + .await + .unwrap(); + assert_eq!(saved_turns.len(), 3); + assert_eq!( + saved_turns + .iter() + .map(|turn| turn.status.clone()) + .collect::>(), + vec![ + TurnStatus::Completed, + TurnStatus::Cancelled, + TurnStatus::Error + ] + ); + for (index, turn) in saved_turns.iter().enumerate() { + assert_eq!(turn.user_message.content, format!("Question {index}")); + assert_eq!( + turn.model_rounds[0].text_items[0].content, + format!("Answer {index}") + ); + } + let saved_context = persistence + .load_turn_context_snapshot(&storage_path, &saved.session_id, 2) + .await + .unwrap() + .unwrap(); + assert_eq!( + serde_json::to_value(saved_context).unwrap(), + serde_json::to_value(compacted).unwrap() + ); + let saved_metadata = persistence + .load_session_metadata(&storage_path, &saved.session_id) + .await + .unwrap() + .unwrap(); + assert_eq!(saved_metadata.relationship, None); + assert_eq!( + session_manager + .get_session(&temporary.session_id) + .unwrap() + .kind, + SessionKind::EphemeralChild + ); + assert!(!storage_path.join(&temporary.session_id).exists()); } assert!(!Path::new(&remote_path).exists()); } diff --git a/src/web-ui/src/app/components/NavPanel/components/PersistentFooterActions.tsx b/src/web-ui/src/app/components/NavPanel/components/PersistentFooterActions.tsx index 4535dedef2..84b86eedfd 100644 --- a/src/web-ui/src/app/components/NavPanel/components/PersistentFooterActions.tsx +++ b/src/web-ui/src/app/components/NavPanel/components/PersistentFooterActions.tsx @@ -242,6 +242,7 @@ const PersistentFooterActions: React.FC = () => { onOpenAppearanceSettings={handleOpenAppearanceSettings} /> + -