Skip to main content

soma_provider_adapters/python/
supervisor.rs

1//! Serial persistent Python worker supervision.
2
3use std::{
4    collections::VecDeque,
5    path::PathBuf,
6    sync::{
7        Arc, Mutex as StdMutex,
8        atomic::{AtomicBool, AtomicU32, AtomicU64, AtomicUsize, Ordering},
9    },
10    time::{Duration, Instant, SystemTime, UNIX_EPOCH},
11};
12
13use serde::Serialize;
14use serde_json::Value;
15use tokio::{
16    io::{AsyncRead, AsyncReadExt, AsyncWrite},
17    net::TcpListener,
18    process::{Child, Command},
19    sync::{Mutex, Notify, OwnedSemaphorePermit},
20    task::JoinHandle,
21    time::timeout,
22};
23
24use crate::{
25    python::{
26        PythonInterpreter,
27        containment::{BrokeredLaunch, CgroupGuard},
28        host::{PythonExecutionProfile, PythonHostAuditEvent, PythonHostBroker},
29    },
30    python_protocol::{
31        PythonInvocationRequest, PythonInvocationState, PythonRequestState, PythonRunnerFeature,
32        PythonRunnerHostCall, PythonRunnerHostMessage, PythonRunnerHostRequest,
33        PythonRunnerProtocolVersion, PythonRunnerReply, PythonRunnerWorkerMessage,
34        negotiate_runner_features,
35    },
36    sidecar::{resolve_sidecar_command, sidecar_base_env},
37};
38
39#[cfg(test)]
40#[path = "supervisor_cancel_tests.rs"]
41mod cancel_tests;
42#[path = "supervisor/cancellation.rs"]
43mod cancellation;
44#[path = "supervisor_frames.rs"]
45mod frames;
46#[path = "supervisor_logs.rs"]
47mod logs;
48#[path = "supervisor_state.rs"]
49mod state;
50#[path = "supervisor/status.rs"]
51mod status;
52#[path = "supervisor_termination.rs"]
53mod termination;
54use frames::{host_call_invocation_id, host_call_request_id, read_frame, write_frame};
55use logs::drain_stderr;
56#[cfg(test)]
57use state::worker_budget_keys_are_live;
58use state::{
59    BusyGuard, candidate_budget, invalid_output, map_worker_error, protocol_error, start_error,
60    worker_budget,
61};
62pub use state::{PythonInvocationOptions, PythonWorkerLogEntry};
63use termination::{JobGuard, ProcessTreeStartupGuard, terminate_process_tree, terminate_worker};
64
65/// Product-neutral limits for one persistent worker.
66#[derive(Debug, Clone, PartialEq, Eq)]
67pub struct PythonSupervisorConfig {
68    pub startup_timeout: Duration,
69    pub request_timeout: Duration,
70    pub shutdown_grace: Duration,
71    pub max_restarts: u32,
72    pub restart_window: Duration,
73    pub restart_backoff: Duration,
74    pub max_stderr_bytes: usize,
75    pub max_pending_bytes: usize,
76    pub max_workers: usize,
77    pub max_candidate_starts: usize,
78    pub execution_profile: PythonExecutionProfile,
79}
80
81impl Default for PythonSupervisorConfig {
82    fn default() -> Self {
83        Self {
84            startup_timeout: Duration::from_secs(10),
85            request_timeout: Duration::from_secs(10),
86            shutdown_grace: Duration::from_secs(2),
87            max_restarts: 3,
88            restart_window: Duration::from_secs(60),
89            restart_backoff: Duration::from_millis(250),
90            max_stderr_bytes: 64 * 1024,
91            max_pending_bytes: 512 * 1024,
92            max_workers: 32,
93            max_candidate_starts: 4,
94            execution_profile: PythonExecutionProfile::Trusted,
95        }
96    }
97}
98
99/// Immutable identity of a worker process.
100#[derive(Debug, Clone, PartialEq, Eq)]
101pub struct PythonWorkerIdentity {
102    pub path: PathBuf,
103    pub generation_id: String,
104    pub worker_group: String,
105    pub source_digest: String,
106    pub catalog_fingerprint: String,
107}
108
109/// Operator-facing state for one persistent Python provider worker.
110#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
111pub struct PythonWorkerStatus {
112    pub provider_source: PathBuf,
113    pub generation_id: String,
114    pub running: bool,
115    pub accepting: bool,
116    pub busy: bool,
117    pub quarantined: bool,
118    pub restart_count: usize,
119    pub logs: Vec<PythonWorkerLogEntry>,
120    pub execution_profile: PythonExecutionProfile,
121    pub host_audit: Vec<PythonHostAuditEvent>,
122}
123
124#[derive(Default)]
125struct WorkerLogBuffer {
126    entries: VecDeque<PythonWorkerLogEntry>,
127    retained_bytes: usize,
128    next_sequence: u64,
129}
130
131#[derive(Debug)]
132pub struct PythonSupervisorError {
133    code: &'static str,
134    message: String,
135}
136
137impl PythonSupervisorError {
138    pub fn new(code: &'static str, message: impl Into<String>) -> Self {
139        Self {
140            code,
141            message: message.into(),
142        }
143    }
144
145    #[must_use]
146    pub const fn code(&self) -> &'static str {
147        self.code
148    }
149}
150
151pub(super) struct Worker {
152    pub(super) child: Child,
153    pub(super) child_pid: Option<u32>,
154    _job_guard: JobGuard,
155    _cgroup_guard: CgroupGuard,
156    _worker_permit: OwnedSemaphorePermit,
157    stdin: Box<dyn AsyncWrite + Unpin + Send>,
158    stdout: Box<dyn AsyncRead + Unpin + Send>,
159    pub(super) stderr_task: JoinHandle<()>,
160    described: bool,
161    provider_path: PathBuf,
162}
163
164/// One persistent worker per Python provider. Invocations are deliberately
165/// serial; callers receive a stable busy error instead of entering a queue.
166pub struct PythonWorkerSupervisor {
167    identity: PythonWorkerIdentity,
168    interpreter: PythonInterpreter,
169    config: PythonSupervisorConfig,
170    worker: Mutex<Option<Worker>>,
171    busy: AtomicBool,
172    accepting: AtomicBool,
173    dispatch_leases: AtomicUsize,
174    leases_released: Notify,
175    request_id: AtomicU64,
176    restarts: StdMutex<VecDeque<Instant>>,
177    quarantined: AtomicBool,
178    started_once: AtomicBool,
179    discard_worker: AtomicBool,
180    cancel_epoch: AtomicU64,
181    active_pid: Arc<AtomicU32>,
182    logs: Arc<StdMutex<WorkerLogBuffer>>,
183    host: Arc<PythonHostBroker>,
184}
185
186impl PythonWorkerSupervisor {
187    #[must_use]
188    pub fn new(
189        identity: PythonWorkerIdentity,
190        interpreter: PythonInterpreter,
191        config: PythonSupervisorConfig,
192    ) -> Arc<Self> {
193        Self::new_with_capabilities(
194            identity,
195            interpreter,
196            config,
197            &soma_provider_core::HostCapabilities::default(),
198        )
199    }
200
201    #[must_use]
202    pub fn new_with_capabilities(
203        identity: PythonWorkerIdentity,
204        interpreter: PythonInterpreter,
205        config: PythonSupervisorConfig,
206        capabilities: &soma_provider_core::HostCapabilities,
207    ) -> Arc<Self> {
208        let host = PythonHostBroker::new(
209            config.execution_profile,
210            capabilities,
211            Arc::new(AtomicBool::new(false)),
212        );
213        Arc::new(Self {
214            identity,
215            interpreter,
216            config,
217            worker: Mutex::new(None),
218            busy: AtomicBool::new(false),
219            accepting: AtomicBool::new(true),
220            dispatch_leases: AtomicUsize::new(0),
221            leases_released: Notify::new(),
222            request_id: AtomicU64::new(1),
223            restarts: StdMutex::new(VecDeque::new()),
224            quarantined: AtomicBool::new(false),
225            started_once: AtomicBool::new(false),
226            discard_worker: AtomicBool::new(false),
227            cancel_epoch: AtomicU64::new(0),
228            active_pid: Arc::new(AtomicU32::new(0)),
229            logs: Arc::new(StdMutex::new(WorkerLogBuffer::default())),
230            host,
231        })
232    }
233
234    /// Clears a crash-loop quarantine after an explicit operator action.
235    pub async fn reset_quarantine(&self) {
236        self.quarantined.store(false, Ordering::Release);
237        self.started_once.store(false, Ordering::Release);
238        self.restarts
239            .lock()
240            .expect("Python worker restart lock should not be poisoned")
241            .clear();
242        self.discard_worker.store(true, Ordering::Release);
243    }
244
245    /// Stops new work while allowing an invocation that already owns the
246    /// worker to complete on this generation.
247    pub fn deactivate(&self) {
248        self.accepting.store(false, Ordering::Release);
249    }
250
251    /// Re-enables a retained generation during atomic rollback.
252    pub fn activate(&self) {
253        self.accepting.store(true, Ordering::Release);
254    }
255
256    /// Reserves a call while a registry generation is still active.
257    pub fn acquire_dispatch(&self) -> bool {
258        if !self.accepting.load(Ordering::Acquire) {
259            return false;
260        }
261        self.dispatch_leases.fetch_add(1, Ordering::AcqRel);
262        if self.accepting.load(Ordering::Acquire) {
263            return true;
264        }
265        self.release_dispatch();
266        false
267    }
268
269    /// Releases a call reservation after the routed invocation completes.
270    pub fn release_dispatch(&self) {
271        let previous = self.dispatch_leases.fetch_sub(1, Ordering::AcqRel);
272        debug_assert!(previous > 0, "dispatch lease count must not underflow");
273        if previous == 1 {
274            self.leases_released.notify_waiters();
275        }
276    }
277
278    /// Permanently parks this generation after all already-routed calls drain.
279    pub async fn suspend(&self) {
280        self.deactivate();
281        while self.dispatch_leases.load(Ordering::Acquire) != 0 {
282            let notified = self.leases_released.notified();
283            if self.dispatch_leases.load(Ordering::Acquire) == 0 {
284                break;
285            }
286            notified.await;
287        }
288        self.shutdown().await;
289    }
290
291    pub async fn preflight(&self) -> Result<Value, PythonSupervisorError> {
292        let _candidate_permit = candidate_budget(self.config.max_candidate_starts)
293            .acquire_owned()
294            .await
295            .map_err(|_| start_error())?;
296        let mut worker = self.worker.lock().await;
297        self.ensure_worker(&mut worker).await
298    }
299
300    pub async fn invoke(
301        &self,
302        provider: &str,
303        action: &str,
304        arguments: Value,
305        surface: soma_provider_core::ProviderSurface,
306        snapshot_id: &str,
307        timeout_override: Duration,
308    ) -> Result<Value, PythonSupervisorError> {
309        let context = soma_provider_core::ProviderInvocationContext::default();
310        self.invoke_with_context(
311            provider,
312            action,
313            arguments,
314            PythonInvocationOptions {
315                surface,
316                snapshot_id,
317                timeout: timeout_override,
318                context: &context,
319            },
320        )
321        .await
322    }
323
324    pub async fn invoke_with_context(
325        &self,
326        provider: &str,
327        action: &str,
328        arguments: Value,
329        options: PythonInvocationOptions<'_>,
330    ) -> Result<Value, PythonSupervisorError> {
331        if self.config.execution_profile == PythonExecutionProfile::Disabled {
332            return Err(PythonSupervisorError::new(
333                "python_execution_disabled",
334                "Python provider execution is disabled by policy",
335            ));
336        }
337        if !self.accepting.load(Ordering::Acquire)
338            && self.dispatch_leases.load(Ordering::Acquire) == 0
339        {
340            return Err(PythonSupervisorError::new(
341                "python_worker_draining",
342                "Python provider generation is draining",
343            ));
344        }
345        let cancel_epoch = self.cancel_epoch.load(Ordering::Acquire);
346        if self
347            .busy
348            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
349            .is_err()
350        {
351            return Err(PythonSupervisorError::new(
352                "python_provider_busy",
353                "Python provider is busy",
354            ));
355        }
356        let mut busy = BusyGuard::new(&self.busy, &self.discard_worker);
357        let result = self
358            .invoke_inner((provider, action), arguments, options, cancel_epoch)
359            .await;
360        busy.complete();
361        result
362    }
363
364    async fn invoke_inner(
365        &self,
366        target: (&str, &str),
367        arguments: Value,
368        options: PythonInvocationOptions<'_>,
369        cancel_epoch: u64,
370    ) -> Result<Value, PythonSupervisorError> {
371        let (provider, action) = target;
372        let PythonInvocationOptions {
373            surface,
374            snapshot_id,
375            timeout: timeout_override,
376            context,
377        } = options;
378        let encoded_len = serde_json::to_vec(&arguments)
379            .map_err(|_| invalid_output())?
380            .len();
381        if encoded_len > self.config.max_pending_bytes {
382            return Err(PythonSupervisorError::new(
383                "python_input_too_large",
384                "Python provider input exceeds the persistent runner limit",
385            ));
386        }
387        let mut slot = self.worker.lock().await;
388        if self.discard_worker.swap(false, Ordering::AcqRel) {
389            self.active_pid.store(0, Ordering::Release);
390            terminate_worker(slot.take()).await;
391        }
392        self.ensure_worker(&mut slot).await?;
393        if self.cancel_epoch.load(Ordering::Acquire) != cancel_epoch {
394            self.active_pid.store(0, Ordering::Release);
395            terminate_worker(slot.take()).await;
396            return Err(PythonSupervisorError::new(
397                "python_provider_cancelled",
398                "Python provider invocation was cancelled",
399            ));
400        }
401        let worker = slot.as_mut().expect("worker was ensured");
402        let request_id = self.next_request_id();
403        let invocation_id = if context.request_id.is_empty() {
404            format!("{}-{request_id}", self.identity.generation_id)
405        } else {
406            context.request_id.clone()
407        };
408        self.host.begin_invocation();
409        let wait = self.config.request_timeout.min(timeout_override);
410        let deadline = SystemTime::now()
411            .duration_since(UNIX_EPOCH)
412            .unwrap_or_default()
413            .saturating_add(wait)
414            .as_millis()
415            .min(u128::from(u64::MAX)) as u64;
416        let actor_context =
417            context
418                .actor_id
419                .as_ref()
420                .map(|actor_id| crate::python_protocol::PythonActorContext {
421                    actor_id: actor_id.clone(),
422                    scopes: context.actor_scopes.clone(),
423                });
424        let request = PythonRunnerHostMessage::Request {
425            request: PythonRunnerHostRequest::Invoke {
426                request_id,
427                invocation: Box::new(PythonInvocationRequest {
428                    invocation_id: invocation_id.clone(),
429                    request_id: invocation_id.clone(),
430                    provider: provider.to_owned(),
431                    action: action.to_owned(),
432                    arguments,
433                    surface,
434                    snapshot_id: snapshot_id.to_owned(),
435                    deadline_unix_ms: deadline,
436                    trace: context.traceparent.as_ref().map(|traceparent| {
437                        crate::python_protocol::PythonTraceContext {
438                            traceparent: traceparent.clone(),
439                            tracestate: context.tracestate.clone(),
440                        }
441                    }),
442                    actor: actor_context.clone(),
443                    cancellation_token_id: format!("cancel-{request_id}"),
444                    generation_id: self.identity.generation_id.clone(),
445                }),
446            },
447        };
448        let exchange = async {
449            write_frame(&mut worker.stdin, &request).await?;
450            let mut state = PythonRequestState::Written;
451            loop {
452                match read_frame::<PythonRunnerWorkerMessage, _>(&mut worker.stdout).await? {
453                    PythonRunnerWorkerMessage::Reply {
454                        reply:
455                            PythonRunnerReply::Accepted {
456                                request_id: actual,
457                                invocation_id: actual_invocation_id,
458                                state: PythonInvocationState::Accepted,
459                            },
460                    } if actual == request_id
461                        && actual_invocation_id == invocation_id
462                        && state == PythonRequestState::Written =>
463                    {
464                        state = PythonRequestState::Accepted;
465                        continue;
466                    }
467                    PythonRunnerWorkerMessage::Reply {
468                        reply:
469                            PythonRunnerReply::Ok {
470                                request_id: actual,
471                                result,
472                            },
473                    } if actual == request_id && state == PythonRequestState::Accepted => {
474                        return Ok(result);
475                    }
476                    PythonRunnerWorkerMessage::Reply {
477                        reply:
478                            PythonRunnerReply::Error {
479                                request_id: actual,
480                                error,
481                            },
482                    } if actual == request_id && state == PythonRequestState::Accepted => {
483                        return Err(map_worker_error(error.code));
484                    }
485                    PythonRunnerWorkerMessage::HostCall { call } => {
486                        let host_request_id = host_call_request_id(&call);
487                        if host_call_invocation_id(&call) != invocation_id {
488                            return Err(protocol_error());
489                        }
490                        let (result, error) =
491                            match self.host.execute(&call, actor_context.as_ref()).await {
492                                Ok(result) => {
493                                    if let PythonRunnerHostCall::Progress {
494                                        current,
495                                        total,
496                                        message,
497                                        ..
498                                    } = &call
499                                    {
500                                        context.progress.report(
501                                            *current,
502                                            *total,
503                                            message.as_deref(),
504                                        );
505                                    }
506                                    (Some(result), None)
507                                }
508                                Err(error) => (None, Some(*error)),
509                            };
510                        write_frame(
511                            &mut worker.stdin,
512                            &PythonRunnerHostMessage::HostReply {
513                                request_id: host_request_id,
514                                result,
515                                error,
516                            },
517                        )
518                        .await?;
519                    }
520                    _ => return Err(protocol_error()),
521                }
522            }
523        };
524        match timeout(wait, exchange).await {
525            Ok(Ok(_)) if self.cancel_epoch.load(Ordering::Acquire) != cancel_epoch => {
526                self.active_pid.store(0, Ordering::Release);
527                terminate_worker(slot.take()).await;
528                Err(PythonSupervisorError::new(
529                    "python_provider_cancelled",
530                    "Python provider invocation was cancelled",
531                ))
532            }
533            Ok(Ok(result)) => Ok(result),
534            Ok(Err(_)) if self.cancel_epoch.load(Ordering::Acquire) != cancel_epoch => {
535                self.active_pid.store(0, Ordering::Release);
536                terminate_worker(slot.take()).await;
537                Err(PythonSupervisorError::new(
538                    "python_provider_cancelled",
539                    "Python provider invocation was cancelled",
540                ))
541            }
542            Ok(Err(error))
543                if error.code() == "python_provider_failed"
544                    || error.code() == "python_provider_cancelled"
545                    || error.code() == "python_output_too_large" =>
546            {
547                Err(error)
548            }
549            Ok(Err(error)) => {
550                self.active_pid.store(0, Ordering::Release);
551                terminate_worker(slot.take()).await;
552                Err(error)
553            }
554            Err(_) => {
555                self.active_pid.store(0, Ordering::Release);
556                terminate_worker(slot.take()).await;
557                Err(PythonSupervisorError::new(
558                    "python_provider_timeout",
559                    "Python provider exceeded its timeout",
560                ))
561            }
562        }
563    }
564
565    async fn ensure_worker(
566        &self,
567        slot: &mut Option<Worker>,
568    ) -> Result<Value, PythonSupervisorError> {
569        if self.quarantined.load(Ordering::Acquire) {
570            return Err(PythonSupervisorError::new(
571                "python_provider_quarantined",
572                "Python provider is quarantined after repeated worker failures",
573            ));
574        }
575        let worker_exited = slot.as_mut().is_some_and(|worker| {
576            worker
577                .child
578                .try_wait()
579                .map_or(true, |status| status.is_some())
580        });
581        if worker_exited {
582            self.active_pid.store(0, Ordering::Release);
583            terminate_worker(slot.take()).await;
584        }
585        if slot.is_none() {
586            self.verify_source_digest()?;
587            let restarting = self.started_once.swap(true, Ordering::AcqRel);
588            if restarting {
589                self.record_restart()?;
590            }
591            if restarting && !self.config.restart_backoff.is_zero() {
592                tokio::time::sleep(self.config.restart_backoff).await;
593            }
594            *slot = Some(self.spawn_worker().await?);
595        }
596        let worker = slot.as_mut().expect("worker exists");
597        if worker.described {
598            return Ok(Value::Null);
599        }
600        let request_id = self.next_request_id();
601        let describe = PythonRunnerHostMessage::Request {
602            request: PythonRunnerHostRequest::Describe {
603                request_id,
604                path: worker.provider_path.clone(),
605                generation_id: self.identity.generation_id.clone(),
606            },
607        };
608        write_frame(&mut worker.stdin, &describe).await?;
609        let described = match timeout(
610            self.config.startup_timeout,
611            read_frame::<PythonRunnerWorkerMessage, _>(&mut worker.stdout),
612        )
613        .await
614        {
615            Ok(Ok(PythonRunnerWorkerMessage::Reply {
616                reply:
617                    PythonRunnerReply::Ok {
618                        request_id: actual,
619                        result,
620                    },
621            })) if actual == request_id => result,
622            _ => {
623                self.active_pid.store(0, Ordering::Release);
624                terminate_worker(slot.take()).await;
625                return Err(PythonSupervisorError::new(
626                    "python_worker_start_failed",
627                    "Python worker failed provider preflight",
628                ));
629            }
630        };
631        if let Err(error) = self.verify_source_digest() {
632            self.active_pid.store(0, Ordering::Release);
633            terminate_worker(slot.take()).await;
634            return Err(error);
635        }
636        let manifest = match soma_provider_core::validate_provider_manifest_value(&described) {
637            Ok(manifest) => manifest,
638            Err(_) => {
639                self.active_pid.store(0, Ordering::Release);
640                terminate_worker(slot.take()).await;
641                return Err(protocol_error());
642            }
643        };
644        let actual_catalog =
645            super::python_catalog_fingerprint(&manifest).map_err(|_| protocol_error())?;
646        if !self.identity.catalog_fingerprint.is_empty()
647            && actual_catalog != self.identity.catalog_fingerprint
648        {
649            self.active_pid.store(0, Ordering::Release);
650            terminate_worker(slot.take()).await;
651            return Err(PythonSupervisorError::new(
652                "python_catalog_changed",
653                "Python provider catalog changed during worker activation",
654            ));
655        }
656        if self.health_check(worker).await.is_err() {
657            self.active_pid.store(0, Ordering::Release);
658            terminate_worker(slot.take()).await;
659            return Err(start_error());
660        }
661        worker.described = true;
662        Ok(described)
663    }
664
665    async fn health_check(&self, worker: &mut Worker) -> Result<(), PythonSupervisorError> {
666        let request_id = self.next_request_id();
667        write_frame(
668            &mut worker.stdin,
669            &PythonRunnerHostMessage::Request {
670                request: PythonRunnerHostRequest::Health { request_id },
671            },
672        )
673        .await?;
674        match timeout(
675            self.config.startup_timeout,
676            read_frame::<PythonRunnerWorkerMessage, _>(&mut worker.stdout),
677        )
678        .await
679        {
680            Ok(Ok(PythonRunnerWorkerMessage::Reply {
681                reply:
682                    PythonRunnerReply::Health {
683                        request_id: actual,
684                        health: crate::python_protocol::PythonWorkerHealth::Ready,
685                        generation_id,
686                    },
687            })) if actual == request_id && generation_id == self.identity.generation_id => Ok(()),
688            _ => Err(start_error()),
689        }
690    }
691
692    async fn spawn_worker(&self) -> Result<Worker, PythonSupervisorError> {
693        let command = match &self.interpreter {
694            PythonInterpreter::Ambient => crate::python::default_python_command().to_owned(),
695            PythonInterpreter::Prepared(path) => path.to_string_lossy().into_owned(),
696        };
697        let listener = TcpListener::bind("127.0.0.1:0")
698            .await
699            .map_err(|_| start_error())?;
700        let address = listener.local_addr().map_err(|_| start_error())?;
701        let brokered = self.config.execution_profile == PythonExecutionProfile::Brokered;
702        let (mut process, brokered_launch, worker_provider_path) = if brokered {
703            let (process, launch, worker_provider_path) = BrokeredLaunch::prepare(
704                resolve_sidecar_command(&command)
705                    .to_str()
706                    .ok_or_else(start_error)?,
707                &self.identity.path,
708            )?;
709            (process, Some(launch), worker_provider_path)
710        } else {
711            (
712                Command::new(resolve_sidecar_command(&command)),
713                None,
714                self.identity.path.clone(),
715            )
716        };
717        let token = {
718            let mut token = [0_u8; 32];
719            getrandom::fill(&mut token).map_err(|_| start_error())?;
720            token
721                .iter()
722                .map(|byte| format!("{byte:02x}"))
723                .collect::<String>()
724        };
725        if !brokered {
726            process.args(["-I", "-m", "soma_provider.runner"]);
727        }
728        process
729            .kill_on_drop(true)
730            .env_clear()
731            .env("SOMA_PYTHON_RUNNER_TOKEN", &token)
732            .stdin(std::process::Stdio::null())
733            .stdout(std::process::Stdio::null())
734            .stderr(std::process::Stdio::piped());
735        if brokered {
736            process.env("SOMA_PYTHON_RUNNER_ADDR", "unix:/run/soma/control.sock");
737        } else {
738            process.env("SOMA_PYTHON_RUNNER_ADDR", address.to_string());
739        }
740        #[cfg(unix)]
741        process.process_group(0);
742        for (key, value) in sidecar_base_env() {
743            process.env(key, value);
744        }
745        // Atomic publication may briefly own both the active and replacement
746        // generations. `max_workers` remains the per-generation bound.
747        let worker_permit = timeout(
748            self.config.startup_timeout,
749            worker_budget(&self.identity.worker_group, self.config.max_workers).acquire_owned(),
750        )
751        .await
752        .map_err(|_| start_error())?
753        .map_err(|_| start_error())?;
754        if let Some(launch) = &brokered_launch {
755            launch.retain_until_spawn();
756        }
757        let mut child = crate::sidecar::spawn_retrying_busy_image(&mut process)
758            .await
759            .map_err(|_| {
760                PythonSupervisorError::new(
761                    "python_worker_start_failed",
762                    "Python worker could not be started",
763                )
764            })?;
765        let child_pid = child.id();
766        let mut startup_guard = ProcessTreeStartupGuard::new(child_pid);
767        let stderr = child.stderr.take().ok_or_else(protocol_error)?;
768        let stderr_task = tokio::spawn(drain_stderr(
769            stderr,
770            self.logs.clone(),
771            self.config.max_stderr_bytes,
772        ));
773        let cgroup_guard = match CgroupGuard::attach(child_pid, self.config.execution_profile) {
774            Ok(guard) => guard,
775            Err(error) => {
776                terminate_process_tree(child_pid);
777                let _ = child.kill().await;
778                return Err(error);
779            }
780        };
781        let job_guard = JobGuard::new(&child)?;
782        let (mut stdout, stdin): (
783            Box<dyn AsyncRead + Unpin + Send>,
784            Box<dyn AsyncWrite + Unpin + Send>,
785        ) = if let Some(launch) = &brokered_launch {
786            #[cfg(target_os = "linux")]
787            {
788                let stream = match timeout(self.config.startup_timeout, launch.accept()).await {
789                    Ok(Ok(stream)) => stream,
790                    result => {
791                        terminate_process_tree(child_pid);
792                        let _ = child.kill().await;
793                        let _ = stderr_task.await;
794                        let detail = self
795                            .logs
796                            .lock()
797                            .expect("Python worker log lock should not be poisoned")
798                            .entries
799                            .back()
800                            .map(|entry| entry.message.clone())
801                            .unwrap_or_else(|| match result {
802                                Ok(Err(error)) => error.to_string(),
803                                Err(_) => "control connection timed out".to_owned(),
804                                Ok(Ok(_)) => unreachable!(),
805                            });
806                        return Err(PythonSupervisorError::new(
807                            "python_worker_start_failed",
808                            format!(
809                                "Python worker could not connect to its control socket: {detail}"
810                            ),
811                        ));
812                    }
813                };
814                let (read, write) = stream.into_split();
815                (Box::new(read), Box::new(write))
816            }
817            #[cfg(not(target_os = "linux"))]
818            {
819                let _ = launch;
820                return Err(start_error());
821            }
822        } else {
823            let (stream, _) = timeout(self.config.startup_timeout, listener.accept())
824                .await
825                .map_err(|_| start_error())?
826                .map_err(|_| start_error())?;
827            let (read, write) = stream.into_split();
828            (Box::new(read), Box::new(write))
829        };
830        drop(brokered_launch);
831        let mut actual_token = vec![0_u8; token.len()];
832        timeout(
833            self.config.startup_timeout,
834            stdout.read_exact(&mut actual_token),
835        )
836        .await
837        .map_err(|_| start_error())?
838        .map_err(|_| start_error())?;
839        if actual_token != token.as_bytes() {
840            terminate_process_tree(child_pid);
841            let _ = child.kill().await;
842            return Err(protocol_error());
843        }
844        let mut worker = Worker {
845            child,
846            child_pid,
847            _job_guard: job_guard,
848            _cgroup_guard: cgroup_guard,
849            _worker_permit: worker_permit,
850            stdin,
851            stdout,
852            stderr_task,
853            described: false,
854            provider_path: worker_provider_path,
855        };
856        startup_guard.disarm();
857        let hello = timeout(
858            self.config.startup_timeout,
859            read_frame::<PythonRunnerWorkerMessage, _>(&mut worker.stdout),
860        )
861        .await
862        .map_err(|_| {
863            PythonSupervisorError::new(
864                "python_worker_start_failed",
865                "Python worker handshake timed out",
866            )
867        })?
868        .map_err(|_| protocol_error())?;
869        let (protocol, features) = match hello {
870            PythonRunnerWorkerMessage::Hello {
871                protocol, features, ..
872            } => (protocol, features),
873            _ => return Err(protocol_error()),
874        };
875        let protocol = PythonRunnerProtocolVersion::current()
876            .negotiate(protocol)
877            .map_err(|_| protocol_error())?;
878        let requested = [
879            PythonRunnerFeature::Describe,
880            PythonRunnerFeature::Invoke,
881            PythonRunnerFeature::Health,
882            PythonRunnerFeature::Drain,
883            PythonRunnerFeature::Shutdown,
884            PythonRunnerFeature::HostCalls,
885        ];
886        let features = negotiate_runner_features(&requested, &features);
887        if features != requested {
888            self.active_pid.store(0, Ordering::Release);
889            terminate_worker(Some(worker)).await;
890            return Err(protocol_error());
891        }
892        write_frame(
893            &mut worker.stdin,
894            &PythonRunnerHostMessage::Initialize {
895                protocol,
896                features: features.clone(),
897                generation_id: self.identity.generation_id.clone(),
898            },
899        )
900        .await?;
901        match timeout(
902            self.config.startup_timeout,
903            read_frame::<PythonRunnerWorkerMessage, _>(&mut worker.stdout),
904        )
905        .await
906        {
907            Ok(Ok(PythonRunnerWorkerMessage::Ready {
908                protocol: ready_protocol,
909                features: ready_features,
910                generation_id,
911            })) if ready_protocol == protocol
912                && ready_features == features
913                && generation_id == self.identity.generation_id =>
914            {
915                self.active_pid
916                    .store(worker.child_pid.unwrap_or_default(), Ordering::Release);
917                Ok(worker)
918            }
919            _ => {
920                self.active_pid.store(0, Ordering::Release);
921                terminate_worker(Some(worker)).await;
922                Err(protocol_error())
923            }
924        }
925    }
926
927    fn record_restart(&self) -> Result<(), PythonSupervisorError> {
928        let now = Instant::now();
929        let mut restarts = self
930            .restarts
931            .lock()
932            .expect("Python worker restart lock should not be poisoned");
933        while restarts
934            .front()
935            .is_some_and(|started| now.duration_since(*started) > self.config.restart_window)
936        {
937            restarts.pop_front();
938        }
939        if restarts.len() >= self.config.max_restarts as usize {
940            self.quarantined.store(true, Ordering::Release);
941            return Err(PythonSupervisorError::new(
942                "python_provider_quarantined",
943                "Python provider is quarantined after repeated worker failures",
944            ));
945        }
946        restarts.push_back(now);
947        Ok(())
948    }
949
950    fn current_restart_count(&self) -> usize {
951        let now = Instant::now();
952        let mut restarts = self
953            .restarts
954            .lock()
955            .expect("Python worker restart lock should not be poisoned");
956        while restarts
957            .front()
958            .is_some_and(|started| now.duration_since(*started) > self.config.restart_window)
959        {
960            restarts.pop_front();
961        }
962        restarts.len()
963    }
964
965    fn worker_running(&self) -> bool {
966        let Ok(mut slot) = self.worker.try_lock() else {
967            return self.active_pid.load(Ordering::Acquire) != 0;
968        };
969        let Some(worker) = slot.as_mut() else {
970            self.active_pid.store(0, Ordering::Release);
971            return false;
972        };
973        match worker.child.try_wait() {
974            Ok(None) => true,
975            Ok(Some(_)) | Err(_) => {
976                self.active_pid.store(0, Ordering::Release);
977                false
978            }
979        }
980    }
981
982    fn next_request_id(&self) -> u64 {
983        self.request_id.fetch_add(1, Ordering::Relaxed)
984    }
985
986    fn verify_source_digest(&self) -> Result<(), PythonSupervisorError> {
987        use sha2::{Digest, Sha256};
988        let bytes = std::fs::read(&self.identity.path).map_err(|_| {
989            PythonSupervisorError::new(
990                "python_source_changed",
991                "Python provider source could not be revalidated",
992            )
993        })?;
994        let actual = Sha256::digest(bytes)
995            .iter()
996            .map(|byte| format!("{byte:02x}"))
997            .collect::<String>();
998        if actual != self.identity.source_digest {
999            return Err(PythonSupervisorError::new(
1000                "python_source_changed",
1001                "Python provider source changed during worker activation",
1002            ));
1003        }
1004        Ok(())
1005    }
1006
1007    pub async fn shutdown(&self) {
1008        let mut slot = self.worker.lock().await;
1009        let Some(mut worker) = slot.take() else {
1010            self.started_once.store(false, Ordering::Release);
1011            return;
1012        };
1013        let request_id = self.next_request_id();
1014        let request = PythonRunnerHostMessage::Request {
1015            request: PythonRunnerHostRequest::Shutdown { request_id },
1016        };
1017        let _ = write_frame(&mut worker.stdin, &request).await;
1018        let _ = timeout(self.config.shutdown_grace, worker.child.wait()).await;
1019        terminate_worker(Some(worker)).await;
1020        self.active_pid.store(0, Ordering::Release);
1021        self.started_once.store(false, Ordering::Release);
1022    }
1023
1024    /// Stop accepting work, then shut the worker down within the configured
1025    /// grace period. Used when a registry generation is retired.
1026    pub async fn drain_and_shutdown(&self) {
1027        let mut slot = self.worker.lock().await;
1028        let Some(worker) = slot.as_mut() else {
1029            return;
1030        };
1031        let request_id = self.next_request_id();
1032        let _ = write_frame(
1033            &mut worker.stdin,
1034            &PythonRunnerHostMessage::Request {
1035                request: PythonRunnerHostRequest::Drain { request_id },
1036            },
1037        )
1038        .await;
1039        drop(slot);
1040        self.shutdown().await;
1041    }
1042}
1043
1044#[cfg(test)]
1045mod tests {
1046    use super::*;
1047    use serde_json::json;
1048    use sha2::{Digest, Sha256};
1049    use std::{fs, path::Path};
1050    use tokio::io::AsyncWriteExt;
1051
1052    /// Deadline for anything these tests expect to *succeed*. Booting a Python
1053    /// worker on a machine running the rest of the suite in parallel is not
1054    /// bounded by any interesting constant, and a tight deadline here only
1055    /// converts scheduler pressure into a false failure. Deliberately short
1056    /// deadlines belong on the calls that are asserted to time out.
1057    const GENEROUS: Duration = Duration::from_secs(60);
1058
1059    /// Polls `condition` until it reports true or `deadline` elapses, yielding
1060    /// between attempts. Returns whether the condition was observed.
1061    ///
1062    /// Prefer this over sleeping a fixed interval and hoping the state has
1063    /// settled: the sleep encodes a guess about machine speed, this encodes the
1064    /// actual precondition.
1065    async fn wait_for(deadline: Duration, mut condition: impl FnMut() -> bool) -> bool {
1066        let expiry = tokio::time::Instant::now() + deadline;
1067        while tokio::time::Instant::now() < expiry {
1068            if condition() {
1069                return true;
1070            }
1071            tokio::time::sleep(Duration::from_millis(5)).await;
1072        }
1073        condition()
1074    }
1075
1076    fn installed_test_python() -> PathBuf {
1077        let root = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../../packages/python/.venv");
1078        let path = if cfg!(windows) {
1079            root.join("Scripts/python.exe")
1080        } else {
1081            root.join("bin/python")
1082        };
1083        assert!(
1084            path.is_file(),
1085            "persistent-runner tests require `uv sync --project packages/python --frozen`"
1086        );
1087        path
1088    }
1089
1090    fn identity(path: &Path) -> PythonWorkerIdentity {
1091        let source_digest = Sha256::digest(fs::read(path).expect("provider source"))
1092            .iter()
1093            .map(|byte| format!("{byte:02x}"))
1094            .collect();
1095        PythonWorkerIdentity {
1096            path: path.to_owned(),
1097            generation_id: "supervisor-test-generation".to_owned(),
1098            worker_group: "supervisor-test-generation".to_owned(),
1099            source_digest,
1100            catalog_fingerprint: String::new(),
1101        }
1102    }
1103
1104    #[tokio::test]
1105    async fn installed_runner_preflights_invokes_times_out_and_restarts() {
1106        let python = installed_test_python();
1107        let temp = tempfile::tempdir().expect("tempdir");
1108        let provider = temp.path().join("persistent.py");
1109        fs::write(
1110            &provider,
1111            r#"
1112import time
1113PROVIDER = {"name": "persistent-test", "kind": "python"}
1114
1115def execute(value: str, delay_ms: int = 0) -> dict:
1116    if delay_ms:
1117        time.sleep(delay_ms / 1000)
1118    return {"value": value}
1119"#,
1120        )
1121        .expect("write provider");
1122        // `invoke` waits `config.request_timeout.min(per_call_timeout)`, so the
1123        // config value caps every call including the two that must succeed.
1124        // Keep it generous and let each call's own override decide: the slow
1125        // call below passes a deliberately short one to force the timeout.
1126        let config = PythonSupervisorConfig {
1127            request_timeout: GENEROUS,
1128            startup_timeout: GENEROUS,
1129            restart_backoff: Duration::ZERO,
1130            ..PythonSupervisorConfig::default()
1131        };
1132        let supervisor = PythonWorkerSupervisor::new(
1133            identity(&provider),
1134            PythonInterpreter::Prepared(python),
1135            config,
1136        );
1137
1138        let catalog = supervisor.preflight().await.expect("preflight");
1139        assert_eq!(catalog["provider"]["name"], "persistent-test");
1140        let output = supervisor
1141            .invoke(
1142                "persistent-test",
1143                "execute",
1144                json!({"value": "first"}),
1145                soma_provider_core::ProviderSurface::Mcp,
1146                "snapshot-a",
1147                GENEROUS,
1148            )
1149            .await
1150            .expect("first invocation");
1151        assert_eq!(output, json!({"value": "first"}));
1152
1153        let timeout = supervisor
1154            .invoke(
1155                "persistent-test",
1156                "execute",
1157                json!({"value": "slow", "delay_ms": 300}),
1158                soma_provider_core::ProviderSurface::Mcp,
1159                "snapshot-a",
1160                Duration::from_millis(100),
1161            )
1162            .await
1163            .expect_err("slow invocation must time out");
1164        assert_eq!(timeout.code(), "python_provider_timeout");
1165
1166        let restarted = supervisor
1167            .invoke(
1168                "persistent-test",
1169                "execute",
1170                json!({"value": "restarted"}),
1171                soma_provider_core::ProviderSurface::Mcp,
1172                "snapshot-a",
1173                GENEROUS,
1174            )
1175            .await
1176            .expect("later invocation restarts without replay");
1177        assert_eq!(restarted, json!({"value": "restarted"}));
1178        supervisor.drain_and_shutdown().await;
1179    }
1180
1181    #[cfg(target_os = "linux")]
1182    #[tokio::test]
1183    #[ignore = "requires a delegated cgroup-v2 root in SOMA_PYTHON_BROKER_CGROUP_ROOT"]
1184    async fn brokered_worker_launches_inside_the_enforced_boundary() {
1185        let python = installed_test_python();
1186        let temp = tempfile::tempdir().expect("tempdir");
1187        let provider = temp.path().join("brokered.py");
1188        fs::write(temp.path().join("host-sentinel.txt"), "must remain hidden")
1189            .expect("write sentinel");
1190        fs::write(
1191            &provider,
1192            r#"
1193from pathlib import Path
1194
1195try:
1196    Path(__file__).with_name("host-sentinel.txt").read_text()
1197    SENTINEL_VISIBLE = True
1198except OSError:
1199    SENTINEL_VISIBLE = False
1200
1201PROVIDER = {"name": "brokered-test", "kind": "python"}
1202def execute(value: str) -> dict:
1203    return {"value": value, "sentinel_visible": SENTINEL_VISIBLE}
1204"#,
1205        )
1206        .expect("write provider");
1207        let supervisor = PythonWorkerSupervisor::new_with_capabilities(
1208            identity(&provider),
1209            PythonInterpreter::Prepared(python),
1210            PythonSupervisorConfig {
1211                execution_profile: PythonExecutionProfile::Brokered,
1212                ..PythonSupervisorConfig::default()
1213            },
1214            &soma_provider_core::HostCapabilities::default(),
1215        );
1216        supervisor.preflight().await.expect("brokered preflight");
1217        let output = supervisor
1218            .invoke(
1219                "brokered-test",
1220                "execute",
1221                json!({"value": "contained"}),
1222                soma_provider_core::ProviderSurface::Mcp,
1223                "snapshot-a",
1224                GENEROUS,
1225            )
1226            .await
1227            .expect("brokered invocation");
1228        assert_eq!(
1229            output,
1230            json!({"value": "contained", "sentinel_visible": false})
1231        );
1232        supervisor.shutdown().await;
1233    }
1234
1235    #[tokio::test]
1236    async fn concurrent_invocation_is_rejected_before_queueing() {
1237        let python = installed_test_python();
1238        let temp = tempfile::tempdir().expect("tempdir");
1239        let provider = temp.path().join("busy.py");
1240        fs::write(
1241            &provider,
1242            r#"
1243import time
1244PROVIDER = {"name": "busy-test", "kind": "python"}
1245def wait(delay_ms: int) -> dict:
1246    time.sleep(delay_ms / 1000)
1247    return {"ok": True}
1248"#,
1249        )
1250        .expect("write provider");
1251        let supervisor = PythonWorkerSupervisor::new(
1252            identity(&provider),
1253            PythonInterpreter::Prepared(python),
1254            PythonSupervisorConfig::default(),
1255        );
1256        supervisor.preflight().await.expect("preflight");
1257        let first = {
1258            let supervisor = supervisor.clone();
1259            tokio::spawn(async move {
1260                supervisor
1261                    .invoke(
1262                        "busy-test",
1263                        "wait",
1264                        // Wide enough that the busy window below cannot close
1265                        // before the contending call is issued.
1266                        json!({"delay_ms": 3_000}),
1267                        soma_provider_core::ProviderSurface::Mcp,
1268                        "snapshot-a",
1269                        GENEROUS,
1270                    )
1271                    .await
1272            })
1273        };
1274        // Wait for the first call to actually occupy the worker instead of
1275        // assuming a fixed interval was long enough to get there.
1276        assert!(
1277            wait_for(GENEROUS, || supervisor.status().busy).await,
1278            "first invocation never occupied the worker"
1279        );
1280        let busy = supervisor
1281            .invoke(
1282                "busy-test",
1283                "wait",
1284                json!({"delay_ms": 0}),
1285                soma_provider_core::ProviderSurface::Mcp,
1286                "snapshot-a",
1287                GENEROUS,
1288            )
1289            .await
1290            .expect_err("second invocation must not queue");
1291        assert_eq!(busy.code(), "python_provider_busy");
1292        first.await.expect("join").expect("first call");
1293        supervisor.shutdown().await;
1294    }
1295
1296    #[tokio::test]
1297    async fn active_invocation_cancels_process_tree_and_later_work_restarts() {
1298        let python = installed_test_python();
1299        let temp = tempfile::tempdir().expect("tempdir");
1300        let provider = temp.path().join("cancel.py");
1301        fs::write(
1302            &provider,
1303            r#"
1304import time
1305PROVIDER = {"name": "cancel-test", "kind": "python"}
1306def wait(delay_ms: int) -> dict:
1307    time.sleep(delay_ms / 1000)
1308    return {"ok": True}
1309"#,
1310        )
1311        .expect("write provider");
1312        let supervisor = PythonWorkerSupervisor::new(
1313            identity(&provider),
1314            PythonInterpreter::Prepared(python),
1315            PythonSupervisorConfig {
1316                restart_backoff: Duration::ZERO,
1317                ..PythonSupervisorConfig::default()
1318            },
1319        );
1320        supervisor.preflight().await.expect("preflight");
1321        let active = {
1322            let supervisor = supervisor.clone();
1323            tokio::spawn(async move {
1324                supervisor
1325                    .invoke(
1326                        "cancel-test",
1327                        "wait",
1328                        json!({"delay_ms": 30_000}),
1329                        soma_provider_core::ProviderSurface::Mcp,
1330                        "snapshot-a",
1331                        GENEROUS,
1332                    )
1333                    .await
1334            })
1335        };
1336        // Same precondition race as
1337        // `closing_stderr_does_not_revoke_active_cancellation`: poll until the
1338        // invocation is genuinely cancellable rather than sleeping a guess.
1339        assert!(
1340            wait_for(GENEROUS, || supervisor.cancel_active()).await,
1341            "active invocation never became cancellable"
1342        );
1343        let error = active
1344            .await
1345            .expect("join")
1346            .expect_err("active call must cancel");
1347        assert_eq!(error.code(), "python_provider_cancelled");
1348
1349        let restarted = supervisor
1350            .invoke(
1351                "cancel-test",
1352                "wait",
1353                json!({"delay_ms": 0}),
1354                soma_provider_core::ProviderSurface::Mcp,
1355                "snapshot-a",
1356                GENEROUS,
1357            )
1358            .await
1359            .expect("later work starts a clean worker");
1360        assert_eq!(restarted, json!({"ok": true}));
1361        supervisor.deactivate();
1362        let draining = supervisor
1363            .invoke(
1364                "cancel-test",
1365                "wait",
1366                json!({"delay_ms": 0}),
1367                soma_provider_core::ProviderSurface::Mcp,
1368                "snapshot-a",
1369                GENEROUS,
1370            )
1371            .await
1372            .expect_err("retained generation must reject new work");
1373        assert_eq!(draining.code(), "python_worker_draining");
1374        supervisor.activate();
1375        supervisor
1376            .invoke(
1377                "cancel-test",
1378                "wait",
1379                json!({"delay_ms": 0}),
1380                soma_provider_core::ProviderSurface::Mcp,
1381                "snapshot-a",
1382                GENEROUS,
1383            )
1384            .await
1385            .expect("rollback activation permits new work");
1386        supervisor.shutdown().await;
1387    }
1388
1389    #[tokio::test]
1390    async fn closing_stderr_does_not_revoke_active_cancellation() {
1391        let python = installed_test_python();
1392        let temp = tempfile::tempdir().expect("tempdir");
1393        let provider = temp.path().join("close_stderr.py");
1394        fs::write(
1395            &provider,
1396            r#"
1397import os
1398import time
1399PROVIDER = {"name": "close-stderr-test", "kind": "python"}
1400def close_and_wait(delay_ms: int) -> dict:
1401    os.close(2)
1402    time.sleep(delay_ms / 1000)
1403    return {"ok": True}
1404"#,
1405        )
1406        .expect("write provider");
1407        let supervisor = PythonWorkerSupervisor::new(
1408            identity(&provider),
1409            PythonInterpreter::Prepared(python),
1410            PythonSupervisorConfig::default(),
1411        );
1412        supervisor.preflight().await.expect("preflight");
1413        let active = {
1414            let supervisor = supervisor.clone();
1415            tokio::spawn(async move {
1416                supervisor
1417                    .invoke(
1418                        "close-stderr-test",
1419                        "close_and_wait",
1420                        // Long enough that returning within CANCEL_BOUND below
1421                        // can only mean cancellation cut the call short.
1422                        json!({"delay_ms": 30_000}),
1423                        soma_provider_core::ProviderSurface::Mcp,
1424                        "snapshot-a",
1425                        GENEROUS,
1426                    )
1427                    .await
1428            })
1429        };
1430        // `cancel_active` reports false until the invocation has both marked
1431        // the supervisor busy and recorded the worker pid, and both of those
1432        // happen after the spawn above returns. Sleeping a fixed interval
1433        // races that setup under load; poll the real precondition instead.
1434        // A false return is side-effect free (it bails before touching any
1435        // state), so retrying it is safe.
1436        assert!(
1437            wait_for(GENEROUS, || supervisor.status().busy).await,
1438            "invocation never reached the worker"
1439        );
1440        assert!(supervisor.status().running);
1441        assert!(
1442            wait_for(GENEROUS, || supervisor.cancel_active()).await,
1443            "active invocation never became cancellable"
1444        );
1445
1446        // Well under the 30s the provider would otherwise sleep, so this still
1447        // proves cancellation short-circuited the call rather than waiting it
1448        // out, but with enough headroom to survive a loaded machine.
1449        const CANCEL_BOUND: Duration = Duration::from_secs(10);
1450        let error = tokio::time::timeout(CANCEL_BOUND, active)
1451            .await
1452            .expect("cancellation must not wait for the invocation to finish")
1453            .expect("join")
1454            .expect_err("active invocation is cancelled");
1455        assert_eq!(error.code(), "python_provider_cancelled");
1456        supervisor.shutdown().await;
1457    }
1458
1459    #[tokio::test]
1460    async fn dead_idle_worker_is_restarted_before_next_dispatch() {
1461        let python = installed_test_python();
1462        let temp = tempfile::tempdir().expect("tempdir");
1463        let provider = temp.path().join("idle_crash.py");
1464        fs::write(
1465            &provider,
1466            r#"
1467PROVIDER = {"name": "idle-crash-test", "kind": "python"}
1468def value() -> dict:
1469    return {"ok": True}
1470"#,
1471        )
1472        .expect("write provider");
1473        let supervisor = PythonWorkerSupervisor::new(
1474            identity(&provider),
1475            PythonInterpreter::Prepared(python),
1476            PythonSupervisorConfig {
1477                restart_backoff: Duration::ZERO,
1478                ..PythonSupervisorConfig::default()
1479            },
1480        );
1481        supervisor.preflight().await.expect("preflight");
1482        let original_pid = supervisor.active_pid.load(Ordering::Acquire);
1483        terminate_process_tree(Some(original_pid));
1484        for _ in 0..100 {
1485            if !supervisor.status().running {
1486                break;
1487            }
1488            tokio::time::sleep(Duration::from_millis(10)).await;
1489        }
1490        assert!(
1491            !supervisor.status().running,
1492            "status observes idle child death before another dispatch"
1493        );
1494
1495        let output = supervisor
1496            .invoke(
1497                "idle-crash-test",
1498                "value",
1499                json!({}),
1500                soma_provider_core::ProviderSurface::Mcp,
1501                "snapshot-a",
1502                GENEROUS,
1503            )
1504            .await
1505            .expect("first post-crash invocation starts a replacement");
1506        assert_eq!(output, json!({"ok": true}));
1507        assert_ne!(supervisor.active_pid.load(Ordering::Acquire), original_pid);
1508        supervisor.shutdown().await;
1509    }
1510
1511    #[tokio::test]
1512    async fn suspension_drains_routed_calls_and_does_not_consume_restart_budget() {
1513        let python = installed_test_python();
1514        let temp = tempfile::tempdir().expect("tempdir");
1515        let provider = temp.path().join("generation.py");
1516        fs::write(
1517            &provider,
1518            r#"
1519PROVIDER = {"name": "generation-test", "kind": "python"}
1520def value() -> dict:
1521    return {"value": "ok"}
1522"#,
1523        )
1524        .expect("write provider");
1525        let supervisor = PythonWorkerSupervisor::new(
1526            identity(&provider),
1527            PythonInterpreter::Prepared(python),
1528            PythonSupervisorConfig {
1529                restart_backoff: Duration::ZERO,
1530                ..PythonSupervisorConfig::default()
1531            },
1532        );
1533        supervisor.preflight().await.expect("preflight");
1534        assert!(supervisor.acquire_dispatch());
1535        let suspending = {
1536            let supervisor = supervisor.clone();
1537            tokio::spawn(async move { supervisor.suspend().await })
1538        };
1539        tokio::time::sleep(Duration::from_millis(20)).await;
1540        assert!(!supervisor.acquire_dispatch());
1541        let routed = supervisor
1542            .invoke(
1543                "generation-test",
1544                "value",
1545                json!({}),
1546                soma_provider_core::ProviderSurface::Mcp,
1547                "snapshot-a",
1548                GENEROUS,
1549            )
1550            .await
1551            .expect("already-routed call drains");
1552        assert_eq!(routed, json!({"value": "ok"}));
1553        supervisor.release_dispatch();
1554        suspending.await.expect("suspension task");
1555
1556        supervisor.activate();
1557        assert!(supervisor.acquire_dispatch());
1558        supervisor
1559            .invoke(
1560                "generation-test",
1561                "value",
1562                json!({}),
1563                soma_provider_core::ProviderSurface::Mcp,
1564                "snapshot-b",
1565                GENEROUS,
1566            )
1567            .await
1568            .expect("rollback starts a planned fresh worker");
1569        supervisor.release_dispatch();
1570        assert_eq!(supervisor.status().restart_count, 0);
1571        supervisor.shutdown().await;
1572    }
1573
1574    #[tokio::test]
1575    async fn worker_logs_are_bounded_structured_and_redacted() {
1576        let python = installed_test_python();
1577        let temp = tempfile::tempdir().expect("tempdir");
1578        let provider = temp.path().join("logs.py");
1579        fs::write(
1580            &provider,
1581            r#"
1582PROVIDER = {"name": "logs-test", "kind": "python"}
1583def emit() -> dict:
1584    print("token=super-secret")
1585    print("credential: unmarked-private-data")
1586    print("safe diagnostic")
1587    return {"ok": True}
1588"#,
1589        )
1590        .expect("write provider");
1591        let supervisor = PythonWorkerSupervisor::new(
1592            identity(&provider),
1593            PythonInterpreter::Prepared(python),
1594            PythonSupervisorConfig {
1595                max_stderr_bytes: 128,
1596                ..PythonSupervisorConfig::default()
1597            },
1598        );
1599        supervisor.preflight().await.expect("preflight");
1600        supervisor
1601            .invoke(
1602                "logs-test",
1603                "emit",
1604                json!({}),
1605                soma_provider_core::ProviderSurface::Mcp,
1606                "snapshot-a",
1607                Duration::from_secs(1),
1608            )
1609            .await
1610            .expect("invoke");
1611        tokio::time::sleep(Duration::from_millis(20)).await;
1612        let status = supervisor.status();
1613        assert!(
1614            status
1615                .logs
1616                .iter()
1617                .all(|entry| entry.message == "[redacted provider diagnostic]")
1618        );
1619        let encoded = serde_json::to_string(&status).unwrap();
1620        assert!(!encoded.contains("super-secret"));
1621        assert!(!encoded.contains("unmarked-private-data"));
1622        assert!(!encoded.contains("safe diagnostic"));
1623        assert!(
1624            status
1625                .logs
1626                .iter()
1627                .map(|entry| entry.message.len())
1628                .sum::<usize>()
1629                <= 128
1630        );
1631        supervisor.shutdown().await;
1632    }
1633
1634    #[test]
1635    fn worker_budget_is_bounded_per_immutable_generation() {
1636        let first = worker_budget("generation-a", 1);
1637        let _active = first.clone().try_acquire_owned().expect("first permit");
1638        assert!(first.try_acquire_owned().is_err());
1639
1640        let replacement = worker_budget("generation-b", 1);
1641        assert!(
1642            replacement.try_acquire_owned().is_ok(),
1643            "a replacement generation has bounded overlap capacity"
1644        );
1645    }
1646
1647    #[test]
1648    fn worker_budget_prunes_retired_generation_keys() {
1649        for generation in 0..128 {
1650            drop(worker_budget(&format!("retired-{generation}"), 1));
1651        }
1652        let _live = worker_budget("live-generation", 1);
1653        assert!(
1654            worker_budget_keys_are_live(),
1655            "retired generation keys must not accumulate"
1656        );
1657    }
1658
1659    #[tokio::test]
1660    async fn restart_count_expires_without_another_restart() {
1661        let temp = tempfile::tempdir().expect("tempdir");
1662        let provider = temp.path().join("provider.py");
1663        fs::write(
1664            &provider,
1665            "PROVIDER = {\"name\": \"restart-window\", \"kind\": \"python\"}\n",
1666        )
1667        .expect("provider");
1668        let supervisor = PythonWorkerSupervisor::new(
1669            identity(&provider),
1670            PythonInterpreter::Prepared(PathBuf::from("/unused")),
1671            PythonSupervisorConfig {
1672                restart_window: Duration::from_millis(10),
1673                ..PythonSupervisorConfig::default()
1674            },
1675        );
1676        supervisor.record_restart().expect("record restart");
1677        assert_eq!(supervisor.status().restart_count, 1);
1678        tokio::time::sleep(Duration::from_millis(20)).await;
1679        assert_eq!(supervisor.status().restart_count, 0);
1680    }
1681
1682    #[cfg(unix)]
1683    #[tokio::test]
1684    async fn failed_startup_terminates_the_entire_spawned_process_group() {
1685        use std::os::unix::fs::PermissionsExt;
1686
1687        let temp = tempfile::tempdir().expect("tempdir");
1688        let provider = temp.path().join("provider.py");
1689        fs::write(
1690            &provider,
1691            "PROVIDER = {\"name\": \"startup-tree\", \"kind\": \"python\"}\n",
1692        )
1693        .expect("provider");
1694        let descendant_pid = temp.path().join("descendant.pid");
1695        let fake_python = temp.path().join("fake-python");
1696        fs::write(
1697            &fake_python,
1698            format!(
1699                "#!/bin/sh\nsleep 30 &\necho $! > '{}'\nsleep 30\n",
1700                descendant_pid.display()
1701            ),
1702        )
1703        .expect("fake python");
1704        fs::set_permissions(&fake_python, fs::Permissions::from_mode(0o700))
1705            .expect("executable fake python");
1706
1707        let supervisor = PythonWorkerSupervisor::new(
1708            identity(&provider),
1709            PythonInterpreter::Prepared(fake_python),
1710            // The fake worker never connects, so startup fails at *any*
1711            // timeout — this value is not the property under test. It does
1712            // gate how long the fake shell has to fork its descendant and
1713            // record the pid the assertions below read, and a 150ms budget
1714            // lost that race whenever the machine was busy, leaving an empty
1715            // pid file.
1716            PythonSupervisorConfig {
1717                startup_timeout: Duration::from_secs(5),
1718                ..PythonSupervisorConfig::default()
1719            },
1720        );
1721        let error = supervisor
1722            .preflight()
1723            .await
1724            .expect_err("fake worker never connects");
1725        assert_eq!(error.code(), "python_worker_start_failed");
1726        let pid: u32 = fs::read_to_string(&descendant_pid)
1727            .expect("descendant pid")
1728            .trim()
1729            .parse()
1730            .expect("numeric pid");
1731        for _ in 0..100 {
1732            if !Path::new("/proc").join(pid.to_string()).exists() {
1733                return;
1734            }
1735            tokio::time::sleep(Duration::from_millis(10)).await;
1736        }
1737        panic!("descendant process {pid} survived failed worker startup");
1738    }
1739
1740    #[tokio::test]
1741    async fn quarantine_exhaustion_is_visible_and_operator_reset_recovers() {
1742        let python = installed_test_python();
1743        let temp = tempfile::tempdir().expect("tempdir");
1744        let provider = temp.path().join("quarantine.py");
1745        fs::write(
1746            &provider,
1747            r#"
1748import time
1749PROVIDER = {"name": "quarantine-test", "kind": "python"}
1750def wait(delay_ms: int) -> dict:
1751    time.sleep(delay_ms / 1000)
1752    return {"ok": True}
1753"#,
1754        )
1755        .expect("write provider");
1756        let supervisor = PythonWorkerSupervisor::new(
1757            identity(&provider),
1758            PythonInterpreter::Prepared(python),
1759            // As in `installed_runner_preflights_invokes_times_out_and_restarts`,
1760            // `invoke` waits `request_timeout.min(per_call)`. A short config
1761            // value here also capped the post-reset call, which has to boot a
1762            // brand-new interpreter — the timeout below supplies its own short
1763            // per-call deadline instead.
1764            PythonSupervisorConfig {
1765                request_timeout: GENEROUS,
1766                max_restarts: 0,
1767                restart_backoff: Duration::ZERO,
1768                ..PythonSupervisorConfig::default()
1769            },
1770        );
1771        supervisor.preflight().await.expect("preflight");
1772        let timeout = supervisor
1773            .invoke(
1774                "quarantine-test",
1775                "wait",
1776                json!({"delay_ms": 200}),
1777                soma_provider_core::ProviderSurface::Mcp,
1778                "snapshot-a",
1779                Duration::from_millis(50),
1780            )
1781            .await
1782            .expect_err("slow call times out");
1783        assert_eq!(timeout.code(), "python_provider_timeout");
1784        let quarantined = supervisor
1785            .invoke(
1786                "quarantine-test",
1787                "wait",
1788                json!({"delay_ms": 0}),
1789                soma_provider_core::ProviderSurface::Mcp,
1790                "snapshot-a",
1791                GENEROUS,
1792            )
1793            .await
1794            .expect_err("restart budget is exhausted");
1795        assert_eq!(quarantined.code(), "python_provider_quarantined");
1796        assert!(supervisor.status().quarantined);
1797
1798        supervisor.reset_quarantine().await;
1799        assert!(!supervisor.status().quarantined);
1800        let output = supervisor
1801            .invoke(
1802                "quarantine-test",
1803                "wait",
1804                json!({"delay_ms": 0}),
1805                soma_provider_core::ProviderSurface::Mcp,
1806                "snapshot-a",
1807                GENEROUS,
1808            )
1809            .await
1810            .expect("operator reset permits a fresh worker");
1811        assert_eq!(output, json!({"ok": true}));
1812        supervisor.shutdown().await;
1813    }
1814
1815    #[tokio::test]
1816    async fn source_substitution_and_missing_runner_fail_closed() {
1817        let temp = tempfile::tempdir().expect("tempdir");
1818        let provider = temp.path().join("source.py");
1819        fs::write(
1820            &provider,
1821            "PROVIDER = {\"name\": \"source-test\", \"kind\": \"python\"}\n",
1822        )
1823        .expect("write provider");
1824        let source_identity = identity(&provider);
1825        fs::write(
1826            &provider,
1827            "PROVIDER = {\"name\": \"substituted\", \"kind\": \"python\"}\n",
1828        )
1829        .expect("substitute provider");
1830        let substituted = PythonWorkerSupervisor::new(
1831            source_identity,
1832            PythonInterpreter::Prepared(PathBuf::from("/missing/soma-python")),
1833            PythonSupervisorConfig::default(),
1834        );
1835        assert_eq!(
1836            substituted.preflight().await.unwrap_err().code(),
1837            "python_source_changed"
1838        );
1839
1840        let missing = PythonWorkerSupervisor::new(
1841            identity(&provider),
1842            PythonInterpreter::Prepared(PathBuf::from("/missing/soma-python")),
1843            PythonSupervisorConfig::default(),
1844        );
1845        assert_eq!(
1846            missing.preflight().await.unwrap_err().code(),
1847            "python_worker_start_failed"
1848        );
1849    }
1850
1851    #[tokio::test]
1852    async fn production_control_reader_rejects_malformed_frames() {
1853        let listener = TcpListener::bind("127.0.0.1:0").await.expect("listener");
1854        let address = listener.local_addr().expect("address");
1855        let client = tokio::spawn(async move {
1856            let mut stream = tokio::net::TcpStream::connect(address)
1857                .await
1858                .expect("connect");
1859            stream
1860                .write_all(&5_u32.to_be_bytes())
1861                .await
1862                .expect("header");
1863            stream.write_all(b"{nope").await.expect("payload");
1864        });
1865        let (stream, _) = listener.accept().await.expect("accept");
1866        let (mut reader, _) = stream.into_split();
1867        let error = read_frame::<PythonRunnerWorkerMessage, _>(&mut reader)
1868            .await
1869            .expect_err("malformed production frame must fail closed");
1870        assert_eq!(error.code(), "python_protocol_mismatch");
1871        client.await.expect("client task");
1872    }
1873}