Compare commits
4 commits
master
...
pool-v0-re
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1418aaf110 | ||
|
|
db22de6793 | ||
|
|
2c0bd55fa7 | ||
|
|
1797ebadfa |
9 changed files with 331 additions and 44 deletions
|
|
@ -236,3 +236,93 @@ a wrong-size cache and mismatched silently); a pool on a class v4 network needs
|
||||||
Building the pool crate on igneum-build-1: the pool reads the fork through the `vendor/igneum-node` symlink; lib.sh syncs
|
Building the pool crate on igneum-build-1: the pool reads the fork through the `vendor/igneum-node` symlink; lib.sh syncs
|
||||||
the fork worktree it points at (`vendor/igneum-node-pr`) as a whole repository, and the symlink itself is created once on
|
the fork worktree it points at (`vendor/igneum-node-pr`) as a whole repository, and the symlink itself is created once on
|
||||||
the box by hand (`ln -s igneum-node-pr /srv/builds/<worktree>/vendor/igneum-node`), the one step the scripts do not do.
|
the box by hand (`ln -s igneum-node-pr /srv/builds/<worktree>/vendor/igneum-node`), the one step the scripts do not do.
|
||||||
|
|
||||||
|
### 9.1 7 October 2026: the daemon refuses to start half-alive
|
||||||
|
|
||||||
|
The fleet agent's night (6 October, 23:01Z to 23:32Z): the rebased pair never ran as a test. The daemon was started while
|
||||||
|
the node still held its member port, then four members (the wind-down's weight rule freed four pods, not ten) sat
|
||||||
|
connected with no job while every template fetch timed out; nothing in the daemon's own log said so except one line per
|
||||||
|
member. The row from it, now code (pool `269364a`+): the member and API listeners are bound in `main` before anything
|
||||||
|
else and a port that cannot be bound ends the start with exit code 2 and `POOL NOT STARTED: cannot bind the members
|
||||||
|
listener on <addr>: <os error>`; the node must answer a template with its `pow_epoch` before the pool serves anyone,
|
||||||
|
retried for `--node-wait-secs` (default 60, each attempt logged), else exit code 3; a `STATUS` line every 30 s
|
||||||
|
(`members= jobs_issued= shares_accepted= shares_rejected= blocks= template_failures=; node ok | NODE NOT ANSWERING |
|
||||||
|
NO TEMPLATE YET`). Test: `server::bind_tests` (a held port refused with the address in the message, the freed port binds).
|
||||||
|
The 10-member run is still owed: a clean 30-minute window, the daemon bound and its `node ... answers templates` line
|
||||||
|
printed before the first member, on the next pods the weight rule frees (the table must read over 75 percent signed first).
|
||||||
|
|
||||||
|
### 9.2 7 October 2026, 06:00Z to 06:33Z: the re-run (five members), what it settled and what it found
|
||||||
|
|
||||||
|
Settled: the class question. pm-1 and pm-2 (RTX 3070, grpc url `none`) at +7 min: 36,917 and 34,323 shares accepted, 0
|
||||||
|
rejected, 0 WORKER MISMATCH, 0 refused, on the class v3 program of the pool's node's epoch (`5530a50d...`, era = the devnet
|
||||||
|
genesis hash), the genesis day (20729) and dataset size (2^28) installed from the daemon's `seeds` line. The 6 October
|
||||||
|
failure does not reproduce.
|
||||||
|
|
||||||
|
Found: the daemon stalls on a real chain. From 06:01Z (the first block found) no job was issued for seven minutes, pm-3
|
||||||
|
and pm-4 never got one, every block found on the stale templates (55, then 105, 480 DAA behind the virtual) was an
|
||||||
|
orphan, vardiff could not act (a new share target rides with a job) so the two hashing members sent about 90 shares a
|
||||||
|
second each and the share check read 9.8 ms on pool-1's cores (1.35 ms on the Mac); after a swap to a 20 s template
|
||||||
|
timeout (f808b3f3) with the data dir kept, the stall returned within a minute. The node built templates in 0.6 ms all
|
||||||
|
the while (its prewarm line) and served its own solo miner.
|
||||||
|
|
||||||
|
The cause, from the daemon's code: `confirm_loop` restarted its chain walk from the PRUNING POINT on any failed
|
||||||
|
`getVirtualChainFromBlock`, and then called `get_block` for every chain block since (about 120,000 on the devnet) over
|
||||||
|
the ONE gRPC connection the templates used (the client is one request stream per connection, 5 s request timeout); one
|
||||||
|
timed-out request under four parallel template fetches started the walk, every later request waited behind it and
|
||||||
|
timed out, which restarted the walk again. The kept state file (105 pending blocks) restarted it after the swap. The
|
||||||
|
private measurement run (section 5) never saw it: a 600-block chain walks in a moment.
|
||||||
|
|
||||||
|
Fixed (pool commit after bd49c2a9): the walk and the network numbers run on a second gRPC connection (`pool.walker`); a
|
||||||
|
failed chain call keeps its cursor (moved to the sink only when the node no longer knows it, never to the pruning
|
||||||
|
point; test `node::walk_tests`); a tick walks at most 600 chain blocks and says how many; a start with pending blocks
|
||||||
|
walks from the sink and says that older ones resolve by the orphan rule. `--template-timeout-s` stays (default 20).
|
||||||
|
|
||||||
|
Still open from the night, for the next window: the node log showed a new gRPC connection about every 0.8 s with the
|
||||||
|
count steady at 5 (something opens and closes one each time; the daemon's re-subscribe loop is the suspect, the
|
||||||
|
subscribe-line count in the daemon log decides it); vardiff's new target should apply to the member's current work
|
||||||
|
without waiting for a job (a one-line member change, needs a miner rebuild); the 9.8 ms share check on the pool box's
|
||||||
|
cores against 1.35 ms on the Mac (the verifier's day cache is 1 GiB and random reads on a rented box's memory are the
|
||||||
|
cost; two verify threads saturate at about 200 shares a second, which vardiff must keep far away); and the 10-member
|
||||||
|
load figure itself, not taken.
|
||||||
|
|
||||||
|
Consequences per tier: a pool operator needs a box whose memory serves the 1 GiB cache at speed (a rented 2 vCPU box
|
||||||
|
verifies about 100 shares a second per core at 9.8 ms; at one share per 10 s per member that is 1,000 members per
|
||||||
|
core, so the verifier is not the limit once vardiff holds); a member on any card is unaffected by any of this; the
|
||||||
|
chain walk now costs the node at most 600 `get_block` calls per 5 s on its own connection.
|
||||||
|
|
||||||
|
### 9.3 The re-run's close (06:33:33Z) and the member rows from its logs
|
||||||
|
|
||||||
|
After 2f6c0358 (06:20:34Z to 06:33:33Z, 13 minutes, four members): 207 blocks found, 144 confirmed, 69 orphaned (13
|
||||||
|
inherited pending at the swap, 56 on its own jobs: a 27 percent residual at 1 bps), 20,404 shares accepted, 0 rejected,
|
||||||
|
0 mismatches, 0 refused, every STATUS `node ok`, the gRPC connection churn gone (ids #3 to #5 once each in a minute; the
|
||||||
|
193 connections in 150 s were under f808b3f3's pruning-point walk), re-subscribe lines 2 (one per daemon start: the
|
||||||
|
re-subscribe loop was never the churn). Before it, f808b3f3 had worked through the walk by about 06:14Z: 245 blocks,
|
||||||
|
69 confirmed, 163 orphaned. First CONFIRMED line 06:24:32Z, 4.80 IGN to the finder. Logs: `~/Desktop/fleet/pool-run-0607/`.
|
||||||
|
|
||||||
|
The residual orphans are template latency: the daemon's STATUS reads `last template 1 s ago` and the solo miner on the
|
||||||
|
same node reads `template_ms=1632`; pool-1's node answers a template request in 1 to 1.6 s although its mining manager
|
||||||
|
answers a per-member request from its cache by a coinbase rewrite (`modify_block_template`) and its prewarm builds in
|
||||||
|
0.6 ms. At 1 bps a template that arrives 1.6 s old loses about a quarter of the blocks found on it, pool or solo. That
|
||||||
|
latency is the node's RPC path, a row for the node lane (measure `get_block_template` round trips on a devnet node under
|
||||||
|
a miner and a pool; the suspect is the RPC service's `pow_epoch` derivation per request). 82 "RPC request timeout" lines
|
||||||
|
after the 2f6c0358 swap are the gRPC client's own 5 s request timeout on that path; `--template-timeout-s` cannot
|
||||||
|
lengthen it.
|
||||||
|
|
||||||
|
From pm-1's log (1.09 GB): 4,342,515 lines are `worker: error N epoch seed mismatch: this worker holds epoch eea66ce7...`,
|
||||||
|
one per re-queued job, for the 40 s the worker compiled the new epoch's pack (NVRTC 40,327 ms) and again after every
|
||||||
|
reconnect, because a session respawned the worker (a second process beside the first, the same compile again) and
|
||||||
|
the member fed jobs for a pair the worker did not hold yet. Vardiff from the same log: shift 10 to 11 (idle easing
|
||||||
|
while the worker compiled, which consumed the sized first correction), then 11 down to 0 one step per 30 s over 5.5
|
||||||
|
minutes at about 190 shares a second, the load that put the share check at 10 to 11 ms on pool-1's cores.
|
||||||
|
|
||||||
|
Fixed in the member (fork commit after b8070476): one worker process per run, jobs held for a pair whose prepare is
|
||||||
|
pending and never re-prepared once held (`WorkerMemory`, test `jobs_are_held_while_a_pair_is_being_prepared`), worker
|
||||||
|
error lines summarised (one per class, then a count per 1,000), `set_target` applied to the current work at once.
|
||||||
|
Fixed in the pool (`vardiff.rs`): the idle easing no longer consumes the sized first correction (test
|
||||||
|
`an_idle_easing_does_not_consume_the_sized_first_correction`).
|
||||||
|
|
||||||
|
Consequences per tier: a member on a 3070-class card loses 40 s of hash at every epoch roll to the NVRTC compile unless
|
||||||
|
its pack is prepared ahead (the pool's `seeds` line carries `next_epoch_seed` for that, the member's prepare-ahead is
|
||||||
|
the next member row); a member on a reconnect now keeps its worker and its compiled pair; a pool operator's verifier
|
||||||
|
sees the sized correction within 10 s of a member's first shares (one share per 10 s per member from then on) instead of
|
||||||
|
5 minutes of a 190-share-a-second flood; the 10-member load figure is still owed and now has a clean form to run in.
|
||||||
|
|
|
||||||
|
|
@ -170,3 +170,7 @@ Every Linux build and suite runs on the box through `tools/build-remote.sh` from
|
||||||
crate reads the fork through the `vendor/igneum-node` symlink; the build library syncs the fork worktree it points at as a
|
crate reads the fork through the `vendor/igneum-node` symlink; the build library syncs the fork worktree it points at as a
|
||||||
whole repository, and the symlink is made once on the box by hand: `ln -s <fork worktree name> /srv/builds/<worktree>/vendor/igneum-node`.
|
whole repository, and the symlink is made once on the box by hand: `ln -s <fork worktree name> /srv/builds/<worktree>/vendor/igneum-node`.
|
||||||
Artefacts land in `pool/target-remote/release/igneum-pool` (x86_64 Linux); the Mac builds no Linux binary.
|
Artefacts land in `pool/target-remote/release/igneum-pool` (x86_64 Linux); the Mac builds no Linux binary.
|
||||||
|
|
||||||
|
Start-up contract (7 October 2026): the daemon exits 2 when `--listen` or `--http` cannot be bound, exits 3 when the node
|
||||||
|
answers no template with `pow_epoch` within `--node-wait-secs` (default 60), and prints `node <url> answers templates: ...`
|
||||||
|
before it accepts a member; a `STATUS` line every 30 s names the member count, jobs issued, shares and the node's state.
|
||||||
|
|
|
||||||
|
|
@ -199,11 +199,7 @@ pub fn route(pool: &Pool, path: &str) -> Vec<u8> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn serve(pool: Arc<Pool>) {
|
pub async fn serve(pool: Arc<Pool>, listener: tokio::net::TcpListener) {
|
||||||
let listener = tokio::net::TcpListener::bind(&pool.cfg.http).await.unwrap_or_else(|e| {
|
|
||||||
eprintln!("cannot listen on {}: {e}", pool.cfg.http);
|
|
||||||
std::process::exit(1)
|
|
||||||
});
|
|
||||||
println!("{} pool: stats API and page on http://{}/", crate::state::unix_ms(), pool.cfg.http);
|
println!("{} pool: stats API and page on http://{}/", crate::state::unix_ms(), pool.cfg.http);
|
||||||
loop {
|
loop {
|
||||||
let Ok((sock, _)) = listener.accept().await else { continue };
|
let Ok((sock, _)) = listener.accept().await else { continue };
|
||||||
|
|
|
||||||
|
|
@ -33,12 +33,18 @@ pub struct Config {
|
||||||
/// Shares sampled at shifts above this are still verified in v0 (spec sample_shift 8 is an allowance, not a duty)
|
/// Shares sampled at shifts above this are still verified in v0 (spec sample_shift 8 is an allowance, not a duty)
|
||||||
pub verify_threads: usize,
|
pub verify_threads: usize,
|
||||||
pub snapshot_interval_s: u64,
|
pub snapshot_interval_s: u64,
|
||||||
|
/// How long the start-up waits for the node to answer a template before the daemon refuses to start (exit 3)
|
||||||
|
pub node_wait_secs: u64,
|
||||||
|
/// Per-member template fetch timeout. 7 October 2026, 06:08Z: pool-1's node built a template in 1.6 to 1.8 s under
|
||||||
|
/// its own miner plus four members, the fixed 5 s gave up on the parallel fetches, no job was issued for 7 minutes,
|
||||||
|
/// every block the stale templates found was an orphan, and vardiff could not act (a new target rides with a job)
|
||||||
|
pub template_timeout_s: u64,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Config {
|
impl Config {
|
||||||
pub fn usage() -> ! {
|
pub fn usage() -> ! {
|
||||||
eprintln!(
|
eprintln!(
|
||||||
"usage: igneum-pool [--node grpc://127.0.0.1:26610]... [--evm-rpc http://127.0.0.1:26790] [--listen 0.0.0.0:4463]\n\
|
"usage: igneum-pool [--node grpc://127.0.0.1:26610]... [--evm-rpc http://127.0.0.1:26790] [--listen 0.0.0.0:4463] [--node-wait-secs 60] [--template-timeout-s 20]\n\
|
||||||
\x20 [--http 127.0.0.1:4480] [--data-dir ./data] [--payout-key <file>] [--fee-percent 1] [--min-payout 1.0]\n\
|
\x20 [--http 127.0.0.1:4480] [--data-dir ./data] [--payout-key <file>] [--fee-percent 1] [--min-payout 1.0]\n\
|
||||||
\x20 [--pplns-window 2.0] [--payout-interval-s 60] [--dry-run] [--network devnet|testnet|mainnet]\n\
|
\x20 [--pplns-window 2.0] [--payout-interval-s 60] [--dry-run] [--network devnet|testnet|mainnet]\n\
|
||||||
\x20 [--share-interval-s 10] [--min-shift 0] [--max-shift 60] [--stale-grace-ms 2000] [--orphan-after-daa 120]\n\
|
\x20 [--share-interval-s 10] [--min-shift 0] [--max-shift 60] [--stale-grace-ms 2000] [--orphan-after-daa 120]\n\
|
||||||
|
|
@ -70,6 +76,8 @@ impl Config {
|
||||||
public_url: String::new(),
|
public_url: String::new(),
|
||||||
verify_threads: 2,
|
verify_threads: 2,
|
||||||
snapshot_interval_s: 15,
|
snapshot_interval_s: 15,
|
||||||
|
node_wait_secs: 60,
|
||||||
|
template_timeout_s: 20,
|
||||||
};
|
};
|
||||||
let mut i = 0;
|
let mut i = 0;
|
||||||
while i < args.len() {
|
while i < args.len() {
|
||||||
|
|
@ -100,6 +108,8 @@ impl Config {
|
||||||
"--public-url" => c.public_url = val(),
|
"--public-url" => c.public_url = val(),
|
||||||
"--verify-threads" => c.verify_threads = val().parse::<usize>().unwrap_or_else(|_| Config::usage()).max(1),
|
"--verify-threads" => c.verify_threads = val().parse::<usize>().unwrap_or_else(|_| Config::usage()).max(1),
|
||||||
"--snapshot-interval-s" => c.snapshot_interval_s = val().parse().unwrap_or_else(|_| Config::usage()),
|
"--snapshot-interval-s" => c.snapshot_interval_s = val().parse().unwrap_or_else(|_| Config::usage()),
|
||||||
|
"--node-wait-secs" => c.node_wait_secs = val().parse().unwrap_or_else(|_| Config::usage()),
|
||||||
|
"--template-timeout-s" => c.template_timeout_s = val().parse::<u64>().unwrap_or_else(|_| Config::usage()).max(1),
|
||||||
"-h" | "--help" => Config::usage(),
|
"-h" | "--help" => Config::usage(),
|
||||||
_ => Config::usage(),
|
_ => Config::usage(),
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -27,6 +27,7 @@ use config::Config;
|
||||||
use kaspa_addresses::{Address, Prefix, Version};
|
use kaspa_addresses::{Address, Prefix, Version};
|
||||||
use kaspa_grpc_client::GrpcClient;
|
use kaspa_grpc_client::GrpcClient;
|
||||||
use kaspa_pow::igneum::IgneumEngine;
|
use kaspa_pow::igneum::IgneumEngine;
|
||||||
|
use kaspa_rpc_core::api::rpc::RpcApi;
|
||||||
use pool::{NetInfo, Pool};
|
use pool::{NetInfo, Pool};
|
||||||
use state::{State, WEI_PER_IGN};
|
use state::{State, WEI_PER_IGN};
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
|
|
@ -38,6 +39,17 @@ use std::time::{Duration, Instant};
|
||||||
async fn main() {
|
async fn main() {
|
||||||
let args: Vec<String> = std::env::args().skip(1).collect();
|
let args: Vec<String> = std::env::args().skip(1).collect();
|
||||||
let cfg = Config::from_args(&args);
|
let cfg = Config::from_args(&args);
|
||||||
|
// The two listeners first (7 October 2026, the fleet agent's row from the 6 October night: a daemon whose member
|
||||||
|
// port another process still held ran 20 minutes as a process with no socket): a port that cannot be bound ends
|
||||||
|
// the start here, exit code 2, before a node is contacted or a member is let down.
|
||||||
|
let member_listener = server::bind_listener("members", &cfg.listen).await.unwrap_or_else(|e| {
|
||||||
|
eprintln!("POOL NOT STARTED: {e}");
|
||||||
|
std::process::exit(2)
|
||||||
|
});
|
||||||
|
let api_listener = server::bind_listener("stats API", &cfg.http).await.unwrap_or_else(|e| {
|
||||||
|
eprintln!("POOL NOT STARTED: {e}");
|
||||||
|
std::process::exit(2)
|
||||||
|
});
|
||||||
std::fs::create_dir_all(&cfg.data_dir).unwrap_or_else(|e| {
|
std::fs::create_dir_all(&cfg.data_dir).unwrap_or_else(|e| {
|
||||||
eprintln!("data dir {}: {e}", cfg.data_dir.display());
|
eprintln!("data dir {}: {e}", cfg.data_dir.display());
|
||||||
std::process::exit(1)
|
std::process::exit(1)
|
||||||
|
|
@ -68,6 +80,44 @@ async fn main() {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
let node = clients.remove(0);
|
let node = clients.remove(0);
|
||||||
|
let walker = match tokio::time::timeout(Duration::from_secs(10), GrpcClient::connect(cfg.nodes[0].clone())).await {
|
||||||
|
Ok(Ok(c)) => Arc::new(c),
|
||||||
|
Ok(Err(e)) => {
|
||||||
|
eprintln!("node {} (second connection): {e}", cfg.nodes[0]);
|
||||||
|
std::process::exit(1)
|
||||||
|
}
|
||||||
|
Err(_) => {
|
||||||
|
eprintln!("node {} (second connection): connect timed out", cfg.nodes[0]);
|
||||||
|
std::process::exit(1)
|
||||||
|
}
|
||||||
|
};
|
||||||
|
// The node must answer a template with its pow_epoch before the pool serves anyone: a pool whose node answers no
|
||||||
|
// template issues no job (the 6 October night: four members connected, "template timed out" per member, no job
|
||||||
|
// for 20 minutes). Retried for --node-wait-secs, saying so; then exit code 3.
|
||||||
|
{
|
||||||
|
let probe_address = Address::new(Prefix::Devnet, Version::PubKey, &[0u8; 32]);
|
||||||
|
let deadline = Instant::now() + Duration::from_secs(cfg.node_wait_secs);
|
||||||
|
let mut attempt = 0u32;
|
||||||
|
loop {
|
||||||
|
attempt += 1;
|
||||||
|
let why = match tokio::time::timeout(Duration::from_secs(10), node.get_block_template(probe_address.clone(), Vec::new())).await {
|
||||||
|
Ok(Ok(t)) if t.pow_epoch.is_some() => {
|
||||||
|
let i = t.pow_epoch.unwrap();
|
||||||
|
println!("{} node {} answers templates: epoch {} class {} era {} daa {}", state::unix_ms(), cfg.nodes[0], i.epoch_index, i.class().name(), i.era_seed.map(|h| h.to_string()).unwrap_or_else(|| "-".into()), i.virtual_daa_score);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
Ok(Ok(_)) => "the template carries no pow_epoch (a node before the 0.3.x line); the pool cannot name a program".to_string(),
|
||||||
|
Ok(Err(e)) => e.to_string(),
|
||||||
|
Err(_) => "template timed out after 10 s".to_string(),
|
||||||
|
};
|
||||||
|
if Instant::now() >= deadline {
|
||||||
|
eprintln!("POOL NOT STARTED: node {} answered no usable template in {} s ({attempt} attempts; last: {why})", cfg.nodes[0], cfg.node_wait_secs);
|
||||||
|
std::process::exit(3)
|
||||||
|
}
|
||||||
|
eprintln!("{} waiting for node {} to answer a template (attempt {attempt}: {why}); giving up at {} s", state::unix_ms(), cfg.nodes[0], cfg.node_wait_secs);
|
||||||
|
tokio::time::sleep(Duration::from_secs(5)).await;
|
||||||
|
}
|
||||||
|
}
|
||||||
// the UTXO-side coinbase address: derived from the pool's EVM address; the execution layer pays the EVM side
|
// the UTXO-side coinbase address: derived from the pool's EVM address; the execution layer pays the EVM side
|
||||||
// (igneum/exec/src/executor.rs), the UTXO output is not spent by v0 (docs/plans/pool.md)
|
// (igneum/exec/src/executor.rs), the UTXO output is not spent by v0 (docs/plans/pool.md)
|
||||||
let prefix = match cfg.network.as_str() {
|
let prefix = match cfg.network.as_str() {
|
||||||
|
|
@ -90,6 +140,7 @@ async fn main() {
|
||||||
next_template_id: AtomicU64::new(0),
|
next_template_id: AtomicU64::new(0),
|
||||||
next_job_id: AtomicU64::new(0),
|
next_job_id: AtomicU64::new(0),
|
||||||
node,
|
node,
|
||||||
|
walker,
|
||||||
extra_nodes: clients,
|
extra_nodes: clients,
|
||||||
pool_address,
|
pool_address,
|
||||||
pool_address_hex: pool_address_hex.clone(),
|
pool_address_hex: pool_address_hex.clone(),
|
||||||
|
|
@ -99,6 +150,8 @@ async fn main() {
|
||||||
verify_permits: tokio::sync::Semaphore::new(cfg.verify_threads),
|
verify_permits: tokio::sync::Semaphore::new(cfg.verify_threads),
|
||||||
started: Instant::now(),
|
started: Instant::now(),
|
||||||
cfg: cfg.clone(),
|
cfg: cfg.clone(),
|
||||||
|
template_failures: AtomicU64::new(0),
|
||||||
|
last_template_ok_ms: AtomicU64::new(0),
|
||||||
});
|
});
|
||||||
let chain = evm.call("eth_chainId", serde_json::json!([])).await.ok().and_then(|v| v.as_str().map(|s| s.to_string()));
|
let chain = evm.call("eth_chainId", serde_json::json!([])).await.ok().and_then(|v| v.as_str().map(|s| s.to_string()));
|
||||||
println!(
|
println!(
|
||||||
|
|
@ -119,8 +172,9 @@ async fn main() {
|
||||||
tokio::spawn(node::confirm_loop(pool.clone()));
|
tokio::spawn(node::confirm_loop(pool.clone()));
|
||||||
tokio::spawn(node::net_loop(pool.clone()));
|
tokio::spawn(node::net_loop(pool.clone()));
|
||||||
tokio::spawn(node::vardiff_loop(pool.clone()));
|
tokio::spawn(node::vardiff_loop(pool.clone()));
|
||||||
tokio::spawn(server::listen(pool.clone()));
|
tokio::spawn(node::status_loop(pool.clone()));
|
||||||
tokio::spawn(api::serve(pool.clone()));
|
tokio::spawn(server::listen(pool.clone(), member_listener));
|
||||||
|
tokio::spawn(api::serve(pool.clone(), api_listener));
|
||||||
{
|
{
|
||||||
let (pool, payer) = (pool.clone(), payer.clone());
|
let (pool, payer) = (pool.clone(), payer.clone());
|
||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
|
|
|
||||||
139
pool/src/node.rs
139
pool/src/node.rs
|
|
@ -39,10 +39,18 @@ pub fn extra_data(member: &Member, pool_address: &[u8; 20]) -> Vec<u8> {
|
||||||
/// cache build), so it runs under `spawn_blocking`.
|
/// cache build), so it runs under `spawn_blocking`.
|
||||||
pub async fn fetch_job(pool: &Arc<Pool>, member: &Arc<Member>) -> Result<(), String> {
|
pub async fn fetch_job(pool: &Arc<Pool>, member: &Arc<Member>) -> Result<(), String> {
|
||||||
let t0 = Instant::now();
|
let t0 = Instant::now();
|
||||||
let tmpl = tokio::time::timeout(Duration::from_secs(5), pool.node.get_block_template(pool.pay_address.clone(), extra_data(member, &pool.pool_address)))
|
let tmpl = match tokio::time::timeout(Duration::from_secs(pool.cfg.template_timeout_s), pool.node.get_block_template(pool.pay_address.clone(), extra_data(member, &pool.pool_address))).await {
|
||||||
.await
|
Ok(Ok(t)) => t,
|
||||||
.map_err(|_| "template timed out".to_string())?
|
Ok(Err(e)) => {
|
||||||
.map_err(|e| e.to_string())?;
|
pool.template_failures.fetch_add(1, Ordering::Relaxed);
|
||||||
|
return Err(e.to_string());
|
||||||
|
}
|
||||||
|
Err(_) => {
|
||||||
|
pool.template_failures.fetch_add(1, Ordering::Relaxed);
|
||||||
|
return Err(format!("template timed out after {} s (--template-timeout-s)", pool.cfg.template_timeout_s));
|
||||||
|
}
|
||||||
|
};
|
||||||
|
pool.last_template_ok_ms.store(unix_ms(), Ordering::Relaxed);
|
||||||
let mut raw = tmpl.block;
|
let mut raw = tmpl.block;
|
||||||
raw.header.vote_key_hash = member.key_hash;
|
raw.header.vote_key_hash = member.key_hash;
|
||||||
let block: Block = raw.clone().try_into().map_err(|e| format!("block convert: {e}"))?;
|
let block: Block = raw.clone().try_into().map_err(|e| format!("block convert: {e}"))?;
|
||||||
|
|
@ -276,9 +284,22 @@ pub async fn submit_block(pool: Arc<Pool>, member: Arc<Member>, mut raw: RpcRawB
|
||||||
pool.want_templates.notify_one();
|
pool.want_templates.notify_one();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Where the confirmation walk continues from after a failed `getVirtualChainFromBlock`: the same cursor, retried
|
||||||
|
/// next tick. The first version restarted from the pruning point, and on a 120,000-block devnet chain that is a
|
||||||
|
/// `get_block` per chain block through the connection the templates shared: one RPC timeout under load (7 October 2026,
|
||||||
|
/// 06:01Z) started a walk that starved every template request for the rest of the window. A cursor the node no
|
||||||
|
/// longer knows (pruned, or from a state file of another chain) moves to the sink, and the pending blocks older than
|
||||||
|
/// it resolve by the orphan rule.
|
||||||
|
pub fn cursor_after_failure(cursor: Hash, sink: Hash, cursor_known: bool) -> Hash {
|
||||||
|
if cursor_known { cursor } else { sink }
|
||||||
|
}
|
||||||
|
|
||||||
|
/// At most this many chain blocks (and `get_block` calls) per 5 s tick; the rest continue next tick.
|
||||||
|
pub const WALK_MAX_PER_TICK: usize = 600;
|
||||||
|
|
||||||
/// Confirms pending blocks by walking the selected chain: a block is paid when it is a chain block or in the
|
/// Confirms pending blocks by walking the selected chain: a block is paid when it is a chain block or in the
|
||||||
/// mergeset blues of one (the execution layer pays exactly those, executor.rs `b.is_blue`); a pending block more
|
/// mergeset blues of one (the execution layer pays exactly those, executor.rs `b.is_blue`); a pending block more
|
||||||
/// than `orphan_after_daa` behind the virtual with no blue merge is an orphan.
|
/// than `orphan_after_daa` behind the virtual with no blue merge is an orphan. On its own connection (`pool.walker`).
|
||||||
pub async fn confirm_loop(pool: Arc<Pool>) {
|
pub async fn confirm_loop(pool: Arc<Pool>) {
|
||||||
let mut last_chain: Option<Hash> = None;
|
let mut last_chain: Option<Hash> = None;
|
||||||
loop {
|
loop {
|
||||||
|
|
@ -287,42 +308,64 @@ pub async fn confirm_loop(pool: Arc<Pool>) {
|
||||||
let s = pool.state.lock().unwrap();
|
let s = pool.state.lock().unwrap();
|
||||||
s.blocks.iter().filter(|b| b.status == "pending").map(|b| (b.hash.clone(), b.daa_score)).collect()
|
s.blocks.iter().filter(|b| b.status == "pending").map(|b| (b.hash.clone(), b.daa_score)).collect()
|
||||||
};
|
};
|
||||||
let info = match pool.node.get_block_dag_info().await {
|
let info = match pool.walker.get_block_dag_info().await {
|
||||||
Ok(i) => i,
|
Ok(i) => i,
|
||||||
Err(_) => continue,
|
Err(e) => {
|
||||||
|
eprintln!("{} confirm: getBlockDagInfo failed ({e})", now());
|
||||||
|
continue;
|
||||||
|
}
|
||||||
};
|
};
|
||||||
if last_chain.is_none() {
|
|
||||||
last_chain = Some(info.pruning_point_hash);
|
|
||||||
}
|
|
||||||
if pending.is_empty() {
|
if pending.is_empty() {
|
||||||
// keep the cursor near the tip so the first pending block costs one short walk
|
// keep the cursor near the tip so the first pending block costs one short walk
|
||||||
last_chain = Some(info.sink);
|
last_chain = Some(info.sink);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
let pending_set: HashSet<String> = pending.iter().map(|p| p.0.clone()).collect();
|
// a first tick with pending blocks (a state file kept across a restart): from the sink, never the pruning
|
||||||
let low = last_chain.unwrap();
|
// point; blocks older than the sink that were blue are a payout lost to the restart, said once
|
||||||
let chain = match pool.node.get_virtual_chain_from_block(low, false, None).await {
|
let low = match last_chain {
|
||||||
|
Some(h) => h,
|
||||||
|
None => {
|
||||||
|
println!("{} confirm: {} pending block(s) at start; the walk begins at the sink, older ones resolve by the orphan rule", now(), pending.len());
|
||||||
|
last_chain = Some(info.sink);
|
||||||
|
info.sink
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let t0 = Instant::now();
|
||||||
|
let chain = match pool.walker.get_virtual_chain_from_block(low, false, None).await {
|
||||||
Ok(c) => c,
|
Ok(c) => c,
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
eprintln!("{} confirm: getVirtualChainFromBlock failed ({e}); restarting from the pruning point", now());
|
let known = pool.walker.get_block(low, false).await.is_ok();
|
||||||
last_chain = Some(info.pruning_point_hash);
|
let next = cursor_after_failure(low, info.sink, known);
|
||||||
|
eprintln!("{} confirm: getVirtualChainFromBlock from {} failed ({e}); cursor {}", now(), low, if next == low { "kept, retried next tick".to_string() } else { format!("unknown to the node, moved to the sink {next}") });
|
||||||
|
last_chain = Some(next);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
let added = &chain.added_chain_block_hashes;
|
||||||
|
let take = added.len().min(WALK_MAX_PER_TICK);
|
||||||
let mut blues: HashSet<String> = HashSet::new();
|
let mut blues: HashSet<String> = HashSet::new();
|
||||||
for h in &chain.added_chain_block_hashes {
|
let mut fetched = 0usize;
|
||||||
|
for h in &added[..take] {
|
||||||
blues.insert(h.to_string());
|
blues.insert(h.to_string());
|
||||||
if let Ok(b) = pool.node.get_block(*h, false).await
|
match pool.walker.get_block(*h, false).await {
|
||||||
&& let Some(v) = b.verbose_data
|
Ok(b) => {
|
||||||
{
|
fetched += 1;
|
||||||
for m in v.merge_set_blues_hashes {
|
if let Some(v) = b.verbose_data {
|
||||||
blues.insert(m.to_string());
|
for m in v.merge_set_blues_hashes {
|
||||||
|
blues.insert(m.to_string());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
eprintln!("{} confirm: getBlock {h} failed ({e}); the walk stops here and continues next tick", now());
|
||||||
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
|
||||||
if let Some(h) = chain.added_chain_block_hashes.last() {
|
|
||||||
last_chain = Some(*h);
|
last_chain = Some(*h);
|
||||||
}
|
}
|
||||||
|
if added.len() > WALK_MAX_PER_TICK || t0.elapsed() > Duration::from_secs(2) {
|
||||||
|
println!("{} confirm: walked {fetched} of {} chain blocks in {:.0} ms ({} pending)", now(), added.len(), t0.elapsed().as_secs_f64() * 1e3, pending.len());
|
||||||
|
}
|
||||||
for (hash, daa) in pending {
|
for (hash, daa) in pending {
|
||||||
if blues.contains(&hash) {
|
if blues.contains(&hash) {
|
||||||
let r = pool.state.lock().unwrap().confirm_block(&hash, pool.cfg.fee_percent);
|
let r = pool.state.lock().unwrap().confirm_block(&hash, pool.cfg.fee_percent);
|
||||||
|
|
@ -340,7 +383,29 @@ pub async fn confirm_loop(pool: Arc<Pool>) {
|
||||||
println!("{} ORPHAN {} (daa {}, virtual {}): not blue within {} DAA", now(), &hash[..16], daa, info.virtual_daa_score, pool.cfg.orphan_after_daa);
|
println!("{} ORPHAN {} (daa {}, virtual {}): not blue within {} DAA", now(), &hash[..16], daa, info.virtual_daa_score, pool.cfg.orphan_after_daa);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
let _ = pending_set;
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The daemon's own status line every 30 s (7 October 2026: a night where four members sat connected with no job for
|
||||||
|
/// 20 minutes was visible only as per-member "template timed out" lines).
|
||||||
|
pub async fn status_loop(pool: Arc<Pool>) {
|
||||||
|
loop {
|
||||||
|
tokio::time::sleep(Duration::from_secs(30)).await;
|
||||||
|
let members = pool.members_snapshot().len();
|
||||||
|
let failures = pool.template_failures.load(Ordering::Relaxed);
|
||||||
|
let last_ok = pool.last_template_ok_ms.load(Ordering::Relaxed);
|
||||||
|
let jobs = pool.next_job_id.load(Ordering::Relaxed);
|
||||||
|
let (accepted, rejected, blocks) = {
|
||||||
|
let s = pool.state.lock().unwrap();
|
||||||
|
(s.accepted, s.rejected, s.blocks_found)
|
||||||
|
};
|
||||||
|
let node_state = if last_ok == 0 {
|
||||||
|
"NO TEMPLATE YET from the node: no job has been issued".to_string()
|
||||||
|
} else {
|
||||||
|
let ago = unix_ms().saturating_sub(last_ok) / 1000;
|
||||||
|
if ago > 30 { format!("NODE NOT ANSWERING: last template {ago} s ago") } else { format!("node ok, last template {ago} s ago") }
|
||||||
|
};
|
||||||
|
println!("{} STATUS members={members} jobs_issued={jobs} shares_accepted={accepted} shares_rejected={rejected} blocks={blocks} template_failures={failures}; {node_state}", now());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -348,17 +413,17 @@ pub async fn confirm_loop(pool: Arc<Pool>) {
|
||||||
pub async fn net_loop(pool: Arc<Pool>) {
|
pub async fn net_loop(pool: Arc<Pool>) {
|
||||||
loop {
|
loop {
|
||||||
let mut n = NetInfo { network: pool.cfg.network.clone(), chain_id: pool.cfg.chain_id(), ..Default::default() };
|
let mut n = NetInfo { network: pool.cfg.network.clone(), chain_id: pool.cfg.chain_id(), ..Default::default() };
|
||||||
if let Ok(i) = pool.node.get_block_dag_info().await {
|
if let Ok(i) = pool.walker.get_block_dag_info().await {
|
||||||
n.difficulty = i.difficulty;
|
n.difficulty = i.difficulty;
|
||||||
n.daa_score = i.virtual_daa_score;
|
n.daa_score = i.virtual_daa_score;
|
||||||
n.block_count = i.block_count;
|
n.block_count = i.block_count;
|
||||||
n.network = i.network.to_string();
|
n.network = i.network.to_string();
|
||||||
}
|
}
|
||||||
n.hashrate = pool.node.estimate_network_hashes_per_second(1000, None).await.ok().map(|h| h as f64);
|
n.hashrate = pool.walker.estimate_network_hashes_per_second(1000, None).await.ok().map(|h| h as f64);
|
||||||
if let Ok(b) = pool.node.get_sink_blue_score().await {
|
if let Ok(b) = pool.walker.get_sink_blue_score().await {
|
||||||
n.blue_score = b;
|
n.blue_score = b;
|
||||||
}
|
}
|
||||||
if let Ok(i) = pool.node.get_info().await {
|
if let Ok(i) = pool.walker.get_info().await {
|
||||||
n.synced = i.is_synced;
|
n.synced = i.is_synced;
|
||||||
n.node_version = i.server_version;
|
n.node_version = i.server_version;
|
||||||
}
|
}
|
||||||
|
|
@ -407,3 +472,21 @@ pub async fn vardiff_loop(pool: Arc<Pool>) {
|
||||||
let _ = pool.next_member_id.load(Ordering::Relaxed);
|
let _ = pool.next_member_id.load(Ordering::Relaxed);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod walk_tests {
|
||||||
|
use super::*;
|
||||||
|
|
||||||
|
/// The cursor rule after a failed chain call: known cursor kept (the 6 October shape restarted from the pruning
|
||||||
|
/// point and walked 120,000 blocks), unknown cursor moved to the sink.
|
||||||
|
#[test]
|
||||||
|
fn a_failed_chain_call_keeps_a_known_cursor_and_never_goes_to_the_pruning_point() {
|
||||||
|
let cursor = Hash::from_bytes([1u8; 32]);
|
||||||
|
let sink = Hash::from_bytes([2u8; 32]);
|
||||||
|
let pruning = Hash::from_bytes([3u8; 32]);
|
||||||
|
assert_eq!(cursor_after_failure(cursor, sink, true), cursor);
|
||||||
|
assert_eq!(cursor_after_failure(cursor, sink, false), sink);
|
||||||
|
assert_ne!(cursor_after_failure(cursor, sink, false), pruning);
|
||||||
|
assert!(WALK_MAX_PER_TICK <= 1000, "a tick's walk stays bounded");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -100,6 +100,10 @@ pub struct Pool {
|
||||||
pub next_template_id: AtomicU64,
|
pub next_template_id: AtomicU64,
|
||||||
pub next_job_id: AtomicU64,
|
pub next_job_id: AtomicU64,
|
||||||
pub node: Arc<GrpcClient>,
|
pub node: Arc<GrpcClient>,
|
||||||
|
/// A second connection to the same node for the confirmation walk and the network numbers (7 October 2026,
|
||||||
|
/// 06:01Z and 06:12:55Z: the walk's `get_block` calls shared the template connection and starved every template
|
||||||
|
/// request; the gRPC client is one request stream per connection)
|
||||||
|
pub walker: Arc<GrpcClient>,
|
||||||
pub extra_nodes: Vec<Arc<GrpcClient>>,
|
pub extra_nodes: Vec<Arc<GrpcClient>>,
|
||||||
/// The pool's EVM coinbase address, lowercase 0x hex, named in every template's `IGNA` field
|
/// The pool's EVM coinbase address, lowercase 0x hex, named in every template's `IGNA` field
|
||||||
pub pool_address: [u8; 20],
|
pub pool_address: [u8; 20],
|
||||||
|
|
@ -109,6 +113,10 @@ pub struct Pool {
|
||||||
pub want_templates: tokio::sync::Notify,
|
pub want_templates: tokio::sync::Notify,
|
||||||
pub verify_permits: tokio::sync::Semaphore,
|
pub verify_permits: tokio::sync::Semaphore,
|
||||||
pub started: Instant,
|
pub started: Instant,
|
||||||
|
/// Template fetches that failed or timed out since start, and the unix ms of the last one that succeeded (0 =
|
||||||
|
/// never): a daemon whose node answers no template issues no job, and the STATUS line in the log says so
|
||||||
|
pub template_failures: AtomicU64,
|
||||||
|
pub last_template_ok_ms: AtomicU64,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Pool {
|
impl Pool {
|
||||||
|
|
|
||||||
|
|
@ -19,11 +19,14 @@ fn now() -> String {
|
||||||
format!("{}.{:03}", t.as_secs(), t.subsec_millis())
|
format!("{}.{:03}", t.as_secs(), t.subsec_millis())
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn listen(pool: Arc<Pool>) {
|
/// Binds a listener or says exactly why not. Called from `main` BEFORE anything else starts (7 October 2026, the fleet
|
||||||
let listener = tokio::net::TcpListener::bind(&pool.cfg.listen).await.unwrap_or_else(|e| {
|
/// agent's row from the 6 October night: a daemon whose member port another process still held ran for 20 minutes as a
|
||||||
eprintln!("cannot listen on {}: {e}", pool.cfg.listen);
|
/// process with no socket; now a bind that fails ends the process with exit code 2 before a member can be let down).
|
||||||
std::process::exit(1)
|
pub async fn bind_listener(what: &str, addr: &str) -> Result<tokio::net::TcpListener, String> {
|
||||||
});
|
tokio::net::TcpListener::bind(addr).await.map_err(|e| format!("cannot bind the {what} listener on {addr}: {e} (another process holds the port, or the address is not this machine's)"))
|
||||||
|
}
|
||||||
|
|
||||||
|
pub async fn listen(pool: Arc<Pool>, listener: tokio::net::TcpListener) {
|
||||||
println!("{} pool: members on {} (chain id {}, {})", now(), pool.cfg.listen, pool.cfg.chain_id(), pool.cfg.network);
|
println!("{} pool: members on {} (chain id {}, {})", now(), pool.cfg.listen, pool.cfg.chain_id(), pool.cfg.network);
|
||||||
loop {
|
loop {
|
||||||
match listener.accept().await {
|
match listener.accept().await {
|
||||||
|
|
@ -323,3 +326,20 @@ async fn handle_share(pool: &Arc<Pool>, m: &Arc<Member>, job_id: u64, nonce: u64
|
||||||
tokio::spawn(submit_block(pool.clone(), m.clone(), raw, nonce, shift, daa, blue));
|
tokio::spawn(submit_block(pool.clone(), m.clone(), raw, nonce, shift, daa, blue));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod bind_tests {
|
||||||
|
use super::bind_listener;
|
||||||
|
|
||||||
|
/// Known good: a free port binds. Known failed: the same port, held, is refused with the address in the message
|
||||||
|
/// (the shape the fleet saw on 6 October: a node still holding 4463).
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_held_port_is_refused_loudly_and_a_free_one_binds() {
|
||||||
|
let held = bind_listener("members", "127.0.0.1:0").await.expect("a free port binds");
|
||||||
|
let addr = held.local_addr().unwrap().to_string();
|
||||||
|
let err = bind_listener("members", &addr).await.err().expect("the held port is refused");
|
||||||
|
assert!(err.contains(&addr) && err.contains("members"), "{err}");
|
||||||
|
drop(held);
|
||||||
|
bind_listener("members", &addr).await.expect("released, it binds again");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -79,8 +79,11 @@ impl Vardiff {
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
if n == 0 {
|
if n == 0 {
|
||||||
// nothing yet after an interval: wait up to three intervals, then one easier step
|
// nothing yet after an interval: wait up to three intervals, then one easier step. The sized
|
||||||
if (now_s - self.started_s) < 3.0 * self.interval_s {
|
// correction stays owed (first_done stays false): 7 October 2026, pm-1's worker compiled its pack for
|
||||||
|
// 40 s, the idle easing consumed the first correction, and the measured 190 shares a second then
|
||||||
|
// walked down one step per 30 s for 5.5 minutes, saturating the verifier
|
||||||
|
if now_s - self.last_change_s.max(self.started_s) < 3.0 * self.interval_s {
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
self.shift = (self.shift + 1).min(cap);
|
self.shift = (self.shift + 1).min(cap);
|
||||||
|
|
@ -92,8 +95,8 @@ impl Vardiff {
|
||||||
} else if per_interval < 0.66 {
|
} else if per_interval < 0.66 {
|
||||||
self.shift = (self.shift + steps).min(cap);
|
self.shift = (self.shift + steps).min(cap);
|
||||||
}
|
}
|
||||||
|
self.first_done = true;
|
||||||
}
|
}
|
||||||
self.first_done = true;
|
|
||||||
} else {
|
} else {
|
||||||
if since_change < self.min_change_s {
|
if since_change < self.min_change_s {
|
||||||
return None;
|
return None;
|
||||||
|
|
@ -226,4 +229,23 @@ mod tests {
|
||||||
}
|
}
|
||||||
assert_eq!(v.retarget(t, 1 << 20), Some(2), "floor at min_shift");
|
assert_eq!(v.retarget(t, 1 << 20), Some(2), "floor at min_shift");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// 7 October 2026, pm-1: a member idle for four intervals (its worker compiling) gets eased one step, and the
|
||||||
|
/// first measured rate must still size the jump (190 shares a second at shift 11 is 11 steps away from one per
|
||||||
|
/// 10 s; the sized step takes 8 at once), not walk down one step per 30 s.
|
||||||
|
#[test]
|
||||||
|
fn an_idle_easing_does_not_consume_the_sized_first_correction() {
|
||||||
|
let t64 = 1u64 << 36;
|
||||||
|
let mut v = Vardiff::new(10, 0, 20, 10.0, 0.0);
|
||||||
|
assert_eq!(v.retarget(31.0, t64), Some(11), "idle three intervals: one easier step");
|
||||||
|
assert_eq!(v.retarget(40.0, t64), None, "still idle, nothing within the next three intervals");
|
||||||
|
// the worker is ready: 190 shares a second for 10 s
|
||||||
|
let mut t = 40.0;
|
||||||
|
for _ in 0..1900 {
|
||||||
|
t += 10.0 / 1900.0;
|
||||||
|
v.on_share(t);
|
||||||
|
}
|
||||||
|
let s = v.retarget(t + 0.1, t64).expect("the first measured rate sizes the jump");
|
||||||
|
assert!(s <= 3, "sized first correction from 11 took at most 8 steps at once, got shift {s}");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue