- Add tasks/cloud_tasks.py: durable refresh_cloud_folder task with 3-retry bounded backoff (30s/90s/270s +jitter), credential decryption in worker only - Register tasks.cloud_tasks.* on documents queue in celery_app.py - Add stale-while-revalidate staleness trigger in browse.py (5-min threshold) - Add 4 Task 3 tests: idempotency, cached-row retention on failure, task structure, no-byte-download contract; add background-refresh scheduling integration test - Bump backend version 0.1.4 → 0.1.5, frontend package.json 0.1.4 → 0.1.5 - Update AGENTS.md with Phase 12 Plan 02 state and new shared module map entries - Update README with connection-ID browse API table and v0.1.5
213 lines
8.5 KiB
Python
213 lines
8.5 KiB
Python
"""
|
|
Celery tasks for cloud metadata refresh in DocuVault — Phase 12.
|
|
|
|
refresh_cloud_folder — called via .delay(user_id, connection_id, parent_ref)
|
|
by the browse endpoint when durable cached rows are stale.
|
|
|
|
The task is a plain sync def (Celery workers have no asyncio event loop); it
|
|
bridges into the async service layer via asyncio.run().
|
|
|
|
Retry harness (D-13/D-14 — bounded transient retries only):
|
|
CRITICAL: self.retry() MUST be raised from the outer sync task, NOT inside
|
|
asyncio.run(). _TransientProviderError is a sentinel raised by _run() to signal
|
|
a retryable provider/network failure. The outer task catches it and calls
|
|
self.retry(countdown=...) in the sync layer.
|
|
|
|
Terminal errors (auth failure, invalid_grant, scope error):
|
|
- Written as CloudFolderState warning state with controlled reason/remedy
|
|
- NOT retried — retrying auth failures wastes network calls and may trigger
|
|
provider rate limits or OAuth token revocation
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import random
|
|
|
|
from celery.exceptions import MaxRetriesExceededError
|
|
|
|
from celery_app import celery_app
|
|
|
|
|
|
# ── Retry sentinels ───────────────────────────────────────────────────────────
|
|
|
|
class _TransientProviderError(Exception):
|
|
"""Raised by _run() to signal a retryable network/provider failure.
|
|
|
|
Escapes asyncio.run() and is caught by refresh_cloud_folder's sync layer,
|
|
which calls self.retry() in the Celery context (not inside asyncio.run).
|
|
"""
|
|
|
|
|
|
class _TerminalProviderError(Exception):
|
|
"""Raised by _run() for non-retryable auth/scope errors.
|
|
|
|
The task writes a warning state and does not retry.
|
|
"""
|
|
|
|
|
|
# ── Async worker ─────────────────────────────────────────────────────────────
|
|
|
|
async def _run(user_id: str, connection_id: str, parent_ref: str | None) -> dict:
|
|
"""Open a fresh DB session, revalidate owner, refresh and reconcile.
|
|
|
|
D-18/CACHE-01: never downloads file bytes, never alters quota.
|
|
D-13/D-14: on success → fresh state; on transient failure → raise sentinel
|
|
for Celery retry; on auth/scope failure → warning state + return.
|
|
"""
|
|
import uuid as _uuid
|
|
|
|
from db.session import AsyncSessionLocal
|
|
from services.cloud_items import (
|
|
ConnectionNotFound,
|
|
get_or_create_folder_state,
|
|
reconcile_cloud_listing,
|
|
resolve_owned_connection,
|
|
update_folder_state,
|
|
)
|
|
from storage.cloud_backend_factory import build_cloud_resource_adapter
|
|
from storage.cloud_utils import decrypt_credentials
|
|
from config import settings
|
|
|
|
master_key = settings.cloud_creds_key.encode()
|
|
|
|
async with AsyncSessionLocal() as session:
|
|
# Revalidate ownership — worker never trusts broker payload alone
|
|
try:
|
|
conn = await resolve_owned_connection(
|
|
session,
|
|
connection_id=connection_id,
|
|
user_id=user_id,
|
|
)
|
|
except ConnectionNotFound:
|
|
# Connection deleted or user changed — silently drop, do not retry
|
|
return {"status": "skipped", "reason": "connection_not_found"}
|
|
|
|
# Mark refreshing so the browse endpoint can show spinner state
|
|
await update_folder_state(
|
|
session,
|
|
user_id=user_id,
|
|
connection_id=connection_id,
|
|
parent_ref=parent_ref or "",
|
|
refresh_state="refreshing",
|
|
)
|
|
await session.commit()
|
|
|
|
# Decrypt credentials inside the worker — never put them in broker payload
|
|
try:
|
|
credentials = decrypt_credentials(master_key, user_id, conn.credentials_enc)
|
|
except Exception as exc:
|
|
# Decryption failure is terminal — key mismatch, corruption, etc.
|
|
await update_folder_state(
|
|
session,
|
|
user_id=user_id,
|
|
connection_id=connection_id,
|
|
parent_ref=parent_ref or "",
|
|
refresh_state="warning",
|
|
error_code="credential_error",
|
|
error_message="Credential decryption failed. Re-connect the account.",
|
|
)
|
|
await session.commit()
|
|
raise _TerminalProviderError(f"Credential decryption failed: {exc}") from exc
|
|
|
|
# Build adapter and fetch listing
|
|
try:
|
|
adapter = build_cloud_resource_adapter(conn.provider, credentials)
|
|
conn_uuid = conn.id if isinstance(conn.id, _uuid.UUID) else _uuid.UUID(str(conn.id))
|
|
user_uuid = _uuid.UUID(str(user_id))
|
|
listing = await adapter.list_folder(
|
|
conn_uuid,
|
|
user_uuid,
|
|
parent_ref=parent_ref,
|
|
)
|
|
except Exception as exc:
|
|
err_str = str(exc).lower()
|
|
# Detect auth/scope errors — these are terminal
|
|
if any(kw in err_str for kw in ("invalid_grant", "unauthorized", "401", "403", "scope", "revoked")):
|
|
await update_folder_state(
|
|
session,
|
|
user_id=user_id,
|
|
connection_id=connection_id,
|
|
parent_ref=parent_ref or "",
|
|
refresh_state="warning",
|
|
error_code="auth_error",
|
|
error_message="Authentication failed. Re-connect the account.",
|
|
)
|
|
await session.commit()
|
|
raise _TerminalProviderError(f"Auth error: {exc}") from exc
|
|
# Transient failure — let Celery retry
|
|
raise _TransientProviderError(f"Provider error: {exc}") from exc
|
|
|
|
# Reconcile listing into durable cloud_items rows
|
|
await reconcile_cloud_listing(
|
|
session,
|
|
user_id=user_id,
|
|
connection_id=connection_id,
|
|
parent_ref=parent_ref,
|
|
listing=listing,
|
|
)
|
|
|
|
await update_folder_state(
|
|
session,
|
|
user_id=user_id,
|
|
connection_id=connection_id,
|
|
parent_ref=parent_ref or "",
|
|
refresh_state="fresh",
|
|
)
|
|
await session.commit()
|
|
|
|
return {"status": "ok", "items_fetched": len(listing.items)}
|
|
|
|
|
|
# ── Celery task ───────────────────────────────────────────────────────────────
|
|
|
|
@celery_app.task(
|
|
bind=True,
|
|
max_retries=3,
|
|
name="tasks.cloud_tasks.refresh_cloud_folder",
|
|
serializer="json",
|
|
acks_late=True,
|
|
reject_on_worker_lost=True,
|
|
)
|
|
def refresh_cloud_folder(self, user_id: str, connection_id: str, parent_ref: str | None) -> dict:
|
|
"""Refresh cloud metadata for (user_id, connection_id, parent_ref) idempotently.
|
|
|
|
D-13/D-14: Bounded transient retries with increasing countdown + jitter.
|
|
Terminal auth/scope failures mark warning state without retry.
|
|
D-18/CACHE-01: Never downloads file bytes; never alters quota.
|
|
"""
|
|
try:
|
|
return asyncio.run(_run(user_id, connection_id, parent_ref))
|
|
except _TerminalProviderError:
|
|
# Already wrote warning state in _run — do not retry
|
|
return {"status": "terminal_error"}
|
|
except _TransientProviderError as exc:
|
|
try:
|
|
# Bounded exponential backoff with jitter: 30s, 90s, 270s (± 10s)
|
|
jitter = random.randint(-10, 10)
|
|
countdown = (30 * (3 ** self.request.retries)) + jitter
|
|
raise self.retry(exc=exc, countdown=countdown)
|
|
except MaxRetriesExceededError:
|
|
# Write permanent warning state after all retries exhausted
|
|
asyncio.run(_write_final_warning(user_id, connection_id, parent_ref))
|
|
return {"status": "max_retries_exceeded"}
|
|
|
|
|
|
async def _write_final_warning(
|
|
user_id: str, connection_id: str, parent_ref: str | None
|
|
) -> None:
|
|
"""Write a terminal warning state after all Celery retries are exhausted."""
|
|
from db.session import AsyncSessionLocal
|
|
from services.cloud_items import update_folder_state
|
|
|
|
async with AsyncSessionLocal() as session:
|
|
await update_folder_state(
|
|
session,
|
|
user_id=user_id,
|
|
connection_id=connection_id,
|
|
parent_ref=parent_ref or "",
|
|
refresh_state="warning",
|
|
error_code="refresh_failed",
|
|
error_message="Background refresh failed after 3 attempts. Will retry on next browse.",
|
|
)
|
|
await session.commit()
|