Skip to main content

atman_runtime/
permission_audit.rs

1use std::collections::BTreeSet;
2
3use chrono::{DateTime, Utc};
4use serde::{Deserialize, Serialize};
5
6use crate::event::FlowRunId;
7use crate::permission::{
8    ApprovalTarget, DecisionActor, EscalationHop, ExecutionBoundary, GrantScope, GroupOwner,
9    PermissionDecision, PermissionGrant, PermissionGroup, PermissionGroupId, PermissionRequest,
10    PermissionRequestId, PermissionRequestState,
11};
12use crate::tool::Tier;
13use crate::trust::{ExecutionPolicy, PolicyResolution, RiskKind, TrustConfig};
14
15#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
16#[serde(tag = "kind", rename_all = "snake_case")]
17pub enum PermissionProjectionActor {
18    Policy {
19        policy_version: String,
20        rule_id: String,
21    },
22    Flow {
23        session_id: String,
24        run_id: FlowRunId,
25    },
26    User {
27        session_id: String,
28        principal_id: Option<String>,
29    },
30    System {
31        component: String,
32    },
33    UnknownLegacy {
34        label: String,
35    },
36}
37
38impl From<&DecisionActor> for PermissionProjectionActor {
39    fn from(actor: &DecisionActor) -> Self {
40        match actor {
41            DecisionActor::Policy {
42                policy_version,
43                rule_id,
44            } => Self::Policy {
45                policy_version: policy_version.clone(),
46                rule_id: rule_id.clone(),
47            },
48            DecisionActor::Flow { session_id, run_id } => Self::Flow {
49                session_id: session_id.clone(),
50                run_id: run_id.clone(),
51            },
52            DecisionActor::User {
53                session_id,
54                principal_id,
55            } => Self::User {
56                session_id: session_id.clone(),
57                principal_id: principal_id.clone(),
58            },
59            DecisionActor::System { component } => Self::System {
60                component: component.clone(),
61            },
62        }
63    }
64}
65
66#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
67#[serde(tag = "kind", rename_all = "snake_case")]
68pub enum PermissionAuditTarget {
69    Flow { run_id: FlowRunId },
70    User,
71}
72
73impl From<&ApprovalTarget> for PermissionAuditTarget {
74    fn from(target: &ApprovalTarget) -> Self {
75        match target {
76            ApprovalTarget::Flow(run_id) => Self::Flow {
77                run_id: run_id.clone(),
78            },
79            ApprovalTarget::User => Self::User,
80        }
81    }
82}
83
84#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
85pub struct PermissionPolicyReference {
86    pub snapshot_id: String,
87    pub rule_id: String,
88}
89
90impl PermissionPolicyReference {
91    pub fn capture(
92        policy: &TrustConfig,
93        tier: Tier,
94        risks: &BTreeSet<RiskKind>,
95        execution_policy: ExecutionPolicy,
96        effective: PolicyResolution,
97    ) -> Self {
98        let bytes = serde_json::to_vec(policy).expect("TrustConfig is serializable");
99        let digest = blake3::hash(&bytes).to_hex().to_string();
100        let risk_names = risks
101            .iter()
102            .map(|risk| format!("{risk:?}"))
103            .collect::<Vec<_>>();
104        let policy_resolution = policy.resolve_policy_resolution(tier, risks.iter().copied());
105        Self {
106            snapshot_id: format!("blake3:{digest}"),
107            rule_id: format!(
108                "mode={:?};tier={tier:?};risks={};escalation={:?};escalation_step={:?};policy_result={:?};effective_result={:?};execution={execution_policy:?}",
109                policy.mode,
110                risk_names.join(","),
111                policy.escalation,
112                policy_resolution.escalation,
113                policy_resolution.action,
114                effective.action,
115            ),
116        }
117    }
118}
119
120#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq, Eq)]
121pub struct PermissionProvenanceSummary {
122    pub cwd: Option<String>,
123    pub path: Option<String>,
124    pub path_origin: Option<String>,
125    pub workspace_id: Option<String>,
126    pub workspace_root: Option<String>,
127    pub repository_root: Option<String>,
128    pub network: bool,
129    pub risks: BTreeSet<String>,
130    pub targets: Vec<String>,
131}
132
133#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
134pub struct PermissionEscalationAuditHop {
135    pub target: PermissionAuditTarget,
136    pub actor: Option<PermissionProjectionActor>,
137    pub action: Option<String>,
138    pub reason: Option<String>,
139    pub at: DateTime<Utc>,
140}
141
142impl From<&EscalationHop> for PermissionEscalationAuditHop {
143    fn from(hop: &EscalationHop) -> Self {
144        Self {
145            target: (&hop.target).into(),
146            actor: hop.actor.as_ref().map(Into::into),
147            action: hop
148                .action
149                .map(|action| format!("{action:?}").to_lowercase()),
150            reason: hop.reason.clone(),
151            at: hop.at,
152        }
153    }
154}
155
156#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
157#[serde(tag = "kind", rename_all = "snake_case")]
158pub enum PermissionAuditScope {
159    CurrentCall,
160    ChildRunSameTool {
161        run_id: FlowRunId,
162        tool_name: String,
163    },
164    ChildRunSamePathRule {
165        run_id: FlowRunId,
166        tool_name: String,
167        workspace_relative_path: String,
168    },
169}
170
171impl From<&GrantScope> for PermissionAuditScope {
172    fn from(scope: &GrantScope) -> Self {
173        match scope {
174            GrantScope::CurrentCall => Self::CurrentCall,
175            GrantScope::ChildRunSameTool { run_id, tool_name } => Self::ChildRunSameTool {
176                run_id: run_id.clone(),
177                tool_name: tool_name.clone(),
178            },
179            GrantScope::ChildRunSamePathRule {
180                run_id,
181                tool_name,
182                workspace_relative_path,
183            } => Self::ChildRunSamePathRule {
184                run_id: run_id.clone(),
185                tool_name: tool_name.clone(),
186                workspace_relative_path: workspace_relative_path.clone(),
187            },
188        }
189    }
190}
191
192#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
193pub struct PermissionRequestAudit {
194    pub request_id: Option<PermissionRequestId>,
195    #[serde(default)]
196    pub revision: u64,
197    pub session_id: String,
198    pub requesting_run_id: FlowRunId,
199    pub parent_run_id: Option<FlowRunId>,
200    pub root_run_id: FlowRunId,
201    pub tool_use_id: String,
202    pub tool: String,
203    #[serde(default, skip_serializing_if = "Option::is_none")]
204    pub call_intent: Option<crate::message::ToolCallIntent>,
205    pub tier: Tier,
206    /// Process boundary selected for this approval. Pending user approvals
207    /// report the boundary that an approval will grant; denied and non-process
208    /// requests have no execution boundary.
209    #[serde(default, skip_serializing_if = "Option::is_none")]
210    pub execution_boundary: Option<ExecutionBoundary>,
211    #[serde(default)]
212    pub provenance: PermissionProvenanceSummary,
213    pub target: PermissionAuditTarget,
214    pub group_ids: Vec<PermissionGroupId>,
215    pub policy: PermissionPolicyReference,
216    pub escalation_path: Vec<PermissionEscalationAuditHop>,
217    pub decision_id: Option<String>,
218    pub actor: Option<PermissionProjectionActor>,
219    pub scope: Option<PermissionAuditScope>,
220    pub reason: Option<String>,
221    pub at: DateTime<Utc>,
222}
223
224impl PermissionRequestAudit {
225    pub fn from_request(
226        request: &PermissionRequest,
227        group_ids: Vec<PermissionGroupId>,
228        decision: Option<&PermissionDecision>,
229        at: DateTime<Utc>,
230    ) -> Self {
231        let provenance = &request.intent.provenance;
232        let target = match &request.state {
233            PermissionRequestState::Pending { target } => target,
234            _ => request
235                .escalation_path
236                .last()
237                .map(|hop| &hop.target)
238                .expect("request has target"),
239        };
240        Self {
241            request_id: Some(request.request_id.clone()),
242            revision: request.revision,
243            session_id: request.session_id.clone(),
244            requesting_run_id: request.requesting_run_id.clone(),
245            parent_run_id: request.parent_run_id.clone(),
246            root_run_id: request.root_run_id.clone(),
247            tool_use_id: request.intent.tool_use_id.clone(),
248            tool: request.intent.tool_name.clone(),
249            call_intent: request.intent.call_intent.clone(),
250            tier: request.intent.tier,
251            execution_boundary: request
252                .intent
253                .risks
254                .contains(&RiskKind::ProcessSpawn)
255                .then(|| match decision {
256                    Some(decision)
257                        if decision.action == crate::permission::PermissionAction::Approve =>
258                    {
259                        Some(decision.execution_boundary)
260                    }
261                    Some(_) => None,
262                    None if matches!(request.state, PermissionRequestState::Pending { .. }) => {
263                        Some(if matches!(target, ApprovalTarget::User) {
264                            ExecutionBoundary::Direct
265                        } else {
266                            ExecutionBoundary::Sandboxed
267                        })
268                    }
269                    None => None,
270                })
271                .flatten(),
272            provenance: PermissionProvenanceSummary {
273                cwd: provenance
274                    .cwd
275                    .as_ref()
276                    .map(|path| path.display().to_string()),
277                path: provenance
278                    .path
279                    .as_ref()
280                    .map(|path| path.display().to_string()),
281                path_origin: provenance.path_origin.map(|origin| format!("{origin:?}")),
282                workspace_id: provenance.workspace_id.clone(),
283                workspace_root: provenance
284                    .workspace_root
285                    .as_ref()
286                    .map(|path| path.display().to_string()),
287                repository_root: provenance
288                    .repository_root
289                    .as_ref()
290                    .map(|path| path.display().to_string()),
291                network: provenance.network,
292                risks: provenance
293                    .risks
294                    .iter()
295                    .map(|risk| format!("{risk:?}"))
296                    .collect(),
297                targets: provenance
298                    .authorized_targets()
299                    .map(|path| path.display().to_string())
300                    .collect(),
301            },
302            target: target.into(),
303            group_ids,
304            policy: request.policy_reference.clone(),
305            escalation_path: request.escalation_path.iter().map(Into::into).collect(),
306            decision_id: decision.map(|decision| decision.decision_id.to_string()),
307            actor: decision
308                .map(|decision| (&decision.actor).into())
309                .or_else(|| {
310                    request
311                        .escalation_path
312                        .iter()
313                        .rev()
314                        .find_map(|hop| hop.actor.as_ref().map(Into::into))
315                }),
316            scope: decision.and_then(|decision| decision.grant_scope.as_ref().map(Into::into)),
317            reason: decision
318                .and_then(|decision| decision.reason.clone())
319                .or_else(|| {
320                    request
321                        .escalation_path
322                        .iter()
323                        .rev()
324                        .find_map(|hop| hop.reason.clone())
325                })
326                .or_else(|| match &request.state {
327                    PermissionRequestState::Cancelled { reason } => Some(reason.clone()),
328                    _ => None,
329                }),
330            at,
331        }
332    }
333}
334
335#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
336#[serde(tag = "kind", rename_all = "snake_case")]
337pub enum PermissionGroupAuditOwner {
338    Flow { run_id: FlowRunId },
339    User { session_id: String },
340    System,
341}
342
343#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
344pub struct PermissionGroupAudit {
345    pub group_id: PermissionGroupId,
346    pub owner: PermissionGroupAuditOwner,
347    pub label: String,
348    pub request_ids: Vec<PermissionRequestId>,
349    pub revision: u64,
350    pub at: DateTime<Utc>,
351}
352
353impl PermissionGroupAudit {
354    pub fn from_group(group: &PermissionGroup, session_id: &str, at: DateTime<Utc>) -> Self {
355        let owner = match &group.owner {
356            GroupOwner::Flow(run_id) => PermissionGroupAuditOwner::Flow {
357                run_id: run_id.clone(),
358            },
359            GroupOwner::User => PermissionGroupAuditOwner::User {
360                session_id: session_id.to_owned(),
361            },
362            GroupOwner::System => PermissionGroupAuditOwner::System,
363        };
364        Self {
365            group_id: group.group_id.clone(),
366            owner,
367            label: group.label.clone(),
368            request_ids: group.request_ids.iter().cloned().collect(),
369            revision: group.revision,
370            at,
371        }
372    }
373}
374
375#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
376pub struct PermissionGrantAudit {
377    pub grant_id: crate::permission::PermissionGrantId,
378    pub request_id: PermissionRequestId,
379    pub requesting_run_id: FlowRunId,
380    pub actor: PermissionProjectionActor,
381    #[serde(default, skip_serializing_if = "Option::is_none")]
382    pub execution_boundary: Option<ExecutionBoundary>,
383    pub scope: PermissionAuditScope,
384    pub reason: Option<String>,
385    pub at: DateTime<Utc>,
386}
387
388impl PermissionGrantAudit {
389    pub fn from_grant(grant: &PermissionGrant, reason: Option<String>, at: DateTime<Utc>) -> Self {
390        Self::from_grant_with_actor(grant, (&grant.actor).into(), reason, at)
391    }
392
393    pub fn from_grant_with_actor(
394        grant: &PermissionGrant,
395        actor: PermissionProjectionActor,
396        reason: Option<String>,
397        at: DateTime<Utc>,
398    ) -> Self {
399        Self {
400            grant_id: grant.grant_id.clone(),
401            request_id: grant.request_id.clone(),
402            requesting_run_id: grant.requesting_run_id.clone(),
403            actor,
404            execution_boundary: grant
405                .requirement
406                .risks
407                .contains(&RiskKind::ProcessSpawn)
408                .then_some(grant.execution_boundary),
409            scope: (&grant.scope).into(),
410            reason,
411            at,
412        }
413    }
414}
415
416#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
417#[serde(tag = "kind", content = "payload", rename_all = "snake_case")]
418pub enum PermissionAuditRecord {
419    RequestCreated(PermissionRequestAudit),
420    RequestTargeted(PermissionRequestAudit),
421    RequestDeferred(PermissionRequestAudit),
422    RequestApproved(PermissionRequestAudit),
423    RequestDenied(PermissionRequestAudit),
424    RequestCancelled(PermissionRequestAudit),
425    GroupCreated(PermissionGroupAudit),
426    GroupUpdated(PermissionGroupAudit),
427    GroupResolved(PermissionGroupAudit),
428    GrantCreated(PermissionGrantAudit),
429    GrantExpired(PermissionGrantAudit),
430    UnrestrictedExecution(PermissionRequestAudit),
431}
432
433#[derive(Clone)]
434pub struct PermissionAuditProjector {
435    sink: crate::event::EventSink,
436    stream: tokio::sync::broadcast::Sender<crate::stream::StreamFrame>,
437}
438
439impl PermissionAuditProjector {
440    pub fn new(
441        sink: crate::event::EventSink,
442        stream: tokio::sync::broadcast::Sender<crate::stream::StreamFrame>,
443    ) -> Self {
444        Self { sink, stream }
445    }
446
447    pub fn emit(&self, record: PermissionAuditRecord) {
448        self.sink.emit(record.clone().into());
449        let _ = self.stream.send(record.into());
450    }
451}
452
453trait AuditAnchor {
454    fn stream_anchor(&self) -> String;
455}
456
457impl AuditAnchor for PermissionRequestAudit {
458    fn stream_anchor(&self) -> String {
459        self.requesting_run_id.to_string()
460    }
461}
462
463impl AuditAnchor for PermissionGroupAudit {
464    fn stream_anchor(&self) -> String {
465        match &self.owner {
466            PermissionGroupAuditOwner::Flow { run_id } => run_id.to_string(),
467            PermissionGroupAuditOwner::User { session_id } => format!("user:{session_id}"),
468            PermissionGroupAuditOwner::System => "system".into(),
469        }
470    }
471}
472
473impl AuditAnchor for PermissionGrantAudit {
474    fn stream_anchor(&self) -> String {
475        self.requesting_run_id.to_string()
476    }
477}
478
479macro_rules! audit_conversions {
480    ($(($record:ident, $event:ident, $frame:ident)),+ $(,)?) => {
481        impl From<PermissionAuditRecord> for crate::event::Event {
482            fn from(record: PermissionAuditRecord) -> Self {
483                match record {
484                    $(PermissionAuditRecord::$record(payload) => Self::$event { payload },)+
485                }
486            }
487        }
488
489        impl From<PermissionAuditRecord> for crate::stream::StreamFrame {
490            fn from(record: PermissionAuditRecord) -> Self {
491                match record {
492                    $(PermissionAuditRecord::$record(payload) => Self::$frame {
493                        run_id: payload.stream_anchor(),
494                        payload,
495                    },)+
496                }
497            }
498        }
499    };
500}
501
502audit_conversions!(
503    (
504        RequestCreated,
505        PermissionRequestCreated,
506        PermissionRequestCreated
507    ),
508    (
509        RequestTargeted,
510        PermissionRequestTargeted,
511        PermissionRequestTargeted
512    ),
513    (
514        RequestDeferred,
515        PermissionRequestDeferred,
516        PermissionRequestDeferred
517    ),
518    (
519        RequestApproved,
520        PermissionRequestApproved,
521        PermissionRequestApproved
522    ),
523    (
524        RequestDenied,
525        PermissionRequestDenied,
526        PermissionRequestDenied
527    ),
528    (
529        RequestCancelled,
530        PermissionRequestCancelled,
531        PermissionRequestCancelled
532    ),
533    (GroupCreated, PermissionGroupCreated, PermissionGroupCreated),
534    (GroupUpdated, PermissionGroupUpdated, PermissionGroupUpdated),
535    (
536        GroupResolved,
537        PermissionGroupResolved,
538        PermissionGroupResolved
539    ),
540    (GrantCreated, PermissionGrantCreated, PermissionGrantCreated),
541    (GrantExpired, PermissionGrantExpired, PermissionGrantExpired),
542    (
543        UnrestrictedExecution,
544        UnrestrictedExecution,
545        UnrestrictedExecution
546    ),
547);
548
549#[cfg(test)]
550mod tests {
551    use super::*;
552    use crate::event::{EventEnvelope, FlowRunId};
553    use crate::permission::{PermissionGrantId, PermissionGroupId, PermissionRequestId};
554    use crate::stream::frame_run_id;
555
556    fn request(run_id: &FlowRunId) -> PermissionRequestAudit {
557        PermissionRequestAudit {
558            request_id: Some(PermissionRequestId::now()),
559            revision: 1,
560            session_id: "session".into(),
561            requesting_run_id: run_id.clone(),
562            parent_run_id: None,
563            root_run_id: run_id.clone(),
564            tool_use_id: "tool-use".into(),
565            tool: "fs.read".into(),
566            call_intent: None,
567            tier: Tier::Two,
568            execution_boundary: None,
569            provenance: PermissionProvenanceSummary::default(),
570            target: PermissionAuditTarget::User,
571            group_ids: Vec::new(),
572            policy: PermissionPolicyReference {
573                snapshot_id: "snapshot".into(),
574                rule_id: "rule".into(),
575            },
576            escalation_path: Vec::new(),
577            decision_id: None,
578            actor: None,
579            scope: None,
580            reason: None,
581            at: Utc::now(),
582        }
583    }
584
585    #[test]
586    fn legacy_request_audit_has_unknown_execution_boundary() {
587        let mut value = serde_json::to_value(request(&FlowRunId::now())).unwrap();
588        value.as_object_mut().unwrap().remove("execution_boundary");
589
590        let decoded: PermissionRequestAudit = serde_json::from_value(value).unwrap();
591
592        assert_eq!(decoded.execution_boundary, None);
593    }
594
595    #[test]
596    fn s8_actor_origin_matrix_remains_distinguishable_after_roundtrip() {
597        let run_id = FlowRunId::now();
598        let actors = [
599            PermissionProjectionActor::Policy {
600                policy_version: "snapshot".into(),
601                rule_id: "automatic:tier".into(),
602            },
603            PermissionProjectionActor::Flow {
604                session_id: "session".into(),
605                run_id: run_id.clone(),
606            },
607            PermissionProjectionActor::User {
608                session_id: "session".into(),
609                principal_id: Some("operator".into()),
610            },
611            PermissionProjectionActor::Policy {
612                policy_version: "snapshot".into(),
613                rule_id: "policy:deny".into(),
614            },
615            PermissionProjectionActor::System {
616                component: "permission.run_cleanup".into(),
617            },
618            PermissionProjectionActor::UnknownLegacy {
619                label: "legacy approver unavailable".into(),
620            },
621        ];
622        let roundtripped = actors
623            .iter()
624            .map(|actor| {
625                serde_json::from_str::<PermissionProjectionActor>(
626                    &serde_json::to_string(actor).unwrap(),
627                )
628                .unwrap()
629            })
630            .collect::<Vec<_>>();
631        assert_eq!(roundtripped, actors);
632        assert!(
633            matches!(roundtripped[0], PermissionProjectionActor::Policy { ref rule_id, .. } if rule_id == "automatic:tier")
634        );
635        assert!(matches!(
636            roundtripped[1],
637            PermissionProjectionActor::Flow { .. }
638        ));
639        assert!(matches!(
640            roundtripped[2],
641            PermissionProjectionActor::User { .. }
642        ));
643        assert!(
644            matches!(roundtripped[3], PermissionProjectionActor::Policy { ref rule_id, .. } if rule_id == "policy:deny")
645        );
646        assert!(matches!(
647            roundtripped[4],
648            PermissionProjectionActor::System { .. }
649        ));
650        assert!(matches!(
651            roundtripped[5],
652            PermissionProjectionActor::UnknownLegacy { .. }
653        ));
654    }
655
656    #[test]
657    fn s8_all_permission_audit_variants_roundtrip_and_keep_event_frame_mappings() {
658        let run_id = FlowRunId::now();
659        let request = request(&run_id);
660        let group = PermissionGroupAudit {
661            group_id: PermissionGroupId(uuid::Uuid::now_v7()),
662            owner: PermissionGroupAuditOwner::Flow {
663                run_id: run_id.clone(),
664            },
665            label: "group".into(),
666            request_ids: vec![request.request_id.clone().unwrap()],
667            revision: 1,
668            at: Utc::now(),
669        };
670        let grant = PermissionGrantAudit {
671            grant_id: PermissionGrantId(uuid::Uuid::now_v7()),
672            request_id: request.request_id.clone().unwrap(),
673            requesting_run_id: run_id.clone(),
674            actor: PermissionProjectionActor::System {
675                component: "test".into(),
676            },
677            execution_boundary: None,
678            scope: PermissionAuditScope::CurrentCall,
679            reason: None,
680            at: Utc::now(),
681        };
682        let cases = [
683            (
684                PermissionAuditRecord::RequestCreated(request.clone()),
685                "permission_request_created",
686            ),
687            (
688                PermissionAuditRecord::RequestTargeted(request.clone()),
689                "permission_request_targeted",
690            ),
691            (
692                PermissionAuditRecord::RequestDeferred(request.clone()),
693                "permission_request_deferred",
694            ),
695            (
696                PermissionAuditRecord::RequestApproved(request.clone()),
697                "permission_request_approved",
698            ),
699            (
700                PermissionAuditRecord::RequestDenied(request.clone()),
701                "permission_request_denied",
702            ),
703            (
704                PermissionAuditRecord::RequestCancelled(request.clone()),
705                "permission_request_cancelled",
706            ),
707            (
708                PermissionAuditRecord::GroupCreated(group.clone()),
709                "permission_group_created",
710            ),
711            (
712                PermissionAuditRecord::GroupUpdated(group.clone()),
713                "permission_group_updated",
714            ),
715            (
716                PermissionAuditRecord::GroupResolved(group),
717                "permission_group_resolved",
718            ),
719            (
720                PermissionAuditRecord::GrantCreated(grant.clone()),
721                "permission_grant_created",
722            ),
723            (
724                PermissionAuditRecord::GrantExpired(grant),
725                "permission_grant_expired",
726            ),
727            (
728                PermissionAuditRecord::UnrestrictedExecution(request),
729                "unrestricted_execution",
730            ),
731        ];
732
733        for (seq, (record, expected_kind)) in cases.into_iter().enumerate() {
734            let event: crate::event::Event = record.clone().into();
735            let frame: crate::stream::StreamFrame = record.into();
736            assert_eq!(crate::event_writer::event_kind(&event), expected_kind);
737            let run_id_text = run_id.to_string();
738            assert_eq!(
739                crate::event_writer::extract_anchors(&event).1.as_deref(),
740                Some(run_id_text.as_str())
741            );
742            assert_eq!(frame_run_id(&frame), Some(run_id_text.as_str()));
743            let json = serde_json::to_string(&EventEnvelope::new(seq as u64 + 1, event)).unwrap();
744            let roundtrip: EventEnvelope = serde_json::from_str(&json).unwrap();
745            assert_eq!(
746                crate::event_writer::event_kind(&roundtrip.event),
747                expected_kind
748            );
749        }
750    }
751}