panurus

Transaction Recovery Service

The Transaction Recovery Service provides automatic re-registration of finality listeners for pending transactions that may have lost their listeners due to node restarts, network interruptions, or other failures. This ensures that transactions eventually reach finality even after system disruptions.

Architecture

The recovery system consists of three main components:

  1. Manager: Orchestrates the recovery process with periodic scanning and distributed coordination
  2. Handler: Implements the actual recovery logic for individual transactions
  3. Storage: Provides database operations for claiming and tracking recovery state

Components

Recovery Manager

The Manager runs in the background and periodically scans for pending transactions that are eligible for recovery. It uses distributed locking (PostgreSQL advisory locks) to ensure only one replica in a multi-instance deployment performs recovery at a time.

Key features:

Recovery Handler

The Handler interface defines how individual transactions are recovered. The TTX service provides a concrete implementation (TTXRecoveryHandler) that:

Storage Interface

The Storage interface abstracts database operations needed for recovery:

Database Support

PostgreSQL is the recommended database for production multi-instance deployments:

A recovery manager is started per TMS for both owner and audit storage. Both stores use the same atomic claim: at any instant, a pending row is claimed by at most one replica, and ReleaseRecoveryClaim frees it again as soon as the sweep is done with it. Claims live in the recovery_claimed_by and recovery_claim_expires_at columns of each store’s own requests table, so an audit claim never hides a row from the owner sweep or the other way round.

The claim does not by itself reduce how many times a still-pending row gets processed: a row that is not yet finalized is released unchanged at the end of the sweep that claimed it, so the next tick claims and processes it again, on this replica or another. What the atomic claim removes is concurrent processing of the same row: no two replicas ever run finality logic for the same transaction at the same moment. That matters once a sweep’s processing time approaches or exceeds the scan interval, or a backlog spans more than one tick.

Leader election currently differs between the two. The owner store elects a leader through a PostgreSQL advisory lock; the audit store is built without a leader factory, so AcquireRecoveryLeadership grants leadership locally to every replica. Audit sweeps therefore run everywhere at once, and it is the atomic claim rather than leader election that keeps each pending audit transaction from being claimed by more than one replica at a time.

The claim query also treats a row as free if it is already claimed by the calling replica’s own instanceID (recovery_claimed_by). Nothing in the recovery manager renews a lease or re-claims a row it already holds: runSweep claims, processes, and releases a row sequentially within one manager, so this case does not occur in normal operation. It exists so the store-level claim API stays safe to call twice with the same owner; the risk it creates is that if two replicas ever presented the same instanceID, this clause would let both claim the same rows, and the audit path has no advisory-lock backstop to catch it. Manager.Start() closes this: it appends a fresh per-process id to instanceID every time it starts, whether the value was left empty (auto-generated) or configured explicitly, so a shared config value can never collide across replicas. The owner string actually stored is therefore always <instanceID>-<generated suffix>, not the configured value verbatim, so expect that suffix when reading recovery_claimed_by off a row.

SQLite (Development and Single-Node)

SQLite is supported for single-node deployments and development:

Configuration

Recovery behavior is controlled via configuration (see Configuration):

recovery:
  enabled: true              # Enable/disable recovery
  ttl: 30s                   # Minimum age before recovery
  scanInterval: 5s           # How often to scan
  batchSize: 16              # Max transactions per scan
  workerCount: 8             # Parallel workers
  leaseDuration: 5m          # Claim lease duration
  transactionTimeout: 60s    # Deadline for a single transaction's recovery attempt (0 = unbounded, or >= 10s)
  instanceID: ""             # Instance identifier
  notFoundGracePeriod: 30m   # Promote NotFound rows to Orphan after this age (0 disables)
  stuckTransactionAlertThreshold: 5  # Escalate to Error log after this many consecutive timeouts (0 disables)

Usage Example

Creating a recovery manager:

config := recovery.Config{
    Enabled:                        true,
    TTL:                            30 * time.Second,
    ScanInterval:                   5 * time.Second,
    BatchSize:                      16,
    WorkerCount:                    8,
    LeaseDuration:                  5 * time.Minute,
    TransactionTimeout:             60 * time.Second,
    NotFoundGracePeriod:            30 * time.Minute,
    StuckTransactionAlertThreshold: 5,
}

manager := recovery.NewManager(
    logger,
    storage,  // Implements Storage interface
    handler,  // Implements Handler interface
    config,
)

// Start recovery
if err := manager.Start(); err != nil {
    return err
}
defer manager.Stop()

Implementing a Custom Handler

To implement a custom recovery handler:

type MyHandler struct {
    // your dependencies
}

func (h *MyHandler) Recover(ctx context.Context, txID string) error {
    // 1. Query transaction status from your backend
    // 2. Apply finality logic based on status
    // 3. Update local database state
    // 4. Return nil on success, error on failure
    return nil
}

Recovery Process Flow

  1. Manager acquires leadership (PostgreSQL advisory lock). Normally released at the end of each sweep, but held across sweeps for as long as any transaction is still abandoned in the background from a previous sweep (see step 5): this instance’s own next ClaimPendingTransactions call is what keeps that transaction’s claim lease alive (step 5), and that only reliably happens if this instance keeps winning leadership, rather than leaving it to the next non-blocking lock attempt succeeding by chance. Stop() always releases leadership when called, regardless of whether anything is still abandoned
  2. Manager queries for pending transactions older than TTL
  3. Manager atomically claims a batch of transactions, each returned as a RecoveryClaim (TxID + StoredAt)
  4. Manager distributes claimed transactions to worker pool
  5. Each worker calls Handler.Recover() for its transactions, bounded by transactionTimeout: the call runs in its own goroutine, and the worker moves on as soon as the deadline fires rather than waiting for the call to return, so a peer that hangs without honouring the context (as some ledger calls do) still cannot block the sweep indefinitely, and since leadership is held for as long as anything is abandoned (step 1), neither can every other replica in the meantime. The abandoned goroutine keeps running in the background until the underlying call eventually completes; its claim stays held (not released) in the meantime, and its eventual result is used, not discarded (see step 8). If the same transaction is still claimed and re-dispatched to a worker before that abandoned goroutine returns, the manager skips it rather than starting a second concurrent Recover call for it
  6. Handler queries network and applies finality logic
  7. If the handler reports NotFound and the row was stored more than notFoundGracePeriod ago, the manager promotes the row to Orphan via SetStatus so it exits the eligible scan range, unless this is a step-5 background finalization arriving after Stop() has already run, in which case the promotion is skipped: leadership was released by Stop() (step 1), so another replica may have since legitimately resolved the same transaction, and an unconditional SetStatus could overwrite that outcome
  8. Manager releases each claim with a success/failure message, as soon as a final result for it is known: synchronously, for a Recover() call that returns within transactionTimeout, or later in the background, for one that did not (step 5), whichever comes first, and only once. The in-flight guard for a transaction is cleared only after this release actually completes, not before, so a sweep that reclaims the same transaction in between cannot start a second attempt out from under the release still in progress
  9. Process repeats on next scan interval

Transaction Status Lifecycle

A token request transitions through the following statuses as the recovery loop interacts with it:

All three terminal statuses (Confirmed, Deleted, Orphan) are excluded from subsequent recovery sweeps by virtue of the status = Pending filter on the claim query.

Error Handling

Performance Tuning

For High-Throughput Environments

For Resource-Constrained Environments

For Long-Running Transaction Assembly

Transaction Timeout

Thread Safety

The Manager is thread-safe and can be safely started/stopped from multiple goroutines. The Handler implementation must also be thread-safe as it will be called concurrently by multiple workers.

Shutdown Behaviour

Stop() cancels the manager context and waits for the recovery loop to return. If a sweep is mid-batch, the fan-out to the worker pool aborts on cancellation: workers exit as soon as they observe the cancellation, and any claims not yet dispatched are simply left undispatched. Those rows keep their Pending status, so they become eligible again once their lease (leaseDuration) expires and are picked up by the next sweep — on this or another replica. The aborted sweep reports a recovery fan-out cancelled error. Because this is an ordinary shutdown rather than a failure, the loop logs it at debug level and only warns for genuine sweep errors. The sweep summary counts successes against the claims actually dispatched, so a partial sweep logs claimed=N, dispatched=M, ... at warn level instead of crediting the undispatched tail as succeeded.

With a transactionTimeout configured, a worker mid-Handler.Recover() does not wait out that call on shutdown: it abandons the call once the deadline fires and moves on (see Transaction Timeout), so Stop() returning does not guarantee every in-flight recovery attempt has actually stopped. An abandoned call can keep running after Stop() returns and eventually reach Commit, writing into a store the SDK may already be tearing down as part of node shutdown. This is the same accepted tradeoff the timeout mechanism makes everywhere else: bounding the worker takes priority over waiting for a ctx-blind call to finish. With transactionTimeout: 0 there is no deadline to abandon the call at, so Stop() instead blocks on Handler.Recover() directly, potentially indefinitely if the handler ignores its context.

The claim on a transaction abandoned this way is not released by Stop() either: it is released only once the abandoned call actually returns and the manager finalizes it in the background, which can be well after Stop() has already returned. That finalization uses its own independent context rather than the (by then long since cancelled) sweep or manager context, specifically so the release can still succeed after shutdown.

Recovery leadership, unlike the claim above, is always released by Stop(), even if a transaction is still abandoned in the background: this instance can no longer productively use it once its own sweep loop has exited, and holding onto it would only block a live peer that could otherwise take over. That means a transaction still abandoned when Stop() runs is no longer guaranteed to have its lease renewed by this instance, and a live peer can legitimately reclaim and finish it while this instance’s own attempt is still finishing up in the background. For that reason, if the background finalization’s handler eventually reports the transaction as NotFound, the manager does not promote it to Orphan once it observes that Stop() already ran: another replica may have since resolved it to something else, and unlike the owner-scoped claim release, SetStatus has no way to tell that has happened.