From ca0fc5a628d12e096e5c3aae98f5e32cd0f4b5e3 Mon Sep 17 00:00:00 2001 From: Kenneth Pernyer Date: Thu, 2 Jul 2026 10:55:06 +0200 Subject: [PATCH 1/2] fix: add multi-user-consensus crate and smoke tests missing from the wiring commit Commit 1d16939 wired crates/multi-user-consensus into the workspace members and workspace.dependencies, but the crate directory and the new cross-extension-smoke tests were never committed, leaving main unbuildable (cargo cannot load the workspace). Add the crate and the three smoke tests that exercise it and the atelier/helm scenarios. Co-Authored-By: Claude Fable 5 --- .../tests/atelier_showcase_examples.rs | 161 +++ .../tests/helm_coordination_headless.rs | 505 ++++++++ .../tests/multi_user_consensus.rs | 145 +++ crates/multi-user-consensus/Cargo.toml | 13 + crates/multi-user-consensus/src/lib.rs | 1056 +++++++++++++++++ 5 files changed, 1880 insertions(+) create mode 100644 crates/cross-extension-smoke/tests/atelier_showcase_examples.rs create mode 100644 crates/cross-extension-smoke/tests/helm_coordination_headless.rs create mode 100644 crates/cross-extension-smoke/tests/multi_user_consensus.rs create mode 100644 crates/multi-user-consensus/Cargo.toml create mode 100644 crates/multi-user-consensus/src/lib.rs diff --git a/crates/cross-extension-smoke/tests/atelier_showcase_examples.rs b/crates/cross-extension-smoke/tests/atelier_showcase_examples.rs new file mode 100644 index 0000000..61e4f62 --- /dev/null +++ b/crates/cross-extension-smoke/tests/atelier_showcase_examples.rs @@ -0,0 +1,161 @@ +//! Atelier showcase cross-module examples. +//! +//! Atelier owns the example agents locally. Arena proves those public example +//! agents still compose with the current Bedrock/Mosaic contract graph. + +use atelier_domain::{ + AvailabilityRetrievalAgent, ConflictDetectionAgent, DomainRecordPayload, DomainTextPayload, + RequireParticipantAvailability, RequirePositiveDuration, RequireValidSlot, + SlotOptimizationAgent, TimeZoneNormalizationAgent, WorkingHoursConstraintAgent, domain_text, + packs::InvoiceIssuerAgent, payload_contains, +}; +use converge_core::suggestors::SeedSuggestor; +use converge_core::{ContextKey, ContextState, Engine, FactPayload}; + +fn projected(facts: &[converge_core::ContextFact]) -> Vec<(ContextKey, String, serde_json::Value)> { + facts + .iter() + .map(|fact| { + ( + fact.key(), + fact.id().as_str().to_string(), + fact.to_wire().expect("fact serializes").payload.payload, + ) + }) + .collect() +} + +#[tokio::test] +async fn atelier_meeting_scheduler_example_converges_through_converge_engine() { + let run = || async { + let mut engine = Engine::new(); + engine.register_suggestor(SeedSuggestor::new("participants", "Alice, Bob, Carol")); + engine.register_suggestor(SeedSuggestor::new("duration", "60")); + engine.register_suggestor(AvailabilityRetrievalAgent); + engine.register_suggestor(TimeZoneNormalizationAgent); + engine.register_suggestor(WorkingHoursConstraintAgent); + engine.register_suggestor(SlotOptimizationAgent); + engine.register_suggestor(ConflictDetectionAgent); + engine.register_invariant(RequireParticipantAvailability); + engine.register_invariant(RequirePositiveDuration); + engine.register_invariant(RequireValidSlot); + engine + .run(ContextState::new()) + .await + .expect("atelier meeting scheduler should converge") + }; + + let first = run().await; + let second = run().await; + + assert!(first.converged); + assert_eq!(first.cycles, second.cycles); + assert_eq!( + projected(first.context.get(ContextKey::Strategies)), + projected(second.context.get(ContextKey::Strategies)), + "candidate slots should be deterministic across runs" + ); + assert_eq!( + projected(first.context.get(ContextKey::Evaluations)), + projected(second.context.get(ContextKey::Evaluations)), + "slot evaluations should be deterministic across runs" + ); + + let evaluations = first.context.get(ContextKey::Evaluations); + assert_eq!(evaluations.len(), 7); + assert!( + evaluations.iter().any(|fact| { + fact.id() == "eval:1" + && domain_text(fact).is_some_and(|text| text.contains("RECOMMENDED")) + }), + "Atelier's first recommended slot should survive Converge promotion" + ); +} + +#[tokio::test] +async fn atelier_money_example_requires_then_accepts_finance_hitl_gate() { + let ready_invoice = r#"{"type":"invoice","state":"ready_to_issue","customer_id":"cust_123","line_items":[{"sku":"svc","amount":12500}],"amount":12500,"currency":"USD"}"#; + + let mut without_approval = ContextState::new(); + without_approval + .add_input( + ContextKey::Proposals, + "invoice:draft:deal_123", + ready_invoice, + ) + .expect("ready invoice fixture should stage"); + + let mut gated_engine = Engine::new(); + gated_engine.register_suggestor(InvoiceIssuerAgent::default()); + let gated = gated_engine + .run(without_approval) + .await + .expect("atelier money gate should run without approval"); + + assert!(gated.converged); + assert!( + gated.context.get(ContextKey::Proposals).iter().any(|fact| { + fact.id() == "invoice:issue_request:invoice:draft:deal_123" + && payload_contains(fact, "\"required_role\":\"finance_manager\"") + }), + "Atelier money pack should request finance approval before issuing" + ); + + let mut with_approval = ContextState::new(); + with_approval + .add_input( + ContextKey::Proposals, + "invoice:draft:deal_123", + ready_invoice, + ) + .expect("ready invoice fixture should stage"); + with_approval + .add_input( + ContextKey::Proposals, + "approval:invoice:invoice:draft:deal_123", + r#"{"target_id":"invoice:draft:deal_123","required_role":"finance_manager"}"#, + ) + .expect("approval fixture should stage"); + + let mut approved_engine = Engine::new(); + approved_engine.register_suggestor(InvoiceIssuerAgent::default()); + let approved = approved_engine + .run(with_approval) + .await + .expect("atelier money gate should run with approval"); + + assert!(approved.converged); + assert!( + approved + .context + .get(ContextKey::Proposals) + .iter() + .any(|fact| { + fact.id() == "invoice:issued:invoice:draft:deal_123" + && payload_contains(fact, "\"state\":\"issued\"") + }), + "finance approval should let the Atelier money example issue the invoice" + ); +} + +#[test] +fn atelier_typed_payloads_keep_public_converge_fact_contract_shape() { + assert_eq!(DomainRecordPayload::FAMILY, "atelier.domain_record"); + assert_eq!(DomainTextPayload::FAMILY, "atelier.domain_text"); + + let provenance = atelier_domain::ATELIER_DOMAIN_PROVENANCE.provenance(); + assert_eq!(provenance.as_str(), "atelier-domain"); + + let record = DomainRecordPayload::new( + "vendor", + serde_json::json!({ + "id": "vendor-a", + "compliant": true, + "years_in_business": 7 + }), + ); + + assert_eq!(record.record_type(), "vendor"); + assert!(record.bool_field_is("compliant", true)); + assert!(record.has_field("years_in_business")); +} diff --git a/crates/cross-extension-smoke/tests/helm_coordination_headless.rs b/crates/cross-extension-smoke/tests/helm_coordination_headless.rs new file mode 100644 index 0000000..d2a9ae4 --- /dev/null +++ b/crates/cross-extension-smoke/tests/helm_coordination_headless.rs @@ -0,0 +1,505 @@ +//! Cross-module proof for the Helm Coordination Layer M0. +//! +//! Atelier owns the narrated example. Arena imports it and asserts the boundary +//! against existing Helm contracts. + +use std::{collections::BTreeSet, sync::Arc}; + +use application_storage::AppConfig; +use helm_coordination_headless::{ + AGENT_COUNT, CROWD_USER_COUNTS, CoordinationEventKind, CrowdConsensusRun, DYNAMIC_USER_LIMIT, + GateOutcome, HUMAN_COUNT, HeadlessCoordinationRun, ParticipantKind, StaticCoordinationFeed, + dynamic_user_script, readiness_payload_hash, readiness_receipt_family, readiness_record_kind, +}; +use helm_module_contracts::HelmModuleState; +use helm_operator_control::{ + AdapterReceiptStatus, AuthorityEffect, EvidenceReadinessStatus, JobVerdict, + OperatorControlModule, OperatorLedgerRecordKind, ReceiptFamily, +}; +use serde_json::Value; + +#[test] +fn atelier_headless_coordination_run_builds_live_helm_readiness() { + let first = HeadlessCoordinationRun::scripted_burst(); + let second = HeadlessCoordinationRun::scripted_burst(); + + assert_eq!(first.event_trace(), second.event_trace()); + assert_eq!(first.participants, second.participants); + assert_eq!(first.contributions, second.contributions); + assert_eq!(first.prepared_proposals, second.prepared_proposals); + assert_eq!(first.admitted_fact_ids, second.admitted_fact_ids); + assert_eq!(first.readiness_snapshot(), second.readiness_snapshot()); + assert_eq!(first.markdown_report(), second.markdown_report()); + assert_eq!(first.jsonl_timeline(), second.jsonl_timeline()); + assert_eq!(first.gate_outcome, GateOutcome::Approved); + assert_eq!(first.participants.len(), HUMAN_COUNT + AGENT_COUNT); + assert_eq!( + first + .participants + .iter() + .filter(|participant| participant.kind == ParticipantKind::Human) + .count(), + HUMAN_COUNT + ); + assert_eq!( + first + .participants + .iter() + .filter(|participant| participant.kind == ParticipantKind::Agent) + .count(), + AGENT_COUNT + ); + assert_eq!(first.contributions.len(), HUMAN_COUNT + AGENT_COUNT); + assert!( + first + .contributions + .iter() + .all(|contribution| contribution.signed + && contribution.provenance == format!("signed://{}", contribution.participant_id)) + ); + + for (expected_sequence, event) in (1_u64..).zip(first.events.iter()) { + assert_eq!(event.sequence, expected_sequence); + } + + let event_kinds: Vec<_> = first.events.iter().map(|event| event.kind).collect(); + assert_eq!( + event_kinds.first(), + Some(&CoordinationEventKind::SessionOpened) + ); + let tail_start = 1 + HUMAN_COUNT + AGENT_COUNT; + assert!( + event_kinds[1..tail_start] + .iter() + .all(|kind| *kind == CoordinationEventKind::ContributionRecorded) + ); + assert_eq!( + &event_kinds[tail_start..], + [ + CoordinationEventKind::Mixed, + CoordinationEventKind::Matched, + CoordinationEventKind::FormationCheckpoint, + CoordinationEventKind::GateRequested, + CoordinationEventKind::GateDecision, + CoordinationEventKind::ConvergeProposalPrepared, + CoordinationEventKind::ConvergeAdmissionRecorded, + CoordinationEventKind::ReadinessSnapshot, + ] + .as_slice() + ); + + let contribution_events: Vec<_> = first + .events + .iter() + .filter(|event| event.kind == CoordinationEventKind::ContributionRecorded) + .collect(); + assert_eq!(contribution_events.len(), HUMAN_COUNT + AGENT_COUNT); + assert!(contribution_events.iter().all(|event| { + event.payload.get("signed").and_then(Value::as_bool) == Some(true) + && event + .payload + .get("suggestion_only") + .and_then(Value::as_bool) + == Some(true) + })); + let arrival_batches: BTreeSet<_> = contribution_events + .iter() + .filter_map(|event| event.payload.get("arrival_batch").and_then(Value::as_u64)) + .collect(); + assert!( + arrival_batches.len() > 1, + "fixture should simulate more than one concurrent arrival batch" + ); + let ordering_keys: Vec<_> = contribution_events + .iter() + .map(|event| { + event + .payload + .get("ordering_key") + .and_then(Value::as_str) + .expect("contribution event records an ordering key") + .to_string() + }) + .collect(); + let mut sorted_ordering_keys = ordering_keys.clone(); + sorted_ordering_keys.sort(); + assert_eq!( + ordering_keys, sorted_ordering_keys, + "concurrent arrivals must replay in deterministic ordering-key order" + ); + assert_eq!( + contribution_events + .iter() + .filter(|event| event + .payload + .get("participant_kind") + .and_then(Value::as_str) + == Some("human")) + .count(), + HUMAN_COUNT + ); + assert_eq!( + contribution_events + .iter() + .filter(|event| event + .payload + .get("participant_kind") + .and_then(Value::as_str) + == Some("agent")) + .count(), + AGENT_COUNT + ); + + let proposal_event = first + .events + .iter() + .find(|event| event.kind == CoordinationEventKind::ConvergeProposalPrepared) + .expect("timeline prepares proposals for Converge admission"); + assert_eq!( + proposal_event + .payload + .get("proposal_ids") + .and_then(Value::as_array) + .map(Vec::len), + Some(first.prepared_proposals.len()) + ); + assert_eq!( + proposal_event + .payload + .get("promoted_facts") + .and_then(Value::as_array) + .map(Vec::len), + Some(0) + ); + + let admission_event = first + .events + .iter() + .find(|event| event.kind == CoordinationEventKind::ConvergeAdmissionRecorded) + .expect("approved gate records a Converge admission outcome"); + assert_eq!( + admission_event + .payload + .get("admitted_fact_ids") + .and_then(Value::as_array) + .map(Vec::len), + Some(first.admitted_fact_ids.len()) + ); + assert!( + !first.admitted_fact_ids.is_empty(), + "approved scenario should expose final facts admitted through Converge" + ); + assert_eq!( + admission_event.payload.get("promoted_facts"), + None, + "admission records Converge fact ids without giving Helm a fact mutation field" + ); + + let snapshot = first.readiness_snapshot(); + let timeline_receipt = first + .events + .last() + .and_then(|event| event.payload.get("receipt_id")) + .and_then(Value::as_str) + .expect("readiness event records a receipt id"); + assert_eq!(timeline_receipt, snapshot.packet.adapter_receipt_id); + assert!(!snapshot.packet.authorizes_domain_action); + assert_eq!( + snapshot.packet.adapter_status, + AdapterReceiptStatus::Succeeded + ); + assert_eq!(snapshot.packet.subject_ref, first.session_id); + assert_eq!(snapshot.packet.domain_hint, "helm-coordination.headless"); + assert_eq!(snapshot.packet.verdict, Some(JobVerdict::Satisfied)); + assert_eq!(snapshot.packet.evidence_status.len(), 6); + assert!( + snapshot + .packet + .evidence_status + .iter() + .all(|evidence| { evidence.status == EvidenceReadinessStatus::Present }) + ); + assert!( + snapshot + .packet + .evidence_status + .iter() + .filter(|evidence| evidence.clause_key != "converge_admission_outcome") + .all(|evidence| evidence.fact_ids.is_empty()) + ); + let admission_evidence = snapshot + .packet + .evidence_status + .iter() + .find(|evidence| evidence.clause_key == "converge_admission_outcome") + .expect("readiness packet records Converge admission outcome"); + assert_eq!(admission_evidence.fact_ids, first.admitted_fact_ids); + assert!( + snapshot + .packet + .verifier_forbidden_actions + .iter() + .any(|action| action == "do not infer app domain authority from Helm readiness") + ); + assert_eq!(snapshot.ledger_entries.len(), 1); + let ledger_entry = &snapshot.ledger_entries[0]; + assert_eq!( + readiness_record_kind(&snapshot), + OperatorLedgerRecordKind::JobReadinessPacket + ); + assert_eq!(ledger_entry.payload_hash, readiness_payload_hash(&snapshot)); + assert_eq!(ledger_entry.sequence, 1); + assert_eq!(ledger_entry.source_ref, snapshot.packet.packet_id); + assert_eq!( + ledger_entry.backlink_ids, + vec![timeline_receipt.to_string()] + ); + assert_eq!(ledger_entry.authority_effect, AuthorityEffect::None); + assert_eq!(ledger_entry.package_id, snapshot.packet.package_id); + assert_eq!(ledger_entry.truth_version, snapshot.packet.truth_version); + assert_eq!(ledger_entry.domain_hint, snapshot.packet.domain_hint); + assert_eq!(readiness_receipt_family(&snapshot), ReceiptFamily::Common); + + let shell_module = OperatorControlModule::new(AppConfig::default()); + let shell_status = shell_module.readiness_status(); + assert_eq!(shell_module.module_state(), HelmModuleState::ShellDefault); + assert_eq!(shell_status.state, HelmModuleState::ShellDefault); + assert!(!shell_module.module_state().is_live()); + assert_eq!( + shell_status.missing_live_requirements, + vec![ + "process_receipt".to_string(), + "integrity_proof".to_string(), + "adapter_receipt".to_string(), + "axiom_report".to_string() + ] + ); + assert_eq!( + HeadlessCoordinationRun::shell_module_state(), + HelmModuleState::ShellDefault + ); + + let live_module = OperatorControlModule::new(AppConfig::default()) + .with_live_readiness_feed(Arc::new(StaticCoordinationFeed::from_run(&first))); + let live_status = live_module.readiness_status(); + assert_eq!(live_module.module_state(), HelmModuleState::Live); + assert_eq!(live_status.state, HelmModuleState::Live); + assert_eq!(first.live_module_state(), HelmModuleState::Live); + assert!( + live_status + .live_requirements + .contains(&"readiness_feed".to_string()) + ); + assert!(live_status.missing_live_requirements.is_empty()); +} + +#[test] +fn humans_as_suggestors_rejects_direct_fact_mutation() { + let mut run = HeadlessCoordinationRun::scripted_burst(); + let before_rejection = run.event_trace(); + let error = run + .reject_direct_fact_mutation("human-001") + .expect_err("coordination must not directly mutate governed context"); + + assert!(error.contains("Converge alone promotes facts")); + assert_eq!( + run.event_trace(), + before_rejection, + "rejected direct mutation must not append side effects" + ); + assert!( + run.events.iter().all(|event| { + !matches!(event.payload.get("promoted_facts"), Some(Value::Array(facts)) if !facts.is_empty()) + }), + "coordination events may prepare proposals, but must not promote facts" + ); +} + +#[test] +fn atelier_outputs_explain_the_same_receipted_timeline() { + let run = HeadlessCoordinationRun::scripted_burst(); + let markdown = run.markdown_report(); + let jsonl = run.jsonl_timeline(); + let parsed_jsonl: Vec = jsonl + .lines() + .map(|line| serde_json::from_str(line).expect("timeline line is JSON")) + .collect(); + + assert!(markdown.contains("Helm Coordination Headless Scenario")); + assert!(markdown.contains("shell module state: ShellDefault")); + assert!(markdown.contains("live module state with feed: Live")); + assert!( + markdown.contains("Converge handoff: proposals prepared, facts admitted through Converge") + ); + assert!(markdown.contains("## Admitted Facts")); + assert!(markdown.contains(&format!( + "receipt: `{}`", + run.readiness_snapshot().packet.adapter_receipt_id + ))); + assert_eq!(parsed_jsonl.len(), run.events.len()); + for (parsed, event) in parsed_jsonl.iter().zip(run.events.iter()) { + assert_eq!( + parsed.get("sequence").and_then(Value::as_u64), + Some(event.sequence) + ); + assert_eq!(parsed.get("payload"), Some(&event.payload)); + } + assert!(jsonl.contains("\"kind\":\"formation-checkpoint\"")); + assert!(jsonl.contains("\"kind\":\"gate-decision\"")); + assert!(jsonl.contains("\"kind\":\"converge-admission-recorded\"")); + assert!(jsonl.contains("\"authorizes_domain_action\":false")); +} + +#[test] +fn gate_reject_and_timeout_variants_block_converge_admission() { + for (gate_outcome, verdict) in [ + (GateOutcome::Rejected, JobVerdict::Blocked), + (GateOutcome::TimedOut, JobVerdict::Exhausted), + ] { + let run = HeadlessCoordinationRun::scripted_burst_with_gate(gate_outcome); + let snapshot = run.readiness_snapshot(); + + assert_eq!(run.gate_outcome, gate_outcome); + assert!(run.prepared_proposals.is_empty()); + assert!(run.admitted_fact_ids.is_empty()); + assert!( + !run.events + .iter() + .any(|event| event.kind == CoordinationEventKind::ConvergeAdmissionRecorded) + ); + + let proposal_event = run + .events + .iter() + .find(|event| event.kind == CoordinationEventKind::ConvergeProposalPrepared) + .expect("blocked gate still records a Converge handoff decision"); + assert_eq!( + proposal_event + .payload + .get("withheld_reason") + .and_then(Value::as_str), + Some(gate_outcome.as_str()) + ); + + assert!(!snapshot.packet.authorizes_domain_action); + assert_eq!( + snapshot.packet.adapter_status, + AdapterReceiptStatus::Rejected + ); + assert_eq!(snapshot.packet.verdict, Some(verdict)); + assert!(snapshot.packet.evidence_status.iter().any(|evidence| { + evidence.clause_key == "converge_admission_outcome" + && evidence.status == EvidenceReadinessStatus::Blocked + && evidence.fact_ids.is_empty() + })); + assert!( + run.markdown_report() + .contains("Converge handoff: blocked before admission, no facts admitted") + ); + } +} + +#[test] +fn crowd_consensus_usecase_scales_and_compares_with_strict_script() { + let strict = HeadlessCoordinationRun::scripted_burst(); + let full_script = dynamic_user_script(7_001); + let alternate_script = dynamic_user_script(7_002); + + assert_eq!(full_script.len(), DYNAMIC_USER_LIMIT); + assert_ne!( + full_script, alternate_script, + "dynamic user scripts should vary when the seed changes" + ); + + for user_count in CROWD_USER_COUNTS { + let first = CrowdConsensusRun::scripted(user_count, 7_001).expect("supported crowd scale"); + let second = CrowdConsensusRun::scripted(user_count, 7_001).expect("supported crowd scale"); + let varied = CrowdConsensusRun::scripted(user_count, 7_002).expect("supported crowd scale"); + + assert_eq!(first.users, second.users); + assert_eq!(first.event_trace(), second.event_trace()); + assert_ne!(first.event_trace(), varied.event_trace()); + assert_ne!( + first.event_trace(), + strict.event_trace(), + "dynamic crowd run should not collapse into the strict Helm readiness script" + ); + assert_eq!(first.users.len(), user_count); + + let opinion_events: Vec<_> = first + .events + .iter() + .filter(|event| event.kind == CoordinationEventKind::OpinionGathered) + .collect(); + assert_eq!(opinion_events.len(), user_count); + + let speeds: BTreeSet<_> = first + .users + .iter() + .map(|user| user.contribution_speed_ms) + .collect(); + assert!( + speeds.len() > 1, + "crowd run should include different contribution speeds" + ); + assert!(first.fastest_speed_ms() < first.slowest_speed_ms()); + + let opinions: BTreeSet<_> = first + .users + .iter() + .map(|user| user.opinion.as_str().to_string()) + .collect(); + assert!( + opinions.len() > 1, + "crowd run should gather different opinions before measuring consensus" + ); + + let contribution_speeds: Vec<_> = opinion_events + .iter() + .map(|event| { + event + .payload + .get("contribution_speed_ms") + .and_then(Value::as_u64) + .expect("opinion event records contribution speed") + }) + .collect(); + let mut sorted_speeds = contribution_speeds.clone(); + sorted_speeds.sort(); + assert_eq!( + contribution_speeds, sorted_speeds, + "opinion events should replay in contribution-speed order" + ); + + assert!( + first + .events + .iter() + .any(|event| event.kind == CoordinationEventKind::ConsensusMeasured) + ); + assert!( + first + .events + .iter() + .any(|event| event.kind == CoordinationEventKind::ConclusionDrafted) + ); + let decision_event = first + .events + .iter() + .find(|event| event.kind == CoordinationEventKind::DecisionRecorded) + .expect("crowd run records a decision"); + assert_eq!( + decision_event + .payload + .get("authorizes_domain_action") + .and_then(Value::as_bool), + Some(false) + ); + assert!(!first.summary.conclusion.is_empty()); + assert!(!first.summary.decision.is_empty()); + assert_eq!( + first.summary.support_count + first.summary.dissent_count, + user_count + ); + } +} diff --git a/crates/cross-extension-smoke/tests/multi_user_consensus.rs b/crates/cross-extension-smoke/tests/multi_user_consensus.rs new file mode 100644 index 0000000..adfd70a --- /dev/null +++ b/crates/cross-extension-smoke/tests/multi_user_consensus.rs @@ -0,0 +1,145 @@ +//! Deterministic multi-user consensus use case. +//! +//! This proves Arena can start the same static coordination script at 5, 30, +//! and 100 users while preserving contribution speeds, opinion diversity, and +//! deterministic consensus decisions. + +use std::collections::BTreeSet; + +use arena_multi_user_consensus::{ + CONSENSUS_THRESHOLD_BPS, ConsensusEventKind, ConsensusLevel, ConsensusUseCase, + ContributionSpeed, Decision, Opinion, STANDARD_USER_COUNTS, STATIC_USER_COUNT, + STATIC_USERS_100, +}; + +#[test] +fn static_user_script_exposes_one_hundred_predictable_users() { + assert_eq!(STATIC_USERS_100.len(), STATIC_USER_COUNT); + assert_eq!(STATIC_USERS_100[0].id, "user-001"); + assert_eq!(STATIC_USERS_100[99].id, "user-100"); + + let ids: BTreeSet<_> = STATIC_USERS_100.iter().map(|user| user.id).collect(); + assert_eq!(ids.len(), STATIC_USER_COUNT); + + let speeds: BTreeSet<_> = STATIC_USERS_100.iter().map(|user| user.speed).collect(); + assert_eq!( + speeds, + BTreeSet::from([ + ContributionSpeed::Fast, + ContributionSpeed::Steady, + ContributionSpeed::Slow, + ContributionSpeed::Delayed, + ]) + ); + + let opinions: BTreeSet<_> = STATIC_USERS_100.iter().map(|user| user.opinion).collect(); + assert_eq!( + opinions, + BTreeSet::from([ + Opinion::Proceed, + Opinion::ProceedWithGuardrails, + Opinion::Defer, + Opinion::Block, + ]) + ); +} + +#[test] +fn consensus_usecase_starts_at_five_thirty_and_one_hundred_users() { + let starters = [ + (5, ConsensusUseCase::five_users()), + (30, ConsensusUseCase::thirty_users()), + (100, ConsensusUseCase::hundred_users()), + ]; + + for (expected_count, usecase) in starters { + let first = usecase.run(); + let second = ConsensusUseCase::with_user_count(expected_count) + .expect("standard count is supported") + .run(); + + assert_eq!(first.event_trace(), second.event_trace()); + assert_eq!(first.user_count, expected_count); + assert_eq!(first.users.len(), expected_count); + assert_eq!(first.contributions.len(), expected_count); + assert_eq!( + first.events.len(), + expected_count + 5, + "one event per contribution plus the session/mix/consensus/conclusion/decision events" + ); + } +} + +#[test] +fn standard_runs_preserve_parallel_speeds_opinions_and_consensus() { + for user_count in STANDARD_USER_COUNTS { + let run = ConsensusUseCase::with_user_count(user_count) + .expect("standard count is supported") + .run(); + + assert_eq!( + run.events.first().map(|event| event.kind), + Some(ConsensusEventKind::SessionStarted) + ); + assert_eq!( + run.events.last().map(|event| event.kind), + Some(ConsensusEventKind::DecisionRecorded) + ); + assert!( + run.events + .iter() + .skip(1) + .take(user_count) + .all(|event| event.kind == ConsensusEventKind::ContributionAccepted) + ); + + let mut sorted = run.contributions.clone(); + sorted.sort_by(|left, right| { + left.arrival_tick + .cmp(&right.arrival_tick) + .then_with(|| left.participant_id.cmp(right.participant_id)) + }); + assert_eq!(run.contributions, sorted); + assert!( + run.arrival_batches + .iter() + .any(|batch| batch.contribution_count > 1), + "fixture should include parallel arrivals in the same deterministic batch" + ); + assert!( + run.speed_set().len() >= 3, + "standard runs should preserve different contribution speeds" + ); + assert!( + run.opinion_set().len() >= 3, + "standard runs should gather different opinions" + ); + assert!(run.contributions.iter().all(|contribution| { + contribution.signed + && contribution.suggestion_only + && contribution.provenance + == format!("static-roster://{}", contribution.participant_id) + })); + + assert!(matches!( + run.conclusion.level, + ConsensusLevel::Working | ConsensusLevel::Strong + )); + assert!(matches!( + run.conclusion.decision, + Decision::ProceedWithGuardrails + )); + assert!(run.conclusion.consensus_basis_points >= CONSENSUS_THRESHOLD_BPS); + assert!(!run.conclusion.statement.is_empty()); + assert!( + run.conclusion + .next_actions + .contains(&"publish consensus checkpoint") + ); + assert!( + run.conclusion + .next_actions + .contains(&"admit final decision only through Converge") + ); + } +} diff --git a/crates/multi-user-consensus/Cargo.toml b/crates/multi-user-consensus/Cargo.toml new file mode 100644 index 0000000..a13b737 --- /dev/null +++ b/crates/multi-user-consensus/Cargo.toml @@ -0,0 +1,13 @@ +[package] +name = "arena-multi-user-consensus" +version.workspace = true +edition.workspace = true +rust-version.workspace = true +license.workspace = true +publish = false + +[lib] +path = "src/lib.rs" + +[lints] +workspace = true diff --git a/crates/multi-user-consensus/src/lib.rs b/crates/multi-user-consensus/src/lib.rs new file mode 100644 index 0000000..8620288 --- /dev/null +++ b/crates/multi-user-consensus/src/lib.rs @@ -0,0 +1,1056 @@ +//! Deterministic multi-user consensus fixture for Arena tests. +//! +//! The use case models a large coordination room where users contribute at +//! different speeds, offer different opinions, and converge to a decision. +//! It is intentionally dependency-free and static so tests can replay 5, 30, +//! and 100 user sessions without wall-clock timing or generated identities. + +use std::collections::{BTreeMap, BTreeSet}; + +pub const STATIC_USER_COUNT: usize = 100; +pub const STANDARD_USER_COUNTS: [usize; 3] = [5, 30, 100]; +pub const CONSENSUS_THRESHOLD_BPS: u16 = 6_000; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub enum ContributionSpeed { + Fast, + Steady, + Slow, + Delayed, +} + +impl ContributionSpeed { + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + ContributionSpeed::Fast => "fast", + ContributionSpeed::Steady => "steady", + ContributionSpeed::Slow => "slow", + ContributionSpeed::Delayed => "delayed", + } + } + + const fn base_tick(self) -> u64 { + match self { + ContributionSpeed::Fast => 0, + ContributionSpeed::Steady => 4, + ContributionSpeed::Slow => 8, + ContributionSpeed::Delayed => 12, + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +pub enum Opinion { + Proceed, + ProceedWithGuardrails, + Defer, + Block, +} + +impl Opinion { + #[must_use] + pub const fn as_str(self) -> &'static str { + match self { + Opinion::Proceed => "proceed", + Opinion::ProceedWithGuardrails => "proceed-with-guardrails", + Opinion::Defer => "defer", + Opinion::Block => "block", + } + } + + const fn supports_decision(self) -> bool { + matches!(self, Opinion::Proceed | Opinion::ProceedWithGuardrails) + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct StaticUser { + pub id: &'static str, + pub display_name: &'static str, + pub role: &'static str, + pub speed: ContributionSpeed, + pub opinion: Opinion, + pub rationale: &'static str, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Contribution { + pub sequence: u64, + pub arrival_tick: u64, + pub arrival_batch: u64, + pub ordering_key: String, + pub participant_id: &'static str, + pub participant_role: &'static str, + pub speed: ContributionSpeed, + pub opinion: Opinion, + pub signed: bool, + pub suggestion_only: bool, + pub provenance: String, + pub body: String, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ConsensusLevel { + Strong, + Working, + NoConsensus, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Decision { + ProceedWithGuardrails, + ReopenForMoreInput, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ConsensusConclusion { + pub decision: Decision, + pub level: ConsensusLevel, + pub support_count: usize, + pub concern_count: usize, + pub block_count: usize, + pub consensus_basis_points: u16, + pub statement: String, + pub next_actions: Vec<&'static str>, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ConsensusEventKind { + SessionStarted, + ContributionAccepted, + OpinionsMixed, + ConsensusCalculated, + ConclusionDrafted, + DecisionRecorded, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ConsensusEvent { + pub sequence: u64, + pub kind: ConsensusEventKind, + pub payload: String, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ArrivalBatch { + pub arrival_tick: u64, + pub contribution_count: usize, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ConsensusRun { + pub user_count: usize, + pub users: Vec, + pub contributions: Vec, + pub arrival_batches: Vec, + pub opinion_counts: BTreeMap, + pub conclusion: ConsensusConclusion, + pub events: Vec, +} + +impl ConsensusRun { + #[must_use] + pub fn event_trace(&self) -> Vec<(u64, ConsensusEventKind, String)> { + self.events + .iter() + .map(|event| (event.sequence, event.kind, event.payload.clone())) + .collect() + } + + #[must_use] + pub fn speed_set(&self) -> BTreeSet { + self.contributions + .iter() + .map(|contribution| contribution.speed) + .collect() + } + + #[must_use] + pub fn opinion_set(&self) -> BTreeSet { + self.contributions + .iter() + .map(|contribution| contribution.opinion) + .collect() + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub struct ConsensusUseCase { + user_count: usize, +} + +impl ConsensusUseCase { + pub fn with_user_count(user_count: usize) -> Result { + if !(1..=STATIC_USER_COUNT).contains(&user_count) { + return Err(format!( + "user_count must be between 1 and {STATIC_USER_COUNT}; got {user_count}" + )); + } + Ok(Self { user_count }) + } + + #[must_use] + pub fn five_users() -> Self { + Self::with_user_count(5).expect("5 users is a supported consensus scale") + } + + #[must_use] + pub fn thirty_users() -> Self { + Self::with_user_count(30).expect("30 users is a supported consensus scale") + } + + #[must_use] + pub fn hundred_users() -> Self { + Self::with_user_count(100).expect("100 users is a supported consensus scale") + } + + #[must_use] + pub fn run(self) -> ConsensusRun { + let users = STATIC_USERS_100[..self.user_count].to_vec(); + let contributions = contributions_for(&users); + let arrival_batches = arrival_batches_for(&contributions); + let opinion_counts = opinion_counts_for(&contributions); + let conclusion = conclusion_for(self.user_count, &opinion_counts); + let events = events_for( + self.user_count, + &contributions, + &opinion_counts, + &conclusion, + ); + + ConsensusRun { + user_count: self.user_count, + users, + contributions, + arrival_batches, + opinion_counts, + conclusion, + events, + } + } +} + +#[derive(Clone)] +struct ScheduledContribution { + arrival_tick: u64, + ordering_key: String, + user: StaticUser, +} + +fn contributions_for(users: &[StaticUser]) -> Vec { + let mut scheduled: Vec<_> = users + .iter() + .copied() + .enumerate() + .map(|(index, user)| { + let arrival_tick = user.speed.base_tick() + (index as u64 % 3); + ScheduledContribution { + arrival_tick, + ordering_key: format!("{arrival_tick:03}:{}", user.id), + user, + } + }) + .collect(); + + scheduled.sort_by(|left, right| { + left.arrival_tick + .cmp(&right.arrival_tick) + .then_with(|| left.user.id.cmp(right.user.id)) + }); + + scheduled + .into_iter() + .enumerate() + .map(|(index, scheduled)| Contribution { + sequence: index as u64 + 1, + arrival_tick: scheduled.arrival_tick, + arrival_batch: scheduled.arrival_tick, + ordering_key: scheduled.ordering_key, + participant_id: scheduled.user.id, + participant_role: scheduled.user.role, + speed: scheduled.user.speed, + opinion: scheduled.user.opinion, + signed: true, + suggestion_only: true, + provenance: format!("static-roster://{}", scheduled.user.id), + body: format!( + "{} contributes {}: {}", + scheduled.user.display_name, + scheduled.user.opinion.as_str(), + scheduled.user.rationale + ), + }) + .collect() +} + +fn arrival_batches_for(contributions: &[Contribution]) -> Vec { + let mut counts = BTreeMap::::new(); + for contribution in contributions { + *counts.entry(contribution.arrival_batch).or_default() += 1; + } + counts + .into_iter() + .map(|(arrival_tick, contribution_count)| ArrivalBatch { + arrival_tick, + contribution_count, + }) + .collect() +} + +fn opinion_counts_for(contributions: &[Contribution]) -> BTreeMap { + let mut counts = BTreeMap::new(); + for contribution in contributions { + *counts.entry(contribution.opinion).or_default() += 1; + } + counts +} + +fn conclusion_for( + user_count: usize, + opinion_counts: &BTreeMap, +) -> ConsensusConclusion { + let support_count = opinion_counts + .iter() + .filter_map(|(opinion, count)| opinion.supports_decision().then_some(*count)) + .sum::(); + let concern_count = opinion_counts.get(&Opinion::Defer).copied().unwrap_or(0); + let block_count = opinion_counts.get(&Opinion::Block).copied().unwrap_or(0); + let consensus_basis_points = ((support_count * 10_000) / user_count) as u16; + let level = if consensus_basis_points >= 7_500 { + ConsensusLevel::Strong + } else if consensus_basis_points >= CONSENSUS_THRESHOLD_BPS { + ConsensusLevel::Working + } else { + ConsensusLevel::NoConsensus + }; + let decision = if matches!(level, ConsensusLevel::Strong | ConsensusLevel::Working) { + Decision::ProceedWithGuardrails + } else { + Decision::ReopenForMoreInput + }; + + ConsensusConclusion { + decision, + level, + support_count, + concern_count, + block_count, + consensus_basis_points, + statement: format!( + "{support_count}/{user_count} users support proceeding; concerns remain from {concern_count} defer and {block_count} block opinions." + ), + next_actions: match decision { + Decision::ProceedWithGuardrails => vec![ + "publish consensus checkpoint", + "carry blockers as explicit gate follow-up", + "admit final decision only through Converge", + ], + Decision::ReopenForMoreInput => vec![ + "open another contribution round", + "assign probes to unresolved objections", + ], + }, + } +} + +fn events_for( + user_count: usize, + contributions: &[Contribution], + opinion_counts: &BTreeMap, + conclusion: &ConsensusConclusion, +) -> Vec { + let mut events = Vec::with_capacity(contributions.len() + 5); + push_event( + &mut events, + ConsensusEventKind::SessionStarted, + format!("users={user_count};topic=multi-user-consensus"), + ); + for contribution in contributions { + push_event( + &mut events, + ConsensusEventKind::ContributionAccepted, + format!( + "participant={};tick={};speed={};opinion={};suggestion_only={}", + contribution.participant_id, + contribution.arrival_tick, + contribution.speed.as_str(), + contribution.opinion.as_str(), + contribution.suggestion_only + ), + ); + } + push_event( + &mut events, + ConsensusEventKind::OpinionsMixed, + format!( + "proceed={};guardrails={};defer={};block={}", + opinion_counts.get(&Opinion::Proceed).copied().unwrap_or(0), + opinion_counts + .get(&Opinion::ProceedWithGuardrails) + .copied() + .unwrap_or(0), + opinion_counts.get(&Opinion::Defer).copied().unwrap_or(0), + opinion_counts.get(&Opinion::Block).copied().unwrap_or(0) + ), + ); + push_event( + &mut events, + ConsensusEventKind::ConsensusCalculated, + format!( + "support={};bps={};level={:?}", + conclusion.support_count, conclusion.consensus_basis_points, conclusion.level + ), + ); + push_event( + &mut events, + ConsensusEventKind::ConclusionDrafted, + conclusion.statement.clone(), + ); + push_event( + &mut events, + ConsensusEventKind::DecisionRecorded, + format!("decision={:?}", conclusion.decision), + ); + events +} + +fn push_event(events: &mut Vec, kind: ConsensusEventKind, payload: String) { + events.push(ConsensusEvent { + sequence: events.len() as u64 + 1, + kind, + payload, + }); +} + +macro_rules! user { + ($num:literal, $role:literal, $speed:ident, $opinion:ident, $rationale:literal) => { + StaticUser { + id: concat!("user-", $num), + display_name: concat!("Static User ", $num), + role: $role, + speed: ContributionSpeed::$speed, + opinion: Opinion::$opinion, + rationale: $rationale, + } + }; +} + +pub const STATIC_USERS_100: [StaticUser; STATIC_USER_COUNT] = [ + user!( + "001", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "002", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("003", "risk", Slow, Defer, "needs rollback detail"), + user!( + "004", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "005", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "006", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "007", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("008", "legal", Delayed, Defer, "needs sign-off language"), + user!("009", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "010", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "011", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "012", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("013", "risk", Slow, Defer, "needs rollback detail"), + user!( + "014", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "015", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "016", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "017", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("018", "legal", Delayed, Defer, "needs sign-off language"), + user!("019", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "020", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "021", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "022", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("023", "risk", Slow, Defer, "needs rollback detail"), + user!( + "024", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "025", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "026", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "027", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("028", "legal", Delayed, Defer, "needs sign-off language"), + user!("029", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "030", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "031", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "032", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("033", "risk", Slow, Defer, "needs rollback detail"), + user!( + "034", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "035", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "036", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "037", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("038", "legal", Delayed, Defer, "needs sign-off language"), + user!("039", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "040", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "041", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "042", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("043", "risk", Slow, Defer, "needs rollback detail"), + user!( + "044", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "045", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "046", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "047", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("048", "legal", Delayed, Defer, "needs sign-off language"), + user!("049", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "050", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "051", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "052", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("053", "risk", Slow, Defer, "needs rollback detail"), + user!( + "054", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "055", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "056", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "057", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("058", "legal", Delayed, Defer, "needs sign-off language"), + user!("059", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "060", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "061", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "062", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("063", "risk", Slow, Defer, "needs rollback detail"), + user!( + "064", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "065", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "066", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "067", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("068", "legal", Delayed, Defer, "needs sign-off language"), + user!("069", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "070", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "071", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "072", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("073", "risk", Slow, Defer, "needs rollback detail"), + user!( + "074", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "075", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "076", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "077", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("078", "legal", Delayed, Defer, "needs sign-off language"), + user!("079", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "080", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "081", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "082", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("083", "risk", Slow, Defer, "needs rollback detail"), + user!( + "084", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "085", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "086", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "087", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("088", "legal", Delayed, Defer, "needs sign-off language"), + user!("089", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "090", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), + user!( + "091", + "strategy", + Fast, + ProceedWithGuardrails, + "support with explicit checkpoints" + ), + user!( + "092", + "operations", + Steady, + Proceed, + "capacity is available" + ), + user!("093", "risk", Slow, Defer, "needs rollback detail"), + user!( + "094", + "customer", + Fast, + ProceedWithGuardrails, + "customer communications need an owner" + ), + user!( + "095", + "security", + Delayed, + Block, + "block until access review is complete" + ), + user!( + "096", + "finance", + Steady, + Proceed, + "budget impact is acceptable" + ), + user!( + "097", + "support", + Slow, + ProceedWithGuardrails, + "support playbook is required" + ), + user!("098", "legal", Delayed, Defer, "needs sign-off language"), + user!("099", "product", Fast, Proceed, "unblocks roadmap decision"), + user!( + "100", + "engineering", + Steady, + ProceedWithGuardrails, + "requires rollout guardrails" + ), +]; + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn static_script_has_one_hundred_predictable_users() { + assert_eq!(STATIC_USERS_100.len(), STATIC_USER_COUNT); + assert_eq!(STATIC_USERS_100[0].id, "user-001"); + assert_eq!(STATIC_USERS_100[99].id, "user-100"); + + let ids: BTreeSet<_> = STATIC_USERS_100.iter().map(|user| user.id).collect(); + assert_eq!(ids.len(), STATIC_USER_COUNT); + + let speeds: BTreeSet<_> = STATIC_USERS_100.iter().map(|user| user.speed).collect(); + assert_eq!( + speeds, + BTreeSet::from([ + ContributionSpeed::Fast, + ContributionSpeed::Steady, + ContributionSpeed::Slow, + ContributionSpeed::Delayed, + ]) + ); + + let opinions: BTreeSet<_> = STATIC_USERS_100.iter().map(|user| user.opinion).collect(); + assert_eq!( + opinions, + BTreeSet::from([ + Opinion::Proceed, + Opinion::ProceedWithGuardrails, + Opinion::Defer, + Opinion::Block, + ]) + ); + } + + #[test] + fn standard_runs_are_deterministic() { + for user_count in STANDARD_USER_COUNTS { + let first = ConsensusUseCase::with_user_count(user_count) + .expect("standard count is supported") + .run(); + let second = ConsensusUseCase::with_user_count(user_count) + .expect("standard count is supported") + .run(); + + assert_eq!(first, second); + assert_eq!(first.user_count, user_count); + assert_eq!(first.contributions.len(), user_count); + assert!(matches!( + first.conclusion.decision, + Decision::ProceedWithGuardrails + )); + } + } + + #[test] + fn standard_runs_preserve_parallelism_and_consensus() { + for user_count in STANDARD_USER_COUNTS { + let run = ConsensusUseCase::with_user_count(user_count) + .expect("standard count is supported") + .run(); + + assert!( + run.arrival_batches + .iter() + .any(|batch| batch.contribution_count > 1), + "fixture should include deterministic same-tick arrivals" + ); + assert!(run.speed_set().len() >= 3); + assert!(run.opinion_set().len() >= 3); + assert!(run.contributions.iter().all(|contribution| { + contribution.signed + && contribution.suggestion_only + && contribution.provenance + == format!("static-roster://{}", contribution.participant_id) + })); + assert!(matches!( + run.conclusion.level, + ConsensusLevel::Working | ConsensusLevel::Strong + )); + assert!(run.conclusion.consensus_basis_points >= CONSENSUS_THRESHOLD_BPS); + assert!( + run.conclusion + .next_actions + .contains(&"publish consensus checkpoint") + ); + assert!( + run.conclusion + .next_actions + .contains(&"admit final decision only through Converge") + ); + } + } +} From 2ecce5cadea9667207fd38821b297bcf141535da Mon Sep 17 00:00:00 2001 From: Kenneth Pernyer Date: Thu, 2 Jul 2026 10:56:29 +0200 Subject: [PATCH 2/2] ci: canonical just ci + first CI workflow RP-CI-PARITY: add the canonical `just ci` aggregate (fmt-check, check, lint, test) and a thin single-job GitHub Actions workflow that runs exactly `just ci`, so the local verdict equals the CI verdict. - Justfile: add `ci`, `check`, and `test` recipes; normalize `fmt-check`/`fmt` to `--all` form. `report`, `build`, `contracts`, and `test-metrics` are preserved. - .github/workflows/ci.yml: checkout + sibling checkout + zsh (the Justfile shell) + pinned Rust 1.96.0 toolchain + rust-cache + just. - scripts/ci/checkout-reflective-siblings.sh: clone the Reflective-Lab sibling repos this workspace consumes via relative path deps, mirroring the local reflective/ tree topology (adapted from arbiter-policy). - .github/dependabot.yml: weekly cargo + github-actions updates. Co-Authored-By: Claude Fable 5 --- .github/dependabot.yml | 10 ++++ .github/workflows/ci.yml | 33 ++++++++++++ Justfile | 15 +++++- scripts/ci/checkout-reflective-siblings.sh | 58 ++++++++++++++++++++++ 4 files changed, 114 insertions(+), 2 deletions(-) create mode 100644 .github/dependabot.yml create mode 100644 .github/workflows/ci.yml create mode 100755 scripts/ci/checkout-reflective-siblings.sh diff --git a/.github/dependabot.yml b/.github/dependabot.yml new file mode 100644 index 0000000..c6d4824 --- /dev/null +++ b/.github/dependabot.yml @@ -0,0 +1,10 @@ +version: 2 +updates: + - package-ecosystem: cargo + directory: "/" + schedule: + interval: weekly + - package-ecosystem: github-actions + directory: "/" + schedule: + interval: weekly diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..0bf2019 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,33 @@ +name: CI + +on: + push: + branches: ["main"] + pull_request: + branches: ["main"] + workflow_dispatch: + +env: + CARGO_TERM_COLOR: always + RUST_VERSION: "1.96.0" + +jobs: + ci: + name: just ci + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v6 + - name: Checkout Reflective sibling dependencies + run: bash scripts/ci/checkout-reflective-siblings.sh + # zsh is the Justfile shell (`set shell`). + - name: Install system dependencies + run: | + sudo apt-get update + sudo apt-get install -y zsh + - uses: dtolnay/rust-toolchain@stable + with: + toolchain: ${{ env.RUST_VERSION }} + components: clippy, rustfmt + - uses: Swatinem/rust-cache@v2 + - uses: extractions/setup-just@v3 + - run: just ci diff --git a/Justfile b/Justfile index 2459531..740d896 100644 --- a/Justfile +++ b/Justfile @@ -3,6 +3,9 @@ set shell := ["zsh", "-cu"] # Run every quality dimension and print the scoreboard. default: report +# Canonical CI aggregate (RP-CI-PARITY): CI runs exactly `just ci`. +ci: fmt-check check lint test + # Print the scoreboard for every default dimension. report: cargo run --quiet --bin arena -- report @@ -24,6 +27,14 @@ contracts: build: cargo build --workspace --all-targets +# Type-check the workspace. +check: + cargo check --workspace --all-targets + +# Run the full test suite. +test: + cargo test --workspace --all-targets + # Run the arena-metrics unit tests (aggregation/Verdict invariants). test-metrics: cargo test -p arena-metrics --all-targets @@ -34,8 +45,8 @@ lint: # Format check. fmt-check: - cargo fmt --check + cargo fmt --all -- --check # rustfmt in place. fmt: - cargo fmt + cargo fmt --all diff --git a/scripts/ci/checkout-reflective-siblings.sh b/scripts/ci/checkout-reflective-siblings.sh new file mode 100755 index 0000000..d67d06f --- /dev/null +++ b/scripts/ci/checkout-reflective-siblings.sh @@ -0,0 +1,58 @@ +#!/usr/bin/env bash +# Check out the Reflective-Lab sibling repos that arena-tests consumes via +# relative path dependencies (see Cargo.toml [workspace.dependencies] and +# [patch.crates-io], plus the transitive closure through helms, +# atelier-showcase, and runtime-runway). +# +# Local layout (arena-tests lives at reflective/arena-tests): +# ../bedrock-platform/ -> platform repos +# ../mosaic-extensions/ -> extension repos +# ../ -> reflective-root siblings +# +# In CI, GITHUB_WORKSPACE (/home/runner/work/arena-tests/arena-tests) plays +# the role of reflective/arena-tests, so its parent acts as the reflective +# root and the topology below mirrors the local tree exactly. Adapted from +# mosaic-extensions/arbiter-policy/scripts/ci/checkout-reflective-siblings.sh. +set -euo pipefail + +workspace="${GITHUB_WORKSPACE:-$(git rev-parse --show-toplevel)}" + +checkout_reflective_repo() { + local repo="$1" + local relative_path="$2" + local dest="${workspace}/${relative_path}" + + if [[ -d "$dest/.git" ]]; then + echo "ok: ${relative_path} already checked out" + return + fi + + if [[ -e "$dest" ]]; then + echo "error: ${dest} exists but is not a git checkout" >&2 + exit 1 + fi + + mkdir -p "$(dirname "$dest")" + echo "==> checkout Reflective-Lab/${repo} -> ${relative_path}" + GIT_TERMINAL_PROMPT=0 git clone --depth=1 --quiet "https://github.com/Reflective-Lab/${repo}.git" "$dest" +} + +# Platform repos. +checkout_reflective_repo axiom ../bedrock-platform/axiom +checkout_reflective_repo converge ../bedrock-platform/converge +checkout_reflective_repo helms ../bedrock-platform/helms +checkout_reflective_repo organism ../bedrock-platform/organism + +# Extension repos. +checkout_reflective_repo arbiter-policy ../mosaic-extensions/arbiter-policy +checkout_reflective_repo embassy-ports ../mosaic-extensions/embassy-ports +checkout_reflective_repo ferrox-solvers ../mosaic-extensions/ferrox-solvers +checkout_reflective_repo manifold-adapters ../mosaic-extensions/manifold-adapters +checkout_reflective_repo mnemos-knowledge ../mosaic-extensions/mnemos-knowledge +checkout_reflective_repo prism-analytics ../mosaic-extensions/prism-analytics + +# Reflective-root siblings. +checkout_reflective_repo atelier-showcase ../atelier-showcase +checkout_reflective_repo runtime-runway ../runtime-runway +# commerce-rails backs runtime-runway's workspace.dependencies entry. +checkout_reflective_repo commerce-rails ../commerce-rails