避免误淘汰实时 DirectThread 订阅
仅在订阅游标落后于队列头时驱逐订阅,并补充队头被 pin 的回归测试。
This commit is contained in:
@@ -364,9 +364,15 @@ impl DirectThreadManager {
|
||||
let Some(thread) = self.threads.get_mut(thread_id) else {
|
||||
return;
|
||||
};
|
||||
let oldest_seq = thread
|
||||
.events
|
||||
.get(thread.head)
|
||||
.map(|stored| stored.event.seq)
|
||||
.unwrap_or(thread.next_seq.saturating_add(1));
|
||||
let slowest = thread
|
||||
.subscribers
|
||||
.iter()
|
||||
.filter(|(_, subscriber)| subscriber.cursor < oldest_seq)
|
||||
.min_by_key(|(_, subscriber)| subscriber.cursor)
|
||||
.map(|(id, _)| id.clone());
|
||||
if let Some(subscription_id) = slowest {
|
||||
@@ -530,6 +536,20 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn current_subscriber_is_not_expired_by_pinned_queue_head() {
|
||||
let mut manager = DirectThreadManager::with_limits(2, 100_000);
|
||||
manager.append("thread-1", draft("item.started", "turn-1", Some("item-1")));
|
||||
let subscription = manager.subscribe("thread-1");
|
||||
manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1")));
|
||||
manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1")));
|
||||
manager.append("thread-1", draft("item.delta", "turn-1", Some("item-1")));
|
||||
assert_ne!(
|
||||
manager.consume(&subscription.subscription_id),
|
||||
Err(SUBSCRIPTION_EXPIRED.to_string())
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unresolved_approval_is_kept_in_bootstrap_until_resolved() {
|
||||
let mut manager = DirectThreadManager::with_limits(100, 100_000);
|
||||
|
||||
Reference in New Issue
Block a user