1use std::collections::{HashMap, HashSet};
4use std::sync::{Arc, LazyLock};
5
6use arrow::array::{
7 Array, FixedSizeBinaryArray, FixedSizeBinaryBuilder, StringArray, StringBuilder,
8 TimestampMicrosecondArray, TimestampMicrosecondBuilder, UInt32Array, UInt32Builder,
9};
10use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
11use arrow::record_batch::RecordBatch;
12use graphforge_core::canonical::{
13 CANONICAL_CONTRACT_VERSION, CanonicalDomain, CanonicalWriter, fingerprint,
14};
15use unicode_normalization::UnicodeNormalization;
16use uuid::{Uuid, Version};
17
18use crate::{
19 EPISTEMIC_CAPABILITY_VERSION, KnowledgeError, MAX_KNOWLEDGE_ROWS, SchemaRegistryEntry,
20};
21
22pub const HYPOTHESIS_GROUP_CONTRACT_VERSION: u32 = 1;
24pub const HYPOTHESIS_MEMBERSHIP_CONTRACT_VERSION: u32 = 1;
26pub const HYPOTHESIS_SELECTION_CONTRACT_VERSION: u32 = 1;
28pub const HYPOTHESIS_KEY_POLICY_VERSION: u32 = 1;
30pub const HYPOTHESIS_STATE_POLICY_VERSION: u32 = 1;
32pub const MAX_HYPOTHESIS_QUESTION_KEY_BYTES: usize = 1_024;
34
35pub static HYPOTHESIS_GROUP_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
37 Arc::new(Schema::new(vec![
38 uuid_field("group_uuid", false),
39 Field::new("question_key", DataType::Utf8, false),
40 uuid_field("provenance_uuid", false),
41 timestamp_field("recorded_at"),
42 Field::new("contract_version", DataType::UInt32, false),
43 ]))
44});
45
46pub static HYPOTHESIS_MEMBERSHIP_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
48 Arc::new(Schema::new(vec![
49 uuid_field("membership_event_uuid", false),
50 uuid_field("operation_uuid", false),
51 uuid_field("group_uuid", false),
52 uuid_field("assertion_uuid", false),
53 Field::new("action", DataType::Utf8, false),
54 uuid_field("reasoning_uuid", false),
55 uuid_field("provenance_uuid", false),
56 timestamp_field("recorded_at"),
57 Field::new("contract_version", DataType::UInt32, false),
58 ]))
59});
60
61pub static HYPOTHESIS_SELECTION_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
63 Arc::new(Schema::new(vec![
64 uuid_field("selection_event_uuid", false),
65 uuid_field("operation_uuid", false),
66 uuid_field("group_uuid", false),
67 uuid_field("selected_assertion_uuid", true),
68 uuid_field("reasoning_uuid", false),
69 uuid_field("provenance_uuid", false),
70 timestamp_field("recorded_at"),
71 Field::new("contract_version", DataType::UInt32, false),
72 ]))
73});
74
75static GROUP_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
76 fingerprint(
77 CanonicalDomain::Schema,
78 CANONICAL_CONTRACT_VERSION,
79 b"hypothesis_group/1|group_uuid:fixed[16]:required|question_key:utf8:required|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
80 )
81 .expect("registered hypothesis-group schema is within canonical bounds")
82});
83
84static MEMBERSHIP_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
85 fingerprint(
86 CanonicalDomain::Schema,
87 CANONICAL_CONTRACT_VERSION,
88 b"hypothesis_membership/1|membership_event_uuid:fixed[16]:required|operation_uuid:fixed[16]:required|group_uuid:fixed[16]:required|assertion_uuid:fixed[16]:required|action:utf8:required|reasoning_uuid:fixed[16]:required|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
89 )
90 .expect("registered hypothesis-membership schema is within canonical bounds")
91});
92
93static SELECTION_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
94 fingerprint(
95 CanonicalDomain::Schema,
96 CANONICAL_CONTRACT_VERSION,
97 b"hypothesis_selection/1|selection_event_uuid:fixed[16]:required|operation_uuid:fixed[16]:required|group_uuid:fixed[16]:required|selected_assertion_uuid:fixed[16]:nullable|reasoning_uuid:fixed[16]:required|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
98 )
99 .expect("registered hypothesis-selection schema is within canonical bounds")
100});
101
102#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
104pub enum HypothesisMembershipAction {
105 Added,
107 Removed,
109}
110
111impl HypothesisMembershipAction {
112 #[must_use]
114 pub const fn as_str(self) -> &'static str {
115 match self {
116 Self::Added => "added",
117 Self::Removed => "removed",
118 }
119 }
120
121 fn parse(value: &str) -> Result<Self, KnowledgeError> {
122 match value {
123 "added" => Ok(Self::Added),
124 "removed" => Ok(Self::Removed),
125 _ => Err(invalid(
126 "hypothesis_membership.action",
127 "unknown registry value",
128 )),
129 }
130 }
131}
132
133#[derive(Clone, Debug, Eq, PartialEq)]
135pub struct HypothesisGroup {
136 pub group_uuid: Uuid,
138 pub question_key: String,
140 pub provenance_uuid: Uuid,
142 pub recorded_at_micros: i64,
144 pub contract_version: u32,
146}
147
148impl HypothesisGroup {
149 pub fn new(
151 group_uuid: Uuid,
152 question_key: String,
153 provenance_uuid: Uuid,
154 recorded_at_micros: i64,
155 ) -> Result<Self, KnowledgeError> {
156 require_v7(group_uuid, "group_uuid")?;
157 require_uuid(provenance_uuid, "provenance_uuid")?;
158 validate_question_key(&question_key)?;
159 Ok(Self {
160 group_uuid,
161 question_key,
162 provenance_uuid,
163 recorded_at_micros,
164 contract_version: HYPOTHESIS_GROUP_CONTRACT_VERSION,
165 })
166 }
167}
168
169#[derive(Clone, Debug, Eq, PartialEq)]
171pub struct HypothesisMembershipEvent {
172 pub membership_event_uuid: Uuid,
174 pub operation_uuid: Uuid,
176 pub group_uuid: Uuid,
178 pub assertion_uuid: Uuid,
180 pub action: HypothesisMembershipAction,
182 pub reasoning_uuid: Uuid,
184 pub provenance_uuid: Uuid,
186 pub recorded_at_micros: i64,
188 pub contract_version: u32,
190}
191
192impl HypothesisMembershipEvent {
193 #[allow(clippy::too_many_arguments)]
195 pub fn new(
196 membership_event_uuid: Uuid,
197 operation_uuid: Uuid,
198 group_uuid: Uuid,
199 assertion_uuid: Uuid,
200 action: HypothesisMembershipAction,
201 reasoning_uuid: Uuid,
202 provenance_uuid: Uuid,
203 recorded_at_micros: i64,
204 ) -> Result<Self, KnowledgeError> {
205 for (uuid, field) in [
206 (membership_event_uuid, "membership_event_uuid"),
207 (operation_uuid, "operation_uuid"),
208 (group_uuid, "group_uuid"),
209 (assertion_uuid, "assertion_uuid"),
210 (reasoning_uuid, "reasoning_uuid"),
211 ] {
212 require_v7(uuid, field)?;
213 }
214 require_uuid(provenance_uuid, "provenance_uuid")?;
215 Ok(Self {
216 membership_event_uuid,
217 operation_uuid,
218 group_uuid,
219 assertion_uuid,
220 action,
221 reasoning_uuid,
222 provenance_uuid,
223 recorded_at_micros,
224 contract_version: HYPOTHESIS_MEMBERSHIP_CONTRACT_VERSION,
225 })
226 }
227}
228
229#[derive(Clone, Debug, Eq, PartialEq)]
231pub struct HypothesisSelectionEvent {
232 pub selection_event_uuid: Uuid,
234 pub operation_uuid: Uuid,
236 pub group_uuid: Uuid,
238 pub selected_assertion_uuid: Option<Uuid>,
240 pub reasoning_uuid: Uuid,
242 pub provenance_uuid: Uuid,
244 pub recorded_at_micros: i64,
246 pub contract_version: u32,
248}
249
250impl HypothesisSelectionEvent {
251 #[allow(clippy::too_many_arguments)]
253 pub fn new(
254 selection_event_uuid: Uuid,
255 operation_uuid: Uuid,
256 group_uuid: Uuid,
257 selected_assertion_uuid: Option<Uuid>,
258 reasoning_uuid: Uuid,
259 provenance_uuid: Uuid,
260 recorded_at_micros: i64,
261 ) -> Result<Self, KnowledgeError> {
262 for (uuid, field) in [
263 (selection_event_uuid, "selection_event_uuid"),
264 (operation_uuid, "operation_uuid"),
265 (group_uuid, "group_uuid"),
266 (reasoning_uuid, "reasoning_uuid"),
267 ] {
268 require_v7(uuid, field)?;
269 }
270 if let Some(assertion_uuid) = selected_assertion_uuid {
271 require_v7(assertion_uuid, "selected_assertion_uuid")?;
272 }
273 require_uuid(provenance_uuid, "provenance_uuid")?;
274 Ok(Self {
275 selection_event_uuid,
276 operation_uuid,
277 group_uuid,
278 selected_assertion_uuid,
279 reasoning_uuid,
280 provenance_uuid,
281 recorded_at_micros,
282 contract_version: HYPOTHESIS_SELECTION_CONTRACT_VERSION,
283 })
284 }
285}
286
287#[derive(Clone, Debug, Default, Eq, PartialEq)]
289pub struct HypothesisLedger {
290 groups: Vec<HypothesisGroup>,
291 membership_events: Vec<HypothesisMembershipEvent>,
292 selection_events: Vec<HypothesisSelectionEvent>,
293}
294
295impl HypothesisLedger {
296 pub fn new(
298 mut groups: Vec<HypothesisGroup>,
299 mut membership_events: Vec<HypothesisMembershipEvent>,
300 mut selection_events: Vec<HypothesisSelectionEvent>,
301 ) -> Result<Self, KnowledgeError> {
302 for (participant, observed) in [
303 ("hypothesis_groups", groups.len()),
304 ("hypothesis_membership_events", membership_events.len()),
305 ("hypothesis_selection_events", selection_events.len()),
306 ] {
307 if observed > MAX_KNOWLEDGE_ROWS {
308 return Err(KnowledgeError::Limit {
309 participant,
310 observed,
311 limit: MAX_KNOWLEDGE_ROWS,
312 });
313 }
314 }
315 validate_groups(&groups)?;
316 validate_event_ids(&membership_events, &selection_events)?;
317 groups.sort_by_key(|row| (row.recorded_at_micros, row.group_uuid));
318 membership_events.sort_by_key(|row| (row.recorded_at_micros, row.membership_event_uuid));
319 selection_events.sort_by_key(|row| (row.recorded_at_micros, row.selection_event_uuid));
320 validate_state(&groups, &membership_events, &selection_events)?;
321 Ok(Self {
322 groups,
323 membership_events,
324 selection_events,
325 })
326 }
327
328 #[must_use]
329 pub fn groups(&self) -> &[HypothesisGroup] {
331 &self.groups
332 }
333
334 #[must_use]
335 pub fn membership_events(&self) -> &[HypothesisMembershipEvent] {
337 &self.membership_events
338 }
339
340 #[must_use]
341 pub fn selection_events(&self) -> &[HypothesisSelectionEvent] {
343 &self.selection_events
344 }
345
346 pub fn group_fingerprint(&self, group_uuid: Uuid) -> Result<[u8; 32], KnowledgeError> {
348 let row = self
349 .groups
350 .iter()
351 .find(|row| row.group_uuid == group_uuid)
352 .ok_or(KnowledgeError::Dangling("group_uuid"))?;
353 let mut writer = CanonicalWriter::new();
354 writer.raw(row.group_uuid.as_bytes())?;
355 writer.text(&row.question_key)?;
356 writer.raw(row.provenance_uuid.as_bytes())?;
357 writer.i64(row.recorded_at_micros)?;
358 writer.u32(row.contract_version)?;
359 fingerprint(
360 CanonicalDomain::HypothesisGroup,
361 CANONICAL_CONTRACT_VERSION,
362 &writer.finish(),
363 )
364 .map_err(Into::into)
365 }
366
367 pub fn membership_fingerprint(
369 &self,
370 membership_event_uuid: Uuid,
371 ) -> Result<[u8; 32], KnowledgeError> {
372 let row = self
373 .membership_events
374 .iter()
375 .find(|row| row.membership_event_uuid == membership_event_uuid)
376 .ok_or(KnowledgeError::Dangling("membership_event_uuid"))?;
377 let mut writer = CanonicalWriter::new();
378 for value in [
379 row.membership_event_uuid,
380 row.operation_uuid,
381 row.group_uuid,
382 row.assertion_uuid,
383 ] {
384 writer.raw(value.as_bytes())?;
385 }
386 writer.text(row.action.as_str())?;
387 writer.raw(row.reasoning_uuid.as_bytes())?;
388 writer.raw(row.provenance_uuid.as_bytes())?;
389 writer.i64(row.recorded_at_micros)?;
390 writer.u32(row.contract_version)?;
391 fingerprint(
392 CanonicalDomain::HypothesisMembership,
393 CANONICAL_CONTRACT_VERSION,
394 &writer.finish(),
395 )
396 .map_err(Into::into)
397 }
398
399 pub fn selection_fingerprint(
401 &self,
402 selection_event_uuid: Uuid,
403 ) -> Result<[u8; 32], KnowledgeError> {
404 let row = self
405 .selection_events
406 .iter()
407 .find(|row| row.selection_event_uuid == selection_event_uuid)
408 .ok_or(KnowledgeError::Dangling("selection_event_uuid"))?;
409 let mut writer = CanonicalWriter::new();
410 for value in [row.selection_event_uuid, row.operation_uuid, row.group_uuid] {
411 writer.raw(value.as_bytes())?;
412 }
413 match row.selected_assertion_uuid {
414 Some(value) => {
415 writer.u8(1)?;
416 writer.raw(value.as_bytes())?;
417 }
418 None => writer.u8(0)?,
419 }
420 writer.raw(row.reasoning_uuid.as_bytes())?;
421 writer.raw(row.provenance_uuid.as_bytes())?;
422 writer.i64(row.recorded_at_micros)?;
423 writer.u32(row.contract_version)?;
424 fingerprint(
425 CanonicalDomain::HypothesisSelection,
426 CANONICAL_CONTRACT_VERSION,
427 &writer.finish(),
428 )
429 .map_err(Into::into)
430 }
431
432 pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
434 Self::new(
435 merge_rows(
436 &self.groups,
437 &staged.groups,
438 |row| row.group_uuid,
439 "group_uuid",
440 )?,
441 merge_rows(
442 &self.membership_events,
443 &staged.membership_events,
444 |row| row.membership_event_uuid,
445 "membership_event_uuid",
446 )?,
447 merge_rows(
448 &self.selection_events,
449 &staged.selection_events,
450 |row| row.selection_event_uuid,
451 "selection_event_uuid",
452 )?,
453 )
454 }
455
456 pub fn group_batch(&self) -> Result<RecordBatch, KnowledgeError> {
458 let len = self.groups.len();
459 let mut ids = FixedSizeBinaryBuilder::with_capacity(len, 16);
460 let mut keys = StringBuilder::with_capacity(
461 len,
462 self.groups.iter().map(|r| r.question_key.len()).sum(),
463 );
464 let mut provenance = FixedSizeBinaryBuilder::with_capacity(len, 16);
465 let mut times = TimestampMicrosecondBuilder::with_capacity(len).with_timezone("UTC");
466 let mut versions = UInt32Builder::with_capacity(len);
467 for row in &self.groups {
468 append_uuid(&mut ids, row.group_uuid, "group_uuid")?;
469 keys.append_value(&row.question_key);
470 append_uuid(&mut provenance, row.provenance_uuid, "provenance_uuid")?;
471 times.append_value(row.recorded_at_micros);
472 versions.append_value(row.contract_version);
473 }
474 record_batch(
475 Arc::clone(&HYPOTHESIS_GROUP_SCHEMA),
476 vec![
477 Arc::new(ids.finish()),
478 Arc::new(keys.finish()),
479 Arc::new(provenance.finish()),
480 Arc::new(times.finish()),
481 Arc::new(versions.finish()),
482 ],
483 )
484 }
485
486 pub fn membership_batch(&self) -> Result<RecordBatch, KnowledgeError> {
488 let len = self.membership_events.len();
489 let mut ids = FixedSizeBinaryBuilder::with_capacity(len, 16);
490 let mut operations = FixedSizeBinaryBuilder::with_capacity(len, 16);
491 let mut groups = FixedSizeBinaryBuilder::with_capacity(len, 16);
492 let mut assertions = FixedSizeBinaryBuilder::with_capacity(len, 16);
493 let mut actions = StringBuilder::with_capacity(len, len * 7);
494 let mut reasoning = FixedSizeBinaryBuilder::with_capacity(len, 16);
495 let mut provenance = FixedSizeBinaryBuilder::with_capacity(len, 16);
496 let mut times = TimestampMicrosecondBuilder::with_capacity(len).with_timezone("UTC");
497 let mut versions = UInt32Builder::with_capacity(len);
498 for row in &self.membership_events {
499 append_uuid(&mut ids, row.membership_event_uuid, "membership_event_uuid")?;
500 append_uuid(&mut operations, row.operation_uuid, "operation_uuid")?;
501 append_uuid(&mut groups, row.group_uuid, "group_uuid")?;
502 append_uuid(&mut assertions, row.assertion_uuid, "assertion_uuid")?;
503 actions.append_value(row.action.as_str());
504 append_uuid(&mut reasoning, row.reasoning_uuid, "reasoning_uuid")?;
505 append_uuid(&mut provenance, row.provenance_uuid, "provenance_uuid")?;
506 times.append_value(row.recorded_at_micros);
507 versions.append_value(row.contract_version);
508 }
509 record_batch(
510 Arc::clone(&HYPOTHESIS_MEMBERSHIP_SCHEMA),
511 vec![
512 Arc::new(ids.finish()),
513 Arc::new(operations.finish()),
514 Arc::new(groups.finish()),
515 Arc::new(assertions.finish()),
516 Arc::new(actions.finish()),
517 Arc::new(reasoning.finish()),
518 Arc::new(provenance.finish()),
519 Arc::new(times.finish()),
520 Arc::new(versions.finish()),
521 ],
522 )
523 }
524
525 pub fn selection_batch(&self) -> Result<RecordBatch, KnowledgeError> {
527 let len = self.selection_events.len();
528 let mut ids = FixedSizeBinaryBuilder::with_capacity(len, 16);
529 let mut operations = FixedSizeBinaryBuilder::with_capacity(len, 16);
530 let mut groups = FixedSizeBinaryBuilder::with_capacity(len, 16);
531 let mut selected = FixedSizeBinaryBuilder::with_capacity(len, 16);
532 let mut reasoning = FixedSizeBinaryBuilder::with_capacity(len, 16);
533 let mut provenance = FixedSizeBinaryBuilder::with_capacity(len, 16);
534 let mut times = TimestampMicrosecondBuilder::with_capacity(len).with_timezone("UTC");
535 let mut versions = UInt32Builder::with_capacity(len);
536 for row in &self.selection_events {
537 append_uuid(&mut ids, row.selection_event_uuid, "selection_event_uuid")?;
538 append_uuid(&mut operations, row.operation_uuid, "operation_uuid")?;
539 append_uuid(&mut groups, row.group_uuid, "group_uuid")?;
540 if let Some(value) = row.selected_assertion_uuid {
541 append_uuid(&mut selected, value, "selected_assertion_uuid")?;
542 } else {
543 selected.append_null();
544 }
545 append_uuid(&mut reasoning, row.reasoning_uuid, "reasoning_uuid")?;
546 append_uuid(&mut provenance, row.provenance_uuid, "provenance_uuid")?;
547 times.append_value(row.recorded_at_micros);
548 versions.append_value(row.contract_version);
549 }
550 record_batch(
551 Arc::clone(&HYPOTHESIS_SELECTION_SCHEMA),
552 vec![
553 Arc::new(ids.finish()),
554 Arc::new(operations.finish()),
555 Arc::new(groups.finish()),
556 Arc::new(selected.finish()),
557 Arc::new(reasoning.finish()),
558 Arc::new(provenance.finish()),
559 Arc::new(times.finish()),
560 Arc::new(versions.finish()),
561 ],
562 )
563 }
564
565 pub fn from_batches(
567 group_batches: &[RecordBatch],
568 membership_batches: &[RecordBatch],
569 selection_batches: &[RecordBatch],
570 ) -> Result<Self, KnowledgeError> {
571 let mut groups = Vec::new();
572 for batch in group_batches {
573 require_schema(batch, &HYPOTHESIS_GROUP_SCHEMA, "hypothesis_group.schema")?;
574 let ids = fixed(batch, "group_uuid")?;
575 let keys = strings(batch, "question_key")?;
576 let provenance = fixed(batch, "provenance_uuid")?;
577 let times = timestamps(batch)?;
578 let versions = versions(batch)?;
579 for row in 0..batch.num_rows() {
580 groups.push(HypothesisGroup {
581 group_uuid: uuid_at(ids, row, "group_uuid")?,
582 question_key: keys.value(row).to_owned(),
583 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
584 recorded_at_micros: times.value(row),
585 contract_version: versions.value(row),
586 });
587 }
588 }
589 let mut membership = Vec::new();
590 for batch in membership_batches {
591 require_schema(
592 batch,
593 &HYPOTHESIS_MEMBERSHIP_SCHEMA,
594 "hypothesis_membership.schema",
595 )?;
596 let ids = fixed(batch, "membership_event_uuid")?;
597 let operations = fixed(batch, "operation_uuid")?;
598 let group_ids = fixed(batch, "group_uuid")?;
599 let assertions = fixed(batch, "assertion_uuid")?;
600 let actions = strings(batch, "action")?;
601 let reasoning = fixed(batch, "reasoning_uuid")?;
602 let provenance = fixed(batch, "provenance_uuid")?;
603 let times = timestamps(batch)?;
604 let versions = versions(batch)?;
605 for row in 0..batch.num_rows() {
606 membership.push(HypothesisMembershipEvent {
607 membership_event_uuid: uuid_at(ids, row, "membership_event_uuid")?,
608 operation_uuid: uuid_at(operations, row, "operation_uuid")?,
609 group_uuid: uuid_at(group_ids, row, "group_uuid")?,
610 assertion_uuid: uuid_at(assertions, row, "assertion_uuid")?,
611 action: HypothesisMembershipAction::parse(actions.value(row))?,
612 reasoning_uuid: uuid_at(reasoning, row, "reasoning_uuid")?,
613 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
614 recorded_at_micros: times.value(row),
615 contract_version: versions.value(row),
616 });
617 }
618 }
619 let mut selection = Vec::new();
620 for batch in selection_batches {
621 require_schema(
622 batch,
623 &HYPOTHESIS_SELECTION_SCHEMA,
624 "hypothesis_selection.schema",
625 )?;
626 let ids = fixed(batch, "selection_event_uuid")?;
627 let operations = fixed(batch, "operation_uuid")?;
628 let group_ids = fixed(batch, "group_uuid")?;
629 let selected = fixed(batch, "selected_assertion_uuid")?;
630 let reasoning = fixed(batch, "reasoning_uuid")?;
631 let provenance = fixed(batch, "provenance_uuid")?;
632 let times = timestamps(batch)?;
633 let versions = versions(batch)?;
634 for row in 0..batch.num_rows() {
635 selection.push(HypothesisSelectionEvent {
636 selection_event_uuid: uuid_at(ids, row, "selection_event_uuid")?,
637 operation_uuid: uuid_at(operations, row, "operation_uuid")?,
638 group_uuid: uuid_at(group_ids, row, "group_uuid")?,
639 selected_assertion_uuid: (!selected.is_null(row))
640 .then(|| uuid_at(selected, row, "selected_assertion_uuid"))
641 .transpose()?,
642 reasoning_uuid: uuid_at(reasoning, row, "reasoning_uuid")?,
643 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
644 recorded_at_micros: times.value(row),
645 contract_version: versions.value(row),
646 });
647 }
648 }
649 Self::new(groups, membership, selection)
650 }
651
652 #[must_use]
654 pub fn current_members(&self, group_uuid: Uuid) -> Vec<Uuid> {
655 let mut state = HashSet::new();
656 for event in self
657 .membership_events
658 .iter()
659 .filter(|row| row.group_uuid == group_uuid)
660 {
661 match event.action {
662 HypothesisMembershipAction::Added => {
663 state.insert(event.assertion_uuid);
664 }
665 HypothesisMembershipAction::Removed => {
666 state.remove(&event.assertion_uuid);
667 }
668 }
669 }
670 let mut members = state.into_iter().collect::<Vec<_>>();
671 members.sort_unstable();
672 members
673 }
674
675 #[must_use]
677 pub fn current_selection(&self, group_uuid: Uuid) -> Option<Uuid> {
678 self.selection_events
679 .iter()
680 .rfind(|row| row.group_uuid == group_uuid)
681 .and_then(|row| row.selected_assertion_uuid)
682 }
683}
684
685fn validate_question_key(value: &str) -> Result<(), KnowledgeError> {
686 if value.is_empty()
687 || value.len() > MAX_HYPOTHESIS_QUESTION_KEY_BYTES
688 || value.trim() != value
689 || !value.nfc().eq(value.chars())
690 {
691 return Err(invalid(
692 "hypothesis_group.question_key",
693 "must be non-empty bounded NFC without surrounding whitespace",
694 ));
695 }
696 Ok(())
697}
698
699fn validate_groups(groups: &[HypothesisGroup]) -> Result<(), KnowledgeError> {
700 let mut ids = HashSet::new();
701 let mut keys = HashSet::new();
702 for row in groups {
703 HypothesisGroup::new(
704 row.group_uuid,
705 row.question_key.clone(),
706 row.provenance_uuid,
707 row.recorded_at_micros,
708 )?;
709 if row.contract_version != HYPOTHESIS_GROUP_CONTRACT_VERSION {
710 return Err(invalid(
711 "hypothesis_group.contract_version",
712 "unsupported version",
713 ));
714 }
715 if !ids.insert(row.group_uuid) {
716 return Err(KnowledgeError::Duplicate("group_uuid"));
717 }
718 if !keys.insert(row.question_key.as_str()) {
719 return Err(KnowledgeError::Duplicate("question_key"));
720 }
721 }
722 Ok(())
723}
724
725fn validate_event_ids(
726 membership: &[HypothesisMembershipEvent],
727 selection: &[HypothesisSelectionEvent],
728) -> Result<(), KnowledgeError> {
729 let mut ids = HashSet::new();
730 for row in membership {
731 HypothesisMembershipEvent::new(
732 row.membership_event_uuid,
733 row.operation_uuid,
734 row.group_uuid,
735 row.assertion_uuid,
736 row.action,
737 row.reasoning_uuid,
738 row.provenance_uuid,
739 row.recorded_at_micros,
740 )?;
741 if row.contract_version != HYPOTHESIS_MEMBERSHIP_CONTRACT_VERSION {
742 return Err(invalid(
743 "hypothesis_membership.contract_version",
744 "unsupported version",
745 ));
746 }
747 if !ids.insert(row.membership_event_uuid) {
748 return Err(KnowledgeError::Duplicate("membership_event_uuid"));
749 }
750 }
751 ids.clear();
752 for row in selection {
753 HypothesisSelectionEvent::new(
754 row.selection_event_uuid,
755 row.operation_uuid,
756 row.group_uuid,
757 row.selected_assertion_uuid,
758 row.reasoning_uuid,
759 row.provenance_uuid,
760 row.recorded_at_micros,
761 )?;
762 if row.contract_version != HYPOTHESIS_SELECTION_CONTRACT_VERSION {
763 return Err(invalid(
764 "hypothesis_selection.contract_version",
765 "unsupported version",
766 ));
767 }
768 if !ids.insert(row.selection_event_uuid) {
769 return Err(KnowledgeError::Duplicate("selection_event_uuid"));
770 }
771 }
772 Ok(())
773}
774
775fn validate_state(
776 groups: &[HypothesisGroup],
777 membership: &[HypothesisMembershipEvent],
778 selection: &[HypothesisSelectionEvent],
779) -> Result<(), KnowledgeError> {
780 let group_times = groups
781 .iter()
782 .map(|row| (row.group_uuid, row.recorded_at_micros))
783 .collect::<HashMap<_, _>>();
784 let mut members = HashMap::<Uuid, HashSet<Uuid>>::new();
785 let mut selected = HashMap::<Uuid, Option<Uuid>>::new();
786 for (time, _, operation_uuid) in operation_keys(membership, selection) {
787 let selected_before = selected.clone();
788 let membership_at_time = membership
789 .iter()
790 .filter(|row| row.recorded_at_micros == time && row.operation_uuid == operation_uuid)
791 .collect::<Vec<_>>();
792 let selection_at_time = selection
793 .iter()
794 .filter(|row| row.recorded_at_micros == time && row.operation_uuid == operation_uuid)
795 .collect::<Vec<_>>();
796 for event in &membership_at_time {
797 let group_time = group_times
798 .get(&event.group_uuid)
799 .ok_or(KnowledgeError::Dangling("group_uuid"))?;
800 if event.recorded_at_micros < *group_time {
801 return Err(invalid(
802 "hypothesis_membership.recorded_at",
803 "cannot predate its hypothesis group",
804 ));
805 }
806 let state = members.entry(event.group_uuid).or_default();
807 match event.action {
808 HypothesisMembershipAction::Added if !state.insert(event.assertion_uuid) => {
809 return Err(invalid(
810 "hypothesis_membership.action",
811 "cannot add a current member",
812 ));
813 }
814 HypothesisMembershipAction::Removed if !state.remove(&event.assertion_uuid) => {
815 return Err(invalid(
816 "hypothesis_membership.action",
817 "cannot remove a non-member",
818 ));
819 }
820 _ => {}
821 }
822 }
823 for event in &selection_at_time {
824 let group_time = group_times
825 .get(&event.group_uuid)
826 .ok_or(KnowledgeError::Dangling("group_uuid"))?;
827 if event.recorded_at_micros < *group_time {
828 return Err(invalid(
829 "hypothesis_selection.recorded_at",
830 "cannot predate its hypothesis group",
831 ));
832 }
833 if let Some(assertion_uuid) = event.selected_assertion_uuid
834 && !members
835 .get(&event.group_uuid)
836 .is_some_and(|state| state.contains(&assertion_uuid))
837 {
838 return Err(invalid(
839 "hypothesis_selection.selected_assertion_uuid",
840 "selected assertion must be a current member",
841 ));
842 }
843 selected.insert(event.group_uuid, event.selected_assertion_uuid);
844 }
845 for removal in membership_at_time
846 .iter()
847 .filter(|row| row.action == HypothesisMembershipAction::Removed)
848 {
849 if selected_before
850 .get(&removal.group_uuid)
851 .copied()
852 .flatten()
853 .is_some_and(|id| id == removal.assertion_uuid)
854 {
855 let paired = selection_at_time.iter().any(|event| {
856 event.group_uuid == removal.group_uuid
857 && event.operation_uuid == removal.operation_uuid
858 && event.selected_assertion_uuid != Some(removal.assertion_uuid)
859 });
860 if !paired {
861 return Err(invalid(
862 "hypothesis_membership.action",
863 "selected-member removal requires a paired explicit selection event",
864 ));
865 }
866 }
867 }
868 }
869 Ok(())
870}
871
872fn operation_keys(
873 membership: &[HypothesisMembershipEvent],
874 selection: &[HypothesisSelectionEvent],
875) -> Vec<(i64, Uuid, Uuid)> {
876 let mut operations = HashMap::<(i64, Uuid), Uuid>::new();
877 for (time, operation, event) in membership
878 .iter()
879 .map(|row| {
880 (
881 row.recorded_at_micros,
882 row.operation_uuid,
883 row.membership_event_uuid,
884 )
885 })
886 .chain(selection.iter().map(|row| {
887 (
888 row.recorded_at_micros,
889 row.operation_uuid,
890 row.selection_event_uuid,
891 )
892 }))
893 {
894 operations
895 .entry((time, operation))
896 .and_modify(|first| *first = (*first).min(event))
897 .or_insert(event);
898 }
899 let mut keys = operations
900 .into_iter()
901 .map(|((time, operation), first_event)| (time, first_event, operation))
902 .collect::<Vec<_>>();
903 keys.sort_unstable();
904 keys
905}
906
907fn merge_rows<T, F>(
908 existing: &[T],
909 staged: &[T],
910 id: F,
911 field: &'static str,
912) -> Result<Vec<T>, KnowledgeError>
913where
914 T: Clone + Eq,
915 F: Fn(&T) -> Uuid,
916{
917 let mut rows = existing.to_vec();
918 let mut by_id = rows
919 .iter()
920 .cloned()
921 .map(|row| (id(&row), row))
922 .collect::<HashMap<_, _>>();
923 for row in staged {
924 if let Some(current) = by_id.get(&id(row)) {
925 if current != row {
926 return Err(KnowledgeError::Conflict(field));
927 }
928 } else {
929 rows.push(row.clone());
930 by_id.insert(id(row), row.clone());
931 }
932 }
933 Ok(rows)
934}
935
936fn append_uuid(
937 builder: &mut FixedSizeBinaryBuilder,
938 value: Uuid,
939 field: &'static str,
940) -> Result<(), KnowledgeError> {
941 builder
942 .append_value(value.as_bytes())
943 .map_err(|_| invalid(field, "invalid UUID width"))
944}
945
946fn record_batch(
947 schema: SchemaRef,
948 columns: Vec<Arc<dyn arrow::array::Array>>,
949) -> Result<RecordBatch, KnowledgeError> {
950 RecordBatch::try_new(schema, columns)
951 .map_err(|_| invalid("hypothesis", "Arrow batch construction failed"))
952}
953
954fn require_schema(
955 batch: &RecordBatch,
956 schema: &SchemaRef,
957 field: &'static str,
958) -> Result<(), KnowledgeError> {
959 if batch.schema().as_ref() != schema.as_ref() {
960 return Err(invalid(field, "schema mismatch"));
961 }
962 Ok(())
963}
964
965fn fixed<'a>(
966 batch: &'a RecordBatch,
967 name: &'static str,
968) -> Result<&'a FixedSizeBinaryArray, KnowledgeError> {
969 batch
970 .column_by_name(name)
971 .and_then(|value| value.as_any().downcast_ref())
972 .ok_or_else(|| invalid(name, "column type mismatch"))
973}
974
975fn strings<'a>(
976 batch: &'a RecordBatch,
977 name: &'static str,
978) -> Result<&'a StringArray, KnowledgeError> {
979 batch
980 .column_by_name(name)
981 .and_then(|value| value.as_any().downcast_ref())
982 .ok_or_else(|| invalid(name, "column type mismatch"))
983}
984
985fn timestamps(batch: &RecordBatch) -> Result<&TimestampMicrosecondArray, KnowledgeError> {
986 batch
987 .column_by_name("recorded_at")
988 .and_then(|value| value.as_any().downcast_ref())
989 .ok_or_else(|| invalid("recorded_at", "column type mismatch"))
990}
991
992fn versions(batch: &RecordBatch) -> Result<&UInt32Array, KnowledgeError> {
993 batch
994 .column_by_name("contract_version")
995 .and_then(|value| value.as_any().downcast_ref())
996 .ok_or_else(|| invalid("contract_version", "column type mismatch"))
997}
998
999fn uuid_at(
1000 values: &FixedSizeBinaryArray,
1001 row: usize,
1002 field: &'static str,
1003) -> Result<Uuid, KnowledgeError> {
1004 if values.is_null(row) {
1005 return Err(invalid(field, "unexpected null UUID"));
1006 }
1007 Uuid::from_slice(values.value(row)).map_err(|_| invalid(field, "invalid UUID bytes"))
1008}
1009
1010pub(crate) fn schema_registry_entries() -> [SchemaRegistryEntry; 3] {
1011 [
1012 entry(
1013 "hypothesis_groups",
1014 HYPOTHESIS_GROUP_CONTRACT_VERSION,
1015 Arc::clone(&HYPOTHESIS_GROUP_SCHEMA),
1016 *GROUP_SCHEMA_FINGERPRINT,
1017 &["recorded_at", "group_uuid"],
1018 (&["group_uuid"], "group_uuid"),
1019 &[("hypothesis_key_policy", HYPOTHESIS_KEY_POLICY_VERSION)],
1020 ),
1021 entry(
1022 "hypothesis_membership_events",
1023 HYPOTHESIS_MEMBERSHIP_CONTRACT_VERSION,
1024 Arc::clone(&HYPOTHESIS_MEMBERSHIP_SCHEMA),
1025 *MEMBERSHIP_SCHEMA_FINGERPRINT,
1026 &["recorded_at", "membership_event_uuid"],
1027 (&["membership_event_uuid"], "membership_event_uuid"),
1028 &[
1029 ("hypothesis_state_policy", HYPOTHESIS_STATE_POLICY_VERSION),
1030 ("membership_action", 1),
1031 ],
1032 ),
1033 entry(
1034 "hypothesis_selection_events",
1035 HYPOTHESIS_SELECTION_CONTRACT_VERSION,
1036 Arc::clone(&HYPOTHESIS_SELECTION_SCHEMA),
1037 *SELECTION_SCHEMA_FINGERPRINT,
1038 &["recorded_at", "selection_event_uuid"],
1039 (&["selection_event_uuid"], "selection_event_uuid"),
1040 &[("hypothesis_state_policy", HYPOTHESIS_STATE_POLICY_VERSION)],
1041 ),
1042 ]
1043}
1044
1045fn entry(
1046 family: &'static str,
1047 version: u32,
1048 schema: SchemaRef,
1049 schema_fingerprint: [u8; 32],
1050 sort_key: &'static [&'static str],
1051 diff_identity: (&'static [&'static str], &'static str),
1052 enums: &'static [(&'static str, u32)],
1053) -> SchemaRegistryEntry {
1054 SchemaRegistryEntry {
1055 capability_id: "epistemic",
1056 capability_version: EPISTEMIC_CAPABILITY_VERSION,
1057 record_family: family,
1058 record_version: version,
1059 schema,
1060 schema_fingerprint,
1061 enum_registry_versions: enums,
1062 sort_key,
1063 diff_identity_fields: diff_identity.0,
1064 diff_record_uuid_field: Some(diff_identity.1),
1065 fingerprint_domain: CanonicalDomain::Schema,
1066 owner: "graphforge-knowledge",
1067 implementation_issue: 779,
1068 max_rows: MAX_KNOWLEDGE_ROWS,
1069 }
1070}
1071
1072fn uuid_field(name: &str, nullable: bool) -> Field {
1073 Field::new(name, DataType::FixedSizeBinary(16), nullable)
1074}
1075
1076fn timestamp_field(name: &str) -> Field {
1077 Field::new(
1078 name,
1079 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1080 false,
1081 )
1082}
1083
1084fn require_v7(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
1085 if value.get_version() != Some(Version::SortRand) {
1086 return Err(invalid(field, "must be UUIDv7"));
1087 }
1088 require_uuid(value, field)
1089}
1090
1091fn require_uuid(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
1092 if value.is_nil() {
1093 return Err(invalid(field, "must not be nil"));
1094 }
1095 Ok(())
1096}
1097
1098const fn invalid(field: &'static str, message: &'static str) -> KnowledgeError {
1099 KnowledgeError::Invalid { field, message }
1100}
1101
1102#[cfg(test)]
1103mod tests {
1104 use super::*;
1105
1106 fn invalid_code(error: &KnowledgeError) -> &'static str {
1107 match error {
1108 KnowledgeError::Invalid { .. } => "invalid",
1109 KnowledgeError::Duplicate(_) => "duplicate",
1110 KnowledgeError::Dangling(_) => "dangling",
1111 KnowledgeError::Conflict(_) => "conflict",
1112 _ => "other",
1113 }
1114 }
1115
1116 fn uuid7(seed: u8) -> Uuid {
1117 let mut bytes = [seed; 16];
1118 bytes[6] = (bytes[6] & 0x0f) | 0x70;
1119 bytes[8] = (bytes[8] & 0x3f) | 0x80;
1120 Uuid::from_bytes(bytes)
1121 }
1122
1123 fn group() -> HypothesisGroup {
1124 HypothesisGroup::new(uuid7(1), "cause.primary".into(), uuid7(2), 1).unwrap()
1125 }
1126
1127 fn member(
1128 id: u8,
1129 operation: u8,
1130 assertion: u8,
1131 action: HypothesisMembershipAction,
1132 time: i64,
1133 ) -> HypothesisMembershipEvent {
1134 HypothesisMembershipEvent::new(
1135 uuid7(id),
1136 uuid7(operation),
1137 uuid7(1),
1138 uuid7(assertion),
1139 action,
1140 uuid7(id.wrapping_add(40)),
1141 uuid7(id.wrapping_add(80)),
1142 time,
1143 )
1144 .unwrap()
1145 }
1146
1147 fn selection(
1148 id: u8,
1149 operation: u8,
1150 selected: Option<u8>,
1151 time: i64,
1152 ) -> HypothesisSelectionEvent {
1153 HypothesisSelectionEvent::new(
1154 uuid7(id),
1155 uuid7(operation),
1156 uuid7(1),
1157 selected.map(uuid7),
1158 uuid7(id.wrapping_add(40)),
1159 uuid7(id.wrapping_add(80)),
1160 time,
1161 )
1162 .unwrap()
1163 }
1164
1165 #[test]
1166 fn keys_are_exact_nfc_bounded_and_unique() {
1167 assert!(HypothesisGroup::new(uuid7(1), String::new(), uuid7(2), 1).is_err());
1168 assert!(HypothesisGroup::new(uuid7(1), " key".into(), uuid7(2), 1).is_err());
1169 assert!(HypothesisGroup::new(uuid7(1), "e\u{301}".into(), uuid7(2), 1).is_err());
1170 let upper = HypothesisGroup::new(uuid7(3), "Cause".into(), uuid7(4), 1).unwrap();
1171 let lower = HypothesisGroup::new(uuid7(5), "cause".into(), uuid7(6), 1).unwrap();
1172 assert!(HypothesisLedger::new(vec![upper, lower], vec![], vec![]).is_ok());
1173 let duplicate =
1174 HypothesisGroup::new(uuid7(7), "cause.primary".into(), uuid7(8), 1).unwrap();
1175 assert!(HypothesisLedger::new(vec![group(), duplicate], vec![], vec![]).is_err());
1176 }
1177
1178 #[test]
1179 fn add_select_change_clear_and_selected_removal_are_explicit() {
1180 let ledger = HypothesisLedger::new(
1181 vec![group()],
1182 vec![
1183 member(10, 20, 30, HypothesisMembershipAction::Added, 2),
1184 member(11, 21, 31, HypothesisMembershipAction::Added, 3),
1185 member(12, 22, 30, HypothesisMembershipAction::Removed, 5),
1186 ],
1187 vec![
1188 selection(13, 23, Some(30), 4),
1189 selection(14, 22, Some(31), 5),
1190 selection(15, 24, None, 6),
1191 ],
1192 )
1193 .unwrap();
1194 assert_eq!(ledger.current_members(uuid7(1)), vec![uuid7(31)]);
1195 assert_eq!(ledger.current_selection(uuid7(1)), None);
1196 assert_eq!(
1197 HypothesisLedger::from_batches(
1198 &[ledger.group_batch().unwrap()],
1199 &[ledger.membership_batch().unwrap()],
1200 &[ledger.selection_batch().unwrap()],
1201 )
1202 .unwrap(),
1203 ledger
1204 );
1205 assert_eq!(ledger.merge(&ledger).unwrap(), ledger);
1206
1207 assert!(
1208 HypothesisLedger::new(
1209 vec![group()],
1210 vec![
1211 member(10, 20, 30, HypothesisMembershipAction::Added, 2),
1212 member(12, 22, 30, HypothesisMembershipAction::Removed, 4),
1213 ],
1214 vec![selection(13, 23, Some(30), 3)],
1215 )
1216 .is_err()
1217 );
1218 }
1219
1220 #[test]
1221 fn canonical_fingerprints_are_exact_order_independent_and_round_trip_stable() {
1222 assert_eq!(
1223 CanonicalDomain::HypothesisGroup.as_str(),
1224 "graphforge/hypothesis-group"
1225 );
1226 assert_eq!(
1227 CanonicalDomain::HypothesisMembership.as_str(),
1228 "graphforge/hypothesis-membership"
1229 );
1230 assert_eq!(
1231 CanonicalDomain::HypothesisSelection.as_str(),
1232 "graphforge/hypothesis-selection"
1233 );
1234 let membership = member(10, 20, 30, HypothesisMembershipAction::Added, 2);
1235 let other = member(11, 21, 31, HypothesisMembershipAction::Added, 3);
1236 let selected = selection(12, 22, Some(30), 4);
1237 let cleared = selection(13, 23, None, 5);
1238 let ledger = HypothesisLedger::new(
1239 vec![group()],
1240 vec![membership.clone(), other.clone()],
1241 vec![selected.clone(), cleared.clone()],
1242 )
1243 .unwrap();
1244 let reordered = HypothesisLedger::new(
1245 vec![group()],
1246 vec![other, membership.clone()],
1247 vec![cleared.clone(), selected.clone()],
1248 )
1249 .unwrap();
1250 let decoded = HypothesisLedger::from_batches(
1251 &[ledger.group_batch().unwrap()],
1252 &[ledger.membership_batch().unwrap()],
1253 &[ledger.selection_batch().unwrap()],
1254 )
1255 .unwrap();
1256 let fingerprints = [
1257 ledger.group_fingerprint(uuid7(1)).unwrap(),
1258 ledger
1259 .membership_fingerprint(membership.membership_event_uuid)
1260 .unwrap(),
1261 ledger
1262 .selection_fingerprint(selected.selection_event_uuid)
1263 .unwrap(),
1264 ledger
1265 .selection_fingerprint(cleared.selection_event_uuid)
1266 .unwrap(),
1267 ];
1268 let hex = |value: [u8; 32]| {
1269 value
1270 .iter()
1271 .map(|byte| format!("{byte:02x}"))
1272 .collect::<String>()
1273 };
1274 assert_eq!(
1275 hex(fingerprints[0]),
1276 "336fe80f4e94c5758b5e77de06c9c5b6ebf51ff20355c5d6950ab3e648991265"
1277 );
1278 assert_eq!(
1279 hex(fingerprints[1]),
1280 "fc4fe4882c1469a3f4606c5352028a63d43d369301e26ac0216373889945e403"
1281 );
1282 assert_eq!(
1283 hex(fingerprints[2]),
1284 "ba0ec00b86950fd8b60afaa40327f240377f8fd1695a54dad1d51cbf45665615"
1285 );
1286 assert_eq!(
1287 hex(fingerprints[3]),
1288 "fc220def4f2a61fb134f5a414a257d50cb29a1aea364464831f89b7747e44f9b"
1289 );
1290 assert_eq!(
1291 reordered.group_fingerprint(uuid7(1)).unwrap(),
1292 fingerprints[0]
1293 );
1294 assert_eq!(
1295 reordered
1296 .membership_fingerprint(membership.membership_event_uuid)
1297 .unwrap(),
1298 fingerprints[1]
1299 );
1300 assert_eq!(
1301 reordered
1302 .selection_fingerprint(selected.selection_event_uuid)
1303 .unwrap(),
1304 fingerprints[2]
1305 );
1306 assert_eq!(
1307 reordered
1308 .selection_fingerprint(cleared.selection_event_uuid)
1309 .unwrap(),
1310 fingerprints[3]
1311 );
1312 assert_eq!(
1313 decoded.group_fingerprint(uuid7(1)).unwrap(),
1314 fingerprints[0]
1315 );
1316 assert_eq!(
1317 decoded
1318 .membership_fingerprint(membership.membership_event_uuid)
1319 .unwrap(),
1320 fingerprints[1]
1321 );
1322 assert_eq!(
1323 decoded
1324 .selection_fingerprint(selected.selection_event_uuid)
1325 .unwrap(),
1326 fingerprints[2]
1327 );
1328 assert_eq!(
1329 decoded
1330 .selection_fingerprint(cleared.selection_event_uuid)
1331 .unwrap(),
1332 fingerprints[3]
1333 );
1334 }
1335
1336 #[test]
1337 fn selecting_nonmember_and_invalid_membership_transitions_fail() {
1338 assert!(
1339 HypothesisLedger::new(vec![group()], vec![], vec![selection(10, 20, Some(30), 2)],)
1340 .is_err()
1341 );
1342 assert!(
1343 HypothesisLedger::new(
1344 vec![group()],
1345 vec![member(10, 20, 30, HypothesisMembershipAction::Removed, 2)],
1346 vec![],
1347 )
1348 .is_err()
1349 );
1350 assert!(
1351 HypothesisLedger::new(
1352 vec![group()],
1353 vec![
1354 member(10, 20, 30, HypothesisMembershipAction::Added, 2),
1355 member(11, 21, 30, HypothesisMembershipAction::Added, 3),
1356 ],
1357 vec![],
1358 )
1359 .is_err()
1360 );
1361 }
1362
1363 #[test]
1364 fn wave12_contract_versions_and_duplicate_identities_fail_closed() {
1365 let mut bad_group_version = group();
1366 bad_group_version.contract_version = 2;
1367 assert_eq!(
1368 invalid_code(
1369 &HypothesisLedger::new(vec![bad_group_version], vec![], vec![]).unwrap_err()
1370 ),
1371 "invalid"
1372 );
1373 assert_eq!(
1374 invalid_code(
1375 &HypothesisLedger::new(vec![group(), group()], vec![], vec![]).unwrap_err()
1376 ),
1377 "duplicate"
1378 );
1379
1380 let membership = member(10, 20, 30, HypothesisMembershipAction::Added, 2);
1381 let mut bad_membership_version = membership.clone();
1382 bad_membership_version.contract_version = 2;
1383 assert_eq!(
1384 invalid_code(
1385 &HypothesisLedger::new(vec![group()], vec![bad_membership_version], vec![])
1386 .unwrap_err()
1387 ),
1388 "invalid"
1389 );
1390 assert_eq!(
1391 invalid_code(
1392 &HypothesisLedger::new(vec![group()], vec![membership.clone(), membership], vec![])
1393 .unwrap_err()
1394 ),
1395 "duplicate"
1396 );
1397
1398 let selected = selection(11, 21, None, 3);
1399 let mut bad_selection_version = selected.clone();
1400 bad_selection_version.contract_version = 2;
1401 assert_eq!(
1402 invalid_code(
1403 &HypothesisLedger::new(vec![group()], vec![], vec![bad_selection_version])
1404 .unwrap_err()
1405 ),
1406 "invalid"
1407 );
1408 assert_eq!(
1409 invalid_code(
1410 &HypothesisLedger::new(vec![group()], vec![], vec![selected.clone(), selected])
1411 .unwrap_err()
1412 ),
1413 "duplicate"
1414 );
1415 }
1416
1417 #[test]
1418 fn wave12_dangling_and_temporally_invalid_events_are_rejected() {
1419 let dangling_member = HypothesisMembershipEvent::new(
1420 uuid7(10),
1421 uuid7(20),
1422 uuid7(99),
1423 uuid7(30),
1424 HypothesisMembershipAction::Added,
1425 uuid7(40),
1426 uuid7(50),
1427 2,
1428 )
1429 .unwrap();
1430 assert_eq!(
1431 invalid_code(
1432 &HypothesisLedger::new(vec![group()], vec![dangling_member], vec![]).unwrap_err()
1433 ),
1434 "dangling"
1435 );
1436 let early_member = member(11, 21, 31, HypothesisMembershipAction::Added, 0);
1437 assert_eq!(
1438 invalid_code(
1439 &HypothesisLedger::new(vec![group()], vec![early_member], vec![]).unwrap_err()
1440 ),
1441 "invalid"
1442 );
1443
1444 let dangling_selection = HypothesisSelectionEvent::new(
1445 uuid7(12),
1446 uuid7(22),
1447 uuid7(99),
1448 None,
1449 uuid7(42),
1450 uuid7(52),
1451 2,
1452 )
1453 .unwrap();
1454 assert_eq!(
1455 invalid_code(
1456 &HypothesisLedger::new(vec![group()], vec![], vec![dangling_selection])
1457 .unwrap_err()
1458 ),
1459 "dangling"
1460 );
1461 let early_selection = selection(13, 23, None, 0);
1462 assert_eq!(
1463 invalid_code(
1464 &HypothesisLedger::new(vec![group()], vec![], vec![early_selection]).unwrap_err()
1465 ),
1466 "invalid"
1467 );
1468 }
1469
1470 #[test]
1471 fn wave12_merge_conflicts_and_arrow_decode_errors_are_structured() {
1472 let base = HypothesisLedger::new(vec![group()], vec![], vec![]).unwrap();
1473 let conflicting_group =
1474 HypothesisGroup::new(uuid7(1), "different.question".into(), uuid7(2), 1).unwrap();
1475 let staged = HypothesisLedger::new(vec![conflicting_group], vec![], vec![]).unwrap();
1476 assert_eq!(invalid_code(&base.merge(&staged).unwrap_err()), "conflict");
1477
1478 let wrong_schema = RecordBatch::new_empty(Arc::new(Schema::empty()));
1479 assert_eq!(
1480 invalid_code(&HypothesisLedger::from_batches(&[wrong_schema], &[], &[]).unwrap_err()),
1481 "invalid"
1482 );
1483
1484 let ledger = HypothesisLedger::new(
1485 vec![group()],
1486 vec![member(10, 20, 30, HypothesisMembershipAction::Added, 2)],
1487 vec![],
1488 )
1489 .unwrap();
1490 let batch = ledger.membership_batch().unwrap();
1491 let mut columns = batch.columns().to_vec();
1492 columns[4] = Arc::new(StringArray::from(vec!["unknown"]));
1493 let unknown_action = RecordBatch::try_new(batch.schema(), columns).unwrap();
1494 assert_eq!(
1495 invalid_code(
1496 &HypothesisLedger::from_batches(
1497 &[ledger.group_batch().unwrap()],
1498 &[unknown_action],
1499 &[]
1500 )
1501 .unwrap_err()
1502 ),
1503 "invalid"
1504 );
1505 }
1506}