Files
poimen-memory/tasks/M3.5.4-query-federation.md
Story Crater Bot 603c2b681f feat: M3.5.8 complete - all endpoints, rate limiting, and deployment (253 tests)
Changes:
- Queue cleanup: Deleted 17 poisoned CI runs from database
- Code: All M3.5 endpoints implemented and tested
- Tests: 253 total, all passing
- Deployment: K8s manifests and ArgoCD configured
- CI: Forgejo Actions dispatcher issue (image not built yet)

Next: Manual image build or CI dispatcher fix
2026-08-23 17:19:42 -07:00

157 lines
6.4 KiB
Markdown
Raw Permalink 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.4 — Federation: single query across multiple projects
| 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.3 (query endpoint exists) |
## Goal
Extend query endpoint to support multi-project search. When `project` param is omitted, a single query searches all projects concurrently, deduplicates results, and merges scores.
## Design
**Single-project query (no change):**
```
GET /memory/query?query=Kong+body&project=poimen
→ results from poimen only
```
**Multi-project query (federation):**
```
GET /memory/query?query=Kong+body
→ results from all projects, merged by rerank_score
```
Response is the same shape; add optional `_federation` metadata:
```json
{
"query": "Kong body",
"projects_searched": ["poimen", "agent-rust"],
"results": [...],
"latency_ms": 512,
"notes": "Searched 2 projects in parallel; 3 results after dedup"
}
```
## Behavior
**Deduplication:** Same `sha256` across projects is impossible (sha256 includes project name in provenance), so no dedup needed. If two projects happen to have identical text:
- Treat as separate nodes (different projects, different provenance)
- Return both in results (may both rank high)
- Ensure test coverage catches this edge case
**Concurrency:** Query all projects in parallel using `tokio::join_all()` or `futures::stream`:
```rust
let futures: Vec<_> = projects.iter()
.map(|proj| query_single_project(query_text, proj, limit))
.collect();
let results: Vec<_> = futures::future::join_all(futures).await;
```
**Merging:** After all projects return, merge result vectors:
- Collect all results from all projects into one vec
- Re-sort by `rerank_score` descending (global order)
- Take top `limit` (e.g., if poimen returns [a,b,c] and agent-rust returns [d,e], merge gives [a,b,c,d,e] → sorted globally → top 5 might be [b,d,a,c,e])
**Timeout:** Per-project timeout is min(timeout_seconds / projects.len(), 2s). If one project is slow, others complete faster and we still return results from fast projects after global timeout.
- E.g., timeout=10s, 2 projects → 5s per project
- If project-a completes in 3s, project-b in 8s, and global timeout is 10s:
- Return results from both (8s < 10s)
- If project-a completes in 3s, project-b in 12s, and global timeout is 10s:
- After 10s, cancel project-b, return results from project-a only
- Note in response: `"warnings": ["project 'agent-rust' timed out"]`
## Steps
1. Parse `project` param:
- If provided, single-project path (M3.5.3 unchanged)
- If omitted, multi-project path
2. List all known projects (from queries YAML):
```rust
let projects = load_standing_queries()?.projects();
```
3. Spawn concurrent query tasks:
```rust
let futures: Vec<_> = projects.into_iter()
.map(|proj| {
let params = params.clone();
params.project = Some(proj);
query_handler_impl(&params, store, llm)
})
.collect();
```
4. Race with timeout:
```rust
let deadline = Instant::now() + Duration::from_secs(timeout_seconds);
let results = match tokio::time::timeout_at(deadline, futures::future::join_all(futures)).await {
Ok(vec) => vec.into_iter().flatten().collect(), // flatten per-project results
Err(_) => { /* partial results + warning */ }
};
```
5. Merge and sort:
```rust
results.sort_by(|a, b| b.rerank_score.partial_cmp(&a.rerank_score).unwrap());
results.truncate(limit);
```
6. Assemble response with federation metadata:
```rust
let response = QueryResponse {
projects_searched: /* only projects that completed */,
warnings: /* projects that timed out */,
results,
latency_ms: start.elapsed().as_millis() as u64,
..
};
```
## Acceptance
- Single project specified: no federation, same result as M3.5.3
- No project specified: all projects queried
- Results merged and globally sorted by rerank_score
- Partial results returned if one project times out
## Verify
**Harness:** Integration tests with two projects in test pgvector DB.
**Integration test** — `tests/it_query_federation.rs`:
1. `a1_single_project_no_federation` — GET /query?project=poimen returns single-project results only.
2. `a2_multi_project_searches_all` — GET /query (no project) with >1 project in DB returns results from all.
3. `a3_global_sort_order` — two projects return results, merge sorts by rerank_score globally (not per-project).
4. `a4_federation_metadata_present` — response includes `projects_searched` array with all completed projects.
5. `a5_partial_results_on_timeout` — slow one project (mock 10s delay), set timeout_seconds=2, GET /query returns results from fast project only with warning.
6. `a6_limit_applied_after_merge` — project-a returns [a1,a2,a3], project-b returns [b1,b2,b3], limit=4, global merge returns 4 results (not 6).
7. `a7_no_project_filter_in_response` — response.project_filter is null or omitted (unlike single-project which sets it).
8. `a8_concurrent_execution` — spy on timing: timestamp project-a query start, project-b query start, both should be ~simultaneous (not sequential).
**Command:** `cargo test -p mem-cli query_federation`
**False pass:**
- Testing only with one project in DB. Federation always "works" if there is nothing to federate.
- Timeout never exercised. Mock a slow project and assert results are partial.
- Per-project sorting instead of global sort. Results look reasonable but violate the contract (should be global top-k).
- Concurrency not verified. Queries can be sequential (slow) and still return correct results; only timing proves concurrency.
## Traps
- Timeout math: if you do `timeout_per_project = timeout_total / num_projects`, a project that completes in 1s uses the full allocated time before returning. Should be `remaining_time = deadline - now()`.
- Partial results: if project-a returns 5 results and project-b times out, you have 5 results but may have wanted 10 (limit=10). Document whether partial results truncate or stay over-limit.
- Clone overhead: cloning `QueryParams` for each project is small; cloning a large result vec is not. Use references/Arc where possible.
- Flatten after join_all: `join_all` returns `Vec<Result>`, must flatten errors (either as partial results or early exit).
---
Background: [DESIGN.md § Distributed API Layer](../DESIGN.md#distributed-api-layer-homelab-frontend)