mirror of
https://tangled.org/evan.jarrett.net/at-container-registry
synced 2026-09-29 13:35:35 +00:00
`go fix` carries the modernize analyzers now, and the tree had drifted behind
them. This is the mechanical result, reviewed rather than trusted: the tool is
capable of rewriting code into something that no longer tests or does what it
did, so every non-test change was read individually and the concurrency-bearing
packages were re-run under -race.
Production code, four changes, all semantics-preserving:
- leases/manager.go: wg.Add(1) + go + defer wg.Done() becomes wg.Go. The
comment above that function turns on Add happening before the goroutine
starts, so that a Wait cannot return before the worker has run. wg.Go does
the Add synchronously on the calling goroutine, so the invariant it
describes still holds.
- auth/token/handler.go: strings.Fields -> strings.FieldsSeq, same splitting,
iterated rather than allocated.
- hold/gc/gc.go: a hand-written map copy -> maps.Copy.
- hold/pds/scan_broadcaster.go: three-clause loop -> range over int.
The rest are tests. The one worth naming is carstore_contention_test.go, where a
careless rewrite could have quietly stopped exercising contention: go fix
converted the reader and side-table goroutines to loopWG.Go but correctly
declined to touch the writer loop, which passes its index as a parameter. The
writer/reader/side-table shape and the stop channel are unchanged, so the test
still contends over the same carstore transactions.
Verified: go build for hold and appview, `make lint` 0 issues, the deploy and
credential-helper modules 0 issues, `make test` green across all 43 packages,
and -race green on leases, hold/pds, hold/gc and auth/token. The scanner module's
two lint findings are unchanged from HEAD and are in files go fix never touched.
Kept separate from the HTTP/2 commit so that one stays readable, and so this can
be reverted on its own if a modernization turns out to matter.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01TA9D4DjaLZTvzQ7dJbu4eg
2919 lines
100 KiB
Go
2919 lines
100 KiB
Go
package pds
|
|
|
|
import (
|
|
"context"
|
|
"crypto/rand"
|
|
"database/sql"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"net"
|
|
"net/http"
|
|
"net/url"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"atcr.io/pkg/atproto"
|
|
holddb "atcr.io/pkg/hold/db"
|
|
"atcr.io/pkg/s3"
|
|
lexutil "github.com/bluesky-social/indigo/lex/util"
|
|
"github.com/gorilla/websocket"
|
|
)
|
|
|
|
const (
|
|
// pendingReclaimAfter is how long a job may sit in 'pending' before the
|
|
// re-dispatch loop offers it to a scanner again. Enqueue and
|
|
// drainPendingJobs are both one-shot, so without a periodic re-offer a row
|
|
// whose only dispatch attempt failed waits for the next scanner connect.
|
|
pendingReclaimAfter = 1 * time.Minute
|
|
|
|
// pendingStaleAfter is how long a pending job counts as active for dispatch
|
|
// throttling. Assigned and processing jobs need no such bound — the ack and
|
|
// processing timeouts always resolve them — but 'pending' is otherwise
|
|
// unbounded, and one row stuck there froze proactive scanning for nine days.
|
|
pendingStaleAfter = 15 * time.Minute
|
|
|
|
// capacityStallWarnAfter is how long the dispatch loop waits with no
|
|
// capacity before it says so. A silent stall is what made this class of
|
|
// failure invisible until users noticed missing scans.
|
|
capacityStallWarnAfter = 10 * time.Minute
|
|
|
|
// scanningTimeout is how long a scan may run before the hold gives up on
|
|
// it. Measured from started_at — the moment a worker told us it picked the
|
|
// job up — so it budgets scanning and nothing else.
|
|
//
|
|
// It used to be measured from assigned_at, which the ack does not refresh.
|
|
// The scanner acks on receipt, off its WebSocket reader, and the job then
|
|
// waits in its own queue behind its workers, so the budget covered
|
|
// queueing: a 100-deep queue of no-op jobs drains in 16m40s and crosses
|
|
// ten minutes at position 59, and every job past that point was failed
|
|
// underneath a perfectly healthy scanner.
|
|
scanningTimeout = 10 * time.Minute
|
|
|
|
// queuedTimeout is the fallback budget for a job a scanner acked but never
|
|
// reported starting. A scanner built before the 'started' message exists
|
|
// never sends one, so for those the hold can only observe dispatch and this
|
|
// is the honest bound: long enough for a full default queue (100 jobs) of
|
|
// real scans to drain ahead of it, short enough that a connected-but-wedged
|
|
// scanner does not hold its share of dispatch capacity for the life of the
|
|
// process.
|
|
queuedTimeout = 60 * time.Minute
|
|
|
|
// reconnectGrace is how long a disconnected scanner's in-flight jobs stay
|
|
// its own.
|
|
//
|
|
// A dropped WebSocket is not evidence that a scanner stopped scanning: its
|
|
// worker pool never learns the socket went away and keeps going. Handing
|
|
// that work to another process immediately means two processes scan the
|
|
// same image and both report a verdict. A scanner keeps one identity for
|
|
// the life of its process and resumes its own rows on reconnect, so this
|
|
// only has to outlast a reconnect — the client retries every five seconds.
|
|
// A scanner that actually restarted comes back with a new identity and its
|
|
// old rows are reclaimed here instead, which is right: that work is gone.
|
|
reconnectGrace = 2 * time.Minute
|
|
|
|
// drainSendTimeout is how long the drain waits for room in a scanner's
|
|
// send buffer before giving the row back.
|
|
drainSendTimeout = 5 * time.Second
|
|
|
|
// activeJobsErrorBudget is how many consecutive database failures
|
|
// activeProactiveJobs answers with "assume busy" before it starts
|
|
// answering "assume idle" instead. See activeProactiveJobs.
|
|
activeJobsErrorBudget = 3
|
|
|
|
// defaultScannerCapacity is how many concurrent scans a scanner that
|
|
// declares nothing is assumed to run. Every scanner built before the
|
|
// workers parameter existed lands here, and one is exactly the depth the
|
|
// hold used to allow hold-wide, so an old scanner against a new hold
|
|
// behaves as it always did.
|
|
defaultScannerCapacity = 1
|
|
|
|
// maxScannerCapacity caps what a single connection may declare. The value
|
|
// arrives over the wire behind nothing but a shared secret and is used as
|
|
// a dispatch budget, so it is clamped rather than trusted.
|
|
maxScannerCapacity = 32
|
|
|
|
// scannerPingInterval is how often the hold pings an idle scanner
|
|
// connection.
|
|
scannerPingInterval = 30 * time.Second
|
|
|
|
// scannerPongWait is how long a connection may go without a pong (or any
|
|
// other frame) before the hold treats it as dead. Three missed pings.
|
|
scannerPongWait = 90 * time.Second
|
|
|
|
// scannerWriteWait bounds a single write to a scanner.
|
|
scannerWriteWait = 30 * time.Second
|
|
|
|
// scannerMaxMessageSize caps one frame from a scanner. gorilla's default is
|
|
// unlimited, so a broken or hostile peer can make the hold allocate
|
|
// whatever it likes.
|
|
//
|
|
// The ceiling has to be generous, because a result message carries the
|
|
// whole SBOM and the whole Grype report inline as JSON strings: exceeding
|
|
// the limit is not a truncation but a connection close, and a legitimate
|
|
// result that trips it becomes a permanent retry loop for that image. The
|
|
// largest report measured in this repo is 1.8 MB for 125 matches; 128 MiB
|
|
// is roughly seventy times the largest frame anyone has documented here,
|
|
// while still bounding a single allocation to something the host survives.
|
|
scannerMaxMessageSize = 128 << 20
|
|
|
|
// scannerResultTimeout bounds the blob uploads one result message costs.
|
|
scannerResultTimeout = 5 * time.Minute
|
|
|
|
// scannerRecordTimeout bounds the PDS record write a terminal message
|
|
// costs. It is deliberately a separate budget from scannerResultTimeout: a
|
|
// stalled S3 endpoint must not be able to spend the time the record write
|
|
// needs. The record is the only thing that stops the appview reporting the
|
|
// manifest as never scanned, so it is the last thing that may be skipped.
|
|
scannerRecordTimeout = 60 * time.Second
|
|
|
|
// scannerReuseTimeout bounds the "has this content already been stored"
|
|
// lookup one result message costs: one record read out of the CAR store and
|
|
// up to two S3 HEADs. It is carved out of scannerResultTimeout rather than
|
|
// given its own budget, because it is a substitute for the uploads that
|
|
// follow it. Kept short: every path out of it falls back to uploading, so
|
|
// giving up early costs bytes, not correctness.
|
|
scannerReuseTimeout = 15 * time.Second
|
|
|
|
// scannerStorageQueueDepth is how many terminal messages one connection may
|
|
// have waiting on storage before the reader has to wait for room. It is
|
|
// small on purpose: each queued message holds an entire SBOM and Grype
|
|
// report in memory, and the hold shares a 1 GiB host with the scanner.
|
|
scannerStorageQueueDepth = 8
|
|
|
|
// maxScannerInstanceID bounds the scanner-supplied identity that becomes
|
|
// assigned_to. Same reasoning: it is client input that ends up in a
|
|
// database column.
|
|
maxScannerInstanceID = 64
|
|
)
|
|
|
|
// Job origins. The proactive dispatch gate counts only proactive rows: push
|
|
// scans bypass the gate entirely (oci/xrpc.go calls Enqueue directly), so
|
|
// counting them meant a hold with steady pushes never dispatched a proactive
|
|
// scan at all. Rows that predate the column read as push, which errs toward
|
|
// dispatching rather than toward the stall this whole gate exists to avoid.
|
|
const (
|
|
originPush = "push"
|
|
originProactive = "proactive"
|
|
)
|
|
|
|
// ScanBroadcaster manages scanner WebSocket connections and dispatches scan jobs
|
|
// using a competing-consumer pattern. Jobs are persisted in SQLite and dispatched
|
|
// round-robin to connected scanners.
|
|
type ScanBroadcaster struct {
|
|
mu sync.RWMutex
|
|
subscribers []*ScanSubscriber
|
|
nextIdx int // Round-robin index for dispatch
|
|
db *sql.DB
|
|
holdDID string
|
|
holdEndpoint string
|
|
s3 *s3.S3Service
|
|
pds *HoldPDS
|
|
ackTimeout time.Duration
|
|
secret string // Shared secret for scanner authentication
|
|
ownsDB bool // true when this broadcaster opened the connection itself
|
|
|
|
// Proactive scan scheduling
|
|
rescanInterval time.Duration // Minimum interval between re-scans (0 = disabled)
|
|
stopCh chan struct{} // Signal to stop background goroutines
|
|
wg sync.WaitGroup // Wait for background goroutines to finish
|
|
// predecessorCache answers "has this hold been migrated into us?" without
|
|
// re-dialling. Only definitive answers belong here. The cache is never
|
|
// reset, so a false recorded from an unreachable hold would outlive the
|
|
// outage and stop that hold's manifests being scanned for the life of the
|
|
// process.
|
|
predecessorCache map[string]bool
|
|
|
|
// predecessorUnresolved holds the DIDs whose predecessor status could not
|
|
// be determined during the current discovery pass. It exists only so that
|
|
// one unreachable hold costs a single 5s timeout per pass rather than one
|
|
// per manifest, and it is cleared at the start of every pass so a hold that
|
|
// was down once is re-checked next time instead of being written off.
|
|
predecessorUnresolved map[string]bool
|
|
|
|
// predecessorMu guards both maps. Only the discovery goroutine reaches them
|
|
// today; the lock is what makes that a property of the code rather than of
|
|
// the current call graph.
|
|
predecessorMu sync.Mutex
|
|
relayEndpoints []string // Relay URLs for listReposByCollection (failover order)
|
|
relayStartIdx int // Rotates per discovery pass so load is shared across endpoints
|
|
relayStartMu sync.Mutex
|
|
|
|
// Work queues for proactive scanning (populated by discovery/stale goroutines)
|
|
unscannedQueue chan *scanCandidate // Medium priority: manifests with no scan record
|
|
staleQueue chan *scanCandidate // Low priority: scan records older than rescanInterval
|
|
inflight map[string]struct{} // Manifest digests currently queued or being scanned
|
|
inflightMu sync.Mutex
|
|
completionSignal chan struct{} // Signaled when a scan job completes (wakes dispatchLoop)
|
|
capacityFreed chan struct{} // Signaled when a scanner frees a slot (wakes reDispatchLoop)
|
|
discoverNow chan struct{} // Signaled to trigger an early discovery pass
|
|
|
|
// activeJobsErrs counts consecutive hasActiveJobs query failures, so a
|
|
// persistent database fault can fail open instead of freezing dispatch.
|
|
activeJobsErrs atomic.Int64
|
|
|
|
// WebSocket liveness knobs. Zero means "use the package default"; tests
|
|
// shrink them so a half-open connection can be provoked in milliseconds
|
|
// rather than in the minute and a half production waits.
|
|
pingInterval time.Duration
|
|
pongWait time.Duration
|
|
writeWait time.Duration
|
|
|
|
// resultTimeout bounds the off-reader work a terminal message triggers.
|
|
// Zero means the package default.
|
|
resultTimeout time.Duration
|
|
}
|
|
|
|
// pingEvery, pongDeadline, writeDeadline and resultDeadline apply the package
|
|
// defaults to the per-broadcaster knobs above.
|
|
func (sb *ScanBroadcaster) pingEvery() time.Duration {
|
|
if sb.pingInterval > 0 {
|
|
return sb.pingInterval
|
|
}
|
|
return scannerPingInterval
|
|
}
|
|
|
|
func (sb *ScanBroadcaster) pongDeadline() time.Duration {
|
|
if sb.pongWait > 0 {
|
|
return sb.pongWait
|
|
}
|
|
return scannerPongWait
|
|
}
|
|
|
|
func (sb *ScanBroadcaster) writeDeadline() time.Duration {
|
|
if sb.writeWait > 0 {
|
|
return sb.writeWait
|
|
}
|
|
return scannerWriteWait
|
|
}
|
|
|
|
func (sb *ScanBroadcaster) resultDeadline() time.Duration {
|
|
if sb.resultTimeout > 0 {
|
|
return sb.resultTimeout
|
|
}
|
|
return scannerResultTimeout
|
|
}
|
|
|
|
// ScanSubscriber represents a connected scanner WebSocket client
|
|
type ScanSubscriber struct {
|
|
conn *websocket.Conn
|
|
send chan *ScanJobEvent
|
|
id string // Scanner instance identity; also the assigned_to value
|
|
done chan struct{}
|
|
|
|
// capacity is how many scans this scanner runs at once, as declared on
|
|
// connect. It is the unit both the proactive dispatch depth and
|
|
// per-scanner admission control are counted in.
|
|
capacity int
|
|
|
|
// storage carries the work a terminal message triggers off the reader
|
|
// goroutine, and readerDone says when nothing more will be submitted.
|
|
//
|
|
// One goroutine drains storage, so messages are handled in the order they
|
|
// arrived. That matters: two messages about one job must not be applied out
|
|
// of order, and the alternative — a pool that hashes seq to a worker — buys
|
|
// parallelism the PDS cannot use anyway, since CreateScanRecord serialises
|
|
// on the hold's single per-uid repo lock.
|
|
storage chan func()
|
|
readerDone chan struct{}
|
|
}
|
|
|
|
// effectiveCapacity is capacity with the pre-declaration default applied.
|
|
func (s *ScanSubscriber) effectiveCapacity() int {
|
|
if s.capacity <= 0 {
|
|
return defaultScannerCapacity
|
|
}
|
|
if s.capacity > maxScannerCapacity {
|
|
return maxScannerCapacity
|
|
}
|
|
return s.capacity
|
|
}
|
|
|
|
// ScanJobEvent is the message sent from hold to scanner over WebSocket
|
|
type ScanJobEvent struct {
|
|
Type string `json:"type"` // Always "job"
|
|
Seq int64 `json:"seq"`
|
|
ManifestDigest string `json:"manifestDigest"`
|
|
Repository string `json:"repository"`
|
|
Tag string `json:"tag"`
|
|
UserDID string `json:"userDid"`
|
|
UserHandle string `json:"userHandle"`
|
|
HoldDID string `json:"holdDid"`
|
|
HoldEndpoint string `json:"holdEndpoint"`
|
|
Tier string `json:"tier"`
|
|
Config json.RawMessage `json:"config"`
|
|
Layers json.RawMessage `json:"layers"`
|
|
}
|
|
|
|
// ScannerMessage is a message received from scanner over WebSocket
|
|
type ScannerMessage struct {
|
|
Type string `json:"type"` // "ack", "started", "result", "error", "skipped"
|
|
Seq int64 `json:"seq"` // Job sequence number
|
|
SBOM string `json:"sbom,omitempty"`
|
|
VulnReport string `json:"vulnReport,omitempty"`
|
|
Summary *VulnerabilitySummary `json:"summary,omitempty"`
|
|
Error string `json:"error,omitempty"`
|
|
Reason string `json:"reason,omitempty"` // Populated for "skipped" messages
|
|
}
|
|
|
|
// VulnerabilitySummary contains counts of vulnerabilities by severity
|
|
type VulnerabilitySummary struct {
|
|
Critical int `json:"critical"`
|
|
High int `json:"high"`
|
|
Medium int `json:"medium"`
|
|
Low int `json:"low"`
|
|
Total int `json:"total"`
|
|
}
|
|
|
|
// NewScanBroadcaster creates a new scan job broadcaster
|
|
// dbPath should point to a SQLite database file (e.g., "/path/to/pds/db.sqlite3")
|
|
func NewScanBroadcaster(holdDID, holdEndpoint, secret string, relayEndpoints []string, dbPath string, s3svc *s3.S3Service, holdPDS *HoldPDS, rescanInterval time.Duration) (*ScanBroadcaster, error) {
|
|
// holddb.OpenLocalDB sets WAL journal mode and applies busy_timeout to
|
|
// every pooled connection. The one-shot "PRAGMA busy_timeout = 5000" that
|
|
// used to live here only configured whichever connection served it, so
|
|
// every other connection in the pool still failed immediately on a busy
|
|
// lock.
|
|
db, err := holddb.OpenLocalDB(dbPath)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to open scan jobs database: %w", err)
|
|
}
|
|
if err := db.Ping(); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("failed to ping scan jobs database: %w", err)
|
|
}
|
|
|
|
relayEndpoints = normalizeRelayEndpoints(relayEndpoints)
|
|
|
|
sb := &ScanBroadcaster{
|
|
subscribers: make([]*ScanSubscriber, 0),
|
|
db: db,
|
|
holdDID: holdDID,
|
|
holdEndpoint: holdEndpoint,
|
|
s3: s3svc,
|
|
pds: holdPDS,
|
|
ackTimeout: 5 * time.Minute,
|
|
secret: secret,
|
|
ownsDB: true,
|
|
rescanInterval: rescanInterval,
|
|
stopCh: make(chan struct{}),
|
|
predecessorCache: make(map[string]bool),
|
|
relayEndpoints: relayEndpoints,
|
|
unscannedQueue: make(chan *scanCandidate, 500),
|
|
staleQueue: make(chan *scanCandidate, 200),
|
|
inflight: make(map[string]struct{}),
|
|
completionSignal: make(chan struct{}, 1),
|
|
capacityFreed: make(chan struct{}, 1),
|
|
discoverNow: make(chan struct{}, 1),
|
|
}
|
|
|
|
if err := sb.initSchema(); err != nil {
|
|
db.Close()
|
|
return nil, fmt.Errorf("failed to initialize scan_jobs schema: %w", err)
|
|
}
|
|
sb.reconcileOnBoot()
|
|
|
|
// Start re-dispatch loop for timed-out jobs
|
|
sb.wg.Add(1)
|
|
go sb.reDispatchLoop()
|
|
|
|
// Start proactive scan loops if rescan interval is configured
|
|
if rescanInterval > 0 {
|
|
sb.wg.Add(3)
|
|
go sb.discoveryLoop()
|
|
go sb.staleScanLoop()
|
|
go sb.dispatchLoop()
|
|
slog.Info("Proactive scan scheduler started", "rescanInterval", rescanInterval, "relayEndpoints", relayEndpoints)
|
|
}
|
|
|
|
return sb, nil
|
|
}
|
|
|
|
// NewScanBroadcasterWithDB creates a scan job broadcaster using an existing *sql.DB connection.
|
|
// The caller is responsible for the DB lifecycle.
|
|
func NewScanBroadcasterWithDB(holdDID, holdEndpoint, secret string, relayEndpoints []string, db *sql.DB, s3svc *s3.S3Service, holdPDS *HoldPDS, rescanInterval time.Duration) (*ScanBroadcaster, error) {
|
|
relayEndpoints = normalizeRelayEndpoints(relayEndpoints)
|
|
|
|
sb := &ScanBroadcaster{
|
|
subscribers: make([]*ScanSubscriber, 0),
|
|
db: db,
|
|
holdDID: holdDID,
|
|
holdEndpoint: holdEndpoint,
|
|
s3: s3svc,
|
|
pds: holdPDS,
|
|
ackTimeout: 5 * time.Minute,
|
|
secret: secret,
|
|
ownsDB: false,
|
|
rescanInterval: rescanInterval,
|
|
stopCh: make(chan struct{}),
|
|
predecessorCache: make(map[string]bool),
|
|
relayEndpoints: relayEndpoints,
|
|
unscannedQueue: make(chan *scanCandidate, 500),
|
|
staleQueue: make(chan *scanCandidate, 200),
|
|
inflight: make(map[string]struct{}),
|
|
completionSignal: make(chan struct{}, 1),
|
|
capacityFreed: make(chan struct{}, 1),
|
|
discoverNow: make(chan struct{}, 1),
|
|
}
|
|
|
|
if err := sb.initSchema(); err != nil {
|
|
return nil, fmt.Errorf("failed to initialize scan_jobs schema: %w", err)
|
|
}
|
|
sb.reconcileOnBoot()
|
|
|
|
sb.wg.Add(1)
|
|
go sb.reDispatchLoop()
|
|
|
|
if rescanInterval > 0 {
|
|
sb.wg.Add(3)
|
|
go sb.discoveryLoop()
|
|
go sb.staleScanLoop()
|
|
go sb.dispatchLoop()
|
|
slog.Info("Proactive scan scheduler started", "rescanInterval", rescanInterval, "relayEndpoints", relayEndpoints)
|
|
}
|
|
|
|
return sb, nil
|
|
}
|
|
|
|
// normalizeRelayEndpoints drops empty entries and falls back to a sensible default
|
|
// if the resulting list is empty. The default mirrors the appview backfill default.
|
|
func normalizeRelayEndpoints(endpoints []string) []string {
|
|
out := make([]string, 0, len(endpoints))
|
|
for _, e := range endpoints {
|
|
if e != "" {
|
|
out = append(out, e)
|
|
}
|
|
}
|
|
if len(out) == 0 {
|
|
return []string{
|
|
"https://relay1.us-east.bsky.network",
|
|
"https://relay1.us-west.bsky.network",
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// initSchema creates the scan_jobs table if it doesn't exist
|
|
func (sb *ScanBroadcaster) initSchema() error {
|
|
// Executed individually for go-libsql compatibility
|
|
if _, err := sb.db.Exec(`CREATE TABLE IF NOT EXISTS scan_jobs (
|
|
seq INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
manifest_digest TEXT NOT NULL,
|
|
repository TEXT NOT NULL,
|
|
tag TEXT,
|
|
user_did TEXT NOT NULL,
|
|
user_handle TEXT,
|
|
hold_did TEXT NOT NULL,
|
|
hold_endpoint TEXT NOT NULL,
|
|
tier TEXT NOT NULL DEFAULT 'deckhand',
|
|
config_json TEXT NOT NULL,
|
|
layers_json TEXT NOT NULL,
|
|
status TEXT NOT NULL DEFAULT 'pending',
|
|
assigned_to TEXT,
|
|
assigned_at TIMESTAMP,
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
|
|
completed_at TIMESTAMP
|
|
)`); err != nil {
|
|
return err
|
|
}
|
|
|
|
// Columns added after the table shipped. CREATE TABLE IF NOT EXISTS does
|
|
// nothing for a database that already has the old shape, so every one of
|
|
// these has to be added on its own.
|
|
added := []struct{ name, def string }{
|
|
// Which dispatcher created the row. The proactive gate counts only its
|
|
// own work; existing rows default to push so they never throttle it.
|
|
{"origin", "TEXT NOT NULL DEFAULT 'push'"},
|
|
// When a worker reported it actually started scanning, as opposed to
|
|
// when the row was handed out. NULL means the scanner never said.
|
|
{"started_at", "TIMESTAMP"},
|
|
// When the scanner holding this row dropped its connection. NULL means
|
|
// it is connected, or the row was never assigned.
|
|
{"disconnected_at", "TIMESTAMP"},
|
|
}
|
|
for _, col := range added {
|
|
if err := sb.ensureColumn("scan_jobs", col.name, col.def); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
indexes := []string{
|
|
`CREATE INDEX IF NOT EXISTS idx_scan_jobs_status ON scan_jobs(status)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_scan_jobs_assigned ON scan_jobs(assigned_to, status)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_scan_jobs_origin_status ON scan_jobs(origin, status)`,
|
|
}
|
|
for _, stmt := range indexes {
|
|
if _, err := sb.db.Exec(stmt); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// reconcileOnBoot treats every job still assigned or processing as belonging to
|
|
// a disconnected scanner, because it does: whatever connections held them died
|
|
// with the previous process.
|
|
//
|
|
// Nothing used to reconcile at boot at all. Rows left mid-flight by a restart
|
|
// sat holding dispatch capacity until their own deadlines fired, which is up
|
|
// to an hour for a job that was queued inside a scanner. Marking them puts
|
|
// them under the same two-minute grace a live disconnect gets: a scanner whose
|
|
// workers are still running them redials with the same identity and resumes
|
|
// them, and one that is not comes back to find them re-offered.
|
|
func (sb *ScanBroadcaster) reconcileOnBoot() {
|
|
res, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET disconnected_at = ?
|
|
WHERE status IN ('assigned', 'processing') AND disconnected_at IS NULL
|
|
`, time.Now())
|
|
if err != nil {
|
|
slog.Error("Failed to reconcile in-flight scan jobs at boot", "error", err)
|
|
return
|
|
}
|
|
if n, err := res.RowsAffected(); err == nil && n > 0 {
|
|
slog.Info("Marked in-flight scan jobs from the previous process",
|
|
"jobs", n, "reclaimAfter", reconnectGrace)
|
|
}
|
|
}
|
|
|
|
// ensureColumn adds a column if the table does not already have it. SQLite has
|
|
// no ADD COLUMN IF NOT EXISTS, and the error text for a duplicate is not
|
|
// something worth matching on.
|
|
func (sb *ScanBroadcaster) ensureColumn(table, column, definition string) error {
|
|
rows, err := sb.db.Query(fmt.Sprintf("PRAGMA table_info(%s)", table))
|
|
if err != nil {
|
|
return fmt.Errorf("inspect %s: %w", table, err)
|
|
}
|
|
present := false
|
|
for rows.Next() {
|
|
var (
|
|
cid int
|
|
name string
|
|
ctype sql.NullString
|
|
notNull sql.NullInt64
|
|
defaultVal sql.NullString
|
|
pk sql.NullInt64
|
|
)
|
|
if err := rows.Scan(&cid, &name, &ctype, ¬Null, &defaultVal, &pk); err != nil {
|
|
rows.Close()
|
|
return fmt.Errorf("inspect %s: %w", table, err)
|
|
}
|
|
if name == column {
|
|
present = true
|
|
}
|
|
}
|
|
rows.Close()
|
|
if err := rows.Err(); err != nil {
|
|
return fmt.Errorf("inspect %s: %w", table, err)
|
|
}
|
|
if present {
|
|
return nil
|
|
}
|
|
|
|
if _, err := sb.db.Exec(fmt.Sprintf("ALTER TABLE %s ADD COLUMN %s %s", table, column, definition)); err != nil {
|
|
return fmt.Errorf("add %s.%s: %w", table, column, err)
|
|
}
|
|
slog.Info("Added scan job column", "table", table, "column", column)
|
|
return nil
|
|
}
|
|
|
|
// Enqueue inserts a push-triggered scan job and dispatches it.
|
|
//
|
|
// Push scans are not subject to the proactive dispatch gate — a user who just
|
|
// pushed is waiting for the answer — but they are still subject to
|
|
// per-scanner admission control, so a job that nobody has room for now waits
|
|
// on the hold's queue rather than in a scanner's.
|
|
func (sb *ScanBroadcaster) Enqueue(job *ScanJobEvent) error {
|
|
return sb.enqueue(job, originPush)
|
|
}
|
|
|
|
// enqueue inserts a scan job into SQLite and dispatches to an available scanner.
|
|
func (sb *ScanBroadcaster) enqueue(job *ScanJobEvent, origin string) error {
|
|
job.Type = "job"
|
|
job.HoldDID = sb.holdDID
|
|
job.HoldEndpoint = sb.holdEndpoint
|
|
|
|
// Track in-flight to prevent duplicate proactive scans
|
|
sb.addInflight(job.ManifestDigest)
|
|
|
|
// Insert into database
|
|
result, err := sb.db.Exec(`
|
|
INSERT INTO scan_jobs (manifest_digest, repository, tag, user_did, user_handle, hold_did, hold_endpoint, tier, config_json, layers_json, status, origin)
|
|
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending', ?)
|
|
`, job.ManifestDigest, job.Repository, job.Tag, job.UserDID, job.UserHandle, job.HoldDID, job.HoldEndpoint, job.Tier, string(job.Config), string(job.Layers), origin)
|
|
if err != nil {
|
|
sb.removeInflight(job.ManifestDigest)
|
|
return fmt.Errorf("failed to insert scan job: %w", err)
|
|
}
|
|
|
|
seq, err := result.LastInsertId()
|
|
if err != nil {
|
|
sb.removeInflight(job.ManifestDigest)
|
|
return fmt.Errorf("failed to get job seq: %w", err)
|
|
}
|
|
job.Seq = seq
|
|
|
|
slog.Info("Scan job enqueued",
|
|
"seq", seq,
|
|
"repository", job.Repository,
|
|
"tag", job.Tag,
|
|
"tier", job.Tier,
|
|
"origin", origin)
|
|
|
|
// Try to dispatch immediately
|
|
sb.dispatchJob(job)
|
|
|
|
return nil
|
|
}
|
|
|
|
// Subscribe adds a new scanner WebSocket subscriber and drains pending jobs to it.
|
|
//
|
|
// instanceID is the scanner's own identity, stable for the life of its process
|
|
// and sent on every connect. It becomes the subscriber id and therefore the
|
|
// assigned_to value, which is what lets a scanner that briefly lost its socket
|
|
// resume the jobs its workers never stopped running. A scanner that declares
|
|
// nothing gets a per-connection id, which is the old behaviour: its in-flight
|
|
// work is not resumable, only reclaimable.
|
|
//
|
|
// capacity is how many scans the scanner runs at once (its worker count).
|
|
func (sb *ScanBroadcaster) Subscribe(conn *websocket.Conn, cursor int64, instanceID string, capacity int) *ScanSubscriber {
|
|
id := sb.subscriberID(instanceID)
|
|
sub := &ScanSubscriber{
|
|
conn: conn,
|
|
send: make(chan *ScanJobEvent, 20),
|
|
id: id,
|
|
done: make(chan struct{}),
|
|
storage: make(chan func(), scannerStorageQueueDepth),
|
|
readerDone: make(chan struct{}),
|
|
}
|
|
sub.capacity = capacity
|
|
|
|
// Before anything is dispatched: reclaim whatever this instance was
|
|
// holding when it dropped, so the drain below counts it against the
|
|
// scanner's capacity rather than treating it as idle.
|
|
sb.resumeInstance(id)
|
|
|
|
sb.mu.Lock()
|
|
sb.subscribers = append(sb.subscribers, sub)
|
|
total := len(sb.subscribers)
|
|
sb.mu.Unlock()
|
|
|
|
slog.Info("Scanner subscribed",
|
|
"id", id,
|
|
"remote", conn.RemoteAddr(),
|
|
"cursor", cursor,
|
|
"capacity", sub.effectiveCapacity(),
|
|
"totalSubscribers", total)
|
|
|
|
// Start writer goroutine (sends jobs to scanner, and pings it)
|
|
go sb.handleWriter(sub)
|
|
|
|
// Start the storage goroutine before the reader, so the reader always has
|
|
// somewhere to hand a terminal message.
|
|
go sb.handleStorage(sub)
|
|
|
|
// Start reader goroutine (receives acks/results/errors from scanner)
|
|
go sb.handleReader(sub)
|
|
|
|
// Drain pending and timed-out jobs from database
|
|
go sb.drainPendingJobs(sub, cursor)
|
|
|
|
// Trigger early discovery pass so missed scans are found quickly
|
|
sb.triggerDiscovery()
|
|
|
|
return sub
|
|
}
|
|
|
|
// subscriberID turns a scanner-declared instance identity into the id used for
|
|
// assigned_to, falling back to a random per-connection id.
|
|
//
|
|
// The value arrives over the wire, so it is bounded and restricted to
|
|
// characters that read cleanly in a log line and a database column. A
|
|
// duplicate is refused rather than shared: two processes answering to one id
|
|
// would each accept the other's jobs, which is precisely the confusion the
|
|
// ownership guards exist to prevent.
|
|
func (sb *ScanBroadcaster) subscriberID(instanceID string) string {
|
|
if instanceID == "" {
|
|
return generateSubscriberID()
|
|
}
|
|
if len(instanceID) > maxScannerInstanceID || !isSafeInstanceID(instanceID) {
|
|
slog.Warn("Scanner declared an unusable instance id, assigning one",
|
|
"declared", instanceID)
|
|
return generateSubscriberID()
|
|
}
|
|
|
|
sb.mu.RLock()
|
|
defer sb.mu.RUnlock()
|
|
for _, existing := range sb.subscribers {
|
|
if existing.id == instanceID {
|
|
slog.Warn("Two scanners declared the same instance id; assigning one",
|
|
"instanceId", instanceID)
|
|
return generateSubscriberID()
|
|
}
|
|
}
|
|
return instanceID
|
|
}
|
|
|
|
func isSafeInstanceID(id string) bool {
|
|
for _, r := range id {
|
|
switch {
|
|
case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9':
|
|
case r == '-', r == '_', r == '.':
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
// resumeInstance hands a reconnecting scanner back the jobs it was holding
|
|
// when its connection dropped.
|
|
//
|
|
// Clearing disconnected_at is all it takes: the rows never left the scanner,
|
|
// because Unsubscribe marks a disconnect rather than acting on it. The
|
|
// scanner's workers carried on through the outage, so returning the rows to
|
|
// the pool would have had another process scan the same images.
|
|
func (sb *ScanBroadcaster) resumeInstance(id string) {
|
|
res, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET disconnected_at = NULL
|
|
WHERE assigned_to = ? AND disconnected_at IS NOT NULL
|
|
AND status IN ('assigned', 'processing')
|
|
`, id)
|
|
if err != nil {
|
|
slog.Error("Failed to resume jobs for reconnecting scanner",
|
|
"subscriberId", id, "error", err)
|
|
return
|
|
}
|
|
if n, err := res.RowsAffected(); err == nil && n > 0 {
|
|
slog.Info("Scanner reconnected and resumed its in-flight jobs",
|
|
"subscriberId", id, "jobs", n)
|
|
}
|
|
}
|
|
|
|
// Unsubscribe removes a scanner subscriber and makes its jobs re-dispatchable
|
|
func (sb *ScanBroadcaster) Unsubscribe(sub *ScanSubscriber) {
|
|
sb.mu.Lock()
|
|
defer sb.mu.Unlock()
|
|
|
|
found := false
|
|
for i, s := range sb.subscribers {
|
|
if s == sub {
|
|
sb.subscribers = append(sb.subscribers[:i], sb.subscribers[i+1:]...)
|
|
found = true
|
|
break
|
|
}
|
|
}
|
|
|
|
// A dropped scanner unwinds both handleWriter (on write error) and
|
|
// handleReader (in its defer), and each calls Unsubscribe. Everything below
|
|
// must happen exactly once: re-running the UPDATE would unassign jobs a
|
|
// replacement scanner had already picked up, and closing done twice panics.
|
|
if !found {
|
|
return
|
|
}
|
|
|
|
// Mark this scanner's jobs as belonging to a disconnected scanner, and stop
|
|
// there.
|
|
//
|
|
// This used to flip them straight back to 'pending'. A dropped WebSocket
|
|
// tells the hold nothing about what the scanner is doing: its worker pool
|
|
// never learns the socket went away, so it keeps downloading layers and
|
|
// running Syft on jobs the hold has just put back in the pool. With one
|
|
// scanner that was a duplicate against itself. With several it is a
|
|
// duplicate nothing can dedupe — the next process to connect drains the
|
|
// rows out from under a scanner that is still mid-scan, and both report a
|
|
// verdict for the same image.
|
|
//
|
|
// So the disconnect is recorded and the work is left alone. The same
|
|
// instance reconnecting resumes it (resumeInstance); a scanner that does
|
|
// not come back inside reconnectGrace has it reclaimed by
|
|
// reDispatchTimedOut. Rows that were never handed over — still 'pending'
|
|
// but stamped with this subscriber — are released outright, since nothing
|
|
// is running them.
|
|
if _, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'pending', assigned_to = NULL, assigned_at = NULL, disconnected_at = NULL
|
|
WHERE assigned_to = ? AND status = 'pending'
|
|
`, sub.id); err != nil {
|
|
slog.Error("Failed to release undispatched jobs from disconnected scanner",
|
|
"subscriberId", sub.id,
|
|
"error", err)
|
|
}
|
|
if _, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET disconnected_at = ?
|
|
WHERE assigned_to = ? AND status IN ('assigned', 'processing') AND disconnected_at IS NULL
|
|
`, time.Now(), sub.id); err != nil {
|
|
slog.Error("Failed to mark jobs from disconnected scanner",
|
|
"subscriberId", sub.id,
|
|
"error", err)
|
|
}
|
|
|
|
// Close done, never send. drainPendingJobs runs in its own goroutine
|
|
// without holding sb.mu, so closing send here would race an in-flight
|
|
// send and panic the process unrecoverably.
|
|
close(sub.done)
|
|
|
|
slog.Info("Scanner unsubscribed",
|
|
"id", sub.id,
|
|
"totalSubscribers", len(sb.subscribers))
|
|
}
|
|
|
|
// dispatchJob hands a job to the connected scanner with the most room for it.
|
|
//
|
|
// Two things changed here when more than one scanner became a supported
|
|
// deployment. Selection is by spare capacity rather than by position, because
|
|
// plain round-robin hands work to a saturated scanner as readily as an idle
|
|
// one, and with heterogeneous processes that is most of the fleet's work going
|
|
// to the wrong place. And a job nobody has room for stays 'pending' instead of
|
|
// being pushed into a scanner's own queue: the hold cannot see into that
|
|
// queue, so anything sitting in it is work it cannot schedule, cannot
|
|
// re-route to a process that freed up first, and cannot put a meaningful
|
|
// deadline on.
|
|
func (sb *ScanBroadcaster) dispatchJob(job *ScanJobEvent) {
|
|
sb.mu.Lock()
|
|
defer sb.mu.Unlock()
|
|
|
|
if len(sb.subscribers) == 0 {
|
|
slog.Debug("No scanners connected, job will wait in queue", "seq", job.Seq)
|
|
return
|
|
}
|
|
|
|
sub := sb.selectSubscriberLocked()
|
|
if sub == nil {
|
|
slog.Debug("Every scanner is at capacity, job stays pending", "seq", job.Seq)
|
|
return
|
|
}
|
|
|
|
// Mark as assigned in database
|
|
res, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'assigned', assigned_to = ?, assigned_at = ?
|
|
WHERE seq = ? AND status = 'pending'
|
|
`, sub.id, time.Now(), job.Seq)
|
|
if err != nil {
|
|
slog.Error("Failed to assign scan job", "seq", job.Seq, "error", err)
|
|
return
|
|
}
|
|
// Two paths hand out pending rows now (drainPendingJobs on connect, and the
|
|
// re-dispatch loop), so a row that is no longer pending was claimed by the
|
|
// other one. Sending it anyway would scan it twice.
|
|
if n, err := res.RowsAffected(); err == nil && n == 0 {
|
|
slog.Debug("Scan job no longer pending, skipping dispatch", "seq", job.Seq)
|
|
return
|
|
}
|
|
|
|
// Send to subscriber
|
|
select {
|
|
case sub.send <- job:
|
|
slog.Info("Scan job dispatched",
|
|
"seq", job.Seq,
|
|
"repository", job.Repository,
|
|
"subscriberId", sub.id)
|
|
default:
|
|
slog.Warn("Scanner buffer full, re-marking job as pending",
|
|
"seq", job.Seq,
|
|
"subscriberId", sub.id)
|
|
// Guarded, for the same reason the assign above is. Between our UPDATE
|
|
// and here another dispatcher can have taken the row; an unguarded
|
|
// reset would return a job that scanner is already holding to the
|
|
// pool, and a third scanner would be handed it.
|
|
sb.releaseUndelivered(sub, job.Seq)
|
|
}
|
|
}
|
|
|
|
// selectSubscriberLocked picks the connected scanner with the most spare
|
|
// capacity, or nil when every one of them is full. Caller holds sb.mu.
|
|
//
|
|
// Ties resolve in round-robin order from nextIdx, so a fleet of equal, idle
|
|
// scanners still gets work spread evenly across it.
|
|
func (sb *ScanBroadcaster) selectSubscriberLocked() *ScanSubscriber {
|
|
n := len(sb.subscribers)
|
|
if n == 0 {
|
|
return nil
|
|
}
|
|
|
|
load, ok := sb.subscriberLoads()
|
|
if !ok {
|
|
// The load query is the only thing that can say a scanner is full, so
|
|
// without it admission control has nothing to stand on. Fall back to
|
|
// plain round-robin: over-dispatching is worse than it was, but it is
|
|
// far better than dispatching nothing while the database misbehaves.
|
|
sub := sb.subscribers[sb.nextIdx%n]
|
|
sb.nextIdx++
|
|
return sub
|
|
}
|
|
|
|
best := -1
|
|
bestScore := 0.0
|
|
for i := range n {
|
|
idx := (sb.nextIdx + i) % n
|
|
sub := sb.subscribers[idx]
|
|
capacity := sub.effectiveCapacity()
|
|
outstanding := load[sub.id]
|
|
if outstanding >= capacity {
|
|
continue
|
|
}
|
|
// Fraction of the scanner's capacity already committed, so a
|
|
// four-worker process holding one job outranks a one-worker process
|
|
// holding none only when it is genuinely emptier.
|
|
score := float64(outstanding) / float64(capacity)
|
|
if best < 0 || score < bestScore {
|
|
best, bestScore = idx, score
|
|
}
|
|
}
|
|
if best < 0 {
|
|
return nil
|
|
}
|
|
|
|
sb.nextIdx = best + 1
|
|
return sb.subscribers[best]
|
|
}
|
|
|
|
// subscriberLoads counts the jobs each scanner is currently holding. The
|
|
// second return reports whether the count is usable; callers must not treat a
|
|
// failed query as "everyone is idle".
|
|
func (sb *ScanBroadcaster) subscriberLoads() (map[string]int, bool) {
|
|
rows, err := sb.db.Query(`
|
|
SELECT assigned_to, COUNT(*) FROM scan_jobs
|
|
WHERE status IN ('assigned', 'processing') AND assigned_to IS NOT NULL
|
|
GROUP BY assigned_to
|
|
`)
|
|
if err != nil {
|
|
slog.Error("Failed to count per-scanner load", "error", err)
|
|
return nil, false
|
|
}
|
|
defer rows.Close()
|
|
|
|
load := make(map[string]int)
|
|
for rows.Next() {
|
|
var id string
|
|
var n int
|
|
if err := rows.Scan(&id, &n); err != nil {
|
|
slog.Error("Failed to scan per-scanner load row", "error", err)
|
|
return nil, false
|
|
}
|
|
load[id] = n
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
slog.Error("Failed to read per-scanner load", "error", err)
|
|
return nil, false
|
|
}
|
|
return load, true
|
|
}
|
|
|
|
// handleWriter sends jobs to a scanner over its WebSocket connection
|
|
func (sb *ScanBroadcaster) handleWriter(sub *ScanSubscriber) {
|
|
defer sub.conn.Close()
|
|
|
|
// The writer owns every write on this connection, which is what makes the
|
|
// keepalive possible: gorilla permits one writer at a time, so the ping has
|
|
// to come from here rather than from a timer of its own.
|
|
ping := time.NewTicker(sb.pingEvery())
|
|
defer ping.Stop()
|
|
|
|
// done is closed by Unsubscribe, not here — ranging over send would block
|
|
// forever now that nothing closes it. handleReader always unsubscribes on
|
|
// its way out, so this goroutine is guaranteed to be released.
|
|
for {
|
|
var job *ScanJobEvent
|
|
select {
|
|
case <-sub.done:
|
|
return
|
|
case <-ping.C:
|
|
// Every write now carries a deadline. Without one a write into a
|
|
// socket whose peer has stopped reading blocks until the kernel
|
|
// gives up, which for a half-open connection is on the order of
|
|
// hours, and this goroutine is what the whole connection's
|
|
// liveness rests on.
|
|
if err := sub.conn.WriteControl(websocket.PingMessage, nil,
|
|
time.Now().Add(sb.writeDeadline())); err != nil {
|
|
slog.Warn("Failed to ping scanner, dropping the connection",
|
|
"subscriberId", sub.id, "error", err)
|
|
sb.Unsubscribe(sub)
|
|
return
|
|
}
|
|
continue
|
|
case job = <-sub.send:
|
|
}
|
|
|
|
data, err := json.Marshal(job)
|
|
if err != nil {
|
|
slog.Error("Failed to marshal scan job", "seq", job.Seq, "error", err)
|
|
continue
|
|
}
|
|
|
|
if err := sub.conn.SetWriteDeadline(time.Now().Add(sb.writeDeadline())); err != nil {
|
|
slog.Error("Failed to set write deadline for scan job",
|
|
"seq", job.Seq, "subscriberId", sub.id, "error", err)
|
|
sb.Unsubscribe(sub)
|
|
return
|
|
}
|
|
if err := sub.conn.WriteMessage(websocket.TextMessage, data); err != nil {
|
|
slog.Error("Failed to write scan job to WebSocket",
|
|
"seq", job.Seq,
|
|
"subscriberId", sub.id,
|
|
"error", err)
|
|
sb.Unsubscribe(sub)
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
// handleReader receives ack/result/error messages from a scanner.
|
|
//
|
|
// Two rules hold here, and they are the same rule seen from two sides: this
|
|
// goroutine must never stop reading, and what it reads must prove the scanner
|
|
// is alive.
|
|
//
|
|
// It never stops reading because the storage a terminal message triggers — two
|
|
// S3 uploads and a CAR commit — is handed to handleStorage instead of being run
|
|
// inline. It used to run here, and while it ran every other job's ack, started,
|
|
// result and error on that connection sat unread in the socket buffer, with no
|
|
// deadline on any of it.
|
|
//
|
|
// Liveness comes from the read deadline, refreshed by a pong and by any other
|
|
// frame. Nothing used to ask: a connection that was open at this end and gone
|
|
// at the other stayed in sb.subscribers for the life of the process, holding
|
|
// its advertised worker count out of the dispatch budget and winning jobs that
|
|
// could only ever time out. The deadline is three ping intervals, so it takes
|
|
// three unanswered pings to condemn a connection.
|
|
func (sb *ScanBroadcaster) handleReader(sub *ScanSubscriber) {
|
|
defer func() {
|
|
// Before Unsubscribe, so handleStorage learns that nothing more will be
|
|
// submitted only once this goroutine has genuinely stopped submitting.
|
|
close(sub.readerDone)
|
|
sb.Unsubscribe(sub)
|
|
}()
|
|
|
|
sub.conn.SetReadLimit(scannerMaxMessageSize)
|
|
refreshRead := func() {
|
|
if err := sub.conn.SetReadDeadline(time.Now().Add(sb.pongDeadline())); err != nil {
|
|
slog.Debug("Failed to set scanner read deadline",
|
|
"subscriberId", sub.id, "error", err)
|
|
}
|
|
}
|
|
refreshRead()
|
|
sub.conn.SetPongHandler(func(string) error {
|
|
refreshRead()
|
|
return nil
|
|
})
|
|
|
|
for {
|
|
_, data, err := sub.conn.ReadMessage()
|
|
if err != nil {
|
|
switch {
|
|
case isReadTimeout(err):
|
|
slog.Warn("Scanner stopped answering, dropping the connection",
|
|
"subscriberId", sub.id,
|
|
"silentFor", sb.pongDeadline(),
|
|
"error", err)
|
|
case websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseNormalClosure):
|
|
slog.Error("Scanner WebSocket read error",
|
|
"subscriberId", sub.id,
|
|
"error", err)
|
|
}
|
|
return
|
|
}
|
|
// Any frame is proof of life, not just a pong.
|
|
refreshRead()
|
|
|
|
var msg ScannerMessage
|
|
if err := json.Unmarshal(data, &msg); err != nil {
|
|
slog.Error("Failed to unmarshal scanner message",
|
|
"subscriberId", sub.id,
|
|
"error", err)
|
|
continue
|
|
}
|
|
|
|
switch msg.Type {
|
|
case "ack":
|
|
// One guarded UPDATE. Cheap enough to stay on the reader, and
|
|
// keeping it here means an ack is never delayed behind a scan
|
|
// record write.
|
|
sb.handleAck(sub, msg.Seq)
|
|
case "started":
|
|
sb.handleStarted(sub, msg.Seq)
|
|
case "result":
|
|
sb.submitStorage(sub, func() { sb.handleResult(sub, msg) })
|
|
case "error":
|
|
sb.submitStorage(sub, func() { sb.handleError(sub, msg) })
|
|
case "skipped":
|
|
sb.submitStorage(sub, func() { sb.handleSkipped(sub, msg) })
|
|
default:
|
|
slog.Warn("Unknown scanner message type",
|
|
"type", msg.Type,
|
|
"subscriberId", sub.id)
|
|
continue
|
|
}
|
|
|
|
// Handing work over can block when storage is backed up, and time
|
|
// spent waiting for our own storage is not evidence against the
|
|
// scanner. Restart the clock rather than counting that wait against it.
|
|
refreshRead()
|
|
}
|
|
}
|
|
|
|
// isReadTimeout reports whether a read failed because the deadline expired
|
|
// rather than because the peer or the network did something.
|
|
func isReadTimeout(err error) bool {
|
|
var netErr net.Error
|
|
return errors.As(err, &netErr) && netErr.Timeout()
|
|
}
|
|
|
|
// submitStorage hands a terminal message's work to handleStorage.
|
|
//
|
|
// The send is a plain blocking one, and safe: handleReader is the only sender,
|
|
// and handleStorage does not exit until handleReader has closed readerDone and
|
|
// the queue has drained. A verdict is never dropped — it is the only thing that
|
|
// stops the appview reporting the manifest as never scanned.
|
|
func (sb *ScanBroadcaster) submitStorage(sub *ScanSubscriber, fn func()) {
|
|
sub.storage <- fn
|
|
}
|
|
|
|
// handleStorage runs the storage work terminal messages trigger, one at a time
|
|
// and in the order the messages arrived.
|
|
//
|
|
// Serial and ordered is the point, not a limitation. Two messages about one job
|
|
// applied out of order would let a stale verdict overwrite a fresh one, and the
|
|
// concurrency there was worth nothing anyway: every one of these writes ends in
|
|
// CreateScanRecord, which serialises on the hold's single per-uid repo lock.
|
|
func (sb *ScanBroadcaster) handleStorage(sub *ScanSubscriber) {
|
|
for {
|
|
select {
|
|
case fn := <-sub.storage:
|
|
fn()
|
|
case <-sub.readerDone:
|
|
// The reader has stopped. Finish what it already handed over: those
|
|
// are verdicts the scanner really produced, and dropping one leaves
|
|
// its manifest looking unscanned until the next stale pass.
|
|
for {
|
|
select {
|
|
case fn := <-sub.storage:
|
|
fn()
|
|
default:
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// handleAck marks a job as processing (scanner received and queued it).
|
|
//
|
|
// The ack means "I have it", nothing more: the scanner sends it off its
|
|
// WebSocket reader the moment a frame arrives, before the job is even queued.
|
|
// It deliberately does not touch assigned_at or started_at — a job can sit
|
|
// acked behind a scanner's workers for a long time without anything being
|
|
// wrong, and the deadlines have to be able to tell that apart from a wedge.
|
|
// The signal for "a worker is on it" is 'started'.
|
|
func (sb *ScanBroadcaster) handleAck(sub *ScanSubscriber, seq int64) {
|
|
_, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'processing'
|
|
WHERE seq = ? AND assigned_to = ? AND status = 'assigned'
|
|
`, seq, sub.id)
|
|
if err != nil {
|
|
slog.Error("Failed to update job status to processing",
|
|
"seq", seq,
|
|
"subscriberId", sub.id,
|
|
"error", err)
|
|
return
|
|
}
|
|
|
|
slog.Info("Scan job acknowledged",
|
|
"seq", seq,
|
|
"subscriberId", sub.id)
|
|
}
|
|
|
|
// handleStarted records that a worker has actually begun this scan, which is
|
|
// the moment the scanning deadline is measured from.
|
|
//
|
|
// Without it the hold could only observe dispatch, and the ten-minute budget
|
|
// covered however long the job spent queued inside the scanner — so a healthy
|
|
// scanner working through a backlog had its work cancelled underneath it, each
|
|
// cancellation now writing a failed scan record the user can see. A scanner
|
|
// too old to send this message is not broken by its absence: started_at stays
|
|
// NULL and the job falls under queuedTimeout instead, a budget sized for what
|
|
// the hold can actually see.
|
|
//
|
|
// The assigned_to guard is what makes the message safe with several scanners
|
|
// connected, and started_at is stamped once: a job has one beginning, and a
|
|
// repeat must not roll the deadline forward.
|
|
func (sb *ScanBroadcaster) handleStarted(sub *ScanSubscriber, seq int64) {
|
|
res, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'processing', started_at = ?
|
|
WHERE seq = ? AND assigned_to = ?
|
|
AND status IN ('assigned', 'processing') AND started_at IS NULL
|
|
`, time.Now(), seq, sub.id)
|
|
if err != nil {
|
|
slog.Error("Failed to record scan start",
|
|
"seq", seq, "subscriberId", sub.id, "error", err)
|
|
return
|
|
}
|
|
if n, err := res.RowsAffected(); err == nil && n == 0 {
|
|
slog.Debug("Ignoring 'started' for a job this scanner does not hold",
|
|
"seq", seq, "subscriberId", sub.id)
|
|
return
|
|
}
|
|
|
|
slog.Info("Scan job started",
|
|
"seq", seq,
|
|
"subscriberId", sub.id)
|
|
}
|
|
|
|
// claimJobForTerminal reads the row a terminal message refers to and refuses
|
|
// it unless this scanner is the one holding it.
|
|
//
|
|
// Only handleAck used to check. With one scanner that was harmless; with
|
|
// several it means any scanner can complete, fail or skip another's job —
|
|
// writing a scan record for an image it never looked at, releasing a digest a
|
|
// different process is still scanning, and freeing dispatch capacity that is
|
|
// still in use. The verdict of whichever message lands first wins.
|
|
//
|
|
// A row whose assigned_to no longer matches has moved on: reclaimed after a
|
|
// disconnect, or handed to another process. Its old holder's answer is stale
|
|
// by definition and is dropped rather than applied.
|
|
func (sb *ScanBroadcaster) claimJobForTerminal(sub *ScanSubscriber, seq int64, kind string) bool {
|
|
var owner sql.NullString
|
|
err := sb.db.QueryRow(`SELECT assigned_to FROM scan_jobs WHERE seq = ?`, seq).Scan(&owner)
|
|
if err != nil {
|
|
slog.Error("Failed to check scan job ownership",
|
|
"seq", seq, "type", kind, "subscriberId", sub.id, "error", err)
|
|
return false
|
|
}
|
|
if !owner.Valid || owner.String != sub.id {
|
|
slog.Warn("Ignoring scan message for a job assigned to another scanner",
|
|
"seq", seq,
|
|
"type", kind,
|
|
"subscriberId", sub.id,
|
|
"assignedTo", owner.String)
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
// handleResult processes a completed scan result: uploads SBOM blob + stores scan record in PDS.
|
|
//
|
|
// It runs on handleStorage's goroutine, not on the reader's, and on a context
|
|
// with a deadline. It used to be neither: an S3 endpoint that accepted the
|
|
// connection and then stalled held the entire scanner connection hostage
|
|
// indefinitely, and because the scanner was still nominally connected nothing
|
|
// requeued the work either.
|
|
//
|
|
// The two budgets are separate on purpose. The uploads get scannerResultTimeout
|
|
// between them; the record write gets its own fresh scannerRecordTimeout, so a
|
|
// stalled upload cannot spend the time the record needs. The record is written
|
|
// whichever way the uploads went — a blob already in S3 with no record pointing
|
|
// at it is orphaned, and a manifest with no record at all reads as never
|
|
// scanned.
|
|
//
|
|
// The uploads are skipped entirely when the rescan produced materially the same
|
|
// result as the record already on file; see reusableScanBlobs. The record is
|
|
// still rewritten in that case, pointing at the blobs already stored, so
|
|
// scannedAt advances and the stale loop moves on.
|
|
func (sb *ScanBroadcaster) handleResult(sub *ScanSubscriber, msg ScannerMessage) {
|
|
ctx, cancel := context.WithTimeout(context.Background(), sb.resultDeadline())
|
|
defer cancel()
|
|
|
|
// Before the S3 uploads, not after: a scanner that does not hold this job
|
|
// should not get its payload stored either.
|
|
if !sb.claimJobForTerminal(sub, msg.Seq, "result") {
|
|
return
|
|
}
|
|
|
|
slog.Info("Scan result received",
|
|
"seq", msg.Seq,
|
|
"subscriberId", sub.id,
|
|
"hasSBOM", msg.SBOM != "",
|
|
"hasVulnReport", msg.VulnReport != "",
|
|
"summary", msg.Summary)
|
|
|
|
// Get job details from database
|
|
var (
|
|
manifestDigest string
|
|
repository string
|
|
tag string
|
|
userDID string
|
|
userHandle string
|
|
)
|
|
err := sb.db.QueryRow(`
|
|
SELECT manifest_digest, repository, tag, user_did, COALESCE(user_handle, '')
|
|
FROM scan_jobs WHERE seq = ?
|
|
`, msg.Seq).Scan(&manifestDigest, &repository, &tag, &userDID, &userHandle)
|
|
if err != nil {
|
|
slog.Error("Failed to get job details for result storage",
|
|
"seq", msg.Seq,
|
|
"error", err)
|
|
return
|
|
}
|
|
|
|
// What, if anything, this result may keep from the record already on file.
|
|
// Empty unless the rescan produced materially the same thing.
|
|
reuse := sb.reusableScanBlobs(ctx, manifestDigest, scannerVersion, msg)
|
|
|
|
// Upload SBOM as a blob to the hold's PDS blob storage (like manifest blobs)
|
|
var sbomBlob *lexutil.LexBlob
|
|
switch {
|
|
case msg.SBOM == "":
|
|
case reuse.sbom != nil:
|
|
sbomBlob = reuse.sbom
|
|
default:
|
|
blob, err := uploadBlobToStorage(ctx, sb.s3, sb.holdDID, []byte(msg.SBOM), sbomMimeType)
|
|
if err != nil {
|
|
slog.Error("Failed to upload SBOM blob to PDS storage",
|
|
"seq", msg.Seq,
|
|
"error", err)
|
|
} else {
|
|
sbomBlob = blob
|
|
}
|
|
}
|
|
|
|
// Upload vulnerability report as a blob (full Grype JSON with CVE details)
|
|
var vulnReportBlob *lexutil.LexBlob
|
|
switch {
|
|
case msg.VulnReport == "":
|
|
case reuse.vuln != nil:
|
|
vulnReportBlob = reuse.vuln
|
|
default:
|
|
blob, err := uploadBlobToStorage(ctx, sb.s3, sb.holdDID, []byte(msg.VulnReport), vulnReportMimeType)
|
|
if err != nil {
|
|
slog.Error("Failed to upload VulnReport blob to PDS storage",
|
|
"seq", msg.Seq,
|
|
"error", err)
|
|
} else {
|
|
vulnReportBlob = blob
|
|
}
|
|
}
|
|
|
|
if reuse.sbom != nil || reuse.vuln != nil {
|
|
slog.Info("Rescan produced the stored result; keeping the existing blobs",
|
|
"seq", msg.Seq,
|
|
"manifest", manifestDigest,
|
|
"sbomReused", reuse.sbom != nil,
|
|
"vulnReportReused", reuse.vuln != nil)
|
|
}
|
|
|
|
// Store scan result as a record in the hold's embedded PDS.
|
|
//
|
|
// A result with no summary is a completed scan from a scanner running with
|
|
// vulnerability scanning off: it produced an SBOM, Grype never ran. The
|
|
// record is written either way, because the SBOM blob is already in S3 by
|
|
// this point and nothing but the record would reference it. The counts stay
|
|
// zero and vulnReportBlob stays nil, which is how the appview tells "not
|
|
// scanned for vulnerabilities" apart from "scanned, found none" — see
|
|
// classifyScanRecord in pkg/appview/handlers/scan_result.go. Do not fill in
|
|
// a zeroed summary here: that would report every image as clean.
|
|
var critical, high, medium, low, total int
|
|
if msg.Summary != nil {
|
|
critical = msg.Summary.Critical
|
|
high = msg.Summary.High
|
|
medium = msg.Summary.Medium
|
|
low = msg.Summary.Low
|
|
total = msg.Summary.Total
|
|
}
|
|
|
|
scanRecord := atproto.NewScanRecord(
|
|
manifestDigest, repository, userDID,
|
|
sbomBlob, vulnReportBlob,
|
|
critical, high, medium, low, total,
|
|
scannerVersion,
|
|
)
|
|
|
|
recordCtx, recordCancel := context.WithTimeout(context.Background(), scannerRecordTimeout)
|
|
defer recordCancel()
|
|
|
|
rpath, _, err := sb.pds.CreateScanRecord(recordCtx, scanRecord)
|
|
if err != nil {
|
|
slog.Error("Failed to store scan record in PDS",
|
|
"seq", msg.Seq,
|
|
"error", err)
|
|
} else if msg.Summary != nil {
|
|
slog.Info("Scan record stored in PDS",
|
|
"rpath", rpath,
|
|
"manifest", scanRecord.Manifest,
|
|
"critical", critical,
|
|
"high", high,
|
|
"total", total)
|
|
} else {
|
|
slog.Info("Scan record stored in PDS",
|
|
"rpath", rpath,
|
|
"manifest", scanRecord.Manifest,
|
|
"vulnerabilities", "not scanned")
|
|
}
|
|
|
|
// Mark job as completed
|
|
_, err = sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'completed', completed_at = ?
|
|
WHERE seq = ?
|
|
`, time.Now(), msg.Seq)
|
|
if err != nil {
|
|
slog.Error("Failed to mark scan job as completed",
|
|
"seq", msg.Seq,
|
|
"error", err)
|
|
}
|
|
|
|
// Remove from in-flight tracking and wake dispatch loop
|
|
sb.removeInflight(manifestDigest)
|
|
sb.signalCompletion()
|
|
|
|
if msg.Summary != nil {
|
|
slog.Info("Scan job completed",
|
|
"seq", msg.Seq,
|
|
"repository", repository,
|
|
"tag", tag,
|
|
"critical", msg.Summary.Critical,
|
|
"high", msg.Summary.High,
|
|
"total", msg.Summary.Total)
|
|
} else {
|
|
slog.Info("Scan job completed",
|
|
"seq", msg.Seq,
|
|
"repository", repository,
|
|
"tag", tag,
|
|
"vulnerabilities", "not scanned")
|
|
}
|
|
}
|
|
|
|
// Blob media types for the two artifacts a scan produces. Named because the
|
|
// vulnerability report's type is also needed to compute the reference a fresh
|
|
// report would be stored under, and a mismatch there would silently defeat
|
|
// deduplication.
|
|
const (
|
|
sbomMimeType = "application/spdx+json"
|
|
vulnReportMimeType = "application/vnd.atcr.vulnerabilities+json"
|
|
|
|
// scannerVersion stamps every record this build writes, and gates blob
|
|
// reuse: see reusableScanBlobs.
|
|
scannerVersion = "atcr-scanner-v1.0.0"
|
|
)
|
|
|
|
// scanBlobReuse is what a result may keep from the scan record already on file
|
|
// instead of uploading its own copy. A nil field means "upload it".
|
|
type scanBlobReuse struct {
|
|
sbom *lexutil.LexBlob
|
|
vuln *lexutil.LexBlob
|
|
}
|
|
|
|
// reusableScanBlobs decides whether a rescan produced materially the same
|
|
// result as the record already on file, and therefore may keep its blobs.
|
|
//
|
|
// The rule: keep the stored blobs when the record on file was written by this
|
|
// same scanner version and carries a vulnerability report whose reference
|
|
// equals the one this report would be stored under. Each blob is then kept
|
|
// individually, and only if it is still in S3.
|
|
//
|
|
// Why the report and not the SBOM. The stale-scan loop rescans on a timer, and
|
|
// an unchanged image rescanned a week later yields a byte-identical Grype
|
|
// report — it carries no timestamp and no database build date, only matches,
|
|
// source, distro and the Grype version (scanner/internal/scan/grype.go). The
|
|
// SBOM does not: Syft stamps creationInfo.created from the wall clock and
|
|
// appends a random UUID to documentNamespace, so its bytes, and therefore its
|
|
// content-addressed reference, move on every single run. The SBOM can never
|
|
// recognise itself; the report has to do it for both.
|
|
//
|
|
// Why the scanner version is part of the rule. The report lists matched
|
|
// packages, not every package, so it is possible in principle for the SBOM to
|
|
// change while the report does not — a Syft upgrade that starts cataloguing an
|
|
// ecosystem with no known vulnerabilities would do it. It is only possible
|
|
// across a scanner change, because the content itself is pinned by the manifest
|
|
// digest, so refusing to reuse across versions closes the hole completely. The
|
|
// version is a build constant today, which makes this clause inert until it
|
|
// starts moving; it is cheap, and the alternative is remembering to add it at
|
|
// the moment it is first needed.
|
|
//
|
|
// A vulnerability database update is not a case this suppresses: new matches
|
|
// mean a different report, no reference match, and both blobs are written
|
|
// fresh. That is correct — the report genuinely changed.
|
|
//
|
|
// Every failure path here returns "reuse nothing", which costs an upload and
|
|
// never costs correctness. That includes the timeout: the lookup runs on a
|
|
// short budget carved out of the caller's, so a slow CAR read cannot eat the
|
|
// upload deadline it is standing in for.
|
|
func (sb *ScanBroadcaster) reusableScanBlobs(ctx context.Context, manifestDigest, version string, msg ScannerMessage) scanBlobReuse {
|
|
var none scanBlobReuse
|
|
|
|
if sb.pds == nil || sb.s3 == nil || msg.VulnReport == "" {
|
|
return none
|
|
}
|
|
|
|
fresh, err := blobRefForBytes([]byte(msg.VulnReport), vulnReportMimeType)
|
|
if err != nil {
|
|
return none
|
|
}
|
|
|
|
ctx, cancel := context.WithTimeout(ctx, scannerReuseTimeout)
|
|
defer cancel()
|
|
|
|
_, prev, err := sb.pds.GetScanRecord(ctx, manifestDigest)
|
|
if err != nil || prev == nil {
|
|
// No record yet: the overwhelmingly common first-scan path.
|
|
return none
|
|
}
|
|
if prev.ScannerVersion != version {
|
|
return none
|
|
}
|
|
if prev.VulnReportBlob == nil || prev.VulnReportBlob.Ref != fresh.Ref {
|
|
// Either the previous upload failed and left nothing to compare
|
|
// against, or the result really is different.
|
|
return none
|
|
}
|
|
|
|
// Materially unchanged. Each blob is kept only if the object it names is
|
|
// still there: a record can outlive its blob (a lifecycle rule, a restore,
|
|
// a GC pass), and reusing a dangling reference would make that permanent,
|
|
// since every later rescan reaches this same decision.
|
|
var reuse scanBlobReuse
|
|
if sb.blobStillStored(ctx, prev.VulnReportBlob) {
|
|
reuse.vuln = prev.VulnReportBlob
|
|
}
|
|
if msg.SBOM != "" && sb.blobStillStored(ctx, prev.SbomBlob) {
|
|
reuse.sbom = prev.SbomBlob
|
|
}
|
|
return reuse
|
|
}
|
|
|
|
// blobStillStored reports whether the object a blob reference names is still in
|
|
// the hold's bucket. Any error is treated as "gone", which re-uploads.
|
|
func (sb *ScanBroadcaster) blobStillStored(ctx context.Context, blob *lexutil.LexBlob) bool {
|
|
if blob == nil {
|
|
return false
|
|
}
|
|
_, err := sb.s3.Stat(ctx, atprotoBlobPath(sb.holdDID, blob.Ref.String()))
|
|
return err == nil
|
|
}
|
|
|
|
// handleError marks a job as failed and creates a scan record so the stale
|
|
// loop won't immediately retry. Failed records still get retried on the
|
|
// rescan interval since failures may be transient (network, OOM, etc.).
|
|
func (sb *ScanBroadcaster) handleError(sub *ScanSubscriber, msg ScannerMessage) {
|
|
// Bounded, and off the reader goroutine: this writes a CAR delta under the
|
|
// hold's per-uid repo lock, which every layer record, stats increment and
|
|
// Bluesky post also contends for.
|
|
ctx, cancel := context.WithTimeout(context.Background(), scannerRecordTimeout)
|
|
defer cancel()
|
|
|
|
if !sb.claimJobForTerminal(sub, msg.Seq, "error") {
|
|
return
|
|
}
|
|
|
|
var manifestDigest, repository, userDID string
|
|
err := sb.db.QueryRow(`
|
|
SELECT manifest_digest, repository, user_did
|
|
FROM scan_jobs WHERE seq = ?
|
|
`, msg.Seq).Scan(&manifestDigest, &repository, &userDID)
|
|
if err != nil {
|
|
slog.Error("Failed to get job details for failure record",
|
|
"seq", msg.Seq, "error", err)
|
|
} else {
|
|
scanRecord := atproto.NewFailedScanRecord(
|
|
manifestDigest, repository, userDID,
|
|
msg.Error,
|
|
"atcr-scanner-v1.0.0",
|
|
)
|
|
if _, _, err := sb.pds.CreateScanRecord(ctx, scanRecord); err != nil {
|
|
slog.Error("Failed to store failure scan record",
|
|
"seq", msg.Seq, "error", err)
|
|
}
|
|
}
|
|
|
|
_, err = sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'failed', completed_at = ?
|
|
WHERE seq = ?
|
|
`, time.Now(), msg.Seq)
|
|
if err != nil {
|
|
slog.Error("Failed to mark scan job as failed",
|
|
"seq", msg.Seq,
|
|
"error", err)
|
|
}
|
|
|
|
sb.removeInflight(manifestDigest)
|
|
sb.signalCompletion()
|
|
|
|
slog.Warn("Scan job failed",
|
|
"seq", msg.Seq,
|
|
"subscriberId", sub.id,
|
|
"error", msg.Error)
|
|
}
|
|
|
|
// handleSkipped marks a job complete and creates a scan record with
|
|
// status="skipped". The stale-scan loop will leave these records alone — the
|
|
// outcome won't change until the scanner gains support for the artifact type.
|
|
func (sb *ScanBroadcaster) handleSkipped(sub *ScanSubscriber, msg ScannerMessage) {
|
|
ctx, cancel := context.WithTimeout(context.Background(), scannerRecordTimeout)
|
|
defer cancel()
|
|
|
|
if !sb.claimJobForTerminal(sub, msg.Seq, "skipped") {
|
|
return
|
|
}
|
|
|
|
var manifestDigest, repository, userDID string
|
|
err := sb.db.QueryRow(`
|
|
SELECT manifest_digest, repository, user_did
|
|
FROM scan_jobs WHERE seq = ?
|
|
`, msg.Seq).Scan(&manifestDigest, &repository, &userDID)
|
|
if err != nil {
|
|
slog.Error("Failed to get job details for skip record",
|
|
"seq", msg.Seq, "error", err)
|
|
} else {
|
|
scanRecord := atproto.NewSkippedScanRecord(
|
|
manifestDigest, repository, userDID,
|
|
msg.Reason,
|
|
"atcr-scanner-v1.0.0",
|
|
)
|
|
if _, _, err := sb.pds.CreateScanRecord(ctx, scanRecord); err != nil {
|
|
slog.Error("Failed to store skipped scan record",
|
|
"seq", msg.Seq, "error", err)
|
|
}
|
|
}
|
|
|
|
_, err = sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'completed', completed_at = ?
|
|
WHERE seq = ?
|
|
`, time.Now(), msg.Seq)
|
|
if err != nil {
|
|
slog.Error("Failed to mark scan job as completed (skipped)",
|
|
"seq", msg.Seq,
|
|
"error", err)
|
|
}
|
|
|
|
sb.removeInflight(manifestDigest)
|
|
sb.signalCompletion()
|
|
|
|
slog.Info("Scan job skipped",
|
|
"seq", msg.Seq,
|
|
"subscriberId", sub.id,
|
|
"reason", msg.Reason)
|
|
}
|
|
|
|
// drainPendingJobs sends pending/timed-out jobs to a newly connected scanner.
|
|
// Collects all pending rows first, closes cursor, then assigns and dispatches
|
|
// to avoid holding a SELECT cursor open during UPDATEs (prevents SQLite BUSY).
|
|
//
|
|
// It stops at the scanner's capacity. It used to walk every pending row, so
|
|
// with a backlog and several scanner processes the first one to connect took
|
|
// all of it and the rest stayed idle — horizontal scaling defeated at the
|
|
// connect path rather than at the dispatch gate. What it leaves behind is
|
|
// picked up by offerPendingJobs as soon as anything frees up.
|
|
func (sb *ScanBroadcaster) drainPendingJobs(sub *ScanSubscriber, cursor int64) {
|
|
rows, err := sb.db.Query(`
|
|
SELECT seq, manifest_digest, repository, tag, user_did, user_handle, hold_did, hold_endpoint, tier, config_json, layers_json
|
|
FROM scan_jobs
|
|
WHERE status = 'pending' AND seq > ?
|
|
ORDER BY seq ASC
|
|
`, cursor)
|
|
if err != nil {
|
|
slog.Error("Failed to drain pending scan jobs", "error", err)
|
|
return
|
|
}
|
|
|
|
var jobs []*ScanJobEvent
|
|
for rows.Next() {
|
|
job := &ScanJobEvent{Type: "job"}
|
|
var configJSON, layersJSON string
|
|
|
|
err := rows.Scan(
|
|
&job.Seq, &job.ManifestDigest, &job.Repository, &job.Tag,
|
|
&job.UserDID, &job.UserHandle, &job.HoldDID, &job.HoldEndpoint,
|
|
&job.Tier, &configJSON, &layersJSON,
|
|
)
|
|
if err != nil {
|
|
slog.Error("Failed to scan pending job row", "error", err)
|
|
continue
|
|
}
|
|
|
|
job.Config = json.RawMessage(configJSON)
|
|
job.Layers = json.RawMessage(layersJSON)
|
|
jobs = append(jobs, job)
|
|
}
|
|
rows.Close()
|
|
|
|
// Anything this scanner is already holding — jobs it resumed after a
|
|
// reconnect, most often — counts against what it can take now.
|
|
budget := sub.effectiveCapacity()
|
|
if load, ok := sb.subscriberLoads(); ok {
|
|
budget -= load[sub.id]
|
|
}
|
|
|
|
count := 0
|
|
for _, job := range jobs {
|
|
if count >= budget {
|
|
slog.Debug("Drain reached the scanner's capacity, leaving the rest pending",
|
|
"subscriberId", sub.id, "capacity", sub.effectiveCapacity())
|
|
break
|
|
}
|
|
res, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'assigned', assigned_to = ?, assigned_at = ?
|
|
WHERE seq = ? AND status = 'pending'
|
|
`, sub.id, time.Now(), job.Seq)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
// The rows were selected up front, so dispatchJob or the re-dispatch
|
|
// loop can have claimed one in between. Sending it anyway would scan
|
|
// it twice — and the losing scanner's ack is dropped by handleAck's
|
|
// own guard, so nothing downstream would notice the duplicate.
|
|
if n, err := res.RowsAffected(); err == nil && n == 0 {
|
|
slog.Debug("Scan job no longer pending, skipping drain", "seq", job.Seq)
|
|
continue
|
|
}
|
|
|
|
select {
|
|
case sub.send <- job:
|
|
count++
|
|
case <-sub.done:
|
|
sb.releaseUndelivered(sub, job.Seq)
|
|
return
|
|
case <-time.After(drainSendTimeout):
|
|
slog.Warn("Drain timeout for scanner", "subscriberId", sub.id, "seq", job.Seq)
|
|
sb.releaseUndelivered(sub, job.Seq)
|
|
return
|
|
}
|
|
}
|
|
|
|
if count > 0 {
|
|
slog.Info("Drained pending jobs to scanner",
|
|
"subscriberId", sub.id,
|
|
"jobsDrained", count)
|
|
}
|
|
}
|
|
|
|
// releaseUndelivered puts back a row this drain claimed but could not hand to
|
|
// the scanner. Without it the row stays 'assigned' to a subscriber that was
|
|
// never sent it, counting as active dispatch capacity until the five-minute
|
|
// ack timeout reclaims it.
|
|
//
|
|
// The assigned_to guard is what makes this safe to run after the subscriber is
|
|
// already gone: if Unsubscribe's bulk reset or another dispatcher has since
|
|
// taken the row, this matches nothing.
|
|
func (sb *ScanBroadcaster) releaseUndelivered(sub *ScanSubscriber, seq int64) {
|
|
_, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'pending', assigned_to = NULL, assigned_at = NULL
|
|
WHERE seq = ? AND assigned_to = ? AND status = 'assigned'
|
|
`, seq, sub.id)
|
|
if err != nil {
|
|
slog.Error("Failed to release undelivered scan job",
|
|
"seq", seq, "subscriberId", sub.id, "error", err)
|
|
}
|
|
}
|
|
|
|
// reDispatchLoop periodically checks for timed-out jobs and re-dispatches them,
|
|
// and fills freed scanner capacity as soon as it is freed rather than on the
|
|
// next tick.
|
|
func (sb *ScanBroadcaster) reDispatchLoop() {
|
|
defer sb.wg.Done()
|
|
|
|
ticker := time.NewTicker(30 * time.Second)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
case <-ticker.C:
|
|
sb.reDispatchTimedOut()
|
|
case <-sb.capacityFreed:
|
|
sb.offerPendingJobs()
|
|
}
|
|
}
|
|
}
|
|
|
|
// offerPendingJobs hands waiting rows to whichever scanners have room, oldest
|
|
// first, and stops as soon as nobody does.
|
|
//
|
|
// This is the other half of admission control. dispatchJob now leaves a row
|
|
// pending rather than pushing it into a saturated scanner's own queue, which
|
|
// is only affordable if the row is offered again the moment something frees
|
|
// up. Waiting for the thirty-second re-dispatch tick would have cost more
|
|
// throughput than the scanner-side queue ever bought.
|
|
//
|
|
// No age guard is needed here, unlike the reclaim in reDispatchTimedOut:
|
|
// dispatchJob's assign is conditional on the row still being pending and reads
|
|
// RowsAffected, so racing another dispatcher costs a skipped row, not a
|
|
// double send.
|
|
func (sb *ScanBroadcaster) offerPendingJobs() {
|
|
if !sb.hasConnectedScanners() {
|
|
return
|
|
}
|
|
|
|
rows, err := sb.db.Query(`
|
|
SELECT seq, manifest_digest, repository, tag, user_did, user_handle, hold_did, hold_endpoint, tier, config_json, layers_json
|
|
FROM scan_jobs
|
|
WHERE status = 'pending'
|
|
ORDER BY seq ASC
|
|
`)
|
|
if err != nil {
|
|
slog.Error("Failed to query pending scan jobs", "error", err)
|
|
return
|
|
}
|
|
|
|
var jobs []*ScanJobEvent
|
|
for rows.Next() {
|
|
job := &ScanJobEvent{Type: "job"}
|
|
var configJSON, layersJSON string
|
|
if err := rows.Scan(
|
|
&job.Seq, &job.ManifestDigest, &job.Repository, &job.Tag,
|
|
&job.UserDID, &job.UserHandle, &job.HoldDID, &job.HoldEndpoint,
|
|
&job.Tier, &configJSON, &layersJSON,
|
|
); err != nil {
|
|
slog.Error("Failed to scan pending job row", "error", err)
|
|
continue
|
|
}
|
|
job.Config = json.RawMessage(configJSON)
|
|
job.Layers = json.RawMessage(layersJSON)
|
|
jobs = append(jobs, job)
|
|
}
|
|
rows.Close()
|
|
|
|
for _, job := range jobs {
|
|
if !sb.hasFreeScannerCapacity() {
|
|
return
|
|
}
|
|
sb.dispatchJob(job)
|
|
}
|
|
}
|
|
|
|
// hasFreeScannerCapacity reports whether any connected scanner has room for
|
|
// another job right now.
|
|
func (sb *ScanBroadcaster) hasFreeScannerCapacity() bool {
|
|
sb.mu.Lock()
|
|
defer sb.mu.Unlock()
|
|
return sb.selectSubscriberLocked() != nil
|
|
}
|
|
|
|
// reDispatchTimedOut finds jobs that were assigned but not acked/completed within timeout,
|
|
// re-offers jobs that have been sitting in 'pending' with nobody to hand them to,
|
|
// reclaims work from a scanner that disconnected and did not come back inside
|
|
// reconnectGrace, and also marks stuck processing jobs as failed.
|
|
// Collects timed-out rows first, closes cursor, then resets and re-dispatches
|
|
// to avoid holding a SELECT cursor open during UPDATEs (prevents SQLite BUSY).
|
|
//
|
|
// Pending rows matter as much as assigned ones: Enqueue and drainPendingJobs are
|
|
// both one-shot dispatchers, so before this loop covered 'pending' a row whose
|
|
// single dispatch attempt failed (or that a disconnect unassigned while no
|
|
// scanner reconnected afterwards) stayed pending forever — and hasActiveJobs()
|
|
// counted it, which stopped proactive scanning deployment-wide.
|
|
func (sb *ScanBroadcaster) reDispatchTimedOut() {
|
|
timeout := time.Now().Add(-sb.ackTimeout)
|
|
|
|
// Fail processing jobs stuck past the deadline (scanner likely crashed mid-scan)
|
|
sb.failStuckProcessingJobs()
|
|
|
|
rows, err := sb.db.Query(`
|
|
SELECT seq, manifest_digest, repository, tag, user_did, user_handle, hold_did, hold_endpoint, tier, config_json, layers_json, status
|
|
FROM scan_jobs
|
|
WHERE (status = 'assigned' AND disconnected_at IS NULL AND assigned_at < ?)
|
|
OR (status = 'pending' AND datetime(created_at) < datetime('now', ?))
|
|
OR (status IN ('assigned', 'processing') AND disconnected_at IS NOT NULL AND disconnected_at < ?)
|
|
ORDER BY seq ASC
|
|
`, timeout, sqliteAgoModifier(pendingReclaimAfter), time.Now().Add(-reconnectGrace))
|
|
if err != nil {
|
|
slog.Error("Failed to query timed-out scan jobs", "error", err)
|
|
return
|
|
}
|
|
|
|
type reclaimable struct {
|
|
job *ScanJobEvent
|
|
status string
|
|
}
|
|
|
|
var jobs []reclaimable
|
|
for rows.Next() {
|
|
job := &ScanJobEvent{Type: "job"}
|
|
var configJSON, layersJSON, status string
|
|
|
|
err := rows.Scan(
|
|
&job.Seq, &job.ManifestDigest, &job.Repository, &job.Tag,
|
|
&job.UserDID, &job.UserHandle, &job.HoldDID, &job.HoldEndpoint,
|
|
&job.Tier, &configJSON, &layersJSON, &status,
|
|
)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
|
|
job.Config = json.RawMessage(configJSON)
|
|
job.Layers = json.RawMessage(layersJSON)
|
|
jobs = append(jobs, reclaimable{job: job, status: status})
|
|
}
|
|
rows.Close()
|
|
|
|
for _, r := range jobs {
|
|
job := r.job
|
|
|
|
// Already-pending rows need no reset — and skipping the UPDATE keeps a
|
|
// row another dispatcher claimed between the SELECT and here assigned,
|
|
// so dispatchJob's status guard can drop it instead of double-sending.
|
|
if r.status != "pending" {
|
|
_, err = sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'pending', assigned_to = NULL,
|
|
assigned_at = NULL, started_at = NULL, disconnected_at = NULL
|
|
WHERE seq = ?
|
|
`, job.Seq)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
}
|
|
|
|
slog.Info("Re-dispatching scan job",
|
|
"seq", job.Seq,
|
|
"repository", job.Repository,
|
|
"previousStatus", r.status)
|
|
|
|
sb.dispatchJob(job)
|
|
}
|
|
}
|
|
|
|
// failStuckProcessingJobs gives up on jobs a scanner acknowledged but never
|
|
// answered for, and retires them the same way every other terminal transition
|
|
// does: a scan record, the in-flight digest released, dispatch capacity
|
|
// signalled.
|
|
//
|
|
// It used to be a bare UPDATE. That made it the only terminal transition that
|
|
// wrote no record and never called removeInflight, and both halves of that
|
|
// were user-visible. Without the record the appview cannot tell a job the hold
|
|
// gave up on from one enqueued thirty seconds ago, so the image shows a grey
|
|
// "Not scanned" forever. Without removeInflight the digest stays in sb.inflight
|
|
// for the life of the process, and discoverUnscannedForUser and runStalePass
|
|
// both skip any manifest already in that set — so every timeout permanently
|
|
// retired one image from scanning.
|
|
//
|
|
// The record is a failure rather than a skip because a hung or crashed scanner
|
|
// is a transient condition: the stale loop should retry it on the rescan
|
|
// interval, which is exactly what it does for failed records and not for
|
|
// skipped ones.
|
|
func (sb *ScanBroadcaster) failStuckProcessingJobs() {
|
|
type stuckJob struct {
|
|
seq int64
|
|
manifestDigest string
|
|
repository string
|
|
userDID string
|
|
started bool
|
|
}
|
|
|
|
// Select first, close the cursor, then write — the same shape as the
|
|
// sibling loops below, and required here because the digests are needed
|
|
// after the UPDATE has already erased which rows were affected.
|
|
// Two deadlines, because there are two things a 'processing' row can be.
|
|
//
|
|
// A row with started_at is being scanned right now and gets scanningTimeout
|
|
// from that moment. A row without it was acked and is queued inside the
|
|
// scanner — either because it genuinely has not reached a worker yet, or
|
|
// because the scanner is too old to say — and gets queuedTimeout from
|
|
// dispatch, which is the only clock the hold has for it.
|
|
scanDeadline := time.Now().Add(-scanningTimeout)
|
|
queueDeadline := time.Now().Add(-queuedTimeout)
|
|
|
|
rows, err := sb.db.Query(`
|
|
SELECT seq, manifest_digest, repository, user_did, started_at IS NOT NULL
|
|
FROM scan_jobs
|
|
WHERE status = 'processing'
|
|
AND ((started_at IS NOT NULL AND started_at < ?)
|
|
OR (started_at IS NULL AND assigned_at < ?))
|
|
`, scanDeadline, queueDeadline)
|
|
if err != nil {
|
|
slog.Error("Failed to query stuck processing jobs", "error", err)
|
|
return
|
|
}
|
|
|
|
var jobs []stuckJob
|
|
for rows.Next() {
|
|
var j stuckJob
|
|
if err := rows.Scan(&j.seq, &j.manifestDigest, &j.repository, &j.userDID, &j.started); err != nil {
|
|
slog.Error("Failed to scan stuck processing job row", "error", err)
|
|
continue
|
|
}
|
|
jobs = append(jobs, j)
|
|
}
|
|
rows.Close()
|
|
|
|
for _, j := range jobs {
|
|
res, err := sb.db.Exec(`
|
|
UPDATE scan_jobs SET status = 'failed', completed_at = ?
|
|
WHERE seq = ? AND status = 'processing'
|
|
AND ((started_at IS NOT NULL AND started_at < ?)
|
|
OR (started_at IS NULL AND assigned_at < ?))
|
|
`, time.Now(), j.seq, scanDeadline, queueDeadline)
|
|
if err != nil {
|
|
slog.Error("Failed to fail stuck processing job", "seq", j.seq, "error", err)
|
|
continue
|
|
}
|
|
// The scanner answered — or started, which moves the row onto the
|
|
// other deadline — between the SELECT and here.
|
|
if n, err := res.RowsAffected(); err == nil && n == 0 {
|
|
continue
|
|
}
|
|
|
|
reason := fmt.Sprintf("scanner did not report a result within %s of starting the scan", scanningTimeout)
|
|
if !j.started {
|
|
reason = fmt.Sprintf("scanner acknowledged the job but never started it within %s", queuedTimeout)
|
|
}
|
|
record := atproto.NewFailedScanRecord(
|
|
j.manifestDigest, j.repository, j.userDID,
|
|
reason,
|
|
"atcr-scanner-v1.0.0",
|
|
)
|
|
if _, _, err := sb.pds.CreateScanRecord(context.Background(), record); err != nil {
|
|
slog.Error("Failed to store timeout scan record", "seq", j.seq, "error", err)
|
|
}
|
|
|
|
sb.removeInflight(j.manifestDigest)
|
|
sb.signalCompletion()
|
|
|
|
slog.Warn("Scan job timed out in processing",
|
|
"seq", j.seq,
|
|
"repository", j.repository,
|
|
"manifest", j.manifestDigest)
|
|
}
|
|
}
|
|
|
|
// sqliteAgoModifier renders a duration as a SQLite datetime() modifier that
|
|
// walks backwards from 'now', e.g. 15m becomes "-900 seconds".
|
|
func sqliteAgoModifier(d time.Duration) string {
|
|
return fmt.Sprintf("-%d seconds", int64(d.Seconds()))
|
|
}
|
|
|
|
// Close stops background goroutines and closes the scan broadcaster's database connection
|
|
func (sb *ScanBroadcaster) Close() error {
|
|
if sb.stopCh != nil {
|
|
close(sb.stopCh)
|
|
sb.wg.Wait()
|
|
}
|
|
if sb.db != nil && sb.ownsDB {
|
|
return sb.db.Close()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Secret returns the scanner shared secret for use in blob read authorization
|
|
func (sb *ScanBroadcaster) Secret() string {
|
|
return sb.secret
|
|
}
|
|
|
|
// ValidateScannerSecret checks if the provided secret matches
|
|
func (sb *ScanBroadcaster) ValidateScannerSecret(secret string) bool {
|
|
return sb.secret != "" && secret == sb.secret
|
|
}
|
|
|
|
// scanCandidate is a manifest that needs scanning, with its scan freshness.
|
|
type scanCandidate struct {
|
|
manifest atproto.ManifestRecord // May be zero-value for stale candidates (resolved lazily)
|
|
manifestDigest string // Always set; used for dedup and lazy resolution
|
|
userDID string
|
|
userHandle string
|
|
scannedAt time.Time // zero value = never scanned
|
|
}
|
|
|
|
// discoveryLoop fetches DIDs from the relay, walks each user's PDS to find manifests
|
|
// with no scan record, and pushes them to the unscannedQueue. Runs on startup (after
|
|
// settle), then every 4 hours. Scanner reconnect triggers an early pass via discoverNow.
|
|
func (sb *ScanBroadcaster) discoveryLoop() {
|
|
defer sb.wg.Done()
|
|
|
|
// Wait for system to settle
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
case <-time.After(45 * time.Second):
|
|
}
|
|
|
|
slog.Info("Discovery loop started")
|
|
sb.runDiscoveryPass()
|
|
|
|
ticker := time.NewTicker(4 * time.Hour)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-sb.stopCh:
|
|
slog.Info("Discovery loop stopped")
|
|
return
|
|
case <-ticker.C:
|
|
sb.runDiscoveryPass()
|
|
case <-sb.discoverNow:
|
|
slog.Info("Discovery loop: early pass triggered (scanner reconnect)")
|
|
sb.runDiscoveryPass()
|
|
}
|
|
}
|
|
}
|
|
|
|
// runDiscoveryPass fetches the DID list from relay and walks each user's PDS.
|
|
func (sb *ScanBroadcaster) runDiscoveryPass() {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
|
|
defer cancel()
|
|
|
|
// A hold that could not be reached last pass gets asked again this one.
|
|
// predecessorUnresolved is a within-pass memo, not a verdict.
|
|
sb.predecessorMu.Lock()
|
|
sb.predecessorUnresolved = nil
|
|
sb.predecessorMu.Unlock()
|
|
|
|
// Fetch DID list from relay
|
|
userDIDs := sb.fetchManifestDIDs(ctx)
|
|
if len(userDIDs) == 0 {
|
|
slog.Debug("Discovery: no manifest DIDs from relay")
|
|
return
|
|
}
|
|
|
|
slog.Info("Discovery: starting pass", "users", len(userDIDs))
|
|
found := 0
|
|
|
|
for _, userDID := range userDIDs {
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
default:
|
|
}
|
|
|
|
n := sb.discoverUnscannedForUser(ctx, userDID)
|
|
found += n
|
|
}
|
|
|
|
slog.Info("Discovery: pass complete", "users", len(userDIDs), "unscannedFound", found)
|
|
}
|
|
|
|
// fetchManifestDIDs queries relays for all DIDs with io.atcr.manifest records.
|
|
// Uses failover: starts from a rotated index (so load is shared across passes)
|
|
// and falls through to the next relay if one fails. Returns DIDs from the first
|
|
// relay that produces a complete result; partial results from a failing relay
|
|
// are discarded so we don't mix incomplete responses.
|
|
func (sb *ScanBroadcaster) fetchManifestDIDs(ctx context.Context) []string {
|
|
endpoints := sb.relayEndpoints
|
|
if len(endpoints) == 0 {
|
|
return nil
|
|
}
|
|
|
|
sb.relayStartMu.Lock()
|
|
start := sb.relayStartIdx % len(endpoints)
|
|
sb.relayStartIdx = (sb.relayStartIdx + 1) % len(endpoints)
|
|
sb.relayStartMu.Unlock()
|
|
|
|
for i := range endpoints {
|
|
relay := endpoints[(start+i)%len(endpoints)]
|
|
dids, err := sb.fetchManifestDIDsFrom(ctx, relay)
|
|
if err != nil {
|
|
slog.Warn("Discovery: relay failed, trying next",
|
|
"relay", relay, "error", err)
|
|
continue
|
|
}
|
|
if i > 0 {
|
|
slog.Info("Discovery: succeeded after failover", "relay", relay, "attemptsTried", i+1)
|
|
}
|
|
return dids
|
|
}
|
|
|
|
slog.Warn("Discovery: all relays failed", "relays", endpoints)
|
|
return nil
|
|
}
|
|
|
|
// fetchManifestDIDsFrom paginates through io.atcr.manifest repos on a single relay.
|
|
// Any error during pagination invalidates the partial result.
|
|
func (sb *ScanBroadcaster) fetchManifestDIDsFrom(ctx context.Context, relay string) ([]string, error) {
|
|
client := atproto.NewClient(relay, "", "")
|
|
var allDIDs []string
|
|
var cursor string
|
|
|
|
for {
|
|
result, err := client.ListReposByCollection(ctx, atproto.ManifestCollection, 1000, cursor)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
for _, repo := range result.Repos {
|
|
allDIDs = append(allDIDs, repo.DID)
|
|
}
|
|
|
|
if result.Cursor == "" || len(result.Repos) == 0 {
|
|
break
|
|
}
|
|
cursor = result.Cursor
|
|
}
|
|
|
|
return allDIDs, nil
|
|
}
|
|
|
|
// discoverUnscannedForUser fetches manifests from a user's PDS and pushes any
|
|
// without scan records to the unscannedQueue. Returns count of candidates found.
|
|
func (sb *ScanBroadcaster) discoverUnscannedForUser(ctx context.Context, userDID string) int {
|
|
did, userHandle, pdsEndpoint, err := atproto.ResolveIdentity(ctx, userDID)
|
|
if err != nil {
|
|
slog.Debug("Discovery: failed to resolve user identity",
|
|
"userDID", userDID, "error", err)
|
|
return 0
|
|
}
|
|
|
|
found := 0
|
|
client := atproto.NewClient(pdsEndpoint, did, "")
|
|
var cursor string
|
|
for {
|
|
records, nextCursor, err := client.ListRecordsForRepo(ctx, did, atproto.ManifestCollection, 100, cursor)
|
|
if err != nil {
|
|
slog.Debug("Discovery: failed to list manifest records",
|
|
"userDID", did, "pds", pdsEndpoint, "error", err)
|
|
return found
|
|
}
|
|
|
|
for _, record := range records {
|
|
var manifest atproto.ManifestRecord
|
|
if err := json.Unmarshal(record.Value, &manifest); err != nil {
|
|
continue
|
|
}
|
|
|
|
// Check if this manifest belongs to us
|
|
holdDID := manifest.HoldDID
|
|
if holdDID == "" {
|
|
holdDID = manifest.HoldEndpoint // Legacy field
|
|
}
|
|
if !sb.isOurManifest(ctx, holdDID) {
|
|
continue
|
|
}
|
|
|
|
// Skip manifest lists and attestations (no scannable content)
|
|
if len(manifest.Layers) == 0 || manifest.Subject != nil || manifest.Config == nil {
|
|
continue
|
|
}
|
|
|
|
// Check if already in-flight
|
|
if !sb.addInflight(manifest.Digest) {
|
|
continue
|
|
}
|
|
|
|
// Check scan status — only interested in never-scanned
|
|
_, _, err := sb.pds.GetScanRecord(ctx, manifest.Digest)
|
|
if err == nil {
|
|
// Has a scan record — not our concern (stale loop handles rescans)
|
|
sb.removeInflight(manifest.Digest)
|
|
continue
|
|
}
|
|
|
|
// Never scanned — push to queue
|
|
candidate := &scanCandidate{
|
|
manifest: manifest,
|
|
manifestDigest: manifest.Digest,
|
|
userDID: did,
|
|
userHandle: userHandle,
|
|
}
|
|
|
|
select {
|
|
case sb.unscannedQueue <- candidate:
|
|
found++
|
|
case <-sb.stopCh:
|
|
sb.removeInflight(manifest.Digest)
|
|
return found
|
|
}
|
|
}
|
|
|
|
if nextCursor == "" || len(records) == 0 {
|
|
break
|
|
}
|
|
cursor = nextCursor
|
|
}
|
|
|
|
return found
|
|
}
|
|
|
|
// staleScanLoop walks local scan records to find stale scans (older than rescanInterval)
|
|
// and pushes them to the staleQueue. No PDS calls needed — all data is local.
|
|
// After a full pass, sleeps for rescanInterval/2 before repeating.
|
|
func (sb *ScanBroadcaster) staleScanLoop() {
|
|
defer sb.wg.Done()
|
|
|
|
// Short initial delay to let system settle
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
case <-time.After(10 * time.Second):
|
|
}
|
|
|
|
slog.Info("Stale scan loop started")
|
|
|
|
for {
|
|
sb.runStalePass()
|
|
|
|
sleepDuration := max(sb.rescanInterval/2, 1*time.Hour)
|
|
|
|
select {
|
|
case <-sb.stopCh:
|
|
slog.Info("Stale scan loop stopped")
|
|
return
|
|
case <-time.After(sleepDuration):
|
|
}
|
|
}
|
|
}
|
|
|
|
// runStalePass walks all scan records and pushes stale ones to the staleQueue.
|
|
func (sb *ScanBroadcaster) runStalePass() {
|
|
ri := sb.pds.RecordsIndex()
|
|
if ri == nil {
|
|
slog.Debug("Stale scan: no records index available")
|
|
return
|
|
}
|
|
|
|
ctx := context.Background()
|
|
found := 0
|
|
var cursor string
|
|
|
|
for {
|
|
records, nextCursor, err := ri.ListRecords(atproto.ScanCollection, 100, cursor, true) // oldest first
|
|
if err != nil {
|
|
slog.Error("Stale scan: failed to list scan records", "error", err)
|
|
return
|
|
}
|
|
|
|
for _, record := range records {
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
default:
|
|
}
|
|
|
|
// The rkey is the manifest digest (without sha256: prefix)
|
|
manifestDigest := "sha256:" + record.Rkey
|
|
|
|
// Check if already in-flight
|
|
if !sb.addInflight(manifestDigest) {
|
|
continue
|
|
}
|
|
|
|
// Fetch the actual scan record to check staleness
|
|
_, scanRecord, err := sb.pds.GetScanRecord(ctx, manifestDigest)
|
|
if err != nil {
|
|
sb.removeInflight(manifestDigest)
|
|
continue
|
|
}
|
|
|
|
// Permanently-skipped records (helm charts, in-toto, etc.) won't
|
|
// change outcome on retry — leave them alone. Failed records still
|
|
// get retried since failures may be transient.
|
|
if scanRecord.Status == atproto.ScanStatusSkipped {
|
|
sb.removeInflight(manifestDigest)
|
|
continue
|
|
}
|
|
|
|
scannedAt, err := time.Parse(time.RFC3339, scanRecord.ScannedAt)
|
|
if err != nil {
|
|
sb.removeInflight(manifestDigest)
|
|
continue
|
|
}
|
|
|
|
// Skip if scanned recently
|
|
if time.Since(scannedAt) < sb.rescanInterval {
|
|
sb.removeInflight(manifestDigest)
|
|
continue
|
|
}
|
|
|
|
candidate := &scanCandidate{
|
|
manifestDigest: manifestDigest,
|
|
userDID: scanRecord.UserDID,
|
|
scannedAt: scannedAt,
|
|
}
|
|
|
|
select {
|
|
case sb.staleQueue <- candidate:
|
|
found++
|
|
case <-sb.stopCh:
|
|
sb.removeInflight(manifestDigest)
|
|
return
|
|
}
|
|
}
|
|
|
|
if nextCursor == "" || len(records) == 0 {
|
|
break
|
|
}
|
|
cursor = nextCursor
|
|
}
|
|
|
|
if found > 0 {
|
|
slog.Info("Stale scan: pass complete", "staleCandidates", found)
|
|
}
|
|
}
|
|
|
|
// dispatchLoop pops candidates from the work queues with strict priority
|
|
// (unscanned before stale) and enqueues them as scan jobs, throttled by
|
|
// waitForProactiveCapacity to one proactive job per connected scanner worker.
|
|
func (sb *ScanBroadcaster) dispatchLoop() {
|
|
defer sb.wg.Done()
|
|
|
|
// Wait for system to settle
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
case <-time.After(15 * time.Second):
|
|
}
|
|
|
|
slog.Info("Dispatch loop started")
|
|
|
|
for {
|
|
select {
|
|
case <-sb.stopCh:
|
|
slog.Info("Dispatch loop stopped")
|
|
return
|
|
default:
|
|
}
|
|
|
|
// Wait until at least one scanner is connected. This comes first
|
|
// because the capacity gate is derived from what is connected: with
|
|
// nothing there the budget is zero and the gate has nothing to wait
|
|
// for.
|
|
if !sb.hasConnectedScanners() {
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
case <-time.After(5 * time.Second):
|
|
}
|
|
continue
|
|
}
|
|
|
|
// Wait until the connected scanners have room for another proactive
|
|
// job. Returns false on shutdown and on the fleet emptying out, both
|
|
// of which are handled by looping back to the top.
|
|
if !sb.waitForProactiveCapacity() {
|
|
continue
|
|
}
|
|
|
|
// Pop from highest-priority non-empty queue
|
|
var candidate *scanCandidate
|
|
|
|
select {
|
|
case c := <-sb.unscannedQueue:
|
|
candidate = c
|
|
default:
|
|
// No unscanned; try stale
|
|
select {
|
|
case c := <-sb.staleQueue:
|
|
candidate = c
|
|
default:
|
|
// Both queues empty — wait for completion signal or timeout
|
|
select {
|
|
case <-sb.stopCh:
|
|
return
|
|
case <-sb.completionSignal:
|
|
case <-time.After(30 * time.Second):
|
|
}
|
|
continue
|
|
}
|
|
}
|
|
|
|
sb.dispatchCandidate(candidate)
|
|
}
|
|
}
|
|
|
|
// waitForProactiveCapacity blocks until the connected scanners have room for
|
|
// another proactive job. Returns false when the caller should re-evaluate from
|
|
// the top instead — shutdown, or every scanner having gone away.
|
|
//
|
|
// It must not block on zero capacity: the budget is derived from the
|
|
// subscriber list, and a gate parked inside its own loop cannot notice a
|
|
// scanner arriving or the loop being asked to stop for up to five seconds.
|
|
func (sb *ScanBroadcaster) waitForProactiveCapacity() bool {
|
|
blockedSince := time.Now()
|
|
var lastWarn time.Time
|
|
|
|
for {
|
|
select {
|
|
case <-sb.stopCh:
|
|
return false
|
|
default:
|
|
}
|
|
|
|
limit := sb.proactiveDispatchLimit()
|
|
if limit == 0 {
|
|
return false
|
|
}
|
|
if n, ok := sb.activeProactiveJobs(); ok && n < limit {
|
|
return true
|
|
}
|
|
|
|
if blocked := time.Since(blockedSince); blocked >= capacityStallWarnAfter &&
|
|
time.Since(lastWarn) >= capacityStallWarnAfter {
|
|
sb.logStalledCapacity(blocked)
|
|
lastWarn = time.Now()
|
|
}
|
|
|
|
select {
|
|
case <-sb.stopCh:
|
|
return false
|
|
case <-sb.completionSignal:
|
|
// Job completed, check again
|
|
case <-time.After(5 * time.Second):
|
|
// Periodic check as fallback
|
|
}
|
|
}
|
|
}
|
|
|
|
// proactiveDispatchLimit is how many proactive scan jobs may be in flight at
|
|
// once: one per worker, summed over every connected scanner.
|
|
//
|
|
// The depth used to be one, hold-wide, which defeated both ways of scaling at
|
|
// the same time — a second worker in a scanner and a second scanner process
|
|
// were equally unable to receive proactive work. Deriving it from declared
|
|
// capacity is what makes `scanner.workers: 2` and a second scanner process
|
|
// mean something.
|
|
//
|
|
// One per worker and not more. The scanner acks on receipt and queues
|
|
// internally, so anything beyond one per worker is backlog the hold cannot
|
|
// see, cannot re-route to a process that frees up first, and cannot put an
|
|
// honest deadline on. Zero when nothing is connected, which is what stops the
|
|
// loop from manufacturing work for a fleet that is not there.
|
|
func (sb *ScanBroadcaster) proactiveDispatchLimit() int {
|
|
sb.mu.RLock()
|
|
defer sb.mu.RUnlock()
|
|
|
|
total := 0
|
|
for _, sub := range sb.subscribers {
|
|
total += sub.effectiveCapacity()
|
|
}
|
|
return total
|
|
}
|
|
|
|
// hasProactiveCapacity is the non-blocking form of the gate, for callers that
|
|
// want an answer rather than a wait.
|
|
func (sb *ScanBroadcaster) hasProactiveCapacity() bool {
|
|
limit := sb.proactiveDispatchLimit()
|
|
if limit == 0 {
|
|
return false
|
|
}
|
|
n, ok := sb.activeProactiveJobs()
|
|
return ok && n < limit
|
|
}
|
|
|
|
// dispatchCandidate resolves manifest details if needed and enqueues a scan job.
|
|
func (sb *ScanBroadcaster) dispatchCandidate(candidate *scanCandidate) {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
|
defer cancel()
|
|
|
|
// For stale candidates, verify the scan is still stale (may have been
|
|
// scanned by a push-triggered job while sitting in the queue)
|
|
if !candidate.scannedAt.IsZero() {
|
|
_, scanRecord, err := sb.pds.GetScanRecord(ctx, candidate.manifestDigest)
|
|
if err == nil {
|
|
scannedAt, parseErr := time.Parse(time.RFC3339, scanRecord.ScannedAt)
|
|
if parseErr == nil && time.Since(scannedAt) < sb.rescanInterval {
|
|
// Recently scanned, skip
|
|
sb.removeInflight(candidate.manifestDigest)
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
// Resolve manifest details if not already present (stale candidates lack Config/Layers)
|
|
if candidate.manifest.Config == nil {
|
|
if !sb.resolveManifestForCandidate(ctx, candidate) {
|
|
sb.removeInflight(candidate.manifestDigest)
|
|
return
|
|
}
|
|
}
|
|
|
|
configJSON, _ := json.Marshal(candidate.manifest.Config)
|
|
layersJSON, _ := json.Marshal(candidate.manifest.Layers)
|
|
|
|
reason := "never scanned"
|
|
if !candidate.scannedAt.IsZero() {
|
|
reason = fmt.Sprintf("last scanned %s ago", time.Since(candidate.scannedAt).Truncate(time.Minute))
|
|
}
|
|
|
|
slog.Info("Dispatching proactive scan",
|
|
"manifestDigest", candidate.manifestDigest,
|
|
"repository", candidate.manifest.Repository,
|
|
"userDID", candidate.userDID,
|
|
"reason", reason)
|
|
|
|
if err := sb.enqueue(&ScanJobEvent{
|
|
ManifestDigest: candidate.manifestDigest,
|
|
Repository: candidate.manifest.Repository,
|
|
UserDID: candidate.userDID,
|
|
UserHandle: candidate.userHandle,
|
|
Tier: "deckhand",
|
|
Config: configJSON,
|
|
Layers: layersJSON,
|
|
}, originProactive); err != nil {
|
|
slog.Error("Dispatch: failed to enqueue",
|
|
"manifest", candidate.manifestDigest, "error", err)
|
|
// removeInflight not needed — Enqueue already cleans up on error
|
|
}
|
|
}
|
|
|
|
// resolveManifestForCandidate fetches the manifest from the user's PDS to populate
|
|
// Config and Layers fields needed for the scan job.
|
|
func (sb *ScanBroadcaster) resolveManifestForCandidate(ctx context.Context, candidate *scanCandidate) bool {
|
|
did, userHandle, pdsEndpoint, err := atproto.ResolveIdentity(ctx, candidate.userDID)
|
|
if err != nil {
|
|
slog.Debug("Dispatch: failed to resolve user identity",
|
|
"userDID", candidate.userDID, "error", err)
|
|
return false
|
|
}
|
|
if candidate.userHandle == "" {
|
|
candidate.userHandle = userHandle
|
|
}
|
|
|
|
client := atproto.NewClient(pdsEndpoint, did, "")
|
|
var cursor string
|
|
for {
|
|
records, nextCursor, err := client.ListRecordsForRepo(ctx, did, atproto.ManifestCollection, 100, cursor)
|
|
if err != nil {
|
|
return false
|
|
}
|
|
|
|
for _, record := range records {
|
|
var manifest atproto.ManifestRecord
|
|
if err := json.Unmarshal(record.Value, &manifest); err != nil {
|
|
continue
|
|
}
|
|
|
|
if manifest.Digest == candidate.manifestDigest {
|
|
candidate.manifest = manifest
|
|
return true
|
|
}
|
|
}
|
|
|
|
if nextCursor == "" || len(records) == 0 {
|
|
break
|
|
}
|
|
cursor = nextCursor
|
|
}
|
|
|
|
slog.Debug("Dispatch: manifest not found on user's PDS",
|
|
"manifestDigest", candidate.manifestDigest,
|
|
"userDID", candidate.userDID)
|
|
return false
|
|
}
|
|
|
|
// isOurManifest checks if a manifest's holdDID matches this hold directly,
|
|
// or if the manifest's hold has been migrated (has a successor label set).
|
|
//
|
|
// An answer we could not get is not an answer of "no". A predecessor is by
|
|
// definition a retired hold, so being briefly unreachable is its normal
|
|
// condition, and this used to cache the resulting false forever: one slow reply
|
|
// during the first discovery pass after boot meant every manifest naming that
|
|
// hold went unscanned for the life of the process, silently, at Debug level.
|
|
// pkg/hold/gc hit the same defect with a worse blast radius (deleted blobs
|
|
// rather than missed scans) and fixed it with the definitive/unresolved split
|
|
// this follows.
|
|
func (sb *ScanBroadcaster) isOurManifest(ctx context.Context, holdDID string) bool {
|
|
if holdDID == "" {
|
|
return false
|
|
}
|
|
|
|
// Direct match
|
|
if holdDID == sb.holdDID {
|
|
return true
|
|
}
|
|
|
|
// Held across the fetch. The lock is uncontended today, and holding it also
|
|
// means two callers can never dial the same unreachable hold at once.
|
|
sb.predecessorMu.Lock()
|
|
defer sb.predecessorMu.Unlock()
|
|
|
|
if isPredecessor, cached := sb.predecessorCache[holdDID]; cached {
|
|
return isPredecessor
|
|
}
|
|
|
|
// Already unreachable earlier in this pass. Answer the same way without
|
|
// paying another timeout; the next pass starts fresh and re-checks.
|
|
if sb.predecessorUnresolved[holdDID] {
|
|
return false
|
|
}
|
|
|
|
// Fetch captain record from the other hold's PDS to check successor
|
|
isPredecessor, definitive := sb.checkPredecessor(ctx, holdDID)
|
|
if !definitive {
|
|
if sb.predecessorUnresolved == nil {
|
|
sb.predecessorUnresolved = make(map[string]bool)
|
|
}
|
|
sb.predecessorUnresolved[holdDID] = true
|
|
slog.Warn("Proactive scan: predecessor status unresolved, not adopting this hold's manifests this pass",
|
|
"holdDID", holdDID)
|
|
return false
|
|
}
|
|
|
|
if sb.predecessorCache == nil {
|
|
sb.predecessorCache = make(map[string]bool)
|
|
}
|
|
sb.predecessorCache[holdDID] = isPredecessor
|
|
return isPredecessor
|
|
}
|
|
|
|
// checkPredecessor fetches a hold's captain record to check if it has a successor label
|
|
// (meaning the hold has been migrated/retired and its manifests should be scanned by us).
|
|
//
|
|
// The second return value reports whether the answer is definitive. It is false
|
|
// whenever the hold could not be reached or its reply could not be understood:
|
|
// those are cases where the hold may well be a predecessor and we simply cannot
|
|
// tell. Callers must not read an inconclusive result as "not a predecessor".
|
|
func (sb *ScanBroadcaster) checkPredecessor(ctx context.Context, holdDID string) (isPredecessor, definitive bool) {
|
|
fetchCtx, cancel := context.WithTimeout(ctx, 5*time.Second)
|
|
defer cancel()
|
|
|
|
holdURL, err := atproto.ResolveHoldURL(fetchCtx, holdDID)
|
|
if err != nil {
|
|
slog.Debug("Proactive scan: failed to resolve predecessor hold URL",
|
|
"holdDID", holdDID, "error", err)
|
|
return false, false
|
|
}
|
|
|
|
return sb.checkPredecessorAt(fetchCtx, holdDID, holdURL)
|
|
}
|
|
|
|
// checkPredecessorAt is checkPredecessor with the hold's base URL already
|
|
// resolved, split out so the fetch-and-parse half can be exercised against a
|
|
// local server. It carries the same contract: the second return value is false
|
|
// whenever the answer is inconclusive rather than negative.
|
|
//
|
|
// A non-200 is inconclusive rather than negative on the same reasoning as
|
|
// pkg/hold/gc: a reachable service that cannot produce its own captain record
|
|
// is malfunctioning, not answering.
|
|
func (sb *ScanBroadcaster) checkPredecessorAt(ctx context.Context, holdDID, holdURL string) (isPredecessor, definitive bool) {
|
|
// Fetch captain record: com.atproto.repo.getRecord
|
|
recordURL := fmt.Sprintf("%s/xrpc/com.atproto.repo.getRecord?repo=%s&collection=%s&rkey=self",
|
|
holdURL,
|
|
url.QueryEscape(holdDID),
|
|
url.QueryEscape(atproto.CaptainCollection),
|
|
)
|
|
|
|
req, err := http.NewRequestWithContext(ctx, "GET", recordURL, nil)
|
|
if err != nil {
|
|
return false, false
|
|
}
|
|
|
|
resp, err := http.DefaultClient.Do(req)
|
|
if err != nil {
|
|
slog.Debug("Proactive scan: failed to fetch predecessor captain record",
|
|
"holdDID", holdDID, "error", err)
|
|
return false, false
|
|
}
|
|
defer resp.Body.Close()
|
|
|
|
if resp.StatusCode != http.StatusOK {
|
|
slog.Debug("Proactive scan: predecessor captain record fetch returned non-200",
|
|
"holdDID", holdDID, "status", resp.StatusCode)
|
|
return false, false
|
|
}
|
|
|
|
body, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) // 1MB limit
|
|
if err != nil {
|
|
slog.Debug("Proactive scan: failed to read predecessor captain record",
|
|
"holdDID", holdDID, "error", err)
|
|
return false, false
|
|
}
|
|
|
|
var envelope struct {
|
|
Value json.RawMessage `json:"value"`
|
|
}
|
|
if err := json.Unmarshal(body, &envelope); err != nil {
|
|
slog.Debug("Proactive scan: failed to parse predecessor captain envelope",
|
|
"holdDID", holdDID, "error", err)
|
|
return false, false
|
|
}
|
|
|
|
var captain atproto.CaptainRecord
|
|
if err := json.Unmarshal(envelope.Value, &captain); err != nil {
|
|
slog.Debug("Proactive scan: failed to parse predecessor captain record",
|
|
"holdDID", holdDID, "error", err)
|
|
return false, false
|
|
}
|
|
|
|
// The hold answered and declares no successor. This is the one negative we
|
|
// are entitled to cache.
|
|
if captain.Successor == "" {
|
|
return false, true
|
|
}
|
|
|
|
if captain.Successor != sb.holdDID {
|
|
slog.Debug("Proactive scan: hold has successor, but it is not us",
|
|
"holdDID", holdDID, "successor", captain.Successor, "ourHoldDID", sb.holdDID)
|
|
return false, true
|
|
}
|
|
|
|
slog.Info("Proactive scan: discovered migrated hold pointing at us as successor",
|
|
"holdDID", holdDID, "successor", captain.Successor)
|
|
return true, true
|
|
}
|
|
|
|
// hasConnectedScanners returns true if at least one scanner is connected.
|
|
func (sb *ScanBroadcaster) hasConnectedScanners() bool {
|
|
sb.mu.RLock()
|
|
defer sb.mu.RUnlock()
|
|
return len(sb.subscribers) > 0
|
|
}
|
|
|
|
// activeProactiveJobs counts the proactive scan jobs currently holding
|
|
// dispatch capacity. The second return reports whether the count can be acted
|
|
// on; a caller must not read (0, false) as "idle".
|
|
//
|
|
// Only proactive rows are counted. Push-triggered scans do not go through this
|
|
// gate at all — oci/xrpc.go calls Enqueue directly, because a user who just
|
|
// pushed is waiting — so counting them meant a hold with steady pushes never
|
|
// dispatched a proactive scan, and "one proactive job at a time" was really
|
|
// "none while anyone is pushing".
|
|
//
|
|
// Pending rows older than pendingStaleAfter are deliberately not counted. A job
|
|
// that has been pending that long is one no scanner can be given (dispatch
|
|
// failed, or no scanner is connected), and counting it blocks the proactive
|
|
// dispatch loop for as long as the row exists — which is how a single
|
|
// undispatchable job stopped scanning deployment-wide. The re-dispatch loop
|
|
// keeps re-offering such rows, so ignoring them here costs nothing when a
|
|
// scanner is available.
|
|
//
|
|
// A query failure is not evidence of activity. The first few are answered
|
|
// "busy" anyway, because guessing "idle" during a blip piles another job onto
|
|
// a scanner that may already have one. A persistent failure is different: this
|
|
// is the only gate on proactive dispatch and waitForProactiveCapacity spins on
|
|
// it, so answering "busy" forever halts scanning for the life of the process
|
|
// with no recovery and nothing but one log line every five seconds to show for
|
|
// it — logStalledCapacity runs the same database and takes its own error
|
|
// branch. After activeJobsErrorBudget consecutive failures it therefore fails
|
|
// open, reporting an idle fleet so dispatch resumes.
|
|
func (sb *ScanBroadcaster) activeProactiveJobs() (int, bool) {
|
|
var count int
|
|
err := sb.db.QueryRow(`
|
|
SELECT COUNT(*) FROM scan_jobs
|
|
WHERE origin = ?
|
|
AND (status IN ('assigned', 'processing')
|
|
OR (status = 'pending' AND datetime(created_at) > datetime('now', ?)))
|
|
`, originProactive, sqliteAgoModifier(pendingStaleAfter)).Scan(&count)
|
|
if err != nil {
|
|
failures := sb.activeJobsErrs.Add(1)
|
|
if failures <= activeJobsErrorBudget {
|
|
slog.Error("Failed to check active scan jobs, assuming busy",
|
|
"error", err, "consecutiveFailures", failures)
|
|
return 0, false
|
|
}
|
|
slog.Error("Failed to check active scan jobs; proceeding as if idle so a "+
|
|
"database fault does not halt proactive scanning outright",
|
|
"error", err, "consecutiveFailures", failures)
|
|
return 0, true
|
|
}
|
|
sb.activeJobsErrs.Store(0)
|
|
return count, true
|
|
}
|
|
|
|
// logStalledCapacity reports what is holding the dispatch loop back, so a stall
|
|
// shows up in the hold's logs instead of only as an absence of scan results.
|
|
func (sb *ScanBroadcaster) logStalledCapacity(blockedFor time.Duration) {
|
|
var (
|
|
count int
|
|
oldestSeq sql.NullInt64
|
|
statuses sql.NullString
|
|
)
|
|
err := sb.db.QueryRow(`
|
|
SELECT COUNT(*), MIN(seq), GROUP_CONCAT(DISTINCT status)
|
|
FROM scan_jobs
|
|
WHERE origin = ? AND status IN ('pending', 'assigned', 'processing')
|
|
`, originProactive).Scan(&count, &oldestSeq, &statuses)
|
|
if err != nil {
|
|
slog.Warn("Proactive scan dispatch stalled; could not inspect active jobs",
|
|
"blockedFor", blockedFor.Truncate(time.Minute), "error", err)
|
|
return
|
|
}
|
|
|
|
slog.Warn("Proactive scan dispatch stalled waiting on active jobs",
|
|
"blockedFor", blockedFor.Truncate(time.Minute),
|
|
"activeProactiveJobs", count,
|
|
"dispatchLimit", sb.proactiveDispatchLimit(),
|
|
"oldestSeq", oldestSeq.Int64,
|
|
"statuses", statuses.String)
|
|
}
|
|
|
|
func generateSubscriberID() string {
|
|
b := make([]byte, 8)
|
|
_, _ = rand.Read(b)
|
|
return hex.EncodeToString(b)
|
|
}
|
|
|
|
// addInflight marks a manifest digest as in-flight. Returns false if already present.
|
|
func (sb *ScanBroadcaster) addInflight(digest string) bool {
|
|
sb.inflightMu.Lock()
|
|
defer sb.inflightMu.Unlock()
|
|
if _, ok := sb.inflight[digest]; ok {
|
|
return false
|
|
}
|
|
sb.inflight[digest] = struct{}{}
|
|
return true
|
|
}
|
|
|
|
// removeInflight removes a manifest digest from the in-flight set.
|
|
func (sb *ScanBroadcaster) removeInflight(digest string) {
|
|
sb.inflightMu.Lock()
|
|
defer sb.inflightMu.Unlock()
|
|
delete(sb.inflight, digest)
|
|
}
|
|
|
|
// signalCompletion non-blocking signal to wake the dispatch loop, and the
|
|
// re-dispatch loop with it: a finished job frees a slot on some scanner, and
|
|
// what is waiting for that slot may be a push-triggered row that the proactive
|
|
// dispatch loop will never look at.
|
|
func (sb *ScanBroadcaster) signalCompletion() {
|
|
select {
|
|
case sb.completionSignal <- struct{}{}:
|
|
default:
|
|
}
|
|
select {
|
|
case sb.capacityFreed <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// triggerDiscovery non-blocking signal to trigger an early discovery pass.
|
|
func (sb *ScanBroadcaster) triggerDiscovery() {
|
|
select {
|
|
case sb.discoverNow <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|