Introduction
Semantic caching is one of the best cost levers you have in an LLM-powered product. Instead of matching prompts character by character, you embed each incoming question, look for a previously answered question that sits close in vector space, and return the stored answer. Done well, it cuts inference spend dramatically and shaves seconds off latency.
But there is a quiet failure mode that most tutorials skip.
Imagine a support agent that answers "What is the refund window for annual plans?" with "30 days". Three weeks later your product team changes the policy to 14 days and updates a row in PostgreSQL. The cache has no idea. The next customer asks "How long do I have to get my money back on a yearly subscription?", the cache finds a near-perfect semantic match, and your agent confidently promises 30 days. It may even take action based on that promise, such as issuing a refund or drafting a contract clause.
That is context poisoning: your agent's memory has drifted away from your source of truth, and nothing in the system knows it.
In this guide, we will build a real-time defence against it. You will learn how to:
- Capture row-level changes from PostgreSQL using Change Data Capture (CDC) through logical replication.
- Stream those events into an isolated Node.js worker.
- Use both exact lineage tracking and an embedding similarity radius to find compromised cache entries.
- Flush them before any user or agent can retrieve them.

What Is Context Poisoning, Really?
Most people hear "poisoning" and think of adversarial attacks, such as prompt injection or malicious documents slipped into a knowledge base. Those are real problems, but context poisoning in the caching sense is far more mundane, and far more common.
It happens whenever three conditions line up:
- Your agent answers from retrieved context, either from a vector store or from a cached previous answer.
- That context was valid at write time.
- The authoritative data changed afterwards, and the cache was not told.
The danger is not that the answer is wrong. Wrong answers happen. The danger is that the answer is wrong and looks authoritative, because it came from your own trusted infrastructure. LLMs are very good at weaving a stale fact into a fluent, convincing response. Users rarely double-check. Downstream tools, such as agents that call APIs, update tickets, or send emails, act on that fluent output without hesitation.
Why Traditional Invalidation Fails Here
Classic cache invalidation relies on keys. You change user 42, so you delete the key user:42. That works when the cache key is a deterministic function of the data.
Semantic caches break that model in three ways:
- The key is a vector, not a name. There is no
refund_policykey to delete. The same fact may be embedded inside hundreds of differently phrased questions. - The answer is a synthesis. An LLM response may blend several rows. A change to any one of them can silently invalidate the whole answer.
- Paraphrases multiply. One underlying fact can live in many cache entries, all of which need to die together.
TTL-based expiry is the usual shortcut, and it forces an ugly trade-off. Short TTLs wipe out your hit rate, so you pay for inference again. Long TTLs leave a freshness gap that can last hours or days. Neither is acceptable when the cache feeds an autonomous agent.
The Architecture at a Glance
The solution has four moving parts:
- PostgreSQL is the source of truth. It emits every committed row change through its write-ahead log (WAL).
- A logical replication slot exposes those changes as a stream of structured events using the built-in
pgoutputprotocol. - An isolated Node.js worker consumes the stream, groups events by transaction, and decides which cache entries are compromised.
- A vector-capable cache, here Redis with a vector index, stores query embeddings, answers, and lineage metadata.
The key design decision is isolation. The invalidation worker runs as its own process, with its own credentials and its own failure domain. Your request-serving path never waits on it, and a crash in the worker cannot take down your API. It also gives you one clean place to monitor freshness.
The worker uses a two-layer strategy to find victims:
- Layer 1, lineage invalidation (exact). When a cache entry is written, you record which table and primary keys contributed to the answer. A change to one of those rows deletes the entry directly.
- Layer 2, semantic radius invalidation (fuzzy). For changes the lineage graph cannot see, such as a brand new row that should now influence old answers, you embed the changed row and flush every cached entry within a calculated radius.
Layer 1 is cheap and precise. Layer 2 is your safety net for everything lineage cannot know about.
Step 1: Prepare PostgreSQL for Logical Replication
First, enable logical decoding. On a self-managed instance, update postgresql.conf:
wal_level = logical
max_replication_slots = 10
max_wal_senders = 10
# Protect your disk from a stuck consumer
max_slot_wal_keep_size = 10GB
On managed services (RDS, Cloud SQL, Neon, Supabase), you enable the equivalent through a parameter group or dashboard flag, for example rds.logical_replication = 1 on AWS.
Next, create a dedicated role, a publication, and a slot. Keep the publication scoped to only the tables your agent actually reads from.
-- A least-privilege role for the invalidation worker
CREATE ROLE cache_invalidator WITH LOGIN REPLICATION PASSWORD 'change-me';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO cache_invalidator;
-- Only publish what the agent depends on
CREATE PUBLICATION agent_kb_pub
FOR TABLE policies, products, faq_articles
WITH (publish = 'insert, update, delete');
-- A named slot so the worker can resume exactly where it stopped
SELECT pg_create_logical_replication_slot('agent_cache_slot', 'pgoutput');
There is one setting people forget, and it matters a lot here. By default, an UPDATE event only carries the old values of the primary key. For semantic comparison you want the full before and after picture:
ALTER TABLE policies REPLICA IDENTITY FULL;
ALTER TABLE products REPLICA IDENTITY FULL;
ALTER TABLE faq_articles REPLICA IDENTITY FULL;
With REPLICA IDENTITY FULL, each update event includes the complete old row and the complete new row. That lets you embed both versions and measure exactly how far the meaning moved, which we will use later.
Step 2: Set Up the Semantic Cache with Lineage Metadata
We will use Redis with a vector index. The schema below stores the query embedding, the generated answer, and a few fields we need for invalidation.
Install the dependencies:
npm install redis pg-logical-replication openai pino
npm install -D typescript tsx @types/node
Create the index once at startup:
// cache/schema.ts
import { createClient } from "redis";
export const redis = createClient({ url: process.env.REDIS_URL });
await redis.connect();
export async function ensureIndex() {
try {
await redis.sendCommand([
"FT.CREATE", "idx:semcache",
"ON", "HASH",
"PREFIX", "1", "semcache:",
"SCHEMA",
"answer", "TEXT",
"created_at", "NUMERIC", "SORTABLE",
"embedding", "VECTOR", "HNSW", "10",
"TYPE", "FLOAT32",
"DIM", "1536",
"DISTANCE_METRIC", "COSINE",
"M", "16",
"EF_CONSTRUCTION", "200",
]);
} catch (err: any) {
if (!String(err.message).includes("Index already exists")) throw err;
}
}
Now the important part: writing a cache entry that remembers where its answer came from. Whenever your RAG pipeline generates a response, it already knows which rows it retrieved. Capture that.
// cache/write.ts
import { redis } from "./schema";
import { randomUUID } from "node:crypto";
export interface SourceRef {
table: string;
pk: string;
}
export async function writeCacheEntry(opts: {
embedding: Float32Array;
answer: string;
sources: SourceRef[];
ttlSeconds?: number;
}) {
const id = randomUUID();
const key = `semcache:${id}`;
await redis
.multi()
.hSet(key, {
answer: opts.answer,
created_at: Date.now(),
embedding: Buffer.from(opts.embedding.buffer),
sources: JSON.stringify(opts.sources),
})
// Safety-net TTL: CDC is the primary mechanism, TTL is the backstop
.expire(key, opts.ttlSeconds ?? 60 * 60 * 24)
// Reverse index: which cache entries depend on this row?
.exec();
for (const s of opts.sources) {
await redis.sAdd(`lineage:${s.table}:${s.pk}`, key);
}
return key;
}
Notice the reverse index. For every source row, we keep a Redis set of cache keys that depend on it. That is what makes Layer 1 invalidation an O(1) lookup instead of a scan.

Step 3: Stream Row-Level Changes into a Node.js Worker
Now we build the isolated worker. We use pg-logical-replication, which speaks the pgoutput protocol natively and gives us typed events for begin, insert, update, delete, and commit.
There is one design rule that will save you from subtle bugs: never invalidate on individual row events. Buffer them per transaction and process on commit. If you flush on an uncommitted row and the transaction later rolls back, you have invalidated for nothing. Worse, if you process mid-transaction, a concurrent reader can repopulate the cache with data that is about to change again.
// worker/cdc.ts
import {
LogicalReplicationService,
PgoutputPlugin,
} from "pg-logical-replication";
import pino from "pino";
import { processChangeBatch, type RowChange } from "./invalidate";
const log = pino({ name: "cache-invalidator" });
const service = new LogicalReplicationService(
{
connectionString: process.env.PG_REPLICATION_URL,
},
{
acknowledge: { auto: false, timeoutSeconds: 0 },
}
);
const plugin = new PgoutputPlugin({
protoVersion: 1,
publicationNames: ["agent_kb_pub"],
});
let buffer: RowChange[] = [];
service.on("data", async (lsn: string, msg: any) => {
switch (msg.tag) {
case "begin":
buffer = [];
break;
case "insert":
buffer.push({
op: "insert",
table: msg.relation.name,
newRow: msg.new,
});
break;
case "update":
buffer.push({
op: "update",
table: msg.relation.name,
oldRow: msg.old,
newRow: msg.new,
});
break;
case "delete":
buffer.push({
op: "delete",
table: msg.relation.name,
oldRow: msg.old,
});
break;
case "commit": {
const batch = buffer;
buffer = [];
try {
await processChangeBatch(batch);
// Only acknowledge once the cache is actually safe
await service.acknowledge(lsn);
} catch (err) {
log.error({ err }, "invalidation failed, will retry from last LSN");
// Do not acknowledge: Postgres will redeliver after restart
process.exit(1);
}
break;
}
}
});
service.on("error", (err) => {
log.error({ err }, "replication stream error");
process.exit(1);
});
log.info("starting CDC subscription");
await service.subscribe(plugin, "agent_cache_slot");
Two details deserve attention here.
First, we acknowledge the LSN (log sequence number) only after invalidation succeeds. If the worker crashes mid-flight, PostgreSQL redelivers the unacknowledged transactions on reconnect. That gives you at-least-once delivery. Because invalidation is idempotent (deleting an already-deleted key is harmless), at-least-once is exactly what you want.
Second, we exit on failure rather than limping along. A process manager like systemd, Kubernetes, or PM2 restarts the worker and it resumes from the last acknowledged position. A silently broken invalidator is far more dangerous than a loudly restarting one.
Step 4: Layer 1, Exact Lineage Invalidation
The first layer is straightforward. For every changed row, look up its lineage set and delete everything in it.
// worker/invalidate.ts
import { redis } from "../cache/schema";
export interface RowChange {
op: "insert" | "update" | "delete";
table: string;
oldRow?: Record<string, unknown>;
newRow?: Record<string, unknown>;
}
function pkOf(row: Record<string, unknown>): string {
return String(row.id);
}
async function invalidateByLineage(change: RowChange): Promise<number> {
// Inserts have no prior lineage, so Layer 2 handles them
if (change.op === "insert") return 0;
const row = change.oldRow!;
const setKey = `lineage:${change.table}:${pkOf(row)}`;
const keys = await redis.sMembers(setKey);
if (keys.length === 0) return 0;
await redis.multi().del(keys).del(setKey).exec();
return keys.length;
}
Lineage invalidation handles updates and deletes of rows that your pipeline knowingly used. It is fast, exact, and has zero false positives.
But it has a blind spot, and it is a big one. Consider these cases:
- A new row is inserted that should now influence answers to existing questions.
- An answer was cached with no retrieved sources, because the first retrieval found nothing relevant.
- A row that was previously irrelevant is edited so that it now is relevant.
Lineage cannot see any of those. That is where the second layer comes in.
Step 5: Layer 2, The Embedding Similarity Radius
This is the part that turns a basic cache buster into something genuinely robust.
The idea: when a row changes, convert it to a vector, then find every cached query whose embedding is close enough that the changed content could plausibly alter its answer. Flush all of them.
The only real question is: how close is close enough? Pick a radius too small and stale answers slip through. Pick one too large and you destroy your hit rate. Let's derive a principled number instead of guessing.
The Math Behind a Safe Radius
Work in angular distance, because it obeys the triangle inequality, unlike raw cosine distance. For two unit-normalised embeddings a and b:
cos_sim(a, b) = a Β· b
angle(a, b) = arccos(cos_sim(a, b))
cosine_dist = 1 - cos_sim(a, b)
Define three quantities:
q = an incoming user query embedding
c = the cached query embedding that q hits
d = the embedding of a changed document or row
alpha = max angle at which a cache hit is allowed
(alpha = arccos(hit_threshold))
rho = max angle at which a document is considered
relevant to a query
(rho = arccos(retrieval_threshold))
A cached entry c is dangerous if some future query q could hit it (so angle(q, c) <= alpha) and that same query would have retrieved the changed document (so angle(q, d) <= rho).
By the triangle inequality on angular distance:
angle(c, d) <= angle(c, q) + angle(q, d)
<= alpha + rho
So any cache entry that could ever be poisoned by d lies within:
r_invalidate = alpha + rho
of d. That is your safe invalidation radius. Anything outside it cannot possibly be reached by a query that is both a cache hit and relevant to the changed row.
Finally, convert it back to a cosine distance, which is what the vector index expects:
cosine_dist_radius = 1 - cos(alpha + rho)
A Concrete Example
Suppose your cache serves a hit when cosine similarity is at least 0.92, and your retriever treats documents as relevant at similarity 0.80 or higher.
alpha = arccos(0.92) β 0.4027 rad (about 23.1Β°)
rho = arccos(0.80) β 0.6435 rad (about 36.9Β°)
r_invalidate = 0.4027 + 0.6435 = 1.0462 rad (about 59.9Β°)
cosine_dist_radius = 1 - cos(1.0462) β 0.499
That is a wide net, and it is intentionally conservative. It guarantees safety, but it will flush more than strictly necessary. In practice you tighten it using real traffic, which we cover in the best practices section. Think of this value as your provably safe ceiling, with tuned values sitting below it.

Implementing the Radius Calculation
// worker/radius.ts
export function angularRadius(hitThreshold: number, retrievalThreshold: number) {
const alpha = Math.acos(hitThreshold);
const rho = Math.acos(retrievalThreshold);
return alpha + rho; // radians
}
// Redis COSINE metric returns distance = 1 - cos_sim
export function toCosineDistance(angleRad: number, safetyMargin = 0): number {
const clamped = Math.min(angleRad + safetyMargin, Math.PI);
return 1 - Math.cos(clamped);
}
export const INVALIDATION_DISTANCE = toCosineDistance(
angularRadius(
Number(process.env.CACHE_HIT_THRESHOLD ?? 0.92),
Number(process.env.RETRIEVAL_THRESHOLD ?? 0.8)
)
);
Embedding the Changed Row and Running a Range Query
For updates, embed both the old and new versions. A row can leave one semantic neighbourhood and enter another, and the cache entries near either location may be stale.
// worker/semantic.ts
import OpenAI from "openai";
import { redis } from "../cache/schema";
import { INVALIDATION_DISTANCE } from "./radius";
const openai = new OpenAI();
function rowToText(table: string, row: Record<string, unknown>): string {
// Keep this deterministic and aligned with how your RAG pipeline
// serialises rows before embedding them for retrieval.
return `${table}: ` + Object.entries(row)
.filter(([, v]) => v !== null && typeof v !== "object")
.map(([k, v]) => `${k}=${v}`)
.join("; ");
}
async function embed(text: string): Promise<Buffer> {
const res = await openai.embeddings.create({
model: "text-embedding-3-small",
input: text,
});
return Buffer.from(new Float32Array(res.data[0].embedding).buffer);
}
export async function invalidateBySemanticRadius(
table: string,
rows: Array<Record<string, unknown> | undefined>
): Promise<number> {
const victims = new Set<string>();
for (const row of rows) {
if (!row) continue;
const vec = await embed(rowToText(table, row));
const result: any = await redis.sendCommand([
"FT.SEARCH", "idx:semcache",
"@embedding:[VECTOR_RANGE $r $vec]=>{$YIELD_DISTANCE_AS: dist}",
"PARAMS", "4", "r", String(INVALIDATION_DISTANCE), "vec", vec,
"RETURN", "1", "dist",
"LIMIT", "0", "10000",
"DIALECT", "2",
]);
// result[0] is the count, followed by key, fields, key, fields...
for (let i = 1; i < result.length; i += 2) victims.add(result[i]);
}
if (victims.size === 0) return 0;
await redis.del([...victims]);
return victims.size;
}
Wiring Both Layers Together
// worker/invalidate.ts (continued)
import { invalidateBySemanticRadius } from "./semantic";
export async function processChangeBatch(batch: RowChange[]) {
let lineageHits = 0;
let semanticHits = 0;
for (const change of batch) {
lineageHits += await invalidateByLineage(change);
semanticHits += await invalidateBySemanticRadius(change.table, [
change.oldRow,
change.newRow,
]);
}
console.log(
JSON.stringify({
msg: "invalidation complete",
changes: batch.length,
lineageHits,
semanticHits,
})
);
}
At this point, every committed change in PostgreSQL results in two things happening within milliseconds: precise removal of known dependents, and a wide sweep of anything semantically nearby. Your agent can no longer read from a cache that disagrees with the database for longer than the pipeline's end-to-end latency.
A Real-World Walkthrough
Let's trace the refund policy example through the finished system.
- A customer asks, "What is the refund window for annual plans?" The agent retrieves
policies.id = 17, answers "30 days", and writes a cache entry. The entry stores lineagepolicies:17. - A product manager runs
UPDATE policies SET refund_days = 14 WHERE id = 17. - PostgreSQL commits and emits the change. The worker receives
begin,update,commit. - Layer 1 looks up
lineage:policies:17, finds the cached entry, and deletes it. - Layer 2 embeds the old and new policy text. It finds another cached entry, "Can I get my money back on a yearly subscription?", which was cached without lineage because of an earlier retrieval miss. It is inside the radius, so it is flushed too.
- The worker acknowledges the LSN.
- The next customer asks about refunds. The cache misses, the agent retrieves the fresh row, and answers "14 days".
Without Layer 2, step 5 would have left a poisoned entry sitting in the cache, waiting to hurt someone.
Best Practices
Operational Safety
- Monitor replication slot lag. An unconsumed slot makes PostgreSQL retain WAL indefinitely. Alert on
pg_replication_slotslag and setmax_slot_wal_keep_sizeso a stuck worker cannot fill your disk. - Run at least one hot standby worker. Only one consumer can attach to a slot at a time, so use a leader election mechanism or a supervisor that restarts quickly.
- Make every invalidation idempotent. At-least-once delivery means duplicates will happen. Deleting a missing key must be harmless.
- Keep a TTL as a backstop. CDC is your primary defence, and a generous TTL is the seatbelt for the day something unexpected breaks.
Accuracy and Tuning
- Replay real traffic to tune the radius. Take a week of production queries and changes, then measure how many flushed entries would actually have returned a different answer. Tighten the radius until false flushes drop without letting any true positives escape.
- Keep embedding models consistent. The worker must use the same model and the same row serialisation as your retrieval pipeline. Mixing models makes distances meaningless.
- Batch embeddings. When a transaction touches many rows, embed in batches to stay under rate limits and reduce latency.
- Version your cache namespace. When you change embedding models or prompt templates, bump a namespace version so old vectors never mix with new ones.
Security
- Use a least-privilege replication role. The worker needs
REPLICATIONand read access, and nothing else. - Treat CDC payloads as sensitive. They contain raw row data. Do not log full rows in production, and redact personal information before embedding if your policy requires it.
Common Mistakes to Avoid
- Invalidating on uncommitted rows. Always buffer by transaction and act on
commit. Acting early creates race conditions and needless flushes after rollbacks. - Forgetting
REPLICA IDENTITY FULL. Without it, update events lack the old row values, and Layer 2 can only embed the new version. You then miss entries that were close to the old meaning. - Using raw cosine distance in the triangle inequality. Cosine distance is not a true metric. Do the bound in angular space, then convert for the query.
- Acknowledging the LSN too early. If you acknowledge before the cache is cleaned and the worker crashes, those changes are lost forever and poisoned entries survive.
- Publishing every table. A broad publication floods the worker with irrelevant events. Scope it to the tables your agent actually reads.
- Relying only on TTL. It feels simple and safe, but it is neither. It either kills your hit rate or leaves a long freshness gap.
- Ignoring agent-written data. If agents can write back to the database, their own writes will trigger invalidations. Make sure your pipeline tolerates that feedback loop without thrashing.
π Pro Tips
- Add a freshness fence for critical queries. For high-stakes intents like pricing, legal, or medical content, store the source row's
updated_atin the cache entry. On a hit, do a cheap primary-key lookup and compare timestamps before serving. It costs one indexed read and gives you a hard guarantee. - Weight the radius per table. A policy table deserves a wide radius. A table of blog comments does not. Store a per-table radius multiplier and tune each independently.
- Emit invalidation metrics. Track flushes per change, lineage versus semantic ratio, and end-to-end lag from commit to cache clean. These numbers tell you whether the radius is too loose or too tight.
- Debounce chatty tables. If a row updates dozens of times a minute, collapse events within a short window so you embed once instead of dozens of times.
- Consider Debezium when you scale out. If many services need the same change stream, Debezium and Kafka let multiple consumers share it without competing for a single slot.
- Test with chaos. Kill the worker mid-transaction in staging and confirm no poisoned entry survives the restart. That one test catches most acknowledgement bugs.
π Key Takeaways
- Context poisoning is a freshness problem disguised as an AI problem. The model is fine, but the memory it reads from is out of date.
- PostgreSQL logical replication with
pgoutputgives you committed, ordered, row-level change events without touching your application code. - Buffer events per transaction and acknowledge only after the cache is clean, so crashes never lose an invalidation.
- Use two layers together: exact lineage for known dependencies and an embedding radius for everything lineage cannot see.
- A safe radius is not a guess. It is the angular hit threshold plus the angular retrieval threshold, converted back to cosine distance, and then tuned down using real traffic.
- Keep a TTL as a backstop, and monitor slot lag so your safety net does not become an outage.
Conclusion
Semantic caching earns its place in modern AI stacks because it saves real money and real time. But a cache that cannot tell when the world has changed is not a performance optimisation, it is a liability. As agents take on more autonomy, answering questions, filing tickets, and triggering workflows, the cost of one confidently stale answer climbs quickly.
The pattern in this guide is small enough to run as a single worker and strong enough to give you a real freshness guarantee. PostgreSQL already knows when your data changes. CDC lets you listen. Lineage tracking handles the entries you know about, and the similarity radius handles the ones you do not. Together they keep your agent's memory aligned with your source of truth, measured in milliseconds rather than hours.
Start small. Publish one table, wire up the worker, and watch the invalidation logs for a few days. You will quickly see how many stale answers your old TTL was quietly letting through.
References
- PostgreSQL Documentation, Logical Replication: https://www.postgresql.org/docs/current/logical-replication.html
- PostgreSQL Documentation, Logical Decoding and the
pgoutputProtocol: https://www.postgresql.org/docs/current/logicaldecoding.html - PostgreSQL Documentation, Replica Identity: https://www.postgresql.org/docs/current/sql-altertable.html#SQL-ALTERTABLE-REPLICA-IDENTITY
pg-logical-replicationon npm: https://www.npmjs.com/package/pg-logical-replication- Redis Documentation, Vector Search and Range Queries: https://redis.io/docs/latest/develop/interact/search-and-query/advanced-concepts/vectors/
- Debezium Documentation, PostgreSQL Connector: https://debezium.io/documentation/reference/stable/connectors/postgresql.html
- OpenAI Documentation, Embeddings: https://platform.openai.com/docs/guides/embeddings