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