1use 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#[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#[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#[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
164pub 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 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 pub fn deactivate(&self) {
248 self.accepting.store(false, Ordering::Release);
249 }
250
251 pub fn activate(&self) {
253 self.accepting.store(true, Ordering::Release);
254 }
255
256 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 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 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 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 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 const GENEROUS: Duration = Duration::from_secs(60);
1058
1059 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 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 json!({"delay_ms": 3_000}),
1267 soma_provider_core::ProviderSurface::Mcp,
1268 "snapshot-a",
1269 GENEROUS,
1270 )
1271 .await
1272 })
1273 };
1274 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 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 json!({"delay_ms": 30_000}),
1423 soma_provider_core::ProviderSurface::Mcp,
1424 "snapshot-a",
1425 GENEROUS,
1426 )
1427 .await
1428 })
1429 };
1430 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 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 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 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}