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/// Durable, idempotent external-delivery identity derived from committed
277/// boundary evidence. The engine creates one entry for every emission; a host
278/// may deliver it at least once and use `delivery_id` as its idempotency key.
279#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
280pub struct OutboxEntry {
281    pub delivery_id: String,
282    pub boundary: BoundaryId,
283    pub event: EventId,
284    pub emission_index: u64,
285    pub plugin: String,
286    pub system: String,
287}
288
289#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
290pub struct BoundaryIngressGeneration {
291    pub ingress: IngressId,
292    pub plugin: String,
293    pub system: String,
294    pub phase: crate::BoundaryPhase,
295    pub visibility: StateVisibility,
296}
297
298#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
299pub struct BoundaryKnowledgeChange {
300    pub plugin: String,
301    pub system: String,
302    pub phase: crate::BoundaryPhase,
303    pub holder: KnowledgeHolderRef,
304    pub producer_correlation: Option<String>,
305    pub records: Vec<KnowledgeRecord>,
306    pub visibility: StateVisibility,
307    pub summary: String,
308}
309
310#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
311pub struct BoundaryRecord {
312    pub id: BoundaryId,
313    pub at: SimTime,
314    pub correlation_id: u64,
315    pub cadences: Vec<SystemCadence>,
316    #[serde(default, skip_serializing_if = "Vec::is_empty")]
317    pub admitted_attempts: Vec<CommandAttemptId>,
318    pub admitted_commands: Vec<CommandId>,
319    #[serde(default, skip_serializing_if = "Vec::is_empty")]
320    pub admitted_ingress: Vec<IngressId>,
321    #[serde(default, skip_serializing_if = "Vec::is_empty")]
322    pub generated_ingress: Vec<BoundaryIngressGeneration>,
323    pub admitted_events: Vec<EventId>,
324    pub reservation_offers: Vec<ReservationOfferRecord>,
325    pub reservation_requests: Vec<ReservationRequestRecord>,
326    pub allocations: Vec<ReservationAllocation>,
327    #[serde(default)]
328    pub random_draws: Vec<RandomDrawId>,
329    pub changes: Vec<BoundaryChange>,
330    #[serde(default, skip_serializing_if = "Vec::is_empty")]
331    pub record_changes: Vec<DomainRecordChange>,
332    #[serde(default, skip_serializing_if = "Vec::is_empty")]
333    pub knowledge_changes: Vec<BoundaryKnowledgeChange>,
334    pub emissions: Vec<BoundaryEmission>,
335    #[serde(default)]
336    /// Untagged legacy full-state hash or a `v1:` incremental state commitment.
337    pub state_hash: Option<String>,
338    #[serde(default)]
339    pub previous_hash: String,
340    #[serde(default)]
341    pub hash: String,
342}
343
344#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
345pub struct BoundaryReceipt {
346    pub boundary_id: BoundaryId,
347    pub settled_at: SimTime,
348    pub emitted_events: Vec<EventId>,
349    pub generated_ingress: Vec<IngressId>,
350    pub random_draws: Vec<RandomDrawId>,
351    pub boundary_hash: String,
352    pub change_count: usize,
353    pub record_change_count: usize,
354    pub knowledge_batch_count: usize,
355    pub knowledge_record_count: usize,
356    pub allocations: Vec<ReservationAllocation>,
357}