added FIP Cooldown before start

This commit is contained in:
ayurishchev committed 2026-08-24 10:29:08 +03:00
1 parent 291eea3eb8
commit f22ad569b6
36 files changed
+1182 -25

No files matched your search

+18
View File
@@ -33,6 +33,9 @@ func (d *DB) BootstrapFromConfig(ctx context.Context, cfg *config.ControlAPI) er
if err := d.bootstrapCheckTypes(ctx, cfg.CheckTypes); err != nil {
return fmt.Errorf("bootstrap check types: %w", err)
}
if err := d.bootstrapSettings(ctx, cfg.Orchestrator.FIPSettleSeconds); err != nil {
return fmt.Errorf("bootstrap settings: %w", err)
}
if err := d.SeedQueue(ctx, cfg.IPAddresses); err != nil {
return fmt.Errorf("seed ip queue: %w", err)
}
@@ -102,3 +105,18 @@ func (d *DB) bootstrapCheckTypes(ctx context.Context, checkTypes []config.CheckT
}
return nil
}
func (d *DB) bootstrapSettings(ctx context.Context, fipSettleSeconds int) error {
var count int
if err := d.QueryRowContext(ctx, `SELECT COUNT(*) FROM settings`).Scan(&count); err != nil {
return err
}
if count > 0 {
return nil
}
now := timeToDB(Now())
_, err := d.ExecContext(ctx, `
INSERT INTO settings (id, fip_settle_seconds, created_at, updated_at) VALUES (1, ?, ?, ?)
`, fipSettleSeconds, now, now)
return err
}
+4
View File
@@ -19,6 +19,9 @@ var initSchema string
//go:embed migrations/0002_dynamic_config.sql
var dynamicConfigSchema string
//go:embed migrations/0003_fip_settle_delay.sql
var fipSettleDelaySchema 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
@@ -29,6 +32,7 @@ var migrations = []struct {
}{
{1, initSchema},
{2, dynamicConfigSchema},
{3, fipSettleDelaySchema},
}
type DB struct {
@@ -0,0 +1,15 @@
-- Support for a configurable pause between FIP association and the start
-- of self-check (see docs/USAGE.md).
ALTER TABLE ip_queue ADD COLUMN fip_associated_at TIMESTAMP;
-- Singleton row for the one dynamic scalar setting we have so far.
-- Deliberately a typed single-row table, not a generic key-value settings
-- table: overengineering for exactly one field. If more scalar knobs show
-- up, add columns here first.
CREATE TABLE settings (
id INTEGER PRIMARY KEY CHECK (id = 1),
fip_settle_seconds INTEGER NOT NULL DEFAULT 0,
created_at TIMESTAMP NOT NULL,
updated_at TIMESTAMP NOT NULL
);
+9
View File
@@ -84,6 +84,7 @@ type IPQueueItem struct {
Site3Complete bool
OverallResult string
AssignedAt *time.Time
FIPAssociatedAt *time.Time
AggregatedAt *time.Time
FIPReleasedAt *time.Time
CreatedAt time.Time
@@ -162,3 +163,11 @@ type DeleteIPsResult struct {
Deleted []string
NotFound []string
}
// Settings is the singleton row of orchestrator-wide scalar knobs that are
// admin-configurable at runtime (see queries_settings.go).
type Settings struct {
FIPSettleSeconds int
CreatedAt time.Time
UpdatedAt time.Time
}
+10 -7
View File
@@ -110,9 +110,9 @@ func (d *DB) ClaimNextQueued(ctx context.Context, validatorID string, leaseTTL t
func (d *DB) SetFIPAssociated(ctx context.Context, ipID int64, fipID string, leaseTTL time.Duration) error {
now := Now()
_, err := d.ExecContext(ctx, `
UPDATE ip_queue SET state=?, fip_id=?, lease_expires_at=?, updated_at=?
UPDATE ip_queue SET state=?, fip_id=?, fip_associated_at=?, lease_expires_at=?, updated_at=?
WHERE id=?
`, IPAwaitingSelfCheck, fipID, timeToDB(now.Add(leaseTTL)), timeToDB(now), ipID)
`, IPAwaitingSelfCheck, fipID, timeToDB(now), timeToDB(now.Add(leaseTTL)), timeToDB(now), ipID)
return err
}
@@ -197,7 +197,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,
lease_expires_at=NULL, egress_complete=0, site1_complete=0, site2_complete=0, site3_complete=0,
overall_result='', assigned_at=NULL, updated_at=?
overall_result='', assigned_at=NULL, fip_associated_at=NULL, updated_at=?
WHERE id=?
`, nextState, retryCount, now, ipID)
} else {
@@ -283,7 +283,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, lease_expires_at=NULL, egress_complete=0,
site1_complete=0, site2_complete=0, site3_complete=0, overall_result='',
assigned_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=?
assigned_at=NULL, fip_associated_at=NULL, aggregated_at=NULL, fip_released_at=NULL, updated_at=?
WHERE ip_address=?
`, IPQueued, seq, now, addr); err != nil {
return result, fmt.Errorf("requeue %s: %w", addr, err)
@@ -480,7 +480,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, site1_complete, site2_complete, site3_complete, overall_result,
assigned_at, aggregated_at, fip_released_at, created_at, updated_at
assigned_at, fip_associated_at, aggregated_at, fip_released_at, created_at, updated_at
FROM ip_queue
`
@@ -499,13 +499,13 @@ func scanIPQueueItems(rows *sql.Rows) ([]IPQueueItem, error) {
func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) {
var item IPQueueItem
var owner sql.NullString
var leaseExpires, assignedAt, aggregatedAt, fipReleasedAt sql.NullString
var leaseExpires, assignedAt, fipAssociatedAt, aggregatedAt, fipReleasedAt sql.NullString
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.Site1Complete, &item.Site2Complete, &item.Site3Complete,
&item.OverallResult, &assignedAt, &aggregatedAt, &fipReleasedAt, &createdAt, &updatedAt,
&item.OverallResult, &assignedAt, &fipAssociatedAt, &aggregatedAt, &fipReleasedAt, &createdAt, &updatedAt,
); err != nil {
return nil, err
}
@@ -519,6 +519,9 @@ func scanIPQueueItem(row rowScanner) (*IPQueueItem, error) {
if item.AssignedAt, err = nullStringToTimePtr(assignedAt); err != nil {
return nil, err
}
if item.FIPAssociatedAt, err = nullStringToTimePtr(fipAssociatedAt); err != nil {
return nil, err
}
if item.AggregatedAt, err = nullStringToTimePtr(aggregatedAt); err != nil {
return nil, err
}
+71
View File
@@ -0,0 +1,71 @@
package db
import (
"testing"
"time"
)
func TestSetFIPAssociatedStampsTimestamp(t *testing.T) {
d, ctx := newTestDB(t)
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatalf("seed queue: %v", err)
}
if err := d.AdminCreateValidator(ctx, "validator-1", "port-1"); err != nil {
t.Fatalf("create validator: %v", err)
}
claimed, err := d.ClaimNextQueued(ctx, "validator-1", time.Minute)
if err != nil || claimed == nil {
t.Fatalf("claim: item=%+v err=%v", claimed, err)
}
if claimed.FIPAssociatedAt != nil {
t.Fatalf("expected FIPAssociatedAt nil before association, got %v", claimed.FIPAssociatedAt)
}
before := Now()
if err := d.SetFIPAssociated(ctx, claimed.ID, "fip-1", time.Minute); err != nil {
t.Fatalf("set fip associated: %v", err)
}
item, err := d.GetIP(ctx, claimed.ID)
if err != nil {
t.Fatalf("get ip: %v", err)
}
if item.FIPAssociatedAt == nil {
t.Fatalf("expected FIPAssociatedAt to be set")
}
if item.FIPAssociatedAt.Before(before) {
t.Fatalf("expected FIPAssociatedAt >= %v, got %v", before, *item.FIPAssociatedAt)
}
}
func TestRequeueClearsFIPAssociatedAt(t *testing.T) {
d, ctx := newTestDB(t)
if err := d.SeedQueue(ctx, []string{"1.2.3.4"}); err != nil {
t.Fatalf("seed queue: %v", err)
}
if err := d.AdminCreateValidator(ctx, "validator-1", "port-1"); err != nil {
t.Fatalf("create validator: %v", err)
}
claimed, err := d.ClaimNextQueued(ctx, "validator-1", time.Minute)
if err != nil || claimed == nil {
t.Fatalf("claim: item=%+v err=%v", claimed, err)
}
if err := d.SetFIPAssociated(ctx, claimed.ID, "fip-1", time.Minute); err != nil {
t.Fatalf("set fip associated: %v", err)
}
if err := d.RequeueOrFail(ctx, claimed.ID, "validator-1", 3); err != nil {
t.Fatalf("requeue or fail: %v", err)
}
item, err := d.GetIP(ctx, claimed.ID)
if err != nil {
t.Fatalf("get ip after requeue: %v", err)
}
if item.State != IPQueued {
t.Fatalf("expected requeued, got state=%s", item.State)
}
if item.FIPAssociatedAt != nil {
t.Fatalf("expected FIPAssociatedAt cleared on requeue, got %v", item.FIPAssociatedAt)
}
}
+44
View File
@@ -0,0 +1,44 @@
package db
import (
"context"
"fmt"
)
// GetSettings returns the singleton settings row. BootstrapFromConfig
// guarantees it exists before any other code path can observe it, so
// sql.ErrNoRows here would indicate a bootstrap bug, not a normal
// condition.
func (d *DB) GetSettings(ctx context.Context) (Settings, error) {
var s Settings
var createdAt, updatedAt string
err := d.QueryRowContext(ctx, `
SELECT fip_settle_seconds, created_at, updated_at FROM settings WHERE id=1
`).Scan(&s.FIPSettleSeconds, &createdAt, &updatedAt)
if err != nil {
return Settings{}, err
}
if s.CreatedAt, err = dbToTime(createdAt); err != nil {
return Settings{}, err
}
if s.UpdatedAt, err = dbToTime(updatedAt); err != nil {
return Settings{}, err
}
return s, nil
}
// SetFIPSettleSeconds persists a new fip_settle_seconds value. Only
// enforces seconds >= 0 — the cross-field rule against
// lease_ttl_seconds/self_check_timeout_seconds lives in
// orchestrator.Orchestrator.SetFIPSettleSeconds, which has access to the
// static orchestrator config this package doesn't.
func (d *DB) SetFIPSettleSeconds(ctx context.Context, seconds int) error {
if seconds < 0 {
return fmt.Errorf("fip_settle_seconds must be >= 0: %w", ErrValidation)
}
now := timeToDB(Now())
_, err := d.ExecContext(ctx, `
UPDATE settings SET fip_settle_seconds=?, updated_at=? WHERE id=1
`, seconds, now)
return err
}
+77
View File
@@ -0,0 +1,77 @@
package db
import (
"errors"
"testing"
"cloudipvalidator/internal/config"
)
func TestBootstrapSettingsSeedsOnceFromConfig(t *testing.T) {
d, ctx := newTestDB(t)
cfg := &config.ControlAPI{Orchestrator: config.OrchestratorConfig{FIPSettleSeconds: 30}}
if err := d.BootstrapFromConfig(ctx, cfg); err != nil {
t.Fatalf("first bootstrap: %v", err)
}
settings, err := d.GetSettings(ctx)
if err != nil {
t.Fatalf("get settings: %v", err)
}
if settings.FIPSettleSeconds != 30 {
t.Fatalf("expected seeded value 30, got %d", settings.FIPSettleSeconds)
}
// A second bootstrap with a different YAML value must not overwrite the
// now-non-empty settings table — same "DB is source of truth once
// seeded" semantics as validators/sites/targets/check_types.
cfg2 := &config.ControlAPI{Orchestrator: config.OrchestratorConfig{FIPSettleSeconds: 99}}
if err := d.BootstrapFromConfig(ctx, cfg2); err != nil {
t.Fatalf("second bootstrap: %v", err)
}
settings, err = d.GetSettings(ctx)
if err != nil {
t.Fatalf("get settings after second bootstrap: %v", err)
}
if settings.FIPSettleSeconds != 30 {
t.Fatalf("expected YAML to be ignored on non-empty settings table, got %d", settings.FIPSettleSeconds)
}
}
func TestSetFIPSettleSecondsRejectsNegative(t *testing.T) {
d, ctx := newTestDB(t)
if err := d.BootstrapFromConfig(ctx, &config.ControlAPI{}); err != nil {
t.Fatalf("bootstrap: %v", err)
}
if err := d.SetFIPSettleSeconds(ctx, -1); !errors.Is(err, ErrValidation) {
t.Fatalf("expected ErrValidation for negative seconds, got %v", err)
}
}
func TestGetSetFIPSettleSecondsRoundTrip(t *testing.T) {
d, ctx := newTestDB(t)
if err := d.BootstrapFromConfig(ctx, &config.ControlAPI{}); err != nil {
t.Fatalf("bootstrap: %v", err)
}
before, err := d.GetSettings(ctx)
if err != nil {
t.Fatalf("get settings: %v", err)
}
if before.FIPSettleSeconds != 0 {
t.Fatalf("expected default 0, got %d", before.FIPSettleSeconds)
}
if err := d.SetFIPSettleSeconds(ctx, 45); err != nil {
t.Fatalf("set fip settle seconds: %v", err)
}
after, err := d.GetSettings(ctx)
if err != nil {
t.Fatalf("get settings after set: %v", err)
}
if after.FIPSettleSeconds != 45 {
t.Fatalf("expected 45, got %d", after.FIPSettleSeconds)
}
if after.UpdatedAt.Before(before.UpdatedAt) {
t.Fatalf("expected updated_at not to go backwards, before=%v after=%v", before.UpdatedAt, after.UpdatedAt)
}
}