Ship live alerts, pool status, AI monitor, remote agent commands, build manager, and uninstall flow; wire Calibrate settings (WS ping, pool traffic log, retention limits) at runtime and exclude server/data from git.
352 lines
12 KiB
Go
352 lines
12 KiB
Go
package db
|
|
|
|
import (
|
|
"database/sql"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"time"
|
|
|
|
"crypto-miner-server/internal/models"
|
|
_ "modernc.org/sqlite"
|
|
)
|
|
|
|
type Database struct {
|
|
*sql.DB
|
|
}
|
|
|
|
func New(dataDir string) (*Database, error) {
|
|
dbPath := filepath.Join(dataDir, "miner.db")
|
|
|
|
// Ensure directory exists
|
|
os.MkdirAll(filepath.Dir(dbPath), 0755)
|
|
|
|
db, err := sql.Open("sqlite", dbPath+"?_pragma=journal_mode(WAL)&_pragma=busy_timeout(5000)")
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to open database: %w", err)
|
|
}
|
|
|
|
d := &Database{db}
|
|
if err := d.migrate(); err != nil {
|
|
return nil, fmt.Errorf("failed to migrate database: %w", err)
|
|
}
|
|
|
|
return d, nil
|
|
}
|
|
|
|
func (d *Database) migrate() error {
|
|
migrations := []string{
|
|
`CREATE TABLE IF NOT EXISTS agents (
|
|
id TEXT PRIMARY KEY,
|
|
name TEXT NOT NULL,
|
|
wallet TEXT NOT NULL DEFAULT '',
|
|
ip TEXT NOT NULL DEFAULT '',
|
|
version TEXT NOT NULL DEFAULT '',
|
|
status TEXT NOT NULL DEFAULT 'offline',
|
|
cpu_cores INTEGER NOT NULL DEFAULT 0,
|
|
memory_gb INTEGER NOT NULL DEFAULT 0,
|
|
last_seen DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
hashrate_15s REAL NOT NULL DEFAULT 0,
|
|
hashrate_1m REAL NOT NULL DEFAULT 0,
|
|
hashrate_15m REAL NOT NULL DEFAULT 0,
|
|
shares_total INTEGER NOT NULL DEFAULT 0,
|
|
shares_good INTEGER NOT NULL DEFAULT 0,
|
|
shares_bad INTEGER NOT NULL DEFAULT 0,
|
|
cpu_usage_pct REAL NOT NULL DEFAULT 0,
|
|
memory_usage_pct REAL NOT NULL DEFAULT 0,
|
|
uptime_seconds INTEGER NOT NULL DEFAULT 0
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS shares (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
agent_id TEXT NOT NULL,
|
|
job_id TEXT NOT NULL,
|
|
difficulty INTEGER NOT NULL DEFAULT 0,
|
|
accepted INTEGER NOT NULL DEFAULT 0,
|
|
hash TEXT NOT NULL DEFAULT '',
|
|
nonce TEXT NOT NULL DEFAULT '',
|
|
error TEXT NOT NULL DEFAULT '',
|
|
timestamp DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS hashrate_samples (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
agent_id TEXT NOT NULL,
|
|
hashrate REAL NOT NULL DEFAULT 0,
|
|
timestamp DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS jobs (
|
|
id TEXT PRIMARY KEY,
|
|
height INTEGER NOT NULL DEFAULT 0,
|
|
difficulty INTEGER NOT NULL DEFAULT 0,
|
|
block_template TEXT NOT NULL DEFAULT '',
|
|
seed_hash TEXT NOT NULL DEFAULT '',
|
|
target TEXT NOT NULL DEFAULT '',
|
|
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
|
|
)`,
|
|
`CREATE TABLE IF NOT EXISTS builds (
|
|
id TEXT PRIMARY KEY,
|
|
worker_name TEXT NOT NULL,
|
|
server_url TEXT NOT NULL,
|
|
wallet TEXT NOT NULL,
|
|
threads INTEGER NOT NULL DEFAULT 0,
|
|
file_size INTEGER NOT NULL DEFAULT 0,
|
|
file_path TEXT NOT NULL DEFAULT '',
|
|
created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
|
|
pool_host TEXT NOT NULL DEFAULT '',
|
|
pool_port INTEGER NOT NULL DEFAULT 0,
|
|
pool_tls INTEGER NOT NULL DEFAULT 0,
|
|
pool_pass TEXT NOT NULL DEFAULT ''
|
|
)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_shares_agent ON shares(agent_id)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_shares_timestamp ON shares(timestamp)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_hashrate_agent ON hashrate_samples(agent_id)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_hashrate_timestamp ON hashrate_samples(timestamp)`,
|
|
}
|
|
|
|
for _, m := range migrations {
|
|
if _, err := d.Exec(m); err != nil {
|
|
return fmt.Errorf("migration failed: %w\nSQL: %s", err, m)
|
|
}
|
|
}
|
|
|
|
// Best-effort schema upgrades for existing databases.
|
|
_, _ = d.Exec(`ALTER TABLE builds ADD COLUMN file_path TEXT NOT NULL DEFAULT ''`)
|
|
|
|
return nil
|
|
}
|
|
|
|
// Agent operations
|
|
|
|
func (d *Database) UpsertAgent(a *models.Agent) error {
|
|
query := `INSERT INTO agents (id, name, wallet, ip, version, status, cpu_cores, memory_gb, last_seen, created_at)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, COALESCE((SELECT created_at FROM agents WHERE id = ?), CURRENT_TIMESTAMP))
|
|
ON CONFLICT(id) DO UPDATE SET
|
|
name = excluded.name,
|
|
wallet = excluded.wallet,
|
|
ip = excluded.ip,
|
|
version = excluded.version,
|
|
status = excluded.status,
|
|
cpu_cores = excluded.cpu_cores,
|
|
memory_gb = excluded.memory_gb,
|
|
last_seen = excluded.last_seen`
|
|
_, err := d.Exec(query, a.ID, a.Name, a.Wallet, a.IP, a.Version, a.Status, a.CPUCores, a.MemoryGB, a.LastSeen, a.ID)
|
|
return err
|
|
}
|
|
|
|
func (d *Database) UpdateAgentStats(id string, hashrate15s, hashrate1m, hashrate15m float64, sharesTotal, sharesGood, sharesBad int, cpuPct, memPct float64, uptime int) error {
|
|
query := `UPDATE agents SET
|
|
hashrate_15s = ?, hashrate_1m = ?, hashrate_15m = ?,
|
|
shares_total = ?, shares_good = ?, shares_bad = ?,
|
|
cpu_usage_pct = ?, memory_usage_pct = ?, uptime_seconds = ?,
|
|
last_seen = CURRENT_TIMESTAMP, status = 'online'
|
|
WHERE id = ?`
|
|
_, err := d.Exec(query, hashrate15s, hashrate1m, hashrate15m, sharesTotal, sharesGood, sharesBad, cpuPct, memPct, uptime, id)
|
|
return err
|
|
}
|
|
|
|
func (d *Database) SetAgentOffline(id string) error {
|
|
_, err := d.Exec("UPDATE agents SET status = 'offline' WHERE id = ?", id)
|
|
return err
|
|
}
|
|
|
|
func (d *Database) GetAgent(id string) (*models.Agent, error) {
|
|
a := &models.Agent{}
|
|
query := `SELECT id, name, wallet, ip, version, status, cpu_cores, memory_gb, last_seen, created_at,
|
|
hashrate_15s, hashrate_1m, hashrate_15m, shares_total, shares_good, shares_bad, cpu_usage_pct, memory_usage_pct, uptime_seconds
|
|
FROM agents WHERE id = ?`
|
|
err := d.QueryRow(query, id).Scan(&a.ID, &a.Name, &a.Wallet, &a.IP, &a.Version, &a.Status,
|
|
&a.CPUCores, &a.MemoryGB, &a.LastSeen, &a.CreatedAt,
|
|
&a.Hashrate15s, &a.Hashrate1m, &a.Hashrate15m, &a.SharesTotal, &a.SharesGood, &a.SharesBad,
|
|
&a.CPUUsagePct, &a.MemoryUsagePct, &a.UptimeSeconds)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return a, nil
|
|
}
|
|
|
|
func (d *Database) ListAgents() ([]*models.Agent, error) {
|
|
query := `SELECT id, name, wallet, ip, version, status, cpu_cores, memory_gb, last_seen, created_at,
|
|
hashrate_15s, hashrate_1m, hashrate_15m, shares_total, shares_good, shares_bad, cpu_usage_pct, memory_usage_pct, uptime_seconds
|
|
FROM agents ORDER BY last_seen DESC`
|
|
rows, err := d.Query(query)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var agents []*models.Agent
|
|
for rows.Next() {
|
|
a := &models.Agent{}
|
|
if err := rows.Scan(&a.ID, &a.Name, &a.Wallet, &a.IP, &a.Version, &a.Status,
|
|
&a.CPUCores, &a.MemoryGB, &a.LastSeen, &a.CreatedAt,
|
|
&a.Hashrate15s, &a.Hashrate1m, &a.Hashrate15m, &a.SharesTotal, &a.SharesGood, &a.SharesBad,
|
|
&a.CPUUsagePct, &a.MemoryUsagePct, &a.UptimeSeconds); err != nil {
|
|
return nil, err
|
|
}
|
|
agents = append(agents, a)
|
|
}
|
|
return agents, nil
|
|
}
|
|
|
|
// Share operations
|
|
|
|
func (d *Database) InsertShare(s *models.Share) (int64, error) {
|
|
query := `INSERT INTO shares (agent_id, job_id, difficulty, accepted, hash, nonce, error, timestamp) VALUES (?, ?, ?, ?, ?, ?, ?, ?)`
|
|
res, err := d.Exec(query, s.AgentID, s.JobID, s.Difficulty, boolToInt(s.Accepted), s.Hash, s.Nonce, s.Error, s.Timestamp)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return res.LastInsertId()
|
|
}
|
|
|
|
func (d *Database) UpdateShareResult(id int64, accepted bool, errMsg string) error {
|
|
query := `UPDATE shares SET accepted = ?, error = ? WHERE id = ?`
|
|
_, err := d.Exec(query, boolToInt(accepted), errMsg, id)
|
|
return err
|
|
}
|
|
|
|
func (d *Database) GetRecentShares(limit int) ([]*models.Share, error) {
|
|
query := `SELECT id, agent_id, job_id, difficulty, accepted, hash, nonce, error, timestamp FROM shares ORDER BY timestamp DESC LIMIT ?`
|
|
rows, err := d.Query(query, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var shares []*models.Share
|
|
for rows.Next() {
|
|
s := &models.Share{}
|
|
var accepted int
|
|
if err := rows.Scan(&s.ID, &s.AgentID, &s.JobID, &s.Difficulty, &accepted, &s.Hash, &s.Nonce, &s.Error, &s.Timestamp); err != nil {
|
|
return nil, err
|
|
}
|
|
s.Accepted = accepted == 1
|
|
shares = append(shares, s)
|
|
}
|
|
return shares, nil
|
|
}
|
|
|
|
// Hashrate operations
|
|
|
|
func (d *Database) InsertHashrateSample(agentID string, hashrate float64) error {
|
|
_, err := d.Exec("INSERT INTO hashrate_samples (agent_id, hashrate, timestamp) VALUES (?, ?, ?)", agentID, hashrate, time.Now())
|
|
return err
|
|
}
|
|
|
|
func (d *Database) GetHashrateHistory(agentID string, limit int) ([]*models.HashrateSample, error) {
|
|
query := `SELECT id, agent_id, hashrate, timestamp FROM hashrate_samples WHERE agent_id = ? ORDER BY timestamp DESC LIMIT ?`
|
|
rows, err := d.Query(query, agentID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var samples []*models.HashrateSample
|
|
for rows.Next() {
|
|
s := &models.HashrateSample{}
|
|
if err := rows.Scan(&s.ID, &s.AgentID, &s.Hashrate, &s.Timestamp); err != nil {
|
|
return nil, err
|
|
}
|
|
samples = append(samples, s)
|
|
}
|
|
return samples, nil
|
|
}
|
|
|
|
// Build operations
|
|
|
|
func (d *Database) InsertBuild(b *models.BuildRecord) error {
|
|
_, err := d.Exec("INSERT INTO builds (id, worker_name, server_url, wallet, threads, file_size, file_path, created_at, pool_host, pool_port, pool_tls, pool_pass) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
|
|
b.ID, b.WorkerName, b.ServerURL, b.Wallet, b.Threads, b.FileSize, b.FilePath, b.CreatedAt,
|
|
b.PoolHost, b.PoolPort, b.PoolTLS, b.PoolPass)
|
|
return err
|
|
}
|
|
|
|
func (d *Database) GetBuild(id string) (*models.BuildRecord, error) {
|
|
b := &models.BuildRecord{}
|
|
err := d.QueryRow(`SELECT id, worker_name, server_url, wallet, threads, file_size, file_path, created_at, pool_host, pool_port, pool_tls, pool_pass FROM builds WHERE id = ?`, id).
|
|
Scan(&b.ID, &b.WorkerName, &b.ServerURL, &b.Wallet, &b.Threads, &b.FileSize, &b.FilePath, &b.CreatedAt, &b.PoolHost, &b.PoolPort, &b.PoolTLS, &b.PoolPass)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return b, nil
|
|
}
|
|
|
|
func (d *Database) ListBuilds(limit int) ([]*models.BuildRecord, error) {
|
|
query := `SELECT id, worker_name, server_url, wallet, threads, file_size, file_path, created_at, pool_host, pool_port, pool_tls, pool_pass FROM builds ORDER BY created_at DESC LIMIT ?`
|
|
rows, err := d.Query(query, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var builds []*models.BuildRecord
|
|
for rows.Next() {
|
|
b := &models.BuildRecord{}
|
|
if err := rows.Scan(&b.ID, &b.WorkerName, &b.ServerURL, &b.Wallet, &b.Threads, &b.FileSize, &b.FilePath, &b.CreatedAt,
|
|
&b.PoolHost, &b.PoolPort, &b.PoolTLS, &b.PoolPass); err != nil {
|
|
return nil, err
|
|
}
|
|
builds = append(builds, b)
|
|
}
|
|
return builds, nil
|
|
}
|
|
|
|
// Stats
|
|
|
|
type FleetStats struct {
|
|
TotalAgents int `json:"total_agents"`
|
|
OnlineAgents int `json:"online_agents"`
|
|
TotalHashrate float64 `json:"total_hashrate"`
|
|
TotalShares int `json:"total_shares"`
|
|
AcceptedShares int `json:"accepted_shares"`
|
|
RejectedShares int `json:"rejected_shares"`
|
|
AcceptRate float64 `json:"accept_rate"`
|
|
}
|
|
|
|
func (d *Database) GetFleetStats() (*FleetStats, error) {
|
|
stats := &FleetStats{}
|
|
|
|
err := d.QueryRow("SELECT COUNT(*) FROM agents").Scan(&stats.TotalAgents)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
err = d.QueryRow("SELECT COUNT(*) FROM agents WHERE status = 'online'").Scan(&stats.OnlineAgents)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
err = d.QueryRow("SELECT COALESCE(SUM(hashrate_15m), 0) FROM agents WHERE status = 'online'").Scan(&stats.TotalHashrate)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
err = d.QueryRow("SELECT COALESCE(SUM(shares_total), 0) FROM agents").Scan(&stats.TotalShares)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
err = d.QueryRow("SELECT COALESCE(SUM(shares_good), 0) FROM agents").Scan(&stats.AcceptedShares)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
err = d.QueryRow("SELECT COALESCE(SUM(shares_bad), 0) FROM agents").Scan(&stats.RejectedShares)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
if stats.TotalShares > 0 {
|
|
stats.AcceptRate = float64(stats.AcceptedShares) / float64(stats.TotalShares) * 100
|
|
}
|
|
|
|
return stats, nil
|
|
}
|
|
|
|
func boolToInt(b bool) int {
|
|
if b {
|
|
return 1
|
|
}
|
|
return 0
|
|
}
|