177 lines
6.4 KiB
Go
177 lines
6.4 KiB
Go
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
|
|
}
|