# Background Jobs and Polling Implementation Plan

**Goal:** Return a `job_id` immediately from the three long catalogue jobs, and let a client poll one endpoint for status and a real percentage.

**Architecture:** One module owns job state in Redis. The existing job functions gain an optional `progress` callback and are otherwise untouched. The endpoints hand work to a thread and return `202`. Nothing moves to a queue.

**Spec:** `docs/specs/2026-08-11-job-polling-design.md`

## Global Constraints

- **Never commit, never push, never stage.** Every task ends with `git status --short`.
- **No assistant attribution** in any file, message, or commit.
- **Comments explain WHY, not WHAT.**
- **The `progress` callback must be optional.** Called with no callback, every job function behaves exactly as today — that is what keeps every existing test valid.
- **Percent must never go backwards and must end at 100.** A bar that jumps or reverses teaches a user to distrust it.
- **A failed job reports the exception class name only**, never the message — the existing rule for every endpoint here.
- **Losing Redis mid-job must not abort the job.** Progress ticks are best-effort. The work is spending money on model calls; the reporting channel is not worth killing it for.
- Match conventions: module-level `logger`, snake_case, double quotes, 4-space indent, `redis.asyncio` via the existing pool.
- Tests: `.venv/bin/python -m pytest`. Baselines: `tests/unit/` **650**, `tests/integration/` **152**.
- **Run the two suites in the FOREGROUND and SEPARATELY.** Never in one pytest process, never in the background.

---

## File Structure

| File | Responsibility |
|---|---|
| `app/services/infra/jobs.py` (new) | Job state in Redis: create, tick, finish, fetch, lock |
| `app/services/enrichment/job.py` (modify) | Accept `progress`, report per batch |
| `app/services/pairing/job.py` (modify) | Accept `progress`, report per phase |
| `app/services/catalog/sync.py` (modify) | Accept `progress`, report per source |
| `app/api/catalog.py` (modify) | `202` + `job_id`, `?wait=true`, and the two job endpoints |
| `app/api/endpoints.py` (modify) | Same for `/catalog/sync` and `/catalog/enrich` |
| `tests/integration/test_jobs.py` (new) | State machine, locking, heartbeat |
| `tests/unit/test_job_progress.py` (new) | Percent monotonicity, optional callback |

---

### Task 1: Job state

**Files:**
- Create: `app/services/infra/jobs.py`
- Test: `tests/integration/test_jobs.py`

**Interfaces:**
- Produces:
  - `STALE_AFTER_SECONDS = 300`, `JOB_TTL_SECONDS = 86400`
  - `async create(tenant_id, kind) -> str` — returns a job id, or raises `JobAlreadyRunning(job_id)`
  - `async tick(job_id, percent, step) -> None` — best-effort
  - `async finish(job_id, result=None, error=None) -> None` — releases the lock
  - `async get(job_id) -> dict | None` — computes `lost` from the heartbeat
  - `async recent(tenant_id, limit=20) -> list[dict]`
  - `JobAlreadyRunning` carrying `.job_id`

Store one Redis hash per job at `job:<id>`, a lock at `job:lock:<tenant>:<kind>`
holding the job id, and a per-tenant list at `job:tenant:<tenant>` for `recent`.
All three carry the TTL, so a dead process cannot hold a tenant's lock forever.

`get` derives status rather than trusting it: a job stored as `running` whose
`updated_at` is older than `STALE_AFTER_SECONDS` is reported `lost`.

- [ ] **Step 1: Write the failing test**

Create `tests/integration/test_jobs.py`:

```python
import asyncio
import time

import pytest

from app.services.infra import jobs

TENANT = "org_job_test"


@pytest.fixture(autouse=True)
async def clean():
    yield
    await jobs.clear_tenant(TENANT)


@pytest.mark.asyncio
async def test_a_new_job_starts_queued():
    job_id = await jobs.create(TENANT, "enrich")
    state = await jobs.get(job_id)
    assert state["status"] == "queued"
    assert state["percent"] == 0
    assert state["kind"] == "enrich"
    await jobs.finish(job_id, result={"ok": True})


@pytest.mark.asyncio
async def test_progress_moves_the_job_to_running():
    job_id = await jobs.create(TENANT, "enrich")
    await jobs.tick(job_id, 45, "batch 5 of 11")
    state = await jobs.get(job_id)
    assert state["status"] == "running"
    assert state["percent"] == 45
    assert state["step"] == "batch 5 of 11"
    await jobs.finish(job_id, result={})


@pytest.mark.asyncio
async def test_a_finished_job_carries_its_report():
    # The report must be the same object the endpoint used to return, so
    # nothing downstream has to learn a new shape.
    job_id = await jobs.create(TENANT, "enrich")
    await jobs.finish(job_id, result={"products": 218, "extracted": 218})
    state = await jobs.get(job_id)
    assert state["status"] == "done"
    assert state["percent"] == 100
    assert state["result"]["extracted"] == 218


@pytest.mark.asyncio
async def test_a_failed_job_carries_only_the_exception_class():
    job_id = await jobs.create(TENANT, "pair")
    await jobs.finish(job_id, error="RuntimeError")
    state = await jobs.get(job_id)
    assert state["status"] == "failed"
    assert state["error"] == "RuntimeError"


@pytest.mark.asyncio
async def test_a_second_job_of_the_same_kind_is_refused():
    # Two concurrent pair runs both rewrite the neighbour table.
    first = await jobs.create(TENANT, "pair")
    with pytest.raises(jobs.JobAlreadyRunning) as caught:
        await jobs.create(TENANT, "pair")
    assert caught.value.job_id == first
    await jobs.finish(first, result={})


@pytest.mark.asyncio
async def test_a_different_kind_may_run_concurrently():
    pair = await jobs.create(TENANT, "pair")
    enrich = await jobs.create(TENANT, "enrich")
    assert pair != enrich
    await jobs.finish(pair, result={})
    await jobs.finish(enrich, result={})


@pytest.mark.asyncio
async def test_finishing_releases_the_lock():
    first = await jobs.create(TENANT, "pair")
    await jobs.finish(first, result={})
    second = await jobs.create(TENANT, "pair")
    assert second != first
    await jobs.finish(second, result={})


@pytest.mark.asyncio
async def test_a_silent_job_is_reported_lost(monkeypatch):
    # An in-process job dies with the process. Without this a client polls
    # "running" forever and never learns.
    monkeypatch.setattr(jobs, "STALE_AFTER_SECONDS", 0)
    job_id = await jobs.create(TENANT, "pair")
    await jobs.tick(job_id, 10, "embedding")
    time.sleep(1)
    assert (await jobs.get(job_id))["status"] == "lost"
    await jobs.finish(job_id, error="Lost")


@pytest.mark.asyncio
async def test_a_finished_job_is_never_lost(monkeypatch):
    monkeypatch.setattr(jobs, "STALE_AFTER_SECONDS", 0)
    job_id = await jobs.create(TENANT, "enrich")
    await jobs.finish(job_id, result={})
    time.sleep(1)
    assert (await jobs.get(job_id))["status"] == "done"


@pytest.mark.asyncio
async def test_an_unknown_job_is_none():
    assert await jobs.get("job_does_not_exist") is None


@pytest.mark.asyncio
async def test_recent_jobs_are_listed_newest_first():
    first = await jobs.create(TENANT, "enrich")
    await jobs.finish(first, result={})
    second = await jobs.create(TENANT, "pair")
    await jobs.finish(second, result={})
    listed = [j["job_id"] for j in await jobs.recent(TENANT)]
    assert listed[:2] == [second, first]


@pytest.mark.asyncio
async def test_a_tick_on_an_unknown_job_does_not_raise():
    # Progress reporting is best-effort: it must never take down work that is
    # already spending money on model calls.
    await jobs.tick("job_gone", 50, "step")
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/integration/test_jobs.py -v`
Expected: FAIL with `ModuleNotFoundError`

- [ ] **Step 3: Write the implementation**

Create `app/services/infra/jobs.py`. Use `get_redis_client()` from
`app.services.infra.redis`, `json` for `result`, and store `percent` as an int.
`create` sets the lock with `SET NX` and raises `JobAlreadyRunning(existing_id)`
when it is already held. `tick` and `finish` update `updated_at` on every write.
`get` reads the hash, parses `result`, and substitutes `lost` when a `running`
job's `updated_at` is older than `STALE_AFTER_SECONDS`.

Add `clear_tenant(tenant_id)` for the test fixture — it deletes the tenant's job
list, its jobs, and any held locks.

Every Redis call in `tick` is wrapped so a failure logs and returns rather than
raising.

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/integration/test_jobs.py -v`
Expected: PASS, 12 tests

- [ ] **Step 5: Report changed files**

```bash
git status --short
```

---

### Task 2: Progress from the job functions

**Files:**
- Modify: `app/services/enrichment/job.py`, `app/services/pairing/job.py`, `app/services/catalog/sync.py`
- Test: `tests/unit/test_job_progress.py`

**Interfaces:**
- `enrich_tenant(tenant_id, force=False, progress=None)`
- `pair_tenant(tenant_id, progress=None)`
- `sync_source(tenant_id, source, progress=None)`

`progress` is called as `progress(percent: int, step: str)`.

**The callback is optional and defaults to None.** Called without it, each
function behaves exactly as it does today — which is what keeps every existing
test valid without modification.

- [ ] **Step 1: Write the failing test**

Create `tests/unit/test_job_progress.py`:

```python
import pytest

import app.services.enrichment.job as enrichment
import app.services.pairing.job as pairing


class Recorder:
    def __init__(self):
        self.calls = []

    def __call__(self, percent, step):
        self.calls.append((percent, step))

    @property
    def percents(self):
        return [p for p, _ in self.calls]


def _products(n):
    return [{"product_key": f"k{i}", "name": f"Product Number {i}",
             "description": "A description.", "attributes": [],
             "price_cents": 100, "taxonomy_path": [], "brand": None,
             "content_hash": f"h{i}"} for i in range(n)]


@pytest.fixture
def stub_enrichment(monkeypatch):
    monkeypatch.setattr(enrichment, "migrate_products_table", lambda t: None)
    monkeypatch.setattr(enrichment, "compute_bands", lambda t: (100, 200))
    monkeypatch.setattr(enrichment, "persist_attributes", lambda *a, **kw: None)
    monkeypatch.setattr(enrichment, "flag_unextractable", lambda *a, **kw: None)
    monkeypatch.setattr(enrichment, "catalog_stats", lambda t: (45, 45))
    monkeypatch.setattr(enrichment, "extract_batch",
                        lambda batch: {p["product_key"]: {"color": "white"}
                                       for p in batch})
    monkeypatch.setattr(enrichment, "select_products",
                        lambda t, f: _products(45))


def test_enrichment_reports_progress(stub_enrichment):
    recorder = Recorder()
    enrichment.enrich_tenant("org_x", progress=recorder)
    assert recorder.calls


def test_percent_never_goes_backwards(stub_enrichment):
    # A bar that reverses teaches a user to distrust it.
    recorder = Recorder()
    enrichment.enrich_tenant("org_x", progress=recorder)
    assert recorder.percents == sorted(recorder.percents)


def test_percent_ends_at_one_hundred(stub_enrichment):
    recorder = Recorder()
    enrichment.enrich_tenant("org_x", progress=recorder)
    assert recorder.percents[-1] == 100


def test_percent_stays_in_range(stub_enrichment):
    recorder = Recorder()
    enrichment.enrich_tenant("org_x", progress=recorder)
    assert all(0 <= p <= 100 for p in recorder.percents)


def test_every_step_is_a_readable_string(stub_enrichment):
    # The step is shown to a person; "phase 2" tells them nothing.
    recorder = Recorder()
    enrichment.enrich_tenant("org_x", progress=recorder)
    assert all(isinstance(s, str) and s.strip() for _, s in recorder.calls)


def test_enrichment_works_with_no_callback(stub_enrichment):
    # Every existing test calls it this way.
    report = enrichment.enrich_tenant("org_x")
    assert report["extracted"] == 45


def test_a_failing_callback_does_not_kill_the_job(stub_enrichment):
    # Reporting must never abort work that is spending money.
    def boom(percent, step):
        raise RuntimeError("redis is down")

    report = enrichment.enrich_tenant("org_x", progress=boom)
    assert report["extracted"] == 45


def test_pairing_reports_progress_through_its_phases(monkeypatch):
    monkeypatch.setattr(pairing, "migrate_pairing_tables", lambda t: None)
    monkeypatch.setattr(pairing, "load_eligible", lambda t: [])
    recorder = Recorder()
    pairing.pair_tenant("org_x", progress=recorder)
    assert recorder.percents == sorted(recorder.percents)
    assert recorder.percents[-1] == 100
```

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/unit/test_job_progress.py -v`
Expected: FAIL — `enrich_tenant` takes no `progress` argument.

- [ ] **Step 3: Write the implementation**

Add to each module a small helper so a broken callback cannot take down the job:

```python
def _report(progress, percent: int, step: str) -> None:
    """Progress is best-effort. The job is spending real money on model calls;
    a reporting failure must never abort it."""
    if progress is None:
        return
    try:
        progress(int(percent), step)
    except Exception:
        logger.warning("Progress callback failed", exc_info=True)
```

**Enrichment** reports once per batch, and 100 at the end:
`_report(progress, done * 100 // total, f"extracting attributes — batch {n} of {total}")`

**Pairing** reports at the phase boundaries from spec §3 — 5 after loading, 45
after embedding, 90 after scoring, 100 after writing — with a readable step each
time. Where a phase is itself looped, interpolate inside its weight rather than
jumping.

**Sync** reports per source: `f"syncing {external_ref} ({i} of {n})"`.

- [ ] **Step 4: Run test to verify it passes**

Run: `.venv/bin/python -m pytest tests/unit/test_job_progress.py -v`
Expected: PASS, 8 tests

- [ ] **Step 5: Confirm nothing regressed**

The whole point of the optional callback is that existing callers are unaffected:

```bash
.venv/bin/python -m pytest tests/unit/ -q
.venv/bin/python -m pytest tests/integration/ -q
```

Expected: no regressions against 650 / 152.

- [ ] **Step 6: Report changed files**

```bash
git status --short
```

---

### Task 3: The endpoints

**Files:**
- Modify: `app/api/catalog.py`, `app/api/endpoints.py`
- Test: `tests/integration/test_job_routes.py`

**Interfaces:**
- `POST /catalog/pair`, `/catalog/enrich`, `/catalog/sync` → `202 {"job_id", "status"}`, or the report inline with `?wait=true`
- `GET /catalog/jobs/{job_id}` → job state, `404` if unknown
- `GET /catalog/jobs?tenant_id=…` → recent jobs

- [ ] **Step 1: Write the failing test**

Create `tests/integration/test_job_routes.py` covering:

```python
def test_pair_returns_202_and_a_job_id(monkeypatch): ...
def test_the_job_can_then_be_polled(monkeypatch): ...
def test_wait_true_returns_the_report_inline(monkeypatch): ...
def test_a_second_job_of_the_same_kind_is_409_with_the_running_id(monkeypatch): ...
def test_an_unknown_job_id_is_404(): ...
def test_recent_jobs_are_listed(monkeypatch): ...
def test_a_failed_job_exposes_only_the_exception_class(monkeypatch): ...
```

Write each in full, monkeypatching `catalog.pair_tenant` so no real work runs.

- [ ] **Step 2: Run test to verify it fails**

Run: `.venv/bin/python -m pytest tests/integration/test_job_routes.py -v`
Expected: FAIL — routes return 200 with a report.

- [ ] **Step 3: Write the implementation**

Each long endpoint becomes:

```python
@router.post("/pair", status_code=202)
async def pair_endpoint(tenant_id: str = Form(...), wait: bool = Form(False)):
    if wait:
        # The old contract, kept so existing callers add one parameter rather
        # than adopting polling. The admin page uses it: its jobs are small.
        try:
            return await anyio.to_thread.run_sync(lambda: pair_tenant(tenant_id))
        except Exception as ex:
            logger.error("Pairing failed for %s", tenant_id, exc_info=True)
            raise HTTPException(status_code=502, detail=type(ex).__name__)

    try:
        job_id = await jobs.create(tenant_id, "pair")
    except jobs.JobAlreadyRunning as running:
        raise HTTPException(status_code=409, detail={
            "error": "already_running", "job_id": running.job_id})

    asyncio.create_task(_run(job_id, tenant_id, "pair", pair_tenant))
    return {"job_id": job_id, "status": "queued"}
```

with one shared runner:

```python
async def _run(job_id, tenant_id, kind, fn):
    """Runs a blocking job off the event loop, reporting progress as it goes."""
    loop = asyncio.get_running_loop()

    def progress(percent, step):
        # The job runs in a worker thread; the Redis client is async and lives
        # on the loop, so the tick has to be handed back to it.
        asyncio.run_coroutine_threadsafe(
            jobs.tick(job_id, percent, step), loop)

    try:
        result = await anyio.to_thread.run_sync(
            lambda: fn(tenant_id, progress=progress))
        await jobs.finish(job_id, result=result)
    except Exception as ex:
        logger.error("%s job failed for %s", kind, tenant_id, exc_info=True)
        await jobs.finish(job_id, error=type(ex).__name__)
```

Then the two read routes. `GET /catalog/jobs/{job_id}` 404s on `None`. Register
both **above** `/products/{product_key}`, as the existing route-ordering comment
in that file requires.

`/catalog/sync` and `/catalog/enrich` live in `app/api/endpoints.py` — give them
the same treatment, importing the shared runner rather than duplicating it.

- [ ] **Step 4: Run the tests**

Run: `.venv/bin/python -m pytest tests/integration/test_job_routes.py -v`
Expected: PASS

- [ ] **Step 5: Update the admin page**

`app/static/catalog_admin.html` calls `/catalog/pair`. Add `wait=true` to that
form post so the button keeps working unchanged.

- [ ] **Step 6: Run both suites**

Sequentially, in the foreground. Expected: no regressions against 650 / 152.

- [ ] **Step 7: Report changed files**

```bash
git status --short
```

---

### Task 4: Live verification

- [ ] **Step 1: Start a real enrichment and poll it**

```bash
JOB=$(curl -s -X POST http://localhost:8001/catalog/enrich \
  -F "tenant_id=org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f" -F "force=true" \
  | python3 -c "import json,sys; print(json.load(sys.stdin)['job_id'])")
echo "job: $JOB"
```

Then poll every 5 seconds and print `percent` and `step` until the status is
terminal. Report the sequence of percentages actually observed — it must rise
monotonically and reach 100.

- [ ] **Step 2: Confirm the lock**

While that job runs, post `/catalog/enrich` again for the same tenant. Expect
**409** carrying the running job's id. Then post `/catalog/pair` for the same
tenant and confirm it is **allowed**, since it is a different kind.

- [ ] **Step 3: Confirm the old contract still works**

```bash
curl -s -X POST http://localhost:8001/catalog/pair \
  -F "tenant_id=org_8c32bf3e-6a18-4739-9b1c-94c0cf11125f" -F "wait=true"
```

Expect the full report inline, no job id — and the admin page's re-run button
still working.

- [ ] **Step 4: Confirm a failure is reported honestly**

Point a job at a tenant whose schema does not exist, poll it, and confirm
`status: failed` with `error` carrying only the exception class name.

- [ ] **Step 5: Report**

Changed files, test counts, the observed percentage sequence, the 409, and the
failure case. Remind the user to commit.

---

## Follow-on work

- **A real queue (arq fits: async, Redis already present)** if jobs must survive a restart or run across replicas. The API does not change.
- **Cancellation**, which needs a cooperative check inside the job loops.
- **Websocket push** instead of polling — `publish_event` already exists.
- **Progress for CSV import**, currently synchronous and fast enough not to need it.
