Skip to content

fix(job): resume paused in-memory tasks - #3083

Open
OldFriendWenjianjian wants to merge 4 commits into
spacedriveapp:mainfrom
OldFriendWenjianjian:fix/job-resume-task
Open

OldFriendWenjianjian wants to merge 4 commits into
spacedriveapp:mainfrom
OldFriendWenjianjian:fix/job-resume-task

Conversation

@OldFriendWenjianjian

Copy link
Copy Markdown

Summary

  • resume the task-system handle for paused jobs that remain in memory
  • wait for paused state persistence before re-queueing the same executor
  • expose the job as running only after the task system accepts the resume request

Testing

  • RUSTFLAGS="--cap-lints warn" cargo test -p sd-core --lib infra::job::manager::tests::in_memory_resume_waits_for_persistence_and_resumes_task --no-default-features -- --exact --nocapture
  • repeated the focused race test 20 times (20/20 passed)
  • RUSTFLAGS="--cap-lints warn" cargo test -p sd-task-system --test integration_test pause_test -- --exact --nocapture
  • RUSTFLAGS="--cap-lints warn" cargo check -p sd-core --no-default-features

The lint cap is test-command-only because current main forbids a deprecated_in_future warning for AtomicUsize::fetch_update in sd-task-system on the installed Rust toolchain.

Fixes #2968

@coderabbitai

coderabbitai Bot commented Aug 3, 2026 •

Copy link
Copy Markdown

Review Change Stack

No actionable comments were generated in the recent review. 🎉

ℹ️ Recent review info
⚙️ Run configuration

Configuration used: Repository UI

Review profile: CHILL

Plan: Pro Plus

Run ID: 6b2330f4-f2f1-46e1-aed7-873a0dcfd5c8

📥 Commits

Reviewing files that changed from the base of the PR and between 6bd8bab and 7ec0201.

📒 Files selected for processing (3)
  • core/src/infra/job/executor.rs
  • core/src/infra/job/manager.rs
  • core/src/infra/job/types.rs
🚧 Files skipped from review as they are similar to previous changes (3)
  • core/src/infra/job/types.rs
  • core/src/infra/job/executor.rs
  • core/src/infra/job/manager.rs

Walkthrough

The job system now uses generation-based watch channels for persistence completion. The manager waits for paused-state persistence, retries task resumption, prevents concurrent resumes, validates task state, and updates status only after successful resumption.

Changes

Job resumption

Layer / File(s) Summary
Persistence signal contract
core/src/infra/job/types.rs, crates/job-derive/src/lib.rs, core/src/infra/job/executor.rs
Executor creation now uses reusable watch::Sender<u64> persistence signals. The counter advances once for each settled execution outcome.
Executor persistence lifecycle
core/src/infra/job/executor.rs
Execution preserves paused status. Completed, canceled, and failed jobs signal persistence. Paused jobs defer signaling until paused-state persistence completes.
Persistence-aware task resume
core/src/infra/job/manager.rs
The manager waits for persistence, retries TaskNotFound, prevents concurrent resumes, validates task identity and paused state, and updates job state after successful resumption. Shutdown continues after persistence timeouts. Tests cover blocking, retries, closed channels, and task-system errors.

Estimated code review effort: 4 (Complex) | ~45 minutes

Sequence Diagram(s)

sequenceDiagram
  participant JobExecutor
  participant PersistenceWatch
  participant JobManager
  participant TaskController
  participant JobDatabase
  JobExecutor->>PersistenceWatch: publish persistence generation
  JobManager->>PersistenceWatch: wait for generation change
  PersistenceWatch-->>JobManager: paused-state persistence completed
  JobManager->>TaskController: resume paused task
  TaskController-->>JobManager: success or TaskNotFound
  JobManager->>TaskController: retry TaskNotFound
  JobManager->>JobDatabase: update status after successful resume
Loading

Poem

A rabbit watched the counter glow,
And waited for persistence to show.
A missing task was tried once more,
Then status changed as tasks restored.
Hop—the paused job ran.

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 77.78% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Title check ✅ Passed The title clearly describes the main change: resuming paused in-memory tasks.
Description check ✅ Passed The description covers the change, testing, and linked issue, with sufficient implementation and validation details.
Linked Issues check ✅ Passed The changes address issue #2968 by resuming paused in-memory task handles after persisted paused state is confirmed.
Out of Scope Changes check ✅ Passed The persistence signaling, retry, locking, logging, and tests directly support reliable paused-job resumption.
✨ Finishing Touches
🧪 Generate unit tests (beta)
  • Create PR with unit tests

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

Comment thread core/src/infra/job/manager.rs Outdated

@coderabbitai coderabbitai Bot 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.

Actionable comments posted: 2

🧹 Nitpick comments (1)
core/src/infra/job/manager.rs (1)

2516-2548: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win

Add coverage for the error paths of resume_in_memory_task.

The test verifies only the happy path: the persistence signal fires and resume proceeds. It does not exercise the branch where persistence_complete_rx is dropped without sending (Line 59-67 error path) or where task_handle.resume() fails (Line 70-73). Adding tests for these branches would guard the new error handling against regressions.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@core/src/infra/job/manager.rs` around lines 2516 - 2548, Add focused tests
for the error branches of resume_in_memory_task: drop persistence_complete_rx
without sending and assert the persistence-wait error is returned, and use a
task handle whose resume operation fails to assert that error is propagated.
Keep the existing in_memory_resume_waits_for_persistence_and_resumes_task
happy-path test unchanged.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@core/src/infra/job/manager.rs`:
- Around line 51-75: Update resume_in_memory_task to bound
persistence_complete_rx waiting with the existing persistence_timeout used by
shutdown. Preserve the current JobError mapping for timeout or channel failure,
and ensure the resume path proceeds only when persistence completes within that
configured duration.
- Around line 56-68: Re-arm the pause persistence signal in pause_job() before
each pause/resume cycle by creating a fresh oneshot sender/receiver pair and
storing the receiver in RunningJob.persistence_complete_rx. Ensure the sender is
passed to the pause persistence workflow, while resume_in_memory_task()
continues consuming and awaiting the stored receiver so repeated pauses do not
race resumed execution.

---

Nitpick comments:
In `@core/src/infra/job/manager.rs`:
- Around line 2516-2548: Add focused tests for the error branches of
resume_in_memory_task: drop persistence_complete_rx without sending and assert
the persistence-wait error is returned, and use a task handle whose resume
operation fails to assert that error is propagated. Keep the existing
in_memory_resume_waits_for_persistence_and_resumes_task happy-path test
unchanged.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository UI

Review profile: CHILL

Plan: Pro Plus

Run ID: e611a57c-43ee-41b8-9fd6-8f0b604d2de1

📥 Commits

Reviewing files that changed from the base of the PR and between 6dfeccf and 5fe5b80.

📒 Files selected for processing (1)
  • core/src/infra/job/manager.rs

Comment thread core/src/infra/job/manager.rs
Comment thread core/src/infra/job/manager.rs Outdated

@coderabbitai coderabbitai Bot 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.

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
core/src/infra/job/manager.rs (1)

2198-2212: 🩺 Stability & Availability | 🟠 Major | ⚡ Quick win

Do not await resume while holding running_jobs write lock.

resume_job acquires a write guard, then awaits resume_in_memory_task, which can await persistence and retry task_handle.resume().await up to 1 second. At the same time, pause_job must acquire running_jobs.read().await, so a single resuming job can block pause operations while this lock is held. Hold the write guard only long enough to mark this entry as resuming, then drop it before awaiting resume/retries. Reacquire and validate the job state before publishing JobStatus::Running.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@core/src/infra/job/manager.rs` around lines 2198 - 2212, Update resume_job
around running_jobs and resume_in_memory_task so the write lock is held only
while marking the matching entry as resuming, then release it before awaiting
resume or retry operations. After resume_in_memory_task completes, reacquire the
lock, validate that the same job is still in the expected resuming state, and
only then publish JobStatus::Running.
🧹 Nitpick comments (1)
core/src/infra/job/manager.rs (1)

51-57: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Document the resume contract.

Add /// documentation to resume_in_memory_task. Use a second paragraph to explain why it consumes the persistence receiver and retries TaskSystemError::TaskNotFound.

Proposed documentation
+/// Resume an in-memory task after its paused state is persisted.
+///
+/// The helper consumes the pause-persistence receiver before it resumes the task.
+/// It retries `TaskSystemError::TaskNotFound` because task-system registration can
+/// complete after persistence.
 async fn resume_in_memory_task(

As per coding guidelines, “Function documentation should explain design rationale and non-obvious behavior in the second paragraph, not just restate the code”.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@core/src/infra/job/manager.rs` around lines 51 - 57, Add Rustdoc comments to
resume_in_memory_task with a first paragraph documenting its resume contract and
a second paragraph explaining that it consumes persistence_complete_rx to avoid
racing serialized state, then retries TaskSystemError::TaskNotFound because the
task may not yet be visible to the executor.

Source: Coding guidelines

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Outside diff comments:
In `@core/src/infra/job/manager.rs`:
- Around line 2198-2212: Update resume_job around running_jobs and
resume_in_memory_task so the write lock is held only while marking the matching
entry as resuming, then release it before awaiting resume or retry operations.
After resume_in_memory_task completes, reacquire the lock, validate that the
same job is still in the expected resuming state, and only then publish
JobStatus::Running.

---

Nitpick comments:
In `@core/src/infra/job/manager.rs`:
- Around line 51-57: Add Rustdoc comments to resume_in_memory_task with a first
paragraph documenting its resume contract and a second paragraph explaining that
it consumes persistence_complete_rx to avoid racing serialized state, then
retries TaskSystemError::TaskNotFound because the task may not yet be visible to
the executor.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository UI

Review profile: CHILL

Plan: Pro Plus

Run ID: 94c43e45-cadb-4e4e-b54a-a310d2fc0b78

📥 Commits

Reviewing files that changed from the base of the PR and between 5fe5b80 and b3bf0ef.

📒 Files selected for processing (1)
  • core/src/infra/job/manager.rs

@OldFriendWenjianjian

Copy link
Copy Markdown
Author

Addressed the remaining review feedback in 6bd8bab.

  • Replaced the one-shot persistence notification with a generation-based watch channel so every pause/resume cycle waits for its own persisted state.
  • Kept the 10-second persistence timeout and bounded TaskNotFound retry.
  • Dropped the running_jobs write lock before waiting or retrying, then reacquired it and verified the same task is still resuming before publishing Running.
  • Added the requested contract documentation and error-path coverage.

Validation on Windows with Rust 1.97.1:

  • RUSTFLAGS="--cap-lints warn" cargo test -p sd-core in_memory_resume --lib — 3 passed, 0 failed
  • RUSTFLAGS="--cap-lints warn" cargo check -p sd-core --lib — passed
  • targeted rustfmt --edition 2021 --check on all four changed Rust files — passed

The command-scoped lint cap is only for the existing mainline sd-task-system future-deprecation lint; no unrelated source was changed.

@coderabbitai coderabbitai Bot 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.

Actionable comments posted: 1

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
core/src/infra/job/manager.rs (1)

2416-2434: 🩺 Stability & Availability | 🟡 Minor | ⚡ Quick win

Handle receivers that have no pending generation during shutdown.

rx.changed() waits for a generation newer than the one the receiver last observed. The pause loop at Lines 2318-2330 only pauses jobs whose status is Running. A job that was already Paused before shutdown therefore produces no new generation, so this changed() call blocks for the full 10 seconds and the break at Line 2431 then skips every remaining receiver.

Check rx.has_changed() first, and continue to the next receiver when no new generation is pending. Consider replacing the break with continue so that one slow job does not skip the rest.

🛠️ Proposed change
 		for (job_id, mut rx) in persistence_receivers {
+			if !matches!(rx.has_changed(), Ok(true)) {
+				debug!(job_id = %job_id, "No pending persistence generation; skipping wait");
+				continue;
+			}
 			tokio::select! {
 				result = rx.changed() => {
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@core/src/infra/job/manager.rs` around lines 2416 - 2434, Update the
persistence receiver loop around rx.changed() to check rx.has_changed() before
waiting, and continue immediately when no newer generation is pending so
already-paused jobs do not block shutdown. Replace the timeout branch’s break
with continue so a slow receiver does not prevent processing the remaining
persistence_receivers.
🧹 Nitpick comments (2)
core/src/infra/job/types.rs (1)

150-150: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Document the persistence generation contract on the trait method.

persistence_complete_tx changed from a one-shot signal to a monotonic generation counter. Implementers must now increment the value once per persisted outcome, and must not increment twice for one pause. Add a /// note that states this rule, because external implementers such as WASM jobs cannot infer it from the type.

As per coding guidelines: "Use //! for module documentation, /// for public items, and include examples".

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@core/src/infra/job/types.rs` at line 150, Add a public-item doc comment on
the trait method associated with persistence_complete_tx explaining that it is a
monotonic generation counter, must increment exactly once for each persisted
outcome, and must not increment twice for a single pause; include a concise
example as required by the documentation guidelines.

Source: Coding guidelines

core/src/infra/job/executor.rs (1)

218-241: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low value

Use debug! instead of warn! for these diagnostic messages.

Lines 221, 227, and 237 log routine flow at warn! level and prefix the text with DEBUG:. The coding guidelines require the correct tracing level and structured context fields. Use debug! with job_id = %self.state.job_id.

♻️ Proposed change
-		if *self.state.status_tx.borrow() != JobStatus::Paused {
-			warn!(
-				"DEBUG: JobExecutor setting status to Running for job {}",
-				self.state.job_id
-			);
+		if *self.state.status_tx.borrow() != JobStatus::Paused {
+			debug!(job_id = %self.state.job_id, "Publishing Running status from executor");
 			let _ = self.state.status_tx.send(super::types::JobStatus::Running);
-
-			warn!(
-				"DEBUG: JobExecutor updating database status to Running for job {}",
-				self.state.job_id
-			);
 			if let Err(e) = self
 				.update_job_status_in_db(super::types::JobStatus::Running)
 				.await
 			{
 				error!("Failed to update job status in database: {}", e);
-			} else {
-				warn!(
-					"DEBUG: JobExecutor successfully updated database status to Running for job {}",
-					self.state.job_id
-				);
 			}
 		}

As per coding guidelines: "Use info!, warn!, error!, debug! macros from tracing crate instead of println!" and "Include relevant context fields in structured logging, e.g., debug!(job_id = %self.id, "message")".

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@core/src/infra/job/executor.rs` around lines 218 - 241, The diagnostic flow
logs in the resume handling block should use the debug level and structured job
context. Replace the three routine `warn!` calls around status publication and
database update in the `JobExecutor` logic with `debug!`, remove the `DEBUG:`
prefixes, and provide `job_id = %self.state.job_id` as structured context while
preserving their existing messages.

Source: Coding guidelines

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@core/src/infra/job/manager.rs`:
- Around line 2282-2297: In the resume path after publishing JobStatus::Running,
update the database write handling around job_model.update in the job manager to
log database failures instead of propagating them with ?. Align it with the
existing failure-handling behavior near the equivalent resume logic, while
preserving the already-accepted running state and JobResumed flow.

---

Outside diff comments:
In `@core/src/infra/job/manager.rs`:
- Around line 2416-2434: Update the persistence receiver loop around
rx.changed() to check rx.has_changed() before waiting, and continue immediately
when no newer generation is pending so already-paused jobs do not block
shutdown. Replace the timeout branch’s break with continue so a slow receiver
does not prevent processing the remaining persistence_receivers.

---

Nitpick comments:
In `@core/src/infra/job/executor.rs`:
- Around line 218-241: The diagnostic flow logs in the resume handling block
should use the debug level and structured job context. Replace the three routine
`warn!` calls around status publication and database update in the `JobExecutor`
logic with `debug!`, remove the `DEBUG:` prefixes, and provide `job_id =
%self.state.job_id` as structured context while preserving their existing
messages.

In `@core/src/infra/job/types.rs`:
- Line 150: Add a public-item doc comment on the trait method associated with
persistence_complete_tx explaining that it is a monotonic generation counter,
must increment exactly once for each persisted outcome, and must not increment
twice for a single pause; include a concise example as required by the
documentation guidelines.
🪄 Autofix (Beta)

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository UI

Review profile: CHILL

Plan: Pro Plus

Run ID: 6f2069d7-9f18-460d-9f91-cc895491e15e

📥 Commits

Reviewing files that changed from the base of the PR and between b3bf0ef and 6bd8bab.

📒 Files selected for processing (4)
  • core/src/infra/job/executor.rs
  • core/src/infra/job/manager.rs
  • core/src/infra/job/types.rs
  • crates/job-derive/src/lib.rs

Comment thread core/src/infra/job/manager.rs
@OldFriendWenjianjian

Copy link
Copy Markdown
Author

Follow-up review fixes are in 7ec0201:

  • A database write failure after the task has already resumed is now logged, while the accepted resume and JobResumed event continue.
  • A persistence timeout during shutdown now continues checking remaining jobs instead of skipping them.
  • Routine executor flow uses structured debug! logging.
  • ErasedJob::create_executor documents the monotonic persistence-generation contract.

I did not add the suggested has_changed() == false early skip. pause_job publishes Paused before the executor finishes serializing and writing state, so no pending generation can mean persistence is still in flight. Skipping in that state could close the database during the save. The code now documents why it waits for changed() and still bounds each wait.

Validation:

  • RUSTFLAGS="--cap-lints warn" cargo test -p sd-core in_memory_resume --lib — 3 passed, 0 failed
  • targeted rustfmt check on all four PR files — passed

This branch has not been deployed

No deployments
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.

Unpausing index job not resuming

1 participant