//! Bot SDK — import this, get powers. //! //! A [`Bot`] is a named distribution endpoint with an iroh identity, optional //! TLS config, and a set of [`Powers`] that gate high-level operations. use crate::error::{Error, Result}; use crate::powers::Powers; use crate::session::{PeerId, Session, SessionHandle, SessionRegistry}; use crate::term; use crate::transport::iroh::IrohTransport; use crate::transport::tcp::TcpTransport; use crate::transport::tls::{TlsConfig, TlsTransport}; use erl_dist::term::{Atom, List, Term}; use iroh::{EndpointId, SecretKey}; use std::sync::Arc; use tokio::net::TcpListener; use tokio::sync::Mutex; use tracing::{info, warn}; /// A distribution-capable bot with gated powers. pub struct Bot { name: String, cookie: String, powers: Powers, iroh: IrohTransport, sessions: SessionRegistry, tls: Option, accept_tasks: Arc>>>, } /// Builder for [`Bot`]. pub struct BotBuilder { name: String, cookie: Option, powers: Powers, secret: Option, tls: Option, host: String, } impl Bot { /// Start building a bot named `name` (short name; host defaults to `iroh`). /// /// Full node name becomes `{name}@{host}` (host overridable via builder). pub fn builder(name: impl Into) -> BotBuilder { BotBuilder { name: name.into(), cookie: None, powers: Powers::default(), secret: None, tls: None, host: "iroh".into(), } } /// Full Erlang-style node name (`scout@iroh`). pub fn node_name(&self) -> &str { &self.name } /// Granted powers. pub fn powers(&self) -> Powers { self.powers } /// This bot's iroh [`EndpointId`] (share to let peers dial you). pub fn endpoint_id(&self) -> EndpointId { self.iroh.endpoint_id() } /// iroh endpoint id as a string. pub fn endpoint_id_str(&self) -> String { self.iroh.endpoint_id().to_string() } /// Underlying iroh transport. pub fn iroh(&self) -> &IrohTransport { &self.iroh } // ── cluster powers ────────────────────────────────────────────────── /// Dial another bot / node over iroh by endpoint id (base32 string). pub async fn connect_iroh(&self, peer_endpoint_id: &str) -> Result { self.powers.require(Powers::CLUSTER)?; let stream = self.iroh.connect(peer_endpoint_id).await?; let peer_id = PeerId::Endpoint(peer_endpoint_id.to_string()); let session = Session::connect(stream, &self.name, &self.cookie, peer_id).await?; let handle = session.handle(); info!(peer = %handle.peer_name(), "connected over iroh"); self.sessions.insert(&session).await; // Keep session alive by leaking into accept_tasks list via spawn of empty future... // Session's JoinHandle is inside Session; we store the Session so the runner lives. self.keep(session).await; Ok(handle) } /// Dial by full iroh [`iroh::EndpointAddr`]. pub async fn connect_iroh_addr(&self, addr: iroh::EndpointAddr) -> Result { self.powers.require(Powers::CLUSTER)?; let id = addr.id.to_string(); let stream = self.iroh.connect_addr(addr).await?; let peer_id = PeerId::Endpoint(id); let session = Session::connect(stream, &self.name, &self.cookie, peer_id).await?; let handle = session.handle(); self.sessions.insert(&session).await; self.keep(session).await; Ok(handle) } /// Dial a classic BEAM node via EPMD + TCP (no TLS). pub async fn connect_tcp(&self, node: &str, cookie: Option<&str>) -> Result { self.powers.require(Powers::CLUSTER)?; let stream = TcpTransport::connect_node(node).await?; let cookie = cookie.unwrap_or(&self.cookie); let peer_id = PeerId::Node(node.to_string()); let session = Session::connect(stream, &self.name, cookie, peer_id).await?; let handle = session.handle(); info!(peer = %handle.peer_name(), "connected over tcp"); self.sessions.insert(&session).await; self.keep(session).await; Ok(handle) } /// Dial `host:port` with TLS, then run the Erlang dist handshake. /// /// `node` is the Erlang node name used in the handshake (`app@host`). /// `server_name` is the TLS SNI / certificate name. pub async fn connect_tls( &self, node: &str, host: &str, port: u16, server_name: &str, cookie: Option<&str>, ) -> Result { self.powers.require(Powers::CLUSTER)?; let tls = self .tls .as_ref() .ok_or_else(|| Error::tls("no TlsConfig configured on bot"))?; let stream = TlsTransport::connect(host, port, server_name, tls).await?; let cookie = cookie.unwrap_or(&self.cookie); let peer_id = PeerId::Node(node.to_string()); let session = Session::connect(stream, &self.name, cookie, peer_id).await?; let handle = session.handle(); info!(peer = %handle.peer_name(), "connected over tls"); self.sessions.insert(&session).await; self.keep(session).await; Ok(handle) } /// Start accepting inbound iroh distribution connections in the background. pub async fn listen_iroh(&self) -> Result<()> { self.powers.require(Powers::LISTEN)?; let iroh = self.iroh.clone(); let name = self.name.clone(); let cookie = self.cookie.clone(); let sessions = self.sessions.clone(); let keep = self.accept_tasks.clone(); let task = tokio::spawn(async move { loop { match iroh.accept().await { Ok((stream, remote)) => { let peer_id = PeerId::Endpoint(remote.to_string()); match Session::accept(stream, &name, &cookie, peer_id).await { Ok(session) => { info!(peer = %session.peer_node().name, remote = %remote, "accepted iroh peer"); sessions.insert(&session).await; // Park the session so its runner stays alive. let parked = Arc::new(session); let k = keep.clone(); k.lock().await.push(tokio::spawn(async move { // Hold Arc until cancelled. std::future::pending::<()>().await; drop(parked); })); } Err(e) => warn!(error = %e, "accept handshake failed"), } } Err(e) => { warn!(error = %e, "iroh accept error"); break; } } } }); self.accept_tasks.lock().await.push(task); info!(id = %self.endpoint_id(), "listening for iroh peers"); Ok(()) } /// Accept one TLS distribution connection on `listener`. pub async fn accept_tls(&self, listener: &TcpListener) -> Result { self.powers.require(Powers::LISTEN)?; let tls = self .tls .as_ref() .ok_or_else(|| Error::tls("no TlsConfig configured on bot"))?; let stream = TlsTransport::accept(listener, tls).await?; let peer_id = PeerId::Node("tls-peer".into()); let session = Session::accept(stream, &self.name, &self.cookie, peer_id).await?; let handle = session.handle(); self.sessions.insert(&session).await; self.keep(session).await; Ok(handle) } /// List connected peer keys. pub async fn peers(&self) -> Vec { self.sessions.list().await } /// Get a session handle by peer name or endpoint id key. pub async fn session(&self, peer: &str) -> Result { self.sessions .get(peer) .await .ok_or_else(|| Error::NotConnected(peer.to_string())) } // ── rpc / send powers ─────────────────────────────────────────────── /// RPC to a connected peer: `module:function(args)`. pub async fn rpc( &self, peer: &str, module: &str, function: &str, args: List, ) -> Result { self.powers.require(Powers::RPC)?; let s = self.session(peer).await?; s.rpc(Atom::from(module), Atom::from(function), args).await } /// Convenience: `erlang:node/0` on peer. pub async fn remote_node_name(&self, peer: &str) -> Result { self.rpc(peer, "erlang", "node", term::nil()).await } /// Send a term to a registered process on a peer. pub async fn send(&self, peer: &str, to: &str, msg: Term) -> Result<()> { self.powers.require(Powers::SEND)?; let s = self.session(peer).await?; s.send(Atom::from(to), msg).await } // ── introspect powers (thin RPC wrappers) ─────────────────────────── /// List processes on a peer (`erlang:processes/0`). pub async fn processes(&self, peer: &str) -> Result { self.powers.require(Powers::INTROSPECT)?; self.rpc(peer, "erlang", "processes", term::nil()).await } /// Which applications are running (`application:which_applications/0`). pub async fn applications(&self, peer: &str) -> Result { self.powers.require(Powers::INTROSPECT)?; self.rpc(peer, "application", "which_applications", term::nil()) .await } /// Process info for a pid term. pub async fn process_info(&self, peer: &str, pid: Term) -> Result { self.powers.require(Powers::INTROSPECT)?; self.rpc(peer, "erlang", "process_info", term::list([pid])) .await } /// `ets:all/0`. pub async fn ets_tables(&self, peer: &str) -> Result { self.powers.require(Powers::INTROSPECT)?; self.rpc(peer, "ets", "all", term::nil()).await } /// Ping a peer via `net_adm:ping/1`. Requires the peer to be a real BEAM node. pub async fn ping(&self, via: &str, target_node: &str) -> Result { self.powers.require(Powers::RPC)?; self.rpc( via, "net_adm", "ping", term::list([term::atom(target_node).into()]), ) .await } // ── eval power ────────────────────────────────────────────────────── /// Evaluate an Erlang expression string via `erl_eval` / `erl_parse`. /// /// Uses: `erl_eval:exprs(element(2, erl_parse:parse_exprs(element(2, erl_scan:string(S)))), [])` /// simplified through a small helper MFA when available; falls back to /// `rpc:call` style via `erl_eval:expr`. pub async fn eval(&self, peer: &str, expr: &str) -> Result { self.powers.require(Powers::EVAL)?; // {ok, Tokens, _} = erl_scan:string(Expr ++ "."). // {ok, Parsed} = erl_parse:parse_exprs(Tokens). // {value, Value, _} = erl_eval:exprs(Parsed, []). // We ship a multi-step via erpc by calling a one-liner through `erlang:apply`. // Simpler approach: use `erl_eval:eval_str` is not std — compose: let code = format!("{expr}."); let scan = self .rpc( peer, "erl_scan", "string", term::list([term::string(&code)]), ) .await?; let tokens = extract_ok_tuple(&scan, 1)?; let parsed = self .rpc(peer, "erl_parse", "parse_exprs", term::list([tokens])) .await?; let forms = extract_ok_tuple(&parsed, 1)?; let value = self .rpc( peer, "erl_eval", "exprs", term::list([forms, term::nil().into()]), ) .await?; // {value, V, Bindings} if let Term::Tuple(t) = &value { if t.elements.len() >= 2 { if term::as_atom(&t.elements[0]) == Some("value") { return Ok(t.elements[1].clone()); } } } Ok(value) } /// Gracefully shut down accept loops and the iroh endpoint. pub async fn shutdown(self) { let mut tasks = self.accept_tasks.lock().await; for t in tasks.drain(..) { t.abort(); } drop(tasks); self.iroh.close().await; } // ── internals ─────────────────────────────────────────────────────── async fn keep(&self, session: Session) { let parked = Arc::new(session); self.accept_tasks.lock().await.push(tokio::spawn(async move { std::future::pending::<()>().await; drop(parked); })); } } impl BotBuilder { /// Erlang cookie shared with peers. pub fn cookie(mut self, cookie: impl Into) -> Self { self.cookie = Some(cookie.into()); self } /// Grant powers. pub fn powers(mut self, powers: Powers) -> Self { self.powers = powers; self } /// Stable iroh identity. pub fn secret_key(mut self, key: SecretKey) -> Self { self.secret = Some(key); self } /// TLS config for classic TCP+TLS paths. pub fn tls(mut self, tls: TlsConfig) -> Self { self.tls = Some(tls); self } /// Host part of the node name (default `iroh` → `name@iroh`). pub fn host(mut self, host: impl Into) -> Self { self.host = host.into(); self } /// Build and bind the iroh endpoint. pub async fn build(self) -> Result { let cookie = self .cookie .ok_or_else(|| Error::InvalidName("cookie is required".into()))?; let full_name = if self.name.contains('@') { self.name } else { format!("{}@{}", self.name, self.host) }; let iroh = if let Some(secret) = self.secret { IrohTransport::bind_with_key(secret).await? } else { IrohTransport::bind().await? }; info!( node = %full_name, endpoint = %iroh.endpoint_id(), powers = %self.powers, "bot ready" ); Ok(Bot { name: full_name, cookie, powers: self.powers, iroh, sessions: SessionRegistry::default(), tls: self.tls, accept_tasks: Arc::new(Mutex::new(Vec::new())), }) } } fn extract_ok_tuple(term: &Term, idx: usize) -> Result { if let Term::Tuple(t) = term { if let Some(Term::Atom(tag)) = t.elements.first() { if tag.name == "ok" && t.elements.len() > idx { return Ok(t.elements[idx].clone()); } if tag.name == "error" { return Err(Error::rpc(format!("{term}"))); } } } Err(Error::rpc(format!("expected {{ok, …}}, got {term}"))) }