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