From 56aa04cc7c77b5d3d13f54bbca531cfe8a9d7fd6 Mon Sep 17 00:00:00 2001 From: limityan Date: Sat, 18 Jul 2026 16:03:27 +0800 Subject: [PATCH] fix(cli): harden terminal process reliability --- docs/architecture/cli-product-line-design.md | 18 +- docs/plans/core-decomposition-plan.md | 9 +- .../product-architecture-evolution-plan.md | 7 +- src/apps/cli/src/ui/mod.rs | 160 +++++-- src/apps/cli/tests/exec_cli_contracts.rs | 160 +++++++ src/apps/cli/tests/support/mod.rs | 431 ++++++++++++++++++ src/apps/cli/tests/tui_terminal_process.rs | 351 +++++++++----- 7 files changed, 975 insertions(+), 161 deletions(-) create mode 100644 src/apps/cli/tests/support/mod.rs diff --git a/docs/architecture/cli-product-line-design.md b/docs/architecture/cli-product-line-design.md index fac00bc34f..68e1c6ceae 100644 --- a/docs/architecture/cli-product-line-design.md +++ b/docs/architecture/cli-product-line-design.md @@ -110,7 +110,8 @@ BitFun CLI 应成为可独立安装和发布的 Agent 产品,而不是 Desktop - TUI 终端句柄由恢复守卫持有;初始化中途失败、正常返回、错误返回或 panic 展开都会尽力退出 alternate screen、 关闭输入捕获、关闭 raw mode 并显示光标。真实 PTY/ConPTY 启动页进程冒烟测试已验证 resize 后仍可交互、 多行输入、空闲 Ctrl+C 和可观察的终端清理序列;Chat 活动 turn 的 resize 静默期已有状态单测,窄屏流式 - reflow 已有 TestBackend 回归,真实 PTY 活动 turn、初始化失败与异常退出等仍需独立验收。 + reflow 已有 TestBackend 回归,Linux PTY 与 Windows ConPTY 活动 turn 的 resize/取消已有本地确定性流式模型夹具进程测试; + OS 级初始化失败与异常退出仍需独立验收。 - Startup 与 Chat 共用 CLI 私有输入读取器;一次读取同时受 256 个事件和 50ms 限制,跨批次仅延续快速文本尾部, 短批次普通按键保持原有路由。被识别为粘贴的文本按批次写入输入缓冲,每批只刷新一次命令菜单;粘贴内容中的 Tab 明确转换为四个空格。 @@ -147,7 +148,7 @@ BitFun CLI 应成为可独立安装和发布的 Agent 产品,而不是 Desktop | OpenCode 来源发现与真实执行尚未形成完整闭环 | “来源可识别”容易被误解为“插件可执行” | 第一条闭环只完成一个无外部依赖的契约样例;取得真实 `execute` 并注册到 Tool Runtime 后才显示可用。 | | 当前 CLI 使用 `product-full`,OHOS target 图包含多组未验证的平台依赖 | 不能据依赖可解析、`hdc shell` 或移动 Remote App 推导 PC 本地 CLI/TUI 可用 | 问题与风险统一记录在平台规约;具体工作另立专题,HAP 不作为替代。 | | Product Capability 已有,但品牌、资源、默认策略和发行配置没有统一产品定义 | 白标需要修改多处常量和工作流,能力隐藏不等于后端禁用 | 产品定义只在组装/构建边界选择身份、资源、能力包、默认策略和发行事实。 | -| CLI 已有独立 Linux 测试,参数互斥、结果/envelope 序列化、前置失败和组装有 focused contract;Linux 与 Windows 分别运行启动页 PTY/ConPTY 生命周期冒烟,发布归档上传前完成 SHA-256 与解压执行验证 | 真实模型审批/取消、真实 PTY 活动 turn 的 resize、Patch I/O 失败和终端故障注入仍可能晚于 PR 发现 | 继续补剩余进程级故障契约;避免为同一依赖图重复建立三平台编译矩阵。 | +| CLI 已有独立 Linux 测试,参数互斥、结果/envelope 序列化、前置失败和组装有 focused contract;启动页及本地确定性流式模型夹具驱动的活动 turn 由 PTY/ConPTY 进程测试覆盖 resize、取消和恢复可编辑状态,`stream-json` 覆盖 Patch 写入失败且不会泄漏成功终态;发布归档上传前完成 SHA-256 与解压执行验证 | 真实供应商审批流和 OS 级终端初始化故障注入仍可能晚于 PR 发现 | 只按剩余真实故障补进程契约;避免为同一依赖图重复建立三平台编译矩阵。 | ## 3. 分阶段产品需求 @@ -158,18 +159,19 @@ CLI-P0 的目标是建立后续功能补齐所需的稳定边界,不改变现 CLI-P0 不是一个统一重构 PR。静态 profile、真实 Runtime Services、Runtime Parts、调用级审批、共享事件源和 本地 Agent 纵向入口已接入;旧门面仅在后续 owner 迁移的行为等价成立后退出。配置解释、产品定制消费和 TUI 进一步拆分仍需独立交付。CLI 托管的 ACP 服务端已独立切换到 ACP profile 与组装后的 SDK runtime;启动页 -PTY/ConPTY 生命周期冒烟测试与发布归档冒烟测试已存在;Chat 活动 turn 的 resize 静默期已有确定性状态单测, -窄屏流式 reflow 已有 TestBackend 回归,真实模型、真实 PTY 活动 turn、终端故障注入与权限失败等完整进程级验收仍需另行完成。 +PTY/ConPTY 生命周期、Chat 活动 turn 的 resize/取消和发布归档冒烟测试已存在;resize 静默期与窄屏流式 reflow +分别有确定性状态单测和 TestBackend 回归,Patch 写入失败也由真实 `stream-json` 进程保护。真实供应商审批流、 +OS 级终端初始化故障注入与权限失败等完整进程级验收仍需另行完成。 其余工作独立立项,不能与 profile 迁移互相充当完成条件: | 切片 | 范围 | 退出条件 | |---|---|---| | 调用级审批 | TUI、`exec` 与 ACP 已使用各自调用级策略且不写全局配置 | Runtime-context `Allow always`、审批规划、`exec` 安全默认值和显式 `--auto` 有 focused test;真实模型/PTY 审批流与 ACP 仍需另行验收 | -| 输出协议 | 保持通用 `text/json/stream-json` 心智,复用现有 Agentic envelope,不新建 CLI schema | 已覆盖结果/envelope 序列化、参数与前置 JSON 失败、失败完成、同会话跨 turn 隔离和 stream-json/Patch stdout 冲突;真实信号、模型权限失败与 Patch I/O 故障注入仍需进程级契约 | +| 输出协议 | 保持通用 `text/json/stream-json` 心智,复用现有 Agentic envelope,不新建 CLI schema | 已覆盖结果/envelope 序列化、参数与前置 JSON 失败、失败完成、同会话跨 turn 隔离、stream-json/Patch stdout 冲突及 Patch 写入失败;写入失败返回非零并只发送结构化错误,不提前发送成功终态;真实信号与模型权限失败仍需进程级契约 | | 配置解释 | Canonical Config 层级、全局/项目持续来源、加载状态和兼容导入 dry-run | 不自动写入;冲突、未知字段、待确认能力和凭据引用可解释 | | 产品定制 | 消费最小产品定义、组装结果和已注册 TUI layout/theme ID | 第二个真实 CLI 产品复用后再提升公共字段 | -| TUI 边界 | 增量提取终端恢复守卫、命令分发和副作用边界 | 不改版视觉设计,恢复/取消回归可单独验证 | +| TUI 边界 | 增量提取终端恢复守卫、命令分发和副作用边界 | 不改版视觉设计;Linux PTY 与 Windows ConPTY 活动 turn 的 resize/取消、恢复可编辑状态和正常退出清理可单独验证,macOS 活动 turn 与 OS 级初始化失败注入另行补齐 | CLI-P0 不包含插件 JS/TS 执行、完整 checkpoint/rewind 或大规模 TUI 重写。 CLI-P0/P1/P2 在 Windows、macOS、Linux 完成不表示 HarmonyOS PC 已支持;HarmonyOS PC 的具体适配由未来独立专题 @@ -603,8 +605,8 @@ CLI Agent 能力加强必须落在共享 Agent Runtime、Tool Runtime 或 Harnes 通用 `cargo check --workspace` 负责三平台 CLI 编译保护;独立 CLI CI 运行 `cargo test --locked -p bitfun-cli -p bitfun-acp -p bitfun-agent-runtime`。Linux 启动页 PTY 生命周期冒烟随独立 CLI 测试运行,Windows 启动页 ConPTY 生命周期冒烟复用通用 Windows job;发布归档在上传前完成 SHA-256 与解压执行 -验证。真实模型进程级交互、真实 PTY 中的 Chat 活动 turn resize 与故障进程矩阵仍按对应切片补入门禁,不能由序列化单测或基础冒烟 -测试代替。 +验证。真实供应商模型进程级交互、macOS 活动 PTY 与 OS 级终端故障进程矩阵仍按对应切片补入门禁;Linux PTY +与 Windows ConPTY 的 Chat 活动 turn resize/取消已由本地确定性流式模型夹具覆盖,不能用它替代上述验收。 ### 10.2 阶段退出条件 diff --git a/docs/plans/core-decomposition-plan.md b/docs/plans/core-decomposition-plan.md index 1017d74310..67fc77d085 100644 --- a/docs/plans/core-decomposition-plan.md +++ b/docs/plans/core-decomposition-plan.md @@ -28,7 +28,7 @@ | Agent Runtime SDK | 已有无 `bitfun-core` 依赖的 v1 preview 门面和 smoke test | 发布边界仍需真实嵌入方证明 | | 插件运行时 | 现有路径只覆盖 BitFun 原生包和 OpenCode custom tool 静态名称预览 | 不能据通用 envelope 或静态候选扩张稳定 ABI | | Relay | room/device 状态、account/sync 存储、asset store 与 HTTP/WebSocket router 已归属 `services/relay-service`,standalone 与 embedded 入口同向消费;embedded 宿主逻辑仍在 assembly 兼容路径 | Cargo metadata 门禁覆盖 workspace、独立 manifest、normal/build/dev 依赖及 optional/target 变体;宿主归位是独立后续工作 | -| CLI CI | 独立 Linux job 运行 CLI test,通用三平台 workspace check 覆盖 CLI 编译;Linux PTY 与 Windows ConPTY 有启动页生命周期进程冒烟,发布归档上传前校验 SHA-256 并解压执行 | 参数/序列化/前置失败和组装已有 focused contract;resize 静默期已有状态单测,活动流式内容窄屏 reflow 已有 TestBackend 回归,真实模型/PTY 活动 turn、终端故障注入与 Patch I/O 失败仍需补齐 | +| CLI CI | 独立 Linux job 运行 CLI test,通用三平台 workspace check 覆盖 CLI 编译;Linux PTY 与 Windows ConPTY 有启动页生命周期及本地确定性流式模型夹具驱动的活动 turn 进程测试,发布归档上传前校验 SHA-256 并解压执行 | 参数/序列化/前置失败和组装已有 focused contract;resize 静默期、活动流式内容窄屏 reflow、Linux PTY/Windows ConPTY 活动 turn resize/取消及 Patch I/O 失败已有分层回归,真实供应商审批流、macOS 活动 PTY 与 OS 级终端故障注入仍需补齐 | ## 3. 目标依赖与归属 @@ -75,11 +75,12 @@ Peer Host 的 Runtime 接入和跨 Relay/Desktop/Web 的协议切换保持独立 1. 以真实调用方和行为等价测试逐项缩小快照及 Peer Host/ACP 持久化维护兼容面;远程分支另行定义身份和存储语义,模型目录与配置仍保留在产品入口。 2. 继续迁移 ACP 尚未接入 SDK 的持久化历史、模型目录/模式和 MCP 操作;ACP stdio 与协议投影生命周期保留在接口入口。 -3. 继续按真实故障样例拆分 TUI 副作用边界,不以大规模重写替代现有回归保护。 +3. 继续按真实故障样例拆分 TUI 副作用边界;当前切片已覆盖本地确定性流式模型夹具驱动的 Linux PTY/Windows ConPTY 活动 turn resize/取消、 + `stream-json` Patch 写入失败和终端恢复错误聚合,不以大规模重写替代现有回归保护。 当前 assembly 切换条件已经满足:CLI 生产入口消费真实组装结果,目标链路没有第二套状态,独立测试与三平台 -编译门禁存在,启动页 PTY/ConPTY 生命周期与发布归档冒烟测试已接入门禁。CLI-P0 整体退出条件尚未满足; -真实模型交互、真实 PTY 活动 turn 的 resize、终端故障注入、兼容门面退出以及 ACP/Desktop 切换仍需分别验收。 +编译门禁存在,启动页 PTY/ConPTY 生命周期、活动 turn resize/取消、Patch I/O 失败与发布归档冒烟测试已接入门禁。 +CLI-P0 整体退出条件尚未满足;真实供应商审批流、OS 级终端初始化故障注入、兼容门面退出以及 ACP/Desktop 切换仍需分别验收。 ### 4.3 依次切换 ACP 与 Desktop diff --git a/docs/plans/product-architecture-evolution-plan.md b/docs/plans/product-architecture-evolution-plan.md index 1fe9a5789a..3b8607adc1 100644 --- a/docs/plans/product-architecture-evolution-plan.md +++ b/docs/plans/product-architecture-evolution-plan.md @@ -29,7 +29,7 @@ |---|---|---| | 编译依赖 | `assembly/core -> apps/relay-server` 已移除;通用检查覆盖 normal/build/dev 依赖及 optional/target 变体 | 后续反向依赖和未知 crate 层级直接失败 | | 公开面 | `bitfun-core` 仍有迁移期 re-export;CLI 主会话客户端已仅消费 Runtime SDK,其他产品入口仍保留兼容路径 | 按入口逐项迁移,不做全仓逐 symbol 台账或批量删除 | -| CLI/TUI | `ShortcutsConfig` 已加载但真实按键分发仍硬编码;Slash、Palette、帮助和执行不是同一来源 | 先统一宿主 action 声明和键位解析,不重写 renderer | +| CLI/TUI | 宿主 `ACTION_SPECS` 已统一 Slash、Palette、Help、Keymap 与 dispatch;启动页及活动 turn 的 Linux PTY / Windows ConPTY 行为由进程级契约保护 | 保持现有 renderer 与交互规格,只按真实故障样例补可靠性契约;macOS 活动 PTY 另行验收 | | OpenCode | Prompt Command、受支持的单文件 JavaScript Tool 和 Subagent 安全子集已分别通过能力专属 provider 接入;受管 package plugin 仍只有静态预览 | 先收敛三条已交付路径的诊断、运行时提示和配置失败语义,再按真实阻塞样例评估下一能力切片 | | HarmonyOS PC | 未来平台目标,当前未实现 | 目标、问题、风险和旧设计闭环见平台规约;具体工作后续分别立项 | | 入口迁移 | CLI 已消费 Runtime Parts;Desktop 主交互消费由现有 owner 构造的窄口径 Runtime SDK 门面,完整 Desktop Runtime Parts 尚未组装;CLI/ACP/Desktop 仍按需保留 `bitfun-core/product-full` 兼容 owner | 保持单一 owner,按真实端口逐项迁移,不批量删除兼容门面或用桩服务提前声明能力 | @@ -58,10 +58,13 @@ - contracts 和 Agent Runtime 不再新增环境或生态来源探测; - 本工作流没有新增无调用方端口、空 registry 或第二个 Runtime owner。 -## 4. 工作流二:CLI action 与快捷键一致 +## 4. 工作流二:CLI action 与快捷键一致(已完成) 用户结果:持久化快捷键真正生效,Slash、命令面板、帮助、快捷键展示和执行不会互相漂移。 +当前状态:CLI 宿主已建立统一 action registry,Slash、Palette、Help、Keymap 与 dispatch 共用同一组稳定条目; +显式旧快捷键、冲突配置、宿主安全 fallback 和帮助尺寸均有回归保护。后续终端可靠性工作不再重复设计 action 层。 + 交付: - 在 CLI 宿主内建立一个 action registry。条目只包含稳定 action id、名称/别名、适用上下文、可用性、处理器、 diff --git a/src/apps/cli/src/ui/mod.rs b/src/apps/cli/src/ui/mod.rs index 4097730c9e..e5dca4aede 100644 --- a/src/apps/cli/src/ui/mod.rs +++ b/src/apps/cli/src/ui/mod.rs @@ -87,28 +87,16 @@ pub(crate) fn init_terminal() -> Result { EnableMouseCapture, EnableBracketedPaste ) { - let _ = disable_raw_mode(); - let _ = execute!( - stdout, - DisableBracketedPaste, - DisableMouseCapture, - LeaveAlternateScreen - ); - return Err(error.into()); + let cleanup = cleanup_partial_terminal(&mut stdout); + return Err(merge_terminal_failure(error, cleanup)); } let backend = CrosstermBackend::new(stdout); let terminal = match Terminal::new(backend) { Ok(terminal) => terminal, Err(error) => { let mut stdout = io::stdout(); - let _ = disable_raw_mode(); - let _ = execute!( - stdout, - DisableBracketedPaste, - DisableMouseCapture, - LeaveAlternateScreen - ); - return Err(error.into()); + let cleanup = cleanup_partial_terminal(&mut stdout); + return Err(merge_terminal_failure(error, cleanup)); } }; Ok(TerminalGuard { @@ -128,22 +116,40 @@ pub(crate) fn restore_terminal(mut guard: TerminalGuard) -> Result<()> { } fn restore_terminal_inner(terminal: &mut CliTerminal) -> Result<()> { - let mut errors = Vec::new(); - if let Err(error) = disable_raw_mode() { - errors.push(format!("disable raw mode: {error}")); - } - if let Err(error) = execute!( - terminal.backend_mut(), - DisableBracketedPaste, - DisableMouseCapture, - LeaveAlternateScreen - ) { - errors.push(format!("restore terminal screen: {error}")); - } - if let Err(error) = terminal.show_cursor() { - errors.push(format!("show terminal cursor: {error}")); - } + let disable_raw = disable_raw_mode(); + let disable_bracketed_paste = execute!(terminal.backend_mut(), DisableBracketedPaste); + let disable_mouse_capture = execute!(terminal.backend_mut(), DisableMouseCapture); + let leave_alternate_screen = execute!(terminal.backend_mut(), LeaveAlternateScreen); + let show_cursor = terminal.show_cursor(); + finish_terminal_cleanup([ + ("disable raw mode", disable_raw), + ("disable bracketed paste", disable_bracketed_paste), + ("disable mouse capture", disable_mouse_capture), + ("leave alternate screen", leave_alternate_screen), + ("show terminal cursor", show_cursor), + ]) +} + +fn cleanup_partial_terminal(stdout: &mut io::Stdout) -> Result<()> { + let disable_raw = disable_raw_mode(); + let disable_bracketed_paste = execute!(stdout, DisableBracketedPaste); + let disable_mouse_capture = execute!(stdout, DisableMouseCapture); + let leave_alternate_screen = execute!(stdout, LeaveAlternateScreen); + finish_terminal_cleanup([ + ("disable raw mode", disable_raw), + ("disable bracketed paste", disable_bracketed_paste), + ("disable mouse capture", disable_mouse_capture), + ("leave alternate screen", leave_alternate_screen), + ]) +} +fn finish_terminal_cleanup( + results: [(&'static str, std::io::Result<()>); N], +) -> Result<()> { + let errors = results + .into_iter() + .filter_map(|(operation, result)| result.err().map(|error| format!("{operation}: {error}"))) + .collect::>(); if errors.is_empty() { Ok(()) } else { @@ -151,6 +157,16 @@ fn restore_terminal_inner(terminal: &mut CliTerminal) -> Result<()> { } } +fn merge_terminal_failure(primary: std::io::Error, cleanup: Result<()>) -> anyhow::Error { + match cleanup { + Ok(()) => primary.into(), + Err(cleanup_error) => { + let context = format!("{primary}; failed to restore the terminal: {cleanup_error}"); + anyhow::Error::new(primary).context(context) + } + } +} + /// Render a loading/status message on the terminal (stays in alternate screen) pub(crate) fn render_loading( terminal: &mut Terminal>, @@ -180,3 +196,85 @@ pub(crate) fn render_loading( })?; Ok(()) } + +#[cfg(test)] +mod terminal_lifecycle_tests { + use super::{finish_terminal_cleanup, merge_terminal_failure}; + + #[test] + fn terminal_cleanup_reports_every_failed_step_in_order() { + let error = finish_terminal_cleanup([ + ( + "disable raw mode", + Err(std::io::Error::other("raw failure")), + ), + ( + "disable bracketed paste", + Err(std::io::Error::other("paste failure")), + ), + ("disable mouse capture", Ok(())), + ( + "leave alternate screen", + Err(std::io::Error::other("screen failure")), + ), + ( + "show terminal cursor", + Err(std::io::Error::other("cursor failure")), + ), + ]) + .expect_err("cleanup failures must be reported") + .to_string(); + + assert_eq!( + error, + "disable raw mode: raw failure; disable bracketed paste: paste failure; leave alternate screen: screen failure; show terminal cursor: cursor failure" + ); + } + + #[test] + fn initialization_failure_without_cleanup_error_keeps_primary_io_error() { + let error = merge_terminal_failure( + std::io::Error::new(std::io::ErrorKind::NotConnected, "terminal unavailable"), + Ok(()), + ); + + assert_eq!( + error + .downcast_ref::() + .map(std::io::Error::kind), + Some(std::io::ErrorKind::NotConnected) + ); + } + + #[test] + fn initialization_failure_keeps_primary_and_cleanup_diagnostics() { + let cleanup = finish_terminal_cleanup([( + "disable raw mode", + Err(std::io::Error::other("cleanup failure")), + )]); + let error = merge_terminal_failure( + std::io::Error::new( + std::io::ErrorKind::PermissionDenied, + "terminal initialization failure", + ), + cleanup, + ); + let message = error.to_string(); + + assert!( + message.contains("terminal initialization failure"), + "{message}" + ); + assert!( + message.contains("disable raw mode: cleanup failure"), + "{message}" + ); + assert_eq!( + error + .downcast_ref::() + .map(std::io::Error::kind), + Some(std::io::ErrorKind::PermissionDenied), + "primary io::Error must remain in the anyhow source chain" + ); + } +} diff --git a/src/apps/cli/tests/exec_cli_contracts.rs b/src/apps/cli/tests/exec_cli_contracts.rs index 6b2b3dbe62..845167eee1 100644 --- a/src/apps/cli/tests/exec_cli_contracts.rs +++ b/src/apps/cli/tests/exec_cli_contracts.rs @@ -1,4 +1,7 @@ +mod support; + use std::process::{Command, Output}; +use support::{CliTestEnvironment, MockOpenAiServer, STREAM_COMPLETED_MARKER}; fn run_cli(args: &[&str]) -> Output { Command::new(env!("CARGO_BIN_EXE_bitfun-cli")) @@ -15,6 +18,23 @@ fn stderr(output: &Output) -> String { String::from_utf8_lossy(&output.stderr).into_owned() } +fn jsonl_events(output: &str) -> Vec { + output + .lines() + .map(|line| { + serde_json::from_str::(line) + .unwrap_or_else(|error| panic!("invalid JSONL line {line:?}: {error}")) + }) + .collect() +} + +fn is_terminal_event(value: &serde_json::Value) -> bool { + matches!( + value["event"]["type"].as_str(), + Some("DialogTurnCompleted" | "DialogTurnCancelled" | "DialogTurnFailed" | "SystemError") + ) +} + #[test] fn exec_help_uses_competitor_aligned_output_and_approval_flags() { let output = run_cli(&["exec", "--help"]); @@ -154,3 +174,143 @@ fn stream_json_rejects_stdout_patch_before_starting_runtime() { stderr(&output) ); } + +#[test] +fn stream_json_patch_write_failure_emits_error_without_success_terminal() { + let server = MockOpenAiServer::immediate(); + let environment = CliTestEnvironment::new(); + environment.initialize_git_repository(); + environment.configure_mock_model(server.base_url()); + let output_target = environment.workspace().to_string_lossy().into_owned(); + let output = environment + .std_command() + .args([ + "exec", + "exercise patch settlement", + "--output-format", + "stream-json", + "--output-patch", + &output_target, + ]) + .output() + .expect("run stream-json patch failure contract"); + + let stdout = stdout(&output); + assert!(!output.status.success(), "{stdout}"); + assert_eq!(output.status.code(), Some(1), "{}", stderr(&output)); + assert!( + stderr(&output) + .lines() + .any(|line| line.starts_with("BITFUN_EXIT: patch_write_failed:")), + "missing stable patch failure diagnostic: {}", + stderr(&output) + ); + assert!(!stdout.trim().is_empty(), "missing stream-json events"); + let events = jsonl_events(&stdout); + let completed_marker_index = events + .iter() + .position(|value| { + value["event"]["type"] == "TextChunk" + && value["event"]["text"] + .as_str() + .is_some_and(|text| text.contains(STREAM_COMPLETED_MARKER)) + }) + .unwrap_or_else(|| { + panic!("model stream did not complete before patch settlement: {stdout}") + }); + let patch_error_index = events + .iter() + .position(|value| { + value["event"]["type"] == "SystemError" + && value["event"]["error"] + .as_str() + .is_some_and(|error| error.contains("Failed to save requested patch")) + }) + .unwrap_or_else(|| { + panic!("patch failure did not emit a structured system error: {stdout}") + }); + assert!( + completed_marker_index < patch_error_index, + "patch failure was emitted before model stream completion: {stdout}" + ); + assert!( + events + .iter() + .all(|value| value["event"]["type"] != "DialogTurnCompleted"), + "successful terminal event leaked before patch settlement: {stdout}" + ); + assert_eq!( + events + .iter() + .filter(|value| is_terminal_event(value)) + .count(), + 1, + "patch failure must emit exactly one terminal envelope: {stdout}" + ); + let terminal_event = events.last().expect("stream-json terminal event"); + assert_eq!(terminal_event["event"]["type"], "SystemError", "{stdout}"); + assert_eq!(terminal_event["event"]["recoverable"], false, "{stdout}"); + assert!( + terminal_event["event"]["error"] + .as_str() + .is_some_and(|error| error.contains("Failed to save requested patch")), + "unexpected terminal patch failure: {stdout}" + ); +} + +#[test] +fn stream_json_patch_success_emits_one_success_terminal() { + let server = MockOpenAiServer::immediate(); + let environment = CliTestEnvironment::new(); + environment.initialize_git_repository(); + environment.configure_mock_model(server.base_url()); + let output_patch = environment + .workspace() + .parent() + .expect("workspace parent") + .join("result.patch"); + let output_target = output_patch.to_string_lossy().into_owned(); + let output = environment + .std_command() + .args([ + "exec", + "exercise successful patch settlement", + "--output-format", + "stream-json", + "--output-patch", + &output_target, + ]) + .output() + .expect("run stream-json patch success contract"); + + let stdout = stdout(&output); + assert!(output.status.success(), "{}\n{stdout}", stderr(&output)); + assert_eq!( + std::fs::read_to_string(&output_patch).expect("read generated patch"), + "", + "clean workspace must produce an explicit empty patch" + ); + let events = jsonl_events(&stdout); + assert!( + events.iter().any(|value| { + value["event"]["type"] == "TextChunk" + && value["event"]["text"] + .as_str() + .is_some_and(|text| text.contains(STREAM_COMPLETED_MARKER)) + }), + "model stream did not complete: {stdout}" + ); + assert_eq!( + events + .iter() + .filter(|value| is_terminal_event(value)) + .count(), + 1, + "success must emit exactly one terminal envelope: {stdout}" + ); + let final_event = events.last().expect("stream-json success terminal event"); + assert_eq!( + final_event["event"]["type"], "DialogTurnCompleted", + "success terminal must be the final envelope: {stdout}" + ); +} diff --git a/src/apps/cli/tests/support/mod.rs b/src/apps/cli/tests/support/mod.rs new file mode 100644 index 0000000000..98326d2034 --- /dev/null +++ b/src/apps/cli/tests/support/mod.rs @@ -0,0 +1,431 @@ +#![allow(dead_code)] + +use portable_pty::CommandBuilder; +use serde_json::json; +use std::io::{Read, Write}; +use std::net::{TcpListener, TcpStream}; +use std::path::{Path, PathBuf}; +use std::process::Command; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::{mpsc, Arc}; +use std::thread::{self, JoinHandle}; +use std::time::{Duration, Instant}; + +pub(crate) const STREAM_START_MARKER: &str = "ACTIVE_TURN_STREAM_MARKER"; +pub(crate) const STREAM_RESIZED_MARKER: &str = "RESIZED_OK"; +pub(crate) const STREAM_COMPLETED_MARKER: &str = "ACTIVE_TURN_STREAM_COMPLETED"; + +pub(crate) struct CliTestEnvironment { + _temp: tempfile::TempDir, + workspace: PathBuf, + user_root: PathBuf, + home_root: PathBuf, + config_root: PathBuf, +} + +impl CliTestEnvironment { + pub(crate) fn new() -> Self { + let temp = tempfile::tempdir().expect("create isolated CLI environment"); + let workspace = temp.path().join("workspace"); + let user_root = temp.path().join("user-root"); + let home_root = temp.path().join("home"); + let config_root = temp.path().join("config-root"); + for path in [&workspace, &user_root, &home_root, &config_root] { + std::fs::create_dir_all(path).expect("create isolated CLI directory"); + } + + Self { + _temp: temp, + workspace, + user_root, + home_root, + config_root, + } + } + + pub(crate) fn workspace(&self) -> &Path { + &self.workspace + } + + pub(crate) fn configure_mock_model(&self, server_base_url: &str) { + let config_dir = self.user_root.join("config"); + std::fs::create_dir_all(&config_dir).expect("create model config directory"); + let base_url = format!("{}/v1", server_base_url.trim_end_matches('/')); + let request_url = format!("{base_url}/chat/completions"); + let config = json!({ + "app": { + "ai_experience": { + "enable_session_title_generation": false + } + }, + "ai": { + "models": [{ + "id": "cli-e2e-model", + "name": "CLI E2E Model", + "provider": "openai", + "model_name": "cli-e2e-model", + "base_url": base_url, + "request_url": request_url, + "api_key": "cli-e2e-key", + "enabled": true, + "category": "general_chat", + "capabilities": ["text_chat", "function_calling"] + }], + "default_models": { + "primary": "cli-e2e-model" + }, + "agent_model_defaults": { + "mode": "cli-e2e-model" + }, + "max_rounds": 1, + "stream_idle_timeout_secs": 10, + "stream_ttft_timeout_secs": 10 + } + }); + std::fs::write( + config_dir.join("app.json"), + serde_json::to_vec_pretty(&config).expect("serialize model config"), + ) + .expect("write model config"); + } + + pub(crate) fn initialize_git_repository(&self) { + self.run_git(&["init", "--quiet"]); + self.run_git(&["config", "user.email", "cli-tests@example.invalid"]); + self.run_git(&["config", "user.name", "CLI Tests"]); + std::fs::write(self.workspace.join("seed.txt"), "seed\n").expect("write git seed"); + self.run_git(&["add", "seed.txt"]); + self.run_git(&["commit", "--quiet", "-m", "seed"]); + } + + pub(crate) fn std_command(&self) -> Command { + let mut command = Command::new(env!("CARGO_BIN_EXE_bitfun-cli")); + command.current_dir(&self.workspace); + self.apply_std_environment(&mut command); + command + } + + pub(crate) fn pty_command(&self) -> CommandBuilder { + let mut command = CommandBuilder::new(env!("CARGO_BIN_EXE_bitfun-cli")); + command.cwd(&self.workspace); + command.env_remove("BITFUN_USER_ROOT"); + command.env_remove("BITFUN_HOME"); + command.env("BITFUN_E2E_STORAGE_GUARD", "1"); + command.env("BITFUN_E2E_USER_ROOT", &self.user_root); + command.env("BITFUN_E2E_HOME", &self.home_root); + command.env("APPDATA", &self.config_root); + command.env("XDG_CONFIG_HOME", &self.config_root); + command.env("HOME", &self.home_root); + command.env("USERPROFILE", &self.home_root); + command.env("TERM", "xterm-256color"); + command + } + + fn apply_std_environment(&self, command: &mut Command) { + command + .env_remove("BITFUN_USER_ROOT") + .env_remove("BITFUN_HOME") + .env("BITFUN_E2E_STORAGE_GUARD", "1") + .env("BITFUN_E2E_USER_ROOT", &self.user_root) + .env("BITFUN_E2E_HOME", &self.home_root) + .env("APPDATA", &self.config_root) + .env("XDG_CONFIG_HOME", &self.config_root) + .env("HOME", &self.home_root) + .env("USERPROFILE", &self.home_root) + .env("TERM", "xterm-256color"); + } + + fn run_git(&self, args: &[&str]) { + let output = Command::new("git") + .args(args) + .current_dir(&self.workspace) + .output() + .expect("run git for CLI test"); + assert!( + output.status.success(), + "git {args:?} failed: {}", + String::from_utf8_lossy(&output.stderr) + ); + } +} + +pub(crate) struct MockOpenAiServer { + base_url: String, + release_stream: mpsc::Sender<()>, + stream_disconnected: mpsc::Receiver<()>, + stop: Arc, + thread: Option>, +} + +impl MockOpenAiServer { + pub(crate) fn gated() -> Self { + Self::spawn(true) + } + + pub(crate) fn immediate() -> Self { + Self::spawn(false) + } + + pub(crate) fn base_url(&self) -> &str { + &self.base_url + } + + pub(crate) fn release(&self) { + let _ = self.release_stream.send(()); + } + + pub(crate) fn expect_stream_disconnect(&self, timeout: Duration) { + self.stream_disconnected + .recv_timeout(timeout) + .expect("model stream remained connected after cancellation"); + } + + fn spawn(gated: bool) -> Self { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock model server"); + listener + .set_nonblocking(true) + .expect("configure mock model listener"); + let address = listener.local_addr().expect("mock model address"); + let stop = Arc::new(AtomicBool::new(false)); + let stop_for_thread = Arc::clone(&stop); + let (release_tx, release_rx) = mpsc::channel(); + let (disconnect_tx, disconnect_rx) = mpsc::channel(); + let thread = thread::spawn(move || { + let deadline = Instant::now() + Duration::from_secs(30); + loop { + match listener.accept() { + Ok((mut stream, _)) => { + stream + .set_nonblocking(false) + .expect("configure accepted mock model connection"); + serve_model_response(&mut stream, gated, &release_rx, &disconnect_tx); + break; + } + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { + if stop_for_thread.load(Ordering::Relaxed) || Instant::now() >= deadline { + break; + } + thread::sleep(Duration::from_millis(10)); + } + Err(error) => panic!("accept mock model request: {error}"), + } + } + }); + + Self { + base_url: format!("http://{address}"), + release_stream: release_tx, + stream_disconnected: disconnect_rx, + stop, + thread: Some(thread), + } + } +} + +impl Drop for MockOpenAiServer { + fn drop(&mut self) { + self.stop.store(true, Ordering::Relaxed); + self.release(); + self.release(); + if let Some(thread) = self.thread.take() { + if let Err(panic) = thread.join() { + if !std::thread::panicking() { + std::panic::resume_unwind(panic); + } + } + } + } +} + +fn serve_model_response( + stream: &mut TcpStream, + gated: bool, + release_stream: &mpsc::Receiver<()>, + stream_disconnected: &mpsc::Sender<()>, +) { + stream + .set_read_timeout(Some(Duration::from_secs(5))) + .expect("configure mock request timeout"); + read_http_request(stream).expect("read mock model request"); + stream + .write_all( + b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\nTransfer-Encoding: chunked\r\nConnection: close\r\n\r\n", + ) + .expect("write mock response headers"); + + write_sse_chunk( + stream, + &json!({ + "id": "chatcmpl_cli_e2e", + "object": "chat.completion.chunk", + "created": 1, + "model": "cli-e2e-model", + "choices": [{ + "index": 0, + "delta": {"role": "assistant", "content": ""}, + "finish_reason": null + }] + }) + .to_string(), + ) + .expect("write mock role chunk"); + write_sse_chunk( + stream, + &json!({ + "id": "chatcmpl_cli_e2e", + "object": "chat.completion.chunk", + "created": 2, + "model": "cli-e2e-model", + "choices": [{ + "index": 0, + "delta": {"content": STREAM_START_MARKER}, + "finish_reason": null + }] + }) + .to_string(), + ) + .expect("write mock streaming marker"); + + if gated + && release_stream + .recv_timeout(Duration::from_secs(30)) + .is_err() + { + return; + } + + if gated { + if write_sse_chunk( + stream, + &json!({ + "id": "chatcmpl_cli_e2e", + "object": "chat.completion.chunk", + "created": 3, + "model": "cli-e2e-model", + "choices": [{ + "index": 0, + "delta": {"content": STREAM_RESIZED_MARKER}, + "finish_reason": null + }] + }) + .to_string(), + ) + .is_err() + { + return; + } + if !wait_for_release_or_disconnect(stream, release_stream, stream_disconnected) { + return; + } + } + + if write_sse_chunk( + stream, + &json!({ + "id": "chatcmpl_cli_e2e", + "object": "chat.completion.chunk", + "created": 4, + "model": "cli-e2e-model", + "choices": [{ + "index": 0, + "delta": {"content": STREAM_COMPLETED_MARKER}, + "finish_reason": "stop" + }], + "usage": {"prompt_tokens": 3, "completion_tokens": 5, "total_tokens": 8} + }) + .to_string(), + ) + .is_err() + { + return; + } + if write_chunk(stream, b"data: [DONE]\n\n").is_err() { + return; + } + let _ = stream.write_all(b"0\r\n\r\n"); + let _ = stream.flush(); +} + +fn wait_for_release_or_disconnect( + stream: &TcpStream, + release_stream: &mpsc::Receiver<()>, + stream_disconnected: &mpsc::Sender<()>, +) -> bool { + stream + .set_read_timeout(Some(Duration::from_millis(25))) + .expect("configure mock disconnect observation"); + let deadline = Instant::now() + Duration::from_secs(30); + let mut probe = [0_u8; 1]; + while Instant::now() < deadline { + if release_stream.try_recv().is_ok() { + return true; + } + match stream.peek(&mut probe) { + Ok(0) => { + let _ = stream_disconnected.send(()); + return false; + } + Ok(_) => {} + Err(error) + if matches!( + error.kind(), + std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut + ) => {} + Err(_) => { + let _ = stream_disconnected.send(()); + return false; + } + } + } + false +} + +fn write_sse_chunk(stream: &mut TcpStream, data: &str) -> std::io::Result<()> { + write_chunk(stream, format!("data: {data}\n\n").as_bytes()) +} + +fn write_chunk(stream: &mut TcpStream, bytes: &[u8]) -> std::io::Result<()> { + write!(stream, "{:X}\r\n", bytes.len())?; + stream.write_all(bytes)?; + stream.write_all(b"\r\n")?; + stream.flush() +} + +fn read_http_request(stream: &mut TcpStream) -> std::io::Result> { + let mut request = Vec::new(); + let mut buffer = [0_u8; 4096]; + let mut expected_len = None; + loop { + let read = stream.read(&mut buffer)?; + if read == 0 { + break; + } + request.extend_from_slice(&buffer[..read]); + if expected_len.is_none() { + if let Some(header_end) = find_header_end(&request) { + let content_length = parse_content_length(&request[..header_end]); + expected_len = Some(header_end + 4 + content_length); + } + } + if expected_len.is_some_and(|expected| request.len() >= expected) { + break; + } + } + Ok(request) +} + +fn find_header_end(request: &[u8]) -> Option { + request.windows(4).position(|window| window == b"\r\n\r\n") +} + +fn parse_content_length(headers: &[u8]) -> usize { + String::from_utf8_lossy(headers) + .lines() + .find_map(|line| { + let (name, value) = line.split_once(':')?; + name.eq_ignore_ascii_case("content-length") + .then(|| value.trim().parse::().ok()) + .flatten() + }) + .unwrap_or_default() +} diff --git a/src/apps/cli/tests/tui_terminal_process.rs b/src/apps/cli/tests/tui_terminal_process.rs index b14eabe16e..a6227b9ece 100644 --- a/src/apps/cli/tests/tui_terminal_process.rs +++ b/src/apps/cli/tests/tui_terminal_process.rs @@ -1,8 +1,11 @@ -use portable_pty::{native_pty_system, CommandBuilder, PtyPair, PtySize}; +mod support; + +use portable_pty::{native_pty_system, CommandBuilder, PtySize}; use std::io::{Read, Write}; use std::sync::{Arc, Mutex}; use std::thread; use std::time::{Duration, Instant}; +use support::{CliTestEnvironment, MockOpenAiServer, STREAM_RESIZED_MARKER, STREAM_START_MARKER}; const INITIAL_SIZE: PtySize = PtySize { rows: 30, @@ -16,127 +19,56 @@ const RESIZED_SIZE: PtySize = PtySize { pixel_width: 0, pixel_height: 0, }; +const STARTUP_INPUT: &[u8] = b"exercise active turn resize Q7Z9"; +const STARTUP_INPUT_SENTINEL: &str = "Q7Z9"; +const RECOVERY_INPUT: &[u8] = b"READY_AFTER_CANCEL K4W8"; +const RECOVERY_INPUT_SENTINEL: &str = "K4W8"; #[test] fn interactive_startup_survives_resize_multiline_input_and_emits_cleanup() { - let storage = tempfile::tempdir().expect("create isolated CLI storage"); - let user_root = storage.path().join("user-root"); - let home_root = storage.path().join("home"); - std::fs::create_dir_all(&user_root).expect("create isolated user root"); - std::fs::create_dir_all(&home_root).expect("create isolated home root"); - - let pair = native_pty_system() - .openpty(INITIAL_SIZE) - .expect("open native PTY"); - let PtyPair { master, slave } = pair; - - let mut command = CommandBuilder::new(env!("CARGO_BIN_EXE_bitfun-cli")); - command.cwd(storage.path()); - command.env("BITFUN_E2E_STORAGE_GUARD", "1"); - command.env("BITFUN_E2E_USER_ROOT", &user_root); - command.env("BITFUN_E2E_HOME", &home_root); - command.env("HOME", &home_root); - command.env("USERPROFILE", &home_root); - command.env("TERM", "xterm-256color"); - - let mut child = slave - .spawn_command(command) - .expect("spawn bitfun-cli in native PTY"); - drop(slave); - - let mut reader = master.try_clone_reader().expect("clone PTY reader"); - let captured = Arc::new(Mutex::new(Vec::new())); - let reader_capture = Arc::clone(&captured); - let reader_thread = thread::spawn(move || { - let mut chunk = [0_u8; 4096]; - while let Ok(read) = reader.read(&mut chunk) { - if read == 0 { - break; - } - reader_capture - .lock() - .expect("lock captured PTY output") - .extend_from_slice(&chunk[..read]); - } - }); + let environment = CliTestEnvironment::new(); + let mut process = PtyProcess::spawn(environment.pty_command(), INITIAL_SIZE); - wait_for_output(&captured, "\x1b[?2004h", Duration::from_secs(30)).unwrap_or_else(|| { - terminate(&mut child); - panic!( - "interactive startup did not enable bracketed paste; output:\n{}", - captured_output(&captured) - ); - }); + process.expect_output( + "\x1b[?2004h", + Duration::from_secs(30), + "interactive startup did not enable bracketed paste", + ); #[cfg(unix)] assert!( - captured_output(&captured).contains("\x1b[?1049h"), + process.output().contains("\x1b[?1049h"), "interactive TUI must enter the alternate screen" ); - master.resize(RESIZED_SIZE).expect("resize native PTY"); - assert_eq!( - master.get_size().expect("read resized PTY dimensions"), - RESIZED_SIZE - ); + process.resize(RESIZED_SIZE); - let mut writer = master.take_writer().expect("take PTY writer"); #[cfg(unix)] - writer - .write_all(b"\x1b[200~alpha\r\nbeta\x1b[201~") - .expect("send bracketed paste"); + process.write(b"\x1b[200~alpha\r\nbeta\x1b[201~"); #[cfg(windows)] { let mut rapid_input = b"alpha".to_vec(); rapid_input.extend(std::iter::repeat_n(b'a', 251)); rapid_input.extend_from_slice(b"\rbeta"); - writer - .write_all(&rapid_input) - .expect("send rapid multiline key input across the batch boundary"); + process.write(&rapid_input); } - writer.flush().expect("flush terminal input"); - wait_for_output(&captured, "alpha", Duration::from_secs(15)).unwrap_or_else(|| { - terminate(&mut child); - panic!( - "interactive startup did not render multiline input; output:\n{}", - captured_output(&captured) - ); - }); - wait_for_output(&captured, "beta", Duration::from_secs(15)).unwrap_or_else(|| { - terminate(&mut child); - panic!( - "interactive startup did not render the multiline input tail; output:\n{}", - captured_output(&captured) - ); - }); + process.expect_output( + "alpha", + Duration::from_secs(15), + "interactive startup did not render multiline input", + ); + process.expect_output( + "beta", + Duration::from_secs(15), + "interactive startup did not render the multiline input tail", + ); assert!( - !captured_output(&captured).contains("Welcome to BitFun CLI!"), + !process.output().contains("Welcome to BitFun CLI!"), "multiline input was submitted instead of remaining in the startup editor" ); - writer.write_all(&[0x03]).expect("send Ctrl+C"); - writer.flush().expect("flush Ctrl+C"); - - let deadline = Instant::now() + Duration::from_secs(15); - let status = loop { - if let Some(status) = child.try_wait().expect("poll bitfun-cli process") { - break status; - } - if Instant::now() >= deadline { - terminate(&mut child); - panic!( - "interactive startup did not exit after Ctrl+C; output:\n{}", - captured_output(&captured) - ); - } - thread::sleep(Duration::from_millis(25)); - }; - - drop(writer); - drop(master); - reader_thread.join().expect("join PTY reader"); - - let output = captured_output(&captured); + process.write(&[0x03]); + let (status, output) = process.finish(Duration::from_secs(15)); assert!( status.success(), "unexpected process status {status}:\n{output}" @@ -172,26 +104,213 @@ fn interactive_startup_survives_resize_multiline_input_and_emits_cleanup() { ); } -fn wait_for_output( - captured: &Arc>>, - expected: &str, - timeout: Duration, -) -> Option<()> { - let deadline = Instant::now() + timeout; - while Instant::now() < deadline { - if captured_output(captured).contains(expected) { - return Some(()); - } - thread::sleep(Duration::from_millis(25)); - } - None +#[test] +fn active_turn_resize_can_be_cancelled_and_returns_to_editable_input() { + let server = MockOpenAiServer::gated(); + let environment = CliTestEnvironment::new(); + environment.initialize_git_repository(); + environment.configure_mock_model(server.base_url()); + let mut process = PtyProcess::spawn(environment.pty_command(), INITIAL_SIZE); + + process.expect_output( + "\x1b[?2004h", + Duration::from_secs(30), + "interactive startup did not enable bracketed paste", + ); + process.write(STARTUP_INPUT); + process.expect_output( + STARTUP_INPUT_SENTINEL, + Duration::from_secs(15), + "startup prompt sentinel was not rendered before submission", + ); + process.write(b"\r"); + process.expect_output( + STREAM_START_MARKER, + Duration::from_secs(30), + "active model stream was not rendered", + ); + + let active_turn_size = PtySize { + rows: 24, + cols: 52, + pixel_width: 0, + pixel_height: 0, + }; + process.resize(active_turn_size); + server.release(); + process.expect_output( + STREAM_RESIZED_MARKER, + Duration::from_secs(15), + "active model stream did not remain renderable after resize", + ); + process.write(&[0x03]); + process.expect_output( + "Cancelled", + Duration::from_secs(15), + "active turn did not reach the cancelled state after resize", + ); + server.expect_stream_disconnect(Duration::from_secs(5)); + + process.write(RECOVERY_INPUT); + process.expect_output( + RECOVERY_INPUT_SENTINEL, + Duration::from_secs(15), + "recovery input sentinel was not rendered after cancellation", + ); + process.write(&[0x03]); + + let (status, output) = process.finish(Duration::from_secs(15)); + assert!( + status.success(), + "unexpected process status {status}:\n{output}" + ); + assert!(output.contains(STREAM_START_MARKER), "{output}"); + assert!(output.contains(STREAM_RESIZED_MARKER), "{output}"); + assert!(output.contains("Cancelled"), "{output}"); + assert!(output.contains(RECOVERY_INPUT_SENTINEL), "{output}"); + assert!(output.contains("\x1b[?2004l"), "{output}"); + assert!(output.contains("\x1b[?25h"), "{output}"); + assert!(output.contains("Goodbye!"), "{output}"); } fn captured_output(captured: &Arc>>) -> String { String::from_utf8_lossy(&captured.lock().expect("lock captured PTY output")).into_owned() } -fn terminate(child: &mut Box) { - let _ = child.kill(); - let _ = child.wait(); +struct PtyProcess { + master: Option>, + writer: Option>, + child: Option>, + captured: Arc>>, + reader_thread: Option>, +} + +impl PtyProcess { + fn spawn(command: CommandBuilder, size: PtySize) -> Self { + let pair = native_pty_system().openpty(size).expect("open native PTY"); + let mut child = pair + .slave + .spawn_command(command) + .expect("spawn bitfun-cli in native PTY"); + drop(pair.slave); + + let mut reader = pair.master.try_clone_reader().expect("clone PTY reader"); + let writer = pair.master.take_writer().expect("take PTY writer"); + let captured = Arc::new(Mutex::new(Vec::new())); + let reader_capture = Arc::clone(&captured); + let reader_thread = thread::spawn(move || { + let mut chunk = [0_u8; 4096]; + while let Ok(read) = reader.read(&mut chunk) { + if read == 0 { + break; + } + reader_capture + .lock() + .expect("lock captured PTY output") + .extend_from_slice(&chunk[..read]); + } + }); + + // Detect an immediate startup failure before handing the process to the test. + if let Some(status) = child.try_wait().expect("poll initial CLI process") { + panic!("bitfun-cli exited during PTY startup: {status}"); + } + + Self { + master: Some(pair.master), + writer: Some(writer), + child: Some(child), + captured, + reader_thread: Some(reader_thread), + } + } + + fn output(&self) -> String { + captured_output(&self.captured) + } + + fn expect_output(&mut self, expected: &str, timeout: Duration, context: &str) { + let deadline = Instant::now() + timeout; + while Instant::now() < deadline { + if self.output().contains(expected) { + return; + } + if let Some(status) = self + .child + .as_mut() + .expect("PTY process child") + .try_wait() + .expect("poll bitfun-cli process") + { + let output = self.output(); + self.close_io(); + panic!("{context}; process exited with {status}; output:\n{output}"); + } + thread::sleep(Duration::from_millis(25)); + } + let output = self.output(); + self.terminate(); + panic!("{context}; output:\n{output}"); + } + + fn resize(&self, size: PtySize) { + let master = self.master.as_ref().expect("PTY master"); + master.resize(size).expect("resize native PTY"); + assert_eq!( + master.get_size().expect("read resized PTY dimensions"), + size + ); + } + + fn write(&mut self, bytes: &[u8]) { + let writer = self.writer.as_mut().expect("PTY writer"); + writer.write_all(bytes).expect("write terminal input"); + writer.flush().expect("flush terminal input"); + } + + fn finish(mut self, timeout: Duration) -> (portable_pty::ExitStatus, String) { + let deadline = Instant::now() + timeout; + let status = loop { + if let Some(status) = self + .child + .as_mut() + .expect("PTY process child") + .try_wait() + .expect("poll bitfun-cli process") + { + break status; + } + if Instant::now() >= deadline { + let output = self.output(); + self.terminate(); + panic!("interactive process did not exit; output:\n{output}"); + } + thread::sleep(Duration::from_millis(25)); + }; + self.child.take(); + self.close_io(); + (status, self.output()) + } + + fn close_io(&mut self) { + self.writer.take(); + self.master.take(); + if let Some(reader_thread) = self.reader_thread.take() { + reader_thread.join().expect("join PTY reader"); + } + } + + fn terminate(&mut self) { + if let Some(mut child) = self.child.take() { + let _ = child.kill(); + let _ = child.wait(); + } + self.close_io(); + } +} + +impl Drop for PtyProcess { + fn drop(&mut self) { + self.terminate(); + } }