You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
An extension mechanism, such that Opsqueue can be extended while preserving separation of concerns.
A first extension: delegation.
This 'delegation' extension can be used to delegate the responsibility of starting/stopping submissions to another service, and keeping that other service up-to-date on a submission's status. An example use-case is a job scheduling system that keeps tracks of dependencies between jobs (which are not necessarily Opsqueue submissions), and schedules them accordingly. Enabling cross-system dependencies to be expressed and handled.
See below for a sequence diagram how this works in practice:
The happy path:
sequenceDiagram
autonumber
participant CL as Client
participant D as Delegator
participant OD as Opsqueue (Delegation module)
participant C as Opsqueue (Core module)
participant DB@{"type": "database"} as Database
CL ->>+ OD: insert_submission(chunks)
OD -->>- CL: submission_id
CL ->> D: submit(submission_id)
Note over D: Delegator is free to wait. (e.g. to ensure all dependencies are satisfied)
D ->>+ OD: delegate(task_id, submission_id)
OD ->> DB: insert_external_task(task_id, submission_id)
OD -->> C: unpause(submission_id)
OD -->>- D: 202 ACCEPTED
C -)+ OD: submission_status_changed(submission_id)
OD ->>+ DB: select_external_tasks(submission_id)
DB -->>- OD: [{task_id, submission_id, last_status_sent}]
OD ->>+ C: submission_status(submission_id)
C -->>- OD: "IN_PROGRESS"
OD ->> D: submit(Updated(task_id, "RUNNING"))
OD ->>- DB: update_last_status_sent(task_id, "RUNNING")
Note over OD,C: Some time passes until the submission is completed.
C -)+ OD: submission_status_changed(submission_id)
OD ->>+ DB: select_external_tasks(submission_id)
DB -->>- OD: [{task_id, submission_id, last_status_sent}]
OD ->>+ C: submission_status(submission_id)
C -->>- OD: "COMPLETED"
OD ->> D: submit(Completion(task_id, "SUCCESS"))
OD ->>- DB: delete_external_task(task_id)
Note over D: Delegator can respond to status change (e.g. start processing jobs that depended on this submission)
Loading
The path when a submission is cancelled mid-way:
sequenceDiagram
autonumber
participant CL as Client
participant D as Delegator
participant OD as Opsqueue (Delegation module)
participant C as Opsqueue (Core module)
participant DB@{"type": "database"} as Database
CL ->>+ OD: insert_submission(chunks)
OD -->>- CL: submission_id
CL ->>+ D: submit(submission_id)
D -->>- CL: task_id
Note over D: Delegator is free to wait. (e.g. to ensure all dependencies are satisfied)
D ->>+ OD: delegate(task_id, submission_id)
OD ->> DB: insert_external_task(task_id, submission_id)
OD -->> C: unpause(submission_id)
OD -->>- D: 202 ACCEPTED
C -)+ OD: submission_status_changed(submission_id)
OD ->>+ DB: select_external_tasks(submission_id)
DB -->>- OD: [{task_id, submission_id, last_status_sent}]
OD ->>+ C: submission_status(submission_id)
C -->>- OD: "IN_PROGRESS"
OD ->> D: update(task_id, "RUNNING")
OD ->>- DB: update_last_status_sent(task_id, "RUNNING")
Note over OD,C: Some time passes but the submission is not completed yet.
CL -) D: kill(task_id)
D -) OD: kill(task_id)
OD -->> C: cancel(submission_id)
C -)+ OD: submission_status_changed(submission_id)
OD ->>+ DB: select_external_tasks(submission_id)
DB -->>- OD: [{task_id, submission_id, last_status_sent}]
OD ->>+ C: submission_status(submission_id)
C -->>- OD: "CANCELLED"
OD ->> D: submit(Failure(task_id, "FORCED"))
OD ->>- DB: delete_external_task(task_id)
Note over D: Delegator can respond to status change (e.g. block processing jobs that depend on this submission)
Loading
The path when an error occurs while running the submission is similar to the path for cancellation, without the client explicitly canceling the submission.
Cancellation log mislabels completed or failed submissions
opsqueue/src/delegation/server.rs:270
This branch also handles completed and failed submissions, so the message incorrectly states they were cancelled when a kill races with terminal completion. Preserve the error value and log that the submission is not cancellable (including its actual reason).
This issue also appears in the following locations of the same file:
Status reporting can deadlock with a one-connection reader pool
Our solution here is to not use a one-connection reader pool.
and introduces costly periodic N+1 queries.
For the reviewer, I'm also still a little bit worried about this. If this proves to be a problem in the future, we might just want to do the following (even though that violates the separation of concerns a little bit):
WITH out_of_date_tasks AS (
SELECT
submission_id,
task_id
FROM submissions_external_task as t
WHEREt.last_status_sent IS NULLOR (t.last_status_sent='paused'AND NOT EXISTS(SELECT*FROM submissions_paused AS s WHEREs.id=t.submission_id))
OR (t.last_status_sent='in_progress'AND NOT EXISTS(SELECT*FROM submissions AS s WHEREs.id=t.submission_id))
OR (t.last_status_sent='completed'AND NOT EXISTS(SELECT*FROM submissions_completed AS s WHEREs.id=t.submission_id))
OR (t.last_status_sent='failed'AND NOT EXISTS(SELECT*FROM submissions_failed AS s WHEREs.id=t.submission_id))
OR (t.last_status_sent='cancelled'AND NOT EXISTS(SELECT*FROM submissions_cancelled AS s WHEREs.id=t.submission_id))
)
SELECT
task_id,
coalesce(
(SELECT'paused'FROM submissions_paused AS s WHEREs.id=t.submission_id),
(SELECT'in_progress'FROM submissions AS s WHEREs.id=t.submission_id),
(SELECT'completed'FROM submissions_completed AS s WHEREs.id=t.submission_id),
(SELECT'failed'FROM submissions_failed AS s WHEREs.id=t.submission_id),
(SELECT'cancelled'FROM submissions_cancelled AS s WHEREs.id=t.submission_id)
) AS"current_status!: DelegatedJobStatus"FROM out_of_date_tasks AS t
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
This PR adds the following:
This 'delegation' extension can be used to delegate the responsibility of starting/stopping submissions to another service, and keeping that other service up-to-date on a submission's status. An example use-case is a job scheduling system that keeps tracks of dependencies between jobs (which are not necessarily Opsqueue submissions), and schedules them accordingly. Enabling cross-system dependencies to be expressed and handled.
See below for a sequence diagram how this works in practice:
The happy path:
sequenceDiagram autonumber participant CL as Client participant D as Delegator participant OD as Opsqueue (Delegation module) participant C as Opsqueue (Core module) participant DB@{"type": "database"} as Database CL ->>+ OD: insert_submission(chunks) OD -->>- CL: submission_id CL ->> D: submit(submission_id) Note over D: Delegator is free to wait. (e.g. to ensure all dependencies are satisfied) D ->>+ OD: delegate(task_id, submission_id) OD ->> DB: insert_external_task(task_id, submission_id) OD -->> C: unpause(submission_id) OD -->>- D: 202 ACCEPTED C -)+ OD: submission_status_changed(submission_id) OD ->>+ DB: select_external_tasks(submission_id) DB -->>- OD: [{task_id, submission_id, last_status_sent}] OD ->>+ C: submission_status(submission_id) C -->>- OD: "IN_PROGRESS" OD ->> D: submit(Updated(task_id, "RUNNING")) OD ->>- DB: update_last_status_sent(task_id, "RUNNING") Note over OD,C: Some time passes until the submission is completed. C -)+ OD: submission_status_changed(submission_id) OD ->>+ DB: select_external_tasks(submission_id) DB -->>- OD: [{task_id, submission_id, last_status_sent}] OD ->>+ C: submission_status(submission_id) C -->>- OD: "COMPLETED" OD ->> D: submit(Completion(task_id, "SUCCESS")) OD ->>- DB: delete_external_task(task_id) Note over D: Delegator can respond to status change (e.g. start processing jobs that depended on this submission)The path when a submission is cancelled mid-way:
sequenceDiagram autonumber participant CL as Client participant D as Delegator participant OD as Opsqueue (Delegation module) participant C as Opsqueue (Core module) participant DB@{"type": "database"} as Database CL ->>+ OD: insert_submission(chunks) OD -->>- CL: submission_id CL ->>+ D: submit(submission_id) D -->>- CL: task_id Note over D: Delegator is free to wait. (e.g. to ensure all dependencies are satisfied) D ->>+ OD: delegate(task_id, submission_id) OD ->> DB: insert_external_task(task_id, submission_id) OD -->> C: unpause(submission_id) OD -->>- D: 202 ACCEPTED C -)+ OD: submission_status_changed(submission_id) OD ->>+ DB: select_external_tasks(submission_id) DB -->>- OD: [{task_id, submission_id, last_status_sent}] OD ->>+ C: submission_status(submission_id) C -->>- OD: "IN_PROGRESS" OD ->> D: update(task_id, "RUNNING") OD ->>- DB: update_last_status_sent(task_id, "RUNNING") Note over OD,C: Some time passes but the submission is not completed yet. CL -) D: kill(task_id) D -) OD: kill(task_id) OD -->> C: cancel(submission_id) C -)+ OD: submission_status_changed(submission_id) OD ->>+ DB: select_external_tasks(submission_id) DB -->>- OD: [{task_id, submission_id, last_status_sent}] OD ->>+ C: submission_status(submission_id) C -->>- OD: "CANCELLED" OD ->> D: submit(Failure(task_id, "FORCED")) OD ->>- DB: delete_external_task(task_id) Note over D: Delegator can respond to status change (e.g. block processing jobs that depend on this submission)The path when an error occurs while running the submission is similar to the path for cancellation, without the client explicitly canceling the submission.