Stream Duo Workflow progress to messaging adapters
What does this MR do and why?
Adds a generic, surface-agnostic progress streaming path for messaging-triggered Duo Workflows (e.g. the internal Slack surface). Today these flows show an acknowledgement and the final answer, with minutes of silence in between. This lets adapters render live updates (plan / tool calls / agent reasoning) while a flow runs.
The path is dormant in this MR: no adapter opts in yet (supports_live_progress? is false for all adapters), so nothing is enqueued in production. The follow-up that implements Slack#on_progress flips the opt-in and activates it. This MR ships the tested plumbing and the correct abstract contract.
How it works
checkpoint write (messaging workflow only)
-> CreateCheckpointService enqueues ProgressDeliveryWorker.perform_in(2s)
(debounced per workflow; skipped unless the adapter opts in via
supports_live_progress?)
-> worker reads ProgressReader#delta_since(workflow, cursor)
(cursor-based delta of the flow's own ui_chat_log entries)
-> adapter.on_progress(delta:) renders on its surface
-> cursor advances in messaging_callback_context (atomic jsonb merge)Key decisions:
- Debounce, not one-job-per-checkpoint.
perform_in+deduplicate :until_executed, including_scheduled: true, if_deduplicated: :reschedule_oncecollapses a burst into one delivery and bounds the per-channel update rate. On real runs this coalesced 18 checkpoints into 4 deliveries; a quiet workflow does zero work. - Read the durable plane. Backend sinks like Slack are reached via Rails and only see persisted checkpoints, so progress is read from the checkpoint, not the gRPC stream that feeds directly-connected clients.
- The flow defines what to show.
ProgressReaderreturns the flow's ownui_chat_logentries (message_type,tool_info, ...); the sink decides how to render them -- same model as the web UI. - Single-reader seam.
Checkpoint#ui_chat_logcentralises the rawchannel_values.ui_chat_logdig-path. When delta/incremental checkpoints land, only this method (and the reader) change. - Concurrency-safe context writes.
merge_messaging_callback_context!merges atomically in Postgres (jsonb||) so the progress cursor and the adapters'status_ts/session_urlwrites can't clobber each other.
How to test locally
The path is dormant, so to see it work you apply a throwaway patch that (a) opts an easy-to-trigger adapter in and (b) gives it a temporary on_progress that logs (the base method is abstract). The gitlab_duo_note adapter is easiest.
-
Apply this diff locally (do not commit it):
--- a/ee/app/services/ai/messaging/adapters/gitlab_duo_note.rb +++ b/ee/app/services/ai/messaging/adapters/gitlab_duo_note.rb @@ def self.adapter_key 'gitlab_duo_note' end + # TEST ONLY -- do not commit + def self.supports_live_progress? + true + end + + # TEST ONLY -- do not commit + def on_progress(delta:, callback_context:) + Gitlab::AppLogger.info( + message: 'Duo progress (local test)', + new_messages: delta.messages.size, + message_types: delta.messages.map { |m| m['message_type'] }.tally, + tools: delta.messages.filter_map { |m| m.dig('tool_info', 'name') }.tally + ) + end -
Ensure this feature flag is on so that duo developer goes through the adapter:
Feature.enable(:ai_use_messaging_adapter_for_mentions) -
Restart Rails + Sidekiq so the change loads (
gdk restart rails-web rails-background-jobs). -
Trigger a flow:
@duo-developer-gitlab-duomention in an MR comment in a project with the Duo Agent Platform enabled. (e.g.@duo-developer-gitlab-duo what is this MR about) -
Watch the deliveries (debounced, coalesced):
tail -f log/application_json.log | grep --line-buffered "Duo progress (local test)"
Follow-ups (intentionally not here)
- Migrate the remaining 6
ui_chat_logconsumers (presenter, mailer, GraphQL type, REST API, model dig, rake task) ontoCheckpoint#ui_chat_log. Slack#on_progressrendering + flippingSlack.supports_live_progress?totrue, behind a feature flag (rollout + kill-switch).- Slack
Retry-After/ timeout / circuit-breaker hardening onSlack::API.
References
- Relates to #602540 (closed)
Database review
This MR adds one new statement, in Ai::DuoWorkflows::Workflow#merge_messaging_callback_context!: an atomic JSONB merge so concurrent writers (the progress cursor vs an adapter's status_ts/session_url) can't clobber each other's keys. It's a single-row update scoped by primary key.
Raw SQL (verbatim, as captured from the executed statement):
UPDATE "duo_workflows_workflows"
SET messaging_callback_context =
COALESCE(messaging_callback_context, '{}'::jsonb) || '{"progress_cursor": {...}}'::jsonb
WHERE "duo_workflows_workflows"."id" = $ID;Query plan (postgres.ai, gitlab-production-main): https://postgres.ai/console/gitlab/gitlab-production-main/sessions/53074/commands/154802
ModifyTable on public.duo_workflows_workflows (cost=0.43..3.45 rows=0 width=0) (actual time=18.127..18.128 rows=0 loops=1)
Buffers: shared hit=115 read=24 dirtied=2
-> Index Scan using duo_workflows_workflows_pkey on public.duo_workflows_workflows (cost=0.43..3.45 rows=1 width=38) (actual time=4.794..4.797 rows=1 loops=1)
Index Cond: (id = 4603815)Single-row Index Scan on the primary key (the ModifyTable reports rows=0 as UPDATE always does; the inner scan's rows=1 shows exactly one row updated). The ~18ms is cold-clone I/O (read=24, dirtied=2); with a warm cache it is sub-millisecond. Comfortably within query timing guidelines.