fix(partman): self-heal default-partition conflicts instead of crash-looping (#139)

* fix(partman): self-heal default-partition conflicts instead of crash-looping

EnsureFuturePartitions died with log.Fatalf when CREATE ... PARTITION OF
failed with SQLSTATE 23514 (check_violation) — Postgres rejects creating a
daily/weekly partition while the default partition still holds a row that
belongs in the new partition's range. A single stray row therefore turned
every startup into a crash loop with no in-app recovery (operator had to
TRUNCATE the default partition by hand).

On 23514, drain exactly the conflicting rows out of the default partition and
attach a fresh partition for the range, in a single transaction: create a
standalone table (without copying identity so original ids re-insert),
DELETE ... RETURNING the in-range rows into it, then ATTACH. If any step
fails the transaction rolls back atomically — rows return to default
untouched, no partition is created, and the next cleanup tick retries. No row
is ever destroyed.

Both startup call sites (operational_logs, activity_log) downgraded from
log.Fatalf to a warning so a partition hiccup degrades to writing into the
default partition rather than taking the server down.

* fix(partman): lock default partition during default-conflict heal

Without the lock, live writers keep routing rows for the missing (current)
period into the default partition while the heal runs; a row committed
between the drain and the ATTACH re-triggers the same 23514 check violation
and rolls the whole heal back, so the retry tick can spin under steady write
load. Locking only the default leaf (not the parent) blocks those inserts
for the heal's short duration while the rest of the table stays readable
and writable.

Addresses CodeRabbit review on #139.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Quick <31828688+Quick104@users.noreply.github.com>
Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
CoffeeKnyte
2026-06-12 15:04:31 -04:00
committed by GitHub
co-authored by Quick Claude Fable 5
parent f4594e902d
commit 91c7eb364d
2 changed files with 113 additions and 4 deletions
+7 -2
View File
@@ -166,7 +166,10 @@ func configureOperationalLogging(
}
opsPM := partman.NewManager(pool, "operational_logs", partman.Daily, 3)
if err := opsPM.EnsureFuturePartitions(ctx); err != nil {
log.Fatalf("ensure operational log partitions: %v", err)
// Non-fatal: a partition hiccup must not crash-loop the server (see the
// operational_logs partition incident). Writes fall back to the default
// partition and the periodic cleanup retries EnsureFuturePartitions.
slog.Warn("ensure operational log partitions; continuing in degraded mode", "error", err)
}
var operationalWriter opslog.Writer
@@ -1529,7 +1532,9 @@ func main() {
}
activityPM := partman.NewManager(pool, "activity_log", partman.Weekly, 2)
if err := activityPM.EnsureFuturePartitions(appCtx); err != nil {
log.Fatalf("ensure activity log partitions: %v", err)
// Non-fatal: see the operational_logs partition incident. Writes fall
// back to the default partition and periodic cleanup retries.
slog.Warn("ensure activity log partitions; continuing in degraded mode", "error", err)
}
var activityWriter activitylog.Writer
activityConsumer := activitylog.NewConsumer(pool, nil, logStreamHub)
+106 -2
View File
@@ -2,15 +2,23 @@ package partman
import (
"context"
"errors"
"fmt"
"regexp"
"sort"
"strings"
"time"
"github.com/jackc/pgx/v5/pgconn"
"github.com/jackc/pgx/v5/pgxpool"
)
// pgCheckViolation is the SQLSTATE Postgres returns when CREATE … PARTITION OF
// (or ATTACH PARTITION) would move rows out of the default partition — i.e. the
// default partition already holds a row that belongs in the new partition's
// range. See incident docs/continuum-to-silo-postgres-migration.md.
const pgCheckViolation = "23514"
const deleteBatchSize = 10000
var partitionBoundsRE = regexp.MustCompile(`FROM \('([^']+)'\) TO \('([^']+)'\)`)
@@ -48,15 +56,111 @@ func (m *Manager) EnsureFuturePartitions(ctx context.Context) error {
lower := m.granularity.addPeriods(start, i)
upper := m.granularity.next(lower)
name := m.partitionName(lower)
if _, err := m.pool.Exec(ctx, fmt.Sprintf(
_, err := m.pool.Exec(ctx, fmt.Sprintf(
`CREATE TABLE IF NOT EXISTS public.%s PARTITION OF public.%s FOR VALUES FROM (%s) TO (%s)`,
quoteIdent(name),
quoteIdent(m.table),
quoteLiteralTimestamp(lower),
quoteLiteralTimestamp(upper),
)); err != nil {
))
if err == nil {
continue
}
// The create fails with a check_violation only when the default
// partition already holds rows that belong in this partition's range.
// Rather than treating that as fatal (the crash-loop in the incident),
// drain those rows out of default and attach the partition so the rows
// land where they belong. Any other error is genuine and propagates.
var pgErr *pgconn.PgError
if !errors.As(err, &pgErr) || pgErr.Code != pgCheckViolation {
return fmt.Errorf("create partition %s: %w", name, err)
}
if healErr := m.healDefaultConflict(ctx, name, lower, upper); healErr != nil {
return fmt.Errorf("create partition %s: heal default conflict: %w", name, healErr)
}
}
return nil
}
// healDefaultConflict recovers from the case where CREATE … PARTITION OF failed
// because the default partition holds rows belonging in [lower, upper). It moves
// exactly those rows out of the default partition and attaches a fresh partition
// for the range, preserving every row (including its original id).
//
// The entire operation runs in a single transaction: a standalone table is
// created (without copying the parent's identity, so original ids re-insert
// cleanly), the conflicting rows are moved into it with a DELETE … RETURNING,
// and it is then ATTACHed — Postgres re-scans the now-drained default partition
// and the attach succeeds. If any step fails (including the re-insert), the
// transaction rolls back atomically: the rows return to the default partition
// untouched and no partition is created, so the caller can safely retry later.
// No row is ever destroyed.
func (m *Manager) healDefaultConflict(ctx context.Context, name string, lower, upper time.Time) error {
defaultTable := m.defaultPartitionName()
tx, err := m.pool.Begin(ctx)
if err != nil {
return fmt.Errorf("begin: %w", err)
}
defer func() { _ = tx.Rollback(ctx) }()
// Block concurrent writes to the default partition for the duration of the
// heal. The missing partition is the current period, so live writers are
// actively routing rows into default; without this lock a row committed
// between the drain and the attach re-triggers the same check violation and
// rolls the whole heal back. Locking only the default leaf (not the parent)
// keeps the rest of the table readable and writable; tuple routing must
// lock the leaf to insert, so this is sufficient.
if _, err := tx.Exec(ctx, fmt.Sprintf(
`LOCK TABLE public.%s IN ACCESS EXCLUSIVE MODE`,
quoteIdent(defaultTable),
)); err != nil {
return fmt.Errorf("lock default partition: %w", err)
}
// Standalone table matching the parent's columns. INCLUDING DEFAULTS is
// deliberately the only option: it must NOT copy the parent's identity, so
// the preserved rows insert with their original ids instead of regenerating.
if _, err := tx.Exec(ctx, fmt.Sprintf(
`CREATE TABLE public.%s (LIKE public.%s INCLUDING DEFAULTS)`,
quoteIdent(name),
quoteIdent(m.table),
)); err != nil {
return fmt.Errorf("create standalone table: %w", err)
}
// Move the conflicting rows out of default into the standalone table in one
// statement. RETURNING * preserves parent column order, matching the LIKE.
if _, err := tx.Exec(ctx, fmt.Sprintf(`
WITH moved AS (
DELETE FROM public.%s
WHERE "timestamp" >= $1 AND "timestamp" < $2
RETURNING *
)
INSERT INTO public.%s SELECT * FROM moved
`,
quoteIdent(defaultTable),
quoteIdent(name),
), lower.UTC(), upper.UTC()); err != nil {
return fmt.Errorf("drain default rows: %w", err)
}
// Attach. Postgres re-scans the default partition to confirm no row still
// falls in [lower, upper); the drain above makes that pass.
if _, err := tx.Exec(ctx, fmt.Sprintf(
`ALTER TABLE public.%s ATTACH PARTITION public.%s FOR VALUES FROM (%s) TO (%s)`,
quoteIdent(m.table),
quoteIdent(name),
quoteLiteralTimestamp(lower),
quoteLiteralTimestamp(upper),
)); err != nil {
return fmt.Errorf("attach partition: %w", err)
}
if err := tx.Commit(ctx); err != nil {
return fmt.Errorf("commit: %w", err)
}
return nil