package postgresadmin import ( "context" "encoding/json" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgtype" "proxy-pool/internal/domain/adminstate" ) func (adapter *adapter) Claim( ctx context.Context, command adminstate.ClaimCommand, ) ([]adminstate.Event, error) { if err := contextError(ctx); err != nil { return nil, err } if command.Validate() != nil || !adapter.valid() { return nil, adminstate.ErrInvalidCommand } tx, err := adapter.begin(ctx, pgx.TxOptions{AccessMode: pgx.ReadWrite}, "begin outbox claim") if err != nil { return nil, err } defer rollback(tx) now := utc(command.Now) claimUntil := now.Add(command.Lease) rows, err := tx.Query(ctx, ` WITH candidates AS ( SELECT id FROM admin_outbox WHERE published_at IS NULL AND (claim_until IS NULL OR claim_until <= $1) ORDER BY id FOR UPDATE SKIP LOCKED LIMIT $2 ), claimed AS ( UPDATE admin_outbox AS outbox SET claim_owner = $3, claim_until = $4 FROM candidates WHERE outbox.id = candidates.id RETURNING outbox.id, outbox.revision, outbox.event_type, outbox.aggregate_type, outbox.aggregate_id, outbox.payload, outbox.occurred_at, outbox.claim_owner, outbox.claim_until, outbox.published_at ) SELECT id, revision, event_type, aggregate_type, aggregate_id, payload, occurred_at, claim_owner, claim_until, published_at FROM claimed ORDER BY id`, now, command.Limit, command.ConsumerID, claimUntil) if err != nil { return nil, databaseError(ctx, "claim outbox events", err) } defer rows.Close() events := make([]adminstate.Event, 0, command.Limit) for rows.Next() { event := adminstate.Event{} var id int64 var revision int64 var payload []byte var claimedBy pgtype.Text var claimedUntil pgtype.Timestamptz var publishedAt pgtype.Timestamptz if err := rows.Scan( &id, &revision, &event.Type, &event.AggregateType, &event.AggregateID, &payload, &event.OccurredAt, &claimedBy, &claimedUntil, &publishedAt, ); err != nil { return nil, databaseError(ctx, "decode claimed outbox events", err) } convertedID, idOK := domainID(id) convertedRevision, revisionOK := domainID(revision) if !idOK || !revisionOK || !claimedBy.Valid || !claimedUntil.Valid { return nil, unavailable("decode claimed outbox events") } event.ID = convertedID event.Revision = convertedRevision event.Payload = append(json.RawMessage(nil), payload...) event.OccurredAt = utc(event.OccurredAt) event.ClaimedBy = claimedBy.String value := utc(claimedUntil.Time) event.ClaimUntil = &value if publishedAt.Valid { value := utc(publishedAt.Time) event.PublishedAt = &value } events = append(events, event) } if err := rows.Err(); err != nil { return nil, databaseError(ctx, "claim outbox events", err) } if err := commit(ctx, tx, "commit outbox claim"); err != nil { return nil, err } return events, nil } func (adapter *adapter) Acknowledge( ctx context.Context, command adminstate.AcknowledgeCommand, ) error { if err := contextError(ctx); err != nil { return err } if command.Validate() != nil || !adapter.valid() { return adminstate.ErrInvalidCommand } eventIDs, ok := eventDatabaseIDs(command.EventIDs) if !ok { return adminstate.ErrNotFound } tx, err := adapter.begin(ctx, pgx.TxOptions{AccessMode: pgx.ReadWrite}, "begin outbox acknowledge") if err != nil { return err } defer rollback(tx) rows, err := tx.Query(ctx, ` SELECT id, published_at, claim_owner, claim_until FROM admin_outbox WHERE id = ANY($1::bigint[]) ORDER BY id FOR UPDATE`, eventIDs) if err != nil { return databaseError(ctx, "lock acknowledged outbox events", err) } count := 0 now := utc(command.Now) conflict := false for rows.Next() { var id int64 var publishedAt pgtype.Timestamptz var claimOwner pgtype.Text var claimUntil pgtype.Timestamptz if err := rows.Scan(&id, &publishedAt, &claimOwner, &claimUntil); err != nil { rows.Close() return databaseError(ctx, "decode acknowledged outbox events", err) } count++ if publishedAt.Valid || !claimOwner.Valid || claimOwner.String != command.ConsumerID || !claimUntil.Valid || !now.Before(claimUntil.Time) { conflict = true } } if err := rows.Err(); err != nil { rows.Close() return databaseError(ctx, "lock acknowledged outbox events", err) } rows.Close() if count != len(eventIDs) { return adminstate.ErrNotFound } if conflict { return adminstate.ErrConflict } commandTag, err := tx.Exec(ctx, ` UPDATE admin_outbox SET published_at = $1 WHERE id = ANY($2::bigint[])`, now, eventIDs) if err != nil { return databaseError(ctx, "publish outbox events", err) } if commandTag.RowsAffected() != int64(len(eventIDs)) { return unavailable("publish outbox events") } return commit(ctx, tx, "commit outbox acknowledge") }