Skip to main content

silicon_apps_client/
events.rs

1//! Event feeds, server-sent event streams, subscriptions and webhook
2//! signature checks.
3
4use anyhow::{Context, Result, bail, ensure};
5use hmac::{Hmac, Mac};
6use serde::{Deserialize, Serialize};
7use serde_json::{Value, json};
8
9/// Which event feed to read.
10#[derive(Debug, Clone, PartialEq, Eq)]
11pub enum Feed {
12    /// The signed-in account: invitations to it, apps it authors and public
13    /// events of apps it installed.
14    Account,
15    /// One app's events, for its authors.
16    App(String),
17    /// One of your subscriptions, through its filters.
18    Subscription(String),
19}
20impl Feed {
21    pub(crate) fn path(&self, stream: bool) -> Vec<String> {
22        let mut path = match self {
23            Feed::App(app) => vec!["v1".into(), "apps".into(), app.clone(), "events".into()],
24            _ => vec!["v1".into(), "events".into()],
25        };
26        if stream {
27            path.push("stream".into());
28        }
29        path
30    }
31    pub(crate) fn query(&self) -> Vec<(&'static str, String)> {
32        match self {
33            Feed::Subscription(id) => vec![("subscription", id.clone())],
34            _ => vec![],
35        }
36    }
37}
38
39/// Where a subscription's events go.
40#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
41#[serde(tag = "mode", rename_all = "lowercase")]
42pub enum Delivery {
43    /// POST each event to this URL, signed with the subscription's secret.
44    Webhook { url: String },
45    /// Read events at `/v1/events/stream?subscription=ID`.
46    Stream,
47}
48
49/// A new subscription. `types` and `channels` may be empty for the defaults.
50#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
51pub struct NewSubscription {
52    #[serde(skip_serializing_if = "Option::is_none")]
53    pub app_id: Option<String>,
54    #[serde(skip_serializing_if = "Vec::is_empty")]
55    pub types: Vec<String>,
56    #[serde(skip_serializing_if = "Vec::is_empty")]
57    pub channels: Vec<String>,
58    pub delivery: Delivery,
59    #[serde(skip_serializing_if = "Option::is_none")]
60    pub description: Option<String>,
61}
62impl NewSubscription {
63    /// Follow one app's production releases, delivered to a webhook.
64    pub fn production_releases(app_id: &str, url: &str) -> Self {
65        Self {
66            app_id: Some(app_id.into()),
67            types: vec!["release.promoted".into()],
68            channels: vec![],
69            delivery: Delivery::Webhook { url: url.into() },
70            description: None,
71        }
72    }
73    pub fn body(&self) -> Value {
74        serde_json::to_value(self).unwrap_or_else(|_| json!({}))
75    }
76}
77
78/// One event from a stream.
79#[derive(Debug, Clone, PartialEq)]
80pub struct StreamEvent {
81    /// The event's seq. Pass it as `last_event_id` to resume after it.
82    pub id: Option<String>,
83    /// The event type, such as `release.promoted`.
84    pub event: String,
85    /// The event as JSON: seq, id, type, app_id, actor_uuid, occurred_at, data.
86    pub data: Value,
87}
88
89/// A server-sent event stream. Comments (the ready line and heartbeats) are
90/// skipped. `next` returns `None` when the server closes the stream; open it
91/// again with the last id to continue.
92pub struct EventStream {
93    response: reqwest::Response,
94    buffer: String,
95    last_event_id: Option<String>,
96}
97impl EventStream {
98    pub(crate) fn new(response: reqwest::Response) -> Self {
99        Self {
100            response,
101            buffer: String::new(),
102            last_event_id: None,
103        }
104    }
105    /// The id of the last event returned, for resuming.
106    pub fn last_event_id(&self) -> Option<&str> {
107        self.last_event_id.as_deref()
108    }
109    pub async fn next(&mut self) -> Result<Option<StreamEvent>> {
110        loop {
111            while let Some(end) = self.buffer.find("\n\n") {
112                let block: String = self.buffer.drain(..end + 2).collect();
113                if let Some(event) = parse_block(&block)? {
114                    if event.id.is_some() {
115                        self.last_event_id = event.id.clone();
116                    }
117                    return Ok(Some(event));
118                }
119            }
120            match self
121                .response
122                .chunk()
123                .await
124                .context("the event stream was interrupted")?
125            {
126                Some(chunk) => {
127                    self.buffer
128                        .push_str(&String::from_utf8_lossy(&chunk).replace("\r\n", "\n"));
129                    ensure!(
130                        self.buffer.len() <= 4 * 1024 * 1024,
131                        "an event stream message exceeded 4 MiB"
132                    );
133                }
134                None => return Ok(None),
135            }
136        }
137    }
138}
139fn parse_block(block: &str) -> Result<Option<StreamEvent>> {
140    let mut id = None;
141    let mut event = None;
142    let mut data = Vec::new();
143    for line in block.lines() {
144        if line.is_empty() || line.starts_with(':') {
145            continue;
146        }
147        let (field, value) = line.split_once(':').unwrap_or((line, ""));
148        let value = value.strip_prefix(' ').unwrap_or(value);
149        match field {
150            "id" => id = Some(value.to_owned()),
151            "event" => event = Some(value.to_owned()),
152            "data" => data.push(value),
153            _ => {}
154        }
155    }
156    if data.is_empty() {
157        return Ok(None);
158    }
159    let data: Value =
160        serde_json::from_str(&data.join("\n")).context("an event's data was not JSON")?;
161    Ok(Some(StreamEvent {
162        id,
163        event: event.unwrap_or_else(|| "message".into()),
164        data,
165    }))
166}
167
168/// Check a subscription webhook delivery before trusting it.
169///
170/// `body` must be the exact bytes received. `signature` is the
171/// `X-Apps-Signature` header, a comma-separated list of `v1=` values; any
172/// matching value is accepted. Timestamps further than `tolerance_seconds`
173/// from `now` (unix seconds) are refused, so a captured delivery cannot be
174/// replayed later. Five minutes is the recommended tolerance.
175pub fn verify_webhook(
176    secret: &str,
177    timestamp: &str,
178    signature: &str,
179    body: &[u8],
180    now: i64,
181    tolerance_seconds: i64,
182) -> Result<()> {
183    let sent: i64 = timestamp
184        .trim()
185        .parse()
186        .context("X-Apps-Timestamp is not a unix timestamp")?;
187    ensure!(
188        (now - sent).abs() <= tolerance_seconds,
189        "X-Apps-Timestamp is {} seconds from this clock, more than the {tolerance_seconds} second tolerance",
190        (now - sent).abs()
191    );
192    let expected = webhook_signature(secret, sent, body);
193    if signature
194        .split(',')
195        .map(str::trim)
196        .filter_map(|v| v.strip_prefix("v1="))
197        .any(|candidate| constant_time_eq(candidate.as_bytes(), &expected.as_bytes()[3..]))
198    {
199        return Ok(());
200    }
201    bail!("X-Apps-Signature does not match this body and secret")
202}
203/// The `v1=` signature Apps sends for a body at a timestamp.
204pub fn webhook_signature(secret: &str, timestamp: i64, body: &[u8]) -> String {
205    let mut mac =
206        Hmac::<sha2::Sha256>::new_from_slice(secret.as_bytes()).expect("HMAC takes any key");
207    mac.update(format!("{timestamp}.").as_bytes());
208    mac.update(body);
209    format!("v1={}", hex::encode(mac.finalize().into_bytes()))
210}
211fn constant_time_eq(a: &[u8], b: &[u8]) -> bool {
212    a.len() == b.len() && a.iter().zip(b).fold(0u8, |acc, (x, y)| acc | (x ^ y)) == 0
213}