1#![allow(clippy::too_many_lines, clippy::unnecessary_wraps)]
2
3use 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] = [¢ral, &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}