//! the dev/debug cli harness. //! //! NOT the product. the product is the gui; this is how we see inside the //! machine during development. it exists to exercise the async paths (oauth, //! repo publish, iroh handshake, message io) without egui in the way, and to //! smoke-test against a real pds with two terminals. //! //! it is gated behind the `debug` subcommand; running iwakura with no args (or //! an unknown one) goes to the gui. when the gui lands, these subcommands //! become the instrumented driver for it. use std::sync::Arc; use atrium_api::types::string::Did; use thiserror::Error; use crate::auth::flow; use crate::lexicon::{DidKey, Message}; use crate::repo::RepoClient; use crate::runtime; #[derive(Error, Debug)] pub enum Error { #[error("runtime error: {0}")] Runtime(#[from] runtime::Error), #[error("auth flow error: {0}")] Flow(#[from] flow::Error), #[error("repo error: {0}")] Repo(#[from] crate::repo::Error), #[error("net error: {0}")] Net(#[from] crate::net::NodeError), #[error("handshake error: {0}")] Handshake(#[from] crate::net::HandshakeError), #[error("bad argument: {0}")] Arg(String), } /// the parsed debug subcommand. `[account]` args are optional and default to /// the first session in `auth.jsonl`; they accept a handle OR a did. pub enum Command { /// `debug login ` — full oauth flow, persist the session. Login { handle: String }, /// `debug whoami [account]` — restore a session and print the identity. Whoami { account: Option }, /// `debug publish [account]` — publish this device's endpoint record. Publish { account: Option }, /// `debug endpoint-id` — print this device's iroh endpoint id / did:key. EndpointId, /// `debug listen [account]` — accept one peer, print identity + messages. Listen { account: Option }, /// `debug connect [account]` — connect to a peer, /// handshake, send one message. Connect { account: Option, endpoint: String, message: String }, /// `debug accounts` — list persisted sessions and the default account. Accounts, } const USAGE: &str = "\ iwakura debug harness (dev tool; the product is the gui) usage: iwakura launch the gui iwakura debug login sign in (oauth), persist the session iwakura debug accounts list sessions + the default account iwakura debug whoami [account] print the active identity iwakura debug endpoint-id print this device's endpoint id iwakura debug publish [account] publish this device's endpoint record iwakura debug listen [account] accept peers and chat with them iwakura debug connect [account] connect to a peer, send a message, then keep chatting interactively in a listen session, type a line and hit enter to message the CURRENT peer: /peers list connected peers /peer switch which peer you're talking to notes: - [account] is a handle (mfzx.net) or a did (did:plc:…). handles are resolved with the bidirectional handle↔did check. - [account] defaults to the first session in auth.jsonl (the earliest sign-in), so after `login` you can omit it. - is the peer's endpoint: a did:key, or a handle/did whose endpoint record we look up and dial via its advertised relay/addr. - a peer's session is dropped if their endpoint record lapses (ttl re-check). - the endpoint record's rkey is the device did:key. - two terminals, same machine: `listen` in one, `connect` in the other. "; /// parse argv. returns `None` if the gui should run (no/unknown subcommand), /// `Some(Ok(cmd))` for a debug command, `Some(Err(_))` for a bad invocation. pub fn parse(args: &[String]) -> Option> { // expect: ["iwakura", "debug", , ...] if args.len() < 2 || args[1] != "debug" { return None; } let sub = args.get(2).map(String::as_str).unwrap_or(""); let arg = |i: usize, name: &str| -> Result { args.get(i).cloned().ok_or_else(|| Error::Arg(format!("missing <{name}>\n\n{USAGE}"))) }; // optional positional account arg. let opt = |i: usize| -> Option { args.get(i).cloned() }; // build the command, propagating a missing-arg error as Some(Err(_)). let built = (|| -> Result { Ok(match sub { "login" => Command::Login { handle: arg(3, "handle")? }, "accounts" => Command::Accounts, "whoami" => Command::Whoami { account: opt(3) }, "endpoint-id" => Command::EndpointId, "publish" => Command::Publish { account: opt(3) }, "listen" => Command::Listen { account: opt(3) }, "connect" => Command::Connect { endpoint: arg(3, "endpoint-id")?, message: arg(4, "message")?, account: opt(5), }, "help" | "-h" | "--help" | "" => { eprint!("{USAGE}"); std::process::exit(if sub.is_empty() { 1 } else { 0 }); } other => { return Err(Error::Arg(format!("unknown debug subcommand {other:?}\n\n{USAGE}"))); } }) })(); Some(built) } /// resolve an optional account argument to a did. /// /// with `Some(handle_or_did)`, resolves via the handle↔did bidirectional /// check. with `None`, falls back to the first session in `auth.jsonl` (the /// default account set by the earliest `login`). async fn resolve_account(account: Option<&str>) -> Result { match account { Some(input) => flow::resolve_account(input).await.map_err(|e| Error::Arg(e.to_string())), None => { let store = crate::auth::store::JsonlSessionStore::open(crate::auth::auth_file()) .await .map_err(|e| Error::Arg(format!("could not open session store: {e}")))?; store.first_did().await.ok_or_else(|| { Error::Arg( "no account given and no default (auth.jsonl is empty; run `debug login ` first)".into(), ) }) } } } /// run a parsed debug command to completion. pub async fn run(cmd: Command) -> Result<(), Error> { match cmd { Command::Login { handle } => login(&handle).await, Command::Accounts => accounts().await, Command::Whoami { account } => whoami(account.as_deref()).await, Command::EndpointId => endpoint_id(), Command::Publish { account } => publish(account.as_deref()).await, Command::Listen { account } => listen(account.as_deref()).await, Command::Connect { account, endpoint, message } => { connect(account.as_deref(), &endpoint, &message).await } } } async fn login(handle: &str) -> Result<(), Error> { let client = runtime::oauth_client().await?; let (did, _agent) = flow::authorize_agent(&client, handle, flow::open_in_browser).await?; println!("signed in as {}", did.as_str()); println!("session persisted to {}", crate::auth::auth_file().display()); println!("this is now the default account (first entry in auth.jsonl)"); Ok(()) } async fn accounts() -> Result<(), Error> { let store = crate::auth::store::JsonlSessionStore::open(crate::auth::auth_file()) .await .map_err(|e| Error::Arg(format!("could not open session store: {e}")))?; let dids = store.list_dids().await; let default = store.first_did().await; if dids.is_empty() { println!("no sessions (run `debug login ` first)"); return Ok(()); } for did in &dids { let marker = if Some(did) == default.as_ref() { " (default)" } else { "" }; println!("{}{marker}", did.as_str()); } Ok(()) } async fn whoami(account: Option<&str>) -> Result<(), Error> { let did = resolve_account(account).await?; let client = runtime::oauth_client().await?; let _agent = flow::restore_agent(&client, &did).await?; println!("restored session for {}", did.as_str()); Ok(()) } fn endpoint_id() -> Result<(), Error> { let key = runtime::load_or_create_device_key()?; let did_key = DidKey::from_verifying_key(&key.verifying_key()); println!("endpoint id (did:key): {did_key}"); println!("device key file: {}", runtime::device_key_file().display()); Ok(()) } async fn publish(account: Option<&str>) -> Result<(), Error> { let did = resolve_account(account).await?; let client = runtime::oauth_client().await?; let agent = flow::restore_agent(&client, &did).await?; let repo = RepoClient::new(agent).await?; let device_key = runtime::load_or_create_device_key()?; let record = runtime::make_endpoint_record(did.as_str(), &device_key)?; let uri = repo.publish_endpoint(did.as_str(), &record).await?; println!("published endpoint record: {uri}"); println!(" sub (endpoint): {}", record.sub); println!(" via: {:?}", record.via); println!(" ttl: {}s", record.ttl); Ok(()) } async fn listen(account: Option<&str>) -> Result<(), Error> { let did = resolve_account(account).await?; let client = runtime::oauth_client().await?; let agent = flow::restore_agent(&client, &did).await?; let repo = Arc::new(RepoClient::new(agent).await?); let device_key = runtime::load_or_create_device_key()?; let our_record = Arc::new(runtime::make_endpoint_record(did.as_str(), &device_key)?); let node = Arc::new(runtime::node(&device_key).await?); let our_handle = flow::did_to_handle(&did).await; println!("listening as {}", our_handle.as_deref().unwrap_or(did.as_str())); println!("endpoint id: {}", DidKey::from_verifying_key(&device_key.verifying_key())); println!("waiting for peers (ctrl+d on an empty line to quit)…"); // the hub tracks active sessions and broadcasts stdin to all of them. let hub = SessionHub::new(Arc::clone(&repo)); // accept loop: accept a raw connection, then hand it to a spawned task for // the handshake + session setup. a slow/failing handshake never wedges the // loop. the loop ends when the endpoint closes (we quit via stdin below). let accept_repo = Arc::clone(&repo); let accept_node = Arc::clone(&node); let accept_record = Arc::clone(&our_record); let accept_hub = hub.clone(); tokio::spawn(async move { loop { let incoming = match accept_node.accept_connection().await { Ok(c) => c, Err(_) => break, // endpoint closed }; let record = Arc::clone(&accept_record); let verify = runtime::presence_check(Arc::clone(&accept_repo)); let node2 = Arc::clone(&accept_node); let hub2 = accept_hub.clone(); // per-connection task: handshake, register session, run ttl // re-check + recv printer. failures here are isolated to this peer. tokio::spawn(async move { let session = match node2.handshake_responder(incoming, &record, &verify).await { Ok(s) => std::sync::Arc::new(s), Err(e) => { eprintln!("(rejected a peer: {e})"); return; } }; hub2.add_and_run(session).await; }); } }); // foreground: read stdin. plain lines go to the CURRENT peer; slash // commands manage which peer that is (/peers, /peer ). use tokio::io::AsyncBufReadExt; let mut lines = tokio::io::BufReader::new(tokio::io::stdin()).lines(); let our_name = our_handle.unwrap_or_else(|| "me".to_string()); while let Ok(Some(line)) = lines.next_line().await { let line = line.trim_end(); if line.is_empty() { continue; } if let Some(cmd) = line.strip_prefix('/') { handle_peer_command(&hub, cmd).await; continue; } match hub.current().await { None => eprintln!("(no current peer; /peers to list, /peer to pick)"), Some((peer_name, session)) => match Message::new(line) { Some(msg) => { if let Err(e) = session.send_message(&msg).await { eprintln!("(send to {peer_name} failed: {e})"); } else { println!("[{our_name} -> {peer_name}] {line}"); } } None => eprintln!("(message too long; max {} bytes)", crate::lexicon::message::MAX_BYTES), }, } } node.close().await; Ok(()) } /// handle a slash command in the listener's stdin loop. /// /// - `/peers` — list connected peers, marking the current one. /// - `/peer ` — make the n-th (1-based) connected peer current. /// /// anything else prints a hint. async fn handle_peer_command(hub: &std::sync::Arc, cmd: &str) { let mut parts = cmd.split_whitespace(); match parts.next() { Some("peers") => { let peers = hub.list().await; if peers.is_empty() { println!("(no peers connected)"); } for (i, (name, is_current)) in peers.iter().enumerate() { println!(" {} {}{}", i + 1, name, if *is_current { " (current)" } else { "" }); } } Some("peer") => { match parts.next().and_then(|n| n.parse::().ok()) { Some(n) => match hub.set_current(n).await { Ok(name) => println!("(now chatting with {name})"), Err(e) => eprintln!("({e})"), }, None => eprintln!("(usage: /peer )"), } } _ => eprintln!("(commands: /peers, /peer )"), } } /// one tracked peer: its session plus a display name (handle, or did). type PeerEntry = (String, std::sync::Arc); /// tracks active peer sessions for the persistent listener. /// /// each added session gets a recv printer and a ttl re-check task. stdin is /// directed at a chosen CURRENT peer (not broadcast); /peers and /peer /// manage which one. this is the gui's session-manager shape minus the /// window — kept deliberately simple (removal happens when a recv loop or the /// ttl re-check notices a close). struct SessionHub { /// (name, session) entries in connection order. sessions: std::sync::Arc>>, /// index into `sessions` of the peer stdin is directed at, if any. current: std::sync::Arc>>, repo: std::sync::Arc, } impl SessionHub { fn new(repo: std::sync::Arc) -> std::sync::Arc { std::sync::Arc::new(Self { sessions: std::sync::Arc::new(tokio::sync::Mutex::new(Vec::new())), current: std::sync::Arc::new(tokio::sync::Mutex::new(None)), repo, }) } /// the current peer's (name, session), if one is set and still present. async fn current(&self) -> Option { let sessions = self.sessions.lock().await; let idx = *self.current.lock().await; idx.and_then(|i| sessions.get(i)).map(|(n, s)| (n.clone(), std::sync::Arc::clone(s))) } /// (name, is_current) for every tracked peer, in order. async fn list(&self) -> Vec<(String, bool)> { let sessions = self.sessions.lock().await; let idx = *self.current.lock().await; sessions .iter() .enumerate() .map(|(i, (name, _))| (name.clone(), Some(i) == idx)) .collect() } /// make the n-th (1-based) peer current. returns its name on success. async fn set_current(&self, n: usize) -> Result { let sessions = self.sessions.lock().await; if n == 0 || n > sessions.len() { return Err(format!("no peer #{n}; /peers to list")); } *self.current.lock().await = Some(n - 1); Ok(sessions[n - 1].0.clone()) } /// remove a session from tracking, fixing up the current index. async fn remove(&self, session: &std::sync::Arc) { let mut sessions = self.sessions.lock().await; let mut current = self.current.lock().await; let pos = sessions.iter().position(|(_, s)| std::sync::Arc::ptr_eq(s, session)); if let Some(pos) = pos { sessions.remove(pos); *current = match *current { // the current peer left: pick whoever's first now, if anyone. Some(c) if c == pos => if sessions.is_empty() { None } else { Some(0) }, // a peer before the current one left: shift the index down. Some(c) if c > pos => Some(c - 1), other => other, }; } } /// register a session and spawn its recv printer + ttl re-check. async fn add_and_run(self: &std::sync::Arc, session: std::sync::Arc) { let peer_did = session.peer.user_did.clone(); let peer_name = match peer_did.parse::() { Ok(d) => flow::did_to_handle(&d).await.unwrap_or_else(|| peer_did.clone()), Err(_) => peer_did.clone(), }; println!("peer connected: {peer_name}"); { let mut sessions = self.sessions.lock().await; let mut current = self.current.lock().await; sessions.push((peer_name.clone(), std::sync::Arc::clone(&session))); // if there's no current peer yet, this one becomes current. if current.is_none() { *current = Some(sessions.len() - 1); println!("(now chatting with {peer_name})"); } } // recv printer: print incoming messages until the peer closes. { let session = std::sync::Arc::clone(&session); let name = peer_name.clone(); let hub = std::sync::Arc::clone(self); tokio::spawn(async move { while let Ok(msg) = session.recv_message().await { println!("[{name}] {}", msg.text); } println!("({name} disconnected)"); hub.remove(&session).await; }); } // ttl re-check: periodically confirm the peer's endpoint record is // still present+valid (expiry-is-revocation). if it lapses, close the // session. this is the spec's 'periodically checking for // freshness/revocation according to the record's ttl'. { let session = std::sync::Arc::clone(&session); let name = peer_name.clone(); let hub = std::sync::Arc::clone(self); let ttl = (session.peer.record.ttl.max(1)) as u64; let peer_iss = session.peer.user_did.clone(); let peer_endpoint = session.peer.endpoint; tokio::spawn(async move { let mut interval = tokio::time::interval(std::time::Duration::from_secs(ttl)); interval.tick().await; // skip the immediate tick loop { interval.tick().await; // stop re-checking once the session is gone. let still_active = hub .sessions .lock() .await .iter() .any(|(_, s)| std::sync::Arc::ptr_eq(s, &session)); if !still_active { break; } // re-run the presence check against the peer's repo. match hub.repo.check_endpoint_presence(&peer_iss, &peer_endpoint).await { Ok(crate::repo::Presence::Active) => { tracing::debug!(peer = %name, "ttl re-check: still active") } Ok(other) => { println!("({name}'s endpoint record is no longer active: {other:?}; closing)"); session.close().await; hub.remove(&session).await; break; } Err(e) => { // a transient resolution/network error is not a // revocation; log and try again next interval. tracing::warn!(peer = %name, error = %e, "ttl re-check errored; will retry") } } } }); } } } /// a symmetric two-way chat loop over an established session. /// /// both ends of a connection run this after the handshake: a background task /// prints incoming messages, while the foreground reads stdin and sends each /// line. the connection stays open for the whole conversation. exits when /// stdin closes (ctrl+d) or the peer disconnects. async fn chat( session: std::sync::Arc, peer_handle: Option, our_handle: Option, ) -> Result<(), Error> { let peer_name = peer_handle.unwrap_or_else(|| session.peer.user_did.clone()); println!("— chatting with {peer_name}. type a line and hit enter; ctrl+d to quit. —"); // background: print incoming messages until the peer closes. let recv_session: std::sync::Arc = std::sync::Arc::clone(&session); let recv_name = peer_name.clone(); let recv_task = tokio::spawn(async move { // print incoming messages until the peer closes / stream errors. while let Ok(msg) = recv_session.recv_message().await { println!("[{recv_name}] {}", msg.text); } }); // foreground: read stdin lines and send them. use tokio::io::AsyncBufReadExt; let mut lines = tokio::io::BufReader::new(tokio::io::stdin()).lines(); let our_name = our_handle.unwrap_or_else(|| "me".to_string()); loop { match lines.next_line().await { Ok(Some(line)) => { let line = line.trim_end(); if line.is_empty() { continue; } match Message::new(line) { Some(msg) => { if let Err(e) = session.send_message(&msg).await { eprintln!("send failed ({e}); connection closed"); break; } // echo our own line locally so the transcript reads whole. println!("[{our_name}] {line}"); } None => eprintln!("(message too long; max {} bytes)", crate::lexicon::message::MAX_BYTES), } } Ok(None) => break, // stdin closed (ctrl+d) Err(e) => { eprintln!("stdin error: {e}"); break; } } } recv_task.abort(); Ok(()) } /// resolve an endpoint argument to a dialable iroh address. /// /// accepts: /// - a did:key (`did:key:z6Mk…`) — a bare endpoint id, dialed via iroh's /// default address discovery. /// - a handle or account did (`mfzx.net`, `did:plc:…`) — we look up the /// account's endpoint record(s) and dial via the advertised relay/addr. /// /// returns the dialable address and the did:key of the endpoint we chose. async fn resolve_endpoint( repo: &RepoClient, endpoint: &str, ) -> Result<(iroh::EndpointAddr, DidKey), Error> { // case 1: a bare did:key endpoint id. if let Ok(dk) = DidKey::parse(endpoint) { let id = iroh::EndpointId::from_bytes(dk.public_key_bytes()) .map_err(|e| Error::Arg(format!("bad endpoint id: {e}")))?; return Ok((iroh::EndpointAddr::from(id), dk)); } // case 2: a handle or account did. resolve to a did, list its active // endpoint records, and dial the first via its advertised reachability. let account = flow::resolve_account(endpoint).await.map_err(|e| Error::Arg(e.to_string()))?; let endpoints = repo .list_active_endpoints(account.as_str()) .await .map_err(|e| Error::Arg(format!("could not list endpoints for {endpoint}: {e}")))?; let record = endpoints.into_iter().next().ok_or_else(|| { Error::Arg(format!( "{endpoint} ({}) has no active endpoint records;\n\ have them run `debug publish` on a device first.", account.as_str() )) })?; let dk = record.subject().map_err(|e| Error::Arg(format!("bad endpoint record: {e}")))?; let addr = record .to_endpoint_addr() .map_err(|e| Error::Arg(format!("could not build dial address: {e}")))?; Ok((addr, dk)) } async fn connect(account: Option<&str>, endpoint: &str, message: &str) -> Result<(), Error> { let did = resolve_account(account).await?; let client = runtime::oauth_client().await?; let agent = flow::restore_agent(&client, &did).await?; let repo = Arc::new(RepoClient::new(agent).await?); let verify = runtime::presence_check(Arc::clone(&repo)); let (peer_addr, peer_endpoint) = resolve_endpoint(&repo, endpoint).await?; let device_key = runtime::load_or_create_device_key()?; let our_record = runtime::make_endpoint_record(did.as_str(), &device_key)?; let node = runtime::node(&device_key).await?; println!("connecting to {endpoint} ({peer_endpoint})…"); let session = node.connect(peer_addr, &our_record, &verify).await?; let peer_did: Did = session .peer .user_did .parse() .map_err(|_| Error::Arg("peer presented a bad did".into()))?; let peer_handle = flow::did_to_handle(&peer_did).await; let our_handle = flow::did_to_handle(&did).await; println!( "handshake ok. peer: {}", peer_handle.as_deref().unwrap_or(session.peer.user_did.as_str()) ); // send the opening message from the arguments, then drop into the // interactive chat so the conversation can continue. let msg = Message::new(message).ok_or_else(|| Error::Arg("message too long".into()))?; session.send_message(&msg).await?; println!( "[{}] {}", our_handle.as_deref().unwrap_or("me"), msg.text ); chat(std::sync::Arc::new(session), peer_handle, our_handle).await?; node.close().await; Ok(()) }