use std::sync::Arc;
use ydb_grpc::ydb_proto::topic::TransactionIdentity;
use tracing::instrument;
use crate::client_query::Transaction;
use crate::client_topic::compression::Executor;
use crate::client_topic::topicwriter::message::TopicWriterMessage;
use crate::client_topic::topicwriter::message_write_status::{
MessageSkipReason, MessageWriteStatus,
};
use crate::client_topic::topicwriter::writer::TopicWriter;
use crate::grpc_connection_manager::GrpcConnectionManager;
use crate::{YdbError, YdbResult};
use super::writer_tx_options::TopicWriterTxOptions;
pub struct TopicWriterTx {
inner: TopicWriter,
}
impl TopicWriterTx {
pub(crate) async fn new(
options: TopicWriterTxOptions,
connection_manager: GrpcConnectionManager,
executor: Arc<dyn Executor>,
tx: &mut Transaction,
) -> YdbResult<Self> {
let (session_id, transaction_id) = tx.identity().await?;
let tx_identity = TransactionIdentity {
id: transaction_id,
session: session_id,
};
let options = options.into_non_tx_options();
let inner =
TopicWriter::with_tx_identity(options, connection_manager, executor, tx_identity)
.await?;
Ok(Self { inner })
}
#[instrument(name = "ydb.TopicWriterTx.Write", skip_all, fields(db.system.name = "ydb"), err)]
pub async fn write(&mut self, message: TopicWriterMessage) -> YdbResult<()> {
match self.inner.write_with_ack(message).await? {
MessageWriteStatus::WrittenInTx(_) => Ok(()),
MessageWriteStatus::Skipped(MessageSkipReason::AlreadyWritten) => Ok(()),
other_message_status => Err(YdbError::custom(format!(
"expected WrittenInTx or AlreadyWritten ack from server, got: {other_message_status:?}"
))),
}
}
#[instrument(name = "ydb.TopicWriterTx.Stop", skip_all, fields(db.system.name = "ydb"), err)]
pub async fn stop(self) -> YdbResult<()> {
self.inner.stop().await
}
}