Name each database target's archive for its webhook and target (closes #376)
check / check (push) Successful in 3m16s
check / check (push) Successful in 3m16s
Each database target now has its own archive file, archive-WEBHOOKNAME-TARGETNAME-TARGETID.db, instead of one archive-WEBHOOKID.db per webhook. delivery.ArchiveFileName builds the name. A change of webhook or target name renames its archive files under the archive writer's lock, before the new name is saved, and back again if the save fails. Webhook edits, target edits and target creation run one at a time, so no two of them interleave. A rename never replaces a file, and one that fails part way moves back what it moved. Deleting a target evicts only that target's writer. Archive files are never deleted, and nothing looks for files under the old name. Model: opus-5-5
This commit is contained in:
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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",
|
||||
)
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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",
|
||||
)
|
||||
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user