use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll};
use std::time::Duration;
use crate::error::Error;
use crate::topic::subscription::RecvDelivery;
use crate::topic::topic::TopicInner;
use crate::topic::types::{
PublishOutcome, SubscriberOptions, TopicOptions, TrackedPublishPermit, TrackedPublishReceipt,
TrackedTryPublishOutcome,
};
use otel_arrow_dfe_config::topic::TopicBroadcastOnLagPolicy;
use otel_arrow_dfe_config::{SubscriptionGroupName, TopicName};
pub type PublishFuture<'a> = Pin<Box<dyn Future<Output = Result<(), Error>> + Send + 'a>>;
pub type PublishTrackedFuture<'a> =
Pin<Box<dyn Future<Output = Result<TrackedPublishReceipt, Error>> + Send + 'a>>;
pub trait TopicBackend<T: Send + Sync + 'static>: Send + Sync + 'static {
fn create_topic(&self, name: TopicName, opts: TopicOptions) -> Arc<dyn TopicState<T>>;
}
pub trait TopicState<T: Send + Sync + 'static>: Send + Sync {
fn name(&self) -> &TopicName;
fn publish(&self, msg: Arc<T>) -> PublishFuture<'_>;
fn publish_tracked(
&self,
msg: Arc<T>,
timeout: Duration,
permit: TrackedPublishPermit,
) -> PublishTrackedFuture<'_>;
fn try_publish(&self, msg: Arc<T>) -> Result<PublishOutcome, Error>;
fn try_publish_tracked(
&self,
msg: Arc<T>,
timeout: Duration,
permit: TrackedPublishPermit,
) -> Result<TrackedTryPublishOutcome, Error>;
fn subscribe_balanced(
&self,
group: SubscriptionGroupName,
opts: SubscriberOptions,
) -> Result<Box<dyn SubscriptionBackend<T>>, Error>;
fn subscribe_broadcast(
&self,
opts: SubscriberOptions,
) -> Result<Box<dyn SubscriptionBackend<T>>, Error>;
fn broadcast_on_lag_policy(&self) -> TopicBroadcastOnLagPolicy;
fn close(&self);
#[cfg(test)]
fn debug_balanced_available_permits(&self) -> Vec<(SubscriptionGroupName, usize)>;
}
pub trait SubscriptionBackend<T: Send + Sync + 'static>: Send {
fn poll_recv_delivery(&mut self, cx: &mut Context<'_>) -> Poll<Result<RecvDelivery<T>, Error>>;
fn ack(&self, id: u64) -> Result<(), Error>;
fn nack(&self, id: u64, reason: Arc<str>) -> Result<(), Error>;
}
pub struct InMemoryBackend;
impl<T: Send + Sync + 'static> TopicBackend<T> for InMemoryBackend {
fn create_topic(&self, name: TopicName, opts: TopicOptions) -> Arc<dyn TopicState<T>> {
Arc::new(TopicInner::new(name, opts))
}
}