Distributed Erlang over iroh P2P + TLS, with a bot powers SDK
1//! Bot SDK — import this, get powers.
2//!
3//! A [`Bot`] is a named distribution endpoint with an iroh identity, optional
4//! TLS config, and a set of [`Powers`] that gate high-level operations.
5
6use crate::error::{Error, Result};
7use crate::powers::Powers;
8use crate::session::{PeerId, Session, SessionHandle, SessionRegistry};
9use crate::term;
10use crate::transport::iroh::IrohTransport;
11use crate::transport::tcp::TcpTransport;
12use crate::transport::tls::{TlsConfig, TlsTransport};
13use erl_dist::term::{Atom, List, Term};
14use iroh::{EndpointId, SecretKey};
15use std::sync::Arc;
16use tokio::net::TcpListener;
17use tokio::sync::Mutex;
18use tracing::{info, warn};
19
20/// A distribution-capable bot with gated powers.
21pub struct Bot {
22 name: String,
23 cookie: String,
24 powers: Powers,
25 iroh: IrohTransport,
26 sessions: SessionRegistry,
27 tls: Option<TlsConfig>,
28 accept_tasks: Arc<Mutex<Vec<tokio::task::JoinHandle<()>>>>,
29}
30
31/// Builder for [`Bot`].
32pub struct BotBuilder {
33 name: String,
34 cookie: Option<String>,
35 powers: Powers,
36 secret: Option<SecretKey>,
37 tls: Option<TlsConfig>,
38 host: String,
39}
40
41impl Bot {
42 /// Start building a bot named `name` (short name; host defaults to `iroh`).
43 ///
44 /// Full node name becomes `{name}@{host}` (host overridable via builder).
45 pub fn builder(name: impl Into<String>) -> BotBuilder {
46 BotBuilder {
47 name: name.into(),
48 cookie: None,
49 powers: Powers::default(),
50 secret: None,
51 tls: None,
52 host: "iroh".into(),
53 }
54 }
55
56 /// Full Erlang-style node name (`scout@iroh`).
57 pub fn node_name(&self) -> &str {
58 &self.name
59 }
60
61 /// Granted powers.
62 pub fn powers(&self) -> Powers {
63 self.powers
64 }
65
66 /// This bot's iroh [`EndpointId`] (share to let peers dial you).
67 pub fn endpoint_id(&self) -> EndpointId {
68 self.iroh.endpoint_id()
69 }
70
71 /// iroh endpoint id as a string.
72 pub fn endpoint_id_str(&self) -> String {
73 self.iroh.endpoint_id().to_string()
74 }
75
76 /// Underlying iroh transport.
77 pub fn iroh(&self) -> &IrohTransport {
78 &self.iroh
79 }
80
81 // ── cluster powers ──────────────────────────────────────────────────
82
83 /// Dial another bot / node over iroh by endpoint id (base32 string).
84 pub async fn connect_iroh(&self, peer_endpoint_id: &str) -> Result<SessionHandle> {
85 self.powers.require(Powers::CLUSTER)?;
86 let stream = self.iroh.connect(peer_endpoint_id).await?;
87 let peer_id = PeerId::Endpoint(peer_endpoint_id.to_string());
88 let session = Session::connect(stream, &self.name, &self.cookie, peer_id).await?;
89 let handle = session.handle();
90 info!(peer = %handle.peer_name(), "connected over iroh");
91 self.sessions.insert(&session).await;
92 // Keep session alive by leaking into accept_tasks list via spawn of empty future...
93 // Session's JoinHandle is inside Session; we store the Session so the runner lives.
94 self.keep(session).await;
95 Ok(handle)
96 }
97
98 /// Dial by full iroh [`iroh::EndpointAddr`].
99 pub async fn connect_iroh_addr(&self, addr: iroh::EndpointAddr) -> Result<SessionHandle> {
100 self.powers.require(Powers::CLUSTER)?;
101 let id = addr.id.to_string();
102 let stream = self.iroh.connect_addr(addr).await?;
103 let peer_id = PeerId::Endpoint(id);
104 let session = Session::connect(stream, &self.name, &self.cookie, peer_id).await?;
105 let handle = session.handle();
106 self.sessions.insert(&session).await;
107 self.keep(session).await;
108 Ok(handle)
109 }
110
111 /// Dial a classic BEAM node via EPMD + TCP (no TLS).
112 pub async fn connect_tcp(&self, node: &str, cookie: Option<&str>) -> Result<SessionHandle> {
113 self.powers.require(Powers::CLUSTER)?;
114 let stream = TcpTransport::connect_node(node).await?;
115 let cookie = cookie.unwrap_or(&self.cookie);
116 let peer_id = PeerId::Node(node.to_string());
117 let session = Session::connect(stream, &self.name, cookie, peer_id).await?;
118 let handle = session.handle();
119 info!(peer = %handle.peer_name(), "connected over tcp");
120 self.sessions.insert(&session).await;
121 self.keep(session).await;
122 Ok(handle)
123 }
124
125 /// Dial `host:port` with TLS, then run the Erlang dist handshake.
126 ///
127 /// `node` is the Erlang node name used in the handshake (`app@host`).
128 /// `server_name` is the TLS SNI / certificate name.
129 pub async fn connect_tls(
130 &self,
131 node: &str,
132 host: &str,
133 port: u16,
134 server_name: &str,
135 cookie: Option<&str>,
136 ) -> Result<SessionHandle> {
137 self.powers.require(Powers::CLUSTER)?;
138 let tls = self
139 .tls
140 .as_ref()
141 .ok_or_else(|| Error::tls("no TlsConfig configured on bot"))?;
142 let stream = TlsTransport::connect(host, port, server_name, tls).await?;
143 let cookie = cookie.unwrap_or(&self.cookie);
144 let peer_id = PeerId::Node(node.to_string());
145 let session = Session::connect(stream, &self.name, cookie, peer_id).await?;
146 let handle = session.handle();
147 info!(peer = %handle.peer_name(), "connected over tls");
148 self.sessions.insert(&session).await;
149 self.keep(session).await;
150 Ok(handle)
151 }
152
153 /// Start accepting inbound iroh distribution connections in the background.
154 pub async fn listen_iroh(&self) -> Result<()> {
155 self.powers.require(Powers::LISTEN)?;
156 let iroh = self.iroh.clone();
157 let name = self.name.clone();
158 let cookie = self.cookie.clone();
159 let sessions = self.sessions.clone();
160 let keep = self.accept_tasks.clone();
161
162 let task = tokio::spawn(async move {
163 loop {
164 match iroh.accept().await {
165 Ok((stream, remote)) => {
166 let peer_id = PeerId::Endpoint(remote.to_string());
167 match Session::accept(stream, &name, &cookie, peer_id).await {
168 Ok(session) => {
169 info!(peer = %session.peer_node().name, remote = %remote, "accepted iroh peer");
170 sessions.insert(&session).await;
171 // Park the session so its runner stays alive.
172 let parked = Arc::new(session);
173 let k = keep.clone();
174 k.lock().await.push(tokio::spawn(async move {
175 // Hold Arc until cancelled.
176 std::future::pending::<()>().await;
177 drop(parked);
178 }));
179 }
180 Err(e) => warn!(error = %e, "accept handshake failed"),
181 }
182 }
183 Err(e) => {
184 warn!(error = %e, "iroh accept error");
185 break;
186 }
187 }
188 }
189 });
190 self.accept_tasks.lock().await.push(task);
191 info!(id = %self.endpoint_id(), "listening for iroh peers");
192 Ok(())
193 }
194
195 /// Accept one TLS distribution connection on `listener`.
196 pub async fn accept_tls(&self, listener: &TcpListener) -> Result<SessionHandle> {
197 self.powers.require(Powers::LISTEN)?;
198 let tls = self
199 .tls
200 .as_ref()
201 .ok_or_else(|| Error::tls("no TlsConfig configured on bot"))?;
202 let stream = TlsTransport::accept(listener, tls).await?;
203 let peer_id = PeerId::Node("tls-peer".into());
204 let session = Session::accept(stream, &self.name, &self.cookie, peer_id).await?;
205 let handle = session.handle();
206 self.sessions.insert(&session).await;
207 self.keep(session).await;
208 Ok(handle)
209 }
210
211 /// List connected peer keys.
212 pub async fn peers(&self) -> Vec<String> {
213 self.sessions.list().await
214 }
215
216 /// Get a session handle by peer name or endpoint id key.
217 pub async fn session(&self, peer: &str) -> Result<SessionHandle> {
218 self.sessions
219 .get(peer)
220 .await
221 .ok_or_else(|| Error::NotConnected(peer.to_string()))
222 }
223
224 // ── rpc / send powers ───────────────────────────────────────────────
225
226 /// RPC to a connected peer: `module:function(args)`.
227 pub async fn rpc(
228 &self,
229 peer: &str,
230 module: &str,
231 function: &str,
232 args: List,
233 ) -> Result<Term> {
234 self.powers.require(Powers::RPC)?;
235 let s = self.session(peer).await?;
236 s.rpc(Atom::from(module), Atom::from(function), args).await
237 }
238
239 /// Convenience: `erlang:node/0` on peer.
240 pub async fn remote_node_name(&self, peer: &str) -> Result<Term> {
241 self.rpc(peer, "erlang", "node", term::nil()).await
242 }
243
244 /// Send a term to a registered process on a peer.
245 pub async fn send(&self, peer: &str, to: &str, msg: Term) -> Result<()> {
246 self.powers.require(Powers::SEND)?;
247 let s = self.session(peer).await?;
248 s.send(Atom::from(to), msg).await
249 }
250
251 // ── introspect powers (thin RPC wrappers) ───────────────────────────
252
253 /// List processes on a peer (`erlang:processes/0`).
254 pub async fn processes(&self, peer: &str) -> Result<Term> {
255 self.powers.require(Powers::INTROSPECT)?;
256 self.rpc(peer, "erlang", "processes", term::nil()).await
257 }
258
259 /// Which applications are running (`application:which_applications/0`).
260 pub async fn applications(&self, peer: &str) -> Result<Term> {
261 self.powers.require(Powers::INTROSPECT)?;
262 self.rpc(peer, "application", "which_applications", term::nil())
263 .await
264 }
265
266 /// Process info for a pid term.
267 pub async fn process_info(&self, peer: &str, pid: Term) -> Result<Term> {
268 self.powers.require(Powers::INTROSPECT)?;
269 self.rpc(peer, "erlang", "process_info", term::list([pid]))
270 .await
271 }
272
273 /// `ets:all/0`.
274 pub async fn ets_tables(&self, peer: &str) -> Result<Term> {
275 self.powers.require(Powers::INTROSPECT)?;
276 self.rpc(peer, "ets", "all", term::nil()).await
277 }
278
279 /// Ping a peer via `net_adm:ping/1`. Requires the peer to be a real BEAM node.
280 pub async fn ping(&self, via: &str, target_node: &str) -> Result<Term> {
281 self.powers.require(Powers::RPC)?;
282 self.rpc(
283 via,
284 "net_adm",
285 "ping",
286 term::list([term::atom(target_node).into()]),
287 )
288 .await
289 }
290
291 // ── eval power ──────────────────────────────────────────────────────
292
293 /// Evaluate an Erlang expression string via `erl_eval` / `erl_parse`.
294 ///
295 /// Uses: `erl_eval:exprs(element(2, erl_parse:parse_exprs(element(2, erl_scan:string(S)))), [])`
296 /// simplified through a small helper MFA when available; falls back to
297 /// `rpc:call` style via `erl_eval:expr`.
298 pub async fn eval(&self, peer: &str, expr: &str) -> Result<Term> {
299 self.powers.require(Powers::EVAL)?;
300 // {ok, Tokens, _} = erl_scan:string(Expr ++ ".").
301 // {ok, Parsed} = erl_parse:parse_exprs(Tokens).
302 // {value, Value, _} = erl_eval:exprs(Parsed, []).
303 // We ship a multi-step via erpc by calling a one-liner through `erlang:apply`.
304 // Simpler approach: use `erl_eval:eval_str` is not std — compose:
305 let code = format!("{expr}.");
306 let scan = self
307 .rpc(
308 peer,
309 "erl_scan",
310 "string",
311 term::list([term::string(&code)]),
312 )
313 .await?;
314 let tokens = extract_ok_tuple(&scan, 1)?;
315 let parsed = self
316 .rpc(peer, "erl_parse", "parse_exprs", term::list([tokens]))
317 .await?;
318 let forms = extract_ok_tuple(&parsed, 1)?;
319 let value = self
320 .rpc(
321 peer,
322 "erl_eval",
323 "exprs",
324 term::list([forms, term::nil().into()]),
325 )
326 .await?;
327 // {value, V, Bindings}
328 if let Term::Tuple(t) = &value {
329 if t.elements.len() >= 2 {
330 if term::as_atom(&t.elements[0]) == Some("value") {
331 return Ok(t.elements[1].clone());
332 }
333 }
334 }
335 Ok(value)
336 }
337
338 /// Gracefully shut down accept loops and the iroh endpoint.
339 pub async fn shutdown(self) {
340 let mut tasks = self.accept_tasks.lock().await;
341 for t in tasks.drain(..) {
342 t.abort();
343 }
344 drop(tasks);
345 self.iroh.close().await;
346 }
347
348 // ── internals ───────────────────────────────────────────────────────
349
350 async fn keep(&self, session: Session) {
351 let parked = Arc::new(session);
352 self.accept_tasks.lock().await.push(tokio::spawn(async move {
353 std::future::pending::<()>().await;
354 drop(parked);
355 }));
356 }
357}
358
359impl BotBuilder {
360 /// Erlang cookie shared with peers.
361 pub fn cookie(mut self, cookie: impl Into<String>) -> Self {
362 self.cookie = Some(cookie.into());
363 self
364 }
365
366 /// Grant powers.
367 pub fn powers(mut self, powers: Powers) -> Self {
368 self.powers = powers;
369 self
370 }
371
372 /// Stable iroh identity.
373 pub fn secret_key(mut self, key: SecretKey) -> Self {
374 self.secret = Some(key);
375 self
376 }
377
378 /// TLS config for classic TCP+TLS paths.
379 pub fn tls(mut self, tls: TlsConfig) -> Self {
380 self.tls = Some(tls);
381 self
382 }
383
384 /// Host part of the node name (default `iroh` → `name@iroh`).
385 pub fn host(mut self, host: impl Into<String>) -> Self {
386 self.host = host.into();
387 self
388 }
389
390 /// Build and bind the iroh endpoint.
391 pub async fn build(self) -> Result<Bot> {
392 let cookie = self
393 .cookie
394 .ok_or_else(|| Error::InvalidName("cookie is required".into()))?;
395 let full_name = if self.name.contains('@') {
396 self.name
397 } else {
398 format!("{}@{}", self.name, self.host)
399 };
400
401 let iroh = if let Some(secret) = self.secret {
402 IrohTransport::bind_with_key(secret).await?
403 } else {
404 IrohTransport::bind().await?
405 };
406
407 info!(
408 node = %full_name,
409 endpoint = %iroh.endpoint_id(),
410 powers = %self.powers,
411 "bot ready"
412 );
413
414 Ok(Bot {
415 name: full_name,
416 cookie,
417 powers: self.powers,
418 iroh,
419 sessions: SessionRegistry::default(),
420 tls: self.tls,
421 accept_tasks: Arc::new(Mutex::new(Vec::new())),
422 })
423 }
424}
425
426fn extract_ok_tuple(term: &Term, idx: usize) -> Result<Term> {
427 if let Term::Tuple(t) = term {
428 if let Some(Term::Atom(tag)) = t.elements.first() {
429 if tag.name == "ok" && t.elements.len() > idx {
430 return Ok(t.elements[idx].clone());
431 }
432 if tag.name == "error" {
433 return Err(Error::rpc(format!("{term}")));
434 }
435 }
436 }
437 Err(Error::rpc(format!("expected {{ok, …}}, got {term}")))
438}