diff --git a/.changeset/archive-read-handoff.md b/.changeset/archive-read-handoff.md new file mode 100644 index 00000000..3709d309 --- /dev/null +++ b/.changeset/archive-read-handoff.md @@ -0,0 +1,5 @@ +--- +"ftw": patch +--- + +Keep history queries and live commits responsive during archive building. Protect only file publication and bounded pruning against readers, and release archive write locks before retrying database contention. diff --git a/go/internal/state/archive_contention_test.go b/go/internal/state/archive_contention_test.go new file mode 100644 index 00000000..936f9dfa --- /dev/null +++ b/go/internal/state/archive_contention_test.go @@ -0,0 +1,534 @@ +package state + +import ( + "context" + "database/sql" + "database/sql/driver" + "errors" + "fmt" + "os" + "path/filepath" + "sync" + "sync/atomic" + "testing" + "time" + + "modernc.org/sqlite" +) + +func TestPartialHourQueriesAndLiveWritesDuringArchiveBuild(t *testing.T) { + s := freshStore(t) + base := time.Now().UTC().Truncate(time.Hour).Add(-7 * 24 * time.Hour).UnixMilli() + if err := s.RecordSamples([]Sample{ + {Driver: "meter", Metric: "power", TsMs: base, Value: 999}, + {Driver: "meter", Metric: "power", TsMs: base + 1, Value: 10}, + {Driver: "meter", Metric: "power", TsMs: base + 24*seriesHourMs, Value: 20}, + {Driver: "meter", Metric: "power", TsMs: base + 72*seriesHourMs + 1, Value: 30}, + {Driver: "meter", Metric: "power", TsMs: base + 72*seriesHourMs + 2, Value: 999}, + }); err != nil { + t.Fatal(err) + } + if err := s.ensureSeriesHours(context.Background()); err != nil { + t.Fatal(err) + } + // The slow staging/compression phase owns the archive operation lock. + // It must not exclude queries, including exact partial-hour boundaries. + s.archiveMu.Lock() + defer s.archiveMu.Unlock() + for i := 0; i < 8; i++ { + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + points, err := s.LoadSeriesBucketsContext(ctx, "meter", "power", base+1, base+72*seriesHourMs+1, 400) + cancel() + if err != nil { + t.Fatal("archive build blocked a completed hourly query:", err) + } + var n int64 + var sum float64 + for _, p := range points { + n += p.N + sum += p.V * float64(p.N) + } + if n != 3 || sum != 60 { + t.Fatalf("partial boundaries included or lost samples: n=%d sum=%v", n, sum) + } + if err := s.EnqueueTelemetryTick(nil, []Sample{{Driver: "meter", Metric: "power", TsMs: base + 96*seriesHourMs + int64(i), Value: float64(i)}}, nil); err != nil { + t.Fatal(err) + } + ctx, cancel = context.WithTimeout(context.Background(), time.Second) + err = s.FlushHistory(ctx) + cancel() + if err != nil { + t.Fatal("archive build blocked live commit:", err) + } + } + if st := s.HistoryWriterStatus(); st.Accepted != 8 || st.Committed != 8 || st.Pending != 0 || st.Rejected != 0 || st.LastError != "" { + t.Fatalf("live writes did not remain durable: %+v", st) + } +} + +func TestDenseArchiveHourReadsOutsideWriteBudgetAndRetriesSnapshot(t *testing.T) { + reading, release := make(chan struct{}), make(chan struct{}) + var entered, released sync.Once + defer released.Do(func() { close(release) }) + fn := fmt.Sprintf("test_archive_slow_hour_%d", time.Now().UnixNano()) + if err := sqlite.RegisterScalarFunction(fn, 1, func(_ *sqlite.FunctionContext, args []driver.Value) (driver.Value, error) { + entered.Do(func() { close(reading); <-release }) + return args[0], nil + }); err != nil { + t.Fatal(err) + } + s := freshStore(t) + base := time.Now().UTC().Truncate(time.Hour).UnixMilli() + if err := s.RecordSamples([]Sample{{Driver: "meter", Metric: "power", TsMs: base + 1, Value: 10}}); err != nil { + t.Fatal(err) + } + d, err := s.driverID("meter") + if err != nil { + t.Fatal(err) + } + m, err := s.metricID("power", "") + if err != nil { + t.Fatal(err) + } + for _, q := range []string{ + `ALTER TABLE ts_samples RENAME TO slow_archive_source`, + `CREATE VIEW ts_samples AS SELECT driver_id,metric_id,ts_ms,` + fn + `(value) AS value FROM slow_archive_source`, + } { + if _, err := s.history.Exec(q); err != nil { + t.Fatal(err) + } + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + done := make(chan error, 1) + go func() { done <- s.mergeArchivedHour(ctx, d, m, base, map[int64]float64{base + 2: 20}) }() + select { + case <-reading: + case <-ctx.Done(): + t.Fatal("hour read did not start") + } + // A dense/slow scan may exceed the write budget. It must not own the + // live mutex, and its snapshot must be retried if raw data changes. + lockCtx, stopLock := context.WithTimeout(ctx, 200*time.Millisecond) + err = lockContext(lockCtx, s.historyWriteMu.TryLock) + stopLock() + if err != nil { + t.Fatal("slow hour read held the live writer mutex:", err) + } + s.historyWriteMu.Unlock() + if _, err := s.history.ExecContext(ctx, `INSERT INTO slow_archive_source(driver_id,metric_id,ts_ms,value) VALUES(?,?,?,?)`, d, m, base+3, 30); err != nil { + t.Fatal(err) + } + time.Sleep(archiveWriteTimeout + 100*time.Millisecond) + released.Do(func() { close(release) }) + if err := <-done; err != nil { + t.Fatal("dense hour could not resume:", err) + } + var n int64 + var sum float64 + if err := s.history.QueryRow(`SELECT n,sum_value FROM ts_series_hour WHERE driver_id=? AND metric_id=? AND hour_ms=?`, d, m, base).Scan(&n, &sum); err != nil || n != 3 || sum != 60 { + t.Fatalf("snapshot retry lost the late raw sample: n=%d sum=%v err=%v", n, sum, err) + } +} + +func TestPreparedArchiveHourRetriesInnerWriteDeadline(t *testing.T) { + s := freshStore(t) + ctx, cancel := context.WithTimeout(context.Background(), 4*time.Second) + defer cancel() + attempts := 0 + err := s.archiveTransaction(ctx, func(context.Context, *sql.Tx) error { + attempts++ + return nil + }, func(ctx context.Context, tx *sql.Tx) error { + if attempts == 1 { + <-ctx.Done() + return ctx.Err() + } + _, err := tx.ExecContext(ctx, `INSERT INTO history_migrations(name) VALUES ('prepared-hour-retried')`) + return err + }) + if err != nil || attempts != 2 { + t.Fatalf("inner deadline abandoned the prepared hour: attempts=%d err=%v", attempts, err) + } + var n int + if err := s.history.QueryRow(`SELECT COUNT(*) FROM history_migrations WHERE name='prepared-hour-retried'`).Scan(&n); err != nil || n != 1 { + t.Fatalf("retried write missing: %d %v", n, err) + } +} + +func TestArchivePruneKeepsReaderViewAndLateCorrections(t *testing.T) { + s := freshStore(t) + s.coldDir = t.TempDir() + base := time.Now().UTC().Add(-20 * 24 * time.Hour).Truncate(24 * time.Hour) + path := filepath.Join(s.coldDir, base.Format("2006/01/02.parquet")) + if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil { + t.Fatal(err) + } + var archive []parquetSampleRow + var copied, current []Sample + for i := 1; i <= 3; i++ { + ts := base.UnixMilli() + int64(i)*1000 + archive = append(archive, parquetSampleRow{TsMs: ts, Driver: "meter", Metric: "power", Value: float64(i * 10)}) + copied = append(copied, Sample{TsMs: ts, Driver: "meter", Metric: "power", Value: float64(i * 10)}) + current = append(current, copied[i-1]) + } + current[1].Value = 99 // A correction arrived after the archive copy. + if err := s.RecordSamples(current); err != nil { + t.Fatal(err) + } + if err := writeParquetDay(path, archive); err != nil { + t.Fatal(err) + } + resolved, err := s.resolveSamples(copied) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + reading, release := make(chan struct{}), make(chan struct{}) + var entered, released sync.Once + defer released.Do(func() { close(release) }) + readDone := make(chan error, 1) + seen := map[int64]float64{} + go func() { + readDone <- s.walkMergedSeries(ctx, s.coldDir, "meter", "power", base.UnixMilli(), base.Add(time.Hour).UnixMilli(), func(ts int64, v float64) error { + entered.Do(func() { close(reading); <-release }) + seen[ts] = v + return nil + }) + }() + select { + case <-reading: + case <-ctx.Done(): + t.Fatal("reader did not start") + } + type pruneResult struct { + n int64 + err error + } + pruned := make(chan pruneResult, 1) + go func() { n, err := s.pruneArchivedSamples(ctx, resolved); pruned <- pruneResult{n, err} }() + select { + case r := <-pruned: + t.Fatalf("pruning crossed an active reader: %+v", r) + case <-time.After(30 * time.Millisecond): + } + if err := s.EnqueueTelemetryTick(nil, []Sample{{Driver: "meter", Metric: "power", TsMs: time.Now().UnixMilli(), Value: 7000}}, nil); err != nil { + t.Fatal(err) + } + live, done := context.WithTimeout(ctx, 500*time.Millisecond) + err = s.FlushHistory(live) + done() + if err != nil { + t.Fatal("pruning waited for a reader while holding the live writer:", err) + } + released.Do(func() { close(release) }) + if err := <-readDone; err != nil { + t.Fatal(err) + } + r := <-pruned + if r.err != nil || r.n != 2 { + t.Fatalf("prune must preserve the late correction: %+v", r) + } + points, err := s.LoadSeries("meter", "power", base.UnixMilli(), base.Add(time.Hour).UnixMilli(), 0) + if err != nil || len(points) != 3 || len(seen) != 3 { + t.Fatalf("handover lost or duplicated samples: before=%v after=%v err=%v", seen, points, err) + } + for _, p := range points { + if seen[p.TsMs] != p.Value { + t.Fatalf("handover changed %d: %v -> %v", p.TsMs, seen[p.TsMs], p.Value) + } + } + if points[1].Value != 99 { + t.Fatal("archive replaced a later raw value") + } +} + +func TestArchivePublicationAndRemovalRespectReaderCancellation(t *testing.T) { + s := freshStore(t) + path, tmp := filepath.Join(t.TempDir(), "day.parquet"), filepath.Join(t.TempDir(), "verified.tmp") + if err := os.WriteFile(path, []byte("original"), 0600); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(tmp, []byte("verified replacement"), 0600); err != nil { + t.Fatal(err) + } + s.archiveViewMu.RLock() + for _, operation := range []func(context.Context) error{ + func(ctx context.Context) error { return s.replaceArchive(ctx, tmp, path) }, + func(ctx context.Context) error { return s.removeArchive(ctx, path) }, + } { + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Millisecond) + err := operation(ctx) + cancel() + if !errors.Is(err, context.DeadlineExceeded) { + s.archiveViewMu.RUnlock() + t.Fatalf("blocked operation: %v", err) + } + } + s.archiveViewMu.RUnlock() + if b, err := os.ReadFile(path); err != nil || string(b) != "original" { + t.Fatalf("canceled operation changed source: %q %v", b, err) + } + if err := s.replaceArchive(context.Background(), tmp, path); err != nil { + t.Fatal(err) + } + if b, err := os.ReadFile(path); err != nil || string(b) != "verified replacement" { + t.Fatalf("replacement missing: %q %v", b, err) + } +} + +func TestArchiveSQLiteBusyDoesNotConsumeLiveCommitBudget(t *testing.T) { + s := freshStore(t) + if err := s.RecordSamples([]Sample{{Driver: "meter", Metric: "power", TsMs: 1, Value: 1}}); err != nil { + t.Fatal(err) + } + s.historyWriter.commitInterval = 0 + s.historyWriter.commitTimeout = time.Second + var expired atomic.Int32 + s.historyWriter.commitFn = func(ctx context.Context, batches []historyBatch, ack int64) (historyBatchCommit, error) { + out, err := s.recordHistoryBatches(ctx, batches, ack) + if errors.Is(err, context.DeadlineExceeded) { + expired.Add(1) + } + return out, err + } + ctx, cancel := context.WithTimeout(context.Background(), 4*time.Second) + defer cancel() + blocker, err := s.history.Conn(ctx) + if err != nil { + t.Fatal(err) + } + defer blocker.Close() + if _, err := blocker.ExecContext(ctx, `BEGIN IMMEDIATE`); err != nil { + t.Fatal(err) + } + defer blocker.ExecContext(context.Background(), `ROLLBACK`) + archiveCtx, stopArchive := context.WithTimeout(ctx, 150*time.Millisecond) + defer stopArchive() + attempted := make(chan struct{}) + var once sync.Once + archived := make(chan error, 1) + go func() { + archived <- s.writeArchiveBatch(archiveCtx, func(ctx context.Context, tx *sql.Tx) error { + once.Do(func() { close(attempted) }) + _, err := tx.ExecContext(ctx, `DELETE FROM ts_samples WHERE ts_ms=1`) + return err + }) + }() + select { + case <-attempted: + case <-ctx.Done(): + t.Fatal("archive did not attempt the write") + } + started := time.Now() + for i := 0; i < 24; i++ { + if err := s.EnqueueTelemetryTick(nil, []Sample{{Driver: "meter", Metric: "power", TsMs: int64(i + 2), Value: float64(i)}}, nil); err != nil { + t.Fatal(err) + } + time.Sleep(5 * time.Millisecond) + } + // The archive must finish cancellation while the external SQL lock is + // still held, not wait for SQLite's default five-second busy timeout. + select { + case err := <-archived: + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("archive cancellation: %v", err) + } + case <-time.After(350 * time.Millisecond): + t.Fatal("archive kept the writer mutex inside SQLite busy handling") + } + if _, err := blocker.ExecContext(ctx, `ROLLBACK`); err != nil { + t.Fatal(err) + } + if err := s.FlushHistory(ctx); err != nil { + t.Fatal(err) + } + if st := s.HistoryWriterStatus(); st.Committed != 24 || st.Pending != 0 || st.Rejected != 0 || st.LastError != "" { + t.Fatalf("temporary archive/database lock exhausted live writes: %+v", st) + } + if time.Since(started) > time.Second { + t.Fatal("live commit consumed its full budget") + } + if expired.Load() != 0 { + t.Fatal("database contention exhausted a live transaction budget") + } + // Reserve the entire pool to verify no zero-busy connection leaked. + if err := blocker.Close(); err != nil { + t.Fatal(err) + } + var conns []*sql.Conn + defer func() { + for _, c := range conns { + c.Close() + } + }() + for i := 0; i < 4; i++ { + c, err := s.history.Conn(ctx) + if err != nil { + t.Fatal(err) + } + conns = append(conns, c) + var busy int + if err := c.QueryRowContext(ctx, `PRAGMA busy_timeout`).Scan(&busy); err != nil || busy != 5000 { + t.Fatalf("connection busy timeout leaked: %d %v", busy, err) + } + } +} + +func TestSlowArchiveTransactionYieldsBeforeLiveDeadline(t *testing.T) { + s := freshStore(t) + if err := s.RecordSamples([]Sample{{Driver: "meter", Metric: "power", TsMs: 1, Value: 1}}); err != nil { + t.Fatal(err) + } + s.historyWriter.commitInterval = 0 + s.historyWriter.commitTimeout = 2 * time.Second + ctx, cancel := context.WithTimeout(context.Background(), 4*time.Second) + defer cancel() + writing := make(chan struct{}) + archived := make(chan error, 1) + go func() { + archived <- s.writeArchiveBatch(ctx, func(ctx context.Context, tx *sql.Tx) error { + if _, err := tx.ExecContext(ctx, `DELETE FROM ts_samples WHERE ts_ms=1`); err != nil { + return err + } + close(writing) + <-ctx.Done() // Synthetic slow IO while the archive owns SQLite's writer. + return ctx.Err() + }) + }() + select { + case <-writing: + case <-ctx.Done(): + t.Fatal("archive did not start") + } + for i := 0; i < 40; i++ { + if err := s.EnqueueTelemetryTick(nil, []Sample{{Driver: "meter", Metric: "power", TsMs: int64(i + 2), Value: float64(i)}}, nil); err != nil { + t.Fatal(err) + } + time.Sleep(5 * time.Millisecond) + } + if err := <-archived; !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("slow archive was not bounded: %v", err) + } + if err := s.FlushHistory(ctx); err != nil { + t.Fatal(err) + } + if st := s.HistoryWriterStatus(); st.Accepted != 40 || st.Committed != 40 || st.Pending != 0 || st.Rejected != 0 || st.LastError != "" { + t.Fatalf("slow archive exhausted the live commit budget: %+v", st) + } + var n int + if err := s.history.QueryRow(`SELECT COUNT(*) FROM ts_samples`).Scan(&n); err != nil || n != 41 { + t.Fatalf("rollback or live samples lost: %d %v", n, err) + } +} + +func TestSlowParquetStagingDoesNotBlockPartialHours(t *testing.T) { + s := freshStore(t) + s.coldDir = t.TempDir() + base := time.Now().UTC().Truncate(24 * time.Hour).Add(-7 * 24 * time.Hour) + first, last := base.UnixMilli()+1, base.Add(72*time.Hour).UnixMilli()+1 + samples := []Sample{{Driver: "meter", Metric: "power", TsMs: first, Value: 10}, {Driver: "meter", Metric: "power", TsMs: last, Value: 20}} + if err := s.RecordSamples(samples); err != nil { + t.Fatal(err) + } + if err := s.ensureSeriesHours(context.Background()); err != nil { + t.Fatal(err) + } + path := filepath.Join(s.coldDir, base.Format("2006/01/02.parquet")) + if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil { + t.Fatal(err) + } + stage, err := openArchiveStage(filepath.Join(t.TempDir(), "staging.db")) + if err != nil { + t.Fatal(err) + } + defer stage.Close() + if err := insertArchiveRows(context.Background(), stage, []parquetSampleRow{{Driver: "meter", Metric: "power", TsMs: first, Value: 10}}); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + // Reserving the stage's sole connection pauses the real publication path + // before compression/readback, without slowing the primary history DB. + blocked, err := stage.Conn(ctx) + if err != nil { + t.Fatal(err) + } + defer blocked.Close() + s.archiveMu.Lock() + defer s.archiveMu.Unlock() + published := make(chan error, 1) + go func() { published <- s.publishStagedSamples(ctx, path, stage) }() + for stage.Stats().WaitCount == 0 { + select { + case err := <-published: + t.Fatalf("publisher did not wait for staging: %v", err) + case <-ctx.Done(): + t.Fatal(ctx.Err()) + case <-time.After(time.Millisecond): + } + } + readCtx, stopRead := context.WithTimeout(ctx, 500*time.Millisecond) + points, err := s.LoadSeriesBucketsContext(readCtx, "meter", "power", first, last, 400) + stopRead() + if err != nil || len(points) != 2 { + t.Fatalf("staging blocked exact boundary query: %v %v", points, err) + } + if err := s.EnqueueTelemetryTick(nil, []Sample{{Driver: "meter", Metric: "power", TsMs: time.Now().UnixMilli(), Value: 7000}}, nil); err != nil { + t.Fatal(err) + } + live, stopLive := context.WithTimeout(ctx, 500*time.Millisecond) + err = s.FlushHistory(live) + stopLive() + if err != nil { + t.Fatal("staging blocked live commit:", err) + } + if err := blocked.Close(); err != nil { + t.Fatal(err) + } + if err := <-published; err != nil { + t.Fatal(err) + } + if err := VerifyParquetFile(ctx, path); err != nil { + t.Fatal(err) + } + points, err = s.LoadSeriesBucketsContext(ctx, "meter", "power", first, last, 400) + var n int64 + for _, p := range points { + n += p.N + } + if err != nil || n != 2 { + t.Fatalf("published overlap duplicated or lost data: %v %v", points, err) + } +} + +func TestHistoryBatchLockWaitHonorsAttemptDeadline(t *testing.T) { + s := freshStore(t) + if err := s.RecordSamples([]Sample{{Driver: "meter", Metric: "power", TsMs: 1, Value: 1}}); err != nil { + t.Fatal(err) + } + s.historyWriteMu.Lock() + var released sync.Once + defer released.Do(s.historyWriteMu.Unlock) + ctx, cancel := context.WithTimeout(context.Background(), 40*time.Millisecond) + defer cancel() + done := make(chan error, 1) + go func() { + _, err := s.recordHistoryBatches(ctx, []historyBatch{{id: "canceled-before-write", hash: "same-payload", payload: historyPayload{Samples: []Sample{{Driver: "meter", Metric: "power", TsMs: 2, Value: 2}}}}}, 0) + done <- err + }() + select { + case err := <-done: + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("lock cancellation: %v", err) + } + case <-time.After(500 * time.Millisecond): + t.Fatal("writer ignored its deadline while waiting for the mutex") + } + released.Do(s.historyWriteMu.Unlock) + var n int + if err := s.history.QueryRow(`SELECT COUNT(*) FROM ts_samples WHERE ts_ms=2`).Scan(&n); err != nil || n != 0 { + t.Fatalf("canceled attempt wrote data: %d %v", n, err) + } +} diff --git a/go/internal/state/archive_write.go b/go/internal/state/archive_write.go new file mode 100644 index 00000000..5973fb1c --- /dev/null +++ b/go/internal/state/archive_write.go @@ -0,0 +1,105 @@ +package state + +import ( + "context" + "database/sql" + "database/sql/driver" + "errors" + "fmt" + "time" +) + +const archiveWriteTimeout = time.Second + +// Archive transactions must yield the writer mutex instead of spending the +// live commit's entire budget in SQLite's busy handler. Only the reserved +// connection uses a zero busy timeout; restore it before returning to the pool. +func (s *Store) writeArchiveBatch(ctx context.Context, write func(context.Context, *sql.Tx) error) error { + return s.archiveTransaction(ctx, nil, write) +} + +// Preparation reads a SQLite snapshot without owning either writer lock. A +// live commit makes a later snapshot upgrade return BUSY; retry preparation +// as well as the write so a new sample cannot be lost from an hourly total. +func (s *Store) archiveTransaction(ctx context.Context, prepare, write func(context.Context, *sql.Tx) error) error { + for { + if err := ctx.Err(); err != nil { + return err + } + if s.HistoryWriterStatus().Pending < historyCommitMaxTicks/2 { + err := s.tryArchiveBatch(ctx, prepare, write) + // A prepared hour stays in its current file job on a transient + // write/lock deadline. Pruning returns deadlines to split its batch. + retryPrepared := prepare != nil && ctx.Err() == nil && errors.Is(err, context.DeadlineExceeded) + if !historyWriteBusy(err) && !retryPrepared { + return err + } + } + // No transaction, reader lock or writer mutex survives this pause. + if err := pauseMaintenance(ctx); err != nil { + return err + } + } +} + +func (s *Store) tryArchiveBatch(parent context.Context, prepare, write func(context.Context, *sql.Tx) error) error { + conn, err := s.history.Conn(parent) + if err != nil { + return err + } + defer conn.Close() + var busyMS int + if err := conn.QueryRowContext(parent, `PRAGMA busy_timeout`).Scan(&busyMS); err != nil { + return err + } + defer func() { + restore, done := context.WithTimeout(context.Background(), time.Second) + defer done() + if _, err := conn.ExecContext(restore, fmt.Sprintf("PRAGMA busy_timeout=%d", busyMS)); err != nil { + // Never lend a connection with changed lock behaviour to live work. + _ = conn.Raw(func(any) error { return driver.ErrBadConn }) + } + }() + if _, err := conn.ExecContext(parent, `PRAGMA busy_timeout=0`); err != nil { + return err + } + txCtx, cancelTx := context.WithCancel(parent) + defer cancelTx() + tx, err := conn.BeginTx(txCtx, nil) + if err != nil { + return err + } + defer tx.Rollback() + if prepare != nil { + readCtx, cancel := context.WithTimeout(parent, 30*time.Second) + err := prepare(readCtx, tx) + cancel() + if err != nil { + return err + } + } + ctx, cancel := context.WithTimeout(parent, archiveWriteTimeout) + defer cancel() + // Tx.Commit uses the context from BeginTx. Arm its deadline only after + // preparation, so the same short budget also covers the durable commit. + stopDeadline := context.AfterFunc(ctx, cancelTx) + defer stopDeadline() + // Readers may postpone archive writes, but must never make an archive + // operation wait while holding the live writer mutex. + if err := lockContext(ctx, s.archiveViewMu.TryLock); err != nil { + return err + } + defer s.archiveViewMu.Unlock() + if err := lockContext(ctx, s.historyWriteMu.TryLock); err != nil { + return err + } + defer s.historyWriteMu.Unlock() + if err := write(ctx, tx); err != nil { + return err + } + err = tx.Commit() + if err != nil && ctx.Err() != nil { + return ctx.Err() + } + return err +} diff --git a/go/internal/state/history_writer.go b/go/internal/state/history_writer.go index bfa7c138..3020d2d8 100644 --- a/go/internal/state/history_writer.go +++ b/go/internal/state/history_writer.go @@ -12,6 +12,7 @@ import ( "time" "github.com/google/uuid" + "modernc.org/sqlite" ) const ( @@ -132,6 +133,11 @@ func historyCommitInterrupted(err error) bool { return false } +func historyWriteBusy(err error) bool { + var busy *sqlite.Error + return errors.As(err, &busy) && (busy.Code()&0xff == 5 || busy.Code()&0xff == 6) +} + // EnqueueTelemetryTick copies a whole tick without waiting on disk. The caller // must report an error as a collection gap. Successful admission is volatile // until HistoryWriterStatus reports the commit; FlushHistory waits for that. @@ -285,7 +291,13 @@ func (w *historyWriter) run() { maxAttempt = next continue } - timer := time.NewTimer(time.Second) + retryDelay := time.Second + if historyWriteBusy(err) { + // A competing transaction can release its lock quickly. Keep the + // accepted batch and retry without a full second of dead time. + retryDelay = 100 * time.Millisecond + } + timer := time.NewTimer(retryDelay) select { case <-w.ctx.Done(): timer.Stop() diff --git a/go/internal/state/parquet_stream.go b/go/internal/state/parquet_stream.go index 9b3195b2..2fffaafb 100644 --- a/go/internal/state/parquet_stream.go +++ b/go/internal/state/parquet_stream.go @@ -194,7 +194,7 @@ func (s *Store) archiveSampleDay(ctx context.Context, coldDir string, from, to i if copied == 0 { return 0, "", nil } - if err := publishStagedSamples(ctx, path, stage); err != nil { + if err := s.publishStagedSamples(ctx, path, stage); err != nil { return copied, "", err } // Hourly totals already include live samples and any prior archive. Do @@ -248,7 +248,7 @@ func (s *Store) archiveSampleDay(ctx context.Context, coldDir string, from, to i return deleted, path, nil } -func publishStagedSamples(ctx context.Context, path string, stage *sql.DB) error { +func (s *Store) publishStagedSamples(ctx context.Context, path string, stage *sql.DB) error { f, err := os.CreateTemp(filepath.Dir(path), ".ftw-parquet-*.tmp") if err != nil { return err @@ -335,41 +335,54 @@ func publishStagedSamples(ctx context.Context, path string, stage *sql.DB) error if err := verify(tmp); err != nil { return err } - if err := os.Rename(tmp, path); err != nil { - return err - } - if err := syncDir(filepath.Dir(path)); err != nil { + if err := s.replaceArchive(ctx, tmp, path); err != nil { return err } return verify(path) } func (s *Store) pruneArchivedSamples(ctx context.Context, batch []resolvedSample) (int64, error) { - s.historyWriteMu.Lock() - defer s.historyWriteMu.Unlock() - tx, err := s.history.BeginTx(ctx, nil) - if err != nil { - return 0, err - } - defer tx.Rollback() - stmt, err := tx.PrepareContext(ctx, `DELETE FROM ts_samples WHERE driver_id=? AND metric_id=? AND ts_ms=? AND value=?`) - if err != nil { - return 0, err - } - defer stmt.Close() - var n int64 - for _, r := range batch { - res, err := stmt.ExecContext(ctx, r.dID, r.mID, r.ts, r.v) + var deleted int64 + limit := 64 + for len(batch) > 0 { + n := min(len(batch), limit) + var removed int64 + err := s.writeArchiveBatch(ctx, func(ctx context.Context, tx *sql.Tx) error { + removed = 0 // A busy attempt rolls back and retries this same prefix. + stmt, err := tx.PrepareContext(ctx, `DELETE FROM ts_samples WHERE driver_id=? AND metric_id=? AND ts_ms=? AND value=?`) + if err != nil { + return err + } + defer stmt.Close() + for _, r := range batch[:n] { + res, err := stmt.ExecContext(ctx, r.dID, r.mID, r.ts, r.v) + if err != nil { + return err + } + count, err := res.RowsAffected() + if err != nil { + return err + } + removed += count + } + return nil + }) if err != nil { - return n, err + if ctx.Err() == nil && errors.Is(err, context.DeadlineExceeded) && n > 1 { + limit = max(1, n/2) + continue + } + return deleted, err } - v, err := res.RowsAffected() - if err != nil { - return n, err + deleted += removed + batch = batch[n:] + if len(batch) > 0 && s.HistoryWriterStatus().Pending > 0 { + if err := pauseMaintenance(ctx); err != nil { + return deleted, err + } } - n += v } - return n, tx.Commit() + return deleted, nil } func (s *Store) archiveDayHours(ctx context.Context, stage *sql.DB) error { @@ -436,43 +449,46 @@ func (s *Store) archiveDayHours(ctx context.Context, stage *sql.DB) error { } func (s *Store) mergeArchivedHour(ctx context.Context, d, m, hour int64, values map[int64]float64) error { - s.historyWriteMu.Lock() - defer s.historyWriteMu.Unlock() - tx, err := s.history.BeginTx(ctx, nil) - if err != nil { - return err - } - defer tx.Rollback() - rows, err := tx.QueryContext(ctx, `SELECT ts_ms,value FROM ts_samples WHERE driver_id=? AND metric_id=? AND ts_ms>=? AND ts_ms=? AND ts_ms= maxRawSeriesPoints { - rows.Close() - return ErrHistoryQueryLimit + for rows.Next() { + var ts int64 + var v float64 + if err := rows.Scan(&ts, &v); err != nil { + rows.Close() + return err + } + if _, exists := merged[ts]; !exists && len(merged) >= maxRawSeriesPoints { + rows.Close() + return ErrHistoryQueryLimit + } + merged[ts] = v } - values[ts] = v - } - err = errors.Join(rows.Err(), rows.Close()) - if err != nil { - return err - } - var a seriesBucketAcc - for ts, v := range values { - a.add(1, v, v, v, ts) - } - _, err = tx.ExecContext(ctx, `INSERT INTO ts_series_hour VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(driver_id,metric_id,hour_ms) DO UPDATE SET sum_value=excluded.sum_value,min_value=excluded.min_value,max_value=excluded.max_value,n=excluded.n,last_ts_ms=excluded.last_ts_ms`, d, m, hour, a.sum, a.min, a.max, a.n, a.last) - if err != nil { - return err - } - return tx.Commit() + err = errors.Join(rows.Err(), rows.Close()) + if err != nil { + return err + } + a = seriesBucketAcc{} + for ts, v := range merged { + a.add(1, v, v, v, ts) + } + return nil + }, func(ctx context.Context, tx *sql.Tx) error { + _, err := tx.ExecContext(ctx, `INSERT INTO ts_series_hour VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(driver_id,metric_id,hour_ms) DO UPDATE SET sum_value=excluded.sum_value,min_value=excluded.min_value,max_value=excluded.max_value,n=excluded.n,last_ts_ms=excluded.last_ts_ms`, d, m, hour, a.sum, a.min, a.max, a.n, a.last) + if err != nil { + return err + } + return nil + }) } func pauseMaintenance(ctx context.Context) error { diff --git a/go/internal/state/series_archive.go b/go/internal/state/series_archive.go index 476efa74..0c0d6d65 100644 --- a/go/internal/state/series_archive.go +++ b/go/internal/state/series_archive.go @@ -46,12 +46,12 @@ func parquetPaths(coldDir string, since, until int64) ([]string, error) { // walkMergedSeries only holds one selected series/day in memory. Raw SQLite // wins on matching timestamps during a retry between file publish and prune. func (s *Store) walkMergedSeries(ctx context.Context, coldDir, driver, metric string, since, until int64, visit func(int64, float64) error) error { - // The day list and its raw overlap must stay on one side of archive - // publication/pruning. This lock never blocks the live history writer. - if err := s.lockArchive(ctx); err != nil { + // Keep the files and their raw overlap on one side of publication/pruning. + // Staging and compression do not need this lock; live writes never take it. + if err := lockContext(ctx, s.archiveViewMu.TryRLock); err != nil { return err } - defer s.archiveMu.Unlock() + defer s.archiveViewMu.RUnlock() if until < since { return nil } @@ -217,10 +217,7 @@ func (s *Store) retainSampleHistory(ctx context.Context, days int, now time.Time if err := s.summarizeParquetDay(ctx, path); err != nil { return err } - if err := os.Remove(path); err != nil { - return err - } - if err := syncDir(filepath.Dir(path)); err != nil { + if err := s.removeArchive(ctx, path); err != nil { return err } } @@ -319,11 +316,15 @@ func (s *Store) markParquetSummary(ctx context.Context, path string) error { } func (s *Store) lockArchive(ctx context.Context) error { + return lockContext(ctx, s.archiveMu.TryLock) +} + +func lockContext(ctx context.Context, tryLock func() bool) error { for { if err := ctx.Err(); err != nil { return err } - if s.archiveMu.TryLock() { + if tryLock() { return nil } timer := time.NewTimer(10 * time.Millisecond) @@ -335,3 +336,25 @@ func (s *Store) lockArchive(ctx context.Context) error { } } } + +func (s *Store) removeArchive(ctx context.Context, path string) error { + if err := lockContext(ctx, s.archiveViewMu.TryLock); err != nil { + return err + } + defer s.archiveViewMu.Unlock() + if err := os.Remove(path); err != nil { + return err + } + return syncDir(filepath.Dir(path)) +} + +func (s *Store) replaceArchive(ctx context.Context, tmp, path string) error { + if err := lockContext(ctx, s.archiveViewMu.TryLock); err != nil { + return err + } + defer s.archiveViewMu.Unlock() + if err := os.Rename(tmp, path); err != nil { + return err + } + return syncDir(filepath.Dir(path)) +} diff --git a/go/internal/state/storage_recovery_test.go b/go/internal/state/storage_recovery_test.go index 586387e2..259a2fd1 100644 --- a/go/internal/state/storage_recovery_test.go +++ b/go/internal/state/storage_recovery_test.go @@ -644,7 +644,7 @@ func writeParquetDay(path string, rows []parquetSampleRow) error { if err := insertArchiveRows(context.Background(), stage, rows); err != nil { return err } - return publishStagedSamples(context.Background(), path, stage) + return new(Store).publishStagedSamples(context.Background(), path, stage) } // Opt-in admission fixture: records kernel peak RSS on the target. Also use diff --git a/go/internal/state/store.go b/go/internal/state/store.go index 38504373..197398ff 100644 --- a/go/internal/state/store.go +++ b/go/internal/state/store.go @@ -53,11 +53,12 @@ type Store struct { hotPath string hotWriteMu sync.Mutex - archiveMu sync.Mutex - coldDir string - db *sql.DB - cache *sql.DB - ts *internCache + archiveMu sync.Mutex // Serializes archive construction and retention. + archiveViewMu sync.RWMutex // Protects file publication/pruning against readers. + coldDir string + db *sql.DB + cache *sql.DB + ts *internCache healEvents []HealEvent diff --git a/go/internal/state/store_ts.go b/go/internal/state/store_ts.go index 9a099ac5..2b3cadd0 100644 --- a/go/internal/state/store_ts.go +++ b/go/internal/state/store_ts.go @@ -489,7 +489,9 @@ func (s *Store) recordHistoryBatches(ctx context.Context, batches []historyBatch prep = append(prep, prepared{b: b, p: p, rs: rs}) } - s.historyWriteMu.Lock() + if err := lockContext(ctx, s.historyWriteMu.TryLock); err != nil { + return out, err + } defer s.historyWriteMu.Unlock() tx, err := s.history.BeginTx(ctx, nil) if err != nil {