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 }