Skip to main content

canwu_sim/runtime/
boundary.rs

1use super::{
2    CanwuError, DecisionOptionWeight, DomainRecordChange, DomainRecordMutation, RandomSample,
3    RandomStreamKey, SimulationView, StateKey, StateVisibility, SystemCadence,
4};
5use canwu_core::{
6    BoundaryId, CommandAttemptId, CommandId, CommandRequestId, DecisionRequestId, DecisionTicketId,
7    EntityRef, EventId, IngressId, KnowledgeHolderRef, KnowledgeSchemaId, RandomDrawId,
8};
9use canwu_knowledge::{KnowledgeRecord, KnowledgeRecordDraft};
10use canwu_time::{SimDuration, SimTime};
11use serde::{Deserialize, Serialize};
12use serde_json::Value;
13
14#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
15pub struct ReservationPoolKey {
16    pub state: StateKey,
17    pub entity: EntityRef,
18    pub resource: String,
19}
20
21impl ReservationPoolKey {
22    #[must_use]
23    pub fn new(state: StateKey, entity: EntityRef, resource: impl Into<String>) -> Self {
24        Self {
25            state,
26            entity,
27            resource: resource.into(),
28        }
29    }
30}
31
32#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
33pub struct ReservationRef {
34    pub plugin: String,
35    pub system: String,
36    pub request: String,
37}
38
39impl ReservationRef {
40    #[must_use]
41    pub fn new(
42        plugin: impl Into<String>,
43        system: impl Into<String>,
44        request: impl Into<String>,
45    ) -> Self {
46        Self {
47            plugin: plugin.into(),
48            system: system.into(),
49            request: request.into(),
50        }
51    }
52}
53
54#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
55pub struct ReservationOffer {
56    pub pool: ReservationPoolKey,
57    pub capacity: u64,
58}
59
60#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
61pub struct ReservationRequest {
62    pub request: String,
63    pub pool: ReservationPoolKey,
64    pub quantity: u64,
65    pub priority: i32,
66    pub tie_break: String,
67}
68
69#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
70pub struct ReservationOfferRecord {
71    pub plugin: String,
72    pub system: String,
73    pub offer: ReservationOffer,
74}
75
76#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
77pub struct ReservationRequestRecord {
78    pub reservation: ReservationRef,
79    pub request: ReservationRequest,
80}
81
82#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
83#[serde(rename_all = "snake_case")]
84pub enum ReservationDisposition {
85    Fulfilled,
86    Partial,
87    Rejected,
88}
89
90#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
91pub struct ReservationAllocation {
92    pub reservation: ReservationRef,
93    pub pool: ReservationPoolKey,
94    pub requested: u64,
95    pub granted: u64,
96    pub remaining_after: u64,
97    pub disposition: ReservationDisposition,
98}
99
100#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
101#[serde(tag = "type", rename_all = "snake_case")]
102pub enum BoundaryDirective {
103    SetComponent {
104        state: StateKey,
105        entity: EntityRef,
106        component: String,
107        value: Value,
108        summary: String,
109    },
110    MutateRecord {
111        mutation: DomainRecordMutation,
112        summary: String,
113    },
114    Emit {
115        event_type: String,
116        summary: String,
117        affected: Vec<EntityRef>,
118    },
119    ScheduleIngress {
120        after: SimDuration,
121        packet_type: String,
122        priority: i32,
123        payload: Value,
124        affected: Vec<EntityRef>,
125    },
126    SchedulePluginIngress {
127        target_plugin: String,
128        after: SimDuration,
129        packet_type: String,
130        priority: i32,
131        payload: Value,
132        affected: Vec<EntityRef>,
133    },
134    ResolveDecisionRandomly {
135        resolution: RandomDecisionResolution,
136    },
137    PublishKnowledge {
138        holder: KnowledgeHolderRef,
139        visibility: StateVisibility,
140        producer_correlation: Option<String>,
141        records: Vec<KnowledgeRecordDraft>,
142        summary: String,
143    },
144}
145
146#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
147pub struct RandomDecisionResolution {
148    pub priority: i32,
149    pub decision_request_id: DecisionRequestId,
150    #[serde(default, skip_serializing_if = "Option::is_none")]
151    pub command_request_id: Option<CommandRequestId>,
152    pub ticket_id: DecisionTicketId,
153    pub expected_version: u64,
154    pub controller_id: String,
155    pub sample: RandomSample,
156    pub option_weights: Vec<DecisionOptionWeight>,
157}
158
159#[derive(Clone, Debug, Default, Deserialize, PartialEq, Serialize)]
160pub struct BoundaryProposal {
161    pub offers: Vec<ReservationOffer>,
162    pub requests: Vec<ReservationRequest>,
163    pub directives: Vec<BoundaryDirective>,
164}
165
166#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
167pub struct KnowledgeWriteGrant {
168    pub schema: KnowledgeSchemaId,
169    pub visibilities: Vec<StateVisibility>,
170}
171
172#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
173pub struct PluginIngressTarget {
174    pub target_plugin: String,
175    pub packet_type: String,
176}
177
178#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
179pub struct BoundarySystemContract {
180    pub name: String,
181    pub phase: crate::BoundaryPhase,
182    pub cadence: SystemCadence,
183    pub reads: Vec<StateKey>,
184    pub writes: Vec<StateKey>,
185    pub emits: Vec<String>,
186    pub reservation_offers: Vec<StateKey>,
187    pub reservation_requests: Vec<StateKey>,
188    pub reservation_reads: Vec<ReservationRef>,
189    #[serde(default)]
190    pub random_streams: Vec<RandomStreamKey>,
191    #[serde(default, skip_serializing_if = "Vec::is_empty")]
192    pub knowledge_writes: Vec<KnowledgeWriteGrant>,
193    #[serde(default, skip_serializing_if = "Vec::is_empty")]
194    pub plugin_ingress_targets: Vec<PluginIngressTarget>,
195    pub visibility: StateVisibility,
196}
197
198impl BoundarySystemContract {
199    #[must_use]
200    pub fn new(
201        name: impl Into<String>,
202        phase: crate::BoundaryPhase,
203        cadence: SystemCadence,
204    ) -> Self {
205        Self {
206            name: name.into(),
207            phase,
208            cadence,
209            reads: Vec::new(),
210            writes: Vec::new(),
211            emits: Vec::new(),
212            reservation_offers: Vec::new(),
213            reservation_requests: Vec::new(),
214            reservation_reads: Vec::new(),
215            random_streams: Vec::new(),
216            knowledge_writes: Vec::new(),
217            plugin_ingress_targets: Vec::new(),
218            visibility: StateVisibility::NextBoundary,
219        }
220    }
221}
222
223#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
224pub struct BoundaryContext {
225    pub boundary_id: BoundaryId,
226    pub at: SimTime,
227    pub phase: crate::BoundaryPhase,
228    pub plugin: String,
229    pub system: String,
230    pub admitted_attempts: Vec<CommandAttemptId>,
231    pub admitted_commands: Vec<CommandId>,
232    pub admitted_ingress: Vec<IngressId>,
233    pub admitted_events: Vec<EventId>,
234    pub emitted_events: Vec<EventId>,
235}
236
237pub type BoundarySystemHandler =
238    fn(&SimulationView<'_>, &BoundaryContext) -> Result<BoundaryProposal, CanwuError>;
239
240#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
241pub struct BoundaryRequest {
242    pub at: SimTime,
243    pub cadences: Vec<SystemCadence>,
244}
245
246impl BoundaryRequest {
247    #[must_use]
248    pub const fn at(at: SimTime) -> Self {
249        Self {
250            at,
251            cadences: Vec::new(),
252        }
253    }
254
255    #[must_use]
256    pub fn with_cadence(mut self, cadence: SystemCadence) -> Self {
257        self.cadences.push(cadence);
258        self
259    }
260}
261
262#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
263pub struct BoundaryChange {
264    pub plugin: String,
265    pub system: String,
266    pub state: StateKey,
267    pub entity: EntityRef,
268    pub component: String,
269    pub previous: Option<Value>,
270    pub value: Value,
271    pub visibility: StateVisibility,
272    pub summary: String,
273}
274
275#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
276#[serde(tag = "kind", rename_all = "snake_case")]
277pub enum BoundaryEmissionKind {
278    Change { change_index: u64 },
279    RecordChange { change_index: u64 },
280    KnowledgeChange { change_index: u64 },
281    Explicit,
282}
283
284#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
285pub struct BoundaryEmission {
286    pub plugin: String,
287    pub system: String,
288    pub event: EventId,
289    pub kind: BoundaryEmissionKind,
290}
291
292/// Durable, idempotent external-delivery identity derived from committed
293/// boundary evidence. The engine creates one entry for every emission; a host
294/// may deliver it at least once and use `delivery_id` as its idempotency key.
295#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
296pub struct OutboxEntry {
297    pub delivery_id: String,
298    pub boundary: BoundaryId,
299    pub event: EventId,
300    pub emission_index: u64,
301    pub plugin: String,
302    pub system: String,
303}
304
305#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
306pub struct BoundaryIngressGeneration {
307    pub ingress: IngressId,
308    pub plugin: String,
309    pub system: String,
310    pub phase: crate::BoundaryPhase,
311    pub visibility: StateVisibility,
312}
313
314#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
315pub struct BoundaryKnowledgeChange {
316    pub plugin: String,
317    pub system: String,
318    pub phase: crate::BoundaryPhase,
319    pub holder: KnowledgeHolderRef,
320    pub producer_correlation: Option<String>,
321    pub records: Vec<KnowledgeRecord>,
322    pub visibility: StateVisibility,
323    pub summary: String,
324}
325
326#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
327pub struct BoundaryRecord {
328    pub id: BoundaryId,
329    pub at: SimTime,
330    pub correlation_id: u64,
331    pub cadences: Vec<SystemCadence>,
332    #[serde(default, skip_serializing_if = "Vec::is_empty")]
333    pub admitted_attempts: Vec<CommandAttemptId>,
334    pub admitted_commands: Vec<CommandId>,
335    #[serde(default, skip_serializing_if = "Vec::is_empty")]
336    pub admitted_ingress: Vec<IngressId>,
337    #[serde(default, skip_serializing_if = "Vec::is_empty")]
338    pub generated_ingress: Vec<BoundaryIngressGeneration>,
339    pub admitted_events: Vec<EventId>,
340    pub reservation_offers: Vec<ReservationOfferRecord>,
341    pub reservation_requests: Vec<ReservationRequestRecord>,
342    pub allocations: Vec<ReservationAllocation>,
343    #[serde(default)]
344    pub random_draws: Vec<RandomDrawId>,
345    pub changes: Vec<BoundaryChange>,
346    #[serde(default, skip_serializing_if = "Vec::is_empty")]
347    pub record_changes: Vec<DomainRecordChange>,
348    #[serde(default, skip_serializing_if = "Vec::is_empty")]
349    pub knowledge_changes: Vec<BoundaryKnowledgeChange>,
350    #[serde(default, skip_serializing_if = "Vec::is_empty")]
351    pub maintenance_changes: Vec<crate::MaintenanceChangeRecord>,
352    #[serde(default, skip_serializing_if = "Option::is_none")]
353    pub maintenance_terminal_root: Option<String>,
354    pub emissions: Vec<BoundaryEmission>,
355    #[serde(default)]
356    /// Untagged legacy full-state hash or a `v1:` incremental state commitment.
357    pub state_hash: Option<String>,
358    #[serde(default)]
359    pub previous_hash: String,
360    #[serde(default)]
361    pub hash: String,
362}
363
364#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
365pub struct BoundaryReceipt {
366    pub boundary_id: BoundaryId,
367    pub settled_at: SimTime,
368    pub emitted_events: Vec<EventId>,
369    pub generated_ingress: Vec<IngressId>,
370    pub random_draws: Vec<RandomDrawId>,
371    pub boundary_hash: String,
372    pub change_count: usize,
373    pub record_change_count: usize,
374    pub knowledge_batch_count: usize,
375    pub knowledge_record_count: usize,
376    pub allocations: Vec<ReservationAllocation>,
377}