Compare commits

...

11 Commits

Author SHA1 Message Date
b25c15fefa docs(claude): record plan 002 phase 2 in the project status
Why:
- The status section described plan 002 as Phase 1 only. The /v1/points read
  routes, the query service, and the points scopes had shipped but were listed
  under "not built yet", so a fresh session would start from a wrong map of the
  codebase.

Changes:
- Record the five read routes, the points:read gating, and the 404-not-403
  mapping.
- Call out the three route-level rules that fail quietly when broken: /count and
  /search must precede /{point_id}, file_id is required on the listing, and the
  search query is Persian-folded before matching.
- Point at tests/support/point_contract.py as the place new repository behaviour
  belongs, since it runs against both the fake and real Qdrant.
- Narrow the "not built yet" list to Phases 3-6.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 16:04:59 +03:30
ba7921dd4e test(points): cover the read paths and the query service
Why:
- The isolation and pagination guarantees are the ones that fail silently, so
  they need tests that would actually notice.

Changes:
- API tests over real Postgres and Qdrant together, because the invariant worth
  testing spans both: the tenant Postgres derived is the only one Qdrant is ever
  queried with.
- Pagination holds when a point is inserted behind the cursor mid-listing --
  the defect an offset cursor would have.
- A route-order guard, since /{point_id} declared first turns count into a 422
  and nothing else in the suite would catch it.
- Unit tests for the query service against the fake, including the Arabic to
  Persian letterform fold and raising rather than returning None.

Impact:
- Suite goes to 325 passed, 3 skipped.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 15:09:31 +03:30
802bae429d docs(claude): record plan 002 phase 1 in the project status 2026-08-22 15:09:18 +03:30
43932b6562 docs(adr): record query normalization and file-scoped pagination
Why:
- Both are contracts callers depend on, not implementation details, and neither
  was written down. This repo treats the ADR as the source of truth rather than
  letting code diverge from it silently.

Changes:
- Record that the search query is folded the same way ingested content was, and
  why the alternative fails in the worst available way: an exact-looking query
  returning nothing, with no error and nothing in the logs to distinguish it
  from a genuine miss.
- Record that listing requires file_id and paginates by order_id value, and why
  an offset cursor repeats an already-served row under a concurrent insert.
- State that results carry no relevance score and no ranked order, so callers
  cannot read array position as relevance.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 15:09:07 +03:30
3b9434faf4 feat(points): add the /v1/points read paths
Why:
- Ingestion writes points in bulk but nothing could read one back. Plan 002
  Phase 2 opens the read surface an admin frontend needs.

Changes:
- GET /v1/points/{point_id}, /v1/points?file_id=..., /v1/points/count,
  /v1/points/search, and /v1/files/{file_id}/points, all under points:read --
  the scope follows the data, so an upload key does not become a way to read
  every chunk of every file.
- The keyword query is Persian-normalized before matching, because ingestion
  letter-folds content at ingest and an unfolded query would return an empty
  result set silently rather than an error.
- file_id is required on the listing: the cursor is an order_id value and
  order_id is only unique within one file.
- PointNotFoundError maps to 404, never 403, so a cross-tenant point id is
  indistinguishable from a nonexistent one.
- Route order is load-bearing: /count and /search precede /{point_id}, or
  "count" is parsed as a UUID and fails 422.

Impact:
- Requires the content/is_active/chunk_index payload indexes, so a deployed
  environment needs qdrant_bootstrap re-run before search works.
- Keyword search returns no relevance score and no ranked order; callers must
  not read array position as relevance.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 15:08:54 +03:30
5251990444 feat(tenant): grant points scopes on provisioned keys
Why:
- points:read and points:write are named in ADR-0009 but were absent from
  DEFAULT_SCOPES, so a provisioned key could not reach the read paths that
  follow. The runbook documented only files:write and domains:write, which
  understated what a default key can now do.

Changes:
- Add points:read and points:write to DEFAULT_SCOPES.
- Document the full scope table in the runbook, calling out that points:read
  grants the text of every chunk of every file -- so an upload-only key gets
  files:write alone.

Impact:
- Keys issued before this change keep their existing scopes; provisioning does
  not backfill. Reissue or widen an existing key explicitly.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 15:08:39 +03:30
062fdd7ac1 refactor(test): share the Qdrant point-seeding helper
Why:
- The helper turning a SeedSpec into an upsertable ChunkPoint lived inside one
  integration test. A second suite now needs to seed real Qdrant the same way,
  and a copy would let the two drift.

Changes:
- Move `_chunk_point` into tests/support/point_contract.py as `chunk_point_for`,
  beside the `build_point` read model it derives its payload from.

Impact:
- Pure move. No behaviour change; enables the API suite that follows.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 15:08:25 +03:30
923ac8e5d6 test(points): hold the fake and the Qdrant adapter to one shared contract
Why:
- Two parallel test files let a fake drift more permissive than the store it
  stands in for, so unit tests stay green while production diverges. Plan 002
  Phase 1's exit criterion is precisely that the two agree.

Changes:
- One scenario suite in tests/support/point_contract.py, run against
  FakePointRepository (unit) and QdrantPointRepository (integration). A
  divergence fails one of the two runs rather than hiding.
- The fake models the behaviours services branch on: the implied is_active read
  filter, value-based cursor pagination, and a stale version guard that matches
  nothing rather than raising -- the no-op Qdrant's filtered set_payload actually
  has, and the reason a service must read back to know its write landed.
- Patched points are re-validated rather than model_copy'd, so the fake holds a
  datetime where a read from real Qdrant returns one.
- The seeded corpus gives each tenant its own file: point IDs derive from
  file_id plus chunk_index alone, so two tenants in one file would collide on a
  single ID and the fixture would assert an impossible state.

Impact:
- 15 scenarios pass against both implementations.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 13:09:50 +03:30
4da30f9983 feat(points): add the point read/edit port and its Qdrant adapter
Why:
- PointStorage is deliberately the two bulk operations ingestion performs. Reads,
  single-point edits, and keyword search have a different caller, a different
  failure vocabulary, and a different tenant-filter obligation, so they get their
  own port rather than accreting onto the ingestion one.

Changes:
- tenant_id is a required keyword argument on every port method, making a
  forgotten tenant filter a type error rather than a review question.
- Reads go through scroll with a HasIdCondition, not retrieve: retrieve takes no
  filter and would push the tenant check into Python after Qdrant already
  answered -- the shape ADR-0002's isolation rule exists to prevent.
- Ordered listing paginates by order_id value, not offset. Qdrant returns no page
  offset under order_by, and an offset cursor skips or repeats rows when a
  concurrent insert shifts positions underneath the reader.
- Point.from_payload takes a Mapping, not a dict: dict is invariant in its value
  type, so the SDK's concrete vector union is not a dict[str, object].
- Request schemas forbid extra keys and omit server-owned fields, so a client
  sending tenant_id or version gets 422 rather than having it silently ignored.

Impact:
- No route uses this yet; the /v1/points surface is Phase 2.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 13:09:37 +03:30
5e935e5895 feat(qdrant): index content, is_active, and chunk_index on the chunks collection
Why:
- ADR-0002's keyword search needs a full-text index on content, which
  collection.py deliberately deferred to plan 002. is_active and chunk_index
  were unindexed while ingestion was the only reader; every /v1/points read path
  filters on them.

Changes:
- content gets a TEXT index with the multilingual tokenizer, which segments
  Persian correctly where the word tokenizer mishandles ZWNJ-joined compounds.
  No stemmer or stopword list: content is already letter-folded by
  normalize_persian_text at ingest, and the ranked Farsi lexical path is the
  benchmarked BM25 sparse vector, not this index.
- Tests assert content is TEXT rather than KEYWORD -- a keyword index would only
  match an entire chunk verbatim, which never happens and fails silently.
- Adds a test that a missing index is added to an already-live collection.

Impact:
- Requires re-running `python -m src.cli.qdrant_bootstrap`. Payload indexes are
  additive, so no collection rebuild and no re-embedding.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 13:09:26 +03:30
ac3810182d docs(adr): record the re-ingestion rule and close plan 002's open decisions
Why:
- ADR-0002 already answered three of the four questions plan 002 listed as
  "decisions needed"; the fourth -- what happens to a manually edited point when
  its file is re-uploaded -- was left for a Phase 6 test to force. Deciding it in
  code rather than in the ADR would invert this repo's rule.

Changes:
- ADR-0002 gains "Re-ingestion versus manual edits": the new file wins,
  surviving points are overwritten with an incremented version, absent points
  are flagged inactive rather than removed, and manually created points sit past
  the ingested chunk_index range so the existing sweep covers them.
- Plan 002's stale decisions section becomes a pointer table; its audit scope is
  pinned to both ADR-0009 tables.

Impact:
- Clobbered edits are recoverable from point_audit_events, not from Qdrant: the
  deterministic point ID cannot hold both versions.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-22 13:09:15 +03:30
24 changed files with 2350 additions and 49 deletions

View File

@@ -34,9 +34,56 @@ e2e suite in the default pytest run (`tests/e2e/test_ingestion_slice.py`:
duplicate upload, retry after failure, tenant isolation, capacity, timeout,
parse and Qdrant failure), and the one Compose-based test — `scripts/smoke.sh`
driving `tests/e2e/test_compose_smoke.py` against a real uvicorn process, which
skips itself unless `SMOKE_BASE_URL` is set. Not built yet: `/v1/points` CRUD
and keyword search (plan 002), and `src/agent/`. That maps to plan 001 Phases
1-6 done.
skips itself unless `SMOKE_BASE_URL` is set. That maps to plan 001 Phases 1-6
done.
Plan 002 (`/v1/points` CRUD and keyword search) is **Phases 1-2 done**. Phase 1
landed the `PointRepository` port (`src/application/ports/point_repository.py`)
with its `Point` read model (`src/application/points/point.py`), the Qdrant adapter
(`src/infrastructure/qdrant/point_repository.py`), request/response schemas
(`src/api/schemas/points.py`), and lifespan wiring. This port is **separate from
`PointStorage`**, which stays exactly the two bulk operations ingestion
performs — reads, single-point edits, and keyword search have a different caller
and a different tenant-filter obligation, so do not accrete them onto the
ingestion port. `tenant_id` is a required keyword argument on every
`PointRepository` method by design; keep it that way, because it is what turns a
forgotten tenant filter into a type error. The `chunks` collection also gained
full-text `content`, `is_active`, and `chunk_index` payload indexes, so a
deployed environment needs `qdrant_bootstrap` re-run (indexes are additive — no
rebuild, no re-embedding).
Two adapter mechanics there are load-bearing and easy to "simplify" into bugs:
reads go through `scroll` with a `HasIdCondition` rather than `retrieve` (which
takes no filter, and would move the tenant check to *after* Qdrant answered),
and ordered listing paginates by `order_id` value rather than offset (Qdrant
returns no page offset under `order_by`, and an offset cursor skips or repeats
rows under a concurrent insert).
Phase 2 added the **read routes**: `GET /v1/points/{point_id}`,
`GET /v1/points?file_id=...`, `GET /v1/points/count`, `GET /v1/points/search`,
and `GET /v1/files/{file_id}/points`, over `src/application/points/queries.py`
(`src/api/routers/points.py`). All are gated on `points:read`, which — with
`points:write` — is now in `DEFAULT_SCOPES`; `GET /v1/files/{file_id}/points`
uses `points:read` rather than `files:write`, so the scope follows the data
rather than the URL prefix. `PointNotFoundError` maps to `404` in
`src/api/errors.py`, never `403`. Three route-level rules are load-bearing:
`/count` and `/search` are declared **before** `/{point_id}` (FastAPI matches in
declaration order, so reordering them makes `/v1/points/count` a `422`),
`file_id` is **required** on the listing (the cursor is an `order_id` value and
`order_id` is unique only within one file), and `search_points` folds the query
with `normalize_persian_text` before matching, because ingestion letter-folds
content and an unfolded Arabic-keyboard query would return an empty result set
silently rather than erroring (ADR-0002).
Also worth knowing before touching the points tests: `tests/support/point_contract.py`
holds **one** scenario suite run against both `FakePointRepository` (unit) and
`QdrantPointRepository` (integration), so new repository behaviour belongs there
rather than in one of the two runners — that is what keeps the fake from drifting
more permissive than the real store.
Not built yet: plan 002 Phases 3-6 — soft delete with neighbour relinking,
create/replace/patch, reorder and batch, and the
`api_request_logs`/`point_audit_events` tables — and `src/agent/`.
Architecture decisions live in `docs/adr/` (18 ADRs plus the 0000 template;
0001–0004 are `Accepted` — 0004 amended by 0018; 0014 is `Superseded by 0017`;

View File

@@ -75,6 +75,51 @@ retrieval used by the AI agent in ADR-0003; the two "search" concepts serve
different callers (a human/admin managing chunks vs. an agent retrieving
context) and should not be conflated in the API or in future discussion.
Two properties follow from the index being a *filter*: results carry no
relevance score, and their order is unspecified. The API therefore returns
neither a score field nor a ranked list, and callers must not read the array
order as relevance. A caller that wants ranking wants ADR-0003's path.
#### The query is normalized the way ingested content was
`normalize_persian_text` (ADR-0018) folds Arabic letterforms to their Persian
equivalents — U+064A to U+06CC, U+0643 to U+06A9 — on every text block before
chunking, so stored `content` is uniformly Persian-formed. A query string is
not chunk content and never passes through that path, so a term typed on an
Arabic keyboard reaches the index as a different codepoint sequence than the
document it should match.
The service therefore applies the same folding to the query before matching.
Without it the endpoint fails in the worst available way: an exact-looking
query returns an empty result set, with no error, no warning, and nothing in
the logs to distinguish "no such term" from "the term is spelled with the
other yeh". Note this is a *query-side* transformation only — it changes what
is compared, never what is stored.
This does not extend to stemming or synonyms. Qdrant's full-text index offers
neither, and adding a Farsi analyzer here would duplicate the benchmarked BM25
sparse pipeline (ADR-0005) in a code path that is not benchmarked against
anything.
#### Listing is scoped to one file, and paginates by `order_id`
`GET /points?file_id=...` requires `file_id` rather than treating it as one
optional filter among several, and its pagination cursor is an `order_id`
value rather than an offset. Both follow from `order_id` being per-file:
- A cursor is only meaningful against a totally ordered key. `order_id` orders
points within one file and says nothing across files, so an unscoped listing
has no stable sort to paginate along.
- An offset cursor is wrong even within one file. Insert, reorder, and delete
all shift positions, so a page-two request issued after a concurrent insert
ahead of the cursor would repeat a row already returned — silently. Ranging
on `order_id > cursor` is unaffected: the reader has passed that value, and a
point inserted behind it was already served.
The second point depends on `order_id` being unique within a file, which the
gap-exhaustion rule below preserves by rejecting a reorder whose computed gap
would collapse onto a neighbour value.
### Delete is soft by default
`DELETE /points/{point_id}` and `DELETE /points?file_id=...` set
@@ -105,6 +150,41 @@ Qdrant's `update_filter`, giving an optimistic-concurrency-style guard
against races between a concurrent ingestion re-run (ADR-0001) and a manual
edit through this API.
### Re-ingestion versus manual edits
A file can be re-uploaded after someone has hand-edited one of its points
through this API. **The newly ingested file wins.** Ingestion is authoritative
for the content of the file it ingested; a manual edit is a correction that
survives only until the source document is replaced.
Concretely:
- A point that still exists in the new version (same `file_id` +
`chunk_index`, hence the same deterministic point ID) is **overwritten in
place**. Ingestion performs a read-check-write so `version` is incremented
from whatever the manual edit left it at, rather than reset to `1`.
- A point from the previous ingestion that is **absent** from the new version
is flagged `is_active: false` with `deleted_at` set. It is never removed
from Qdrant — the soft-delete rule above applies to re-ingestion exactly as
it applies to `DELETE`.
- A manually created point (`POST /points`) is assigned a `chunk_index` past
the ingested range, so the same sweep deactivates it on the next upload of
its file. This is the intended consequence of "the new file wins", not an
accident of the sweep's bounds.
Because the point ID is derived from the immutable `chunk_index`, an
overwritten point cannot hold both the manual edit and the new file's content.
The clobbered content is therefore recorded in `point_audit_events`
(ADR-0009) as a `reingest_overwrite` operation carrying `before_version`, so
the edit is recoverable from the audit trail even though it is no longer a
live point.
Rejected alternative: preserving manual edits by having ingestion skip points
with `version > 1`. It breaks the guarantee that a successful upload leaves
Qdrant matching the uploaded document, and it needs a second, separate rule
for edited points that no longer exist in the new version — two divergent
notions of authority over one file.
### Re-embedding on content edit
`PUT /points/{point_id}` can change `content`, which leaves the stored

View File

@@ -15,8 +15,9 @@ are; this document defines order, scope, and verification criteria.
## Prerequisite
Plan 001 must be complete through **Phase 5** before Phase 3 of this plan
starts. Specifically this plan depends on: the `chunks` collection and its
Plan 001 is complete through Phase 6, so this prerequisite is satisfied. It
required plan 001 through **Phase 5** before Phase 3 of this plan starts.
Specifically this plan depends on: the `chunks` collection and its
payload indexes actually existing, API-key authentication and `AuthContext`
tenant derivation, the application-lifetime Qdrant client from the FastAPI
lifespan, and the request-lifetime `AsyncSession` wiring. Phases 1–2 below
@@ -71,7 +72,9 @@ project owner accepts them, and update the ADR rather than diverging silently.
(bulk soft delete of a file's points), from ADR-0008.
- Soft delete as the default for every delete path, with neighbor relinking.
- Optimistic concurrency on every mutating path via the `version` payload field.
- Audit rows in Postgres for mutating operations.
- Audit rows in Postgres for mutating operations: both ADR-0009 tables,
`api_request_logs` (one row per API call, written from the request middleware)
and `point_audit_events` with the real `api_request_log_id` foreign key.
- Automated tests for tenant isolation, pointer integrity, concurrency
conflicts, and pagination.
@@ -113,40 +116,22 @@ project owner accepts them, and update the ADR rather than diverging silently.
9. Routers contain no Qdrant SDK calls and no filter construction. The Qdrant
client is injected from the lifespan (ADR-0012).
## Decisions needed before the affected phase
## Decisions resolved before implementation
### Re-embedding on content edit (blocks Phase 4)
An earlier revision of this plan listed three open decisions here. All are now
settled, and one further question this plan deferred to a Phase 6 test has been
settled too. They are recorded in the ADRs — these lines are a pointer, not a
second source of truth.
`PUT /v1/points/{point_id}` can change `content`. The stored vectors then no
longer match the text. Three options, in order of preference:
| Question | Resolution | Recorded in |
|---|---|---|
| Re-embedding on content edit | Re-embed inline, reusing ingestion's ports and bounds and its `502`/`504` codes. The re-embed happens *before* the version-guarded write, so a stale edit still `409`s rather than re-embedding for nothing. | ADR-0002, "Re-embedding on content edit" |
| Fractional-key exhaustion | No renormalize endpoint in this slice. Log `points.order_id.gap_low` under a safety threshold; reject with `409` and a distinct error code if the gap would collapse onto a neighbor value. Recovery is a runbook operation. | ADR-0002, "`order_id` gap exhaustion" |
| Batch semantics | All-or-nothing, capped at 100 operations. Every operation's `version` precondition is validated before any is applied; one failure rejects the whole request and nothing reaches Qdrant. | ADR-0002, "`POST /points/batch` semantics" |
| Re-ingestion versus manual edits | The newly uploaded file wins. Surviving points are overwritten in place with an incremented `version`; points absent from the new version are flagged inactive, never removed; manually created points sit past the ingested `chunk_index` range and are swept by the same rule. Clobbered content is recorded in `point_audit_events` as `reingest_overwrite`. | ADR-0002, "Re-ingestion versus manual edits" |
1. **Re-embed inline** on content change, reusing plan 001's embedding ports and
bounds. Consistent, but puts embedder latency and `502`/`504` failure modes
on an admin edit path.
2. **Require caller-supplied vectors** when content changes, and reject the edit
otherwise. Simple and honest, but pushes model knowledge to the client.
3. **Mark the point stale** (a payload flag) and re-embed later. Needs
background work, which ADR-0017 currently rules out.
Default to (1) for parity with ingestion, with the same batch/semaphore bounds
and the same status codes. Record whichever is chosen in ADR-0002 before
implementing Phase 4 — this is a real behavioral contract, not an
implementation detail.
### Fractional-key exhaustion
ADR-0001 notes float keys eventually need renormalization. Decide now whether
this slice ships a renormalize path (an internal operation rewriting a file's
`order_id` values to `1000, 2000, 3000, ...`) or explicitly defers it with a
logged warning when the gap between neighbors falls under a threshold. Deferring
is acceptable; silently producing unrepresentable gaps is not.
### Batch semantics
`POST /v1/points/batch` must define, in the API schema and the tests: whether
operations are all-or-nothing, what happens when operation 3 of 5 fails a
version check, and the maximum operation count per request. Decide before
Phase 5; do not let the answer be "whatever Qdrant happened to do."
Phase 6's cross-slice end-to-end test therefore *verifies* the re-ingestion rule
rather than forcing the decision.
## Build order
@@ -251,8 +236,8 @@ of mutations.
ordering behavior, isolated per test by unique collection or tenant keys.
2. An end-to-end test crossing plan 001 and this slice: ingest a CSV, list its
points, reorder one, soft-delete another, re-upload the same file, and assert
the manual edits interact with re-ingestion exactly as ADR-0001/0002 specify.
If that interaction is not yet decided, this test is what forces the decision.
the manual edits interact with re-ingestion exactly as ADR-0002's
"Re-ingestion versus manual edits" specifies.
3. Structured logging at the mutation boundary with stable event names
(`points.updated`, `points.reordered`, `points.soft_deleted`) carrying
`request_id`, `tenant_id`, `file_id`, and the resulting version.

View File

@@ -80,10 +80,22 @@ Re-running with the same `--slug` reuses the tenant and any domains it already
has, and issues an **additional** key. Both keys stay valid; this adds a key, it
does not rotate one.
Scopes are the security boundary between uploading and managing the allowlist.
Give an upload client `files:write` only. `domains:write` lets its holder create
new domains, which is exactly what the allowlist exists to prevent an upload key
from doing.
Scopes are the security boundary between uploading, reading chunks, and
managing the allowlist. Give an upload client `files:write` only. `domains:write`
lets its holder create new domains, which is exactly what the allowlist exists to
prevent an upload key from doing, and `points:read` lets its holder read the text
of every chunk of every file — so an upload-only key gets neither.
| Scope | Grants |
|---|---|
| `files:write` | Upload a document and read its ingestion status. |
| `points:read` | Read, list, count, and keyword-search this tenant's points, including `GET /v1/files/{file_id}/points`. |
| `points:write` | Create, edit, reorder, and soft-delete points (plan 002 Phases 3-5; no route uses it yet). |
| `domains:read` / `domains:write` | Inspect and manage the domain allowlist. |
| `admin` | Satisfies every scope check. |
The command's `--scopes` default issues all of the above except `admin`, which
suits a first operator key; narrow it explicitly for per-client keys.
### Domains after the first one

View File

@@ -27,6 +27,7 @@ from src.application.ingestion.errors import (
PointIndexingError,
UnsupportedSourceTypeError,
)
from src.application.points.point import PointNotFoundError
logger = structlog.get_logger(__name__)
@@ -41,6 +42,9 @@ _MAPPING: tuple[tuple[type[Exception], int, str], ...] = (
(TenantInactiveError, status.HTTP_401_UNAUTHORIZED, "tenant_not_found"),
(MissingScopeError, status.HTTP_403_FORBIDDEN, "missing_scope"),
(InvalidUploadError, status.HTTP_400_BAD_REQUEST, "validation_error"),
# 404, never 403: a cross-tenant point id must be indistinguishable from a
# nonexistent one, or the API becomes an existence oracle (ADR-0016).
(PointNotFoundError, status.HTTP_404_NOT_FOUND, "not_found"),
(UnknownDomainError, status.HTTP_400_BAD_REQUEST, "unknown_domain"),
(DomainAlreadyExistsError, status.HTTP_409_CONFLICT, "conflict"),
(DocumentParseError, status.HTTP_400_BAD_REQUEST, "validation_error"),

View File

@@ -2,7 +2,9 @@ from fastapi import APIRouter
from src.api.routers.domains import router as domains_router
from src.api.routers.files import router as files_router
from src.api.routers.points import router as points_router
router = APIRouter()
router.include_router(domains_router)
router.include_router(files_router)
router.include_router(points_router)

View File

@@ -14,17 +14,21 @@ from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from src.api.dependencies.auth import require_scope
from src.api.schemas.files import FileStatusResponse, FileUploadResponse
from src.api.schemas.points import DEFAULT_PAGE_SIZE, LimitQuery, PointListResponse
from src.application.auth.context import AuthContext
from src.application.files.status import get_file_status
from src.application.files.upload import upload_source_file
from src.application.points.queries import list_file_points
from src.application.ports.embedding import DenseEmbedder, SparseEmbedder
from src.application.ports.object_storage import ObjectStorage
from src.application.ports.point_repository import PointRepository
from src.application.ports.point_storage import PointStorage
from src.bootstrap.dependencies import (
get_dense_embedders,
get_ingestion_concurrency_limiter,
get_ingestion_limiter,
get_object_storage,
get_point_repository,
get_point_storage,
get_sessionmaker,
get_settings,
@@ -35,9 +39,11 @@ from src.config import Settings
router = APIRouter(prefix="/files", tags=["files"])
_RequireFilesWrite = Annotated[AuthContext, Depends(require_scope("files:write"))]
_RequirePointsRead = Annotated[AuthContext, Depends(require_scope("points:read"))]
_SessionmakerDep = Annotated[async_sessionmaker[AsyncSession], Depends(get_sessionmaker)]
_ObjectStorageDep = Annotated[ObjectStorage, Depends(get_object_storage)]
_PointStorageDep = Annotated[PointStorage, Depends(get_point_storage)]
_PointRepositoryDep = Annotated[PointRepository, Depends(get_point_repository)]
_SettingsDep = Annotated[Settings, Depends(get_settings)]
_IngestionLimiterDep = Annotated[CapacityLimiter, Depends(get_ingestion_limiter)]
_ConcurrencyLimiterDep = Annotated[Semaphore, Depends(get_ingestion_concurrency_limiter)]
@@ -92,3 +98,29 @@ async def get_file(
if result is None:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="file not found")
return FileStatusResponse.from_result(result)
@router.get("/{file_id}/points")
async def list_points_for_file(
file_id: uuid.UUID,
auth: _RequirePointsRead,
repository: _PointRepositoryDep,
limit: LimitQuery = DEFAULT_PAGE_SIZE,
cursor: str | None = None,
include_inactive: bool = False,
) -> PointListResponse:
"""The same listing as `GET /v1/points?file_id=...`, addressed by file.
Gated on `points:read`, not `files:write`: the resource being read is the
file's chunks, so the scope follows the data rather than the URL prefix. An
upload-only key must not become a way to read every chunk of every file.
"""
page = await list_file_points(
repository,
tenant_id=auth.tenant_id,
file_id=file_id,
limit=limit,
cursor=cursor,
include_inactive=include_inactive,
)
return PointListResponse.from_page(page)

134
src/api/routers/points.py Normal file
View File

@@ -0,0 +1,134 @@
"""`/v1/points` read paths (ADR-0002, ADR-0008).
Routes adapt HTTP to `application/points` calls. They build no Qdrant filters
and hold no CRUD semantics (ADR-0015), and they never read a tenant from the
request — `auth.tenant_id` is the only source, which is what makes ADR-0002's
isolation rule structural rather than a habit.
**Route order is load-bearing.** `/count` and `/search` are declared before
`/{point_id}`. FastAPI matches in declaration order, so with `/{point_id}` first
a request for `/v1/points/count` would try to parse `"count"` as a UUID and
fail with `422` instead of counting anything. The failure is loud but confusing,
and it comes back the moment someone reorders these for tidiness.
Gated on `points:read`, separately from `files:write`: a key that can upload
documents should not thereby be able to read every chunk of every file, and
plan 002's mutating paths will want `points:write` distinct again.
"""
import uuid
from typing import Annotated
from fastapi import APIRouter, Depends, Query
from src.api.dependencies.auth import require_scope
from src.api.schemas.points import (
DEFAULT_PAGE_SIZE,
LimitQuery,
PointCountResponse,
PointListResponse,
PointResponse,
PointSearchResponse,
)
from src.application.auth.context import AuthContext
from src.application.points.queries import (
count_points,
get_point,
list_file_points,
search_points,
)
from src.application.ports.point_repository import PointRepository
from src.bootstrap.dependencies import get_point_repository
router = APIRouter(prefix="/points", tags=["points"])
_RequirePointsRead = Annotated[AuthContext, Depends(require_scope("points:read"))]
_PointRepositoryDep = Annotated[PointRepository, Depends(get_point_repository)]
@router.get("/count")
async def count_tenant_points(
auth: _RequirePointsRead,
repository: _PointRepositoryDep,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> PointCountResponse:
count = await count_points(
repository,
tenant_id=auth.tenant_id,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)
return PointCountResponse(count=count)
@router.get("/search")
async def search_tenant_points(
auth: _RequirePointsRead,
repository: _PointRepositoryDep,
q: Annotated[str, Query(min_length=1)],
limit: LimitQuery = DEFAULT_PAGE_SIZE,
cursor: str | None = None,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> PointSearchResponse:
"""Keyword search over point content — **not** semantic retrieval.
Matches Qdrant's full-text payload index on `content`, combined with the
structured filters below. Results are unranked: the index filters rather
than scores, so there is no relevance order and no score to return. Callers
wanting ranked answers want the agent retrieval path (plan 003), not this.
"""
page = await search_points(
repository,
tenant_id=auth.tenant_id,
query=q,
limit=limit,
cursor=cursor,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)
return PointSearchResponse.from_search(page, query=q)
@router.get("")
async def list_tenant_points(
auth: _RequirePointsRead,
repository: _PointRepositoryDep,
file_id: uuid.UUID,
limit: LimitQuery = DEFAULT_PAGE_SIZE,
cursor: str | None = None,
include_inactive: bool = False,
) -> PointListResponse:
"""A file's points in `order_id` order.
`file_id` is required rather than optional: the pagination cursor is an
`order_id` value, and `order_id` is only unique within one file. Listing
across files would silently drop or repeat rows at every page boundary.
"""
page = await list_file_points(
repository,
tenant_id=auth.tenant_id,
file_id=file_id,
limit=limit,
cursor=cursor,
include_inactive=include_inactive,
)
return PointListResponse.from_page(page)
@router.get("/{point_id}")
async def get_tenant_point(
point_id: uuid.UUID,
auth: _RequirePointsRead,
repository: _PointRepositoryDep,
with_vectors: bool = False,
) -> PointResponse:
point = await get_point(
repository, tenant_id=auth.tenant_id, point_id=point_id, with_vectors=with_vectors
)
return PointResponse.from_point(point)

165
src/api/schemas/points.py Normal file
View File

@@ -0,0 +1,165 @@
"""Public request/response models for `/v1/points` (ADR-0002, ADR-0008).
The shape callers see, kept separate from `application/points`' domain models
(ADR-0015). Two rules are encoded here rather than left to route code:
- **Vectors are opt-in.** `PointResponse` omits them unless the caller asked,
so a listing does not ship megabytes of floats nobody reads (ADR-0008).
- **Server-owned fields are not accepted on input.** The request models simply
do not declare `tenant_id`, `version`, or `chunk_index`, and forbid extra
keys, so a client that sends one gets `422` from Pydantic instead of having
it silently ignored — ADR-0002's isolation rule enforced at the boundary.
"""
import uuid
from datetime import datetime
from typing import Annotated
from fastapi import Query
from pydantic import BaseModel, ConfigDict, Field
from src.application.points.point import Point
from src.application.ports.point_repository import PointPage
# A page ceiling the caller cannot raise. Scroll pages are materialized in
# memory both here and in Qdrant, so an unbounded `limit` is a cheap way for one
# request to hurt every other tenant sharing the process. Declared once because
# two routers paginate points -- `/v1/points` and `/v1/files/{file_id}/points` --
# and a ceiling that differs between them is a ceiling in only one of them.
DEFAULT_PAGE_SIZE = 50
MAX_PAGE_SIZE = 200
LimitQuery = Annotated[int, Query(ge=1, le=MAX_PAGE_SIZE)]
class PointResponse(BaseModel):
point_id: uuid.UUID
domain: str
file_id: uuid.UUID
chunk_id: uuid.UUID
content: str
content_type: str
source_filename: str
source_type: str
order_id: float
chunk_index: int
previous_chunk_id: uuid.UUID | None
next_chunk_id: uuid.UUID | None
is_active: bool
deleted_at: datetime | None
created_at: datetime
updated_at: datetime
created_by: str
updated_by: str
version: int
content_hash: str
embedding_model_version: str
vectors: dict[str, object] | None = None
@classmethod
def from_point(cls, point: Point) -> "PointResponse":
# `tenant_id` is present on `Point` and deliberately absent here: the
# caller already knows which tenant it authenticated as, and echoing it
# back invites clients to start sending it.
return cls.model_validate(point.model_dump(exclude={"tenant_id"}))
class PointListResponse(BaseModel):
"""A page of points plus the cursor for the next one.
Cursor-based rather than `limit`/`offset`: an offset cursor silently skips
or repeats rows when a concurrent insert shifts positions, which is exactly
the pagination defect plan 002 requires a test for.
"""
points: list[PointResponse]
next_cursor: str | None = None
@classmethod
def from_page(cls, page: PointPage) -> "PointListResponse":
return cls(
points=[PointResponse.from_point(point) for point in page.points],
next_cursor=page.next_cursor,
)
class PointCountResponse(BaseModel):
count: int
class PointSearchResponse(PointListResponse):
"""Results of a **keyword** match, not of semantic retrieval.
Named and documented so it cannot be mistaken for ADR-0003's hybrid
retrieval: these points matched a full-text filter on `content`, they are
not ranked by relevance, and there is no score to report. Anything that
wants ranked results wants the agent retrieval path in plan 003.
"""
query: str
@classmethod
def from_search(cls, page: PointPage, *, query: str) -> "PointSearchResponse":
return cls(
query=query,
points=[PointResponse.from_point(point) for point in page.points],
next_cursor=page.next_cursor,
)
class PointCreateRequest(BaseModel):
"""Create one point. The server assigns identity, ordering, and provenance.
`after_point_id` positions the new point rather than a raw `order_id`: the
caller says where in the sequence it goes and the server computes the
fractional key and relinks neighbours, which a client-supplied `order_id`
could not do correctly (ADR-0002). `None` means "at the start of the file".
"""
model_config = ConfigDict(extra="forbid")
file_id: uuid.UUID
content: str = Field(min_length=1)
content_type: str = "paragraph"
after_point_id: uuid.UUID | None = None
class PointReplaceRequest(BaseModel):
"""Replace a point's content under a version guard.
`version` here is the *expected* version, not a value being written — the
optimistic-concurrency precondition. A mismatch is `409`.
"""
model_config = ConfigDict(extra="forbid")
content: str = Field(min_length=1)
content_type: str | None = None
version: int
class PointPayloadPatchRequest(BaseModel):
"""Payload-only update of caller-writable fields, under a version guard."""
model_config = ConfigDict(extra="forbid")
payload: dict[str, object]
version: int
class PointReorderRequest(BaseModel):
"""Move a point to sit immediately after `after_point_id`.
`None` moves it to the front of the file. Expressed as a neighbour rather
than an `order_id` for the same reason as `PointCreateRequest`.
"""
model_config = ConfigDict(extra="forbid")
after_point_id: uuid.UUID | None = None
version: int

View File

@@ -0,0 +1,126 @@
"""The caller-facing point model for `/v1/points` (ADR-0001, ADR-0002).
`ChunkPoint` in `models.py` is the *write* shape ingestion upserts: an id, its
named vectors, and an opaque payload dict. This module is the *read/edit* shape
the `/v1/points` surface works in, where the payload's individual fields matter
and the distinction between what a caller may write and what the server owns is
a security boundary rather than a convention.
That split is the reason this is a model and not a dict. ADR-0002's isolation
rule ("never accepted as client-supplied input") and its optimistic-concurrency
guard both fail open if a caller can smuggle `tenant_id` or `version` through a
payload update, so the writable field set is enumerated in one place here and
every write path validates against it.
"""
import uuid
from collections.abc import Mapping
from datetime import datetime
from typing import Self
from pydantic import BaseModel, ConfigDict, Field
# Fields the server derives and a caller may never set, patch, or override.
# `tenant_id` is authority, `version` is the concurrency guard, `chunk_index`
# derives the point id, and the rest are provenance the server timestamps.
SERVER_OWNED_FIELDS: frozenset[str] = frozenset(
{
"tenant_id",
"version",
"chunk_index",
"chunk_id",
"is_active",
"deleted_at",
"created_at",
"created_by",
"updated_at",
"updated_by",
"content_hash",
"embedding_model_version",
}
)
# Fields a caller may supply on create, replace, or payload patch. `order_id`
# is writable on create but moves only through `PATCH /v1/points/{id}/order`
# afterwards, because a bare `order_id` write would not relink neighbours.
CALLER_WRITABLE_FIELDS: frozenset[str] = frozenset(
{
"content",
"content_type",
"domain",
"source_filename",
"source_type",
"order_id",
"previous_chunk_id",
"next_chunk_id",
}
)
class Point(BaseModel):
"""One Qdrant point, read back with its ADR-0001 payload fields typed.
Vectors are deliberately absent: ADR-0008 returns them only when explicitly
requested, and every read path that does not ask for them should not pay to
deserialize them.
"""
model_config = ConfigDict(frozen=True)
point_id: uuid.UUID
tenant_id: uuid.UUID
domain: str
file_id: uuid.UUID
chunk_id: uuid.UUID
content: str
content_type: str
source_filename: str
source_type: str
order_id: float
chunk_index: int
previous_chunk_id: uuid.UUID | None = None
next_chunk_id: uuid.UUID | None = None
is_active: bool = True
deleted_at: datetime | None = None
created_at: datetime
updated_at: datetime
created_by: str
updated_by: str
version: int
content_hash: str
embedding_model_version: str
# Only populated when the caller explicitly asked for vectors.
vectors: dict[str, object] | None = Field(default=None)
@classmethod
def from_payload(
cls,
point_id: uuid.UUID,
payload: Mapping[str, object],
*,
vectors: Mapping[str, object] | None = None,
) -> Self:
"""Build a `Point` from a raw Qdrant payload dict.
Lives here rather than in the Qdrant adapter so the payload field names
are declared once, next to the model that mirrors them. The adapter
stays responsible for talking to the SDK, not for knowing ADR-0001's
schema twice.
"""
return cls.model_validate({**payload, "point_id": point_id, "vectors": vectors})
class PointNotFoundError(LookupError):
"""No such point *within the requesting tenant*.
Routes map this to `404`, never `403` — a caller must not be able to probe
for the existence of another tenant's point ids (ADR-0016). The error
deliberately carries no hint about which of the two cases occurred.
"""

View File

@@ -0,0 +1,115 @@
"""Read paths for `/v1/points` (ADR-0002, ADR-0008).
The caller-facing entry points for point reads. Routers call these; they never
touch the `PointRepository` directly, and never build a filter.
This module is thin on purpose but not empty, and the two things it does own are
exactly the ones a route would otherwise get wrong:
- **`tenant_id` always comes from the caller's `AuthContext`.** Every function
takes it as a required keyword and hands it to the repository. Nothing here
reads a tenant from a query string or body.
- **A keyword query is Persian-normalized before it reaches the index.**
Ingestion letter-folds chunk content (`normalize_persian_text`, ADR-0018), so
stored text contains Persian yeh/keheh. A query typed on an Arabic keyboard
carries U+064A/U+0643 and would match nothing at all — a silent empty result,
not an error. Folding the query the same way is what makes the two comparable.
Reads emit no log events. The request middleware already records every call, and
ADR-0011 reserves `INFO` for lifecycle events rather than per-read volume;
mutations get their own events when those paths land.
"""
import uuid
from src.application.ingestion.normalization import normalize_persian_text
from src.application.points.point import Point, PointNotFoundError
from src.application.ports.point_repository import PointPage, PointRepository
async def get_point(
repository: PointRepository,
*,
tenant_id: uuid.UUID,
point_id: uuid.UUID,
with_vectors: bool = False,
) -> Point:
"""One point, or `PointNotFoundError` if it is not this tenant's.
Raises rather than returning `None` so a route cannot forget the check and
serve `200 null`. "Absent" and "another tenant's" are the same outcome by
design (ADR-0016: cross-tenant access is `404`, never `403`).
"""
point = await repository.get(tenant_id=tenant_id, point_id=point_id, with_vectors=with_vectors)
if point is None:
raise PointNotFoundError(f"point {point_id} not found")
return point
async def list_file_points(
repository: PointRepository,
*,
tenant_id: uuid.UUID,
file_id: uuid.UUID,
limit: int,
cursor: str | None = None,
include_inactive: bool = False,
) -> PointPage:
"""One file's points in display (`order_id`) order.
An unknown or foreign `file_id` yields an empty page rather than an error:
the two are indistinguishable to the caller, which is the same
non-disclosure property `get_point` gets from raising.
"""
return await repository.list_by_file(
tenant_id=tenant_id,
file_id=file_id,
limit=limit,
cursor=cursor,
include_inactive=include_inactive,
)
async def count_points(
repository: PointRepository,
*,
tenant_id: uuid.UUID,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> int:
return await repository.count(
tenant_id=tenant_id,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)
async def search_points(
repository: PointRepository,
*,
tenant_id: uuid.UUID,
query: str,
limit: int,
cursor: str | None = None,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> PointPage:
"""Keyword match on `content`, within this tenant.
**Not semantic retrieval.** Qdrant's full-text index filters rather than
scores, so results carry no relevance ranking and their order is
unspecified. Ranked retrieval is ADR-0003's hybrid path in plan 003; this
function must not grow a semantic mode (ADR-0002).
"""
return await repository.keyword_search(
tenant_id=tenant_id,
query=normalize_persian_text(query),
limit=limit,
cursor=cursor,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)

View File

@@ -0,0 +1,131 @@
"""The point read/edit port for `/v1/points` (ADR-0002, ADR-0015).
Separate from `PointStorage`, which stays exactly the two bulk operations
ingestion performs. Reads, single-point edits, reordering, and keyword search
have a different caller, a different failure vocabulary, and a different
tenant-filter obligation, so they get their own port rather than accreting onto
the ingestion one.
`tenant_id` is a required keyword argument on **every** method. That is not
style: ADR-0002's isolation rule has to hold on every code path that touches the
collection, and an optional tenant filter is one forgotten argument away from a
cross-tenant read. Making it required moves that from a review question to a
type error.
`src/infrastructure/qdrant/point_repository.py` is the production adapter;
`tests.fakes.FakePointRepository` is the test double.
"""
import uuid
from collections.abc import Sequence
from typing import Protocol
from pydantic import BaseModel
from src.application.points.point import Point
class PointPage(BaseModel):
"""One page of points plus the cursor that continues it.
`next_cursor` is opaque to callers and encoded by the adapter: ordered
scrolls and keyword searches paginate by different Qdrant mechanisms, and
neither is a plain integer offset. `None` means the listing is exhausted.
"""
points: tuple[Point, ...]
next_cursor: str | None = None
class PayloadPatch(BaseModel):
"""Set these payload fields on one point, optionally guarded by `version`.
When `expected_version` is set, the adapter attaches it to the operation's
filter, so a concurrent write that has already moved the version on means
this patch matches nothing rather than clobbering it. The guard is what
makes a lost update impossible; detecting that it fired is the service's
job (see `apply_patches`).
"""
point_id: uuid.UUID
payload: dict[str, object]
expected_version: int | None = None
class PointRepository(Protocol):
async def get(
self, *, tenant_id: uuid.UUID, point_id: uuid.UUID, with_vectors: bool = False
) -> Point | None:
"""One point, or `None` if it does not exist *under this tenant*.
The two cases are deliberately indistinguishable — the route maps both
to `404` so a caller cannot probe for another tenant's point ids.
"""
...
async def get_many(
self, *, tenant_id: uuid.UUID, point_ids: Sequence[uuid.UUID]
) -> tuple[Point, ...]:
"""The subset of `point_ids` that exists under this tenant.
Order is not guaranteed and missing ids are silently absent: callers are
neighbour-relinking and batch precondition checks, both of which match
on id rather than position.
"""
...
async def list_by_file(
self,
*,
tenant_id: uuid.UUID,
file_id: uuid.UUID,
limit: int,
cursor: str | None = None,
include_inactive: bool = False,
) -> PointPage:
"""One file's points in `order_id` order (ADR-0008's `scroll`).
Scoped to a single file because the cursor is an `order_id` value, and
`order_id` is only unique within a file.
"""
...
async def count(
self,
*,
tenant_id: uuid.UUID,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> int: ...
async def keyword_search(
self,
*,
tenant_id: uuid.UUID,
query: str,
limit: int,
cursor: str | None = None,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> PointPage:
"""Full-text payload match on `content`, plus structured filters.
Keyword matching, **not** semantic retrieval (ADR-0002). Qdrant's
full-text index is a filter, not a scorer, so results carry no relevance
ranking and their order is unspecified.
"""
...
async def apply_patches(self, *, tenant_id: uuid.UUID, patches: Sequence[PayloadPatch]) -> None:
"""Apply every patch in one Qdrant `points/batch` request.
Qdrant has no multi-point transaction, so this is not atomic and does
not pretend to be. ADR-0002's all-or-nothing rule is implemented one
layer up as validate-every-precondition-then-apply; the per-patch
`expected_version` guard here is what makes the residual window safe,
turning a lost update into a no-op the service can detect rather than a
silent clobber.
"""
...

View File

@@ -28,7 +28,13 @@ from src.infrastructure.postgres.repositories import tenants as tenants_repo
logger = structlog.get_logger(__name__)
DEFAULT_SCOPES = ("files:write", "domains:read", "domains:write")
DEFAULT_SCOPES = (
"files:write",
"domains:read",
"domains:write",
"points:read",
"points:write",
)
@dataclass(frozen=True)

View File

@@ -9,6 +9,7 @@ from sqlalchemy.ext.asyncio import AsyncEngine, AsyncSession, async_sessionmaker
from src.application.ports.embedding import DenseEmbedder, SparseEmbedder
from src.application.ports.object_storage import ObjectStorage
from src.application.ports.point_repository import PointRepository
from src.application.ports.point_storage import PointStorage
from src.config import Settings
@@ -22,6 +23,7 @@ class AppResources:
qdrant_client: AsyncQdrantClient
object_storage: ObjectStorage
point_storage: PointStorage
point_repository: PointRepository
ingestion_limiter: CapacityLimiter
dense_embedders: Sequence[DenseEmbedder]
sparse_embedder: SparseEmbedder
@@ -52,6 +54,10 @@ def get_point_storage(request: Request) -> PointStorage:
return _resources(request).point_storage
def get_point_repository(request: Request) -> PointRepository:
return _resources(request).point_repository
def get_ingestion_limiter(request: Request) -> CapacityLimiter:
return _resources(request).ingestion_limiter

View File

@@ -20,6 +20,7 @@ from src.infrastructure.minio.storage import MinioObjectStorage
from src.infrastructure.observability.logging import configure_logging
from src.infrastructure.postgres.database import create_engine, create_sessionmaker
from src.infrastructure.qdrant.client import create_client as create_qdrant_client
from src.infrastructure.qdrant.point_repository import QdrantPointRepository
from src.infrastructure.qdrant.points import QdrantPointStorage
logger = structlog.get_logger(__name__)
@@ -85,6 +86,9 @@ def create_lifespan(
point_storage = QdrantPointStorage(
qdrant_client, collection=resolved_settings.qdrant.collection
)
point_repository = QdrantPointRepository(
qdrant_client, collection=resolved_settings.qdrant.collection
)
logger.info("lifespan.qdrant.client.created")
nomic_settings = resolved_settings.embedding.nomic
@@ -145,6 +149,7 @@ def create_lifespan(
qdrant_client=qdrant_client,
object_storage=object_storage,
point_storage=point_storage,
point_repository=point_repository,
ingestion_limiter=ingestion_limiter,
dense_embedders=dense_embedders,
sparse_embedder=sparse_embedder,

View File

@@ -80,14 +80,25 @@ def _sparse_vectors_config() -> dict[str, models.SparseVectorParams]:
return {SPARSE_VECTOR: models.SparseVectorParams(modifier=models.Modifier.IDF)}
# `content`'s full-text index backs ADR-0002's keyword search. The
# `multilingual` tokenizer is the one that segments Persian correctly; `word`
# splits on non-alphanumerics, which mis-handles ZWNJ-joined compounds. No
# stemmer or stopword list is configured: `content` is already letter-folded by
# `normalize_persian_text` at ingest (ADR-0018), and the *ranked* Farsi lexical
# path is the benchmarked BM25 sparse vector, not this index. This one exists
# for exact keyword/filter matching, which ADR-0002 keeps deliberately distinct
# from retrieval.
_CONTENT_INDEX = models.TextIndexParams(
type=models.TextIndexType.TEXT,
tokenizer=models.TokenizerType.MULTILINGUAL,
lowercase=True,
)
# (field name, schema). `order_id` is float because Qdrant's `Range` conditions
# and `order_by` only support numeric/datetime payloads -- a keyword key could
# only be sorted client-side after fetching every chunk (ADR-0001).
#
# The full-text index on `content` is deliberately absent: it belongs to plan
# 002's keyword search, and payload indexes -- unlike vector config -- can be
# added to a live collection later.
_PAYLOAD_INDEXES: tuple[tuple[str, models.PayloadSchemaType | models.KeywordIndexParams], ...] = (
_PayloadIndexSchema = models.PayloadSchemaType | models.KeywordIndexParams | models.TextIndexParams
_PAYLOAD_INDEXES: tuple[tuple[str, _PayloadIndexSchema], ...] = (
# `is_tenant` co-locates a tenant's vectors on disk for sequential reads,
# which is the whole point of payload-partitioned multitenancy.
(
@@ -99,6 +110,13 @@ _PAYLOAD_INDEXES: tuple[tuple[str, models.PayloadSchemaType | models.KeywordInde
("order_id", models.PayloadSchemaType.FLOAT),
("previous_chunk_id", models.PayloadSchemaType.KEYWORD),
("next_chunk_id", models.PayloadSchemaType.KEYWORD),
("content", _CONTENT_INDEX),
# Every read path filters on `is_active` (ADR-0002 implies `is_active: true`
# unless the caller opts in), and the re-ingestion sweep and `/v1/points`
# both range over `chunk_index`. Both were unindexed while ingestion was the
# only reader; plan 002 makes them hot.
("is_active", models.PayloadSchemaType.BOOL),
("chunk_index", models.PayloadSchemaType.INTEGER),
)

View File

@@ -0,0 +1,247 @@
"""Qdrant adapter for the `PointRepository` port (ADR-0002, ADR-0015).
Every Qdrant filter for the `/v1/points` surface is built here. Routers and
application services never import `qdrant_client` — they pass a `tenant_id` and
get `Point` models back.
Two mechanics are worth reading before changing anything:
**Reads go through `scroll`, not `retrieve`.** `retrieve` fetches by id and
takes no filter, which would force the tenant check to happen *after* Qdrant
answered — exactly the "check it in Python afterwards" shape ADR-0002's
isolation rule exists to prevent. `scroll` with a `HasIdCondition` plus the
tenant condition pushes the check server-side, so a foreign id returns an empty
page rather than a row this code has to remember to reject.
**Ordered listing paginates by `order_id` value, not by offset.** Qdrant does
not return a `next_page_offset` when `order_by` is set, and an offset-based
cursor would skip or repeat rows when a concurrent insert shifts positions
underneath the reader. Ranging on `order_id > cursor` is stable instead: a point
inserted ahead of the cursor was already returned, and one inserted after it
shows up on a later page. This relies on `order_id` being unique within a file,
which ADR-0002 guarantees by rejecting a reorder whose gap would collapse onto a
neighbour value.
"""
import uuid
from collections.abc import Sequence
from qdrant_client import AsyncQdrantClient, models
from src.application.points.point import Point
from src.application.ports.point_repository import PayloadPatch, PointPage
# Qdrant's scroll returns `(records, next_page_offset)`; a record's id is a
# `str | int` union in the SDK, and every id this service writes is a UUID.
type _ScrollRecord = models.Record
def _uuid_of(record: _ScrollRecord) -> uuid.UUID:
return uuid.UUID(str(record.id))
def _to_point(record: _ScrollRecord, *, with_vectors: bool) -> Point:
vectors = record.vector if with_vectors and isinstance(record.vector, dict) else None
return Point.from_payload(_uuid_of(record), record.payload or {}, vectors=vectors)
def _conditions(
*,
tenant_id: uuid.UUID,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> list[models.Condition]:
"""The base filter every point query carries.
`tenant_id` is unconditional and first. `is_active` is added unless the
caller explicitly opted into inactive points, which is ADR-0002's
"`is_active: true` implied" rule expressed once rather than per method.
"""
conditions: list[models.Condition] = [
models.FieldCondition(key="tenant_id", match=models.MatchValue(value=str(tenant_id)))
]
if domain is not None:
conditions.append(
models.FieldCondition(key="domain", match=models.MatchValue(value=domain))
)
if file_id is not None:
conditions.append(
models.FieldCondition(key="file_id", match=models.MatchValue(value=str(file_id)))
)
if not include_inactive:
conditions.append(
models.FieldCondition(key="is_active", match=models.MatchValue(value=True))
)
return conditions
class QdrantPointRepository:
"""A `PointRepository` (see `src/application/ports/point_repository.py`)."""
def __init__(self, client: AsyncQdrantClient, *, collection: str) -> None:
self._client = client
self._collection = collection
async def get(
self, *, tenant_id: uuid.UUID, point_id: uuid.UUID, with_vectors: bool = False
) -> Point | None:
records, _ = await self._client.scroll(
collection_name=self._collection,
scroll_filter=models.Filter(
must=[
*_conditions(tenant_id=tenant_id, include_inactive=True),
models.HasIdCondition(has_id=[str(point_id)]),
]
),
limit=1,
with_payload=True,
with_vectors=with_vectors,
)
if not records:
return None
return _to_point(records[0], with_vectors=with_vectors)
async def get_many(
self, *, tenant_id: uuid.UUID, point_ids: Sequence[uuid.UUID]
) -> tuple[Point, ...]:
if not point_ids:
return ()
records, _ = await self._client.scroll(
collection_name=self._collection,
scroll_filter=models.Filter(
must=[
*_conditions(tenant_id=tenant_id, include_inactive=True),
models.HasIdCondition(has_id=[str(point_id) for point_id in point_ids]),
]
),
limit=len(point_ids),
with_payload=True,
with_vectors=False,
)
return tuple(_to_point(record, with_vectors=False) for record in records)
async def list_by_file(
self,
*,
tenant_id: uuid.UUID,
file_id: uuid.UUID,
limit: int,
cursor: str | None = None,
include_inactive: bool = False,
) -> PointPage:
conditions = _conditions(
tenant_id=tenant_id, file_id=file_id, include_inactive=include_inactive
)
if cursor is not None:
conditions.append(
models.FieldCondition(key="order_id", range=models.Range(gt=float(cursor)))
)
records, _ = await self._client.scroll(
collection_name=self._collection,
scroll_filter=models.Filter(must=conditions),
order_by=models.OrderBy(key="order_id", direction=models.Direction.ASC),
limit=limit,
with_payload=True,
with_vectors=False,
)
points = tuple(_to_point(record, with_vectors=False) for record in records)
# A short page means the listing is exhausted. A full page might be, but
# claiming so would need an extra round trip; handing back a cursor that
# yields an empty final page is the cheaper honest answer.
next_cursor = repr(points[-1].order_id) if len(points) == limit else None
return PointPage(points=points, next_cursor=next_cursor)
async def count(
self,
*,
tenant_id: uuid.UUID,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> int:
result = await self._client.count(
collection_name=self._collection,
count_filter=models.Filter(
must=_conditions(
tenant_id=tenant_id,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)
),
exact=True,
)
return result.count
async def keyword_search(
self,
*,
tenant_id: uuid.UUID,
query: str,
limit: int,
cursor: str | None = None,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> PointPage:
conditions = _conditions(
tenant_id=tenant_id,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)
conditions.append(models.FieldCondition(key="content", match=models.MatchText(text=query)))
# No `order_by` here, so Qdrant does return a page offset: the full-text
# index filters rather than scores, and imposing `order_id` ordering
# across files would be meaningless (`order_id` is per-file).
records, next_offset = await self._client.scroll(
collection_name=self._collection,
scroll_filter=models.Filter(must=conditions),
limit=limit,
offset=cursor,
with_payload=True,
with_vectors=False,
)
return PointPage(
points=tuple(_to_point(record, with_vectors=False) for record in records),
next_cursor=str(next_offset) if next_offset is not None else None,
)
async def apply_patches(self, *, tenant_id: uuid.UUID, patches: Sequence[PayloadPatch]) -> None:
if not patches:
return
await self._client.batch_update_points(
collection_name=self._collection,
update_operations=[
models.SetPayloadOperation(
set_payload=models.SetPayload(
payload=patch.payload,
filter=models.Filter(must=self._patch_conditions(tenant_id, patch)),
)
)
for patch in patches
],
)
@staticmethod
def _patch_conditions(tenant_id: uuid.UUID, patch: PayloadPatch) -> list[models.Condition]:
"""Address one point by id, under this tenant, at an expected version.
`include_inactive=True`: soft-deleting and relinking both have to reach
points the default read filter hides, and a patch is addressed by an
explicit id rather than discovered by a listing.
"""
conditions: list[models.Condition] = [
*_conditions(tenant_id=tenant_id, include_inactive=True),
models.HasIdCondition(has_id=[str(patch.point_id)]),
]
if patch.expected_version is not None:
conditions.append(
models.FieldCondition(
key="version", match=models.MatchValue(value=patch.expected_version)
)
)
return conditions

View File

@@ -8,6 +8,8 @@ from datetime import datetime
from src.application.ingestion.models import SparseVector
from src.application.points.models import ChunkPoint
from src.application.points.point import Point
from src.application.ports.point_repository import PayloadPatch, PointPage
@dataclass
@@ -139,3 +141,161 @@ class FakePointStorage:
point.payload["is_active"] = False
point.payload["deleted_at"] = deleted_at.isoformat()
return len(stale)
@dataclass
class FakePointRepository:
"""An in-memory `PointRepository`.
Stores `Point` models keyed by point id. The filter, ordering, and version
semantics below are reimplemented in Python rather than stubbed, because
every service this fake stands behind is being tested *for* those semantics:
a stub that returned points in insertion order would pass a reorder test
that real Qdrant fails.
Three behaviours are modelled deliberately:
- `is_active` is filtered out unless `include_inactive`, matching ADR-0002's
implied read filter.
- `list_by_file` sorts by `order_id` and paginates by value, the same cursor
mechanism the Qdrant adapter uses, so a test can catch a cursor that skips
or repeats.
- `apply_patches` honours `expected_version` by silently matching nothing on
a mismatch — the no-op-not-error behaviour Qdrant's filtered `set_payload`
actually has, which is what makes the service's post-check necessary.
"""
points: dict[str, Point] = field(default_factory=dict)
patch_batches: list[list[PayloadPatch]] = field(default_factory=list)
fail_next: bool = False
def add(self, point: Point) -> None:
"""Seed a point. Test-only helper, not part of the port."""
self.points[str(point.point_id)] = point
def _visible(
self,
*,
tenant_id: uuid.UUID,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> list[Point]:
return [
point
for point in self.points.values()
if point.tenant_id == tenant_id
and (domain is None or point.domain == domain)
and (file_id is None or point.file_id == file_id)
and (include_inactive or point.is_active)
]
async def get(
self, *, tenant_id: uuid.UUID, point_id: uuid.UUID, with_vectors: bool = False
) -> Point | None:
point = self.points.get(str(point_id))
if point is None or point.tenant_id != tenant_id:
return None
return point if with_vectors else point.model_copy(update={"vectors": None})
async def get_many(
self, *, tenant_id: uuid.UUID, point_ids: Sequence[uuid.UUID]
) -> tuple[Point, ...]:
wanted = {str(point_id) for point_id in point_ids}
return tuple(
point
for key, point in self.points.items()
if key in wanted and point.tenant_id == tenant_id
)
async def list_by_file(
self,
*,
tenant_id: uuid.UUID,
file_id: uuid.UUID,
limit: int,
cursor: str | None = None,
include_inactive: bool = False,
) -> PointPage:
ordered = sorted(
self._visible(tenant_id=tenant_id, file_id=file_id, include_inactive=include_inactive),
key=lambda point: point.order_id,
)
if cursor is not None:
ordered = [point for point in ordered if point.order_id > float(cursor)]
page = tuple(ordered[:limit])
next_cursor = repr(page[-1].order_id) if len(page) == limit else None
return PointPage(points=page, next_cursor=next_cursor)
async def count(
self,
*,
tenant_id: uuid.UUID,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> int:
return len(
self._visible(
tenant_id=tenant_id,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)
)
async def keyword_search(
self,
*,
tenant_id: uuid.UUID,
query: str,
limit: int,
cursor: str | None = None,
domain: str | None = None,
file_id: uuid.UUID | None = None,
include_inactive: bool = False,
) -> PointPage:
# Whole-token containment, lowercased: an approximation of Qdrant's
# full-text index good enough to tell "matched" from "did not", which is
# all the service layer branches on. Ranking is not modelled because the
# real index does not rank either (ADR-0002).
terms = query.lower().split()
matches = [
point
for point in self._visible(
tenant_id=tenant_id,
domain=domain,
file_id=file_id,
include_inactive=include_inactive,
)
if all(term in point.content.lower().split() for term in terms)
]
start = int(cursor) if cursor is not None else 0
page = tuple(matches[start : start + limit])
next_start = start + limit
return PointPage(
points=page, next_cursor=str(next_start) if next_start < len(matches) else None
)
async def apply_patches(self, *, tenant_id: uuid.UUID, patches: Sequence[PayloadPatch]) -> None:
self.patch_batches.append(list(patches))
if self.fail_next:
self.fail_next = False
raise RuntimeError("simulated point repository failure")
for patch in patches:
point = self.points.get(str(patch.point_id))
if point is None or point.tenant_id != tenant_id:
continue
if patch.expected_version is not None and point.version != patch.expected_version:
continue
# Re-validated rather than `model_copy`d, because a patch payload
# carries wire values (an ISO `deleted_at`, a stringified
# `previous_chunk_id`) exactly as it would reach Qdrant. Copying
# without validation would leave the fake holding a `str` where a
# read from real Qdrant returns a `datetime`, and a service bug that
# depends on the difference would pass here and fail in production.
updated = point.model_dump()
updated.update(patch.payload)
self.points[str(patch.point_id)] = Point.model_validate(updated)

View File

@@ -65,10 +65,42 @@ async def test_ensure_chunks_collection_creates_the_payload_indexes(
"order_id",
"previous_chunk_id",
"next_chunk_id",
"content",
"is_active",
"chunk_index",
}
# order_id must be numeric: Qdrant's Range/order_by reject keyword payloads.
assert schema["order_id"].data_type == models.PayloadSchemaType.FLOAT
assert schema["tenant_id"].data_type == models.PayloadSchemaType.KEYWORD
# `content` must be TEXT, not KEYWORD: ADR-0002's keyword search is a
# full-text match on it, and a keyword index would only match the entire
# chunk verbatim -- which never happens and would fail silently.
assert schema["content"].data_type == models.PayloadSchemaType.TEXT
assert schema["is_active"].data_type == models.PayloadSchemaType.BOOL
assert schema["chunk_index"].data_type == models.PayloadSchemaType.INTEGER
async def test_ensure_chunks_collection_adds_a_missing_index_to_a_live_collection(
qdrant_client: AsyncQdrantClient, qdrant_settings: QdrantSettings
) -> None:
"""Payload indexes are additive, unlike vector config.
This is the operational claim the runbook makes when a new index ships:
re-running the bootstrap against an existing collection adds it in place, so
plan 002's `content`/`is_active`/`chunk_index` indexes do not require
recreating a collection that already holds a tenant's points.
"""
await ensure_chunks_collection(qdrant_client, collection=qdrant_settings.collection)
await qdrant_client.delete_payload_index(
collection_name=qdrant_settings.collection, field_name="content"
)
info = await qdrant_client.get_collection(qdrant_settings.collection)
assert "content" not in info.payload_schema
assert not await ensure_chunks_collection(qdrant_client, collection=qdrant_settings.collection)
info = await qdrant_client.get_collection(qdrant_settings.collection)
assert info.payload_schema["content"].data_type == models.PayloadSchemaType.TEXT
async def test_ensure_chunks_collection_is_idempotent(

View File

@@ -0,0 +1,86 @@
"""`QdrantPointRepository` against the shared `PointRepository` contract.
The identical scenarios run against the in-memory fake in
`tests/unit/application/points/test_point_repository_contract.py`. This file is
the half that decides whether the fake is telling the truth: filter
construction, `order_by: order_id` scroll, cursor pagination, full-text
matching, and the filtered `set_payload` version guard are all Qdrant
behaviours a fake can only approximate, and plan 002 Phase 1's exit criterion is
that the approximation holds.
Seeding differs from the fake's (real points need vectors), so it lives here;
everything asserted lives in `tests/support/point_contract.py`.
"""
import uuid
from collections.abc import Awaitable, Callable
import pytest
from qdrant_client import AsyncQdrantClient, models
from src.config import QdrantSettings
from src.infrastructure.qdrant.collection import ensure_chunks_collection
from src.infrastructure.qdrant.point_repository import QdrantPointRepository
from src.infrastructure.qdrant.points import QdrantPointStorage
from tests.support import point_contract
from tests.support.point_contract import chunk_point_for, seed_specs
pytestmark = [
pytest.mark.integration,
pytest.mark.qdrant,
pytest.mark.asyncio(loop_scope="session"),
]
type Scenario = Callable[..., Awaitable[None]]
TENANT_A = uuid.UUID("11111111-1111-4111-8111-111111111111")
TENANT_B = uuid.UUID("22222222-2222-4222-8222-222222222222")
FILE_A = uuid.UUID("33333333-3333-4333-8333-333333333333")
FILE_B = uuid.UUID("44444444-4444-4444-8444-444444444444")
FILE_C = uuid.UUID("55555555-5555-4555-8555-555555555555")
IDS = {
"tenant_a": TENANT_A,
"tenant_b": TENANT_B,
"file_a": FILE_A,
"file_b": FILE_B,
"file_c": FILE_C,
}
async def _seeded(
qdrant_client: AsyncQdrantClient, qdrant_settings: QdrantSettings
) -> QdrantPointRepository:
await ensure_chunks_collection(qdrant_client, collection=qdrant_settings.collection)
storage = QdrantPointStorage(qdrant_client, collection=qdrant_settings.collection)
await storage.upsert_points(
[chunk_point_for(spec) for spec in seed_specs(TENANT_A, TENANT_B, FILE_A, FILE_B, FILE_C)]
)
# Payload indexes are built asynchronously; without waiting, a full-text or
# ordered query can run against a half-built index and return short.
await _await_indexing(qdrant_client, qdrant_settings.collection)
return QdrantPointRepository(qdrant_client, collection=qdrant_settings.collection)
async def _await_indexing(client: AsyncQdrantClient, collection: str) -> None:
for _ in range(100):
info = await client.get_collection(collection)
if info.status == models.CollectionStatus.GREEN and info.indexed_vectors_count is not None:
return
raise AssertionError(f"collection {collection!r} did not finish indexing")
def _arguments(scenario: Scenario) -> dict[str, uuid.UUID]:
return {name: value for name, value in IDS.items() if name in scenario.__annotations__}
@pytest.mark.parametrize(
"scenario",
[*point_contract.READ_SCENARIOS, *point_contract.WRITE_SCENARIOS],
ids=lambda scenario: scenario.__name__.removeprefix("scenario_"),
)
async def test_qdrant_point_repository_satisfies_the_contract(
scenario: Scenario, qdrant_client: AsyncQdrantClient, qdrant_settings: QdrantSettings
) -> None:
repository = await _seeded(qdrant_client, qdrant_settings)
await scenario(repository, **_arguments(scenario))

View File

@@ -0,0 +1,395 @@
"""`/v1/points` read paths over HTTP (ADR-0002, ADR-0008).
Needs both containers: Postgres authenticates the key and derives the tenant,
Qdrant holds the points. Testing them together is the point — the invariant this
file exists for is that the tenant Postgres derived is the only one Qdrant is
ever queried with.
Everything below goes through the real app: real routing, real `require_scope`,
real error envelope. Only the session factory is overridden, so a route that
forgot its scope or read a tenant from the query string fails here.
"""
import uuid
from collections.abc import AsyncIterator
import pytest
import pytest_asyncio
from asgi_lifespan import LifespanManager
from httpx import ASGITransport, AsyncClient
from qdrant_client import AsyncQdrantClient
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from src.bootstrap.dependencies import get_sessionmaker
from src.config import QdrantSettings, Settings
from src.infrastructure.qdrant.collection import ensure_chunks_collection
from src.infrastructure.qdrant.points import QdrantPointStorage
from src.main import create_app
from tests.support.factories import create_api_key, create_tenant
from tests.support.point_contract import SeedSpec, chunk_point_for, point_id_for
pytestmark = [
pytest.mark.integration,
pytest.mark.qdrant,
pytest.mark.postgres,
pytest.mark.asyncio(loop_scope="session"),
]
FILE_ID = uuid.UUID("aaaaaaaa-aaaa-4aaa-8aaa-aaaaaaaaaaaa")
OTHER_FILE_ID = uuid.UUID("bbbbbbbb-bbbb-4bbb-8bbb-bbbbbbbbbbbb")
@pytest.fixture
def api_settings(settings: Settings, qdrant_settings: QdrantSettings) -> Settings:
"""The unit-test settings with Qdrant pointed at the real container.
Everything else keeps its closed-port default, so this exercises the point
paths without the lifespan reaching a colleague's Ollama box or OpenAI.
"""
return settings.model_copy(update={"qdrant": qdrant_settings})
@pytest_asyncio.fixture(loop_scope="session")
async def api_client(
api_settings: Settings, db_sessionmaker: async_sessionmaker[AsyncSession]
) -> AsyncIterator[AsyncClient]:
app = create_app(api_settings)
app.dependency_overrides[get_sessionmaker] = lambda: db_sessionmaker
async with (
LifespanManager(app) as manager,
AsyncClient(transport=ASGITransport(app=manager.app), base_url="http://test") as client,
):
yield client
def _auth(token: str) -> dict[str, str]:
return {"Authorization": f"Bearer {token}"}
async def _seed(
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
specs: tuple[SeedSpec, ...],
) -> None:
await ensure_chunks_collection(qdrant_client, collection=qdrant_settings.collection)
storage = QdrantPointStorage(qdrant_client, collection=qdrant_settings.collection)
await storage.upsert_points([chunk_point_for(spec) for spec in specs])
async def _tenant_with_key(
db_session: AsyncSession, *, scopes: list[str] | None = None
) -> tuple[uuid.UUID, str]:
tenant = await create_tenant(db_session)
_, token = await create_api_key(
db_session, tenant=tenant, scopes=scopes if scopes is not None else ["points:read"]
)
await db_session.commit()
return tenant.id, token
async def test_list_points_returns_a_files_points_in_order_id_order(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
tenant_id, token = await _tenant_with_key(db_session)
await _seed(
qdrant_client,
qdrant_settings,
(
SeedSpec(tenant_id, FILE_ID, 0, "first", 1.0),
SeedSpec(tenant_id, FILE_ID, 2, "third", 3.0),
SeedSpec(tenant_id, FILE_ID, 1, "second", 2.0),
),
)
response = await api_client.get(
"/v1/points", params={"file_id": str(FILE_ID)}, headers=_auth(token)
)
assert response.status_code == 200
assert [point["content"] for point in response.json()["points"]] == [
"first",
"second",
"third",
]
async def test_list_points_omits_the_tenant_id_from_the_response(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
"""Echoing `tenant_id` back invites clients to start sending it."""
tenant_id, token = await _tenant_with_key(db_session)
await _seed(qdrant_client, qdrant_settings, (SeedSpec(tenant_id, FILE_ID, 0, "only", 1.0),))
response = await api_client.get(
"/v1/points", params={"file_id": str(FILE_ID)}, headers=_auth(token)
)
assert "tenant_id" not in response.json()["points"][0]
async def test_get_point_returns_404_for_another_tenants_point(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
"""`404`, not `403`: the API must not confirm the id exists (ADR-0016)."""
owner_id, _ = await _tenant_with_key(db_session)
_, intruder_token = await _tenant_with_key(db_session)
await _seed(qdrant_client, qdrant_settings, (SeedSpec(owner_id, FILE_ID, 0, "secret", 1.0),))
response = await api_client.get(
f"/v1/points/{point_id_for(FILE_ID, 0)}", headers=_auth(intruder_token)
)
assert response.status_code == 404
assert response.json()["error"]["code"] == "not_found"
async def test_list_points_returns_nothing_for_another_tenants_file(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
owner_id, _ = await _tenant_with_key(db_session)
_, intruder_token = await _tenant_with_key(db_session)
await _seed(qdrant_client, qdrant_settings, (SeedSpec(owner_id, FILE_ID, 0, "secret", 1.0),))
response = await api_client.get(
"/v1/points", params={"file_id": str(FILE_ID)}, headers=_auth(intruder_token)
)
assert response.status_code == 200
assert response.json()["points"] == []
async def test_list_points_excludes_inactive_points_unless_asked(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
tenant_id, token = await _tenant_with_key(db_session)
await _seed(
qdrant_client,
qdrant_settings,
(
SeedSpec(tenant_id, FILE_ID, 0, "live", 1.0),
SeedSpec(tenant_id, FILE_ID, 1, "deleted", 2.0, is_active=False),
),
)
default = await api_client.get(
"/v1/points", params={"file_id": str(FILE_ID)}, headers=_auth(token)
)
opted_in = await api_client.get(
"/v1/points",
params={"file_id": str(FILE_ID), "include_inactive": True},
headers=_auth(token),
)
assert [point["content"] for point in default.json()["points"]] == ["live"]
assert [point["content"] for point in opted_in.json()["points"]] == ["live", "deleted"]
async def test_list_points_pagination_does_not_skip_or_repeat_under_a_concurrent_insert(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
"""The defect an offset cursor would have.
Between page one and page two, a point is inserted *behind* the cursor
(`order_id` 0.5, ahead of everything already returned). An offset-based
cursor would shift every later row down by one and hand back `b` a second
time. A value-based cursor is unaffected: the reader has passed 2.0 and the
new row is behind it.
"""
tenant_id, token = await _tenant_with_key(db_session)
await _seed(
qdrant_client,
qdrant_settings,
(
SeedSpec(tenant_id, FILE_ID, 0, "a", 1.0),
SeedSpec(tenant_id, FILE_ID, 1, "b", 2.0),
SeedSpec(tenant_id, FILE_ID, 2, "c", 3.0),
SeedSpec(tenant_id, FILE_ID, 3, "d", 4.0),
),
)
first = await api_client.get(
"/v1/points",
params={"file_id": str(FILE_ID), "limit": 2},
headers=_auth(token),
)
await _seed(qdrant_client, qdrant_settings, (SeedSpec(tenant_id, FILE_ID, 9, "inserted", 0.5),))
second = await api_client.get(
"/v1/points",
params={
"file_id": str(FILE_ID),
"limit": 2,
"cursor": first.json()["next_cursor"],
},
headers=_auth(token),
)
assert [point["content"] for point in first.json()["points"]] == ["a", "b"]
assert [point["content"] for point in second.json()["points"]] == ["c", "d"]
async def test_count_points_respects_tenant_and_domain_filters(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
tenant_id, token = await _tenant_with_key(db_session)
other_id, _ = await _tenant_with_key(db_session)
await _seed(
qdrant_client,
qdrant_settings,
(
SeedSpec(tenant_id, FILE_ID, 0, "one", 1.0),
SeedSpec(tenant_id, OTHER_FILE_ID, 0, "two", 1.0, domain="ops"),
SeedSpec(other_id, FILE_ID, 5, "theirs", 1.0),
),
)
everything = await api_client.get("/v1/points/count", headers=_auth(token))
scoped = await api_client.get(
"/v1/points/count", params={"domain": "ops"}, headers=_auth(token)
)
assert everything.json()["count"] == 2
assert scoped.json()["count"] == 1
async def test_search_points_matches_content_within_the_tenant_only(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
tenant_id, token = await _tenant_with_key(db_session)
other_id, _ = await _tenant_with_key(db_session)
await _seed(
qdrant_client,
qdrant_settings,
(
SeedSpec(tenant_id, FILE_ID, 0, "annual leave policy", 1.0),
SeedSpec(tenant_id, FILE_ID, 1, "expense policy", 2.0),
SeedSpec(tenant_id, FILE_ID, 2, "office address", 3.0),
SeedSpec(other_id, OTHER_FILE_ID, 0, "their leave policy", 1.0),
),
)
response = await api_client.get(
"/v1/points/search", params={"q": "policy"}, headers=_auth(token)
)
body = response.json()
assert body["query"] == "policy"
assert {point["content"] for point in body["points"]} == {
"annual leave policy",
"expense policy",
}
async def test_search_points_folds_arabic_letterforms_in_the_query(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
"""A query typed on an Arabic keyboard must match Persian-folded content.
Ingestion folds U+064A/U+0643 to U+06CC/U+06A9 before storing, so an
unfolded query would match nothing — an empty result rather than an error,
which is exactly the kind of silent failure nobody reports as a bug.
Codepoints, not literals: the two forms are visually identical here.
"""
persian = "مدیریت" # stored, Persian yeh
arabic = "مديريت" # queried, Arabic yeh
assert persian != arabic
tenant_id, token = await _tenant_with_key(db_session)
await _seed(qdrant_client, qdrant_settings, (SeedSpec(tenant_id, FILE_ID, 0, persian, 1.0),))
response = await api_client.get("/v1/points/search", params={"q": arabic}, headers=_auth(token))
assert [point["content"] for point in response.json()["points"]] == [persian]
async def test_point_read_routes_require_the_points_read_scope(
api_client: AsyncClient, db_session: AsyncSession
) -> None:
"""An upload key must not double as a way to read every chunk."""
_, token = await _tenant_with_key(db_session, scopes=["files:write"])
response = await api_client.get(
"/v1/points", params={"file_id": str(FILE_ID)}, headers=_auth(token)
)
assert response.status_code == 403
assert response.json()["error"]["code"] == "missing_scope"
async def test_count_and_search_paths_are_not_parsed_as_point_ids(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
"""Route order regression guard.
`/{point_id}` declared before `/count` would make this `422` — Pydantic
failing to read "count" as a UUID — rather than a count. The bug is one
tidy-up reorder away, and nothing else in the suite would catch it.
Seeds nothing but still bootstraps the collection: querying a collection
that does not exist is a `500`, which `/readyz` is responsible for
preventing, and which would mask the `422` this test is watching for.
"""
_, token = await _tenant_with_key(db_session)
await _seed(qdrant_client, qdrant_settings, ())
count = await api_client.get("/v1/points/count", headers=_auth(token))
search = await api_client.get(
"/v1/points/search", params={"q": "anything"}, headers=_auth(token)
)
assert count.status_code == 200
assert search.status_code == 200
async def test_file_points_route_lists_the_same_points_as_the_query_form(
api_client: AsyncClient,
db_session: AsyncSession,
qdrant_client: AsyncQdrantClient,
qdrant_settings: QdrantSettings,
) -> None:
tenant_id, token = await _tenant_with_key(db_session)
await _seed(
qdrant_client,
qdrant_settings,
(
SeedSpec(tenant_id, FILE_ID, 0, "first", 1.0),
SeedSpec(tenant_id, FILE_ID, 1, "second", 2.0),
),
)
by_query = await api_client.get(
"/v1/points", params={"file_id": str(FILE_ID)}, headers=_auth(token)
)
by_path = await api_client.get(f"/v1/files/{FILE_ID}/points", headers=_auth(token))
assert by_path.status_code == 200
assert by_path.json() == by_query.json()

View File

@@ -0,0 +1,353 @@
"""One behavioural contract, run against both `PointRepository` implementations.
Plan 002 Phase 1's exit criterion is that the Qdrant adapter and the in-memory
fake agree on filtering, ordering, and pagination. Asserting that by writing two
parallel test files invites them to drift — the usual outcome is a fake that
quietly grows more forgiving than the real store, so unit tests keep passing
while production does something else.
So the scenarios live here once. `tests/unit/application/points/` runs them
against `FakePointRepository`; `tests/integration/qdrant/` runs the identical
functions against real Qdrant. A divergence fails one of the two runs rather
than hiding.
Each scenario takes an already-seeded repository plus the ids it was seeded
with, so seeding (which genuinely differs — one writes a dict, the other upserts
vectors) stays outside the shared code.
"""
import uuid
from dataclasses import dataclass
from datetime import UTC, datetime
from src.application.ingestion.chunking import chunk_id_for
from src.application.ingestion.models import SparseVector
from src.application.points.models import ChunkPoint
from src.application.points.point import Point
from src.application.ports.point_repository import PayloadPatch, PointRepository
from src.infrastructure.qdrant.collection import (
DENSE_NOMIC_DIMENSIONS,
DENSE_OPENAI_DIMENSIONS,
)
# Two tenants and two files, so every scenario can assert both that the right
# rows come back and that the wrong ones do not.
SEEDED_AT = datetime(2026, 1, 1, 12, 0, tzinfo=UTC)
@dataclass(frozen=True)
class SeedSpec:
"""One point to seed. Deliberately flat — the fixture builds the model."""
tenant_id: uuid.UUID
file_id: uuid.UUID
chunk_index: int
content: str
order_id: float
is_active: bool = True
version: int = 1
domain: str = "support"
def point_id_for(file_id: uuid.UUID, chunk_index: int) -> uuid.UUID:
"""The deterministic point id (ADR-0001), re-exported for readability."""
return chunk_id_for(file_id, chunk_index)
def build_point(spec: SeedSpec) -> Point:
point_id = point_id_for(spec.file_id, spec.chunk_index)
return Point(
point_id=point_id,
tenant_id=spec.tenant_id,
domain=spec.domain,
file_id=spec.file_id,
chunk_id=point_id,
content=spec.content,
content_type="paragraph",
source_filename="handbook.docx",
source_type="docx",
order_id=spec.order_id,
chunk_index=spec.chunk_index,
previous_chunk_id=None,
next_chunk_id=None,
is_active=spec.is_active,
deleted_at=None if spec.is_active else SEEDED_AT,
created_at=SEEDED_AT,
updated_at=SEEDED_AT,
created_by="ingestion",
updated_by="ingestion",
version=spec.version,
content_hash="0" * 64,
embedding_model_version="test-model",
)
def chunk_point_for(spec: SeedSpec) -> ChunkPoint:
"""A `SeedSpec` as something upsertable into real Qdrant.
The payload is taken straight off `build_point`'s read model rather than
hand-written, so the write shape and the read shape cannot drift apart. The
vectors are constant filler: nothing in plan 002's read paths scores by
similarity, so their values are irrelevant and their dimensions are not.
"""
point = build_point(spec)
return ChunkPoint(
point_id=point.point_id,
dense={
"dense_nomic": [0.1] * DENSE_NOMIC_DIMENSIONS,
"dense_openai": [0.2] * DENSE_OPENAI_DIMENSIONS,
},
sparse=SparseVector(indices=[1, 2], values=[0.5, 0.25]),
payload=point.model_dump(mode="json", exclude={"point_id", "vectors"}),
)
def seed_specs(
tenant_a: uuid.UUID,
tenant_b: uuid.UUID,
file_a: uuid.UUID,
file_b: uuid.UUID,
file_c: uuid.UUID,
) -> tuple[SeedSpec, ...]:
"""The fixed corpus every scenario below assumes.
Ordering is deliberately not insertion order: `order_id` 3.0 is seeded
before 2.0 so a repository that returns rows in write order fails
`scenario_list_orders_by_order_id` instead of passing by accident.
Tenant B owns its own file. Sharing a `file_id` across tenants would be a
fixture that cannot occur: point ids are derived from `file_id` plus
`chunk_index` alone, so two tenants in one file would collide on a single
id and the corpus would be testing an impossible state.
"""
return (
SeedSpec(tenant_a, file_a, 0, "onboarding checklist for new staff", 1.0),
SeedSpec(tenant_a, file_a, 2, "payroll schedule and bank details", 3.0),
SeedSpec(tenant_a, file_a, 1, "expense policy and receipts", 2.0),
SeedSpec(tenant_a, file_a, 3, "retired parking policy", 4.0, is_active=False),
SeedSpec(tenant_a, file_b, 0, "supplier onboarding contacts", 1.0, domain="ops"),
SeedSpec(tenant_b, file_c, 0, "other tenant onboarding secrets", 1.0),
)
async def scenario_get_returns_point_within_tenant(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
point = await repository.get(tenant_id=tenant_a, point_id=point_id_for(file_a, 0))
assert point is not None
assert point.content == "onboarding checklist for new staff"
assert point.chunk_index == 0
assert point.order_id == 1.0
async def scenario_get_hides_other_tenants_point(
repository: PointRepository, *, tenant_b: uuid.UUID, file_a: uuid.UUID
) -> None:
"""Tenant B probing a real point id that belongs to tenant A gets nothing.
The point exists, so this distinguishes a tenant filter from a plain
existence check — the difference between `404` and a leak.
"""
assert await repository.get(tenant_id=tenant_b, point_id=point_id_for(file_a, 0)) is None
async def scenario_get_reaches_inactive_points(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
"""A direct `get` is not subject to the implied `is_active` read filter.
Soft delete and relinking both address points by id and must be able to see
a deactivated one; only listings hide them.
"""
point = await repository.get(tenant_id=tenant_a, point_id=point_id_for(file_a, 3))
assert point is not None
assert point.is_active is False
async def scenario_get_many_returns_only_this_tenants_points(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
found = await repository.get_many(
tenant_id=tenant_a,
point_ids=[point_id_for(file_a, 0), point_id_for(file_a, 1), uuid.uuid4()],
)
assert {point.chunk_index for point in found} == {0, 1}
async def scenario_list_orders_by_order_id(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
page = await repository.list_by_file(tenant_id=tenant_a, file_id=file_a, limit=10)
assert [point.order_id for point in page.points] == [1.0, 2.0, 3.0]
assert page.next_cursor is None
async def scenario_list_excludes_inactive_by_default(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
default = await repository.list_by_file(tenant_id=tenant_a, file_id=file_a, limit=10)
assert 4.0 not in [point.order_id for point in default.points]
opted_in = await repository.list_by_file(
tenant_id=tenant_a, file_id=file_a, limit=10, include_inactive=True
)
assert [point.order_id for point in opted_in.points] == [1.0, 2.0, 3.0, 4.0]
async def scenario_list_paginates_without_skipping_or_repeating(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
first = await repository.list_by_file(tenant_id=tenant_a, file_id=file_a, limit=2)
assert [point.order_id for point in first.points] == [1.0, 2.0]
assert first.next_cursor is not None
second = await repository.list_by_file(
tenant_id=tenant_a, file_id=file_a, limit=2, cursor=first.next_cursor
)
assert [point.order_id for point in second.points] == [3.0]
assert second.next_cursor is None
async def scenario_list_is_scoped_to_one_file(
repository: PointRepository, *, tenant_a: uuid.UUID, file_b: uuid.UUID
) -> None:
page = await repository.list_by_file(tenant_id=tenant_a, file_id=file_b, limit=10)
assert [point.content for point in page.points] == ["supplier onboarding contacts"]
async def scenario_count_respects_tenant_and_filters(
repository: PointRepository,
*,
tenant_a: uuid.UUID,
tenant_b: uuid.UUID,
file_a: uuid.UUID,
) -> None:
assert await repository.count(tenant_id=tenant_a) == 4
assert await repository.count(tenant_id=tenant_a, include_inactive=True) == 5
assert await repository.count(tenant_id=tenant_a, file_id=file_a) == 3
assert await repository.count(tenant_id=tenant_a, domain="ops") == 1
assert await repository.count(tenant_id=tenant_b) == 1
async def scenario_keyword_search_matches_content_within_tenant(
repository: PointRepository, *, tenant_a: uuid.UUID, tenant_b: uuid.UUID
) -> None:
found = await repository.keyword_search(tenant_id=tenant_a, query="onboarding", limit=10)
assert {point.content for point in found.points} == {
"onboarding checklist for new staff",
"supplier onboarding contacts",
}
# The same term matches a different tenant's point, which must not leak.
other = await repository.keyword_search(tenant_id=tenant_b, query="onboarding", limit=10)
assert {point.content for point in other.points} == {"other tenant onboarding secrets"}
async def scenario_keyword_search_excludes_inactive_by_default(
repository: PointRepository, *, tenant_a: uuid.UUID
) -> None:
""" "retired parking policy" is seeded inactive, so only the opt-in sees it."""
default = await repository.keyword_search(tenant_id=tenant_a, query="parking", limit=10)
assert default.points == ()
opted_in = await repository.keyword_search(
tenant_id=tenant_a, query="parking", limit=10, include_inactive=True
)
assert {point.content for point in opted_in.points} == {"retired parking policy"}
async def scenario_keyword_search_filters_by_domain(
repository: PointRepository, *, tenant_a: uuid.UUID
) -> None:
found = await repository.keyword_search(
tenant_id=tenant_a, query="onboarding", limit=10, domain="ops"
)
assert {point.content for point in found.points} == {"supplier onboarding contacts"}
async def scenario_patch_applies_when_version_matches(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
target = point_id_for(file_a, 0)
await repository.apply_patches(
tenant_id=tenant_a,
patches=[
PayloadPatch(
point_id=target,
payload={"content": "amended onboarding checklist", "version": 2},
expected_version=1,
)
],
)
point = await repository.get(tenant_id=tenant_a, point_id=target)
assert point is not None
assert point.content == "amended onboarding checklist"
assert point.version == 2
async def scenario_patch_is_a_noop_when_version_is_stale(
repository: PointRepository, *, tenant_a: uuid.UUID, file_a: uuid.UUID
) -> None:
"""A stale guard must match nothing rather than raise or clobber.
This is the behaviour the service's post-apply check depends on: Qdrant
reports success for a filtered `set_payload` that matched zero points, so
"did my write land?" is only answerable by reading back.
"""
target = point_id_for(file_a, 1)
await repository.apply_patches(
tenant_id=tenant_a,
patches=[
PayloadPatch(
point_id=target,
payload={"content": "should not land"},
expected_version=99,
)
],
)
point = await repository.get(tenant_id=tenant_a, point_id=target)
assert point is not None
assert point.content == "expense policy and receipts"
async def scenario_patch_cannot_cross_tenants(
repository: PointRepository, *, tenant_a: uuid.UUID, tenant_b: uuid.UUID, file_a: uuid.UUID
) -> None:
"""Tenant B patching tenant A's point id must change nothing.
An unguarded patch (no `expected_version`) addressed by a known-good id is
the sharpest form of the probe: only the tenant condition can stop it, so a
repository that filtered on id alone would rewrite another tenant's row.
"""
target = point_id_for(file_a, 0)
await repository.apply_patches(
tenant_id=tenant_b,
patches=[PayloadPatch(point_id=target, payload={"content": "hijacked"})],
)
victim = await repository.get(tenant_id=tenant_a, point_id=target)
assert victim is not None
assert victim.content == "onboarding checklist for new staff"
# Scenarios that only read, safe to run in any order against one seeded corpus.
READ_SCENARIOS = (
scenario_get_returns_point_within_tenant,
scenario_get_hides_other_tenants_point,
scenario_get_reaches_inactive_points,
scenario_get_many_returns_only_this_tenants_points,
scenario_list_orders_by_order_id,
scenario_list_excludes_inactive_by_default,
scenario_list_paginates_without_skipping_or_repeating,
scenario_list_is_scoped_to_one_file,
scenario_count_respects_tenant_and_filters,
scenario_keyword_search_matches_content_within_tenant,
scenario_keyword_search_excludes_inactive_by_default,
scenario_keyword_search_filters_by_domain,
)
# Scenarios that mutate, and therefore need a freshly seeded corpus each.
WRITE_SCENARIOS = (
scenario_patch_applies_when_version_matches,
scenario_patch_is_a_noop_when_version_is_stale,
scenario_patch_cannot_cross_tenants,
)

View File

@@ -0,0 +1,57 @@
"""`FakePointRepository` against the shared `PointRepository` contract.
The same scenarios run against real Qdrant in
`tests/integration/qdrant/test_point_repository.py`. Keeping both runs green is
what stops the fake from drifting into something more permissive than the store
it stands in for — see `tests/support/point_contract.py`.
"""
import uuid
from collections.abc import Awaitable, Callable
import pytest
from src.application.ports.point_repository import PointRepository
from tests.fakes import FakePointRepository
from tests.support import point_contract
from tests.support.point_contract import build_point, seed_specs
pytestmark = [pytest.mark.unit, pytest.mark.asyncio]
type Scenario = Callable[..., Awaitable[None]]
TENANT_A = uuid.UUID("11111111-1111-4111-8111-111111111111")
TENANT_B = uuid.UUID("22222222-2222-4222-8222-222222222222")
FILE_A = uuid.UUID("33333333-3333-4333-8333-333333333333")
FILE_B = uuid.UUID("44444444-4444-4444-8444-444444444444")
FILE_C = uuid.UUID("55555555-5555-4555-8555-555555555555")
IDS = {
"tenant_a": TENANT_A,
"tenant_b": TENANT_B,
"file_a": FILE_A,
"file_b": FILE_B,
"file_c": FILE_C,
}
def _seeded() -> FakePointRepository:
repository = FakePointRepository()
for spec in seed_specs(TENANT_A, TENANT_B, FILE_A, FILE_B, FILE_C):
repository.add(build_point(spec))
return repository
def _arguments(scenario: Scenario) -> dict[str, uuid.UUID]:
"""Pass only the ids a scenario declares, so each one names its own inputs."""
return {name: value for name, value in IDS.items() if name in scenario.__annotations__}
@pytest.mark.parametrize(
"scenario",
[*point_contract.READ_SCENARIOS, *point_contract.WRITE_SCENARIOS],
ids=lambda scenario: scenario.__name__.removeprefix("scenario_"),
)
async def test_fake_point_repository_satisfies_the_contract(scenario: Scenario) -> None:
repository: PointRepository = _seeded()
await scenario(repository, **_arguments(scenario))

View File

@@ -0,0 +1,103 @@
"""`application/points/queries.py` against the fake repository (ADR-0016).
The HTTP-level behaviour is covered in `tests/integration/qdrant/`. What is
worth isolating here is the small amount of logic the service actually owns —
raising instead of returning `None`, and folding a query's Arabic letterforms —
because both are cheap to test directly and expensive to notice when broken.
"""
import uuid
import pytest
from src.application.points.point import PointNotFoundError
from src.application.points.queries import (
count_points,
get_point,
list_file_points,
search_points,
)
from tests.fakes import FakePointRepository
from tests.support.point_contract import SeedSpec, build_point, point_id_for
pytestmark = [pytest.mark.unit, pytest.mark.asyncio]
TENANT = uuid.UUID("11111111-1111-4111-8111-111111111111")
OTHER_TENANT = uuid.UUID("22222222-2222-4222-8222-222222222222")
FILE = uuid.UUID("33333333-3333-4333-8333-333333333333")
def _repository(*specs: SeedSpec) -> FakePointRepository:
repository = FakePointRepository()
for spec in specs:
repository.add(build_point(spec))
return repository
async def test_get_point_returns_the_point_for_its_owner() -> None:
repository = _repository(SeedSpec(TENANT, FILE, 0, "hello", 1.0))
point = await get_point(repository, tenant_id=TENANT, point_id=point_id_for(FILE, 0))
assert point.content == "hello"
async def test_get_point_raises_for_another_tenants_point() -> None:
"""Raising, not returning `None`, so a route cannot serve `200 null`."""
repository = _repository(SeedSpec(TENANT, FILE, 0, "hello", 1.0))
with pytest.raises(PointNotFoundError):
await get_point(repository, tenant_id=OTHER_TENANT, point_id=point_id_for(FILE, 0))
async def test_get_point_raises_for_an_unknown_point() -> None:
repository = _repository()
with pytest.raises(PointNotFoundError):
await get_point(repository, tenant_id=TENANT, point_id=uuid.uuid4())
async def test_search_points_folds_arabic_letterforms_before_matching() -> None:
"""The query is normalized the same way ingestion normalized the content.
Stored content carries Persian yeh (U+06CC) because `normalize_persian_text`
folded it at ingest. A query typed with Arabic yeh (U+064A) is a different
codepoint and would match nothing — silently, with no error to notice.
"""
persian = "مدیریت"
arabic = "مديريت"
assert persian != arabic
repository = _repository(SeedSpec(TENANT, FILE, 0, persian, 1.0))
found = await search_points(repository, tenant_id=TENANT, query=arabic, limit=10)
assert [point.content for point in found.points] == [persian]
async def test_search_points_does_not_reach_another_tenant() -> None:
repository = _repository(SeedSpec(OTHER_TENANT, FILE, 0, "secret", 1.0))
found = await search_points(repository, tenant_id=TENANT, query="secret", limit=10)
assert found.points == ()
async def test_list_file_points_returns_an_empty_page_for_a_foreign_file() -> None:
"""Indistinguishable from an unknown file — the same non-disclosure rule."""
repository = _repository(SeedSpec(OTHER_TENANT, FILE, 0, "secret", 1.0))
page = await list_file_points(repository, tenant_id=TENANT, file_id=FILE, limit=10)
assert page.points == ()
assert page.next_cursor is None
async def test_count_points_excludes_inactive_unless_asked() -> None:
repository = _repository(
SeedSpec(TENANT, FILE, 0, "live", 1.0),
SeedSpec(TENANT, FILE, 1, "gone", 2.0, is_active=False),
)
assert await count_points(repository, tenant_id=TENANT) == 1
assert await count_points(repository, tenant_id=TENANT, include_inactive=True) == 2