Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -197,8 +197,8 @@ When modifying a module with a spec, read the spec first. Code follows spec; spe

```
Pending -> InProgress -> Completed -> Submitted -> Accepted
\-> Failed
\-> Stuck -> InProgress (recovery)
\ \-> Failed
\-> Failed \-> Stuck -> InProgress (recovery)
\-> Failed
```

Expand Down
2 changes: 1 addition & 1 deletion channels-src/telegram/Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

37 changes: 27 additions & 10 deletions src/agent/agent_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -482,6 +482,11 @@ impl Agent {
let repair_channels = self.channels.clone();
let repair_owner_id = self.owner_id().to_string();
let repair_handle = tokio::spawn(async move {
// Track jobs that have already been escalated to ManualRequired
// to prevent sending duplicate notifications every repair cycle.
let mut notified_manual: std::collections::HashSet<uuid::Uuid> =
std::collections::HashSet::new();

loop {
tokio::time::sleep(repair_interval).await;

Expand All @@ -502,19 +507,31 @@ impl Agent {
}
Ok(RepairResult::Failed { message }) => {
tracing::error!("Repair failed: {}", message);
Some(format!(
"Job {} was stuck for {}s, recovery failed permanently: {}",
job.job_id,
job.stuck_duration.as_secs(),
message
))
// Dedup: only notify once per job (same pattern as ManualRequired)
if notified_manual.insert(job.job_id) {
Some(format!(
"Job {} was stuck for {}s, recovery failed permanently: {}",
job.job_id,
job.stuck_duration.as_secs(),
message
))
} else {
None
}
}
Ok(RepairResult::ManualRequired { message }) => {
tracing::warn!("Manual intervention needed: {}", message);
Some(format!(
"Job {} needs manual intervention: {}",
job.job_id, message
))
// Only notify once per job to prevent notification spam.
// The job should have been transitioned to Failed by
// repair_stuck_job, but guard against that failing too.
if notified_manual.insert(job.job_id) {
Some(format!(
"Job {} needs manual intervention: {}",
job.job_id, message
))
} else {
None
}
}
Ok(RepairResult::Retry { message }) => {
tracing::warn!("Repair needs retry: {}", message);
Expand Down
63 changes: 60 additions & 3 deletions src/agent/self_repair.rs
Original file line number Diff line number Diff line change
Expand Up @@ -201,10 +201,44 @@ impl SelfRepair for DefaultSelfRepair {
async fn repair_stuck_job(&self, job: &StuckJob) -> Result<RepairResult, RepairError> {
// Check if we've exceeded max repair attempts
if job.repair_attempts >= self.max_repair_attempts {
// Transition to Failed so detect_stuck_jobs() stops finding this job.
// Without this, the repair loop re-detects the job every cycle and
// sends a ManualRequired notification each time (notification spam).
// update_context returns Result<Result<(), String>, JobError>.
// Outer Err = job not found. Inner Err = invalid state transition.
// Both mean the job was NOT transitioned to Failed.
let transition_ok = matches!(
self.context_manager
.update_context(job.job_id, |ctx| {
ctx.transition_to(
JobState::Failed,
Some(format!(
"exceeded max repair attempts ({})",
self.max_repair_attempts
)),
)
})
.await,
Ok(Ok(()))
);

if !transition_ok {
tracing::error!(
job = %job.job_id,
"Failed to transition job to Failed state after exceeding max repair attempts"
);
}

let status = if transition_ok {
"and has been marked failed"
} else {
"but could not be marked failed (will be suppressed by dedup)"
};

return Ok(RepairResult::ManualRequired {
message: format!(
"Job {} has exceeded maximum repair attempts ({})",
job.job_id, self.max_repair_attempts
"Job {} has exceeded maximum repair attempts ({}) {}",
job.job_id, self.max_repair_attempts, status
),
});
}
Expand Down Expand Up @@ -524,7 +558,19 @@ mod tests {
let cm = Arc::new(ContextManager::new(10));
let job_id = cm.create_job("Unrepairable", "desc").await.unwrap();

let repair = DefaultSelfRepair::new(cm, Duration::from_secs(60), 2);
// Transition through the production path: Pending → InProgress → Stuck
cm.update_context(job_id, |ctx| ctx.transition_to(JobState::InProgress, None))
.await
.unwrap()
.unwrap();
cm.update_context(job_id, |ctx| {
ctx.transition_to(JobState::Stuck, Some("test".into()))
})
.await
.unwrap()
.unwrap();

let repair = DefaultSelfRepair::new(cm.clone(), Duration::from_secs(60), 2);

let stuck_job = StuckJob {
job_id,
Expand All @@ -540,6 +586,17 @@ mod tests {
"Expected ManualRequired, got: {:?}",
result
);

// Regression: the job must be transitioned to Failed so
// detect_stuck_jobs() stops finding it. Without this, the repair
// loop re-detects the job every cycle and sends ManualRequired
// notifications forever (notification spam bug).
let ctx = cm.get_context(job_id).await.unwrap();
assert_eq!(
ctx.state,
JobState::Failed,
"Job should be Failed after exceeding max repair attempts"
);
}

#[tokio::test]
Expand Down
5 changes: 3 additions & 2 deletions src/context/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,8 +59,9 @@ impl JobState {

matches!(
(self, target),
// From Pending
(Pending, InProgress) | (Pending, Cancelled) |
// From Pending (Failed added for self-repair: stuck Pending jobs
// that exhaust repair attempts must be terminable)
(Pending, InProgress) | (Pending, Failed) | (Pending, Cancelled) |
// From InProgress
(InProgress, Completed) | (InProgress, Failed) |
(InProgress, Stuck) | (InProgress, Cancelled) |
Expand Down
Loading