use std::collections::HashMap;
use std::time::Duration;
use crate::errors::{YdbError, YdbResult};
use crate::grpc_wrapper::raw_query_service::stream::ExecuteQueryStream;
use crate::result::ResultSet;
use crate::types::Value;
use super::exec::{
apply_stream_tx_id, finish_pooled_query_stream, resolve_commit_tx,
transaction_finish_committed_via_query, transaction_mark_invalidated_on_query_error,
CallOptions,
};
use super::internal::ExecCoreRef;
#[must_use = "QueryStream must be fully consumed; call close() when using with_commit(true)"]
pub struct QueryStream<'a> {
pub(crate) core: ExecCoreRef<'a>,
pub(crate) stream: ExecuteQueryStream,
pub(crate) commit_tx: bool,
}
impl Drop for QueryStream<'_> {
fn drop(&mut self) {
if let Some(tx_id) = self.stream.take_captured_tx_id() {
if let ExecCoreRef::Transaction(ctx) = &mut self.core {
apply_stream_tx_id(ctx, Some(tx_id));
}
}
let dropped_mid_stream = self.stream.in_progress();
self.stream.cancel();
if let ExecCoreRef::Transaction(ctx) = &mut self.core {
if let Some(lease) = &mut ctx.pooled_lease {
if dropped_mid_stream {
lease.invalidate_session();
}
lease.end_use();
}
}
}
}
impl QueryStream<'_> {
pub async fn next_result_set(&mut self) -> YdbResult<Option<ResultSet>> {
let next = match self.stream.next_result_set().await {
Ok(v) => v,
Err(err) => {
let ydb_err = YdbError::from(err);
if let ExecCoreRef::Transaction(ctx) = &mut self.core {
transaction_mark_invalidated_on_query_error(ctx, &ydb_err);
}
return Err(ydb_err);
}
};
let Some((raw, tx_id)) = next else {
return Ok(None);
};
if let ExecCoreRef::Transaction(ctx) = &mut self.core {
apply_stream_tx_id(ctx, tx_id);
}
Ok(Some(ResultSet::try_from(raw)?))
}
pub fn stats(&self) -> Option<QueryStats> {
self.stream
.stats()
.map(|total_duration| QueryStats { total_duration })
}
pub async fn close(mut self) -> YdbResult<()> {
match self.stream.close().await {
Ok(meta) => {
finish_pooled_query_stream(&mut self.stream);
if let ExecCoreRef::Transaction(ctx) = &mut self.core {
apply_stream_tx_id(ctx, meta.tx_id);
if self.commit_tx {
transaction_finish_committed_via_query(ctx).await;
}
}
Ok(())
}
Err(err) => {
let ydb_err = YdbError::from(err);
if let ExecCoreRef::Transaction(ctx) = &mut self.core {
transaction_mark_invalidated_on_query_error(ctx, &ydb_err);
}
Err(ydb_err)
}
}
}
}
pub(crate) async fn materialize_query(
core: &mut ExecCoreRef<'_>,
text: String,
params: HashMap<String, Value>,
opts: CallOptions,
) -> YdbResult<Vec<ResultSet>> {
let commit_tx = resolve_commit_tx(core, &opts);
let result: YdbResult<Vec<ResultSet>> = async {
let mut stream = core.begin_stream(text, params, opts, true).await?;
let raw_sets = match stream.materialize_all_result_sets().await {
Ok(v) => v,
Err(err) => {
let ydb_err = YdbError::from(err);
if let ExecCoreRef::Transaction(ctx) = core {
transaction_mark_invalidated_on_query_error(ctx, &ydb_err);
}
return Err(ydb_err);
}
};
let mut sets = Vec::with_capacity(raw_sets.len());
for raw in raw_sets {
sets.push(ResultSet::try_from(raw)?);
}
match stream.close().await {
Ok(meta) => {
finish_pooled_query_stream(&mut stream);
if let ExecCoreRef::Transaction(ctx) = core {
apply_stream_tx_id(ctx, meta.tx_id);
if commit_tx {
transaction_finish_committed_via_query(ctx).await;
}
}
}
Err(err) => {
let ydb_err = YdbError::from(err);
if let ExecCoreRef::Transaction(ctx) = core {
transaction_mark_invalidated_on_query_error(ctx, &ydb_err);
}
return Err(ydb_err);
}
}
Ok(sets)
}
.await;
if let ExecCoreRef::Transaction(ctx) = core {
if let Some(lease) = &mut ctx.pooled_lease {
lease.end_use();
}
}
result
}
#[derive(Debug, Default)]
pub struct QueryStats {
pub total_duration: Duration,
}