Skip to content

Track workflow runs in a database for scalable per-user listing #30

Description

@mehalter

Summary

Per-user workflow run listing (GET /workflows/runs, surfaced in the Workflows
tab) currently filters Airflow DAG runs in the API Lambda because Airflow has no
server-side filter for the attribute we key ownership on. This works but does
not scale with total run volume. We want to move run ownership/indexing into a
database so per-user listing is a bounded, paginated query instead of a scan.

This issue is filed on the frontend repo for tracking, but the core change is
backend (cape-ph/cape-cod). The frontend consumes the endpoint and will only
need to adopt pagination once the backend supports it.

Background: how per-user runs work today

Airflow runs in MWAA behind a single IAM principal, so Airflow cannot natively
know which Cognito user triggered a DAG. Ownership is therefore attributed by
CAPE, not by Airflow:

  • The API Gateway authorizer resolves the caller from the Cognito token and
    exposes their identity to handlers.
  • On trigger, post_workflow_run.py stamps the caller into the Airflow DAG run
    config as conf.cape = { triggering_user_id, triggering_user_name } (and a
    human-readable note). Any client-supplied conf.cape is stripped first to
    prevent spoofing.
  • On listing, get_workflow_runs.py fetches recent runs across all DAGs via
    Airflow's cross-DAG endpoint GET /dags/~/dagRuns/list (paged), then filters
    in the Lambda to the runs whose conf.cape.triggering_user_id matches the
    caller.
  • The frontend (src/lib/workflowStatus.ts -> getMyWorkflowRuns) just calls
    GET /workflows/runs and renders the result. There is no client-side run
    tracking anymore (the old workflow_runs cookie was removed).

The problem

Airflow (all versions, including the 3.0.6 we run) has no server-side filter on
a conf value. The native filterable triggering_user_name field is exactly
what MWAA cannot populate, which is why ownership lives in conf.cape. As a
result:

  • Listing a single user's runs requires fetching recent runs across all DAGs and
    filtering in the Lambda. Cost grows with total run volume, not with the number
    of runs the user actually owns.
  • There is a scan cap (MAX_RUNS_SCANNED, currently 1000). Once total runs
    exceed that window, a user's older runs can silently drop off the list even
    though they still exist in Airflow.
  • Latency and payload grow as the deployment gets busier and as Airflow
    retention grows.
  • True pagination for a single user is not possible, because pages are defined
    over all runs, not over the filtered per-user set.

Desired future solution

Once the CAPE environment database is more developed, back workflow-run
ownership with a database index:

  • At trigger time, write a row keyed by user: at minimum
    { triggering_user_id, dag_id, dag_run_id, submitted_at }, optionally a small
    submission snapshot (or just rely on Airflow conf for details).
  • Change GET /workflows/runs to query that table by triggering_user_id with
    real pagination and ordering, returning dag_id / dag_run_id pairs.
  • Keep Airflow as the source of truth for run state: hydrate live state/task
    info from Airflow per listed run (or lazily), so the DB only owns the
    ownership/index, not the run lifecycle.
  • Keep conf.cape as durable, in-Airflow attribution (backfill/repair source and
    a hedge against DB drift).

Open design questions to settle when picking this up:

  • Which datastore (the environment DB being developed vs. a dedicated table) and
    where the write happens (in post_workflow_run.py after a successful trigger,
    transactionally with the Airflow trigger call as best-effort).
  • Backfill strategy for runs triggered before the table existed (one-time scan of
    Airflow conf.cape to seed the table).
  • Retention alignment between the index table and Airflow run retention.
  • Pagination contract for GET /workflows/runs and the matching frontend change
    (getMyWorkflowRuns -> paged fetch / infinite scroll).

Acceptance criteria

  • GET /workflows/runs returns a user's runs via a bounded DB query, not a
    full-scan-and-filter, with pagination.
  • No silent truncation of a user's older runs (the MAX_RUNS_SCANNED cap is
    removed or made irrelevant).
  • Run state is still accurate (hydrated from Airflow).
  • Existing runs remain visible (backfill or dual-read during migration).

References

  • Backend (cape-ph/cape-cod):
    • assets/api/capi/handlers/get_workflow_runs.py (current list + filter,
      MAX_RUNS_SCANNED, /dags/~/dagRuns/list)
    • assets/api/capi/handlers/post_workflow_run.py (conf.cape + note)
    • assets/api/authz/default_apigw_authorizer.py (identity resolution)
  • Frontend (this repo):
    • src/lib/workflowStatus.ts (getMyWorkflowRuns)
    • src/lib/components/Status/Status.svelte

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions