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

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: each name is lowercased, keeps ASCII letters and digits, turns
every other run of characters into one dash, and is cut to 40
characters.

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. A rename never replaces a file: if one already has the
new name, the edit is refused. 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:
2026-10-02 07:54:37 +00:00
committed by clawbot
parent c23ffbac65
commit 6257c6ec23
21 changed files with 1545 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)
+81 -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,62 @@ 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.
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 _, suffix := range suffixes {
err := os.Rename(w.path+suffix, path+suffix)
if err != nil && !errors.Is(err, fs.ErrNotExist) {
return fmt.Errorf(
"renaming archive %s to %s: %w", w.path, path, 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",
)
+275 -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,258 @@ 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())
}
+3 -3
View File
@@ -63,7 +63,7 @@ type HandlersParams struct {
Session *session.Session
Middleware *middleware.Middleware
Notifier delivery.Notifier
Evictor delivery.WebhookEvictor
Archives delivery.Archives
SSRFGuard *delivery.Guard
Metrics *metrics.Set
Registry *prometheus.Registry
@@ -80,7 +80,7 @@ type Handlers struct {
session *session.Session
mw *middleware.Middleware
notifier delivery.Notifier
evictor delivery.WebhookEvictor
archives delivery.Archives
mtr *metrics.Set
templates map[string]*template.Template
@@ -131,7 +131,7 @@ func New(
s.session = params.Session
s.mw = params.Middleware
s.notifier = params.Notifier
s.evictor = params.Evictor
s.archives = params.Archives
s.mtr = params.Metrics
s.ssrf = params.SSRFGuard
+84 -11
View File
@@ -3,6 +3,7 @@ package handlers_test
import (
"context"
"errors"
"fmt"
"html/template"
"net/http"
"net/http/httptest"
@@ -52,23 +53,73 @@ func (n *recordingNotifier) Tasks() []delivery.Task {
return out
}
// recordingEvictor is a delivery.WebhookEvictor that records
// the webhook ids it was asked to evict, so a test can prove
// that a deletion path reached the delivery engine.
type recordingEvictor struct {
mu sync.Mutex
evicted []string
// recordingArchives is a delivery.Archives that records what it
// was asked to do, so a test can prove that a deletion or rename
// path reached the delivery engine. After FailRenames, every
// rename fails with the given error.
type recordingArchives struct {
mu sync.Mutex
evicted []string
evictedTargets []string
renames []archiveRename
renameErr error
}
func (r *recordingEvictor) EvictWebhook(webhookID string) {
// errInjectedRename is the failure a test hands FailRenames.
var errInjectedRename = errors.New("injected rename failure")
// errNameTaken is what the delivery engine returns when a file
// already has an archive's new name, here archive-taken.db.
var errNameTaken = fmt.Errorf(
"%w: archive-taken.db", delivery.ErrArchiveNameTaken,
)
// archiveRename is one recorded Rename call.
type archiveRename struct {
TargetID string
WebhookName string
TargetName string
}
func (r *recordingArchives) EvictWebhook(webhookID string) {
r.mu.Lock()
defer r.mu.Unlock()
r.evicted = append(r.evicted, webhookID)
}
func (r *recordingArchives) EvictTarget(targetID string) {
r.mu.Lock()
defer r.mu.Unlock()
r.evictedTargets = append(r.evictedTargets, targetID)
}
func (r *recordingArchives) Rename(
targetID, webhookName, targetName string,
) error {
r.mu.Lock()
defer r.mu.Unlock()
r.renames = append(r.renames, archiveRename{
TargetID: targetID,
WebhookName: webhookName,
TargetName: targetName,
})
return r.renameErr
}
// FailRenames makes every later rename fail with err.
func (r *recordingArchives) FailRenames(err error) {
r.mu.Lock()
defer r.mu.Unlock()
r.renameErr = err
}
// Evicted returns a copy of the recorded webhook ids.
func (r *recordingEvictor) Evicted() []string {
func (r *recordingArchives) Evicted() []string {
r.mu.Lock()
defer r.mu.Unlock()
@@ -78,6 +129,28 @@ func (r *recordingEvictor) Evicted() []string {
return out
}
// EvictedTargets returns a copy of the recorded target ids.
func (r *recordingArchives) EvictedTargets() []string {
r.mu.Lock()
defer r.mu.Unlock()
out := make([]string, len(r.evictedTargets))
copy(out, r.evictedTargets)
return out
}
// Renames returns a copy of the recorded renames.
func (r *recordingArchives) Renames() []archiveRename {
r.mu.Lock()
defer r.mu.Unlock()
out := make([]archiveRename, len(r.renames))
copy(out, r.renames)
return out
}
func newTestApp(
t *testing.T,
targets ...any,
@@ -104,10 +177,10 @@ func newTestApp(
func(n *recordingNotifier) delivery.Notifier {
return n
},
func() *recordingEvictor {
return &recordingEvictor{}
func() *recordingArchives {
return &recordingArchives{}
},
func(r *recordingEvictor) delivery.WebhookEvictor {
func(r *recordingArchives) delivery.Archives {
return r
},
metrics.NewRegistry,
+59 -80
View File
@@ -15,6 +15,7 @@ import (
"gorm.io/gorm"
"gorm.io/gorm/clause"
"sneak.berlin/go/webhooker/internal/database"
"sneak.berlin/go/webhooker/internal/delivery"
"sneak.berlin/go/webhooker/internal/handlers"
"sneak.berlin/go/webhooker/internal/session"
)
@@ -79,6 +80,10 @@ func seedTarget(
// from a delete statement.
var errInjectedDelete = errors.New("injected delete failure")
// errInjectedSave is the failure failSaveOnTable reports from a
// save of an existing row.
var errInjectedSave = errors.New("injected save failure")
// seedEntrypoint inserts an entrypoint for a webhook.
func seedEntrypoint(
t *testing.T,
@@ -146,19 +151,42 @@ func failDeleteOnTable(
)
}
// failSaveOnTable is failDeleteOnTable for saves: every update of
// an existing row in the named table fails.
func failSaveOnTable(
t *testing.T,
db *database.Database,
table string,
) {
t.Helper()
require.NoError(t, db.DB().Callback().Update().
Before("gorm:update").
Register(
"test:fail_save_"+table,
func(tx *gorm.DB) {
if tx.Statement.Table == table {
_ = tx.AddError(errInjectedSave)
}
},
),
)
}
// archivePathFor returns the archive database path the
// delivery engine would use for a webhook: beside the webhook's
// event database in the data directory.
// delivery engine would use for a database target: beside the
// webhook's event database in the data directory.
func archivePathFor(
t *testing.T,
mgr *database.WebhookDBManager,
webhookID string,
wh *database.Webhook,
tgt *database.Target,
) string {
t.Helper()
return filepath.Join(
filepath.Dir(mgr.DBPath(webhookID)),
"archive-"+webhookID+".db",
filepath.Dir(mgr.DBPath(wh.ID)),
delivery.ArchiveFileName(wh.Name, tgt.Name, tgt.ID),
)
}
@@ -195,8 +223,8 @@ func postRequest(
// TestHandleSourceDelete_EvictsArchiveWriter proves that
// deleting a webhook reaches the delivery engine and releases
// the webhook's archive writer, exercised through the real
// deletion handler rather than by calling the evictor directly.
// the webhook's archive writers, exercised through the real
// deletion handler rather than by calling the engine directly.
func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) {
t.Parallel()
@@ -204,7 +232,7 @@ func TestHandleSourceDelete_EvictsArchiveWriter(t *testing.T) {
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
ev *recordingArchives
)
app := newTestApp(t, &h, &sess, &db, &ev)
@@ -254,9 +282,10 @@ func TestHandleSourceDelete_KeepsArchiveFile(t *testing.T) {
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
tgt := seedTarget(t, db, wh.ID, database.TargetTypeDatabase)
// Place an archive file where the delivery engine would.
archivePath := archivePathFor(t, mgr, wh.ID)
archivePath := archivePathFor(t, mgr, wh, tgt)
require.NoError(
t,
writeArchivePlaceholder(archivePath),
@@ -435,68 +464,17 @@ func TestHandleSourceDelete_RemovesConfigAndEventDatabase(
)
}
// TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone
// proves that removing the last database target releases the
// archive writer.
func TestHandleTargetDelete_EvictsWhenLastDatabaseTargetGone(
t *testing.T,
) {
// TestHandleTargetDelete_EvictsThatTarget proves that deleting a
// database target releases that target's archive writer and no
// other: the webhook's other database target keeps its own.
func TestHandleTargetDelete_EvictsThatTarget(t *testing.T) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
)
app := newTestApp(t, &h, &sess, &db, &ev)
app.RequireStart()
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
tgt := seedTarget(
t, db, wh.ID, database.TargetTypeDatabase,
)
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
req := postRequest(
"/hook/"+wh.ID+"/targets/"+tgt.ID+"/delete",
cookies,
map[string]string{
paramSourceID: wh.ID,
paramTargetID: tgt.ID,
},
)
w := httptest.NewRecorder()
h.HandleTargetDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Equal(
t, []string{wh.ID}, ev.Evicted(),
"removing the last database target should evict",
)
}
// TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains
// proves that deleting one of several database targets leaves
// the still-needed archive writer alone: the surviving target
// keeps archiving to the same file, so the writer must stay.
func TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains(
t *testing.T,
) {
t.Parallel()
var (
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
ev *recordingArchives
)
app := newTestApp(t, &h, &sess, &db, &ev)
@@ -527,17 +505,17 @@ func TestHandleTargetDelete_KeepsWriterWhenDatabaseTargetRemains(
h.HandleTargetDelete().ServeHTTP(w, req)
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Empty(
t, ev.Evicted(),
"a second database target still needs the writer",
assert.Equal(
t, []string{doomed.ID}, ev.EvictedTargets(),
"deleting a database target should evict its writer",
)
assert.Empty(t, ev.Evicted(), "the webhook is not deleted")
}
// TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted proves
// that deleting a target of an unrelated type leaves a
// still-needed archive writer alone: the webhook's database
// target is untouched, so its writer must stay.
func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
// TestHandleTargetDelete_IgnoresAnotherWebhooksTarget proves that
// a target id from the URL that is not a target of the webhook
// deletes nothing and so evicts nothing.
func TestHandleTargetDelete_IgnoresAnotherWebhooksTarget(
t *testing.T,
) {
t.Parallel()
@@ -546,7 +524,7 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
h *handlers.Handlers
sess *session.Session
db *database.Database
ev *recordingEvictor
ev *recordingArchives
)
app := newTestApp(t, &h, &sess, &db, &ev)
@@ -555,19 +533,20 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
t.Cleanup(app.RequireStop)
wh := seedWebhook(t, db)
seedTarget(t, db, wh.ID, database.TargetTypeDatabase)
other := seedTarget(t, db, wh.ID, database.TargetTypeLog)
elsewhere := seedTarget(
t, db, seedWebhook(t, db).ID, database.TargetTypeDatabase,
)
cookies := authenticatedCookies(
t, sess, deleteTestUserID, deleteTestUsername,
)
req := postRequest(
"/hook/"+wh.ID+"/targets/"+other.ID+"/delete",
"/hook/"+wh.ID+"/targets/"+elsewhere.ID+"/delete",
cookies,
map[string]string{
paramSourceID: wh.ID,
paramTargetID: other.ID,
paramTargetID: elsewhere.ID,
},
)
w := httptest.NewRecorder()
@@ -576,7 +555,7 @@ func TestHandleTargetDelete_KeepsWriterWhenOtherTypeDeleted(
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Empty(
t, ev.Evicted(),
"a surviving database target must keep its writer",
t, ev.EvictedTargets(),
"another webhook's target must not be evicted",
)
}
+86 -44
View File
@@ -549,6 +549,7 @@ func (h *Handlers) applyWebhookEdit(
return
}
oldName := webhook.Name
webhook.Name = name
webhook.Description = r.PostFormValue("description")
@@ -571,8 +572,40 @@ func (h *Handlers) applyWebhookEdit(
webhook.RetentionDays = retentionDays
err := h.db.DB().Save(webhook).Error
// A new name renames the archive files before it is saved (see
// delivery.Engine.Rename). If either step fails, they go back to
// the name that is still stored.
err := h.renameWebhookArchives(webhook.ID, oldName, webhook.Name)
if err == nil {
err = h.db.DB().Save(webhook).Error
}
if err != nil {
restoreErr := h.renameWebhookArchives(
webhook.ID, webhook.Name, oldName,
)
if restoreErr != nil {
h.log.Error(
"failed to rename archives back",
"webhook_id", webhook.ID,
"error", restoreErr,
)
}
if errors.Is(err, delivery.ErrArchiveNameTaken) {
data := map[string]any{
tmplKeyWebhook: webhook,
tmplKeyError: "Not saved: " + err.Error() +
". Move that file out of the data directory, " +
"then save again.",
}
w.WriteHeader(http.StatusConflict)
h.renderTemplate(w, r, "source_edit.html", data)
return
}
h.serverError(w, r, "failed to update webhook", err)
return
@@ -707,11 +740,11 @@ func (h *Handlers) commitWebhookDeletion(
return tx.Commit().Error
}
// evictArchiveWriter asks the delivery engine to drop its
// cached archive writer for a webhook, closing the archive file
// handle.
// evictArchiveWriter asks the delivery engine to drop the cached
// archive writers of a webhook's database targets, closing their
// archive file handles.
//
// The archive database file is NOT deleted. Unlike the event
// The archive database files are NOT deleted. Unlike the event
// database — which is per-webhook working storage and is
// hard-deleted with the webhook — an archive is explicitly
// long-term storage that an operator may want to keep or move
@@ -719,50 +752,58 @@ func (h *Handlers) commitWebhookDeletion(
// deleting a webhook would be a surprising and unrecoverable
// data loss, so the file is left for the operator to handle.
func (h *Handlers) evictArchiveWriter(webhookID string) {
if h.evictor == nil {
if h.archives == nil {
return
}
h.evictor.EvictWebhook(webhookID)
h.archives.EvictWebhook(webhookID)
}
// evictArchiveWriterIfUnused releases a webhook's archive
// writer once the webhook has no database target left to feed
// it.
//
// It is called after any child resource of a webhook is
// deleted, and is correct without knowing which kind was: it
// evicts only when no database target remains, so deleting one
// of several database targets — or deleting an unrelated
// target type — leaves a still-needed writer alone. When no
// database target ever existed there is no writer and eviction
// is a no-op. Soft-deleted targets are excluded by GORM's
// default scope, so the row just deleted is not counted.
func (h *Handlers) evictArchiveWriterIfUnused(webhookID string) {
var remaining int64
// evictTargetArchiveWriter is evictArchiveWriter for one deleted
// target, and leaves its archive file on disk for the same reason.
// A target that is not a database target has no writer, and
// evicting it does nothing.
func (h *Handlers) evictTargetArchiveWriter(targetID string) {
if h.archives == nil {
return
}
h.archives.EvictTarget(targetID)
}
// renameWebhookArchives renames the archive file of every database
// target of a webhook from the webhook name oldName to newName,
// keeping each target's own name. It does nothing when the name is
// unchanged, and stops at the first failure.
func (h *Handlers) renameWebhookArchives(
webhookID, oldName, newName string,
) error {
if h.archives == nil || oldName == newName {
return nil
}
var targets []database.Target
err := h.db.DB().
Model(&database.Target{}).
Where(
"webhook_id = ? AND type = ?",
webhookID, database.TargetTypeDatabase,
).
Count(&remaining).Error
Find(&targets).Error
if err != nil {
h.log.Error(
"failed to count remaining database targets",
"webhook_id", webhookID,
"error", err,
return err
}
for i := range targets {
err = h.archives.Rename(
targets[i].ID, newName, targets[i].Name,
)
return
if err != nil {
return err
}
}
if remaining > 0 {
return
}
h.evictArchiveWriter(webhookID)
return nil
}
// ownedWebhook resolves the request's sourceID parameter to a
@@ -1646,27 +1687,26 @@ func (h *Handlers) HandleEntrypointDelete() http.HandlerFunc {
)
}
// HandleTargetDelete handles deleting a target. Deleting the
// last database target of a webhook leaves its archive writer
// with nothing to write, so the writer is evicted and its
// handle closed; the archive file is left on disk.
// HandleTargetDelete handles deleting a target. A deleted
// database target's archive writer is evicted and its handle
// closed; the archive file is left on disk.
func (h *Handlers) HandleTargetDelete() http.HandlerFunc {
return h.deleteChildResource(
"targetID", &database.Target{},
"failed to delete target",
h.evictArchiveWriterIfUnused,
h.evictTargetArchiveWriter,
)
}
// deleteChildResource returns a handler that deletes a child
// resource (entrypoint or target) belonging to a webhook. The
// optional afterDelete hook runs with the webhook's id once the
// delete has succeeded, before the redirect.
// optional afterDelete hook runs with the child's id once the
// delete has removed it, before the redirect.
func (h *Handlers) deleteChildResource(
idParam string,
model any,
errMsg string,
afterDelete func(webhookID string),
afterDelete func(childID string),
) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
userID, ok := h.getUserID(r)
@@ -1702,8 +1742,10 @@ func (h *Handlers) deleteChildResource(
return
}
if afterDelete != nil {
afterDelete(webhook.ID)
// Only for a row this webhook really had: the id came from
// the URL and may name another webhook's child.
if afterDelete != nil && result.RowsAffected > 0 {
afterDelete(childID)
}
http.Redirect(
+141 -1
View File
@@ -187,6 +187,7 @@ func storedRetentionDays(
type sourceTestEnv struct {
handlers *handlers.Handlers
db *database.Database
archives *recordingArchives
cookies []*http.Cookie
}
@@ -199,7 +200,9 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
var db *database.Database
app := newTestApp(t, &h, &sess, &db)
var archives *recordingArchives
app := newTestApp(t, &h, &sess, &db, &archives)
app.RequireStart()
t.Cleanup(app.RequireStop)
@@ -207,6 +210,7 @@ func setupSourceTest(t *testing.T) *sourceTestEnv {
return &sourceTestEnv{
handlers: h,
db: db,
archives: archives,
cookies: authenticatedCookies(
t, sess, sourceTestUserID, "sourceuser",
),
@@ -498,6 +502,142 @@ func TestHandleSourceEditSubmit_EmptyRetentionLeavesValueUnchanged(
assert.Equal(t, 7, storedRetentionDays(t, env.db, wh.ID))
}
// renamedWebhookName is the name the rename tests give a webhook.
const renamedWebhookName = "Renamed"
// TestHandleSourceEditSubmit_RenamesArchives proves that a save
// that keeps the webhook's name renames nothing, and that renaming a
// webhook renames the archive of each of its database targets and
// asks nothing of its other targets.
func TestHandleSourceEditSubmit_RenamesArchives(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
first := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
second := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
seedTarget(t, env.db, wh.ID, database.TargetTypeLog)
w := submitEdit(t, env, wh, "")
require.Equal(t, http.StatusSeeOther, w.Code)
assert.Empty(t, env.archives.Renames())
wh.Name = renamedWebhookName
w = submitEdit(t, env, wh, "")
require.Equal(t, http.StatusSeeOther, w.Code)
assert.ElementsMatch(
t,
[]archiveRename{
{first.ID, renamedWebhookName, first.Name},
{second.ID, renamedWebhookName, second.Name},
},
env.archives.Renames(),
)
}
// TestHandleSourceEditSubmit_FailedRenameKeepsTheName proves that a
// webhook whose archive cannot be renamed keeps its stored name, so
// the name on disk and the name in the UI do not part, and that the
// handler puts back what it may already have moved.
func TestHandleSourceEditSubmit_FailedRenameKeepsTheName(
t *testing.T,
) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
tgt := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
env.archives.FailRenames(errInjectedRename)
oldName := wh.Name
wh.Name = renamedWebhookName
w := submitEdit(t, env, wh, "")
require.Equal(t, http.StatusInternalServerError, w.Code)
var stored database.Webhook
require.NoError(
t, env.db.DB().First(&stored, "id = ?", wh.ID).Error,
)
assert.Equal(t, oldName, stored.Name)
assert.Equal(
t,
[]archiveRename{
{tgt.ID, renamedWebhookName, tgt.Name},
{tgt.ID, oldName, tgt.Name},
},
env.archives.Renames(),
)
}
// TestHandleSourceEditSubmit_FailedSaveRenamesBack proves that when
// the archive is renamed but the new name cannot be saved, the
// archive is renamed back to the stored name and the stored name
// stays.
func TestHandleSourceEditSubmit_FailedSaveRenamesBack(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
tgt := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
failSaveOnTable(t, env.db, "webhooks")
oldName := wh.Name
wh.Name = renamedWebhookName
w := submitEdit(t, env, wh, "")
require.Equal(t, http.StatusInternalServerError, w.Code)
var stored database.Webhook
require.NoError(
t, env.db.DB().First(&stored, "id = ?", wh.ID).Error,
)
assert.Equal(t, oldName, stored.Name)
assert.Equal(
t,
[]archiveRename{
{tgt.ID, renamedWebhookName, tgt.Name},
{tgt.ID, oldName, tgt.Name},
},
env.archives.Renames(),
)
}
// TestHandleSourceEditSubmit_ArchiveNameTaken proves that when a file
// already has an archive's new name, the edit is refused with an
// error naming that file, and the webhook keeps its stored name.
func TestHandleSourceEditSubmit_ArchiveNameTaken(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
env.archives.FailRenames(errNameTaken)
oldName := wh.Name
wh.Name = renamedWebhookName
w := submitEdit(t, env, wh, "")
require.Equal(t, http.StatusConflict, w.Code)
assert.Contains(t, w.Body.String(), "archive-taken.db")
var stored database.Webhook
require.NoError(
t, env.db.DB().First(&stored, "id = ?", wh.ID).Error,
)
assert.Equal(t, oldName, stored.Name)
}
// TestSourceEditForm_ForeverWebhookRoundTrips walks the exact path that
// the removed max="365" cap used to break: render the edit form for a
// retain-forever webhook, confirm the pre-filled sentinel is not capped
+48 -1
View File
@@ -1,6 +1,7 @@
package handlers
import (
"errors"
"net/http"
"github.com/go-chi/chi"
@@ -150,11 +151,42 @@ func (h *Handlers) applyTargetEdit(
target.MaxRetries = retries
}
oldName := target.Name
target.Name = name
target.Config = configJSON
err = h.db.DB().Save(target).Error
// A new name renames the archive file before it is saved (see
// delivery.Engine.Rename). If either step fails, it goes back to
// the name that is still stored.
err = h.renameTargetArchive(target, webhook.Name, oldName, name)
if err == nil {
err = h.db.DB().Save(target).Error
}
if err != nil {
restoreErr := h.renameTargetArchive(
target, webhook.Name, name, oldName,
)
if restoreErr != nil {
h.log.Error(
"failed to rename archive back",
"target_id", target.ID,
"error", restoreErr,
)
}
if errors.Is(err, delivery.ErrArchiveNameTaken) {
http.Error(
w,
"Not saved: "+err.Error()+
". Move that file out of the data directory, "+
"then save again.",
http.StatusConflict,
)
return
}
h.serverError(w, r, "failed to update target", err)
return
@@ -165,6 +197,21 @@ func (h *Handlers) applyTargetEdit(
)
}
// renameTargetArchive renames a database target's archive file from
// the target name oldName to newName. It does nothing when the name
// is unchanged; other target types have no archive.
func (h *Handlers) renameTargetArchive(
target *database.Target,
webhookName, oldName, newName string,
) error {
if h.archives == nil || oldName == newName ||
target.Type != database.TargetTypeDatabase {
return nil
}
return h.archives.Rename(target.ID, webhookName, newName)
}
// renderTargetEdit renders the target edit page with an optional
// error message.
func (h *Handlers) renderTargetEdit(
+97
View File
@@ -635,3 +635,100 @@ func assertWebhookOfAnotherUser404s(
assert.Equal(t, http.StatusNotFound, w.Code)
}
// renamedTargetName is the name the rename tests give a target.
const renamedTargetName = "Long Term"
// TestHandleTargetEditSubmit_RenamesArchive proves that renaming a
// database target renames its archive, that a save that keeps the
// name renames nothing, that a target of another type has no archive
// to rename, and that a target whose archive cannot be renamed keeps
// its stored name. When a file already has the archive's new name,
// the edit is refused with an error naming that file.
func TestHandleTargetEditSubmit_RenamesArchive(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
rename := url.Values{"name": {renamedTargetName}}
w := submitTargetEdit(env, wh.ID, archive.ID, rename)
require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String())
assert.Equal(
t,
[]archiveRename{{archive.ID, wh.Name, renamedTargetName}},
env.archives.Renames(),
)
w = submitTargetEdit(env, wh.ID, archive.ID, rename)
require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String())
assert.Len(
t, env.archives.Renames(), 1,
"a save that keeps the name renames nothing",
)
httpWebhook, httpTarget := seedHTTPTarget(t, env, "", "")
w = submitTargetEdit(
env, httpWebhook.ID, httpTarget.ID,
editForm(editOriginalURL, "", ""),
)
require.Equal(t, http.StatusSeeOther, w.Code, w.Body.String())
assert.Len(
t, env.archives.Renames(), 1,
"an HTTP target has no archive to rename",
)
again := url.Values{"name": {"Again"}}
env.archives.FailRenames(errInjectedRename)
w = submitTargetEdit(env, wh.ID, archive.ID, again)
require.Equal(t, http.StatusInternalServerError, w.Code)
assert.Equal(
t, renamedTargetName, storedTarget(t, env, archive.ID).Name,
"a target whose archive was not renamed keeps its name",
)
env.archives.FailRenames(errNameTaken)
w = submitTargetEdit(env, wh.ID, archive.ID, again)
require.Equal(t, http.StatusConflict, w.Code)
assert.Contains(t, w.Body.String(), "archive-taken.db")
assert.Equal(
t, renamedTargetName, storedTarget(t, env, archive.ID).Name,
)
}
// TestHandleTargetEditSubmit_FailedSaveRenamesBack proves that when a
// database target's archive is renamed but the new name cannot be
// saved, the archive is renamed back to the stored name and the
// stored name stays.
func TestHandleTargetEditSubmit_FailedSaveRenamesBack(t *testing.T) {
t.Parallel()
env := setupSourceTest(t)
wh := seedWebhookWithRetention(t, env.db, 7)
archive := seedTarget(t, env.db, wh.ID, database.TargetTypeDatabase)
failSaveOnTable(t, env.db, "targets")
form := url.Values{}
form.Set("name", renamedTargetName)
w := submitTargetEdit(env, wh.ID, archive.ID, form)
require.Equal(t, http.StatusInternalServerError, w.Code)
assert.Equal(
t, archive.Name, storedTarget(t, env, archive.ID).Name,
)
assert.Equal(
t,
[]archiveRename{
{archive.ID, wh.Name, renamedTargetName},
{archive.ID, wh.Name, archive.Name},
},
env.archives.Renames(),
)
}
+9 -3
View File
@@ -132,9 +132,15 @@ type noopNotifier struct{}
func (n *noopNotifier) Notify([]delivery.Task) {}
type noopEvictor struct{}
type noopArchives struct{}
func (n *noopEvictor) EvictWebhook(string) {}
func (n *noopArchives) EvictWebhook(string) {}
func (n *noopArchives) EvictTarget(string) {}
func (n *noopArchives) Rename(_, _, _ string) error {
return nil
}
// newServerApp starts the real login path against dir: the handlers,
// the middleware that bounds password verification, the session store
@@ -163,7 +169,7 @@ func newServerApp(
healthcheck.New,
session.New,
func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} },
func() delivery.Archives { return &noopArchives{} },
metrics.NewRegistry,
metrics.New,
middleware.New,
+12 -6
View File
@@ -47,12 +47,18 @@ type noopNotifier struct{}
func (n *noopNotifier) Notify([]delivery.Task) {}
// noopEvictor satisfies handlers.New's delivery.WebhookEvictor
// dependency. No test here checks what gets evicted, so it records
// nothing.
type noopEvictor struct{}
// noopArchives satisfies handlers.New's delivery.Archives
// dependency. No test here checks what gets evicted or renamed, so
// it records nothing.
type noopArchives struct{}
func (e *noopEvictor) EvictWebhook(string) {}
func (e *noopArchives) EvictWebhook(string) {}
func (e *noopArchives) EvictTarget(string) {}
func (e *noopArchives) Rename(_, _, _ string) error {
return nil
}
// testEnv is the real router from routes.go plus the collaborators
// tests need to seed users and forge sessions.
@@ -113,7 +119,7 @@ func newTestEnvWithConfig(
healthcheck.New,
session.New,
func() delivery.Notifier { return &noopNotifier{} },
func() delivery.WebhookEvictor { return &noopEvictor{} },
func() delivery.Archives { return &noopArchives{} },
metrics.NewRegistry,
metrics.New,
middleware.New,