zeph-subagent 0.22.2

Subagent management: spawning, grants, transcripts, and lifecycle hooks for Zeph
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
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
// SPDX-FileCopyrightText: 2026 Andrei G <bug-ops>
// SPDX-License-Identifier: MIT OR Apache-2.0

//! Sub-agent lifecycle management: spawn, cancel, collect, and resume.

mod collect;
mod secrets;
mod spawn;
mod worktree;

use std::collections::{HashMap, HashSet};
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Instant;

use tokio::sync::{mpsc, watch};
use tokio::task::JoinSet;
use tokio_util::sync::CancellationToken;
use zeph_common::task_supervisor::BlockingHandle;
use zeph_common::{SkillTrustLevel, TaskSupervisor};
use zeph_config::{ContentIsolationConfig, McpServerConfig};
use zeph_llm::provider::Message;

use crate::def::{PermissionMode, SubAgentDef};
use crate::durable::DurableResolverSeat;
use crate::error::SubAgentError;
use crate::fleet::SharedFleetRegistry;
use crate::forward::ForwardSurfaces;
use crate::grants::{GrantedSecret, PermissionGrants, SecretRequest};
use crate::state::SubAgentState;

/// Parent-derived state propagated to a spawned sub-agent at spawn time.
///
/// All fields default to empty/`None`, preserving existing behavior when callers
/// pass `SpawnContext::default()`.
///
/// # Constraint propagation
///
/// [`max_trust_level`][Self::max_trust_level] and
/// [`inherited_tool_allowlist`][Self::inherited_tool_allowlist] implement transitive
/// constraint propagation: safety constraints set at orchestration time are enforced on
/// every sub-agent in the spawn chain, regardless of nesting depth.
///
/// When a sub-agent spawns its own sub-agents it must forward these fields downward so
/// that grandchild agents cannot silently receive more privileges than the original
/// orchestration policy allowed.
///
/// # Examples
///
/// ```rust
/// use zeph_subagent::manager::SpawnContext;
///
/// // Minimal context — all fields use their defaults.
/// let ctx = SpawnContext::default();
/// assert!(ctx.parent_messages.is_empty());
/// assert_eq!(ctx.spawn_depth, 0);
/// assert!(ctx.max_trust_level.is_none());
/// assert!(ctx.inherited_tool_allowlist.is_none());
/// ```
#[derive(Default)]
pub struct SpawnContext {
    /// Recent parent conversation messages (last N turns).
    pub parent_messages: Vec<Message>,
    /// Parent's cancellation token for linked cancellation (foreground spawns).
    pub parent_cancel: Option<CancellationToken>,
    /// Parent's active provider name (for context propagation).
    pub parent_provider_name: Option<String>,
    /// Current spawn depth (0 = top-level agent).
    pub spawn_depth: u32,
    /// MCP tool names available in the parent's tool executor (for diagnostics).
    pub mcp_tool_names: Vec<String>,
    /// Seeded trajectory risk score from the parent sentinel (spec 050 §4).
    ///
    /// When `Some`, the subagent's `TrajectorySentinel` starts with this pre-seeded score
    /// rather than `0.0`, preventing a subagent spawn from acting as a free risk reset.
    /// The subagent loop applies this via `TrajectorySentinel::seed_score` after build.
    pub seed_trajectory_score: Option<f32>,
    /// Parent's content isolation config, propagated so the subagent loop can run the
    /// same sanitizer settings on hook-replaced tool output.
    pub content_isolation: ContentIsolationConfig,
    /// Name of the orchestrator that spawned this subagent.
    ///
    /// When set, the subagent's system prompt includes an identity header naming the
    /// orchestrator, so the subagent can validate that instructions are consistent with
    /// the expected authority.
    pub orchestrator_name: Option<String>,
    /// Role or task label of the orchestrating agent (e.g., `"planner"`, `"tool-router"`).
    ///
    /// Injected alongside [`orchestrator_name`][Self::orchestrator_name] when both are set.
    /// Omitted from the identity header when only `orchestrator_name` is provided.
    pub orchestrator_role: Option<String>,
    /// Per-session MCP servers to inject into this subagent's tool name annotations.
    ///
    /// The parent is responsible for connecting these servers and including them in the
    /// `tool_executor` passed to [`SubAgentManager::spawn`]. This field only carries the
    /// server metadata so the subagent's system prompt lists the additional tool names.
    pub session_mcp_servers: Vec<McpServerConfig>,
    /// Maximum trust level cap inherited from the parent agent or orchestration policy.
    ///
    /// When `Some(cap)`, the spawned sub-agent's effective trust level is clamped to
    /// `min(own_trust, cap)` so that sub-agents can never receive higher privileges than
    /// the orchestration policy originally allowed.
    ///
    /// # Caller responsibility for nested spawns
    ///
    /// This field does **not** propagate automatically. When a sub-agent itself spawns a
    /// grandchild, it must copy this field from its own received `SpawnContext` into the
    /// grandchild's `SpawnContext`. Passing `None` (the default) at that point means the
    /// grandchild receives **no cap**, which is a privilege escalation if the parent was
    /// constrained. Only the top-level session (spawned by `build_spawn_context`) correctly
    /// leaves this `None` — that represents an unconstrained top-level entry point.
    ///
    /// `None` means no cap is imposed by the parent (the sub-agent's own definition
    /// determines its trust level).
    pub max_trust_level: Option<SkillTrustLevel>,
    /// Tool names that this sub-agent is allowed to invoke, inherited from the parent.
    ///
    /// When `Some(set)`, the effective tool allowlist for the spawned agent is the
    /// intersection of `set` and the agent's own definition policy. This prevents a
    /// sub-agent from accessing tools that the parent is itself not allowed to use.
    ///
    /// # Caller responsibility for nested spawns
    ///
    /// Like [`max_trust_level`][Self::max_trust_level], this field does **not** propagate
    /// automatically. When a constrained sub-agent spawns its own children, it must copy
    /// this field from its received `SpawnContext` into the child's `SpawnContext`.
    /// Passing `None` at that point would grant the grandchild unrestricted tool access,
    /// defeating the original orchestration policy.
    ///
    /// `None` means no additional allowlist restriction is imposed by the parent
    /// (the agent's definition policy applies without narrowing).
    pub inherited_tool_allowlist: Option<HashSet<String>>,

    /// Durable resolver seat for promise-based subagent spawn/await (spec-064 §P4, INV-9).
    ///
    /// When `Some`, the spawned background task resolves the parent's durable promise after the
    /// agent loop terminates. The seat carries the resolver token and MUST NOT be forwarded to
    /// the child's tool executor or LLM surface — only the background task wrapper consumes it.
    ///
    /// `None` when `durable.enabled && durable.subagent` is false (plain spawn/collect path).
    pub durable_resolver: Option<DurableResolverSeat>,

    /// Deny network egress for this sub-agent's `bash` tool calls.
    ///
    /// Set by the orchestration layer when the spawning `TaskNode` carries
    /// `network_scope: NetworkScope::Deny` (spec `069-threat-model` OQ-1). When `true`,
    /// `build_filtered_executor` wraps the tool executor with
    /// [`NetworkDenyToolExecutor`](crate::NetworkDenyToolExecutor), which blocks `bash`
    /// invocations of `curl`, `wget`, `nc`, `ncat`, and `netcat` for this spawn only —
    /// sibling tasks and the parent agent's own executor are unaffected.
    ///
    /// # Caller responsibility for nested spawns
    ///
    /// Like [`max_trust_level`][Self::max_trust_level], this field does **not** propagate
    /// automatically. A sub-agent that itself spawns a grandchild must copy this field
    /// from its own received `SpawnContext` into the grandchild's `SpawnContext`, or the
    /// grandchild spawns with network access regardless of the original task's scope.
    ///
    /// `false` (the default) imposes no restriction beyond the executor/global
    /// `allow_network` default.
    pub network_denied: bool,

    /// Shared progress heartbeat for idle-timeout detection (issue #6245).
    ///
    /// Set by the orchestration driver (`handle_scheduler_spawn_action` in `zeph-core`'s
    /// `scheduler_loop.rs`) alongside [`network_denied`][Self::network_denied] — same
    /// post-construction assignment pattern, not part of `build_spawn_context`'s base
    /// literal. The driver creates the `Arc`, clones it in here, and keeps the original for
    /// `zeph_orchestration::DagScheduler::record_spawn`'s `last_progress_at` parameter so
    /// both the running loop and the scheduler observe the same counter.
    ///
    /// `None` (the default) for spawns not tracked by a `DagScheduler` — e.g. the standalone
    /// `/agent run` command — which are never idle-tracked.
    pub progress_at: Option<Arc<std::sync::atomic::AtomicU64>>,

    /// Cross-crate debug-dump sink, threaded down so sub-agent LLM calls are captured
    /// through the same pipeline as the top-level agent loop's `--debug-dump` output (#6391).
    ///
    /// Set by `zeph-core`'s `build_spawn_context` from `DebugState::debug_dumper`. `None`
    /// when debug dumps are disabled — no sub-agent dump is written in that case, mirroring
    /// the top-level `debug_dumper: None` behavior.
    ///
    /// # Caller responsibility for nested spawns
    ///
    /// Like [`max_trust_level`][Self::max_trust_level], this does **not** propagate
    /// automatically — a sub-agent that spawns its own children must copy this field from
    /// its received `SpawnContext` into the child's, or grandchild LLM calls go undumped.
    pub debug_dump_sink: Option<Arc<dyn zeph_llm::debug_dump::DebugDumpSink>>,
}

/// Live status snapshot of a running sub-agent.
///
/// Values are updated by the background agent loop via a [`tokio::sync::watch`] channel.
/// Callers receive snapshots via [`SubAgentManager::statuses`].
#[derive(Debug, Clone)]
pub struct SubAgentStatus {
    /// Current lifecycle state of the agent task.
    pub state: SubAgentState,
    /// Last message content from the agent (trimmed for display).
    pub last_message: Option<String>,
    /// Number of LLM turns consumed so far.
    pub turns_used: u32,
    /// Monotonic timestamp recorded at spawn time.
    pub started_at: Instant,
}

/// Handle to a spawned sub-agent task, owned by [`SubAgentManager`].
///
/// Fields are public to allow test harnesses in downstream crates to construct handles
/// without going through the full spawn lifecycle. Production code must not mutate
/// grants or the cancellation state directly — use the [`SubAgentManager`] API instead.
///
/// The `Drop` implementation cancels the task and revokes all grants as a safety net.
pub struct SubAgentHandle {
    /// Short display ID (same as `task_id` for non-resumed sessions).
    pub id: String,
    /// The definition that was used to spawn this agent.
    pub def: SubAgentDef,
    /// UUID assigned at spawn time (currently identical to `id`; separated for future use).
    pub task_id: String,
    /// Cached state — may lag the background task by one watch broadcast.
    pub state: SubAgentState,
    /// Supervised handle for the background agent loop task.
    pub join_handle: Option<BlockingHandle<Result<String, SubAgentError>>>,
    /// Cancellation token; cancelled on [`SubAgentManager::cancel`] or drop.
    pub cancel: CancellationToken,
    /// Watch receiver for live status updates from the agent loop.
    pub status_rx: watch::Receiver<SubAgentStatus>,
    /// Zero-trust TTL-bounded grants for this agent session.
    pub grants: PermissionGrants,
    /// Receives secret requests from the sub-agent loop.
    pub pending_secret_rx: mpsc::Receiver<SecretRequest>,
    /// Delivers the approval outcome to the sub-agent loop: `None` = denied,
    /// `Some(value)` = approved, carrying the resolved vault secret value and its
    /// grant expiry so the loop can re-validate the TTL locally on every tool call.
    pub secret_tx: mpsc::Sender<Option<GrantedSecret>>,
    /// ISO 8601 UTC timestamp recorded when the agent was spawned or resumed.
    pub started_at_str: String,
    /// Resolved transcript directory at spawn time; `None` if transcripts were disabled.
    pub transcript_dir: Option<PathBuf>,
    /// MCP tool names available at spawn time, persisted for transcript meta on collect.
    pub mcp_tool_names: Vec<String>,
}

impl SubAgentHandle {
    /// Construct a minimal [`SubAgentHandle`] for use in unit tests.
    ///
    /// The returned handle has a no-op cancel token, closed channels, and no grants.
    /// It must not be spawned or collected — it is only valid for inspection logic
    /// that operates on the handle's metadata fields (id, def, state, etc.).
    #[cfg(test)]
    pub fn for_test(id: impl Into<String>, def: SubAgentDef) -> Self {
        let initial_status = SubAgentStatus {
            state: SubAgentState::Working,
            last_message: None,
            turns_used: 0,
            started_at: Instant::now(),
        };
        let (status_tx, status_rx) = watch::channel(initial_status);
        drop(status_tx);
        let (pending_secret_rx_tx, pending_secret_rx) = mpsc::channel(1);
        drop(pending_secret_rx_tx);
        let (secret_tx, _) = mpsc::channel(1);
        let id_str = id.into();
        Self {
            task_id: id_str.clone(),
            id: id_str,
            def,
            state: SubAgentState::Working,
            join_handle: None,
            cancel: CancellationToken::new(),
            status_rx,
            grants: PermissionGrants::default(),
            pending_secret_rx,
            secret_tx,
            started_at_str: String::new(),
            transcript_dir: None,
            mcp_tool_names: Vec::new(),
        }
    }
}

impl std::fmt::Debug for SubAgentHandle {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.debug_struct("SubAgentHandle")
            .field("id", &self.id)
            .field("task_id", &self.task_id)
            .field("state", &self.state)
            .field("def_name", &self.def.name)
            .finish_non_exhaustive()
    }
}

impl Drop for SubAgentHandle {
    fn drop(&mut self) {
        // Defense-in-depth: cancel the task and revoke grants on drop even if
        // cancel() or collect() was not called (e.g., on panic or early return).
        self.cancel.cancel();
        if !self.grants.is_empty_grants() {
            tracing::warn!(
                id = %self.id,
                "SubAgentHandle dropped without explicit cleanup — revoking grants"
            );
        }
        self.grants.revoke_all();
    }
}

/// Manages sub-agent lifecycle: definitions, spawning, cancellation, and result collection.
///
/// `SubAgentManager` is the central coordinator for all sub-agent tasks. It tracks active
/// [`SubAgentHandle`]s, enforces the global concurrency limit, and stores loaded
/// [`SubAgentDef`]s.
///
/// # Concurrency model
///
/// The concurrency limit counts agents whose [`SubAgentState`] is `Submitted` or `Working`.
/// Reserved slots (via [`reserve_slots`][Self::reserve_slots]) also count against this limit
/// to allow orchestration schedulers to guarantee capacity before spawning.
///
/// # Examples
///
/// ```rust
/// use zeph_subagent::SubAgentManager;
///
/// let manager = SubAgentManager::new(4);
/// assert_eq!(manager.definitions().len(), 0);
/// ```
pub struct SubAgentManager {
    definitions: Vec<SubAgentDef>,
    agents: HashMap<String, SubAgentHandle>,
    max_concurrent: usize,
    /// Number of slots soft-reserved by the orchestration scheduler.
    ///
    /// Reserved slots count against the concurrency limit so that the scheduler can
    /// guarantee capacity for tasks it is about to spawn, preventing a planning-phase
    /// sub-agent from exhausting the pool and causing a deadlock.
    reserved_slots: usize,
    /// Config-level `SubagentStop` hooks, cached so `cancel()` and `collect()` can fire them.
    stop_hooks: Vec<super::hooks::HookDef>,
    /// Directory for JSONL transcripts and meta sidecars.
    transcript_dir: Option<PathBuf>,
    /// Maximum number of transcript files to keep (0 = unlimited).
    transcript_max_files: usize,
    /// Optional fleet registry for registering sub-agents in the fleet dashboard.
    ///
    /// When `None`, fleet registration is skipped silently. Inject via
    /// [`set_fleet_registry`][Self::set_fleet_registry].
    fleet_registry: Option<SharedFleetRegistry>,
    /// Tracks fire-and-forget hook and fleet-registry tasks to prevent silent panic swallowing.
    ///
    /// Completed and panicked tasks are drained before each new spawn. On graceful shutdown,
    /// [`shutdown_all`][Self::shutdown_all] aborts all outstanding tasks via
    /// [`JoinSet::shutdown`].
    hook_tasks: JoinSet<()>,
    /// Maximum number of concurrent hook tasks allowed in [`hook_tasks`][Self::hook_tasks].
    ///
    /// When the limit is reached, new fire-and-forget tasks are dropped with a warning instead
    /// of growing the set unboundedly under high-throughput spawning.
    max_hook_tasks: usize,
    /// Optional worktree manager; `Some` iff `worktree.enabled = true` in config.
    ///
    /// When set, every [`spawn`][Self::spawn] acquires [`cwd_lock`][Self::cwd_lock] for
    /// its full run so that plain agents cannot observe a stale cwd mutated by a worktree
    /// agent (INV-1). Only agents with `permissions.worktree = true` and a non-`None`
    /// `bg_isolation` actually get a dedicated worktree.
    ///
    /// This is the single live instance shared by the running agent's `/worktree`
    /// slash command (see [`worktree_manager`][Self::worktree_manager]) — distinct from
    /// the CLI's `zeph worktree list`/`clean`, which constructs its own fresh manager
    /// per invocation (`src/commands/worktree.rs`).
    worktree_manager: Option<Arc<zeph_worktree::DefaultWorktreeManager>>,
    /// Process-level serialisation mutex for working-directory mutations (INV-1).
    ///
    /// Acquired by every spawned task when `worktree_manager.is_some()`.  The
    /// `OwnedMutexGuard` is held for the full duration of `run_agent_loop` via the
    /// `CwdRestoreGuard` RAII wrapper.
    cwd_lock: Arc<tokio::sync::Mutex<()>>,
    /// Optional supervisor for subagent lifecycle tasks.
    ///
    /// When set, each spawned agent loop task is registered under its task ID so it is
    /// visible to TUI status panels and shutdown is coordinated through the supervisor.
    task_supervisor: Option<TaskSupervisor>,
    /// Which forwarding consumer surfaces are active for this session (issue #6359).
    ///
    /// Fixed at session start via [`set_forward_surfaces`][Self::set_forward_surfaces].
    /// `ForwardSurfaces::default()` (all `false`) until then, matching "forwarding
    /// disabled" byte-for-byte (FR-001/NFR-003).
    forward_surfaces: ForwardSurfaces,
    /// Ring buffer of recent sanitized display lines per task ID, owned by each task's
    /// forwarding drain (issue #6359). Read by [`forwarded_tail`][Self::forwarded_tail]
    /// for the TUI runtime detail view; entries are evicted by the drain itself shortly
    /// after that task's terminal chunk.
    forward_buffer: Arc<crate::forward::ForwardBuffer>,
    /// Optional secret-mask registry applied to forwarded chunks (issue #6359, security
    /// Finding 1 / NFR-005) — closes the consistency gap with the analogous outbound-LLM
    /// egress path (`apply_secret_masking`), which the forwarding pipeline would otherwise
    /// bypass. `None` (the default) means forwarded content is not masked for known vault
    /// secrets — set via [`set_secret_registry`][Self::set_secret_registry].
    secret_registry: Option<Arc<zeph_sanitizer::secret_mask::SecretMaskRegistry>>,
    /// Optional PII filter applied to forwarded chunks (issue #6359, security Finding 1 /
    /// NFR-005) — mirrors the optional `PiiFilter` layer sub-agent debug dumps get via
    /// `PiiScrubbingDumpSink` (#6407). `None` (the default) means no PII scrubbing beyond the
    /// baseline `ContentSanitizer` pass — set via [`set_pii_filter`][Self::set_pii_filter].
    pii_filter: Option<zeph_sanitizer::pii::PiiFilter>,
}

impl std::fmt::Debug for SubAgentManager {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.debug_struct("SubAgentManager")
            .field("definitions_count", &self.definitions.len())
            .field("active_agents", &self.agents.len())
            .field("max_concurrent", &self.max_concurrent)
            .field("reserved_slots", &self.reserved_slots)
            .field("stop_hooks_count", &self.stop_hooks.len())
            .field("transcript_dir", &self.transcript_dir)
            .field("transcript_max_files", &self.transcript_max_files)
            .field("fleet_registry", &self.fleet_registry.is_some())
            .field("hook_tasks_len", &self.hook_tasks.len())
            .field("max_hook_tasks", &self.max_hook_tasks)
            .field("worktree_manager", &self.worktree_manager.is_some())
            .field("cwd_lock", &"<Mutex>")
            .field("task_supervisor", &self.task_supervisor.is_some())
            .field("forward_surfaces", &self.forward_surfaces)
            .field("forward_buffer", &"<Mutex>")
            .field("secret_registry", &self.secret_registry.is_some())
            .field("pii_filter", &self.pii_filter.is_some())
            .finish()
    }
}

impl SubAgentManager {
    /// Create a new manager with the given concurrency limit.
    #[must_use]
    pub fn new(max_concurrent: usize) -> Self {
        Self {
            definitions: Vec::new(),
            agents: HashMap::new(),
            max_concurrent,
            reserved_slots: 0,
            stop_hooks: Vec::new(),
            transcript_dir: None,
            transcript_max_files: 50,
            fleet_registry: None,
            hook_tasks: JoinSet::new(),
            max_hook_tasks: 64,
            worktree_manager: None,
            cwd_lock: Arc::new(tokio::sync::Mutex::new(())),
            task_supervisor: None,
            forward_surfaces: ForwardSurfaces::default(),
            forward_buffer: crate::forward::new_buffer(),
            secret_registry: None,
            pii_filter: None,
        }
    }

    /// Inject a [`TaskSupervisor`] so subagent lifecycle tasks are registered and visible.
    ///
    /// Must be called before the first [`spawn`][Self::spawn]. When set, each spawned agent
    /// loop task is registered under its task ID and is observable in TUI status panels and
    /// [`TaskSupervisor::snapshot`].
    pub fn set_task_supervisor(&mut self, supervisor: TaskSupervisor) {
        self.task_supervisor = Some(supervisor);
    }

    /// Declare which forwarding consumer surfaces are active for this session (issue #6359).
    ///
    /// Fixed at session start — call once during bootstrap, before the first
    /// [`spawn`][Self::spawn]. When `surfaces.any()` is `false` (the default) or
    /// `SubAgentConfig::forward_transcript` is `false`, no forwarding sender or drain is ever
    /// constructed for any subagent (FR-007). A [`TaskSupervisor`] must also be wired via
    /// [`set_task_supervisor`][Self::set_task_supervisor] — unlike other subagent lifecycle
    /// tasks, the forward drain never falls back to an untracked spawn (NFR-002).
    pub fn set_forward_surfaces(&mut self, surfaces: ForwardSurfaces) {
        self.forward_surfaces = surfaces;
    }

    /// Wire the bootstrap-level secret-mask registry into the forwarding pipeline (issue
    /// #6359, security Finding 1 / NFR-005).
    ///
    /// When set, every forwarded `Text`/`Thinking` chunk has known vault secrets replaced
    /// with opaque placeholders before reaching any sink (TUI ring, `--bare` stdout) — the
    /// same registry already applied to the outbound-LLM path via `apply_secret_masking`.
    /// Call during bootstrap, before the first [`spawn`][Self::spawn]; has no effect on
    /// already-spawned drains.
    pub fn set_secret_registry(
        &mut self,
        registry: Arc<zeph_sanitizer::secret_mask::SecretMaskRegistry>,
    ) {
        self.secret_registry = Some(registry);
    }

    /// Wire a PII filter into the forwarding pipeline (issue #6359, security Finding 1 /
    /// NFR-005).
    ///
    /// When set, every forwarded `Text`/`Thinking` chunk is scrubbed for emails, phone
    /// numbers, SSNs, etc. before reaching any sink — mirroring the optional `PiiFilter`
    /// layer sub-agent debug dumps already get via `PiiScrubbingDumpSink` (#6407). The filter
    /// itself is unconditionally constructible and self-gates on its own `enabled` config
    /// field, matching the top-level agent's own `PiiFilter` construction convention. Call
    /// during bootstrap, before the first [`spawn`][Self::spawn].
    pub fn set_pii_filter(&mut self, filter: zeph_sanitizer::pii::PiiFilter) {
        self.pii_filter = Some(filter);
    }

    /// Read the current forwarded-transcript tail for `task_id` (up to the last `n` lines).
    ///
    /// Returns an empty vector when forwarding is disabled, no surface is active, or the
    /// task has not forwarded any lines yet. Backing store for
    /// [`crate::manager::SubAgentManager`]'s TUI runtime detail view integration
    /// (FR-005) — see `zeph-core`'s `refresh_subagent_metrics`.
    #[must_use]
    pub fn forwarded_tail(&self, task_id: &str, n: usize) -> Vec<String> {
        crate::forward::forwarded_tail(&self.forward_buffer, task_id, n)
    }

    /// Build a forwarding sender + spawn its per-task drain for `task_id`, if forwarding is
    /// active for this session.
    ///
    /// Returns `None` (no-op) unless `config.forward_transcript` is set, at least one
    /// consumer surface is active, and a [`TaskSupervisor`] is wired — the forward drain must
    /// never fall back to an untracked `tokio::spawn` (NFR-002, P-new-2). The returned
    /// [`crate::forward::ForwardSender`] must be threaded into `AgentLoopArgs::forward` and
    /// owned exclusively by that subagent's own turn loop for the run's lifetime (P-new-3).
    pub(crate) fn maybe_spawn_forward(
        &self,
        task_id: &str,
        def_name: &str,
        forward_transcript: bool,
        content_isolation: &ContentIsolationConfig,
    ) -> Option<crate::forward::ForwardSender> {
        if !forward_transcript || !self.forward_surfaces.any() {
            return None;
        }
        let Some(ref supervisor) = self.task_supervisor else {
            tracing::warn!(
                task_id,
                "subagent transcript forwarding is enabled but no TaskSupervisor is wired — \
                 skipping forwarding for this run rather than spawning an untracked drain \
                 (NFR-002)"
            );
            return None;
        };

        let task_id_arc: Arc<str> = Arc::from(task_id);
        let def_name_arc: Arc<str> = Arc::from(def_name);
        let (sender, rx) =
            crate::forward::new_channel(Arc::clone(&task_id_arc), Arc::clone(&def_name_arc));

        let surfaces = self.forward_surfaces;
        let buffer = Arc::clone(&self.forward_buffer);
        let layers = crate::forward::SanitizeLayers {
            sanitizer: zeph_sanitizer::ContentSanitizer::new(content_isolation),
            secret_registry: self.secret_registry.clone(),
            pii_filter: self.pii_filter.clone(),
        };
        let span = tracing::info_span!(
            "subagent.forward.drain",
            task_id = %task_id_arc,
            tui = surfaces.tui,
            bare = surfaces.bare,
        );
        let drain_task_id = Arc::clone(&task_id_arc);
        let drain_def_name = Arc::clone(&def_name_arc);
        let drain_name: Arc<str> = Arc::from(format!("subagent-forward-drain-{task_id}").as_str());
        let _handle = supervisor.spawn_oneshot(drain_name, move || {
            use tracing::Instrument as _;
            crate::forward::run_forward_drain(
                drain_task_id,
                drain_def_name,
                rx,
                layers,
                surfaces,
                buffer,
            )
            .instrument(span)
        });

        Some(sender)
    }

    /// Inject a [`DefaultWorktreeManager`][zeph_worktree::DefaultWorktreeManager] into the
    /// manager.
    ///
    /// Must be called at most once, before the first [`spawn`][Self::spawn].  When set,
    /// every spawned task acquires the process-level cwd mutex (INV-1) and agents with
    /// `permissions.worktree = true` receive a dedicated git worktree.
    pub fn set_worktree_manager(&mut self, wm: Arc<zeph_worktree::DefaultWorktreeManager>) {
        self.worktree_manager = Some(wm);
    }

    /// Returns the live worktree manager, if the worktree subsystem is enabled for this
    /// session.
    ///
    /// This is the same instance [`spawn`][Self::spawn] uses to create per-subagent
    /// worktrees, so callers (e.g. the `/worktree` slash command) observe this session's
    /// actual live state rather than a fresh disk scan. Its own
    /// `prune_branch_on_remove()` reflects `WorktreeConfig::prune_branch_on_remove` — no
    /// need to retain a separate copy on `SubAgentManager`.
    #[must_use]
    pub fn worktree_manager(&self) -> Option<&Arc<zeph_worktree::DefaultWorktreeManager>> {
        self.worktree_manager.as_ref()
    }

    /// Drain completed hook tasks and spawn a new one if below the limit.
    ///
    /// Polls [`hook_tasks`][Self::hook_tasks] for finished entries so the set does not
    /// accumulate stale handles. When the set is at capacity, logs a warning and skips
    /// the spawn rather than growing unboundedly.
    fn spawn_hook_task<F>(&mut self, future: F)
    where
        F: std::future::Future<Output = ()> + Send + 'static,
    {
        // Drain completed/panicked tasks before checking capacity.
        while self.hook_tasks.try_join_next().is_some() {}
        if self.hook_tasks.len() >= self.max_hook_tasks {
            tracing::warn!(
                limit = self.max_hook_tasks,
                "hook task limit reached — dropping fire-and-forget task"
            );
            return;
        }
        self.hook_tasks.spawn(future);
    }

    /// Spawns a named subagent task under the session [`TaskSupervisor`] if one is configured,
    /// making the task visible in TUI status and abortable on shutdown via
    /// [`TaskSupervisor::shutdown_all`].
    ///
    /// Falls back to a transient local supervisor when no session supervisor has been wired via
    /// [`SubAgentManager::set_task_supervisor`] — the task runs but is not tracked globally.
    /// The returned [`BlockingHandle`] type is identical in both cases so call sites are uniform.
    ///
    /// Every agent loop task's future resolves to a `Result<T, E>` (in practice always
    /// `Result<String, SubAgentError>`), so this classifies via
    /// [`TaskSupervisor::spawn_oneshot_classified`] rather than plain `spawn_oneshot` — an
    /// `Err` produced by a genuinely completed task (e.g. a worktree-quota or cwd-guard setup
    /// failure returned before the agent loop ever starts) is thus classified and logged as a
    /// supervisor-level failure instead of a normal completion (#6257).
    pub(crate) fn spawn_agent_task<F, Fut, T, E>(
        &self,
        name: Arc<str>,
        factory: F,
    ) -> BlockingHandle<Result<T, E>>
    where
        F: FnOnce() -> Fut + Send + 'static,
        Fut: std::future::Future<Output = Result<T, E>> + Send + 'static,
        T: Send + 'static,
        E: Send + 'static,
    {
        if let Some(ref sup) = self.task_supervisor {
            sup.spawn_oneshot_classified(name, factory, Result::is_ok)
        } else {
            let local = TaskSupervisor::new(CancellationToken::new());
            local.spawn_oneshot_classified(name, factory, Result::is_ok)
        }
    }

    /// Reserve `n` concurrency slots for the orchestration scheduler.
    ///
    /// Reserved slots count against the concurrency limit in [`spawn`](Self::spawn) so that
    /// the scheduler can guarantee capacity for tasks it is about to launch. Call
    /// [`release_reservation`](Self::release_reservation) when the scheduler finishes.
    pub fn reserve_slots(&mut self, n: usize) {
        self.reserved_slots = self.reserved_slots.saturating_add(n);
    }

    /// Release `n` previously reserved concurrency slots.
    pub fn release_reservation(&mut self, n: usize) {
        self.reserved_slots = self.reserved_slots.saturating_sub(n);
    }

    /// Configure transcript storage settings.
    pub fn set_transcript_config(&mut self, dir: Option<PathBuf>, max_files: usize) {
        self.transcript_dir = dir;
        self.transcript_max_files = max_files;
    }

    /// Set config-level lifecycle stop hooks (fired when any agent finishes or is cancelled).
    pub fn set_stop_hooks(&mut self, hooks: Vec<super::hooks::HookDef>) {
        self.stop_hooks = hooks;
    }

    /// Inject a fleet registry so spawned sub-agents appear in the fleet dashboard.
    ///
    /// When set, [`spawn`][Self::spawn] registers the session as `Active` and
    /// [`collect`][Self::collect] / [`cancel`][Self::cancel] mark it terminal.
    /// Errors from the registry are logged at `warn` level and never propagate to callers.
    pub fn set_fleet_registry(&mut self, registry: SharedFleetRegistry) {
        self.fleet_registry = Some(registry);
    }

    /// Load sub-agent definitions from the given directories.
    ///
    /// Higher-priority directories should appear first. Name conflicts are resolved
    /// by keeping the first occurrence. Non-existent directories are silently skipped.
    ///
    /// # Errors
    ///
    /// Returns [`SubAgentError`] if any definition file fails to parse.
    pub fn load_definitions(&mut self, dirs: &[PathBuf]) -> Result<(), SubAgentError> {
        let defs = SubAgentDef::load_all(dirs)?;

        // Security gate: non-Default permission_mode is forbidden when the user-level
        // agents directory (~/.zeph/agents/) is one of the load sources. This prevents
        // a crafted agent file from escalating its own privileges.
        // Validation happens here (in the manager) because this is the only place
        // that has full context about which directories were searched.
        //
        // FIX-5: fail-closed — if user_agents_dir is in dirs and a definition has
        // non-Default permission_mode, we cannot verify it did not originate from the
        // user-level dir (SubAgentDef no longer stores source_path), so we reject it.
        let user_agents_dir = dirs::home_dir().map(|h| h.join(".zeph").join("agents"));
        let loads_user_dir = user_agents_dir.as_ref().is_some_and(|user_dir| {
            // FIX-8: log and treat as non-user-level if canonicalize fails.
            match std::fs::canonicalize(user_dir) {
                Ok(canonical_user) => dirs
                    .iter()
                    .filter_map(|d| std::fs::canonicalize(d).ok())
                    .any(|d| d == canonical_user),
                Err(e) => {
                    tracing::warn!(
                        dir = %user_dir.display(),
                        error = %e,
                        "could not canonicalize user agents dir, treating as non-user-level"
                    );
                    false
                }
            }
        });

        if loads_user_dir {
            for def in &defs {
                if def.permissions.permission_mode != PermissionMode::Default {
                    return Err(SubAgentError::Invalid(format!(
                        "sub-agent '{}': non-default permission_mode is not allowed for \
                         user-level definitions (~/.zeph/agents/)",
                        def.name
                    )));
                }
            }
        }

        self.definitions = defs;
        tracing::info!(
            count = self.definitions.len(),
            "sub-agent definitions loaded"
        );
        Ok(())
    }

    /// Load definitions with full scope context for source tracking and security checks.
    ///
    /// The blocking filesystem scan runs on a dedicated thread via
    /// `tokio::task::spawn_blocking` so the tokio worker thread is not stalled (#5108).
    ///
    /// # Errors
    ///
    /// Returns [`SubAgentError`] if a CLI-sourced definition file fails to parse.
    #[tracing::instrument(name = "subagent.manager.load_definitions_with_sources", skip_all)]
    pub async fn load_definitions_with_sources(
        &mut self,
        ordered_paths: &[PathBuf],
        cli_agents: &[PathBuf],
        config_user_dir: Option<&PathBuf>,
        extra_dirs: &[PathBuf],
    ) -> Result<(), SubAgentError> {
        // Clone inputs so they can be moved into spawn_blocking ('static bound).
        let ordered = ordered_paths.to_vec();
        let cli = cli_agents.to_vec();
        let user_dir = config_user_dir.cloned();
        let extra = extra_dirs.to_vec();

        let defs = tokio::task::spawn_blocking(move || {
            SubAgentDef::load_all_with_sources(&ordered, &cli, user_dir.as_ref(), &extra)
        })
        .await
        .map_err(|e| SubAgentError::TaskPanic(format!("load_definitions_with_sources: {e}")))?;

        self.definitions = defs?;
        tracing::info!(
            count = self.definitions.len(),
            "sub-agent definitions loaded"
        );
        Ok(())
    }

    /// Return all loaded definitions.
    #[must_use]
    pub fn definitions(&self) -> &[SubAgentDef] {
        &self.definitions
    }

    /// Return mutable access to the loaded definitions list.
    ///
    /// Intended for test harnesses and dynamic definition registration. Production code
    /// should prefer [`load_definitions`][Self::load_definitions].
    pub fn definitions_mut(&mut self) -> &mut Vec<SubAgentDef> {
        &mut self.definitions
    }

    /// Insert a pre-built handle directly into the active agents map.
    ///
    /// Used in tests to simulate an agent that has already run and left a pending secret
    /// request in its channel without going through the full spawn lifecycle.
    pub fn insert_handle_for_test(&mut self, id: String, handle: SubAgentHandle) {
        self.agents.insert(id, handle);
    }
}

#[cfg(test)]
mod tests;