1use serde::Deserialize;
11
12use super::{ConflictResolutionMode, OpRecord, RecordedHead, ThreadUpdateSnapshots};
13use crate::{
14 error::{HeddleError, Result},
15 object::{Attribution, ChangeId, ContentHash, StateId, VisibilityTier},
16};
17
18pub const CURRENT_OP_RECORD_SCHEMA_VERSION: u32 = 4;
19const CURRENT_OP_RECORD_SCHEMA_NAME: &str = "state-id-v4";
20const OP_RECORD_STORAGE: &str = "oplog record schema";
21
22pub fn validate_op_record_schema_version(version: u32) -> Result<()> {
23 if version < CURRENT_OP_RECORD_SCHEMA_VERSION {
24 return Err(HeddleError::StorageFormatTooOld {
25 storage: OP_RECORD_STORAGE.to_string(),
26 found: version,
27 required: CURRENT_OP_RECORD_SCHEMA_VERSION,
28 });
29 }
30 if version > CURRENT_OP_RECORD_SCHEMA_VERSION {
31 return Err(HeddleError::StorageFormatTooNew {
32 storage: OP_RECORD_STORAGE.to_string(),
33 found: version,
34 supported: CURRENT_OP_RECORD_SCHEMA_VERSION,
35 });
36 }
37 Ok(())
38}
39
40pub fn decode_current_record(bytes: &[u8]) -> Result<OpRecord> {
41 let record: StrictCurrentOpRecord = decode_rmp(bytes, CURRENT_OP_RECORD_SCHEMA_NAME)?;
42 Ok(record.into_current())
43}
44
45pub fn encode_current_record(record: &OpRecord) -> Result<Vec<u8>> {
46 rmp_serde::to_vec(record).map_err(|e| HeddleError::Serialization(e.to_string()))
47}
48
49fn decode_rmp<T>(bytes: &[u8], schema_name: &str) -> Result<T>
50where
51 T: for<'de> Deserialize<'de>,
52{
53 rmp_serde::from_slice(bytes).map_err(|e| {
54 HeddleError::Serialization(format!(
55 "failed to decode OpRecord payload as {schema_name}: {e}"
56 ))
57 })
58}
59
60#[derive(Debug, Clone, Deserialize)]
65enum StrictCurrentOpRecord {
66 Snapshot {
67 new_state: StateId,
68 prev_head: Option<StateId>,
69 head: Option<StateId>,
70 thread: Option<String>,
71 },
72 Goto {
73 target: StateId,
74 prev_head: Option<StateId>,
75 head: StateId,
76 },
77 ThreadCreate {
78 name: String,
79 state: StateId,
80 manager_snapshot: Option<Vec<u8>>,
81 },
82 ThreadDelete {
83 name: String,
84 state: StateId,
85 },
86 ThreadUpdate {
87 name: String,
88 old_state: StateId,
89 new_state: StateId,
90 #[serde(default)]
91 manager_snapshots: Option<ThreadUpdateSnapshots>,
92 },
93 Fork {
94 from: StateId,
95 new_state: StateId,
96 thread: Option<String>,
97 head: Option<StateId>,
98 },
99 Collapse {
100 sources: Vec<StateId>,
101 result: StateId,
102 thread: Option<String>,
103 pre_thread_state: Option<StateId>,
104 },
105 MarkerCreate {
106 name: String,
107 state: StateId,
108 },
109 MarkerDelete {
110 name: String,
111 state: StateId,
112 },
113 Checkpoint {
114 parent: Option<StateId>,
115 state: StateId,
116 thread: Option<String>,
117 },
118 TransactionAbort {
119 transaction_id: String,
120 reason: String,
121 },
122 EphemeralThreadCollapse {
123 thread: String,
124 final_state: StateId,
125 },
126 ConflictResolved {
127 conflict_id: String,
128 resolution: String,
129 resolver: Attribution,
130 mode: ConflictResolutionMode,
131 },
132 TransactionCommit {
133 transaction_id: String,
134 op_count: u32,
135 },
136 Redact {
137 redaction_id: ContentHash,
138 blob: ContentHash,
139 state: StateId,
140 path: String,
141 },
142 Purge {
143 redaction_id: ContentHash,
144 blob: ContentHash,
145 },
146 FastForward {
147 source_thread: String,
148 target_thread: String,
149 pre_target_id: StateId,
150 post_target_id: StateId,
151 },
152 GitCheckpoint {
153 branch: String,
154 state: StateId,
155 previous_git_oid: Option<String>,
156 new_git_oid: String,
157 },
158 RemoteThreadUpdate {
159 remote: String,
160 thread: String,
161 state: StateId,
162 },
163 RemoteThreadDelete {
164 remote: String,
165 thread: String,
166 state: StateId,
167 },
168 UndoRecoveryUpdate {
169 state: StateId,
170 },
171 StateVisibilitySet {
172 state: StateId,
173 record_id: ContentHash,
174 tier: VisibilityTier,
175 #[serde(default)]
176 prior_sidecar: Option<Vec<u8>>,
177 #[serde(default)]
178 new_sidecar: Option<Vec<u8>>,
179 },
180 StateVisibilityPromote {
181 state: StateId,
182 superseded: ContentHash,
183 record_id: ContentHash,
184 tier: VisibilityTier,
185 #[serde(default)]
186 prior_sidecar: Option<Vec<u8>>,
187 #[serde(default)]
188 new_sidecar: Option<Vec<u8>>,
189 },
190 HeadUpdate {
191 previous: RecordedHead,
192 new: RecordedHead,
193 },
194 EntryVisibilitySet {
195 change_id: ChangeId,
196 record_id: ContentHash,
197 #[serde(default)]
198 prior_sidecar: Option<Vec<u8>>,
199 #[serde(default)]
200 new_sidecar: Option<Vec<u8>>,
201 },
202}
203
204impl StrictCurrentOpRecord {
205 fn into_current(self) -> OpRecord {
206 match self {
207 Self::Snapshot {
208 new_state,
209 prev_head,
210 head,
211 thread,
212 } => OpRecord::Snapshot {
213 new_state,
214 prev_head,
215 head,
216 thread,
217 },
218 Self::Goto {
219 target,
220 prev_head,
221 head,
222 } => OpRecord::Goto {
223 target,
224 prev_head,
225 head,
226 },
227 Self::ThreadCreate {
228 name,
229 state,
230 manager_snapshot,
231 } => OpRecord::ThreadCreate {
232 name,
233 state,
234 manager_snapshot,
235 },
236 Self::ThreadDelete { name, state } => OpRecord::ThreadDelete { name, state },
237 Self::ThreadUpdate {
238 name,
239 old_state,
240 new_state,
241 manager_snapshots,
242 } => OpRecord::ThreadUpdate {
243 name,
244 old_state,
245 new_state,
246 manager_snapshots,
247 },
248 Self::Fork {
249 from,
250 new_state,
251 thread,
252 head,
253 } => OpRecord::Fork {
254 from,
255 new_state,
256 thread,
257 head,
258 },
259 Self::Collapse {
260 sources,
261 result,
262 thread,
263 pre_thread_state,
264 } => OpRecord::Collapse {
265 sources,
266 result,
267 thread,
268 pre_thread_state,
269 },
270 Self::MarkerCreate { name, state } => OpRecord::MarkerCreate { name, state },
271 Self::MarkerDelete { name, state } => OpRecord::MarkerDelete { name, state },
272 Self::Checkpoint {
273 parent,
274 state,
275 thread,
276 } => OpRecord::Checkpoint {
277 parent,
278 state,
279 thread,
280 },
281 Self::TransactionAbort {
282 transaction_id,
283 reason,
284 } => OpRecord::TransactionAbort {
285 transaction_id,
286 reason,
287 },
288 Self::EphemeralThreadCollapse {
289 thread,
290 final_state,
291 } => OpRecord::EphemeralThreadCollapse {
292 thread,
293 final_state,
294 },
295 Self::ConflictResolved {
296 conflict_id,
297 resolution,
298 resolver,
299 mode,
300 } => OpRecord::ConflictResolved {
301 conflict_id,
302 resolution,
303 resolver,
304 mode,
305 },
306 Self::TransactionCommit {
307 transaction_id,
308 op_count,
309 } => OpRecord::TransactionCommit {
310 transaction_id,
311 op_count,
312 },
313 Self::Redact {
314 redaction_id,
315 blob,
316 state,
317 path,
318 } => OpRecord::Redact {
319 redaction_id,
320 blob,
321 state,
322 path,
323 },
324 Self::Purge { redaction_id, blob } => OpRecord::Purge { redaction_id, blob },
325 Self::FastForward {
326 source_thread,
327 target_thread,
328 pre_target_id,
329 post_target_id,
330 } => OpRecord::FastForward {
331 source_thread,
332 target_thread,
333 pre_target_id,
334 post_target_id,
335 },
336 Self::GitCheckpoint {
337 branch,
338 state,
339 previous_git_oid,
340 new_git_oid,
341 } => OpRecord::GitCheckpoint {
342 branch,
343 state,
344 previous_git_oid,
345 new_git_oid,
346 },
347 Self::RemoteThreadUpdate {
348 remote,
349 thread,
350 state,
351 } => OpRecord::RemoteThreadUpdate {
352 remote,
353 thread,
354 state,
355 },
356 Self::RemoteThreadDelete {
357 remote,
358 thread,
359 state,
360 } => OpRecord::RemoteThreadDelete {
361 remote,
362 thread,
363 state,
364 },
365 Self::UndoRecoveryUpdate { state } => OpRecord::UndoRecoveryUpdate { state },
366 Self::StateVisibilitySet {
367 state,
368 record_id,
369 tier,
370 prior_sidecar,
371 new_sidecar,
372 } => OpRecord::StateVisibilitySet {
373 state,
374 record_id,
375 tier,
376 prior_sidecar,
377 new_sidecar,
378 },
379 Self::StateVisibilityPromote {
380 state,
381 superseded,
382 record_id,
383 tier,
384 prior_sidecar,
385 new_sidecar,
386 } => OpRecord::StateVisibilityPromote {
387 state,
388 superseded,
389 record_id,
390 tier,
391 prior_sidecar,
392 new_sidecar,
393 },
394 Self::HeadUpdate { previous, new } => OpRecord::HeadUpdate { previous, new },
395 Self::EntryVisibilitySet {
396 change_id,
397 record_id,
398 prior_sidecar,
399 new_sidecar,
400 } => OpRecord::EntryVisibilitySet {
401 change_id,
402 record_id,
403 prior_sidecar,
404 new_sidecar,
405 },
406 }
407 }
408}
409
410#[cfg(test)]
411mod tests {
412 use super::*;
413 use crate::object::{Agent, Principal};
414
415 fn state(byte: u8) -> StateId {
416 StateId::from_bytes([byte; 32])
417 }
418
419 fn hash(byte: u8) -> ContentHash {
420 ContentHash::from_bytes([byte; 32])
421 }
422
423 fn assert_round_trip(record: OpRecord) {
424 let bytes = encode_current_record(&record).unwrap();
425 let decoded = decode_current_record(&bytes).unwrap();
426 assert_eq!(format!("{decoded:?}"), format!("{record:?}"));
427 }
428
429 fn canonical_current_records() -> Vec<OpRecord> {
430 vec![
431 OpRecord::Snapshot {
432 new_state: state(1),
433 prev_head: Some(state(2)),
434 head: None,
435 thread: Some("main".into()),
436 },
437 OpRecord::Goto {
438 target: state(3),
439 prev_head: Some(state(2)),
440 head: state(3),
441 },
442 OpRecord::ThreadCreate {
443 name: "topic".into(),
444 state: state(4),
445 manager_snapshot: Some(vec![1, 2, 3]),
446 },
447 OpRecord::ThreadDelete {
448 name: "old".into(),
449 state: state(5),
450 },
451 OpRecord::ThreadUpdate {
452 name: "main".into(),
453 old_state: state(6),
454 new_state: state(7),
455 manager_snapshots: ThreadUpdateSnapshots::from_record_sets(
456 Some(vec![6]),
457 Some(vec![7]),
458 vec![vec![60], vec![61]],
459 vec![vec![70]],
460 true,
461 ),
462 },
463 OpRecord::Fork {
464 from: state(8),
465 new_state: state(9),
466 thread: Some("topic".into()),
467 head: None,
468 },
469 OpRecord::Collapse {
470 sources: vec![state(8), state(9)],
471 result: state(10),
472 thread: Some("main".into()),
473 pre_thread_state: Some(state(7)),
474 },
475 OpRecord::MarkerCreate {
476 name: "release".into(),
477 state: state(11),
478 },
479 OpRecord::MarkerDelete {
480 name: "draft".into(),
481 state: state(12),
482 },
483 OpRecord::Checkpoint {
484 parent: Some(state(12)),
485 state: state(13),
486 thread: Some("main".into()),
487 },
488 OpRecord::TransactionAbort {
489 transaction_id: "abort".into(),
490 reason: "reason".into(),
491 },
492 OpRecord::EphemeralThreadCollapse {
493 thread: "ephemeral".into(),
494 final_state: state(14),
495 },
496 OpRecord::ConflictResolved {
497 conflict_id: "conflict".into(),
498 resolution: "ours".into(),
499 resolver: Attribution::with_agent(
500 Principal::new("Resolver", "resolver@example.com"),
501 Agent::new("openai", "gpt-5-codex"),
502 ),
503 mode: ConflictResolutionMode::Ours,
504 },
505 OpRecord::TransactionCommit {
506 transaction_id: "tx".into(),
507 op_count: 2,
508 },
509 OpRecord::Redact {
510 redaction_id: hash(1),
511 blob: hash(2),
512 state: state(15),
513 path: "secret.txt".into(),
514 },
515 OpRecord::Purge {
516 redaction_id: hash(3),
517 blob: hash(4),
518 },
519 OpRecord::FastForward {
520 source_thread: "feature".into(),
521 target_thread: "main".into(),
522 pre_target_id: state(17),
523 post_target_id: state(18),
524 },
525 OpRecord::GitCheckpoint {
526 branch: "main".into(),
527 state: state(20),
528 previous_git_oid: Some("abc".into()),
529 new_git_oid: "def".into(),
530 },
531 OpRecord::RemoteThreadUpdate {
532 remote: "origin".into(),
533 thread: "main".into(),
534 state: state(21),
535 },
536 OpRecord::RemoteThreadDelete {
537 remote: "origin".into(),
538 thread: "old".into(),
539 state: state(22),
540 },
541 OpRecord::UndoRecoveryUpdate { state: state(23) },
542 OpRecord::StateVisibilitySet {
543 state: state(24),
544 record_id: hash(5),
545 tier: VisibilityTier::Internal,
546 prior_sidecar: None,
547 new_sidecar: Some(vec![1, 2, 3]),
548 },
549 OpRecord::StateVisibilityPromote {
550 state: state(25),
551 superseded: hash(6),
552 record_id: hash(7),
553 tier: VisibilityTier::Restricted {
554 scope_label: "embargo".into(),
555 },
556 prior_sidecar: Some(vec![4]),
557 new_sidecar: Some(vec![5]),
558 },
559 OpRecord::HeadUpdate {
560 previous: RecordedHead::Detached { state: state(26) },
561 new: RecordedHead::Attached {
562 thread: "main".into(),
563 },
564 },
565 OpRecord::EntryVisibilitySet {
566 change_id: crate::object::ChangeId::from_bytes([7u8; 16]),
567 record_id: hash(8),
568 prior_sidecar: None,
569 new_sidecar: Some(vec![9, 9, 9]),
570 },
571 ]
572 }
573
574 fn variant_name(record: &OpRecord) -> &'static str {
575 match record {
576 OpRecord::Snapshot { .. } => "Snapshot",
577 OpRecord::Goto { .. } => "Goto",
578 OpRecord::ThreadCreate { .. } => "ThreadCreate",
579 OpRecord::ThreadDelete { .. } => "ThreadDelete",
580 OpRecord::ThreadUpdate { .. } => "ThreadUpdate",
581 OpRecord::Fork { .. } => "Fork",
582 OpRecord::Collapse { .. } => "Collapse",
583 OpRecord::MarkerCreate { .. } => "MarkerCreate",
584 OpRecord::MarkerDelete { .. } => "MarkerDelete",
585 OpRecord::Checkpoint { .. } => "Checkpoint",
586 OpRecord::TransactionAbort { .. } => "TransactionAbort",
587 OpRecord::EphemeralThreadCollapse { .. } => "EphemeralThreadCollapse",
588 OpRecord::ConflictResolved { .. } => "ConflictResolved",
589 OpRecord::TransactionCommit { .. } => "TransactionCommit",
590 OpRecord::Redact { .. } => "Redact",
591 OpRecord::Purge { .. } => "Purge",
592 OpRecord::FastForward { .. } => "FastForward",
593 OpRecord::GitCheckpoint { .. } => "GitCheckpoint",
594 OpRecord::RemoteThreadUpdate { .. } => "RemoteThreadUpdate",
595 OpRecord::RemoteThreadDelete { .. } => "RemoteThreadDelete",
596 OpRecord::UndoRecoveryUpdate { .. } => "UndoRecoveryUpdate",
597 OpRecord::StateVisibilitySet { .. } => "StateVisibilitySet",
598 OpRecord::StateVisibilityPromote { .. } => "StateVisibilityPromote",
599 OpRecord::HeadUpdate { .. } => "HeadUpdate",
600 OpRecord::EntryVisibilitySet { .. } => "EntryVisibilitySet",
601 }
602 }
603
604 #[test]
605 fn schema_four_is_current_and_legacy_versions_are_refused() {
606 assert_eq!(CURRENT_OP_RECORD_SCHEMA_VERSION, 4);
607 validate_op_record_schema_version(4).unwrap();
608 for legacy in 1..=3 {
609 let error = validate_op_record_schema_version(legacy).unwrap_err();
610 assert!(matches!(
611 error,
612 HeddleError::StorageFormatTooOld {
613 found,
614 required: 4,
615 ..
616 } if found == legacy
617 ));
618 }
619 assert!(matches!(
620 validate_op_record_schema_version(5).unwrap_err(),
621 HeddleError::StorageFormatTooNew {
622 found: 5,
623 supported: 4,
624 ..
625 }
626 ));
627 }
628
629 #[test]
630 fn every_current_variant_round_trips() {
631 let records = canonical_current_records();
632 assert_eq!(
633 records.iter().map(variant_name).collect::<Vec<_>>(),
634 [
635 "Snapshot",
636 "Goto",
637 "ThreadCreate",
638 "ThreadDelete",
639 "ThreadUpdate",
640 "Fork",
641 "Collapse",
642 "MarkerCreate",
643 "MarkerDelete",
644 "Checkpoint",
645 "TransactionAbort",
646 "EphemeralThreadCollapse",
647 "ConflictResolved",
648 "TransactionCommit",
649 "Redact",
650 "Purge",
651 "FastForward",
652 "GitCheckpoint",
653 "RemoteThreadUpdate",
654 "RemoteThreadDelete",
655 "UndoRecoveryUpdate",
656 "StateVisibilitySet",
657 "StateVisibilityPromote",
658 "HeadUpdate",
659 "EntryVisibilitySet",
660 ]
661 );
662 for record in records {
663 assert_round_trip(record);
664 }
665 }
666
667 #[test]
668 fn state_id_v4_visibility_tail_bytes_are_frozen() {
669 let record = OpRecord::StateVisibilityPromote {
670 state: state(1),
671 superseded: hash(2),
672 record_id: hash(3),
673 tier: VisibilityTier::Internal,
674 prior_sidecar: Some(vec![4]),
675 new_sidecar: Some(vec![5]),
676 };
677
678 let expected = [
679 &[
680 129, 182, 83, 116, 97, 116, 101, 86, 105, 115, 105, 98, 105, 108, 105, 116, 121,
681 80, 114, 111, 109, 111, 116, 101, 150, 220, 0, 32,
682 ][..],
683 &[1; 32],
684 &[220, 0, 32],
685 &[2; 32],
686 &[220, 0, 32],
687 &[3; 32],
688 &[168, 73, 110, 116, 101, 114, 110, 97, 108, 145, 4, 145, 5],
689 ]
690 .concat();
691
692 assert_eq!(encode_current_record(&record).unwrap(), expected);
693 }
694
695 #[test]
696 fn historical_sixteen_byte_payload_is_not_a_state_id_record() {
697 let historical = [
698 129, 168, 67, 111, 108, 108, 97, 112, 115, 101, 147, 146, 220, 0, 16, 10, 10, 10, 10,
699 10, 10, 10, 10, 10, 10, 10, 10, 10, 10, 10, 10, 220, 0, 16, 11, 11, 11, 11, 11, 11, 11,
700 11, 11, 11, 11, 11, 11, 11, 11, 11, 220, 0, 16, 12, 12, 12, 12, 12, 12, 12, 12, 12, 12,
701 12, 12, 12, 12, 12, 12, 164, 109, 97, 105, 110,
702 ];
703 let error = decode_current_record(&historical)
704 .expect_err("16-byte ChangeIds must not decode as StateIds");
705 assert!(error.to_string().contains("expected an array of length 32"));
706 }
707}