package mining import ( "context" "fmt" "log" "sync" "time" "forge-mesh/internal/api/types" "forge-mesh/internal/policy" ) const ( TierStateProbing = "probing" TierStateActive = "active" TierStateFailed = "failed" TierStateIdle = "idle" ) // ChainConfig controls tier timeouts and hashrate gates. type ChainConfig struct { HostID string StratumHost string StratumXMRPort string StratumRVNPort string MinHashrateHps float64 GateWindow time.Duration PollInterval time.Duration DefaultProbe time.Duration } func DefaultChainConfig() ChainConfig { return ChainConfig{ StratumHost: "127.0.0.1", StratumXMRPort: "3333", StratumRVNPort: "3388", MinHashrateHps: 100, GateWindow: 30 * time.Second, PollInterval: 5 * time.Second, DefaultProbe: 5 * time.Minute, } } // Status is reported in agent heartbeats. type Status struct { CurrentTier int TierType string TierState string HashrateHps float64 Wallet string Simulated bool } // Chain runs the ordered mining tier list with wallet pinning. type Chain struct { cfg ChainConfig profile types.MiningProfile mu sync.RWMutex status Status runner tierRunner cancel context.CancelFunc } type tierRunner interface { Start(ctx context.Context) error Stop() HashrateHps() float64 Simulated() bool } func NewChain(profile types.MiningProfile, cfg ChainConfig) *Chain { if len(profile.Tiers) == 0 { profile = policy.DefaultMiningProfile(profile.WalletAddress) } return &Chain{ cfg: cfg, profile: profile, status: Status{ TierState: TierStateIdle, Wallet: profile.WalletAddress, }, } } func (c *Chain) Profile() types.MiningProfile { return c.profile } func (c *Chain) UpdateProfile(profile types.MiningProfile) { c.mu.Lock() if profile.WalletAddress != "" { c.profile.WalletAddress = profile.WalletAddress c.status.Wallet = profile.WalletAddress } if len(profile.Tiers) > 0 { c.profile.Tiers = profile.Tiers } running := c.cancel != nil c.mu.Unlock() if running { c.Stop() } } func (c *Chain) Status() Status { c.mu.RLock() defer c.mu.RUnlock() return c.status } // Run executes tiers in order until one passes the hashrate gate or all fail. func (c *Chain) Run(ctx context.Context) error { if len(c.profile.Tiers) == 0 { return fmt.Errorf("mining profile has no tiers") } runCtx, cancel := context.WithCancel(ctx) c.mu.Lock() c.cancel = cancel c.mu.Unlock() defer cancel() for i, spec := range c.profile.Tiers { tierNum := i + 1 c.setStatus(tierNum, spec.Type, TierStateProbing, 0, false) runner, err := c.buildRunner(spec) if err != nil { log.Printf("mining: tier %d (%s) build failed: %v", tierNum, spec.Type, err) c.setStatus(tierNum, spec.Type, TierStateFailed, 0, false) continue } c.mu.Lock() c.runner = runner c.mu.Unlock() if err := runner.Start(runCtx); err != nil { log.Printf("mining: tier %d (%s) start failed: %v", tierNum, spec.Type, err) runner.Stop() c.setStatus(tierNum, spec.Type, TierStateFailed, 0, runner.Simulated()) continue } c.setStatus(tierNum, spec.Type, TierStateActive, runner.HashrateHps(), runner.Simulated()) duration := time.Duration(spec.Duration) * time.Minute if duration <= 0 && spec.Type == "stratum" { // Final fallback runs until context cancelled. return c.holdTier(runCtx, tierNum, spec.Type, runner) } if duration <= 0 { duration = c.cfg.DefaultProbe } passed, peak := c.waitGate(runCtx, duration, runner) runner.Stop() if passed { c.setStatus(tierNum, spec.Type, TierStateActive, peak, runner.Simulated()) log.Printf("mining: tier %d (%s) passed hashrate gate at %.0f H/s", tierNum, spec.Type, peak) return c.holdTier(runCtx, tierNum, spec.Type, nil) } log.Printf("mining: tier %d (%s) timed out (peak %.0f H/s), advancing", tierNum, spec.Type, peak) c.setStatus(tierNum, spec.Type, TierStateFailed, peak, runner.Simulated()) } c.setStatus(0, "", TierStateFailed, 0, false) return fmt.Errorf("all mining tiers exhausted") } func (c *Chain) holdTier(ctx context.Context, tierNum int, tierType string, runner tierRunner) error { ticker := time.NewTicker(c.cfg.PollInterval) defer ticker.Stop() for { select { case <-ctx.Done(): c.Stop() return ctx.Err() case <-ticker.C: hps := c.readHashrate(runner) c.setStatus(tierNum, tierType, TierStateActive, hps, c.isSimulated(runner)) } } } func (c *Chain) readHashrate(runner tierRunner) float64 { if runner != nil { return runner.HashrateHps() } c.mu.RLock() r := c.runner c.mu.RUnlock() if r != nil { return r.HashrateHps() } return 0 } func (c *Chain) isSimulated(runner tierRunner) bool { if runner != nil { return runner.Simulated() } c.mu.RLock() r := c.runner c.mu.RUnlock() if r != nil { return r.Simulated() } return false } func (c *Chain) waitGate(ctx context.Context, maxWait time.Duration, runner tierRunner) (bool, float64) { deadline := time.Now().Add(maxWait) gateEnd := time.Now().Add(c.cfg.GateWindow) var peak float64 for time.Now().Before(deadline) { select { case <-ctx.Done(): return false, peak case <-time.After(c.cfg.PollInterval): } hps := runner.HashrateHps() if hps > peak { peak = hps } c.mu.Lock() c.status.HashrateHps = hps c.mu.Unlock() if hps >= c.cfg.MinHashrateHps && time.Now().After(gateEnd) { return true, peak } } return peak >= c.cfg.MinHashrateHps, peak } func (c *Chain) buildRunner(spec types.MiningTierSpec) (tierRunner, error) { wallet := c.profile.WalletAddress switch spec.Type { case "oci", "podman": return newOCITier(wallet, spec, c.cfg) case "xmrig": return newXMRigTier(wallet, spec, c.cfg) case "gpu", "lolminer": return newGPUTier(wallet, spec, c.cfg) case "stratum": return newStratumTier(wallet, spec, c.cfg) default: return nil, fmt.Errorf("unknown tier type %q", spec.Type) } } func (c *Chain) setStatus(tierNum int, tierType, state string, hps float64, simulated bool) { c.mu.Lock() defer c.mu.Unlock() c.status.CurrentTier = tierNum c.status.TierType = tierType c.status.TierState = state c.status.HashrateHps = hps c.status.Simulated = simulated } func (c *Chain) Stop() { c.mu.Lock() cancel := c.cancel runner := c.runner c.mu.Unlock() if runner != nil { runner.Stop() } if cancel != nil { cancel() } }