diff --git a/docs/architecture/agent-runtime-deployment-design.md b/docs/architecture/agent-runtime-deployment-design.md index bb2d417f9f..1ddc4904ed 100644 --- a/docs/architecture/agent-runtime-deployment-design.md +++ b/docs/architecture/agent-runtime-deployment-design.md @@ -42,8 +42,8 @@ flowchart TB | Session 写入 | BitFun Runtime 的持久化 Session 由 `SessionManager` 管理;同一存储位置中的同一 Session 同时只允许一个本机进程写入,list/view 等只读操作不受影响 | | 当前 HTTP Server | 只提供 health/info/WebSocket 外壳,未装配 Agent Runtime,因此不取得 workspace ownership;`bootstrap.rs` 仅保持 agent-enabled composition 的一致边界,不由当前入口启动 | | Shared local IPC | 未发布的本机协议已有 discovery、实例锁、严格握手、Session 控制权、有界事件流和 cleanup;唯一 consumer 是第一方交互式 TUI adapter | -| Shared TUI | `bitfun --shared` / `bitfun chat --shared` 可列出、创建、恢复 Session,删除未被控制的空闲非当前 Session,重命名当前 Session,读取 transcript,切换当前 Session 的 Agent mode/model,通过 `/reload [skills|instructions]` 刷新声明式上下文,通过 `/compact` 或 `/summarize` 压缩当前 Session 上下文,提交/取消 Turn,处理 Permission 和 UserInput;默认仍是 Embedded | -| Shared GUI/Headless/ACP/SDK Host/Remote | 未交付,也不会由 `--shared` 隐式启用;Replay、Observer、Controller transfer、Session archive/fork 同样不在当前协议中 | +| Shared TUI | `bitfun --shared` / `bitfun chat --shared` 可列出、创建、恢复 Session,删除未被控制的空闲非当前 Session,通过 `/fork` 从完整历史或选中提示词之前创建分支,重命名当前 Session,读取 transcript,切换当前 Session 的 Agent mode/model,通过 `/reload [skills|instructions]` 刷新声明式上下文,通过 `/compact` 或 `/summarize` 压缩当前 Session 上下文,提交/取消 Turn,处理 Permission 和 UserInput;默认仍是 Embedded | +| Shared GUI/Headless/ACP/SDK Host/Remote | 未交付,也不会由 `--shared` 隐式启用;Replay、Observer、通用 Controller transfer 和 Session archive 同样不在当前协议中 | 因此当前交付的是一条窄的、显式启用的 Shared TUI deployment,不是通用本机 Server。具体 `EventQueue` 仍由 Core 产品装配;IPC 只把当前 TUI 必需的强类型操作和事件映射到同一个 Runtime owner,没有事件重放或公开协议承诺。 @@ -207,14 +207,14 @@ sequenceDiagram end ``` -当前私有协议(v8)只覆盖 TUI 已有用户旅程需要的窄操作: +当前私有协议(v9)只覆盖 TUI 已有用户旅程需要的窄操作: | 已支持 | 明确不支持 | |---|---| -| Health、Session list/create、原子 restore(含 transcript 与 pending Permission)、删除未被控制的空闲 Session、当前 Session rename、Agent mode/model update、声明式上下文 reload | Session archive/fork、跨 workspace attach、transcript 分页、模型目录/默认值和 Agent/Subagent 管理 | +| Health、Session list/create、原子 restore(含 transcript 与 pending Permission)、删除未被控制的空闲 Session、当前 Session fork(含 transcript)、rename、Agent mode/model update、声明式上下文 reload | Session archive、跨 workspace attach、transcript 分页、模型目录/默认值和 Agent/Subagent 管理 | | Turn submit/cancel、当前 Session 手动 context compaction | replay、cursor、resume event stream | -| pending/respond Permission、submit UserInput answers | observer、controller transfer、多 Session multiplex | -| 连接断开清理、Session-filtered events | detach/observer/controller transfer、SDK callbacks、GUI/Remote/Peer/ACP/Headless wire | +| pending/respond Permission、submit UserInput answers | observer、通用 controller transfer、多 Session multiplex | +| 连接断开清理、Session-filtered events | detach/observer/通用 controller transfer、SDK callbacks、GUI/Remote/Peer/ACP/Headless wire | 这些操作先满足以下本机 IPC 地基,而不把协议升级为公开 SDK: @@ -228,7 +228,7 @@ sequenceDiagram - JSON frame 使用 4-byte 长度前缀;request 在发送前执行 128 KiB 上限(覆盖 TUI 已有的 64 KiB 粘贴输入及类型化信封),response/event 在序列化时执行 8 MiB 上限。超限返回类型化错误,不能进行无界分配;超过该上限的历史 Session 暂由 Embedded TUI 打开,不在本阶段引入分页协议; - 未认证连接也计入有界 connection budget,单个客户端不能无限制造 server task; - 未知 frame/operation 信封字段、未知 operation、错误身份和不兼容版本 fail closed;复用的 Runtime DTO 按其既有反序列化契约处理字段; -- 一个连接最多控制一个 Session、同时最多提交一个活动 Turn;一个 Session 同时只有一个 controller。create/restore 在完整结果通过大小检查后才原子切换控制权,失败时保留原 Session。活动 Turn 期间不能切换 Session,也不能修改其名称、Agent mode 或 model;删除只作用于非当前且未被任何连接控制的 Session。 +- 一个连接最多控制一个 Session、同时最多提交一个活动 Turn;一个 Session 同时只有一个 controller。create/restore/fork 在完整结果通过大小检查后才原子切换控制权,失败时保留原 Session。fork 只接受当前 controller 的空闲 Session;无选中 Turn 时复制到最新持久化 Turn,指定 `before_turn_id` 时只复制该 Turn 之前的历史。活动 Turn 期间不能切换或 fork Session,也不能修改其名称、Agent mode 或 model;删除只作用于非当前且未被任何连接控制的 Session。 - Submit 与手动 context compaction 都使用调用方已有的 `turn_id` 标识不确定结果;若操作超时,返回 `outcome_unknown`、关闭连接并按该 ID 取消。手动 compaction 要求当前 controller 且 Session 空闲,由 Core 通过与普通对话 Turn 共用的原子准入路径创建一个可审计 maintenance Turn,并在取得所有权后读取压缩上下文:planning 阶段允许取消,atomic commit 开始后忽略晚到取消并保持 Processing 直至终态持久化完成。maintenance Turn 保留在权威 transcript 中但不进入模型上下文,live/restored payload 使用同一 compression ID 和 `applied` 事实;commit 后的持久化故障发布明确失败终态而不是遗留 Processing。断连取消只有得到确认后才释放 Session 控制权;无法确认时继续隔离该 Session,直到 Runtime 进程退出。 - Session delete/rename 和 Agent mode/model update 复用既有 Runtime 端口和校验,Runtime 对最终结果保持权威并拒绝无效目标。它们都是有副作用操作;发送前编码或 frame 上限失败表示请求未执行,连接仍可使用。rename 写入失败时恢复旧 metadata:确认恢复后返回明确失败,无法确认时返回 `outcome_unknown`。Shared Client 在请求写入后响应超时或丢失连接时也返回 `outcome_unknown` 并断开连接。两种情况都不自动重试:rename 由用户恢复 Session 并核对当前值;delete 由用户重新打开 `/sessions` 核对目标是否仍存在。模式与模型目录仍是同版本第一方产品事实,不加入 IPC。 - 声明式上下文 reload 只失效当前 Session 的 instructions 缓存,并按目标复用 Skill Registry 刷新;它可在活动 Turn 中执行但不改写该 Turn,generation 保护保证下一条消息重建上下文。它不引入 watcher、热替换或第二套 Runtime owner。 diff --git a/docs/architecture/cli-product-line-design.md b/docs/architecture/cli-product-line-design.md index 3b09138ae2..71cda26dee 100644 --- a/docs/architecture/cli-product-line-design.md +++ b/docs/architecture/cli-product-line-design.md @@ -256,12 +256,12 @@ Headless CLI 和公开 Agent SDK 都调用同一 Agent Runtime API,但交付 | 形态 | 默认部署 | 当前 Shared 范围 | |---|---|---| -| 交互式 TUI | Embedded | 显式 `--shared` 后支持 Session list/create/restore、transcript、当前 Session rename/Agent mode/model、声明式上下文 reload、当前 Session 手动 context compaction、Turn submit/cancel、Permission 和 UserInput | +| 交互式 TUI | Embedded | 显式 `--shared` 后支持 Session list/create/restore/delete/fork、transcript、当前 Session rename/Agent mode/model、声明式上下文 reload、当前 Session 手动 context compaction、Turn submit/cancel、Permission 和 UserInput | | `bitfun exec` / CI | Embedded | 不接受 Shared;保持独立进程、stdout/stderr 和退出码语义 | | ACP / SDK Host / GUI / Remote / Peer | 各自既有部署 | 不消费 TUI IPC,也不因本开关改变生命周期 | -Shared TUI 不提供 Session delete/fork、模型目录/默认值、Agent/Subagent 管理、MCP/扩展、账号同步、用量、observer、replay 或 controller transfer;对应入口给出明确的 Embedded 恢复建议,不在 Client 进程初始化第二套 Core owner。 -Shared 模式的斜杠命令、快捷键帮助和底部提示使用同一能力投影:`/rename ` 修改当前 Session 名称;`/agent`、Tab 和 Shift+Tab 只切换当前 Session 的 Agent mode;`/models` 只切换当前 Session 的 model;`/reload [skills|instructions]` 刷新下一条消息使用的声明式上下文;OpenCode 对齐的 `/compact` 及其 `/summarize` alias 以一个可取消的 maintenance Turn 压缩当前 Session 上下文,不增加自创命令或快捷键。该 Turn 与普通对话共用 Session 原子准入,取得所有权后再读取待压缩上下文;权威 transcript 保留完整 tool payload,但重建模型上下文时排除该 maintenance Turn。Embedded 与 Shared 的 `/help` 都从 Action Registry 展示这些入口;在 slash menu 中选择 rename 只预填命令并等待用户输入名称。若外部来源使用相同命令名,用户明确选择的 BitFun 命令可完成这一次参数提交,即使偏好保存失败也不会重新弹出来源选择。它们不进入管理页面,也不修改未来 Session 的默认值。其他不支持动作不显示为可执行入口。Session 切换失败保留原控制权;单个连接已有活动 Turn 时拒绝重复提交、manual compaction 以及 Session rename/mode/model update,但允许 reload 只影响下一条消息;事件订阅失效后当前视图立即失效并要求重启 Shared TUI。 +Shared TUI 不提供 Session archive、模型目录/默认值、Agent/Subagent 管理、MCP/扩展、账号同步、用量、observer、replay 或通用 controller transfer;对应入口给出明确的 Embedded 恢复建议,不在 Client 进程初始化第二套 Core owner。 +Shared 模式的斜杠命令、快捷键帮助和底部提示使用同一能力投影:OpenCode 对齐的 `/fork` 以 `Full session` 或历史用户提示词选择分支边界;选择提示词时 fork 只复制该 Turn 之前的历史,并把提示词放回 composer 而不自动发送。`/rename ` 修改当前 Session 名称;`/agent`、Tab 和 Shift+Tab 只切换当前 Session 的 Agent mode;`/models` 只切换当前 Session 的 model;`/reload [skills|instructions]` 刷新下一条消息使用的声明式上下文;OpenCode 对齐的 `/compact` 及其 `/summarize` alias 以一个可取消的 maintenance Turn 压缩当前 Session 上下文,不增加自创命令或快捷键。该 Turn 与普通对话共用 Session 原子准入,取得所有权后再读取待压缩上下文;权威 transcript 保留完整 tool payload,但重建模型上下文时排除该 maintenance Turn。Embedded 与 Shared 的 `/help` 都从 Action Registry 展示这些入口;在 slash menu 中选择 rename 只预填命令并等待用户输入名称。若外部来源使用相同命令名,用户明确选择的 BitFun 命令可完成这一次参数提交,即使偏好保存失败也不会重新弹出来源选择。它们不进入管理页面,也不修改未来 Session 的默认值。其他不支持动作不显示为可执行入口。Session 切换或 fork 失败保留原控制权;Shared fork 只有在新 Session 与 transcript 的响应可编码后才原子转移 controller。单个连接已有活动 Turn 时拒绝重复提交、fork、manual compaction 以及 Session rename/mode/model update,但允许 reload 只影响下一条消息;事件订阅失效后当前视图立即失效并要求重启 Shared TUI。 部署差异由 CLI Runtime client 封装。Embedded 以 Rust 类型直接调用 `AgentRuntime`,不初始化 IPC 或执行 JSON 编解码;Shared 将同一业务请求映射为一个有界本机 frame,Client/Server 各自只编码一次,再交给同一 Runtime owner。多 TUI 复用一个 Runtime 进程,连接和队列保持有界,不按 TUI 数量复制 Session owner。详细的 4+1 视图、帧上限和并发边界见 [`agent-runtime-deployment-design.md`](agent-runtime-deployment-design.md)。 diff --git a/scripts/core-boundaries/rules/source/forbidden-rules.mjs b/scripts/core-boundaries/rules/source/forbidden-rules.mjs index b7f8287cc0..773bb39911 100644 --- a/scripts/core-boundaries/rules/source/forbidden-rules.mjs +++ b/scripts/core-boundaries/rules/source/forbidden-rules.mjs @@ -6,9 +6,9 @@ export const forbiddenContentRules = [ reason: 'agent-runtime-ipc operation scope is frozen to the reviewed Shared TUI slice', patterns: [ { - regex: /^\s+(?!(?:Health|ListSessions|CreateSession|RestoreSession|DeleteSession|RenameSession|UpdateSessionMode|UpdateSessionModel|ReloadSessionContext|CompactSession|SubmitTurn|CancelTurn|PendingPermissions|RespondPermission|SubmitUserAnswers|Unit|Sessions|SessionCreated|SessionRestored|TurnAccepted|TurnCancelled|None|CurrentController|AttachExisting|UncontrolledTarget|Self|RuntimeIpcSessionRequirement|RuntimeIpcOperationRules|AgentContextReloadRequest|AgentDialogTurnRequest|AgentSessionCompactionRequest|AgentSessionCreateRequest|AgentSessionCreateResult|AgentSessionListRequest|AgentSessionModeUpdateRequest|AgentSessionModelUpdateRequest|AgentSessionSummary|AgentTurnCancellationRequest|AgentTurnCancellationResult|SessionTranscript)\b)[A-Z][A-Za-z0-9_]*\b/, + regex: /^\s+(?!(?:Health|ListSessions|CreateSession|RestoreSession|DeleteSession|ForkSession|RenameSession|UpdateSessionMode|UpdateSessionModel|ReloadSessionContext|CompactSession|SubmitTurn|CancelTurn|PendingPermissions|RespondPermission|SubmitUserAnswers|Unit|Sessions|SessionCreated|SessionRestored|SessionForked|TurnAccepted|TurnCancelled|None|CurrentController|AttachExisting|UncontrolledTarget|Self|RuntimeIpcSessionRequirement|RuntimeIpcOperationRules|RuntimeSessionForkRequest|AgentContextReloadRequest|AgentDialogTurnRequest|AgentSessionCompactionRequest|AgentSessionCreateRequest|AgentSessionCreateResult|AgentSessionListRequest|AgentSessionModeUpdateRequest|AgentSessionModelUpdateRequest|AgentSessionSummary|AgentTurnCancellationRequest|AgentTurnCancellationResult|SessionTranscript)\b)[A-Z][A-Za-z0-9_]*\b/, message: - 'agent-runtime-ipc may not add archive, replay, observer, controller-transfer, fork, or other operations beyond the reviewed Shared TUI slice', + 'agent-runtime-ipc may not add archive, replay, observer, general controller-transfer, or other operations beyond the reviewed Shared TUI slice', }, ], }, diff --git a/scripts/core-boundaries/self-test.mjs b/scripts/core-boundaries/self-test.mjs index 2a782878f4..60e06f636d 100644 --- a/scripts/core-boundaries/self-test.mjs +++ b/scripts/core-boundaries/self-test.mjs @@ -4859,10 +4859,12 @@ export function runManifestParserSelfTest({ ].every((name) => runtimeIpcOperationPattern.test(` ${name},`)) || runtimeIpcOperationPattern.test(' Health,') || runtimeIpcOperationPattern.test(' DeleteSession {') || + runtimeIpcOperationPattern.test(' ForkSession {') || runtimeIpcOperationPattern.test(' RenameSession {') || runtimeIpcOperationPattern.test(' UpdateSessionMode {') || runtimeIpcOperationPattern.test(' UpdateSessionModel {') || - runtimeIpcOperationPattern.test(' SubmitTurn {') + runtimeIpcOperationPattern.test(' SubmitTurn {') || + runtimeIpcOperationPattern.test(' SessionForked {') ) { throw new Error('agent-runtime-ipc operation guard must preserve the Shared TUI operation budget'); } diff --git a/src/apps/cli/README.md b/src/apps/cli/README.md index 5ab94d6a63..e44ac2cfd8 100644 --- a/src/apps/cli/README.md +++ b/src/apps/cli/README.md @@ -87,6 +87,9 @@ The Embedded and Shared TUI use the same session command names: - `/sessions` opens the session browser; `/resume`, `/continue`, and `/history` are aliases. - `/new` starts a fresh conversation session; `/clear` is its OpenCode-compatible alias and does not merely clear the terminal display. +- `/fork` opens an OpenCode-compatible fork dialog. `Full session` copies through the latest + persisted turn; choosing a previous user prompt forks immediately before that turn and copies the + prompt into the composer without sending it. Forking requires an idle session. - `/status` opens a transient view of current session, runtime, workspace, approval, and latest primary-model request facts observed by this TUI. It is not a cumulative usage report; use `/usage` for cumulative session usage in Embedded TUI. diff --git a/src/apps/cli/src/actions.rs b/src/apps/cli/src/actions.rs index 0106c29591..26c101dadc 100644 --- a/src/apps/cli/src/actions.rs +++ b/src/apps/cli/src/actions.rs @@ -74,6 +74,7 @@ pub(crate) enum ActionHandler { AddModel, NewSession, Sessions, + ForkSession, RenameSession, Skills, Reload, @@ -115,7 +116,7 @@ pub(crate) enum ActionHandler { pub(crate) const SHARED_TUI_EMBEDDED_HANDOFF: &str = "Exit all Shared TUI clients, wait up to 30 seconds for their Runtime to stop, then use default Embedded `bitfun chat`"; pub(crate) const SHARED_TUI_HELP_NOTE: &str = - "Shared TUI: start with `bitfun chat --shared`. Multiple TUI processes reuse one workspace Runtime, while each TUI controls at most one Session and each Session has one controller. Use `/sessions` and Ctrl+D to delete an idle, non-current Session; use `/rename ` to rename the current Session, `/compact` to compact its context, `/agent`, Tab, or Shift+Tab to change its Agent mode, `/models` to change its model, and `/reload [skills|instructions]` to refresh declarative context for the next message. Model configuration, Agent/Subagent management, MCP, extension, account-sync, usage, and other management remain Embedded. Exit all Shared TUI clients and wait up to 30 seconds before returning to default Embedded `bitfun chat`."; + "Shared TUI: start with `bitfun chat --shared`. Multiple TUI processes reuse one workspace Runtime, while each TUI controls at most one Session and each Session has one controller. Use `/sessions` and Ctrl+D to delete an idle, non-current Session; use `/fork` to branch the current idle Session, `/rename ` to rename it, `/compact` to compact its context, `/agent`, Tab, or Shift+Tab to change its Agent mode, `/models` to change its model, and `/reload [skills|instructions]` to refresh declarative context for the next message. Model configuration, Agent/Subagent management, MCP, extension, account-sync, usage, and other management remain Embedded. Exit all Shared TUI clients and wait up to 30 seconds before returning to default Embedded `bitfun chat`."; impl ActionHandler { pub(crate) const fn available_in_shared_tui(self, context: ActionContext) -> bool { @@ -126,6 +127,7 @@ impl ActionHandler { | Self::SelectTheme | Self::NewSession | Self::Sessions + | Self::ForkSession | Self::RenameSession | Self::AcpHelp | Self::Init @@ -374,6 +376,21 @@ static ACTION_SPECS: &[ActionSpec] = &[ shortcut_label: None, slash_on_startup: false, }, + ActionSpec { + id: "fork_session", + name: "Fork session", + aliases: &["/fork"], + description: "Fork from the full session or a previous prompt", + contexts: CHAT, + availability: ActionAvailability::Idle, + handler: ActionHandler::ForkSession, + default_bindings: &[], + fallback_bindings: &[], + shortcut_field: None, + palette: palette("Session", false), + shortcut_label: None, + slash_on_startup: false, + }, ActionSpec { id: "skills", name: "Skills", @@ -1838,6 +1855,7 @@ mod tests { assert!(SHARED_TUI_HELP_NOTE.contains("bitfun chat --shared")); assert!(SHARED_TUI_HELP_NOTE.contains("one Session")); assert!(SHARED_TUI_HELP_NOTE.contains("`/models`")); + assert!(SHARED_TUI_HELP_NOTE.contains("`/fork`")); assert!(SHARED_TUI_HELP_NOTE.contains("`/rename `")); assert!(SHARED_TUI_HELP_NOTE.contains("`/reload [skills|instructions]`")); assert!(SHARED_TUI_HELP_NOTE.contains("Ctrl+D")); @@ -1876,6 +1894,22 @@ mod tests { assert!(action_by_id("rename_session", ActionContext::Startup).is_none()); } + #[test] + fn fork_uses_only_the_opencode_command_and_is_idle_in_both_deployments() { + let action = action_by_id("fork_session", ActionContext::Chat) + .expect("OpenCode-compatible session fork action"); + + assert_eq!(action.aliases, &["/fork"]); + assert_eq!(action.handler, ActionHandler::ForkSession); + assert_eq!(action.availability, ActionAvailability::Idle); + assert!(action.default_bindings.is_empty()); + assert!(action.available(ActionState::chat(false, false))); + assert!(action.available(ActionState::chat(false, false).for_shared_tui())); + assert!(!action.available(ActionState::chat(true, false))); + assert!(action_by_id("fork_session", ActionContext::Startup).is_none()); + assert!(action_for_alias("/branch", ActionContext::Chat).is_none()); + } + #[test] fn compact_uses_opencode_commands_without_an_invented_shortcut() { let action = action_by_id("compact_session", ActionContext::Chat) diff --git a/src/apps/cli/src/agent/runtime_client.rs b/src/apps/cli/src/agent/runtime_client.rs index 04e33db967..d72b7e2b65 100644 --- a/src/apps/cli/src/agent/runtime_client.rs +++ b/src/apps/cli/src/agent/runtime_client.rs @@ -14,18 +14,18 @@ use tokio::sync::{broadcast, Mutex}; use bitfun_agent_runtime::sdk::{ AgentDialogTurnRequest, AgentEventReceiver, AgentLocalCommandTurnRecordRequest, AgentRuntime, AgentSessionCompactionRequest, AgentSessionCreateRequest, AgentSessionDeleteRequest, - AgentSessionForkRequest, AgentSessionForkResult, AgentSessionListRequest, - AgentSessionModeUpdateRequest, AgentSessionModelUpdateRequest, AgentSessionRenameRequest, - AgentSessionRestoreRequest, AgentSessionUsageRequest, AgentTurnCancellationRequest, - AgentTurnSettlementRequest, AgentUserAnswersRequest, PermissionReply, PermissionRequest, - PermissionRequestEventReceiver, PortError, PortErrorKind, RuntimeError, SessionTranscript, - SessionTranscriptRequest, SessionUsageReport, + AgentSessionForkBeforeTurnRequest, AgentSessionForkRequest, AgentSessionForkResult, + AgentSessionListRequest, AgentSessionModeUpdateRequest, AgentSessionModelUpdateRequest, + AgentSessionRenameRequest, AgentSessionRestoreRequest, AgentSessionUsageRequest, + AgentTurnCancellationRequest, AgentTurnSettlementRequest, AgentUserAnswersRequest, + PermissionReply, PermissionRequest, PermissionRequestEventReceiver, PortError, PortErrorKind, + RuntimeError, SessionTranscript, SessionTranscriptRequest, SessionUsageReport, }; use bitfun_agent_runtime_ipc::{ RuntimeIpcClient, RuntimeIpcClientError, RuntimeIpcClientEvent, RuntimeIpcErrorCode, RuntimeIpcEvent, RuntimeIpcOperation, RuntimeIpcOperationResult, - RuntimeIpcStreamInvalidationReason, RuntimeSessionRenameRequest, RuntimeSessionRestoreRequest, - RuntimeUserAnswersRequest, + RuntimeIpcStreamInvalidationReason, RuntimeSessionForkRequest, RuntimeSessionRenameRequest, + RuntimeSessionRestoreRequest, RuntimeUserAnswersRequest, }; use bitfun_events::{AgenticEvent, AgenticEventEnvelope}; use bitfun_runtime_ports::{ @@ -731,6 +731,90 @@ impl CliAgentRuntimeClient { .map_err(|error| anyhow::anyhow!(error.into_message())) } + pub(crate) async fn fork_current_session( + &self, + before_turn_id: Option<&str>, + ) -> Result<( + AgentSessionSummary, + AgentSessionWorkspaceBinding, + SessionTranscript, + )> { + let source_session_id = self.require_session_id().await?; + let workspace_path = self.project_workspace_path_string(); + let (session, transcript) = match &self.backend { + CliAgentRuntimeBackend::Embedded(runtime) => { + let forked = match before_turn_id { + Some(source_turn_id) => { + runtime + .fork_session_before_turn(AgentSessionForkBeforeTurnRequest { + workspace_path: workspace_path.clone(), + source_session_id, + source_turn_id: source_turn_id.to_string(), + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + } + None => { + runtime + .fork_session(AgentSessionForkRequest { + workspace_path: workspace_path.clone(), + source_session_id, + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + } + } + .map_err(|error| anyhow::anyhow!(error.into_message()))?; + let restored = runtime + .restore_session(AgentSessionRestoreRequest { + workspace_path: workspace_path.clone(), + session_id: forked.session_id.clone(), + include_internal: false, + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + .map_err(|error| anyhow::anyhow!(error.into_message()))?; + let transcript = runtime + .read_session_transcript(SessionTranscriptRequest { + session_id: forked.session_id, + turn_id: None, + }) + .await + .map_err(|error| anyhow::anyhow!(error.into_message()))?; + (restored.session, transcript) + } + CliAgentRuntimeBackend::Shared(client) => match client + .request(RuntimeIpcOperation::ForkSession { + request: RuntimeSessionForkRequest { + session_id: source_session_id, + before_turn_id: before_turn_id.map(str::to_string), + }, + }) + .await? + { + RuntimeIpcOperationResult::SessionForked { + session, + transcript, + } => (session, transcript), + _ => return Err(unexpected_shared_result("fork_session")), + }, + }; + + let binding = self + .resolve_session_workspace_binding(&session.session_id, Path::new(&workspace_path)) + .await?; + *self.session_id.lock().await = Some(session.session_id.clone()); + *self.current_turn_id.lock().await = None; + self.shared_pending_permissions + .write() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + .clear(); + Ok((session, binding, transcript)) + } + pub(crate) async fn generate_session_usage_report( &self, request: AgentSessionUsageRequest, @@ -1553,6 +1637,25 @@ mod tests { assert!(!compact.contains("serde_json::from_value")); } + #[test] + fn interactive_session_fork_uses_the_same_runtime_boundary_in_both_deployments() { + let source = include_str!("runtime_client.rs").replace("\r\n", "\n"); + let fork = source + .split_once("pub(crate) async fn fork_current_session(") + .expect("interactive fork method") + .1 + .split_once("pub(crate) async fn generate_session_usage_report(") + .expect("interactive fork method boundary") + .0; + + assert!(fork.contains("CliAgentRuntimeBackend::Embedded(runtime)")); + assert!(fork.contains(".fork_session_before_turn(")); + assert!(fork.contains(".fork_session(AgentSessionForkRequest")); + assert!(fork.contains("CliAgentRuntimeBackend::Shared(client)")); + assert!(fork.contains("RuntimeIpcOperation::ForkSession")); + assert!(fork.contains("RuntimeIpcOperationResult::SessionForked")); + } + #[test] fn session_delete_uses_direct_runtime_or_private_shared_ipc() { let source = include_str!("runtime_client.rs").replace("\r\n", "\n"); diff --git a/src/apps/cli/src/chat_state.rs b/src/apps/cli/src/chat_state.rs index 73e0a1883a..3d4df90fdd 100644 --- a/src/apps/cli/src/chat_state.rs +++ b/src/apps/cli/src/chat_state.rs @@ -163,6 +163,8 @@ pub(crate) enum FlowItem { #[derive(Debug, Clone)] pub(crate) struct ChatMessage { pub id: String, + /// Stable persisted DialogTurn identity used for history operations. + pub turn_id: Option, pub role: MessageRole, pub timestamp: SystemTime, pub flow_items: Vec, @@ -272,6 +274,7 @@ impl ChatMessage { .id .clone() .unwrap_or_else(|| format!("transcript-message-{index}")), + turn_id: msg.turn_id.clone(), role, timestamp: UNIX_EPOCH .checked_add(Duration::from_millis(msg.timestamp_ms.unwrap_or_default())) @@ -283,6 +286,13 @@ impl ChatMessage { } } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct SessionForkPoint { + pub turn_id: String, + pub prompt: String, + pub timestamp: SystemTime, +} + // ============ Chat Metadata ============ /// Statistics for the current chat session @@ -460,6 +470,32 @@ impl ChatState { self.metadata.message_count > 0 } + /// User prompts eligible for `/fork`, newest first like OpenCode's fork dialog. + pub(crate) fn session_fork_points(&self) -> Vec { + self.messages + .iter() + .rev() + .filter(|message| message.role == MessageRole::User) + .filter_map(|message| { + let turn_id = message.turn_id.clone()?; + let prompt = message + .flow_items + .iter() + .filter_map(|item| match item { + FlowItem::Text { content, .. } => Some(content.as_str()), + FlowItem::Thinking { .. } | FlowItem::Tool { .. } => None, + }) + .collect::>() + .join("\n"); + (!prompt.is_empty()).then_some(SessionForkPoint { + turn_id, + prompt, + timestamp: message.timestamp, + }) + }) + .collect() + } + pub(crate) fn set_worktree_control_available(&mut self, available: bool) { self.worktree_control_available = available; } @@ -720,6 +756,7 @@ impl ChatState { // Add user message self.messages.push(ChatMessage { id: uuid::Uuid::new_v4().to_string(), + turn_id: Some(turn_id.to_string()), role: MessageRole::User, timestamp: SystemTime::now(), flow_items: vec![FlowItem::Text { @@ -734,6 +771,7 @@ impl ChatState { // Add empty assistant message (will be filled by streaming) self.messages.push(ChatMessage { id: uuid::Uuid::new_v4().to_string(), + turn_id: Some(turn_id.to_string()), role: MessageRole::Assistant, timestamp: SystemTime::now(), flow_items: Vec::new(), @@ -1172,6 +1210,7 @@ impl ChatState { pub(crate) fn add_system_message(&mut self, content: String) { self.messages.push(ChatMessage { id: uuid::Uuid::new_v4().to_string(), + turn_id: None, role: MessageRole::System, timestamp: SystemTime::now(), flow_items: vec![FlowItem::Text { @@ -1187,6 +1226,7 @@ impl ChatState { pub(crate) fn add_assistant_message(&mut self, content: String) { self.messages.push(ChatMessage { id: uuid::Uuid::new_v4().to_string(), + turn_id: None, role: MessageRole::Assistant, timestamp: SystemTime::now(), flow_items: vec![FlowItem::Text { @@ -1752,6 +1792,51 @@ mod tests { ); } + #[test] + fn session_fork_points_keep_stable_turn_ids_and_newest_prompt_first() { + let transcript = SessionTranscript { + session_id: "session-1".to_string(), + messages: vec![ + TranscriptMessage { + id: Some("user-1".to_string()), + role: "user".to_string(), + turn_id: Some("turn-1".to_string()), + timestamp_ms: Some(1_000), + content: TranscriptContent::Text("First prompt".to_string()), + }, + TranscriptMessage { + id: Some("assistant-1".to_string()), + role: "assistant".to_string(), + turn_id: Some("turn-1".to_string()), + timestamp_ms: Some(1_100), + content: TranscriptContent::Text("First answer".to_string()), + }, + TranscriptMessage { + id: Some("user-2".to_string()), + role: "user".to_string(), + turn_id: Some("turn-2".to_string()), + timestamp_ms: Some(2_000), + content: TranscriptContent::Text("Second\nprompt".to_string()), + }, + ], + }; + let state = ChatState::from_session_transcript( + "session-1".to_string(), + "Session".to_string(), + "agentic".to_string(), + None, + &transcript, + ); + + let points = state.session_fork_points(); + + assert_eq!(points.len(), 2); + assert_eq!(points[0].turn_id, "turn-2"); + assert_eq!(points[0].prompt, "Second\nprompt"); + assert_eq!(points[1].turn_id, "turn-1"); + assert_eq!(points[1].prompt, "First prompt"); + } + #[test] fn transcript_history_merges_tool_results_into_the_rendered_tool_card() { let transcript = SessionTranscript { diff --git a/src/apps/cli/src/modes/chat.rs b/src/apps/cli/src/modes/chat.rs index ff9e749573..418db49875 100644 --- a/src/apps/cli/src/modes/chat.rs +++ b/src/apps/cli/src/modes/chat.rs @@ -36,6 +36,7 @@ use crate::ui::agent_selector::{AgentItem, AgentSelectorAction}; use crate::ui::chat::{session_status_text, ChatView, MouseGestureOutcome}; use crate::ui::command_menu::{ExternalCommandProjection, NativeCommandCollisionProjection}; use crate::ui::command_palette::PaletteAction; +use crate::ui::fork_selector::{ForkAction, ForkTarget}; use crate::ui::login_form::LoginFormAction; use crate::ui::mcp_add_dialog::McpAddAction; use crate::ui::mcp_selector::{McpItem, McpItemAction}; @@ -136,6 +137,8 @@ pub(crate) enum ChatExitReason { Quit, /// Switch to a different session SwitchSession(String), + /// Fork the current session at the selected OpenCode-compatible boundary. + ForkSession(crate::ui::fork_selector::ForkTarget), /// Create a new session NewSession, } diff --git a/src/apps/cli/src/modes/chat/account.rs b/src/apps/cli/src/modes/chat/account.rs index 72df1c65f1..aab4bf94d0 100644 --- a/src/apps/cli/src/modes/chat/account.rs +++ b/src/apps/cli/src/modes/chat/account.rs @@ -187,6 +187,7 @@ impl ChatMode { || chat_view.model_selector_visible() || chat_view.agent_selector_visible() || chat_view.session_selector_visible() + || chat_view.fork_selector_visible() || chat_view.skill_selector_visible() || chat_view.subagent_selector_visible() || chat_view.mcp_selector_visible() @@ -208,6 +209,7 @@ impl ChatMode { chat_view.hide_model_selector(); chat_view.hide_agent_selector(); chat_view.hide_session_selector(); + chat_view.hide_fork_selector(); chat_view.hide_skill_selector(); chat_view.hide_subagent_selector(); chat_view.hide_mcp_selector(); @@ -230,6 +232,7 @@ impl ChatMode { crate::ui::chat::PopupType::ModelSelector => chat_view.hide_model_selector(), crate::ui::chat::PopupType::AgentSelector => chat_view.hide_agent_selector(), crate::ui::chat::PopupType::SessionSelector => chat_view.hide_session_selector(), + crate::ui::chat::PopupType::ForkSelector => chat_view.hide_fork_selector(), crate::ui::chat::PopupType::SkillSelector => chat_view.hide_skill_selector(), crate::ui::chat::PopupType::SubagentSelector => chat_view.hide_subagent_selector(), crate::ui::chat::PopupType::McpSelector => chat_view.hide_mcp_selector(), @@ -255,6 +258,7 @@ impl ChatMode { crate::ui::chat::PopupType::SessionSelector => { chat_view.reshow_session_selector() } + crate::ui::chat::PopupType::ForkSelector => chat_view.reshow_fork_selector(), crate::ui::chat::PopupType::SkillSelector => chat_view.reshow_skill_selector(), crate::ui::chat::PopupType::SubagentSelector => { chat_view.reshow_subagent_selector() diff --git a/src/apps/cli/src/modes/chat/commands.rs b/src/apps/cli/src/modes/chat/commands.rs index fe277c6273..777a8c695e 100644 --- a/src/apps/cli/src/modes/chat/commands.rs +++ b/src/apps/cli/src/modes/chat/commands.rs @@ -42,6 +42,7 @@ fn pending_session_operation_blocks_runtime_action( && matches!( handler, ActionHandler::Sessions + | ActionHandler::ForkSession | ActionHandler::RenameSession | ActionHandler::CompactSession | ActionHandler::Init @@ -81,15 +82,20 @@ fn builtin_arguments_route(route: CommandRoute, handler: ActionHandler) -> bool route == CommandRoute::Builtin && handler == ActionHandler::RenameSession } -fn compact_arguments_error( +fn builtin_arguments_error( route: CommandRoute, handler: ActionHandler, arguments: &str, ) -> Option<&'static str> { - (route == CommandRoute::Builtin - && handler == ActionHandler::CompactSession - && !arguments.trim().is_empty()) - .then_some("Usage: /compact") + if route != CommandRoute::Builtin || arguments.trim().is_empty() { + return None; + } + + match handler { + ActionHandler::CompactSession => Some("Usage: /compact"), + ActionHandler::ForkSession => Some("Usage: /fork"), + _ => None, + } } fn selected_command_prefill(handler: ActionHandler) -> Option<&'static str> { @@ -143,7 +149,12 @@ fn clear_selected_native_command_prefill( fn session_command_help_note() -> String { let rename = action_for_alias("/rename", ActionContext::Chat) .expect("current session rename action must remain registered"); - format!("Session Commands\n {}", rename.description) + let fork = action_for_alias("/fork", ActionContext::Chat) + .expect("current session fork action must remain registered"); + format!( + "Session Commands\n {}\n {}", + fork.description, rename.description + ) } impl ChatMode { @@ -358,7 +369,7 @@ impl ChatMode { if let Some(action) = builtin_action { let state = self.action_state(chat_state.is_processing, false); if let Some(usage) = - compact_arguments_error(CommandRoute::Builtin, action.handler, arguments) + builtin_arguments_error(CommandRoute::Builtin, action.handler, arguments) { chat_view.set_status(Some(usage.to_string())); return Ok(None); @@ -432,7 +443,7 @@ impl ChatMode { ) }; if let Some(action) = builtin_action { - if let Some(usage) = compact_arguments_error(route, action.handler, arguments) { + if let Some(usage) = builtin_arguments_error(route, action.handler, arguments) { chat_view.set_status(Some(usage.to_string())); return Ok(None); } @@ -903,6 +914,9 @@ impl ChatMode { ActionHandler::Sessions => { self.show_session_selector(chat_view, chat_state, rt_handle); } + ActionHandler::ForkSession => { + self.show_fork_selector(chat_view, chat_state); + } ActionHandler::RenameSession => { return self.start_session_rename("", chat_view, chat_state, rt_handle); } diff --git a/src/apps/cli/src/modes/chat/input.rs b/src/apps/cli/src/modes/chat/input.rs index c606081bdb..4dab9ff9de 100644 --- a/src/apps/cli/src/modes/chat/input.rs +++ b/src/apps/cli/src/modes/chat/input.rs @@ -195,6 +195,16 @@ impl ChatMode { return Ok(None); } + if chat_view.fork_selector_visible() { + match chat_view.fork_selector_handle_key(key) { + ForkAction::Select(target) => { + return Ok(Some(ChatExitReason::ForkSession(target))); + } + ForkAction::Close | ForkAction::None => {} + } + return Ok(None); + } + if chat_view.skill_selector_visible() { match key.code { KeyCode::Up => chat_view.skill_selector_up(), @@ -403,7 +413,9 @@ impl ChatMode { } = context; if matches!( &reason, - ChatExitReason::SwitchSession(_) | ChatExitReason::NewSession + ChatExitReason::SwitchSession(_) + | ChatExitReason::ForkSession(_) + | ChatExitReason::NewSession ) && shared_session_change_is_blocked( this.agent.is_shared(), this.pending_session_operation.is_some(), @@ -446,6 +458,15 @@ impl ChatMode { } } } + ChatExitReason::ForkSession(target) => { + match this.fork_session(target, session_id, chat_state, chat_view, rt_handle) { + Ok(()) => tracing::info!("Forked current session: {}", session_id), + Err(error) => { + chat_view.set_status(Some(format!("Failed to fork session: {error}"))); + tracing::error!("Failed to fork session: {error}"); + } + } + } ChatExitReason::Quit => { if let Some(pending) = this.pending_session_operation.as_mut() { if !pending.exit_warning_shown { diff --git a/src/apps/cli/src/modes/chat/sessions.rs b/src/apps/cli/src/modes/chat/sessions.rs index 3cae71d3da..5224e5a911 100644 --- a/src/apps/cli/src/modes/chat/sessions.rs +++ b/src/apps/cli/src/modes/chat/sessions.rs @@ -1,4 +1,61 @@ impl ChatMode { + fn fork_session( + &mut self, + target: ForkTarget, + session_id: &mut String, + chat_state: &mut ChatState, + chat_view: &mut ChatView, + rt_handle: &tokio::runtime::Handle, + ) -> Result<()> { + let (before_turn_id, prefill) = match target { + ForkTarget::FullSession => (None, None), + ForkTarget::BeforeTurn { turn_id, prompt } => (Some(turn_id), Some(prompt)), + }; + chat_view.set_status(Some("Forking session...".to_string())); + self.close_all_popups(chat_view); + let agent = self.agent.clone(); + let (summary, workspace_binding, transcript) = tokio::task::block_in_place(|| { + rt_handle.block_on(agent.fork_current_session(before_turn_id.as_deref())) + })?; + let new_session_id = summary.session_id.clone(); + let restored_agent_type = summary.agent_type.clone(); + let mut new_state = ChatState::from_session_transcript( + new_session_id.clone(), + summary.session_name.clone(), + restored_agent_type.clone(), + Some(workspace_binding.workspace_path.clone()), + &transcript, + ); + new_state.current_model_id = summary.model_id; + new_state.apply_workspace_binding(workspace_binding); + + *session_id = new_session_id.clone(); + *chat_state = new_state; + chat_state.set_worktree_control_available(!self.agent.is_shared()); + self.agent_type = restored_agent_type; + self.workspace = chat_state.workspace.clone(); + self.refresh_workspace_git_status(chat_state, rt_handle); + self.auto_approve_ask_override = None; + clear_selected_native_command_prefill(&mut self.selected_native_command_once, chat_view); + chat_state.auto_approve_ask = self.auto_approve_ask_default; + self.agent + .set_approval_policy(crate::runtime::approval::CliApprovalPolicy::Ask); + self.load_current_model_name(chat_state, rt_handle); + + chat_view.clear_screen(); + chat_view.scroll_to_bottom(); + if let Some(prompt) = prefill { + chat_view.set_input(&prompt); + chat_view.set_status(Some( + "Forked before the selected prompt; review the copied input before sending." + .to_string(), + )); + } else { + chat_view.set_status(Some(format!("Forked session: {}", summary.session_name))); + } + Ok(()) + } + /// Switch to a different session: restore it from core, reload messages, update state fn switch_to_session( &mut self, @@ -222,6 +279,20 @@ impl ChatMode { ); } + fn show_fork_selector(&self, chat_view: &mut ChatView, chat_state: &ChatState) { + let points = chat_state.session_fork_points(); + if points.is_empty() { + chat_view.set_status(Some( + "No persisted prompts are available to fork.".to_string(), + )); + return; + } + chat_view.show_fork_selector(points); + chat_view.set_status(Some( + "Choose Full session or a prompt to fork immediately before it.".to_string(), + )); + } + /// Handle session deletion from the session selector fn handle_session_delete( &mut self, diff --git a/src/apps/cli/src/modes/chat/tests.rs b/src/apps/cli/src/modes/chat/tests.rs index 9e76066986..4558e1c2e0 100644 --- a/src/apps/cli/src/modes/chat/tests.rs +++ b/src/apps/cli/src/modes/chat/tests.rs @@ -5,10 +5,10 @@ mod tests { use super::{ action_opens_extension_management, agent_event_stream_failure, apply_agent_mode_feedback, apply_model_selection_feedback, apply_session_model_migration, - apply_session_rename_feedback, begin_slash_menu_selection, builtin_arguments_route, - builtin_command_reconfirmation, clear_selected_native_command_prefill, - cli_native_prompt_command_descriptors, command_route, compact_arguments_error, - consume_selected_native_command_once, context_compression_tool_event, + apply_session_rename_feedback, begin_slash_menu_selection, builtin_arguments_error, + builtin_arguments_route, builtin_command_reconfirmation, + clear_selected_native_command_prefill, cli_native_prompt_command_descriptors, + command_route, consume_selected_native_command_once, context_compression_tool_event, extension_command_help_request, external_agent_attention, external_agent_diagnostic_lines, external_agent_pending_notice_key, external_agent_result_is_stale, external_agent_review_text, external_command_projections, external_control_review_text, @@ -1168,9 +1168,9 @@ mod tests { } #[test] - fn compact_rejects_arguments_only_after_the_builtin_route_wins() { + fn session_actions_reject_arguments_only_after_the_builtin_route_wins() { assert_eq!( - compact_arguments_error( + builtin_arguments_error( CommandRoute::Builtin, ActionHandler::CompactSession, "unexpected" @@ -1178,17 +1178,25 @@ mod tests { Some("Usage: /compact") ); assert_eq!( - compact_arguments_error(CommandRoute::Builtin, ActionHandler::CompactSession, " "), + builtin_arguments_error(CommandRoute::Builtin, ActionHandler::CompactSession, " "), None ); assert_eq!( - compact_arguments_error( + builtin_arguments_error( CommandRoute::External, ActionHandler::CompactSession, "unexpected" ), None ); + assert_eq!( + builtin_arguments_error( + CommandRoute::Builtin, + ActionHandler::ForkSession, + "unexpected" + ), + Some("Usage: /fork") + ); } #[test] diff --git a/src/apps/cli/src/shared_runtime.rs b/src/apps/cli/src/shared_runtime.rs index e09429913d..5963b876c9 100644 --- a/src/apps/cli/src/shared_runtime.rs +++ b/src/apps/cli/src/shared_runtime.rs @@ -1,7 +1,8 @@ use anyhow::{anyhow, Context, Result}; use async_trait::async_trait; use bitfun_agent_runtime::sdk::{ - AgentRuntime, AgentSessionDeleteRequest, AgentSessionRenameRequest, AgentSessionRestoreRequest, + AgentRuntime, AgentSessionDeleteRequest, AgentSessionForkBeforeTurnRequest, + AgentSessionForkRequest, AgentSessionRenameRequest, AgentSessionRestoreRequest, AgentUserAnswersRequest, DialogSubmitOutcome, PermissionRequest, PermissionRequestEvent, PortErrorKind, RuntimeError, SessionTranscriptRequest, }; @@ -265,6 +266,56 @@ impl RuntimeIpcRequestHandler for SharedRuntimeHandler { pending_permissions, }) } + RuntimeIpcOperation::ForkSession { request } => { + let workspace_path = self.workspace.to_string_lossy().into_owned(); + let forked = match request.before_turn_id { + Some(source_turn_id) => { + self.runtime + .fork_session_before_turn(AgentSessionForkBeforeTurnRequest { + workspace_path: workspace_path.clone(), + source_session_id: request.session_id, + source_turn_id, + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + } + None => { + self.runtime + .fork_session(AgentSessionForkRequest { + workspace_path: workspace_path.clone(), + source_session_id: request.session_id, + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + } + } + .map_err(runtime_ipc_error)?; + let restored = self + .runtime + .restore_session(AgentSessionRestoreRequest { + workspace_path, + session_id: forked.session_id.clone(), + include_internal: false, + remote_connection_id: None, + remote_ssh_host: None, + }) + .await + .map_err(runtime_ipc_error)?; + let transcript = self + .runtime + .read_session_transcript(SessionTranscriptRequest { + session_id: forked.session_id, + turn_id: None, + }) + .await + .map_err(runtime_ipc_error)?; + Ok(RuntimeIpcOperationResult::SessionForked { + session: restored.session, + transcript, + }) + } RuntimeIpcOperation::DeleteSession { session_id } => { delete_owned_session(&self.runtime, &self.workspace, session_id).await?; Ok(RuntimeIpcOperationResult::Unit) diff --git a/src/apps/cli/src/ui/chat/popups.rs b/src/apps/cli/src/ui/chat/popups.rs index b8d6226d6c..e56037131f 100644 --- a/src/apps/cli/src/ui/chat/popups.rs +++ b/src/apps/cli/src/ui/chat/popups.rs @@ -436,6 +436,35 @@ impl ChatView { self.session_selector.remove_item(session_id); } + // ============ Fork selector methods ============ + + pub(crate) fn show_fork_selector( + &mut self, + points: Vec, + ) { + self.fork_selector.show(points); + self.popup_stack.push(PopupType::ForkSelector); + } + + pub(crate) fn fork_selector_visible(&self) -> bool { + self.fork_selector.is_visible() + } + + pub(crate) fn hide_fork_selector(&mut self) { + self.fork_selector.hide(); + } + + pub(crate) fn reshow_fork_selector(&mut self) { + self.fork_selector.reshow(); + } + + pub(crate) fn fork_selector_handle_key( + &mut self, + key: crossterm::event::KeyEvent, + ) -> ForkAction { + self.fork_selector.handle_key_event(key) + } + // ============ Provider selector methods (add model step 1) ============ pub(crate) fn show_provider_selector(&mut self) { diff --git a/src/apps/cli/src/ui/chat/render.rs b/src/apps/cli/src/ui/chat/render.rs index 6164a1d93d..45e09de897 100644 --- a/src/apps/cli/src/ui/chat/render.rs +++ b/src/apps/cli/src/ui/chat/render.rs @@ -97,6 +97,7 @@ impl ChatView { self.render_model_selector(frame, chunks[1]); self.render_agent_selector(frame, chunks[1]); self.render_session_selector(frame, chunks[1]); + self.render_fork_selector(frame, chunks[1]); self.render_skill_selector(frame, chunks[1]); self.render_subagent_selector(frame, chunks[1]); self.render_mcp_selector(frame, chunks[1]); @@ -887,6 +888,10 @@ impl ChatView { self.session_selector.render(frame, area, &self.theme); } + fn render_fork_selector(&mut self, frame: &mut Frame, area: Rect) { + self.fork_selector.render(frame, area, &self.theme); + } + fn render_skill_selector(&mut self, frame: &mut Frame, area: Rect) { self.skill_selector.render(frame, area, &self.theme); } diff --git a/src/apps/cli/src/ui/chat/state.rs b/src/apps/cli/src/ui/chat/state.rs index 79466c5417..fa9458bc56 100644 --- a/src/apps/cli/src/ui/chat/state.rs +++ b/src/apps/cli/src/ui/chat/state.rs @@ -11,6 +11,7 @@ use unicode_width::{UnicodeWidthChar, UnicodeWidthStr}; use super::agent_selector::{AgentItem, AgentSelectorAction, AgentSelectorState}; use super::command_menu::{CommandMenuSelection, CommandMenuState}; use super::command_palette::{CommandPaletteState, PaletteAction}; +use super::fork_selector::{ForkAction, ForkSelectorState}; use super::login_form::{LoginFormAction, LoginFormState}; use super::markdown::MarkdownRenderer; use super::mcp_add_dialog::{McpAddAction, McpAddDialogState}; @@ -37,6 +38,7 @@ pub(crate) enum PopupType { ModelSelector, AgentSelector, SessionSelector, + ForkSelector, SkillSelector, SubagentSelector, McpSelector, @@ -145,6 +147,8 @@ pub(crate) struct ChatView { agent_selector: AgentSelectorState, /// Session selector popup state session_selector: SessionSelectorState, + /// OpenCode-compatible session fork-point selector. + fork_selector: ForkSelectorState, /// Skill selector popup state skill_selector: SkillSelectorState, /// Subagent selector popup state @@ -260,6 +264,7 @@ impl ChatView { model_selector: ModelSelectorState::new(), agent_selector: AgentSelectorState::new(), session_selector: SessionSelectorState::new(), + fork_selector: ForkSelectorState::new(), skill_selector: SkillSelectorState::new(), subagent_selector: SubagentSelectorState::new(), mcp_selector: McpSelectorState::new(), diff --git a/src/apps/cli/src/ui/command_palette.rs b/src/apps/cli/src/ui/command_palette.rs index 7e90c751cb..2e608608ae 100644 --- a/src/apps/cli/src/ui/command_palette.rs +++ b/src/apps/cli/src/ui/command_palette.rs @@ -42,6 +42,7 @@ pub(crate) enum PaletteAction { const DEFAULT_ITEM_ORDER: &[&str] = &[ "new_session", "sessions", + "fork_session", "compact_session", "usage", "toggle_auto_approve", diff --git a/src/apps/cli/src/ui/fork_selector.rs b/src/apps/cli/src/ui/fork_selector.rs new file mode 100644 index 0000000000..50ca8f18c5 --- /dev/null +++ b/src/apps/cli/src/ui/fork_selector.rs @@ -0,0 +1,233 @@ +use crossterm::event::{KeyCode, KeyEvent}; +use ratatui::{ + layout::Rect, + style::{Modifier, Style}, + text::{Line, Span}, + widgets::{Block, Borders, Clear, List, ListItem, ListState, Paragraph}, + Frame, +}; +use std::time::SystemTime; + +use crate::chat_state::SessionForkPoint; +use crate::ui::theme::{StyleKind, Theme}; + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum ForkTarget { + FullSession, + BeforeTurn { turn_id: String, prompt: String }, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) enum ForkAction { + None, + Select(ForkTarget), + Close, +} + +pub(super) struct ForkSelectorState { + points: Vec, + list_state: ListState, + visible: bool, +} + +impl ForkSelectorState { + pub(super) fn new() -> Self { + Self { + points: Vec::new(), + list_state: ListState::default(), + visible: false, + } + } + + pub(super) fn show(&mut self, points: Vec) { + self.points = points; + self.list_state.select(Some(0)); + self.visible = true; + } + + pub(super) fn hide(&mut self) { + self.visible = false; + } + + pub(super) fn reshow(&mut self) { + self.visible = true; + } + + pub(super) fn is_visible(&self) -> bool { + self.visible + } + + pub(super) fn handle_key_event(&mut self, key: KeyEvent) -> ForkAction { + if !self.visible { + return ForkAction::None; + } + match key.code { + KeyCode::Up => { + self.move_selection(-1); + ForkAction::None + } + KeyCode::Down => { + self.move_selection(1); + ForkAction::None + } + KeyCode::Enter => { + let target = self.selected_target(); + if target.is_some() { + self.hide(); + } + target.map(ForkAction::Select).unwrap_or(ForkAction::None) + } + KeyCode::Esc => { + self.hide(); + ForkAction::Close + } + _ => ForkAction::None, + } + } + + fn selected_target(&self) -> Option { + match self.list_state.selected()? { + 0 => Some(ForkTarget::FullSession), + index => self + .points + .get(index - 1) + .map(|point| ForkTarget::BeforeTurn { + turn_id: point.turn_id.clone(), + prompt: point.prompt.clone(), + }), + } + } + + fn move_selection(&mut self, delta: isize) { + let len = self.points.len() + 1; + let selected = self.list_state.selected().unwrap_or(0) as isize; + let next = (selected + delta).rem_euclid(len as isize) as usize; + self.list_state.select(Some(next)); + } + + pub(super) fn render(&mut self, frame: &mut Frame, area: Rect, theme: &Theme) { + if !self.visible { + return; + } + let width = area.width.saturating_sub(4).min(78); + let height = (self.points.len() as u16 + 5).min(area.height.saturating_sub(2)); + if width < 24 || height < 5 { + return; + } + let popup = Rect { + x: area.x + area.width.saturating_sub(width) / 2, + y: area.y + area.height.saturating_sub(height) / 2, + width, + height, + }; + let preview_width = width.saturating_sub(16) as usize; + let mut items = Vec::with_capacity(self.points.len() + 1); + items.push(ListItem::new(Line::from(vec![ + Span::styled( + "Full session", + theme.style(StyleKind::Primary).add_modifier(Modifier::BOLD), + ), + Span::styled(" Fork from the latest turn", theme.style(StyleKind::Muted)), + ]))); + items.extend(self.points.iter().map(|point| { + let preview = one_line_preview(&point.prompt, preview_width); + ListItem::new(Line::from(vec![ + Span::styled(preview, theme.style(StyleKind::Primary)), + Span::styled( + format!(" {}", relative_time(point.timestamp)), + theme.style(StyleKind::Muted), + ), + ])) + })); + + let list = List::new(items) + .block( + Block::default() + .borders(Borders::ALL) + .border_style(theme.style(StyleKind::Primary)) + .style(Style::default().bg(theme.background)) + .title(" Fork Session "), + ) + .highlight_style( + Style::default() + .bg(theme.primary) + .fg(theme.selection_foreground()) + .add_modifier(Modifier::BOLD), + ); + frame.render_widget(Clear, popup); + frame.render_stateful_widget(list, popup, &mut self.list_state); + + let hint_y = popup.y + popup.height; + if hint_y < area.y + area.height { + frame.render_widget( + Paragraph::new(" Up/Down: Navigate Enter: Fork Esc: Close ") + .style(theme.style(StyleKind::Muted)), + Rect { + x: popup.x, + y: hint_y, + width: popup.width, + height: 1, + }, + ); + } + } +} + +fn one_line_preview(prompt: &str, max_chars: usize) -> String { + let normalized = prompt.split_whitespace().collect::>().join(" "); + if normalized.chars().count() <= max_chars { + return normalized; + } + let mut preview = normalized + .chars() + .take(max_chars.saturating_sub(1)) + .collect::(); + preview.push('…'); + preview +} + +fn relative_time(timestamp: SystemTime) -> String { + let seconds = timestamp.elapsed().unwrap_or_default().as_secs(); + match seconds { + 0..=59 => "now".to_string(), + 60..=3_599 => format!("{}m ago", seconds / 60), + 3_600..=86_399 => format!("{}h ago", seconds / 3_600), + _ => format!("{}d ago", seconds / 86_400), + } +} + +#[cfg(test)] +mod tests { + use super::{ForkAction, ForkSelectorState, ForkTarget}; + use crate::chat_state::SessionForkPoint; + use crossterm::event::{KeyCode, KeyEvent, KeyModifiers}; + use std::time::SystemTime; + + #[test] + fn full_session_is_first_then_prompts_keep_their_supplied_order() { + let mut selector = ForkSelectorState::new(); + selector.show(vec![SessionForkPoint { + turn_id: "turn-newest".to_string(), + prompt: "Newest prompt".to_string(), + timestamp: SystemTime::now(), + }]); + + assert_eq!( + selector.handle_key_event(KeyEvent::new(KeyCode::Enter, KeyModifiers::NONE)), + ForkAction::Select(ForkTarget::FullSession) + ); + selector.show(vec![SessionForkPoint { + turn_id: "turn-newest".to_string(), + prompt: "Newest prompt".to_string(), + timestamp: SystemTime::now(), + }]); + selector.handle_key_event(KeyEvent::new(KeyCode::Down, KeyModifiers::NONE)); + assert_eq!( + selector.handle_key_event(KeyEvent::new(KeyCode::Enter, KeyModifiers::NONE)), + ForkAction::Select(ForkTarget::BeforeTurn { + turn_id: "turn-newest".to_string(), + prompt: "Newest prompt".to_string(), + }) + ); + } +} diff --git a/src/apps/cli/src/ui/mod.rs b/src/apps/cli/src/ui/mod.rs index 40e307ee29..5a212bb9ae 100644 --- a/src/apps/cli/src/ui/mod.rs +++ b/src/apps/cli/src/ui/mod.rs @@ -6,6 +6,7 @@ pub(crate) mod chat; pub(crate) mod command_menu; pub(crate) mod command_palette; mod diff_render; +pub(crate) mod fork_selector; pub(crate) mod input; pub(crate) mod login_form; mod markdown; diff --git a/src/apps/cli/src/ui/startup.rs b/src/apps/cli/src/ui/startup.rs index 46edeaf7fe..35d8da9d5b 100644 --- a/src/apps/cli/src/ui/startup.rs +++ b/src/apps/cli/src/ui/startup.rs @@ -1067,6 +1067,7 @@ impl StartupPage { ActionHandler::ClosePopups => self.close_all_popups(), ActionHandler::NavigateBack => self.navigate_back(), ActionHandler::RenameSession + | ActionHandler::ForkSession | ActionHandler::Reload | ActionHandler::Tools | ActionHandler::Extensions diff --git a/src/crates/adapters/agent-runtime-ipc/AGENTS-CN.md b/src/crates/adapters/agent-runtime-ipc/AGENTS-CN.md index 78cff8b520..3a680b345e 100644 --- a/src/crates/adapters/agent-runtime-ipc/AGENTS-CN.md +++ b/src/crates/adapters/agent-runtime-ipc/AGENTS-CN.md @@ -15,7 +15,7 @@ ## 边界 - 只导出 CLI adapter 实际使用的 workspace-private API,且 crate 不得发布,也不得把 wire 作为 SDK 合同。 -- 封闭 operation 范围为 Health、Session list/create/restore/delete(restore 结果包含 transcript)、当前 Session rename、Agent mode/model update 和手动 context compaction、声明式上下文 reload、Turn submit/cancel、pending/respond Permission 和 UserInput answers。delete 只允许作用于未被任何 Client 控制的空闲 Session。手动 compaction 要求当前 controller 且 Session 空闲;Client 在准入前提供精确 Turn ID,使超时或断连 cleanup 可以取消同一个 owned task;Core 开始原子 context commit 后,晚到取消不能暴露错误的空闲状态。上下文 reload 可在活动 Turn 中执行,不改写该 Turn,并通过缓存保护保证下一条消息重新读取已失效的 instructions。断连 cleanup 属于内部生命周期,不是 detach operation。模型目录和默认值仍是 wire 之外的产品配置;禁止顺带加入 archive、fork、replay、observer、controller transfer、Tool/MCP/Hook 管理或其他产品配置。 +- 封闭 operation 范围为 Health、Session list/create/restore/delete/fork(restore/fork 结果包含 transcript)、当前 Session rename、Agent mode/model update 和手动 context compaction、声明式上下文 reload、Turn submit/cancel、pending/respond Permission 和 UserInput answers。delete 只允许作用于未被任何 Client 控制的空闲 Session。fork 要求当前 controller 且 Session 空闲:可以复制到最新持久化 Turn,也可以停在显式选中 Turn 之前;只有包含新 Session 与 transcript 的成功结果完成编码后,Server 才能把连接 lease 从源 Session 原子切换到 fork。手动 compaction 要求当前 controller 且 Session 空闲;Client 在准入前提供精确 Turn ID,使超时或断连 cleanup 可以取消同一个 owned task;Core 开始原子 context commit 后,晚到取消不能暴露错误的空闲状态。上下文 reload 可在活动 Turn 中执行,不改写该 Turn,并通过缓存保护保证下一条消息重新读取已失效的 instructions。断连 cleanup 属于内部生命周期,不是 detach operation。模型目录和默认值仍是 wire 之外的产品配置;禁止顺带加入 archive、replay、observer、通用 controller transfer、Tool/MCP/Hook 管理或其他产品配置。 - 可以复用稳定 Event、Product Domain 和 Runtime Port DTO。禁止依赖 `bitfun-core`、Agent Runtime 实现、SDK Host、services、Tauri、terminal、tool runtime 或远程 transport。 - 只使用 Windows Named Pipe 或 Unix Domain Socket;禁止 TCP、HTTP、WebSocket、浏览器访问或远程 fallback。 - 这是本机同用户隔离,不是沙箱。未来产品 composition 必须提供当前用户私有 runtime 目录。 diff --git a/src/crates/adapters/agent-runtime-ipc/AGENTS.md b/src/crates/adapters/agent-runtime-ipc/AGENTS.md index 2df15f6f79..9bc72bf7b4 100644 --- a/src/crates/adapters/agent-runtime-ipc/AGENTS.md +++ b/src/crates/adapters/agent-runtime-ipc/AGENTS.md @@ -22,13 +22,14 @@ session controller leases, event delivery, connection bounds, and cleanup. It is - Export only the exact workspace-private API needed by the CLI adapter. Do not publish this crate or expose its wire as an SDK contract. -- The closed operation budget is Health, Session list/create/restore/delete (including transcript on restore), current-Session rename, Agent mode/model update and manual context compaction, +- The closed operation budget is Health, Session list/create/restore/delete/fork (including transcript on restore/fork), current-Session rename, Agent mode/model update and manual context compaction, declarative context reload, Turn submit/cancel, pending/respond Permission, and UserInput answers. Delete is limited to an idle Session not controlled by any client. + Fork is a current-controller, idle-only operation. It either copies through the latest persisted Turn or stops immediately before an explicitly selected Turn. The encoded success result carries the authoritative new Session and transcript; only then may the server atomically switch the connection lease from the source Session to the fork. Manual compaction is a current-controller, idle-only Turn operation. The client supplies its exact Turn ID before admission so timeout or disconnect cleanup can cancel the same owned task; once Core begins the atomic context commit, a late cancellation does not expose a false idle state. Context reload may run during an active Turn, does not rewrite that Turn, and guards the cache so the next message reads invalidated instructions. Disconnect cleanup is internal lifecycle, not a detach operation. - Model catalogs and defaults remain product configuration outside this wire. Do not add archive, fork, replay, observer, - controller transfer, Tool/MCP/Hook management, or other product configuration incidentally. + Model catalogs and defaults remain product configuration outside this wire. Do not add archive, replay, observer, + general controller transfer, Tool/MCP/Hook management, or other product configuration incidentally. - Stable Event, Product Domain, and Runtime Port DTOs may be reused. Do not depend on `bitfun-core`, Agent Runtime implementations, SDK Host, services, Tauri, terminal, tool runtime, or remote transports. diff --git a/src/crates/adapters/agent-runtime-ipc/src/lib.rs b/src/crates/adapters/agent-runtime-ipc/src/lib.rs index 8830ad908f..39a991ca4b 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/lib.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/lib.rs @@ -25,8 +25,8 @@ pub use handler::RuntimeIpcRequestHandler; pub use ipc::RuntimeIpcTransportError; pub(crate) use ipc::{LocalIpcEndpoint, LocalIpcListener, LocalIpcStream}; pub use operation::{ - RuntimeIpcOperation, RuntimeIpcOperationResult, RuntimeSessionRenameRequest, - RuntimeSessionRestoreRequest, RuntimeUserAnswersRequest, + RuntimeIpcOperation, RuntimeIpcOperationResult, RuntimeSessionForkRequest, + RuntimeSessionRenameRequest, RuntimeSessionRestoreRequest, RuntimeUserAnswersRequest, }; pub use protocol::{ HealthResult, InitializeRequest, InitializeResult, RuntimeIpcCapabilities, RuntimeIpcError, diff --git a/src/crates/adapters/agent-runtime-ipc/src/operation.rs b/src/crates/adapters/agent-runtime-ipc/src/operation.rs index ed148727d7..1a347a5186 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/operation.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/operation.rs @@ -21,6 +21,16 @@ pub struct RuntimeSessionRenameRequest { pub session_name: String, } +/// Forks the controlled Session at its latest persisted turn, or immediately +/// before `before_turn_id` when the TUI selected a historical user prompt. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase", deny_unknown_fields)] +pub struct RuntimeSessionForkRequest { + pub session_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub before_turn_id: Option, +} + #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(rename_all = "camelCase", deny_unknown_fields)] pub struct RuntimeUserAnswersRequest { @@ -59,6 +69,9 @@ pub enum RuntimeIpcOperation { RenameSession { request: RuntimeSessionRenameRequest, }, + ForkSession { + request: RuntimeSessionForkRequest, + }, ReloadSessionContext { request: AgentContextReloadRequest, }, @@ -92,6 +105,7 @@ impl RuntimeIpcOperation { Self::UpdateSessionMode { request } => Some(&request.session_id), Self::UpdateSessionModel { request } => Some(&request.session_id), Self::RenameSession { request } => Some(&request.session_id), + Self::ForkSession { request } => Some(&request.session_id), Self::ReloadSessionContext { request } => Some(&request.session_id), Self::CompactSession { request } => Some(&request.session_id), Self::SubmitTurn { request } => Some(&request.session_id), @@ -126,6 +140,9 @@ impl RuntimeIpcOperation { | Self::SubmitTurn { .. } => { RuntimeIpcOperationRules::new(CurrentController, true, false, true) } + Self::ForkSession { .. } => { + RuntimeIpcOperationRules::new(CurrentController, true, true, true) + } Self::ReloadSessionContext { .. } | Self::CancelTurn { .. } | Self::RespondPermission { .. } @@ -195,6 +212,10 @@ pub enum RuntimeIpcOperationResult { transcript: SessionTranscript, pending_permissions: Vec, }, + SessionForked { + session: AgentSessionSummary, + transcript: SessionTranscript, + }, TurnAccepted { session_id: String, turn_id: String, diff --git a/src/crates/adapters/agent-runtime-ipc/src/protocol.rs b/src/crates/adapters/agent-runtime-ipc/src/protocol.rs index 8e67b98576..0c6cb94fd6 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/protocol.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/protocol.rs @@ -5,7 +5,7 @@ use crate::{RuntimeIpcOperation, RuntimeIpcOperationResult}; use bitfun_events::AgenticEventEnvelope; use bitfun_product_domains::tool_permissions::PermissionRequestEvent; -pub const PROTOCOL_VERSION: u32 = 8; +pub const PROTOCOL_VERSION: u32 = 9; #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] #[serde(tag = "type", rename_all = "snake_case", deny_unknown_fields)] diff --git a/src/crates/adapters/agent-runtime-ipc/src/server.rs b/src/crates/adapters/agent-runtime-ipc/src/server.rs index 3a4f92a144..432bf04aa6 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/server.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/server.rs @@ -548,21 +548,30 @@ async fn run_initialized_connection( let RuntimeIpcFrame::Response { result, .. } = &response else { unreachable!("response frame was just constructed") }; - if let RuntimeIpcOperationResult::SessionCreated { session } = result { - lease_transition = - match config.leases.switch(connection_id, &session.session_id) { - Ok(transition) => transition, - Err(error) => { - send_runtime_error( - stream, - config.request_timeout, - Some(request_id), - error, - ) - .await?; - continue; - } - }; + let created_session_id = match result { + RuntimeIpcOperationResult::SessionCreated { session } => { + Some(session.session_id.as_str()) + } + RuntimeIpcOperationResult::SessionForked { session, .. } => { + Some(session.session_id.as_str()) + } + _ => None, + }; + if let Some(created_session_id) = created_session_id { + lease_transition = match config.leases.switch(connection_id, created_session_id) + { + Ok(transition) => transition, + Err(error) => { + send_runtime_error( + stream, + config.request_timeout, + Some(request_id), + error, + ) + .await?; + continue; + } + }; } let mut event_stream_unavailable = false; if let Some(session_id) = match result { @@ -572,6 +581,9 @@ async fn run_initialized_connection( RuntimeIpcOperationResult::SessionRestored { session, .. } => { Some(session.session_id.as_str()) } + RuntimeIpcOperationResult::SessionForked { session, .. } => { + Some(session.session_id.as_str()) + } _ => None, } { match handler.subscribe_events(session_id) { diff --git a/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs b/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs index c91e2c9789..e34900b276 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/tests/protocol_contracts.rs @@ -1,8 +1,8 @@ use crate::operation::RuntimeIpcSessionRequirement; use crate::{ serialize_frame_with_limit, InitializeRequest, RuntimeIpcError, RuntimeIpcErrorCode, - RuntimeIpcFrame, RuntimeIpcOperation, RuntimeSessionRenameRequest, RuntimeUserAnswersRequest, - MAX_REQUEST_FRAME_BYTES, PROTOCOL_VERSION, + RuntimeIpcFrame, RuntimeIpcOperation, RuntimeSessionForkRequest, RuntimeSessionRenameRequest, + RuntimeUserAnswersRequest, MAX_REQUEST_FRAME_BYTES, PROTOCOL_VERSION, }; use bitfun_product_domains::tool_permissions::PermissionReply; @@ -118,7 +118,7 @@ fn protocol_round_trips_the_reviewed_session_model_operation() { #[test] fn protocol_round_trips_the_current_session_rename_operation() { - assert_eq!(PROTOCOL_VERSION, 8); + assert_eq!(PROTOCOL_VERSION, 9); let operation = RuntimeIpcOperation::RenameSession { request: RuntimeSessionRenameRequest { @@ -149,6 +149,41 @@ fn protocol_round_trips_the_current_session_rename_operation() { ); } +#[test] +fn protocol_round_trips_fork_as_an_atomic_idle_controller_transition() { + let operation = RuntimeIpcOperation::ForkSession { + request: RuntimeSessionForkRequest { + session_id: "session-1".to_string(), + before_turn_id: Some("turn-2".to_string()), + }, + }; + + let encoded = serde_json::to_value(&operation).expect("serialize session fork"); + assert_eq!( + encoded, + json!({ + "operation": "fork_session", + "request": { + "sessionId": "session-1", + "beforeTurnId": "turn-2" + } + }) + ); + let decoded: RuntimeIpcOperation = + serde_json::from_value(encoded).expect("deserialize session fork"); + + assert_eq!(decoded, operation); + assert_eq!(decoded.session_id(), Some("session-1")); + let rules = decoded.rules(); + assert_eq!( + rules.session_requirement, + RuntimeIpcSessionRequirement::CurrentController + ); + assert!(rules.requires_idle); + assert!(rules.serializes_session_selection); + assert!(rules.side_effecting); +} + #[test] fn protocol_round_trips_manual_compaction_as_an_idle_controller_turn() { let operation = RuntimeIpcOperation::CompactSession { diff --git a/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs b/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs index 96faefab89..c44b2bec12 100644 --- a/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs +++ b/src/crates/adapters/agent-runtime-ipc/src/tests/shared_controller.rs @@ -2,8 +2,8 @@ use crate::{ read_frame, write_frame, InitializeRequest, LocalIpcStream, RuntimeInstanceIdentity, RuntimeIpcClient, RuntimeIpcClientError, RuntimeIpcError, RuntimeIpcErrorCode, RuntimeIpcEvent, RuntimeIpcFrame, RuntimeIpcOperation, RuntimeIpcOperationResult, RuntimeIpcRequestHandler, - RuntimeIpcServer, RuntimeIpcServerConfig, RuntimeSessionRenameRequest, - RuntimeSessionRestoreRequest, PROTOCOL_VERSION, + RuntimeIpcServer, RuntimeIpcServerConfig, RuntimeSessionForkRequest, + RuntimeSessionRenameRequest, RuntimeSessionRestoreRequest, PROTOCOL_VERSION, }; use async_trait::async_trait; use bitfun_events::{AgenticEvent, AgenticEventEnvelope, AgenticEventPriority}; @@ -203,6 +203,15 @@ impl RuntimeIpcRequestHandler for FakeHandler { } match operation { RuntimeIpcOperation::RestoreSession { request } => Ok(restored(&request.session_id)), + RuntimeIpcOperation::ForkSession { .. } => { + Ok(RuntimeIpcOperationResult::SessionForked { + session: summary("session-fork"), + transcript: SessionTranscript { + session_id: "session-fork".to_string(), + messages: Vec::new(), + }, + }) + } RuntimeIpcOperation::SubmitTurn { request } => { if let Some(delay) = self.submit_delay { tokio::time::sleep(delay).await; @@ -555,6 +564,15 @@ fn rename_operation(session_id: &str, session_name: &str) -> RuntimeIpcOperation } } +fn fork_operation(session_id: &str, before_turn_id: Option<&str>) -> RuntimeIpcOperation { + RuntimeIpcOperation::ForkSession { + request: RuntimeSessionForkRequest { + session_id: session_id.to_string(), + before_turn_id: before_turn_id.map(str::to_string), + }, + } +} + fn delete_operation(session_id: &str) -> RuntimeIpcOperation { RuntimeIpcOperation::DeleteSession { session_id: session_id.to_string(), @@ -709,6 +727,53 @@ async fn session_switching_is_exclusive_and_disconnect_releases_control() { server.finish().await; } +#[tokio::test] +async fn successful_fork_atomically_transfers_control_and_releases_the_source() { + let server = TestServer::start(server_config(), Arc::new(FakeHandler::default())).await; + let mut forker = server.connect("fork-controller").await; + let mut observer = server.connect("source-controller").await; + + expect_response( + &mut forker, + 2, + restore_operation(server.workspace.path(), "session-a"), + ) + .await; + assert!(matches!( + request(&mut forker, 3, fork_operation("session-a", Some("turn-2"))).await, + RuntimeIpcFrame::Response { + result: RuntimeIpcOperationResult::SessionForked { session, .. }, + .. + } if session.session_id == "session-fork" + )); + + expect_response( + &mut observer, + 2, + restore_operation(server.workspace.path(), "session-a"), + ) + .await; + expect_error( + &mut forker, + 4, + update_mode_operation("session-a", "ask"), + RuntimeIpcErrorCode::SessionMismatch, + ) + .await; + expect_response(&mut forker, 5, update_mode_operation("session-fork", "ask")).await; + expect_error( + &mut observer, + 3, + restore_operation(server.workspace.path(), "session-fork"), + RuntimeIpcErrorCode::SessionInUse, + ) + .await; + + drop(forker); + drop(observer); + server.finish().await; +} + #[tokio::test] async fn one_connection_rejects_a_second_turn_until_the_first_finishes() { let handler = Arc::new(FakeHandler::default()); @@ -740,6 +805,13 @@ async fn one_connection_rejects_a_second_turn_until_the_first_finishes() { RuntimeIpcErrorCode::SessionInUse, ) .await; + expect_error( + &mut client, + 6, + fork_operation("session-a", None), + RuntimeIpcErrorCode::SessionInUse, + ) + .await; drop(client); wait_for_calls(&handler, |calls| { @@ -755,7 +827,10 @@ async fn one_connection_rejects_a_second_turn_until_the_first_finishes() { && request.turn_id.as_deref() == Some("turn-a") ) }); - submitted == 1 && cancelled_first + let forked = calls + .iter() + .any(|call| matches!(call, RuntimeIpcOperation::ForkSession { .. })); + submitted == 1 && !forked && cancelled_first }) .await; server.finish().await; 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 52490c387e..89ede815fb 100644 --- a/src/crates/assembly/core/src/agentic/persistence/session_branch.rs +++ b/src/crates/assembly/core/src/agentic/persistence/session_branch.rs @@ -5,6 +5,7 @@ use bitfun_services_core::session::{ build_branched_session_metadata, format_branch_session_name, resolve_branch_session_lineage, BranchSessionMetadataFacts, }; +use bitfun_services_core::session::SessionBranchBoundary; pub use bitfun_services_core::session::{SessionBranchRequest, SessionBranchResult}; use std::path::Path; use std::time::{SystemTime, UNIX_EPOCH}; @@ -64,6 +65,7 @@ impl PersistenceManager { request.source_turn_id )) })?; + let copied_turn_count = request.boundary.copied_turn_count(source_turn_index); let target_session_name = format_branch_session_name(&branch_lineage.base_session_name, branch_lineage.ordinal); @@ -77,7 +79,9 @@ impl PersistenceManager { target_session.created_by = None; target_session.kind = SessionKind::Standard; target_session.snapshot_session_id = None; - target_session.compression_state = source_session.compression_state.clone(); + if copied_turn_count > 0 && request.boundary == SessionBranchBoundary::ThroughTurn { + target_session.compression_state = source_session.compression_state.clone(); + } let target_session_id = target_session.session_id.clone(); let _target_session_write = self.lock_session_write_operation(workspace_path, &target_session_id)?; @@ -87,7 +91,7 @@ impl PersistenceManager { let branch_result = async { let branched_turns = source_turns .iter() - .take(source_turn_index + 1) + .take(copied_turn_count) .enumerate() .map(|(new_index, turn)| { let mut branched_turn = turn.clone(); @@ -97,8 +101,7 @@ impl PersistenceManager { }) .collect::>(); - for (new_index, source_turn) in - source_turns.iter().take(source_turn_index + 1).enumerate() + for (new_index, source_turn) in source_turns.iter().take(copied_turn_count).enumerate() { if let Some(messages) = self .load_turn_context_snapshot( @@ -138,13 +141,15 @@ impl PersistenceManager { self.save_dialog_turn(workspace_path, turn).await?; } - self.copy_compression_transcripts_through( - workspace_path, - &request.source_session_id, - &target_session_id, - source_turn_index, - ) - .await?; + if let Some(last_copied_turn_index) = copied_turn_count.checked_sub(1) { + self.copy_compression_transcripts_through( + workspace_path, + &request.source_session_id, + &target_session_id, + last_copied_turn_index, + ) + .await?; + } if let Some(cache) = source_prompt_cache.as_ref() { self.save_prompt_cache(workspace_path, &target_session_id, cache) @@ -177,6 +182,7 @@ impl PersistenceManager { source_session_id: &request.source_session_id, source_turn_id: &request.source_turn_id, source_turn_index, + boundary: request.boundary, branched_turns: &branched_turns, branch_lineage: &branch_lineage, now_ms, @@ -206,7 +212,7 @@ impl PersistenceManager { #[cfg(test)] mod tests { - use super::{PersistenceManager, SessionBranchRequest}; + use super::{PersistenceManager, SessionBranchBoundary, SessionBranchRequest}; use crate::agentic::core::{Message, Session, SessionKind}; use crate::agentic::session::{ CachedSystemPrompt, CachedUserContext, SessionPromptCache, SystemPromptCacheIdentity, @@ -424,6 +430,7 @@ mod tests { &SessionBranchRequest { source_session_id: source_session.session_id.clone(), source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await @@ -521,6 +528,148 @@ mod tests { ); } + #[tokio::test] + async fn branch_session_before_first_turn_creates_an_empty_replayable_session() { + let workspace = TestWorkspace::new(); + let manager = + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"); + let source_session = Session::new( + "Source Title".to_string(), + "agentic".to_string(), + Default::default(), + ); + manager + .save_session(workspace.path(), &source_session) + .await + .expect("source session should save"); + manager + .save_dialog_turn( + workspace.path(), + &build_turn(&source_session.session_id, "turn-0", 0, "first"), + ) + .await + .expect("source turn should save"); + + let result = manager + .branch_session( + workspace.path(), + &SessionBranchRequest { + source_session_id: source_session.session_id.clone(), + source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::BeforeTurn, + }, + ) + .await + .expect("before-turn branch should succeed"); + + let branched_turns = manager + .load_session_turns(workspace.path(), &result.session_id) + .await + .expect("branched turns should load"); + assert!(branched_turns.is_empty()); + + let branched_metadata = manager + .load_session_metadata(workspace.path(), &result.session_id) + .await + .expect("branched metadata should load") + .expect("branched metadata should exist"); + assert_eq!(branched_metadata.turn_count, 0); + assert_eq!( + branched_metadata.custom_metadata.unwrap()["forkOrigin"], + serde_json::json!({ + "sessionId": source_session.session_id, + "turnId": "turn-0", + "turnIndex": 0, + "boundary": "before_turn", + "baseTitle": "Source Title" + }) + ); + } + + #[tokio::test] + async fn branch_session_before_middle_turn_copies_only_earlier_turn_state() { + let workspace = TestWorkspace::new(); + let manager = + PersistenceManager::new(workspace.path_manager()).expect("persistence manager"); + let mut source_session = Session::new( + "Source Title".to_string(), + "agentic".to_string(), + Default::default(), + ); + source_session.compression_state.compression_count = 3; + manager + .save_session(workspace.path(), &source_session) + .await + .expect("source session should save"); + for (index, turn_id) in ["turn-0", "turn-1", "turn-2"].into_iter().enumerate() { + manager + .save_dialog_turn( + workspace.path(), + &build_turn( + &source_session.session_id, + turn_id, + index, + &format!("prompt-{index}"), + ), + ) + .await + .expect("source turn should save"); + manager + .save_turn_context_snapshot( + workspace.path(), + &source_session.session_id, + index, + &[Message::user(format!("snapshot-{index}"))], + ) + .await + .expect("source snapshot should save"); + } + + let result = manager + .branch_session( + workspace.path(), + &SessionBranchRequest { + source_session_id: source_session.session_id.clone(), + source_turn_id: "turn-1".to_string(), + boundary: SessionBranchBoundary::BeforeTurn, + }, + ) + .await + .expect("before-turn branch should succeed"); + + let turns = manager + .load_session_turns(workspace.path(), &result.session_id) + .await + .expect("branched turns should load"); + assert_eq!(turns.len(), 1); + assert_eq!(turns[0].turn_id, "turn-0"); + assert!(manager + .load_turn_context_snapshot(workspace.path(), &result.session_id, 0) + .await + .expect("snapshot load") + .is_some()); + assert!(manager + .load_turn_context_snapshot(workspace.path(), &result.session_id, 1) + .await + .expect("snapshot load") + .is_none()); + let metadata = manager + .load_session_metadata(workspace.path(), &result.session_id) + .await + .expect("metadata load") + .expect("metadata exists"); + assert_eq!(metadata.turn_count, 1); + assert_eq!( + metadata.custom_metadata.unwrap()["forkOrigin"]["turnIndex"], + 1 + ); + let forked_session = manager + .load_session(workspace.path(), &result.session_id) + .await + .expect("forked session load"); + assert_eq!(forked_session.compression_state.compression_count, 0); + } + #[tokio::test] async fn branch_session_advances_the_family_title_without_growing_suffixes() { let workspace = TestWorkspace::new(); @@ -550,6 +699,7 @@ mod tests { &SessionBranchRequest { source_session_id: source_session.session_id.clone(), source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await @@ -562,6 +712,7 @@ mod tests { &SessionBranchRequest { source_session_id: first_branch.session_id, source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await @@ -574,6 +725,7 @@ mod tests { &SessionBranchRequest { source_session_id: source_session.session_id.clone(), source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await @@ -626,6 +778,7 @@ mod tests { &SessionBranchRequest { source_session_id: source_session.session_id.clone(), source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await @@ -645,6 +798,7 @@ mod tests { &SessionBranchRequest { source_session_id: first_branch.session_id.clone(), source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await @@ -664,6 +818,7 @@ mod tests { &SessionBranchRequest { source_session_id: first_branch.session_id, source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await @@ -722,10 +877,12 @@ mod tests { let first_request = SessionBranchRequest { source_session_id: source_session_id.clone(), source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }; let second_request = SessionBranchRequest { source_session_id: second_source_session.session_id, source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }; let (first_result, second_result) = tokio::join!( first_manager.branch_session(workspace.path(), &first_request), @@ -796,6 +953,7 @@ mod tests { &SessionBranchRequest { source_session_id: source_session.session_id.clone(), source_turn_id: "turn-0".to_string(), + boundary: SessionBranchBoundary::ThroughTurn, }, ) .await diff --git a/src/crates/assembly/core/src/product_runtime.rs b/src/crates/assembly/core/src/product_runtime.rs index c338255df8..aa2d766e94 100644 --- a/src/crates/assembly/core/src/product_runtime.rs +++ b/src/crates/assembly/core/src/product_runtime.rs @@ -14,8 +14,9 @@ use std::time::{Duration, SystemTime, UNIX_EPOCH}; use bitfun_agent_runtime::permission::PermissionRequestManager; use bitfun_agent_runtime::sdk::{ AgentEventReceiver, AgentEventSource, AgentRuntime, AgentSessionForkAtTurnRequest, - AgentSessionForkPort, AgentSessionForkRequest, AgentSessionForkResult, AgentSessionUsagePort, - AgentSessionUsageRequest, AgentTurnSettlementPort, AgentTurnSettlementRequest, + AgentSessionForkBeforeTurnRequest, AgentSessionForkPort, AgentSessionForkRequest, + AgentSessionForkResult, AgentSessionUsagePort, AgentSessionUsageRequest, + AgentTurnSettlementPort, AgentTurnSettlementRequest, }; use bitfun_harness::HarnessRegistry; use bitfun_runtime_ports::{ @@ -27,6 +28,7 @@ use bitfun_runtime_ports::{ }; use bitfun_runtime_services::RuntimeServices; use bitfun_services_core::permission_store::ProjectPermissionSqliteStore; +use bitfun_services_core::session::SessionBranchBoundary; use crate::agentic::coordination::{ ConversationCoordinator, DialogScheduler, DialogSteerOutcome, SessionMaintenancePermit, @@ -1029,6 +1031,7 @@ impl CoreSessionOperationsPort { storage_path: &Path, source_session_id: String, source_turn_id: String, + boundary: SessionBranchBoundary, ) -> PortResult { if source_turn_id.trim().is_empty() { return Err(PortError::new( @@ -1044,6 +1047,7 @@ impl CoreSessionOperationsPort { &SessionBranchRequest { source_session_id, source_turn_id, + boundary, }, ) .await @@ -1091,7 +1095,17 @@ fn runtime_port_error(error: BitFunError) -> PortError { } fn validate_latest_turn_fork_scope(request: &AgentSessionForkRequest) -> PortResult<()> { - if request.remote_connection_id.is_some() || request.remote_ssh_host.is_some() { + validate_local_fork_scope( + request.remote_connection_id.as_deref(), + request.remote_ssh_host.as_deref(), + ) +} + +fn validate_local_fork_scope( + remote_connection_id: Option<&str>, + remote_ssh_host: Option<&str>, +) -> PortResult<()> { + if remote_connection_id.is_some() || remote_ssh_host.is_some() { return Err(PortError::new( PortErrorKind::NotAvailable, "Remote session fork is not supported by the local CLI runtime", @@ -1129,14 +1143,54 @@ impl AgentSessionForkPort for CoreSessionOperationsPort { .await .map_err(runtime_port_error)?; let source_turn_id = latest_persisted_turn_id(&turns).map_err(runtime_port_error)?; - self.fork_at_persisted_turn(&storage_path, source_session_id, source_turn_id) - .await + self.fork_at_persisted_turn( + &storage_path, + source_session_id, + source_turn_id, + SessionBranchBoundary::ThroughTurn, + ) + .await } async fn fork_session_at_turn( &self, request: AgentSessionForkAtTurnRequest, ) -> PortResult { + validate_local_fork_scope( + request.remote_connection_id.as_deref(), + request.remote_ssh_host.as_deref(), + )?; + self.coordinator + .ensure_workspace_runtime_ownership( + Path::new(&request.workspace_path), + request.remote_connection_id.as_deref(), + request.remote_ssh_host.as_deref(), + ) + .map_err(runtime_port_error)?; + let storage_path = self + .resolve_fork_storage_path( + request.workspace_path, + request.remote_connection_id, + request.remote_ssh_host, + ) + .await?; + self.fork_at_persisted_turn( + &storage_path, + request.source_session_id, + request.source_turn_id, + SessionBranchBoundary::ThroughTurn, + ) + .await + } + + async fn fork_session_before_turn( + &self, + request: AgentSessionForkBeforeTurnRequest, + ) -> PortResult { + validate_local_fork_scope( + request.remote_connection_id.as_deref(), + request.remote_ssh_host.as_deref(), + )?; self.coordinator .ensure_workspace_runtime_ownership( Path::new(&request.workspace_path), @@ -1155,6 +1209,7 @@ impl AgentSessionForkPort for CoreSessionOperationsPort { &storage_path, request.source_session_id, request.source_turn_id, + SessionBranchBoundary::BeforeTurn, ) .await } @@ -1416,9 +1471,10 @@ mod tests { fork_impl .matches("ensure_workspace_runtime_ownership") .count(), - 2, - "latest-turn and explicit-turn forks must share the Coordinator ownership gate" + 3, + "latest-turn, explicit-turn, and before-turn forks must share the Coordinator ownership gate" ); + assert!(fork_impl.contains("fork_session_before_turn")); assert!(!fork_impl.contains("RuntimeOwnershipKey")); assert!(!fork_impl.contains("try_acquire")); } diff --git a/src/crates/contracts/runtime-ports/src/lib.rs b/src/crates/contracts/runtime-ports/src/lib.rs index 52dd4e6b75..4d0e9a6237 100644 --- a/src/crates/contracts/runtime-ports/src/lib.rs +++ b/src/crates/contracts/runtime-ports/src/lib.rs @@ -1243,6 +1243,23 @@ pub struct AgentSessionForkAtTurnRequest { pub remote_ssh_host: Option, } +/// Forks a session immediately before an explicitly selected persisted turn. +/// +/// The selected turn is the replayable prompt boundary and is not copied into +/// the fork. This stays separate from [`AgentSessionForkAtTurnRequest`] so its +/// inclusive behavior remains source- and behavior-compatible. +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub struct AgentSessionForkBeforeTurnRequest { + pub workspace_path: String, + pub source_session_id: String, + pub source_turn_id: String, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub remote_connection_id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub remote_ssh_host: Option, +} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct AgentSessionForkResult { @@ -2072,6 +2089,17 @@ pub trait AgentSessionForkPort: Send + Sync { "exact-turn session fork is not supported by this provider", )) } + + async fn fork_session_before_turn( + &self, + request: AgentSessionForkBeforeTurnRequest, + ) -> PortResult { + let _ = request; + Err(PortError::new( + PortErrorKind::NotAvailable, + "before-turn session fork is not supported by this provider", + )) + } } #[async_trait::async_trait] @@ -2521,6 +2549,25 @@ mod tests { assert_eq!(error.kind, PortErrorKind::NotAvailable); } + #[tokio::test] + async fn before_turn_fork_default_preserves_existing_provider_compatibility() { + let provider = LatestTurnForkOnlyProvider; + let error = AgentSessionForkPort::fork_session_before_turn( + &provider, + AgentSessionForkBeforeTurnRequest { + workspace_path: "/workspace/project".to_string(), + source_session_id: "session_1".to_string(), + source_turn_id: "turn_1".to_string(), + remote_connection_id: None, + remote_ssh_host: None, + }, + ) + .await + .expect_err("existing providers must reject before-turn fork by default"); + + assert_eq!(error.kind, PortErrorKind::NotAvailable); + } + #[test] fn agent_session_create_request_keeps_rust_literal_compatible() { let request = AgentSessionCreateRequest { @@ -3327,6 +3374,13 @@ mod tests { remote_connection_id: Some("conn-1".to_string()), remote_ssh_host: Some("host-1".to_string()), }; + let fork_before_turn_request = AgentSessionForkBeforeTurnRequest { + workspace_path: "/workspace/project".to_string(), + source_session_id: "session_1".to_string(), + source_turn_id: "turn_2".to_string(), + remote_connection_id: None, + remote_ssh_host: None, + }; let model_request = AgentSessionModelUpdateRequest { session_id: "session_1".to_string(), model_id: "provider/model".to_string(), @@ -3358,6 +3412,8 @@ mod tests { let fork_json = serde_json::to_value(fork_request).expect("serialize fork request"); let fork_at_turn_json = serde_json::to_value(fork_at_turn_request).expect("serialize exact-turn fork request"); + let fork_before_turn_json = serde_json::to_value(fork_before_turn_request) + .expect("serialize before-turn fork request"); let model_json = serde_json::to_value(model_request).expect("serialize model request"); let mode_json = serde_json::to_value(mode_request).expect("serialize mode request"); let workspace_json = @@ -3385,6 +3441,8 @@ mod tests { assert_eq!(archive_state_json["archived"], false); assert!(fork_json.get("sourceTurnId").is_none()); assert_eq!(fork_at_turn_json["sourceTurnId"], "turn_2"); + assert_eq!(fork_before_turn_json["sourceTurnId"], "turn_2"); + assert_eq!(fork_before_turn_json["sourceSessionId"], "session_1"); assert_eq!(model_json["sessionId"], "session_1"); assert_eq!(model_json["modelId"], "provider/model"); assert_eq!(mode_json["sessionId"], "session_1"); diff --git a/src/crates/execution/agent-runtime/src/runtime.rs b/src/crates/execution/agent-runtime/src/runtime.rs index 90a80dcd7c..7d94741a51 100644 --- a/src/crates/execution/agent-runtime/src/runtime.rs +++ b/src/crates/execution/agent-runtime/src/runtime.rs @@ -16,19 +16,20 @@ use bitfun_runtime_ports::{ AgentSessionArchiveStateRequest, AgentSessionClosePort, AgentSessionCompactionPort, AgentSessionCompactionRequest, AgentSessionCompactionResult, AgentSessionCreateRequest, AgentSessionCreateResult, AgentSessionDeleteRequest, AgentSessionForkAtTurnRequest, - AgentSessionForkPort, AgentSessionForkRequest, AgentSessionForkResult, AgentSessionListRequest, - AgentSessionManagementPort, AgentSessionModePort, AgentSessionModeUpdateRequest, - AgentSessionModelPort, AgentSessionModelUpdateRequest, AgentSessionRenameRequest, - AgentSessionSummary, AgentSessionUsagePort, AgentSessionUsageRequest, - AgentSessionWorkspaceBinding, AgentSessionWorkspaceRequest, AgentSubmissionPort, - AgentSubmissionRequest, AgentSubmissionResult, AgentSubmissionSource, - AgentThreadGoalCreateRequest, AgentThreadGoalDeliveryRequest, AgentThreadGoalGetRequest, - AgentThreadGoalManagementPort, AgentThreadGoalUpdateStatusRequest, - AgentTransientSessionDiscardRequest, AgentTurnCancellationPort, AgentTurnCancellationRequest, - AgentTurnCancellationResult, AgentTurnSettlementPort, AgentTurnSettlementRequest, - DialogSubmitOutcome, PermissionAuditRecord, PermissionGrant, PermissionGrantKey, - PluginRuntimeBinding, PortError, PortErrorKind, PortResult, RuntimeEventEnvelope, - SessionTranscript, SessionTranscriptReader, SessionTranscriptRequest, ThreadGoal, + AgentSessionForkBeforeTurnRequest, AgentSessionForkPort, AgentSessionForkRequest, + AgentSessionForkResult, AgentSessionListRequest, AgentSessionManagementPort, + AgentSessionModePort, AgentSessionModeUpdateRequest, AgentSessionModelPort, + AgentSessionModelUpdateRequest, AgentSessionRenameRequest, AgentSessionSummary, + AgentSessionUsagePort, AgentSessionUsageRequest, AgentSessionWorkspaceBinding, + AgentSessionWorkspaceRequest, AgentSubmissionPort, AgentSubmissionRequest, + AgentSubmissionResult, AgentSubmissionSource, AgentThreadGoalCreateRequest, + AgentThreadGoalDeliveryRequest, AgentThreadGoalGetRequest, AgentThreadGoalManagementPort, + AgentThreadGoalUpdateStatusRequest, AgentTransientSessionDiscardRequest, + AgentTurnCancellationPort, AgentTurnCancellationRequest, AgentTurnCancellationResult, + AgentTurnSettlementPort, AgentTurnSettlementRequest, DialogSubmitOutcome, + PermissionAuditRecord, PermissionGrant, PermissionGrantKey, PluginRuntimeBinding, PortError, + PortErrorKind, PortResult, RuntimeEventEnvelope, SessionTranscript, SessionTranscriptReader, + SessionTranscriptRequest, ThreadGoal, }; use bitfun_runtime_services::RuntimeServices; @@ -1121,6 +1122,21 @@ impl AgentRuntime { .map_err(RuntimeError::from) } + pub async fn fork_session_before_turn( + &self, + request: AgentSessionForkBeforeTurnRequest, + ) -> Result { + let port = self.session_fork.as_ref().ok_or_else(|| { + RuntimeError::Port(PortError::new( + PortErrorKind::NotAvailable, + "agent session fork port is not registered", + )) + })?; + port.fork_session_before_turn(request) + .await + .map_err(RuntimeError::from) + } + pub async fn generate_session_usage( &self, request: AgentSessionUsageRequest, @@ -1424,6 +1440,7 @@ mod tests { submitted_messages: Mutex>, cancelled_turns: Mutex>, compaction_requests: Mutex>, + before_turn_fork_requests: Mutex>, listed_sessions: Mutex>, deleted_sessions: Mutex>, renamed_sessions: Mutex>, @@ -1594,6 +1611,32 @@ mod tests { } } + #[async_trait::async_trait] + impl AgentSessionForkPort for FakeAgentRuntimePorts { + async fn fork_session( + &self, + request: AgentSessionForkRequest, + ) -> PortResult { + Ok(AgentSessionForkResult { + session_id: format!("{}-fork", request.source_session_id), + session_name: "Forked session".to_string(), + agent_type: "agentic".to_string(), + }) + } + + async fn fork_session_before_turn( + &self, + request: AgentSessionForkBeforeTurnRequest, + ) -> PortResult { + self.before_turn_fork_requests.lock().unwrap().push(request); + Ok(AgentSessionForkResult { + session_id: "session-fork".to_string(), + session_name: "Forked session".to_string(), + agent_type: "agentic".to_string(), + }) + } + } + #[async_trait::async_trait] impl AgentLocalCommandTurnPort for FakeAgentRuntimePorts { async fn record_completed_local_command_turn( @@ -1862,6 +1905,34 @@ mod tests { ); } + #[tokio::test] + async fn session_fork_before_turn_forwards_exact_boundary_identity() { + let ports = Arc::new(FakeAgentRuntimePorts::default()); + let runtime = AgentRuntimeBuilder::new() + .with_submission_port(ports.clone()) + .with_session_fork_port(ports.clone()) + .build() + .expect("runtime"); + let request = AgentSessionForkBeforeTurnRequest { + workspace_path: "/workspace/project".to_string(), + source_session_id: "session-1".to_string(), + source_turn_id: "turn-2".to_string(), + remote_connection_id: None, + remote_ssh_host: None, + }; + + let result = runtime + .fork_session_before_turn(request.clone()) + .await + .expect("fork before turn"); + + assert_eq!(result.session_id, "session-fork"); + assert_eq!( + ports.before_turn_fork_requests.lock().unwrap().as_slice(), + &[request] + ); + } + #[tokio::test] async fn builder_keeps_plugin_runtime_disabled_by_default() { let ports = Arc::new(FakeAgentRuntimePorts::default()); diff --git a/src/crates/execution/agent-runtime/src/sdk.rs b/src/crates/execution/agent-runtime/src/sdk.rs index 7c06aee5fd..a4a4940b6a 100644 --- a/src/crates/execution/agent-runtime/src/sdk.rs +++ b/src/crates/execution/agent-runtime/src/sdk.rs @@ -64,22 +64,23 @@ pub use bitfun_runtime_ports::{ AgentSessionArchiveStateRequest, AgentSessionClosePort, AgentSessionCompactionPort, AgentSessionCompactionRequest, AgentSessionCompactionResult, AgentSessionCreateRequest, AgentSessionCreateResult, AgentSessionDeleteRequest, AgentSessionForkAtTurnRequest, - AgentSessionForkPort, AgentSessionForkRequest, AgentSessionForkResult, AgentSessionListRequest, - AgentSessionManagementPort, AgentSessionModePort, AgentSessionModeUpdateRequest, - AgentSessionModelPort, AgentSessionModelUpdateRequest, AgentSessionRenameRequest, - AgentSessionSummary, AgentSessionUsagePort, AgentSessionUsageRequest, - AgentSessionWorkspaceBinding, AgentSessionWorkspaceRequest, AgentSubmissionPort, - AgentSubmissionRequest, AgentSubmissionResult, AgentSubmissionSource, - AgentThreadGoalCreateRequest, AgentThreadGoalDeliveryRequest, AgentThreadGoalGetRequest, - AgentThreadGoalManagementPort, AgentThreadGoalUpdateStatusRequest, - AgentTransientSessionDiscardRequest, AgentTurnCancellationPort, AgentTurnCancellationRequest, - AgentTurnCancellationResult, AgentTurnSettlementPort, AgentTurnSettlementRequest, ClockPort, - DialogSubmissionPolicy, DialogSubmitOutcome, FileSystemPort, GitPort, McpCatalogPort, - NetworkPort, PermissionAuditRecord, PermissionDelegationContext, PermissionGrant, - PermissionGrantKey, PermissionReply, PermissionReplySource, PermissionRequest, - PermissionRequestEvent, PermissionRequestSource, PermissionRequestSourceKind, PortError, - PortErrorKind, PortResult, RemoteAssistantWorkspaceFacts, RemoteCapabilityPort, - RemoteConnectionPort, RemoteProjectionPort, RemoteRecentWorkspaceFacts, RemoteWorkspaceFacts, + AgentSessionForkBeforeTurnRequest, AgentSessionForkPort, AgentSessionForkRequest, + AgentSessionForkResult, AgentSessionListRequest, AgentSessionManagementPort, + AgentSessionModePort, AgentSessionModeUpdateRequest, AgentSessionModelPort, + AgentSessionModelUpdateRequest, AgentSessionRenameRequest, AgentSessionSummary, + AgentSessionUsagePort, AgentSessionUsageRequest, AgentSessionWorkspaceBinding, + AgentSessionWorkspaceRequest, AgentSubmissionPort, AgentSubmissionRequest, + AgentSubmissionResult, AgentSubmissionSource, AgentThreadGoalCreateRequest, + AgentThreadGoalDeliveryRequest, AgentThreadGoalGetRequest, AgentThreadGoalManagementPort, + AgentThreadGoalUpdateStatusRequest, AgentTransientSessionDiscardRequest, + AgentTurnCancellationPort, AgentTurnCancellationRequest, AgentTurnCancellationResult, + AgentTurnSettlementPort, AgentTurnSettlementRequest, ClockPort, DialogSubmissionPolicy, + DialogSubmitOutcome, FileSystemPort, GitPort, McpCatalogPort, NetworkPort, + PermissionAuditRecord, PermissionDelegationContext, PermissionGrant, PermissionGrantKey, + PermissionReply, PermissionReplySource, PermissionRequest, PermissionRequestEvent, + PermissionRequestSource, PermissionRequestSourceKind, PortError, PortErrorKind, PortResult, + RemoteAssistantWorkspaceFacts, RemoteCapabilityPort, RemoteConnectionPort, + RemoteProjectionPort, RemoteRecentWorkspaceFacts, RemoteWorkspaceFacts, RemoteWorkspaceFileRuntimeHost, RemoteWorkspaceKind, RemoteWorkspacePort, RemoteWorkspaceRuntimeHost, RemoteWorkspaceUpdate, RuntimeEventEnvelope, RuntimeEventSink, RuntimeEventType, RuntimeServiceCapability, RuntimeServicePort, SessionStorageKind, @@ -491,6 +492,13 @@ impl AgentRuntime { self.inner.fork_session_at_turn(request).await } + pub async fn fork_session_before_turn( + &self, + request: AgentSessionForkBeforeTurnRequest, + ) -> Result { + self.inner.fork_session_before_turn(request).await + } + pub async fn generate_session_usage( &self, request: AgentSessionUsageRequest, diff --git a/src/crates/services/services-core/src/session/lineage.rs b/src/crates/services/services-core/src/session/lineage.rs index ede951a252..5418c05d89 100644 --- a/src/crates/services/services-core/src/session/lineage.rs +++ b/src/crates/services/services-core/src/session/lineage.rs @@ -32,6 +32,32 @@ struct SubagentRelationshipFacts { pub struct SessionBranchRequest { pub source_session_id: String, pub source_turn_id: String, + #[serde( + default, + skip_serializing_if = "SessionBranchBoundary::is_through_turn" + )] + pub boundary: SessionBranchBoundary, +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum SessionBranchBoundary { + #[default] + ThroughTurn, + BeforeTurn, +} + +impl SessionBranchBoundary { + pub const fn is_through_turn(&self) -> bool { + matches!(self, Self::ThroughTurn) + } + + pub const fn copied_turn_count(self, source_turn_index: usize) -> usize { + match self { + Self::ThroughTurn => source_turn_index + 1, + Self::BeforeTurn => source_turn_index, + } + } } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -58,6 +84,7 @@ pub struct BranchSessionMetadataFacts<'a> { pub source_session_id: &'a str, pub source_turn_id: &'a str, pub source_turn_index: usize, + pub boundary: SessionBranchBoundary, pub branched_turns: &'a [DialogTurnData], pub branch_lineage: &'a BranchSessionLineage, pub now_ms: u64, @@ -271,6 +298,7 @@ pub fn build_branched_session_metadata(facts: BranchSessionMetadataFacts<'_>) -> facts.source_session_id, facts.source_turn_id, facts.source_turn_index, + facts.boundary, facts.branch_lineage, ); metadata.relationship = None; @@ -314,6 +342,7 @@ fn build_branch_custom_metadata( source_session_id: &str, source_turn_id: &str, source_turn_index: usize, + boundary: SessionBranchBoundary, branch_lineage: &BranchSessionLineage, ) -> Option { let mut base = match strip_child_session_metadata(source_metadata) { @@ -321,15 +350,16 @@ fn build_branch_custom_metadata( _ => JsonMap::new(), }; - base.insert( - "forkOrigin".to_string(), - serde_json::json!({ - "sessionId": source_session_id, - "turnId": source_turn_id, - "turnIndex": source_turn_index + 1, - "baseTitle": branch_lineage.base_session_name, - }), - ); + let mut fork_origin = serde_json::json!({ + "sessionId": source_session_id, + "turnId": source_turn_id, + "turnIndex": boundary.copied_turn_count(source_turn_index), + "baseTitle": branch_lineage.base_session_name, + }); + if boundary == SessionBranchBoundary::BeforeTurn { + fork_origin["boundary"] = JsonValue::String("before_turn".to_string()); + } + base.insert("forkOrigin".to_string(), fork_origin); Some(JsonValue::Object(base)) } @@ -588,6 +618,7 @@ mod tests { source_session_id: "source", source_turn_id: "turn-2", source_turn_index: 1, + boundary: SessionBranchBoundary::ThroughTurn, branched_turns: &turns, branch_lineage: &branch_lineage, now_ms: 42, diff --git a/src/crates/services/services-core/src/session/mod.rs b/src/crates/services/services-core/src/session/mod.rs index bb9f4e3088..f3bf057ab2 100644 --- a/src/crates/services/services-core/src/session/mod.rs +++ b/src/crates/services/services-core/src/session/mod.rs @@ -13,7 +13,7 @@ pub use layout::SessionStorageLayout; pub use lineage::{ apply_session_lineage, build_branched_session_metadata, collect_hidden_subagent_cascade, format_branch_session_name, resolve_branch_session_lineage, BranchSessionLineage, - BranchSessionMetadataFacts, SessionBranchRequest, SessionBranchResult, + BranchSessionMetadataFacts, SessionBranchBoundary, SessionBranchRequest, SessionBranchResult, }; pub use memory_workspace::{ ensure_memory_workspace_git_baseline, memory_workspace_diff, render_memory_workspace_diff_file,