-
Notifications
You must be signed in to change notification settings - Fork 34
feat(orchestrator): a seam for taking queued executions off the launch path #346
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -39,6 +39,21 @@ class OrchestratorError(RuntimeError): | |
| pass | ||
|
|
||
|
|
||
| class QueuedExecutionInterceptor(typing.Protocol): | ||
| """Given a chance to take a queued execution off the launch path. | ||
|
|
||
| Implemented downstream. Called on the orchestrator's session once the execution is | ||
| known to be launchable -- inputs present, not conditionally skipped, no cache hit, not | ||
| cancelled. An implementation that returns True owns the execution from that point: it | ||
| sets whatever status it wants and commits. The orchestrator makes no assumption about | ||
| which status that is. | ||
| """ | ||
|
|
||
| def intercept(self, *, session: orm.Session, execution: bts.ExecutionNode) -> bool: | ||
| """True if this execution was taken over and must not launch; False to continue.""" | ||
| ... | ||
|
|
||
|
|
||
|
|
||
| class OrchestratorService_Sql: | ||
| def __init__( | ||
| self, | ||
|
|
@@ -57,6 +72,7 @@ def __init__( | |
| _max_container_execution_refresh_error_retries: int = 3, | ||
| _max_queue_batch_size: int = 1, | ||
| _max_queue_batch_duration: datetime.timedelta = datetime.timedelta(), | ||
| queued_execution_interceptor: QueuedExecutionInterceptor | None = None, | ||
| ): | ||
| self._session_factory = session_factory | ||
| self._launcher = launcher | ||
|
|
@@ -75,6 +91,7 @@ def __init__( | |
|
|
||
| self._max_queue_batch_size = _max_queue_batch_size | ||
| self._max_queue_batch_duration = _max_queue_batch_duration | ||
| self._queued_execution_interceptor = queued_execution_interceptor | ||
|
|
||
| def run_loop(self): | ||
| while True: | ||
|
|
@@ -124,12 +141,8 @@ def internal_process_queued_executions_queue(self, session: orm.Session): | |
| query_start_timestamp = time.monotonic_ns() | ||
| query = ( | ||
| sql.select(bts.ExecutionNode).where( | ||
| bts.ExecutionNode.container_execution_status.in_( | ||
| ( | ||
| bts.ContainerExecutionStatus.UNINITIALIZED, | ||
| bts.ContainerExecutionStatus.QUEUED, | ||
| ) | ||
| ) | ||
| bts.ExecutionNode.container_execution_status | ||
| == bts.ContainerExecutionStatus.QUEUED | ||
|
Comment on lines
142
to
+145
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. (AI-assisted) Overloading Three reasons the reuse bothers me:
The honest counter-argument, which I don't think is fatal: So: either add the state now, or — if the cost wins — please record the decision in the enum comment ( |
||
| ) | ||
| # TODO: Maybe add last_processed_at | ||
| # .order_by(bts.ExecutionNode.last_processed_at) | ||
|
|
@@ -610,6 +623,15 @@ def internal_process_one_queued_execution( | |
| session.commit() | ||
| return | ||
|
|
||
| # Give the interceptor a chance to take this execution off the launch path. | ||
| # If it returns True it has taken ownership: it decided what state the execution is | ||
|
Comment on lines
+626
to
+627
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. (AI-assisted) The seam's contract is positional, and no test pins the position — which matters more than usual because a consumer vendoring this repo at a pinned SHA cannot detect it moving. The docstring's promise is where this runs: "once the execution is known to be launchable — inputs present, not conditionally skipped, no cache hit, not cancelled." The failure that allows is quiet. Move it above the cache lookup and a cache hit now consults the interceptor — an implementation that gates on capacity would charge a slot for work that is about to be satisfied from cache, and could park a node waiting for a slot it never needed. Move it above the cancel check and cancelled executions do the same. Neither errors; both surface much later as "the gate is mysteriously saturated." Worth pinning here specifically because this is a library seam. A downstream implementation typically vendors this repo at a fixed revision, so its own CI compiles against a frozen copy — a move here cannot fail any test it owns until someone advances the pin, at which point the behaviour change looks like it came from the bump rather than the refactor. These tests are the only place the guarantee can be stated where it runs on every push. Cheap, given # interceptor.calls == [] for each of:
# 1. missing input artifact -> WAITING_FOR_UPSTREAM
# 2. cache hit -> reuses cached execution
# 3. desired_state = TERMINATED -> CANCELLED (node-level and run-level)
# 4. is_enabled: false -> SKIPPEDIf only one is worth it, take the cache-hit case: it is the one whose failure mode is silent. The cancel case at least ends in a terminal status somebody notices. |
||
| # in and committed that itself. We stop here and do not launch. | ||
| if self._queued_execution_interceptor is not None: | ||
| if self._queued_execution_interceptor.intercept( | ||
| session=session, execution=execution | ||
|
Comment on lines
+626
to
+631
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. (AI-assisted) An implementation that returns The contract is prose only: "it sets whatever status it wants and commits." If an implementation returns
That is the exact failure your own description gives as the reason for narrowing the selector, reintroduced through the seam — but silent this time, because a buggy implementation looks identical to a working one from here. Cheap fix — make the violation loud rather than fatal: if self._queued_execution_interceptor.intercept(session=session, execution=execution):
session.refresh(execution)
if execution.container_execution_status == bts.ContainerExecutionStatus.QUEUED:
_logger.error(
"Interceptor claimed execution %s but left it QUEUED; it will be re-swept.",
execution.id,
)
returnStronger alternative: make the contract enforceable by construction — have |
||
| ): | ||
|
Comment on lines
+629
to
+632
Collaborator
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. (AI-assisted) Building on the docstring comment above (#discussion_r3919337398) rather than repeating it — that one asks you to document that raising here is fatal. This is the argument that the call site should instead make it survivable. To restate only the conclusion: an exception out of The reason that is worth more than a docstring: this is an optional seam whose implementation is, by construction, foreign code talking to storage the orchestrator knows nothing about. As written, installing an admission gate silently couples every pipeline's survival to that gate's availability — a lock-wait timeout or a deploy-time blip during a sweep is enough to kill a run outright. That is a much stronger commitment than an opt-in hook implies, and an implementer reading the signature has no way to infer it. Fail-open at the call site: if self._queued_execution_interceptor is not None:
try:
if self._queued_execution_interceptor.intercept(session=session, execution=execution):
return
except Exception as exc:
_logger.exception("Queued execution interceptor raised; launching anyway.")
bugsnag_instrumentation.notify(exception=exc)
session.rollback()A broken gate then degrades to "no gating" rather than "no pipelines". The cost is a bounded, self-correcting overshoot while the implementation is down; the benefit is that an optional component cannot take the orchestrator's core job with it. If you would rather fail closed, one caveat worth stating explicitly: swallowing the exception and falling through without launching leaves the row |
||
| return | ||
|
|
||
| # Creating new container execution | ||
| container_execution_uuid = _generate_random_id() | ||
|
|
||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
(AI-assisted)
The Protocol licenses
commit()on the orchestrator's own session, but doesn't state the invariant that makes that safe — or what happens if the implementation raises.The
commit()is only safe because of something invisible from here: at the call site (line 626) the session has no pending orchestrator writes.session.rollback()at line 597 clears everything, and the only work between that rollback andinterceptis the twoextra_datareads for the cancellation check. Nothing pins that invariant. Add alast_processed_atbump or a status-history append between 597 and 626 later, and an implementation'scommit()silently flushes orchestrator state that was never meant to be durable — presenting as half-applied execution rows far from this diff.Separately, raising out of
interceptis a much bigger deal than the docstring implies: it propagates to the handler ininternal_process_queued_executions_queue(lines 165-185), which marks the executionSYSTEM_ERRORand_mark_all_downstream_executions_as_skipped. So a transient error in downstream code — which by design talks to its own tables — permanently kills the run.Both belong in this docstring:
Worth a one-line comment above
session.rollback()at line 597 too, noting the clean-session invariant is load-bearing for the seam.