stratum-server 3.0.0-beta-9

The server code for the Rust Stratum (v1) implementation
Documentation
use async_std::sync::{Arc, Mutex};
use chrono::{NaiveDateTime, Utc};
use extended_primitives::Buffer;
use log::warn;
use uuid::Uuid;

use crate::connection::MinerOptions;

//A miner is essentially an individual worker unit. There can be multiple Miners on a single
//connection which is why we needed to break it into these primitives.
#[derive(Debug, Clone)]
pub struct Miner {
    pub id: Uuid,
    pub sid: Buffer,
    pub client: Option<String>,
    pub name: Option<String>,
    pub difficulty: Arc<Mutex<u64>>,
    pub next_difficulty: Arc<Mutex<Option<u64>>>,
    pub stats: Arc<Mutex<MinerStats>>,
    pub job_stats: Arc<Mutex<JobStats>>,
    pub options: Arc<MinerOptions>,
    pub needs_ban: Arc<Mutex<bool>>,
}

//@todo don't forget to update job_stats in update_difficulty
impl Miner {
    //@todo going to need to add a difficulty to this guy.
    pub fn new(
        id: Uuid,
        client: Option<String>,
        name: Option<String>,
        sid: Buffer,
        options: Arc<MinerOptions>,
        difficulty: u64,
    ) -> Self {
        Miner {
            id,
            sid,
            client,
            name,
            difficulty: Arc::new(Mutex::new(difficulty)),
            next_difficulty: Arc::new(Mutex::new(None)),
            stats: Arc::new(Mutex::new(MinerStats {
                accepted_shares: 0,
                rejected_shares: 0,
                last_active: Utc::now().naive_utc(),
            })),
            job_stats: Arc::new(Mutex::new(JobStats {
                last_timestamp: Utc::now().naive_utc().timestamp(),
                last_retarget: Utc::now().naive_utc().timestamp()
                    - options.retarget_time as i64 / 2,
                times: Vec::new(),
                current_difficulty: difficulty,
            })),
            options,
            needs_ban: Arc::new(Mutex::new(false)),
        }
    }

    pub async fn ban(&self) {
        *self.needs_ban.lock().await = true;
        // self.disconnect().await;
    }

    pub async fn consider_ban(&self) {
        let accepted = self.stats.lock().await.accepted_shares;
        let rejected = self.stats.lock().await.rejected_shares;

        let total = accepted + rejected;

        //@todo come from options.
        let check_threshold = 500;
        let invalid_percent = 50.0;

        if total >= check_threshold {
            let percent_bad: f64 = (rejected as f64 / total as f64) * 100.0;

            if percent_bad < invalid_percent {
                //@todo make this possible. Reset stats to 0.
                // self.stats.lock().await = MinerStats::default();
            } else {
                warn!(
                    "Miner: {} banned. {} out of the last {} shares were invalid",
                    self.id, rejected, total
                );
                // self.ban().await; @todo
            }
        }
    }

    pub async fn valid_share(&self) {
        let mut stats = self.stats.lock().await;
        stats.accepted_shares += 1;
        stats.last_active = Utc::now().naive_utc();
        drop(stats);
        // self.consider_ban().await; @todo
        // @todo if we want to wrap this in an option, lets make it options.
        // @todo don't retarget until new job has been added.
        self.retarget().await;
    }

    pub async fn invalid_share(&self) {
        self.stats.lock().await.rejected_shares += 1;
        // self.consider_ban().await;
        //@todo see below
        //I don't think we want to retarget on invalid shares, but let's double check later.
        // self.retarget().await;
    }

    //@todo note, this only can be sent over ExMessage when it's hit a certain threshold.
    //@todo self.set_difficulty
    //@todo self.set_next_difficulty
    //@todo does this need to return a result? Ideally not, but if we send difficulty, then maybe.
    //@todo see if we can solve a lot of these recasting issues.
    //@todo wrap u64 with a custom difficulty type.
    async fn retarget(&self) {
        let now = Utc::now().naive_utc().timestamp();

        let mut job_stats = self.job_stats.lock().await;

        let since_last = now - job_stats.last_timestamp;

        job_stats.times.push(since_last);
        job_stats.last_timestamp = now;

        //@todo in the futuer we may want to support non-full window difficulty adjustments. For now, this
        //works.
        if now - job_stats.last_retarget < self.options.retarget_time as i64 {
            return;
        }

        // let variance = self.options.target_time * (self.options.variance_percent as f64 / 100.0);
        let time_min = self.options.target_time as f64 * 0.40;
        let time_max = self.options.target_time as f64 * 1.40;
        job_stats.last_retarget = now;

        let mut avg: i64 = job_stats.times.iter().sum::<i64>() / job_stats.times.len() as i64;
        // let mut d_dif = self.options.target_time as f64 / avg as f64;

        let mut new_difficulty = job_stats.current_difficulty.clone();
        // Too Fast
        if (avg as f64) < time_min {
            while (avg as f64) < time_min && new_difficulty < self.options.max_diff {
                new_difficulty *= 2;
                avg *= 2;
            }

            *self.next_difficulty.lock().await = Some(new_difficulty);
        }

        // Too SLow
        if (avg as f64) > time_max && new_difficulty >= self.options.min_diff * 2 {
            while (avg as f64) > time_max && new_difficulty >= self.options.min_diff * 2 {
                new_difficulty /= 2;
                avg /= 2;
            }
            *self.next_difficulty.lock().await = Some(new_difficulty);
        }

        job_stats.times.clear();
    }

    pub async fn update_difficulty(&self) -> Option<u64> {
        let next_difficulty = *self.next_difficulty.lock().await;

        if let Some(next_difficulty) = next_difficulty {
            *self.difficulty.lock().await = next_difficulty;
            self.job_stats.lock().await.current_difficulty = next_difficulty;

            *self.next_difficulty.lock().await = None;

            Some(next_difficulty)
        } else {
            None
        }
    }
}

#[derive(Debug, Clone)]
pub struct MinerStats {
    accepted_shares: u64,
    rejected_shares: u64,
    last_active: NaiveDateTime,
}

//@todo probably move these over to types.
//@todo maybe rename this as vardiff stats.
#[derive(Debug, Default)]
pub struct JobStats {
    last_timestamp: i64,
    last_retarget: i64,
    times: Vec<i64>,
    current_difficulty: u64,
}