Skip to content

feat(orchestrator): system/pipeline_run.ended_at run annotation - #391

Open
morgan-wowk wants to merge 1 commit into
masterfrom
run-ended-annotation
Open

morgan-wowk wants to merge 1 commit into
masterfrom
run-ended-annotation

Conversation

@morgan-wowk

Copy link
Copy Markdown
Collaborator

Adds a system-owned run annotation, system/pipeline_run.ended_at. It's written when a run's last container execution ends, so "is this run still running?" becomes an indexed annotation lookup instead of a status aggregation. The first consumer is the trigger max_concurrency cap in oasis-backend: https://github.com/Shopify/oasis-backend/pull/710

sequenceDiagram
    participant O as Orchestrator
    participant S as Session (after_flush)
    participant DB as pipeline_run_annotation
    O->>S: last execution → SUCCEEDED / FAILED / …
    S->>S: every container execution in the run ended?
    S->>DB: INSERT system/pipeline_run.ended_at = <ISO time>
    Note over O,DB: same transaction as the status change
    O->>S: execution leaves an ended status
    S->>DB: DELETE system/pipeline_run.ended_at
Loading

Changes

  • The container_execution_status set listener flags executions that cross into or out of an ended status. A new after_flush listener reconciles the annotation for the affected runs using Core SQL on the session connection, so it's atomic with the node change and idempotent.
  • Adds a new (key, value) index on pipeline_run_annotation, created on startup for existing databases. Every existing index leads with pipeline_run_id, so key = … AND value = … was a full index scan on MySQL (EXPLAIN type=index → ref with this index).
  • key_exists filtering is allowed on the new key. not key_exists system/pipeline_run.ended_at lists the runs that are still running.

Notes

  • No backfill. Runs that ended before this deploys have no annotation.
  • The listener runs wherever orchestrator_sql is imported, which means the orchestrator process. A terminal status set through the admin status route on the API process doesn't write it.
  • The key is under system/, so users can't set or delete it through the annotation API.

Testing

  • tests/test_pipeline_run_ended_at.py covers:

    • written only when the last execution ends;
    • every ended status counts;
    • removed when an execution reopens;
    • the first timestamp is kept;
    • flush-then-commit and rollback;
    • a single-container root run;
    • other runs untouched;
    • users can't set the key;
    • the index is created on an existing database.

    With the listener removed, the five positive tests fail.

  • The full suite passes on Python 3.13 and 3.10 (540 tests).

🤖 Generated with Claude Code

Mark a pipeline run as ended with a system-owned annotation so callers can
ask "is this run still running?" with an indexed annotation lookup instead
of aggregating execution statuses.

- An after_flush listener writes system/pipeline_run.ended_at (ISO
  timestamp) in the same transaction as the execution status change that
  ends the run's last container execution, and removes it if an execution
  leaves an ended status.
- Add a (key, value) index on pipeline_run_annotation, created on startup
  for existing databases, so key/value reverse lookups use an index.
- Allow key_exists filtering on the new system key.

Runs that ended before this change carry no annotation (no backfill).
@morgan-wowk
morgan-wowk requested a review from a team October 9, 2026 04:20
@morgan-wowk
morgan-wowk requested a review from Ark-kun as a code owner October 9, 2026 04:20
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant