Files
poimen-memory/tasks/M3.5.2-ingest-endpoint.md
T

144 lines
6.7 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# M3.5.2 — POST /ingest endpoint: async queue, idempotency, job polling
| Field | Value |
|---|---|
| Phase | M3.5 — Distributed API Layer |
| Size | M — 13 days |
| Status | ✅ Done |
| Flags | — |
| Spec | inlined below |
| Blocks | M3.5.8 |
| Depends | M3.5.1, M1.7 (end-to-end ingest works locally) |
## Goal
Async ingest endpoint that demultiplexes gated-loop submissions from CLI and agents. Idempotent by batch content hash (`ingest_id`). Prevent duplicate L0 evidence in the log. Enrich records with git context (file, commit, blame) if repo.git available.
## Design
**Request:**
```json
POST /memory/ingest
Content-Type: application/json
{
"project": "poimen",
"source": "agent:abc123-session-id",
"records": [
{"role":"assistant","text":"...","timestamp":"2026-08-20T...","source_position":0},
...
],
"ingest_id": "sha256(all_record_texts)",
"git_repo_path": "/path/to/repo/.git",
"git_head": "abc123def789"
}
```
**Git enrichment (optional):** If `git_repo_path` and `git_head` provided:
- Walk repo blame for timestamps matching record timestamps
- Correlate evidence text with recent commits touching files
- Populate `git_context` on each L0 node (file, line, commit, author)
**Response (accepted):**
```
HTTP 202 Accepted
{
"job_id": "ingest-<uuid>",
"ingest_id": "sha256(...)",
"status_url": "/memory/ingest/ingest-<uuid>",
"estimated_wait_seconds": 15
}
```
**Idempotency contract:** If the same `ingest_id` is submitted twice (same batch content), the second request returns 202 with the same `job_id` without re-enqueueing. If `ingest_id` differs but project overlaps, both are enqueued separately (ordering is per-project FIFO after dedup).
**Job status (polling):**
```
GET /memory/ingest/ingest-<job-id>
→ 200 {
"job_id": "...",
"ingest_id": "...",
"project": "poimen",
"status": "running|completed|failed",
"chunks_seen": 42,
"chunks_used": 7,
"error": null,
"created_at": "2026-08-20T...",
"completed_at": null
}
```
## Steps
1. Ingest queue — choose **local in-memory (BTreeMap keyed by ingest_id) or Redis**. For M3.5, start in-memory; scaling to Redis is P2-deferred.
- Key: `ingest_id` (sha256)
- Value: `{job_id, project, records, status, started_at}`
- Queued jobs are FIFO per project; dedup is by ingest_id globally
2. `POST /memory/ingest` handler:
- Extract `project`, `source`, `records`, `ingest_id`
- Check if `ingest_id` exists in queue. If yes, return 202 with existing `job_id` (no duplicate enqueue).
- If new, generate `job_id = format!("ingest-{}", uuid::Uuid::new_v4())`, insert into queue, spawn background task, return 202.
- Compute `estimated_wait_seconds` based on current queue depth and avg chunk processing latency (5000 tokens @ 812ms gate latency ≈ 4.2s per chunk).
3. Background task (tokio::spawn):
- Dequeue from project queue (FIFO per project)
- **Git enrichment (if git_repo_path provided):**
- Open repo.git with `git2::Repository`
- For each record, find blame line by timestamp + closest file match (via commit log)
- Populate `git_context: {file, line, commit_sha, commit_msg, author, author_date}`
- Call the M1.7 `mem::ingest()` function with enriched records
- Update status to `completed` with `chunks_seen` and `chunks_used` from the log
- On error, update status to `failed` with error message
4. `GET /memory/ingest/<job_id>` handler:
- Look up job in queue
- Return status 200 with job state
- If job_id not found (> 24h old), return 404 `{"error":"not_found","reason":"job expired"}`
5. Validation:
- `ingest_id` must be a hex string of length 64 (sha256); malformed → 400
- `project` must be a known project (loaded from queries/); unknown → 400
- `records` array must not be empty; empty → 400
## Acceptance
- POST returns 202 with a job_id
- Same ingest_id resubmitted returns same job_id (idempotent)
- Job status is pollable
- Two different ingest_ids for the same project are both queued (not deduplicated by project)
- Background task completes without blocking the request
- Malformed request (bad ingest_id, unknown project) returns 400
## Verify
**Harness:** Integration tests + one manual queue inspection.
**Integration test**`tests/it_ingest_endpoint.rs`:
1. `a1_ingest_accepted` — POST /ingest with valid payload returns 202 and body contains `job_id` field.
2. `a2_ingest_id_is_idempotent` — POST twice with same `ingest_id`, same `project` — both return 202 with identical `job_id`.
3. `a3_status_polling_works` — POST /ingest, GET /ingest/<job_id> immediately returns `status: "running"` or `status: "completed"`.
4. `a4_different_ingest_ids_both_queued` — POST /ingest (id_a), POST /ingest (id_b), GET status of both — both in queue.
5. `a5_bad_ingest_id_returns_400` — POST with `ingest_id: "xyz"` (not 64 hex chars) returns 400.
6. `a6_unknown_project_returns_400` — POST with `project: "nonexistent"` returns 400.
7. `a7_async_task_runs` — POST /ingest with a small test batch, poll /ingest/<job_id> repeatedly, verify status transitions from `running` to `completed`.
8. `a8_empty_records_returns_400` — POST with `records: []` returns 400.
**Manual verification:**
- Run the server, ingest two batches with different ingest_ids for the same project, verify they are queued in order by checking JSONL log — both should be present after ingest completes, in the order submitted.
**Command:** `cargo test -p mem-cli ingest_endpoint`
**False pass:**
- Testing with one project only. Multi-project FIFO ordering is the hard part; a single project always looks correct.
- Job status never actually transitions from `running` to `completed`. A mock status endpoint can always return `running` and pass the test if the test only polls once.
- Idempotency checked for `ingest_id` but not for `project` — two requests with same `ingest_id` but different `project` must be treated as different (they are).
- Latency estimate never validated. Estimated wait can be any number; test should assert it is > 0 and < 1 hour.
## Traps
- Using a simple Vec for the queue. FIFO per project requires either a per-project queue map or a global queue with project filtering. Per-project is cheaper.
- Job expiry: in-memory queue will grow unbounded if jobs are never pruned. Set an eviction policy (e.g., remove jobs older than 24h on every ingest request).
- Tokio task panic in the background task. Spawn with `.spawn()` which detaches on panic; use a panic hook or `.spawn_blocking()` with error handling.
- Reusing the M1.7 function directly without error wrapping. If it panics (log write fails, db timeout), the background task crashes and the job status never updates. Wrap in a Result type and catch panics.
---
Background: [DESIGN.md § Distributed API Layer](../DESIGN.md#distributed-api-layer-homelab-frontend)