Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand All @@ -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. 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.

## 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`。
Expand Down
13 changes: 9 additions & 4 deletions crates/fkst-framework/src/external_command.rs
Original file line number Diff line number Diff line change
Expand Up @@ -533,13 +533,18 @@ fn run_audited_inner(spec: &CommandSpec) -> Result<AuditedOutput> {

fn run_audited_with_timeout(spec: &CommandSpec, timeout: Duration) -> Result<AuditedOutput> {
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()
Expand Down
164 changes: 139 additions & 25 deletions crates/fkst-framework/src/process_tree.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand All @@ -10,31 +11,51 @@ 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<ProcessGroupRegistry> = OnceLock::new();
static SIGNAL_WATCH_INSTALLED: OnceLock<()> = OnceLock::new();
static SHUTDOWN_REQUESTED: AtomicBool = AtomicBool::new(false);
static SHUTDOWN_SIGNAL: AtomicI32 = AtomicI32::new(0);

#[derive(Clone, Default)]
pub(crate) struct ProcessGroupRegistry {
pgids: Arc<Mutex<BTreeSet<u32>>>,
state: Arc<Mutex<ProcessGroupState>>,
}

#[derive(Default)]
struct ProcessGroupState {
pgids: BTreeSet<u32>,
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<ProcessGroupSpawnGuard> {
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<u32> {
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!(
Expand All @@ -47,60 +68,137 @@ 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.is_idle() {
return;
}
tokio::time::sleep(POLL_INTERVAL).await;
}
self.warn_survivors(label);
}

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 {
send_group_signal(*pgid, Signal::SIGTERM, label);
}
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.is_idle() {
return;
}
std::thread::sleep(POLL_INTERVAL);
}
self.warn_survivors(label);
}

fn snapshot(&self) -> Vec<u32> {
self.pgids
self.state
.lock()
.expect("process group registry poisoned")
.pgids
.iter()
.copied()
.collect()
}

fn warn_survivors(&self, label: &str) {
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,
"process group survived SIGKILL grace"
);
}
}
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<Mutex<ProcessGroupState>>,
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<Mutex<BTreeSet<u32>>>,
state: Arc<Mutex<ProcessGroupState>>,
}

impl Drop for ProcessGroupRegistration {
fn drop(&mut self) {
self.registry
self.state
.lock()
.expect("process group registry poisoned")
.pgids
.remove(&self.pgid);
}
}
Expand Down Expand Up @@ -133,23 +231,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<u32> {
std::env::var(SUPERVISOR_PID_ENV)
.ok()
.and_then(|value| value.parse::<u32>().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) {
Expand Down
Loading
Loading