Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2,983 changes: 1,591 additions & 1,392 deletions Pipfile.lock

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion docs/open-api-docs.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ openapi: 3.0.3
info:
title: The Agent's user-facing API
description: The user-facing parts of The Agent's API service (excluding system-level endpoints, chat completion, maintenance endpoints, etc.)
version: 5.34.1
version: 5.35.0
license:
name: MIT
url: https://opensource.org/licenses/MIT
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
schema: spec-driven
created: 2026-09-07
Original file line number Diff line number Diff line change
@@ -0,0 +1,123 @@
## Context

Before this change, Telegram and WhatsApp webhook endpoints enqueued synchronous responders through FastAPI `BackgroundTasks`. Each responder created one detached SQLAlchemy session, ingested the message, constructed `ChatAgent`, and performed command, debounce, reply-decision, LLM, and send work. `ChatAgent` slept for the configured debounce delay and independently searched message history for a newer message from the same invoker.

The application runs multiple service instances against one PostgreSQL database. Process-local cancellation or locking therefore cannot establish a single burst winner. Database sessions must not remain checked out while waiting or while external services block. See `proposal.md` for the behavioral motivation and `specs/message-burst-processing/spec.md` for the required behavior.

## Goals / Non-Goals

**Goals:**

- Make PostgreSQL the authority for the active burst state and deadline.
- Keep webhook responses immediate while using a non-blocking 1 second delayed attempt.
- Preserve a strict session boundary around each short database phase.
- Pass a stable history cutoff and aggregate addressing decision into conversational processing.
- Keep the coordination mechanism small enough to share between Telegram and WhatsApp responders.

**Non-Goals:**

- Automatic retries, periodic recovery workers, or exactly-once guarantees across process failure and external platform sends.
- Platform-specific media album grouping.
- Changing command behavior, group reply probability, or LLM/tool behavior beyond invoking reply logic once per settled burst.
- Introducing Redis, a message broker, or another external dependency.

## Decisions

### Store one active burst mailbox per chat and author

Add a coordination table with a composite primary key of `(chat_id, author_id)`. A row represents outstanding conversational work rather than a permanent event log. It contains:

- a monotonically increasing `message_count`;
- `process_after`, calculated from database time plus 1 second;
- the final logical message cutoff;
- an aggregate `is_addressed` flag;
- an `is_processing` flag.

On each newly ingested non-command message, an upsert increments the message count, resets `process_after`, advances the logical cutoff when appropriate, and ORs the addressing flag. The operation returns an immutable scheduled-burst value containing the mailbox identity, message count, and wait duration.

This model is preferred over processing status on every chat message because the claimable unit is a coalesced burst, not each fragment. It also avoids accumulating completed queue rows. The mailbox row is deleted after successful processing when its message count has not advanced.

Alternative considered: add `received`, `processing`, and `completed` to every chat message. This would require grouping and transitioning several history rows for one invocation, including a separate representation for coalesced fragments, while still needing a per-author coordination lock.

### Use asynchronous delayed attempts initiated by webhooks

After the ingress transaction commits and its detached session closes, the background responder creates a DI instance containing the burst invoker and chat identifiers but no database session, then invokes its `MessageBurstService`. The suspended coroutine retains that contextual service and the immutable schedule, but no SQLAlchemy session, transaction, or checked-out connection.

After waking, it opens a new detached session and performs one atomic conditional claim. The claim succeeds only when:

- the stored message count equals the attempt's expected message count;
- `process_after` is no later than database time; and
- the mailbox is not currently being processed.

An attempt that does not satisfy those conditions closes its fresh session and exits. Every message may therefore create a small suspended coroutine, but older attempts become inexpensive no-ops. FastAPI continues returning the webhook response before these background attempts execute.

The delayed orchestration function remains asynchronous. Existing synchronous ingestion and conversational processing execute through the thread pool rather than on the event-loop thread.

Both platform responders now use the same per-scheduled-burst helper and contextual DI construction. Telegram awaits its single scheduled burst, while WhatsApp gathers the helper once for each scheduled message because one webhook may contain multiple messages.

Alternative considered: a periodic database poller. It would make recovery stronger, but it is unnecessary for the current scope and adds a continuously running subsystem. Outstanding rows can be inspected directly in the database when needed.

### Separate database phases from waiting and external work

The flow uses explicit resource boundaries:

1. Open a detached session, ingest the message, classify commands/addressing, upsert the burst mailbox, commit, and close the session.
2. Create a DI instance containing the burst invoker and chat identifiers but no database session, then await the 1 second timer.
3. Clone that contextual DI with a fresh detached session, conditionally claim the burst and its cutoff, commit, and close the session.
4. Load the bounded inputs needed for processing with fresh database access, then release that transaction before LLM, media, or platform network calls.
5. Clone the contextual DI with another short fresh transaction to finalize the mailbox after the response path completes.

The claim session is never passed into `ChatAgent`. A `MessageBurstService` created from the cloned contextual DI loads the claimed message and owns the shared processing path, preserving the resource-release discipline around long-running external operations.

### Preserve pending work that arrives during processing

A new message may arrive after an earlier message count has been claimed. Its upsert still increments the current message count and establishes a new quiet deadline without modifying the claimed cutoff.

When processing finishes, finalization compares the current message count with the claimed message count:

- If they match, delete the mailbox row.
- If they differ, clear the processing flag while retaining the newer messages and return a new scheduled-burst value for the existing deadline, with no delay if that deadline has already passed.

This completion handoff prevents newer messages from being stranded when their own delayed attempt woke while the preceding burst was still processing.

### Determine reply eligibility once per settled burst

Commands are recognized before burst upsert and continue through their immediate path. They neither wait nor change an existing conversational burst deadline.

For non-command messages, ingress contributes whether that individual message explicitly addressed the bot. The mailbox retains the logical OR across the active burst. On claim, `ChatAgent` evaluates the existing `should_reply` behavior once using the claimed cutoff and addressed state:

- private chats are addressed;
- group bursts with any explicit tag are addressed even if the final fragment is untagged;
- untagged group bursts retain the existing probabilistic/non-mention decision behavior.

Responders classify commands and record conversational messages; `MessageBurstService` owns burst coordination, claimed-message hydration, and the shared processing path. Conversational decision logic remains in `ChatAgent`.

### Add a deterministic history cutoff

Add a database-generated monotonic order column to `chat_messages`. Logical history order uses platform `sent_at` followed by this database order as a tie-breaker. The mailbox stores the final cutoff tuple for the burst, and history loading for a claimed burst is bounded by that tuple.

Every arrival resets the quiet deadline, including a late-delivered message with an earlier platform timestamp. The stored cutoff advances only to the greatest logical message tuple, allowing earlier late-delivered fragments to be included without moving the invocation past logically later chat messages.

This replaces the current ambiguous same-timestamp behavior, where differing message IDs can each be treated as newer depending on which handler evaluates them.

### Keep successful mailbox state ephemeral

Successful finalization deletes the mailbox row when its message count has not advanced.

No `completed` queue state is retained, and outstanding mailbox rows remain directly inspectable.

## Risks / Trade-offs

- [A platform delivery gap exceeds 1 second] → The later delivery forms another burst; the threshold may need adjustment if this is observed operationally.
- [One coroutine is scheduled per delivered message] → Coroutines sleep without threads or database resources, and obsolete schedules perform one short conditional check before exiting.
- [A process exits after accepting a webhook] → Outstanding mailbox state remains in the database, but automatic recovery is outside this change.
- [A process exits after sending a reply but before finalizing the mailbox] → Exactly-once recovery across the external send boundary is explicitly outside scope.
- [Active processing overlaps newer messages] → Retain the newer messages in the mailbox and schedule them when active processing finalizes.

## Migration Plan

1. Add the deterministic message-order column and active burst mailbox model, including model imports in `src/db/alembic/env.py`.
2. Have the user generate the Alembic migration with `./tools/db_generate_migration -y` and review its backfill, constraints, foreign keys, and indexes.
3. Deploy the schema before application instances that reference the new coordination table.
4. Deploy application instances with shared burst coordination and remove the old debounce path in the same application release.
5. If rollback is required, restore the old application version before removing the new table or message-order column; the additional schema is otherwise harmless to the old code.
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
## Why

Concurrent webhook handlers currently make independent timing and history decisions for messages that may belong to one logical user burst. Across multiple service instances, this can drop the intended final response, answer an earlier fragment, or process incomplete photo-and-prompt context.

## What Changes

- Group non-command messages into bursts per chat and author using a 1 second quiet period.
- Persist the active burst deadline and message count in the shared database so all service instances make one coordinated processing decision.
- Schedule non-blocking asynchronous wake-ups without retaining a database session or transaction during the quiet period.
- Atomically allow only the wake-up matching the current settled message count to claim processing.
- In group chats, reply once when any message in the settled burst addressed the bot; bursts from different authors remain independent.
- In private chats, treat every settled burst as addressed to the bot and reply once using all messages through the burst cutoff.
- Keep commands outside burst processing so they receive an immediate response.
- Remove the existing per-message sleep-and-supersession behavior.

## Capabilities

### New Capabilities

- `message-burst-processing`: Defines quiet-period burst formation, cross-instance claiming, command bypass, reply eligibility, context boundaries, and resource-release behavior.

### Modified Capabilities

None.

## Impact

- Affects Telegram and WhatsApp inbound message handling, chat-agent invocation, chat-history loading, dependency injection, configuration, and database models/repositories.
- Requires a database migration for shared burst coordination and deterministic message cutoff data.
- Changes non-command reply timing to begin after 1 second without a message from the same author in the same chat.
- Does not add automatic retry or continuous recovery processing; outstanding work remains visible in the database.
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
## Purpose

Define how incoming conversational messages are combined into user bursts and processed exactly once across service instances without delaying commands or retaining database resources while waiting.

## ADDED Requirements

### Requirement: Messages form per-author quiet-period bursts
The system SHALL group non-command messages by chat and author. A burst SHALL become eligible for processing only after 1 second has elapsed since the most recently received non-command message for that chat and author.

#### Scenario: Several messages arrive within the quiet period
- **WHEN** multiple non-command messages from one author in one chat arrive with less than 1 second of silence between them
- **THEN** the system treats them as one burst and starts no response processing until 1 second after the last message arrives

#### Scenario: A later message arrives after the quiet period
- **WHEN** a message arrives more than 1 second after the preceding burst became eligible and was claimed
- **THEN** the system treats the later message as part of a new burst

#### Scenario: Different authors send messages concurrently
- **WHEN** messages from different authors arrive in the same group chat
- **THEN** each author's quiet period and burst are tracked independently

### Requirement: Burst claiming is coordinated across instances
The system SHALL use shared database state to ensure that no more than one service instance claims a settled burst for response processing.

#### Scenario: Multiple delayed attempts wake for one burst
- **WHEN** delayed attempts created by several messages or service instances wake for the same chat and author
- **THEN** only the attempt matching the current settled message count claims response processing

#### Scenario: An obsolete attempt wakes
- **WHEN** a delayed attempt wakes after a newer message has extended the burst deadline
- **THEN** the obsolete attempt performs no response processing

### Requirement: Waiting does not retain database resources
The system SHALL commit and release the database session and transaction used to record a message before waiting for the quiet period. The delayed attempt SHALL acquire a fresh database session only after its asynchronous wait completes.

#### Scenario: A burst is waiting for additional messages
- **WHEN** a delayed attempt is suspended during the 1 second quiet period
- **THEN** that attempt holds no database transaction or checked-out database connection

### Requirement: Commands bypass burst processing
The system SHALL process each recognized command immediately without waiting for the burst quiet period. A command SHALL NOT extend, replace, or be claimed as part of a conversational burst.

#### Scenario: A command arrives without other messages
- **WHEN** the system receives a recognized command
- **THEN** command processing begins without a burst delay

#### Scenario: A command arrives during a conversational burst
- **WHEN** a recognized command arrives while non-command messages from the same author are waiting for their quiet period
- **THEN** the command is processed immediately and the existing conversational burst retains its own deadline

### Requirement: Group-chat reply eligibility is evaluated per burst
For a group chat, the system SHALL preserve the existing reply-decision behavior while evaluating it once for each settled author burst. If any message in the burst explicitly addresses the bot, the burst SHALL be treated as explicitly addressed even when the final burst message does not repeat the tag.

#### Scenario: A tagged message is followed by untagged fragments
- **WHEN** an author sends a tagged message followed within the same burst by untagged attachments or text
- **THEN** the system evaluates and processes one explicitly addressed burst containing all of those messages

#### Scenario: Two authors tag the bot
- **WHEN** two authors independently send tagged bursts in one group chat
- **THEN** each author's settled burst remains eligible for its own reply

#### Scenario: A burst does not tag the bot
- **WHEN** no message in a settled group-chat burst explicitly addresses the bot
- **THEN** the existing non-mention reply-decision behavior is applied once to that burst

### Requirement: Private-chat messages are evaluated as one addressed burst
The system SHALL treat every settled private-chat burst as explicitly addressed and SHALL produce at most one response for that burst.

#### Scenario: Photos and a prompt arrive as separate deliveries
- **WHEN** a private-chat user sends multiple photos and prompt text within one burst
- **THEN** the system invokes conversational processing once with all burst messages available as context

### Requirement: Processing uses a stable burst cutoff
The system SHALL evaluate a burst using chat history through the burst's final logical message and SHALL exclude messages logically after that cutoff. Message ordering SHALL be deterministic when platform timestamps are equal or deliveries arrive out of order.

#### Scenario: Another message arrives after a burst is claimed
- **WHEN** a burst has been claimed and a later message is received
- **THEN** the claimed processing uses its original cutoff and the later message is assigned to subsequent processing

#### Scenario: Messages have equal platform timestamps
- **WHEN** two stored messages have equal platform timestamps
- **THEN** the system applies a stable database-backed ordering to determine the burst cutoff and history order

Loading
Loading