feat(vector-sync): scan provisioned users immediately

Background vector sync discovered newly provisioned users only on the
periodic user-manager poll (VECTOR_SYNC_USER_POLL_INTERVAL, default 60s),
delaying first indexing by up to a minute. Add a ProvisionSignal doorbell
that provisioning paths ring after storing a user's app password, waking
user_manager_task to re-poll and spawn the user's scanner at once. The
periodic poll remains the backstop (covers cross-replica provisioning).

- ProvisionSignal (stable reference, wait-and-re-arm) held on
  VectorSyncState; closes the lost-wakeup window (no await between observing
  the ring and re-arming; anyio.Event stickiness covers a mid-poll ring)
- user_manager_task races its poll timeout against the doorbell + shutdown
- notify_user_provisioned() rung from the three app-password provisioning
  sites: Login Flow v2 web, MCP provisioning tool, management/BasicAuth API

Note: the pre-existing scanner_wake_event was never .set() and only wakes
existing scanners; a brand-new user has none, so the manager is what must
be nudged.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Chris Coutinho
2026-06-12 09:48:56 +02:00
co-authored by Claude Opus 4.8
parent 4688f2f95a
commit d8e3e9bc33
7 changed files with 503 additions and 2 deletions
+6
View File
@@ -407,6 +407,12 @@ async def provision_app_password(request: Request) -> JSONResponse:
username, app_password, scopes=scopes, username=nc_username
)
invalidate_scope_cache(username)
# Wake the background sync user manager so this user's scanner starts
# now instead of after the next poll. Local import avoids an app <->
# api-module import cycle.
from nextcloud_mcp_server.app import notify_user_provisioned # noqa: PLC0415
notify_user_provisioned()
_record_rate_limit_attempt(path_user_id, success=True)
logger.info("Provisioned app password for user: %s", username)
+32
View File
@@ -128,6 +128,7 @@ from nextcloud_mcp_server.server.auth_tools import register_auth_tools
from nextcloud_mcp_server.server.oauth_tools import register_oauth_tools
from nextcloud_mcp_server.vector.metrics_publisher import vector_sync_metrics_task
from nextcloud_mcp_server.vector.oauth_sync import (
ProvisionSignal,
oauth_processor_task,
user_manager_task,
)
@@ -344,6 +345,12 @@ class VectorSyncState:
task_producer: "TaskProducer | None" = None
shutdown_event: anyio.Event | None = None
scanner_wake_event: anyio.Event | None = None
# Rung by a provisioning request to wake ``user_manager_task`` immediately so
# a just-provisioned user's scanner is spawned without waiting out the
# ``VECTOR_SYNC_USER_POLL_INTERVAL`` poll. ``None`` when no user manager is
# running (single-user mode or vector sync disabled), in which case
# ``notify_user_provisioned`` is a no-op.
provision_signal: "ProvisionSignal | None" = None
# Long-lived task group used for fire-and-forget background work spawned
# from the request path (e.g. ADR-019 verify-on-read eviction). Set by the
# starlette lifespan after entering its task group; cleared on shutdown.
@@ -354,6 +361,23 @@ class VectorSyncState:
_vector_sync_state = VectorSyncState()
def notify_user_provisioned() -> None:
"""Wake the user manager to discover a just-provisioned user immediately.
Provisioning call sites invoke this after a successful app-password store so
``user_manager_task`` re-polls at once instead of waiting out
``VECTOR_SYNC_USER_POLL_INTERVAL``. The 60s poll remains the backstop, so a
missed signal (e.g. provisioning handled on a different replica than the
manager) only delays the scan, never skips it.
No-op when ``provision_signal`` is ``None`` — single-user mode or vector
sync disabled, where no user manager is running.
"""
signal = _vector_sync_state.provision_signal
if signal is not None:
signal.ring()
def _wire_vector_sync_state(
app: Starlette,
transport: IngestTransport,
@@ -416,6 +440,7 @@ def _clear_vector_sync_state() -> None:
# closed lifespan; the next startup's _wire_vector_sync_state rebinds them.
_vector_sync_state.shutdown_event = None
_vector_sync_state.scanner_wake_event = None
_vector_sync_state.provision_signal = None
# =============================================================================
@@ -2060,6 +2085,12 @@ def get_app(transport: str = "streamable-http", enabled_apps: list[str] | None =
# that choice — this path is now backend-agnostic.
shutdown_event = anyio.Event()
scanner_wake_event = anyio.Event()
# Doorbell the provisioning request path rings (via
# notify_user_provisioned) to wake the user manager immediately
# for a newly provisioned user. Held on the singleton only — both
# the manager and the signal helper reach it there.
provision_signal = ProvisionSignal()
_vector_sync_state.provision_signal = provision_signal
# User state tracking for user manager
user_states: dict = {}
@@ -2095,6 +2126,7 @@ def get_app(transport: str = "streamable-http", enabled_apps: list[str] | None =
nextcloud_host_for_sync,
user_states,
tg,
provision_signal,
)
# In-process consumer pool. ``run_consumers`` is a no-op for
@@ -118,6 +118,14 @@ async def _poll_and_store(provision_id: str) -> None:
username=result.login_name,
)
invalidate_scope_cache(effective_user_id)
# Wake the background sync user manager so this user's scanner
# starts now instead of after the next poll. Local import avoids an
# app <-> route-module import cycle.
from nextcloud_mcp_server.app import ( # noqa: PLC0415
notify_user_provisioned,
)
notify_user_provisioned()
session = _provision_sessions.get(provision_id)
if session:
session["status"] = "completed"
@@ -288,6 +288,14 @@ def register_auth_tools(mcp: FastMCP) -> None:
username=poll_result.login_name,
)
invalidate_scope_cache(user_id)
# Wake the background sync user manager so this user's scanner
# starts now instead of after the next poll. Local import avoids an
# app <-> server-module import cycle.
from nextcloud_mcp_server.app import ( # noqa: PLC0415
notify_user_provisioned,
)
notify_user_provisioned()
# Clean up the flow session
await storage.delete_login_flow_session(user_id)
+57 -2
View File
@@ -44,6 +44,40 @@ class NotProvisionedError(Exception):
pass
class ProvisionSignal:
"""One-shot doorbell that wakes ``user_manager_task`` on demand.
A provisioning request rings this (``ring()``) right after storing a new
user's app password so the manager re-polls immediately instead of waiting
out ``VECTOR_SYNC_USER_POLL_INTERVAL``. The manager parks on ``wait()``;
each ring releases exactly one wait, after which the underlying event is
re-armed for the next cycle.
The reference is stable for the life of the lifespan (stored once on the
``VectorSyncState`` singleton), so the manager never has to republish a new
event back to shared state — avoiding any ``app`` ↔ ``vector`` import cycle.
Concurrency: ``anyio.Event`` is sticky, so a ``ring()`` that lands before
``wait()`` is still observed. ``wait()`` re-arms with no ``await`` between
observing the set and swapping the event, so under cooperative scheduling a
concurrent ``ring()`` cannot slip into that window and be lost.
"""
def __init__(self) -> None:
self._event = anyio.Event()
def ring(self) -> None:
"""Signal a pending wait (or the next one to arrive)."""
self._event.set()
async def wait(self) -> None:
"""Block until the next ring, then re-arm for the following cycle."""
await self._event.wait()
# No await before the swap: a concurrent ring() cannot interleave here,
# so it lands on the fresh event and the next wait() observes it.
self._event = anyio.Event()
# Process-wide app-password storage for the BasicAuth client path.
#
# get_user_client_basic_auth is on the search hot path (Unified Search and the
@@ -408,6 +442,7 @@ async def user_manager_task(
nextcloud_host: str,
user_states: dict[str, UserSyncState],
tg: TaskGroup,
provision_signal: "ProvisionSignal",
*,
task_status: TaskStatus = anyio.TASK_STATUS_IGNORED,
) -> None:
@@ -417,6 +452,12 @@ async def user_manager_task(
- New users who have provisioned access -> start scanner
- Users who have revoked access -> cancel their scanner
Polls every ``VECTOR_SYNC_USER_POLL_INTERVAL`` seconds, but also wakes
early whenever ``provision_signal`` is rung (by a provisioning request via
``notify_user_provisioned``) so a just-provisioned user's scanner starts at
once rather than after up to a full poll interval. The poll remains the
backstop for any missed ring.
Args:
send_stream: Stream to send documents to processors
shutdown_event: Event signaling shutdown
@@ -425,6 +466,7 @@ async def user_manager_task(
nextcloud_host: Nextcloud base URL
user_states: Shared dict tracking active user scanners
tg: Task group for spawning scanner tasks
provision_signal: Doorbell rung on provisioning to force an early re-poll
task_status: Status object for signaling task readiness
"""
settings = get_settings()
@@ -491,10 +533,23 @@ async def user_manager_task(
exc_info=True,
)
# Sleep until next poll
# Sleep until the next poll tick, but wake early on shutdown or a
# provisioning signal so a just-provisioned user is discovered at once.
# Race both waits in a child task group; whichever fires first cancels
# the scope, ending the sleep. move_on_after caps it at poll_interval.
async def _wake_on(wait_fn, scope: anyio.CancelScope) -> None:
await wait_fn()
scope.cancel()
try:
with anyio.move_on_after(poll_interval):
await shutdown_event.wait()
async with anyio.create_task_group() as wake_tg:
wake_tg.start_soon(
_wake_on, shutdown_event.wait, wake_tg.cancel_scope
)
wake_tg.start_soon(
_wake_on, provision_signal.wait, wake_tg.cancel_scope
)
except anyio.get_cancelled_exc_class():
break
@@ -0,0 +1,179 @@
"""Integration test: Login Flow v2 provisioning wakes the background sync
user manager immediately.
Wires the real login-flow web-provision path end to end at the component level
(no browser / container):
_poll_and_store (provision_routes)
-> RefreshTokenStorage.store_app_password_with_scopes (real, temp DB)
-> notify_user_provisioned -> ProvisionSignal.ring
-> user_manager_task wakes, re-polls the same storage, spawns the scanner
Only the Nextcloud-facing Login Flow v2 poll is mocked (returning "completed");
everything in between is the real code. With a deliberately long poll interval,
the new user's scanner must still be spawned promptly — proving the wake came
from the provisioning signal, not the periodic poll.
This is the Login Flow v2 deployment-mode counterpart to the multi-user
BasicAuth coverage in ``test_app_password_provisioning.py``.
"""
import tempfile
import time
from pathlib import Path
from unittest.mock import AsyncMock, MagicMock, patch
import anyio
import pytest
from cryptography.fernet import Fernet
from nextcloud_mcp_server.auth.login_flow import LoginFlowPollResult
from nextcloud_mcp_server.auth.provision_routes import (
_poll_and_store,
_provision_sessions,
)
from nextcloud_mcp_server.auth.storage import RefreshTokenStorage
from nextcloud_mcp_server.vector.oauth_sync import ProvisionSignal, user_manager_task
pytestmark = pytest.mark.integration
@pytest.fixture
def encryption_key():
return Fernet.generate_key().decode()
@pytest.fixture
async def temp_storage(encryption_key):
"""Real RefreshTokenStorage backed by a temporary SQLite DB."""
with tempfile.TemporaryDirectory() as tmpdir:
db_path = Path(tmpdir) / "test_login_flow_wake.db"
storage = RefreshTokenStorage(
db_path=str(db_path), encryption_key=encryption_key
)
await storage.initialize()
yield storage
async def test_login_flow_provision_wakes_user_manager(temp_storage, mocker):
"""A completed Login Flow v2 web provision spawns the user's scanner at
once via the provision signal, well inside a long poll interval."""
# ── user_manager: long poll interval + stubbed per-user scanner ──────────
manager_settings = MagicMock()
manager_settings.vector_sync_user_poll_interval = 1000 # never fires here
mocker.patch(
"nextcloud_mcp_server.vector.oauth_sync.get_settings",
return_value=manager_settings,
)
spawned: set[str] = set()
alice_spawned = anyio.Event()
async def fake_scanner(
user_id,
cancel_scope,
send_stream,
shutdown_event,
wake_event,
nextcloud_host,
user_states,
):
spawned.add(user_id)
if user_id == "alice":
alice_spawned.set()
with cancel_scope:
await shutdown_event.wait()
user_states.pop(user_id, None)
mocker.patch(
"nextcloud_mcp_server.vector.oauth_sync._run_user_scanner_with_scope",
fake_scanner,
)
# ── wire the doorbell exactly as the lifespan does ───────────────────────
import nextcloud_mcp_server.app as app_module
provision_signal = ProvisionSignal()
mocker.patch.object(
app_module._vector_sync_state, "provision_signal", provision_signal
)
# ── mock only the Nextcloud Login Flow v2 poll ───────────────────────────
completed = LoginFlowPollResult(
status="completed",
server="https://cloud.example.com",
login_name="alice",
app_password="aaaaa-bbbbb-ccccc-ddddd-eeeee",
)
flow_client = AsyncMock()
flow_client.poll.return_value = completed
provision_settings = MagicMock()
provision_settings.nextcloud_host = "https://cloud.example.com"
provision_settings.nextcloud_public_issuer_url = None
provision_id = "login-flow-wake"
_provision_sessions[provision_id] = {
"status": "pending",
"poll_endpoint": "https://cloud.example.com/login/v2/poll",
"poll_token": "secret-token",
"user_id": "alice",
"created_at": time.time(),
"expires_at": time.time() + 1200,
}
try:
shutdown_event = anyio.Event()
user_states: dict = {}
async with anyio.create_task_group() as tg:
await tg.start(
user_manager_task,
None, # send_stream — unused by the stubbed scanner
shutdown_event,
anyio.Event(), # scanner wake_event
temp_storage,
"https://cloud.example.com",
user_states,
tg,
provision_signal,
)
# First poll: no users provisioned yet → no scanner.
await anyio.sleep(0.1)
assert not spawned
# Run the real login-flow web-provision background task. It stores
# the app password into temp_storage and rings the doorbell.
with (
patch(
"nextcloud_mcp_server.auth.provision_routes.get_settings",
return_value=provision_settings,
),
patch(
"nextcloud_mcp_server.auth.provision_routes.get_nextcloud_ssl_verify",
return_value=False,
),
patch(
"nextcloud_mcp_server.auth.provision_routes.LoginFlowV2Client",
return_value=flow_client,
),
patch(
"nextcloud_mcp_server.auth.provision_routes.get_shared_storage",
new_callable=AsyncMock,
return_value=temp_storage,
),
):
await _poll_and_store(provision_id)
# The app password was really stored …
assert "alice" in await temp_storage.get_all_app_password_user_ids()
# … and the manager woke and spawned alice's scanner far inside the
# 1000s poll interval (i.e. because of the signal, not the poll).
with anyio.fail_after(3):
await alice_spawned.wait()
assert spawned == {"alice"}
shutdown_event.set()
finally:
_provision_sessions.pop(provision_id, None)
@@ -0,0 +1,213 @@
"""Unit tests for immediate-on-provision scanner spawning.
Covers the doorbell that lets a provisioning request wake ``user_manager_task``
at once instead of waiting out ``VECTOR_SYNC_USER_POLL_INTERVAL``:
- ``ProvisionSignal`` semantics (sticky ring, wake-a-parked-wait, re-arm).
- ``user_manager_task`` re-polls early when the signal is rung, spawning a
newly provisioned user's scanner well before the next poll tick.
- ``notify_user_provisioned`` is a no-op when no manager is running.
"""
from unittest.mock import MagicMock
import anyio
import pytest
from nextcloud_mcp_server.vector.oauth_sync import ProvisionSignal, user_manager_task
pytestmark = pytest.mark.unit
# ── ProvisionSignal primitive ────────────────────────────────────────────────
async def test_provision_signal_ring_before_wait_is_observed():
"""A ring that lands before wait() is sticky and returns immediately."""
signal = ProvisionSignal()
signal.ring()
with anyio.fail_after(1):
await signal.wait()
async def test_provision_signal_wakes_parked_waiter():
"""ring() releases a wait() that is already parked."""
signal = ProvisionSignal()
woke = anyio.Event()
async def waiter():
await signal.wait()
woke.set()
async with anyio.create_task_group() as tg:
tg.start_soon(waiter)
await anyio.sleep(0.05) # let waiter park
assert not woke.is_set()
signal.ring()
with anyio.fail_after(1):
await woke.wait()
async def test_provision_signal_rearms_for_next_cycle():
"""After a ring is consumed, the next wait() blocks until the next ring."""
signal = ProvisionSignal()
signal.ring()
await signal.wait() # consumes first ring, re-arms
# Second wait must block (no pending ring) then release on the next ring.
second = anyio.Event()
async def waiter():
await signal.wait()
second.set()
async with anyio.create_task_group() as tg:
tg.start_soon(waiter)
await anyio.sleep(0.05)
assert not second.is_set() # proves the first ring did not carry over
signal.ring()
with anyio.fail_after(1):
await second.wait()
# ── user_manager_task wake-on-provision ──────────────────────────────────────
class _FakeStorage:
"""Storage stub whose provisioned-user set the test mutates between polls."""
def __init__(self, users: set[str]):
self.users = users
async def get_all_app_password_user_ids(self) -> list[str]:
return list(self.users)
async def test_user_manager_wakes_on_provision_signal(mocker):
"""Ringing the signal makes the manager re-poll and spawn the new user's
scanner well before the (long) poll interval elapses."""
# Long poll interval so any prompt spawn proves it was the signal, not poll.
settings = MagicMock()
settings.vector_sync_user_poll_interval = 1000
mocker.patch(
"nextcloud_mcp_server.vector.oauth_sync.get_settings", return_value=settings
)
spawned: set[str] = set()
spawn_events: dict[str, anyio.Event] = {
"alice": anyio.Event(),
"bob": anyio.Event(),
}
async def fake_scanner(
user_id,
cancel_scope,
send_stream,
shutdown_event,
wake_event,
nextcloud_host,
user_states,
):
spawned.add(user_id)
if user_id in spawn_events:
spawn_events[user_id].set()
# Stay alive (keeps user_states populated) until shutdown.
with cancel_scope:
await shutdown_event.wait()
user_states.pop(user_id, None)
mocker.patch(
"nextcloud_mcp_server.vector.oauth_sync._run_user_scanner_with_scope",
fake_scanner,
)
storage = _FakeStorage({"alice"})
provision_signal = ProvisionSignal()
shutdown_event = anyio.Event()
scanner_wake_event = anyio.Event()
user_states: dict = {}
async with anyio.create_task_group() as tg:
await tg.start(
user_manager_task,
None, # send_stream — unused by the stubbed scanner
shutdown_event,
scanner_wake_event,
storage,
"http://nextcloud",
user_states,
tg,
provision_signal,
)
# First poll discovers the already-provisioned user.
with anyio.fail_after(2):
await spawn_events["alice"].wait()
assert "bob" not in spawned
# Provision a new user, then ring — the manager must pick bob up fast.
storage.users.add("bob")
provision_signal.ring()
with anyio.fail_after(2): # << 1000s poll interval
await spawn_events["bob"].wait()
assert spawned == {"alice", "bob"}
shutdown_event.set()
async def test_user_manager_shutdown_still_breaks_sleep(mocker):
"""Setting shutdown wakes the manager out of its sleep promptly even with a
long poll interval (the doorbell race must not regress shutdown latency)."""
settings = MagicMock()
settings.vector_sync_user_poll_interval = 1000
mocker.patch(
"nextcloud_mcp_server.vector.oauth_sync.get_settings", return_value=settings
)
mocker.patch(
"nextcloud_mcp_server.vector.oauth_sync._run_user_scanner_with_scope",
# No users provisioned, so this is never called; keep a harmless stub.
lambda *a, **k: anyio.sleep(0),
)
storage = _FakeStorage(set())
shutdown_event = anyio.Event()
async with anyio.create_task_group() as tg:
await tg.start(
user_manager_task,
None,
shutdown_event,
anyio.Event(),
storage,
"http://nextcloud",
{},
tg,
ProvisionSignal(),
)
await anyio.sleep(0.05) # let it enter the sleep
shutdown_event.set()
# If shutdown didn't break the 1000s sleep, fail_after would trip.
with anyio.fail_after(2):
await anyio.sleep(0) # task group exit below is the real assertion
# ── notify_user_provisioned no-op guard ──────────────────────────────────────
def test_notify_user_provisioned_noop_without_manager(mocker):
"""When no manager is running, the helper must not raise."""
import nextcloud_mcp_server.app as app_module
mocker.patch.object(app_module._vector_sync_state, "provision_signal", None)
# Should be a silent no-op.
app_module.notify_user_provisioned()
def test_notify_user_provisioned_rings_when_present(mocker):
"""When a manager is running, the helper rings its signal."""
import nextcloud_mcp_server.app as app_module
signal = ProvisionSignal()
mocker.patch.object(app_module._vector_sync_state, "provision_signal", signal)
app_module.notify_user_provisioned()
assert signal._event.is_set()