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