Skip to main content

governance_transition/
governance_transition.rs

1#![allow(clippy::too_many_lines, clippy::unnecessary_wraps)]
2
3//! A public-API-only governance case.
4//!
5//! The case models a Southern Ming relief order without putting court,
6//! treasury, county, or granary vocabulary into Canwu itself:
7//!
8//! 1. The central relief office publishes a manifest.
9//! 2. The manifest schedules one request for the treasury and one for a county.
10//! 3. Both owners prepare their own records in the same boundary.
11//! 4. A read-only central audit commits the order only when both records match
12//!    the manifest.
13//!
14//! A zero-delay plugin ingress is intentionally deferred to the next boundary.
15//! The final snapshot is then restored and replayed to demonstrate that the
16//! governance transition is deterministic and persistence-safe.
17
18use canwu_api::{
19    BoundaryContext, BoundaryDirective, BoundaryPhase, BoundaryProposal, BoundaryRequest,
20    BoundarySystemContract, Canwu, CanwuError, DomainEntityKindClass, DomainRecordDraft,
21    DomainRecordMutation, DomainRecordSchema, DomainRecordType, DomainValueKindClass, ErrorCode,
22    IngressClass, IngressPayload, PayloadProperty, PayloadSchema, PayloadValueType,
23    PluginIngressDescriptor, PluginRegistrar, Scenario, SimDuration, SimTime, SimulationPlugin,
24    SimulationView, StateKey, StateVisibility, SystemCadence, TypedDomainRecordRef, canonical_hash,
25};
26use serde::{Deserialize, Serialize};
27use serde_json::{Value, json};
28use std::collections::BTreeMap;
29
30const CENTRAL_PLUGIN: &str = "case-relief-central";
31const TREASURY_PLUGIN: &str = "case-relief-treasury";
32const COUNTY_PLUGIN: &str = "case-relief-county";
33const ISSUE_INGRESS: &str = "issue-relief-order";
34const EXECUTE_INGRESS: &str = "execute-relief-order";
35const ORDER_ID: &str = "relief-order-1646";
36const TREASURY_RECORD_ID: &str = "relief-order-1646-treasury";
37const COUNTY_RECORD_ID: &str = "relief-order-1646-county";
38const ACTION_HASH_DOMAIN: &str = "case.relief-action.v1";
39
40#[derive(Debug, Deserialize, PartialEq, Serialize)]
41struct ReliefOrder {
42    order_id: String,
43    issued_by: String,
44    treasury_system: String,
45    treasury_version: u64,
46    treasury_disposition: String,
47    treasury_hash: String,
48    county_system: String,
49    county_version: u64,
50    county_disposition: String,
51    county_hash: String,
52}
53
54struct ReliefOrderRecord;
55
56impl DomainRecordType for ReliefOrderRecord {
57    type Payload = ReliefOrder;
58    type Class = DomainValueKindClass;
59
60    const NAMESPACE: &'static str = "case.relief";
61    const NAME: &'static str = "order";
62}
63
64#[derive(Debug, Deserialize, PartialEq, Serialize)]
65struct ReliefAction {
66    status: String,
67    grain_units: u64,
68}
69
70struct TreasuryActionRecord;
71
72impl DomainRecordType for TreasuryActionRecord {
73    type Payload = ReliefAction;
74    type Class = DomainEntityKindClass;
75
76    const NAMESPACE: &'static str = "case.relief.treasury";
77    const NAME: &'static str = "action";
78}
79
80struct CountyActionRecord;
81
82impl DomainRecordType for CountyActionRecord {
83    type Payload = ReliefAction;
84    type Class = DomainEntityKindClass;
85
86    const NAMESPACE: &'static str = "case.relief.county";
87    const NAME: &'static str = "action";
88}
89
90fn order_reference() -> TypedDomainRecordRef<ReliefOrderRecord> {
91    TypedDomainRecordRef::new(ORDER_ID)
92}
93
94fn treasury_reference() -> TypedDomainRecordRef<TreasuryActionRecord> {
95    TypedDomainRecordRef::new(TREASURY_RECORD_ID)
96}
97
98fn county_reference() -> TypedDomainRecordRef<CountyActionRecord> {
99    TypedDomainRecordRef::new(COUNTY_RECORD_ID)
100}
101
102fn object_schema(fields: &[(&str, PayloadValueType)]) -> PayloadSchema {
103    PayloadSchema::Object {
104        properties: fields
105            .iter()
106            .map(|(name, value_type)| {
107                (
108                    (*name).to_owned(),
109                    PayloadProperty {
110                        value_type: value_type.clone(),
111                        required: true,
112                    },
113                )
114            })
115            .collect::<BTreeMap<_, _>>(),
116        allow_additional: false,
117    }
118}
119
120fn order_state() -> StateKey {
121    DomainRecordSchema::for_record::<ReliefOrderRecord>().state_key()
122}
123
124fn treasury_state() -> StateKey {
125    DomainRecordSchema::for_entity::<TreasuryActionRecord>().state_key()
126}
127
128fn county_state() -> StateKey {
129    DomainRecordSchema::for_entity::<CountyActionRecord>().state_key()
130}
131
132fn action_payload(grain_units: u64) -> ReliefAction {
133    ReliefAction {
134        status: "committed".to_owned(),
135        grain_units,
136    }
137}
138
139fn action_hash(grain_units: u64) -> Result<String, CanwuError> {
140    let payload = serde_json::to_value(action_payload(grain_units)).map_err(|error| {
141        CanwuError::new(
142            ErrorCode::InvalidPayload,
143            format!("relief action payload cannot be encoded: {error}"),
144        )
145    })?;
146    canonical_hash(ACTION_HASH_DOMAIN, &payload)
147}
148
149fn issue_descriptor() -> PluginIngressDescriptor {
150    PluginIngressDescriptor {
151        name: ISSUE_INGRESS.to_owned(),
152        description: "Issue a central relief order".to_owned(),
153        class: IngressClass::ScheduledSystem,
154        payload_schema: object_schema(&[("order_id", PayloadValueType::String)]),
155    }
156}
157
158fn execute_descriptor() -> PluginIngressDescriptor {
159    PluginIngressDescriptor {
160        name: EXECUTE_INGRESS.to_owned(),
161        description: "Admit an owner execution of a relief order".to_owned(),
162        class: IngressClass::ScheduledSystem,
163        payload_schema: object_schema(&[("order_id", PayloadValueType::String)]),
164    }
165}
166
167fn owned_order_ingress(
168    view: &SimulationView<'_>,
169    context: &BoundaryContext,
170    owner: &str,
171    packet_type: &str,
172) -> Result<Option<String>, CanwuError> {
173    let mut order_id = None;
174    for ingress_id in &context.admitted_ingress {
175        let Some(record) = view.ingress(*ingress_id)? else {
176            continue;
177        };
178        let IngressPayload::Plugin {
179            plugin,
180            packet_type: admitted_packet_type,
181            payload,
182            ..
183        } = &record.payload
184        else {
185            continue;
186        };
187        if plugin != owner || admitted_packet_type != packet_type {
188            continue;
189        }
190        if order_id.is_some() {
191            return Err(CanwuError::new(
192                ErrorCode::InvalidBoundary,
193                format!("{owner} received duplicate {packet_type} ingress"),
194            ));
195        }
196        let value = payload
197            .get("order_id")
198            .and_then(Value::as_str)
199            .ok_or_else(|| {
200                CanwuError::new(
201                    ErrorCode::InvalidPayload,
202                    format!("{owner}.{packet_type} is missing order_id"),
203                )
204            })?;
205        order_id = Some(value.to_owned());
206    }
207    Ok(order_id)
208}
209
210fn order_payload(view: &SimulationView<'_>) -> Result<ReliefOrder, CanwuError> {
211    view.typed_domain_record(&order_reference())?
212        .ok_or_else(|| CanwuError::new(ErrorCode::InvalidBoundary, "relief order is missing"))?
213        .decode_payload::<ReliefOrderRecord>()
214}
215
216fn validate_order(
217    order: &ReliefOrder,
218    expected_system: &str,
219    expected_plugin: &str,
220    expected_hash: &str,
221) -> Result<(), CanwuError> {
222    let valid = order.order_id == ORDER_ID
223        && order.issued_by == CENTRAL_PLUGIN
224        && order.treasury_system == "prepare-treasury"
225        && order.treasury_version == 1
226        && order.treasury_disposition == "committed"
227        && order.county_system == "prepare-county"
228        && order.county_version == 1
229        && order.county_disposition == "committed"
230        && ((expected_plugin == TREASURY_PLUGIN
231            && expected_system == "prepare-treasury"
232            && order.treasury_hash == expected_hash)
233            || (expected_plugin == COUNTY_PLUGIN
234                && expected_system == "prepare-county"
235                && order.county_hash == expected_hash));
236    if !valid {
237        return Err(CanwuError::new(
238            ErrorCode::InvalidBoundary,
239            format!("{expected_plugin} received an invalid relief manifest"),
240        ));
241    }
242    Ok(())
243}
244
245fn publish_order(
246    view: &SimulationView<'_>,
247    context: &BoundaryContext,
248) -> Result<BoundaryProposal, CanwuError> {
249    let Some(order_id) = owned_order_ingress(view, context, CENTRAL_PLUGIN, ISSUE_INGRESS)? else {
250        return Ok(BoundaryProposal::default());
251    };
252    if order_id != ORDER_ID {
253        return Err(CanwuError::new(
254            ErrorCode::InvalidBoundary,
255            "the issue ingress names an unexpected relief order",
256        ));
257    }
258    if view.typed_domain_record(&order_reference())?.is_some() {
259        return Err(CanwuError::new(
260            ErrorCode::InvalidBoundary,
261            "the relief order already exists; duplicate issue ingress is invalid",
262        ));
263    }
264    let treasury_hash = action_hash(600)?;
265    let county_hash = action_hash(600)?;
266    let manifest = DomainRecordDraft::from_typed(
267        order_reference(),
268        &ReliefOrder {
269            order_id: ORDER_ID.to_owned(),
270            issued_by: CENTRAL_PLUGIN.to_owned(),
271            treasury_system: "prepare-treasury".to_owned(),
272            treasury_version: 1,
273            treasury_disposition: "committed".to_owned(),
274            treasury_hash,
275            county_system: "prepare-county".to_owned(),
276            county_version: 1,
277            county_disposition: "committed".to_owned(),
278            county_hash,
279        },
280    )?;
281    Ok(BoundaryProposal {
282        directives: vec![
283            BoundaryDirective::MutateRecord {
284                mutation: DomainRecordMutation::Create { record: manifest },
285                summary: "Publish the central relief order manifest".to_owned(),
286            },
287            BoundaryDirective::SchedulePluginIngress {
288                target_plugin: TREASURY_PLUGIN.to_owned(),
289                after: SimDuration::ZERO,
290                packet_type: EXECUTE_INGRESS.to_owned(),
291                priority: 0,
292                payload: json!({"order_id": ORDER_ID}),
293                affected: Vec::new(),
294            },
295            BoundaryDirective::SchedulePluginIngress {
296                target_plugin: COUNTY_PLUGIN.to_owned(),
297                after: SimDuration::ZERO,
298                packet_type: EXECUTE_INGRESS.to_owned(),
299                priority: 0,
300                payload: json!({"order_id": ORDER_ID}),
301                affected: Vec::new(),
302            },
303        ],
304        ..BoundaryProposal::default()
305    })
306}
307
308fn prepare_treasury(
309    view: &SimulationView<'_>,
310    context: &BoundaryContext,
311) -> Result<BoundaryProposal, CanwuError> {
312    prepare_owner(
313        view,
314        context,
315        TREASURY_PLUGIN,
316        "prepare-treasury",
317        treasury_reference(),
318        &action_payload(600),
319    )
320}
321
322fn prepare_county(
323    view: &SimulationView<'_>,
324    context: &BoundaryContext,
325) -> Result<BoundaryProposal, CanwuError> {
326    prepare_owner(
327        view,
328        context,
329        COUNTY_PLUGIN,
330        "prepare-county",
331        county_reference(),
332        &action_payload(600),
333    )
334}
335
336fn prepare_owner<T: DomainRecordType<Payload = ReliefAction>>(
337    view: &SimulationView<'_>,
338    context: &BoundaryContext,
339    plugin: &str,
340    system: &str,
341    reference: TypedDomainRecordRef<T>,
342    payload: &ReliefAction,
343) -> Result<BoundaryProposal, CanwuError> {
344    let Some(order_id) = owned_order_ingress(view, context, plugin, EXECUTE_INGRESS)? else {
345        return Ok(BoundaryProposal::default());
346    };
347    if order_id != ORDER_ID {
348        return Err(CanwuError::new(
349            ErrorCode::InvalidBoundary,
350            format!("{plugin} received an ingress for the wrong order"),
351        ));
352    }
353    let order = order_payload(view)?;
354    let hash = action_hash(payload.grain_units)?;
355    validate_order(&order, system, plugin, &hash)?;
356    let record = DomainRecordDraft::from_typed(reference, payload)?;
357    Ok(BoundaryProposal {
358        directives: vec![BoundaryDirective::MutateRecord {
359            mutation: DomainRecordMutation::Create { record },
360            summary: format!("{plugin} prepares its guarded relief disposition"),
361        }],
362        ..BoundaryProposal::default()
363    })
364}
365
366fn audit_order(
367    view: &SimulationView<'_>,
368    context: &BoundaryContext,
369) -> Result<BoundaryProposal, CanwuError> {
370    if owned_order_ingress(view, context, CENTRAL_PLUGIN, ISSUE_INGRESS)?.is_some() {
371        return Ok(BoundaryProposal::default());
372    }
373    let order = order_payload(view)?;
374    for (label, reference, expected_plugin, expected_system) in [
375        (
376            "treasury",
377            treasury_reference().as_untyped(),
378            TREASURY_PLUGIN,
379            "prepare-treasury",
380        ),
381        (
382            "county",
383            county_reference().as_untyped(),
384            COUNTY_PLUGIN,
385            "prepare-county",
386        ),
387    ] {
388        let record = view.domain_record(reference)?.ok_or_else(|| {
389            CanwuError::new(
390                ErrorCode::InvalidBoundary,
391                format!("central audit is missing the {label} relief disposition"),
392            )
393        })?;
394        let expected_hash = action_hash(600)?;
395        validate_order(&order, expected_system, expected_plugin, &expected_hash)?;
396        let actual_hash = canonical_hash(ACTION_HASH_DOMAIN, &record.payload)?;
397        if record.version != 1
398            || !record.is_active()
399            || record.owner != expected_plugin
400            || actual_hash != expected_hash
401            || record.payload != json!({"status": "committed", "grain_units": 600})
402        {
403            return Err(CanwuError::new(
404                ErrorCode::InvalidBoundary,
405                format!("central audit found an invalid {label} disposition"),
406            ));
407        }
408    }
409    Ok(BoundaryProposal::default())
410}
411
412struct CentralPlugin;
413
414impl SimulationPlugin for CentralPlugin {
415    fn name(&self) -> &'static str {
416        CENTRAL_PLUGIN
417    }
418
419    fn version(&self) -> &'static str {
420        "1"
421    }
422
423    fn semantic_hash(&self) -> &'static str {
424        "0000000000000000000000000000000000000000000000000000000000000301"
425    }
426
427    fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
428        let mut schema = DomainRecordSchema::for_record::<ReliefOrderRecord>();
429        schema.payload_schema = object_schema(&[
430            ("order_id", PayloadValueType::String),
431            ("issued_by", PayloadValueType::String),
432            ("treasury_system", PayloadValueType::String),
433            ("treasury_version", PayloadValueType::Integer),
434            ("treasury_disposition", PayloadValueType::String),
435            ("treasury_hash", PayloadValueType::String),
436            ("county_system", PayloadValueType::String),
437            ("county_version", PayloadValueType::Integer),
438            ("county_disposition", PayloadValueType::String),
439            ("county_hash", PayloadValueType::String),
440        ]);
441        registrar.register_record_schema(schema)?;
442        registrar.register_ingress(issue_descriptor())?;
443
444        let mut publish = BoundarySystemContract::new(
445            "publish-order",
446            BoundaryPhase::DomainDeltaProposal,
447            SystemCadence::EventDriven,
448        );
449        publish.reads = vec![order_state(), StateKey::core_ingress()];
450        publish.writes = vec![order_state()];
451        publish.plugin_ingress_targets = vec![
452            canwu_api::PluginIngressTarget {
453                target_plugin: TREASURY_PLUGIN.to_owned(),
454                packet_type: EXECUTE_INGRESS.to_owned(),
455            },
456            canwu_api::PluginIngressTarget {
457                target_plugin: COUNTY_PLUGIN.to_owned(),
458                packet_type: EXECUTE_INGRESS.to_owned(),
459            },
460        ];
461        publish.visibility = StateVisibility::SameBoundary;
462        registrar.register_boundary_system(publish, publish_order)?;
463
464        let mut audit = BoundarySystemContract::new(
465            "audit-order",
466            BoundaryPhase::StrategicAggregation,
467            SystemCadence::EventDriven,
468        );
469        audit.reads = vec![
470            order_state(),
471            treasury_state(),
472            county_state(),
473            StateKey::core_ingress(),
474        ];
475        registrar.register_boundary_system(audit, audit_order)
476    }
477}
478
479struct TreasuryPlugin;
480
481impl SimulationPlugin for TreasuryPlugin {
482    fn name(&self) -> &'static str {
483        TREASURY_PLUGIN
484    }
485
486    fn version(&self) -> &'static str {
487        "1"
488    }
489
490    fn semantic_hash(&self) -> &'static str {
491        "0000000000000000000000000000000000000000000000000000000000000302"
492    }
493
494    fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
495        let mut schema = DomainRecordSchema::for_entity::<TreasuryActionRecord>();
496        schema.payload_schema = object_schema(&[
497            ("status", PayloadValueType::String),
498            ("grain_units", PayloadValueType::Integer),
499        ]);
500        registrar.register_record_schema(schema)?;
501        registrar.register_ingress(execute_descriptor())?;
502        let mut prepare = BoundarySystemContract::new(
503            "prepare-treasury",
504            BoundaryPhase::HistoricalCandidateEvaluation,
505            SystemCadence::EventDriven,
506        );
507        prepare.reads = vec![order_state(), StateKey::core_ingress()];
508        prepare.writes = vec![treasury_state()];
509        prepare.visibility = StateVisibility::SameBoundary;
510        registrar.register_boundary_system(prepare, prepare_treasury)
511    }
512}
513
514struct CountyPlugin;
515
516impl SimulationPlugin for CountyPlugin {
517    fn name(&self) -> &'static str {
518        COUNTY_PLUGIN
519    }
520
521    fn version(&self) -> &'static str {
522        "1"
523    }
524
525    fn semantic_hash(&self) -> &'static str {
526        "0000000000000000000000000000000000000000000000000000000000000303"
527    }
528
529    fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
530        let mut schema = DomainRecordSchema::for_entity::<CountyActionRecord>();
531        schema.payload_schema = object_schema(&[
532            ("status", PayloadValueType::String),
533            ("grain_units", PayloadValueType::Integer),
534        ]);
535        registrar.register_record_schema(schema)?;
536        registrar.register_ingress(execute_descriptor())?;
537        let mut prepare = BoundarySystemContract::new(
538            "prepare-county",
539            BoundaryPhase::HistoricalCandidateEvaluation,
540            SystemCadence::EventDriven,
541        );
542        prepare.reads = vec![order_state(), StateKey::core_ingress()];
543        prepare.writes = vec![county_state()];
544        prepare.visibility = StateVisibility::SameBoundary;
545        registrar.register_boundary_system(prepare, prepare_county)
546    }
547}
548
549fn main() -> Result<(), CanwuError> {
550    let central = CentralPlugin;
551    let treasury = TreasuryPlugin;
552    let county = CountyPlugin;
553    let plugins: [&dyn SimulationPlugin; 3] = [&central, &treasury, &county];
554    let mut canwu =
555        Canwu::new_with_plugins(11, Scenario::new(SimTime::EPOCH, Vec::new()), &plugins)?;
556
557    canwu.enqueue_plugin_ingress(canwu_api::PluginIngressRequest::new(
558        CENTRAL_PLUGIN,
559        ISSUE_INGRESS,
560        SimTime::EPOCH,
561        json!({"order_id": ORDER_ID}),
562    ))?;
563
564    let first = canwu.settle_boundary(BoundaryRequest::at(SimTime::EPOCH))?;
565    assert_eq!(first.generated_ingress.len(), 2);
566    assert!(canwu.typed_domain_record(&order_reference()).is_some());
567    assert!(canwu.typed_domain_record(&treasury_reference()).is_none());
568    assert!(canwu.typed_domain_record(&county_reference()).is_none());
569
570    let second = canwu.settle_boundary(BoundaryRequest::at(SimTime::EPOCH))?;
571    assert_eq!(second.record_change_count, 2);
572    assert!(canwu.typed_domain_record(&treasury_reference()).is_some());
573    assert!(canwu.typed_domain_record(&county_reference()).is_some());
574
575    let snapshot = canwu.snapshot_json()?;
576    let restored = Canwu::from_snapshot_json_with_plugins(&snapshot, &plugins)?;
577    let replayed = Canwu::replay_from_journal(&plugins, &canwu.replay_journal())?;
578    assert_eq!(restored.snapshot(), canwu.snapshot());
579    assert_eq!(replayed.snapshot(), canwu.snapshot());
580
581    println!(
582        "relief_order={} treasury_grain={} county_grain={} exact_replay=ok",
583        ORDER_ID, 600, 600
584    );
585    Ok(())
586}