Skip to main content

lash_remote_protocol/
triggers.rs

1//! Trigger envelopes: occurrence emission, subscriptions, and registrations.
2
3use std::collections::BTreeMap;
4
5use schemars::JsonSchema;
6use serde::{Deserialize, Serialize};
7
8use crate::processes::{
9    RemoteProcessDefinitionIdentity, RemoteProcessEventType, RemoteProcessExecutionEnvRef,
10    RemoteProcessIdentity, RemoteProcessInput, RemoteProcessOriginator, RemoteSessionScope,
11};
12use crate::registry_errors::{RemoteProtocolError, require_non_empty};
13use crate::{REMOTE_PROTOCOL_VERSION, ensure_protocol_version};
14
15#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
16pub struct RemoteTriggerOccurrenceRequest {
17    pub protocol_version: u32,
18    pub source_type: String,
19    pub source_key: String,
20    #[serde(default)]
21    pub payload: serde_json::Value,
22    pub idempotency_key: String,
23    #[serde(default, skip_serializing_if = "Option::is_none")]
24    pub source: Option<serde_json::Value>,
25    #[serde(default, skip_serializing_if = "Option::is_none")]
26    pub session_id: Option<String>,
27}
28
29impl RemoteTriggerOccurrenceRequest {
30    pub fn new(
31        source_type: impl Into<String>,
32        source_key: impl Into<String>,
33        payload: serde_json::Value,
34        idempotency_key: impl Into<String>,
35    ) -> Self {
36        Self {
37            protocol_version: REMOTE_PROTOCOL_VERSION,
38            source_type: source_type.into(),
39            source_key: source_key.into(),
40            payload,
41            idempotency_key: idempotency_key.into(),
42            source: None,
43            session_id: None,
44        }
45    }
46
47    pub fn with_source(mut self, source: serde_json::Value) -> Self {
48        self.source = Some(source);
49        self
50    }
51
52    pub fn for_session(mut self, session_id: impl Into<String>) -> Self {
53        self.session_id = Some(session_id.into());
54        self
55    }
56
57    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
58        ensure_protocol_version(self.protocol_version)?;
59        require_non_empty(
60            "RemoteTriggerOccurrenceRequest",
61            "source_type",
62            &self.source_type,
63        )?;
64        require_non_empty(
65            "RemoteTriggerOccurrenceRequest",
66            "source_key",
67            &self.source_key,
68        )?;
69        require_non_empty(
70            "RemoteTriggerOccurrenceRequest",
71            "idempotency_key",
72            &self.idempotency_key,
73        )?;
74        if let Some(session_id) = &self.session_id {
75            require_non_empty("RemoteTriggerOccurrenceRequest", "session_id", session_id)?;
76        }
77        Ok(())
78    }
79}
80
81#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
82pub struct RemoteTriggerOccurrenceRecord {
83    pub occurrence_id: String,
84    pub source_type: String,
85    pub source_key: String,
86    #[serde(default)]
87    pub payload: serde_json::Value,
88    pub idempotency_key: String,
89    #[serde(default, skip_serializing_if = "Option::is_none")]
90    pub source: Option<serde_json::Value>,
91    #[serde(default, skip_serializing_if = "Option::is_none")]
92    pub session_id: Option<String>,
93    pub occurred_at_ms: u64,
94}
95
96#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
97#[serde(rename_all = "snake_case")]
98pub enum RemoteTriggerDeliveryEmitOutcome {
99    Started,
100    AlreadyReserved,
101    Failed { reason: String },
102}
103
104#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
105pub struct RemoteTriggerDeliveryEmitReport {
106    pub occurrence_id: String,
107    pub subscription_id: String,
108    pub process_id: String,
109    pub outcome: RemoteTriggerDeliveryEmitOutcome,
110}
111
112#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
113pub struct RemoteTriggerEmitReport {
114    pub protocol_version: u32,
115    #[serde(default, skip_serializing_if = "String::is_empty")]
116    pub occurrence_id: String,
117    #[serde(default, skip_serializing_if = "Vec::is_empty")]
118    pub deliveries: Vec<RemoteTriggerDeliveryEmitReport>,
119}
120
121impl RemoteTriggerEmitReport {
122    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
123        ensure_protocol_version(self.protocol_version)
124    }
125}
126
127#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
128pub struct RemoteTriggerSubscriptionFilter {
129    pub protocol_version: u32,
130    #[serde(default, skip_serializing_if = "Option::is_none")]
131    pub registrant_scope_id: Option<String>,
132    #[serde(default, skip_serializing_if = "Option::is_none")]
133    pub session_id: Option<String>,
134    #[serde(default, skip_serializing_if = "Option::is_none")]
135    pub subscription_key: Option<String>,
136    #[serde(default, skip_serializing_if = "Option::is_none")]
137    pub name: Option<String>,
138    #[serde(default, skip_serializing_if = "Option::is_none")]
139    pub source_type: Option<String>,
140    #[serde(default, skip_serializing_if = "Option::is_none")]
141    pub source_key: Option<String>,
142    #[serde(default, skip_serializing_if = "Option::is_none")]
143    pub target: Option<RemoteProcessDefinitionIdentity>,
144    #[serde(default, skip_serializing_if = "Option::is_none")]
145    pub enabled: Option<bool>,
146}
147
148impl Default for RemoteTriggerSubscriptionFilter {
149    fn default() -> Self {
150        Self {
151            protocol_version: REMOTE_PROTOCOL_VERSION,
152            registrant_scope_id: None,
153            session_id: None,
154            subscription_key: None,
155            name: None,
156            source_type: None,
157            source_key: None,
158            target: None,
159            enabled: None,
160        }
161    }
162}
163
164impl RemoteTriggerSubscriptionFilter {
165    pub fn for_session(session_id: impl Into<String>) -> Self {
166        Self {
167            protocol_version: REMOTE_PROTOCOL_VERSION,
168            session_id: Some(session_id.into()),
169            ..Self::default()
170        }
171    }
172
173    pub fn for_registrant_scope(scope_id: impl Into<String>) -> Self {
174        Self {
175            protocol_version: REMOTE_PROTOCOL_VERSION,
176            registrant_scope_id: Some(scope_id.into()),
177            ..Self::default()
178        }
179    }
180
181    pub fn for_source_type(source_type: impl Into<String>) -> Self {
182        Self {
183            protocol_version: REMOTE_PROTOCOL_VERSION,
184            source_type: Some(source_type.into()),
185            ..Self::default()
186        }
187    }
188
189    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
190        ensure_protocol_version(self.protocol_version)
191    }
192}
193
194#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
195pub struct RemoteTriggerRegistration {
196    pub subscription_key: String,
197    pub incarnation: String,
198    pub revision: u64,
199    pub registrant: RemoteProcessOriginator,
200    pub manifest_membership: RemoteTriggerManifestMembership,
201    pub source_key: String,
202    #[serde(default, skip_serializing_if = "Option::is_none")]
203    pub name: Option<String>,
204    pub source_type: String,
205    #[serde(default)]
206    pub source: serde_json::Value,
207    pub target: RemoteTriggerTargetSummary,
208    #[serde(default = "default_true")]
209    pub enabled: bool,
210}
211
212#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
213#[serde(rename_all = "snake_case")]
214pub enum RemoteTriggerManifestMembership {
215    PresentInCurrentArtifact,
216    Orphaned,
217    #[default]
218    Unknown,
219}
220
221#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
222pub struct RemoteTriggerTargetSummary {
223    #[serde(default, skip_serializing_if = "Option::is_none")]
224    pub label: Option<String>,
225    pub identity: RemoteProcessIdentity,
226    pub input: RemoteProcessInput,
227    #[serde(default)]
228    pub inputs: RemoteTriggerInputTemplate,
229}
230
231fn default_true() -> bool {
232    true
233}
234
235#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize, JsonSchema)]
236#[serde(transparent)]
237pub struct RemoteTriggerInputTemplate {
238    pub entries: BTreeMap<String, RemoteTriggerInputBinding>,
239}
240
241impl RemoteTriggerInputTemplate {
242    pub fn new(entries: BTreeMap<String, RemoteTriggerInputBinding>) -> Self {
243        Self { entries }
244    }
245
246    pub fn validate(&self, type_name: &'static str) -> Result<(), RemoteProtocolError> {
247        for name in self.entries.keys() {
248            require_non_empty(type_name, "input_template key", name)?;
249        }
250        Ok(())
251    }
252}
253
254#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
255#[serde(tag = "kind", rename_all = "snake_case")]
256pub enum RemoteTriggerInputBinding {
257    Event,
258    Fixed { value: serde_json::Value },
259}
260
261#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
262#[serde(tag = "type", rename_all = "snake_case")]
263pub enum RemoteTriggerOwnerScope {
264    Session { session_id: String },
265    Host { binding_id: String },
266    Platform,
267}
268
269#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
270pub struct RemoteTriggerSubscriptionDraft {
271    pub protocol_version: u32,
272    pub subscription_key: String,
273    pub env_ref: RemoteProcessExecutionEnvRef,
274    #[serde(default, skip_serializing_if = "Option::is_none")]
275    pub wake_target: Option<RemoteSessionScope>,
276    #[serde(default, skip_serializing_if = "Option::is_none")]
277    pub name: Option<String>,
278    pub source_type: String,
279    pub source_key: String,
280    #[serde(default)]
281    pub source: serde_json::Value,
282    #[serde(default)]
283    pub payload_schema: serde_json::Value,
284    pub target: RemoteProcessInput,
285    pub target_identity: RemoteProcessIdentity,
286    #[serde(default, skip_serializing_if = "Vec::is_empty")]
287    pub event_types: Vec<RemoteProcessEventType>,
288    #[serde(default)]
289    pub input_template: RemoteTriggerInputTemplate,
290    #[serde(default, skip_serializing_if = "Option::is_none")]
291    pub target_label: Option<String>,
292}
293
294impl RemoteTriggerSubscriptionDraft {
295    pub fn for_process(
296        subscription_key: impl Into<String>,
297        env_ref: RemoteProcessExecutionEnvRef,
298        source_type: impl Into<String>,
299        source_key: impl Into<String>,
300        target: RemoteProcessInput,
301        target_identity: RemoteProcessIdentity,
302    ) -> Self {
303        let target_label = target_identity.label.clone();
304        Self {
305            protocol_version: REMOTE_PROTOCOL_VERSION,
306            subscription_key: subscription_key.into(),
307            env_ref,
308            wake_target: None,
309            name: None,
310            source_type: source_type.into(),
311            source_key: source_key.into(),
312            source: serde_json::Value::Object(serde_json::Map::new()),
313            payload_schema: serde_json::Value::Object(serde_json::Map::new()),
314            target,
315            target_identity,
316            event_types: Vec::new(),
317            input_template: RemoteTriggerInputTemplate::default(),
318            target_label,
319        }
320    }
321
322    pub fn with_name(mut self, name: impl Into<String>) -> Self {
323        self.name = Some(name.into());
324        self
325    }
326
327    pub fn with_source(mut self, source: serde_json::Value) -> Self {
328        self.source = source;
329        self
330    }
331
332    pub fn with_payload_schema(mut self, payload_schema: serde_json::Value) -> Self {
333        self.payload_schema = payload_schema;
334        self
335    }
336
337    pub fn with_wake_target(mut self, wake_target: RemoteSessionScope) -> Self {
338        self.wake_target = Some(wake_target);
339        self
340    }
341
342    pub fn with_event_types(
343        mut self,
344        event_types: impl IntoIterator<Item = RemoteProcessEventType>,
345    ) -> Self {
346        self.event_types = event_types.into_iter().collect();
347        self
348    }
349
350    pub fn with_input_template(mut self, input_template: RemoteTriggerInputTemplate) -> Self {
351        self.input_template = input_template;
352        self
353    }
354
355    pub fn with_target_label(mut self, target_label: impl Into<String>) -> Self {
356        self.target_label = Some(target_label.into());
357        self
358    }
359
360    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
361        ensure_protocol_version(self.protocol_version)?;
362        require_non_empty(
363            "RemoteTriggerSubscriptionDraft",
364            "subscription_key",
365            &self.subscription_key,
366        )?;
367        self.env_ref.validate("RemoteTriggerSubscriptionDraft")?;
368        if let Some(wake_target) = &self.wake_target {
369            wake_target.validate("RemoteTriggerSubscriptionDraft")?;
370        }
371        require_non_empty(
372            "RemoteTriggerSubscriptionDraft",
373            "source_type",
374            &self.source_type,
375        )?;
376        require_non_empty(
377            "RemoteTriggerSubscriptionDraft",
378            "source_key",
379            &self.source_key,
380        )?;
381        self.target.validate("RemoteTriggerSubscriptionDraft")?;
382        self.target_identity
383            .validate("RemoteTriggerSubscriptionDraft")?;
384        for event_type in &self.event_types {
385            event_type.validate("RemoteTriggerSubscriptionDraft")?;
386        }
387        validate_remote_trigger_target_label(
388            "RemoteTriggerSubscriptionDraft",
389            self.target_label.as_deref(),
390            self.target_identity.label.as_deref(),
391        )?;
392        self.input_template
393            .validate("RemoteTriggerSubscriptionDraft")
394    }
395}
396
397#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
398pub struct RemoteTriggerSubscriptionRecord {
399    pub subscription_id: String,
400    pub owner_scope: RemoteTriggerOwnerScope,
401    pub subscription_key: String,
402    pub incarnation: String,
403    pub revision: u64,
404    pub definition_hash: String,
405    pub registrant: RemoteProcessOriginator,
406    pub env_ref: RemoteProcessExecutionEnvRef,
407    #[serde(default, skip_serializing_if = "Option::is_none")]
408    pub wake_target: Option<RemoteSessionScope>,
409    #[serde(default, skip_serializing_if = "Option::is_none")]
410    pub name: Option<String>,
411    pub source_type: String,
412    pub source_key: String,
413    #[serde(default)]
414    pub source: serde_json::Value,
415    #[serde(default)]
416    pub payload_schema: serde_json::Value,
417    pub target: RemoteProcessInput,
418    pub target_identity: RemoteProcessIdentity,
419    #[serde(default, skip_serializing_if = "Vec::is_empty")]
420    pub event_types: Vec<RemoteProcessEventType>,
421    #[serde(default)]
422    pub input_template: RemoteTriggerInputTemplate,
423    #[serde(default, skip_serializing_if = "Option::is_none")]
424    pub target_label: Option<String>,
425    #[serde(default = "default_true")]
426    pub enabled: bool,
427    #[serde(default)]
428    pub tombstoned: bool,
429    #[serde(default, skip_serializing_if = "Option::is_none")]
430    pub deleted_at_ms: Option<u64>,
431    pub created_at_ms: u64,
432    pub updated_at_ms: u64,
433}
434
435impl RemoteTriggerSubscriptionRecord {
436    pub fn validate(&self, type_name: &'static str) -> Result<(), RemoteProtocolError> {
437        require_non_empty(type_name, "subscription_id", &self.subscription_id)?;
438        require_non_empty(type_name, "subscription_key", &self.subscription_key)?;
439        require_non_empty(type_name, "incarnation", &self.incarnation)?;
440        require_non_empty(type_name, "definition_hash", &self.definition_hash)?;
441        self.registrant.validate(type_name)?;
442        self.env_ref.validate(type_name)?;
443        if let Some(wake_target) = &self.wake_target {
444            wake_target.validate(type_name)?;
445        }
446        require_non_empty(type_name, "source_type", &self.source_type)?;
447        require_non_empty(type_name, "source_key", &self.source_key)?;
448        self.target.validate(type_name)?;
449        self.target_identity.validate(type_name)?;
450        for event_type in &self.event_types {
451            event_type.validate(type_name)?;
452        }
453        validate_remote_trigger_target_label(
454            type_name,
455            self.target_label.as_deref(),
456            self.target_identity.label.as_deref(),
457        )?;
458        self.input_template.validate(type_name)
459    }
460}
461
462fn validate_remote_trigger_target_label(
463    type_name: &'static str,
464    target_label: Option<&str>,
465    identity_label: Option<&str>,
466) -> Result<(), RemoteProtocolError> {
467    match (target_label, identity_label) {
468        (Some(target_label), Some(identity_label)) if target_label != identity_label => {
469            Err(RemoteProtocolError::InvalidEnvelope {
470                type_name,
471                message: "target_label must match target_identity.label when both are present"
472                    .to_string(),
473            })
474        }
475        _ => Ok(()),
476    }
477}
478
479#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
480pub struct RemoteTriggerRegisterSubscriptionRequest {
481    pub protocol_version: u32,
482    pub draft: RemoteTriggerSubscriptionDraft,
483}
484
485impl RemoteTriggerRegisterSubscriptionRequest {
486    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
487        ensure_protocol_version(self.protocol_version)?;
488        if self.draft.protocol_version != self.protocol_version {
489            return Err(RemoteProtocolError::MismatchedNestedProtocolVersion {
490                parent: "RemoteTriggerRegisterSubscriptionRequest",
491                child: "draft",
492                parent_version: self.protocol_version,
493                child_version: self.draft.protocol_version,
494            });
495        }
496        self.draft.validate()
497    }
498}
499
500#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
501pub struct RemoteTriggerRegisterSubscriptionResult {
502    pub protocol_version: u32,
503    pub record: RemoteTriggerSubscriptionRecord,
504}
505
506impl RemoteTriggerRegisterSubscriptionResult {
507    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
508        ensure_protocol_version(self.protocol_version)?;
509        self.record
510            .validate("RemoteTriggerRegisterSubscriptionResult")
511    }
512}
513
514#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
515pub struct RemoteTriggerListSubscriptionsResponse {
516    pub protocol_version: u32,
517    #[serde(default)]
518    pub subscriptions: Vec<RemoteTriggerSubscriptionRecord>,
519}
520
521impl RemoteTriggerListSubscriptionsResponse {
522    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
523        ensure_protocol_version(self.protocol_version)?;
524        for record in &self.subscriptions {
525            record.validate("RemoteTriggerListSubscriptionsResponse")?;
526        }
527        Ok(())
528    }
529}