Show egress/ingress levels in the registry; freeze checks at the verdict

Registry: the "last result" column now also shows, per level (egress,
ingress), how many of the recorded checks of the latest cycle succeeded, split
by check family (tcp-22 and tcp-443 are both "tcp"). One grouped query per
chunk of addresses; new fields last_cycle_id, egress, ingress in
GET /admin/registry; the dashboard renders them under the verdict.

Verdict integrity (migration 0010):
- the prober is handed an address once per site and attempt, not on every
  poll, so results are no longer overwritten by later probe rounds;
- UpsertCheckIfOpen refuses writes once the address is aggregating or has its
  verdict, or for an older attempt; senders get {"ok":true,"ignored":N} and a
  result_dropped event is recorded;
- the checking window counts from checking_started_at, not from assigned_at;
- checks.recorded_at (server clock) and checks.after_verdict (flag for rows
  written after the verdict in existing data);
- the verdict rule is a pure function (computeVerdict) and the aggregated
  event carries the egress/ingress check counts.

Rebuilt bin/control-api and bin/admin-dashboard to match. Plans and summaries
are in docs/changes; README, API, USAGE, DASHBOARD and DIAGRAMS are updated.

Co-Authored-By: Claude Sonnet 5.5 <noreply@anthropic.com>
This commit is contained in:
ayurishchevandClaude Sonnet 5.5 committed 2026-10-03 17:59:52 +03:00
1 parent db73409e8f
commit 864208238f
34 files changed
+1570 -72

No files matched your search

+4
View File
@@ -40,6 +40,9 @@ var autoCycleSchema string
//go:embed migrations/0009_scale_indexes.sql
var scaleIndexesSchema string
//go:embed migrations/0010_verdict_integrity.sql
var verdictIntegritySchema string
// migrations is the ordered list of schema versions. Each entry's SQL is
// applied, in order, for any version greater than the database's current
// PRAGMA user_version — so a fresh database walks the whole list and an
@@ -57,6 +60,7 @@ var migrations = []struct {
{7, ipRegistrySchema},
{8, autoCycleSchema},
{9, scaleIndexesSchema},
{10, verdictIntegritySchema},
}
type DB struct {
@@ -0,0 +1,31 @@
-- Verdict integrity (see docs/changes/2026-10-03_17-21_verdict-no-late-results-plan.md).
--
-- ip_queue.checking_started_at: when the address entered the "checking" state.
-- The aggregation window (checking_window_seconds) is counted from here, not
-- from assigned_at, which also covers floating-IP association, the settle
-- pause and the self-check. NULL on rows that predate this migration; the
-- orchestrator falls back to assigned_at for them.
--
-- checks.recorded_at: when control-api last wrote the row, by its own clock.
-- checked_at comes from the probing machine's clock, so it cannot be compared
-- reliably with aggregated_at. Existing rows get created_at (their first
-- write); a later overwrite is not recoverable.
--
-- checks.after_verdict: 1 when the row was written after the address's
-- verdict (checked_at later than aggregated_at). Set here for existing data
-- only; new writes after the verdict are rejected, so it stays 0 afterwards.
ALTER TABLE ip_queue ADD COLUMN checking_started_at TIMESTAMP;
ALTER TABLE checks ADD COLUMN recorded_at TIMESTAMP;
UPDATE checks SET recorded_at = created_at;
ALTER TABLE checks ADD COLUMN after_verdict INTEGER NOT NULL DEFAULT 0;
UPDATE checks SET after_verdict = 1
WHERE EXISTS (
SELECT 1 FROM ip_queue q
WHERE q.registry_id = checks.registry_id
AND q.cycle_id = checks.cycle_id
AND q.aggregated_at IS NOT NULL
AND julianday(checks.checked_at) > julianday(q.aggregated_at)
);
+45 -8
View File
@@ -1,6 +1,9 @@
package db
import "time"
import (
"strings"
"time"
)
// Validator and IP lifecycle states. Kept as typed string constants rather
// than a Go enum type so they round-trip through SQLite TEXT columns and
@@ -71,6 +74,37 @@ func InboundSource(siteIndex int) string {
return "inbound-site-" + itoa(siteIndex)
}
// Check levels: the two directions a check can run in, derived from
// checks.source (there is no separate direction column).
const (
LevelEgress = "egress"
LevelIngress = "ingress"
)
// CheckLevel maps a checks.source value to its level: "egress" for the
// validator's own outbound checks, "ingress" for any prober site
// ("inbound-site-N"). Any other source yields "".
func CheckLevel(source string) string {
switch {
case source == SourceEgress:
return LevelEgress
case strings.HasPrefix(source, "inbound-site-"):
return LevelIngress
}
return ""
}
// CheckFamily maps a checks.check_type value to its family for grouping:
// the part before the first "-", so tcp-22 and tcp-443 are both "tcp" while
// https, icmp, ssh and tls-443 ("tls") stay distinct. A check type added in
// the future is grouped by its own name without code changes.
func CheckFamily(checkType string) string {
if i := strings.IndexByte(checkType, '-'); i > 0 {
return checkType[:i]
}
return checkType
}
func itoa(n int) string {
if n == 0 {
return "0"
@@ -118,13 +152,16 @@ type IPQueueItem struct {
EgressComplete bool
OverallResult string
AssignedAt *time.Time
FIPAssociatedAt *time.Time
AggregatedAt *time.Time
FIPReleasedAt *time.Time
RegistryID int64
CycleID int
CreatedAt time.Time
UpdatedAt time.Time
// CheckingStartedAt is when the address entered the checking state; the
// aggregation window counts from it. nil on rows that predate it.
CheckingStartedAt *time.Time
FIPAssociatedAt *time.Time
AggregatedAt *time.Time
FIPReleasedAt *time.Time
RegistryID int64
CycleID int
CreatedAt time.Time
UpdatedAt time.Time
}
type Check struct {
+29 -8
View File
@@ -15,20 +15,41 @@ import (
// per-address cycle rather than the ip_queue row's attempt_number, so it
// survives that row being deleted and the address later resubmitted.
func (d *DB) UpsertCheck(ctx context.Context, c Check) error {
_, err := d.ExecContext(ctx, `
_, err := d.UpsertCheckIfOpen(ctx, c)
return err
}
// UpsertCheckIfOpen is UpsertCheck for the live write path. It stores the
// check only while the address can still take results: the queue row exists,
// c.AttemptNumber is its current attempt, and it has not reached the
// aggregating state or a terminal one. Once the verdict is being computed the
// checks are frozen, so the verdict always matches the stored rows and a late
// probe of an already released floating IP cannot change them. It returns
// false (and writes nothing) when the result was dropped. The state test and
// the write are one statement, so they cannot interleave with SetAggregating.
func (d *DB) UpsertCheckIfOpen(ctx context.Context, c Check) (bool, error) {
now := timeToDB(Now())
res, err := d.ExecContext(ctx, `
INSERT INTO checks (registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id,
source, check_type, target, success, latency_ms, detail, checked_at, created_at)
VALUES ((SELECT registry_id FROM ip_queue WHERE id=?), (SELECT cycle_id FROM ip_queue WHERE id=?),
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
source, check_type, target, success, latency_ms, detail, checked_at, created_at, recorded_at)
SELECT q.registry_id, q.cycle_id, q.id, ?, q.attempt_number, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
FROM ip_queue q
WHERE q.id=? AND q.attempt_number=? AND q.state NOT IN (?, ?, ?, ?)
ON CONFLICT(registry_id, cycle_id, source, check_type, target) DO UPDATE SET
validator_id=excluded.validator_id,
success=excluded.success,
latency_ms=excluded.latency_ms,
detail=excluded.detail,
checked_at=excluded.checked_at
`, c.IPID, c.IPID, c.IPID, c.IPAddress, c.AttemptNumber, c.ValidatorID, c.Source, c.CheckType, c.Target,
c.Success, c.LatencyMS, c.Detail, timeToDB(c.CheckedAt), timeToDB(Now()))
return err
checked_at=excluded.checked_at,
recorded_at=excluded.recorded_at
`, c.IPAddress, c.ValidatorID, c.Source, c.CheckType, c.Target,
c.Success, c.LatencyMS, c.Detail, timeToDB(c.CheckedAt), now, now,
c.IPID, c.AttemptNumber, IPAggregating, IPDone, IPFailed, IPOccupied)
if err != nil {
return false, err
}
n, err := res.RowsAffected()
return n > 0, err
}
const checksSelect = `
+28 -7
View File
@@ -180,9 +180,9 @@ func (d *DB) KnownAddresses(ctx context.Context, addresses []string) (map[string
func (d *DB) SetChecking(ctx context.Context, ipID int64, leaseTTL time.Duration) error {
now := Now()
_, err := d.ExecContext(ctx, `
UPDATE ip_queue SET state=?, lease_expires_at=?, updated_at=?
UPDATE ip_queue SET state=?, lease_expires_at=?, checking_started_at=?, updated_at=?
WHERE id=?
`, IPChecking, timeToDB(now.Add(leaseTTL)), timeToDB(now), ipID)
`, IPChecking, timeToDB(now.Add(leaseTTL)), timeToDB(now), timeToDB(now), ipID)
return err
}
@@ -295,7 +295,7 @@ func (d *DB) RequeueOrFail(ctx context.Context, ipID int64, validatorID string,
UPDATE ip_queue SET
state=?, owner_validator_id=NULL, fip_id='', retry_count=?, attempt_number=attempt_number+1,
cycle_id=?, lease_expires_at=NULL, egress_complete=0,
overall_result='', assigned_at=NULL, fip_associated_at=NULL, updated_at=?
overall_result='', assigned_at=NULL, checking_started_at=NULL, fip_associated_at=NULL, updated_at=?
WHERE id=?
`, nextState, retryCount, cycle, now, ipID)
} else {
@@ -394,7 +394,7 @@ func (d *DB) SubmitIPs(ctx context.Context, addresses []string) (SubmitIPsResult
state=?, sequence=?, owner_validator_id=NULL, fip_id='', retry_count=0,
attempt_number=attempt_number+1, cycle_id=?, lease_expires_at=NULL, egress_complete=0,
overall_result='',
assigned_at=NULL, fip_associated_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=?
assigned_at=NULL, checking_started_at=NULL, fip_associated_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=?
WHERE ip_address=?
`, IPQueued, seq, cycle, now, addr); err != nil {
return result, fmt.Errorf("requeue %s: %w", addr, err)
@@ -835,6 +835,24 @@ func (d *DB) ListChecking(ctx context.Context) ([]IPQueueItem, error) {
return scanIPQueueItems(rows)
}
// ListCheckingForSite returns the addresses in the checking state that the
// given prober site still has to probe: those for which it has not yet
// reported completion in the current attempt. Without this filter a prober
// would re-probe every checking address on every poll until the verdict.
func (d *DB) ListCheckingForSite(ctx context.Context, siteIndex int) ([]IPQueueItem, error) {
rows, err := d.QueryContext(ctx, ipQueueSelect+`
WHERE state=? AND NOT EXISTS (
SELECT 1 FROM ip_site_checks s
WHERE s.ip_id=ip_queue.id AND s.attempt_number=ip_queue.attempt_number
AND s.site_idx=? AND s.complete=1)
ORDER BY sequence`, IPChecking, siteIndex)
if err != nil {
return nil, err
}
defer rows.Close()
return scanIPQueueItems(rows)
}
// ListExpiredLeases returns non-terminal IPs whose lease has expired —
// candidates for the lease sweep (crash recovery + stuck-validator reclaim).
func (d *DB) ListExpiredLeases(ctx context.Context, now time.Time) ([]IPQueueItem, error) {
@@ -851,7 +869,7 @@ func (d *DB) ListExpiredLeases(ctx context.Context, now time.Time) ([]IPQueueIte
const ipQueueSelect = `
SELECT id, ip_address, sequence, state, owner_validator_id, fip_id, attempt_number, retry_count,
lease_expires_at, egress_complete, overall_result,
assigned_at, fip_associated_at, aggregated_at, fip_released_at,
assigned_at, checking_started_at, fip_associated_at, aggregated_at, fip_released_at,
registry_id, cycle_id, created_at, updated_at
FROM ip_queue
`
@@ -871,14 +889,14 @@ func scanIPQueueItems(rows *sql.Rows) ([]IPQueueItem, error) {
func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) {
var item IPQueueItem
var owner sql.NullString
var leaseExpires, assignedAt, fipAssociatedAt, aggregatedAt, fipReleasedAt sql.NullString
var leaseExpires, assignedAt, checkingStartedAt, fipAssociatedAt, aggregatedAt, fipReleasedAt sql.NullString
var registryID sql.NullInt64
var createdAt, updatedAt string
if err := row.Scan(
&item.ID, &item.IPAddress, &item.Sequence, &item.State, &owner, &item.FIPID,
&item.AttemptNumber, &item.RetryCount, &leaseExpires,
&item.EgressComplete,
&item.OverallResult, &assignedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt,
&item.OverallResult, &assignedAt, &checkingStartedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt,
&registryID, &item.CycleID, &createdAt, &updatedAt,
); err != nil {
return nil, err
@@ -896,6 +914,9 @@ func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) {
if item.AssignedAt, err = nullStringToTimePtr(assignedAt); err != nil {
return nil, err
}
if item.CheckingStartedAt, err = nullStringToTimePtr(checkingStartedAt); err != nil {
return nil, err
}
if item.FIPAssociatedAt, err = nullStringToTimePtr(fipAssociatedAt); err != nil {
return nil, err
}
+126
View File
@@ -4,6 +4,7 @@ import (
"context"
"database/sql"
"fmt"
"sort"
"strings"
"time"
)
@@ -56,6 +57,28 @@ type RegistrySummary struct {
LastCheckedAt *time.Time
InQueue bool
CurrentState string
// LastCycleID is the newest cycle_id with recorded checks (0 if none).
// Egress and Ingress count that cycle's recorded checks per level.
LastCycleID int
Egress LevelResult
Ingress LevelResult
}
// TypeStat counts the recorded checks of one check family (see CheckFamily)
// within a level: Total checks, OK of them successful.
type TypeStat struct {
Type string
Total int
OK int
}
// LevelResult is the "OK of Total" rollup of one level (egress or ingress) of
// a cycle, with the same counts split by check family, sorted by Type.
type LevelResult struct {
Total int
OK int
ByType []TypeStat
}
// ListRegistry returns every address ever submitted, newest first-seen
@@ -87,6 +110,9 @@ func (d *DB) ListRegistry(ctx context.Context) ([]RegistrySummary, error) {
}
out[i] = s
}
if err := d.fillRegistryLevels(ctx, out); err != nil {
return nil, err
}
return out, nil
}
@@ -165,6 +191,9 @@ func (d *DB) ListRegistryPage(ctx context.Context, f RegistryFilter, limit, offs
}
out[i] = s
}
if err := d.fillRegistryLevels(ctx, out); err != nil {
return nil, 0, err
}
return out, total, nil
}
@@ -186,6 +215,11 @@ func (d *DB) GetRegistryByAddress(ctx context.Context, address string) (*Registr
if err := d.fillRegistrySummary(ctx, s); err != nil {
return nil, err
}
one := []RegistrySummary{*s}
if err := d.fillRegistryLevels(ctx, one); err != nil {
return nil, err
}
*s = one[0]
return s, nil
}
@@ -253,6 +287,98 @@ func (d *DB) fillRegistrySummary(ctx context.Context, s *RegistrySummary) error
return nil
}
// registryLevelsChunk bounds the number of registry ids per query, well under
// SQLite's bound-variable limit.
const registryLevelsChunk = 500
// registryLevelsQuery is the grouped query of fillRegistryLevels for n
// registry ids: per address, the counts of its newest cycle by source and
// check type.
func registryLevelsQuery(n int) string {
return `
SELECT c.registry_id, m.cid, c.source, c.check_type, COUNT(*), COALESCE(SUM(c.success), 0)
FROM checks c
JOIN (SELECT registry_id, MAX(cycle_id) AS cid FROM checks
WHERE registry_id IN (` + strings.TrimSuffix(strings.Repeat("?,", n), ",") + `)
GROUP BY registry_id) m
ON m.registry_id = c.registry_id AND m.cid = c.cycle_id
GROUP BY c.registry_id, m.cid, c.source, c.check_type
`
}
// fillRegistryLevels sets LastCycleID, Egress and Ingress on every summary in
// sums: the recorded checks of each address's newest cycle (the same cycle
// whose time is LastCheckedAt), counted per level and per check family. It
// runs one grouped query per chunk of addresses, not one per address, over
// idx_checks_registry_cycle. Checks whose source is neither egress nor an
// inbound site are not counted. The counts follow the recorded rows only, so
// they can differ from LastResult, which also treats missing results as
// failures.
func (d *DB) fillRegistryLevels(ctx context.Context, sums []RegistrySummary) error {
pos := make(map[int64]int, len(sums))
for i := range sums {
pos[sums[i].ID] = i
}
type key struct {
id int64
level, family string
}
for start := 0; start < len(sums); start += registryLevelsChunk {
end := min(start+registryLevelsChunk, len(sums))
args := make([]any, 0, end-start)
for _, s := range sums[start:end] {
args = append(args, s.ID)
}
rows, err := d.QueryContext(ctx, registryLevelsQuery(len(args)), args...)
if err != nil {
return err
}
stats := map[key]*TypeStat{}
for rows.Next() {
var id int64
var cid, total, ok int
var source, checkType string
if err := rows.Scan(&id, &cid, &source, &checkType, &total, &ok); err != nil {
rows.Close()
return err
}
sums[pos[id]].LastCycleID = cid
level := CheckLevel(source)
if level == "" {
continue
}
k := key{id, level, CheckFamily(checkType)}
st := stats[k]
if st == nil {
st = &TypeStat{Type: k.family}
stats[k] = st
}
st.Total += total
st.OK += ok
}
err = rows.Err()
rows.Close()
if err != nil {
return err
}
for k, st := range stats {
lr := &sums[pos[k.id]].Egress
if k.level == LevelIngress {
lr = &sums[pos[k.id]].Ingress
}
lr.Total += st.Total
lr.OK += st.OK
lr.ByType = append(lr.ByType, *st)
}
}
for i := range sums {
for _, lr := range []*LevelResult{&sums[i].Egress, &sums[i].Ingress} {
sort.Slice(lr.ByType, func(a, b int) bool { return lr.ByType[a].Type < lr.ByType[b].Type })
}
}
return nil
}
// lastCycleResultFromChecks classifies the most recent cycle recorded for
// registryID directly from its checks rows: pass if every recorded check
// succeeded, fail if every one failed, partial on a mix. Returns "" if no
+169
View File
@@ -0,0 +1,169 @@
package db
import (
"reflect"
"testing"
)
func TestCheckLevelAndFamily(t *testing.T) {
for source, want := range map[string]string{
"egress": LevelEgress, "inbound-site-1": LevelIngress, "inbound-site-12": LevelIngress,
"": "", "other": "",
} {
if got := CheckLevel(source); got != want {
t.Errorf("CheckLevel(%q) = %q, want %q", source, got, want)
}
}
for ct, want := range map[string]string{
"https": "https", "icmp": "icmp", "ssh": "ssh",
"tcp-22": "tcp", "tcp-443": "tcp", "tls-443": "tls", "dns": "dns", "-x": "-x",
} {
if got := CheckFamily(ct); got != want {
t.Errorf("CheckFamily(%q) = %q, want %q", ct, got, want)
}
}
}
// addCheck records one check for the address's current queue row.
func addCheck(t *testing.T, d *DB, addr, source, checkType, target string, success bool) {
t.Helper()
ctx := t.Context()
ip, err := d.GetIPByAddress(ctx, addr)
if err != nil {
t.Fatalf("get %s: %v", addr, err)
}
if err := d.UpsertCheck(ctx, Check{
IPID: ip.ID, IPAddress: addr, AttemptNumber: ip.AttemptNumber,
Source: source, CheckType: checkType, Target: target,
Success: success, CheckedAt: Now(),
}); err != nil {
t.Fatalf("upsert check: %v", err)
}
}
func TestRegistryLevelsGroupByTypeAndLevel(t *testing.T) {
d, ctx := newTestDB(t)
if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatal(err)
}
a := "1.2.3.4"
// Egress: https 2 of 3, icmp 1 of 1.
addCheck(t, d, a, SourceEgress, "https", "https://a.test", true)
addCheck(t, d, a, SourceEgress, "https", "https://b.test", true)
addCheck(t, d, a, SourceEgress, "https", "https://c.test", false)
addCheck(t, d, a, SourceEgress, "icmp", "a.test", true)
// Ingress from two sites: tcp-22 and tcp-443 are one family; tls, ssh,
// icmp and a type unknown today ("dns") are listed on their own.
for site := 1; site <= 2; site++ {
src := InboundSource(site)
addCheck(t, d, a, src, "tcp-22", a, true)
addCheck(t, d, a, src, "tcp-443", a, site == 1)
addCheck(t, d, a, src, "tls-443", a, true)
addCheck(t, d, a, src, "ssh", a, false)
addCheck(t, d, a, src, "icmp", a, true)
addCheck(t, d, a, src, "dns", a, true)
}
// A source that is neither egress nor an inbound site is not counted.
addCheck(t, d, a, "manual", "https", "x", true)
s, err := d.GetRegistryByAddress(ctx, a)
if err != nil {
t.Fatal(err)
}
wantEgress := LevelResult{Total: 4, OK: 3, ByType: []TypeStat{
{"https", 3, 2}, {"icmp", 1, 1},
}}
wantIngress := LevelResult{Total: 12, OK: 9, ByType: []TypeStat{
{"dns", 2, 2}, {"icmp", 2, 2}, {"ssh", 2, 0}, {"tcp", 4, 3}, {"tls", 2, 2},
}}
if !reflect.DeepEqual(s.Egress, wantEgress) {
t.Errorf("egress = %+v, want %+v", s.Egress, wantEgress)
}
if !reflect.DeepEqual(s.Ingress, wantIngress) {
t.Errorf("ingress = %+v, want %+v", s.Ingress, wantIngress)
}
if s.LastCycleID != 1 {
t.Errorf("LastCycleID = %d, want 1", s.LastCycleID)
}
}
func TestRegistryLevelsNoChecksAreZero(t *testing.T) {
d, ctx := newTestDB(t)
if _, err := d.SubmitIPs(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatal(err)
}
s, err := d.GetRegistryByAddress(ctx, "1.2.3.4")
if err != nil {
t.Fatal(err)
}
if s.LastCycleID != 0 || s.Egress.Total != 0 || s.Ingress.Total != 0 || len(s.Egress.ByType) != 0 {
t.Fatalf("expected empty levels, got %+v", s)
}
}
// The counts follow the newest cycle, also after the queue row is deleted,
// and one grouped query serves several addresses without mixing them up.
func TestRegistryLevelsLatestCycleAndManyAddresses(t *testing.T) {
d, ctx := newTestDB(t)
if _, err := d.SubmitIPs(ctx, []string{"1.1.1.1", "2.2.2.2"}); err != nil {
t.Fatal(err)
}
addCheck(t, d, "1.1.1.1", SourceEgress, "https", "t1", false)
addCheck(t, d, "2.2.2.2", SourceEgress, "https", "t1", true)
addCheck(t, d, "2.2.2.2", InboundSource(1), "tcp-22", "2.2.2.2", true)
// New cycle for 1.1.1.1: delete and submit again.
ip, err := d.GetIPByAddress(ctx, "1.1.1.1")
if err != nil {
t.Fatal(err)
}
if err := d.DeleteIP(ctx, ip.ID); err != nil {
t.Fatal(err)
}
if _, err := d.SubmitIPs(ctx, []string{"1.1.1.1"}); err != nil {
t.Fatal(err)
}
addCheck(t, d, "1.1.1.1", SourceEgress, "icmp", "t2", true)
addCheck(t, d, "1.1.1.1", SourceEgress, "https", "t2", true)
all, err := d.ListRegistry(ctx)
if err != nil {
t.Fatal(err)
}
got := map[string]RegistrySummary{}
for _, s := range all {
got[s.IPAddress] = s
}
a := got["1.1.1.1"]
if a.LastCycleID != 2 || a.Egress.Total != 2 || a.Egress.OK != 2 || a.Ingress.Total != 0 {
t.Errorf("1.1.1.1: %+v", a)
}
b := got["2.2.2.2"]
if b.LastCycleID != 1 || b.Egress.Total != 1 || b.Ingress.Total != 1 || b.Ingress.OK != 1 {
t.Errorf("2.2.2.2: %+v", b)
}
// Page and single lookups agree with the full list.
page, _, err := d.ListRegistryPage(ctx, RegistryFilter{}, 10, 0)
if err != nil {
t.Fatal(err)
}
for _, s := range page {
if !reflect.DeepEqual(s.Egress, got[s.IPAddress].Egress) || !reflect.DeepEqual(s.Ingress, got[s.IPAddress].Ingress) {
t.Errorf("page differs from list for %s", s.IPAddress)
}
}
// Delete the queue row of 2.2.2.2: the counts stay (history outlives it).
ip2, _ := d.GetIPByAddress(ctx, "2.2.2.2")
if err := d.DeleteIP(ctx, ip2.ID); err != nil {
t.Fatal(err)
}
s, err := d.GetRegistryByAddress(ctx, "2.2.2.2")
if err != nil {
t.Fatal(err)
}
if s.Egress.Total != 1 || s.Ingress.Total != 1 {
t.Errorf("after delete: %+v", s)
}
}
+53
View File
@@ -5,6 +5,7 @@ import (
"fmt"
"reflect"
"sort"
"strings"
"testing"
"time"
)
@@ -350,6 +351,33 @@ func TestMigration0009Indexes(t *testing.T) {
}
}
// TestRegistryLevelsQueryUsesIndex guards against a full scan of checks: both
// the per-address MAX(cycle_id) and the join back must go through
// idx_checks_registry_cycle.
func TestRegistryLevelsQueryUsesIndex(t *testing.T) {
d, ctx := newTestDB(t)
rows, err := d.QueryContext(ctx, "EXPLAIN QUERY PLAN "+registryLevelsQuery(3), 1, 2, 3)
if err != nil {
t.Fatal(err)
}
defer rows.Close()
var plan string
for rows.Next() {
var id, parent, unused int
var detail string
if err := rows.Scan(&id, &parent, &unused, &detail); err != nil {
t.Fatal(err)
}
plan += detail + "\n"
if strings.HasPrefix(detail, "SCAN") && strings.Contains(detail, "checks") {
t.Errorf("full scan of checks in plan:\n%s", plan)
}
}
if !strings.Contains(plan, "idx_checks_registry_cycle") {
t.Errorf("expected idx_checks_registry_cycle in plan:\n%s", plan)
}
}
// TestScaleSmoke6440 pushes a realistic project size through the hot paths
// with a loose time bound: the point is the absence of O(n^2) / N+1 work, not
// a benchmark.
@@ -369,11 +397,36 @@ func TestScaleSmoke6440(t *testing.T) {
}
submitDur := time.Since(start)
// A realistic check set for the addresses of the registry page below:
// 12 egress and 18 ingress checks each.
for _, addr := range addrs[3000:3100] {
ip, err := d.GetIPByAddress(ctx, addr)
if err != nil {
t.Fatalf("get %s: %v", addr, err)
}
for i := 0; i < 30; i++ {
source, checkType := SourceEgress, []string{"https", "icmp"}[i%2]
if i >= 12 {
source, checkType = InboundSource(i%3+1), fmt.Sprintf("tcp-%d", 20+i)
}
if err := d.UpsertCheck(ctx, Check{
IPID: ip.ID, IPAddress: addr, AttemptNumber: ip.AttemptNumber,
Source: source, CheckType: checkType, Target: fmt.Sprintf("t%d", i),
Success: i%5 != 0, CheckedAt: Now(),
}); err != nil {
t.Fatalf("upsert check: %v", err)
}
}
}
start = time.Now()
page, total, err := d.ListRegistryPage(ctx, RegistryFilter{}, 100, 3000)
if err != nil || total != 6440 || len(page) != 100 {
t.Fatalf("registry page: total=%d len=%d err=%v", total, len(page), err)
}
if e, i := page[0].Egress.Total, page[0].Ingress.Total; e != 12 || i != 18 {
t.Fatalf("levels of the first page row: egress=%d ingress=%d", e, i)
}
if _, total, err = d.ListRegistryPage(ctx, RegistryFilter{LastResult: ResultPass}, 100, 0); err != nil || total != 0 {
t.Fatalf("registry last_result filter: total=%d err=%v", total, err)
}
@@ -0,0 +1,246 @@
package db
import (
"context"
"database/sql"
"path/filepath"
"testing"
"time"
)
// checkingIP submits one address and moves it to the checking state.
func checkingIP(t *testing.T, d *DB, addr string) *IPQueueItem {
t.Helper()
ctx := context.Background()
if _, err := d.SubmitIPs(ctx, []string{addr}); err != nil {
t.Fatal(err)
}
ip, err := d.GetIPByAddress(ctx, addr)
if err != nil {
t.Fatal(err)
}
if err := d.SetChecking(ctx, ip.ID, time.Minute); err != nil {
t.Fatal(err)
}
ip, err = d.GetIP(ctx, ip.ID)
if err != nil {
t.Fatal(err)
}
return ip
}
func checkOf(ip *IPQueueItem, ct string, ok bool) Check {
return Check{IPID: ip.ID, IPAddress: ip.IPAddress, AttemptNumber: ip.AttemptNumber,
Source: InboundSource(1), CheckType: ct, Target: ip.IPAddress, Success: ok, CheckedAt: Now()}
}
func storedSuccess(t *testing.T, d *DB, ip *IPQueueItem, ct string) (success bool, found bool) {
t.Helper()
err := d.QueryRowContext(context.Background(),
`SELECT success FROM checks WHERE ip_id=? AND check_type=?`, ip.ID, ct).Scan(&success)
if err == sql.ErrNoRows {
return false, false
}
if err != nil {
t.Fatal(err)
}
return success, true
}
// Results are accepted while the address is being checked (also repeated ones,
// idempotently) and refused from the moment the verdict is being computed.
func TestUpsertCheckIfOpenFreezesAtVerdict(t *testing.T) {
d, ctx := newTestDB(t)
ip := checkingIP(t, d, "1.2.3.4")
for i := 0; i < 2; i++ {
if ok, err := d.UpsertCheckIfOpen(ctx, checkOf(ip, "tcp-22", true)); err != nil || !ok {
t.Fatalf("write %d while checking: ok=%v err=%v", i, ok, err)
}
}
var n int
if err := d.QueryRowContext(ctx, `SELECT COUNT(*) FROM checks WHERE ip_id=?`, ip.ID).Scan(&n); err != nil || n != 1 {
t.Fatalf("expected one row after repeated write, got %d err=%v", n, err)
}
// An older attempt's result does not touch the current attempt.
old := checkOf(ip, "icmp", true)
old.AttemptNumber = ip.AttemptNumber - 1
if ok, err := d.UpsertCheckIfOpen(ctx, old); err != nil || ok {
t.Fatalf("stale attempt must be dropped: ok=%v err=%v", ok, err)
}
if _, found := storedSuccess(t, d, ip, "icmp"); found {
t.Fatal("stale attempt wrote a row")
}
if err := d.SetAggregating(ctx, ip.ID); err != nil {
t.Fatal(err)
}
// Neither a new check nor a change of an existing one gets through.
if ok, err := d.UpsertCheckIfOpen(ctx, checkOf(ip, "ssh", false)); err != nil || ok {
t.Fatalf("new check while aggregating must be dropped: ok=%v err=%v", ok, err)
}
if ok, err := d.UpsertCheckIfOpen(ctx, checkOf(ip, "tcp-22", false)); err != nil || ok {
t.Fatalf("overwrite while aggregating must be dropped: ok=%v err=%v", ok, err)
}
if err := d.FinishIP(ctx, ip.ID, ResultPass); err != nil {
t.Fatal(err)
}
if ok, err := d.UpsertCheckIfOpen(ctx, checkOf(ip, "tcp-22", false)); err != nil || ok {
t.Fatalf("overwrite after the verdict must be dropped: ok=%v err=%v", ok, err)
}
if s, found := storedSuccess(t, d, ip, "tcp-22"); !found || !s {
t.Fatalf("stored tcp-22 changed after the verdict: found=%v success=%v", found, s)
}
if _, found := storedSuccess(t, d, ip, "ssh"); found {
t.Fatal("a check written after the verdict")
}
}
// recorded_at is the server's own write time and moves on every accepted write.
func TestUpsertCheckIfOpenSetsRecordedAt(t *testing.T) {
d, ctx := newTestDB(t)
ip := checkingIP(t, d, "1.2.3.4")
if _, err := d.UpsertCheckIfOpen(ctx, checkOf(ip, "icmp", true)); err != nil {
t.Fatal(err)
}
var first string
if err := d.QueryRowContext(ctx, `SELECT recorded_at FROM checks WHERE ip_id=?`, ip.ID).Scan(&first); err != nil || first == "" {
t.Fatalf("recorded_at not set: %q %v", first, err)
}
time.Sleep(5 * time.Millisecond)
if _, err := d.UpsertCheckIfOpen(ctx, checkOf(ip, "icmp", false)); err != nil {
t.Fatal(err)
}
var second, created string
if err := d.QueryRowContext(ctx, `SELECT recorded_at, created_at FROM checks WHERE ip_id=?`, ip.ID).Scan(&second, &created); err != nil {
t.Fatal(err)
}
if second <= first || created != first {
t.Fatalf("recorded_at must advance, created_at stay: first=%s second=%s created=%s", first, second, created)
}
}
// A prober site is handed an address until it reports it complete in the
// current attempt; other sites are unaffected; a new attempt hands it out again.
func TestListCheckingForSite(t *testing.T) {
d, ctx := newTestDB(t)
a := checkingIP(t, d, "1.1.1.1")
b := checkingIP(t, d, "2.2.2.2")
names := func(items []IPQueueItem) string {
s := ""
for _, it := range items {
s += it.IPAddress + " "
}
return s
}
for site := 1; site <= 2; site++ {
items, err := d.ListCheckingForSite(ctx, site)
if err != nil || len(items) != 2 {
t.Fatalf("site %d before any report: %q err=%v", site, names(items), err)
}
}
if err := d.SetSiteComplete(ctx, a.ID, 1); err != nil {
t.Fatal(err)
}
if items, _ := d.ListCheckingForSite(ctx, 1); names(items) != "2.2.2.2 " {
t.Fatalf("site 1 after completing 1.1.1.1: %q", names(items))
}
if items, _ := d.ListCheckingForSite(ctx, 2); len(items) != 2 {
t.Fatalf("site 2 must still get both: %q", names(items))
}
// A retry starts a new attempt: site 1 has to probe the address again.
if err := d.RequeueOrFail(ctx, a.ID, "", 3); err != nil {
t.Fatal(err)
}
if err := d.SetChecking(ctx, a.ID, time.Minute); err != nil {
t.Fatal(err)
}
if items, _ := d.ListCheckingForSite(ctx, 1); len(items) != 2 {
t.Fatalf("site 1 after a new attempt: %q", names(items))
}
_ = b
}
// The aggregation window counts from the start of checking; a retry clears it.
func TestCheckingStartedAt(t *testing.T) {
d, ctx := newTestDB(t)
ip := checkingIP(t, d, "1.2.3.4")
if ip.CheckingStartedAt == nil || time.Since(*ip.CheckingStartedAt) > time.Minute {
t.Fatalf("checking_started_at not set: %v", ip.CheckingStartedAt)
}
if err := d.RequeueOrFail(ctx, ip.ID, "", 3); err != nil {
t.Fatal(err)
}
ip, err := d.GetIP(ctx, ip.ID)
if err != nil {
t.Fatal(err)
}
if ip.CheckingStartedAt != nil {
t.Fatalf("checking_started_at must be cleared on requeue, got %v", ip.CheckingStartedAt)
}
}
// Migration 0010 flags the rows of an existing database that were written
// after their address's verdict, and leaves the others alone.
func TestMigration0010MarksRowsAfterVerdict(t *testing.T) {
ctx := context.Background()
path := filepath.Join(t.TempDir(), "old.db")
raw, err := sql.Open("sqlite", path)
if err != nil {
t.Fatal(err)
}
raw.SetMaxOpenConns(1)
for _, m := range migrations {
if m.version > 9 {
break
}
if _, err := raw.ExecContext(ctx, m.sql); err != nil {
t.Fatalf("migration %d: %v", m.version, err)
}
}
if _, err := raw.ExecContext(ctx, `PRAGMA user_version=9`); err != nil {
t.Fatal(err)
}
for _, q := range []string{
`INSERT INTO ip_registry (id, ip_address, first_seen_at, last_seen_at, next_cycle, created_at, updated_at)
VALUES (1, '1.2.3.4', '2026-10-02T13:00:00Z', '2026-10-02T13:00:00Z', 2, '2026-10-02T13:00:00Z', '2026-10-02T13:00:00Z')`,
`INSERT INTO ip_queue (id, ip_address, sequence, state, overall_result, aggregated_at, registry_id, cycle_id, created_at, updated_at)
VALUES (1, '1.2.3.4', 1, 'done', 'pass', '2026-10-02T13:48:45.659Z', 1, 1, '2026-10-02T13:00:00Z', '2026-10-02T13:00:00Z')`,
// before the verdict, and (with a longer fraction) after it
`INSERT INTO checks (registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id, source, check_type, target, success, checked_at, created_at)
VALUES (1, 1, 1, '1.2.3.4', 1, '', 'inbound-site-1', 'icmp', '1.2.3.4', 1, '2026-10-02T13:48:40.100000000Z', '2026-10-02T13:48:40.2Z')`,
`INSERT INTO checks (registry_id, cycle_id, ip_id, ip_address, attempt_number, validator_id, source, check_type, target, success, checked_at, created_at)
VALUES (1, 1, 1, '1.2.3.4', 1, '', 'inbound-site-1', 'ssh', '1.2.3.4', 0, '2026-10-02T13:48:45.730314288Z', '2026-10-02T13:48:40.3Z')`,
} {
if _, err := raw.ExecContext(ctx, q); err != nil {
t.Fatalf("seed: %v\n%s", err, q)
}
}
raw.Close()
d, err := Open(ctx, path)
if err != nil {
t.Fatalf("open (runs migration 10): %v", err)
}
defer d.Close()
flag := func(ct string) (after int, recorded, created string) {
if err := d.QueryRowContext(ctx, `SELECT after_verdict, recorded_at, created_at FROM checks WHERE check_type=?`, ct).Scan(&after, &recorded, &created); err != nil {
t.Fatal(err)
}
return
}
if a, rec, cr := flag("icmp"); a != 0 || rec != cr {
t.Errorf("icmp: after_verdict=%d recorded=%s created=%s", a, rec, cr)
}
if a, rec, cr := flag("ssh"); a != 1 || rec != cr {
t.Errorf("ssh: after_verdict=%d recorded=%s created=%s", a, rec, cr)
}
var ver int
if err := d.QueryRowContext(ctx, `PRAGMA user_version`).Scan(&ver); err != nil || ver != 10 {
t.Errorf("user_version=%d err=%v", ver, err)
}
}