fix(job): resume paused in-memory tasks - #3083
OldFriendWenjianjian wants to merge 4 commits into
Conversation
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Repository UI Review profile: CHILL Plan: Pro Plus Run ID: 📒 Files selected for processing (3)
🚧 Files skipped from review as they are similar to previous changes (3)
WalkthroughThe 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. ChangesJob resumption
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
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches🧪 Generate unit tests (beta)
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. Comment |
There was a problem hiding this comment.
Actionable comments posted: 2
🧹 Nitpick comments (1)
core/src/infra/job/manager.rs (1)
2516-2548: 📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick winAdd 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_rxis dropped without sending (Line 59-67 error path) or wheretask_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
📒 Files selected for processing (1)
core/src/infra/job/manager.rs
There was a problem hiding this comment.
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 winDo not await resume while holding
running_jobswrite lock.
resume_jobacquires a write guard, then awaitsresume_in_memory_task, which can await persistence and retrytask_handle.resume().awaitup to 1 second. At the same time,pause_jobmust acquirerunning_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 publishingJobStatus::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 valueDocument the resume contract.
Add
///documentation toresume_in_memory_task. Use a second paragraph to explain why it consumes the persistence receiver and retriesTaskSystemError::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
📒 Files selected for processing (1)
core/src/infra/job/manager.rs
|
Addressed the remaining review feedback in 6bd8bab.
Validation on Windows with Rust 1.97.1:
The command-scoped lint cap is only for the existing mainline |
There was a problem hiding this comment.
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 winHandle 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 isRunning. A job that was alreadyPausedbefore shutdown therefore produces no new generation, so thischanged()call blocks for the full 10 seconds and thebreakat 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 thebreakwithcontinueso 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 valueDocument the persistence generation contract on the trait method.
persistence_complete_txchanged 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 valueUse
debug!instead ofwarn!for these diagnostic messages.Lines 221, 227, and 237 log routine flow at
warn!level and prefix the text withDEBUG:. The coding guidelines require the correcttracinglevel and structured context fields. Usedebug!withjob_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 fromtracingcrate instead ofprintln!" 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
📒 Files selected for processing (4)
core/src/infra/job/executor.rscore/src/infra/job/manager.rscore/src/infra/job/types.rscrates/job-derive/src/lib.rs
|
Follow-up review fixes are in 7ec0201:
I did not add the suggested Validation:
|
Summary
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 --nocaptureRUSTFLAGS="--cap-lints warn" cargo test -p sd-task-system --test integration_test pause_test -- --exact --nocaptureRUSTFLAGS="--cap-lints warn" cargo check -p sd-core --no-default-featuresThe lint cap is test-command-only because current main forbids a
deprecated_in_futurewarning forAtomicUsize::fetch_updateinsd-task-systemon the installed Rust toolchain.Fixes #2968