MCPcopy Create free account
hub / github.com/devaccuracy/ledgerforge / ClaimPendingOutboxEntries

Method ClaimPendingOutboxEntries

database/lineage.go:279–349  ·  view source on GitHub ↗

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)

Source from the content-addressed store, hash-verified

277// ClaimPendingOutboxEntries claims a batch of pending outbox entries for processing.
278// It uses SELECT FOR UPDATE SKIP LOCKED to allow concurrent processors.
279func (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 }

Callers

nothing calls this directly

Calls 2

NewAPIErrorFunction · 0.92
CloseMethod · 0.45

Tested by

no test coverage detected