Skip to main content

harn_vm/triggers/event/
payloads.rs

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/// Payload emitted by `emit_channel(...)` to channel-source triggers.
50#[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/// Package-owned payload emitted by a Harn connector.
68#[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// Keep the public payload hierarchy PartialEq-only; deriving Eq here would
93// implicitly widen TriggerEvent's public trait contract.
94#[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}