1use std::collections::HashSet;
2
3use futures::StreamExt;
4use made_proto::v1::{stream_ceremony_response, StreamCeremonyEndReason, StreamCeremonyRequest};
5
6use crate::{MadeClient, MadeClientError, ProgressBatch, ProgressCheckpoint};
7
8impl MadeClient {
9 pub async fn watch_once(
10 &self,
11 checkpoint: &ProgressCheckpoint,
12 max_events: u32,
13 wait_timeout_ms: Option<u32>,
14 ) -> Result<ProgressBatch, MadeClientError> {
15 let mut rpc = self.rpc();
16 let response = rpc
17 .stream_ceremony(Self::request(
18 &self.context(),
19 "/underpass.made.v1.MadeService/StreamCeremony",
20 StreamCeremonyRequest {
21 ceremony_id: checkpoint.ceremony_id().to_owned(),
22 after_sequence: checkpoint.after_sequence(),
23 max_events,
24 wait_timeout_ms,
25 },
26 ))
27 .await
28 .map_err(MadeClientError::from_status)?;
29 let mut stream = response.into_inner();
30 let mut expected_sequence = checkpoint.after_sequence().saturating_add(1);
31 let mut event_ids = HashSet::new();
32 let mut records = Vec::new();
33
34 while let Some(frame) = stream.next().await {
35 let frame = frame.map_err(MadeClientError::from_status)?;
36 match frame.frame {
37 Some(stream_ceremony_response::Frame::Record(record)) => {
38 if record.ceremony_id != checkpoint.ceremony_id() {
39 return Err(MadeClientError::ProtocolViolation(format!(
40 "stream returned ceremony {} for scope {}",
41 record.ceremony_id,
42 checkpoint.ceremony_id()
43 )));
44 }
45 if !event_ids.insert(record.event_id.clone()) {
46 continue;
47 }
48 if record.sequence != expected_sequence {
49 return Err(MadeClientError::ProtocolViolation(format!(
50 "expected event sequence {expected_sequence}, got {}",
51 record.sequence
52 )));
53 }
54 expected_sequence = expected_sequence.saturating_add(1);
55 records.push(record);
56 }
57 Some(stream_ceremony_response::Frame::End(end)) => {
58 let last_seen = expected_sequence.saturating_sub(1);
59 if end.resume_after_sequence != last_seen {
60 return Err(MadeClientError::ProtocolViolation(format!(
61 "end cursor {} does not match last sequence {last_seen}",
62 end.resume_after_sequence
63 )));
64 }
65 let reason = StreamCeremonyEndReason::try_from(end.reason).map_err(|_| {
66 MadeClientError::ProtocolViolation(format!(
67 "unknown stream end reason {}",
68 end.reason
69 ))
70 })?;
71 return Ok(ProgressBatch::new(
72 records,
73 ProgressCheckpoint::new(
74 checkpoint.ceremony_id(),
75 end.resume_after_sequence,
76 ),
77 end.head_sequence,
78 reason,
79 ));
80 }
81 None => {
82 return Err(MadeClientError::ProtocolViolation(
83 "stream returned an empty frame".to_owned(),
84 ));
85 }
86 }
87 }
88
89 Err(MadeClientError::ProtocolViolation(
90 "stream ended without a resume cursor".to_owned(),
91 ))
92 }
93}