use std::time::{Duration, Instant};
use anyhow::{Context, Result};
#[cfg(unix)]
use libc;
use serde_json::json;
use crate::commands::dev;
use crate::commands::dev::agent::client::AgentHttpClient;
use crate::commands::dev::agent::session::AgentSession;
use crate::commands::dev::host::{self, DaemonHost, InstanceProfile, Mode};
use crate::commands::harness::bitcoind::Bitcoind;
use crate::commands::harness::state::{BitcoindState, ChannelState, HarnessState, InstanceState, state_path};
use crate::commands::harness::BitcoindMode;
pub fn up(
monorepo_path: std::path::PathBuf,
bitcoind_mode: BitcoindMode,
with_channel: bool,
clean: bool,
) -> Result<()> {
if clean {
if let Ok(path) = state_path() {
if path.exists() {
let _ = std::fs::remove_file(&path);
}
}
for instance in ["alice", "bob"] {
if let Ok(env_dir) = crate::commands::dev::host::monorepo::monorepo_dev_dir(
&monorepo_path,
instance,
None,
).map(|dev_dir| dev_dir.parent().unwrap_or(&dev_dir).to_path_buf()) {
for name in ["dev.db", "dev.db-shm", "dev.db-wal", "lightning.db"] {
let _ = std::fs::remove_file(env_dir.join(name));
}
for name in ["ldk_data", "ldk_node_data", "ldk_node_data_backup"] {
let _ = std::fs::remove_dir_all(env_dir.join(name));
}
let dev_dir = env_dir.join("dev-apps");
let _ = std::fs::remove_file(
crate::commands::dev::agent::session::AgentSession::file_path(&dev_dir, instance)
);
}
}
}
for instance in ["alice", "bob"] {
if let Ok(dev_dir) = crate::commands::dev::host::monorepo::monorepo_dev_dir(
&monorepo_path,
instance,
None,
) {
let session_exists =
crate::commands::dev::agent::session::AgentSession::file_path(&dev_dir, instance).exists();
let env_dir = dev_dir.parent().unwrap_or(&dev_dir).to_path_buf();
let db_exists = env_dir.join("dev.db").exists();
if db_exists && !session_exists {
eprintln!(
"harness: {instance} DB exists but session missing — \
wiping DB for clean re-onboard"
);
for name in ["dev.db", "dev.db-shm", "dev.db-wal", "lightning.db"] {
let _ = std::fs::remove_file(env_dir.join(name));
}
for name in ["ldk_data", "ldk_node_data", "ldk_node_data_backup"] {
let _ = std::fs::remove_dir_all(env_dir.join(name));
}
}
}
}
dev::install_shutdown_handler();
let btc = Bitcoind::ensure(bitcoind_mode).context("bring up regtest bitcoind")?;
let mode = Mode::Monorepo { path: monorepo_path.clone() };
let profiles = [InstanceProfile::alice(), InstanceProfile::bob()];
let hosts: Vec<Box<dyn DaemonHost>> = profiles
.iter()
.map(|p| host::for_mode(mode.clone(), p.clone(), None, None, None, false))
.collect();
let mut handles = Vec::new();
for h in &hosts {
handles.push(h.ensure_running().context("start platform daemon")?);
}
dev::agent::run_agent_setup(&handles, None);
let mut instances = Vec::new();
for profile in &profiles {
let name = profile.name.as_str();
let ldk_addr = format!("127.0.0.1:{}", profile.p2p_port);
let session = AgentSession::load_for_harness(&monorepo_path, name)
.with_context(|| format!("load {name} agent session"))?;
let dev_dir = crate::commands::dev::host::monorepo::monorepo_dev_dir(
&monorepo_path,
name,
None,
)?;
let session_path = AgentSession::file_path(&dev_dir, name);
instances.push(InstanceState {
name: name.into(),
session_path,
base_url: session.base_url.clone(),
ldk_addr,
node_id: session.node_id.clone(),
pid: None,
});
}
let mut state = HarnessState {
created_at: now_iso8601(),
bitcoind: BitcoindState {
mode: format!("{bitcoind_mode:?}").to_lowercase(),
rpc_url: btc.rpc_url.clone(),
rpc_user: crate::commands::harness::bitcoind::RPC_USER.into(),
container_id: btc.container_id.clone(),
},
instances,
channel: None,
supervisor_pid: Some(std::process::id()),
};
if with_channel {
let ch = open_channel_flow(&state, &btc, "alice", "bob", 100_000)?;
state.channel = Some(ch);
}
let path = state.save()?;
println!("{}", serde_json::to_string_pretty(&state)?);
eprintln!(
"HARNESS READY — supervising alice+bob; state at {}. \
Send SIGTERM (or `node-app harness down`) to stop.",
path.display()
);
while !dev::shutdown_requested() {
std::thread::sleep(Duration::from_millis(500));
}
eprintln!("harness: shutdown signal received — stopping alice+bob…");
for h in &hosts {
h.shutdown();
}
let _ = btc.stop();
eprintln!("harness: teardown complete.");
Ok(())
}
pub fn open_channel_flow(
state: &HarnessState,
btc: &Bitcoind,
from: &str,
to: &str,
sats: u64,
) -> Result<ChannelState> {
let from_i = state.instance(from)?;
let to_i = state.instance(to)?;
let from_tok = AgentSession::read_token(&from_i.session_path)?;
let client = AgentHttpClient::new(from_i.base_url.clone());
eprintln!("harness: funding {from} on-chain (1 BTC)…");
let addr = client.new_onchain_address(&from_tok)?;
btc.send_to_address(&addr, 1.0)?;
eprintln!("harness: mining 6 confirmation blocks…");
btc.mine(6)?;
std::thread::sleep(Duration::from_secs(3));
let peer_str = format!("{}@{}", to_i.node_id, to_i.ldk_addr);
eprintln!("harness: connecting {from} → {to} ({peer_str})…");
client.connect_peer(&from_tok, &peer_str)?;
let push_msat = (sats / 10) * 1_000;
eprintln!(
"harness: opening channel {from}→{to}: {sats} sats, pushing {} sats to {to}…",
push_msat / 1_000
);
client.open_channel(&from_tok, &peer_str, sats, push_msat)?;
eprintln!("harness: mining 6 channel confirmation blocks…");
btc.mine(6)?;
std::thread::sleep(Duration::from_secs(5));
eprintln!("harness: waiting for channel to become usable (120s deadline)…");
let deadline = Instant::now() + Duration::from_secs(120);
loop {
let channels = client.list_channels(&from_tok)?;
let is_usable = if let Some(arr) = channels.as_array() {
arr.iter().any(|ch| {
ch.get("is_usable").and_then(|v| v.as_bool()).unwrap_or(false)
|| ch.get("is_channel_ready").and_then(|v| v.as_bool()).unwrap_or(false)
})
} else {
false
};
if is_usable {
eprintln!("harness: channel {from}→{to} is usable");
return Ok(ChannelState {
from: from.into(),
to: to.into(),
capacity_sats: sats,
status: "usable".into(),
});
}
if Instant::now() >= deadline {
anyhow::bail!(
"channel {from}→{to} not usable within 120s; last: {channels}"
);
}
let _ = btc.mine(1);
std::thread::sleep(Duration::from_secs(2));
}
}
fn now_iso8601() -> String {
chrono::Utc::now().to_rfc3339()
}
pub fn status() -> Result<()> {
let state = load_state()?;
let mut result = serde_json::Map::new();
for inst in &state.instances {
let token = AgentSession::read_token(&inst.session_path)
.with_context(|| format!("read token for {}", inst.name))?;
let client = AgentHttpClient::new(inst.base_url.clone());
let balance = client.get_balance(&token)
.unwrap_or_else(|e| json!({ "error": e.to_string() }));
let channels = client.list_channels(&token)
.unwrap_or_else(|e| json!({ "error": e.to_string() }));
let peers = client.list_peers(&token)
.unwrap_or_else(|e| json!({ "error": e.to_string() }));
result.insert(inst.name.clone(), json!({
"balance": balance,
"channels": channels,
"peers": peers,
}));
}
println!("{}", serde_json::to_string_pretty(&serde_json::Value::Object(result))?);
Ok(())
}
pub fn mine(blocks: u32) -> Result<()> {
let state = load_state()?;
let btc = Bitcoind::from_state(&state.bitcoind);
btc.mine(blocks)?;
let height = btc.block_count()?;
println!("{}", json!({ "mined": blocks, "height": height }));
Ok(())
}
pub fn fund(node: &str, btc_amount: f64) -> Result<()> {
let state = load_state()?;
let inst = state.instance(node)?;
let token = AgentSession::read_token(&inst.session_path)
.with_context(|| format!("read token for {node}"))?;
let client = AgentHttpClient::new(inst.base_url.clone());
let addr = client.new_onchain_address(&token)?;
let btc_handle = Bitcoind::from_state(&state.bitcoind);
let txid = btc_handle.send_to_address(&addr, btc_amount)?;
btc_handle.mine(6)?;
println!("{}", json!({
"address": addr,
"txid": txid,
"confirmations": 6,
}));
Ok(())
}
pub fn channel_open(from: &str, to: &str, sats: u64) -> Result<()> {
let mut state = load_state()?;
let btc = Bitcoind::from_state(&state.bitcoind);
let ch = open_channel_flow(&state, &btc, from, to, sats)?;
let ch_json = serde_json::to_value(&ch)?;
state.channel = Some(ch);
state.save()?;
println!("{}", serde_json::to_string_pretty(&ch_json)?);
Ok(())
}
pub fn invoke(node: &str, capability: &str, payload_str: &str) -> Result<()> {
let state = load_state()?;
let inst = state.instance(node)?;
let token = AgentSession::read_token(&inst.session_path)
.with_context(|| format!("read token for {node}"))?;
let client = AgentHttpClient::new(inst.base_url.clone());
let payload: serde_json::Value = serde_json::from_str(payload_str)
.with_context(|| format!("parse payload JSON: {payload_str}"))?;
let response = client.invoke_capability(&token, capability, payload)?;
println!("{}", serde_json::to_string_pretty(&response)?);
Ok(())
}
pub fn logs(node: &str, tail: bool) -> Result<()> {
let state = load_state()?;
let inst = state.instance(node)?;
let instance_root = inst.session_path
.parent() .and_then(|p| p.parent()) .ok_or_else(|| anyhow::anyhow!("session_path has no grandparent"))?;
let log_path = instance_root.join("daemon.log");
if !log_path.exists() {
anyhow::bail!("log not found: {}", log_path.display());
}
let content = std::fs::read_to_string(&log_path)
.with_context(|| format!("read {}", log_path.display()))?;
if tail {
let lines: Vec<&str> = content.lines().collect();
let start = lines.len().saturating_sub(200);
for line in &lines[start..] {
println!("{}", line);
}
} else {
print!("{}", content);
}
Ok(())
}
pub fn down(clean: bool) -> Result<()> {
let state = match HarnessState::load() {
Ok(s) => s,
Err(_) => {
println!("{}", json!({ "stopped": false, "reason": "no running harness" }));
return Ok(());
}
};
#[cfg(unix)]
{
if let Some(pid) = state.supervisor_pid {
eprintln!("harness down: sending SIGTERM to supervisor PID {pid}");
let rc = unsafe { libc::kill(pid as libc::pid_t, libc::SIGTERM) };
if rc != 0 {
let still_alive = unsafe { libc::kill(pid as libc::pid_t, 0) } == 0;
if !still_alive {
eprintln!("harness down: supervisor {pid} is no longer running");
} else {
eprintln!("harness down: kill({pid}, SIGTERM) returned non-zero");
}
} else {
let deadline = Instant::now() + Duration::from_secs(40);
let mut timed_out = false;
loop {
std::thread::sleep(Duration::from_millis(300));
let probe = unsafe { libc::kill(pid as libc::pid_t, 0) };
if probe != 0 {
break;
}
if Instant::now() >= deadline {
eprintln!(
"harness down: supervisor {pid} did not exit within 40s — \
escalating to SIGKILL"
);
timed_out = true;
break;
}
}
if timed_out {
unsafe { libc::kill(pid as libc::pid_t, libc::SIGKILL) };
let _ = std::process::Command::new("pkill")
.args(["-f", "node-server --env"])
.status();
}
}
} else {
eprintln!("harness down: no supervisor_pid in state, falling back to best-effort teardown");
let _ = std::process::Command::new("pkill")
.args(["-f", "node-server --env"])
.status();
}
}
if let Some(id) = &state.bitcoind.container_id {
let _ = std::process::Command::new("docker").args(["stop", id]).status();
}
if clean {
let cache_root = state_path()
.ok()
.and_then(|p| {
let mut cur = p.parent()?.to_path_buf();
loop {
if cur.file_name().map(|n| n == "node-app").unwrap_or(false) {
return Some(cur);
}
if !cur.pop() {
return None;
}
}
});
for inst in &state.instances {
let instance_root = match inst.session_path
.parent() .and_then(|p| p.parent()) {
Some(d) => d.to_path_buf(),
None => continue,
};
let instance_root = instance_root.canonicalize().unwrap_or(instance_root);
let under_cache = cache_root.as_ref().map(|root| instance_root.starts_with(root)).unwrap_or(false);
if !under_cache {
eprintln!(
"harness down: skipping clean for {} (not under cache root)",
instance_root.display()
);
continue;
}
for name in ["dev.db", "dev.db-shm", "dev.db-wal", "lightning.db"] {
let _ = std::fs::remove_file(instance_root.join(name));
}
for name in ["ldk_data", "ldk_node_data", "ldk_node_data_backup"] {
let p = instance_root.join(name);
if p.is_dir() {
let _ = std::fs::remove_dir_all(&p);
}
}
}
}
let _ = std::fs::remove_file(state_path()?);
println!("{}", json!({ "stopped": true, "cleaned": clean }));
Ok(())
}
pub fn pay(from: &str, to: &str, sats: u64) -> Result<()> {
let state = load_state()?;
let from_i = state.instance(from)?;
let to_i = state.instance(to)?;
let from_tok = AgentSession::read_token(&from_i.session_path)
.with_context(|| format!("read token for {from}"))?;
let to_tok = AgentSession::read_token(&to_i.session_path)
.with_context(|| format!("read token for {to}"))?;
let from_client = AgentHttpClient::new(from_i.base_url.clone());
let to_client = AgentHttpClient::new(to_i.base_url.clone());
let bolt11 = to_client
.create_invoice(&to_tok, sats * 1_000, "harness pay")
.with_context(|| format!("{to} create_invoice failed"))?;
let pay_resp = from_client
.pay_invoice(&from_tok, &bolt11)
.map_err(|e| {
let msg = e.to_string();
if msg.to_lowercase().contains("route")
|| msg.to_lowercase().contains("channel")
|| msg.to_lowercase().contains("liquidity")
|| msg.to_lowercase().contains("no path")
{
anyhow::anyhow!(
"payment failed — no usable route/channel: {msg}\n\
hint: run `node-app harness channel-open {from} {to} 100000` first"
)
} else {
anyhow::anyhow!("pay_invoice failed: {msg}")
}
})?;
let payment_hash = pay_resp
.get("payment_hash")
.and_then(|v| v.as_str())
.unwrap_or("")
.to_string();
let initial_status = pay_resp
.get("status")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let initial_preimage = pay_resp
.get("preimage")
.and_then(|v| v.as_str())
.map(str::to_owned);
let (final_status, final_preimage) = if initial_status == "pending"
&& !payment_hash.is_empty()
{
eprintln!("harness pay: payment pending — polling status (30s deadline)…");
let deadline = Instant::now() + Duration::from_secs(30);
let mut status = initial_status.clone();
let mut preimage = initial_preimage.clone();
loop {
std::thread::sleep(Duration::from_millis(500));
match from_client.invoke_capability(
&from_tok,
"core.lightning.get_payment_status",
json!({ "payment_hash": payment_hash }),
) {
Ok(v) => {
status = v
.get("status")
.and_then(|s| s.as_str())
.unwrap_or("unknown")
.to_string();
preimage = v
.get("preimage")
.and_then(|p| p.as_str())
.map(str::to_owned);
if status == "succeeded" || status == "failed" {
break;
}
}
Err(e) => {
eprintln!("harness pay: get_payment_status error: {e}");
}
}
if Instant::now() >= deadline {
eprintln!("harness pay: payment still pending after 30s");
break;
}
}
(status, preimage)
} else {
(initial_status, initial_preimage)
};
println!(
"{}",
serde_json::to_string_pretty(&json!({
"invoice": bolt11,
"status": final_status,
"payment_hash": payment_hash,
"preimage": final_preimage,
}))?
);
if !matches!(final_status.as_str(), "succeeded" | "settled") {
anyhow::bail!(
"payment did not settle (status: {final_status}); \
check that a usable channel exists: `node-app harness channel-open {from} {to} 100000`"
);
}
Ok(())
}
pub fn l402(from: &str, to: &str, route: &str, payload_str: &str) -> Result<()> {
let state = load_state()?;
let from_i = state.instance(from)?;
let _to_i = state.instance(to)?;
let from_tok = AgentSession::read_token(&from_i.session_path)
.with_context(|| format!("read token for {from}"))?;
let payload: serde_json::Value = serde_json::from_str(payload_str)
.with_context(|| format!("parse payload JSON: {payload_str}"))?;
let mut payload = payload;
if let Some(obj) = payload.as_object_mut() {
obj.entry("execution_preference".to_string())
.or_insert(serde_json::Value::String("remote".to_string()));
}
let from_client = AgentHttpClient::new(from_i.base_url.clone());
let (status_code, body) = from_client
.post_raw_with_status(&from_tok, route, &payload)
.with_context(|| format!("POST {route} on {from}"))?;
if status_code == 200 {
println!(
"{}",
serde_json::to_string_pretty(&json!({
"route": route,
"settled": true,
"note": "settled=true means the remote route returned 200; \
in a dev env with no paid route configured the daemon \
may handle locally",
"response": body,
}))?
);
Ok(())
} else {
let hint = if status_code == 402 {
format!(
"received 402 — no channel or budget / no paid route configured \
between {from} and {to}; \
open a channel first: `node-app harness channel-open {from} {to} 100000`"
)
} else {
format!(
"route returned HTTP {status_code} (may be unconfigured or require \
additional setup)"
)
};
println!(
"{}",
serde_json::to_string_pretty(&json!({
"route": route,
"settled": false,
"http_status": status_code,
"hint": hint,
"body": body,
}))?
);
anyhow::bail!("l402 probe: route {route} did not return 200 (got {status_code})");
}
}
fn load_state() -> Result<HarnessState> {
HarnessState::load().map_err(|_| {
anyhow::anyhow!(
"no harness state found — run `node-app harness up` first"
)
})
}