ClaimPendingOutboxEntries claims a batch of pending outbox entries for processing. It uses SELECT FOR UPDATE SKIP LOCKED to allow concurrent processors.
(ctx context.Context, batchSize int, lockDuration time.Duration)
| 277 | // ClaimPendingOutboxEntries claims a batch of pending outbox entries for processing. |
| 278 | // It uses SELECT FOR UPDATE SKIP LOCKED to allow concurrent processors. |
| 279 | func (d Datasource) ClaimPendingOutboxEntries(ctx context.Context, batchSize int, lockDuration time.Duration) ([]model.LineageOutbox, error) { |
| 280 | lockedUntil := time.Now().Add(lockDuration) |
| 281 | |
| 282 | query := ` |
| 283 | UPDATE ledgerforge.lineage_outbox |
| 284 | SET status = $1, locked_until = $2 |
| 285 | WHERE id IN ( |
| 286 | SELECT id FROM ledgerforge.lineage_outbox |
| 287 | WHERE status = 'pending' |
| 288 | AND (locked_until IS NULL OR locked_until < NOW()) |
| 289 | AND attempts < max_attempts |
| 290 | ORDER BY created_at ASC |
| 291 | LIMIT $3 |
| 292 | FOR UPDATE SKIP LOCKED |
| 293 | ) |
| 294 | RETURNING id, transaction_id, source_balance_id, destination_balance_id, provider, lineage_type, payload, status, attempts, max_attempts, last_error, created_at, processed_at, locked_until, inflight |
| 295 | ` |
| 296 | |
| 297 | rows, err := d.Conn.QueryContext(ctx, query, model.OutboxStatusProcessing, lockedUntil, batchSize) |
| 298 | if err != nil { |
| 299 | return nil, apierror.NewAPIError(apierror.ErrInternalServer, "Failed to claim pending outbox entries", err) |
| 300 | } |
| 301 | defer func() { _ = rows.Close() }() |
| 302 | |
| 303 | var entries []model.LineageOutbox |
| 304 | for rows.Next() { |
| 305 | var entry model.LineageOutbox |
| 306 | var sourceBalanceID, destinationBalanceID, provider, lastError sql.NullString |
| 307 | var processedAt, lockedUntilVal sql.NullTime |
| 308 | |
| 309 | err := rows.Scan( |
| 310 | &entry.ID, |
| 311 | &entry.TransactionID, |
| 312 | &sourceBalanceID, |
| 313 | &destinationBalanceID, |
| 314 | &provider, |
| 315 | &entry.LineageType, |
| 316 | &entry.Payload, |
| 317 | &entry.Status, |
| 318 | &entry.Attempts, |
| 319 | &entry.MaxAttempts, |
| 320 | &lastError, |
| 321 | &entry.CreatedAt, |
| 322 | &processedAt, |
| 323 | &lockedUntilVal, |
| 324 | &entry.Inflight, |
| 325 | ) |
| 326 | if err != nil { |
| 327 | return nil, apierror.NewAPIError(apierror.ErrInternalServer, "Failed to scan outbox entry", err) |
| 328 | } |
| 329 | |
| 330 | entry.SourceBalanceID = sourceBalanceID.String |
| 331 | entry.DestinationBalanceID = destinationBalanceID.String |
| 332 | entry.Provider = provider.String |
| 333 | entry.LastError = lastError.String |
| 334 | if processedAt.Valid { |
| 335 | entry.ProcessedAt = &processedAt.Time |
| 336 | } |
nothing calls this directly
no test coverage detected