Orchestrator
Event-driven conversation processing pipeline (RabbitMQ + three layers)
Orchestrator Contributor Guide#
The orchestrator turns an incoming customer message into a bot reply. It is event-driven: messages are enqueued to RabbitMQ and consumed by a worker that runs a single pipeline — Perception → Retrieval → Execution — with tool calling, human handoff, and confidence guardrails along the way.
This guide covers how the pieces fit together and where to make changes. For code-organization rules (repository vs. entity vs. service responsibilities), read server/orchestrator/ARCHITECTURE.md — that document defines the layering conventions; this one describes the runtime flow.
Note: Older docs described an 8-phase polling architecture with cooldowns, runtime agent routing, and "plan/ender" modules. That design is gone. Agents are assigned at conversation creation, cooldowns no longer gate processing (
Conversation.addMessage()still writescooldown_until, but only the disabled legacy polling path and the unusedfindReadyForProcessing()repository method read it), and processing is triggered by queue messages, not polling.
File Map#
| Path | What lives there |
|---|---|
server/orchestrator/run.ts |
runConversation() — the main pipeline and the execution loop |
server/orchestrator/index.ts |
Orchestrator class: processConversation(), checkInactivity(), legacy loop() |
server/orchestrator/perception.layer.ts |
PerceptionLayer — intent / sentiment / language analysis |
server/orchestrator/retrieval.layer.ts |
RetrievalLayer — playbook selection and vector document search |
server/orchestrator/execution.layer.ts |
ExecutionLayer — planner LLM call, guardrails, prompt assembly |
server/orchestrator/playbook-gate.ts |
shouldReselectPlaybook() — pure predicate gating the per-turn playbook selector |
server/orchestrator/conversation-utils.ts |
Closure validation, title generation, handoff summary, inactivity helpers |
server/orchestrator/types.ts |
ConversationContext, ProcessingPhase, guardrail log types |
server/workers/orchestrator.worker.ts |
OrchestratorWorker — RabbitMQ consumer, retry policy |
server/services/orchestrator-queue.service.ts |
Queue names, enqueue()/enqueueRetry()/enqueueDead(), message shape |
server/services/rabbitmq.service.ts |
amqplib wrapper: reconnect, publish, consume |
server/services/scheduled-jobs.registry.ts |
Sweep, inactivity check, stale-message recovery jobs |
server/services/core/tool-execution.service.ts |
Executes planner tool calls |
server/services/core/action-claim-guardrail.service.ts |
Stage 0 guardrail: action-claim vs. tool-call consistency |
server/services/core/company-interest-guardrail.service.ts |
Stage 1 guardrail: company interest |
server/services/core/confidence-guardrail.service.ts |
Stage 2 guardrail: fact grounding |
server/prompts/en/{perception,retrieval,execution,conversation,playbook}/ |
All LLM prompt templates (also pt/, es/) |
End-to-End Flow#
customer message ──> Conversation.addMessage()
│ sets needs_processing = true
▼
orchestratorQueueService.enqueue(id, orgId, trigger)
│ publishes JSON to RabbitMQ
▼
orchestrator.process queue
│
▼
OrchestratorWorker.startConsuming() (prefetch: 2)
│
▼
Orchestrator.processConversation() ──> runConversation()
│
lock ─> perception ─> retrieval ─> execution loop
│
▼
bot message(s) via conversation.addMessage()
(delivered to channels via message hooks / WebSocket)
Enqueue triggers#
OrchestratorTrigger (server/services/orchestrator-queue.service.ts) enumerates every reason a conversation gets processed:
| Trigger | Fired from |
|---|---|
customer_message |
Conversation.addMessage() for customer messages (server/database/entities/conversation.entity.ts) |
creation |
server/services/conversation.service.ts when a conversation is created |
ai_return |
Conversation entity, when a human hands the conversation back to the AI |
recovery |
server/services/message-recovery.service.ts (stale-message recovery) |
inactivity |
Orchestrator.checkInactivity() after sending an inactivity warning |
sweep |
The orchestrator-sweep scheduled job (safety net, see below) |
The queue message is an OrchestratorMessage: { messageId, conversationId, organizationId, trigger, timestamp, attempt }.
Queue Topology and Retries#
Three durable queues, declared in OrchestratorQueueService.declareQueues():
orchestrator.process— main queue. Unhandled consumer nacks dead-letter toorchestrator.process.dead.orchestrator.process.retry— holds failed messages with amessageTtlof 10 s; expired messages dead-letter back toorchestrator.process. Retries are explicit: on failure the worker republishes withattempt + 1, up toMAX_RETRY_ATTEMPTS = 3.orchestrator.process.dead— terminal. Messages land here after max retries (withfailedAtanderrorattached) or via nack fallback.
Two error messages are treated as non-retryable and acked immediately (isNonRetryableError() in the worker): "Conversation does not need processing" and "Last customer message not found". If you add a new benign early-exit to runConversation(), add its message there too — otherwise it will burn three retries per occurrence.
RabbitMQService reconnects with exponential backoff (max 10 attempts, capped at 30 s), re-declares queues and re-attaches consumers on reconnect, and publishes with persistent: true.
Safety nets#
Registered in server/services/scheduled-jobs.registry.ts:
orchestrator-sweep(every 30 s) — re-enqueues conversations that have hadneeds_processing = truefor more than 30 s, with trigger"sweep". This is whyenqueue()can safely no-op when RabbitMQ is down: the sweep catches up once it reconnects.orchestrator-inactivity-check(every 5 min) — runsOrchestrator.checkInactivity().orchestrator-stale-message-check(every 60 s) — detects and recovers stale/lost messages viamessage-recovery.service.ts.orchestrator-worker-tick— the legacy 1-second polling loop. It still exists (Orchestrator.loop()/orchestratorWorker.tick()) but is registered withenabled: false; the RabbitMQ consumer replaced it. Don't build on it.
runConversation() Step by Step#
The pipeline in server/orchestrator/run.ts is numbered 00–04 in code comments:
- Skip check — conversations with
status === "human-took-over"or anassigned_user_idare skipped entirely. Human takeover works by setting those fields;"ai_return"re-enqueues when handed back. - Lock (00) —
conversation.lock()returns alockerIdor bails (lock contention is not an error, just "someone else is on it"). The lock lasts 15 seconds (LOCK_DURATION_MS) and a heartbeat refreshes it every 5 s (LOCK_HEARTBEAT_INTERVAL_MS) so long LLM/tool calls don't look abandoned. Thefinallyblock does an ownership-checked unlock. - Bootstrap — seeds initial system, agent-instruction, and bot messages if absent. Throws the two non-retryable errors if
needs_processingis false or there is no customer message. - Perception (01) — intent/sentiment/language on the last customer message, saved via
message.savePerception(). If the intent isclose_satisfied/close_unsatisfied(orgreetplus a gratitude regex),validateConversationClosure()re-checks against the full transcript; a confirmed closure generates a title, resolves the conversation, purges Redis secrets (conversationSecretService), sends an LLM-generated closing message, and exits early. - Retrieval (02) — the active-playbook fetch and vector search run in parallel (
Promise.all). Playbook selection is gated byshouldReselectPlaybook()(playbook-gate.ts): the selector LLM runs only when no playbook is set yet or the planner's last verdict was that the current one no longer fits (conversation.metadata.playbookFits !== true) — otherwise the previous selection stands. Selected playbooks are stored withconversation.updatePlaybook(), which populatesenabled_tools. Matched documents are attached viaconversation.addDocument(). - Context init (03) —
orchestration_statusis initialized ({ version: "v1", lastTurn: 0, toolLog: [] }) if missing. - Execution loop (04) —
handleExecutionLoop(), see below. - On success — error counters reset, title generated asynchronously once there are ≥ 2 customer messages,
setProcessed(true). On failure —processing_error_countincrements; at ≥ 3 the conversation is flaggedis_stuckwithstuck_reason: "repeated_processing_failures".
Throughout the run, the processing phase (ProcessingPhase in types.ts: perceiving → retrieving → executing → idle) is persisted in orchestration_status.processingState and broadcast as a conversation_status_changed WebSocket event — via Redis pub/sub channel websocket:events when Redis is up, direct websocketService otherwise (publishStatusChange() in run.ts). The dashboard's "typing" indicators are driven by these events.
The Three Layers#
All LLM calls go through LLMService (server/services/core/llm.service.ts). Instead of hardcoding models, callers pass a task-complexity tier ("easy" | "medium" | "hard", default "hard") which resolves to a concrete model through the organization's provider config (server/services/llm/model-catalog.ts, tier-maps.ts). Env overrides: LLM_TIER_HARD, LLM_TIER_MEDIUM, LLM_TIER_EASY. Never document or assume a single model name.
Prompts are file-based markdown templates resolved by PromptService.getPrompt() from server/prompts/{lang}/ — to change what a layer says to the LLM, edit the prompt file, not the layer code.
Perception (perception.layer.ts)#
perceive(message, conversationId, organizationId)— one structured-output LLM call (tier"medium", promptperception/intent-analysis) returning{ intent: {label, score}, sentiment: {label, score}, language }. Intent/sentiment labels are schema-enum-constrained toMessageIntent/MessageSentiment; language is an ISO 639-1 code enforced by the JSON-schema pattern^[a-z]{2}$.getAgentCandidate()exists (promptperception/agent-selection, score threshold > 0.7) but is not called fromrunConversation()— agents are assigned at conversation creation with a fallback to the organization default. Don't wire new features through it without reconsidering that decision.
Retrieval (retrieval.layer.ts)#
getPlaybookCandidate(messages, playbooks, orgId)— an LLM (tier"medium") scores active playbooks against the conversation (promptretrieval/playbook-selection); accepted only if score > 0.7, otherwisenull. Because of the playbook gate, this ideally runs once per conversation rather than every turn.getRelevantDocuments(messages, orgId)— builds the query from the last 3 customer messages, runsvectorStoreService.search(orgId, query, 5)(top-5 chunks, pgvector), and keeps results with similarity > 0.4. Errors return[]— retrieval failure never fails the run.
Execution (execution.layer.ts)#
execute(conversation, customerLanguage, options) assembles a single planner prompt from:
- the base planner prompt (
execution/planner), including the combined tool name list (playbookenabled_tools+ always-available core tools fromcoreToolRegistryinserver/services/core/core-tools); - a language-enforcement block (only while the conversation has ≤ 3 customer messages);
- a user-context block (
buildContextPrompt()) — merged customerexternal_metadata+ conversationcontext, plus Redis secret key names to be referenced as<<secret.key>>in tool args (values are never exposed to the model); - a knowledge block (
buildKnowledgePrompt()) — the last 5 attached documents, each truncated to 8,000 chars; - optionally, corrective
plannerFeedbackfrom a guardrail-triggered retry (appended last).
options (ExecuteOptions) also carries per-turn state owned by the run loop and shared by reference across re-plans: toolsCalledThisTurn (the tool ledger) and turnGuardrailState (retry budgets).
One structured LLM call against buildPlannerSchema(allToolNames) returns { step, userMessage, toolName, toolArgs, handoffReason, closeReason, rationale, playbookFits } with step ∈ ASK | RESPOND | CALL_TOOL | HANDOFF | CLOSE. playbookFits is the planner's verdict on whether the active playbook still matches the conversation (null when none is active) — it feeds the playbook gate. When tools exist, toolName is enum-constrained in the schema to the allowed set, so the model can't emit a malformed name. The layer then auto-corrects common LLM mistakes: missing step with a tool → CALL_TOOL; RESPOND with a tool but no message → CALL_TOOL; CALL_TOOL with a userMessage → the message is stripped; invalid tool name or ASK/RESPOND without a userMessage → return null, which makes the loop retry the planner.
The Execution Loop (run.ts::handleExecutionLoop)#
Caps: MAX_ITERATIONS = 15, MAX_EMPTY_RETRIES = 3 (after which a canned fallback message is sent). The loop owns the per-turn tool ledger (toolsCalledThisTurn) and the guardrail retry budgets (turnGuardrailState). Each iteration calls executionLayer.execute() and dispatches:
playbookFits === false— the planner says the customer moved on (e.g. order status → refund): the verdict is persisted toconversation.metadata.playbookFits, the selector re-runs, and the plan is discarded and re-made without spending an iteration. Bounded to one re-selection per turn so the planner can't ping-pong between playbooks.CALL_TOOL— aTOOL-type message is written withtoolStatus: "RUNNING", thenToolExecutionService.handleToolExecution()runs the tool and updates the message in place; the outcome is recorded in the tool ledger and the loop continues so the planner can read the result. A guard inrun.tsblocks non-core tools not present inenabled_tools— the blocked call is fed back to the model as aTOOLerror message and does not count against the iteration cap. After a successfulrecommend_productscall,maybeEmitProductRecommendation()emits a separatePRODUCT_RECOMMENDATIONmessage carrying the product payload for the webchat/dashboard cards.HANDOFF— status becomespending-human(broadcast via WebSocket) andgenerateHandoffSummary()is fired asynchronously (promptconversation/handoff-summary, last 15 public messages, saved toconversation.summaryfor the human agent). If online humans exist (userRepository.findOnlineByOrganization), the agent'shuman_handoff_available_instructionsare injected and the loop continues, or a default/LLM-generated transfer message is sent; with nobody online,human_handoff_unavailable_instructionsor a default apology. AhandoffProcessedflag deduplicates repeated HANDOFF steps.CLOSE— optional final message, Redis secrets deleted, loop ends. (The conversation status itself is resolved by the closure path in perception or by inactivity — CLOSE only ends the turn.)RESPOND/ASK— the bot message is sent with guardrail metadata attached. If the result also carries an unexecuted tool, the message is flaggedpotentialHallucinationin metadata and an error is logged.
Guardrails (Three-Stage)#
Applied by ExecutionLayer.applyConfidenceGuardrails() only to RESPOND steps:
- Action-claim consistency (
action-claim-guardrail.service.ts, promptexecution/action-claim-check) — runs first, before the intent exemption, so a false "I've cancelled it, glad I could help!" on a closing turn is still caught. An LLM extracts state-changing action claims from the drafted response and checks them against the turn's tool ledger (failed calls don't back a claim). On violation the planner gets one corrective re-plan with explicit feedback (maxRetries: 1per turn); if the retry still can't back its claim, the result becomes aHANDOFF(or the fallback message whenescalateOnFailureis off). Defaults inDEFAULT_ACTION_CLAIM_GUARDRAIL_CONFIG(server/types/organization-settings.types.ts):{ enabled: true, maxRetries: 1, escalateOnFailure: true }, org-overridable viasettings.actionClaimGuardrail.
Stages 1–2 are skipped for GREET / CLOSE_SATISFIED / CLOSE_UNSATISFIED intents:
- Company interest (
company-interest-guardrail.service.ts) — an LLM assesses the drafted response and returns{ violationType, severity, requiresFactCheck }. Critical severity blocks the response — but first the planner gets one corrective re-plan carrying the reviewer's reasoning (maxRetries: 1per turn, same philosophy as Stage 0); only if that fails does it become aHANDOFFwith a fallback message, the original text preserved in metadata. - Fact grounding (
confidence-guardrail.service.ts) — runs only when Stage 1 setsrequiresFactCheck. Weighted score = grounding 0.6 + retrieval 0.3 + certainty 0.1; default tiershigh ≥ 0.8,medium ≥ 0.5,low < 0.5(org-overridable via organization settings, merged bymergeConfig). Recent successful tool results (last 3 non-errorTOOLmessages) and the active playbook's rendered instructions count as authoritative synthetic documents (similarity 0.95) — without the latter, playbook-driven policy statements would be flagged as ungrounded. Medium tier triggers oneperformRecheck()with relaxed retrieval (threshold 0.3, up to 10 docs), kept only if the score improves. Low tier converts toHANDOFF(ifenableEscalation) or a fallback message.
Fallback messages are no longer only canned: composeContextualFallbackMessage() asks an LLM (prompt execution/contextual-fallback, tier "medium") to acknowledge the customer's actual request from the last few customer-visible messages — deliberately without seeing the blocked response, so unverified claims can't leak through — and falls back to the static configured message on error.
All Stage 1/2 outcomes are appended to orchestration_status.guardrailLog (plus a legacy confidenceLog) by saveConfidenceLog() in run.ts; Stage 0 results ride on message metadata (actionClaim, actionClaimRetryAttempted). The full guardrail design is documented in docs/technical/guardrails.md.
Inactivity Lifecycle#
Orchestrator.checkInactivity() (server/orchestrator/index.ts), helpers in conversation-utils.ts. Thresholds derive from config.conversation.inactivityInterval:
- ½× — LLM-generated warning message (marked
metadata.isInactivityWarning), conversation re-enqueued with trigger"inactivity". - 1× — close with a message.
- 2× — close silently; empty conversations older than 2× are deleted outright.
Closures set status: "resolved", fire the conversation.resolved hook, and trigger title generation (prompt conversation/title-generation, first 10 public messages, 255-char cap).
Extending the Orchestrator#
- Change LLM behavior — edit the prompt file under
server/prompts/en/…(and translations underpt/,es/); keep the JSON schema in the layer in sync if the output shape changes. - New planner step type — extend the
stepenum inbuildPlannerSchema()andExecutionResult(execution.layer.ts), then add a branch inhandleExecutionLoop()(run.ts). - New processing trigger — add it to
OrchestratorTriggerand callorchestratorQueueService.enqueue()from wherever the event originates. Never callrunConversation()directly from request handlers; go through the queue so locking, retries, and the sweep apply. - New always-available tool — register it in the core tool registry (
server/services/core/core-tools); it will be merged into the planner's tool list and exempted from theenabled_toolsgate automatically. - New early-exit condition — if it should not be retried, add its error message to
isNonRetryableError()inorchestrator.worker.ts. - New guardrail stage — follow the Stage 0 pattern: a service with
getDefaultConfig()/mergeConfig(), a per-turn retry budget inturnGuardrailState, and correctiveplannerFeedbackbefore escalating.
Follow the layering rules in server/orchestrator/ARCHITECTURE.md: DB access through repositories, entity state changes through entity methods, business logic in the layers.
Testing and Debugging#
- Orchestrator tests live in
server/tests/orchestrator/(Jest):cd server && npm test. Good entry points:playbook-gating.test.ts,execution-action-claim.test.ts,execution-company-interest-retry.test.ts,execution-confidence-context.test.ts,execution-knowledge.test.ts. - Set
LOG_LEVEL=debugand filter withDEBUG_MODULES="perception,retrieval,execution"— the layers log under those module names; the pipeline logs underorchestrator-run, the consumer underorchestrator-worker, the queue underorchestrator-queue, Stage 0 underaction-claim-guardrail. - Every log line is a child logger carrying
organizationIdandconversationId, so you can trace a single conversation across layers. - Stuck conversations are visible in the DB:
is_stuck,stuck_reason,processing_error_count,last_processing_error. Terminally failed queue messages sit inorchestrator.process.deadwith the error attached.