Skip to main content

made_client/
progress.rs

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}