fix(#937): narrow skip to same-body pending — different schedules can still enqueue
This commit is contained in:
parent
bf4937d116
commit
ce27ef8b9e
1 changed files with 22 additions and 0 deletions
|
|
@ -278,6 +278,27 @@ impl Broker {
|
||||||
Ok(n > 0)
|
Ok(n > 0)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Returns true when the recipient already has at least one
|
||||||
|
/// undelivered message from `sender` with exactly `body` in the
|
||||||
|
/// broker. Used by the scheduler to skip re-delivery of the same
|
||||||
|
/// scheduled prompt without blocking distinct schedules whose
|
||||||
|
/// bodies differ.
|
||||||
|
pub fn has_pending_with_body(
|
||||||
|
&self,
|
||||||
|
recipient: &str,
|
||||||
|
sender: &str,
|
||||||
|
body: &str,
|
||||||
|
) -> Result<bool> {
|
||||||
|
let conn = self.conn.lock().unwrap();
|
||||||
|
let n: i64 = conn.query_row(
|
||||||
|
"SELECT COUNT(*) FROM messages
|
||||||
|
WHERE recipient = ?1 AND sender = ?2 AND body = ?3 AND delivered_at IS NULL",
|
||||||
|
params![recipient, sender, body],
|
||||||
|
|row| row.get(0),
|
||||||
|
)?;
|
||||||
|
Ok(n > 0)
|
||||||
|
}
|
||||||
|
|
||||||
/// Long-poll variant of `recv_batch`: returns immediately if any
|
/// Long-poll variant of `recv_batch`: returns immediately if any
|
||||||
/// row is pending (popping up to `max`); otherwise waits up to
|
/// row is pending (popping up to `max`); otherwise waits up to
|
||||||
/// `timeout` for the broker to emit a `Sent { to: recipient }`
|
/// `timeout` for the broker to emit a `Sent { to: recipient }`
|
||||||
|
|
@ -1269,3 +1290,4 @@ mod tests {
|
||||||
assert!(pop_one(broker, "bob").is_none());
|
assert!(pop_one(broker, "bob").is_none());
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue