use std::collections::HashMap;
use crate::error::{Result, RiftError, SessionReject};
use crate::session::offset_tracker::ResumeDecision;
use crate::session::offset_tracker::decide;
use crate::session::session::Session;
use crate::topic::TopicStore;
pub struct ResumeManager {}
impl Default for ResumeManager {
fn default() -> Self {
Self::new()
}
}
impl ResumeManager {
pub fn new() -> Self {
Self {}
}
pub fn evaluate(
&self,
session: &Session,
last_offsets: &HashMap<String, i64>,
topic_offsets: &HashMap<String, i64>,
) -> Result<ResumeDecision> {
if !session.is_alive() {
return Err(RiftError::Session(SessionReject::Expired));
}
Ok(decide(last_offsets, topic_offsets))
}
pub fn topic_offsets(&self, store: &TopicStore, topics: &[String]) -> HashMap<String, i64> {
let mut out = HashMap::new();
for t in topics {
if let Some(entry) = store.get(t) {
out.insert(t.clone(), entry.head_offset());
}
}
out
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::session::session::ClientId;
use crate::topic::profile::TopicProfile;
#[test]
fn evaluate_checks_liveness() {
let m = ResumeManager::new();
let s = Session::new(
crate::session::session::SessionId::new(),
ClientId::new("c"),
);
s.bump_epoch();
let mut last = HashMap::new();
last.insert("t".into(), 1);
let mut head = HashMap::new();
head.insert("t".into(), 5);
let r = m.evaluate(&s, &last, &head);
assert!(r.is_ok());
}
#[test]
fn happy_resume() {
let m = ResumeManager::new();
let s = Session::new(
crate::session::session::SessionId::new(),
ClientId::new("c"),
);
let mut last = HashMap::new();
last.insert("t".into(), 4);
let mut head = HashMap::new();
head.insert("t".into(), 5);
let r = m.evaluate(&s, &last, &head).unwrap();
assert_eq!(r, ResumeDecision::Replaying);
}
#[tokio::test]
async fn topic_offsets_from_store() {
use crate::storage::{MemoryOffsetStore, OffsetStore};
let m = ResumeManager::new();
let store = TopicStore::new();
let offsets = MemoryOffsetStore::new();
let entry = store.get_or_create("t", TopicProfile::default()).unwrap();
let o1 = offsets.alloc("t").await;
let o2 = offsets.alloc("t").await;
entry.append(crate::topic::store::LogEntry {
offset: o1,
publisher_session: None,
message_id: "m1".into(),
class: "event".into(),
event: Some("e".into()),
payload: bytes::Bytes::from_static(b"x"),
timestamp: 0,
appended_at: None,
});
entry.append(crate::topic::store::LogEntry {
offset: o2,
publisher_session: None,
message_id: "m2".into(),
class: "event".into(),
event: Some("e".into()),
payload: bytes::Bytes::from_static(b"x"),
timestamp: 0,
appended_at: None,
});
let heads = m.topic_offsets(&store, &["t".to_string()]);
assert_eq!(heads.get("t").copied(), Some(2));
}
}