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