dpc-tau-ext-shell 0.2.1

A minimal Unix-first coding agent.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
//! Bounded lifecycle state for admitted model tool calls.

#[cfg(test)]
mod tests;

use std::collections::HashMap;
use std::sync::{Arc, Mutex, Weak, mpsc};

use tau_proto::{AgentId, Event, ToolCallId, ToolCancelled, ToolName, ToolType};

use crate::Output;

/// Shared cancellation state for pre-effect and actively cancellable model
/// calls.
#[derive(Clone, Default)]
pub(crate) struct ToolCancellationState {
    /// Bounded lifecycle authority for every admitted scheduled model call.
    pub(crate) lifecycles: ToolLifecycleRegistry,
    /// Cancellation senders for shell and search effects currently executing.
    pub(crate) running_calls: Arc<Mutex<HashMap<ToolCallId, mpsc::Sender<()>>>>,
}

/// Registry that keeps cancellation authoritative across scheduler and lock
/// handoffs.
#[derive(Clone, Default)]
pub(crate) struct ToolLifecycleRegistry {
    /// Live calls indexed by their harness call identifier.
    inner: Arc<Mutex<HashMap<ToolCallId, Arc<Entry>>>>,
    #[cfg(test)]
    /// Deterministic handoff barriers installed by focused race tests.
    hooks: Arc<Mutex<TestHooks>>,
}

/// One admitted call's shared lifecycle authority.
#[derive(Clone)]
pub(crate) struct ToolLifecycle {
    /// State shared with cancellation processing.
    entry: Arc<Entry>,
}

/// Result of processing a cancellation against a live lifecycle.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum CancelOutcome {
    /// Cancellation won before the call crossed the effect-start boundary.
    PreventedEffect,
    /// The effect had started, so existing active cancellation must be
    /// signalled.
    EffectStarted,
}

/// Shared state and report metadata for one admitted model call.
struct Entry {
    /// Call identifier used for registry cleanup and cancellation reports.
    call_id: ToolCallId,
    /// Local tool name; the scoped output maps it back to the wire name.
    tool_name: ToolName,
    /// Agent that owns the admitted call.
    agent_id: AgentId,
    /// Output scope captured when the call was admitted.
    tx: Output,
    /// Registry owning this entry while the call remains live.
    registry: Weak<Mutex<HashMap<ToolCallId, Arc<Entry>>>>,
    /// Atomic lifecycle winner protected by one short critical section.
    state: Mutex<State>,
    #[cfg(test)]
    /// Registry-scoped deterministic handoff barriers.
    hooks: Arc<Mutex<TestHooks>>,
}

#[derive(Clone, Copy)]
enum State {
    /// The call has not started an externally visible effect.
    BeforeEffect,
    /// The effect started; the flag remembers cancellation during sender
    /// handoff.
    EffectStarted { cancel_requested: bool },
    /// A terminal path won; successful publication or shutdown removes the
    /// registry entry.
    Terminal,
}

#[cfg(test)]
#[derive(Default)]
/// Registry-scoped deterministic barriers used only by lifecycle race tests.
struct TestHooks {
    /// Barrier after scheduler dequeue and before ordinary dispatch.
    after_dequeue: Option<TestHandoff>,
    /// Barrier after lock acquisition and before effect start.
    after_lock: Option<TestHandoff>,
    /// Barrier after effect start and before active sender registration.
    before_active_registration: Option<TestHandoff>,
    /// Barrier after effect start and before direct lock-waiter registration.
    before_lock_waiter_registration: Option<TestHandoff>,
}

#[cfg(test)]
/// One deterministic worker-to-test handoff barrier.
struct TestHandoff {
    /// Notification that the worker reached the handoff.
    reached: mpsc::SyncSender<()>,
    /// Permission for the worker to leave the handoff.
    resume: mpsc::Receiver<()>,
}

impl ToolLifecycleRegistry {
    /// Admit one scheduled model call into the bounded live-call registry.
    pub(crate) fn admit(
        &self,
        call_id: ToolCallId,
        tool_name: ToolName,
        agent_id: AgentId,
        tx: Output,
    ) -> ToolLifecycle {
        let entry = Arc::new(Entry {
            call_id: call_id.clone(),
            tool_name,
            agent_id,
            tx,
            registry: Arc::downgrade(&self.inner),
            state: Mutex::new(State::BeforeEffect),
            #[cfg(test)]
            hooks: Arc::clone(&self.hooks),
        });
        self.inner
            .lock()
            .expect("tool lifecycle registry lock poisoned")
            .insert(call_id, Arc::clone(&entry));
        ToolLifecycle { entry }
    }

    /// Install a deterministic scheduler-dequeue handoff barrier for one test.
    #[cfg(test)]
    pub(crate) fn pause_after_dequeue(
        &self,
        reached: mpsc::SyncSender<()>,
        resume: mpsc::Receiver<()>,
    ) {
        self.hooks
            .lock()
            .expect("tool lifecycle test hooks poisoned")
            .after_dequeue = Some(TestHandoff { reached, resume });
    }

    /// Install a deterministic lock-acquired handoff barrier for one test.
    #[cfg(test)]
    pub(crate) fn pause_after_lock(
        &self,
        reached: mpsc::SyncSender<()>,
        resume: mpsc::Receiver<()>,
    ) {
        self.hooks
            .lock()
            .expect("tool lifecycle test hooks poisoned")
            .after_lock = Some(TestHandoff { reached, resume });
    }

    /// Install a deterministic active-sender registration barrier for one test.
    #[cfg(test)]
    pub(crate) fn pause_before_active_registration(
        &self,
        reached: mpsc::SyncSender<()>,
        resume: mpsc::Receiver<()>,
    ) {
        self.hooks
            .lock()
            .expect("tool lifecycle test hooks poisoned")
            .before_active_registration = Some(TestHandoff { reached, resume });
    }

    /// Install a deterministic direct lock-waiter registration barrier.
    #[cfg(test)]
    pub(crate) fn pause_before_lock_waiter_registration(
        &self,
        reached: mpsc::SyncSender<()>,
        resume: mpsc::Receiver<()>,
    ) {
        self.hooks
            .lock()
            .expect("tool lifecycle test hooks poisoned")
            .before_lock_waiter_registration = Some(TestHandoff { reached, resume });
    }

    /// Process cancellation against the call's single lifecycle authority.
    pub(crate) fn cancel(&self, call_id: &ToolCallId) -> Option<CancelOutcome> {
        let entry = self
            .inner
            .lock()
            .expect("tool lifecycle registry lock poisoned")
            .get(call_id)
            .cloned()?;
        let outcome = {
            let mut state = entry.state.lock().expect("tool lifecycle state poisoned");
            match *state {
                State::BeforeEffect => {
                    *state = State::Terminal;
                    CancelOutcome::PreventedEffect
                }
                State::EffectStarted { .. } => {
                    *state = State::EffectStarted {
                        cancel_requested: true,
                    };
                    CancelOutcome::EffectStarted
                }
                State::Terminal => return None,
            }
        };
        if outcome == CancelOutcome::PreventedEffect
            && entry
                .tx
                .report_tool_terminal(Event::ToolCancelled(ToolCancelled {
                    presentation: Default::default(),
                    call_id: entry.call_id.clone(),
                    tool_name: entry.tool_name.clone(),
                    tool_type: ToolType::Function,
                    display: None,
                }))
                .is_ok()
        {
            entry.remove_from_registry();
        }
        Some(outcome)
    }

    /// Prevent queued or pre-effect work owned by an agent that is leaving.
    ///
    /// Active work retains its entry until its ordinary terminal path because
    /// agent unload did not previously cancel already-running model tools.
    pub(crate) fn remove_agent(&self, agent_id: &AgentId) {
        let entries = {
            self.inner
                .lock()
                .expect("tool lifecycle registry lock poisoned")
                .values()
                .filter(|entry| &entry.agent_id == agent_id)
                .cloned()
                .collect::<Vec<_>>()
        };
        let mut removed = Vec::new();
        for entry in entries {
            let mut state = entry.state.lock().expect("tool lifecycle state poisoned");
            if matches!(*state, State::BeforeEffect) {
                *state = State::Terminal;
                drop(state);
                removed.push(entry);
            }
        }
        let mut registry = self
            .inner
            .lock()
            .expect("tool lifecycle registry lock poisoned");
        registry.retain(|_, entry| !removed.iter().any(|removed| Arc::ptr_eq(removed, entry)));
    }
    /// Stop pre-effect calls and preserve cancellation across active sender
    /// handoff.
    pub(crate) fn prepare_shutdown(&self) {
        let entries = {
            self.inner
                .lock()
                .expect("tool lifecycle registry lock poisoned")
                .values()
                .cloned()
                .collect::<Vec<_>>()
        };
        for entry in entries {
            let remove = {
                let mut state = entry.state.lock().expect("tool lifecycle state poisoned");
                match *state {
                    State::BeforeEffect => {
                        *state = State::Terminal;
                        true
                    }
                    State::EffectStarted { .. } => {
                        *state = State::EffectStarted {
                            cancel_requested: true,
                        };
                        false
                    }
                    State::Terminal => true,
                }
            };
            if remove {
                entry.remove_from_registry();
            }
        }
    }
}

impl ToolLifecycle {
    /// Report cancellation if this call is still before effect start.
    ///
    /// A prior explicit cancellation or shutdown has already made the state
    /// terminal, so this method suppresses duplicate reports on those paths.
    pub(crate) fn report_cancelled_before_effect(&self) {
        let should_report = {
            let mut state = self
                .entry
                .state
                .lock()
                .expect("tool lifecycle state poisoned");
            if matches!(*state, State::BeforeEffect) {
                *state = State::Terminal;
                true
            } else {
                false
            }
        };
        if should_report
            && self
                .entry
                .tx
                .report_tool_terminal(Event::ToolCancelled(ToolCancelled {
                    presentation: Default::default(),
                    call_id: self.entry.call_id.clone(),
                    tool_name: self.entry.tool_name.clone(),
                    tool_type: ToolType::Function,
                    display: None,
                }))
                .is_ok()
        {
            self.entry.remove_from_registry();
        }
    }

    /// Claim a terminal error before effect start, unless cancellation won
    /// first.
    pub(crate) fn claim_terminal_before_effect(&self) -> bool {
        {
            let mut state = self
                .entry
                .state
                .lock()
                .expect("tool lifecycle state poisoned");
            match *state {
                State::BeforeEffect => {
                    *state = State::Terminal;
                    true
                }
                State::EffectStarted { .. } => {
                    // ast-grep-ignore: debug-assert-expression-must-not-mutate
                    debug_assert!(false, "pre-effect terminal claimed after effect start");
                    false
                }
                State::Terminal => false,
            }
        }
    }

    /// Pause at the scheduler-dequeue handoff when a focused test installed a
    /// barrier.
    #[cfg(test)]
    pub(crate) fn test_pause_after_dequeue(&self) {
        self.entry
            .pause_test_handoff(TestHandoffPoint::AfterDequeue);
    }

    /// Pause at the lock-acquired handoff when a focused test installed a
    /// barrier.
    #[cfg(test)]
    pub(crate) fn test_pause_after_lock(&self) {
        self.entry.pause_test_handoff(TestHandoffPoint::AfterLock);
    }

    /// Pause before active sender registration when a focused test installed a
    /// barrier.
    #[cfg(test)]
    pub(crate) fn test_pause_before_active_registration(&self) {
        self.entry
            .pause_test_handoff(TestHandoffPoint::BeforeActiveRegistration);
    }

    /// Pause before direct lock-waiter registration when a focused test
    /// installed a barrier.
    #[cfg(test)]
    pub(crate) fn test_pause_before_lock_waiter_registration(&self) {
        self.entry
            .pause_test_handoff(TestHandoffPoint::BeforeLockWaiterRegistration);
    }

    /// Atomically cross the effect-start boundary if cancellation has not won.
    pub(crate) fn start_effect(&self) -> bool {
        let mut state = self
            .entry
            .state
            .lock()
            .expect("tool lifecycle state poisoned");
        match *state {
            State::BeforeEffect => {
                *state = State::EffectStarted {
                    cancel_requested: false,
                };
                true
            }
            State::EffectStarted { .. } => true,
            State::Terminal => false,
        }
    }

    /// Return whether cancellation arrived after effect start.
    pub(crate) fn effect_cancel_requested(&self) -> bool {
        matches!(
            *self
                .entry
                .state
                .lock()
                .expect("tool lifecycle state poisoned"),
            State::EffectStarted {
                cancel_requested: true
            }
        )
    }

    /// Mark the call terminal and remove its bounded registry entry.
    pub(crate) fn finish(&self) {
        *self
            .entry
            .state
            .lock()
            .expect("tool lifecycle state poisoned") = State::Terminal;
        self.entry.remove_from_registry();
    }
}

impl Entry {
    /// Take and execute one deterministic test handoff without holding hook
    /// state.
    #[cfg(test)]
    fn pause_test_handoff(&self, point: TestHandoffPoint) {
        let handoff = {
            let mut hooks = self
                .hooks
                .lock()
                .expect("tool lifecycle test hooks poisoned");
            match point {
                TestHandoffPoint::AfterDequeue => hooks.after_dequeue.take(),
                TestHandoffPoint::AfterLock => hooks.after_lock.take(),
                TestHandoffPoint::BeforeActiveRegistration => {
                    hooks.before_active_registration.take()
                }
                TestHandoffPoint::BeforeLockWaiterRegistration => {
                    hooks.before_lock_waiter_registration.take()
                }
            }
        };
        if let Some(handoff) = handoff {
            handoff.reached.send(()).expect("handoff observer");
            handoff.resume.recv().expect("handoff resume");
        }
    }

    /// Remove this exact admission without deleting a reused call identifier.
    fn remove_from_registry(self: &Arc<Self>) {
        let Some(registry) = self.registry.upgrade() else {
            return;
        };
        let mut registry = registry
            .lock()
            .expect("tool lifecycle registry lock poisoned");
        if registry
            .get(&self.call_id)
            .is_some_and(|entry| Arc::ptr_eq(entry, self))
        {
            registry.remove(&self.call_id);
        }
    }
}

#[cfg(test)]
enum TestHandoffPoint {
    /// Scheduler removed the call from its queue.
    AfterDequeue,
    /// Automatic directory-lock acquisition returned a held guard.
    AfterLock,
    /// Effect start won but the active cancellation sender is not registered.
    BeforeActiveRegistration,
    /// Effect start won but a direct directory-lock waiter is not registered.
    BeforeLockWaiterRegistration,
}