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.
The recovery system consists of three main components:
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:
The Handler interface defines how individual transactions are recovered. The TTX service provides a concrete implementation (TTXRecoveryHandler) that:
The Storage interface abstracts database operations needed for recovery:
AcquireRecoveryLeadership: Obtains distributed lock for leader electionClaimPendingTransactions: Atomically claims a batch of pending transactions, returning a lightweight RecoveryClaim (TxID + StoredAt) for each row — the recovery loop only needs these two fields, so the SQL projection is kept narrowReleaseRecoveryClaim: Releases claim after processingSetStatus: Promotes a transaction to a terminal status. Used by the recovery loop to mark NotFound-past-grace-period rows as Orphan so they exit the eligible scan range without being conflated with ledger-rejected transactions (Deleted)PostgreSQL is the recommended database for production multi-instance deployments:
UPDATE...RETURNING ensures no duplicate claimsA 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 is supported for single-node deployments and development:
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)
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()
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
}
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 abandonedRecoveryClaim (TxID + StoredAt)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 itNotFound 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 outcomeRecover() 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 progressA token request transitions through the following statuses as the recovery loop interacts with it:
ClaimPendingTransactions; the claim query and its supporting partial index filter on status = Pending.network.Invalid) or by local validation (token request hash mismatch via the finality listener). Terminal.NotFound from the network past notFoundGracePeriod. Terminal in this version, and intentionally distinct from Deleted so operators (and future replay tooling) can identify broadcast failures separately from ledger-rejected transactions.All three terminal statuses (Confirmed, Deleted, Orphan) are excluded from subsequent recovery sweeps by virtue of the status = Pending filter on the claim query.
Deleted in the databaseNotFound past notFoundGracePeriod): Marked as Orphan to indicate the transaction never reached the ledger; distinct from Deleted so operators can distinguish broadcast failures from ledger-rejected transactionsbatchSize (200-500)workerCount (8-16)scanInterval (2-3s)batchSize (50)workerCount (2)scanInterval (10-15s)ttl (60s or more)leaseDuration > expected processing timetransactionTimeout bounds a single Handler.Recover() call (ledger status query plus finality logic), not the whole sweep. It is not a general no-stall guarantee for the sweep: ClaimPendingTransactions, ReleaseRecoveryClaim, and SetStatus still run on the sweep’s own context, which has no deadline, so a wedged database connection in any of those can still stall recovery on every replica of the TMS, the same failure mode transactionTimeout was added to fix on the ledger sidebatchSize claims split across workerCount workers, each taking the full transactionTimeout, takes up to (batchSize / workerCount) × transactionTimeout worst case. The invariant to size leaseDuration against is leaseDuration > (batchSize / workerCount) × transactionTimeout. With the shipped defaults (batchSize: 16, workerCount: 8, transactionTimeout: 60s) that worst case is 2 × 60s = 120s, comfortably inside the default leaseDuration: 5m. The manager Warns at startup, rather than refusing to start, when this does not hold: an attempt still in flight when this happens keeps its claim rather than releasing it (see the “stuck transaction” bullet below), so this is a throughput concern, not a correctness oneleaseDuration should also comfortably exceed scanInterval. A transaction still recovering past its transactionTimeout keeps its claim alive only because this instance’s own next sweep reclaims the same still-Pending row, which refreshes the claim’s lease as a side effect, and this instance holds recovery leadership across sweeps for as long as that transaction stays abandoned specifically so that reclaim keeps happening on schedule (see step 1 of the process flow above) rather than depending on this instance winning leadership again by chance. If scanInterval is not shorter than leaseDuration, that renewal can still arrive too late and the claim can lapse between this instance’s own sweeps. The manager Warns at startup when this does not holdtransactionTimeout below 10s is rejected when the configuration is loaded, since a deadline that tight is likely to abandon recoveries that were merely slow rather than genuinely stuck. Setting it to 0 disables the per-transaction deadline entirely, but that also removes the only thing that can interrupt a Recover call that ignores its context: Stop() itself can then block indefinitely on a hung call, and every later Start/Stop wedges behind it. The manager logs a Warn at startup when the timeout is disabled; it does not stop you from doing itPending and is re-claimed on every sweep, but the manager will not start a second concurrent Recover call for it while an earlier attempt is still running past its own transactionTimeout: it skips the attempt instead. Unlike earlier versions of this manager, the claim is not released when this happens: releasing it while the earlier attempt might still be running is what let a second replica claim and run a concurrent Recover (and Commit) against the same transaction. The claim stays held, kept alive by the leadership-holding and reclaim-renewal described above, until the original stuck attempt actually returns, at which point the manager finalizes it in the background, releasing the claim with the real outcome and only then clearing the in-flight guard, whenever that turns out to be, including well after the sweep (or even the process’s Stop() call) that first hit the timeout has already moved on. A recoverCtx.Done() firing because Stop() cancelled the manager, rather than because transactionTimeout actually elapsed, is not counted or alerted on as a timeout: it is an ordinary shutdown catching an attempt mid-flight, not evidence the transaction itself is stuckstuckTransactionAlertThreshold escalates the per-attempt failure log from Warn to Error once the same transaction has timed out this many times in a row, so a persistently-unresponsive peer is visible to alerting instead of blending into routine sweep failures. Past the threshold the log backs off by doubling (threshold, 2×, 4×, 8×, …) rather than firing every sweep, so a transaction stuck for hours does not flood alerting. It is purely a log-severity signal: a timeout means the status query didn’t answer in time, not that the transaction failed. Unlike a persistent NotFound, the manager does not promote it to Orphan. Set to 0 to disable the escalationThe 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.
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.