use std::fs::{self, OpenOptions};
use std::io::{BufReader, BufWriter, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
use std::{fs::File, io::Read};
use bstr::BString;
use chrono::{DateTime, Utc};
use rayon::iter::{IndexedParallelIterator, ParallelIterator};
use rayon::slice::ParallelSlice;
use thiserror::Error;
use crate::util::crc32::{self, CRC_SEED};
use crate::util::echomail::EchomailAddress;
use self::jhr_header::JhrHeaderInfo;
use self::last_read_storage::JamLastReadStorage;
use self::msg_header::{JamMessageHeader, MessageSubfield, SubfieldType};
use self::pack::EMPTY_SLOT;
pub mod jhr_header;
pub mod last_read_storage;
pub mod msg_header;
pub mod pack;
pub mod raw;
pub mod verify;
#[derive(Error, Debug)]
#[non_exhaustive]
pub enum JamError {
#[error("Invalid header signature (needs to start with 'JAM\\0')")]
InvalidHeaderSignature,
#[error("Index file corrupted")]
IndexFileCorrupted,
#[error("Unsupported message header revision: {0}")]
UnsupportedMessageHeaderRevision(u16),
#[error("Invalid subfield length {0} for sub field {1}")]
InvalidSubfieldLength(u32, usize),
#[error("Message number {0} out of range. Valid range is {1}..={2}")]
MessageNumberOutOfRange(u32, u32, u32),
#[error("Message was deleted")]
MessageDeleted,
#[error("Index file corrupt at record {0} (file length: {1})")]
IndexFileCorrupt(u64, u64),
#[error("Index record {0} is half empty")]
InvalidIndexRecord(u64),
#[error("Index record for message {0} points to a header numbered {1}")]
IndexMessageNumberMismatch(u32, u32),
#[error("Message base is full, offset {0} does not fit the 32 bit field the format provides")]
MessageBaseFull(u64),
#[error("Message numbers are exhausted")]
MessageNumbersExhausted,
#[error("Message text of {0} bytes exceeds the 4 GB a JAM message base can address")]
MessageTooLarge(usize),
#[error("Message text at offset {0} runs {1} bytes past the end of the text file")]
TextOutOfBounds(u64, u64),
#[error("Lastread file corrupted, {0} bytes is not a whole number of records")]
LastReadFileCorrupted(u64),
#[error(
"a shared message base lock cannot be upgraded, take the exclusive lock before reading"
)]
LockUpgrade,
#[error("subfield block of {0} bytes exceeds the {1} byte limit")]
SubfieldBlockTooLarge(u32, u32),
#[error("subfield block truncated, expected {0} bytes, got {1}")]
SubfieldBlockTruncated(u32, usize),
#[error("subfield of {0} bytes exceeds the {1} byte limit")]
SubfieldTooLarge(u32, u32),
}
fn offset_u32(value: u64) -> crate::Result<u32> {
u32::try_from(value).map_err(|_| JamError::MessageBaseFull(value).into())
}
mod extensions {
pub const HEADER_DATA: &str = "jhr";
pub const TEXT_DATA: &str = "jdt";
pub const MESSAGE_INDEX: &str = "jdx";
pub const LASTREAD_INFO: &str = "jlr";
pub const LOCK_FILE: &str = "jamlock";
}
const JAM_SIGNATURE: [u8; 4] = [b'J', b'A', b'M', 0];
pub(crate) const INDEX_RECORD_SIZE: usize = 8;
pub(crate) const LASTREAD_RECORD_SIZE: usize = 16;
pub(crate) const LASTREAD_DELETED: [u8; 8] = [0xFF; 8];
const PARALLEL_SEARCH_THRESHOLD: usize = 8192;
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum LockMode {
Shared,
Exclusive,
}
pub struct JamMessageBase {
file_name: PathBuf,
header_info: JhrHeaderInfo,
index_records: u32,
lock_file: File,
lock_depth: u32,
lock_mode: LockMode,
}
impl JamMessageBase {
pub fn open<P: AsRef<Path>>(file_name: P) -> crate::Result<Self> {
let file_name = file_name.as_ref();
let lock_path = file_name.with_extension(extensions::LOCK_FILE);
let lock_file = OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(lock_path)?;
lock_file.lock()?;
pack::recover_interrupted_pack(file_name)?;
let header_file_name = file_name.with_extension(extensions::HEADER_DATA);
let mut file = File::open(&header_file_name).inspect_err(|err| {
log::error!("Error opening message base {}: {err}", file_name.display());
})?;
let header_info = JhrHeaderInfo::load(&mut file)?;
let mut base = Self {
file_name: file_name.into(),
header_info,
index_records: 0,
lock_file,
lock_depth: 1,
lock_mode: LockMode::Exclusive,
};
base.index_records = base.count_index_records()?;
base.unlock();
Ok(base)
}
fn refresh_state(&mut self) -> crate::Result<()> {
let header_file_name = self.file_name.with_extension(extensions::HEADER_DATA);
let mut header = File::open(header_file_name)?;
self.header_info = JhrHeaderInfo::load(&mut header)?;
self.index_records = self.count_index_records()?;
Ok(())
}
fn count_index_records(&self) -> crate::Result<u32> {
let path = self.file_name.with_extension(extensions::MESSAGE_INDEX);
let len = match fs::metadata(&path) {
Ok(meta) => meta.len(),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => 0,
Err(err) => return Err(err.into()),
};
if !len.is_multiple_of(INDEX_RECORD_SIZE as u64) {
return Err(JamError::IndexFileCorrupted.into());
}
Ok((len / INDEX_RECORD_SIZE as u64) as u32)
}
pub fn path(&self) -> &Path {
&self.file_name
}
pub fn info(&self) -> &JhrHeaderInfo {
&self.header_info
}
pub fn mod_counter(&self) -> u32 {
self.header_info.mod_counter
}
pub fn lowest_message_number(&self) -> u32 {
self.header_info.base_msg_num
}
pub fn highest_message_number(&self) -> u32 {
self.header_info
.base_msg_num
.saturating_add(self.index_records)
.saturating_sub(1)
}
pub fn index_records(&self) -> u32 {
self.index_records
}
pub fn active_messages(&self) -> u32 {
self.header_info.active_msgs
}
pub fn needs_password(&self) -> bool {
self.header_info.password_crc != CRC_SEED
}
pub fn is_password_valid(&self, password: &BString) -> bool {
self.header_info.password_crc == CRC_SEED
|| self.header_info.password_crc == Self::crc(password)
}
pub fn create<P: AsRef<Path>>(file_name: P) -> crate::Result<Self> {
Self::create_with_password_crc(file_name, CRC_SEED)
}
pub fn create_with_password<P: AsRef<Path>>(
file_name: P,
password: &BString,
) -> crate::Result<Self> {
Self::create_with_password_crc(file_name, Self::crc(password))
}
pub fn create_with_password_crc<P: AsRef<Path>>(
file_name: P,
passwordcrc: u32,
) -> crate::Result<Self> {
let file_name = file_name.as_ref();
let lock_file = OpenOptions::new()
.create(true)
.read(true)
.write(true)
.truncate(false)
.open(file_name.with_extension(extensions::LOCK_FILE))?;
lock_file.lock()?;
let header_path = file_name.with_extension(extensions::HEADER_DATA);
let text_path = file_name.with_extension(extensions::TEXT_DATA);
let index_path = file_name.with_extension(extensions::MESSAGE_INDEX);
let lastread_path = file_name.with_extension(extensions::LASTREAD_INFO);
let mut created = Vec::new();
let result: crate::Result<()> = (|| {
JhrHeaderInfo::create(&header_path, passwordcrc)?;
created.push(header_path.clone());
for path in [&text_path, &index_path, &lastread_path] {
OpenOptions::new().create_new(true).write(true).open(path)?;
created.push(path.clone());
}
Ok(())
})();
if result.is_err() {
for path in created.iter().rev() {
let _ = fs::remove_file(path);
}
}
let unlock_result = lock_file.unlock();
drop(lock_file);
result?;
unlock_result?;
Self::open(file_name)
}
pub fn delete_message_base(mut self) -> crate::Result<()> {
self.lock()?;
let file_name = self.file_name.clone();
for extension in [
extensions::HEADER_DATA,
extensions::TEXT_DATA,
extensions::MESSAGE_INDEX,
extensions::LASTREAD_INFO,
] {
fs::remove_file(file_name.with_extension(extension))?;
}
self.unlock();
Ok(())
}
pub fn sync(&self) -> crate::Result<()> {
for extension in [
extensions::TEXT_DATA,
extensions::HEADER_DATA,
extensions::MESSAGE_INDEX,
extensions::LASTREAD_INFO,
] {
let path = self.file_name.with_extension(extension);
match OpenOptions::new().write(true).truncate(false).open(&path) {
Ok(file) => file.sync_all()?,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {}
Err(err) => return Err(err.into()),
}
}
pack::sync_dir(&self.file_name);
Ok(())
}
pub fn lock(&mut self) -> crate::Result<()> {
self.acquire(LockMode::Exclusive, true).map(|_| ())
}
pub fn lock_shared(&mut self) -> crate::Result<()> {
self.acquire(LockMode::Shared, true).map(|_| ())
}
pub fn unlock(&mut self) {
if self.lock_depth == 0 {
return;
}
self.lock_depth -= 1;
if self.lock_depth == 0
&& let Err(err) = self.lock_file.unlock()
{
log::error!("Could not release the message base lock: {err}");
}
}
pub fn try_lock(&mut self) -> crate::Result<bool> {
self.acquire(LockMode::Exclusive, false)
}
pub fn try_lock_shared(&mut self) -> crate::Result<bool> {
self.acquire(LockMode::Shared, false)
}
fn acquire(&mut self, mode: LockMode, blocking: bool) -> crate::Result<bool> {
if self.lock_depth > 0 {
if mode == LockMode::Exclusive && self.lock_mode == LockMode::Shared {
return Err(JamError::LockUpgrade.into());
}
} else {
if !self.take_file_lock(mode, blocking)? {
return Ok(false);
}
if let Err(err) = self.refresh_state() {
let _ = self.lock_file.unlock();
return Err(err);
}
self.lock_mode = mode;
}
self.lock_depth = self
.lock_depth
.checked_add(1)
.ok_or_else(|| crate::Error::jam(0, "message base lock depth overflow"))?;
Ok(true)
}
fn take_file_lock(&self, mode: LockMode, blocking: bool) -> crate::Result<bool> {
if blocking {
match mode {
LockMode::Shared => self.lock_file.lock_shared()?,
LockMode::Exclusive => self.lock_file.lock()?,
}
return Ok(true);
}
let result = match mode {
LockMode::Shared => self.lock_file.try_lock_shared(),
LockMode::Exclusive => self.lock_file.try_lock(),
};
match result {
Ok(()) => Ok(true),
Err(std::fs::TryLockError::WouldBlock) => Ok(false),
Err(std::fs::TryLockError::Error(err)) => Err(err.into()),
}
}
pub fn transaction<T>(
&mut self,
f: impl FnOnce(&mut Self) -> crate::Result<T>,
) -> crate::Result<T> {
self.lock()?;
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| f(self)));
self.unlock();
match result {
Ok(result) => result,
Err(payload) => std::panic::resume_unwind(payload),
}
}
pub fn read_transaction<T>(
&mut self,
f: impl FnOnce(&Self) -> crate::Result<T>,
) -> crate::Result<T> {
self.lock_shared()?;
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| f(self)));
self.unlock();
match result {
Ok(result) => result,
Err(payload) => std::panic::resume_unwind(payload),
}
}
pub fn crc(str: &BString) -> u32 {
let mut str = str.clone();
str.make_ascii_lowercase();
let crc = crc32::checksum(&str);
crc ^ CRC_SEED
}
pub fn write_message(&mut self, message: &JamMessage) -> crate::Result<u32> {
self.transaction(|base| {
let text_path = base.file_name.with_extension(extensions::TEXT_DATA);
let header_path = base.file_name.with_extension(extensions::HEADER_DATA);
let index_path = base.file_name.with_extension(extensions::MESSAGE_INDEX);
let text_len = file_len(&text_path)?;
let header_len = file_len(&header_path)?;
let index_len = file_len(&index_path)?;
let original_header_info = base.header_info.clone();
let original_index_records = base.index_records;
let result = (|| {
let msg_number = base
.highest_message_number()
.checked_add(1)
.ok_or(JamError::MessageNumbersExhausted)?;
Self::append_message(
message,
msg_number,
&text_path,
&header_path,
&index_path,
text_len,
header_len,
)?;
base.index_records = base
.index_records
.checked_add(1)
.ok_or(JamError::MessageNumbersExhausted)?;
if !message.is_deleted() {
base.header_info.active_msgs = base
.header_info
.active_msgs
.checked_add(1)
.ok_or(JamError::MessageNumbersExhausted)?;
}
base.write_jhr_header()?;
Ok(msg_number)
})();
match result {
Ok(msg_number) => Ok(msg_number),
Err(err) => {
rollback(&text_path, text_len);
rollback(&header_path, header_len);
rollback(&index_path, index_len);
base.header_info = original_header_info;
base.index_records = original_index_records;
if let Err(restore_err) = base.store_jhr_header() {
log::error!("Could not restore the JAM base header after a failed write: {restore_err}");
}
Err(err)
}
}
})
}
fn append_message(
message: &JamMessage,
msg_number: u32,
text_path: &Path,
header_path: &Path,
index_path: &Path,
text_offset: u64,
header_offset: u64,
) -> crate::Result<()> {
let mut header = message.create_jam_header();
header.message_number = msg_number;
header.offset = offset_u32(text_offset)?;
header.txt_len = u32::try_from(message.text().len())
.map_err(|_| JamError::MessageTooLarge(message.text().len()))?;
offset_u32(text_offset + header.txt_len as u64)?;
let index_offset = offset_u32(header_offset)?;
let mut text_file = OpenOptions::new()
.create(true)
.append(true)
.open(text_path)?;
text_file.write_all(message.text())?;
text_file.flush()?;
let header_file = OpenOptions::new()
.create(true)
.append(true)
.open(header_path)?;
let mut writer = BufWriter::new(header_file);
header.write(&mut writer)?;
writer.flush()?;
let mut index_file = OpenOptions::new()
.create(true)
.append(true)
.open(index_path)?;
let crc = header.to().map_or(CRC_SEED, Self::crc);
index_file.write_all(&crc.to_le_bytes())?;
index_file.write_all(&index_offset.to_le_bytes())?;
index_file.flush()?;
Ok(())
}
pub fn write_jhr_header(&mut self) -> crate::Result<()> {
self.transaction(|base| base.write_jhr_header_locked())
}
fn write_jhr_header_locked(&mut self) -> crate::Result<()> {
let header_path = self.file_name.with_extension(extensions::HEADER_DATA);
let header_file = OpenOptions::new()
.write(true)
.truncate(false)
.open(header_path)?;
let mut writer = BufWriter::new(header_file);
self.header_info.update(&mut writer)?;
writer.flush()?;
Ok(())
}
fn store_jhr_header(&self) -> crate::Result<()> {
let header_path = self.file_name.with_extension(extensions::HEADER_DATA);
let header_file = OpenOptions::new()
.write(true)
.truncate(false)
.open(header_path)?;
let mut writer = BufWriter::new(header_file);
self.header_info.store(&mut writer)?;
writer.flush()?;
Ok(())
}
pub fn read_jhr_header(&mut self) -> crate::Result<()> {
self.refresh_state()
}
pub fn read_message_text(&self, header: &JamMessageHeader) -> crate::Result<BString> {
let text_file_name = self.file_name.with_extension(extensions::TEXT_DATA);
let mut text_file = File::open(text_file_name)?;
let file_len = text_file.metadata()?.len();
let end = header.offset as u64 + header.txt_len as u64;
if end > file_len {
return Err(JamError::TextOutOfBounds(header.offset as u64, end - file_len).into());
}
text_file.seek(SeekFrom::Start(header.offset as u64))?;
let mut buffer = vec![0; header.txt_len as usize];
text_file.read_exact(&mut buffer)?;
Ok(BString::new(buffer))
}
pub fn messages(&self) -> impl Iterator<Item = crate::Result<JamMessageHeader>> + use<'_> {
let mut reader = IndexedReader::open(self).map_err(Some);
let mut numbers = self.lowest_message_number()..=self.highest_message_number();
std::iter::from_fn(move || {
let reader = match &mut reader {
Ok(reader) => reader,
Err(err) => return err.take().map(Err),
};
loop {
let number = numbers.next()?;
match reader.header(number) {
Ok(Some(header)) => return Some(Ok(header)),
Ok(None) => continue,
Err(err) => return Some(Err(err)),
}
}
})
}
pub fn messages_full(&self) -> impl Iterator<Item = crate::Result<JamMessage>> + use<'_> {
let mut readers = IndexedReader::open(self)
.and_then(|headers| Ok((headers, TextReader::open(self)?)))
.map_err(Some);
let mut numbers = self.lowest_message_number()..=self.highest_message_number();
std::iter::from_fn(move || {
let (headers, text) = match &mut readers {
Ok((headers, text)) => (headers, text),
Err(err) => return err.take().map(Err),
};
loop {
let number = numbers.next()?;
let header = match headers.header(number) {
Ok(Some(header)) => header,
Ok(None) => continue,
Err(err) => return Some(Err(err)),
};
return Some(
text.text(&header)
.map(|text| JamMessage::from_stored(header, text)),
);
}
})
}
pub fn read_message(&self, msg_number: u32) -> crate::Result<JamMessage> {
let header = self.read_header(msg_number)?;
let text = self.read_message_text(&header)?;
Ok(JamMessage::from_stored(header, text))
}
pub fn read_header(&self, msg_number: u32) -> crate::Result<JamMessageHeader> {
let offset = self.header_offset(msg_number)?;
let header_file_name = self.file_name.with_extension(extensions::HEADER_DATA);
let mut header_file = File::open(header_file_name)?;
header_file.seek(SeekFrom::Start(offset))?;
let mut reader = BufReader::new(header_file);
let header = JamMessageHeader::read(&mut reader)?;
if header.message_number != msg_number {
return Err(
JamError::IndexMessageNumberMismatch(msg_number, header.message_number).into(),
);
}
if header.is_deleted() {
return Err(JamError::MessageDeleted.into());
}
Ok(header)
}
fn index_record(&self, msg_number: u32) -> crate::Result<u64> {
let low = self.lowest_message_number();
let high = self.highest_message_number();
if self.index_records == 0 || msg_number < low || msg_number > high {
return Err(JamError::MessageNumberOutOfRange(msg_number, low, high).into());
}
Ok((msg_number - low) as u64)
}
fn header_offset(&self, msg_number: u32) -> crate::Result<u64> {
let record = self.index_record(msg_number)?;
let index_file_name = self.file_name.with_extension(extensions::MESSAGE_INDEX);
let mut index_file = File::open(index_file_name)?;
let len = index_file.metadata()?.len();
if index_file
.seek(SeekFrom::Start(record * INDEX_RECORD_SIZE as u64 + 4))
.is_err()
{
return Err(JamError::IndexFileCorrupt(record, len).into());
}
let mut offset = [0; 4];
if let Err(err) = index_file.read_exact(&mut offset) {
log::error!("Error reading index file: {err}");
return Err(JamError::IndexFileCorrupt(record, len).into());
}
let offset = u32::from_le_bytes(offset);
if offset == EMPTY_SLOT {
return Err(JamError::MessageDeleted.into());
}
Ok(offset as u64)
}
pub(crate) fn set_attributes(
&mut self,
msg_number: u32,
set: u32,
clear: u32,
) -> crate::Result<bool> {
self.transaction(|base| base.set_attributes_locked(msg_number, set, clear))
}
fn set_attributes_locked(
&mut self,
msg_number: u32,
set: u32,
clear: u32,
) -> crate::Result<bool> {
let offset = self.header_offset(msg_number)?;
let old_header = self.read_header_at(offset)?;
let attributes = (old_header.attributes | set) & !clear;
if attributes == old_header.attributes {
return Ok(false);
}
let original_info = self.header_info.clone();
let old_deleted = old_header.is_deleted();
let mut new_header = old_header.clone();
new_header.attributes = attributes;
let result = (|| {
self.write_header_at(offset, &new_header)?;
self.adjust_active_messages(old_deleted, new_header.is_deleted())?;
self.write_jhr_header()
})();
if let Err(err) = result {
self.header_info = original_info;
if let Err(restore_err) = self.write_header_at(offset, &old_header) {
log::error!(
"Could not restore a JAM header after an attribute update failed: {restore_err}"
);
}
if let Err(restore_err) = self.store_jhr_header() {
log::error!(
"Could not restore the JAM base header after an attribute update failed: {restore_err}"
);
}
return Err(err);
}
Ok(true)
}
pub(crate) fn update_header(
&mut self,
msg_number: u32,
header: &JamMessageHeader,
) -> crate::Result<()> {
self.transaction(|base| base.update_header_locked(msg_number, header))
}
fn update_header_locked(
&mut self,
msg_number: u32,
header: &JamMessageHeader,
) -> crate::Result<()> {
let old_offset = self.header_offset(msg_number)?;
let record = self.index_record(msg_number)?;
let old_header = self.read_header_at(old_offset)?;
let original_info = self.header_info.clone();
let header_path = self.file_name.with_extension(extensions::HEADER_DATA);
let original_header_len = file_len(&header_path)?;
let index_path = self.file_name.with_extension(extensions::MESSAGE_INDEX);
let original_index = read_index_record(&index_path, record)?;
let mut replacement = header.clone();
replacement.message_number = msg_number;
let result = (|| {
let header_file = OpenOptions::new().append(true).open(&header_path)?;
let offset = offset_u32(header_file.metadata()?.len())?;
let mut writer = BufWriter::new(header_file);
replacement.write(&mut writer)?;
writer.flush()?;
write_index_record(&index_path, record, &replacement, offset)?;
self.retire_header(old_offset)?;
self.adjust_active_messages(old_header.is_deleted(), replacement.is_deleted())?;
self.write_jhr_header()
})();
if let Err(err) = result {
rollback(&header_path, original_header_len);
self.header_info = original_info;
if let Err(restore_err) = self.write_header_at(old_offset, &old_header) {
log::error!("Could not restore a superseded JAM header: {restore_err}");
}
if let Err(restore_err) = restore_index_record(&index_path, record, original_index) {
log::error!("Could not restore a JAM index record: {restore_err}");
}
if let Err(restore_err) = self.store_jhr_header() {
log::error!(
"Could not restore the JAM base header after an update failed: {restore_err}"
);
}
return Err(err);
}
Ok(())
}
fn read_header_at(&self, offset: u64) -> crate::Result<JamMessageHeader> {
let header_path = self.file_name.with_extension(extensions::HEADER_DATA);
let mut file = BufReader::new(File::open(header_path)?);
file.seek(SeekFrom::Start(offset))?;
JamMessageHeader::read(&mut file)
}
fn write_header_at(&self, offset: u64, header: &JamMessageHeader) -> crate::Result<()> {
let header_path = self.file_name.with_extension(extensions::HEADER_DATA);
let mut file = OpenOptions::new()
.write(true)
.truncate(false)
.open(header_path)?;
file.seek(SeekFrom::Start(offset))?;
let mut writer = BufWriter::new(file);
header.write(&mut writer)?;
writer.flush()?;
Ok(())
}
fn adjust_active_messages(&mut self, was_deleted: bool, is_deleted: bool) -> crate::Result<()> {
match (was_deleted, is_deleted) {
(false, true) => {
self.header_info.active_msgs = self
.header_info
.active_msgs
.checked_sub(1)
.ok_or_else(|| crate::Error::jam(0, "ActiveMsgs underflow"))?;
}
(true, false) => {
self.header_info.active_msgs = self
.header_info
.active_msgs
.checked_add(1)
.ok_or_else(|| crate::Error::jam(0, "ActiveMsgs overflow"))?;
}
_ => {}
}
Ok(())
}
fn retire_header(&self, offset: u64) -> crate::Result<()> {
let header_path = self.file_name.with_extension(extensions::HEADER_DATA);
let mut file = File::open(&header_path)?;
file.seek(SeekFrom::Start(offset))?;
let mut reader = BufReader::new(file);
let mut old = JamMessageHeader::read(&mut reader)?;
old.set_deleted(true);
old.txt_len = 0;
let mut file = OpenOptions::new()
.write(true)
.truncate(false)
.open(&header_path)?;
file.seek(SeekFrom::Start(offset))?;
let mut writer = BufWriter::new(file);
old.write(&mut writer)?;
writer.flush()?;
Ok(())
}
pub fn delete_message(&mut self, msg_number: u32) -> crate::Result<()> {
self.set_attributes(msg_number, attributes::MSG_DELETED, 0)
.map(|_| ())
}
pub fn restore_message(&mut self, msg_number: u32) -> crate::Result<()> {
self.set_attributes(msg_number, 0, attributes::MSG_DELETED)
.map(|_| ())
}
pub fn read_last_read_file(&self) -> crate::Result<Vec<JamLastReadStorage>> {
let data = self.read_last_read_records()?;
data.chunks_exact(LASTREAD_RECORD_SIZE)
.map(|mut record| JamLastReadStorage::load(&mut record))
.collect()
}
fn read_last_read_records(&self) -> crate::Result<Vec<u8>> {
let path = self.file_name.with_extension(extensions::LASTREAD_INFO);
let data = match fs::read(&path) {
Ok(data) => data,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => Vec::new(),
Err(err) => return Err(err.into()),
};
if !data.len().is_multiple_of(LASTREAD_RECORD_SIZE) {
return Err(JamError::LastReadFileCorrupted(data.len() as u64).into());
}
Ok(data)
}
fn last_read_record(&self, user_crc: u32, user_id: u32) -> crate::Result<Option<u64>> {
let data = self.read_last_read_records()?;
Ok(find_last_read_record(&data, user_crc, user_id).map(|record| record as u64))
}
pub(crate) fn store_last_read_records(&self, data: &[u8]) -> crate::Result<()> {
let path = self.file_name.with_extension(extensions::LASTREAD_INFO);
let mut tmp = path.as_os_str().to_os_string();
tmp.push(".writing");
let tmp = PathBuf::from(tmp);
fs::write(&tmp, data)?;
pack::sync_file(&tmp)?;
fs::rename(&tmp, &path)?;
pack::sync_dir(&self.file_name);
Ok(())
}
pub fn write_last_read(&mut self, storage: &JamLastReadStorage) -> crate::Result<()> {
self.transaction(|base| {
let mut data = base.read_last_read_records()?;
let mut record = Vec::with_capacity(LASTREAD_RECORD_SIZE);
storage.write(&mut record)?;
match find_last_read_record(&data, storage.user_crc, storage.user_id) {
Some(position) => {
let start = position * LASTREAD_RECORD_SIZE;
data[start..start + LASTREAD_RECORD_SIZE].copy_from_slice(&record);
}
None => data.extend_from_slice(&record),
}
base.store_last_read_records(&data)
})
}
pub fn create_last_read(
&mut self,
user_crc: u32,
user_id: u32,
) -> crate::Result<JamLastReadStorage> {
self.transaction(|base| {
if let Some(existing) = base.find_last_read(user_crc, user_id)? {
return Ok(existing);
}
let storage = JamLastReadStorage {
user_crc,
user_id,
..Default::default()
};
base.write_last_read(&storage)?;
Ok(storage)
})
}
pub fn find_last_read(
&self,
user_crc: u32,
user_id: u32,
) -> crate::Result<Option<JamLastReadStorage>> {
let Some(record) = self.last_read_record(user_crc, user_id)? else {
return Ok(None);
};
let data = self.read_last_read_records()?;
let start = record as usize * LASTREAD_RECORD_SIZE;
let mut slice = &data[start..start + LASTREAD_RECORD_SIZE];
Ok(Some(JamLastReadStorage::load(&mut slice)?))
}
pub fn delete_last_read(&mut self, user_crc: u32, user_id: u32) -> crate::Result<bool> {
self.transaction(|base| {
let mut data = base.read_last_read_records()?;
let Some(position) = find_last_read_record(&data, user_crc, user_id) else {
return Ok(false);
};
let start = position * LASTREAD_RECORD_SIZE;
data[start..start + LASTREAD_DELETED.len()].copy_from_slice(&LASTREAD_DELETED);
base.store_last_read_records(&data)?;
Ok(true)
})
}
pub fn search_message_index(&self, crc: u32) -> crate::Result<Vec<u32>> {
let index_file_name = self.file_name.with_extension(extensions::MESSAGE_INDEX);
let index_file = fs::read(index_file_name)?;
if !index_file.len().is_multiple_of(INDEX_RECORD_SIZE) {
return Err(JamError::IndexFileCorrupted.into());
}
let needle = crc.to_le_bytes();
let empty = EMPTY_SLOT.to_le_bytes();
let base = self.header_info.base_msg_num;
let hit = move |(record, data): (usize, &[u8])| -> Option<u32> {
if data[0..4] == needle && data[4..8] != empty {
Some(base + record as u32)
} else {
None
}
};
let res = if index_file.len() / INDEX_RECORD_SIZE < PARALLEL_SEARCH_THRESHOLD {
index_file
.chunks_exact(INDEX_RECORD_SIZE)
.enumerate()
.filter_map(hit)
.collect()
} else {
index_file
.par_chunks_exact(INDEX_RECORD_SIZE)
.enumerate()
.filter_map(hit)
.collect()
};
Ok(res)
}
pub fn search_to(&self, name: &BString) -> crate::Result<Vec<u32>> {
self.search_message_index(Self::crc(name))
}
pub fn find_by_msgid(&self, id: &BString) -> crate::Result<Vec<u32>> {
self.find_by_msgid_crc(Self::crc(id))
}
pub fn find_by_msgid_crc(&self, crc: u32) -> crate::Result<Vec<u32>> {
self.messages()
.filter_map(|header| match header {
Ok(header) if header.msgid_crc == crc => Some(Ok(header.message_number)),
Ok(_) => None,
Err(err) => Some(Err(err)),
})
.collect()
}
pub(crate) fn read_physical_headers(&self) -> crate::Result<Vec<JamMessageHeader>> {
self.physical_headers()?.collect()
}
pub(crate) fn physical_headers(
&self,
) -> crate::Result<impl Iterator<Item = crate::Result<JamMessageHeader>> + use<>> {
let header_file_name = self.file_name.with_extension(extensions::HEADER_DATA);
let mut f = File::open(header_file_name)?;
let size = f.metadata()?.len();
f.seek(SeekFrom::Start(JhrHeaderInfo::JHR_HEADER_SIZE))?;
Ok(JamBaseMessageIter {
reader: BufReader::new(f),
size,
})
}
}
struct JamBaseMessageIter {
reader: BufReader<File>,
size: u64,
}
impl Iterator for JamBaseMessageIter {
type Item = crate::Result<JamMessageHeader>;
fn next(&mut self) -> Option<Self::Item> {
match self.reader.stream_position() {
Ok(pos) if pos >= self.size => None,
Ok(_) => Some(JamMessageHeader::read(&mut self.reader)),
Err(err) => Some(Err(err.into())),
}
}
}
struct IndexedReader {
index: Vec<u8>,
headers: BufReader<File>,
base_msg_num: u32,
}
impl IndexedReader {
fn open(base: &JamMessageBase) -> crate::Result<Self> {
let index = fs::read(base.file_name.with_extension(extensions::MESSAGE_INDEX))?;
if !index.len().is_multiple_of(INDEX_RECORD_SIZE) {
return Err(JamError::IndexFileCorrupted.into());
}
let headers = File::open(base.file_name.with_extension(extensions::HEADER_DATA))?;
Ok(Self {
index,
headers: BufReader::new(headers),
base_msg_num: base.lowest_message_number(),
})
}
fn header(&mut self, msg_number: u32) -> crate::Result<Option<JamMessageHeader>> {
let Some(record) = msg_number.checked_sub(self.base_msg_num) else {
return Ok(None);
};
let start = record as usize * INDEX_RECORD_SIZE + 4;
let Some(bytes) = self.index.get(start..start + 4) else {
return Ok(None);
};
let offset = u32::from_le_bytes(bytes.try_into().unwrap_or_default());
if offset == EMPTY_SLOT {
return Ok(None);
}
seek_buffered(&mut self.headers, offset as u64)?;
let header = JamMessageHeader::read(&mut self.headers)?;
if header.message_number != msg_number {
return Err(
JamError::IndexMessageNumberMismatch(msg_number, header.message_number).into(),
);
}
Ok((!header.is_deleted()).then_some(header))
}
}
struct TextReader {
file: BufReader<File>,
len: u64,
}
impl TextReader {
fn open(base: &JamMessageBase) -> crate::Result<Self> {
let file = File::open(base.file_name.with_extension(extensions::TEXT_DATA))?;
let len = file.metadata()?.len();
Ok(Self {
file: BufReader::new(file),
len,
})
}
fn text(&mut self, header: &JamMessageHeader) -> crate::Result<BString> {
let end = header.offset as u64 + header.txt_len as u64;
if end > self.len {
return Err(JamError::TextOutOfBounds(header.offset as u64, end - self.len).into());
}
seek_buffered(&mut self.file, header.offset as u64)?;
let mut buffer = vec![0; header.txt_len as usize];
self.file.read_exact(&mut buffer)?;
Ok(BString::new(buffer))
}
}
fn seek_buffered(reader: &mut BufReader<File>, offset: u64) -> crate::Result<()> {
let current = reader.stream_position()?;
match (i64::try_from(offset), i64::try_from(current)) {
(Ok(offset), Ok(current)) => reader.seek_relative(offset - current)?,
_ => {
reader.seek(SeekFrom::Start(offset))?;
}
}
Ok(())
}
pub(crate) fn file_len(path: &Path) -> crate::Result<u64> {
match fs::metadata(path) {
Ok(meta) => Ok(meta.len()),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(0),
Err(err) => Err(err.into()),
}
}
fn find_last_read_record(data: &[u8], user_crc: u32, user_id: u32) -> Option<usize> {
let mut needle = [0; 8];
needle[..4].copy_from_slice(&user_crc.to_le_bytes());
needle[4..].copy_from_slice(&user_id.to_le_bytes());
if needle == LASTREAD_DELETED {
return None;
}
data.chunks_exact(LASTREAD_RECORD_SIZE)
.position(|record| record[..8] == needle)
}
fn rollback(path: &Path, len: u64) {
if let Err(err) = OpenOptions::new()
.write(true)
.truncate(false)
.open(path)
.and_then(|f| f.set_len(len))
{
log::error!(
"Could not roll back {} to {len} bytes: {err}",
path.display()
);
}
}
fn read_index_record(path: &Path, record: u64) -> crate::Result<[u8; INDEX_RECORD_SIZE]> {
let mut file = File::open(path)?;
file.seek(SeekFrom::Start(record * INDEX_RECORD_SIZE as u64))?;
let mut data = [0; INDEX_RECORD_SIZE];
file.read_exact(&mut data)?;
Ok(data)
}
fn restore_index_record(
path: &Path,
record: u64,
data: [u8; INDEX_RECORD_SIZE],
) -> crate::Result<()> {
let mut file = OpenOptions::new().write(true).truncate(false).open(path)?;
file.seek(SeekFrom::Start(record * INDEX_RECORD_SIZE as u64))?;
file.write_all(&data)?;
Ok(())
}
fn write_index_record(
path: &Path,
record: u64,
header: &JamMessageHeader,
offset: u32,
) -> crate::Result<()> {
let mut file = OpenOptions::new().write(true).truncate(false).open(path)?;
file.seek(SeekFrom::Start(record * INDEX_RECORD_SIZE as u64))?;
let crc = header.to().map_or(CRC_SEED, JamMessageBase::crc);
file.write_all(&crc.to_le_bytes())?;
file.write_all(&offset.to_le_bytes())?;
Ok(())
}
#[derive(Default)]
pub struct JamMessage {
header: JamMessageHeader,
text: BString,
}
impl JamMessage {
pub fn msgid_crc(&self) -> u32 {
self.header.msgid_crc
}
pub fn from_stored(header: JamMessageHeader, text: BString) -> Self {
Self { header, text }
}
pub fn new(aka: &EchomailAddress) -> Self {
let now = SystemTime::now();
let date_written = if let Ok(unix_time) = now.duration_since(UNIX_EPOCH) {
unix_time.as_secs() as u32
} else {
0
};
let rnd: u32 = fastrand::u32(..);
let id = BString::from(format!("{} {:08x}", aka, rnd));
let msgid_crc = JamMessageBase::crc(&id);
JamMessage {
header: JamMessageHeader {
msgid_crc,
date_written,
sub_fields: vec![MessageSubfield::new(SubfieldType::MsgID, id)],
..Default::default()
},
text: BString::default(),
}
}
pub fn with_reply_to(mut self, reply_to: u32) -> Self {
self.header.reply_to = reply_to;
self
}
pub fn with_msg_id(mut self, id: BString) -> Self {
self.header.msgid_crc = JamMessageBase::crc(&id);
self.header
.sub_fields
.push(MessageSubfield::new(SubfieldType::MsgID, id));
self
}
pub fn with_reply_id(mut self, id: BString) -> Self {
self.header.reply_crc = JamMessageBase::crc(&id);
self.header
.sub_fields
.push(MessageSubfield::new(SubfieldType::ReplyID, id));
self
}
pub fn with_date_time(mut self, time: DateTime<Utc>) -> Self {
self.header.date_written = time.timestamp() as u32;
self.header.sub_fields.push(MessageSubfield::new(
SubfieldType::DateWritten,
BString::from(time.to_rfc3339()),
));
self
}
pub fn with_packout_date(mut self, time: DateTime<Utc>) -> Self {
self.header.sub_fields.push(MessageSubfield::new(
SubfieldType::PackoutDate,
BString::from(time.to_rfc3339()),
));
self
}
pub fn with_text(mut self, text: BString) -> Self {
self.text = text;
self
}
pub fn with_attributes(mut self, attributes: u32) -> Self {
self.header.attributes = attributes;
self
}
pub fn with_password(mut self, password: &BString) -> Self {
self.header.password_crc = JamMessageBase::crc(password);
self
}
pub fn with_from(mut self, name: BString) -> Self {
self.header
.sub_fields
.push(MessageSubfield::new(SubfieldType::SenderName, name));
self
}
pub fn with_to(mut self, name: BString) -> Self {
self.header
.sub_fields
.push(MessageSubfield::new(SubfieldType::RecvName, name));
self
}
pub fn with_subject(mut self, subject: BString) -> Self {
self.header
.sub_fields
.push(MessageSubfield::new(SubfieldType::Subject, subject));
self
}
pub fn with_sub_field(mut self, sub_field: MessageSubfield) -> Self {
self.header.sub_fields.push(sub_field);
self
}
pub fn with_is_deleted(mut self, deleted: bool) -> Self {
if deleted {
self.header.attributes |= attributes::MSG_DELETED;
} else {
self.header.attributes &= !attributes::MSG_DELETED;
}
self
}
pub fn text(&self) -> &BString {
&self.text
}
pub fn header(&self) -> &JamMessageHeader {
&self.header
}
pub(crate) fn create_jam_header(&self) -> JamMessageHeader {
self.header.clone()
}
pub fn reply_to(&self) -> u32 {
self.header.reply_to
}
pub fn reply_first(&self) -> u32 {
self.header.reply_first
}
pub fn reply_next(&self) -> u32 {
self.header.reply_next
}
pub fn from(&self) -> Option<&BString> {
self.header.from()
}
pub fn to(&self) -> Option<&BString> {
self.header.to()
}
pub fn set_reply_crc(&mut self, crc: u32) {
self.header.reply_crc = crc;
}
pub fn set_reply_first(&mut self, reply_first: u32) {
self.header.reply_first = reply_first;
}
pub fn set_reply_next(&mut self, reply_next: u32) {
self.header.reply_next = reply_next;
}
pub(crate) fn is_deleted(&self) -> bool {
self.header.is_deleted()
}
}
pub mod attributes {
pub const MSG_LOCAL: u32 = 0x00000001;
pub const MSG_INTRANSIT: u32 = 0x00000002;
pub const MSG_PRIVATE: u32 = 0x00000004;
pub const MSG_READ: u32 = 0x00000008;
pub const MSG_SENT: u32 = 0x00000010;
pub const MSG_KILLSENT: u32 = 0x00000020;
pub const MSG_ARCHIVESENT: u32 = 0x00000040;
pub const MSG_HOLD: u32 = 0x00000080;
pub const MSG_CRASH: u32 = 0x00000100;
pub const MSG_IMMEDIATE: u32 = 0x00000200;
pub const MSG_DIRECT: u32 = 0x00000400;
pub const MSG_GATE: u32 = 0x00000800;
pub const MSG_FILEREQUEST: u32 = 0x00001000;
pub const MSG_FILEATTACH: u32 = 0x00002000;
pub const MSG_TRUNCFILE: u32 = 0x00004000;
pub const MSG_KILLFILE: u32 = 0x00008000;
pub const MSG_RECEIPTREQ: u32 = 0x00010000;
pub const MSG_CONFIRMREQ: u32 = 0x00020000;
pub const MSG_ORPHAN: u32 = 0x00040000;
pub const MSG_ENCRYPT: u32 = 0x00080000;
pub const MSG_COMPRESS: u32 = 0x00100000;
pub const MSG_ESCAPED: u32 = 0x00200000;
pub const MSG_FPU: u32 = 0x00400000;
pub const MSG_TYPELOCAL: u32 = 0x00800000;
pub const MSG_TYPEECHO: u32 = 0x01000000;
pub const MSG_TYPENET: u32 = 0x02000000;
pub const MSG_NODISP: u32 = 0x20000000;
pub const MSG_LOCKED: u32 = 0x40000000;
pub const MSG_DELETED: u32 = 0x80000000;
}