fix: supervise the Wiki.js change listener so it survives a database restart
The listener opened one asyncpg connection, called add_listener, and set running = True. Nothing watched that connection afterwards. When it dropped, the subscription was gone for good while running still reported True, so the service stayed healthy in every way anything could observe and silently stopped indexing page edits. Recovery needed a manual container restart. That happened on 2026-08-08 when postgres-shared was redeployed. The sibling settings_client survived the same event because it uses asyncpg.create_pool, which replaces dead connections; a bare LISTEN connection has no such recovery. A supervisor task now waits on asyncpg's termination callback and reconnects with bounded exponential backoff, 1s doubling to a 60s cap. It retries forever rather than giving up after N attempts: a database under maintenance does come back, and a listener that stopped trying would reproduce exactly the silent deafness this exists to prevent. The termination listener is re-registered on every new connection because asyncpg clears its listener list as soon as it fires them, so a one-time registration survives exactly one drop. running is now derived from the connection rather than assigned, and stop() sets a flag the termination callback and supervisor both check so a deliberate shutdown cannot race into a reconnect. NOTIFY is fire-and-forget, so events emitted during an outage are lost and cannot be replayed. The reconnect logs the gap and names POST /maintenance/integrity-check rather than reporting a clean recovery. Reconciling automatically is left out on purpose: deriving the tenant for a changed page is subtle here, and getting it wrong writes into the wrong user's namespace. Verified against the real database by terminating the listener's backend with pg_terminate_backend. Old code: running=True with is_closed()=True, dead forever. New code: reconnects on its own onto a new server pid. The same probe was run against both implementations so the check is known to discriminate. One existing test mocked the connection with a bare AsyncMock, which models asyncpg's synchronous is_closed() as a coroutine — always truthy, so the connection read as closed once running started deriving from it. Corrected. Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
+15
-1
@@ -123,10 +123,19 @@ async def health(settings: Settings = Depends(get_settings)) -> HealthResponse:
|
||||
# Check service connectivity
|
||||
service_health = await check_service_health()
|
||||
|
||||
# Overall status is healthy if at least Neo4j and Qdrant are up
|
||||
# Overall status is healthy if at least Neo4j and Qdrant are up.
|
||||
#
|
||||
# The Wiki.js change listener is reported below but deliberately excluded
|
||||
# from this decision. It supervises and reconnects itself, and a database
|
||||
# restart would otherwise flip the container unhealthy for the duration of
|
||||
# an outage it is already recovering from. It is reported so the state is
|
||||
# observable at all — previously nothing anywhere exposed it, which is how a
|
||||
# dead listener went unnoticed while this endpoint answered "healthy".
|
||||
all_healthy = service_health.get("neo4j", False) and service_health.get("qdrant", False)
|
||||
overall_status = "healthy" if all_healthy else "degraded"
|
||||
|
||||
wiki_listener = getattr(app.state, "wiki_listener", None)
|
||||
|
||||
return HealthResponse(
|
||||
status=overall_status,
|
||||
app_name=settings.app_name,
|
||||
@@ -152,6 +161,11 @@ async def health(settings: Settings = Depends(get_settings)) -> HealthResponse:
|
||||
"url": settings.ollama_url,
|
||||
"model": settings.ollama_llm_model,
|
||||
"healthy": service_health.get("ollama", False)
|
||||
},
|
||||
"wiki_listener": {
|
||||
"subscribed": bool(wiki_listener and wiki_listener.running),
|
||||
"reconnects": getattr(wiki_listener, "reconnects", 0),
|
||||
"last_gap_seconds": getattr(wiki_listener, "last_gap_seconds", None)
|
||||
}
|
||||
}
|
||||
)
|
||||
|
||||
@@ -27,21 +27,59 @@ class WikiChangeListener:
|
||||
NOTIFY events on INSERT/UPDATE/DELETE to the pages table.
|
||||
"""
|
||||
|
||||
CHANNEL = 'wiki_page_changes'
|
||||
|
||||
# Bounded exponential backoff between reconnect attempts.
|
||||
BACKOFF_INITIAL_SECONDS = 1.0
|
||||
BACKOFF_MAX_SECONDS = 60.0
|
||||
|
||||
def __init__(self):
|
||||
self.settings = get_settings()
|
||||
self.connection: Optional[asyncpg.Connection] = None
|
||||
self.running = False
|
||||
|
||||
# Loop prevention: Track recently processed pages
|
||||
# Key: page_id, Value: timestamp of last processing
|
||||
self._recent_notifications = {}
|
||||
self._debounce_seconds = self.settings.wikijs_change_listener_debounce_seconds
|
||||
|
||||
async def start(self):
|
||||
"""Start listening to database changes."""
|
||||
logger.info("Starting Wiki.js database change listener")
|
||||
# Supervision state. A single LISTEN connection does not heal itself the
|
||||
# way an asyncpg pool does, so the drop has to be detected and repaired
|
||||
# explicitly — see _supervise().
|
||||
self._stopping = False
|
||||
self._disconnected = asyncio.Event()
|
||||
self._supervisor_task: Optional[asyncio.Task] = None
|
||||
self._disconnected_at: Optional[datetime] = None
|
||||
self.reconnects = 0
|
||||
self.last_gap_seconds: Optional[float] = None
|
||||
|
||||
# Connect to Wiki.js PostgreSQL database
|
||||
@property
|
||||
def running(self) -> bool:
|
||||
"""
|
||||
Whether a live subscription actually exists.
|
||||
|
||||
Derived rather than assigned. The previous implementation set a flag once
|
||||
in start() and never revisited it, so after the connection dropped the
|
||||
listener reported itself as running while being deaf to every event.
|
||||
"""
|
||||
return (
|
||||
not self._stopping
|
||||
and self.connection is not None
|
||||
and not self.connection.is_closed()
|
||||
)
|
||||
|
||||
async def start(self):
|
||||
"""Start listening to database changes, and keep listening."""
|
||||
logger.info("Starting Wiki.js database change listener")
|
||||
self._stopping = False
|
||||
self._disconnected.clear()
|
||||
|
||||
await self._connect()
|
||||
|
||||
self._supervisor_task = asyncio.create_task(self._supervise())
|
||||
logger.info("Listening for Wiki.js page changes via PostgreSQL NOTIFY")
|
||||
|
||||
async def _connect(self):
|
||||
"""Open a connection and subscribe. Raises if the database is unreachable."""
|
||||
self.connection = await asyncpg.connect(
|
||||
host=self.settings.wikijs_db_host,
|
||||
port=self.settings.wikijs_db_port,
|
||||
@@ -50,18 +88,136 @@ class WikiChangeListener:
|
||||
database=self.settings.wikijs_db_name
|
||||
)
|
||||
|
||||
# Listen to the wiki_page_changes channel
|
||||
await self.connection.add_listener('wiki_page_changes', self._handle_notification)
|
||||
await self.connection.add_listener(self.CHANNEL, self._handle_notification)
|
||||
|
||||
self.running = True
|
||||
logger.info("Listening for Wiki.js page changes via PostgreSQL NOTIFY")
|
||||
# Must be re-registered on every connection: asyncpg clears its
|
||||
# termination listeners as soon as it fires them, so this is one-shot.
|
||||
self.connection.add_termination_listener(self._on_connection_lost)
|
||||
|
||||
def _on_connection_lost(self, connection):
|
||||
"""
|
||||
Called by asyncpg when the connection terminates.
|
||||
|
||||
Dispatched through loop.call_soon, so it must stay synchronous — the work
|
||||
of reconnecting belongs to _supervise(), which this only wakes.
|
||||
"""
|
||||
if self._stopping:
|
||||
return
|
||||
|
||||
self._disconnected_at = datetime.now()
|
||||
logger.error(
|
||||
"Wiki.js change listener lost its database connection — "
|
||||
"page changes are NOT being processed until it reconnects"
|
||||
)
|
||||
self._disconnected.set()
|
||||
|
||||
async def _supervise(self, max_iterations: Optional[int] = None) -> int:
|
||||
"""
|
||||
Reconnect whenever the subscription drops.
|
||||
|
||||
Args:
|
||||
max_iterations: Stop after N reconnect cycles (None = run forever;
|
||||
used by tests)
|
||||
|
||||
Returns:
|
||||
Number of completed reconnect cycles
|
||||
"""
|
||||
iterations = 0
|
||||
while max_iterations is None or iterations < max_iterations:
|
||||
await self._disconnected.wait()
|
||||
if self._stopping:
|
||||
break
|
||||
|
||||
self._disconnected.clear()
|
||||
await self._reconnect_with_backoff()
|
||||
iterations += 1
|
||||
|
||||
return iterations
|
||||
|
||||
async def _reconnect_with_backoff(self, max_attempts: Optional[int] = None) -> bool:
|
||||
"""
|
||||
Re-establish the subscription, backing off between failures.
|
||||
|
||||
Keeps trying indefinitely by default: a database that is down for
|
||||
maintenance will come back, and giving up would recreate exactly the
|
||||
silent-deafness this supervision exists to prevent.
|
||||
"""
|
||||
delay = self.BACKOFF_INITIAL_SECONDS
|
||||
attempts = 0
|
||||
|
||||
while not self._stopping and (max_attempts is None or attempts < max_attempts):
|
||||
attempts += 1
|
||||
await self._close_connection()
|
||||
|
||||
try:
|
||||
await self._connect()
|
||||
except Exception as e:
|
||||
logger.warning(
|
||||
f"Wiki.js change listener reconnect attempt {attempts} failed: {e}; "
|
||||
f"retrying in {delay:.0f}s"
|
||||
)
|
||||
await asyncio.sleep(delay)
|
||||
delay = min(delay * 2, self.BACKOFF_MAX_SECONDS)
|
||||
continue
|
||||
|
||||
self.reconnects += 1
|
||||
gap = None
|
||||
if self._disconnected_at is not None:
|
||||
gap = (datetime.now() - self._disconnected_at).total_seconds()
|
||||
self.last_gap_seconds = gap
|
||||
self._disconnected_at = None
|
||||
|
||||
# NOTIFY is fire-and-forget: anything emitted while we were gone was
|
||||
# delivered to nobody and cannot be replayed. Say so, and say what
|
||||
# closes the gap, rather than reporting a clean recovery.
|
||||
outage = f" after {gap:.0f}s" if gap is not None else ""
|
||||
logger.warning(
|
||||
f"Wiki.js change listener reconnected{outage} "
|
||||
f"(reconnect #{self.reconnects}). NOTIFY events emitted during the "
|
||||
f"outage were lost and cannot be replayed — run "
|
||||
f"POST /maintenance/integrity-check to reconcile pages that "
|
||||
f"changed while the listener was down."
|
||||
)
|
||||
return True
|
||||
|
||||
return False
|
||||
|
||||
async def _close_connection(self):
|
||||
"""Drop the current connection, tolerating one that is already dead."""
|
||||
if not self.connection:
|
||||
return
|
||||
|
||||
try:
|
||||
if not self.connection.is_closed():
|
||||
await self.connection.remove_listener(
|
||||
self.CHANNEL, self._handle_notification
|
||||
)
|
||||
await self.connection.close()
|
||||
except Exception as e:
|
||||
# A terminated connection raises on both calls; that is expected here.
|
||||
logger.debug(f"Error closing Wiki.js listener connection: {e}")
|
||||
finally:
|
||||
self.connection = None
|
||||
|
||||
async def stop(self):
|
||||
"""Stop listening and close connection."""
|
||||
if self.connection:
|
||||
await self.connection.remove_listener('wiki_page_changes', self._handle_notification)
|
||||
await self.connection.close()
|
||||
self.running = False
|
||||
self._stopping = True
|
||||
|
||||
# Wake the supervisor so it observes _stopping and exits rather than
|
||||
# racing us to reconnect the connection we are about to close.
|
||||
self._disconnected.set()
|
||||
|
||||
if self._supervisor_task:
|
||||
self._supervisor_task.cancel()
|
||||
try:
|
||||
await self._supervisor_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
except Exception as e:
|
||||
logger.debug(f"Wiki.js listener supervisor ended with: {e}")
|
||||
self._supervisor_task = None
|
||||
|
||||
await self._close_connection()
|
||||
logger.info("Stopped Wiki.js change listener")
|
||||
|
||||
async def _handle_notification(self, connection, pid, channel, payload):
|
||||
|
||||
Reference in New Issue
Block a user