use std::collections::HashMap;
use std::path::Path;
use std::process::Command;
use std::sync::Arc;
use once_cell::sync::Lazy;
use time::OffsetDateTime;
use tokio::sync::{Mutex, OwnedMutexGuard};
use crate::engine::git;
use crate::engine::worktrees::main_repo_root;
use crate::lfd::attention::{
attention_id_for_queue_block, queue_block_attention_item_from_existing,
};
use crate::lfd::config::GitHubConfig;
use crate::lfd::events::EventHub;
use crate::lfd::id::LfdId;
use crate::lfd::live_pr::{build_live_pr_snapshot, run_live_pr_key, LivePrSnapshot};
use crate::lfd::types::{
AttentionStatus, Event, LivePrState, LivePullRequestState, QueueBlock, QueueBlockReason, Run,
RunStackStatus,
};
use crate::lfdb::SharedStore;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QueueRole {
Ready,
Draft,
Blocked,
Merged,
Superseded,
}
impl QueueRole {
pub fn as_str(self) -> &'static str {
match self {
Self::Ready => "ready",
Self::Draft => "draft",
Self::Blocked => "blocked",
Self::Merged => "merged",
Self::Superseded => "superseded",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum QueueNextAction {
OpenPr,
ResolveConflict,
CombinePrs,
AwaitMerge,
}
impl QueueNextAction {
pub fn as_str(self) -> &'static str {
match self {
Self::OpenPr => "open_pr",
Self::ResolveConflict => "resolve_conflict",
Self::CombinePrs => "combine_prs",
Self::AwaitMerge => "await_merge",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueueRunView {
pub role: QueueRole,
pub block_reason: Option<QueueBlockReason>,
pub blocked_at: Option<OffsetDateTime>,
pub next_action: QueueNextAction,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct QueueRebaseConflict {
pub files: Vec<String>,
}
pub(crate) trait QueueOps: Send + Sync {
fn ensure_branch_checked_out(&self, worktree: &Path, branch: &str) -> Result<(), String>;
fn mark_ready(&self, worktree: &Path, pr_number: u32) -> Result<(), String>;
fn mark_draft(&self, worktree: &Path, pr_number: u32) -> Result<(), String>;
fn rebase_onto_default(
&self,
worktree: &Path,
default_branch: &str,
) -> Result<(), QueueRebaseConflict>;
fn scratch_clean(&self, worktree: &Path) -> Result<bool, String>;
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct RealQueueOps;
static QUEUE_RECONCILE_LOCKS: Lazy<Mutex<HashMap<String, Arc<Mutex<()>>>>> =
Lazy::new(|| Mutex::new(HashMap::new()));
impl QueueOps for RealQueueOps {
fn ensure_branch_checked_out(&self, worktree: &Path, branch: &str) -> Result<(), String> {
git::checkout(worktree, branch).map_err(|err| err.to_string())
}
fn mark_ready(&self, worktree: &Path, _pr_number: u32) -> Result<(), String> {
crate::ops::mark_ready(worktree).map_err(|err| err.to_string())
}
fn mark_draft(&self, worktree: &Path, pr_number: u32) -> Result<(), String> {
let output = Command::new("gh")
.arg("pr")
.arg("ready")
.arg("--undo")
.arg(pr_number.to_string())
.current_dir(worktree)
.output()
.map_err(|err| err.to_string())?;
if output.status.success() {
return Ok(());
}
let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string();
if stderr.to_ascii_lowercase().contains("already a draft") {
return Ok(());
}
Err(if stderr.is_empty() {
"failed to mark PR draft".to_string()
} else {
stderr
})
}
fn rebase_onto_default(
&self,
worktree: &Path,
default_branch: &str,
) -> Result<(), QueueRebaseConflict> {
let main_repo = main_repo_root(worktree).unwrap_or_else(|_| worktree.to_path_buf());
git::fetch(&main_repo, "origin", default_branch).map_err(|err| QueueRebaseConflict {
files: vec![err.to_string()],
})?;
let rebase_result = git::rebase(worktree, &format!("origin/{default_branch}"), None)
.map_err(|err| QueueRebaseConflict {
files: vec![err.to_string()],
})?;
if !rebase_result.success {
return Err(QueueRebaseConflict {
files: rebase_result
.conflicts
.unwrap_or_default()
.into_iter()
.map(|path| path.to_string_lossy().to_string())
.collect(),
});
}
git::push(worktree, true).map_err(|err| QueueRebaseConflict {
files: vec![err.to_string()],
})?;
Ok(())
}
fn scratch_clean(&self, worktree: &Path) -> Result<bool, String> {
let output = Command::new("git")
.arg("-C")
.arg(worktree)
.args(["status", "--porcelain", "--", "scratch/"])
.output()
.map_err(|err| err.to_string())?;
if !output.status.success() {
return Err(String::from_utf8_lossy(&output.stderr).trim().to_string());
}
Ok(String::from_utf8_lossy(&output.stdout).trim().is_empty())
}
}
pub async fn remove_reconcile_lock(wave_id: &LfdId) {
let mut locks = QUEUE_RECONCILE_LOCKS.lock().await;
locks.remove(&wave_id.to_string());
}
pub(crate) async fn acquire_reconcile_lock(wave_id: &LfdId) -> OwnedMutexGuard<()> {
let wave_key = wave_id.to_string();
let lock = {
let mut locks = QUEUE_RECONCILE_LOCKS.lock().await;
locks
.entry(wave_key)
.or_insert_with(|| Arc::new(Mutex::new(())))
.clone()
};
lock.lock_owned().await
}
pub(crate) async fn reconcile_wave_queue_with_ops(
store: &SharedStore,
github_config: &GitHubConfig,
wave_id: &LfdId,
ops: &dyn QueueOps,
event_hub: Option<&EventHub>,
) -> Result<(), String> {
let wave_id_for_log = wave_id.clone();
let mut runs = store
.list_stack_runs(wave_id)
.await
.map_err(|err| format!("list_stack_runs failed: {err}"))?;
if runs.is_empty() {
return Ok(());
}
let mut live_snapshot = build_live_pr_snapshot(store, github_config, &runs)
.await
.map_err(|err| format!("build_live_pr_snapshot failed: {err}"))?;
let mut status_changed = false;
for run in &mut runs {
let Some(live_state) = live_snapshot.state_for_run(run) else {
continue;
};
let inferred = inferred_stack_status(run.stack_status, Some(live_state));
if inferred != run.stack_status {
run.stack_status = inferred;
store
.update_run(run)
.await
.map_err(|err| format!("update_run failed: {err}"))?;
status_changed = true;
}
}
if status_changed {
runs = store
.list_stack_runs(wave_id)
.await
.map_err(|err| format!("list_stack_runs refresh failed: {err}"))?;
live_snapshot = build_live_pr_snapshot(store, github_config, &runs)
.await
.map_err(|err| format!("build_live_pr_snapshot refresh failed: {err}"))?;
}
let head_index = find_queue_head_index(&runs, &live_snapshot);
let Some(head_index) = head_index else {
tracing::debug!(wave_id = %wave_id_for_log, "queue reconcile: no active queue head");
return Ok(());
};
for (index, run) in runs.iter().enumerate() {
if index == head_index {
continue;
}
let Some(pr_number) = pr_number(run) else {
continue;
};
let Some(state) = live_snapshot.state_for_run(run) else {
continue;
};
if state.state == LivePrState::Open
&& !state.is_draft
&& ops
.ensure_branch_checked_out(Path::new(&run.worktree), &run.branch)
.and_then(|_| ops.mark_draft(Path::new(&run.worktree), pr_number))
.is_ok()
{
let mut updated = state.clone();
updated.is_draft = true;
updated.synced_at = OffsetDateTime::now_utc();
let _ = store.upsert_live_pr_state(&updated).await;
if let Some(key) = run_live_pr_key(run) {
live_snapshot.live_states.insert(key, updated);
}
}
}
let head = runs[head_index].clone();
let Some(head_pr_number) = pr_number(&head) else {
set_queue_block(
store,
&head,
QueueBlockReason::MissingPr,
Vec::new(),
None,
event_hub,
)
.await?;
return Ok(());
};
let Some(head_live_state) = live_snapshot.state_for_run(&head) else {
set_queue_block(
store,
&head,
QueueBlockReason::MissingPr,
Vec::new(),
None,
event_hub,
)
.await?;
return Ok(());
};
if head_live_state.state != LivePrState::Open {
return Ok(());
}
if let Some(active_run) = store
.get_active_run(wave_id)
.await
.map_err(|err| format!("get_active_run failed: {err}"))?
{
if active_run.id != head.id {
set_queue_block(
store,
&head,
QueueBlockReason::WaveRunning,
Vec::new(),
None,
event_hub,
)
.await?;
return Ok(());
}
}
let worktree = Path::new(&head.worktree);
if !ops.scratch_clean(worktree)? {
set_queue_block(
store,
&head,
QueueBlockReason::ScratchDirty,
Vec::new(),
None,
event_hub,
)
.await?;
return Ok(());
}
if head.stack_position > 0 {
let main_repo = main_repo_root(worktree).unwrap_or_else(|_| worktree.to_path_buf());
let default_branch =
git::get_default_branch(&main_repo).unwrap_or_else(|_| "main".to_string());
if let Err(conflict) = ops
.ensure_branch_checked_out(worktree, &head.branch)
.map_err(|err| QueueRebaseConflict { files: vec![err] })
.and_then(|_| ops.rebase_onto_default(worktree, &default_branch))
{
set_queue_block(
store,
&head,
QueueBlockReason::RebaseConflict,
conflict.files.clone(),
Some("lazy rebase failed".to_string()),
event_hub,
)
.await?;
return Ok(());
}
}
if head_live_state.is_draft {
ops.ensure_branch_checked_out(worktree, &head.branch)?;
if let Err(err) = ops.mark_ready(worktree, head_pr_number) {
set_queue_block(
store,
&head,
QueueBlockReason::PromotionFailed,
Vec::new(),
Some(err),
event_hub,
)
.await?;
return Ok(());
}
let mut promoted_state = head_live_state.clone();
promoted_state.is_draft = false;
promoted_state.synced_at = OffsetDateTime::now_utc();
let _ = store.upsert_live_pr_state(&promoted_state).await;
}
clear_queue_block(store, wave_id, &head.id, event_hub).await?;
Ok(())
}
pub fn project_queue_views<F>(
runs: &[Run],
mut live_state_for: F,
blocks: &HashMap<LfdId, QueueBlock>,
) -> HashMap<LfdId, QueueRunView>
where
F: FnMut(&Run) -> Option<LivePullRequestState>,
{
let mut live_by_run = HashMap::new();
for run in runs {
live_by_run.insert(run.id.clone(), live_state_for(run));
}
let head_index = runs.iter().position(|run| {
let live = live_by_run.get(&run.id).and_then(|value| value.as_ref());
inferred_stack_status(run.stack_status, live) == RunStackStatus::Active
});
let mut result = HashMap::with_capacity(runs.len());
for (index, run) in runs.iter().enumerate() {
let block = blocks.get(&run.id);
let live = live_by_run.get(&run.id).and_then(|value| value.as_ref());
let role = match inferred_stack_status(run.stack_status, live) {
RunStackStatus::Merged => QueueRole::Merged,
RunStackStatus::Superseded => QueueRole::Superseded,
RunStackStatus::Active => {
if block.is_some() {
QueueRole::Blocked
} else if head_index == Some(index) {
QueueRole::Ready
} else {
QueueRole::Draft
}
}
};
let next_action = queue_next_action(role, block, pr_number(run).is_some());
result.insert(
run.id.clone(),
QueueRunView {
role,
block_reason: block.map(|value| value.reason),
blocked_at: block.map(|value| value.attempted_at),
next_action,
},
);
}
result
}
fn queue_next_action(role: QueueRole, block: Option<&QueueBlock>, has_pr: bool) -> QueueNextAction {
match role {
QueueRole::Ready | QueueRole::Merged => QueueNextAction::AwaitMerge,
QueueRole::Superseded => QueueNextAction::CombinePrs,
QueueRole::Draft => {
if has_pr {
QueueNextAction::AwaitMerge
} else {
QueueNextAction::OpenPr
}
}
QueueRole::Blocked => match block.map(|value| value.reason) {
Some(QueueBlockReason::ScratchDirty | QueueBlockReason::RebaseConflict) => {
QueueNextAction::ResolveConflict
}
Some(QueueBlockReason::MissingPr) => QueueNextAction::OpenPr,
_ => QueueNextAction::AwaitMerge,
},
}
}
fn find_queue_head_index(runs: &[Run], live_snapshot: &LivePrSnapshot) -> Option<usize> {
runs.iter().position(|run| {
let live = live_snapshot.state_for_run(run);
inferred_stack_status(run.stack_status, live) == RunStackStatus::Active
})
}
fn inferred_stack_status(
durable: RunStackStatus,
live: Option<&LivePullRequestState>,
) -> RunStackStatus {
if durable != RunStackStatus::Active {
return durable;
}
match live.map(|state| state.state) {
Some(LivePrState::Merged) => RunStackStatus::Merged,
Some(LivePrState::Closed) => RunStackStatus::Superseded,
_ => RunStackStatus::Active,
}
}
fn pr_number(run: &Run) -> Option<u32> {
run.pr.as_ref()?.number
}
async fn set_queue_block(
store: &SharedStore,
run: &Run,
reason: QueueBlockReason,
conflict_files: Vec<String>,
error: Option<String>,
event_hub: Option<&EventHub>,
) -> Result<(), String> {
let block = QueueBlock {
wave_id: run.wave_id.clone(),
run_id: run.id.clone(),
reason,
attempted_at: OffsetDateTime::now_utc(),
conflict_files,
error,
};
let attention_id = attention_id_for_queue_block(&block.run_id);
let existing = store
.get_attention_item(&attention_id)
.await
.map_err(|err| format!("get queue attention item failed: {err}"))?;
let item = queue_block_attention_item_from_existing(&block, existing.as_ref());
store
.upsert_attention_item(&item)
.await
.map_err(|err| format!("upsert_queue_block failed: {err}"))?;
if let Some(event_hub) = event_hub {
match existing {
None => event_hub.send(Event::attention_created(item)),
Some(existing) if existing.status == AttentionStatus::Resolved => {
event_hub.send(Event::attention_created(item));
}
Some(existing)
if existing.status != item.status
|| existing.title != item.title
|| existing.summary != item.summary
|| existing.context != item.context =>
{
event_hub.send(Event::attention_updated(item));
}
Some(_) => {}
}
}
Ok(())
}
async fn clear_queue_block(
store: &SharedStore,
_wave_id: &LfdId,
run_id: &LfdId,
event_hub: Option<&EventHub>,
) -> Result<(), String> {
let attention_id = attention_id_for_queue_block(run_id);
let Some(mut item) = store
.get_attention_item(&attention_id)
.await
.map_err(|err| format!("get queue attention item failed: {err}"))?
else {
return Ok(());
};
if item.status == AttentionStatus::Resolved {
return Ok(());
}
item.status = AttentionStatus::Resolved;
item.resolved_at = Some(OffsetDateTime::now_utc());
store
.upsert_attention_item(&item)
.await
.map_err(|err| format!("resolve queue attention item failed: {err}"))?;
if let Some(event_hub) = event_hub {
event_hub.send(Event::attention_resolved(item));
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Mutex};
use crate::lfd::events::EventHub;
use crate::lfd::id::LfdId;
use crate::lfd::types::{
PullRequest, QueueBlockReason, RepoWork, Run, RunStatus, Wave, WaveStatus,
};
#[derive(Debug, Default)]
struct MockOps {
scratch_clean: bool,
rebase_fail_for_branch: Option<String>,
ready_calls: Mutex<Vec<String>>,
}
impl QueueOps for MockOps {
fn ensure_branch_checked_out(&self, _worktree: &Path, _branch: &str) -> Result<(), String> {
Ok(())
}
fn mark_ready(&self, _worktree: &Path, _pr_number: u32) -> Result<(), String> {
self.ready_calls
.lock()
.expect("mutex")
.push("ready".to_string());
Ok(())
}
fn mark_draft(&self, _worktree: &Path, _pr_number: u32) -> Result<(), String> {
Ok(())
}
fn rebase_onto_default(
&self,
_worktree: &Path,
default_branch: &str,
) -> Result<(), QueueRebaseConflict> {
if self
.rebase_fail_for_branch
.as_ref()
.is_some_and(|value| value == default_branch)
{
return Err(QueueRebaseConflict {
files: vec!["src/lib.rs".to_string()],
});
}
Ok(())
}
fn scratch_clean(&self, _worktree: &Path) -> Result<bool, String> {
Ok(self.scratch_clean)
}
}
async fn sqlite_store() -> SharedStore {
let db_path = std::env::temp_dir().join(format!("lfd-queue-test-{}.db", LfdId::new()));
let config = crate::lfdb::StorageConfig::sqlite(db_path);
Arc::new(
crate::lfdb::open_store(&config)
.await
.expect("sqlite store should initialize"),
)
}
fn make_wave(repo: &str) -> Wave {
Wave {
id: LfdId::new(),
name: "queue-wave".to_string(),
primary_flow: "ship-roadmap".to_string(),
goal: "ship-roadmap".to_string(),
metrics: Vec::new(),
repos: vec![RepoWork {
repo: repo.to_string(),
worktree: String::new(),
branch: String::new(),
status: WaveStatus::Idle,
iteration: 0,
cycle_start_iteration: 0,
position: 0,
}],
direction: Vec::new(),
area: Vec::new(),
paused: false,
created_at: Some(OffsetDateTime::now_utc()),
workers: 1,
parent_wave_id: None,
}
}
fn make_run(wave: &Wave, stack_position: u32, pr_number: u32) -> Run {
Run {
id: LfdId::new(),
wave_id: wave.id().clone(),
repo: wave.repo().to_string(),
flow: wave.primary_flow().clone(),
task: None,
direction: wave.direction().clone(),
area: wave.area().clone(),
iteration: stack_position,
step_index: 0,
status: RunStatus::Completed,
worktree: ".".to_string(),
branch: format!("feature-{pr_number}"),
started_at: Some(OffsetDateTime::now_utc()),
ended_at: Some(OffsetDateTime::now_utc()),
error: None,
flow_parents: Vec::new(),
execution_cursor: None,
parent_run_id: None,
parent_pr_number: None,
stack_position,
stack_group_id: wave.id().to_string(),
stack_status: RunStackStatus::Active,
lineage_inferred: false,
target_branch: "main".to_string(),
repair_of: None,
pr: Some(PullRequest {
url: format!("https://example.test/pr/{pr_number}"),
number: Some(pr_number),
state: Some("open".to_string()),
title: Some(format!("run-{pr_number}")),
branch: Some(format!("feature-{pr_number}")),
}),
}
}
async fn set_live_open(store: &SharedStore, run: &Run, is_draft: bool) {
let pr_number = pr_number(run).expect("pr number");
store
.upsert_live_pr_state(&LivePullRequestState {
repo_id: run.repo.clone(),
pr_number,
state: LivePrState::Open,
is_draft,
head_ref: run.branch.clone(),
head_sha: "abc123".to_string(),
base_ref: "main".to_string(),
updated_at: OffsetDateTime::now_utc(),
merged_at: None,
synced_at: OffsetDateTime::now_utc(),
})
.await
.expect("live state");
}
#[tokio::test]
async fn reconcile_promotes_only_oldest_unmerged() {
let store = sqlite_store().await;
let wave = make_wave(".");
store.create_wave(&wave).await.expect("wave");
let run1 = make_run(&wave, 0, 11);
let run2 = make_run(&wave, 1, 12);
store.create_run(&run1).await.expect("run1");
store.create_run(&run2).await.expect("run2");
set_live_open(&store, &run1, true).await;
set_live_open(&store, &run2, true).await;
let ops = MockOps {
scratch_clean: true,
..Default::default()
};
reconcile_wave_queue_with_ops(&store, &GitHubConfig::default(), wave.id(), &ops, None)
.await
.expect("reconcile");
let blocks = store
.list_queue_blocks(wave.id())
.await
.expect("list queue blocks")
.into_iter()
.map(|block| (block.run_id.clone(), block))
.collect::<HashMap<_, _>>();
let runs = store.list_stack_runs(wave.id()).await.expect("runs");
let live_snapshot = build_live_pr_snapshot(&store, &GitHubConfig::default(), &runs)
.await
.expect("live snapshot");
let projected = project_queue_views(
&runs,
|run| live_snapshot.state_for_run(run).cloned(),
&blocks,
);
let ready_count = projected
.values()
.filter(|view| view.role == QueueRole::Ready)
.count();
assert_eq!(ready_count, 1);
assert_eq!(
projected.get(&run1.id).map(|view| view.role),
Some(QueueRole::Ready)
);
assert_eq!(
projected.get(&run2.id).map(|view| view.role),
Some(QueueRole::Draft)
);
}
#[tokio::test]
async fn scratch_dirty_marks_blocked_with_resolve_conflict_action() {
let store = sqlite_store().await;
let wave = make_wave(".");
store.create_wave(&wave).await.expect("wave");
let run = make_run(&wave, 0, 31);
store.create_run(&run).await.expect("run");
set_live_open(&store, &run, true).await;
let ops = MockOps {
scratch_clean: false,
..Default::default()
};
reconcile_wave_queue_with_ops(&store, &GitHubConfig::default(), wave.id(), &ops, None)
.await
.expect("reconcile");
let block = store
.list_queue_blocks(wave.id())
.await
.expect("blocks")
.into_iter()
.next()
.expect("block");
assert_eq!(block.reason, QueueBlockReason::ScratchDirty);
let live_snapshot =
build_live_pr_snapshot(&store, &GitHubConfig::default(), std::slice::from_ref(&run))
.await
.expect("live snapshot");
let projection = project_queue_views(
std::slice::from_ref(&run),
|r| live_snapshot.state_for_run(r).cloned(),
&HashMap::from([(run.id.clone(), block)]),
);
assert_eq!(
projection.get(&run.id).map(|view| view.next_action),
Some(QueueNextAction::ResolveConflict)
);
}
#[tokio::test]
async fn repeated_queue_block_preserves_age_and_emits_only_once() {
let store = sqlite_store().await;
let wave = make_wave(".");
store.create_wave(&wave).await.expect("wave");
let run = make_run(&wave, 0, 41);
store.create_run(&run).await.expect("run");
set_live_open(&store, &run, true).await;
let event_hub = EventHub::new(8);
let mut rx = event_hub.subscribe();
let ops = MockOps {
scratch_clean: false,
..Default::default()
};
reconcile_wave_queue_with_ops(
&store,
&GitHubConfig::default(),
wave.id(),
&ops,
Some(&event_hub),
)
.await
.expect("first reconcile");
let created = rx.try_recv().expect("created event");
let first_item = match created {
Event::AttentionCreated { item, .. } => item,
other => panic!("expected attention created, got {other:?}"),
};
reconcile_wave_queue_with_ops(
&store,
&GitHubConfig::default(),
wave.id(),
&ops,
Some(&event_hub),
)
.await
.expect("second reconcile");
let queue_item = store
.get_attention_item(&first_item.id)
.await
.expect("get attention item")
.expect("attention item exists");
assert_eq!(
queue_item.surfaced_at.unix_timestamp(),
first_item.surfaced_at.unix_timestamp()
);
assert!(rx.try_recv().is_err(), "no duplicate attention event");
}
#[tokio::test]
async fn clearing_queue_block_emits_attention_resolved() {
let store = sqlite_store().await;
let wave = make_wave(".");
store.create_wave(&wave).await.expect("wave");
let run = make_run(&wave, 0, 51);
store.create_run(&run).await.expect("run");
set_live_open(&store, &run, true).await;
let event_hub = EventHub::new(8);
let mut rx = event_hub.subscribe();
let blocked_ops = MockOps {
scratch_clean: false,
..Default::default()
};
reconcile_wave_queue_with_ops(
&store,
&GitHubConfig::default(),
wave.id(),
&blocked_ops,
Some(&event_hub),
)
.await
.expect("block reconcile");
let _ = rx.try_recv().expect("created event");
let cleared_ops = MockOps {
scratch_clean: true,
..Default::default()
};
reconcile_wave_queue_with_ops(
&store,
&GitHubConfig::default(),
wave.id(),
&cleared_ops,
Some(&event_hub),
)
.await
.expect("clear reconcile");
let resolved = rx.try_recv().expect("resolved event");
let resolved_item = match resolved {
Event::AttentionResolved { item, .. } => item,
other => panic!("expected attention resolved, got {other:?}"),
};
assert_eq!(resolved_item.status, AttentionStatus::Resolved);
}
#[tokio::test]
async fn reconcile_lock_serializes_per_wave() {
let wave_id = LfdId::new();
let guard = acquire_reconcile_lock(&wave_id).await;
let locked = Arc::new(Mutex::new(false));
let locked_clone = Arc::clone(&locked);
let wave_id_clone = wave_id.clone();
let waiter = tokio::spawn(async move {
let _guard = acquire_reconcile_lock(&wave_id_clone).await;
*locked_clone.lock().expect("mutex") = true;
});
tokio::time::sleep(std::time::Duration::from_millis(25)).await;
assert!(!*locked.lock().expect("mutex"));
drop(guard);
waiter.await.expect("waiter task");
assert!(*locked.lock().expect("mutex"));
}
}