Draft: Drive risk assessment lifecycle from session events

Proposal to move the risk assessment row's lifecycle off a 30-minute timeout clock and onto the Duo Agent Platform session's own events. The review of !253762 (merged) surfaced race and stale-clock edges where the timer could overwrite a completed row as failed, or leave a second run or late result wrongly failed. Targets !253762 (merged)'s source branch so the diff is only the redesign; retarget to master once that merges. Everything stays behind the disabled duo_mr_risk_classification flag, and this is deliberately unfinished (no new specs) - posted to discuss the shape first.

Detailed context for AI agents

Context: related to epic &23131 (MR Risk Assessment - Customer Zero) and issue #609292 (closed). Builds on !253762 (merged) ("Run merge request risk classification end to end"), which introduced one merge_requests_risk_assessments row per MR written by four actors with separate guards: a before_start lambda in the flow definition (queues the row, schedules Ai::RiskClassification::TimeoutWorker 30 minutes out), the results REST endpoint via UpdateResultsService#refresh, CalculateScoreWorker (score + complete), and TimeoutWorker (failed, based on idle time via updated_at).

Problems found in !253762 (merged)'s design:

  • A stale in-memory timer instance can overwrite an already-completed row as failed.
  • updated_at is not bumped when enqueue is refused on an already-queued row, or when refresh loops queued-to-queued, so a second run or a late result can be wrongly failed by the timer.
  • The timer is a proxy for "did the agent session die", but the Duo Agent Platform session already tracks its own status.

New design: the row follows its session. Ai::RiskClassification::SessionLifecycleWorker (EventStore subscriber, data_consistency :sticky, feature_category :duo_code_review) is the sole owner of the lifecycle. It subscribes to Ai::DuoWorkflows::WorkflowStartedEvent, WorkflowFinishedEvent, and the new WorkflowDroppedEvent, and ignores any workflow whose workflow_definition is not risk_classification/v1 or that has no merge_request.

Duo Agent Platform session (Ai::DuoWorkflows::Workflow)
        |
        |-- Started         --> WorkflowStartedEvent  ---\
        |-- Finished        --> WorkflowFinishedEvent ---+--> EventStore --> SessionLifecycleWorker
        |-- Dropped/Stopped --> WorkflowDroppedEvent  ---/                        |
                                                                                  v
                                              RiskAssessment row (merge_requests_risk_assessments)
                                                - Started:  ensure_for!(mr), duo_workflow_id = session.id, enqueue
                                                - Dropped:  mark_failed if tracking?(session) and can_mark_failed?
                                                - Finished, still queued and tracking?(session):
                                                    < 5 min since session finished -> re-dispatch later
                                                    otherwise                       -> mark_failed

Agent claims POST (unchanged) --> UpdateResultsService#refresh --> CalculateScoreWorker
                                                                   update!(score) + finish, re-checking queued? under with_lock

Worker logic:

  • Started: RiskAssessment.ensure_for!(mr), then under with_lock set duo_workflow_id to the session id and call enqueue (pending/complete/failed -> queued; no-op if already queued). The row always tracks the latest session that started for the MR.
  • Dropped (workflow drop or stop transition): if the row's duo_workflow_id still matches this session and the row can be failed, mark_failed under lock. Score fields are left untouched so a previous result stays on the row.
  • Finished: if the row still tracks this session and is queued, and the session finished less than FINISH_GRACE (5 minutes) ago, re-dispatch the same event via handle_event_in(FINISH_GRACE); otherwise mark_failed as "finished without reporting". The grace exists because the claims POST enqueues CalculateScoreWorker asynchronously, and the session can finish before that job runs.
    • Known soft spot: the grace is measured from workflow.updated_at, which could keep moving for unrelated reasons and delay the re-check.
  • New model helper RiskAssessment#tracking?(workflow) = duo_workflow_id == workflow.id.
  • duo_workflow_id semantics shift: while running it means "the session the row currently follows"; CalculateScoreWorker still writes it with the score, so after completion it is the session that produced the score. The widget's "View session" link therefore points at the in-flight session while running.

Removed: before_start lambda from the risk classification flow definition, TimeoutWorker and its spec, RiskAssessment#enqueue_timeout, and the after_transition on: :enqueue hook. ensure_for! is kept because two sessions can still start concurrently for the same MR.

Platform-level changes (small, generic):

  • Ai::DuoWorkflows::Workflow#publishes_lifecycle_events? = messaging_callback_context present, OR the workflow's foundational flow has the new boolean attribute session_lifecycle_events (added to Ai::Catalog::FoundationalFlow::Attributes, default false). Started/Finished publishing is now gated on this predicate instead of messaging_callback_context directly.
  • New WorkflowDroppedEvent (CloudEvent, category ai_duo_workflows, type workflow_dropped, data workflow_id), published after_transition on drop and stop under the same gate. Event docs added/updated under data/events/ai/duo_workflows, doc/development/eventstore/events.md regenerated.
  • Side effect: Ai::Messaging::CallbackWorker now also receives Started/Finished for risk classification sessions; it returns early because they have no messaging_callback_context, so one extra no-op Sidekiq job per session start and finish.

CalculateScoreWorker: keeps the early queued? check, then re-checks queued? under with_lock before update! + finish, so the lifecycle worker cannot fail the row mid-write.

Unchanged: UpdateResultsService and the results endpoint.

Queue config: config/sidekiq_queues.yml and ee/app/workers/all_queues.yml regenerated for the new worker and the removal of TimeoutWorker.

Alternatives considered:

  • Storing claims on the row at POST time and scoring on Finished, dropping the grace entirely - needs a new column for pending claims, deferred.
  • Scoring synchronously in the results request - avoids the ordering issue but puts signal extraction on the API request path.
  • Deduping the two triggers (MR created, draft-to-ready) at the trigger layer - a separate question, out of scope here.

Tradeoffs: failure detection latency now depends on the session platform's own reconciliation (Ci::Workloads::WorkloadFinishedEvent handling plus the stale-session cron), up to about an hour, versus a flat 30 minutes before. Relies on EventStore delivery, as the existing messaging callbacks do.

Verification done: rubocop clean on all changed files; rails runner boot check of the event build, the new definition attribute, and the worker's data consistency; ee/spec/workers/ai/risk_classification/calculate_score_worker_spec.rb passes (7 examples).

Deliberately not done (draft):

  • No new specs for SessionLifecycleWorker or the platform changes.
  • 10 existing examples in risk_assessment_spec.rb and definition_spec.rb fail because they test the removed TimeoutWorker / before_start behaviour; they need rewriting if this direction is accepted.
  • No spec for the ai_subscriptions change.
  • Widget not updated to show the previous score on a failed row (isClassified still excludes FAILED).

Feature flag: everything stays behind duo_mr_risk_classification (wip, default off). No changelog trailer since the feature is unreleased.

Merge request reports

Loading
Loading