igneum/pool/src/payout.rs

341 lines
16 KiB
Rust

//! Payouts on the EVM side. The execution layer credits 80% of every blue block's subsidy to the coinbase's `IGNA`
//! address (igneum/exec/src/executor.rs, "Rewards as a state transition applied by rule"), and the 20% proving share
//! to the proving pool escrow by the same rule: the pool's templates name the pool's address, so the pool receives
//! the 80% and never touches the 20%. Members are paid from that address by EIP-1559 transfers through the node's
//! eth_ JSON-RPC, one transaction per payee, at most 16 in flight (the node's mempool refuses a nonce more than 16
//! ahead). `--dry-run` computes and records every payment and sends nothing.
use crate::state::{PaymentRec, State, WEI_PER_IGN};
use alloy_consensus::{SignableTransaction, TxEip1559, TxEnvelope};
use alloy_eips::eip2718::Encodable2718;
use alloy_primitives::{Address, Bytes, Signature, TxKind, B256, U256};
use alloy_signer_local::PrivateKeySigner;
use serde_json::{json, Value};
use std::path::Path;
use std::sync::Mutex;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
/// The payout key file: `{"private_key": "0x..64 hex", "address": "0x..40 hex"}`, created with a fresh random key
/// when missing, mode 0600. The operator backs it up and never commits it (pool/README.md).
pub fn load_or_create_key(path: &Path) -> Result<PrivateKeySigner, String> {
if let Ok(bytes) = std::fs::read(path) {
let v: Value = serde_json::from_slice(&bytes).map_err(|e| format!("{}: {e}", path.display()))?;
let pk = v.get("private_key").and_then(|x| x.as_str()).ok_or("payout key file has no private_key")?;
let raw = hex::decode(pk.trim_start_matches("0x")).map_err(|e| format!("private_key: {e}"))?;
let b = B256::try_from(raw.as_slice()).map_err(|_| "private_key is not 32 bytes".to_string())?;
return PrivateKeySigner::from_bytes(&b).map_err(|e| e.to_string());
}
let signer = PrivateKeySigner::random();
let key = format!("0x{}", hex::encode(signer.to_bytes()));
let body = serde_json::to_vec_pretty(&json!({ "private_key": key, "address": format!("{:?}", signer.address()), "created": "igneum-pool" }))
.map_err(|e| e.to_string())?;
if let Some(dir) = path.parent() {
std::fs::create_dir_all(dir).map_err(|e| e.to_string())?;
}
std::fs::write(path, body).map_err(|e| e.to_string())?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let _ = std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o600));
}
eprintln!("payout key: created {} for address {:?} (back it up; it holds the pool's coinbase income)", path.display(), signer.address());
Ok(signer)
}
/// A minimal JSON-RPC over HTTP/1.1 client for the node's eth_ endpoint: one request per connection.
pub struct EvmRpc {
host: String,
port: u16,
path: String,
id: Mutex<u64>,
}
impl EvmRpc {
pub fn new(url: &str) -> Result<Self, String> {
let rest = url.strip_prefix("http://").ok_or("the EVM RPC URL must start with http://")?;
let (hostport, path) = match rest.find('/') {
Some(i) => (&rest[..i], &rest[i..]),
None => (rest, "/"),
};
let (host, port) = match hostport.rsplit_once(':') {
Some((h, p)) => (h.to_string(), p.parse::<u16>().map_err(|_| "bad port")?),
None => (hostport.to_string(), 80),
};
Ok(Self { host, port, path: path.to_string(), id: Mutex::new(0) })
}
pub async fn call(&self, method: &str, params: Value) -> Result<Value, String> {
let id = {
let mut g = self.id.lock().unwrap();
*g += 1;
*g
};
let body = serde_json::to_vec(&json!({"jsonrpc": "2.0", "id": id, "method": method, "params": params})).unwrap();
let req = format!(
"POST {} HTTP/1.1\r\nHost: {}:{}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
self.path,
self.host,
self.port,
body.len()
);
let fut = async {
let mut s = tokio::net::TcpStream::connect((self.host.as_str(), self.port)).await.map_err(|e| format!("connect {}:{}: {e}", self.host, self.port))?;
s.write_all(req.as_bytes()).await.map_err(|e| e.to_string())?;
s.write_all(&body).await.map_err(|e| e.to_string())?;
let mut buf = Vec::new();
s.read_to_end(&mut buf).await.map_err(|e| e.to_string())?;
Ok::<Vec<u8>, String>(buf)
};
let raw = tokio::time::timeout(std::time::Duration::from_secs(15), fut).await.map_err(|_| "EVM RPC timed out".to_string())??;
let body = http_body(&raw)?;
let v: Value = serde_json::from_slice(&body).map_err(|e| format!("{method}: bad JSON: {e}"))?;
if let Some(err) = v.get("error") {
return Err(format!("{method}: {}", err.get("message").and_then(|m| m.as_str()).unwrap_or(&err.to_string())));
}
Ok(v.get("result").cloned().unwrap_or(Value::Null))
}
}
/// The body of an HTTP/1.1 response: by Content-Length, by chunked transfer encoding, or to the close.
pub fn http_body(raw: &[u8]) -> Result<Vec<u8>, String> {
let sep = raw.windows(4).position(|w| w == b"\r\n\r\n").ok_or("no HTTP header end")?;
let head = String::from_utf8_lossy(&raw[..sep]).to_string();
let rest = &raw[sep + 4..];
let status = head.lines().next().unwrap_or("").to_string();
let mut chunked = false;
let mut len: Option<usize> = None;
for l in head.lines().skip(1) {
let lower = l.to_ascii_lowercase();
if let Some(v) = lower.strip_prefix("content-length:") {
len = v.trim().parse().ok();
}
if lower.starts_with("transfer-encoding:") && lower.contains("chunked") {
chunked = true;
}
}
let body = if chunked {
let mut out = Vec::new();
let mut p = 0;
loop {
let nl = rest[p..].windows(2).position(|w| w == b"\r\n").ok_or("bad chunk")? + p;
let size = usize::from_str_radix(String::from_utf8_lossy(&rest[p..nl]).trim().split(';').next().unwrap_or("0"), 16).map_err(|_| "bad chunk size")?;
if size == 0 {
break;
}
let start = nl + 2;
out.extend_from_slice(rest.get(start..start + size).ok_or("short chunk")?);
p = start + size + 2;
}
out
} else if let Some(n) = len {
rest.get(..n).ok_or("short body")?.to_vec()
} else {
rest.to_vec()
};
if !status.contains(" 200") {
return Err(format!("HTTP {status}: {}", String::from_utf8_lossy(&body).chars().take(200).collect::<String>()));
}
Ok(body)
}
fn hex_u128(v: &Value) -> Option<u128> {
u128::from_str_radix(v.as_str()?.trim_start_matches("0x"), 16).ok()
}
pub struct Payer {
pub signer: PrivateKeySigner,
pub rpc: EvmRpc,
pub chain_id: u64,
pub dry_run: bool,
pub min_payout_wei: u128,
}
impl Payer {
pub fn address_hex(&self) -> String {
format!("{:?}", self.signer.address()).to_lowercase()
}
/// The pool address's EVM balance in wei, or an error text.
pub async fn balance_of(&self, address: &str) -> Result<u128, String> {
let v = self.rpc.call("eth_getBalance", json!([address, "latest"])).await?;
hex_u128(&v).ok_or_else(|| "eth_getBalance: not a quantity".into())
}
/// Signs and sends one transfer; returns the transaction hash.
pub async fn send(&self, to: &str, amount_wei: u128, nonce: u64) -> Result<String, String> {
let from = self.address_hex();
let block = self.rpc.call("eth_getBlockByNumber", json!(["latest", false])).await?;
let base_fee = block.get("baseFeePerGas").and_then(hex_u128).unwrap_or(1_000_000_000);
let tip = self.rpc.call("eth_maxPriorityFeePerGas", json!([])).await.ok().and_then(|v| hex_u128(&v)).unwrap_or(1_000_000_000);
let gas = self
.rpc
.call("eth_estimateGas", json!([{ "from": from, "to": to, "value": format!("{:#x}", amount_wei) }]))
.await
.ok()
.and_then(|v| hex_u128(&v))
.map(|g| (g as u64).max(21_000))
.unwrap_or(30_000);
let to_addr: Address = to.parse().map_err(|_| format!("bad payee address {to}"))?;
let tx = TxEip1559 {
chain_id: self.chain_id,
nonce,
gas_limit: gas,
max_fee_per_gas: base_fee * 2 + tip,
max_priority_fee_per_gas: tip,
to: TxKind::Call(to_addr),
value: U256::from(amount_wei),
access_list: Default::default(),
input: Bytes::new(),
};
let sig: Signature = alloy_signer::SignerSync::sign_hash_sync(&self.signer, &tx.signature_hash()).map_err(|e| e.to_string())?;
let env: TxEnvelope = tx.into_signed(sig).into();
let raw = format!("0x{}", hex::encode(env.encoded_2718()));
let h = self.rpc.call("eth_sendRawTransaction", json!([raw])).await?;
h.as_str().map(|s| s.to_string()).ok_or_else(|| "eth_sendRawTransaction: no hash".into())
}
/// One payout round over the ledger: every balance at or above the minimum, at most 16 transfers.
pub async fn round(&self, state: &Mutex<State>) -> Vec<String> {
let mut lines = Vec::new();
let due: Vec<(String, u128)> = {
let s = state.lock().unwrap();
let mut v: Vec<(String, u128)> = s.balances.iter().filter(|(_, w)| **w >= self.min_payout_wei).map(|(a, w)| (a.clone(), *w)).collect();
v.sort();
v.truncate(16);
v
};
if due.is_empty() {
return lines;
}
if self.dry_run {
let mut s = state.lock().unwrap();
for (a, wei) in due {
s.payments.push(PaymentRec {
at_ms: crate::state::unix_ms(),
address: a.clone(),
amount_wei: wei,
tx_hash: None,
status: "dry-run".into(),
dry_run: true,
nonce: None,
note: "computed, not sent (--dry-run); the balance stays".into(),
});
lines.push(format!("DRY RUN payout {} IGN to {a}", wei as f64 / WEI_PER_IGN as f64));
}
s.dirty = true;
return lines;
}
let from = self.address_hex();
let nonce0 = match self.rpc.call("eth_getTransactionCount", json!([from, "pending"])).await.and_then(|v| hex_u128(&v).ok_or("nonce".into())) {
Ok(n) => n as u64,
Err(e) => {
lines.push(format!("payout: nonce unavailable: {e}"));
return lines;
}
};
let pool_balance = self.balance_of(&from).await.unwrap_or(0);
let mut nonce = nonce0;
let mut spent = 0u128;
for (a, wei) in due {
if spent + wei + 100_000_000_000_000 > pool_balance {
lines.push(format!("payout: pool balance {} IGN cannot cover {} IGN to {a}; waiting", pool_balance as f64 / WEI_PER_IGN as f64, wei as f64 / WEI_PER_IGN as f64));
break;
}
match self.send(&a, wei, nonce).await {
Ok(tx) => {
let mut s = state.lock().unwrap();
let left = s.balances.get(&a).copied().unwrap_or(0).saturating_sub(wei);
s.balances.insert(a.clone(), left);
*s.paid.entry(a.clone()).or_insert(0) += wei;
s.payments.push(PaymentRec { at_ms: crate::state::unix_ms(), address: a.clone(), amount_wei: wei, tx_hash: Some(tx.clone()), status: "sent".into(), dry_run: false, nonce: Some(nonce), note: String::new() });
s.dirty = true;
lines.push(format!("PAYOUT {} IGN to {a} tx {tx} nonce {nonce}", wei as f64 / WEI_PER_IGN as f64));
nonce += 1;
spent += wei;
}
Err(e) => {
lines.push(format!("payout to {a} failed: {e}"));
break;
}
}
}
lines
}
/// Receipts for sent payments: confirmed, failed (balance re-credited), or still pending.
pub async fn receipts(&self, state: &Mutex<State>) -> Vec<String> {
let mut lines = Vec::new();
let sent: Vec<(usize, String, String, u128)> = {
let s = state.lock().unwrap();
s.payments.iter().enumerate().filter(|(_, p)| p.status == "sent").filter_map(|(i, p)| p.tx_hash.clone().map(|h| (i, h, p.address.clone(), p.amount_wei))).collect()
};
for (i, h, a, wei) in sent {
let r = match self.rpc.call("eth_getTransactionReceipt", json!([h])).await {
Ok(r) => r,
Err(e) => {
lines.push(format!("receipt {h}: {e}"));
continue;
}
};
if r.is_null() {
continue;
}
let ok = r.get("status").and_then(|s| s.as_str()).map(|s| s == "0x1").unwrap_or(false);
let mut s = state.lock().unwrap();
if let Some(p) = s.payments.get_mut(i) {
p.status = if ok { "confirmed".into() } else { "failed".into() };
}
if !ok {
*s.balances.entry(a.clone()).or_insert(0) += wei;
let paid = s.paid.get(&a).copied().unwrap_or(0).saturating_sub(wei);
s.paid.insert(a.clone(), paid);
lines.push(format!("payment {h} to {a} FAILED on chain; {} IGN re-credited", wei as f64 / WEI_PER_IGN as f64));
} else {
lines.push(format!("payment {h} to {a} confirmed ({} IGN)", wei as f64 / WEI_PER_IGN as f64));
}
s.dirty = true;
}
lines
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn http_bodies_by_length_and_by_chunks() {
let r = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhello".to_vec();
assert_eq!(http_body(&r).unwrap(), b"hello");
let r = b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n3\r\nabc\r\n2\r\nde\r\n0\r\n\r\n".to_vec();
assert_eq!(http_body(&r).unwrap(), b"abcde");
let r = b"HTTP/1.1 500 Internal\r\nContent-Length: 3\r\n\r\nbad".to_vec();
assert!(http_body(&r).unwrap_err().contains("500"));
}
#[test]
fn the_key_file_round_trips_and_names_its_address() {
let dir = std::env::temp_dir().join(format!("igneum-pool-key-{}", std::process::id()));
let path = dir.join("payout-key.json");
let a = load_or_create_key(&path).unwrap();
let b = load_or_create_key(&path).unwrap();
assert_eq!(a.address(), b.address());
let v: Value = serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap();
assert_eq!(v["address"].as_str().unwrap().to_lowercase(), format!("{:?}", a.address()).to_lowercase());
let _ = std::fs::remove_dir_all(dir);
}
#[test]
fn a_signed_transfer_decodes_to_its_fields() {
let signer = PrivateKeySigner::from_bytes(&B256::from([7u8; 32])).unwrap();
let tx = TxEip1559 { chain_id: 4463, nonce: 3, gas_limit: 21_000, max_fee_per_gas: 3_000_000_000, max_priority_fee_per_gas: 1_000_000_000, to: TxKind::Call(Address::from([1u8; 20])), value: U256::from(5u64), access_list: Default::default(), input: Bytes::new() };
let sig: Signature = alloy_signer::SignerSync::sign_hash_sync(&signer, &tx.signature_hash()).unwrap();
let env: TxEnvelope = tx.into_signed(sig).into();
let raw = env.encoded_2718();
assert_eq!(raw[0], 2, "EIP-1559 type byte");
let back = alloy_eips::eip2718::Decodable2718::decode_2718(&mut raw.as_slice()).unwrap();
let TxEnvelope::Eip1559(s) = back else { panic!("not 1559") };
assert_eq!(s.tx().nonce, 3);
assert_eq!(s.recover_signer().unwrap(), signer.address());
}
}