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_once collapses 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. ProgressReader returns the flow's own ui_chat_log entries (message_type, tool_info, ...); the sink decides how to render them -- same model as the web UI.
  • Single-reader seam. Checkpoint#ui_chat_log centralises the raw channel_values.ui_chat_log dig-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_url writes 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.

  1. 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
  2. Ensure this feature flag is on so that duo developer goes through the adapter: Feature.enable(:ai_use_messaging_adapter_for_mentions)

  3. Restart Rails + Sidekiq so the change loads (gdk restart rails-web rails-background-jobs).

  4. Trigger a flow: @duo-developer-gitlab-duo mention in an MR comment in a project with the Duo Agent Platform enabled. (e.g. @duo-developer-gitlab-duo what is this MR about )

  5. 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_log consumers (presenter, mailer, GraphQL type, REST API, model dig, rake task) onto Checkpoint#ui_chat_log.
  • Slack#on_progress rendering + flipping Slack.supports_live_progress? to true, behind a feature flag (rollout + kill-switch).
  • Slack Retry-After / timeout / circuit-breaker hardening on Slack::API.

References

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.

Edited by Thomas Schmidt

Merge request reports

Loading
Loading