package store import ( "context" "database/sql" "errors" "fmt" "os" "strings" "sync" "testing" "time" "yovision/sense/internal/device" ) const ( postgresTestDSNEnv = "YOVISION_TEST_POSTGRES_DSN" postgresTestAdminDSNEnv = "YOVISION_TEST_POSTGRES_ADMIN_DSN" ) func TestPostgresInvalidDSNDoesNotLeakInput(t *testing.T) { secret := "do-not-echo-this-value" _, err := OpenPostgres(context.Background(), "postgres://sense:"+secret+"@%zz") if err == nil { t.Fatal("invalid PostgreSQL DSN was accepted") } if strings.Contains(err.Error(), secret) { t.Fatal("PostgreSQL configuration error leaked DSN input") } } func TestOpenRepositoryRejectsUnknownDriver(t *testing.T) { _, err := OpenRepository(context.Background(), "mysql", "unused") if err == nil { t.Fatal("unknown repository driver was accepted") } } func TestPostgresDefaultAndMaximumQuota(t *testing.T) { store, admin := openPostgresTestStore(t) ctx := context.Background() insertBellSite(t, admin, "tenant-default", "site-default", 0) for index := 1; index <= device.DefaultVideoChannels; index++ { if err := store.CreateDevice(ctx, videoDevice(index, "tenant-default", "site-default")); err != nil { t.Fatalf("create default channel %d: %v", index, err) } } radar := device.Device{ ID: "radar-default", TenantID: "tenant-default", SiteID: "site-default", SerialNumber: "radar-default", Name: "Radar", Modality: device.ModalityRadar, Capabilities: []device.Capability{device.CapabilityTelemetry}, DesiredState: device.DesiredEnabled, ActualState: device.ActualPending, } if err := store.CreateDevice(ctx, radar); err != nil { t.Fatalf("non-video device must not consume video quota: %v", err) } err := store.CreateDevice(ctx, videoDevice(17, "tenant-default", "site-default")) var quotaError *device.QuotaExceededError if !errors.As(err, "aError) || quotaError.Limit != 16 { t.Fatalf("expected default quota error, got %v", err) } insertBellSite(t, admin, "tenant-max", "site-max", 128) for index := 1; index <= device.MaximumVideoChannels; index++ { if err := store.CreateDevice(ctx, videoDevice(index+1000, "tenant-max", "site-max")); err != nil { t.Fatalf("create maximum channel %d: %v", index, err) } } err = store.CreateDevice(ctx, videoDevice(1129, "tenant-max", "site-max")) if !errors.As(err, "aError) || quotaError.Limit != 128 { t.Fatalf("expected maximum quota error, got %v", err) } } func TestPostgresConcurrentAdmissionCannotExceedQuota(t *testing.T) { store, admin := openPostgresTestStore(t) insertBellSite(t, admin, "tenant", "site", 1) ctx := context.Background() start := make(chan struct{}) errorsFound := make(chan error, 2) var wait sync.WaitGroup for index := 1; index <= 2; index++ { wait.Add(1) go func(index int) { defer wait.Done() <-start errorsFound <- store.CreateDevice(ctx, videoDevice(index, "tenant", "site")) }(index) } close(start) wait.Wait() close(errorsFound) var succeeded, rejected int for err := range errorsFound { if err == nil { succeeded++ continue } var quotaError *device.QuotaExceededError if errors.As(err, "aError) { rejected++ continue } t.Fatalf("unexpected concurrent admission error: %v", err) } if succeeded != 1 || rejected != 1 { t.Fatalf("expected one success and one quota rejection, got success=%d rejected=%d", succeeded, rejected) } } func TestPostgresProjectionFailureAndRollbackFailClosed(t *testing.T) { store, admin := openPostgresTestStore(t) ctx := context.Background() err := store.CreateDevice(ctx, videoDevice(1, "missing-tenant", "missing-site")) if !errors.Is(err, ErrQuotaProjectionUnavailable) { t.Fatalf("missing projection must fail closed, got %v", err) } insertBellSite(t, admin, "tenant", "site", 2) if _, err := store.db.ExecContext(ctx, `INSERT INTO sense.site_quota_projection_state( tenant_id, site_id, source_version, synced_at ) VALUES ('tenant', 'site', 99, clock_timestamp())`); err != nil { t.Fatal(err) } err = store.CreateDevice(ctx, videoDevice(2, "tenant", "site")) if !errors.Is(err, ErrQuotaProjectionInvalid) { t.Fatalf("projection version rollback must fail closed, got %v", err) } var devices int if err := store.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM sense.devices`).Scan(&devices); err != nil { t.Fatal(err) } if devices != 0 { t.Fatalf("failed admissions must not persist devices, got %d", devices) } } func TestPostgresLowerQuotaDoesNotDisableAndEnableStillChecks(t *testing.T) { store, admin := openPostgresTestStore(t) ctx := context.Background() insertBellSite(t, admin, "tenant", "site", 2) for index := 1; index <= 2; index++ { if err := store.CreateDevice(ctx, videoDevice(index, "tenant", "site")); err != nil { t.Fatal(err) } } disabled := videoDevice(3, "tenant", "site") disabled.DesiredState = device.DesiredDisabled if err := store.CreateDevice(ctx, disabled); err != nil { t.Fatal(err) } if _, err := admin.ExecContext(ctx, `UPDATE bell.sites SET max_video_channels = 1 WHERE tenant_id = 'tenant' AND id = 'site'`); err != nil { t.Fatal(err) } for index := 1; index <= 2; index++ { value, err := store.GetDevice(ctx, fmt.Sprintf("camera-%03d", index)) if err != nil { t.Fatal(err) } if value.DesiredState != device.DesiredEnabled { t.Fatalf("quota decrease disabled existing device %d", index) } } err := store.SetDesiredState(ctx, "camera-003", device.DesiredEnabled) var quotaError *device.QuotaExceededError if !errors.As(err, "aError) || quotaError.Limit != 1 { t.Fatalf("new enable after quota decrease must be rejected, got %v", err) } } func TestPostgresTenantIsolationAndConvergenceParity(t *testing.T) { store, admin := openPostgresTestStore(t) ctx := context.Background() insertBellSite(t, admin, "tenant-a", "site", 1) insertBellSite(t, admin, "tenant-b", "site", 1) left := videoDevice(1, "tenant-a", "site") right := videoDevice(2, "tenant-b", "site") right.SerialNumber = left.SerialNumber if err := store.CreateDevice(ctx, left); err != nil { t.Fatal(err) } if err := store.CreateDevice(ctx, right); err != nil { t.Fatal(err) } now := time.Date(2026, 8, 7, 0, 0, 0, 0, time.UTC) if err := store.MarkReconciled(ctx, left.ID, 1, now); err != nil { t.Fatal(err) } if err := store.UpdateActualState(ctx, left.ID, device.ActualOnline, now); err != nil { t.Fatal(err) } snapshot, err := store.ConvergenceSnapshot(ctx) if err != nil { t.Fatal(err) } if snapshot.Total != 2 || snapshot.Unconverged != 1 { t.Fatalf("unexpected convergence snapshot: %+v", snapshot) } if err := store.RequestReconcile(ctx, left.ID, now.Add(time.Second)); err != nil { t.Fatal(err) } snapshot, err = store.ConvergenceSnapshot(ctx) if err != nil { t.Fatal(err) } if snapshot.Unconverged != 2 { t.Fatalf("runtime reconcile request was not persisted: %+v", snapshot) } } func TestPostgresRecordsQuotaSourceVersion(t *testing.T) { store, admin := openPostgresTestStore(t) ctx := context.Background() insertBellSite(t, admin, "tenant", "site", 2) if err := store.CreateDevice(ctx, videoDevice(1, "tenant", "site")); err != nil { t.Fatal(err) } var firstVersion int64 if err := store.db.QueryRowContext(ctx, `SELECT quota_source_version FROM sense.devices WHERE id = 'camera-001'`).Scan(&firstVersion); err != nil { t.Fatal(err) } if _, err := admin.ExecContext(ctx, `UPDATE bell.sites SET name = name WHERE tenant_id = 'tenant' AND id = 'site'`); err != nil { t.Fatal(err) } second := videoDevice(2, "tenant", "site") if err := store.CreateDevice(ctx, second); err != nil { t.Fatal(err) } var secondVersion int64 if err := store.db.QueryRowContext(ctx, `SELECT quota_source_version FROM sense.devices WHERE id = 'camera-002'`).Scan(&secondVersion); err != nil { t.Fatal(err) } if firstVersion != 1 || secondVersion != 2 { t.Fatalf("expected recorded projection versions 1 and 2, got %d and %d", firstVersion, secondVersion) } } func TestPostgresReconcilePortParity(t *testing.T) { store, admin := openPostgresTestStore(t) ctx := context.Background() insertBellSite(t, admin, "tenant", "site", 1) if err := store.CreateDevice(ctx, videoDevice(1, "tenant", "site")); err != nil { t.Fatal(err) } now := time.Now().UTC() due, err := store.ListDueReconcile(ctx, now.Add(time.Second), 10) if err != nil || len(due) != 1 || len(due[0].Device.Capabilities) != 2 { t.Fatalf("unexpected due reconcile list: values=%+v error=%v", due, err) } enabled, err := store.ListEnabledVideoDevices(ctx, 10) if err != nil || len(enabled) != 1 || enabled[0].ID != "camera-001" { t.Fatalf("unexpected enabled video list: values=%+v error=%v", enabled, err) } if err := store.SetDesiredState(ctx, "camera-001", device.DesiredEnabled); err != nil { t.Fatal(err) } value, err := store.GetDevice(ctx, "camera-001") if err != nil || value.Generation != 1 { t.Fatalf("idempotent desired state changed generation: value=%+v error=%v", value, err) } nextAttempt := now.Add(time.Minute) if err := store.MarkReconcileFailure(ctx, "camera-001", 1, nextAttempt, "test_failure", now); err != nil { t.Fatal(err) } due, err = store.ListDueReconcile(ctx, now.Add(30*time.Second), 10) if err != nil || len(due) != 0 { t.Fatalf("backoff device became due early: values=%+v error=%v", due, err) } due, err = store.ListDueReconcile(ctx, now.Add(2*time.Minute), 10) if err != nil || len(due) != 1 || due[0].FailureCount != 1 { t.Fatalf("backoff device did not become due: values=%+v error=%v", due, err) } value, err = store.GetDevice(ctx, "camera-001") if err != nil || value.ActualState != device.ActualFailed { t.Fatalf("reconcile failure did not persist actual state: value=%+v error=%v", value, err) } } func TestPostgresOpenRejectsOverprivilegedRuntimeRole(t *testing.T) { _, admin := openPostgresTestStore(t) ctx := context.Background() if _, err := admin.ExecContext(ctx, `GRANT UPDATE ON bell.site_quota_v1 TO yovision_t009_sense`); err != nil { t.Fatal(err) } defer func() { _, _ = admin.ExecContext(context.Background(), `REVOKE UPDATE ON bell.site_quota_v1 FROM yovision_t009_sense`) }() value, err := OpenPostgres(ctx, os.Getenv(postgresTestDSNEnv)) if value != nil { _ = value.Close() t.Fatal("overprivileged runtime role was accepted") } if err == nil || !strings.Contains(err.Error(), "privilege boundary") { t.Fatalf("expected privilege-boundary error, got %v", err) } } func openPostgresTestStore(t *testing.T) (*Postgres, *sql.DB) { t.Helper() dsn := os.Getenv(postgresTestDSNEnv) adminDSN := os.Getenv(postgresTestAdminDSNEnv) if dsn == "" || adminDSN == "" { t.Skip("PostgreSQL integration DSNs are not configured") } admin, err := sql.Open("pgx", adminDSN) if err != nil { t.Fatal(err) } if err := admin.PingContext(context.Background()); err != nil { admin.Close() t.Fatal("connect PostgreSQL test administrator") } if _, err := admin.ExecContext(context.Background(), `TRUNCATE sense.device_capabilities, sense.reconcile_state, sense.devices, sense.site_quota_projection_state, bell.sites CASCADE`); err != nil { admin.Close() t.Fatal(err) } store, err := OpenPostgres(context.Background(), dsn) if err != nil { admin.Close() t.Fatal(err) } t.Cleanup(func() { _ = store.Close() _ = admin.Close() }) return store, admin } func insertBellSite(t *testing.T, admin *sql.DB, tenantID, siteID string, quota int) { t.Helper() query := `INSERT INTO bell.sites(tenant_id, id, name, max_video_channels) VALUES ($1, $2, $3, $4)` arguments := []any{tenantID, siteID, "Test Site", quota} if quota == 0 { query = `INSERT INTO bell.sites(tenant_id, id, name) VALUES ($1, $2, $3)` arguments = arguments[:3] } if _, err := admin.ExecContext(context.Background(), query, arguments...); err != nil { t.Fatal(err) } }