Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 at `<data_dir>/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.

Expand Down
14 changes: 14 additions & 0 deletions internal/app/main_test.go
Original file line number Diff line number Diff line change
@@ -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()
}
13 changes: 12 additions & 1 deletion internal/mq/embedded.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"log/slog"
"maps"
"math"
"os"
"slices"
"strings"
"sync"
Expand Down Expand Up @@ -131,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.
Expand All @@ -146,11 +152,16 @@ 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, 0o700); err != nil {
return nil, fmt.Errorf("nats store: %w", err)
}
opts := &natsserver.Options{
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
Expand Down
Loading
Loading