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