(
&self,
reservation: JournalTurnReservation,
envelopes: Vec<ProvenanceJournalEnvelope>,
now: i64,
)
| 286 | } |
| 287 | |
| 288 | fn append( |
| 289 | &self, |
| 290 | reservation: JournalTurnReservation, |
| 291 | envelopes: Vec<ProvenanceJournalEnvelope>, |
| 292 | now: i64, |
| 293 | ) -> Result<Vec<JournalAppendAck>, String> { |
| 294 | let repository = self.repository.clone(); |
| 295 | let envelopes = envelopes |
| 296 | .into_iter() |
| 297 | .map(|envelope| { |
| 298 | Ok(WireEnvelope { |
| 299 | event_id: envelope.event_id.clone(), |
| 300 | bytes: envelope.to_json_bytes()?, |
| 301 | }) |
| 302 | }) |
| 303 | .collect::<Result<Vec<_>, atomic_agent::ProvenanceJournalError>>() |
| 304 | .map_err(|error| error.to_string())?; |
| 305 | run_outside_async_runtime(move || { |
| 306 | let mut acknowledgements = Vec::new(); |
| 307 | for chunk in chunk_wire_envelopes(envelopes) { |
| 308 | let response = request_with_reconnect( |
| 309 | &repository, |
| 310 | OwnerRequest::AppendProvenanceEnvelopes { |
| 311 | provenance_id: reservation.provenance_id, |
| 312 | expected_generation: reservation.generation, |
| 313 | envelopes: chunk, |
| 314 | now, |
| 315 | }, |
| 316 | )?; |
| 317 | match response { |
| 318 | OwnerResponse::ProvenanceEnvelopesCommitted { |
| 319 | acknowledgements: acks, |
| 320 | } => acknowledgements.extend(acks), |
| 321 | other => return Err(unexpected_response("append provenance", other)), |
| 322 | } |
| 323 | } |
| 324 | Ok(acknowledgements) |
| 325 | }) |
| 326 | .map_err(|error| error.to_string()) |
| 327 | } |
| 328 | |
| 329 | fn prepare_checkpoint( |
| 330 | &self, |
no test coverage detected