Files
yovision/Sense/internal/store/audit_relay_postgres.go
T

177 lines
6.4 KiB
Go
Raw Normal View History

2026-08-11 00:24:32 +08:00
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,
&quota, &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 = &quota.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
}