Skip to main content

macula_rust/
direct_dial.rs

1//! Direct-dial resolve-and-call: resolving a signed `procedure_advertisement`
2//! DHT record and its serving station's own signed `station_endpoint`, then
3//! dialing that station in one hop — instead of depending on ordinary
4//! advertise-gossip having propagated a route between whichever two
5//! stations happen to be involved.
6//!
7//! Ported from `macula-io/macula`'s `macula_direct_dial.erl`, cross-checked
8//! against `macula-go`'s own port of the same reference
9//! (`directdial/directdial.go`) — see that file's doc for the fuller
10//! reasoning behind each design choice made here.
11//!
12//! **Trust model** (see `macula_direct_dial.erl`'s module doc for the full
13//! reasoning): every candidate `procedure_advertisement` must carry a valid
14//! Ed25519 signature before its `serving_station` is trusted at all, and
15//! the resolved `station_endpoint` must be signed by the station itself.
16//! The actual QUIC dial trusts neither the TLS certificate (a production
17//! station's TLS is terminated by an unrelated PKI) nor nothing — trust is
18//! enforced at the application layer, by checking the freshly dialed
19//! session's own signature-verified HELLO identity against the exact
20//! pubkey the signed DHT chain resolved.
21//!
22//! `cert_chain`-based org/realm authorization (Slice 7c Direction B,
23//! `macula_record:verify_advertisement_cert_chain/3` on the Erlang side) is
24//! opt-in here too, matching the reference and `macula-go`'s own port —
25//! see [`resolve_with_cert_chain`]/[`call_with_cert_chain`]/
26//! [`advertise_direct_with_cert_chain`]. Plain [`resolve`]/[`call`]/
27//! [`advertise_direct`] are completely unaffected.
28
29use std::future::Future;
30use std::time::Duration;
31
32use crate::cbor::Value;
33use crate::cert_chain::{self, CertChainError};
34use crate::connection::{self, Session};
35use crate::content;
36use crate::dht::{self, DhtError, Record};
37use crate::frame::{CallResponse, StreamMode};
38use crate::identity::KeyPair;
39use crate::manifest::Mcid;
40use crate::stream::{self, StreamHandle};
41use crate::transport::Trust;
42
43fn now_ms() -> i128 {
44    use std::time::{SystemTime, UNIX_EPOCH};
45    SystemTime::now()
46        .duration_since(UNIX_EPOCH)
47        .expect("system clock before 1970")
48        .as_millis() as i128
49}
50
51/// Matches `macula_direct_dial.erl`'s `?RESOLVE_RETRIES`/`?RESOLVE_RETRY_MS`
52/// — a record just published on the provider's station has not necessarily
53/// replicated to the resolving station yet, so the first miss is not
54/// treated as failure.
55const RESOLVE_RETRIES: u32 = 50;
56const RESOLVE_RETRY_DELAY: Duration = Duration::from_millis(100);
57
58#[derive(Debug)]
59pub enum ResolveError {
60    /// Every `find_records` attempt came back empty after retrying past
61    /// DHT propagation lag.
62    ProcedureNotAdvertised,
63    /// Records were found, but none had a valid signature.
64    NoTrustedAdvertisement,
65    /// A resolved station published no reachable (or no longer valid)
66    /// `station_endpoint` after retrying.
67    StationEndpointNotFound,
68    /// A `station_endpoint` record was found under the right key, but its
69    /// signer didn't match the station it's supposed to describe.
70    StationEndpointSignerMismatch,
71    Dht(DhtError),
72    /// [`resolve_with_cert_chain`] only: at least one candidate
73    /// advertisement's envelope signature verified (otherwise
74    /// [`ResolveError::NoTrustedAdvertisement`] would apply instead), but
75    /// none passed cert-chain authorization for the expected org — carries
76    /// the specific [`CertChainError`] from the LAST candidate tried
77    /// (absent chain, wrong org, untrusted chain, etc.).
78    NoAuthorizedAdvertisement(CertChainError),
79}
80
81impl std::fmt::Display for ResolveError {
82    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
83        match self {
84            ResolveError::ProcedureNotAdvertised => {
85                write!(
86                    f,
87                    "direct_dial: procedure has no direct-dial advertisement in the DHT"
88                )
89            }
90            ResolveError::NoTrustedAdvertisement => write!(
91                f,
92                "direct_dial: every candidate advertisement failed signature verification"
93            ),
94            ResolveError::StationEndpointNotFound => write!(
95                f,
96                "direct_dial: resolved station published no reachable station_endpoint"
97            ),
98            ResolveError::StationEndpointSignerMismatch => {
99                write!(f, "direct_dial: station_endpoint signer mismatch")
100            }
101            ResolveError::Dht(e) => write!(f, "direct_dial: {e}"),
102            ResolveError::NoAuthorizedAdvertisement(e) => write!(
103                f,
104                "direct_dial: no candidate advertisement is cert-chain-authorized for the expected org: {e}"
105            ),
106        }
107    }
108}
109
110impl std::error::Error for ResolveError {}
111
112/// One resolved direct-dial target: the station's own node id plus a
113/// dialable host/port.
114#[derive(Debug, Clone)]
115pub struct Resolved {
116    pub station: [u8; 32],
117    pub host: String,
118    pub port: u16,
119}
120
121/// Finds `procedure`'s currently-advertised serving station and its
122/// dialable host/port, retrying past DHT propagation lag. `realm` and
123/// `procedure` must match exactly what the provider passed to
124/// [`advertise_direct`] (or the Erlang equivalent) — the discovery URI they
125/// derive must agree. `session` is used only to query the DHT; it does not
126/// need to be connected to the same station that will end up serving the
127/// call.
128pub async fn resolve(
129    session: &mut Session,
130    id: &KeyPair,
131    realm: [u8; 32],
132    procedure: &str,
133) -> Result<Resolved, ResolveError> {
134    let uri = dht::discovery_uri(realm, procedure);
135    let key = dht::procedure_key(&uri);
136
137    let mut recs: Vec<Record> = Vec::new();
138    for _ in 0..RESOLVE_RETRIES {
139        match dht::find_records(session, id, key).await {
140            Ok(found) if !found.is_empty() => {
141                recs = found;
142                break;
143            }
144            Ok(_) => {}
145            Err(e) => return Err(ResolveError::Dht(e)),
146        }
147        tokio::time::sleep(RESOLVE_RETRY_DELAY).await;
148    }
149    if recs.is_empty() {
150        return Err(ResolveError::ProcedureNotAdvertised);
151    }
152
153    let adv = first_trusted_advertisement(&recs).ok_or(ResolveError::NoTrustedAdvertisement)?;
154    resolve_station_endpoint(session, id, adv.serving_station).await
155}
156
157fn first_trusted_advertisement(recs: &[Record]) -> Option<dht::ProcedureAdvertisement> {
158    recs.iter().find_map(|rec| {
159        dht::verify(rec).ok()?;
160        dht::read_procedure_advertisement(rec).ok()
161    })
162}
163
164/// [`resolve`] plus Slice 7c Direction B managed-realm authorization: only
165/// an advertisement whose embedded cert chain validates to `realm_ca_pem`
166/// and names `expected_org` is trusted. Opt-in — [`resolve`] itself is
167/// unaffected and remains the right choice for unmanaged realms.
168pub async fn resolve_with_cert_chain(
169    session: &mut Session,
170    id: &KeyPair,
171    realm: [u8; 32],
172    procedure: &str,
173    realm_ca_pem: &[u8],
174    expected_org: &str,
175) -> Result<Resolved, ResolveError> {
176    let uri = dht::discovery_uri(realm, procedure);
177    let key = dht::procedure_key(&uri);
178
179    let mut recs: Vec<Record> = Vec::new();
180    for _ in 0..RESOLVE_RETRIES {
181        match dht::find_records(session, id, key).await {
182            Ok(found) if !found.is_empty() => {
183                recs = found;
184                break;
185            }
186            Ok(_) => {}
187            Err(e) => return Err(ResolveError::Dht(e)),
188        }
189        tokio::time::sleep(RESOLVE_RETRY_DELAY).await;
190    }
191    if recs.is_empty() {
192        return Err(ResolveError::ProcedureNotAdvertised);
193    }
194
195    let adv = first_authorized_advertisement(&recs, realm_ca_pem, expected_org)?;
196    resolve_station_endpoint(session, id, adv.serving_station).await
197}
198
199/// [`first_trusted_advertisement`] plus the cert-chain check. Matches Go's
200/// `firstAuthorizedAdvertisement`: if every candidate fails even the plain
201/// envelope-signature check, report [`ResolveError::NoTrustedAdvertisement`]
202/// (same as the plain path); only report
203/// [`ResolveError::NoAuthorizedAdvertisement`] once at least one candidate's
204/// signature verified but none passed cert-chain authorization.
205fn first_authorized_advertisement(
206    recs: &[Record],
207    realm_ca_pem: &[u8],
208    expected_org: &str,
209) -> Result<dht::ProcedureAdvertisement, ResolveError> {
210    let mut last_cert_err: Option<CertChainError> = None;
211    for rec in recs {
212        if dht::verify(rec).is_err() {
213            continue;
214        }
215        match cert_chain::verify_advertisement_cert_chain(realm_ca_pem, rec, expected_org) {
216            Ok(()) => {
217                if let Ok(adv) = dht::read_procedure_advertisement(rec) {
218                    return Ok(adv);
219                }
220            }
221            Err(e) => last_cert_err = Some(e),
222        }
223    }
224    match last_cert_err {
225        Some(e) => Err(ResolveError::NoAuthorizedAdvertisement(e)),
226        None => Err(ResolveError::NoTrustedAdvertisement),
227    }
228}
229
230/// Retries past a resolved-but-stale record, not just an absent one — the
231/// DHT can hand back a replica that hasn't been evicted yet even though the
232/// station's own current publish is live. Giving up on the first stale hit
233/// would make an otherwise healthy station unreachable via direct-dial
234/// until that one replica ages out.
235async fn resolve_station_endpoint(
236    session: &mut Session,
237    id: &KeyPair,
238    station: [u8; 32],
239) -> Result<Resolved, ResolveError> {
240    let key = dht::station_endpoint_key(station);
241    for _ in 0..RESOLVE_RETRIES {
242        let rec = match dht::find_record(session, id, key).await {
243            Ok(rec) => rec,
244            Err(DhtError::NotFound) => {
245                tokio::time::sleep(RESOLVE_RETRY_DELAY).await;
246                continue;
247            }
248            Err(e) => return Err(ResolveError::Dht(e)),
249        };
250        // The station_endpoint record for `station` must be SIGNED BY
251        // `station` itself — checking the signature and that the signer is
252        // exactly `station`, not just any valid signature, is what makes
253        // pinning the dial's expected identity meaningful below.
254        if rec.key != station {
255            return Err(ResolveError::StationEndpointSignerMismatch);
256        }
257        match dht::verify(&rec) {
258            Ok(()) => {}
259            Err(dht::VerifyError::Expired) => {
260                tokio::time::sleep(RESOLVE_RETRY_DELAY).await;
261                continue;
262            }
263            Err(_) => return Err(ResolveError::NoTrustedAdvertisement),
264        }
265        let ep =
266            dht::read_station_endpoint(&rec).map_err(|_| ResolveError::StationEndpointNotFound)?;
267        let Some(host) = ep.host_advertised.into_iter().next() else {
268            return Err(ResolveError::StationEndpointNotFound);
269        };
270        return Ok(Resolved {
271            station,
272            host,
273            port: ep.quic_port,
274        });
275    }
276    Err(ResolveError::StationEndpointNotFound)
277}
278
279#[derive(Debug)]
280pub enum CallError {
281    Resolve(ResolveError),
282    Dial(connection::HandshakeError),
283    /// The dialed peer's own signature-verified HELLO identity didn't
284    /// match the pubkey the signed DHT chain resolved — a trust violation,
285    /// not a retryable error.
286    TrustViolation {
287        resolved: [u8; 32],
288        dialed: [u8; 32],
289    },
290    Call(connection::CallError),
291}
292
293impl std::fmt::Display for CallError {
294    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
295        match self {
296            CallError::Resolve(e) => write!(f, "{e}"),
297            CallError::Dial(e) => write!(f, "direct_dial: dialing resolved station: {e}"),
298            CallError::TrustViolation { resolved, dialed } => write!(
299                f,
300                "direct_dial: trust violation -- resolved station {} but the dialed peer proved identity {}",
301                hex_of(resolved),
302                hex_of(dialed)
303            ),
304            CallError::Call(e) => write!(f, "direct_dial: {e}"),
305        }
306    }
307}
308
309impl std::error::Error for CallError {}
310
311fn hex_of(b: &[u8; 32]) -> String {
312    b.iter().map(|byte| format!("{byte:02x}")).collect()
313}
314
315/// Resolves `procedure`'s provider via direct-dial (through `resolve_via`,
316/// used only to query the DHT) and calls it there, in one hop, in a
317/// SEPARATE connection from `resolve_via`. The provider must have
318/// advertised via [`advertise_direct`] (or the Erlang
319/// `macula_response:advertise_direct/6,7`) — a plain `advertise` publishes
320/// no discoverable record and [`resolve`] will return
321/// [`ResolveError::ProcedureNotAdvertised`].
322///
323/// The dial itself uses [`Trust::Insecure`] (no TLS verification) because
324/// trust is enforced at the application layer instead — see the module
325/// doc's "Trust model". After the dial, the freshly connected session's own
326/// signature-verified HELLO identity is checked against the exact pubkey
327/// the signed DHT chain resolved; a mismatch is
328/// [`CallError::TrustViolation`], and the call is refused.
329pub async fn call(
330    resolve_via: &mut Session,
331    id: &KeyPair,
332    realm: [u8; 32],
333    procedure: &str,
334    payload: Value,
335    timeout: Duration,
336) -> Result<CallResponse, CallError> {
337    let resolved = resolve(resolve_via, id, realm, procedure)
338        .await
339        .map_err(CallError::Resolve)?;
340
341    let mut target = tokio::time::timeout(
342        timeout,
343        connection::connect(&resolved.host, resolved.port, Trust::Insecure, id),
344    )
345    .await
346    .unwrap_or(Err(connection::HandshakeError::Timeout))
347    .map_err(CallError::Dial)?;
348
349    if target.station.node_id != resolved.station {
350        let dialed = target.station.node_id;
351        target.close("trust_violation", None, id).await;
352        return Err(CallError::TrustViolation {
353            resolved: resolved.station,
354            dialed,
355        });
356    }
357
358    let deadline_ms = now_ms() + timeout.as_millis() as i128;
359    let result = target
360        .call(procedure, realm, payload, deadline_ms, id, timeout)
361        .await
362        .map_err(CallError::Call);
363    target.close("normal", None, id).await;
364    result
365}
366
367/// [`call`], presenting `ucan_token` to a provider gated with
368/// `{ucan_required, Issuer}`. Every hecate-om capability is advertised via
369/// [`advertise_direct`], so this is the only way a UCAN-gated capability
370/// is reachable through this crate at all -- [`call`] itself has no token
371/// parameter, and [`Session::call_with_ucan`] is the plain, non-direct
372/// path, which cannot resolve a direct-dial-only advertisement to begin
373/// with.
374pub async fn call_with_ucan(
375    resolve_via: &mut Session,
376    id: &KeyPair,
377    realm: [u8; 32],
378    procedure: &str,
379    payload: Value,
380    timeout: Duration,
381    ucan_token: Vec<u8>,
382) -> Result<CallResponse, CallError> {
383    let resolved = resolve(resolve_via, id, realm, procedure)
384        .await
385        .map_err(CallError::Resolve)?;
386
387    let mut target = tokio::time::timeout(
388        timeout,
389        connection::connect(&resolved.host, resolved.port, Trust::Insecure, id),
390    )
391    .await
392    .unwrap_or(Err(connection::HandshakeError::Timeout))
393    .map_err(CallError::Dial)?;
394
395    if target.station.node_id != resolved.station {
396        let dialed = target.station.node_id;
397        target.close("trust_violation", None, id).await;
398        return Err(CallError::TrustViolation {
399            resolved: resolved.station,
400            dialed,
401        });
402    }
403
404    let deadline_ms = now_ms() + timeout.as_millis() as i128;
405    let result = target
406        .call_with_ucan(
407            procedure,
408            realm,
409            payload,
410            deadline_ms,
411            id,
412            timeout,
413            ucan_token,
414        )
415        .await
416        .map_err(CallError::Call);
417    target.close("normal", None, id).await;
418    result
419}
420
421/// [`call`], resolved via [`resolve_with_cert_chain`] instead of
422/// [`resolve`] — see both for the full contract. Opt-in managed-realm
423/// authorization; [`call`] itself is unaffected.
424#[allow(clippy::too_many_arguments)]
425pub async fn call_with_cert_chain(
426    resolve_via: &mut Session,
427    id: &KeyPair,
428    realm: [u8; 32],
429    procedure: &str,
430    realm_ca_pem: &[u8],
431    expected_org: &str,
432    payload: Value,
433    timeout: Duration,
434) -> Result<CallResponse, CallError> {
435    let resolved = resolve_with_cert_chain(
436        resolve_via,
437        id,
438        realm,
439        procedure,
440        realm_ca_pem,
441        expected_org,
442    )
443    .await
444    .map_err(CallError::Resolve)?;
445
446    let mut target = tokio::time::timeout(
447        timeout,
448        connection::connect(&resolved.host, resolved.port, Trust::Insecure, id),
449    )
450    .await
451    .unwrap_or(Err(connection::HandshakeError::Timeout))
452    .map_err(CallError::Dial)?;
453
454    if target.station.node_id != resolved.station {
455        let dialed = target.station.node_id;
456        target.close("trust_violation", None, id).await;
457        return Err(CallError::TrustViolation {
458            resolved: resolved.station,
459            dialed,
460        });
461    }
462
463    let deadline_ms = now_ms() + timeout.as_millis() as i128;
464    let result = target
465        .call(procedure, realm, payload, deadline_ms, id, timeout)
466        .await
467        .map_err(CallError::Call);
468    target.close("normal", None, id).await;
469    result
470}
471
472/// Publishes a signed `procedure_advertisement` naming `session`'s own
473/// currently-connected station (`session.station.node_id`) as `procedure`'s
474/// server, discoverable by any caller's [`resolve`]/[`call`]. Mirrors
475/// `macula_response:advertise_direct/6,7` +
476/// `macula_direct_dial:publish_advertisement/4,5` — unlike the Erlang
477/// reference's pool (many links, one chosen by `connected_station/1`), a
478/// [`Session`] is always exactly one connection, so there is no
479/// link-selection step: the session's own verified HELLO identity IS the
480/// serving station.
481///
482/// **Sends the ordinary ADVERTISE frame first, then publishes the DHT
483/// record** — matching `macula_response:advertise_direct/7`'s own body
484/// exactly (`case advertise(Pool, Realm, Procedure, Module, Args, Opts) of
485/// {ok, Sup} -> ... macula_direct_dial:publish_advertisement(...)`). The
486/// DHT record is an ADDITIONAL discovery path for a caller on a different
487/// station to skip inter-station gossip propagation — it is not a
488/// substitute for the station actually knowing to route inbound CALLs
489/// here. **Found live, 2026-08-30**: an earlier version of this function
490/// (and its `macula-go` port, same gap, not yet fixed there as of this
491/// writing) published only the DHT record — a direct-dial caller could
492/// resolve and dial the right station, but the station itself had never
493/// been told to route the call anywhere, so every call still failed with
494/// `unknown_next_peer` despite a perfectly valid, resolvable, trusted
495/// advertisement. Caught by a live test that, unlike the earlier
496/// direct-dial verification, actually tried to get a real RESULT back
497/// instead of accepting `unknown_next_peer` as the expected terminal state.
498///
499/// Unlike the Erlang SDK's supervised `macula_response`, this does not
500/// itself keep anything alive — it does not spawn a responder process, so
501/// a caller still needs its own [`Session::serve_one_call`](crate::connection::Session::serve_one_call)
502/// loop to actually answer what gets routed here. A station's registration
503/// for a procedure does not survive the connection that sent it being
504/// replaced, so a long-lived server needs to call this again on its own
505/// schedule; see [`keep_advertised_direct`] for that loop.
506pub async fn advertise_direct(
507    session: &mut Session,
508    id: &KeyPair,
509    realm: [u8; 32],
510    procedure: &str,
511    ttl: Duration,
512) -> Result<(), AdvertiseDirectError> {
513    let advertise_spec = crate::frame::AdvertiseSpec::new(realm, procedure, id.node_id());
514    session
515        .advertise(&advertise_spec, id)
516        .await
517        .map_err(AdvertiseDirectError::Advertise)?;
518
519    let uri = dht::discovery_uri(realm, procedure);
520    let rec = dht::new_procedure_advertisement(id.node_id(), uri, session.station.node_id, ttl);
521    let rec = dht::sign(rec, id);
522    dht::put_record(session, id, &rec)
523        .await
524        .map_err(AdvertiseDirectError::Dht)
525}
526
527/// [`advertise_direct`] plus an embedded X.509 service-cert chain, for
528/// Slice 7c Direction B managed-realm authorization — see
529/// [`resolve_with_cert_chain`]/[`call_with_cert_chain`] for the
530/// corresponding checks. Opt-in: plain [`advertise_direct`] is unaffected.
531pub async fn advertise_direct_with_cert_chain(
532    session: &mut Session,
533    id: &KeyPair,
534    realm: [u8; 32],
535    procedure: &str,
536    ttl: Duration,
537    cert_chain_pem: Vec<u8>,
538) -> Result<(), AdvertiseDirectError> {
539    let advertise_spec = crate::frame::AdvertiseSpec::new(realm, procedure, id.node_id());
540    session
541        .advertise(&advertise_spec, id)
542        .await
543        .map_err(AdvertiseDirectError::Advertise)?;
544
545    let uri = dht::discovery_uri(realm, procedure);
546    let rec = dht::new_procedure_advertisement_with_cert_chain(
547        id.node_id(),
548        uri,
549        session.station.node_id,
550        ttl,
551        cert_chain_pem,
552    );
553    let rec = dht::sign(rec, id);
554    dht::put_record(session, id, &rec)
555        .await
556        .map_err(AdvertiseDirectError::Dht)
557}
558
559#[derive(Debug)]
560pub enum AdvertiseDirectError {
561    /// The ordinary station-side ADVERTISE frame failed to send.
562    Advertise(connection::SendFrameError),
563    /// The ordinary ADVERTISE succeeded, but publishing the direct-dial
564    /// DHT record failed — the procedure IS now reachable via ordinary
565    /// advertise-gossip, just not via direct-dial resolution.
566    Dht(DhtError),
567}
568
569impl std::fmt::Display for AdvertiseDirectError {
570    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
571        match self {
572            AdvertiseDirectError::Advertise(e) => write!(f, "direct_dial: sending ADVERTISE: {e}"),
573            AdvertiseDirectError::Dht(e) => write!(f, "direct_dial: {e}"),
574        }
575    }
576}
577
578impl std::error::Error for AdvertiseDirectError {}
579
580/// Calls [`advertise_direct`] immediately, then again every `interval`,
581/// until `stop` resolves. Rust has nothing equivalent to
582/// `macula_response`'s `reuse_sup` to worry about here, because
583/// [`advertise_direct`] (unlike Erlang's `advertise/5`, which spawns a real
584/// per-call OTP supervisor) is already a stateless, side-effect-free-on-
585/// repeat async function: nothing is created per tick that could leak —
586/// same reasoning `macula-go`'s `KeepAdvertisedDirect` already applied
587/// and verified live.
588///
589/// `interval` should leave real margin before `ttl` expires — production
590/// practice in `hecate-om`'s own capability re-advertise loop (the actual
591/// consumer of `advertise_direct`'s `reuse_sup` option on the Erlang side)
592/// uses a 4x margin: a 30s republish interval against a 120s record TTL.
593///
594/// A failed tick (network blip, connection genuinely dead, etc.) is
595/// reported via `on_error` but does NOT stop the loop; it tries again at
596/// the next interval regardless, matching `hecate-om`'s own log-and-continue
597/// practice around every DHT publish. This loop cannot detect or repair a
598/// dead `session` on its own — if its underlying connection has actually
599/// gone down, every tick will keep failing the same way until `stop`
600/// resolves; reconnecting a dead session is a separate, larger concern this
601/// does not attempt to solve.
602///
603/// 8 parameters: a target (`session`/`realm`/`procedure`), a re-advertise
604/// schedule (`id`/`ttl`/`interval`), and two independent callbacks
605/// (`stop`/`on_error`) with no natural sub-grouping — folding any of them
606/// into a synthetic struct would relocate the count, not reduce it.
607#[allow(clippy::too_many_arguments)]
608pub async fn keep_advertised_direct<F>(
609    session: &mut Session,
610    id: &KeyPair,
611    realm: [u8; 32],
612    procedure: &str,
613    ttl: Duration,
614    interval: Duration,
615    stop: F,
616    on_error: impl Fn(AdvertiseDirectError),
617) where
618    F: Future<Output = ()>,
619{
620    tokio::pin!(stop);
621    let mut ticker = tokio::time::interval(interval);
622    loop {
623        tokio::select! {
624            _ = &mut stop => return,
625            _ = ticker.tick() => {
626                if let Err(e) = advertise_direct(session, id, realm, procedure, ttl).await {
627                    on_error(e);
628                }
629            }
630        }
631    }
632}
633
634/// The dial-then-pin sequence every direct-dial call shape needs after
635/// resolving: dial `resolved`'s host:port, then check the freshly
636/// connected session's own signature-verified HELLO identity against
637/// `resolved.station` — factored out here (unlike [`call`]/
638/// [`call_with_cert_chain`], which had it inline before this existed)
639/// because [`open_stream_direct`]/[`put_direct`]/[`get_direct`] all need
640/// the identical sequence against a station identity that isn't
641/// necessarily reached via [`resolve`].
642#[derive(Debug)]
643pub enum DialAndVerifyError {
644    Dial(connection::HandshakeError),
645    /// The dialed peer's own signature-verified HELLO identity didn't
646    /// match the pubkey the signed DHT chain resolved — a trust
647    /// violation, not a retryable error.
648    TrustViolation {
649        resolved: [u8; 32],
650        dialed: [u8; 32],
651    },
652}
653
654impl std::fmt::Display for DialAndVerifyError {
655    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
656        match self {
657            DialAndVerifyError::Dial(e) => write!(f, "direct_dial: dialing resolved station: {e}"),
658            DialAndVerifyError::TrustViolation { resolved, dialed } => write!(
659                f,
660                "direct_dial: trust violation -- resolved station {} but the dialed peer proved identity {}",
661                hex_of(resolved),
662                hex_of(dialed)
663            ),
664        }
665    }
666}
667
668impl std::error::Error for DialAndVerifyError {}
669
670async fn dial_and_verify(
671    host: &str,
672    port: u16,
673    station: [u8; 32],
674    id: &KeyPair,
675    timeout: Duration,
676) -> Result<Session, DialAndVerifyError> {
677    let target = tokio::time::timeout(
678        timeout,
679        connection::connect(host, port, Trust::Insecure, id),
680    )
681    .await
682    .unwrap_or(Err(connection::HandshakeError::Timeout))
683    .map_err(DialAndVerifyError::Dial)?;
684
685    if target.station.node_id != station {
686        let dialed = target.station.node_id;
687        target.close("trust_violation", None, id).await;
688        return Err(DialAndVerifyError::TrustViolation {
689            resolved: station,
690            dialed,
691        });
692    }
693    Ok(target)
694}
695
696#[derive(Debug)]
697pub enum OpenStreamDirectError {
698    Resolve(ResolveError),
699    Dial(DialAndVerifyError),
700    Open(stream::OpenError),
701}
702
703impl std::fmt::Display for OpenStreamDirectError {
704    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
705        match self {
706            OpenStreamDirectError::Resolve(e) => write!(f, "{e}"),
707            OpenStreamDirectError::Dial(e) => write!(f, "{e}"),
708            OpenStreamDirectError::Open(e) => write!(f, "direct_dial: open stream: {e}"),
709        }
710    }
711}
712
713impl std::error::Error for OpenStreamDirectError {}
714
715/// Resolves `procedure`'s provider via direct-dial (through `resolve_via`,
716/// used only to query the DHT) and opens a stream there, in one hop, in a
717/// SEPARATE connection from `resolve_via` — the streaming-RPC counterpart
718/// to [`call`]. The provider must have advertised via [`advertise_direct`]:
719/// streaming's provider side (`macula_streamer.erl`) shares the identical
720/// `procedure_advertisement` mechanism RPC uses (confirmed against
721/// `macula_streamer.erl`/`macula_stream_sink.erl`'s own `advertise_direct`/
722/// `start_link_direct` — both are `macula_response:advertise_direct`/
723/// `macula_direct_dial:call_stream` under the hood, nothing stream-specific
724/// added), so no separate stream-shaped advertise function exists or is
725/// needed.
726///
727/// The caller owns the returned [`Session`] (and must close it once the
728/// stream and any other work on it is done) alongside the
729/// [`StreamHandle`] itself, since — unlike [`call`], which owns its dial
730/// for exactly one request/reply — a stream outlives the single function
731/// call that opens it.
732#[allow(clippy::too_many_arguments)]
733pub async fn open_stream_direct(
734    resolve_via: &mut Session,
735    id: &KeyPair,
736    realm: [u8; 32],
737    procedure: &str,
738    mode: StreamMode,
739    args: Value,
740    deadline_ms: i128,
741    timeout: Duration,
742) -> Result<(Session, StreamHandle), OpenStreamDirectError> {
743    let resolved = resolve(resolve_via, id, realm, procedure)
744        .await
745        .map_err(OpenStreamDirectError::Resolve)?;
746    let mut target = dial_and_verify(&resolved.host, resolved.port, resolved.station, id, timeout)
747        .await
748        .map_err(OpenStreamDirectError::Dial)?;
749    match StreamHandle::open(&mut target, procedure, realm, mode, args, deadline_ms, id).await {
750        Ok(handle) => Ok((target, handle)),
751        Err(e) => {
752            target.close("normal", None, id).await;
753            Err(OpenStreamDirectError::Open(e))
754        }
755    }
756}
757
758/// [`open_stream_direct`], resolved via [`resolve_with_cert_chain`]
759/// instead of [`resolve`] — see both for the full contract. Opt-in
760/// managed-realm authorization; [`open_stream_direct`] itself is
761/// unaffected.
762#[allow(clippy::too_many_arguments)]
763pub async fn open_stream_direct_with_cert_chain(
764    resolve_via: &mut Session,
765    id: &KeyPair,
766    realm: [u8; 32],
767    procedure: &str,
768    realm_ca_pem: &[u8],
769    expected_org: &str,
770    mode: StreamMode,
771    args: Value,
772    deadline_ms: i128,
773    timeout: Duration,
774) -> Result<(Session, StreamHandle), OpenStreamDirectError> {
775    let resolved = resolve_with_cert_chain(
776        resolve_via,
777        id,
778        realm,
779        procedure,
780        realm_ca_pem,
781        expected_org,
782    )
783    .await
784    .map_err(OpenStreamDirectError::Resolve)?;
785    let mut target = dial_and_verify(&resolved.host, resolved.port, resolved.station, id, timeout)
786        .await
787        .map_err(OpenStreamDirectError::Dial)?;
788    match StreamHandle::open(&mut target, procedure, realm, mode, args, deadline_ms, id).await {
789        Ok(handle) => Ok((target, handle)),
790        Err(e) => {
791            target.close("normal", None, id).await;
792            Err(OpenStreamDirectError::Open(e))
793        }
794    }
795}
796
797#[derive(Debug)]
798pub enum PutDirectError {
799    Resolve(ResolveError),
800    Dial(DialAndVerifyError),
801    Put(content::PutError),
802}
803
804impl std::fmt::Display for PutDirectError {
805    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
806        match self {
807            PutDirectError::Resolve(e) => write!(f, "{e}"),
808            PutDirectError::Dial(e) => write!(f, "{e}"),
809            PutDirectError::Put(e) => write!(f, "direct_dial: {e}"),
810        }
811    }
812}
813
814impl std::error::Error for PutDirectError {}
815
816/// Stores `data` at a KNOWN `station` directly, in one hop, instead of
817/// going through whatever station `resolve_via` happens to be connected
818/// to. Mirrors `macula_feeder:start_link_direct/5,6`, which — unlike
819/// procedure/stream direct-dial — takes the target station's pubkey
820/// directly rather than resolving one via a `procedure_advertisement`:
821/// content has no "procedure" to advertise, so there is nothing to
822/// resolve here beyond the station's own `station_endpoint`
823/// (`resolve_station_endpoint`). `resolve_via` is used only to query the
824/// DHT for `station`'s `station_endpoint`; it does not need to already be
825/// connected to `station`.
826///
827/// **Caveat found live in `macula-go`'s port of this same function**:
828/// if `resolve_via` happens to already be connected to `station` (the
829/// common case when the caller doesn't have a separate resolver session),
830/// this call's own internal dial reuses `id` against the SAME station
831/// `resolve_via` is on — this fleet enforces one connection per identity
832/// and kicks whichever connects second, so `resolve_via`'s own connection
833/// can be closed out from under the caller by this call. Use a different
834/// identity for `resolve_via` than for `id` if the caller needs
835/// `resolve_via` to keep working afterward against that same station.
836pub async fn put_direct(
837    resolve_via: &mut Session,
838    id: &KeyPair,
839    station: [u8; 32],
840    data: &[u8],
841    name: impl Into<String>,
842    timeout: Duration,
843) -> Result<Mcid, PutDirectError> {
844    let resolved = resolve_station_endpoint(resolve_via, id, station)
845        .await
846        .map_err(PutDirectError::Resolve)?;
847    let mut target = dial_and_verify(&resolved.host, resolved.port, resolved.station, id, timeout)
848        .await
849        .map_err(PutDirectError::Dial)?;
850    let result = content::put(&mut target, data, name, id)
851        .await
852        .map_err(PutDirectError::Put);
853    target.close("normal", None, id).await;
854    result
855}
856
857/// `mcid` has no live, verifiable `content_announcement` in the DHT —
858/// either nobody announced it (common: a single-block content put alone is
859/// never announced, matching `macula_content_transfer:put_single_block/3`),
860/// or every candidate found failed signature/self-consistency
861/// verification.
862#[derive(Debug)]
863pub struct ContentNotAnnounced;
864
865impl std::fmt::Display for ContentNotAnnounced {
866    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
867        write!(
868            f,
869            "direct_dial: content has no verifiable announcement in the DHT"
870        )
871    }
872}
873
874impl std::error::Error for ContentNotAnnounced {}
875
876#[derive(Debug)]
877pub enum GetDirectError {
878    Dht(DhtError),
879    NotAnnounced(ContentNotAnnounced),
880    /// A `content_announcement`'s `endpoint` field wasn't a dialable
881    /// `host:port` or URL.
882    EndpointParse(String),
883    Dial(DialAndVerifyError),
884    Get(content::GetError),
885}
886
887impl std::fmt::Display for GetDirectError {
888    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
889        match self {
890            GetDirectError::Dht(e) => write!(f, "direct_dial: find content providers: {e}"),
891            GetDirectError::NotAnnounced(e) => write!(f, "{e}"),
892            GetDirectError::EndpointParse(endpoint) => {
893                write!(
894                    f,
895                    "direct_dial: content provider endpoint {endpoint:?}: not a URL or host:port"
896                )
897            }
898            GetDirectError::Dial(e) => write!(f, "{e}"),
899            GetDirectError::Get(e) => write!(f, "direct_dial: {e}"),
900        }
901    }
902}
903
904impl std::error::Error for GetDirectError {}
905
906/// Fetches and verifies the content addressed by `mcid` from whichever
907/// station a signed `content_announcement` names as its host, dialing
908/// that station in one hop instead of relaying through `resolve_via`'s own
909/// station. Mirrors `macula_direct_dial:get_content/3`.
910///
911/// **Architectural note this module's other direct-dial functions don't
912/// need**: a `content_announcement`'s `endpoint` is the FINAL dial target
913/// directly (see `macula_record:read_content_announcement/1`'s `endpoint`
914/// field and `macula:get_content_station/5`'s use of it as-is) — unlike
915/// `procedure_advertisement`, there is no station-relay indirection, so
916/// the announcer must genuinely BE independently dialable there. A plain
917/// outbound-only leaf (everything this SDK's own identity/session model
918/// supports) cannot legitimately publish one of these about itself — only
919/// something with its own listening identity (`macula-station`, or a
920/// dedicated content-serving relay) can; confirmed directly against
921/// `macula.erl`, which states a `content_announcement` is made
922/// "automatically by the station on receipt," not by an arbitrary
923/// publisher. This crate therefore does not expose a client-facing
924/// "announce content direct": [`dht::new_content_announcement`] stays a
925/// low-level primitive (mirroring `macula_record.erl`'s own export) for
926/// that kind of infrastructure-tier code, not ordinary leaf use.
927/// [`get_direct`] itself has no such limitation — resolving and fetching
928/// FROM an already-announced provider is a perfectly ordinary leaf
929/// operation.
930pub async fn get_direct(
931    resolve_via: &mut Session,
932    id: &KeyPair,
933    mcid: Mcid,
934    timeout: Duration,
935) -> Result<Vec<u8>, GetDirectError> {
936    let recs = dht::find_records(resolve_via, id, dht::content_key(mcid))
937        .await
938        .map_err(GetDirectError::Dht)?;
939    let adv = first_trusted_content_provider(&recs)
940        .ok_or(GetDirectError::NotAnnounced(ContentNotAnnounced))?;
941    let (host, port) = parse_seed_url(&adv.endpoint)
942        .ok_or_else(|| GetDirectError::EndpointParse(adv.endpoint.clone()))?;
943    let mut target = dial_and_verify(&host, port, adv.announcer_node, id, timeout)
944        .await
945        .map_err(GetDirectError::Dial)?;
946    let result = content::get(&mut target, mcid, id)
947        .await
948        .map_err(GetDirectError::Get);
949    target.close("normal", None, id).await;
950    result
951}
952
953/// Mirrors `macula.erl`'s `decode_provider/1`: the record's OWN signature
954/// must verify, AND the payload's claimed `announcer_node` must equal the
955/// record's own envelope key — a record merely stored under the right key
956/// but self-signed by a different identity would otherwise still be
957/// trusted.
958fn first_trusted_content_provider(recs: &[Record]) -> Option<dht::ContentAnnouncement> {
959    recs.iter().find_map(|rec| {
960        dht::verify(rec).ok()?;
961        let adv = dht::read_content_announcement(rec).ok()?;
962        (adv.announcer_node == rec.key).then_some(adv)
963    })
964}
965
966/// Splits a `content_announcement`'s `endpoint` (a dialable seed URL, e.g.
967/// `"https://host:4433"` — `macula_client:seed()`'s own format) into the
968/// host/port pair [`connection::connect`] wants. Distinct from
969/// `station_endpoint`'s already-split `host_advertised`/`quic_port`
970/// fields — `content_announcement` embeds a single ready-to-dial URL
971/// instead. Tolerates a bare `host:port` with no scheme too, matching this
972/// crate's own tolerance elsewhere for a station config given without one.
973fn parse_seed_url(seed: &str) -> Option<(String, u16)> {
974    if let Some(rest) = seed
975        .strip_prefix("https://")
976        .or_else(|| seed.strip_prefix("http://"))
977    {
978        let hostport = rest.split('/').next().unwrap_or(rest);
979        let (host, port_str) = hostport.rsplit_once(':')?;
980        return Some((host.to_string(), port_str.parse().ok()?));
981    }
982    let (host, port_str) = seed.rsplit_once(':')?;
983    Some((host.to_string(), port_str.parse().ok()?))
984}