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#[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 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}