1use 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}