Skip to main content

yrs_kafka/
error.rs

1use thiserror::Error;
2
3/// Non-recoverable errors during initialisation.
4#[derive(Error, Debug)]
5pub enum InitError {
6    /// The [rdkafka] producer could not be properly initialised.
7    #[error("failed to create kafka producer: {0}")]
8    CreateProducer(rdkafka::error::KafkaError),
9    /// The [rocksdb] database could not be opened, consider deleting it.
10    #[error("failed to open rocksdb instance: {0}")]
11    OpenRocksDb(rocksdb::Error),
12}
13
14/// Errors thrown during `yrs-kafka` runtime.
15#[derive(Error, Debug)]
16pub enum Error {
17    /// A message could not be sent by [rdkafka].
18    #[error("failed to send message to kafka producer: {0}")]
19    SendProducer(rdkafka::error::KafkaError),
20    /// State could not be successfully read by [rocksdb].
21    #[error("failed to read state from rocksdb: {0}")]
22    ReadRocksDb(rocksdb::Error),
23    /// A blocking task failed to be ran in the background.
24    #[error("failed to join on spawned task: {0}")]
25    SpawnBlocking(tokio::task::JoinError),
26}
27
28/// Recoverable internal errors thrown by background tasks.
29#[derive(Error, Debug)]
30pub(super) enum InternalError {
31    #[error("failed to create topic stream reader: {0}")]
32    TopicReader(rdkafka::error::KafkaError),
33    #[error("failed to subscribe to topic: {0}")]
34    TopicSubscribe(rdkafka::error::KafkaError),
35    #[error("failed to read from topic: {0}")]
36    ReadTopic(rdkafka::error::KafkaError),
37    #[error("failed to merge update into rocksdb store: {0}")]
38    MergeUpdate(rocksdb::Error),
39    #[error("failed to join on tokio task: {0}")]
40    Join(tokio::task::JoinError),
41    #[error("failed to read from rocksdb store after merge: {0}")]
42    ReadAfterMerge(rocksdb::Error),
43    #[error("bad state, data being unavailable after rocksdb merge")]
44    MissingUnexpected,
45    #[error("failed to update compacted topic: {0}")]
46    UpdateCompacted(rdkafka::error::KafkaError),
47    #[error("failed to commit offset: {0}")]
48    CommitOffset(rdkafka::error::KafkaError),
49}