Moving AI analysis out of an HTTP request improves the request path. Queueing that work does not make it correct.
In my AI Support Assistant, a new message triggers BullMQ analysis of the ticket conversation. The request waits for database writes and queue publication, then returns without waiting for the AI results.
But a retryable job can still run more than once, reload different context, race with another job, repeat provider work or overwrite newer results. The message can also be saved without its analysis job reaching Redis.
Examining the current implementation makes those boundaries concrete. The scenarios below are correctness risks exposed by the code, not production incidents; the proposed fixes are not implemented features.
The System This Came From
The first article described the broader React, Node.js, PostgreSQL, Redis and OpenAI architecture.
This article focuses on one workflow: a message is added to a support ticket, and background processing generates a summary, sentiment and priority. Those values are written back to the ticket. The worker also records an AIInteraction containing the conversation, summary response and summary token usage.
The queue is named ai-processing. Its job type is generate-ticket-summary, although the processor does more than summarize. It analyzes the whole ticket conversation, rather than just the message that triggered publication.
Why the AI Work Went to BullMQ
Saving a message and enriching a ticket have different completion requirements. The message is the user-authored record. AI metadata can arrive later.
The current boundary looks like this:
React client → Node.js API → PostgreSQL: Message + Ticket timestamp
→ Redis/BullMQ: publish analysis job
→ HTTP 201 after publication
Redis/BullMQ → separate worker → load ticket messages
→ summary → sentiment → priority
→ update Ticket → insert AIInteraction
Publication makes the job available to the worker. The worker may start before HTTP 201 is sent; the diagram shows two execution paths, not a guarantee that background processing begins after the response.
The important boundary is what the request awaits: persistence and publication, followed by a return to the caller. Summary, sentiment and priority belong to the worker’s lifecycle.
The Actual Producer Flow
From backend/src/modules/messages/message.service.ts, lines 26–58. Authorization precedes this excerpt. The excerpt is otherwise complete; blank lines are condensed.
const message = await createMessageRepository({
content: data.content,
ticketId: data.ticketId,
senderId: data.senderId,
});
await updateTicketRepository(data.ticketId, {
updatedAt: new Date(),
});
await aiQueue.add(
GENERATE_TICKET_SUMMARY_JOB,
{
ticketId: data.ticketId,
messageId: message.id,
},
{
attempts: 3,
backoff: {
type: "exponential",
delay: 5000,
},
removeOnComplete: true,
removeOnFail: false,
},
);
return message;
The repository helpers execute prisma.message.create and prisma.ticket.update. These are separate writes, followed by publication to Redis. There is no encompassing transaction or outbox in this flow.
The controller awaits this service before returning HTTP 201. Consequently, a message can already exist even if the service does not reach its successful return. Queue publication is part of request completion; AI processing is not.
The payload includes messageId, but the worker’s processing logic only consumes ticketId. The message ID currently provides neither a processing boundary nor an idempotency key.
What the Worker Actually Does
The worker dispatches generate-ticket-summary to processTicketSummaryJob using only the ticket ID. Its configuration in backend/src/queues/ai.worker.ts, lines 71–75, is:
{
connection: redis as any,
concurrency: 5,
},
This is the complete options object from the worker constructor, with blank lines condensed.
Five jobs can run concurrently within this worker instance, including jobs for the same ticket. This is a local worker limit, not a global or per-ticket limit. See BullMQ’s worker concurrency guide.
The processor loads messages with where: { ticketId }, includes sender role and name, and sorts by createdAt: "asc". It returns early when there are no messages. Otherwise it formats the conversation and executes this sequence from lines 113–127:
const aiResult = await generateTicketSummary(
{
conversation,
},
aiContext,
);
const sentimentResult = await generateSentiment(
{
conversation,
},
aiContext,
);
const priorityResult = await detectTicketPriority(conversation, aiContext);
These operations are sequential within one job. Different jobs can still overlap. The service routes generation through an OpenAI-primary orchestrator with optional Gemini fallback for eligible failures, so three logical operations do not necessarily mean exactly three external HTTP requests.
After sentiment and priority validation, lines 133–155 perform two separate writes:
await prisma.ticket.update({
where: {
id: ticketId,
},
data: {
aiSummary: aiResult.summary,
sentiment,
priority,
},
});
await createAIInteraction({
ticketId,
prompt: conversation,
response: aiResult.summary || "",
tokensUsed: aiResult.tokensUsed,
});
Blank lines are condensed. createAIInteraction is a plain Prisma insert. There is no job identity or conversation version attached to that record, and no conditional version check on the ticket update.
The worker runs through a separate entry point in backend/src/worker.ts. That separates its execution lifecycle from HTTP handling; it does not add coordination between jobs for one ticket.
Retries Solve Failure Recovery, Not Correctness
The producer configures three total attempts, including the initial execution, with exponential backoff and a five-second base. For ordinary thrown-error retries, that gives nominal delays of five and ten seconds before the remaining attempts; scheduling can add waiting time. See BullMQ’s retry documentation.
Completed jobs are removed. Failed jobs are retained. Those choices affect recovery and diagnosis, but they do not make processing idempotent.
Retry asks, “Should I try this work again after failure?” Idempotency asks, “If I perform it again, will the resulting state remain correct?” Deduplication asks, “Should equivalent work be admitted or executed again at all?”
Here, retries restart a processor that calls providers and performs database writes. The retry configuration supplies no rule for deciding whether those effects have already happened.
Failure Mode 1: The Database Commits but the Queue Does Not
Consider the producer’s actual order:
- The message insert succeeds.
- The ticket timestamp update succeeds.
- Publication rejects, or the process exits before publication completes.
The database now contains the conversation change without assured corresponding queue work. Retrying a worker cannot recover a job that was never published.
A publication error can also leave an ambiguous outcome: Redis might have accepted the job before the caller lost its acknowledgement. Recovery must handle both missing work and repeated publication.
The earlier boundary also matters: if the ticket update fails after message creation, the message remains persisted and publication is never reached. Wrapping those two database writes in a transaction would address that database-only partial state. It would not atomically commit a Redis job with PostgreSQL data.
A proposed transactional outbox changes where publication intent becomes durable:
BEFORE — CURRENT
Message insert → Ticket update → Redis publication
database writes complete can fail separately
IMPROVED — PROPOSED
PostgreSQL transaction:
Message insert + Ticket update + Outbox event
↓ committed publication intent
Outbox dispatcher → Redis/BullMQ → worker
↓
mark event published after successful queue publication
The dispatcher retries unpublished events. If it publishes and crashes before marking the event, it may publish again. An outbox therefore needs stable event identity and duplicate-safe consumption; it does not confer exactly-once execution.
This adds a table, dispatcher lifecycle, backlog monitoring and retention policy. For optional enrichment in an early application, explicit reconciliation may be an interim choice. If every accepted message must eventually be analyzed, I would prioritize an outbox regardless of traffic volume. Reliability requirements, rather than scale alone, determine when it is needed.
Failure Mode 2: Retrying AI Work Is Not Automatically Idempotent
Suppose summary generation succeeds and sentiment generation fails with an error that escapes the provider layer. BullMQ can retry the processor. It starts again by loading messages and generating the summary. There is no checkpoint that reuses the previous summary.
Now move the failure later. Summary, sentiment and priority succeed; the ticket update commits; the interaction insert fails. The ticket already contains AI metadata, but the job has failed. A retry repeats the AI sequence and attempts both database writes again.
The resulting values need not match the previous output. More messages may also have arrived before the retry reloads the conversation. The same job payload therefore does not necessarily identify the same input snapshot.
A duplicate execution that reaches the insert can append another AIInteraction: its schema uses a generated UUID and has no unique analysis key. Interruption after application writes but before queue completion is recorded creates another window for repeated effects. The relevant effects here are provider requests, ticket metadata writes and interaction records.
Deterministic job identity helps at publication
A proposed identity such as ticket-analysis-<messageId> could suppress repeated publication for the same saved message while that queue job exists. It would not suppress jobs triggered by two different messages, nor deduplicate message creation if a repeated HTTP request creates a new message ID.
BullMQ’s job ID documentation explains that an existing custom ID prevents another job with that ID from being added. Once the job is removed, the ID no longer protects against duplicates. That is directly relevant to this application’s removeOnComplete: true setting.
Durable idempotency belongs at the application effect
For this workflow, I would define a logical analysis key from ticket identity, conversation revision and analysis/prompt version. A durable record with a unique constraint could track whether that analysis has already been applied.
The ticket update, interaction insert and successful analysis marker should commit together, with concurrent duplicate attempts resolved by database constraints and conditional writes. A separate initial “already done?” read is insufficient because two workers can both pass it.
That would protect database effects. Provider requests made before commit could still repeat. Reusing stored generation results would require checkpointing and invalidation rules; I would start with atomic result persistence and durable identity.
Failure Mode 3: Ordering Is a Business Constraint
Two messages arrive quickly on ticket T:
- Message A is saved and job A starts. It reads the conversation containing A.
- Message B is saved and job B starts in another available worker slot. It reads A and B.
- Job B finishes first and writes metadata based on A and B.
- Job A finishes later and overwrites that metadata using its older context.
The current ticket update matches only id: ticketId. Nothing prevents step four. This is a possible stale, last-writer-wins result visible from the architecture.
Both jobs could also start after B exists, read overlapping or identical context, and perform redundant analysis. Because messageId is unused, neither job is restricted to the conversation as it stood when its triggering message was saved.
Sorting messages inside each job does not order jobs or their writes. The timestamp sort also has no explicit tie-breaker for equal timestamps.
Global throughput and per-ticket correctness are separate constraints. Lowering one worker’s concurrency to one would reduce overlap there, but also serialize unrelated tickets and would not establish a cross-process ordering policy.
Idempotency, Deduplication and Ordering Are Different
| Concept | Question it answers | Example in this system |
|---|---|---|
| Retry | Should failed work be attempted again? | Retry a ticket analysis after a transient failure. |
| Deduplication | Should equivalent work be admitted again? | Proposed: suppress publication while the same job ID exists. |
| Idempotency | Will repeating the work preserve correct state? | Proposed: commit one analysis record for a ticket revision. |
| Ordering | In what sequence may work take effect? | Proposed: apply ticket analyses in conversation-revision order. |
| Concurrency | How much work may run at the same time? | Allow five active jobs within this worker instance. |
Only retries and worker concurrency are configured in the current path. The other examples describe proposed behavior. A version check can reject stale results without executing every job in strict order; serialization still needs a policy for retries and earlier failures.
A Better Per-Ticket Processing Model
The following is a proposed model, not current code. For summary, sentiment and priority, the business goal is usually a current conversation view. It may not require publishing every intermediate analysis.
Messages A, B, C → atomically advance ticket conversation revision
↓
record latest requested revision
↓
coalesce pending work for this ticket
↓
one active analysis per ticket, if needed
↓
read consistent snapshot at revision R
↓
generate AI results
↓
atomic apply only if revision is still R
↙ ↘
stale: discard current: commit
↓ results + analysis record
ensure latest revision remains scheduled
Version checks are the first correctness protection I would add. A dedicated conversation revision should advance atomically with message persistence. The worker must read the revision and messages as a consistent snapshot, then conditionally apply results in a short database transaction. A separate check followed by an unconditional update leaves a race.
I would avoid using Ticket.updatedAt as the conversation revision. Prisma marks it @updatedAt, so AI writes and unrelated ticket edits also change it. A dedicated revision makes the input contract explicit. The tradeoff is schema and transaction work; discarded stale analyses still consume processing. This protection is useful now because the current concurrency already permits stale writes.
Coalescing can reduce unnecessary work. Several pending triggers can represent one request to analyze the latest conversation. A short debounce period trades freshness for fewer intermediate analyses. BullMQ offers deduplication modes, but their lifetime semantics must match the application.
In particular, simply suppressing every new job while one analysis is active can miss updates arriving after that analysis loaded its messages. The design must durably track the latest requested revision and arrange another pass when needed. Checking once at the end without an atomic handoff can still lose a concurrent update. Coalescing is a pragmatic next step when message bursts cause redundant work; it is not a substitute for version checks.
Per-ticket serialization can limit overlap. A coordinated lease or partitioned processing scheme can allow different tickets to progress concurrently while restricting one ticket’s active analysis. An in-memory mutex covers only one process. Distributed leases introduce expiry, renewal and crash-recovery concerns; conditional version checks remain valuable when a lease expires while work is still running.
If a future job must act on every message in business order, it needs explicit sequence tracking and a policy for earlier failures. Coalescing intermediate analyses would no longer satisfy that requirement.
What I Would Change First
These are recommendations for this application, not features already implemented.
- Define analysis identity and protect freshness together. Introduce a conversation revision, conditional result application and an atomic transaction for ticket metadata, interaction and an analysis marker. Then use deterministic publication identity for the chosen unit of work. This addresses correctness exposed by today’s concurrency; it costs a schema change and careful transaction design.
- Coalesce pending ticket analysis where the product permits it. Treat the job as “refresh this ticket” rather than “produce every intermediate result,” while preserving a durable follow-up for changes during active processing. This costs scheduling complexity and may delay freshness slightly. It becomes more useful when bursts create redundant work. Add coordinated per-ticket serialization if overlap remains significant; version protection comes first.
- Close the database-to-queue gap when eventual analysis is required. Add an outbox and replay-safe dispatcher. Also review producer connection policy: the shared Redis module sets
maxRetriesPerRequest: null, and the message service has no explicit publication deadline. BullMQ’s connection guidance distinguishes worker recovery from bounded waiting in HTTP producers. The operational cost is real, but the need depends on the reliability contract, not a hypothetical high-traffic milestone. - Build recovery on the visibility already present. The worker logs receipt, completion, failures, ticket/job IDs and duration. Before using the attempt fields for alerts, I would verify their meaning at receipt and failure events. Add exhausted-attempt alerts, queue age/backlog visibility, stale-result counters and revision correlation. These are useful now; a larger telemetry platform can wait.
- Add a deliberate replay policy. Retained failed jobs are a starting point, not an implemented dead-letter workflow. An operator should distinguish transient failures from invalid input, replay after repair, and know whether replay means the original snapshot or latest conversation. Controlled replay tooling and retention add maintenance work. A separate dead-letter queue is optional until volume or operational ownership warrants it.
I would evaluate these changes through duplicate work, result freshness and publication failures. Infrastructure changes should follow a demonstrated constraint.
What I Would Not Add Yet
I would not introduce Kafka, Kubernetes or decompose the application into microservices solely because it has a queue.
The concrete concerns are durable publication intent, safe repeated writes and valid application of conversation-derived results. They can be addressed within PostgreSQL, Redis, BullMQ and the existing API/worker boundary.
More infrastructure may become appropriate under specific throughput, ownership or deployment constraints. It does not remove the need to define what makes an analysis current or safe to repeat.
Why AI Workloads Make This More Visible
In this worker, three sequential generation operations separate the initial conversation read from the final ticket update. While those external requests are in flight, another message can arrive and another job can complete. The queue can function normally while the result becomes obsolete.
Provider calls can be relatively slow, and successful generation consumes usage billed under the provider’s terms; OpenAI documents token-based API pricing. Repeating a completed summary because sentiment failed can therefore repeat meaningful work, even when no database row has yet changed. This is a reason to control repetition, not a measured cost or latency claim about this application.
Generated text can also differ between attempts. Overwriting aiSummary with a second output is not automatically idempotent just because both writes target the same row. Meanwhile, this processor reloads messages on every attempt, so a retry can change both its input and its output.
For this workload, the useful contract is: which conversation revision does this analysis describe, and when may it take effect? The queue schedules execution. The application must answer that question.
Final Takeaway
BullMQ gives this application a clear boundary between message creation and background AI enrichment. Correctness still needs an application-level design: durable publication intent, duplicate-safe effects and protection against stale conversation results.
Retries make another attempt possible. Idempotency and ordering rules determine whether that attempt is safe to apply.
Originally part of the engineering work behind my AI Support Assistant architecture.