package delivery_test import ( "context" "database/sql" "fmt" "net/http" "os" "path/filepath" "sync" "testing" "time" "github.com/google/uuid" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "go.uber.org/fx" "gorm.io/driver/sqlite" "gorm.io/gorm" "gorm.io/gorm/clause" _ "modernc.org/sqlite" // Pure Go SQLite driver. "sneak.berlin/go/webhooker/internal/database" "sneak.berlin/go/webhooker/internal/delivery" ) const ( // sweepRowOld and sweepRowNew are the event ids // seedArchiveRows assigns to the first and second seeded // rows. sweepRowOld = "ev-0" sweepRowNew = "ev-1" // sweepConcurrentWrites is how many deliveries the // concurrent write-plus-sweep test races against the sweep. 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 { sweeper *delivery.ArchiveSweeper eng *delivery.Engine mainDB *database.Database dataDir string } func setupSweeperTest(t *testing.T) *sweeperEnv { t.Helper() dataDir := t.TempDir() log := archiveTestLogger() sqlDB, err := sql.Open( "sqlite", fmt.Sprintf( "file:%s?mode=rwc", filepath.Join(dataDir, "main.db"), ), ) require.NoError(t, err) t.Cleanup(func() { _ = sqlDB.Close() }) gdb, err := gorm.Open( sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, ) require.NoError(t, err) mainDB := database.NewTestDatabase(gdb) require.NoError(t, mainDB.Migrate()) eng := delivery.NewTestEngineWithDB( mainDB, database.NewTestWebhookDBManager(dataDir), log, &http.Client{Timeout: 5 * time.Second}, 1, ) return &sweeperEnv{ sweeper: delivery.NewTestArchiveSweeper( mainDB, eng, log, ), eng: eng, mainDB: mainDB, dataDir: dataDir, } } // archivePath returns where the engine keeps a webhook's // archive file. func (env *sweeperEnv) archivePath(webhookID string) string { return filepath.Join( env.dataDir, fmt.Sprintf("archive-%s.db", webhookID), ) } // seedDatabaseTarget creates a webhook with one database target // carrying the given target config JSON, and returns the // webhook id. func (env *sweeperEnv) seedDatabaseTarget( t *testing.T, configJSON string, ) string { t.Helper() wh := &database.Webhook{ UserID: uuid.New().String(), Name: "sweep-test", } require.NoError( t, env.mainDB.DB(). Omit(clause.Associations). Create(wh).Error, ) tgt := &database.Target{ WebhookID: wh.ID, Name: "archive", Type: database.TargetTypeDatabase, Active: true, Config: configJSON, } require.NoError( t, env.mainDB.DB(). Omit(clause.Associations). Create(tgt).Error, ) return wh.ID } // seedArchiveRows creates the archive file for a webhook 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, ) string { t.Helper() path := env.archivePath(webhookID) sqlDB, err := sql.Open( "sqlite", fmt.Sprintf("file:%s?mode=rwc", path), ) require.NoError(t, err) gdb, err := gorm.Open( sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, ) require.NoError(t, err) require.NoError( t, gdb.AutoMigrate(&delivery.ExportArchivedEvent{}), ) for i, at := range archivedAt { row := delivery.ExportArchivedEvent{ EventID: fmt.Sprintf("ev-%d", i), WebhookID: webhookID, Method: http.MethodPost, Body: `{"seeded":true}`, ArchivedAt: at, } require.NoError(t, gdb.Create(&row).Error) } require.NoError(t, sqlDB.Close()) return path } // archivedEventIDs returns the event ids currently stored in an // archive file, read through a separate read-only handle. func archivedEventIDs( t *testing.T, path string, ) []string { t.Helper() var rows []delivery.ExportArchivedEvent rdb := openArchiveDBForRead(t, path) require.NoError(t, rdb.Order("event_id").Find(&rows).Error) ids := make([]string, 0, len(rows)) for i := range rows { ids = append(ids, rows[i].EventID) } return ids } // countArchivedRows counts the rows in an archive file without // asserting anything, so it is safe to poll from an // assert.Eventually condition (which runs off the test // goroutine, where testify assertions must not be used). func countArchivedRows(path string) (int64, error) { sqlDB, err := sql.Open( "sqlite", fmt.Sprintf("file:%s?mode=ro", path), ) if err != nil { return 0, err } defer func() { _ = sqlDB.Close() }() gdb, err := gorm.Open( sqlite.Dialector{Conn: sqlDB}, &gorm.Config{}, ) if err != nil { return 0, err } var count int64 err = gdb.Model(&delivery.ExportArchivedEvent{}). Count(&count).Error if err != nil { return 0, err } return count, nil } // captureLifecycle is a minimal fx.Lifecycle that records the // hooks a component registers, so a test can invoke the real // OnStart/OnStop functions with a context of its choosing. type captureLifecycle struct { hooks []fx.Hook } func (l *captureLifecycle) Append(h fx.Hook) { l.hooks = append(l.hooks, h) } // TestArchiveSweeper_LoopOutlivesStartHookContext is the // regression test for a sweeper that never swept. fx calls // OnStart with a context carrying the application's start // timeout (15 seconds by default), so a background loop whose // context is derived from it is cancelled 15 seconds into the // process — three quarters of an hour before the first tick // under the default one-hour sweep interval. // // The hook context here is already cancelled, which is the same // defect taken to its limit: a loop that inherits it never runs // a single tick, while a correctly rooted loop keeps sweeping // for as long as the process lives. Handing the hook a plain // context.Background() would assert nothing at all. func TestArchiveSweeper_LoopOutlivesStartHookContext( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) now := time.Now() path := env.seedArchiveRows( t, webhookID, now.Add(-48*time.Hour), now.Add(-time.Minute), ) env.sweeper.ExportSetInterval(10 * time.Millisecond) // Drive the genuine fx hooks the application registers, // rather than a test-only entry point. lc := &captureLifecycle{} env.sweeper.ExportRegisterHooks(lc) require.Len(t, lc.hooks, 1) hookCtx, cancel := context.WithCancel(context.Background()) cancel() require.NoError(t, lc.hooks[0].OnStart(hookCtx)) t.Cleanup(func() { _ = lc.hooks[0].OnStop(context.Background()) }) assert.Eventually( t, func() bool { count, err := countArchivedRows(path) return err == nil && count == 1 }, 5*time.Second, 10*time.Millisecond, "the sweep loop must keep running after the start "+ "hook's context is done; it pruned nothing, so it "+ "inherited the hook context and died", ) } // 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. func TestArchiveSweep_DoesNotResurrectEvictedWriter( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( t, webhookID, 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 // the deletion committed. _, err := env.eng.ExportEnsureArchiveWriter(webhookID) require.NoError(t, err) env.eng.EvictWebhook(webhookID) require.False(t, env.eng.ExportHasArchiveWriter(webhookID)) env.sweeper.ExportSweep(context.Background()) assert.False( t, env.eng.ExportHasArchiveWriter(webhookID), "a sweep must never re-register a writer for a webhook "+ "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 // registry keeps holding only writers a delivery created and an // eviction can reach. func TestArchiveSweep_LeavesNoRegistryEntry(t *testing.T) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path := env.seedArchiveRows( t, webhookID, time.Now().Add(-48*time.Hour), time.Now().Add(-time.Minute), ) require.False(t, env.eng.ExportHasArchiveWriter(webhookID)) env.sweeper.ExportSweep(context.Background()) assert.Equal( t, []string{sweepRowNew}, archivedEventIDs(t, path), "the sweep must still prune an idle archive", ) assert.False( t, env.eng.ExportHasArchiveWriter(webhookID), "the sweep must release the registry entry it created", ) } // TestArchiveSweep_KeepsWriterAdoptedByDelivery is the other // half of that invariant: an entry the sweep created but a // delivery then claimed belongs to the registry and must survive // the sweep, or the delivery would be left holding a detached // writer with an open handle that no eviction can reach. func TestArchiveSweep_KeepsWriterAdoptedByDelivery( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( t, webhookID, 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"}`, ) env.sweeper.ExportSweep(context.Background()) require.False(t, env.eng.ExportHasArchiveWriter(webhookID)) env.eng.ExportDeliverDatabase(webhookDB, d) assert.True( t, env.eng.ExportHasArchiveWriter(webhookID), "a delivery's writer must stay registered", ) env.sweeper.ExportSweep(context.Background()) assert.True( t, env.eng.ExportHasArchiveWriter(webhookID), "a sweep must not drop a writer a delivery owns", ) } // TestArchiveSweep_KeepsWriterAdoptedDuringSweep covers the one // interleaving the sweepOwned flag exists for, which // TestArchiveSweep_KeepsWriterAdoptedByDelivery cannot reach: a // delivery adopting the sweep's own entry WHILE that sweep is // still running. // // The registry operations are driven directly, in the order the // sweep and a concurrent delivery perform them, so the window is // exercised deterministically rather than hoped for: // // 1. the sweep finds no cached writer and registers one of its // own, marked sweep-owned; // 2. a delivery arrives, is handed that very writer, clears the // flag and opens the archive handle; // 3. the sweep finishes and releases what it created. // // Step 3 must leave the entry alone. Dropping it would detach a // writer that is holding an open archive handle inside its // debounce window, and no eviction could ever reach it again — // exactly the process-lifetime handle leak this change exists to // close. The eviction at the end proves the entry is still // reachable. func TestArchiveSweep_KeepsWriterAdoptedDuringSweep( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( t, webhookID, time.Now().Add(-48*time.Hour), ) sweepWriter, created, err := env.eng.ExportSweepWriterFor( webhookID, ) require.NoError(t, err) require.True( t, created, "the sweep must have created the registry entry itself", ) // 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"}`, ) env.eng.ExportDeliverDatabase(webhookDB, d) adopted := env.eng.ExportArchiveWriterFor(webhookID) 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), "the delivery leaves the archive handle open", ) // The sweep finishes. env.eng.ExportReleaseSweepWriter(webhookID, sweepWriter) require.True( t, env.eng.ExportHasArchiveWriter(webhookID), "a writer adopted by a delivery during a sweep must "+ "stay registered, or its open handle is unreachable", ) env.eng.EvictWebhook(webhookID) assert.False( t, env.eng.ExportHasArchiveWriter(webhookID), "the adopted writer must still be evictable", ) assert.False( t, sweepWriter.HandleOpen(), "eviction must have closed the adopted writer's handle", ) } // TestArchiveSweep_ContinuesAfterPerWebhookFailure proves a // failure for one webhook 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( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) // Seeded first so the sweep reaches them before the healthy // webhook: targets come back in insertion order. badConfigID := env.seedDatabaseTarget(t, `{"expiry":"!!!"}`) env.seedArchiveRows( t, badConfigID, time.Now().Add(-48*time.Hour), ) corruptID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) require.NoError(t, os.WriteFile( env.archivePath(corruptID), []byte("this is not a sqlite database"), 0o600, )) healthyID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) healthyPath := env.seedArchiveRows( t, healthyID, time.Now().Add(-48*time.Hour), time.Now().Add(-time.Minute), ) env.sweeper.ExportSweep(context.Background()) assert.Equal( t, []string{sweepRowNew}, archivedEventIDs(t, healthyPath), "a failure for an earlier webhook 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 // protects the window between that stat and the open. Flipping // the sweep's mode to create-if-missing makes this fail. func TestArchiveSweep_OpenExistingDoesNotCreateFile( t *testing.T, ) { t.Parallel() dir := t.TempDir() path := filepath.Join(dir, "archive-absent.db") w := delivery.NewExportArchiveWriter( path, archiveTestLogger(), 0, ) err := w.OpenExisting(time.Hour) require.Error( t, err, "opening a missing archive without create permission "+ "must fail rather than conjure the file", ) for _, suffix := range archiveFileSuffixes() { assert.NoFileExists(t, path+suffix) } } // TestArchiveSweep_PrunesIdleArchive is the core regression // test for this issue: an archive that receives no further // writes must still lose its expired rows. Before the sweeper // existed, pruning only ever ran on a write-triggered reopen, // so an idle archive kept expired rows forever. func TestArchiveSweep_PrunesIdleArchive(t *testing.T) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) now := time.Now() path := env.seedArchiveRows( t, webhookID, now.Add(-48*time.Hour), now.Add(-time.Minute), ) require.Equal( t, []string{sweepRowOld, sweepRowNew}, archivedEventIDs(t, path), ) env.sweeper.ExportSweep(context.Background()) assert.Equal( t, []string{sweepRowNew}, archivedEventIDs(t, path), "the sweep should prune rows older than the expiry "+ "from an idle archive and keep the rest", ) } // TestArchiveSweep_LeavesArchiveClosed proves the sweep does // not hold the archive open afterwards, so an operator can // still move the file away for offline retention. // // The assertion is made on a writer the test holds a reference // to, and the handle is proven OPEN before the sweep runs, so the // test observes the sweep closing it rather than a writer that // merely never opened anything. Asking the registry instead would // be vacuous here: the sweep releases an entry it created, and a // missing entry reports "not open" whether or not anything was // closed. func TestArchiveSweep_LeavesArchiveClosed(t *testing.T) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path := env.seedArchiveRows( t, webhookID, time.Now().Add(-48*time.Hour), ) w := delivery.NewExportArchiveWriter( path, archiveTestLogger(), 0, ) require.NoError(t, w.OpenExisting(time.Hour)) require.True( t, w.HandleOpen(), "the writer must hold an open handle before the sweep", ) require.NoError(t, w.SweepExpired(time.Hour)) assert.False( t, w.HandleOpen(), "an idle archive must end the sweep closed", ) } // TestArchiveSweep_ClosesHandleOfRegisteredWriter states the same // guarantee end to end, through the real sweeper and a writer the // registry keeps. // // The delivery leaves the archive handle open inside its debounce // window and makes the entry delivery-owned, so the sweep finds a // cached writer (created is false, nothing is released) and the // registry query afterwards is answered by a writer that really // exists. A handle left open here would be doubly wrong: it also // blocks the operator's move-the-file-away workflow. func TestArchiveSweep_ClosesHandleOfRegisteredWriter( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( t, webhookID, 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"}`, ) env.eng.ExportDeliverDatabase(webhookDB, d) require.True( t, env.eng.ExportArchiveHandleOpen(webhookID), "the delivery must leave the archive handle open", ) env.sweeper.ExportSweep(context.Background()) require.True( t, env.eng.ExportHasArchiveWriter(webhookID), "the delivery's registry entry must survive the sweep", ) assert.False( t, env.eng.ExportArchiveHandleOpen(webhookID), "the sweep must leave the archive closed", ) } // TestArchiveSweep_NeverExpiryUntouched proves the sweep is a // no-op for the default retention policy, so archives with no // expiry (or the literal "never") behave exactly as before. func TestArchiveSweep_NeverExpiryUntouched(t *testing.T) { t.Parallel() for _, configJSON := range []string{ `{"expiry":"never"}`, `{"expiry":""}`, "", } { env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, configJSON) path := env.seedArchiveRows( t, webhookID, time.Now().Add(-10000*time.Hour), ) env.sweeper.ExportSweep(context.Background()) assert.Equal( t, []string{sweepRowOld}, archivedEventIDs(t, path), "config %q must keep rows forever", configJSON, ) assert.False( t, env.eng.ExportHasArchiveWriter(webhookID), "config %q must leave no registry entry behind", configJSON, ) } } // TestArchiveSweep_NeverExpirySkipsBeforeOpening pins the // expiry <= 0 boundary in sweepTarget, which the row assertions // above cannot reach: pruning is separately gated on a positive // expiry, so a "never" archive keeps its rows even if the sweep // does open it. // // The spec is stronger than that — a "never" archive is skipped // before any file is touched — so the archive here exists but has // never been migrated. Opening it at all would run AutoMigrate // and create the archive table, which is exactly what must not // happen. func TestArchiveSweep_NeverExpirySkipsBeforeOpening( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"never"}`) path := env.archivePath(webhookID) seedUnmigratedArchive(t, path) require.False(t, archiveTableExists(t, path)) env.sweeper.ExportSweep(context.Background()) assert.False( t, archiveTableExists(t, path), "a never-expiry archive must not be opened at all", ) } // seedUnmigratedArchive creates an archive file that exists but // carries no archive schema, so any open of it is observable: the // archive table appears only if something ran AutoMigrate. func seedUnmigratedArchive(t *testing.T, path string) { t.Helper() sqlDB, err := sql.Open( "sqlite", fmt.Sprintf("file:%s?mode=rwc", path), ) require.NoError(t, err) _, err = sqlDB.ExecContext( t.Context(), "CREATE TABLE placeholder (id INTEGER)", ) require.NoError(t, err) require.NoError(t, sqlDB.Close()) } // archiveTableExists reports whether an archive file has had the // archive schema migrated into it. func archiveTableExists(t *testing.T, path string) bool { t.Helper() return openArchiveDBForRead(t, path). Migrator(). HasTable(&delivery.ExportArchivedEvent{}) } // 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. func TestArchiveSweep_DoesNotCreateArchiveFile(t *testing.T) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path := env.archivePath(webhookID) require.NoFileExists(t, path) env.sweeper.ExportSweep(context.Background()) for _, suffix := range archiveFileSuffixes() { assert.NoFileExists( t, path+suffix, "the sweep must not create an archive file", ) } } // TestArchiveSweep_DoesNotCreateAfterWriterExists covers the // same guarantee once a writer is cached in the registry but // the file itself is still absent (for instance because the // operator moved the archive away). func TestArchiveSweep_DoesNotCreateAfterWriterExists( t *testing.T, ) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path, err := env.eng.ExportEnsureArchiveWriter(webhookID) require.NoError(t, err) require.NoFileExists(t, path) env.sweeper.ExportSweep(context.Background()) assert.NoFileExists(t, path) } // TestArchiveSweep_SkipsDeletedWebhookTargets proves that the // sweep ignores targets soft-deleted along with their webhook, // so a deleted webhook's archive is never reopened. func TestArchiveSweep_SkipsDeletedWebhookTargets(t *testing.T) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) path := env.seedArchiveRows( t, webhookID, time.Now().Add(-48*time.Hour), ) require.NoError( t, env.mainDB.DB(). Where("webhook_id = ?", webhookID). Delete(&database.Target{}).Error, ) env.sweeper.ExportSweep(context.Background()) assert.Equal( t, []string{sweepRowOld}, archivedEventIDs(t, path), "a deleted target's archive must be left alone", ) } // TestArchiveSweep_ConcurrentWrites proves the sweep serialises // against writes through the per-webhook writer mutex. Run // under -race, an unsynchronised sweep would be caught here. func TestArchiveSweep_ConcurrentWrites(t *testing.T) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) webhookDB := testWebhookDB(t) // The deliveries are seeded up front, on the test's own // goroutine: the seed helpers assert, and testify assertions // must not run off the test goroutine. deliveries := make( []*database.Delivery, 0, sweepConcurrentWrites, ) for range sweepConcurrentWrites { event := seedEvent(t, webhookDB, `{"n":1}`) event.WebhookID = webhookID deliveries = append( deliveries, seedDatabaseTargetDelivery( t, webhookDB, event, `{"expiry":"1h"}`, ), ) } var wg sync.WaitGroup wg.Add(2) go func() { defer wg.Done() for _, d := range deliveries { env.eng.ExportDeliverDatabase(webhookDB, d) } }() go func() { defer wg.Done() for range sweepConcurrentWrites { env.sweeper.ExportSweep(context.Background()) } }() wg.Wait() assert.FileExists(t, env.archivePath(webhookID)) } // TestArchiveSweeper_StopsCleanly proves the background loop // exits on OnStop rather than leaking a goroutine. func TestArchiveSweeper_StopsCleanly(t *testing.T) { t.Parallel() env := setupSweeperTest(t) webhookID := env.seedDatabaseTarget(t, `{"expiry":"1h"}`) env.seedArchiveRows( t, webhookID, time.Now().Add(-48*time.Hour), ) env.sweeper.ExportSetInterval(time.Millisecond) env.sweeper.ExportStart() // stop blocks on the loop's WaitGroup, so returning at all // proves the loop observed the cancellation and exited. env.sweeper.ExportStop() }