Skip to main content

agentplane/keyring/
push.rs

1//! Webhook registrations whose credentials are sealed at rest.
2//!
3//! A push registration is two things fused: a **destination** — task id, URL,
4//! delivery cursor — and the **credentials** for it, A2A's opaque correlation
5//! token and the receiver's own HTTP authentication. The concept spec is blunt
6//! about what the second half means in the wrong hands: a leaked registration
7//! is "a destination and a bearer token for it". The stores kept both halves
8//! as they arrived, so every backup of the push table was a list of endpoints
9//! anybody holding it could authenticate to.
10//!
11//! So the credentials are sealed and the destination is not, by the same
12//! dividing rule every other decorator here follows: *what a store is asked
13//! questions about stays readable; what it merely holds is sealed.* The due
14//! query orders on `(next_attempt_at, task, id)` and filters on the id prefix;
15//! the worker routes on the URL; nothing anywhere asks a question about a
16//! token. Sealing the routing fields would leave a delivery table that cannot
17//! deliver, and would buy nothing the sealed credentials do not already buy.
18//!
19//! # The erasure unit is the tenant
20//!
21//! Everything case-shaped in this crate seals under the case, because the case
22//! is what an erasure request names. A webhook credential is not case data: it
23//! belongs to the tenant's relationship with a receiver, it is reused across
24//! every task that receiver registers for, and no erasure request about a
25//! matter has ever meant "and stop being able to authenticate to my webhook".
26//! The one erasure that does reach it is the one that erases the tenant — a
27//! webhook token outlives no tenant — so the scope is `<tenant>/push`, one key
28//! for the unit that goes away together.
29//!
30//! # What an erased credential reads back as
31//!
32//! Absent. A registration whose credentials no longer open comes back with no
33//! token and no authentication rather than as an error, because the rows share
34//! a store with every other tenant's: a sweep that died on the first erased
35//! row would stop delivering for everyone still here. Delivery then goes out
36//! unauthenticated, the receiver refuses it, and the registration ages out
37//! through the ordinary retry ceiling — the failure lands on exactly the rows
38//! whose tenant was erased, and nowhere else.
39
40use std::sync::Arc;
41
42use async_trait::async_trait;
43
44use crate::core::{RunId, Secret, Seq, StoreError, TenantId};
45use crate::journal::payload;
46use crate::push::{DueBatch, PushConfig, PushNamespace, PushRegistration, PushStore};
47
48use super::KeyRing;
49
50/// A [`PushStore`] that seals webhook credentials under a key ring.
51#[derive(Debug)]
52pub struct SealedPush {
53    inner: Arc<dyn PushStore>,
54    keys: Arc<dyn KeyRing>,
55    tenant: TenantId,
56}
57
58impl SealedPush {
59    /// Seal this store's webhook credentials under `keys`.
60    ///
61    /// `tenant` must be the tenant the wrapped store serves — see
62    /// [`SealedCases::wrap`](super::SealedCases::wrap) for why this argument
63    /// exists and what a mismatch costs.
64    ///
65    /// # Panics
66    ///
67    /// If `tenant` is not the tenant `inner` serves — see
68    /// [`SealedCases::wrap`](super::SealedCases::wrap) for why that pair is
69    /// checked rather than trusted.
70    #[must_use]
71    pub fn wrap(inner: Arc<dyn PushStore>, keys: Arc<dyn KeyRing>, tenant: TenantId) -> Arc<Self> {
72        super::assert_serves(inner.tenant(), &tenant, "push");
73        Arc::new(Self {
74            inner,
75            keys,
76            tenant,
77        })
78    }
79
80    /// One scope for the whole tenant's credentials — the module docs say why
81    /// the tenant, and not the task or a case, is the unit that goes away
82    /// together.
83    fn scope(&self) -> String {
84        super::scope(&self.tenant, "push")
85    }
86
87    /// Bound to tenant, purpose label and registration, so a sealed credential
88    /// lifted onto another row fails to authenticate rather than opening as
89    /// authority to a destination it was never given for. The `push:{tenant}:`
90    /// prefix follows the journal and case decorators: the registration id is
91    /// caller-supplied, so a bare `"{task}/{id}"` was a string a caller could
92    /// shape to collide with another decorator's AAD. What the AAD does not do
93    /// alone: the sealing scope also separates envelopes, and the two are
94    /// deliberately redundant.
95    /// `pub(super)` so the keyring's own tests can hold the three decorators'
96    /// derivations side by side and prove colliding identifiers never share
97    /// an AAD.
98    pub(super) fn aad(tenant: &TenantId, task: RunId, id: &str) -> String {
99        format!("push:{tenant}:{task}/{id}")
100    }
101
102    async fn sealed_secret(&self, aad: &str, secret: &Secret) -> Result<Secret, StoreError> {
103        let envelope = super::envelope::seal(
104            self.keys.as_ref(),
105            &self.scope(),
106            aad.as_bytes(),
107            secret.expose().as_bytes(),
108        )
109        .await
110        .map_err(|e| StoreError::Backend(format!("sealing a webhook credential failed: {e}")))?;
111        Ok(Secret::new(payload::wrap_text(&envelope)))
112    }
113
114    /// The stored shape: credentials sealed, routing untouched.
115    async fn sealed(&self, config: &PushConfig) -> Result<PushConfig, StoreError> {
116        let aad = Self::aad(&self.tenant, config.task, &config.id);
117        let mut sealed = config.clone();
118        if let Some(token) = &config.token {
119            sealed.token = Some(self.sealed_secret(&aad, token).await?);
120        }
121        if let Some(authentication) = &mut sealed.authentication {
122            authentication.credentials = self
123                .sealed_secret(&aad, &authentication.credentials)
124                .await?;
125        }
126        Ok(sealed)
127    }
128
129    /// A stored credential back in the clear, or `None` once its key is
130    /// **destroyed**.
131    ///
132    /// A secret that never was sealed passes through unchanged — the rows a
133    /// deployment wrote before it configured a key ring must stay deliverable,
134    /// exactly as the other decorators leave pre-sealing payloads readable.
135    ///
136    /// # Errors
137    ///
138    /// Every key failure that is not an erasure. This is the sharpest place in
139    /// the crate where those differ: `None` drops the credential and delivers
140    /// the notification *without it*, which is the right answer for a
141    /// credential that has been erased and a silent downgrade of an
142    /// authenticated channel for a key service that is merely unreachable. A
143    /// caller that cannot get a credential it was given must not send the
144    /// request; it must come back.
145    async fn opened_secret(&self, aad: &str, secret: Secret) -> Result<Option<Secret>, StoreError> {
146        let Some(envelope) = payload::unwrap_text(secret.expose()) else {
147            return Ok(Some(secret));
148        };
149        let opened = super::envelope::open_or_erased(self.keys.as_ref(), aad.as_bytes(), &envelope)
150            .await
151            .map_err(|e| StoreError::Backend(e.to_string()))?;
152        Ok(opened
153            .and_then(|plain| String::from_utf8(plain).ok())
154            .map(Secret::new))
155    }
156
157    async fn opened(&self, mut config: PushConfig) -> Result<PushConfig, StoreError> {
158        let aad = Self::aad(&self.tenant, config.task, &config.id);
159        if let Some(token) = config.token.take() {
160            config.token = self.opened_secret(&aad, token).await?;
161        }
162        // The scheme stays; the credentials decide. An authentication whose
163        // credentials were erased is no authentication at all, not a header
164        // carrying ciphertext to a receiver.
165        if let Some(authentication) = config.authentication.take() {
166            config.authentication = self
167                .opened_secret(&aad, authentication.credentials)
168                .await?
169                .map(|credentials| crate::push::PushAuthentication {
170                    scheme: authentication.scheme,
171                    credentials,
172                });
173        }
174        Ok(config)
175    }
176
177    async fn opened_all(
178        &self,
179        rows: Vec<PushRegistration>,
180    ) -> Result<Vec<PushRegistration>, StoreError> {
181        let mut out = Vec::with_capacity(rows.len());
182        for mut registration in rows {
183            registration.config = self.opened(registration.config).await?;
184            out.push(registration);
185        }
186        Ok(out)
187    }
188}
189
190#[async_trait]
191impl PushStore for SealedPush {
192    fn tenant(&self) -> &str {
193        self.tenant.as_str()
194    }
195
196    async fn put(&self, config: &PushConfig, next_seq: Seq) -> Result<(), StoreError> {
197        self.inner.put(&self.sealed(config).await?, next_seq).await
198    }
199
200    async fn get(&self, task: RunId, id: &str) -> Result<Option<PushConfig>, StoreError> {
201        let found = self.inner.get(task, id).await?;
202        Ok(match found {
203            Some(config) => Some(self.opened(config).await?),
204            None => None,
205        })
206    }
207
208    async fn list(&self, task: RunId) -> Result<Vec<PushConfig>, StoreError> {
209        let configs = self.inner.list(task).await?;
210        let mut out = Vec::with_capacity(configs.len());
211        for config in configs {
212            out.push(self.opened(config).await?);
213        }
214        Ok(out)
215    }
216
217    async fn due(&self, at: u64, limit: usize) -> Result<Vec<PushRegistration>, StoreError> {
218        // The delivery worker's read: this is where the credentials come back,
219        // because the POST that needs them happens on the other side of it.
220        let rows = self.inner.due(at, limit).await?;
221        self.opened_all(rows).await
222    }
223
224    async fn due_in(
225        &self,
226        at: u64,
227        limit: usize,
228        namespace: PushNamespace,
229    ) -> Result<DueBatch, StoreError> {
230        // Forwarded rather than left to the trait default, which would page
231        // over `self.due` and pay the decode twice — and would lose whatever
232        // native filter the wrapped store implements.
233        let mut batch = self.inner.due_in(at, limit, namespace).await?;
234        batch.rows = self.opened_all(batch.rows).await?;
235        Ok(batch)
236    }
237
238    async fn advance(&self, task: RunId, id: &str, next_seq: Seq) -> Result<(), StoreError> {
239        self.inner.advance(task, id, next_seq).await
240    }
241
242    async fn retry(
243        &self,
244        task: RunId,
245        id: &str,
246        next_attempt_at: u64,
247        error: &str,
248    ) -> Result<(), StoreError> {
249        self.inner.retry(task, id, next_attempt_at, error).await
250    }
251
252    async fn park(&self, task: RunId, id: &str, error: &str) -> Result<(), StoreError> {
253        self.inner.park(task, id, error).await
254    }
255
256    async fn parked(&self, limit: usize) -> Result<Vec<PushRegistration>, StoreError> {
257        // The credentials come back here for the same reason `due` opens them:
258        // an operator reading a parked row is about to decide whether to
259        // re-arm it, and a sealed URL or scheme tells them nothing.
260        let rows = self.inner.parked(limit).await?;
261        self.opened_all(rows).await
262    }
263
264    async fn unpark(&self, task: RunId, id: &str, at: u64) -> Result<bool, StoreError> {
265        self.inner.unpark(task, id, at).await
266    }
267
268    async fn delete(&self, task: RunId, id: &str) -> Result<(), StoreError> {
269        self.inner.delete(task, id).await
270    }
271}