Pin the close before the reopen in the archive sweep (closes #103) #445
@@ -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