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 handle: 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            handle: 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 handle: String,
197    pub source_key: String,
198    #[serde(default, skip_serializing_if = "Option::is_none")]
199    pub name: Option<String>,
200    pub source_type: String,
201    #[serde(default)]
202    pub source: serde_json::Value,
203    pub target: RemoteTriggerTargetSummary,
204    #[serde(default = "default_true")]
205    pub enabled: bool,
206}
207
208#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
209pub struct RemoteTriggerTargetSummary {
210    #[serde(default, skip_serializing_if = "Option::is_none")]
211    pub label: Option<String>,
212    pub identity: RemoteProcessIdentity,
213    pub input: RemoteProcessInput,
214    #[serde(default)]
215    pub inputs: RemoteTriggerInputTemplate,
216}
217
218fn default_true() -> bool {
219    true
220}
221
222#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize, JsonSchema)]
223#[serde(transparent)]
224pub struct RemoteTriggerInputTemplate {
225    pub entries: BTreeMap<String, RemoteTriggerInputBinding>,
226}
227
228impl RemoteTriggerInputTemplate {
229    pub fn new(entries: BTreeMap<String, RemoteTriggerInputBinding>) -> Self {
230        Self { entries }
231    }
232
233    pub fn validate(&self, type_name: &'static str) -> Result<(), RemoteProtocolError> {
234        for name in self.entries.keys() {
235            require_non_empty(type_name, "input_template key", name)?;
236        }
237        Ok(())
238    }
239}
240
241#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
242#[serde(tag = "kind", rename_all = "snake_case")]
243pub enum RemoteTriggerInputBinding {
244    Event,
245    Fixed { value: serde_json::Value },
246}
247
248#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
249pub struct RemoteTriggerSubscriptionDraft {
250    pub protocol_version: u32,
251    pub registrant: RemoteProcessOriginator,
252    pub env_ref: RemoteProcessExecutionEnvRef,
253    #[serde(default, skip_serializing_if = "Option::is_none")]
254    pub wake_target: Option<RemoteSessionScope>,
255    #[serde(default, skip_serializing_if = "Option::is_none")]
256    pub name: Option<String>,
257    pub source_type: String,
258    pub source_key: String,
259    #[serde(default)]
260    pub source: serde_json::Value,
261    #[serde(default)]
262    pub payload_schema: serde_json::Value,
263    pub target: RemoteProcessInput,
264    pub target_identity: RemoteProcessIdentity,
265    #[serde(default, skip_serializing_if = "Vec::is_empty")]
266    pub event_types: Vec<RemoteProcessEventType>,
267    #[serde(default)]
268    pub input_template: RemoteTriggerInputTemplate,
269    #[serde(default, skip_serializing_if = "Option::is_none")]
270    pub target_label: Option<String>,
271}
272
273impl RemoteTriggerSubscriptionDraft {
274    pub fn for_process(
275        registrant: RemoteProcessOriginator,
276        env_ref: RemoteProcessExecutionEnvRef,
277        source_type: impl Into<String>,
278        source_key: impl Into<String>,
279        target: RemoteProcessInput,
280        target_identity: RemoteProcessIdentity,
281    ) -> Self {
282        let target_label = target_identity.label.clone();
283        Self {
284            protocol_version: REMOTE_PROTOCOL_VERSION,
285            registrant,
286            env_ref,
287            wake_target: None,
288            name: None,
289            source_type: source_type.into(),
290            source_key: source_key.into(),
291            source: serde_json::Value::Object(serde_json::Map::new()),
292            payload_schema: serde_json::Value::Object(serde_json::Map::new()),
293            target,
294            target_identity,
295            event_types: Vec::new(),
296            input_template: RemoteTriggerInputTemplate::default(),
297            target_label,
298        }
299    }
300
301    pub fn with_name(mut self, name: impl Into<String>) -> Self {
302        self.name = Some(name.into());
303        self
304    }
305
306    pub fn with_source(mut self, source: serde_json::Value) -> Self {
307        self.source = source;
308        self
309    }
310
311    pub fn with_payload_schema(mut self, payload_schema: serde_json::Value) -> Self {
312        self.payload_schema = payload_schema;
313        self
314    }
315
316    pub fn with_wake_target(mut self, wake_target: RemoteSessionScope) -> Self {
317        self.wake_target = Some(wake_target);
318        self
319    }
320
321    pub fn with_event_types(
322        mut self,
323        event_types: impl IntoIterator<Item = RemoteProcessEventType>,
324    ) -> Self {
325        self.event_types = event_types.into_iter().collect();
326        self
327    }
328
329    pub fn with_input_template(mut self, input_template: RemoteTriggerInputTemplate) -> Self {
330        self.input_template = input_template;
331        self
332    }
333
334    pub fn with_target_label(mut self, target_label: impl Into<String>) -> Self {
335        self.target_label = Some(target_label.into());
336        self
337    }
338
339    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
340        ensure_protocol_version(self.protocol_version)?;
341        self.registrant.validate("RemoteTriggerSubscriptionDraft")?;
342        self.env_ref.validate("RemoteTriggerSubscriptionDraft")?;
343        if let Some(wake_target) = &self.wake_target {
344            wake_target.validate("RemoteTriggerSubscriptionDraft")?;
345        }
346        require_non_empty(
347            "RemoteTriggerSubscriptionDraft",
348            "source_type",
349            &self.source_type,
350        )?;
351        require_non_empty(
352            "RemoteTriggerSubscriptionDraft",
353            "source_key",
354            &self.source_key,
355        )?;
356        self.target.validate("RemoteTriggerSubscriptionDraft")?;
357        self.target_identity
358            .validate("RemoteTriggerSubscriptionDraft")?;
359        for event_type in &self.event_types {
360            event_type.validate("RemoteTriggerSubscriptionDraft")?;
361        }
362        validate_remote_trigger_target_label(
363            "RemoteTriggerSubscriptionDraft",
364            self.target_label.as_deref(),
365            self.target_identity.label.as_deref(),
366        )?;
367        self.input_template
368            .validate("RemoteTriggerSubscriptionDraft")
369    }
370}
371
372#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
373pub struct RemoteTriggerSubscriptionRecord {
374    pub subscription_id: String,
375    pub registrant: RemoteProcessOriginator,
376    pub env_ref: RemoteProcessExecutionEnvRef,
377    #[serde(default, skip_serializing_if = "Option::is_none")]
378    pub wake_target: Option<RemoteSessionScope>,
379    pub handle: String,
380    #[serde(default, skip_serializing_if = "Option::is_none")]
381    pub name: Option<String>,
382    pub source_type: String,
383    pub source_key: String,
384    #[serde(default)]
385    pub source: serde_json::Value,
386    #[serde(default)]
387    pub payload_schema: serde_json::Value,
388    pub target: RemoteProcessInput,
389    pub target_identity: RemoteProcessIdentity,
390    #[serde(default, skip_serializing_if = "Vec::is_empty")]
391    pub event_types: Vec<RemoteProcessEventType>,
392    #[serde(default)]
393    pub input_template: RemoteTriggerInputTemplate,
394    #[serde(default, skip_serializing_if = "Option::is_none")]
395    pub target_label: Option<String>,
396    #[serde(default = "default_true")]
397    pub enabled: bool,
398    pub created_at_ms: u64,
399    pub updated_at_ms: u64,
400}
401
402impl RemoteTriggerSubscriptionRecord {
403    pub fn validate(&self, type_name: &'static str) -> Result<(), RemoteProtocolError> {
404        require_non_empty(type_name, "subscription_id", &self.subscription_id)?;
405        self.registrant.validate(type_name)?;
406        self.env_ref.validate(type_name)?;
407        if let Some(wake_target) = &self.wake_target {
408            wake_target.validate(type_name)?;
409        }
410        require_non_empty(type_name, "handle", &self.handle)?;
411        require_non_empty(type_name, "source_type", &self.source_type)?;
412        require_non_empty(type_name, "source_key", &self.source_key)?;
413        self.target.validate(type_name)?;
414        self.target_identity.validate(type_name)?;
415        for event_type in &self.event_types {
416            event_type.validate(type_name)?;
417        }
418        validate_remote_trigger_target_label(
419            type_name,
420            self.target_label.as_deref(),
421            self.target_identity.label.as_deref(),
422        )?;
423        self.input_template.validate(type_name)
424    }
425}
426
427fn validate_remote_trigger_target_label(
428    type_name: &'static str,
429    target_label: Option<&str>,
430    identity_label: Option<&str>,
431) -> Result<(), RemoteProtocolError> {
432    match (target_label, identity_label) {
433        (Some(target_label), Some(identity_label)) if target_label != identity_label => {
434            Err(RemoteProtocolError::InvalidEnvelope {
435                type_name,
436                message: "target_label must match target_identity.label when both are present"
437                    .to_string(),
438            })
439        }
440        _ => Ok(()),
441    }
442}
443
444#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
445pub struct RemoteTriggerRegisterSubscriptionRequest {
446    pub protocol_version: u32,
447    pub draft: RemoteTriggerSubscriptionDraft,
448}
449
450impl RemoteTriggerRegisterSubscriptionRequest {
451    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
452        ensure_protocol_version(self.protocol_version)?;
453        if self.draft.protocol_version != self.protocol_version {
454            return Err(RemoteProtocolError::MismatchedNestedProtocolVersion {
455                parent: "RemoteTriggerRegisterSubscriptionRequest",
456                child: "draft",
457                parent_version: self.protocol_version,
458                child_version: self.draft.protocol_version,
459            });
460        }
461        self.draft.validate()
462    }
463}
464
465#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
466pub struct RemoteTriggerRegisterSubscriptionResult {
467    pub protocol_version: u32,
468    pub record: RemoteTriggerSubscriptionRecord,
469}
470
471impl RemoteTriggerRegisterSubscriptionResult {
472    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
473        ensure_protocol_version(self.protocol_version)?;
474        self.record
475            .validate("RemoteTriggerRegisterSubscriptionResult")
476    }
477}
478
479#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, JsonSchema)]
480pub struct RemoteTriggerListSubscriptionsResponse {
481    pub protocol_version: u32,
482    #[serde(default)]
483    pub subscriptions: Vec<RemoteTriggerSubscriptionRecord>,
484}
485
486impl RemoteTriggerListSubscriptionsResponse {
487    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
488        ensure_protocol_version(self.protocol_version)?;
489        for record in &self.subscriptions {
490            record.validate("RemoteTriggerListSubscriptionsResponse")?;
491        }
492        Ok(())
493    }
494}
495
496#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
497pub struct RemoteTriggerCancelSubscriptionRequest {
498    pub protocol_version: u32,
499    pub session_id: String,
500    pub handle: String,
501}
502
503impl RemoteTriggerCancelSubscriptionRequest {
504    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
505        ensure_protocol_version(self.protocol_version)?;
506        require_non_empty(
507            "RemoteTriggerCancelSubscriptionRequest",
508            "session_id",
509            &self.session_id,
510        )?;
511        require_non_empty(
512            "RemoteTriggerCancelSubscriptionRequest",
513            "handle",
514            &self.handle,
515        )
516    }
517}
518
519#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
520pub struct RemoteTriggerCancelSubscriptionResult {
521    pub protocol_version: u32,
522    pub session_id: String,
523    pub handle: String,
524    pub cancelled: bool,
525}
526
527impl RemoteTriggerCancelSubscriptionResult {
528    pub fn validate(&self) -> Result<(), RemoteProtocolError> {
529        ensure_protocol_version(self.protocol_version)?;
530        require_non_empty(
531            "RemoteTriggerCancelSubscriptionResult",
532            "session_id",
533            &self.session_id,
534        )?;
535        require_non_empty(
536            "RemoteTriggerCancelSubscriptionResult",
537            "handle",
538            &self.handle,
539        )
540    }
541}