Skip to content

Feature/async connection keepalive heartbeat - #735

Open
abdullahilateefat03-boop wants to merge 2 commits into
StellarFlow-Network:mainfrom
abdullahilateefat03-boop:feature/async-connection-keepalive-heartbeat
Open

Feature/async connection keepalive heartbeat#735
abdullahilateefat03-boop wants to merge 2 commits into
StellarFlow-Network:mainfrom
abdullahilateefat03-boop:feature/async-connection-keepalive-heartbeat

Conversation

@abdullahilateefat03-boop

Copy link
Copy Markdown
Contributor

closes #609

Problem

Async database connections (asyncpg-backed ingestion sinks) had no keep-alive coverage. During quiet low-volume market windows, PgBouncer and serverless Postgres instances drop idle TCP connections after their timeout threshold. When the next price record arrives, the ingestion sink stalls waiting for a fresh TCP handshake and re-authentication before the write can proceed.

The existing ConnectionKeepAlive class solves this for synchronous DB-API 2.0 connections using a background thread. Async connections expose an awaitable execute() coroutine, not a cursor() — so the sync class cannot be used against them.

The test suite in
test_connection_keepalive.py
was already importing AsyncConnectionKeepAlive and failing at import time because the class did not exist.

Solution

Added AsyncConnectionKeepAlive to
connection.py
How it works:

start() schedules a background asyncio.Task via asyncio.ensure_future(). Calling start() on an already-running instance is a no-op.
ping() awaits connection.execute("SELECT 1;") using the same HEARTBEAT_QUERY constant as the sync variant. Any exception is caught, logged as a warning, and swallowed — a single transient drop never kills the loop.
_run() is the internal coroutine loop: asyncio.sleep(interval) → ping() → repeat, until cancelled via asyncio.CancelledError.
stop() cancels the task and awaits its exit — returns promptly regardless of how large the interval is.
is_running property returns True while the underlying task is alive and not done.
Default interval is 30.0 seconds (DEFAULT_PING_INTERVAL), consistent with the sync class.
Constructor validates inputs: None connection and non-positive interval both raise ValueError.
Compatible with any async connection exposing an awaitable execute(query: str) method — asyncpg, aiopg, databases, or test fakes.
Files changed

connection.py

… overwrite sweep

- Remove duplicate SigningError class definition that silently shadowed
  the canonical exception type defined near the other exception classes.

- Fix SecureKeyHandle._do_wipe to call _munlock_buffer (respecting the
  _locked flag set by _mlock_buffer) instead of the legacy _unlock_memory
  helper that bypassed lock-state tracking. Unlock now always runs AFTER
  the zero-wipe, preventing the OS from evicting live key material to swap
  during the cleanup window. Adds audit log entry on scope close, matching
  SecureSessionCredentials and SecureVariableWrapper behavior.

- Fix _sign_internal transient copy wipe: replace _wipe_bytes_view (which
  zeroed only a copy-of-a-copy via from_buffer_copy) with _wipe_bytes_object
  (which uses ctypes.memset on id(obj)+offset to overwrite the CPython bytes
  object internal data buffer in-place). This closes the window where the
  transient key_bytes bytes object remained readable in heap memory after
  the signing routine returned.

Addresses: memory-dump extraction vulnerability where active private signing
key material lingered in process memory after transaction signing completed.
…n sinks

The existing ConnectionKeepAlive (threading-based) had no async counterpart,
leaving asyncpg-backed multi-corridor ingestion sinks without keep-alive
coverage. During quiet low-volume market windows these async connections were
dropped by PgBouncer / serverless Postgres, forcing incoming price records to
stall on fresh TCP handshakes.

Changes
-------
- Add AsyncConnectionKeepAlive to src/database/connection.py alongside the
  existing synchronous ConnectionKeepAlive class.

Implementation details
----------------------
- start() schedules a background asyncio.Task via asyncio.ensure_future().
  Calling start() on an already-running instance is a no-op (idempotent).
- ping() is a coroutine that awaits connection.execute(SELECT 1;) — the same
  HEARTBEAT_QUERY constant used by the sync variant. Failures are caught,
  logged, and swallowed so a single transient drop never terminates the loop.
- _run() is the internal coroutine loop: asyncio.sleep(interval) then ping(),
  repeating until cancelled. Handles asyncio.CancelledError cleanly.
- stop() cancels the task and awaits its exit — returns promptly even with a
  large interval (no blocking wait for the full interval to elapse).
- is_running property returns True while the task is alive and not done.
- Accepts any async connection exposing an awaitable execute(query) method
  (asyncpg, aiopg, databases, or test fakes).
- Default interval is 30.0 seconds (DEFAULT_PING_INTERVAL), matching the
  sync class and the 30-second requirement.
- Validates constructor args: None connection and non-positive interval both
  raise ValueError.

All cases are covered by the pre-existing test suite in
tests/test_connection_keepalive.py which was already importing
AsyncConnectionKeepAlive and failing at import time.
@drips-wave

drips-wave Bot commented Jul 29, 2026

Copy link
Copy Markdown

@abdullahilateefat03-boop Great news! 🎉 Based on an automated assessment of this PR, the linked Wave issue(s) no longer count against your application limits.

You can now already apply to more issues while waiting for a review of this PR. Keep up the great work! 🚀

Learn more about application limits

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

🗄️ Database-Pipeline | Asynchronous Connection Ping Anchors for Multi-Corridor Ingestion Sinks

1 participant