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