Files
AetherForge/server/internal/pool/proxy.go

904 lines
20 KiB
Go

package pool
import (
"bufio"
"crypto/tls"
"encoding/binary"
"encoding/hex"
"encoding/json"
"fmt"
"log"
"math/big"
"net"
"strconv"
"strings"
"sync"
"time"
"crypto-miner-server/internal/models"
)
// Stratum protocol message types
type StratumRequest struct {
ID int `json:"id"`
Method string `json:"method"`
Params json.RawMessage `json:"params"`
}
type StratumResponse struct {
ID int `json:"id"`
Result json.RawMessage `json:"result"`
Error interface{} `json:"error"`
}
type StratumNotification struct {
Method string `json:"method"`
Params json.RawMessage `json:"params"`
}
// Job represents a mining job from the pool
type Job struct {
ID string `json:"job_id"`
Height int64 `json:"height"`
BlockTemplate string `json:"blocktemplate"`
Difficulty int64 `json:"difficulty"`
SeedHash string `json:"seed_hash"`
Target string `json:"target"`
Blob string `json:"blob"`
Algo string `json:"algo"`
}
// ShareSubmit represents a share submission to the pool
type ShareSubmit struct {
ID int `json:"id"`
Method string `json:"method"`
Params []string `json:"params"`
}
// Proxy connects to a Monero mining pool via Stratum protocol
// and acts as a bridge between the pool and our agents
type Proxy struct {
mu sync.RWMutex
config *Config
conn net.Conn
reader *bufio.Reader
connected bool
requestID int
loginRequestID int
currentJob *Job
jobSubscribed bool
stopCh chan struct{}
wg sync.WaitGroup
running bool
// Callbacks
onJob func(job *Job)
onShare func(accepted bool, agentID string, jobID string)
onError func(err error)
// Agent share submissions queue
shareQueue chan *PendingShare
pendingMu sync.Mutex
pendingResults map[int]*pendingShareResult
reconnecting bool
reconnectDelay time.Duration
verboseTraffic bool
// Pool failover: if the active config fails on reconnect, cycle through
// these alternatives before giving up. The primary Config is always index -1
// (the original); backups are tried in order. primaryConfig preserves the
// original so we can reset after a full round fails.
primaryConfig Config
backupConfigs []Config
backupIdx int // next backup to try
}
type PendingShare struct {
AgentID string
JobID string
Nonce string
Hash string
Wallet string
OnResult func(accepted bool, errMsg string)
}
type pendingShareResult struct {
AgentID string
JobID string
OnResult func(accepted bool, errMsg string)
}
type Config struct {
Host string
Port int
UseTLS bool
Wallet string
Password string
// PaymentID is an optional Monero integrated-address payment ID.
// When non-empty it is appended to the wallet login as "wallet.paymentID"
// so the pool can credit payouts to the correct sub-account.
PaymentID string
}
func NewProxy(cfg *Config) *Proxy {
return &Proxy{
config: cfg,
primaryConfig: *cfg,
stopCh: make(chan struct{}),
shareQueue: make(chan *PendingShare, 100),
pendingResults: make(map[int]*pendingShareResult),
reconnectDelay: 10 * time.Second,
}
}
// SetBackupConfigs registers fallback pool configs tried in order when the
// active config fails on reconnect. Call before Start().
func (p *Proxy) SetBackupConfigs(backups []Config) {
p.mu.Lock()
p.backupConfigs = backups
p.mu.Unlock()
}
func (p *Proxy) SetReconnectDelay(d time.Duration) {
if d > 0 {
p.mu.Lock()
p.reconnectDelay = d
p.mu.Unlock()
}
}
func (p *Proxy) SetVerboseTraffic(enabled bool) {
p.mu.Lock()
p.verboseTraffic = enabled
p.mu.Unlock()
}
func (p *Proxy) trafficLog(format string, args ...interface{}) {
p.mu.RLock()
v := p.verboseTraffic
p.mu.RUnlock()
if v {
log.Printf(format, args...)
}
}
// SetCallbacks sets the callbacks for job updates and share results
func (p *Proxy) SetCallbacks(onJob func(job *Job), onShare func(accepted bool, agentID string, jobID string), onError func(err error)) {
p.mu.Lock()
defer p.mu.Unlock()
p.onJob = onJob
p.onShare = onShare
p.onError = onError
}
// Start connects to the pool and begins processing.
func (p *Proxy) Start() error {
p.mu.Lock()
if !p.running {
p.running = true
p.wg.Add(2)
go p.readLoop()
go p.shareSubmitLoop()
}
p.mu.Unlock()
return p.connect()
}
func (p *Proxy) connect() error {
addr := net.JoinHostPort(p.config.Host, strconv.Itoa(p.config.Port))
log.Printf("[Pool] Connecting to %s (TLS: %v)...", addr, p.config.UseTLS)
var conn net.Conn
var err error
if p.config.UseTLS {
netDialer := &net.Dialer{Timeout: 30 * time.Second}
tlsConn, tlsErr := tls.DialWithDialer(netDialer, "tcp", addr, &tls.Config{})
conn = tlsConn
err = tlsErr
} else {
dialer := net.Dialer{Timeout: 30 * time.Second}
conn, err = dialer.Dial("tcp", addr)
}
if err != nil {
return fmt.Errorf("failed to connect to pool: %w", err)
}
p.mu.Lock()
if p.conn != nil {
_ = p.conn.Close()
}
p.conn = conn
p.reader = bufio.NewReader(conn)
p.connected = true
p.mu.Unlock()
log.Printf("[Pool] Connected to %s", addr)
if err := p.authenticate(); err != nil {
return fmt.Errorf("failed to authenticate with pool: %w", err)
}
return nil
}
// Stop disconnects from the pool
func (p *Proxy) Stop() {
p.mu.Lock()
if p.stopCh != nil {
select {
case <-p.stopCh:
default:
close(p.stopCh)
}
}
if p.conn != nil {
p.conn.Close()
p.connected = false
}
p.mu.Unlock()
p.wg.Wait()
log.Println("[Pool] Disconnected from pool")
}
// IsConnected returns whether the proxy is connected to the pool
func (p *Proxy) IsConnected() bool {
p.mu.RLock()
defer p.mu.RUnlock()
return p.connected
}
func (p *Proxy) Config() Config {
p.mu.RLock()
defer p.mu.RUnlock()
if p.config == nil {
return Config{}
}
return *p.config
}
// GetCurrentJob returns the current mining job
func (p *Proxy) GetCurrentJob() *Job {
p.mu.RLock()
defer p.mu.RUnlock()
if p.currentJob == nil {
return nil
}
jobCopy := *p.currentJob
return &jobCopy
}
// SubmitShare queues a share for submission to the pool using the forged wallet.
func (p *Proxy) SubmitShare(agentID, wallet, jobID, nonce, hash string, onResult func(accepted bool, errMsg string)) {
p.shareQueue <- &PendingShare{
AgentID: agentID,
JobID: jobID,
Nonce: nonce,
Hash: hash,
Wallet: wallet,
OnResult: onResult,
}
}
func (p *Proxy) authenticate() error {
p.mu.Lock()
p.requestID++
loginID := p.requestID
p.loginRequestID = loginID
p.mu.Unlock()
// Build the login wallet string. If a payment ID is configured, the pool
// receives "wallet.paymentID" which routes payouts to the correct account.
walletLogin := p.config.Wallet
if p.config.PaymentID != "" {
walletLogin = walletLogin + "." + p.config.PaymentID
}
// Login request
loginParams := []interface{}{
walletLogin,
p.config.Password,
"crypto-miner-server/1.0",
}
paramsData, _ := json.Marshal(loginParams)
loginReq := StratumRequest{
ID: loginID,
Method: "login",
Params: paramsData,
}
data, _ := json.Marshal(loginReq)
p.trafficLog("[Pool] Sending login request...")
if err := p.writeLine(data); err != nil {
return fmt.Errorf("failed to send login: %w", err)
}
return nil
}
func (p *Proxy) readLoop() {
defer p.wg.Done()
for {
select {
case <-p.stopCh:
return
default:
}
p.mu.RLock()
reader := p.reader
p.mu.RUnlock()
if reader == nil {
time.Sleep(100 * time.Millisecond)
continue
}
line, err := reader.ReadString('\n')
if err != nil {
log.Printf("[Pool] Read error: %v", err)
p.mu.Lock()
p.connected = false
p.mu.Unlock()
if p.onError != nil {
p.onError(fmt.Errorf("pool connection lost: %w", err))
}
p.scheduleReconnect()
// Wait for reconnect before resuming reads (readLoop stays alive).
for i := 0; i < 600; i++ {
select {
case <-p.stopCh:
return
default:
}
p.mu.RLock()
ok := p.connected && p.reader != nil
p.mu.RUnlock()
if ok {
break
}
time.Sleep(100 * time.Millisecond)
}
continue
}
line = strings.TrimSpace(line)
if line == "" {
continue
}
p.handleMessage([]byte(line))
}
}
func (p *Proxy) handleMessage(data []byte) {
// Try to parse as response first
var resp StratumResponse
if err := json.Unmarshal(data, &resp); err == nil && resp.ID > 0 {
p.handleResponse(resp)
return
}
// Try to parse as notification
var notif StratumNotification
if err := json.Unmarshal(data, &notif); err == nil && notif.Method != "" {
p.handleNotification(notif)
return
}
p.trafficLog("[Pool] Unhandled message: %s", string(data))
}
func (p *Proxy) handleResponse(resp StratumResponse) {
p.trafficLog("[Pool] Response ID=%d: %s", resp.ID, string(resp.Result))
// Share submit responses (any ID > 1 that we tracked)
p.pendingMu.Lock()
pending, tracked := p.pendingResults[resp.ID]
if tracked {
delete(p.pendingResults, resp.ID)
}
p.pendingMu.Unlock()
if tracked && pending != nil && pending.OnResult != nil {
accepted, errMsg := parseSubmitResult(resp)
pending.OnResult(accepted, errMsg)
if p.onShare != nil {
p.onShare(accepted, pending.AgentID, pending.JobID)
}
return
}
p.mu.RLock()
loginID := p.loginRequestID
p.mu.RUnlock()
if resp.ID == loginID {
// Login response
var loginResult struct {
ID string `json:"id"`
Job json.RawMessage `json:"job"`
Status string `json:"status"`
}
if err := json.Unmarshal(resp.Result, &loginResult); err != nil {
log.Printf("[Pool] Failed to parse login result: %v", err)
return
}
log.Printf("[Pool] Login successful! Pool ID: %s, Status: %s", loginResult.ID, loginResult.Status)
// Parse initial job if provided
if len(loginResult.Job) > 0 {
p.parseAndSetJob(loginResult.Job)
}
// Subscribe for jobs
p.subscribe()
}
}
func (p *Proxy) handleNotification(notif StratumNotification) {
switch notif.Method {
case "job":
p.trafficLog("[Pool] New job received")
p.parseAndSetJob(notif.Params)
case "submit":
// Share submission result
var submitResult struct {
ID int `json:"id"`
Result string `json:"result"`
Status string `json:"status"`
}
if err := json.Unmarshal(notif.Params, &submitResult); err != nil {
log.Printf("[Pool] Failed to parse submit result: %v", err)
return
}
p.trafficLog("[Pool] Share submission result: %s", submitResult.Status)
default:
p.trafficLog("[Pool] Unknown notification method: %s", notif.Method)
}
}
func (p *Proxy) parseAndSetJob(data json.RawMessage) {
var rawJob struct {
ID string `json:"job_id"`
Height int64 `json:"height"`
BlockTemplate string `json:"blocktemplate"`
Difficulty int64 `json:"difficulty"`
SeedHash string `json:"seed_hash"`
Target string `json:"target"`
Blob string `json:"blob"`
Algo string `json:"algo"`
}
// Try different field name variations that pools use
if err := json.Unmarshal(data, &rawJob); err != nil {
// Try flat params
var flatParams []json.RawMessage
if err2 := json.Unmarshal(data, &flatParams); err2 == nil && len(flatParams) >= 1 {
json.Unmarshal(flatParams[0], &rawJob)
} else {
// Try as array of params
var params [][]json.RawMessage
if err3 := json.Unmarshal(data, &params); err3 == nil && len(params) >= 1 && len(params[0]) >= 1 {
json.Unmarshal(params[0][0], &rawJob)
} else {
log.Printf("[Pool] Failed to parse job: %s", string(data))
return
}
}
}
job := &Job{
ID: rawJob.ID,
Height: rawJob.Height,
BlockTemplate: rawJob.BlockTemplate,
Difficulty: rawJob.Difficulty,
SeedHash: rawJob.SeedHash,
Target: rawJob.Target,
Blob: rawJob.Blob,
Algo: rawJob.Algo,
}
// Calculate target from difficulty if not provided
if job.Target == "" && job.Difficulty > 0 {
job.Target = p.difficultyToTarget(job.Difficulty)
}
p.mu.Lock()
p.currentJob = job
p.mu.Unlock()
p.trafficLog("[Pool] New job: ID=%s, Height=%d, Difficulty=%d, Algo=%s",
job.ID, job.Height, job.Difficulty, job.Algo)
if p.onJob != nil {
p.onJob(job)
}
}
func (p *Proxy) subscribe() {
p.mu.Lock()
p.requestID++
subID := p.requestID
p.mu.Unlock()
subParams := []string{}
paramsData, _ := json.Marshal(subParams)
subReq := StratumRequest{
ID: subID,
Method: "subscribe",
Params: paramsData,
}
data, _ := json.Marshal(subReq)
p.trafficLog("[Pool] Subscribing for jobs...")
if err := p.writeLine(data); err != nil {
log.Printf("[Pool] Failed to subscribe: %v", err)
}
}
func (p *Proxy) shareSubmitLoop() {
defer p.wg.Done()
for {
select {
case <-p.stopCh:
return
case share := <-p.shareQueue:
p.submitShareToPool(share)
}
}
}
func (p *Proxy) submitShareToPool(share *PendingShare) {
p.mu.RLock()
connected := p.connected
p.mu.RUnlock()
if !connected {
log.Printf("[Pool] Cannot submit share - not connected to pool")
if share.OnResult != nil {
share.OnResult(false, "pool not connected")
}
return
}
// Increment and read requestID under the write lock to avoid data race (H17)
p.mu.Lock()
p.requestID++
reqID := p.requestID
p.mu.Unlock()
wallet := share.Wallet
if wallet == "" {
wallet = p.config.Wallet
}
p.pendingMu.Lock()
p.pendingResults[reqID] = &pendingShareResult{
AgentID: share.AgentID,
JobID: share.JobID,
OnResult: share.OnResult,
}
p.pendingMu.Unlock()
submitParams := []string{
wallet,
share.JobID,
share.Nonce,
share.Hash,
}
paramsData, _ := json.Marshal(submitParams)
submitReq := StratumRequest{
ID: reqID,
Method: "submit",
Params: paramsData,
}
data, _ := json.Marshal(submitReq)
p.trafficLog("[Pool] Submitting share for agent %s (job: %s)...", share.AgentID[:min(8, len(share.AgentID))], share.JobID)
if err := p.writeLine(data); err != nil {
log.Printf("[Pool] Failed to submit share: %v", err)
p.pendingMu.Lock()
delete(p.pendingResults, reqID)
p.pendingMu.Unlock()
if share.OnResult != nil {
share.OnResult(false, err.Error())
}
if p.onShare != nil {
p.onShare(false, share.AgentID, share.JobID)
}
return
}
// Timeout fallback if pool never responds
go func(id int, ps *PendingShare) {
time.Sleep(30 * time.Second)
p.pendingMu.Lock()
pending, ok := p.pendingResults[id]
if ok {
delete(p.pendingResults, id)
}
p.pendingMu.Unlock()
if ok && pending != nil && pending.OnResult != nil {
log.Printf("[Pool] Share response timeout for agent %s job %s", pending.AgentID, pending.JobID)
pending.OnResult(false, "pool response timeout")
}
}(reqID, share)
}
func parseSubmitResult(resp StratumResponse) (accepted bool, errMsg string) {
if resp.Error != nil {
switch v := resp.Error.(type) {
case string:
return false, v
case []interface{}:
if len(v) > 1 {
if s, ok := v[1].(string); ok {
return false, s
}
}
case map[string]interface{}:
if msg, ok := v["message"].(string); ok {
return false, msg
}
}
return false, "pool rejected share"
}
if len(resp.Result) == 0 {
return true, ""
}
var status struct {
Status string `json:"status"`
}
if err := json.Unmarshal(resp.Result, &status); err == nil && status.Status != "" {
if strings.EqualFold(status.Status, "OK") || strings.EqualFold(status.Status, "ACCEPTED") {
return true, ""
}
return false, status.Status
}
var boolResult bool
if err := json.Unmarshal(resp.Result, &boolResult); err == nil {
if boolResult {
return true, ""
}
return false, "pool rejected share"
}
return true, ""
}
func (p *Proxy) scheduleReconnect() {
p.mu.Lock()
if p.reconnecting {
p.mu.Unlock()
return
}
p.reconnecting = true
p.mu.Unlock()
go func() {
defer func() {
p.mu.Lock()
p.reconnecting = false
p.mu.Unlock()
}()
p.reconnect()
}()
}
func (p *Proxy) reconnect() {
p.mu.RLock()
delay := p.reconnectDelay
if delay <= 0 {
delay = 10 * time.Second
}
p.mu.RUnlock()
log.Printf("[Pool] Attempting reconnect in %s...", delay)
time.Sleep(delay)
select {
case <-p.stopCh:
return
default:
}
// Build the full candidate list: [current primary] + backup configs.
// If the primary fails we rotate through backups so the fleet keeps hashing
// even when the configured pool is down.
p.mu.RLock()
primary := p.primaryConfig
backups := append([]Config(nil), p.backupConfigs...) // copy
startIdx := p.backupIdx
p.mu.RUnlock()
candidates := make([]Config, 0, 1+len(backups))
candidates = append(candidates, primary)
candidates = append(candidates, backups...)
total := len(candidates)
for i := 0; i < total; i++ {
// Rotate: try startIdx (last successful backup) before falling back to primary
idx := (startIdx + i) % total
cfg := candidates[idx]
p.mu.Lock()
*p.config = cfg
p.mu.Unlock()
if err := p.connect(); err != nil {
label := "primary"
if idx > 0 {
label = fmt.Sprintf("backup#%d (%s:%d)", idx, cfg.Host, cfg.Port)
}
log.Printf("[Pool] Reconnect to %s failed: %v", label, err)
continue
}
// Succeeded — remember where we are so next reconnect starts here
p.mu.Lock()
p.backupIdx = idx
p.mu.Unlock()
if idx > 0 {
log.Printf("[Pool] Now connected to backup#%d (%s:%d)", idx, cfg.Host, cfg.Port)
}
return
}
// All candidates failed
if p.onError != nil {
p.onError(fmt.Errorf("all pool endpoints unreachable (%d tried)", total))
}
// Restore primary so next cycle starts fresh
p.mu.Lock()
*p.config = primary
p.backupIdx = 0
p.mu.Unlock()
time.Sleep(30 * time.Second)
select {
case <-p.stopCh:
return
default:
p.scheduleReconnect()
}
}
func (p *Proxy) writeLine(data []byte) error {
p.mu.RLock()
conn := p.conn
p.mu.RUnlock()
if conn == nil {
return fmt.Errorf("not connected")
}
line := append(data, '\n')
_, err := conn.Write(line)
return err
}
func (p *Proxy) difficultyToTarget(difficulty int64) string {
// Convert difficulty to target hex string
// target = 0xFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFF / difficulty
maxTarget := new(big.Int)
maxTarget.SetString("FFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFF", 16)
diff := big.NewInt(difficulty)
target := new(big.Int).Div(maxTarget, diff)
// Convert to 32-byte hex (little-endian for Monero)
bytes := target.Bytes()
padded := make([]byte, 32)
copy(padded[32-len(bytes):], bytes)
// Reverse for little-endian
for i, j := 0, len(padded)-1; i < j; i, j = i+1, j-1 {
padded[i], padded[j] = padded[j], padded[i]
}
return hex.EncodeToString(padded)
}
// Helper to convert models.Job to pool.Job
func FromModelJob(job *models.Job) *Job {
if job == nil {
return nil
}
return &Job{
ID: job.ID,
Height: job.Height,
BlockTemplate: job.BlockTemplate,
Difficulty: job.Difficulty,
SeedHash: job.SeedHash,
Target: job.Target,
}
}
// Helper to convert pool.Job to models.Job
func (j *Job) ToModelJob() *models.Job {
return &models.Job{
ID: j.ID,
Height: j.Height,
Difficulty: j.Difficulty,
BlockTemplate: j.BlockTemplate,
SeedHash: j.SeedHash,
Target: j.Target,
CreatedAt: time.Now(),
}
}
// Helper to convert target hex to difficulty
func targetToDifficulty(targetHex string) int64 {
bytes, err := hex.DecodeString(targetHex)
if err != nil || len(bytes) == 0 {
return 0
}
// Reverse from little-endian
for i, j := 0, len(bytes)-1; i < j; i, j = i+1, j-1 {
bytes[i], bytes[j] = bytes[j], bytes[i]
}
target := new(big.Int).SetBytes(bytes)
if target.Sign() == 0 {
return 0
}
maxTarget := new(big.Int)
maxTarget.SetString("FFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFFF", 16)
diff := new(big.Int).Div(maxTarget, target)
return diff.Int64()
}
// ParseBlob extracts fields from a Monero mining blob
func ParseBlob(blobHex string) (map[string]interface{}, error) {
blob, err := hex.DecodeString(blobHex)
if err != nil {
return nil, fmt.Errorf("invalid blob hex: %w", err)
}
if len(blob) < 43 {
return nil, fmt.Errorf("blob too short: %d bytes", len(blob))
}
result := make(map[string]interface{})
// Monero blob structure (simplified):
// [0:1] - Reserved (1 byte)
// [1:9] - Block ID (8 bytes, little-endian)
// [9:17] - Nonce (8 bytes, little-endian) - miners fill this
// [17:43] - Merkle root + extra data
result["reserved"] = blob[0]
result["block_id"] = binary.LittleEndian.Uint64(blob[1:9])
result["nonce_offset"] = 9
result["nonce_size"] = 4 // Standard nonce is 4 bytes for most pools
return result, nil
}
func min(a, b int) int {
if a < b {
return a
}
return b
}