Djaouad

Building a Shopping Assistant AI Agent from Scratch: MCP, RAG and Conversations in a DDD TypeScript Backend

October 6, 202665 min read
typescriptai-agentsllmragmcpsemantic-searchpgvectorgeminidddonion-architecturebackendnodejs

An agent loop, an MCP tool server, an event driven RAG pipeline and persisted conversations, all built as ports and adapters on top of the same DDD + Onion Architecture codebase.

This post builds on Dependency Injection from Scratch, the container, scopes and composition roots described there are used here without re-explaining them. also, the full application is in this repository


Part I, Agent Fundamentals

1. What is an AI Agent

  • an AI agent is a system that has the ability to take autonomous decisions to achieve certain task. in coding terms, instead of hard coding the control flow of a system, you let an LLM decide it at request time, this is useful when the task can not be achieved by a predictable workflow with exact sequence of steps.
  • a large language model is just a function that takes text as an input and returns text as an output.
  • an AI agent is essentially a loop that keeps feeding the LLM a larger context(text), so it gets more informed to take the right action or generate the right response. for each LLM call, it could return a final response, or ask for a tool to be executed.
  • the context we feed to the LLM is the combination of:
    • system prompt + user prompt + conversation history + tools declarations + any tool call request by the LLM + the tool call result, the context initially starts with only system & user prompts + tools declarations and then conversation history + tool call requests and responses are incrementally added as the agent runs.

Here is the whole loop, as it runs in this codebase:

export type AssistantAgentConfig = {
  maxSteps: number;
};

export class AssistantAgent {
  constructor(
    private mcpClientGateway: McpClientGateway,
    private chatModel: ChatModelPort,
    private systemPrompt: string,
    private config: AssistantAgentConfig,
  ) {}

  async run(
    messages: ChatMessage[],
  ): Promise<{ response: string; newMessages: ChatMessage[] }> {
    const maxSteps = this.config.maxSteps ?? 8;
    const startIndex = messages.length;

    for (let step = 0; step < maxSteps; step++) {
      const response = await this.chatModel.generate({
        systemInstruction: this.systemPrompt,
        messages,
        tools: await this.mcpClientGateway.listTools(),
      });

      messages.push(response.message);

      if (response.functionCalls.length === 0) {
        return {
          response: response.text,
          newMessages: messages.slice(startIndex),
        };
      }

      for (const call of response.functionCalls) {
        const callResult = await this.mcpClientGateway.safeToolCall(
          call.name,
          call.args,
        );

        messages.push({
          role: "user",
          parts: [
            {
              functionResponse: {
                name: call.name,
                response: { result: callResult },
              },
            },
          ],
        });
      }
    }

    throw new MaxStepsExceededError(
      `Assistant agent exceeded max steps: ${maxSteps}`,
    );
  }
}
Rendering diagram…

A few things worth noticing in the loop

  • the loop never knows which LLM provider or which tool server it talks to, it only holds two contracts (McpClientGateway and ChatModelPort), a system prompt string and a config object.
  • maxSteps is the safety net against an LLM that keeps asking for tools forever, every iteration costs tokens and latency, so the loop is bounded and fails loudly with a MaxStepsExceededError.
  • the messages array is the context, it only ever grows, each model message and each tool result is pushed to it, so the next generate() call sees everything that happened before.
  • startIndex remembers where the new messages begin, so the caller gets back only the messages produced by this run (newMessages), which is exactly what we need later to persist the conversation.
  • the tool declarations are requested from the gateway on every step, it's cheap because the gateway serves them from its cache.

2. How to Build an AI Agent

  • the core agent functionality depends on a couple of things:
    1. a gateway that provides tools declarations + the tool's executing function.
    2. a chat model that takes the context as an input and generates a response (text or tool call request)
    3. a system prompt

Notice how the first two are exactly the shape of a port from the onion architecture. The agent class lives in the application layer, it depends on contracts owned by the application layer, and the infrastructure layer supplies the adapters (a Gemini client, an MCP client), the same DIP story as the repositories in the DI post.

Rendering diagram…

Swapping Gemini for another provider, or MCP for locally defined tools, means writing one new adapter and changing one registration in the composition root, the agent loop doesn't change.


Part II, Constructing the Building Blocks first

3. The Tools Provider Gateway

  • a tools provider gateway is an abstraction over an MCP client or locally defined tools, its main purpose is to abstract the logic behind tools loading, listing and execution.
  • a tools provider gateway as an abstraction over an MCP client, could expose 3 methods:
    1. loadTools(): called at server startup, it uses the MCP client to fetch the tools declarations from one or more MCP servers, and possibly cache them internally for future agent calls.
    2. listTools(): called inside the agent loop to provide available tools declarations (names, descriptions, arguments, outputs,... etc.)
    3. safeToolCall(): given a tool name and arguments, it safely executes that tool, either returning tool call result or any thrown error inside the call (doesn't allow errors to propagate to the agent loop)
export type ToolDeclaration = {
  name: string;
  description: string;
  parameters: Record<string, unknown>;
};

export type McpToolsMap = Map<string, ToolDeclaration>;

export type McpClientGateway = {
  /** load and store mcp tools declarations in adapter cache, called at bootstrap */
  loadTools(): Promise<void>;

  /** return cached tools declarations, if none exist, load them first */
  listTools(): Promise<ToolDeclaration[]>;

  safeToolCall(name: string, args: unknown): Promise<unknown>;

  close(): Promise<void>;
};

The contract also has a close() method, the entrypoints that own the gateway (the local REPL for example) use it to release the MCP connection on shutdown.

The adapter is the only place in the codebase that knows the MCP SDK exists:

export class McpClientGatwayAdapter implements McpClientGateway {
  private toolsMap: Map<string, ToolDeclaration> | null = null;

  constructor(
    private client: Client,
    private config: McpClientGatewayConfig,
  ) {}

  async loadTools(): Promise<void> {
    const transport = new StreamableHTTPClientTransport(
      new URL(this.config.SERVER_URL),
      {
        requestInit: {
          headers: {
            authorization: `Bearer ${this.config.API_KEY}`,
          },
        },
      },
    );

    try {
      await this.client.connect(transport as Transport);

      const { tools } = await this.client.listTools();

      this.toolsMap = new Map();

      tools.map((tool) => {
        this.toolsMap!.set(tool.name, {
          name: tool.name,
          description: tool.description ?? "",
          parameters: tool.inputSchema ?? {},
        });
      });
    } catch (error) {
      throw new GatewayError("MCP", error);
    }
  }

  async listTools(): Promise<ToolDeclaration[]> {
    if (!this.toolsMap) {
      await this.loadTools();
    }

    return [...(this.toolsMap ? this.toolsMap.values() : [])];
  }

  async safeToolCall(name: string, args: unknown): Promise<unknown> {
    try {
      // safeToolCall shouldn't throw when called in the middle of the agent loop
      if (!this.toolsMap) {
        await this.loadTools();
      }

      const tool = this.toolsMap?.get(name);

      if (!tool) {
        throw Error(`Tool ${name} not found`);
      }

      if (!this.isObject(args)) {
        throw new Error(
          `Invalid arguments for tool ${name}: expected an object`,
        );
      }

      const result = await this.client.callTool({
        name: tool.name,
        arguments: args,
      });

      if ("toolResult" in result) {
        return result.toolResult;
      }

      return {
        content: result.content,
        ...(result.structuredContent !== undefined && {
          structuredContent: result.structuredContent,
        }),
        ...(result.isError !== undefined && {
          isError: result.isError,
        }),
      };
    } catch (error) {
      return {
        content:
          error instanceof Error
            ? error.message
            : "tool call failed, unknown error",
        isError: true,
      };
    }
  }

  async close(): Promise<void> {
    await this.client.close();
  }

  // ...
}
  • loadTools() opens a Streamable HTTP connection to the MCP server (authenticated with the API key as a bearer token), lists the tools and normalizes each one into the port's ToolDeclaration shape, so the rest of the application never sees an MCP SDK type.
  • listTools() serves the cached map, and lazily loads it when nobody called loadTools() yet.
  • safeToolCall() is where "safe" is earned: an unknown tool name, malformed arguments or a transport failure are all converted into a normal result with isError: true, so the LLM reads the failure as a tool result and can react to it (retry with different arguments, or apologize to the user), instead of the whole agent run crashing.

4. The Chat Model Port

  • a chat model port is an abstraction over an actual LLM API call (using fetch or an SDK). it defines the contracts that govern how the agent loop interacts with the LLM by exposing a generate() method, it:
    • takes as parameters:
      1. the system prompts
      2. the tools declarations
      3. a "chat messages" array.
      4. an optional "temperature" number that configures the agent's creativity.
    • it returns:
      1. a "chat message" we add to the messages history array.
      2. a "text" field if the LLM decided a final response.
      3. a tool calls array that defines tool names and the call arguments.
  • a chat message is a combination of an "actor" + "message type" + "optional provider state"
    • we have two actors, either a "user" or a "model"
    • we have 3 message types:
      1. a text message: either the initial user prompt or the final LLM response.
      2. a function call message: contains a function name and arguments, returned by the LLM throughout the loop before generating the final message.
      3. a function response: contains a function name and its returned call result.
    • the optional provider state: some LLM providers like Google's Gemini, return extra fields in their response and they require you to send them back (like thought signatures), these fields are not accessed by our agent loop "it uses the 3 chat model port defined message shapes after normalizing the raw LLM response". so we store the raw LLM response in the chat message as a "provider state" field that holds extra information irrelevant to our agent loop and that will always be accessed only by the adapter implementing the chat model port.
export type ChatTextPart = {
  text: string;
};

export type ChatFunctionCallPart = {
  functionCall: {
    name: string;
    args: Record<string, unknown>;
  };
};

export type ChatFunctionResponsePart = {
  functionResponse: {
    name: string;
    response: { result: unknown };
  };
};

export type ChatPart =
  | ChatTextPart
  | ChatFunctionCallPart
  | ChatFunctionResponsePart;

export type ChatMessage = {
  role: "user" | "model";
  parts: ChatPart[];
  /**
   * Opaque, adapter-owned continuation state. Some providers return reasoning
   * state (signatures, encrypted reasoning items) that must be sent back
   * unmodified on the next turn. The application never inspects this; it only
   * keeps it attached to the message. Must be JSON-serializable.
   */
  providerState?: unknown;
};

export type GenerateParams = {
  systemInstruction: string;
  messages: ChatMessage[];
  tools?: ToolDeclaration[];
  temperature?: number;
};

export type GenerateResult = {
  message: ChatMessage;
  text: string;
  functionCalls: { name: string; args: Record<string, unknown> }[];
};

export type ChatModelPort = {
  generate(params: GenerateParams): Promise<GenerateResult>;
};

The Gemini adapter is where the provider state idea becomes concrete, it normalizes on the way out and de-normalizes on the way in:

export class GemeniChatModelAdapter implements ChatModelPort {
  constructor(
    private chatClient: GoogleGenAI,
    private config: GemeniChatModelAdapterConfig,
  ) {}

  async generate(params: GenerateParams): Promise<GenerateResult> {
    const { messages, systemInstruction, temperature, tools } = params;
    try {
      const response = await this.chatClient.models.generateContent({
        model: this.config.GEMENI_CHAT_MODEL,
        contents: messages.map((m) => this.toContent(m)), // if a message was generated by the model, we will send the raw gemeni parts with their thought signatures, if by user, we will send the normalized ChatPart[]
        config: {
          temperature: temperature !== undefined ? temperature : 0.2,
          systemInstruction,
          ...(tools &&
            tools.length > 0 && {
              tools: [
                {
                  functionDeclarations: tools.map((t) => ({
                    name: t.name,
                    description: t.description,
                    parametersJsonSchema: t.parameters,
                  })),
                },
              ],
            }),
        },
      });

      return this.normalizeResponse(response.candidates);
    } catch (error) {
      handleGeminiClientErrors(error, "GemeniChatModelAdapter.generate");
    }
  }

  private normalizeResponse(
    candidates: Candidate[] | undefined,
  ): GenerateResult {
    const candidate = this.requireCandidate(candidates);
    const content = this.requireContent(candidate);
    const parts = content.parts ?? []; // contains text, functionCalls, gemeni's thought signature per part,... etc.

    const normalizedParts: ChatPart[] = parts.flatMap((part) =>
      this.normalizePart(part),
    );

    const message: ChatMessage = {
      role: "model",
      parts: normalizedParts,
      providerState: parts, // save raw gemeni parts with thought signatures, our application doesn't care about them so it will never access the providerState field, we will send the raw parts however to the generate() method of the adapter if the ChatMessage was generated by the model. since they contain all the previous gemeni messages with their thought signatures(which is required by the model)
    };

    return {
      message,
      text: normalizedParts
        .filter(
          (part): part is Extract<ChatPart, { text: string }> => "text" in part,
        )
        .map((part) => part.text)
        .join(""),

      functionCalls: normalizedParts.flatMap((part) =>
        "functionCall" in part ? [part.functionCall] : [],
      ),
    };
  }

  private toContent(message: ChatMessage): Content {
    if (Array.isArray(message.providerState)) {
      // message.providerState is only defined when the message was generated by the model
      // it contains the raw gemeni parts with their thought signatures, text, functionCalls, ...etc.
      return { role: message.role, parts: message.providerState as Part[] };
    }
    return { role: message.role, parts: message.parts as Part[] };
  }

  // ...
}
  • normalizeResponse() turns the raw Gemini candidate into the port's GenerateResult: the agent loop gets a flat text, a list of functionCalls and a message it can push to the history.
  • the same raw parts are stored untouched in providerState, and toContent() sends them back on the next call for any message generated by the model, this is how the thought signatures survive without a single line of the agent loop knowing they exist.
  • the adapter owns the provider specific failure translation too (handleGeminiClientErrors), anything that goes wrong inside the SDK call comes out as one of our own error types, keep that in mind for the context window section later.

5. The System Prompt

  • A system prompt is a set of high-priority instructions that defines how an AI should behave, what rules it must follow, and what constraints it operates under. Think of it as the AI's operating rules, which guide how it interprets and responds to user requests. it could affect the order in which the LLM use tools and describe some workflows for specific scenarios, and how to respond in edge cases.
export function buildAssistantAgentSystemPrompt(deps: {
  storeName: string;
}): string {
  return `You are ${deps.storeName}'s AI shopping assistant.

YOUR TOOLS
- product-semantic-search: semantic product discovery with optional filters
  (colors, sizes, minPrice, maxPrice, inStock). ALWAYS use this first when the
  user asks about products, wants recommendations, or describes what they're
  looking for.
- get-product-full-details: full details of ONE product by productId. Use
  only to expand a candidate from search — never for every search result.

HARD RULES
1. Extract structured constraints from the user's message into the search
   filters: color words, sizes, price bounds ("under 8000"), availability
   ("in stock"). Put only free-text intent into the query field.
2. ALL product facts (names, prices, variations stock and available qty,
   colors, sizes, weight and descriptions) come ONLY from tool
   results. NEVER invent them. If search returns nothing, say so.
3. Quote prices EXACTLY as returned, in DZD. Price filters use the effective
   (discounted) price.
4. Answer concisely (under 120 words unless asked for detail). Mention product
   name and price when recommending.
5. If the user asks about something other than shopping in this store, politely redirect.

You are speaking with customers in Algeria.`;
}

The prompt is a function, not a constant: the store name is injected from the environment by the composition root, so the same agent code can serve any store. The prompt also encodes a workflow (search first, expand only one candidate) and the grounding rule that makes RAG trustworthy: product facts come only from tool results.


Part III, RAG

  • retrieval augmented generation is the process of retrieving external information at runtime and providing it to the model as a context to get answers tailored and grounded to that knowledge. the LLM doesn't know your products catalog or your order states,... etc. instead of training the model on your data (which could be quite expensive and inconvenient for frequently changing products for example) you can give it tools to access fresh up to date data.
  • the tools to access the data vary, for example a "find product by ID" tool could retrieve the full product information if the agent has the ID. but this tool will not be very useful if the agent doesn't know what products even match the user's intent. a user searching for a specific product in a store but they don't have the exact name of the product, don't know the category or the material, has a blurry description with a price range or they want to compare two products and the store contains thousands of products even after using the traditional filters. in this case filtering rows in the database to find words matching the user's query will not be very useful, we could be selling leather boots and a user searches for "winter shoes", this will not match any tokens in the database and the client will think you don't sell what they're looking for, a loss for the business. the solution here is a hybrid "semantic search" with range filters tool that finds the top K products matching the users intent both "semantically" and within price/stock ranges,... etc.
  • this pipeline is essentially the workflow of how we store our data semantically, and how we retrieve it later by "meaning matching". this pipeline can be split into two workflows:

6.1 Offline RAG

  • offline RAG is the process of keeping an up to date knowledge base made of "indexed products", this process can be done in an event driven way by:
    1. storing the product created/updated/deleted events in the DB.
    2. a processor worker reads and publishes those events to an "embeddings queue".
    3. a handler worker subscribed to the embedding queue listens to events and does:
      1. in case of a "product deleted" event:
        1. delete all the related chunks of the target product from DB so any future semantic search call will never match that product.
      2. in case of a "product created/updated" event:
        1. read the new product details from DB.
        2. chunk them semantically to preserve valid meaning.
        3. embed all the chunks using a text embedding model
        4. upsert all the new chunks to DB
Rendering diagram…

Step 1 and 2 reuse the outbox machinery that already exists for the rest of the system, which means the embeddings queue gets the same delivery guarantees as the other queues, the RAG pipeline is just another consumer.

The knowledge base lives in a Postgres table with a vector column (pgvector), one row per chunk:

export const productEmbeddings = pgTable(
  "product_embeddings",
  {
    id: varchar("id", { length: 40 }).notNull().primaryKey(),
    product_id: varchar("product_id", { length: 40 }).notNull(),
    chunk_index: smallint("chunk_index").notNull(),
    content: text("content").notNull(),
    embedding: vector("embedding", { dimensions: 768 }).notNull(),
    created_at: timestamp("created_at", { withTimezone: true })
      .notNull()
      .defaultNow(),
    updated_at: timestamp("updated_at", { withTimezone: true })
      .notNull()
      .$onUpdate(() => new Date())
      .defaultNow(),
  },
  (t) => [
    index("product_embeddings_hnsw_idx").using(
      "hnsw",
      t.embedding.op("vector_cosine_ops"),
    ),
    uniqueIndex("product_embeddings_product_chunk_idx").on(
      t.product_id,
      t.chunk_index,
    ),
  ],
);
  • the hnsw index with vector_cosine_ops is what makes nearest neighbour search by cosine distance fast, instead of scanning every vector.
  • the unique index on (product_id, chunk_index) guarantees a product can't have two chunks claiming the same position.
  • the dimensions: 768 must match the dimensions requested from the embedding model, otherwise inserts fail.

The two ports the pipeline depends on, the embedding model and the repository that stores the chunks:

export type TextEmbeddingModelPort = {
  // batch embedding by default
  embed(text: string[]): Promise<number[][]>;
};
export type EmbeddingUpsertParams = {
  batch: {
    embedding: number[];
    content: string;
    chunkIndex: number;
  }[];
  productId: string;
};

export type ProductEmbeddingRepository = {
  upsert(params: EmbeddingUpsertParams, tx: TransactionClient): Promise<void>;
  deleteByProductId(productId: string, tx: TransactionClient): Promise<void>;
};

embed() is batch by default, one network call embeds all the chunks of a product, and the returned vectors are in the same order as the provided texts. Here are the Gemini embedding adapter and the Postgres repository implementing those ports:

export class GemeniTextEmbeddingModelAdapter implements TextEmbeddingModelPort {
  constructor(
    private gemeniClient: GoogleGenAI,
    private config: GemeniTextEmbeddingModelAdapterConfig,
  ) {}

  async embed(text: string[]): Promise<number[][]> {
    try {
      if (text.length >= 100)
        throw new ValidationError(
          "text",
          "text array length exceeds the gemeni limit for synchronous batched embeddings",
        );

      const response = await this.gemeniClient.models.embedContent({
        contents: text,
        model: this.config.GEMINI_EMBEDDING_MODEL,
        config: {
          outputDimensionality: this.config.EMBEDDING_DIMENSIONS,
        },
      });

      if (!response.embeddings)
        throw new GatewayError("Gemeni", new Error("No embeddings returned"));

      const embeddings = response.embeddings.map((embedding, index) => {
        if (!embedding.values)
          throw new GatewayError(
            "Gemeni",
            new Error(
              `no embedding values found for ${text[index] ?? `text[${index}]`}`,
            ),
          );

        return embedding.values;
      });

      // in the same order as provided in the batch request.
      return embeddings;
    } catch (error) {
      handleGeminiClientErrors(error, "GemeniTextEmbeddingModelAdapter.embed");
    }
  }
}
export class PostgresProductEmbeddingRepository implements ProductEmbeddingRepository {
  async upsert(
    params: EmbeddingUpsertParams,
    tx: TransactionClient,
  ): Promise<void> {
    const db = tx as DrizzleTransactionClient;

    try {
      // next queries are atomic in postgres, they use the same tx client

      await this.deleteByProductId(params.productId, tx);

      await db.insert(productEmbeddings).values(
        params.batch.map((b) => ({
          id: generateProductEmbeddingId(),
          product_id: params.productId,
          embedding: b.embedding,
          content: b.content,
          chunk_index: b.chunkIndex,
        })),
      );
    } catch (error) {
      handleDrizzleErrors(error, "PostgresProductEmbeddingRepository.upsert");
    }
  }

  async deleteByProductId(
    productId: string,
    tx: TransactionClient,
  ): Promise<void> {
    try {
      const db = tx as DrizzleTransactionClient;

      await db
        .delete(productEmbeddings)
        .where(eq(productEmbeddings.product_id, productId));
    } catch (error) {
      handleDrizzleErrors(
        error,
        "PostgresProductEmbeddingRepository.deleteByProductId",
      );
    }
  }
}
  • the "upsert" here is delete then insert inside the caller's transaction: a product whose description got shorter might now have fewer chunks than before, so updating rows in place would leave stale chunks behind. replacing the whole set is simpler and always correct.
  • both methods take the transaction client from the caller, which is what lets the handler services below make the idempotency key insert and the embedding write one atomic unit.

Now the two handler services the worker calls. The deleted event is the simplest one:

export class EmbeddingQueueProductDeletedEventHandlerService {
  constructor(
    private db: DBClient,
    private productEmbeddingRepository: ProductEmbeddingRepository,
    private idempotencyKeysRepository: IdempotencyKeysRepository,
  ) {}

  async execute(
    command: EmbeddingQueueProductDeletedEventHandlerCommand,
    jobId: string,
  ) {
    await this.db.transaction(async (tx) => {
      await this.idempotencyKeysRepository.create(
        jobId,
        "EmbeddingQueueProductDeletedEventHandlerService",
        tx,
      );

      await this.productEmbeddingRepository.deleteByProductId(
        command.productId,
        tx,
      );
    });
  }
}

The created/updated event does the real work:

export class EmbeddingQueueProductUpsertedEventsHandlerService {
  constructor(
    private db: DBClient,
    private embeddingModel: TextEmbeddingModelPort,
    private productEmbeddingRepository: ProductEmbeddingRepository,
    private productQueries: ProductQueries,
    private idempotencyKeysRepository: IdempotencyKeysRepository,
  ) {}

  async execute(
    command: EmbeddingQueueProductUpsertedEventsHandlerCommand,
    jobId: string,
  ) {
    const productDto = await this.productQueries.getStaticData(
      ProductId.of(command.productId),
    );

    if (!productDto) throw new NotFoundError("product", command.productId);

    const chunks = chunkProduct({
      name: productDto.name,
      material: productDto.material,
      categoryName: productDto.category?.name ?? null,
      brand: productDto.brand,
      description: productDto.description,
    });

    const embeddings = await this.embeddingModel.embed(
      chunks.map((c) => c.content),
    );

    // shape the batch for the embedding model's upsert() method
    const batch = chunks.map((chunk, index) => {
      const embedding = embeddings[index];
      // throw if no embedding generated for a specific chunk
      if (!embedding)
        throw new GatewayError("Gemeni", new Error("No embeddings returned"));

      return {
        chunkIndex: chunk.index,
        content: chunk.content,
        embedding,
      };
    });

    await this.db.transaction(async (tx) => {
      await this.idempotencyKeysRepository.create(
        jobId,
        "EmbeddingQueueProductUpsertedEventsHandlerService",
        tx,
      );

      await this.productEmbeddingRepository.upsert(
        {
          productId: productDto.id,
          batch,
        },
        tx,
      );
    });
  }
}
  • the slow and fallible work (reading the product, chunking, calling the embedding API) happens before the transaction opens, the transaction itself is short: an idempotency key insert + the embedding write.
  • queues deliver at least once, so a job can run twice. the idempotency key is keyed by the job id, a second run fails on the key insert and rolls back, which makes the handlers safe to retry without duplicating or corrupting chunks.
  • the handler asks for the product's data through ProductQueries and not through a repository, it only needs a read model to build the text, not an aggregate.

6.2 Online RAG

  • online RAG is the process of matching a user's search query by meaning instead of token overlap, it works by converting the user's text query into an embedding (a vector representing the query's meaning), this process is done by a specialized "Text Embedding Model" (must be the same one we used to embed the product chunks), we then use that query embedding to look for the semantically closest top K chunk vectors in the database using cosine similarity, and with the vector search we could also add optional hardcoded filters covering information we didn't embed earlier, things like: available colors, sizes, stock quantity, prices range, availability... etc.
export class ProductSemanticSearchService {
  constructor(
    private productQueries: ProductQueries,
    private embeddingModel: TextEmbeddingModelPort,
  ) {}

  async execute(query: ProductSemanticSearchQuery) {
    const [embedding] = await this.embeddingModel.embed([query.query]);

    if (!embedding)
      throw new GatewayError(
        "embeddingModel",
        new Error("No embedding returned"),
      );

    return this.productQueries.semanticSearch({
      queryVector: embedding,
      limit: query.limit ?? 5,
      filters: query.filters ?? {},
    });
  }
}

The read side contract carries the filters in domain language, nothing about SQL or pgvector:

export interface SemanticSearchFilters {
  minPrice?: number;
  maxPrice?: number;
  inStock?: boolean; // at least one variation with available qty > 0
  colors?: Color[]; // at least one variation matching ANY of these
  sizes?: Size[]; // at least one variation matching ANY of these
}

export type ProductQueries = {
  // ...

  semanticSearch(params: {
    queryVector: number[];
    limit: number;
    filters: SemanticSearchFilters;
  }): Promise<SemanticProductHit[]>;
};

And the Postgres implementation of the hybrid search, vector distance plus hard filters in one query:

async semanticSearch(params: {
  queryVector: number[];
  limit: number;
  filters: SemanticSearchFilters;
}): Promise<SemanticProductHit[]> {
  const { queryVector, limit, filters } = params;

  try {
    const conditions: SQL[] = [];

    if (filters.minPrice != null || filters.maxPrice != null) {
      const displayPrice = sql`coalesce(${product.discount_price}, ${product.price})`;

      if (filters.minPrice != null)
        conditions.push(sql`${displayPrice} >= ${filters.minPrice}`);

      if (filters.maxPrice != null)
        conditions.push(sql`${displayPrice} <= ${filters.maxPrice}`);
    }

    if (filters.colors && filters.colors.length > 0) {
      conditions.push(sql`exists (select 1 from ${variation}
    where ${variation.product_id} = ${product.id}
      and ${inArray(variation.color, filters.colors)})`);
    }

    if (filters.sizes && filters.sizes.length > 0) {
      conditions.push(sql`exists (select 1 from ${variation}
    where ${variation.product_id} = ${product.id}
      and ${inArray(variation.size, filters.sizes)})`);
    }

    if (filters.inStock) {
      conditions.push(sql`exists (select 1 from ${variation}
    where ${variation.product_id} = ${product.id}
      and ${variation.total_qty} - ${variation.reserved_qty} > 0)`);
    }

    const distance = cosineDistance(productEmbeddings.embedding, queryVector);

    const rows = await this.db
      .select({
        productId: product.id,
        name: product.name,
        slug: product.slug,
        price: product.price,
        discountedPrice: product.discount_price,
        similarityDistance: sql<number>`min(${distance})`,
      })
      .from(productEmbeddings)
      .innerJoin(product, eq(productEmbeddings.product_id, product.id))
      .where(and(...conditions))
      .groupBy(product.id)
      .orderBy(asc(sql`min(${distance})`))
      .limit(limit);

    return rows.map((r) => ({
      ...r,
      currency: Currency.DZD,
      similarityDistance: Number(r.similarityDistance.toFixed(4)),
    }));
  } catch (error) {
    handleDrizzleErrors(error, "PostgresProductQueries.semanticSearch");
  }
}
Rendering diagram…
  • a product has many chunks, so the query takes the min distance per product and groups by product: the best matching chunk decides the product's rank, and the agent never sees the same product twice.
  • the structured constraints (price, colors, sizes, stock) are plain SQL conditions, the semantic part only has to answer "what does the user mean", the filters answer "what is actually allowed". this is what makes it hybrid.
  • the price filters use the effective price, coalesce(discount_price, price), which is the same price the customer pays, and the tool description tells the LLM exactly that.
  • inStock is computed from the variations (total_qty - reserved_qty), stock information is deliberately not embedded, because it changes constantly and a vector can't be cheaply kept in sync with it, a filter always reads fresh data.

Part IV, MCP

  • MCP (Model Context Protocol) is a standardized protocol that lets an AI model/agent discover and use external tools and data through a consistent interface.
  • building an MCP server that exposes our application-specific tools will make them a "build once, use by any agent", and this MCP server could run over the network for remote communication via StreamableHTTPServerTransport, or run locally via stdio for things like local agents.
  • we can protect our MCP endpoint by an API key middleware, and build a fresh scope per request to avoid leaking different tool executions to different MCP calls.
  • for each tool we register on the MCP server, we store its name, description, input/output schemas + a handler function executed on tool calls, it's important to expose all of these details to give the agent maximum clarity on how a specific tool behaves.

The server is just a list of tool registrations, built from the request's scope:

export function createMcpServer(scope: Scope): McpServer {
  const server = new McpServer({
    name: "ecommerce-assistant",
    version: "1.0.0",
  });

  productSemanticSearchToolRegistration(scope, server);

  getProductFullDetailsToolRegistration(scope, server);

  return server;
}

Each tool is a thin shell around an application service. It validates through its schemas, resolves the service from the scope, builds the query object, and turns any exception into an isError result:

export function productSemanticSearchToolRegistration(
  scope: Scope,
  server: McpServer,
) {
  server.registerTool(
    "product-semantic-search",
    {
      description: toolDescription,
      inputSchema: productSemanticSearchInputSchema,
      outputSchema: productSemanticSearchOutputSchema,
    },
    async ({ query, limit, filters }) => {
      try {
        const service = scope.resolve(PRODUCT_SEMANTIC_SEARCH_SERVICE);
        const semanticSearchQuery = new ProductSemanticSearchQuery(
          query,
          limit,
          filters,
        );

        const result = await service.execute(semanticSearchQuery);

        return {
          content: [
            {
              type: "text",
              text: JSON.stringify(result),
            },
          ],
          structuredContent: { products: result },
        };
      } catch (error) {
        return {
          isError: true,
          content: [
            {
              type: "text",
              text:
                error instanceof Error
                  ? error.message
                  : "Unknown semantic search error",
            },
          ],
          structuredContent: { products: [] },
        };
      }
    },
  );
}

const toolDescription = `
Find products matching a user's request using natural-language semantic
search, with optional filters for color, size, price, and stock.
Use this tool FIRST to discover and shortlist relevant products.
Returns lightweight product records, including productId, name, price,
and similarityDistance.
When you need more information about a candidate, call
get-product-static-details with its productId. Do not fetch full details
for every result unnecessarily.
Results are ordered by ascending cosine distance (smaller is a closer
semantic match). Semantic relevance does not guarantee that a product
satisfies every requested attribute; verify specific requirements
using product details. Price filters use the effective price
(discounted price when available, otherwise regular price), in DZD.
`;

The second tool follows the same shape, with a description that tells the agent when to reach for it:

export function getProductFullDetailsToolRegistration(
  scope: Scope,
  server: McpServer,
) {
  server.registerTool(
    "get-product-full-details",
    {
      description: toolDescription,
      inputSchema: getProductFullDetailsInputSchema,
      outputSchema: getProductFullDetailsOutputSchema,
    },
    async ({ productId }) => {
      try {
        const service = scope.resolve(GET_PRODUCT_FULL_DETAILS_SERVICE);

        const result = await service.execute(
          new GetProductStaticDataQuery(productId),
        );

        return {
          content: [{ type: "text", text: JSON.stringify(result) }],
          structuredContent: result,
        };
      } catch (error) {
        return {
          isError: true,
          content: [
            {
              type: "text",
              text:
                error instanceof Error
                  ? error.message
                  : "Unknown product full details search error",
            },
          ],
          structuredContent: {},
        };
      }
    },
  );
}

const toolDescription = `
Retrieve the full static details of a specific product with all of it's variation details
using its productId.

Use this after product-semantic-search identifies a relevant product,
when additional information is needed to answer the user's question, such as questions
about the product's variations availability, colors, sizes, weight, ...etc.

Returns the product's description, brand, material, pricing, category,
rating, images, variations data and timestamps. The product must exist.
`;
  • the tool descriptions are part of the prompt the LLM sees, they are not documentation for humans, so they say when to use the tool, what comes back, and what not to do (don't expand every search result).
  • a tool never throws, errors come back as isError: true with the message, and the agent loop sees a regular tool result. that's the server side twin of the client's safeToolCall().
  • tools resolve services from the request's scope, which is the same scope pattern as the Express routes in the DI post, the MCP server is just another delivery mechanism over the same application services.

The transport exposes the server over HTTP, with an API key check and a scope per request:

export function createMcpTransport(container: Container): express.Express {
  const app = express();
  app.use(requestTimerMiddleware);
  app.use(scopeMiddleware(container));
  app.use(contextMiddleware);
  app.use(requestLogger);

  app.post(
    "/mcp",
    requireMcpApiKey(mcpEnv.MCP_API_KEY),
    express.json(),
    async (req, res) => {
      const scope = req.scope;
      const mcpServer = createMcpServer(scope);

      const transport = new StreamableHTTPServerTransport({
        enableJsonResponse: true,
      });

      try {
        await mcpServer.connect(transport as Transport);

        await transport.handleRequest(req, res, req.body);
      } finally {
        // we use a nested try-finally to ensure that all clean up steps are executed even if one of them fails
        try {
          await transport.close();
        } finally {
          try {
            await mcpServer.close();
          } finally {
            await scope.dispose();
          }
        }
      }
    },
  );

  app.use(errorHandlingMiddleware);

  return app;
}
Rendering diagram…
  • a new McpServer and a new transport are created for every request, nothing is shared between two tool calls except the singletons from the container, so one caller can't observe another caller's state.
  • the cleanup is a nested try/finally chain: closing the transport, closing the server and disposing the scope are three independent steps, and each one must run even if the previous one throws.
  • the bootstrap is tiny, the entrypoint builds the MCP container, builds the transport app and listens, everything else is wiring that lives in the composition root.
function bootstrap() {
  const container = buildMcpContainer();
  const app = createMcpTransport(container);

  const port = mcpEnv.MCP_PORT || 8000;

  app.listen(port, () => console.log(`MCP server is running on port ${port}`));
}

bootstrap();

Part V, Run Agent Use Case

  • our current agent loop, runs once with a fresh context and returns the agent response or throws a MaxStepsExceededError, then reset the context for the next prompt, however in a real use case, we would want the user to have a multi prompt conversation with the agent, since one question is rarely enough for someone to decide what product to buy.

The loop itself should stay stateless, so the multi prompt behavior is added around it by a use case: load the conversation, run the agent with the history, save what happened. That use case needs three things the loop doesn't have: a place to persist messages, a way to stop two requests from running the same conversation at the same time, and a way to stop a conversation that outgrew the model.

7. Conversation Persistence

  • conversation persistence means storing the full message history of each conversation, so the next prompt can rebuild the agent's context exactly as it was, including the tool calls, the tool results and the provider state.
  • a conversation belongs to a user, and has a few extra columns that make the lifecycle explicit: is_processing and processing_started_at (is some request currently running this conversation), and max_context_window_reached (has the conversation outgrown the model).
  • every message is a row with a sequence number that gives it its position, and its parts and provider_state are stored as jsonb, because they are exactly the shapes defined by the chat model port.
export const conversation = pgTable(
  "conversation",
  {
    id: varchar("id", { length: 40 }).notNull().primaryKey(),
    user_id: text("user_id")
      .notNull()
      .references(() => user.id, { onDelete: "cascade" }),
    title: varchar("title", { length: 200 }),
    is_processing: boolean("is_processing").notNull().default(false),
    processing_started_at: timestamp("processing_started_at"),
    max_context_window_reached: boolean("max_context_window_reached")
      .notNull()
      .default(false),
    model_id: varchar("model_id", { length: 100 }).notNull(),
    created_at: timestamp("created_at").notNull().defaultNow(),
    updated_at: timestamp("updated_at")
      .notNull()
      .$onUpdate(() => new Date()),
  },
  (t) => [index("conversation_user_id_idx").on(t.user_id)],
);

export const conversationMessage = pgTable(
  "conversation_message",
  {
    id: varchar("id", { length: 40 }).notNull().primaryKey(),
    conversation_id: varchar("conversation_id", { length: 40 })
      .notNull()
      .references(() => conversation.id, { onDelete: "cascade" }),
    sequence: integer("sequence").notNull(), // ordering within the conversation
    role: conversationRolesEnum("role").notNull(),
    parts: jsonb("parts").notNull().$type<ChatPart[]>(),
    provider_state: jsonb("provider_state").$type<unknown>(), // opaque, e.g. Gemini's raw parts
    created_at: timestamp("created_at").notNull().defaultNow(),
  },
  (t) => [
    uniqueIndex("conversation_message_conv_seq_idx").on(
      t.conversation_id,
      t.sequence,
    ),
    index("conversation_message_conversation_id_idx").on(t.conversation_id),
  ],
);
  • the provider_state column is the reason Gemini's thought signatures survive between HTTP requests: the raw parts are saved next to the normalized ones, and the adapter sends them back when the conversation is loaded again.
  • the unique index on (conversation_id, sequence) is a database level guard, two writers can never append different messages at the same position of the same conversation.

The application layer talks to this through a port, the domain style Conversation type and the ConversationRepository contract:

export type Conversation = {
  id: string;
  userId: string;
  title: string | null;
  isProcessing: boolean;
  processingStartedAt: Date | null;
  maxContextWindowReached: boolean;
  modelId: string;
  messages: ChatMessage[];
  createdAt: Date;
  updatedAt: Date;
};

export type ConversationRepository = {
  /** returns an existing convo or null if not found */
  find(conversationId: string): Promise<Conversation | null>;
  /** creates and claims convo by default, throw if creation fails */
  create(userId: string, modelId: string): Promise<Conversation>;
  /** appends new messages to the existing convo */
  appendMessages(
    conversationId: string,
    newMessages: ChatMessage[],
    startIndex: number,
  ): Promise<void>;
  /** isProcessing (claimed) && processingStartedAt > x minutes */
  findStuckConversations(
    batchSize: number,
    stuckforMs: number, // number of milliseconds since the convo's isProcessing nad processingStartedAt were set
    tx?: TransactionClient,
  ): Promise<Conversation[]>;
  /** sets isProcessing to true and processingStartedAt to now, if failed to claim throws error */
  claimConversation(conversationId: string): Promise<void>;
  /** sets isProcessing to false and processingStartedAt to null  if failed to release throws error*/
  releaseConversation(conversationId: string): Promise<void>;
  /** deletes conversation forever, throws if convo is still processing*/
  deleteConversation(
    conversationId: string,
    tx: TransactionClient,
  ): Promise<void>;
  /** sets maxContextWindowReached flag to true, throw if update failsS */
  setConversationCtxLimitAsReached(conversationId: string): Promise<void>;
};

Notice that Conversation.messages is a ChatMessage[], the exact type the agent loop consumes, loading a conversation and calling agent.run() needs no translation. The Postgres implementation, trimmed to the interesting methods:

export class PostgresConversationRepository implements ConversationRepository {
  constructor(private db: DrizzleDBClient) {}

  async find(conversationId: string): Promise<Conversation | null> {
    const row = await this.db.query.conversation.findFirst({
      where: (conversation, { eq }) => eq(conversation.id, conversationId),
      with: {
        messages: {
          orderBy: (message, { asc }) => [asc(message.sequence)],
        },
      },
    });

    if (!row) return null;

    return {
      id: row.id,
      title: row.title,
      userId: row.user_id,
      isProcessing: row.is_processing,
      processingStartedAt: row.processing_started_at,
      maxContextWindowReached: row.max_context_window_reached,
      modelId: row.model_id,
      messages: row.messages.map((message) => ({
        role: message.role,
        parts: message.parts,
        providerState: message.provider_state,
      })), // ordered by sequence number
      createdAt: row.created_at,
      updatedAt: row.updated_at,
    };
  }

  async create(userId: string, modelId: string): Promise<Conversation> {
    const [row] = await this.db
      .insert(conversation)
      .values({
        id: generateConversationId(),
        user_id: userId,
        model_id: modelId,
        is_processing: true, // claimed by default
        processing_started_at: new Date(),
        max_context_window_reached: false,
        title: "New conversation",
      })
      .returning();

    // ... throws a DatabaseError if no row came back, then maps the row
  }

  async claimConversation(conversationId: string): Promise<void> {
    const [row] = await this.db
      .update(conversation)
      .set({ is_processing: true, processing_started_at: new Date() })
      .where(
        and(
          eq(conversation.id, conversationId),
          eq(conversation.is_processing, false),
        ),
      )
      .returning({ id: conversation.id });

    if (!row) {
      throw new ConflictError(
        "conversation",
        conversationId,
        "can't claim a processing conversation",
      );
    }
  }

  async releaseConversation(conversationId: string): Promise<void> {
    const [row] = await this.db
      .update(conversation)
      .set({ is_processing: false, processing_started_at: null })
      .where(
        and(
          eq(conversation.id, conversationId),
          eq(conversation.is_processing, true),
        ),
      )
      .returning({ id: conversation.id });

    if (!row) {
      throw new ConflictError(
        "conversation",
        conversationId,
        "can't release a not processing conversation",
      );
    }
  }

  async appendMessages(
    conversationId: string,
    newMessages: ChatMessage[],
    startIndex: number,
  ): Promise<void> {
    const rows = newMessages.map((message, index) => {
      return {
        conversation_id: conversationId,
        id: generateConversationMessageId(),
        role: message.role,
        parts: message.parts,
        provider_state: message.providerState,
        sequence: startIndex + index,
      };
    });

    await this.db.insert(conversationMessage).values(rows);
  }

  // ...
}
  • find() loads the messages ordered by sequence, so what comes out is the history in the exact order the agent produced it.
  • claimConversation() is a compare and set: the where is_processing = false clause makes the update itself the lock. if two requests race for the same conversation, Postgres lets exactly one update match a row, the other gets no row back and a ConflictError. no explicit lock and no read-then-write race.
  • releaseConversation() is the mirror image, and it also fails loudly when the conversation wasn't claimed, a release without a claim means something is wrong in the flow.
  • create() inserts the conversation already claimed, so a brand new conversation can't be grabbed by a second request between "created" and "claimed".
  • appendMessages() is a single insert statement, so either every message of the turn is saved or none is.

8. Agent Running

export class RunAssistantAgentService {
  constructor(
    private agent: AssistantAgent,
    private conversationRepository: ConversationRepository,
  ) {}

  async execute(query: RunAssistantAgentQuery): Promise<{
    conversationId: string;
    response: string;
  }> {
    let conversationId: string;
    let conversation: Conversation | null;

    if (query.conversationId) {
      conversationId = query.conversationId;
      conversation = await this.conversationRepository.find(conversationId);

      if (!conversation)
        throw new NotFoundError("conversation", conversationId);

      if (conversation.userId !== query.userId)
        throw new ForbiddenError(
          "continue conversation with assistant agent",
          query.userId,
        );

      if (conversation.maxContextWindowReached)
        throw new MaxContextWindowReachedError(conversationId);

      if (conversation.isProcessing)
        throw new ConflictError(
          "conversation",
          conversationId,
          "conversation is already processing",
        );

      await this.conversationRepository.claimConversation(conversationId); // a claim failure will throw
    } else {
      conversation = await this.conversationRepository.create(
        query.userId,
        "gemini-3.6-flash", // hardcoded for now, will be configurable later
      ); // conversationRepository.create claims convo by default

      conversationId = conversation.id;
    }

    const userMessage: ChatMessage = {
      role: "user",
      parts: [{ text: query.prompt }],
    };

    const messages: ChatMessage[] = [...conversation.messages, userMessage];

    try {
      const { response, newMessages } = await this.agent.run([...messages]);
      // returned newMessages don't include userMessage

      await this.conversationRepository.appendMessages(
        conversationId,
        [userMessage, ...newMessages],
        conversation.messages.length, // last existing convo message index + 1, so that new messages are gonna have a sequence starting right after the last existing convo message in DB
      );

      return { conversationId, response };
    } catch (error) {
      if (error instanceof MaxContextWindowReachedError) {
        await this.conversationRepository.setConversationCtxLimitAsReached(
          conversationId,
        );
      }

      throw error; // rethrow error
    } finally {
      await this.conversationRepository.releaseConversation(conversationId);
    }
  }
}
Rendering diagram…
  • the checks before the claim (isProcessing, ownership, context window) give the caller precise errors cheaply, but they are not what protects the conversation, the atomic claimConversation() is. two simultaneous requests can both pass the isProcessing check, only one wins the claim.
  • the finally block releases the conversation no matter what happened in the run, success, a MaxStepsExceededError, a provider failure, so a failed turn doesn't lock the conversation.
  • the agent receives a copy of the history ([...messages]) and the service only persists newMessages plus the user message, with startIndex set to the number of messages that already existed, so the new rows land exactly after the old ones.
  • persistence happens only after a successful run. if the run throws halfway through a tool loop, nothing from that turn is saved, and the conversation stays consistent: it never contains a function call without its response.
  • the agent loop didn't change at all, it still receives an array of messages and returns the new ones.

9. Max Context Window For a Model Reached

  • the whole conversation is fed to the model on every prompt, so the context only grows: each turn adds the user message, the model's tool calls, the tool results (product JSON can be big) and the final answer. every model has a hard limit on how much text it accepts, and a long shopping conversation will eventually cross it.
  • when that happens the provider rejects the request, and the Gemini adapter's error handler (handleGeminiClientErrors) translates that provider specific failure into our own MaxContextWindowReachedError, the agent loop doesn't catch it, so it travels up to the run service.
  • the service reacts in its catch block by calling setConversationCtxLimitAsReached(), then rethrows, and the finally block still releases the conversation.
  • the next prompt on that conversation is rejected at the top of execute() by the maxContextWindowReached check, before claiming anything and before paying for a model call that is guaranteed to fail again. the client is expected to start a new conversation.
  • the flag is cheap and explicit, and because messages are saved only after a successful run, the history stays valid and readable, only the ability to continue it is closed.

Truncating old messages or summarizing them are the two classic ways to keep going instead of stopping, both are future work here, they need care to never cut between a function call and its response.

10. Stuck Conversations

  • claiming is a lock stored in the database, and a lock has the classic failure mode: if the process dies after the claim and before the release (a crash, a deploy, an out of memory kill), the finally block never runs and the conversation stays is_processing forever, the user can never continue it.
  • this is why processing_started_at exists, and why the repository exposes findStuckConversations(batchSize, stuckforMs): every conversation that has been processing for longer than a threshold is, practically, a dead claim.
async findStuckConversations(
  batchSize: number,
  stuckforMs: number,
  tx?: TransactionClient,
): Promise<Conversation[]> {
  const db = tx ? (tx as DrizzleTransactionClient) : this.db;

  const dateToDeleteBefore = new Date(Date.now() - stuckforMs);

  const rows = await db.query.conversation.findMany({
    where: and(
      eq(conversation.is_processing, true),
      lte(conversation.processing_started_at, dateToDeleteBefore),
    ),
    limit: batchSize,
    with: {
      messages: true,
    },
  });

  // ... maps rows to Conversation[]
}
export class ResetStuckConvosCommand {
  constructor(
    public readonly batchSize: number,
    public readonly stuckforMs: number,
  ) {
    this.validate();
  }

  private validate() {
    if (this.batchSize <= 0) {
      throw new ValidationError("batchSize", "must be greater than 0");
    }

    if (this.stuckforMs <= 0) {
      throw new ValidationError("stuckforMs", "must be greater than 0");
    }
  }
}

A cron style worker process (reset-stuck-convos) runs a service with this command on an interval, it's the same pattern as the stuck outbox rows resetter from the DI post: a scope per iteration, a command carrying the batch size and the threshold, and its own composition root.


Part VI, Agent Entry Points

The agent has two ways in, and neither of them contains agent logic: an entry point only builds the container, wires the process specific things (an HTTP server, a readline loop) and delegates. This is the IoC at the edges idea from the DI post, applied to the agent.

11. HTTP Server

async function bootStrap() {
  const container = buildAssistantAgentContainer();

  const mcpClientGateway = container.resolveSingleton(MCP_CLIENT_GATEWAY);

  // load mcp tools at startup
  await mcpClientGateway.loadTools();

  const server = createAssistantAgentServer(container);

  const port = agentEnv.PORT || 8080;

  server.listen(port, () => {
    console.log(`Assistant Agent Server is running on port ${port}`);
  });
}

bootStrap().catch((e) => console.error(e));
export function createAssistantAgentServer(container: Container) {
  const app = express();

  app.use(express.json());
  app.use(cors());

  app.use(requestTimerMiddleware);
  app.use(scopeMiddleware(container));
  app.use(attachUserMiddleware);
  app.use(contextMiddleware);
  app.use(requestLogger);

  app.post("/api/assistant-agent/chat", authMiddleware, async (req, res) => {
    const safeBody = validate(assistantAgentChatBodySchema, req.body);
    const userId = req.user!.id; //  auth middleware ensures req.user is defined

    const service = req.scope.resolve(RUN_ASSISTANT_AGENT_SERVICE);
    const query = new RunAssistantAgentQuery(
      safeBody.query,
      userId,
      safeBody.conversationId,
    );

    const { response } = await service.execute(query);

    res.status(200).json({ response });
  });

  app.use(errorHandlingMiddleware);

  return app;
}
Rendering diagram…
  • the bootstrap resolves the MCP client gateway as a singleton and loads the tools once, at startup, listTools() serves them from the cache for every request after that. a bad MCP URL or API key fails at boot, not on the first customer message.
  • the middlewares are the same family used by the main API (scope per request, request timer, context, request logger, error handling), the agent server is a second Express app over the same building blocks.
  • the route identifies the user through the auth middleware, and passes the user id into the query, this is what the service uses for the ownership check.
  • the route resolves the service from the request scope, and the route itself stays about a dozen lines: validate the body, build the query, call the service, answer.

12. Local REPL

async function bootstrap() {
  const container = buildAssistantAgentContainer();
  const scope = container.createScope();

  const mcpClientGateway = scope.resolve(MCP_CLIENT_GATEWAY);
  const assistantAgent = scope.resolve(ASSISTANT_AGENT);

  await mcpClientGateway.loadTools();

  const rl = readLine.createInterface({
    input: process.stdin,
    output: process.stdout,
  });

  while (true) {
    const input = (await rl.question("you> ")).trim();
    if (input === "exit") break;

    const { response } = await assistantAgent.run([
      { role: "user", parts: [{ text: input }] },
    ]);

    console.log(`agent> ${response}\n`);
  }

  await mcpClientGateway.close();
  rl.close();
}

bootstrap().catch((e) => {
  logger.error("agent crashed", e as Error);
  process.exit(1);
});
  • the REPL resolves AssistantAgent directly and skips the run service, so there is no conversation, no user and no database involved in the loop: every line you type starts from a fresh context, which is the "runs once with a fresh context" behavior from the beginning of this section.
  • that makes it a very fast way to iterate on the system prompt and the tool descriptions: change a sentence, restart, ask the same question, and watch which tools the agent picks (the agent logs every step).
  • it's also the place where close() earns its spot in the gateway contract, on exit the MCP connection is closed explicitly.

Part VII, One Process Per Concern

Everything in this post runs as its own process, the same process shaped composition idea from the DI post: the MCP server, the agent HTTP server, the embedding queue handler and the stuck conversations resetter are four different entrypoints over one image, each with its own composition root and its own environment file.

x-service-defaults: &service-defaults
  image: ghcr.io/djaouad10/ddd-e-commerce-backend:latest
  restart: unless-stopped
  env_file: .env
  # ...

services:
  mcp:
    <<: *service-defaults
    env_file:
      - .env
      - .env.mcp
    command: ["node", "./dist/entrypoints/mcp/index.js"]
    ports:
      - "8000:8000"
    # healthcheck: ...

  assistant-agent:
    <<: *service-defaults
    env_file:
      - .env
      - .env.agent
    command: ["node", "./dist/entrypoints/agents/assistant-agent.js"]
    ports:
      - "8080:8080"
    # healthcheck: ...

  embedding-queue-handler:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/embedding-queue-handler.js"]

  reset-stuck-convos:
    <<: *service-defaults
    command: ["node", "./dist/entrypoints/workers/reset-stuck-convos.js"]
    scale: 1
Rendering diagram…
  • the agent and the MCP server are separate on purpose: the tools can scale, deploy and fail independently of the agent, and any other agent (not only this one) can connect to the same MCP endpoint.
  • the agent process needs the chat model and the conversations table, the MCP process needs the embedding model and the product tables, each env file only contains what its process uses.
  • the workers keep the knowledge base fresh and the conversations unlocked without ever sitting on the request path.

Part VIII, Close

Recap in one breath: a stateless loop (agent) → two ports (tools gateway and chat model) → an MCP server exposing application services as tools → an event driven RAG pipeline keeping the semantic index fresh → persisted conversations with an atomic claim → entry points that only wire.

The repo: github.com/djaouad10/DDD-E-commerce-Backend.

Thanks for reading 👋More posts →