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}