1use std::cell::RefCell;
44use std::collections::BTreeMap;
45use std::path::{Path, PathBuf};
46use std::sync::atomic::{AtomicU64, AtomicU8, Ordering};
47use std::sync::{Arc, Mutex};
48
49use serde::{Deserialize, Serialize};
50
51use crate::clock_mock;
52
53pub const TAPE_FORMAT_VERSION: u32 = 1;
56
57pub const MAX_INLINE_BYTES: usize = 4 * 1024;
62
63#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
66pub struct TapeHeader {
67 pub version: u32,
68 pub harn_version: String,
72 #[serde(default)]
76 pub started_at_unix_ms: Option<i64>,
77 #[serde(default)]
80 pub script_path: Option<String>,
81 #[serde(default)]
84 pub argv: Vec<String>,
85}
86
87impl TapeHeader {
88 pub fn current(
89 started_at_unix_ms: Option<i64>,
90 script_path: Option<String>,
91 argv: Vec<String>,
92 ) -> Self {
93 Self {
94 version: TAPE_FORMAT_VERSION,
95 harn_version: env!("CARGO_PKG_VERSION").to_string(),
96 started_at_unix_ms,
97 script_path,
98 argv,
99 }
100 }
101}
102
103#[derive(Debug, Clone, Serialize, Deserialize)]
107#[serde(tag = "type", rename_all = "snake_case")]
108enum TapeLine {
109 Header(TapeHeader),
110 Record(TapeRecord),
111}
112
113#[derive(Debug, Clone, Serialize, Deserialize)]
117pub struct TapeRecord {
118 pub seq: u64,
120 #[serde(default)]
124 pub phase: TapePhase,
125 pub virtual_time_ms: i64,
128 pub monotonic_ms: i64,
132 pub kind: TapeRecordKind,
134}
135
136#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
138#[serde(rename_all = "snake_case")]
139pub enum TapePhase {
140 #[default]
142 UserScript,
143 RuntimeFinalize,
146}
147
148impl TapePhase {
149 fn as_u8(self) -> u8 {
150 match self {
151 Self::UserScript => 0,
152 Self::RuntimeFinalize => 1,
153 }
154 }
155
156 fn from_u8(value: u8) -> Self {
157 match value {
158 1 => Self::RuntimeFinalize,
159 _ => Self::UserScript,
160 }
161 }
162
163 pub fn label(self) -> &'static str {
164 match self {
165 Self::UserScript => "user_script",
166 Self::RuntimeFinalize => "runtime_finalize",
167 }
168 }
169}
170
171#[derive(Debug, Clone, Serialize, Deserialize)]
176#[serde(tag = "kind", rename_all = "snake_case")]
177pub enum TapeRecordKind {
178 ClockRead { source: ClockSource, value_ms: i64 },
183 ClockSleep { duration_ms: u64 },
186 LlmCall {
191 request_digest: String,
192 response: TapePayload,
193 },
194 FileRead {
198 path: String,
199 content_hash: String,
200 len_bytes: u64,
201 },
202 FileWrite {
204 path: String,
205 content_hash: String,
206 len_bytes: u64,
207 },
208 FileDelete { path: String },
210 ProcessSpawn {
214 program: String,
215 args: Vec<String>,
216 cwd: Option<String>,
217 exit_code: i32,
218 duration_ms: u64,
219 stdout_payload: TapePayload,
220 stderr_payload: TapePayload,
221 },
222 McpJsonRpc {
226 server: String,
227 method: String,
228 request_digest: String,
229 response_digest: String,
230 latency_ms: u64,
231 request_payload: TapePayload,
232 response_payload: TapePayload,
233 },
234 ModelJob {
239 job_id: String,
240 request_id: String,
241 backend: String,
242 state: String,
243 event_kind: String,
244 event: TapePayload,
245 asset_digests: Vec<String>,
246 },
247 #[serde(other)]
251 Unknown,
252}
253
254impl TapeRecordKind {
255 pub fn label(&self) -> &'static str {
260 match self {
261 Self::ClockRead { .. } => "clock_read",
262 Self::ClockSleep { .. } => "clock_sleep",
263 Self::LlmCall { .. } => "llm_call",
264 Self::FileRead { .. } => "file_read",
265 Self::FileWrite { .. } => "file_write",
266 Self::FileDelete { .. } => "file_delete",
267 Self::ProcessSpawn { .. } => "process_spawn",
268 Self::McpJsonRpc { .. } => "mcp_json_rpc",
269 Self::ModelJob { .. } => "model_job",
270 Self::Unknown => "unknown",
271 }
272 }
273}
274
275#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
278#[serde(rename_all = "snake_case")]
279pub enum ClockSource {
280 Wall,
281 Monotonic,
282}
283
284#[derive(Debug, Clone, Serialize, Deserialize)]
287#[serde(untagged)]
288pub enum TapePayload {
289 Inline { content_hash: String, text: String },
292 Cas {
294 content_hash: String,
295 len_bytes: u64,
296 },
297}
298
299impl TapePayload {
300 pub fn content_hash(&self) -> &str {
301 match self {
302 Self::Inline { content_hash, .. } | Self::Cas { content_hash, .. } => content_hash,
303 }
304 }
305
306 pub fn len_bytes(&self) -> u64 {
307 match self {
308 Self::Inline { text, .. } => text.len() as u64,
309 Self::Cas { len_bytes, .. } => *len_bytes,
310 }
311 }
312}
313
314pub fn content_hash(bytes: &[u8]) -> String {
317 blake3::hash(bytes).to_hex().to_string()
318}
319
320fn build_payload(bytes: Vec<u8>, cas: &mut BTreeMap<String, Vec<u8>>) -> TapePayload {
323 let hash = content_hash(&bytes);
324 if bytes.len() > MAX_INLINE_BYTES {
325 let len_bytes = bytes.len() as u64;
326 cas.entry(hash.clone()).or_insert(bytes);
327 TapePayload::Cas {
328 content_hash: hash,
329 len_bytes,
330 }
331 } else {
332 let text = match String::from_utf8(bytes) {
333 Ok(text) => text,
334 Err(error) => {
335 let bytes = error.into_bytes();
339 let len_bytes = bytes.len() as u64;
340 cas.entry(hash.clone()).or_insert(bytes);
341 return TapePayload::Cas {
342 content_hash: hash,
343 len_bytes,
344 };
345 }
346 };
347 TapePayload::Inline {
348 content_hash: hash,
349 text,
350 }
351 }
352}
353
354#[derive(Debug, Clone)]
358pub struct EventTape {
359 pub header: TapeHeader,
360 pub records: Vec<TapeRecord>,
361 cas: BTreeMap<String, Vec<u8>>,
365}
366
367impl EventTape {
368 pub fn new(header: TapeHeader) -> Self {
369 Self {
370 header,
371 records: Vec::new(),
372 cas: BTreeMap::new(),
373 }
374 }
375
376 pub fn resolve_payload(&self, payload: &TapePayload) -> Result<Vec<u8>, String> {
379 match payload {
380 TapePayload::Inline { text, .. } => Ok(text.as_bytes().to_vec()),
381 TapePayload::Cas { content_hash, .. } => self
382 .cas
383 .get(content_hash)
384 .cloned()
385 .ok_or_else(|| format!("tape CAS missing entry for {content_hash}")),
386 }
387 }
388
389 pub fn cas_len(&self) -> usize {
391 self.cas.len()
392 }
393
394 pub fn persist(&self, path: &Path) -> Result<(), String> {
397 if let Some(parent) = path.parent() {
398 if !parent.as_os_str().is_empty() {
399 std::fs::create_dir_all(parent)
400 .map_err(|err| format!("mkdir {}: {err}", parent.display()))?;
401 }
402 }
403
404 let mut body = String::new();
405 let header_line = serde_json::to_string(&TapeLine::Header(self.header.clone()))
406 .map_err(|err| format!("serialize tape header: {err}"))?;
407 body.push_str(&header_line);
408 body.push('\n');
409 for record in &self.records {
410 let line = serde_json::to_string(&TapeLine::Record(record.clone()))
411 .map_err(|err| format!("serialize tape record: {err}"))?;
412 body.push_str(&line);
413 body.push('\n');
414 }
415 std::fs::write(path, body).map_err(|err| format!("write {}: {err}", path.display()))?;
416
417 if !self.cas.is_empty() {
418 let cas_dir = cas_dir_for(path);
419 std::fs::create_dir_all(&cas_dir)
420 .map_err(|err| format!("mkdir {}: {err}", cas_dir.display()))?;
421 for (hash, bytes) in &self.cas {
422 let entry = cas_dir.join(hash);
423 std::fs::write(&entry, bytes)
424 .map_err(|err| format!("write {}: {err}", entry.display()))?;
425 }
426 }
427 Ok(())
428 }
429
430 pub fn load(path: &Path) -> Result<Self, String> {
433 let body = std::fs::read_to_string(path)
434 .map_err(|err| format!("read {}: {err}", path.display()))?;
435 let mut lines = body.lines();
436 let first_line = lines
437 .next()
438 .ok_or_else(|| format!("empty tape file: {}", path.display()))?;
439 let header_line: TapeLine = serde_json::from_str(first_line)
440 .map_err(|err| format!("parse tape header in {}: {err}", path.display()))?;
441 let header = match header_line {
442 TapeLine::Header(header) => header,
443 TapeLine::Record(_) => {
444 return Err(format!(
445 "tape {} is missing its header (first line is a record)",
446 path.display()
447 ))
448 }
449 };
450 if header.version > TAPE_FORMAT_VERSION {
451 return Err(format!(
452 "tape {} declares version {} but this runtime supports up to {TAPE_FORMAT_VERSION}",
453 path.display(),
454 header.version
455 ));
456 }
457 let mut records = Vec::new();
458 for (idx, line) in lines.enumerate() {
459 let trimmed = line.trim();
460 if trimmed.is_empty() {
461 continue;
462 }
463 let parsed: TapeLine = serde_json::from_str(trimmed).map_err(|err| {
464 format!(
465 "parse tape record at line {} in {}: {err}",
466 idx + 2,
467 path.display()
468 )
469 })?;
470 match parsed {
471 TapeLine::Record(record) => records.push(record),
472 TapeLine::Header(_) => {
473 return Err(format!(
474 "tape {} contains a second header at line {}",
475 path.display(),
476 idx + 2
477 ))
478 }
479 }
480 }
481
482 let mut cas = BTreeMap::new();
483 let cas_dir = cas_dir_for(path);
484 if cas_dir.is_dir() {
485 for record in &records {
486 visit_payloads(&record.kind, |payload| {
487 if let TapePayload::Cas { content_hash, .. } = payload {
488 if cas.contains_key(content_hash) {
489 return;
490 }
491 let entry = cas_dir.join(content_hash);
492 if let Ok(bytes) = std::fs::read(&entry) {
493 cas.insert(content_hash.clone(), bytes);
494 }
495 }
496 });
497 }
498 }
499 Ok(Self {
500 header,
501 records,
502 cas,
503 })
504 }
505}
506
507fn cas_dir_for(tape_path: &Path) -> PathBuf {
508 let mut buf = tape_path.as_os_str().to_owned();
509 buf.push(".cas");
510 PathBuf::from(buf)
511}
512
513fn visit_payloads(kind: &TapeRecordKind, mut visit: impl FnMut(&TapePayload)) {
514 match kind {
515 TapeRecordKind::LlmCall { response, .. } => visit(response),
516 TapeRecordKind::ProcessSpawn {
517 stdout_payload,
518 stderr_payload,
519 ..
520 } => {
521 visit(stdout_payload);
522 visit(stderr_payload);
523 }
524 TapeRecordKind::McpJsonRpc {
525 request_payload,
526 response_payload,
527 ..
528 } => {
529 visit(request_payload);
530 visit(response_payload);
531 }
532 TapeRecordKind::ModelJob { event, .. } => visit(event),
533 TapeRecordKind::ClockRead { .. }
534 | TapeRecordKind::ClockSleep { .. }
535 | TapeRecordKind::FileRead { .. }
536 | TapeRecordKind::FileWrite { .. }
537 | TapeRecordKind::FileDelete { .. }
538 | TapeRecordKind::Unknown => {}
539 }
540}
541
542#[derive(Debug)]
547pub struct TapeRecorder {
548 next_seq: AtomicU64,
549 phase: AtomicU8,
550 started_at: clock_mock::ClockInstant,
551 inner: Mutex<RecorderInner>,
552}
553
554#[derive(Debug, Default)]
555struct RecorderInner {
556 records: Vec<TapeRecord>,
557 cas: BTreeMap<String, Vec<u8>>,
558}
559
560impl Default for TapeRecorder {
561 fn default() -> Self {
562 Self::new()
563 }
564}
565
566impl TapeRecorder {
567 pub fn new() -> Self {
568 Self {
569 next_seq: AtomicU64::new(0),
570 phase: AtomicU8::new(TapePhase::UserScript.as_u8()),
571 started_at: clock_mock::instant_now(),
572 inner: Mutex::new(RecorderInner::default()),
573 }
574 }
575
576 pub fn record(&self, kind: TapeRecordKind) {
579 let virtual_time_ms = clock_mock::now_ms();
580 let monotonic_ms = clock_mock::instant_now()
581 .duration_since(self.started_at)
582 .as_millis()
583 .min(i64::MAX as u128) as i64;
584 self.record_at(kind, virtual_time_ms, monotonic_ms);
585 }
586
587 fn record_at(&self, kind: TapeRecordKind, virtual_time_ms: i64, monotonic_ms: i64) {
588 let seq = self.next_seq.fetch_add(1, Ordering::SeqCst);
589 let record = TapeRecord {
590 seq,
591 phase: TapePhase::from_u8(self.phase.load(Ordering::SeqCst)),
592 virtual_time_ms,
593 monotonic_ms,
594 kind,
595 };
596 self.inner
597 .lock()
598 .expect("tape recorder mutex poisoned")
599 .records
600 .push(record);
601 }
602
603 fn swap_phase(&self, phase: TapePhase) -> TapePhase {
604 TapePhase::from_u8(self.phase.swap(phase.as_u8(), Ordering::SeqCst))
605 }
606
607 pub fn payload_from_bytes(&self, bytes: Vec<u8>) -> TapePayload {
612 let mut inner = self.inner.lock().expect("tape recorder mutex poisoned");
613 build_payload(bytes, &mut inner.cas)
614 }
615
616 pub fn snapshot(&self, header: TapeHeader) -> EventTape {
621 let inner = self.inner.lock().expect("tape recorder mutex poisoned");
622 EventTape {
623 header,
624 records: inner.records.clone(),
625 cas: inner.cas.clone(),
626 }
627 }
628}
629
630thread_local! {
631 static ACTIVE_RECORDER: RefCell<Option<Arc<TapeRecorder>>> = const { RefCell::new(None) };
632}
633
634pub struct TapeRecorderGuard {
637 previous: Option<Arc<TapeRecorder>>,
638}
639
640impl Drop for TapeRecorderGuard {
641 fn drop(&mut self) {
642 let prev = self.previous.take();
643 ACTIVE_RECORDER.with(|slot| {
644 *slot.borrow_mut() = prev;
645 });
646 }
647}
648
649pub fn install_recorder(recorder: Arc<TapeRecorder>) -> TapeRecorderGuard {
650 let previous = ACTIVE_RECORDER.with(|slot| slot.replace(Some(recorder)));
651 TapeRecorderGuard { previous }
652}
653
654pub fn active_recorder() -> Option<Arc<TapeRecorder>> {
657 ACTIVE_RECORDER.with(|slot| slot.borrow().clone())
658}
659
660pub fn record_model_job_event(payload: &serde_json::Value) {
665 let Some(recorder) = active_recorder() else {
666 return;
667 };
668 let policy = crate::redact::current_policy();
669 let redacted = policy.redact_json(payload);
670 let event_bytes = serde_json::to_vec(&redacted).unwrap_or_default();
671 let event = recorder.payload_from_bytes(event_bytes);
672 let job_id = redacted
673 .get("job_id")
674 .and_then(|value| value.as_str())
675 .unwrap_or("")
676 .to_string();
677 let request_id = redacted
678 .get("request_id")
679 .and_then(|value| value.as_str())
680 .unwrap_or("")
681 .to_string();
682 let backend = redacted
683 .get("backend")
684 .and_then(|value| value.as_str())
685 .unwrap_or("")
686 .to_string();
687 let state = redacted
688 .get("state")
689 .and_then(|value| value.as_str())
690 .unwrap_or("")
691 .to_string();
692 let event_kind = redacted
693 .get("kind")
694 .and_then(|value| value.as_str())
695 .unwrap_or("")
696 .to_string();
697 let mut asset_digests = Vec::new();
698 if let Some(assets) = redacted.get("assets").and_then(|value| value.as_array()) {
699 for asset in assets {
700 if let Some(digest) = asset
701 .get("sha256")
702 .and_then(|value| value.as_str())
703 .filter(|value| !value.is_empty())
704 {
705 asset_digests.push(digest.to_string());
706 } else if let Some(uri) = asset.get("uri").and_then(|value| value.as_str()) {
707 if let Some(digest) = uri.strip_prefix("asset://sha256/") {
708 if !digest.is_empty() {
709 asset_digests.push(digest.to_string());
710 }
711 }
712 }
713 }
714 }
715 if let Some(digest) = redacted
716 .get("asset_digest")
717 .and_then(|value| value.as_str())
718 .filter(|value| !value.is_empty())
719 {
720 asset_digests.push(digest.to_string());
721 }
722 asset_digests.sort();
723 asset_digests.dedup();
724 recorder.record(TapeRecordKind::ModelJob {
725 job_id,
726 request_id,
727 backend,
728 state,
729 event_kind,
730 event,
731 asset_digests,
732 });
733}
734
735pub fn record_mcp_json_rpc(
739 server: &str,
740 method: &str,
741 request: &serde_json::Value,
742 response: &serde_json::Value,
743 latency_ms: u64,
744) {
745 let Some(recorder) = active_recorder() else {
746 return;
747 };
748 let policy = crate::redact::current_policy();
749 let request = policy.redact_json(request);
750 let response = policy.redact_json(response);
751 let request_bytes = serde_json::to_vec(&request).unwrap_or_default();
752 let response_bytes = serde_json::to_vec(&response).unwrap_or_default();
753 let request_digest = content_hash(&request_bytes);
754 let response_digest = content_hash(&response_bytes);
755 let request_payload = recorder.payload_from_bytes(request_bytes);
756 let response_payload = recorder.payload_from_bytes(response_bytes);
757 recorder.record(TapeRecordKind::McpJsonRpc {
758 server: server.to_string(),
759 method: method.to_string(),
760 request_digest,
761 response_digest,
762 latency_ms,
763 request_payload,
764 response_payload,
765 });
766}
767
768pub struct TapePhaseGuard {
771 recorder: Arc<TapeRecorder>,
772 previous: TapePhase,
773}
774
775impl Drop for TapePhaseGuard {
776 fn drop(&mut self) {
777 self.recorder.swap_phase(self.previous);
778 }
779}
780
781pub fn enter_phase(phase: TapePhase) -> Option<TapePhaseGuard> {
784 let recorder = active_recorder()?;
785 let previous = recorder.swap_phase(phase);
786 Some(TapePhaseGuard { recorder, previous })
787}
788
789pub fn with_active_recorder<F>(build: F)
792where
793 F: FnOnce(&Arc<TapeRecorder>) -> Option<TapeRecordKind>,
794{
795 let Some(recorder) = active_recorder() else {
796 return;
797 };
798 if let Some(kind) = build(&recorder) {
799 recorder.record(kind);
800 }
801}
802
803pub fn with_active_recorder_clock<F>(clock: &dyn harn_clock::Clock, build: F)
809where
810 F: FnOnce(&Arc<TapeRecorder>) -> Option<TapeRecordKind>,
811{
812 let Some(recorder) = active_recorder() else {
813 return;
814 };
815 if let Some(kind) = build(&recorder) {
816 recorder.record_at(kind, harn_clock::now_wall_ms(clock), clock.monotonic_ms());
817 }
818}
819
820#[cfg(test)]
821mod tests {
822 use super::*;
823 use tempfile::TempDir;
824
825 fn small_record(seq: u64, dur: u64) -> TapeRecord {
826 TapeRecord {
827 seq,
828 phase: TapePhase::UserScript,
829 virtual_time_ms: seq as i64 * 1000,
830 monotonic_ms: seq as i64 * 1000,
831 kind: TapeRecordKind::ClockSleep { duration_ms: dur },
832 }
833 }
834
835 #[test]
836 fn round_trip_inline_records() {
837 let temp = TempDir::new().unwrap();
838 let path = temp.path().join("run.tape");
839 let mut tape = EventTape::new(TapeHeader::current(
840 Some(1_700_000_000_000),
841 Some("script.harn".to_string()),
842 vec!["a".into()],
843 ));
844 tape.records.push(small_record(0, 250));
845 tape.records.push(small_record(1, 750));
846 tape.persist(&path).unwrap();
847
848 let loaded = EventTape::load(&path).unwrap();
849 assert_eq!(loaded.header.version, TAPE_FORMAT_VERSION);
850 assert_eq!(loaded.header.argv, vec!["a".to_string()]);
851 assert_eq!(loaded.records.len(), 2);
852 match &loaded.records[0].kind {
853 TapeRecordKind::ClockSleep { duration_ms } => assert_eq!(*duration_ms, 250),
854 other => panic!("unexpected: {other:?}"),
855 }
856 }
857
858 #[test]
859 fn recorder_phase_guard_stamps_and_restores() {
860 let recorder = Arc::new(TapeRecorder::new());
861 let _recorder_guard = install_recorder(Arc::clone(&recorder));
862
863 with_active_recorder(|_| Some(TapeRecordKind::ClockSleep { duration_ms: 1 }));
864 {
865 let _phase_guard = enter_phase(TapePhase::RuntimeFinalize).unwrap();
866 with_active_recorder(|_| Some(TapeRecordKind::ClockSleep { duration_ms: 2 }));
867 }
868 with_active_recorder(|_| Some(TapeRecordKind::ClockSleep { duration_ms: 3 }));
869
870 let tape = recorder.snapshot(TapeHeader::current(None, None, Vec::new()));
871 let phases = tape
872 .records
873 .iter()
874 .map(|record| record.phase)
875 .collect::<Vec<_>>();
876 assert_eq!(
877 phases,
878 vec![
879 TapePhase::UserScript,
880 TapePhase::RuntimeFinalize,
881 TapePhase::UserScript
882 ]
883 );
884 }
885
886 #[test]
887 fn large_payloads_spill_to_cas_and_round_trip() {
888 let temp = TempDir::new().unwrap();
889 let path = temp.path().join("run.tape");
890 let mut tape = EventTape::new(TapeHeader::current(None, None, Vec::new()));
891 let big = vec![b'x'; MAX_INLINE_BYTES + 32];
892 let payload = build_payload(big.clone(), &mut tape.cas);
893 let hash = payload.content_hash().to_string();
894 let kind = TapeRecordKind::ProcessSpawn {
895 program: "/bin/echo".to_string(),
896 args: vec!["x".to_string()],
897 cwd: None,
898 exit_code: 0,
899 duration_ms: 1,
900 stdout_payload: payload,
901 stderr_payload: build_payload(Vec::new(), &mut tape.cas),
902 };
903 tape.records.push(TapeRecord {
904 seq: 0,
905 phase: TapePhase::UserScript,
906 virtual_time_ms: 0,
907 monotonic_ms: 0,
908 kind,
909 });
910 tape.persist(&path).unwrap();
911
912 assert!(path.with_extension("tape.cas").exists() || cas_dir_for(&path).exists());
914 let cas_dir = cas_dir_for(&path);
915 assert!(cas_dir.join(&hash).exists());
916
917 let loaded = EventTape::load(&path).unwrap();
918 let resolved = match &loaded.records[0].kind {
919 TapeRecordKind::ProcessSpawn { stdout_payload, .. } => {
920 loaded.resolve_payload(stdout_payload).unwrap()
921 }
922 other => panic!("unexpected: {other:?}"),
923 };
924 assert_eq!(resolved.len(), big.len());
925 }
926
927 #[test]
928 fn rejects_newer_version() {
929 let temp = TempDir::new().unwrap();
930 let path = temp.path().join("future.tape");
931 std::fs::write(
932 &path,
933 r#"{"type":"header","version":99,"harn_version":"x","started_at_unix_ms":null,"script_path":null,"argv":[]}
934"#,
935 )
936 .unwrap();
937 let err = EventTape::load(&path).unwrap_err();
938 assert!(err.contains("version 99"), "{err}");
939 }
940
941 #[test]
942 fn recorder_assigns_monotonic_seq() {
943 let recorder = Arc::new(TapeRecorder::new());
944 recorder.record(TapeRecordKind::ClockSleep { duration_ms: 1 });
945 recorder.record(TapeRecordKind::ClockSleep { duration_ms: 2 });
946 let snapshot = recorder.snapshot(TapeHeader::current(None, None, Vec::new()));
947 assert_eq!(snapshot.records[0].seq, 0);
948 assert_eq!(snapshot.records[1].seq, 1);
949 }
950
951 #[test]
952 fn records_model_job_events_with_asset_digests() {
953 let recorder = Arc::new(TapeRecorder::new());
954 let _guard = install_recorder(Arc::clone(&recorder));
955 record_model_job_event(&serde_json::json!({
956 "schema": "harn.model_job_event.v1",
957 "kind": "output",
958 "job_id": "job-1",
959 "request_id": "req-1",
960 "backend": "fixture",
961 "state": "succeeded",
962 "assets": [
963 {
964 "uri": "asset://sha256/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
965 "sha256": "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
966 }
967 ]
968 }));
969 let snapshot = recorder.snapshot(TapeHeader::current(None, None, Vec::new()));
970 assert_eq!(snapshot.records.len(), 1);
971 match &snapshot.records[0].kind {
972 TapeRecordKind::ModelJob {
973 job_id,
974 event_kind,
975 asset_digests,
976 ..
977 } => {
978 assert_eq!(job_id, "job-1");
979 assert_eq!(event_kind, "output");
980 assert_eq!(
981 asset_digests,
982 &vec![
983 "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
984 .to_string()
985 ]
986 );
987 }
988 other => panic!("unexpected tape kind: {other:?}"),
989 }
990 }
991}