Skip to main content

heddle_object_model/op_record/
codec.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Versioned codec for current OpRecord payloads stored in the packed oplog.
3//!
4//! Repository format v3 changed state identity from 16-byte logical ChangeIds
5//! to 32-byte content-addressed StateIds. That boundary is intentionally not
6//! migratable: repository open refuses older formats, and the oplog codec
7//! refuses record schema versions 1 through 3 without rewriting their bytes.
8//! Future incompatible record changes must allocate a new schema version.
9
10use 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/// Strict snapshot of record schema version 4.
61///
62/// This mirror has no defaults for structural fields: malformed or older
63/// positional payloads fail instead of being interpreted as current records.
64#[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}