Files
QiuSW 12857fdf32
Harness governance / validate (push) Has been cancelled
Harness governance / validate (pull_request) Has been cancelled
feat(sense): add reconciliation safety controls [T-012]
2026-08-07 23:00:03 +08:00

1241 lines
45 KiB
Go

package store
import (
"context"
"crypto/sha256"
"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",
AreaID: "area-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, &quotaError) || 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, &quotaError) || 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, &quotaError) {
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, ErrAreaPolicyUnavailable) {
t.Fatalf("missing Area projection must fail closed first, 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, &quotaError) || 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_t012_sense`); err != nil {
t.Fatal(err)
}
defer func() {
_, _ = admin.ExecContext(context.Background(),
`REVOKE UPDATE ON bell.site_quota_v1 FROM yovision_t012_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 TestPostgresOpenRejectsPublicReconciliationStatePrivilege(t *testing.T) {
_, admin := openPostgresTestStore(t)
ctx := context.Background()
if _, err := admin.ExecContext(ctx,
`GRANT SELECT ON sense.orphan_scan_runs TO PUBLIC`); err != nil {
t.Fatal(err)
}
defer func() {
_, _ = admin.ExecContext(context.Background(),
`REVOKE SELECT ON sense.orphan_scan_runs FROM PUBLIC`)
}()
value, err := OpenPostgres(ctx, os.Getenv(postgresTestDSNEnv))
if value != nil {
_ = value.Close()
t.Fatal("PUBLIC reconciliation state privilege was accepted")
}
if err == nil || !strings.Contains(err.Error(), "reconciliation safety privilege boundary") {
t.Fatalf("expected reconciliation privilege-boundary error, got %v", err)
}
}
func TestPostgresAreaPolicyAllowsNonImagingAndDeniesImagingCreate(t *testing.T) {
store, admin := openPostgresTestStore(t)
ctx := context.Background()
insertBellSite(t, admin, "tenant", "site", 2)
if _, err := admin.ExecContext(ctx, `UPDATE bell.areas SET capture_policy = 'non_imaging_only'
WHERE tenant_id = 'tenant' AND id = 'area-default'`); err != nil {
t.Fatal(err)
}
blocked := videoDevice(1, "tenant", "site")
blocked.DesiredState = device.DesiredDisabled
err := store.CreateDevice(ctx, blocked)
if !errors.Is(err, ErrAreaPolicyDenied) {
t.Fatalf("disabled imaging create must still be denied, got %v", err)
}
radar := device.Device{
ID: "radar-001", TenantID: "tenant", SiteID: "site", AreaID: "area-default",
SerialNumber: "radar-001", 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-imaging device must be allowed: %v", err)
}
stored, err := store.GetDevice(ctx, radar.ID)
if err != nil {
t.Fatal(err)
}
if stored.AreaID != "area-default" || stored.AreaPolicySourceVersion != 2 {
t.Fatalf("Area projection evidence was not stored: %+v", stored)
}
var devices, audits int
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.devices`).Scan(&devices); err != nil {
t.Fatal(err)
}
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.device_operation_outbox`).Scan(&audits); err != nil {
t.Fatal(err)
}
if devices != 1 || audits != 1 {
t.Fatalf("denied create left partial state: devices=%d audits=%d", devices, audits)
}
}
func TestPostgresAreaProjectionMissingAndRollbackFailClosed(t *testing.T) {
store, admin := openPostgresTestStore(t)
ctx := context.Background()
insertBellSite(t, admin, "tenant", "site", 2)
missing := videoDevice(1, "tenant", "site")
missing.AreaID = "missing-area"
if err := store.CreateDevice(ctx, missing); !errors.Is(err, ErrAreaPolicyUnavailable) {
t.Fatalf("missing Area projection must fail closed, got %v", err)
}
if _, err := store.db.ExecContext(ctx, `INSERT INTO sense.area_policy_projection_state(
tenant_id, site_id, area_id, source_version, synced_at
) VALUES ('tenant', 'site', 'area-default', 99, clock_timestamp())`); err != nil {
t.Fatal(err)
}
if err := store.CreateDevice(ctx, videoDevice(2, "tenant", "site")); !errors.Is(err, ErrAreaPolicyInvalid) {
t.Fatalf("Area source-version rollback must fail closed, got %v", err)
}
var devices, audits int
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.devices`).Scan(&devices); err != nil {
t.Fatal(err)
}
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.device_operation_outbox`).Scan(&audits); err != nil {
t.Fatal(err)
}
if devices != 0 || audits != 0 {
t.Fatalf("failed Area admissions persisted state: devices=%d audits=%d", devices, audits)
}
}
func TestPostgresConcurrentAreaObservationIsMonotonic(t *testing.T) {
store, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 2)
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()
value := videoDevice(index, "tenant", "site")
value.DesiredState = device.DesiredDisabled
<-start
errorsFound <- store.CreateDevice(ctx, value)
}(index)
}
close(start)
wait.Wait()
close(errorsFound)
for err := range errorsFound {
if err != nil {
t.Fatalf("concurrent Area observation failed: %v", err)
}
}
var sourceVersion int64
if err := admin.QueryRowContext(ctx, `SELECT source_version
FROM sense.area_policy_projection_state
WHERE tenant_id = 'tenant' AND site_id = 'site' AND area_id = 'area-default'`).
Scan(&sourceVersion); err != nil {
t.Fatal(err)
}
if sourceVersion != 1 {
t.Fatalf("concurrent observation recorded version %d", sourceVersion)
}
}
func TestPostgresAreaProjectionCannotCrossTenantBoundary(t *testing.T) {
store, admin := openPostgresTestStore(t)
ctx := context.Background()
insertBellSite(t, admin, "tenant-a", "site", 2)
insertBellSite(t, admin, "tenant-b", "site", 2)
insertBellArea(t, admin, "tenant-b", "site", "private-area", "video_allowed")
value := videoDevice(1, "tenant-a", "site")
value.AreaID = "private-area"
if err := store.CreateDevice(ctx, value); !errors.Is(err, ErrAreaPolicyUnavailable) {
t.Fatalf("cross-tenant Area was not hidden as unavailable: %v", err)
}
}
func TestPostgresEnableRechecksAreaWithoutStoppingExistingDevice(t *testing.T) {
store, admin := openPostgresTestStore(t)
ctx := context.Background()
insertBellSite(t, admin, "tenant", "site", 2)
disabled := videoDevice(1, "tenant", "site")
disabled.DesiredState = device.DesiredDisabled
if err := store.CreateDevice(ctx, disabled); err != nil {
t.Fatal(err)
}
if err := store.CreateDevice(ctx, videoDevice(2, "tenant", "site")); err != nil {
t.Fatal(err)
}
if _, err := admin.ExecContext(ctx, `UPDATE bell.areas SET capture_policy = 'non_imaging_only'
WHERE tenant_id = 'tenant' AND id = 'area-default'`); err != nil {
t.Fatal(err)
}
if err := store.SetDesiredState(ctx, disabled.ID, device.DesiredEnabled); !errors.Is(err, ErrAreaPolicyDenied) {
t.Fatalf("enable under non-imaging policy must be denied, got %v", err)
}
stillDisabled, err := store.GetDevice(ctx, disabled.ID)
if err != nil {
t.Fatal(err)
}
stillEnabled, err := store.GetDevice(ctx, "camera-002")
if err != nil {
t.Fatal(err)
}
if stillDisabled.DesiredState != device.DesiredDisabled ||
stillEnabled.DesiredState != device.DesiredEnabled {
t.Fatalf("policy change altered existing state: disabled=%s enabled=%s",
stillDisabled.DesiredState, stillEnabled.DesiredState)
}
}
func TestPostgresAuditContextRedactionAndNoopDesiredState(t *testing.T) {
store, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 2)
ctx := WithAuditContext(context.Background(), AuditContext{
ActorType: AuditActorUser, ActorID: "operator-7", Reason: "approved change", TraceID: "trace-7",
})
value := videoDevice(1, "tenant", "site")
if err := store.CreateDevice(ctx, value); err != nil {
t.Fatal(err)
}
if err := store.SetDesiredState(ctx, value.ID, device.DesiredEnabled); err != nil {
t.Fatal(err)
}
stored, err := store.GetDevice(ctx, value.ID)
if err != nil {
t.Fatal(err)
}
if stored.Generation != 1 {
t.Fatalf("no-op desired-state request changed generation to %d", stored.Generation)
}
rows, err := admin.QueryContext(ctx, `SELECT event_id, event_type, actor_type, actor_id,
COALESCE(reason, ''), COALESCE(trace_id, ''), payload::text
FROM sense.device_operation_outbox ORDER BY occurred_at, event_id`)
if err != nil {
t.Fatal(err)
}
defer rows.Close()
var count int
for rows.Next() {
var eventID, eventType, actorType, actorID, reason, traceID, payload string
if err := rows.Scan(&eventID, &eventType, &actorType, &actorID, &reason, &traceID, &payload); err != nil {
t.Fatal(err)
}
count++
if !strings.HasPrefix(eventID, "audit_") || len(eventID) != 38 {
t.Fatalf("invalid audit event ID %q", eventID)
}
if actorType != "user" || actorID != "operator-7" || reason != "approved change" || traceID != "trace-7" {
t.Fatalf("audit principal/context drift: %s/%s %s %s", actorType, actorID, reason, traceID)
}
for _, secret := range []string{value.EndpointRef, value.CredentialRef, value.PathName} {
if strings.Contains(payload, secret) {
t.Fatalf("audit payload leaked sensitive runtime data for %s", eventType)
}
}
}
if err := rows.Err(); err != nil {
t.Fatal(err)
}
if count != 2 {
t.Fatalf("expected create and no-op audit facts, got %d", count)
}
}
func TestPostgresOutboxFailureRollsBackAdmission(t *testing.T) {
store, admin := openPostgresTestStore(t)
ctx := context.Background()
insertBellSite(t, admin, "tenant", "site", 2)
if _, err := admin.ExecContext(ctx, `CREATE FUNCTION sense.t010_reject_outbox()
RETURNS trigger LANGUAGE plpgsql AS $function$
BEGIN RAISE EXCEPTION 'synthetic outbox failure'; END
$function$;
CREATE TRIGGER t010_reject_outbox BEFORE INSERT ON sense.device_operation_outbox
FOR EACH ROW EXECUTE FUNCTION sense.t010_reject_outbox()`); err != nil {
t.Fatal(err)
}
defer func() {
_, _ = admin.ExecContext(context.Background(),
`DROP TRIGGER IF EXISTS t010_reject_outbox ON sense.device_operation_outbox;
DROP FUNCTION IF EXISTS sense.t010_reject_outbox()`)
}()
if err := store.CreateDevice(ctx, videoDevice(1, "tenant", "site")); err == nil {
t.Fatal("synthetic Outbox failure did not reject device creation")
}
var devices, projections, audits int
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.devices`).Scan(&devices); err != nil {
t.Fatal(err)
}
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.area_policy_projection_state`).Scan(&projections); err != nil {
t.Fatal(err)
}
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.device_operation_outbox`).Scan(&audits); err != nil {
t.Fatal(err)
}
if devices != 0 || projections != 0 || audits != 0 {
t.Fatalf("Outbox failure left partial transaction: devices=%d projections=%d audits=%d",
devices, projections, audits)
}
}
func TestPostgresOutboxFailureRollsBackDesiredState(t *testing.T) {
store, admin := openPostgresTestStore(t)
ctx := context.Background()
insertBellSite(t, admin, "tenant", "site", 2)
value := videoDevice(1, "tenant", "site")
value.DesiredState = device.DesiredDisabled
if err := store.CreateDevice(ctx, value); err != nil {
t.Fatal(err)
}
if _, err := admin.ExecContext(ctx, `CREATE FUNCTION sense.t010_reject_state_audit()
RETURNS trigger LANGUAGE plpgsql AS $function$
BEGIN RAISE EXCEPTION 'synthetic state-audit failure'; END
$function$;
CREATE TRIGGER t010_reject_state_audit BEFORE INSERT ON sense.device_operation_outbox
FOR EACH ROW EXECUTE FUNCTION sense.t010_reject_state_audit()`); err != nil {
t.Fatal(err)
}
defer func() {
_, _ = admin.ExecContext(context.Background(),
`DROP TRIGGER IF EXISTS t010_reject_state_audit ON sense.device_operation_outbox;
DROP FUNCTION IF EXISTS sense.t010_reject_state_audit()`)
}()
if err := store.SetDesiredState(ctx, value.ID, device.DesiredEnabled); err == nil {
t.Fatal("synthetic Outbox failure did not reject desired-state change")
}
stored, err := store.GetDevice(ctx, value.ID)
if err != nil {
t.Fatal(err)
}
if stored.DesiredState != device.DesiredDisabled || stored.Generation != 1 {
t.Fatalf("Outbox failure committed desired state: %+v", stored)
}
var audits int
if err := admin.QueryRowContext(ctx, `SELECT count(*) FROM sense.device_operation_outbox`).Scan(&audits); err != nil {
t.Fatal(err)
}
if audits != 1 {
t.Fatalf("failed desired-state transaction changed Outbox count to %d", audits)
}
}
func TestPostgresOpenRejectsAreaSourcePrivilege(t *testing.T) {
_, admin := openPostgresTestStore(t)
ctx := context.Background()
if _, err := admin.ExecContext(ctx,
`GRANT SELECT ON bell.areas TO yovision_t012_sense`); err != nil {
t.Fatal(err)
}
defer func() {
_, _ = admin.ExecContext(context.Background(),
`REVOKE SELECT ON bell.areas FROM yovision_t012_sense`)
}()
value, err := OpenPostgres(ctx, os.Getenv(postgresTestDSNEnv))
if value != nil {
_ = value.Close()
t.Fatal("Area-source privilege was accepted")
}
if err == nil || !strings.Contains(err.Error(), "privilege boundary") {
t.Fatalf("expected Area privilege-boundary error, got %v", err)
}
}
func TestPostgresOpenRejectsPublicControlStatePrivilege(t *testing.T) {
_, admin := openPostgresTestStore(t)
ctx := context.Background()
if _, err := admin.ExecContext(ctx, `GRANT SELECT ON sense.batch_operations TO PUBLIC`); err != nil {
t.Fatal(err)
}
defer func() {
_, _ = admin.ExecContext(context.Background(), `REVOKE SELECT ON sense.batch_operations FROM PUBLIC`)
}()
value, err := OpenPostgres(ctx, os.Getenv(postgresTestDSNEnv))
if value != nil {
_ = value.Close()
t.Fatal("PUBLIC Control API table privilege was accepted")
}
if err == nil || !strings.Contains(err.Error(), "Control API state privilege boundary") {
t.Fatalf("expected Control API privilege-boundary error, got %v", err)
}
}
func TestPostgresControlCreateIdempotencyAndRedactedSnapshot(t *testing.T) {
postgres, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 2)
value := videoDevice(1, "tenant", "site")
value.ID = "dev_control_create_1"
value.DesiredState = device.DesiredDisabled
hash := sha256.Sum256([]byte("canonical-create"))
ctx := WithAuditContext(context.Background(), AuditContext{
ActorType: AuditActorService, ActorID: "bell-control", TraceID: "trace-control-create",
})
request := ControlCreateRequest{
Scope: IdempotencyScope{
PrincipalID: "bell-control", TenantID: "tenant", SiteID: "site",
Operation: "createDevice", Key: "create-control-0001",
RequestHash: hash, TraceID: "trace-control-create",
},
Device: value,
}
first, err := postgres.CreateControlDevice(ctx, request)
if err != nil {
t.Fatal(err)
}
request.Device.ID = "dev_control_create_retry"
replayed, err := postgres.CreateControlDevice(ctx, request)
if err != nil {
t.Fatal(err)
}
if !replayed.Replay || replayed.Device.ID != first.Device.ID || replayed.TraceID != first.TraceID {
t.Fatalf("create replay drifted: first=%+v replay=%+v", first, replayed)
}
conflictHash := sha256.Sum256([]byte("different-create"))
request.Scope.RequestHash = conflictHash
if _, err := postgres.CreateControlDevice(ctx, request); !errors.Is(err, ErrIdempotencyConflict) {
t.Fatalf("same key with different body was not rejected: %v", err)
}
var devices, receipts int
var responseBody string
if err := admin.QueryRow(`SELECT count(*) FROM sense.devices`).Scan(&devices); err != nil {
t.Fatal(err)
}
if err := admin.QueryRow(`SELECT count(*), min(response_body::text)
FROM sense.control_idempotency_receipts`).Scan(&receipts, &responseBody); err != nil {
t.Fatal(err)
}
if devices != 1 || receipts != 1 {
t.Fatalf("idempotent create counts drifted: devices=%d receipts=%d", devices, receipts)
}
for _, forbidden := range []string{value.EndpointRef, value.CredentialRef, "profile_token", "path_name"} {
if strings.Contains(responseBody, forbidden) {
t.Fatalf("idempotency response snapshot leaked %q", forbidden)
}
}
}
func TestPostgresControlListIsStableFilteredAndTenantScoped(t *testing.T) {
postgres, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 4)
insertBellSite(t, admin, "other", "site", 4)
for index := 1; index <= 3; index++ {
value := videoDevice(index, "tenant", "site")
value.ID = fmt.Sprintf("dev_list_%d", index)
value.DesiredState = device.DesiredDisabled
if index == 3 {
value.DesiredState = device.DesiredEnabled
}
if err := postgres.CreateDevice(context.Background(), value); err != nil {
t.Fatal(err)
}
}
other := videoDevice(9, "other", "site")
other.ID, other.DesiredState = "dev_list_other", device.DesiredDisabled
if err := postgres.CreateDevice(context.Background(), other); err != nil {
t.Fatal(err)
}
first, err := postgres.ListControlDevices(context.Background(), "tenant", "site", ControlListFilter{Limit: 1})
if err != nil || len(first.Items) != 1 || !first.HasMore || first.Quota.UsedVideoChannels != 1 {
t.Fatalf("unexpected first page: %+v %v", first, err)
}
second, err := postgres.ListControlDevices(context.Background(), "tenant", "site", ControlListFilter{
Limit: 1, AfterCreated: &first.Items[0].CreatedAt, AfterDeviceID: first.Items[0].ID,
})
if err != nil || len(second.Items) != 1 || second.Items[0].ID == first.Items[0].ID {
t.Fatalf("stable cursor position failed: %+v %v", second, err)
}
desired := device.DesiredEnabled
filtered, err := postgres.ListControlDevices(context.Background(), "tenant", "site", ControlListFilter{
Limit: 100, DesiredState: &desired,
})
if err != nil || len(filtered.Items) != 1 || filtered.Items[0].ID != "dev_list_3" {
t.Fatalf("desired-state filter or tenant scope failed: %+v %v", filtered, err)
}
if _, err := postgres.ListControlDevices(context.Background(), "tenant", "missing", ControlListFilter{Limit: 50}); !errors.Is(err, ErrNotFound) {
t.Fatalf("missing Site did not return not found: %v", err)
}
}
func TestPostgresControlPatchUsesETagAndAuditV2(t *testing.T) {
postgres, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 2)
value := videoDevice(1, "tenant", "site")
value.ID = "dev_control_patch_1"
value.DesiredState = device.DesiredDisabled
if err := postgres.CreateDevice(context.Background(), value); err != nil {
t.Fatal(err)
}
current, err := postgres.GetControlDevice(context.Background(), "tenant", "site", value.ID)
if err != nil {
t.Fatal(err)
}
etag := DeviceETag(current.ID, current.ResourceVersion)
newName, newProfile := "Updated camera", "profile-main"
ctx := WithAuditContext(context.Background(), AuditContext{
ActorType: AuditActorUser, ActorID: "operator-1", TraceID: "trace-control-patch",
})
updated, err := postgres.PatchControlDevice(ctx, "tenant", "site", value.ID, etag, ControlPatch{
Name: &newName, ProfileToken: &newProfile,
})
if err != nil {
t.Fatal(err)
}
if updated.Device.ResourceVersion != current.ResourceVersion+1 ||
updated.Device.Generation != current.Generation+1 || updated.ETag == etag {
t.Fatalf("patch did not advance versions: before=%+v after=%+v", current, updated)
}
if _, err := postgres.PatchControlDevice(ctx, "tenant", "site", value.ID, etag, ControlPatch{Name: &newName}); !errors.Is(err, ErrETagMismatch) {
t.Fatalf("stale ETag was accepted: %v", err)
}
var profileToken, eventType, payload string
if err := admin.QueryRow(`SELECT profile_token FROM sense.devices WHERE id = $1`, value.ID).Scan(&profileToken); err != nil {
t.Fatal(err)
}
if err := admin.QueryRow(`SELECT event_type, payload::text
FROM sense.device_operation_outbox WHERE event_type = 'device.configuration.accepted'`).
Scan(&eventType, &payload); err != nil {
t.Fatal(err)
}
if profileToken != newProfile || eventType != "device.configuration.accepted" ||
strings.Contains(payload, newProfile) || !strings.Contains(payload, "profile_token") {
t.Fatalf("configuration persistence/audit mismatch: profile=%q event=%q payload=%s", profileToken, eventType, payload)
}
}
func TestPostgresControlBatchIsPerItemDurableAndReplayable(t *testing.T) {
postgres, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 4)
first := videoDevice(1, "tenant", "site")
first.ID, first.DesiredState = "dev_batch_1", device.DesiredDisabled
second := videoDevice(2, "tenant", "site")
second.ID, second.DesiredState = "dev_batch_2", device.DesiredDisabled
for _, value := range []device.Device{first, second} {
if err := postgres.CreateDevice(context.Background(), value); err != nil {
t.Fatal(err)
}
}
hash := sha256.Sum256([]byte("canonical-batch"))
ctx := WithAuditContext(context.Background(), AuditContext{
ActorType: AuditActorUser, ActorID: "operator-1", Reason: "approved", TraceID: "trace-batch",
})
request := ControlBatchRequest{
Scope: IdempotencyScope{
PrincipalID: "operator-1", TenantID: "tenant", SiteID: "site",
Operation: "batchSetDeviceDesiredState", Key: "batch-control-0001",
RequestHash: hash, TraceID: "trace-batch",
},
Reason: "approved",
Items: []ControlBatchItem{
{DeviceID: first.ID, ETag: DeviceETag(first.ID, 1), DesiredState: device.DesiredEnabled},
{DeviceID: first.ID, ETag: DeviceETag(first.ID, 1), DesiredState: device.DesiredEnabled},
{DeviceID: second.ID, ETag: DeviceETag(second.ID, 1), DesiredState: device.DesiredEnabled},
},
}
operation, err := postgres.BatchSetControlDesiredState(ctx, request)
if err != nil {
t.Fatal(err)
}
if operation.Status != "partially_succeeded" || len(operation.Results) != 3 ||
operation.Results[0].Status != "rejected" || operation.Results[1].Status != "rejected" ||
operation.Results[2].Status != "succeeded" {
t.Fatalf("unexpected batch result: %+v", operation)
}
replayed, err := postgres.BatchSetControlDesiredState(ctx, request)
if err != nil || !replayed.Replay || replayed.ID != operation.ID {
t.Fatalf("batch replay drifted: %+v %v", replayed, err)
}
read, err := postgres.GetControlOperation(context.Background(), "tenant", operation.ID)
if err != nil || read.SiteID != "site" || len(read.Results) != 3 {
t.Fatalf("stored operation could not be read: %+v %v", read, err)
}
if _, err := postgres.GetControlOperation(context.Background(), "other-tenant", operation.ID); !errors.Is(err, ErrNotFound) {
t.Fatalf("cross-tenant operation was visible: %v", err)
}
var enabled, operations, receipts int
if err := admin.QueryRow(`SELECT count(*) FROM sense.devices WHERE desired_state = 'enabled'`).Scan(&enabled); err != nil {
t.Fatal(err)
}
if err := admin.QueryRow(`SELECT count(*) FROM sense.batch_operations`).Scan(&operations); err != nil {
t.Fatal(err)
}
if err := admin.QueryRow(`SELECT count(*) FROM sense.control_idempotency_receipts`).Scan(&receipts); err != nil {
t.Fatal(err)
}
if enabled != 1 || operations != 1 || receipts != 1 {
t.Fatalf("batch durability counts drifted: enabled=%d operations=%d receipts=%d", enabled, operations, receipts)
}
}
func TestPostgresConcurrentControlCreateExecutesOnce(t *testing.T) {
postgres, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 2)
hash := sha256.Sum256([]byte("concurrent-create"))
start := make(chan struct{})
results := make(chan ControlCreateResult, 2)
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()
value := videoDevice(index, "tenant", "site")
value.ID = fmt.Sprintf("dev_concurrent_%d", index)
value.SerialNumber = "same-semantic-serial"
value.DesiredState = device.DesiredDisabled
request := ControlCreateRequest{
Scope: IdempotencyScope{
PrincipalID: "service", TenantID: "tenant", SiteID: "site",
Operation: "createDevice", Key: "concurrent-create-0001",
RequestHash: hash, TraceID: fmt.Sprintf("trace-concurrent-%d", index),
}, Device: value,
}
<-start
result, err := postgres.CreateControlDevice(context.Background(), request)
if err != nil {
errorsFound <- err
return
}
results <- result
}(index)
}
close(start)
wait.Wait()
close(results)
close(errorsFound)
for err := range errorsFound {
t.Fatalf("concurrent idempotent create failed: %v", err)
}
ids := make(map[string]struct{})
for result := range results {
ids[result.Device.ID] = struct{}{}
}
if len(ids) != 1 {
t.Fatalf("concurrent create returned multiple resources: %+v", ids)
}
var devices, audits, receipts int
if err := admin.QueryRow(`SELECT count(*) FROM sense.devices`).Scan(&devices); err != nil {
t.Fatal(err)
}
if err := admin.QueryRow(`SELECT count(*) FROM sense.device_operation_outbox`).Scan(&audits); err != nil {
t.Fatal(err)
}
if err := admin.QueryRow(`SELECT count(*) FROM sense.control_idempotency_receipts`).Scan(&receipts); err != nil {
t.Fatal(err)
}
if devices != 1 || audits != 1 || receipts != 1 {
t.Fatalf("concurrent create executed more than once: devices=%d audits=%d receipts=%d", devices, audits, receipts)
}
}
func TestPostgresConcurrentControlBatchesUseStableDeviceLockOrder(t *testing.T) {
postgres, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 4)
ids := []string{"dev_lock_a", "dev_lock_b"}
for index, id := range ids {
value := videoDevice(index+1, "tenant", "site")
value.ID, value.DesiredState = id, device.DesiredDisabled
if err := postgres.CreateDevice(context.Background(), value); err != nil {
t.Fatal(err)
}
}
start := make(chan struct{})
operations := make(chan ControlBatchOperation, 2)
errorsFound := make(chan error, 2)
var wait sync.WaitGroup
for index := 0; index < 2; index++ {
wait.Add(1)
go func(index int) {
defer wait.Done()
order := ids
if index == 1 {
order = []string{ids[1], ids[0]}
}
hash := sha256.Sum256([]byte(fmt.Sprintf("batch-order-%d", index)))
request := ControlBatchRequest{
Scope: IdempotencyScope{
PrincipalID: "operator", TenantID: "tenant", SiteID: "site",
Operation: "batchSetDeviceDesiredState",
Key: fmt.Sprintf("batch-lock-order-%04d", index), RequestHash: hash,
TraceID: fmt.Sprintf("trace-lock-order-%d", index),
},
Reason: "concurrency test",
Items: []ControlBatchItem{
{DeviceID: order[0], ETag: DeviceETag(order[0], 1), DesiredState: device.DesiredEnabled},
{DeviceID: order[1], ETag: DeviceETag(order[1], 1), DesiredState: device.DesiredEnabled},
},
}
<-start
operation, err := postgres.BatchSetControlDesiredState(context.Background(), request)
if err != nil {
errorsFound <- err
return
}
operations <- operation
}(index)
}
close(start)
wait.Wait()
close(operations)
close(errorsFound)
for err := range errorsFound {
t.Fatalf("opposite-order batch failed or deadlocked: %v", err)
}
var succeeded, failed int
for operation := range operations {
switch operation.Status {
case "succeeded":
succeeded++
case "failed":
failed++
default:
t.Fatalf("unexpected concurrent batch status: %+v", operation)
}
}
if succeeded != 1 || failed != 1 {
t.Fatalf("expected one winner and one stale loser, got succeeded=%d failed=%d", succeeded, failed)
}
}
func TestPostgresConcurrentReconcileClaimHasOneWinner(t *testing.T) {
first, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 2)
if err := first.CreateDevice(context.Background(), videoDevice(1, "tenant", "site")); err != nil {
t.Fatal(err)
}
second, err := OpenPostgres(context.Background(), os.Getenv(postgresTestDSNEnv))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = second.Close() })
now := time.Date(2026, 8, 7, 0, 0, 0, 0, time.UTC)
start := make(chan struct{})
counts := make(chan int, 2)
errorsFound := make(chan error, 2)
var wait sync.WaitGroup
for index, repository := range []*Postgres{first, second} {
wait.Add(1)
go func(index int, repository *Postgres) {
defer wait.Done()
<-start
values, err := repository.ClaimDueReconcile(context.Background(), ReconcileClaim{
Owner: fmt.Sprintf("ins-%d", index), Token: fmt.Sprintf("token-%d", index),
Now: now, LeaseDuration: 30 * time.Second, Limit: 1,
})
if err != nil {
errorsFound <- err
return
}
counts <- len(values)
}(index, repository)
}
close(start)
wait.Wait()
close(counts)
close(errorsFound)
for err := range errorsFound {
t.Fatal(err)
}
total, winners := 0, 0
for count := range counts {
total += count
if count == 1 {
winners++
}
}
if total != 1 || winners != 1 {
t.Fatalf("due row was not exclusively claimed: total=%d winners=%d", total, winners)
}
}
func TestPostgresExpiredReconcileLeaseFencesOldWorkerAndRecordsOwnership(t *testing.T) {
postgres, admin := openPostgresTestStore(t)
insertBellSite(t, admin, "tenant", "site", 2)
value := videoDevice(1, "tenant", "site")
if err := postgres.CreateDevice(context.Background(), value); err != nil {
t.Fatal(err)
}
now := time.Date(2026, 8, 7, 0, 0, 0, 0, time.UTC)
first, err := postgres.ClaimDueReconcile(context.Background(), ReconcileClaim{
Owner: "ins-a", Token: "token-a", Now: now, LeaseDuration: 30 * time.Second, Limit: 1,
})
if err != nil || len(first) != 1 {
t.Fatalf("first claim failed: %+v %v", first, err)
}
early, err := postgres.ClaimDueReconcile(context.Background(), ReconcileClaim{
Owner: "ins-b", Token: "token-b", Now: now.Add(10 * time.Second),
LeaseDuration: 30 * time.Second, Limit: 1,
})
if err != nil || len(early) != 0 {
t.Fatalf("live lease was stolen: %+v %v", early, err)
}
if _, err := admin.Exec(`UPDATE sense.reconcile_state
SET lease_until = clock_timestamp() - interval '1 second'
WHERE device_id = $1`, value.ID); err != nil {
t.Fatal(err)
}
second, err := postgres.ClaimDueReconcile(context.Background(), ReconcileClaim{
Owner: "ins-b", Token: "token-b", Now: now.Add(31 * time.Second),
LeaseDuration: 30 * time.Second, Limit: 1,
})
if err != nil || len(second) != 1 {
t.Fatalf("expired lease was not recoverable: %+v %v", second, err)
}
if err := postgres.CompleteReconcile(
context.Background(), value.ID, value.Generation, "ins-a", "token-a", now.Add(32*time.Second),
); !errors.Is(err, ErrReconcileLeaseLost) {
t.Fatalf("old worker was not fenced: %v", err)
}
if err := postgres.CompleteReconcile(
context.Background(), value.ID, value.Generation, "ins-b", "token-b", now.Add(32*time.Second),
); err != nil {
t.Fatal(err)
}
var ownershipDevice string
if err := admin.QueryRow(`SELECT device_id FROM sense.media_path_ownership WHERE path_name = $1`,
value.PathName).Scan(&ownershipDevice); err != nil {
t.Fatal(err)
}
if ownershipDevice != value.ID {
t.Fatalf("wrong ownership was recorded: %q", ownershipDevice)
}
var leaseToken sql.NullString
if err := admin.QueryRow(`SELECT lease_token FROM sense.reconcile_state WHERE device_id = $1`,
value.ID).Scan(&leaseToken); err != nil {
t.Fatal(err)
}
if leaseToken.Valid {
t.Fatal("completion did not release the reconcile lease")
}
}
func TestPostgresOrphanReportLeaseAndCleanupAuditAreFencedAndIdempotent(t *testing.T) {
postgres, _ := openPostgresTestStore(t)
ctx := context.Background()
now := time.Date(2026, 8, 7, 0, 0, 0, 0, time.UTC)
acquired, err := postgres.AcquireOperationalLease(
ctx, OperationalLeaseOrphanScan, "ins-a", "scan-token-a", now, 30*time.Second,
)
if err != nil || !acquired {
t.Fatalf("scan lease failed: %v %v", acquired, err)
}
scan := OrphanScan{
ID: "scan_" + strings.Repeat("0", 26), InstanceID: "ins-a",
ObservedCount: 10, OwnedStaleCount: 1, UnownedCount: 1,
SafetyAllowed: true, SafetyReason: "allowed",
CompletedAt: now.Add(time.Second), ExpiresAt: now.Add(15 * time.Minute),
Findings: []OrphanFinding{
{PathName: "stale", Classification: OrphanOwnedStale, DeviceID: "old-device"},
{PathName: "unknown", Classification: OrphanUnowned},
},
}
if err := postgres.SaveOrphanScan(ctx, scan, "ins-a", "scan-token-a"); err != nil {
t.Fatal(err)
}
loaded, err := postgres.GetOrphanScan(ctx, scan.ID)
if err != nil || len(loaded.Findings) != 2 || !loaded.SafetyAllowed {
t.Fatalf("stored scan mismatch: %+v %v", loaded, err)
}
if err := postgres.RecordOrphanCleanup(
ctx, scan.ID, "stale", "operator", "deleted", "", now.Add(2*time.Second),
); err != nil {
t.Fatal(err)
}
if err := postgres.RecordOrphanCleanup(
ctx, scan.ID, "unknown", "operator", "deleted", "", now.Add(2*time.Second),
); err == nil {
t.Fatal("unowned path accepted a cleanup audit record")
}
if err := postgres.RecordOrphanCleanup(
ctx, scan.ID, "stale", "operator-2", "failed", "media_error", now.Add(3*time.Second),
); err != nil {
t.Fatal(err)
}
loaded, err = postgres.GetOrphanScan(ctx, scan.ID)
if err != nil || !loaded.Findings[0].Deleted {
t.Fatalf("successful cleanup was downgraded: %+v %v", loaded, err)
}
acquired, err = postgres.AcquireOperationalLease(
ctx, OperationalLeaseOrphanScan, "ins-b", "scan-token-b", now.Add(31*time.Second), 30*time.Second,
)
if err != nil || !acquired {
t.Fatalf("expired scan lease was not recoverable: %v %v", acquired, err)
}
staleScan := scan
staleScan.ID = "scan_" + strings.Repeat("1", 26)
staleScan.CompletedAt = now.Add(32 * time.Second)
staleScan.ExpiresAt = staleScan.CompletedAt.Add(15 * time.Minute)
if err := postgres.SaveOrphanScan(ctx, staleScan, "ins-a", "scan-token-a"); !errors.Is(err, ErrOperationalLeaseLost) {
t.Fatalf("stale scan worker was not fenced: %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.orphan_cleanup_actions,
sense.orphan_scan_findings,
sense.orphan_scan_runs,
sense.operational_leases,
sense.media_path_ownership,
sense.control_idempotency_receipts,
sense.batch_operation_items,
sense.batch_operations,
sense.device_operation_outbox,
sense.device_capabilities,
sense.reconcile_state,
sense.devices,
sense.site_quota_projection_state,
sense.area_policy_projection_state,
bell.areas,
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)
}
insertBellArea(t, admin, tenantID, siteID, "area-default", "video_allowed")
}
func insertBellArea(
t *testing.T,
admin *sql.DB,
tenantID, siteID, areaID, capturePolicy string,
) {
t.Helper()
if _, err := admin.ExecContext(context.Background(), `INSERT INTO bell.areas(
tenant_id, site_id, id, name, capture_policy
) VALUES ($1, $2, $3, $4, $5)`, tenantID, siteID, areaID, "Test Area", capturePolicy); err != nil {
t.Fatal(err)
}
}