Distributed Erlang over iroh P2P + TLS, with a bot powers SDK
0

Configure Feed

Select the types of activity you want to include in your feed.

erl-iroh / src / bot.rs
16 kB 438 lines
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}