From dafab5c439a626aeb3d3016113264b901cf4e1a2 Mon Sep 17 00:00:00 2001 From: AuricStudio Date: Sat, 22 Aug 2026 04:54:03 +0800 Subject: [PATCH 1/3] Fix Codex ownership across runtime restart MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Retain and reap adoption workers, cascade parent-loss cleanup through Codex process groups, and fence process creation against concurrent shutdown. Keep recovery on durable lease redelivery instead of predecessor-runtime scratch discovery. Tests: cargo test -p fkst-framework --bin fkst-framework shutdown_fences_an_in_progress_process_group_spawn; cargo test -p fkst-framework --test supervise_smoke supervise_restart_terminates_codex_tree_and_redelivers_expired_lease ⟦AI:FKST⟧ --- SPEC.md | 6 +- crates/fkst-framework/src/process_tree.rs | 64 +++++- crates/fkst-framework/src/sdk_codex.rs | 39 +++- .../fkst-framework/src/supervise/spawner.rs | 2 +- .../fkst-framework/tests/supervise_smoke.rs | 214 ++++++++++++++++++ docs/architecture.md | 8 +- 6 files changed, 321 insertions(+), 12 deletions(-) diff --git a/SPEC.md b/SPEC.md index f204749..edb5a99 100644 --- a/SPEC.md +++ b/SPEC.md @@ -88,7 +88,7 @@ The outcome matrix is normative. “Catchable” means `SIGINT` or `SIGTERM` rea | Shipped process root | Uncatchable death of `fkst-supervisor` | Running event runtime and Department invocation | The event runtime and its child may continue, and their already committed external effects survive | The surviving runtime retains delivery ownership; an unacked reliable lease becomes eligible after lease expiry if that runtime also dies | Possible between the surviving runtime generation and an externally started replacement generation | The process root's ability to wait for or clean up the event-runtime process group | | Shipped process root | Catchable signal directly to `fkst-framework supervise`; after runtime cleanup, `fkst-supervisor` observes its exit and exits without replacement | Running Department invocation | Only effects already committed to an authoritative external system | An unacked reliable lease becomes eligible after lease expiry | Possible if process-group cleanup does not complete before an external manager starts a replacement | In-memory execution state and unaccepted output | | Direct event runtime | Catchable signal to `fkst-framework supervise` | Running Department invocation | Only effects already committed to an authoritative external system | An unacked reliable lease becomes eligible after lease expiry | Possible if cleanup does not complete before an external manager starts a replacement | In-memory execution state and unaccepted output | -| Either | Uncatchable event-runtime death | Running Department invocation | A child process may continue, and its already committed external effects survive | An unacked reliable lease becomes eligible after lease expiry | Possible between the old execution and lease-expiry replay | The old owner's ability to ACK or accept further child output | +| Either | Uncatchable event-runtime death | Running Department invocation | The invocation detects owner loss, terminates and reaps its registered descendants, and exits; already committed external effects survive | An unacked reliable lease becomes eligible after lease expiry | Possible only until parent-loss cleanup completes if an external manager starts a replacement immediately | The old owner's ability to ACK or accept further child output | | Either | Catchable or uncatchable | Pending or leased reliable delivery, or DLQ row | The record in `FKST_DURABLE_ROOT` | Pending work remains due; an unacked lease is eligible after expiry | A stale execution may overlap replay, but lease-generation fencing rejects its ACK | No acknowledged durable delivery is reconstructed from memory | | Either | Catchable or uncatchable | Ephemeral event, fanout/wake, source timer, Lua state, or raise buffer | Nothing in this class | None; a later event exists only if a source re-derives it | None guaranteed | The in-memory value | | Either | Catchable or uncatchable | Authenticated raised output | A raised event already published to a reliable destination survives there; an ephemeral destination survives only while its runtime remains alive | Reliable destination rules apply after publication | Publication before parent ACK can cause replayed work, with deterministic delivery identity providing destination deduplication | Any frame not accepted by its owning runtime | @@ -97,6 +97,10 @@ The outcome matrix is normative. “Catchable” means `SIGINT` or `SIGTERM` rea The engine guarantees no graph reload, no atomic engine/package version pair, and no coherence between the startup-scanned graph and package or runner bytes read by later invocations. Authenticated child output is bound to the event runtime that created its token and pipe; it is never handed to another runtime generation. The engine also does not guarantee that catchable cleanup completes before an external lifecycle manager starts a replacement. Deployment-manager replacement, ordering, canary, and rollback are outside this repository. Mechanism and source locations are described in `docs/architecture.md`. +Each supervised Department invocation receives a private expected-parent PID and monitors it for its full lifetime. A parent change is an ownership loss: the invocation terminates all registered SDK process groups, waits for their direct-child reapers, and exits without accepting further output. A worktree-backed Codex adoption worker repeats this boundary for `codex exec`; the Department retains a worker reaper and registers the worker process group before continuing. This chain makes both catchable shutdown and uncatchable event-runtime death converge on descendant termination and reaping. + +Codex adoption records remain current-runtime scratch. A replacement runtime does not scan, migrate, or adopt predecessor runtime directories. Reliable recovery comes from the durable delivery record: after the old owner disappears, its unacknowledged lease expires, `lease_generation` advances, and the replacement starts fresh work under its own runtime layout. Existing lease-generation fencing rejects a stale ACK if cleanup and externally managed replacement briefly overlap. + ## SDK surface - 固定 Lua SDK surface 锚点是 `fixed-lua-sdk-surface`;允许 surface 是 `pipeline`、`source`、`raise`、`spawn_codex_sync`、`spawn_codex`、`fkst.codex_runs`、`fkst.observe`、`exec_sync`、`env_read`、`await_all`、`with_lock`、`once`、`cache_set`、`cache_get`、`cache_expire`、`graph_json`、`t`、`restricted_lua_load`、`git_log_count`、`git_log_grep`、`count_worktrees`、`list_orphan_worktrees`、`setup_worktree`、`file`、`json.decode`、`toml.decode`、`log.info`、`log.warn`、`log.error`、`now`。 diff --git a/crates/fkst-framework/src/process_tree.rs b/crates/fkst-framework/src/process_tree.rs index 4e5698c..b634055 100644 --- a/crates/fkst-framework/src/process_tree.rs +++ b/crates/fkst-framework/src/process_tree.rs @@ -2,6 +2,7 @@ use nix::errno::Errno; use nix::sys::signal::{killpg, SaFlags, SigAction, SigHandler, SigSet, Signal}; use nix::unistd::Pid; use std::collections::BTreeSet; +use std::io::Write; use std::sync::atomic::{AtomicBool, AtomicI32, Ordering}; use std::sync::OnceLock; use std::sync::{Arc, Mutex}; @@ -10,6 +11,7 @@ use tracing::{info, warn}; const TERMINATION_GRACE: Duration = Duration::from_secs(2); const POLL_INTERVAL: Duration = Duration::from_millis(25); +pub(crate) const SUPERVISOR_PID_ENV: &str = "FKST_SUPERVISOR_PID"; static SDK_PROCESS_GROUPS: OnceLock = OnceLock::new(); static SIGNAL_WATCH_INSTALLED: OnceLock<()> = OnceLock::new(); static SHUTDOWN_REQUESTED: AtomicBool = AtomicBool::new(false); @@ -57,6 +59,18 @@ impl ProcessGroupRegistry { send_group_signal(pgid, Signal::SIGKILL, label); } } + let deadline = Instant::now() + TERMINATION_GRACE; + while Instant::now() < deadline { + if self + .snapshot() + .iter() + .all(|pgid| !process_group_exists(*pgid)) + { + return; + } + tokio::time::sleep(POLL_INTERVAL).await; + } + self.warn_survivors(label); } pub(crate) fn terminate_all_blocking(&self, label: &str) { @@ -79,6 +93,18 @@ impl ProcessGroupRegistry { send_group_signal(pgid, Signal::SIGKILL, label); } } + let deadline = Instant::now() + TERMINATION_GRACE; + while Instant::now() < deadline { + if self + .snapshot() + .iter() + .all(|pgid| !process_group_exists(*pgid)) + { + return; + } + std::thread::sleep(POLL_INTERVAL); + } + self.warn_survivors(label); } fn snapshot(&self) -> Vec { @@ -89,6 +115,18 @@ impl ProcessGroupRegistry { .copied() .collect() } + + fn warn_survivors(&self, label: &str) { + for pgid in self.snapshot() { + if process_group_exists(pgid) { + warn!( + pgid = pgid, + label = label, + "process group survived SIGKILL grace" + ); + } + } + } } pub(crate) struct ProcessGroupRegistration { @@ -133,23 +171,39 @@ pub(crate) fn install_sdk_shutdown_watch() { SIGNAL_WATCH_INSTALLED.get_or_init(|| { install_signal_handler(Signal::SIGTERM); install_signal_handler(Signal::SIGINT); - std::thread::spawn(|| loop { + let expected_parent_pid = supervisor_parent_pid(); + std::thread::spawn(move || loop { if SHUTDOWN_REQUESTED.load(Ordering::SeqCst) { sdk_process_groups().terminate_all_blocking("sdk child"); let signal = SHUTDOWN_SIGNAL.load(Ordering::SeqCst); std::process::exit(128 + signal); } + if let Some(expected_parent_pid) = + expected_parent_pid.filter(|pid| parent_changed(*pid)) + { + let _ = writeln!( + std::io::stderr().lock(), + "[framework] process owner parent lost: expected parent pid {}, current parent pid {}", + expected_parent_pid, + current_parent_pid() + ); + sdk_process_groups().terminate_all_blocking("sdk child after parent loss"); + std::process::exit(125); + } std::thread::sleep(POLL_INTERVAL); }); }); } -#[cfg(not(test))] -pub(crate) fn ensure_supervisor_parent_alive() -> mlua::Result<()> { - let Some(expected_parent_pid) = std::env::var("FKST_SUPERVISOR_PID") +fn supervisor_parent_pid() -> Option { + std::env::var(SUPERVISOR_PID_ENV) .ok() .and_then(|value| value.parse::().ok()) - else { +} + +#[cfg(not(test))] +pub(crate) fn ensure_supervisor_parent_alive() -> mlua::Result<()> { + let Some(expected_parent_pid) = supervisor_parent_pid() else { return Ok(()); }; if parent_changed(expected_parent_pid) { diff --git a/crates/fkst-framework/src/sdk_codex.rs b/crates/fkst-framework/src/sdk_codex.rs index 15c093f..097d0e5 100644 --- a/crates/fkst-framework/src/sdk_codex.rs +++ b/crates/fkst-framework/src/sdk_codex.rs @@ -1666,13 +1666,49 @@ fn start_adoption_worker( command.stdin(Stdio::null()); command.stdout(Stdio::null()); command.stderr(Stdio::null()); + command.env( + crate::process_tree::SUPERVISOR_PID_ENV, + crate::process_tree::current_pid().to_string(), + ); #[cfg(unix)] { command.process_group(0); inherit_fd_on_exec(&mut command, run_lock_fd); } let child = command.spawn().map_err(mlua::Error::external)?; - drop(child); + let worker_pid = child.id(); + let registration = crate::process_tree::sdk_process_groups().register(worker_pid); + let child_slot = Arc::new(Mutex::new(Some(child))); + let reaper_child_slot = Arc::clone(&child_slot); + if let Err(err) = std::thread::Builder::new() + .name("fkst-codex-worker-reaper".to_string()) + .spawn(move || { + if let Some(mut child) = reaper_child_slot + .lock() + .expect("codex adoption worker child slot poisoned") + .take() + { + if let Err(err) = child.wait() { + eprintln!( + "WARN: codex adoption worker wait failed pid={worker_pid} error={err}" + ); + } + } + drop(registration); + }) + { + kill_process_group_by_pid(worker_pid); + if let Some(mut child) = child_slot + .lock() + .expect("codex adoption worker child slot poisoned") + .take() + { + let _ = child.wait(); + } + return Err(mlua::Error::external(format!( + "codex adoption worker reaper spawn failed: {err}" + ))); + } Ok(()) } @@ -2017,6 +2053,7 @@ pub(crate) fn parse_worker_args(args: Vec) -> anyhow::Result anyhow::Result { + crate::process_tree::install_sdk_shutdown_watch(); let run_lock = inherited_adoption_run_lock(&options)?; std::fs::create_dir_all(&options.work_dir)?; let prompt = std::fs::read_to_string(&options.prompt_file)?; diff --git a/crates/fkst-framework/src/supervise/spawner.rs b/crates/fkst-framework/src/supervise/spawner.rs index 7f5d7d3..0122f33 100644 --- a/crates/fkst-framework/src/supervise/spawner.rs +++ b/crates/fkst-framework/src/supervise/spawner.rs @@ -158,7 +158,7 @@ pub async fn spawn_framework_with_stdout_observer( codex_permit_slots.to_string(), ); cmd.env( - "FKST_SUPERVISOR_PID", + crate::process_tree::SUPERVISOR_PID_ENV, crate::process_tree::current_pid().to_string(), ); if let Some(token) = raised_auth_token { diff --git a/crates/fkst-framework/tests/supervise_smoke.rs b/crates/fkst-framework/tests/supervise_smoke.rs index e60b2a4..50b5036 100644 --- a/crates/fkst-framework/tests/supervise_smoke.rs +++ b/crates/fkst-framework/tests/supervise_smoke.rs @@ -116,6 +116,19 @@ fn process_exists(pid: i32) -> bool { nix::sys::signal::kill(nix::unistd::Pid::from_raw(pid), None).is_ok() } +fn process_parent_pid(pid: i32) -> i32 { + let output = Command::new("ps") + .args(["-o", "ppid=", "-p", &pid.to_string()]) + .output() + .unwrap(); + assert!(output.status.success(), "ps status={}", output.status); + String::from_utf8(output.stdout) + .unwrap() + .trim() + .parse() + .unwrap() +} + fn read_single_supervisor_journal(runtime_root: &Path) -> String { let entries = fs::read_dir(runtime_root.join("logs")) .unwrap() @@ -712,6 +725,207 @@ return M ); } +#[test] +fn supervise_restart_terminates_codex_tree_and_redelivers_expired_lease() { + let _lock = supervise_smoke_lock(); + let tmp = tempfile::tempdir().unwrap(); + let root = tmp.path(); + let bin_dir = root.join("bin"); + let capture_dir = root.join("capture"); + let worktree = root.join("worktree"); + let runtime_one = root.join("runtime-one"); + let runtime_two = root.join("runtime-two"); + let durable_root = root.join("durable"); + let input = root.join("input.txt"); + let completed = root.join("completed.txt"); + let second_release = root.join("second-release"); + fs::create_dir_all(root.join("departments/worker")).unwrap(); + fs::create_dir_all(root.join("raisers")).unwrap(); + fs::create_dir_all(&bin_dir).unwrap(); + fs::create_dir_all(&capture_dir).unwrap(); + fs::create_dir_all(&worktree).unwrap(); + write_fkst_env(root); + fs::write(&input, "ready").unwrap(); + fs::write( + root.join("raisers/input.lua"), + format!( + r#"return {{ type = "file_watch", glob = {}, produces = "jobs" }}"#, + lua_string(&input) + ), + ) + .unwrap(); + fs::write( + root.join("departments/worker/main.lua"), + format!( + r#" +local M = {{}} +M.spec = {{ consumes = {{ "jobs" }}, stall_window = "1s" }} +function pipeline(event) + local result = spawn_codex_sync({{ + prompt = "restart-owned-work", + worktree = {}, + dedup_key = "restart-owned-work", + timeout = 60, + }}) + assert(result.exit_code == 0, result.error or result.stderr) + local f = assert(io.open({}, "w")) + f:write(result.stdout) + f:close() +end +return M +"#, + lua_string(&worktree), + lua_string(&completed) + ), + ) + .unwrap(); + write_executable( + &bin_dir.join("codex"), + r#"#!/bin/sh +set -eu +cat >/dev/null +count_file="$CAPTURE_DIR/spawn-count" +count=0 +if [ -f "$count_file" ]; then + count=$(cat "$count_file") +fi +count=$((count + 1)) +printf '%s' "$count" > "$count_file" +printf '%s\n' "$$" > "$CAPTURE_DIR/codex-$count.pid" +printf '%s\n' "$PPID" > "$CAPTURE_DIR/worker-$count.pid" +if [ "$count" -eq 1 ]; then + printf 'started' > "$CAPTURE_DIR/first-started" + while :; do sleep 1; done +fi +printf 'started' > "$CAPTURE_DIR/second-started" +while [ ! -f "$SECOND_RELEASE" ]; do sleep 0.05; done +printf '%s\n' '{"type":"item.completed","item":{"type":"agent_message","text":"restart-complete"}}' +"#, + ); + write_single_package_workspace(root); + + let mut search_path = vec![bin_dir.clone()]; + search_path.extend(std::env::split_paths( + &std::env::var_os("PATH").unwrap_or_default(), + )); + let search_path = std::env::join_paths(search_path).unwrap(); + let start_runtime = |runtime_root: &Path| { + Command::new(env!("CARGO_BIN_EXE_fkst-framework")) + .current_dir(root) + .arg("supervise") + .arg("--project-root") + .arg(root) + .arg("--package-root") + .arg(root) + .arg("--framework-bin") + .arg(env!("CARGO_BIN_EXE_fkst-framework")) + .env("PATH", &search_path) + .env("CAPTURE_DIR", &capture_dir) + .env("SECOND_RELEASE", &second_release) + .env("FKST_RUNTIME_ROOT", runtime_root) + .env("FKST_RUNTIME_LOG_DIR", runtime_root.join("codex-logs")) + .env("FKST_DURABLE_ROOT", &durable_root) + .spawn() + .unwrap() + }; + + let mut first = start_runtime(&runtime_one); + wait_for_file_containing( + &capture_dir.join("first-started"), + "started", + Duration::from_secs(10), + ) + .unwrap_or_else(|| { + let _ = first.kill(); + let _ = first.wait(); + panic!("timed out waiting for first Codex invocation"); + }); + let worker_pid: i32 = fs::read_to_string(capture_dir.join("worker-1.pid")) + .unwrap() + .trim() + .parse() + .unwrap(); + let codex_pid: i32 = fs::read_to_string(capture_dir.join("codex-1.pid")) + .unwrap() + .trim() + .parse() + .unwrap(); + let department_pid = process_parent_pid(worker_pid); + + first.kill().unwrap(); + let first_status = first.wait().unwrap(); + assert!(!first_status.success(), "status={first_status}"); + let survivors = [ + ("Department", department_pid), + ("adoption worker", worker_pid), + ("Codex", codex_pid), + ] + .into_iter() + .filter(|(_, pid)| !wait_for_process_exit(*pid, Duration::from_secs(8))) + .collect::>(); + for (_, pid) in &survivors { + let _ = nix::sys::signal::killpg( + nix::unistd::Pid::from_raw(*pid), + nix::sys::signal::Signal::SIGKILL, + ); + } + assert!(survivors.is_empty(), "surviving processes={survivors:?}"); + + let mut second = start_runtime(&runtime_two); + wait_for_file_containing( + &capture_dir.join("second-started"), + "started", + Duration::from_secs(15), + ) + .unwrap_or_else(|| { + let _ = second.kill(); + let _ = second.wait(); + panic!("timed out waiting for redelivered Codex invocation"); + }); + let observe = Command::new(env!("CARGO_BIN_EXE_fkst-framework")) + .arg("observe") + .arg("--durable-root") + .arg(&durable_root) + .arg("--json") + .env_remove(process_tree::SUPERVISOR_PID_ENV) + .output() + .unwrap(); + assert!( + observe.status.success(), + "observe stderr={}", + String::from_utf8_lossy(&observe.stderr) + ); + let snapshot: serde_json::Value = serde_json::from_slice(&observe.stdout).unwrap(); + assert_eq!(snapshot["deliveries"][0]["lease_generation"], 2); + + fs::write(&second_release, "release").unwrap(); + let result = wait_for_file_containing(&completed, "restart-complete", Duration::from_secs(10)) + .unwrap_or_else(|| { + let _ = second.kill(); + let _ = second.wait(); + panic!("timed out waiting for redelivered result"); + }); + wait_for_journal_event(&runtime_two, "acked", "worker", Duration::from_secs(5)).unwrap_or_else( + || { + let _ = second.kill(); + let _ = second.wait(); + panic!("timed out waiting for redelivered delivery ACK"); + }, + ); + nix::sys::signal::kill( + nix::unistd::Pid::from_raw(second.id() as i32), + nix::sys::signal::Signal::SIGTERM, + ) + .unwrap(); + let second_status = second.wait().unwrap(); + assert!(second_status.success(), "status={second_status}"); + assert_eq!(result, "restart-complete"); + assert_eq!( + fs::read_to_string(capture_dir.join("spawn-count")).unwrap(), + "2" + ); +} + #[test] fn supervise_delivers_cross_package_raise_from_composed_child() { let _lock = supervise_smoke_lock(); diff --git a/docs/architecture.md b/docs/architecture.md index 1645a6b..4af8246 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -223,7 +223,7 @@ queue 是包内命名空间。多 graph-root 组合时,裸 queue 名按 owner - `with_lock`、`once` 与 `cache` 共用 runtime key 合约:key / name 必须是非空相对 filesystem path,`/` 表示目录;每个 segment 非空、最长 255 bytes、匹配 `[A-Za-z0-9._-]+`,且不是全点 segment(如 `.` / `..` / `...`);禁止 leading / trailing `/`、`//`、反斜杠、NUL 与绝对路径。校验后的 key 保持为 `/{locks,marks,cache}//` 目录路径,engine 在该目录下写 reserved leaf file(`=lock` / `=mark` / `=value`),形成可人工浏览的目录树,不做 byte hex 编码。`=` 不在合法 key segment 字符集内,因此不会与有效 key 冲突。`locks/once/` 是 `once` 内部锁的保留子目录,不属于 `with_lock` 用户锁命名空间。 - `file_watch` 只接受 host-root 相对或绝对 glob;不支持 runtime scheme。 - codex log **不属** `RuntimeKind`/``:`sdk_codex` 把它落到 `FKST_RUNTIME_LOG_DIR` 或平台默认目录(如 `~/Library/Logs/fkst`)下的 `codex/`。它与 `/logs` 同属 process-trace scratch(可 grep、非事实源),但落点不同,`supervise` 也不给 framework child 注入 `FKST_RUNTIME_LOG_DIR`。每次 `spawn_codex_sync` / `spawn_codex` 会在写入本次 log 前 best-effort 修剪同一个 `codex/` 目录中的旧 `.log` 文件:`FKST_CODEX_LOG_MAX_AGE` 默认 `48h`,`0` 或空值表示关闭年龄修剪;`FKST_CODEX_LOG_MAX_BYTES` 为空或 `0` 表示不启用容量上限,启用时优先删除最旧 log;当前请求的 log path 永远豁免,删除或扫描失败只写 warning。 -- `fkst.codex_runs()` remains observational: append order is the lifecycle authority within each generic run log, while wall-clock fields only order distinct projected runs. A generic log-backed `running` record is projected as running only while its per-`run_id` log-file lifetime witness is locked, and an unlocked record is projected as failed without rewriting the log or creating keyed redrive semantics. The reader accepts only status records whose `run_id` matches the ULID-backed identity encoded in the engine-owned log filename. Engine status snapshots form the leading control header; raw stdout/stderr is a byte-counted frame between `CODEX_OUTPUT_BEGIN:` and `CODEX_OUTPUT_END`, and a final `completed` snapshot is accepted only after the complete frame. A torn frame therefore cannot promote output-controlled text into lifecycle state or borrow the witness. With `worktree`, `spawn_codex_sync` writes prompt, stdout, stderr, a trusted effect receipt log, effect marker, result marker, and completion status into writable runtime scratch at `/logs/codex-adoption//`; `key` is derived from optional `dedup_key`, `worktree`, and prompt/context identity. While holding the key's `run.lock`, the framework writes and rereads a visible `intent` status before allowing the detached `fkst-framework __codex-worker` to claim `running` and execute the real `codex exec`; the worker must also claim visible intent under the same lock, preventing duplicate effects if the parent exits after spawn but before the running record is visible. The worker passes its per-`run_id` locked witness through the `codex exec` boundary, so a live codex process tree keeps that incarnation authoritative even if the worker exits; `run_id` carries a process-independent ULID token so witness names do not repeat across framework process incarnations. After `codex exec` returns, the worker first appends a `CODEX_EFFECT:` keyed receipt to the adoption-local trusted receipt log, then publishes the adoption-local keyed effect marker and result marker; redrive can recover completed status from the result marker, effect marker, or trusted `effect_key` receipt without repeating the effect. The human/debug Codex `log_path` remains observational and is not a recovery receipt source. Re-delivery of the same work unit checks this handoff first: completed reads the result, running waits for the same worker/codex, and stale intent or dead running owner redrives and claims before effect without storm-spawning `codex`. Different worktrees do not adopt each other even with the same `dedup_key`, prompt, and context; different `dedup_key` or prompt/context under the same worktree also remain distinct. This directory is `RuntimeLayout`-owned process-trace scratch used for supervisor-generation adoption; it is not written into the worktree and does not carry package business schema, accepted-state, or rollback state. Without `worktree`, `spawn_codex_sync` keeps pipe/wait semantics. +- `fkst.codex_runs()` remains observational: append order is the lifecycle authority within each generic run log, while wall-clock fields only order distinct projected runs. A generic log-backed `running` record is projected as running only while its per-`run_id` log-file lifetime witness is locked, and an unlocked record is projected as failed without rewriting the log or creating keyed redrive semantics. The reader accepts only status records whose `run_id` matches the ULID-backed identity encoded in the engine-owned log filename. Engine status snapshots form the leading control header; raw stdout/stderr is a byte-counted frame between `CODEX_OUTPUT_BEGIN:` and `CODEX_OUTPUT_END`, and a final `completed` snapshot is accepted only after the complete frame. A torn frame therefore cannot promote output-controlled text into lifecycle state or borrow the witness. With `worktree`, `spawn_codex_sync` writes prompt, stdout, stderr, a trusted effect receipt log, effect marker, result marker, and completion status into writable runtime scratch at `/logs/codex-adoption//`; `key` is derived from optional `dedup_key`, `worktree`, and prompt/context identity. While holding the key's `run.lock`, the framework writes and rereads a visible `intent` status before allowing the separate `fkst-framework __codex-worker` to claim `running` and execute the real `codex exec`; the worker must also claim visible intent under the same lock, preventing duplicate effects if the parent exits after spawn but before the running record is visible. The Department registers that worker process group and retains a direct-child reaper. The worker monitors its Department parent, registers the `codex exec` process group, and retains the Codex reaper. The worker passes its per-`run_id` locked witness through the `codex exec` boundary, so a live codex process tree keeps that incarnation authoritative while its owning invocation remains alive; `run_id` carries a process-independent ULID token so witness names do not repeat across framework process incarnations. After `codex exec` returns, the worker first appends a `CODEX_EFFECT:` keyed receipt to the adoption-local trusted receipt log, then publishes the adoption-local keyed effect marker and result marker; redelivery in the same runtime can recover completed status from the result marker, effect marker, or trusted `effect_key` receipt without repeating the effect. The human/debug Codex `log_path` remains observational and is not a recovery receipt source. Another invocation for the same work unit in the same runtime checks this handoff first: completed reads the result, running waits for the same worker/codex, and stale intent or dead running owner redrives and claims before effect without storm-spawning `codex`. Different worktrees do not adopt each other even with the same `dedup_key`, prompt, and context; different `dedup_key` or prompt/context under the same worktree also remain distinct. This directory is `RuntimeLayout`-owned process-trace scratch for current-runtime duplicate suppression and recovery. A replacement runtime does not discover or adopt it; reliable work is redelivered from `FKST_DURABLE_ROOT` after lease expiry. It is not written into the worktree and does not carry package business schema, accepted-state, or rollback state. Without `worktree`, `spawn_codex_sync` keeps pipe/wait semantics. - engine **不写** runtime 持久状态;accepted-state / rollback 是外部 release pipeline 的事实,见 §13。 ## 7. 运行态数据流 @@ -348,15 +348,15 @@ validation 规则: The engine supports two entry topologies. `fkst-supervisor` is the shipped process root: it starts one `fkst-framework supervise` child as a process-group leader, inherits its stdout/stderr, and waits. Direct `fkst-framework supervise` is also supported; in that topology its lifecycle manager is external. The engine does not identify which topology an operator uses. The normative outcome matrix is in `SPEC.md` under “Runtime generations and termination outcomes.” -On `SIGINT` or `SIGTERM`, `fkst-supervisor` sends `SIGTERM` to the event-runtime process group, waits up to two seconds, then sends `SIGKILL` if it has not exited. The event runtime has its own `SIGINT`/`SIGTERM` handlers. A caught signal makes it send `SIGTERM` to every registered Department invocation process group, wait up to two seconds, send `SIGKILL` to survivors, abort runtime tasks, and exit. The supervisor does not launch a replacement. The event runtime is placed in an independent process group, and the supervisor installs no parent-death signal or other handoff. `SIGKILL` or a crash bypasses the relevant handler. These sequences are implemented in `crates/fkst-supervisor/src/main.rs`, `crates/fkst-supervisor/src/process_tree.rs`, `supervise/mod.rs`, and `process_tree.rs`; their normative outcomes are defined only by the matrix in `SPEC.md`. +On `SIGINT` or `SIGTERM`, `fkst-supervisor` sends `SIGTERM` to the event-runtime process group, waits up to two seconds, then sends `SIGKILL` if it has not exited. The event runtime has its own `SIGINT`/`SIGTERM` handlers. A caught signal makes it send `SIGTERM` to every registered Department invocation process group, wait up to two seconds, send `SIGKILL` to survivors, abort runtime tasks, and exit. The supervisor does not launch a replacement. The event runtime is placed in an independent process group, and the supervisor installs no parent-death signal or replacement handoff for itself. `SIGKILL` or a crash bypasses the event-runtime handler, but each Department invocation independently detects loss of its expected event-runtime parent and tears down its registered descendants before exiting. These sequences are implemented in `crates/fkst-supervisor/src/main.rs`, `crates/fkst-supervisor/src/process_tree.rs`, `supervise/mod.rs`, and `process_tree.rs`; their normative outcomes are defined only by the matrix in `SPEC.md`. -Each Department event spawns a `fkst-framework run` child as a new process-group leader. The event runtime registers that group before observing pipes and exit. The child installs a catchable-signal watcher for SDK subprocess groups, but an uncatchable death does not execute it. Department `M.spec.stall_window` is the reliable delivery lease and renewal window, not a child kill deadline. +Each Department event spawns a `fkst-framework run` child as a new process-group leader. The event runtime registers that group before observing pipes and exit. The child installs one watcher for catchable signals and expected-parent loss. Both paths terminate registered SDK process groups before exit. Worktree-backed Codex calls register their separate adoption-worker group and retain a reaper; the worker applies the same watcher and reaper ownership to its `codex exec` group. If the event runtime dies uncatchably, this ownership chain terminates and reaps old Codex work while the unacknowledged reliable delivery remains eligible after lease expiry. Department `M.spec.stall_window` is the reliable delivery lease and renewal window, not a child kill deadline. Graph and bytes bind at different times. The event runtime scans and validates one composed `Config` at startup. It retains the configured runner path, which the OS resolves at every invocation spawn. Each invocation then creates a fresh Lua state and reads its owner-scoped Department and modules. These are the bind-time mechanisms for the graph and later runner/package bytes; `SPEC.md` alone defines the resulting coherence guarantee and replacement outcome. The supervise process boundary follows process-supervision practice: each `fkst-framework run` child has one engine owner that observes exit and reaps it, and that observation is the source for delivery ack or retry. Blocking `wait` and blocking pipe I/O are not allowed on the core async runtime thread; stdout/stderr capture and child exit observation must not starve cron ticks, file-watch ticks, reliable wakes, lease renewal, retry, or dead-letter maintenance. Department spawn admission is bounded before a reliable lease is consumed, so lack of capacity leaves delivery due for a later dispatch pass instead of creating an unobservable leased backlog. -Codex SDK 也把 `codex exec` 放入 process group。`spawn_codex_sync` 与 `spawn_codex` 使用整体 wall-clock `timeout`,默认 3600 秒;只有总运行时间超过 timeout 时才 kill process group,stdout/stderr 输出只被捕获,不延长 timeout。permit 池使用 fcntl lock file,不是内存 semaphore。permit 数来自 registry 的 `codex_permit_slots`:env 或 host `fkst.env` 可覆盖,未设置时默认 `20`。Worktree-backed synchronous calls use the adoption protocol described in §6: a request-derived key, locked intent, detached worker, per-`run_id` lifetime witness, and effect/result records allow a later invocation to observe or recover that same work. This is the only cross-invocation codex-adoption mechanism; it does not transfer Department stdout ownership or make runtime scratch a business fact. +Codex SDK 也把 `codex exec` 放入 process group。`spawn_codex_sync` 与 `spawn_codex` 使用整体 wall-clock `timeout`,默认 3600 秒;只有总运行时间超过 timeout 时才 kill process group,stdout/stderr 输出只被捕获,不延长 timeout。permit 池使用 fcntl lock file,不是内存 semaphore。permit 数来自 registry 的 `codex_permit_slots`:env 或 host `fkst.env` 可覆盖,未设置时默认 `20`。Worktree-backed synchronous calls use the adoption protocol described in §6: a request-derived key, locked intent, owned worker, per-`run_id` lifetime witness, and effect/result records allow another invocation in the same runtime to observe or recover that same work. This is the only cross-invocation Codex adoption mechanism. It does not cross a runtime-generation boundary, transfer Department stdout ownership, or make runtime scratch a business fact. `with_lock(name, fn)` 是跨 pipeline 互斥 primitive。它把校验后的锁名解析到 `/locks//=lock` 并打开,获取 exclusive flock,执行 Lua function,释放 file handle。进程死时 lock 自动释放。 From 53a40eb35f923f7d3ffacbccb3abc7718038f9d9 Mon Sep 17 00:00:00 2001 From: AuricStudio Date: Sat, 22 Aug 2026 05:06:09 +0800 Subject: [PATCH 2/3] Fence process group spawns during shutdown MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Hold a registry guard across process creation and registration so termination closes the spawn gate before taking ownership snapshots. Late groups are killed and remain tracked until their direct children are reaped. Tests: cargo test -p fkst-framework --bin fkst-framework shutdown_fences_an_in_progress_process_group_spawn; cargo test -p fkst-framework --test supervise_smoke supervise_restart_terminates_codex_tree_and_redelivers_expired_lease ⟦AI:FKST⟧ --- SPEC.md | 2 +- crates/fkst-framework/src/external_command.rs | 13 +- crates/fkst-framework/src/process_tree.rs | 124 +++++++++++++----- crates/fkst-framework/src/sdk_codex.rs | 39 +++++- crates/fkst-framework/src/supervise/mod.rs | 45 ++++++- .../fkst-framework/src/supervise/spawner.rs | 6 +- docs/architecture.md | 2 +- 7 files changed, 187 insertions(+), 44 deletions(-) diff --git a/SPEC.md b/SPEC.md index edb5a99..514f5d6 100644 --- a/SPEC.md +++ b/SPEC.md @@ -97,7 +97,7 @@ The outcome matrix is normative. “Catchable” means `SIGINT` or `SIGTERM` rea The engine guarantees no graph reload, no atomic engine/package version pair, and no coherence between the startup-scanned graph and package or runner bytes read by later invocations. Authenticated child output is bound to the event runtime that created its token and pipe; it is never handed to another runtime generation. The engine also does not guarantee that catchable cleanup completes before an external lifecycle manager starts a replacement. Deployment-manager replacement, ordering, canary, and rollback are outside this repository. Mechanism and source locations are described in `docs/architecture.md`. -Each supervised Department invocation receives a private expected-parent PID and monitors it for its full lifetime. A parent change is an ownership loss: the invocation terminates all registered SDK process groups, waits for their direct-child reapers, and exits without accepting further output. A worktree-backed Codex adoption worker repeats this boundary for `codex exec`; the Department retains a worker reaper and registers the worker process group before continuing. This chain makes both catchable shutdown and uncatchable event-runtime death converge on descendant termination and reaping. +Each supervised Department invocation receives a private expected-parent PID and monitors it for its full lifetime. A parent change is an ownership loss: the invocation terminates all registered SDK process groups, waits for their direct-child reapers, and exits without accepting further output. Process-group owners close a spawn gate before termination and hold a spawn guard across process creation and registration; a group whose spawn completes after the gate closes is killed and remains registered until its direct child is reaped. A worktree-backed Codex adoption worker repeats this boundary for `codex exec`; the Department retains a worker reaper and registers the worker process group before continuing. This chain makes both catchable shutdown and uncatchable event-runtime death converge on descendant termination and reaping. Codex adoption records remain current-runtime scratch. A replacement runtime does not scan, migrate, or adopt predecessor runtime directories. Reliable recovery comes from the durable delivery record: after the old owner disappears, its unacknowledged lease expires, `lease_generation` advances, and the replacement starts fresh work under its own runtime layout. Existing lease-generation fencing rejects a stale ACK if cleanup and externally managed replacement briefly overlap. diff --git a/crates/fkst-framework/src/external_command.rs b/crates/fkst-framework/src/external_command.rs index d40d993..bb1329f 100644 --- a/crates/fkst-framework/src/external_command.rs +++ b/crates/fkst-framework/src/external_command.rs @@ -533,13 +533,18 @@ fn run_audited_inner(spec: &CommandSpec) -> Result { fn run_audited_with_timeout(spec: &CommandSpec, timeout: Duration) -> Result { let mut command = build_command(spec); - let mut child = command.spawn()?; - let child_pid = child.id(); - let _registration = if spec.process_group { - Some(crate::process_tree::sdk_process_groups().register(child_pid)) + let spawn_guard = if spec.process_group { + Some( + crate::process_tree::sdk_process_groups() + .begin_spawn() + .context("external command spawn rejected: process shutdown in progress")?, + ) } else { None }; + let mut child = command.spawn()?; + let child_pid = child.id(); + let _registration = spawn_guard.map(|guard| guard.register(child_pid)); let stdout_reader = child .stdout .take() diff --git a/crates/fkst-framework/src/process_tree.rs b/crates/fkst-framework/src/process_tree.rs index b634055..5fb9aca 100644 --- a/crates/fkst-framework/src/process_tree.rs +++ b/crates/fkst-framework/src/process_tree.rs @@ -19,24 +19,43 @@ static SHUTDOWN_SIGNAL: AtomicI32 = AtomicI32::new(0); #[derive(Clone, Default)] pub(crate) struct ProcessGroupRegistry { - pgids: Arc>>, + state: Arc>, +} + +#[derive(Default)] +struct ProcessGroupState { + pgids: BTreeSet, + spawns_in_progress: usize, + terminating: bool, } impl ProcessGroupRegistry { - pub(crate) fn register(&self, pgid: u32) -> ProcessGroupRegistration { - self.pgids - .lock() - .expect("process group registry poisoned") - .insert(pgid); - ProcessGroupRegistration { - pgid, - registry: self.pgids.clone(), + pub(crate) fn begin_spawn(&self) -> Option { + let mut state = self.state.lock().expect("process group registry poisoned"); + if state.terminating { + return None; } + state.spawns_in_progress += 1; + Some(ProcessGroupSpawnGuard { + state: self.state.clone(), + active: true, + }) + } + + fn begin_termination(&self) -> Vec { + let mut state = self.state.lock().expect("process group registry poisoned"); + state.terminating = true; + state.pgids.iter().copied().collect() + } + + fn is_idle(&self) -> bool { + let state = self.state.lock().expect("process group registry poisoned"); + state.pgids.is_empty() && state.spawns_in_progress == 0 } pub(crate) async fn terminate_all(&self, label: &str) { - let pgids = self.snapshot(); - if pgids.is_empty() { + let pgids = self.begin_termination(); + if pgids.is_empty() && self.is_idle() { return; } info!( @@ -49,23 +68,19 @@ impl ProcessGroupRegistry { } let deadline = Instant::now() + TERMINATION_GRACE; while Instant::now() < deadline { - if pgids.iter().all(|pgid| !process_group_exists(*pgid)) { + if self.is_idle() { return; } tokio::time::sleep(POLL_INTERVAL).await; } - for pgid in pgids { + for pgid in self.snapshot() { if process_group_exists(pgid) { send_group_signal(pgid, Signal::SIGKILL, label); } } let deadline = Instant::now() + TERMINATION_GRACE; while Instant::now() < deadline { - if self - .snapshot() - .iter() - .all(|pgid| !process_group_exists(*pgid)) - { + if self.is_idle() { return; } tokio::time::sleep(POLL_INTERVAL).await; @@ -74,8 +89,8 @@ impl ProcessGroupRegistry { } pub(crate) fn terminate_all_blocking(&self, label: &str) { - let pgids = self.snapshot(); - if pgids.is_empty() { + let pgids = self.begin_termination(); + if pgids.is_empty() && self.is_idle() { return; } for pgid in &pgids { @@ -83,23 +98,19 @@ impl ProcessGroupRegistry { } let deadline = Instant::now() + TERMINATION_GRACE; while Instant::now() < deadline { - if pgids.iter().all(|pgid| !process_group_exists(*pgid)) { + if self.is_idle() { return; } std::thread::sleep(POLL_INTERVAL); } - for pgid in pgids { + for pgid in self.snapshot() { if process_group_exists(pgid) { send_group_signal(pgid, Signal::SIGKILL, label); } } let deadline = Instant::now() + TERMINATION_GRACE; while Instant::now() < deadline { - if self - .snapshot() - .iter() - .all(|pgid| !process_group_exists(*pgid)) - { + if self.is_idle() { return; } std::thread::sleep(POLL_INTERVAL); @@ -108,17 +119,19 @@ impl ProcessGroupRegistry { } fn snapshot(&self) -> Vec { - self.pgids + self.state .lock() .expect("process group registry poisoned") + .pgids .iter() .copied() .collect() } fn warn_survivors(&self, label: &str) { - for pgid in self.snapshot() { - if process_group_exists(pgid) { + let state = self.state.lock().expect("process group registry poisoned"); + for pgid in &state.pgids { + if process_group_exists(*pgid) { warn!( pgid = pgid, label = label, @@ -126,19 +139,66 @@ impl ProcessGroupRegistry { ); } } + if state.spawns_in_progress > 0 { + warn!( + spawns = state.spawns_in_progress, + label = label, + "process group spawns did not finish during termination grace" + ); + } + } +} + +pub(crate) struct ProcessGroupSpawnGuard { + state: Arc>, + active: bool, +} + +impl ProcessGroupSpawnGuard { + pub(crate) fn register(mut self, pgid: u32) -> ProcessGroupRegistration { + let terminating = { + let mut state = self.state.lock().expect("process group registry poisoned"); + state.spawns_in_progress = state + .spawns_in_progress + .checked_sub(1) + .expect("process group spawn guard underflow"); + state.pgids.insert(pgid); + state.terminating + }; + self.active = false; + if terminating { + send_group_signal(pgid, Signal::SIGKILL, "process spawned during shutdown"); + } + ProcessGroupRegistration { + pgid, + state: self.state.clone(), + } + } +} + +impl Drop for ProcessGroupSpawnGuard { + fn drop(&mut self) { + if self.active { + let mut state = self.state.lock().expect("process group registry poisoned"); + state.spawns_in_progress = state + .spawns_in_progress + .checked_sub(1) + .expect("process group spawn guard underflow"); + } } } pub(crate) struct ProcessGroupRegistration { pgid: u32, - registry: Arc>>, + state: Arc>, } impl Drop for ProcessGroupRegistration { fn drop(&mut self) { - self.registry + self.state .lock() .expect("process group registry poisoned") + .pgids .remove(&self.pgid); } } diff --git a/crates/fkst-framework/src/sdk_codex.rs b/crates/fkst-framework/src/sdk_codex.rs index 097d0e5..ff14e5b 100644 --- a/crates/fkst-framework/src/sdk_codex.rs +++ b/crates/fkst-framework/src/sdk_codex.rs @@ -1675,11 +1675,20 @@ fn start_adoption_worker( command.process_group(0); inherit_fd_on_exec(&mut command, run_lock_fd); } + let spawn_guard = crate::process_tree::sdk_process_groups() + .begin_spawn() + .ok_or_else(|| { + mlua::Error::external( + "codex adoption worker spawn rejected: process shutdown in progress", + ) + })?; let child = command.spawn().map_err(mlua::Error::external)?; let worker_pid = child.id(); - let registration = crate::process_tree::sdk_process_groups().register(worker_pid); + let registration = spawn_guard.register(worker_pid); let child_slot = Arc::new(Mutex::new(Some(child))); + let registration_slot = Arc::new(Mutex::new(Some(registration))); let reaper_child_slot = Arc::clone(&child_slot); + let reaper_registration_slot = Arc::clone(®istration_slot); if let Err(err) = std::thread::Builder::new() .name("fkst-codex-worker-reaper".to_string()) .spawn(move || { @@ -1694,7 +1703,12 @@ fn start_adoption_worker( ); } } - drop(registration); + drop( + reaper_registration_slot + .lock() + .expect("codex adoption worker registration slot poisoned") + .take(), + ); }) { kill_process_group_by_pid(worker_pid); @@ -1705,6 +1719,12 @@ fn start_adoption_worker( { let _ = child.wait(); } + drop( + registration_slot + .lock() + .expect("codex adoption worker registration slot poisoned") + .take(), + ); return Err(mlua::Error::external(format!( "codex adoption worker reaper spawn failed: {err}" ))); @@ -2866,13 +2886,24 @@ fn run_codex_request_with_permit( request.timeout_seconds ); + let Some(spawn_guard) = crate::process_tree::sdk_process_groups().begin_spawn() else { + return logged_failure( + &request, + &cmd_line, + "shutdown", + "codex spawn rejected: process shutdown in progress".to_string(), + String::new(), + String::new(), + -1, + ); + }; let child = match spawn_codex_child(&request, &mut cmd, &cmd_line) { Ok(child) => child, Err(result) => return *result, }; let child_pid = child.id(); write_adoption_child_pid(&request, child_pid); - let registration = crate::process_tree::sdk_process_groups().register(child_pid); + let registration = spawn_guard.register(child_pid); let prompt = std::mem::take(&mut request.prompt); let (child, stdin_writer) = match spawn_stdin_writer(child, prompt, &request, &cmd_line) { Ok(parts) => parts, @@ -2938,6 +2969,8 @@ fn spawn_stdin_writer( cmd_line: &str, ) -> std::result::Result<(Child, JoinHandle>), Box> { let Some(mut stdin) = child.stdin.take() else { + kill_process_group_by_pid(child.id()); + let _ = child.wait(); return Err(Box::new(logged_failure( request, cmd_line, diff --git a/crates/fkst-framework/src/supervise/mod.rs b/crates/fkst-framework/src/supervise/mod.rs index 120c97d..ab31ea1 100644 --- a/crates/fkst-framework/src/supervise/mod.rs +++ b/crates/fkst-framework/src/supervise/mod.rs @@ -631,6 +631,7 @@ mod tests { use std::os::unix::process::CommandExt; let process_groups = ProcessGroupRegistry::default(); + let spawn_guard = process_groups.begin_spawn().unwrap(); let mut child = Command::new("/bin/sh") .arg("-c") .arg("trap '' TERM; sleep 60") @@ -639,9 +640,8 @@ mod tests { .unwrap(); let child_pid = child.id(); let (registered_tx, registered_rx) = oneshot::channel(); - let process_groups_for_task = process_groups.clone(); let handle = tokio::spawn(async move { - let _registration = process_groups_for_task.register(child_pid); + let _registration = spawn_guard.register(child_pid); let _ = registered_tx.send(()); std::future::pending::<()>().await; }); @@ -660,6 +660,47 @@ mod tests { ); } + #[tokio::test] + async fn shutdown_fences_an_in_progress_process_group_spawn() { + use std::os::unix::process::CommandExt; + + let process_groups = ProcessGroupRegistry::default(); + let spawn_guard = process_groups.begin_spawn().unwrap(); + let process_groups_for_shutdown = process_groups.clone(); + let shutdown = tokio::spawn(async move { + process_groups_for_shutdown + .terminate_all("late department") + .await; + }); + + let deadline = Instant::now() + Duration::from_secs(2); + loop { + if process_groups.begin_spawn().is_none() { + break; + } + assert!( + Instant::now() < deadline, + "process group spawn gate did not close" + ); + tokio::task::yield_now().await; + } + + let mut child = Command::new("/bin/sh") + .arg("-c") + .arg("trap '' TERM; sleep 60") + .process_group(0) + .spawn() + .unwrap(); + let child_pid = child.id(); + let registration = spawn_guard.register(child_pid); + let status = child.wait().unwrap(); + drop(registration); + shutdown.await.unwrap(); + + assert!(!status.success(), "late process group was not terminated"); + assert!(process_groups.begin_spawn().is_none()); + } + fn wait_for_child_exit(child: &mut std::process::Child, timeout: Duration) -> bool { let deadline = Instant::now() + timeout; loop { diff --git a/crates/fkst-framework/src/supervise/spawner.rs b/crates/fkst-framework/src/supervise/spawner.rs index 0122f33..c4de8a0 100644 --- a/crates/fkst-framework/src/supervise/spawner.rs +++ b/crates/fkst-framework/src/supervise/spawner.rs @@ -178,6 +178,10 @@ pub async fn spawn_framework_with_stdout_observer( // tokio::process exposes `process_group(0)` to call setpgid(0,0); equivalent for our purposes. cmd.process_group(0); + let spawn_guard = process_groups.begin_spawn().ok_or_else(|| { + log.write_line("SPAWN_ERROR=event runtime shutdown in progress"); + anyhow::anyhow!("framework spawn rejected: event runtime shutdown in progress") + })?; let (child, spawn_return_ms) = match cmd.spawn() { Ok(child) => (child, start.elapsed().as_millis()), Err(err) => { @@ -190,7 +194,7 @@ pub async fn spawn_framework_with_stdout_observer( .ok_or_else(|| anyhow::anyhow!("no pid after spawn"))?; log.write_line(&format!("PID={pid}")); info!(pid = pid, lua = %lua_path.display(), "framework spawned"); - let registration = process_groups.register(pid); + let registration = spawn_guard.register(pid); wait_for_framework_child( child, diff --git a/docs/architecture.md b/docs/architecture.md index 4af8246..bcf8dc8 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -350,7 +350,7 @@ The engine supports two entry topologies. `fkst-supervisor` is the shipped proce On `SIGINT` or `SIGTERM`, `fkst-supervisor` sends `SIGTERM` to the event-runtime process group, waits up to two seconds, then sends `SIGKILL` if it has not exited. The event runtime has its own `SIGINT`/`SIGTERM` handlers. A caught signal makes it send `SIGTERM` to every registered Department invocation process group, wait up to two seconds, send `SIGKILL` to survivors, abort runtime tasks, and exit. The supervisor does not launch a replacement. The event runtime is placed in an independent process group, and the supervisor installs no parent-death signal or replacement handoff for itself. `SIGKILL` or a crash bypasses the event-runtime handler, but each Department invocation independently detects loss of its expected event-runtime parent and tears down its registered descendants before exiting. These sequences are implemented in `crates/fkst-supervisor/src/main.rs`, `crates/fkst-supervisor/src/process_tree.rs`, `supervise/mod.rs`, and `process_tree.rs`; their normative outcomes are defined only by the matrix in `SPEC.md`. -Each Department event spawns a `fkst-framework run` child as a new process-group leader. The event runtime registers that group before observing pipes and exit. The child installs one watcher for catchable signals and expected-parent loss. Both paths terminate registered SDK process groups before exit. Worktree-backed Codex calls register their separate adoption-worker group and retain a reaper; the worker applies the same watcher and reaper ownership to its `codex exec` group. If the event runtime dies uncatchably, this ownership chain terminates and reaps old Codex work while the unacknowledged reliable delivery remains eligible after lease expiry. Department `M.spec.stall_window` is the reliable delivery lease and renewal window, not a child kill deadline. +Each Department event spawns a `fkst-framework run` child as a new process-group leader. The event runtime registers that group before observing pipes and exit. The child installs one watcher for catchable signals and expected-parent loss. Both paths close the process-group spawn gate, terminate registered SDK process groups, and wait for their direct-child reapers before exit. A spawn guard spans process creation and registration, so a group that finishes spawning concurrently with shutdown is killed and remains tracked until reaped. Worktree-backed Codex calls register their separate adoption-worker group and retain a reaper; the worker applies the same watcher, spawn gate, and reaper ownership to its `codex exec` group. If the event runtime dies uncatchably, this ownership chain terminates and reaps old Codex work while the unacknowledged reliable delivery remains eligible after lease expiry. Department `M.spec.stall_window` is the reliable delivery lease and renewal window, not a child kill deadline. Graph and bytes bind at different times. The event runtime scans and validates one composed `Config` at startup. It retains the configured runner path, which the OS resolves at every invocation spawn. Each invocation then creates a fresh Lua state and reads its owner-scoped Department and modules. These are the bind-time mechanisms for the graph and later runner/package bytes; `SPEC.md` alone defines the resulting coherence guarantee and replacement outcome. From 537063bfe45bb93e2ed72aeb6f6c4ca088af8822 Mon Sep 17 00:00:00 2001 From: AuricStudio Date: Sat, 22 Aug 2026 05:14:23 +0800 Subject: [PATCH 3/3] fkst: implementation result v1 5f382efc1cb40b284ff010ddc52421d1b9bd7511750aa2e0129ad25ae267d7f3