1use std::collections::BTreeMap;
2
3use serde::{Deserialize, Serialize};
4use serde_json::Value as JsonValue;
5use time::OffsetDateTime;
6
7#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
8pub struct CronEventPayload {
9 pub cron_id: Option<String>,
10 pub schedule: Option<String>,
11 #[serde(with = "time::serde::rfc3339")]
12 pub tick_at: OffsetDateTime,
13 pub raw: JsonValue,
14}
15
16#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
17pub struct GenericWebhookPayload {
18 pub source: Option<String>,
19 pub content_type: Option<String>,
20 pub raw: JsonValue,
21}
22
23#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
24pub struct A2aPushPayload {
25 pub task_id: Option<String>,
26 pub task_state: Option<String>,
27 pub artifact: Option<JsonValue>,
28 pub sender: Option<String>,
29 #[serde(default, skip_serializing_if = "Option::is_none")]
30 pub actor_chain: Option<JsonValue>,
31 pub raw: JsonValue,
32 pub kind: String,
33}
34
35#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
36pub struct StreamEventPayload {
37 pub event: String,
38 pub source: Option<String>,
39 pub stream: Option<String>,
40 pub partition: Option<String>,
41 pub offset: Option<String>,
42 pub key: Option<String>,
43 pub timestamp: Option<String>,
44 #[serde(default)]
45 pub headers: BTreeMap<String, String>,
46 pub raw: JsonValue,
47}
48
49#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
51pub struct ChannelEventPayload {
52 pub id: String,
53 pub name: String,
54 pub name_resolved: String,
55 pub scope: String,
56 pub scope_id: String,
57 pub payload: JsonValue,
58 pub emitted_by: String,
59 #[serde(skip_serializing_if = "Option::is_none")]
60 pub tenant_id: Option<String>,
61 #[serde(skip_serializing_if = "Option::is_none")]
62 pub session_id: Option<String>,
63 #[serde(skip_serializing_if = "Option::is_none")]
64 pub pipeline_id: Option<String>,
65}
66
67#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
69pub struct ExtensionProviderPayload {
70 pub provider: String,
71 pub schema_name: String,
72 pub raw: JsonValue,
73}
74
75#[allow(clippy::large_enum_variant)]
76#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
77#[serde(untagged)]
78pub enum ProviderPayload {
79 Extension(ExtensionProviderPayload),
80 Known(KnownProviderPayload),
81}
82
83impl ProviderPayload {
84 pub fn provider(&self) -> &str {
85 match self {
86 Self::Known(known) => known.provider(),
87 Self::Extension(payload) => payload.provider.as_str(),
88 }
89 }
90}
91
92#[allow(clippy::derive_partial_eq_without_eq)]
95#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
96#[serde(tag = "provider")]
97pub enum KnownProviderPayload {
98 #[serde(rename = "cron")]
99 Cron(CronEventPayload),
100 #[serde(rename = "webhook")]
101 Webhook(GenericWebhookPayload),
102 #[serde(rename = "a2a-push")]
103 A2aPush(A2aPushPayload),
104 #[serde(rename = "kafka")]
105 Kafka(StreamEventPayload),
106 #[serde(rename = "nats")]
107 Nats(StreamEventPayload),
108 #[serde(rename = "pulsar")]
109 Pulsar(StreamEventPayload),
110 #[serde(rename = "postgres-cdc")]
111 PostgresCdc(StreamEventPayload),
112 #[serde(rename = "email")]
113 Email(StreamEventPayload),
114 #[serde(rename = "websocket")]
115 Websocket(StreamEventPayload),
116 #[serde(rename = "channel")]
117 Channel(ChannelEventPayload),
118}
119
120impl KnownProviderPayload {
121 pub fn provider(&self) -> &str {
122 match self {
123 Self::Cron(_) => "cron",
124 Self::Webhook(_) => "webhook",
125 Self::A2aPush(_) => "a2a-push",
126 Self::Kafka(_) => "kafka",
127 Self::Nats(_) => "nats",
128 Self::Pulsar(_) => "pulsar",
129 Self::PostgresCdc(_) => "postgres-cdc",
130 Self::Email(_) => "email",
131 Self::Websocket(_) => "websocket",
132 Self::Channel(_) => "channel",
133 }
134 }
135}