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 like IDLE, FETCHING_CONTEXT, or STREAMING_RESPONSE.
  • runStatus: V1 API-specific lifecycle states including WAITING or RUNNING.
  • 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 threads table in packages/db/src/schema.ts uses Drizzle ORM to define type-safe columns for generation stages, run status, and metadata, supported by strategic indexes.
  • Centralized Writes: All modifications flow through updateThread in packages/db/src/operations/thread.ts, ensuring updatedAt freshness and type safety.
  • Transactional Safety: Service-level functions in apps/api/src/threads/util/thread-state.ts wrap state changes in PostgreSQL transactions with read committed isolation.
  • Concurrency Control: addUserMessage enforces single-writer semantics by checking generationStage before accepting new input, preventing overlapping LLM runs.
  • Lifecycle Management: finishInProgressMessage atomically 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →