Name each database target's archive for its webhook and target (closes #376)
check / check (push) Waiting to run

Each database target now writes its own archive file, archive-WEBHOOKNAME-TARGETNAME-TARGETID.db, named by delivery.ArchiveFileName, in place of one archive per webhook keyed on its UUID. Renaming a webhook or a target renames its archive files (with any -wal and -shm) before the new name is saved, never over an existing file, and moves every one back if a rename or the save fails. The webhook edit, the target edit and target creation share one lock so no two interleave. Deleting a webhook or target leaves its files on disk. Nothing looks for the old archive-WEBHOOKID.db files. The README gives the naming and the recovery steps.

Model: opus-5-5
This commit was merged in pull request #418.
This commit is contained in:
2026-10-02 13:28:49 +02:00
parent 21aafbf928
commit fd036774f9
21 changed files with 1825 additions and 663 deletions
+14 -12
View File
@@ -8,6 +8,7 @@ import (
"time"
"go.uber.org/fx"
"gorm.io/gorm"
"sneak.berlin/go/webhooker/internal/config"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/lifecycle"
@@ -25,14 +26,14 @@ type ArchiveSweeperParams struct {
Logger *logger.Logger
}
// ArchiveSweeper periodically prunes expired rows from
// per-webhook archive databases whose database target carries a
// positive expiry.
// ArchiveSweeper periodically prunes expired rows from the
// archive databases of database targets that carry a positive
// expiry.
//
// Without it, pruning happens only when an archive is
// (re)opened, and archives are only ever reopened by writes: an
// archive belonging to a webhook that has stopped receiving
// events would keep its expired rows forever. The sweep closes
// archive whose target has stopped receiving events would keep
// its expired rows forever. The sweep closes
// that gap without changing anything for archives whose expiry
// is unset or "never".
//
@@ -155,7 +156,7 @@ func (s *ArchiveSweeper) run(ctx context.Context) {
// soft-deleted along with it, so GORM's default scope already
// excludes them.
//
// A failure for one webhook is logged and the sweep continues,
// A failure for one target is logged and the sweep continues,
// matching how the write path already treats a prune error as
// non-fatal.
func (s *ArchiveSweeper) sweep(ctx context.Context) {
@@ -210,19 +211,20 @@ func (s *ArchiveSweeper) sweepTarget(target *database.Target) {
return
}
err = s.eng.dbTarget.sweepWebhook(target.WebhookID, expiry)
err = s.eng.dbTarget.sweepArchive(target.ID, expiry)
if err == nil {
return
}
// A writer evicted underneath the sweep means the operator
// deleted the webhook (or its last database target) while the
// sweep was walking the target list. That is an ordinary
// A writer evicted, or a target row gone, underneath the sweep
// means the operator deleted the target or its webhook while
// the sweep was walking the target list. That is an ordinary
// interleaving, not a failure, so it must not produce an
// error line.
if errors.Is(err, errArchiveWriterEvicted) {
if errors.Is(err, errArchiveWriterEvicted) ||
errors.Is(err, gorm.ErrRecordNotFound) {
s.log.Debug(
"archive sweep: writer evicted mid-sweep",
"archive sweep: target deleted mid-sweep",
"webhook_id", target.WebhookID,
"target_id", target.ID,
)
+143 -133
View File
@@ -34,18 +34,23 @@ const (
sweepConcurrentWrites = 20
)
// sweeperEnv bundles the pieces an archive sweep test drives:
// a main configuration database holding webhooks and targets, a
// delivery engine owning the archive writer registry, and the
// data directory the archive files live in.
type sweeperEnv struct {
// archiveTestWebhookName is the name of every webhook
// seedDatabaseTarget creates. It is not safe in a file name as it
// stands, so every archive test goes through archiveNamePart.
const archiveTestWebhookName = "Sweep Test!"
// archiveEnv bundles the pieces an archive test drives: a main
// configuration database holding webhooks and targets, a delivery
// engine owning the archive writer registry, the archive sweeper,
// and the data directory the archive files live in.
type archiveEnv struct {
sweeper *delivery.ArchiveSweeper
eng *delivery.Engine
mainDB *database.Database
dataDir string
}
func setupSweeperTest(t *testing.T) *sweeperEnv {
func setupArchiveTest(t *testing.T) *archiveEnv {
t.Helper()
dataDir := t.TempDir()
@@ -78,7 +83,7 @@ func setupSweeperTest(t *testing.T) *sweeperEnv {
1,
)
return &sweeperEnv{
return &archiveEnv{
sweeper: delivery.NewTestArchiveSweeper(
mainDB, eng, log,
),
@@ -88,25 +93,27 @@ func setupSweeperTest(t *testing.T) *sweeperEnv {
}
}
// archivePath returns where the engine keeps a webhook's
// archive file.
func (env *sweeperEnv) archivePath(webhookID string) string {
// archivePath returns where the engine keeps a database target's
// archive file, for the names seedDatabaseTarget gave it.
func (env *archiveEnv) archivePath(tgt *database.Target) string {
return filepath.Join(
env.dataDir, fmt.Sprintf("archive-%s.db", webhookID),
env.dataDir,
delivery.ArchiveFileName(
archiveTestWebhookName, tgt.Name, tgt.ID,
),
)
}
// seedDatabaseTarget creates a webhook with one database target
// carrying the given target config JSON, and returns the
// webhook id.
func (env *sweeperEnv) seedDatabaseTarget(
// carrying the given target config JSON, and returns the target.
func (env *archiveEnv) seedDatabaseTarget(
t *testing.T, configJSON string,
) string {
) *database.Target {
t.Helper()
wh := &database.Webhook{
UserID: uuid.New().String(),
Name: "sweep-test",
Name: archiveTestWebhookName,
}
require.NoError(
t,
@@ -115,9 +122,19 @@ func (env *sweeperEnv) seedDatabaseTarget(
Create(wh).Error,
)
return env.addDatabaseTarget(t, wh.ID, configJSON)
}
// addDatabaseTarget creates one more database target on an
// existing webhook and returns it.
func (env *archiveEnv) addDatabaseTarget(
t *testing.T, webhookID, configJSON string,
) *database.Target {
t.Helper()
tgt := &database.Target{
WebhookID: wh.ID,
Name: "archive",
WebhookID: webhookID,
Name: "Archive",
Type: database.TargetTypeDatabase,
Active: true,
Config: configJSON,
@@ -129,19 +146,19 @@ func (env *sweeperEnv) seedDatabaseTarget(
Create(tgt).Error,
)
return wh.ID
return tgt
}
// seedArchiveRows creates the archive file for a webhook and
// seedArchiveRows creates the archive file for a target and
// inserts one row per supplied archived-at timestamp, returning
// the archive path. The handle is closed before returning, so
// the archive is idle exactly as it would be with no traffic.
func (env *sweeperEnv) seedArchiveRows(
t *testing.T, webhookID string, archivedAt ...time.Time,
func (env *archiveEnv) seedArchiveRows(
t *testing.T, tgt *database.Target, archivedAt ...time.Time,
) string {
t.Helper()
path := env.archivePath(webhookID)
path := env.archivePath(tgt)
sqlDB, err := sql.Open(
"sqlite", fmt.Sprintf("file:%s?mode=rwc", path),
@@ -160,7 +177,7 @@ func (env *sweeperEnv) seedArchiveRows(
for i, at := range archivedAt {
row := delivery.ExportArchivedEvent{
EventID: fmt.Sprintf("ev-%d", i),
WebhookID: webhookID,
WebhookID: tgt.WebhookID,
Method: http.MethodPost,
Body: `{"seeded":true}`,
ArchivedAt: at,
@@ -243,13 +260,13 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext(
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
now := time.Now()
path := env.seedArchiveRows(
t, webhookID,
t, tgt,
now.Add(-48*time.Hour),
now.Add(-time.Minute),
)
@@ -287,60 +304,60 @@ func TestArchiveSweeper_LoopOutlivesStartHookContext(
}
// TestArchiveSweep_DoesNotResurrectEvictedWriter covers the
// interleaving where a sweep tick has already listed a webhook's
// target when the webhook is deleted and its writer evicted. The
// sweep must not put a writer back into the registry: nothing
// would ever evict it again, which is precisely the leak this
// change exists to close.
// interleaving where a sweep tick has already listed a target
// when the target is deleted and its writer evicted. The sweep
// must not put a writer back into the registry: nothing would
// ever evict it again, which is precisely the leak this change
// exists to close.
func TestArchiveSweep_DoesNotResurrectEvictedWriter(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
t, tgt, time.Now().Add(-48*time.Hour),
)
// Prime the registry the way a delivery would, then evict as
// the deletion path does. The target row is deliberately left
// in place: this is the tick that listed the webhook before
// in place: this is the tick that listed the target before
// the deletion committed.
_, err := env.eng.ExportEnsureArchiveWriter(webhookID)
_, err := env.eng.ExportEnsureArchiveWriter(tgt.ID)
require.NoError(t, err)
env.eng.EvictWebhook(webhookID)
require.False(t, env.eng.ExportHasArchiveWriter(webhookID))
env.eng.EvictTarget(tgt.ID)
require.False(t, env.eng.ExportHasArchiveWriter(tgt.ID))
env.sweeper.ExportSweep(context.Background())
assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID),
"a sweep must never re-register a writer for a webhook "+
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a sweep must never re-register a writer for a target "+
"whose registry entry has already been released",
)
}
// TestArchiveSweep_LeavesNoRegistryEntry states the same
// invariant in its general form: sweeping an archive whose
// webhook has no cached writer must not leave one behind, so the
// target has no cached writer must not leave one behind, so the
// registry keeps holding only writers a delivery created and an
// eviction can reach.
func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows(
t, webhookID,
t, tgt,
time.Now().Add(-48*time.Hour),
time.Now().Add(-time.Minute),
)
require.False(t, env.eng.ExportHasArchiveWriter(webhookID))
require.False(t, env.eng.ExportHasArchiveWriter(tgt.ID))
env.sweeper.ExportSweep(context.Background())
@@ -349,7 +366,7 @@ func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) {
"the sweep must still prune an idle archive",
)
assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the sweep must release the registry entry it created",
)
}
@@ -364,34 +381,31 @@ func TestArchiveSweep_KeepsWriterAdoptedByDelivery(
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
t, tgt, time.Now().Add(-48*time.Hour),
)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
)
d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
env.sweeper.ExportSweep(context.Background())
require.False(t, env.eng.ExportHasArchiveWriter(webhookID))
require.False(t, env.eng.ExportHasArchiveWriter(tgt.ID))
env.eng.ExportDeliverDatabase(webhookDB, d)
assert.True(
t, env.eng.ExportHasArchiveWriter(webhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a delivery's writer must stay registered",
)
env.sweeper.ExportSweep(context.Background())
assert.True(
t, env.eng.ExportHasArchiveWriter(webhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a sweep must not drop a writer a delivery owns",
)
}
@@ -423,15 +437,15 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep(
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
t, tgt, time.Now().Add(-48*time.Hour),
)
sweepWriter, created, err := env.eng.ExportSweepWriterFor(
webhookID,
tgt.ID,
)
require.NoError(t, err)
require.True(
@@ -442,37 +456,34 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep(
// The delivery lands mid-sweep and adopts the entry.
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
)
d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
env.eng.ExportDeliverDatabase(webhookDB, d)
adopted := env.eng.ExportArchiveWriterFor(webhookID)
adopted := env.eng.ExportArchiveWriterFor(tgt.ID)
require.NotNil(t, adopted)
require.True(
t, sweepWriter.Same(adopted),
"the delivery must have adopted the sweep's writer",
)
require.True(
t, env.eng.ExportArchiveHandleOpen(webhookID),
t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the delivery leaves the archive handle open",
)
// The sweep finishes.
env.eng.ExportReleaseSweepWriter(webhookID, sweepWriter)
env.eng.ExportReleaseSweepWriter(tgt.ID, sweepWriter)
require.True(
t, env.eng.ExportHasArchiveWriter(webhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a writer adopted by a delivery during a sweep must "+
"stay registered, or its open handle is unreachable",
)
env.eng.EvictWebhook(webhookID)
env.eng.EvictTarget(tgt.ID)
assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the adopted writer must still be evictable",
)
assert.False(
@@ -481,34 +492,34 @@ func TestArchiveSweep_KeepsWriterAdoptedDuringSweep(
)
}
// TestArchiveSweep_ContinuesAfterPerWebhookFailure proves a
// failure for one webhook does not abort the sweep for the
// TestArchiveSweep_ContinuesAfterPerTargetFailure proves a
// failure for one target does not abort the sweep for the
// others: an unparseable expiry and an unreadable archive both
// have to be logged and stepped over.
func TestArchiveSweep_ContinuesAfterPerWebhookFailure(
func TestArchiveSweep_ContinuesAfterPerTargetFailure(
t *testing.T,
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
// Seeded first so the sweep reaches them before the healthy
// webhook: targets come back in insertion order.
badConfigID := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`)
// target: targets come back in insertion order.
badConfig := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`)
env.seedArchiveRows(
t, badConfigID, time.Now().Add(-48*time.Hour),
t, badConfig, time.Now().Add(-48*time.Hour),
)
corruptID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
corrupt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
require.NoError(t, os.WriteFile(
env.archivePath(corruptID),
env.archivePath(corrupt),
[]byte("this is not a sqlite database"),
0o600,
))
healthyID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
healthy := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
healthyPath := env.seedArchiveRows(
t, healthyID,
t, healthy,
time.Now().Add(-48*time.Hour),
time.Now().Add(-time.Minute),
)
@@ -518,14 +529,14 @@ func TestArchiveSweep_ContinuesAfterPerWebhookFailure(
assert.Equal(
t, []string{sweepRowNew},
archivedEventIDs(t, healthyPath),
"a failure for an earlier webhook must not stop the "+
"a failure for an earlier target must not stop the "+
"sweep from pruning the ones after it",
)
}
// TestArchiveSweep_OpenExistingDoesNotCreateFile pins the second
// of the two no-create guards. The first is the stat in
// sweepWebhook; this one is the SQLite open mode, which is what
// sweepExpired; this one is the SQLite open mode, which is what
// protects the window between that stat and the open. Flipping
// the sweep's mode to create-if-missing makes this fail.
func TestArchiveSweep_OpenExistingDoesNotCreateFile(
@@ -561,13 +572,13 @@ func TestArchiveSweep_OpenExistingDoesNotCreateFile(
func TestArchiveSweep_PrunesIdleArchive(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
now := time.Now()
path := env.seedArchiveRows(
t, webhookID,
t, tgt,
now.Add(-48*time.Hour),
now.Add(-time.Minute),
)
@@ -600,11 +611,11 @@ func TestArchiveSweep_PrunesIdleArchive(t *testing.T) {
func TestArchiveSweep_LeavesArchiveClosed(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
t, tgt, time.Now().Add(-48*time.Hour),
)
w := delivery.NewExportArchiveWriter(
@@ -640,35 +651,32 @@ func TestArchiveSweep_ClosesHandleOfRegisteredWriter(
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
t, tgt, time.Now().Add(-48*time.Hour),
)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
)
d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
env.eng.ExportDeliverDatabase(webhookDB, d)
require.True(
t, env.eng.ExportArchiveHandleOpen(webhookID),
t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the delivery must leave the archive handle open",
)
env.sweeper.ExportSweep(context.Background())
require.True(
t, env.eng.ExportHasArchiveWriter(webhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the delivery's registry entry must survive the sweep",
)
assert.False(
t, env.eng.ExportArchiveHandleOpen(webhookID),
t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the sweep must leave the archive closed",
)
}
@@ -684,11 +692,11 @@ func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) {
`{"expiry":""}`,
"",
} {
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, configJSON)
tgt := env.seedDatabaseTarget(t, configJSON)
path := env.seedArchiveRows(
t, webhookID,
t, tgt,
time.Now().Add(-10000*time.Hour),
)
@@ -699,7 +707,7 @@ func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) {
"config %q must keep rows forever", configJSON,
)
assert.False(
t, env.eng.ExportHasArchiveWriter(webhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"config %q must leave no registry entry behind",
configJSON,
)
@@ -722,10 +730,10 @@ func TestArchiveSweep_NeverExpirySkipsBeforeOpening(
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"never"}`)
path := env.archivePath(webhookID)
tgt := env.seedDatabaseTarget(t, `{"expiry":"never"}`)
path := env.archivePath(tgt)
seedUnmigratedArchive(t, path)
require.False(t, archiveTableExists(t, path))
@@ -768,16 +776,16 @@ func archiveTableExists(t *testing.T, path string) bool {
}
// TestArchiveSweep_DoesNotCreateArchiveFile proves the sweep
// never conjures an archive: a webhook with a database target
// that has never received an event must still have no archive
// file (nor SQLite sidecar) after a sweep.
// never conjures an archive: a database target that has never
// received an event must still have no archive file (nor SQLite
// sidecar) after a sweep, and no registry entry either.
func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.archivePath(webhookID)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.archivePath(tgt)
require.NoFileExists(t, path)
@@ -789,6 +797,11 @@ func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) {
"the sweep must not create an archive file",
)
}
assert.False(
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the sweep must leave no registry entry behind",
)
}
// TestArchiveSweep_DoesNotCreateAfterWriterExists covers the
@@ -800,11 +813,11 @@ func TestArchiveSweep_DoesNotCreateAfterWriterExists(
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path, err := env.eng.ExportEnsureArchiveWriter(webhookID)
path, err := env.eng.ExportEnsureArchiveWriter(tgt.ID)
require.NoError(t, err)
require.NoFileExists(t, path)
@@ -819,17 +832,17 @@ func TestArchiveSweep_DoesNotCreateAfterWriterExists(
func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
path := env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
t, tgt, time.Now().Add(-48*time.Hour),
)
require.NoError(
t,
env.mainDB.DB().
Where("webhook_id = ?", webhookID).
Where("webhook_id = ?", tgt.WebhookID).
Delete(&database.Target{}).Error,
)
@@ -842,14 +855,14 @@ func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) {
}
// TestArchiveSweep_ConcurrentWrites proves the sweep serialises
// against writes through the per-webhook writer mutex. Run
// under -race, an unsynchronised sweep would be caught here.
// against writes through the target's writer mutex. Run under
// -race, an unsynchronised sweep would be caught here.
func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
webhookDB := testWebhookDB(t)
@@ -862,13 +875,10 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
for range sweepConcurrentWrites {
event := seedEvent(t, webhookDB, `{"n":1}`)
event.WebhookID = webhookID
deliveries = append(
deliveries,
seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"1h"}`,
),
seedDatabaseTargetDelivery(t, webhookDB, event, tgt),
)
}
@@ -894,7 +904,7 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
wg.Wait()
assert.FileExists(t, env.archivePath(webhookID))
assert.FileExists(t, env.archivePath(tgt))
}
// TestArchiveSweeper_StopsCleanly proves the background loop
@@ -902,11 +912,11 @@ func TestArchiveSweep_ConcurrentWrites(t *testing.T) {
func TestArchiveSweeper_StopsCleanly(t *testing.T) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
tgt := env.seedDatabaseTarget(t, `{"expiry":"1h"}`)
env.seedArchiveRows(
t, webhookID, time.Now().Add(-48*time.Hour),
t, tgt, time.Now().Add(-48*time.Hour),
)
env.sweeper.ExportSetInterval(time.Millisecond)
@@ -930,7 +940,7 @@ func TestArchiveSweeper_StopHookHonoursStopTimeout(
) {
t.Parallel()
env := setupSweeperTest(t)
env := setupArchiveTest(t)
lc := &recordingLifecycle{}
env.sweeper.ExportRegisterHooks(lc)
+51 -20
View File
@@ -123,21 +123,24 @@ type Notifier interface {
Notify(tasks []Task)
}
// WebhookEvictor releases the delivery engine's per-webhook
// state for a webhook that no longer needs it — currently the
// cached archive writer of the database target, whose open
// file handle would otherwise outlive the webhook.
// Archives is how the handlers keep the database targets' archive
// files in step with the configuration. Deleting a webhook or a
// target releases the cached archive writers, whose open file
// handles would otherwise outlive them; renaming one renames the
// archive files, which are named for the webhook and the target
// (see ArchiveFileName).
//
// It is deliberately separate from Notifier and deliberately
// one method wide: archiving lifecycle is not notification, and
// a single-method interface keeps the handlers package free of
// any dependency on the engine's internals while staying
// trivially fakeable in tests.
// It is deliberately separate from Notifier: archiving lifecycle
// is not notification, and a small interface keeps the handlers
// package free of any dependency on the engine's internals while
// staying trivially fakeable in tests.
//
// EvictWebhook never deletes an archive file. It is idempotent
// and is a no-op for a webhook with no engine state.
type WebhookEvictor interface {
// Neither eviction deletes an archive file. Both are idempotent
// and are no-ops for a webhook or target with no engine state.
type Archives interface {
EvictWebhook(webhookID string)
EvictTarget(targetID string)
Rename(targetID, webhookName, targetName string) error
}
// EngineParams are the fx dependencies for the delivery
@@ -188,7 +191,7 @@ type Engine struct {
httpTarget *httpTarget
// dbTarget is retained so the engine can reach the archive
// writer registry for webhook eviction and the idle sweep.
// writer registry for eviction, renames and the idle sweep.
dbTarget *databaseTarget
// inflight is the set of deliveries this engine currently owns.
@@ -257,17 +260,44 @@ func (e *Engine) Notify(tasks []Task) {
}
}
// EvictWebhook implements WebhookEvictor. It releases the
// engine's per-webhook archiving state: the database target's
// cached archive writer is dropped from the registry and its
// file handle closed. The archive file itself is left on disk
// — it is long-term storage the operator owns.
// EvictWebhook implements Archives. The cached archive writer of
// every database target of the webhook is dropped from the
// registry and its file handle closed. The archive files
// themselves are left on disk — they are long-term storage the
// operator owns.
func (e *Engine) EvictWebhook(webhookID string) {
if e.dbTarget == nil {
return
}
e.dbTarget.evict(webhookID)
e.dbTarget.evictWebhook(webhookID)
}
// EvictTarget implements Archives. It is EvictWebhook for a single
// database target, and leaves the archive file on disk the same
// way.
func (e *Engine) EvictTarget(targetID string) {
if e.dbTarget == nil {
return
}
e.dbTarget.evict(targetID)
}
// Rename implements Archives. It renames a database target's
// archive file to ArchiveFileName(webhookName, targetName,
// targetID), under the lock the target's archive writes and the
// idle sweep take. It never replaces a file: if one already has the
// new name, the error is ErrArchiveNameTaken. The caller renames
// before it saves the new name: see databaseTarget.rename.
func (e *Engine) Rename(
targetID, webhookName, targetName string,
) error {
if e.dbTarget == nil {
return nil
}
return e.dbTarget.rename(targetID, webhookName, targetName)
}
// ScheduleRetry schedules a task to be re-enqueued onto the
@@ -381,7 +411,8 @@ func (e *Engine) start() {
// Once the pool has drained it closes the archive writers, so a
// clean stop leaves no archive -wal behind. Nothing else holds a
// writer for long by then: the archive sweeper stops before the
// engine, and deleting a webhook only closes one. If the pool did
// engine, and deleting or renaming a webhook or target only closes
// or moves one. If the pool did
// not drain in time, the writers are left open, as a kill would
// leave them. Closing them would wait for any write in progress,
// and a worker still running would then open new writers that
+23 -10
View File
@@ -2,7 +2,6 @@ package delivery_test
import (
"context"
"fmt"
"path/filepath"
"testing"
"time"
@@ -10,6 +9,7 @@ import (
"github.com/google/uuid"
"github.com/stretchr/testify/require"
"go.uber.org/fx"
"gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
)
@@ -272,22 +272,35 @@ func TestEngine_StopHookHonoursStopTimeout(t *testing.T) {
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
}
// deliverToArchive runs one delivery to a database target through
// the running engine and returns the webhook's archive file path.
// The archive writer holds the file open afterwards.
func deliverToArchive(t *testing.T, s iSetup) string {
// deliverToArchive gives the setup's webhook a database target,
// runs one delivery to it through the running engine, and returns
// the target's ID and archive file path. The archive writer holds
// the file open afterwards.
func deliverToArchive(t *testing.T, s iSetup) (string, string) {
t.Helper()
iCreateWebhook(t, s.MainDB, s.WebhookID, "hook")
tgt := &database.Target{
WebhookID: s.WebhookID,
Name: "archive",
Type: database.TargetTypeDatabase,
}
require.NoError(
t, s.MainDB.Omit(clause.Associations).Create(tgt).Error,
)
deliveryID, task := seedLogTask(t, s)
task.TargetID = tgt.ID
task.TargetType = database.TargetTypeDatabase
s.Engine.Notify([]delivery.Task{task})
iWaitForDelivered(t, s.WebhookDB, deliveryID)
return filepath.Join(
return tgt.ID, filepath.Join(
filepath.Dir(s.DBMgr.DBPath(s.WebhookID)),
fmt.Sprintf("archive-%s.db", s.WebhookID),
"archive-hook-archive-"+tgt.ID+".db",
)
}
@@ -304,7 +317,7 @@ func TestEngine_StopHookClosesArchives(t *testing.T) {
lc := startEngineViaHook(t, s.Engine)
path := deliverToArchive(t, s)
_, path := deliverToArchive(t, s)
require.FileExists(
t, path+"-wal",
"an open archive should have a -wal for the stop to remove",
@@ -338,7 +351,7 @@ func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) {
lc := startEngineViaHook(t, s.Engine)
deliverToArchive(t, s)
targetID, _ := deliverToArchive(t, s)
release := make(chan struct{})
@@ -352,7 +365,7 @@ func TestEngine_StopHookTimeoutLeavesArchivesOpen(t *testing.T) {
requireStopHookExpires(t, lc.hooks[0], "delivery engine")
require.True(
t, s.Engine.ExportArchiveHandleOpen(s.WebhookID),
t, s.Engine.ExportArchiveHandleOpen(targetID),
"a stop that timed out must not close archive writers",
)
}
+16 -29
View File
@@ -351,23 +351,15 @@ func TestDeliverDatabase_ImmediateSuccess(
db := testWebhookDB(t)
// The database target archives for real now, so the engine
// needs a webhook DB manager to locate the data directory.
e := delivery.NewTestEngineWithDB(
nil,
database.NewTestWebhookDBManager(t.TempDir()),
slog.New(slog.NewTextHandler(
os.Stderr,
&slog.HandlerOptions{Level: slog.LevelDebug},
)),
&http.Client{Timeout: 5 * time.Second},
1,
)
// The database target archives for real, so the engine needs
// the target in the main database and a data directory.
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
event := seedEvent(t, db, `{"db":"target"}`)
d := seedDatabaseTargetDelivery(t, db, event, "")
d := seedDatabaseTargetDelivery(t, db, event, tgt)
e.ExportDeliverDatabase(db, d)
env.eng.ExportDeliverDatabase(db, d)
var updated database.Delivery
@@ -1332,32 +1324,27 @@ func TestProcessDelivery_RoutesToCorrectHandler(
db := testWebhookDB(t)
// The database target archives for real now, so the engine
// needs a webhook DB manager to locate the data directory.
e := delivery.NewTestEngineWithDB(
nil,
database.NewTestWebhookDBManager(t.TempDir()),
slog.New(slog.NewTextHandler(
os.Stderr,
&slog.HandlerOptions{Level: slog.LevelDebug},
)),
&http.Client{Timeout: 5 * time.Second},
1,
)
// The database target archives for real, so the engine needs
// the target in the main database and a data directory.
env := setupArchiveTest(t)
archive := env.seedDatabaseTarget(t, "")
tests := []struct {
name string
targetType database.TargetType
targetID string
wantStatus database.DeliveryStatus
}{
{
"database target",
database.TargetTypeDatabase,
archive.ID,
database.DeliveryStatusDelivered,
},
{
"log target",
database.TargetTypeLog,
uuid.New().String(),
database.DeliveryStatusDelivered,
},
}
@@ -1367,7 +1354,7 @@ func TestProcessDelivery_RoutesToCorrectHandler(
t.Parallel()
runRoutingSubtest(
t, db, e, tt.targetType,
t, db, env.eng, tt.targetType, tt.targetID,
tt.wantStatus,
)
})
@@ -1379,6 +1366,7 @@ func runRoutingSubtest(
db *gorm.DB,
e *delivery.Engine,
targetType database.TargetType,
targetID string,
wantStatus database.DeliveryStatus,
) {
t.Helper()
@@ -1386,8 +1374,7 @@ func runRoutingSubtest(
event := seedEvent(t, db, `{"routing":"test"}`)
dlv := seedDelivery(
t, db, event.ID,
uuid.New().String(),
t, db, event.ID, targetID,
database.DeliveryStatusPending,
)
+26 -20
View File
@@ -474,7 +474,7 @@ func NewTestCircuitBreaker(
type ExportArchivedEvent = archivedEvent
// ExportArchiveWriter wraps an archiveWriter so black-box tests
// can exercise the per-webhook archive file mechanics.
// can exercise the archive file mechanics.
type ExportArchiveWriter struct {
w *archiveWriter
}
@@ -549,6 +549,12 @@ func (e *ExportArchiveWriter) Evict() {
e.w.evict()
}
// Rename gives the archive file a new name in the same directory,
// as a rename of the webhook or target does.
func (e *ExportArchiveWriter) Rename(name string) error {
return e.w.rename(name)
}
// HandleOpen reports whether the writer currently holds an open
// archive handle.
func (e *ExportArchiveWriter) HandleOpen() bool {
@@ -568,16 +574,16 @@ func (e *ExportArchiveWriter) Same(
}
// ExportArchiveWriterFor returns the archive writer the registry
// currently caches for a webhook, or nil when none is cached. It
// never creates one, so a test can hold a reference to the very
// writer an eviction is about to detach.
// currently caches for a database target, or nil when none is
// cached. It never creates one, so a test can hold a reference to
// the very writer an eviction is about to detach.
func (e *Engine) ExportArchiveWriterFor(
webhookID string,
targetID string,
) *ExportArchiveWriter {
e.dbTarget.mu.Lock()
defer e.dbTarget.mu.Unlock()
w, ok := e.dbTarget.writers[webhookID]
w, ok := e.dbTarget.writers[targetID]
if !ok {
return nil
}
@@ -586,26 +592,26 @@ func (e *Engine) ExportArchiveWriterFor(
}
// ExportHasArchiveWriter reports whether the database target
// currently caches an archive writer for a webhook.
// type currently caches an archive writer for a target.
func (e *Engine) ExportHasArchiveWriter(
webhookID string,
targetID string,
) bool {
e.dbTarget.mu.Lock()
defer e.dbTarget.mu.Unlock()
_, ok := e.dbTarget.writers[webhookID]
_, ok := e.dbTarget.writers[targetID]
return ok
}
// ExportArchiveHandleOpen reports whether the cached archive
// writer for a webhook holds an open database handle. It
// writer for a target holds an open database handle. It
// returns false when no writer is cached.
func (e *Engine) ExportArchiveHandleOpen(
webhookID string,
targetID string,
) bool {
e.dbTarget.mu.Lock()
w, ok := e.dbTarget.writers[webhookID]
w, ok := e.dbTarget.writers[targetID]
e.dbTarget.mu.Unlock()
if !ok {
@@ -619,12 +625,12 @@ func (e *Engine) ExportArchiveHandleOpen(
}
// ExportEnsureArchiveWriter creates (if needed) and returns the
// archive file path of the cached writer for a webhook, so a
// archive file path of the cached writer for a target, so a
// test can prime the registry the way a delivery would.
func (e *Engine) ExportEnsureArchiveWriter(
webhookID string,
targetID string,
) (string, error) {
w, err := e.dbTarget.writerFor(webhookID)
w, err := e.dbTarget.writerFor(targetID)
if err != nil {
return "", err
}
@@ -632,14 +638,14 @@ func (e *Engine) ExportEnsureArchiveWriter(
return w.path, nil
}
// ExportSweepWriterFor takes a webhook's registry writer exactly
// ExportSweepWriterFor takes a target's registry writer exactly
// as the idle sweep does, reporting whether the sweep had to
// create the entry. It lets a test drive the registry through the
// sweep's own entry point instead of choreographing goroutines.
func (e *Engine) ExportSweepWriterFor(
webhookID string,
targetID string,
) (*ExportArchiveWriter, bool, error) {
w, created, err := e.dbTarget.sweepWriterFor(webhookID)
w, created, err := e.dbTarget.sweepWriterFor(targetID)
if err != nil {
return nil, false, err
}
@@ -650,9 +656,9 @@ func (e *Engine) ExportSweepWriterFor(
// ExportReleaseSweepWriter releases a sweep-created registry entry
// exactly as a finished sweep does.
func (e *Engine) ExportReleaseSweepWriter(
webhookID string, w *ExportArchiveWriter,
targetID string, w *ExportArchiveWriter,
) {
e.dbTarget.releaseSweepWriter(webhookID, w.w)
e.dbTarget.releaseSweepWriter(targetID, w.w)
}
// NewTestArchiveSweeper builds an ArchiveSweeper backed by the
+200 -93
View File
@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"path/filepath"
"strings"
"sync"
"time"
@@ -11,22 +12,75 @@ import (
"sneak.berlin/go/webhooker/internal/database"
)
// databaseTarget is a no-retry target that archives the
// full inbound event into a per-webhook archive SQLite file,
// separate from the per-webhook event database. The event is
// already persisted in the per-webhook event DB by the time
// delivery runs; the database target additionally writes a
// durable long-term copy into archive-{webhookID}.db and then
// records a single attempt whose outcome reflects whether the
// archive write succeeded. See archiveWriter for the
// close/reopen, auto-recreate, and expiry semantics.
// archiveNameMaxLen is how many characters of a webhook or target
// name an archive file name keeps.
const archiveNameMaxLen = 40
// databaseTarget is a no-retry target that archives the full
// inbound event into the target's own archive SQLite file, separate
// from the per-webhook event database. The event is already
// persisted in the per-webhook event DB by the time delivery runs;
// the database target additionally writes a durable long-term copy
// into the file ArchiveFileName names and then records a single
// attempt whose outcome reflects whether the archive write
// succeeded. See archiveWriter for the close/reopen, auto-recreate,
// and expiry semantics.
type databaseTarget struct {
eng *Engine
// writers holds one archive writer per database target, keyed
// by target ID.
mu sync.Mutex
writers map[string]*archiveWriter
}
// ArchiveFileName returns the file name of a database target's
// archive: archive-WEBHOOKNAME-TARGETNAME-TARGETID.db, with both
// names passed through archiveNamePart. The target ID keeps the
// name unique when two targets' names come out the same.
func ArchiveFileName(webhookName, targetName, targetID string) string {
return "archive-" + archiveNamePart(webhookName) + "-" +
archiveNamePart(targetName) + "-" + targetID + ".db"
}
// archiveNamePart makes a webhook or target name safe to put in a
// file name. It is lowercased; ASCII letters and digits are kept,
// every other run of characters becomes a single "-", and no "-" is
// left at either end. It is cut to archiveNameMaxLen characters, and
// a name with nothing left is "unnamed".
func archiveNamePart(name string) string {
var b strings.Builder
dash := false
for _, r := range strings.ToLower(name) {
if (r < 'a' || r > 'z') && (r < '0' || r > '9') {
dash = b.Len() > 0
continue
}
if dash {
b.WriteByte('-')
dash = false
}
b.WriteRune(r)
}
part := b.String()
if len(part) > archiveNameMaxLen {
part = strings.TrimRight(part[:archiveNameMaxLen], "-")
}
if part == "" {
return "unnamed"
}
return part
}
// Deliver implements Target. It archives the event, then
// records one successful attempt and marks the delivery
// delivered. An archiving error fails the delivery: the
@@ -92,7 +146,7 @@ func (t *databaseTarget) Deliver(
)
}
// archive writes the full event as a row into the webhook's
// archive writes the full event as a row into the target's
// archive database, honouring the optional per-target expiry
// parsed from the target config JSON.
func (t *databaseTarget) archive(d *database.Delivery) error {
@@ -106,7 +160,7 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
return err
}
w, err := t.writerFor(webhookID)
w, err := t.writerFor(d.TargetID)
if err != nil {
return err
}
@@ -124,30 +178,31 @@ func (t *databaseTarget) archive(d *database.Delivery) error {
return w.write(row, expiry)
}
// writerFor returns the archiveWriter for a webhook, creating
// and caching it on first use. Each webhook has one writer so
// its close/reopen debounce state is shared across concurrent
// deliveries. The archive file lives beside the per-webhook
// event database in the data directory.
// writerFor returns the archive writer for a database target,
// creating and caching it on first use. Each target has one writer
// so its close/reopen debounce state is shared across concurrent
// deliveries, and so a rename and the idle sweep take the same lock
// as its writes.
func (t *databaseTarget) writerFor(
webhookID string,
targetID string,
) (*archiveWriter, error) {
path, err := t.archivePath(webhookID)
if err != nil {
return nil, err
}
t.mu.Lock()
defer t.mu.Unlock()
if t.writers == nil {
t.writers = make(map[string]*archiveWriter)
}
w, ok := t.writers[webhookID]
w, ok := t.writers[targetID]
if !ok {
w = newArchiveWriter(path, t.eng.log)
t.writers[webhookID] = w
var err error
w, err = t.newWriter(targetID)
if err != nil {
return nil, err
}
if t.writers == nil {
t.writers = make(map[string]*archiveWriter)
}
t.writers[targetID] = w
}
// A delivery claims the entry: even if the idle sweep created
@@ -159,40 +214,39 @@ func (t *databaseTarget) writerFor(
}
// sweepWriterFor returns the archive writer the idle sweep should
// prune a webhook through, together with whether the sweep itself
// created the registry entry.
// prune a target's archive through, together with whether the sweep
// itself created the registry entry.
//
// The sweep must route its prune through the registered writer so
// the writer's mutex orders it against concurrent writes, but it
// must never leave a registry entry behind: a sweep that ran
// concurrently with the webhook's deletion would otherwise
// concurrently with the target's deletion would otherwise
// re-create an entry that nothing will ever evict again, which is
// exactly the leak eviction exists to prevent. An entry the sweep
// creates is therefore marked sweep-owned and handed back to
// releaseSweepWriter when the sweep is done.
func (t *databaseTarget) sweepWriterFor(
webhookID string,
targetID string,
) (*archiveWriter, bool, error) {
path, err := t.archivePath(webhookID)
t.mu.Lock()
defer t.mu.Unlock()
w, ok := t.writers[targetID]
if ok {
return w, false, nil
}
w, err := t.newWriter(targetID)
if err != nil {
return nil, false, err
}
t.mu.Lock()
defer t.mu.Unlock()
if t.writers == nil {
t.writers = make(map[string]*archiveWriter)
}
w, ok := t.writers[webhookID]
if ok {
return w, false, nil
}
w = newArchiveWriter(path, t.eng.log)
w.sweepOwned = true
t.writers[webhookID] = w
t.writers[targetID] = w
return w, true, nil
}
@@ -209,57 +263,95 @@ func (t *databaseTarget) sweepWriterFor(
// delivery that adopted the writer keeps a registered, evictable
// one.
func (t *databaseTarget) releaseSweepWriter(
webhookID string, w *archiveWriter,
targetID string, w *archiveWriter,
) {
t.mu.Lock()
defer t.mu.Unlock()
cur, ok := t.writers[webhookID]
cur, ok := t.writers[targetID]
if !ok || cur != w || !cur.sweepOwned {
return
}
delete(t.writers, webhookID)
delete(t.writers, targetID)
}
// archivePath returns the archive file path for a webhook: it
// lives beside the per-webhook event database in the data
// directory. It does not touch the filesystem.
func (t *databaseTarget) archivePath(
webhookID string,
) (string, error) {
// newWriter builds the writer for a database target's archive. The
// file lives beside the webhook's event database in the data
// directory and is named for the webhook and the target as the main
// database has them now; from then on only rename changes the name
// the writer uses. It does not touch the archive file.
func (t *databaseTarget) newWriter(
targetID string,
) (*archiveWriter, error) {
if t.eng.dbManager == nil {
return "", errArchiveNoDataDir
return nil, errArchiveNoDataDir
}
dir := filepath.Dir(t.eng.dbManager.DBPath(webhookID))
var target database.Target
return filepath.Join(
dir, fmt.Sprintf("archive-%s.db", webhookID),
), nil
err := t.eng.database.DB().
Preload("Webhook").
First(&target, "id = ?", targetID).Error
if err != nil {
return nil, fmt.Errorf(
"loading database target %s: %w", targetID, err,
)
}
dir := filepath.Dir(t.eng.dbManager.DBPath(target.WebhookID))
name := ArchiveFileName(
target.Webhook.Name, target.Name, target.ID,
)
w := newArchiveWriter(filepath.Join(dir, name), t.eng.log)
w.webhookID = target.WebhookID
return w, nil
}
// evict drops a webhook's archive writer from the registry and
// closes its handle, so a deleted webhook does not leave a
// writer (and an open archive handle within its debounce
// window) alive for the process lifetime.
// rename moves a database target's archive file to the name for
// webhookName and targetName. It goes through the target's writer,
// so the move holds the lock that writes and the idle sweep take,
// and later writes use the new name.
//
// The writer is created if there is none, and it stays cached. The
// handlers rename before they save the new name, so until the save
// the main database still has the old one; a delivery in that window
// must find this writer rather than build one from the old name.
func (t *databaseTarget) rename(
targetID, webhookName, targetName string,
) error {
w, err := t.writerFor(targetID)
if err != nil {
return err
}
return w.rename(ArchiveFileName(webhookName, targetName, targetID))
}
// evict drops a database target's archive writer from the registry
// and closes its handle, so a deleted target does not leave a
// writer (and an open archive handle within its debounce window)
// alive for the process lifetime.
//
// The map entry is removed under the registry lock, which is
// then released before the handle is closed under the writer's
// own lock: that ordering keeps the registry available to other
// webhooks while an in-flight write on this one drains, and
// targets while an in-flight write on this one drains, and
// closing under the writer's lock means eviction can never race
// a write.
//
// Eviction is idempotent and silent for a webhook with no
// writer, which is the common case: a webhook with no database
// target never creates one. It never deletes the archive file.
func (t *databaseTarget) evict(webhookID string) {
// Eviction is idempotent and silent for a target with no writer,
// which is the common case: only a database target that has
// received an event or been renamed has one. It never deletes the
// archive file.
func (t *databaseTarget) evict(targetID string) {
t.mu.Lock()
w, ok := t.writers[webhookID]
w, ok := t.writers[targetID]
if ok {
delete(t.writers, webhookID)
delete(t.writers, targetID)
}
t.mu.Unlock()
@@ -272,13 +364,41 @@ func (t *databaseTarget) evict(webhookID string) {
t.eng.log.Info(
"evicted archive writer",
"webhook_id", webhookID,
"target_id", targetID,
"path", w.path,
)
}
// evictWebhook evicts, exactly as evict does, the writer of every
// database target of a webhook.
func (t *databaseTarget) evictWebhook(webhookID string) {
t.mu.Lock()
var gone []*archiveWriter
for targetID, w := range t.writers {
if w.webhookID == webhookID {
delete(t.writers, targetID)
gone = append(gone, w)
}
}
t.mu.Unlock()
for _, w := range gone {
w.evict()
t.eng.log.Info(
"evicted archive writer",
"webhook_id", webhookID,
"path", w.path,
)
}
}
// evictAll evicts every cached archive writer, exactly as evict
// does for one webhook. The engine calls it at shutdown, once its
// does for one target. The engine calls it at shutdown, once its
// workers have returned. Closing the last handle on an archive
// moves the contents of its -wal into the .db and removes the
// -wal, so a clean stop leaves each archive as a single file.
@@ -295,38 +415,25 @@ func (t *databaseTarget) evictAll() {
}
}
// sweepWebhook prunes one webhook's archive of rows older than
// expiry, without requiring a write. It returns nil (nothing to
// do) when the archive file does not exist, so a sweep never
// creates an archive for a webhook that has a database target
// but has never received an event.
// sweepArchive prunes one database target's archive of rows older
// than expiry, without requiring a write. A missing archive file is
// left missing (see sweepExpired), so a sweep never creates an
// archive for a target that has never received an event.
//
// It also never leaves a registry entry behind: an entry it had
// to create to reach the writer's mutex is released again once
// the prune is done, so a sweep racing a webhook deletion cannot
// the prune is done, so a sweep racing a target deletion cannot
// resurrect the writer the eviction just dropped.
func (t *databaseTarget) sweepWebhook(
webhookID string, expiry time.Duration,
func (t *databaseTarget) sweepArchive(
targetID string, expiry time.Duration,
) error {
path, err := t.archivePath(webhookID)
if err != nil {
return err
}
// Check before taking a writer at all: a webhook whose
// archive has never been created gets no writer, no handle,
// and no file.
if !fileExists(path) {
return nil
}
w, created, err := t.sweepWriterFor(webhookID)
w, created, err := t.sweepWriterFor(targetID)
if err != nil {
return err
}
if created {
defer t.releaseSweepWriter(webhookID, w)
defer t.releaseSweepWriter(targetID, w)
}
return w.sweepExpired(expiry)
+97 -17
View File
@@ -4,8 +4,10 @@ import (
"encoding/json"
"errors"
"fmt"
"io/fs"
"log/slog"
"os"
"path/filepath"
"sync"
"time"
@@ -41,7 +43,7 @@ const (
var (
// errArchiveMissingWebhookID is returned when an event to
// archive has no webhook id to key its archive file on.
// archive has no webhook id to record in its archive row.
errArchiveMissingWebhookID = errors.New(
"cannot archive event without a webhook id",
)
@@ -61,13 +63,19 @@ var (
)
// errArchiveWriterEvicted is returned when a writer that has
// been evicted (its webhook was deleted, or its last database
// target was removed) is used again. An evicted writer is no
// longer in the registry, so reopening its file would leak a
// handle nothing owns.
// been evicted (its target or its webhook was deleted) is used
// again. An evicted writer is no longer in the registry, so
// reopening its file would leak a handle nothing owns.
errArchiveWriterEvicted = errors.New(
"archive writer has been evicted",
)
// ErrArchiveNameTaken is returned when an archive cannot be
// renamed because a file already has the new name. That file may
// be an archive with rows of its own, so it is never replaced.
ErrArchiveNameTaken = errors.New(
"a file already has the archive's new name",
)
)
// databaseTargetConfig is the optional per-target JSON config
@@ -80,7 +88,7 @@ type databaseTargetConfig struct {
}
// archivedEvent is one fully captured webhook event stored in a
// per-webhook archive database for long-term retention. It is a
// database target's archive for long-term retention. It is a
// self-contained copy — independent of the per-webhook event
// database, which may prune events under its own retention.
type archivedEvent struct {
@@ -170,8 +178,8 @@ func ValidateArchiveExpiry(expiry string) error {
return nil
}
// archiveWriter owns one per-webhook archive SQLite file. It
// serialises writes, and after each write closes and reopens
// archiveWriter owns one database target's archive SQLite file.
// It serialises writes, and after each write closes and reopens
// the file (debounced to at most once per debounce window) so
// an operator can move the file away for offline archiving. The
// next write recreates a moved or removed file, because the
@@ -187,16 +195,21 @@ type archiveWriter struct {
reopens int
// evicted marks a writer that has been removed from the
// per-webhook registry. Its handle is closed and it must
// never open the file again: nothing holds it any more, so a
// reopen would leak the handle for the process lifetime.
// registry. Its handle is closed and it must never open the
// file again: nothing holds it any more, so a reopen would
// leak the handle for the process lifetime.
evicted bool
// webhookID is the webhook the archive's target belongs to,
// so deleting the webhook can find its writers. It is set
// when the writer is created and never changes.
webhookID string
// sweepOwned marks a registry entry that the idle sweep
// created because no writer was cached for the webhook. The
// created because no writer was cached for the target. The
// sweep removes such an entry again when it is done, so a
// sweep can never leave — or resurrect — a registry entry
// for a webhook that has been deleted. A delivery that adopts
// for a target that has been deleted. A delivery that adopts
// the writer clears the flag, handing the entry to the
// registry proper.
//
@@ -385,11 +398,78 @@ func (w *archiveWriter) sweepExpired(expiry time.Duration) error {
return nil
}
// rename gives the archive file a new name in the same directory,
// and the writer uses the file under that name from now on. The
// handle is closed first, which folds the -wal into the .db; any
// -wal or -shm still beside the file (left by a crash) is moved with
// it, because SQLite finds them by name. A missing file is not an
// error: the operator may have moved it away, and the next write
// creates it under the new name.
//
// If a file already has the new name, nothing is moved and the
// error is ErrArchiveNameTaken. If one file fails to move, those
// already moved are moved back before the error is returned, so the
// archive is never split across two names.
func (w *archiveWriter) rename(name string) error {
w.mu.Lock()
defer w.mu.Unlock()
if w.evicted {
return fmt.Errorf(
"%w: %s", errArchiveWriterEvicted, w.path,
)
}
path := filepath.Join(filepath.Dir(w.path), name)
if path == w.path {
return nil
}
suffixes := []string{"", "-wal", "-shm"}
for _, suffix := range suffixes {
if fileExists(path + suffix) {
return fmt.Errorf(
"%w: %s", ErrArchiveNameTaken, name+suffix,
)
}
}
w.close()
for i, suffix := range suffixes {
err := os.Rename(w.path+suffix, path+suffix)
if err == nil || errors.Is(err, fs.ErrNotExist) {
continue
}
for _, moved := range suffixes[:i] {
backErr := os.Rename(path+moved, w.path+moved)
if backErr != nil && !errors.Is(backErr, fs.ErrNotExist) {
w.log.Error(
"failed to move archive file back",
"from", path+moved,
"to", w.path+moved,
"error", backErr,
)
}
}
return fmt.Errorf(
"renaming archive %s to %s: %w", w.path+suffix, path+suffix, err,
)
}
w.path = path
return nil
}
// evict closes the writer's handle and marks it unusable. It is
// called when the writer leaves the registry, either because the
// webhook was deleted or because its last database target was
// removed. The archive FILE is deliberately left on disk: it is
// long-term storage an operator may still want.
// called when the writer leaves the registry, because its target
// or its webhook was deleted, or at shutdown. The archive FILE is
// deliberately left on disk: it is long-term storage an operator
// may still want.
func (w *archiveWriter) evict() {
w.mu.Lock()
defer w.mu.Unlock()
+100 -90
View File
@@ -17,85 +17,109 @@ import (
"sneak.berlin/go/webhooker/internal/delivery"
)
// evictTestEngine builds an engine backed by a temporary data
// directory and returns it along with that directory.
func evictTestEngine(t *testing.T) (*delivery.Engine, string) {
// deliverTo archives one event to a database target, leaving the
// target's writer cached with its handle open.
func deliverTo(
t *testing.T, env *archiveEnv, tgt *database.Target,
) {
t.Helper()
dataDir := t.TempDir()
eng := delivery.NewTestEngineWithDB(
nil,
database.NewTestWebhookDBManager(dataDir),
archiveTestLogger(),
&http.Client{Timeout: 5 * time.Second},
1,
)
return eng, dataDir
}
// TestEvictWebhook_ClosesAndRemovesWriter proves that evicting
// a webhook drops its archive writer from the registry and
// closes the open archive handle, rather than leaving both
// alive for the process lifetime.
func TestEvictWebhook_ClosesAndRemovesWriter(t *testing.T) {
t.Parallel()
eng, dataDir := evictTestEngine(t)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
eng.ExportDeliverDatabase(webhookDB, d)
webhookID := event.WebhookID
require.True(
t, eng.ExportHasArchiveWriter(webhookID),
"a delivery should have cached an archive writer",
)
require.True(
t, eng.ExportArchiveHandleOpen(webhookID),
"the writer should hold an open handle after a write",
env.eng.ExportDeliverDatabase(
webhookDB, seedDatabaseTargetDelivery(t, webhookDB, event, tgt),
)
}
eng.EvictWebhook(webhookID)
// TestEvictWebhook_ClosesAndRemovesWriter proves that evicting
// a webhook drops the archive writers of its database targets
// from the registry and closes their open handles, rather than
// leaving them alive for the process lifetime, and leaves another
// webhook's writer alone.
func TestEvictWebhook_ClosesAndRemovesWriter(t *testing.T) {
t.Parallel()
assert.False(
t, eng.ExportHasArchiveWriter(webhookID),
"eviction should remove the registry entry",
)
assert.False(
t, eng.ExportArchiveHandleOpen(webhookID),
"eviction should close the archive handle",
)
env := setupArchiveTest(t)
first := env.seedDatabaseTarget(t, "")
second := env.addDatabaseTarget(t, first.WebhookID, "")
other := env.seedDatabaseTarget(t, "")
archivePath := filepath.Join(
dataDir, fmt.Sprintf("archive-%s.db", webhookID),
for _, tgt := range []*database.Target{first, second, other} {
deliverTo(t, env, tgt)
require.True(
t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"the writer should hold an open handle after a write",
)
}
env.eng.EvictWebhook(first.WebhookID)
for _, tgt := range []*database.Target{first, second} {
assert.False(
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"eviction should remove the registry entry",
)
assert.False(
t, env.eng.ExportArchiveHandleOpen(tgt.ID),
"eviction should close the archive handle",
)
assert.FileExists(
t, env.archivePath(tgt),
"eviction must not delete the archive file",
)
}
assert.True(
t, env.eng.ExportArchiveHandleOpen(other.ID),
"another webhook's writer must be left alone",
)
}
// TestEvictTarget_LeavesOtherTargets proves that evicting one
// database target leaves the writer of another target of the same
// webhook in place.
func TestEvictTarget_LeavesOtherTargets(t *testing.T) {
t.Parallel()
env := setupArchiveTest(t)
doomed := env.seedDatabaseTarget(t, "")
kept := env.addDatabaseTarget(t, doomed.WebhookID, "")
deliverTo(t, env, doomed)
deliverTo(t, env, kept)
env.eng.EvictTarget(doomed.ID)
assert.False(t, env.eng.ExportHasArchiveWriter(doomed.ID))
assert.FileExists(
t, archivePath,
t, env.archivePath(doomed),
"eviction must not delete the archive file",
)
assert.True(
t, env.eng.ExportArchiveHandleOpen(kept.ID),
"the other target's writer must be left alone",
)
}
// TestEvictWebhook_UnknownWebhookIsNoOp proves eviction is safe
// for the common case of a webhook that never had a database
// target, and that repeating it does not panic.
// for the common case of a webhook or target that never had an
// archive writer, and that repeating it does not panic.
func TestEvictWebhook_UnknownWebhookIsNoOp(t *testing.T) {
t.Parallel()
eng, _ := evictTestEngine(t)
env := setupArchiveTest(t)
assert.NotPanics(t, func() {
eng.EvictWebhook("no-such-webhook")
eng.EvictWebhook("no-such-webhook")
env.eng.EvictWebhook("no-such-webhook")
env.eng.EvictWebhook("no-such-webhook")
env.eng.EvictTarget("no-such-target")
env.eng.EvictTarget("no-such-target")
})
assert.False(
t, eng.ExportHasArchiveWriter("no-such-webhook"),
t, env.eng.ExportHasArchiveWriter("no-such-target"),
"eviction must not create a writer",
)
}
@@ -289,17 +313,14 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
) {
t.Parallel()
eng, _ := evictTestEngine(t)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
// Prime the registry so the test can hold the very writer the
// eviction is about to detach.
eng.ExportDeliverDatabase(webhookDB, d)
deliverTo(t, env, tgt)
w := eng.ExportArchiveWriterFor(event.WebhookID)
w := env.eng.ExportArchiveWriterFor(tgt.ID)
require.NotNil(t, w)
require.True(t, w.HandleOpen())
@@ -309,7 +330,7 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
// eviction has to contend for the writer's mutex.
race.awaitFirstWrite()
eng.EvictWebhook(event.WebhookID)
env.eng.EvictWebhook(tgt.WebhookID)
sawEvicted, otherErr := race.wait()
@@ -324,41 +345,33 @@ func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
"been evicted",
)
assert.False(
t, eng.ExportHasArchiveWriter(event.WebhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the registry entry must stay gone",
)
}
// TestEvictWebhook_LaterDeliveryRecreatesWriter proves eviction
// does not break archiving for a webhook that is still alive: a
// does not break archiving for a target that is still alive: a
// subsequent delivery gets a brand new writer from the registry.
// It says nothing about the evicted writer itself — that is what
// TestEvictedWriter_WriteDoesNotReopenFile covers.
func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) {
t.Parallel()
eng, _ := evictTestEngine(t)
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
deliverTo(t, env, tgt)
require.True(t, env.eng.ExportHasArchiveWriter(tgt.ID))
eng.ExportDeliverDatabase(webhookDB, d)
require.True(
t, eng.ExportHasArchiveWriter(event.WebhookID),
)
env.eng.EvictWebhook(tgt.WebhookID)
eng.EvictWebhook(event.WebhookID)
// A fresh delivery for the same webhook gets a brand new
// A fresh delivery for the same target gets a brand new
// writer from the registry, so archiving keeps working.
second := seedDatabaseTargetDelivery(
t, webhookDB, event, "",
)
eng.ExportDeliverDatabase(webhookDB, second)
deliverTo(t, env, tgt)
assert.True(
t, eng.ExportHasArchiveWriter(event.WebhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"a later delivery should recreate the writer",
)
}
@@ -370,19 +383,16 @@ func TestEvictWebhook_LaterDeliveryRecreatesWriter(t *testing.T) {
func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) {
t.Parallel()
eng, _ := evictTestEngine(t)
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
deliverTo(t, env, tgt)
eng.ExportDeliverDatabase(webhookDB, d)
w := eng.ExportArchiveWriterFor(event.WebhookID)
w := env.eng.ExportArchiveWriterFor(tgt.ID)
require.NotNil(t, w)
require.True(t, w.HandleOpen())
require.NoError(t, eng.ExportStop(context.Background()))
require.NoError(t, env.eng.ExportStop(context.Background()))
err := w.Write(evictTestRow("ev-after-stop"), 0)
@@ -395,7 +405,7 @@ func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) {
"a refused write must not reopen the archive",
)
assert.False(
t, eng.ExportHasArchiveWriter(event.WebhookID),
t, env.eng.ExportHasArchiveWriter(tgt.ID),
"the stop should empty the registry",
)
+307 -40
View File
@@ -4,13 +4,12 @@ import (
"database/sql"
"fmt"
"log/slog"
"net/http"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/driver/sqlite"
@@ -74,25 +73,18 @@ func removeArchiveFiles(t *testing.T, path string) {
// TestDeliverDatabase_ArchivesEvent verifies that delivering to
// a database target marks the delivery delivered and archives
// the full event into a separate per-webhook archive file.
// the full event into the target's own archive file.
func TestDeliverDatabase_ArchivesEvent(t *testing.T) {
t.Parallel()
dataDir := t.TempDir()
dbMgr := database.NewTestWebhookDBManager(dataDir)
e := delivery.NewTestEngineWithDB(
nil, dbMgr,
archiveTestLogger(),
&http.Client{Timeout: 5 * time.Second},
1,
)
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":true}`)
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
e.ExportDeliverDatabase(webhookDB, d)
env.eng.ExportDeliverDatabase(webhookDB, d)
var updated database.Delivery
@@ -105,8 +97,7 @@ func TestDeliverDatabase_ArchivesEvent(t *testing.T) {
)
archivePath := filepath.Join(
dataDir,
fmt.Sprintf("archive-%s.db", event.WebhookID),
env.dataDir, "archive-sweep-test-archive-"+tgt.ID+".db",
)
assert.FileExists(t, archivePath)
@@ -288,31 +279,31 @@ func TestParseArchiveExpiry(t *testing.T) {
}
}
// seedDatabaseTargetDelivery seeds a pending delivery for a
// database target with the given config JSON and returns the
// in-memory delivery the target handler is invoked with.
// seedDatabaseTargetDelivery seeds a pending delivery of an event
// to a database target and returns the in-memory delivery the
// target handler is invoked with.
func seedDatabaseTargetDelivery(
t *testing.T,
webhookDB *gorm.DB,
event database.Event,
config string,
tgt *database.Target,
) *database.Delivery {
t.Helper()
dlv := seedDelivery(
t, webhookDB, event.ID, uuid.New().String(),
t, webhookDB, event.ID, tgt.ID,
database.DeliveryStatusPending,
)
d := &database.Delivery{
EventID: event.ID,
TargetID: dlv.TargetID,
TargetID: tgt.ID,
Status: database.DeliveryStatusPending,
Event: event,
Target: database.Target{
Name: "test-db",
Name: tgt.Name,
Type: database.TargetTypeDatabase,
Config: config,
Config: tgt.Config,
},
}
d.ID = dlv.ID
@@ -330,22 +321,14 @@ func TestDeliverDatabase_ArchiveFailureFailsDelivery(
) {
t.Parallel()
dataDir := t.TempDir()
e := delivery.NewTestEngineWithDB(
nil, database.NewTestWebhookDBManager(dataDir),
archiveTestLogger(),
&http.Client{Timeout: 5 * time.Second},
1,
)
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, `{"expiry":"nonsense"}`)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"archived":false}`)
d := seedDatabaseTargetDelivery(
t, webhookDB, event, `{"expiry":"nonsense"}`,
)
d := seedDatabaseTargetDelivery(t, webhookDB, event, tgt)
e.ExportDeliverDatabase(webhookDB, d)
env.eng.ExportDeliverDatabase(webhookDB, d)
var updated database.Delivery
@@ -373,10 +356,7 @@ func TestDeliverDatabase_ArchiveFailureFailsDelivery(
)
assert.NoFileExists(t,
filepath.Join(
dataDir,
fmt.Sprintf("archive-%s.db", event.WebhookID),
),
env.archivePath(tgt),
"no archive file should exist for a failed config",
)
}
@@ -400,3 +380,290 @@ func TestValidateArchiveExpiry(t *testing.T) {
)
}
}
// TestArchiveFileName pins the archive file name and the rules
// that make a webhook or target name safe to put in it.
func TestArchiveFileName(t *testing.T) {
t.Parallel()
const id = "3f2a1c9e-8d4b-4c1a-9e2f-0a1b2c3d4e5f"
cases := []struct {
name string
webhook string
target string
want string
}{
{
"plain names", "orders", "archive",
"archive-orders-archive-" + id + ".db",
},
{
"lowercased", "Orders", "Main Archive",
"archive-orders-main-archive-" + id + ".db",
},
{
"a run of other characters is one dash",
`a /\..b`, "c__--d",
"archive-a-b-c-d-" + id + ".db",
},
{
"no dash at either end", " --orders!! ", "(archive)",
"archive-orders-archive-" + id + ".db",
},
{
"path separators", "../../etc/passwd", "a/b",
"archive-etc-passwd-a-b-" + id + ".db",
},
{
"letters outside ASCII are dropped",
"Bestellungen Größe", "café",
"archive-bestellungen-gr-e-caf-" + id + ".db",
},
{
"nothing left is unnamed", "", "!!!",
"archive-unnamed-unnamed-" + id + ".db",
},
{
"cut to 40 characters", strings.Repeat("a", 50), "x",
"archive-" + strings.Repeat("a", 40) + "-x-" + id + ".db",
},
{
"no dash left by the cut",
strings.Repeat("a", 39) + " b", "x",
"archive-" + strings.Repeat("a", 39) + "-x-" + id + ".db",
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()
assert.Equal(
t, tc.want,
delivery.ArchiveFileName(tc.webhook, tc.target, id),
)
})
}
}
// TestDeliverDatabase_EachTargetHasItsOwnArchive proves two
// database targets of one webhook archive into separate files.
func TestDeliverDatabase_EachTargetHasItsOwnArchive(t *testing.T) {
t.Parallel()
env := setupArchiveTest(t)
first := env.seedDatabaseTarget(t, "")
second := env.addDatabaseTarget(t, first.WebhookID, "")
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`)
for _, tgt := range []*database.Target{first, second} {
env.eng.ExportDeliverDatabase(
webhookDB,
seedDatabaseTargetDelivery(t, webhookDB, event, tgt),
)
}
require.NotEqual(
t, env.archivePath(first), env.archivePath(second),
)
assert.Equal(
t, []string{event.ID},
archivedEventIDs(t, env.archivePath(first)),
)
assert.Equal(
t, []string{event.ID},
archivedEventIDs(t, env.archivePath(second)),
)
}
// TestRename_MovesTheFile proves a rename moves the archive, rows
// and all, and that later writes go to the new name.
func TestRename_MovesTheFile(t *testing.T) {
t.Parallel()
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
oldPath := env.archivePath(tgt)
webhookDB := testWebhookDB(t)
first := seedEvent(t, webhookDB, `{"n":1}`)
env.eng.ExportDeliverDatabase(
webhookDB, seedDatabaseTargetDelivery(t, webhookDB, first, tgt),
)
require.FileExists(t, oldPath)
require.NoError(
t, env.eng.Rename(tgt.ID, "Orders", "Long Term"),
)
newPath := filepath.Join(
env.dataDir, "archive-orders-long-term-"+tgt.ID+".db",
)
assert.NoFileExists(t, oldPath)
assert.Equal(t, []string{first.ID}, archivedEventIDs(t, newPath))
second := seedEvent(t, webhookDB, `{"n":2}`)
env.eng.ExportDeliverDatabase(
webhookDB,
seedDatabaseTargetDelivery(t, webhookDB, second, tgt),
)
assert.ElementsMatch(
t, []string{first.ID, second.ID},
archivedEventIDs(t, newPath),
)
assert.NoFileExists(
t, oldPath, "a write after the rename must use the new name",
)
}
// TestRename_NeverReplacesAFile plants a file at the new name, once
// the .db alone, once a lone -wal and once a lone -shm, and proves
// each time that the rename is refused, the planted file survives,
// and the archive keeps its name and its rows.
func TestRename_NeverReplacesAFile(t *testing.T) {
t.Parallel()
for _, suffix := range archiveFileSuffixes() {
t.Run("planted .db"+suffix, func(t *testing.T) {
t.Parallel()
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
oldPath := env.archivePath(tgt)
webhookDB := testWebhookDB(t)
first := seedEvent(t, webhookDB, `{"n":1}`)
env.eng.ExportDeliverDatabase(
webhookDB,
seedDatabaseTargetDelivery(t, webhookDB, first, tgt),
)
newPath := filepath.Join(
env.dataDir, "archive-orders-long-term-"+tgt.ID+".db",
)
plantedPath := newPath + suffix
require.NoError(
t, os.WriteFile(plantedPath, []byte("planted"), 0o600),
)
require.ErrorIs(
t, env.eng.Rename(tgt.ID, "Orders", "Long Term"),
delivery.ErrArchiveNameTaken,
)
//nolint:gosec // reads the file the test planted under t.TempDir()
planted, err := os.ReadFile(plantedPath)
require.NoError(t, err)
assert.Equal(t, "planted", string(planted))
second := seedEvent(t, webhookDB, `{"n":2}`)
env.eng.ExportDeliverDatabase(
webhookDB,
seedDatabaseTargetDelivery(t, webhookDB, second, tgt),
)
assert.ElementsMatch(
t, []string{first.ID, second.ID},
archivedEventIDs(t, oldPath),
)
})
}
}
// TestRename_BeforeTheNameIsSaved covers the order the handlers
// use: they rename before they save the new name, so a delivery in
// between must write under the new name although the main database
// still has the old one. It also shows that renaming an archive that
// does not exist yet is not an error.
func TestRename_BeforeTheNameIsSaved(t *testing.T) {
t.Parallel()
env := setupArchiveTest(t)
tgt := env.seedDatabaseTarget(t, "")
require.NoError(
t, env.eng.Rename(tgt.ID, "Orders", "Archive"),
)
webhookDB := testWebhookDB(t)
event := seedEvent(t, webhookDB, `{"n":1}`)
env.eng.ExportDeliverDatabase(
webhookDB, seedDatabaseTargetDelivery(t, webhookDB, event, tgt),
)
assert.FileExists(
t,
filepath.Join(
env.dataDir, "archive-orders-archive-"+tgt.ID+".db",
),
)
assert.NoFileExists(t, env.archivePath(tgt))
}
// TestArchiveWriter_RenameMovesSidecars proves a rename carries
// the -wal and -shm a crash can leave beside an archive no handle
// has opened since. SQLite finds them by name, so a -wal left
// behind would lose the transactions it holds.
func TestArchiveWriter_RenameMovesSidecars(t *testing.T) {
t.Parallel()
dir := t.TempDir()
oldPath := filepath.Join(dir, "archive-old.db")
newPath := filepath.Join(dir, "archive-new.db")
for _, suffix := range archiveFileSuffixes() {
require.NoError(
t, os.WriteFile(oldPath+suffix, []byte(suffix), 0o600),
)
}
w := delivery.NewExportArchiveWriter(
oldPath, archiveTestLogger(), 0,
)
require.NoError(t, w.Rename("archive-new.db"))
for _, suffix := range archiveFileSuffixes() {
assert.NoFileExists(t, oldPath+suffix)
assert.FileExists(t, newPath+suffix)
}
assert.Equal(t, newPath, w.Path())
}
// TestArchiveWriter_RenameMovesBackOnFailure makes the -wal fail to
// move after the .db has moved, and proves the .db is moved back, so
// the archive is never split across two names. The new name is 255
// bytes, the longest a file name may be, so the .db can take it but
// the -wal, four bytes longer, cannot.
func TestArchiveWriter_RenameMovesBackOnFailure(t *testing.T) {
t.Parallel()
dir := t.TempDir()
oldPath := filepath.Join(dir, "archive-old.db")
newName := strings.Repeat("a", 252) + ".db"
for _, suffix := range archiveFileSuffixes() {
require.NoError(
t, os.WriteFile(oldPath+suffix, []byte(suffix), 0o600),
)
}
w := delivery.NewExportArchiveWriter(
oldPath, archiveTestLogger(), 0,
)
require.Error(t, w.Rename(newName))
for _, suffix := range archiveFileSuffixes() {
assert.FileExists(t, oldPath+suffix)
}
assert.NoFileExists(t, filepath.Join(dir, newName))
assert.Equal(t, oldPath, w.Path())
}