1use serde::de::DeserializeOwned;
20use serde::{Deserialize, Serialize};
21use serde_json::Value;
22
23use crate::events::Dispatch;
24
25pub trait TypedEvent {
31 type Payload: Serialize + DeserializeOwned + Clone + Send + Sync + 'static;
33 const NAME: &'static str;
35 const MODE: Dispatch;
37 const AROUND: bool;
39}
40
41#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
48pub struct AgentAdmitPayload {
49 pub tenant_id: String,
50 #[serde(default)]
51 pub monthly: u64,
52 #[serde(default)]
53 pub daily: u64,
54 #[serde(default)]
55 pub requests_per_month: Option<u64>,
56 #[serde(default)]
57 pub requests_per_day: Option<u64>,
58 pub tier: String,
59}
60
61#[derive(Debug, Clone, Copy)]
64pub struct AgentAdmitEvent;
65impl TypedEvent for AgentAdmitEvent {
66 type Payload = AgentAdmitPayload;
67 const NAME: &'static str = crate::events_catalog::ev::AGENT_ADMIT;
68 const MODE: Dispatch = Dispatch::Bail;
69 const AROUND: bool = false;
70}
71
72#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
78pub struct AgentStartedPayload {
79 pub agent_name: String,
80 #[serde(default)]
81 pub run_id: String,
82 #[serde(default)]
83 pub tenant: String,
84 #[serde(default)]
85 pub event: String,
86}
87
88#[derive(Debug, Clone, Copy)]
90pub struct AgentStartedEvent;
91impl TypedEvent for AgentStartedEvent {
92 type Payload = AgentStartedPayload;
93 const NAME: &'static str = crate::events_catalog::ev::AGENT_STARTED;
94 const MODE: Dispatch = Dispatch::Parallel;
95 const AROUND: bool = false;
96}
97
98#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
105pub struct AgentUsagePayload {
106 #[serde(default)]
107 pub tenant: Option<String>,
108 #[serde(default)]
109 pub prompt: i64,
110 #[serde(default)]
111 pub completion: i64,
112 #[serde(default)]
113 pub total: i64,
114}
115
116#[derive(Debug, Clone, Copy)]
118pub struct AgentUsageEvent;
119impl TypedEvent for AgentUsageEvent {
120 type Payload = AgentUsagePayload;
121 const NAME: &'static str = crate::events_catalog::ev::AGENT_USAGE;
122 const MODE: Dispatch = Dispatch::Emit;
123 const AROUND: bool = false;
124}
125
126#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
132pub struct AgentCompletedPayload {
133 pub agent_name: String,
134 #[serde(default)]
135 pub run_id: String,
136 #[serde(default = "default_status")]
137 pub status: String,
138 #[serde(default)]
139 pub event: String,
140}
141
142fn default_status() -> String {
143 "unknown".to_string()
144}
145
146#[derive(Debug, Clone, Copy)]
148pub struct AgentCompletedEvent;
149impl TypedEvent for AgentCompletedEvent {
150 type Payload = AgentCompletedPayload;
151 const NAME: &'static str = crate::events_catalog::ev::AGENT_COMPLETED;
152 const MODE: Dispatch = Dispatch::Emit;
153 const AROUND: bool = false;
154}
155
156#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
158pub struct AgentFailedPayload {
159 pub agent_name: String,
160 #[serde(default)]
161 pub run_id: String,
162 #[serde(default)]
163 pub tenant: String,
164 #[serde(default)]
165 pub event: String,
166}
167
168#[derive(Debug, Clone, Copy)]
171pub struct AgentFailedEvent;
172impl TypedEvent for AgentFailedEvent {
173 type Payload = AgentFailedPayload;
174 const NAME: &'static str = crate::events_catalog::ev::AGENT_FAILED;
175 const MODE: Dispatch = Dispatch::Emit;
176 const AROUND: bool = false;
177}
178
179#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
186pub struct AgentRunRequest {
187 #[serde(default)]
188 pub agent_name: String,
189 #[serde(default)]
190 pub message: String,
191}
192
193#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
197pub struct AgentRunResult {
198 #[serde(default)]
199 pub content: Value,
200 #[serde(default)]
201 pub source: Value,
202 #[serde(default)]
203 pub agent_name: String,
204 #[serde(default)]
205 pub run_id: String,
206 #[serde(default)]
207 pub usage: Option<Value>,
208 #[serde(default)]
209 pub metadata: Option<Value>,
210}
211
212#[derive(Debug, Clone, Copy)]
215pub struct AgentRunEvent;
216impl TypedEvent for AgentRunEvent {
217 type Payload = AgentRunRequest;
218 const NAME: &'static str = crate::events_catalog::ev::AGENT_RUN;
219 const MODE: Dispatch = Dispatch::Waterfall;
220 const AROUND: bool = true;
221}
222
223#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
229pub struct LlmCompleteRequest {
230 #[serde(default)]
231 pub prompt: String,
232}
233
234#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
236pub struct LlmCompleteResult {
237 #[serde(default)]
238 pub prompt: String,
239 #[serde(default)]
240 pub content: Value,
241}
242
243#[derive(Debug, Clone, Copy)]
247pub struct LlmCompleteEvent;
248impl TypedEvent for LlmCompleteEvent {
249 type Payload = LlmCompleteRequest;
250 const NAME: &'static str = crate::events_catalog::ev::LLM_COMPLETE;
251 const MODE: Dispatch = Dispatch::Waterfall;
252 const AROUND: bool = true;
253}
254
255#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
262pub struct LlmGetClientPayload {
263 pub capability: String,
264 #[serde(default)]
265 pub deny: Option<bool>,
266 #[serde(default)]
267 pub model: Option<String>,
268}
269
270#[derive(Debug, Clone, Copy)]
273pub struct LlmGetClientEvent;
274impl TypedEvent for LlmGetClientEvent {
275 type Payload = LlmGetClientPayload;
276 const NAME: &'static str = crate::events_catalog::ev::LLM_GET_CLIENT;
277 const MODE: Dispatch = Dispatch::Waterfall;
278 const AROUND: bool = true;
279}
280
281#[derive(Debug, Clone, PartialEq, Default, Serialize, Deserialize)]
289pub struct LlmMessage {
290 pub role: String,
291 #[serde(default)]
292 pub content: Value,
293 #[serde(default, skip_serializing_if = "Vec::is_empty")]
296 pub parts: Vec<Value>,
297}
298
299#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
301pub struct LlmGeneratePayload {
302 #[serde(default)]
303 pub messages: Vec<LlmMessage>,
304}
305
306#[derive(Debug, Clone, Copy)]
309pub struct LlmGenerateEvent;
310impl TypedEvent for LlmGenerateEvent {
311 type Payload = LlmGeneratePayload;
312 const NAME: &'static str = crate::events_catalog::ev::LLM_GENERATE;
313 const MODE: Dispatch = Dispatch::Waterfall;
314 const AROUND: bool = true;
315}
316
317#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
321pub struct LlmGenerateToolsPayload {
322 #[serde(default)]
323 pub messages: Vec<Value>,
324 #[serde(default)]
325 pub tools: Vec<Value>,
326}
327
328#[derive(Debug, Clone, Copy)]
331pub struct LlmGenerateToolsEvent;
332impl TypedEvent for LlmGenerateToolsEvent {
333 type Payload = LlmGenerateToolsPayload;
334 const NAME: &'static str = crate::events_catalog::ev::LLM_GENERATE_TOOLS;
335 const MODE: Dispatch = Dispatch::Waterfall;
336 const AROUND: bool = true;
337}
338
339#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
345pub struct LlmEmbedRequest {
346 #[serde(default)]
347 pub inputs: Vec<String>,
348}
349
350#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
352pub struct LlmEmbedResponse {
353 #[serde(default)]
354 pub inputs: Vec<String>,
355 #[serde(default)]
356 pub embeddings: Vec<Vec<f32>>,
357}
358
359#[derive(Debug, Clone, Copy)]
363pub struct LlmEmbedEvent;
364impl TypedEvent for LlmEmbedEvent {
365 type Payload = LlmEmbedRequest;
366 const NAME: &'static str = crate::events_catalog::ev::LLM_EMBED;
367 const MODE: Dispatch = Dispatch::Waterfall;
368 const AROUND: bool = true;
369}
370
371#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
377pub struct ToolsExecutePayload {
378 pub name: String,
379 #[serde(default)]
380 pub args: Value,
381}
382
383#[derive(Debug, Clone, Copy)]
386pub struct ToolsExecuteEvent;
387impl TypedEvent for ToolsExecuteEvent {
388 type Payload = ToolsExecutePayload;
389 const NAME: &'static str = crate::events_catalog::ev::TOOLS_EXECUTE;
390 const MODE: Dispatch = Dispatch::Waterfall;
391 const AROUND: bool = true;
392}
393
394#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
396pub struct ToolsListRequest {
397 #[serde(default)]
398 pub tenant: Option<String>,
399}
400
401#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
404pub struct ToolsListResult {
405 #[serde(default)]
406 pub tenant: Option<String>,
407 #[serde(default)]
408 pub tools: Vec<Value>,
409}
410
411#[derive(Debug, Clone, Copy)]
414pub struct ToolsListEvent;
415impl TypedEvent for ToolsListEvent {
416 type Payload = ToolsListRequest;
417 const NAME: &'static str = crate::events_catalog::ev::TOOLS_LIST;
418 const MODE: Dispatch = Dispatch::Waterfall;
419 const AROUND: bool = true;
420}
421
422#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
425pub struct ToolsResolveRequest {
426 pub name: String,
427 #[serde(default)]
428 pub tenant: Option<String>,
429}
430
431#[derive(Debug, Clone, Copy)]
434pub struct ToolsResolveEvent;
435impl TypedEvent for ToolsResolveEvent {
436 type Payload = ToolsResolveRequest;
437 const NAME: &'static str = crate::events_catalog::ev::TOOLS_RESOLVE;
438 const MODE: Dispatch = Dispatch::Waterfall;
439 const AROUND: bool = true;
440}
441
442#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
449pub struct SchedulerBeforeRunPayload {
450 #[serde(default)]
451 pub agent_name: String,
452 #[serde(default)]
453 pub run_id: String,
454 #[serde(default)]
455 pub tenant: Option<String>,
456}
457
458#[derive(Debug, Clone, Copy)]
461pub struct SchedulerBeforeRunEvent;
462impl TypedEvent for SchedulerBeforeRunEvent {
463 type Payload = SchedulerBeforeRunPayload;
464 const NAME: &'static str = crate::events_catalog::ev::SCHEDULER_BEFORE_RUN;
465 const MODE: Dispatch = Dispatch::Waterfall;
466 const AROUND: bool = true;
467}
468
469#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
472pub struct SchedulerAdmitPayload {
473 #[serde(default)]
474 pub agent_name: String,
475 #[serde(default)]
476 pub deny: Option<bool>,
477}
478
479#[derive(Debug, Clone, Copy)]
481pub struct SchedulerAdmitEvent;
482impl TypedEvent for SchedulerAdmitEvent {
483 type Payload = SchedulerAdmitPayload;
484 const NAME: &'static str = crate::events_catalog::ev::SCHEDULER_ADMIT;
485 const MODE: Dispatch = Dispatch::Bail;
486 const AROUND: bool = false;
487}
488
489#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
496pub struct ServiceChangedPayload {
497 pub type_id: String,
498 #[serde(default)]
499 pub event: String,
500}
501
502#[derive(Debug, Clone, Default, Serialize, Deserialize)]
508pub struct SchedulerTickPayload {
509 #[serde(default)]
510 pub due_count: u64,
511 #[serde(default)]
512 pub catchup_count: u64,
513}
514
515#[derive(Debug, Clone, Copy)]
518pub struct SchedulerTickEvent;
519impl TypedEvent for SchedulerTickEvent {
520 type Payload = SchedulerTickPayload;
521 const NAME: &'static str = crate::events_catalog::ev::SCHEDULER_TICK;
522 const MODE: Dispatch = Dispatch::Emit;
523 const AROUND: bool = false;
524}
525
526#[derive(Debug, Clone, Default, Serialize, Deserialize)]
530pub struct ScheduleDispatchedPayload {
531 pub schedule_id: String,
532 #[serde(default)]
533 pub agent_name: String,
534 #[serde(default)]
535 pub tenant_id: String,
536 #[serde(default)]
537 pub is_catchup: bool,
538 #[serde(default)]
539 pub ok: bool,
540 #[serde(default)]
541 pub denied: bool,
542 #[serde(default)]
543 pub error: Option<String>,
544}
545
546#[derive(Debug, Clone, Copy)]
549pub struct ScheduleDispatchedEvent;
550impl TypedEvent for ScheduleDispatchedEvent {
551 type Payload = ScheduleDispatchedPayload;
552 const NAME: &'static str = crate::events_catalog::ev::SCHEDULER_SCHEDULE_DISPATCHED;
553 const MODE: Dispatch = Dispatch::Emit;
554 const AROUND: bool = false;
555}
556
557#[derive(Debug, Clone, Serialize, Deserialize)]
559pub struct PipelineStepStartedPayload {
560 pub pipeline_id: String,
561 pub target_agent: String,
562 #[serde(default)]
563 pub tenant_id: String,
564 #[serde(default)]
565 pub run_id: String,
566}
567
568#[derive(Debug, Clone, Copy)]
571pub struct PipelineStepStartedEvent;
572impl TypedEvent for PipelineStepStartedEvent {
573 type Payload = PipelineStepStartedPayload;
574 const NAME: &'static str = crate::events_catalog::ev::PIPELINE_STEP_STARTED;
575 const MODE: Dispatch = Dispatch::Emit;
576 const AROUND: bool = false;
577}
578
579#[derive(Debug, Clone, Serialize, Deserialize)]
581pub struct PipelineStepFinishedPayload {
582 pub pipeline_id: String,
583 pub target_agent: String,
584 #[serde(default)]
585 pub tenant_id: String,
586 #[serde(default)]
587 pub status: String,
588 #[serde(default)]
589 pub duration_ms: u64,
590 #[serde(default)]
591 pub error: Option<String>,
592}
593
594#[derive(Debug, Clone, Copy)]
597pub struct PipelineStepFinishedEvent;
598impl TypedEvent for PipelineStepFinishedEvent {
599 type Payload = PipelineStepFinishedPayload;
600 const NAME: &'static str = crate::events_catalog::ev::PIPELINE_STEP_FINISHED;
601 const MODE: Dispatch = Dispatch::Emit;
602 const AROUND: bool = false;
603}
604
605#[derive(Debug, Clone, Default, Serialize, Deserialize)]
607pub struct PipelineFanoutCompletedPayload {
608 #[serde(default)]
609 pub source_agent: String,
610 #[serde(default)]
611 pub tenant_id: String,
612 #[serde(default)]
613 pub triggered: Vec<String>,
614}
615
616#[derive(Debug, Clone, Copy)]
619pub struct PipelineFanoutCompletedEvent;
620impl TypedEvent for PipelineFanoutCompletedEvent {
621 type Payload = PipelineFanoutCompletedPayload;
622 const NAME: &'static str = crate::events_catalog::ev::PIPELINE_FANOUT_COMPLETED;
623 const MODE: Dispatch = Dispatch::Emit;
624 const AROUND: bool = false;
625}
626
627#[derive(Debug, Clone, Serialize, Deserialize)]
629pub struct TriggerFiredPayload {
630 pub trigger_id: String,
631 #[serde(default)]
632 pub event_type: String,
633 #[serde(default)]
634 pub target_agent: String,
635 #[serde(default)]
636 pub tenant_id: String,
637}
638
639#[derive(Debug, Clone, Copy)]
642pub struct TriggerFiredEvent;
643impl TypedEvent for TriggerFiredEvent {
644 type Payload = TriggerFiredPayload;
645 const NAME: &'static str = crate::events_catalog::ev::TRIGGER_FIRED;
646 const MODE: Dispatch = Dispatch::Emit;
647 const AROUND: bool = false;
648}
649
650#[derive(Debug, Clone, Copy)]
653pub struct ServiceChangedEvent;
654impl TypedEvent for ServiceChangedEvent {
655 type Payload = ServiceChangedPayload;
656 const NAME: &'static str = crate::events_catalog::ev::SERVICE_CHANGED;
657 const MODE: Dispatch = Dispatch::Emit;
658 const AROUND: bool = false;
659}
660
661#[cfg(test)]
662mod tests {
663 use super::*;
664 use crate::events_catalog::{contract_for, CONTRACTS};
665
666 #[test]
669 fn typed_events_match_catalog_contracts() {
670 let bindings: &[(&'static str, Dispatch, bool)] = &[
671 (
672 AgentAdmitEvent::NAME,
673 AgentAdmitEvent::MODE,
674 AgentAdmitEvent::AROUND,
675 ),
676 (
677 AgentCompletedEvent::NAME,
678 AgentCompletedEvent::MODE,
679 AgentCompletedEvent::AROUND,
680 ),
681 (
682 AgentFailedEvent::NAME,
683 AgentFailedEvent::MODE,
684 AgentFailedEvent::AROUND,
685 ),
686 (
687 AgentRunEvent::NAME,
688 AgentRunEvent::MODE,
689 AgentRunEvent::AROUND,
690 ),
691 (
692 AgentStartedEvent::NAME,
693 AgentStartedEvent::MODE,
694 AgentStartedEvent::AROUND,
695 ),
696 (
697 AgentUsageEvent::NAME,
698 AgentUsageEvent::MODE,
699 AgentUsageEvent::AROUND,
700 ),
701 (
702 LlmCompleteEvent::NAME,
703 LlmCompleteEvent::MODE,
704 LlmCompleteEvent::AROUND,
705 ),
706 (
707 LlmGetClientEvent::NAME,
708 LlmGetClientEvent::MODE,
709 LlmGetClientEvent::AROUND,
710 ),
711 (
712 LlmGenerateEvent::NAME,
713 LlmGenerateEvent::MODE,
714 LlmGenerateEvent::AROUND,
715 ),
716 (
717 LlmGenerateToolsEvent::NAME,
718 LlmGenerateToolsEvent::MODE,
719 LlmGenerateToolsEvent::AROUND,
720 ),
721 (
722 LlmEmbedEvent::NAME,
723 LlmEmbedEvent::MODE,
724 LlmEmbedEvent::AROUND,
725 ),
726 (
727 SchedulerAdmitEvent::NAME,
728 SchedulerAdmitEvent::MODE,
729 SchedulerAdmitEvent::AROUND,
730 ),
731 (
732 SchedulerBeforeRunEvent::NAME,
733 SchedulerBeforeRunEvent::MODE,
734 SchedulerBeforeRunEvent::AROUND,
735 ),
736 (
737 ServiceChangedEvent::NAME,
738 ServiceChangedEvent::MODE,
739 ServiceChangedEvent::AROUND,
740 ),
741 (
742 ToolsExecuteEvent::NAME,
743 ToolsExecuteEvent::MODE,
744 ToolsExecuteEvent::AROUND,
745 ),
746 (
747 ToolsListEvent::NAME,
748 ToolsListEvent::MODE,
749 ToolsListEvent::AROUND,
750 ),
751 (
752 ToolsResolveEvent::NAME,
753 ToolsResolveEvent::MODE,
754 ToolsResolveEvent::AROUND,
755 ),
756 (
757 SchedulerTickEvent::NAME,
758 SchedulerTickEvent::MODE,
759 SchedulerTickEvent::AROUND,
760 ),
761 (
762 ScheduleDispatchedEvent::NAME,
763 ScheduleDispatchedEvent::MODE,
764 ScheduleDispatchedEvent::AROUND,
765 ),
766 (
767 PipelineStepStartedEvent::NAME,
768 PipelineStepStartedEvent::MODE,
769 PipelineStepStartedEvent::AROUND,
770 ),
771 (
772 PipelineStepFinishedEvent::NAME,
773 PipelineStepFinishedEvent::MODE,
774 PipelineStepFinishedEvent::AROUND,
775 ),
776 (
777 PipelineFanoutCompletedEvent::NAME,
778 PipelineFanoutCompletedEvent::MODE,
779 PipelineFanoutCompletedEvent::AROUND,
780 ),
781 (
782 TriggerFiredEvent::NAME,
783 TriggerFiredEvent::MODE,
784 TriggerFiredEvent::AROUND,
785 ),
786 ];
787 for (name, mode, around) in bindings {
788 let contract = contract_for(name)
789 .unwrap_or_else(|| panic!("typed binding {name} missing from catalog"));
790 assert_eq!(&contract.mode, mode, "mode drift for {name}");
791 assert_eq!(&contract.around, around, "around drift for {name}");
792 }
793 assert_eq!(
794 bindings.len(),
795 CONTRACTS.len(),
796 "every catalog event must have exactly one typed binding"
797 );
798 }
799
800 #[test]
804 fn payload_round_trip_and_defaults() {
805 let started = AgentStartedPayload {
806 agent_name: "a".into(),
807 run_id: "r".into(),
808 tenant: "t".into(),
809 event: crate::events_catalog::ev::AGENT_STARTED.into(),
810 };
811 let v = serde_json::to_value(&started).unwrap();
812 assert_eq!(v.get("agent_name").and_then(Value::as_str), Some("a"));
813 let back: AgentStartedPayload = serde_json::from_value(v).unwrap();
814 assert_eq!(back, started);
815
816 let minimal: AgentCompletedPayload =
817 serde_json::from_value(serde_json::json!({ "agent_name": "x" })).unwrap();
818 assert_eq!(minimal.status, "unknown");
819 assert_eq!(minimal.run_id, "");
820
821 let usage = AgentUsagePayload {
822 tenant: None,
823 prompt: 3,
824 completion: 4,
825 total: 7,
826 };
827 let v = serde_json::to_value(&usage).unwrap();
828 assert!(v.get("tenant").map(Value::is_null).unwrap_or(false));
829 let back: AgentUsagePayload = serde_json::from_value(v).unwrap();
830 assert_eq!(back.prompt, 3);
831 assert!(back.tenant.is_none());
832 }
833}