Files
flowdeck/app/services/importers/jobs.py
T
bruno 3b00cbc371 feat(import): v5.6.0 unified data import (phases 0-5)
- unified importer framework (app/services/importers/): normalized model,
  registry, common pipeline (hierarchy, attachments, collections, dedup),
  async jobs, dry-run preview, column->type mapping
- Phase 1: Obsidian, Notion, Logseq/Roam, HTML (Apple Notes/Bear/Ulysses/
  OneNote), Google Keep, generic Markdown
- Phase 2: typed CSV/TSV, Excel (openpyxl), generic JSON
- Phase 3: Word .docx (python-docx), PDF (pypdf), HTML folders
- Phase 4: Raindrop, Pocket, Readwise, Shaarli, Netscape bookmarks, .ics,
  OPML, Standard Notes, Gitea/GitHub issues (+labels/milestones)
- Phase 5: incremental re-sync (skip/update/duplicate), partial-error resume,
  forge repo files, URL web clipper (SSRF guard), batch multi-file + UI queue,
  Notion relation resolution, exportable JSON reports
- /import wizard, API /api/import/*, migration 9 (import_items, import_jobs)
- fix: property values stored by property id (correct DB view rendering)
- deps: openpyxl, beautifulsoup4, PyYAML, python-docx, pypdf
- 43 import tests; full suite 491 green; ruff clean
- bump version 5.11.5
2026-09-13 10:43:20 -04:00

151 lines
4.6 KiB
Python

"""FlowDeck — background import jobs (v5.6.0, Phase 0).
Small in-process job manager used for large uploads (vaults, zips): the upload
is parsed and persisted in a worker thread while the UI polls job status.
"""
from __future__ import annotations
import threading
import time
import traceback
import uuid
from typing import Any
from app.db import get_conn
from app.services.importers.base import (
Importer,
ImportResult,
detect_importer,
get_importer,
)
from app.services.importers.pipeline import run_import
_JOBS: dict[str, dict[str, Any]] = {}
_LOCK = threading.Lock()
def parse_upload(filename: str, data: bytes, source_id: str | None = None) -> tuple[Importer | None, ImportResult]:
"""Detect (or use) an importer and parse the upload synchronously."""
imp = get_importer(source_id) if source_id else None
if imp is None:
imp = detect_importer(filename, data)
if imp is None:
return None, ImportResult(source=source_id or "unknown", warnings=["Format non reconnu"])
return imp, imp.parse(filename, data)
def _record(job: dict, *, status: str | None = None, error: str = "",
report: dict | None = None, progress: int | None = None) -> None:
with _LOCK:
if status:
job["status"] = status
if error:
job["error"] = error
if report is not None:
job["report"] = report
if progress is not None:
job["progress"] = progress
job["updated_at"] = time.time()
_persist(job)
def _persist(job: dict) -> None:
try:
with get_conn() as conn:
conn.execute(
"INSERT INTO import_jobs (id, source, filename, status, error, report_json, created_at, updated_at) "
"VALUES (?,?,?,?,?,?,?,?) "
"ON CONFLICT(id) DO UPDATE SET status=excluded.status, error=excluded.error, "
"report_json=excluded.report_json, updated_at=excluded.updated_at",
(job["id"], job["source"], job["filename"], job["status"], job.get("error", ""),
_json(job.get("report")), job["created_at"], job["updated_at"]),
)
conn.commit()
except Exception: # noqa: BLE001 - persistence is best-effort
pass
def _json(value: Any) -> str:
import json
try:
return json.dumps(value, ensure_ascii=False)
except (TypeError, ValueError):
return "{}"
def create_job(source: str, filename: str) -> dict:
job = {
"id": uuid.uuid4().hex[:16],
"source": source,
"filename": filename,
"status": "queued",
"progress": 0,
"error": "",
"report": None,
"created_at": time.time(),
"updated_at": time.time(),
}
with _LOCK:
_JOBS[job["id"]] = job
_persist(job)
return job
def get_job(job_id: str) -> dict | None:
with _LOCK:
job = _JOBS.get(job_id)
return dict(job) if job else None
def list_jobs(limit: int = 50) -> list[dict]:
with _LOCK:
jobs = sorted(_JOBS.values(), key=lambda j: j["created_at"], reverse=True)
return [dict(j) for j in jobs[:limit]]
def start_import_job(
*,
filename: str,
data: bytes,
source_id: str | None,
workspace_id: int | None,
workspace_name: str | None,
user_login: str,
parent_page_id: int | None,
target_collection_id: int | None,
dedup: bool = True,
mapping: dict[str, str] | None = None,
mode: str | None = None,
) -> dict:
"""Create a job and run parse + persist in a background thread."""
job = create_job(source_id or "auto", filename)
_record(job, status="running", progress=5)
def worker() -> None:
try:
imp, result = parse_upload(filename, data, source_id)
if imp is None:
_record(job, status="error", error="Format non reconnu")
return
_record(job, progress=40)
report = run_import(
result,
workspace_id=workspace_id,
workspace_name=workspace_name,
user_login=user_login,
parent_page_id=parent_page_id,
target_collection_id=target_collection_id,
dedup=dedup,
mapping=mapping,
mode=mode,
)
_record(job, status="done", progress=100, report=report)
except Exception as exc: # noqa: BLE001 - surface the error to the UI
_record(job, status="error", error=f"{exc}", report={
"traceback": traceback.format_exc()[-2000:],
})
threading.Thread(target=worker, name=f"import-{job['id']}", daemon=True).start()
return job