Skip to main content

clickhouse_c/
ioless.rs

1//! Transport-independent native protocol client.
2//!
3//! [`IolessClient`] processes protocol bytes without reading or writing a
4//! transport. Callers submit received bytes and consume bytes queued for
5//! sending. Operations return [`Step::NeedsInput`] when more data is required.
6
7use core::pin::Pin;
8use core::ptr::NonNull;
9use core::slice;
10use std::ffi::c_char;
11
12use crate::alloc::Allocator;
13use crate::builder::BlockBuilder;
14use crate::client::{ClientOpts, Event, ServerInfo, take_handshake_exception};
15use crate::codec::Codec;
16use crate::error::{Result, check};
17use crate::sys;
18
19/// Result of a protocol operation that may require more input.
20#[derive(Debug, Clone, Copy, PartialEq, Eq)]
21pub enum Step<T> {
22    Ready(T),
23    /// Operation requires more bytes. Submit them with
24    /// [`IolessClient::submit`] and repeat operation.
25    NeedsInput,
26}
27
28impl<T> Step<T> {
29    pub fn ready(self) -> Option<T> {
30        match self {
31            Self::Ready(v) => Some(v),
32            Self::NeedsInput => None,
33        }
34    }
35
36    pub fn is_ready(&self) -> bool {
37        matches!(self, Self::Ready(_))
38    }
39}
40
41/// Native protocol client without transport I/O.
42///
43/// ```no_run
44/// use clickhouse_c::{Allocator, ClientOpts, Event, IolessClient, Step};
45/// use std::io::{Read, Write};
46/// use std::net::TcpStream;
47///
48/// # fn main() -> clickhouse_c::Result<()> {
49/// let mut sock = TcpStream::connect("localhost:9000")?;
50/// let mut core = IolessClient::new(&ClientOpts::new(), Allocator::stdlib(), None)?;
51/// let mut buf = [0u8; 8192];
52///
53/// // Send queued bytes before waiting for more input.
54/// let mut pump = |core: &mut IolessClient, sock: &mut TcpStream| -> clickhouse_c::Result<()> {
55///     while !core.pending_out().is_empty() {
56///         let n = sock.write(core.pending_out())?;
57///         core.consume_out(n);
58///     }
59///     let n = sock.read(&mut buf)?;
60///     core.submit(&buf[..n])
61/// };
62///
63/// while !core.handshake()?.is_ready() {
64///     pump(&mut core, &mut sock)?;
65/// }
66///
67/// core.send_query("SELECT 1", None)?;
68/// loop {
69///     match core.recv_event()? {
70///         Step::Ready(Event::EndOfStream) => break,
71///         Step::Ready(_) => {}
72///         Step::NeedsInput => pump(&mut core, &mut sock)?,
73///     }
74/// }
75/// # Ok(())
76/// # }
77/// ```
78pub struct IolessClient {
79    raw: NonNull<sys::chc_async_client>,
80    // C client retains allocator address until destruction
81    alloc: Box<Allocator>,
82    _codec: Option<Pin<Box<Codec>>>,
83}
84
85impl IolessClient {
86    /// Creates protocol client. Call [`handshake`](Self::handshake) before
87    /// sending queries.
88    pub fn new(
89        opts: &ClientOpts,
90        alloc: Allocator,
91        codec: Option<Pin<Box<Codec>>>,
92    ) -> Result<Self> {
93        opts.validate_codec(codec.as_ref().map(|codec| codec.as_ref()))?;
94        let codec_ptr = codec.as_ref().map(|c| c.as_ref().as_ptr());
95        let raw_opts = opts.to_raw(codec_ptr)?;
96        let alloc = Box::new(alloc);
97        let mut out: *mut sys::chc_async_client = core::ptr::null_mut();
98        let mut err = sys::chc_err::zeroed();
99        let rc = unsafe {
100            sys::chc_async_client_init(&mut out, raw_opts.as_ptr(), alloc.as_ptr(), &mut err)
101        };
102        check(rc, &err)?;
103        Ok(Self {
104            raw: NonNull::new(out).expect("chc_async_client_init returned OK with NULL"),
105            alloc,
106            _codec: codec,
107        })
108    }
109
110    /// Copies received transport bytes into protocol input buffer.
111    pub fn submit(&mut self, bytes: &[u8]) -> Result<()> {
112        let mut err = sys::chc_err::zeroed();
113        let rc = unsafe {
114            sys::chc_async_submit(
115                self.raw.as_ptr(),
116                bytes.as_ptr().cast(),
117                bytes.len(),
118                &mut err,
119            )
120        };
121        check(rc, &err)
122    }
123
124    /// Returns bytes waiting to be written to transport.
125    ///
126    /// Send methods append to this buffer without backpressure. Callers should
127    /// limit queued operations according to their memory requirements.
128    pub fn pending_out(&self) -> &[u8] {
129        let mut ptr: *const u8 = core::ptr::null();
130        let mut len = 0usize;
131        unsafe { sys::chc_async_pending_out(self.raw.as_ptr(), &mut ptr, &mut len) };
132        if ptr.is_null() || len == 0 {
133            return &[];
134        }
135        // SAFETY: C buffer remains unchanged until next mutable operation
136        unsafe { slice::from_raw_parts(ptr, len) }
137    }
138
139    /// Removes first `n` bytes after transport accepts them.
140    ///
141    /// Values larger than queued length remove all queued bytes.
142    pub fn consume_out(&mut self, n: usize) {
143        unsafe { sys::chc_async_consume_out(self.raw.as_ptr(), n) };
144    }
145
146    /// Advances Hello exchange. Repeat after sending output and submitting
147    /// input until method returns [`Step::Ready`].
148    ///
149    /// Server rejection returns
150    /// [`ErrorKind::Server`](crate::ErrorKind::Server) carrying exception
151    /// code, class, and untruncated message.
152    pub fn handshake(&mut self) -> Result<Step<()>> {
153        let mut exc: *mut sys::chc_exception = core::ptr::null_mut();
154        let mut err = sys::chc_err::zeroed();
155        let rc = unsafe { sys::chc_async_handshake(self.raw.as_ptr(), &mut exc, &mut err) };
156        if let Some(e) = take_handshake_exception(exc, *self.alloc) {
157            return Err(e);
158        }
159        step(rc, &err).map(|s| s.map_ready(|()| ()))
160    }
161
162    /// Queues a query.
163    ///
164    /// I/O-independent C API does not support query settings or parameters.
165    pub fn send_query(&mut self, sql: &str, query_id: Option<&str>) -> Result<()> {
166        let (qid, qid_len) = query_id
167            .map(|q| (q.as_ptr().cast::<c_char>(), q.len()))
168            .unwrap_or((core::ptr::null(), 0));
169        let mut err = sys::chc_err::zeroed();
170        let rc = unsafe {
171            sys::chc_async_send_query(
172                self.raw.as_ptr(),
173                sql.as_ptr().cast::<c_char>(),
174                sql.len(),
175                qid,
176                qid_len,
177                &mut err,
178            )
179        };
180        check(rc, &err)
181    }
182
183    /// Queues a Data block, or empty terminator when `builder` is `None`.
184    pub fn send_data(&mut self, builder: Option<&BlockBuilder<'_>>) -> Result<()> {
185        let bb_ptr = builder.map(|b| b.as_ptr()).unwrap_or(core::ptr::null());
186        let mut err = sys::chc_err::zeroed();
187        let rc = unsafe { sys::chc_async_send_data(self.raw.as_ptr(), bb_ptr, &mut err) };
188        check(rc, &err)
189    }
190
191    /// Queues empty Data block that ends INSERT input.
192    pub fn send_data_end(&mut self) -> Result<()> {
193        let mut err = sys::chc_err::zeroed();
194        let rc = unsafe { sys::chc_async_send_data_end(self.raw.as_ptr(), &mut err) };
195        check(rc, &err)
196    }
197
198    /// Decodes next server event from submitted bytes.
199    ///
200    /// Returned event owns block or exception payload.
201    pub fn recv_event(&mut self) -> Result<Step<Event>> {
202        let mut raw = sys::chc_packet::zeroed();
203        let mut err = sys::chc_err::zeroed();
204        let rc = unsafe { sys::chc_async_recv_packet(self.raw.as_ptr(), &mut raw, &mut err) };
205        if rc == sys::CHC_WOULD_BLOCK {
206            return Ok(Step::NeedsInput);
207        }
208        if let Err(e) = check(rc, &err) {
209            unsafe { sys::chc_async_packet_clear(self.raw.as_ptr(), &mut raw) };
210            return Err(e);
211        }
212        let event = Event::from_raw(&mut raw, *self.alloc);
213        unsafe { sys::chc_async_packet_clear(self.raw.as_ptr(), &mut raw) };
214        event.map(Step::Ready)
215    }
216
217    /// Returns server identity information.
218    ///
219    /// Before handshake completes, value contains empty name and requested
220    /// revision. Read value after [`handshake`](Self::handshake) returns
221    /// [`Step::Ready`] to obtain server-provided information.
222    pub fn server_info(&self) -> Option<ServerInfo> {
223        let p = unsafe { sys::chc_async_server_info(self.raw.as_ptr().cast_const()) };
224        (!p.is_null()).then(|| ServerInfo::from_raw(unsafe { &*p }))
225    }
226}
227
228impl<T> Step<T> {
229    fn map_ready<U>(self, f: impl FnOnce(T) -> U) -> Step<U> {
230        match self {
231            Self::Ready(v) => Step::Ready(f(v)),
232            Self::NeedsInput => Step::NeedsInput,
233        }
234    }
235}
236
237fn step(rc: i32, err: &sys::chc_err) -> Result<Step<()>> {
238    if rc == sys::CHC_WOULD_BLOCK {
239        Ok(Step::NeedsInput)
240    } else {
241        check(rc, err).map(Step::Ready)
242    }
243}
244
245impl Drop for IolessClient {
246    fn drop(&mut self) {
247        unsafe { sys::chc_async_client_free(self.raw.as_ptr()) };
248    }
249}
250
251// C client is uniquely owned and used from one thread at a time
252unsafe impl Send for IolessClient {}