Skip to main content

cassandra_cpp/cassandra/
batch.rs

1use crate::cassandra::consistency::Consistency;
2use crate::cassandra::custom_payload::CustomPayload;
3use crate::cassandra::error::*;
4use crate::cassandra::future::CassFuture;
5use crate::cassandra::policy::retry::RetryPolicy;
6use crate::cassandra::statement::Statement;
7use crate::cassandra::util::{Protected, ProtectedInner, ProtectedWithSession};
8use crate::cassandra_sys::cass_session_execute_batch;
9use crate::{CassResult, Session};
10
11use crate::cassandra_sys::cass_batch_add_statement;
12use crate::cassandra_sys::cass_batch_free;
13use crate::cassandra_sys::cass_batch_new;
14use crate::cassandra_sys::cass_batch_set_consistency;
15use crate::cassandra_sys::cass_batch_set_custom_payload;
16use crate::cassandra_sys::cass_batch_set_request_timeout;
17use crate::cassandra_sys::cass_batch_set_retry_policy;
18use crate::cassandra_sys::cass_batch_set_serial_consistency;
19use crate::cassandra_sys::cass_batch_set_timestamp;
20use crate::cassandra_sys::cass_batch_set_tracing;
21use crate::cassandra_sys::CassBatch as _Batch;
22use crate::cassandra_sys::CassBatchType_;
23use crate::cassandra_sys::CASS_UINT64_MAX;
24use crate::cassandra_sys::{cass_false, cass_true};
25use std::time::Duration;
26
27#[derive(Debug)]
28struct BatchInner(*mut _Batch);
29
30/// A group of statements that are executed as a single batch.
31/// <b>Note:</b> Batches are not supported by the binary protocol version 1.
32#[derive(Debug)]
33pub struct Batch(BatchInner, Session);
34
35// The underlying C type has no thread-local state, and forbids only concurrent
36// mutation/free: https://datastax.github.io/cpp-driver/topics/#thread-safety
37unsafe impl Send for BatchInner {}
38unsafe impl Sync for BatchInner {}
39
40impl ProtectedInner<*mut _Batch> for BatchInner {
41    #[inline(always)]
42    fn inner(&self) -> *mut _Batch {
43        self.0
44    }
45}
46
47impl Protected<*mut _Batch> for BatchInner {
48    #[inline(always)]
49    fn build(inner: *mut _Batch) -> Self {
50        if inner.is_null() {
51            panic!("Unexpected null pointer")
52        };
53        Self(inner)
54    }
55}
56
57impl ProtectedInner<*mut _Batch> for Batch {
58    #[inline(always)]
59    fn inner(&self) -> *mut _Batch {
60        self.0.inner()
61    }
62}
63
64impl ProtectedWithSession<*mut _Batch> for Batch {
65    #[inline(always)]
66    fn build(inner: *mut _Batch, session: Session) -> Self {
67        Self(BatchInner::build(inner), session)
68    }
69
70    #[inline(always)]
71    fn session(&self) -> &Session {
72        &self.1
73    }
74}
75
76impl Drop for BatchInner {
77    /// Frees a batch instance. Batches can be immediately freed after being
78    /// executed.
79    fn drop(&mut self) {
80        unsafe { cass_batch_free(self.0) }
81    }
82}
83
84impl Batch {
85    /// Creates a new batch statement with batch type.
86    pub(crate) fn new(batch_type: BatchType, session: Session) -> Batch {
87        unsafe { Batch(BatchInner(cass_batch_new(batch_type.inner())), session) }
88    }
89
90    /// Returns the session of which this batch is bound to.
91    pub fn session(&self) -> &Session {
92        ProtectedWithSession::session(self)
93    }
94
95    /// Executes this batch.
96    pub async fn execute(self) -> Result<CassResult> {
97        let (batch, session) = (self.0, self.1);
98        let execute_future = {
99            let execute_batch =
100                unsafe { cass_session_execute_batch(session.inner(), batch.inner()) };
101            CassFuture::build(session, execute_batch)
102        };
103        execute_future.await
104    }
105
106    /// Sets the batch's consistency level
107    pub fn set_consistency(&mut self, consistency: Consistency) -> Result<&mut Self> {
108        unsafe { cass_batch_set_consistency(self.inner(), consistency.inner()).to_result(self) }
109    }
110
111    /// Sets the batch's serial consistency level.
112    ///
113    /// <b>Default:</b> Not set
114    pub fn set_serial_consistency(&mut self, consistency: Consistency) -> Result<&mut Self> {
115        unsafe {
116            cass_batch_set_serial_consistency(self.inner(), consistency.inner()).to_result(self)
117        }
118    }
119
120    /// Sets the batch's timestamp.
121    pub fn set_timestamp(&mut self, timestamp: i64) -> Result<&Self> {
122        unsafe { cass_batch_set_timestamp(self.inner(), timestamp).to_result(self) }
123    }
124
125    /// Sets the batch's request timeout — how long the driver waits for a
126    /// response before failing the batch with `LIB_REQUEST_TIMED_OUT`.
127    /// `Some(Duration::from_millis(0))` means no timeout; `None` disables the
128    /// per-batch override so the cluster-level request timeout applies.
129    ///
130    /// This is the ONLY timeout that governs batch execution: the DataStax
131    /// C++ driver ignores the per-statement request timeouts set on the
132    /// statements added to a batch — `cass_session_execute_batch` reads the
133    /// timeout from the batch itself. Without this, a batch always uses the
134    /// cluster default (12s) regardless of any member-statement timeout. See
135    /// CHANGES.local.md.
136    pub fn set_request_timeout(&mut self, timeout: Option<Duration>) -> Result<&mut Self> {
137        let timeout_millis = match timeout {
138            None => CASS_UINT64_MAX as u64,
139            Some(time) => time.as_millis() as u64,
140        };
141        unsafe { cass_batch_set_request_timeout(self.inner(), timeout_millis).to_result(self) }
142    }
143
144    /// Enables or disables server-side query tracing for this batch
145    /// (`cass_batch_set_tracing`). The batch is a statement type like any
146    /// other, so it carries the same tracing aspect a single statement does.
147    /// See CHANGES.local.md.
148    pub fn set_tracing(&mut self, enabled: bool) -> Result<&mut Self> {
149        let flag = if enabled { cass_true } else { cass_false };
150        unsafe { cass_batch_set_tracing(self.inner(), flag).to_result(self) }
151    }
152
153    /// Sets the batch's retry policy.
154    pub fn set_retry_policy(&mut self, retry_policy: RetryPolicy) -> Result<&mut Self> {
155        unsafe { cass_batch_set_retry_policy(self.inner(), retry_policy.inner()).to_result(self) }
156    }
157
158    /// Sets the batch's custom payload.
159    pub fn set_custom_payload(&mut self, custom_payload: CustomPayload) -> Result<&mut Self> {
160        unsafe {
161            cass_batch_set_custom_payload(self.inner(), custom_payload.inner()).to_result(self)
162        }
163    }
164
165    /// Adds a statement to a batch.
166    pub fn add_statement(&mut self, statement: Statement) -> Result<&Self> {
167        // If their sessions are not the same, we can reject at this level.
168        if self.session() != statement.session() {
169            return Err(ErrorKind::BatchSessionMismatch(
170                self.session().clone(),
171                statement.session().clone(),
172            )
173            .into());
174        }
175        unsafe { cass_batch_add_statement(self.inner(), statement.inner()).to_result(self) }
176    }
177}
178
179/// A type of batch.
180#[derive(Debug, Eq, PartialEq, Copy, Clone, Hash)]
181#[allow(missing_docs)] // Meanings are defined in CQL documentation.
182#[allow(non_camel_case_types)] // Names are traditional.
183pub enum BatchType {
184    LOGGED,
185    UNLOGGED,
186    COUNTER,
187}
188
189enhance_nullary_enum!(BatchType, CassBatchType_, {
190    (LOGGED, CASS_BATCH_TYPE_LOGGED, "LOGGED"),
191    (UNLOGGED, CASS_BATCH_TYPE_UNLOGGED, "UNLOGGED"),
192    (COUNTER, CASS_BATCH_TYPE_COUNTER, "COUNTER"),
193});