use super::{LookupQueue, PrefixLookupQuery, PrimaryLookupQuery, QueuedLookup};
use crate::client::lookup::lookup_sender::LookupSender;
use crate::client::metadata::Metadata;
use crate::config::Config;
use crate::error::{Error, Result};
use crate::metadata::{TableBucket, TablePath};
use bytes::Bytes;
use log::{debug, error};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use tokio::sync::{mpsc, watch};
use tokio::task::JoinHandle;
pub struct LookupClient {
lookup_tx: mpsc::Sender<QueuedLookup>,
sender_handle: Option<JoinHandle<()>>,
shutdown_tx: watch::Sender<bool>,
closed: AtomicBool,
}
impl LookupClient {
pub fn new(config: &Config, metadata: Arc<Metadata>) -> Self {
let queue_size = config.lookup_queue_size;
let max_batch_size = config.lookup_max_batch_size;
let batch_timeout_ms = config.lookup_batch_timeout_ms;
let max_inflight = config.lookup_max_inflight_requests;
let max_retries = config.lookup_max_retries;
let cluster_rx = metadata.subscribe_cluster_changes();
let (queue, lookup_tx, re_enqueue_tx) =
LookupQueue::new(queue_size, max_batch_size, batch_timeout_ms, cluster_rx);
let (shutdown_tx, shutdown_rx) = watch::channel(false);
let mut sender = LookupSender::new(
metadata,
queue,
re_enqueue_tx,
max_inflight,
max_retries,
shutdown_rx,
);
let sender_handle = tokio::spawn(async move {
sender.run().await;
debug!("Lookup sender completed");
});
Self {
lookup_tx,
sender_handle: Some(sender_handle),
shutdown_tx,
closed: AtomicBool::new(false),
}
}
pub async fn lookup(
&self,
table_path: TablePath,
table_bucket: TableBucket,
key_bytes: Bytes,
) -> Result<Option<Vec<u8>>> {
if self.closed.load(Ordering::Acquire) {
return Err(Error::UnexpectedError {
message: "Lookup client is closed".to_string(),
source: None,
});
}
let (result_tx, result_rx) = tokio::sync::oneshot::channel();
let query = QueuedLookup::Primary(PrimaryLookupQuery::new(
table_path,
table_bucket,
key_bytes,
result_tx,
));
self.enqueue(query).await?;
result_rx.await.map_err(|_| Error::UnexpectedError {
message: "Lookup result channel closed".to_string(),
source: None,
})?
}
pub async fn prefix_lookup(
&self,
table_path: TablePath,
table_bucket: TableBucket,
key_bytes: Bytes,
) -> Result<Vec<Vec<u8>>> {
if self.closed.load(Ordering::Acquire) {
return Err(Error::UnexpectedError {
message: "Lookup client is closed".to_string(),
source: None,
});
}
let (result_tx, result_rx) = tokio::sync::oneshot::channel();
let query = QueuedLookup::Prefix(PrefixLookupQuery::new(
table_path,
table_bucket,
key_bytes,
result_tx,
));
self.enqueue(query).await?;
result_rx.await.map_err(|_| Error::UnexpectedError {
message: "Lookup result channel closed".to_string(),
source: None,
})?
}
async fn enqueue(&self, query: QueuedLookup) -> Result<()> {
self.lookup_tx.send(query).await.map_err(|e| {
let failed_query = e.0;
error!(
"Failed to queue lookup: channel closed. table_path: {}, table_bucket: {:?}, key_len: {}",
failed_query.table_path(),
failed_query.table_bucket(),
failed_query.key().len()
);
Error::UnexpectedError {
message: "Failed to queue lookup: channel closed".to_string(),
source: None,
}
})
}
pub async fn close(mut self, timeout: Duration) {
debug!("Closing lookup client");
self.closed.store(true, Ordering::Release);
let _ = self.shutdown_tx.send(true);
if let Some(handle) = self.sender_handle.take() {
debug!("Waiting for sender task to complete...");
let abort_handle = handle.abort_handle();
match tokio::time::timeout(timeout, handle).await {
Ok(Ok(())) => {
debug!("Lookup sender task completed gracefully.");
}
Ok(Err(join_error)) => {
error!("Lookup sender task panicked: {:?}", join_error);
}
Err(_elapsed) => {
error!("Lookup sender task did not complete within timeout. Forcing shutdown.");
abort_handle.abort();
}
}
} else {
debug!("Lookup client was already closed or never initialized properly.");
}
debug!("Lookup client closed");
}
}
impl Drop for LookupClient {
fn drop(&mut self) {
if let Some(handle) = self.sender_handle.take() {
handle.abort();
}
}
}