From ddeee8918c928e4c493dd6f9bbc4ae82d100e7f7 Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Tue, 18 Aug 2026 22:08:55 +0530 Subject: [PATCH 1/3] Auto-committed by item done: uncommitted changes at completion --- crates/agentflare-backend/src/claim.rs | 24 +++++++ src/mcp_server/item.rs | 40 +++++++++-- src/mcp_server/tests/action_tests.rs | 40 +++++++++++ src/work_item_pipeline.rs | 35 ++++++--- src/work_item_pipeline/tests.rs | 99 ++++++++++++++++++++++++++ 5 files changed, 223 insertions(+), 15 deletions(-) diff --git a/crates/agentflare-backend/src/claim.rs b/crates/agentflare-backend/src/claim.rs index 4352361f..63b9b895 100644 --- a/crates/agentflare-backend/src/claim.rs +++ b/crates/agentflare-backend/src/claim.rs @@ -86,6 +86,30 @@ pub fn current_owner(conn: &Connection, item_id: &str) -> Option { .map(|c| c.owner) } +/// A live (non-stale) `claimed` lease on an item, if any. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct LiveClaimOnItem { + pub owner: String, + pub age_secs: i64, +} + +/// Returns the live claim holder on `item_id`, if one exists. +pub fn live_claim_on_item( + conn: &Connection, + item_id: &str, + now: i64, + ttl_secs: i64, +) -> rusqlite::Result> { + let claims = LEDGER.list(conn, false, now, ttl_secs)?; + Ok(claims + .iter() + .find(|c| c.key == [item_id]) + .map(|c| LiveClaimOnItem { + owner: c.owner.clone(), + age_secs: now - c.heartbeat_at, + })) +} + /// Returns true if there is an active (live, non-stale) claim on this item /// whose owner differs from `owner`. Used by the comment edit/delete gates /// to prevent modifying a comment when another agent has started work. diff --git a/src/mcp_server/item.rs b/src/mcp_server/item.rs index 097005bf..8a065ffa 100644 --- a/src/mcp_server/item.rs +++ b/src/mcp_server/item.rs @@ -1150,6 +1150,8 @@ impl AgentflareMcp { return Err(ErrorData::invalid_params("id is required", None)); } let assignee_agent = req.assignee_agent.clone(); + let now = crate::claims::now(); + let ttl = crate::mcp_server::types::backend_claim_ttl_secs(); self.with_backend_db(|conn| { let item_id = self.resolve_item_id(conn, &raw)?; let outcome = agentflare_backend::item::redispatch( @@ -1159,14 +1161,42 @@ impl AgentflareMcp { ) .map_err(map_backend_err)?; match outcome { - agentflare_backend::item::RedispatchOutcome::Ready { assignee_agent } => Ok( - serde_json::json!({ + agentflare_backend::item::RedispatchOutcome::Ready { assignee_agent } => { + let effective_ttl = + agentflare_backend::claim::effective_ttl_secs(conn, &item_id, ttl); + let dispatchable = !agentflare_backend::claim::has_active_claim_by_other( + conn, + &item_id, + &assignee_agent, + now, + effective_ttl, + ) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + let mut resp = serde_json::json!({ "item_id": item_id, "ready": true, + "dispatchable": dispatchable, "assignee_agent": assignee_agent, - }) - .to_string(), - ), + }); + if !dispatchable { + let live = agentflare_backend::claim::live_claim_on_item( + conn, + &item_id, + now, + effective_ttl, + ) + .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + if let Some(live) = live { + resp["blocked_by_live_claim"] = serde_json::json!({ + "owner": live.owner, + "age_secs": live.age_secs, + "ttl_secs": effective_ttl, + "reason": "another agent already holds a live claim on this item", + }); + } + } + Ok(resp.to_string()) + }, agentflare_backend::item::RedispatchOutcome::NoAssignee => { Err(ErrorData::invalid_params( format!( diff --git a/src/mcp_server/tests/action_tests.rs b/src/mcp_server/tests/action_tests.rs index c8833302..83b63669 100644 --- a/src/mcp_server/tests/action_tests.rs +++ b/src/mcp_server/tests/action_tests.rs @@ -1047,6 +1047,46 @@ fn item_release_errors_when_a_different_owner_holds_a_live_claim() { ); } +#[test] +fn item_redispatch_reports_dispatch_blocked_when_a_live_claim_remains() { + let (_tmp, s) = harness(); + let created: serde_json::Value = + serde_json::from_str(&s.item(Parameters(empty_item_create("Test"))).unwrap()).unwrap(); + let item_id = created["id"].as_str().unwrap().to_string(); + + s.item(Parameters(ItemRequest { + action: "update".into(), + id: Some(item_id.clone()), + assignee_agent: Some("opencode".into()), + ..Default::default() + })) + .unwrap(); + + seed_claim(&s, &item_id, "opencode:dead-job", 60); + + let resp: serde_json::Value = serde_json::from_str( + &s.item(Parameters(ItemRequest { + action: "redispatch".into(), + id: Some(item_id.clone()), + ..Default::default() + })) + .unwrap(), + ) + .unwrap(); + + assert_eq!(resp["ready"], serde_json::Value::Bool(true)); + assert_eq!(resp["dispatchable"], serde_json::Value::Bool(false)); + assert_eq!( + resp["blocked_by_live_claim"]["owner"].as_str(), + Some("opencode:dead-job") + ); + assert_eq!( + resp["blocked_by_live_claim"]["reason"].as_str(), + Some("another agent already holds a live claim on this item") + ); + assert!(resp["blocked_by_live_claim"]["age_secs"].as_i64().unwrap() >= 60); +} + #[test] fn item_release_reclaims_and_releases_a_stale_claim_from_an_abandoned_owner() { let (_tmp, s) = harness(); diff --git a/src/work_item_pipeline.rs b/src/work_item_pipeline.rs index 206ff8cf..b636c548 100644 --- a/src/work_item_pipeline.rs +++ b/src/work_item_pipeline.rs @@ -468,6 +468,8 @@ pub(crate) fn build_sdd_loop_step( /// human with a comment instead of opening a PR on unreviewed code, since /// this step has no access to `supervisor`'s label-id lookups for a real /// relabel (that stays the supervisor's job on its next discovery tick). +/// The job is finished either way — release the claim so redispatch / +/// supervisor discovery can pick the item back up. /// 4. Otherwise — the success path: `item_done`, then the same /// `cap_reply_for_comment`/`format_success_comment`/comment/notify /// sequence `execute_work` runs today. @@ -477,6 +479,24 @@ pub(crate) fn build_sdd_loop_step( /// way `coder`/`review_or_fix`'s agent dispatch can, and unlike those two, /// a failure here has already done the real work and just needs to land the /// result. +/// +/// Best-effort claim release on every terminal success except when +/// `item_done` deliberately left the lease held for an open PR (`in_review`). +fn finalize_release_claim_best_effort( + mcp: &crate::mcp_server::AgentflareMcp, + item_id: &str, + leave_claim_held: bool, +) { + if leave_claim_held { + return; + } + let _ = mcp.item_release(ItemRequest { + action: "release".into(), + id: Some(item_id.to_string()), + ..Default::default() + }); +} + pub(crate) fn build_finalize_step( mcp: std::sync::Arc, item_id: String, @@ -492,11 +512,7 @@ pub(crate) fn build_finalize_step( Box::pin(async move { crate::claims::with_owner_override(owner, || { if let Some(reason) = ctx.data.hold_reason.clone() { - let _ = mcp.item_release(ItemRequest { - action: "release".into(), - id: Some(item_id.clone()), - ..Default::default() - }); + finalize_release_claim_best_effort(&mcp, &item_id, false); let body = format!("## agentflare work — on hold\n\n{reason}"); let _ = mcp.comment_impl(CommentRequest { action: "create".into(), @@ -520,11 +536,7 @@ pub(crate) fn build_finalize_step( } else { ctx.data.review_findings.join("\n\n---\n\n") }; - let _ = mcp.item_release(ItemRequest { - action: "release".into(), - id: Some(item_id.clone()), - ..Default::default() - }); + finalize_release_claim_best_effort(&mcp, &item_id, false); let body = format!("## agentflare work — review findings\n\n{findings}"); let _ = mcp.comment_impl(CommentRequest { action: "create".into(), @@ -540,6 +552,7 @@ pub(crate) fn build_finalize_step( if ctx.data.review_issues.is_some() { let issues = ctx.data.review_issues.clone().unwrap_or_default(); + finalize_release_claim_best_effort(&mcp, &item_id, false); let _ = mcp.comment_impl(CommentRequest { action: "create".into(), item_id: Some(item_id.clone()), @@ -567,6 +580,8 @@ pub(crate) fn build_finalize_step( let done_val: serde_json::Value = serde_json::from_str(&done_resp).unwrap_or(serde_json::Value::Null); ctx.data.pr_url = done_val["pr_url"].as_str().map(str::to_string); + let leave_claim_held = done_val["status"].as_str() == Some("in_review"); + finalize_release_claim_best_effort(&mcp, &item_id, leave_claim_held); let comment_reply = crate::cli::work::cap_reply_for_comment( &mcp, diff --git a/src/work_item_pipeline/tests.rs b/src/work_item_pipeline/tests.rs index 5a124d43..47529e14 100644 --- a/src/work_item_pipeline/tests.rs +++ b/src/work_item_pipeline/tests.rs @@ -292,6 +292,105 @@ async fn finalize_step_uses_accumulated_review_findings_when_last_report_was_cle panic!("finalize step did not complete"); } +#[tokio::test] +async fn finalize_step_releases_claim_when_human_review_gate_is_hit() { + let (mcp, _backend_tmp, _repo_tmp, item_id, _project_id, _worktree_path) = + crate::mcp_server::tests::mcp_with_claimed_item("Human-review finalize test item"); + let mcp = Arc::new(mcp); + let owner = crate::claims::owner_id(); + + let data = WorkItemData { + review_issues: Some("- still broken".into()), + ..Default::default() + }; + let step = build_finalize_step( + mcp.clone(), + item_id.clone(), + None, + owner.clone(), + ); + let wf = WorkflowDefinition::new(WORKFLOW_ID, "work item").add_step(step); + let engine = WorkflowEngine::>::new(); + engine.register_workflow(wf).unwrap(); + let run_id = engine + .start_workflow(WorkflowId::new(WORKFLOW_ID), data, String::new()) + .await + .unwrap(); + + for _ in 0..50 { + let state = engine.get_status(run_id).await.unwrap(); + if state.status == flare_workflow::WorkflowStatus::Completed { + let still_claimed = mcp + .with_backend_db(|conn| { + agentflare_backend::claim::is_owner(conn, &item_id, &owner) + .map_err(|e| e.to_string()) + }) + .unwrap() + .unwrap(); + assert!( + !still_claimed, + "finalize must release the claim when gating for human review" + ); + return; + } + if state.status == flare_workflow::WorkflowStatus::Failed { + panic!("finalize must succeed for human-review gate: {:?}", state.error); + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + panic!("finalize step did not complete"); +} + +#[tokio::test] +async fn finalize_step_releases_claim_after_review_only_success() { + let (mcp, _backend_tmp, _repo_tmp, item_id, _project_id, _worktree_path) = + crate::mcp_server::tests::mcp_with_claimed_item("Review-only claim release test item"); + let mcp = Arc::new(mcp); + let owner = crate::claims::owner_id(); + + let data = WorkItemData { + review_only: true, + last_report: Some("Found an issue.".to_string()), + ..Default::default() + }; + let step = build_finalize_step( + mcp.clone(), + item_id.clone(), + None, + owner.clone(), + ); + let wf = WorkflowDefinition::new(WORKFLOW_ID, "work item").add_step(step); + let engine = WorkflowEngine::>::new(); + engine.register_workflow(wf).unwrap(); + let run_id = engine + .start_workflow(WorkflowId::new(WORKFLOW_ID), data, String::new()) + .await + .unwrap(); + + for _ in 0..50 { + let state = engine.get_status(run_id).await.unwrap(); + if state.status == flare_workflow::WorkflowStatus::Completed { + let still_claimed = mcp + .with_backend_db(|conn| { + agentflare_backend::claim::is_owner(conn, &item_id, &owner) + .map_err(|e| e.to_string()) + }) + .unwrap() + .unwrap(); + assert!( + !still_claimed, + "finalize must release the claim after a review-only run" + ); + return; + } + if state.status == flare_workflow::WorkflowStatus::Failed { + panic!("finalize must not fail for review-only: {:?}", state.error); + } + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + } + panic!("finalize step did not complete"); +} + // requires a real headless agent binary; run manually / in an // environment with one installed — the mock-sender variant right below // covers the same metadata-persistence assertion unconditionally. From 4f3ecdc0bbb073ecd3a3b0ec4218340a4a82fae1 Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Wed, 19 Aug 2026 00:09:55 +0530 Subject: [PATCH 2/3] style: rustfmt finalize claim-release tests Agentflare-Branch: task/508-work-item-pipeline-finalize-doesn-t-rele Agentflare-Item: 508 --- src/work_item_pipeline/tests.rs | 19 ++++++------------- 1 file changed, 6 insertions(+), 13 deletions(-) diff --git a/src/work_item_pipeline/tests.rs b/src/work_item_pipeline/tests.rs index 47529e14..62d14d7b 100644 --- a/src/work_item_pipeline/tests.rs +++ b/src/work_item_pipeline/tests.rs @@ -303,12 +303,7 @@ async fn finalize_step_releases_claim_when_human_review_gate_is_hit() { review_issues: Some("- still broken".into()), ..Default::default() }; - let step = build_finalize_step( - mcp.clone(), - item_id.clone(), - None, - owner.clone(), - ); + let step = build_finalize_step(mcp.clone(), item_id.clone(), None, owner.clone()); let wf = WorkflowDefinition::new(WORKFLOW_ID, "work item").add_step(step); let engine = WorkflowEngine::>::new(); engine.register_workflow(wf).unwrap(); @@ -334,7 +329,10 @@ async fn finalize_step_releases_claim_when_human_review_gate_is_hit() { return; } if state.status == flare_workflow::WorkflowStatus::Failed { - panic!("finalize must succeed for human-review gate: {:?}", state.error); + panic!( + "finalize must succeed for human-review gate: {:?}", + state.error + ); } tokio::time::sleep(std::time::Duration::from_millis(10)).await; } @@ -353,12 +351,7 @@ async fn finalize_step_releases_claim_after_review_only_success() { last_report: Some("Found an issue.".to_string()), ..Default::default() }; - let step = build_finalize_step( - mcp.clone(), - item_id.clone(), - None, - owner.clone(), - ); + let step = build_finalize_step(mcp.clone(), item_id.clone(), None, owner.clone()); let wf = WorkflowDefinition::new(WORKFLOW_ID, "work item").add_step(step); let engine = WorkflowEngine::>::new(); engine.register_workflow(wf).unwrap(); From 87b4f50cff60e92843aadfab32f90d39523de2ef Mon Sep 17 00:00:00 2001 From: Shivakumar Date: Wed, 19 Aug 2026 09:48:31 +0530 Subject: [PATCH 3/3] fix(mcp): dedupe live-claim lookup in item redispatch response CodeRabbit finding on PR #555: has_active_claim_by_other and live_claim_on_item both scanned the same claim ledger for the same item -- one query to compute dispatchable, a second to fetch details for blocked_by_live_claim. Call live_claim_on_item once and derive both from the same snapshot. Agentflare-Agent: claude-code_2-1-234_agent Agentflare-Branch: task/508-work-item-pipeline-finalize-doesn-t-rele Agentflare-Item: 508 --- src/mcp_server/item.rs | 28 ++++++++++------------------ 1 file changed, 10 insertions(+), 18 deletions(-) diff --git a/src/mcp_server/item.rs b/src/mcp_server/item.rs index 8a065ffa..cc24896c 100644 --- a/src/mcp_server/item.rs +++ b/src/mcp_server/item.rs @@ -1164,36 +1164,28 @@ impl AgentflareMcp { agentflare_backend::item::RedispatchOutcome::Ready { assignee_agent } => { let effective_ttl = agentflare_backend::claim::effective_ttl_secs(conn, &item_id, ttl); - let dispatchable = !agentflare_backend::claim::has_active_claim_by_other( + let live = agentflare_backend::claim::live_claim_on_item( conn, &item_id, - &assignee_agent, now, effective_ttl, ) .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; + let blocked_by = live.filter(|c| c.owner != assignee_agent); + let dispatchable = blocked_by.is_none(); let mut resp = serde_json::json!({ "item_id": item_id, "ready": true, "dispatchable": dispatchable, "assignee_agent": assignee_agent, }); - if !dispatchable { - let live = agentflare_backend::claim::live_claim_on_item( - conn, - &item_id, - now, - effective_ttl, - ) - .map_err(|e| ErrorData::internal_error(e.to_string(), None))?; - if let Some(live) = live { - resp["blocked_by_live_claim"] = serde_json::json!({ - "owner": live.owner, - "age_secs": live.age_secs, - "ttl_secs": effective_ttl, - "reason": "another agent already holds a live claim on this item", - }); - } + if let Some(live) = blocked_by { + resp["blocked_by_live_claim"] = serde_json::json!({ + "owner": live.owner, + "age_secs": live.age_secs, + "ttl_secs": effective_ttl, + "reason": "another agent already holds a live claim on this item", + }); } Ok(resp.to_string()) },