mirror of
https://tangled.org/evan.jarrett.net/at-container-registry
synced 2026-09-29 21:45:33 +00:00
scanner: fix five crash and halt classes found by a pipeline audit
An audit of the scan pipeline and the hold side of scanning found several ways scanning stops without saying so. Each fix here was written test-first: a test expressing the wanted behaviour, confirmed failing for the right reason, then the change. A summary-less result crash-looped both processes. worker.go dereferenced result.Summary unconditionally, but processJob only sets it when Grype runs, and SendResult puts the nil on the wire before the scanner dies on it, so handleResult's unguarded log killed the hold too. A nil Summary now means "not scanned for vulnerabilities", deliberately distinct from "scanned, found zero" — inventing a zeroed summary would report every image as clean when Grype never ran. The hold writes a record rather than orphaning the uploaded SBOM, and the appview renders an "SBOM only" state instead of a green Clean badge. The Grype database could wedge with no way back short of a restart. All three throttles in loadVulnDatabase were guarded by vulnDB != nil, so a scanner holding no provider retried a full download on every scan under the exclusive lock. Two earlier attempts at this bug each added one more condition to the same chain; this replaces the chain with a single decision function over a state snapshot, consulted by both call sites so they cannot disagree. That disagreement was itself a bug: the 50-scan reload had never once executed. Two independent halts. An unparseable frame was dropped in silence, stranding a row that held the hold's only dispatch slot forever; it is now answered "skipped" on first delivery. The 10-minute sweep leaked the in-flight digest and wrote no record, permanently retiring one image per timeout. A digest went unvalidated into filepath.Join and os.Create, so a layer digest of sha256:../../../x wrote outside the scan directory, and nothing verified that downloaded bytes hashed to the digest naming them. Digests come from records in a user's own PDS. Both are fixed together: verification is what makes an escaping write self-defeating. Concurrency did not work on either axis. The proactive capacity gate was depth-one hold-wide, so neither extra workers nor extra scanner processes received work. Depth is now the sum of the worker counts scanners advertise on connect, the gate is scoped to proactive work, and dispatch prefers the least-loaded scanner. Disconnects no longer hand a running scan to someone else: a scanner keeps a stable per-process identity and reclaims its own rows within a grace window, while a process that truly restarted returns with a new identity and has its work reclaimed, which is correct because the restart did lose it. The hold's scanning deadline measured queueing rather than scanning, because the scanner acks on receipt and handleAck never refreshed assigned_at. A new "started" message, sent by the worker that dequeues the job, separates the two budgets. An older scanner never sends it and falls under the queueing budget, which is more forgiving than the deadline it gets today. Adds an in-process mock hold and an e2e harness that runs the real client, queue and worker pool, seeded with 84 real manifest records fetched from a live PDS. Real image layouts and the Grype database are fetched by scripts and gitignored; suites needing them skip cleanly, so the default run stays offline and fast. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01U1Km3N3uUmeGaj7VbaM8PF
This commit is contained in:
co-authored by
Claude Opus 5
parent
f16a8eaa82
commit
a63f668de0
@@ -469,8 +469,15 @@ func generateAdvisorPrompt(w io.Writer, r *advisorReportData) {
|
||||
// Vulnerability summary
|
||||
if r.ScanRecord != nil {
|
||||
sr := r.ScanRecord
|
||||
fmt.Fprintf(w, "vulns: {critical: %d, high: %d, medium: %d, low: %d, total: %d}\n",
|
||||
sr.Critical, sr.High, sr.Medium, sr.Low, sr.Total)
|
||||
if vulnScanDidNotRun(sr) {
|
||||
// Counts of zero here would be read as "no vulnerabilities", but
|
||||
// this record was written by a scan that never ran a vulnerability
|
||||
// database against the image.
|
||||
fmt.Fprintf(w, "vulns: not scanned (SBOM only, no vulnerability data)\n")
|
||||
} else {
|
||||
fmt.Fprintf(w, "vulns: {critical: %d, high: %d, medium: %d, low: %d, total: %d}\n",
|
||||
sr.Critical, sr.High, sr.Medium, sr.Low, sr.Total)
|
||||
}
|
||||
}
|
||||
|
||||
// Fixable critical/high vulns
|
||||
|
||||
@@ -25,50 +25,74 @@ type ScanResultHandler struct {
|
||||
}
|
||||
|
||||
// vulnBadgeData is the template data for the vuln-badge partial.
|
||||
// The badge renders one of five states, in priority order:
|
||||
// 1. Error — we couldn't reach the hold at all (network/5xx)
|
||||
// 2. NotScanned — hold reachable, no scan record for this digest (404)
|
||||
// 3. Skipped — scan record explicitly marks this artifact as not-scannable
|
||||
// 4. ScanFailed — scan record exists but the scanner errored
|
||||
// 5. Found — scan succeeded; render tier counts (or "Clean" when zero)
|
||||
// The badge renders one of six states, in priority order:
|
||||
// 1. Error — we couldn't reach the hold at all (network/5xx)
|
||||
// 2. NotScanned — hold reachable, no scan record for this digest (404)
|
||||
// 3. Skipped — scan record explicitly marks this artifact as not-scannable
|
||||
// 4. ScanFailed — scan record exists but the scanner errored
|
||||
// 5. VulnsNotScanned — scan succeeded but no vulnerability data was produced
|
||||
// 6. Found — scan succeeded; render tier counts (or "Clean" when zero)
|
||||
//
|
||||
// These states must stay distinct so users can tell "hold is down" from
|
||||
// "this hasn't been scanned yet" from "scanner errored on this image" from
|
||||
// "this artifact type is intentionally not scanned".
|
||||
// "this artifact type is intentionally not scanned" from "we catalogued the
|
||||
// image but never matched it against a vulnerability database".
|
||||
type vulnBadgeData struct {
|
||||
Critical int64
|
||||
High int64
|
||||
Medium int64
|
||||
Low int64
|
||||
Total int64
|
||||
ScannedAt string
|
||||
Found bool // true if scan record exists and succeeded
|
||||
Error bool // true if hold unreachable (network/5xx)
|
||||
NotScanned bool // true if hold is up but no scan record (404)
|
||||
ScanFailed bool // true if scan record exists but scan failed
|
||||
Skipped bool // true if scan record marks the artifact as intentionally not scanned (helm, in-toto, etc.)
|
||||
Digest string // for the detail modal link
|
||||
HoldEndpoint string // for the detail modal link
|
||||
Critical int64
|
||||
High int64
|
||||
Medium int64
|
||||
Low int64
|
||||
Total int64
|
||||
ScannedAt string
|
||||
Found bool // true if scan record exists and succeeded
|
||||
Error bool // true if hold unreachable (network/5xx)
|
||||
NotScanned bool // true if hold is up but no scan record (404)
|
||||
ScanFailed bool // true if scan record exists but scan failed
|
||||
Skipped bool // true if scan record marks the artifact as intentionally not scanned (helm, in-toto, etc.)
|
||||
// VulnsNotScanned means the scan produced an SBOM but no vulnerability
|
||||
// data, which is what a scanner running with vulnerability scanning off
|
||||
// reports. Its zero counts are the absence of a measurement, not a finding,
|
||||
// so the badge must not render them as "Clean".
|
||||
VulnsNotScanned bool
|
||||
Digest string // for the detail modal link
|
||||
HoldEndpoint string // for the detail modal link
|
||||
}
|
||||
|
||||
// vulnScanDidNotRun reports whether a successful scan record carries no
|
||||
// vulnerability data at all: no report blob and no counts. That is what the
|
||||
// hold writes when the scanner ran with vulnerability scanning disabled, and
|
||||
// its zeros must never be presented as "no vulnerabilities found".
|
||||
//
|
||||
// It deliberately requires an explicit "ok" status. Records written before the
|
||||
// vulnReportBlob field existed carry counts and an SBOM but no status, and
|
||||
// those really were scanned, so the legacy shape stays out of this branch.
|
||||
func vulnScanDidNotRun(scanRecord *atproto.ScanRecord) bool {
|
||||
return scanRecord.Status == atproto.ScanStatusOK &&
|
||||
scanRecord.VulnReportBlob == nil &&
|
||||
scanRecord.Total == 0
|
||||
}
|
||||
|
||||
// classifyScanRecord maps a scan record's Status field to badge data flags.
|
||||
// An empty Status is treated as a legacy record from before the status field
|
||||
// existed: nil-blob + zero-counts = treat as failed (preserves the prior badge
|
||||
// for un-backfilled holds); otherwise treat as success.
|
||||
func classifyScanRecord(scanRecord *atproto.ScanRecord) (found, skipped, failed bool) {
|
||||
func classifyScanRecord(scanRecord *atproto.ScanRecord) (found, skipped, failed, vulnsNotScanned bool) {
|
||||
switch scanRecord.Status {
|
||||
case atproto.ScanStatusSkipped:
|
||||
return false, true, false
|
||||
return false, true, false, false
|
||||
case atproto.ScanStatusFailed:
|
||||
return false, false, true
|
||||
return false, false, true, false
|
||||
case atproto.ScanStatusOK:
|
||||
return true, false, false
|
||||
if vulnScanDidNotRun(scanRecord) {
|
||||
return false, false, false, true
|
||||
}
|
||||
return true, false, false, false
|
||||
default:
|
||||
// Legacy record (status field didn't exist when this was written).
|
||||
if scanRecord.SbomBlob == nil && scanRecord.Total == 0 {
|
||||
return false, false, true
|
||||
return false, false, true, false
|
||||
}
|
||||
return true, false, false
|
||||
return true, false, false, false
|
||||
}
|
||||
}
|
||||
|
||||
@@ -146,19 +170,21 @@ func (h *ScanResultHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
found, skipped, failed := classifyScanRecord(&scanRecord)
|
||||
found, skipped, failed, vulnsNotScanned := classifyScanRecord(&scanRecord)
|
||||
h.renderBadge(w, vulnBadgeData{
|
||||
Critical: scanRecord.Critical,
|
||||
High: scanRecord.High,
|
||||
Medium: scanRecord.Medium,
|
||||
Low: scanRecord.Low,
|
||||
Total: scanRecord.Total,
|
||||
ScannedAt: scanRecord.ScannedAt,
|
||||
Found: found,
|
||||
Skipped: skipped,
|
||||
ScanFailed: failed,
|
||||
Digest: digest,
|
||||
HoldEndpoint: holdDID,
|
||||
Critical: scanRecord.Critical,
|
||||
High: scanRecord.High,
|
||||
Medium: scanRecord.Medium,
|
||||
Low: scanRecord.Low,
|
||||
Total: scanRecord.Total,
|
||||
ScannedAt: scanRecord.ScannedAt,
|
||||
Found: found,
|
||||
Skipped: skipped,
|
||||
ScanFailed: failed,
|
||||
|
||||
VulnsNotScanned: vulnsNotScanned,
|
||||
Digest: digest,
|
||||
HoldEndpoint: holdDID,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -211,19 +237,21 @@ func fetchScanRecord(ctx context.Context, holdEndpoint, holdDID, hexDigest strin
|
||||
return vulnBadgeData{Error: true}
|
||||
}
|
||||
|
||||
found, skipped, failed := classifyScanRecord(&scanRecord)
|
||||
found, skipped, failed, vulnsNotScanned := classifyScanRecord(&scanRecord)
|
||||
return vulnBadgeData{
|
||||
Critical: scanRecord.Critical,
|
||||
High: scanRecord.High,
|
||||
Medium: scanRecord.Medium,
|
||||
Low: scanRecord.Low,
|
||||
Total: scanRecord.Total,
|
||||
ScannedAt: scanRecord.ScannedAt,
|
||||
Found: found,
|
||||
Skipped: skipped,
|
||||
ScanFailed: failed,
|
||||
Digest: fullDigest,
|
||||
HoldEndpoint: holdDID,
|
||||
Critical: scanRecord.Critical,
|
||||
High: scanRecord.High,
|
||||
Medium: scanRecord.Medium,
|
||||
Low: scanRecord.Low,
|
||||
Total: scanRecord.Total,
|
||||
ScannedAt: scanRecord.ScannedAt,
|
||||
Found: found,
|
||||
Skipped: skipped,
|
||||
ScanFailed: failed,
|
||||
|
||||
VulnsNotScanned: vulnsNotScanned,
|
||||
Digest: fullDigest,
|
||||
HoldEndpoint: holdDID,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -503,3 +503,95 @@ func TestBatchScanResult_SingleDigest(t *testing.T) {
|
||||
t.Error("Expected critical count of 1")
|
||||
}
|
||||
}
|
||||
|
||||
// mockSBOMOnlyScanRecord is what the hold writes when the scanner finished a
|
||||
// scan with vulnerability scanning turned off: status "ok", an SBOM blob, no
|
||||
// vulnerability report, and zero counts because Grype never ran.
|
||||
//
|
||||
// The zero counts are not a finding. Rendering them as "Clean" would tell every
|
||||
// user of a vuln-disabled scanner that their images have no vulnerabilities,
|
||||
// which is the misreading this record shape has to avoid.
|
||||
func mockSBOMOnlyScanRecord() string {
|
||||
record := map[string]any{
|
||||
"$type": "io.atcr.hold.scan",
|
||||
"manifest": "at://did:plc:test/io.atcr.manifest/abc123",
|
||||
"repository": "myapp",
|
||||
"userDid": "did:plc:test",
|
||||
"status": "ok",
|
||||
"critical": 0,
|
||||
"high": 0,
|
||||
"medium": 0,
|
||||
"low": 0,
|
||||
"total": 0,
|
||||
"sbomBlob": map[string]any{
|
||||
"$type": "blob",
|
||||
"ref": map[string]any{"$link": "bafkreigv3xw47pk7cbeahkmttetf4smxyluwlu3jmteo2nzke2oa7dbhhm"},
|
||||
"mimeType": "application/spdx+json",
|
||||
"size": 1234,
|
||||
},
|
||||
"scannerVersion": "atcr-scanner-v1.0.0",
|
||||
"scannedAt": "2025-01-15T10:30:00Z",
|
||||
}
|
||||
envelope := map[string]any{
|
||||
"uri": "at://did:web:hold.example.com/io.atcr.hold.scan/abc123",
|
||||
"cid": "bafyreiabc123",
|
||||
"value": record,
|
||||
}
|
||||
b, _ := json.Marshal(envelope)
|
||||
return string(b)
|
||||
}
|
||||
|
||||
func TestScanResult_SBOMWithoutVulnScanIsNotClean(t *testing.T) {
|
||||
hold := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if handleMockDID(w, r) {
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.Write([]byte(mockSBOMOnlyScanRecord()))
|
||||
}))
|
||||
defer hold.Close()
|
||||
|
||||
handler := setupScanResultHandler(t, hold.URL)
|
||||
|
||||
req := httptest.NewRequest("GET", "/api/scan-result?digest=sha256:abc123&holdEndpoint="+hold.URL, nil)
|
||||
rr := httptest.NewRecorder()
|
||||
handler.ServeHTTP(rr, req)
|
||||
|
||||
body := rr.Body.String()
|
||||
|
||||
if strings.Contains(body, "Clean") || strings.Contains(body, "badge-success") {
|
||||
t.Errorf("a record with no vulnerability data was rendered as clean: %s", body)
|
||||
}
|
||||
if strings.Contains(body, "vuln-strip") {
|
||||
t.Errorf("a record with no vulnerability data rendered severity counts: %s", body)
|
||||
}
|
||||
if !strings.Contains(body, "SBOM only") {
|
||||
t.Errorf("expected the SBOM-only badge, got: %s", body)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanResult_LegacyCleanRecordStillReadsClean guards the discriminator from
|
||||
// the other side. Records written before the vulnReportBlob field existed carry
|
||||
// an SBOM, zero counts and no status, and they really were scanned clean, so
|
||||
// the "no vulnerability data" rule must key on the explicit status "ok" and not
|
||||
// swallow them.
|
||||
func TestScanResult_LegacyCleanRecordStillReadsClean(t *testing.T) {
|
||||
hold := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if handleMockDID(w, r) {
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.Write([]byte(mockScanRecord(0, 0, 0, 0, 0)))
|
||||
}))
|
||||
defer hold.Close()
|
||||
|
||||
handler := setupScanResultHandler(t, hold.URL)
|
||||
|
||||
req := httptest.NewRequest("GET", "/api/scan-result?digest=sha256:abc123&holdEndpoint="+hold.URL, nil)
|
||||
rr := httptest.NewRecorder()
|
||||
handler.ServeHTTP(rr, req)
|
||||
|
||||
if body := rr.Body.String(); !strings.Contains(body, "Clean") {
|
||||
t.Errorf("legacy zero-count record no longer reads clean: %s", body)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -60,12 +60,17 @@ type vulnDetailsData struct {
|
||||
// NotScanned means the hold answered but holds no scan record for this
|
||||
// manifest. Distinct from Error: nothing failed, the image was simply
|
||||
// never scanned, and the UI must not present it as a failure.
|
||||
NotScanned bool
|
||||
Status string // scan record's status field (ok | failed | skipped); empty for legacy records
|
||||
Reason string // scan record's reason field (only meaningful when Status != ok)
|
||||
ScannedAt string
|
||||
Digest string // image digest (for download URLs)
|
||||
HoldEndpoint string // hold DID (for download URLs)
|
||||
NotScanned bool
|
||||
// VulnsNotScanned means the scan succeeded and produced an SBOM, but no
|
||||
// vulnerability data: the scanner ran with vulnerability scanning off. The
|
||||
// zero counts are an absence of measurement, so the panel must say that
|
||||
// rather than report zero findings or blame a failed fetch.
|
||||
VulnsNotScanned bool
|
||||
Status string // scan record's status field (ok | failed | skipped); empty for legacy records
|
||||
Reason string // scan record's reason field (only meaningful when Status != ok)
|
||||
ScannedAt string
|
||||
Digest string // image digest (for download URLs)
|
||||
HoldEndpoint string // hold DID (for download URLs)
|
||||
}
|
||||
|
||||
type vulnMatch struct {
|
||||
@@ -171,6 +176,18 @@ func (h *VulnDetailsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
Total: scanRecord.Total,
|
||||
}
|
||||
|
||||
// A successful record with no report blob and no counts is a scan that
|
||||
// never ran Grype. Distinct from a missing report: there is nothing to
|
||||
// fetch and nothing went wrong.
|
||||
if vulnScanDidNotRun(&scanRecord) {
|
||||
h.renderDetails(w, vulnDetailsData{
|
||||
ScannedAt: scanRecord.ScannedAt,
|
||||
Status: scanRecord.Status,
|
||||
VulnsNotScanned: true,
|
||||
})
|
||||
return
|
||||
}
|
||||
|
||||
// Step 2: Fetch the vulnerability report blob
|
||||
if scanRecord.VulnReportBlob == nil || scanRecord.VulnReportBlob.Ref.String() == "" {
|
||||
h.renderDetails(w, vulnDetailsData{
|
||||
@@ -343,6 +360,17 @@ func FetchVulnDetails(ctx context.Context, holdEndpoint, digest string) vulnDeta
|
||||
}
|
||||
}
|
||||
|
||||
// A successful record with no report blob and no counts is a scan that
|
||||
// never ran Grype. Distinct from a missing report: there is nothing to
|
||||
// fetch and nothing went wrong.
|
||||
if vulnScanDidNotRun(&scanRecord) {
|
||||
return vulnDetailsData{
|
||||
ScannedAt: scanRecord.ScannedAt,
|
||||
Status: scanRecord.Status,
|
||||
VulnsNotScanned: true,
|
||||
}
|
||||
}
|
||||
|
||||
// Fetch the vulnerability report blob
|
||||
if scanRecord.VulnReportBlob == nil || scanRecord.VulnReportBlob.Ref.String() == "" {
|
||||
return vulnDetailsData{
|
||||
|
||||
@@ -356,3 +356,63 @@ func TestVulnDetails_MissingParams(t *testing.T) {
|
||||
t.Error("Expected error message for missing parameters")
|
||||
}
|
||||
}
|
||||
|
||||
// mockSBOMOnlyRecordEnvelope is the detail-modal counterpart of the badge
|
||||
// fixture in scan_result_test.go: status "ok", an SBOM, no vulnerability report
|
||||
// and zero counts, which is what a scanner running with vuln.enabled=false
|
||||
// produces.
|
||||
func mockSBOMOnlyRecordEnvelope() string {
|
||||
record := map[string]any{
|
||||
"$type": "io.atcr.hold.scan",
|
||||
"manifest": "at://did:plc:test/io.atcr.manifest/abc123",
|
||||
"repository": "myapp",
|
||||
"userDid": "did:plc:test",
|
||||
"status": "ok",
|
||||
"critical": 0,
|
||||
"high": 0,
|
||||
"medium": 0,
|
||||
"low": 0,
|
||||
"total": 0,
|
||||
"scannerVersion": "atcr-scanner-v1.0.0",
|
||||
"scannedAt": "2025-01-15T10:30:00Z",
|
||||
}
|
||||
envelope := map[string]any{
|
||||
"uri": "at://did:web:hold.example.com/io.atcr.hold.scan/abc123",
|
||||
"cid": "bafyreiabc123",
|
||||
"value": record,
|
||||
}
|
||||
b, _ := json.Marshal(envelope)
|
||||
return string(b)
|
||||
}
|
||||
|
||||
// TestVulnDetails_VulnScanDidNotRun is the modal's half of the same rule: a
|
||||
// record with no vulnerability data must say so, not report zero findings and
|
||||
// not blame the hold for a fetch that never happened.
|
||||
func TestVulnDetails_VulnScanDidNotRun(t *testing.T) {
|
||||
hold := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if handleMockDID(w, r) {
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
w.Write([]byte(mockSBOMOnlyRecordEnvelope()))
|
||||
}))
|
||||
defer hold.Close()
|
||||
|
||||
handler := setupVulnDetailsHandler(t)
|
||||
|
||||
req := httptest.NewRequest("GET", "/api/vuln-details?digest=sha256:abc123&holdEndpoint="+hold.URL, nil)
|
||||
rr := httptest.NewRecorder()
|
||||
handler.ServeHTTP(rr, req)
|
||||
|
||||
body := rr.Body.String()
|
||||
|
||||
if !strings.Contains(body, "Vulnerability scanning did not run") {
|
||||
t.Errorf("expected copy saying vulnerability scanning did not run, got: %s", body)
|
||||
}
|
||||
if strings.Contains(body, "0 vulnerabilities") {
|
||||
t.Errorf("a record with no vulnerability data claimed zero vulnerabilities: %s", body)
|
||||
}
|
||||
if strings.Contains(body, "No detailed vulnerability report") {
|
||||
t.Errorf("a scan that never ran Grype was reported as a missing report: %s", body)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,6 +14,11 @@
|
||||
<span></span>
|
||||
{{ else if .ScanFailed }}
|
||||
<span class="badge badge-sm badge-warning" title="Scanner ran but produced no SBOM">{{ icon "alert-triangle" "size-3" }} Scan failed</span>
|
||||
{{ else if .VulnsNotScanned }}
|
||||
{{/* The image was catalogued but never matched against a vulnerability
|
||||
database, so its zero counts mean "unmeasured", not "clean". Ghost, like
|
||||
"Not scanned": this is an absence of data, not a good result. */}}
|
||||
<span class="badge badge-sm badge-ghost" title="An SBOM was generated, but this image was not checked for vulnerabilities">{{ icon "file-text" "size-3" }} SBOM only</span>
|
||||
{{ else if eq .Total 0 }}
|
||||
<span class="badge badge-sm badge-success" title="No vulnerabilities found (scanned {{ .ScannedAt }})">{{ icon "shield-check" "size-3" }} Clean</span>
|
||||
{{ else }}
|
||||
|
||||
@@ -4,6 +4,12 @@
|
||||
<p class="font-medium text-base-content">No vulnerability scan available yet</p>
|
||||
<p class="mt-1">Scans run automatically shortly after a push. Check back in a few minutes, or push a new tag to trigger a scan.</p>
|
||||
</div>
|
||||
{{ else if .VulnsNotScanned }}
|
||||
<div class="py-8 text-sm text-base-content/70 max-w-prose">
|
||||
<p class="font-medium text-base-content">Vulnerability scanning did not run for this image</p>
|
||||
<p class="mt-1">The scanner catalogued the image contents, so an SBOM is available, but it was never matched against a vulnerability database. No result here means no data, not a clean bill of health.</p>
|
||||
{{ if .ScannedAt }}<p class="mt-2 text-xs text-base-content/60">Scanned: {{ .ScannedAt }}</p>{{ end }}
|
||||
</div>
|
||||
{{ else if .Error }}
|
||||
{{ if gt .Summary.Total 0 }}
|
||||
<!-- Summary available but no detailed report -->
|
||||
|
||||
+923
-104
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,962 @@
|
||||
package pds
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"atcr.io/pkg/atproto"
|
||||
)
|
||||
|
||||
// These tests describe a hold that can keep more than one scan running at a
|
||||
// time — across the workers of a single scanner process (vertical) and across
|
||||
// several scanner processes (horizontal).
|
||||
//
|
||||
// Both were defeated by the same thing: dispatchLoop called waitForCapacity()
|
||||
// before every dispatch, that gate blocked until no row anywhere was
|
||||
// 'assigned' or 'processing', and then exactly one candidate was dispatched.
|
||||
// Proactive scanning was therefore depth one hold-wide, so a second worker
|
||||
// could never be given proactive work and neither could a second scanner
|
||||
// process, even though dispatchJob already round-robins across subscribers.
|
||||
//
|
||||
// Two separate defects live in that one gate. It counted every row rather than
|
||||
// only the proactive ones, so a push-triggered scan starved proactive dispatch
|
||||
// and vice versa; and the depth was hardcoded at one rather than derived from
|
||||
// how much scanning capacity is actually connected.
|
||||
|
||||
// newConcurrencyBroadcaster is newRecordingScanBroadcaster plus the channels
|
||||
// the dispatch gate selects on. The bare helpers leave stopCh nil, which makes
|
||||
// every select in waitForProactiveCapacity block forever.
|
||||
func newConcurrencyBroadcaster(t *testing.T) *ScanBroadcaster {
|
||||
t.Helper()
|
||||
|
||||
sb := newRecordingScanBroadcaster(t)
|
||||
sb.stopCh = make(chan struct{})
|
||||
sb.completionSignal = make(chan struct{}, 1)
|
||||
t.Cleanup(func() {
|
||||
select {
|
||||
case <-sb.stopCh:
|
||||
default:
|
||||
close(sb.stopCh)
|
||||
}
|
||||
})
|
||||
return sb
|
||||
}
|
||||
|
||||
// seedJob inserts one pending job with an explicit origin, which is what the
|
||||
// proactive capacity gate keys on.
|
||||
func seedJob(t *testing.T, sb *ScanBroadcaster, digest, origin string) int64 {
|
||||
t.Helper()
|
||||
|
||||
res, 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', ?)
|
||||
`, digest, "repo", "latest", "did:plc:user", "user.example.com",
|
||||
sb.holdDID, sb.holdEndpoint, "deckhand", "{}", "[]", origin)
|
||||
if err != nil {
|
||||
t.Fatalf("seed %s job: %v", origin, err)
|
||||
}
|
||||
seq, err := res.LastInsertId()
|
||||
if err != nil {
|
||||
t.Fatalf("seq: %v", err)
|
||||
}
|
||||
return seq
|
||||
}
|
||||
|
||||
func setStatus(t *testing.T, sb *ScanBroadcaster, seq int64, status string) {
|
||||
t.Helper()
|
||||
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET status = ? WHERE seq = ?`, status, seq); err != nil {
|
||||
t.Fatalf("set status %q on %d: %v", status, seq, err)
|
||||
}
|
||||
}
|
||||
|
||||
func assignedTo(t *testing.T, sb *ScanBroadcaster, seq int64) string {
|
||||
t.Helper()
|
||||
|
||||
var to sql.NullString
|
||||
if err := sb.db.QueryRow(`SELECT assigned_to FROM scan_jobs WHERE seq = ?`, seq).Scan(&to); err != nil {
|
||||
t.Fatalf("query assigned_to for %d: %v", seq, err)
|
||||
}
|
||||
return to.String
|
||||
}
|
||||
|
||||
// markStarted sets started_at the given number of minutes in the past, as a
|
||||
// 'started' message from a worker would have done at that time.
|
||||
func markStarted(t *testing.T, sb *ScanBroadcaster, seq int64, agoMinutes int) {
|
||||
t.Helper()
|
||||
|
||||
at := time.Now().Add(-time.Duration(agoMinutes) * time.Minute)
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET started_at = ? WHERE seq = ?`, at, seq); err != nil {
|
||||
t.Fatalf("mark started %d: %v", seq, err)
|
||||
}
|
||||
}
|
||||
|
||||
// jobIsDisconnected reports whether the row is marked as belonging to a
|
||||
// scanner that has dropped its connection but may still be running it.
|
||||
func jobIsDisconnected(t *testing.T, sb *ScanBroadcaster, seq int64) bool {
|
||||
t.Helper()
|
||||
|
||||
var at sql.NullTime
|
||||
if err := sb.db.QueryRow(`SELECT disconnected_at FROM scan_jobs WHERE seq = ?`, seq).Scan(&at); err != nil {
|
||||
t.Fatalf("query disconnected_at for %d: %v", seq, err)
|
||||
}
|
||||
return at.Valid
|
||||
}
|
||||
|
||||
func startedAt(t *testing.T, sb *ScanBroadcaster, seq int64) sql.NullTime {
|
||||
t.Helper()
|
||||
|
||||
var at sql.NullTime
|
||||
if err := sb.db.QueryRow(`SELECT started_at FROM scan_jobs WHERE seq = ?`, seq).Scan(&at); err != nil {
|
||||
t.Fatalf("query started_at for %d: %v", seq, err)
|
||||
}
|
||||
return at
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Dispatch depth
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// TestScanProactiveDispatchLimit_TracksConnectedScannerWorkers pins what the
|
||||
// proactive dispatch depth is derived from: the scanning capacity actually
|
||||
// connected, summed over subscribers.
|
||||
//
|
||||
// A scanner declares its worker count when it subscribes. One that declares
|
||||
// nothing — every scanner built before the parameter existed — counts as one
|
||||
// worker, which is exactly the depth the hold had before, so an old scanner
|
||||
// against a new hold behaves as it always did.
|
||||
func TestScanProactiveDispatchLimit_TracksConnectedScannerWorkers(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
|
||||
if got := sb.proactiveDispatchLimit(); got != 0 {
|
||||
t.Errorf("limit with no scanner connected = %d, want 0: there is "+
|
||||
"nothing to dispatch to", got)
|
||||
}
|
||||
|
||||
old := newTestScanSubscriber(t, sb, 4)
|
||||
old.id = "legacy"
|
||||
if got := sb.proactiveDispatchLimit(); got != 1 {
|
||||
t.Errorf("limit for a scanner that declares no worker count = %d, want 1", got)
|
||||
}
|
||||
|
||||
old.capacity = 2 // the same process reconnecting with workers: 2
|
||||
if got := sb.proactiveDispatchLimit(); got != 2 {
|
||||
t.Errorf("limit for one two-worker scanner = %d, want 2: workers above "+
|
||||
"one must buy proactive throughput", got)
|
||||
}
|
||||
|
||||
second := newTestScanSubscriber(t, sb, 4)
|
||||
second.id = "second-process"
|
||||
second.capacity = 3
|
||||
if got := sb.proactiveDispatchLimit(); got != 5 {
|
||||
t.Errorf("limit across two scanner processes = %d, want 5", got)
|
||||
}
|
||||
|
||||
sb.Unsubscribe(second)
|
||||
if got := sb.proactiveDispatchLimit(); got != 2 {
|
||||
t.Errorf("limit after a scanner disconnected = %d, want 2", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProactiveCapacity_AllowsOneJobPerWorker is the vertical scaling case:
|
||||
// `scanner.workers: 2` must genuinely produce two concurrent proactive scans.
|
||||
//
|
||||
// Depth is one proactive job per connected worker. Not more: the scanner acks
|
||||
// on receipt and queues internally, so anything beyond one per worker is a
|
||||
// backlog the hold cannot see into, which is what made the ten-minute deadline
|
||||
// fire under a healthy scanner. Not fewer: one worker would sit idle.
|
||||
func TestScanProactiveCapacity_AllowsOneJobPerWorker(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
sub.capacity = 2
|
||||
|
||||
first := seedJob(t, sb, "sha256:first", originProactive)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: first, Repository: "repo"})
|
||||
|
||||
if n, ok := sb.activeProactiveJobs(); !ok || n != 1 {
|
||||
t.Fatalf("active proactive jobs = %d (ok=%v), want 1", n, ok)
|
||||
}
|
||||
if !sb.hasProactiveCapacity() {
|
||||
t.Fatal("a two-worker scanner with one job running has capacity for a " +
|
||||
"second; depth-one dispatch is what defeats scanner.workers")
|
||||
}
|
||||
|
||||
second := seedJob(t, sb, "sha256:second", originProactive)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: second, Repository: "repo"})
|
||||
|
||||
if n, _ := sb.activeProactiveJobs(); n != 2 {
|
||||
t.Fatalf("active proactive jobs = %d, want 2", n)
|
||||
}
|
||||
if sb.hasProactiveCapacity() {
|
||||
t.Error("both workers are busy; the hold must stop dispatching rather " +
|
||||
"than pile a backlog into the scanner's own queue")
|
||||
}
|
||||
|
||||
// Both jobs really are out with the scanner, not merely counted.
|
||||
if len(sub.send) != 2 {
|
||||
t.Errorf("scanner received %d jobs, want 2", len(sub.send))
|
||||
}
|
||||
if got := assignedTo(t, sb, first); got != sub.id {
|
||||
t.Errorf("job %d assigned_to = %q, want %q", first, got, sub.id)
|
||||
}
|
||||
if got := assignedTo(t, sb, second); got != sub.id {
|
||||
t.Errorf("job %d assigned_to = %q, want %q", second, got, sub.id)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProactiveCapacity_ScalesAcrossScannerProcesses is the horizontal
|
||||
// case. Two single-worker scanners are two units of capacity, and dispatchJob
|
||||
// already round-robins, so the two jobs must land on different scanners.
|
||||
func TestScanProactiveCapacity_ScalesAcrossScannerProcesses(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
a := newTestScanSubscriber(t, sb, 4)
|
||||
a.id = "scanner-a"
|
||||
b := newTestScanSubscriber(t, sb, 4)
|
||||
b.id = "scanner-b"
|
||||
|
||||
first := seedJob(t, sb, "sha256:first", originProactive)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: first, Repository: "repo"})
|
||||
|
||||
if !sb.hasProactiveCapacity() {
|
||||
t.Fatal("a second scanner process is a second unit of capacity and " +
|
||||
"could never receive proactive work under a depth-one gate")
|
||||
}
|
||||
|
||||
second := seedJob(t, sb, "sha256:second", originProactive)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: second, Repository: "repo"})
|
||||
|
||||
if len(a.send) != 1 || len(b.send) != 1 {
|
||||
t.Errorf("jobs per scanner = a:%d b:%d, want 1 each", len(a.send), len(b.send))
|
||||
}
|
||||
if sb.hasProactiveCapacity() {
|
||||
t.Error("both scanners are busy; dispatch must stop")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProactiveCapacity_IgnoresPushTriggeredJobs covers the second defect
|
||||
// in the gate. Push-triggered scans bypass it entirely — oci/xrpc.go calls
|
||||
// Enqueue directly — but they were counted by it, so on a hold with steady
|
||||
// pushes "one proactive job at a time" was really "none while anyone pushes".
|
||||
func TestScanProactiveCapacity_IgnoresPushTriggeredJobs(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
newTestScanSubscriber(t, sb, 4) // one worker: depth one
|
||||
|
||||
push := seedJob(t, sb, "sha256:pushed", originPush)
|
||||
setStatus(t, sb, push, "processing")
|
||||
|
||||
if n, ok := sb.activeProactiveJobs(); !ok || n != 0 {
|
||||
t.Fatalf("active proactive jobs = %d (ok=%v), want 0: a push-triggered "+
|
||||
"scan is not proactive work", n, ok)
|
||||
}
|
||||
if !sb.hasProactiveCapacity() {
|
||||
t.Error("a push-triggered scan is throttling proactive dispatch; the " +
|
||||
"two paths must not consume each other's budget")
|
||||
}
|
||||
|
||||
// And the mirror: a proactive job in flight fills the proactive budget.
|
||||
proactive := seedJob(t, sb, "sha256:proactive", originProactive)
|
||||
setStatus(t, sb, proactive, "assigned")
|
||||
|
||||
if n, _ := sb.activeProactiveJobs(); n != 1 {
|
||||
t.Errorf("active proactive jobs = %d, want 1", n)
|
||||
}
|
||||
if sb.hasProactiveCapacity() {
|
||||
t.Error("the one worker is busy with a proactive job already")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanWaitForProactiveCapacity_BlocksAtTheLimitAndReleasesOnCompletion
|
||||
// checks the gate the dispatch loop actually calls: it must park while the
|
||||
// budget is full, and wake on the completion signal rather than on a timer.
|
||||
func TestScanWaitForProactiveCapacity_BlocksAtTheLimitAndReleasesOnCompletion(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
newTestScanSubscriber(t, sb, 4) // one worker
|
||||
|
||||
seq := seedJob(t, sb, "sha256:running", originProactive)
|
||||
setStatus(t, sb, seq, "processing")
|
||||
|
||||
got := make(chan bool, 1)
|
||||
go func() { got <- sb.waitForProactiveCapacity() }()
|
||||
|
||||
select {
|
||||
case <-got:
|
||||
t.Fatal("the gate returned while the only worker was busy")
|
||||
case <-time.After(250 * time.Millisecond):
|
||||
}
|
||||
|
||||
setStatus(t, sb, seq, "completed")
|
||||
sb.signalCompletion()
|
||||
|
||||
select {
|
||||
case ok := <-got:
|
||||
if !ok {
|
||||
t.Error("the gate reported no capacity after the job completed")
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("the gate did not wake on the completion signal")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanWaitForProactiveCapacity_YieldsWhenNoScannerIsConnected keeps the
|
||||
// gate safe at zero capacity. Blocking inside it forever would be correct-ish
|
||||
// today but leaves the dispatch loop unable to notice a scanner arriving, and
|
||||
// unable to re-check anything else; it returns instead so the loop can wait on
|
||||
// the connection.
|
||||
func TestScanWaitForProactiveCapacity_YieldsWhenNoScannerIsConnected(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
|
||||
done := make(chan bool, 1)
|
||||
go func() { done <- sb.waitForProactiveCapacity() }()
|
||||
|
||||
select {
|
||||
case ok := <-done:
|
||||
if ok {
|
||||
t.Error("the gate granted capacity with no scanner connected")
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("the gate blocked with no scanner connected; the dispatch loop " +
|
||||
"cannot re-check the subscriber list from in there")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanWaitForProactiveCapacity_YieldsWhenTheLastScannerDisconnects is the
|
||||
// mid-flight version: capacity that vanishes while the gate is parked must
|
||||
// release it, not strand the loop.
|
||||
func TestScanWaitForProactiveCapacity_YieldsWhenTheLastScannerDisconnects(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
seq := seedJob(t, sb, "sha256:running", originProactive)
|
||||
setStatus(t, sb, seq, "processing")
|
||||
|
||||
got := make(chan bool, 1)
|
||||
go func() { got <- sb.waitForProactiveCapacity() }()
|
||||
|
||||
select {
|
||||
case <-got:
|
||||
t.Fatal("the gate returned while the only worker was busy")
|
||||
case <-time.After(250 * time.Millisecond):
|
||||
}
|
||||
|
||||
sb.Unsubscribe(sub)
|
||||
|
||||
select {
|
||||
case ok := <-got:
|
||||
if ok {
|
||||
t.Error("the gate granted capacity after the last scanner left")
|
||||
}
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("the gate stayed parked after the last scanner disconnected")
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// F4: the deadline must measure scanning, not queueing
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// TestScanStarted_StartsTheScanningClock covers the new signal.
|
||||
//
|
||||
// The ack means "I have it": the scanner sends it off the WebSocket reader the
|
||||
// moment a frame arrives, before the job is even queued. The job then waits in
|
||||
// the scanner's own 100-deep queue behind its workers. Measuring the scanning
|
||||
// deadline from dispatch therefore budgets queueing, and a 100-deep queue of
|
||||
// no-op jobs crosses ten minutes at position 59.
|
||||
//
|
||||
// 'started' is sent by the worker that dequeues the job, so it marks the one
|
||||
// moment the hold could not otherwise observe.
|
||||
func TestScanStarted_StartsTheScanningClock(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
seq := seedJob(t, sb, "sha256:queued", originProactive)
|
||||
assignJob(t, sb, seq, sub, 0)
|
||||
sb.handleAck(sub, seq)
|
||||
|
||||
if at := startedAt(t, sb, seq); at.Valid {
|
||||
t.Fatal("the ack started the scanning clock; it means \"I have it\", " +
|
||||
"not \"a worker is on it\"")
|
||||
}
|
||||
|
||||
sb.handleStarted(sub, seq)
|
||||
|
||||
at := startedAt(t, sb, seq)
|
||||
if !at.Valid {
|
||||
t.Fatal("'started' did not stamp started_at, so the deadline still " +
|
||||
"measures queueing rather than scanning")
|
||||
}
|
||||
if d := time.Since(at.Time); d > time.Minute || d < -time.Minute {
|
||||
t.Errorf("started_at is %s away from now, want ~0", d)
|
||||
}
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Errorf("status = %q, want processing", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanStarted_KeepsAQueuedJobInsideItsScanningDeadline is the behaviour
|
||||
// the whole change is for: a job that sat in a scanner's queue for longer than
|
||||
// the scanning deadline, and has only just started, must not be cancelled out
|
||||
// from under the worker now running it.
|
||||
func TestScanStarted_KeepsAQueuedJobInsideItsScanningDeadline(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:longqueue"
|
||||
seq := seedJob(t, sb, digest, originProactive)
|
||||
assignJob(t, sb, seq, sub, 20) // dispatched twenty minutes ago
|
||||
sb.handleAck(sub, seq)
|
||||
sb.addInflight(digest)
|
||||
|
||||
sb.handleStarted(sub, seq) // a worker picked it up just now
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("status = %q, want processing: a scan that started seconds "+
|
||||
"ago was cancelled because it had queued for twenty minutes", got)
|
||||
}
|
||||
if _, _, err := sb.pds.GetScanRecord(context.Background(), digest); err == nil {
|
||||
t.Error("a failure record was written for a scan that is actively running")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProcessingTimeout_FailsAJobStuckSinceItStarted keeps the deadline
|
||||
// real. Once a worker has said it started, the scanning budget applies from
|
||||
// that moment and a wedged scan is still retired.
|
||||
func TestScanProcessingTimeout_FailsAJobStuckSinceItStarted(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:wedgedscan"
|
||||
seq := seedJob(t, sb, digest, originProactive)
|
||||
assignJob(t, sb, seq, sub, 30)
|
||||
sb.handleAck(sub, seq)
|
||||
sb.addInflight(digest)
|
||||
markStarted(t, sb, seq, 11) // started eleven minutes ago, still nothing
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "failed" {
|
||||
t.Fatalf("status = %q, want failed", got)
|
||||
}
|
||||
_, record, err := sb.pds.GetScanRecord(context.Background(), digest)
|
||||
if err != nil {
|
||||
t.Fatalf("no scan record for a job the hold gave up on: %v", err)
|
||||
}
|
||||
if record.Status != atproto.ScanStatusFailed {
|
||||
t.Errorf("record status = %q, want %q", record.Status, atproto.ScanStatusFailed)
|
||||
}
|
||||
if !sb.addInflight(digest) {
|
||||
t.Error("the timed-out digest is still in flight")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProcessingTimeout_ToleratesAScannerThatNeverReportsStarts is the
|
||||
// compatibility half. A scanner built before 'started' existed never sends it,
|
||||
// so started_at stays NULL and the hold can only observe dispatch. The budget
|
||||
// for that case is the one the hold can actually justify: long enough for a
|
||||
// full scanner queue to drain, so a healthy backlogged scanner is not killed,
|
||||
// and still bounded so a wedged-but-connected scanner does not hold capacity
|
||||
// for the life of the process.
|
||||
func TestScanProcessingTimeout_ToleratesAScannerThatNeverReportsStarts(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:oldscanner"
|
||||
seq := seedJob(t, sb, digest, originProactive)
|
||||
assignJob(t, sb, seq, sub, 15) // well past the ten-minute scanning deadline
|
||||
sb.handleAck(sub, seq)
|
||||
sb.addInflight(digest)
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("status = %q, want processing: an acked job queued fifteen "+
|
||||
"minutes behind a busy scanner is not evidence of anything wrong", got)
|
||||
}
|
||||
|
||||
// Past the queue budget it is retired like any other stuck job.
|
||||
if _, err := sb.db.Exec(
|
||||
`UPDATE scan_jobs SET assigned_at = ? WHERE seq = ?`,
|
||||
time.Now().Add(-queuedTimeout-time.Minute), seq,
|
||||
); err != nil {
|
||||
t.Fatalf("age the row: %v", err)
|
||||
}
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "failed" {
|
||||
t.Errorf("status = %q, want failed: an unstarted job must still have a "+
|
||||
"bound, or a connected-but-wedged scanner holds capacity forever", got)
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Ownership of terminal transitions
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// TestScanTerminalHandlers_IgnoreAnotherScannersJob closes the hole that opens
|
||||
// the moment horizontal scaling works. handleAck guards on assigned_to; the
|
||||
// handlers that write records and retire the row did not, so any scanner could
|
||||
// complete, fail or skip a job belonging to another one — writing a scan
|
||||
// record for an image it never looked at and releasing a digest a different
|
||||
// scanner is still working on.
|
||||
func TestScanTerminalHandlers_IgnoreAnotherScannersJob(t *testing.T) {
|
||||
cases := []struct {
|
||||
name string
|
||||
send func(sb *ScanBroadcaster, sub *ScanSubscriber, seq int64)
|
||||
}{
|
||||
{"result", func(sb *ScanBroadcaster, sub *ScanSubscriber, seq int64) {
|
||||
sb.handleResult(sub, ScannerMessage{Type: "result", Seq: seq, SBOM: testSBOM})
|
||||
}},
|
||||
{"error", func(sb *ScanBroadcaster, sub *ScanSubscriber, seq int64) {
|
||||
sb.handleError(sub, ScannerMessage{Type: "error", Seq: seq, Error: "boom"})
|
||||
}},
|
||||
{"skipped", func(sb *ScanBroadcaster, sub *ScanSubscriber, seq int64) {
|
||||
sb.handleSkipped(sub, ScannerMessage{Type: "skipped", Seq: seq, Reason: "nope"})
|
||||
}},
|
||||
{"started", func(sb *ScanBroadcaster, sub *ScanSubscriber, seq int64) {
|
||||
sb.handleStarted(sub, seq)
|
||||
}},
|
||||
}
|
||||
|
||||
for _, tc := range cases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
owner := newTestScanSubscriber(t, sb, 4)
|
||||
owner.id = "owner"
|
||||
intruder := newTestScanSubscriber(t, sb, 4)
|
||||
intruder.id = "intruder"
|
||||
|
||||
digest := "sha256:owned" + tc.name
|
||||
seq := seedJob(t, sb, digest, originProactive)
|
||||
assignJob(t, sb, seq, owner, 0)
|
||||
sb.handleAck(owner, seq)
|
||||
sb.addInflight(digest)
|
||||
|
||||
tc.send(sb, intruder, seq)
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Errorf("status = %q, want processing: a scanner that does not "+
|
||||
"own the job retired it", got)
|
||||
}
|
||||
if sb.addInflight(digest) {
|
||||
t.Error("another scanner's message released the in-flight digest")
|
||||
}
|
||||
if _, _, err := sb.pds.GetScanRecord(context.Background(), digest); err == nil {
|
||||
t.Error("a scan record was written by a scanner that never had the job")
|
||||
}
|
||||
if tc.name == "started" {
|
||||
if startedAt(t, sb, seq).Valid {
|
||||
t.Error("another scanner restarted the scanning clock")
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanTerminalHandlers_AcceptTheOwningScanner is the other half: the guard
|
||||
// must not reject the scanner that legitimately holds the job.
|
||||
func TestScanTerminalHandlers_AcceptTheOwningScanner(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
owner := newTestScanSubscriber(t, sb, 4)
|
||||
owner.id = "owner"
|
||||
|
||||
const digest = "sha256:legitimate"
|
||||
seq := seedJob(t, sb, digest, originProactive)
|
||||
assignJob(t, sb, seq, owner, 0)
|
||||
sb.handleAck(owner, seq)
|
||||
sb.handleStarted(owner, seq)
|
||||
sb.addInflight(digest)
|
||||
|
||||
sb.handleResult(owner, ScannerMessage{Type: "result", Seq: seq, SBOM: testSBOM})
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "completed" {
|
||||
t.Fatalf("status = %q, want completed", got)
|
||||
}
|
||||
if !sb.addInflight(digest) {
|
||||
t.Error("the owning scanner's result did not release the digest")
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// The buffer-full reset
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// TestScanDispatchJob_BufferFullResetLeavesAnotherScannersClaimAlone is the
|
||||
// mirror of the double-dispatch window closed in drainPendingJobs.
|
||||
//
|
||||
// dispatchJob assigns the row, finds the scanner's send buffer full, and puts
|
||||
// the row back to 'pending' — with no assigned_to or status guard. If another
|
||||
// dispatcher claimed the row in between, that reset hands a job a second
|
||||
// scanner is already holding back to the pool, and a third scanner gets it.
|
||||
//
|
||||
// The interleaving is built with a SQLite trigger rather than a race: the
|
||||
// trigger fires on dispatchJob's own assign UPDATE and reassigns the row to
|
||||
// another scanner, so by the time the buffer-full branch runs, the row is
|
||||
// provably not ours. SQLite does not recurse triggers by default, so the reset
|
||||
// UPDATE does not re-fire it.
|
||||
func TestScanDispatchJob_BufferFullResetLeavesAnotherScannersClaimAlone(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 0) // unbuffered, nobody receiving: always full
|
||||
sub.id = "loser"
|
||||
|
||||
seq := seedJob(t, sb, "sha256:contended", originProactive)
|
||||
|
||||
if _, err := sb.db.Exec(`
|
||||
CREATE TRIGGER steal_assignment AFTER UPDATE OF status ON scan_jobs
|
||||
WHEN NEW.status = 'assigned' AND NEW.assigned_to = 'loser'
|
||||
BEGIN
|
||||
UPDATE scan_jobs SET assigned_to = 'winner' WHERE seq = NEW.seq;
|
||||
END
|
||||
`); err != nil {
|
||||
t.Fatalf("create trigger: %v", err)
|
||||
}
|
||||
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: seq, Repository: "repo"})
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "assigned" {
|
||||
t.Errorf("status = %q, want assigned: the buffer-full reset returned a "+
|
||||
"row another dispatcher holds to the pool, so a third scanner will "+
|
||||
"be handed a job that is already out", got)
|
||||
}
|
||||
if got := assignedTo(t, sb, seq); got != "winner" {
|
||||
t.Errorf("assigned_to = %q, want \"winner\": the claim was stolen back", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanDispatchJob_BufferFullReturnsOurOwnRowToPending is the guard's other
|
||||
// side: when the row really is ours, a full buffer must still release it so
|
||||
// the re-dispatch loop can offer it again rather than leaving it assigned to a
|
||||
// scanner that was never sent it.
|
||||
func TestScanDispatchJob_BufferFullReturnsOurOwnRowToPending(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 0) // unbuffered, nobody receiving
|
||||
sub.id = "solo"
|
||||
|
||||
seq := seedJob(t, sb, "sha256:nobuffer", originProactive)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: seq, Repository: "repo"})
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "pending" {
|
||||
t.Errorf("status = %q, want pending", got)
|
||||
}
|
||||
if got := assignedTo(t, sb, seq); got != "" {
|
||||
t.Errorf("assigned_to = %q, want empty", got)
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Fairness and admission control across N scanner processes
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// TestScanDispatch_PrefersTheLeastLoadedScanner covers plain round-robin's
|
||||
// blind spot. nextIdx hands the next job to whichever subscriber is next in
|
||||
// the slice, saturated or idle, which with heterogeneous processes (different
|
||||
// worker counts, different hosts, one mid-scan and one just connected) piles
|
||||
// work onto a scanner that cannot take it while another sits idle.
|
||||
func TestScanDispatch_PrefersTheLeastLoadedScanner(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
busy := newTestScanSubscriber(t, sb, 4)
|
||||
busy.id = "busy"
|
||||
busy.capacity = 2
|
||||
idle := newTestScanSubscriber(t, sb, 4)
|
||||
idle.id = "idle"
|
||||
idle.capacity = 2
|
||||
|
||||
// "busy" is already running one scan.
|
||||
running := seedJob(t, sb, "sha256:running", originPush)
|
||||
if _, err := sb.db.Exec(
|
||||
`UPDATE scan_jobs SET status='processing', assigned_to=?, assigned_at=? WHERE seq=?`,
|
||||
busy.id, time.Now(), running,
|
||||
); err != nil {
|
||||
t.Fatalf("prime load: %v", err)
|
||||
}
|
||||
|
||||
seq := seedJob(t, sb, "sha256:next", originProactive)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: seq, Repository: "repo"})
|
||||
|
||||
if got := assignedTo(t, sb, seq); got != idle.id {
|
||||
t.Errorf("job went to %q, want %q: dispatch must weigh what each "+
|
||||
"scanner is already holding", got, idle.id)
|
||||
}
|
||||
if len(idle.send) != 1 || len(busy.send) != 0 {
|
||||
t.Errorf("sends: idle=%d busy=%d, want 1/0", len(idle.send), len(busy.send))
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanDispatch_LeavesTheRowPendingWhenEveryScannerIsSaturated is the
|
||||
// admission control the deadline depends on.
|
||||
//
|
||||
// The hold cannot see into a scanner's own queue: the scanner acks on receipt
|
||||
// and the job then waits behind its workers, which is exactly why a deadline
|
||||
// measured from dispatch cancels healthy work. Keeping the queue on the hold
|
||||
// instead of pushing it into the scanner keeps the row observable, keeps it
|
||||
// available to whichever process frees up first, and bounds what any one
|
||||
// process is holding.
|
||||
func TestScanDispatch_LeavesTheRowPendingWhenEveryScannerIsSaturated(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 20)
|
||||
sub.capacity = 1
|
||||
|
||||
first := seedJob(t, sb, "sha256:one", originPush)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: first, Repository: "repo"})
|
||||
if got := jobStatus(t, sb, first); got != "assigned" {
|
||||
t.Fatalf("first job status = %q, want assigned", got)
|
||||
}
|
||||
|
||||
second := seedJob(t, sb, "sha256:two", originPush)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: second, Repository: "repo"})
|
||||
|
||||
if got := jobStatus(t, sb, second); got != "pending" {
|
||||
t.Errorf("second job status = %q, want pending: the only worker is "+
|
||||
"busy, so the row belongs on the hold's queue, not the scanner's", got)
|
||||
}
|
||||
if len(sub.send) != 1 {
|
||||
t.Errorf("scanner received %d jobs, want 1", len(sub.send))
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanDispatch_OffersPendingWorkAsSoonAsCapacityFrees is the other side of
|
||||
// admission control: a row held back must not wait for the thirty-second
|
||||
// re-dispatch tick, or holding it back would cost more throughput than the
|
||||
// scanner-side queue ever did.
|
||||
func TestScanDispatch_OffersPendingWorkAsSoonAsCapacityFrees(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 20)
|
||||
sub.capacity = 1
|
||||
|
||||
first := seedJob(t, sb, "sha256:one", originPush)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: first, Repository: "repo"})
|
||||
second := seedJob(t, sb, "sha256:two", originPush)
|
||||
sb.dispatchJob(&ScanJobEvent{Type: "job", Seq: second, Repository: "repo"})
|
||||
<-sub.send // drop the first job so the assertion below is unambiguous
|
||||
|
||||
setStatus(t, sb, first, "completed")
|
||||
sb.offerPendingJobs()
|
||||
|
||||
if got := jobStatus(t, sb, second); got != "assigned" {
|
||||
t.Errorf("held job status = %q, want assigned once the worker freed up", got)
|
||||
}
|
||||
select {
|
||||
case job := <-sub.send:
|
||||
if job.Seq != second {
|
||||
t.Errorf("scanner received seq %d, want %d", job.Seq, second)
|
||||
}
|
||||
default:
|
||||
t.Error("the freed worker was given nothing")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanDrain_StopsAtTheSubscriberCapacity stops the first scanner to
|
||||
// connect from swallowing a whole backlog. drainPendingJobs walked every
|
||||
// pending row and pushed it at the new subscriber, so with a backlog and two
|
||||
// scanner processes the first one took all of it and the second stayed idle —
|
||||
// horizontal scaling defeated at the connect path rather than at the gate.
|
||||
func TestScanDrain_StopsAtTheSubscriberCapacity(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 20)
|
||||
sub.capacity = 2
|
||||
|
||||
var seqs []int64
|
||||
for i := 0; i < 5; i++ {
|
||||
seqs = append(seqs, seedJob(t, sb, "sha256:backlog", originProactive))
|
||||
}
|
||||
|
||||
sb.drainPendingJobs(sub, 0)
|
||||
|
||||
var assigned, pending int
|
||||
for _, seq := range seqs {
|
||||
switch got := jobStatus(t, sb, seq); got {
|
||||
case "assigned":
|
||||
assigned++
|
||||
case "pending":
|
||||
pending++
|
||||
default:
|
||||
t.Errorf("job %d in unexpected status %q", seq, got)
|
||||
}
|
||||
}
|
||||
if assigned != 2 || pending != 3 {
|
||||
t.Errorf("assigned=%d pending=%d, want 2/3: a two-worker scanner takes "+
|
||||
"two jobs, and the rest stay available to other processes",
|
||||
assigned, pending)
|
||||
}
|
||||
if len(sub.send) != 2 {
|
||||
t.Errorf("scanner received %d jobs, want 2", len(sub.send))
|
||||
}
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// Disconnect, reconnect, and who owns the work in between
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
// TestScanUnsubscribe_DoesNotImmediatelyHandOffWorkStillRunning is the N-process
|
||||
// version of the duplicate-scan problem.
|
||||
//
|
||||
// Unsubscribe flipped every assigned and processing row belonging to the
|
||||
// dropped subscriber straight back to 'pending'. The scanner's worker pool
|
||||
// never learns the socket dropped, so it keeps scanning. With one scanner that
|
||||
// produced a duplicate against itself. With several it is worse and cannot be
|
||||
// deduped anywhere: scanner A blips, B claims A's rows out of drainPendingJobs,
|
||||
// and both processes scan the same images and both report a verdict.
|
||||
//
|
||||
// A disconnect is not evidence that a scanner is gone. It is marked as
|
||||
// disconnected and its work is left alone for a grace period instead.
|
||||
func TestScanUnsubscribe_DoesNotImmediatelyHandOffWorkStillRunning(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
first := newTestScanSubscriber(t, sb, 4)
|
||||
first.id = "scanner-a"
|
||||
|
||||
seq := seedJob(t, sb, "sha256:midscan", originProactive)
|
||||
assignJob(t, sb, seq, first, 0)
|
||||
sb.handleAck(first, seq)
|
||||
sb.handleStarted(first, seq)
|
||||
|
||||
sb.Unsubscribe(first)
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("status = %q, want processing: a five-second blip must not "+
|
||||
"put a running scan back in the pool", got)
|
||||
}
|
||||
if got := assignedTo(t, sb, seq); got != first.id {
|
||||
t.Errorf("assigned_to = %q, want %q", got, first.id)
|
||||
}
|
||||
|
||||
// A second scanner process connecting must not be handed it.
|
||||
second := newTestScanSubscriber(t, sb, 4)
|
||||
second.id = "scanner-b"
|
||||
sb.drainPendingJobs(second, 0)
|
||||
|
||||
select {
|
||||
case job := <-second.send:
|
||||
t.Fatalf("seq %d was handed to a second scanner process while the first "+
|
||||
"is still scanning it", job.Seq)
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanReconnect_LetsAScannerResumeItsOwnWork is why the disconnect is only
|
||||
// marked rather than acted on. A scanner keeps one identity for the life of the
|
||||
// process and sends it on every connect, so a reconnection inside the grace
|
||||
// window reclaims the jobs its workers never stopped running.
|
||||
//
|
||||
// A scanner that actually restarted comes back with a new identity, so its old
|
||||
// rows are not resumed and are reclaimed by the grace timeout instead — which
|
||||
// is right, because a restarted process really did lose that work.
|
||||
func TestScanReconnect_LetsAScannerResumeItsOwnWork(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
sub.id = "instance-1"
|
||||
|
||||
seq := seedJob(t, sb, "sha256:resumed", originProactive)
|
||||
assignJob(t, sb, seq, sub, 0)
|
||||
sb.handleAck(sub, seq)
|
||||
sb.handleStarted(sub, seq)
|
||||
sb.Unsubscribe(sub)
|
||||
|
||||
if !jobIsDisconnected(t, sb, seq) {
|
||||
t.Fatal("the row was not marked as belonging to a disconnected scanner")
|
||||
}
|
||||
|
||||
sb.resumeInstance("instance-1")
|
||||
|
||||
if jobIsDisconnected(t, sb, seq) {
|
||||
t.Error("reconnecting did not clear the disconnect mark, so the job " +
|
||||
"will be reclaimed from a scanner that never stopped running it")
|
||||
}
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Errorf("status = %q, want processing", got)
|
||||
}
|
||||
|
||||
// A different process must not adopt it.
|
||||
sb.resumeInstance("instance-2")
|
||||
if got := assignedTo(t, sb, seq); got != "instance-1" {
|
||||
t.Errorf("assigned_to = %q, want instance-1", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanDisconnect_ReclaimsWorkOnceTheGraceExpires bounds the wait. A scanner
|
||||
// that does not come back has its work returned to the pool, so a genuinely
|
||||
// dead process costs one grace period rather than a permanent hole.
|
||||
func TestScanDisconnect_ReclaimsWorkOnceTheGraceExpires(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
gone := newTestScanSubscriber(t, sb, 4)
|
||||
gone.id = "departed"
|
||||
|
||||
seq := seedJob(t, sb, "sha256:abandoned", originProactive)
|
||||
assignJob(t, sb, seq, gone, 0)
|
||||
sb.handleAck(gone, seq)
|
||||
sb.handleStarted(gone, seq)
|
||||
sb.Unsubscribe(gone)
|
||||
|
||||
// Inside the grace: still theirs.
|
||||
sb.reDispatchTimedOut()
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("status = %q, want processing inside the grace window", got)
|
||||
}
|
||||
|
||||
if _, err := sb.db.Exec(
|
||||
`UPDATE scan_jobs SET disconnected_at = ? WHERE seq = ?`,
|
||||
time.Now().Add(-reconnectGrace-time.Minute), seq,
|
||||
); err != nil {
|
||||
t.Fatalf("age the disconnect: %v", err)
|
||||
}
|
||||
|
||||
replacement := newTestScanSubscriber(t, sb, 4)
|
||||
replacement.id = "replacement"
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := assignedTo(t, sb, seq); got != replacement.id {
|
||||
t.Errorf("assigned_to = %q, want %q: a scanner that never came back "+
|
||||
"must not hold its work forever", got, replacement.id)
|
||||
}
|
||||
if len(replacement.send) != 1 {
|
||||
t.Errorf("replacement received %d jobs, want 1", len(replacement.send))
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanReconcileOnBoot_TreatsInFlightJobsAsDisconnected covers the restart
|
||||
// case, which nothing reconciled before.
|
||||
//
|
||||
// Every connection the previous process held died with it, so a row still
|
||||
// 'assigned' or 'processing' at boot belongs to a scanner the hold has no
|
||||
// connection to — whether or not that scanner is still running it. Without the
|
||||
// mark those rows sat holding dispatch capacity until their own deadlines
|
||||
// fired, which is up to an hour for a job that was queued inside a scanner.
|
||||
// With it they fall under the same grace as a live disconnect: resumed by the
|
||||
// scanner that redials with the same identity, reclaimed otherwise.
|
||||
func TestScanReconcileOnBoot_TreatsInFlightJobsAsDisconnected(t *testing.T) {
|
||||
sb := newConcurrencyBroadcaster(t)
|
||||
|
||||
assigned := seedJob(t, sb, "sha256:assigned", originProactive)
|
||||
setStatus(t, sb, assigned, "assigned")
|
||||
processing := seedJob(t, sb, "sha256:processing", originProactive)
|
||||
setStatus(t, sb, processing, "processing")
|
||||
done := seedJob(t, sb, "sha256:done", originProactive)
|
||||
setStatus(t, sb, done, "completed")
|
||||
|
||||
sb.reconcileOnBoot()
|
||||
|
||||
for _, seq := range []int64{assigned, processing} {
|
||||
if !jobIsDisconnected(t, sb, seq) {
|
||||
t.Errorf("job %d was left holding dispatch capacity across a restart", seq)
|
||||
}
|
||||
}
|
||||
if jobIsDisconnected(t, sb, done) {
|
||||
t.Error("a completed job was marked as in flight")
|
||||
}
|
||||
|
||||
// And the mark is what the grace acts on, so a scanner that comes back
|
||||
// with the same identity still keeps its work.
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET assigned_to = ? WHERE seq = ?`,
|
||||
"instance-1", processing); err != nil {
|
||||
t.Fatalf("set owner: %v", err)
|
||||
}
|
||||
sb.resumeInstance("instance-1")
|
||||
if jobIsDisconnected(t, sb, processing) {
|
||||
t.Error("a scanner that redialed did not get its own job back")
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
package pds
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
@@ -28,62 +29,85 @@ func jobStatus(t *testing.T, sb *ScanBroadcaster, seq int64) string {
|
||||
return status
|
||||
}
|
||||
|
||||
// TestScanHasActiveJobs_IgnoresLongPendingJob is the regression test for the
|
||||
// nine-day deployment-wide scanning outage. A single job sat in 'pending' with
|
||||
// nothing left to dispatch it, hasActiveJobs() counted it forever, and
|
||||
// waitForCapacity() therefore never let the proactive dispatch loop enqueue
|
||||
// another job — so discovery kept finding unscanned images and creating none.
|
||||
func TestScanHasActiveJobs_IgnoresLongPendingJob(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
seedPendingJobs(t, sb, 1)
|
||||
// activeProactive is the count half of activeProactiveJobs, failing the test if
|
||||
// the answer is not usable.
|
||||
func activeProactive(t *testing.T, sb *ScanBroadcaster) int {
|
||||
t.Helper()
|
||||
|
||||
if !sb.hasActiveJobs() {
|
||||
n, ok := sb.activeProactiveJobs()
|
||||
if !ok {
|
||||
t.Fatal("activeProactiveJobs reported an unusable answer")
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
// TestScanActiveProactiveJobs_IgnoresLongPendingJob is the regression test for
|
||||
// the nine-day deployment-wide scanning outage. A single job sat in 'pending'
|
||||
// with nothing left to dispatch it, the capacity check counted it forever, and
|
||||
// the dispatch gate therefore never let the proactive loop enqueue another job
|
||||
// — so discovery kept finding unscanned images and creating none.
|
||||
func TestScanActiveProactiveJobs_IgnoresLongPendingJob(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
seq := seedJob(t, sb, "sha256:pendingforever", originProactive)
|
||||
|
||||
if activeProactive(t, sb) != 1 {
|
||||
t.Fatal("a freshly enqueued pending job must count as active")
|
||||
}
|
||||
|
||||
backdateJob(t, sb, 1, 60)
|
||||
backdateJob(t, sb, seq, 60)
|
||||
|
||||
if sb.hasActiveJobs() {
|
||||
t.Error("a job pending for an hour must not block dispatch capacity")
|
||||
if n := activeProactive(t, sb); n != 0 {
|
||||
t.Errorf("active = %d, want 0: a job pending for an hour must not block "+
|
||||
"dispatch capacity", n)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanHasActiveJobs_CountsAssignedAndProcessing guards the other half:
|
||||
// assigned and processing jobs have their own reclaim timeouts, so they must
|
||||
// still hold capacity no matter how old the row is.
|
||||
func TestScanHasActiveJobs_CountsAssignedAndProcessing(t *testing.T) {
|
||||
// TestScanActiveProactiveJobs_CountsAssignedAndProcessing guards the other
|
||||
// half: assigned and processing jobs have their own reclaim timeouts, so they
|
||||
// must still hold capacity no matter how old the row is.
|
||||
func TestScanActiveProactiveJobs_CountsAssignedAndProcessing(t *testing.T) {
|
||||
for _, status := range []string{"assigned", "processing"} {
|
||||
t.Run(status, func(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
seedPendingJobs(t, sb, 1)
|
||||
backdateJob(t, sb, 1, 60)
|
||||
seq := seedJob(t, sb, "sha256:inflight", originProactive)
|
||||
backdateJob(t, sb, seq, 60)
|
||||
setStatus(t, sb, seq, status)
|
||||
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET status = ? WHERE seq = 1`, status); err != nil {
|
||||
t.Fatalf("set status: %v", err)
|
||||
}
|
||||
|
||||
if !sb.hasActiveJobs() {
|
||||
if activeProactive(t, sb) != 1 {
|
||||
t.Errorf("%s job must hold dispatch capacity", status)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanHasActiveJobs_IgnoresTerminalJobs keeps completed and failed rows out
|
||||
// of the capacity check — the table holds tens of thousands of them.
|
||||
func TestScanHasActiveJobs_IgnoresTerminalJobs(t *testing.T) {
|
||||
// TestScanActiveProactiveJobs_IgnoresTerminalJobs keeps completed and failed
|
||||
// rows out of the capacity check — the table holds tens of thousands of them.
|
||||
func TestScanActiveProactiveJobs_IgnoresTerminalJobs(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
seedPendingJobs(t, sb, 2)
|
||||
setStatus(t, sb, seedJob(t, sb, "sha256:done", originProactive), "completed")
|
||||
setStatus(t, sb, seedJob(t, sb, "sha256:dead", originProactive), "failed")
|
||||
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET status = 'completed' WHERE seq = 1`); err != nil {
|
||||
t.Fatalf("complete: %v", err)
|
||||
}
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET status = 'failed' WHERE seq = 2`); err != nil {
|
||||
t.Fatalf("fail: %v", err)
|
||||
if n := activeProactive(t, sb); n != 0 {
|
||||
t.Errorf("active = %d, want 0: terminal jobs must not hold capacity", n)
|
||||
}
|
||||
}
|
||||
|
||||
if sb.hasActiveJobs() {
|
||||
t.Error("terminal jobs must not hold dispatch capacity")
|
||||
// TestScanActiveProactiveJobs_IgnoresPushTriggeredWork is the second half of
|
||||
// the gate's original defect, alongside its hardcoded depth of one. Push scans
|
||||
// never pass through this gate — a user who just pushed is waiting for the
|
||||
// answer, so oci/xrpc.go enqueues directly — but they were counted by it, so
|
||||
// on a hold with steady pushes the proactive loop had capacity approximately
|
||||
// never.
|
||||
func TestScanActiveProactiveJobs_IgnoresPushTriggeredWork(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
setStatus(t, sb, seedJob(t, sb, "sha256:pushed", originPush), "processing")
|
||||
// Rows written before the column existed read as push for the same reason:
|
||||
// unknown provenance must not throttle proactive dispatch.
|
||||
seedPendingJobs(t, sb, 1)
|
||||
|
||||
if n := activeProactive(t, sb); n != 0 {
|
||||
t.Errorf("active = %d, want 0: push-triggered work must not spend the "+
|
||||
"proactive budget", n)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -165,3 +189,140 @@ func TestScanDispatchJob_SkipsClaimedJob(t *testing.T) {
|
||||
t.Errorf("assignment stolen from the first dispatcher, assigned_to=%q", assignedTo)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanDrainPendingJobs_SkipsClaimedJob is the sibling of
|
||||
// TestScanDispatchJob_SkipsClaimedJob for the other dispatcher.
|
||||
//
|
||||
// drainPendingJobs collects every pending row up front, then assigns and sends
|
||||
// them one at a time, so a row can be claimed by dispatchJob or the re-dispatch
|
||||
// loop in between. Its UPDATE carries the same `AND status = 'pending'` guard,
|
||||
// but the result was never read: the job was pushed onto this scanner's queue
|
||||
// whether or not this scanner had won it, and both scanners then scanned the
|
||||
// same image. The losing scanner's ack is silently dropped by handleAck (which
|
||||
// does guard on assigned_to), so nothing downstream notices either.
|
||||
//
|
||||
// The interleaving is built deterministically: an unbuffered send channel
|
||||
// parks the drain inside the send for job 1, which is after job 1's UPDATE and
|
||||
// before job 2's, and job 2 is claimed from under it there.
|
||||
func TestScanDrainPendingJobs_SkipsClaimedJob(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 0) // unbuffered: the drain parks in the send
|
||||
sub.capacity = 2 // enough room that the drain would reach job 2
|
||||
seedPendingJobs(t, sb, 2)
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
sb.drainPendingJobs(sub, 0)
|
||||
}()
|
||||
|
||||
// Wait until job 1 is assigned, which means the drain is now parked in the
|
||||
// send for it and has not yet looked at job 2.
|
||||
deadline := time.Now().Add(10 * time.Second)
|
||||
for jobStatus(t, sb, 1) != "assigned" {
|
||||
if time.Now().After(deadline) {
|
||||
t.Fatal("drain never assigned the first job")
|
||||
}
|
||||
time.Sleep(time.Millisecond)
|
||||
}
|
||||
|
||||
// Another dispatcher claims job 2 while the drain is stuck on job 1.
|
||||
if _, err := sb.db.Exec(
|
||||
`UPDATE scan_jobs SET status='assigned', assigned_to='other' WHERE seq = 2 AND status = 'pending'`,
|
||||
); err != nil {
|
||||
t.Fatalf("claim: %v", err)
|
||||
}
|
||||
|
||||
// Release the drain. Job 1 is legitimately ours.
|
||||
select {
|
||||
case job := <-sub.send:
|
||||
if job.Seq != 1 {
|
||||
t.Fatalf("first send was seq %d, want 1", job.Seq)
|
||||
}
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatal("drain never sent the first job")
|
||||
}
|
||||
|
||||
// Job 2 must not follow it. The receive is what makes this decisive: a
|
||||
// drain that sends anyway is parked in that send right now, and would
|
||||
// otherwise give up after its own timeout and look identical to a drain
|
||||
// that correctly skipped it.
|
||||
select {
|
||||
case job := <-sub.send:
|
||||
t.Fatalf("job %d was sent to a scanner that did not claim it; both "+
|
||||
"scanners now scan the same image", job.Seq)
|
||||
case <-time.After(time.Second):
|
||||
}
|
||||
|
||||
<-done
|
||||
|
||||
var assignedTo string
|
||||
if err := sb.db.QueryRow(`SELECT assigned_to FROM scan_jobs WHERE seq = 2`).Scan(&assignedTo); err != nil {
|
||||
t.Fatalf("query assigned_to: %v", err)
|
||||
}
|
||||
if assignedTo != "other" {
|
||||
t.Errorf("assignment stolen from the first dispatcher, assigned_to=%q", assignedTo)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanHasActiveJobs_FailsOpenAfterPersistentDBErrors covers the last way
|
||||
// proactive scanning halts for the life of the process with no recovery.
|
||||
//
|
||||
// hasActiveJobs returned true on any query error ("assume busy") and
|
||||
// waitForCapacity spins on it, so a persistently failing query — a locked
|
||||
// database, a handle closed under a shared connection — stopped proactive
|
||||
// dispatch entirely. logStalledCapacity could not report it either: it runs
|
||||
// the same database and takes its own error branch.
|
||||
//
|
||||
// A transient error should still be treated as busy, since guessing "idle"
|
||||
// piles work onto a scanner that may already have some. A persistent one must
|
||||
// not: after a small budget of consecutive failures the check fails open, and
|
||||
// says so distinctly in the log.
|
||||
func TestScanHasActiveJobs_FailsOpenAfterPersistentDBErrors(t *testing.T) {
|
||||
dbPath := "file:" + t.TempDir() + "/scan.db"
|
||||
sb := newTestScanBroadcaster(t)
|
||||
seedPendingJobs(t, sb, 1)
|
||||
|
||||
// Closing the handle is the cheapest persistent query failure there is.
|
||||
if err := sb.db.Close(); err != nil {
|
||||
t.Fatalf("close db: %v", err)
|
||||
}
|
||||
|
||||
if _, ok := sb.activeProactiveJobs(); ok {
|
||||
t.Error("the first query failure must be treated as busy; a blip is " +
|
||||
"no reason to pile another job on a scanner")
|
||||
}
|
||||
for i := 1; i < activeJobsErrorBudget; i++ {
|
||||
if _, ok := sb.activeProactiveJobs(); ok {
|
||||
t.Errorf("failure %d is still within the budget and must read as busy", i+1)
|
||||
}
|
||||
}
|
||||
if _, ok := sb.activeProactiveJobs(); !ok {
|
||||
t.Fatalf("the capacity check still reports busy after %d consecutive "+
|
||||
"database errors: proactive dispatch is halted for the life of the "+
|
||||
"process with no recovery", activeJobsErrorBudget+1)
|
||||
}
|
||||
|
||||
// A working database restores the budget, so a later blip is absorbed
|
||||
// rather than landing on an already-exhausted counter.
|
||||
db, err := sql.Open("libsql", dbPath)
|
||||
if err != nil {
|
||||
t.Fatalf("reopen db: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = db.Close() })
|
||||
sb.db = db
|
||||
if err := sb.initSchema(); err != nil {
|
||||
t.Fatalf("initSchema: %v", err)
|
||||
}
|
||||
seedJob(t, sb, "sha256:fresh", originProactive)
|
||||
|
||||
if activeProactive(t, sb) != 1 {
|
||||
t.Fatal("a fresh pending job must count as active")
|
||||
}
|
||||
if err := sb.db.Close(); err != nil {
|
||||
t.Fatalf("close db again: %v", err)
|
||||
}
|
||||
if _, ok := sb.activeProactiveJobs(); ok {
|
||||
t.Error("the error budget was not reset by a successful query")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,616 @@
|
||||
package pds
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"atcr.io/pkg/atproto"
|
||||
)
|
||||
|
||||
// These tests cover the ways a scan job wedges on the hold side: the clock the
|
||||
// processing deadline is measured from, what happens to a job that blows it,
|
||||
// and what the hold does with a job whose scanner vanished mid-scan.
|
||||
//
|
||||
// They are characterization tests. Where the behaviour they pin is wrong, the
|
||||
// comment says so and the assertion still describes what the code does today,
|
||||
// so the suite stays green until a fix lands. A test that can only pass after a
|
||||
// fix is marked with t.Skip and names the bug.
|
||||
//
|
||||
// The scanner-side half of the disconnect story lives in
|
||||
// scanner/internal/e2e/stuck_test.go; where a scenario here encodes an
|
||||
// assumption about what the scanner does with a re-offered job, that file pins
|
||||
// it against the real scanner.
|
||||
|
||||
// newStuckBroadcaster is the broadcaster these scenarios need: the in-flight
|
||||
// digest set and the ack timeout the constructor fills in but the bare helper
|
||||
// leaves zeroed, plus the PDS and S3 stand-in behind every terminal
|
||||
// transition. The processing timeout writes a scan record now, so a
|
||||
// broadcaster with a nil pds is no longer a usable stand-in for one.
|
||||
func newStuckBroadcaster(t *testing.T) *ScanBroadcaster {
|
||||
t.Helper()
|
||||
|
||||
sb := newRecordingScanBroadcaster(t)
|
||||
// A buffered signal so a test can assert dispatch capacity was actually
|
||||
// released; production wires the same channel to the dispatch loop.
|
||||
sb.completionSignal = make(chan struct{}, 1)
|
||||
return sb
|
||||
}
|
||||
|
||||
// seedJobWithDigest inserts one pending job carrying a specific manifest
|
||||
// digest. seedPendingJobs gives every row the same digest, which is fine for
|
||||
// status bookkeeping but useless for anything that keys on the digest.
|
||||
func seedJobWithDigest(t *testing.T, sb *ScanBroadcaster, digest string) int64 {
|
||||
t.Helper()
|
||||
|
||||
res, 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)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending')
|
||||
`, digest, "repo", "latest", "did:plc:user", "user.example.com",
|
||||
sb.holdDID, sb.holdEndpoint, "deckhand", "{}", "[]")
|
||||
if err != nil {
|
||||
t.Fatalf("seed job: %v", err)
|
||||
}
|
||||
seq, err := res.LastInsertId()
|
||||
if err != nil {
|
||||
t.Fatalf("seq: %v", err)
|
||||
}
|
||||
return seq
|
||||
}
|
||||
|
||||
// assignJob puts a row in 'assigned' as dispatchJob would, with assigned_at set
|
||||
// the given number of minutes in the past.
|
||||
func assignJob(t *testing.T, sb *ScanBroadcaster, seq int64, sub *ScanSubscriber, agoMinutes int) {
|
||||
t.Helper()
|
||||
|
||||
at := time.Now().Add(-time.Duration(agoMinutes) * time.Minute)
|
||||
_, err := sb.db.Exec(`
|
||||
UPDATE scan_jobs SET status = 'assigned', assigned_to = ?, assigned_at = ?
|
||||
WHERE seq = ?
|
||||
`, sub.id, at, seq)
|
||||
if err != nil {
|
||||
t.Fatalf("assign job %d: %v", seq, err)
|
||||
}
|
||||
}
|
||||
|
||||
func assignedAt(t *testing.T, sb *ScanBroadcaster, seq int64) time.Time {
|
||||
t.Helper()
|
||||
|
||||
var at time.Time
|
||||
if err := sb.db.QueryRow(`SELECT assigned_at FROM scan_jobs WHERE seq = ?`, seq).Scan(&at); err != nil {
|
||||
t.Fatalf("query assigned_at for %d: %v", seq, err)
|
||||
}
|
||||
return at
|
||||
}
|
||||
|
||||
// TestScanAck_DoesNotStartTheScanningClock is the hold-side half of the ack
|
||||
// timing mismatch.
|
||||
//
|
||||
// The scanner acks the moment a job comes off the WebSocket, before it is even
|
||||
// queued (client/hold.go handleFrame sends the ack, then Enqueue). handleAck
|
||||
// moves the row 'assigned' → 'processing' and touches nothing else, which is
|
||||
// correct — the ack means "I have it", not "a worker is on it".
|
||||
//
|
||||
// What was wrong was the deadline. The ten-minute processing budget was
|
||||
// measured from assigned_at, so it covered however long the job spent queued
|
||||
// inside the scanner: a 100-deep queue of no-op jobs drains in 16m40s and
|
||||
// crosses the deadline at position 59, and with 16-second scans at position
|
||||
// 23. Every job past that point was failed underneath a scanner that was
|
||||
// working perfectly, and since the timeout now writes a scan record, each one
|
||||
// is a "scan failed" the user can see.
|
||||
//
|
||||
// The scanning deadline is measured from 'started' instead, and a job that has
|
||||
// only been acked falls under the much larger queue budget.
|
||||
func TestScanAck_DoesNotStartTheScanningClock(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
seq := seedJobWithDigest(t, sb, "sha256:aaaa")
|
||||
assignJob(t, sb, seq, sub, 11)
|
||||
before := assignedAt(t, sb, seq)
|
||||
|
||||
sb.handleAck(sub, seq)
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("after ack: status = %q, want processing", got)
|
||||
}
|
||||
if after := assignedAt(t, sb, seq); !after.Equal(before) {
|
||||
t.Errorf("ack moved assigned_at from %s to %s; it is the dispatch "+
|
||||
"timestamp and the queue budget is measured from it", before, after)
|
||||
}
|
||||
if startedAt(t, sb, seq).Valid {
|
||||
t.Error("the ack stamped started_at; only a worker picking the job up does that")
|
||||
}
|
||||
|
||||
// The scanner acked eleven minutes after dispatch and is still holding the
|
||||
// job. Nothing here is evidence of a problem.
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("status = %q, want processing: a job acked and queued for "+
|
||||
"eleven minutes was cancelled underneath a working scanner", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProcessingTimeout_RetiresTheJobProperly is the compounding half of
|
||||
// the stuck-scan story, and the reason a single hung scan used to remove an
|
||||
// image from scanning until the hold process restarted.
|
||||
//
|
||||
// Every proactive enqueue path adds the manifest digest to sb.inflight and
|
||||
// relies on a terminal transition to take it out again. The ten-minute
|
||||
// processing sweep in reDispatchTimedOut was the one terminal transition that
|
||||
// did neither of the two things all the others do: it wrote no scan record
|
||||
// and never called removeInflight. So the digest stayed in the set —
|
||||
// discoverUnscannedForUser and runStalePass both skip any manifest whose
|
||||
// addInflight returns false — and nothing in the system recorded that the
|
||||
// manifest had ever been attempted, which the appview renders as a grey "Not
|
||||
// scanned" indefinitely.
|
||||
//
|
||||
// The sweep is now a terminal transition like handleError: a failed scan
|
||||
// record with a reason, the digest released, and dispatch capacity signalled.
|
||||
// "Failed" rather than "skipped" is deliberate — a hung scanner is a
|
||||
// transient condition, so the stale loop should retry it on the rescan
|
||||
// interval.
|
||||
func TestScanProcessingTimeout_RetiresTheJobProperly(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:hungscan"
|
||||
seq := seedJobWithDigest(t, sb, digest)
|
||||
assignJob(t, sb, seq, sub, 11)
|
||||
sb.handleAck(sub, seq)
|
||||
markStarted(t, sb, seq, 11) // a worker took it and then went quiet
|
||||
|
||||
if !sb.addInflight(digest) {
|
||||
// The enqueue paths do this; do it here so the state matches.
|
||||
t.Fatal("digest was already in flight before the test started")
|
||||
}
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "failed" {
|
||||
t.Fatalf("status = %q, want failed", got)
|
||||
}
|
||||
|
||||
if !sb.addInflight(digest) {
|
||||
t.Error("the timed-out digest is still in flight: discovery and the " +
|
||||
"stale loop will skip this manifest for the life of the process")
|
||||
}
|
||||
sb.removeInflight(digest)
|
||||
|
||||
_, record, err := sb.pds.GetScanRecord(context.Background(), digest)
|
||||
if err != nil {
|
||||
t.Fatalf("no scan record was written for a job the hold gave up on, so "+
|
||||
"the appview cannot tell it apart from one enqueued a minute ago: %v", err)
|
||||
}
|
||||
if record.Status != atproto.ScanStatusFailed {
|
||||
t.Errorf("record status = %q, want %q", record.Status, atproto.ScanStatusFailed)
|
||||
}
|
||||
if record.Reason == "" {
|
||||
t.Error("failed record carries no reason; it is the only thing a user sees")
|
||||
}
|
||||
|
||||
select {
|
||||
case <-sb.completionSignal:
|
||||
default:
|
||||
t.Error("no completion signal: the dispatch loop sleeps up to 5s longer than it needs to")
|
||||
}
|
||||
|
||||
// And nothing re-offers the row to the scanner that is still connected.
|
||||
select {
|
||||
case job := <-sub.send:
|
||||
t.Fatalf("timed-out processing job %d was re-dispatched", job.Seq)
|
||||
default:
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProcessingTimeout_LeavesAJobInsideTheDeadlineAlone is the guard on
|
||||
// the other side of the sweep: it must only touch rows that have actually
|
||||
// blown the ten minutes, and it must not write records for jobs a scanner is
|
||||
// still legitimately working on.
|
||||
func TestScanProcessingTimeout_LeavesAJobInsideTheDeadlineAlone(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:stillworking"
|
||||
seq := seedJobWithDigest(t, sb, digest)
|
||||
assignJob(t, sb, seq, sub, 2)
|
||||
sb.handleAck(sub, seq)
|
||||
markStarted(t, sb, seq, 2)
|
||||
sb.addInflight(digest)
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Errorf("status = %q, want processing", got)
|
||||
}
|
||||
if sb.addInflight(digest) {
|
||||
t.Error("digest was released while the scan is still inside its deadline")
|
||||
}
|
||||
if _, _, err := sb.pds.GetScanRecord(context.Background(), digest); err == nil {
|
||||
t.Error("a failure record was written for a scan that is still running")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanProcessingTimeout_ReleasesCapacityWhileTheScannerIsStillWedged is
|
||||
// what turns one hung scan into a slow leak rather than a single lost job.
|
||||
//
|
||||
// hasActiveJobs counts 'processing' rows with no age bound, so a wedged scan
|
||||
// holds the proactive dispatch loop still — until the ten-minute timeout marks
|
||||
// it 'failed', at which point capacity is free again and dispatchLoop picks the
|
||||
// next candidate and hands it to the same scanner, whose only worker is still
|
||||
// stuck on the first one. Repeat every ten minutes: each new job is acked,
|
||||
// queued behind the wedge, failed by the timeout, and leaks its digest out of
|
||||
// the in-flight set for good.
|
||||
//
|
||||
// Before dfd604b the same wedge froze dispatch outright, which is the outage
|
||||
// that commit was written for. It bounded 'pending' but left 'processing'
|
||||
// alone, so the freeze became this drip instead.
|
||||
func TestScanProcessingTimeout_ReleasesCapacityWhileTheScannerIsStillWedged(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
seq := seedJob(t, sb, "sha256:wedged", originProactive)
|
||||
assignJob(t, sb, seq, sub, 11)
|
||||
sb.handleAck(sub, seq)
|
||||
markStarted(t, sb, seq, 11)
|
||||
|
||||
if n := activeProactive(t, sb); n != 1 {
|
||||
t.Fatalf("active proactive jobs = %d, want 1: a processing job must "+
|
||||
"hold dispatch capacity", n)
|
||||
}
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
|
||||
if activeProactive(t, sb) != 0 {
|
||||
t.Fatal("capacity is still held after the processing timeout; the " +
|
||||
"drip this test describes cannot happen, update it")
|
||||
}
|
||||
t.Logf("job %d is failed and capacity is free, but the scanner that never "+
|
||||
"answered for it is unchanged: the next candidate goes to the same "+
|
||||
"wedged worker. The digest is released and a failed record written, so "+
|
||||
"the manifest is at least retried on the rescan interval", seq)
|
||||
}
|
||||
|
||||
// TestScanResult_ArrivingAfterTheTimeoutHealsTheRow bounds the previous test:
|
||||
// the leak is permanent only when the scanner never speaks for that seq again.
|
||||
// handleResult has no status guard, so a result that arrives after the deadline
|
||||
// flips 'failed' back to 'completed' and releases the digest.
|
||||
//
|
||||
// That is what makes the leak hard to see in production. It bites exactly in
|
||||
// the cases where the scanner is wedged (an unbounded Syft extraction) or where
|
||||
// its terminal message was written to a dead socket and dropped
|
||||
// (client/hold.go sendJSON) — the same cases where scanning was already stuck.
|
||||
func TestScanResult_ArrivingAfterTheTimeoutHealsTheRow(t *testing.T) {
|
||||
sb := newRecordingScanBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:lateresult"
|
||||
seq := seedJobWithDigest(t, sb, digest)
|
||||
assignJob(t, sb, seq, sub, 11)
|
||||
sb.handleAck(sub, seq)
|
||||
markStarted(t, sb, seq, 11)
|
||||
sb.addInflight(digest)
|
||||
|
||||
sb.reDispatchTimedOut()
|
||||
if got := jobStatus(t, sb, seq); got != "failed" {
|
||||
t.Fatalf("status = %q, want failed", got)
|
||||
}
|
||||
|
||||
// handleResult writes a scan record, so this needs the broadcaster that
|
||||
// has a PDS behind it. The row bookkeeping under test happens after that
|
||||
// write, which is exactly why the write must not be able to panic.
|
||||
sb.handleResult(sub, ScannerMessage{Type: "result", Seq: seq, SBOM: testSBOM})
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "completed" {
|
||||
t.Errorf("status = %q, want completed: a late result overwrites the "+
|
||||
"failure with no status guard", got)
|
||||
}
|
||||
if !sb.addInflight(digest) {
|
||||
t.Error("late result did not release the in-flight digest")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanHandleResult_SurvivesAResultWithoutSummary covers what used to be a
|
||||
// hold crash reachable from any scanner running with vuln.enabled=false.
|
||||
//
|
||||
// The scanner only fills ScanResult.Summary when Grype ran (worker.go
|
||||
// processJob, step 3), and SendResult copies it straight through, so with
|
||||
// vulnerability scanning disabled every successful scan sends a result with no
|
||||
// summary. handleResult guarded the record-writing branch with
|
||||
// `if msg.Summary != nil` and then dereferenced msg.Summary unguarded in its
|
||||
// final log line. That panic was not in an HTTP handler; it was in the
|
||||
// subscriber's reader goroutine, so it took the whole hold process down, and
|
||||
// the job was re-dispatched on restart into the same crash.
|
||||
//
|
||||
// Two things have to hold now. The obvious one is that nothing panics. The
|
||||
// less obvious one is that the SBOM still lands in a scan record: the blob was
|
||||
// uploaded to S3 before the record write, so skipping the record (as the old
|
||||
// `if msg.Summary != nil` guard did) leaves that blob orphaned with nothing
|
||||
// referencing it.
|
||||
//
|
||||
// A nil summary means "not scanned for vulnerabilities", which is NOT the same
|
||||
// as "scanned, found zero". The record therefore carries no vulnerability
|
||||
// report blob, and the appview keys off that to avoid claiming the image is
|
||||
// clean (see classifyScanRecord in pkg/appview/handlers/scan_result.go).
|
||||
func TestScanHandleResult_SurvivesAResultWithoutSummary(t *testing.T) {
|
||||
sb := newRecordingScanBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:nosummary"
|
||||
seq := seedJobWithDigest(t, sb, digest)
|
||||
assignJob(t, sb, seq, sub, 0)
|
||||
sb.handleAck(sub, seq)
|
||||
|
||||
func() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
t.Fatalf("handleResult panicked on a summary-less result: %v", r)
|
||||
}
|
||||
}()
|
||||
// Exactly the message a vuln-disabled scanner sends: an SBOM, no
|
||||
// summary, no vulnerability report.
|
||||
sb.handleResult(sub, ScannerMessage{Type: "result", Seq: seq, SBOM: testSBOM})
|
||||
}()
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "completed" {
|
||||
t.Errorf("status = %q, want completed", got)
|
||||
}
|
||||
|
||||
_, record, err := sb.pds.GetScanRecord(context.Background(), digest)
|
||||
if err != nil {
|
||||
t.Fatalf("no scan record was written for a summary-less result, so the "+
|
||||
"SBOM blob already in S3 is orphaned: %v", err)
|
||||
}
|
||||
if record.SbomBlob == nil {
|
||||
t.Error("scan record carries no SBOM blob; the uploaded blob is orphaned")
|
||||
}
|
||||
if record.VulnReportBlob != nil {
|
||||
t.Error("scan record carries a vulnerability report that was never produced")
|
||||
}
|
||||
if record.Status != atproto.ScanStatusOK {
|
||||
t.Errorf("status = %q, want %q: the scan itself succeeded", record.Status, atproto.ScanStatusOK)
|
||||
}
|
||||
if record.Total != 0 || record.Critical != 0 || record.High != 0 || record.Medium != 0 || record.Low != 0 {
|
||||
t.Errorf("counts = %d/%d/%d/%d total %d, want all zero: Grype never ran",
|
||||
record.Critical, record.High, record.Medium, record.Low, record.Total)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanHandleResult_RecordsSummaryCountsWhenGrypeRan is the other half of
|
||||
// the pair: with a summary present the counts must reach the record, so the
|
||||
// nil-tolerant path above cannot be satisfied by dropping them everywhere.
|
||||
func TestScanHandleResult_RecordsSummaryCountsWhenGrypeRan(t *testing.T) {
|
||||
sb := newRecordingScanBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:withsummary"
|
||||
seq := seedJobWithDigest(t, sb, digest)
|
||||
assignJob(t, sb, seq, sub, 0)
|
||||
sb.handleAck(sub, seq)
|
||||
|
||||
sb.handleResult(sub, ScannerMessage{
|
||||
Type: "result",
|
||||
Seq: seq,
|
||||
SBOM: testSBOM,
|
||||
VulnReport: `{"matches":[]}`,
|
||||
Summary: &VulnerabilitySummary{Critical: 1, High: 2, Medium: 3, Low: 4, Total: 10},
|
||||
})
|
||||
|
||||
_, record, err := sb.pds.GetScanRecord(context.Background(), digest)
|
||||
if err != nil {
|
||||
t.Fatalf("get scan record: %v", err)
|
||||
}
|
||||
if record.Critical != 1 || record.High != 2 || record.Medium != 3 || record.Low != 4 || record.Total != 10 {
|
||||
t.Errorf("counts = %d/%d/%d/%d total %d, want 1/2/3/4 total 10",
|
||||
record.Critical, record.High, record.Medium, record.Low, record.Total)
|
||||
}
|
||||
if record.SbomBlob == nil || record.VulnReportBlob == nil {
|
||||
t.Errorf("blobs = sbom:%v vuln:%v, want both", record.SbomBlob != nil, record.VulnReportBlob != nil)
|
||||
}
|
||||
}
|
||||
|
||||
// testSBOM is a stand-in for the SPDX document the scanner sends. Only its
|
||||
// bytes matter here: the hold hashes them into a blob CID and never parses it.
|
||||
const testSBOM = `{"spdxVersion":"SPDX-2.3","packages":[]}`
|
||||
|
||||
// TestScanUnsubscribe_HoldsAJobForAReconnectingScanner is the hold-side
|
||||
// assumption behind the duplicate-scan scenario in
|
||||
// scanner/internal/e2e/stuck_test.go, updated for a hold that expects several
|
||||
// scanner processes.
|
||||
//
|
||||
// Unsubscribe used to flip every assigned and processing row belonging to the
|
||||
// dropped subscriber back to 'pending'. Nothing tells the scanner: its worker
|
||||
// pool still has the job queued, or a worker is halfway through downloading
|
||||
// blobs for it. drainPendingJobs then handed the same seq straight to whoever
|
||||
// connected next. Against a single scanner that was a duplicate against
|
||||
// itself, which a local dedupe could in principle catch. Against N processes
|
||||
// nothing can catch it, because the two copies are in different processes.
|
||||
//
|
||||
// A disconnect is now recorded rather than acted on, and a scanner keeps one
|
||||
// identity for the life of its process, so the reconnecting scanner gets its
|
||||
// own work back and a different process is never offered it.
|
||||
func TestScanUnsubscribe_HoldsAJobForAReconnectingScanner(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
first := newTestScanSubscriber(t, sb, 4)
|
||||
first.id = "instance-1"
|
||||
|
||||
seq := seedJobWithDigest(t, sb, "sha256:duplicated")
|
||||
assignJob(t, sb, seq, first, 0)
|
||||
sb.handleAck(first, seq)
|
||||
sb.handleStarted(first, seq)
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("status = %q, want processing", got)
|
||||
}
|
||||
|
||||
sb.Unsubscribe(first)
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Fatalf("after disconnect: status = %q, want processing", got)
|
||||
}
|
||||
|
||||
// A different process connecting must not be offered it.
|
||||
other := newTestScanSubscriber(t, sb, 4)
|
||||
other.id = "instance-2"
|
||||
sb.drainPendingJobs(other, 0)
|
||||
|
||||
select {
|
||||
case job := <-other.send:
|
||||
t.Fatalf("seq %d was handed to a second scanner process while the first "+
|
||||
"is still scanning it", job.Seq)
|
||||
default:
|
||||
}
|
||||
|
||||
// The original process redialing resumes it. Nothing is re-sent: the
|
||||
// scanner never lost the job, only the socket.
|
||||
sb.resumeInstance("instance-1")
|
||||
if jobIsDisconnected(t, sb, seq) {
|
||||
t.Error("the reconnecting scanner did not get its own job back")
|
||||
}
|
||||
if got := jobStatus(t, sb, seq); got != "processing" {
|
||||
t.Errorf("status = %q, want processing", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanUnsubscribe_DoesNotReleaseTheInFlightDigest is a smaller sibling of
|
||||
// the timeout leak: a disconnect leaves the digest in flight. That is harmless
|
||||
// — the job is still the disconnected scanner's, and will either be resumed or
|
||||
// reclaimed — but it means the set cannot be read as "jobs a scanner is
|
||||
// working on".
|
||||
func TestScanUnsubscribe_DoesNotReleaseTheInFlightDigest(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:disconnected"
|
||||
seq := seedJobWithDigest(t, sb, digest)
|
||||
assignJob(t, sb, seq, sub, 0)
|
||||
sb.addInflight(digest)
|
||||
|
||||
sb.Unsubscribe(sub)
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "assigned" {
|
||||
t.Fatalf("status = %q, want assigned: the row stays the disconnected "+
|
||||
"scanner's until it resumes or the grace expires", got)
|
||||
}
|
||||
if !jobIsDisconnected(t, sb, seq) {
|
||||
t.Error("the row was not marked as belonging to a disconnected scanner")
|
||||
}
|
||||
if sb.addInflight(digest) {
|
||||
t.Error("Unsubscribe released the in-flight digest: behaviour changed")
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanDispatchQueue_ReturnsUndeliverableJobsToPending covers what the
|
||||
// drain does when it cannot hand a scanner everything it claimed.
|
||||
//
|
||||
// sub.send is 20 deep in production. dispatchJob's default branch resets the
|
||||
// row to 'pending' when the buffer is full; drainPendingJobs used to block for
|
||||
// up to five seconds and then simply return, leaving every row it had already
|
||||
// marked 'assigned' owned by a subscriber that was never sent them. Those rows
|
||||
// sat in 'assigned' — counted as active dispatch capacity — until the
|
||||
// five-minute ack timeout reclaimed them.
|
||||
//
|
||||
// The drain now puts the row it could not deliver back to 'pending', so the
|
||||
// re-dispatch loop can offer it again on its next tick.
|
||||
func TestScanDispatchQueue_ReturnsUndeliverableJobsToPending(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 2) // tiny buffer, nobody draining it
|
||||
sub.capacity = 4 // the send buffer, not capacity, is the constraint here
|
||||
|
||||
var seqs []int64
|
||||
for i := 0; i < 4; i++ {
|
||||
seqs = append(seqs, seedJobWithDigest(t, sb, "sha256:burst"))
|
||||
}
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
sb.drainPendingJobs(sub, 0)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(30 * time.Second):
|
||||
t.Fatal("drainPendingJobs never returned")
|
||||
}
|
||||
|
||||
var assigned, pending int
|
||||
for _, seq := range seqs {
|
||||
switch got := jobStatus(t, sb, seq); got {
|
||||
case "assigned":
|
||||
assigned++
|
||||
case "pending":
|
||||
pending++
|
||||
default:
|
||||
t.Errorf("job %d in unexpected status %q", seq, got)
|
||||
}
|
||||
}
|
||||
|
||||
// Two made it into the buffer and are genuinely assigned. The third blocked
|
||||
// out the five-second window and must have been handed back, and the fourth
|
||||
// was never claimed at all.
|
||||
if assigned != 2 || pending != 2 {
|
||||
t.Errorf("assigned=%d pending=%d, want 2 assigned / 2 pending: a row the "+
|
||||
"scanner was never sent must not stay assigned to it", assigned, pending)
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanSkipped_RetiresAnUndecodableFrame is the hold-side half of the
|
||||
// scanner's answer to a frame it cannot parse
|
||||
// (TestUnparseableFramesAreAnsweredWithSkipped in
|
||||
// scanner/internal/e2e/protocol_test.go).
|
||||
//
|
||||
// The scanner now replies "skipped" for any job frame carrying a usable seq,
|
||||
// which is only worth doing if it actually retires the row. It does: the job
|
||||
// goes to 'completed', a skipped scan record lands in the PDS so the stale
|
||||
// loop leaves it alone, and the in-flight digest is released. Before that
|
||||
// reply existed, the row stayed 'assigned', timed out after five minutes, was
|
||||
// re-offered to the same scanner within thirty seconds, and was dropped again
|
||||
// — while hasActiveJobs counted it and no proactive scan was dispatched
|
||||
// anywhere in the deployment.
|
||||
func TestScanSkipped_RetiresAnUndecodableFrame(t *testing.T) {
|
||||
sb := newStuckBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
|
||||
const digest = "sha256:undecodable"
|
||||
seq := seedJobWithDigest(t, sb, digest)
|
||||
assignJob(t, sb, seq, sub, 0)
|
||||
sb.addInflight(digest)
|
||||
|
||||
// Exactly what connectOnce sends when the config sub-document does not
|
||||
// decode: no ack ever arrives, so the row is still 'assigned'.
|
||||
sb.handleSkipped(sub, ScannerMessage{
|
||||
Type: "skipped",
|
||||
Seq: seq,
|
||||
Reason: "malformed job config: json: cannot unmarshal string into Go value of type scanner.BlobDescriptor",
|
||||
})
|
||||
|
||||
if got := jobStatus(t, sb, seq); got != "completed" {
|
||||
t.Errorf("status = %q, want completed: a skip must be terminal", got)
|
||||
}
|
||||
if !sb.addInflight(digest) {
|
||||
t.Error("skip did not release the in-flight digest")
|
||||
}
|
||||
sb.removeInflight(digest)
|
||||
|
||||
_, record, err := sb.pds.GetScanRecord(context.Background(), digest)
|
||||
if err != nil {
|
||||
t.Fatalf("no scan record for a skipped job: %v", err)
|
||||
}
|
||||
if record.Status != atproto.ScanStatusSkipped {
|
||||
t.Errorf("record status = %q, want %q", record.Status, atproto.ScanStatusSkipped)
|
||||
}
|
||||
|
||||
// Terminal means terminal: the re-dispatch loop must not pick it back up.
|
||||
sb.reDispatchTimedOut()
|
||||
if got := jobStatus(t, sb, seq); got != "completed" {
|
||||
t.Errorf("after a re-dispatch tick: status = %q, want completed", got)
|
||||
}
|
||||
select {
|
||||
case job := <-sub.send:
|
||||
t.Fatalf("a skipped job (%d) was re-offered to the scanner", job.Seq)
|
||||
default:
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,8 @@ import (
|
||||
"database/sql"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"atcr.io/pkg/s3"
|
||||
)
|
||||
|
||||
// newTestScanBroadcaster builds a ScanBroadcaster with just a database, no
|
||||
@@ -30,6 +32,26 @@ func newTestScanBroadcaster(t *testing.T) *ScanBroadcaster {
|
||||
return sb
|
||||
}
|
||||
|
||||
// newRecordingScanBroadcaster is newTestScanBroadcaster plus the two
|
||||
// dependencies handleResult needs to finish its work: an embedded PDS to write
|
||||
// the scan record into, and an S3 stand-in to take the SBOM blob. The bare
|
||||
// helper leaves both nil, which is fine for row bookkeeping and fatal for
|
||||
// anything that asserts on what was stored.
|
||||
func newRecordingScanBroadcaster(t *testing.T) *ScanBroadcaster {
|
||||
t.Helper()
|
||||
|
||||
sb := newTestScanBroadcaster(t)
|
||||
sb.inflight = make(map[string]struct{})
|
||||
sb.ackTimeout = 5 * time.Minute
|
||||
|
||||
pds, _ := setupTestPDS(t)
|
||||
sb.pds = pds
|
||||
sb.holdDID = pds.did
|
||||
sb.s3 = &s3.S3Service{Client: s3.NewMockS3Client(""), Bucket: "test-bucket"}
|
||||
|
||||
return sb
|
||||
}
|
||||
|
||||
// newTestScanSubscriber mirrors what Subscribe builds, registered with the
|
||||
// broadcaster so Unsubscribe finds it.
|
||||
func newTestScanSubscriber(t *testing.T, sb *ScanBroadcaster, bufSize int) *ScanSubscriber {
|
||||
@@ -40,6 +62,9 @@ func newTestScanSubscriber(t *testing.T, sb *ScanBroadcaster, bufSize int) *Scan
|
||||
send: make(chan *ScanJobEvent, bufSize),
|
||||
id: "test-subscriber",
|
||||
done: make(chan struct{}),
|
||||
// One worker unless a test says otherwise, which is what a scanner
|
||||
// that declares nothing is treated as.
|
||||
capacity: 1,
|
||||
}
|
||||
|
||||
sb.mu.Lock()
|
||||
@@ -93,41 +118,66 @@ func TestScanUnsubscribe_IsIdempotent(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestScanUnsubscribe_UnassignsJobsOnce verifies the idempotency guard protects
|
||||
// the job-reassignment UPDATE too. Re-running it would unassign jobs that a
|
||||
// replacement scanner had already been given.
|
||||
func TestScanUnsubscribe_UnassignsJobsOnce(t *testing.T) {
|
||||
// TestScanUnsubscribe_MarksItsOwnJobsOnce verifies the idempotency guard
|
||||
// protects the disconnect bookkeeping too, and that the bookkeeping is scoped
|
||||
// to this subscriber's rows.
|
||||
//
|
||||
// Unsubscribe used to flip every assigned and processing row straight back to
|
||||
// 'pending', which handed a running scan to whichever process connected next
|
||||
// (see TestScanUnsubscribe_DoesNotImmediatelyHandOffWorkStillRunning). It now
|
||||
// marks the disconnect and leaves the work where it is. Either way the guard
|
||||
// matters for the same reason: a dropped scanner unwinds both handleWriter and
|
||||
// handleReader, and each calls Unsubscribe.
|
||||
func TestScanUnsubscribe_MarksItsOwnJobsOnce(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 4)
|
||||
seedPendingJobs(t, sb, 1)
|
||||
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET status='assigned', assigned_to=?`, sub.id); err != nil {
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET status='processing', assigned_to=?`, sub.id); err != nil {
|
||||
t.Fatalf("assign: %v", err)
|
||||
}
|
||||
|
||||
sb.Unsubscribe(sub)
|
||||
|
||||
var status string
|
||||
if err := sb.db.QueryRow(`SELECT status FROM scan_jobs LIMIT 1`).Scan(&status); err != nil {
|
||||
t.Fatalf("query: %v", err)
|
||||
}
|
||||
if status != "pending" {
|
||||
t.Errorf("expected job returned to pending, got %q", status)
|
||||
var (
|
||||
status string
|
||||
assignedTo sql.NullString
|
||||
disconnectedAt sql.NullTime
|
||||
)
|
||||
row := func() {
|
||||
t.Helper()
|
||||
if err := sb.db.QueryRow(
|
||||
`SELECT status, assigned_to, disconnected_at FROM scan_jobs LIMIT 1`,
|
||||
).Scan(&status, &assignedTo, &disconnectedAt); err != nil {
|
||||
t.Fatalf("query: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
row()
|
||||
if status != "processing" || assignedTo.String != sub.id {
|
||||
t.Errorf("job = %q/%q, want processing/%q: a disconnect is not evidence "+
|
||||
"the scanner stopped scanning", status, assignedTo.String, sub.id)
|
||||
}
|
||||
if !disconnectedAt.Valid {
|
||||
t.Error("the row was not marked as belonging to a disconnected scanner")
|
||||
}
|
||||
firstMark := disconnectedAt.Time
|
||||
|
||||
// Hand the job to a "replacement" scanner, then unsubscribe the dead one
|
||||
// again. The guard must stop it from stealing the job back.
|
||||
if _, err := sb.db.Exec(`UPDATE scan_jobs SET status='assigned', assigned_to=?`, "replacement"); err != nil {
|
||||
// again. The guard must stop it from touching a row that has moved on.
|
||||
if _, err := sb.db.Exec(
|
||||
`UPDATE scan_jobs SET status='assigned', assigned_to=?, disconnected_at=NULL`, "replacement",
|
||||
); err != nil {
|
||||
t.Fatalf("reassign: %v", err)
|
||||
}
|
||||
sb.Unsubscribe(sub)
|
||||
|
||||
if err := sb.db.QueryRow(`SELECT status FROM scan_jobs LIMIT 1`).Scan(&status); err != nil {
|
||||
t.Fatalf("query: %v", err)
|
||||
}
|
||||
if status != "assigned" {
|
||||
t.Errorf("second Unsubscribe stole the replacement's job, status=%q", status)
|
||||
row()
|
||||
if status != "assigned" || assignedTo.String != "replacement" || disconnectedAt.Valid {
|
||||
t.Errorf("second Unsubscribe touched the replacement's job: %q/%q "+
|
||||
"disconnected=%v", status, assignedTo.String, disconnectedAt.Valid)
|
||||
}
|
||||
_ = firstMark
|
||||
}
|
||||
|
||||
// TestScanDrainPendingJobs_ConcurrentUnsubscribe is the regression test for
|
||||
@@ -143,6 +193,7 @@ func TestScanDrainPendingJobs_ConcurrentUnsubscribe(t *testing.T) {
|
||||
func() {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
sub := newTestScanSubscriber(t, sb, 1)
|
||||
sub.capacity = 50 // the drain must walk rows, not stop at capacity
|
||||
seedPendingJobs(t, sb, 50)
|
||||
|
||||
done := make(chan struct{})
|
||||
|
||||
@@ -67,7 +67,7 @@ func TestScanBroadcaster_UnsubscribeIsIdempotent(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
conn, _ := wsPair(t)
|
||||
|
||||
sub := sb.Subscribe(conn, 0)
|
||||
sub := sb.Subscribe(conn, 0, "", 0)
|
||||
if got := subscriberCount(sb); got != 1 {
|
||||
t.Fatalf("subscriber count after Subscribe = %d, want 1", got)
|
||||
}
|
||||
@@ -88,7 +88,7 @@ func TestScanBroadcaster_WriterExitsOnUnsubscribe(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
conn, _ := wsPair(t)
|
||||
|
||||
sub := sb.Subscribe(conn, 0)
|
||||
sub := sb.Subscribe(conn, 0, "", 0)
|
||||
sb.Unsubscribe(sub)
|
||||
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
@@ -108,7 +108,7 @@ func TestScanBroadcaster_DroppedScannerUnwindsOnce(t *testing.T) {
|
||||
sb := newTestScanBroadcaster(t)
|
||||
conn, closeClient := wsPair(t)
|
||||
|
||||
sub := sb.Subscribe(conn, 0)
|
||||
sub := sb.Subscribe(conn, 0, "", 0)
|
||||
closeClient()
|
||||
|
||||
deadline := time.Now().Add(5 * time.Second)
|
||||
|
||||
+19
-1
@@ -1101,6 +1101,24 @@ func (h *XRPCHandler) HandleSubscribeScanJobs(w http.ResponseWriter, r *http.Req
|
||||
}
|
||||
}
|
||||
|
||||
// How much scanning this process can run at once, and who it is.
|
||||
//
|
||||
// Both are optional: a scanner built before they existed sends neither,
|
||||
// and the broadcaster then treats it as one worker with a
|
||||
// per-connection identity — exactly the behaviour it had before. A
|
||||
// malformed workers value is ignored rather than rejected, since the
|
||||
// connection is still perfectly usable at the default.
|
||||
workers := 0
|
||||
if raw := r.URL.Query().Get("workers"); raw != "" {
|
||||
if n, err := strconv.Atoi(raw); err == nil && n > 0 {
|
||||
workers = n
|
||||
} else {
|
||||
slog.Warn("Scanner declared an unusable worker count, using the default",
|
||||
"workers", raw)
|
||||
}
|
||||
}
|
||||
instanceID := r.URL.Query().Get("instance")
|
||||
|
||||
// Upgrade to WebSocket
|
||||
conn, err := upgrader.Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
@@ -1108,7 +1126,7 @@ func (h *XRPCHandler) HandleSubscribeScanJobs(w http.ResponseWriter, r *http.Req
|
||||
return
|
||||
}
|
||||
|
||||
h.scanBroadcaster.Subscribe(conn, cursor)
|
||||
h.scanBroadcaster.Subscribe(conn, cursor, instanceID, workers)
|
||||
}
|
||||
|
||||
// ScanBroadcasterRef returns the scan broadcaster (used by OCI handler to enqueue jobs)
|
||||
|
||||
Reference in New Issue
Block a user