Welcome back to The Substack Brain Course!
This is Lesson 2 of our free six-week course, where we build a living knowledge base from the AI engineering newsletters you actually want to read, one production problem at a time. In Lesson 1 we built the foundation: a durable pipeline that turns a newsletter's RSS feed into searchable passages in Postgres.
Every step of that pipeline runs on Inngest, so if a run crashes halfway through, it resumes from the last completed step instead of paying OpenAI twice for the same embeddings. If you missed it, you can read the Lesson 1 article or watch the video (yes, I'm back on YouTube!)
Last week we tested the pipeline on one newsletter and five articles, which is a very forgiving workload.
Today we point it at three, The Neural Maze, Paul Iusztin's Decoding AI and Sebastian Raschka, PhD's Ahead of AI, and fix everything that breaks along the way. None of it is a crash: the websites we scrape start blocking us, OpenAI starts answering with 429 errors, urgent articles get stuck behind old ones, and when an answer is wrong, we have no idea why.
By the end of this lesson, the pipeline handles all four with concurrency keys, throttling, priority, hybrid retrieval and context snapshots.
💻 The Substack Brain code is open-source. Support our work by dropping a friendly ⭐ on the repo!
The system for Week 2
To fix those four problems, we add five pieces to last week's system:
🔁 backfill-publication, a new Inngest function that ingests every post in a feed instead of only the latest five.
🚦 Flow control on ingest-article: a concurrency limit per website so we don't get blocked, and a global throttle so we stay under OpenAI’s rate limit.
⚡ POST /articles, a new endpoint that pushes an urgent article ahead of everything in the queue.
🔎 Hybrid search, which combines keyword and vector search with Reciprocal Rank Fusion, so we find both exact terms and related ideas.
💬 POST /ask, which answers questions with citations and saves a snapshot of every answer, so we can tell why one went wrong.
The Asymmetry Principle
Before looking at any code, we want to talk about a temptation that catches a lot of teams once they start using a durable execution framework.
You discover steps, automatic retries and that very satisfying run timeline in the dashboard, and suddenly every endpoint looks like it should be a workflow. If ingestion is an Inngest function, why not search too? Why not the question-answering endpoint? The answer comes from looking at what each kind of work actually needs.
Ingesting an article is slow, it costs money, and nobody is sitting in front of a spinner waiting for it.
If it fails at the embed step, resuming from the embed step is exactly what you want. Answering a question is the opposite. Someone is waiting, right now, for a response. There is nothing expensive to save halfway through, and if the LLM provider has a hiccup, retrying three minutes later is pointless because the user has already closed the tab. Routing that request through an orchestrator only adds network hops and scheduling delay.
We call this the Asymmetry Principle:
Durable orchestration belongs where there is expensive, stateful work to lose. Fast, read-only, latency-sensitive paths should skip the workflow engine and query the database directly.
That raises a fair question. If POST /ask bypasses Inngest, how do we keep track of answer quality over time? (check the diagram below). The trick is to separate answering from evaluating. The endpoint answers directly from Postgres, saves a full record of what it did (more on that in later sections), and on its way out sends a small event, kb/answer.produced, to Inngest. The user gets the answer without waiting for anything else, and in Week 5 that event will trigger the evaluators that do the slow scoring work in the background, where durability actually helps.
You can see the principle in the code itself. The write routes in kb/main.py (POST /publications and the new POST /articles) only send an event and return. The read routes (GET /search and the new POST /ask) never touch Inngest. This is the whole /ask route: it opens a database session and calls ask_question directly, with no event and no run.
class AskRequest(BaseModel):
question: str
limit: int | None = None
@app.post("/ask")
async def ask(payload: AskRequest) -> dict[str, Any]:
"""Answer a question with inline citations. A fast read: it never creates an
Inngest run. It saves a snapshot and sends `kb/answer.produced` on its way out.
"""
async with get_session() as session:
return await ask_question(session, question=payload.question, limit=payload.limit)See it on GitHub: kb/main.py, lines 182–193.
Concurrency, Throttling, and Priority Queues
The first change this week is on the write path. In Week 1, POST /publications triggered add-publication, which took the five most recent entries of a feed. This week the same event triggers a new function, backfill-publication in kb/functions/backfill.py, which takes every entry in the feed, skips the articles we already have and fans out the rest as kb/article.discovered events tagged with via="backfill". Both functions listen to the same event, so kb/functions/__init__.py now registers only the new one.
The last step of the backfill is where the fan-out happens. It records a job first, so GET /jobs/{event_id} can track progress from the very first article, and then sends the new articles to Inngest in batches:
async def _record_and_fan_out() -> dict[str, int]:
# The job row goes first: GET /jobs counts articles written after the job
# started, so it must exist before the first article can land.
async with get_session() as session:
await create_job(
session,
event_id=ctx.event.id,
kind="backfill",
total=len(new_articles),
publication_id=publication_id,
)
await session.commit()
batch_size = settings.backfill_batch_size
for start in range(0, len(new_articles), batch_size):
batch = new_articles[start : start + batch_size]
await client.send([inngest.Event(name=ARTICLE_DISCOVERED, data=d) for d in batch])
return {"discovered": len(new_articles)}
return await ctx.step.run("record-and-fan-out", _record_and_fan_out)See it on GitHub: kb/functions/backfill.py, lines 97–116
One honest note about the numbers: a Substack RSS feed only carries a publication's most recent posts, around twenty, so this isn't the complete archive of each newsletter. It's still enough to break the Week 1 pipeline, because three backfills now push dozens of articles into ingest-article within a few seconds. That's where the three flow-control primitives come in: concurrency, throttling and priority.
Concurrency keys
When three backfills start at the same time, each one wants to fetch its articles as fast as it can. To a website's server, a burst of simultaneous requests from a single IP looks a lot like an attack, and that's how you end up with Cloudflare challenges or a ban. The obvious fix is a global limit, something like "never process more than three articles at once", but that creates a different problem: all three publications now share those three slots, and whichever backfill fans out first takes all of them while the other two wait.
What we really want is a limit per website, and that's what a concurrency key gives us:
@client.create_function(
fn_id="ingest-article",
trigger=inngest.TriggerEvent(event=ARTICLE_DISCOVERED),
# A separate limit for each publication, so we never hit one website with
# more than N requests at once, and one big backfill can't starve the others.
concurrency=[
inngest.Concurrency(
key="event.data.publication_id",
limit=settings.concurrency_per_publication,
)
],
...
)See it on GitHub: kb/functions/ingest.py, lines 46–56
With key="event.data.publication_id", Inngest keeps a separate limit for every publication. The Neural Maze gets three slots, Decoding AI gets three, Ahead of AI gets three, and none of them waits on the others. Everything above the limit stays in Inngest's queue; nothing is dropped. The limit lives in kb/config.py, so you can change it from .env with CONCURRENCY_PER_PUBLICATION, which will be useful for the demo later.
Throttling
Concurrency keeps the websites happy, but there's a second shared resource we haven't protected yet: the embedding provider. OpenAI limits how many requests you can send per minute, and that limit applies to your whole account, not to each publication. With nine articles in flight and new ones starting every few seconds, it doesn't take many newsletters to hit it.
For this we use a throttle, which limits how many runs can start per minute. Runs over the limit are queued and started later, so nothing is lost:
# No key: the OpenAI quota is shared by every publication, so the limit is too.
# Throttling limits run STARTS, not steps; embed_texts sends a whole article
# in one request, which keeps "runs per minute" close to "requests per minute".
throttle=inngest.Throttle(
limit=settings.ingest_runs_per_minute,
period=timedelta(minutes=1),
),See it on GitHub: kb/functions/ingest.py, lines 57–63
Notice that this one has no key. The concurrency limit protects each website separately, but the throttle protects a single quota that every publication shares.
There's a detail here that's easy to miss. Inngest's throttle counts run starts, not the steps inside a run. If each article made one embedding call per chunk, a throttle of 50 runs per minute could quietly turn into thousands of requests per minute. Our embedder avoids that because it sends an article's chunks in batches of up to 96 per request (embedding_batch_size in the config), so for a typical newsletter post, one article means one embedding request:
for start in range(0, len(texts), settings.embedding_batch_size):
batch = texts[start : start + settings.embedding_batch_size]
response = await _get_client().embeddings.create(
input=batch,
model=settings.embedding_model,
dimensions=settings.embedding_dimensions,
)
# The API does not guarantee response order matches input order.
ordered = sorted(response.data, key=lambda item: item.index)
embeddings.extend(item.embedding for item in ordered)
total_tokens += response.usage.total_tokensSee it on GitHub: kb/sources/embedder.py, lines 41–51
The loop only runs more than once for articles with more than 96 chunks, which is rare for a newsletter post. That's what makes the throttle honest: 50 runs per minute is roughly 50 embedding requests per minute. We keep it comfortably below the real quota because POST /ask also embeds every question users type.
Priority
The third problem is about order. Imagine the three backfills are running and The Neural Maze publishes a new post that needs to be searchable right away. In a first-in, first-out queue, that article would wait behind every queued The Neural Maze article.
So this week adds a second way into the pipeline. POST /articles takes a single URL, checks it against the allowlist and sends a kb/article.discovered event tagged with via="manual". It builds exactly the same event the backfill sends, so ingest-article doesn't need to know where an article came from. The only difference is the tag:
event = inngest.Event(
name=ARTICLE_DISCOVERED,
data=ArticleDiscovered(
canonical_id=canonical_id(pub.slug, payload.url),
publication_id=publication_id,
url=payload.url,
title=payload.title or "Manual ingestion",
# Only a fallback: the parse step uses the article's real date if it finds one.
published_at=datetime.now(UTC),
via="manual",
).model_dump(mode="json"),
)
event_ids = await client.send(event)
event_id = event_ids[0]
return {
"event_id": event_id,
"trace_url": f"{settings.inngest_base_url}/event/{event_id}",
"status": "queued",
"priority_offset_s": settings.priority_manual_s,
}See it on GitHub: kb/main.py, lines 96–115
Inngest can then use that tag to prioritise the run.
This is where we have to be careful, because Inngest's priority doesn't work the way most people expect (it didn't work the way we expected in our first draft either). The priority expression doesn't return a rank. It returns a time offset in seconds, between -600 and 600 by default. A value of 120 means "schedule this run as if it had arrived 120 seconds earlier than it really did", and a negative value pushes the run back.
Here's why that matters. Say the backfill queued its articles at 10:00:00 and the urgent article arrives at 10:02:00. With a priority of 100, Inngest treats it as if it had arrived at 10:00:20, which is still after every backfill article, so it waits like everything else. With a priority of 600, it's treated as if it had arrived at 9:52:00, ahead of the whole backfill. A value that sounds generous can do nothing at all, depending on how old the queue is.
So we push on both ends. Manual runs get the maximum boost and backfill runs get the maximum delay, which gives an urgent article up to twenty minutes of head start over anything already queued:
# Priority is a time shift in seconds, not a rank: +600 schedules a manual run
# as if it had arrived 10 minutes earlier, and -600 pushes backfill runs 10
# minutes later, so a manual article gets up to 20 minutes of head start.
priority=inngest.Priority(
run=(
f"event.data.via == 'manual' ? {settings.priority_manual_s} "
f": {settings.priority_backfill_s}"
)
),See it on GitHub: kb/functions/ingest.py, lines 64–72
Two things priority won't do are worth keeping in mind. It doesn't interrupt a run that's already executing, so the urgent article takes the next free slot rather than kicking someone out. And it still respects the concurrency limit: an urgent Neural Maze article waits for a The Neural Maze slot, not a Decoding AI one, which is exactly what we want, since the rule about not overloading a website shouldn't disappear just because something is urgent.
Idempotency
Sooner or later you'll run the same backfill twice. Maybe you fixed a bug in the chunker, maybe a run failed overnight, maybe someone simply clicked twice. If the pipeline isn't idempotent, every replay writes duplicate passages that pollute search results, and you pay for every embedding again.
Our protection comes in three layers. The first is a stable ID for every article. Every ingestion path builds it with the same function, canonical_id in kb/schemas/ids.py, which normalises the URL (https, lowercase host, no query string, no trailing slash) and combines the publication slug with the post slug, giving IDs like pub:the-neural-maze:some-post. The same post on www.theneuralmaze.com and on theneuralmaze.substack.com ends up with the same ID.
def canonical_id(publication_slug: str, url: str) -> str:
"""`pub:<publication_slug>:<post_slug>` — the cross-pipeline dedupe key.
`publication_slug` is passed in rather than derived from the URL's host:
every producer (discovery, backfill, freshness, manual add) already
knows which publication it is working from the allowlist, so this stays
a pure string function instead of doing a lookup.
"""
normalized = normalize_url(url)
segments = [s for s in urlsplit(normalized).path.split("/") if s]
if not segments:
raise ValueError(f"URL has no path segment to use as a post slug: {url!r}")
post_slug = segments[-1]
return f"pub:{publication_slug}:{post_slug}"See it on GitHub: kb/schemas/ids.py, lines 28–41
The second layer is skipping work before it starts. The backfill checks which articles are already stored and only fans out the new ones:
async with get_session() as session:
stored = await get_existing_canonical_ids(session, [c.canonical_id for c in candidates])
return [c.model_dump(mode="json") for c in candidates if c.canonical_id not in stored]See it on GitHub: kb/functions/backfill.py, lines 91–93
The third layer is the database. Every write is an upsert: articles are keyed by canonical_id and passages by (article_id, chunk_index), so even if two runs race each other, Postgres won’t let a duplicate in.
stmt = stmt.on_conflict_do_update(
index_elements=[Passage.article_id, Passage.chunk_index],
set_={"text": stmt.excluded.text, "embedding": stmt.excluded.embedding},
)See it on GitHub: kb/db/queries.py, lines 114–117
You might wonder why we don't also use Inngest's built-in idempotency option, which skips any event whose key has been seen in the last 24 hours. The first version of this code used it, and it had two problems. If the key was the article's ID, a manual submission for an article that was already queued by a backfill got silently discarded, so the priority from the previous section never kicked in. Worse, it broke something from Week 1: re-adding a publication is how you retry the articles that failed, and a 24-hour idempotency window blocks exactly that. With the "skip what's stored" filter and the upserts already preventing duplicates, the extra layer was costing us more than it gave us, so we removed it.
To prove that replays really are free, we can lean on something we built in Week 1. Every step that calls OpenAI returns a usage block, and a small Inngest middleware in kb/inngest_client.py writes it to the costs table:
class CostCaptureMiddleware(inngest.Middleware):
async def transform_output(self, result: inngest.TransformOutputResult) -> None:
if result.step is None:
return # not a newly-executed step's output
usage = _usage_from_output(result.output)
if usage is None:
return
raw_article_id = usage.get("article_id")
article_id = uuid.UUID(raw_article_id) if raw_article_id else None
async with get_session() as session:
await insert_cost(
session,
model=usage["model"],
input_tokens=usage["input_tokens"],
output_tokens=usage["output_tokens"],
est_usd=usage["est_usd"],
article_id=article_id,
)
await session.commit()See it on GitHub: kb/inngest_client.py, lines 48–67
The check on result.step is what makes this safe on replays: a memoised step returns its saved output without running again, so it never gets here and never records a second cost. That makes the costs table an honest record of every embedding call we actually paid for. The new scripts/replay_backfill.py uses that table as one of its four counts:
async def counts() -> dict[str, int]:
async with get_session() as session:
async def count(stmt: object) -> int:
return int((await session.execute(stmt)).scalar_one()) # type: ignore[call-overload]
return {
"publications": await count(select(func.count(Publication.id))),
"articles": await count(select(func.count(Article.id))),
"passages": await count(select(func.count(Passage.id))),
"embedding_calls": await count(
select(func.count(Cost.id)).where(Cost.model == settings.embedding_model)
),
}See it on GitHub: scripts/replay_backfill.py, lines 32–45
Run it before and after a replay:
uv run python scripts/replay_backfill.py --verify # note the counts
uv run python scripts/replay_backfill.py --all # re-send kb/publication.added for every publication
uv run python scripts/replay_backfill.py --verify # once Inngest is done: same counts--verify prints the number of publications, articles, passages and embedding calls. If ingestion is idempotent, all four numbers are identical before and after the replay.
Hybrid retrieval with Reciprocal Rank Fusion
With the write path under control, we can move to the read path, starting with retrieval. In Week 1 we built two search engines and made you pick one through GET /search?mode=sparse or mode=dense. This week they work together.
Each engine fails in a different way.
Vector search, with 1536-dimensional embeddings, is very good at meaning. It knows that "process crash" is related to "step memoization" even though the two phrases share no words. But it's surprisingly weak on exact tokens such as function names like create_function, version strings, acronyms and niche jargon.
Postgres full-text search, built on tsvector and ts_rank, is the mirror image: it finds exact matches every time, but if you search for "resilience" and the author wrote "fault tolerance", it returns nothing. (You'll often see this kind of search called BM25, including in our own code comments. Strictly speaking ts_rank is a simpler scoring function, but it does the job for hybrid search).
The new kb/retrieval/hybrid.py asks each engine for 15 candidates and merges them. The first version of this code ran both searches with asyncio.gather on the same database session, and it seemed to work fine. It only worked because the vector search has to wait for OpenAI to embed the question before touching the database, so the two queries almost never overlapped. Almost never isn't never, and SQLAlchemy’s AsyncSession refuses to run two queries at the same time. The final version overlaps the slow part, the network call, with the full-text query, and runs the vector query afterwards on the same session:
candidate_fetch_limit = max(limit * fetch_multiplier, 15)
# One AsyncSession cannot run two queries at the same time. So we overlap the
# slow network call (embedding the question with OpenAI) with the full-text
# query, and only run the vector query once the embedding is back.
sparse_results, embedded = await asyncio.gather(
search_passages_bm25(session, query, candidate_fetch_limit),
embed_texts([query]),
)
(query_embedding,) = embedded["embeddings"]
dense_results = await search_passages_dense(session, query_embedding, candidate_fetch_limit)See it on GitHub: kb/retrieval/hybrid.py, lines 98–108
Now we have two ranked lists, and the obvious way to merge them is to add each passage's two scores. That doesn't work, because the scores live on completely different scales: ts_rank returns small positive numbers with no fixed maximum, while cosine similarity tops out at 1. Adding them is like adding someone's height in centimetres to their age in years. Whichever engine happens to produce larger numbers wins, for no good reason.
Reciprocal Rank Fusion avoids the problem by ignoring the scores and using only positions. Each engine contributes 1 / (k + rank) for every passage it found, where rank is the passage's position in that engine’s list and k is a smoothing constant, 60 by convention, that stops the first result from completely overwhelming the second and third. A passage final score is the sum of what it gets from each engine.
A small example shows why this works. Say passage A is first in keyword search and second in vector search, so it scores 1/61 + 1/62, about 0.0325. Passage C is third in keyword search and fifth in vector search, so it scores 1/63 + 1/65, about 0.0313. Passage B is the vector engine's first result, but keyword search never found it, so it only gets 1/61, about 0.0164. B was one engine’s favourite and still finishes last, while C, which was never anybody's top pick, beats it comfortably because both engines thought it was relevant. When two very different engines agree, that's a strong signal, and RRF rewards it.
The implementation is short:
def compute_rrf(
sparse_results: list[dict[str, Any]],
dense_results: list[dict[str, Any]],
k: int = 60,
) -> list[RetrievedCandidate]:
"""score(passage) = sum over engines of 1 / (k + rank in that engine).
A passage found by both engines collects points twice, so agreement between
two very different engines beats being one engine's favourite. Ties are
broken by passage_id, so the same question always gives the same order.
"""
candidates: dict[str, RetrievedCandidate] = {}
for rank, item in enumerate(sparse_results, start=1):
cand = _candidate(item, 1.0 / (k + rank))
cand.sparse_rank = rank
cand.sparse_score = float(item.get("rank", 0.0))
cand.source = "bm25"
candidates[cand.passage_id] = cand
for rank, item in enumerate(dense_results, start=1):
pid = str(item["passage_id"])
if pid in candidates:
cand = candidates[pid]
cand.rrf_score += 1.0 / (k + rank)
cand.source = "both"
else:
cand = _candidate(item, 1.0 / (k + rank))
cand.source = "vector"
candidates[pid] = cand
cand.dense_rank = rank
cand.dense_score = float(item.get("rank", 0.0))
return sorted(candidates.values(), key=lambda c: (-c.rrf_score, c.passage_id))See it on GitHub: kb/retrieval/hybrid.py, lines 53–86
Two details are worth pointing out. Every candidate remembers its rank in each engine and where it came from (bm25, vector or both), and ties are broken by passage_id, so the same question always produces the same order. Both will matter in the next section. GET /search now uses hybrid mode by default, and you can still pass mode=sparse or mode=dense to compare the three on the same query.
Grounded answers and the Context Snapshot
The last piece is POST /ask, implemented in kb/functions/ask.py. It takes the top five fused passages, numbers them from [1] to [5] in the prompt, and asks Claude Sonnet 5.5 to answer using only those passages and to cite them inline. After generation, the code checks which markers actually appear in the answer, and only those passages become citations. It also sets an as_of date, the publication date of the newest cited source, so readers know how fresh the answer is:
# 2. Answer with Claude, citing passages as [1], [2], ...
llm = get_anthropic_client()
if not top_passages:
answer_text = "No relevant passages were found in the knowledge base for this question."
citations: list[dict[str, Any]] = []
elif llm is None:
answer_text = (
f"Retrieved {len(top_passages)} relevant passages with hybrid search. "
"Set ANTHROPIC_API_KEY to get an answer with inline citations."
)
citations = [_citation(i, p) for i, p in enumerate(top_passages, start=1)]
else:
answer_text = await _synthesize(llm, question, top_passages)
# Only the passages the model actually cited become citations.
citations = [
_citation(i, p) for i, p in enumerate(top_passages, start=1) if f"[{i}]" in answer_text
]
# 3. as_of: publication date of the newest CITED source (falls back to all
# retrieved passages, then to now), so readers know how fresh the answer is.
cited_ids = {c["passage_id"] for c in citations}
cited = [p for p in top_passages if p["passage_id"] in cited_ids]
as_of = _newest_date(cited) or _newest_date(top_passages) or datetime.now(tz=UTC)See it on GitHub: kb/functions/ask.py, lines 142–164
The first two branches matter more than they look. With no passages, the endpoint says so instead of letting the model improvise, and without an Anthropic key it still returns the retrieved passages, so retrieval can be tested on its own.
The most important design decision in this endpoint, though, is what it stores. When a RAG system gives a wrong answer, there are only three possible causes.
The right article may never have been ingested. It may be in the database, but retrieval didn't surface it. Or retrieval surfaced the right passage and the model ignored or misread it. Most implementations only log the answer and the passages that were cited, and with that data you can't tell the second case from the third. If the right passage isn't among the citations, was it never retrieved, or retrieved and ignored?
Those problems have completely different fixes, one in retrieval and one in the prompt or the model, and guessing wrong can cost weeks.
So for every answer we save a snapshot with the question, the answer, the citations, the latency, the retriever variant and every candidate returned by both engines, cited or not, with its ranks and fused score. That's why compute_rrf keeps all that information on each candidate. The snapshots live in a new table created by this week's migration:
class Snapshot(Base):
"""One row per answer from `POST /ask`: what was asked, what was answered, and
EVERY candidate passage the retriever surfaced (not only the cited ones), with
its rank in each search engine. That is what lets us tell, later, whether a bad
answer came from ingestion, retrieval or generation.
"""
__tablename__ = "snapshots"
snapshot_id: Mapped[str] = mapped_column(String, primary_key=True)
question: Mapped[str] = mapped_column(Text, nullable=False)
answer_text: Mapped[str] = mapped_column(Text, nullable=False)
citations: Mapped[list[dict[str, object]]] = mapped_column(JSONB, nullable=False, default=list)
retrieved: Mapped[list[dict[str, object]]] = mapped_column(JSONB, nullable=False, default=list)
retriever_variant: Mapped[str] = mapped_column(String, nullable=False)
as_of: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
# Filled in by the evaluators in week 5; always NULL in week 2.
confidence: Mapped[float | None] = mapped_column(Numeric(5, 4), nullable=True)
latency_ms: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)See it on GitHub: kb/db/models.py, lines 137–158
The retriever_variant column (hybrid_rrf_k60 for now) will let us compare answer quality across retriever versions using real questions once we start experimenting, and confidence stays empty until the evaluators of Week 5 fill it in. Here is where ask_question writes the snapshot and sends the event:
# 4. Snapshot first, and committed, BEFORE anything is sent anywhere.
snapshot_id = await record_snapshot(
session,
question=question,
answer_text=answer_text,
citations=citations,
retrieved=all_candidates,
retriever_variant=retriever_variant,
as_of=as_of,
latency_ms=latency_ms,
)
await session.commit()
# 5. Fire-and-forget: the user doesn't wait for the event to be delivered.
task = asyncio.create_task(_emit_answer_produced(snapshot_id, retriever_variant))
_background_tasks.add(task)
task.add_done_callback(_background_tasks.discard)See it on GitHub: kb/functions/ask.py, lines 168–184
There are three small decisions in those lines that we'd recommend copying into your own projects:
The snapshot is committed before the event is sent, so if the event gets lost because the server restarts or Inngest is briefly unreachable, nothing is lost: the database is the source of truth and the event is only a notification.
The task created by asyncio.create_task is stored in a set until it finishes, because Python's event loop only keeps weak references to tasks, and an unreferenced fire-and-forget task can be garbage-collected before it runs.
And if sending the event fails, _emit_answer_produced logs the error instead of swallowing it. Telemetry should never break the user's request, but it shouldn't fail silently either.
Running it end to end
Time to put everything together!
If you're coming from Week 1, pull the latest code, add your ANTHROPIC_API_KEY to .env and run make start, which applies this week's migration automatically. Before you start, set CONCURRENCY_PER_PUBLICATION=1 in .env. Since each feed only has around twenty recent posts, three slots per publication empty the queue quickly, and one slot gives you enough time to watch the priority in action.
With the stack running, trigger the three backfills from a second terminal:
curl -X POST http://localhost:8000/publications -H "Content-Type: application/json" \
-d '{"feed_url": "https://theneuralmaze.com/feed"}'
curl -X POST http://localhost:8000/publications -H "Content-Type: application/json" \
-d '{"feed_url": "https://www.decodingai.com/feed"}'
curl -X POST http://localhost:8000/publications -H "Content-Type: application/json" \
-d '{"feed_url": "https://magazine.sebastianraschka.com/feed"}'Each call returns immediately with an event_id and a trace_url, and you can follow progress with GET /jobs/{event_id}. In the Inngest dashboard at http://localhost:8288 you'll see three backfill-publication runs fan out their articles, followed by ingest-article runs executing one per publication, in parallel, while the rest wait in the queue.
Right after that, while there are still The Neural Maze articles waiting, submit one by hand:
curl -X POST http://localhost:8000/articles -H "Content-Type: application/json" \
-d '{"url": "[FILL: URL of a real Neural Maze post]"}'The new run, tagged via: "manual", starts as soon as the current The Neural Maze run finishes, ahead of every The Neural Maze backfill article still in the queue. The other two publications keep going in their own lanes.
Once everything has finished, run the three replay_backfill.py commands from section 3 and check that the counts match. Then ask a question:
curl -X POST http://localhost:8000/ask -H "Content-Type: application/json" \
-d '{"question": "How do durable execution engines handle failure recovery compared to traditional task queues?"}'Here's what our run returned, shortened:
{
"answer": "**What it is**\n\nThe passages describe durable execution through Inngest, where the **step** is \"the core unit of durable execution\" [1][3]. A step acts as an atomic checkpoint...",
"citations": [
{
"citation_idx": 1,
"passage_id": "87912a1c-6d69-4b26-bd13-9ccdc8f9763c",
"title": "Durable harnesses for graph-powered agents",
"author": "Miguel Otero Pedrido; Antonio Zarauz Moreno",
"url": "https://www.theneuralmaze.com/p/durable-harnesses-for-graph-powered"
},
{
"citation_idx": 3,
"passage_id": "46aa313a-26e5-488a-969c-22af1f2ecc1a",
"title": "Durable harnesses for graph-powered agents",
"author": "Miguel Otero Pedrido; Antonio Zarauz Moreno",
"url": "https://www.theneuralmaze.com/p/durable-harnesses-for-graph-powered"
}
],
"as_of": "2026-09-30T00:00:00+00:00",
"snapshot_id": "snap_3jh6_2nxHqmw",
"latency_ms": 6154
}The citation indexes skip a number, [1] and [3] but no [2], because only the passages Claude actually cited are listed. You might also be surprised by the latency after all the talk about fast paths. Almost none of those six seconds is retrieval: embedding the question, running both searches, fusing them and writing the snapshot together take well under a second.
The rest is Claude writing a long, structured answer. Routing this request through a workflow engine would only have made it slower. If you want it to feel faster, the next step is streaming the answer token by token.
Finally, take a look at what got stored:
SELECT snapshot_id, latency_ms,
jsonb_array_length(retrieved) AS candidates,
jsonb_array_length(citations) AS cited
FROM snapshots
ORDER BY created_at DESC
LIMIT 5;You'll always see far more candidates than citations, and that gap is the whole point. It's the information that will let us tell a retrieval problem from a generation problem when we start evaluating answers.
What's next
At the start of the week we had a pipeline that worked on one newsletter and five articles. Now it ingests three publications in parallel without overloading their servers or our OpenAI quota, lets urgent articles skip the queue, survives replays without spending a cent, searches with two engines at once and keeps a detailed record of every answer it gives.
What it can't do yet is tell us whether those answers are any good. Every snapshot we store fires a kb/answer.produced event that nobody is listening to. In Week 5 we'll plug evaluators into that event and turn the snapshots table into a real measurement of retrieval quality and hallucinations, so we find out about problems before our users do. Before that, Week 3 keeps the knowledge base fresh with a poller and exposes it through an MCP server.
All the code from this lesson is in the course repository, under the week-2 tag, with the full walkthrough in docs/week-2.md.
The Substack Brain, lesson by lesson
Lesson 1: A durable ingestion pipeline (article · video)
→ Lesson 2: Flow control and hybrid retrieval (you are here)
Lesson 3: A freshness poller and an MCP server (next week)
Lesson 4: A knowledge graph in Memgraph
Lesson 5: Evals built from real failures
Lesson 6: A research agent in production
Our sponsor
The Substack Brain is made possible by Inngest, the durable execution engine behind this course. Every step's result is recorded, so a crashed run resumes instead of restarting, and the flow control we used this week (concurrency keys, throttling and priority) is configured right on the function, with no queues or workers of your own to run.
If you're building AI pipelines or agents that need to survive the real world, give Inngest a try. Thanks to the Inngest team for supporting this course 💙
See you next week!















An honor to be in your Substack brain. You are crushing it with this new series 🫶