1use super::{
2 BoundaryPersonAvailabilityChange, BoundaryPersonCreation, CanwuError, CreatedPerson,
3 DecisionOptionWeight, DomainRecordChange, DomainRecordMutation, PersonAvailability,
4 PersonDraft, PolicyDecision, RandomSample, RandomStreamKey, SimulationView, StateKey,
5 StateVisibility, SystemCadence,
6};
7use canwu_core::{
8 BoundaryId, CommandAttemptId, CommandId, CommandRequestId, DecisionRequestId, DecisionTicketId,
9 EntityRef, EvaluationTraceRecord, EventId, IngressId, KnowledgeHolderRef, KnowledgeSchemaId,
10 PersonId, RandomDrawId,
11};
12use canwu_knowledge::{KnowledgeRecord, KnowledgeRecordDraft};
13use canwu_time::{SimDuration, SimTime};
14use serde::{Deserialize, Serialize};
15use serde_json::Value;
16
17#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
18pub struct ReservationPoolKey {
19 pub state: StateKey,
20 pub entity: EntityRef,
21 pub resource: String,
22}
23
24impl ReservationPoolKey {
25 #[must_use]
26 pub fn new(state: StateKey, entity: EntityRef, resource: impl Into<String>) -> Self {
27 Self {
28 state,
29 entity,
30 resource: resource.into(),
31 }
32 }
33}
34
35#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
36pub struct ReservationRef {
37 pub plugin: String,
38 pub system: String,
39 pub request: String,
40}
41
42impl ReservationRef {
43 #[must_use]
44 pub fn new(
45 plugin: impl Into<String>,
46 system: impl Into<String>,
47 request: impl Into<String>,
48 ) -> Self {
49 Self {
50 plugin: plugin.into(),
51 system: system.into(),
52 request: request.into(),
53 }
54 }
55}
56
57#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
58pub struct ReservationOffer {
59 pub pool: ReservationPoolKey,
60 pub capacity: u64,
61}
62
63#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
64pub struct ReservationRequest {
65 pub request: String,
66 pub pool: ReservationPoolKey,
67 pub quantity: u64,
68 pub priority: i32,
69 pub tie_break: String,
70}
71
72#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
73pub struct ReservationOfferRecord {
74 pub plugin: String,
75 pub system: String,
76 pub offer: ReservationOffer,
77}
78
79#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
80pub struct ReservationRequestRecord {
81 pub reservation: ReservationRef,
82 pub request: ReservationRequest,
83}
84
85#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
86#[serde(rename_all = "snake_case")]
87pub enum ReservationDisposition {
88 Fulfilled,
89 Partial,
90 Rejected,
91}
92
93#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
94pub struct ReservationAllocation {
95 pub reservation: ReservationRef,
96 pub pool: ReservationPoolKey,
97 pub requested: u64,
98 pub granted: u64,
99 pub remaining_after: u64,
100 pub disposition: ReservationDisposition,
101}
102
103#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
104#[serde(tag = "type", rename_all = "snake_case")]
105pub enum BoundaryDirective {
106 SetComponent {
107 state: StateKey,
108 entity: EntityRef,
109 component: String,
110 value: Value,
111 summary: String,
112 },
113 MutateRecord {
114 mutation: DomainRecordMutation,
115 summary: String,
116 },
117 Emit {
118 event_type: String,
119 summary: String,
120 affected: Vec<EntityRef>,
121 },
122 ScheduleIngress {
123 after: SimDuration,
124 packet_type: String,
125 priority: i32,
126 payload: Value,
127 affected: Vec<EntityRef>,
128 },
129 SchedulePluginIngress {
130 target_plugin: String,
131 after: SimDuration,
132 packet_type: String,
133 priority: i32,
134 payload: Value,
135 affected: Vec<EntityRef>,
136 },
137 ResolveDecisionRandomly {
163 resolution: RandomDecisionResolution,
164 },
165 PublishKnowledge {
166 holder: KnowledgeHolderRef,
167 visibility: StateVisibility,
168 producer_correlation: Option<String>,
169 records: Vec<KnowledgeRecordDraft>,
170 summary: String,
171 },
172 SetPersonAvailability {
190 person: PersonId,
191 availability: PersonAvailability,
192 summary: String,
193 },
194 CreatePerson {
200 draft: PersonDraft,
201 correlation: String,
202 summary: String,
203 },
204 CancelPluginIngress {
216 ingress_id: IngressId,
217 reason: String,
218 },
219 RecordEvaluationTrace { trace: EvaluationTraceRecord },
234 RegisterTransitionManifest { manifest: crate::TransitionManifest },
252 StageTransitionWrite {
271 manifest_id: crate::TransitionManifestId,
272 writes: Vec<BoundaryDirective>,
273 },
274}
275
276#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
277pub struct RandomDecisionResolution {
278 pub priority: i32,
279 pub decision_request_id: DecisionRequestId,
280 #[serde(default, skip_serializing_if = "Option::is_none")]
281 pub command_request_id: Option<CommandRequestId>,
282 pub ticket_id: DecisionTicketId,
283 pub expected_version: u64,
284 pub controller_id: String,
285 pub sample: RandomSample,
286 pub option_weights: Vec<DecisionOptionWeight>,
287 #[serde(default, skip_serializing_if = "Option::is_none")]
294 pub tie_break: Option<Box<PolicyDecision>>,
295}
296
297#[derive(Clone, Debug, Default, Deserialize, PartialEq, Serialize)]
298pub struct BoundaryProposal {
299 pub offers: Vec<ReservationOffer>,
300 pub requests: Vec<ReservationRequest>,
301 pub directives: Vec<BoundaryDirective>,
302}
303
304#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
305pub struct KnowledgeWriteGrant {
306 pub schema: KnowledgeSchemaId,
307 pub visibilities: Vec<StateVisibility>,
308}
309
310#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
311pub struct PluginIngressTarget {
312 pub target_plugin: String,
313 pub packet_type: String,
314}
315
316#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
317pub struct BoundarySystemContract {
318 pub name: String,
319 pub phase: crate::BoundaryPhase,
320 pub cadence: SystemCadence,
321 pub reads: Vec<StateKey>,
322 pub writes: Vec<StateKey>,
323 pub emits: Vec<String>,
324 pub reservation_offers: Vec<StateKey>,
325 pub reservation_requests: Vec<StateKey>,
326 pub reservation_reads: Vec<ReservationRef>,
327 #[serde(default)]
328 pub random_streams: Vec<RandomStreamKey>,
329 #[serde(default, skip_serializing_if = "Vec::is_empty")]
330 pub knowledge_writes: Vec<KnowledgeWriteGrant>,
331 #[serde(default, skip_serializing_if = "Vec::is_empty")]
332 pub plugin_ingress_targets: Vec<PluginIngressTarget>,
333 pub visibility: StateVisibility,
334}
335
336impl BoundarySystemContract {
337 #[must_use]
338 pub fn new(
339 name: impl Into<String>,
340 phase: crate::BoundaryPhase,
341 cadence: SystemCadence,
342 ) -> Self {
343 Self {
344 name: name.into(),
345 phase,
346 cadence,
347 reads: Vec::new(),
348 writes: Vec::new(),
349 emits: Vec::new(),
350 reservation_offers: Vec::new(),
351 reservation_requests: Vec::new(),
352 reservation_reads: Vec::new(),
353 random_streams: Vec::new(),
354 knowledge_writes: Vec::new(),
355 plugin_ingress_targets: Vec::new(),
356 visibility: StateVisibility::NextBoundary,
357 }
358 }
359}
360
361#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
362pub struct BoundaryContext {
363 pub boundary_id: BoundaryId,
364 pub at: SimTime,
365 pub phase: crate::BoundaryPhase,
366 pub plugin: String,
367 pub system: String,
368 pub admitted_attempts: Vec<CommandAttemptId>,
369 pub admitted_commands: Vec<CommandId>,
370 pub admitted_ingress: Vec<IngressId>,
371 pub admitted_events: Vec<EventId>,
372 pub emitted_events: Vec<EventId>,
373}
374
375pub type BoundarySystemHandler =
376 fn(&SimulationView<'_>, &BoundaryContext) -> Result<BoundaryProposal, CanwuError>;
377
378#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
379pub struct BoundaryRequest {
380 pub at: SimTime,
381 pub cadences: Vec<SystemCadence>,
382}
383
384impl BoundaryRequest {
385 #[must_use]
386 pub const fn at(at: SimTime) -> Self {
387 Self {
388 at,
389 cadences: Vec::new(),
390 }
391 }
392
393 #[must_use]
394 pub fn with_cadence(mut self, cadence: SystemCadence) -> Self {
395 self.cadences.push(cadence);
396 self
397 }
398}
399
400#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
401pub struct BoundaryChange {
402 pub plugin: String,
403 pub system: String,
404 pub state: StateKey,
405 pub entity: EntityRef,
406 pub component: String,
407 pub previous: Option<Value>,
408 pub value: Value,
409 pub visibility: StateVisibility,
410 pub summary: String,
411}
412
413#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
414#[serde(tag = "kind", rename_all = "snake_case")]
415pub enum BoundaryEmissionKind {
416 Change { change_index: u64 },
417 RecordChange { change_index: u64 },
418 KnowledgeChange { change_index: u64 },
419 Explicit,
420}
421
422#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
423pub struct BoundaryEmission {
424 pub plugin: String,
425 pub system: String,
426 pub event: EventId,
427 pub kind: BoundaryEmissionKind,
428}
429
430#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
434pub struct OutboxEntry {
435 pub delivery_id: String,
436 pub boundary: BoundaryId,
437 pub event: EventId,
438 pub emission_index: u64,
439 pub plugin: String,
440 pub system: String,
441}
442
443#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
444pub struct BoundaryIngressGeneration {
445 pub ingress: IngressId,
446 pub plugin: String,
447 pub system: String,
448 pub phase: crate::BoundaryPhase,
449 pub visibility: StateVisibility,
450}
451
452#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
453pub struct BoundaryKnowledgeChange {
454 pub plugin: String,
455 pub system: String,
456 pub phase: crate::BoundaryPhase,
457 pub holder: KnowledgeHolderRef,
458 pub producer_correlation: Option<String>,
459 pub records: Vec<KnowledgeRecord>,
460 pub visibility: StateVisibility,
461 pub summary: String,
462}
463
464#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
465pub struct BoundaryRecord {
466 pub id: BoundaryId,
467 pub at: SimTime,
468 pub correlation_id: u64,
469 pub cadences: Vec<SystemCadence>,
470 #[serde(default, skip_serializing_if = "Vec::is_empty")]
471 pub admitted_attempts: Vec<CommandAttemptId>,
472 pub admitted_commands: Vec<CommandId>,
473 #[serde(default, skip_serializing_if = "Vec::is_empty")]
474 pub admitted_ingress: Vec<IngressId>,
475 #[serde(default, skip_serializing_if = "Vec::is_empty")]
476 pub generated_ingress: Vec<BoundaryIngressGeneration>,
477 pub admitted_events: Vec<EventId>,
478 pub reservation_offers: Vec<ReservationOfferRecord>,
479 pub reservation_requests: Vec<ReservationRequestRecord>,
480 pub allocations: Vec<ReservationAllocation>,
481 #[serde(default)]
482 pub random_draws: Vec<RandomDrawId>,
483 pub changes: Vec<BoundaryChange>,
484 #[serde(default, skip_serializing_if = "Vec::is_empty")]
485 pub record_changes: Vec<DomainRecordChange>,
486 #[serde(default, skip_serializing_if = "Vec::is_empty")]
487 pub knowledge_changes: Vec<BoundaryKnowledgeChange>,
488 #[serde(default, skip_serializing_if = "Vec::is_empty")]
489 pub maintenance_changes: Vec<crate::MaintenanceChangeRecord>,
490 #[serde(default, skip_serializing_if = "Option::is_none")]
491 pub maintenance_terminal_root: Option<String>,
492 #[serde(default, skip_serializing_if = "Vec::is_empty")]
493 pub person_availability_changes: Vec<BoundaryPersonAvailabilityChange>,
494 #[serde(default, skip_serializing_if = "Vec::is_empty")]
495 pub created_persons: Vec<BoundaryPersonCreation>,
496 #[serde(default, skip_serializing_if = "Vec::is_empty")]
499 pub evaluation_traces: Vec<crate::BoundaryEvaluationTrace>,
500 #[serde(default, skip_serializing_if = "Vec::is_empty")]
502 pub transition_manifests: Vec<crate::PendingTransitionManifest>,
503 #[serde(default, skip_serializing_if = "Vec::is_empty")]
505 pub transition_audits: Vec<crate::TransitionAuditRecord>,
506 pub emissions: Vec<BoundaryEmission>,
507 #[serde(default)]
508 pub state_hash: Option<String>,
510 #[serde(default)]
511 pub previous_hash: String,
512 #[serde(default)]
513 pub hash: String,
514}
515
516#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
517pub struct BoundaryReceipt {
518 pub boundary_id: BoundaryId,
519 pub settled_at: SimTime,
520 pub emitted_events: Vec<EventId>,
521 pub generated_ingress: Vec<IngressId>,
522 pub random_draws: Vec<RandomDrawId>,
523 pub boundary_hash: String,
524 pub change_count: usize,
525 pub record_change_count: usize,
526 pub knowledge_batch_count: usize,
527 pub knowledge_record_count: usize,
528 pub allocations: Vec<ReservationAllocation>,
529 #[serde(default, skip_serializing_if = "Vec::is_empty")]
532 pub created_persons: Vec<CreatedPerson>,
533 #[serde(default, skip_serializing_if = "Vec::is_empty")]
535 pub transition_audits: Vec<crate::TransitionAuditRecord>,
536}