diff --git a/docs/content/docs/development/configuration.mdx b/docs/content/docs/development/configuration.mdx index 70439be5..7fbe3ae2 100644 --- a/docs/content/docs/development/configuration.mdx +++ b/docs/content/docs/development/configuration.mdx @@ -291,6 +291,9 @@ On `filesystem`, a remote worker writes blobs to its own disk rather than a volu |---|---|---|---| | `EVENTBUS_PROVIDER` | `nats` or `kafka`. Kafka needs images built with `GO_TAGS=kafka` | `nats` under compose, `kafka` for a bare binary | yes | | `NATS_URL` | JetStream address. Credentials in the URL are honored by every service, including the Rust tracking publisher: `nats://user:pass@host:4222` for a user, `nats://token@host:4222` for a token, `tls://` for TLS | `nats://nats:4222` | yes | +| `NATS_CREDS` | Path to a NATS credentials file (user JWT plus nkey seed), for a bus that authenticates with JWT rather than a token. Synadia Cloud and any nsc-managed account issue one | unset | yes | +| `NATS_CREDS_B64` | The same file, base64 encoded, as a single line. This is the form the fleet uses: a node receives environment variables rather than files, and the env file docker reads cannot express a multi-line value. Takes precedence over `NATS_CREDS` | unset | yes | +| `NATS_MAX_BYTES` | Ceiling on the stream's size on disk. Accepts a plain byte count or a size (`2GiB`, `512MB`, `1G`). Unset leaves the stream bounded only by `NATS_STREAM_MAX_AGE` and the account's own quota, which a managed bus may refuse: Synadia's "Max Bytes Required" rejects any stream created without one | unset (unlimited) | yes | | `NATS_STREAM_NAME`, `NATS_SUBJECT_PREFIX` | Stream and subject naming | `warmbly` | yes | | `KAFKA_BOOTSTRAP_SERVERS` | Broker list when `EVENTBUS_PROVIDER=kafka` | unset | yes | | `KAFKA_SASL_USERNAME`, `KAFKA_SASL_PASSWORD` | Broker credentials | unset | yes | diff --git a/go.mod b/go.mod index a8542128..76d5252d 100644 --- a/go.mod +++ b/go.mod @@ -32,8 +32,10 @@ require ( github.com/meszmate/apple-go v0.0.0-20250828163208-7fea48c91b32 github.com/microcosm-cc/bluemonday v1.0.27 github.com/mileusna/useragent v1.3.5 + github.com/nats-io/jwt/v2 v2.8.1 github.com/nats-io/nats-server/v2 v2.14.1 github.com/nats-io/nats.go v1.52.0 + github.com/nats-io/nkeys v0.4.15 github.com/openai/openai-go/v2 v2.7.1 github.com/oschwald/geoip2-golang/v2 v2.0.0 github.com/posthog/posthog-go v1.25.1 @@ -227,8 +229,6 @@ require ( github.com/moricho/tparallel v0.3.2 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/nakabonne/nestif v0.3.1 // indirect - github.com/nats-io/jwt/v2 v2.8.1 // indirect - github.com/nats-io/nkeys v0.4.15 // indirect github.com/nats-io/nuid v1.0.1 // indirect github.com/nishanths/exhaustive v0.12.0 // indirect github.com/nishanths/predeclared v0.2.2 // indirect diff --git a/internal/api/handler/fleet_nodes.go b/internal/api/handler/fleet_nodes.go index 753ab35c..bb79744b 100644 --- a/internal/api/handler/fleet_nodes.go +++ b/internal/api/handler/fleet_nodes.go @@ -217,6 +217,10 @@ var nodeEnvKeys = []string{ "APP_ENV", "EVENTBUS_PROVIDER", "NATS_URL", + // Base64, not a path: a node has no file to point at, and the env file + // docker reads cannot hold a multi-line value. + "NATS_CREDS_B64", + "NATS_MAX_BYTES", "KAFKA_BOOTSTRAP_SERVERS", "KAFKA_SASL_USERNAME", "KAFKA_SASL_PASSWORD", diff --git a/internal/infrastructure/eventbus/factory.go b/internal/infrastructure/eventbus/factory.go index ec4bc16d..bb4f6507 100644 --- a/internal/infrastructure/eventbus/factory.go +++ b/internal/infrastructure/eventbus/factory.go @@ -46,10 +46,15 @@ func FromEnv(bootstrap string, sasl *kafka.SASLConfig) (EventBus, error) { SASL: sasl, }) case "nats": + maxBytes, err := parseByteSize(os.Getenv("NATS_MAX_BYTES")) + if err != nil { + return nil, err + } return NewNATS(NATSConfig{ URL: natsURLFromEnv(), StreamName: os.Getenv("NATS_STREAM_NAME"), SubjectPrefix: os.Getenv("NATS_SUBJECT_PREFIX"), + MaxBytes: maxBytes, }) default: return nil, fmt.Errorf("eventbus: unknown EVENTBUS_PROVIDER %q (want: kafka, nats)", provider) diff --git a/internal/infrastructure/eventbus/nats.go b/internal/infrastructure/eventbus/nats.go index 87a8e1f3..e839e0ac 100644 --- a/internal/infrastructure/eventbus/nats.go +++ b/internal/infrastructure/eventbus/nats.go @@ -34,6 +34,11 @@ type NATSBus struct { stream string prefix string + // Held rather than passed: ensureStream retries after a failure, and the + // retry callers have no config to hand it. + maxAge time.Duration + maxBytes int64 + mu sync.Mutex subscribers []jetstream.ConsumeContext streamEnsure sync.Once @@ -59,6 +64,11 @@ type NATSConfig struct { // "use the stream's existing setting or 7 days for new streams". MaxAge time.Duration + // MaxBytes caps the stream on disk. Zero leaves it unbounded, which a + // managed account may refuse: Synadia's "Max Bytes Required" rejects any + // stream created without one. + MaxBytes int64 + // Options passed to nats.Connect (auth, TLS, etc). Options []nats.Option } @@ -86,6 +96,15 @@ func NewNATS(cfg NATSConfig) (*NATSBus, error) { nats.ReconnectWait(2 * time.Second), }, cfg.Options...) + // A JWT credential cannot travel in the URL, so it is resolved separately. + creds, err := credsOption() + if err != nil { + return nil, err + } + if creds != nil { + opts = append(opts, creds) + } + nc, err := nats.Connect(cfg.URL, opts...) if err != nil { return nil, fmt.Errorf("eventbus nats: connect: %w", err) @@ -101,11 +120,14 @@ func NewNATS(cfg NATSConfig) (*NATSBus, error) { js: js, stream: cfg.StreamName, prefix: cfg.SubjectPrefix, + + maxAge: cfg.MaxAge, + maxBytes: cfg.MaxBytes, } // Eagerly ensure the stream so misconfiguration surfaces at boot rather // than on the first publish. Failures are non-fatal here — the lazy // retry inside Publish/Subscribe will surface them to callers. - if err := b.ensureStream(context.Background(), cfg.MaxAge); err != nil { + if err := b.ensureStream(context.Background()); err != nil { log.Warn().Err(err).Msg("eventbus nats: deferred stream setup") } return b, nil @@ -143,7 +165,7 @@ func (b *NATSBus) durable(group string, topics []string) string { return sb.String() } -func (b *NATSBus) ensureStream(ctx context.Context, maxAge time.Duration) error { +func (b *NATSBus) ensureStream(ctx context.Context) error { b.streamEnsure.Do(func() { cctx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() @@ -153,11 +175,12 @@ func (b *NATSBus) ensureStream(ctx context.Context, maxAge time.Duration) error Subjects: []string{filter}, Retention: jetstream.LimitsPolicy, Storage: jetstream.FileStorage, - MaxAge: maxAge, + MaxAge: b.maxAge, + MaxBytes: maxBytesOrUnlimited(b.maxBytes), Discard: jetstream.DiscardOld, }) if err != nil { - b.streamErr = fmt.Errorf("eventbus nats: ensure stream %q: %w", b.stream, err) + b.streamErr = fmt.Errorf("eventbus nats: ensure stream %q: %w%s", b.stream, err, maxBytesHint(err, b.maxBytes)) b.streamEnsure = sync.Once{} // allow retry on next call } }) @@ -175,7 +198,7 @@ func (b *NATSBus) Publish(ctx context.Context, topic, key string, payload []byte } b.mu.Unlock() - if err := b.ensureStream(ctx, 0); err != nil { + if err := b.ensureStream(ctx); err != nil { return err } @@ -209,7 +232,7 @@ func (b *NATSBus) Subscribe(ctx context.Context, topics []string, group string, if handler == nil { return errors.New("eventbus nats: handler required") } - if err := b.ensureStream(ctx, 0); err != nil { + if err := b.ensureStream(ctx); err != nil { return err } @@ -325,3 +348,24 @@ func natsURLFromEnv() string { // Compile-time interface check. var _ EventBus = (*NATSBus)(nil) + +// maxBytesOrUnlimited maps an unset ceiling onto JetStream's own spelling for +// it. Zero in a StreamConfig means zero bytes, not unlimited. +func maxBytesOrUnlimited(v int64) int64 { + if v <= 0 { + return -1 + } + return v +} + +// maxBytesHint turns the server's refusal into the thing to do about it. The +// bare error names a policy the operator has probably never heard of. +func maxBytesHint(err error, maxBytes int64) string { + if maxBytes > 0 || err == nil { + return "" + } + if !strings.Contains(strings.ToLower(err.Error()), "max bytes") { + return "" + } + return " (this account requires every stream to declare a size; set NATS_MAX_BYTES, e.g. NATS_MAX_BYTES=1GiB)" +} diff --git a/internal/infrastructure/eventbus/nats_creds.go b/internal/infrastructure/eventbus/nats_creds.go new file mode 100644 index 00000000..a34d1a7d --- /dev/null +++ b/internal/infrastructure/eventbus/nats_creds.go @@ -0,0 +1,60 @@ +package eventbus + +import ( + "encoding/base64" + "fmt" + "os" + "strings" + + "github.com/nats-io/jwt/v2" + "github.com/nats-io/nats.go" +) + +// credsOption resolves the NATS user credential from the environment. +// +// A managed bus (Synadia Cloud and anything else using JWT auth) authenticates +// with a .creds file holding a user JWT and an nkey seed, which no URL can +// carry. Two forms, because the two places this runs cannot use the same one: +// +// NATS_CREDS path to the file. Containers and local development. +// NATS_CREDS_B64 the file, base64 encoded. A worker joins with an env file +// that docker passes via --env-file, which cannot express a +// multi-line value, so the fleet needs a single-line form. +// +// Base64 rather than separate JWT and seed variables: it is one value, it is +// exactly what Synadia hands you, and there is no chance of pasting the two +// halves the wrong way round. It is the shape NATS's own NEX project uses. +// +// Returns nil when neither is set, so token and user-password URLs keep working. +func credsOption() (nats.Option, error) { + if raw := strings.TrimSpace(os.Getenv("NATS_CREDS_B64")); raw != "" { + decoded, err := base64.StdEncoding.DecodeString(raw) + if err != nil { + return nil, fmt.Errorf("eventbus nats: NATS_CREDS_B64 is not valid base64: %w", err) + } + return credsFromContents(decoded) + } + if path := strings.TrimSpace(os.Getenv("NATS_CREDS")); path != "" { + return nats.UserCredentials(path), nil + } + return nil, nil +} + +// credsFromContents builds the option without writing the secret to disk. +// nats.UserCredentials only takes a path, and a temp file would leave the +// seed readable to anything else on the host for the life of the process. +func credsFromContents(contents []byte) (nats.Option, error) { + userJWT, err := jwt.ParseDecoratedJWT(contents) + if err != nil { + return nil, fmt.Errorf("eventbus nats: no user JWT in the credentials: %w", err) + } + kp, err := jwt.ParseDecoratedUserNKey(contents) + if err != nil { + return nil, fmt.Errorf("eventbus nats: no nkey seed in the credentials: %w", err) + } + seed, err := kp.Seed() + if err != nil { + return nil, fmt.Errorf("eventbus nats: credentials carry no private seed: %w", err) + } + return nats.UserJWTAndSeed(userJWT, string(seed)), nil +} diff --git a/internal/infrastructure/eventbus/nats_creds_test.go b/internal/infrastructure/eventbus/nats_creds_test.go new file mode 100644 index 00000000..068d9ea9 --- /dev/null +++ b/internal/infrastructure/eventbus/nats_creds_test.go @@ -0,0 +1,82 @@ +package eventbus + +import ( + "encoding/base64" + "os" + "testing" + + "github.com/nats-io/nkeys" +) + +// Built at run time from a real throwaway keypair: an nkey seed carries a +// checksum, so a hand-written one cannot parse, and embedding a valid seed in +// the repo would put key-shaped material in git for no reason. +func sampleCredsFile(t *testing.T) string { + t.Helper() + kp, err := nkeys.CreateUser() + if err != nil { + t.Fatal(err) + } + seed, err := kp.Seed() + if err != nil { + t.Fatal(err) + } + return "-----BEGIN NATS USER JWT-----\n" + + "eyJ0eXAiOiJKV1QiLCJhbGciOiJlZDI1NTE5LW5rZXkifQ.eyJzdWIiOiJVQUEifQ.c2ln\n" + + "------END NATS USER JWT------\n\n" + + "-----BEGIN USER NKEY SEED-----\n" + + string(seed) + "\n" + + "------END USER NKEY SEED------\n" +} + +// The fleet ships env files through docker --env-file, which cannot express a +// multi-line value, so the base64 form is the one that has to work. +func TestCredsOptionFromBase64(t *testing.T) { + t.Setenv("NATS_CREDS_B64", base64.StdEncoding.EncodeToString([]byte(sampleCredsFile(t)))) + t.Setenv("NATS_CREDS", "") + opt, err := credsOption() + if err != nil { + t.Fatalf("credsOption: %v", err) + } + if opt == nil { + t.Fatal("NATS_CREDS_B64 was set but produced no option") + } +} + +func TestCredsOptionFromFile(t *testing.T) { + f := t.TempDir() + "/u.creds" + if err := os.WriteFile(f, []byte(sampleCredsFile(t)), 0o600); err != nil { + t.Fatal(err) + } + t.Setenv("NATS_CREDS_B64", "") + t.Setenv("NATS_CREDS", f) + opt, err := credsOption() + if err != nil || opt == nil { + t.Fatalf("file form: opt=%v err=%v", opt, err) + } +} + +// Neither set must stay silent: token and user-password URLs still work, and +// erroring here would break every existing deployment. +func TestCredsOptionAbsent(t *testing.T) { + t.Setenv("NATS_CREDS_B64", "") + t.Setenv("NATS_CREDS", "") + opt, err := credsOption() + if err != nil || opt != nil { + t.Fatalf("expected no option and no error, got opt=%v err=%v", opt, err) + } +} + +// A truncated or mis-pasted value must say so at boot rather than fail as an +// unexplained authorization error against the bus. +func TestCredsOptionRejectsGarbage(t *testing.T) { + t.Setenv("NATS_CREDS", "") + t.Setenv("NATS_CREDS_B64", "not-base64!!") + if _, err := credsOption(); err == nil { + t.Error("invalid base64 was accepted") + } + t.Setenv("NATS_CREDS_B64", base64.StdEncoding.EncodeToString([]byte("no jwt here"))) + if _, err := credsOption(); err == nil { + t.Error("credentials with no JWT were accepted") + } +} diff --git a/internal/infrastructure/eventbus/nats_maxbytes.go b/internal/infrastructure/eventbus/nats_maxbytes.go new file mode 100644 index 00000000..a111340d --- /dev/null +++ b/internal/infrastructure/eventbus/nats_maxbytes.go @@ -0,0 +1,57 @@ +package eventbus + +import ( + "fmt" + "strconv" + "strings" +) + +// parseByteSize reads NATS_MAX_BYTES. Plain bytes, or a size with a suffix: +// 2GiB, 512MB, 1G. Binary and decimal are both accepted because operators +// reach for either, and a stream ceiling does not need the distinction to be +// load bearing. +// +// Empty means unset, which leaves the stream bounded only by MaxAge and the +// account's own quota. +func parseByteSize(raw string) (int64, error) { + s := strings.TrimSpace(raw) + if s == "" { + return 0, nil + } + upper := strings.ToUpper(s) + + // Longest suffix first: GIB has to match before GB, and GB before G. + units := []struct { + suffix string + mult int64 + }{ + {"TIB", 1 << 40}, {"GIB", 1 << 30}, {"MIB", 1 << 20}, {"KIB", 1 << 10}, + {"TB", 1e12}, {"GB", 1e9}, {"MB", 1e6}, {"KB", 1e3}, + {"T", 1 << 40}, {"G", 1 << 30}, {"M", 1 << 20}, {"K", 1 << 10}, + {"B", 1}, + } + for _, u := range units { + if !strings.HasSuffix(upper, u.suffix) { + continue + } + num := strings.TrimSpace(upper[:len(upper)-len(u.suffix)]) + // Fractional sizes are the natural way to write 2.5GiB. + f, err := strconv.ParseFloat(num, 64) + if err != nil { + return 0, fmt.Errorf("eventbus nats: NATS_MAX_BYTES %q is not a size", raw) + } + if f < 0 { + return 0, fmt.Errorf("eventbus nats: NATS_MAX_BYTES %q is negative", raw) + } + return int64(f * float64(u.mult)), nil + } + + n, err := strconv.ParseInt(upper, 10, 64) + if err != nil { + return 0, fmt.Errorf("eventbus nats: NATS_MAX_BYTES %q is not a size", raw) + } + if n < 0 { + return 0, fmt.Errorf("eventbus nats: NATS_MAX_BYTES %q is negative", raw) + } + return n, nil +} diff --git a/internal/infrastructure/eventbus/nats_maxbytes_test.go b/internal/infrastructure/eventbus/nats_maxbytes_test.go new file mode 100644 index 00000000..4e96ff15 --- /dev/null +++ b/internal/infrastructure/eventbus/nats_maxbytes_test.go @@ -0,0 +1,53 @@ +package eventbus + +import "testing" + +// The ceiling is written by hand in a dashboard or an env file, so the shapes +// an operator actually types have to parse: a bare count, binary and decimal +// suffixes, and the fractional form a console shows (2.5 GiB). +func TestParseByteSize(t *testing.T) { + cases := map[string]int64{ + "": 0, + "1024": 1024, + "1KiB": 1024, + "1 KiB": 1024, + "2GiB": 2 << 30, + "2.5GiB": int64(2.5 * float64(1<<30)), + "512MB": 512_000_000, + "1G": 1 << 30, + "1B": 1, + " 4MiB ": 4 << 20, + "1gib": 1 << 30, + } + for in, want := range cases { + got, err := parseByteSize(in) + if err != nil { + t.Errorf("%q: unexpected error %v", in, err) + continue + } + if got != want { + t.Errorf("%q = %d, want %d", in, got, want) + } + } +} + +// A typo must stop the process at boot, not silently leave the stream +// unbounded on an account that will then refuse to create it. +func TestParseByteSizeRejectsNonsense(t *testing.T) { + for _, in := range []string{"big", "1XB", "-5", "-1GiB", "1.2.3GiB"} { + if _, err := parseByteSize(in); err == nil { + t.Errorf("%q was accepted", in) + } + } +} + +// JetStream spells "no ceiling" as -1; a literal 0 in a StreamConfig means +// zero bytes, which would reject every message. +func TestMaxBytesOrUnlimited(t *testing.T) { + if got := maxBytesOrUnlimited(0); got != -1 { + t.Errorf("unset = %d, want -1", got) + } + if got := maxBytesOrUnlimited(1 << 30); got != 1<<30 { + t.Errorf("set = %d, want %d", got, 1<<30) + } +} diff --git a/tracking/src/nats.rs b/tracking/src/nats.rs index 5b42e97d..df79148e 100644 --- a/tracking/src/nats.rs +++ b/tracking/src/nats.rs @@ -1,4 +1,6 @@ use async_nats::jetstream; +use base64::engine::general_purpose::STANDARD as BASE64; +use base64::Engine as _; use crate::config::Config; use crate::events::TrackingEvent; @@ -27,6 +29,22 @@ impl NatsProducer { None => opts.token(user.to_string()), }; } + // A managed bus authenticates with a user JWT and nkey seed, which no + // URL can carry. NATS_CREDS_B64 is the single-line form the fleet needs + // (docker --env-file cannot express a multi-line value); NATS_CREDS is + // a path, for containers and local development. + if let Ok(b64) = std::env::var("NATS_CREDS_B64") { + if !b64.trim().is_empty() { + let raw = BASE64.decode(b64.trim())?; + opts = opts.credentials(std::str::from_utf8(&raw)?)?; + } + } else if let Ok(path) = std::env::var("NATS_CREDS") { + if !path.trim().is_empty() { + let contents = std::fs::read_to_string(path.trim())?; + opts = opts.credentials(&contents)?; + } + } + let client = opts.connect(addr).await?; let js = jetstream::new(client); let subject = format!("{}.{}", config.nats_subject_prefix, config.kafka_topic);