1use std::collections::HashMap;
10use std::error::Error;
11use std::io::{Read, Write};
12use std::path::{Path, PathBuf};
13use std::sync::{Arc, Mutex, mpsc};
14use std::time::{Duration, Instant};
15
16use tau_proto::{
17 ActionError, ActionInvoke, ActionOutput, ActionResult, AgentContextKey, AgentContextValue,
18 CborValue, DiscoveryAgentsFile, DiscoveryModifiedMicros, DiscoverySkillCandidate, Event,
19 ExtAgentContextPublish, ExtensionAgentDiscoverySnapshotDeclared, ExtensionContextReady,
20 ExtensionSessionContextReady, ExtensionSessionDiscoverySnapshotDeclared, HarnessInputMessage,
21 PromptContent, PromptFragment, PromptPriority, SessionAgentLoaded, SessionStarted,
22 ToolCancelled, ToolExample, ToolExampleSelector, ToolResult, ToolResultKind, ToolSpec, ToolTag,
23};
24use tracing::{debug, trace};
25
26#[cfg(test)]
27static DETACHED_OUTPUT_OVERLOAD_NOTIFY: Mutex<Option<mpsc::Sender<()>>> = Mutex::new(None);
28
29use crate::tools::{shell as path_crate_tools_shell, world as path_crate_tools_world};
30use crate::{
31 dir_lock as path_crate_dir_lock, display as path_crate_display, tools as path_crate_tools,
32};
33
34mod agents;
35mod argument;
36mod artifact_transfer;
37mod config;
38mod cwd_state;
39mod diff;
40mod dir_lock;
41mod display;
42mod isolation;
43#[cfg(any(target_os = "android", target_os = "linux", target_os = "macos"))]
44mod pty_stdio;
45mod runtime;
46mod scheduler;
47mod shell_output_spool;
48mod shell_process;
49mod terminal_frame;
50mod tool_lifecycle;
51mod tool_started_identity;
52mod tools;
53mod truncate;
54mod ui_shell_shutdown_generation;
55
56#[cfg(test)]
57mod tests;
58
59use crate::agents::{ancestor_dirs, discover_session_agents_files};
60use crate::artifact_transfer::{ArtifactTransferControl, ArtifactTransferManager};
61use crate::config::{ExtConfig, ShellConfig};
62use crate::cwd_state::{CwdState, WorkdirSnapshot};
63use crate::dir_lock::{DIR_LOCK_TOOL_NAME, DirLockManager};
64use crate::runtime::ShellRuntime;
65use crate::scheduler::{WorkMeta, WorkPriority, WorkScheduler};
66use crate::tool_lifecycle::{ToolCancellationState, ToolLifecycle};
67#[cfg(any(test, feature = "echo-agent"))]
68use crate::tools::ECHO_TOOL_NAME;
69use crate::tools::shell::{ShellAccessMode, ShellCommandMode};
70use crate::tools::{
71 APPLY_PATCH_TOOL_NAME, EDIT_TOOL_NAME, EXPORT_TOOL_NAME, FIND_TOOL_NAME, GPT_SHELL_TOOL_NAME,
72 GREP_TOOL_NAME, IMPORT_TOOL_NAME, LS_TOOL_NAME, READ_IMAGE_TOOL_NAME, READ_TOOL_NAME,
73 REPLACE_TOOL_NAME, SHELL_TOOL_NAME, WORKDIR_TOOL_NAME, execute_tool,
74};
75use crate::ui_shell_shutdown_generation::{
76 UiShellShutdownGeneration, UiShellShutdownGenerationCounter,
77};
78
79#[derive(Clone)]
89pub(crate) struct Output {
90 inner: OutputInner,
92 tool_name_scope: Option<(tau_proto::ToolName, tau_proto::ToolName)>,
94 failure: Arc<Mutex<MandatoryOutputFailure>>,
96}
97
98#[derive(Default)]
100struct MandatoryOutputFailure {
101 message: Option<String>,
103 failed: bool,
105 waker: Option<tau_client::ManualRuntimeWaker>,
107}
108
109#[derive(Clone)]
111enum OutputInner {
112 Client(tau_client::ClientHandle),
114 #[cfg(test)]
115 Channel(mpsc::Sender<HarnessInputMessage>),
117}
118
119impl Output {
120 fn client(handle: tau_client::ClientHandle) -> Self {
121 Self {
122 inner: OutputInner::Client(handle),
123 tool_name_scope: None,
124 failure: Arc::default(),
125 }
126 }
127
128 #[cfg(test)]
129 fn channel(tx: mpsc::Sender<HarnessInputMessage>) -> Self {
130 Self {
131 inner: OutputInner::Channel(tx),
132 tool_name_scope: None,
133 failure: Arc::default(),
134 }
135 }
136
137 fn scoped_tool(&self, local: tau_proto::ToolName, wire: tau_proto::ToolName) -> Self {
138 Self {
139 inner: self.inner.clone(),
140 tool_name_scope: Some((local, wire)),
141 failure: Arc::clone(&self.failure),
142 }
143 }
144
145 fn scope_tool_name(&self, tool_name: &mut tau_proto::ToolName) {
146 if let Some((local, wire)) = &self.tool_name_scope
147 && tool_name == local
148 {
149 *tool_name = wire.clone();
150 }
151 }
152
153 fn send(&self, mut message: HarnessInputMessage) -> tau_client::ClientResult<()> {
154 self.scope_message(&mut message);
155 let result = match &self.inner {
156 OutputInner::Client(handle) => handle.send_detached(message),
157 #[cfg(test)]
158 OutputInner::Channel(tx) => tx
159 .send(message)
160 .map_err(|_| tau_client::ClientError::WriterClosed),
161 };
162 #[cfg(test)]
163 if matches!(result, Err(tau_client::ClientError::Overloaded))
164 && let Some(notify) = DETACHED_OUTPUT_OVERLOAD_NOTIFY
165 .lock()
166 .expect("detached overload notification")
167 .as_ref()
168 {
169 let _ = notify.send(());
170 }
171 result
172 }
173
174 fn send_checked(&self, mut message: HarnessInputMessage) -> tau_client::ClientResult<()> {
176 self.scope_message(&mut message);
177 let result = match &self.inner {
178 OutputInner::Client(handle) => handle.send(message),
179 #[cfg(test)]
180 OutputInner::Channel(tx) => tx
181 .send(message)
182 .map_err(|_| tau_client::ClientError::WriterClosed),
183 };
184 self.retain_mandatory_failure(result)
185 }
186
187 fn scope_message(&self, message: &mut HarnessInputMessage) {
188 if let HarnessInputMessage::Emit(emit) = message {
189 let tool_name = match emit.event.as_mut() {
190 Event::ToolProgressReported(event) => Some(&mut event.tool_name),
191 Event::ToolResultReported(event) => Some(&mut event.tool_name),
192 Event::ToolResult(event) => Some(&mut event.tool_name),
193 Event::ToolErrorReported(event) => Some(&mut event.tool_name),
194 Event::ToolError(event) => Some(&mut event.tool_name),
195 Event::ToolCancelledReported(event) => Some(&mut event.tool_name),
196 Event::ToolCancelled(event) => Some(&mut event.tool_name),
197 _ => None,
198 };
199 if let Some(tool_name) = tool_name {
200 self.scope_tool_name(tool_name);
201 }
202 }
203 }
204
205 fn report_tool_progress(
207 &self,
208 progress: tau_proto::ToolProgress,
209 ) -> tau_client::ClientResult<()> {
210 self.send(HarnessInputMessage::emit_with_persist(
211 Event::ToolProgressReported(progress),
212 false,
213 ))
214 }
215
216 fn report_tool_terminal(&self, event: Event) -> tau_client::ClientResult<()> {
218 let mut outcome = tau_client::ToolTerminalOutcome::try_from(event).map_err(|event| {
219 tau_client::ClientError::handler(format!(
220 "terminal report helper received {}",
221 event.name()
222 ))
223 })?;
224 self.scope_tool_name(outcome.tool_name_mut());
225 let message = terminal_frame::budget_terminal_report(outcome.into_reported_event())?;
226 let result = match &self.inner {
227 OutputInner::Client(handle) => handle.send(message),
228 #[cfg(test)]
229 OutputInner::Channel(tx) => tx
230 .send(message)
231 .map_err(|_| tau_client::ClientError::WriterClosed),
232 };
233 self.retain_mandatory_failure(result)
234 }
235
236 fn install_waker(&self, waker: tau_client::ManualRuntimeWaker) {
239 self.failure
240 .lock()
241 .expect("mandatory output failure lock poisoned")
242 .waker = Some(waker);
243 }
244
245 fn take_mandatory_failure(&self) -> tau_client::ClientResult<()> {
247 let message = self
248 .failure
249 .lock()
250 .expect("mandatory output failure lock poisoned")
251 .message
252 .take();
253 match message {
254 Some(message) => Err(tau_client::ClientError::handler(message)),
255 None => Ok(()),
256 }
257 }
258
259 fn mandatory_output_failed(&self) -> bool {
261 self.failure
262 .lock()
263 .expect("mandatory output failure lock poisoned")
264 .failed
265 }
266
267 fn retain_mandatory_failure(
269 &self,
270 result: tau_client::ClientResult<()>,
271 ) -> tau_client::ClientResult<()> {
272 if let Err(error) = &result {
273 let waker = {
274 let mut failure = self
275 .failure
276 .lock()
277 .expect("mandatory output failure lock poisoned");
278 failure.message.get_or_insert_with(|| error.to_string());
279 failure.failed = true;
280 failure.waker.clone()
281 };
282 if let Some(waker) = waker {
283 waker.wake();
284 }
285 }
286 result
287 }
288
289 fn register_local_tool(
290 &self,
291 registration: tau_proto::ToolRegistrationDeclared,
292 ) -> tau_client::ClientResult<()> {
293 match &self.inner {
294 OutputInner::Client(handle) => handle.register_local_tool(registration),
295 #[cfg(test)]
296 OutputInner::Channel(tx) => tx
297 .send(HarnessInputMessage::emit_with_persist(
298 Event::ToolRegistrationDeclared(registration),
299 false,
300 ))
301 .map_err(|_| tau_client::ClientError::WriterClosed),
302 }
303 }
304}
305
306fn tool_tags(tags: &[&str]) -> Vec<ToolTag> {
307 tags.iter().map(|tag| ToolTag::new(*tag)).collect()
308}
309
310fn example_field(name: &str, value: CborValue) -> (CborValue, CborValue) {
311 (CborValue::Text(name.to_owned()), value)
312}
313
314fn example_text(value: &str) -> CborValue {
315 CborValue::Text(value.to_owned())
316}
317
318fn example_int(value: i64) -> CborValue {
319 CborValue::Integer(value.into())
320}
321
322const SHELL_DIR_FORCE_UNLOCK_ACTION_ID: &str = "shell.dir.force_unlock";
323
324const SLOW_LOCK_WAIT_THRESHOLD_SECS: u64 = 5;
325const LOCK_WAIT_DURATION_SECONDS_HEADER: &str = "lock_wait_duration_seconds";
326const XDG_USER_SKILL_SOURCE_PRECEDENCE: u32 = 0;
327const LEGACY_USER_SKILL_SOURCE_PRECEDENCE: u32 = 1;
328
329#[derive(Clone, Copy)]
330enum DiscoverySourcePolicy {
331 Environment,
332 #[cfg(any(test, feature = "echo-agent"))]
333 EmptyFixture,
334}
335
336impl DiscoverySourcePolicy {
337 const fn reads_environment(self) -> bool {
338 matches!(self, Self::Environment)
339 }
340}
341
342enum RuntimeCwdSource {
343 Process,
344 #[cfg(any(test, feature = "echo-agent"))]
345 Fixture(PathBuf),
346}
347
348struct DiscoveryScan {
351 snapshot: ExtensionSessionDiscoverySnapshotDeclared,
353 diagnostics: Vec<HarnessInputMessage>,
355}
356
357pub fn run_stdio() -> Result<(), Box<dyn Error>> {
359 tau_client::init_logging_for("tau_ext_shell");
360 run_impl(
361 std::io::stdin(),
362 std::io::stdout(),
363 DiscoverySourcePolicy::Environment,
364 RuntimeCwdSource::Process,
365 )
366}
367
368pub fn run<R, W>(reader: R, writer: W) -> Result<(), Box<dyn Error>>
373where
374 R: Read + Send + 'static,
375 W: Write + Send + 'static,
376{
377 run_impl(
378 reader,
379 writer,
380 DiscoverySourcePolicy::Environment,
381 RuntimeCwdSource::Process,
382 )
383}
384
385#[cfg(any(test, feature = "echo-agent"))]
393pub fn run_for_test_harness<R, W>(
394 reader: R,
395 writer: W,
396 fixture_cwd: PathBuf,
397) -> Result<(), Box<dyn Error>>
398where
399 R: Read + Send + 'static,
400 W: Write + Send + 'static,
401{
402 run_impl(
403 reader,
404 writer,
405 DiscoverySourcePolicy::EmptyFixture,
406 RuntimeCwdSource::Fixture(fixture_cwd),
407 )
408}
409
410fn registered_tool_specs(dir_lock_enabled: bool) -> Vec<ToolSpec> {
411 #[cfg(any(test, feature = "echo-agent"))]
412 let echo_tool = Some(ToolSpec {
413 provider_scope: None,
414 name: tau_proto::ToolName::new(ECHO_TOOL_NAME),
415 model_visible_name: None,
416 description: Some("Echo the provided payload unchanged".to_owned()),
417 tool_type: tau_proto::ToolType::Function,
418 parameters: None,
419 format: None,
420 tags: tool_tags(&["test:echo"]),
421 enabled_by_default: false,
422 background_support: None,
423 examples: Vec::new(),
424 });
425 #[cfg(not(any(test, feature = "echo-agent")))]
426 let echo_tool: Option<ToolSpec> = None;
427 let mut tools = Vec::new();
428 if let Some(echo_tool) = echo_tool {
429 tools.push(echo_tool);
430 }
431 let read_tool = ToolSpec {
432 provider_scope: None,
433 name: tau_proto::ToolName::new(READ_TOOL_NAME),
434 model_visible_name: None,
435 description: Some(
436 "Reads a file. Defaults to reading the whole file in one call — \
437 output is capped at 2000 lines / 10 KiB. Truncated output keeps \
438 the first 1000 and last 1000 lines separated by a literal `...` line. \
439 Files over 10 MiB are rejected by an input safety cap before output truncation. \
440 Prefer one full read. Pass inclusive `start_line`/`end_line` only to \
441 fetch one specific known slice, or `ranges` for up to 100 slices; \
442 range chunks are separated by one empty line and may overlap, but large overlapping \
443 multi-range expansions can be rejected before rendering to keep memory bounded. `start_line` past EOF errors, \
444 while `end_line` past EOF returns available lines. Returned content lines are prefixed \
445 by their 1-based line number and a space; \
446 CRLF, CR, and missing final line endings are marked after the number, e.g. \
447 `2(crlf)`, `3(cr)`, or `4(no_nl)`. Invalid UTF-8 is shown with \
448 Unicode replacement characters and an `invalid-utf8` line flag. Lines that would exceed \
449 the 10 KiB visible output budget are marker-only, e.g. `1(truncated)`. Truncated results include `truncated: true`, `total_lines`, \
450 and `total_bytes`, plus a private path to bounded saved output (or `saved_output_unavailable: true` when private storage fails); `valid_utf8: false` is included only when applicable."
451 .to_owned(),
452 ),
453 tool_type: tau_proto::ToolType::Function,
454 parameters: Some(serde_json::json!({
455 "type": "object",
456 "properties": {
457 "path": {
458 "type": "string",
459 "description": "Path to the file"
460 },
461 "start_line": {
462 "type": "integer",
463 "minimum": 1,
464 "description": "Optional, 1-based inclusive. Omit to start at line 1 (the default)."
465 },
466 "end_line": {
467 "type": "integer",
468 "minimum": 1,
469 "description": "Optional, 1-based inclusive. Omit to read to end of file (the default and preferred mode). Set this only to continue past a previous truncation, or to fetch a known specific slice of a large file — do NOT pre-slice an ordinary file you haven't already established is large."
470 },
471 "ranges": {
472 "type": "array",
473 "description": "Optional list of inclusive line ranges to read. Cannot be combined with top-level start_line or end_line. Each chunk is separated by one empty line in the output, and overlapping ranges are returned redundantly. Requests whose overlapping ranges would expand into too much rendered content are rejected before rendering.",
474 "minItems": 1,
475 "maxItems": 100,
476 "items": {
477 "type": "object",
478 "properties": {
479 "start_line": {
480 "type": "integer",
481 "minimum": 1,
482 "description": "1-based inclusive start line to read."
483 },
484 "end_line": {
485 "type": "integer",
486 "minimum": 1,
487 "description": "1-based inclusive end line to read."
488 }
489 },
490 "required": ["start_line", "end_line"],
491 "additionalProperties": false
492 }
493 }
494 },
495 "required": ["path"],
496 "additionalProperties": false
497 })),
498 format: None,
499 tags: tool_tags(&["shell:read", tau_proto::TURN_DATA_FETCH_TOOL_TAG]),
500 enabled_by_default: true,
501 background_support: None,
502 examples: vec![ToolExample {
503 id: "read-file".to_owned(),
504 title: Some("Read a file".to_owned()),
505 arguments: CborValue::Map(vec![example_field("path", example_text("src/main.rs"))]),
506 note: Some("Use only the path field for a full-file read.".to_owned()),
507 subcommand: None,
508 }],
509 };
510 let read_image_tool = ToolSpec {
511 provider_scope: None,
512 name: tau_proto::ToolName::new(READ_IMAGE_TOOL_NAME),
513 model_visible_name: None,
514 description: Some("Read one local image for visual inspection.".to_owned()),
515 tool_type: tau_proto::ToolType::Function,
516 parameters: Some(serde_json::json!({
517 "type": "object",
518 "properties": {
519 "path": {
520 "type": "string",
521 "description": "Path to one local PNG, JPEG, or WebP image"
522 },
523 "mode": {
524 "type": "string",
525 "enum": ["high", "overview"],
526 "default": "high",
527 "description": "Local preparation profile. `high` (the default) preserves the existing 2048-side/2500-patch bounds. `overview` is experimental, intended only for coarse inspection, and is bounded to 1024 pixels on a side and 600 32px patches."
528 },
529 "region": {
530 "type": "object",
531 "description": "Optional crop in pixels of the EXIF-oriented source raster, before mode resizing. Uses a top-left origin and half-open extents.",
532 "properties": {
533 "x": {
534 "type": "integer",
535 "minimum": 0,
536 "maximum": u32::MAX
537 },
538 "y": {
539 "type": "integer",
540 "minimum": 0,
541 "maximum": u32::MAX
542 },
543 "width": {
544 "type": "integer",
545 "minimum": 1,
546 "maximum": u32::MAX
547 },
548 "height": {
549 "type": "integer",
550 "minimum": 1,
551 "maximum": u32::MAX
552 }
553 },
554 "required": ["x", "y", "width", "height"],
555 "additionalProperties": false
556 }
557 },
558 "required": ["path"],
559 "additionalProperties": false
560 })),
561 format: None,
562 tags: tool_tags(&[
563 "shell:read",
564 "shell:read:image",
565 "provider-content:image",
566 tau_proto::TURN_DATA_FETCH_TOOL_TAG,
567 ]),
568 enabled_by_default: true,
569 background_support: Some(tau_proto::BackgroundSupport::Never),
570 examples: vec![ToolExample {
571 id: "read-image".to_owned(),
572 title: Some("Inspect a screenshot".to_owned()),
573 arguments: CborValue::Map(vec![example_field("path", example_text("screenshot.png"))]),
574 note: None,
575 subcommand: None,
576 }],
577 };
578 let export_tool = ToolSpec {
579 provider_scope: None,
580 name: tau_proto::ToolName::new(EXPORT_TOOL_NAME),
581 model_visible_name: None,
582 description: Some(
583 "Export one local regular file to the shared content-addressed artifact store. \
584 Originals are limited to 16 MiB. A successful export returns the BLAKE3 key and \
585 size and renews shared artifact age, including for duplicate bytes. Original bytes \
586 persist independently of ephemeral session transcripts."
587 .to_owned(),
588 ),
589 tool_type: tau_proto::ToolType::Function,
590 parameters: Some(serde_json::json!({
591 "type": "object",
592 "properties": {
593 "path": {"type": "string", "description": "Path to one local regular file"}
594 },
595 "required": ["path"],
596 "additionalProperties": false
597 })),
598 format: None,
599 tags: tool_tags(&[
600 "shell:read",
601 "artifact:write",
602 tau_proto::TURN_DATA_FETCH_TOOL_TAG,
603 ]),
604 enabled_by_default: true,
605 background_support: None,
606 examples: vec![ToolExample {
607 id: "export-artifact".to_owned(),
608 title: Some("Export an original".to_owned()),
609 arguments: CborValue::Map(vec![example_field("path", example_text("output.png"))]),
610 note: Some("The returned key can be imported by another session.".to_owned()),
611 subcommand: None,
612 }],
613 };
614 let import_tool = ToolSpec {
615 provider_scope: None,
616 name: tau_proto::ToolName::new(IMPORT_TOOL_NAME),
617 model_visible_name: None,
618 description: Some(
619 "Import one artifact key to a private, unpredictable, non-executable temporary file \
620 on this shell host. Size and digest are verified before success. Import does not \
621 renew retention age; pass the returned local path to read_image or filesystem tools."
622 .to_owned(),
623 ),
624 tool_type: tau_proto::ToolType::Function,
625 parameters: Some(serde_json::json!({
626 "type": "object",
627 "properties": {
628 "key": {
629 "type": "string",
630 "pattern": "^blake3:[0-9a-f]{64}$",
631 "description": "Canonical artifact key returned by export"
632 }
633 },
634 "required": ["key"],
635 "additionalProperties": false
636 })),
637 format: None,
638 tags: tool_tags(&[
639 "shell:read",
640 "artifact:read",
641 tau_proto::TURN_DATA_FETCH_TOOL_TAG,
642 ]),
643 enabled_by_default: true,
644 background_support: None,
645 examples: Vec::new(),
646 };
647 let edit_tool = ToolSpec {
648 provider_scope: None,
649 name: tau_proto::ToolName::new(EDIT_TOOL_NAME),
650 model_visible_name: None,
651 description: Some(
652 "Edit a file using line-oriented replacements. Each edit fully replaces \
653 the 1-based half-open `start_line`..`end_line_exclusive` range \
654 with `newText`. `start_line` is included and `end_line_exclusive` \
655 is excluded. Empty insertion ranges use \
656 `start_line == end_line_exclusive`; for example, `1..<1` inserts \
657 at the start of the file and `total_lines + 1 ..< total_lines + 1` \
658 appends at EOF. All ranges use the original file numbering as if \
659 applied simultaneously. Non-empty replacements are kept as whole \
660 lines. Ranges must be non-overlapping. Missing files are treated as \
661 empty and missing parent directories are created. Per-edit `context_line` \
662 must exactly match the original content of `start_line`. Use an empty \
663 context_line when `start_line` is the append slot past the end of the \
664 file."
665 .to_owned(),
666 ),
667 tool_type: tau_proto::ToolType::Function,
668 parameters: Some(serde_json::json!({
669 "type": "object",
670 "properties": {
671 "path": {
672 "type": "string",
673 "description": "Path to the file"
674 },
675 "edits": {
676 "type": "array",
677 "description": "One or more line ranges to replace in the original file",
678 "minItems": 1,
679 "maxItems": 100,
680 "items": {
681 "type": "object",
682 "properties": {
683 "start_line": {
684 "type": "integer",
685 "minimum": 1,
686 "description": "1-based included start line or insertion slot. Use 1 for the start of the file. To append at EOF, use total_lines + 1. Use together with end_line_exclusive."
687 },
688 "end_line_exclusive": {
689 "type": "integer",
690 "minimum": 1,
691 "description": "1-based excluded end line or insertion slot. Empty insertion ranges have end_line_exclusive == start_line. To replace read output lines A through B, use start_line A and end_line_exclusive B + 1. Use together with start_line."
692 },
693 "newText": {
694 "type": "string",
695 "description": "Replacement text. Non-empty replacements stay whole-line."
696 },
697 "context_line": {
698 "type": "string",
699 "description": "Exact expected content of the original start_line, including spaces and tabs. Use an empty context_line when start_line is the append slot past the end of the file. If it does not match, the edit fails and returns current line-numbered context around the expected context line."
700 }
701 },
702 "required": ["start_line", "end_line_exclusive", "newText", "context_line"],
703 "additionalProperties": false
704 }
705 }
706 },
707 "required": ["path", "edits"],
708 "additionalProperties": false
709 })),
710 format: None,
711 tags: tool_tags(&[
712 "shell:edit",
713 "shell:edit:line",
714 "shell:mutates-files",
715 tau_proto::TURN_MANIPULATOR_TOOL_TAG,
716 ]),
717 enabled_by_default: true,
718 background_support: None,
719 examples: vec![ToolExample {
720 id: "replace-lines".to_owned(),
721 title: Some("Replace one line".to_owned()),
722 arguments: CborValue::Map(vec![
723 example_field("path", example_text("src/main.rs")),
724 (
725 CborValue::Text("edits".to_owned()),
726 CborValue::Array(vec![CborValue::Map(vec![
727 example_field("start_line", example_int(10)),
728 example_field("end_line_exclusive", example_int(11)),
729 example_field("newText", example_text("replacement line")),
730 example_field("context_line", example_text("line being replaced")),
731 ])]),
732 ),
733 ]),
734 note: Some("end_line_exclusive is one past the last line replaced.".to_owned()),
735 subcommand: None,
736 }],
737 };
738 let apply_patch_tool = ToolSpec {
739 provider_scope: None,
740 name: tau_proto::ToolName::new(APPLY_PATCH_TOOL_NAME),
741 model_visible_name: None,
742 description: Some("Use the `apply_patch` tool to edit files.".to_owned()),
743 tool_type: tau_proto::ToolType::Custom,
744 parameters: None,
745 format: Some(tau_proto::ToolFormat::Text),
746 tags: tool_tags(&[
747 "shell:edit",
748 "shell:edit:apply_patch",
749 "shell:mutates-files",
750 tau_proto::TURN_MANIPULATOR_TOOL_TAG,
751 ]),
752 enabled_by_default: false,
753 background_support: None,
754 examples: Vec::new(),
755 };
756 let replace_tool = ToolSpec {
757 provider_scope: None,
758 name: tau_proto::ToolName::new(REPLACE_TOOL_NAME),
759 model_visible_name: Some(tau_proto::ToolName::new(EDIT_TOOL_NAME)),
760 description: Some(
761 "Replace exact text in one existing UTF-8 file. Each oldText must occur exactly \
762 once in the same original file snapshot; all edits apply atomically. Matching \
763 ignores only an initial UTF-8 BOM and normalizes CRLF/CR to LF. Use newText \
764 as an empty string to delete text."
765 .to_owned(),
766 ),
767 tool_type: tau_proto::ToolType::Function,
768 parameters: Some(serde_json::json!({
769 "type": "object",
770 "properties": {
771 "path": { "type": "string", "minLength": 1 },
772 "edits": {
773 "type": "array",
774 "minItems": 1,
775 "maxItems": 100,
776 "items": {
777 "type": "object",
778 "properties": {
779 "oldText": { "type": "string", "minLength": 1 },
780 "newText": { "type": "string" }
781 },
782 "required": ["oldText", "newText"],
783 "additionalProperties": false
784 }
785 }
786 },
787 "required": ["path", "edits"],
788 "additionalProperties": false
789 })),
790 format: None,
791 tags: tool_tags(&[
792 "shell:edit",
793 "shell:edit:replace",
794 "shell:mutates-files",
795 tau_proto::TURN_MANIPULATOR_TOOL_TAG,
796 ]),
797 enabled_by_default: false,
798 background_support: None,
799 examples: Vec::new(),
800 };
801 let dir_lock_tool = dir_lock_tool_spec(dir_lock_enabled);
802 let grep_tool = ToolSpec {
803 provider_scope: None,
804 name: tau_proto::ToolName::new(GREP_TOOL_NAME),
805 model_visible_name: None,
806 description: Some(
807 "Search file contents for a pattern using ripgrep. Patterns are literal by default; \
808 regex metacharacters like `|` require `regex: true`. Returns matching lines \
809 with file paths and line numbers. Respects .gitignore. Output is truncated at \
810 `limit` matches or 10 KiB of visible output. Visible-cap truncation provides a private saved-output path, or explicit unavailable metadata when storage fails; limit-only and per-line truncation retain native metadata. Long lines are truncated to 500 chars."
811 .to_owned(),
812 ),
813 tool_type: tau_proto::ToolType::Function,
814 parameters: Some(serde_json::json!({
815 "type": "object",
816 "properties": {
817 "pattern": {
818 "type": "string",
819 "description": "Search pattern. Treated as a literal string by default. Set `regex: true` to interpret as a regex."
820 },
821 "path": {
822 "type": "string",
823 "description": "Directory or file to search (default: current directory)"
824 },
825 "glob": {
826 "type": "string",
827 "description": "Filter files by glob pattern, e.g. '*.ts' or '**/*.rs'"
828 },
829 "ignoreCase": {
830 "type": "boolean",
831 "description": "Case-insensitive search (default: false)"
832 },
833 "regex": {
834 "type": "boolean",
835 "description": "Interpret `pattern` as a regex instead of a literal string (default: false)"
836 },
837 "context": {
838 "type": "integer",
839 "description": "Number of lines to show before and after each match (default: 0, max: 20)"
840 },
841 "limit": {
842 "type": "integer",
843 "description": "Maximum number of matches to return (default: 100, max: 2000)"
844 }
845 },
846 "required": ["pattern"],
847 "additionalProperties": false
848 })),
849 format: None,
850 tags: tool_tags(&[
851 "shell:read",
852 "shell:search",
853 tau_proto::TURN_DATA_FETCH_TOOL_TAG,
854 ]),
855 enabled_by_default: true,
856 background_support: None,
857 examples: vec![ToolExample {
858 id: "search-literal".to_owned(),
859 title: Some("Search literal text".to_owned()),
860 arguments: CborValue::Map(vec![
861 example_field("pattern", example_text("TODO")),
862 example_field("path", example_text("src")),
863 example_field("glob", example_text("**/*.rs")),
864 ]),
865 note: Some("Set regex=true only when pattern is a regular expression.".to_owned()),
866 subcommand: None,
867 }],
868 };
869 let find_tool = ToolSpec {
870 provider_scope: None,
871 name: tau_proto::ToolName::new(FIND_TOOL_NAME),
872 model_visible_name: None,
873 description: Some(
874 "Search for files by glob pattern. Returns only file paths (directories are \
875 never included, even with '**/*') relative to the search directory. Respects \
876 .gitignore. Output is truncated at `limit` results or 10 KiB of visible output. Visible-cap truncation provides a private saved-output path, or explicit unavailable metadata when storage fails; limit-only truncation retains native metadata. Use the ls tool \
877 if you want to see directory entries."
878 .to_owned(),
879 ),
880 tool_type: tau_proto::ToolType::Function,
881 parameters: Some(serde_json::json!({
882 "type": "object",
883 "properties": {
884 "pattern": {
885 "type": "string",
886 "description": "Glob pattern matched against file paths relative to `path`. `**` matches any number of intermediate directories, including zero — so `**/*.rs` finds both top-level `a.rs` and nested `src/a.rs`. Directories are not returned, even with `**/*`."
887 },
888 "path": {
889 "type": "string",
890 "description": "Directory to search (default: current directory)"
891 },
892 "limit": {
893 "type": "integer",
894 "description": "Maximum number of results to return (default: 1000, max: 2000)"
895 }
896 },
897 "required": ["pattern"],
898 "additionalProperties": false
899 })),
900 format: None,
901 tags: tool_tags(&[
902 "shell:read",
903 "shell:search",
904 tau_proto::TURN_DATA_FETCH_TOOL_TAG,
905 ]),
906 enabled_by_default: true,
907 background_support: None,
908 examples: vec![ToolExample {
909 id: "find-rust-files".to_owned(),
910 title: Some("Find files by glob".to_owned()),
911 arguments: CborValue::Map(vec![
912 example_field("pattern", example_text("**/*.rs")),
913 example_field("path", example_text("crates")),
914 ]),
915 note: None,
916 subcommand: None,
917 }],
918 };
919 let ls_tool = ToolSpec {
920 provider_scope: None,
921 name: tau_proto::ToolName::new(LS_TOOL_NAME),
922 model_visible_name: None,
923 description: Some(
924 "List directory contents. Returns entries sorted alphabetically, with '/' suffix \
925 for directories. Includes dotfiles. Output lines are prefixed with 1-based \
926 entry numbers plus flags such as `escaped`, `invalid-utf8`, or `truncated`; \
927 output is capped at `limit` entries, 2000 lines, or 10 KiB of visible output with saved-output metadata and standard truncation headers. \
928 When `limit_reached` is true, entries are a bounded filesystem-order sample sorted \
929 for display, not a complete alphabetic prefix."
930 .to_owned(),
931 ),
932 tool_type: tau_proto::ToolType::Function,
933 parameters: Some(serde_json::json!({
934 "type": "object",
935 "properties": {
936 "path": {
937 "type": "string",
938 "description": "Directory to list (default: current directory)"
939 },
940 "limit": {
941 "type": "integer",
942 "minimum": 1,
943 "description": "Maximum number of entries to return (default: 500, max: 2001)"
944 }
945 },
946 "additionalProperties": false
947 })),
948 format: None,
949 tags: tool_tags(&[
950 "shell:read",
951 "shell:list",
952 tau_proto::TURN_DATA_FETCH_TOOL_TAG,
953 ]),
954 enabled_by_default: true,
955 background_support: None,
956 examples: vec![ToolExample {
957 id: "list-directory".to_owned(),
958 title: Some("List a directory".to_owned()),
959 arguments: CborValue::Map(vec![example_field("path", example_text("src"))]),
960 note: None,
961 subcommand: None,
962 }],
963 };
964 let workdir_tool = ToolSpec {
965 provider_scope: None,
966 name: tau_proto::ToolName::new(WORKDIR_TOOL_NAME),
967 model_visible_name: None,
968 description: Some(
969 "Read or change your durable workdir. Omit `path` \
970 to read the current path and availability. A provided path is resolved from the \
971 last committed workdir, validated, canonicalized, and persisted. Do not combine a \
972 workdir change with shell or filesystem calls that rely on the new directory."
973 .to_owned(),
974 ),
975 tool_type: tau_proto::ToolType::Function,
976 parameters: Some(serde_json::json!({
977 "type": "object",
978 "properties": { "path": { "type": "string", "minLength": 1, "description": "Optional directory to persist as this instance's workdir" } },
979 "additionalProperties": false
980 })),
981 format: None,
982 tags: tool_tags(&["shell:workdir", tau_proto::TURN_MANIPULATOR_TOOL_TAG]),
983 enabled_by_default: true,
984 background_support: None,
985 examples: vec![ToolExample {
986 id: "change-directory".to_owned(),
987 title: Some("Change directory".to_owned()),
988 arguments: CborValue::Map(vec![example_field("path", example_text("crates/tau"))]),
989 note: None,
990 subcommand: None,
991 }],
992 };
993 let shell_tool = ToolSpec {
994 provider_scope: None,
995 name: tau_proto::ToolName::new(SHELL_TOOL_NAME),
996 model_visible_name: None,
997 description: Some(
998 "Execute a shell command via `sh -c`. When directory locking is enabled, commands \
999 are inferred read-write only while the agent holds a matching `dir_lock`; otherwise \
1000 they are read-only. When directory locking is disabled, shell commands run read-write. \
1001 Non-zero exits and timeouts are returned as structured command results with output details. \
1002 The native `output` body is capped at 2000 lines / 15 KiB; small result \
1003 metadata is deliberately outside that budget, which does not cap fully \
1004 rendered provider text. Truncated output keeps the first 1000 and last 1000 lines \
1005 separated by a literal `...` line. Output lines are prefixed with `out ` \
1006 for stdout or `err ` for stderr; missing trailing newlines are marked, e.g. \
1007 `out(no_nl)`; CRLF and CR line endings are marked as `out(crlf)` \
1008 or `out(cr)`. Invalid UTF-8 is shown with Unicode replacement characters and \
1009 an `invalid-utf8` line flag. Lines that would exceed the 15 KiB output budget \
1010 are marker-only, e.g. `err(truncated)`. Truncated results include complete totals, a warning, and normally an exact temporary path to up to 16 MiB of rendered output; output beyond that saved cap is explicitly marked incomplete, while platforms or filesystems that cannot enforce private storage report `saved_output_unavailable: true`. Byte totals and artifacts count the complete rendered UTF-8 records, including stream prefixes, flags, and separators, rather than raw process bytes. \
1011 Stdin is closed and commands cannot receive interactive input. Stdout and stderr may be TTY-backed even though no controlling terminal exists. Use explicit noninteractive flags/messages; do not launch prompts, pagers, or editors. \
1012 Commands taking longer than 5 seconds include duration metadata. Prefer dedicated \
1013 tools like `read`, `grep`, and `find` when they fit."
1014 .to_owned(),
1015 ),
1016 tool_type: tau_proto::ToolType::Function,
1017 parameters: Some(serde_json::json!({
1018 "type": "object",
1019 "properties": {
1020 "command": {
1021 "type": "string",
1022 "description": "The shell command to execute"
1023 },
1024 "timeout": {
1025 "type": "integer",
1026 "minimum": 0,
1027 "description": "Timeout in seconds. The command is killed if it exceeds this. Default: 300"
1028 },
1029 "cwd": {
1030 "type": "string",
1031 "description": "Working directory for this invocation only. Relative paths resolve from this shell instance's remembered workdir; omission uses the remembered workdir. This does not change later calls; use workdir in an earlier turn to change later calls."
1032 }
1033 },
1034 "required": ["command"],
1035 "additionalProperties": false
1036 })),
1037 format: None,
1038 tags: tool_tags(&[
1039 "shell:exec",
1040 "shell:exec:generic",
1041 tau_proto::TURN_MANIPULATOR_TOOL_TAG,
1042 ]),
1043 enabled_by_default: true,
1044 background_support: None,
1045 examples: vec![ToolExample {
1046 id: "run-command".to_owned(),
1047 title: Some("Run a command".to_owned()),
1048 arguments: CborValue::Map(vec![
1049 example_field("command", example_text("cargo test -p dpc-tau-core")),
1050 example_field("timeout", example_int(300)),
1051 ]),
1052 note: Some("For file edits, prefer apply_patch when available.".to_owned()),
1053 subcommand: None,
1054 }],
1055 };
1056 let gpt_shell_tool = ToolSpec {
1057 provider_scope: None,
1058 name: tau_proto::ToolName::new(GPT_SHELL_TOOL_NAME),
1059 model_visible_name: Some(tau_proto::ToolName::new("shell_command")),
1060 description: Some(
1061 "Run a shell command. The native `output` body is capped at 2000 lines / 15 KiB; \
1062 small result metadata is deliberately outside that budget, which does not cap fully \
1063 rendered provider text. \
1064 truncated results normally provide an exact temporary path to up to 16 MiB of rendered output and mark an incomplete saved artifact honestly; private-storage failures instead report `saved_output_unavailable: true`. \
1065 Output lines are prefixed with `out ` for stdout or `err ` for stderr; missing \
1066 trailing newlines are marked with `(no_nl)`. Byte totals and artifacts count the complete rendered UTF-8 records, including stream prefixes, flags, and separators, rather than raw process bytes. Stdin is closed and commands cannot receive interactive input. Stdout and stderr may be TTY-backed even though no controlling terminal exists. Use explicit noninteractive flags/messages; do not launch prompts, pagers, or editors. For file changes, prefer apply_patch."
1067 .to_owned(),
1068 ),
1069 tool_type: tau_proto::ToolType::Function,
1070 parameters: Some(serde_json::json!({
1071 "type": "object",
1072 "properties": {
1073 "command": {
1074 "type": "string",
1075 "description": "The shell command to execute"
1076 },
1077 "timeout": {
1078 "type": "integer",
1079 "description": "Timeout in seconds. The command is killed if it exceeds this. Default: 300"
1080 },
1081 "workdir": {
1082 "type": "string",
1083 "description": "Optional working directory for this shell_command invocation only. Relative paths resolve from this shell instance's remembered persistent workdir; omission uses that remembered workdir. This does not change later calls; use the separate top-level workdir(path) tool in an earlier turn to change later calls."
1084 }
1085 },
1086 "required": ["command"],
1087 "additionalProperties": false
1088 })),
1089 format: None,
1090 tags: tool_tags(&[
1091 "shell:exec",
1092 "shell:exec:shell_command",
1093 tau_proto::TURN_MANIPULATOR_TOOL_TAG,
1094 ]),
1095 enabled_by_default: false,
1096 background_support: None,
1097 examples: vec![ToolExample {
1098 id: "run-command".to_owned(),
1099 title: Some("Run a command".to_owned()),
1100 arguments: CborValue::Map(vec![
1101 example_field("command", example_text("cargo test -p dpc-tau-core")),
1102 example_field("timeout", example_int(300)),
1103 ]),
1104 note: Some("For file edits, prefer apply_patch when available.".to_owned()),
1105 subcommand: None,
1106 }],
1107 };
1108 let builtin_tools = [
1109 read_tool,
1110 read_image_tool,
1111 export_tool,
1112 import_tool,
1113 edit_tool,
1114 replace_tool,
1115 apply_patch_tool,
1116 dir_lock_tool,
1117 grep_tool,
1118 find_tool,
1119 ls_tool,
1120 workdir_tool,
1121 shell_tool,
1122 gpt_shell_tool,
1123 ];
1124 tools.extend(builtin_tools);
1125 tools
1126}
1127
1128fn run_impl<R, W>(
1129 reader: R,
1130 writer: W,
1131 discovery_policy: DiscoverySourcePolicy,
1132 runtime_cwd_source: RuntimeCwdSource,
1133) -> Result<(), Box<dyn Error>>
1134where
1135 R: Read + Send + 'static,
1136 W: Write + Send + 'static,
1137{
1138 let initial_config = ExtConfig::default();
1139 let mut runtime = tau_client::TauExtensionRunner::new(ShellExtension {
1140 initial_config: initial_config.clone(),
1141 })
1142 .start_manual_loop_with_state(reader, writer, |handle| match runtime_cwd_source {
1143 RuntimeCwdSource::Process => ShellRuntime::new_with_artifacts(
1144 Output::client(handle.clone()),
1145 ArtifactTransferManager::new(tau_client::ArtifactClient::new(handle)),
1146 initial_config,
1147 discovery_policy,
1148 ),
1149 #[cfg(any(test, feature = "echo-agent"))]
1150 RuntimeCwdSource::Fixture(fixture_cwd) => ShellRuntime::new_for_test_harness(
1151 Output::client(handle),
1152 initial_config,
1153 discovery_policy,
1154 fixture_cwd,
1155 ),
1156 })?;
1157
1158 let waker = runtime.waker();
1159 runtime.state().install_waker(waker);
1160 let loop_result = run_shell_manual_loop(&mut runtime);
1161
1162 runtime.state_mut().final_shutdown();
1168 let finish_result = runtime.finish().map(|_| ());
1169 match (loop_result, finish_result) {
1170 (Ok(()), Ok(())) => Ok(()),
1171 (_, Err(error)) => Err(Box::new(error)),
1172 (Err(error), _) => Err(Box::new(error)),
1173 }
1174}
1175
1176fn run_shell_manual_loop(
1177 runtime: &mut tau_client::ManualExtensionRuntime<ShellRuntime>,
1178) -> tau_client::ClientResult<()> {
1179 loop {
1180 runtime.state_mut().drain_artifact_commands()?;
1181 runtime.state().take_mandatory_output_failure()?;
1182 match runtime.try_recv()? {
1183 tau_client::ManualRuntimePoll::Message(message) => {
1184 if let tau_proto::HarnessOutputMessage::ArtifactResult(result) = message {
1185 runtime.state_mut().handle_artifact_result(*result)?;
1186 continue;
1187 }
1188 match runtime.dispatch_one(message)? {
1189 tau_client::DispatchOutcome::Continue => {}
1190 tau_client::DispatchOutcome::StopRequested
1191 | tau_client::DispatchOutcome::Disconnect(_) => return Ok(()),
1192 }
1193 }
1194 tau_client::ManualRuntimePoll::InputClosed => return Ok(()),
1195 tau_client::ManualRuntimePoll::Empty => runtime.wait_for_wake(),
1196 }
1197 }
1198}
1199
1200struct ShellExtension {
1201 initial_config: ExtConfig,
1202}
1203
1204impl tau_client::TauExtension for ShellExtension {
1205 type State = ShellRuntime;
1206
1207 fn name(&self) -> &'static str {
1208 "tau-ext-shell"
1209 }
1210
1211 fn register(self, builder: &mut tau_client::ExtensionBuilder<Self::State>) {
1212 let tools = registered_tool_specs(self.initial_config.dir_lock.enable);
1213
1214 let shell_tool_group = tau_proto::ToolGroup {
1218 name: tau_proto::ToolGroupName::new("shell"),
1219 prompt_fragment: None,
1220 };
1221 let test_tool_group = tau_proto::ToolGroup {
1222 name: tau_proto::ToolGroupName::new("test"),
1223 prompt_fragment: None,
1224 };
1225
1226 for tool in tools {
1227 let tool_group = if tool.name.as_str() == "echo" {
1228 test_tool_group.clone()
1229 } else {
1230 shell_tool_group.clone()
1231 };
1232 builder.tool_with_group_and_prompt_fragment(tool, Some(tool_group), None, |cx| {
1233 let local_tool_name = cx.local_tool_name().clone();
1234 cx.state
1235 .handle_scoped_tool_started(cx.invoke.clone(), &local_tool_name)
1236 });
1237 }
1238 builder
1239 .register_context_provider()
1240 .register_session_context_provider()
1241 .publish_prompt_fragment(shell_workdir_prompt_fragment(&self.initial_config.shell))
1242 .publish_actions(shell_action_schema())
1243 .on_live::<tau_proto::ToolCancelRequest>(|cx| {
1244 cx.state
1245 .handle_event(Event::ToolCancelRequest(cx.event.clone()), false)
1246 })
1247 .on_raw_live(
1248 tau_proto::EventSelector::Exact(tau_proto::EventName::ACTION_INVOKE),
1249 |cx| cx.state.handle_event(cx.event().clone(), false),
1250 )
1251 .on_restore::<tau_proto::SessionStarted>(|cx| {
1252 cx.state
1253 .handle_event(Event::SessionStarted(cx.event.clone()), true)
1254 })
1255 .on_live::<tau_proto::SessionStarted>(|cx| {
1256 cx.state
1257 .handle_event(Event::SessionStarted(cx.event.clone()), false)
1258 })
1259 .on_restore::<tau_proto::SessionAgentLoaded>(|cx| {
1260 cx.state
1261 .handle_event(Event::SessionAgentLoaded(cx.event.clone()), true)
1262 })
1263 .on_live::<tau_proto::SessionAgentLoaded>(|cx| {
1264 cx.state
1265 .handle_event(Event::SessionAgentLoaded(cx.event.clone()), false)
1266 })
1267 .on_restore::<tau_proto::SessionAgentUnloaded>(|cx| {
1268 cx.state
1269 .handle_event(Event::SessionAgentUnloaded(cx.event.clone()), true)
1270 })
1271 .on_live::<tau_proto::SessionAgentUnloaded>(|cx| {
1272 cx.state
1273 .handle_event(Event::SessionAgentUnloaded(cx.event.clone()), false)
1274 })
1275 .on_live::<tau_proto::AgentReplayComplete>(|cx| {
1276 cx.state
1277 .handle_event(Event::AgentReplayComplete(cx.event.clone()), false)
1278 })
1279 .on_restore::<tau_proto::AgentMetadataSet>(|cx| {
1280 cx.state
1281 .handle_event(Event::AgentMetadataSet(cx.event.clone()), true)
1282 })
1283 .on_live::<tau_proto::AgentMetadataSet>(|cx| {
1284 cx.state
1285 .handle_event(Event::AgentMetadataSet(cx.event.clone()), false)
1286 })
1287 .on_restore::<tau_proto::AgentMetadataUnset>(|cx| {
1288 cx.state
1289 .handle_event(Event::AgentMetadataUnset(cx.event.clone()), true)
1290 })
1291 .on_live::<tau_proto::AgentMetadataUnset>(|cx| {
1292 cx.state
1293 .handle_event(Event::AgentMetadataUnset(cx.event.clone()), false)
1294 })
1295 .on_live::<tau_proto::SessionShutdown>(|cx| {
1296 cx.state
1297 .handle_event(Event::SessionShutdown(cx.event.clone()), false)
1298 })
1299 .on_live::<tau_proto::StartAgentAccepted>(|cx| {
1300 cx.state
1301 .handle_event(Event::StartAgentAccepted(cx.event.clone()), false)
1302 })
1303 .on_live::<tau_proto::StartAgentResult>(|cx| {
1304 cx.state
1305 .handle_event(Event::StartAgentResult(cx.event.clone()), false)
1306 })
1307 .on_live::<tau_proto::UiShellCommand>(|cx| {
1308 cx.state
1309 .handle_event(Event::UiShellCommand(cx.event.clone()), false)
1310 })
1311 .configure_raw(|cx| {
1312 let cfg = cx.parse_config::<ExtConfig>()?;
1313 cx.state.apply_config(
1314 cx.configure.instance_name.clone(),
1315 cx.configure.tool_prefix.clone(),
1316 cfg,
1317 )
1318 })
1319 .ready_message("filesystem and shell tools ready");
1320 }
1321}
1322
1323fn apply_working_directory(
1324 current: &ExtConfig,
1325 next: &ExtConfig,
1326 runtime_started: bool,
1327) -> Result<(), String> {
1328 match (¤t.working_directory, &next.working_directory) {
1329 (None, Some(_)) if runtime_started => Err(
1330 "ext-shell working_directory cannot be set after runtime events have started"
1331 .to_owned(),
1332 ),
1333 (None, Some(working_directory)) => set_process_working_directory(working_directory),
1334 (Some(current), Some(next)) if current == next => Ok(()),
1335 (Some(current), Some(next)) => Err(format!(
1336 "ext-shell working_directory cannot be changed after startup (current: {}, requested: {})",
1337 current.display(),
1338 next.display()
1339 )),
1340 _ => Ok(()),
1341 }
1342}
1343
1344fn set_process_working_directory(working_directory: &Path) -> Result<(), String> {
1345 std::env::set_current_dir(working_directory).map_err(|err| {
1346 format!(
1347 "failed to set ext-shell working_directory to {}: {err}",
1348 working_directory.display()
1349 )
1350 })
1351}
1352
1353fn dir_lock_tool_spec(enabled_by_default: bool) -> ToolSpec {
1354 let tags = if enabled_by_default {
1355 tool_tags(&["shell:lock", tau_proto::TURN_WAIT_TOOL_TAG])
1356 } else {
1357 tool_tags(&[tau_proto::TURN_WAIT_TOOL_TAG])
1358 };
1359 ToolSpec {
1360 provider_scope: None,
1361 name: tau_proto::ToolName::new(DIR_LOCK_TOOL_NAME),
1362 model_visible_name: None,
1363 description: Some(
1364 "Lock or unlock a directory and its contents for updates. Waits for the lock when \
1365 necessary."
1366 .to_owned(),
1367 ),
1368 tool_type: tau_proto::ToolType::Function,
1369 parameters: Some(serde_json::json!({
1370 "type": "object",
1371 "properties": {
1372 "command": {
1373 "type": "string",
1374 "enum": ["update", "unlock"],
1375 "description": "Lock or unlock the directory for updates"
1376 },
1377 "directory": {
1378 "type": "string",
1379 "description": "Existing directory to canonicalize before locking"
1380 },
1381 "owner_agent_id": {
1382 "type": "string",
1383 "description": "Optional owner agent id for force-unlocking a manual lock held by another agent"
1384 }
1385 },
1386 "required": ["command", "directory"],
1387 "additionalProperties": false
1388 })),
1389 format: None,
1390 tags,
1391 enabled_by_default,
1392 background_support: None,
1393 examples: vec![
1394 ToolExample {
1395 id: "update-lock".to_owned(),
1396 title: Some("Acquire update lock".to_owned()),
1397 arguments: CborValue::Map(vec![
1398 example_field("command", example_text("update")),
1399 example_field("directory", example_text(".")),
1400 ]),
1401 note: Some(
1402 "Acquire before making file changes when directory locking is enabled."
1403 .to_owned(),
1404 ),
1405 subcommand: Some(ToolExampleSelector {
1406 path: vec!["command".to_owned()],
1407 value: example_text("update"),
1408 }),
1409 },
1410 ToolExample {
1411 id: "unlock".to_owned(),
1412 title: Some("Release update lock".to_owned()),
1413 arguments: CborValue::Map(vec![
1414 example_field("command", example_text("unlock")),
1415 example_field("directory", example_text(".")),
1416 ]),
1417 note: None,
1418 subcommand: Some(ToolExampleSelector {
1419 path: vec!["command".to_owned()],
1420 value: example_text("unlock"),
1421 }),
1422 },
1423 ],
1424 }
1425}
1426
1427fn shell_action_schema() -> tau_actions::ActionSchema {
1428 tau_actions::ActionSchema {
1429 version: tau_actions::ACTION_SCHEMA_VERSION,
1430 roots: vec![tau_actions::ActionCommand {
1431 name: ":shell-dir-force-unlock".to_owned(),
1432 description: "Force-release ext-shell manual directory locks overlapping a directory"
1433 .to_owned(),
1434 action_id: Some(SHELL_DIR_FORCE_UNLOCK_ACTION_ID.to_owned()),
1435 args: vec![tau_actions::ActionArg {
1436 name: "directory".to_owned(),
1437 description: "Existing directory whose overlapping manual locks should be released"
1438 .to_owned(),
1439 required: true,
1440 suggestions: Vec::new(),
1441 kind: tau_actions::ActionArgKind::RestString,
1442 }],
1443 children: Vec::new(),
1444 }],
1445 }
1446}
1447
1448fn dispatch_action_invoke(invoke: ActionInvoke, lock_manager: &DirLockManager) -> Event {
1449 if invoke.action_id != SHELL_DIR_FORCE_UNLOCK_ACTION_ID {
1450 return action_error(invoke, "unknown shell action".to_owned());
1451 }
1452 let Some(directory) = invoke.argv.first().map(String::as_str) else {
1453 return action_error(invoke, "missing directory argument".to_owned());
1454 };
1455 let dir = match crate::dir_lock::canonical_existing_dir(Path::new(directory)) {
1456 Ok(dir) => dir,
1457 Err(message) => return action_error(invoke, message),
1458 };
1459 let removed = match lock_manager.force_unlock_overlapping(&dir) {
1460 Ok(removed) => removed,
1461 Err(message) => {
1462 return action_error(invoke, format!("dir_lock backend error: {message}"));
1463 }
1464 };
1465 if removed.is_empty() {
1466 return action_error(
1467 invoke,
1468 format!("no manual directory locks overlap {}", dir.display()),
1469 );
1470 }
1471
1472 let mut lines = vec![format!(
1473 "Force-unlocked {} manual directory lock(s) overlapping {}.",
1474 removed.len(),
1475 dir.display()
1476 )];
1477 for entry in removed {
1478 lines.push(format!("{} owner={}", entry.dir.display(), entry.owner));
1479 }
1480 Event::ActionResultReported(ActionResult {
1481 invocation_id: invoke.invocation_id,
1482 action_id: invoke.action_id,
1483 output: ActionOutput::Text {
1484 text: lines.join("\n"),
1485 },
1486 })
1487}
1488
1489fn action_error(invoke: ActionInvoke, message: String) -> Event {
1490 Event::ActionErrorReported(ActionError {
1491 invocation_id: invoke.invocation_id,
1492 action_id: invoke.action_id,
1493 message,
1494 details: None,
1495 })
1496}
1497
1498fn rewrite_invoke_for_cwd(
1499 mut invoke: tau_proto::ToolStarted,
1500 base: &Path,
1501) -> tau_proto::ToolStarted {
1502 if invoke.tool_name == WORKDIR_TOOL_NAME {
1503 return invoke;
1504 }
1505 let field = match invoke.tool_name.as_str() {
1506 SHELL_TOOL_NAME => path_crate_tools::ShellSurface::Generic.directory_argument(),
1507 GPT_SHELL_TOOL_NAME => path_crate_tools::ShellSurface::ChatGpt.directory_argument(),
1508 READ_TOOL_NAME | READ_IMAGE_TOOL_NAME | EXPORT_TOOL_NAME | EDIT_TOOL_NAME
1509 | REPLACE_TOOL_NAME | FIND_TOOL_NAME | GREP_TOOL_NAME | LS_TOOL_NAME => "path",
1510 DIR_LOCK_TOOL_NAME => "directory",
1511 _ => return invoke,
1512 };
1513 let explicit_path = cbor_optional_text(&invoke.arguments, field);
1514 if explicit_path.is_none() && cbor_has_field(&invoke.arguments, field) {
1515 return invoke;
1517 }
1518 let Some(path) = explicit_path
1519 .clone()
1520 .or_else(|| matches!(field, "path").then(|| ".".to_owned()))
1521 .or_else(|| {
1522 matches!(
1523 invoke.tool_name.as_str(),
1524 SHELL_TOOL_NAME | GPT_SHELL_TOOL_NAME
1525 )
1526 .then(|| base.display().to_string())
1527 })
1528 else {
1529 return invoke;
1530 };
1531 let path = PathBuf::from(path);
1532 let absolute = if path.is_absolute() {
1533 path
1534 } else {
1535 base.join(path)
1536 };
1537 if let Some(canonical) = canonicalize_existing_dir_for_cwd_field(&absolute, field) {
1538 set_cbor_text_field(
1539 &mut invoke.arguments,
1540 field,
1541 canonical.display().to_string(),
1542 );
1543 } else {
1544 set_cbor_text_field(&mut invoke.arguments, field, absolute.display().to_string());
1545 }
1546 invoke
1547}
1548
1549fn canonicalize_existing_dir_for_cwd_field(path: &Path, field: &str) -> Option<PathBuf> {
1550 (field == "cwd" || field == "workdir" || field == "directory" || field == "path")
1551 .then(|| path.canonicalize().ok())
1552 .flatten()
1553 .filter(|path| path.is_dir())
1554}
1555
1556fn cbor_optional_text(arguments: &CborValue, field: &str) -> Option<String> {
1557 let CborValue::Map(entries) = arguments else {
1558 return None;
1559 };
1560 entries.iter().find_map(|(key, value)| match (key, value) {
1561 (CborValue::Text(key), CborValue::Text(value)) if key == field => Some(value.clone()),
1562 _ => None,
1563 })
1564}
1565
1566fn cbor_has_field(arguments: &CborValue, field: &str) -> bool {
1567 let CborValue::Map(entries) = arguments else {
1568 return false;
1569 };
1570 entries
1571 .iter()
1572 .any(|(key, _)| matches!(key, CborValue::Text(key) if key == field))
1573}
1574
1575fn set_cbor_text_field(arguments: &mut CborValue, field: &str, value: String) {
1576 let CborValue::Map(entries) = arguments else {
1577 return;
1578 };
1579 if let Some((_, existing)) = entries
1580 .iter_mut()
1581 .find(|(key, _)| matches!(key, CborValue::Text(key) if key == field))
1582 {
1583 *existing = CborValue::Text(value);
1584 } else {
1585 entries.push((CborValue::Text(field.to_owned()), CborValue::Text(value)));
1586 }
1587}
1588
1589#[expect(
1590 clippy::too_many_arguments,
1591 reason = "admission receives independently owned scheduler, policy, lifecycle, cwd, and Artifact routes"
1592)]
1593fn schedule_tool_started(
1594 (invoke, local_tool_name): (tau_proto::ToolStarted, &tau_proto::ToolName),
1595 scheduler: &WorkScheduler,
1596 tx: &Output,
1597 config: ExtConfig,
1598 lock_manager: DirLockManager,
1599 cancellation: ToolCancellationState,
1600 cwd_state: CwdState,
1601 artifact_control: ArtifactTransferControl,
1602) -> Result<
1603 (),
1604 Box<(
1605 tool_started_identity::ToolStartedIdentity,
1606 crate::display::ToolFailure,
1607 )>,
1608> {
1609 let (identity, arguments) =
1610 tool_started_identity::ToolStartedIdentity::split(invoke, local_tool_name.clone());
1611 let tx = tx.scoped_tool(
1612 identity.local_tool_name.clone(),
1613 identity.wire_tool_name.clone(),
1614 );
1615 let workdir_snapshot = cwd_state.snapshot(&identity.agent_id).map_err(|message| {
1616 Box::new((
1617 identity.clone(),
1618 path_crate_display::ToolFailure::new(message),
1619 ))
1620 })?;
1621 if matches!(workdir_snapshot, WorkdirSnapshot::Invalid)
1622 && identity.local_tool_name != WORKDIR_TOOL_NAME
1623 {
1624 return Err(Box::new((
1625 identity,
1626 path_crate_display::ToolFailure::new(
1627 "remembered workdir metadata is invalid; repair it with an absolute workdir path",
1628 ),
1629 )));
1630 }
1631 if matches!(workdir_snapshot, WorkdirSnapshot::ReplayFailed) {
1632 return Err(Box::new((
1633 identity,
1634 path_crate_display::ToolFailure::new(
1635 "workdir replay failed for this agent; reload the agent before retrying",
1636 ),
1637 )));
1638 }
1639 if matches!(workdir_snapshot, WorkdirSnapshot::Invalid) {
1640 let requested = cbor_optional_text(&arguments, "path");
1641 if !requested
1642 .as_deref()
1643 .is_none_or(|path| Path::new(path).is_absolute())
1644 {
1645 return Err(Box::new((
1646 identity,
1647 path_crate_display::ToolFailure::new(
1648 "remembered workdir metadata is invalid; repair it with an absolute workdir path",
1649 ),
1650 )));
1651 }
1652 }
1653 let mut invoke = identity.clone().into_local_started(arguments);
1654 invoke = match &workdir_snapshot {
1655 WorkdirSnapshot::Valid(cwd) => rewrite_invoke_for_cwd(invoke, cwd),
1656 WorkdirSnapshot::Invalid => invoke,
1657 WorkdirSnapshot::ReplayFailed => unreachable!("replay failures return above"),
1658 };
1659 if invoke.tool_name == WORKDIR_TOOL_NAME
1660 && cbor_optional_text(&invoke.arguments, "path").is_some()
1661 {
1662 let base = match &workdir_snapshot {
1663 WorkdirSnapshot::Valid(path) => Some(path.as_path()),
1664 WorkdirSnapshot::Invalid => None,
1665 WorkdirSnapshot::ReplayFailed => unreachable!("replay failures return above"),
1666 };
1667 let path = path_crate_tools::workdir::target_dir(&invoke.arguments, base)
1668 .map_err(|failure| Box::new((identity.clone(), failure)))?;
1669 cwd_state
1670 .start_pending_workdir_result(
1671 invoke.agent_id.clone(),
1672 path,
1673 identity.clone(),
1674 None,
1675 )
1676 .map_err(|_| {
1677 Box::new((
1678 identity.clone(),
1679 path_crate_display::ToolFailure::new(
1680 "another workdir change is already pending for this agent and shell instance",
1681 ),
1682 ))
1683 })?;
1684 cwd_state.mark_pending_workdir_awaiting_echo(&invoke.agent_id, &identity.call_id);
1685 let path = cwd_state
1686 .pending_workdir_target(&invoke.agent_id, &identity.call_id)
1687 .expect("newly reserved workdir target");
1688 let mutation_id =
1689 cwd_state.pending_workdir_mutation_id(&invoke.agent_id, &identity.call_id);
1690 if tx
1691 .send_checked(HarnessInputMessage::emit_transient(
1692 Event::AgentMetadataSetRequest(tau_proto::AgentMetadataSet {
1693 agent_id: invoke.agent_id,
1694 key: cwd_state.key(),
1695 value: CborValue::Text(path.display().to_string()),
1696 mutation_id,
1697 inheritable: true,
1698 }),
1699 ))
1700 .is_err()
1701 {
1702 let failure =
1703 path_crate_display::ToolFailure::new("failed to request workdir metadata commit");
1704 if send_identity_failure(identity.clone(), failure, &tx).is_ok() {
1705 cwd_state.take_pending_workdir_by_call(&identity.call_id);
1706 }
1707 return Ok(());
1708 }
1709 return Ok(());
1710 }
1711 let priority = priority_for_tool(&invoke, &config);
1712 let meta = WorkMeta {
1713 call_id: Some(invoke.call_id.clone()),
1714 agent_id: Some(invoke.agent_id.clone()),
1715 queued_bytes: approximate_tool_bytes(&invoke, scheduler.queued_bytes_limit()),
1716 };
1717 #[cfg(test)]
1718 tool_started_identity::ownership_probe::record_queued_bytes(
1719 &identity.call_id,
1720 meta.queued_bytes,
1721 );
1722 let tx_for_job = tx.clone();
1723 let lifecycle = cancellation.lifecycles.admit(
1724 invoke.call_id.clone(),
1725 invoke.tool_name.clone(),
1726 invoke.agent_id.clone(),
1727 tx_for_job.clone(),
1728 );
1729 let lifecycle_for_error = lifecycle.clone();
1730 let identity_for_error = identity;
1731 let cwd_state_for_error = cwd_state.clone();
1732 scheduler
1733 .enqueue(priority, meta, move || {
1734 #[cfg(test)]
1735 lifecycle.test_pause_after_dequeue();
1736 if invoke.tool_name == DIR_LOCK_TOOL_NAME {
1737 if lifecycle.start_effect() {
1738 crate::dir_lock::dispatch_dir_lock_tool(
1739 invoke,
1740 &lock_manager,
1741 config.dir_lock.enable,
1742 &tx_for_job,
1743 lifecycle.clone(),
1744 );
1745 }
1746 } else if config.dir_lock.enable && is_dir_lock_update_tool(invoke.tool_name.as_str()) {
1747 dispatch_locked_tool_invoke(
1748 invoke,
1749 ToolDispatchContext {
1750 shell_config: config.shell,
1751 tx: tx_for_job.clone(),
1752 running_calls: Arc::clone(&cancellation.running_calls),
1753 enforce_ro_bind: config.dir_lock.enforce_ro_bind,
1754 cwd_state: cwd_state.clone(),
1755 lifecycle: lifecycle.clone(),
1756 },
1757 &lock_manager,
1758 match &workdir_snapshot {
1759 WorkdirSnapshot::Valid(cwd) => cwd.clone(),
1760 WorkdirSnapshot::Invalid => {
1761 unreachable!("only workdir admits invalid state")
1762 }
1763 WorkdirSnapshot::ReplayFailed => {
1764 unreachable!("replay failures return above")
1765 }
1766 },
1767 );
1768 } else if invoke.tool_name == EXPORT_TOOL_NAME || invoke.tool_name == IMPORT_TOOL_NAME {
1769 if lifecycle.start_effect()
1770 && artifact_control.prepare(
1771 invoke,
1772 lifecycle.clone(),
1773 match &workdir_snapshot {
1774 WorkdirSnapshot::Valid(cwd) => cwd,
1775 WorkdirSnapshot::Invalid | WorkdirSnapshot::ReplayFailed => {
1776 unreachable!("artifact tools require a valid workdir")
1777 }
1778 },
1779 &tx_for_job,
1780 )
1781 {
1782 return;
1783 }
1784 } else {
1785 if lifecycle.start_effect() {
1786 dispatch_tool_invoke(
1787 invoke,
1788 ToolDispatchContext {
1789 shell_config: config.shell,
1790 tx: tx_for_job.clone(),
1791 running_calls: Arc::clone(&cancellation.running_calls),
1792 enforce_ro_bind: config.dir_lock.enforce_ro_bind,
1793 cwd_state: cwd_state.clone(),
1794 lifecycle: lifecycle.clone(),
1795 },
1796 None,
1797 config
1798 .dir_lock
1799 .enable
1800 .then_some(ShellCommandMode::visible(ShellAccessMode::ReadOnly)),
1801 workdir_snapshot.clone(),
1802 );
1803 }
1804 }
1805 if !tx_for_job.mandatory_output_failed() {
1806 lifecycle.finish();
1807 }
1808 })
1809 .map_err(|error| {
1810 lifecycle_for_error.finish();
1811 cwd_state_for_error.take_pending_workdir_by_call(&identity_for_error.call_id);
1812 Box::new((
1813 identity_for_error,
1814 path_crate_display::ToolFailure::new(error.message),
1815 ))
1816 })
1817}
1818
1819struct UiShellScheduleContext<'a> {
1821 scheduler: &'a WorkScheduler,
1823 tx: &'a Output,
1825 shell_config: ShellConfig,
1827 running_ui_commands: Arc<Mutex<HashMap<tau_proto::ShellCommandId, mpsc::Sender<()>>>>,
1829 shutdown_generation_counter: Arc<UiShellShutdownGenerationCounter>,
1831 scheduled_generation: UiShellShutdownGeneration,
1833 cwd: PathBuf,
1835}
1836
1837fn schedule_ui_shell_command(
1838 cmd: tau_proto::UiShellCommand,
1839 context: UiShellScheduleContext<'_>,
1840) -> Result<(), Box<(tau_proto::UiShellCommand, String)>> {
1841 let UiShellScheduleContext {
1842 scheduler,
1843 tx,
1844 shell_config,
1845 running_ui_commands,
1846 shutdown_generation_counter,
1847 scheduled_generation,
1848 cwd,
1849 } = context;
1850 let meta = WorkMeta {
1851 call_id: None,
1852 agent_id: cmd.target_agent_id.clone(),
1853 queued_bytes: cmd.command.len(),
1854 };
1855 let tx_for_job = tx.clone();
1856 let cmd_for_error = cmd.clone();
1857 let command_id = cmd.command_id.clone();
1858 scheduler
1859 .enqueue(WorkPriority::User, meta, move || {
1860 let (cancel_tx, cancel_rx) = mpsc::channel();
1861 running_ui_commands
1862 .lock()
1863 .expect("running ui shell registry lock poisoned")
1864 .insert(command_id.clone(), cancel_tx.clone());
1865 if shutdown_generation_counter.current() != scheduled_generation {
1866 let _ = cancel_tx.send(());
1867 }
1868 path_crate_tools::shell::dispatch_user_shell_command(
1869 cmd,
1870 shell_config,
1871 &tx_for_job,
1872 cancel_rx,
1873 cwd,
1874 );
1875 running_ui_commands
1876 .lock()
1877 .expect("running ui shell registry lock poisoned")
1878 .remove(&command_id);
1879 })
1880 .map_err(|error| Box::new((cmd_for_error, error.message)))
1881}
1882
1883fn priority_for_tool(invoke: &tau_proto::ToolStarted, config: &ExtConfig) -> WorkPriority {
1884 if invoke.tool_name == DIR_LOCK_TOOL_NAME {
1885 if is_dir_lock_update_invocation(&invoke.arguments) {
1886 return WorkPriority::Bulk;
1887 }
1888 return WorkPriority::Control;
1889 }
1890 if matches!(
1891 invoke.tool_name.as_str(),
1892 READ_TOOL_NAME
1893 | EXPORT_TOOL_NAME
1894 | IMPORT_TOOL_NAME
1895 | GREP_TOOL_NAME
1896 | FIND_TOOL_NAME
1897 | LS_TOOL_NAME
1898 ) {
1899 return WorkPriority::Cheap;
1900 }
1901 if config.dir_lock.enable && is_dir_lock_update_tool(invoke.tool_name.as_str()) {
1902 return WorkPriority::Bulk;
1903 }
1904 WorkPriority::Bulk
1905}
1906
1907fn approximate_tool_bytes(invoke: &tau_proto::ToolStarted, queued_bytes_limit: usize) -> usize {
1908 let cap = queued_bytes_limit.saturating_add(1);
1909 let base = invoke
1910 .call_id
1911 .as_str()
1912 .len()
1913 .saturating_add(invoke.tool_name.as_str().len())
1914 .saturating_add(invoke.agent_id.as_str().len());
1915 saturating_add_capped(base, estimate_cbor_bytes(&invoke.arguments, cap), cap)
1916}
1917
1918fn estimate_cbor_bytes(value: &CborValue, cap: usize) -> usize {
1919 if cap == 0 {
1920 return 0;
1921 }
1922 match value {
1923 CborValue::Integer(_) | CborValue::Float(_) | CborValue::Bool(_) | CborValue::Null => {
1924 8.min(cap)
1925 }
1926 CborValue::Bytes(bytes) => bytes.len().min(cap),
1927 CborValue::Text(text) => text.len().min(cap),
1928 CborValue::Tag(_, inner) => saturating_add_capped(8, estimate_cbor_bytes(inner, cap), cap),
1929 CborValue::Array(values) => estimate_cbor_sequence(values.iter(), cap),
1930 CborValue::Map(entries) => {
1931 let mut total = 1usize;
1932 for (key, value) in entries {
1933 total = saturating_add_capped(total, estimate_cbor_bytes(key, cap - total), cap);
1934 if cap <= total {
1935 return cap;
1936 }
1937 total = saturating_add_capped(total, estimate_cbor_bytes(value, cap - total), cap);
1938 if cap <= total {
1939 return cap;
1940 }
1941 }
1942 total
1943 }
1944 _ => 8.min(cap),
1945 }
1946}
1947
1948fn estimate_cbor_sequence<'a>(values: impl Iterator<Item = &'a CborValue>, cap: usize) -> usize {
1949 let mut total = 1usize;
1950 for value in values {
1951 total = saturating_add_capped(total, estimate_cbor_bytes(value, cap - total), cap);
1952 if cap <= total {
1953 return cap;
1954 }
1955 }
1956 total
1957}
1958
1959fn saturating_add_capped(lhs: usize, rhs: usize, cap: usize) -> usize {
1960 lhs.saturating_add(rhs).min(cap)
1961}
1962
1963struct ToolDispatchContext {
1965 shell_config: ShellConfig,
1967 tx: Output,
1969 running_calls: Arc<Mutex<HashMap<tau_proto::ToolCallId, mpsc::Sender<()>>>>,
1971 enforce_ro_bind: bool,
1973 cwd_state: CwdState,
1975 lifecycle: ToolLifecycle,
1977}
1978
1979fn dispatch_locked_tool_invoke(
1980 invoke: tau_proto::ToolStarted,
1981 context: ToolDispatchContext,
1982 lock_manager: &DirLockManager,
1983 cwd: PathBuf,
1984) {
1985 let ToolDispatchContext {
1986 shell_config,
1987 tx,
1988 running_calls,
1989 enforce_ro_bind,
1990 cwd_state,
1991 lifecycle,
1992 } = context;
1993 let dirs = match crate::dir_lock::automatic_lock_dirs_for_tool_in_dir(
1994 invoke.tool_name.as_str(),
1995 &invoke.arguments,
1996 &cwd,
1997 ) {
1998 Ok(dirs) => crate::dir_lock::normalize_lock_dirs(dirs),
1999 Err(error) => {
2000 if lifecycle.claim_terminal_before_effect() {
2001 let _ = send_tool_failure(invoke, error, &tx);
2002 }
2003 return;
2004 }
2005 };
2006 let shell_command_mode = is_shell_command_tool(invoke.tool_name.as_str())
2007 .then_some(ShellCommandMode::visible(ShellAccessMode::ReadWrite));
2008
2009 let lock_wait_started = Instant::now();
2010 let wait_progress = crate::dir_lock::waiting_progress(&invoke, &dirs, shell_command_mode);
2011 let wait_tx = tx.clone();
2012 let on_wait = move || {
2013 let _ = wait_tx.report_tool_progress(wait_progress);
2014 };
2015 let guard = match if shell_command_mode.is_some() {
2016 lock_manager.acquire_auto_if_manual_covers(
2017 invoke.call_id.clone(),
2018 invoke.agent_id.clone(),
2019 dirs,
2020 on_wait,
2021 )
2022 } else {
2023 lock_manager.acquire_auto(
2024 invoke.call_id.clone(),
2025 invoke.agent_id.clone(),
2026 dirs,
2027 on_wait,
2028 )
2029 } {
2030 Ok(guard) => guard,
2031 Err(path_crate_dir_lock::LockAcquireError::NotCovered) => {
2032 if lifecycle.start_effect() {
2033 dispatch_tool_invoke(
2034 invoke,
2035 ToolDispatchContext {
2036 shell_config,
2037 tx,
2038 running_calls,
2039 enforce_ro_bind,
2040 cwd_state,
2041 lifecycle: lifecycle.clone(),
2042 },
2043 None,
2044 Some(ShellCommandMode::visible(ShellAccessMode::ReadOnly)),
2045 WorkdirSnapshot::Valid(cwd),
2046 );
2047 }
2048 return;
2049 }
2050 Err(path_crate_dir_lock::LockAcquireError::Cancelled) => {
2051 lifecycle.report_cancelled_before_effect();
2052 return;
2053 }
2054 Err(path_crate_dir_lock::LockAcquireError::Abandoned(lock)) => {
2055 if lifecycle.claim_terminal_before_effect() {
2056 let _ = send_tool_failure(invoke, lock.tool_failure(), &tx);
2057 }
2058 return;
2059 }
2060 Err(path_crate_dir_lock::LockAcquireError::SelfConflict {
2061 uncovered_dir,
2062 held_dir,
2063 }) => {
2064 if lifecycle.claim_terminal_before_effect() {
2065 let _ = send_tool_failure(
2066 invoke,
2067 path_crate_display::ToolFailure::new(format!(
2068 "automatic directory lock is outside your manual lock coverage: requested {}; held {}",
2069 uncovered_dir.display(),
2070 held_dir.display()
2071 )),
2072 &tx,
2073 );
2074 }
2075 return;
2076 }
2077 Err(path_crate_dir_lock::LockAcquireError::Backend(message)) => {
2078 if lifecycle.claim_terminal_before_effect() {
2079 let _ = send_tool_failure(
2080 invoke,
2081 path_crate_display::ToolFailure::new(format!(
2082 "dir_lock backend error: {message}"
2083 )),
2084 &tx,
2085 );
2086 }
2087 return;
2088 }
2089 };
2090
2091 let lock_wait_duration_seconds =
2092 reported_lock_wait_duration_seconds(lock_wait_started.elapsed());
2093 #[cfg(test)]
2094 lifecycle.test_pause_after_lock();
2095 if lifecycle.start_effect() {
2096 dispatch_tool_invoke(
2097 invoke,
2098 ToolDispatchContext {
2099 shell_config,
2100 tx,
2101 running_calls,
2102 enforce_ro_bind,
2103 cwd_state,
2104 lifecycle: lifecycle.clone(),
2105 },
2106 lock_wait_duration_seconds,
2107 shell_command_mode,
2108 WorkdirSnapshot::Valid(cwd),
2109 );
2110 }
2111 drop(guard);
2112}
2113
2114fn send_ui_shell_saturated_failure(cmd: tau_proto::UiShellCommand, message: String, tx: &Output) {
2115 let _ = tx.send_checked(HarnessInputMessage::emit(
2116 Event::ShellCommandFinishedReported(tau_proto::ShellCommandFinished {
2117 command_id: cmd.command_id,
2118 session_id: cmd.session_id,
2119 command: cmd.command,
2120 include_in_context: cmd.include_in_context,
2121 target_agent_id: cmd.target_agent_id,
2122 output: message,
2123 exit_code: None,
2124 cancelled: false,
2125 }),
2126 ));
2127}
2128
2129fn send_identity_failure(
2130 identity: tool_started_identity::ToolStartedIdentity,
2131 failure: crate::display::ToolFailure,
2132 tx: &Output,
2133) -> tau_client::ClientResult<()> {
2134 let crate::display::ToolFailure {
2135 message,
2136 details,
2137 display,
2138 } = failure;
2139 tx.report_tool_terminal(Event::ToolError(tau_proto::ToolError {
2140 presentation: Default::default(),
2141 call_id: identity.call_id,
2142 tool_name: identity.wire_tool_name,
2143 tool_type: tau_proto::ToolType::Function,
2144 message,
2145 details: details.map(|details| *details),
2146 display: Some(*display),
2147 originator: identity.originator,
2148 }))
2149}
2150
2151fn send_tool_failure(
2152 invoke: tau_proto::ToolStarted,
2153 failure: crate::display::ToolFailure,
2154 tx: &Output,
2155) -> tau_client::ClientResult<()> {
2156 send_identity_failure(invoke.into(), failure, tx)
2157}
2158
2159fn reported_lock_wait_duration_seconds(elapsed: Duration) -> Option<u64> {
2160 if elapsed <= Duration::from_secs(SLOW_LOCK_WAIT_THRESHOLD_SECS) {
2161 return None;
2162 }
2163
2164 let whole_seconds = elapsed.as_secs();
2165 if Duration::from_secs(whole_seconds) < elapsed {
2166 Some(whole_seconds.saturating_add(1))
2167 } else {
2168 Some(whole_seconds)
2169 }
2170}
2171
2172fn with_lock_wait_duration(event: Event, lock_wait_duration_seconds: Option<u64>) -> Event {
2173 let Some(seconds) = lock_wait_duration_seconds else {
2174 return event;
2175 };
2176
2177 match event {
2178 Event::ToolResult(mut result) => {
2179 result.result = cbor_value_with_lock_wait_duration(result.result, seconds, "output");
2180 Event::ToolResult(result)
2181 }
2182 Event::ToolError(mut error) => {
2183 error.details = Some(match error.details {
2184 Some(details) => cbor_value_with_lock_wait_duration(details, seconds, "details"),
2185 None => CborValue::Map(vec![lock_wait_duration_entry(seconds)]),
2186 });
2187 Event::ToolError(error)
2188 }
2189 event => event,
2190 }
2191}
2192
2193fn cbor_value_with_lock_wait_duration(
2194 value: CborValue,
2195 seconds: u64,
2196 non_map_payload_key: &str,
2197) -> CborValue {
2198 match value {
2199 CborValue::Map(mut entries) => {
2200 prepend_lock_wait_duration(&mut entries, seconds);
2201 CborValue::Map(entries)
2202 }
2203 value => CborValue::Map(vec![
2204 lock_wait_duration_entry(seconds),
2205 (CborValue::Text(non_map_payload_key.to_owned()), value),
2206 ]),
2207 }
2208}
2209
2210fn prepend_lock_wait_duration(entries: &mut Vec<(CborValue, CborValue)>, seconds: u64) {
2211 entries.retain(|(key, _)| match key {
2212 CborValue::Text(key) => key != LOCK_WAIT_DURATION_SECONDS_HEADER,
2213 _ => true,
2214 });
2215 entries.insert(0, lock_wait_duration_entry(seconds));
2216}
2217
2218fn lock_wait_duration_entry(seconds: u64) -> (CborValue, CborValue) {
2219 let seconds = i64::try_from(seconds).unwrap_or(i64::MAX);
2220 (
2221 CborValue::Text(LOCK_WAIT_DURATION_SECONDS_HEADER.to_owned()),
2222 CborValue::Integer(seconds.into()),
2223 )
2224}
2225
2226fn dispatch_tool_invoke(
2228 mut invoke: tau_proto::ToolStarted,
2229 context: ToolDispatchContext,
2230 lock_wait_duration_seconds: Option<u64>,
2231 shell_command_mode: Option<ShellCommandMode>,
2232 workdir_snapshot: WorkdirSnapshot,
2233) {
2234 let ToolDispatchContext {
2235 shell_config,
2236 tx,
2237 running_calls,
2238 enforce_ro_bind,
2239 cwd_state,
2240 lifecycle,
2241 } = context;
2242 if invoke.tool_name == WORKDIR_TOOL_NAME {
2243 if cbor_optional_text(&invoke.arguments, "path").is_none() {
2244 let output = path_crate_tools::workdir::status_output(match &workdir_snapshot {
2245 WorkdirSnapshot::Valid(path) => Some(path.as_path()),
2246 WorkdirSnapshot::Invalid => None,
2247 WorkdirSnapshot::ReplayFailed => unreachable!("replay failures return above"),
2248 });
2249 let _ = tx.report_tool_terminal(Event::ToolResult(ToolResult {
2250 presentation: Default::default(),
2251 call_id: invoke.call_id,
2252 tool_name: invoke.tool_name,
2253 tool_type: tau_proto::ToolType::Function,
2254 result: output.result,
2255 provider_content: output.provider_content,
2256 kind: ToolResultKind::Final,
2257 display: Some(output.display),
2258 originator: invoke.originator,
2259 }));
2260 return;
2261 }
2262 let agent_id = invoke.agent_id.clone();
2263 if let Some(path) = cwd_state.pending_workdir_target(&agent_id, &invoke.call_id) {
2264 if cwd_state.mark_pending_workdir_awaiting_echo(&agent_id, &invoke.call_id) {
2265 let metadata = Event::AgentMetadataSetRequest(tau_proto::AgentMetadataSet {
2266 agent_id,
2267 key: cwd_state.key(),
2268 value: CborValue::Text(path.display().to_string()),
2269 mutation_id: None,
2270 inheritable: true,
2271 });
2272 let _ = tx.send_checked(HarnessInputMessage::emit_transient(metadata));
2273 }
2274 return;
2275 }
2276 return;
2279 }
2280 let tool_cwd = match workdir_snapshot {
2281 WorkdirSnapshot::Valid(cwd) => cwd,
2282 WorkdirSnapshot::Invalid => unreachable!("non-workdir calls reject invalid state"),
2283 WorkdirSnapshot::ReplayFailed => unreachable!("replay failures return above"),
2284 };
2285 if matches!(
2286 invoke.tool_name.as_str(),
2287 READ_TOOL_NAME
2288 | EDIT_TOOL_NAME
2289 | GREP_TOOL_NAME
2290 | FIND_TOOL_NAME
2291 | LS_TOOL_NAME
2292 | SHELL_TOOL_NAME
2293 | GPT_SHELL_TOOL_NAME
2294 ) {
2295 crate::shell_output_spool::note_call();
2296 }
2297 let (world, authorized_cwd) = match world_after_shell_authorization(
2298 &mut invoke,
2299 &shell_config,
2300 tau_vcr::VcrConfig::from_env(),
2301 tool_cwd,
2302 ) {
2303 Ok(world) => world,
2304 Err(crate::display::ToolFailure {
2305 message,
2306 details,
2307 display,
2308 }) => {
2309 let event = Event::ToolError(tau_proto::ToolError {
2310 presentation: Default::default(),
2311 call_id: invoke.call_id.clone(),
2312 tool_name: invoke.tool_name.clone(),
2313 tool_type: tau_proto::ToolType::Function,
2314 message,
2315 details: details.map(|details| *details),
2316 display: Some(*display),
2317 originator: invoke.originator.clone(),
2318 });
2319 let event = with_lock_wait_duration(event, lock_wait_duration_seconds);
2320 let _ = tx.report_tool_terminal(event);
2321 return;
2322 }
2323 };
2324
2325 if invoke.tool_name == SHELL_TOOL_NAME || invoke.tool_name == GPT_SHELL_TOOL_NAME {
2326 dispatch_cancellable_shell_tool(CancellableShellDispatch {
2327 invoke,
2328 shell_config,
2329 tx: &tx,
2330 running_calls: &running_calls,
2331 lifecycle: &lifecycle,
2332 lock_wait_duration_seconds,
2333 shell_command_mode: shell_command_mode.unwrap_or(ShellCommandMode::READ_WRITE_HIDDEN),
2334 enforce_ro_bind,
2335 world,
2336 authorized_cwd,
2337 });
2338 return;
2339 }
2340
2341 if invoke.tool_name == GREP_TOOL_NAME || invoke.tool_name == FIND_TOOL_NAME {
2342 dispatch_cancellable_non_shell_tool(
2343 invoke,
2344 &tx,
2345 &running_calls,
2346 &lifecycle,
2347 lock_wait_duration_seconds,
2348 world,
2349 );
2350 return;
2351 }
2352
2353 if let Some(display) = crate::tools::initial_display(&invoke) {
2354 let _ = tx.report_tool_progress(tau_proto::ToolProgress {
2355 call_id: invoke.call_id.clone(),
2356 tool_name: invoke.tool_name.clone(),
2357 message: None,
2358 progress: None,
2359 display: Some(display),
2360 });
2361 }
2362
2363 let events = execute_tool(invoke, world);
2364 for event in events {
2365 let event = with_lock_wait_duration(event, lock_wait_duration_seconds);
2366 let _ = tx.report_tool_terminal(event);
2367 }
2368}
2369
2370fn world_after_shell_authorization(
2373 invoke: &mut tau_proto::ToolStarted,
2374 shell_config: &ShellConfig,
2375 vcr_config: Option<tau_vcr::VcrConfig>,
2376 tool_cwd: PathBuf,
2377) -> Result<(path_crate_tools_world::ShellWorld, Option<PathBuf>), crate::display::ToolFailure> {
2378 let authorized_cwd = if let Some(surface) =
2379 path_crate_tools::ShellSurface::for_tool_name(invoke.tool_name.as_str())
2380 {
2381 path_crate_tools::shell::prepare_tool_invocation(surface, &invoke.arguments, shell_config)?
2382 } else {
2383 None
2384 };
2385 if let Some(surface) = path_crate_tools::ShellSurface::for_tool_name(invoke.tool_name.as_str())
2386 && let Some(canonical_cwd) = authorized_cwd.as_ref()
2387 {
2388 set_cbor_text_field(
2389 &mut invoke.arguments,
2390 surface.directory_argument(),
2391 canonical_cwd.display().to_string(),
2392 );
2393 }
2394 let world = path_crate_tools_world::ShellWorld::for_tool_in_dir(
2395 invoke.tool_name.as_str(),
2396 invoke.call_id.as_str(),
2397 &invoke.arguments,
2398 vcr_config,
2399 tool_cwd,
2400 )?;
2401 Ok((world, authorized_cwd))
2402}
2403
2404fn dispatch_cancellable_non_shell_tool(
2405 invoke: tau_proto::ToolStarted,
2406 tx: &Output,
2407 running_calls: &Arc<Mutex<HashMap<tau_proto::ToolCallId, mpsc::Sender<()>>>>,
2408 lifecycle: &ToolLifecycle,
2409 lock_wait_duration_seconds: Option<u64>,
2410 world: path_crate_tools::world::ShellWorld,
2411) {
2412 #[cfg(test)]
2413 lifecycle.test_pause_before_active_registration();
2414 let (cancel_tx, cancel_rx) = mpsc::channel();
2415 running_calls
2416 .lock()
2417 .expect("running call registry lock poisoned")
2418 .insert(invoke.call_id.clone(), cancel_tx.clone());
2419 if lifecycle.effect_cancel_requested() {
2420 let _ = cancel_tx.send(());
2421 }
2422
2423 if let Some(display) = crate::tools::initial_display(&invoke) {
2424 let _ = tx.report_tool_progress(tau_proto::ToolProgress {
2425 call_id: invoke.call_id.clone(),
2426 tool_name: invoke.tool_name.clone(),
2427 message: None,
2428 progress: None,
2429 display: Some(display),
2430 });
2431 }
2432
2433 let call_id = invoke.call_id.clone();
2434 let tool_name = invoke.tool_name.clone();
2435 let outcome = crate::tools::execute_cancellable_tool(invoke, world, cancel_rx);
2436
2437 running_calls
2438 .lock()
2439 .expect("running call registry lock poisoned")
2440 .remove(&call_id);
2441
2442 match outcome {
2443 path_crate_tools::CancellableToolOutcome::Finished(events) => {
2444 for event in events {
2445 let event = with_lock_wait_duration(event, lock_wait_duration_seconds);
2446 let _ = tx.report_tool_terminal(event);
2447 }
2448 }
2449 path_crate_tools::CancellableToolOutcome::Cancelled => {
2450 let event = Event::ToolCancelled(ToolCancelled {
2451 presentation: Default::default(),
2452 call_id,
2453 tool_name,
2454 tool_type: tau_proto::ToolType::Function,
2455 display: None,
2456 });
2457 let event = with_lock_wait_duration(event, lock_wait_duration_seconds);
2458 let _ = tx.report_tool_terminal(event);
2459 }
2460 }
2461}
2462
2463struct CancellableShellDispatch<'a> {
2465 invoke: tau_proto::ToolStarted,
2467 shell_config: ShellConfig,
2469 tx: &'a Output,
2471 running_calls: &'a Arc<Mutex<HashMap<tau_proto::ToolCallId, mpsc::Sender<()>>>>,
2474 lifecycle: &'a ToolLifecycle,
2476 lock_wait_duration_seconds: Option<u64>,
2478 shell_command_mode: ShellCommandMode,
2480 enforce_ro_bind: bool,
2483 world: path_crate_tools::world::ShellWorld,
2485 authorized_cwd: Option<PathBuf>,
2487}
2488
2489fn dispatch_cancellable_shell_tool(params: CancellableShellDispatch<'_>) {
2490 let CancellableShellDispatch {
2491 invoke,
2492 shell_config,
2493 tx,
2494 running_calls,
2495 lifecycle,
2496 lock_wait_duration_seconds,
2497 shell_command_mode,
2498 enforce_ro_bind,
2499 mut world,
2500 authorized_cwd,
2501 } = params;
2502 #[cfg(test)]
2503 lifecycle.test_pause_before_active_registration();
2504 let (cancel_tx, cancel_rx) = mpsc::channel();
2505 debug!(
2506 call_id = %invoke.call_id,
2507 tool_name = %invoke.tool_name,
2508 "registering cancellable shell call"
2509 );
2510 running_calls
2511 .lock()
2512 .expect("running call registry lock poisoned")
2513 .insert(invoke.call_id.clone(), cancel_tx.clone());
2514 if lifecycle.effect_cancel_requested() {
2515 let _ = cancel_tx.send(());
2516 }
2517
2518 let _ = tx.report_tool_progress(tau_proto::ToolProgress {
2519 call_id: invoke.call_id.clone(),
2520 tool_name: invoke.tool_name.clone(),
2521 message: None,
2522 progress: None,
2523 display: Some(path_crate_tools::shell::initial_display(
2524 &invoke.arguments,
2525 shell_command_mode,
2526 )),
2527 });
2528 let result = path_crate_tools::shell::run_command_cancellable_for_tool(
2529 path_crate_tools::shell::ShellInvocation {
2530 surface: path_crate_tools::ShellSurface::for_tool_name(invoke.tool_name.as_str())
2531 .expect("shell dispatch accepts only known shell tools"),
2532 call_id: invoke.call_id.as_str(),
2533 arguments: &invoke.arguments,
2534 authorized_cwd: authorized_cwd.as_deref(),
2535 },
2536 &shell_config,
2537 shell_command_mode,
2538 enforce_ro_bind,
2539 Some(cancel_rx),
2540 &mut world,
2541 );
2542 let outcome = match (result, world.finish()) {
2543 (Ok(outcome), Ok(())) => Ok(outcome),
2544 (Ok(_), Err(failure)) | (Err(failure), Ok(())) | (Err(failure), Err(_)) => Err(failure),
2545 };
2546 let event = match outcome {
2547 Ok(path_crate_tools_shell::CommandOutcome::Finished(output)) => {
2548 debug!(call_id = %invoke.call_id, tool_name = %invoke.tool_name, "cancellable shell call finished");
2549 Event::ToolResult(ToolResult {
2550 presentation: Default::default(),
2551 call_id: invoke.call_id.clone(),
2552 tool_name: invoke.tool_name.clone(),
2553 tool_type: tau_proto::ToolType::Function,
2554 result: output.result,
2555 provider_content: Vec::new(),
2556 kind: ToolResultKind::Final,
2557 display: Some(output.display),
2558 originator: invoke.originator.clone(),
2559 })
2560 }
2561 Ok(path_crate_tools_shell::CommandOutcome::Cancelled) => {
2562 debug!(call_id = %invoke.call_id, tool_name = %invoke.tool_name, "cancellable shell call cancelled");
2563 Event::ToolCancelled(ToolCancelled {
2564 presentation: Default::default(),
2565 call_id: invoke.call_id.clone(),
2566 tool_name: invoke.tool_name.clone(),
2567 tool_type: tau_proto::ToolType::Function,
2568 display: None,
2569 })
2570 }
2571 Err(crate::display::ToolFailure {
2572 message,
2573 details,
2574 display,
2575 }) => {
2576 debug!(
2577 call_id = %invoke.call_id,
2578 tool_name = %invoke.tool_name,
2579 message,
2580 "cancellable shell call failed"
2581 );
2582 Event::ToolError(tau_proto::ToolError {
2583 presentation: Default::default(),
2584 call_id: invoke.call_id.clone(),
2585 tool_name: invoke.tool_name.clone(),
2586 tool_type: tau_proto::ToolType::Function,
2587 message,
2588 details: details.map(|details| *details),
2589 display: Some(*display),
2590 originator: invoke.originator.clone(),
2591 })
2592 }
2593 };
2594
2595 running_calls
2596 .lock()
2597 .expect("running call registry lock poisoned")
2598 .remove(&invoke.call_id);
2599 trace!(call_id = %invoke.call_id, "removed shell call from cancellation registry");
2600 let event = with_lock_wait_duration(event, lock_wait_duration_seconds);
2601 if tx.report_tool_terminal(event).is_err() {
2602 debug!(call_id = %invoke.call_id, "failed to send terminal shell event to harness");
2603 }
2604}
2605
2606fn dispatch_session_started(
2607 started: SessionStarted,
2608 tx: &Output,
2609 discovery_policy: DiscoverySourcePolicy,
2610) -> tau_client::ClientResult<()> {
2611 let session_id = started.session_id.clone();
2612 let scan = build_discovery_snapshot(started, discovery_policy);
2613 for diagnostic in scan.diagnostics {
2614 let _ = tx.send(diagnostic);
2615 }
2616 dispatch_session_discovery_messages(
2617 session_id,
2618 vec![HarnessInputMessage::emit_transient(
2619 Event::ExtensionSessionDiscoverySnapshotDeclared(scan.snapshot),
2620 )],
2621 tx,
2622 )
2623}
2624
2625fn dispatch_session_discovery_messages(
2628 session_id: tau_proto::SessionId,
2629 messages: Vec<HarnessInputMessage>,
2630 tx: &Output,
2631) -> tau_client::ClientResult<()> {
2632 for message in messages {
2633 tx.send_checked(message)?;
2634 }
2635 tx.send_checked(HarnessInputMessage::emit_transient(
2636 Event::ExtensionSessionContextReady(ExtensionSessionContextReady { session_id }),
2637 ))
2638}
2639
2640fn apply_started_cwd_metadata(
2641 started: tau_proto::AgentStarted,
2642 tx: &Output,
2643 cwd_state: &CwdState,
2644 is_replay: bool,
2645) -> tau_client::ClientResult<()> {
2646 for item in started.metadata {
2647 if item.key == cwd_state.key() {
2648 if let CborValue::Text(path) = item.value {
2649 let cwd = PathBuf::from(path);
2650 if cwd_state.set_metadata_text(started.agent_id.clone(), cwd.clone())
2651 && !is_replay
2652 && let Some((session_id, initialization_id)) =
2653 cwd_state.initialization(&started.agent_id)
2654 {
2655 tx.send_checked(HarnessInputMessage::emit_transient(cwd_context_event(
2656 session_id,
2657 started.agent_id.clone(),
2658 initialization_id,
2659 &cwd,
2660 cwd_state,
2661 )))?;
2662 }
2663 } else {
2664 cwd_state.set_invalid(started.agent_id.clone());
2665 }
2666 }
2667 }
2668 Ok(())
2669}
2670
2671fn dispatch_session_agent_loaded(
2672 loaded: SessionAgentLoaded,
2673 tx: &Output,
2674 cwd_state: &CwdState,
2675 defer_default_until_replay_complete: bool,
2676 discovery_policy: DiscoverySourcePolicy,
2677) -> tau_client::ClientResult<()> {
2678 if defer_default_until_replay_complete {
2679 cwd_state.set_pending_ready(
2680 loaded.agent_id,
2681 loaded.session_id,
2682 loaded.agent_initialization_id,
2683 );
2684 return Ok(());
2685 }
2686 publish_agent_discovery_snapshot(&loaded, tx, discovery_policy)?;
2687 if let Some(cwd) = cwd_state.get(&loaded.agent_id) {
2688 tx.send_checked(HarnessInputMessage::emit_transient(cwd_context_event(
2689 loaded.session_id.clone(),
2690 loaded.agent_id.clone(),
2691 loaded.agent_initialization_id.clone(),
2692 &cwd,
2693 cwd_state,
2694 )))?;
2695 tx.send_checked(HarnessInputMessage::emit_transient(
2696 Event::ExtensionContextReady(ExtensionContextReady {
2697 session_id: loaded.session_id,
2698 agent_id: loaded.agent_id,
2699 agent_initialization_id: loaded.agent_initialization_id,
2700 }),
2701 ))?;
2702 return Ok(());
2703 }
2704
2705 cwd_state.set_pending_ready(
2706 loaded.agent_id.clone(),
2707 loaded.session_id,
2708 loaded.agent_initialization_id,
2709 );
2710 let Ok(cwd) = cwd_state.process_default() else {
2711 return Ok(());
2712 };
2713 tx.send_checked(HarnessInputMessage::emit_transient(
2714 Event::AgentMetadataSetRequest(tau_proto::AgentMetadataSet {
2715 agent_id: loaded.agent_id,
2716 key: cwd_state.key(),
2717 value: CborValue::Text(cwd.display().to_string()),
2718 mutation_id: None,
2719 inheritable: true,
2720 }),
2721 ))
2722}
2723
2724fn cwd_context_event(
2725 session_id: tau_proto::SessionId,
2726 agent_id: tau_proto::AgentId,
2727 agent_initialization_id: tau_proto::AgentInitializationId,
2728 cwd: &Path,
2729 cwd_state: &CwdState,
2730) -> Event {
2731 let status = if cwd.is_dir() {
2732 "available"
2733 } else {
2734 "unavailable"
2735 };
2736 Event::ExtAgentContextPublish(ExtAgentContextPublish {
2737 session_id,
2738 agent_id,
2739 agent_initialization_id,
2740 key: AgentContextKey::new("workdir"),
2741 value: AgentContextValue(serde_json::json!({
2742 "label": cwd_state.context_label(),
2743 "path": cwd.display().to_string(),
2744 "status": status,
2745 })),
2746 })
2747}
2748
2749fn invalid_cwd_context_event(
2750 session_id: tau_proto::SessionId,
2751 agent_id: tau_proto::AgentId,
2752 agent_initialization_id: tau_proto::AgentInitializationId,
2753 cwd_state: &CwdState,
2754) -> Event {
2755 Event::ExtAgentContextPublish(ExtAgentContextPublish {
2756 session_id,
2757 agent_id,
2758 agent_initialization_id,
2759 key: AgentContextKey::new("workdir"),
2760 value: AgentContextValue(serde_json::json!({
2761 "label": cwd_state.context_label(),
2762 "path": "<invalid>",
2763 "status": "invalid",
2764 })),
2765 })
2766}
2767
2768fn cwd_notice_event(agent_id: tau_proto::AgentId, cwd: &Path) -> Event {
2769 Event::AgentUserMessageInjected(tau_proto::AgentUserMessageInjected {
2770 inference_activation: false,
2771 agent_id,
2772 text: format!("Your working directory changed to {}.", cwd.display()),
2773 message_class: tau_proto::PromptMessageClass::Internal,
2774 })
2775}
2776
2777fn is_shell_tool(name: &str) -> bool {
2778 matches!(
2779 name,
2780 READ_TOOL_NAME
2781 | READ_IMAGE_TOOL_NAME
2782 | EXPORT_TOOL_NAME
2783 | IMPORT_TOOL_NAME
2784 | EDIT_TOOL_NAME
2785 | REPLACE_TOOL_NAME
2786 | APPLY_PATCH_TOOL_NAME
2787 | GREP_TOOL_NAME
2788 | FIND_TOOL_NAME
2789 | LS_TOOL_NAME
2790 | WORKDIR_TOOL_NAME
2791 | SHELL_TOOL_NAME
2792 | GPT_SHELL_TOOL_NAME
2793 | DIR_LOCK_TOOL_NAME
2794 ) || is_echo_tool(name)
2795}
2796
2797fn is_dir_lock_update_invocation(arguments: &CborValue) -> bool {
2798 crate::argument::optional_argument_text(arguments, "command")
2799 .ok()
2800 .flatten()
2801 .as_deref()
2802 == Some("update")
2803}
2804
2805fn is_dir_lock_update_tool(name: &str) -> bool {
2806 matches!(
2807 name,
2808 EDIT_TOOL_NAME
2809 | REPLACE_TOOL_NAME
2810 | APPLY_PATCH_TOOL_NAME
2811 | SHELL_TOOL_NAME
2812 | GPT_SHELL_TOOL_NAME
2813 )
2814}
2815
2816fn is_shell_command_tool(name: &str) -> bool {
2817 matches!(name, SHELL_TOOL_NAME | GPT_SHELL_TOOL_NAME)
2818}
2819
2820#[cfg(any(test, feature = "echo-agent"))]
2821fn is_echo_tool(name: &str) -> bool {
2822 name == ECHO_TOOL_NAME
2823}
2824
2825#[cfg(not(any(test, feature = "echo-agent")))]
2826fn is_echo_tool(_name: &str) -> bool {
2827 false
2828}
2829
2830fn build_discovery_snapshot(
2831 _started: SessionStarted,
2832 discovery_policy: DiscoverySourcePolicy,
2833) -> DiscoveryScan {
2834 let mut diagnostics = Vec::new();
2835 let (skills, agents_files) = if discovery_policy.reads_environment() {
2836 let skill_dirs = session_skill_dirs(std::env::current_dir().ok(), dirs::home_dir());
2837 let result = tau_skills::load_skills_from_skill_dirs(&skill_dirs);
2838 push_skill_diagnostic_requests(&mut diagnostics, result.diagnostics);
2839 let skills = result
2840 .skills
2841 .into_iter()
2842 .map(discovery_skill_candidate)
2843 .collect();
2844 let agents_files = discover_session_agents_files()
2845 .into_iter()
2846 .map(|file| DiscoveryAgentsFile {
2847 file_path: file.file_path,
2848 content: file.content,
2849 })
2850 .collect();
2851 (skills, agents_files)
2852 } else {
2853 (Vec::new(), Vec::new())
2854 };
2855 DiscoveryScan {
2856 snapshot: ExtensionSessionDiscoverySnapshotDeclared {
2857 session_id: _started.session_id,
2858 skills,
2859 agents_files,
2860 },
2861 diagnostics,
2862 }
2863}
2864
2865fn discovery_skill_candidate(skill: tau_skills::Skill) -> DiscoverySkillCandidate {
2866 let file_path = skill.file_path.canonicalize().unwrap_or(skill.file_path);
2867 let sampled_modified = std::fs::metadata(&file_path)
2868 .and_then(|metadata| metadata.modified())
2869 .ok()
2870 .and_then(system_time_to_discovery_micros);
2871 DiscoverySkillCandidate {
2872 name: skill.name,
2873 description: skill.description,
2874 file_path,
2875 add_to_prompt: skill.add_to_prompt,
2876 user_invocable: skill.user_invocable,
2877 disable_model_invocation: skill.disable_model_invocation,
2878 argument_hint: skill.argument_hint,
2879 sampled_modified,
2880 }
2881}
2882
2883fn system_time_to_discovery_micros(time: std::time::SystemTime) -> Option<DiscoveryModifiedMicros> {
2884 match time.duration_since(std::time::UNIX_EPOCH) {
2885 Ok(duration) => i64::try_from(duration.as_micros())
2886 .ok()
2887 .map(DiscoveryModifiedMicros::new),
2888 Err(error) => i64::try_from(error.duration().as_micros())
2889 .ok()
2890 .and_then(i64::checked_neg)
2891 .map(DiscoveryModifiedMicros::new),
2892 }
2893}
2894
2895fn publish_agent_discovery_snapshot(
2896 loaded: &SessionAgentLoaded,
2897 tx: &Output,
2898 discovery_policy: DiscoverySourcePolicy,
2899) -> tau_client::ClientResult<()> {
2900 publish_agent_discovery_snapshot_for(
2901 loaded.session_id.clone(),
2902 loaded.agent_id.clone(),
2903 loaded.agent_initialization_id.clone(),
2904 tx,
2905 discovery_policy,
2906 )
2907}
2908
2909fn publish_agent_discovery_snapshot_for(
2910 session_id: tau_proto::SessionId,
2911 agent_id: tau_proto::AgentId,
2912 agent_initialization_id: tau_proto::AgentInitializationId,
2913 tx: &Output,
2914 discovery_policy: DiscoverySourcePolicy,
2915) -> tau_client::ClientResult<()> {
2916 let session = SessionStarted {
2917 session_id: session_id.clone(),
2918 reason: tau_proto::SessionStartReason::Resume,
2919 };
2920 publish_agent_discovery_scan(
2921 build_discovery_snapshot(session, discovery_policy),
2922 agent_id,
2923 agent_initialization_id,
2924 tx,
2925 )
2926}
2927
2928fn publish_agent_discovery_scan(
2930 scan: DiscoveryScan,
2931 agent_id: tau_proto::AgentId,
2932 agent_initialization_id: tau_proto::AgentInitializationId,
2933 tx: &Output,
2934) -> tau_client::ClientResult<()> {
2935 tx.send_checked(agent_discovery_message(
2936 scan.snapshot,
2937 agent_id,
2938 agent_initialization_id,
2939 ))
2940}
2941
2942fn agent_discovery_message(
2944 snapshot: ExtensionSessionDiscoverySnapshotDeclared,
2945 agent_id: tau_proto::AgentId,
2946 agent_initialization_id: tau_proto::AgentInitializationId,
2947) -> HarnessInputMessage {
2948 HarnessInputMessage::emit_transient(Event::ExtensionAgentDiscoverySnapshotDeclared(
2949 ExtensionAgentDiscoverySnapshotDeclared {
2950 session_id: snapshot.session_id,
2951 agent_id,
2952 agent_initialization_id,
2953 skills: snapshot.skills,
2954 agents_files: snapshot.agents_files,
2955 },
2956 ))
2957}
2958
2959fn shell_workdir_prompt_fragment(shell: &config::ShellConfig) -> PromptFragment {
2960 let mut template = String::from(
2961 "{{#if agent_context.workdir}}### Shell workdirs\n\nEach shell extension instance \
2962 has its own persistent workdir; there is no global shell cwd.\n\
2963 {{#each agent_context.workdir}}- {{#if (eq value.label \"default\")}}default shell \
2964 tools (`workdir`){{else}}`{{value.label}}_*` shell tools \
2965 (`{{value.label}}_workdir`){{/if}}: `{{value.path}}` \
2966 [{{value.status}}]\n{{/each}}\nNormally set the matching workdir tool to the project \
2967 root before project work. It sets the cwd/base for later shell and filesystem calls \
2968 in that same instance. The cwd can select configured directory-scoped wrappers, \
2969 notably `direnv exec .`, and affect other cwd-sensitive wrappers/tools. After \
2970 changing it, make dependent calls only in a later tool turn after success; sibling \
2971 calls have no workdir-first ordering.{{/if}}",
2972 );
2973 if let Some(allowlist) = shell.allowlist_prompt_fragment() {
2974 template.push_str(&allowlist);
2975 }
2976 PromptFragment::new(
2977 "shell.workdir",
2978 PromptPriority::new(900),
2979 PromptContent::new(template),
2980 )
2981}
2982
2983fn push_skill_diagnostic_requests(
2984 messages: &mut Vec<HarnessInputMessage>,
2985 diagnostics: Vec<tau_skills::SkillDiagnostic>,
2986) {
2987 for diagnostic in diagnostics {
2988 let (kind, level) = match diagnostic.kind {
2989 tau_skills::DiagnosticKind::Warning => ("warning", tau_proto::NoticeLevel::Info),
2990 tau_skills::DiagnosticKind::Collision => ("collision", tau_proto::NoticeLevel::Trace),
2991 tau_skills::DiagnosticKind::Skipped => ("skipped", tau_proto::NoticeLevel::Warning),
2992 };
2993 messages.push(HarnessInputMessage::ExtensionNoticeRequest(
2994 tau_proto::ExtensionNoticeRequest {
2995 message: format!(
2996 "skill {kind}: {}\n{}",
2997 diagnostic.path.display(),
2998 diagnostic.message
2999 ),
3000 level,
3001 },
3002 ));
3003 }
3004}
3005
3006fn session_skill_dirs(
3007 cwd: Option<std::path::PathBuf>,
3008 home: Option<std::path::PathBuf>,
3009) -> Vec<tau_skills::SkillDir> {
3010 let mut skill_dirs = Vec::new();
3011 if let Some(cwd) = cwd.as_deref() {
3012 for project_dir in project_skill_ancestor_dirs(cwd, home.as_deref()) {
3013 push_existing_project_skill_dir(
3014 &mut skill_dirs,
3015 project_dir.join(".agents").join("skills"),
3016 );
3017 push_existing_project_skill_dir(
3018 &mut skill_dirs,
3019 project_dir.join(".agents.local").join("skills"),
3020 );
3021 }
3022 }
3023 if let Some(home) = home {
3024 skill_dirs.push(user_skill_dir_precedence(
3025 home.join(".config").join("agents").join("skills"),
3026 XDG_USER_SKILL_SOURCE_PRECEDENCE,
3027 ));
3028 skill_dirs.push(user_skill_dir_precedence(
3029 home.join(".config").join("agents.local").join("skills"),
3030 XDG_USER_SKILL_SOURCE_PRECEDENCE,
3031 ));
3032 skill_dirs.push(user_skill_dir_precedence(
3033 home.join(".agents").join("skills"),
3034 LEGACY_USER_SKILL_SOURCE_PRECEDENCE,
3035 ));
3036 skill_dirs.push(user_skill_dir_precedence(
3037 home.join(".agents.local").join("skills"),
3038 LEGACY_USER_SKILL_SOURCE_PRECEDENCE,
3039 ));
3040 }
3041 skill_dirs
3042}
3043
3044fn project_skill_ancestor_dirs(
3045 cwd: &std::path::Path,
3046 home: Option<&std::path::Path>,
3047) -> Vec<std::path::PathBuf> {
3048 ancestor_dirs(cwd)
3049 .into_iter()
3050 .filter(|dir| dir.parent().is_some())
3051 .filter(|dir| {
3052 let Some(home) = home else {
3053 return true;
3054 };
3055 !cwd.starts_with(home) || (dir.starts_with(home) && dir != home)
3056 })
3057 .collect()
3058}
3059
3060fn push_existing_project_skill_dir(
3061 skill_dirs: &mut Vec<tau_skills::SkillDir>,
3062 path: std::path::PathBuf,
3063) {
3064 if path.is_dir() {
3065 skill_dirs.push(project_skill_dir(path));
3066 }
3067}
3068
3069fn project_skill_dir(path: std::path::PathBuf) -> tau_skills::SkillDir {
3070 tau_skills::SkillDir {
3071 path,
3072 add_to_prompt_by_default: true,
3073 source_precedence: None,
3074 }
3075}
3076
3077fn user_skill_dir_precedence(
3078 path: std::path::PathBuf,
3079 source_precedence: u32,
3080) -> tau_skills::SkillDir {
3081 tau_skills::SkillDir {
3082 path,
3083 add_to_prompt_by_default: false,
3084 source_precedence: Some(source_precedence),
3085 }
3086}