check / check (push) Successful in 3m32s
Brings `prod`, which upaas deploys, up to `main` at `9cf9cdd`, the merge of #321. `prod` was cut from `main` at `251cb3d` (1.0.0b1). What it deploys is everything listed in #321. For running it: - With `WEBHOOKER_ENVIRONMENT` unset, the instance runs as `prod` and sends no `Access-Control-Allow-Origin: *`. - Each event database gains its new indexes the first time it is opened after the upgrade. - `webhooker_delivery_retries_total` no longer counts a circuit breaker holding back a delivery that is already `retrying`. Not in this PR yet: #340, in which the container sets its own data directory owner and mode before start. It is in progress on `next`. Once it reaches `main`, this PR carries it, because the PR follows `main`. Model: opus-5-5 Co-authored-by: Jeffrey Paul <1+sneak@noreply.example.org> Reviewed-on: #343
408 lines
10 KiB
Go
408 lines
10 KiB
Go
package delivery_test
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"net/http"
|
|
"os"
|
|
"path/filepath"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
"sneak.berlin/go/webhooker/internal/database"
|
|
"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) {
|
|
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",
|
|
)
|
|
|
|
eng.EvictWebhook(webhookID)
|
|
|
|
assert.False(
|
|
t, eng.ExportHasArchiveWriter(webhookID),
|
|
"eviction should remove the registry entry",
|
|
)
|
|
assert.False(
|
|
t, eng.ExportArchiveHandleOpen(webhookID),
|
|
"eviction should close the archive handle",
|
|
)
|
|
|
|
archivePath := filepath.Join(
|
|
dataDir, fmt.Sprintf("archive-%s.db", webhookID),
|
|
)
|
|
assert.FileExists(
|
|
t, archivePath,
|
|
"eviction must not delete the archive file",
|
|
)
|
|
}
|
|
|
|
// 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.
|
|
func TestEvictWebhook_UnknownWebhookIsNoOp(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
eng, _ := evictTestEngine(t)
|
|
|
|
assert.NotPanics(t, func() {
|
|
eng.EvictWebhook("no-such-webhook")
|
|
eng.EvictWebhook("no-such-webhook")
|
|
})
|
|
|
|
assert.False(
|
|
t, eng.ExportHasArchiveWriter("no-such-webhook"),
|
|
"eviction must not create a writer",
|
|
)
|
|
}
|
|
|
|
// evictTestRow builds an archive row for the eviction tests.
|
|
func evictTestRow(eventID string) delivery.ExportArchivedEvent {
|
|
return delivery.ExportArchivedEvent{
|
|
EventID: eventID,
|
|
WebhookID: "wh-evict",
|
|
Method: http.MethodPost,
|
|
Body: `{"seeded":true}`,
|
|
}
|
|
}
|
|
|
|
// TestEvictedWriter_WriteDoesNotReopenFile is the direct test of
|
|
// the evicted guard on the write path. A writer that has left
|
|
// the registry is held by nobody, so a handle it opened could
|
|
// never be closed again: it must refuse the write outright
|
|
// rather than recreate the archive behind the registry's back.
|
|
//
|
|
// The archive file is removed before the eviction, so an
|
|
// unguarded write is unmistakable — it recreates the file.
|
|
func TestEvictedWriter_WriteDoesNotReopenFile(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
path := filepath.Join(t.TempDir(), "archive-evicted.db")
|
|
|
|
w := delivery.NewExportArchiveWriter(
|
|
path, archiveTestLogger(), 0,
|
|
)
|
|
|
|
require.NoError(t, w.Write(evictTestRow("ev-1"), 0))
|
|
require.FileExists(t, path)
|
|
|
|
// The operator moves the archive away for offline retention,
|
|
// which the write path would ordinarily undo on the next
|
|
// write by recreating the file.
|
|
require.NoError(t, os.Remove(path))
|
|
|
|
w.Evict()
|
|
|
|
err := w.Write(evictTestRow("ev-2"), 0)
|
|
|
|
require.ErrorIs(
|
|
t, err, delivery.ErrExportArchiveWriterEvicted,
|
|
"an evicted writer must refuse writes",
|
|
)
|
|
assert.NoFileExists(
|
|
t, path,
|
|
"an evicted writer must not reopen (or recreate) the "+
|
|
"archive file",
|
|
)
|
|
assert.False(
|
|
t, w.HandleOpen(),
|
|
"an evicted writer must hold no handle",
|
|
)
|
|
}
|
|
|
|
// TestEvictedWriter_SweepDoesNotReopenFile is the same test for
|
|
// the sweep path: an idle sweep that reaches a writer already
|
|
// evicted underneath it must return the sentinel rather than
|
|
// reopen a file nothing owns.
|
|
func TestEvictedWriter_SweepDoesNotReopenFile(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
path := filepath.Join(t.TempDir(), "archive-evicted.db")
|
|
|
|
w := delivery.NewExportArchiveWriter(
|
|
path, archiveTestLogger(), 0,
|
|
)
|
|
|
|
require.NoError(t, w.Write(evictTestRow("ev-1"), 0))
|
|
require.FileExists(t, path)
|
|
|
|
w.Evict()
|
|
|
|
err := w.SweepExpired(time.Hour)
|
|
|
|
require.ErrorIs(
|
|
t, err, delivery.ErrExportArchiveWriterEvicted,
|
|
"an evicted writer must refuse an idle sweep",
|
|
)
|
|
assert.False(
|
|
t, w.HandleOpen(),
|
|
"a refused sweep must not leave a handle open",
|
|
)
|
|
}
|
|
|
|
// racingWrites drives a pack of goroutines writing to one
|
|
// archive writer until each is refused, so an eviction on the
|
|
// test goroutine has to take the writer's mutex away from writes
|
|
// that are already contending for it.
|
|
type racingWrites struct {
|
|
wg sync.WaitGroup
|
|
mu sync.Mutex
|
|
sawEvicted bool
|
|
otherErr error
|
|
started chan struct{}
|
|
}
|
|
|
|
// racingWriteGoroutines is how many goroutines contend for the
|
|
// writer's mutex while the eviction lands.
|
|
const racingWriteGoroutines = 4
|
|
|
|
// startRacingWrites launches the writing goroutines. Each writes
|
|
// in a loop and stops at its first error, recording whether that
|
|
// error was the eviction sentinel. The deadline is a backstop
|
|
// against a hang, not a timing assumption: the first write after
|
|
// the eviction is refused.
|
|
func startRacingWrites(
|
|
w *delivery.ExportArchiveWriter,
|
|
) *racingWrites {
|
|
r := &racingWrites{
|
|
started: make(chan struct{}, racingWriteGoroutines),
|
|
}
|
|
|
|
deadline := time.Now().Add(10 * time.Second)
|
|
|
|
r.wg.Add(racingWriteGoroutines)
|
|
|
|
for i := range racingWriteGoroutines {
|
|
go func() {
|
|
defer r.wg.Done()
|
|
|
|
first := true
|
|
|
|
for time.Now().Before(deadline) {
|
|
err := w.Write(
|
|
evictTestRow(fmt.Sprintf("ev-%d", i)), 0,
|
|
)
|
|
|
|
if first {
|
|
r.started <- struct{}{}
|
|
|
|
first = false
|
|
}
|
|
|
|
if err == nil {
|
|
continue
|
|
}
|
|
|
|
r.record(err)
|
|
|
|
return
|
|
}
|
|
}()
|
|
}
|
|
|
|
return r
|
|
}
|
|
|
|
// record classifies the error that stopped one goroutine.
|
|
func (r *racingWrites) record(err error) {
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
|
|
if errors.Is(err, delivery.ErrExportArchiveWriterEvicted) {
|
|
r.sawEvicted = true
|
|
|
|
return
|
|
}
|
|
|
|
r.otherErr = err
|
|
}
|
|
|
|
// awaitFirstWrite blocks until at least one write has run, so
|
|
// the eviction that follows is a genuine race.
|
|
func (r *racingWrites) awaitFirstWrite() {
|
|
<-r.started
|
|
}
|
|
|
|
// wait joins the goroutines and reports whether any write was
|
|
// refused with the eviction sentinel, plus any unexpected error.
|
|
func (r *racingWrites) wait() (bool, error) {
|
|
r.wg.Wait()
|
|
|
|
r.mu.Lock()
|
|
defer r.mu.Unlock()
|
|
|
|
return r.sawEvicted, r.otherErr
|
|
}
|
|
|
|
// TestEvictWebhook_RacingWriteDoesNotReopenHandle exercises the
|
|
// interleaving the evicted flag exists for: writes already
|
|
// contending for the writer's mutex when the eviction takes it.
|
|
// The write that wins the mutex after the eviction must abandon
|
|
// its work rather than reopen the archive, leaving the writer
|
|
// permanently handle-free. Run under -race.
|
|
func TestEvictWebhook_RacingWriteDoesNotReopenHandle(
|
|
t *testing.T,
|
|
) {
|
|
t.Parallel()
|
|
|
|
eng, _ := evictTestEngine(t)
|
|
|
|
webhookDB := testWebhookDB(t)
|
|
event := seedEvent(t, webhookDB, `{"archived":true}`)
|
|
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
|
|
|
|
// Prime the registry so the test can hold the very writer the
|
|
// eviction is about to detach.
|
|
eng.ExportDeliverDatabase(webhookDB, d)
|
|
|
|
w := eng.ExportArchiveWriterFor(event.WebhookID)
|
|
require.NotNil(t, w)
|
|
require.True(t, w.HandleOpen())
|
|
|
|
race := startRacingWrites(w)
|
|
|
|
// Evict only once writes are genuinely in flight, so the
|
|
// eviction has to contend for the writer's mutex.
|
|
race.awaitFirstWrite()
|
|
|
|
eng.EvictWebhook(event.WebhookID)
|
|
|
|
sawEvicted, otherErr := race.wait()
|
|
|
|
require.NoError(t, otherErr)
|
|
assert.True(
|
|
t, sawEvicted,
|
|
"a write after eviction must be refused",
|
|
)
|
|
assert.False(
|
|
t, w.HandleOpen(),
|
|
"no write may reopen the archive once the writer has "+
|
|
"been evicted",
|
|
)
|
|
assert.False(
|
|
t, eng.ExportHasArchiveWriter(event.WebhookID),
|
|
"the registry entry must stay gone",
|
|
)
|
|
}
|
|
|
|
// TestEvictWebhook_LaterDeliveryRecreatesWriter proves eviction
|
|
// does not break archiving for a webhook 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)
|
|
|
|
webhookDB := testWebhookDB(t)
|
|
event := seedEvent(t, webhookDB, `{"archived":true}`)
|
|
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
|
|
|
|
eng.ExportDeliverDatabase(webhookDB, d)
|
|
require.True(
|
|
t, eng.ExportHasArchiveWriter(event.WebhookID),
|
|
)
|
|
|
|
eng.EvictWebhook(event.WebhookID)
|
|
|
|
// A fresh delivery for the same webhook gets a brand new
|
|
// writer from the registry, so archiving keeps working.
|
|
second := seedDatabaseTargetDelivery(
|
|
t, webhookDB, event, "",
|
|
)
|
|
eng.ExportDeliverDatabase(webhookDB, second)
|
|
|
|
assert.True(
|
|
t, eng.ExportHasArchiveWriter(event.WebhookID),
|
|
"a later delivery should recreate the writer",
|
|
)
|
|
}
|
|
|
|
// TestEngineStop_WriteAfterStopIsRefused proves the engine's stop
|
|
// closes each archive writer the way deleting its webhook does: a
|
|
// write that reaches a writer after the stop is refused, reopens
|
|
// nothing and adds no row.
|
|
func TestEngineStop_WriteAfterStopIsRefused(t *testing.T) {
|
|
t.Parallel()
|
|
|
|
eng, _ := evictTestEngine(t)
|
|
|
|
webhookDB := testWebhookDB(t)
|
|
event := seedEvent(t, webhookDB, `{"archived":true}`)
|
|
d := seedDatabaseTargetDelivery(t, webhookDB, event, "")
|
|
|
|
eng.ExportDeliverDatabase(webhookDB, d)
|
|
|
|
w := eng.ExportArchiveWriterFor(event.WebhookID)
|
|
require.NotNil(t, w)
|
|
require.True(t, w.HandleOpen())
|
|
|
|
require.NoError(t, eng.ExportStop(context.Background()))
|
|
|
|
err := w.Write(evictTestRow("ev-after-stop"), 0)
|
|
|
|
require.ErrorIs(
|
|
t, err, delivery.ErrExportArchiveWriterEvicted,
|
|
"a write after the stop must be refused",
|
|
)
|
|
assert.False(
|
|
t, w.HandleOpen(),
|
|
"a refused write must not reopen the archive",
|
|
)
|
|
assert.False(
|
|
t, eng.ExportHasArchiveWriter(event.WebhookID),
|
|
"the stop should empty the registry",
|
|
)
|
|
|
|
count, err := countArchivedRows(w.Path())
|
|
require.NoError(t, err)
|
|
assert.Equal(
|
|
t, int64(1), count, "the refused row must not be written",
|
|
)
|
|
}
|