Skip to content

MarketerAI Module Architecture - Information Flow ​

Table of Contents ​

  1. Overview
  2. Conversation Initialization
  3. AI Response Streaming
  4. Function Calling
  5. Subagent Mechanism
  6. Confirmation System
  7. Sequence Diagram
  8. Example Scenarios

Overview ​

This document describes the detailed information flow from the moment a user sends a message, through AI processing and function calls, to receiving a response.


Conversation Initialization ​

Scenario 1: New Conversation ​

Endpoint: POST /api/projects/{project}/agents/{agent}/chat

Request Flow ​

  1. Frontend sends: POST request with body containing the user message ("Show me top 10 products")

  2. Controller (SembotChatController.create()):

    • Starts a database transaction
    • Creates a new ChatThread with user_id, project_id, agent_id, and title (truncated to 200 characters)
    • Updates LastAgentUsage (for agent last-use statistics)
    • Commits the transaction
    • Dispatches the ChatStream job to the asynchronous queue
    • Returns 200 OK immediately (empty response)
  3. Frontend:

    • Receives 200 OK
    • Listens on the private WebSocket channel: App.Chat.{userId}.{projectId}
    • Waits for events: ThreadStart, StreamMessageChunk, ToolCall, ToolResponse, MessageEnd, ThreadEnd

Key point: The HTTP response is immediate — we do not wait for AI processing. All further communication continues via WebSocket in real time.


Scenario 2: Continuing a Conversation ​

Endpoint: POST /api/projects/{project}/chat-threads/{thread}

Flow: Identical to Scenario 1, but instead of creating a new ChatThread, we use an existing one. Conversation history is preserved — the AI sees all previous messages, tool calls, and tool responses.

Example: If the user asked "Show top 10 products" in the first message and the AI returned a list, the user can then write "Change labels for first 5" in the next message — the AI knows which products were referenced in the previous response.


AI Response Streaming ​

Step 1: ChatStream Job Starts ​

What happens:

  1. The job is pulled from the Redis queue by a queue worker
  2. Broadcasts a ThreadStart event (frontend shows loading indicator)
  3. Initializes ChatService and ChatMessageRepository
  4. Checks whether this is a resume after a confirmation (if so → different logic)
  5. If this is a UserMessage object (predefined message) → generates its content via AI
  6. Auto-reject pending confirmations: If a pending confirmation exists in this thread or its subthreads, it is automatically rejected and the AI is informed that the user sent a new message
  7. Creates a USER_MESSAGE in the database
  8. Calls ChatService.handleStreamedAgent()

Step 2: ChatService Selects API Format ​

Decision:

  • If provider=OpenAI AND config('ai.default_api')='responses' → Responses API
  • Otherwise → Completions API

Why it matters: The Responses API is newer (2024), has a better event structure, and supports structured outputs. The Completions API is older but more universal (also works with Groq, OpenRouter).


Step 3: OpenAiService Creates a Stream ​

What happens:

  1. Fetches message history from ChatMessageRepository (last N user/assistant/tool_call/tool_response messages)
  2. Prepends the agent's system prompt
  3. Fetches available tools (functions from the database + runtime functions like searchKnowledgeBase)
  4. Moderates the user's last message (checks for hate speech/violence)
  5. Calls the OpenAI API with stream=true
  6. Returns a StreamResponse — an iterator that produces chunk by chunk

Step 4: ChatService Processes the Stream ​

For each chunk:

If it is text (delta):

  • Accumulates in a local variable
  • Broadcasts StreamMessageChunk event (frontend appends to the displayed message)

If it is a tool call delta:

  • Accumulates ID, function name, and JSON arguments
  • Broadcasts FunctionGenerationStart (at the beginning)
  • Broadcasts FunctionArgumentsChunk for each fragment of arguments

If finish_reason='stop':

  • Saves ASSISTANT_MESSAGE in the database
  • Broadcasts MessageEnd
  • Ends iteration

If finish_reason='tool_calls':

  • Saves ASSISTANT_TOOL_CALL in the database
  • Calls handleFunctionCalling()
  • If the function returned that another iteration is needed → creates a new stream with tool_response in context

Function Calling ​

Tool Call Flow ​

Example: AI calls getTopProducts(limit=10, sortBy="revenue")

  1. handleFunctionCalling receives:

    • toolCallId: "call_abc123"
    • functionName: "getTopProducts"
    • arguments:
  2. Broadcast ToolCall event: Frontend shows "Calling getTopProducts..." with arguments

  3. callChatFunction() decides:

    • Is this a runtime function? (isRuntimeFunction("getTopProducts") → false)
    • Fetches ChatFunction from the database for this agent
    • Does it have a subagent_id? (no)
    • Calls executeChatFunctionClass()
  4. executeChatFunctionClass():

    • Laravel DI creates an instance of App\Services\Ai\Chat\ChatFunctions\Products\GetTopProducts
    • Injects User, Project, Agent via the constructor
    • Validates arguments according to the rules() methods (limit: integer|min:1|max:100, sortBy: in:revenue,units,clicks)
    • If validation fails → returns "Validation failed: limit must be between 1 and 100"
    • If OK → calls handle($validated, $additionalData)
  5. GetTopProducts.handle():

    • Executes a database query: fetches the project's products, sorts by revenue DESC, limit 10
    • Formats the result as text: "Product 1: iPhone 15 Pro, Revenue: $52,450, Units: 124\nProduct 2: ..."
    • Optionally adds $additionalData['table_data'] with a table structure for the frontend
    • Returns a text response
  6. Save and broadcast:

    • Creates a TOOL_RESPONSE message in the database with toolCallId="call_abc123" and the response
    • Broadcasts ToolResponse event
    • Frontend can display a table (if additionalData is present)
  7. Next iteration:

    • handleFunctionCalling returns true (another iteration is needed)
    • handleCompletionsStream creates a new stream
    • The AI now sees the tool_response in context and can respond to the user: "Here are the top 10 products by revenue: [formatted results]..."

Subagent Mechanism ​

When It Is Used ​

When a function has a subagent_id set — instead of executing PHP code, we delegate to a specialized AI agent.

Example use case: The main agent (General Assistant) has a "manageProducts" function that delegates to the "Product Manager" subagent. The Product Manager has its own system prompt, its own functions (getTopProducts, changeCustomLabel, duplicateProducts), and its own AI model.

Subagent Flow ​

User → Main Agent: "Manage products: change labels for top 10"

  1. Main Agent calls: manageProducts(user_request="change labels for top 10")

  2. handleSubagentFunction():

    • Fetches the subagent (Product Manager) from function.subagent
    • Prepares a custom system prompt (if present in function.subagent_system_prompt)
    • Prepares the user message with placeholders (function.subagent_user_message → "User wants to: {user_request}")
    • Checks the preserve_previous_context parameter
  3. SubagentThreadService:

    • If preserve_previous_context=true → looks for an existing subthread
    • If found → returns it (continuation with context)
    • If not found or preserve=false → creates a new subthread with parent_thread_id
  4. Creating the structure:

    • Saves an AGENT_SUBTHREAD message in the parent thread (link to the subthread)
    • Broadcasts AgentSubthreadStart (frontend creates a nested container)
    • Creates a USER_MESSAGE in the subthread
    • Broadcasts SubthreadUserMessage
  5. Recursion:

    • ChatService.handleStreamedAgent() is called RECURSIVELY for the subagent
    • The subagent has its own streaming, its own tool calls, and its own events
    • All events include subthread_id so the frontend knows where to display them
  6. Subagent works:

    • Calls getTopProducts(limit=10) → gets the list
    • Calls changeCustomLabel(products=[1,2,3,4,5,6,7,8,9,10], label="Sale")
    • STOP — the function requires confirmation!
  7. Pending in subthread:

    • Creates PENDING_USER_CONFIRMATION in the subthread
    • Broadcasts ActionRequiresConfirmation with details
    • Job ends (return)
  8. User approves:

    • Frontend sends POST /confirmations/{id} with approved=true
    • Dispatches a new ChatStream job with pendingConfirmation + approved flag
  9. Subthread resumes:

    • handleConfirmationResume() executes changeCustomLabel
    • Broadcasts ToolResponse for the subthread
    • Subagent streaming continues
    • Subagent responds: "Changed custom_label_0 to 'Sale' for 10 products"
    • Ends (no more tool calls, finish_reason='stop')
  10. Return to parent:

    • continueParentThreads() is called
    • createToolResponseInParentThread() — extracts the subagent's last response and saves it as a tool_response in the parent thread
    • Continues main agent streaming with tool_response in context
    • Main agent responds: "Done! I've managed your products - changed labels for top 10."

Hierarchy:

ChatThread (id=1, main thread)
  ├─ USER_MESSAGE: "Manage products..."
  ├─ ASSISTANT_TOOL_CALL: manageProducts()
  ├─ AGENT_SUBTHREAD: → subthread_id=2
  │
  └─ ChatThread (id=2, subthread, parent_thread_id=1)
       ├─ USER_MESSAGE: "User wants to: change labels for top 10"
       ├─ ASSISTANT_MESSAGE: "Let me get top products..."
       ├─ ASSISTANT_TOOL_CALL: getTopProducts()
       ├─ TOOL_RESPONSE: "Product 1: ..., Product 2: ..."
       ├─ ASSISTANT_TOOL_CALL: changeCustomLabel()
       ├─ PENDING_USER_CONFIRMATION: [PAUSE]
       ├─ TOOL_RESPONSE: "Changed labels for 10 products"
       └─ ASSISTANT_MESSAGE: "Changed custom_label_0 to 'Sale' for 10 products"
  
  ├─ TOOL_RESPONSE: "Changed custom_label_0 to 'Sale' for 10 products" (from subthread)
  └─ ASSISTANT_MESSAGE: "Done! I've managed your products..."

Confirmation System ​

Trigger Confirmation ​

Condition: ChatFunction.requires_user_confirmation = true

Examples: changeCustomLabel, editCampaignStatus, addNegativeKeywords, createTask

Flow ​

  1. AI calls a function that requires confirmation
  2. handleFunctionCalling() detects: function.requires_user_confirmation === true
  3. Creates PENDING_USER_CONFIRMATION message:
    • Contains: tool_call_id, function_name, arguments, thread_context (thread_id, parent_thread_id)
  4. Broadcasts ActionRequiresConfirmation:
    • Frontend shows a modal: "Change custom_label_0 to 'Sale' for 5 products?" [Cancel] [Confirm]
  5. Job ends — return (function is not executed, streaming does not continue)

User Decision ​

Cancel:

  • POST /confirmations/{id} with approved=false
  • New job with pendingConfirmation + approved=false
  • handleConfirmationResume() creates a tool_response: "User declined this action."
  • Broadcasts ActionRejected
  • Streaming continues — the AI knows the user declined and responds e.g. "Okay, I won't change the labels."

Confirm:

  • POST /confirmations/{id} with approved=true
  • New job with pendingConfirmation + approved=true
  • handleConfirmationResume() EXECUTES the function
  • Broadcasts ActionConfirmed
  • Broadcasts ToolResponse with the result
  • Streaming continues — the AI sees the result and responds e.g. "Done! Changed labels for 5 products."

Auto-Reject ​

Situation: The user has a pending confirmation but sends a new message instead of confirming or rejecting.

Resolution:

  • ChatStream.handle() checks getActivePendingConfirmationInThreadTree() at startup
  • If found → automatically rejects (declinePendingAction)
  • Creates a tool_response: "Previous action was automatically declined because user sent a new message."
  • Broadcasts ActionRejected with automaticRejection=true
  • Continues with the current thread (the one with the old pending) — the AI is informed that the user declined
  • Then processes the new user message normally

Why: The user changed their mind — instead of confirming "change labels" they sent "show me analytics." We do not leave hanging confirmations.


Sequence Diagram ​

Full Conversation with Subagent and Confirmation ​


Example Scenarios ​

Scenario 1: Simple Function Without Confirmation ​

User: "Show me top 5 products"

Timeline:

  • 0ms: ThreadStart
  • 50ms: StreamMessageChunk "Let"
  • 55ms: StreamMessageChunk " me"
  • 100ms: FunctionGenerationStart
  • 150ms: FunctionArgumentsChunk (JSON fragments)
  • 200ms: ToolCall getTopProducts
  • 1500ms: ToolResponse (with product data)
  • 1600ms: StreamMessageChunk "Here"
  • 1700ms: MessageEnd
  • 1750ms: ThreadEnd
  • 1800ms: ProposeActions

Result: User sees a list of 5 products + suggestions ("Analyze trends", "Change prices")


Scenario 2: Function With Confirmation ​

User: "Change custom_label_0 to 'Sale' for products 1, 2, 3"

Timeline:

  • 0ms: ThreadStart
  • 100ms: FunctionGenerationStart
  • 200ms: ToolCall changeCustomLabel
  • 250ms: ActionRequiresConfirmation → PAUSE
  • [User interaction 5000ms]
  • 5000ms: User clicks Confirm
  • 5010ms: ThreadStart (resume)
  • 5200ms: ToolResponse "Changed labels for 3 products"
  • 5300ms: StreamMessageChunk "Done!"
  • 5400ms: ThreadEnd

Result: Labels changed after confirmation


Scenario 3: Subagent With Nested Confirmation ​

User: "Manage products: change label for top 10"

Flow:

  1. Main thread → calls manageProducts
  2. Subthread created → Product Manager
  3. Subagent calls getTopProducts → success
  4. Subagent calls changeCustomLabel → requires confirmation
  5. PAUSE in subthread
  6. User confirms
  7. Subthread ends
  8. Parent thread continues
  9. Main agent responds

Thread Structure:

  • Main Thread (id=1)
    • Subthread (id=2, parent=1) ← confirmation happened here
  • Result saved in both threads

Created: 2025-11-26 Version: 1.0