Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
56 commits
Select commit Hold shift + click to select a range
f41088f
feat(chat): land shared-agent chat bridge and supporting runtime infr…
yordis Jul 28, 2026
4f77567
feat(chat): let a chat surface reset its agent session
yordis Aug 2, 2026
ac5dbff
refactor(channel): name the crate after the vocabulary it already uses
yordis Aug 2, 2026
876e58e
refactor(channel): finish the vocabulary move through the deployment …
yordis Aug 2, 2026
5e0600f
refactor(channel): mint conversation ids the way the rest of the work…
yordis Aug 2, 2026
a4adc05
fix(repo): restore the ignore rules a wholesale block rewrite dropped
yordis Aug 2, 2026
9e8a548
docs(adr): record where inbound media redemption belongs
yordis Aug 2, 2026
cad6254
docs(channel): define the endpoint vocabulary before the topology tha…
yordis Aug 2, 2026
ef6c9ec
docs(channel): write down that the chat, not the speaker, is the unit…
yordis Aug 3, 2026
ae36482
docs(channel): say whose perspective the subject direction token is from
yordis Aug 3, 2026
5741837
docs(channel): describe one architecture instead of a phase plan
yordis Aug 3, 2026
d12b521
docs(channel): correct the rendering claims and record the claim-chec…
yordis Aug 4, 2026
7a32ca3
fix(channel): stop losing Telegram updates the gateway had to offload
yordis Aug 4, 2026
194a1b5
chore(repo): merge main into telegram
yordis Aug 4, 2026
f05f5d0
fix(channel): stop a rejected prompt from costing a conversation its …
yordis Aug 4, 2026
21a7fe9
fix(channel): reject a blank bot token instead of booting on it
yordis Aug 4, 2026
b86f1d7
fix(docs): link the one ADR reference that CI rejects
yordis Aug 4, 2026
a491fd5
fix(deps): keep teloxide off native-tls
yordis Aug 4, 2026
26e89d6
chore(channel): declare the license both new crates were missing
yordis Aug 4, 2026
98b42d4
fix(channel): satisfy the repo Rust policy lints
yordis Aug 4, 2026
20858cf
fix(channel): keep chunked replies inside the limit Telegram actually…
yordis Aug 4, 2026
157bbb8
fix(devops): let a deployment actually ask for no session triggers
yordis Aug 4, 2026
53df691
fix(channel): let the coverage build compile the Telegram bridge
yordis Aug 4, 2026
06b6ebe
fix(channel): stop a wrong lost-session guess from orphaning a live s…
yordis Aug 4, 2026
4f5d80b
fix(channel): harden domain validation and prompt-failure cleanup
yordis Aug 4, 2026
14c4a2c
fix(repo): clear the lint and coverage gates the branch was failing
yordis Aug 4, 2026
70729c4
fix(channel): stop a failed pointer write from stranding a live session
yordis Aug 4, 2026
122764e
fix(channel): let a trigger outside ASCII be reachable at all
yordis Aug 4, 2026
5842fce
refactor(channel): say why a media type was rejected rather than echo it
yordis Aug 4, 2026
953b14a
fix(channel): keep a media type's parameters out of its normalization
yordis Aug 4, 2026
5a60b78
fix(channel): stop a second conversation from burying an endpoint's f…
yordis Aug 4, 2026
8d29d07
fix(nats): stop a claim's bucket header from naming a bucket it was n…
yordis Aug 4, 2026
9521b28
refactor(nats): drop the claim publisher accessor nothing reads
yordis Aug 4, 2026
bcc0b30
refactor(channel): keep what an endpoint's binding leads to in one place
yordis Aug 4, 2026
f6d745a
test(channel): cover the records a contested claim can still leave be…
yordis Aug 4, 2026
6858891
fix(channel): read a media type spelled the way the standard writes it
yordis Aug 5, 2026
ec0c8d0
chore(devops): hold the Telegram bridge's container surface
yordis Aug 5, 2026
bc047fd
fix(repo): ignore a Compose override again, under the name Compose lo…
yordis Aug 5, 2026
3644299
fix(channel): scope a media handle to the credential that can redeem it
yordis Aug 5, 2026
7e1f69c
refactor(channel): let a boot failure name the variable an operator m…
yordis Aug 5, 2026
db3be20
fix(channel): settle a contested endpoint claim by reading the claim …
yordis Aug 5, 2026
59c7f9b
refactor(nats): tell a foreign bucket apart from a name no bucket cou…
yordis Aug 5, 2026
d39f1a2
fix(channel): stop orphaning a session the agent has already opened
yordis Aug 5, 2026
05e335a
fix(devops): give the Compose service the image it says it builds
yordis Aug 5, 2026
1071adb
refactor(channel): declare test modules without a path attribute
yordis Aug 5, 2026
e14f657
chore(devops): keep the Telegram bridge out of the container surface
yordis Aug 5, 2026
b09d4bd
refactor(nats): let a test carry what a compile-time check was carrying
yordis Aug 5, 2026
3a9aebb
docs(nats): say what an unnamable bucket header actually reports
yordis Aug 5, 2026
c47faca
docs(adr): stop promising a terminal record a failed write can withhold
yordis Aug 5, 2026
75b2072
refactor(channel): reach the claim read through a single unwind
yordis Aug 5, 2026
7d86940
test(channel): pin the length a session id may not exceed
yordis Aug 5, 2026
3d47759
test(channel): assert the whole line an operator has to read
yordis Aug 5, 2026
55f96ba
docs(adr): stop promising a readiness state the bucket never writes
yordis Aug 5, 2026
2c82b0f
fix(channel): refuse at boot an account no endpoint could carry
yordis Aug 5, 2026
188716d
fix(channel): hand back a session the conversation never recorded
yordis Aug 5, 2026
15fc0d2
fix(channel): keep the words a photo arrived with
yordis Aug 5, 2026
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
9 changes: 9 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -26,18 +26,24 @@ docs/.vitepress/cache
mise.local.toml

# Docker
compose.override.yml
compose.override.yaml
docker-compose.override.yml

# IDE
.idea/
.vscode/
*.swp
*.swo
*~

# OS
.DS_Store
Thumbs.db

# Logs
*.log

# trogonai internal
.trogonai/
*.internal.trogonai.md
Expand All @@ -50,6 +56,9 @@ coverage-*.xml
*.profraw
*.profdata

# NATS
*.creds

# Misc
tmp

Expand Down
4 changes: 4 additions & 0 deletions docs/.vitepress/config.mts
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,10 @@ export default async () => {
{ text: "Key Custody", link: "/architecture/key-custody" },
{ text: "Key Management", link: "/architecture/key-management" },
{ text: "Key States", link: "/architecture/key-states" },
{
text: "Multi-Channel Agent Routing",
link: "/architecture/multi-channel-agent-routing",
},
{
text: "OpenTelemetry Transport Context",
link: "/architecture/opentelemetry-transport-context",
Expand Down
1 change: 1 addition & 0 deletions docs/.vitepress/helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ const GLOSSARY_SECTION_ORDER = [
"Protocols and transports",
"Event sourcing and the decider",
"Agent execution model",
"Channels and conversations",
"Messaging and storage infrastructure",
"WebAssembly execution",
"Wire contracts and serialization",
Expand Down
254 changes: 254 additions & 0 deletions docs/adr/0044-inbound-media-fetch-out-of-band.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,254 @@
---
number: "0044"
slug: inbound-media-fetch-out-of-band
status: draft
date: 2026-08-02
---

# ADR#0044: Inbound Media Is Fetched Out of Band by a Dedicated Consumer

## Context

A chat platform does not deliver media inline. It delivers an opaque handle
(a Telegram `file_id`, a Slack `file.url_private`, a Discord attachment URL)
that only the holder of the platform credential can redeem. Something in this
system has to do that redemption, and where it happens determines which
processes need platform credentials, which components stop being generic, and
what a conversational turn has to wait for.

[The multi-channel routing design](../architecture/multi-channel-agent-routing.md)
originally recorded "eager claim-check": the bridge downloads media at
normalize time, before dispatching the prompt. That decision was never
implemented. `channel-bridge-telegram`'s parser carries no media: `parse.rs`
takes the words a message came with, from `text` or from a media `caption`, and
hardcodes `attachments: Vec::new()`, so a photo reaches the agent as whatever was
said about it and never as a handle. Nothing is sunk and the decision is open.

Three properties of the problem constrain the answer:

1. **Redemption needs the platform credential.** Whoever fetches holds the bot
token. This is the whole reason a handle cannot simply be passed downstream
to a credential-free consumer.
2. **Redemption is slow and optional.** A turn that is pure text, or where the
agent never opens the attached file, pays nothing for media it does not
read. Fetching before dispatch makes every turn pay for the worst case.
3. **Handles are durable; redemption URLs are not.** A Telegram `file_id` is
permanent and redeemable by the token holder indefinitely. What expires
(roughly one hour) is the `file_path` that `getFile` returns. The original
design rejected lazy fetch partly on "platform download URLs expire," which
is true of the second step only and does not rule out deferring the first.

Three placements were considered.

**In the gateway.** `trogon-gateway` already holds a `TelegramBotToken` for
webhook registration, and the claim-check machinery with `ObjectStorePut` and
`ObjectStoreGet` already exists generically in `trogon-nats`. Both pieces are
in place, so the objection is not capability. The objection is scope: the
gateway is a verbatim transport shared across GitHub, GitLab, Linear, Slack,
Discord, and Telegram sources, and its stated contract is raw fidelity. Media
fetching would make one source materially smarter than its siblings, and would
put the bot token in the ingress. It also cannot be done in the request path at
all: Telegram retries webhooks that
do not return promptly, so a download inside the handler trades a fast ingress
for a slow one.

**In the bridge, before dispatch.** The original decision. It concentrates the
credential in a component that already has it and needs no coordination, but it
makes property 2 impossible: every turn waits for media the agent may never
read, on a path that is already fully serial.

**In a dedicated consumer.** JetStream already supports multiple independent
durables on one stream, so a second consumer of the raw stream costs no new
transport and no gateway change.

## Decision

### 1. The gateway stays a verbatim transport

`trogon-gateway` does not fetch media and does not gain per-source
intelligence. Its contract remains raw fidelity from webhook to stream. This is
a deliberate purchase: we accept a third credential holder (below) to keep the
ingress generic.

The gateway does already depend on an object store, through
`ClaimCheckPublisher`, which offloads any body over the NATS max payload and
publishes claim headers in its place. That is transport plumbing applied
identically to every source, not knowledge of what Telegram media is, so it does
not weaken this decision. It does mean any consumer of a raw stream must redeem
the claim before deserializing, the downloader below included. A consumer that
skips it gets no error, only an empty body.

### 2. A dedicated downloader consumes the raw stream on its own durable

A per-platform downloader (`channel-downloader-telegram` first) takes its own
durable consumer on the same raw stream the bridge reads, redeems the platform
handle with the bot token, and writes bytes through the existing
`ObjectStorePut` in `trogon-nats`. It is size-capped, and the cap is its own
configuration rather than the bridge's. For Telegram that cap is bounded from
above by the platform: the public Bot API refuses `getFile` for anything over
20 MB, so a larger configured cap only means something against a self-hosted Bot
API server, which is what lifts the limit.

Redemption is gated by the same identity check the bridge applies. The
downloader resolves the message's endpoint against `channel_endpoints_{prefix}`
and drops the update when that endpoint names no principal, before it calls
`getFile` and before anything reaches `ObjectStorePut`. Reading the raw stream
directly buys independence from the bridge, not exemption from its access
control. Without this check an unlinked chat could spend our bot token and our
object store by sending a file, and the bridge's later decision to drop the
message would arrive after the work was already done.

This is deliberately not a request/reply service. It is driven by the same
stream as the bridge, so a fetch begins at ingestion whether or not any agent
ever asks for the file, and the freshness window for `getFile` is entered
immediately rather than at some later moment of the agent's choosing.
Comment thread
coderabbitai[bot] marked this conversation as resolved.

### 3. Readiness is a KV status record, not an object-store lookup

An object store cannot answer "not yet." A `get` on a missing key returns
not-found, which is indistinguishable from a permanent failure and from a dead
downloader. Readiness therefore lives in its own JetStream KV bucket, keyed by
the credential that received the handle together with the handle itself:

```text
channel_media_{prefix}:
{channel}.{account}.{platform_ref} ->
{ state: ready | failed, object_ref, mime, size, error }
```
Comment thread
coderabbitai[bot] marked this conversation as resolved.

The credential belongs in the key because a handle is not globally meaningful. A
Telegram `file_id` is issued per bot and redeemable only by the token that
received it, so the same string seen by two of our bot accounts denotes two
different files, and the bucket is channel-neutral besides. A bare handle key
would let one account's record answer for another's, and would point redemption
at a credential that cannot honor it. Readers already have both tokens: they are
the leading part of the endpoint on the event the attachment arrived with, and
they are what selects the token used to redeem.

The bucket holds outcomes only. Nothing is written when a handle is first seen,
so there is no `pending` state and no record whose job is to say that work
started: a `pending` write would be a second thing the downloader must do before
it may fail, and it would still be missing in exactly the case a reader has to
survive, which is a downloader that never ran.

Absence is therefore the unresolved state, and it says nothing about whether work
is under way. The record is absent before the downloader's durable has reached the
message, while a download is running, and for as long as the downloader is down
or behind. A reader cannot tell those apart and does not need to. What it needs
is that absence is never permanent by accident, and two rules give it that.

The downloader writes a `failed` record for every permanent error, and also on
its last delivery attempt, which it recognizes from the delivery count JetStream
puts on the message. The write comes before the acknowledgement, and the
acknowledgement is what the write earns: a downloader that cannot reach the KV
bucket leaves the message unacknowledged, so JetStream redelivers and the
terminal record is attempted again on the next delivery. A handle whose
deliveries are exhausted therefore ends as a terminal record rather than as
silence, since a consumer that has stopped redelivering will never speak again on
its own.

Redelivery is a bounded number of attempts, so ordering the write before the ack
narrows the window in which a handle ends absent without closing it. A KV bucket
unreachable for the whole life of a message, or a downloader that dies between
its final read and its final write, exhausts the deliveries with nothing written.
That is the same shape as a downloader that never runs at all, and it has the
same backstop: the reader's deadline, not the terminal record, is what bounds how
long absence can last. The terminal record is what turns the failures a
downloader survives into an explanation the agent can be given instead of a
timeout, and it is not load-bearing for liveness.

Readers await readiness with a KV watch and a deadline, not a poll. Finding
nothing is where a reader starts, not a failure it reports: it watches for the
first record the key ever gets and stops on its own deadline. A late reader
observes current state directly with no replay concern, and a deadline expiry is
reported to the agent as an unavailable attachment rather than as a turn failure.

### 4. The inbound event carries the handle, never the object reference

`Attachment` on the inbound event carries `platform_ref` and drops
`object_ref`. The event states that a photo exists and gives its handle;
resolving that handle to bytes is a lookup performed later, not a field
populated earlier. An event that carried `object_ref` would be asserting the
presence of bytes that may not exist yet.

Outbound is not symmetric and does not change. `RenderCommand::SendAttachment`
keeps its `object_ref` because the agent produced that file and already put it
in the object store; there is no handle to redeem and nothing to wait for.

### 5. The turn does not block on media; the agent's tool does

The bridge builds the inbound event and dispatches the prompt without waiting.
Waiting happens inside the agent-facing download tool, at the moment the agent
actually opens the file, and it is the reader described above: absence means keep
waiting until the deadline, `failed` is an explanation to hand the agent, and
`ready` is the object reference. Text-only turns and turns that ignore an
attachment pay nothing.

## Invariants

- The gateway never interprets a source's payload beyond what publishing
requires. Its object-store use is claim-check transport, identical for every
source.
- No component blocks a conversational turn on media the agent has not asked
for.
- Any component that redeems a platform handle holds that platform's
credential; no credential-free component is ever handed a handle it is
expected to resolve.
- A handle is redeemed with the credential of the account that received it, and
is never keyed or cached in a way that lets one account's handle be resolved
by another's.
- No handle is redeemed for an endpoint that resolves to no principal.
Authorization precedes credential use, in every component that holds a
credential.
- The readiness bucket records outcomes and nothing else: `ready` and `failed`
are the only states ever written, and absence is the unresolved state. It is
never read as an assertion that a download is running, has not started, or ever
will. "Bytes absent from the object store" is not a lifecycle signal either,
because readiness is asked of the bucket and never of the store.
- Every handle a downloader stops working on leaves a terminal record, written
before the message is acknowledged. Giving up is written down, not expressed by
falling silent.
- No wait for readiness depends on a record arriving. A reader's deadline bounds
absence on its own, so a downloader that cannot write its terminal record
degrades the explanation an agent receives, never the reader's ability to stop
Comment thread
coderabbitai[bot] marked this conversation as resolved.
waiting.
- An inbound event never asserts the existence of bytes that have not been
written.

## Consequences

- **The bot token lives in three processes**: the gateway (webhook
registration), the bridge (Bot API sends), and the downloader (`getFile`).
This is the direct cost of keeping the gateway generic, and it is accepted.
The evolution path already recorded in the routing design, a generic
gateway **sink** concept that would centralize outbound token custody, is
the eventual consolidation and remains out of scope here.
- **A third worker appears, but only when media does.** The routing design's
"two workers total" claim holds until the first platform that carries media
is supported. Nothing needs to be built before then.
- **A new KV bucket** (`channel_media_{prefix}`) joins the four the channel
store already provisions.
- **Failure is legible.** A download that fails permanently is a `failed`
record with a reason, distinguishable from one still in flight, so an agent
can be told the difference. The cost is that the downloader has to write that
record on the way out, including on its final delivery attempt and before it
acknowledges, rather than letting the consumer's own give-up be the ending.
Redelivery covers a write that fails while deliveries remain, and the reader's
deadline covers the rest, which is why readers keep a deadline instead of
trusting that a record always arrives.
- **The endpoints bucket gains a second reader.** `channel_endpoints_{prefix}`
stays the bridge's to write, but the downloader reads it to authorize before
redeeming, so identity is one registry consulted by every component that acts
on a message rather than a check the bridge performs on everyone's behalf.
- **The downloader can be restarted or backfilled independently.** Because it
is a durable consumer of a retained raw stream rather than a request/reply
service, a downloader that was down comes back and works through what it
missed without the bridge participating.
- **Two consumers now read the same raw stream.** This is ordinary JetStream
usage, but it does mean the bridge no longer has exclusive knowledge of what
arrived, and the two consumers' positions can differ.

## References

- [Multi-Channel Agent Routing](../architecture/multi-channel-agent-routing.md)
- [ADR#0024: Agent Platform Stream Topology](./0024-agent-platform-stream-topology.md)
1 change: 1 addition & 0 deletions docs/adr/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,3 +49,4 @@ future implementation work.
- [ADR#0041: Canonical MCP JSON-RPC Bodies over NATS (Draft)](./0041-canonical-mcp-jsonrpc-bodies-over-nats.md)
- [ADR#0042: NATS Trace Context and Message Path Tracing (Draft)](./0042-nats-trace-context-and-message-path-tracing.md)
- [ADR#0043: Agent Instructions Ownership and Shape (Draft)](./0043-agent-instructions-ownership-and-shape.md)
- [ADR#0044: Inbound Media Is Fetched Out of Band by a Dedicated Consumer (Draft)](./0044-inbound-media-fetch-out-of-band.md)
Loading
Loading