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 provider = Arc::new(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    let tools = if no_tools {
1182        Arc::new(ToolRegistry::default())
1183    } else {
1184        Arc::new(builtin_registry(
1185            config.tools(),
1186            skills.map,
1187            skills.roots,
1188            config.skills.max_skill_bytes,
1189            config.adapters(),
1190        )?)
1191    };
1192    let context = Arc::new(BudgetContextPolicy::new((&config.context).into())?);
1193    let runtime = Arc::new(AgentRuntime::new(
1194        provider,
1195        tools,
1196        context,
1197        config.core_agent(system_prompt),
1198        workspace.clone(),
1199    ));
1200    Ok(Session {
1201        id: Uuid::new_v4().to_string(),
1202        workspace,
1203        config,
1204        runtime,
1205        history: Arc::new(Mutex::new(Vec::new())),
1206        seq: Arc::new(AtomicU64::new(0)),
1207        queue: Arc::new(Mutex::new(VecDeque::new())),
1208        paused: Arc::new(std::sync::atomic::AtomicBool::new(false)),
1209    })
1210}
1211
1212fn build_system_prompt(
1213    workspace: &Path,
1214    config: &Config,
1215    skills: &DiscoveredSkills,
1216) -> Result<String> {
1217    let mut prompt = config.agent.system_prompt.clone();
1218    prompt.push_str(&format!(
1219        "\nCurrent working directory: {}\n",
1220        workspace.display()
1221    ));
1222    let agents_path = workspace.join("AGENTS.md");
1223    if agents_path.is_file() {
1224        let canonical = std::fs::canonicalize(&agents_path).context("resolve project AGENTS.md")?;
1225        if !canonical.starts_with(workspace) {
1226            return Err(anyhow!("project AGENTS.md escaped workspace"));
1227        }
1228        let (bytes, truncated) = read_prefix(&canonical, config.tools.max_read_bytes)
1229            .context("read project AGENTS.md")?;
1230        let instructions = std::str::from_utf8(&bytes).context("project AGENTS.md is not UTF-8")?;
1231        prompt.push_str("\n# Project instructions\n");
1232        prompt.push_str(instructions);
1233        if truncated {
1234            prompt.push_str("\n[AGENTS.md truncated by configured read limit]\n");
1235        }
1236    }
1237    if !skills.listing.is_empty() {
1238        prompt.push_str("\n# Available skills\n");
1239        prompt.push_str(&skills.listing);
1240        prompt.push_str("\nUse read_skill with a skill name when its workflow applies.\n");
1241    }
1242    if !skills.project_listing.is_empty() {
1243        prompt.push_str("\n# Project skills\n");
1244        prompt.push_str(
1245            "Projects in this workspace provide these skills to agents working in them:\n",
1246        );
1247        prompt.push_str(&skills.project_listing);
1248        prompt.push_str(
1249            "\nTo use one, delegate with an agent_* tool such as agent_codex or agent_claude, \
1250             set its cwd to the skill's project, and name the skill in the prompt: that agent \
1251             then loads the project's instructions and skills itself. read_skill loads a \
1252             skill for reference.\n",
1253        );
1254    }
1255    Ok(prompt)
1256}
1257
1258/// Skills found at session start: the names `read_skill` serves, the roots it
1259/// revalidates them against, and their system-prompt listings.
1260struct DiscoveredSkills {
1261    map: SkillMap,
1262    roots: Vec<PathBuf>,
1263    listing: String,
1264    project_listing: String,
1265}
1266
1267/// Agent-native skill directories, relative to a project, that Codex and
1268/// Claude Code load from their working directory.
1269const PROJECT_SKILL_DIRS: [&str; 2] = [".agents/skills", ".claude/skills"];
1270/// Workspace entries and child projects inspected for project skills, so a
1271/// large workspace such as a home directory costs bounded lookups.
1272const MAX_WORKSPACE_ENTRIES: usize = 4096;
1273const MAX_SKILL_PROJECTS: usize = 256;
1274/// Bytes read to find a project skill's description, and its listed length.
1275const PROJECT_SKILL_HEADER_BYTES: usize = 16 * 1024;
1276const MAX_PROJECT_SKILL_DESCRIPTION: usize = 400;
1277
1278fn discover_skills(workspace: &Path, config: &Config, tools: bool) -> Result<DiscoveredSkills> {
1279    let mut skills = SkillMap::new();
1280    let mut roots = Vec::new();
1281    let project_root = workspace.join(&config.skills.project_dir);
1282    for (root, must_be_workspace) in [(&project_root, true), (&config.skills.user_dir, false)] {
1283        if !root.is_dir() {
1284            continue;
1285        }
1286        let canonical = std::fs::canonicalize(root)
1287            .with_context(|| format!("resolve skill root {}", root.display()))?;
1288        if must_be_workspace && !canonical.starts_with(workspace) {
1289            return Err(anyhow!("project skill root escaped workspace"));
1290        }
1291        roots.push(canonical.clone());
1292        let mut entries: Vec<_> = std::fs::read_dir(&canonical)
1293            .with_context(|| format!("read skill root {}", canonical.display()))?
1294            .filter_map(Result::ok)
1295            .collect();
1296        entries.sort_by_key(|entry| entry.file_name());
1297        for entry in entries {
1298            if skills.len() >= config.skills.max_skills {
1299                break;
1300            }
1301            let path = entry.path().join("SKILL.md");
1302            if !path.is_file() {
1303                continue;
1304            }
1305            let canonical_file = std::fs::canonicalize(&path)
1306                .with_context(|| format!("resolve skill {}", path.display()))?;
1307            if !canonical_file.starts_with(&canonical) {
1308                continue;
1309            }
1310            let name = entry.file_name().to_string_lossy().to_string();
1311            skills.entry(name).or_insert(canonical_file);
1312        }
1313    }
1314    let mut names: Vec<_> = skills.keys().cloned().collect();
1315    names.sort();
1316    let mut listing = String::new();
1317    for name in names {
1318        let path = &skills[&name];
1319        let bytes = read_prefix(path, config.skills.max_skill_bytes)
1320            .map(|(bytes, _)| bytes)
1321            .unwrap_or_default();
1322        let content = String::from_utf8_lossy(&bytes);
1323        let description = skill_description(&content);
1324        listing.push_str(&format!("- {name}: {description}\n"));
1325    }
1326    // Project skills are only actionable by delegating, so tool-free sessions
1327    // neither list them nor learn the workspace's project names.
1328    let project_listing = if tools && config.skills.scan_projects {
1329        discover_project_skills(workspace, config, &mut skills, &mut roots)
1330    } else {
1331        String::new()
1332    };
1333    Ok(DiscoveredSkills {
1334        map: skills,
1335        roots,
1336        listing,
1337        project_listing,
1338    })
1339}
1340
1341/// List the agent skills of the workspace and its immediate, non-hidden child
1342/// projects. A child's skills are named `<project>:<skill>`. Everything
1343/// resolves inside the workspace, SKILL.md files that resolve to the same file
1344/// (such as a `.claude/skills` link to `.agents/skills`) count once, and
1345/// unreadable entries are skipped so one broken project cannot stop a session.
1346fn discover_project_skills(
1347    workspace: &Path,
1348    config: &Config,
1349    skills: &mut SkillMap,
1350    roots: &mut Vec<PathBuf>,
1351) -> String {
1352    let mut projects = vec![(None, workspace.to_path_buf())];
1353    let mut names: Vec<_> = std::fs::read_dir(workspace)
1354        .into_iter()
1355        .flatten()
1356        .filter_map(Result::ok)
1357        .take(MAX_WORKSPACE_ENTRIES)
1358        .map(|entry| entry.file_name().to_string_lossy().into_owned())
1359        .filter(|name| !name.starts_with('.'))
1360        .collect();
1361    names.sort();
1362    for name in names {
1363        if projects.len() > MAX_SKILL_PROJECTS {
1364            break;
1365        }
1366        let Ok(directory) = std::fs::canonicalize(workspace.join(&name)) else {
1367            continue;
1368        };
1369        if directory.is_dir()
1370            && directory.starts_with(workspace)
1371            && !projects.iter().any(|(_, seen)| seen == &directory)
1372        {
1373            projects.push((Some(name), directory));
1374        }
1375    }
1376    let mut seen_files = std::collections::HashSet::new();
1377    let mut listing = String::new();
1378    'projects: for (project, directory) in projects {
1379        let project_roots: Vec<PathBuf> = PROJECT_SKILL_DIRS
1380            .iter()
1381            .filter_map(|relative| std::fs::canonicalize(directory.join(relative)).ok())
1382            .filter(|root| root.is_dir() && root.starts_with(workspace))
1383            .collect();
1384        for root in &project_roots {
1385            if !roots.contains(root) {
1386                roots.push(root.clone());
1387            }
1388        }
1389        for root in &project_roots {
1390            let mut entries: Vec<_> = std::fs::read_dir(root)
1391                .into_iter()
1392                .flatten()
1393                .filter_map(Result::ok)
1394                .take(MAX_WORKSPACE_ENTRIES)
1395                .collect();
1396            entries.sort_by_key(|entry| entry.file_name());
1397            for entry in entries {
1398                if skills.len() >= config.skills.max_skills {
1399                    break 'projects;
1400                }
1401                let Ok(file) = std::fs::canonicalize(entry.path().join("SKILL.md")) else {
1402                    continue;
1403                };
1404                if !file.is_file()
1405                    || !project_roots.iter().any(|root| file.starts_with(root))
1406                    || !seen_files.insert(file.clone())
1407                {
1408                    continue;
1409                }
1410                let skill = entry.file_name().to_string_lossy().into_owned();
1411                let (name, location) = match &project {
1412                    Some(project) => (format!("{project}:{skill}"), format!("project {project}")),
1413                    None => (skill, "workspace root".to_owned()),
1414                };
1415                // SCV's own and the user's skills keep their names.
1416                if skills.contains_key(&name) {
1417                    continue;
1418                }
1419                let header = read_prefix(
1420                    &file,
1421                    config
1422                        .skills
1423                        .max_skill_bytes
1424                        .min(PROJECT_SKILL_HEADER_BYTES),
1425                )
1426                .map(|(bytes, _)| bytes)
1427                .unwrap_or_default();
1428                let description: String = skill_description(&String::from_utf8_lossy(&header))
1429                    .chars()
1430                    .take(MAX_PROJECT_SKILL_DESCRIPTION)
1431                    .collect();
1432                listing.push_str(&format!("- {name} ({location}): {description}\n"));
1433                skills.insert(name, file);
1434            }
1435        }
1436    }
1437    listing
1438}
1439
1440fn read_prefix(path: &Path, max_bytes: usize) -> std::io::Result<(Vec<u8>, bool)> {
1441    let file = std::fs::File::open(path)?;
1442    let mut bytes = Vec::with_capacity(max_bytes.min(8192));
1443    file.take(
1444        u64::try_from(max_bytes)
1445            .unwrap_or(u64::MAX)
1446            .saturating_add(1),
1447    )
1448    .read_to_end(&mut bytes)?;
1449    let truncated = bytes.len() > max_bytes;
1450    bytes.truncate(max_bytes);
1451    Ok((bytes, truncated))
1452}
1453
1454fn skill_description(content: &str) -> String {
1455    if let Some(frontmatter) = content.strip_prefix("---\n")
1456        && let Some((header, _)) = frontmatter.split_once("\n---")
1457    {
1458        for line in header.lines() {
1459            if let Some(description) = line.strip_prefix("description:") {
1460                return description.trim().trim_matches('"').to_owned();
1461            }
1462        }
1463    }
1464    content
1465        .lines()
1466        .map(str::trim)
1467        .find(|line| !line.is_empty() && !line.starts_with('#'))
1468        .unwrap_or("No description provided")
1469        .chars()
1470        .take(240)
1471        .collect()
1472}
1473
1474struct ProtocolSink {
1475    meta: TurnMeta,
1476    output: OutboundSender,
1477    cancellation: CancellationToken,
1478}
1479
1480#[async_trait]
1481impl EventSink for ProtocolSink {
1482    async fn emit(&self, event: CoreEvent) -> Result<(), AgentError> {
1483        let seq = next_seq(&self.meta.seq);
1484        let event = match event {
1485            CoreEvent::AssistantDelta { content } => ServerEvent::AssistantDelta {
1486                request_id: self.meta.request_id.clone(),
1487                session_id: self.meta.session_id.clone(),
1488                turn_id: self.meta.turn_id.clone(),
1489                seq,
1490                content,
1491            },
1492            CoreEvent::AssistantCompleted { content } => ServerEvent::AssistantCompleted {
1493                request_id: self.meta.request_id.clone(),
1494                session_id: self.meta.session_id.clone(),
1495                turn_id: self.meta.turn_id.clone(),
1496                seq,
1497                content,
1498            },
1499            CoreEvent::ToolProposed {
1500                call_id,
1501                name,
1502                arguments,
1503            } => ServerEvent::ToolProposed {
1504                request_id: self.meta.request_id.clone(),
1505                session_id: self.meta.session_id.clone(),
1506                turn_id: self.meta.turn_id.clone(),
1507                seq,
1508                call_id,
1509                name,
1510                arguments,
1511            },
1512            CoreEvent::ToolStarted { call_id, name } => ServerEvent::ToolStarted {
1513                request_id: self.meta.request_id.clone(),
1514                session_id: self.meta.session_id.clone(),
1515                turn_id: self.meta.turn_id.clone(),
1516                seq,
1517                call_id,
1518                name,
1519            },
1520            CoreEvent::ToolCompleted {
1521                call_id,
1522                name,
1523                output,
1524            } => ServerEvent::ToolCompleted {
1525                request_id: self.meta.request_id.clone(),
1526                session_id: self.meta.session_id.clone(),
1527                turn_id: self.meta.turn_id.clone(),
1528                seq,
1529                call_id,
1530                name,
1531                success: !output.is_error,
1532                output: output.content,
1533                truncated: output.truncated,
1534            },
1535            CoreEvent::ContextCompacted {
1536                before_tokens,
1537                after_tokens,
1538                removed_messages,
1539            } => ServerEvent::ContextCompacted {
1540                request_id: self.meta.request_id.clone(),
1541                session_id: self.meta.session_id.clone(),
1542                turn_id: self.meta.turn_id.clone(),
1543                seq,
1544                before_tokens,
1545                after_tokens,
1546                removed_messages,
1547            },
1548            CoreEvent::SessionTrimmed {
1549                removed_messages,
1550                history_bytes,
1551            } => ServerEvent::SessionTrimmed {
1552                request_id: self.meta.request_id.clone(),
1553                session_id: self.meta.session_id.clone(),
1554                seq,
1555                removed_messages,
1556                history_bytes,
1557            },
1558        };
1559        send_turn_event(
1560            &self.output,
1561            event,
1562            self.meta.max_server_frame,
1563            &self.cancellation,
1564        )
1565        .await
1566    }
1567}
1568
1569#[derive(Default)]
1570struct ApprovalBroker {
1571    pending: Mutex<HashMap<String, oneshot::Sender<bool>>>,
1572}
1573
1574impl ApprovalBroker {
1575    async fn insert(&self, id: String, sender: oneshot::Sender<bool>) {
1576        self.pending.lock().await.insert(id, sender);
1577    }
1578
1579    async fn remove(&self, id: &str) {
1580        self.pending.lock().await.remove(id);
1581    }
1582
1583    async fn resolve(&self, id: &str, approved: bool) -> bool {
1584        let sender = self.pending.lock().await.remove(id);
1585        sender.is_some_and(|sender| sender.send(approved).is_ok())
1586    }
1587}
1588
1589struct ProtocolApprovalGate {
1590    policy: ApprovalPolicy,
1591    broker: Arc<ApprovalBroker>,
1592    meta: TurnMeta,
1593    output: OutboundSender,
1594}
1595
1596#[async_trait]
1597impl ApprovalGate for ProtocolApprovalGate {
1598    async fn approve(
1599        &self,
1600        request: ApprovalRequest,
1601        cancellation: CancellationToken,
1602    ) -> Result<bool, AgentError> {
1603        match self.policy {
1604            ApprovalPolicy::OnRisk if request.risk == ToolRisk::ReadOnly => return Ok(true),
1605            ApprovalPolicy::Never => return Ok(request.risk == ToolRisk::ReadOnly),
1606            ApprovalPolicy::Always | ApprovalPolicy::OnRisk => {}
1607        }
1608        let approval_id = Uuid::new_v4().to_string();
1609        let (sender, receiver) = oneshot::channel();
1610        self.broker.insert(approval_id.clone(), sender).await;
1611        let event = ServerEvent::ApprovalRequested {
1612            request_id: self.meta.request_id.clone(),
1613            session_id: self.meta.session_id.clone(),
1614            turn_id: self.meta.turn_id.clone(),
1615            seq: next_seq(&self.meta.seq),
1616            approval_id: approval_id.clone(),
1617            call_id: request.call_id,
1618            name: request.name,
1619            risk: request.risk.as_str().into(),
1620            cwd: request.cwd.display().to_string(),
1621            summary: request.summary,
1622        };
1623        if let Err(error) = send_turn_event(
1624            &self.output,
1625            event,
1626            self.meta.max_server_frame,
1627            &cancellation,
1628        )
1629        .await
1630        {
1631            self.broker.remove(&approval_id).await;
1632            return Err(error);
1633        }
1634        tokio::select! {
1635            result = receiver => result.map_err(|_| AgentError::Cancelled),
1636            _ = cancellation.cancelled() => {
1637                self.broker.remove(&approval_id).await;
1638                Err(AgentError::Cancelled)
1639            }
1640        }
1641    }
1642}
1643
1644async fn send_event(output: &OutboundSender, event: ServerEvent, max_bytes: usize) -> Result<()> {
1645    let bytes = encode_event(&event, max_bytes)?;
1646    output.send(bytes, None).await.map_err(anyhow::Error::new)
1647}
1648
1649async fn send_turn_event(
1650    output: &OutboundSender,
1651    event: ServerEvent,
1652    max_bytes: usize,
1653    cancellation: &CancellationToken,
1654) -> Result<(), AgentError> {
1655    let bytes = encode_event(&event, max_bytes)
1656        .map_err(|error| AgentError::ResponseLimit(error.to_string()))?;
1657    match output.send(bytes, Some(cancellation)).await {
1658        Ok(()) => Ok(()),
1659        Err(OutboundSendError::Cancelled) => Err(AgentError::Cancelled),
1660        Err(error) => Err(AgentError::Internal(error.to_string())),
1661    }
1662}
1663
1664fn encode_event(event: &ServerEvent, max_bytes: usize) -> Result<Vec<u8>> {
1665    let bytes = serde_json::to_vec(event).context("serialize protocol event")?;
1666    if bytes.len() > max_bytes {
1667        return Err(anyhow!("server event exceeds configured frame limit"));
1668    }
1669    Ok(bytes)
1670}
1671
1672async fn send_error(
1673    output: &OutboundSender,
1674    request_id: &str,
1675    code: &str,
1676    message: &str,
1677    fatal: bool,
1678    max_bytes: usize,
1679) -> Result<()> {
1680    send_event(
1681        output,
1682        ServerEvent::Error {
1683            request_id: (!request_id.is_empty()).then(|| request_id.to_owned()),
1684            code: code.into(),
1685            message: message.into(),
1686            fatal,
1687        },
1688        max_bytes,
1689    )
1690    .await
1691}
1692
1693#[cfg(test)]
1694mod tests {
1695    use std::{
1696        future::pending,
1697        io::Cursor,
1698        sync::atomic::{AtomicBool, Ordering},
1699    };
1700
1701    use super::*;
1702
1703    struct DropSignal(Arc<AtomicBool>);
1704
1705    impl Drop for DropSignal {
1706        fn drop(&mut self) {
1707            self.0.store(true, Ordering::Release);
1708        }
1709    }
1710
1711    #[tokio::test]
1712    async fn nonreading_management_client_does_not_hold_component_lock() {
1713        let (mut input, server_input) = tokio::io::duplex(65536);
1714        let (server_output, _blocked_output) = tokio::io::duplex(1);
1715        let tasks = TaskTracker::new();
1716        let components = Arc::new(Mutex::new(components::Components::new(
1717            PathBuf::from("/unused.sock"),
1718            PathBuf::from("/"),
1719        )));
1720        let cancel = CancellationToken::new();
1721        let handler = tokio::spawn(run_managed(
1722            server_input,
1723            server_output,
1724            ConfigOverrides::default(),
1725            Some(components.clone()),
1726            cancel.clone(),
1727            tasks.clone(),
1728        ));
1729        input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
1730        for _ in 0..300 {
1731            input.write_all(b"{\"type\":\"daemon.control\",\"request_id\":\"s\",\"command\":{\"action\":\"status\"}}\n").await.unwrap();
1732        }
1733        tokio::time::sleep(Duration::from_millis(50)).await;
1734        let status = tokio::time::timeout(Duration::from_millis(100), async {
1735            components.lock().await.status()
1736        })
1737        .await
1738        .unwrap();
1739        assert_eq!(status.pid, std::process::id());
1740        cancel.cancel();
1741        handler.abort();
1742        let _ = handler.await;
1743        tasks.close();
1744        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1745            .await
1746            .unwrap();
1747    }
1748
1749    #[tokio::test]
1750    async fn forced_connection_abort_drops_and_joins_writer_descendants() {
1751        let (mut input, server_input) = tokio::io::duplex(512);
1752        let (server_output, _blocked_output) = tokio::io::duplex(1);
1753        let tasks = TaskTracker::new();
1754        let handler = tokio::spawn(run_managed(
1755            server_input,
1756            server_output,
1757            ConfigOverrides::default(),
1758            None,
1759            CancellationToken::new(),
1760            tasks.clone(),
1761        ));
1762        input.write_all(b"{\"type\":\"initialize\",\"request_id\":\"init\",\"protocol_version\":2,\"client\":{\"name\":\"test\",\"version\":\"0\"}}\n").await.unwrap();
1763        tokio::time::timeout(Duration::from_secs(1), async {
1764            while tasks.is_empty() {
1765                tokio::task::yield_now().await;
1766            }
1767        })
1768        .await
1769        .unwrap();
1770        handler.abort();
1771        let _ = handler.await;
1772        tasks.close();
1773        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1774            .await
1775            .unwrap();
1776        assert!(tasks.is_empty());
1777    }
1778
1779    #[tokio::test]
1780    async fn forced_handler_abort_cancels_and_joins_active_turn() {
1781        let tasks = TaskTracker::new();
1782        let cancellation = CancellationToken::new();
1783        let child_cancel = cancellation.child_token();
1784        let observed_cancel = child_cancel.clone();
1785        let dropped = Arc::new(AtomicBool::new(false));
1786        let (ready_tx, ready_rx) = oneshot::channel();
1787        let task = tasks.spawn({
1788            let dropped = dropped.clone();
1789            async move {
1790                let _guard = DropSignal(dropped);
1791                let _ = ready_tx.send(());
1792                pending::<()>().await;
1793            }
1794        });
1795        ready_rx.await.unwrap();
1796        let (owned_tx, owned_rx) = oneshot::channel();
1797        let handler = tokio::spawn(async move {
1798            let _active = ActiveTurn {
1799                turn_id: "test".into(),
1800                cancellation: child_cancel,
1801                task,
1802            };
1803            let _ = owned_tx.send(());
1804            pending::<()>().await;
1805        });
1806        owned_rx.await.unwrap();
1807        handler.abort();
1808        let _ = handler.await;
1809        tasks.close();
1810        tokio::time::timeout(Duration::from_secs(1), tasks.wait())
1811            .await
1812            .unwrap();
1813        assert!(observed_cancel.is_cancelled());
1814        assert!(dropped.load(Ordering::Acquire));
1815    }
1816
1817    #[tokio::test]
1818    async fn frame_buffer_preserves_partial_and_discard_state_across_cancellation() {
1819        let (mut input, output) = tokio::io::duplex(64);
1820        let mut reader = BufReader::new(output);
1821        let mut frames = FrameBuffer::default();
1822        input.write_all(b"12").await.unwrap();
1823        assert!(
1824            tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
1825                .await
1826                .is_err()
1827        );
1828        input.write_all(b"34\n").await.unwrap();
1829        assert!(
1830            matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"1234")
1831        );
1832        input.write_all(b"123456789").await.unwrap();
1833        assert!(
1834            tokio::time::timeout(Duration::from_millis(10), frames.read(&mut reader, 4))
1835                .await
1836                .is_err()
1837        );
1838        input.write_all(b"\n{}\n").await.unwrap();
1839        assert!(matches!(
1840            frames.read(&mut reader, 4).await.unwrap(),
1841            FrameRead::TooLarge
1842        ));
1843        assert!(
1844            matches!(frames.read(&mut reader, 4).await.unwrap(), FrameRead::Frame(value) if value == b"{}")
1845        );
1846    }
1847
1848    #[tokio::test]
1849    async fn bounded_reader_discards_an_oversized_line() {
1850        let input = format!("{}\n{{}}\n", "x".repeat(10));
1851        let mut reader = BufReader::new(Cursor::new(input.into_bytes()));
1852        assert!(matches!(
1853            read_bounded_frame(&mut reader, 4).await.unwrap(),
1854            FrameRead::TooLarge
1855        ));
1856        match read_bounded_frame(&mut reader, 4).await.unwrap() {
1857            FrameRead::Frame(frame) => assert_eq!(frame, b"{}"),
1858            _ => panic!("expected the frame following the oversized line"),
1859        }
1860    }
1861
1862    #[tokio::test]
1863    async fn bounded_reader_accepts_exact_crlf_limit() {
1864        let mut reader = BufReader::new(Cursor::new(b"1234\r\n".to_vec()));
1865        match read_bounded_frame(&mut reader, 4).await.unwrap() {
1866            FrameRead::Frame(frame) => assert_eq!(frame, b"1234"),
1867            _ => panic!("expected an exact-limit frame"),
1868        }
1869    }
1870
1871    #[tokio::test]
1872    async fn outbound_byte_backpressure_is_cancellation_aware() {
1873        let (output, mut receiver) = outbound_channel(5);
1874        output.send(vec![0; 4], None).await.unwrap();
1875
1876        let cancellation = CancellationToken::new();
1877        let blocked = tokio::spawn({
1878            let output = output.clone();
1879            let cancellation = cancellation.clone();
1880            async move { output.send(vec![1; 4], Some(&cancellation)).await }
1881        });
1882        tokio::task::yield_now().await;
1883        assert!(!blocked.is_finished());
1884
1885        cancellation.cancel();
1886        assert_eq!(blocked.await.unwrap(), Err(OutboundSendError::Cancelled));
1887
1888        drop(receiver.recv().await.unwrap());
1889        output.send(vec![2; 4], None).await.unwrap();
1890    }
1891
1892    #[tokio::test]
1893    async fn outbound_control_send_times_out_under_byte_backpressure() {
1894        let (output, _receiver) = outbound_channel(5);
1895        output.send(vec![0; 4], None).await.unwrap();
1896        let result = output
1897            .send_with_timeout(vec![1; 4], None, Duration::from_millis(10))
1898            .await;
1899        assert_eq!(result, Err(OutboundSendError::TimedOut));
1900    }
1901
1902    #[tokio::test]
1903    async fn active_turn_shutdown_aborts_after_grace_period() {
1904        let cancellation = CancellationToken::new();
1905        let dropped = Arc::new(AtomicBool::new(false));
1906        let (started_tx, started_rx) = oneshot::channel();
1907        let task = tokio::spawn({
1908            let dropped = Arc::clone(&dropped);
1909            async move {
1910                let _signal = DropSignal(dropped);
1911                let _ = started_tx.send(());
1912                pending::<()>().await;
1913            }
1914        });
1915        started_rx.await.unwrap();
1916
1917        let graceful = shutdown_active_turn(
1918            ActiveTurn {
1919                turn_id: "turn".into(),
1920                cancellation,
1921                task,
1922            },
1923            Duration::from_millis(10),
1924        )
1925        .await;
1926
1927        assert!(!graceful);
1928        assert!(dropped.load(Ordering::Acquire));
1929    }
1930
1931    #[tokio::test]
1932    async fn writer_shutdown_aborts_after_grace_period() {
1933        let dropped = Arc::new(AtomicBool::new(false));
1934        let (started_tx, started_rx) = oneshot::channel();
1935        let writer = tokio::spawn({
1936            let dropped = Arc::clone(&dropped);
1937            async move {
1938                let _signal = DropSignal(dropped);
1939                let _ = started_tx.send(());
1940                pending::<std::io::Result<()>>().await
1941            }
1942        });
1943        started_rx.await.unwrap();
1944
1945        let result = shutdown_writer(writer, Duration::from_millis(10)).await;
1946
1947        assert!(result.is_err());
1948        assert!(dropped.load(Ordering::Acquire));
1949    }
1950
1951    #[tokio::test]
1952    async fn workspace_projects_list_their_agent_skills_for_delegation() {
1953        use std::os::unix::fs::symlink;
1954        let temporary = tempfile::tempdir().unwrap();
1955        let workspace = temporary.path().canonicalize().unwrap();
1956        let outside = tempfile::tempdir().unwrap();
1957        let outside = outside.path().canonicalize().unwrap();
1958        let write_skill = |directory: &Path, description: &str| {
1959            std::fs::create_dir_all(directory).unwrap();
1960            std::fs::write(
1961                directory.join("SKILL.md"),
1962                format!(
1963                    "---\nname: skill\ndescription: {description}\n---\nBody of {description}\n"
1964                ),
1965            )
1966            .unwrap();
1967        };
1968        write_skill(
1969            &workspace.join("scv/.agents/skills/feature-flow"),
1970            "Land SCV",
1971        );
1972        std::fs::create_dir_all(workspace.join("scv/.claude/skills")).unwrap();
1973        symlink(
1974            "../../.agents/skills/feature-flow",
1975            workspace.join("scv/.claude/skills/feature-flow"),
1976        )
1977        .unwrap();
1978        write_skill(
1979            &workspace.join("web/.claude/skills/deploy"),
1980            "Deploy the site",
1981        );
1982        write_skill(&workspace.join(".agents/skills/triage"), "Root triage");
1983        write_skill(&workspace.join(".agents/skills/notes"), "Root notes");
1984        write_skill(&workspace.join(".scv/skills/triage"), "SCV triage");
1985        write_skill(&workspace.join(".hidden/.agents/skills/secret"), "Hidden");
1986        write_skill(&outside.join(".agents/skills/evil"), "Outside");
1987        symlink(&outside, workspace.join("escape")).unwrap();
1988        std::fs::create_dir_all(workspace.join("rogue/.agents")).unwrap();
1989        symlink(
1990            outside.join(".agents/skills"),
1991            workspace.join("rogue/.agents/skills"),
1992        )
1993        .unwrap();
1994        std::fs::create_dir_all(workspace.join("sneaky/.agents/skills/leak")).unwrap();
1995        symlink(
1996            outside.join(".agents/skills/evil/SKILL.md"),
1997            workspace.join("sneaky/.agents/skills/leak/SKILL.md"),
1998        )
1999        .unwrap();
2000        std::fs::write(workspace.join("file"), "not a project").unwrap();
2001        let mut config = Config::default();
2002        config.skills.user_dir = workspace.join("no-user-skills");
2003
2004        let skills = discover_skills(&workspace, &config, true).unwrap();
2005        let mut names: Vec<_> = skills.map.keys().cloned().collect();
2006        names.sort();
2007        assert_eq!(names, ["notes", "scv:feature-flow", "triage", "web:deploy"]);
2008        assert_eq!(
2009            skills.map["triage"],
2010            workspace.join(".scv/skills/triage/SKILL.md")
2011        );
2012        assert_eq!(
2013            skills.project_listing,
2014            "- notes (workspace root): Root notes\n\
2015             - scv:feature-flow (project scv): Land SCV\n\
2016             - web:deploy (project web): Deploy the site\n"
2017        );
2018        let prompt = build_system_prompt(&workspace, &config, &skills).unwrap();
2019        assert!(prompt.contains("# Project skills"));
2020        assert!(prompt.contains("set its cwd to the skill's project"));
2021
2022        let registry = builtin_registry(
2023            config.tools(),
2024            skills.map,
2025            skills.roots,
2026            config.skills.max_skill_bytes,
2027            HashMap::new(),
2028        )
2029        .unwrap();
2030        let read_skill = registry.get("read_skill").unwrap();
2031        let loaded = read_skill
2032            .execute(
2033                serde_json::json!({"name":"scv:feature-flow"}),
2034                scv_core::ToolContext {
2035                    workspace: workspace.clone(),
2036                    cancellation: CancellationToken::new(),
2037                },
2038            )
2039            .await
2040            .unwrap();
2041        assert!(loaded.content.contains("Body of Land SCV"));
2042
2043        let tool_free = discover_skills(&workspace, &config, false).unwrap();
2044        assert!(tool_free.project_listing.is_empty());
2045        assert!(!tool_free.map.contains_key("scv:feature-flow"));
2046        config.skills.scan_projects = false;
2047        let disabled = discover_skills(&workspace, &config, true).unwrap();
2048        assert!(disabled.project_listing.is_empty());
2049        config.skills.scan_projects = true;
2050        config.skills.max_skills = 3;
2051        let capped = discover_skills(&workspace, &config, true).unwrap();
2052        assert_eq!(capped.map.len(), 3);
2053        assert!(capped.map.contains_key("scv:feature-flow"));
2054        assert!(!capped.map.contains_key("web:deploy"));
2055    }
2056}