Series overview
Part 5 of 1729% complete
2026-08-21•7 min read

Build the runbook ingestion pipeline: Markdown to pgvector

Checkpoint tag: chapter-04-knowledge-ingestion — docs/runbooks is ingested idempotently, chunks carry tenant/service/version metadata, and you can inspect them with SQL.

What will be built

knowledge-ingestion becomes a batch job: it streams Markdown files from docs/runbooks, splits them on heading boundaries before applying token limits, computes deterministic chunk IDs, embeds each chunk through the configured EmbeddingModel, and upserts into opsdb.knowledge. Unchanged documents are skipped by checksum; superseded versions are retired; --ingest.dry-run=true reports the plan without writing.

Why it matters

RAG quality is decided here, not at query time. Chunks that split a procedure mid-step, IDs that change between runs, or metadata that can’t filter by tenant all show up downstream as wrong answers and cross-tenant leaks. And the incremental design matters operationally: re-embedding every document on every run turns a 30-second job into a 30-minute one and makes the embedding provider a billing dependency of your CI.

Production note: the prompt-sized files here could be read with Files.readString, but the parser streams line-by-line anyway (BufferedReader). The day someone drops a 200 MB export into the runbook directory, the job should degrade into “slow but bounded memory,” not OutOfMemoryError. Splitting an already-loaded giant string does not fix that — the load already happened.

Prerequisites and starting tag

Starting tag: chapter-03-basic-agent-api. Postgres (with pgvector) and Ollama running; nomic-embed-text pulled (the Compose Ollama entrypoint from Chapter 3 does this).

Concepts explained

Semantic before token. A Markdown document already declares its structure in headings. We split on ##-level sections first — a chunk never crosses a procedure boundary — then subdivide oversized sections at token boundaries. Reversing that order produces chunks that start mid-command.

Deterministic chunk IDs. chunk_id = sha256(documentId + version + headingPath + index). Re-ingesting an unchanged file produces identical IDs — the upsert is a no-op, not a delete/insert storm — and a chunk ID in a citation is stable enough to audit later.

Why our own tables, not PgVectorStore. Spring AI’s PgVectorStore is a good default, but it stores metadata as a JSON blob and manages its own schema. We need typed columns for tenant_id, doc_version, and service_id — the Chapter 5 retrieval filter is a plain WHERE clause with an index behind it, not a JSON predicate. The trade-off is real (we own the similarity SQL and the HNSW index), and it’s stated as an ADR in docs/adr/.

Files added or changed

docs/runbooks/*.md (seeded runbooks)
knowledge-ingestion/build.gradle.kts
knowledge-ingestion/src/main/resources/application.yml
knowledge-ingestion/src/main/resources/db/migration/V1__knowledge_init.sql
knowledge-ingestion/src/main/java/in/o612/eng/opsagent/ingest/
IngestionApplication.java (exists), IngestionRunner.java
parse/MarkdownSectionParser.java, parse/RunbookFrontMatter.java
chunk/Chunker.java
embed/EmbeddingGateway.java
store/KnowledgeStore.java
IngestionProperties.java
knowledge-ingestion/src/test/java/...

Complete code

application.yml

knowledge-ingestion/src/main/resources/application.yml
spring:
application:
name: knowledge-ingestion
datasource:
url: ${OPS_DB_URL:jdbc:postgresql://localhost:5432/opsdb}
username: ${OPS_DB_USER:agent}
password: ${OPS_DB_PASSWORD:agent-dev-password}
flyway:
schemas: knowledge
default-schema: knowledge
locations: classpath:db/migration
ai:
ollama:
base-url: ${OLLAMA_BASE_URL:http://localhost:11434}
embedding:
options:
model: nomic-embed-text
main:
web-application-type: none
ingest:
source-dir: ${RUNBOOK_DIR:docs/runbooks}
dry-run: false
max-chunk-tokens: 512
overlap-tokens: 64

V1__knowledge_init.sql

knowledge-ingestion/src/main/resources/db/migration/V1__knowledge_init.sql
CREATE SCHEMA IF NOT EXISTS knowledge;
CREATE EXTENSION IF NOT EXISTS vector;
CREATE TABLE knowledge.documents (
document_id TEXT NOT NULL,
version TEXT NOT NULL,
tenant_id TEXT NOT NULL,
service_id TEXT,
title TEXT NOT NULL,
effective_date DATE,
checksum TEXT NOT NULL,
source_uri TEXT NOT NULL,
superseded BOOLEAN NOT NULL DEFAULT FALSE,
PRIMARY KEY (document_id, version)
);
CREATE TABLE knowledge.document_chunks (
chunk_id TEXT PRIMARY KEY,
document_id TEXT NOT NULL,
version TEXT NOT NULL,
chunk_index INT NOT NULL,
heading_path TEXT NOT NULL,
content TEXT NOT NULL,
tenant_id TEXT NOT NULL,
service_id TEXT,
embedding vector(768) NOT NULL
);
CREATE INDEX chunks_tenant_service_idx
ON knowledge.document_chunks (tenant_id, service_id);
CREATE INDEX chunks_embedding_hnsw
ON knowledge.document_chunks
USING hnsw (embedding vector_cosine_ops);
CREATE TABLE knowledge.ingestion_runs (
id BIGINT GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
started_at TIMESTAMPTZ NOT NULL DEFAULT now(),
dry_run BOOLEAN NOT NULL,
scanned INT NOT NULL,
embedded INT NOT NULL,
skipped INT NOT NULL,
superseded INT NOT NULL
);

RunbookFrontMatter + MarkdownSectionParser — front matter is key: value between --- fences; sections split on ATX headings, tracking the heading path (Deployment > Rollback):

knowledge-ingestion/src/main/java/in/o612/eng/opsagent/ingest/parse/MarkdownSectionParser.java
package in.o612.eng.opsagent.ingest.parse;
import java.io.BufferedReader;
import java.io.IOException;
import java.io.Reader;
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Deque;
import java.util.List;
public class MarkdownSectionParser {
public record Section(List<String> headingPath, String body) {}
public List<Section> parse(Reader source) throws IOException {
var sections = new ArrayList<Section>();
var path = new ArrayDeque<String>();
var body = new StringBuilder();
boolean inFence = false;
String currentHeading = null;
try (var reader = new BufferedReader(source)) {
String line;
while ((line = reader.readLine()) != null) {
if (line.strip().startsWith("```")) {
inFence = !inFence; // headings inside code are not headings
body.append(line).append('\n');
continue;
}
var m = HEADING.matcher(line);
if (!inFence && m.matches()) {
flush(sections, path, body, currentHeading);
int level = m.group(1).length();
String text = m.group(2).strip();
while (path.size() >= level) path.pollLast();
path.addLast(text);
currentHeading = text;
continue;
}
body.append(line).append('\n');
}
}
flush(sections, path, body, currentHeading);
return sections;
}
private static final java.util.regex.Pattern HEADING =
java.util.regex.Pattern.compile("^(#{1,6})\\s+(.+)$");
private void flush(List<Section> out, Deque<String> path, StringBuilder body, String heading) {
var text = body.toString().strip();
if (!text.isEmpty()) {
out.add(new Section(List.copyOf(path), text));
}
body.setLength(0);
}
}

Add the tokenizer dependency — jtokkit = "1.1.0" in the catalog’s [versions], jtokkit = { module = "com.knuddels:jtokkit", version.ref = "jtokkit" } under [libraries], then implementation(libs.jtokkit) in knowledge-ingestion/build.gradle.kts.

Chunker — subdivide oversized sections on paragraph boundaries using a token count (jtokkit’s cl100k_base, close enough to every embedding model’s splitter for budgeting purposes):

knowledge-ingestion/src/main/java/in/o612/eng/opsagent/ingest/chunk/Chunker.java
package in.o612.eng.opsagent.ingest.chunk;
import com.knuddels.jtokkit.Encodings;
import com.knuddels.jtokkit.api.Encoding;
import in.o612.eng.opsagent.ingest.parse.MarkdownSectionParser.Section;
import java.util.ArrayList;
import java.util.List;
public class Chunker {
public record Chunk(List<String> headingPath, int index, String content) {}
private final Encoding encoding = Encodings.newDefaultEncodingRegistry()
.getEncoding(com.knuddels.jtokkit.api.EncodingType.CL100K_BASE);
private final int maxTokens;
private final int overlapTokens;
public Chunker(int maxTokens, int overlapTokens) {
this.maxTokens = maxTokens;
this.overlapTokens = overlapTokens;
}
public List<Chunk> chunk(Section section) {
if (count(section.body()) <= maxTokens) {
return List.of(new Chunk(section.headingPath(), 0, section.body()));
}
var chunks = new ArrayList<Chunk>();
var current = new StringBuilder();
int index = 0;
for (String para : section.body().split("\n\n")) {
if (count(current + para) > maxTokens && current.length() > 0) {
chunks.add(new Chunk(section.headingPath(), index++, current.toString().strip()));
current = new StringBuilder(tailTokens(current.toString())).append(para).append("\n\n");
} else {
current.append(para).append("\n\n");
}
}
if (current.length() > 0) {
chunks.add(new Chunk(section.headingPath(), index, current.toString().strip()));
}
return chunks;
}
private int count(String s) { return encoding.countTokens(s); }
private String tailTokens(String s) {
// carry the last ~overlapTokens worth of text into the next chunk
var ids = encoding.encode(s);
if (ids.size() <= overlapTokens) return s;
return encoding.decode(ids.subList(ids.size() - overlapTokens, ids.size()));
}
}

IngestionRunner — the job itself. Plan first, write second; dry-run stops after the plan:

knowledge-ingestion/src/main/java/in/o612/eng/opsagent/ingest/IngestionRunner.java
package in.o612.eng.opsagent.ingest;
import in.o612.eng.opsagent.ingest.chunk.Chunker;
import in.o612.eng.opsagent.ingest.embed.EmbeddingGateway;
import in.o612.eng.opsagent.ingest.parse.MarkdownSectionParser;
import in.o612.eng.opsagent.ingest.parse.RunbookFrontMatter;
import in.o612.eng.opsagent.ingest.store.KnowledgeStore;
import org.springframework.boot.ApplicationArguments;
import org.springframework.boot.ApplicationRunner;
import org.springframework.stereotype.Component;
import java.io.FileReader;
import java.nio.charset.StandardCharsets;
import java.nio.file.Files;
import java.nio.file.Path;
import java.security.MessageDigest;
import java.util.HexFormat;
import java.util.List;
import java.util.stream.Stream;
@Component
public class IngestionRunner implements ApplicationRunner {
private final MarkdownSectionParser parser = new MarkdownSectionParser();
private final Chunker chunker;
private final EmbeddingGateway embeddings;
private final KnowledgeStore store;
private final IngestionProperties props;
public IngestionRunner(Chunker chunker, EmbeddingGateway embeddings,
KnowledgeStore store, IngestionProperties props) {
this.chunker = chunker;
this.embeddings = embeddings;
this.store = store;
this.props = props;
}
@Override
public void run(ApplicationArguments args) throws Exception {
int scanned = 0, embedded = 0, skipped = 0, superseded = 0;
try (Stream<Path> files = Files.list(Path.of(props.sourceDir()))) {
List<Path> runbooks = files.filter(p -> p.toString().endsWith(".md")).sorted().toList();
for (Path file : runbooks) {
scanned++;
String checksum = sha256(file);
var fm = RunbookFrontMatter.read(file);
if (store.isCurrent(fm.documentId(), fm.version(), checksum)) {
skipped++;
continue;
}
if (props.dryRun()) {
System.out.printf("[dry-run] would ingest %s v%s (%s)%n",
fm.documentId(), fm.version(), file);
continue;
}
try (var reader = new FileReader(file.toFile(), StandardCharsets.UTF_8)) {
int idx = 0;
for (var section : parser.parse(reader)) {
for (var chunk : chunker.chunk(section)) {
String chunkId = chunkId(fm.documentId(), fm.version(),
String.join(">", chunk.headingPath()), idx);
float[] vector = embeddings.embed(chunk.content());
store.upsertChunk(chunkId, fm, chunk, vector, idx++);
}
}
}
superseded += store.supersedeOlderVersions(fm.documentId(), fm.version(), checksum);
store.recordDocument(fm, checksum, file.toUri().toString());
embedded++;
}
}
store.recordRun(props.dryRun(), scanned, embedded, skipped, superseded);
System.out.printf("ingest complete: scanned=%d embedded=%d skipped=%d superseded=%d%n",
scanned, embedded, skipped, superseded);
}
private static String chunkId(String doc, String version, String headingPath, int index) {
return "chk-" + sha256(doc + "|" + version + "|" + headingPath + "|" + index).substring(0, 24);
}
private static String sha256(Path file) throws Exception {
try (var in = Files.newInputStream(file)) {
return HexFormat.of().formatHex(
MessageDigest.getInstance("SHA-256").digest(in.readAllBytes()));
}
}
private static String sha256(String s) {
try {
return HexFormat.of().formatHex(MessageDigest.getInstance("SHA-256")
.digest(s.getBytes(StandardCharsets.UTF_8)));
} catch (java.security.NoSuchAlgorithmException e) { throw new IllegalStateException(e); }
}
}

EmbeddingGateway wraps EmbeddingModel.embedForResponse(List.of(text)) behind a port so tests can substitute a fixed-vector stub — one more place where tests never touch a model.

KnowledgeStore upserts with embedding = cast(:v as vector), and supersedeOlderVersions marks previous documents rows and deletes their chunks — a runbook’s old procedure must stop being retrievable the moment the new one lands.

A seeded runbook (docs/runbooks/payment-gateway-degraded.md) with front matter document_id: rb-payment-gateway-degraded, version: "3", tenant_id: acme, service_id: payment-gateway, effective_date: 2026-08-01, then ## sections for Symptoms, Diagnosis, Mitigation, Rollback. The repo carries three runbooks covering both tenants and one deliberately stale version: "2" twin of the same document to exercise supersede.

Commands to build and run

terminal
./gradlew :knowledge-ingestion:bootRun --args='--ingest.dry-run=true'
./gradlew :knowledge-ingestion:bootRun
psql "postgresql://agent:agent-dev-password@localhost:5432/opsdb" \
-c "select chunk_id, heading_path, left(content,60) from knowledge.document_chunks"
./gradlew :knowledge-ingestion:bootRun # again -> skipped=N, embedded=0

Automated tests

  • MarkdownSectionParserTest: headings inside ``` fences are not split; heading paths nest correctly; a document with no headings yields one section.
  • ChunkerTest: sections under the token limit pass through untouched; oversized sections split on paragraph boundaries with overlap; chunk count is deterministic.
  • IngestionIT (Testcontainers pgvector/pgvector:0.8.6-pg17 + stub EmbeddingGateway returning fixed vectors): first run inserts N chunks; second run inserts zero; editing a runbook re-embeds only that document; superseded versions return no chunks.
  • DeterminismTest: same file → same chunk IDs on repeated runs.

Failure-injection lab

  1. Point OLLAMA_BASE_URL at a dead port mid-run: the job fails with a partial batch committed. Re-run — checksums make the retry safe; only unwritten chunks embed. This is why recordDocument happens after all its chunks.
  2. Flip a runbook’s tenant_id front matter and re-ingest: the new chunks carry the new tenant while the old version’s chunks are gone — verify no acme chunk survives for the globex-owned doc.

Security considerations

Runbook content is untrusted input even though we wrote it — Chapter 5 wraps retrieved text in delimiters and Chapter 13 includes a poisoned-document eval case. document_chunks never stores secrets: ingestion logs counts and IDs, never content. The HNSW index is built on the column itself; a WHERE tenant_id = ? filter in the query keeps cross-tenant recall impossible rather than merely unlikely.

Troubleshooting

  • type "vector" does not exist: pgvector extension missing — CREATE EXTENSION vector is in V1; confirm the migration ran against opsdb, not simdb.
  • Dimension mismatch (expected 768): you swapped embedding models without updating the column; embedding dimensions are part of the schema contract — changing the model is a migration, not a config tweak.
  • Dry-run prints nothing: --ingest.dry-run=true with an unchanged checksum set still exits quietly; delete a chunk row to force work.

Checkpoint verification checklist

  • dry-run reports the plan and writes nothing (verify ingestion_runs row).
  • Second full run: skipped=3, embedded=0.
  • Chunk IDs identical across runs; document_chunks queryable by tenant_id.
  • Superseded version contributes zero chunks.

Commit message and Git tag

feat(ingestion): streaming markdown-to-pgvector pipeline with incremental re-ingest

git tag chapter-04-knowledge-ingestion

What comes next

Chapter 5 turns these chunks into answers: the retrieval port, tenant-filtered similarity search, citation mapping, and — just as important — the abstention path when evidence isn’t there.

Project State Ledger — chapter-04-knowledge-ingestion

  • Tables: knowledge.documents, knowledge.document_chunks (vector(768), HNSW), knowledge.ingestion_runs (V1__knowledge_init.sql)
  • Pipeline: stream → front matter → heading sections → token-boundary chunker → embed → upsert; checksum skip; supersede on new version; ingestion_runs audit row
  • Ports: EmbeddingGateway (stubbable)
  • Chunk ID: chk- + sha256(docId|version|headingPath|index)[:24]
  • Config: ingest.source-dir, ingest.dry-run, ingest.max-chunk-tokens=512, ingest.overlap-tokens=64
  • Test doubles: fixed-vector embedding stub; pgvector Testcontainer
  • Next: chapter-05-rag-citations
JavaPostgresAISpring Boot

Type to search the site.

↑↓ navigate⏎ openPowered by Pagefind