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