Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,10 @@ All notable changes to Sheaf are documented here. The format is based on [Keep a

## [Unreleased]

### Fixed

- **Imports start immediately again, instead of sitting for a few seconds first.** The runner keeps a database connection open to be told the instant an import is queued, so a file you upload starts processing in milliseconds rather than waiting for the next poll. That connection was opening a database transaction it never closed, and Postgres only hands a notification to a connection that is between transactions, so a minute after the server started the notifications quietly stopped arriving. Nothing was lost or stuck, because the timed poll behind it always picked the job up anyway, but every import paid a short wait it was not supposed to. Two things made it invisible: the poll is a deliberate safety net, so the symptom was "slightly slower" rather than "broken", and the test covering it finished before the transaction ever opened. It now also no longer holds the oldest transaction on the database, which was hiding genuinely stuck transactions from monitoring. This is the same fault that was fixed in the leader election earlier; these were the only two places with the pattern.

### Added

- **An opt-in extended metrics tier, starting with active accounts by client version.** Self-hosters with metrics enabled only. `METRICS_EXTENDED=true` turns on a second tier of metrics for the questions whose answers need more series or short-lived per-account state, off by default so no instance inherits the scrape cost or the data-handling posture without asking. Every metric in it is named `sheaf_ext_*` and every Redis key `sheaf:ext:*`, so one regex routes or drops the lot in a pipeline; per-account state is folded under a day-salted token that cannot be joined across days and nothing outlives 48 hours. The first metric is `sheaf_ext_active_accounts_by_version{client_family, version}`: distinct accounts active today per client family and `major.minor` client version, the number that decides how long an API compatibility shim has to stay. The tier's own settings share the `METRICS_EXTENDED_` prefix; the first, `METRICS_EXTENDED_VERSION_PAIRS_PER_DAY` (default 64), bounds how many distinct versions a day may hold before new ones fold into `other`. Documented in `docs/METRICS.md`.
Expand Down
30 changes: 28 additions & 2 deletions sheaf/services/import_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -332,6 +332,13 @@ async def run_import_tick(db: AsyncSession) -> dict:
return {"items_processed": 1, "details": f"complete: {job.id}"}


# How often the LISTEN connection probes that it is still alive. A module
# constant rather than a literal so the regression test can drive a heartbeat
# without waiting a minute: the bug this guards against only appears AFTER the
# first probe, so a test that never heartbeats cannot see it.
_LISTEN_HEARTBEAT_SECONDS = 60


async def _listen_for_enqueues(wake: asyncio.Event) -> None:
"""Hold a LISTEN connection and set `wake` on every enqueue NOTIFY.

Expand All @@ -347,9 +354,28 @@ async def _listen_for_enqueues(wake: asyncio.Event) -> None:

from sheaf.database import engine

# AUTOCOMMIT, because this connection lives for the whole process. In the
# default mode SQLAlchemy autobegins a transaction on the first statement
# and the heartbeat below never ends it, so the session sat "idle in
# transaction" forever. That is the same bug the leader election had (see
# sheaf/services/leader.py), with the same two costs: it pinned
# pg_stat_activity's max transaction age, blinding any long-transaction
# alert, and it kept resetting idle_in_transaction_session_timeout.
#
# Here it also broke the feature outright. Postgres delivers NOTIFY to a
# listening backend only when that backend is idle, so a listener parked
# inside a never-ending transaction stops receiving notifications
# altogether - which meant that from 60 seconds after startup this
# accelerator silently did nothing and every enqueued import waited out
# the poll interval instead. The LISTEN registration itself belongs to the
# SESSION, not the transaction, so nothing else about the listener
# changes: it still dies with the connection, which is what the reconnect
# path relies on.
listen_engine = engine.execution_options(isolation_level="AUTOCOMMIT")

while True:
try:
async with engine.connect() as conn:
async with listen_engine.connect() as conn:
raw = await conn.get_raw_connection()
driver = raw.driver_connection # asyncpg connection
await driver.add_listener(
Expand All @@ -362,7 +388,7 @@ async def _listen_for_enqueues(wake: asyncio.Event) -> None:
while True:
# Liveness probe: raises when the connection has died,
# dropping us to the reconnect path.
await asyncio.sleep(60)
await asyncio.sleep(_LISTEN_HEARTBEAT_SECONDS)
await conn.execute(text("SELECT 1"))
except asyncio.CancelledError:
raise
Expand Down
49 changes: 49 additions & 0 deletions tests/test_job_runner_rework.py
Original file line number Diff line number Diff line change
Expand Up @@ -583,6 +583,55 @@ async def run() -> None:
asyncio.run(run())


def test_listener_still_wakes_after_a_heartbeat():
"""A NOTIFY that arrives AFTER the liveness probe has run still wakes the
listener.

The regression this pins: the listener connection used to run in the
default isolation level, so the probe's `SELECT 1` autobegan a
transaction that nothing ever ended. Postgres only delivers NOTIFY to a
backend that is idle, so a listener parked inside an open transaction
stops receiving notifications entirely - the accelerator silently died 60
seconds after startup and every import fell back to the poll interval.

The pre-existing test above cannot catch this, because it notifies within
a few seconds and the first heartbeat has not run yet. This one drives a
heartbeat first, which is why the interval is a patchable constant.
"""
from unittest.mock import patch

import sheaf.database as database_module
from sheaf.services import import_runner as runner_module
from sheaf.services.import_runner import _listen_for_enqueues

async def run() -> None:
engine = _test_engine()
wake = asyncio.Event()
try:
with (
patch.object(database_module, "engine", engine),
patch.object(runner_module, "_LISTEN_HEARTBEAT_SECONDS", 0.2),
):
listener = asyncio.create_task(_listen_for_enqueues(wake))
# Long enough for several heartbeats to have run, so the
# connection is well past the point the old code broke at.
await asyncio.sleep(1.5)
assert not wake.is_set()

async with engine.connect() as conn:
await conn.execute(text("NOTIFY sheaf_import_enqueued"))
await conn.commit()

await asyncio.wait_for(wake.wait(), timeout=5)
listener.cancel()
with pytest.raises(asyncio.CancelledError):
await listener
finally:
await engine.dispose()

asyncio.run(run())


# ---------------------------------------------------------------------------
# Operator kill switch for data-deleting jobs

Expand Down
Loading