use super::gateway::{DIRECT_MESSAGES, GUILD_MESSAGES, MESSAGE_CONTENT};
use crate::discord::{Discord, Source, DISCORD_EPOCH_MS};
use anyhow::{Context, Result};
use serde_json::Value;
use std::collections::BTreeMap;
use std::num::NonZeroU64;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc::{Receiver, RecvTimeoutError};
use std::sync::Arc;
use std::time::{Duration, Instant, SystemTime};
const RETRY_FIRST: Duration = Duration::from_secs(60);
const RETRY_LAST: Duration = Duration::from_secs(15 * 60);
const PULLS: usize = 5;
const UNSTORED: &str = "unstored";
const NEW_DM: &str = "new";
pub struct Intake {
discord: Discord,
source: Box<dyn Source + Send>,
channels: Vec<NonZeroU64>,
dms: bool,
directory: PathBuf,
floor: u64,
account: Option<String>,
}
#[derive(Debug)]
pub enum Work {
Message(Value),
Update(Value),
Account(String),
Backfill,
}
#[derive(Debug, PartialEq, Eq)]
pub enum Done {
Stored,
Ignored,
Recorded,
Backfilled {
channels: usize,
failed: usize,
more: bool,
unstored: usize,
},
}
pub fn floor_at(time: SystemTime) -> u64 {
let milliseconds = time
.duration_since(SystemTime::UNIX_EPOCH)
.map_or(0, |since| since.as_millis());
let since_epoch =
u64::try_from(milliseconds.saturating_sub(u128::from(DISCORD_EPOCH_MS))).unwrap_or(0);
(since_epoch << 22).saturating_sub(1)
}
impl Intake {
pub fn new(
discord: Discord,
source: Box<dyn Source + Send>,
channels: Vec<NonZeroU64>,
dms: bool,
directory: PathBuf,
floor: u64,
) -> Self {
Self {
discord,
source,
channels,
dms,
directory,
floor,
account: None,
}
}
pub fn intents(&self) -> u64 {
let mut intents = 0;
if !self.channels.is_empty() {
intents |= GUILD_MESSAGES | MESSAGE_CONTENT;
}
if self.dms {
intents |= DIRECT_MESSAGES;
}
intents
}
pub fn preflight(&self) -> Result<()> {
self.discord.preflight_write()
}
pub fn handle(&mut self, work: Work) -> Result<Done> {
match work {
Work::Message(message) => self.message(message, true),
Work::Update(message) => self.message(message, false),
Work::Account(user) => {
crate::discord::validate_snowflake(&user).context("the bot's user id")?;
self.account = Some(user);
self.record_account()?;
Ok(Done::Recorded)
}
Work::Backfill => Ok(self.backfill()),
}
}
fn record_account(&mut self) -> Result<()> {
if let Some(user) = &self.account {
self.discord.record_bot_account(user)?;
self.account = None;
}
Ok(())
}
fn message(&mut self, message: Value, created: bool) -> Result<Done> {
let channel = snowflake(&message["channel_id"]).context("the message names its channel")?;
let id = snowflake(&message["id"]).context("the message has an id")?;
let dm = message["guild_id"].is_null();
let new_dm = dm && created;
let first = created.then(|| id.saturating_sub(1));
let floor = match self.floor_for(channel, dm, first) {
Ok(Some(floor)) => floor,
Ok(None) => return Ok(Done::Ignored),
Err(error) => {
self.keep(channel, id, new_dm);
return Err(error);
}
};
self.store(message, channel, id, floor, new_dm)
}
fn store(
&mut self,
message: Value,
channel: u64,
id: u64,
floor: u64,
new_dm: bool,
) -> Result<Done> {
if id <= floor {
return Ok(Done::Ignored);
}
match self.discord.observe(message, floor, &mut *self.source) {
Ok(_) => {
self.unkeep(channel, id);
Ok(Done::Stored)
}
Err(error) => {
self.keep(channel, id, new_dm);
Err(error)
}
}
}
fn marker(&self, channel: u64, id: u64, new_dm: bool) -> PathBuf {
let name = if new_dm {
format!("{channel}-{id}-{NEW_DM}")
} else {
format!("{channel}-{id}")
};
self.directory.join(UNSTORED).join(name)
}
fn keep(&self, channel: u64, id: u64, new_dm: bool) {
let marker = self.marker(channel, id, new_dm);
let kept = std::fs::create_dir_all(self.directory.join(UNSTORED)).and_then(|()| {
match std::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&marker)
{
Ok(_) => Ok(()),
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
Err(error) => Err(error),
}
});
if let Err(keeping) = kept {
eprintln!(
"[discord] message {id} in channel {channel} is not kept to be fetched again \
({}: {keeping}); only a backfill's pages may bring it",
marker.display()
);
}
}
fn unkeep(&self, channel: u64, id: u64) {
for new_dm in [false, true] {
let marker = self.marker(channel, id, new_dm);
match std::fs::remove_file(&marker) {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => eprintln!(
"[discord] message {id} is no longer to be fetched again, but {} stays: \
{error}",
marker.display()
),
}
}
}
fn floor_for(&self, channel: u64, dm: bool, first: Option<u64>) -> Result<Option<u64>> {
if self.channels.iter().any(|id| id.get() == channel) {
return self.keep_floor("channels", channel, self.floor).map(Some);
}
if !(self.dms && dm) {
return Ok(None);
}
match (self.floor("dms", channel)?, first) {
(Some(floor), Some(first)) if first < floor => {
eprintln!(
"[discord] DM channel {channel} begins earlier, at message {}, heard as new \
after a later one",
first + 1
);
self.write_floor("dms", channel, first).map(Some)
}
(Some(floor), _) => Ok(Some(floor)),
(None, Some(first)) => self.write_floor("dms", channel, first).map(Some),
(None, None) => Ok(None),
}
}
fn floor(&self, kind: &str, channel: u64) -> Result<Option<u64>> {
let path = self.directory.join(kind).join(channel.to_string());
match std::fs::read_to_string(&path) {
Ok(floor) => floor
.trim()
.parse()
.map(Some)
.with_context(|| format!("{} holds no message id", path.display())),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(error) => Err(error).with_context(|| format!("read {}", path.display())),
}
}
fn keep_floor(&self, kind: &str, channel: u64, first: u64) -> Result<u64> {
if let Some(floor) = self.floor(kind, channel)? {
return Ok(floor);
}
self.write_floor(kind, channel, first)
}
fn write_floor(&self, kind: &str, channel: u64, floor: u64) -> Result<u64> {
let directory = self.directory.join(kind);
let path = directory.join(channel.to_string());
std::fs::create_dir_all(&directory)
.with_context(|| format!("create {}", directory.display()))?;
let staging = directory.join(format!(".{channel}"));
std::fs::write(&staging, floor.to_string())
.with_context(|| format!("write {}", staging.display()))?;
std::fs::rename(&staging, &path)
.with_context(|| format!("keep the floor in {}", path.display()))?;
Ok(floor)
}
fn backfill_channels(&self) -> (Vec<(u64, Result<u64>)>, bool) {
let mut channels: Vec<(u64, Result<u64>)> = self
.channels
.iter()
.map(|id| (id.get(), self.keep_floor("channels", id.get(), self.floor)))
.collect();
let mut unlisted = false;
if self.dms {
let directory = self.directory.join("dms");
let entries = match std::fs::read_dir(&directory) {
Ok(entries) => entries.flatten().collect(),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Vec::new(),
Err(error) => {
eprintln!(
"[discord] the DM channels seen before are unknown ({}: {error}); \
none is backfilled",
directory.display()
);
unlisted = true;
Vec::new()
}
};
for entry in entries {
let Some(id) = entry.file_name().to_str().and_then(|n| n.parse().ok()) else {
continue;
};
if channels.iter().any(|(channel, _)| *channel == id) {
continue;
}
match self.floor("dms", id) {
Ok(Some(floor)) => channels.push((id, Ok(floor))),
Ok(None) => {}
Err(error) => channels.push((id, Err(error))),
}
}
}
(channels, unlisted)
}
fn kept(&self) -> Result<BTreeMap<(u64, u64), bool>> {
let directory = self.directory.join(UNSTORED);
let entries = match std::fs::read_dir(&directory) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
return Ok(BTreeMap::new())
}
Err(error) => return Err(error).with_context(|| directory.display().to_string()),
};
let mut kept: BTreeMap<(u64, u64), bool> = BTreeMap::new();
for entry in entries.flatten() {
let name = entry.file_name();
let Some(name) = name.to_str() else {
continue;
};
let (name, new_dm) = match name
.strip_suffix(NEW_DM)
.and_then(|name| name.strip_suffix('-'))
{
Some(name) => (name, true),
None => (name, false),
};
let Some((channel, id)) = name.split_once('-') else {
continue;
};
let (Ok(channel), Ok(id)) = (channel.parse(), id.parse()) else {
continue;
};
*kept.entry((channel, id)).or_default() |= new_dm;
}
Ok(kept)
}
fn retry_unstored(&mut self) -> usize {
let kept = match self.kept() {
Ok(kept) => kept,
Err(error) => {
eprintln!("[discord] the messages to fetch again are unknown ({error:#})");
return 1;
}
};
let mut unstored = 0;
for ((channel, id), new_dm) in kept {
match self.refetch(channel, id, new_dm) {
Ok(Done::Stored) => {
eprintln!("[discord] stored Discord message {id}, fetched again");
}
Ok(_) => {
eprintln!(
"[discord] Discord message {id} in channel {channel} does not exist or \
is not intake's to store; not fetched again"
);
self.unkeep(channel, id);
}
Err(error) => {
unstored += 1;
eprintln!(
"[discord] fetching Discord message {id} in channel {channel} again \
failed: {error:#}"
);
}
}
}
unstored
}
fn refetch(&mut self, channel: u64, id: u64, new_dm: bool) -> Result<Done> {
let first = new_dm.then(|| id.saturating_sub(1));
let Some(floor) = self.floor_for(channel, true, first)? else {
return Ok(Done::Ignored);
};
if id <= floor {
return Ok(Done::Ignored);
}
let Some(message) = self.source.message(&channel.to_string(), id)? else {
return Ok(Done::Ignored);
};
anyhow::ensure!(
snowflake(&message["id"]) == Some(id),
"Discord answered message {id} with message {}",
message["id"]
);
self.store(message, channel, id, floor, new_dm)
}
fn backfill(&mut self) -> Done {
if let Err(error) = self.record_account() {
eprintln!(
"[discord] recording the bot account failed, so nothing is pulled: {error:#}"
);
let (channels, unlisted) = self.backfill_channels();
let count = channels.len() + usize::from(unlisted) + 1;
return Done::Backfilled {
channels: count,
failed: count,
more: false,
unstored: self.kept().map_or(1, |kept| kept.len()),
};
}
let unstored = self.retry_unstored();
let (channels, unlisted) = self.backfill_channels();
let count = channels.len() + usize::from(unlisted);
let mut failed = usize::from(unlisted);
let mut more = false;
for (channel, floor) in channels {
let floor = match floor {
Ok(floor) => floor,
Err(error) => {
failed += 1;
eprintln!(
"[discord] where Discord channel {channel} begins is unknown, so it \
is passed over: {error:#}"
);
continue;
}
};
let mut stored = 0;
for pull in 1..=PULLS {
match self.discord.pull_channel(
&channel.to_string(),
Some(floor),
&mut *self.source,
) {
Ok(receipt) => {
stored += receipt.observations;
if !receipt.more {
break;
}
more |= pull == PULLS;
}
Err(error) => {
failed += 1;
eprintln!(
"[discord] backfilling Discord channel {channel} failed: {error:#}"
);
break;
}
}
}
if stored > 0 {
eprintln!("[discord] backfilled {stored} Discord messages in channel {channel}");
}
}
Done::Backfilled {
channels: count,
failed,
more,
unstored,
}
}
}
#[derive(Debug, Default)]
struct Retry {
due: Option<Instant>,
wait: Duration,
}
impl Retry {
fn soon(&mut self, now: Instant) {
self.due.get_or_insert(now + self.wait.max(RETRY_FIRST));
}
fn again(&mut self, now: Instant) {
self.wait = (self.wait * 2).clamp(RETRY_FIRST, RETRY_LAST);
self.due = Some(now + self.wait);
}
fn caught_up(&mut self) {
*self = Self::default();
}
fn next(&mut self, inbox: &Receiver<Work>) -> Option<Work> {
let Some(due) = self.due else {
return inbox.recv().ok();
};
match inbox.recv_timeout(due.saturating_duration_since(Instant::now())) {
Ok(work) => Some(work),
Err(RecvTimeoutError::Timeout) => {
self.due = None;
Some(Work::Backfill)
}
Err(RecvTimeoutError::Disconnected) => None,
}
}
}
pub struct Worker {
sender: Option<std::sync::mpsc::Sender<Work>>,
stopping: Arc<AtomicBool>,
finished: Option<tokio::sync::oneshot::Receiver<()>>,
stopped: bool,
}
pub fn start(mut intake: Intake) -> Worker {
let (sender, inbox) = std::sync::mpsc::channel::<Work>();
let stopping = Arc::new(AtomicBool::new(false));
let (done, finished) = tokio::sync::oneshot::channel();
let passing = stopping.clone();
std::thread::Builder::new()
.name("discord-intake".to_owned())
.spawn(move || {
let mut retry = Retry::default();
while let Some(work) = retry.next(&inbox) {
if passing.load(Ordering::SeqCst) && matches!(work, Work::Backfill) {
continue;
}
let what = match &work {
Work::Message(message) => format!(
"Discord message {}",
message["id"].as_str().unwrap_or("without an id")
),
Work::Update(message) => format!(
"the edit of Discord message {}",
message["id"].as_str().unwrap_or("without an id")
),
Work::Account(user) => format!("the bot account {user}"),
Work::Backfill => "the backfill".to_owned(),
};
let handled =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| intake.handle(work)));
match handled {
Ok(Ok(Done::Stored)) => eprintln!("[discord] stored {what}"),
Ok(Ok(Done::Backfilled {
channels,
failed,
more,
unstored,
})) => {
if failed > 0 {
eprintln!("[discord] backfill: {failed} of {channels} channels failed");
}
if unstored > 0 {
eprintln!(
"[discord] backfill: {unstored} messages whose live write \
failed are still not stored"
);
}
if failed > 0 || more || unstored > 0 {
retry.again(Instant::now());
} else {
retry.caught_up();
}
}
Ok(Ok(_)) => {}
Ok(Err(error)) => {
eprintln!(
"[discord] storing {what} failed: {error:#}; a backfill fetches it \
again"
);
retry.soon(Instant::now());
}
Err(_) => {
eprintln!("[discord] storing {what} panicked; intake goes on");
retry.soon(Instant::now());
}
}
}
let _ = done.send(());
})
.expect("spawn the intake thread");
Worker {
sender: Some(sender),
stopping,
finished: Some(finished),
stopped: false,
}
}
impl Worker {
#[cfg(test)]
pub fn from_sender(sender: std::sync::mpsc::Sender<Work>) -> Self {
Self {
sender: Some(sender),
stopping: Arc::new(AtomicBool::new(false)),
finished: None,
stopped: false,
}
}
pub fn send(&mut self, work: Work) {
let sent = self
.sender
.as_ref()
.is_some_and(|sender| sender.send(work).is_ok());
if !sent && !self.stopped {
eprintln!("[discord] the intake thread has stopped; Discord messages are not stored");
self.stopped = true;
}
}
pub async fn stop(mut self, limit: Duration) {
self.stopping.store(true, Ordering::SeqCst);
drop(self.sender.take());
let Some(finished) = self.finished.take() else {
return;
};
if tokio::time::timeout(limit, finished).await.is_err() {
eprintln!("[discord] intake still busy after {limit:?}; leaving it");
}
}
}
fn snowflake(value: &Value) -> Option<u64> {
value.as_str().and_then(|id| id.parse().ok())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::discord::{PageRequest, ReadOptions};
use anyhow::anyhow;
use serde_json::json;
use std::collections::{BTreeMap, BTreeSet};
use std::sync::{Arc, Mutex};
const CHANNEL: u64 = 100000000000000200;
const GUILD: &str = "100000000000000300";
const AUTHOR: &str = "100000000000000400";
const FLOOR: u64 = 100000000000000000;
#[derive(Clone, Default)]
struct Fake {
messages: Arc<Mutex<BTreeMap<String, Vec<Value>>>>,
attachments: Arc<Mutex<BTreeMap<String, Vec<u8>>>>,
pages: Arc<Mutex<Vec<(String, PageRequest)>>>,
fetched: Arc<Mutex<Vec<(String, u64)>>>,
unreadable: Arc<Mutex<BTreeSet<String>>>,
}
impl Source for Fake {
fn page(&mut self, channel_id: &str, request: PageRequest) -> Result<Vec<Value>> {
self.pages
.lock()
.unwrap()
.push((channel_id.to_owned(), request));
if self.unreadable.lock().unwrap().contains(channel_id) {
return Ok(Vec::new());
}
let mut page: Vec<Value> = self
.messages
.lock()
.unwrap()
.get(channel_id)
.cloned()
.unwrap_or_default()
.into_iter()
.filter(|message| {
let id = id_of(message);
request.after.is_none_or(|after| id > after)
&& request.before.is_none_or(|before| id < before)
})
.collect();
page.sort_by_key(id_of);
let limit = request.limit as usize;
let mut page = if request.after.is_some() {
page.into_iter().take(limit).collect()
} else {
page.split_off(page.len().saturating_sub(limit))
};
page.reverse();
Ok(page)
}
fn message(&mut self, channel_id: &str, message_id: u64) -> Result<Option<Value>> {
self.fetched
.lock()
.unwrap()
.push((channel_id.to_owned(), message_id));
anyhow::ensure!(
!self.unreadable.lock().unwrap().contains(channel_id),
"discord read of message {message_id} failed (403 Forbidden): \
{{\"message\": \"Missing Permissions\", \"code\": 50013}}"
);
Ok(self
.messages
.lock()
.unwrap()
.get(channel_id)
.and_then(|messages| {
messages
.iter()
.find(|message| id_of(message) == message_id)
.cloned()
}))
}
fn attachment(&mut self, url: &str, _limit: u64) -> Result<Vec<u8>> {
self.attachments
.lock()
.unwrap()
.get(url)
.cloned()
.ok_or_else(|| anyhow!("no such attachment {url}"))
}
}
fn id_of(message: &Value) -> u64 {
message["id"].as_str().unwrap().parse().unwrap()
}
fn rest(id: &str, channel: u64, content: &str, attachment_url: Option<&str>) -> Value {
let attachments = match attachment_url {
Some(url) => json!([{
"id": format!("{id}9"), "url": url, "filename": "photo.png",
"content_type": "image/png", "size": 5
}]),
None => json!([]),
};
json!({
"id": id,
"channel_id": channel.to_string(),
"type": 0,
"content": content,
"author": {"id": AUTHOR, "username": "ada", "global_name": "Ada"},
"timestamp": format!("2026-09-26T08:00:{:02}Z", id.parse::<u64>().unwrap() % 60),
"edited_timestamp": null,
"attachments": attachments,
"referenced_message": null,
})
}
fn gateway(message: &Value) -> Value {
let mut message = message.clone();
message["guild_id"] = json!(GUILD);
message["member"] = json!({"roles": [], "joined_at": "2026-01-01T00:00:00Z"});
message["mentions"] = json!([]);
if let Some(attachment) = message["attachments"].get_mut(0) {
attachment["url"] = json!("https://cdn.example/photo.png?ex=gateway");
}
message
}
struct Fixture {
_directory: tempfile::TempDir,
discord: Discord,
fake: Fake,
state: PathBuf,
}
impl Fixture {
fn new() -> Self {
crate::test_support::clear_ambient_environment();
let directory = tempfile::tempdir().unwrap();
let pile = directory.path().join("intake.pile");
let key = directory.path().join("intake.key");
std::fs::File::create(&pile).unwrap();
crate::storage::initialize_signer(&pile, Some(&key)).unwrap();
let fake = Fake::default();
for (url, bytes) in [
("https://cdn.example/photo.png?ex=gateway", b"image"),
("https://cdn.example/photo.png?ex=rest", b"image"),
] {
fake.attachments
.lock()
.unwrap()
.insert(url.to_owned(), bytes.to_vec());
}
Self {
discord: Discord::new(pile, Some(key)),
fake,
state: directory.path().join("intake"),
_directory: directory,
}
}
fn intake(&self, channels: &[u64], dms: bool) -> Intake {
Intake::new(
self.discord.clone(),
Box::new(self.fake.clone()),
channels
.iter()
.map(|id| NonZeroU64::new(*id).unwrap())
.collect(),
dms,
self.state.clone(),
FLOOR,
)
}
fn serve(&self, channel: u64, messages: Vec<Value>) {
self.fake
.messages
.lock()
.unwrap()
.insert(channel.to_string(), messages);
}
fn stored(&self, channel: u64) -> Vec<(String, usize, usize)> {
self.discord
.read(ReadOptions {
channel_id: Some(channel.to_string()),
limit: 0,
..ReadOptions::default()
})
.unwrap()
.messages
.into_iter()
.map(|message| {
(
message.content,
message.attachments.len(),
message.variant_count,
)
})
.collect()
}
}
#[test]
fn a_gateway_message_is_stored_once_under_replay_and_an_overlapping_pull() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], false);
assert_eq!(intake.intents(), GUILD_MESSAGES | MESSAGE_CONTENT);
let first = rest(
"100000000000000501",
CHANNEL,
"look at this",
Some("https://cdn.example/photo.png?ex=rest"),
);
let bytes = |_: &str, _: u64| Ok(b"image".to_vec());
assert_eq!(
crate::discord::operations::observed_fragment(gateway(&first), bytes).unwrap(),
crate::discord::operations::observed_fragment(first.clone(), bytes).unwrap()
);
assert_eq!(
intake.handle(Work::Message(gateway(&first))).unwrap(),
Done::Stored
);
assert_eq!(fixture.stored(CHANNEL), [("look at this".to_owned(), 1, 1)]);
assert_eq!(
intake.handle(Work::Message(gateway(&first))).unwrap(),
Done::Stored
);
assert_eq!(fixture.stored(CHANNEL), [("look at this".to_owned(), 1, 1)]);
let second = rest("100000000000000502", CHANNEL, "and this", None);
fixture.serve(CHANNEL, vec![first.clone(), second]);
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0
}
);
assert_eq!(
fixture.stored(CHANNEL),
[
("look at this".to_owned(), 1, 1),
("and this".to_owned(), 0, 1)
]
);
intake.handle(Work::Backfill).unwrap();
let pages = fixture.fake.pages.lock().unwrap().clone();
assert_eq!(
pages[0].1.after,
Some(FLOOR),
"the first pull reads forward from the channel's floor"
);
assert!(pages
.iter()
.any(|(_, request)| request.after == Some(100000000000000502)));
assert_eq!(fixture.stored(CHANNEL).len(), 2);
}
#[test]
fn only_configured_channels_and_asked_for_dms_are_stored() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], false);
let elsewhere = rest("100000000000000601", CHANNEL + 1, "elsewhere", None);
assert_eq!(
intake.handle(Work::Message(gateway(&elsewhere))).unwrap(),
Done::Ignored
);
let dm = rest("100000000000000602", CHANNEL + 2, "a DM", None);
assert_eq!(
intake.handle(Work::Message(dm.clone())).unwrap(),
Done::Ignored
);
assert!(fixture.stored(CHANNEL + 2).is_empty());
let mut intake = fixture.intake(&[CHANNEL], true);
assert_eq!(
intake.intents(),
GUILD_MESSAGES | MESSAGE_CONTENT | DIRECT_MESSAGES
);
assert_eq!(intake.handle(Work::Message(dm)).unwrap(), Done::Stored);
assert_eq!(fixture.stored(CHANNEL + 2), [("a DM".to_owned(), 0, 1)]);
let mut restarted = fixture.intake(&[CHANNEL], true);
restarted.handle(Work::Backfill).unwrap();
let pulled: Vec<String> = fixture
.fake
.pages
.lock()
.unwrap()
.iter()
.map(|(channel, _)| channel.clone())
.collect();
assert!(pulled.contains(&CHANNEL.to_string()));
assert!(pulled.contains(&(CHANNEL + 2).to_string()));
assert!(!pulled.contains(&(CHANNEL + 1).to_string()));
assert_eq!(fixture.intake(&[], true).intents(), DIRECT_MESSAGES);
}
#[test]
fn a_message_whose_attachment_cannot_be_fetched_waits_for_the_backfill() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], false);
let message = rest(
"100000000000000701",
CHANNEL,
"with a file",
Some("https://cdn.example/photo.png?ex=rest"),
);
let mut live = gateway(&message);
live["attachments"][0]["url"] = json!("https://cdn.example/expired");
assert!(intake.handle(Work::Message(live)).is_err());
assert!(fixture.stored(CHANNEL).is_empty());
fixture.serve(CHANNEL, vec![message]);
intake.handle(Work::Backfill).unwrap();
assert_eq!(fixture.stored(CHANNEL), [("with a file".to_owned(), 1, 1)]);
}
#[test]
fn a_channel_is_stored_from_its_floor_and_a_dm_channel_from_its_first_message() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], true);
fixture.serve(
CHANNEL,
vec![
rest(&(FLOOR - 2).to_string(), CHANNEL, "history", None),
rest(&(FLOOR - 1).to_string(), CHANNEL, "more history", None),
rest(&(FLOOR + 1).to_string(), CHANNEL, "while down", None),
],
);
let dm = CHANNEL + 2;
let heard = rest(&(FLOOR + 10).to_string(), dm, "hello bot", None);
fixture.serve(
dm,
vec![
rest(&(FLOOR + 5).to_string(), dm, "an old DM", None),
heard.clone(),
rest(&(FLOOR + 11).to_string(), dm, "while down", None),
],
);
assert_eq!(intake.handle(Work::Message(heard)).unwrap(), Done::Stored);
let mut restarted = fixture.intake(&[CHANNEL], true);
restarted.handle(Work::Backfill).unwrap();
assert_eq!(fixture.stored(CHANNEL), [("while down".to_owned(), 0, 1)]);
assert_eq!(
fixture.stored(dm),
[
("hello bot".to_owned(), 0, 1),
("while down".to_owned(), 0, 1)
]
);
}
#[test]
fn nothing_at_or_below_a_floor_is_stored_on_any_path() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], true);
let dm = CHANNEL + 2;
fixture.serve(
CHANNEL,
vec![
rest(&(FLOOR - 2).to_string(), CHANNEL, "history", None),
rest(&(FLOOR - 1).to_string(), CHANNEL, "more history", None),
rest(&(FLOOR + 1).to_string(), CHANNEL, "after", None),
],
);
let heard = rest(&(FLOOR + 10).to_string(), dm, "hello bot", None);
fixture.serve(
dm,
vec![
rest(&(FLOOR + 5).to_string(), dm, "an old DM", None),
heard.clone(),
],
);
assert_eq!(intake.handle(Work::Message(heard)).unwrap(), Done::Stored);
for _ in 0..3 {
intake.handle(Work::Backfill).unwrap();
}
let reconciled = |channel: u64| {
fixture
.fake
.pages
.lock()
.unwrap()
.iter()
.any(|(id, request)| {
*id == channel.to_string()
&& request.after.is_none()
&& request.before.is_none()
})
};
assert!(reconciled(CHANNEL) && reconciled(dm));
assert_eq!(fixture.stored(CHANNEL), [("after".to_owned(), 0, 1)]);
assert_eq!(fixture.stored(dm), [("hello bot".to_owned(), 0, 1)]);
let edited = |id: u64, channel: u64, content: &str| {
let mut message = rest(&id.to_string(), channel, content, None);
message["edited_timestamp"] = json!("2026-09-26T09:00:00Z");
message
};
let history = gateway(&edited(FLOOR - 1, CHANNEL, "more history, edited"));
assert_eq!(intake.handle(Work::Update(history)).unwrap(), Done::Ignored);
let old_dm = edited(FLOOR + 5, dm, "an old DM, edited");
assert_eq!(intake.handle(Work::Update(old_dm)).unwrap(), Done::Ignored);
assert_eq!(fixture.stored(CHANNEL), [("after".to_owned(), 0, 1)]);
assert_eq!(fixture.stored(dm), [("hello bot".to_owned(), 0, 1)]);
let other = CHANNEL + 3;
fixture.serve(
other,
vec![
rest(&(FLOOR + 12).to_string(), other, "long ago", None),
rest(&(FLOOR + 13).to_string(), other, "later", None),
],
);
let old = edited(FLOOR + 12, other, "long ago, edited");
assert_eq!(intake.handle(Work::Update(old)).unwrap(), Done::Ignored);
intake.handle(Work::Backfill).unwrap();
assert!(fixture.stored(other).is_empty());
assert!(!fixture.state.join("dms").join(other.to_string()).exists());
}
#[test]
fn older_coverage_does_not_lower_a_floor() {
let fixture = Fixture::new();
let pulled = vec![
rest(&(FLOOR - 5).to_string(), CHANNEL, "pulled", None),
rest(&(FLOOR - 4).to_string(), CHANNEL, "pulled too", None),
];
fixture.serve(CHANNEL, pulled.clone());
fixture
.discord
.pull_channel(&CHANNEL.to_string(), None, &mut fixture.fake.clone())
.unwrap();
let mut later = pulled;
later.extend([
rest(&(FLOOR - 2).to_string(), CHANNEL, "history", None),
rest(&(FLOOR - 1).to_string(), CHANNEL, "more history", None),
rest(&(FLOOR + 1).to_string(), CHANNEL, "after", None),
]);
fixture.serve(CHANNEL, later);
let mut intake = fixture.intake(&[CHANNEL], false);
intake.handle(Work::Backfill).unwrap();
intake.handle(Work::Backfill).unwrap();
assert_eq!(
fixture.stored(CHANNEL),
[
("pulled".to_owned(), 0, 1),
("pulled too".to_owned(), 0, 1),
("after".to_owned(), 0, 1)
]
);
}
#[test]
fn a_floored_pull_records_where_its_channel_is_heard_from_once() {
let fixture = Fixture::new();
let pull = || {
fixture
.discord
.pull_channel(&CHANNEL.to_string(), Some(FLOOR), &mut fixture.fake.clone())
.unwrap()
};
let first = pull();
assert_eq!(first.observations, 0);
assert!(first.commit.is_some(), "the channel is recorded as heard");
assert!(pull().commit.is_none(), "and only once");
let history = rest(&FLOOR.to_string(), CHANNEL, "history", None);
let error = fixture
.discord
.observe(gateway(&history), FLOOR, &mut fixture.fake.clone())
.unwrap_err();
assert!(
format!("{error:#}").contains("at or below where intake reads channel"),
"{error:#}"
);
assert!(fixture.stored(CHANNEL).is_empty());
}
#[test]
fn an_unreadable_floor_passes_its_channel_over() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], true);
let dm = CHANNEL + 2;
for (kind, channel, floor) in [("channels", CHANNEL, "not a message id"), ("dms", dm, "")] {
let directory = fixture.state.join(kind);
std::fs::create_dir_all(&directory).unwrap();
std::fs::write(directory.join(channel.to_string()), floor).unwrap();
}
for channel in [CHANNEL, dm] {
fixture.serve(
channel,
vec![
rest(&(FLOOR - 1).to_string(), channel, "history", None),
rest(&(FLOOR + 1).to_string(), channel, "after", None),
],
);
}
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 2,
failed: 2,
more: false,
unstored: 0
}
);
let live = gateway(&rest(&(FLOOR + 2).to_string(), CHANNEL, "live", None));
assert!(intake.handle(Work::Message(live)).is_err());
let live_dm = rest(&(FLOOR + 3).to_string(), dm, "live DM", None);
assert!(intake.handle(Work::Message(live_dm)).is_err());
assert!(fixture.fake.pages.lock().unwrap().is_empty());
assert!(fixture.stored(CHANNEL).is_empty());
assert!(fixture.stored(dm).is_empty());
assert_eq!(
std::fs::read_to_string(fixture.state.join("channels").join(CHANNEL.to_string()))
.unwrap(),
"not a message id"
);
}
#[test]
fn a_failed_live_write_is_fetched_again_by_its_id() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], false);
let mut messages: Vec<Value> = (1..=160)
.map(|n| rest(&(FLOOR + n).to_string(), CHANNEL, &format!("m{n}"), None))
.collect();
fixture.serve(CHANNEL, messages.clone());
let caught_up = Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0,
};
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up);
let mut edited = rest(
&(FLOOR + 101).to_string(),
CHANNEL,
"m101, edited",
Some("https://cdn.example/photo.png?ex=rest"),
);
edited["edited_timestamp"] = json!("2026-09-26T09:00:00Z");
let mut live = gateway(&edited);
live["attachments"][0]["url"] = json!("https://cdn.example/expired");
assert!(intake.handle(Work::Update(live)).is_err());
messages[100] = edited;
fixture.serve(CHANNEL, messages);
let photo = fixture
.fake
.attachments
.lock()
.unwrap()
.remove("https://cdn.example/photo.png?ex=rest")
.unwrap();
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 1
}
);
fixture
.fake
.attachments
.lock()
.unwrap()
.insert("https://cdn.example/photo.png?ex=rest".to_owned(), photo);
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up);
assert!(fixture
.stored(CHANNEL)
.contains(&("m101, edited".to_owned(), 1, 1)));
let fetched = || {
fixture
.fake
.fetched
.lock()
.unwrap()
.iter()
.filter(|(channel, id)| *channel == CHANNEL.to_string() && *id == FLOOR + 101)
.count()
};
assert_eq!(fetched(), 2);
intake.handle(Work::Backfill).unwrap();
assert_eq!(fetched(), 2, "a stored message is not fetched again");
let mut gone = gateway(&rest(
&(FLOOR + 161).to_string(),
CHANNEL,
"gone",
Some("https://cdn.example/expired"),
));
gone["attachments"][0]["url"] = json!("https://cdn.example/expired");
assert!(intake.handle(Work::Message(gone)).is_err());
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up);
assert!(std::fs::read_dir(fixture.state.join(UNSTORED))
.unwrap()
.next()
.is_none());
}
#[test]
fn a_kept_message_in_a_channel_no_longer_served_is_let_go() {
struct Refusing {
inner: Fake,
refused: String,
}
impl Source for Refusing {
fn page(&mut self, channel_id: &str, request: PageRequest) -> Result<Vec<Value>> {
anyhow::ensure!(channel_id != self.refused, "403 Missing Access");
self.inner.page(channel_id, request)
}
fn message(&mut self, channel_id: &str, message_id: u64) -> Result<Option<Value>> {
anyhow::ensure!(channel_id != self.refused, "403 Missing Access");
self.inner.message(channel_id, message_id)
}
fn attachment(&mut self, url: &str, limit: u64) -> Result<Vec<u8>> {
self.inner.attachment(url, limit)
}
}
let fixture = Fixture::new();
let other = CHANNEL + 7;
let mut intake = fixture.intake(&[CHANNEL, other], false);
fixture.serve(
CHANNEL,
vec![rest(&(FLOOR + 1).to_string(), CHANNEL, "a", None)],
);
let mut live = gateway(&rest(
&(FLOOR + 2).to_string(),
other,
"b",
Some("https://cdn.example/expired"),
));
live["attachments"][0]["url"] = json!("https://cdn.example/expired");
assert!(intake.handle(Work::Message(live)).is_err());
assert!(fixture
.state
.join(UNSTORED)
.read_dir()
.unwrap()
.next()
.is_some());
let mut restarted = Intake::new(
fixture.discord.clone(),
Box::new(Refusing {
inner: fixture.fake.clone(),
refused: other.to_string(),
}),
vec![NonZeroU64::new(CHANNEL).unwrap()],
false,
fixture.state.clone(),
FLOOR,
);
let caught_up = Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0,
};
assert_eq!(restarted.handle(Work::Backfill).unwrap(), caught_up);
assert!(fixture
.state
.join(UNSTORED)
.read_dir()
.unwrap()
.next()
.is_none());
assert!(!fixture
.fake
.pages
.lock()
.unwrap()
.iter()
.any(|(channel, _)| *channel == other.to_string()));
assert!(!fixture
.fake
.fetched
.lock()
.unwrap()
.iter()
.any(|(channel, _)| *channel == other.to_string()));
assert_eq!(restarted.handle(Work::Backfill).unwrap(), caught_up);
}
#[test]
fn a_live_event_whose_floor_cannot_be_had_is_fetched_again_by_its_id() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], true);
let mut messages: Vec<Value> = (1..=160)
.map(|n| rest(&(FLOOR + n).to_string(), CHANNEL, &format!("m{n}"), None))
.collect();
fixture.serve(CHANNEL, messages.clone());
intake.handle(Work::Backfill).unwrap();
let floor_file = fixture.state.join("channels").join(CHANNEL.to_string());
let good = std::fs::read_to_string(&floor_file).unwrap();
std::fs::write(&floor_file, "").unwrap();
let mut edited = rest(&(FLOOR + 101).to_string(), CHANNEL, "m101, edited", None);
edited["edited_timestamp"] = json!("2026-09-26T09:00:00Z");
assert!(intake.handle(Work::Update(gateway(&edited))).is_err());
messages[100] = edited;
fixture.serve(CHANNEL, messages);
let dm = CHANNEL + 2;
let first = rest(&(FLOOR + 170).to_string(), dm, "hello bot", None);
fixture.serve(
dm,
vec![
rest(&(FLOOR + 165).to_string(), dm, "an old DM", None),
first.clone(),
],
);
let dms = fixture.state.join("dms");
std::fs::write(&dms, "").unwrap();
assert!(intake.handle(Work::Message(first)).is_err());
assert_eq!(
fixture.state.join(UNSTORED).read_dir().unwrap().count(),
2,
"both are kept by id"
);
let Done::Backfilled {
failed, unstored, ..
} = intake.handle(Work::Backfill).unwrap()
else {
panic!("a backfill");
};
assert_eq!((failed, unstored), (2, 2));
assert!(fixture
.stored(CHANNEL)
.iter()
.all(|(content, _, _)| content != "m101, edited"));
std::fs::write(&floor_file, good).unwrap();
std::fs::remove_file(&dms).unwrap();
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 2,
failed: 0,
more: false,
unstored: 0
}
);
assert!(fixture
.stored(CHANNEL)
.contains(&("m101, edited".to_owned(), 0, 1)));
assert_eq!(fixture.stored(dm), [("hello bot".to_owned(), 0, 1)]);
assert_eq!(
std::fs::read_to_string(dms.join(dm.to_string())).unwrap(),
(FLOOR + 169).to_string()
);
}
#[test]
fn the_earliest_new_message_of_a_dm_channel_begins_it_in_any_order() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[], true);
let dms = fixture.state.join("dms");
let unstored = fixture.state.join(UNSTORED);
let floor_of = |dm: u64| {
std::fs::read_to_string(dms.join(dm.to_string()))
.unwrap()
.parse::<u64>()
.unwrap()
};
let serve_dm = |dm: u64, base: u64, contents: [&str; 2]| {
let mut messages = vec![rest(&(FLOOR + base - 5).to_string(), dm, "an old DM", None)];
messages.extend(
(base..)
.zip(contents)
.map(|(id, content)| rest(&(FLOOR + id).to_string(), dm, content, None)),
);
fixture.serve(dm, messages);
};
let caught_up = |channels| Done::Backfilled {
channels,
failed: 0,
more: false,
unstored: 0,
};
let both =
|first: &str, second: &str| vec![(first.to_owned(), 0, 1), (second.to_owned(), 0, 1)];
let dm = CHANNEL + 2;
serve_dm(dm, 170, ["hello bot", "are you there"]);
std::fs::create_dir_all(&fixture.state).unwrap();
std::fs::write(&dms, "").unwrap();
for id in [170, 171] {
let message = rest(&(FLOOR + id).to_string(), dm, "live", None);
assert!(intake.handle(Work::Message(message)).is_err());
}
std::fs::remove_file(&dms).unwrap();
let earlier = unstored.join(format!("{dm}-{}-{NEW_DM}", FLOOR + 170));
let aside = fixture.state.join("aside");
std::fs::rename(&earlier, &aside).unwrap();
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up(1));
assert_eq!(fixture.stored(dm), [("are you there".to_owned(), 0, 1)]);
assert_eq!(floor_of(dm), FLOOR + 170);
std::fs::rename(&aside, &earlier).unwrap();
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up(1));
assert_eq!(fixture.stored(dm), both("hello bot", "are you there"));
assert_eq!(floor_of(dm), FLOOR + 169);
assert!(unstored.read_dir().unwrap().next().is_none());
let dm = CHANNEL + 3;
serve_dm(dm, 270, ["hello bot", "are you there"]);
std::fs::rename(&dms, &aside).unwrap();
std::fs::write(&dms, "").unwrap();
let first = rest(&(FLOOR + 270).to_string(), dm, "hello bot", None);
assert!(intake.handle(Work::Message(first)).is_err());
std::fs::remove_file(&dms).unwrap();
std::fs::rename(&aside, &dms).unwrap();
let second = rest(&(FLOOR + 271).to_string(), dm, "are you there", None);
assert_eq!(intake.handle(Work::Message(second)).unwrap(), Done::Stored);
assert_eq!(floor_of(dm), FLOOR + 270);
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up(2));
assert_eq!(fixture.stored(dm), both("hello bot", "are you there"));
assert_eq!(floor_of(dm), FLOOR + 269);
let dm = CHANNEL + 4;
serve_dm(dm, 380, ["one", "two"]);
for (id, content) in [(381, "two"), (380, "one")] {
let message = rest(&(FLOOR + id).to_string(), dm, content, None);
assert_eq!(intake.handle(Work::Message(message)).unwrap(), Done::Stored);
}
assert_eq!(floor_of(dm), FLOOR + 379);
assert_eq!(fixture.stored(dm), both("one", "two"));
let mut edited = rest(&(FLOOR + 375).to_string(), dm, "an old DM, edited", None);
edited["edited_timestamp"] = json!("2026-09-26T09:00:00Z");
assert_eq!(intake.handle(Work::Update(edited)).unwrap(), Done::Ignored);
assert_eq!(floor_of(dm), FLOOR + 379);
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up(3));
for dm in [CHANNEL + 2, CHANNEL + 3, CHANNEL + 4] {
assert_eq!(fixture.stored(dm).len(), 2, "no history of {dm}");
}
}
#[test]
fn a_kept_new_dm_is_marked_by_its_name_alone() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[], true);
let dms = fixture.state.join("dms");
let unstored = fixture.state.join(UNSTORED);
let dm = CHANNEL + 2;
let first = rest(&(FLOOR + 170).to_string(), dm, "hello bot", None);
let second = rest(&(FLOOR + 171).to_string(), dm, "are you there", None);
fixture.serve(
dm,
vec![
rest(&(FLOOR + 165).to_string(), dm, "an old DM", None),
first.clone(),
second.clone(),
],
);
std::fs::create_dir_all(&fixture.state).unwrap();
std::fs::write(&dms, "").unwrap();
assert!(intake.handle(Work::Message(first.clone())).is_err());
assert!(intake.handle(Work::Update(first)).is_err());
let new = unstored.join(format!("{dm}-{}-{NEW_DM}", FLOOR + 170));
let edit = unstored.join(format!("{dm}-{}", FLOOR + 170));
assert_eq!(std::fs::metadata(&new).unwrap().len(), 0);
assert_eq!(std::fs::metadata(&edit).unwrap().len(), 0);
std::fs::remove_file(&dms).unwrap();
assert_eq!(intake.handle(Work::Message(second)).unwrap(), Done::Stored);
let floor_of = || {
std::fs::read_to_string(dms.join(dm.to_string()))
.unwrap()
.parse::<u64>()
.unwrap()
};
assert_eq!(floor_of(), FLOOR + 170);
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0
}
);
assert_eq!(
fixture.stored(dm),
[
("hello bot".to_owned(), 0, 1),
("are you there".to_owned(), 0, 1)
]
);
assert_eq!(floor_of(), FLOOR + 169);
assert!(unstored.read_dir().unwrap().next().is_none());
}
fn serve_gap(fixture: &Fixture, dm: u64, through: u64) {
let mut messages = vec![rest(&(FLOOR + 165).to_string(), dm, "an old DM", None)];
messages.extend(
(170..=through).map(|n| rest(&(FLOOR + n).to_string(), dm, &format!("m{n}"), None)),
);
fixture.serve(dm, messages);
}
fn gap(ids: std::ops::RangeInclusive<u64>) -> Vec<(String, usize, usize)> {
ids.map(|n| (format!("m{n}"), 0, 1)).collect()
}
fn pulled_after(fixture: &Fixture, channel: u64, after: u64) -> bool {
fixture
.fake
.pages
.lock()
.unwrap()
.iter()
.any(|(id, request)| *id == channel.to_string() && request.after == Some(after))
}
fn keep_first(fixture: &Fixture, intake: &mut Intake, dm: u64) {
let dms = fixture.state.join("dms");
std::fs::create_dir_all(&fixture.state).unwrap();
std::fs::write(&dms, "").unwrap();
let first = rest(&(FLOOR + 170).to_string(), dm, "m170", None);
assert!(intake.handle(Work::Message(first)).is_err());
std::fs::remove_file(&dms).unwrap();
}
#[test]
fn a_dm_channel_a_kept_message_begins_is_pulled_before_caught_up() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[], true);
let dm = CHANNEL + 2;
serve_gap(&fixture, dm, 179);
keep_first(&fixture, &mut intake, dm);
assert!(!fixture.state.join("dms").exists());
let done = intake.handle(Work::Backfill).unwrap();
assert_eq!(fixture.stored(dm), gap(170..=179));
assert_eq!(
done,
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0
}
);
assert_eq!(
std::fs::read_to_string(fixture.state.join("dms").join(dm.to_string())).unwrap(),
(FLOOR + 169).to_string()
);
assert!(pulled_after(&fixture, dm, FLOOR + 169));
}
#[test]
fn a_dm_floor_a_kept_message_lowers_is_pulled_before_caught_up() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[], true);
let dm = CHANNEL + 2;
serve_gap(&fixture, dm, 180);
keep_first(&fixture, &mut intake, dm);
let later = rest(&(FLOOR + 180).to_string(), dm, "m180", None);
assert_eq!(intake.handle(Work::Message(later)).unwrap(), Done::Stored);
let dms = fixture.state.join("dms");
let floor = dms.join(dm.to_string());
assert_eq!(
std::fs::read_to_string(&floor).unwrap(),
(FLOOR + 179).to_string()
);
let staging = dms.join(format!(".{dm}"));
std::fs::create_dir(&staging).unwrap();
let done = intake.handle(Work::Backfill);
std::fs::remove_dir(&staging).unwrap();
assert_eq!(
done.unwrap(),
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 1
}
);
assert!(pulled_after(&fixture, dm, FLOOR + 179));
assert_eq!(fixture.stored(dm), gap(180..=180));
let done = intake.handle(Work::Backfill).unwrap();
assert_eq!(fixture.stored(dm), gap(170..=180));
assert_eq!(
done,
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0
}
);
assert_eq!(
std::fs::read_to_string(&floor).unwrap(),
(FLOOR + 169).to_string()
);
assert!(pulled_after(&fixture, dm, FLOOR + 169));
}
#[test]
fn a_backfill_that_cannot_record_the_account_has_not_caught_up() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[], true);
let dm = CHANNEL + 2;
serve_gap(&fixture, dm, 171);
let key = fixture._directory.path().join("intake.key");
let aside = fixture._directory.path().join("intake.key.aside");
std::fs::rename(&key, &aside).unwrap();
let account = intake.handle(Work::Account("100000000000000800".to_owned()));
let alone = intake.handle(Work::Backfill);
keep_first(&fixture, &mut intake, dm);
let kept = intake.handle(Work::Backfill);
std::fs::rename(&aside, &key).unwrap();
assert!(account.is_err());
assert_eq!(
alone.unwrap(),
Done::Backfilled {
channels: 1,
failed: 1,
more: false,
unstored: 0
}
);
assert_eq!(
kept.unwrap(),
Done::Backfilled {
channels: 1,
failed: 1,
more: false,
unstored: 1
}
);
assert!(fixture.fake.fetched.lock().unwrap().is_empty());
assert!(fixture.fake.pages.lock().unwrap().is_empty());
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0
}
);
assert_eq!(fixture.stored(dm), gap(170..=171));
}
#[test]
fn a_kept_message_outlives_a_channel_whose_history_cannot_be_read() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], false);
let mut messages: Vec<Value> = (1..=160)
.map(|n| rest(&(FLOOR + n).to_string(), CHANNEL, &format!("m{n}"), None))
.collect();
fixture.serve(CHANNEL, messages.clone());
let caught_up = Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0,
};
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up);
let mut edited = rest(
&(FLOOR + 101).to_string(),
CHANNEL,
"m101, edited",
Some("https://cdn.example/photo.png?ex=rest"),
);
edited["edited_timestamp"] = json!("2026-09-26T09:00:00Z");
let mut live = gateway(&edited);
live["attachments"][0]["url"] = json!("https://cdn.example/expired");
assert!(intake.handle(Work::Update(live)).is_err());
messages[100] = edited;
fixture.serve(CHANNEL, messages);
fixture
.fake
.unreadable
.lock()
.unwrap()
.insert(CHANNEL.to_string());
let marker = fixture
.state
.join(UNSTORED)
.join(format!("{CHANNEL}-{}", FLOOR + 101));
for _ in 0..2 {
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 1
}
);
assert!(marker.exists(), "the kept message stays kept");
}
fixture.fake.unreadable.lock().unwrap().clear();
assert_eq!(intake.handle(Work::Backfill).unwrap(), caught_up);
assert!(fixture
.stored(CHANNEL)
.contains(&("m101, edited".to_owned(), 1, 1)));
assert!(!marker.exists());
}
#[test]
fn dm_channels_that_cannot_be_listed_fail_the_backfill() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[], true);
let dm = CHANNEL + 2;
let heard = rest(&(FLOOR + 10).to_string(), dm, "hello bot", None);
fixture.serve(dm, vec![heard.clone()]);
assert_eq!(intake.handle(Work::Message(heard)).unwrap(), Done::Stored);
let dms = fixture.state.join("dms");
let aside = fixture.state.join("dms.aside");
std::fs::rename(&dms, &aside).unwrap();
std::fs::write(&dms, "").unwrap();
let done = intake.handle(Work::Backfill);
std::fs::remove_file(&dms).unwrap();
std::fs::rename(&aside, &dms).unwrap();
assert_eq!(
done.unwrap(),
Done::Backfilled {
channels: 1,
failed: 1,
more: false,
unstored: 0
}
);
}
#[test]
fn a_long_gap_is_closed_by_one_backfill() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], false);
let messages: Vec<Value> = (1..=350)
.map(|n| rest(&(FLOOR + n).to_string(), CHANNEL, &format!("m{n}"), None))
.collect();
fixture.serve(CHANNEL, messages[..1].to_vec());
intake.handle(Work::Backfill).unwrap();
assert_eq!(fixture.stored(CHANNEL).len(), 1);
fixture.serve(CHANNEL, messages);
assert_eq!(
intake.handle(Work::Backfill).unwrap(),
Done::Backfilled {
channels: 1,
failed: 0,
more: false,
unstored: 0
}
);
assert_eq!(fixture.stored(CHANNEL).len(), 350);
}
#[test]
fn failures_bring_the_next_backfill_forward_and_back_off() {
let start = Instant::now();
let mut retry = Retry::default();
retry.soon(start);
assert_eq!(retry.due, Some(start + RETRY_FIRST));
retry.soon(start + Duration::from_secs(30));
assert_eq!(retry.due, Some(start + RETRY_FIRST));
let mut waits = Vec::new();
for _ in 0..6 {
retry.again(start);
waits.push((retry.due.unwrap() - start).as_secs());
}
assert_eq!(waits, [60, 120, 240, 480, 900, 900]);
retry.caught_up();
assert_eq!(retry.due, None);
let (sender, inbox) = std::sync::mpsc::channel();
retry.due = Some(Instant::now());
assert!(matches!(retry.next(&inbox), Some(Work::Backfill)));
assert_eq!(retry.due, None);
sender.send(Work::Account("1".to_owned())).unwrap();
assert!(matches!(retry.next(&inbox), Some(Work::Account(_))));
drop(sender);
assert!(retry.next(&inbox).is_none());
}
struct Slow;
impl Source for Slow {
fn page(&mut self, _: &str, _: PageRequest) -> Result<Vec<Value>> {
std::thread::sleep(Duration::from_secs(3));
Ok(Vec::new())
}
fn message(&mut self, _: &str, _: u64) -> Result<Option<Value>> {
std::thread::sleep(Duration::from_secs(3));
Ok(None)
}
fn attachment(&mut self, _: &str, _: u64) -> Result<Vec<u8>> {
Ok(Vec::new())
}
}
#[tokio::test]
async fn stopping_is_bounded_while_a_backfill_runs() {
let fixture = Fixture::new();
let intake = Intake::new(
fixture.discord.clone(),
Box::new(Slow),
vec![NonZeroU64::new(CHANNEL).unwrap()],
false,
fixture.state.clone(),
FLOOR,
);
let mut worker = start(intake);
worker.send(Work::Backfill);
worker.send(Work::Backfill);
tokio::time::sleep(Duration::from_millis(100)).await;
let stopping = Instant::now();
worker.stop(Duration::from_millis(200)).await;
assert!(stopping.elapsed() < Duration::from_secs(1));
}
#[test]
fn the_bot_account_is_recorded() {
let fixture = Fixture::new();
let mut intake = fixture.intake(&[CHANNEL], false);
assert_eq!(
intake
.handle(Work::Account("100000000000000800".to_owned()))
.unwrap(),
Done::Recorded
);
assert!(intake
.handle(Work::Account("not a snowflake".to_owned()))
.is_err());
}
}