Skip to main content

canwu_sim/runtime/
random.rs

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    /// Current unbiased `SplitMix64` range reduction.
20    #[default]
21    SplitMix64V2,
22    /// Historical modulo-reduction behavior retained for existing journals.
23    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/// Stable application target for an operation-addressed random draw.
104///
105/// Format-7 target for the enabled byte-exact keyed algorithm.
106#[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/// Version-one stable entropy address for a future keyed random draw.
123#[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/// Persisted address of a random draw.
133#[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}