mod batch;
mod condition;
mod coordinator;
mod handle;
mod tips;
use std::sync::Arc;
use flume::Sender;
use thiserror::Error;
use crate::Position;
use crate::event::{DecodeError, Event};
use crate::log::set::{LogError, PositionRange};
use crate::query::AppendCondition;
use crate::read::ReadConfig;
pub use coordinator::WriteCoordinator;
pub use handle::WriteHandle;
#[derive(Clone, Copy, Debug)]
pub struct WriterConfig {
pub queue_capacity: usize,
pub max_batch_records: usize,
pub max_batch_bytes: usize,
pub tips_window: u64,
pub verify_tips: bool,
pub condition_force_scan: bool,
pub read: ReadConfig,
}
impl Default for WriterConfig {
fn default() -> Self {
WriterConfig {
queue_capacity: 1024,
max_batch_records: 1024,
max_batch_bytes: 8 * 1024 * 1024,
tips_window: 1_000_000,
verify_tips: false,
condition_force_scan: false,
read: ReadConfig::default(),
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ConflictSite {
Durable(Position),
SameBatch,
}
#[derive(Clone, Debug, Error)]
pub enum AppendError {
#[error("append condition conflict ({at:?})")]
Conflict { at: ConflictSite },
#[error("after {after} is beyond the durable tip {tip}")]
AfterBeyondTip { after: Position, tip: Position },
#[error("append with no events")]
Empty,
#[error("event batch of {size} bytes exceeds the segment capacity")]
TooLarge { size: usize },
#[error("log write failed: {0}")]
Log(Arc<LogError>),
#[error("event on the log failed to decode: {0}")]
Corrupt(DecodeError),
#[error("write coordinator has shut down")]
Shutdown,
}
type Reply = Sender<Result<PositionRange, AppendError>>;
struct Request {
events: Vec<Event>,
condition: Option<AppendCondition>,
reply: Reply,
}
enum Message {
Append(Request),
Shutdown,
}