Message Lifecycle¶
How a message travels from Telegram to a response and back. These flows are Telegram-specific — WhatsApp follows the same logical path (intent routing → agent → log) but via a webhook handler (app/interfaces/whatsapp/webhook.py) without streaming or TTS. Each Telegram message type has a slightly different path.
Text Message (Primary Flow)¶
Text messages are not streamed. A "…" placeholder goes out immediately, the agent runs to completion, and the placeholder is edited once with the final answer.
Why streaming was removed
run_stream returns the first model response and treats any text in it as the final result — even when that text is followed by tool calls (pydantic_ai/agent/abstract.py). On Gemini that routinely made a "Let me check…" preamble the user-visible answer while the actual post-tool synthesis was silently dropped. The flow now uses agent.run via run_with_resilience.
sequenceDiagram
participant You
participant TG as Telegram
participant H as handle_message
participant R as Intent Router
participant S as Storage
participant A as Agent (domain or full)
participant MR as run_with_resilience
participant LLM as LLM
participant T as Tools
participant GW as approval_gate
You->>TG: Send text message
TG->>H: Update event (long-polling)
H->>H: Check allowlist (ALLOWED_TELEGRAM_USER_IDS)
H->>H: Check edit state (user_data["pending_edit_action_id"])
H->>S: fetch_message_history(user_id, message, settings)
S-->>H: chronological 10 OR (3 recent + 3 semantic) if ENABLE_SEMANTIC_HISTORY → build_message_history (drops oldest if >6,000 est. tokens; scrubs persona-drift turns to "[earlier response omitted]")
H->>R: classify_intent(text)
Note over R: keyword → semantic top-K (≥0.65 single, 0.45–0.65 top-K, <0.45 empty) → context_hint
R-->>H: categories e.g. {"email", "utility"} or {"github","health","utility"} (composed)
H->>R: select_agent(categories)
R-->>H: email_agent / composed_agent / full_agent
H->>H: request_deps = _deps_replace(deps, active_categories, chat_id, pending_action_user_message=text)
H->>TG: send_chat_action("typing")
Note over H,S: Context injection (shared 1,000-token budget)
H->>S: find_relevant_notes (≥0.6 similarity + recency boost)
H->>S: find_relevant_summaries (≥0.6 similarity + recency boost, if budget remains)
H->>S: find_relevant_read_later (tag overlap, if budget remains)
H->>H: Wrap blocks in XML tags + prepend datetime → enriched_text
Note over H: Planning gate (Spec 008) — runs here, after routing and enrichment
H->>H: classify_complexity(text) — regex on connectives, ≥30 chars
opt Complex → generate_plan returns 2+ steps
H->>S: save PendingAction(action_type="plan")
H->>TG: Plan preview + [✅ Confirm] [❌ Cancel] [✏️ Edit]
Note over H: return — the agent does not run on this turn
end
H->>TG: Send "…" placeholder message
H->>MR: run_with_resilience(run_agent, enriched_text, model, fallback_model=mini, message_history)
MR->>A: agent.run(...) — retries 429/5xx with backoff, then falls back to the mini model
A->>S: get_context("global")
S-->>A: UserContext (reflection profile or empty)
A->>A: build_system_prompt() via @agent.instructions — injects active_categories focus hint
A->>LLM: instructions + message_history + user_message
loop 0 or more tool calls
LLM->>T: tool_name(args)
alt Gated tool (create_task, send_email, etc.)
T->>GW: approval_gate(tool_name, payload, deps)
GW->>S: save_pending_action(PendingAction)
GW-->>T: "[APPROVAL_PENDING:<uuid>]" sentinel
else Read-only tool
T-->>LLM: result
end
T-->>LLM: result (sentinel or direct)
end
LLM-->>A: final response ("action is pending your approval")
A-->>MR: RunResult(output, all_messages, usage)
MR-->>H: RunResult — or re-raises the primary error → "I'm having trouble thinking right now."
opt Voice-trigger phrase matched and response short enough
H->>TG: Send a voice note instead of a text edit
end
H->>TG: Edit placeholder with the complete response
Note over H,TG: Splits at newline if >4096 chars
H->>H: Scan tool calls for sentinels
opt Sentinel(s) found
H->>S: get_pending_action(uuid)
H->>TG: Edit sent_msg with preview + [✅ Confirm] [❌ Cancel] [✏️ Edit]
H->>S: update_pending_action(message_id=sent_msg.id)
end
H->>S: log_interaction(user_msg, response, channel="telegram")
H->>S: _write_audit_entries() — sanitised tool call records
Note over H: Fire-and-forget post-turn tasks
H-)S: extract_facts_from_exchange (mini-model)
H-)S: detect_and_save_capability_gap — regex prefilter → mini-classify → grep_source verify → capability_gaps row (gap | available | unknown)
Notes on the streaming flow and approval gate¶
- Intent routing runs early —
classify_intent()keyword-matches the message in microseconds; on a miss it falls back to a semantic classifier with three confidence bands (≥0.65 single best domain, 0.45–0.65 top-K composed agent, <0.45 full agent).select_agent()then picks one of 13 domain agents (email, calendar, memory, github, news, slack, utility, meetings, jira, drive, diagnostics, health, database) or returns the full agent. active_categoriesis set onAgentDepsbefore the agent run —build_system_promptinjects a focus hint telling the model which tool categories to prefer and only includes tool guidance for those domains- Context injection — before the agent run, the handler runs three retrieval layers sharing a
CONTEXT_TOKEN_BUDGET = 1000token budget: (1)find_relevant_notes(≥0.6 similarity), (2)find_relevant_summaries(≥0.6), (3)find_relevant_read_later(tag overlap). Each layer's output is wrapped in XML tags (<context type="notes">,<context type="summaries">,<context type="read_later">) and prepended to the user message so the model can distinguish retrieved memory from instructions. - Datetime in user turn —
[Thursday, April 24, 2026 — 09:15 AM Europe/Paris]is prepended to the enriched message (not the system prompt) to keep the stable system prompt prefix identical across requests for Gemini implicit cache hits. - Planning runs after routing, not before.
classify_complexityis evaluated once history, intent routing, and context injection are already done — so a message that turns out to be a plan has paid for that work, and one that doesn't reaches the agent with everything in place. (The plan flow below picks up from here.) - Model resilience wraps the run.
run_with_resilienceretries transient429/5xxagainst the primary model with backoff, then falls back toMINI_MODEL_NAME. Only when every attempt fails does the handler'sexceptproduce "I'm having trouble thinking right now." — see Model Resilience. - The placeholder
"…"message is edited in-place rather than sending a new message — this prevents the chat from jumping around while you wait - If the final response exceeds 4096 characters, it's split at the last newline before the limit and sent as sequential messages
- Approval gate — when the agent calls a consequential tool (create task, send email, etc.),
approval_gate()saves aPendingActionto the DB and returns a sentinel token instead of executing. After the run completes,handle_messagescans for sentinels and edits the sent message to add the inline keyboard. The agent is instructed via_APPROVAL_INSTRUCTIONSto tell the user "action pending approval" when it sees a sentinel — never to show the raw token. - Edit state — at the very start of
handle_message, before intent routing, the handler checkscontext.user_dataforpending_edit_action_id. If set, the message is treated as a correction and routed to_handle_edit_correction()instead of the normal flow. - The reply is registered for feedback. After the answer is sent, a
system:msgtrace:<chat_id>:<message_id>row records the turn's trace id and a genre (chat, or the loop that sent it). A 👍/👎 reaction arriving minutes later — long after the trace closed — is scored against the right trace from that row. Telegram only deliversmessage_reactionto bots that name it inallowed_updates, so the bot polls withUpdate.ALL_TYPES. See Reaction feedback. - Token logging — each
telegram/messagespan recordsest_tokens_history,est_tokens_context,est_tokens_user_turn,context_layers, andhistory_interactionsas Logfire span attributes for per-request token observability.
Reaction Feedback¶
Not a message flow — an out-of-band one. A reaction arrives as its own update type, with no message text and often long after the turn it refers to.
sequenceDiagram
participant You
participant TG as Telegram
participant RH as handle_message_reaction
participant S as Storage (context KV)
participant LFS as Langfuse
You->>TG: React 👍 / 👎 to any Kwasi message
TG->>RH: message_reaction update
Note over TG,RH: Delivered only because the bot polls with<br/>allowed_updates=Update.ALL_TYPES — it is<br/>excluded from Telegram's default set.
RH->>RH: Allowlist check on reaction.user.id
RH->>S: read system:msgtrace:{chat}:{message}
alt Mapping found and emoji carries a verdict
RH-)LFS: score_trace(trace_id, "user_feedback", 1.0 / 0.0, BOOLEAN)
RH->>S: write the reaction back onto the row
else No mapping, or an ambiguous emoji
RH->>RH: Log and stop — nothing scored
end
Why it matters. user_approval only exists on turns that reach the approval
gate — a small minority. user_feedback covers everything else, at one tap and
no typing.
Ambiguous reactions score nothing. 🤔 on an answer could mean anything; recording it either way would be noise in the metric this exists to produce. Removing a reaction scores nothing either — un-reacting is not a verdict, and the earlier score stands rather than being silently contradicted.
Genre is the second signal, and it must stay stable. Every sent message records what produced it, so "which proactive genres earn any engagement" is one query — the measurement that decides which silenced background loops deserve reactivation.
A loop's label names this send ("meeting prep: Weekly Standup") because it
becomes the Langfuse span and session name, where per-event uniqueness is what
you want when opening one trace. Its genre names the loop. Conflating them
gave every meeting its own genre and made engagement unanswerable for the loop
that fires most; loops whose label varies now pass an explicit genre.
SELECT content::json->>'genre' AS genre,
count(*) AS sent,
count(content::json->>'reaction') AS reacted
FROM context WHERE user_id LIKE 'system:msgtrace:%' GROUP BY 1 ORDER BY 2 DESC;
Multi-step Plan Flow (Spec 008)¶
Before a Telegram text message hits the agent, a planning gate decides whether it should be split into ordered steps. Most messages are simple ("what's on my calendar?") and skip planning entirely. Complex messages ("check email, then post a summary in #standup, then add a follow-up task") get decomposed, previewed for approval, and executed step-by-step with live progress.
sequenceDiagram
participant You
participant TG as Telegram
participant H as handle_message
participant CC as classify_complexity
participant GP as generate_plan
participant S as Storage
participant EX as execute_plan
participant R as Intent Router (per step)
participant A as Domain Agent
You->>TG: "Check email, then post a summary in #standup, then add a task"
TG->>H: Update event
Note over H: Standard pre-checks (allowlist, edit-state, history, context injection)
H->>CC: classify_complexity(text)
CC-->>H: True (regex matched "then" connectives, len ≥ 30)
H->>GP: generate_plan(text, deps)
GP->>GP: _planner_agent.run(text)\n→ ExecutionPlan(goal, steps[], needs_planning)
alt needs_planning=False or <2 steps
GP-->>H: None (fall through to normal flow)
Note over H,A: Continue to standard text-message flow
else 2+ steps
GP-->>H: ExecutionPlan
H->>H: format_plan_preview(plan) → MarkdownV1 message
H->>S: save PendingAction(action_type="plan", payload=plan.json, trace_id)
H->>TG: Send preview + [✅ Confirm] [❌ Cancel] [✏️ Edit]
Note over You,TG: Wait for user
You->>TG: Tap ✅ Confirm
TG->>H: CallbackQuery (handle_callback_query loads plan)
H->>EX: execute_plan(plan, deps, send_progress)
loop For each step in plan.steps
EX->>TG: Edit progress message\n(✓ done · ▶ current · _ pending)
EX->>R: classify_intent(step_message + scratchpad)
R-->>EX: domain agent
EX->>A: agent.run(step_message_with_scratchpad)
alt Step succeeds
A-->>EX: output text
EX->>EX: scratchpad.append(step output, truncated 300 chars)
else Step throws
A-->>EX: Exception
EX->>S: save PendingAction(action_type="plan_resume",\npayload=PlanResumePayload(plan, remaining_steps, scratchpad,\nfailed_step, failure_reason))
EX->>TG: Send "Step N failed: {reason}\nContinue or abort?"\n+ inline keyboard
EX-->>H: partial summary (return)
end
end
EX->>TG: Final message — per-step outputs concatenated under "*Done!* {goal}"
H->>S: log_interaction
end
Notes on the plan flow¶
- Pre-filter is cheap.
classify_complexityis a single regex match on connectives (and then,also,after that) with a 30-char minimum — single-action messages cost zero LLM tokens for planning. The planner agent itself is tool-less; per-step execution uses the existing router + domain agents. - Scratchpad threading. Step N+1 receives a prepended block
"Context from previous steps:\n- Step N (description): output[:300]"so later steps build on earlier results without re-fetching. - The plan is always gated, even when none of its constituent tool calls would normally be gated.
action_type="plan"andtool_name="__plan__"distinguish plan actions from regular tool approvals; plan-resume actions use"plan_resume"/"__plan_resume__". - Resume preserves work. Confirm restarts
execute_planwithsteps_override=remaining_stepsandinitial_scratchpad=payload.scratchpad— completed steps are not re-run. Trace scoring carries through to plan and plan_resume actions viaPendingAction.trace_id(Spec 009).
Inline Approval (Button Callbacks)¶
When the user taps a button under a pending action preview, Telegram fires a CallbackQuery event handled by handle_callback_query.
sequenceDiagram
participant You
participant TG as Telegram
participant CB as handle_callback_query
participant S as Storage
participant CA as confirm_action (app/approval.py)
participant EX as executor (ACTION_REGISTRY)
participant TD as Microsoft To Do
participant LFS as Langfuse
You->>TG: Tap ✅ Confirm / ❌ Cancel / ✏️ Edit
TG->>CB: CallbackQuery (data="approve:{verb}:{uuid}")
CB->>CB: query.answer() — dismiss Telegram loading spinner
CB->>S: get_pending_action(uuid)
CB->>CB: Guard: status == "pending" and now < expires_at
alt Confirm
CB->>CA: confirm_action(action, deps)
CA->>EX: execute_approved_action(tool_name, payload, deps)
opt Task executor (create / edit / complete / delete)
EX->>TD: push to Microsoft To Do default list (best-effort)
Note over EX,TD: No-op unless todo_sync_active. A failure never<br/>breaks the local task op. Delete captures todo_task_id first.
end
EX-->>CA: result text
CA->>S: update_pending_action(status="confirmed")
CA->>S: log_audit_entry(channel="telegram")
CA-)LFS: score_trace(trace_id, "user_approval", 1.0)
CA-->>CB: result text
CB->>TG: edit_message_text("✅ Done!\n{result}")
Note over CA: The Mini App calls the same function,<br/>so the two surfaces cannot diverge.
else Cancel
CB->>CA: cancel_action(action, deps)
CA->>S: update_pending_action(status="cancelled")
CA-)LFS: score_trace(trace_id, "user_approval", 0.0)
CB->>TG: edit_message_text("❌ Cancelled.")
else Edit
CB->>S: user_data["pending_edit_action_id"] = uuid
CB->>TG: edit_message_text("What would you like to change?")
CB-)LFS: score_trace(trace_id, "user_edit", 1.0)
Note over You,TG: Next message from user → _handle_edit_correction()
else Already resolved or expired
CB->>TG: edit_message_text("This action was already {status}.")
end
_handle_edit_correction() — called when the user sends a correction after tapping Edit:
- Loads the original
PendingAction, verifies it's still pending and unexpired - Marks the old action
cancelled, incrementsedit_count - Re-runs the agent with a
[REVISION]prompt:"Original: {preview_text}\nCorrection: {user_text}\nPlease redo with this change." - The agent produces a new sentinel → new preview with a fresh keyboard
- Old
message_idis no longer relevant (previous message was edited to show "Cancelled")
To Do push rides the executor, not the loop. The four task executors mirror to the Microsoft To Do default list inline — best-effort, exactly like embeddings, so a Graph failure never breaks the local task write. It is a no-op unless todo_sync_active. The reverse direction (completions and app-created tasks) is pulled by the To Do Sync Loop.
Approval taps are trace scores. Confirm / Cancel write a user_approval score (1.0 / 0.0) and Edit writes user_edit, against the trace captured in PendingAction.trace_id at gate time — so a decision made minutes later still lands on the right Langfuse trace (Spec 009).
One resolution path, two surfaces. Execute / status / audit / score live in app.approval.confirm_action and cancel_action; the callback handler and the Mini App both call them, and each owns only its own presentation. The audit row records which surface approved (channel), so the log cannot be quietly wrong about where a consequential action was authorised. The Mini App additionally offers a payload editor — the card here truncates code at 400 characters, so what you approve in the chat is not always what you can read.
Non-Telegram interfaces (CLI, WhatsApp): approval_gate() detects interface != "telegram" and executes the action immediately without saving a PendingAction. The approval UI is Telegram-only.
Voice Message¶
Voice goes through an extra STT step before hitting the agent. No streaming — the full transcript is needed first.
sequenceDiagram
participant You
participant TG as Telegram
participant H as handle_voice_message
participant G as Gemini STT
participant S as Storage
participant A as Agent
You->>TG: Send voice note
TG->>H: Update event (voice/audio)
H->>H: Check allowlist
H->>TG: Download audio bytes
TG-->>H: audio_bytes
H->>G: transcribe_audio(audio_bytes, model_id)
G-->>H: transcript text
H->>S: get_interactions_by_user(limit=5)
S-->>H: message_history
H->>A: agent.run(transcript, deps, message_history)
A-->>H: response_text
alt response ≤ 500 words
H->>H: synthesize_speech(response_text) via edge-tts
H->>TG: send_voice(audio_bytes)
else response > 500 words
H->>TG: send_message(response_text)
end
H->>S: log_interaction(transcript, response)
The model ID for Gemini STT is derived from MODEL_NAME by stripping the provider prefix (e.g. google-gla:gemini-2.5-flash → gemini-2.5-flash). TTS uses edge-tts with the voice configured via TTS_VOICE (default en-GB-RyanNeural). WhatsApp voice messages receive a text reply only — no TTS.
Photo Message¶
Photos are analyzed by Gemini Vision. The analysis result becomes the user message fed to the agent.
sequenceDiagram
participant You
participant TG as Telegram
participant H as handle_photo_message
participant GV as Gemini Vision
participant A as Agent
participant S as Storage
You->>TG: Send photo (+ optional caption)
TG->>H: Update event (photo)
H->>H: Check allowlist
H->>TG: Download highest-resolution photo
TG-->>H: image_bytes
H->>GV: analyze_image(image_bytes, "image/jpeg", caption)
GV-->>H: image analysis text
H->>H: Build user_message:\n"[Image shared]\nQuestion: {caption}\nAnalysis: {analysis}"
H->>S: get_interactions_by_user(limit=5)
S-->>H: message_history
H->>A: agent.run(user_message, deps, message_history)
A-->>H: response_text
H->>TG: Reply with response
H->>S: log_interaction(user_message, response)
Document Message (PDF / Text)¶
PDFs go through Gemini Vision. Text files are decoded directly and passed to the agent as-is.
sequenceDiagram
participant You
participant TG as Telegram
participant H as handle_document_message
participant GV as Gemini Vision
participant A as Agent
participant S as Storage
You->>TG: Send document (+ optional caption)
TG->>H: Update event (document)
H->>H: Check allowlist
H->>H: Check mime_type (PDF / text / other)
alt Unsupported type
H->>TG: "I can analyze PDFs and text files only"
else Text file (text/*)
H->>TG: Download file bytes
TG-->>H: file_bytes
H->>H: Decode UTF-8, truncate at 50k chars
H->>A: agent.run("[Document: name]\nContents:\n{text}", ...)
else PDF (application/pdf)
H->>TG: Download file bytes
TG-->>H: file_bytes
H->>GV: analyze_image(file_bytes, "application/pdf", caption)
GV-->>H: document analysis
H->>A: agent.run("[PDF: name]\nAnalysis:\n{analysis}", ...)
end
A-->>H: response_text
H->>TG: Reply with response
H->>S: log_interaction(user_message, response)
Slash Commands¶
Slash commands (/notes, /tasks, /reminders) bypass the agent entirely and query storage directly. They're fast, deterministic, and don't consume LLM tokens.
sequenceDiagram
participant You
participant TG as Telegram
participant H as Command Handler
participant S as Storage
You->>TG: /tasks
TG->>H: CommandHandler("tasks")
H->>H: Check allowlist
H->>S: list_tasks(status="todo")
S-->>H: list[Task]
H->>TG: Formatted task list (Markdown)
| Command | Storage call | Returns |
|---|---|---|
/notes |
get_notes() |
5 most recent notes |
/tasks |
list_tasks(status="todo") |
All pending tasks |
/reminders |
list_reminders() |
All pending reminders |
/readlater |
delegates to handle_message |
Runs list_read_later via agent |
/briefing |
delegates to handle_message |
Runs the morning briefing prompt on demand |
/start, /help |
none | Static help text |
Voice Reply Trigger (Text Messages)¶
Text messages that contain specific phrases cause Kwasi to reply with a voice note instead of text, provided the response is ≤500 words. This works even though the input is typed text, not a voice message.
Trigger phrases (case-insensitive, word-boundary matched):
| Phrase pattern | Examples |
|---|---|
tell me |
"tell me the weather today" |
read me / read that / read this / read it (out) |
"read me that article", "read it out" |
say that / say it |
"say that again" |
speak (to me) |
"speak to me", "speak" |
flowchart TD
A[Text message arrives] --> B{_VOICE_TRIGGER_RE\nmatches phrase?}
B -- No --> C[Normal text reply\n_to_telegram_md + edit flow]
B -- Yes --> D{Response ≤ 500 words\nAND TTS_VOICE set?}
D -- No --> C
D -- Yes --> E[synthesize_speech via edge-tts]
E --> F{Audio generated?}
F -- No, exception --> C
F -- Yes --> G[Delete placeholder message]
G --> H[reply_voice with audio_bytes]
If TTS synthesis fails for any reason, it falls back silently to the normal text response. The voice reply trigger only applies to Telegram text messages — voice input, WhatsApp, and CLI are unaffected.
External API Message (POST /message)¶
For Android HTTP Shortcuts and other external clients. Requires API_TOKEN and BRIEFING_CHAT_ID.
flowchart TD
A[POST /message] --> B{X-API-Token\nvalid?}
B -- No --> B1[401 Unauthorized]
B -- Yes --> C{Body type?}
C -- multipart/form-data --> D[Read text field\nRead file field if present]
C -- application/json --> E[Parse text\nDecode image_b64 if present]
D --> F{Image provided?}
E --> F
F -- Yes --> G[analyze_image via Gemini Vision\nResult prepended to user_message]
F -- No --> H[user_message = text]
G --> H
H --> I[Load last 10 interactions\nusing briefing_chat_id as user_id]
I --> J[classify_intent → select_agent]
J --> K[Inject read-later context]
K --> L[agent.run with message_history]
L --> M[log_interaction\nuser_id = briefing_chat_id]
M --> N[Send response to Telegram\nsplit at 4096 chars]
N --> O[Return JSON\n{response, status: ok}]
Key design decisions:
- user_id is set to BRIEFING_CHAT_ID — the same ID used by the Telegram bot. This means Android shares and Telegram messages share the same conversation history and context window.
- Images are analysed via Gemini Vision before being passed to the agent — the agent sees the analysis text, not raw bytes.
- The response is always delivered to Telegram. The JSON response body also contains the full response text for use in HTTP Shortcuts toasts or Tasker variables.
- The approval gate is bypassed for consequential actions — interface is set to "api", which is not "telegram", so approval-gated tools execute immediately (same as CLI and briefing contexts).