1mod 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
35pub 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
60pub 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
73pub 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
89pub 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
94pub fn agent_stored_status(agent: &str, store: adapters::KeyStore) -> Result<(bool, Vec<String>)> {
97 agents::stored_status(store, &key_store_home(agent)?)
98}
99
100pub fn remove_agent_credentials(agent: &str, store: adapters::KeyStore) -> Result<Vec<String>> {
102 agents::remove_stored(store, &key_store_home(agent)?)
103}
104
105pub fn configure_pi_endpoint(endpoint: &Endpoint, key: &str) -> Result<Vec<String>> {
107 agents::configure_pi_endpoint(&pi_agent_dir()?, endpoint, key)
108}
109
110pub 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
156pub 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
165pub 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
226pub fn default_socket_path() -> Result<PathBuf> {
228 scv_client::default_socket_path()
229}
230
231pub 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(¤t.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(¤t.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(¤t.seq),
519 }, current.config.protocol.max_server_frame_bytes).await?;
520 let runtime = Arc::clone(¤t.runtime);
521 let history = Arc::clone(¤t.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(¤t.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(¤t.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(¤t.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(¤t.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(¤t.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(¤t.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(¤t.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(¤t.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(¤t.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(¤t.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(¤t.runtime);
678 let history = Arc::clone(¤t.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#[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
927struct 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 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 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 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
1258struct DiscoveredSkills {
1261 map: SkillMap,
1262 roots: Vec<PathBuf>,
1263 listing: String,
1264 project_listing: String,
1265}
1266
1267const PROJECT_SKILL_DIRS: [&str; 2] = [".agents/skills", ".claude/skills"];
1270const MAX_WORKSPACE_ENTRIES: usize = 4096;
1273const MAX_SKILL_PROJECTS: usize = 256;
1274const 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 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
1341fn 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 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}