use pushkin_core::delivery::{
Delivery, DeliveryError, DeliveryIndex, DeliveryRequest, DELIVERY_HORIZON,
};
use pushkin_core::events::SessionId;
fn open(dir: &tempfile::TempDir) -> Result<DeliveryIndex, DeliveryError> {
DeliveryIndex::open(dir.path().join("events.db"))
}
fn request<'a>(session: &'a SessionId, slice_key: &'a str, complete: bool) -> DeliveryRequest<'a> {
DeliveryRequest {
session,
cwd: "/repo",
slice_key,
complete,
}
}
#[test]
fn full_delivery_recorded_per_session_and_cwd() {
let dir = tempfile::tempdir().unwrap();
let index = open(&dir).unwrap();
let session = SessionId::from_name("s1");
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Full,
"first reference delivers in full"
);
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Pointer,
"recorded for this (session, cwd)"
);
let other_session = SessionId::from_name("s2");
assert_eq!(
index
.decide(&request(&other_session, "contract/user", true))
.unwrap(),
Delivery::Full,
"a different session has its own state"
);
let other_cwd = DeliveryRequest {
session: &session,
cwd: "/elsewhere",
slice_key: "contract/user",
complete: true,
};
assert_eq!(
index.decide(&other_cwd).unwrap(),
Delivery::Full,
"a different cwd has its own state"
);
}
#[test]
fn repeat_reference_collapses_to_pointer_line() {
let dir = tempfile::tempdir().unwrap();
let index = open(&dir).unwrap();
let session = SessionId::from_name("s1");
index
.decide(&request(&session, "contract/user", true))
.unwrap();
for _ in 0..3 {
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Pointer,
"every repeat collapses to a pointer"
);
}
}
#[test]
fn truncated_first_emission_never_deduped() {
let dir = tempfile::tempdir().unwrap();
let index = open(&dir).unwrap();
let session = SessionId::from_name("s1");
assert_eq!(
index
.decide(&request(&session, "contract/user", false))
.unwrap(),
Delivery::Full,
"truncated emission still emits what fits"
);
assert_eq!(
index
.decide(&request(&session, "contract/user", false))
.unwrap(),
Delivery::Full,
"a truncated emission was never recorded — the agent never saw it in full"
);
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Full,
"still undelivered until a complete emission happens"
);
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Pointer,
"the complete emission was recorded"
);
}
#[test]
fn compaction_event_clears_index_and_reemits_full() {
let dir = tempfile::tempdir().unwrap();
let index = open(&dir).unwrap();
let session = SessionId::from_name("s1");
let untouched = SessionId::from_name("s2");
index
.decide(&request(&session, "contract/user", true))
.unwrap();
index
.decide(&request(&untouched, "contract/user", true))
.unwrap();
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Pointer
);
index.clear(&session, "/repo").unwrap();
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Full,
"compaction cleared the scope — re-emit in full"
);
assert_eq!(
index
.decide(&request(&untouched, "contract/user", true))
.unwrap(),
Delivery::Pointer,
"clearing one scope leaves other sessions' state intact"
);
}
#[test]
fn horizon_expired_slice_reemits() {
let dir = tempfile::tempdir().unwrap();
let index = open(&dir).unwrap();
let session = SessionId::from_name("s1");
index
.decide(&request(&session, "contract/user", true))
.unwrap();
let noise: Vec<String> = (0..DELIVERY_HORIZON)
.map(|counter| format!("contract/noise-{counter}"))
.collect();
for key in &noise {
index.decide(&request(&session, key, true)).unwrap();
}
assert_eq!(
index
.decide(&request(&session, "contract/user", true))
.unwrap(),
Delivery::Full,
"content past the horizon has scrolled out of context — re-emit in full"
);
}