Skip to main content

scv_server/
lib.rs

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