1use super::{CanwuError, ErrorCode};
2use canwu_core::{
3 ArmyId, BoundaryId, DecisionTicketId, DeterministicRng, DomainRecordRef, EntityRef, EventId,
4 EvidenceRef, KnowledgeHolderRef, PersonId, RandomDrawId,
5};
6use canwu_event::CauseRef;
7use canwu_time::SimTime;
8use serde::{Deserialize, Serialize};
9use std::collections::BTreeMap;
10
11const STREAM_DERIVATION_DOMAIN: &[u8] = b"canwu.random-stream.v1";
12const OPERATION_DOMAIN: &[u8] = b"canwu.random.operation.v1";
13const PURPOSE_DOMAIN: &[u8] = b"canwu.random.purpose.v1";
14const OPERATION_TEXT_BYTES: usize = 256;
15
16#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
17#[serde(rename_all = "snake_case")]
18pub enum RandomAlgorithm {
19 #[default]
21 SplitMix64V2,
22 SplitMix64V1,
24}
25
26#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
27pub struct RandomStreamKey {
28 pub namespace: String,
29 pub name: String,
30 pub version: u32,
31}
32
33impl RandomStreamKey {
34 #[must_use]
35 pub fn new(namespace: impl Into<String>, name: impl Into<String>, version: u32) -> Self {
36 Self {
37 namespace: namespace.into(),
38 name: name.into(),
39 version,
40 }
41 }
42}
43
44#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
45pub struct RandomStreamState {
46 pub key: RandomStreamKey,
47 pub algorithm: RandomAlgorithm,
48 pub seed: u64,
49 pub position: u64,
50 pub generator_state: u64,
51}
52
53impl RandomStreamState {
54 pub(crate) fn initial(root_seed: u64, key: RandomStreamKey) -> Self {
55 let seed = derive_stream_seed(root_seed, &key);
56 Self {
57 key,
58 algorithm: RandomAlgorithm::SplitMix64V2,
59 seed,
60 position: 0,
61 generator_state: seed,
62 }
63 }
64
65 pub(crate) fn is_coherent(&self, root_seed: u64) -> bool {
66 matches!(
67 self.algorithm,
68 RandomAlgorithm::SplitMix64V1 | RandomAlgorithm::SplitMix64V2
69 ) && self.seed == derive_stream_seed(root_seed, &self.key)
70 }
71}
72
73#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
74#[serde(tag = "type", rename_all = "snake_case")]
75pub enum RandomDrawProducer {
76 BoundarySystem {
77 boundary: BoundaryId,
78 plugin: String,
79 system: String,
80 },
81 CoreSystem {
82 system: String,
83 },
84}
85
86#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
87#[serde(tag = "type", rename_all = "snake_case")]
88pub enum RandomDrawOutcome {
89 BoundarySystemDecision,
90 DecisionSelection {
91 ticket_id: DecisionTicketId,
92 ticket_version: u64,
93 option_id: String,
94 },
95 KnowledgeReportDelivery {
96 recipient: PersonId,
97 army: ArmyId,
98 dispatch_event: EventId,
99 arrives_at: SimTime,
100 },
101}
102
103#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
107#[serde(tag = "type", content = "value", rename_all = "snake_case")]
108pub enum RandomOperationTarget {
109 Entity(EntityRef),
110 DecisionTicket {
111 ticket_id: DecisionTicketId,
112 ticket_version: u64,
113 },
114 DomainRecord {
115 record: DomainRecordRef,
116 version: u64,
117 },
118 KnowledgeHolder(KnowledgeHolderRef),
119 CanonicalKey(String),
120}
121
122#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
124pub struct RandomOperationAddressV1 {
125 pub producer_plugin: String,
126 pub operation_kind: String,
127 pub application_operation_id: String,
128 pub target: RandomOperationTarget,
129 pub draw_slot: u32,
130}
131
132#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
134#[serde(tag = "type", content = "value", rename_all = "snake_case")]
135pub enum RandomDrawAddress {
136 Sequential { position: u64 },
137 OperationV1(RandomOperationAddressV1),
138}
139
140#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
141pub struct RandomSample {
142 pub stream: RandomStreamKey,
143 pub address: RandomDrawAddress,
144 pub upper_exclusive: u64,
145 pub value: u64,
146}
147
148impl RandomDrawAddress {
149 #[must_use]
150 pub const fn sequential_position(&self) -> Option<u64> {
151 match self {
152 Self::Sequential { position } => Some(*position),
153 Self::OperationV1(_) => None,
154 }
155 }
156}
157#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
158pub struct RandomDrawRecord {
159 pub id: RandomDrawId,
160 pub at: SimTime,
161 pub stream: RandomStreamKey,
162 pub address: RandomDrawAddress,
163 #[serde(default, skip_serializing_if = "Option::is_none")]
164 pub operation_evidence: Option<EvidenceRef>,
165 pub upper_exclusive: u64,
166 pub value: u64,
167 pub purpose: String,
168 pub producer: RandomDrawProducer,
169 #[serde(default)]
170 pub outcome: Option<RandomDrawOutcome>,
171 pub cause: CauseRef,
172 pub correlation_id: u64,
173}
174
175#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
176pub struct KeyedDrawReservation {
177 pub stream: RandomStreamKey,
178 pub address: RandomOperationAddressV1,
179 pub upper_exclusive: u64,
180 pub purpose_hash: String,
181 pub result: u64,
182 pub draw_id: RandomDrawId,
183 pub operation_evidence: EvidenceRef,
184 pub draw_receipt: crate::ArchivedEvidenceReceipt,
185}
186
187#[derive(Clone, Debug)]
188pub(crate) struct PendingRandomDraw {
189 pub stream: RandomStreamKey,
190 pub address: RandomDrawAddress,
191 pub operation_evidence: Option<EvidenceRef>,
192 pub upper_exclusive: u64,
193 pub value: u64,
194 pub purpose: String,
195}
196
197#[derive(Clone, Debug, Eq, PartialEq)]
198pub(crate) struct KeyedDrawMemo {
199 pub operation_evidence: EvidenceRef,
200 pub upper_exclusive: u64,
201 pub value: u64,
202 pub purpose_hash: String,
203}
204
205pub(crate) type KeyedDrawIndex =
206 BTreeMap<(RandomStreamKey, RandomOperationAddressV1), KeyedDrawMemo>;
207
208pub(crate) struct RandomExecution {
209 pub states: BTreeMap<RandomStreamKey, RandomStreamState>,
210 pub draws: Vec<PendingRandomDraw>,
211}
212
213pub(crate) struct RandomSession {
214 states: BTreeMap<RandomStreamKey, RandomStreamState>,
215 draws: Vec<PendingRandomDraw>,
216 root_seed: u64,
217 producer_plugin: String,
218 keyed: KeyedDrawIndex,
219}
220
221impl RandomSession {
222 pub(crate) fn new(
223 available: &BTreeMap<RandomStreamKey, RandomStreamState>,
224 allowed: &[RandomStreamKey],
225 root_seed: u64,
226 producer_plugin: &str,
227 keyed: &KeyedDrawIndex,
228 ) -> Result<Self, CanwuError> {
229 let mut states = BTreeMap::new();
230 for key in allowed {
231 let Some(state) = available.get(key) else {
232 return Err(CanwuError::new(
233 ErrorCode::InvalidRandomStream,
234 format!(
235 "declared random stream {}.{}@{} is not initialized",
236 key.namespace, key.name, key.version
237 ),
238 ));
239 };
240 states.insert(key.clone(), state.clone());
241 }
242 Ok(Self {
243 states,
244 draws: Vec::new(),
245 root_seed,
246 producer_plugin: producer_plugin.to_owned(),
247 keyed: keyed.clone(),
248 })
249 }
250
251 pub(crate) fn range(
252 &mut self,
253 key: &RandomStreamKey,
254 upper_exclusive: u64,
255 purpose: &str,
256 ) -> Result<u64, CanwuError> {
257 if upper_exclusive == 0 || purpose.trim().is_empty() || purpose != purpose.trim() {
258 return Err(CanwuError::new(
259 ErrorCode::InvalidRandomDraw,
260 "random draws require a positive bound and canonical purpose",
261 ));
262 }
263 let Some(state) = self.states.get_mut(key) else {
264 return Err(CanwuError::new(
265 ErrorCode::UndeclaredRandomStream,
266 format!(
267 "random stream {}.{}@{} was not declared by this system",
268 key.namespace, key.name, key.version
269 ),
270 ));
271 };
272 let next_position = state.position.checked_add(1).ok_or_else(|| {
273 CanwuError::new(
274 ErrorCode::IdentifierExhausted,
275 "random stream position is exhausted",
276 )
277 })?;
278 let position = state.position;
279 let mut generator = DeterministicRng::from_seed(state.generator_state);
280 let value = match state.algorithm {
281 RandomAlgorithm::SplitMix64V1 => generator.range_modulo(upper_exclusive),
282 RandomAlgorithm::SplitMix64V2 => generator.range(upper_exclusive),
283 };
284 state.position = next_position;
285 state.generator_state = generator.state();
286 self.draws.push(PendingRandomDraw {
287 stream: key.clone(),
288 address: RandomDrawAddress::Sequential { position },
289 operation_evidence: None,
290 upper_exclusive,
291 value,
292 purpose: purpose.to_owned(),
293 });
294 Ok(value)
295 }
296
297 #[cfg(test)]
298 #[allow(clippy::too_many_arguments)]
299 pub(crate) fn range_for_operation(
300 &mut self,
301 key: &RandomStreamKey,
302 evidence: EvidenceRef,
303 operation_kind: &str,
304 application_operation_id: &str,
305 target: RandomOperationTarget,
306 draw_slot: u32,
307 upper_exclusive: u64,
308 purpose: &str,
309 ) -> Result<u64, CanwuError> {
310 self.sample_for_operation(
311 key,
312 evidence,
313 operation_kind,
314 application_operation_id,
315 target,
316 draw_slot,
317 upper_exclusive,
318 purpose,
319 )
320 .map(|sample| sample.value)
321 }
322
323 #[allow(clippy::too_many_arguments)]
324 pub(crate) fn sample_for_operation(
325 &mut self,
326 key: &RandomStreamKey,
327 evidence: EvidenceRef,
328 operation_kind: &str,
329 application_operation_id: &str,
330 target: RandomOperationTarget,
331 draw_slot: u32,
332 upper_exclusive: u64,
333 purpose: &str,
334 ) -> Result<RandomSample, CanwuError> {
335 if !self.states.contains_key(key) {
336 return Err(CanwuError::new(
337 ErrorCode::UndeclaredRandomStream,
338 format!(
339 "random stream {}.{}@{} was not declared by this system",
340 key.namespace, key.name, key.version
341 ),
342 ));
343 }
344 let address = RandomOperationAddressV1 {
345 producer_plugin: self.producer_plugin.clone(),
346 operation_kind: operation_kind.to_owned(),
347 application_operation_id: application_operation_id.to_owned(),
348 target,
349 draw_slot,
350 };
351 validate_operation_inputs(key, &address, upper_exclusive, purpose)?;
352 let index_key = (key.clone(), address.clone());
353 let purpose_hash = purpose_hash_hex_v1(purpose)?;
354 if let Some(existing) = self.keyed.get(&index_key) {
355 if existing.operation_evidence == evidence
356 && existing.upper_exclusive == upper_exclusive
357 && existing.purpose_hash == purpose_hash
358 {
359 return Ok(RandomSample {
360 stream: key.clone(),
361 address: RandomDrawAddress::OperationV1(address),
362 upper_exclusive,
363 value: existing.value,
364 });
365 }
366 return Err(CanwuError::new(
367 ErrorCode::RandomOperationConflict,
368 "operation-keyed random address was reused with different evidence, bound, or purpose",
369 ));
370 }
371 let value = operation_value_v1(self.root_seed, key, &address, upper_exclusive, purpose)?;
372 let memo = KeyedDrawMemo {
373 operation_evidence: evidence.clone(),
374 upper_exclusive,
375 value,
376 purpose_hash,
377 };
378 self.keyed.insert(index_key, memo);
379 self.draws.push(PendingRandomDraw {
380 stream: key.clone(),
381 address: RandomDrawAddress::OperationV1(address.clone()),
382 operation_evidence: Some(evidence),
383 upper_exclusive,
384 value,
385 purpose: purpose.to_owned(),
386 });
387 Ok(RandomSample {
388 stream: key.clone(),
389 address: RandomDrawAddress::OperationV1(address),
390 upper_exclusive,
391 value,
392 })
393 }
394
395 pub(crate) fn finish(self) -> RandomExecution {
396 RandomExecution {
397 states: self.states,
398 draws: self.draws,
399 }
400 }
401}
402
403pub(crate) fn retained_keyed_draws(
404 draws: &[RandomDrawRecord],
405) -> Result<KeyedDrawIndex, CanwuError> {
406 let mut index = BTreeMap::new();
407 for draw in draws {
408 let RandomDrawAddress::OperationV1(address) = &draw.address else {
409 continue;
410 };
411 let evidence = draw.operation_evidence.clone().ok_or_else(|| {
412 CanwuError::new(
413 ErrorCode::InvalidRandomDraw,
414 "operation-keyed draw is missing its evidence reference",
415 )
416 })?;
417 validate_operation_inputs(&draw.stream, address, draw.upper_exclusive, &draw.purpose)?;
418 let key = (draw.stream.clone(), address.clone());
419 let memo = KeyedDrawMemo {
420 operation_evidence: evidence,
421 upper_exclusive: draw.upper_exclusive,
422 value: draw.value,
423 purpose_hash: purpose_hash_hex_v1(&draw.purpose)?,
424 };
425 if index.insert(key, memo).is_some() {
426 return Err(CanwuError::new(
427 ErrorCode::RandomOperationConflict,
428 "random journal contains a duplicate operation-keyed address",
429 ));
430 }
431 }
432 Ok(index)
433}
434
435pub(crate) fn keyed_draws_with_reservations(
436 draws: &[RandomDrawRecord],
437 reservations: &[KeyedDrawReservation],
438) -> Result<KeyedDrawIndex, CanwuError> {
439 let mut index = retained_keyed_draws(draws)?;
440 for reservation in reservations {
441 validate_operation_address(
442 &reservation.stream,
443 &reservation.address,
444 reservation.upper_exclusive,
445 )?;
446 if reservation.purpose_hash.len() != 64
447 || reservation
448 .purpose_hash
449 .bytes()
450 .any(|byte| !byte.is_ascii_digit() && !(b'a'..=b'f').contains(&byte))
451 || reservation.draw_receipt.evidence != EvidenceRef::RandomDraw(reservation.draw_id)
452 {
453 return Err(CanwuError::new(
454 ErrorCode::InvalidRandomDraw,
455 "keyed draw reservation has an invalid purpose hash or draw receipt",
456 ));
457 }
458 let key = (reservation.stream.clone(), reservation.address.clone());
459 let memo = KeyedDrawMemo {
460 operation_evidence: reservation.operation_evidence.clone(),
461 upper_exclusive: reservation.upper_exclusive,
462 value: reservation.result,
463 purpose_hash: reservation.purpose_hash.clone(),
464 };
465 if index.insert(key, memo).is_some() {
466 return Err(CanwuError::new(
467 ErrorCode::RandomOperationConflict,
468 "keyed draw reservation overlaps retained or reserved evidence",
469 ));
470 }
471 }
472 Ok(index)
473}
474
475pub(crate) fn extend_keyed_draws(
476 index: &mut KeyedDrawIndex,
477 draws: &[PendingRandomDraw],
478) -> Result<(), CanwuError> {
479 for draw in draws {
480 let RandomDrawAddress::OperationV1(address) = &draw.address else {
481 continue;
482 };
483 let evidence = draw.operation_evidence.clone().ok_or_else(|| {
484 CanwuError::new(
485 ErrorCode::InvalidRandomDraw,
486 "pending operation-keyed draw is missing evidence",
487 )
488 })?;
489 let memo = KeyedDrawMemo {
490 operation_evidence: evidence,
491 upper_exclusive: draw.upper_exclusive,
492 value: draw.value,
493 purpose_hash: purpose_hash_hex_v1(&draw.purpose)?,
494 };
495 if index
496 .insert((draw.stream.clone(), address.clone()), memo)
497 .is_some()
498 {
499 return Err(CanwuError::new(
500 ErrorCode::RandomOperationConflict,
501 "pending random execution duplicated an operation-keyed address",
502 ));
503 }
504 }
505 Ok(())
506}
507
508pub(crate) fn validate_operation_draw(
509 root_seed: u64,
510 draw: &RandomDrawRecord,
511) -> Result<(), CanwuError> {
512 match (&draw.address, &draw.operation_evidence) {
513 (RandomDrawAddress::Sequential { .. }, None) => Ok(()),
514 (RandomDrawAddress::Sequential { .. }, Some(_))
515 | (RandomDrawAddress::OperationV1(_), None) => Err(CanwuError::new(
516 ErrorCode::InvalidRandomDraw,
517 "random address and operation evidence are inconsistent",
518 )),
519 (RandomDrawAddress::OperationV1(address), Some(_)) => {
520 let expected = operation_value_v1(
521 root_seed,
522 &draw.stream,
523 address,
524 draw.upper_exclusive,
525 &draw.purpose,
526 )?;
527 if expected != draw.value {
528 return Err(CanwuError::new(
529 ErrorCode::InvalidRandomDraw,
530 "operation-keyed random value does not match its exact V1 address",
531 ));
532 }
533 Ok(())
534 }
535 }
536}
537
538fn operation_value_v1(
539 root_seed: u64,
540 key: &RandomStreamKey,
541 address: &RandomOperationAddressV1,
542 upper_exclusive: u64,
543 purpose: &str,
544) -> Result<u64, CanwuError> {
545 validate_operation_inputs(key, address, upper_exclusive, purpose)?;
546 let purpose_hash = purpose_hash_v1(purpose)?;
547 for candidate_index in 0..=u32::MAX {
548 let bytes = operation_input_v1(
549 root_seed,
550 key,
551 address,
552 upper_exclusive,
553 &purpose_hash,
554 candidate_index,
555 )?;
556 let digest = blake3::hash(&bytes);
557 let mut candidate_bytes = [0_u8; 8];
558 candidate_bytes.copy_from_slice(&digest.as_bytes()[..8]);
559 let candidate = u64::from_le_bytes(candidate_bytes);
560 let range = 1_u128 << 64;
561 let bound = u128::from(upper_exclusive);
562 let accept_limit = (range / bound) * bound;
563 if u128::from(candidate) < accept_limit {
564 return Ok(candidate % upper_exclusive);
565 }
566 }
567 Err(CanwuError::new(
568 ErrorCode::IdentifierExhausted,
569 "operation-keyed random candidate space is exhausted",
570 ))
571}
572
573fn validate_operation_inputs(
574 key: &RandomStreamKey,
575 address: &RandomOperationAddressV1,
576 upper_exclusive: u64,
577 purpose: &str,
578) -> Result<(), CanwuError> {
579 validate_operation_address(key, address, upper_exclusive)?;
580 validate_operation_text(purpose)
581}
582
583fn validate_operation_address(
584 key: &RandomStreamKey,
585 address: &RandomOperationAddressV1,
586 upper_exclusive: u64,
587) -> Result<(), CanwuError> {
588 if upper_exclusive == 0 || key.version == 0 {
589 return Err(CanwuError::new(
590 ErrorCode::InvalidRandomDraw,
591 "operation-keyed draws require positive stream version and bound",
592 ));
593 }
594 for value in [
595 key.namespace.as_str(),
596 key.name.as_str(),
597 address.producer_plugin.as_str(),
598 address.operation_kind.as_str(),
599 address.application_operation_id.as_str(),
600 ] {
601 validate_operation_text(value)?;
602 }
603 validate_target(&address.target)
604}
605
606fn validate_target(target: &RandomOperationTarget) -> Result<(), CanwuError> {
607 match target {
608 RandomOperationTarget::Entity(EntityRef::Domain(reference)) => {
609 validate_operation_text(&reference.kind.namespace)?;
610 validate_operation_text(&reference.kind.name)?;
611 validate_operation_text(&reference.id)?;
612 }
613 RandomOperationTarget::DomainRecord { record, version } => {
614 if *version == 0 {
615 return Err(CanwuError::new(
616 ErrorCode::InvalidRandomDraw,
617 "operation-keyed domain record target requires a positive version",
618 ));
619 }
620 validate_operation_text(&record.kind.namespace)?;
621 validate_operation_text(&record.kind.name)?;
622 validate_operation_text(&record.id)?;
623 }
624 RandomOperationTarget::DecisionTicket {
625 ticket_id,
626 ticket_version,
627 } => {
628 if ticket_id.get() == 0 || *ticket_version == 0 {
629 return Err(CanwuError::new(
630 ErrorCode::InvalidRandomDraw,
631 "operation-keyed decision target requires positive ticket identity and version",
632 ));
633 }
634 }
635 RandomOperationTarget::CanonicalKey(value) => validate_operation_text(value)?,
636 RandomOperationTarget::Entity(_) | RandomOperationTarget::KnowledgeHolder(_) => {}
637 }
638 if let RandomOperationTarget::KnowledgeHolder(KnowledgeHolderRef::Entity(EntityRef::Domain(
639 reference,
640 ))) = target
641 {
642 validate_operation_text(&reference.kind.namespace)?;
643 validate_operation_text(&reference.kind.name)?;
644 validate_operation_text(&reference.id)?;
645 }
646 Ok(())
647}
648
649fn validate_operation_text(value: &str) -> Result<(), CanwuError> {
650 if value.is_empty()
651 || value != value.trim()
652 || value.len() > OPERATION_TEXT_BYTES
653 || u32::try_from(value.len()).is_err()
654 {
655 return Err(CanwuError::new(
656 ErrorCode::InvalidRandomDraw,
657 "operation-keyed random text is empty, non-canonical, or too long",
658 ));
659 }
660 Ok(())
661}
662
663fn purpose_hash_v1(purpose: &str) -> Result<[u8; 32], CanwuError> {
664 validate_operation_text(purpose)?;
665 let mut bytes = Vec::new();
666 bytes.extend_from_slice(PURPOSE_DOMAIN);
667 bytes.push(0);
668 put_text(&mut bytes, purpose)?;
669 Ok(*blake3::hash(&bytes).as_bytes())
670}
671
672pub(crate) fn purpose_hash_hex_v1(purpose: &str) -> Result<String, CanwuError> {
673 Ok(blake3::Hash::from_bytes(purpose_hash_v1(purpose)?)
674 .to_hex()
675 .to_string())
676}
677
678fn operation_input_v1(
679 root_seed: u64,
680 key: &RandomStreamKey,
681 address: &RandomOperationAddressV1,
682 upper_exclusive: u64,
683 purpose_hash: &[u8; 32],
684 candidate_index: u32,
685) -> Result<Vec<u8>, CanwuError> {
686 let mut bytes = Vec::new();
687 bytes.extend_from_slice(OPERATION_DOMAIN);
688 bytes.push(0);
689 bytes.push(1);
690 bytes.extend_from_slice(&root_seed.to_le_bytes());
691 put_text(&mut bytes, &key.namespace)?;
692 put_text(&mut bytes, &key.name)?;
693 bytes.extend_from_slice(&key.version.to_le_bytes());
694 put_text(&mut bytes, &address.producer_plugin)?;
695 put_text(&mut bytes, &address.operation_kind)?;
696 put_text(&mut bytes, &address.application_operation_id)?;
697 encode_target(&mut bytes, &address.target)?;
698 bytes.extend_from_slice(&address.draw_slot.to_le_bytes());
699 bytes.extend_from_slice(&upper_exclusive.to_le_bytes());
700 bytes.extend_from_slice(purpose_hash);
701 bytes.extend_from_slice(&candidate_index.to_le_bytes());
702 Ok(bytes)
703}
704
705fn put_text(bytes: &mut Vec<u8>, value: &str) -> Result<(), CanwuError> {
706 validate_operation_text(value)?;
707 let length = u32::try_from(value.len()).map_err(|_| {
708 CanwuError::new(
709 ErrorCode::InvalidRandomDraw,
710 "operation-keyed text length exceeds u32",
711 )
712 })?;
713 bytes.extend_from_slice(&length.to_le_bytes());
714 bytes.extend_from_slice(value.as_bytes());
715 Ok(())
716}
717
718fn encode_target(bytes: &mut Vec<u8>, target: &RandomOperationTarget) -> Result<(), CanwuError> {
719 match target {
720 RandomOperationTarget::Entity(entity) => {
721 bytes.push(1);
722 encode_entity(bytes, entity)?;
723 }
724 RandomOperationTarget::DomainRecord { record, version } => {
725 if *version == 0 {
726 return Err(CanwuError::new(
727 ErrorCode::InvalidRandomDraw,
728 "operation-keyed domain record target requires a positive version",
729 ));
730 }
731 bytes.push(2);
732 put_text(bytes, &record.kind.namespace)?;
733 put_text(bytes, &record.kind.name)?;
734 put_text(bytes, &record.id)?;
735 bytes.extend_from_slice(&version.to_le_bytes());
736 }
737 RandomOperationTarget::KnowledgeHolder(holder) => {
738 bytes.push(3);
739 match holder {
740 KnowledgeHolderRef::Person(person) => {
741 bytes.push(1);
742 bytes.extend_from_slice(&person.get().to_le_bytes());
743 }
744 KnowledgeHolderRef::Entity(entity) => {
745 bytes.push(2);
746 encode_entity(bytes, entity)?;
747 }
748 }
749 }
750 RandomOperationTarget::CanonicalKey(value) => {
751 bytes.push(4);
752 put_text(bytes, value)?;
753 }
754 RandomOperationTarget::DecisionTicket {
755 ticket_id,
756 ticket_version,
757 } => {
758 if ticket_id.get() == 0 || *ticket_version == 0 {
759 return Err(CanwuError::new(
760 ErrorCode::InvalidRandomDraw,
761 "operation-keyed decision target requires positive ticket identity and version",
762 ));
763 }
764 bytes.push(5);
765 bytes.extend_from_slice(&ticket_id.get().to_le_bytes());
766 bytes.extend_from_slice(&ticket_version.to_le_bytes());
767 }
768 }
769 Ok(())
770}
771
772fn encode_entity(bytes: &mut Vec<u8>, entity: &EntityRef) -> Result<(), CanwuError> {
773 match entity {
774 EntityRef::Army(id) => {
775 bytes.push(1);
776 bytes.extend_from_slice(&id.get().to_le_bytes());
777 }
778 EntityRef::Government(id) => {
779 bytes.push(2);
780 bytes.extend_from_slice(&id.get().to_le_bytes());
781 }
782 EntityRef::Organization(id) => {
783 bytes.push(3);
784 bytes.extend_from_slice(&id.get().to_le_bytes());
785 }
786 EntityRef::Person(id) => {
787 bytes.push(4);
788 bytes.extend_from_slice(&id.get().to_le_bytes());
789 }
790 EntityRef::Resource(id) => {
791 bytes.push(5);
792 bytes.extend_from_slice(&id.get().to_le_bytes());
793 }
794 EntityRef::Route(id) => {
795 bytes.push(6);
796 bytes.extend_from_slice(&id.get().to_le_bytes());
797 }
798 EntityRef::Territory(id) => {
799 bytes.push(7);
800 bytes.extend_from_slice(&id.get().to_le_bytes());
801 }
802 EntityRef::Domain(reference) => {
803 bytes.push(8);
804 put_text(bytes, &reference.kind.namespace)?;
805 put_text(bytes, &reference.kind.name)?;
806 put_text(bytes, &reference.id)?;
807 }
808 }
809 Ok(())
810}
811
812pub(crate) fn derive_stream_seed(root_seed: u64, key: &RandomStreamKey) -> u64 {
813 if key.namespace == "canwu.core" && key.name == "knowledge-report-delay" && key.version == 1 {
814 return root_seed;
815 }
816 let mut hasher = blake3::Hasher::new();
817 hasher.update(STREAM_DERIVATION_DOMAIN);
818 hasher.update(&root_seed.to_le_bytes());
819 update_text(&mut hasher, &key.namespace);
820 update_text(&mut hasher, &key.name);
821 hasher.update(&key.version.to_le_bytes());
822 let mut seed = [0_u8; 8];
823 seed.copy_from_slice(&hasher.finalize().as_bytes()[..8]);
824 u64::from_le_bytes(seed)
825}
826
827pub(crate) fn core_report_delay_stream() -> RandomStreamKey {
828 RandomStreamKey::new("canwu.core", "knowledge-report-delay", 1)
829}
830
831fn update_text(hasher: &mut blake3::Hasher, value: &str) {
832 let length = u64::try_from(value.len()).unwrap_or(u64::MAX);
833 hasher.update(&length.to_le_bytes());
834 hasher.update(value.as_bytes());
835}
836
837#[cfg(test)]
838mod tests {
839 use super::*;
840 use canwu_core::{GovernmentId, KnowledgeHolderRef, PersonId};
841
842 fn fixture_stream() -> RandomStreamKey {
843 RandomStreamKey::new("fixture.random", "resolution", 1)
844 }
845
846 fn fixture_address(target: RandomOperationTarget) -> RandomOperationAddressV1 {
847 RandomOperationAddressV1 {
848 producer_plugin: "fixture-random".to_owned(),
849 operation_kind: "resolve".to_owned(),
850 application_operation_id: "operation-alpha".to_owned(),
851 target,
852 draw_slot: 3,
853 }
854 }
855
856 #[test]
857 fn operation_v1_golden_vectors_cover_every_target_encoding() {
858 let stream = fixture_stream();
859 let targets = [
860 RandomOperationTarget::Entity(EntityRef::Government(GovernmentId::new(9))),
861 RandomOperationTarget::DomainRecord {
862 record: DomainRecordRef {
863 kind: canwu_core::DomainRecordKind::new("fixture", "record"),
864 id: "r-7".to_owned(),
865 },
866 version: 4,
867 },
868 RandomOperationTarget::KnowledgeHolder(KnowledgeHolderRef::Person(PersonId::new(5))),
869 RandomOperationTarget::CanonicalKey("键-α".to_owned()),
870 ];
871 let values = targets
872 .into_iter()
873 .map(|target| {
874 let address = fixture_address(target);
875 let purpose_hash = purpose_hash_v1("stable outcome").expect("purpose should hash");
876 let input = operation_input_v1(
877 0x0102_0304_0506_0708,
878 &stream,
879 &address,
880 10_000,
881 &purpose_hash,
882 0,
883 )
884 .expect("golden input should encode");
885 let digest = blake3::hash(&input).to_hex().to_string();
886 let value = operation_value_v1(
887 0x0102_0304_0506_0708,
888 &stream,
889 &address,
890 10_000,
891 "stable outcome",
892 )
893 .expect("golden vector should encode");
894 (input, digest, value)
895 })
896 .collect::<Vec<_>>();
897 assert_eq!(
898 values
899 .iter()
900 .map(|(_, digest, _)| digest.as_str())
901 .collect::<Vec<_>>(),
902 vec![
903 "55b21978711c8d81a42bdacef84c8b22e16bf5e2a36135097f850e658ed86a74",
904 "4ac40330c89fc0bce339b16c59fa25166f714b664f4fd3e9e252cb9960b979d7",
905 "c23814e1a5343aba6e00b88c1267aef583d6adf2385173b80268f14c6d34f594",
906 "8f84ed6f5700db025dff8e59966c8f14d9c2a262904726688ea08d9fd043816d",
907 ]
908 );
909 assert_eq!(
910 values
911 .iter()
912 .map(|(_, _, value)| *value)
913 .collect::<Vec<_>>(),
914 vec![8389, 6730, 1186, 1231]
915 );
916 assert_eq!(
917 values[0].0.iter().fold(String::new(), |mut output, byte| {
918 use std::fmt::Write as _;
919 write!(&mut output, "{byte:02x}").expect("writing to a string cannot fail");
920 output
921 }),
922 "63616e77752e72616e646f6d2e6f7065726174696f6e2e7631000108070605040302010e000000666978747572652e72616e646f6d0a0000007265736f6c7574696f6e010000000e000000666978747572652d72616e646f6d070000007265736f6c76650f0000006f7065726174696f6e2d616c70686101020900000000000000030000001027000000000000784ea13c085b19a2d0420898b8c3848517c24f30b1d82658617e80f99edff43a00000000"
923 );
924 }
925
926 #[test]
927 fn keyed_retry_is_idempotent_conflicts_fail_and_sequential_state_does_not_move() {
928 let stream = fixture_stream();
929 let state = RandomStreamState::initial(41, stream.clone());
930 let mut session = RandomSession::new(
931 &BTreeMap::from([(stream.clone(), state)]),
932 std::slice::from_ref(&stream),
933 41,
934 "fixture-random",
935 &BTreeMap::new(),
936 )
937 .expect("session should initialize");
938 let evidence = EvidenceRef::Ingress(canwu_core::IngressId::new(2));
939 let first = session
940 .range_for_operation(
941 &stream,
942 evidence.clone(),
943 "resolve",
944 "operation-alpha",
945 RandomOperationTarget::CanonicalKey("target".to_owned()),
946 0,
947 100,
948 "resolution",
949 )
950 .expect("first keyed draw should succeed");
951 let retry = session
952 .range_for_operation(
953 &stream,
954 evidence.clone(),
955 "resolve",
956 "operation-alpha",
957 RandomOperationTarget::CanonicalKey("target".to_owned()),
958 0,
959 100,
960 "resolution",
961 )
962 .expect("exact retry should reuse the result");
963 assert_eq!(retry, first);
964 let conflict = session
965 .range_for_operation(
966 &stream,
967 evidence,
968 "resolve",
969 "operation-alpha",
970 RandomOperationTarget::CanonicalKey("target".to_owned()),
971 0,
972 101,
973 "resolution",
974 )
975 .expect_err("changed bound must conflict");
976 assert_eq!(conflict.code, ErrorCode::RandomOperationConflict);
977 let purpose_conflict = session
978 .range_for_operation(
979 &stream,
980 EvidenceRef::Ingress(canwu_core::IngressId::new(2)),
981 "resolve",
982 "operation-alpha",
983 RandomOperationTarget::CanonicalKey("target".to_owned()),
984 0,
985 100,
986 "different-resolution",
987 )
988 .expect_err("changed purpose must conflict");
989 assert_eq!(purpose_conflict.code, ErrorCode::RandomOperationConflict);
990 let execution = session.finish();
991 assert_eq!(execution.draws.len(), 1);
992 assert_eq!(execution.states[&stream].position, 0);
993 }
994
995 #[test]
996 fn evidence_renumbering_and_unrelated_operations_do_not_change_keyed_entropy() {
997 let stream = fixture_stream();
998 let run = |evidence, include_unrelated| {
999 let state = RandomStreamState::initial(91, stream.clone());
1000 let mut session = RandomSession::new(
1001 &BTreeMap::from([(stream.clone(), state)]),
1002 std::slice::from_ref(&stream),
1003 91,
1004 "fixture-random",
1005 &BTreeMap::new(),
1006 )
1007 .expect("session should initialize");
1008 if include_unrelated {
1009 session
1010 .range_for_operation(
1011 &stream,
1012 EvidenceRef::Ingress(canwu_core::IngressId::new(1)),
1013 "resolve",
1014 "unrelated-operation",
1015 RandomOperationTarget::CanonicalKey("unrelated".to_owned()),
1016 0,
1017 1_000,
1018 "unrelated-purpose",
1019 )
1020 .expect("unrelated keyed operation should succeed");
1021 }
1022 session
1023 .range_for_operation(
1024 &stream,
1025 EvidenceRef::Ingress(canwu_core::IngressId::new(evidence)),
1026 "resolve",
1027 "operation-alpha",
1028 RandomOperationTarget::CanonicalKey("target".to_owned()),
1029 0,
1030 1_000,
1031 "resolution",
1032 )
1033 .expect("target keyed operation should succeed")
1034 };
1035 let baseline = run(2, false);
1036 assert_eq!(run(99, false), baseline);
1037 assert_eq!(run(2, true), baseline);
1038 }
1039
1040 #[test]
1041 fn producer_namespace_is_encoded_and_rejection_reduction_retries() {
1042 let stream = fixture_stream();
1043 let mut first = fixture_address(RandomOperationTarget::CanonicalKey("target".to_owned()));
1044 let mut second = first.clone();
1045 second.producer_plugin = "fixture-random-b".to_owned();
1046 assert_ne!(first, second);
1047 let purpose_hash = purpose_hash_v1("resolution").expect("purpose should hash");
1048 assert_ne!(
1049 operation_input_v1(17, &stream, &first, 100, &purpose_hash, 0)
1050 .expect("first producer input"),
1051 operation_input_v1(17, &stream, &second, 100, &purpose_hash, 0)
1052 .expect("second producer input")
1053 );
1054
1055 first.application_operation_id = "rejection-reduction".to_owned();
1056 let upper_exclusive = (1_u64 << 63) + 1;
1057 let range = 1_u128 << 64;
1058 let bound = u128::from(upper_exclusive);
1059 let accept_limit = (range / bound) * bound;
1060 let purpose = (0_u32..10_000)
1061 .map(|index| format!("retry-purpose-{index}"))
1062 .find(|purpose| {
1063 let purpose_hash = purpose_hash_v1(purpose).expect("purpose should hash");
1064 let bytes =
1065 operation_input_v1(17, &stream, &first, upper_exclusive, &purpose_hash, 0)
1066 .expect("candidate zero should encode");
1067 let digest = blake3::hash(&bytes);
1068 let mut candidate_bytes = [0_u8; 8];
1069 candidate_bytes.copy_from_slice(&digest.as_bytes()[..8]);
1070 u128::from(u64::from_le_bytes(candidate_bytes)) >= accept_limit
1071 })
1072 .expect("fixture search should find a rejected candidate zero");
1073 let value = operation_value_v1(17, &stream, &first, upper_exclusive, &purpose)
1074 .expect("rejection reduction must find a later candidate");
1075 assert!(value < upper_exclusive);
1076 }
1077
1078 #[test]
1079 fn sequential_rejection_sampling_replays_the_actual_generator_state() {
1080 let stream = fixture_stream();
1081 let upper_exclusive = (1_u64 << 63) + 1;
1082 let rejection_threshold = upper_exclusive.wrapping_neg() % upper_exclusive;
1083 let (root_seed, initial) = (1_u64..)
1084 .find_map(|root_seed| {
1085 let initial = RandomStreamState::initial(root_seed, stream.clone());
1086 let mut probe = DeterministicRng::from_seed(initial.generator_state);
1087 (probe.next_u64() < rejection_threshold).then_some((root_seed, initial))
1088 })
1089 .expect("fixture search should find a sequential rejection");
1090 let mut session = RandomSession::new(
1091 &BTreeMap::from([(stream.clone(), initial.clone())]),
1092 std::slice::from_ref(&stream),
1093 root_seed,
1094 "fixture-random",
1095 &BTreeMap::new(),
1096 )
1097 .expect("session should initialize");
1098 let value = session
1099 .range(&stream, upper_exclusive, "sequential rejection")
1100 .expect("sequential rejection draw should succeed");
1101 let execution = session.finish();
1102 let final_state = &execution.states[&stream];
1103
1104 let mut replay = DeterministicRng::from_seed(initial.generator_state);
1105 assert_eq!(replay.range(upper_exclusive), value);
1106 assert_eq!(replay.state(), final_state.generator_state);
1107 assert_eq!(final_state.position, 1);
1108 assert_ne!(
1109 final_state.generator_state,
1110 DeterministicRng::state_after(initial.seed, final_state.position)
1111 );
1112 }
1113}