Files
krow_backend/go-api/internal/knowledge/ingest.go
2026-08-28 12:21:44 +05:30

447 lines
16 KiB
Go

package knowledge
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"strings"
"github.com/jackc/pgx/v5"
"github.com/krow/krow-backend/go-api/internal/repo"
)
// Ingest: getting a document into the index, permissioned.
//
// The gate here is §5's — "chunks without ACL metadata are rejected at ingest"
// — and it is a gate rather than a default because the alternative fails
// silently in both directions. Defaulting to tenant-wide over-shares a document
// somebody meant to restrict; defaulting to empty indexes it into invisibility.
// Neither raises anything. So an ingest that does not say who a document is for
// is refused, and the caller has to decide.
// Document is what an ingester supplies.
type Document struct {
// Source is the corpus. An agent spec names sources, and retrieval filters
// by them, so this is part of the permission story: an agent granted the
// policy library does not thereby gain the incident log.
Source string
// ExternalID is this document's id in wherever it came from. A re-ingest
// with the same id replaces rather than duplicates.
ExternalID string
Title string
URI string
Body string
// Audience is who may read it. Required — see TagsFor.
Audience Audience
// Metadata is provenance the surface renders beside a citation. Never
// interpolated into a prompt: I7 covers everything on this table.
Metadata map[string]any
}
// IngestResult reports what an ingest did.
type IngestResult struct {
DocumentID string
Chunks int
Embedded int
// Unchanged is set when the body hashed identically to what was already
// stored and nothing was re-chunked or re-embedded. Worth reporting because
// re-embedding an unchanged corpus is the most expensive no-op available.
Unchanged bool
// EmbeddingDeferred is set when chunks were written but not embedded,
// because the embedder was unavailable. The document is retrievable by
// keyword in the meantime, and a backfill can finish the job.
//
// Reported rather than swallowed: a corpus that is silently keyword-only is
// a retrieval quality problem that presents as "the agent seems worse than
// it was" months later.
EmbeddingDeferred bool
}
// Ingester writes documents into the index.
type Ingester struct {
db repo.Querier
embedder Embedder
}
// NewIngester builds an ingester. A nil embedder is allowed: chunks are written
// and left unembedded for a backfill, which is better than refusing the
// document outright.
func NewIngester(db repo.Querier, e Embedder) *Ingester {
return &Ingester{db: db, embedder: e}
}
// Ingest writes one document and its chunks.
//
// Ordering matters and is deliberate:
//
// 1. Derive the ACL. Refuse before touching the database if it reaches nobody.
// 2. Upsert the document, hash-checked, so an unchanged body is a no-op.
// 3. Replace its chunks wholesale.
// 4. Embed, and tolerate failure — a document that is keyword-searchable today
// and dense-searchable after a backfill is better than one that was
// rejected because a provider was rate-limiting.
func (i *Ingester) Ingest(ctx context.Context, orgID string, doc Document) (*IngestResult, error) {
if strings.TrimSpace(orgID) == "" {
return nil, &Error{Code: ErrIngestFailed, Message: "a document needs an organization"}
}
if strings.TrimSpace(doc.Source) == "" || strings.TrimSpace(doc.ExternalID) == "" {
return nil, &Error{Code: ErrIngestFailed, Message: "a document needs a source and an external id"}
}
// §5's gate. Before any write, so a refused document leaves no trace.
tags, err := TagsFor(doc.Audience)
if err != nil {
return nil, &Error{Code: ErrNoAudience, Message: err.Error(), Cause: err}
}
body := strings.TrimSpace(doc.Body)
if body == "" {
return nil, &Error{Code: ErrIngestFailed, Message: "a document needs a body"}
}
// The hash covers the ACL as well as the text. A document whose audience
// changed has not changed its words, but it HAS changed what a retrieval
// may return — and the chunks carry a denormalised copy of the tags, so
// they must be rewritten.
hash := contentHash(doc.Title, body, tags)
metadata := doc.Metadata
if metadata == nil {
metadata = map[string]any{}
}
encodedMeta, err := json.Marshal(metadata)
if err != nil {
return nil, &Error{Code: ErrIngestFailed, Message: "the document metadata could not be encoded", Cause: err}
}
// The PREVIOUS hash, read before the upsert overwrites it. This is the
// whole of the unchanged check, and it has to happen first: once the
// document row carries the new hash there is nothing left to compare
// against, and every ingest would look like a change.
previous, chunksIntact, err := i.priorState(ctx, orgID, doc.Source, doc.ExternalID)
if err != nil {
return nil, err
}
var documentID string
err = i.db.QueryRow(ctx, `
INSERT INTO knowledge_documents
(org_id, source, external_id, title, uri, acl, acl_version, metadata, content_hash)
VALUES ($1::uuid, $2, $3, $4, $5, $6, $7, $8::jsonb, $9)
ON CONFLICT (org_id, source, external_id) DO UPDATE
SET title = EXCLUDED.title,
uri = EXCLUDED.uri,
acl = EXCLUDED.acl,
acl_version = EXCLUDED.acl_version,
metadata = EXCLUDED.metadata,
content_hash = EXCLUDED.content_hash,
ingested_at = now(),
updated_date = now()
RETURNING id::text`,
orgID, doc.Source, doc.ExternalID, doc.Title, doc.URI, tags, ACLVersion, encodedMeta, hash,
).Scan(&documentID)
if err != nil {
return nil, &Error{Code: ErrIngestFailed, Message: "the document could not be written", Cause: err}
}
// Unchanged means BOTH that the content hashed the same AND that the chunks
// actually made it into the table last time. A document whose ingest died
// between writing the document row and writing its chunks would otherwise
// be permanently "unchanged" and permanently unretrievable.
if previous == hash && chunksIntact {
var count int
if err := i.db.QueryRow(ctx,
`SELECT count(*) FROM knowledge_chunks WHERE document_id = $1::uuid`, documentID,
).Scan(&count); err != nil {
return nil, &Error{Code: ErrIngestFailed, Message: "the chunk count could not be read", Cause: err}
}
return &IngestResult{DocumentID: documentID, Chunks: count, Embedded: count, Unchanged: true}, nil
}
chunks := Split(doc.Title, body)
if len(chunks) == 0 {
return nil, &Error{Code: ErrIngestFailed, Message: "the document produced no chunks"}
}
// Replaced wholesale rather than diffed. A diff would save writes on a
// small edit and would have to reason about ordinals shifting, which is
// exactly the kind of cleverness that leaves an orphaned chunk carrying an
// old ACL. Delete-then-insert cannot.
if _, err := i.db.Exec(ctx,
`DELETE FROM knowledge_chunks WHERE document_id = $1::uuid`, documentID); err != nil {
return nil, &Error{Code: ErrIngestFailed, Message: "the old chunks could not be removed", Cause: err}
}
// Embed before inserting, so a chunk row is written with its vector in one
// statement rather than inserted and then updated.
vectors, embedErr := i.embed(ctx, chunks)
model := ""
if i.embedder != nil {
model = i.embedder.Model()
}
if err := i.insertChunks(ctx, documentID, orgID, doc.Source, tags, chunks, vectors, model); err != nil {
return nil, err
}
if _, err := i.db.Exec(ctx,
`UPDATE knowledge_documents SET chunk_count = $2 WHERE id = $1::uuid`,
documentID, len(chunks)); err != nil {
return nil, &Error{Code: ErrIngestFailed, Message: "the chunk count could not be recorded", Cause: err}
}
result := &IngestResult{DocumentID: documentID, Chunks: len(chunks)}
if vectors == nil {
result.EmbeddingDeferred = true
_ = embedErr // reported through the flag; the document is still usable
} else {
result.Embedded = len(vectors)
}
return result, nil
}
// priorState reads what was already stored for this document.
//
// Called BEFORE the upsert, because the upsert destroys the answer. Returns the
// hash the previous ingest recorded and whether that ingest's chunks are all
// still present — the second half matters because an ingest that died halfway
// leaves a document row claiming a chunk count it does not have, and comparing
// hashes alone would decline to fix it forever.
//
// A document that has never been ingested returns ("", false), which compares
// unequal to every hash and therefore always chunks.
func (i *Ingester) priorState(ctx context.Context, orgID, source, externalID string) (hash string, chunksIntact bool, err error) {
var (
claimed int
actual int
)
scanErr := i.db.QueryRow(ctx, `
SELECT d.content_hash, d.chunk_count,
(SELECT count(*) FROM knowledge_chunks c WHERE c.document_id = d.id)
FROM knowledge_documents d
WHERE d.org_id = $1::uuid AND d.source = $2 AND d.external_id = $3`,
orgID, source, externalID,
).Scan(&hash, &claimed, &actual)
if scanErr != nil {
if errors.Is(scanErr, pgx.ErrNoRows) {
return "", false, nil
}
return "", false, &Error{Code: ErrIngestFailed, Message: "the document could not be read", Cause: scanErr}
}
return hash, claimed > 0 && actual == claimed, nil
}
// insertChunks writes a document's chunks in as few statements as possible.
//
// One multi-row INSERT rather than a statement per chunk: a 40-chunk document
// is 40 round trips otherwise, and ingest is the path that runs over a whole
// corpus. Batched at insertBatch rows because Postgres caps a statement at
// 65535 bind parameters and this uses ten per chunk.
func (i *Ingester) insertChunks(ctx context.Context, documentID, orgID, source string,
tags []string, chunks []Chunk, vectors [][]float32, model string) error {
for start := 0; start < len(chunks); start += insertBatch {
end := start + insertBatch
if end > len(chunks) {
end = len(chunks)
}
var (
values []string
args []any
)
for n := start; n < end; n++ {
c := chunks[n]
var vec any
var vecModel string
if vectors != nil && n < len(vectors) && len(vectors[n]) > 0 {
vec, vecModel = vectors[n], model
}
base := len(args)
values = append(values, fmt.Sprintf(
"($%d::uuid, $%d::uuid, $%d, $%d, $%d, $%d, $%d, $%d, $%d, $%d)",
base+1, base+2, base+3, base+4, base+5, base+6, base+7, base+8, base+9, base+10))
args = append(args, documentID, orgID, source, tags, c.Ordinal,
c.Text, c.Heading, vec, vecModel, c.TokenEstimate)
}
if _, err := i.db.Exec(ctx, `
INSERT INTO knowledge_chunks
(document_id, org_id, source, acl, ordinal,
text, heading, embedding, embedding_model, token_estimate)
VALUES `+strings.Join(values, ", "), args...); err != nil {
return &Error{Code: ErrIngestFailed, Message: "the chunks could not be written", Cause: err}
}
}
return nil
}
// insertBatch is how many chunks go in one statement. Ten bind parameters each,
// against Postgres's 65535 limit, with room to spare.
const insertBatch = 500
// embed vectors for a set of chunks, tolerating an unavailable provider.
//
// Returns nil vectors rather than an error when embedding could not happen. The
// caller writes the chunks anyway: a document that is keyword-searchable now
// and dense-searchable after a backfill is strictly better than one rejected
// because a rate limit was in force for ninety seconds.
func (i *Ingester) embed(ctx context.Context, chunks []Chunk) ([][]float32, error) {
if i.embedder == nil {
return nil, &Error{Code: ErrNotConfigured, Message: "no embedder is configured"}
}
texts := make([]string, len(chunks))
for n, c := range chunks {
// The heading goes into the embedded text as well as the tsvector. A
// chunk that says "ten minutes" means something different under
// "Lateness" than under "Break entitlement", and the vector should know.
if c.Heading != "" {
texts[n] = c.Heading + "\n\n" + c.Text
} else {
texts[n] = c.Text
}
}
vectors, err := i.embedder.Embed(ctx, texts, KindDocument)
if err != nil {
return nil, err
}
return vectors, nil
}
// contentHash fingerprints what a document's chunks were built from.
func contentHash(title, body string, tags []string) string {
h := sha256.New()
h.Write([]byte(title))
h.Write([]byte{0})
h.Write([]byte(body))
h.Write([]byte{0})
for _, t := range tags {
h.Write([]byte(t))
h.Write([]byte{0})
}
return hex.EncodeToString(h.Sum(nil))
}
/* ── Re-embedding ───────────────────────────────────────────────────────── */
// Reembed gives every chunk in a tenant a vector from the current model.
//
// THE PROBLEM THIS SOLVES IS SILENT. Vectors from two embedding models are not
// comparable, so every chunk records which model produced it and retrieval only
// searches the ones matching the current embedder. Switch provider — or pull a
// newer model — and the old vectors are not wrong, they are simply not looked
// at. Retrieval keeps working, keeps citing, and quietly drops to keyword-only.
// Nothing errors. The only symptom is answers getting worse.
//
// It is also what §5 means by "reindex is required whenever ACL derivation
// logic changes", from the other direction: a corpus whose vectors no longer
// match the reader is a corpus that has stopped being fully searchable.
//
// Works in batches and reports progress, because a real corpus takes long
// enough that a silent command is one an operator kills.
func (i *Ingester) Reembed(ctx context.Context, orgID string, batch int,
progress func(done, total int)) (int, error) {
if i.embedder == nil {
return 0, &Error{Code: ErrNotConfigured, Message: "no embedder is configured"}
}
if strings.TrimSpace(orgID) == "" {
return 0, &Error{Code: ErrIngestFailed, Message: "re-embedding needs an organization"}
}
if batch <= 0 || batch > 128 {
// The provider is the constraint, not this loop. A batch far past what
// a local model holds in memory turns one slow request into one failed
// one.
batch = 32
}
model := i.embedder.Model()
var total int
if err := i.db.QueryRow(ctx, `
SELECT count(*) FROM knowledge_chunks
WHERE org_id = $1::uuid AND (embedding IS NULL OR embedding_model <> $2)`,
orgID, model).Scan(&total); err != nil {
return 0, &Error{Code: ErrIngestFailed, Message: "the corpus could not be counted", Cause: err}
}
if total == 0 {
return 0, nil
}
done := 0
for {
// Re-queried each round rather than paged: the predicate is "still
// needs this model", and rows leave that set as they are written. An
// OFFSET would walk past rows the previous round had just fixed.
rows, err := i.db.Query(ctx, `
SELECT id::text, heading, text
FROM knowledge_chunks
WHERE org_id = $1::uuid AND (embedding IS NULL OR embedding_model <> $2)
ORDER BY created_date
LIMIT $3`, orgID, model, batch)
if err != nil {
return done, &Error{Code: ErrIngestFailed, Message: "the corpus could not be read", Cause: err}
}
var (
ids []string
texts []string
)
for rows.Next() {
var id, heading, text string
if err := rows.Scan(&id, &heading, &text); err != nil {
rows.Close()
return done, &Error{Code: ErrIngestFailed, Message: "a chunk could not be read", Cause: err}
}
ids = append(ids, id)
// The heading goes into the embedded text, exactly as it does on
// first ingest. A re-embed that dropped it would produce vectors
// subtly different from the ones ingest makes, and the difference
// would show up as retrieval quality drifting after a reindex.
if heading != "" {
texts = append(texts, heading+"\n\n"+text)
} else {
texts = append(texts, text)
}
}
rows.Close()
if len(ids) == 0 {
break
}
vectors, err := i.embedder.Embed(ctx, texts, KindDocument)
if err != nil {
return done, err
}
if len(vectors) != len(ids) {
return done, &Error{Code: ErrEmbedFailed, Message: "the embedder returned the wrong number of vectors"}
}
for n, id := range ids {
if _, err := i.db.Exec(ctx, `
UPDATE knowledge_chunks
SET embedding = $2, embedding_model = $3
WHERE id = $1::uuid`, id, vectors[n], model); err != nil {
return done, &Error{Code: ErrIngestFailed, Message: "a chunk could not be updated", Cause: err}
}
done++
}
if progress != nil {
progress(done, total)
}
}
return done, nil
}