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}