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