(
&self,
session_id: &str,
turn_number: u32,
drafts: Vec<JournalDraft>,
)
| 105 | } |
| 106 | |
| 107 | fn commit_journal_drafts( |
| 108 | &self, |
| 109 | session_id: &str, |
| 110 | turn_number: u32, |
| 111 | drafts: Vec<JournalDraft>, |
| 112 | ) -> AgentResult<()> { |
| 113 | let Some(sink) = &self.journal_sink else { |
| 114 | return Ok(()); |
| 115 | }; |
| 116 | if drafts.is_empty() { |
| 117 | return Ok(()); |
| 118 | } |
| 119 | let now = drafts |
| 120 | .iter() |
| 121 | .map(|draft| draft.timestamp_ms.div_euclid(1000)) |
| 122 | .max() |
| 123 | .unwrap_or(0); |
| 124 | let reservation = sink |
| 125 | .reserve_turn(session_id, turn_number, now) |
| 126 | .map_err(|reason| AgentError::ProvenanceJournalFailed { |
| 127 | session_id: session_id.to_string(), |
| 128 | reason, |
| 129 | })?; |
| 130 | let envelopes: Vec<_> = drafts |
| 131 | .into_iter() |
| 132 | .map(|draft| { |
| 133 | ProvenanceJournalEnvelope::new( |
| 134 | draft.event_id, |
| 135 | session_id, |
| 136 | turn_number, |
| 137 | reservation.generation, |
| 138 | draft.timestamp_ms, |
| 139 | draft.event, |
| 140 | ) |
| 141 | .with_causal_parents(draft.causal_parent_ids) |
| 142 | }) |
| 143 | .collect(); |
| 144 | let expected_ids: Vec<_> = envelopes |
| 145 | .iter() |
| 146 | .map(|envelope| envelope.event_id.clone()) |
| 147 | .collect(); |
| 148 | let acknowledgements = sink.append(reservation, envelopes, now).map_err(|reason| { |
| 149 | AgentError::ProvenanceJournalFailed { |
| 150 | session_id: session_id.to_string(), |
| 151 | reason, |
| 152 | } |
| 153 | })?; |
| 154 | let acknowledged_ids: std::collections::HashSet<_> = acknowledgements |
| 155 | .into_iter() |
| 156 | .map(|ack| ack.event_id) |
| 157 | .collect(); |
| 158 | if expected_ids |
| 159 | .iter() |
| 160 | .any(|event_id| !acknowledged_ids.contains(event_id)) |
| 161 | { |
| 162 | return Err(AgentError::ProvenanceJournalFailed { |
| 163 | session_id: session_id.to_string(), |
| 164 | reason: "owner omitted a committed event acknowledgement".to_string(), |
no test coverage detected