package store import ( "context" "database/sql" "encoding/json" "errors" "regexp" "strings" "time" "yovision/sense/internal/auditrelay" ) var relayErrorCode = regexp.MustCompile(`^[a-z][a-z0-9_]{0,63}$`) func (s *Postgres) AuditRelayReady(ctx context.Context) error { var version int64 if err := s.db.QueryRowContext(ctx, `SELECT COALESCE(MAX(version), 0) FROM sense.schema_migrations`).Scan(&version); err != nil || version < 6 { return errors.New("postgres Sense schema migration v6 is required for audit relay") } var canUseOutbox, canReadBell, canWriteBell bool if err := s.db.QueryRowContext(ctx, `SELECT has_table_privilege(current_user, 'sense.device_operation_outbox', 'SELECT,INSERT,UPDATE,DELETE'), has_table_privilege(current_user, 'bell.audit_events', 'SELECT'), has_table_privilege(current_user, 'bell.audit_events', 'INSERT,UPDATE,DELETE')`).Scan( &canUseOutbox, &canReadBell, &canWriteBell, ); err != nil { return errors.New("verify audit relay privileges") } if !canUseOutbox || canReadBell || canWriteBell { return errors.New("Sense audit relay violates schema ownership boundary") } return nil } func (s *Postgres) ClaimAuditRelayBatch( ctx context.Context, owner string, limit int, lease time.Duration, ) ([]auditrelay.QueuedEvent, error) { if strings.TrimSpace(owner) == "" || limit < 1 || limit > auditrelay.MaxBatchSize || lease <= 0 { return nil, errors.New("invalid audit relay claim") } rows, err := s.db.QueryContext(ctx, `WITH due AS ( SELECT event_id FROM sense.device_operation_outbox WHERE delivered_at IS NULL AND dead_lettered_at IS NULL AND COALESCE(next_attempt_at, available_at) <= clock_timestamp() AND (relay_lease_until IS NULL OR relay_lease_until <= clock_timestamp()) ORDER BY COALESCE(next_attempt_at, available_at), event_id FOR UPDATE SKIP LOCKED LIMIT $1 ) UPDATE sense.device_operation_outbox AS outbox SET relay_lease_owner = $2, relay_lease_token = outbox.relay_lease_token + 1, relay_lease_until = clock_timestamp() + ($3 * interval '1 second'), attempt_count = outbox.attempt_count + 1, last_error_code = NULL FROM due WHERE outbox.event_id = due.event_id RETURNING outbox.event_id, outbox.event_type, outbox.tenant_id, outbox.site_id, outbox.device_id, outbox.actor_type, outbox.actor_id, outbox.reason, outbox.trace_id, outbox.aggregate_generation, outbox.quota_source_version, outbox.area_policy_source_version, outbox.payload, outbox.occurred_at, outbox.relay_lease_token, outbox.attempt_count`, limit, owner, lease.Seconds()) if err != nil { return nil, errors.New("claim audit relay batch") } defer rows.Close() values := make([]auditrelay.QueuedEvent, 0) for rows.Next() { var value auditrelay.QueuedEvent var reason, trace sql.NullString var quota, area sql.NullInt64 var payload []byte if err := rows.Scan( &value.Event.EventID, &value.Event.EventType, &value.Event.TenantID, &value.Event.SiteID, &value.Event.DeviceID, &value.Event.Actor.Type, &value.Event.Actor.ID, &reason, &trace, &value.Event.AggregateGeneration, "a, &area, &payload, &value.Event.OccurredAt, &value.LeaseToken, &value.AttemptCount, ); err != nil { return nil, errors.New("scan audit relay claim") } if reason.Valid { value.Event.Reason = &reason.String } if trace.Valid { value.Event.TraceID = &trace.String } if quota.Valid { value.Event.ProjectionVersions.QuotaSourceVersion = "a.Int64 } if area.Valid { value.Event.ProjectionVersions.AreaPolicySourceVersion = &area.Int64 } if !json.Valid(payload) { return nil, errors.New("invalid audit payload in outbox") } value.Event.Data = append(json.RawMessage(nil), payload...) value.SchemaVersion = 1 if value.Event.EventType == "device.configuration.accepted" { value.SchemaVersion = 2 } values = append(values, value) } if err := rows.Err(); err != nil { return nil, errors.New("iterate audit relay claims") } return values, nil } func (s *Postgres) CompleteAuditRelayBatch( ctx context.Context, owner string, values []auditrelay.Completion, ) error { if strings.TrimSpace(owner) == "" || len(values) == 0 || len(values) > auditrelay.MaxBatchSize { return errors.New("invalid audit relay completion") } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return errors.New("begin audit relay completion") } defer tx.Rollback() for _, value := range values { code := value.ErrorCode if code != "" && !relayErrorCode.MatchString(code) { code = "invalid_response" } var result sql.Result switch value.Disposition { case auditrelay.Delivered: result, err = tx.ExecContext(ctx, `UPDATE sense.device_operation_outbox SET delivered_at = clock_timestamp(), last_error_code = NULL, relay_lease_owner = NULL, relay_lease_until = NULL WHERE event_id = $1 AND relay_lease_owner = $2 AND relay_lease_token = $3 AND relay_lease_until > clock_timestamp()`, value.EventID, owner, value.LeaseToken) case auditrelay.DeadLetter: if code == "" { return errors.New("dead-letter completion requires an error code") } result, err = tx.ExecContext(ctx, `UPDATE sense.device_operation_outbox SET dead_lettered_at = clock_timestamp(), last_error_code = $4, relay_lease_owner = NULL, relay_lease_until = NULL WHERE event_id = $1 AND relay_lease_owner = $2 AND relay_lease_token = $3 AND relay_lease_until > clock_timestamp()`, value.EventID, owner, value.LeaseToken, code) case auditrelay.Retry: if code == "" || value.RetryAfter < time.Second || value.RetryAfter > 300*time.Second { return errors.New("invalid audit relay retry completion") } result, err = tx.ExecContext(ctx, `UPDATE sense.device_operation_outbox SET next_attempt_at = clock_timestamp() + ($4 * interval '1 second'), last_error_code = $5, relay_lease_owner = NULL, relay_lease_until = NULL WHERE event_id = $1 AND relay_lease_owner = $2 AND relay_lease_token = $3 AND relay_lease_until > clock_timestamp()`, value.EventID, owner, value.LeaseToken, value.RetryAfter.Seconds(), code) default: return errors.New("invalid audit relay disposition") } if err != nil { return errors.New("persist audit relay completion") } affected, err := result.RowsAffected() if err != nil { return errors.New("read audit relay completion") } if affected != 1 { return auditrelay.ErrLeaseLost } } if err := tx.Commit(); err != nil { return errors.New("commit audit relay completion") } return nil }