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