Skip to main content

atman_runtime/
task_registry.rs

1use std::collections::HashMap;
2use std::sync::Arc;
3use std::time::Instant;
4
5use serde::{Deserialize, Serialize};
6use tokio::sync::broadcast;
7use tokio_util::sync::CancellationToken;
8use uuid::Uuid;
9
10#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq, Hash)]
11#[serde(transparent)]
12pub struct TaskId(pub Uuid);
13
14impl TaskId {
15    pub fn now() -> Self {
16        Self(Uuid::now_v7())
17    }
18}
19
20impl std::fmt::Display for TaskId {
21    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
22        self.0.fmt(f)
23    }
24}
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
27#[serde(rename_all = "snake_case")]
28pub enum TaskKind {
29    Bash,
30    Terminal,
31    Flow,
32}
33
34impl TaskKind {
35    pub fn label(self) -> &'static str {
36        match self {
37            TaskKind::Bash => "Bash",
38            TaskKind::Terminal => "Terminal",
39            TaskKind::Flow => "Flow",
40        }
41    }
42}
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
45#[serde(rename_all = "snake_case")]
46pub enum TaskStatus {
47    Running,
48    Killing,
49    Ok,
50    Err,
51    Killed,
52}
53
54#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
55#[serde(rename_all = "snake_case")]
56pub enum TaskTermination {
57    Killed,
58    Suicide,
59}
60
61#[derive(Debug, Clone, Copy, PartialEq, Eq)]
62pub enum KillOutcome {
63    Killed { termination: TaskTermination },
64    NotFound,
65    NotRunning,
66    SelfKillRejected,
67}
68
69impl TaskStatus {
70    pub fn display_label(self) -> &'static str {
71        match self {
72            TaskStatus::Running => "running",
73            TaskStatus::Killing => "stopping",
74            TaskStatus::Ok => "completed",
75            TaskStatus::Err => "failed",
76            TaskStatus::Killed => "stopped",
77        }
78    }
79
80    pub fn is_terminal(self) -> bool {
81        matches!(self, TaskStatus::Ok | TaskStatus::Err | TaskStatus::Killed)
82    }
83
84    pub fn is_running(self) -> bool {
85        matches!(self, TaskStatus::Running | TaskStatus::Killing)
86    }
87}
88
89pub fn normalize_task_label(value: &str) -> Option<String> {
90    crate::message::ToolCallIntent::new(value).map(|label| label.as_str().to_owned())
91}
92
93#[derive(Debug, Clone)]
94pub struct TaskSnapshot {
95    pub id: TaskId,
96    pub kind: TaskKind,
97    pub label: String,
98    pub command: Option<String>,
99    pub status: TaskStatus,
100    pub started_at: Instant,
101    pub ended_at: Option<Instant>,
102    pub source_handle: String,
103    pub session_id: String,
104    pub workspace_id: Option<String>,
105    pub flow_run_id: Option<crate::event::FlowRunId>,
106    pub termination: Option<TaskTermination>,
107}
108
109#[derive(Debug, Clone, PartialEq, Eq)]
110pub struct TaskDisplay {
111    pub label: String,
112    pub command: Option<String>,
113}
114
115impl From<String> for TaskDisplay {
116    fn from(label: String) -> Self {
117        Self {
118            label,
119            command: None,
120        }
121    }
122}
123
124impl From<&str> for TaskDisplay {
125    fn from(label: &str) -> Self {
126        label.to_owned().into()
127    }
128}
129
130impl TaskSnapshot {
131    pub fn elapsed_ms(&self) -> u64 {
132        self.ended_at
133            .unwrap_or_else(Instant::now)
134            .duration_since(self.started_at)
135            .as_millis() as u64
136    }
137
138    pub fn is_running(&self) -> bool {
139        self.status.is_running()
140    }
141}
142
143#[derive(Debug, Clone)]
144pub enum TaskEvent {
145    Registered(TaskSnapshot),
146    StatusChanged {
147        id: TaskId,
148        kind: TaskKind,
149        old: TaskStatus,
150        new: TaskStatus,
151        termination: Option<TaskTermination>,
152    },
153    Reaped {
154        id: TaskId,
155    },
156}
157
158#[derive(Debug, Clone, Default)]
159pub struct TaskFilter {
160    pub kind: Option<TaskKind>,
161    pub status: Option<TaskStatus>,
162    pub session_id: Option<String>,
163}
164
165impl TaskFilter {
166    pub fn all() -> Self {
167        Self::default()
168    }
169
170    pub fn running() -> Self {
171        Self {
172            status: Some(TaskStatus::Running),
173            ..Default::default()
174        }
175    }
176
177    pub fn matches(&self, snap: &TaskSnapshot) -> bool {
178        if let Some(k) = self.kind
179            && snap.kind != k
180        {
181            return false;
182        }
183        if let Some(s) = self.status
184            && snap.status != s
185        {
186            return false;
187        }
188        if let Some(ref sid) = self.session_id
189            && snap.session_id != *sid
190        {
191            return false;
192        }
193        true
194    }
195}
196
197struct TaskEntry {
198    snapshot: TaskSnapshot,
199    cancel: CancellationToken,
200    kill_hook: Option<std::sync::Arc<dyn Fn() + Send + Sync>>,
201}
202
203/// Central management layer for all observable tasks.
204///
205/// Sub-registries (BgRegistry, TermRegistry, Executor, agent_ctrl) register
206/// tasks here on spawn and call `finish` when the task ends. Typed operations
207/// (bash.output, term.input, term.capture) stay on the sub-registries and are
208/// looked up by `source_handle`.
209#[derive(Clone)]
210pub struct TaskRegistry {
211    inner: Arc<std::sync::Mutex<HashMap<TaskId, TaskEntry>>>,
212    event_tx: broadcast::Sender<TaskEvent>,
213}
214
215impl Default for TaskRegistry {
216    fn default() -> Self {
217        let (event_tx, _) = broadcast::channel(256);
218        Self {
219            inner: Arc::new(std::sync::Mutex::new(HashMap::new())),
220            event_tx,
221        }
222    }
223}
224
225fn running_snapshot(
226    kind: TaskKind,
227    display: TaskDisplay,
228    source_handle: String,
229    session_id: String,
230    workspace_id: Option<String>,
231    flow_run_id: Option<crate::event::FlowRunId>,
232) -> TaskSnapshot {
233    TaskSnapshot {
234        id: TaskId::now(),
235        kind,
236        label: display.label,
237        command: display.command,
238        status: TaskStatus::Running,
239        started_at: Instant::now(),
240        ended_at: None,
241        source_handle,
242        session_id,
243        workspace_id,
244        flow_run_id,
245        termination: None,
246    }
247}
248
249impl TaskRegistry {
250    pub fn new() -> Self {
251        Self::default()
252    }
253
254    pub fn register(
255        &self,
256        kind: TaskKind,
257        display: TaskDisplay,
258        source_handle: String,
259        session_id: String,
260        cancel: CancellationToken,
261    ) -> TaskId {
262        self.register_snapshot(
263            running_snapshot(kind, display, source_handle, session_id, None, None),
264            cancel,
265            None,
266        )
267    }
268
269    pub fn register_flow(
270        &self,
271        label: String,
272        source_handle: String,
273        session_id: String,
274        cancel: CancellationToken,
275        workspace_id: Option<String>,
276    ) -> TaskId {
277        self.register_snapshot(
278            running_snapshot(
279                TaskKind::Flow,
280                label.into(),
281                source_handle,
282                session_id,
283                workspace_id,
284                None,
285            ),
286            cancel,
287            None,
288        )
289    }
290
291    pub fn register_flow_with_run_id(
292        &self,
293        label: String,
294        source_handle: String,
295        session_id: String,
296        cancel: CancellationToken,
297        workspace_id: Option<String>,
298        flow_run_id: crate::event::FlowRunId,
299    ) -> TaskId {
300        self.register_snapshot(
301            running_snapshot(
302                TaskKind::Flow,
303                label.into(),
304                source_handle,
305                session_id,
306                workspace_id,
307                Some(flow_run_id),
308            ),
309            cancel,
310            None,
311        )
312    }
313
314    pub fn register_with_kill_hook(
315        &self,
316        kind: TaskKind,
317        display: TaskDisplay,
318        source_handle: String,
319        session_id: String,
320        cancel: CancellationToken,
321        kill_hook: Option<std::sync::Arc<dyn Fn() + Send + Sync>>,
322    ) -> TaskId {
323        self.register_snapshot(
324            running_snapshot(kind, display, source_handle, session_id, None, None),
325            cancel,
326            kill_hook,
327        )
328    }
329
330    fn register_snapshot(
331        &self,
332        snapshot: TaskSnapshot,
333        cancel: CancellationToken,
334        kill_hook: Option<std::sync::Arc<dyn Fn() + Send + Sync>>,
335    ) -> TaskId {
336        let id = snapshot.id.clone();
337        let entry = TaskEntry {
338            snapshot: snapshot.clone(),
339            cancel,
340            kill_hook,
341        };
342        self.inner.lock().unwrap().insert(id.clone(), entry);
343        let _ = self.event_tx.send(TaskEvent::Registered(snapshot));
344        id
345    }
346
347    pub fn lookup(&self, id: &TaskId) -> Option<TaskSnapshot> {
348        self.inner
349            .lock()
350            .unwrap()
351            .get(id)
352            .map(|e| e.snapshot.clone())
353    }
354
355    /// Find a task by its source handle (e.g. "bg_1", "term_2", run_id).
356    pub fn lookup_by_handle(&self, handle: &str) -> Option<TaskSnapshot> {
357        self.inner
358            .lock()
359            .unwrap()
360            .values()
361            .find(|e| e.snapshot.source_handle == handle)
362            .map(|e| e.snapshot.clone())
363    }
364
365    pub fn lookup_by_handle_in_session(
366        &self,
367        handle: &str,
368        session_id: &str,
369    ) -> Option<TaskSnapshot> {
370        self.inner
371            .lock()
372            .unwrap()
373            .values()
374            .find(|entry| {
375                entry.snapshot.source_handle == handle && entry.snapshot.session_id == session_id
376            })
377            .map(|entry| entry.snapshot.clone())
378    }
379
380    pub fn list(&self, filter: &TaskFilter) -> Vec<TaskSnapshot> {
381        let inner = self.inner.lock().unwrap();
382        let mut out: Vec<TaskSnapshot> = inner
383            .values()
384            .map(|e| e.snapshot.clone())
385            .filter(|s| filter.matches(s))
386            .collect();
387        out.sort_by_key(|s| s.started_at);
388        out
389    }
390
391    /// Kill a task at the request of an external operator, such as the TUI.
392    ///
393    /// Operator actions have no Flow identity, so they can never be mistaken
394    /// for a Flow killing itself.
395    pub fn kill_from_operator(&self, id: &TaskId) -> KillOutcome {
396        self.kill_from(id, None, false)
397    }
398
399    pub fn kill_by_handle_from_operator(&self, handle: &str, session_id: &str) -> KillOutcome {
400        let id = self
401            .lookup_by_handle_in_session(handle, session_id)
402            .map(|snapshot| snapshot.id);
403        id.map_or(KillOutcome::NotFound, |id| self.kill_from_operator(&id))
404    }
405
406    pub fn kill_from(
407        &self,
408        id: &TaskId,
409        caller_flow_run_id: Option<&crate::event::FlowRunId>,
410        suicide: bool,
411    ) -> KillOutcome {
412        let mut inner = self.inner.lock().unwrap();
413        let Some(entry) = inner.get_mut(id) else {
414            return KillOutcome::NotFound;
415        };
416        if entry.snapshot.status.is_terminal() {
417            return KillOutcome::NotRunning;
418        }
419        let self_targeting = caller_flow_run_id.is_some()
420            && entry.snapshot.flow_run_id.as_ref() == caller_flow_run_id;
421        if self_targeting && !suicide {
422            return KillOutcome::SelfKillRejected;
423        }
424        let termination = if self_targeting {
425            TaskTermination::Suicide
426        } else {
427            TaskTermination::Killed
428        };
429        let old_status = entry.snapshot.status;
430        let kind = entry.snapshot.kind;
431        entry.snapshot.status = TaskStatus::Killing;
432        entry.snapshot.termination = Some(termination);
433        let cancel = entry.cancel.clone();
434        let hook = entry.kill_hook.clone();
435        drop(inner);
436        cancel.cancel();
437        if let Some(hook) = hook {
438            hook();
439        }
440        let _ = self.event_tx.send(TaskEvent::StatusChanged {
441            id: id.clone(),
442            kind,
443            old: old_status,
444            new: TaskStatus::Killing,
445            termination: Some(termination),
446        });
447        KillOutcome::Killed { termination }
448    }
449
450    /// Transition a task to a terminal status. Called by the owning
451    /// sub-registry when the task finishes.
452    pub fn finish(&self, id: &TaskId, status: TaskStatus) {
453        let mut inner = self.inner.lock().unwrap();
454        let Some(entry) = inner.get_mut(id) else {
455            return;
456        };
457        if entry.snapshot.status.is_terminal() {
458            return;
459        }
460        let old = entry.snapshot.status;
461        let status = if old == TaskStatus::Killing && entry.snapshot.termination.is_some() {
462            TaskStatus::Killed
463        } else {
464            status
465        };
466        entry.snapshot.status = status;
467        entry.snapshot.ended_at = Some(Instant::now());
468        let kind = entry.snapshot.kind;
469        let termination = entry.snapshot.termination;
470        drop(inner);
471        let _ = self.event_tx.send(TaskEvent::StatusChanged {
472            id: id.clone(),
473            kind,
474            old,
475            new: status,
476            termination,
477        });
478    }
479
480    pub fn reap(&self, id: &TaskId) {
481        let mut inner = self.inner.lock().unwrap();
482        let should_remove = inner
483            .get(id)
484            .map(|e| e.snapshot.status.is_terminal())
485            .unwrap_or(false);
486        if should_remove {
487            inner.remove(id);
488            drop(inner);
489            let _ = self.event_tx.send(TaskEvent::Reaped { id: id.clone() });
490        }
491    }
492
493    pub fn subscribe(&self) -> broadcast::Receiver<TaskEvent> {
494        self.event_tx.subscribe()
495    }
496
497    pub fn running_count(&self) -> usize {
498        self.inner
499            .lock()
500            .unwrap()
501            .values()
502            .filter(|e| e.snapshot.status.is_running())
503            .count()
504    }
505}
506
507#[cfg(test)]
508mod tests {
509    use super::*;
510
511    fn cancel() -> CancellationToken {
512        CancellationToken::new()
513    }
514
515    #[test]
516    fn register_and_lookup() {
517        let reg = TaskRegistry::new();
518        let id = reg.register(
519            TaskKind::Bash,
520            "cargo build".into(),
521            "bg_1".into(),
522            "sess".into(),
523            cancel(),
524        );
525        let snap = reg.lookup(&id).expect("found");
526        assert_eq!(snap.kind, TaskKind::Bash);
527        assert_eq!(snap.status, TaskStatus::Running);
528        assert!(snap.ended_at.is_none());
529        assert!(snap.command.is_none());
530    }
531
532    #[test]
533    fn command_is_independent_from_the_user_facing_label() {
534        let reg = TaskRegistry::new();
535        let id = reg.register(
536            TaskKind::Bash,
537            TaskDisplay {
538                label: "运行项目测试".into(),
539                command: Some("cargo test --workspace".into()),
540            },
541            "bg_1".into(),
542            "sess".into(),
543            cancel(),
544        );
545        let snap = reg.lookup(&id).expect("found");
546        assert_eq!(snap.label, "运行项目测试");
547        assert_eq!(snap.command.as_deref(), Some("cargo test --workspace"));
548    }
549
550    #[test]
551    fn lookup_by_handle() {
552        let reg = TaskRegistry::new();
553        let _id = reg.register(
554            TaskKind::Terminal,
555            "vim".into(),
556            "term_1".into(),
557            "sess".into(),
558            cancel(),
559        );
560        let snap = reg.lookup_by_handle("term_1").expect("found");
561        assert_eq!(snap.kind, TaskKind::Terminal);
562        assert!(reg.lookup_by_handle("nope").is_none());
563    }
564
565    #[test]
566    fn list_filters_by_kind_and_status() {
567        let reg = TaskRegistry::new();
568        let b1 = reg.register(
569            TaskKind::Bash,
570            "a".into(),
571            "bg_1".into(),
572            "s".into(),
573            cancel(),
574        );
575        let _t1 = reg.register(
576            TaskKind::Terminal,
577            "vim".into(),
578            "term_1".into(),
579            "s".into(),
580            cancel(),
581        );
582        let _b2 = reg.register(
583            TaskKind::Bash,
584            "ls".into(),
585            "bg_2".into(),
586            "s".into(),
587            cancel(),
588        );
589
590        let bash_only = reg.list(&TaskFilter {
591            kind: Some(TaskKind::Bash),
592            ..Default::default()
593        });
594        assert_eq!(bash_only.len(), 2);
595
596        reg.finish(&b1, TaskStatus::Ok);
597        let running = reg.list(&TaskFilter::running());
598        assert_eq!(running.len(), 2);
599    }
600
601    #[test]
602    fn kill_cancels_token() {
603        let reg = TaskRegistry::new();
604        let tok = cancel();
605        let id = reg.register(
606            TaskKind::Bash,
607            "x".into(),
608            "bg".into(),
609            "s".into(),
610            tok.clone(),
611        );
612        assert_eq!(
613            reg.kill_from_operator(&id),
614            KillOutcome::Killed {
615                termination: TaskTermination::Killed
616            }
617        );
618        assert!(tok.is_cancelled());
619    }
620
621    #[test]
622    fn self_kill_requires_suicide_confirmation() {
623        let reg = TaskRegistry::new();
624        let run_id = crate::event::FlowRunId::now();
625        let token = cancel();
626        let id = reg.register_flow_with_run_id(
627            "flow".into(),
628            "agent".into(),
629            "s".into(),
630            token.clone(),
631            None,
632            run_id.clone(),
633        );
634
635        assert_eq!(
636            reg.kill_from(&id, Some(&run_id), false),
637            KillOutcome::SelfKillRejected
638        );
639        assert!(!token.is_cancelled());
640        assert_eq!(reg.lookup(&id).unwrap().status, TaskStatus::Running);
641    }
642
643    #[test]
644    fn confirmed_self_kill_records_suicide_cause() {
645        let reg = TaskRegistry::new();
646        let run_id = crate::event::FlowRunId::now();
647        let token = cancel();
648        let id = reg.register_flow_with_run_id(
649            "flow".into(),
650            "agent".into(),
651            "s".into(),
652            token.clone(),
653            None,
654            run_id.clone(),
655        );
656
657        assert_eq!(
658            reg.kill_from(&id, Some(&run_id), true),
659            KillOutcome::Killed {
660                termination: TaskTermination::Suicide
661            }
662        );
663        assert!(token.is_cancelled());
664        assert_eq!(
665            reg.lookup(&id).unwrap().termination,
666            Some(TaskTermination::Suicide)
667        );
668    }
669
670    #[test]
671    fn operator_kill_is_not_classified_as_suicide() {
672        let reg = TaskRegistry::new();
673        let target_run_id = crate::event::FlowRunId::now();
674        let id = reg.register_flow_with_run_id(
675            "flow".into(),
676            "agent".into(),
677            "s".into(),
678            cancel(),
679            None,
680            target_run_id,
681        );
682
683        assert_eq!(
684            reg.kill_from_operator(&id),
685            KillOutcome::Killed {
686                termination: TaskTermination::Killed
687            }
688        );
689        assert_eq!(
690            reg.lookup(&id).unwrap().termination,
691            Some(TaskTermination::Killed)
692        );
693    }
694
695    #[test]
696    fn operator_kill_by_handle_is_scoped_to_the_session() {
697        let reg = TaskRegistry::new();
698        let first = reg.register(
699            TaskKind::Terminal,
700            "first".into(),
701            "term_shared".into(),
702            "session_a".into(),
703            cancel(),
704        );
705        let second = reg.register(
706            TaskKind::Terminal,
707            "second".into(),
708            "term_shared".into(),
709            "session_b".into(),
710            cancel(),
711        );
712
713        assert_eq!(
714            reg.kill_by_handle_from_operator("term_shared", "session_b"),
715            KillOutcome::Killed {
716                termination: TaskTermination::Killed
717            }
718        );
719        assert_eq!(reg.lookup(&first).unwrap().status, TaskStatus::Running);
720        assert_eq!(reg.lookup(&second).unwrap().status, TaskStatus::Killing);
721    }
722
723    #[test]
724    fn operator_kill_cannot_be_finished_as_success() {
725        let reg = TaskRegistry::new();
726        let id = reg.register(
727            TaskKind::Terminal,
728            "terminal".into(),
729            "term_1".into(),
730            "session".into(),
731            cancel(),
732        );
733
734        assert!(matches!(
735            reg.kill_from_operator(&id),
736            KillOutcome::Killed { .. }
737        ));
738        reg.finish(&id, TaskStatus::Ok);
739
740        let snapshot = reg.lookup(&id).unwrap();
741        assert_eq!(snapshot.status, TaskStatus::Killed);
742        assert_eq!(snapshot.termination, Some(TaskTermination::Killed));
743    }
744
745    #[test]
746    fn another_flow_can_kill_without_suicide_confirmation() {
747        let reg = TaskRegistry::new();
748        let target_run_id = crate::event::FlowRunId::now();
749        let caller_run_id = crate::event::FlowRunId::now();
750        let id = reg.register_flow_with_run_id(
751            "flow".into(),
752            "agent".into(),
753            "s".into(),
754            cancel(),
755            None,
756            target_run_id,
757        );
758
759        assert_eq!(
760            reg.kill_from(&id, Some(&caller_run_id), false),
761            KillOutcome::Killed {
762                termination: TaskTermination::Killed
763            }
764        );
765    }
766
767    #[test]
768    fn kill_returns_false_for_terminal() {
769        let reg = TaskRegistry::new();
770        let id = reg.register(
771            TaskKind::Bash,
772            "x".into(),
773            "bg".into(),
774            "s".into(),
775            cancel(),
776        );
777        reg.finish(&id, TaskStatus::Ok);
778        assert_eq!(reg.kill_from_operator(&id), KillOutcome::NotRunning);
779    }
780
781    #[test]
782    fn finish_is_idempotent() {
783        let reg = TaskRegistry::new();
784        let id = reg.register(
785            TaskKind::Bash,
786            "x".into(),
787            "bg".into(),
788            "s".into(),
789            cancel(),
790        );
791        reg.finish(&id, TaskStatus::Ok);
792        reg.finish(&id, TaskStatus::Err);
793        let snap = reg.lookup(&id).unwrap();
794        assert_eq!(snap.status, TaskStatus::Ok);
795    }
796
797    #[test]
798    fn reap_removes_terminal_only() {
799        let reg = TaskRegistry::new();
800        let id = reg.register(
801            TaskKind::Bash,
802            "x".into(),
803            "bg".into(),
804            "s".into(),
805            cancel(),
806        );
807        reg.reap(&id);
808        assert!(reg.lookup(&id).is_some());
809        reg.finish(&id, TaskStatus::Ok);
810        reg.reap(&id);
811        assert!(reg.lookup(&id).is_none());
812    }
813
814    #[test]
815    fn subscribe_receives_registered_event() {
816        let reg = TaskRegistry::new();
817        let mut rx = reg.subscribe();
818        let _id = reg.register(
819            TaskKind::Bash,
820            "x".into(),
821            "bg".into(),
822            "s".into(),
823            cancel(),
824        );
825        let ev = rx.try_recv().expect("got event");
826        match ev {
827            TaskEvent::Registered(s) => assert_eq!(s.kind, TaskKind::Bash),
828            _ => panic!("wrong event"),
829        }
830    }
831
832    #[test]
833    fn subscribe_receives_status_changed() {
834        let reg = TaskRegistry::new();
835        let mut rx = reg.subscribe();
836        let id = reg.register(
837            TaskKind::Bash,
838            "x".into(),
839            "bg".into(),
840            "s".into(),
841            cancel(),
842        );
843        let _ = rx.try_recv();
844        reg.finish(&id, TaskStatus::Ok);
845        let ev = rx.try_recv().expect("got status event");
846        match ev {
847            TaskEvent::StatusChanged { new, .. } => assert_eq!(new, TaskStatus::Ok),
848            _ => panic!("wrong event"),
849        }
850    }
851
852    #[test]
853    fn filter_matches_combines() {
854        let snap = TaskSnapshot {
855            id: TaskId::now(),
856            kind: TaskKind::Terminal,
857            label: "vim".into(),
858            command: None,
859            status: TaskStatus::Running,
860            started_at: Instant::now(),
861            ended_at: None,
862            source_handle: "term_1".into(),
863            session_id: "sess_a".into(),
864            workspace_id: None,
865            flow_run_id: None,
866            termination: None,
867        };
868        let f = TaskFilter {
869            kind: Some(TaskKind::Terminal),
870            status: Some(TaskStatus::Running),
871            session_id: Some("sess_a".into()),
872        };
873        assert!(f.matches(&snap));
874
875        let f2 = TaskFilter {
876            kind: Some(TaskKind::Bash),
877            ..Default::default()
878        };
879        assert!(!f2.matches(&snap));
880    }
881
882    #[test]
883    fn task_status_has_stable_display_labels() {
884        assert_eq!(TaskStatus::Running.display_label(), "running");
885        assert_eq!(TaskStatus::Killing.display_label(), "stopping");
886        assert_eq!(TaskStatus::Ok.display_label(), "completed");
887        assert_eq!(TaskStatus::Err.display_label(), "failed");
888        assert_eq!(TaskStatus::Killed.display_label(), "stopped");
889    }
890
891    #[test]
892    fn task_label_normalization_matches_tool_intent_bounds() {
893        assert_eq!(
894            normalize_task_label("  检查   当前状态  ").as_deref(),
895            Some("检查 当前状态")
896        );
897        assert!(normalize_task_label(" \n\t ").is_none());
898        assert_eq!(
899            normalize_task_label(&"x".repeat(121))
900                .unwrap()
901                .chars()
902                .count(),
903            120
904        );
905    }
906}