1use std::ffi::{c_char, CStr};
4
5use crate::api::Writer;
6
7use super::connection::LaminarConnection;
8use super::error::{
9 clear_last_error, set_last_error, LAMINAR_ERR_INVALID_UTF8, LAMINAR_ERR_NULL_POINTER,
10 LAMINAR_OK,
11};
12use super::query::LaminarRecordBatch;
13
14#[repr(C)]
16pub struct LaminarWriter {
17 pub(super) inner: Option<Writer>,
18}
19
20impl LaminarWriter {
21 fn new(writer: Writer) -> Self {
22 Self {
23 inner: Some(writer),
24 }
25 }
26}
27
28#[no_mangle]
35pub unsafe extern "C" fn laminar_writer_create(
36 conn: *mut LaminarConnection,
37 source_name: *const c_char,
38 out: *mut *mut LaminarWriter,
39) -> i32 {
40 clear_last_error();
41
42 if conn.is_null() || source_name.is_null() || out.is_null() {
43 return LAMINAR_ERR_NULL_POINTER;
44 }
45
46 let Ok(name_str) = (unsafe { CStr::from_ptr(source_name) }).to_str() else {
47 return LAMINAR_ERR_INVALID_UTF8;
48 };
49
50 let conn_ref = unsafe { &(*conn).inner };
51
52 match conn_ref.writer(name_str) {
53 Ok(writer) => {
54 let handle = Box::new(LaminarWriter::new(writer));
55 unsafe { *out = Box::into_raw(handle) };
56 LAMINAR_OK
57 }
58 Err(e) => {
59 let code = e.code();
60 set_last_error(e);
61 code
62 }
63 }
64}
65
66#[no_mangle]
72pub unsafe extern "C" fn laminar_writer_write(
73 writer: *mut LaminarWriter,
74 batch: *mut LaminarRecordBatch,
75) -> i32 {
76 clear_last_error();
77
78 if writer.is_null() || batch.is_null() {
79 return LAMINAR_ERR_NULL_POINTER;
80 }
81
82 let batch_box = unsafe { Box::from_raw(batch) };
83 let record_batch = batch_box.into_inner();
84
85 let writer_ref = unsafe { &mut (*writer).inner };
86
87 if let Some(w) = writer_ref.as_mut() {
88 match w.write(record_batch) {
89 Ok(()) => LAMINAR_OK,
90 Err(e) => {
91 let code = e.code();
92 set_last_error(e);
93 code
94 }
95 }
96 } else {
97 set_last_error(crate::api::ApiError::internal("Writer already closed"));
98 LAMINAR_ERR_NULL_POINTER
99 }
100}
101
102#[no_mangle]
108pub unsafe extern "C" fn laminar_writer_flush(writer: *mut LaminarWriter) -> i32 {
109 clear_last_error();
110
111 if writer.is_null() {
112 return LAMINAR_ERR_NULL_POINTER;
113 }
114
115 let writer_ref = unsafe { &mut (*writer).inner };
116
117 if let Some(w) = writer_ref.as_mut() {
118 match w.flush() {
119 Ok(()) => LAMINAR_OK,
120 Err(e) => {
121 let code = e.code();
122 set_last_error(e);
123 code
124 }
125 }
126 } else {
127 set_last_error(crate::api::ApiError::internal("Writer already closed"));
128 LAMINAR_ERR_NULL_POINTER
129 }
130}
131
132#[no_mangle]
138pub unsafe extern "C" fn laminar_writer_close(writer: *mut LaminarWriter) -> i32 {
139 clear_last_error();
140
141 if writer.is_null() {
142 return LAMINAR_ERR_NULL_POINTER;
143 }
144
145 let writer_ref = unsafe { &mut (*writer).inner };
146
147 match writer_ref.take() {
148 Some(w) => match w.close() {
149 Ok(()) => LAMINAR_OK,
150 Err(e) => {
151 let code = e.code();
152 set_last_error(e);
153 code
154 }
155 },
156 None => LAMINAR_OK,
157 }
158}
159
160#[no_mangle]
166pub unsafe extern "C" fn laminar_writer_free(writer: *mut LaminarWriter) {
167 if !writer.is_null() {
168 drop(unsafe { Box::from_raw(writer) });
169 }
170}
171
172#[cfg(test)]
173#[allow(clippy::borrow_as_ptr)]
174mod tests {
175 use super::*;
176 use crate::ffi::connection::{laminar_close, laminar_execute, laminar_open};
177 use std::ptr;
178
179 #[test]
180 fn test_writer_create() {
181 let mut conn: *mut LaminarConnection = ptr::null_mut();
182 let mut writer: *mut LaminarWriter = ptr::null_mut();
183
184 unsafe {
186 laminar_open(&mut conn);
187
188 let sql = b"CREATE SOURCE writer_test (id BIGINT)\0";
190 laminar_execute(conn, sql.as_ptr().cast(), ptr::null_mut());
191
192 let name = b"writer_test\0";
194 let rc = laminar_writer_create(conn, name.as_ptr().cast(), &mut writer);
195 assert_eq!(rc, LAMINAR_OK);
196 assert!(!writer.is_null());
197
198 let rc = laminar_writer_close(writer);
200 assert_eq!(rc, LAMINAR_OK);
201
202 laminar_writer_free(writer);
203 laminar_close(conn);
204 }
205 }
206
207 #[test]
208 fn test_writer_null_pointer() {
209 let rc = unsafe { laminar_writer_flush(ptr::null_mut()) };
211 assert_eq!(rc, LAMINAR_ERR_NULL_POINTER);
212 }
213
214 #[test]
215 fn test_writer_source_not_found() {
216 let mut conn: *mut LaminarConnection = ptr::null_mut();
217 let mut writer: *mut LaminarWriter = ptr::null_mut();
218
219 unsafe {
221 laminar_open(&mut conn);
222
223 let name = b"nonexistent_source\0";
224 let rc = laminar_writer_create(conn, name.as_ptr().cast(), &mut writer);
225 assert_ne!(rc, LAMINAR_OK);
226 assert!(writer.is_null() || rc != LAMINAR_OK);
227
228 laminar_close(conn);
229 }
230 }
231}