Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 6 additions & 0 deletions lib/vey-icap-client/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,12 @@ vey-http.workspace = true
vey-h2.workspace = true
vey-smtp-proto.workspace = true
vey-yaml = { workspace = true, optional = true, features = ["rustls", "http"] }
futures-util.workspace = true
Comment thread
DanielHaimanot marked this conversation as resolved.
Outdated
arc-swap.workspace = true
log.workspace = true

[dev-dependencies]
tokio = { workspace = true, features = ["rt-multi-thread"] }

[features]
default = []
Expand Down
68 changes: 17 additions & 51 deletions lib/vey-icap-client/src/service/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,80 +5,46 @@

use std::sync::Arc;

use anyhow::anyhow;
use tokio::sync::oneshot;

use super::{
IcapClientConnection, IcapConnector, IcapServiceClientCommand, IcapServiceConfig,
IcapServicePool,
IcapClientConnection, IcapConnectionPool, IcapConnector, IcapServiceConfig, PoolMaintainer,
};
use crate::options::{IcapOptionsRequest, IcapServiceOptions};

Comment thread
DanielHaimanot marked this conversation as resolved.
Outdated
use crate::options::IcapServiceOptions;

pub struct IcapServiceClient {
pub(crate) config: Arc<IcapServiceConfig>,
pub(crate) partial_request_header: Vec<u8>,
cmd_sender: kanal::AsyncSender<IcapServiceClientCommand>,
conn_creator: Arc<IcapConnector>,
conn_pool: Arc<IcapConnectionPool>,
}

impl IcapServiceClient {
pub fn new(config: Arc<IcapServiceConfig>) -> anyhow::Result<Self> {
let (cmd_sender, cmd_receiver) = kanal::unbounded_async();
let conn_creator = IcapConnector::new(config.clone())?;
let conn_creator = Arc::new(conn_creator);
let pool = IcapServicePool::new(config.clone(), cmd_receiver, conn_creator.clone());
tokio::spawn(pool.into_running());
let connector = Arc::new(IcapConnector::new(config.clone())?);
let conn_pool = Arc::new(IcapConnectionPool::new(config.clone(), connector));
let check_interval = config.connection_pool.check_interval();

let maintainer = PoolMaintainer::new(Arc::downgrade(&conn_pool), check_interval);

tokio::spawn(maintainer.into_running());

let partial_request_header = config.build_request_header();
Ok(IcapServiceClient {
config,
partial_request_header,
cmd_sender,
conn_creator,
conn_pool,
})
}

async fn fetch_from_pool(&self) -> Option<(IcapClientConnection, Arc<IcapServiceOptions>)> {
let (rsp_sender, rsp_receiver) = oneshot::channel();
let cmd = IcapServiceClientCommand::FetchConnection(rsp_sender);
if self.cmd_sender.send(cmd).await.is_ok() {
rsp_receiver.await.ok()
} else {
None
}
}

pub async fn fetch_connection(
&self,
) -> anyhow::Result<(IcapClientConnection, Arc<IcapServiceOptions>)> {
if let Some(conn) = self.fetch_from_pool().await {
return Ok(conn);
}

let mut conn = self
.conn_creator
.create()
.await
.map_err(|e| anyhow!("create new connection failed: {e:?}"))?;
let options_req = IcapOptionsRequest::new(self.config.as_ref());

conn.mark_io_inuse();
let options = options_req
.get_options(&mut conn, self.config.icap_max_header_size)
.await
.map_err(|e| anyhow!("failed to get icap service options: {e}"))?;

let mut conn = self.conn_pool.get().await?;
let options = self.conn_pool.get_options();
conn.mark_io_inuse();
Ok((conn, Arc::new(options)))
Ok((conn, options))
}

pub fn save_connection(&self, conn: IcapClientConnection) {
if conn.reusable() {
let pool_sender = self.cmd_sender.clone();
tokio::spawn(async move {
let _ = pool_sender
.send(IcapServiceClientCommand::SaveConnection(conn))
.await;
});
}
self.conn_pool.try_put(conn);
}
}
6 changes: 6 additions & 0 deletions lib/vey-icap-client/src/service/config/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ pub struct IcapServiceConfig {
pub(crate) icap_206_enable: bool,
pub(crate) icap_max_header_size: usize,
pub(crate) disable_preview: bool,
pub(crate) options_timeout: Duration,
pub(crate) preview_data_read_timeout: Duration,
pub(crate) respond_shared_names: BTreeSet<String>,
pub(crate) bypass: bool,
Expand Down Expand Up @@ -87,6 +88,7 @@ impl IcapServiceConfig {
icap_206_enable: false,
icap_max_header_size: 8192,
disable_preview: false,
options_timeout: Duration::from_secs(1),
preview_data_read_timeout: Duration::from_secs(4),
respond_shared_names: BTreeSet::new(),
bypass: false,
Expand All @@ -113,6 +115,10 @@ impl IcapServiceConfig {
self.icap_max_header_size = max_size;
}

pub fn set_options_read_timeout(&mut self, time: Duration) {
self.options_timeout = time;
}

pub fn set_preview_data_read_timeout(&mut self, time: Duration) {
self.preview_data_read_timeout = time;
}
Expand Down
6 changes: 6 additions & 0 deletions lib/vey-icap-client/src/service/config/yaml.rs
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,12 @@ impl IcapServiceConfig {
config.disable_preview = vey_yaml::value::as_bool(v)?;
Ok(())
}
"options_timeout" => {
let time = vey_yaml::humanize::as_duration(v)
.context(format!("invalid humanize duration value for key {k}"))?;
config.set_options_read_timeout(time);
Ok(())
}
"preview_data_read_timeout" => {
let time = vey_yaml::humanize::as_duration(v)
.context(format!("invalid humanize duration value for key {k}"))?;
Expand Down
113 changes: 58 additions & 55 deletions lib/vey-icap-client/src/service/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,21 +4,20 @@
* SPDX-FileCopyrightText: 2026 VEY-OSS Developers.
*/

use futures_util::poll;
Comment thread
DanielHaimanot marked this conversation as resolved.
Outdated
use std::io;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use std::task::Poll;

use anyhow::Context;
use tokio::io::{AsyncRead, AsyncWrite, BufReader};
use tokio::sync::oneshot;
use tokio_rustls::TlsConnector;

use vey_io_ext::{AsyncStream, LimitedBufReadExt};
use vey_types::net::{Host, RustlsClientConfig};

use super::IcapServiceConfig;
use crate::IcapServiceOptions;

pub type IcapClientWriter = Box<dyn AsyncWrite + Send + Sync + Unpin>;
pub type IcapClientReader = BufReader<Box<dyn AsyncRead + Send + Sync + Unpin>>;
Expand All @@ -32,7 +31,7 @@ pub struct IcapClientConnection {
}

impl IcapClientConnection {
fn new<R, W>(reader: R, writer: W) -> Self
pub(super) fn new<R, W>(reader: R, writer: W) -> Self
where
R: AsyncRead + Send + Sync + Unpin + 'static,
W: AsyncWrite + Send + Sync + Unpin + 'static,
Expand All @@ -50,6 +49,14 @@ impl IcapClientConnection {
self.reused_connection
}

pub(super) fn mark_reused(&mut self) {
self.reused_connection = true
}

pub(super) fn reusable(&self) -> bool {
self.reader_clean && self.writer_clean
}

pub fn mark_reader_finished(&mut self) {
self.reader_clean = true;
}
Expand All @@ -63,8 +70,8 @@ impl IcapClientConnection {
self.writer_clean = false;
}

pub(super) fn reusable(&self) -> bool {
self.reader_clean && self.writer_clean
pub(super) async fn probe_idle(&mut self) -> bool {
matches!(poll!(self.reader.fill_wait_data()), Poll::Pending)
}
}

Expand Down Expand Up @@ -146,59 +153,55 @@ impl IcapConnector {
}
}

pub(super) struct IcapConnectionPollRequest {
client_sender: oneshot::Sender<(IcapClientConnection, Arc<IcapServiceOptions>)>,
options: Arc<IcapServiceOptions>,
}
#[allow(unused_imports)]
#[cfg(test)]
mod tests {
use super::*;
use tokio::io::AsyncWriteExt;

impl IcapConnectionPollRequest {
pub(super) fn new(
client_sender: oneshot::Sender<(IcapClientConnection, Arc<IcapServiceOptions>)>,
options: Arc<IcapServiceOptions>,
) -> Self {
IcapConnectionPollRequest {
client_sender,
options,
}
}
}
#[tokio::test]
async fn probe_detects_closed_idle_connection() {
let (client, server) = tokio::io::duplex(1024);

pub(super) struct IcapConnectionEofPoller {
conn: IcapClientConnection,
req_receiver: kanal::AsyncReceiver<IcapConnectionPollRequest>,
}
let (reader, writer) = tokio::io::split(client);

impl IcapConnectionEofPoller {
pub(super) fn new(
conn: IcapClientConnection,
req_receiver: &kanal::AsyncReceiver<IcapConnectionPollRequest>,
) -> Option<Self> {
if conn.reusable() {
Some(IcapConnectionEofPoller {
conn,
req_receiver: req_receiver.clone(),
})
} else {
None
}
let mut conn = IcapClientConnection {
reader: BufReader::new(Box::new(reader)),
writer: Box::new(writer),
reader_clean: true,
writer_clean: true,
reused_connection: false,
};

// Idle live connection: nothing readable.
assert!(conn.probe_idle().await);

// Simulate ICAP server closing its side.
drop(server);

tokio::task::yield_now().await;

// EOF should now be observable.
assert!(!conn.probe_idle().await);
}

pub(super) async fn into_running(mut self, idle_timeout: Duration) {
let idle_sleep = tokio::time::sleep(idle_timeout);

tokio::select! {
_ = self.conn.reader.fill_wait_data() => {}
_ = idle_sleep => {}
r = self.req_receiver.recv() => {
if let Ok(req) = r {
let IcapConnectionPollRequest {
client_sender,
options,
} = req;
self.conn.reused_connection = true;
let _ = client_sender.send((self.conn, options));
}
}
}
#[tokio::test]
async fn probe_rejects_unexpected_data() {
let (client, mut server) = tokio::io::duplex(1024);
let (reader, writer) = tokio::io::split(client);

let mut conn = IcapClientConnection {
reader: BufReader::new(Box::new(reader)),
writer: Box::new(writer),
reader_clean: true,
writer_clean: true,
reused_connection: false,
};

server.write_all(b"garbage").await.unwrap();

tokio::task::yield_now().await;

assert!(!conn.probe_idle().await);
}
}
4 changes: 2 additions & 2 deletions lib/vey-icap-client/src/service/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,14 +8,14 @@ mod config;
pub use config::IcapServiceConfig;

mod connection;
use connection::IcapConnector;
pub(super) use connection::{IcapClientConnection, IcapClientReader, IcapClientWriter};
use connection::{IcapConnectionEofPoller, IcapConnectionPollRequest, IcapConnector};

mod client;
pub use client::IcapServiceClient;

mod pool;
use pool::{IcapServiceClientCommand, IcapServicePool};
use pool::{IcapConnectionPool, PoolMaintainer};

#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum IcapMethod {
Expand Down
Loading
Loading