feat(EVO-2180): forward incoming media to the AI Processor as A2A file parts - #5
Conversation
…e parts
The bot only ever sent a text part, so images/audio the customer sent never
reached the AI (the agent replied "No content to process"). Accept attachments on
the inbound event, carry them through the debounce window, download each and send
it as a base64 A2A file part. Part of EVO-2178 (image end-to-end).
- pipeline/model: MessageEvent.Attachments + Attachment{URL,ContentType,FileType}.
- pipeline/repository: AppendAttachments/GetAttachments on a parallel Redis list
(bot_runtime:attach:{contact}:{conv}), aggregated like the text buffer and cleared
together in ClearState (no stale media leaks into the next turn).
- debounce/service: Start/Reset accept attachments; GetAttachments added.
- pipeline/service: thread event.Attachments through start/skip/reset/advance -> the
A2ARequest (read fresh from Redis at stage launch, like the buffer).
- ai/model: A2ARequest.Attachments; JSONRPCPart.File + JSONRPCFile{Name,MimeType,Bytes}
(tags match the processor's extract_files_from_message).
- ai/service/ai_adapter: download each attachment once (before Marshal, reused across
retries) with a 15 MiB cap; base64-encode; append a file part. A download failure is
logged and skipped so the text-only message always survives.
- tests: adapter forwards a file part with decodable base64 + download-failure sends
text only; repo AppendAttachments/GetAttachments roundtrip + ClearState clears the
attach key. Full suite green (go build/vet/test ./pkg/... ./internal/...).
Note: test/e2e was already incompatible with NewAIAdapter on develop (pre-existing);
repo CI is docker-only.
Reviewer's GuideAdds end-to-end support for forwarding incoming media (attachments) through the debounce/pipeline into the AI adapter, which downloads them once and sends them as base64 A2A file parts, with Redis-backed attachment buffering and clear failure isolation so text-only flow still works. Sequence diagram for forwarding incoming media as A2A file partssequenceDiagram
actor CRM
participant PipelineService
participant DebounceEngine
participant PipelineRepository
participant Redis
participant AIAdapter
participant AttachmentServer
participant AIProcessor
CRM ->> PipelineService: Process(MessageEvent{MessageContent, Attachments})
PipelineService ->> DebounceEngine: Start(ctx, contactID, conversationID, content, attachments, cfg)
DebounceEngine ->> PipelineRepository: AppendToBuffer(ctx, contactID, conversationID, content)
DebounceEngine ->> PipelineRepository: AppendAttachments(ctx, contactID, conversationID, attachments)
PipelineRepository ->> Redis: RPush bot_runtime:buffer / bot_runtime:attach
note over DebounceEngine,PipelineService: Later, when debounce is skipped/expired
PipelineService ->> DebounceEngine: GetBuffer(ctx, contactID, conversationID)
DebounceEngine -->> PipelineService: buffer
PipelineService ->> DebounceEngine: GetAttachments(ctx, contactID, conversationID)
DebounceEngine -->> PipelineService: attachments
PipelineService ->> AIAdapter: Call(ctx, &A2ARequest{Message: buffer, Attachments: attachments})
AIAdapter ->> AIAdapter: buildFileParts(ctx, req)
loop for each Attachment
AIAdapter ->> AIAdapter: downloadAttachment(ctx, att.URL)
AIAdapter ->> AttachmentServer: HTTP GET att.URL (maxAttachmentBytes)
AttachmentServer -->> AIAdapter: data or error
alt success
AIAdapter ->> AIAdapter: base64.StdEncoding.EncodeToString(data)
AIAdapter ->> AIAdapter: append JSONRPCPart{Type:file, File:JSONRPCFile}
else failure
AIAdapter ->> AIAdapter: slog.Warn("pipeline.ai.attachment.download_failed")
end
end
AIAdapter ->> AIProcessor: HTTP POST JSONRPCRequest{Message.Parts: [text + file parts]}
AIProcessor -->> AIAdapter: JSON-RPC response
AIAdapter -->> PipelineService: NormalizedResponse
File-Level Changes
Tips and commandsInteracting with Sourcery
Customizing Your ExperienceAccess your dashboard to:
Getting Help
|
There was a problem hiding this comment.
Hey - I've found 3 issues, and left some high level feedback:
- In
advanceToAIandskipDebounce, a failure inGetAttachmentscurrently aborts launching the AI stage; consider treating attachment retrieval errors like download failures (log and proceed with text-only) so media issues do not block the entire turn. - In
downloadAttachment, the per-download timeout is tied to the adapter's general timeout and downloads are done sequentially, so a message with several slow attachments can exceed the intended overall latency; you may want a shorter per-attachment timeout or an overall cap for all media fetches. - The
get_attachments.skip_malformedwarning only logs the error; addingcontact_idandconversation_idto this log would make debugging malformed attachment entries in Redis significantly easier.
Prompt for AI Agents
Please address the comments from this code review:
## Overall Comments
- In `advanceToAI` and `skipDebounce`, a failure in `GetAttachments` currently aborts launching the AI stage; consider treating attachment retrieval errors like download failures (log and proceed with text-only) so media issues do not block the entire turn.
- In `downloadAttachment`, the per-download timeout is tied to the adapter's general timeout and downloads are done sequentially, so a message with several slow attachments can exceed the intended overall latency; you may want a shorter per-attachment timeout or an overall cap for all media fetches.
- The `get_attachments.skip_malformed` warning only logs the error; adding `contact_id` and `conversation_id` to this log would make debugging malformed attachment entries in Redis significantly easier.
## Individual Comments
### Comment 1
<location path="pkg/ai/service/ai_adapter.go" line_range="28-32" />
<code_context>
+// maxAttachmentBytes caps a single downloaded incoming attachment (EVO-2180).
+// Images routinely exceed maxResponseBytes (1 MiB); base64 inflates ~33%, so keep
+// this conservative relative to the processor's request-body limit.
+const maxAttachmentBytes = 15 << 20 // 15 MiB
+
// maxBackoff caps the exponential backoff between retries so a large
</code_context>
<issue_to_address>
**suggestion (performance):** Unbounded number of attachments could cause high memory usage despite per-attachment cap
The 15 MiB per-attachment cap doesn’t prevent `buildFileParts` from loading an unbounded number of attachments into memory for a single request, which could lead to excessive memory usage for media-heavy conversations. Please consider enforcing a global limit (e.g., max attachment count or total bytes) or aborting once a threshold is reached so overall memory usage stays bounded.
Suggested implementation:
```golang
// maxAttachmentBytes caps a single downloaded incoming attachment (EVO-2180).
// Images routinely exceed maxResponseBytes (1 MiB); base64 inflates ~33%, so keep
// this conservative relative to the processor's request-body limit.
const maxAttachmentBytes = 15 << 20 // 15 MiB
// maxTotalAttachmentBytes caps the aggregate size of all downloaded attachments
// for a single request to keep overall memory usage bounded even for media-heavy
// conversations.
const maxTotalAttachmentBytes = 50 << 20 // 50 MiB
// maxAttachmentCount caps the number of attachments processed per request to
// prevent unbounded growth in memory usage due to many small files.
const maxAttachmentCount = 16
```
```golang
// Message parts: text first, then one file part per downloaded attachment.
// EVO-2180: downloads happen ONCE here (before Marshal), so the byte-identical
// body is reused across retries.
parts := []model.JSONRPCPart{{Type: "text", Text: req.Message}}
fileParts, err := a.buildFileParts(ctx, req, maxAttachmentBytes, maxTotalAttachmentBytes, maxAttachmentCount)
if err != nil {
return nil, brtErrors.Wrap(err, "buildFileParts")
}
parts = append(parts, fileParts...)
```
To fully enforce bounded memory usage for attachments, you’ll need to update the implementation of `buildFileParts` in `pkg/ai/service/ai_adapter.go` (or wherever it’s defined):
1. **Update the signature** to accept the new limits and return an error:
- Change from something like:
```go
func (a *Adapter) buildFileParts(ctx context.Context, req *model.AIRequest) []model.JSONRPCPart
```
to:
```go
func (a *Adapter) buildFileParts(
ctx context.Context,
req *model.AIRequest,
maxAttachmentBytes int,
maxTotalAttachmentBytes int,
maxAttachmentCount int,
) ([]model.JSONRPCPart, error)
```
2. **Enforce per-request limits inside `buildFileParts`:**
- Track `totalBytes := 0` and `count := 0`.
- For each attachment:
- Abort with an error if `count+1 > maxAttachmentCount`.
- After downloading/reading an attachment into memory (or knowing its size), reject it if `len(data) > maxAttachmentBytes`.
- Before appending, check `totalBytes+len(data) <= maxTotalAttachmentBytes`; if exceeded, return an error (or stop processing further attachments).
- Accumulate `totalBytes += len(data)` and `count++` for accepted attachments.
3. **Propagate the error**:
- Ensure all call sites of `buildFileParts` are updated to handle the `(parts, err)` return, similar to the pattern in the edited code above, and that the error is surfaced appropriately to the caller (e.g., as a user-facing error or logged and wrapped).
</issue_to_address>
### Comment 2
<location path="pkg/ai/service/ai_adapter.go" line_range="349-358" />
<code_context>
+func (a *aiAdapter) downloadAttachment(ctx context.Context, url string) ([]byte, error) {
</code_context>
<issue_to_address>
**suggestion (performance):** Per-attachment timeout tied to adapter timeout may allow long overall delays when multiple media items are present
Each attachment download uses `time.Duration(a.timeoutSecs) * time.Second`, the full adapter timeout. With multiple attachments, total latency can scale as `N * timeoutSecs`. Consider a shorter per-media timeout or a shared deadline for all downloads in a request to prevent excessive end-to-end delays when many attachments are present.
</issue_to_address>
### Comment 3
<location path="pkg/pipeline/repository/redis_pipeline_repository.go" line_range="95-70" />
<code_context>
+ return nil, fmt.Errorf("pipeline.repository.get_attachments: %w", err)
+ }
+ atts := make([]model.Attachment, 0, len(raw))
+ for _, s := range raw {
+ var att model.Attachment
+ if err := json.Unmarshal([]byte(s), &att); err != nil {
+ slog.Warn("pipeline.repository.get_attachments.skip_malformed", "error", err)
+ continue
+ }
+ atts = append(atts, att)
+ }
+ return atts, nil
+}
+
func (r *redisPipelineRepository) SetTimer(ctx context.Context, contactID, conversationID int64, ttl time.Duration) error {
return r.rdb.Set(ctx, timerKey(contactID, conversationID), "1", ttl).Err()
</code_context>
<issue_to_address>
**suggestion:** Attachment deserialization logging lacks contact/conversation context, which may hinder debugging
In `GetAttachments`, malformed entries are logged without `contact_id` or `conversation_id`, making it difficult to tie warnings to specific conversations in production. Please include these IDs in the `slog.Warn` fields so bad data can be traced and cleaned up more easily, without changing behavior.
Suggested implementation:
```golang
for _, s := range raw {
var att model.Attachment
if err := json.Unmarshal([]byte(s), &att); err != nil {
slog.Warn(
"pipeline.repository.get_attachments.skip_malformed",
"error", err,
"contact_id", contactID,
"conversation_id", conversationID,
)
continue
}
atts = append(atts, att)
}
```
None required beyond this edit. The function already has `contactID` and `conversationID` parameters, so they can be safely included in the `slog.Warn` fields without changing behavior.
</issue_to_address>Help me be more useful! Please click 👍 or 👎 on each comment and I'll use the feedback to improve your reviews.
| @@ -67,6 +68,41 @@ func (r *redisPipelineRepository) GetBuffer(ctx context.Context, contactID, conv | |||
| return r.rdb.LRange(ctx, bufferKey(contactID, conversationID), 0, -1).Result() | |||
| } | |||
|
|
|||
There was a problem hiding this comment.
suggestion: Attachment deserialization logging lacks contact/conversation context, which may hinder debugging
In GetAttachments, malformed entries are logged without contact_id or conversation_id, making it difficult to tie warnings to specific conversations in production. Please include these IDs in the slog.Warn fields so bad data can be traced and cleaned up more easily, without changing behavior.
Suggested implementation:
for _, s := range raw {
var att model.Attachment
if err := json.Unmarshal([]byte(s), &att); err != nil {
slog.Warn(
"pipeline.repository.get_attachments.skip_malformed",
"error", err,
"contact_id", contactID,
"conversation_id", conversationID,
)
continue
}
atts = append(atts, att)
}None required beyond this edit. The function already has contactID and conversationID parameters, so they can be safely included in the slog.Warn fields without changing behavior.
Review follow-ups on the incoming-media path. The per-file cap was the only bound, so the failure modes it did not cover fell back on the customer losing the whole reply instead of just the media. - Shared byte budget (20 MiB) across every attachment of the call. The debounce window aggregates the media of all its messages, so a photo burst built a body of len(attachments) x 15 MiB; base64 pushed that past the gateway's client_max_body_size and the resulting 413 is not retryable, killing the text reply too. Probe: 20 x 2 MiB went from a 53 MiB request to 26 MiB. - Dedicated download timeouts. Downloads run before the AI call and outside its retry ceiling, but reused AI_CALL_TIMEOUT_SECONDS (30s) per attachment, so an unreachable media host stalled the turn by 30s x len(attachments) with no bound. Now 10s per download and 30s for the whole set. - Resolve the mime type from the bytes in hand: the response Content-Type wins, then the CRM's declared type, then the URL extension. The processor feeds this straight into Blob(mime_type=...), so an HTML error/login page answered with 200 was being forwarded as a valid image, and a missing content_type became application/octet-stream. Both are now dropped or resolved. - A Redis failure on the attachment buffer no longer aborts the turn: media is best-effort everywhere else in this path, and dropping the text reply over it contradicted the card's own acceptance criterion. Tests: the event -> debounce -> Redis -> A2ARequest seam had no coverage (the debounce mock always returned nil attachments), so a refactor could silently drop the media; two pipeline tests now pin it, including aggregation across the debounce window. Adapter tests cover the byte budget, the time budget, HTML responses, oversize files and the mime resolution table. Also repairs test/e2e, which has not compiled since EVO-2167 changed NewAIAdapter/NewDispatchEngine — which is why `go vet ./...` and `go test ./...` could not be run at all. Two assertions had drifted: the message signature moved to a prefix on the first segment in EVO-558, and the state-leak check raced the cleanup goroutine it was asserting on. go build ./... && go vet ./... && go test ./... green, e2e included.
* fix(EVO-2167): retry the AI Processor call on transient failures The bot-runtime made a single call to the AI Processor: any non-200 (401/5xx) or network error aborted the pipeline and the customer's message was dropped with no reply and no retry (only AI_CALL_TIMEOUT_SECONDS existed). A momentary blip — deploy, restart, DB hiccup — meant a permanently lost answer. - ai_adapter.go: extract a single attempt into doOnce() and wrap Call() in a retry loop with exponential backoff + jitter. Retryable: network errors and 429/500/502/503/504. NOT retried: 4xx (permanent), per-attempt timeout, pipeline cancellation. Body is built once and reused per attempt. - config.go: AI_CALL_MAX_RETRIES (default 2) and AI_CALL_RETRY_BASE_MS (default 200). - main.go: wire the new config into NewAIAdapter. - tests: 503->200 retry succeeds (2 calls); persistent 500 exhausts retries (1+2 calls); 400 not retried (1 call); network error then success (2 calls); maxRetries=0 disables retry. Existing tests updated to the new signature (0 retries). Complements EVO-2166: the processor now returns 503 (not a silent 401) on infra errors, so this retry covers the transient auth/infra case. Root cause of the incident is EVO-2141 (pool_pre_ping, already merged); this is defense in depth. Note: test/e2e/e2e_test.go was already incompatible with NewAIAdapter on develop (pre-existing, unrelated) and is left as-is; repo CI is docker-only (no go test lane). * fix(EVO-2167): harden retry — total-time cap + per-attempt timeout test Review follow-ups on the AI Processor retry path: - Add an overall time-budget backstop so the retry loop is provably bounded ((attempts+1) x per-attempt timeout + summed max backoff); the +1 slack keeps a per-attempt timeout surfacing as ErrAITimeout instead of being swallowed by the backstop. AC "teto de tempo total". - Add TestCall_TimeoutIsNotRetried: a per-attempt timeout must return ErrAITimeout and must NOT be retried with retries enabled. AC #5 "timeout por tentativa". - Document the idempotency contract on the retry path (502/504/network replay can re-run an already-processed turn; customer still gets one reply; dedupe of the duplicate server-side turn is the AI Processor's job, tracked in EVO-2166). * feat(EVO-2180): forward incoming media to the AI Processor as A2A file parts The bot only ever sent a text part, so images/audio the customer sent never reached the AI (the agent replied "No content to process"). Accept attachments on the inbound event, carry them through the debounce window, download each and send it as a base64 A2A file part. Part of EVO-2178 (image end-to-end). - pipeline/model: MessageEvent.Attachments + Attachment{URL,ContentType,FileType}. - pipeline/repository: AppendAttachments/GetAttachments on a parallel Redis list (bot_runtime:attach:{contact}:{conv}), aggregated like the text buffer and cleared together in ClearState (no stale media leaks into the next turn). - debounce/service: Start/Reset accept attachments; GetAttachments added. - pipeline/service: thread event.Attachments through start/skip/reset/advance -> the A2ARequest (read fresh from Redis at stage launch, like the buffer). - ai/model: A2ARequest.Attachments; JSONRPCPart.File + JSONRPCFile{Name,MimeType,Bytes} (tags match the processor's extract_files_from_message). - ai/service/ai_adapter: download each attachment once (before Marshal, reused across retries) with a 15 MiB cap; base64-encode; append a file part. A download failure is logged and skipped so the text-only message always survives. - tests: adapter forwards a file part with decodable base64 + download-failure sends text only; repo AppendAttachments/GetAttachments roundtrip + ClearState clears the attach key. Full suite green (go build/vet/test ./pkg/... ./internal/...). Note: test/e2e was already incompatible with NewAIAdapter on develop (pre-existing); repo CI is docker-only. * fix(EVO-2180): bound media forwarding and validate what is forwarded Review follow-ups on the incoming-media path. The per-file cap was the only bound, so the failure modes it did not cover fell back on the customer losing the whole reply instead of just the media. - Shared byte budget (20 MiB) across every attachment of the call. The debounce window aggregates the media of all its messages, so a photo burst built a body of len(attachments) x 15 MiB; base64 pushed that past the gateway's client_max_body_size and the resulting 413 is not retryable, killing the text reply too. Probe: 20 x 2 MiB went from a 53 MiB request to 26 MiB. - Dedicated download timeouts. Downloads run before the AI call and outside its retry ceiling, but reused AI_CALL_TIMEOUT_SECONDS (30s) per attachment, so an unreachable media host stalled the turn by 30s x len(attachments) with no bound. Now 10s per download and 30s for the whole set. - Resolve the mime type from the bytes in hand: the response Content-Type wins, then the CRM's declared type, then the URL extension. The processor feeds this straight into Blob(mime_type=...), so an HTML error/login page answered with 200 was being forwarded as a valid image, and a missing content_type became application/octet-stream. Both are now dropped or resolved. - A Redis failure on the attachment buffer no longer aborts the turn: media is best-effort everywhere else in this path, and dropping the text reply over it contradicted the card's own acceptance criterion. Tests: the event -> debounce -> Redis -> A2ARequest seam had no coverage (the debounce mock always returned nil attachments), so a refactor could silently drop the media; two pipeline tests now pin it, including aggregation across the debounce window. Adapter tests cover the byte budget, the time budget, HTML responses, oversize files and the mime resolution table. Also repairs test/e2e, which has not compiled since EVO-2167 changed NewAIAdapter/NewDispatchEngine — which is why `go vet ./...` and `go test ./...` could not be run at all. Two assertions had drifted: the message signature moved to a prefix on the first segment in EVO-558, and the state-leak check raced the cleanup goroutine it was asserting on. go build ./... && go vet ./... && go test ./... green, e2e included. * fix(EVO-2178): validate incoming media URLs before fetching them Review follow-up to EVO-2180. Attachment URLs arrive inside the /events payload and the adapter fetched them verbatim, so the endpoint doubled as a read primitive aimed by its caller: the bytes of any URL reachable from this service were base64-encoded into the A2A call, whose destination (outgoing_url) comes from the same payload. Reproduced end to end against a local metadata-style endpoint, both directly and through a 302. - checkMediaURL pins the scheme to http/https and requires the host to be one the CRM is known to serve blobs from: the postback URL's host (already mandatory in MessageEvent.Validate, so no new config for the default topology) plus whatever MEDIA_HOST_ALLOWLIST names, for deployments serving blobs off an S3/MinIO/CDN host. Unauthorized media is skipped and logged; the text reply is unaffected, like every other media failure here. - The download client re-runs that check on every redirect hop, so an authorized host cannot walk the fetch onto an internal address. - BOT_RUNTIME_SECRET becomes required. It was read with os.Getenv, and SecretMiddleware compares the header against it, so an empty value authenticated every caller that simply omitted the header. - The media buffer key gets a TTL. ClearState remains the normal cleanup; the TTL only stops a turn that dies before reaching it from leaving media URLs in Redis forever. - Attachment download failures now log the HTTP status: the common production case is a 404 from a signed link that expired while the queue was backed up, and it read identically to an unreachable host. - Adds .github/workflows/ci.yml. Nothing ran the Go suite on a PR, which is how test/e2e stayed non-compiling from EVO-558 until EVO-2180. * style(EVO-2180): trim the comments to what is not obvious from the code * fix(EVO-2178): take the media host allowlist from config, not from the event Fixes a regression I introduced in #6. The allowlist was anchored on the host of the event's postback_url, which never matches the host the CRM actually signs media URLs with: postback_url comes from BOT_RUNTIME_POSTBACK_BASE_URL (internal DNS, "evo-crm" in the shipped compose) while the URL is built from ACTIVE_STORAGE_URL, falling back to BACKEND_URL — which production requires to be a public host. So on develop every attachment was rejected as blocked_url and the agent stopped seeing images: the EVO-2178 bug, back. Anchoring on the event was also the wrong shape for the guard. Whoever sends the event chooses every field in it, including the one being used to decide what that same event may reach, so the check constrained nobody it needed to. Reading MEDIA_HOST_ALLOWLIST only puts the decision with the operator, where it cannot be chosen by the caller. A2ARequest.PostbackURL is dropped again. The scheme check and the per-redirect re-check are unchanged. This makes the variable required wherever media is expected: unset means no attachment is fetched. The deploy surfaces are wired up in the umbrella PR; k8s/configmap.yaml and k8s/deployment.yaml carry it here. * Merge pull request #9 from evolution-foundation/fix/CRM-236-degraded-provider-feedback fix(pipeline): tell the customer when the AI backend fails (CRM-236) --------- Co-authored-by: Matheus Pastorini <matheus.pastorini@etus.com.br> Co-authored-by: Matheus Pastorini <pastorinimatheus@gmail.com>
EVO-2180 — encaminhar mídia de entrada ao AI Processor como A2A file part
Sub-issue de EVO-2178 (agente não processa mídia). O bot só enviava um part de texto → imagem/áudio do cliente nunca chegavam à IA.
O que faz
Aceita
attachmentsno evento de entrada, carrega pela janela de debounce, baixa cada anexo e envia como file part base64{type:"file", file:{name,mimeType,bytes}}.pipeline/modelMessageEvent.Attachments+Attachment{URL,ContentType,FileType}pipeline/repositoryAppendAttachments/GetAttachmentsnuma lista Redis paralela (bot_runtime:attach:{contact}:{conv}), agregada como o buffer de texto e limpa junto noClearState(mídia velha não vaza p/ o próximo turno)debounce/serviceStart/Resetrecebem attachments;GetAttachmentspipeline/serviceevent.Attachmentspor start/skip/reset/advance →A2ARequest(lido fresh do Redis no launch, como o buffer)ai/modelA2ARequest.Attachments;JSONRPCPart.File+JSONRPCFile{Name,MimeType,Bytes}(tags batem com oextract_files_from_messagedo processor)ai/service/ai_adapterTestes
AppendAttachments/GetAttachments(ordem) +ClearStatelimpa a attach-key.go build ./...✅ ·go vet ./pkg/... ./internal/... ./cmd/...✅ ·go test ./pkg/... ./internal/...ok (ai/service, debounce, dispatch, pipeline/handler, pipeline/repository, pipeline/service) — zero regressão.Depende de / relacionado
Ponta-a-ponta da imagem = EVO-2179 (CRM envia attachments) + EVO-2181 (processor deixa a imagem chegar ao modelo, já em PR #44).
Summary by Sourcery
Forward incoming media attachments through the debounce and pipeline to the AI Processor, where they are downloaded, size-limited, base64-encoded, and sent as JSON-RPC file parts alongside the text message.
New Features:
Bug Fixes:
Enhancements:
Tests: