Skip to main content

tau_ext_shell/
lib.rs

1//! Filesystem and shell tool extension.
2//!
3//! Provides `read`, `edit`, `replace`, `apply_patch`, `dir_lock`, `grep`,
4//! `find`, `ls`, `workdir`, `shell`, and `gpt_shell` tools.
5//!
6//! The `echo` tool is available under `cfg(test)` or the
7//! `echo-agent` cargo feature for harness-side echo-agent tests.
8
9use 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/// Cloneable shell output adapter.
80///
81/// Production optional worker, progress, and diagnostic output uses
82/// tau-client's detached enqueue path so shell workers do not block on protocol
83/// flush. Sole terminals, user-shell completion, discovery, context,
84/// prerequisite metadata, and readiness use checked writes instead. The shared
85/// sticky failure state wakes the manual loop and prevents worker cleanup from
86/// releasing ownership after a failed terminal. Tests can use an mpsc-backed
87/// adapter for direct state-machine coverage.
88#[derive(Clone)]
89pub(crate) struct Output {
90    /// Production client or test channel receiving output frames.
91    inner: OutputInner,
92    /// Optional local-to-wire tool-name mapping for one scoped invocation.
93    tool_name_scope: Option<(tau_proto::ToolName, tau_proto::ToolName)>,
94    /// First mandatory-output failure shared with the manual policy loop.
95    failure: Arc<Mutex<MandatoryOutputFailure>>,
96}
97
98/// Cross-thread notification that a checked protocol write failed.
99#[derive(Default)]
100struct MandatoryOutputFailure {
101    /// First failure retained until the policy loop observes it.
102    message: Option<String>,
103    /// Sticky marker preserving worker ownership after loop observation.
104    failed: bool,
105    /// Manual-loop wake handle installed after tau-client startup.
106    waker: Option<tau_client::ManualRuntimeWaker>,
107}
108
109/// Backend used by [`Output`] to publish protocol frames.
110#[derive(Clone)]
111enum OutputInner {
112    /// Production tau-client writer handle.
113    Client(tau_client::ClientHandle),
114    #[cfg(test)]
115    /// Direct unit-test channel.
116    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    /// Sends mandatory lifecycle traffic in order and waits for writer flush.
175    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    /// Submit one transient tool progress observation.
206    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    /// Submit one terminal tool outcome through the typed client report helper.
217    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    /// Installs the wake handle used when a worker observes checked-output
237    /// failure.
238    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    /// Returns the first worker-side mandatory-output failure.
246    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    /// Returns whether any checked output has failed before loop teardown.
260    fn mandatory_output_failed(&self) -> bool {
261        self.failure
262            .lock()
263            .expect("mandatory output failure lock poisoned")
264            .failed
265    }
266
267    /// Records a checked-output failure and wakes the policy loop.
268    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
348/// One discovery pass split into mandatory state and optional session-only
349/// diagnostics.
350struct DiscoveryScan {
351    /// Complete discovery state published for the session or one agent.
352    snapshot: ExtensionSessionDiscoverySnapshotDeclared,
353    /// Best-effort diagnostic notices published only during session discovery.
354    diagnostics: Vec<HarnessInputMessage>,
355}
356
357/// Runs the extension on stdin/stdout.
358pub 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
368/// Runs the extension over arbitrary reader/writer streams.
369///
370/// The test-only `echo` tool is registered when built with
371/// `cfg(test)` or the `echo-agent` cargo feature.
372pub 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/// Runs an in-process harness test extension without discovering caller-owned
386/// skills or `AGENTS.md` files.
387///
388/// Harness unit tests share one process, so they cannot safely rewrite
389/// process-wide HOME, XDG, or working-directory state. This entrypoint keeps
390/// their synthetic extension from reading those ambient discovery inputs while
391/// preserving the ordinary protocol and tool behavior.
392#[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    // EOF/disconnect/errors may arrive without a committed SessionShutdown
1163    // event. Wake lock waiters and drop scheduler workers before tau-client
1164    // shuts down its writer, because active worker jobs may be blocked
1165    // inside DirLockManager and worker-held output handles must not enqueue
1166    // after writer shutdown.
1167    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        // Replay policy is declared per handler below: historical cwd metadata
1215        // is folded, while effectful tools/actions/cancellation/UI commands and
1216        // session lifecycle publications stay live-only.
1217        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 (&current.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        // Preserve malformed present values for the surface parser to reject.
1516        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
1819/// Frozen resources needed to enqueue one UI shell command.
1820struct UiShellScheduleContext<'a> {
1821    /// Scheduler that owns the queued command.
1822    scheduler: &'a WorkScheduler,
1823    /// Extension output used to publish command events.
1824    tx: &'a Output,
1825    /// Shell execution policy captured at admission.
1826    shell_config: ShellConfig,
1827    /// Cancellation senders for commands currently executing.
1828    running_ui_commands: Arc<Mutex<HashMap<tau_proto::ShellCommandId, mpsc::Sender<()>>>>,
1829    /// Shutdown authority shared with the runtime teardown handler.
1830    shutdown_generation_counter: Arc<UiShellShutdownGenerationCounter>,
1831    /// Lifecycle generation captured when the command was admitted.
1832    scheduled_generation: UiShellShutdownGeneration,
1833    /// Canonical workdir captured when the command was admitted.
1834    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
1963/// Frozen resources shared by one model tool dispatch.
1964struct ToolDispatchContext {
1965    /// Shell execution policy captured at admission.
1966    shell_config: ShellConfig,
1967    /// Extension output used to publish tool events.
1968    tx: Output,
1969    /// Cancellation senders for tool calls currently executing.
1970    running_calls: Arc<Mutex<HashMap<tau_proto::ToolCallId, mpsc::Sender<()>>>>,
1971    /// Whether read-only shell workdirs must be bind-mounted.
1972    enforce_ro_bind: bool,
1973    /// Per-instance workdir state used by the persistent workdir tool.
1974    cwd_state: CwdState,
1975    /// Shared effect-start and cancellation authority for this admitted call.
1976    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
2226/// Execute a single tool invocation and send the response event(s).
2227fn 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        // Every setter is reserved and validated at admission. A missing
2277        // reservation means cancellation or lifecycle cleanup won the race.
2278        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
2370/// Authorize shell invocations before opening VCR state, then construct the
2371/// execution world shared by shell and non-shell tools.
2372fn 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
2463/// Parameters needed to run a cancellable shell-like tool invocation.
2464struct CancellableShellDispatch<'a> {
2465    /// Tool invocation emitted by the harness.
2466    invoke: tau_proto::ToolStarted,
2467    /// Effective shell execution configuration for this invocation.
2468    shell_config: ShellConfig,
2469    /// Channel used to send progress and terminal events back to the harness.
2470    tx: &'a Output,
2471    /// Shared registry used by cancel requests to signal running shell
2472    /// processes.
2473    running_calls: &'a Arc<Mutex<HashMap<tau_proto::ToolCallId, mpsc::Sender<()>>>>,
2474    /// Lifecycle authority that bridges effect start to sender registration.
2475    lifecycle: &'a ToolLifecycle,
2476    /// Seconds spent waiting on a directory lock before this invocation ran.
2477    lock_wait_duration_seconds: Option<u64>,
2478    /// Display and access mode chosen for the shell command.
2479    shell_command_mode: ShellCommandMode,
2480    /// Whether read-only commands should run under the native read-only bind
2481    /// guard.
2482    enforce_ro_bind: bool,
2483    /// Tool execution world carrying the cwd and recorded side effects.
2484    world: path_crate_tools::world::ShellWorld,
2485    /// Canonical operational cwd retained from pre-VCR authorization.
2486    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
2625/// Publish one ordered session-discovery batch followed by its readiness
2626/// acknowledgement.
2627fn 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
2928/// Publishes only the mandatory snapshot from one per-agent discovery scan.
2929fn 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
2942/// Converts session discovery data into one correlated per-agent declaration.
2943fn 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}