Skip to main content

canwu_sim/runtime/
boundary.rs

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, EventId, IngressId, KnowledgeHolderRef, KnowledgeSchemaId, PersonId, RandomDrawId,
10};
11use canwu_knowledge::{KnowledgeRecord, KnowledgeRecordDraft};
12use canwu_time::{SimDuration, SimTime};
13use serde::{Deserialize, Serialize};
14use serde_json::Value;
15
16#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
17pub struct ReservationPoolKey {
18    pub state: StateKey,
19    pub entity: EntityRef,
20    pub resource: String,
21}
22
23impl ReservationPoolKey {
24    #[must_use]
25    pub fn new(state: StateKey, entity: EntityRef, resource: impl Into<String>) -> Self {
26        Self {
27            state,
28            entity,
29            resource: resource.into(),
30        }
31    }
32}
33
34#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
35pub struct ReservationRef {
36    pub plugin: String,
37    pub system: String,
38    pub request: String,
39}
40
41impl ReservationRef {
42    #[must_use]
43    pub fn new(
44        plugin: impl Into<String>,
45        system: impl Into<String>,
46        request: impl Into<String>,
47    ) -> Self {
48        Self {
49            plugin: plugin.into(),
50            system: system.into(),
51            request: request.into(),
52        }
53    }
54}
55
56#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
57pub struct ReservationOffer {
58    pub pool: ReservationPoolKey,
59    pub capacity: u64,
60}
61
62#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
63pub struct ReservationRequest {
64    pub request: String,
65    pub pool: ReservationPoolKey,
66    pub quantity: u64,
67    pub priority: i32,
68    pub tie_break: String,
69}
70
71#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
72pub struct ReservationOfferRecord {
73    pub plugin: String,
74    pub system: String,
75    pub offer: ReservationOffer,
76}
77
78#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
79pub struct ReservationRequestRecord {
80    pub reservation: ReservationRef,
81    pub request: ReservationRequest,
82}
83
84#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
85#[serde(rename_all = "snake_case")]
86pub enum ReservationDisposition {
87    Fulfilled,
88    Partial,
89    Rejected,
90}
91
92#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
93pub struct ReservationAllocation {
94    pub reservation: ReservationRef,
95    pub pool: ReservationPoolKey,
96    pub requested: u64,
97    pub granted: u64,
98    pub remaining_after: u64,
99    pub disposition: ReservationDisposition,
100}
101
102#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
103#[serde(tag = "type", rename_all = "snake_case")]
104pub enum BoundaryDirective {
105    SetComponent {
106        state: StateKey,
107        entity: EntityRef,
108        component: String,
109        value: Value,
110        summary: String,
111    },
112    MutateRecord {
113        mutation: DomainRecordMutation,
114        summary: String,
115    },
116    Emit {
117        event_type: String,
118        summary: String,
119        affected: Vec<EntityRef>,
120    },
121    ScheduleIngress {
122        after: SimDuration,
123        packet_type: String,
124        priority: i32,
125        payload: Value,
126        affected: Vec<EntityRef>,
127    },
128    SchedulePluginIngress {
129        target_plugin: String,
130        after: SimDuration,
131        packet_type: String,
132        priority: i32,
133        payload: Value,
134        affected: Vec<EntityRef>,
135    },
136    /// Resolves an open ticket with an operation-keyed random draw bound to
137    /// the ticket and its version, either for a controller with random policy
138    /// identity or as a guarded utility policy's random tie-break. The
139    /// boundary generates the `Resolve` decision ingress, admitted at the next
140    /// boundary.
141    ///
142    /// Before any draw is committed, the directive fails the boundary when
143    /// the ticket's person decision maker
144    /// ([`crate::ErrorCode::DecisionMakerUnavailable`]) or its assigned
145    /// controller's authority person ([`crate::ErrorCode::IssuerUnavailable`])
146    /// is unavailable in the availability committed before this boundary. The
147    /// authority person is the actor of an actor authority or the responsible
148    /// actor of an institution authority. Because the end-of-boundary sweep
149    /// and `Open` admission keep such tickets from staying open, this is a
150    /// safeguard. Any availability change made in the same boundary does not
151    /// fail the directive, because failing would roll the change back and
152    /// repeat on every retry; the draw is then committed, the end-of-boundary
153    /// sweep cancels the ticket (see
154    /// [`BoundaryDirective::SetPersonAvailability`]), and the generated
155    /// resolution is rejected at admission. To avoid that wasted draw,
156    /// tie-break and random-policy systems should skip tickets whose decision
157    /// maker or controller authority person is unavailable, read through
158    /// [`crate::SimulationView::person_availability`] (which requires
159    /// declaring `StateKey::core_person_availability()` in the contract's
160    /// reads).
161    ResolveDecisionRandomly {
162        resolution: RandomDecisionResolution,
163    },
164    PublishKnowledge {
165        holder: KnowledgeHolderRef,
166        visibility: StateVisibility,
167        producer_correlation: Option<String>,
168        records: Vec<KnowledgeRecordDraft>,
169        summary: String,
170    },
171    /// Replaces one person's core life and custody state. Accepted from a
172    /// phase-7 or phase-10 system that declares
173    /// `StateKey::core_person_availability()` as a write; two writes for the
174    /// same person in one boundary fail the boundary.
175    ///
176    /// Making a person unavailable cancels, at the end of the same boundary
177    /// and after the boundary's random decisions are materialized, every open
178    /// decision ticket whose decision maker is that person
179    /// ([`crate::DECISION_MAKER_UNAVAILABLE_REASON`]) and then every remaining
180    /// open ticket whose assigned controller's authority person is that
181    /// person ([`crate::CONTROLLER_AUTHORITY_UNAVAILABLE_REASON`]). The
182    /// authority person is the actor of an actor authority or the responsible
183    /// actor of an institution authority; council and no-responsible-actor
184    /// authorities are never affected. A ticket that qualifies for both
185    /// reasons carries the decision-maker reason. The cancelled IDs are
186    /// recorded in ticket-ID order on the boundary's
187    /// [`crate::BoundaryPersonAvailabilityChange`].
188    SetPersonAvailability {
189        person: PersonId,
190        availability: PersonAvailability,
191        summary: String,
192    },
193    /// Creates a person with an engine-allocated ID. Accepted from a phase-7
194    /// system that declares `StateKey::core_people()` as a write. The person
195    /// is committed at the end of the boundary and becomes visible to systems
196    /// at the next boundary; `correlation` is unique per plugin, system, and
197    /// boundary and binds the receipt's allocated ID.
198    CreatePerson {
199        draft: PersonDraft,
200        correlation: String,
201        summary: String,
202    },
203    /// Withdraws one still-pending plugin ingress item that this system's
204    /// plugin scheduled inside the engine (through `ScheduleIngress`,
205    /// `SchedulePluginIngress`, or a plugin command), strictly before the
206    /// item's due time. The boundary records a terminal
207    /// [`crate::IngressPayload::PluginCancellation`] entry among its generated
208    /// ingress; the withdrawn item is never admitted.
209    ///
210    /// Take targets from [`crate::SimulationView::cancellable_plugin_ingress`]
211    /// in the same boundary. A target that is foreign, already due, admitted,
212    /// or cancelled, or that another proposal already cancels in this
213    /// boundary, fails the whole boundary deterministically.
214    CancelPluginIngress {
215        ingress_id: IngressId,
216        reason: String,
217    },
218}
219
220#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
221pub struct RandomDecisionResolution {
222    pub priority: i32,
223    pub decision_request_id: DecisionRequestId,
224    #[serde(default, skip_serializing_if = "Option::is_none")]
225    pub command_request_id: Option<CommandRequestId>,
226    pub ticket_id: DecisionTicketId,
227    pub expected_version: u64,
228    pub controller_id: String,
229    pub sample: RandomSample,
230    pub option_weights: Vec<DecisionOptionWeight>,
231    /// The pending decision of a utility-policy controller whose
232    /// `PendingRandom` candidates this draw resolves. `None` for a
233    /// random-policy controller, whose weights cover every available option.
234    /// When present, `option_weights` must equal the pending candidates, and
235    /// the generated resolution keeps the pending evaluations, fired guards,
236    /// and random stage.
237    #[serde(default, skip_serializing_if = "Option::is_none")]
238    pub tie_break: Option<Box<PolicyDecision>>,
239}
240
241#[derive(Clone, Debug, Default, Deserialize, PartialEq, Serialize)]
242pub struct BoundaryProposal {
243    pub offers: Vec<ReservationOffer>,
244    pub requests: Vec<ReservationRequest>,
245    pub directives: Vec<BoundaryDirective>,
246}
247
248#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
249pub struct KnowledgeWriteGrant {
250    pub schema: KnowledgeSchemaId,
251    pub visibilities: Vec<StateVisibility>,
252}
253
254#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
255pub struct PluginIngressTarget {
256    pub target_plugin: String,
257    pub packet_type: String,
258}
259
260#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
261pub struct BoundarySystemContract {
262    pub name: String,
263    pub phase: crate::BoundaryPhase,
264    pub cadence: SystemCadence,
265    pub reads: Vec<StateKey>,
266    pub writes: Vec<StateKey>,
267    pub emits: Vec<String>,
268    pub reservation_offers: Vec<StateKey>,
269    pub reservation_requests: Vec<StateKey>,
270    pub reservation_reads: Vec<ReservationRef>,
271    #[serde(default)]
272    pub random_streams: Vec<RandomStreamKey>,
273    #[serde(default, skip_serializing_if = "Vec::is_empty")]
274    pub knowledge_writes: Vec<KnowledgeWriteGrant>,
275    #[serde(default, skip_serializing_if = "Vec::is_empty")]
276    pub plugin_ingress_targets: Vec<PluginIngressTarget>,
277    pub visibility: StateVisibility,
278}
279
280impl BoundarySystemContract {
281    #[must_use]
282    pub fn new(
283        name: impl Into<String>,
284        phase: crate::BoundaryPhase,
285        cadence: SystemCadence,
286    ) -> Self {
287        Self {
288            name: name.into(),
289            phase,
290            cadence,
291            reads: Vec::new(),
292            writes: Vec::new(),
293            emits: Vec::new(),
294            reservation_offers: Vec::new(),
295            reservation_requests: Vec::new(),
296            reservation_reads: Vec::new(),
297            random_streams: Vec::new(),
298            knowledge_writes: Vec::new(),
299            plugin_ingress_targets: Vec::new(),
300            visibility: StateVisibility::NextBoundary,
301        }
302    }
303}
304
305#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
306pub struct BoundaryContext {
307    pub boundary_id: BoundaryId,
308    pub at: SimTime,
309    pub phase: crate::BoundaryPhase,
310    pub plugin: String,
311    pub system: String,
312    pub admitted_attempts: Vec<CommandAttemptId>,
313    pub admitted_commands: Vec<CommandId>,
314    pub admitted_ingress: Vec<IngressId>,
315    pub admitted_events: Vec<EventId>,
316    pub emitted_events: Vec<EventId>,
317}
318
319pub type BoundarySystemHandler =
320    fn(&SimulationView<'_>, &BoundaryContext) -> Result<BoundaryProposal, CanwuError>;
321
322#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
323pub struct BoundaryRequest {
324    pub at: SimTime,
325    pub cadences: Vec<SystemCadence>,
326}
327
328impl BoundaryRequest {
329    #[must_use]
330    pub const fn at(at: SimTime) -> Self {
331        Self {
332            at,
333            cadences: Vec::new(),
334        }
335    }
336
337    #[must_use]
338    pub fn with_cadence(mut self, cadence: SystemCadence) -> Self {
339        self.cadences.push(cadence);
340        self
341    }
342}
343
344#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
345pub struct BoundaryChange {
346    pub plugin: String,
347    pub system: String,
348    pub state: StateKey,
349    pub entity: EntityRef,
350    pub component: String,
351    pub previous: Option<Value>,
352    pub value: Value,
353    pub visibility: StateVisibility,
354    pub summary: String,
355}
356
357#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
358#[serde(tag = "kind", rename_all = "snake_case")]
359pub enum BoundaryEmissionKind {
360    Change { change_index: u64 },
361    RecordChange { change_index: u64 },
362    KnowledgeChange { change_index: u64 },
363    Explicit,
364}
365
366#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
367pub struct BoundaryEmission {
368    pub plugin: String,
369    pub system: String,
370    pub event: EventId,
371    pub kind: BoundaryEmissionKind,
372}
373
374/// Durable, idempotent external-delivery identity derived from committed
375/// boundary evidence. The engine creates one entry for every emission; a host
376/// may deliver it at least once and use `delivery_id` as its idempotency key.
377#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
378pub struct OutboxEntry {
379    pub delivery_id: String,
380    pub boundary: BoundaryId,
381    pub event: EventId,
382    pub emission_index: u64,
383    pub plugin: String,
384    pub system: String,
385}
386
387#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
388pub struct BoundaryIngressGeneration {
389    pub ingress: IngressId,
390    pub plugin: String,
391    pub system: String,
392    pub phase: crate::BoundaryPhase,
393    pub visibility: StateVisibility,
394}
395
396#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
397pub struct BoundaryKnowledgeChange {
398    pub plugin: String,
399    pub system: String,
400    pub phase: crate::BoundaryPhase,
401    pub holder: KnowledgeHolderRef,
402    pub producer_correlation: Option<String>,
403    pub records: Vec<KnowledgeRecord>,
404    pub visibility: StateVisibility,
405    pub summary: String,
406}
407
408#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
409pub struct BoundaryRecord {
410    pub id: BoundaryId,
411    pub at: SimTime,
412    pub correlation_id: u64,
413    pub cadences: Vec<SystemCadence>,
414    #[serde(default, skip_serializing_if = "Vec::is_empty")]
415    pub admitted_attempts: Vec<CommandAttemptId>,
416    pub admitted_commands: Vec<CommandId>,
417    #[serde(default, skip_serializing_if = "Vec::is_empty")]
418    pub admitted_ingress: Vec<IngressId>,
419    #[serde(default, skip_serializing_if = "Vec::is_empty")]
420    pub generated_ingress: Vec<BoundaryIngressGeneration>,
421    pub admitted_events: Vec<EventId>,
422    pub reservation_offers: Vec<ReservationOfferRecord>,
423    pub reservation_requests: Vec<ReservationRequestRecord>,
424    pub allocations: Vec<ReservationAllocation>,
425    #[serde(default)]
426    pub random_draws: Vec<RandomDrawId>,
427    pub changes: Vec<BoundaryChange>,
428    #[serde(default, skip_serializing_if = "Vec::is_empty")]
429    pub record_changes: Vec<DomainRecordChange>,
430    #[serde(default, skip_serializing_if = "Vec::is_empty")]
431    pub knowledge_changes: Vec<BoundaryKnowledgeChange>,
432    #[serde(default, skip_serializing_if = "Vec::is_empty")]
433    pub maintenance_changes: Vec<crate::MaintenanceChangeRecord>,
434    #[serde(default, skip_serializing_if = "Option::is_none")]
435    pub maintenance_terminal_root: Option<String>,
436    #[serde(default, skip_serializing_if = "Vec::is_empty")]
437    pub person_availability_changes: Vec<BoundaryPersonAvailabilityChange>,
438    #[serde(default, skip_serializing_if = "Vec::is_empty")]
439    pub created_persons: Vec<BoundaryPersonCreation>,
440    pub emissions: Vec<BoundaryEmission>,
441    #[serde(default)]
442    /// Untagged legacy full-state hash or a `v1:` incremental state commitment.
443    pub state_hash: Option<String>,
444    #[serde(default)]
445    pub previous_hash: String,
446    #[serde(default)]
447    pub hash: String,
448}
449
450#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
451pub struct BoundaryReceipt {
452    pub boundary_id: BoundaryId,
453    pub settled_at: SimTime,
454    pub emitted_events: Vec<EventId>,
455    pub generated_ingress: Vec<IngressId>,
456    pub random_draws: Vec<RandomDrawId>,
457    pub boundary_hash: String,
458    pub change_count: usize,
459    pub record_change_count: usize,
460    pub knowledge_batch_count: usize,
461    pub knowledge_record_count: usize,
462    pub allocations: Vec<ReservationAllocation>,
463    /// Engine-allocated IDs of persons created in this boundary, keyed by
464    /// producing plugin, system, and correlation.
465    #[serde(default, skip_serializing_if = "Vec::is_empty")]
466    pub created_persons: Vec<CreatedPerson>,
467}