Skip to main content

cassandra_cpp/cassandra/
session.rs

1#![allow(non_camel_case_types)]
2#![allow(dead_code)]
3#![allow(missing_copy_implementations)]
4
5use crate::cassandra::custom_payload::CustomPayloadResponse;
6use crate::cassandra::error::*;
7use crate::cassandra::future::CassFuture;
8use crate::cassandra::metrics::SessionMetrics;
9use crate::cassandra::prepared::PreparedStatement;
10use crate::cassandra::result::CassResult;
11use crate::cassandra::schema::schema_meta::SchemaMeta;
12use crate::cassandra::statement::Statement;
13use crate::cassandra::util::{Protected, ProtectedInner};
14use crate::{cassandra::batch::Batch, BatchType};
15
16use crate::cassandra_sys::cass_session_close;
17use crate::cassandra_sys::cass_session_execute;
18use crate::cassandra_sys::cass_session_execute_batch;
19use crate::cassandra_sys::cass_session_free;
20use crate::cassandra_sys::cass_session_get_metrics;
21use crate::cassandra_sys::cass_session_get_schema_meta;
22use crate::cassandra_sys::cass_session_new;
23use crate::cassandra_sys::cass_session_prepare_n;
24use crate::cassandra_sys::CassSession as _Session;
25
26use std::mem;
27use std::os::raw::c_char;
28use std::sync::Arc;
29
30#[derive(Debug, Eq, PartialEq)]
31pub struct SessionInner(*mut _Session);
32
33// The underlying C type has no thread-local state, and explicitly supports access
34// from multiple threads: https://datastax.github.io/cpp-driver/topics/#thread-safety
35unsafe impl Send for SessionInner {}
36unsafe impl Sync for SessionInner {}
37
38impl SessionInner {
39    fn new(inner: *mut _Session) -> Arc<Self> {
40        Arc::new(Self(inner))
41    }
42}
43
44/// A session object is used to execute queries and maintains cluster state through
45/// the control connection. The control connection is used to auto-discover nodes and
46/// monitor cluster changes (topology and schema). Each session also maintains multiple
47/// pools of connections to cluster nodes which are used to query the cluster.
48///
49/// Instances of the session object are thread-safe to execute queries.
50#[derive(Debug, Clone, Eq, PartialEq)]
51pub struct Session(pub Arc<SessionInner>);
52
53impl ProtectedInner<*mut _Session> for SessionInner {
54    fn inner(&self) -> *mut _Session {
55        self.0
56    }
57}
58
59impl ProtectedInner<*mut _Session> for Session {
60    fn inner(&self) -> *mut _Session {
61        self.0.inner()
62    }
63}
64
65impl Protected<*mut _Session> for Session {
66    fn build(inner: *mut _Session) -> Self {
67        if inner.is_null() {
68            panic!("Unexpected null pointer")
69        };
70        Session(SessionInner::new(inner))
71    }
72}
73
74impl Drop for SessionInner {
75    /// Frees a session instance. If the session is still connected it will be synchronously
76    /// closed before being deallocated.
77    fn drop(&mut self) {
78        unsafe { cass_session_free(self.0) }
79    }
80}
81
82impl Default for Session {
83    fn default() -> Session {
84        Session::new()
85    }
86}
87
88impl Session {
89    pub(crate) fn new() -> Session {
90        unsafe { Session(SessionInner::new(cass_session_new())) }
91    }
92
93    /// Closes the session connection and synchronously waits
94    /// for any in-flight requests to complete before the
95    /// returned future resolves.
96    ///
97    /// Returns a `CassFuture<()>` that resolves once the
98    /// underlying C++ driver has finished its close handshake.
99    /// Callers can compose timeouts (`tokio::time::timeout`)
100    /// around the await to bound how long they'll wait for a
101    /// hung session.
102    ///
103    /// `Drop for SessionInner` still calls `cass_session_free`
104    /// to deallocate; calling `close().await` first is the
105    /// graceful path. A session that's dropped without
106    /// awaiting `close()` is closed synchronously by the C++
107    /// driver inside `cass_session_free`, which can block for
108    /// up to the underlying close timeout.
109    pub fn close(&self) -> CassFuture<()> {
110        let inner_future = unsafe { cass_session_close(self.inner()) };
111        <CassFuture<()>>::build(self.clone(), inner_future)
112    }
113
114    /// Create a prepared statement with the given query.
115    pub async fn prepare(&self, query: impl AsRef<str>) -> Result<PreparedStatement> {
116        let query = query.as_ref();
117        let prepare_future = {
118            let query_ptr = query.as_ptr() as *const c_char;
119            CassFuture::build(self.clone(), unsafe {
120                cass_session_prepare_n(self.inner(), query_ptr, query.len())
121            })
122        };
123        prepare_future.await
124    }
125
126    /// Creates a statement with the given query.
127    pub fn statement(&self, query: impl AsRef<str>) -> Statement {
128        let query = query.as_ref();
129        let param_count = query.matches('?').count();
130        Statement::new(self.clone(), query, param_count)
131    }
132
133    /// Execute a batch statement.
134    pub fn execute_batch(&self, batch: &Batch) -> CassFuture<CassResult> {
135        let inner_future = unsafe { cass_session_execute_batch(self.inner(), batch.inner()) };
136        <CassFuture<CassResult>>::build(self.clone(), inner_future)
137    }
138
139    /// Execute a batch statement and get any custom payloads from the response.
140    pub fn execute_batch_with_payloads(
141        &self,
142        batch: &Batch,
143    ) -> CassFuture<(CassResult, CustomPayloadResponse)> {
144        let inner_future = unsafe { cass_session_execute_batch(self.inner(), batch.inner()) };
145        <CassFuture<(CassResult, CustomPayloadResponse)>>::build(self.clone(), inner_future)
146    }
147
148    /// Execute a batch statement and retrieve the server-side
149    /// trace UUID from the future. The UUID is `Some` when at
150    /// least one statement in the batch had `set_tracing(true)`
151    /// set before being added (the cpp-driver returns
152    /// `CASS_ERROR_LIB_NO_TRACING_ID` otherwise, surfaced here
153    /// as `None`).
154    ///
155    /// Mirrors [`Session::execute_with_tracing`] in shape — for
156    /// batched execution where we want the trace UUID alongside
157    /// the result, without awaiting in the constructor.
158    pub fn execute_batch_with_tracing(
159        &self,
160        batch: &Batch,
161    ) -> CassFuture<(CassResult, Option<crate::cassandra::uuid::Uuid>)> {
162        let inner_future = unsafe { cass_session_execute_batch(self.inner(), batch.inner()) };
163        <CassFuture<(CassResult, Option<crate::cassandra::uuid::Uuid>)>>::build(
164            self.clone(),
165            inner_future,
166        )
167    }
168
169    /// Executes a given query.
170    pub async fn execute(&self, query: impl AsRef<str>) -> Result<CassResult> {
171        let statement = self.statement(query);
172        statement.execute().await
173    }
174
175    /// Creates a new batch that is bound to this session.
176    pub fn batch(&self, batch_type: BatchType) -> Batch {
177        Batch::new(batch_type, self.clone())
178    }
179
180    /// Execute a statement and get any custom payloads from the response.
181    pub fn execute_with_payloads(
182        &self,
183        statement: &Statement,
184    ) -> CassFuture<(CassResult, CustomPayloadResponse)> {
185        let inner_future = unsafe { cass_session_execute(self.inner(), statement.inner()) };
186        <CassFuture<(CassResult, CustomPayloadResponse)>>::build(self.clone(), inner_future)
187    }
188
189    /// Execute a statement and retrieve the server-side trace
190    /// UUID from the future. The UUID is `Some` when the
191    /// statement had `set_tracing(true)` set before execute
192    /// (the cpp-driver returns `CASS_ERROR_LIB_NO_TRACING_ID`
193    /// otherwise, surfaced here as `None`). Pair with
194    /// `system_traces.sessions` / `system_traces.events`
195    /// queries by trace UUID for full client-driven trace
196    /// capture.
197    ///
198    /// Mirrors [`Session::execute_with_payloads`] in shape —
199    /// borrows the statement (so it can be re-used or dropped
200    /// by the caller) and returns the future without awaiting,
201    /// matching the pattern callers expect for the
202    /// auxiliary-data execute variants.
203    pub fn execute_with_tracing(
204        &self,
205        statement: &Statement,
206    ) -> CassFuture<(CassResult, Option<crate::cassandra::uuid::Uuid>)> {
207        let inner_future = unsafe { cass_session_execute(self.inner(), statement.inner()) };
208        <CassFuture<(CassResult, Option<crate::cassandra::uuid::Uuid>)>>::build(
209            self.clone(),
210            inner_future,
211        )
212    }
213
214    /// Gets a snapshot of this session's schema metadata. The returned
215    /// snapshot of the schema metadata is not updated. This function
216    /// must be called again to retrieve any schema changes since the
217    /// previous call.
218    pub fn get_schema_meta(&self) -> SchemaMeta {
219        unsafe { SchemaMeta::build(cass_session_get_schema_meta(self.inner())) }
220    }
221
222    /// Gets a copy of this session's performance/diagnostic metrics.
223    pub fn get_metrics(&self) -> SessionMetrics {
224        unsafe {
225            let mut metrics = mem::zeroed();
226            cass_session_get_metrics(self.inner(), &mut metrics);
227            SessionMetrics::build(&metrics)
228        }
229    }
230
231    //    pub fn get_schema(&self) -> Schema {
232    //        unsafe { Schema(cass_session_get_schema(self.0)) }
233    //    }
234}