agentplane 0.14.0

Durable, replayable agent runtime — the journal is the plan of record
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
//! Push notifications: telling a peer's webhook that its task moved.
//!
//! Streaming holds a connection open; push does not. A peer that asked for a
//! task at 09:00 and disconnected wants to know at 14:00 that it finished, and
//! neither polling forever nor holding a socket for five hours is a reasonable
//! way to arrange that.
//!
//! # The URL comes from the caller, which is the whole problem
//!
//! Every other outbound destination in this crate is granted by an operator. A
//! webhook URL is supplied by whoever created the task — so push is the one
//! feature where an untrusted party names an address this plane will connect to,
//! with a payload describing somebody's task.
//!
//! Three controls, and none of them is sufficient alone:
//!
//! * **An operator grant.** The host must be on an allowlist. A caller may pick
//!   any URL under a host the deployment permits, and no host it does not. This
//!   is the primary control; the rest are the second lock.
//! * **Public addresses only.** Every DNS answer is checked with
//!   [`netguard`](crate::netguard) and the connection is pinned to those
//!   addresses, so a name that resolves inward — or answers differently the
//!   second time — reaches nothing.
//! * **HTTPS only.** The payload describes a task; sending it in clear to an
//!   address chosen by the recipient is a disclosure with extra steps.
//!
//! # What is delivered, and what is not
//!
//! The payload is a `StreamResponse` — the same status/artifact union streaming
//! sends, so a receiver parses one thing and A2A's two delivery mechanisms do
//! not disagree. Registering a destination is therefore authorization to send
//! that task's output there: task-level policy and the operator host grant are
//! both checked rather than treating the allowlist as sufficient authority.
//!
//! # The journal is the outbox
//!
//! A registration stores the first journal sequence it has not acknowledged.
//! Workers derive `StreamResponse` payloads from those records and advance only
//! after HTTP 2xx. A crash after POST but before cursor persistence duplicates
//! an event instead of losing it, which is A2A's at-least-once contract.

use std::fmt::Debug;

use async_trait::async_trait;

use crate::core::{RunId, Secret, Seq, StoreError};

/// Where a peer wants to be told about a task.
#[derive(Debug, Clone)]
pub struct PushConfig {
    /// The configuration's own id, unique within a task.
    pub id: String,
    /// The task this is about.
    pub task: RunId,
    /// Where to POST. HTTPS, and on a granted host.
    pub url: String,
    /// A2A's opaque token for this task/session. It is not HTTP authentication;
    /// those credentials live in [`authentication`](Self::authentication).
    pub token: Option<Secret>,
    /// HTTP authentication for the receiver, distinct from A2A's opaque
    /// per-task/session token.
    pub authentication: Option<PushAuthentication>,
}

/// Authentication information from A2A's push configuration.
#[derive(Clone)]
pub struct PushAuthentication {
    pub scheme: String,
    pub credentials: Secret,
}

impl Debug for PushAuthentication {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        f.debug_struct("PushAuthentication")
            .field("scheme", &self.scheme)
            .field("credentials", &"<redacted>")
            .finish()
    }
}

impl PushAuthentication {
    /// Validate the HTTP authentication scheme and resulting header value.
    ///
    /// Kept on the protocol value rather than only on [`PushSender`], so a
    /// custom transport cannot accidentally make malformed A2A input valid.
    ///
    /// # Errors
    ///
    /// [`PushError::Malformed`] when the required scheme is not an RFC 9110
    /// token or the credentials cannot be represented in a header.
    pub fn validate(&self) -> Result<(), PushError> {
        if self.scheme.is_empty()
            || !self.scheme.bytes().all(|byte| {
                byte.is_ascii_alphanumeric()
                    || matches!(
                        byte,
                        b'!' | b'#'
                            | b'$'
                            | b'%'
                            | b'&'
                            | b'\''
                            | b'*'
                            | b'+'
                            | b'-'
                            | b'.'
                            | b'^'
                            | b'_'
                            | b'`'
                            | b'|'
                            | b'~'
                    )
            })
        {
            return Err(PushError::Malformed(
                "authentication.scheme must be an HTTP authentication token".to_owned(),
            ));
        }
        let value = format!("{} {}", self.scheme, self.credentials.expose());
        reqwest::header::HeaderValue::from_str(&value)
            .map(|_| ())
            .map_err(|error| PushError::Malformed(format!("invalid authentication: {error}")))
    }
}

/// One durable delivery cursor.
///
/// The journal is the outbox: `next_seq` names the first task record not yet
/// acknowledged by this receiver. A crash after POST and before `advance`
/// causes a duplicate, never a loss — the at-least-once direction A2A requires.
#[derive(Debug, Clone)]
pub struct PushRegistration {
    pub config: PushConfig,
    pub next_seq: Seq,
    pub attempts: u32,
    pub next_attempt_at: u64,
    pub last_error: Option<String>,
}

impl PushConfig {
    /// What a caller may see back.
    ///
    /// The token is **not** echoed. A caller that can read a config it did not
    /// create would otherwise learn another party's correlation secret, and the
    /// only party that needs the token already has it.
    #[must_use]
    pub fn redacted(&self) -> serde_json::Value {
        serde_json::json!({
            "id": self.id,
            "taskId": self.task.to_string(),
            "url": self.url,
            "authentication": self.authentication.as_ref().map(|auth| serde_json::json!({
                "scheme": auth.scheme,
            })),
        })
    }
}

/// Why a webhook was not accepted or not delivered.
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum PushError {
    #[error(
        "a webhook URL must be https — the payload describes a task, and \
         sending it in clear to an address the recipient chose is a disclosure"
    )]
    NotHttps,
    #[error("this deployment does not permit webhooks to '{0}'")]
    HostNotGranted(String),
    #[error("the webhook URL is not a URL: {0}")]
    Malformed(String),
    #[error("'{0}'")]
    Unroutable(String),
}

impl PushError {
    /// Whether waiting could change this answer.
    ///
    /// The grant is re-checked at delivery because a registration outlives the
    /// configuration that permitted it — but noticing a refusal and then
    /// scheduling that same refusal again is a decision that will never change,
    /// retried forever. A scheme that is not HTTPS, a URL that does not parse,
    /// and a host the operator has taken off the allowlist are all answers no
    /// backoff improves, so the worker abandons them rather than queueing a
    /// forty-first attempt against an answer it already has.
    ///
    /// [`Unroutable`](Self::Unroutable) is deliberately **not** permanent: it
    /// covers DNS, and DNS changes. A name that resolves inward today may be
    /// repointed tomorrow, and abandoning on the first answer would make a
    /// transient misconfiguration indistinguishable from a revoked grant.
    #[must_use]
    pub const fn is_permanent(&self) -> bool {
        matches!(
            self,
            Self::NotHttps | Self::HostNotGranted(_) | Self::Malformed(_)
        )
    }
}

/// Durable storage for webhook registrations.
///
/// Durable because a task outlives the connection that created it, and a
/// registration that vanished on restart would leave a peer waiting for a
/// notification nobody remembers promising.
#[async_trait]
pub trait PushStore: Send + Sync + Debug {
    /// Register or replace a configuration.
    ///
    /// Replacement preserves the existing acknowledgement cursor. Changing a
    /// URL or credentials is not permission to discard updates that receiver
    /// has not accepted.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn put(&self, config: &PushConfig, next_seq: Seq) -> Result<(), StoreError>;

    /// One configuration.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn get(&self, task: RunId, id: &str) -> Result<Option<PushConfig>, StoreError>;

    /// Every configuration for a task.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn list(&self, task: RunId) -> Result<Vec<PushConfig>, StoreError>;

    /// Registrations whose retry instant has arrived, in stable order.
    async fn due(&self, at: u64, limit: usize) -> Result<Vec<PushRegistration>, StoreError>;

    /// Acknowledge every record before `next_seq`.
    async fn advance(&self, task: RunId, id: &str, next_seq: Seq) -> Result<(), StoreError>;

    /// Record a failed attempt without advancing the cursor.
    async fn retry(
        &self,
        task: RunId,
        id: &str,
        next_attempt_at: u64,
        error: &str,
    ) -> Result<(), StoreError>;

    /// Forget one. Idempotent.
    ///
    /// # Errors
    ///
    /// If the store cannot be reached.
    async fn delete(&self, task: RunId, id: &str) -> Result<(), StoreError>;
}

/// Which webhook destinations this deployment permits.
///
/// Deny-by-default with no `allow_all`, for the reason
/// [`Egress`](crate::core::Egress) gives: a deployment that wants no control
/// should configure nothing, not configure something that looks like a control
/// and is not.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct PushPolicy {
    hosts: std::collections::BTreeSet<String>,
}

impl PushPolicy {
    #[must_use]
    pub fn new() -> Self {
        Self::default()
    }

    /// Permit webhooks to this exact host.
    ///
    /// # Panics
    ///
    /// If the host is not one a URL could name. The grant is matched against
    /// `Url::host_str`, which the URL crate IDNA-encodes to punycode — so a
    /// grant only lowercased would store an internationalised host in a form
    /// no webhook URL ever presents, silently refusing every delivery to it.
    /// Canonicalised through the same helper governed media uses, so the two
    /// host-granting surfaces cannot drift the way they had.
    #[must_use]
    pub fn allow_host(mut self, host: impl AsRef<str>) -> Self {
        let raw = host.as_ref();
        let host = crate::netguard::canonical_host(raw).unwrap_or_else(|| {
            panic!(
                "push host grant '{raw}' is not a host a URL can name — give an \
                 internationalised host in the form the URL parser accepts, or it \
                 would silently never match a webhook"
            )
        });
        self.hosts.insert(host);
        self
    }

    /// Check a caller-supplied URL before anything is stored.
    ///
    /// Checked at **registration**, not only at delivery, so a caller learns
    /// immediately that its webhook will never be called — rather than waiting
    /// for a notification that silently never comes.
    ///
    /// # Errors
    ///
    /// [`PushError::NotHttps`], [`PushError::Malformed`], or
    /// [`PushError::HostNotGranted`].
    pub fn check(&self, url: &str) -> Result<(), PushError> {
        self.check_allowing_loopback(url, false)
    }

    /// The same check, with the two address-shape refusals optionally lifted.
    ///
    /// `allow_loopback` is reachable only through
    /// [`PushSender::allow_plaintext_loopback`], which exists only under
    /// `testkit`. The **host grant is not lifted** — that is the primary
    /// control and it still has to name the host.
    fn check_allowing_loopback(&self, url: &str, allow_loopback: bool) -> Result<(), PushError> {
        let parsed = reqwest::Url::parse(url).map_err(|e| PushError::Malformed(e.to_string()))?;
        let host = parsed
            .host_str()
            .ok_or_else(|| PushError::Malformed("no host".to_owned()))?
            .trim_end_matches('.')
            .to_ascii_lowercase();
        if parsed.scheme() != "https" && !(allow_loopback && is_loopback_name(&host)) {
            return Err(PushError::NotHttps);
        }
        if !self.hosts.contains(&host) {
            return Err(PushError::HostNotGranted(host));
        }
        Ok(())
    }
}

/// Whether a host names this machine, without resolving it.
///
/// Literals only, plus the one name every stack special-cases. Anything that
/// merely *resolves* to loopback is deliberately not covered here — that is
/// [`netguard`](crate::netguard)'s job at delivery, against the answers DNS
/// actually gave, and a name-based guess in front of it would be a second
/// implementation of one rule.
fn is_loopback_name(host: &str) -> bool {
    if host == "localhost" {
        return true;
    }
    let bare = host
        .strip_prefix('[')
        .and_then(|h| h.strip_suffix(']'))
        .unwrap_or(host);
    bare.parse::<std::net::IpAddr>()
        .is_ok_and(|ip| ip.is_loopback())
}

/// Delivers notifications, under the controls in the module docs.
#[derive(Debug, Clone)]
pub struct PushSender {
    policy: PushPolicy,
    timeout: std::time::Duration,
    /// Lift the HTTPS requirement and the public-address check for a webhook
    /// on this machine. `testkit` only, and absent from any other build.
    #[cfg(feature = "testkit")]
    plaintext_loopback: bool,
}

/// Delivery transport used by the durable worker.
#[async_trait]
pub trait PushTransport: Send + Sync + Debug {
    fn validate(&self, config: &PushConfig) -> Result<(), PushError>;

    async fn deliver(
        &self,
        config: &PushConfig,
        payload: &serde_json::Value,
    ) -> Result<Delivered, PushError>;
}

impl PushSender {
    /// The spec recommends 10–30 seconds; a webhook that needs longer is doing
    /// work it should not be doing on our thread.
    pub const DEFAULT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15);

    #[must_use]
    pub fn new(policy: PushPolicy) -> Self {
        Self {
            policy,
            timeout: Self::DEFAULT_TIMEOUT,
            #[cfg(feature = "testkit")]
            plaintext_loopback: false,
        }
    }

    /// Permit `http://` to a webhook on this machine. **`testkit` only.**
    ///
    /// The A2A conformance kit's webhook receiver is an `http://localhost:PORT`
    /// server, because a kit cannot mint a public TLS endpoint for a run on a
    /// laptop. Both of this crate's address controls refuse that, correctly —
    /// and the consequence was that the kit's **ten push MUSTs could not run at
    /// all**, so the one surface where an untrusted party names an address this
    /// plane connects to had no outside-authority evidence behind it. Ten
    /// unrunnable rows is a worse answer than one named exception.
    ///
    /// What this does **not** lift is the part that is the actual control: the
    /// operator's **host grant** still has to name the host, the task-level
    /// authorization still runs, the cursor still advances only on 2xx, and
    /// every non-loopback destination is judged exactly as before — a plaintext
    /// URL to a public host stays refused with the flag set, which is the half
    /// that keeps this from being an off switch.
    ///
    /// It cannot exist in a production build: the field is `cfg(testkit)`, and
    /// `testkit` is documented as never belonging in one.
    #[cfg(feature = "testkit")]
    #[must_use]
    pub const fn allow_plaintext_loopback(mut self) -> Self {
        self.plaintext_loopback = true;
        self
    }

    /// Whether the loopback exception is in force. Always false without
    /// `testkit`, which is what lets the delivery path read one flag.
    const fn loopback_allowed(&self) -> bool {
        #[cfg(feature = "testkit")]
        {
            self.plaintext_loopback
        }
        #[cfg(not(feature = "testkit"))]
        {
            false
        }
    }

    #[must_use]
    pub const fn timeout(mut self, d: std::time::Duration) -> Self {
        self.timeout = d;
        self
    }

    /// The grant this sender enforces, so a registration can be checked against
    /// the same policy that will later be checked at delivery.
    #[must_use]
    pub const fn policy(&self) -> &PushPolicy {
        &self.policy
    }

    /// POST one `StreamResponse` to a registered webhook.
    ///
    /// The grant is re-checked here and not only at registration. A registration
    /// outlives the configuration that permitted it: a host removed from the
    /// allowlist must stop receiving notifications for tasks registered while it
    /// was still granted, and a check performed only at write time cannot do
    /// that.
    ///
    /// # Errors
    ///
    /// [`PushError`] when the URL is not permitted or does not resolve to a
    /// public address. Transport failures are **not** errors here — see the
    /// return type.
    pub async fn deliver(
        &self,
        config: &PushConfig,
        payload: &serde_json::Value,
    ) -> Result<Delivered, PushError> {
        <Self as PushTransport>::validate(self, config)?;

        let url =
            reqwest::Url::parse(&config.url).map_err(|e| PushError::Malformed(e.to_string()))?;
        let host = url
            .host_str()
            .ok_or_else(|| PushError::Malformed("no host".to_owned()))?
            .to_owned();
        let port = url.port_or_known_default().unwrap_or(443);

        // Resolved once, every answer checked, and the connection pinned to
        // exactly those addresses. Without the pin the client resolves again and
        // may be handed a different answer than the one that passed — which is
        // the rebinding attack this check would otherwise only appear to stop.
        let resolved = tokio::net::lookup_host((host.as_str(), port))
            .await
            .map_err(|e| PushError::Unroutable(format!("DNS for '{host}': {e}")))?;
        let addrs = if self.loopback_allowed() && is_loopback_name(&host) {
            // Named rather than inferred: the exception applies to a host that
            // *is* a loopback literal or `localhost`, not to one that merely
            // resolved to one. A name that resolves inward is the rebinding
            // attack, and it stays refused with the flag set.
            let addrs: Vec<_> = resolved.collect();
            if addrs.is_empty() {
                return Err(PushError::Unroutable(format!(
                    "DNS for '{host}' returned no addresses"
                )));
            }
            addrs
        } else {
            crate::netguard::all_public(&host, resolved)
                .map_err(|e| PushError::Unroutable(e.to_string()))?
        };

        let mut client = reqwest::Client::builder()
            .timeout(self.timeout)
            // No ambient authority: a webhook is somebody else's endpoint, and a
            // proxy, a cookie jar or a stored credential would attach this
            // plane's identity to a request it did not authorize.
            .no_proxy()
            .redirect(reqwest::redirect::Policy::none());
        for addr in &addrs {
            client = client.resolve(&host, *addr);
        }
        let client = client
            .build()
            .map_err(|e| PushError::Unroutable(e.to_string()))?;

        let mut request = client
            .post(url)
            // The spec's media type, so a receiver can route on it.
            .header("Content-Type", "application/a2a+json")
            .json(payload);
        if let Some(authentication) = &config.authentication {
            let value = format!(
                "{} {}",
                authentication.scheme,
                authentication.credentials.expose()
            );
            let value = reqwest::header::HeaderValue::from_str(&value).map_err(|error| {
                PushError::Malformed(format!("invalid authentication: {error}"))
            })?;
            request = request.header(reqwest::header::AUTHORIZATION, value);
        }

        // A transport failure is an *outcome*, not an error of this function. A
        // webhook that is down is ordinary, and the caller decides whether to
        // retry or to give up on a configuration — a distinction lost if this
        // returned `Err` for both "you may not" and "it did not answer".
        Ok(match request.send().await {
            Ok(response) if response.status().is_success() => Delivered::Accepted,
            Ok(response) => Delivered::Rejected(response.status().as_u16()),
            Err(e) => Delivered::Unreachable(e.to_string()),
        })
    }
}

#[async_trait]
impl PushTransport for PushSender {
    fn validate(&self, config: &PushConfig) -> Result<(), PushError> {
        self.policy
            .check_allowing_loopback(&config.url, self.loopback_allowed())?;
        if let Some(authentication) = &config.authentication {
            authentication.validate()?;
        }
        Ok(())
    }

    async fn deliver(
        &self,
        config: &PushConfig,
        payload: &serde_json::Value,
    ) -> Result<Delivered, PushError> {
        PushSender::deliver(self, config, payload).await
    }
}

/// What happened to one delivery attempt.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Delivered {
    /// The receiver answered 2xx, which is the acknowledgement the spec asks for.
    Accepted,
    /// It answered something else. Its problem, recorded rather than retried
    /// forever.
    Rejected(u16),
    /// It did not answer.
    Unreachable(String),
}