1use std::fmt::{self, Debug, Formatter};
28use std::sync::atomic::Ordering;
29use std::thread;
30use std::time::{Duration, Instant};
31
32use crate::ErrorCode;
33use crate::ingress::AckLevel;
34use crate::ingress::QwpWsSenderError;
35use crate::ingress::buffer::{Buffer, QwpWsColumnarBuffer, QwpWsEncodeScratch, SymbolGlobalDict};
36use crate::ingress::sender::qwp_ws::{
37 SyncQwpWsHandlerState, publish_qwp_ws_payload_background, qwp_ws_acked_fsn_background,
38 qwp_ws_begin_close_background, qwp_ws_check_error_background,
39 qwp_ws_drain_to_deadline_background, qwp_ws_is_terminal_background, qwp_ws_ok_fsn_background,
40 qwp_ws_poll_sender_error_background, qwp_ws_poll_sender_error_notification_background,
41 qwp_ws_published_fsn_background, qwp_ws_sender_errors_dropped_background,
42};
43use crate::ingress::sender::qwp_ws_sfa_publisher::{SfaForegroundPublisher, SfaPublishOutcome};
44#[cfg(feature = "arrow-ingress")]
45use crate::ingress::{ColumnName, TableName};
46use crate::{Result, error};
47
48#[cfg(feature = "arrow-ingress")]
49use super::arrow_batch::{self, ArrowColumnOverride, ArrowTsSource};
50use super::chunk::Chunk;
51use super::conn::{ColumnConn, PublishError};
52use super::encoder;
53
54#[cfg(feature = "arrow-ingress")]
55use arrow::array::RecordBatch;
56
57fn classify_flush_error(err: crate::Error) -> crate::Error {
58 if err.code() == ErrorCode::SocketError {
59 return crate::Error::new(ErrorCode::FailoverRetry, err.msg().to_owned());
60 }
61 err
62}
63
64enum FrameOutcome {
66 Published,
67 TooLarge(crate::Error),
71 NoSlot(crate::Error),
77}
78
79#[cfg(feature = "arrow-ingress")]
83struct ArrowFrameSpec<'a> {
84 table: TableName<'a>,
85 batch: &'a RecordBatch,
86 ts: ArrowTsSource,
87 overrides: &'a [ArrowColumnOverride<'a>],
88}
89
90fn split_mid(row_count: usize) -> Option<usize> {
95 if row_count <= 8 {
96 return None;
97 }
98 let mid = (row_count / 2) & !7;
99 Some(if mid == 0 { 8 } else { mid })
100}
101
102#[derive(Clone, Copy)]
105struct SfaFrameCaps {
106 hard: usize,
113 soft: usize,
117}
118
119impl SfaFrameCaps {
120 fn for_range(self, row_count: usize) -> usize {
124 if split_mid(row_count).is_some() {
125 self.soft
126 } else {
127 self.hard
128 }
129 }
130}
131
132fn sfa_frame_size_error(encoded_len: usize, frame_cap: usize) -> crate::Error {
136 error::fmt!(
137 BatchTooLarge,
138 "QWP frame ({} bytes) exceeds the store-and-forward per-frame cap ({} bytes, \
139 the smaller of max_buf_size and the sf_max_segment_bytes segment payload capacity)",
140 encoded_len,
141 frame_cap
142 )
143}
144
145#[derive(Clone, Copy, Debug)]
149pub(crate) enum WaitForAck {
150 No,
151 Yes(AckLevel),
152}
153
154#[derive(Debug)]
162#[doc(hidden)]
163#[non_exhaustive]
164pub enum FlushFailure {
165 NotDelivered(crate::Error),
169 DeliveryUnknown(crate::Error),
182}
183
184impl FlushFailure {
185 #[doc(hidden)]
191 pub fn into_error(self) -> crate::Error {
192 match self {
193 FlushFailure::NotDelivered(e) => e,
194 FlushFailure::DeliveryUnknown(e) => e.with_in_doubt(true),
195 }
196 }
197
198 #[doc(hidden)]
200 #[must_use]
201 pub fn is_not_delivered(&self) -> bool {
202 matches!(self, FlushFailure::NotDelivered(_))
203 }
204}
205
206fn direct_not_delivered(e: crate::Error) -> FlushFailure {
209 FlushFailure::NotDelivered(classify_flush_error(e))
210}
211
212fn direct_delivery_unknown(e: crate::Error) -> FlushFailure {
215 FlushFailure::DeliveryUnknown(classify_flush_error(e))
216}
217
218fn deny_retry_after_partial(f: FlushFailure) -> FlushFailure {
225 match f {
226 FlushFailure::NotDelivered(e) => FlushFailure::DeliveryUnknown(e),
227 other => other,
228 }
229}
230
231pub struct PooledSenderCore {
232 backend: Box<SfaBackend>,
233}
234
235#[doc(hidden)]
239pub struct DirectSenderCore {
240 backend: Box<DirectColumnBackend>,
241}
242
243struct DirectColumnBackend {
244 conn: ColumnConn,
245 symbol_dict: SymbolGlobalDict,
246 scratch: encoder::EncodeScratch,
247 first_frame_sent: bool,
248 commit_since_sync: bool,
254}
255
256struct SfaBackend {
257 foreground: SfaForegroundPublisher,
261 state: SyncQwpWsHandlerState,
262 buffer_scratch: QwpWsEncodeScratch,
263 scratch: encoder::EncodeScratch,
264 max_buf_size: usize,
265 request_durable_ack: bool,
266 sync_timeout: Duration,
273 last_ok_sync_boundary: Option<u64>,
274 last_durable_sync_boundary: Option<u64>,
275 drop_on_return: bool,
276}
277
278impl Debug for PooledSenderCore {
279 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
280 f.debug_struct("PooledSenderCore")
281 .field(
282 "must_close",
283 &(self.backend.drop_on_return
284 || qwp_ws_is_terminal_background(&self.backend.state)),
285 )
286 .finish()
287 }
288}
289
290impl Debug for DirectSenderCore {
291 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
292 f.debug_struct("DirectSenderCore")
293 .field("must_close", &self.backend.conn.must_close())
294 .field("in_flight", &self.backend.conn.in_flight())
295 .finish()
296 }
297}
298
299impl PooledSenderCore {
300 pub(crate) fn new_store_and_forward(
301 mut state: SyncQwpWsHandlerState,
302 max_buf_size: usize,
303 request_durable_ack: bool,
304 sync_timeout: Duration,
305 ) -> Result<Self> {
306 let delta_dict_enabled = state.delta_dict_enabled;
315 let persisted_symbol_dict = state.persisted_symbol_dict.take();
316 let mut foreground = SfaForegroundPublisher::new(delta_dict_enabled, persisted_symbol_dict);
317 if delta_dict_enabled {
318 let recovered = std::mem::take(&mut state.recovered_dict_entries);
331 foreground.seed(&recovered, state.recovered_dict_count)?;
332 }
333 state.release_dormant_encoder_dict();
338 Ok(Self {
339 backend: Box::new(SfaBackend {
340 foreground,
341 state,
342 buffer_scratch: QwpWsEncodeScratch::new(),
343 scratch: encoder::EncodeScratch::new(),
344 max_buf_size,
345 request_durable_ack,
346 sync_timeout,
347 last_ok_sync_boundary: None,
348 last_durable_sync_boundary: None,
349 drop_on_return: false,
350 }),
351 })
352 }
353
354 pub(crate) fn rebase_lease_observation(&mut self) {
361 let sfa = &mut self.backend;
362 if let Ok(published) = qwp_ws_published_fsn_background(&sfa.state) {
363 sfa.last_ok_sync_boundary = published;
364 sfa.last_durable_sync_boundary = published;
365 }
366 while let Ok(Some(_)) = qwp_ws_poll_sender_error_background(&sfa.state) {}
367 while let Ok(Some(_)) = qwp_ws_poll_sender_error_notification_background(&sfa.state) {}
368 }
369
370 pub fn poll_error(&self) -> Result<Option<QwpWsSenderError>> {
374 qwp_ws_poll_sender_error_background(&self.backend.state)
375 }
376
377 pub fn error_events_dropped(&self) -> Result<u64> {
379 qwp_ws_sender_errors_dropped_background(&self.backend.state)
380 }
381
382 #[must_use]
383 pub fn must_close(&self) -> bool {
384 self.backend.drop_on_return || qwp_ws_is_terminal_background(&self.backend.state)
385 }
386
387 pub fn mark_must_close(&mut self) {
388 self.backend.drop_on_return = true;
389 }
390
391 pub fn effective_frame_cap(&self) -> (usize, bool) {
396 self.backend.effective_hard_frame_cap()
397 }
398
399 pub(crate) fn sfa_fully_delivered(&self, durable: bool) -> bool {
408 let sfa = &self.backend;
409 if qwp_ws_is_terminal_background(&sfa.state) {
410 return true;
411 }
412 let Ok(Some(published)) = qwp_ws_published_fsn_background(&sfa.state) else {
413 return true;
416 };
417 let watermark = if durable {
418 qwp_ws_acked_fsn_background(&sfa.state)
419 } else {
420 qwp_ws_ok_fsn_background(&sfa.state)
421 };
422 matches!(watermark, Ok(Some(w)) if w >= published)
423 }
424
425 pub(crate) fn begin_close(&self) {
429 qwp_ws_begin_close_background(&self.backend.state);
430 }
431
432 pub(crate) fn drain_to_deadline(&mut self, deadline: Option<Instant>) -> crate::Result<()> {
437 qwp_ws_drain_to_deadline_background(&mut self.backend.state, deadline)
438 }
439
440 pub fn flush(&mut self, chunk: &mut Chunk<'_>) -> Result<()> {
443 self.backend
444 .flush_chunk(chunk, WaitForAck::No)
445 .map_err(FlushFailure::into_error)
446 }
447
448 pub fn flush_buffer(&mut self, buffer: &mut Buffer) -> Result<()> {
452 self.flush_buffer_and_get_fsn(buffer).map(|_| ())
453 }
454
455 pub fn flush_buffer_and_keep(&mut self, buffer: &Buffer) -> Result<()> {
457 self.flush_buffer_and_keep_and_get_fsn(buffer).map(|_| ())
458 }
459
460 pub fn flush_buffer_and_get_fsn(&mut self, buffer: &mut Buffer) -> Result<Option<u64>> {
464 let fsn = self.publish_buffer(buffer, None)?;
465 buffer.clear();
466 Ok(fsn)
467 }
468
469 pub fn flush_buffer_and_keep_and_get_fsn(&mut self, buffer: &Buffer) -> Result<Option<u64>> {
472 self.publish_buffer(buffer, None)
473 }
474
475 pub fn flush_buffer_and_wait(
480 &mut self,
481 buffer: &mut Buffer,
482 ack_level: AckLevel,
483 ) -> Result<()> {
484 let boundary = self.publish_buffer(buffer, Some(ack_level))?;
485 buffer.clear();
486
487 let sfa = &mut self.backend;
488 match boundary {
489 Some(fsn) => sfa
490 .wait_for_boundary(ack_level, fsn, sfa.sync_timeout)
491 .map_err(FlushFailure::DeliveryUnknown)
492 .map_err(FlushFailure::into_error),
493 None => sfa.wait(ack_level, sfa.sync_timeout),
494 }
495 }
496
497 fn publish_buffer(
498 &mut self,
499 buffer: &Buffer,
500 ack_level: Option<AckLevel>,
501 ) -> Result<Option<u64>> {
502 let qwp = buffer.as_qwp_ws().ok_or_else(|| {
503 error::fmt!(
504 InvalidApiCall,
505 "Pooled QWP ingestion requires a QWP/WebSocket buffer created by `QuestDb::new_buffer()`."
506 )
507 })?;
508 self.backend
509 .publish_buffer(qwp, ack_level)
510 .map_err(FlushFailure::into_error)
511 }
512
513 pub fn flush_and_get_fsn(&mut self, chunk: &mut Chunk<'_>) -> Result<Option<u64>> {
518 self.backend
519 .flush_chunk_and_get_fsn(chunk)
520 .map(Some)
521 .map_err(FlushFailure::into_error)
522 }
523
524 pub fn flush_and_wait(&mut self, chunk: &mut Chunk<'_>, ack_level: AckLevel) -> Result<()> {
545 self.backend
546 .flush_chunk(chunk, WaitForAck::Yes(ack_level))
547 .map_err(FlushFailure::into_error)
548 }
549
550 #[cfg(feature = "arrow-ingress")]
560 pub fn flush_arrow_batch_at_now<'t, T>(
561 &mut self,
562 table: T,
563 batch: &RecordBatch,
564 overrides: &[ArrowColumnOverride<'_>],
565 ) -> Result<()>
566 where
567 T: TryInto<TableName<'t>>,
568 crate::Error: From<T::Error>,
569 {
570 let table: TableName<'t> = table.try_into()?;
571 self.flush_arrow_batch_dispatch(
572 table,
573 batch,
574 ArrowTsSource::ServerNow,
575 overrides,
576 WaitForAck::No,
577 )
578 .map_err(FlushFailure::into_error)
579 }
580
581 #[cfg(feature = "arrow-ingress")]
584 pub fn flush_arrow_batch_at_now_and_get_fsn<'t, T>(
585 &mut self,
586 table: T,
587 batch: &RecordBatch,
588 overrides: &[ArrowColumnOverride<'_>],
589 ) -> Result<Option<u64>>
590 where
591 T: TryInto<TableName<'t>>,
592 crate::Error: From<T::Error>,
593 {
594 let table: TableName<'t> = table.try_into()?;
595 self.flush_arrow_batch_dispatch_get_fsn(table, batch, ArrowTsSource::ServerNow, overrides)
596 .map_err(FlushFailure::into_error)
597 }
598
599 #[cfg(feature = "arrow-ingress")]
603 pub fn flush_arrow_batch_at_now_and_wait<'t, T>(
604 &mut self,
605 table: T,
606 batch: &RecordBatch,
607 overrides: &[ArrowColumnOverride<'_>],
608 ack_level: AckLevel,
609 ) -> Result<()>
610 where
611 T: TryInto<TableName<'t>>,
612 crate::Error: From<T::Error>,
613 {
614 let table: TableName<'t> = table.try_into()?;
615 self.flush_arrow_batch_dispatch(
616 table,
617 batch,
618 ArrowTsSource::ServerNow,
619 overrides,
620 WaitForAck::Yes(ack_level),
621 )
622 .map_err(FlushFailure::into_error)
623 }
624
625 #[cfg(feature = "arrow-ingress")]
631 pub fn flush_arrow_batch_at_column<'t, T>(
632 &mut self,
633 table: T,
634 batch: &RecordBatch,
635 ts_column: ColumnName<'_>,
636 overrides: &[ArrowColumnOverride<'_>],
637 ) -> Result<()>
638 where
639 T: TryInto<TableName<'t>>,
640 crate::Error: From<T::Error>,
641 {
642 let table: TableName<'t> = table.try_into()?;
643 let ts_col_idx = arrow_batch::resolve_ts_column(batch, ts_column)?;
644 self.flush_arrow_batch_dispatch(
645 table,
646 batch,
647 ArrowTsSource::Column(ts_col_idx),
648 overrides,
649 WaitForAck::No,
650 )
651 .map_err(FlushFailure::into_error)
652 }
653
654 #[cfg(feature = "arrow-ingress")]
661 pub fn flush_arrow_batch_at_scalar_nanos<'t, T>(
662 &mut self,
663 table: T,
664 batch: &RecordBatch,
665 nanos: i64,
666 overrides: &[ArrowColumnOverride<'_>],
667 ) -> Result<()>
668 where
669 T: TryInto<TableName<'t>>,
670 crate::Error: From<T::Error>,
671 {
672 let table: TableName<'t> = table.try_into()?;
673 self.flush_arrow_batch_dispatch(
674 table,
675 batch,
676 ArrowTsSource::ScalarNanos(nanos),
677 overrides,
678 WaitForAck::No,
679 )
680 .map_err(FlushFailure::into_error)
681 }
682
683 #[cfg(feature = "arrow-ingress")]
686 pub fn flush_arrow_batch_at_column_and_get_fsn<'t, T>(
687 &mut self,
688 table: T,
689 batch: &RecordBatch,
690 ts_column: ColumnName<'_>,
691 overrides: &[ArrowColumnOverride<'_>],
692 ) -> Result<Option<u64>>
693 where
694 T: TryInto<TableName<'t>>,
695 crate::Error: From<T::Error>,
696 {
697 let table: TableName<'t> = table.try_into()?;
698 let ts_col_idx = arrow_batch::resolve_ts_column(batch, ts_column)?;
699 self.flush_arrow_batch_dispatch_get_fsn(
700 table,
701 batch,
702 ArrowTsSource::Column(ts_col_idx),
703 overrides,
704 )
705 .map_err(FlushFailure::into_error)
706 }
707
708 #[cfg(feature = "arrow-ingress")]
712 pub fn flush_arrow_batch_at_column_and_wait<'t, T>(
713 &mut self,
714 table: T,
715 batch: &RecordBatch,
716 ts_column: ColumnName<'_>,
717 overrides: &[ArrowColumnOverride<'_>],
718 ack_level: AckLevel,
719 ) -> Result<()>
720 where
721 T: TryInto<TableName<'t>>,
722 crate::Error: From<T::Error>,
723 {
724 let table: TableName<'t> = table.try_into()?;
725 let ts_col_idx = arrow_batch::resolve_ts_column(batch, ts_column)?;
726 self.flush_arrow_batch_dispatch(
727 table,
728 batch,
729 ArrowTsSource::Column(ts_col_idx),
730 overrides,
731 WaitForAck::Yes(ack_level),
732 )
733 .map_err(FlushFailure::into_error)
734 }
735
736 #[cfg(feature = "arrow-ingress")]
738 fn flush_arrow_batch_dispatch(
739 &mut self,
740 table: TableName<'_>,
741 batch: &RecordBatch,
742 ts: ArrowTsSource,
743 overrides: &[ArrowColumnOverride<'_>],
744 wait: WaitForAck,
745 ) -> std::result::Result<(), FlushFailure> {
746 self.backend
747 .flush_arrow_batch(table, batch, ts, overrides, wait)
748 }
749
750 #[cfg(feature = "arrow-ingress")]
751 fn flush_arrow_batch_dispatch_get_fsn(
752 &mut self,
753 table: TableName<'_>,
754 batch: &RecordBatch,
755 ts: ArrowTsSource,
756 overrides: &[ArrowColumnOverride<'_>],
757 ) -> std::result::Result<Option<u64>, FlushFailure> {
758 self.backend
759 .flush_arrow_batch_and_get_fsn(table, batch, ts, overrides)
760 .map(Some)
761 }
762
763 #[doc(hidden)]
768 pub fn validate_ack_level(&self, ack_level: AckLevel) -> Result<()> {
769 self.backend.validate_ack_level(ack_level)
770 }
771
772 #[doc(hidden)]
776 #[cfg(feature = "arrow-ingress")]
777 pub fn flush_arrow_batch_at_now_and_wait_ffi(
778 &mut self,
779 table: TableName<'_>,
780 batch: &RecordBatch,
781 overrides: &[ArrowColumnOverride<'_>],
782 ack_level: AckLevel,
783 ) -> std::result::Result<(), FlushFailure> {
784 self.flush_arrow_batch_dispatch(
785 table,
786 batch,
787 ArrowTsSource::ServerNow,
788 overrides,
789 WaitForAck::Yes(ack_level),
790 )
791 }
792
793 #[doc(hidden)]
797 #[cfg(feature = "arrow-ingress")]
798 pub fn flush_arrow_batch_at_column_and_wait_ffi(
799 &mut self,
800 table: TableName<'_>,
801 batch: &RecordBatch,
802 ts_column: ColumnName<'_>,
803 overrides: &[ArrowColumnOverride<'_>],
804 ack_level: AckLevel,
805 ) -> std::result::Result<(), FlushFailure> {
806 let ts_col_idx =
807 arrow_batch::resolve_ts_column(batch, ts_column).map_err(FlushFailure::NotDelivered)?;
808 self.flush_arrow_batch_dispatch(
809 table,
810 batch,
811 ArrowTsSource::Column(ts_col_idx),
812 overrides,
813 WaitForAck::Yes(ack_level),
814 )
815 }
816
817 pub fn sync(&mut self, ack_level: AckLevel) -> Result<()> {
818 self.backend.sync(ack_level)
819 }
820
821 pub fn wait(&mut self, ack_level: AckLevel, timeout: Duration) -> Result<()> {
828 self.backend.wait(ack_level, timeout)
829 }
830
831 pub fn published_fsn(&self) -> Result<Option<u64>> {
835 self.backend.published_fsn()
836 }
837
838 pub fn acked_fsn(&self) -> Result<Option<u64>> {
842 self.backend.acked_fsn()
843 }
844}
845
846impl DirectSenderCore {
847 pub(crate) fn new(
848 conn: ColumnConn,
849 symbol_dict: SymbolGlobalDict,
850 scratch: encoder::EncodeScratch,
851 first_frame_sent: bool,
852 ) -> Self {
853 Self {
854 backend: Box::new(DirectColumnBackend {
855 conn,
856 symbol_dict,
857 scratch,
858 first_frame_sent,
859 commit_since_sync: false,
860 }),
861 }
862 }
863
864 #[must_use]
865 pub fn must_close(&self) -> bool {
866 self.backend.conn.must_close()
867 }
868
869 pub fn mark_must_close(&mut self) {
870 self.backend.conn.mark_must_close();
871 }
872
873 pub(crate) fn in_flight(&self) -> u32 {
874 self.backend.conn.in_flight()
875 }
876
877 pub(crate) fn transport_dead(&self) -> bool {
878 self.backend.conn.transport_dead()
879 }
880
881 pub(crate) fn can_drain_in_flight(&self) -> bool {
886 self.backend.conn.can_drain_in_flight()
887 }
888
889 pub(crate) fn endpoint_idx(&self) -> usize {
890 self.backend.conn.endpoint_idx()
891 }
892
893 pub fn flush(&mut self, chunk: &mut Chunk<'_>) -> Result<()> {
894 self.backend
895 .flush_inner(chunk, WaitForAck::No)
896 .map_err(FlushFailure::into_error)
897 }
898
899 pub fn flush_and_wait(&mut self, chunk: &mut Chunk<'_>, ack_level: AckLevel) -> Result<()> {
900 self.backend
901 .flush_inner(chunk, WaitForAck::Yes(ack_level))
902 .map_err(FlushFailure::into_error)
903 }
904
905 #[cfg(feature = "arrow-ingress")]
906 pub fn flush_arrow_batch_at_now<'t, T>(
907 &mut self,
908 table: T,
909 batch: &RecordBatch,
910 overrides: &[ArrowColumnOverride<'_>],
911 ) -> Result<()>
912 where
913 T: TryInto<TableName<'t>>,
914 crate::Error: From<T::Error>,
915 {
916 let table: TableName<'t> = table.try_into()?;
917 self.backend
918 .flush_arrow_batch_inner(
919 table,
920 batch,
921 ArrowTsSource::ServerNow,
922 overrides,
923 WaitForAck::No,
924 )
925 .map_err(FlushFailure::into_error)
926 }
927
928 #[cfg(feature = "arrow-ingress")]
929 pub fn flush_arrow_batch_at_column<'t, T>(
930 &mut self,
931 table: T,
932 batch: &RecordBatch,
933 ts_column: ColumnName<'_>,
934 overrides: &[ArrowColumnOverride<'_>],
935 ) -> Result<()>
936 where
937 T: TryInto<TableName<'t>>,
938 crate::Error: From<T::Error>,
939 {
940 let table: TableName<'t> = table.try_into()?;
941 let ts_col_idx = arrow_batch::resolve_ts_column(batch, ts_column)?;
942 self.backend
943 .flush_arrow_batch_inner(
944 table,
945 batch,
946 ArrowTsSource::Column(ts_col_idx),
947 overrides,
948 WaitForAck::No,
949 )
950 .map_err(FlushFailure::into_error)
951 }
952
953 #[cfg(feature = "arrow-ingress")]
954 pub fn flush_arrow_batch_at_scalar_nanos<'t, T>(
955 &mut self,
956 table: T,
957 batch: &RecordBatch,
958 nanos: i64,
959 overrides: &[ArrowColumnOverride<'_>],
960 ) -> Result<()>
961 where
962 T: TryInto<TableName<'t>>,
963 crate::Error: From<T::Error>,
964 {
965 let table: TableName<'t> = table.try_into()?;
966 self.backend
967 .flush_arrow_batch_inner(
968 table,
969 batch,
970 ArrowTsSource::ScalarNanos(nanos),
971 overrides,
972 WaitForAck::No,
973 )
974 .map_err(FlushFailure::into_error)
975 }
976
977 #[cfg(feature = "arrow-ingress")]
978 pub fn flush_arrow_batch_at_now_and_wait<'t, T>(
979 &mut self,
980 table: T,
981 batch: &RecordBatch,
982 overrides: &[ArrowColumnOverride<'_>],
983 ack_level: AckLevel,
984 ) -> Result<()>
985 where
986 T: TryInto<TableName<'t>>,
987 crate::Error: From<T::Error>,
988 {
989 let table: TableName<'t> = table.try_into()?;
990 self.backend
991 .flush_arrow_batch_inner(
992 table,
993 batch,
994 ArrowTsSource::ServerNow,
995 overrides,
996 WaitForAck::Yes(ack_level),
997 )
998 .map_err(FlushFailure::into_error)
999 }
1000
1001 #[cfg(feature = "arrow-ingress")]
1002 pub fn flush_arrow_batch_at_column_and_wait<'t, T>(
1003 &mut self,
1004 table: T,
1005 batch: &RecordBatch,
1006 ts_column: ColumnName<'_>,
1007 overrides: &[ArrowColumnOverride<'_>],
1008 ack_level: AckLevel,
1009 ) -> Result<()>
1010 where
1011 T: TryInto<TableName<'t>>,
1012 crate::Error: From<T::Error>,
1013 {
1014 let table: TableName<'t> = table.try_into()?;
1015 let ts_col_idx = arrow_batch::resolve_ts_column(batch, ts_column)?;
1016 self.backend
1017 .flush_arrow_batch_inner(
1018 table,
1019 batch,
1020 ArrowTsSource::Column(ts_col_idx),
1021 overrides,
1022 WaitForAck::Yes(ack_level),
1023 )
1024 .map_err(FlushFailure::into_error)
1025 }
1026
1027 #[doc(hidden)]
1028 pub fn validate_ack_level(&self, ack_level: AckLevel) -> Result<()> {
1029 self.backend.conn.validate_ack_level(ack_level)
1030 }
1031
1032 #[doc(hidden)]
1033 #[cfg(feature = "arrow-ingress")]
1034 pub fn flush_arrow_batch_at_now_and_wait_ffi(
1035 &mut self,
1036 table: TableName<'_>,
1037 batch: &RecordBatch,
1038 overrides: &[ArrowColumnOverride<'_>],
1039 ack_level: AckLevel,
1040 ) -> std::result::Result<(), FlushFailure> {
1041 self.backend.flush_arrow_batch_inner(
1042 table,
1043 batch,
1044 ArrowTsSource::ServerNow,
1045 overrides,
1046 WaitForAck::Yes(ack_level),
1047 )
1048 }
1049
1050 #[doc(hidden)]
1051 #[cfg(feature = "arrow-ingress")]
1052 pub fn flush_arrow_batch_at_column_and_wait_ffi(
1053 &mut self,
1054 table: TableName<'_>,
1055 batch: &RecordBatch,
1056 ts_column: ColumnName<'_>,
1057 overrides: &[ArrowColumnOverride<'_>],
1058 ack_level: AckLevel,
1059 ) -> std::result::Result<(), FlushFailure> {
1060 let ts_col_idx =
1061 arrow_batch::resolve_ts_column(batch, ts_column).map_err(FlushFailure::NotDelivered)?;
1062 self.backend.flush_arrow_batch_inner(
1063 table,
1064 batch,
1065 ArrowTsSource::Column(ts_col_idx),
1066 overrides,
1067 WaitForAck::Yes(ack_level),
1068 )
1069 }
1070
1071 pub fn sync(&mut self, ack_level: AckLevel) -> Result<()> {
1072 self.backend.sync(ack_level)
1073 }
1074}
1075
1076impl DirectColumnBackend {
1077 fn sync(&mut self, ack_level: AckLevel) -> Result<()> {
1078 let first_frame_sent = self.first_frame_sent;
1091 let mut commit_chunk = Chunk::new("");
1092 let mut result = self.flush_inner(&mut commit_chunk, WaitForAck::Yes(ack_level));
1093 self.first_frame_sent = first_frame_sent;
1094 if self.commit_since_sync {
1095 result = result.map_err(deny_retry_after_partial);
1096 }
1097 result.map_err(FlushFailure::into_error)
1098 }
1099
1100 fn flush_inner(
1101 &mut self,
1102 chunk: &mut Chunk<'_>,
1103 wait: WaitForAck,
1104 ) -> std::result::Result<(), FlushFailure> {
1105 let defer_commit = match wait {
1111 WaitForAck::No => self.first_frame_sent,
1112 WaitForAck::Yes(level) => {
1113 self.conn
1114 .validate_ack_level(level)
1115 .map_err(direct_not_delivered)?;
1116 false
1117 }
1118 };
1119
1120 self.conn.try_drain_acks().map_err(direct_not_delivered)?;
1121
1122 match self.publish_frame(chunk, None, defer_commit)? {
1127 FrameOutcome::Published => {}
1128 FrameOutcome::NoSlot(err) => return Err(FlushFailure::NotDelivered(err)),
1129 FrameOutcome::TooLarge(err) => {
1130 let row_count = chunk.row_count();
1131 match split_mid(row_count) {
1132 Some(mid) => {
1133 let mut committed = false;
1139 let mut result = self.publish_split(chunk, 0, mid, true, &mut committed);
1140 if result.is_ok() {
1141 result = self.publish_split(
1142 chunk,
1143 mid,
1144 row_count - mid,
1145 defer_commit,
1146 &mut committed,
1147 );
1148 }
1149 if let Err(e) = result {
1150 self.conn.mark_must_close();
1151 return Err(if committed {
1157 deny_retry_after_partial(e)
1158 } else {
1159 e
1160 });
1161 }
1162 }
1163 None => return Err(direct_not_delivered(err)),
1164 }
1165 }
1166 }
1167
1168 chunk.clear();
1171
1172 if let WaitForAck::Yes(level) = wait {
1173 self.conn
1174 .sync_all_acks(level)
1175 .map_err(direct_delivery_unknown)?;
1176 self.commit_since_sync = false;
1177 }
1178 Ok(())
1179 }
1180
1181 fn publish_frame(
1187 &mut self,
1188 chunk: &Chunk<'_>,
1189 range: Option<(usize, usize)>,
1190 defer_commit: bool,
1191 ) -> std::result::Result<FrameOutcome, FlushFailure> {
1192 if defer_commit && !self.conn.has_sync_commit_slot() {
1193 return Ok(FrameOutcome::NoSlot(error::fmt!(
1194 InvalidApiCall,
1195 "column sender deferred flush capacity exhausted; call sync() \
1196 before flushing more chunks."
1197 )));
1198 }
1199
1200 if self.conn.at_in_flight_cap() {
1201 self.conn
1202 .drain_one_ack_blocking()
1203 .map_err(direct_not_delivered)?;
1204 }
1205
1206 let dict_mark = self.symbol_dict.mark();
1207 let result = self.conn.publish_qwp(|out| match range {
1208 None => encoder::encode_chunk_into(
1209 out,
1210 chunk,
1211 &mut self.symbol_dict,
1212 &mut self.scratch,
1213 defer_commit,
1214 ),
1215 Some((offset, count)) => {
1216 let view = unsafe { chunk.slice_rows(offset, count) };
1217 encoder::encode_chunk_into(
1218 out,
1219 &view,
1220 &mut self.symbol_dict,
1221 &mut self.scratch,
1222 defer_commit,
1223 )
1224 }
1225 });
1226
1227 match result {
1228 Ok(published) => {
1229 self.conn.push_pending(published.fsn);
1230 self.first_frame_sent = true;
1231 Ok(FrameOutcome::Published)
1232 }
1233 Err(PublishError::BeforeWrite(e)) if e.code() == ErrorCode::BatchTooLarge => {
1234 self.symbol_dict.rollback(dict_mark);
1235 Ok(FrameOutcome::TooLarge(e))
1236 }
1237 Err(PublishError::BeforeWrite(e)) => {
1238 if e.code() != ErrorCode::SocketError {
1239 self.symbol_dict.rollback(dict_mark);
1240 }
1241 self.latch_if_connection_is_spent(&e);
1242 Err(direct_not_delivered(e))
1243 }
1244 Err(PublishError::DuringWrite(e)) => Err(direct_delivery_unknown(e)),
1247 }
1248 }
1249
1250 fn latch_if_connection_is_spent(&mut self, err: &crate::Error) {
1282 if err.code() == ErrorCode::SymbolDictFull {
1283 self.conn.mark_spent();
1284 }
1285 }
1286
1287 fn publish_split(
1292 &mut self,
1293 chunk: &Chunk<'_>,
1294 row_offset: usize,
1295 row_count: usize,
1296 defer_commit: bool,
1297 committed: &mut bool,
1298 ) -> std::result::Result<(), FlushFailure> {
1299 let outcome =
1300 match self.publish_frame(chunk, Some((row_offset, row_count)), defer_commit)? {
1301 FrameOutcome::NoSlot(_) => {
1302 self.sync(AckLevel::Ok)
1308 .map_err(FlushFailure::DeliveryUnknown)?;
1309 *committed = true;
1313 self.commit_since_sync = true;
1314 self.publish_frame(chunk, Some((row_offset, row_count)), defer_commit)?
1315 }
1316 outcome => outcome,
1317 };
1318 match outcome {
1319 FrameOutcome::Published => Ok(()),
1320 FrameOutcome::NoSlot(err) => Err(FlushFailure::NotDelivered(err)),
1321 FrameOutcome::TooLarge(err) => match split_mid(row_count) {
1322 Some(mid) => {
1323 self.publish_split(chunk, row_offset, mid, true, committed)?;
1324 self.publish_split(
1325 chunk,
1326 row_offset + mid,
1327 row_count - mid,
1328 defer_commit,
1329 committed,
1330 )
1331 }
1332 None => Err(direct_not_delivered(err)),
1333 },
1334 }
1335 }
1336
1337 #[cfg(feature = "arrow-ingress")]
1338 #[allow(clippy::too_many_arguments)]
1339 fn flush_arrow_batch_inner(
1340 &mut self,
1341 table: TableName<'_>,
1342 batch: &RecordBatch,
1343 ts: ArrowTsSource,
1344 overrides: &[ArrowColumnOverride<'_>],
1345 wait: WaitForAck,
1346 ) -> std::result::Result<(), FlushFailure> {
1347 let defer_commit = match wait {
1348 WaitForAck::No => self.first_frame_sent,
1349 WaitForAck::Yes(level) => {
1350 self.conn
1351 .validate_ack_level(level)
1352 .map_err(direct_not_delivered)?;
1353 false
1354 }
1355 };
1356
1357 self.conn.try_drain_acks().map_err(direct_not_delivered)?;
1358
1359 let spec = ArrowFrameSpec {
1360 table,
1361 batch,
1362 ts,
1363 overrides,
1364 };
1365 match self.publish_arrow_frame(&spec, None, defer_commit)? {
1369 FrameOutcome::Published => {}
1370 FrameOutcome::NoSlot(err) => return Err(FlushFailure::NotDelivered(err)),
1371 FrameOutcome::TooLarge(err) => {
1372 let row_count = batch.num_rows();
1373 match split_mid(row_count) {
1374 Some(mid) => {
1375 let mut committed = false;
1376 let mut result =
1377 self.publish_arrow_split(&spec, 0, mid, true, &mut committed);
1378 if result.is_ok() {
1379 result = self.publish_arrow_split(
1380 &spec,
1381 mid,
1382 row_count - mid,
1383 defer_commit,
1384 &mut committed,
1385 );
1386 }
1387 if let Err(e) = result {
1388 self.conn.mark_must_close();
1389 return Err(if committed {
1393 deny_retry_after_partial(e)
1394 } else {
1395 e
1396 });
1397 }
1398 }
1399 None => return Err(direct_not_delivered(err)),
1400 }
1401 }
1402 }
1403
1404 if let WaitForAck::Yes(level) = wait {
1405 self.conn
1406 .sync_all_acks(level)
1407 .map_err(direct_delivery_unknown)?;
1408 self.commit_since_sync = false;
1409 }
1410 Ok(())
1411 }
1412
1413 #[cfg(feature = "arrow-ingress")]
1417 fn publish_arrow_frame(
1418 &mut self,
1419 spec: &ArrowFrameSpec<'_>,
1420 range: Option<(usize, usize)>,
1421 defer_commit: bool,
1422 ) -> std::result::Result<FrameOutcome, FlushFailure> {
1423 if defer_commit && !self.conn.has_sync_commit_slot() {
1424 return Ok(FrameOutcome::NoSlot(error::fmt!(
1425 InvalidApiCall,
1426 "column sender deferred flush capacity exhausted; call sync() \
1427 before flushing more arrow batches."
1428 )));
1429 }
1430
1431 if self.conn.at_in_flight_cap() {
1432 self.conn
1433 .drain_one_ack_blocking()
1434 .map_err(direct_not_delivered)?;
1435 }
1436
1437 let dict_mark = self.symbol_dict.mark();
1438 let sliced;
1439 let batch = match range {
1440 None => spec.batch,
1441 Some((offset, count)) => {
1442 sliced = spec.batch.slice(offset, count);
1443 &sliced
1444 }
1445 };
1446 let result = self.conn.publish_qwp(|out| {
1447 arrow_batch::encode_arrow_batch_into(
1448 out,
1449 spec.table,
1450 batch,
1451 spec.ts,
1452 spec.overrides,
1453 &mut self.symbol_dict,
1454 defer_commit,
1455 )
1456 });
1457
1458 match result {
1459 Ok(published) => {
1460 self.conn.push_pending(published.fsn);
1461 self.first_frame_sent = true;
1462 Ok(FrameOutcome::Published)
1463 }
1464 Err(PublishError::BeforeWrite(e)) if e.code() == ErrorCode::BatchTooLarge => {
1465 self.symbol_dict.rollback(dict_mark);
1466 Ok(FrameOutcome::TooLarge(e))
1467 }
1468 Err(PublishError::BeforeWrite(e)) => {
1469 if e.code() != ErrorCode::SocketError {
1470 self.symbol_dict.rollback(dict_mark);
1471 }
1472 self.latch_if_connection_is_spent(&e);
1473 Err(direct_not_delivered(e))
1474 }
1475 Err(PublishError::DuringWrite(e)) => Err(direct_delivery_unknown(e)),
1476 }
1477 }
1478
1479 #[cfg(feature = "arrow-ingress")]
1481 fn publish_arrow_split(
1482 &mut self,
1483 spec: &ArrowFrameSpec<'_>,
1484 row_offset: usize,
1485 row_count: usize,
1486 defer_commit: bool,
1487 committed: &mut bool,
1488 ) -> std::result::Result<(), FlushFailure> {
1489 let outcome =
1490 match self.publish_arrow_frame(spec, Some((row_offset, row_count)), defer_commit)? {
1491 FrameOutcome::NoSlot(_) => {
1492 self.sync(AckLevel::Ok)
1495 .map_err(FlushFailure::DeliveryUnknown)?;
1496 *committed = true;
1498 self.commit_since_sync = true;
1499 self.publish_arrow_frame(spec, Some((row_offset, row_count)), defer_commit)?
1500 }
1501 outcome => outcome,
1502 };
1503 match outcome {
1504 FrameOutcome::Published => Ok(()),
1505 FrameOutcome::NoSlot(err) => Err(FlushFailure::NotDelivered(err)),
1506 FrameOutcome::TooLarge(err) => match split_mid(row_count) {
1507 Some(mid) => {
1508 self.publish_arrow_split(spec, row_offset, mid, true, committed)?;
1509 self.publish_arrow_split(
1510 spec,
1511 row_offset + mid,
1512 row_count - mid,
1513 defer_commit,
1514 committed,
1515 )
1516 }
1517 None => Err(direct_not_delivered(err)),
1518 },
1519 }
1520 }
1521}
1522
1523impl SfaBackend {
1524 fn validate_ack_level(&self, ack_level: AckLevel) -> Result<()> {
1528 if ack_level == AckLevel::Durable && !self.request_durable_ack {
1529 return Err(error::fmt!(
1530 InvalidApiCall,
1531 "AckLevel::Durable requires the pool to be opened with \
1532 `request_durable_ack=on` in the connect string."
1533 ));
1534 }
1535 Ok(())
1536 }
1537
1538 fn latch_if_connection_is_spent(&mut self, err: &crate::Error) {
1560 if err.code() == ErrorCode::SymbolDictFull {
1561 self.drop_on_return = true;
1562 }
1563 }
1564
1565 fn publish_buffer(
1566 &mut self,
1567 buffer: &QwpWsColumnarBuffer,
1568 ack_level: Option<AckLevel>,
1569 ) -> std::result::Result<Option<u64>, FlushFailure> {
1570 if let Some(level) = ack_level {
1571 self.validate_ack_level(level)
1572 .map_err(FlushFailure::NotDelivered)?;
1573 }
1574 if let Err(err) = qwp_ws_check_error_background(&self.state) {
1575 return Err(FlushFailure::NotDelivered(err));
1576 }
1577 if buffer.is_empty() {
1578 return Ok(None);
1579 }
1580
1581 let frame_cap = self.effective_frame_caps().hard;
1584 let result = {
1585 let Self {
1586 foreground,
1587 state,
1588 buffer_scratch,
1589 ..
1590 } = self;
1591 foreground.encode_persist_publish(
1592 frame_cap,
1593 |payload, symbol_dict, delta_enabled| {
1594 buffer.encode_ws_replay_message_with_defer(
1595 payload,
1596 buffer_scratch,
1597 symbol_dict,
1598 super::wire::QWP_VERSION_1,
1599 false,
1600 delta_enabled,
1601 )
1602 },
1603 |payload| publish_qwp_ws_payload_background(state, payload, frame_cap),
1604 )
1605 };
1606 if let Err(err) = &result {
1607 self.latch_if_connection_is_spent(err);
1608 }
1609 match result.map_err(FlushFailure::NotDelivered)? {
1610 SfaPublishOutcome::Published(fsn) => Ok(Some(fsn)),
1611 SfaPublishOutcome::TooLarge {
1612 encoded_len,
1613 max_buf_size,
1614 } => Err(FlushFailure::NotDelivered(sfa_frame_size_error(
1615 encoded_len,
1616 max_buf_size,
1617 ))),
1618 }
1619 }
1620
1621 fn flush_chunk(
1622 &mut self,
1623 chunk: &mut Chunk<'_>,
1624 wait: WaitForAck,
1625 ) -> std::result::Result<(), FlushFailure> {
1626 self.flush_chunk_boundary(chunk, wait).map(|_| ())
1627 }
1628
1629 fn flush_chunk_and_get_fsn(
1630 &mut self,
1631 chunk: &mut Chunk<'_>,
1632 ) -> std::result::Result<u64, FlushFailure> {
1633 self.flush_chunk_boundary(chunk, WaitForAck::No)
1634 }
1635
1636 fn flush_chunk_boundary(
1637 &mut self,
1638 chunk: &mut Chunk<'_>,
1639 wait: WaitForAck,
1640 ) -> std::result::Result<u64, FlushFailure> {
1641 if let WaitForAck::Yes(level) = wait {
1644 self.validate_ack_level(level)
1645 .map_err(FlushFailure::NotDelivered)?;
1646 }
1647 if let Err(e) = qwp_ws_check_error_background(&self.state) {
1648 return Err(FlushFailure::NotDelivered(e));
1649 }
1650 let caps = self.effective_frame_caps();
1651 let boundary =
1660 match self.publish_chunk_sfa(chunk, None, caps.for_range(chunk.row_count()))? {
1661 SfaPublishOutcome::Published(fsn) => fsn,
1662 SfaPublishOutcome::TooLarge {
1663 encoded_len,
1664 max_buf_size,
1665 } => {
1666 let err = sfa_frame_size_error(encoded_len, max_buf_size);
1667 let row_count = chunk.row_count();
1668 match split_mid(row_count) {
1669 Some(mid) => {
1670 self.publish_split_sfa(chunk, 0, mid, caps)?;
1671 self.publish_split_sfa(chunk, mid, row_count - mid, caps)
1675 .map_err(deny_retry_after_partial)?
1676 }
1677 None => return Err(FlushFailure::NotDelivered(err)),
1678 }
1679 }
1680 };
1681 chunk.clear();
1682 if let WaitForAck::Yes(level) = wait {
1683 self.wait_for_boundary(level, boundary, self.sync_timeout)
1685 .map_err(FlushFailure::DeliveryUnknown)?;
1686 }
1687 Ok(boundary)
1688 }
1689
1690 fn publish_chunk_sfa(
1695 &mut self,
1696 chunk: &Chunk<'_>,
1697 range: Option<(usize, usize)>,
1698 frame_cap: usize,
1699 ) -> std::result::Result<SfaPublishOutcome, FlushFailure> {
1700 let view;
1701 let target = match range {
1702 None => chunk,
1703 Some((offset, count)) => {
1704 view = unsafe { chunk.slice_rows(offset, count) };
1705 &view
1706 }
1707 };
1708 let result = {
1709 let Self {
1710 state,
1711 foreground,
1712 scratch,
1713 ..
1714 } = self;
1715 foreground.encode_persist_publish(
1716 frame_cap,
1717 |payload, symbol_dict, delta_enabled| {
1718 if delta_enabled {
1719 encoder::encode_chunk_into(payload, target, symbol_dict, scratch, false)
1720 } else {
1721 encoder::encode_chunk_replay_into(payload, target, symbol_dict, scratch)
1722 }
1723 },
1724 |encoded| publish_qwp_ws_payload_background(state, encoded, frame_cap),
1725 )
1726 };
1727 if let Err(err) = &result {
1728 self.latch_if_connection_is_spent(err);
1729 }
1730 result.map_err(FlushFailure::NotDelivered)
1731 }
1732
1733 fn publish_split_sfa(
1736 &mut self,
1737 chunk: &Chunk<'_>,
1738 row_offset: usize,
1739 row_count: usize,
1740 caps: SfaFrameCaps,
1741 ) -> std::result::Result<u64, FlushFailure> {
1742 match self.publish_chunk_sfa(
1743 chunk,
1744 Some((row_offset, row_count)),
1745 caps.for_range(row_count),
1746 )? {
1747 SfaPublishOutcome::Published(fsn) => Ok(fsn),
1748 SfaPublishOutcome::TooLarge {
1749 encoded_len,
1750 max_buf_size,
1751 } => match split_mid(row_count) {
1752 Some(mid) => {
1753 self.publish_split_sfa(chunk, row_offset, mid, caps)?;
1754 self.publish_split_sfa(chunk, row_offset + mid, row_count - mid, caps)
1755 .map_err(deny_retry_after_partial)
1756 }
1757 None => Err(FlushFailure::NotDelivered(sfa_frame_size_error(
1758 encoded_len,
1759 max_buf_size,
1760 ))),
1761 },
1762 }
1763 }
1764
1765 #[cfg(feature = "arrow-ingress")]
1766 fn flush_arrow_batch(
1767 &mut self,
1768 table: TableName<'_>,
1769 batch: &RecordBatch,
1770 ts: ArrowTsSource,
1771 overrides: &[ArrowColumnOverride<'_>],
1772 wait: WaitForAck,
1773 ) -> std::result::Result<(), FlushFailure> {
1774 self.flush_arrow_batch_boundary(table, batch, ts, overrides, wait)
1775 .map(|_| ())
1776 }
1777
1778 #[cfg(feature = "arrow-ingress")]
1779 fn flush_arrow_batch_and_get_fsn(
1780 &mut self,
1781 table: TableName<'_>,
1782 batch: &RecordBatch,
1783 ts: ArrowTsSource,
1784 overrides: &[ArrowColumnOverride<'_>],
1785 ) -> std::result::Result<u64, FlushFailure> {
1786 self.flush_arrow_batch_boundary(table, batch, ts, overrides, WaitForAck::No)
1787 }
1788
1789 #[cfg(feature = "arrow-ingress")]
1790 fn flush_arrow_batch_boundary(
1791 &mut self,
1792 table: TableName<'_>,
1793 batch: &RecordBatch,
1794 ts: ArrowTsSource,
1795 overrides: &[ArrowColumnOverride<'_>],
1796 wait: WaitForAck,
1797 ) -> std::result::Result<u64, FlushFailure> {
1798 if let WaitForAck::Yes(level) = wait {
1799 self.validate_ack_level(level)
1800 .map_err(FlushFailure::NotDelivered)?;
1801 }
1802 if let Err(e) = qwp_ws_check_error_background(&self.state) {
1803 return Err(FlushFailure::NotDelivered(e));
1804 }
1805 let caps = self.effective_frame_caps();
1806 let spec = ArrowFrameSpec {
1807 table,
1808 batch,
1809 ts,
1810 overrides,
1811 };
1812 let boundary =
1815 match self.publish_arrow_sfa(&spec, None, caps.for_range(batch.num_rows()))? {
1816 SfaPublishOutcome::Published(fsn) => fsn,
1817 SfaPublishOutcome::TooLarge {
1818 encoded_len,
1819 max_buf_size,
1820 } => {
1821 let err = sfa_frame_size_error(encoded_len, max_buf_size);
1822 let row_count = batch.num_rows();
1823 match split_mid(row_count) {
1824 Some(mid) => {
1825 self.publish_arrow_split_sfa(&spec, 0, mid, caps)?;
1826 self.publish_arrow_split_sfa(&spec, mid, row_count - mid, caps)
1828 .map_err(deny_retry_after_partial)?
1829 }
1830 None => return Err(FlushFailure::NotDelivered(err)),
1831 }
1832 }
1833 };
1834 if let WaitForAck::Yes(level) = wait {
1835 self.wait_for_boundary(level, boundary, self.sync_timeout)
1836 .map_err(FlushFailure::DeliveryUnknown)?;
1837 }
1838 Ok(boundary)
1839 }
1840
1841 #[cfg(feature = "arrow-ingress")]
1843 fn publish_arrow_sfa(
1844 &mut self,
1845 spec: &ArrowFrameSpec<'_>,
1846 range: Option<(usize, usize)>,
1847 frame_cap: usize,
1848 ) -> std::result::Result<SfaPublishOutcome, FlushFailure> {
1849 let sliced;
1850 let batch = match range {
1851 None => spec.batch,
1852 Some((offset, count)) => {
1853 sliced = spec.batch.slice(offset, count);
1854 &sliced
1855 }
1856 };
1857 let result = {
1858 let Self {
1859 state, foreground, ..
1860 } = self;
1861 foreground.encode_persist_publish(
1862 frame_cap,
1863 |payload, symbol_dict, delta_enabled| {
1864 if delta_enabled {
1865 arrow_batch::encode_arrow_batch_into(
1866 payload,
1867 spec.table,
1868 batch,
1869 spec.ts,
1870 spec.overrides,
1871 symbol_dict,
1872 false,
1873 )
1874 } else {
1875 arrow_batch::encode_arrow_batch_replay_into(
1876 payload,
1877 spec.table,
1878 batch,
1879 spec.ts,
1880 spec.overrides,
1881 symbol_dict,
1882 )
1883 }
1884 },
1885 |payload| publish_qwp_ws_payload_background(state, payload, frame_cap),
1886 )
1887 };
1888 if let Err(err) = &result {
1889 self.latch_if_connection_is_spent(err);
1890 }
1891 result.map_err(FlushFailure::NotDelivered)
1892 }
1893
1894 #[cfg(feature = "arrow-ingress")]
1897 fn publish_arrow_split_sfa(
1898 &mut self,
1899 spec: &ArrowFrameSpec<'_>,
1900 row_offset: usize,
1901 row_count: usize,
1902 caps: SfaFrameCaps,
1903 ) -> std::result::Result<u64, FlushFailure> {
1904 match self.publish_arrow_sfa(
1905 spec,
1906 Some((row_offset, row_count)),
1907 caps.for_range(row_count),
1908 )? {
1909 SfaPublishOutcome::Published(fsn) => Ok(fsn),
1910 SfaPublishOutcome::TooLarge {
1911 encoded_len,
1912 max_buf_size,
1913 } => match split_mid(row_count) {
1914 Some(mid) => {
1915 self.publish_arrow_split_sfa(spec, row_offset, mid, caps)?;
1916 self.publish_arrow_split_sfa(spec, row_offset + mid, row_count - mid, caps)
1917 .map_err(deny_retry_after_partial)
1918 }
1919 None => Err(FlushFailure::NotDelivered(sfa_frame_size_error(
1920 encoded_len,
1921 max_buf_size,
1922 ))),
1923 },
1924 }
1925 }
1926
1927 fn sync(&mut self, ack_level: AckLevel) -> Result<()> {
1928 self.wait(ack_level, self.sync_timeout)
1929 }
1930
1931 fn wait(&mut self, ack_level: AckLevel, timeout: Duration) -> Result<()> {
1932 self.validate_ack_level(ack_level)?;
1933 let Some(boundary) = qwp_ws_published_fsn_background(&self.state)? else {
1934 return Ok(());
1935 };
1936 self.wait_for_boundary(ack_level, boundary, timeout)
1937 }
1938
1939 fn published_fsn(&self) -> Result<Option<u64>> {
1940 qwp_ws_published_fsn_background(&self.state)
1941 }
1942
1943 fn acked_fsn(&self) -> Result<Option<u64>> {
1944 qwp_ws_acked_fsn_background(&self.state)
1945 }
1946
1947 fn wait_for_boundary(
1952 &mut self,
1953 ack_level: AckLevel,
1954 boundary: u64,
1955 timeout: Duration,
1956 ) -> Result<()> {
1957 let last_boundary = match ack_level {
1958 AckLevel::Ok => self.last_ok_sync_boundary,
1959 AckLevel::Durable => self.last_durable_sync_boundary,
1960 };
1961 if last_boundary.is_some_and(|last| last >= boundary) {
1962 return Ok(());
1963 }
1964
1965 let mut deadline_anchor = Instant::now();
1971 let mut last_completed: Option<u64> = None;
1972
1973 loop {
1974 let completed = match ack_level {
1975 AckLevel::Ok => qwp_ws_ok_fsn_background(&self.state)?,
1976 AckLevel::Durable => qwp_ws_acked_fsn_background(&self.state)?,
1977 };
1978 if completed.is_some_and(|fsn| fsn >= boundary) {
1979 match ack_level {
1980 AckLevel::Ok => self.last_ok_sync_boundary = Some(boundary),
1981 AckLevel::Durable => self.last_durable_sync_boundary = Some(boundary),
1982 }
1983 return Ok(());
1984 }
1985 if completed != last_completed {
1986 last_completed = completed;
1987 deadline_anchor = Instant::now();
1988 }
1989
1990 qwp_ws_check_error_background(&self.state)?;
1991
1992 if !timeout.is_zero() && deadline_anchor.elapsed() >= timeout {
1993 return Err(sfa_sync_timeout(timeout, ack_level, boundary, completed));
1994 }
1995 thread::sleep(Duration::from_millis(10));
1996 }
1997 }
1998
1999 fn effective_hard_frame_cap(&self) -> (usize, bool) {
2000 let server_max = self.state.server_max_batch_size.load(Ordering::Acquire);
2001 effective_hard_frame_cap(
2002 self.max_buf_size,
2003 server_max,
2004 self.state.sfa_frame_payload_cap,
2005 )
2006 }
2007
2008 fn effective_frame_caps(&self) -> SfaFrameCaps {
2009 let (hard, _) = self.effective_hard_frame_cap();
2010 SfaFrameCaps {
2011 hard,
2012 soft: hard.min(self.state.sfa_frame_split_target),
2013 }
2014 }
2015}
2016
2017fn effective_hard_frame_cap(
2018 max_buf_size: usize,
2019 server_max_batch_size: usize,
2020 sfa_frame_payload_cap: usize,
2021) -> (usize, bool) {
2022 let configured_cap = if server_max_batch_size == 0 {
2023 max_buf_size
2024 } else {
2025 max_buf_size.min(server_max_batch_size)
2026 };
2027 (
2028 configured_cap.min(sfa_frame_payload_cap),
2029 server_max_batch_size != 0,
2030 )
2031}
2032
2033fn sfa_sync_timeout(
2047 sync_timeout: Duration,
2048 ack_level: AckLevel,
2049 boundary: u64,
2050 completed: Option<u64>,
2051) -> crate::Error {
2052 let level = match ack_level {
2053 AckLevel::Ok => "ok",
2054 AckLevel::Durable => "durable",
2055 };
2056 let progress = match completed {
2057 Some(fsn) => format!("reached FSN {}", fsn),
2058 None => "reached no frame".to_string(),
2059 };
2060 crate::Error::new(
2061 ErrorCode::FailoverRetry,
2062 format!(
2063 "QWP/WebSocket store-and-forward wait({}) timed out after {:?} \
2064 with no ack progress (target FSN {}, {}); the connection is alive \
2065 but the server is not advancing the watermark. The frames remain \
2066 queued and the background runner keeps delivering them: retry \
2067 wait() to keep awaiting the ack, or close the pool to drain. Do \
2068 not re-flush the same data, which is already accepted and would \
2069 be delivered twice.",
2070 level, sync_timeout, boundary, progress
2071 ),
2072 )
2073}
2074
2075#[cfg(test)]
2076mod tests {
2077 use super::{effective_hard_frame_cap, split_mid};
2078
2079 #[test]
2080 fn effective_hard_cap_reports_whether_the_server_cap_is_known() {
2081 assert_eq!(effective_hard_frame_cap(1000, 0, 800), (800, false));
2082 assert_eq!(effective_hard_frame_cap(1000, 400, 800), (400, true));
2083 assert_eq!(effective_hard_frame_cap(1000, 1200, 800), (800, true));
2084 }
2085
2086 #[test]
2087 fn split_mid_floors_at_eight_rows() {
2088 assert_eq!(split_mid(0), None);
2089 assert_eq!(split_mid(1), None);
2090 assert_eq!(split_mid(8), None);
2091 }
2092
2093 #[test]
2094 fn split_mid_returns_eight_aligned_point_below_count() {
2095 for count in [9usize, 12, 15, 16, 17, 100, 10_000, 16_384] {
2096 let mid = split_mid(count).unwrap();
2097 assert_eq!(
2098 mid % 8,
2099 0,
2100 "split point must be 8-aligned for count {count}"
2101 );
2102 assert!(mid >= 8, "split point must be at least 8 for count {count}");
2103 assert!(
2104 mid < count,
2105 "split point must make progress for count {count}"
2106 );
2107 }
2108 }
2109
2110 #[test]
2111 fn sfa_frame_caps_use_split_target_only_while_range_can_split() {
2112 use super::SfaFrameCaps;
2113
2114 let caps = SfaFrameCaps {
2115 hard: 1000,
2116 soft: 400,
2117 };
2118 assert_eq!(caps.for_range(1), 1000);
2121 assert_eq!(caps.for_range(8), 1000);
2122 assert_eq!(caps.for_range(9), 400);
2124 assert_eq!(caps.for_range(10_000), 400);
2125 }
2126
2127 #[test]
2128 fn deny_retry_after_partial_downgrades_not_delivered_and_never_upgrades() {
2129 use super::{FlushFailure, deny_retry_after_partial};
2130 use crate::{Error, ErrorCode};
2131
2132 let nd = FlushFailure::NotDelivered(Error::new(ErrorCode::SocketError, "boom"));
2137 assert!(nd.is_not_delivered());
2138 let downgraded = deny_retry_after_partial(nd);
2139 assert!(!downgraded.is_not_delivered());
2140 assert!(
2141 downgraded.into_error().in_doubt(),
2142 "downgraded failure must be flagged in-doubt"
2143 );
2144
2145 let du = FlushFailure::DeliveryUnknown(Error::new(ErrorCode::SocketError, "boom"));
2148 let still = deny_retry_after_partial(du);
2149 assert!(!still.is_not_delivered());
2150 assert!(still.into_error().in_doubt());
2151 }
2152}