use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use crate::gpu::ScanMiner;
use aruminium::Gpu;
use erga_pool::stratum::{Job, PoolEvent, Stratum};
pub const DONATION_ADDRESS: &str = "9f8DEbXprAnTS4yhPjp9BEgqnzThVzBPggw5184RureWaRcoGYM";
pub const DONATION_EVERY_NTH: u64 = 20;
fn donation() -> Option<(String, u64)> {
let addr = match std::env::var("ERGA_DONATION") {
Ok(v) if v.eq_ignore_ascii_case("off") || v.is_empty() => return None,
Ok(v) => v,
Err(_) => DONATION_ADDRESS.to_string(),
};
let every = std::env::var("ERGA_DONATION_EVERY")
.ok()
.and_then(|v| v.parse::<u64>().ok())
.unwrap_or(DONATION_EVERY_NTH);
if every == 0 || addr.is_empty() {
return None;
}
Some((addr, every))
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum Intensity {
Max,
Eco,
Min,
}
impl Intensity {
fn duty(self) -> f64 {
match self {
Intensity::Max => 1.0,
Intensity::Eco => 0.25,
Intensity::Min => 0.10,
}
}
fn batch(self) -> u32 {
let base = match self {
Intensity::Max => 8_388_608,
Intensity::Eco => 4_194_304,
Intensity::Min => 2_097_152,
};
std::env::var("ERGA_BATCH")
.ok()
.and_then(|v| v.parse::<u32>().ok())
.filter(|v| *v >= 65_536)
.unwrap_or(base)
}
pub fn parse(s: &str) -> Intensity {
match s.trim().to_ascii_lowercase().as_str() {
"min" => Intensity::Min,
"eco" => Intensity::Eco,
_ => Intensity::Max,
}
}
pub fn as_str(self) -> &'static str {
match self {
Intensity::Max => "max",
Intensity::Eco => "eco",
Intensity::Min => "min",
}
}
}
fn intensity_path() -> Option<std::path::PathBuf> {
let home = std::env::var_os("HOME")?;
Some(
std::path::PathBuf::from(home)
.join("Library/Application Support/ai.cyber.erga")
.join("intensity"),
)
}
fn current_intensity(last: &mut std::time::Instant, cached: &mut Intensity) -> Intensity {
if last.elapsed().as_millis() >= 500 {
*last = std::time::Instant::now();
if let Some(s) = intensity_path().and_then(|p| std::fs::read_to_string(p).ok()) {
*cached = Intensity::parse(&s);
}
}
*cached
}
fn cadence_path() -> Option<std::path::PathBuf> {
let home = std::env::var_os("HOME")?;
Some(
std::path::PathBuf::from(home)
.join("Library/Application Support/ai.cyber.erga")
.join("cadence"),
)
}
fn load_owed(every: u64) -> u64 {
cadence_path()
.and_then(|p| std::fs::read_to_string(p).ok())
.and_then(|s| s.trim().parse::<u64>().ok())
.unwrap_or(every - 1)
.min(every - 1)
}
fn save_owed(owed: u64) {
let Some(p) = cadence_path() else { return };
if let Some(dir) = p.parent() {
let _ = std::fs::create_dir_all(dir);
}
let _ = std::fs::write(p, owed.to_string());
}
pub const NO_PREFETCH: u64 = 200;
struct SendMiner(ScanMiner);
unsafe impl Send for SendMiner {}
#[derive(Default)]
struct Tables {
current: Option<(u32, ScanMiner)>,
next: Option<Prefetch>,
}
struct Prefetch {
height: u32,
rx: std::sync::mpsc::Receiver<Result<SendMiner, String>>,
urgent: Arc<AtomicBool>,
cancel: Arc<AtomicBool>,
}
impl Drop for Prefetch {
fn drop(&mut self) {
self.cancel.store(true, Ordering::Relaxed);
}
}
fn prefetch_allowed(intensity: Intensity, available: u64, need: u64) -> bool {
const GIB: u64 = 1 << 30;
let headroom = match intensity {
Intensity::Max => 3 * GIB,
Intensity::Eco => 8 * GIB,
Intensity::Min => return false,
};
available > need.saturating_add(headroom)
}
fn available_memory() -> u64 {
let mut sys = sysinfo::System::new();
sys.refresh_memory();
sys.available_memory()
}
fn start_prefetch(
height: u32,
version: u8,
m: &[u8],
p: &Arc<Progress>,
intensity: Intensity,
) -> Option<Prefetch> {
let n = autolykos::calc_big_n(version, height);
let need = n as u64 * 32;
if !prefetch_allowed(intensity, available_memory(), need) {
return None;
}
let (tx, rx) = std::sync::mpsc::channel();
let urgent = Arc::new(AtomicBool::new(false));
let cancel = Arc::new(AtomicBool::new(false));
let (u, c, pp, mm) = (urgent.clone(), cancel.clone(), p.clone(), m.to_vec());
std::thread::spawn(move || {
let result = Gpu::open()
.map_err(|e| format!("{e:?}"))
.and_then(|gpu| {
ScanMiner::new_gpu_built(gpu, n, height, &mm, &|f| {
pp.next_pct.store((f * 100.0) as u64, Ordering::Relaxed);
if c.load(Ordering::Relaxed) {
return false;
}
if f < 1.0 && !u.load(Ordering::Relaxed) {
for _ in 0..30 {
if u.load(Ordering::Relaxed) || c.load(Ordering::Relaxed) {
break;
}
std::thread::sleep(std::time::Duration::from_millis(200));
}
}
!c.load(Ordering::Relaxed)
})
})
.map(SendMiner);
if result.is_err() {
pp.next_pct.store(NO_PREFETCH, Ordering::Relaxed);
}
let _ = tx.send(result);
});
Some(Prefetch { height, rx, urgent, cancel })
}
pub struct Progress {
pub running: AtomicBool,
pub stop: AtomicBool,
pub rate_khs: AtomicU64, pub accepted: AtomicU64,
pub rejected: AtomicU64,
pub height: AtomicU64,
pub hashed: AtomicU64,
pub submitted: AtomicU64, pub donated: AtomicU64, pub build_pct: AtomicU64,
pub next_pct: AtomicU64,
pub device: Mutex<String>,
pub status: Mutex<String>,
}
impl Progress {
pub fn new() -> Arc<Self> {
Arc::new(Progress {
running: AtomicBool::new(false),
stop: AtomicBool::new(false),
rate_khs: AtomicU64::new(0),
accepted: AtomicU64::new(0),
rejected: AtomicU64::new(0),
height: AtomicU64::new(0),
hashed: AtomicU64::new(0),
submitted: AtomicU64::new(0),
donated: AtomicU64::new(0),
build_pct: AtomicU64::new(0),
next_pct: AtomicU64::new(NO_PREFETCH),
device: Mutex::new(String::new()),
status: Mutex::new("idle".into()),
})
}
pub fn mhs(&self) -> f64 {
self.rate_khs.load(Ordering::Relaxed) as f64 / 1000.0
}
pub fn set_status(&self, s: impl Into<String>) {
*self.status.lock().unwrap() = s.into();
}
}
pub struct PoolCfg {
pub host: String,
pub port: u16,
pub address: String,
}
pub fn run(cfg: PoolCfg, p: Arc<Progress>) {
p.running.store(true, Ordering::Relaxed);
p.accepted.store(0, Ordering::Relaxed);
p.rejected.store(0, Ordering::Relaxed);
p.hashed.store(0, Ordering::Relaxed);
p.submitted.store(0, Ordering::Relaxed);
p.donated.store(0, Ordering::Relaxed);
match Gpu::open() {
Ok(g) => *p.device.lock().unwrap() = g.name(),
Err(e) => {
p.set_status(format!("no Metal GPU: {e:?}"));
p.running.store(false, Ordering::Relaxed);
return;
}
}
let m = autolykos::big_m();
let don = donation();
let mut tables = Tables::default();
let mut owed_to_you: u64 = don.as_ref().map(|(_, every)| load_owed(*every)).unwrap_or(u64::MAX);
let mut donating = don.is_some() && owed_to_you == 0;
while !p.stop.load(Ordering::Relaxed) {
let (addr, mut quota) = match (&don, donating) {
(Some((a, _)), true) => (a.clone(), 1u64),
_ => (cfg.address.clone(), owed_to_you.max(1)),
};
p.set_status(if donating { "connectingโฆ (development share)" } else { "connectingโฆ" });
match Stratum::connect(&cfg.host, cfg.port, &addr, "erga") {
Ok(s) => {
let end = mine_session(s, &addr, &p, &m, &mut tables, &mut quota, donating);
match end {
SessionEnd::Stopped => break,
SessionEnd::QuotaMet => {
if donating {
donating = false;
owed_to_you = don.as_ref().map(|(_, e)| e - 1).unwrap_or(u64::MAX);
} else if don.is_some() {
donating = true;
owed_to_you = 0;
}
if don.is_some() {
save_owed(owed_to_you);
}
continue; }
SessionEnd::Closed => {
if !donating {
owed_to_you = quota;
if don.is_some() {
save_owed(owed_to_you);
}
}
p.set_status("pool disconnected โ reconnectingโฆ");
}
}
}
Err(e) => p.set_status(format!("connect failed ({e}) โ retryingโฆ")),
}
if p.stop.load(Ordering::Relaxed) {
break;
}
p.rate_khs.store(0, Ordering::Relaxed);
for _ in 0..30 {
if p.stop.load(Ordering::Relaxed) {
break;
}
std::thread::sleep(std::time::Duration::from_millis(100));
}
}
p.set_status("idle");
p.rate_khs.store(0, Ordering::Relaxed);
p.running.store(false, Ordering::Relaxed);
}
enum SessionEnd {
QuotaMet,
Closed,
Stopped,
}
fn mine_session(
mut s: Stratum,
address: &str,
p: &Arc<Progress>,
m: &[u8],
tables: &mut Tables,
quota: &mut u64,
donating: bool,
) -> SessionEnd {
let en1 = s.extranonce1.clone();
let en1_bits = en1.len() as u32 * 8;
let search_bits = 64 - en1_bits;
let en1_val: u64 = en1.iter().fold(0u64, |a, &b| (a << 8) | b as u64);
let en1_prefix = if en1_bits == 0 { 0 } else { en1_val << search_bits };
let tail_mask: u64 = if search_bits >= 64 { u64::MAX } else { (1u64 << search_bits) - 1 };
let mut job: Option<Job> = None;
let mut cursor: u64 = 0;
let mut intensity = Intensity::parse(&std::env::var("ERGA_INTENSITY").unwrap_or_default());
let mut intensity_checked = std::time::Instant::now();
let mut window_start = std::time::Instant::now();
let mut window_hashed: u64 = 0;
p.set_status("waiting for workโฆ");
while !p.stop.load(Ordering::Relaxed) {
let mut latest_job: Option<Job> = None;
while let Ok(ev) = s.events.try_recv() {
match ev {
PoolEvent::Job(j) => latest_job = Some(j),
PoolEvent::Difficulty(_) => {}
PoolEvent::SubmitResult { accepted, .. } => {
if accepted {
p.accepted.fetch_add(1, Ordering::Relaxed);
} else {
p.rejected.fetch_add(1, Ordering::Relaxed);
}
}
PoolEvent::Closed => return SessionEnd::Closed,
}
}
if let Some(j) = latest_job {
let (table, prefetch) = (&mut tables.current, &mut tables.next);
if table.as_ref().map(|(h, _)| *h) != Some(j.height) {
p.height.store(j.height as u64, Ordering::Relaxed);
let mut swapped = false;
if prefetch.as_ref().is_some_and(|pf| pf.height == j.height) {
let pf = prefetch.take().unwrap();
pf.urgent.store(true, Ordering::Relaxed);
p.set_status("building tableโฆ");
loop {
match pf.rx.try_recv() {
Ok(Ok(mn)) => {
*table = Some((j.height, mn.0));
swapped = true;
break;
}
Ok(Err(_)) | Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
Err(std::sync::mpsc::TryRecvError::Empty) => {
let f = p.next_pct.load(Ordering::Relaxed).min(100);
p.build_pct.store(f, Ordering::Relaxed);
std::thread::sleep(std::time::Duration::from_millis(100));
}
}
}
} else {
*prefetch = None;
}
p.next_pct.store(NO_PREFETCH, Ordering::Relaxed);
if !swapped {
let n = autolykos::calc_big_n(j.version, j.height);
p.set_status("building tableโฆ");
p.rate_khs.store(0, Ordering::Relaxed);
let gpu = match Gpu::open() {
Ok(g) => g,
Err(e) => {
p.set_status(format!("GPU open failed: {e:?}"));
return SessionEnd::Closed;
}
};
*table = None;
match ScanMiner::new_gpu_built(gpu, n, j.height, m, &|f| {
p.build_pct.store((f * 100.0) as u64, Ordering::Relaxed);
true
}) {
Ok(mn) => *table = Some((j.height, mn)),
Err(e) => {
p.set_status(format!("table build retry: {e}"));
*table = None;
continue;
}
}
}
*prefetch =
start_prefetch(j.height + 1, j.version, m, p, intensity);
cursor = 0;
window_start = std::time::Instant::now();
window_hashed = 0;
}
p.height.store(j.height as u64, Ordering::Relaxed);
p.set_status(if donating { "mining ยท development share" } else { "mining" });
job = Some(j);
}
let (Some((_, mn)), Some(j)) = (&tables.current, &job) else {
std::thread::sleep(std::time::Duration::from_millis(80));
continue;
};
let mut msg = [0u8; 32];
if j.msg.len() == 32 {
msg.copy_from_slice(&j.msg);
}
let target = left_pad_32(&j.target_b.to_bytes_be());
intensity = current_intensity(&mut intensity_checked, &mut intensity);
let batch = intensity.batch();
let nonce_base = en1_prefix | (cursor & tail_mask);
let dispatch_started = std::time::Instant::now();
if let Some(nonce) = mn.scan(&msg, &target, nonce_base, batch) {
let nb = nonce.to_be_bytes();
let hit = autolykos::pow_hit(&msg, &nb, &j.height.to_be_bytes(), mn.n, m);
if hit < j.target_b {
let nonce_hex = hex(&nb);
let en2_hex = hex(&nb[en1.len()..]);
let _ = s.submit(address, "erga", &j.job_id, &en2_hex, &nonce_hex, "");
p.submitted.fetch_add(1, Ordering::Relaxed);
if donating {
p.donated.fetch_add(1, Ordering::Relaxed);
}
*quota = quota.saturating_sub(1);
if !donating && *quota != u64::MAX {
save_owed(*quota);
}
if *quota == 0 {
std::thread::sleep(std::time::Duration::from_millis(400));
drain_results(&mut s, p);
return SessionEnd::QuotaMet;
}
}
}
cursor = cursor.wrapping_add(batch as u64);
p.hashed.fetch_add(batch as u64, Ordering::Relaxed);
window_hashed += batch as u64;
let duty = intensity.duty();
if duty < 1.0 {
let worked = dispatch_started.elapsed();
let rest = worked.mul_f64((1.0 / duty) - 1.0);
std::thread::sleep(rest.min(std::time::Duration::from_secs(2)));
}
let dt = window_start.elapsed().as_secs_f64();
if dt > 0.5 {
let mhs = window_hashed as f64 / dt / 1e6;
p.rate_khs.store((mhs * 1000.0) as u64, Ordering::Relaxed);
window_start = std::time::Instant::now();
window_hashed = 0;
}
}
SessionEnd::Stopped
}
fn drain_results(s: &mut Stratum, p: &Arc<Progress>) {
while let Ok(ev) = s.events.try_recv() {
if let PoolEvent::SubmitResult { accepted, .. } = ev {
if accepted {
p.accepted.fetch_add(1, Ordering::Relaxed);
} else {
p.rejected.fetch_add(1, Ordering::Relaxed);
}
}
}
}
fn left_pad_32(b: &[u8]) -> [u8; 32] {
let mut out = [0u8; 32];
let n = b.len().min(32);
out[32 - n..].copy_from_slice(&b[b.len() - n..]);
out
}
fn hex(b: &[u8]) -> String {
b.iter().map(|x| format!("{x:02x}")).collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn intensity_eases_off_monotonically() {
let (max, eco, min) = (Intensity::Max, Intensity::Eco, Intensity::Min);
assert!(max.duty() > eco.duty() && eco.duty() > min.duty());
assert!(max.batch() >= eco.batch() && eco.batch() >= min.batch());
assert_eq!(max.duty(), 1.0, "max must not throttle itself");
}
#[test]
fn intensity_parses_and_round_trips() {
for i in [Intensity::Max, Intensity::Eco, Intensity::Min] {
assert_eq!(Intensity::parse(i.as_str()), i);
}
assert_eq!(Intensity::parse(" ECO \n"), Intensity::Eco);
for junk in ["", "fast", "0", "maximum"] {
assert_eq!(Intensity::parse(junk), Intensity::Max, "{junk:?}");
}
}
}