import { OnModuleInit } from "@nestjs/common"; import { ClsService } from "nestjs-cls"; import { AiStatus } from "../../../common/enums/ai.status"; import { AgentScopeFilterService } from "../../../common/repositories/agent-scope.filter"; import { AiSourceQueryProvider } from "../../../common/repositories/ai-source-query.provider"; import { DataLimits } from "../../../common/types/data.limits"; import { EmbedderAttribution, EmbedderService } from "../../../core"; import { ModelService } from "../../../core/llm/services/model.service"; import { Neo4jService } from "../../../core/neo4j/services/neo4j.service"; import { SecurityService } from "../../../core/security/services/security.service"; import { Chunk } from "../../chunk/entities/chunk.entity"; export declare class ChunkRepository implements OnModuleInit { private readonly neo4j; private readonly modelService; private readonly embedderService; private readonly clsService; private readonly securityService; private readonly aiSourceQuery; private readonly agentScope; /** * Chunks embedded (and written) per round-trip in `enrichContentAndEmbedBatch`. * A whole document's worth of vectors held at once is what pushed the worker past * its heap on large uploads: 760 pages ≈ 2.5k chunks × 3072 floats. Slicing bounds * the peak to one slice's vectors while keeping batching's latency win (one embedder * call per 50 chunks, not per chunk). */ private static readonly EMBED_SLICE; constructor(neo4j: Neo4jService, modelService: ModelService, embedderService: EmbedderService, clsService: ClsService, securityService: SecurityService, aiSourceQuery: AiSourceQueryProvider, agentScope: AgentScopeFilterService); onModuleInit(): Promise; recreateVectorIndex(): Promise; findAllChunks(): Promise; updateEmbedding(params: { chunkId: string; embedding: number[]; }): Promise; /** * `attribution` is OPTIONAL and opt-in: this is a QUERY-time embedding, so * the repository has no entity of its own to bill it to. The retrieval scope * lives with the caller (the contextualiser knows which content the run is * bound to), which is why it is passed down rather than derived here. Absent, * `EmbedderService.persistUsage` records nothing. * * Each returned chunk carries `score`: the cosine similarity of that chunk against * the question embedding. Both retrieval halves have to reach the notebook on ONE * scale — the answer node orders entries best-score-first and fills a character * budget, so an unscored half sorts last however good it is. The RRF score is NOT * usable for that: it is rank-derived (the top hit scores ≈1/61 whether it is a * perfect match or noise) and is not comparable with a cosine from the graph half. */ findPotentialChunks(params: { question: string; dataLimits: DataLimits; attribution?: EmbedderAttribution; /** * Precomputed question embedding. When supplied the repository does NOT * embed again — the same question is otherwise embedded twice per turn, * once here and once in findPotentialKeyConcepts. Optional so existing * callers are unaffected. */ queryEmbedding?: number[]; }): Promise>; /** * Exact cosine over the scoped set. No recall cliff by construction: every * in-scope chunk is scored, so a small tenant can never be crowded out of a * global top-K by a large one. */ private vectorIdsByExactScan; /** * Index-backed fallback for scoped sets too large to scan exactly. Still * post-filters, so it still has a cliff — the over-fetch is what pushes that * cliff out of reach, and it is deliberately far above the previous 1,000. */ private vectorIdsByIndex; /** * Lexical branch, scoped the same way. Previously filtered against the * client-side id list; now joined in the database like the vector branch. */ private lexicalIdsInScope; /** * INPUT ORDER IS PART OF THE CONTRACT here too: `WHERE chunk.id IN $ids` yields rows * in store order, and the fused RRF order is what the caller means by "best first". * * `queryEmbedding` attaches the cosine score of each chunk against the question, on * the same scale `findChunksByIds` puts on the graph half. No floor and no count cap * are applied — the notebook's character budget decides what reaches the answer. */ private findChunksByIdsOrdered; /** * Attaches `vector.similarity.cosine(chunk.embedding, $queryEmbedding)` to already * hydrated chunks, by id, in JS. * * Why a separate read rather than one extra column on the hydration query: the * hydration goes through `readMany`, which maps every row with the descriptor's * generated mapper, and that mapper only ever reads the descriptor's own fields, * computed fields and virtual fields (`define-entity.ts`). A projected column such * as `score` is therefore silently DROPPED before the caller ever sees it. Reading * the id/score pairs on their own and merging them here keeps entity mapping on * `readMany` — this repository never hand-maps raw Neo4j records. * * SCOPE: this read scores only ids that the caller's own scope-gated query already * returned, so it cannot widen the scope by construction. */ private attachCosineScores; findParentName(params: { id: string; nodeType: string; }): Promise; findSubsequentChunkId(params: { chunkId: string; }): Promise; findPreviousChunkId(params: { chunkId: string; }): Promise; /** * `dataLimits` is optional so existing callers keep working, but the * contextualiser MUST pass it: the ids reaching this method come from an LLM * choosing among the chunks it was shown. That is a soft constraint — the * model can echo an id it saw in an earlier hop, or hallucinate one outright * — so the scope root is re-checked here rather than trusted from upstream. */ findChunkById(params: { chunkId: string; dataLimits?: DataLimits; }): Promise; /** * Hydrates many chunks in ONE query, with the same scope gate `findChunkById` * applies to one. The caller previously looped with an `await` inside, * costing one round trip per queued chunk. Returns only the chunks that * exist and are in scope; the caller must not assume a 1:1 mapping with * `chunkIds`. * * ORDER IS PART OF THE CONTRACT. `WHERE chunk.id IN $chunkIds` returns rows in * whatever order the store yields them, but the loop this replaced hydrated in * queue order, and that order reaches the contextualiser's per-chunk fan-out and * the notebook entries it writes. Returning them shuffled changes what the answer * node cites. The input order is therefore restored here, exactly as * `findChunksByIdsOrdered` does for the retrieval path. * * `queryEmbedding` is OPTIONAL and, when given, attaches `score`: the cosine * similarity of the chunk against the question. These chunks arrive from a fact * join with no score of their own, which is precisely why an LLM used to have to * judge them; cosine puts this half on the SAME scale as the document half, so the * answer node can order both together. When it is absent the scoring clause is not * in the Cypher at all. No floor is applied at any point: the notebook's character * budget, not a threshold, decides what reaches the answer. */ findChunksByIds(params: { chunkIds: string[]; dataLimits?: DataLimits; queryEmbedding?: number[]; }): Promise>; /** * Pipeline hydration of a content node's chunks — deliberately WITHOUT `embedding`. * * A :Chunk node carries a full embedding vector (3072 floats ≈ 24 KB raw, far more * once the driver has boxed it). Returning the whole node made every pipeline guard * pull the entire document's vectors into the worker heap even though NO consumer * reads `chunk.embedding` — that is what killed the worker on a 760-page matrix. * * The rows are therefore returned as an explicit `{ labels, properties }` map instead * of a Node: `EntityFactory.createOrMerge` treats any column with a `labels` key as a * node and maps `data.properties`, so a hand-built map of the same shape hydrates * identically (a bare map projection would NOT — it has no `labels`/`properties` and * the factory would drop the row). * * The projected property list is exactly the descriptor's own fields (minus * `embedding`) plus the `Entity` base fields, so nothing is lost: the descriptor's * auto-generated mapper only ever reads those keys anyway. */ findChunks(params: { id: string; nodeType: string; }): Promise; createChunk(params: { id: string; nodeId: string; nodeType: string; previousChunkId?: string; content: string; heading?: string; imagePath?: string; position: number; }): Promise; updateStatus(params: { id: string; aiStatus: AiStatus; }): Promise; updateDates(params: { chunkId: string; dates: string; }): Promise; /** * `attribution` is OPTIONAL (added second, positionally, so existing callers * keep compiling). One usage record is written per embedded slice — see * `EmbedderService.persistUsage`. That is honest here because every item in * the batch is a chunk of the SAME parent entity: the only caller, * `ChunkService.propagateAndEmbedDates`, builds `items` from * `findChunks({ id, nodeType })`. */ enrichContentAndEmbedBatch(items: { chunkId: string; enrichedContent: string; propagatedDates?: string; }[], attribution?: EmbedderAttribution): Promise; markChunksCompleted(params: { id: string; nodeType: string; }): Promise; /** * Count of chunks not yet completed for a content node. Replaces full hydration in * pipeline guards: they only ever asked "are any chunks still pending?", and hydrating * every chunk (embedding vectors included) to answer it is what made the guard * quadratic in a document's chunk count. * * Returns a scalar, so it cannot go through `readOne`/`readMany` (those map entity * columns). Same raw-scalar read as `TokenUsageRepository.findUsageSummary` — the * package's precedent for non-entity aggregate reads. */ countChunksInProgress(params: { id: string; nodeType: string; }): Promise; /** * Single-winner claim on a content node's post-chunking pipeline. * * `countChunksInProgress` alone CANNOT gate finalisation. It answers "are any chunks * still pending?", never "has finalisation already run?", and `ChunkService.generateGraph` * enqueues one finalise job per chunk — so a 25-chunk document gets 25 jobs. Once the last * chunk lands, every remaining queued job passes a pending-count check and re-runs the whole * pipeline (measured 2026-08-12: summariser ran 25× on one document, 100.91 credits; ~70% of * that run's spend was duplicated work). * * The `SET n.finalisationClaimedAt = n.finalisationClaimedAt` is a deliberate no-op write: it * takes the node's write lock BEFORE the predicate is evaluated. A concurrent claimer blocks * there, and re-reads the committed value afterwards (Neo4j reads are read-committed, not * repeatable), so it sees the winner's timestamp and its `WHERE` fails. Without the lock both * transactions evaluate the predicate on the pre-write state and both claim. * * Re-ingestion re-arms the claim via `clearFinalisationClaim` in `ChunkService.createChunks`. */ claimContentFinalisation(params: { id: string; nodeType: string; }): Promise; /** * Re-arms `claimContentFinalisation` for a content node whose chunks are being (re)built. * Called from `ChunkService.createChunks`, the single entry point every ingestion and * rebuild path goes through — without it a re-ingested document would keep the stale claim * and skip finalisation silently. */ clearFinalisationClaim(params: { id: string; nodeType: string; }): Promise; getChunksInProgress(params: { id: string; nodeType: string; }): Promise; createNextRelationship(params: { chunkId: string; nextChunkId: string; }): Promise; deleteChunks(params: { chunkIds: string[]; }): Promise; deleteDisconnectedChunks(): Promise; deleteChunksByNodeType(params: { id: string; nodeType: string; }): Promise; findChunkNeighbors(params: { chunkIds: string[]; window: number; }): Promise<{ chunkId: string; before: string[]; after: string[]; }[]>; findChunkByContentIdAndType(params: { id: string; type: string; }): Promise; } //# sourceMappingURL=chunk.repository.d.ts.map