How Thread State Is Managed in Tambo AI's Backend Using PostgreSQL and Drizzle ORM
Tambo AI persists conversation thread state in a PostgreSQL table managed by Drizzle ORM, using transactional helpers in thread-state.ts to enforce single-writer semantics and atomic stage transitions.
Tambo AI's backend, built in the tambo-ai/tambo repository, relies on PostgreSQL as the source of truth for conversation state. By combining a carefully designed schema with Drizzle ORM's type-safe query builder, the system tracks thread lifecycle stages—from idle to streaming responses—while preventing race conditions during concurrent access.
Thread Schema Design in PostgreSQL and Drizzle ORM
The threads Table Structure
The foundation of thread state management lives in packages/db/src/schema.ts. Here, the threads table is defined using Drizzle's pgTable helper:
// packages/db/src/schema.ts – definition of the `threads` table
export const threads = pgTable(
"threads",
({ text, timestamp }) => ({
id: text("id") .primaryKey().notNull().unique()
.default(sql`generate_custom_id('thr_')`),
name: text("name"),
projectId: text("project_id") .references(() => projects.id).notNull(),
contextKey: text("context_key"), // end‑user identifier
metadata: customJsonb<Record<string, unknown>>("metadata"),
generationStage: text("generation_stage", {
enum: Object.values<string>(GenerationStage) as [GenerationStage],
})
.default(GenerationStage.IDLE).notNull(),
// V1 API run‑lifecycle fields
runStatus: text("run_status", {
enum: Object.values<string>(V1RunStatus) as [V1RunStatus],
})
.default(V1RunStatus.IDLE).notNull(),
currentRunId: text("current_run_id"),
statusMessage: text("status_message"),
lastRunCancelled: boolean("last_run_cancelled"),
lastRunError: customJsonb<V1RunError>("last_run_error"),
pendingToolCallIds: customJsonb<string[]>("pending_tool_call_ids"),
lastCompletedRunId: text("last_completed_run_id"),
createdAt: timestamp("created_at", { withTimezone: true }).defaultNow().notNull(),
updatedAt: timestamp("updated_at", { withTimezone: true }).defaultNow().notNull(),
sdkVersion: text("sdk_version"),
}),
(table) => [
index("threads_context_key_idx").on(table.contextKey),
index("threads_project_updated_idx").on(table.projectId, table.updatedAt),
index("threads_updated_at_idx").on(table.updatedAt),
index("threads_project_id_idx").on(table.projectId),
index("threads_created_at_idx").on(table.createdAt),
index("threads_sdk_version_idx")
.on(table.sdkVersion).where(sql`${table.sdkVersion} IS NOT NULL`),
],
);
State Tracking Columns
Several columns work together to capture the thread's runtime condition:
generationStage: Tracks high-level LLM generation phases likeIDLE,FETCHING_CONTEXT, orSTREAMING_RESPONSE.runStatus: V1 API-specific lifecycle states includingWAITINGorRUNNING.statusMessage: Human-readable descriptions such as "Fetching weather data…" for UI display.pendingToolCallIds: JSON array tracking outstanding tool executions.lastRunError: Serialized error details when runs fail.updatedAt: Automatically refreshed on every state change for ordering and cache invalidation.
Database Indexes for Performance
The schema defines multiple indexes to support high-throughput queries. These include threads_project_id_idx for project filtering, threads_updated_at_idx for recency sorting, and threads_context_key_idx for user-specific lookups. These indexes optimize the query patterns used by the API controllers when fetching thread lists or checking conversation status.
CRUD Operations with Drizzle ORM
Centralized Update Helper in operations/thread.ts
All writes to the threads table flow through packages/db/src/operations/thread.ts. The updateThread function serves as the single entry point for modifications:
// packages/db/src/operations/thread.ts – update a thread
export async function updateThread(
db: HydraDb,
threadId: string,
{
contextKey,
metadata,
generationStage,
statusMessage,
name,
sdkVersion,
}: {
contextKey?: string | null;
metadata?: ThreadMetadata;
generationStage?: GenerationStage;
statusMessage?: string;
name?: string;
sdkVersion?: string;
},
): Promise<schema.DBThreadWithMessages> {
const [updated] = await db
.update(schema.threads)
.set({
contextKey,
metadata,
updatedAt: sql`now()`,
generationStage,
statusMessage,
name,
...(sdkVersion ? { sdkVersion } : {}),
})
.where(eq(schema.threads.id, threadId))
.returning();
…
}
This centralization guarantees that updatedAt is always refreshed via sqlnow()``. It also enforces type safety—only the columns defined in the TypeScript interface may be changed—and prevents callers from bypassing business logic.
Transactional State Management
Higher-level lifecycle logic resides in apps/api/src/threads/util/thread-state.ts. These functions wrap database operations in transactions to maintain consistency across the threads and messages tables.
Preventing Concurrent Processing in addUserMessage
The addUserMessage function implements a critical guard against overlapping LLM runs:
export async function addUserMessage(
db: HydraDb,
threadId: string,
message: MessageRequest,
logger?: Logger,
sdkVersion?: string,
) {
try {
const result = await db.transaction(
async (tx) => {
const currentThread = await tx.query.threads.findFirst({
where: eq(schema.threads.id, threadId),
});
if (!currentThread) {
throw new Error(`Thread ${threadId} not found`);
}
// Prevent overlapping runs
if (isThreadProcessing(currentThread.generationStage)) {
throw new Error(
`Thread is already in processing (${currentThread.generationStage})`,
);
}
// Move the thread into a “fetching context” stage
await updateGenerationStage(
tx,
threadId,
GenerationStage.FETCHING_CONTEXT,
"Starting processing...",
);
// Insert the new user message
return await addMessage(tx, threadId, message, sdkVersion);
},
{ isolationLevel: "read committed" },
);
return result;
} catch (error) {
logger?.error("Transaction failed: Adding user message", (error as Error).stack);
throw error;
}
}
Before accepting a new user message, the function queries the current generationStage. If isThreadProcessing returns true, it throws an error, enforcing single-writer semantics at the application level. The stage transition to FETCHING_CONTEXT and the message insertion occur within the same transaction, ensuring atomicity.
Finalizing Assistant Responses in finishInProgressMessage
When the LLM stream completes, finishInProgressMessage handles the final state transition:
export async function finishInProgressMessage(
db: HydraDb,
threadId: string,
newestMessageId: string,
inProgressMessageId: string,
finalThreadMessage: ThreadMessage,
logger?: Logger,
): Promise<{
resultingGenerationStage: GenerationStage;
resultingStatusMessage: string;
}> {
try {
const result = await db.transaction(
async (tx) => {
// Verify we are finishing the latest message in the thread
await verifyLatestMessageConsistency(tx, threadId, newestMessageId, true);
// Persist the final message content (including component / tool‑call data)
await updateMessage(tx, inProgressMessageId, {
...finalThreadMessage,
component: finalThreadMessage.component as ComponentDecisionV2Dto,
content: contentPartToDbFormat(finalThreadMessage.content) as ChatCompletionContentPartDto[],
});
// Determine the next generation stage
const resultingGenerationStage = finalThreadMessage.toolCallRequest
? GenerationStage.FETCHING_CONTEXT
: GenerationStage.COMPLETE;
const resultingStatusMessage = finalThreadMessage.toolCallRequest
? `Fetching context...`
: `Complete`;
// Persist the stage on the thread row
await updateGenerationStage(
tx,
threadId,
resultingGenerationStage,
resultingStatusMessage,
);
return { resultingGenerationStage, resultingStatusMessage };
},
{ isolationLevel: "read committed" },
);
return result;
} catch (error) {
logger?.error("Transaction failed: Finishing in-progress message", (error as Error).stack);
throw error;
}
}
This function determines the next stage based on whether tool calls remain pending. If finalThreadMessage.toolCallRequest exists, the stage returns to FETCHING_CONTEXT; otherwise, it becomes COMPLETE. The statusMessage updates simultaneously, providing immediate feedback to polling clients.
Concurrency Guarantees and Consistency
Read Committed Isolation
All transactions in thread-state.ts specify { isolationLevel: "read committed" }. This PostgreSQL isolation level prevents dirty reads while allowing concurrent transactions to proceed, striking a balance between data consistency and system throughput for conversation state updates.
Message Consistency Verification
The verifyLatestMessageConsistency helper checks that operations target the most recent message IDs. This prevents race conditions where two workers might attempt to finish the same in-progress message, ensuring that state transitions always apply to the correct conversation version.
Summary
- Schema Design: The
threadstable inpackages/db/src/schema.tsuses Drizzle ORM to define type-safe columns for generation stages, run status, and metadata, supported by strategic indexes. - Centralized Writes: All modifications flow through
updateThreadinpackages/db/src/operations/thread.ts, ensuringupdatedAtfreshness and type safety. - Transactional Safety: Service-level functions in
apps/api/src/threads/util/thread-state.tswrap state changes in PostgreSQL transactions withread committedisolation. - Concurrency Control:
addUserMessageenforces single-writer semantics by checkinggenerationStagebefore accepting new input, preventing overlapping LLM runs. - Lifecycle Management:
finishInProgressMessageatomically transitions threads from streaming to complete (or back to context fetching) based on pending tool calls.
Frequently Asked Questions
How does Tambo AI prevent two LLM runs from processing the same thread simultaneously?
The backend enforces single-writer semantics through the addUserMessage function in apps/api/src/threads/util/thread-state.ts. Before inserting a new user message, it checks the current generationStage column using isThreadProcessing. If the thread is already in an active stage like FETCHING_CONTEXT or STREAMING_RESPONSE, the function throws an error and aborts the transaction.
What happens to the thread state when a tool call is required during a conversation?
When the LLM response includes a tool call request, the finishInProgressMessage function detects this via the toolCallRequest field in the final message. Instead of setting the stage to COMPLETE, it transitions generationStage back to FETCHING_CONTEXT and updates pendingToolCallIds with the tool call identifiers. This signals to the system that the thread is awaiting external tool results before continuing.
Which PostgreSQL isolation level does Tambo AI use for thread state transactions?
All state transition functions in apps/api/src/threads/util/thread-state.ts explicitly use the read committed isolation level. This PostgreSQL setting prevents dirty reads while allowing concurrent transactions to proceed, providing a balance between data consistency and system throughput for conversation state updates.
How does the schema ensure fast lookups of threads by project or recent activity?
The threads table definition in packages/db/src/schema.ts includes multiple Drizzle index declarations. These include threads_project_id_idx for project filtering, threads_updated_at_idx for recency sorting, and threads_context_key_idx for user-specific lookups. These indexes optimize the query patterns used by the API controllers when fetching thread lists or checking conversation status.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →