Skip to main content

scv_server/
lib.rs

1//! SCV's authoritative stdio server.
2
3mod agents;
4pub mod components;
5mod config;
6
7use std::{
8    collections::{HashMap, VecDeque},
9    io::Read as _,
10    path::{Path, PathBuf},
11    sync::{
12        Arc,
13        atomic::{AtomicU64, AtomicUsize, Ordering},
14    },
15    time::Duration,
16};
17
18use anyhow::{Context, Result, anyhow};
19use async_trait::async_trait;
20use config::Config;
21pub use config::{ApprovalPolicy, ConfigOverrides, user_home_path};
22use sha2::{Digest, Sha256};
23pub fn init_user_config() -> anyhow::Result<std::path::PathBuf> {
24    config::Config::init_user_config()
25}
26pub fn update_index_url(workspace: &std::path::Path) -> anyhow::Result<Option<String>> {
27    Ok(config::Config::load(workspace, ConfigOverrides::default())?
28        .update
29        .index_url)
30}
31
32pub use agents::{Endpoint, PI_PROVIDER, WireApi, read_secret};
33pub use scv_tools::{adapters, delegation};
34
35/// Build a command for a native agent's CLI with the same private home and
36/// cleaned environment the daemon's `agent_<name>` tool uses, so the agent's
37/// own sign-in stores credentials where delegated runs will find them.
38/// Project configuration cannot set `[agents]`, so none is read.
39pub fn agent_command(agent: &str) -> Result<std::process::Command> {
40    let config = Config::load_user(ConfigOverrides::default())?;
41    config.prepare_adapter_homes()?;
42    let adapter = config
43        .adapters()
44        .remove(&format!("agent_{agent}"))
45        .ok_or_else(|| anyhow!("unknown agent {agent}"))?;
46    let executable =
47        scv_tools::adapters::resolve_agent_executable(&adapter.command, &adapter.search_dirs)
48            .ok_or_else(|| {
49                anyhow!(
50                    "{agent} is not installed: {:?} was not found on PATH or in ~/.local/bin",
51                    adapter.command
52                )
53            })?;
54    let mut command = std::process::Command::new(executable);
55    command.current_dir(config.instance_home.join("adapters").join(agent));
56    scv_tools::apply_agent_environment(&mut command, &adapter.environment);
57    Ok(command)
58}
59
60/// Where `agent`'s executable resolves, as the daemon would find it.
61pub fn agent_executable(agent: &str) -> Result<Option<PathBuf>> {
62    let config = Config::load_user(ConfigOverrides::default())?;
63    let adapter = config
64        .adapters()
65        .remove(&format!("agent_{agent}"))
66        .ok_or_else(|| anyhow!("unknown agent {agent}"))?;
67    Ok(scv_tools::adapters::resolve_agent_executable(
68        &adapter.command,
69        &adapter.search_dirs,
70    ))
71}
72
73/// The prepared private adapter home for `agent`.
74pub fn agent_home(agent: &str) -> Result<PathBuf> {
75    let config = Config::load_user(ConfigOverrides::default())?;
76    config.prepare_adapter_homes()?;
77    let home = config.instance_home.join("adapters").join(agent);
78    if !home.is_dir() {
79        return Err(anyhow!("unknown agent {agent}"));
80    }
81    Ok(home)
82}
83
84/// Remove delegated-conversation transcripts older than `older_than` from
85/// the adapter homes (`agent`, or every agent that keeps them), keeping any a
86/// live conversation still uses. Returns each agent's report.
87pub fn collect_agent_garbage(
88    agent: Option<&str>,
89    older_than: std::time::Duration,
90    dry_run: bool,
91) -> Result<Vec<(&'static str, scv_tools::conversation::GcReport)>> {
92    let config = Config::load_user(ConfigOverrides::default())?;
93    let markers = config.instance_home.join("run").join("conversations");
94    let mut reports = Vec::new();
95    for adapter in adapters::ADAPTERS {
96        if agent.is_some_and(|agent| agent != adapter.name) {
97            continue;
98        }
99        let Some(files) = adapter.conversation_files else {
100            continue;
101        };
102        let home = config.instance_home.join("adapters").join(adapter.name);
103        if !home.is_dir() {
104            continue;
105        }
106        let report =
107            scv_tools::conversation::collect_garbage(&home, files, &markers, older_than, dry_run)
108                .with_context(|| format!("clean {} transcripts", adapter.name))?;
109        reports.push((adapter.name, report));
110    }
111    Ok(reports)
112}
113
114/// Parse a `scv agents gc --older-than` age such as `30d`.
115pub fn conversation_age(value: &str) -> std::result::Result<std::time::Duration, String> {
116    scv_tools::conversation::parse_age(value)
117}
118
119fn key_store_home(agent: &str) -> Result<PathBuf> {
120    scv_tools::adapters::adapter(agent).ok_or_else(|| anyhow!("unknown agent {agent}"))?;
121    agent_home(agent)
122}
123
124/// Store `key` in the agent's native credential file inside its adapter home.
125pub fn store_agent_key(agent: &str, store: adapters::KeyStore, key: &str) -> Result<Vec<String>> {
126    agents::store_key(store, &key_store_home(agent)?, key)
127}
128
129/// Whether the agent's stored credentials exist, with display lines that
130/// never contain secrets.
131pub fn agent_stored_status(agent: &str, store: adapters::KeyStore) -> Result<(bool, Vec<String>)> {
132    agents::stored_status(store, &key_store_home(agent)?)
133}
134
135/// Remove the agent's stored credentials from its adapter home.
136pub fn remove_agent_credentials(agent: &str, store: adapters::KeyStore) -> Result<Vec<String>> {
137    agents::remove_stored(store, &key_store_home(agent)?)
138}
139
140/// Point SCV's pi at an OpenAI-compatible endpoint.
141pub fn configure_pi_endpoint(endpoint: &Endpoint, key: &str) -> Result<Vec<String>> {
142    agents::configure_pi_endpoint(&pi_agent_dir()?, endpoint, key)
143}
144
145/// Point SCV's pi at SCV's own active provider: its base URL, model, and key
146/// (read from `api_key`, or from the `api_key_env` variable now, since
147/// delegated agents never inherit key variables).
148pub fn import_pi_from_scv_provider() -> Result<Vec<String>> {
149    let config = Config::load_user(ConfigOverrides::default())?;
150    let provider = &config.provider;
151    let key = scv_provider_key(provider)?;
152    let endpoint = Endpoint {
153        base_url: provider.base_url.clone(),
154        api: WireApi::Responses,
155        model: provider.model.clone(),
156    };
157    let mut notes = agents::configure_pi_endpoint(&pi_agent_dir()?, &endpoint, &key)?;
158    if !provider.headers.is_empty() {
159        notes.push(
160            "Note: SCV's provider sends extra headers, which were not copied; add them to \
161             pi's models.json if the endpoint needs them"
162                .into(),
163        );
164    }
165    Ok(notes)
166}
167
168/// The key of SCV's own provider: `api_key`, or the `api_key_env` variable
169/// read now, since delegated agents never inherit key variables.
170fn scv_provider_key(provider: &config::ProviderConfig) -> Result<String> {
171    if provider.kind != "openai-compatible" {
172        return Err(anyhow!("SCV's provider is not openai-compatible"));
173    }
174    match (&provider.api_key, &provider.api_key_env) {
175        (Some(key), _) if !key.trim().is_empty() => Ok(key.trim().to_owned()),
176        (_, Some(variable)) => std::env::var(variable)
177            .ok()
178            .filter(|key| !key.trim().is_empty())
179            .ok_or_else(|| {
180                anyhow!("SCV's provider reads its key from ${variable}, which is not set here")
181            }),
182        _ => Err(anyhow!("SCV's provider has no API key configured")),
183    }
184}
185
186/// Give the nested SCV (`agent_scv`) its own copy of SCV's active provider,
187/// in `$SCV_HOME/adapters/scv/config.toml` (mode 0600).
188pub fn import_scv_from_scv_provider() -> Result<Vec<String>> {
189    let config = Config::load_user(ConfigOverrides::default())?;
190    config.prepare_adapter_homes()?;
191    let provider = &config.provider;
192    let key = scv_provider_key(provider)?;
193    agents::configure_scv_child(
194        &agent_home("scv")?,
195        &agents::ScvChildProvider {
196            wire_api: &provider.wire_api,
197            model: &provider.model,
198            base_url: &provider.base_url,
199            timeout_seconds: provider.timeout_seconds,
200            headers: &provider.headers,
201            hosted_web_search: config.hosted_web_search(),
202        },
203        &key,
204    )
205}
206
207fn pi_agent_dir() -> Result<PathBuf> {
208    let descriptor =
209        scv_tools::adapters::adapter("pi").ok_or_else(|| anyhow!("unknown agent pi"))?;
210    let scv_tools::adapters::Status::Stored(scv_tools::adapters::KeyStore::Pi { dir }) =
211        descriptor.status
212    else {
213        return Err(anyhow!("pi has no SCV-managed store"));
214    };
215    Ok(agent_home("pi")?.join(dir))
216}
217
218/// Copy the user's own Codex setup from `source` into SCV's private Codex
219/// adapter home: `config.toml`, and `auth.json` only when it holds an API key.
220/// Returns display lines that never contain secret values.
221pub fn import_codex(source: &Path) -> Result<Vec<String>> {
222    let config = Config::load_user(ConfigOverrides::default())?;
223    config.prepare_adapter_homes()?;
224    agents::import_codex(source, &config.instance_home.join("adapters").join("codex"))
225}
226
227/// Copy the user's own Grok `config.toml` from `source` (a Grok home) into
228/// SCV's private Grok home, keeping settings only SCV's copy has. Returns
229/// display lines that never contain secret values.
230pub fn import_grok(source: &Path) -> Result<Vec<String>> {
231    let config = Config::load_user(ConfigOverrides::default())?;
232    config.prepare_adapter_homes()?;
233    let descriptor =
234        scv_tools::adapters::adapter("grok").ok_or_else(|| anyhow!("unknown agent grok"))?;
235    let grok_home = descriptor
236        .home_environment
237        .iter()
238        .find(|(variable, _)| *variable == "GROK_HOME")
239        .map(|(_, relative)| *relative)
240        .ok_or_else(|| anyhow!("grok has no GROK_HOME in its adapter home"))?;
241    agents::import_grok(
242        source,
243        &config
244            .instance_home
245            .join("adapters")
246            .join("grok")
247            .join(grok_home),
248    )
249}
250
251/// Return the user service name for the selected SCV instance.
252pub fn service_name() -> anyhow::Result<String> {
253    if std::env::var_os("SCV_HOME").is_none() {
254        return Ok("scv.service".into());
255    }
256    let home = user_home_path().ok_or_else(|| anyhow!("cannot determine SCV instance home"))?;
257    let digest = Sha256::digest(home.to_string_lossy().as_bytes());
258    let suffix = digest[..8]
259        .iter()
260        .map(|byte| format!("{byte:02x}"))
261        .collect::<String>();
262    Ok(format!("scv-{suffix}.service"))
263}
264
265pub fn service_unit_path() -> anyhow::Result<std::path::PathBuf> {
266    let config =
267        dirs::config_dir().ok_or_else(|| anyhow!("cannot determine XDG config directory"))?;
268    Ok(config.join("systemd/user").join(service_name()?))
269}
270use scv_core::{
271    AgentError, AgentRuntime, ApprovalGate, ApprovalRequest, BudgetContextPolicy, CoreEvent,
272    EventSink, Message, ToolRegistry, ToolRisk,
273};
274use scv_protocol::{
275    ClientMessage, DaemonCommand, DaemonStatus, DelegationInfo, DelegationSummary,
276    ORIGIN_BACKGROUND, PROTOCOL_VERSION, PeerInfo, QueueEntry, ServerEvent, TurnOrigin, Usage,
277};
278use scv_provider_openai::OpenAiProvider;
279use scv_tools::{
280    DelegationContext, SkillMap,
281    background::{self, BackgroundJobs},
282    builtin_registry,
283    delegation::{self as delegations, DelegationRegistry},
284};
285use tokio::{
286    io::{AsyncBufRead, AsyncBufReadExt, AsyncWriteExt, BufReader},
287    net::{UnixListener, UnixStream},
288    sync::{Mutex, OwnedSemaphorePermit, Semaphore, mpsc, oneshot},
289    task::JoinHandle,
290};
291use tokio_util::sync::CancellationToken;
292use tokio_util::task::TaskTracker;
293use uuid::Uuid;
294
295const PROMPT_LIMIT_BYTES: usize = 256 * 1024;
296const OUTPUT_QUEUE_CAPACITY: usize = 256;
297const OUTPUT_QUEUE_MIN_BYTES: usize = 16 * 1024 * 1024;
298const SHUTDOWN_GRACE: Duration = Duration::from_secs(3);
299const MAX_QUEUE_ITEMS: usize = 64;
300const MAX_QUEUE_BYTES: usize = 4 * 1024 * 1024;
301
302pub async fn run_stdio(overrides: ConfigOverrides) -> Result<()> {
303    let stdin = tokio::io::stdin();
304    let stdout = tokio::io::stdout();
305    let tasks = TaskTracker::new();
306    let registry = instance_delegations()?;
307    // Without a daemon, a later `scv exec` is what cleans up after an earlier
308    // one that was killed; this runs alongside the session.
309    tokio::spawn(reconcile_delegations(Arc::clone(&registry)));
310    let result = run_managed(
311        stdin,
312        stdout,
313        overrides,
314        None,
315        registry,
316        CancellationToken::new(),
317        tasks.clone(),
318    )
319    .await;
320    tasks.close();
321    tasks.wait().await;
322    result
323}
324
325/// Return the local Unix socket used by the SCV daemon and TUI.
326pub fn default_socket_path() -> Result<PathBuf> {
327    scv_client::default_socket_path()
328}
329
330/// Run the authoritative server on the local Unix socket.
331pub async fn run_socket(path: &Path, overrides: ConfigOverrides) -> Result<()> {
332    if let Some(parent) = path.parent() {
333        tokio::fs::create_dir_all(parent)
334            .await
335            .context("create SCV socket directory")?;
336    }
337    let _lock = SocketLock::acquire(path)?;
338    if path.exists() {
339        if UnixStream::connect(path).await.is_ok() {
340            return Err(anyhow!(
341                "SCV server is already running at {}",
342                path.display()
343            ));
344        }
345        use std::os::unix::fs::FileTypeExt;
346        if !std::fs::symlink_metadata(path)?.file_type().is_socket() {
347            return Err(anyhow!(
348                "refusing to remove a non-socket at SCV socket path"
349            ));
350        }
351        tokio::fs::remove_file(path)
352            .await
353            .with_context(|| format!("remove stale SCV socket {}", path.display()))?;
354    }
355    let listener = UnixListener::bind(path)
356        .with_context(|| format!("bind SCV server socket {}", path.display()))?;
357    #[cfg(unix)]
358    {
359        use std::os::unix::fs::PermissionsExt;
360        std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600))
361            .context("secure SCV socket")?;
362    }
363    let components = Arc::new(Mutex::new(components::Components::new(
364        path.to_owned(),
365        std::env::current_dir()?,
366    )));
367    let registry = instance_delegations()?;
368    // Descendants a delegated agent leaves behind reparent to the daemon, not init.
369    if !delegations::become_child_subreaper() {
370        tracing::debug!("SCV daemon is not a child subreaper on this platform");
371    }
372    let cancellation = CancellationToken::new();
373    let delegation_registry = Arc::clone(&registry);
374    let delegation_cancel = cancellation.clone();
375    let mut delegation_task = tokio::spawn(async move {
376        // The first tick is immediate: orphans from before a restart go first.
377        let mut interval = tokio::time::interval(DELEGATION_RECONCILE_INTERVAL);
378        loop {
379            tokio::select! {
380                biased;
381                _ = delegation_cancel.cancelled() => break,
382                _ = interval.tick() => {
383                    reconcile_delegations(Arc::clone(&delegation_registry)).await;
384                    let zombies = delegations::reap_orphaned_zombies();
385                    if zombies > 0 {
386                        tracing::debug!("Reaped {zombies} exited orphan processes");
387                    }
388                }
389            }
390        }
391    });
392    let _delegation_abort = AbortGuard(delegation_task.abort_handle());
393    let tasks = TaskTracker::new();
394    let mut clients = tokio::task::JoinSet::new();
395    let refresh_components = components.clone();
396    let refresh_cancel = cancellation.clone();
397    let mut refresh_task = tokio::spawn(async move {
398        let mut refresh = tokio::time::interval(Duration::from_secs(2));
399        loop {
400            tokio::select! {
401                biased;
402                _ = refresh_cancel.cancelled() => break,
403                _ = refresh.tick() => {
404                    tokio::select! {
405                        biased;
406                        _ = refresh_cancel.cancelled() => break,
407                        result = async { refresh_components.lock().await.reconcile().await } => {
408                            if result.is_err() { tracing::warn!("Component account discovery failed"); }
409                        }
410                    }
411                }
412            }
413        }
414    });
415    let _refresh_abort = AbortGuard(refresh_task.abort_handle());
416    let mut terminate = tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())?;
417    let result = loop {
418        tokio::select! {
419            accepted = listener.accept() => {
420                let (stream, _) = match accepted { Ok(value) => value, Err(error) => break Err(error.into()) };
421                let child_overrides = overrides.clone();
422                let components = components.clone();
423                let registry = Arc::clone(&registry);
424                let cancellation = cancellation.clone();
425                let tasks = tasks.clone();
426                clients.spawn(async move {
427                    let (reader, writer) = stream.into_split();
428                    if run_managed(reader, writer, child_overrides, Some(components), registry, cancellation, tasks).await.is_err() {
429                        tracing::warn!("SCV socket client stopped");
430                    }
431                });
432            }
433            _ = clients.join_next(), if !clients.is_empty() => {},
434            _ = tokio::signal::ctrl_c() => break Ok(()),
435            _ = terminate.recv() => break Ok(()),
436        }
437    };
438    drop(listener);
439    cancellation.cancel();
440    let _ = (&mut refresh_task).await;
441    let _ = (&mut delegation_task).await;
442    components.lock().await.shutdown().await;
443    if tokio::time::timeout(Duration::from_secs(8), async {
444        while clients.join_next().await.is_some() {}
445    })
446    .await
447    .is_err()
448    {
449        clients.abort_all();
450        while clients.join_next().await.is_some() {}
451    }
452    tasks.close();
453    tasks.wait().await;
454    let _ = tokio::fs::remove_file(path).await;
455    result
456}
457
458const DELEGATION_RECONCILE_INTERVAL: Duration = Duration::from_secs(60);
459
460/// The delegation registry for this process's SCV instance.
461fn instance_delegations() -> Result<Arc<DelegationRegistry>> {
462    let home =
463        config::user_home_path().ok_or_else(|| anyhow!("cannot determine SCV instance home"))?;
464    Ok(Arc::new(DelegationRegistry::new(&home)))
465}
466
467/// Stop orphaned delegations of this instance and log what was stopped.
468async fn reconcile_delegations(registry: Arc<DelegationRegistry>) {
469    let report = registry.reconcile().await;
470    if !report.reaped.is_empty() {
471        tracing::info!(
472            "Reaped {} orphaned delegations: {}",
473            report.reaped.len(),
474            report.reaped.join(", ")
475        );
476    }
477    if report.removed > 0 {
478        tracing::debug!(
479            "Removed {} delegation records whose processes had exited",
480            report.removed
481        );
482    }
483    if report.stale_markers > 0 {
484        tracing::debug!(
485            "Removed {} conversation markers whose SCV process had exited",
486            report.stale_markers
487        );
488    }
489}
490
491/// Why a daemon control request failed.
492enum ControlFailure {
493    /// A delegation request the client can correct; the message is safe to show.
494    Delegation(String),
495    Component,
496}
497
498/// Apply a daemon control command, adding the instance's delegations.
499async fn daemon_control(
500    components: &Arc<Mutex<components::Components>>,
501    registry: &DelegationRegistry,
502    command: DaemonCommand,
503) -> std::result::Result<DaemonStatus, ControlFailure> {
504    let mut killed = Vec::new();
505    let listing = match &command {
506        DaemonCommand::Delegations { all } => Some(*all),
507        DaemonCommand::DelegationKill { handle, orphans } => {
508            if handle.is_none() && !orphans {
509                return Err(ControlFailure::Delegation(
510                    "name a delegation handle or ask for orphans".into(),
511                ));
512            }
513            if *orphans {
514                let report = registry.reconcile().await;
515                killed.extend(report.reaped);
516            }
517            if let Some(handle) = handle {
518                registry
519                    .kill(handle)
520                    .await
521                    .map_err(ControlFailure::Delegation)?;
522                killed.push(handle.clone());
523            }
524            Some(true)
525        }
526        _ => None,
527    };
528    let mut status = components
529        .lock()
530        .await
531        .control(command)
532        .await
533        .map_err(|_| ControlFailure::Component)?;
534    let running = registry.list(false);
535    status.delegations = DelegationSummary {
536        active: running.len() as u64,
537        reaped: registry.reaped_total(),
538        entries: match listing {
539            Some(true) => registry.list(true),
540            Some(false) => running,
541            None => Vec::new(),
542        }
543        .into_iter()
544        .map(|entry| DelegationInfo {
545            handle: entry.record.handle,
546            agent: entry.record.agent,
547            session: entry.record.session,
548            depth: entry.record.depth,
549            pid: entry.record.process.pid,
550            owner_pid: entry.record.owner.pid,
551            processes: u32::try_from(entry.processes).unwrap_or(u32::MAX),
552            cwd: entry.record.cwd.display().to_string(),
553            started_unix_seconds: entry.record.started_unix,
554            orphaned: entry.orphaned,
555            conversation: entry.record.conversation,
556            turn: entry.record.turn,
557        })
558        .collect(),
559        killed,
560    };
561    Ok(status)
562}
563
564async fn run_managed<R, W>(
565    reader: R,
566    writer: W,
567    overrides: ConfigOverrides,
568    components: Option<Arc<Mutex<components::Components>>>,
569    registry: Arc<DelegationRegistry>,
570    cancellation: CancellationToken,
571    tasks: TaskTracker,
572) -> Result<()>
573where
574    R: tokio::io::AsyncRead + Unpin,
575    W: tokio::io::AsyncWrite + Unpin + Send + 'static,
576{
577    let initial_output_bytes =
578        output_queue_bytes(Config::default().protocol.max_server_frame_bytes)?;
579    let (output_tx, mut output_rx) = outbound_channel(initial_output_bytes);
580    let mut writer_task = tasks.spawn(async move {
581        let mut writer = writer;
582        while let Some(frame) = output_rx.recv().await {
583            writer.write_all(&frame.bytes).await?;
584            writer.write_all(b"\n").await?;
585            writer.flush().await?;
586        }
587        Ok::<(), std::io::Error>(())
588    });
589    let _writer_abort = AbortGuard(writer_task.abort_handle());
590    let (done_tx, mut done_rx) = mpsc::channel::<TurnDone>(4);
591    let approvals = Arc::new(ApprovalBroker::default());
592    let mut reader = BufReader::new(reader);
593    let mut frames = FrameBuffer::default();
594    let mut initialized = false;
595    let mut session: Option<Session> = None;
596    let mut active: Option<ActiveTurn> = None;
597    // Woken when a background job finishes; set until its report turn starts.
598    let mut background_rx: Option<mpsc::UnboundedReceiver<()>> = None;
599    let mut background_ready = false;
600    let mut fatal = false;
601    let mut writer_finished = false;
602
603    let loop_result: Result<()> = async {
604        loop {
605        let frame_limit = session.as_ref().map_or_else(
606            || Config::default().protocol.max_client_frame_bytes,
607            |value| value.config.protocol.max_client_frame_bytes,
608        );
609        tokio::select! {
610            _ = cancellation.cancelled() => break,
611            read = frames.read(&mut reader, frame_limit) => {
612                let frame = match read.context("read protocol input")? {
613                    FrameRead::Eof => {
614                        if let Some(active) = &active { active.cancellation.cancel(); }
615                        break;
616                    }
617                    FrameRead::TooLarge => {
618                        send_error(&output_tx, "", "invalid_request", "client frame exceeds configured limit", false, server_frame_limit(&session)).await?;
619                        continue;
620                    }
621                    FrameRead::Frame(frame) => frame,
622                };
623                if frame.is_empty() {
624                    send_error(&output_tx, "", "invalid_json", "protocol frame is empty", false, server_frame_limit(&session)).await?;
625                    continue;
626                }
627                let message = match serde_json::from_slice::<ClientMessage>(&frame) {
628                    Ok(message) => message,
629                    Err(error) => {
630                        send_error(&output_tx, "", "invalid_json", &format!("invalid protocol JSON: {error}"), false, server_frame_limit(&session)).await?;
631                        continue;
632                    }
633                };
634                match message {
635                    ClientMessage::Initialize { request_id, protocol_version, .. } => {
636                        if initialized {
637                            send_error(&output_tx, &request_id, "invalid_request", "connection is already initialized", false, server_frame_limit(&session)).await?;
638                            continue;
639                        }
640                        if protocol_version != PROTOCOL_VERSION {
641                            send_error(&output_tx, &request_id, "version_mismatch", &format!("server supports protocol {PROTOCOL_VERSION}"), true, server_frame_limit(&session)).await?;
642                            fatal = true;
643                            break;
644                        }
645                        initialized = true;
646                        send_event(&output_tx, ServerEvent::Initialized {
647                            request_id,
648                            protocol_version: PROTOCOL_VERSION,
649                            server: PeerInfo { name: "scv-server".into(), version: env!("CARGO_PKG_VERSION").into() },
650                        }, Config::default().protocol.max_server_frame_bytes).await?;
651                    }
652                    other if !initialized => {
653                        send_error(&output_tx, other.request_id(), "not_initialized", "initialize must be the first message", false, server_frame_limit(&session)).await?;
654                    }
655                    ClientMessage::DaemonControl { request_id, command } => {
656                        if let Some(components) = &components {
657                            let result = tokio::select! {
658                                biased;
659                                _ = cancellation.cancelled() => break,
660                                result = daemon_control(components, &registry, command) => result,
661                            };
662                            match result {
663                                Ok(status) => send_event(&output_tx, ServerEvent::DaemonStatus { request_id, status }, server_frame_limit(&session)).await?,
664                                Err(ControlFailure::Delegation(message)) => send_error(&output_tx, &request_id, "delegation_error", &message, false, server_frame_limit(&session)).await?,
665                                Err(ControlFailure::Component) => send_error(&output_tx, &request_id, "component_error", "Component operation failed; check account credentials, private file permissions and absolute workspace", false, server_frame_limit(&session)).await?,
666                            }
667                        } else {
668                            send_error(&output_tx, &request_id, "unsupported", "Component management requires the daemon socket", false, server_frame_limit(&session)).await?;
669                        }
670                    }
671                    ClientMessage::SessionStart { request_id, cwd, provider, model, base_url, no_tools, delegation_depth } => {
672                        if session.is_some() {
673                            send_error(&output_tx, &request_id, "invalid_request", "this connection already has a session", false, server_frame_limit(&session)).await?;
674                            continue;
675                        }
676                        let session_overrides = ConfigOverrides {
677                            provider: provider.or_else(|| overrides.provider.clone()),
678                            model: model.or_else(|| overrides.model.clone()),
679                            base_url: base_url.or_else(|| overrides.base_url.clone()),
680                            approval_policy: overrides.approval_policy,
681                            no_tools: no_tools.unwrap_or(overrides.no_tools),
682                        };
683                        match build_session(&cwd, session_overrides, delegation_depth.unwrap_or(0), &registry).await {
684                            Ok((new_session, finished)) => {
685                                background_rx = finished;
686                                output_tx.ensure_capacity(output_queue_bytes(
687                                    new_session.config.protocol.max_server_frame_bytes,
688                                )?)?;
689                                let event = ServerEvent::SessionStarted {
690                                    request_id,
691                                    session_id: new_session.id.clone(),
692                                    cwd: new_session.workspace.display().to_string(),
693                                    model: new_session.runtime.model().to_owned(),
694                                    context_max_tokens: new_session.config.context.max_tokens,
695                                    max_server_frame_bytes: new_session.config.protocol.max_server_frame_bytes,
696                                    max_transcript_bytes: new_session.config.tui.max_transcript_bytes,
697                                    max_transcript_items: new_session.config.tui.max_transcript_items,
698                                    max_prompt_history_bytes: new_session.config.tui.max_prompt_history_bytes,
699                                    max_prompt_history_items: new_session.config.tui.max_prompt_history_items,
700                                };
701                                send_event(&output_tx, event, new_session.config.protocol.max_server_frame_bytes).await?;
702                                send_event(&output_tx, ServerEvent::QueueSnapshot {
703                                    request_id: None,
704                                    session_id: new_session.id.clone(),
705                                    seq: next_seq(&new_session.seq),
706                                    entries: new_session.queue.lock().await.iter().cloned().collect(),
707                                    paused: new_session.paused.load(Ordering::Acquire),
708                                }, new_session.config.protocol.max_server_frame_bytes).await?;
709                                session = Some(new_session);
710                            }
711                            Err(error) => {
712                                send_error(&output_tx, &request_id, "invalid_request", &error.to_string(), false, server_frame_limit(&session)).await?;
713                            }
714                        }
715                    }
716                    ClientMessage::SessionAttach { request_id, .. } => {
717                        send_error(&output_tx, &request_id, "unsupported", "session attach requires the shared socket server", false, server_frame_limit(&session)).await?;
718                    }
719                    ClientMessage::TurnStart { request_id, session_id, prompt } => {
720                        let Some(current) = session.as_ref() else {
721                            send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?;
722                            continue;
723                        };
724                        if current.id != session_id {
725                            send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
726                            continue;
727                        }
728                        if prompt.trim().is_empty() || prompt.len() > PROMPT_LIMIT_BYTES {
729                            send_error(&output_tx, &request_id, "invalid_request", "prompt must be non-empty and no larger than 256 KiB", false, server_frame_limit(&session)).await?;
730                            continue;
731                        }
732                        if active.is_some() {
733                            let entry = match current.enqueue(prompt, request_id.clone()).await {
734                                Ok(entry) => entry,
735                                Err(code) => { send_error(&output_tx, &request_id, code, "session queue limit reached", false, server_frame_limit(&session)).await?; continue; }
736                            };
737                            let position = current.queue.lock().await.len().saturating_sub(1);
738                            send_event(&output_tx, ServerEvent::QueueEnqueued {
739                                request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), entry, position,
740                            }, current.config.protocol.max_server_frame_bytes).await?;
741                            continue;
742                        }
743                        let starter = TurnStarter { output: &output_tx, approvals: &approvals, done: &done_tx, tasks: &tasks, cancellation: &cancellation };
744                        active = Some(starter.start(current, Uuid::new_v4().to_string(), request_id, prompt, None).await?);
745                    }
746                    ClientMessage::QueueUpdate { request_id, session_id, queue_id, revision, prompt } => {
747                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
748                        if current.id != session_id { send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?; continue; }
749                        if prompt.trim().is_empty() || prompt.len() > PROMPT_LIMIT_BYTES { send_error(&output_tx, &request_id, "invalid_request", "prompt must be non-empty and no larger than 256 KiB", false, server_frame_limit(&session)).await?; continue; }
750                        match current.update_queue(&queue_id, revision, prompt).await {
751                            Ok(entry) => send_event(&output_tx, ServerEvent::QueueUpdated { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), entry }, current.config.protocol.max_server_frame_bytes).await?,
752                            Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
753                        }
754                    }
755                    ClientMessage::QueueMove { request_id, session_id, queue_id, revision, before_queue_id } => {
756                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
757                        match current.move_queue(&session_id, &queue_id, revision, before_queue_id).await {
758                            Ok((id, rev, pos)) => send_event(&output_tx, ServerEvent::QueueMoved { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), queue_id: id, position: pos, revision: rev }, current.config.protocol.max_server_frame_bytes).await?,
759                            Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
760                        }
761                    }
762                    ClientMessage::QueueRemove { request_id, session_id, queue_id, revision } => {
763                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
764                        match current.remove_queue(&session_id, &queue_id, revision).await {
765                            Ok((id, rev)) => send_event(&output_tx, ServerEvent::QueueRemoved { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), queue_id: id, revision: rev }, current.config.protocol.max_server_frame_bytes).await?,
766                            Err(code) => send_error(&output_tx, &request_id, code, "queue entry was not found or revision is stale", false, server_frame_limit(&session)).await?,
767                        }
768                    }
769                    ClientMessage::SessionPause { request_id, session_id, paused } => {
770                        let Some(current) = session.as_ref() else { send_error(&output_tx, &request_id, "session_not_found", "start a session first", false, server_frame_limit(&session)).await?; continue; };
771                        if current.id != session_id { send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?; continue; }
772                        current.paused.store(paused, Ordering::Release);
773                        send_event(&output_tx, ServerEvent::SessionPaused { request_id, session_id: current.id.clone(), seq: next_seq(&current.seq), paused }, current.config.protocol.max_server_frame_bytes).await?;
774                    }
775                    ClientMessage::TurnCancel { request_id, session_id, turn_id } => {
776                        match (&session, &active) {
777                            (Some(current), Some(running)) if current.id == session_id && running.turn_id == turn_id => running.cancellation.cancel(),
778                            _ => send_error(&output_tx, &request_id, "turn_not_found", "active turn was not found", false, server_frame_limit(&session)).await?,
779                        }
780                    }
781                    ClientMessage::ApprovalResolve { request_id, session_id, approval_id, approved } => {
782                        if session.as_ref().is_none_or(|current| current.id != session_id) {
783                            send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
784                        } else if !approvals.resolve(&approval_id, approved).await {
785                            send_error(&output_tx, &request_id, "approval_not_found", "approval was not found or already resolved", false, server_frame_limit(&session)).await?;
786                        }
787                    }
788                    ClientMessage::SessionClear { request_id, session_id } => {
789                        let Some(current) = session.as_ref() else {
790                            send_error(&output_tx, &request_id, "session_not_found", "session was not found", false, server_frame_limit(&session)).await?;
791                            continue;
792                        };
793                        if current.id != session_id {
794                            send_error(&output_tx, &request_id, "session_not_found", "session id does not match", false, server_frame_limit(&session)).await?;
795                        } else if active.is_some() {
796                            send_error(&output_tx, &request_id, "turn_active", "cancel the active turn before clearing", false, server_frame_limit(&session)).await?;
797                        } else {
798                            current.history.lock().await.clear();
799                            current.queue.lock().await.clear();
800                            send_event(&output_tx, ServerEvent::SessionCleared {
801                                request_id,
802                                session_id: current.id.clone(),
803                                seq: next_seq(&current.seq),
804                            }, current.config.protocol.max_server_frame_bytes).await?;
805                            send_event(&output_tx, ServerEvent::QueueSnapshot { request_id: None, session_id: current.id.clone(), seq: next_seq(&current.seq), entries: Vec::new(), paused: current.paused.load(Ordering::Acquire) }, current.config.protocol.max_server_frame_bytes).await?;
806                        }
807                    }
808                }
809            }
810            writer = &mut writer_task => {
811                writer_finished = true;
812                writer.context("join protocol writer")??;
813                break;
814            }
815            Some(()) = recv_background(&mut background_rx) => {
816                background_ready = true;
817                if active.is_none()
818                    && let Some(current) = session.as_ref()
819                    && !current.paused.load(Ordering::Acquire)
820                    && current.queue.lock().await.is_empty()
821                {
822                    let starter = TurnStarter { output: &output_tx, approvals: &approvals, done: &done_tx, tasks: &tasks, cancellation: &cancellation };
823                    active = starter.report_background(current).await?;
824                    background_ready = active.is_some();
825                }
826            }
827            done = done_rx.recv(), if active.is_some() => {
828                if let Some(done) = done {
829                    if let Some(current) = session.as_ref() {
830                        let seq = next_seq(&current.seq);
831                        let event = match done.result {
832                            Ok(outcome) => ServerEvent::TurnCompleted {
833                                request_id: done.request_id,
834                                session_id: done.session_id,
835                                turn_id: done.turn_id,
836                                seq,
837                                steps: outcome.steps,
838                                usage: Usage { input_tokens: outcome.usage.input_tokens, output_tokens: outcome.usage.output_tokens },
839                                origin: done.origin,
840                            },
841                            Err(AgentError::Cancelled) => ServerEvent::TurnCancelled {
842                                request_id: done.request_id,
843                                session_id: done.session_id,
844                                turn_id: done.turn_id,
845                                seq,
846                                origin: done.origin,
847                            },
848                            Err(error) => ServerEvent::TurnFailed {
849                                request_id: done.request_id,
850                                session_id: done.session_id,
851                                turn_id: done.turn_id,
852                                seq,
853                                code: error.code().into(),
854                                message: error.to_string(),
855                                origin: done.origin,
856                            },
857                        };
858                        send_event(&output_tx, event, current.config.protocol.max_server_frame_bytes).await?;
859                    }
860                    if let Some(mut active) = active.take() {
861                        let _ = (&mut active.task).await;
862                    }
863                    if let Some(current) = session.as_ref()
864                        && !current.paused.load(Ordering::Acquire)
865                        && let Some(entry) = current.queue.lock().await.pop_front()
866                    {
867                        let turn_id = Uuid::new_v4().to_string();
868                        send_event(&output_tx, ServerEvent::QueueDequeued {
869                            request_id: entry.submitter.clone(),
870                            session_id: current.id.clone(),
871                            seq: next_seq(&current.seq),
872                            queue_id: entry.queue_id,
873                            turn_id: turn_id.clone(),
874                        }, current.config.protocol.max_server_frame_bytes).await?;
875                        let starter = TurnStarter { output: &output_tx, approvals: &approvals, done: &done_tx, tasks: &tasks, cancellation: &cancellation };
876                        active = Some(starter.start(current, turn_id, entry.submitter, entry.prompt, None).await?);
877                    }
878                    // Report finished background jobs once the user's own work is done.
879                    if active.is_none() && background_ready
880                        && let Some(current) = session.as_ref()
881                    {
882                        let starter = TurnStarter { output: &output_tx, approvals: &approvals, done: &done_tx, tasks: &tasks, cancellation: &cancellation };
883                        active = starter.report_background(current).await?;
884                        background_ready = active.is_some();
885                    }
886                }
887            }
888        }
889        }
890        Ok(())
891    }
892    .await;
893
894    if let Some(active) = active.take() {
895        shutdown_active_turn(active, SHUTDOWN_GRACE).await;
896    }
897    drop(output_tx);
898    let writer_result = if writer_finished {
899        Ok(())
900    } else {
901        shutdown_writer(writer_task, SHUTDOWN_GRACE).await
902    };
903    loop_result?;
904    writer_result?;
905    if fatal {
906        return Err(anyhow!("protocol version mismatch"));
907    }
908    Ok(())
909}
910
911enum FrameRead {
912    Eof,
913    Frame(Vec<u8>),
914    TooLarge,
915}
916
917struct OutboundFrame {
918    bytes: Vec<u8>,
919    _byte_permit: OwnedSemaphorePermit,
920}
921
922#[derive(Clone)]
923struct OutboundSender {
924    frames: mpsc::Sender<OutboundFrame>,
925    budget: Arc<Semaphore>,
926    capacity: Arc<AtomicUsize>,
927}
928
929#[derive(Debug, PartialEq, Eq)]
930enum OutboundSendError {
931    Cancelled,
932    Closed,
933    TimedOut,
934    FrameExceedsQueue { frame_bytes: usize, capacity: usize },
935}
936
937impl std::fmt::Display for OutboundSendError {
938    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
939        match self {
940            Self::Cancelled => formatter.write_str("outbound send cancelled"),
941            Self::Closed => formatter.write_str("protocol client disconnected"),
942            Self::TimedOut => formatter.write_str("outbound send timed out under backpressure"),
943            Self::FrameExceedsQueue {
944                frame_bytes,
945                capacity,
946            } => write!(
947                formatter,
948                "outbound frame uses {frame_bytes} bytes but queue capacity is {capacity} bytes"
949            ),
950        }
951    }
952}
953
954impl std::error::Error for OutboundSendError {}
955
956fn outbound_channel(capacity: usize) -> (OutboundSender, mpsc::Receiver<OutboundFrame>) {
957    let (frames, receiver) = mpsc::channel(OUTPUT_QUEUE_CAPACITY);
958    (
959        OutboundSender {
960            frames,
961            budget: Arc::new(Semaphore::new(capacity)),
962            capacity: Arc::new(AtomicUsize::new(capacity)),
963        },
964        receiver,
965    )
966}
967
968impl OutboundSender {
969    fn ensure_capacity(&self, required: usize) -> Result<()> {
970        if required > Semaphore::MAX_PERMITS {
971            return Err(anyhow!(
972                "outbound queue capacity {required} exceeds runtime limit {}",
973                Semaphore::MAX_PERMITS
974            ));
975        }
976        let current = self.capacity.load(Ordering::Acquire);
977        if required > current {
978            self.budget.add_permits(required - current);
979            self.capacity.store(required, Ordering::Release);
980        }
981        Ok(())
982    }
983
984    async fn send(
985        &self,
986        bytes: Vec<u8>,
987        cancellation: Option<&CancellationToken>,
988    ) -> std::result::Result<(), OutboundSendError> {
989        self.send_with_timeout(bytes, cancellation, SHUTDOWN_GRACE)
990            .await
991    }
992
993    async fn send_with_timeout(
994        &self,
995        bytes: Vec<u8>,
996        cancellation: Option<&CancellationToken>,
997        control_timeout: Duration,
998    ) -> std::result::Result<(), OutboundSendError> {
999        let frame_bytes =
1000            bytes
1001                .len()
1002                .checked_add(1)
1003                .ok_or(OutboundSendError::FrameExceedsQueue {
1004                    frame_bytes: usize::MAX,
1005                    capacity: self.capacity.load(Ordering::Acquire),
1006                })?;
1007        let capacity = self.capacity.load(Ordering::Acquire);
1008        let permits =
1009            u32::try_from(frame_bytes).map_err(|_| OutboundSendError::FrameExceedsQueue {
1010                frame_bytes,
1011                capacity,
1012            })?;
1013        if frame_bytes > capacity {
1014            return Err(OutboundSendError::FrameExceedsQueue {
1015                frame_bytes,
1016                capacity,
1017            });
1018        }
1019
1020        let control_deadline = tokio::time::Instant::now() + control_timeout;
1021        let acquire = Arc::clone(&self.budget).acquire_many_owned(permits);
1022        let permit = if let Some(cancellation) = cancellation {
1023            tokio::select! {
1024                biased;
1025                _ = cancellation.cancelled() => return Err(OutboundSendError::Cancelled),
1026                permit = acquire => permit.map_err(|_| OutboundSendError::Closed)?,
1027            }
1028        } else {
1029            tokio::time::timeout_at(control_deadline, acquire)
1030                .await
1031                .map_err(|_| OutboundSendError::TimedOut)?
1032                .map_err(|_| OutboundSendError::Closed)?
1033        };
1034        let frame = OutboundFrame {
1035            bytes,
1036            _byte_permit: permit,
1037        };
1038        if let Some(cancellation) = cancellation {
1039            tokio::select! {
1040                biased;
1041                _ = cancellation.cancelled() => Err(OutboundSendError::Cancelled),
1042                result = self.frames.send(frame) => result.map_err(|_| OutboundSendError::Closed),
1043            }
1044        } else {
1045            tokio::time::timeout_at(control_deadline, self.frames.send(frame))
1046                .await
1047                .map_err(|_| OutboundSendError::TimedOut)?
1048                .map_err(|_| OutboundSendError::Closed)
1049        }
1050    }
1051}
1052
1053fn output_queue_bytes(max_frame_bytes: usize) -> Result<usize> {
1054    let required = max_frame_bytes
1055        .checked_add(1)
1056        .and_then(|bytes| bytes.checked_mul(2))
1057        .ok_or_else(|| anyhow!("configured server frame limit is too large"))?
1058        .max(OUTPUT_QUEUE_MIN_BYTES);
1059    if required > Semaphore::MAX_PERMITS {
1060        return Err(anyhow!(
1061            "configured server frame limit requires an outbound queue larger than the runtime supports"
1062        ));
1063    }
1064    Ok(required)
1065}
1066
1067/// A persistent decoder keeps consumed partial bytes across select cancellation.
1068#[derive(Default)]
1069struct FrameBuffer {
1070    bytes: Vec<u8>,
1071    oversized: bool,
1072}
1073
1074impl FrameBuffer {
1075    async fn read<R>(&mut self, reader: &mut R, max_bytes: usize) -> std::io::Result<FrameRead>
1076    where
1077        R: AsyncBufRead + Unpin,
1078    {
1079        loop {
1080            let available = reader.fill_buf().await?;
1081            let eof = available.is_empty();
1082            let end = available.iter().position(|b| *b == b'\n');
1083            let take = end.map_or(available.len(), |n| n + 1);
1084            if !self.oversized {
1085                if self.bytes.len().saturating_add(take) > max_bytes.saturating_add(2) {
1086                    self.oversized = true;
1087                    self.bytes.clear();
1088                } else {
1089                    self.bytes.extend_from_slice(&available[..take]);
1090                }
1091            }
1092            reader.consume(take);
1093            if end.is_some() || eof {
1094                if std::mem::take(&mut self.oversized) {
1095                    return Ok(FrameRead::TooLarge);
1096                }
1097                if eof && self.bytes.is_empty() {
1098                    return Ok(FrameRead::Eof);
1099                }
1100                let mut bytes = std::mem::take(&mut self.bytes);
1101                while matches!(bytes.last(), Some(b'\n' | b'\r')) {
1102                    bytes.pop();
1103                }
1104                return Ok(if bytes.len() > max_bytes {
1105                    FrameRead::TooLarge
1106                } else {
1107                    FrameRead::Frame(bytes)
1108                });
1109            }
1110        }
1111    }
1112}
1113
1114#[cfg(test)]
1115async fn read_bounded_frame<R: AsyncBufRead + Unpin>(
1116    reader: &mut R,
1117    max_bytes: usize,
1118) -> std::io::Result<FrameRead> {
1119    FrameBuffer::default().read(reader, max_bytes).await
1120}
1121
1122/// A persistent advisory lock closes the stale-socket unlink/bind race.
1123struct SocketLock(std::fs::File);
1124impl SocketLock {
1125    fn acquire(socket: &Path) -> Result<Self> {
1126        use std::os::unix::{fs::OpenOptionsExt, io::AsRawFd};
1127        let file = std::fs::OpenOptions::new()
1128            .read(true)
1129            .write(true)
1130            .create(true)
1131            .truncate(false)
1132            .mode(0o600)
1133            .custom_flags(libc::O_NOFOLLOW)
1134            .open(socket.with_extension("lock"))?;
1135        // SAFETY: flock operates on this owned, live file descriptor.
1136        if unsafe { libc::flock(file.as_raw_fd(), libc::LOCK_EX | libc::LOCK_NB) } != 0 {
1137            return Err(anyhow!("SCV daemon already owns this socket"));
1138        }
1139        Ok(Self(file))
1140    }
1141}
1142impl Drop for SocketLock {
1143    fn drop(&mut self) {
1144        use std::os::unix::io::AsRawFd;
1145        // SAFETY: the descriptor remains live until this drop returns.
1146        unsafe {
1147            libc::flock(self.0.as_raw_fd(), libc::LOCK_UN);
1148        }
1149    }
1150}
1151
1152fn server_frame_limit(session: &Option<Session>) -> usize {
1153    session.as_ref().map_or_else(
1154        || Config::default().protocol.max_server_frame_bytes,
1155        |value| value.config.protocol.max_server_frame_bytes,
1156    )
1157}
1158
1159struct Session {
1160    id: String,
1161    workspace: PathBuf,
1162    config: Config,
1163    runtime: Arc<AgentRuntime>,
1164    history: Arc<Mutex<Vec<Message>>>,
1165    seq: Arc<AtomicU64>,
1166    queue: Arc<Mutex<VecDeque<QueueEntry>>>,
1167    paused: Arc<std::sync::atomic::AtomicBool>,
1168    /// Background agent jobs, shared with the session's tools.
1169    background: Option<Arc<BackgroundJobs>>,
1170}
1171
1172impl Session {
1173    async fn enqueue(
1174        &self,
1175        prompt: String,
1176        submitter: String,
1177    ) -> std::result::Result<QueueEntry, &'static str> {
1178        let entry = QueueEntry {
1179            queue_id: Uuid::new_v4().to_string(),
1180            revision: 1,
1181            prompt,
1182            submitter,
1183        };
1184        let mut queue = self.queue.lock().await;
1185        let bytes: usize = queue.iter().map(|item| item.prompt.len()).sum();
1186        if queue.len() >= MAX_QUEUE_ITEMS
1187            || bytes.saturating_add(entry.prompt.len()) > MAX_QUEUE_BYTES
1188        {
1189            return Err("queue_limit");
1190        }
1191        queue.push_back(entry.clone());
1192        Ok(entry)
1193    }
1194
1195    async fn update_queue(
1196        &self,
1197        id: &str,
1198        revision: u64,
1199        prompt: String,
1200    ) -> std::result::Result<QueueEntry, &'static str> {
1201        let mut queue = self.queue.lock().await;
1202        let bytes: usize = queue.iter().map(|item| item.prompt.len()).sum();
1203        let entry = queue
1204            .iter_mut()
1205            .find(|entry| entry.queue_id == id)
1206            .ok_or("queue_not_found")?;
1207        if entry.revision != revision {
1208            return Err("queue_conflict");
1209        }
1210        if bytes
1211            .saturating_sub(entry.prompt.len())
1212            .saturating_add(prompt.len())
1213            > MAX_QUEUE_BYTES
1214        {
1215            return Err("queue_limit");
1216        }
1217        entry.prompt = prompt;
1218        entry.revision += 1;
1219        Ok(entry.clone())
1220    }
1221
1222    async fn move_queue(
1223        &self,
1224        session_id: &str,
1225        id: &str,
1226        revision: u64,
1227        before: Option<String>,
1228    ) -> std::result::Result<(String, u64, usize), &'static str> {
1229        if self.id != session_id {
1230            return Err("session_not_found");
1231        }
1232        let mut queue = self.queue.lock().await;
1233        let index = queue
1234            .iter()
1235            .position(|entry| entry.queue_id == id)
1236            .ok_or("queue_not_found")?;
1237        if queue[index].revision != revision {
1238            return Err("queue_conflict");
1239        }
1240        // Validate the destination while the source is still present. This keeps
1241        // the operation atomic and handles a self move as a no-op reorder.
1242        let target_index = match before.as_deref() {
1243            Some(target) if target == id => return Ok((id.to_string(), revision, index)),
1244            Some(target) => Some(
1245                queue
1246                    .iter()
1247                    .position(|item| item.queue_id == target)
1248                    .ok_or("queue_not_found")?,
1249            ),
1250            None => None,
1251        };
1252        let mut entry = queue.remove(index).expect("queue index exists");
1253        let target = target_index.map_or(queue.len(), |target| {
1254            target.saturating_sub(usize::from(target > index))
1255        });
1256        let pos = target.min(queue.len());
1257        let id = entry.queue_id.clone();
1258        let rev = entry.revision + 1;
1259        entry.revision = rev;
1260        queue.insert(pos, entry);
1261        Ok((id, rev, pos))
1262    }
1263
1264    async fn remove_queue(
1265        &self,
1266        session_id: &str,
1267        id: &str,
1268        revision: u64,
1269    ) -> std::result::Result<(String, u64), &'static str> {
1270        if self.id != session_id {
1271            return Err("session_not_found");
1272        }
1273        let mut queue = self.queue.lock().await;
1274        let index = queue
1275            .iter()
1276            .position(|entry| entry.queue_id == id)
1277            .ok_or("queue_not_found")?;
1278        if queue[index].revision != revision {
1279            return Err("queue_conflict");
1280        }
1281        let entry = queue.remove(index).expect("queue index exists");
1282        Ok((entry.queue_id, entry.revision))
1283    }
1284}
1285
1286struct ActiveTurn {
1287    turn_id: String,
1288    cancellation: CancellationToken,
1289    task: JoinHandle<()>,
1290}
1291
1292impl Drop for ActiveTurn {
1293    fn drop(&mut self) {
1294        self.cancellation.cancel();
1295        self.task.abort();
1296    }
1297}
1298
1299struct AbortGuard(tokio::task::AbortHandle);
1300impl Drop for AbortGuard {
1301    fn drop(&mut self) {
1302        self.0.abort();
1303    }
1304}
1305
1306async fn shutdown_active_turn(mut active: ActiveTurn, grace: Duration) -> bool {
1307    active.cancellation.cancel();
1308    if tokio::time::timeout(grace, &mut active.task).await.is_ok() {
1309        true
1310    } else {
1311        active.task.abort();
1312        let _ = (&mut active.task).await;
1313        false
1314    }
1315}
1316
1317async fn shutdown_writer(
1318    mut writer: JoinHandle<std::io::Result<()>>,
1319    grace: Duration,
1320) -> Result<()> {
1321    match tokio::time::timeout(grace, &mut writer).await {
1322        Ok(result) => {
1323            result.context("join protocol writer")??;
1324            Ok(())
1325        }
1326        Err(_) => {
1327            writer.abort();
1328            let _ = writer.await;
1329            Err(anyhow!("protocol writer shutdown timed out"))
1330        }
1331    }
1332}
1333
1334struct TurnDone {
1335    request_id: String,
1336    session_id: String,
1337    turn_id: String,
1338    origin: Option<TurnOrigin>,
1339    result: Result<scv_core::TurnOutcome, AgentError>,
1340}
1341
1342/// What a connection needs to start a turn in its session.
1343struct TurnStarter<'a> {
1344    output: &'a OutboundSender,
1345    approvals: &'a Arc<ApprovalBroker>,
1346    done: &'a mpsc::Sender<TurnDone>,
1347    tasks: &'a TaskTracker,
1348    cancellation: &'a CancellationToken,
1349}
1350
1351impl TurnStarter<'_> {
1352    /// Announce and run one turn of `current` for `prompt`.
1353    async fn start(
1354        &self,
1355        current: &Session,
1356        turn_id: String,
1357        request_id: String,
1358        prompt: String,
1359        origin: Option<TurnOrigin>,
1360    ) -> Result<ActiveTurn> {
1361        let cancellation = self.cancellation.child_token();
1362        send_event(
1363            self.output,
1364            ServerEvent::TurnStarted {
1365                request_id: request_id.clone(),
1366                session_id: current.id.clone(),
1367                turn_id: turn_id.clone(),
1368                seq: next_seq(&current.seq),
1369                origin: origin.clone(),
1370            },
1371            current.config.protocol.max_server_frame_bytes,
1372        )
1373        .await?;
1374        let meta = TurnMeta {
1375            request_id: request_id.clone(),
1376            session_id: current.id.clone(),
1377            turn_id: turn_id.clone(),
1378            seq: Arc::clone(&current.seq),
1379            max_server_frame: current.config.protocol.max_server_frame_bytes,
1380        };
1381        let sink: Arc<dyn EventSink> = Arc::new(ProtocolSink {
1382            meta: meta.clone(),
1383            output: self.output.clone(),
1384            cancellation: cancellation.clone(),
1385        });
1386        let gate: Arc<dyn ApprovalGate> = Arc::new(ProtocolApprovalGate {
1387            policy: current.config.tools.approval_policy,
1388            broker: Arc::clone(self.approvals),
1389            meta,
1390            output: self.output.clone(),
1391        });
1392        let runtime = Arc::clone(&current.runtime);
1393        let history = Arc::clone(&current.history);
1394        let done = self.done.clone();
1395        let session_id = current.id.clone();
1396        let task_turn = turn_id.clone();
1397        let task_cancel = cancellation.clone();
1398        let task = self.tasks.spawn(async move {
1399            let mut history = history.lock().await;
1400            let result = runtime
1401                .run_turn(&mut history, prompt, sink, gate, task_cancel)
1402                .await;
1403            let _ = done
1404                .send(TurnDone {
1405                    request_id,
1406                    session_id,
1407                    turn_id: task_turn,
1408                    origin,
1409                    result,
1410                })
1411                .await;
1412        });
1413        Ok(ActiveTurn {
1414            turn_id,
1415            cancellation,
1416            task,
1417        })
1418    }
1419
1420    /// Start a turn reporting background jobs the model has not seen yet,
1421    /// or `None` when every finished job was already seen.
1422    async fn report_background(&self, current: &Session) -> Result<Option<ActiveTurn>> {
1423        let Some(jobs) = &current.background else {
1424            return Ok(None);
1425        };
1426        let reports = jobs.take_unreported();
1427        if reports.is_empty() {
1428            return Ok(None);
1429        }
1430        let origin = TurnOrigin {
1431            kind: ORIGIN_BACKGROUND.into(),
1432            jobs: reports.iter().map(|report| report.job.clone()).collect(),
1433        };
1434        let prompt = background::report_prompt(&reports);
1435        let request_id = format!("background:{}", Uuid::new_v4());
1436        self.start(
1437            current,
1438            Uuid::new_v4().to_string(),
1439            request_id,
1440            prompt,
1441            Some(origin),
1442        )
1443        .await
1444        .map(Some)
1445    }
1446}
1447
1448/// The next background-job wake-up, or never without a receiver.
1449async fn recv_background(receiver: &mut Option<mpsc::UnboundedReceiver<()>>) -> Option<()> {
1450    match receiver {
1451        Some(receiver) => receiver.recv().await,
1452        None => std::future::pending().await,
1453    }
1454}
1455
1456#[derive(Clone)]
1457struct TurnMeta {
1458    request_id: String,
1459    session_id: String,
1460    turn_id: String,
1461    seq: Arc<AtomicU64>,
1462    max_server_frame: usize,
1463}
1464
1465fn next_seq(sequence: &AtomicU64) -> u64 {
1466    sequence.fetch_add(1, Ordering::Relaxed) + 1
1467}
1468
1469/// `delegation_depth` is the depth the client declared in `session.start`
1470/// (0 for a direct client); the session's delegated runs count from it.
1471async fn build_session(
1472    cwd: &str,
1473    overrides: ConfigOverrides,
1474    delegation_depth: u32,
1475    registry: &Arc<DelegationRegistry>,
1476) -> Result<(Session, Option<mpsc::UnboundedReceiver<()>>)> {
1477    let id = Uuid::new_v4().to_string();
1478    let workspace = std::fs::canonicalize(cwd).with_context(|| format!("resolve cwd {cwd}"))?;
1479    if !workspace.is_dir() {
1480        return Err(anyhow!("cwd is not a directory"));
1481    }
1482    let no_tools = overrides.no_tools;
1483    let config = Config::load(&workspace, overrides)?;
1484    if !no_tools {
1485        config.prepare_adapter_homes()?;
1486    }
1487    let provider_config = config.provider.clone();
1488    let api_key = provider_config.api_key.clone().or_else(|| {
1489        provider_config.api_key_env.as_deref().and_then(|name| std::env::var(name).ok())
1490    }).filter(|key| !key.trim().is_empty()).ok_or_else(|| anyhow!("provider credential is not configured; set provider.api_key or provider.api_key_env"))?;
1491    let skills = discover_skills(&workspace, &config, !no_tools)?;
1492    let system_prompt = build_system_prompt(&workspace, &config, &skills)?;
1493    let mut provider = OpenAiProvider::new(
1494        provider_config.model.clone(),
1495        provider_config.base_url.clone(),
1496        api_key,
1497        Duration::from_secs(provider_config.timeout_seconds),
1498        config.provider_limits(),
1499        provider_config.headers.clone(),
1500    )?;
1501    if !no_tools && config.hosted_web_search() {
1502        provider = provider.with_web_search();
1503    }
1504    let provider = Arc::new(provider);
1505    let (background, finished) = if !no_tools && config.agent.max_background > 0 {
1506        let (finished_tx, finished_rx) = mpsc::unbounded_channel();
1507        (
1508            Some(Arc::new(BackgroundJobs::new(
1509                config.agent.max_background,
1510                Some(finished_tx),
1511            ))),
1512            Some(finished_rx),
1513        )
1514    } else {
1515        (None, None)
1516    };
1517    let tools = if no_tools {
1518        Arc::new(ToolRegistry::default())
1519    } else {
1520        let mut tools = config.tools();
1521        tools.delegation = Some(DelegationContext {
1522            registry: Arc::clone(registry),
1523            session: id.clone(),
1524            depth: delegation_depth,
1525        });
1526        tools.background = background.clone();
1527        let mut registry = builtin_registry(
1528            tools,
1529            skills.map,
1530            skills.roots,
1531            config.skills.max_skill_bytes,
1532            config.adapters(),
1533        )?;
1534        if let Some(web) = config.web_tools() {
1535            scv_tools::web::register(&mut registry, web)?;
1536        }
1537        Arc::new(registry)
1538    };
1539    let context = Arc::new(BudgetContextPolicy::new((&config.context).into())?);
1540    let runtime = Arc::new(AgentRuntime::new(
1541        provider,
1542        tools,
1543        context,
1544        config.core_agent(system_prompt),
1545        workspace.clone(),
1546    ));
1547    Ok((
1548        Session {
1549            id,
1550            workspace,
1551            config,
1552            runtime,
1553            history: Arc::new(Mutex::new(Vec::new())),
1554            seq: Arc::new(AtomicU64::new(0)),
1555            queue: Arc::new(Mutex::new(VecDeque::new())),
1556            paused: Arc::new(std::sync::atomic::AtomicBool::new(false)),
1557            background,
1558        },
1559        finished,
1560    ))
1561}
1562
1563fn build_system_prompt(
1564    workspace: &Path,
1565    config: &Config,
1566    skills: &DiscoveredSkills,
1567) -> Result<String> {
1568    let mut prompt = config.agent.system_prompt.clone();
1569    prompt.push_str(&format!(
1570        "\nCurrent working directory: {}\n",
1571        workspace.display()
1572    ));
1573    let agents_path = workspace.join("AGENTS.md");
1574    if agents_path.is_file() {
1575        let canonical = std::fs::canonicalize(&agents_path).context("resolve project AGENTS.md")?;
1576        if !canonical.starts_with(workspace) {
1577            return Err(anyhow!("project AGENTS.md escaped workspace"));
1578        }
1579        let (bytes, truncated) = read_prefix(&canonical, config.tools.max_read_bytes)
1580            .context("read project AGENTS.md")?;
1581        let instructions = std::str::from_utf8(&bytes).context("project AGENTS.md is not UTF-8")?;
1582        prompt.push_str("\n# Project instructions\n");
1583        prompt.push_str(instructions);
1584        if truncated {
1585            prompt.push_str("\n[AGENTS.md truncated by configured read limit]\n");
1586        }
1587    }
1588    if !skills.listing.is_empty() {
1589        prompt.push_str("\n# Available skills\n");
1590        prompt.push_str(&skills.listing);
1591        prompt.push_str("\nUse read_skill with a skill name when its workflow applies.\n");
1592    }
1593    if !skills.project_listing.is_empty() {
1594        prompt.push_str("\n# Project skills\n");
1595        prompt.push_str(
1596            "Projects in this workspace provide these skills to agents working in them:\n",
1597        );
1598        prompt.push_str(&skills.project_listing);
1599        prompt.push_str(
1600            "\nTo use one, delegate with an agent_* tool such as agent_codex or agent_claude, \
1601             set its cwd to the skill's project, and name the skill in the prompt: that agent \
1602             then loads the project's instructions and skills itself. read_skill loads a \
1603             skill for reference.\n",
1604        );
1605    }
1606    Ok(prompt)
1607}
1608
1609/// Skills found at session start: the names `read_skill` serves, the roots it
1610/// revalidates them against, and their system-prompt listings.
1611struct DiscoveredSkills {
1612    map: SkillMap,
1613    roots: Vec<PathBuf>,
1614    listing: String,
1615    project_listing: String,
1616}
1617
1618/// Agent-native skill directories, relative to a project, that Codex and
1619/// Claude Code load from their working directory.
1620const PROJECT_SKILL_DIRS: [&str; 2] = [".agents/skills", ".claude/skills"];
1621/// Workspace entries and child projects inspected for project skills, so a
1622/// large workspace such as a home directory costs bounded lookups.
1623const MAX_WORKSPACE_ENTRIES: usize = 4096;
1624const MAX_SKILL_PROJECTS: usize = 256;
1625/// Bytes read to find a project skill's description, and its listed length.
1626const PROJECT_SKILL_HEADER_BYTES: usize = 16 * 1024;
1627const MAX_PROJECT_SKILL_DESCRIPTION: usize = 400;
1628
1629fn discover_skills(workspace: &Path, config: &Config, tools: bool) -> Result<DiscoveredSkills> {
1630    let mut skills = SkillMap::new();
1631    let mut roots = Vec::new();
1632    let project_root = workspace.join(&config.skills.project_dir);
1633    for (root, must_be_workspace) in [(&project_root, true), (&config.skills.user_dir, false)] {
1634        if !root.is_dir() {
1635            continue;
1636        }
1637        let canonical = std::fs::canonicalize(root)
1638            .with_context(|| format!("resolve skill root {}", root.display()))?;
1639        if must_be_workspace && !canonical.starts_with(workspace) {
1640            return Err(anyhow!("project skill root escaped workspace"));
1641        }
1642        roots.push(canonical.clone());
1643        let mut entries: Vec<_> = std::fs::read_dir(&canonical)
1644            .with_context(|| format!("read skill root {}", canonical.display()))?
1645            .filter_map(Result::ok)
1646            .collect();
1647        entries.sort_by_key(|entry| entry.file_name());
1648        for entry in entries {
1649            if skills.len() >= config.skills.max_skills {
1650                break;
1651            }
1652            let path = entry.path().join("SKILL.md");
1653            if !path.is_file() {
1654                continue;
1655            }
1656            let canonical_file = std::fs::canonicalize(&path)
1657                .with_context(|| format!("resolve skill {}", path.display()))?;
1658            if !canonical_file.starts_with(&canonical) {
1659                continue;
1660            }
1661            let name = entry.file_name().to_string_lossy().to_string();
1662            skills.entry(name).or_insert(canonical_file);
1663        }
1664    }
1665    let mut names: Vec<_> = skills.keys().cloned().collect();
1666    names.sort();
1667    let mut listing = String::new();
1668    for name in names {
1669        let path = &skills[&name];
1670        let bytes = read_prefix(path, config.skills.max_skill_bytes)
1671            .map(|(bytes, _)| bytes)
1672            .unwrap_or_default();
1673        let content = String::from_utf8_lossy(&bytes);
1674        let description = skill_description(&content);
1675        listing.push_str(&format!("- {name}: {description}\n"));
1676    }
1677    // Project skills are only actionable by delegating, so tool-free sessions
1678    // neither list them nor learn the workspace's project names.
1679    let project_listing = if tools && config.skills.scan_projects {
1680        discover_project_skills(workspace, config, &mut skills, &mut roots)
1681    } else {
1682        String::new()
1683    };
1684    Ok(DiscoveredSkills {
1685        map: skills,
1686        roots,
1687        listing,
1688        project_listing,
1689    })
1690}
1691
1692/// List the agent skills of the workspace and its immediate, non-hidden child
1693/// projects. A child's skills are named `<project>:<skill>`. Everything
1694/// resolves inside the workspace, SKILL.md files that resolve to the same file
1695/// (such as a `.claude/skills` link to `.agents/skills`) count once, and
1696/// unreadable entries are skipped so one broken project cannot stop a session.
1697fn discover_project_skills(
1698    workspace: &Path,
1699    config: &Config,
1700    skills: &mut SkillMap,
1701    roots: &mut Vec<PathBuf>,
1702) -> String {
1703    let mut projects = vec![(None, workspace.to_path_buf())];
1704    let mut names: Vec<_> = std::fs::read_dir(workspace)
1705        .into_iter()
1706        .flatten()
1707        .filter_map(Result::ok)
1708        .take(MAX_WORKSPACE_ENTRIES)
1709        .map(|entry| entry.file_name().to_string_lossy().into_owned())
1710        .filter(|name| !name.starts_with('.'))
1711        .collect();
1712    names.sort();
1713    for name in names {
1714        if projects.len() > MAX_SKILL_PROJECTS {
1715            break;
1716        }
1717        let Ok(directory) = std::fs::canonicalize(workspace.join(&name)) else {
1718            continue;
1719        };
1720        if directory.is_dir()
1721            && directory.starts_with(workspace)
1722            && !projects.iter().any(|(_, seen)| seen == &directory)
1723        {
1724            projects.push((Some(name), directory));
1725        }
1726    }
1727    let mut seen_files = std::collections::HashSet::new();
1728    let mut listing = String::new();
1729    'projects: for (project, directory) in projects {
1730        let project_roots: Vec<PathBuf> = PROJECT_SKILL_DIRS
1731            .iter()
1732            .filter_map(|relative| std::fs::canonicalize(directory.join(relative)).ok())
1733            .filter(|root| root.is_dir() && root.starts_with(workspace))
1734            .collect();
1735        for root in &project_roots {
1736            if !roots.contains(root) {
1737                roots.push(root.clone());
1738            }
1739        }
1740        for root in &project_roots {
1741            let mut entries: Vec<_> = std::fs::read_dir(root)
1742                .into_iter()
1743                .flatten()
1744                .filter_map(Result::ok)
1745                .take(MAX_WORKSPACE_ENTRIES)
1746                .collect();
1747            entries.sort_by_key(|entry| entry.file_name());
1748            for entry in entries {
1749                if skills.len() >= config.skills.max_skills {
1750                    break 'projects;
1751                }
1752                let Ok(file) = std::fs::canonicalize(entry.path().join("SKILL.md")) else {
1753                    continue;
1754                };
1755                if !file.is_file()
1756                    || !project_roots.iter().any(|root| file.starts_with(root))
1757                    || !seen_files.insert(file.clone())
1758                {
1759                    continue;
1760                }
1761                let skill = entry.file_name().to_string_lossy().into_owned();
1762                let (name, location) = match &project {
1763                    Some(project) => (format!("{project}:{skill}"), format!("project {project}")),
1764                    None => (skill, "workspace root".to_owned()),
1765                };
1766                // SCV's own and the user's skills keep their names.
1767                if skills.contains_key(&name) {
1768                    continue;
1769                }
1770                let header = read_prefix(
1771                    &file,
1772                    config
1773                        .skills
1774                        .max_skill_bytes
1775                        .min(PROJECT_SKILL_HEADER_BYTES),
1776                )
1777                .map(|(bytes, _)| bytes)
1778                .unwrap_or_default();
1779                let description: String = skill_description(&String::from_utf8_lossy(&header))
1780                    .chars()
1781                    .take(MAX_PROJECT_SKILL_DESCRIPTION)
1782                    .collect();
1783                listing.push_str(&format!("- {name} ({location}): {description}\n"));
1784                skills.insert(name, file);
1785            }
1786        }
1787    }
1788    listing
1789}
1790
1791fn read_prefix(path: &Path, max_bytes: usize) -> std::io::Result<(Vec<u8>, bool)> {
1792    let file = std::fs::File::open(path)?;
1793    let mut bytes = Vec::with_capacity(max_bytes.min(8192));
1794    file.take(
1795        u64::try_from(max_bytes)
1796            .unwrap_or(u64::MAX)
1797            .saturating_add(1),
1798    )
1799    .read_to_end(&mut bytes)?;
1800    let truncated = bytes.len() > max_bytes;
1801    bytes.truncate(max_bytes);
1802    Ok((bytes, truncated))
1803}
1804
1805fn skill_description(content: &str) -> String {
1806    if let Some(frontmatter) = content.strip_prefix("---\n")
1807        && let Some((header, _)) = frontmatter.split_once("\n---")
1808    {
1809        for line in header.lines() {
1810            if let Some(description) = line.strip_prefix("description:") {
1811                return description.trim().trim_matches('"').to_owned();
1812            }
1813        }
1814    }
1815    content
1816        .lines()
1817        .map(str::trim)
1818        .find(|line| !line.is_empty() && !line.starts_with('#'))
1819        .unwrap_or("No description provided")
1820        .chars()
1821        .take(240)
1822        .collect()
1823}
1824
1825/// Cut progress text to `MAX_PROGRESS_EVENT_BYTES` on a character boundary.
1826fn bounded_progress(mut text: String) -> String {
1827    let limit = scv_core::MAX_PROGRESS_EVENT_BYTES;
1828    if text.len() > limit {
1829        let mut end = limit;
1830        while !text.is_char_boundary(end) {
1831            end -= 1;
1832        }
1833        text.truncate(end);
1834    }
1835    text
1836}
1837
1838struct ProtocolSink {
1839    meta: TurnMeta,
1840    output: OutboundSender,
1841    cancellation: CancellationToken,
1842}
1843
1844#[async_trait]
1845impl EventSink for ProtocolSink {
1846    async fn emit(&self, event: CoreEvent) -> Result<(), AgentError> {
1847        let seq = next_seq(&self.meta.seq);
1848        let event = match event {
1849            CoreEvent::AssistantDelta { content } => ServerEvent::AssistantDelta {
1850                request_id: self.meta.request_id.clone(),
1851                session_id: self.meta.session_id.clone(),
1852                turn_id: self.meta.turn_id.clone(),
1853                seq,
1854                content,
1855            },
1856            CoreEvent::AssistantCompleted { content } => ServerEvent::AssistantCompleted {
1857                request_id: self.meta.request_id.clone(),
1858                session_id: self.meta.session_id.clone(),
1859                turn_id: self.meta.turn_id.clone(),
1860                seq,
1861                content,
1862            },
1863            CoreEvent::ToolProposed {
1864                call_id,
1865                name,
1866                arguments,
1867            } => ServerEvent::ToolProposed {
1868                request_id: self.meta.request_id.clone(),
1869                session_id: self.meta.session_id.clone(),
1870                turn_id: self.meta.turn_id.clone(),
1871                seq,
1872                call_id,
1873                name,
1874                arguments,
1875            },
1876            CoreEvent::ToolStarted { call_id, name } => ServerEvent::ToolStarted {
1877                request_id: self.meta.request_id.clone(),
1878                session_id: self.meta.session_id.clone(),
1879                turn_id: self.meta.turn_id.clone(),
1880                seq,
1881                call_id,
1882                name,
1883            },
1884            CoreEvent::ToolProgress { call_id, text } => ServerEvent::ToolProgress {
1885                request_id: self.meta.request_id.clone(),
1886                session_id: self.meta.session_id.clone(),
1887                turn_id: self.meta.turn_id.clone(),
1888                seq,
1889                call_id,
1890                // Tools report through a bounded sink; bound again here so a
1891                // custom tool can never grow a frame past the documented size.
1892                text: bounded_progress(text),
1893            },
1894            CoreEvent::ToolCompleted {
1895                call_id,
1896                name,
1897                output,
1898            } => ServerEvent::ToolCompleted {
1899                request_id: self.meta.request_id.clone(),
1900                session_id: self.meta.session_id.clone(),
1901                turn_id: self.meta.turn_id.clone(),
1902                seq,
1903                call_id,
1904                name,
1905                success: !output.is_error,
1906                output: output.content,
1907                truncated: output.truncated,
1908            },
1909            CoreEvent::ContextCompacted {
1910                before_tokens,
1911                after_tokens,
1912                removed_messages,
1913            } => ServerEvent::ContextCompacted {
1914                request_id: self.meta.request_id.clone(),
1915                session_id: self.meta.session_id.clone(),
1916                turn_id: self.meta.turn_id.clone(),
1917                seq,
1918                before_tokens,
1919                after_tokens,
1920                removed_messages,
1921            },
1922            CoreEvent::SessionTrimmed {
1923                removed_messages,
1924                history_bytes,
1925            } => ServerEvent::SessionTrimmed {
1926                request_id: self.meta.request_id.clone(),
1927                session_id: self.meta.session_id.clone(),
1928                seq,
1929                removed_messages,
1930                history_bytes,
1931            },
1932        };
1933        send_turn_event(
1934            &self.output,
1935            event,
1936            self.meta.max_server_frame,
1937            &self.cancellation,
1938        )
1939        .await
1940    }
1941}
1942
1943#[derive(Default)]
1944struct ApprovalBroker {
1945    pending: Mutex<HashMap<String, oneshot::Sender<bool>>>,
1946}
1947
1948impl ApprovalBroker {
1949    async fn insert(&self, id: String, sender: oneshot::Sender<bool>) {
1950        self.pending.lock().await.insert(id, sender);
1951    }
1952
1953    async fn remove(&self, id: &str) {
1954        self.pending.lock().await.remove(id);
1955    }
1956
1957    async fn resolve(&self, id: &str, approved: bool) -> bool {
1958        let sender = self.pending.lock().await.remove(id);
1959        sender.is_some_and(|sender| sender.send(approved).is_ok())
1960    }
1961}
1962
1963struct ProtocolApprovalGate {
1964    policy: ApprovalPolicy,
1965    broker: Arc<ApprovalBroker>,
1966    meta: TurnMeta,
1967    output: OutboundSender,
1968}
1969
1970#[async_trait]
1971impl ApprovalGate for ProtocolApprovalGate {
1972    async fn approve(
1973        &self,
1974        request: ApprovalRequest,
1975        cancellation: CancellationToken,
1976    ) -> Result<bool, AgentError> {
1977        match self.policy {
1978            ApprovalPolicy::OnRisk if request.risk == ToolRisk::ReadOnly => return Ok(true),
1979            ApprovalPolicy::Never => return Ok(request.risk == ToolRisk::ReadOnly),
1980            ApprovalPolicy::Always | ApprovalPolicy::OnRisk => {}
1981        }
1982        let approval_id = Uuid::new_v4().to_string();
1983        let (sender, receiver) = oneshot::channel();
1984        self.broker.insert(approval_id.clone(), sender).await;
1985        let event = ServerEvent::ApprovalRequested {
1986            request_id: self.meta.request_id.clone(),
1987            session_id: self.meta.session_id.clone(),
1988            turn_id: self.meta.turn_id.clone(),
1989            seq: next_seq(&self.meta.seq),
1990            approval_id: approval_id.clone(),
1991            call_id: request.call_id,
1992            name: request.name,
1993            risk: request.risk.as_str().into(),
1994            cwd: request.cwd.display().to_string(),
1995            summary: request.summary,
1996        };
1997        if let Err(error) = send_turn_event(
1998            &self.output,
1999            event,
2000            self.meta.max_server_frame,
2001            &cancellation,
2002        )
2003        .await
2004        {
2005            self.broker.remove(&approval_id).await;
2006            return Err(error);
2007        }
2008        tokio::select! {
2009            result = receiver => result.map_err(|_| AgentError::Cancelled),
2010            _ = cancellation.cancelled() => {
2011                self.broker.remove(&approval_id).await;
2012                Err(AgentError::Cancelled)
2013            }
2014        }
2015    }
2016}
2017
2018async fn send_event(output: &OutboundSender, event: ServerEvent, max_bytes: usize) -> Result<()> {
2019    let bytes = encode_event(&event, max_bytes)?;
2020    output.send(bytes, None).await.map_err(anyhow::Error::new)
2021}
2022
2023async fn send_turn_event(
2024    output: &OutboundSender,
2025    event: ServerEvent,
2026    max_bytes: usize,
2027    cancellation: &CancellationToken,
2028) -> Result<(), AgentError> {
2029    let bytes = encode_event(&event, max_bytes)
2030        .map_err(|error| AgentError::ResponseLimit(error.to_string()))?;
2031    match output.send(bytes, Some(cancellation)).await {
2032        Ok(()) => Ok(()),
2033        Err(OutboundSendError::Cancelled) => Err(AgentError::Cancelled),
2034        Err(error) => Err(AgentError::Internal(error.to_string())),
2035    }
2036}
2037
2038fn encode_event(event: &ServerEvent, max_bytes: usize) -> Result<Vec<u8>> {
2039    let bytes = serde_json::to_vec(event).context("serialize protocol event")?;
2040    if bytes.len() > max_bytes {
2041        return Err(anyhow!("server event exceeds configured frame limit"));
2042    }
2043    Ok(bytes)
2044}
2045
2046async fn send_error(
2047    output: &OutboundSender,
2048    request_id: &str,
2049    code: &str,
2050    message: &str,
2051    fatal: bool,
2052    max_bytes: usize,
2053) -> Result<()> {
2054    send_event(
2055        output,
2056        ServerEvent::Error {
2057            request_id: (!request_id.is_empty()).then(|| request_id.to_owned()),
2058            code: code.into(),
2059            message: message.into(),
2060            fatal,
2061        },
2062        max_bytes,
2063    )
2064    .await
2065}
2066
2067#[cfg(test)]
2068mod tests {
2069    use std::{
2070        future::pending,
2071        io::Cursor,
2072        sync::atomic::{AtomicBool, Ordering},
2073    };
2074
2075    /// A delegation registry in a private temporary instance home.
2076    fn test_registry() -> Arc<DelegationRegistry> {
2077        let home = tempfile::tempdir().unwrap().keep();
2078        Arc::new(DelegationRegistry::new(&home))
2079    }
2080
2081    use super::*;
2082
2083    #[test]
2084    fn clients_and_tools_agree_on_the_depth_variable() {
2085        assert_eq!(
2086            scv_client::DELEGATION_DEPTH_VARIABLE,
2087            scv_tools::delegation::DEPTH_VARIABLE
2088        );
2089    }
2090
2091    #[test]
2092    fn progress_text_is_bounded_on_a_character_boundary() {
2093        let text = "é".repeat(400);
2094        let bounded = bounded_progress(text);
2095        assert!(bounded.len() <= scv_core::MAX_PROGRESS_EVENT_BYTES);
2096        assert!(bounded.chars().all(|character| character == 'é'));
2097        assert_eq!(bounded_progress("short".into()), "short");
2098    }
2099
2100    struct DropSignal(Arc<AtomicBool>);
2101
2102    impl Drop for DropSignal {
2103        fn drop(&mut self) {
2104            self.0.store(true, Ordering::Release);
2105        }
2106    }
2107
2108    #[tokio::test]
2109    async fn nonreading_management_client_does_not_hold_component_lock() {
2110        let (mut input, server_input) = tokio::io::duplex(65536);
2111        let (server_output, _blocked_output) = tokio::io::duplex(1);
2112        let tasks = TaskTracker::new();
2113        let components = Arc::new(Mutex::new(components::Components::new(
2114            PathBuf::from("/unused.sock"),
2115            PathBuf::from("/"),
2116        )));
2117        let cancel = CancellationToken::new();
2118        let handler = tokio::spawn(run_managed(
2119            server_input,
2120            server_output,
2121            ConfigOverrides::default(),
2122            Some(components.clone()),
2123            test_registry(),
2124            cancel.clone(),
2125            tasks.clone(),
2126        ));
2127        input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
2128        for _ in 0..300 {
2129            input.write_all(b"{\"type\":\"daemon.control\",\"request_id\":\"s\",\"command\":{\"action\":\"status\"}}\n").await.unwrap();
2130        }
2131        tokio::time::sleep(Duration::from_millis(50)).await;
2132        let status = tokio::time::timeout(Duration::from_millis(100), async {
2133            components.lock().await.status()
2134        })
2135        .await
2136        .unwrap();
2137        assert_eq!(status.pid, std::process::id());
2138        cancel.cancel();
2139        handler.abort();
2140        let _ = handler.await;
2141        tasks.close();
2142        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
2143            .await
2144            .unwrap();
2145    }
2146
2147    #[tokio::test]
2148    async fn forced_connection_abort_drops_and_joins_writer_descendants() {
2149        let (mut input, server_input) = tokio::io::duplex(512);
2150        let (server_output, _blocked_output) = tokio::io::duplex(1);
2151        let tasks = TaskTracker::new();
2152        let handler = tokio::spawn(run_managed(
2153            server_input,
2154            server_output,
2155            ConfigOverrides::default(),
2156            None,
2157            test_registry(),
2158            CancellationToken::new(),
2159            tasks.clone(),
2160        ));
2161        input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
2162        tokio::time::timeout(Duration::from_secs(1), async {
2163            while tasks.is_empty() {
2164                tokio::task::yield_now().await;
2165            }
2166        })
2167        .await
2168        .unwrap();
2169        handler.abort();
2170        let _ = handler.await;
2171        tasks.close();
2172        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
2173            .await
2174            .unwrap();
2175        assert!(tasks.is_empty());
2176    }
2177
2178    #[tokio::test]
2179    async fn forced_handler_abort_cancels_and_joins_active_turn() {
2180        let tasks = TaskTracker::new();
2181        let cancellation = CancellationToken::new();
2182        let child_cancel = cancellation.child_token();
2183        let observed_cancel = child_cancel.clone();
2184        let dropped = Arc::new(AtomicBool::new(false));
2185        let (ready_tx, ready_rx) = oneshot::channel();
2186        let task = tasks.spawn({
2187            let dropped = dropped.clone();
2188            async move {
2189                let _guard = DropSignal(dropped);
2190                let _ = ready_tx.send(());
2191                pending::<()>().await;
2192            }
2193        });
2194        ready_rx.await.unwrap();
2195        let (owned_tx, owned_rx) = oneshot::channel();
2196        let handler = tokio::spawn(async move {
2197            let _active = ActiveTurn {
2198                turn_id: "test".into(),
2199                cancellation: child_cancel,
2200                task,
2201            };
2202            let _ = owned_tx.send(());
2203            pending::<()>().await;
2204        });
2205        owned_rx.await.unwrap();
2206        handler.abort();
2207        let _ = handler.await;
2208        tasks.close();
2209        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
2210            .await
2211            .unwrap();
2212        assert!(observed_cancel.is_cancelled());
2213        assert!(dropped.load(Ordering::Acquire));
2214    }
2215
2216    #[tokio::test]
2217    async fn frame_buffer_preserves_partial_and_discard_state_across_cancellation() {
2218        let (mut input, output) = tokio::io::duplex(64);
2219        let mut reader = BufReader::new(output);
2220        let mut frames = FrameBuffer::default();
2221        input.write_all(b"12").await.unwrap();
2222        assert!(
2223            tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
2224                .await
2225                .is_err()
2226        );
2227        input.write_all(b"34\n").await.unwrap();
2228        assert!(
2229            matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"1234")
2230        );
2231        input.write_all(b"123456789").await.unwrap();
2232        assert!(
2233            tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
2234                .await
2235                .is_err()
2236        );
2237        input.write_all(b"\n{}\n").await.unwrap();
2238        assert!(matches!(
2239            frames.read(&mut reader, 4).await.unwrap(),
2240            FrameRead::TooLarge
2241        ));
2242        assert!(
2243            matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"{}")
2244        );
2245    }
2246
2247    #[tokio::test]
2248    async fn bounded_reader_discards_an_oversized_line() {
2249        let input = format!("{}\n{{}}\n", "x".repeat(10));
2250        let mut reader = BufReader::new(Cursor::new(input.into_bytes()));
2251        assert!(matches!(
2252            read_bounded_frame(&mut reader, 4).await.unwrap(),
2253            FrameRead::TooLarge
2254        ));
2255        match read_bounded_frame(&mut reader, 4).await.unwrap() {
2256            FrameRead::Frame(frame) => assert_eq!(frame, b"{}"),
2257            _ => panic!("expected the frame following the oversized line"),
2258        }
2259    }
2260
2261    #[tokio::test]
2262    async fn bounded_reader_accepts_exact_crlf_limit() {
2263        let mut reader = BufReader::new(Cursor::new(b"1234\r\n".to_vec()));
2264        match read_bounded_frame(&mut reader, 4).await.unwrap() {
2265            FrameRead::Frame(frame) => assert_eq!(frame, b"1234"),
2266            _ => panic!("expected an exact-limit frame"),
2267        }
2268    }
2269
2270    #[tokio::test]
2271    async fn outbound_byte_backpressure_is_cancellation_aware() {
2272        let (output, mut receiver) = outbound_channel(5);
2273        output.send(vec![0; 4], None).await.unwrap();
2274
2275        let cancellation = CancellationToken::new();
2276        let blocked = tokio::spawn({
2277            let output = output.clone();
2278            let cancellation = cancellation.clone();
2279            async move { output.send(vec![1; 4], Some(&cancellation)).await }
2280        });
2281        tokio::task::yield_now().await;
2282        assert!(!blocked.is_finished());
2283
2284        cancellation.cancel();
2285        assert_eq!(blocked.await.unwrap(), Err(OutboundSendError::Cancelled));
2286
2287        drop(receiver.recv().await.unwrap());
2288        output.send(vec![2; 4], None).await.unwrap();
2289    }
2290
2291    #[tokio::test]
2292    async fn outbound_control_send_times_out_under_byte_backpressure() {
2293        let (output, _receiver) = outbound_channel(5);
2294        output.send(vec![0; 4], None).await.unwrap();
2295        let result = output
2296            .send_with_timeout(vec![1; 4], None, Duration::from_millis(10))
2297            .await;
2298        assert_eq!(result, Err(OutboundSendError::TimedOut));
2299    }
2300
2301    #[tokio::test]
2302    async fn active_turn_shutdown_aborts_after_grace_period() {
2303        let cancellation = CancellationToken::new();
2304        let dropped = Arc::new(AtomicBool::new(false));
2305        let (started_tx, started_rx) = oneshot::channel();
2306        let task = tokio::spawn({
2307            let dropped = Arc::clone(&dropped);
2308            async move {
2309                let _signal = DropSignal(dropped);
2310                let _ = started_tx.send(());
2311                pending::<()>().await;
2312            }
2313        });
2314        started_rx.await.unwrap();
2315
2316        let graceful = shutdown_active_turn(
2317            ActiveTurn {
2318                turn_id: "turn".into(),
2319                cancellation,
2320                task,
2321            },
2322            Duration::from_millis(10),
2323        )
2324        .await;
2325
2326        assert!(!graceful);
2327        assert!(dropped.load(Ordering::Acquire));
2328    }
2329
2330    #[tokio::test]
2331    async fn writer_shutdown_aborts_after_grace_period() {
2332        let dropped = Arc::new(AtomicBool::new(false));
2333        let (started_tx, started_rx) = oneshot::channel();
2334        let writer = tokio::spawn({
2335            let dropped = Arc::clone(&dropped);
2336            async move {
2337                let _signal = DropSignal(dropped);
2338                let _ = started_tx.send(());
2339                pending::<std::io::Result<()>>().await
2340            }
2341        });
2342        started_rx.await.unwrap();
2343
2344        let result = shutdown_writer(writer, Duration::from_millis(10)).await;
2345
2346        assert!(result.is_err());
2347        assert!(dropped.load(Ordering::Acquire));
2348    }
2349
2350    #[tokio::test]
2351    async fn workspace_projects_list_their_agent_skills_for_delegation() {
2352        use std::os::unix::fs::symlink;
2353        let temporary = tempfile::tempdir().unwrap();
2354        let workspace = temporary.path().canonicalize().unwrap();
2355        let outside = tempfile::tempdir().unwrap();
2356        let outside = outside.path().canonicalize().unwrap();
2357        let write_skill = |directory: &Path, description: &str| {
2358            std::fs::create_dir_all(directory).unwrap();
2359            std::fs::write(
2360                directory.join("SKILL.md"),
2361                format!(
2362                    "---\nname: skill\ndescription: {description}\n---\nBody of {description}\n"
2363                ),
2364            )
2365            .unwrap();
2366        };
2367        write_skill(
2368            &workspace.join("scv/.agents/skills/feature-flow"),
2369            "Land SCV",
2370        );
2371        std::fs::create_dir_all(workspace.join("scv/.claude/skills")).unwrap();
2372        symlink(
2373            "../../.agents/skills/feature-flow",
2374            workspace.join("scv/.claude/skills/feature-flow"),
2375        )
2376        .unwrap();
2377        write_skill(
2378            &workspace.join("web/.claude/skills/deploy"),
2379            "Deploy the site",
2380        );
2381        write_skill(&workspace.join(".agents/skills/triage"), "Root triage");
2382        write_skill(&workspace.join(".agents/skills/notes"), "Root notes");
2383        write_skill(&workspace.join(".scv/skills/triage"), "SCV triage");
2384        write_skill(&workspace.join(".hidden/.agents/skills/secret"), "Hidden");
2385        write_skill(&outside.join(".agents/skills/evil"), "Outside");
2386        symlink(&outside, workspace.join("escape")).unwrap();
2387        std::fs::create_dir_all(workspace.join("rogue/.agents")).unwrap();
2388        symlink(
2389            outside.join(".agents/skills"),
2390            workspace.join("rogue/.agents/skills"),
2391        )
2392        .unwrap();
2393        std::fs::create_dir_all(workspace.join("sneaky/.agents/skills/leak")).unwrap();
2394        symlink(
2395            outside.join(".agents/skills/evil/SKILL.md"),
2396            workspace.join("sneaky/.agents/skills/leak/SKILL.md"),
2397        )
2398        .unwrap();
2399        std::fs::write(workspace.join("file"), "not a project").unwrap();
2400        let mut config = Config::default();
2401        config.skills.user_dir = workspace.join("no-user-skills");
2402
2403        let skills = discover_skills(&workspace, &config, true).unwrap();
2404        let mut names: Vec<_> = skills.map.keys().cloned().collect();
2405        names.sort();
2406        assert_eq!(names, ["notes", "scv:feature-flow", "triage", "web:deploy"]);
2407        assert_eq!(
2408            skills.map["triage"],
2409            workspace.join(".scv/skills/triage/SKILL.md")
2410        );
2411        assert_eq!(
2412            skills.project_listing,
2413            "- notes (workspace root): Root notes\n\
2414             - scv:feature-flow (project scv): Land SCV\n\
2415             - web:deploy (project web): Deploy the site\n"
2416        );
2417        let prompt = build_system_prompt(&workspace, &config, &skills).unwrap();
2418        assert!(prompt.contains("# Project skills"));
2419        assert!(prompt.contains("set its cwd to the skill's project"));
2420
2421        let registry = builtin_registry(
2422            config.tools(),
2423            skills.map,
2424            skills.roots,
2425            config.skills.max_skill_bytes,
2426            HashMap::new(),
2427        )
2428        .unwrap();
2429        let read_skill = registry.get("read_skill").unwrap();
2430        let loaded = read_skill
2431            .execute(
2432                serde_json::json!({"name":"scv:feature-flow"}),
2433                scv_core::ToolContext::new(workspace.clone(), CancellationToken::new()),
2434            )
2435            .await
2436            .unwrap();
2437        assert!(loaded.content.contains("Body of Land SCV"));
2438
2439        let tool_free = discover_skills(&workspace, &config, false).unwrap();
2440        assert!(tool_free.project_listing.is_empty());
2441        assert!(!tool_free.map.contains_key("scv:feature-flow"));
2442        config.skills.scan_projects = false;
2443        let disabled = discover_skills(&workspace, &config, true).unwrap();
2444        assert!(disabled.project_listing.is_empty());
2445        config.skills.scan_projects = true;
2446        config.skills.max_skills = 3;
2447        let capped = discover_skills(&workspace, &config, true).unwrap();
2448        assert_eq!(capped.map.len(), 3);
2449        assert!(capped.map.contains_key("scv:feature-flow"));
2450        assert!(!capped.map.contains_key("web:deploy"));
2451    }
2452
2453    #[tokio::test]
2454    async fn daemon_control_lists_and_stops_delegations() {
2455        use std::os::unix::process::CommandExt as _;
2456        let home = tempfile::tempdir().unwrap();
2457        let registry = DelegationRegistry::new(home.path());
2458        let components = Arc::new(Mutex::new(components::Components::new(
2459            PathBuf::from("/unused.sock"),
2460            PathBuf::from("/"),
2461        )));
2462        // A run owned by another live SCV process of the same instance.
2463        let mut owner = std::process::Command::new("sleep")
2464            .arg("30")
2465            .spawn()
2466            .unwrap();
2467        let mut agent = std::process::Command::new("sleep")
2468            .arg("30")
2469            .process_group(0)
2470            .spawn()
2471            .unwrap();
2472        let identity = |pid| delegations::ProcessIdentity::of(pid).unwrap();
2473        let record = delegations::DelegationRecord {
2474            handle: "codex-a1b2c3".into(),
2475            agent: "codex".into(),
2476            instance: registry.instance().into(),
2477            session: "session".into(),
2478            owner: identity(owner.id()),
2479            process: identity(agent.id()),
2480            pgid: agent.id(),
2481            cwd: "/work/project\u{7}".into(),
2482            started_unix: 1,
2483            depth: 1,
2484            conversation: Some("codex-2".into()),
2485            turn: Some(3),
2486        };
2487        std::fs::create_dir_all(registry.record_dir()).unwrap();
2488        std::fs::write(
2489            registry.record_dir().join("codex-a1b2c3.json"),
2490            serde_json::to_vec(&record).unwrap(),
2491        )
2492        .unwrap();
2493        let control = |command| daemon_control(&components, &registry, command);
2494        let Ok(status) = control(DaemonCommand::Status).await else {
2495            panic!("status failed");
2496        };
2497        assert_eq!(status.delegations.active, 1);
2498        assert!(status.delegations.entries.is_empty());
2499        let Ok(status) = control(DaemonCommand::Delegations { all: false }).await else {
2500            panic!("listing failed");
2501        };
2502        let [entry] = status.delegations.entries.as_slice() else {
2503            panic!("{:?}", status.delegations);
2504        };
2505        assert_eq!(entry.handle, "codex-a1b2c3");
2506        assert_eq!(entry.conversation.as_deref(), Some("codex-2"));
2507        assert_eq!(entry.turn, Some(3));
2508        assert_eq!(entry.pid, agent.id());
2509        assert_eq!(entry.owner_pid, owner.id());
2510        assert!(!entry.orphaned);
2511        assert_eq!(entry.processes, 1);
2512        for command in [
2513            DaemonCommand::DelegationKill {
2514                handle: Some("codex-nosuch".into()),
2515                orphans: false,
2516            },
2517            DaemonCommand::DelegationKill {
2518                handle: None,
2519                orphans: false,
2520            },
2521        ] {
2522            assert!(matches!(
2523                control(command).await,
2524                Err(ControlFailure::Delegation(_))
2525            ));
2526        }
2527        let Ok(status) = control(DaemonCommand::DelegationKill {
2528            handle: Some("codex-a1b2c3".into()),
2529            orphans: false,
2530        })
2531        .await
2532        else {
2533            panic!("kill failed");
2534        };
2535        assert_eq!(status.delegations.killed, ["codex-a1b2c3"]);
2536        assert!(agent.wait().unwrap().code().is_none());
2537        // Its live owner removes the record itself; once the owner is gone
2538        // the record is an orphan that an orphan sweep removes.
2539        owner.kill().unwrap();
2540        owner.wait().unwrap();
2541        let Ok(status) = control(DaemonCommand::Delegations { all: true }).await else {
2542            panic!("listing failed");
2543        };
2544        assert!(status.delegations.entries[0].orphaned);
2545        assert_eq!(status.delegations.active, 0);
2546        let Ok(_) = control(DaemonCommand::DelegationKill {
2547            handle: None,
2548            orphans: true,
2549        })
2550        .await
2551        else {
2552            panic!("orphan sweep failed");
2553        };
2554        assert!(registry.list(true).is_empty());
2555    }
2556}