Skip to main content

canwu_sim/runtime/
boundary.rs

1use super::{
2    CanwuError, DomainRecordChange, DomainRecordMutation, RandomStreamKey, SimulationView,
3    StateKey, StateVisibility, SystemCadence,
4};
5use canwu_core::{
6    BoundaryId, CommandAttemptId, CommandId, EntityRef, EventId, IngressId, KnowledgeHolderRef,
7    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    PublishKnowledge {
135        holder: KnowledgeHolderRef,
136        visibility: StateVisibility,
137        producer_correlation: Option<String>,
138        records: Vec<KnowledgeRecordDraft>,
139        summary: String,
140    },
141}
142
143#[derive(Clone, Debug, Default, Deserialize, PartialEq, Serialize)]
144pub struct BoundaryProposal {
145    pub offers: Vec<ReservationOffer>,
146    pub requests: Vec<ReservationRequest>,
147    pub directives: Vec<BoundaryDirective>,
148}
149
150#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
151pub struct KnowledgeWriteGrant {
152    pub schema: KnowledgeSchemaId,
153    pub visibilities: Vec<StateVisibility>,
154}
155
156#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
157pub struct PluginIngressTarget {
158    pub target_plugin: String,
159    pub packet_type: String,
160}
161
162#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
163pub struct BoundarySystemContract {
164    pub name: String,
165    pub phase: crate::BoundaryPhase,
166    pub cadence: SystemCadence,
167    pub reads: Vec<StateKey>,
168    pub writes: Vec<StateKey>,
169    pub emits: Vec<String>,
170    pub reservation_offers: Vec<StateKey>,
171    pub reservation_requests: Vec<StateKey>,
172    pub reservation_reads: Vec<ReservationRef>,
173    #[serde(default)]
174    pub random_streams: Vec<RandomStreamKey>,
175    #[serde(default, skip_serializing_if = "Vec::is_empty")]
176    pub knowledge_writes: Vec<KnowledgeWriteGrant>,
177    #[serde(default, skip_serializing_if = "Vec::is_empty")]
178    pub plugin_ingress_targets: Vec<PluginIngressTarget>,
179    pub visibility: StateVisibility,
180}
181
182impl BoundarySystemContract {
183    #[must_use]
184    pub fn new(
185        name: impl Into<String>,
186        phase: crate::BoundaryPhase,
187        cadence: SystemCadence,
188    ) -> Self {
189        Self {
190            name: name.into(),
191            phase,
192            cadence,
193            reads: Vec::new(),
194            writes: Vec::new(),
195            emits: Vec::new(),
196            reservation_offers: Vec::new(),
197            reservation_requests: Vec::new(),
198            reservation_reads: Vec::new(),
199            random_streams: Vec::new(),
200            knowledge_writes: Vec::new(),
201            plugin_ingress_targets: Vec::new(),
202            visibility: StateVisibility::NextBoundary,
203        }
204    }
205}
206
207#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
208pub struct BoundaryContext {
209    pub boundary_id: BoundaryId,
210    pub at: SimTime,
211    pub phase: crate::BoundaryPhase,
212    pub plugin: String,
213    pub system: String,
214    pub admitted_attempts: Vec<CommandAttemptId>,
215    pub admitted_commands: Vec<CommandId>,
216    pub admitted_ingress: Vec<IngressId>,
217    pub admitted_events: Vec<EventId>,
218    pub emitted_events: Vec<EventId>,
219}
220
221pub type BoundarySystemHandler =
222    fn(&SimulationView<'_>, &BoundaryContext) -> Result<BoundaryProposal, CanwuError>;
223
224#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
225pub struct BoundaryRequest {
226    pub at: SimTime,
227    pub cadences: Vec<SystemCadence>,
228}
229
230impl BoundaryRequest {
231    #[must_use]
232    pub const fn at(at: SimTime) -> Self {
233        Self {
234            at,
235            cadences: Vec::new(),
236        }
237    }
238
239    #[must_use]
240    pub fn with_cadence(mut self, cadence: SystemCadence) -> Self {
241        self.cadences.push(cadence);
242        self
243    }
244}
245
246#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
247pub struct BoundaryChange {
248    pub plugin: String,
249    pub system: String,
250    pub state: StateKey,
251    pub entity: EntityRef,
252    pub component: String,
253    pub previous: Option<Value>,
254    pub value: Value,
255    pub visibility: StateVisibility,
256    pub summary: String,
257}
258
259#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
260#[serde(tag = "kind", rename_all = "snake_case")]
261pub enum BoundaryEmissionKind {
262    Change { change_index: u64 },
263    RecordChange { change_index: u64 },
264    KnowledgeChange { change_index: u64 },
265    Explicit,
266}
267
268#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
269pub struct BoundaryEmission {
270    pub plugin: String,
271    pub system: String,
272    pub event: EventId,
273    pub kind: BoundaryEmissionKind,
274}
275
276#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
277pub struct BoundaryIngressGeneration {
278    pub ingress: IngressId,
279    pub plugin: String,
280    pub system: String,
281    pub phase: crate::BoundaryPhase,
282    pub visibility: StateVisibility,
283}
284
285#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
286pub struct BoundaryKnowledgeChange {
287    pub plugin: String,
288    pub system: String,
289    pub phase: crate::BoundaryPhase,
290    pub holder: KnowledgeHolderRef,
291    pub producer_correlation: Option<String>,
292    pub records: Vec<KnowledgeRecord>,
293    pub visibility: StateVisibility,
294    pub summary: String,
295}
296
297#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
298pub struct BoundaryRecord {
299    pub id: BoundaryId,
300    pub at: SimTime,
301    pub correlation_id: u64,
302    pub cadences: Vec<SystemCadence>,
303    #[serde(default, skip_serializing_if = "Vec::is_empty")]
304    pub admitted_attempts: Vec<CommandAttemptId>,
305    pub admitted_commands: Vec<CommandId>,
306    #[serde(default, skip_serializing_if = "Vec::is_empty")]
307    pub admitted_ingress: Vec<IngressId>,
308    #[serde(default, skip_serializing_if = "Vec::is_empty")]
309    pub generated_ingress: Vec<BoundaryIngressGeneration>,
310    pub admitted_events: Vec<EventId>,
311    pub reservation_offers: Vec<ReservationOfferRecord>,
312    pub reservation_requests: Vec<ReservationRequestRecord>,
313    pub allocations: Vec<ReservationAllocation>,
314    #[serde(default)]
315    pub random_draws: Vec<RandomDrawId>,
316    pub changes: Vec<BoundaryChange>,
317    #[serde(default, skip_serializing_if = "Vec::is_empty")]
318    pub record_changes: Vec<DomainRecordChange>,
319    #[serde(default, skip_serializing_if = "Vec::is_empty")]
320    pub knowledge_changes: Vec<BoundaryKnowledgeChange>,
321    pub emissions: Vec<BoundaryEmission>,
322    #[serde(default)]
323    /// Untagged legacy full-state hash or a `v1:` incremental state commitment.
324    pub state_hash: Option<String>,
325    #[serde(default)]
326    pub previous_hash: String,
327    #[serde(default)]
328    pub hash: String,
329}
330
331#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
332pub struct BoundaryReceipt {
333    pub boundary_id: BoundaryId,
334    pub settled_at: SimTime,
335    pub emitted_events: Vec<EventId>,
336    pub generated_ingress: Vec<IngressId>,
337    pub random_draws: Vec<RandomDrawId>,
338    pub boundary_hash: String,
339    pub change_count: usize,
340    pub record_change_count: usize,
341    pub knowledge_batch_count: usize,
342    pub knowledge_record_count: usize,
343    pub allocations: Vec<ReservationAllocation>,
344}