use std::fmt::{self, Debug, Formatter};
use std::marker::PhantomData;
#[cfg(feature = "_egress")]
use std::ops::{Deref, DerefMut};
use std::path::{Path, PathBuf};
use std::rc::Rc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::thread::{self, JoinHandle};
use std::time::{Duration, Instant};
#[cfg(feature = "_egress")]
use crate::egress::Reader;
use crate::ingress::conn_events;
use crate::ingress::rejection_events;
use crate::ingress::sender::is_candidate_orphan;
use crate::ingress::sender::qwp_ws::QwpWsHostHealthTracker;
use crate::ingress::{Buffer, SenderBuilder};
use crate::ingress::{
QwpWsConnector, QwpWsManagedSlotExclusion, RawQwpWsRoundStream, ReconnectReason,
};
#[cfg_attr(
not(any(
feature = "polars-ingress",
feature = "polars-egress",
feature = "ffi-support"
)),
allow(unused_imports)
)]
use crate::ingress::{reconnect_backoff_step, reconnect_error_is_terminal};
use crate::{Result, error};
mod conf;
use crate::ingress::AckLevel;
use crate::ingress::column_sender::conn::ColumnConn;
use crate::ingress::column_sender::{DirectSenderCore, PooledSenderCore};
use conf::PoolReap;
#[cfg(feature = "ffi-support")]
#[doc(hidden)]
pub mod ffi_support;
const REAPER_MIN_TICK: Duration = Duration::from_secs(5);
fn lock_state<S>(m: &Mutex<PoolState<S>>) -> std::sync::MutexGuard<'_, PoolState<S>> {
m.lock().unwrap_or_else(|e| e.into_inner())
}
fn lock_health(
m: &Mutex<QwpWsHostHealthTracker>,
) -> std::sync::MutexGuard<'_, QwpWsHostHealthTracker> {
m.lock().unwrap_or_else(|e| e.into_inner())
}
#[cfg(feature = "_egress")]
fn lock_reader_state(m: &Mutex<ReaderPoolState>) -> std::sync::MutexGuard<'_, ReaderPoolState> {
m.lock().unwrap_or_else(|e| e.into_inner())
}
struct InUseSlot<'a, S> {
state: &'a Mutex<PoolState<S>>,
cv: &'a Condvar,
slot_index: Option<usize>,
armed: bool,
}
impl<S> InUseSlot<'_, S> {
fn commit(mut self) {
self.armed = false;
}
}
impl<S> Drop for InUseSlot<'_, S> {
fn drop(&mut self) {
if self.armed {
let mut state = lock_state(self.state);
state.in_use = state.in_use.saturating_sub(1);
state.free_slot_index(self.slot_index);
self.cv.notify_all();
}
}
}
#[cfg(feature = "_egress")]
struct ReaderInUseSlot<'a> {
inner: &'a DbInner,
armed: bool,
}
#[cfg(feature = "_egress")]
impl ReaderInUseSlot<'_> {
fn commit(mut self) {
self.armed = false;
}
}
#[cfg(feature = "_egress")]
impl Drop for ReaderInUseSlot<'_> {
fn drop(&mut self) {
if self.armed {
{
let mut state = lock_reader_state(&self.inner.reader_state);
state.in_use = state.in_use.saturating_sub(1);
}
self.inner.reader_cv.notify_all();
}
}
}
struct SenderSlotRelease<'a> {
inner: &'a DbInner,
slot_index: Option<usize>,
decrement_in_use: bool,
decrement_closing: bool,
}
impl Drop for SenderSlotRelease<'_> {
fn drop(&mut self) {
if self.slot_index.is_none() && !self.decrement_in_use && !self.decrement_closing {
return;
}
let mut state = lock_state(&self.inner.state);
if self.decrement_in_use {
state.in_use = state.in_use.saturating_sub(1);
}
if self.decrement_closing {
state.closing = state.closing.saturating_sub(1);
}
state.free_slot_index(self.slot_index);
self.inner.cv.notify_all();
}
}
#[derive(Default)]
#[non_exhaustive]
pub struct ConnectHandlers {
pub connection_listener: Option<crate::ingress::ConnectionListener>,
pub connection_event_inbox_capacity: usize,
pub error_handler: Option<crate::ingress::QwpWsErrorHandler>,
pub error_inbox_capacity: usize,
}
pub struct QuestDb {
inner: Arc<DbInner>,
reaper: Option<JoinHandle<()>>,
}
struct DbInner {
#[cfg(feature = "_egress")]
conf: String,
connector: QwpWsConnector,
buffer_max_name_len: usize,
health: Mutex<QwpWsHostHealthTracker>,
sender_pool_min: usize,
sender_pool_max: usize,
#[cfg(feature = "_egress")]
query_pool_min: usize,
#[cfg(feature = "_egress")]
query_pool_max: usize,
acquire_timeout: Duration,
sf_disk: bool,
slot_base_id: String,
managed_slot_exclusion: Option<QwpWsManagedSlotExclusion>,
out_of_range_recovery_candidates: Vec<PathBuf>,
idle_timeout: Duration,
conn_events: Arc<conn_events::ConnectionEventSource>,
state: Mutex<PoolState<PooledSenderCore>>,
direct_state: Mutex<PoolState<DirectSenderCore>>,
#[cfg(feature = "_egress")]
reader_state: Mutex<ReaderPoolState>,
cv: Condvar,
direct_cv: Condvar,
#[cfg(feature = "_egress")]
reader_cv: Condvar,
rejections: Arc<rejection_events::RejectionEventSource>,
shutdown: AtomicBool,
}
#[derive(Default)]
struct SlotReservations(Option<Vec<bool>>);
impl SlotReservations {
fn with_disk_slots(pool_max: usize) -> Self {
Self(Some(vec![false; pool_max]))
}
fn reserved_total(&self, fallback_total: usize) -> usize {
match &self.0 {
Some(slots) => slots.iter().filter(|in_use| **in_use).count(),
None => fallback_total,
}
}
fn allocate(&mut self) -> Option<usize> {
let slots = self.0.as_mut()?;
let index = slots.iter().position(|in_use| !*in_use)?;
slots[index] = true;
Some(index)
}
fn reserve(&mut self, index: usize) -> bool {
let Some(slots) = self.0.as_mut() else {
return false;
};
let Some(slot) = slots.get_mut(index) else {
return false;
};
if *slot {
return false;
}
*slot = true;
true
}
fn free(&mut self, slot_index: Option<usize>) {
if let (Some(slots), Some(index)) = (&mut self.0, slot_index)
&& let Some(slot) = slots.get_mut(index)
{
*slot = false;
}
}
}
struct PoolState<S> {
free: Vec<PoolEntry<S>>,
in_use: usize,
closing: usize,
slots: SlotReservations,
}
impl<S> Default for PoolState<S> {
fn default() -> Self {
Self {
free: Vec::new(),
in_use: 0,
closing: 0,
slots: SlotReservations::default(),
}
}
}
impl<S> PoolState<S> {
fn total(&self) -> usize {
self.free.len() + self.in_use
}
fn with_disk_slots(pool_max: usize) -> Self {
Self {
free: Vec::new(),
in_use: 0,
closing: 0,
slots: SlotReservations::with_disk_slots(pool_max),
}
}
fn reserved_total(&self) -> usize {
self.slots.reserved_total(self.total())
}
fn allocate_slot_index(&mut self) -> Option<usize> {
self.slots.allocate()
}
fn reserve_slot_index(&mut self, index: usize) -> bool {
self.slots.reserve(index)
}
fn free_slot_index(&mut self, slot_index: Option<usize>) {
self.slots.free(slot_index);
}
}
struct PoolEntry<S> {
sender: S,
slot_index: Option<usize>,
last_idle_at: Instant,
}
struct PooledSender<S> {
sender: S,
slot_index: Option<usize>,
}
#[cfg(feature = "_egress")]
#[derive(Default)]
struct ReaderPoolState {
free: Vec<ReaderPoolEntry>,
in_use: usize,
}
#[cfg(feature = "_egress")]
impl ReaderPoolState {
fn total(&self) -> usize {
self.free.len() + self.in_use
}
}
#[cfg(feature = "_egress")]
struct ReaderPoolEntry {
reader: Reader,
last_idle_at: Instant,
}
#[doc(hidden)]
#[non_exhaustive]
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct DbgPoolCount {
pub free: usize,
pub in_use: usize,
pub closing: usize,
}
#[doc(hidden)]
#[non_exhaustive]
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub struct DbgPoolCounts {
pub ingress: DbgPoolCount,
pub column_direct: DbgPoolCount,
pub reader: DbgPoolCount,
}
struct ManagedSlotRecoveryCandidate {
index: usize,
path: PathBuf,
}
#[derive(Default)]
struct ManagedSlotRecoveryScan {
in_range: Vec<ManagedSlotRecoveryCandidate>,
out_of_range: Vec<PathBuf>,
}
fn managed_slot_exclusion(base: &str, pool_max: usize) -> QwpWsManagedSlotExclusion {
QwpWsManagedSlotExclusion::new(managed_slot_prefix(base), pool_max)
}
fn managed_slot_id(base: &str, index: usize) -> String {
managed_slot_exclusion(base, usize::MAX).slot_name(index)
}
fn managed_slot_prefix(base: &str) -> String {
format!("{base}-ingest-")
}
fn parse_managed_slot_id(base: &str, name: &str) -> Option<usize> {
managed_slot_exclusion(base, usize::MAX).parse_index(name)
}
fn managed_slot_recovery_scan_from(
sf_dir: &Path,
base: &str,
pool_max: usize,
) -> ManagedSlotRecoveryScan {
let Ok(entries) = std::fs::read_dir(sf_dir) else {
return ManagedSlotRecoveryScan::default();
};
let mut scan = ManagedSlotRecoveryScan::default();
for entry in entries.flatten() {
let slot_path = entry.path();
if !slot_path.is_dir() {
continue;
}
let Some(name) = entry.file_name().to_str().map(str::to_owned) else {
continue;
};
let Some(index) = parse_managed_slot_id(base, &name) else {
continue;
};
if !is_candidate_orphan(&slot_path) {
continue;
}
if index < pool_max {
scan.in_range.push(ManagedSlotRecoveryCandidate {
index,
path: slot_path,
});
} else {
log::warn!(
"adopting out-of-range store-and-forward slot `{}`; \
`<sender_id>-ingest-*` directories under \
sf_dir belong to the QuestDb pool namespace, so use a unique \
sender_id for pools sharing an sf_dir",
slot_path.display()
);
scan.out_of_range.push(slot_path);
}
}
scan
}
fn preopen_recovery_senders(
inner: &Arc<DbInner>,
in_range_candidates: &[ManagedSlotRecoveryCandidate],
) {
for candidate in in_range_candidates {
preopen_recovery_sender(
inner,
candidate.index,
&candidate.path,
&inner.out_of_range_recovery_candidates,
);
}
}
fn preopen_recovery_sender(
inner: &Arc<DbInner>,
index: usize,
slot_path: &Path,
recovery_candidates: &[PathBuf],
) {
let slot = {
let mut state = lock_state(&inner.state);
if !state.reserve_slot_index(index) {
return;
}
state.in_use += 1;
InUseSlot {
state: &inner.state,
cv: &inner.cv,
slot_index: Some(index),
armed: true,
}
};
match connect_sfa_pool_with_recovery_candidates(inner, Some(index), recovery_candidates, true) {
Ok(sender) => {
let slot_index = slot.slot_index;
{
let mut state = lock_state(&inner.state);
state.in_use = state.in_use.saturating_sub(1);
state.free.push(PoolEntry {
sender,
slot_index,
last_idle_at: Instant::now(),
});
}
slot.commit();
inner.cv.notify_all();
}
Err(err) => {
log::warn!(
"skipping parked store-and-forward ingestion slot `{}` during recovery: {}",
slot_path.display(),
err
);
}
}
}
impl QuestDb {
pub fn connect(conf: &str) -> Result<Self> {
Self::connect_with_handlers(conf, ConnectHandlers::default())
}
pub fn connect_with_listener(
conf: &str,
listener: crate::ingress::ConnectionListener,
inbox_capacity: usize,
) -> Result<Self> {
Self::connect_with_handlers(
conf,
ConnectHandlers {
connection_listener: Some(listener),
connection_event_inbox_capacity: inbox_capacity,
..ConnectHandlers::default()
},
)
}
pub fn connect_with_handlers(conf: &str, handlers: ConnectHandlers) -> Result<Self> {
let conn_events = match handlers.connection_listener {
Some(listener) => conn_events::ConnectionEventSource::new(
listener,
handlers.connection_event_inbox_capacity,
),
None => conn_events::ConnectionEventSource::disabled(),
};
let rejections = match handlers.error_handler {
Some(handler) => rejection_events::RejectionEventSource::with_handler(
handler,
handlers.error_inbox_capacity,
),
None => rejection_events::RejectionEventSource::logging_default(),
};
Self::connect_impl(conf, conn_events, rejections)
}
fn connect_impl(
conf: &str,
conn_events: conn_events::ConnectionEventSource,
rejections: rejection_events::RejectionEventSource,
) -> Result<Self> {
let parsed = conf::parse(conf)?;
let sf_disk = parsed.sf_disk;
let pool_cfg = parsed.pool;
let mut builder = SenderBuilder::from_conf(conf)?;
if pool_cfg.lazy_connect {
builder.force_async_initial_connect();
}
let buffer_max_name_len = builder.configured_max_name_len();
let connector = builder.build_qwp_ws_connector()?;
let health = QwpWsHostHealthTracker::new(connector.endpoint_count());
let slot_base_id = connector.sender_id().to_owned();
let managed_slot_exclusion = if sf_disk {
Some(managed_slot_exclusion(
&slot_base_id,
pool_cfg.sender_pool_max,
))
} else {
None
};
let recovery_scan = if sf_disk {
connector
.sf_dir()
.map(|sf_dir| {
managed_slot_recovery_scan_from(sf_dir, &slot_base_id, pool_cfg.sender_pool_max)
})
.unwrap_or_default()
} else {
ManagedSlotRecoveryScan::default()
};
let ManagedSlotRecoveryScan {
in_range: in_range_recovery_candidates,
out_of_range: out_of_range_recovery_candidates,
} = recovery_scan;
let free = Vec::new();
let inner = Arc::new(DbInner {
#[cfg(feature = "_egress")]
conf: conf.to_owned(),
connector,
buffer_max_name_len,
health: Mutex::new(health),
sender_pool_min: pool_cfg.sender_pool_min,
sender_pool_max: pool_cfg.sender_pool_max,
#[cfg(feature = "_egress")]
query_pool_min: pool_cfg.query_pool_min,
#[cfg(feature = "_egress")]
query_pool_max: pool_cfg.query_pool_max,
acquire_timeout: pool_cfg.acquire_timeout,
sf_disk,
slot_base_id,
managed_slot_exclusion,
out_of_range_recovery_candidates,
idle_timeout: pool_cfg.idle_timeout,
state: Mutex::new(if sf_disk {
PoolState::with_disk_slots(pool_cfg.sender_pool_max)
} else {
PoolState {
free,
..PoolState::default()
}
}),
direct_state: Mutex::new(PoolState::default()),
#[cfg(feature = "_egress")]
reader_state: Mutex::new(ReaderPoolState::default()),
cv: Condvar::new(),
direct_cv: Condvar::new(),
#[cfg(feature = "_egress")]
reader_cv: Condvar::new(),
rejections: Arc::new(rejections),
shutdown: AtomicBool::new(false),
conn_events: Arc::new(conn_events),
});
let reaper = match pool_cfg.pool_reap {
PoolReap::Auto => Some(spawn_reaper(Arc::clone(&inner)).map_err(|err| {
inner.shutdown.store(true, Ordering::SeqCst);
crate::Error::new(
crate::ErrorCode::SocketError,
format!("Failed to spawn pool reaper thread: {err}"),
)
})?),
PoolReap::Manual => None,
};
let db = Self { inner, reaper };
if !pool_cfg.lazy_connect {
prewarm_min_connections(&db)?;
}
preopen_recovery_senders(&db.inner, &in_range_recovery_candidates);
Ok(db)
}
pub fn new_buffer(&self) -> Buffer {
Buffer::qwp_ws_with_max_name_len(self.inner.buffer_max_name_len)
}
#[doc(hidden)]
pub fn buffer_max_name_len(&self) -> usize {
self.inner.buffer_max_name_len
}
pub fn borrow_sender(&self) -> Result<BorrowedSender<'_>> {
let cs = self.pick_sender()?;
Ok(BorrowedSender(SenderHandle::new(self, cs)))
}
#[doc(hidden)]
pub fn borrow_direct_column_sender(&self) -> Result<BorrowedDirectColumnSender<'_>> {
let cs = pick_direct_sender(&self.inner)?;
Ok(BorrowedDirectColumnSender(DirectSenderHandle::new(
self, cs,
)))
}
#[cfg(feature = "arrow-ingress")]
pub fn flush_arrow_batch<'t, T>(
&self,
table: T,
batch: &arrow::array::RecordBatch,
timestamp_column: Option<crate::ingress::ColumnName<'_>>,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
ack_level: Option<AckLevel>,
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
let ack = ack_level.unwrap_or_else(|| self.default_ack_level());
let mut sender = self.borrow_direct_column_sender()?;
match timestamp_column {
Some(ts) => {
sender.flush_arrow_batch_at_column_and_wait(table, batch, ts, overrides, ack)
}
None => sender.flush_arrow_batch_at_now_and_wait(table, batch, overrides, ack),
}
}
#[cfg(feature = "arrow-ingress")]
pub(crate) fn default_ack_level(&self) -> AckLevel {
if self.inner.connector.request_durable_ack() {
AckLevel::Durable
} else {
AckLevel::Ok
}
}
#[cfg(feature = "ffi-support")]
pub(crate) fn borrow_sender_owned(&self) -> Result<OwnedSender> {
let cs = self.pick_sender()?;
Ok(OwnedSender::new(Arc::clone(&self.inner), cs))
}
#[cfg(feature = "ffi-support")]
pub(crate) fn borrow_sender_owned_with_retry(&self, budget: Duration) -> Result<OwnedSender> {
let deadline = Instant::now().checked_add(budget);
let cs = reconnect_pick(&self.inner, deadline, pick_sfa_sender)?;
Ok(OwnedSender::new(Arc::clone(&self.inner), cs))
}
#[cfg(feature = "ffi-support")]
pub(crate) fn borrow_direct_column_sender_owned(&self) -> Result<OwnedDirectColumnSender> {
let cs = pick_direct_sender(&self.inner)?;
Ok(OwnedDirectColumnSender::new(Arc::clone(&self.inner), cs))
}
#[cfg(feature = "ffi-support")]
pub(crate) fn borrow_direct_column_sender_owned_with_retry(
&self,
budget: Duration,
) -> Result<OwnedDirectColumnSender> {
let deadline = Instant::now().checked_add(budget);
let cs = reconnect_pick(&self.inner, deadline, pick_direct_sender)?;
Ok(OwnedDirectColumnSender::new(Arc::clone(&self.inner), cs))
}
fn pick_sender(&self) -> Result<PooledSender<PooledSenderCore>> {
pick_sfa_sender(&self.inner)
}
fn pick_replacement_sender(&self) -> Result<PooledSender<DirectSenderCore>> {
if self.inner.shutdown.load(Ordering::SeqCst) {
return Err(error::fmt!(
InvalidApiCall,
"QuestDb pool is closed; cannot replace sender"
));
}
if let Some(entry) = lock_state(&self.inner.direct_state).free.pop() {
return Ok(PooledSender {
sender: entry.sender,
slot_index: entry.slot_index,
});
}
let conn = connect_conn_pool(&self.inner)?;
Ok(PooledSender {
sender: DirectSenderCore::new(
conn,
crate::ingress::SymbolGlobalDict::new(),
crate::ingress::column_sender::encoder::EncodeScratch::new(),
false,
),
slot_index: None,
})
}
pub fn reap_idle(&self) -> usize {
reap_idle_inner(&self.inner)
}
pub fn connection_events_dropped(&self) -> u64 {
self.inner.conn_events.dropped()
}
pub fn connection_events_delivered(&self) -> u64 {
self.inner.conn_events.delivered()
}
pub fn rejection_events_delivered(&self) -> u64 {
self.inner.rejections.delivered()
}
pub fn rejection_events_dropped(&self) -> u64 {
self.inner.rejections.dropped()
}
#[doc(hidden)]
pub fn dbg_pool_counts(&self) -> DbgPoolCounts {
let ingress = {
let s = lock_state(&self.inner.state);
DbgPoolCount {
free: s.free.len(),
in_use: s.in_use,
closing: s.closing,
}
};
let column_direct = {
let s = lock_state(&self.inner.direct_state);
DbgPoolCount {
free: s.free.len(),
in_use: s.in_use,
closing: s.closing,
}
};
#[cfg(feature = "_egress")]
let reader = {
let s = lock_reader_state(&self.inner.reader_state);
DbgPoolCount {
free: s.free.len(),
in_use: s.in_use,
closing: 0,
}
};
#[cfg(not(feature = "_egress"))]
let reader = DbgPoolCount::default();
DbgPoolCounts {
ingress,
column_direct,
reader,
}
}
pub fn close(self) {
drop(self);
}
#[cfg(any(feature = "polars-ingress", feature = "polars-egress"))]
pub(crate) fn reconnect_policy(&self) -> crate::ingress::ReconnectPolicy {
self.inner.connector.reconnect_policy()
}
#[cfg(feature = "ffi-support")]
pub(crate) fn reconnect_max_duration(&self) -> Duration {
self.inner.connector.reconnect_policy().max_duration()
}
#[cfg(test)]
pub(crate) fn free_count(&self) -> usize {
lock_state(&self.inner.state).free.len()
}
#[cfg(test)]
pub(crate) fn in_use_count(&self) -> usize {
lock_state(&self.inner.state).in_use
}
#[cfg(all(test, feature = "ffi-support"))]
pub(crate) fn closing_count(&self) -> usize {
lock_state(&self.inner.state).closing
}
#[cfg(test)]
pub(crate) fn direct_free_count(&self) -> usize {
lock_state(&self.inner.direct_state).free.len()
}
#[cfg(test)]
pub(crate) fn direct_in_use_count(&self) -> usize {
lock_state(&self.inner.direct_state).in_use
}
#[cfg(feature = "_egress")]
pub fn borrow_reader(&self) -> crate::error::Result<BorrowedReader<'_>> {
let reader = self.pick_reader()?;
Ok(BorrowedReader::new(self, reader))
}
#[cfg(all(feature = "_egress", feature = "ffi-support"))]
pub(crate) fn borrow_reader_owned(&self) -> crate::error::Result<OwnedReader> {
let reader = self.pick_reader()?;
Ok(OwnedReader {
inner: Arc::clone(&self.inner),
reader: Some(reader),
must_close: false,
})
}
#[cfg(all(feature = "_egress", feature = "ffi-support"))]
pub(crate) fn reader_pool_handle(&self) -> ReaderPoolHandle {
ReaderPoolHandle {
inner: Arc::clone(&self.inner),
}
}
#[cfg(feature = "_egress")]
fn pick_reader(&self) -> crate::error::Result<Reader> {
use crate::{Error, ErrorCode};
let slot = {
let mut state = lock_reader_state(&self.inner.reader_state);
let mut acquire_deadline = None;
loop {
if self.inner.shutdown.load(Ordering::SeqCst) {
return Err(Error::new(
ErrorCode::InvalidApiCall,
"QuestDb pool is closed; cannot borrow reader",
));
}
if let Some(entry) = state.free.pop() {
state.in_use += 1;
drop(state);
return Ok(entry.reader);
}
if state.total() < self.inner.query_pool_max {
break;
}
if let Some(wait_for) =
remaining_wait(&mut acquire_deadline, self.inner.acquire_timeout)
{
let (next_state, _) = match self.inner.reader_cv.wait_timeout(state, wait_for) {
Ok((guard, result)) => (guard, result),
Err(poisoned) => poisoned.into_inner(),
};
state = next_state;
continue;
}
return Err(Error::new(
ErrorCode::InvalidApiCall,
format!(
"Reader pool exhausted: {} readers are currently borrowed at \
the `query_pool_max` cap of {} after waiting \
acquire_timeout_ms={}. Release a reader, or raise \
`query_pool_max` / `acquire_timeout_ms`.",
state.in_use,
self.inner.query_pool_max,
self.inner.acquire_timeout.as_millis()
),
));
}
state.in_use += 1;
ReaderInUseSlot {
inner: &self.inner,
armed: true,
}
};
let reader = Reader::from_conf(&self.inner.conf)?;
slot.commit();
Ok(reader)
}
#[cfg(all(feature = "_egress", any(test, feature = "ffi-support")))]
pub(crate) fn reader_free_count(&self) -> usize {
lock_reader_state(&self.inner.reader_state).free.len()
}
#[cfg(all(feature = "_egress", any(test, feature = "ffi-support")))]
pub(crate) fn reader_in_use_count(&self) -> usize {
lock_reader_state(&self.inner.reader_state).in_use
}
}
impl Debug for QuestDb {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
let state = lock_state(&self.inner.state);
let mut s = f.debug_struct("QuestDb");
s.field("sender_pool_min", &self.inner.sender_pool_min)
.field("sender_pool_max", &self.inner.sender_pool_max);
#[cfg(feature = "_egress")]
s.field("query_pool_min", &self.inner.query_pool_min)
.field("query_pool_max", &self.inner.query_pool_max);
s.field("acquire_timeout", &self.inner.acquire_timeout)
.field("free", &state.free.len())
.field("in_use", &state.in_use)
.finish()
}
}
impl Drop for QuestDb {
fn drop(&mut self) {
self.inner.shutdown.store(true, Ordering::SeqCst);
{
let _g = lock_state(&self.inner.state);
self.inner.cv.notify_all();
}
{
let _g = lock_state(&self.inner.direct_state);
self.inner.direct_cv.notify_all();
}
#[cfg(feature = "_egress")]
{
let _g = lock_reader_state(&self.inner.reader_state);
self.inner.reader_cv.notify_all();
}
if let Some(handle) = self.reaper.take() {
let _ = handle.join();
}
drain_idle_inner(&self.inner);
self.inner.conn_events.close();
self.inner.rejections.close();
}
}
struct SenderHandle<'a> {
db: &'a QuestDb,
sender: Option<PooledSenderCore>,
slot_index: Option<usize>,
_not_send: PhantomData<Rc<()>>,
}
impl<'a> SenderHandle<'a> {
fn new(db: &'a QuestDb, sender: PooledSender<PooledSenderCore>) -> Self {
Self {
db,
sender: Some(sender.sender),
slot_index: sender.slot_index,
_not_send: PhantomData,
}
}
fn inner_mut(&mut self) -> &mut PooledSenderCore {
self.sender
.as_mut()
.expect("borrowed sender already returned")
}
fn inner_ref(&self) -> &PooledSenderCore {
self.sender
.as_ref()
.expect("borrowed sender already returned")
}
}
struct DirectSenderHandle<'a> {
db: &'a QuestDb,
sender: Option<DirectSenderCore>,
_not_send: PhantomData<Rc<()>>,
}
impl<'a> DirectSenderHandle<'a> {
fn new(db: &'a QuestDb, sender: PooledSender<DirectSenderCore>) -> Self {
debug_assert!(sender.slot_index.is_none());
Self {
db,
sender: Some(sender.sender),
_not_send: PhantomData,
}
}
fn inner_mut(&mut self) -> &mut DirectSenderCore {
self.sender
.as_mut()
.expect("borrowed direct sender already returned")
}
#[cfg(test)]
fn inner_ref(&self) -> &DirectSenderCore {
self.sender
.as_ref()
.expect("borrowed direct sender already returned")
}
#[cfg(any(feature = "polars-ingress", feature = "polars-egress"))]
pub(crate) fn reconnect_policy(&self) -> crate::ingress::ReconnectPolicy {
self.db.reconnect_policy()
}
pub fn reborrow_from_pool(&mut self) -> Result<()> {
if let Some(sender) = self.sender.as_mut() {
if sender.in_flight() == 0 && !sender.must_close() && !sender.transport_dead() {
return Ok(());
}
if sender.in_flight() > 0 {
log::warn!(
"direct sender failover dropped a connection with un-sync'd \
deferred frame(s); their data is discarded. Re-drive the source \
from the last successful sync(), not from the failing chunk."
);
sender.mark_must_close();
}
record_sender_transport_failure(&self.db.inner, sender);
}
let fresh = self.db.pick_replacement_sender()?;
debug_assert!(fresh.slot_index.is_none());
if let Some(old) = self.sender.replace(fresh.sender) {
finish_replaced_sender(&self.db.inner, old);
}
Ok(())
}
#[cfg(any(feature = "polars-ingress", feature = "polars-egress"))]
pub(crate) fn reborrow_with_retry(&mut self, deadline: Option<Instant>) -> Result<()> {
let policy = self.reconnect_policy();
let mut backoff = policy.initial_backoff();
loop {
match self.reborrow_from_pool() {
Ok(()) => return Ok(()),
Err(e)
if reconnect_error_is_terminal(&e) || reconnect_deadline_expired(deadline) =>
{
return Err(e);
}
Err(e) => {
let (sleep_for, next) = reconnect_backoff_step(
&e,
policy.initial_backoff(),
policy.max_backoff(),
backoff,
);
sleep_until_deadline(sleep_for, deadline);
backoff = next;
}
}
}
}
}
impl Debug for SenderHandle<'_> {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
f.debug_struct("SenderHandle")
.field("sender", &self.sender)
.finish()
}
}
impl Debug for DirectSenderHandle<'_> {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
f.debug_struct("DirectSenderHandle")
.field("sender", &self.sender)
.finish()
}
}
pub struct BorrowedSender<'a>(SenderHandle<'a>);
impl<'a> BorrowedSender<'a> {
#[cfg(test)]
pub(crate) fn effective_frame_cap_for_test(&self) -> (usize, bool) {
self.0.inner_ref().effective_frame_cap()
}
pub fn new_buffer(&self) -> Buffer {
self.0.db.new_buffer()
}
pub fn flush(&mut self, chunk: &mut crate::ingress::column_sender::Chunk<'_>) -> Result<()> {
self.0.inner_mut().flush(chunk)
}
pub fn flush_buffer(&mut self, buffer: &mut Buffer) -> Result<()> {
self.0.inner_mut().flush_buffer(buffer)
}
pub fn flush_buffer_and_keep(&mut self, buffer: &Buffer) -> Result<()> {
self.0.inner_mut().flush_buffer_and_keep(buffer)
}
pub fn flush_buffer_and_get_fsn(&mut self, buffer: &mut Buffer) -> Result<Option<u64>> {
self.0.inner_mut().flush_buffer_and_get_fsn(buffer)
}
pub fn flush_buffer_and_keep_and_get_fsn(&mut self, buffer: &Buffer) -> Result<Option<u64>> {
self.0.inner_mut().flush_buffer_and_keep_and_get_fsn(buffer)
}
pub fn flush_buffer_and_wait(
&mut self,
buffer: &mut Buffer,
ack_level: AckLevel,
) -> Result<()> {
self.0.inner_mut().flush_buffer_and_wait(buffer, ack_level)
}
pub fn flush_and_wait(
&mut self,
chunk: &mut crate::ingress::column_sender::Chunk<'_>,
ack_level: AckLevel,
) -> Result<()> {
self.0.inner_mut().flush_and_wait(chunk, ack_level)
}
pub fn flush_and_get_fsn(
&mut self,
chunk: &mut crate::ingress::column_sender::Chunk<'_>,
) -> Result<Option<u64>> {
self.0.inner_mut().flush_and_get_fsn(chunk)
}
pub fn published_fsn(&self) -> Result<Option<u64>> {
self.0.inner_ref().published_fsn()
}
pub fn acked_fsn(&self) -> Result<Option<u64>> {
self.0.inner_ref().acked_fsn()
}
pub fn wait(&mut self, ack_level: AckLevel, timeout: Duration) -> Result<()> {
self.0.inner_mut().wait(ack_level, timeout)
}
pub fn drop_on_return(&mut self) {
self.0.inner_mut().mark_must_close()
}
#[cfg(test)]
pub(crate) fn must_close_for_test(&self) -> bool {
self.0.inner_ref().must_close()
}
#[cfg(test)]
pub(crate) fn is_store_and_forward(&self) -> bool {
true
}
#[cfg(test)]
pub(crate) fn in_flight(&self) -> u32 {
0
}
#[cfg(feature = "arrow-ingress")]
pub fn flush_arrow_batch_at_now<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_now(table, batch, overrides)
}
#[cfg(feature = "arrow-ingress")]
pub fn flush_arrow_batch_at_now_and_wait<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
ack_level: AckLevel,
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_now_and_wait(table, batch, overrides, ack_level)
}
#[cfg(feature = "arrow-ingress")]
pub fn flush_arrow_batch_at_now_and_get_fsn<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
) -> Result<Option<u64>>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_now_and_get_fsn(table, batch, overrides)
}
#[cfg(feature = "arrow-ingress")]
pub fn flush_arrow_batch_at_column<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
ts_column: crate::ingress::ColumnName<'_>,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_column(table, batch, ts_column, overrides)
}
#[cfg(feature = "arrow-ingress")]
pub fn flush_arrow_batch_at_column_and_wait<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
ts_column: crate::ingress::ColumnName<'_>,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
ack_level: AckLevel,
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_column_and_wait(table, batch, ts_column, overrides, ack_level)
}
#[cfg(feature = "arrow-ingress")]
pub fn flush_arrow_batch_at_column_and_get_fsn<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
ts_column: crate::ingress::ColumnName<'_>,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
) -> Result<Option<u64>>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_column_and_get_fsn(table, batch, ts_column, overrides)
}
}
impl Debug for BorrowedSender<'_> {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
f.debug_tuple("BorrowedSender").field(&self.0).finish()
}
}
pub struct BorrowedDirectColumnSender<'a>(DirectSenderHandle<'a>);
impl<'a> BorrowedDirectColumnSender<'a> {
pub fn flush(&mut self, chunk: &mut crate::ingress::column_sender::Chunk<'_>) -> Result<()> {
self.0.inner_mut().flush(chunk)
}
pub fn flush_and_wait(
&mut self,
chunk: &mut crate::ingress::column_sender::Chunk<'_>,
ack_level: AckLevel,
) -> Result<()> {
self.0.inner_mut().flush_and_wait(chunk, ack_level)
}
pub fn commit(&mut self, ack_level: AckLevel) -> Result<()> {
self.0.inner_mut().sync(ack_level)
}
pub fn reborrow_from_pool(&mut self) -> Result<()> {
self.0.reborrow_from_pool()
}
#[cfg(any(feature = "polars-ingress", feature = "polars-egress"))]
pub(crate) fn reborrow_with_retry(&mut self, deadline: Option<Instant>) -> Result<()> {
self.0.reborrow_with_retry(deadline)
}
#[cfg(any(feature = "polars-ingress", feature = "polars-egress"))]
pub(crate) fn reconnect_policy(&self) -> crate::ingress::ReconnectPolicy {
self.0.reconnect_policy()
}
#[cfg(feature = "polars-ingress")]
pub(crate) fn default_ack_level(&self) -> AckLevel {
self.0.db.default_ack_level()
}
pub fn drop_on_return(&mut self) {
self.0.inner_mut().mark_must_close()
}
#[cfg(test)]
pub(crate) fn must_close_for_test(&self) -> bool {
self.0.inner_ref().must_close()
}
#[cfg(test)]
pub(crate) fn is_store_and_forward(&self) -> bool {
false
}
#[cfg(test)]
pub(crate) fn in_flight(&self) -> u32 {
self.0.inner_ref().in_flight()
}
#[cfg(feature = "polars-ingress")]
pub(crate) fn flush_arrow_batch_at_now<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_now(table, batch, overrides)
}
#[cfg(feature = "polars-ingress")]
pub(crate) fn flush_arrow_batch_at_column<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
ts_column: crate::ingress::ColumnName<'_>,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_column(table, batch, ts_column, overrides)
}
#[cfg(feature = "arrow-ingress")]
pub(crate) fn flush_arrow_batch_at_now_and_wait<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
ack_level: AckLevel,
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_now_and_wait(table, batch, overrides, ack_level)
}
#[cfg(feature = "arrow-ingress")]
pub(crate) fn flush_arrow_batch_at_column_and_wait<'t, T>(
&mut self,
table: T,
batch: &arrow::array::RecordBatch,
ts_column: crate::ingress::ColumnName<'_>,
overrides: &[crate::ingress::column_sender::ArrowColumnOverride<'_>],
ack_level: AckLevel,
) -> Result<()>
where
T: TryInto<crate::ingress::TableName<'t>>,
crate::Error: From<T::Error>,
{
self.0
.inner_mut()
.flush_arrow_batch_at_column_and_wait(table, batch, ts_column, overrides, ack_level)
}
}
impl Debug for BorrowedDirectColumnSender<'_> {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
f.debug_tuple("BorrowedDirectColumnSender")
.field(&self.0)
.finish()
}
}
impl Drop for SenderHandle<'_> {
fn drop(&mut self) {
let Some(sender) = self.sender.take() else {
return;
};
return_sfa_to_pool(&self.db.inner, sender, self.slot_index);
}
}
impl Drop for DirectSenderHandle<'_> {
fn drop(&mut self) {
let Some(mut sender) = self.sender.take() else {
return;
};
commit_in_flight_on_drop(self.db.inner.connector.request_durable_ack(), &mut sender);
return_direct_to_pool(&self.db.inner, sender);
}
}
#[cfg(feature = "ffi-support")]
pub struct OwnedSender {
inner: Arc<DbInner>,
sender: Option<PooledSenderCore>,
slot_index: Option<usize>,
}
#[cfg(feature = "ffi-support")]
impl OwnedSender {
fn new(inner: Arc<DbInner>, sender: PooledSender<PooledSenderCore>) -> Self {
Self {
inner,
sender: Some(sender.sender),
slot_index: sender.slot_index,
}
}
pub fn get_mut(&mut self) -> &mut PooledSenderCore {
self.sender
.as_mut()
.expect("OwnedSender already returned to the pool")
}
pub fn get(&self) -> &PooledSenderCore {
self.sender
.as_ref()
.expect("OwnedSender already returned to the pool")
}
pub fn pool_closed(&self) -> bool {
self.inner.shutdown.load(Ordering::SeqCst)
}
pub fn mark_must_close(&mut self) {
self.get_mut().mark_must_close();
}
pub fn must_close(&self) -> bool {
self.pool_closed() || self.get().must_close()
}
}
#[cfg(feature = "ffi-support")]
impl Drop for OwnedSender {
fn drop(&mut self) {
if let Some(sender) = self.sender.take() {
return_sfa_to_pool(&self.inner, sender, self.slot_index);
}
}
}
#[cfg(feature = "ffi-support")]
enum DirectBacking {
Pool(Arc<DbInner>),
Standalone { request_durable_ack: bool },
}
#[cfg(feature = "ffi-support")]
pub struct OwnedDirectColumnSender {
backing: DirectBacking,
sender: Option<DirectSenderCore>,
}
#[cfg(feature = "ffi-support")]
impl OwnedDirectColumnSender {
fn new(inner: Arc<DbInner>, sender: PooledSender<DirectSenderCore>) -> Self {
debug_assert!(sender.slot_index.is_none());
Self {
backing: DirectBacking::Pool(inner),
sender: Some(sender.sender),
}
}
pub fn from_conf(conf: &str) -> Result<Self> {
Self::from_builder(&SenderBuilder::from_conf(conf)?)
}
pub fn from_builder(builder: &SenderBuilder) -> Result<Self> {
let connector = builder.build_qwp_ws_connector()?;
let health = Mutex::new(QwpWsHostHealthTracker::new(connector.endpoint_count()));
let raw = connector.connect_round_pooled(&health, None)?;
let conn = ColumnConn::from_round_stream(raw)?;
let sender = DirectSenderCore::new(
conn,
crate::ingress::SymbolGlobalDict::new(),
crate::ingress::column_sender::encoder::EncodeScratch::new(),
false,
);
Ok(Self {
backing: DirectBacking::Standalone {
request_durable_ack: connector.request_durable_ack(),
},
sender: Some(sender),
})
}
pub fn get_mut(&mut self) -> &mut DirectSenderCore {
self.sender
.as_mut()
.expect("OwnedDirectColumnSender already released")
}
pub fn get(&self) -> &DirectSenderCore {
self.sender
.as_ref()
.expect("OwnedDirectColumnSender already released")
}
pub fn pool_closed(&self) -> bool {
match &self.backing {
DirectBacking::Pool(inner) => inner.shutdown.load(Ordering::SeqCst),
DirectBacking::Standalone { .. } => false,
}
}
pub fn mark_must_close(&mut self) {
self.get_mut().mark_must_close();
}
pub fn must_close(&self) -> bool {
self.pool_closed() || self.get().must_close()
}
}
#[cfg(feature = "ffi-support")]
impl Drop for OwnedDirectColumnSender {
fn drop(&mut self) {
let Some(mut sender) = self.sender.take() else {
return;
};
match &self.backing {
DirectBacking::Pool(inner) => {
commit_in_flight_on_drop(inner.connector.request_durable_ack(), &mut sender);
return_direct_to_pool(inner, sender);
}
DirectBacking::Standalone {
request_durable_ack,
} => {
commit_in_flight_on_drop(*request_durable_ack, &mut sender);
}
}
}
}
#[cfg(feature = "_egress")]
pub struct BorrowedReader<'a> {
db: &'a QuestDb,
reader: Option<Reader>,
must_close: bool,
_not_send: PhantomData<Rc<()>>,
}
#[cfg(feature = "_egress")]
impl<'a> BorrowedReader<'a> {
fn new(db: &'a QuestDb, reader: Reader) -> Self {
Self {
db,
reader: Some(reader),
must_close: false,
_not_send: PhantomData,
}
}
pub fn drop_on_return(&mut self) {
self.must_close = true;
}
}
#[cfg(feature = "_egress")]
impl Debug for BorrowedReader<'_> {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
f.debug_struct("BorrowedReader")
.field("borrowed", &self.reader.is_some())
.field("must_close", &self.must_close)
.finish()
}
}
#[cfg(feature = "_egress")]
impl Deref for BorrowedReader<'_> {
type Target = Reader;
fn deref(&self) -> &Self::Target {
self.reader
.as_ref()
.expect("borrowed reader already returned")
}
}
#[cfg(feature = "_egress")]
impl DerefMut for BorrowedReader<'_> {
fn deref_mut(&mut self) -> &mut Self::Target {
self.reader
.as_mut()
.expect("borrowed reader already returned")
}
}
#[cfg(feature = "_egress")]
impl Drop for BorrowedReader<'_> {
fn drop(&mut self) {
if let Some(reader) = self.reader.take() {
return_reader_to_pool(&self.db.inner, reader, self.must_close);
}
}
}
#[cfg(all(feature = "_egress", feature = "ffi-support"))]
pub struct OwnedReader {
inner: Arc<DbInner>,
reader: Option<Reader>,
must_close: bool,
}
#[cfg(all(feature = "_egress", feature = "ffi-support"))]
impl OwnedReader {
pub fn get(&self) -> &Reader {
self.reader
.as_ref()
.expect("OwnedReader already returned to the pool")
}
pub fn get_mut(&mut self) -> &mut Reader {
self.reader
.as_mut()
.expect("OwnedReader already returned to the pool")
}
pub fn mark_must_close(&mut self) {
self.must_close = true;
}
pub fn take(mut self) -> Option<Reader> {
self.reader.take()
}
}
#[cfg(all(feature = "_egress", feature = "ffi-support"))]
impl Drop for OwnedReader {
fn drop(&mut self) {
if let Some(reader) = self.reader.take() {
return_reader_to_pool(&self.inner, reader, self.must_close);
}
}
}
#[cfg(all(feature = "_egress", feature = "ffi-support"))]
#[derive(Clone)]
pub struct ReaderPoolHandle {
inner: Arc<DbInner>,
}
#[cfg(all(feature = "_egress", feature = "ffi-support"))]
impl ReaderPoolHandle {
pub fn return_reader(&self, reader: Reader, must_close: bool) {
return_reader_to_pool(&self.inner, reader, must_close);
}
pub fn pool_closed(&self) -> bool {
self.inner.shutdown.load(Ordering::SeqCst)
}
pub fn release_leaked_slot(&self) {
let mut state = lock_reader_state(&self.inner.reader_state);
state.in_use = state.in_use.saturating_sub(1);
}
}
#[cfg(feature = "_egress")]
fn return_reader_to_pool(inner: &Arc<DbInner>, reader: Reader, must_close: bool) {
let must_close = must_close || reader.transport_torn_down();
let mut state = lock_reader_state(&inner.reader_state);
state.in_use = state.in_use.saturating_sub(1);
if !must_close && !inner.shutdown.load(Ordering::SeqCst) {
state.free.push(ReaderPoolEntry {
reader,
last_idle_at: Instant::now(),
});
}
drop(state);
inner.reader_cv.notify_all();
}
trait PoolableSender {
fn is_stale(&self) -> bool;
fn drain_for_retire(&mut self, inner: &DbInner);
}
impl PoolableSender for PooledSenderCore {
fn is_stale(&self) -> bool {
self.must_close()
}
fn drain_for_retire(&mut self, inner: &DbInner) {
drain_sfa_before_drop(inner, self);
}
}
impl PoolableSender for DirectSenderCore {
fn is_stale(&self) -> bool {
self.must_close()
}
fn drain_for_retire(&mut self, _inner: &DbInner) {}
}
fn retire_stale_entry<S: PoolableSender>(inner: &Arc<DbInner>, entry: PoolEntry<S>) {
let _release = entry.slot_index.is_some().then_some(SenderSlotRelease {
inner: inner.as_ref(),
slot_index: entry.slot_index,
decrement_in_use: false,
decrement_closing: true,
});
let mut sender = entry.sender;
sender.drain_for_retire(inner);
drop(sender);
}
fn pick_sender_inner<S: PoolableSender>(
inner: &Arc<DbInner>,
pool: &Mutex<PoolState<S>>,
cv: &Condvar,
sfa: bool,
connect: impl FnOnce(Option<usize>) -> Result<S>,
) -> Result<PooledSender<S>> {
let slot = {
let mut state = lock_state(pool);
let mut close_wait_deadline = None;
let mut acquire_deadline = None;
loop {
if inner.shutdown.load(Ordering::SeqCst) {
return Err(error::fmt!(
InvalidApiCall,
"QuestDb pool is closed; cannot borrow sender"
));
}
if let Some(entry) = state.free.pop() {
if entry.sender.is_stale() {
if entry.slot_index.is_some() {
state.closing += 1;
}
drop(state);
retire_stale_entry(inner, entry);
state = lock_state(pool);
continue;
}
state.in_use += 1;
drop(state);
return Ok(PooledSender {
sender: entry.sender,
slot_index: entry.slot_index,
});
}
if state.reserved_total() < inner.sender_pool_max {
break;
}
let wait_timeout = inner.connector.close_flush_timeout();
if sfa
&& inner.sf_disk
&& state.closing > 0
&& let Some(wait_for) = remaining_wait(&mut close_wait_deadline, wait_timeout)
{
let (next_state, _) = match cv.wait_timeout(state, wait_for) {
Ok((guard, result)) => (guard, result),
Err(poisoned) => poisoned.into_inner(),
};
state = next_state;
continue;
}
if let Some(wait_for) = remaining_wait(&mut acquire_deadline, inner.acquire_timeout) {
let (next_state, _) = match cv.wait_timeout(state, wait_for) {
Ok((guard, result)) => (guard, result),
Err(poisoned) => poisoned.into_inner(),
};
state = next_state;
continue;
}
return Err(error::fmt!(
InvalidApiCall,
"Connection pool exhausted: {} sender(s) in use at the \
sender_pool_max cap of {} after waiting acquire_timeout_ms={}. \
Drop a borrowed sender, or raise sender_pool_max / \
acquire_timeout_ms.",
state.in_use,
inner.sender_pool_max,
inner.acquire_timeout.as_millis()
));
}
let slot_index = state.allocate_slot_index();
debug_assert_eq!(slot_index.is_some(), sfa && inner.sf_disk);
state.in_use += 1;
InUseSlot {
state: pool,
cv,
slot_index,
armed: true,
}
};
let sender = connect(slot.slot_index)?;
let slot_index = slot.slot_index;
slot.commit();
Ok(PooledSender { sender, slot_index })
}
fn pick_sfa_sender(inner: &Arc<DbInner>) -> Result<PooledSender<PooledSenderCore>> {
let mut picked = pick_sender_inner(inner, &inner.state, &inner.cv, true, |slot_index| {
connect_sfa_pool(inner, slot_index)
})?;
picked.sender.rebase_lease_observation();
Ok(picked)
}
fn pick_direct_sender(inner: &Arc<DbInner>) -> Result<PooledSender<DirectSenderCore>> {
pick_sender_inner(
inner,
&inner.direct_state,
&inner.direct_cv,
false,
|_slot_index| {
let conn = connect_conn_pool(inner)?;
Ok(DirectSenderCore::new(
conn,
crate::ingress::SymbolGlobalDict::new(),
crate::ingress::column_sender::encoder::EncodeScratch::new(),
false,
))
},
)
}
fn prewarm_min_connections(db: &QuestDb) -> Result<()> {
let inner = &db.inner;
let mut warm = Vec::new();
let mut outcome: Result<()> = Ok(());
for _ in 0..inner.sender_pool_min {
match pick_sfa_sender(inner) {
Ok(sender) => warm.push(sender),
Err(err) => {
outcome = Err(err);
break;
}
}
}
for picked in warm {
return_sfa_to_pool(inner, picked.sender, picked.slot_index);
}
outcome?;
#[cfg(feature = "_egress")]
{
let mut warm = Vec::new();
let mut outcome: Result<()> = Ok(());
for _ in 0..inner.query_pool_min {
match db.pick_reader() {
Ok(reader) => warm.push(reader),
Err(err) => {
outcome = Err(err);
break;
}
}
}
for reader in warm {
return_reader_to_pool(inner, reader, false);
}
outcome?;
}
Ok(())
}
fn connect_sfa_pool(inner: &Arc<DbInner>, slot_index: Option<usize>) -> Result<PooledSenderCore> {
connect_sfa_pool_with_recovery_candidates(
inner,
slot_index,
&inner.out_of_range_recovery_candidates,
false,
)
}
fn connect_sfa_pool_with_recovery_candidates(
inner: &Arc<DbInner>,
slot_index: Option<usize>,
recovery_candidates: &[PathBuf],
force_async_initial_connect: bool,
) -> Result<PooledSenderCore> {
let sender_id = slot_index.map(|index| managed_slot_id(&inner.slot_base_id, index));
let state = inner
.connector
.connect_sfa_background_with_pool_slot(
sender_id.as_deref(),
inner.managed_slot_exclusion.as_slice(),
recovery_candidates,
Arc::clone(&inner.conn_events),
Arc::clone(&inner.rejections),
force_async_initial_connect,
)
.map_err(|err| {
crate::Error::new(
err.code(),
format!("Failed to open store-and-forward sender: {}", err.msg()),
)
})?;
PooledSenderCore::new_store_and_forward(
state,
inner.connector.max_buf_size(),
inner.connector.request_durable_ack(),
inner.connector.request_timeout(),
)
}
#[cfg(feature = "ffi-support")]
fn reconnect_pick<S>(
inner: &Arc<DbInner>,
deadline: Option<Instant>,
mut pick: impl FnMut(&Arc<DbInner>) -> Result<PooledSender<S>>,
) -> Result<PooledSender<S>> {
let policy = inner.connector.reconnect_policy();
let mut backoff = policy.initial_backoff();
loop {
match pick(inner) {
Ok(cs) => return Ok(cs),
Err(e) if reconnect_error_is_terminal(&e) || reconnect_deadline_expired(deadline) => {
return Err(e);
}
Err(e) => {
let (sleep_for, next) = reconnect_backoff_step(
&e,
policy.initial_backoff(),
policy.max_backoff(),
backoff,
);
sleep_until_deadline(sleep_for, deadline);
backoff = next;
}
}
}
}
#[cfg(any(
feature = "polars-ingress",
feature = "polars-egress",
feature = "ffi-support"
))]
pub(crate) fn reconnect_deadline_expired(deadline: Option<Instant>) -> bool {
deadline.is_some_and(|d| Instant::now() >= d)
}
fn remaining_wait(deadline: &mut Option<Instant>, timeout: Duration) -> Option<Duration> {
if timeout.is_zero() {
return None;
}
let now = Instant::now();
let deadline = deadline.get_or_insert_with(|| now.checked_add(timeout).unwrap_or(now));
let remaining = deadline.saturating_duration_since(now);
if remaining.is_zero() {
None
} else {
Some(remaining)
}
}
#[cfg(any(
feature = "polars-ingress",
feature = "polars-egress",
feature = "ffi-support"
))]
fn sleep_until_deadline(sleep_for: Duration, deadline: Option<Instant>) {
let d = match deadline {
Some(dl) => sleep_for.min(dl.saturating_duration_since(Instant::now())),
None => sleep_for,
};
if !d.is_zero() {
thread::sleep(d);
}
}
fn connect_conn_pool(inner: &Arc<DbInner>) -> Result<ColumnConn> {
let raw: RawQwpWsRoundStream = inner
.connector
.connect_round_pooled(&inner.health, Some(inner.conn_events.as_ref()))?;
ColumnConn::from_round_stream(raw)
}
fn commit_in_flight_on_drop(request_durable_ack: bool, sender: &mut DirectSenderCore) {
if sender.in_flight() == 0 {
return;
}
let ack = if request_durable_ack {
AckLevel::Durable
} else {
AckLevel::Ok
};
let committed = sender.can_drain_in_flight() && sender.sync(ack).is_ok();
if !committed {
log::warn!(
"direct sender dropped with un-sync'd deferred frame(s) that could \
not be committed; their data is discarded. Call sync() (or \
flush_and_wait() on the final chunk) before the handle is dropped."
);
sender.mark_must_close();
}
}
fn drain_sfa_before_drop(inner: &DbInner, sender: &mut PooledSenderCore) {
let timeout = inner.connector.close_flush_timeout();
if timeout.is_zero() {
return;
}
let durable = inner.connector.request_durable_ack();
if sender.sfa_fully_delivered(durable) {
return;
}
sender.begin_close();
if let Err(err) = sender.drain_to_deadline(Instant::now().checked_add(timeout)) {
log::warn!(
"store-and-forward sender dropped with frame(s) that could \
not be delivered within close_flush_timeout; their data is \
discarded. Call wait() before closing the pool, or set sf_dir for \
crash-durable persistence. Cause: {err}"
);
}
}
fn drain_sfa_senders_bounded(inner: &DbInner, senders: &mut [PooledSenderCore]) {
let timeout = inner.connector.close_flush_timeout();
if timeout.is_zero() || senders.is_empty() {
return;
}
let durable = inner.connector.request_durable_ack();
for sender in senders.iter() {
sender.begin_close();
}
let deadline = Instant::now().checked_add(timeout);
for sender in senders.iter_mut() {
if sender.sfa_fully_delivered(durable) {
continue;
}
if let Err(err) = sender.drain_to_deadline(deadline) {
log::warn!(
"store-and-forward sender dropped on pool close with \
frame(s) that could not be delivered within close_flush_timeout; \
their data is discarded. Call wait() before close, or set \
sf_dir for crash-durable persistence. Cause: {err}"
);
}
}
}
fn return_sfa_to_pool(
inner: &Arc<DbInner>,
mut sender: PooledSenderCore,
slot_index: Option<usize>,
) {
let must_close = sender.must_close();
let release_slot;
{
let mut state = lock_state(&inner.state);
if !must_close && !inner.shutdown.load(Ordering::SeqCst) {
state.in_use = state.in_use.saturating_sub(1);
state.free.push(PoolEntry {
sender,
slot_index,
last_idle_at: Instant::now(),
});
inner.cv.notify_all();
return;
}
release_slot = slot_index.is_some();
if release_slot {
state.closing += 1;
} else {
state.in_use = state.in_use.saturating_sub(1);
}
}
let _release = release_slot.then_some(SenderSlotRelease {
inner: inner.as_ref(),
slot_index,
decrement_in_use: true,
decrement_closing: true,
});
drain_sfa_before_drop(inner, &mut sender);
drop(sender);
}
fn return_direct_to_pool(inner: &Arc<DbInner>, sender: DirectSenderCore) {
let must_close = sender.must_close();
record_sender_transport_failure(inner, &sender);
{
let mut state = lock_state(&inner.direct_state);
state.in_use = state.in_use.saturating_sub(1);
if !must_close && !inner.shutdown.load(Ordering::SeqCst) {
state.free.push(PoolEntry {
sender,
slot_index: None,
last_idle_at: Instant::now(),
});
}
}
inner.direct_cv.notify_all();
}
fn finish_replaced_sender(inner: &Arc<DbInner>, sender: DirectSenderCore) {
let must_close = sender.must_close();
record_sender_transport_failure(inner, &sender);
{
let mut state = lock_state(&inner.direct_state);
if !must_close
&& !inner.shutdown.load(Ordering::SeqCst)
&& state.total() < inner.sender_pool_max
{
state.free.push(PoolEntry {
sender,
slot_index: None,
last_idle_at: Instant::now(),
});
}
}
inner.direct_cv.notify_all();
}
fn record_sender_transport_failure(inner: &Arc<DbInner>, sender: &DirectSenderCore) {
if sender.transport_dead() {
let idx = sender.endpoint_idx();
lock_health(&inner.health)
.record_mid_stream_failure(idx, Some(ReconnectReason::RetryableFailure));
if let Some(endpoint) = inner.connector.endpoint(idx) {
inner
.conn_events
.disconnected(&endpoint.host, &endpoint.port);
}
}
}
fn spawn_reaper(inner: Arc<DbInner>) -> std::io::Result<JoinHandle<()>> {
let tick = reaper_tick(inner.idle_timeout);
thread::Builder::new()
.name("questdb-ingress-pool-reaper".to_string())
.spawn(move || reaper_loop(inner, tick))
}
fn reaper_tick(idle_timeout: Duration) -> Duration {
let twelfth = idle_timeout / 12;
if twelfth > REAPER_MIN_TICK {
twelfth
} else {
REAPER_MIN_TICK
}
}
fn reaper_loop(inner: Arc<DbInner>, tick: Duration) {
loop {
let state = lock_state(&inner.state);
if inner.shutdown.load(Ordering::SeqCst) {
break;
}
let (state, _) = inner
.cv
.wait_timeout(state, tick)
.unwrap_or_else(|e| e.into_inner());
if inner.shutdown.load(Ordering::SeqCst) {
break;
}
drop(state);
reap_idle_inner(&inner);
}
}
fn reap_idle_inner(inner: &DbInner) -> usize {
let mut dropped = reap_idle_senders(inner);
dropped += reap_idle_direct_senders(inner);
#[cfg(feature = "_egress")]
{
dropped += reap_idle_readers(inner);
}
dropped
}
fn drain_idle_inner(inner: &DbInner) -> usize {
let mut dropped = drain_idle_senders(inner);
dropped += drain_idle_direct_senders(inner);
#[cfg(feature = "_egress")]
{
dropped += drain_idle_readers(inner);
}
dropped
}
fn drain_idle_senders(inner: &DbInner) -> usize {
let mut to_drop: Vec<PooledSenderCore> = {
let mut state = lock_state(&inner.state);
state.free.drain(..).map(|entry| entry.sender).collect()
};
let dropped = to_drop.len();
drain_sfa_senders_bounded(inner, &mut to_drop);
drop(to_drop);
dropped
}
fn drain_idle_direct_senders(inner: &DbInner) -> usize {
let to_drop: Vec<DirectSenderCore> = {
let mut state = lock_state(&inner.direct_state);
state.free.drain(..).map(|entry| entry.sender).collect()
};
let dropped = to_drop.len();
drop(to_drop);
dropped
}
#[cfg(feature = "_egress")]
fn drain_idle_readers(inner: &DbInner) -> usize {
let to_drop: Vec<Reader> = {
let mut state = lock_reader_state(&inner.reader_state);
state.free.drain(..).map(|entry| entry.reader).collect()
};
let dropped = to_drop.len();
drop(to_drop);
dropped
}
fn reap_idle_senders(inner: &DbInner) -> usize {
let durable = inner.connector.request_durable_ack();
let mut dropped = 0;
while let Some((sender, slot_index)) = take_reapable_column_sender(inner, durable) {
let _release = slot_index.is_some().then_some(SenderSlotRelease {
inner,
slot_index,
decrement_in_use: false,
decrement_closing: true,
});
drop(sender);
dropped += 1;
}
dropped
}
fn take_reapable_column_sender(
inner: &DbInner,
durable: bool,
) -> Option<(PooledSenderCore, Option<usize>)> {
let mut state = lock_state(&inner.state);
let now = Instant::now();
let mut i = 0;
while i < state.free.len() {
if state.total() <= inner.sender_pool_min {
return None;
}
let idle_for = now.saturating_duration_since(state.free[i].last_idle_at);
if idle_for > inner.idle_timeout && state.free[i].sender.sfa_fully_delivered(durable) {
let entry = state.free.remove(i);
if entry.slot_index.is_some() {
state.closing += 1;
}
return Some((entry.sender, entry.slot_index));
}
i += 1;
}
None
}
fn reap_idle_direct_senders(inner: &DbInner) -> usize {
let to_drop: Vec<DirectSenderCore> = {
let mut state = lock_state(&inner.direct_state);
let mut to_drop = Vec::new();
let now = Instant::now();
let mut i = 0;
while i < state.free.len() {
let idle_for = now.saturating_duration_since(state.free[i].last_idle_at);
if idle_for > inner.idle_timeout {
let entry = state.free.remove(i);
to_drop.push(entry.sender);
} else {
i += 1;
}
}
to_drop
};
let dropped = to_drop.len();
drop(to_drop);
dropped
}
#[cfg(feature = "_egress")]
fn reap_idle_readers(inner: &DbInner) -> usize {
let to_drop: Vec<Reader> = {
let mut state = lock_reader_state(&inner.reader_state);
let mut to_drop = Vec::new();
let now = Instant::now();
let mut i = 0;
while i < state.free.len() {
if state.total() <= inner.query_pool_min {
break;
}
let idle_for = now.saturating_duration_since(state.free[i].last_idle_at);
if idle_for > inner.idle_timeout {
let entry = state.free.remove(i);
to_drop.push(entry.reader);
} else {
i += 1;
}
}
to_drop
};
let dropped = to_drop.len();
drop(to_drop);
dropped
}
const _: fn() = || {
fn assert_send_sync<T: Send + Sync>() {}
assert_send_sync::<QuestDb>();
#[cfg(feature = "ffi-support")]
{
fn assert_send<T: Send>() {}
assert_send::<OwnedSender>();
assert_send::<OwnedDirectColumnSender>();
}
};
const _: fn() = || {
trait AmbiguousIfSend<A> {
fn _disambiguate() {}
}
impl<T: ?Sized> AmbiguousIfSend<()> for T {}
impl<T: ?Sized + Send> AmbiguousIfSend<u8> for T {}
fn assert_not_send<T: ?Sized>() {
let _: fn() = <T as AmbiguousIfSend<_>>::_disambiguate;
}
assert_not_send::<BorrowedSender<'_>>();
assert_not_send::<BorrowedDirectColumnSender<'_>>();
#[cfg(feature = "_egress")]
assert_not_send::<BorrowedReader<'_>>();
assert_not_send::<crate::ingress::column_sender::Chunk<'_>>();
};
const _: fn() = || {
trait AmbiguousIfSync<A> {
fn _disambiguate() {}
}
impl<T: ?Sized> AmbiguousIfSync<()> for T {}
impl<T: ?Sized + Sync> AmbiguousIfSync<u8> for T {}
fn assert_not_sync<T: ?Sized>() {
let _: fn() = <T as AmbiguousIfSync<_>>::_disambiguate;
}
assert_not_sync::<BorrowedSender<'_>>();
assert_not_sync::<BorrowedDirectColumnSender<'_>>();
#[cfg(feature = "_egress")]
assert_not_sync::<BorrowedReader<'_>>();
assert_not_sync::<crate::ingress::column_sender::Chunk<'_>>();
};
#[cfg(test)]
mod tests {
use std::fs;
use tempfile::TempDir;
use super::{SlotReservations, managed_slot_recovery_scan_from};
fn dirty_slot(root: &std::path::Path, name: &str) {
let slot = root.join(name);
fs::create_dir(&slot).unwrap();
fs::write(slot.join("sf-0.sfa"), b"queued").unwrap();
}
#[test]
fn managed_slot_recovery_candidates_exclude_live_pool_range() {
let temp = TempDir::new().unwrap();
dirty_slot(temp.path(), "default-ingest-0");
dirty_slot(temp.path(), "default-ingest-1");
dirty_slot(temp.path(), "default-ingest-2");
dirty_slot(temp.path(), &format!("default-{}-2", "col"));
dirty_slot(temp.path(), &format!("default-{}-2", "row"));
let mut actual = managed_slot_recovery_scan_from(temp.path(), "default", 2).out_of_range;
actual.sort();
assert_eq!(actual, vec![temp.path().join("default-ingest-2")]);
}
#[test]
fn slot_reservations_reserve_specific_index() {
let mut disk = SlotReservations::with_disk_slots(2);
assert!(disk.reserve(1));
assert!(!disk.reserve(1), "double-reserve must fail");
assert!(!disk.reserve(2), "out-of-range reserve must fail");
disk.free(Some(1));
assert!(disk.reserve(1), "freed slot can be reserved again");
let mut in_memory = SlotReservations::default();
assert!(!in_memory.reserve(0));
}
}