Files
drjones 3678b199d0
Some checks failed
Test / test (push) Has been cancelled
Initial commit: AetherForge Linux (forge-mesh) v0.1.0-dev
2026-07-04 09:31:23 +00:00

284 lines
6.2 KiB
Go

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()
}
}