package main import ( "context" "errors" "io" "log/slog" "os" "strings" "testing" "time" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/authentication" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/envelope" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/identity" "lingniu-vehicle-ingest/go/vehicle-gateway/internal/metrics" ) func TestNATSSinkConfigFromEnvUsesExplicitSubjects(t *testing.T) { t.Setenv("NATS_URL", "nats://172.17.111.56:4222") t.Setenv("NATS_SUBJECT_GB32960_RAW", "custom.raw.gb32960") t.Setenv("NATS_SUBJECT_JT808_RAW", "custom.raw.jt808") t.Setenv("NATS_SUBJECT_YUTONG_MQTT_RAW", "custom.raw.yutong") t.Setenv("NATS_SUBJECT_GB32960_FIELDS", "custom.fields.gb32960") t.Setenv("NATS_SUBJECT_JT808_FIELDS", "custom.fields.jt808") t.Setenv("NATS_SUBJECT_YUTONG_MQTT_FIELDS", "custom.fields.yutong") t.Setenv("NATS_SUBJECT_UNIFIED", "custom.unified") t.Setenv("NATS_OUTBOX_MAX_INFLIGHT", "12000") t.Setenv("NATS_OUTBOX_ACK_TIMEOUT_MS", "4500") cfg := natsSinkConfigFromEnv() if got, want := cfg.URL, "nats://172.17.111.56:4222"; got != want { t.Fatalf("URL = %q, want %q", got, want) } if got, want := cfg.RawSubjects[envelope.ProtocolGB32960], "custom.raw.gb32960"; got != want { t.Fatalf("gb32960 subject = %q, want %q", got, want) } if got, want := cfg.RawSubjects[envelope.ProtocolJT808], "custom.raw.jt808"; got != want { t.Fatalf("jt808 subject = %q, want %q", got, want) } if got, want := cfg.RawSubjects[envelope.ProtocolYutongMQTT], "custom.raw.yutong"; got != want { t.Fatalf("yutong subject = %q, want %q", got, want) } if got, want := cfg.FieldsSubjects[envelope.ProtocolGB32960], "custom.fields.gb32960"; got != want { t.Fatalf("gb32960 fields subject = %q, want %q", got, want) } if got, want := cfg.FieldsSubjects[envelope.ProtocolJT808], "custom.fields.jt808"; got != want { t.Fatalf("jt808 fields subject = %q, want %q", got, want) } if got, want := cfg.FieldsSubjects[envelope.ProtocolYutongMQTT], "custom.fields.yutong"; got != want { t.Fatalf("yutong fields subject = %q, want %q", got, want) } if got, want := cfg.UnifiedSubject, "custom.unified"; got != want { t.Fatalf("unified subject = %q, want %q", got, want) } if got, want := cfg.AsyncMaxPending, 12000; got != want { t.Fatalf("async max pending = %d, want %d", got, want) } if got, want := cfg.AsyncAckTimeout, 4500*time.Millisecond; got != want { t.Fatalf("async ack timeout = %s, want %s", got, want) } } func TestNATSOutboxConfigFromEnvIsExplicitAndDefaultsToSafeFsync(t *testing.T) { t.Setenv("NATS_DURABLE_OUTBOX_ENABLED", "true") t.Setenv("NATS_OUTBOX_DIR", "/var/lib/lingniu-go/nats-outbox") t.Setenv("NATS_OUTBOX_REPLAY_BATCH_SIZE", "750") t.Setenv("NATS_OUTBOX_REPLAY_INTERVAL_MS", "250") t.Setenv("NATS_OUTBOX_CLOSE_TIMEOUT_MS", "7000") t.Setenv("NATS_OUTBOX_WAL_SEGMENT_BYTES", "8388608") t.Setenv("NATS_OUTBOX_WAL_SEGMENT_AGE_MS", "4000") t.Setenv("NATS_OUTBOX_WAL_APPEND_QUEUE_SIZE", "50000") t.Setenv("NATS_OUTBOX_WAL_COMMIT_BATCH_SIZE", "128") t.Setenv("NATS_OUTBOX_WAL_COMMIT_INTERVAL_MS", "2") t.Setenv("NATS_OUTBOX_FSYNC", "") registry := metrics.NewRegistry() onError := func(error) {} runtime := natsOutboxConfigFromEnv(registry, onError) if !runtime.Enabled { t.Fatal("outbox should be enabled") } if got, want := runtime.Config.Directory, "/var/lib/lingniu-go/nats-outbox"; got != want { t.Fatalf("directory = %q, want %q", got, want) } if got, want := runtime.Config.ReplayBatchSize, 750; got != want { t.Fatalf("replay batch = %d, want %d", got, want) } if !runtime.Config.SyncWrites { t.Fatal("outbox fsync should default to true") } if got, want := runtime.Config.CloseTimeout, 7*time.Second; got != want { t.Fatalf("close timeout = %s, want %s", got, want) } if got, want := runtime.Config.WALSegmentBytes, int64(8<<20); got != want { t.Fatalf("WAL segment bytes = %d, want %d", got, want) } if got, want := runtime.Config.WALSegmentAge, 4*time.Second; got != want { t.Fatalf("WAL segment age = %s, want %s", got, want) } if got, want := runtime.Config.WALAppendQueue, 50_000; got != want { t.Fatalf("WAL append queue = %d, want %d", got, want) } if got, want := runtime.Config.WALCommitBatch, 128; got != want { t.Fatalf("WAL commit batch = %d, want %d", got, want) } if got, want := runtime.Config.WALCommitWait, 2*time.Millisecond; got != want { t.Fatalf("WAL commit wait = %s, want %s", got, want) } if got, want := runtime.ReplayInterval, 250*time.Millisecond; got != want { t.Fatalf("replay interval = %s, want %s", got, want) } if runtime.Config.Metrics != registry || runtime.Config.OnError == nil { t.Fatal("metrics and error callback should be wired") } } func TestNATSOutboxConfigDoesNotReuseLegacySpoolDirectory(t *testing.T) { t.Setenv("NATS_DURABLE_OUTBOX_ENABLED", "") t.Setenv("NATS_OUTBOX_DIR", "") t.Setenv("NATS_SPOOL_DIR", "/var/lib/lingniu-go/nats-spool") t.Setenv("NATS_OUTBOX_FSYNC", "false") runtime := natsOutboxConfigFromEnv(nil, nil) if runtime.Enabled { t.Fatal("outbox must remain opt-in") } if runtime.Config.Directory != "" { t.Fatalf("WAL directory must be explicit, got %q", runtime.Config.Directory) } if runtime.Config.SyncWrites { t.Fatal("explicit false should disable fsync") } } func TestBuildGB32960AuthenticatorUsesConfiguredCredentials(t *testing.T) { t.Setenv("GB32960_AUTH_MODE", "enforce") t.Setenv("GB32960_PLATFORM_CREDENTIALS_FILE", "") t.Setenv("GB32960_PLATFORM_CREDENTIALS_JSON", `{"platform-a":"secret-a","platform-b":["secret-b-old","secret-b"]}`) authenticator, mode, count, err := buildGB32960Authenticator() if err != nil { t.Fatalf("buildGB32960Authenticator() error = %v", err) } if mode != authentication.ModeEnforce || count != 2 { t.Fatalf("mode=%q count=%d", mode, count) } result := authenticator.Authenticate(envelope.FrameEnvelope{ Protocol: envelope.ProtocolGB32960, MessageID: "0x05", Parsed: map[string]any{ "platform_login": map[string]any{"username": "platform-b", "password": "secret-b"}, }, }) if !result.Allowed || result.Status != authentication.StatusAccepted { t.Fatalf("authentication result = %#v", result) } } func TestBuildGB32960AuthenticatorRejectsEnforceWithoutCredentials(t *testing.T) { t.Setenv("GB32960_AUTH_MODE", "enforce") t.Setenv("GB32960_PLATFORM_CREDENTIALS_FILE", "") t.Setenv("GB32960_PLATFORM_CREDENTIALS_JSON", "") if _, _, _, err := buildGB32960Authenticator(); err == nil { t.Fatal("expected enforce mode configuration error") } } func TestGatewayConfiguresIdentityLookupCacheTTL(t *testing.T) { source, err := os.ReadFile("main.go") if err != nil { t.Fatalf("read main.go: %v", err) } for _, want := range []string{"LookupCacheTTL", "StaleLookupTTL", "IDENTITY_LOOKUP_CACHE_TTL_SECONDS", "IDENTITY_STALE_LOOKUP_TTL_SECONDS", "IDENTITY_LOOKUP_CACHE_MAX_ENTRIES", "IDENTITY_LOOKUP_CACHE_CLEANUP_INTERVAL_SECONDS", "vehicle_gateway_identity_cache_entries", "vehicle_gateway_jt808_registration_write_total", "OnRegistrationWriteResult", "IDENTITY_SOURCE_CODE_LOOKUP_ENABLED", "IDENTITY_RESOLVE_TIMEOUT_MS", "IDENTITY_SNAPSHOT_ONLY_ENABLED", "IDENTITY_SNAPSHOT_REFRESH_INTERVAL_SECONDS", "startIdentitySnapshotRefresh", "JT808_REGISTRATION_LOCATION_TOUCH_RETRY_INTERVAL_SECONDS", "JT808_REGISTRATION_WRITE_RETRY_ATTEMPTS", "JT808_REGISTRATION_WRITE_RETRY_DELAY_MS", "JT808_REGISTRATION_ASYNC_WRITE_ENABLED", "JT808_REGISTRATION_WRITE_QUEUE_SIZE", "JT808_REGISTRATION_WRITE_WORKERS", "JT808_REGISTRATION_WRITE_ENQUEUE_TIMEOUT_MS", "TimeoutResolver"} { if !strings.Contains(string(source), want) { t.Fatalf("gateway should expose identity lookup cache ttl, missing %s", want) } } } func TestRecordIdentityCacheStatsIncludesLocationTouchFailures(t *testing.T) { registry := metrics.NewRegistry() recordIdentityCacheStats(registry, identity.CacheStats{ LookupEntries: 1, RegistrationEntries: 2, SourceCodeEntries: 3, LocationTouchEntries: 4, LocationTouchFailureEntries: 5, SnapshotBindingEntries: 6, SnapshotIdentifierEntries: 7, SnapshotRegistrationEntries: 8, SnapshotSourceEntries: 9, SnapshotReady: true, SnapshotRefreshedAt: time.Unix(1234, 0), MaxEntries: 9, }) rendered := registry.Render() for _, want := range []string{ `vehicle_gateway_identity_cache_entries{cache="location_touch_failure"} 5`, `vehicle_gateway_identity_cache_entries{cache="registration_write_queue"} 0`, `vehicle_gateway_identity_cache_entries{cache="snapshot_binding"} 6`, `vehicle_gateway_identity_cache_entries{cache="snapshot_identifier"} 7`, `vehicle_gateway_identity_cache_entries{cache="snapshot_registration"} 8`, `vehicle_gateway_identity_cache_entries{cache="snapshot_source"} 9`, `vehicle_gateway_identity_cache_max_entries{cache="location_touch_failure"} 9`, `vehicle_gateway_identity_snapshot_ready 1`, `vehicle_gateway_identity_snapshot_last_success_unix_seconds 1234`, } { if !strings.Contains(rendered, want) { t.Fatalf("identity cache metrics missing %s in:\n%s", want, rendered) } } } func TestRecordIdentitySnapshotRefresh(t *testing.T) { registry := metrics.NewRegistry() recordIdentitySnapshotRefresh(registry, identity.SnapshotRefreshResult{ BindingEntries: 10, IdentifierEntries: 20, RegistrationEntries: 30, SourceEntries: 40, RefreshedAt: time.Unix(5678, 0), }, "ok") rendered := registry.Render() for _, want := range []string{ `vehicle_gateway_identity_snapshot_refresh_total{status="ok"} 1`, `vehicle_gateway_identity_snapshot_entries{kind="binding"} 10`, `vehicle_gateway_identity_snapshot_entries{kind="identifier"} 20`, `vehicle_gateway_identity_snapshot_entries{kind="registration"} 30`, `vehicle_gateway_identity_snapshot_entries{kind="source"} 40`, `vehicle_gateway_identity_snapshot_last_success_unix_seconds 5678`, } { if !strings.Contains(rendered, want) { t.Fatalf("identity snapshot metrics missing %s in:\n%s", want, rendered) } } } func TestIdentityDatabaseMaintenanceDoesNotBlockIngressStartup(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() db := &blockingIdentityDatabase{ started: make(chan struct{}), release: make(chan struct{}), } logger := slog.New(slog.NewTextHandler(io.Discard, nil)) startedAt := time.Now() startIdentityDatabaseMaintenance(ctx, logger, db, nil, false, time.Second, time.Second) if elapsed := time.Since(startedAt); elapsed > 100*time.Millisecond { t.Fatalf("database maintenance blocked gateway startup for %s", elapsed) } select { case <-db.started: case <-time.After(time.Second): t.Fatal("database maintenance did not start in background") } close(db.release) } func TestIdentitySnapshotRefreshDoesNotBlockIngressStartup(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() refresher := &blockingIdentitySnapshotRefresher{started: make(chan struct{})} logger := slog.New(slog.NewTextHandler(io.Discard, nil)) startedAt := time.Now() startIdentitySnapshotRefresh(ctx, logger, metrics.NewRegistry(), refresher, time.Hour, time.Second) if elapsed := time.Since(startedAt); elapsed > 100*time.Millisecond { t.Fatalf("snapshot refresh blocked gateway startup for %s", elapsed) } select { case <-refresher.started: case <-time.After(time.Second): t.Fatal("snapshot refresh did not start in background") } } func TestRecordJT808RegistrationWriteResult(t *testing.T) { registry := metrics.NewRegistry() recordJT808RegistrationWriteResult(registry, identity.RegistrationWriteResult{ Mode: "async_background", Status: "error", }) rendered := registry.Render() for _, want := range []string{ `vehicle_gateway_jt808_registration_write_total{mode="async_background",status="error"} 1`, `vehicle_gateway_last_jt808_registration_write_unix_seconds{mode="async_background",status="error"}`, } { if !strings.Contains(rendered, want) { t.Fatalf("registration write metric missing %s in:\n%s", want, rendered) } } } type blockingIdentityDatabase struct { started chan struct{} release chan struct{} } func (d *blockingIdentityDatabase) PingContext(ctx context.Context) error { close(d.started) select { case <-ctx.Done(): return ctx.Err() case <-d.release: return nil } } type blockingIdentitySnapshotRefresher struct { started chan struct{} } func (r *blockingIdentitySnapshotRefresher) RefreshSnapshot(ctx context.Context) (identity.SnapshotRefreshResult, error) { close(r.started) <-ctx.Done() return identity.SnapshotRefreshResult{}, errors.New("test refresh stopped") } func TestGatewayDefaultsTo100KConnectionCeiling(t *testing.T) { source, err := os.ReadFile("main.go") if err != nil { t.Fatalf("read main.go: %v", err) } if !strings.Contains(string(source), `envInt("TCP_MAX_CONNECTIONS", 120_000)`) { t.Fatal("gateway should default TCP_MAX_CONNECTIONS to a 100K-ready ceiling") } } func TestGatewayPassesMetricsRegistryToAsyncSink(t *testing.T) { source, err := os.ReadFile("main.go") if err != nil { t.Fatalf("read main.go: %v", err) } for _, want := range []string{ "buildSink(ctx, logger, registry)", "Metrics: registry", "EnqueueTimeout: enqueueTimeout", "NewPartitionedAsyncSink", "NATS_PARTITIONED_ASYNC_ENABLED", "KAFKA_PARTITIONED_ASYNC_ENABLED", "_ASYNC_RAW_QUEUE_SIZE", "_ASYNC_DERIVED_QUEUE_SIZE", "_ASYNC_RAW_ENQUEUE_TIMEOUT_MS", "_ASYNC_DERIVED_ENQUEUE_TIMEOUT_MS", "RawEnqueueTimeout", "DerivedEnqueueTimeout", "NATS_ASYNC_ENQUEUE_TIMEOUT_MS", "KAFKA_ASYNC_ENQUEUE_TIMEOUT_MS", `Name: "nats"`, `Name: "kafka"`, } { if !strings.Contains(string(source), want) { t.Fatalf("gateway async sink metrics wiring missing %s", want) } } } func TestGatewayExposesNATSDurableSpoolConfig(t *testing.T) { source, err := os.ReadFile("main.go") if err != nil { t.Fatalf("read main.go: %v", err) } for _, want := range []string{ "NATS_DURABLE_OUTBOX_ENABLED", "NATS_OUTBOX_DIR", "NATS_OUTBOX_MAX_INFLIGHT", "NATS_OUTBOX_ACK_TIMEOUT_MS", "NATS_OUTBOX_FSYNC", "NATS_OUTBOX_CLOSE_TIMEOUT_MS", "NATS_OUTBOX_REPLAY_BATCH_SIZE", "NATS_OUTBOX_REPLAY_INTERVAL_MS", "NATS_OUTBOX_WAL_SEGMENT_BYTES", "NATS_OUTBOX_WAL_SEGMENT_AGE_MS", "NATS_OUTBOX_WAL_APPEND_QUEUE_SIZE", "NATS_OUTBOX_WAL_COMMIT_BATCH_SIZE", "NATS_OUTBOX_WAL_COMMIT_INTERVAL_MS", "eventbus.NewDurableOutboxSink", "nats durable outbox enabled", "NATS_SPOOL_DIR", "NATS_SPOOL_REPLAY_BATCH_SIZE", "NATS_SPOOL_REPLAY_INTERVAL_MS", "nats durable spool enabled", "nats spool replay failed", "eventbus.NewDurableSink", "Metrics: registry", `Name: "nats"`, `Name: "kafka"`, } { if !strings.Contains(string(source), want) { t.Fatalf("gateway nats spool wiring missing %s", want) } } } func TestNATSSinkConfigFromEnvDefaultsToGoSubjects(t *testing.T) { t.Setenv("NATS_URL", "nats://172.17.111.56:4222") cfg := natsSinkConfigFromEnv() if got, want := cfg.RawSubjects[envelope.ProtocolGB32960], "vehicle.raw.go.gb32960.v1"; got != want { t.Fatalf("gb32960 subject = %q, want %q", got, want) } if got, want := cfg.RawSubjects[envelope.ProtocolJT808], "vehicle.raw.go.jt808.v1"; got != want { t.Fatalf("jt808 subject = %q, want %q", got, want) } if got, want := cfg.RawSubjects[envelope.ProtocolYutongMQTT], "vehicle.raw.go.yutong-mqtt.v1"; got != want { t.Fatalf("yutong subject = %q, want %q", got, want) } if got, want := cfg.FieldsSubjects[envelope.ProtocolGB32960], "vehicle.fields.go.gb32960.v1"; got != want { t.Fatalf("gb32960 fields subject = %q, want %q", got, want) } if got, want := cfg.FieldsSubjects[envelope.ProtocolJT808], "vehicle.fields.go.jt808.v1"; got != want { t.Fatalf("jt808 fields subject = %q, want %q", got, want) } if got, want := cfg.FieldsSubjects[envelope.ProtocolYutongMQTT], "vehicle.fields.go.yutong-mqtt.v1"; got != want { t.Fatalf("yutong fields subject = %q, want %q", got, want) } if got, want := cfg.UnifiedSubject, "vehicle.event.go.unified.v1"; got != want { t.Fatalf("unified subject = %q, want %q", got, want) } }