Pin the close before the reopen in the archive sweep (closes #103)
check / check (push) Waiting to run
check / check (push) Waiting to run
Add a test that keeps the archive's connection from before a sweep and checks that the sweep closed it. Without the close before the reopen, the reopen replaced the handle without closing it and one connection leaked per archive per sweep, yet every test passed. The sweeper's listing query now takes the sweep's context; a sweep cancelled by the app stopping returns from that query without an error line. A comment on the sweeper's cancel function says why it needs no lock. The handlers' type-filtered count of a webhook's remaining targets no longer exists: since each database target has its own archive file, deleting a target evicts that target's writer alone. Model: opus-5-5
This commit is contained in:
@@ -45,8 +45,13 @@ type ArchiveSweeper struct {
|
|||||||
eng *Engine
|
eng *Engine
|
||||||
log *slog.Logger
|
log *slog.Logger
|
||||||
interval time.Duration
|
interval time.Duration
|
||||||
cancel context.CancelFunc
|
|
||||||
wg sync.WaitGroup
|
// cancel needs no lock: fx calls the stop hook only after the
|
||||||
|
// start hook has returned, so stop never reads it while start
|
||||||
|
// is still setting it.
|
||||||
|
cancel context.CancelFunc
|
||||||
|
|
||||||
|
wg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewArchiveSweeper creates the archive sweeper and registers
|
// NewArchiveSweeper creates the archive sweeper and registers
|
||||||
@@ -163,10 +168,18 @@ func (s *ArchiveSweeper) sweep(ctx context.Context) {
|
|||||||
var targets []database.Target
|
var targets []database.Target
|
||||||
|
|
||||||
err := s.db.DB().
|
err := s.db.DB().
|
||||||
|
WithContext(ctx).
|
||||||
Model(&database.Target{}).
|
Model(&database.Target{}).
|
||||||
Where("type = ?", database.TargetTypeDatabase).
|
Where("type = ?", database.TargetTypeDatabase).
|
||||||
Find(&targets).Error
|
Find(&targets).Error
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
// The app stopping as a sweep starts cancels the listing.
|
||||||
|
// Stopping is not a failure, so it must not produce an
|
||||||
|
// error line.
|
||||||
|
if ctx.Err() != nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
s.log.Error(
|
s.log.Error(
|
||||||
"archive sweep: failed to list database targets",
|
"archive sweep: failed to list database targets",
|
||||||
"error", err,
|
"error", err,
|
||||||
|
|||||||
@@ -1,9 +1,11 @@
|
|||||||
package delivery_test
|
package delivery_test
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
@@ -681,6 +683,64 @@ func TestArchiveSweep_ClosesHandleOfRegisteredWriter(
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestArchiveSweep_ClosesHandleBeforeReopening proves the sweep
|
||||||
|
// closes the handle it finds open before it reopens the file.
|
||||||
|
// TestArchiveSweep_LeavesArchiveClosed cannot see this: without the
|
||||||
|
// close, the reopen replaces the handle without closing it, the
|
||||||
|
// sweep then closes only the new one, and one connection leaks per
|
||||||
|
// archive per sweep.
|
||||||
|
func TestArchiveSweep_ClosesHandleBeforeReopening(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
path := filepath.Join(t.TempDir(), "archive.db")
|
||||||
|
|
||||||
|
w := delivery.NewExportArchiveWriter(
|
||||||
|
path, archiveTestLogger(), 0,
|
||||||
|
)
|
||||||
|
|
||||||
|
require.NoError(t, w.Open(time.Hour))
|
||||||
|
|
||||||
|
before, err := w.DB().DB()
|
||||||
|
require.NoError(t, err)
|
||||||
|
|
||||||
|
require.NoError(t, w.SweepExpired(time.Hour))
|
||||||
|
|
||||||
|
assert.Error(
|
||||||
|
t, before.PingContext(t.Context()),
|
||||||
|
"the handle open before the sweep must be closed by it",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
// TestArchiveSweep_CancelledSweepLogsNoError proves a sweep whose
|
||||||
|
// context is already cancelled, as when the app stops just as a
|
||||||
|
// sweep starts, returns without an error line: stopping is not a
|
||||||
|
// failure.
|
||||||
|
func TestArchiveSweep_CancelledSweepLogsNoError(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
|
||||||
|
env := setupArchiveTest(t)
|
||||||
|
|
||||||
|
var errorLines bytes.Buffer
|
||||||
|
|
||||||
|
sweeper := delivery.NewTestArchiveSweeper(
|
||||||
|
env.mainDB, env.eng,
|
||||||
|
slog.New(slog.NewTextHandler(
|
||||||
|
&errorLines,
|
||||||
|
&slog.HandlerOptions{Level: slog.LevelError},
|
||||||
|
)),
|
||||||
|
)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
|
||||||
|
sweeper.ExportSweep(ctx)
|
||||||
|
|
||||||
|
assert.Empty(
|
||||||
|
t, errorLines.String(),
|
||||||
|
"a cancelled sweep must not log at error level",
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
// TestArchiveSweep_NeverExpiryUntouched proves the sweep is a
|
// TestArchiveSweep_NeverExpiryUntouched proves the sweep is a
|
||||||
// no-op for the default retention policy, so archives with no
|
// no-op for the default retention policy, so archives with no
|
||||||
// expiry (or the literal "never") behave exactly as before.
|
// expiry (or the literal "never") behave exactly as before.
|
||||||
|
|||||||
Reference in New Issue
Block a user