feat(revisions): add pruning estimate and cleanup endpoints
CI / changes (push) Successful in 8s
CI / commitlint (push) Has been skipped
CI / openapi (push) Successful in 25s
CI / web (push) Successful in 32s
CI / go (push) Successful in 57s
CI / bird2 (push) Successful in 15s
CI / release (push) Successful in 3m18s
CI / changes (push) Successful in 8s
CI / commitlint (push) Has been skipped
CI / openapi (push) Successful in 25s
CI / web (push) Successful in 32s
CI / go (push) Successful in 57s
CI / bird2 (push) Successful in 15s
CI / release (push) Successful in 3m18s
Implemented new endpoints for estimating and pruning revisions, including detailed schemas for requests and responses. The `RevisionPruneEstimate` and `RevisionPruneResult` components were added to the OpenAPI documentation, enhancing the API's functionality for managing revision retention. Updated the backend to support these operations and integrated them into the tenant settings UI for improved user interaction.
This commit is contained in:
@@ -970,43 +970,6 @@ func (p *Postgres) RevisionDiff(tenantID, aID, bID string) (map[string]any, erro
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (p *Postgres) PruneRevisionsBefore(tenantID string, cutoff time.Time) (int, error) {
|
||||
ctx := context.Background()
|
||||
total := 0
|
||||
const batchSize = 50
|
||||
for {
|
||||
cmd, err := p.pool.Exec(ctx, `
|
||||
DELETE FROM config_revision AS cr
|
||||
WHERE cr.id IN (
|
||||
SELECT id FROM config_revision
|
||||
WHERE tenant_id = $1
|
||||
AND created_at < $2
|
||||
AND id <> (
|
||||
SELECT id FROM config_revision
|
||||
WHERE tenant_id = $1
|
||||
ORDER BY created_at DESC
|
||||
LIMIT 1
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM bgp_speaker AS sp
|
||||
WHERE sp.tenant_id = $1
|
||||
AND (sp.last_applied_revision_id = config_revision.id OR sp.published_revision_id = config_revision.id)
|
||||
)
|
||||
ORDER BY created_at ASC
|
||||
LIMIT $3
|
||||
)`, tenantID, cutoff.UTC(), batchSize)
|
||||
if err != nil {
|
||||
return total, err
|
||||
}
|
||||
n := int(cmd.RowsAffected())
|
||||
total += n
|
||||
if n < batchSize {
|
||||
break
|
||||
}
|
||||
}
|
||||
return total, nil
|
||||
}
|
||||
|
||||
func (p *Postgres) SetLastAppliedRevision(tenantID, speakerID, revisionID string) error {
|
||||
ctx := context.Background()
|
||||
tag, err := p.pool.Exec(ctx, `
|
||||
|
||||
@@ -53,21 +53,8 @@ func (p *Postgres) ensurePrefixSnapshot(ctx context.Context, db execQuerier, con
|
||||
return "", nil
|
||||
}
|
||||
hash := normalizeSnapshotHash(contentHash)
|
||||
var existing string
|
||||
err := db.QueryRow(ctx, `SELECT id::text FROM prefix_snapshot WHERE content_hash = $1`, hash).Scan(&existing)
|
||||
if err == nil && existing != "" {
|
||||
return existing, nil
|
||||
}
|
||||
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
|
||||
return "", err
|
||||
}
|
||||
snapID := uuid.NewString()
|
||||
if _, err := db.Exec(ctx, `
|
||||
INSERT INTO prefix_snapshot (id, content_hash) VALUES ($1::uuid, $2)
|
||||
ON CONFLICT (content_hash) DO NOTHING`, snapID, hash); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := db.QueryRow(ctx, `SELECT id::text FROM prefix_snapshot WHERE content_hash = $1`, hash).Scan(&snapID); err != nil {
|
||||
snapID, err := p.resolvePrefixSnapshotID(ctx, db, hash)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
var rowCount int
|
||||
@@ -94,6 +81,27 @@ func (p *Postgres) ensurePrefixSnapshot(ctx context.Context, db execQuerier, con
|
||||
return snapID, nil
|
||||
}
|
||||
|
||||
func (p *Postgres) resolvePrefixSnapshotID(ctx context.Context, db execQuerier, hash string) (string, error) {
|
||||
var existing string
|
||||
err := db.QueryRow(ctx, `SELECT id::text FROM prefix_snapshot WHERE content_hash = $1`, hash).Scan(&existing)
|
||||
if err == nil && existing != "" {
|
||||
return existing, nil
|
||||
}
|
||||
if err != nil && !errors.Is(err, pgx.ErrNoRows) {
|
||||
return "", err
|
||||
}
|
||||
snapID := uuid.NewString()
|
||||
if _, err := db.Exec(ctx, `
|
||||
INSERT INTO prefix_snapshot (id, content_hash) VALUES ($1::uuid, $2)
|
||||
ON CONFLICT (content_hash) DO NOTHING`, snapID, hash); err != nil {
|
||||
return "", err
|
||||
}
|
||||
if err := db.QueryRow(ctx, `SELECT id::text FROM prefix_snapshot WHERE content_hash = $1`, hash).Scan(&snapID); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return snapID, nil
|
||||
}
|
||||
|
||||
func (p *Postgres) listSnapshotPrefixes(ctx context.Context, snapshotID, cursor string, limit int) ([]store.PrefixRow, string, bool) {
|
||||
afterOrd, off, useOffset := store.ParsePrefixPageCursor(cursor)
|
||||
var rows pgx.Rows
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
package repository
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"evobgp/internal/db"
|
||||
"evobgp/internal/store"
|
||||
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
func TestEnsurePrefixSnapshotFillsEmptyExistingIntegration(t *testing.T) {
|
||||
dsn := os.Getenv("EVOBGP_TEST_DATABASE_URL")
|
||||
if dsn == "" {
|
||||
t.Skip("EVOBGP_TEST_DATABASE_URL not set")
|
||||
}
|
||||
ctx := context.Background()
|
||||
pool, err := db.OpenPostgresPool(ctx, dsn)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer pool.Close()
|
||||
pg, err := NewPostgres(ctx, pool, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
hash := normalizeSnapshotHash("sha256:empty-fill-" + uuid.NewString())
|
||||
snapID := uuid.NewString()
|
||||
if _, err := pool.Exec(ctx, `INSERT INTO prefix_snapshot (id, content_hash) VALUES ($1::uuid, $2)`, snapID, hash); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM prefix_snapshot WHERE id = $1::uuid`, snapID)
|
||||
})
|
||||
|
||||
tx, err := pool.Begin(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = tx.Rollback(ctx) }()
|
||||
|
||||
got, err := pg.ensurePrefixSnapshot(ctx, tx, "sha256:"+hash, []store.PrefixRow{
|
||||
{Prefix: "203.0.113.1/32", Source: "test"},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got != snapID {
|
||||
t.Fatalf("snap id: got %q want %q", got, snapID)
|
||||
}
|
||||
var rowCount int
|
||||
if err := tx.QueryRow(ctx, `SELECT COUNT(*)::int FROM prefix_snapshot_row WHERE snapshot_id = $1::uuid`, snapID).Scan(&rowCount); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if rowCount != 1 {
|
||||
t.Fatalf("row count: got %d", rowCount)
|
||||
}
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEnsurePrefixSnapshotIdempotentIntegration(t *testing.T) {
|
||||
dsn := os.Getenv("EVOBGP_TEST_DATABASE_URL")
|
||||
if dsn == "" {
|
||||
t.Skip("EVOBGP_TEST_DATABASE_URL not set")
|
||||
}
|
||||
ctx := context.Background()
|
||||
pool, err := db.OpenPostgresPool(ctx, dsn)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer pool.Close()
|
||||
pg, err := NewPostgres(ctx, pool, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
contentHash := "sha256:idempotent-" + uuid.NewString()
|
||||
prefixes := []store.PrefixRow{{Prefix: "198.51.100.0/24", Source: "test"}}
|
||||
|
||||
tx1, err := pool.Begin(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
id1, err := pg.ensurePrefixSnapshot(ctx, tx1, contentHash, prefixes)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := tx1.Commit(ctx); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
_, _ = pool.Exec(ctx, `DELETE FROM prefix_snapshot WHERE id = $1::uuid`, id1)
|
||||
})
|
||||
|
||||
tx2, err := pool.Begin(ctx)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = tx2.Rollback(ctx) }()
|
||||
id2, err := pg.ensurePrefixSnapshot(ctx, tx2, contentHash, prefixes)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if id1 != id2 {
|
||||
t.Fatalf("ids differ: %q vs %q", id1, id2)
|
||||
}
|
||||
var rowCount int
|
||||
if err := tx2.QueryRow(ctx, `SELECT COUNT(*)::int FROM prefix_snapshot_row WHERE snapshot_id = $1::uuid`, id1).Scan(&rowCount); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if rowCount != 1 {
|
||||
t.Fatalf("expected 1 row, got %d", rowCount)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,176 @@
|
||||
package repository
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"evobgp/internal/store"
|
||||
)
|
||||
|
||||
const revisionPruneBatchSize = 50
|
||||
|
||||
// sqlPrunableRevisionsWhere appends prunable revision predicates (tenant + cutoff params).
|
||||
func sqlPrunableRevisionsWhere(tenantParam, cutoffParam string) string {
|
||||
return `tenant_id = ` + tenantParam + `
|
||||
AND created_at < ` + cutoffParam + `
|
||||
AND id <> (
|
||||
SELECT id FROM config_revision
|
||||
WHERE tenant_id = ` + tenantParam + `
|
||||
ORDER BY created_at DESC
|
||||
LIMIT 1
|
||||
)
|
||||
AND NOT EXISTS (
|
||||
SELECT 1 FROM bgp_speaker AS sp
|
||||
WHERE sp.tenant_id = ` + tenantParam + `
|
||||
AND (sp.last_applied_revision_id = config_revision.id OR sp.published_revision_id = config_revision.id)
|
||||
)`
|
||||
}
|
||||
|
||||
func (p *Postgres) EstimateRevisionPrune(tenantID string, cutoff time.Time, retentionMinutes int) (store.RevisionPruneEstimate, error) {
|
||||
ctx := context.Background()
|
||||
out := store.RevisionPruneEstimate{
|
||||
RetentionMinutes: retentionMinutes,
|
||||
CutoffAt: cutoff.UTC(),
|
||||
}
|
||||
where := sqlPrunableRevisionsWhere("$1", "$2")
|
||||
var revBytes, previewBytes, rmpBytes, snapRowBytes int64
|
||||
err := p.pool.QueryRow(ctx, `
|
||||
WITH prunable AS (
|
||||
SELECT id, prefix_snapshot_id, content_hash, meta_json
|
||||
FROM config_revision
|
||||
WHERE `+where+`
|
||||
),
|
||||
freed_snaps AS (
|
||||
SELECT DISTINCT pr.prefix_snapshot_id AS snap_id
|
||||
FROM prunable pr
|
||||
WHERE pr.prefix_snapshot_id IS NOT NULL
|
||||
EXCEPT
|
||||
SELECT DISTINCT cr.prefix_snapshot_id
|
||||
FROM config_revision cr
|
||||
WHERE cr.prefix_snapshot_id IS NOT NULL
|
||||
AND cr.id NOT IN (SELECT id FROM prunable)
|
||||
)
|
||||
SELECT
|
||||
(SELECT COUNT(*)::int FROM prunable),
|
||||
COALESCE((SELECT SUM(octet_length(content_hash) + octet_length(meta_json::text))::bigint FROM prunable), 0),
|
||||
COALESCE((
|
||||
SELECT SUM(octet_length(value))::bigint
|
||||
FROM config_revision_preview p
|
||||
JOIN prunable pr ON pr.id = p.revision_id
|
||||
CROSS JOIN LATERAL jsonb_each_text(COALESCE(p.fragments, '{}'::jsonb))
|
||||
), 0),
|
||||
COALESCE((
|
||||
SELECT SUM(octet_length(rmp.prefix) + octet_length(COALESCE(rmp.source, '')))::bigint
|
||||
FROM revision_materialized_prefix rmp
|
||||
WHERE rmp.revision_id IN (SELECT id FROM prunable)
|
||||
), 0),
|
||||
(SELECT COUNT(*)::int FROM freed_snaps),
|
||||
COALESCE((SELECT COUNT(*)::int FROM prefix_snapshot_row psr WHERE psr.snapshot_id IN (SELECT snap_id FROM freed_snaps)), 0),
|
||||
COALESCE((
|
||||
SELECT SUM(octet_length(psr.prefix) + octet_length(COALESCE(psr.source, '')))::bigint
|
||||
FROM prefix_snapshot_row psr
|
||||
WHERE psr.snapshot_id IN (SELECT snap_id FROM freed_snaps)
|
||||
), 0)`,
|
||||
tenantID, cutoff.UTC()).Scan(
|
||||
&out.RevisionCount,
|
||||
&revBytes,
|
||||
&previewBytes,
|
||||
&rmpBytes,
|
||||
&out.OrphanSnapshotCount,
|
||||
&out.PrefixRowCount,
|
||||
&snapRowBytes,
|
||||
)
|
||||
if err != nil {
|
||||
return out, err
|
||||
}
|
||||
out.BytesEstimate = revBytes + previewBytes + rmpBytes + snapRowBytes
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (p *Postgres) PruneRevisionsWithStats(tenantID string, cutoff time.Time) (store.RevisionPruneResult, error) {
|
||||
ctx := context.Background()
|
||||
est, err := p.EstimateRevisionPrune(tenantID, cutoff, 0)
|
||||
if err != nil {
|
||||
return store.RevisionPruneResult{}, err
|
||||
}
|
||||
out := store.RevisionPruneResult{BytesEstimate: est.BytesEstimate}
|
||||
|
||||
where := sqlPrunableRevisionsWhere("$1", "$2")
|
||||
for {
|
||||
cmd, err := p.pool.Exec(ctx, `
|
||||
DELETE FROM config_revision AS cr
|
||||
WHERE cr.id IN (
|
||||
SELECT id FROM config_revision
|
||||
WHERE `+where+`
|
||||
ORDER BY created_at ASC
|
||||
LIMIT $3
|
||||
)`, tenantID, cutoff.UTC(), revisionPruneBatchSize)
|
||||
if err != nil {
|
||||
return out, err
|
||||
}
|
||||
n := int(cmd.RowsAffected())
|
||||
out.DeletedRevisions += n
|
||||
if n < revisionPruneBatchSize {
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
for {
|
||||
snaps, rows, err := p.pruneUnreferencedPrefixSnapshots(ctx, revisionPruneBatchSize)
|
||||
if err != nil {
|
||||
return out, err
|
||||
}
|
||||
out.DeletedPrefixSnapshots += snaps
|
||||
out.DeletedPrefixRows += rows
|
||||
if snaps < revisionPruneBatchSize {
|
||||
break
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
func (p *Postgres) PruneRevisionsBefore(tenantID string, cutoff time.Time) (int, error) {
|
||||
res, err := p.PruneRevisionsWithStats(tenantID, cutoff)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return res.DeletedRevisions, nil
|
||||
}
|
||||
|
||||
func (p *Postgres) pruneUnreferencedPrefixSnapshots(ctx context.Context, batchSize int) (deletedSnapshots int, deletedRows int, err error) {
|
||||
if !prefixSnapshotTableExists(ctx, p.pool) {
|
||||
return 0, 0, nil
|
||||
}
|
||||
var snapCount, rowCount int
|
||||
err = p.pool.QueryRow(ctx, `
|
||||
WITH doomed AS (
|
||||
SELECT ps.id FROM prefix_snapshot ps
|
||||
WHERE NOT EXISTS (
|
||||
SELECT 1 FROM config_revision cr WHERE cr.prefix_snapshot_id = ps.id
|
||||
)
|
||||
LIMIT $1
|
||||
)
|
||||
SELECT
|
||||
(SELECT COUNT(*)::int FROM doomed),
|
||||
(SELECT COUNT(*)::int FROM prefix_snapshot_row psr WHERE psr.snapshot_id IN (SELECT id FROM doomed))`,
|
||||
batchSize).Scan(&snapCount, &rowCount)
|
||||
if err != nil {
|
||||
return 0, 0, err
|
||||
}
|
||||
if snapCount == 0 {
|
||||
return 0, 0, nil
|
||||
}
|
||||
_, err = p.pool.Exec(ctx, `
|
||||
DELETE FROM prefix_snapshot
|
||||
WHERE id IN (
|
||||
SELECT ps.id FROM prefix_snapshot ps
|
||||
WHERE NOT EXISTS (
|
||||
SELECT 1 FROM config_revision cr WHERE cr.prefix_snapshot_id = ps.id
|
||||
)
|
||||
LIMIT $1
|
||||
)`, batchSize)
|
||||
if err != nil {
|
||||
return 0, 0, err
|
||||
}
|
||||
return snapCount, rowCount, nil
|
||||
}
|
||||
@@ -0,0 +1,83 @@
|
||||
package repository
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"evobgp/internal/db"
|
||||
"evobgp/internal/pipeline"
|
||||
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
func TestPostgresPruneUnreferencedPrefixSnapshotsIntegration(t *testing.T) {
|
||||
dsn := os.Getenv("EVOBGP_TEST_DATABASE_URL")
|
||||
if dsn == "" {
|
||||
t.Skip("EVOBGP_TEST_DATABASE_URL not set")
|
||||
}
|
||||
ctx := context.Background()
|
||||
pool, err := db.OpenPostgresPool(ctx, dsn)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer pool.Close()
|
||||
pg, err := NewPostgres(ctx, pool, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
snapID := uuid.NewString()
|
||||
hash := normalizeSnapshotHash("sha256:test-orphan-" + snapID)
|
||||
if _, err := pool.Exec(ctx, `INSERT INTO prefix_snapshot (id, content_hash) VALUES ($1::uuid, $2)`, snapID, hash); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := pool.Exec(ctx, `
|
||||
INSERT INTO prefix_snapshot_row (snapshot_id, ord, prefix, source)
|
||||
VALUES ($1::uuid, 0, '203.0.113.0/24', 'test')`, snapID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
snaps, rows, err := pg.pruneUnreferencedPrefixSnapshots(ctx, 50)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if snaps < 1 || rows < 1 {
|
||||
t.Fatalf("expected orphan cleanup, got snaps=%d rows=%d", snaps, rows)
|
||||
}
|
||||
|
||||
var n int
|
||||
if err := pool.QueryRow(ctx, `SELECT COUNT(*)::int FROM prefix_snapshot WHERE id = $1::uuid`, snapID).Scan(&n); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if n != 0 {
|
||||
t.Fatalf("snapshot still exists")
|
||||
}
|
||||
}
|
||||
|
||||
func TestPostgresEstimateRevisionPruneIntegration(t *testing.T) {
|
||||
dsn := os.Getenv("EVOBGP_TEST_DATABASE_URL")
|
||||
if dsn == "" {
|
||||
t.Skip("EVOBGP_TEST_DATABASE_URL not set")
|
||||
}
|
||||
ctx := context.Background()
|
||||
pool, err := db.OpenPostgresPool(ctx, dsn)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer pool.Close()
|
||||
pg, err := NewPostgres(ctx, pool, false)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
tenant := uuid.NewString()
|
||||
cutoff := pipeline.RevisionCutoffFromMinutes(15)
|
||||
est, err := pg.EstimateRevisionPrune(tenant, cutoff, 15)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if est.RevisionCount != 0 {
|
||||
t.Fatalf("expected 0 revisions for empty tenant, got %d", est.RevisionCount)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user