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