cassandra_cpp/cassandra/
batch.rs1use 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#[derive(Debug)]
33pub struct Batch(BatchInner, Session);
34
35unsafe 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 fn drop(&mut self) {
80 unsafe { cass_batch_free(self.0) }
81 }
82}
83
84impl Batch {
85 pub(crate) fn new(batch_type: BatchType, session: Session) -> Batch {
87 unsafe { Batch(BatchInner(cass_batch_new(batch_type.inner())), session) }
88 }
89
90 pub fn session(&self) -> &Session {
92 ProtectedWithSession::session(self)
93 }
94
95 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 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 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 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 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 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 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 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 pub fn add_statement(&mut self, statement: Statement) -> Result<&Self> {
167 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#[derive(Debug, Eq, PartialEq, Copy, Clone, Hash)]
181#[allow(missing_docs)] #[allow(non_camel_case_types)] pub 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});