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