diff --git a/hive-c0re/src/broker.rs b/hive-c0re/src/broker.rs index 020c07dd..5a701ce8 100644 --- a/hive-c0re/src/broker.rs +++ b/hive-c0re/src/broker.rs @@ -278,6 +278,27 @@ impl Broker { 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 { + 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 /// row is pending (popping up to `max`); otherwise waits up to /// `timeout` for the broker to emit a `Sent { to: recipient }` @@ -1269,3 +1290,4 @@ mod tests { assert!(pop_one(broker, "bob").is_none()); } } +