// Command ingest puts Markdown documents into a tenant's knowledge corpus. // // The knowledge layer had an Ingester and no way to reach it — everything that // had ever been ingested was ingested by a test. This is the missing half. // // make ingest ORG= // // Each document declares its own audience in front matter, and a document that // declares none is REFUSED rather than defaulted. Both directions of a default // are wrong and neither raises: tenant-wide over-shares something somebody // meant to restrict, and empty indexes it into invisibility. See §5. package main import ( "context" "errors" "flag" "fmt" "os" "path/filepath" "sort" "strings" "time" "github.com/jackc/pgx/v5" "github.com/krow/krow-backend/go-api/internal/config" "github.com/krow/krow-backend/go-api/internal/db" "github.com/krow/krow-backend/go-api/internal/domain" "github.com/krow/krow-backend/go-api/internal/knowledge" "github.com/krow/krow-backend/go-api/internal/runtime" ) func main() { var ( dir = flag.String("dir", "./knowledge", "directory of Markdown documents") org = flag.String("org", "", "organization slug to ingest into (required)") dryRun = flag.Bool("dry-run", false, "parse and report, write nothing") timeout = flag.Duration("timeout", 15*time.Minute, "overall timeout") ) flag.Parse() if err := run(*dir, *org, *dryRun, *timeout); err != nil { fmt.Fprintf(os.Stderr, "ingest: %v\n", err) os.Exit(1) } } // parsed is one document, read and validated before anything is opened. type parsed struct { file string doc knowledge.Document } func run(dir, orgSlug string, dryRun bool, timeout time.Duration) error { if strings.TrimSpace(orgSlug) == "" { return errors.New("an organization is required: --org=") } docs, err := readAll(dir) if err != nil { return err } if len(docs) == 0 { return fmt.Errorf("no documents found in %s", dir) } for _, d := range docs { tags, err := knowledge.TagsFor(d.doc.Audience) if err != nil { return fmt.Errorf("%s: %w", d.file, err) } fmt.Printf(" %-28s %-14s %s\n", d.doc.ExternalID, d.doc.Source, strings.Join(tags, " ")) } if dryRun { fmt.Printf("\n%d document(s) parsed; nothing written (--dry-run)\n", len(docs)) return nil } cfg, err := config.Load() if err != nil { return fmt.Errorf("load configuration: %w", err) } embedder := runtime.NewEmbedder(*cfg) if embedder == nil { // Not fatal. Chunks are written and left unembedded for `make reembed`, // so a corpus is keyword-searchable immediately and dense-searchable // once a model exists. Said out loud because a silently keyword-only // corpus is a retrieval problem that surfaces months later as "the // agent seems worse than it was". fmt.Println("\nno embedding model configured — documents will be keyword-searchable only") fmt.Println("set EMBED_PROVIDER and run `make reembed` to finish them") } else { fmt.Printf("\nembedding with %s\n", embedder.Model()) } ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() database, err := db.Open(ctx, cfg.DB) if err != nil { return fmt.Errorf("connect: %w", err) } defer database.Close() var orgID string err = database.Pool.QueryRow(ctx, `SELECT id::text FROM organizations WHERE slug = $1`, orgSlug).Scan(&orgID) if errors.Is(err, pgx.ErrNoRows) { return fmt.Errorf("no organization with slug %q", orgSlug) } if err != nil { return fmt.Errorf("resolve organization: %w", err) } ing := knowledge.NewIngester(database.Pool, embedder) var chunks, unchanged int for _, d := range docs { res, err := ing.Ingest(ctx, orgID, d.doc) if err != nil { return fmt.Errorf("%s: %w", d.file, err) } chunks += res.Chunks if res.Unchanged { unchanged++ fmt.Printf(" %-28s unchanged (%d chunks)\n", d.doc.ExternalID, res.Chunks) continue } note := "" if res.EmbeddingDeferred { note = " [not embedded]" } fmt.Printf(" %-28s %d chunks%s\n", d.doc.ExternalID, res.Chunks, note) } fmt.Printf("\n%d document(s), %d chunk(s), %d unchanged, into %s\n", len(docs), chunks, unchanged, orgSlug) return nil } /* ── Reading the directory ──────────────────────────────────────────────── */ func readAll(dir string) ([]parsed, error) { entries, err := os.ReadDir(dir) if err != nil { return nil, fmt.Errorf("read %s: %w", dir, err) } var out []parsed for _, e := range entries { name := e.Name() if e.IsDir() || !strings.HasSuffix(name, ".md") || name == "README.md" { continue } raw, err := os.ReadFile(filepath.Join(dir, name)) if err != nil { return nil, fmt.Errorf("read %s: %w", name, err) } doc, err := parse(name, string(raw)) if err != nil { return nil, fmt.Errorf("%s: %w", name, err) } out = append(out, parsed{file: name, doc: *doc}) } sort.Slice(out, func(a, b int) bool { return out[a].file < out[b].file }) return out, nil } // parse reads a document's front matter and body. // // A small reader rather than a YAML library: the front matter here is four flat // keys, and the value of a real parser is handling shapes this format does not // have. What matters is that a malformed audience is an error rather than a // silent default. func parse(file, raw string) (*knowledge.Document, error) { body := strings.ReplaceAll(raw, "\r\n", "\n") if !strings.HasPrefix(body, "---\n") { return nil, errors.New("no front matter; a document must declare its source and audience") } end := strings.Index(body[4:], "\n---") if end < 0 { return nil, errors.New("front matter is not closed") } head := body[4 : 4+end] rest := strings.TrimLeft(body[4+end+4:], "\n") fields := map[string]string{} for _, line := range strings.Split(head, "\n") { k, v, ok := strings.Cut(line, ":") if !ok { continue } fields[strings.TrimSpace(k)] = strings.TrimSpace(v) } source := fields["source"] if source == "" { return nil, errors.New("no `source`; an agent's spec names the corpora it may read") } audience, err := parseAudience(fields["audience"]) if err != nil { return nil, err } // The filename is the external id, so re-ingesting the same file updates // rather than duplicating. Stable, obvious, and something a person can // point at. id := strings.TrimSuffix(file, ".md") title := fields["title"] if title == "" { title = id } return &knowledge.Document{ Source: source, ExternalID: id, Title: title, URI: fields["uri"], Body: rest, Audience: audience, }, nil } // parseAudience turns the declared audience into the one the ingester takes. // // An empty or unrecognised value is an ERROR. That is the whole point: §5 // refuses a document that reaches nobody, and a typo'd role silently producing // a tag no principal holds is the same failure wearing better clothes. func parseAudience(raw string) (knowledge.Audience, error) { raw = strings.TrimSpace(raw) if raw == "" { return knowledge.Audience{}, errors.New( "no `audience`; a document that declares none is unreachable, not private") } var a knowledge.Audience for _, part := range strings.Split(raw, ",") { part = strings.TrimSpace(part) switch { case part == "tenant": a.Tenant = true case strings.HasPrefix(part, "role:"): name := strings.TrimPrefix(part, "role:") role, ok := domain.ParseRole(name) if !ok { return knowledge.Audience{}, fmt.Errorf( "%q is not a role; use admin, employer or talent", name) } a.Roles = append(a.Roles, role) case strings.HasPrefix(part, "email:"): a.Emails = append(a.Emails, strings.TrimPrefix(part, "email:")) default: return knowledge.Audience{}, fmt.Errorf( "%q is not an audience; use tenant, role: or email:
", part) } } return a, nil }