From 2a494bdae06825880d778c6dfe0f8c098cfc00bc Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Fri, 25 Sep 2026 08:27:13 -0400 Subject: [PATCH 1/6] fix(mq): refuse an embedded store it cannot create at once A store directory the embedded server cannot create (a regular file in its place) failed JetStream in the background, and NewEmbedded only gave up after ReadyForConnections' full 5s wait, reporting "nats server not ready" instead of the cause. Create the directory first and return its error. TestNew_LateBootFailureReleasesEverything in internal/app used exactly this failure and spent 5s of its package's 15s budget on it. Part of #617. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- internal/mq/embedded.go | 6 ++++++ internal/mq/embedded_test.go | 11 +++++++++++ 2 files changed, 17 insertions(+) diff --git a/internal/mq/embedded.go b/internal/mq/embedded.go index 830be900..ff726497 100644 --- a/internal/mq/embedded.go +++ b/internal/mq/embedded.go @@ -7,6 +7,7 @@ import ( "log/slog" "maps" "math" + "os" "slices" "strings" "sync" @@ -146,6 +147,11 @@ var errNoQueue = errors.New("no queue is open for it yet") // applied, or by a publish or park that finds it missing, at the budget last // asked for it. The server logs through slog's default logger. func NewEmbedded(storeDir string) (*EmbeddedNATS, error) { + // A store the server cannot create fails JetStream in the background, and + // ReadyForConnections would only give up on it after its whole wait. + if err := os.MkdirAll(storeDir, 0o750); err != nil { + return nil, fmt.Errorf("nats store: %w", err) + } opts := &natsserver.Options{ DontListen: true, JetStream: true, diff --git a/internal/mq/embedded_test.go b/internal/mq/embedded_test.go index 2a7c483a..32447943 100644 --- a/internal/mq/embedded_test.go +++ b/internal/mq/embedded_test.go @@ -1290,6 +1290,17 @@ func TestEmbeddedNATS_ReplaySince_IsPerTenant(t *testing.T) { assert.Equal(t, []string{"acme1", "acme2"}, got) } +// A store directory that cannot be created refuses the boot at once, rather +// than after the server's whole wait for a JetStream that will never start. +func TestNewEmbedded_AStoreItCannotCreateFailsAtOnce(t *testing.T) { + file := filepath.Join(t.TempDir(), "nats") + require.NoError(t, os.WriteFile(file, nil, 0o600)) + start := time.Now() + _, err := NewEmbedded(file) + require.Error(t, err) + assert.Less(t, time.Since(start), 3*time.Second) +} + // A boot over a directory an earlier build wrote deletes the pair of streams // it kept for every tenant together: their subjects overlap every tenant's, // so no tenant's queue could open beside them. From 98acc11d49af1b7181bd4e5a248ed215aa00de6e Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Fri, 25 Sep 2026 08:27:44 -0400 Subject: [PATCH 2/6] test(mq): run the embedded broker's tests in parallel, without fsync internal/mq's unit binary took 13s alone and 18s under a parallel `make test-unit`, against the 15s per-package budget (#617). Two costs dominated: - An fsync per JetStream write (SyncAlways), which on macOS is a full flush and was over half the run. It is now EmbeddedSyncAlways, true in production and turned off by TestMain: nothing here asserts anything across a crash. - The embedded tests ran one after another, each on its own in-process server in its own directory. They now call t.Parallel; the ones that capture the default logger stay serial, so no captured log gains another test's lines. Running in parallel exposed two races that load alone had hidden: - A store's TempDir removal racing a consumer's state file written after Close (#442): the store now lives in storeDir, the retrying removal internal/testutil and mqtest already use. - TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen opened globex while the server was still removing the directories acme's failed open had emptied, on a goroutine of its own. The test now waits for that removal. It cannot keep the directory occupied the way the pacing test does: a queue open first would keep the store reservation count above zero, and the test would no longer catch a store limit at the top of the int64 range (checked by mutation). Alone: 13.0s -> ~3.0s. Part of #617. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- internal/mq/embedded.go | 7 ++- internal/mq/embedded_test.go | 104 +++++++++++++++++++++++++++++------ internal/mq/main_test.go | 3 +- 3 files changed, 96 insertions(+), 18 deletions(-) diff --git a/internal/mq/embedded.go b/internal/mq/embedded.go index ff726497..d4ac14b9 100644 --- a/internal/mq/embedded.go +++ b/internal/mq/embedded.go @@ -132,6 +132,11 @@ const ( reopenRetry = 5 * time.Second ) +// EmbeddedSyncAlways is NewEmbedded's SyncAlways. Only a TestMain may turn it +// off, before any broker starts: unit tests assert nothing across a crash, and +// on macOS an fsync per write is most of their run time (#617). +var EmbeddedSyncAlways = true + // errNoQueue is why a publish or park finds no queue it can open: no budget // has been asked for the tenant yet (see SetMaxBytes). Publish reports it as // ErrQueueFull. @@ -156,7 +161,7 @@ func NewEmbedded(storeDir string) (*EmbeddedNATS, error) { DontListen: true, JetStream: true, StoreDir: storeDir, - SyncAlways: true, // fsync every JetStream write — publish ACKs only after data is on disk + SyncAlways: EmbeddedSyncAlways, // fsync every JetStream write — publish ACKs only after data is on disk // Without NoSigs, Start() installs a process-wide SIGINT handler that // races the app's graceful shutdown (double Shutdown → "close of nil // channel" panic) and os.Exit(0)s past its cleanup. WaveHouse owns diff --git a/internal/mq/embedded_test.go b/internal/mq/embedded_test.go index 32447943..cfe9b30c 100644 --- a/internal/mq/embedded_test.go +++ b/internal/mq/embedded_test.go @@ -19,6 +19,26 @@ import ( // testBudget is the byte budget newTestEmbedded opens each queue at. const testBudget = 64 << 20 +// storeDir is a temporary directory for a broker's store whose removal +// retries briefly: a consumer's state file can land after Close has returned, +// which fails t.TempDir's one-shot RemoveAll (#442). The retrying cleanup runs +// first (cleanups are LIFO), leaving t.TempDir an empty directory to remove. +func storeDir(t *testing.T) string { + t.Helper() + dir := filepath.Join(t.TempDir(), "store") + t.Cleanup(func() { + var err error + for range 50 { + if err = os.RemoveAll(dir); err == nil { + return + } + time.Sleep(20 * time.Millisecond) + } + t.Errorf("remove %s: %v", dir, err) + }) + return dir +} + // openEmbedded starts an EmbeddedNATS over dir, closed by the test framework. func openEmbedded(t *testing.T, dir string) *EmbeddedNATS { t.Helper() @@ -33,7 +53,7 @@ func openEmbedded(t *testing.T, dir string) *EmbeddedNATS { // at testBudget. func newTestEmbedded(t *testing.T, tenants ...tenant.ID) *EmbeddedNATS { t.Helper() - e := openEmbedded(t, t.TempDir()) + e := openEmbedded(t, storeDir(t)) if len(tenants) == 0 { tenants = []tenant.ID{tenant.Default} } @@ -73,8 +93,7 @@ func ackAll(t *testing.T, e *EmbeddedNATS, consumer string, n int) { } func TestEmbeddedNATS_PublishSubscribe(t *testing.T) { - // No t.Parallel(): each embedded server uses DontListen+InProcessServer, - // but starting several in parallel still slows tests unnecessarily. + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) @@ -113,6 +132,7 @@ func TestEmbeddedNATS_PublishSubscribe(t *testing.T) { } func TestEmbeddedNATS_Stats(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) stats, err := e.Stats() @@ -123,6 +143,7 @@ func TestEmbeddedNATS_Stats(t *testing.T) { } func TestEmbeddedNATS_PublishHeaders(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -145,7 +166,8 @@ func TestEmbeddedNATS_PublishHeaders(t *testing.T) { // subjects alone at the budget, refusing when full, and a dead-letter stream // at a tenth of it, dropping its oldest when full. No other tenant gets one. func TestEmbeddedNATS_SetMaxBytes_OpensTheTenantsQueue(t *testing.T) { - e := openEmbedded(t, t.TempDir()) + t.Parallel() + e := openEmbedded(t, storeDir(t)) assert.Zero(t, e.MaxBytes("acme"), "no budget applied yet") require.NoError(t, e.SetMaxBytes(t.Context(), "acme", testBudget)) @@ -168,6 +190,7 @@ func TestEmbeddedNATS_SetMaxBytes_OpensTheTenantsQueue(t *testing.T) { } func TestEmbeddedNATS_StreamHandle(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -255,6 +278,7 @@ func TestEmbeddedNATS_StreamHandle(t *testing.T) { // MaxAckPending (ingest backpressure, per tenant) are checkable nowhere else, // and a dropped field would compile and pass every delivery test. func TestEmbeddedNATS_CreateConsumer_Config(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -279,6 +303,7 @@ func TestEmbeddedNATS_CreateConsumer_Config(t *testing.T) { } func TestEmbeddedNATS_ReplaySince(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -314,14 +339,16 @@ func TestEmbeddedNATS_ReplaySince(t *testing.T) { } func TestEmbeddedNATS_DefaultLogger(t *testing.T) { + t.Parallel() // NewEmbedded without a logger should not panic — it falls back to the // default slog logger. - e, err := NewEmbedded(t.TempDir()) + e, err := NewEmbedded(storeDir(t)) require.NoError(t, err) t.Cleanup(func() { _ = e.Close() }) } func TestEmbeddedNATS_SubscribeCancellation(t *testing.T) { + t.Parallel() // When the caller's context is cancelled, the consume loop should stop // cleanly without leaking goroutines or blocking. e := newTestEmbedded(t) @@ -355,6 +382,7 @@ func TestSlogNATSLogger_Levels(t *testing.T) { } func TestEmbeddedNATS_SetMaxBytes(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -384,6 +412,7 @@ func TestEmbeddedNATS_SetMaxBytes(t *testing.T) { } func TestEmbeddedNATS_SetMaxBytes_DLQFailureRollsBackIngest(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -408,6 +437,7 @@ func TestEmbeddedNATS_SetMaxBytes_DLQFailureRollsBackIngest(t *testing.T) { } func TestEmbeddedNATS_SetMaxBytes_IngestFailureChangesNothing(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithCancel(t.Context()) cancel() // a stop caught mid-reload: the first JetStream call gives up @@ -427,7 +457,8 @@ func TestEmbeddedNATS_SetMaxBytes_IngestFailureChangesNothing(t *testing.T) { // cause is gone, a reload opens the queue at the budget last asked for it, // however recently a publish tried. func TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen(t *testing.T) { - dir := t.TempDir() + t.Parallel() + dir := storeDir(t) // The dead-letter stream is the first of the pair to open. A failed open // removes what was in the way, so the obstacle is put back before each // attempt meant to fail. @@ -453,6 +484,16 @@ func TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen(t *testing.T) { require.Error(t, e.SetMaxBytes(ctx, "acme", testBudget)) assert.Zero(t, e.MaxBytes("acme"), "no budget applied") + // The server removes the emptied streams and account directories on a + // goroutine of its own after the failed open, and globex's open must not + // race it (see the pacing test below). No queue may be open first to keep + // them: the reservation count has to be at zero when the failed open + // releases one it never made. + account := filepath.Dir(filepath.Dir(block)) + require.Eventually(t, func() bool { + _, err := os.Stat(account) + return os.IsNotExist(err) + }, 5*time.Second, 5*time.Millisecond, "the failed open's cleanup never removed %s", account) require.NoError(t, e.SetMaxBytes(ctx, "globex", testBudget), "one tenant's failed open costs the next nothing") require.NoError(t, e.Publish(ctx, Topic{Tenant: "globex", Table: "t"}, []byte("x"))) @@ -476,7 +517,8 @@ func TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen(t *testing.T) { // broken queue would otherwise hold the lock that every other tenant's open, // resize and reload takes. Once the window has passed, a publish tries again. func TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen(t *testing.T) { - dir := t.TempDir() + t.Parallel() + dir := storeDir(t) block := filepath.Join(dir, "jetstream", "$G", "streams", dlqStreamName("acme")) obstruct := func() { t.Helper() @@ -540,7 +582,8 @@ func TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen(t *testing.T) { // not by the stream answering: it opens the queue properly first, consumers // joined, so its row reaches them rather than a stream nobody reads. func TestEmbeddedNATS_Publish_OpensAQueueItsOpenGaveUpOn(t *testing.T) { - dir := t.TempDir() + t.Parallel() + dir := storeDir(t) block := filepath.Join(dir, "jetstream", "$G", "streams", ingestStreamName("acme")) require.NoError(t, os.MkdirAll(filepath.Dir(block), 0o750)) require.NoError(t, os.WriteFile(block, nil, 0o600)) @@ -581,9 +624,10 @@ func TestEmbeddedNATS_Publish_OpensAQueueItsOpenGaveUpOn(t *testing.T) { // that found the pair split applied none, and a cap of 0 would leave the // ingest stream with no cap at all. func TestEmbeddedNATS_SetMaxBytes_UndoRestoresTheIngestStreamsCap(t *testing.T) { + t.Parallel() ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() - dir := t.TempDir() + dir := storeDir(t) first, err := NewEmbedded(dir) require.NoError(t, err) require.NoError(t, first.SetMaxBytes(ctx, "acme", 8<<20)) @@ -605,7 +649,8 @@ func TestEmbeddedNATS_SetMaxBytes_UndoRestoresTheIngestStreamsCap(t *testing.T) // otherwise let the tenant's ingest answer 200 for rows nobody reads. The // queue itself is open, so SetMaxBytes succeeds. func TestEmbeddedNATS_Consume_ReportsAQueueItCannotJoin(t *testing.T) { - e := openEmbedded(t, t.TempDir()) + t.Parallel() + e := openEmbedded(t, storeDir(t)) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() // A durable name the client refuses: with no queue yet, nothing checks it. @@ -629,7 +674,8 @@ func TestEmbeddedNATS_Consume_ReportsAQueueItCannotJoin(t *testing.T) { // would have DiscardOld delete the oldest parked rows to fit (#532), so the // stream keeps what it holds, capped at that, and every row survives. func TestEmbeddedNATS_SetMaxBytes_NeverShrinksTheDeadLetterQueueBelowWhatItHolds(t *testing.T) { - e := openEmbedded(t, t.TempDir()) + t.Parallel() + e := openEmbedded(t, storeDir(t)) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() require.NoError(t, e.SetMaxBytes(ctx, "acme", 10<<20)) @@ -662,6 +708,7 @@ func TestEmbeddedNATS_SetMaxBytes_NeverShrinksTheDeadLetterQueueBelowWhatItHolds } func TestEmbeddedNATS_ReplaySince_PullFailureIsAnError(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -683,6 +730,7 @@ func TestEmbeddedNATS_ReplaySince_PullFailureIsAnError(t *testing.T) { } func TestEmbeddedNATS_ReplaySince_StopsWhenContextIsDone(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithCancel(t.Context()) defer cancel() @@ -703,6 +751,7 @@ func TestEmbeddedNATS_ReplaySince_StopsWhenContextIsDone(t *testing.T) { } func TestEmbeddedNATS_DeadLetter(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, tenant.Default, "acme") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -757,6 +806,7 @@ func TestEmbeddedNATS_DeadLetter(t *testing.T) { // A tenant with no queue — one never given a budget on this data directory — // has nothing parked, which is not the same as a failed read. func TestEmbeddedNATS_DeadLetterCounts_NoQueue(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -774,6 +824,7 @@ func TestEmbeddedNATS_DeadLetterCounts_NoQueue(t *testing.T) { } func TestEmbeddedNATS_DeadLetterCounts_BrokerFailureIsNotAnEmptyQueue(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -787,6 +838,7 @@ func TestEmbeddedNATS_DeadLetterCounts_BrokerFailureIsNotAnEmptyQueue(t *testing } func TestEmbeddedNATS_DeadLetter_IsAPrefixSwap(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "a") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -823,6 +875,7 @@ func TestEmbeddedNATS_DeadLetter_IsAPrefixSwap(t *testing.T) { // tenant's queue, at a tenth of the budget last asked for it, rather than // leaving the row to be redelivered. func TestEmbeddedNATS_DeadLetter_ReopensAMissingQueue(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -836,7 +889,8 @@ func TestEmbeddedNATS_DeadLetter_ReopensAMissingQueue(t *testing.T) { } func TestEmbeddedNATS_Publish_QueueFull(t *testing.T) { - e := openEmbedded(t, t.TempDir()) + t.Parallel() + e := openEmbedded(t, storeDir(t)) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() require.NoError(t, e.SetMaxBytes(ctx, "acme", 4<<10)) @@ -862,6 +916,7 @@ func TestEmbeddedNATS_Publish_QueueFull(t *testing.T) { // it missing, and a tenant never given a budget has no queue to publish to: // that is refused as a full queue, and nothing is opened for it. func TestEmbeddedNATS_Publish_OpensTheQueueAtTheLastBudget(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -886,6 +941,7 @@ func TestEmbeddedNATS_Publish_OpensTheQueueAtTheLastBudget(t *testing.T) { // report as its delivery ending. So the reopen — joins included — outlives // the caller's cancellation. func TestEmbeddedNATS_ReopenOutlivesTheCallersCancellation(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -923,6 +979,7 @@ func TestEmbeddedNATS_ReopenOutlivesTheCallersCancellation(t *testing.T) { } func TestEmbeddedNATS_PurgeAcked(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -981,6 +1038,7 @@ func TestEmbeddedNATS_PurgeAcked(t *testing.T) { // goes, and a tenant the cutoffs do not name — one no longer served — keeps // no history at all. func TestEmbeddedNATS_PurgeAcked_EachTenantAtItsOwnCutoff(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex", "initech") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -1023,6 +1081,7 @@ func TestEmbeddedNATS_PurgeAcked_EachTenantAtItsOwnCutoff(t *testing.T) { // after one whose durable is gone and one whose stream is. A sweep whose // context has already ended touches no tenant. func TestEmbeddedNATS_PurgeAcked_OneTenantsFailureStopsNoOther(t *testing.T) { + t.Parallel() ids := []tenant.ID{"acme", "globex", "initech", "umbrella"} e := newTestEmbedded(t, ids...) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) @@ -1072,6 +1131,7 @@ func TestEmbeddedNATS_PurgeAcked_OneTenantsFailureStopsNoOther(t *testing.T) { // whose handler is stuck, holds back its own delivery and no other tenant's — // each tenant's messages arrive on a delivery of their own, in order. func TestEmbeddedNATS_Consume_OneTenantsBacklogDoesNotHoldAnother(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex", "initech") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -1127,6 +1187,7 @@ func TestEmbeddedNATS_Consume_OneTenantsBacklogDoesNotHoldAnother(t *testing.T) // consumer paths deliver its events as they do the queues that were there // first, whether those were opened in this process or found on disk. func TestEmbeddedNATS_ConsumersJoinQueuesOpenedLater(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -1166,6 +1227,7 @@ func TestEmbeddedNATS_ConsumersJoinQueuesOpenedLater(t *testing.T) { } func TestEmbeddedNATS_Consume_ReportsDeliveryEndingOnItsOwn(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex") ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) defer cancel() @@ -1198,6 +1260,7 @@ func TestEmbeddedNATS_Consume_ReportsDeliveryEndingOnItsOwn(t *testing.T) { } func TestEmbeddedNATS_Consume_StopIsNotAFailure(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -1242,6 +1305,7 @@ func TestFanIn_SharesThePrefetch(t *testing.T) { // tenants' queues, like the worker's prefetch, so what it holds client-side // does not grow with the number of tenants. func TestEmbeddedNATS_Subscribe_SharesTheClientDefault(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex") ctx, cancel := context.WithCancel(t.Context()) defer cancel() @@ -1257,6 +1321,7 @@ func TestEmbeddedNATS_Subscribe_SharesTheClientDefault(t *testing.T) { // Nothing lands on the default tenant by omission (#583): the tenant is a // required token, checked against its grammar before anything is sent. func TestEmbeddedNATS_Publish_RefusesATopicWithoutATenant(t *testing.T) { + t.Parallel() e := newTestEmbedded(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -1274,6 +1339,7 @@ func TestEmbeddedNATS_Publish_RefusesATopicWithoutATenant(t *testing.T) { // Two tenants, one table name: a replay of one never carries the other's rows. func TestEmbeddedNATS_ReplaySince_IsPerTenant(t *testing.T) { + t.Parallel() e := newTestEmbedded(t, "acme", "globex") ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() @@ -1293,6 +1359,7 @@ func TestEmbeddedNATS_ReplaySince_IsPerTenant(t *testing.T) { // A store directory that cannot be created refuses the boot at once, rather // than after the server's whole wait for a JetStream that will never start. func TestNewEmbedded_AStoreItCannotCreateFailsAtOnce(t *testing.T) { + t.Parallel() file := filepath.Join(t.TempDir(), "nats") require.NoError(t, os.WriteFile(file, nil, 0o600)) start := time.Now() @@ -1305,7 +1372,8 @@ func TestNewEmbedded_AStoreItCannotCreateFailsAtOnce(t *testing.T) { // it kept for every tenant together: their subjects overlap every tenant's, // so no tenant's queue could open beside them. func TestNewEmbedded_DeletesTheStreamsAnEarlierBuildShared(t *testing.T) { - dir := t.TempDir() + t.Parallel() + dir := storeDir(t) old, err := NewEmbedded(dir) require.NoError(t, err) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) @@ -1333,9 +1401,10 @@ func TestNewEmbedded_DeletesTheStreamsAnEarlierBuildShared(t *testing.T) { // again; a dead-letter stream kept above its tenth because it holds more (the // shrink guard) is at its budget and left as it is. func TestNewEmbedded_ASplitPairIsAppliedAgainAtBoot(t *testing.T) { + t.Parallel() ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() - dir := t.TempDir() + dir := storeDir(t) first, err := NewEmbedded(dir) require.NoError(t, err) for _, id := range []tenant.ID{"split", "gone", "guarded"} { @@ -1372,7 +1441,8 @@ func TestNewEmbedded_ASplitPairIsAppliedAgainAtBoot(t *testing.T) { // tenant no longer served, which is never given a budget again, included — // so what such a tenant had queued still reaches the worker. func TestNewEmbedded_TakesStockOfTheQueuesOnDisk(t *testing.T) { - dir := t.TempDir() + t.Parallel() + dir := storeDir(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() first, err := NewEmbedded(dir) @@ -1413,6 +1483,7 @@ func TestNewEmbedded_TakesStockOfTheQueuesOnDisk(t *testing.T) { // updated in place when they differ; either way delivery resumes past what it // acknowledged before the restart. func TestEmbeddedNATS_ADurableOnDiskIsReusedAcrossARestart(t *testing.T) { + t.Parallel() for _, tt := range []struct { name string maxAckPending int @@ -1421,7 +1492,8 @@ func TestEmbeddedNATS_ADurableOnDiskIsReusedAcrossARestart(t *testing.T) { {"other settings", 20}, } { t.Run(tt.name, func(t *testing.T) { - dir := t.TempDir() + t.Parallel() + dir := storeDir(t) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() topic := Topic{Tenant: "acme", Table: "t"} diff --git a/internal/mq/main_test.go b/internal/mq/main_test.go index f759977d..063e8ab8 100644 --- a/internal/mq/main_test.go +++ b/internal/mq/main_test.go @@ -7,8 +7,9 @@ import ( ) // TestMain silences the default logger, which the embedded server logs -// through. +// through, and turns off the embedded server's fsync per write. func TestMain(m *testing.M) { logtest.Silence() + EmbeddedSyncAlways = false m.Run() } From e441f57aba778456f2272228db333317713f868d Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Fri, 25 Sep 2026 08:28:08 -0400 Subject: [PATCH 3/6] test(app): boot without the embedded broker's fsync per write Every boot here opens its tenants' queues on the embedded broker, each open a handful of fsynced JetStream writes. TestMain turns mq.EmbeddedSyncAlways off, as internal/mq's own tests do: nothing here asserts anything across a crash. About 1.6s of the package's run. Part of #617. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- internal/app/main_test.go | 14 ++++++++++++++ 1 file changed, 14 insertions(+) create mode 100644 internal/app/main_test.go diff --git a/internal/app/main_test.go b/internal/app/main_test.go new file mode 100644 index 00000000..2f21403b --- /dev/null +++ b/internal/app/main_test.go @@ -0,0 +1,14 @@ +package app + +import ( + "testing" + + "github.com/Wave-RF/WaveHouse/internal/mq" +) + +// TestMain turns off the embedded broker's fsync per write, which every boot +// here pays for opening its queues. +func TestMain(m *testing.M) { + mq.EmbeddedSyncAlways = false + m.Run() +} From d46e273a97a72bfde150b499a1c5173a38d41fa2 Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Fri, 25 Sep 2026 08:34:10 -0400 Subject: [PATCH 4/6] fix(mq): create the embedded store at 0700; changelog the fail-fast Review fixes for the fail-fast store commit: nats-server creates its store directory at 0700, so NewEmbedded's MkdirAll now does too rather than widening it to group-readable, and the operator-visible change gets its CHANGELOG entry. Part of #617. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- CHANGELOG.md | 1 + internal/mq/embedded.go | 2 +- 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 70ce51e1..c26ef634 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -78,6 +78,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Fixed +- **An embedded queue store that cannot be created fails boot at once, naming the cause** (`internal/mq/embedded.go` (+ tests)): part of [#617](https://github.com/Wave-RF/WaveHouse/issues/617). A regular file or unwritable path at `/nats` failed JetStream in the background, so boot waited out the server's 5s readiness check and reported only `nats server not ready`. `NewEmbedded` now creates the directory first (at `0700`, as the server does) and refuses boot with the mkdir error. - **An insert invalidates a table's cached results under every tenant the directory holds** (`internal/app/wire.go` (+ tests), `internal/settings/registry.go` (+ tests), `AGENTS.md`, `docs/src/content/docs/{deployment,architecture,ingest-pipeline}.md`): until [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6 gives each tenant its own ClickHouse, every tenant reads the same tables, but the ingest worker — which writes every event as tenant `0`'s until story 5 — bumped only tenant `0`'s cache namespaces after an insert, so another tenant's cached query could answer stale rows for up to its TTL (an hour at most). The cache the worker invalidates through now fans each bumped namespace out to every tenant the registry knows (the new `Registry.Known`), the named one and a rejected one included — a rejected tenant comes back into service with the entries it has, so leaving it out would let a folder repaired inside a TTL serve pre-insert rows; reads are untouched, so a tenant is still never served another's cached rows. The residual, a folder removed and restored inside a TTL, is closed since #610: a tenant back on a pool after an absence has its structured-query results orphaned at once (`Cache.InvalidateTenant`). Raised by CodeRabbit on #602. - **The `?token=` strip no longer repairs a query string that does not parse** (`internal/auth/auth.go`): `bearerToken` removed a query-string token by parsing the query, deleting `token`, and re-encoding what was left — and `url.ParseQuery` skips a pair it cannot read, so the re-encoding erased that pair. A handler that parses the query strictly in order to refuse a malformed one would then see a clean query: `GET /v1/ops/pipes?tenant=acme;x=1&token=…` would have answered `200` with the default tenant's pipes. The token is read exactly as before and a query that parses is rewritten exactly as before; a query that does not parse now loses its token pairs and nothing else, byte for byte. Pinned through `api.NewRouter` with the real authenticator, since a handler-level test never runs the middleware that rewrote the URL. diff --git a/internal/mq/embedded.go b/internal/mq/embedded.go index d4ac14b9..98971492 100644 --- a/internal/mq/embedded.go +++ b/internal/mq/embedded.go @@ -154,7 +154,7 @@ var errNoQueue = errors.New("no queue is open for it yet") func NewEmbedded(storeDir string) (*EmbeddedNATS, error) { // A store the server cannot create fails JetStream in the background, and // ReadyForConnections would only give up on it after its whole wait. - if err := os.MkdirAll(storeDir, 0o750); err != nil { + if err := os.MkdirAll(storeDir, 0o700); err != nil { return nil, fmt.Errorf("nats store: %w", err) } opts := &natsserver.Options{ From f89739447b599c34fca16aced847aa51308b72c2 Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Fri, 25 Sep 2026 08:37:46 -0400 Subject: [PATCH 5/6] docs(changelog): scope the fail-fast store entry to what mkdir catches MkdirAll succeeds on an existing directory whatever its mode, so an unwritable /nats still waits the 5s and reports "nats server not ready". Say so instead of claiming it. Part of #617. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c26ef634..d50d6cae 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -78,7 +78,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), ### Fixed -- **An embedded queue store that cannot be created fails boot at once, naming the cause** (`internal/mq/embedded.go` (+ tests)): part of [#617](https://github.com/Wave-RF/WaveHouse/issues/617). A regular file or unwritable path at `/nats` failed JetStream in the background, so boot waited out the server's 5s readiness check and reported only `nats server not ready`. `NewEmbedded` now creates the directory first (at `0700`, as the server does) and refuses boot with the mkdir error. +- **An embedded queue store that cannot be created fails boot at once, naming the cause** (`internal/mq/embedded.go` (+ tests)): part of [#617](https://github.com/Wave-RF/WaveHouse/issues/617). A regular file at `/nats`, or a `nats` directory that could not be created there, failed JetStream in the background, so boot waited out the server's 5s readiness check and reported only `nats server not ready`. `NewEmbedded` now creates the directory first (at `0700`, as the server does) and refuses boot with the mkdir error. An existing but unwritable `nats` directory still takes the old path. - **An insert invalidates a table's cached results under every tenant the directory holds** (`internal/app/wire.go` (+ tests), `internal/settings/registry.go` (+ tests), `AGENTS.md`, `docs/src/content/docs/{deployment,architecture,ingest-pipeline}.md`): until [#583](https://github.com/Wave-RF/WaveHouse/issues/583) story 6 gives each tenant its own ClickHouse, every tenant reads the same tables, but the ingest worker — which writes every event as tenant `0`'s until story 5 — bumped only tenant `0`'s cache namespaces after an insert, so another tenant's cached query could answer stale rows for up to its TTL (an hour at most). The cache the worker invalidates through now fans each bumped namespace out to every tenant the registry knows (the new `Registry.Known`), the named one and a rejected one included — a rejected tenant comes back into service with the entries it has, so leaving it out would let a folder repaired inside a TTL serve pre-insert rows; reads are untouched, so a tenant is still never served another's cached rows. The residual, a folder removed and restored inside a TTL, is closed since #610: a tenant back on a pool after an absence has its structured-query results orphaned at once (`Cache.InvalidateTenant`). Raised by CodeRabbit on #602. - **The `?token=` strip no longer repairs a query string that does not parse** (`internal/auth/auth.go`): `bearerToken` removed a query-string token by parsing the query, deleting `token`, and re-encoding what was left — and `url.ParseQuery` skips a pair it cannot read, so the re-encoding erased that pair. A handler that parses the query strictly in order to refuse a malformed one would then see a clean query: `GET /v1/ops/pipes?tenant=acme;x=1&token=…` would have answered `200` with the default tenant's pipes. The token is read exactly as before and a query that parses is rewritten exactly as before; a query that does not parse now loses its token pairs and nothing else, byte for byte. Pinned through `api.NewRouter` with the real authenticator, since a handler-level test never runs the middleware that rewrote the URL. From 668c2c4295c72bb91bcb979e472a1b49b72a14b4 Mon Sep 17 00:00:00 2001 From: Eric Andrechek Date: Fri, 25 Sep 2026 09:49:26 -0400 Subject: [PATCH 6/6] test(mq): drop #612's occupied directory; the test waits for $G instead #612 kept the streams directory occupied through acme's failed open so the server would not remove it under globex's open. The parallel-tests change fixes the same race by waiting for the server to remove $G, which never happens while the occupier is there: with both, the test times out every run. Keep the wait. Part of #617. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01EJr5tY4WQUy2sc4MbW67vL --- internal/mq/embedded_test.go | 9 --------- 1 file changed, 9 deletions(-) diff --git a/internal/mq/embedded_test.go b/internal/mq/embedded_test.go index cfe9b30c..55e7fec4 100644 --- a/internal/mq/embedded_test.go +++ b/internal/mq/embedded_test.go @@ -469,15 +469,6 @@ func TestEmbeddedNATS_SetMaxBytes_AQueueThatCannotOpen(t *testing.T) { require.NoError(t, os.WriteFile(block, nil, 0o600)) } obstruct() - // A directory JetStream ignores (no metafile, so recovery skips it) - // keeps the streams directory occupied through acme's failed open, - // which would otherwise leave it empty: the server then removes it on a - // goroutine of its own, and globex's open right after would race that - // inside its own MkdirAll (see - // TestEmbeddedNATS_PacesTheRetriesOfAQueueThatCannotOpen). Not another - // tenant's streams: those reserve bytes, and the refusal guarded - // against below needs the reserved count to have gone negative. - require.NoError(t, os.Mkdir(filepath.Join(filepath.Dir(block), "occupied"), 0o750)) e := openEmbedded(t, dir) ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel()