use std::error::Error;
use std::fmt;
use crate::{
BoxFuture, Checkpoint, ChunkCounts, ExecutionContext, FailureCategory, FaultProgress,
JobExecutionId, SkipCounts, StepExecutionId, StopToken,
};
#[derive(Clone, Copy, Debug)]
pub struct ReadContext<'a> {
stop: &'a StopToken,
}
impl<'a> ReadContext<'a> {
#[must_use]
pub const fn new(stop: &'a StopToken) -> Self {
Self { stop }
}
#[must_use]
pub const fn stop_token(self) -> &'a StopToken {
self.stop
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum ReadOutcome<I> {
Item(I),
EndOfInput,
Stopped,
}
pub trait ItemReader<I>: Send {
fn read<'a>(
&'a mut self,
context: ReadContext<'a>,
) -> BoxFuture<'a, Result<ReadOutcome<I>, ReaderError>>;
}
#[derive(Clone, Copy, Debug)]
pub struct ProcessContext<'a> {
stop: &'a StopToken,
}
impl<'a> ProcessContext<'a> {
#[must_use]
pub const fn new(stop: &'a StopToken) -> Self {
Self { stop }
}
#[must_use]
pub const fn stop_token(self) -> &'a StopToken {
self.stop
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum ProcessOutcome<O> {
Item(O),
Filtered,
Stopped,
}
pub trait ItemProcessor<I, O>: Send + Sync {
fn process<'a>(
&'a self,
item: &'a I,
context: ProcessContext<'a>,
) -> BoxFuture<'a, Result<ProcessOutcome<O>, ProcessorError>>;
}
#[derive(Clone, Copy, Eq, PartialEq)]
pub struct BusinessValue<'a>(BusinessValueInner<'a>);
#[derive(Clone, Copy, Eq, PartialEq)]
enum BusinessValueInner<'a> {
Text(&'a str),
Bytes(&'a [u8]),
I64(i64),
Bool(bool),
Null,
}
impl<'a> BusinessValue<'a> {
#[must_use]
pub const fn text(value: &'a str) -> Self {
Self(BusinessValueInner::Text(value))
}
#[must_use]
pub const fn bytes(value: &'a [u8]) -> Self {
Self(BusinessValueInner::Bytes(value))
}
#[must_use]
pub const fn i64(value: i64) -> Self {
Self(BusinessValueInner::I64(value))
}
#[must_use]
pub const fn boolean(value: bool) -> Self {
Self(BusinessValueInner::Bool(value))
}
#[must_use]
pub const fn null() -> Self {
Self(BusinessValueInner::Null)
}
#[must_use]
pub const fn kind(self) -> BusinessValueKind {
match self.0 {
BusinessValueInner::Text(_) => BusinessValueKind::Text,
BusinessValueInner::Bytes(_) => BusinessValueKind::Bytes,
BusinessValueInner::I64(_) => BusinessValueKind::I64,
BusinessValueInner::Bool(_) => BusinessValueKind::Bool,
BusinessValueInner::Null => BusinessValueKind::Null,
}
}
#[must_use]
pub const fn as_text(self) -> Option<&'a str> {
match self.0 {
BusinessValueInner::Text(value) => Some(value),
_ => None,
}
}
#[must_use]
pub const fn as_bytes(self) -> Option<&'a [u8]> {
match self.0 {
BusinessValueInner::Bytes(value) => Some(value),
_ => None,
}
}
#[must_use]
pub const fn as_i64(self) -> Option<i64> {
match self.0 {
BusinessValueInner::I64(value) => Some(value),
_ => None,
}
}
#[must_use]
pub const fn as_bool(self) -> Option<bool> {
match self.0 {
BusinessValueInner::Bool(value) => Some(value),
_ => None,
}
}
}
impl fmt::Debug for BusinessValue<'_> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("BusinessValue")
.field("kind", &self.kind())
.field("value", &"<redacted>")
.finish()
}
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
#[non_exhaustive]
pub enum BusinessValueKind {
Text,
Bytes,
I64,
Bool,
Null,
}
pub struct BusinessStatement<'a> {
text: &'a str,
values: &'a [BusinessValue<'a>],
}
impl<'a> BusinessStatement<'a> {
#[must_use]
pub const fn new(text: &'a str, values: &'a [BusinessValue<'a>]) -> Self {
Self { text, values }
}
#[must_use]
pub const fn text(&self) -> &'a str {
self.text
}
#[must_use]
pub const fn values(&self) -> &'a [BusinessValue<'a>] {
self.values
}
}
impl fmt::Debug for BusinessStatement<'_> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("BusinessStatement")
.field("text", &"<redacted>")
.field("value_count", &self.values.len())
.finish()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct BusinessWriteResult {
rows_affected: u64,
}
impl BusinessWriteResult {
#[must_use]
pub const fn new(rows_affected: u64) -> Self {
Self { rows_affected }
}
#[must_use]
pub const fn rows_affected(self) -> u64 {
self.rows_affected
}
}
pub trait BusinessTransaction: Send {
fn execute<'a>(
&'a mut self,
statement: BusinessStatement<'a>,
) -> BoxFuture<'a, Result<BusinessWriteResult, BusinessTransactionError>>;
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ChunkCommitReceipt {
checkpoint: Checkpoint,
execution_context: ExecutionContext,
}
impl ChunkCommitReceipt {
#[must_use]
pub const fn new(checkpoint: Checkpoint, execution_context: ExecutionContext) -> Self {
Self {
checkpoint,
execution_context,
}
}
#[must_use]
pub const fn checkpoint(&self) -> &Checkpoint {
&self.checkpoint
}
#[must_use]
pub const fn execution_context(&self) -> &ExecutionContext {
&self.execution_context
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum ChunkTransactionError {
NotCommitted,
CommitOutcomeUnknown,
}
impl fmt::Display for ChunkTransactionError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(match self {
Self::NotCommitted => "chunk transaction did not commit",
Self::CommitOutcomeUnknown => "chunk transaction commit outcome is unknown",
})
}
}
impl Error for ChunkTransactionError {}
#[derive(Clone, Copy, Debug, Default, Eq, Hash, PartialEq)]
pub struct ChunkFaultProgress {
skips: SkipCounts,
no_rollbacks: u64,
}
impl ChunkFaultProgress {
pub const NONE: Self = Self {
skips: SkipCounts::ZERO,
no_rollbacks: 0,
};
#[must_use]
pub const fn new(skips: SkipCounts, no_rollbacks: u64) -> Self {
Self {
skips,
no_rollbacks,
}
}
#[must_use]
pub const fn skips(self) -> SkipCounts {
self.skips
}
#[must_use]
pub const fn no_rollbacks(self) -> u64 {
self.no_rollbacks
}
}
pub trait ChunkTransaction: Send {
fn business_transaction(&mut self) -> Option<&mut dyn BusinessTransaction>;
fn commit(
&mut self,
counts: ChunkCounts,
fault: ChunkFaultProgress,
) -> BoxFuture<'_, Result<ChunkCommitReceipt, ChunkTransactionError>>;
fn rollback(&mut self) -> BoxFuture<'_, Result<(), ChunkTransactionError>>;
}
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
pub struct ChunkTransactionContext {
job_execution_id: JobExecutionId,
step_execution_id: StepExecutionId,
}
impl ChunkTransactionContext {
#[must_use]
pub const fn new(job_execution_id: JobExecutionId, step_execution_id: StepExecutionId) -> Self {
Self {
job_execution_id,
step_execution_id,
}
}
#[must_use]
pub const fn job_execution_id(self) -> JobExecutionId {
self.job_execution_id
}
#[must_use]
pub const fn step_execution_id(self) -> StepExecutionId {
self.step_execution_id
}
}
#[derive(Clone, Copy, Debug, Default, Eq, Hash, PartialEq)]
pub struct InheritedStepProgress {
read_ordinal: u64,
checkpoint_digest: [u8; 32],
fault: FaultProgress,
}
impl InheritedStepProgress {
pub const NONE: Self = Self {
read_ordinal: 0,
checkpoint_digest: [0; 32],
fault: FaultProgress::NONE,
};
#[must_use]
pub const fn new(read_ordinal: u64, checkpoint_digest: [u8; 32], fault: FaultProgress) -> Self {
Self {
read_ordinal,
checkpoint_digest,
fault,
}
}
#[must_use]
pub const fn read_ordinal(self) -> u64 {
self.read_ordinal
}
#[must_use]
pub const fn checkpoint_digest(self) -> [u8; 32] {
self.checkpoint_digest
}
#[must_use]
pub const fn fault(self) -> FaultProgress {
self.fault
}
}
pub trait ChunkTransactionManager: Send + Sync {
fn begin(&self)
-> BoxFuture<'_, Result<Box<dyn ChunkTransaction + '_>, ChunkTransactionError>>;
fn inherited_progress(
&self,
_context: ChunkTransactionContext,
) -> BoxFuture<'_, Result<InheritedStepProgress, ChunkTransactionError>> {
Box::pin(std::future::ready(Ok(InheritedStepProgress::NONE)))
}
fn begin_for(
&self,
_context: ChunkTransactionContext,
) -> BoxFuture<'_, Result<Box<dyn ChunkTransaction + '_>, ChunkTransactionError>> {
self.begin()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum BusinessTransactionError {
Infrastructure,
Rejected,
Cancelled,
}
impl fmt::Display for BusinessTransactionError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(match self {
Self::Infrastructure => "business transaction infrastructure failed",
Self::Rejected => "business transaction statement was rejected",
Self::Cancelled => "business transaction operation was cancelled",
})
}
}
impl Error for BusinessTransactionError {}
pub struct WriteContext<'a> {
stop: &'a StopToken,
transaction: Option<&'a mut dyn BusinessTransaction>,
}
impl<'a> WriteContext<'a> {
#[must_use]
pub const fn non_transactional(stop: &'a StopToken) -> Self {
Self {
stop,
transaction: None,
}
}
#[must_use]
pub fn enlisted(stop: &'a StopToken, transaction: &'a mut dyn BusinessTransaction) -> Self {
Self {
stop,
transaction: Some(transaction),
}
}
#[must_use]
pub const fn stop_token(&self) -> &'a StopToken {
self.stop
}
#[must_use]
pub fn transaction(&mut self) -> Option<&mut (dyn BusinessTransaction + 'a)> {
self.transaction.as_deref_mut()
}
#[must_use]
pub const fn is_enlisted(&self) -> bool {
self.transaction.is_some()
}
}
impl fmt::Debug for WriteContext<'_> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("WriteContext")
.field("stop_requested", &self.stop.is_stop_requested())
.field("enlisted", &self.transaction.is_some())
.finish()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum WriteOutcome {
Written,
Stopped,
}
pub trait ItemWriter<I>: Send + Sync {
fn write<'a>(
&'a self,
items: &'a [I],
context: WriteContext<'a>,
) -> BoxFuture<'a, Result<WriteOutcome, WriterError>>;
}
#[derive(Clone, Copy, Debug)]
pub struct ChunkCompletionContext<'a> {
checkpoint: &'a Checkpoint,
execution_context: &'a ExecutionContext,
counts: ChunkCounts,
stop: &'a StopToken,
}
impl<'a> ChunkCompletionContext<'a> {
#[must_use]
pub const fn new(
checkpoint: &'a Checkpoint,
execution_context: &'a ExecutionContext,
counts: ChunkCounts,
stop: &'a StopToken,
) -> Self {
Self {
checkpoint,
execution_context,
counts,
stop,
}
}
#[must_use]
pub const fn checkpoint(self) -> &'a Checkpoint {
self.checkpoint
}
#[must_use]
pub const fn execution_context(self) -> &'a ExecutionContext {
self.execution_context
}
#[must_use]
pub const fn counts(self) -> ChunkCounts {
self.counts
}
#[must_use]
pub const fn stop_token(self) -> &'a StopToken {
self.stop
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum ChunkCompletionOutcome {
Acknowledged,
StoppedAfterCommit,
}
pub trait ChunkCompletion: Send + Sync {
fn after_commit<'a>(
&'a self,
context: ChunkCompletionContext<'a>,
) -> BoxFuture<'a, Result<ChunkCompletionOutcome, ChunkCompletionError>>;
}
macro_rules! component_error {
(
$name:ident,
$message:literal
$(, $field:ident : $field_type:ty = $field_default:expr, $field_docs:literal)* $(,)?
) => {
#[doc = $message]
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct $name {
category: FailureCategory,
$(
#[doc = $field_docs]
$field: $field_type,
)*
}
impl $name {
#[must_use]
pub const fn new() -> Self {
Self {
category: FailureCategory::UserComponent,
$($field: $field_default,)*
}
}
#[must_use]
pub const fn with_category(category: FailureCategory) -> Self {
Self {
category,
$($field: $field_default,)*
}
}
#[must_use]
pub fn from_error(error: impl Error + Send + Sync + 'static) -> Self {
drop(error);
Self::new()
}
#[must_use]
pub const fn category(self) -> FailureCategory {
self.category
}
}
impl Default for $name {
fn default() -> Self {
Self::new()
}
}
impl fmt::Display for $name {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str($message)
}
}
impl Error for $name {}
};
}
component_error!(
ReaderError,
"item reader failed",
checkpoint_advanced: bool = false,
"Whether the reader proved its checkpoint moved past one failed input.",
);
component_error!(ProcessorError, "item processor failed");
component_error!(
WriterError,
"item writer failed",
rolled_back_output: Option<usize> = None,
"The located, known-rolled-back output index, when the writer supplied one.",
);
component_error!(ChunkCompletionError, "chunk completion callback failed");
impl ReaderError {
#[must_use]
pub const fn with_checkpoint_advanced(mut self, advanced: bool) -> Self {
self.checkpoint_advanced = advanced;
self
}
#[must_use]
pub const fn has_checkpoint_advanced(self) -> bool {
self.checkpoint_advanced
}
}
impl WriterError {
#[must_use]
pub const fn with_rolled_back_output(mut self, index: usize) -> Self {
self.rolled_back_output = Some(index);
self
}
#[must_use]
pub const fn rolled_back_output(self) -> Option<usize> {
self.rolled_back_output
}
}