diff --git a/.github/workflows/cmux-tui.yml b/.github/workflows/cmux-tui.yml index 94e94588800f..0bcbeec3f601 100644 --- a/.github/workflows/cmux-tui.yml +++ b/.github/workflows/cmux-tui.yml @@ -721,6 +721,78 @@ jobs: windows_runner: windows-latest checkout_ref: ${{ inputs.commit }} + bench-interact: + name: bench interact (${{ matrix.os }}) + needs: validate-inputs + if: inputs.mode == 'full' + runs-on: ${{ matrix.runner }} + timeout-minutes: 30 + # Record-only: this job publishes the IX0 interaction-latency baseline as an + # artifact and is intentionally not in hosted-verification's needs, so it is + # never a required check. See plans/cmux-tui-zero-wait-interaction.md (IX0). + strategy: + fail-fast: false + matrix: + include: + - os: macos + runner: blacksmith-6vcpu-macos-15 + - os: linux + runner: blacksmith-4vcpu-ubuntu-2404 + steps: + - uses: actions/checkout@de0fac2e4500dabe0009e67214ff5f5447ce83dd # v6.0.2 + with: + persist-credentials: false + ref: ${{ inputs.commit }} + + - name: Require exact checkout + env: + EXACT_COMMIT: ${{ inputs.commit }} + run: test "$(git rev-parse HEAD)" = "$EXACT_COMMIT" + + - name: Use short temporary directory for macOS socket tests + if: runner.os == 'macOS' + run: echo "TMPDIR=/tmp" >> "$GITHUB_ENV" + + - name: Init ghostty submodule + run: git submodule update --init --depth 1 ghostty + + - name: Install Linux build dependencies + if: runner.os == 'Linux' + run: | + sudo apt-get update + sudo apt-get install -y clang libclang-dev pkg-config + + - name: Install zig + run: ./scripts/install-zig-ci.sh + + - name: Set up pinned Rust + uses: ./.github/actions/setup-cmux-tui-rust + + - name: Build cmux-tui server + working-directory: cmux-tui + run: cargo build -p cmux-tui --locked + + - name: Run interaction benchmark + working-directory: cmux-tui + shell: bash + run: | + set -euo pipefail + # Human-readable table in the job log. + ./target/debug/cmux-tui bench interact --creates 20 --clients 1 --typing-probes 50 + # Machine-readable baseline artifact. + ./target/debug/cmux-tui --json bench interact \ + --creates 20 --clients 1 --typing-probes 50 \ + > "$RUNNER_TEMP/bench-interact.json" + cat "$RUNNER_TEMP/bench-interact.json" + + - name: Upload interaction baseline + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a # v7.0.1 + with: + name: cmux-tui-bench-interact-${{ runner.os }} + path: ${{ runner.temp }}/bench-interact.json + if-no-files-found: error + retention-days: 7 + hosted-verification: name: ${{ inputs.mode == 'full' && 'hosted verification' || 'focused hosted verification' }} if: always() diff --git a/cmux-tui/crates/cmux-tui-core/src/budgets.rs b/cmux-tui/crates/cmux-tui-core/src/budgets.rs new file mode 100644 index 000000000000..bf4097dfdc0e --- /dev/null +++ b/cmux-tui/crates/cmux-tui-core/src/budgets.rs @@ -0,0 +1,371 @@ +//! Named timing and size budgets shared by the daemon, terminal hosts, and +//! clients. +//! +//! Every bounded wait in cmux-tui is a budget with one name, one value, and +//! one stage of the interaction lifecycle it belongs to. The constants here are +//! the single source for the values; the code sites that enforce them import +//! these constants instead of spelling the number again, and `cmux-tui diag +//! budgets` prints [`table`] so an operator or an agent can read every bound in +//! one place. A timeout error should name the budget it exhausted. +//! +//! Stages: +//! - `accept`: the request is validated and applied to in-memory state. +//! - `durable`: the journal batch that carries the request has committed. +//! - `settle`: an external effect (host launch, terminate, first frame) has +//! reached its outcome. +//! - `frame`: a frontend paint cadence. +//! - `client`: a client-side wait on the daemon. +//! - `planned`: reserved by design and not enforced yet. + +use std::time::Duration; + +/// One bounded quantity. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum BudgetValue { + Duration(Duration), + Bytes(u64), +} + +impl BudgetValue { + /// Machine-readable unit name for JSON output. + pub fn unit(self) -> &'static str { + match self { + BudgetValue::Duration(_) => "ms", + BudgetValue::Bytes(_) => "bytes", + } + } + + /// Numeric value in [`Self::unit`]. + pub fn amount(self) -> u64 { + match self { + BudgetValue::Duration(duration) => { + u64::try_from(duration.as_millis()).unwrap_or(u64::MAX) + } + BudgetValue::Bytes(bytes) => bytes, + } + } +} + +/// One named budget and where it is enforced. +#[derive(Clone, Copy, Debug)] +pub struct Budget { + /// Dotted name, unique in [`table`]. + pub name: &'static str, + pub value: BudgetValue, + /// One of `accept`, `durable`, `settle`, `frame`, `client`, `planned`. + pub stage: &'static str, + /// What exhausting the budget means. + pub purpose: &'static str, + /// Code site that enforces it, as `crate::path` or a file name. + pub site: &'static str, +} + +/// Every stage name a budget may carry. +pub const STAGES: &[&str] = &["accept", "durable", "settle", "frame", "client", "planned"]; + +// Terminal host launch, handshake, and teardown. +pub const HOST_CONNECT_WINDOW: Duration = Duration::from_secs(1); +pub const HOST_CONNECT_INTERVAL: Duration = Duration::from_millis(10); +pub const HOST_HANDSHAKE: Duration = Duration::from_secs(2); +pub const HOST_SNAPSHOT_BOUNDARY: Duration = Duration::from_millis(1500); +pub const HOST_TERMINATE_GRACE: Duration = Duration::from_millis(250); +pub const HOST_PTY_DRAIN: Duration = Duration::from_millis(250); +pub const HOST_KILL_WAIT: Duration = Duration::from_secs(2); +pub const HOST_FORCED_DRAIN: Duration = Duration::from_millis(100); +pub const HOST_LAUNCH_ROLLBACK: Duration = Duration::from_secs(4); +pub const HOST_LAUNCH_OWNER: Duration = Duration::from_secs(5); +pub const HOST_CLIENT_WRITE: Duration = Duration::from_secs(2); +pub const HOST_CONTROL_RESPONSE: Duration = Duration::from_secs(2); +pub const TERMINAL_CLOSE_WAIT: Duration = Duration::from_secs(4); + +// Journal writer. +pub const JOURNAL_DURABLE_WAIT: Duration = Duration::from_secs(2); +pub const JOURNAL_COMMIT_RESULT_WAIT: Duration = Duration::from_secs(1); + +// Control server. +pub const SERVER_STREAM_WRITE: Duration = Duration::from_secs(2); +pub const SERVER_CONNECTION_SURFACE_SHUTDOWN: Duration = Duration::from_secs(3); + +// Clients. +pub const CLIENT_REQUEST: Duration = Duration::from_secs(10); +pub const CLIENT_WRITE: Duration = Duration::from_secs(2); +pub const OWNER_ENSURE: Duration = Duration::from_secs(10); +pub const OWNER_POLL: Duration = Duration::from_millis(25); +pub const FRAME: Duration = Duration::from_millis(16); + +// Planned, not enforced. +pub const INPUT_TYPEAHEAD_BYTES: u64 = 64 * 1024; + +const fn duration( + name: &'static str, + value: Duration, + stage: &'static str, + purpose: &'static str, + site: &'static str, +) -> Budget { + Budget { name, value: BudgetValue::Duration(value), stage, purpose, site } +} + +static TABLE: [Budget; 23] = [ + duration( + "host.connect_window", + HOST_CONNECT_WINDOW, + "settle", + "total time the daemon retries connecting to a freshly launched terminal host socket", + "cmux_tui_core::terminal_host_runtime::HOST_CONNECT_RETRY_WINDOW", + ), + duration( + "host.connect_interval", + HOST_CONNECT_INTERVAL, + "settle", + "pause between host socket connect attempts inside host.connect_window", + "cmux_tui_core::terminal_host_runtime::HOST_CONNECT_RETRY_INTERVAL", + ), + duration( + "host.handshake", + HOST_HANDSHAKE, + "settle", + "read and write timeout for one CMTH handshake exchange", + "cmux_tui_core::terminal_host_runtime::HOST_HANDSHAKE_TIMEOUT", + ), + duration( + "host.snapshot_boundary", + HOST_SNAPSHOT_BOUNDARY, + "settle", + "how long a host waits for its VT parser to reach ground before snapshotting for a new client", + "cmux_tui_core::terminal_host_runtime::HOST_SNAPSHOT_BOUNDARY_TIMEOUT", + ), + duration( + "host.terminate_grace", + HOST_TERMINATE_GRACE, + "settle", + "time after SIGHUP before the host escalates to SIGKILL", + "cmux_tui_core::terminal_host_runtime::HOST_TERMINATE_GRACE", + ), + duration( + "host.pty_drain", + HOST_PTY_DRAIN, + "settle", + "time the host waits for the PTY to drain after the child exits", + "cmux_tui_core::terminal_host_runtime::HOST_PTY_DRAIN_GRACE", + ), + duration( + "host.kill_wait", + HOST_KILL_WAIT, + "settle", + "time the host waits for the child to die after SIGKILL", + "cmux_tui_core::terminal_host_runtime::HOST_KILL_WAIT", + ), + duration( + "host.forced_drain", + HOST_FORCED_DRAIN, + "settle", + "final forced PTY drain window when the child ignored SIGKILL escalation", + "cmux_tui_core::terminal_host_runtime::HOST_FORCED_DRAIN_WINDOW", + ), + duration( + "host.launch_rollback", + HOST_LAUNCH_ROLLBACK, + "settle", + "time the daemon waits for a half-launched host to exit while rolling a failed create back", + "cmux_tui_core::terminal_host_runtime::HOST_LAUNCH_ROLLBACK_WAIT", + ), + duration( + "host.launch_owner", + HOST_LAUNCH_OWNER, + "settle", + "time a host waits for its launch owner to send Activate before releasing the PTY reader itself", + "cmux_tui_core::terminal_host_runtime::HOST_LAUNCH_OWNER_TIMEOUT", + ), + duration( + "host.client_write", + HOST_CLIENT_WRITE, + "settle", + "write timeout for daemon-to-host control frames", + "cmux_tui_core::terminal_host_runtime::HOST_CLIENT_WRITE_TIMEOUT", + ), + duration( + "host.control_response", + HOST_CONTROL_RESPONSE, + "settle", + "time the daemon waits for a host control response such as TerminateAck", + "cmux_tui_core::terminal_host_runtime::CONTROL_RESPONSE_TIMEOUT", + ), + duration( + "terminal.close_wait", + TERMINAL_CLOSE_WAIT, + "settle", + "total time a terminal close waits for the host to report Exit before killing it", + "cmux_tui_core::mux::TERMINAL_HOST_CLOSE_WAIT", + ), + duration( + "journal.durable_wait", + JOURNAL_DURABLE_WAIT, + "durable", + "time a producer waits for journal lane admission and the durable receipt", + "cmux_tui_core::journal_ingress::JOURNAL_DURABLE_WAIT", + ), + duration( + "journal.commit_result_wait", + JOURNAL_COMMIT_RESULT_WAIT, + "durable", + "time a producer waits for the writer to report its batch commit result", + "cmux_tui_core::journal_ingress::JOURNAL_COMMIT_RESULT_WAIT", + ), + duration( + "server.stream_write", + SERVER_STREAM_WRITE, + "client", + "write timeout for one response or event line to a control connection", + "cmux_tui_core::server::STREAM_WRITE_TIMEOUT", + ), + duration( + "server.connection_surface_shutdown", + SERVER_CONNECTION_SURFACE_SHUTDOWN, + "client", + "time the server waits for a connection's surface dispatcher to drain on close", + "cmux_tui_core::server::CONNECTION_SURFACE_SHUTDOWN_TIMEOUT", + ), + duration( + "client.request", + CLIENT_REQUEST, + "client", + "time the TUI client waits for one control response", + "cmux_tui::session::remote::REMOTE_REQUEST_TIMEOUT", + ), + duration( + "client.write", + CLIENT_WRITE, + "client", + "time the TUI client waits for an ordered socket write to be accepted", + "cmux_tui::session::remote::remote_write_timeout", + ), + duration( + "owner.ensure", + OWNER_ENSURE, + "client", + "total time a client spends probing, spawning, and waiting for a session owner", + "cmux_tui::local_owner::ENSURE_DEADLINE", + ), + duration( + "owner.poll", + OWNER_POLL, + "client", + "pause between owner readiness probes inside owner.ensure", + "cmux_tui::local_owner::POLL_INTERVAL", + ), + duration( + "frame", + FRAME, + "frame", + "TUI paint cadence; any wait longer than this must be shown as state", + "cmux_tui::app::TERMINAL_PAINT_CADENCE", + ), + Budget { + name: "input.typeahead_bytes", + value: BudgetValue::Bytes(INPUT_TYPEAHEAD_BYTES), + stage: "planned", + purpose: "reserved for IX2: bytes queued for a launching terminal before Activate; not yet enforced", + site: "plans/cmux-tui-zero-wait-interaction.md L4", + }, +]; + +/// Every budget, in declaration order. Names are unique. +pub fn table() -> &'static [Budget] { + &TABLE +} + +/// Look one budget up by name. +pub fn find(name: &str) -> Option<&'static Budget> { + TABLE.iter().find(|budget| budget.name == name) +} + +#[cfg(test)] +mod tests { + use std::collections::BTreeSet; + + use super::*; + + #[test] + fn budget_names_are_unique_and_dotted() { + let mut seen = BTreeSet::new(); + for budget in table() { + assert!(seen.insert(budget.name), "duplicate budget name {}", budget.name); + assert!( + budget.name.chars().all(|c| c.is_ascii_lowercase() || c == '.' || c == '_'), + "budget name {} must be lowercase dotted", + budget.name + ); + assert!(!budget.purpose.is_empty()); + assert!(!budget.site.is_empty()); + } + } + + #[test] + fn every_stage_is_known() { + for budget in table() { + assert!( + STAGES.contains(&budget.stage), + "{} has unknown stage {}", + budget.name, + budget.stage + ); + } + } + + #[test] + fn table_values_match_the_named_constants() { + let expect = |name: &str, value: BudgetValue| { + assert_eq!( + find(name).unwrap_or_else(|| panic!("missing {name}")).value, + value, + "{name}" + ); + }; + expect("host.connect_window", BudgetValue::Duration(HOST_CONNECT_WINDOW)); + expect("host.connect_interval", BudgetValue::Duration(HOST_CONNECT_INTERVAL)); + expect("host.handshake", BudgetValue::Duration(HOST_HANDSHAKE)); + expect("host.snapshot_boundary", BudgetValue::Duration(HOST_SNAPSHOT_BOUNDARY)); + expect("host.terminate_grace", BudgetValue::Duration(HOST_TERMINATE_GRACE)); + expect("host.pty_drain", BudgetValue::Duration(HOST_PTY_DRAIN)); + expect("host.kill_wait", BudgetValue::Duration(HOST_KILL_WAIT)); + expect("host.forced_drain", BudgetValue::Duration(HOST_FORCED_DRAIN)); + expect("host.launch_rollback", BudgetValue::Duration(HOST_LAUNCH_ROLLBACK)); + expect("host.launch_owner", BudgetValue::Duration(HOST_LAUNCH_OWNER)); + expect("host.client_write", BudgetValue::Duration(HOST_CLIENT_WRITE)); + expect("host.control_response", BudgetValue::Duration(HOST_CONTROL_RESPONSE)); + expect("terminal.close_wait", BudgetValue::Duration(TERMINAL_CLOSE_WAIT)); + expect("journal.durable_wait", BudgetValue::Duration(JOURNAL_DURABLE_WAIT)); + expect("journal.commit_result_wait", BudgetValue::Duration(JOURNAL_COMMIT_RESULT_WAIT)); + expect("server.stream_write", BudgetValue::Duration(SERVER_STREAM_WRITE)); + expect( + "server.connection_surface_shutdown", + BudgetValue::Duration(SERVER_CONNECTION_SURFACE_SHUTDOWN), + ); + expect("client.request", BudgetValue::Duration(CLIENT_REQUEST)); + expect("client.write", BudgetValue::Duration(CLIENT_WRITE)); + expect("owner.ensure", BudgetValue::Duration(OWNER_ENSURE)); + expect("owner.poll", BudgetValue::Duration(OWNER_POLL)); + expect("frame", BudgetValue::Duration(FRAME)); + expect("input.typeahead_bytes", BudgetValue::Bytes(INPUT_TYPEAHEAD_BYTES)); + } + + #[test] + fn planned_budgets_are_not_enforced_by_code_sites() { + for budget in table().iter().filter(|budget| budget.stage == "planned") { + assert!( + !budget.site.contains("::"), + "planned budget {} names a code site", + budget.name + ); + } + } + + #[test] + fn values_render_in_one_unit_each() { + assert_eq!(BudgetValue::Duration(Duration::from_millis(1500)).amount(), 1500); + assert_eq!(BudgetValue::Duration(Duration::from_millis(1500)).unit(), "ms"); + assert_eq!(BudgetValue::Bytes(65_536).amount(), 65_536); + assert_eq!(BudgetValue::Bytes(65_536).unit(), "bytes"); + } +} diff --git a/cmux-tui/crates/cmux-tui-core/src/journal_ingress.rs b/cmux-tui/crates/cmux-tui-core/src/journal_ingress.rs index c96f67ba3c66..827e08776df7 100644 --- a/cmux-tui/crates/cmux-tui-core/src/journal_ingress.rs +++ b/cmux-tui/crates/cmux-tui-core/src/journal_ingress.rs @@ -21,8 +21,8 @@ const JOURNAL_DURABLE_BATCH_BYTES: usize = 8 * 1024 * 1024; pub(crate) const TERMINAL_OUTPUT_INGRESS_BYTES: usize = 64 * 1024; const TERMINAL_OUTPUT_BATCH_BYTES: usize = 256 * 1024; const JOURNAL_TERMINAL_FAILURE_RETRY_ATTEMPTS: usize = 6; -const JOURNAL_DURABLE_WAIT: Duration = Duration::from_secs(2); -const JOURNAL_COMMIT_RESULT_WAIT: Duration = Duration::from_secs(1); +const JOURNAL_DURABLE_WAIT: Duration = crate::budgets::JOURNAL_DURABLE_WAIT; +const JOURNAL_COMMIT_RESULT_WAIT: Duration = crate::budgets::JOURNAL_COMMIT_RESULT_WAIT; const JOURNAL_WRITER_SHUTDOWN_WAIT: Duration = Duration::from_secs(1); const JOURNAL_SQLITE_RETRY_SLICE: Duration = Duration::from_millis(100); const JOURNAL_RETRY_INITIAL_DELAY: Duration = Duration::from_millis(100); diff --git a/cmux-tui/crates/cmux-tui-core/src/lib.rs b/cmux-tui/crates/cmux-tui-core/src/lib.rs index 7f3d624e65f5..1f37cf3da27b 100644 --- a/cmux-tui/crates/cmux-tui-core/src/lib.rs +++ b/cmux-tui/crates/cmux-tui-core/src/lib.rs @@ -31,6 +31,7 @@ mod sidebar_resource; mod surface; mod workspace_registry; +pub mod budgets; pub mod layout; pub mod platform; pub mod server; diff --git a/cmux-tui/crates/cmux-tui-core/src/mux.rs b/cmux-tui/crates/cmux-tui-core/src/mux.rs index 0812a3ffe924..b64d5eebc924 100644 --- a/cmux-tui/crates/cmux-tui-core/src/mux.rs +++ b/cmux-tui/crates/cmux-tui-core/src/mux.rs @@ -272,7 +272,7 @@ const CELL_PIXEL_RETRY_MAX_ATTEMPTS: u8 = 4; const KITTY_IMAGE_BUDGET_RETRY_INITIAL: Duration = Duration::from_millis(25); const KITTY_IMAGE_BUDGET_RETRY_MAX: Duration = Duration::from_secs(1); const KITTY_IMAGE_BUDGET_RETRY_MAX_ATTEMPTS: u32 = 4; -const TERMINAL_HOST_CLOSE_WAIT: Duration = Duration::from_secs(4); +const TERMINAL_HOST_CLOSE_WAIT: Duration = crate::budgets::TERMINAL_CLOSE_WAIT; const TERMINAL_READER_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(1); pub(crate) const RENDER_ATTACHMENT_LIMIT: usize = 64; const KITTY_IMAGE_PROCESS_BUDGET_BYTES: u64 = 128 * 1024 * 1024; diff --git a/cmux-tui/crates/cmux-tui-core/src/server.rs b/cmux-tui/crates/cmux-tui-core/src/server.rs index 6af60aeb24f4..49abc968a307 100644 --- a/cmux-tui/crates/cmux-tui-core/src/server.rs +++ b/cmux-tui/crates/cmux-tui-core/src/server.rs @@ -1480,7 +1480,7 @@ impl std::error::Error for DeliveryClassifiedError { } const STREAM_DISCONNECT_POLL: Duration = Duration::from_millis(100); -const STREAM_WRITE_TIMEOUT: Duration = Duration::from_secs(2); +const STREAM_WRITE_TIMEOUT: Duration = crate::budgets::SERVER_STREAM_WRITE; const SHUTDOWN_ACK_FLUSH_TIMEOUT: Duration = Duration::from_secs(5); #[cfg(not(test))] const WEBSOCKET_HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(5); @@ -1529,7 +1529,8 @@ const OUTBOUND_CONNECTION_BYTE_CAPACITY: usize = OUTBOUND_BYTE_CAPACITY * 8; const CLIENT_DETACH_WRITE_TIMEOUT: Duration = Duration::from_millis(100); const CONNECTION_SURFACE_QUEUE_CAPACITY: usize = 256; const CONNECTION_SURFACE_QUEUE_BYTE_CAPACITY: usize = 16 * 1024 * 1024; -const CONNECTION_SURFACE_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(3); +const CONNECTION_SURFACE_SHUTDOWN_TIMEOUT: Duration = + crate::budgets::SERVER_CONNECTION_SURFACE_SHUTDOWN; const SERVER_SURFACE_WORKER_CAPACITY: usize = 16; const SERVER_SURFACE_RETAINED_BYTE_CAPACITY: usize = 16 * 1024 * 1024; const RESOURCE_STREAMS_PER_CLIENT_CAPACITY: usize = 64; diff --git a/cmux-tui/crates/cmux-tui-core/src/terminal_host_runtime.rs b/cmux-tui/crates/cmux-tui-core/src/terminal_host_runtime.rs index 23ad308b25b3..6ed66d1c840d 100644 --- a/cmux-tui/crates/cmux-tui-core/src/terminal_host_runtime.rs +++ b/cmux-tui/crates/cmux-tui-core/src/terminal_host_runtime.rs @@ -47,10 +47,11 @@ const MAX_BLOB: usize = crate::surface::VT_REPLAY_MAX_BYTES; const MAX_ARGV: usize = 256; const MAX_ENV: usize = 1024; const MAX_RENDERER_CAPABILITY_TTL: std::time::Duration = std::time::Duration::from_secs(60); -pub(crate) const CONTROL_RESPONSE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2); -const HOST_HANDSHAKE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2); -const HOST_CONNECT_RETRY_WINDOW: std::time::Duration = std::time::Duration::from_secs(1); -const HOST_CONNECT_RETRY_INTERVAL: std::time::Duration = std::time::Duration::from_millis(10); +pub(crate) const CONTROL_RESPONSE_TIMEOUT: std::time::Duration = + crate::budgets::HOST_CONTROL_RESPONSE; +const HOST_HANDSHAKE_TIMEOUT: std::time::Duration = crate::budgets::HOST_HANDSHAKE; +const HOST_CONNECT_RETRY_WINDOW: std::time::Duration = crate::budgets::HOST_CONNECT_WINDOW; +const HOST_CONNECT_RETRY_INTERVAL: std::time::Duration = crate::budgets::HOST_CONNECT_INTERVAL; const TERMINAL_HOST_PUBLICATION_LOCK_FILE: &str = ".publication.lock"; // Keep live PTY backpressure independent from the extra headroom needed by // one maximum Resized + Colors + targeted acknowledgement transition. @@ -62,7 +63,7 @@ const MAX_HOST_CLIENT_STATE_QUEUED_BYTES: usize = MAX_FRAME_PAYLOAD + 3 * crate::terminal_host_protocol::HEADER_LEN; const MAX_HOST_CLIENT_QUEUED_BYTES: usize = MAX_HOST_CLIENT_OUTPUT_QUEUED_BYTES + MAX_HOST_CLIENT_STATE_QUEUED_BYTES; -const HOST_SNAPSHOT_BOUNDARY_TIMEOUT: std::time::Duration = std::time::Duration::from_millis(1500); +const HOST_SNAPSHOT_BOUNDARY_TIMEOUT: std::time::Duration = crate::budgets::HOST_SNAPSHOT_BOUNDARY; const MAX_SMART_RETAINED_BYTES: usize = 8 * 1024 * 1024; const MAX_SMART_RETAINED_FRAMES: usize = 4096; const HOST_PARSER_QUEUE_CAPACITY: usize = 256; @@ -433,13 +434,13 @@ mod unix { use super::*; static RECORD_TEMP_SEQUENCE: AtomicU64 = AtomicU64::new(1); - const HOST_TERMINATE_GRACE: Duration = Duration::from_millis(250); - const HOST_KILL_WAIT: Duration = Duration::from_secs(2); - const HOST_PTY_DRAIN_GRACE: Duration = Duration::from_millis(250); - const HOST_FORCED_DRAIN_WINDOW: Duration = Duration::from_millis(100); - const HOST_LAUNCH_ROLLBACK_WAIT: Duration = Duration::from_secs(4); - const HOST_LAUNCH_OWNER_TIMEOUT: Duration = Duration::from_secs(5); - const HOST_CLIENT_WRITE_TIMEOUT: Duration = Duration::from_secs(2); + const HOST_TERMINATE_GRACE: Duration = crate::budgets::HOST_TERMINATE_GRACE; + const HOST_KILL_WAIT: Duration = crate::budgets::HOST_KILL_WAIT; + const HOST_PTY_DRAIN_GRACE: Duration = crate::budgets::HOST_PTY_DRAIN; + const HOST_FORCED_DRAIN_WINDOW: Duration = crate::budgets::HOST_FORCED_DRAIN; + const HOST_LAUNCH_ROLLBACK_WAIT: Duration = crate::budgets::HOST_LAUNCH_ROLLBACK; + const HOST_LAUNCH_OWNER_TIMEOUT: Duration = crate::budgets::HOST_LAUNCH_OWNER; + const HOST_CLIENT_WRITE_TIMEOUT: Duration = crate::budgets::HOST_CLIENT_WRITE; const HOST_HANDSHAKE_TRANSIENT_RETRIES: usize = 1; const HOST_EXIT_PERSIST_RETRY_MIN: Duration = Duration::from_millis(100); const HOST_EXIT_PERSIST_RETRY_MAX: Duration = Duration::from_secs(5); diff --git a/cmux-tui/crates/cmux-tui/src/app.rs b/cmux-tui/crates/cmux-tui/src/app.rs index 9e828780ba93..a5e9c8113b1a 100644 --- a/cmux-tui/crates/cmux-tui/src/app.rs +++ b/cmux-tui/crates/cmux-tui/src/app.rs @@ -3855,7 +3855,7 @@ enum RenderAction { Draw, } -const TERMINAL_PAINT_CADENCE: Duration = Duration::from_millis(16); +const TERMINAL_PAINT_CADENCE: Duration = cmux_tui_core::budgets::FRAME; /// Keep terminal parsing lossless while collapsing presentation-only wakes to /// the host's frame cadence. Structural draws remain immediate. diff --git a/cmux-tui/crates/cmux-tui/src/cli.rs b/cmux-tui/crates/cmux-tui/src/cli.rs index 31b876a577b0..3cc8d2d4e5c9 100644 --- a/cmux-tui/crates/cmux-tui/src/cli.rs +++ b/cmux-tui/crates/cmux-tui/src/cli.rs @@ -5,10 +5,14 @@ //! accidentally fall back to the private command protocol. mod command; +mod diag; +mod internal; mod lifecycle; mod raw; mod wire; +use internal::bench; + use std::borrow::Cow; use std::io::{self, Write}; use std::path::PathBuf; @@ -33,6 +37,8 @@ const PUBLIC_SCOPES: &[&str] = &[ "projection", "provider", "raw", + "diag", + "bench", ]; const REMOTE_COMMANDS: &[&str] = &[ @@ -157,6 +163,8 @@ pub fn run(args: &[String], startup_usage: &str) -> i32 { command::run_provider_authority(global, authority) } CommandPlan::RawCommand(command) => raw::run(global, command), + CommandPlan::Diag(plan) => diag::run(global, plan), + CommandPlan::Bench(plan) => bench::run(global, plan), }, Err(failure) => { let message = if matches!(failure.output, OutputMode::Quiet | OutputMode::Human) { @@ -421,6 +429,8 @@ fn scope_help_for( "projection" => Cow::Borrowed(PROJECTION_HELP), "provider" => Cow::Borrowed(PROVIDER_HELP), "raw" => Cow::Borrowed(RAW_HELP), + "diag" => Cow::Borrowed(DIAG_HELP), + "bench" => Cow::Borrowed(BENCH_HELP), _ => Cow::Owned(root_help(&catalog.local_server)), } } @@ -476,6 +486,8 @@ const ROOT_HELP_SCOPES_SUFFIX: &str = "\ projection Read and update frontend projections provider Install private provider authority raw Send an explicit low-level operation + diag Local diagnostics (named budgets) + bench Measure interaction latency against a session Run `cmux --help` for scope-specific paths. "; @@ -677,6 +689,36 @@ USAGE escape for the legacy control protocol and provides no compatibility promise. "; +const BENCH_HELP: &str = "\ +USAGE + cmux bench interact --creates --clients [--typing-probes ] + [--socket | --session ] [--json] + +`bench interact` drives a session as a client and records interaction +latencies: create request to response, request to the tree delta that makes it +visible on a separate subscriber, attach to first frame, close to response, and +one-byte typing latency three ways: on a separate connection while creates are +in flight (typing.separate_conn), on the create connection with one probe +after each create request (typing.same_conn_interleaved: what a keystroke +waits behind 1..K in-flight creates), and on the create connection after the +whole batch (typing.same_conn_after_batch: the wait behind the entire batch). +With no --socket and no --session it starts and stops a throwaway session. +The bench owns the session it runs against: at teardown it closes every +terminal that appeared during the run, because a stopped session keeps its +terminal hosts alive. Run it only against a throwaway or idle session. +Exit status is 1 when any create, close, or probe failed, so a degraded +environment (for example exhausted PTYs) does not pass as a measurement. +It sends only existing commands; it adds no protocol command. +"; + +const DIAG_HELP: &str = "\ +USAGE + cmux diag budgets [--json] + +`diag budgets` prints every named timing and size budget the daemon, terminal +hosts, and clients enforce, with its stage and code site. It needs no session. +"; + #[cfg(test)] mod tests { use super::*; diff --git a/cmux-tui/crates/cmux-tui/src/cli/command.rs b/cmux-tui/crates/cmux-tui/src/cli/command.rs index 318240bcf5e5..f90ecda538ba 100644 --- a/cmux-tui/crates/cmux-tui/src/cli/command.rs +++ b/cmux-tui/crates/cmux-tui/src/cli/command.rs @@ -24,6 +24,8 @@ pub(super) enum CommandPlan { Plugin(PluginPlan), ProviderAuthority(ProviderAuthorityPlan), RawCommand(super::raw::RawCommandPlan), + Diag(super::diag::DiagPlan), + Bench(super::bench::BenchPlan), } #[derive(Clone, Debug)] @@ -171,6 +173,8 @@ pub(super) fn parse(args: &[String]) -> Result { "projection" => parse_projection(&tokens.words[1..], &mut selectors, &mut tokens.flags)?, "provider" => parse_provider(&tokens.words[1..], &mut selectors, &mut tokens.flags)?, "raw" => parse_raw(&tokens.words[1..], &mut tokens.flags)?, + "diag" => parse_diag(&tokens.words[1..])?, + "bench" => parse_bench(&tokens.words[1..], &mut tokens.flags)?, value => return Err(super::unknown_scope(value)), }; tokens.flags.reject_remaining()?; @@ -1653,6 +1657,53 @@ fn parse_provider( } } +fn parse_bench(words: &[String], flags: &mut Flags) -> Result { + match strs(words).as_slice() { + ["interact"] => { + let creates = parse_count(flags, "creates", 20)?; + let clients = parse_count(flags, "clients", 1)?.max(1); + let typing_probes = parse_count(flags, "typing-probes", 0)?; + Ok(CommandPlan::Bench(super::bench::BenchPlan { + creates_per_client: creates, + clients, + typing_probes, + })) + } + [action] => { + Err(UsageError::new(format!("unknown bench action {action:?}; expected interact"))) + } + _ => Err(UsageError::new( + "usage: cmux bench interact --creates --clients [--typing-probes ]", + )), + } +} + +fn parse_count(flags: &mut Flags, name: &str, default: usize) -> Result { + match flags.take(name) { + None => Ok(default), + Some(value) => { + let count = value + .parse::() + .map_err(|_| UsageError::new(format!("--{name} must be a non-negative integer")))?; + const MAX_BENCH_COUNT: usize = 100_000; + if count > MAX_BENCH_COUNT { + return Err(UsageError::new(format!("--{name} must be at most {MAX_BENCH_COUNT}"))); + } + Ok(count) + } +} +} + +fn parse_diag(words: &[String]) -> Result { + match strs(words).as_slice() { + ["budgets"] => Ok(CommandPlan::Diag(super::diag::DiagPlan::Budgets)), + [action] => { + Err(UsageError::new(format!("unknown diag action {action:?}; expected budgets"))) + } + _ => Err(UsageError::new("usage: cmux diag budgets [--json]")), + } +} + fn parse_raw(words: &[String], flags: &mut Flags) -> Result { let refs = strs(words); if refs.as_slice() == ["command"] { @@ -2971,6 +3022,16 @@ mod tests { } } + #[test] + fn diag_budgets_is_a_local_plan() { + assert!(matches!( + parse(&strings(&["diag", "budgets"])).unwrap(), + CommandPlan::Diag(super::super::diag::DiagPlan::Budgets) + )); + assert!(parse(&strings(&["diag", "locks"])).is_err()); + assert!(parse(&strings(&["diag"])).is_err()); + } + #[test] fn coding_agent_hook_management_stays_local() { let CommandPlan::AgentHooks(plan) = diff --git a/cmux-tui/crates/cmux-tui/src/cli/diag.rs b/cmux-tui/crates/cmux-tui/src/cli/diag.rs new file mode 100644 index 000000000000..670dcad94174 --- /dev/null +++ b/cmux-tui/crates/cmux-tui/src/cli/diag.rs @@ -0,0 +1,106 @@ +//! Local diagnostics that need no session connection. + +use std::io::{self, Write}; + +use cmux_tui_core::budgets::{self, Budget}; +use serde_json::{Value, json}; + +use super::{GlobalArgs, OutputMode}; + +#[derive(Clone, Debug, PartialEq, Eq)] +pub(super) enum DiagPlan { + /// Print every named budget. + Budgets, +} + +pub(super) fn run(global: GlobalArgs, plan: DiagPlan) -> i32 { + match plan { + DiagPlan::Budgets => match global.output { + OutputMode::Human => { + let mut stdout = io::stdout().lock(); + let _ = stdout.write_all(budgets_text().as_bytes()); + let _ = stdout.flush(); + 0 + } + output => super::wire::print_local_success(&budgets_json(), output), + }, + } +} + +fn sorted_budgets() -> Vec<&'static Budget> { + let mut rows: Vec<&Budget> = budgets::table().iter().collect(); + rows.sort_by(|left, right| left.name.cmp(right.name)); + rows +} + +pub(super) fn budgets_json() -> Value { + Value::Array( + sorted_budgets() + .into_iter() + .map(|budget| { + json!({ + "name": budget.name, + "value": budget.value.amount(), + "unit": budget.value.unit(), + "stage": budget.stage, + "purpose": budget.purpose, + "site": budget.site, + }) + }) + .collect(), + ) +} + +pub(super) fn budgets_text() -> String { + let rows = sorted_budgets(); + let name_width = rows.iter().map(|b| b.name.len()).max().unwrap_or(4).max(4); + let stage_width = rows.iter().map(|b| b.stage.len()).max().unwrap_or(5).max(5); + let mut out = + format!("{:10} {:10} {: = rows.iter().map(|row| row["name"].as_str().unwrap()).collect(); + let mut sorted = names.clone(); + sorted.sort_unstable(); + assert_eq!(names, sorted); + for row in rows { + let budget = budgets::find(row["name"].as_str().unwrap()).expect("known budget"); + assert_eq!(row["value"].as_u64().unwrap(), budget.value.amount()); + assert_eq!(row["unit"].as_str().unwrap(), budget.value.unit()); + assert_eq!(row["stage"].as_str().unwrap(), budget.stage); + assert_eq!(row["site"].as_str().unwrap(), budget.site); + } + } + + #[test] + fn text_table_lists_every_budget_once() { + let text = budgets_text(); + assert!(text.starts_with("NAME")); + for budget in budgets::table() { + let rows = text + .lines() + .filter(|line| line.split_whitespace().next() == Some(budget.name)) + .count(); + assert_eq!(rows, 1, "{}", budget.name); + } + assert!(text.contains("16 ms")); + assert!(text.contains("65536 bytes")); + } +} diff --git a/cmux-tui/crates/cmux-tui/src/cli/internal/bench.rs b/cmux-tui/crates/cmux-tui/src/cli/internal/bench.rs new file mode 100644 index 000000000000..09a96e2bfe51 --- /dev/null +++ b/cmux-tui/crates/cmux-tui/src/cli/internal/bench.rs @@ -0,0 +1,1158 @@ +//! `cmux-tui bench interact`: a client-side interaction benchmark. +//! +//! It drives a session over the raw control protocol as an ordinary client and +//! records, per user intent, the latencies an interactive frontend or an agent +//! actually feels: request to response, request to the tree delta that makes +//! the new resource visible on a separate subscriber, attach to first frame, +//! and close to response. It adds no protocol command and no resource +//! operation; it only sends existing commands. The output feeds the IX0 +//! baseline of `plans/cmux-tui-zero-wait-interaction.md`. +//! +//! The bench owns the session it runs against: at the end it closes every +//! terminal that appeared during the run (`server stop` keeps terminal hosts +//! alive by design, so a bench that only detached views would leak one host +//! and one shell per create), and it exits non-zero when any create, close, or +//! probe failed so a degraded environment cannot pass as a measurement. + +use std::collections::{HashMap, HashSet}; +use std::io::{self, BufRead, BufReader, Read, Write}; +use std::sync::{Arc, Barrier, Mutex}; +use std::thread; +use std::time::{Duration, Instant}; + +use cmux_tui_core::platform::transport; +use serde_json::{Value, json}; + +use crate::cli::{GlobalArgs, OutputMode}; + +const READ_LIMIT: usize = 16 * 1024 * 1024; +const RPC_TIMEOUT: Duration = Duration::from_secs(20); +/// How long to wait for the visibility delta after a response arrives. +const VISIBILITY_GRACE: Duration = Duration::from_secs(2); +/// How long teardown waits for closed terminals' host processes to exit before +/// reporting them as leaked. Bounded by the host's own SIGKILL escalation. +const HOST_EXIT_AUDIT_WINDOW: Duration = Duration::from_secs(3); + +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) struct BenchPlan { + pub creates_per_client: usize, + pub clients: usize, + pub typing_probes: usize, +} + +pub(crate) fn run(global: GlobalArgs, plan: BenchPlan) -> i32 { + match execute(&global, &plan) { + Ok(report) => { + // Exit 1 when any create, close, or probe errored: percentiles + // over a partial sample are not a measurement. Usage errors are + // 2 and transport failures 3, as in the rest of the CLI. + let failed = !report.errors.is_empty(); + let code = match global.output { + OutputMode::Human => { + let mut out = io::stdout().lock(); + let _ = out.write_all(render_text(&report).as_bytes()); + let _ = out.flush(); + 0 + } + output => crate::cli::wire::print_local_success(&report.to_json(), output), + }; + if failed { 1 } else { code } + } + Err(error) => crate::cli::wire::print_local_error( + &json!({"code":"bench.failed","message":error,"details":{},"retryable":false}), + global.output, + 3, + ), + } +} + +// ---- percentiles -------------------------------------------------------- + +/// Nearest-rank percentile of `samples` in milliseconds (f64). `samples` is +/// sorted in place. Returns `None` for an empty slice. +fn percentile(samples: &mut [f64], quantile: f64) -> Option { + if samples.is_empty() { + return None; + } + samples.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal)); + let rank = (quantile.clamp(0.0, 1.0) * samples.len() as f64).ceil() as usize; + let rank = rank.saturating_sub(1); + Some(samples[rank.min(samples.len() - 1)]) +} + +#[derive(Default)] +struct Metric { + samples: Vec, +} + +impl Metric { + fn record(&mut self, value: Duration) { + self.samples.push(value.as_secs_f64() * 1000.0); + } + + fn summary(&self) -> Option { + if self.samples.is_empty() { + return None; + } + let mut sorted = self.samples.clone(); + Some(MetricSummary { + count: sorted.len(), + p50: percentile(&mut sorted, 0.50).unwrap(), + p90: percentile(&mut sorted, 0.90).unwrap(), + p99: percentile(&mut sorted, 0.99).unwrap(), + max: sorted.iter().copied().fold(f64::MIN, f64::max), + }) + } +} + +#[derive(Clone, Copy)] +struct MetricSummary { + count: usize, + p50: f64, + p90: f64, + p99: f64, + max: f64, +} + +// ---- visibility matching ------------------------------------------------ + +/// One timestamped event seen on the subscriber connection. +#[derive(Clone)] +struct TimedEvent { + at: Instant, + value: Value, +} + +/// Find the earliest event at or after `sent` whose payload references +/// `surface_id`, and return how long after `sent` it arrived. Deltas may +/// arrive before the command response, so callers time from the request write, +/// not from the response. +fn visibility_delay(events: &[TimedEvent], sent: Instant, surface_id: u64) -> Option { + events + .iter() + .filter(|event| event.at >= sent) + .filter(|event| event_references_surface(&event.value, surface_id)) + .map(|event| event.at.duration_since(sent)) + .min() +} + +/// True if `value` mentions `surface_id` as a `surface` field anywhere in the +/// tree-delta payload (the delta carries the created surface in its entity). +fn event_references_surface(value: &Value, surface_id: u64) -> bool { + match value { + Value::Object(map) => { + if map.get("surface").and_then(Value::as_u64) == Some(surface_id) { + return true; + } + map.values().any(|child| event_references_surface(child, surface_id)) + } + Value::Array(items) => { + items.iter().any(|child| event_references_surface(child, surface_id)) + } + _ => false, + } +} + +// ---- raw connection ----------------------------------------------------- + +struct Conn { + reader: BufReader>, + next_id: u64, +} + +impl Conn { + fn open(socket: &std::path::Path) -> Result { + let stream = transport::connect(socket).map_err(|e| format!("connect: {e}"))?; + stream.set_read_timeout(Some(RPC_TIMEOUT)).map_err(|e| format!("timeout: {e}"))?; + stream.set_write_timeout(Some(RPC_TIMEOUT)).map_err(|e| format!("timeout: {e}"))?; + Ok(Self { reader: BufReader::new(stream), next_id: 1 }) + } + + fn send(&mut self, mut request: Value) -> Result { + let id = self.next_id; + self.next_id += 1; + request["id"] = json!(id); + let mut line = serde_json::to_vec(&request).map_err(|e| e.to_string())?; + line.push(b'\n'); + self.reader.get_mut().write_all(&line).map_err(|e| format!("write: {e}"))?; + self.reader.get_mut().flush().map_err(|e| format!("flush: {e}"))?; + Ok(id) + } + + fn read_value(&mut self) -> Result { + let mut bytes = Vec::new(); + let read = self + .reader + .by_ref() + .take((READ_LIMIT + 2) as u64) + .read_until(b'\n', &mut bytes) + .map_err(|e| format!("read: {e}"))?; + if read == 0 { + return Err("connection closed".into()); + } + if !bytes.ends_with(b"\n") { + return Err("partial line".into()); + } + bytes.pop(); + serde_json::from_slice(&bytes).map_err(|e| format!("decode: {e}")) + } + + /// Send a command and return its `data`, ignoring any event lines. + fn request(&mut self, request: Value) -> Result { + let id = self.send(request)?; + loop { + let value = self.read_value()?; + if value.get("event").is_some() { + continue; + } + if value.get("id").and_then(Value::as_u64) != Some(id) { + continue; + } + if value.get("ok").and_then(Value::as_bool) == Some(true) { + return Ok(value.get("data").cloned().unwrap_or(Value::Null)); + } + return Err(value + .get("error") + .and_then(Value::as_str) + .unwrap_or("command failed") + .to_string()); + } + } + + fn identify(&mut self) -> Result { + self.request(json!({"cmd":"identify"})) + } +} + +// ---- execution ---------------------------------------------------------- + +struct SessionGuard { + socket: std::path::PathBuf, + owner: Option, +} + +/// The order in which a create connection submits requests before it drains +/// any response. Two same-connection typing probes exist because they answer +/// different questions: `TypingInterleaved` follows each create request, so +/// its distribution is what one keystroke waits when 1..K creates are in +/// flight ahead of it; `TypingAfterBatch` probes are all submitted after the +/// whole batch, so they share one value, the wait behind the entire batch. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum SubmissionKind { + Create { index: usize, kind: usize }, + TypingInterleaved { index: usize }, + TypingAfterBatch { probe: usize }, +} + +fn same_connection_submission_plan(creates: usize, typing_probes: usize) -> Vec { + let mut submissions = Vec::with_capacity(creates * 2 + typing_probes); + for index in 0..creates { + submissions.push(SubmissionKind::Create { index, kind: index % 3 }); + if typing_probes > 0 { + submissions.push(SubmissionKind::TypingInterleaved { index }); + } + } + submissions.extend((0..typing_probes).map(|probe| SubmissionKind::TypingAfterBatch { probe })); + submissions +} + +/// Terminal ids to `close-terminal` at teardown: every terminal that is not in +/// the pre-bench snapshot and is not already gone. The bench owns the session +/// it runs against, so anything that appeared during the run is its own. +fn teardown_close_plan<'a>( + initial: &HashSet, + current: impl IntoIterator, +) -> Vec { + current + .into_iter() + .filter(|(terminal_id, _)| !initial.contains(*terminal_id)) + .filter(|(_, lifecycle)| !matches!(*lifecycle, "tombstoned" | "exited")) + .map(|(terminal_id, _)| terminal_id.to_string()) + .collect() +} + +struct PendingRequest { + id: u64, + sent: Instant, + kind: SubmissionKind, +} + +/// Keep all create requests unread until the separate-connection probe is on +/// the wire. The second barrier prevents a create worker from draining its +/// connection before that probe has submitted its requests. +struct ProbeGates { + creates_submitted: Barrier, + probes_submitted: Barrier, + release_workers: Barrier, +} + +impl ProbeGates { + fn new(client_count: usize) -> Self { + let parties = client_count + 1; + Self { + creates_submitted: Barrier::new(parties), + probes_submitted: Barrier::new(parties), + release_workers: Barrier::new(parties), + } + } +} + +fn execute(global: &GlobalArgs, plan: &BenchPlan) -> Result { + let (socket, guard) = ensure_session(global)?; + + // Subscriber connection: timestamp every tree delta. + let mut subscriber = Conn::open(&socket)?; + subscriber.identify()?; + subscriber.request(json!({"cmd":"subscribe","tree_events":"deltas"}))?; + let events: Arc>> = Arc::new(Mutex::new(Vec::new())); + let stop = Arc::new(std::sync::atomic::AtomicBool::new(false)); + let subscriber_thread = spawn_subscriber(subscriber, Arc::clone(&events), Arc::clone(&stop)); + + // A baseline terminal to type into. Snapshot the terminal catalog first so + // teardown can close exactly what this run created. + let mut control = Conn::open(&socket)?; + control.identify()?; + let initial_terminals = + list_terminal_ids(&mut control)?.into_iter().map(|(id, _)| id).collect(); + let baseline = control.request(json!({"cmd":"new-workspace"}))?; + let baseline_surface = baseline["surface"].as_u64().ok_or("baseline surface missing")?; + let active_pane = fetch_active_pane(&mut control)?; + + let report = Arc::new(Mutex::new(Report::new(&socket))); + + // Concurrent create loops. Each worker submits its whole create batch + // before reading a response. That gives both typing probes the same + // in-flight create load to compare. + let client_count = plan.clients.max(1); + let gates = Arc::new(ProbeGates::new(client_count)); + let mut handles = Vec::new(); + for client in 0..client_count { + let socket = socket.clone(); + let events = Arc::clone(&events); + let report = Arc::clone(&report); + let creates = plan.creates_per_client; + let pane = active_pane; + let gates = Arc::clone(&gates); + let same_connection = client == 0; + let typing_probes = plan.typing_probes; + handles.push(thread::spawn(move || { + if let Err(error) = run_create_loop( + &socket, + creates, + client, + pane, + same_connection, + typing_probes, + &events, + &report, + &gates, + baseline_surface, + ) { + report.lock().unwrap().errors.push(error); + } + })); + } + + // Wait until every create connection has submitted its batch, then run + // the separate-connection probe while those requests remain unread. + let _ = gates.creates_submitted.wait(); + run_separate_typing_probe(&socket, baseline_surface, plan.typing_probes, &report, &gates); + + for handle in handles { + let _ = handle.join(); + } + stop.store(true, std::sync::atomic::Ordering::Release); + let _ = subscriber_thread.join(); + + close_created_terminals(&mut control, &initial_terminals, &report); + if let Some(owner_pid) = guard.owner.as_ref().map(|owner| owner.pid()) { + // A close-terminal response can precede the host process's own exit + // by a few milliseconds, so poll briefly before calling a host leaked. + let remaining = wait_for_child_terminal_hosts_to_exit(owner_pid, HOST_EXIT_AUDIT_WINDOW); + let mut report = report.lock().unwrap(); + report.hosts_remaining = remaining; + if let Some(remaining) = remaining.filter(|count| *count > 0) { + report.warnings.push(format!( + "{remaining} terminal host process(es) still owned by the bench session at teardown" + )); + } + } + + drop(guard); + Arc::try_unwrap(report) + .map(|m| m.into_inner().unwrap()) + .map_err(|_| "report still shared".into()) +} + +fn run_separate_typing_probe( + socket: &std::path::Path, + surface: u64, + probes: usize, + report: &Arc>, + gates: &ProbeGates, +) { + let mut setup_error = None; + let mut conn = match Conn::open(socket) { + Ok(conn) => Some(conn), + Err(error) => { + setup_error = Some(error); + None + } + }; + let identify_error = conn.as_mut().and_then(|connection| connection.identify().err()); + if let Some(error) = identify_error { + setup_error = Some(error); + conn = None; + } + + // Submit the entire separate-connection probe before releasing workers. + // Responses are drained only after the release barrier, so create response + // timing is not inflated by waiting for every typing probe in sequence. + let mut pending = Vec::new(); + if let Some(connection) = conn.as_mut() { + for _ in 0..probes { + let sent = Instant::now(); + match connection.send(json!({"cmd":"send","surface":surface,"text":"x"})) { + Ok(id) => pending.push((id, sent)), + Err(error) => { + setup_error = Some(format!("typing(separate) send: {error}")); + break; + } + } + } + } + + let _ = gates.probes_submitted.wait(); + let _ = gates.release_workers.wait(); + + if let Some(Err(error)) = + conn.as_mut().map(|connection| drain_separate_typing(connection, pending, report)) + { + setup_error = Some(match setup_error { + Some(previous) => format!("{previous}; {error}"), + None => error, + }); + } + if let Some(error) = setup_error { + report.lock().unwrap().errors.push(format!("typing(separate): {error}")); + } +} + +fn drain_separate_typing( + conn: &mut Conn, + pending: Vec<(u64, Instant)>, + report: &Arc>, +) -> Result<(), String> { + let mut pending_by_id: HashMap = pending.into_iter().collect(); + while !pending_by_id.is_empty() { + let value = conn.read_value()?; + if value.get("event").is_some() { + continue; + } + let Some(id) = value.get("id").and_then(Value::as_u64) else { + continue; + }; + let Some(sent) = pending_by_id.remove(&id) else { + continue; + }; + if value.get("ok").and_then(Value::as_bool) == Some(true) { + report.lock().unwrap().typing_separate.record(sent.elapsed()); + } else { + let error = value.get("error").and_then(Value::as_str).unwrap_or("command failed"); + report.lock().unwrap().errors.push(format!("typing(separate): {error}")); + } + } + Ok(()) +} + +fn command_for_submission(submission: SubmissionKind, pane: u64, surface: u64) -> Value { + match submission { + SubmissionKind::Create { index, .. } => match index % 3 { + 0 => json!({"cmd":"new-workspace"}), + 1 => json!({"cmd":"new-tab"}), + _ => json!({"cmd":"split","pane":pane,"dir":"right"}), + }, + SubmissionKind::TypingInterleaved { .. } | SubmissionKind::TypingAfterBatch { .. } => { + json!({"cmd":"send","surface":surface,"text":"x"}) + } + } +} + +/// `(terminal_id, lifecycle)` for every terminal the daemon knows about. +fn list_terminal_ids(conn: &mut Conn) -> Result, String> { + let data = conn.request(json!({"cmd":"list-terminals"}))?; + Ok(data["terminals"] + .as_array() + .map(|terminals| { + terminals + .iter() + .filter_map(|terminal| { + Some(( + terminal["terminal_id"].as_str()?.to_string(), + terminal["lifecycle"].as_str().unwrap_or("").to_string(), + )) + }) + .collect() + }) + .unwrap_or_default()) +} + +/// Close every terminal this run created, including the baseline typing +/// target and creates that were only `close-surface`d (a view-only close keeps +/// the terminal, and `server stop` keeps its host alive by design). +fn close_created_terminals( + conn: &mut Conn, + initial: &HashSet, + report: &Arc>, +) { + let current = match list_terminal_ids(conn) { + Ok(current) => current, + Err(error) => { + report.lock().unwrap().errors.push(format!("teardown list-terminals: {error}")); + return; + } + }; + let plan = + teardown_close_plan(initial, current.iter().map(|(id, life)| (id.as_str(), life.as_str()))); + for terminal_id in plan { + match conn.request(json!({"cmd":"close-terminal","terminal_id":&terminal_id})) { + Ok(_) => report.lock().unwrap().terminals_closed_at_teardown += 1, + Err(error) => report + .lock() + .unwrap() + .errors + .push(format!("teardown close-terminal {terminal_id}: {error}")), + } + } +} + +/// Count `__terminal-host` processes whose parent is the bench-owned session +/// owner. Hosts do not carry the session in their command line, but the owner +/// that spawned them is still their parent while it runs. `None` when the +/// platform cannot answer. +fn count_child_terminal_hosts(owner_pid: u64) -> Option { + if !cfg!(unix) { + return None; + } + let output = + std::process::Command::new("ps").args(["-axo", "pid=,ppid=,command="]).output().ok()?; + if !output.status.success() { + return None; + } + let listing = String::from_utf8_lossy(&output.stdout); + Some(count_hosts_in_ps_listing(&listing, owner_pid)) +} + +/// Poll `count_child_terminal_hosts` until it reports zero or `window` +/// elapses; returns the final count (`None` where the platform cannot count). +fn wait_for_child_terminal_hosts_to_exit(owner_pid: u64, window: Duration) -> Option { + let deadline = Instant::now() + window; + loop { + let remaining = count_child_terminal_hosts(owner_pid)?; + if remaining == 0 || Instant::now() >= deadline { + return Some(remaining); + } + thread::sleep(Duration::from_millis(50)); + } +} + +fn count_hosts_in_ps_listing(listing: &str, owner_pid: u64) -> u64 { + listing + .lines() + .filter_map(|line| { + let mut fields = line.split_whitespace(); + let _pid = fields.next()?; + let ppid = fields.next()?.parse::().ok()?; + let command = fields.collect::>().join(" "); + (ppid == owner_pid && command.contains("__terminal-host")).then_some(()) + }) + .count() as u64 +} + +// One worker owns one connection and every knob that shapes its submission +// plan; bundling these into a struct would only move the same ten fields. +#[allow(clippy::too_many_arguments)] +fn run_create_loop( + socket: &std::path::Path, + creates: usize, + client: usize, + pane: u64, + same_connection: bool, + typing_probes: usize, + events: &Arc>>, + report: &Arc>, + gates: &ProbeGates, + baseline_surface: u64, +) -> Result<(), String> { + let mut setup_error = None; + let mut conn = match Conn::open(socket) { + Ok(conn) => Some(conn), + Err(error) => { + setup_error = Some(error); + None + } + }; + let identify_error = conn.as_mut().and_then(|connection| connection.identify().err()); + if let Some(error) = identify_error { + setup_error = Some(error); + conn = None; + } + + let mut pending = Vec::new(); + if let Some(connection) = conn.as_mut() { + let submissions = same_connection_submission_plan( + creates, + if same_connection { typing_probes } else { 0 }, + ); + for submission in submissions { + let sent = Instant::now(); + match connection.send(command_for_submission(submission, pane, baseline_surface)) { + Ok(id) => pending.push(PendingRequest { id, sent, kind: submission }), + Err(error) => { + setup_error = Some(format!("create[{client}] send: {error}")); + break; + } + } + } + } + + // All workers reach this point before any response is read. The main + // thread uses this barrier to start the separate-connection probe against + // the same in-flight create load. + let _ = gates.creates_submitted.wait(); + let _ = gates.probes_submitted.wait(); + let _ = gates.release_workers.wait(); + + let Some(mut conn) = conn else { + return Err(setup_error.unwrap_or_else(|| "create connection unavailable".into())); + }; + + let drain_error = drain_pending(&mut conn, pending, client, socket, events, report); + match (setup_error, drain_error) { + (None, result) => result, + (Some(setup_error), Ok(())) => Err(setup_error), + (Some(setup_error), Err(drain_error)) => Err(format!("{setup_error}; {drain_error}")), + } +} + +fn drain_pending( + conn: &mut Conn, + pending: Vec, + client: usize, + socket: &std::path::Path, + events: &Arc>>, + report: &Arc>, +) -> Result<(), String> { + let mut pending_by_id: HashMap = + pending.into_iter().map(|request| (request.id, request)).collect(); + let mut completed_creates = Vec::new(); + + while !pending_by_id.is_empty() { + let value = conn.read_value()?; + if value.get("event").is_some() { + continue; + } + let Some(id) = value.get("id").and_then(Value::as_u64) else { + continue; + }; + let Some(request) = pending_by_id.remove(&id) else { + continue; + }; + if value.get("ok").and_then(Value::as_bool) != Some(true) { + let error = value.get("error").and_then(Value::as_str).unwrap_or("command failed"); + match request.kind { + SubmissionKind::Create { kind, .. } => { + report.lock().unwrap().errors.push(format!("create[{client}:{kind}]: {error}")); + } + SubmissionKind::TypingInterleaved { .. } => { + report.lock().unwrap().errors.push(format!("typing(interleaved): {error}")); + } + SubmissionKind::TypingAfterBatch { .. } => { + report.lock().unwrap().errors.push(format!("typing(after-batch): {error}")); + } + } + continue; + } + + match request.kind { + SubmissionKind::Create { kind, .. } => { + completed_creates.push(( + request.sent, + request.sent.elapsed(), + kind, + value.get("data").cloned().unwrap_or(Value::Null), + )); + } + SubmissionKind::TypingInterleaved { .. } => { + report.lock().unwrap().typing_same_interleaved.record(request.sent.elapsed()); + } + SubmissionKind::TypingAfterBatch { .. } => { + report.lock().unwrap().typing_same_after_batch.record(request.sent.elapsed()); + } + } + } + + // No responses remain on this connection, so the close requests below + // cannot consume another pending request while we process each create. + for (sent, response, _kind, data) in completed_creates { + record_create_result(conn, sent, response, data, socket, events, report); + } + Ok(()) +} + +fn record_create_result( + conn: &mut Conn, + sent: Instant, + response: Duration, + data: Value, + socket: &std::path::Path, + events: &Arc>>, + report: &Arc>, +) { + let surface = data["surface"].as_u64(); + let terminal_id = data.get("terminal_id").and_then(Value::as_str).map(str::to_owned); + { + let mut report = report.lock().unwrap(); + report.create_response.record(response); + report.record_lifecycle(data.get("lifecycle").and_then(Value::as_str)); + } + + if let Some(surface_id) = surface { + // Give the delta a moment; it may already be recorded. + let delay = wait_for_visibility(events, sent, surface_id, VISIBILITY_GRACE); + if let Some(delay) = delay { + report.lock().unwrap().create_visible.record(delay); + } else { + report.lock().unwrap().visibility_misses += 1; + } + + if let Some(first_frame) = measure_first_frame(socket, surface_id) { + report.lock().unwrap().first_frame.record(first_frame); + } + + // View-only close of this surface (default destroy for a tab). + let close_start = Instant::now(); + match conn.request(json!({"cmd":"close-surface","surface":surface_id})) { + Ok(_) => report.lock().unwrap().close_surface.record(close_start.elapsed()), + Err(error) => report.lock().unwrap().errors.push(format!("close-surface: {error}")), + } + } + + // For terminals with a stable id, also measure the process-terminating + // close, which blocks on host exit escalation (terminal.close_wait). + if let Some(terminal_id) = terminal_id { + let close_start = Instant::now(); + match conn.request(json!({"cmd":"close-terminal","terminal_id":&terminal_id})) { + Ok(_) => report.lock().unwrap().close_terminal.record(close_start.elapsed()), + Err(error) => { + // Teardown closes by catalog difference, so a failure here is + // reported, not fatal. + report + .lock() + .unwrap() + .warnings + .push(format!("close-terminal {terminal_id}: {error}")); + } + } + } +} + +fn wait_for_visibility( + events: &Arc>>, + sent: Instant, + surface_id: u64, + grace: Duration, +) -> Option { + let deadline = Instant::now() + grace; + loop { + if let Some(delay) = visibility_delay(&events.lock().unwrap(), sent, surface_id) { + return Some(delay); + } + if Instant::now() >= deadline { + return None; + } + thread::sleep(Duration::from_millis(1)); + } +} + +fn is_first_frame_for_surface(value: &Value, surface_id: u64) -> bool { + // Attach streams share the connection's event channel, so an unrelated + // render-state event must not satisfy this surface's frame measurement. + value.get("event").and_then(Value::as_str) == Some("render-state") + && value.get("surface").and_then(Value::as_u64) == Some(surface_id) +} + +fn measure_first_frame(socket: &std::path::Path, surface_id: u64) -> Option { + let mut conn = Conn::open(socket).ok()?; + conn.identify().ok()?; + let start = Instant::now(); + let id = + conn.send(json!({"cmd":"attach-surface","surface":surface_id,"mode":"render"})).ok()?; + let deadline = Instant::now() + RPC_TIMEOUT; + loop { + if Instant::now() >= deadline { + return None; + } + let value = conn.read_value().ok()?; + if is_first_frame_for_surface(&value, surface_id) { + return Some(start.elapsed()); + } + // A failed attach response ends the attempt. + if value.get("id").and_then(Value::as_u64) == Some(id) + && value.get("ok").and_then(Value::as_bool) == Some(false) + { + return None; + } + } +} + +fn fetch_active_pane(conn: &mut Conn) -> Result { + let tree = conn.request(json!({"cmd":"list-workspaces"}))?; + let workspaces = tree["workspaces"].as_array().ok_or("no workspaces")?; + let workspace = workspaces + .iter() + .find(|ws| ws["active"].as_bool() == Some(true)) + .or_else(|| workspaces.last()) + .ok_or("no active workspace")?; + let screens = workspace["screens"].as_array().ok_or("no screens")?; + let screen = screens + .iter() + .find(|s| s["active"].as_bool() == Some(true)) + .or_else(|| screens.first()) + .ok_or("no screen")?; + screen["active_pane"].as_u64().ok_or_else(|| "no active pane".into()) +} + +fn spawn_subscriber( + mut conn: Conn, + events: Arc>>, + stop: Arc, +) -> thread::JoinHandle<()> { + thread::spawn(move || { + while !stop.load(std::sync::atomic::Ordering::Acquire) { + match conn.read_value() { + Ok(value) if value.get("event").is_some() => { + events.lock().unwrap().push(TimedEvent { at: Instant::now(), value }); + } + Ok(_) => {} + Err(_) => break, + } + } + }) +} + +fn ensure_session(global: &GlobalArgs) -> Result<(std::path::PathBuf, SessionGuard), String> { + if let Some(socket) = &global.socket { + return Ok((socket.clone(), SessionGuard { socket: socket.clone(), owner: None })); + } + if let Some(session) = &global.session { + let socket = cmux_tui_core::server::try_default_socket_path(session) + .map_err(|e| format!("socket path: {e}"))?; + let owner = crate::local_owner::ensure_owner_for_bench(session, &socket)?; + return Ok((socket.clone(), SessionGuard { socket, owner: Some(owner) })); + } + let session = format!("bench-{:08x}", fastrand_u32()); + let socket = cmux_tui_core::server::try_default_socket_path(&session) + .map_err(|e| format!("socket path: {e}"))?; + let owner = crate::local_owner::ensure_owner_for_bench(&session, &socket)?; + Ok((socket.clone(), SessionGuard { socket, owner: Some(owner) })) +} + +impl Drop for SessionGuard { + fn drop(&mut self) { + if let Some(owner) = self.owner.take() { + owner.stop(&self.socket); + } + } +} + +fn fastrand_u32() -> u32 { + let mut buf = [0u8; 4]; + getrandom::fill(&mut buf).ok(); + u32::from_le_bytes(buf) +} + +// ---- report ------------------------------------------------------------- + +struct Report { + socket: String, + create_response: Metric, + create_visible: Metric, + first_frame: Metric, + close_surface: Metric, + close_terminal: Metric, + typing_separate: Metric, + typing_same_interleaved: Metric, + typing_same_after_batch: Metric, + lifecycle_counts: std::collections::BTreeMap, + visibility_misses: u64, + terminals_closed_at_teardown: u64, + hosts_remaining: Option, + warnings: Vec, + errors: Vec, +} + +impl Report { + fn new(socket: &std::path::Path) -> Self { + Self { + socket: socket.display().to_string(), + create_response: Metric::default(), + create_visible: Metric::default(), + first_frame: Metric::default(), + close_surface: Metric::default(), + close_terminal: Metric::default(), + typing_separate: Metric::default(), + typing_same_interleaved: Metric::default(), + typing_same_after_batch: Metric::default(), + lifecycle_counts: std::collections::BTreeMap::new(), + visibility_misses: 0, + terminals_closed_at_teardown: 0, + hosts_remaining: None, + warnings: Vec::new(), + errors: Vec::new(), + } + } + + fn record_lifecycle(&mut self, lifecycle: Option<&str>) { + if let Some(lifecycle) = lifecycle { + *self.lifecycle_counts.entry(lifecycle.to_string()).or_insert(0) += 1; + } + } + + fn metrics(&self) -> [(&'static str, &Metric); 8] { + [ + ("create.response_ms", &self.create_response), + ("create.visible_ms", &self.create_visible), + ("create.first_frame_ms", &self.first_frame), + ("close.surface_response_ms", &self.close_surface), + ("close.terminal_response_ms", &self.close_terminal), + ("typing.separate_conn_ms", &self.typing_separate), + ("typing.same_conn_interleaved_ms", &self.typing_same_interleaved), + ("typing.same_conn_after_batch_ms", &self.typing_same_after_batch), + ] + } + + fn to_json(&self) -> Value { + let mut metrics = serde_json::Map::new(); + for (name, metric) in self.metrics() { + if let Some(summary) = metric.summary() { + metrics.insert( + name.to_string(), + json!({ + "count": summary.count, + "p50": round(summary.p50), + "p90": round(summary.p90), + "p99": round(summary.p99), + "max": round(summary.max), + }), + ); + } + } + json!({ + "commit": option_env!("CMUX_TUI_BUILD_COMMIT").unwrap_or("unknown"), + "platform": std::env::consts::OS, + "socket": self.socket, + "metrics": Value::Object(metrics), + "lifecycle_counts": self.lifecycle_counts, + "visibility_misses": self.visibility_misses, + "terminals_closed_at_teardown": self.terminals_closed_at_teardown, + "hosts_remaining": self.hosts_remaining, + "warnings": self.warnings, + "errors": self.errors, + }) + } +} + +fn round(value: f64) -> f64 { + (value * 1000.0).round() / 1000.0 +} + +fn render_text(report: &Report) -> String { + let mut out = String::new(); + out.push_str(&format!( + "cmux-tui bench interact ({}, socket {})\n", + std::env::consts::OS, + report.socket + )); + // Failures first: a table over a partial sample must not read as a result. + match report.errors.first() { + Some(first) => out.push_str(&format!("errors: {} (first: {first})\n", report.errors.len())), + None => out.push_str("errors: 0\n"), + } + out.push_str(&format!("lifecycle on create response: {:?}\n", report.lifecycle_counts)); + out.push_str(&format!( + "{:<34}{:>6}{:>10}{:>10}{:>10}{:>10}\n", + "metric", "n", "p50", "p90", "p99", "max" + )); + for (name, metric) in report.metrics() { + if let Some(s) = metric.summary() { + out.push_str(&format!( + "{:<34}{:>6}{:>10.2}{:>10.2}{:>10.2}{:>10.2}\n", + name, s.count, s.p50, s.p90, s.p99, s.max + )); + } + } + if report.visibility_misses > 0 { + out.push_str(&format!( + "visibility misses (no delta within grace): {}\n", + report.visibility_misses + )); + } + out.push_str(&format!( + "teardown: closed {} terminal(s); hosts still owned by the bench session: {}\n", + report.terminals_closed_at_teardown, + report.hosts_remaining.map_or("unknown".to_string(), |count| count.to_string()) + )); + for warning in &report.warnings { + out.push_str(&format!("warning: {warning}\n")); + } + if report.errors.len() > 1 { + for error in report.errors.iter().skip(1).take(9) { + out.push_str(&format!(" {error}\n")); + } + } + out +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn percentile_is_nearest_rank() { + let mut data = vec![10.0, 20.0, 30.0, 40.0, 50.0]; + assert_eq!(percentile(&mut data.clone(), 0.0), Some(10.0)); + assert_eq!(percentile(&mut data.clone(), 0.5), Some(30.0)); + assert_eq!(percentile(&mut data, 1.0), Some(50.0)); + assert_eq!(percentile(&mut Vec::::new(), 0.5), None); + let mut pair = vec![10.0, 20.0]; + assert_eq!(percentile(&mut pair, 0.5), Some(10.0)); + } + + #[test] + fn metric_summarizes_counts_and_max() { + let mut metric = Metric::default(); + for ms in [5u64, 1, 3, 9, 7] { + metric.record(Duration::from_millis(ms)); + } + let summary = metric.summary().unwrap(); + assert_eq!(summary.count, 5); + assert_eq!(summary.max, 9.0); + assert_eq!(summary.p50, 5.0); + } + + #[test] + fn visibility_matches_earliest_referencing_event() { + let base = Instant::now(); + let events = vec![ + TimedEvent { at: base, value: json!({"event":"tab-added","entity":{"surface":7}}) }, + TimedEvent { + at: base + Duration::from_millis(5), + value: json!({"event":"tab-added","surface":42,"index":1}), + }, + TimedEvent { + at: base + Duration::from_millis(9), + value: json!({"event":"tab-added","surface":42,"index":2}), + }, + ]; + // Sent one ms before the first matching event at +5ms. + let sent = base + Duration::from_millis(4); + let delay = visibility_delay(&events, sent, 42).unwrap(); + assert_eq!(delay, Duration::from_millis(1)); + // A surface never referenced returns None. + assert!(visibility_delay(&events, sent, 999).is_none()); + // An event before `sent` is ignored. + assert!(visibility_delay(&events, base + Duration::from_millis(6), 7).is_none()); + } + + #[test] + fn deep_entity_reference_is_found() { + let value = json!({ + "event":"workspace-added", + "entity":{"screens":[{"panes":[{"tabs":[{"surface":99}]}]}]} + }); + assert!(event_references_surface(&value, 99)); + assert!(!event_references_surface(&value, 98)); + } + + #[test] + fn same_connection_probe_is_submitted_before_response_drain() { + let submissions = same_connection_submission_plan(3, 2); + assert_eq!( + submissions, + vec![ + SubmissionKind::Create { index: 0, kind: 0 }, + SubmissionKind::TypingInterleaved { index: 0 }, + SubmissionKind::Create { index: 1, kind: 1 }, + SubmissionKind::TypingInterleaved { index: 1 }, + SubmissionKind::Create { index: 2, kind: 2 }, + SubmissionKind::TypingInterleaved { index: 2 }, + SubmissionKind::TypingAfterBatch { probe: 0 }, + SubmissionKind::TypingAfterBatch { probe: 1 }, + ] + ); + } + + #[test] + fn interleaved_probe_follows_each_create_and_needs_probes_enabled() { + let submissions = same_connection_submission_plan(4, 1); + let creates = + submissions.iter().filter(|s| matches!(s, SubmissionKind::Create { .. })).count(); + let interleaved = submissions + .iter() + .filter(|s| matches!(s, SubmissionKind::TypingInterleaved { .. })) + .count(); + assert_eq!((creates, interleaved), (4, 4)); + for pair in submissions.windows(2) { + if let SubmissionKind::Create { index, .. } = pair[0] { + assert_eq!(pair[1], SubmissionKind::TypingInterleaved { index }); + } + } + // With typing probes disabled the plan is creates only. + assert!( + same_connection_submission_plan(2, 0) + .iter() + .all(|s| matches!(s, SubmissionKind::Create { .. })) + ); + } + + #[test] + fn teardown_closes_every_terminal_created_during_the_run() { + let initial: HashSet = + ["pre-a".to_string(), "pre-b".to_string()].into_iter().collect(); + let current = [ + ("pre-a", "running"), + ("pre-b", "exited"), + ("baseline", "running"), + ("new-tab", "running"), + ("split", "launching"), + ("already-gone", "tombstoned"), + ("finished", "exited"), + ]; + let plan = teardown_close_plan(&initial, current.iter().copied()); + assert_eq!(plan, vec!["baseline", "new-tab", "split"]); + // Nothing created: nothing closed, including pre-existing terminals. + assert!(teardown_close_plan(&initial, [("pre-a", "running")]).is_empty()); + } + + #[test] + fn host_count_matches_owner_children_only() { + let listing = "\ + 100 1 /usr/bin/cmux-tui --headless --session bench-1 + 101 100 /usr/bin/cmux-tui __terminal-host --bootstrap-stdio + 102 100 /usr/bin/cmux-tui __terminal-host --bootstrap-stdio + 103 999 /usr/bin/cmux-tui __terminal-host --bootstrap-stdio + 104 100 /bin/zsh -l +"; + assert_eq!(count_hosts_in_ps_listing(listing, 100), 2); + assert_eq!(count_hosts_in_ps_listing(listing, 999), 1); + assert_eq!(count_hosts_in_ps_listing(listing, 7), 0); + } + + #[test] + fn first_frame_requires_the_requested_surface() { + assert!(!is_first_frame_for_surface(&json!({"event":"render-state","surface":41}), 42)); + assert!(is_first_frame_for_surface(&json!({"event":"render-state","surface":42}), 42)); + assert!(!is_first_frame_for_surface(&json!({"event":"render-delta","surface":42}), 42)); + } +} diff --git a/cmux-tui/crates/cmux-tui/src/cli/internal/mod.rs b/cmux-tui/crates/cmux-tui/src/cli/internal/mod.rs new file mode 100644 index 000000000000..439c3e8ef3ef --- /dev/null +++ b/cmux-tui/crates/cmux-tui/src/cli/internal/mod.rs @@ -0,0 +1,3 @@ +//! CLI implementations that speak the private control protocol directly. + +pub(super) mod bench; diff --git a/cmux-tui/crates/cmux-tui/src/local_owner.rs b/cmux-tui/crates/cmux-tui/src/local_owner.rs index 8c2a2d0a9012..fd6fe3f52d65 100644 --- a/cmux-tui/crates/cmux-tui/src/local_owner.rs +++ b/cmux-tui/crates/cmux-tui/src/local_owner.rs @@ -22,13 +22,13 @@ use std::process::{Command, Stdio}; use std::time::{Duration, Instant}; use cmux_tui_core::platform::{self, transport}; -use serde_json::Value; +use serde_json::{Value, json}; /// Total time an ensure may spend probing, spawning, and waiting for the /// owner to accept clients. Matches the lifecycle CLI exchange deadline. -pub(crate) const ENSURE_DEADLINE: Duration = Duration::from_secs(10); +pub(crate) const ENSURE_DEADLINE: Duration = cmux_tui_core::budgets::OWNER_ENSURE; -const POLL_INTERVAL: Duration = Duration::from_millis(25); +const POLL_INTERVAL: Duration = cmux_tui_core::budgets::OWNER_POLL; /// The reaper's only job is clearing a zombie when the owner exits early /// (a lost bind race or a crash), and `terminate` reaps synchronously @@ -60,6 +60,7 @@ pub(crate) enum Ensured { Started(ReadyOwner), } +#[derive(Debug)] pub(crate) enum EnsureError { /// The owner process could not be spawned (or the spawn lock failed). Spawn(io::Error), @@ -84,6 +85,102 @@ enum Attempt { Ready(ReadyOwner), } +/// A throwaway bench-owned session: a spawned owner and the temporary state +/// root to clean up when the bench stops. +pub(crate) struct EnsuredOwnerHandle { + pid: u64, + generation: String, + state_root: Option, +} + +impl EnsuredOwnerHandle { + /// Process id of the owner; terminal hosts it spawned are its children + /// while it runs, which is how the bench audits for leaked hosts. + pub(crate) fn pid(&self) -> u64 { + self.pid + } + + /// Ask the owner to shut down (best effort) and remove the temp state root. + pub(crate) fn stop(self, socket: &Path) { + let deadline = Instant::now() + Duration::from_secs(5); + if let Ok(stream) = transport::connect(socket) { + let _ = stream.set_read_timeout(Some(Duration::from_secs(2))); + let _ = stream.set_write_timeout(Some(Duration::from_secs(2))); + let mut connection = BufReader::new(stream); + let request = json!({ + "id": 1, + "cmd": "shutdown-daemon", + "pid": self.pid, + "generation": self.generation, + }); + if writeln!(connection.get_mut(), "{request}") + .and_then(|()| connection.get_mut().flush()) + .is_ok() + { + // Drain until the socket closes or the deadline passes. + loop { + if Instant::now() >= deadline { + break; + } + let mut bytes = Vec::new(); + match connection.by_ref().take(4096).read_until(b'\n', &mut bytes) { + Ok(0) | Err(_) => break, + Ok(_) => {} + } + } + } + } + if let Some(root) = self.state_root { + let _ = std::fs::remove_dir_all(root); + // `SocketStartLock` deliberately leaves `.spawn-lock` in + // place for durable sessions, because unlinking it reopens the + // stale-socket start race for that session name. A bench session + // name is random and never started again, so removing its lock + // after the owner we spawned has been asked to exit leaves nothing + // behind under the runtime directory. + let mut name = socket.file_name().unwrap_or_default().to_os_string(); + name.push(".spawn-lock"); + let _ = std::fs::remove_file(socket.with_file_name(name)); + } + } +} + +/// Spawn (or adopt) a headless owner for a bench session and return a handle +/// that can stop it. Uses a private temporary state root so a throwaway +/// session never touches the default durable state. +pub(crate) fn ensure_owner_for_bench( + session: &str, + socket: &Path, +) -> Result { + let state_root = std::env::temp_dir().join(format!("cmux-bench-{session}")); + std::fs::create_dir_all(&state_root).map_err(|error| format!("state dir: {error}"))?; + let spec = OwnerSpec { + session: session.to_string(), + socket: socket.to_path_buf(), + socket_is_derived: true, + state: Some(state_root.clone()), + term: None, + }; + let deadline = Instant::now() + ENSURE_DEADLINE; + match ensure_owner(&spec, Some(session), deadline) { + Ok(Ensured::Running(ready)) => Ok(EnsuredOwnerHandle { + pid: ready.pid, + generation: ready.generation, + // The owner predates this bench; do not remove its state root. + state_root: None, + }), + Ok(Ensured::Started(ready)) => Ok(EnsuredOwnerHandle { + pid: ready.pid, + generation: ready.generation, + state_root: Some(state_root), + }), + Err(error) => { + let _ = std::fs::remove_dir_all(&state_root); + Err(format!("ensure bench owner: {error:?}")) + } + } +} + /// Connect to a ready owner for `spec.session` on `spec.socket`, spawning a /// detached headless owner process when none is running. `expected_session` /// is the caller's session constraint for a server that is already running; diff --git a/cmux-tui/crates/cmux-tui/src/session/remote.rs b/cmux-tui/crates/cmux-tui/src/session/remote.rs index fb369f66a2de..45ea81223d0b 100644 --- a/cmux-tui/crates/cmux-tui/src/session/remote.rs +++ b/cmux-tui/crates/cmux-tui/src/session/remote.rs @@ -76,7 +76,7 @@ const INTERACTIVE_LATENCY_BUCKET_UPPER_US: [u64; 18] = [ ]; #[cfg(not(test))] fn remote_write_timeout() -> Duration { - Duration::from_secs(2) + cmux_tui_core::budgets::CLIENT_WRITE } #[cfg(test)] @@ -92,7 +92,7 @@ fn remote_write_timeout() -> Duration { }) } #[cfg(not(test))] -const REMOTE_REQUEST_TIMEOUT: Duration = Duration::from_secs(10); +const REMOTE_REQUEST_TIMEOUT: Duration = cmux_tui_core::budgets::CLIENT_REQUEST; #[cfg(test)] const REMOTE_REQUEST_TIMEOUT: Duration = Duration::from_millis(100); #[cfg(not(test))] diff --git a/cmux-tui/docs/README.md b/cmux-tui/docs/README.md index cfaf226c80bb..adb394bb69d5 100644 --- a/cmux-tui/docs/README.md +++ b/cmux-tui/docs/README.md @@ -15,5 +15,6 @@ - [Raw control protocol](protocol.md): private protocol-v12 JSON-lines commands for cmux frontends and compatibility adapters. - [Journal operations / ジャーナル操作](journal-operations.md): reading `cmux server stats` to find registry lock convoys, journal writer batch shape, and connection refusals. `cmux server stats` でレジストリロックの競合、ジャーナル書き込みのバッチ状況、接続拒否を確認します。 - [Public CLI](../spec/cli.md): noun-first commands and selectors. +- [Journal and interaction operations](journal-operations.md): named budgets (`diag budgets`) and the `bench interact` interaction-latency benchmark. - [SDK contract](../spec/bindings.md): handwritten facades and generated raw layers. - [Browser panes](browser-panes.md): CDP-backed browser tabs, rendering, input, profiles, and current limitations. diff --git a/cmux-tui/docs/journal-operations.md b/cmux-tui/docs/journal-operations.md index 20ef64c83aca..8dcc76378e6c 100644 --- a/cmux-tui/docs/journal-operations.md +++ b/cmux-tui/docs/journal-operations.md @@ -1,70 +1,89 @@ -# Journal operations: reading `server stats` +# Journal and interaction operations -The daemon reports where its own time goes. Ask it before sampling the -process: +How to observe cmux-tui's runtime timing from the outside: the named budgets +every wait is bound by, and the interaction benchmark that measures how those +waits add up for a user or an agent. + +## Named budgets + +Every bounded wait in the daemon, the terminal hosts, and the clients is a +budget with one name, one value, and one interaction stage. `cmux_tui_core::budgets` +holds the values; the code sites that enforce them import those constants, so +there is one value per budget. ```bash -cmux server stats --json # exact object, spec/commands.md "server-stats" -cmux server stats # nested key: value lines +cmux diag budgets # table +cmux diag budgets --json # array of {name,value,unit,stage,purpose,site} ``` -Counters accumulate since daemon start and reading them costs a few atomic -loads, so polling is safe. Everything below names a field of that object. - -## The write path in one paragraph - -Every durable write goes through SQLite in WAL mode with `synchronous=FULL` -and `fullfsync`, about 20 ms per commit on an Apple SSD. Producers (agent -hooks, frontends) enqueue into bounded lanes; one writer thread drains both -lanes into a single transaction, commits, and hands each producer a receipt. -The writer takes the workspace registry mutex for the commit. Request threads -take the same mutex for their own commits, so it is the one lock whose -contention shapes latency under load. `spec/session-journal.md` "Retention and -storage" describes the intended single-writer shape; the plan to finish it is -the cmuxterm-hq journal write path plan. - -## Reading `registry_lock` - -- `holder` names the `file:line` that holds the registry mutex right now and - for how long. Under a healthy daemon this is usually `null`; a holder older - than a few milliseconds is in an fsync. -- `hold_us` is the hold-time histogram. Its p50 near 20,000 means most holds - are one fsync. `top_sites` says which code holds it longest in total; that is - the list to shorten. -- `wait_us` is how long acquirers waited. A p90 far above `hold_us.p50` means a - convoy: many threads queue behind each hold. `contended_acquisitions` (waits - of 1 ms or more) and `stalls` (100 ms or more) count it, and `last_stall` - names the waiter and the site that held the lock when the wait began. - -## Reading `journal_writer` - -- `batch_size.mean` is the group-commit ratio. Under many concurrent producers - it should climb well above 1. A mean near 1 under load means producers are - throttled before they reach the queue, which was the symptom of the resolver - convoy fixed in cmux PR 11630. -- `commit_us` is transaction time including the fsync, without lock wait. - `commit_lock_wait_us` is the writer waiting for the registry mutex; if it - rivals `commit_us`, request threads are starving the writer. -- `receipt_wait_us` is what a producer such as `cmux-tui-hook` waited from - enqueue to receipt. It bounds the hook's detached child lifetime. -- `terminal_queued` and `durable_queued` are live lane depths. - `deadline_expiries` and `commit_failures` should stay at zero; either one - rising means the writer is being interrupted or refused. -- `phase` and `phase_for_us` say what the writer is doing now; a long - `waiting_lock` phase points at `registry_lock.holder`. - -## Reading `connections` - -`limit` is the control-socket cap. `refused` counts sockets dropped at the cap; -each refused hook connection is a lost agent event, so any non-zero value is a -capacity incident, not a statistic. `peak` shows how close normal operation -comes to the cap. - -## Benchmarking - -Until the in-binary load generator lands, use the Python scripts recorded in -the cmuxterm-hq plan (`bench4.py`: N clients over K live terminals through a -blocking hook helper) and read `server stats` before and after. Numbers on the -2026-09-02 main with 16 terminals, dev build, Apple M-series laptop: 16 clients -p50 185 ms at 75 events/s; 48 clients p50 621 ms at 59 events/s; the remaining -ceiling is the per-event projection commit under the registry lock. +`stage` places a budget in the interaction lifecycle: + +- `accept`: validating and applying a mutation to in-memory state. +- `durable`: the journal batch that carries the mutation has committed. +- `settle`: an external effect (host launch, terminate, first frame) reached its + outcome. +- `frame`: a frontend paint cadence. +- `client`: a client-side wait on the daemon. +- `planned`: reserved by design, not enforced yet (for example + `input.typeahead_bytes`, reserved for the launching-terminal typeahead queue). + +A timeout error names the budget it exhausted, so an operator or an agent can +map a stall to one row of this table. + +## Interaction bench + +`cmux bench interact` drives a session as an ordinary client over the raw +control protocol and records the latencies a frontend or an agent actually +feels. It sends only existing commands; it is not a protocol command or a +resource operation. + +```bash +# Throwaway session (started and stopped by the bench): +cmux bench interact --creates 20 --clients 1 --typing-probes 50 +cmux --json bench interact --creates 20 --clients 8 --typing-probes 100 + +# Against a running session: +cmux --socket /path/to/session.sock bench interact --creates 20 --clients 4 +``` + +Metrics, each reported as `count`, `p50`, `p90`, `p99`, `max` in milliseconds: + +| metric | what it times | +| --- | --- | +| `create.response_ms` | create request write to the command response | +| `create.visible_ms` | create request write to the first tree delta on a separate `tree_events:"deltas"` subscriber that references the new surface (deltas may precede the response, so both are timed from the write) | +| `create.first_frame_ms` | `attach-surface` with `mode:"render"` to the first `render-state` | +| `close.surface_response_ms` | `close-surface` (view-only close) request to response | +| `close.terminal_response_ms` | `close-terminal` (process-terminating close) request to response; this is the one that waits on host exit escalation, bounded by `terminal.close_wait` | +| `typing.separate_conn_ms` | one-byte `send` on a connection that issues no creates, while creates are in flight | +| `typing.same_conn_interleaved_ms` | one-byte `send` on the connection that also issues creates, one probe submitted right after each create request; the distribution is what a keystroke waits when 1..K creates are queued ahead of it on the same connection | +| `typing.same_conn_after_batch_ms` | one-byte `send` on the create connection, all probes submitted after the whole create batch and before its responses are drained; every probe waits behind the entire batch, so p50 equals p99 and the value is the batch tail, not per-keystroke latency. A gap between either same-connection number and the separate-connection number is head-of-line blocking | + +`--clients N` runs N concurrent create loops on N connections; `--creates K` is +creates per client; `--typing-probes M` sets the number of typing samples (the +interleaved probe is one per create and is enabled whenever M is non-zero). The +JSON output also carries `lifecycle_counts` (the `lifecycle` field on each +create response), `visibility_misses` (creates whose visibility delta did not +arrive within the grace window), `terminals_closed_at_teardown`, +`hosts_remaining` (terminal host processes still parented by the bench-owned +session owner after teardown; `null` where the platform cannot count them), +`warnings`, `errors`, `commit`, and `platform`. The text output prints the error +count and first error and the lifecycle counts above the table, and the table +carries `n` per metric. + +The bench owns the session it runs against. At teardown it lists the terminal +catalog and closes every terminal that was not present before the run, +including the baseline typing target and creates that were only detached with +`close-surface`, because `server stop` keeps terminal hosts alive by design and +a view-only close leaves the terminal and its shell running. Run it only +against a throwaway session (the default) or a session nobody else is +mutating. It exits 1 when any create, close, or probe failed: percentiles over +a partial sample are not a measurement, and PTY exhaustion on a loaded machine +is the usual cause of that failure. + +Read the numbers against the budget table: a `create.response_ms` far above the +`accept`-stage cost, or a `typing.same_conn_interleaved_ms` far above +`typing.separate_conn_ms`, is the interaction cost the zero-wait work removes. The record-only +`bench interact` job in `.github/workflows/cmux-tui.yml` publishes the JSON per +commit as the `cmux-tui-bench-interact-` artifact. Design and targets: +`plans/cmux-tui-zero-wait-interaction.md` (IX0). diff --git a/cmux-tui/scripts/test_check_resource_api_boundary.py b/cmux-tui/scripts/test_check_resource_api_boundary.py index 9f920ec1de1b..69b05314701f 100644 --- a/cmux-tui/scripts/test_check_resource_api_boundary.py +++ b/cmux-tui/scripts/test_check_resource_api_boundary.py @@ -295,6 +295,10 @@ def matching_contract(tui: Path, operations: list[str] | None = None) -> None: class PublicBoundaryScanTests(unittest.TestCase): + def test_repository_public_boundary_is_clean(self) -> None: + diagnostics, _ = CHECKER.scan_public_boundaries(CHECKER.TUI) + self.assertEqual(diagnostics, []) + def test_raw_internal_and_manifest_generated_occurrences_are_allowed(self) -> None: with tempfile.TemporaryDirectory() as directory: tui = Path(directory) diff --git a/cmux-tui/spec/cli.md b/cmux-tui/spec/cli.md index c4fd674d9cf5..6011b419edcc 100644 --- a/cmux-tui/spec/cli.md +++ b/cmux-tui/spec/cli.md @@ -108,7 +108,7 @@ The public resource roots are: ```text server machine session client workspace screen pane tab terminal browser notification agent sidebar -pairing projection provider raw +pairing projection provider raw diag bench ``` Structural resources may be addressed directly by opaque ID or through their @@ -129,6 +129,10 @@ map every operational one-shot command and parameter in Sensitive renderer grants and connection-owned stream/viewer controls remain SDK and raw-only. +`diag` and `bench` are local diagnostic scopes. `bench interact` uses the +private control protocol internally to measure existing commands; it adds no +resource operation or protocol command. + ## Selectors An instance selector accepts: @@ -363,6 +367,21 @@ plugin names are slugs matching `[a-z0-9-_]+`. installs the credential for an already running provider-managed session and is not a transported resource operation or cross-machine discovery API. +## Diagnostics + +```text +cmux diag budgets [--json] +``` + +`diag budgets` runs locally and needs no session. It prints every named timing +and size budget that the daemon, terminal hosts, and clients enforce: the +dotted `name`, the `value` and `unit` (`ms` or `bytes`), the interaction +`stage` it belongs to (`accept`, `durable`, `settle`, `frame`, `client`, or +`planned` for a budget reserved by design and not yet enforced), a one-line +`purpose`, and the code `site` that holds the value. The constants live in +`cmux_tui_core::budgets`; a timeout error names the budget it exhausted. The +verb is CLI-local and is not a raw protocol command or a resource operation. + ## Raw access ```text