Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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 @@ -24,6 +24,9 @@ rustls-pki-types.workspace = true
http.workspace = true
h2.workspace = true
yaml-rust = { workspace = true, optional = true }
futures-util.workspace = true
arc-swap.workspace = true
log.workspace = true
vey-types = { workspace = true, features = ["rustls"] }
vey-io-ext = { workspace = true, features = ["rustls"] }
vey-socket.workspace = true
Expand All @@ -32,6 +35,9 @@ vey-h2.workspace = true
vey-smtp-proto.workspace = true
vey-yaml = { workspace = true, optional = true, features = ["rustls", "http"] }

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

[features]
default = []
yaml = ["dep:vey-yaml", "dep:yaml-rust"]
67 changes: 16 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,45 @@

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};
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 @@ -7,18 +7,17 @@
use std::io;
use std::net::SocketAddr;
use std::sync::Arc;
use std::time::Duration;
use std::task::Poll;

use anyhow::Context;
use futures_util::poll;
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