Skip to content

Add delegation extension - #187

Draft
SemMulder wants to merge 1 commit into
masterfrom
sm/delegation-2
Draft

SemMulder wants to merge 1 commit into
masterfrom
sm/delegation-2

Conversation

@SemMulder

@SemMulder SemMulder commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

This PR adds the following:

  • 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.

Comment thread opsqueue/src/common/submission.rs
Comment thread opsqueue/src/delegation/server.rs

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Unresolved cleanup data-retention, delegated-reference, and transaction-notification issues block approval.

Get a fresh assessment by requesting another Copilot review.

Review effort: Lite
Findings: 2 High severity · 1 Medium severity

Open (3)
What changed in this PR

Adds optional delegation support with status synchronization, event handling, database tracking, and cleanup integration.

Changes:

  • Adds delegation routes and background reporting.
  • Propagates submission status notifications.
  • Adds external-task migrations and dependency updates.
  • Refactors cleanup and extension APIs.

Review findings:

  • Critical (4 votes): Cleanup can delete submissions still referenced by delegated tasks.
  • Moderate (1 vote): Chunk cleanup leaves completed and failed chunk rows.
  • Moderate (3 votes): Submission cleanup leaves terminal chunk payloads.
  • Critical (2 votes): Delegation notifications may be broadcast before transactions commit.
File Summary
workspace-hack/​Cargo.toml Updates shared dependency features.
opsqueue/​src/​server.rs Registers delegation routes and status channels.
opsqueue/​src/​producer/​server.rs Emits delegation status changes.
opsqueue/​src/​lib.rs Exposes the delegation module.
opsqueue/​src/​delegation/​server.rs Implements delegation APIs and background reporting.
opsqueue/​src/​delegation/​mod.rs Registers the delegation server module.
opsqueue/​src/​db/​mod.rs Adds a test database-pool helper.
opsqueue/​src/​consumer/​server/​mod.rs Emits status changes on chunk completion and failure.
opsqueue/​src/​config.rs Adds delegation server URL configuration.
opsqueue/​src/​common/​submission.rs Adds status notifications and cleanup changes.
opsqueue/​src/​common/​mod.rs Registers extension support.
opsqueue/​src/​common/​extension.rs Adds extension and core APIs.
opsqueue/​src/​common/​chunk.rs Adds status notifications and chunk deletion helpers.
opsqueue/​migrations/​20260811150000_submissions_external_task.up.sql Creates external-task storage.
opsqueue/​migrations/​20260811150000_submissions_external_task.down.sql Removes external-task storage.
opsqueue/​Cargo.toml Adds delegation and test dependencies.
opsqueue/​app/​main.rs Updates cleanup task initialization.
Cargo.toml Bumps the workspace version.
Cargo.lock Updates locked dependencies and package versions.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread opsqueue/app/main.rs Outdated
Comment thread opsqueue/src/delegation/server.rs Outdated
Comment thread opsqueue/src/common/submission.rs
Comment thread opsqueue/src/delegation/server.rs Outdated
@SemMulder
SemMulder force-pushed the sm/delegation-2 branch 2 times, most recently from 7d44741 to fd4e096 Compare September 25, 2026 15:40
@SemMulder
SemMulder requested a balanced review from Copilot September 25, 2026 15:41

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Unimplemented event handling, incorrect terminal-status recovery, and duplicate background workers can produce incorrect external state.

Get a fresh assessment by requesting another Copilot review.

Review effort: Balanced
Findings: 3 High severity · 2 Medium severity

Open (5)
Resolved since last review (3)

Comment thread opsqueue/src/delegation/extension.rs
Comment thread opsqueue/src/delegation/server.rs
Comment thread opsqueue/src/delegation/server.rs
Comment thread opsqueue/src/common/submission.rs Outdated
Comment thread opsqueue/src/common/submission.rs Outdated

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🔵 Needs a closer look

Status reporting can deadlock with a one-connection reader pool and introduces costly periodic N+1 queries.

Review effort: Balanced
Findings: None

Resolved since last review (5)
Previously missed (1)

In code that hasn't changed since last review

Medium severity 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:

  • line 357
  • line 371
  • line 600

@SemMulder

SemMulder commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor Author

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
WHERE
   t.last_status_sent IS NULL
   OR (t.last_status_sent = 'paused' AND NOT EXISTS(SELECT * FROM submissions_paused AS s WHERE s.id = t.submission_id))
   OR (t.last_status_sent = 'in_progress' AND NOT EXISTS(SELECT * FROM submissions AS s WHERE s.id = t.submission_id))
   OR (t.last_status_sent = 'completed' AND NOT EXISTS(SELECT * FROM submissions_completed AS s WHERE s.id = t.submission_id))
   OR (t.last_status_sent = 'failed' AND NOT EXISTS(SELECT * FROM submissions_failed AS s WHERE s.id = t.submission_id))
   OR (t.last_status_sent = 'cancelled' AND NOT EXISTS(SELECT * FROM submissions_cancelled AS s WHERE s.id = t.submission_id))
)
SELECT
   task_id,
   coalesce(
    (SELECT 'paused' FROM submissions_paused AS s WHERE s.id = t.submission_id),
    (SELECT 'in_progress' FROM submissions AS s WHERE s.id = t.submission_id),
    (SELECT 'completed' FROM submissions_completed AS s WHERE s.id = t.submission_id),
    (SELECT 'failed' FROM submissions_failed AS s WHERE s.id = t.submission_id),
    (SELECT 'cancelled' FROM submissions_cancelled AS s WHERE s.id = t.submission_id)
   ) AS "current_status!: DelegatedJobStatus"
FROM out_of_date_tasks AS t

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.

2 participants