2026-08-13 22:43:33 +02:00
|
|
|
"""ARQ background jobs for large import operations.
|
|
|
|
|
|
|
|
|
|
When an import file exceeds BACKGROUND_JOB_THRESHOLD rows, the import is
|
|
|
|
|
enqueued as an ARQ job instead of running synchronously. Job status is
|
|
|
|
|
tracked in Redis with the key pattern 'leocrm:import_job:{job_id}'.
|
|
|
|
|
"""
|
|
|
|
|
|
|
|
|
|
from __future__ import annotations
|
|
|
|
|
|
|
|
|
|
import json
|
|
|
|
|
import logging
|
|
|
|
|
import uuid
|
|
|
|
|
from typing import Any
|
|
|
|
|
|
|
|
|
|
from app.core.job_registry import register_job
|
|
|
|
|
from app.services import import_export_service
|
|
|
|
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
|
|
|
|
_JOB_KEY_PREFIX = "leocrm:import_job"
|
|
|
|
|
_JOB_TTL = 3600 # 1 hour
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def _job_key(job_id: str) -> str:
|
|
|
|
|
return f"{_JOB_KEY_PREFIX}:{job_id}"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def _set_job_status(
|
|
|
|
|
redis_client: Any,
|
|
|
|
|
job_id: str,
|
|
|
|
|
status: str,
|
|
|
|
|
**extra: Any,
|
|
|
|
|
) -> None:
|
|
|
|
|
"""Store job status in Redis with TTL."""
|
|
|
|
|
data = {"status": status, **extra}
|
|
|
|
|
await redis_client.set(_job_key(job_id), json.dumps(data, default=str), ex=_JOB_TTL)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def _get_job_status(redis_client: Any, job_id: str) -> dict[str, Any] | None:
|
|
|
|
|
"""Retrieve job status from Redis."""
|
|
|
|
|
raw = await redis_client.get(_job_key(job_id))
|
|
|
|
|
if raw is None:
|
|
|
|
|
return None
|
|
|
|
|
if isinstance(raw, bytes):
|
|
|
|
|
raw = raw.decode("utf-8")
|
|
|
|
|
return json.loads(raw)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def import_background_job(ctx: dict[str, Any], job_id: str, params: dict[str, Any]) -> dict[str, Any]:
|
|
|
|
|
"""ARQ job function: run import in background with status tracking.
|
|
|
|
|
|
|
|
|
|
Args:
|
|
|
|
|
ctx: ARQ worker context.
|
|
|
|
|
job_id: Unique job identifier.
|
|
|
|
|
params: Dict with keys:
|
|
|
|
|
- entity_type: 'contacts' or 'companies'
|
|
|
|
|
- csv_content: File content as string
|
|
|
|
|
- tenant_id: Tenant UUID string
|
|
|
|
|
- user_id: User UUID string
|
|
|
|
|
- field_mapping: Optional dict[str, str]
|
|
|
|
|
"""
|
|
|
|
|
from app.core.auth import get_redis
|
2026-08-16 01:17:18 +02:00
|
|
|
from app.core.db import get_worker_session_factory, set_tenant_context
|
2026-08-13 22:43:33 +02:00
|
|
|
|
|
|
|
|
redis_client = get_redis()
|
|
|
|
|
entity_type = params["entity_type"]
|
|
|
|
|
csv_content = params["csv_content"]
|
|
|
|
|
tenant_id = uuid.UUID(params["tenant_id"])
|
|
|
|
|
user_id = uuid.UUID(params["user_id"])
|
|
|
|
|
field_mapping = params.get("field_mapping")
|
|
|
|
|
|
|
|
|
|
# Set initial status
|
|
|
|
|
await _set_job_status(redis_client, job_id, "processing", progress=0, total=0)
|
|
|
|
|
|
|
|
|
|
factory = get_worker_session_factory()
|
|
|
|
|
try:
|
|
|
|
|
async with factory() as db:
|
|
|
|
|
# Set tenant context for RLS
|
|
|
|
|
await set_tenant_context(db, tenant_id)
|
|
|
|
|
|
|
|
|
|
result = await import_export_service.import_csv(
|
|
|
|
|
db,
|
|
|
|
|
tenant_id,
|
|
|
|
|
user_id,
|
|
|
|
|
csv_content,
|
|
|
|
|
entity_type=entity_type,
|
|
|
|
|
dry_run=False,
|
|
|
|
|
field_mapping=field_mapping,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
# Determine final status
|
|
|
|
|
status = result.get("status", "completed")
|
|
|
|
|
if status == "success":
|
|
|
|
|
final_status = "completed"
|
|
|
|
|
elif status == "partial_success":
|
|
|
|
|
final_status = "partial_success"
|
|
|
|
|
else:
|
|
|
|
|
final_status = "failed"
|
|
|
|
|
|
|
|
|
|
await _set_job_status(
|
|
|
|
|
redis_client,
|
|
|
|
|
job_id,
|
|
|
|
|
final_status,
|
|
|
|
|
progress=100,
|
|
|
|
|
total=result.get("total", 0),
|
|
|
|
|
succeeded=result.get("succeeded", 0),
|
|
|
|
|
failed=result.get("failed", 0),
|
|
|
|
|
error_report=result.get("error_report", {}),
|
|
|
|
|
created_count=len(result.get("created", [])),
|
|
|
|
|
)
|
|
|
|
|
return result
|
|
|
|
|
|
|
|
|
|
except Exception as exc:
|
|
|
|
|
logger.error("Import background job %s failed: %s", job_id, exc, exc_info=True)
|
|
|
|
|
await _set_job_status(
|
|
|
|
|
redis_client,
|
|
|
|
|
job_id,
|
|
|
|
|
"failed",
|
|
|
|
|
progress=0,
|
|
|
|
|
total=0,
|
|
|
|
|
error=str(exc),
|
|
|
|
|
)
|
|
|
|
|
raise
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
# Register the job so it appears in WorkerSettings.functions
|
|
|
|
|
register_job("import_background", import_background_job)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def get_import_job_status(job_id: str) -> dict[str, Any] | None:
|
|
|
|
|
"""Public helper to get job status from Redis."""
|
|
|
|
|
from app.core.auth import get_redis
|
|
|
|
|
|
|
|
|
|
redis_client = get_redis()
|
|
|
|
|
return await _get_job_status(redis_client, job_id)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def create_import_job(
|
|
|
|
|
entity_type: str,
|
|
|
|
|
csv_content: str,
|
|
|
|
|
tenant_id: uuid.UUID,
|
|
|
|
|
user_id: uuid.UUID,
|
|
|
|
|
field_mapping: dict[str, str] | None = None,
|
|
|
|
|
) -> str:
|
|
|
|
|
"""Enqueue an import as ARQ background job and return job_id."""
|
|
|
|
|
from app.core.jobs import enqueue_job
|
|
|
|
|
|
|
|
|
|
job_id = str(uuid.uuid4())
|
|
|
|
|
params = {
|
|
|
|
|
"entity_type": entity_type,
|
|
|
|
|
"csv_content": csv_content,
|
|
|
|
|
"tenant_id": str(tenant_id),
|
|
|
|
|
"user_id": str(user_id),
|
|
|
|
|
"field_mapping": field_mapping,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
# Pre-set status as pending
|
|
|
|
|
from app.core.auth import get_redis
|
|
|
|
|
|
|
|
|
|
redis_client = get_redis()
|
|
|
|
|
await _set_job_status(redis_client, job_id, "pending", progress=0, total=0)
|
|
|
|
|
|
|
|
|
|
await enqueue_job("import_background", job_id, params)
|
|
|
|
|
return job_id
|