447 lines
16 KiB
Go
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
|
|
}
|