Pipelines
Ingestion and embeddings
When a trainer uploads a resource the API stores the file, writes a row and queues a job. A worker in the Python ai-service extracts the text, splits it into chunks of about 500 tokens, asks Ollama for a 768-dimension vector per chunk, and stores the results in resource_chunks.
Sequence
Ingestion sequence
Job lifecycle
Flow of one job
Functions
The API side is short. Almost everything lives in the worker.
| Function | Runs in | Does |
|---|---|---|
saveUpload(file, courseId) | API | Validates type and size, writes to uploads/<course>/<uuid>.<ext>, returns the relative path. |
createResource(req) | API | Route handler for POST /api/resources. In one transaction: INSERT learning_resources, then enqueueEmbedJob. |
enqueueEmbedJob(db, resourceId) | API | INSERT into the queue table with status queued. |
claim_job(conn) | ai-service | Atomically takes one queued job using FOR UPDATE SKIP LOCKED. |
extract_text(path, mime) | ai-service | Text from pptx, pdf or transcript files. |
chunk_text(text, target_tokens=500, overlap=50) | ai-service | Splits on paragraph, then sentence boundaries. |
embed(text) | ai-service | POST to Ollama, checks the vector has 768 numbers. |
store_chunks(conn, resource_id, chunks, vectors) | ai-service | Deletes old chunks for the resource, inserts the new ones, in one transaction. |
mark_done(conn, job_id) | ai-service | status done. On an exception, mark_failed sets status failed and schedules a retry. |
Worker sketch
import os, requests
OLLAMA = os.environ["OLLAMA_HOST"] # never hardcode localhost
MODEL = os.environ.get("EMBED_MODEL", "nomic-embed-text")
def embed(text: str) -> list[float]:
r = requests.post(f"{OLLAMA}/api/embeddings",
json={"model": MODEL, "prompt": text}, timeout=60)
r.raise_for_status()
v = r.json()["embedding"]
assert len(v) == 768, f"expected 768 dims, got {len(v)}"
return v
def process(conn, job):
path, mime = load_resource_row(conn, job["resource_id"])
text = extract_text(path, mime)
chunks = chunk_text(text, target_tokens=500, overlap=50)
vecs = [embed(c) for c in chunks]
with conn.transaction():
set_role(conn, "super_admin") # set_config('app.current_role', ...)
store_chunks(conn, job["resource_id"], chunks, vecs) # delete old, insert new
mark_done(conn, job["id"])UPDATE embed_jobs
SET status = 'running', started_at = now()
WHERE id = (SELECT id FROM embed_jobs
WHERE status = 'queued'
ORDER BY created_at
FOR UPDATE SKIP LOCKED
LIMIT 1)
RETURNING id, resource_id;Rules and gotchas
store_chunks deletes the resource's old chunks and inserts new ones in one transaction. A retry after a crash cannot leave duplicates.
Workers can die mid-job. A periodic sweep should return jobs that stayed running past a timeout to queued.
Write embedding_model on every chunk. Query time should only compare vectors from the same model.
Nomic embedding models were trained with task prefixes such as search_document: and search_query:. Pick one convention, apply it everywhere, and record it in embedding_model, for example nomic-embed-text/search_document.
Newer Ollama versions accept a list of inputs at /api/embed and return embeddings. Use it if the installed version supports it, otherwise keep one call per chunk.