diff --git a/Cargo.lock b/Cargo.lock index a475f8fce..690ad5218 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -321,9 +321,9 @@ dependencies = [ [[package]] name = "bssl-cmake-sys" -version = "0.1.2607130" +version = "0.1.2607300" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "7935c69a9ec131b051320d33812b6d32b89b333d8929e878a417cb341ff16fd9" +checksum = "2c5b662a5345336135fd6b35209b8e2062db5c4782802d02d7824646a9c8cdd4" dependencies = [ "bindgen", "cc", @@ -522,9 +522,9 @@ dependencies = [ [[package]] name = "clang-sys" -version = "1.8.1" +version = "1.9.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "0b023947811758c97c59bf9d1c188fd619ad4718dcaa767947df1cadb14f39f4" +checksum = "157a8ba7b480713b56f4c09fd13fc3e0a22a5dfab8097ba61cbc5feef950788a" dependencies = [ "glob", "libc", @@ -1911,9 +1911,9 @@ dependencies = [ [[package]] name = "redis" -version = "1.4.1" +version = "1.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b0b9503711b03773e43b31668c7b5bd279ee7cd9b7d18cff7c23a42cc1d08e5a" +checksum = "3257df217f7eab0044627a268c9cc6cdb60c0c421c88f83ac41c4e31520b6b84" dependencies = [ "arcstr", "async-lock", @@ -2042,9 +2042,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.42" +version = "0.23.43" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3c54fcab019b409d04215d3a17cb438fd7fbf192ee61461f20f4fe18704bc138" +checksum = "0283386ce02abc0151e1761d08802dfe86c173b0b494af5cbc086574e453da06" dependencies = [ "aws-lc-rs", "brotli", @@ -2415,13 +2415,13 @@ dependencies = [ [[package]] name = "tokio-macros" -version = "2.7.1" +version = "2.7.2" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6328af13490e73a9b4694030fafd93f8c8c6a9dede33e821c3fc63eddf8042ba" +checksum = "78773a2a397f451582ce068015985c33193cf6dea8b74d2a639fe457b2f07b0e" dependencies = [ "proc-macro2", "quote", - "syn 2.0.119", + "syn 3.0.3", ] [[package]] diff --git a/lib/vey-fluentd/src/config/mod.rs b/lib/vey-fluentd/src/config/mod.rs index 4f2e26160..1075e072f 100644 --- a/lib/vey-fluentd/src/config/mod.rs +++ b/lib/vey-fluentd/src/config/mod.rs @@ -291,3 +291,267 @@ impl FluentdClientConfig { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::handshake::{PongMsgRef, parse_helo}; + use openssl::md::Md; + use openssl::md_ctx::MdCtx; + use std::net::{Ipv4Addr, SocketAddr}; + use std::str::FromStr; + use tokio::io::{AsyncReadExt, AsyncWriteExt, duplex}; + + fn sha512_hex(parts: &[&[u8]]) -> String { + let mut md = MdCtx::new().unwrap(); + md.digest_init(Md::sha512()).unwrap(); + for part in parts { + md.digest_update(part).unwrap(); + } + let mut hash = [0u8; FLUENTD_HASH_SIZE]; + md.digest_final(&mut hash).unwrap(); + hex::encode(hash) + } + + #[test] + fn defaults_and_setters() { + let mut cfg = FluentdClientConfig::default(); + assert_eq!( + cfg.server_addr, + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), FLUENTD_DEFAULT_PORT) + ); + assert_eq!(cfg.connect_timeout, Duration::from_secs(10)); + assert_eq!(cfg.connect_delay, Duration::from_secs(10)); + assert_eq!(cfg.write_timeout, Duration::from_secs(1)); + assert_eq!(cfg.flush_interval, Duration::from_millis(100)); + assert_eq!(cfg.retry_queue_len, 10); + assert!(cfg.shared_key.is_empty()); + assert!(cfg.tls_client.is_none()); + + cfg.set_server_addr("10.0.0.1:24225".parse().unwrap()); + cfg.set_bind_ip("127.0.0.1".parse().unwrap()); + cfg.set_shared_key("sk".into()); + cfg.set_username("u".into()); + cfg.set_password("p".into()); + cfg.set_hostname("host".into()); + cfg.set_connect_timeout(Duration::from_secs(3)); + cfg.set_connect_delay(Duration::from_secs(4)); + cfg.set_write_timeout(Duration::from_millis(250)); + cfg.set_flush_interval(Duration::from_millis(50)); + cfg.set_retry_queue_len(7); + cfg.set_tls_name(Host::from_str("example.com").unwrap()); + + assert_eq!(cfg.server_addr, "10.0.0.1:24225".parse().unwrap()); + assert_eq!(cfg.bind.ip().unwrap().to_string(), "127.0.0.1"); + assert_eq!(cfg.shared_key, "sk"); + assert_eq!(cfg.username, "u"); + assert_eq!(cfg.password, "p"); + assert_eq!(cfg.hostname, "host"); + assert_eq!(cfg.connect_timeout, Duration::from_secs(3)); + assert_eq!(cfg.connect_delay, Duration::from_secs(4)); + assert_eq!(cfg.write_timeout, Duration::from_millis(250)); + assert_eq!(cfg.flush_interval, Duration::from_millis(50)); + assert_eq!(cfg.retry_queue_len, 7); + assert!(cfg.tls_name.is_some()); + } + + #[tokio::test] + async fn handshake_skipped_without_shared_key() { + let cfg = FluentdClientConfig::default(); + let (client, _server) = duplex(64); + cfg.handshake(client).await.unwrap(); + } + + #[test] + fn build_ping_and_verify_pong_roundtrip() { + let mut cfg = FluentdClientConfig::default(); + cfg.set_shared_key("secret".into()); + cfg.set_hostname("client-host".into()); + cfg.set_username("alice".into()); + cfg.set_password("pw".into()); + + let helo = parse_helo(&[ + 0x92, 0xa4, b'H', b'E', b'L', b'O', 0x82, 0xa5, b'n', b'o', b'n', b'c', b'e', 0xab, + b'n', b'o', b'n', b'c', b'e', b'-', b'b', b'y', b't', b'e', b's', 0xa4, b'a', b'u', + b't', b'h', 0xaa, b's', b'a', b'l', b't', b'-', b'b', b'y', b't', b'e', b's', + ]) + .unwrap(); + assert_eq!(helo.nonce, b"nonce-bytes"); + assert_eq!(helo.auth_salt, b"salt-bytes"); + let shared_key_salt = b"0123456789abcdef"; + let ping = cfg.build_ping(&helo, shared_key_salt).unwrap(); + + // PING is ["PING", hostname, salt, shared_key_digest, username, password_digest] + let mut bytes = rmp::decode::Bytes::new(&ping); + assert_eq!(rmp::decode::read_array_len(&mut bytes).unwrap(), 6); + let (msg_type, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(msg_type, "PING"); + bytes = rmp::decode::Bytes::new(rest); + let (hostname, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(hostname, "client-host"); + bytes = rmp::decode::Bytes::new(rest); + let salt_len = rmp::decode::read_bin_len(&mut bytes).unwrap() as usize; + assert_eq!(&bytes.remaining_slice()[..salt_len], shared_key_salt); + bytes = rmp::decode::Bytes::new(&bytes.remaining_slice()[salt_len..]); + let (client_digest, rest) = + rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!( + client_digest, + sha512_hex(&[shared_key_salt, b"client-host", b"nonce-bytes", b"secret"]) + ); + bytes = rmp::decode::Bytes::new(rest); + let (username, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(username, "alice"); + bytes = rmp::decode::Bytes::new(rest); + let (password_digest, _) = + rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!( + password_digest, + sha512_hex(&[b"salt-bytes", b"alice", b"pw"]) + ); + + let server_hostname = "fluentd.example"; + let shared_key_digest = sha512_hex(&[ + shared_key_salt, + server_hostname.as_bytes(), + b"nonce-bytes", + b"secret", + ]); + let pong = PongMsgRef { + auth_result: true, + reason: "", + server_hostname, + shared_key_digest: &shared_key_digest, + }; + cfg.verify_pong(pong, shared_key_salt, b"nonce-bytes") + .unwrap(); + } + + #[test] + fn verify_pong_rejects_auth_failure_and_digest_mismatch() { + let mut cfg = FluentdClientConfig::default(); + cfg.set_shared_key("secret".into()); + + let failed = PongMsgRef { + auth_result: false, + reason: "bad key", + server_hostname: "s", + shared_key_digest: "", + }; + assert!( + cfg.verify_pong(failed, b"salt", b"nonce") + .unwrap_err() + .to_string() + .contains("server auth failed") + ); + + let bad_hex = PongMsgRef { + auth_result: true, + reason: "", + server_hostname: "s", + shared_key_digest: "zz", + }; + assert!( + cfg.verify_pong(bad_hex, b"salt", b"nonce") + .unwrap_err() + .to_string() + .contains("invalid shared_key_hex_digest") + ); + + let mismatch = PongMsgRef { + auth_result: true, + reason: "", + server_hostname: "s", + shared_key_digest: &"ab".repeat(64), + }; + assert!( + cfg.verify_pong(mismatch, b"salt", b"nonce") + .unwrap_err() + .to_string() + .contains("mismatch") + ); + } + + #[test] + fn build_ping_omits_user_digest_when_nonce_empty() { + let mut cfg = FluentdClientConfig::default(); + cfg.set_shared_key("secret".into()); + cfg.set_hostname("h".into()); + cfg.set_username("alice".into()); + cfg.set_password("pw".into()); + + let helo = parse_helo(&[0x92, 0xa4, b'H', b'E', b'L', b'O', 0x80]).unwrap(); + assert_eq!(helo.nonce, b""); + let ping = cfg.build_ping(&helo, b"salt0123456789ab").unwrap(); + let mut bytes = rmp::decode::Bytes::new(&ping); + assert_eq!(rmp::decode::read_array_len(&mut bytes).unwrap(), 6); + let (_, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + bytes = rmp::decode::Bytes::new(rest); + let (_, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + bytes = rmp::decode::Bytes::new(rest); + let salt_len = rmp::decode::read_bin_len(&mut bytes).unwrap() as usize; + bytes = rmp::decode::Bytes::new(&bytes.remaining_slice()[salt_len..]); + let (_, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + bytes = rmp::decode::Bytes::new(rest); + let (username, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(username, ""); + bytes = rmp::decode::Bytes::new(rest); + let (password_digest, _) = + rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(password_digest, ""); + } + + #[tokio::test] + async fn handshake_roundtrip_with_shared_key() { + let mut cfg = FluentdClientConfig::default(); + cfg.set_shared_key("secret".into()); + cfg.set_hostname("client-host".into()); + cfg.set_username("alice".into()); + cfg.set_password("pw".into()); + + let (client, mut server) = duplex(4096); + let server_task = tokio::spawn(async move { + // HELO with string nonce/auth + let helo: &[u8] = &[ + 0x92, 0xa4, b'H', b'E', b'L', b'O', 0x82, 0xa5, b'n', b'o', b'n', b'c', b'e', 0xa5, + b'n', b'o', b'n', b'c', b'e', 0xa4, b'a', b'u', b't', b'h', 0xa4, b's', b'a', b'l', + b't', + ]; + assert_eq!(parse_helo(helo).unwrap().nonce, b"nonce"); + server.write_all(helo).await.unwrap(); + + let mut buf = vec![0u8; 2048]; + let n = server.read(&mut buf).await.unwrap(); + let ping = &buf[..n]; + + let mut bytes = rmp::decode::Bytes::new(ping); + assert_eq!(rmp::decode::read_array_len(&mut bytes).unwrap(), 6); + let (_, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + bytes = rmp::decode::Bytes::new(rest); + let (_, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + bytes = rmp::decode::Bytes::new(rest); + let salt_len = rmp::decode::read_bin_len(&mut bytes).unwrap() as usize; + let shared_key_salt = bytes.remaining_slice()[..salt_len].to_vec(); + + let server_hostname = "fluentd"; + let digest = sha512_hex(&[ + shared_key_salt.as_slice(), + server_hostname.as_bytes(), + b"nonce", + b"secret", + ]); + + let mut pong = Vec::new(); + rmp::encode::write_array_len(&mut pong, 5).unwrap(); + rmp::encode::write_str(&mut pong, "PONG").unwrap(); + rmp::encode::write_bool(&mut pong, true).unwrap(); + rmp::encode::write_str(&mut pong, "").unwrap(); + rmp::encode::write_str(&mut pong, server_hostname).unwrap(); + rmp::encode::write_str(&mut pong, &digest).unwrap(); + server.write_all(&pong).await.unwrap(); + }); + + cfg.handshake(client).await.unwrap(); + server_task.await.unwrap(); + } +} diff --git a/lib/vey-fluentd/src/format.rs b/lib/vey-fluentd/src/format.rs index 7bab0bc78..9913d1a68 100644 --- a/lib/vey-fluentd/src/format.rs +++ b/lib/vey-fluentd/src/format.rs @@ -196,6 +196,7 @@ impl Serializer for FormatterKv<'_> { #[cfg(test)] mod tests { use super::*; + use slog::{OwnedKVList, Record, RecordLocation, RecordStatic, Serializer, b, o}; #[test] fn emit_u64_writes_key_and_value() { @@ -222,4 +223,114 @@ mod tests { assert!(buf.windows(3).any(|w| w == b"msg")); assert!(buf.windows(5).any(|w| w == b"hello")); } + + #[test] + fn emit_scalar_types() { + let mut buf = Vec::new(); + let mut fmt = FormatterKv(&mut buf); + fmt.emit_bool("ok".into(), true).unwrap(); + fmt.emit_i64("neg".into(), -7).unwrap(); + fmt.emit_f64("pi".into(), 3.5).unwrap(); + fmt.emit_char("ch".into(), 'Z').unwrap(); + assert!(buf.windows(2).any(|w| w == b"ok")); + assert!(buf.windows(3).any(|w| w == b"neg")); + assert!(buf.windows(2).any(|w| w == b"pi")); + assert!(buf.windows(2).any(|w| w == b"ch")); + assert!(buf.windows(1).any(|w| w == b"Z")); + assert!(buf.contains(&0xc3)); // true + } + + #[test] + fn emit_arguments_static_and_formatted() { + let mut buf = Vec::new(); + let mut fmt = FormatterKv(&mut buf); + fmt.emit_arguments("static".into(), &format_args!("plain")) + .unwrap(); + fmt.emit_arguments("fmt".into(), &format_args!("n={}", 9)) + .unwrap(); + assert!(buf.windows(5).any(|w| w == b"plain")); + assert!(buf.windows(3).any(|w| w == b"n=9")); + } + + #[test] + fn counter_kv_counts_emits() { + let mut counter = CounterKV(0); + counter + .emit_arguments("a".into(), &format_args!("1")) + .unwrap(); + counter + .emit_arguments("b".into(), &format_args!("2")) + .unwrap(); + assert_eq!(counter.0, 2); + } + + #[test] + fn format_slog_encodes_fluent_forward_shape() { + static LOC: RecordLocation = RecordLocation { + file: file!(), + line: line!(), + column: 0, + module: module_path!(), + function: "", + }; + static RS: RecordStatic = RecordStatic { + location: &LOC, + tag: "", + level: slog::Level::Info, + }; + + let fmt = FluentdFormatter::new("vey.app".to_owned()); + let msg = format_args!("hello {}", "world"); + let kv = b!("count" => 3u64); + let record = Record::new(&RS, &msg, kv); + let owned: OwnedKVList = o!("host" => "local").into(); + let buf = fmt.format_slog(&record, &owned).unwrap(); + + let mut bytes = rmp::decode::Bytes::new(&buf); + assert_eq!(rmp::decode::read_array_len(&mut bytes).unwrap(), 3); + + let (tag, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(tag, "vey.app"); + bytes = rmp::decode::Bytes::new(rest); + + let ext = rmp::decode::read_ext_meta(&mut bytes).unwrap(); + assert_eq!(ext.typeid, 0); + assert_eq!(ext.size, 8); + let rem = bytes.remaining_slice(); + bytes = rmp::decode::Bytes::new(&rem[8..]); + + let map_len = rmp::decode::read_map_len(&mut bytes).unwrap(); + assert_eq!(map_len, 3); // host + count + msg + + let mut saw_msg = false; + let mut saw_host = false; + let mut saw_count = false; + for _ in 0..map_len { + let (key, rest) = rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + bytes = rmp::decode::Bytes::new(rest); + match key { + "msg" => { + let (value, rest) = + rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(value, "hello world"); + bytes = rmp::decode::Bytes::new(rest); + saw_msg = true; + } + "host" => { + let (value, rest) = + rmp::decode::read_str_from_slice(bytes.remaining_slice()).unwrap(); + assert_eq!(value, "local"); + bytes = rmp::decode::Bytes::new(rest); + saw_host = true; + } + "count" => { + let value = rmp::decode::read_int::(&mut bytes).unwrap(); + assert_eq!(value, 3); + saw_count = true; + } + other => panic!("unexpected key {other}"), + } + } + assert!(saw_msg && saw_host && saw_count); + } } diff --git a/lib/vey-fluentd/src/handshake.rs b/lib/vey-fluentd/src/handshake.rs index f355d14aa..bd085c7a8 100644 --- a/lib/vey-fluentd/src/handshake.rs +++ b/lib/vey-fluentd/src/handshake.rs @@ -239,6 +239,32 @@ mod tests { 0x92, 0xa4, b'H', b'E', b'L', b'O', 0x81, 0xa4, b'a', b'u', b't', b'h', 0x01, ]; assert!(parse_helo(buf).is_err()); + + let buf: &[u8] = &[ + 0x92, 0xa4, b'H', b'E', b'L', b'O', 0x81, 0xa9, b'k', b'e', b'e', b'p', b'a', b'l', + b'i', b'v', b'e', 0xa4, b't', b'r', b'u', b'e', + ]; + assert!(parse_helo(buf).is_err()); + } + + #[test] + fn parse_helo_empty_options_defaults() { + let buf: &[u8] = &[0x92, 0xa4, b'H', b'E', b'L', b'O', 0x80]; + let helo = parse_helo(buf).unwrap(); + assert_eq!(helo.nonce, b""); + assert_eq!(helo.auth_salt, b""); + assert!(helo.keepalive); + } + + #[test] + fn parse_helo_empty_bin_nonce() { + let buf: &[u8] = &[ + 0x92, 0xa4, b'H', b'E', b'L', b'O', 0x81, 0xa5, b'n', b'o', b'n', b'c', b'e', 0xc4, + 0x00, + ]; + let helo = parse_helo(buf).unwrap(); + assert_eq!(helo.nonce, b""); + assert!(helo.keepalive); } #[test] diff --git a/lib/vey-fluentd/src/lib.rs b/lib/vey-fluentd/src/lib.rs index 682a112ca..7b152c409 100644 --- a/lib/vey-fluentd/src/lib.rs +++ b/lib/vey-fluentd/src/lib.rs @@ -201,3 +201,33 @@ impl AsyncIoThread { } } } + +#[cfg(test)] +mod tests { + use super::*; + use std::sync::Arc; + use vey_types::log::LogStats; + + #[test] + fn push_to_retry_keeps_and_drops_overflow() { + let mut cfg = FluentdClientConfig::default(); + cfg.set_retry_queue_len(2); + let (_sender, receiver) = kanal::bounded::>(1); + let mut thread = AsyncIoThread { + config: Arc::new(cfg), + receiver: receiver.clone_async(), + stats: Arc::new(LogStats::default()), + retry_queue: VecDeque::new(), + }; + + assert!(thread.push_to_retry(vec![1]).is_none()); + assert!(thread.push_to_retry(vec![2]).is_none()); + assert_eq!(thread.retry_queue.len(), 2); + + let dropped = thread.push_to_retry(vec![3]).unwrap(); + assert_eq!(dropped, vec![1]); + assert_eq!(thread.retry_queue.len(), 2); + assert_eq!(thread.retry_queue[0], vec![2]); + assert_eq!(thread.retry_queue[1], vec![3]); + } +} diff --git a/lib/vey-ftp-client/src/config/mod.rs b/lib/vey-ftp-client/src/config/mod.rs index b32f4d33c..66c33e3ef 100644 --- a/lib/vey-ftp-client/src/config/mod.rs +++ b/lib/vey-ftp-client/src/config/mod.rs @@ -72,3 +72,30 @@ impl FtpTransferConfig { self.list_all_timeout = timeout.min(MAXIMUM_LIST_ALL_TIMEOUT); } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn defaults() { + let cfg = FtpClientConfig::default(); + assert_eq!(cfg.connect_timeout, Duration::from_secs(30)); + assert_eq!(cfg.greeting_timeout, Duration::from_secs(10)); + assert!(cfg.always_try_epsv); + assert_eq!(cfg.control.max_line_len, 2048); + assert_eq!(cfg.control.max_multi_lines, 128); + assert_eq!(cfg.transfer.list_max_entries, 1024); + assert_eq!(cfg.transfer.list_max_line_len, 2048); + } + + #[test] + fn list_all_timeout_is_clamped() { + let mut transfer = FtpTransferConfig::default(); + transfer.set_list_all_timeout(Duration::from_secs(60)); + assert_eq!(transfer.list_all_timeout, Duration::from_secs(60)); + + transfer.set_list_all_timeout(Duration::from_secs(10_000)); + assert_eq!(transfer.list_all_timeout, MAXIMUM_LIST_ALL_TIMEOUT); + } +} diff --git a/lib/vey-ftp-client/src/control/command.rs b/lib/vey-ftp-client/src/control/command.rs index c77ec30f0..bf634e437 100644 --- a/lib/vey-ftp-client/src/control/command.rs +++ b/lib/vey-ftp-client/src/control/command.rs @@ -112,3 +112,38 @@ where self.send_all().await } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::FtpControlConfig; + use tokio::io::{AsyncReadExt, duplex}; + + #[test] + fn command_display() { + assert_eq!(FtpCommand::FEAT.to_string(), "FEAT"); + assert_eq!(FtpCommand::OPTS_UTF8_ON.to_string(), "OPTS UTF8 ON"); + assert_eq!(FtpCommand::TYPE_I.to_string(), "TYPE I"); + } + + #[tokio::test] + async fn send_cmd_and_cmd1() { + let (client, mut server) = duplex(256); + let mut channel = FtpControlChannel::new(client, FtpControlConfig::default()); + + channel.send_cmd(FtpCommand::FEAT).await.unwrap(); + channel + .send_cmd1(FtpCommand::USER, "anonymous") + .await + .unwrap(); + channel + .send_pre_transfer_cmd1(FtpCommand::RETR, "/a") + .await + .unwrap(); + + let mut buf = vec![0u8; 128]; + let n = server.read(&mut buf).await.unwrap(); + let sent = std::str::from_utf8(&buf[..n]).unwrap(); + assert_eq!(sent, "FEAT\r\nUSER anonymous\r\nPRET RETR /a\r\n"); + } +} diff --git a/lib/vey-ftp-client/src/control/mod.rs b/lib/vey-ftp-client/src/control/mod.rs index ebfec983e..cd0e8b854 100644 --- a/lib/vey-ftp-client/src/control/mod.rs +++ b/lib/vey-ftp-client/src/control/mod.rs @@ -624,3 +624,91 @@ where } } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::FtpControlConfig; + use tokio::io::{AsyncReadExt, AsyncWriteExt, duplex}; + + async fn pump_reply(server: &mut (impl AsyncWrite + Unpin), reply: &str) { + server.write_all(reply.as_bytes()).await.unwrap(); + } + + #[tokio::test] + async fn wait_greetings_accepts_220() { + let (client, mut server) = duplex(256); + let mut channel = FtpControlChannel::new(client, FtpControlConfig::default()); + pump_reply(&mut server, "220 Service ready.\r\n").await; + channel.wait_greetings().await.unwrap(); + } + + #[tokio::test] + async fn wait_greetings_skips_120_then_accepts_220() { + let (client, mut server) = duplex(256); + let mut channel = FtpControlChannel::new(client, FtpControlConfig::default()); + pump_reply( + &mut server, + "120 Service ready in n minutes.\r\n220 Service ready.\r\n", + ) + .await; + channel.wait_greetings().await.unwrap(); + } + + #[tokio::test] + async fn wait_greetings_rejects_421() { + let (client, mut server) = duplex(256); + let mut channel = FtpControlChannel::new(client, FtpControlConfig::default()); + pump_reply(&mut server, "421 Service not available.\r\n").await; + let err = channel.wait_greetings().await.unwrap_err(); + assert!(matches!(err, FtpCommandError::ServiceNotAvailable)); + } + + #[tokio::test] + async fn check_server_feature_parses_feat_list() { + let (client, mut server) = duplex(1024); + let mut channel = FtpControlChannel::new(client, FtpControlConfig::default()); + + let reader = tokio::spawn(async move { + let mut buf = vec![0u8; 64]; + let n = server.read(&mut buf).await.unwrap(); + assert_eq!(&buf[..n], b"FEAT\r\n"); + server + .write_all(b"211-Features:\r\n UTF8\r\n SIZE\r\n REST STREAM\r\n211 End\r\n") + .await + .unwrap(); + server + }); + + let feat = channel.check_server_feature().await.unwrap(); + assert!(feat.support_utf8_path()); + assert!(feat.support_file_size()); + assert!(feat.support_rest_stream()); + reader.await.unwrap(); + } + + #[tokio::test] + async fn read_raw_response_single_and_multiline() { + let (client, mut server) = duplex(1024); + let mut channel = FtpControlChannel::new(client, FtpControlConfig::default()); + + pump_reply(&mut server, "200 Command okay.\r\n").await; + let rsp = channel.read_raw_response().await.unwrap(); + assert_eq!(rsp.code(), 200); + assert_eq!(rsp.line_trimmed(), Some("Command okay.")); + + pump_reply(&mut server, "211-Features:\r\n SIZE\r\n211 End\r\n").await; + let rsp = channel.read_raw_response().await.unwrap(); + assert_eq!(rsp.code(), 211); + assert_eq!(rsp.lines().unwrap().len(), 3); + } + + #[tokio::test] + async fn read_raw_response_rejects_short_line() { + let (client, mut server) = duplex(64); + let mut channel = FtpControlChannel::new(client, FtpControlConfig::default()); + pump_reply(&mut server, "22\r\n").await; + let err = channel.read_raw_response().await.unwrap_err(); + assert!(matches!(err, FtpRawResponseError::InvalidLineFormat)); + } +} diff --git a/lib/vey-ftp-client/src/control/response.rs b/lib/vey-ftp-client/src/control/response.rs index c37cf01ed..be0d1fdad 100644 --- a/lib/vey-ftp-client/src/control/response.rs +++ b/lib/vey-ftp-client/src/control/response.rs @@ -310,13 +310,38 @@ mod tests { assert_eq!(rsp.code(), 220); assert_eq!(rsp.line_trimmed(), Some("Service ready")); assert!(rsp.lines().is_none()); + + let rsp = FtpRawResponse::parse_single_line("213 42 ").unwrap(); + assert_eq!(rsp.code(), 213); + assert_eq!(rsp.line_trimmed(), Some("42")); } #[test] fn parse_code_bytes_rejects_out_of_range() { - assert!(FtpRawResponse::parse_code_bytes(b'0', b'0', b'0').is_err()); - assert!(FtpRawResponse::parse_code_bytes(b'6', b'0', b'0').is_err()); - assert!(FtpRawResponse::parse_code_bytes(b'2', b'2', b'0').is_ok()); + assert!(matches!( + FtpRawResponse::parse_code_bytes(b'0', b'0', b'0'), + Err(FtpRawResponseError::InvalidReplyCode(0)) + )); + assert!(matches!( + FtpRawResponse::parse_code_bytes(b'6', b'0', b'0'), + Err(FtpRawResponseError::InvalidReplyCode(600)) + )); + assert!(matches!( + FtpRawResponse::parse_code_bytes(b'a', b'2', b'0'), + Err(FtpRawResponseError::InvalidLineFormat) + )); + assert_eq!( + FtpRawResponse::parse_code_bytes(b'1', b'0', b'0').unwrap(), + 100 + ); + assert_eq!( + FtpRawResponse::parse_code_bytes(b'5', b'9', b'9').unwrap(), + 599 + ); + assert_eq!( + FtpRawResponse::parse_code_bytes(b'2', b'2', b'0').unwrap(), + 220 + ); } #[test] @@ -331,6 +356,24 @@ mod tests { ); } + #[test] + fn parse_pasv_227_rejects_malformed() { + let no_parens = FtpRawResponse::parse_single_line("227 Entering Passive Mode").unwrap(); + assert!(no_parens.parse_pasv_227_reply().is_none()); + + let wrong_count = + FtpRawResponse::parse_single_line("227 Entering Passive Mode (1,2,3,4,5)").unwrap(); + assert!(wrong_count.parse_pasv_227_reply().is_none()); + + let non_numeric = + FtpRawResponse::parse_single_line("227 Entering Passive Mode (a,b,c,d,e,f)").unwrap(); + assert!(non_numeric.parse_pasv_227_reply().is_none()); + + let mut parser = FtpRawResponse::get_multi_line_parser("227-Entering", 4).unwrap(); + assert!(parser.feed_line("227 Done").unwrap()); + assert!(parser.finish().parse_pasv_227_reply().is_none()); + } + #[test] fn parse_epsv_229_reply() { let rsp = @@ -339,10 +382,35 @@ mod tests { assert_eq!(rsp.parse_epsv_229_reply(), Some(6446)); } + #[test] + fn parse_epsv_229_rejects_malformed() { + let wrong_delim = + FtpRawResponse::parse_single_line("229 Entering Extended Passive Mode (|1|2|3|)") + .unwrap(); + assert!(wrong_delim.parse_epsv_229_reply().is_none()); + + let empty_port = + FtpRawResponse::parse_single_line("229 Entering Extended Passive Mode (||||)").unwrap(); + assert!(empty_port.parse_epsv_229_reply().is_none()); + + let missing_close = + FtpRawResponse::parse_single_line("229 Entering Extended Passive Mode (|||6446") + .unwrap(); + assert!(missing_close.parse_epsv_229_reply().is_none()); + + let non_numeric = + FtpRawResponse::parse_single_line("229 Entering Extended Passive Mode (|||abc|)") + .unwrap(); + assert!(non_numeric.parse_epsv_229_reply().is_none()); + } + #[test] fn parse_spsv_227_reply() { let rsp = FtpRawResponse::parse_single_line("227 (abc123)").unwrap(); assert_eq!(rsp.parse_spsv_227_reply(), Some("abc123".to_owned())); + + let no_parens = FtpRawResponse::parse_single_line("227 no identifier").unwrap(); + assert!(no_parens.parse_spsv_227_reply().is_none()); } #[test] @@ -356,6 +424,8 @@ mod tests { let lines = rsp.lines().unwrap(); assert_eq!(lines.len(), 4); assert_eq!(lines[0], "Features:"); + assert_eq!(lines[1], " UTF8"); assert_eq!(lines[3], "End"); + assert!(rsp.line_trimmed().is_none()); } } diff --git a/lib/vey-ftp-client/src/facts/entry_type.rs b/lib/vey-ftp-client/src/facts/entry_type.rs index 6feb5eb8a..661b0865c 100644 --- a/lib/vey-ftp-client/src/facts/entry_type.rs +++ b/lib/vey-ftp-client/src/facts/entry_type.rs @@ -89,14 +89,16 @@ mod tests { } #[test] - fn display_roundtrip() { - for t in [ - FtpFileEntryType::File, - FtpFileEntryType::Directory, - FtpFileEntryType::CurrentDir, - FtpFileEntryType::ParentDir, - ] { - assert_eq!(format!("{t}"), t.as_str()); - } + fn is_dir_and_unknown() { + assert!(FtpFileEntryType::Directory.is_dir()); + assert!(FtpFileEntryType::CurrentDir.is_dir()); + assert!(FtpFileEntryType::ParentDir.is_dir()); + assert!(!FtpFileEntryType::File.is_dir()); + assert!(!FtpFileEntryType::Unknown.is_dir()); + + assert_eq!(FtpFileEntryType::Unknown.as_str(), "unknown"); + assert!(FtpFileEntryType::Unknown.maybe_file()); + assert!(FtpFileEntryType::File.maybe_file()); + assert!(!FtpFileEntryType::Directory.maybe_file()); } } diff --git a/lib/vey-ftp-client/src/facts/mod.rs b/lib/vey-ftp-client/src/facts/mod.rs index 7ef35a6e8..b74a920f4 100644 --- a/lib/vey-ftp-client/src/facts/mod.rs +++ b/lib/vey-ftp-client/src/facts/mod.rs @@ -135,6 +135,79 @@ mod tests { fn parse_line() { let ff = FtpFileFacts::parse_line("type=pdir;sizd=4096;modify=20210525083610;UNIX.mode=0755;UNIX.uid=0;UNIX.gid=0;unique=804g2; /").unwrap(); assert_eq!(ff.entry_type, FtpFileEntryType::ParentDir); + assert_eq!(ff.entry_path(), "/"); assert!(ff.size.is_none()); + assert!(!ff.maybe_file()); + } + + #[test] + fn parse_line_with_common_facts() { + let ff = FtpFileFacts::parse_line( + "type=file;size=1024;modify=20211201102030;create=20211101000000;media-type=text/plain; /docs/readme.txt", + ) + .unwrap(); + assert_eq!(ff.entry_type(), &FtpFileEntryType::File); + assert_eq!(ff.entry_path(), "/docs/readme.txt"); + assert_eq!(ff.size(), Some(1024)); + assert!(ff.maybe_file()); + assert_eq!( + ff.mtime().unwrap().to_rfc3339(), + "2021-12-01T10:20:30+00:00" + ); + assert_eq!( + ff.create_time.as_ref().unwrap().to_rfc3339(), + "2021-11-01T00:00:00+00:00" + ); + assert_eq!(ff.media_type().unwrap().essence_str(), "text/plain"); + } + + #[test] + fn parse_line_skips_empty_facts_and_unknown_keys() { + let ff = FtpFileFacts::parse_line("type=file;;perm=r; /a").unwrap(); + assert_eq!(ff.entry_type(), &FtpFileEntryType::File); + assert_eq!(ff.entry_path(), "/a"); + assert!(ff.size().is_none()); + } + + #[test] + fn parse_line_ignores_invalid_media_type() { + let ff = FtpFileFacts::parse_line("type=file;media-type=@@@; /a").unwrap(); + assert!(ff.media_type().is_none()); + } + + #[test] + fn parse_line_errors() { + assert!(matches!( + FtpFileFacts::parse_line("nospace"), + Err(FtpFileFactsParseError::NoSpaceDelimiter) + )); + assert!(matches!( + FtpFileFacts::parse_line("typefile; /a"), + Err(FtpFileFactsParseError::NoDelimiterInFact(_)) + )); + assert!(matches!( + FtpFileFacts::parse_line("size=abc; /a"), + Err(FtpFileFactsParseError::InvalidSize) + )); + assert!(matches!( + FtpFileFacts::parse_line("modify=not-a-time; /a"), + Err(FtpFileFactsParseError::InvalidModifyTime(_)) + )); + assert!(matches!( + FtpFileFacts::parse_line("create=not-a-time; /a"), + Err(FtpFileFactsParseError::InvalidCreateTime(_)) + )); + } + + #[test] + fn set_size_and_mtime() { + let mut ff = FtpFileFacts::new("/tmp/x"); + assert_eq!(ff.entry_path(), "/tmp/x"); + assert!(ff.maybe_file()); + ff.set_size(9); + assert_eq!(ff.size(), Some(9)); + let dt = time_val::parse_from_str("20211201102030").unwrap(); + ff.set_mtime(dt); + assert_eq!(ff.mtime().unwrap(), &dt); } } diff --git a/lib/vey-ftp-client/src/facts/time_val.rs b/lib/vey-ftp-client/src/facts/time_val.rs index 7a4e0e0ce..7d080cc4a 100644 --- a/lib/vey-ftp-client/src/facts/time_val.rs +++ b/lib/vey-ftp-client/src/facts/time_val.rs @@ -37,4 +37,12 @@ mod tests { let expected = DateTime::parse_from_rfc3339("2021-12-01T10:20:30.123+00:00").unwrap(); assert_eq!(dt, expected.with_timezone(&Utc)); } + + #[test] + fn parse_rejects_invalid() { + assert!(parse_from_str("").is_err()); + assert!(parse_from_str("2021").is_err()); + assert!(parse_from_str("not-a-timestamp").is_err()); + assert!(parse_from_str("20211301102030").is_err()); + } } diff --git a/lib/vey-ftp-client/src/feature.rs b/lib/vey-ftp-client/src/feature.rs index 2ce0692ed..71571d6a4 100644 --- a/lib/vey-ftp-client/src/feature.rs +++ b/lib/vey-ftp-client/src/feature.rs @@ -113,4 +113,34 @@ mod tests { feat.parse_and_set("REST stream"); assert!(feat.support_rest_stream()); } + + #[test] + fn feature_names_are_case_insensitive() { + let mut feat = FtpServerFeature::default(); + feat.parse_and_set("utf8"); + feat.parse_and_set("Size"); + feat.parse_and_set("mDtM"); + feat.parse_and_set("MLST type*;size*;"); + feat.parse_and_set("EpSv"); + assert!(feat.support_utf8_path()); + assert!(feat.support_file_size()); + assert!(feat.support_file_mtime()); + assert!(feat.support_machine_list()); + assert!(feat.support_epsv()); + assert!(!feat.support_spsv()); + assert!(!feat.support_pre_transfer()); + } + + #[test] + fn default_has_no_features() { + let feat = FtpServerFeature::default(); + assert!(!feat.support_utf8_path()); + assert!(!feat.support_file_size()); + assert!(!feat.support_file_mtime()); + assert!(!feat.support_rest_stream()); + assert!(!feat.support_pre_transfer()); + assert!(!feat.support_machine_list()); + assert!(!feat.support_epsv()); + assert!(!feat.support_spsv()); + } } diff --git a/lib/vey-ftp-client/src/transfer/line.rs b/lib/vey-ftp-client/src/transfer/line.rs index 0d4c8d3ab..5a66fec78 100644 --- a/lib/vey-ftp-client/src/transfer/line.rs +++ b/lib/vey-ftp-client/src/transfer/line.rs @@ -85,3 +85,95 @@ where Err(FtpLineDataReadError::TooManyLines) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::FtpTransferConfig; + use std::io::Cursor; + + struct CollectLines { + lines: Vec, + abort_after: Option, + } + + impl FtpLineDataReceiver for CollectLines { + async fn recv_line(&mut self, line: &str) { + self.lines.push(line.to_owned()); + } + + fn should_return_early(&self) -> bool { + self.abort_after.is_some_and(|n| self.lines.len() >= n) + } + } + + fn config(max_entries: usize, max_line_len: usize) -> FtpTransferConfig { + FtpTransferConfig { + list_max_entries: max_entries, + list_max_line_len: max_line_len, + ..FtpTransferConfig::default() + } + } + + #[tokio::test] + async fn read_to_end_collects_lines() { + let io = Cursor::new(b"one\ntwo\nthree\n".to_vec()); + let transfer = FtpLineDataTransfer::new(io, &config(16, 64)); + let mut receiver = CollectLines { + lines: Vec::new(), + abort_after: None, + }; + transfer.read_to_end(&mut receiver).await.unwrap(); + assert_eq!(receiver.lines, ["one\n", "two\n", "three\n"]); + } + + #[tokio::test] + async fn read_to_end_rejects_too_many_lines() { + let io = Cursor::new(b"a\nb\nc\n".to_vec()); + let transfer = FtpLineDataTransfer::new(io, &config(2, 64)); + let mut receiver = CollectLines { + lines: Vec::new(), + abort_after: None, + }; + let err = transfer.read_to_end(&mut receiver).await.unwrap_err(); + assert!(matches!(err, FtpLineDataReadError::TooManyLines)); + assert_eq!(receiver.lines.len(), 2); + } + + #[tokio::test] + async fn read_to_end_rejects_line_too_long() { + let io = Cursor::new(b"abcdefghij\n".to_vec()); + let transfer = FtpLineDataTransfer::new(io, &config(8, 4)); + let mut receiver = CollectLines { + lines: Vec::new(), + abort_after: None, + }; + let err = transfer.read_to_end(&mut receiver).await.unwrap_err(); + assert!(matches!(err, FtpLineDataReadError::LineTooLong(1))); + } + + #[tokio::test] + async fn read_to_end_aborted_by_callback() { + let io = Cursor::new(b"keep\nstop\nmore\n".to_vec()); + let transfer = FtpLineDataTransfer::new(io, &config(16, 64)); + let mut receiver = CollectLines { + lines: Vec::new(), + abort_after: Some(1), + }; + let err = transfer.read_to_end(&mut receiver).await.unwrap_err(); + assert!(matches!(err, FtpLineDataReadError::AbortedByCallback)); + assert_eq!(receiver.lines, ["keep\n"]); + } + + #[tokio::test] + async fn read_to_end_rejects_non_utf8() { + let io = Cursor::new(vec![0xff, 0xfe, b'\n']); + let transfer = FtpLineDataTransfer::new(io, &config(8, 64)); + let mut receiver = CollectLines { + lines: Vec::new(), + abort_after: None, + }; + let err = transfer.read_to_end(&mut receiver).await.unwrap_err(); + assert!(matches!(err, FtpLineDataReadError::UnsupportedEncoding)); + } +} diff --git a/lib/vey-icap-client/src/options/request.rs b/lib/vey-icap-client/src/options/request.rs index 90c7b585b..e5f71a3af 100644 --- a/lib/vey-icap-client/src/options/request.rs +++ b/lib/vey-icap-client/src/options/request.rs @@ -1,6 +1,7 @@ /* * SPDX-License-Identifier: Apache-2.0 * SPDX-FileCopyrightText: 2023-2025 ByteDance and/or its affiliates. + * SPDX-FileCopyrightText: 2026 VEY-OSS Developers. */ use std::io; @@ -20,19 +21,24 @@ impl<'a> IcapOptionsRequest<'a> { IcapOptionsRequest { config } } - async fn send(&self, writer: &mut W) -> io::Result<()> - where - W: AsyncWrite + Unpin, - { + fn build_header(&self) -> Vec { let mut header = self.config.build_options_request(); if self.config.icap_206_enable { header.put_slice(b"Allow: 204, 206\r\n"); } else { header.put_slice(b"Allow: 204\r\n"); } + // RFC 3507 ยง4.4.1: Encapsulated MUST be included in every ICAP message + header.put_slice(b"Encapsulated: null-body=0\r\n"); header.put_slice(b"\r\n"); + header + } - writer.write_all(&header).await + async fn send(&self, writer: &mut W) -> io::Result<()> + where + W: AsyncWrite + Unpin, + { + writer.write_all(&self.build_header()).await } pub(crate) async fn get_options( @@ -56,3 +62,35 @@ impl<'a> IcapOptionsRequest<'a> { Ok(options) } } + +#[cfg(test)] +mod tests { + use super::*; + use crate::IcapMethod; + use url::Url; + + #[test] + fn options_request_includes_encapsulated_null_body() { + let url = Url::parse("icap://icap.example/reqmod").unwrap(); + let config = IcapServiceConfig::new(IcapMethod::Reqmod, url).unwrap(); + let req = IcapOptionsRequest::new(&config); + let text = String::from_utf8(req.build_header()).unwrap(); + + assert!(text.starts_with("OPTIONS icap://icap.example/reqmod ICAP/1.0\r\n")); + assert!(text.contains("Allow: 204\r\n")); + assert!(text.contains("Encapsulated: null-body=0\r\n")); + assert!(text.ends_with("\r\n\r\n")); + } + + #[test] + fn options_request_allow_includes_206_when_enabled() { + let url = Url::parse("icap://icap.example/reqmod").unwrap(); + let mut config = IcapServiceConfig::new(IcapMethod::Reqmod, url).unwrap(); + config.icap_206_enable = true; + let req = IcapOptionsRequest::new(&config); + let text = String::from_utf8(req.build_header()).unwrap(); + + assert!(text.contains("Allow: 204, 206\r\n")); + assert!(text.contains("Encapsulated: null-body=0\r\n")); + } +} diff --git a/lib/vey-icap-client/src/options/response.rs b/lib/vey-icap-client/src/options/response.rs index 3ba8f9f6d..bd9f1d758 100644 --- a/lib/vey-icap-client/src/options/response.rs +++ b/lib/vey-icap-client/src/options/response.rs @@ -280,4 +280,125 @@ mod tests { Err(IcapOptionsParseError::NoServiceTagSet) )); } + + #[test] + fn parse_header_methods_accepts_matching_method() { + let mut options = IcapServiceOptions::new(IcapMethod::Reqmod); + options + .parse_header_line(b"Methods: OPTIONS, REQMOD, RESPMOD\r\n") + .unwrap(); + } + + #[test] + fn parse_header_service_and_istag() { + let mut options = IcapServiceOptions::new(IcapMethod::Respmod); + options + .parse_header_line(b"Service: Example ICAP Server 1.0\r\n") + .unwrap(); + options + .parse_header_line(b"ISTag: \"W3E4R7U9-L3E4R7U9-W3E4R7U9\"\r\n") + .unwrap(); + options + .parse_header_line(b"Service-ID: respmod-scan\r\n") + .unwrap(); + options + .parse_header_line(b"Max-Connections: 100\r\n") + .unwrap(); + assert_eq!(options.server.as_deref(), Some("Example ICAP Server 1.0")); + assert_eq!( + options.service_tag, + "\"W3E4R7U9-L3E4R7U9-W3E4R7U9\"" + ); + assert_eq!(options.service_id.as_deref(), Some("respmod-scan")); + assert_eq!(options.max_connections, Some(100)); + options.check().unwrap(); + } + + #[test] + fn parse_header_preview_rejects_invalid_value() { + let mut options = IcapServiceOptions::new(IcapMethod::Reqmod); + assert!(matches!( + options.parse_header_line(b"Preview: not-a-number\r\n"), + Err(IcapOptionsParseError::InvalidHeaderValue("Preview")) + )); + } + + #[test] + fn parse_header_max_connections_rejects_invalid_value() { + let mut options = IcapServiceOptions::new(IcapMethod::Reqmod); + assert!(matches!( + options.parse_header_line(b"Max-Connections: abc\r\n"), + Err(IcapOptionsParseError::InvalidHeaderValue("Max-Connections")) + )); + } + + #[test] + fn parse_header_encapsulated_rejects_unknown_part() { + let mut options = IcapServiceOptions::new(IcapMethod::Reqmod); + assert!(matches!( + options.parse_header_line(b"Encapsulated: req-hdr=0\r\n"), + Err(IcapOptionsParseError::InvalidHeaderValue("Encapsulated")) + )); + } + + #[test] + fn parse_header_opt_body_type_unsupported() { + let mut options = IcapServiceOptions::new(IcapMethod::Reqmod); + assert!(matches!( + options.parse_header_line(b"Opt-body-type: text/html\r\n"), + Err(IcapOptionsParseError::UnsupportedBody(_)) + )); + } + + #[test] + fn parse_status_line_rejects_non_2xx() { + let mut options = IcapServiceOptions::new(IcapMethod::Reqmod); + assert!(matches!( + options.parse_status_line(b"ICAP/1.0 500 Internal Error\r\n"), + Err(IcapOptionsParseError::RequestFailed(500, _)) + )); + } + + #[test] + fn new_expired_is_expired() { + let options = IcapServiceOptions::new_expired(IcapMethod::Options); + assert!(options.expired()); + } + + #[tokio::test] + async fn parse_full_options_response() { + use std::io::Cursor; + + let data = b"ICAP/1.0 200 OK\r\n\ +Methods: REQMOD\r\n\ +Service: test\r\n\ +ISTag: \"tag-1\"\r\n\ +Allow: 204\r\n\ +Preview: 1024\r\n\ +\r\n"; + let mut reader = Cursor::new(&data[..]); + let options = IcapServiceOptions::parse(&mut reader, IcapMethod::Reqmod, 8192) + .await + .unwrap(); + assert!(options.support_204); + assert!(!options.support_206); + assert_eq!(options.preview_size, Some(1024)); + assert_eq!(options.service_tag, "\"tag-1\""); + assert!(!options.expired()); + } + + #[tokio::test] + async fn parse_full_options_requires_istag() { + use std::io::Cursor; + + let data = b"ICAP/1.0 200 OK\r\n\ +Methods: REQMOD\r\n\ +\r\n"; + let mut reader = Cursor::new(&data[..]); + match IcapServiceOptions::parse(&mut reader, IcapMethod::Reqmod, 8192).await { + Err(IcapOptionsParseError::NoServiceTagSet) => {} + Err(e) => panic!("unexpected error: {e}"), + Ok(_) => panic!("expected NoServiceTagSet"), + } + } } diff --git a/lib/vey-icap-client/src/parse/header_line.rs b/lib/vey-icap-client/src/parse/header_line.rs index f650b0e79..4d317a693 100644 --- a/lib/vey-icap-client/src/parse/header_line.rs +++ b/lib/vey-icap-client/src/parse/header_line.rs @@ -62,4 +62,12 @@ mod tests { assert_eq!(header.name, "Service"); assert_eq!(header.value, "my-icap-server"); } + + #[test] + fn rejects_invalid_utf8() { + match HeaderLine::parse(b"Name: \xff\r\n") { + Err(IcapLineParseError::InvalidUtf8Encoding(_)) => {} + _ => panic!("expected InvalidUtf8Encoding"), + } + } } diff --git a/lib/vey-icap-client/src/parse/status_line.rs b/lib/vey-icap-client/src/parse/status_line.rs index 90dee2817..bcaa6e34b 100644 --- a/lib/vey-icap-client/src/parse/status_line.rs +++ b/lib/vey-icap-client/src/parse/status_line.rs @@ -100,4 +100,23 @@ mod tests { _ => panic!("expected InvalidStatusCode"), } } + + #[test] + fn rejects_four_digit_status_code() { + match StatusLine::parse(b"ICAP/1.0 2000 OK\r\n") { + Err(IcapLineParseError::InvalidStatusCode) => {} + _ => panic!("expected InvalidStatusCode"), + } + } + + #[test] + fn rejects_invalid_utf8_message() { + let mut buf = b"ICAP/1.0 200 ".to_vec(); + buf.push(0xff); + buf.extend_from_slice(b"\r\n"); + match StatusLine::parse(&buf) { + Err(IcapLineParseError::InvalidUtf8Encoding(_)) => {} + _ => panic!("expected InvalidUtf8Encoding"), + } + } } diff --git a/lib/vey-icap-client/src/reqmod/payload.rs b/lib/vey-icap-client/src/reqmod/payload.rs index 12faa443e..8b2759801 100644 --- a/lib/vey-icap-client/src/reqmod/payload.rs +++ b/lib/vey-icap-client/src/reqmod/payload.rs @@ -98,3 +98,88 @@ impl IcapReqmodResponsePayload { } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parse_null_body() { + assert_eq!( + IcapReqmodResponsePayload::parse("null-body=0").unwrap(), + IcapReqmodResponsePayload::NoPayload + ); + } + + #[test] + fn parse_req_hdr_with_body() { + assert_eq!( + IcapReqmodResponsePayload::parse("req-hdr=0, req-body=128").unwrap(), + IcapReqmodResponsePayload::HttpRequestWithBody(128) + ); + } + + #[test] + fn parse_req_hdr_without_body() { + assert_eq!( + IcapReqmodResponsePayload::parse("req-hdr=0, null-body=64").unwrap(), + IcapReqmodResponsePayload::HttpRequestWithoutBody(64) + ); + } + + #[test] + fn parse_res_hdr_with_body() { + assert_eq!( + IcapReqmodResponsePayload::parse("res-hdr=0, res-body=256").unwrap(), + IcapReqmodResponsePayload::HttpResponseWithBody(256) + ); + } + + #[test] + fn parse_res_hdr_without_body() { + assert_eq!( + IcapReqmodResponsePayload::parse("res-hdr=0, null-body=32").unwrap(), + IcapReqmodResponsePayload::HttpResponseWithoutBody(32) + ); + } + + #[test] + fn rejects_non_zero_hdr_offset() { + assert!(matches!( + IcapReqmodResponsePayload::parse("req-hdr=8, req-body=16"), + Err(IcapReqmodParseError::UnsupportedBody(_)) + )); + } + + #[test] + fn rejects_missing_equals() { + assert!(matches!( + IcapReqmodResponsePayload::parse("null-body"), + Err(IcapReqmodParseError::InvalidHeaderValue("Encapsulated")) + )); + } + + #[test] + fn rejects_missing_body_part() { + assert!(matches!( + IcapReqmodResponsePayload::parse("req-hdr=0"), + Err(IcapReqmodParseError::UnsupportedBody(_)) + )); + } + + #[test] + fn rejects_invalid_body_name() { + assert!(matches!( + IcapReqmodResponsePayload::parse("req-hdr=0, opt-body=10"), + Err(IcapReqmodParseError::UnsupportedBody(_)) + )); + } + + #[test] + fn rejects_invalid_body_offset() { + assert!(matches!( + IcapReqmodResponsePayload::parse("req-hdr=0, req-body=abc"), + Err(IcapReqmodParseError::UnsupportedBody(_)) + )); + } +} diff --git a/lib/vey-icap-client/src/reqmod/response.rs b/lib/vey-icap-client/src/reqmod/response.rs index 1cff664c1..9b801db5d 100644 --- a/lib/vey-icap-client/src/reqmod/response.rs +++ b/lib/vey-icap-client/src/reqmod/response.rs @@ -155,3 +155,82 @@ impl ReqmodResponse { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Cursor; + + #[tokio::test] + async fn parse_null_body_response() { + let data = b"ICAP/1.0 204 No Content\r\n\ +Encapsulated: null-body=0\r\n\ +\r\n"; + let mut reader = Cursor::new(&data[..]); + let rsp = ReqmodResponse::parse(&mut reader, 8192, &BTreeSet::new()) + .await + .unwrap(); + assert_eq!(rsp.code, 204); + assert_eq!(rsp.reason, "No Content"); + assert!(rsp.keep_alive); + assert_eq!(rsp.payload, IcapReqmodResponsePayload::NoPayload); + } + + #[tokio::test] + async fn parse_connection_close_and_req_body() { + let data = b"ICAP/1.0 200 OK\r\n\ +Connection: close\r\n\ +Encapsulated: req-hdr=0, req-body=42\r\n\ +\r\n"; + let mut reader = Cursor::new(&data[..]); + let rsp = ReqmodResponse::parse(&mut reader, 8192, &BTreeSet::new()) + .await + .unwrap(); + assert_eq!(rsp.code, 200); + assert!(!rsp.keep_alive); + assert_eq!( + rsp.payload, + IcapReqmodResponsePayload::HttpRequestWithBody(42) + ); + } + + #[tokio::test] + async fn parse_collects_shared_headers() { + let mut shared = BTreeSet::new(); + shared.insert("x-virus-id".to_string()); + let data = b"ICAP/1.0 200 OK\r\n\ +X-Virus-ID: clamav\r\n\ +X-Ignored: skip\r\n\ +Encapsulated: null-body=0\r\n\ +\r\n"; + let mut reader = Cursor::new(&data[..]); + let mut rsp = ReqmodResponse::parse(&mut reader, 8192, &shared) + .await + .unwrap(); + let headers = rsp.take_shared_headers(); + assert!(headers.contains_key(HeaderName::from_static("x-virus-id"))); + assert!(!headers.contains_key(HeaderName::from_static("x-ignored"))); + } + + #[tokio::test] + async fn parse_rejects_too_large_header() { + let data = b"ICAP/1.0 200 OK\r\nEncapsulated: null-body=0\r\n\r\n"; + let mut reader = Cursor::new(&data[..]); + match ReqmodResponse::parse(&mut reader, 8, &BTreeSet::new()).await { + Err(IcapReqmodParseError::TooLargeHeader(8)) => {} + Err(e) => panic!("unexpected error: {e}"), + Ok(_) => panic!("expected TooLargeHeader"), + } + } + + #[tokio::test] + async fn parse_rejects_remote_closed() { + let data = b""; + let mut reader = Cursor::new(&data[..]); + match ReqmodResponse::parse(&mut reader, 8192, &BTreeSet::new()).await { + Err(IcapReqmodParseError::RemoteClosed) => {} + Err(e) => panic!("unexpected error: {e}"), + Ok(_) => panic!("expected RemoteClosed"), + } + } +} diff --git a/lib/vey-icap-client/src/respmod/payload.rs b/lib/vey-icap-client/src/respmod/payload.rs index 40c4f5321..5d17017b4 100644 --- a/lib/vey-icap-client/src/respmod/payload.rs +++ b/lib/vey-icap-client/src/respmod/payload.rs @@ -68,3 +68,72 @@ impl IcapRespmodResponsePayload { } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parse_null_body() { + assert_eq!( + IcapRespmodResponsePayload::parse("null-body=0").unwrap(), + IcapRespmodResponsePayload::NoPayload + ); + } + + #[test] + fn parse_res_hdr_with_body() { + assert_eq!( + IcapRespmodResponsePayload::parse("res-hdr=0, res-body=128").unwrap(), + IcapRespmodResponsePayload::HttpResponseWithBody(128) + ); + } + + #[test] + fn parse_res_hdr_without_body() { + assert_eq!( + IcapRespmodResponsePayload::parse("res-hdr=0, null-body=64").unwrap(), + IcapRespmodResponsePayload::HttpResponseWithoutBody(64) + ); + } + + #[test] + fn rejects_req_hdr() { + assert!(matches!( + IcapRespmodResponsePayload::parse("req-hdr=0, req-body=16"), + Err(IcapRespmodParseError::UnsupportedBody(_)) + )); + } + + #[test] + fn rejects_non_zero_hdr_offset() { + assert!(matches!( + IcapRespmodResponsePayload::parse("res-hdr=8, res-body=16"), + Err(IcapRespmodParseError::UnsupportedBody(_)) + )); + } + + #[test] + fn rejects_missing_equals() { + assert!(matches!( + IcapRespmodResponsePayload::parse("null-body"), + Err(IcapRespmodParseError::InvalidHeaderValue("Encapsulated")) + )); + } + + #[test] + fn rejects_missing_body_part() { + assert!(matches!( + IcapRespmodResponsePayload::parse("res-hdr=0"), + Err(IcapRespmodParseError::UnsupportedBody(_)) + )); + } + + #[test] + fn rejects_invalid_body_offset() { + assert!(matches!( + IcapRespmodResponsePayload::parse("res-hdr=0, res-body=1x"), + Err(IcapRespmodParseError::UnsupportedBody(_)) + )); + } +} diff --git a/lib/vey-icap-client/src/respmod/response.rs b/lib/vey-icap-client/src/respmod/response.rs index b245e0476..55802333e 100644 --- a/lib/vey-icap-client/src/respmod/response.rs +++ b/lib/vey-icap-client/src/respmod/response.rs @@ -125,3 +125,60 @@ impl RespmodResponse { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use std::io::Cursor; + + #[tokio::test] + async fn parse_null_body_response() { + let data = b"ICAP/1.0 204 No Modifications\r\n\ +Encapsulated: null-body=0\r\n\ +\r\n"; + let mut reader = Cursor::new(&data[..]); + let rsp = RespmodResponse::parse(&mut reader, 8192).await.unwrap(); + assert_eq!(rsp.code, 204); + assert_eq!(rsp.reason, "No Modifications"); + assert!(rsp.keep_alive); + assert_eq!(rsp.payload, IcapRespmodResponsePayload::NoPayload); + } + + #[tokio::test] + async fn parse_connection_close_and_res_body() { + let data = b"ICAP/1.0 200 OK\r\n\ +Connection: Keep-Alive, close\r\n\ +Encapsulated: res-hdr=0, res-body=100\r\n\ +\r\n"; + let mut reader = Cursor::new(&data[..]); + let rsp = RespmodResponse::parse(&mut reader, 8192).await.unwrap(); + assert_eq!(rsp.code, 200); + assert!(!rsp.keep_alive); + assert_eq!( + rsp.payload, + IcapRespmodResponsePayload::HttpResponseWithBody(100) + ); + } + + #[tokio::test] + async fn parse_rejects_too_large_header() { + let data = b"ICAP/1.0 200 OK\r\nEncapsulated: null-body=0\r\n\r\n"; + let mut reader = Cursor::new(&data[..]); + match RespmodResponse::parse(&mut reader, 8).await { + Err(IcapRespmodParseError::TooLargeHeader(8)) => {} + Err(e) => panic!("unexpected error: {e}"), + Ok(_) => panic!("expected TooLargeHeader"), + } + } + + #[tokio::test] + async fn parse_rejects_remote_closed() { + let data = b""; + let mut reader = Cursor::new(&data[..]); + match RespmodResponse::parse(&mut reader, 8192).await { + Err(IcapRespmodParseError::RemoteClosed) => {} + Err(e) => panic!("unexpected error: {e}"), + Ok(_) => panic!("expected RemoteClosed"), + } + } +} diff --git a/lib/vey-icap-client/src/serialize/header.rs b/lib/vey-icap-client/src/serialize/header.rs index 17d52d733..2015f5395 100644 --- a/lib/vey-icap-client/src/serialize/header.rs +++ b/lib/vey-icap-client/src/serialize/header.rs @@ -57,6 +57,19 @@ mod tests { assert!(text.contains("X-Client-Port: 8080\r\n")); } + #[test] + fn add_client_addr_serializes_ipv6() { + use std::net::Ipv6Addr; + + let addr = SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 1344); + let mut buf = Vec::new(); + add_client_addr(&mut buf, addr); + + let text = String::from_utf8(buf).unwrap(); + assert!(text.contains("X-Client-IP: ::1\r\n")); + assert!(text.contains("X-Client-Port: 1344\r\n")); + } + #[test] fn add_client_username_url_encodes_and_base64_authenticated_user() { let mut buf = Vec::new(); diff --git a/lib/vey-icap-client/src/service/config/mod.rs b/lib/vey-icap-client/src/service/config/mod.rs index a7889fe97..5f8dd7138 100644 --- a/lib/vey-icap-client/src/service/config/mod.rs +++ b/lib/vey-icap-client/src/service/config/mod.rs @@ -16,7 +16,8 @@ use rustls_pki_types::ServerName; use url::Url; use vey_types::net::{ - ConnectionPoolConfig, HttpAuth, RustlsClientConfigBuilder, TcpKeepAliveConfig, UpstreamAddr, + ConnectionPoolConfig, Host, HttpAuth, RustlsClientConfigBuilder, TcpKeepAliveConfig, + UpstreamAddr, }; #[cfg(feature = "yaml")] @@ -46,9 +47,9 @@ pub struct IcapServiceConfig { impl IcapServiceConfig { pub fn new(method: IcapMethod, mut url: Url) -> anyhow::Result { - let tls_client = match url.scheme().to_ascii_lowercase().as_str() { - "icap" => None, - "icaps" => Some(RustlsClientConfigBuilder::default()), + let (tls_client, default_port) = match url.scheme().to_ascii_lowercase().as_str() { + "icap" => (None, 1344u16), + "icaps" => (Some(RustlsClientConfigBuilder::default()), 11344u16), _ => return Err(anyhow!("unsupported ICAP URL scheme: {}", url.scheme())), }; @@ -61,8 +62,12 @@ impl IcapServiceConfig { url.set_password(None) .map_err(|_| anyhow!("failed to clear password in url"))?; - let upstream = UpstreamAddr::try_from(&url) + let host = url + .host() + .ok_or_else(|| anyhow!("no host found in this url"))?; + let host = Host::try_from(host) .map_err(|e| anyhow!("failed to get upstream address from url: {e}"))?; + let upstream = UpstreamAddr::new(host, url.port().unwrap_or(default_port)); let tls_name = ServerName::try_from(upstream.host()) .map_err(|e| anyhow!("invalid ICAP server name: {e}"))?; Ok(IcapServiceConfig { @@ -128,8 +133,8 @@ impl IcapServiceConfig { fn write_header(&self, header: &mut Vec, method: &str) { let _ = write!(header, "{method} {} ICAP/1.0\r\n", self.url); - if let Some(host) = self.url.host_str() { - let _ = write!(header, "Host: {host}\r\n"); + if let Some(host) = self.url.host() { + self.write_host_header(header, host); } if let Some(user_agent) = &self.user_agent { let _ = write!(header, "User-Agent: {user_agent}\r\n"); @@ -145,4 +150,121 @@ impl IcapServiceConfig { } } } + + fn write_host_header(&self, header: &mut Vec, host: url::Host<&str>) { + let default_port = match self.url.scheme().to_ascii_lowercase().as_str() { + "icaps" => 11344, + _ => 1344, + }; + let include_port = self + .url + .port() + .is_some_and(|port| port != default_port); + + match host { + url::Host::Domain(domain) => { + if include_port { + let _ = write!(header, "Host: {domain}:{}\r\n", self.url.port().unwrap()); + } else { + let _ = write!(header, "Host: {domain}\r\n"); + } + } + url::Host::Ipv4(ip) => { + if include_port { + let _ = write!(header, "Host: {ip}:{}\r\n", self.url.port().unwrap()); + } else { + let _ = write!(header, "Host: {ip}\r\n"); + } + } + url::Host::Ipv6(ip) => { + if include_port { + let _ = write!(header, "Host: [{ip}]:{}\r\n", self.url.port().unwrap()); + } else { + let _ = write!(header, "Host: [{ip}]\r\n"); + } + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn new_rejects_unsupported_scheme() { + let url = Url::parse("http://icap.example/reqmod").unwrap(); + match IcapServiceConfig::new(IcapMethod::Reqmod, url) { + Err(err) => assert!(err.to_string().contains("unsupported ICAP URL scheme")), + Ok(_) => panic!("expected unsupported scheme error"), + } + } + + #[test] + fn new_accepts_icap_and_builds_request_header() { + let url = Url::parse("icap://icap.example:1344/reqmod").unwrap(); + let mut config = IcapServiceConfig::new(IcapMethod::Reqmod, url).unwrap(); + assert!(config.tls_client.is_none()); + assert_eq!(config.upstream.port(), 1344); + config.user_agent = Some("vey-test/1.0".to_string()); + + let header = String::from_utf8(config.build_request_header()).unwrap(); + assert!(header.starts_with("REQMOD icap://icap.example:1344/reqmod ICAP/1.0\r\n")); + assert!(header.contains("Host: icap.example\r\n")); + assert!(header.contains("User-Agent: vey-test/1.0\r\n")); + + let options = String::from_utf8(config.build_options_request()).unwrap(); + assert!(options.starts_with("OPTIONS icap://icap.example:1344/reqmod ICAP/1.0\r\n")); + } + + #[test] + fn new_uses_default_icap_port() { + let url = Url::parse("icap://icap.example/reqmod").unwrap(); + let config = IcapServiceConfig::new(IcapMethod::Reqmod, url).unwrap(); + assert_eq!(config.upstream.port(), 1344); + assert_eq!(config.url.to_string(), "icap://icap.example/reqmod"); + } + + #[test] + fn new_accepts_icaps_with_default_port() { + let url = Url::parse("icaps://secure.example/respmod").unwrap(); + let config = IcapServiceConfig::new(IcapMethod::Respmod, url).unwrap(); + assert!(config.tls_client.is_some()); + assert_eq!(config.upstream.port(), 11344); + } + + #[test] + fn new_with_basic_auth_adds_authorization() { + let url = Url::parse("icap://user:secret@icap.example/reqmod").unwrap(); + let config = IcapServiceConfig::new(IcapMethod::Reqmod, url).unwrap(); + assert_eq!(config.upstream.port(), 1344); + let header = String::from_utf8(config.build_request_header()).unwrap(); + assert!(header.contains("Authorization: Basic ")); + assert!(!header.contains("user:secret@")); + } + + #[test] + fn host_header_includes_non_default_port() { + let url = Url::parse("icap://icap.example:2000/reqmod").unwrap(); + let config = IcapServiceConfig::new(IcapMethod::Reqmod, url).unwrap(); + let header = String::from_utf8(config.build_request_header()).unwrap(); + assert!(header.contains("Host: icap.example:2000\r\n")); + } + + #[test] + fn host_header_omits_default_icaps_port() { + let url = Url::parse("icaps://secure.example:11344/respmod").unwrap(); + let config = IcapServiceConfig::new(IcapMethod::Respmod, url).unwrap(); + let header = String::from_utf8(config.build_request_header()).unwrap(); + assert!(header.contains("Host: secure.example\r\n")); + assert!(!header.contains("Host: secure.example:11344\r\n")); + } + + #[test] + fn host_header_brackets_ipv6_with_non_default_port() { + let url = Url::parse("icap://[2001:db8::1]:2000/reqmod").unwrap(); + let config = IcapServiceConfig::new(IcapMethod::Reqmod, url).unwrap(); + let header = String::from_utf8(config.build_request_header()).unwrap(); + assert!(header.contains("Host: [2001:db8::1]:2000\r\n")); + } } diff --git a/lib/vey-icap-client/src/service/config/yaml.rs b/lib/vey-icap-client/src/service/config/yaml.rs index 80331544d..2208e37ef 100644 --- a/lib/vey-icap-client/src/service/config/yaml.rs +++ b/lib/vey-icap-client/src/service/config/yaml.rs @@ -155,6 +155,14 @@ mod tests { ); assert!(config.tls_client.is_some()); + // Default ports when omitted + let yaml = yaml_str!("icap://example.com/service"); + let config = IcapServiceConfig::parse_reqmod_service_yaml(&yaml, None).unwrap(); + assert_eq!(config.upstream.port(), 1344); + let yaml = yaml_str!("icaps://secure.example.com/service"); + let config = IcapServiceConfig::parse_reqmod_service_yaml(&yaml, None).unwrap(); + assert_eq!(config.upstream.port(), 11344); + // Invalid URL format let yaml = yaml_str!("invalid-url"); assert!(IcapServiceConfig::parse_reqmod_service_yaml(&yaml, None).is_err()); diff --git a/lib/vey-icap-client/src/service/mod.rs b/lib/vey-icap-client/src/service/mod.rs index 2d01b29ff..6ed1b485a 100644 --- a/lib/vey-icap-client/src/service/mod.rs +++ b/lib/vey-icap-client/src/service/mod.rs @@ -32,3 +32,15 @@ impl IcapMethod { } } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn method_as_str() { + assert_eq!(IcapMethod::Options.as_str(), "OPTIONS"); + assert_eq!(IcapMethod::Reqmod.as_str(), "REQMOD"); + assert_eq!(IcapMethod::Respmod.as_str(), "RESPMOD"); + } +} diff --git a/lib/vey-imap-proto/Cargo.toml b/lib/vey-imap-proto/Cargo.toml index 45c2e1468..e78fcf557 100644 --- a/lib/vey-imap-proto/Cargo.toml +++ b/lib/vey-imap-proto/Cargo.toml @@ -14,3 +14,6 @@ ahash.workspace = true log.workspace = true tokio = { workspace = true, features = ["io-util"] } vey-io-ext.workspace = true + +[dev-dependencies] +tokio = { workspace = true, features = ["macros", "io-util", "rt"] } diff --git a/lib/vey-imap-proto/src/command/mod.rs b/lib/vey-imap-proto/src/command/mod.rs index 30079913c..299f6d9e8 100644 --- a/lib/vey-imap-proto/src/command/mod.rs +++ b/lib/vey-imap-proto/src/command/mod.rs @@ -279,9 +279,19 @@ impl Command { mod tests { use super::*; + fn parse(line: &[u8]) -> Command { + Command::parse_line(line).unwrap() + } + + fn assert_parsed(line: &[u8], expected: ParsedCommand) -> Command { + let cmd = parse(line); + assert_eq!(cmd.parsed, expected); + cmd + } + #[test] fn capability() { - let cmd = Command::parse_line(b"a441 CAPABILITY\r\n").unwrap(); + let cmd = parse(b"a441 CAPABILITY\r\n"); assert_eq!(cmd.tag.as_str(), "a441"); assert_eq!(cmd.parsed, ParsedCommand::Capability); assert!(cmd.literal_arg.is_none()); @@ -289,14 +299,14 @@ mod tests { #[test] fn append() { - let cmd = Command::parse_line(b"A003 APPEND saved-messages (\\Seen) {326}\r\n").unwrap(); + let cmd = parse(b"A003 APPEND saved-messages (\\Seen) {326}\r\n"); assert_eq!(cmd.tag.as_str(), "A003"); assert_eq!(cmd.parsed, ParsedCommand::Append); let literal = cmd.literal_arg.unwrap(); assert!(literal.wait_continuation); assert_eq!(literal.size, 326); - let cmd = Command::parse_line(b"A003 APPEND saved-messages (\\Seen) {297+}\r\n").unwrap(); + let cmd = parse(b"A003 APPEND saved-messages (\\Seen) {297+}\r\n"); assert_eq!(cmd.tag.as_str(), "A003"); assert_eq!(cmd.parsed, ParsedCommand::Append); let literal = cmd.literal_arg.unwrap(); @@ -306,7 +316,7 @@ mod tests { #[test] fn enable() { - let cmd = Command::parse_line(b"A001 ENABLE CONDSTORE\r\n").unwrap(); + let cmd = parse(b"A001 ENABLE CONDSTORE\r\n"); assert_eq!(cmd.tag.as_str(), "A001"); assert_eq!(cmd.parsed, ParsedCommand::Enable); assert!(cmd.literal_arg.is_none()); @@ -314,28 +324,211 @@ mod tests { #[test] fn login() { - let cmd = Command::parse_line(b"A002 LOGIN user pass\r\n").unwrap(); - assert_eq!(cmd.parsed, ParsedCommand::Login); + assert_parsed(b"A002 LOGIN user pass\r\n", ParsedCommand::Login); } #[test] fn noop() { - let cmd = Command::parse_line(b"A003 NOOP\r\n").unwrap(); - assert_eq!(cmd.parsed, ParsedCommand::NoOperation); + assert_parsed(b"A003 NOOP\r\n", ParsedCommand::NoOperation); + } + + #[test] + fn commands_without_params() { + assert_parsed(b"A001 LOGOUT\r\n", ParsedCommand::Logout); + assert_parsed(b"A002 STARTTLS\r\n", ParsedCommand::StartTls); + assert_parsed(b"A003 NAMESPACE\r\n", ParsedCommand::Namespace); + assert_parsed(b"A004 IDLE\r\n", ParsedCommand::Idle); + assert_parsed(b"A005 CLOSE\r\n", ParsedCommand::Close); + assert_parsed(b"A006 UNSELECT\r\n", ParsedCommand::Unselect); + assert_parsed(b"A007 EXPUNGE\r\n", ParsedCommand::Expunge); + assert_parsed(b"A008 LANGUAGE\r\n", ParsedCommand::Language); + assert_parsed(b"A009 COMPARATOR\r\n", ParsedCommand::Comparator); + assert_parsed(b"A010 UNAUTHENTICATE\r\n", ParsedCommand::UnAuthenticate); + assert_parsed(b"A011 RESETKEY\r\n", ParsedCommand::ResetKey); + // CONVERSIONS requires source/target MIME args (RFC 5259); bare form is Unknown + // so the proxy can reply tagged BAD rather than forwarding. + assert_parsed(b"A012 CONVERSIONS\r\n", ParsedCommand::Unknown); + } + + #[test] + fn commands_with_params() { + assert_parsed(b"A001 AUTHENTICATE PLAIN\r\n", ParsedCommand::Auth); + assert_parsed(b"A002 SELECT INBOX\r\n", ParsedCommand::Select); + assert_parsed(b"A003 EXAMINE INBOX\r\n", ParsedCommand::Examine); + assert_parsed(b"A004 CREATE foo\r\n", ParsedCommand::Create); + assert_parsed(b"A005 DELETE foo\r\n", ParsedCommand::Delete); + assert_parsed(b"A006 RENAME foo bar\r\n", ParsedCommand::Rename); + assert_parsed(b"A007 SUBSCRIBE foo\r\n", ParsedCommand::Subscribe); + assert_parsed(b"A008 UNSUBSCRIBE foo\r\n", ParsedCommand::Unsubscribe); + assert_parsed(b"A009 LIST \"\" *\r\n", ParsedCommand::List); + assert_parsed(b"A010 LSUB \"\" *\r\n", ParsedCommand::Lsub); + assert_parsed(b"A011 STATUS INBOX (MESSAGES)\r\n", ParsedCommand::Status); + assert_parsed(b"A012 SEARCH ALL\r\n", ParsedCommand::Search); + assert_parsed(b"A013 FETCH 1:* FLAGS\r\n", ParsedCommand::Fetch); + assert_parsed(b"A014 STORE 1 +FLAGS (\\Seen)\r\n", ParsedCommand::Store); + assert_parsed(b"A015 COPY 1:3 archive\r\n", ParsedCommand::Copy); + assert_parsed(b"A016 MOVE 1:3 archive\r\n", ParsedCommand::Move); + assert_parsed(b"A017 UID FETCH 1:* FLAGS\r\n", ParsedCommand::Uid); + assert_parsed(b"A018 ID NIL\r\n", ParsedCommand::Id); + assert_parsed(b"A019 CANCELUPDATE 1\r\n", ParsedCommand::CancelUpdate); + assert_parsed(b"A020 SORT (DATE) UTF-8 ALL\r\n", ParsedCommand::Sort); + assert_parsed( + b"A021 THREAD REFERENCES UTF-8 ALL\r\n", + ParsedCommand::Thread, + ); + assert_parsed(b"A022 LANGUAGE en\r\n", ParsedCommand::Language); + assert_parsed( + b"A023 COMPARATOR \"i;unicode-casemap\"\r\n", + ParsedCommand::Comparator, + ); + assert_parsed( + b"A024 ESEARCH IN (mailboxes) RETURN (ALL) ALL\r\n", + ParsedCommand::Esearch, + ); + assert_parsed(b"A025 GETQUOTA \"\"\r\n", ParsedCommand::GetQuota); + assert_parsed(b"A026 GETQUOTAROOT INBOX\r\n", ParsedCommand::GetQuotaRoot); + assert_parsed( + b"A027 SETQUOTA \"\" (STORAGE 512)\r\n", + ParsedCommand::SetQuota, + ); + assert_parsed(b"A028 GETACL INBOX\r\n", ParsedCommand::GetAcl); + assert_parsed(b"A029 DELETEACL INBOX user\r\n", ParsedCommand::DeleteAcl); + assert_parsed( + b"A030 SETACL INBOX user lrswipkxtecdan\r\n", + ParsedCommand::SetAcl, + ); + assert_parsed(b"A031 LISTRIGHTS INBOX user\r\n", ParsedCommand::ListRights); + assert_parsed(b"A032 MYRIGHTS INBOX\r\n", ParsedCommand::MyRights); + assert_parsed( + b"A033 CONVERSIONS image/jpeg\r\n", + ParsedCommand::Conversions, + ); + assert_parsed( + b"A034 CONVERT 1 image/jpeg image/png\r\n", + ParsedCommand::Convert, + ); + assert_parsed( + b"A035 GETMETADATA INBOX (/shared)\r\n", + ParsedCommand::GetMetadata, + ); + assert_parsed( + b"A036 SETMETADATA INBOX (/shared \"x\")\r\n", + ParsedCommand::SetMetadata, + ); + assert_parsed(b"A037 NOTIFY NONE\r\n", ParsedCommand::Notify); + assert_parsed(b"A038 RESETKEY INBOX\r\n", ParsedCommand::ResetKey); + assert_parsed( + b"A039 GENURLAUTH imap://x INTERNAL\r\n", + ParsedCommand::GenUrlAuth, + ); + assert_parsed(b"A040 URLFETCH imap://x\r\n", ParsedCommand::UrlFetch); + } + + #[test] + fn command_names_are_case_insensitive() { + assert_parsed(b"a1 capability\r\n", ParsedCommand::Capability); + assert_parsed(b"a2 NoOp\r\n", ParsedCommand::NoOperation); + assert_parsed(b"a3 login USER PASS\r\n", ParsedCommand::Login); + assert_parsed(b"a4 UnSubscribe foo\r\n", ParsedCommand::Unsubscribe); + assert_parsed(b"a5 eNaBlE CONDSTORE\r\n", ParsedCommand::Enable); + } + + #[test] + fn unknown_commands() { + assert_parsed(b"A001 FOOBAR\r\n", ParsedCommand::Unknown); + assert_parsed(b"A002 FOOBAR arg\r\n", ParsedCommand::Unknown); + } + + #[test] + fn display_includes_parsed_and_tag() { + let cmd = parse(b"A001 NOOP\r\n"); + assert_eq!(cmd.to_string(), "NoOperation/A001"); + } + + #[test] + fn parse_continue_line_updates_literal() { + let mut cmd = parse(b"A003 APPEND m () {10}\r\n"); + assert_eq!(cmd.literal_arg.unwrap().size, 10); + + cmd.parse_continue_line(b"{20+}\r\n").unwrap(); + let literal = cmd.literal_arg.unwrap(); + assert_eq!(literal.size, 20); + assert!(!literal.wait_continuation); + + cmd.parse_continue_line(b"\r\n").unwrap(); + assert!(cmd.literal_arg.is_none()); } #[test] fn missing_crlf_rejected() { - assert!(Command::parse_line(b"A003 NOOP").is_err()); + assert!(matches!( + Command::parse_line(b"A003 NOOP"), + Err(CommandLineError::NoTrailingSequence) + )); } #[test] fn missing_tag_rejected() { - assert!(Command::parse_line(b"NOOP\r\n").is_err()); + assert!(matches!( + Command::parse_line(b"NOOP\r\n"), + Err(CommandLineError::NotTagPrefixed) + )); + assert!(matches!( + Command::parse_line(b"A001 \r\n"), + Err(CommandLineError::NotTagPrefixed) + )); + assert!(matches!( + Command::parse_line(b"\r\n"), + Err(CommandLineError::NotTagPrefixed) + )); + } + + #[test] + fn invalid_utf8_tag_or_command_rejected() { + assert!(matches!( + Command::parse_line(b"A\xff01 NOOP\r\n"), + Err(CommandLineError::InvalidUtf8Command(_)) + )); + assert!(matches!( + Command::parse_line(b"A001 N\xffOP\r\n"), + Err(CommandLineError::InvalidUtf8Command(_)) + )); + assert!(matches!( + Command::parse_line(b"A001 L\xffGIN user\r\n"), + Err(CommandLineError::InvalidUtf8Command(_)) + )); } #[test] fn invalid_literal_size() { - assert!(Command::parse_line(b"A003 APPEND m () {abc}\r\n").is_err()); + assert!(matches!( + Command::parse_line(b"A003 APPEND m () {abc}\r\n"), + Err(CommandLineError::InvalidLiteralFormat) + )); + assert!(matches!( + Command::parse_line(b"A003 APPEND m () {}\r\n"), + Err(CommandLineError::InvalidLiteralFormat) + )); + assert!(matches!( + Command::parse_line(b"A003 APPEND m () {12x}\r\n"), + Err(CommandLineError::InvalidLiteralFormat) + )); + assert!(matches!( + Command::parse_line(b"A003 APPEND m () {12+x}\r\n"), + Err(CommandLineError::InvalidLiteralFormat) + )); + } + + #[test] + fn continue_line_rejects_bad_input() { + let mut cmd = parse(b"A003 APPEND m () {10}\r\n"); + assert!(matches!( + cmd.parse_continue_line(b"{10}"), + Err(CommandLineError::NoTrailingSequence) + )); + assert!(matches!( + cmd.parse_continue_line(b"{abc}\r\n"), + Err(CommandLineError::InvalidLiteralFormat) + )); } } diff --git a/lib/vey-imap-proto/src/pipeline.rs b/lib/vey-imap-proto/src/pipeline.rs index 12b51f7ac..af1edf5e4 100644 --- a/lib/vey-imap-proto/src/pipeline.rs +++ b/lib/vey-imap-proto/src/pipeline.rs @@ -99,6 +99,17 @@ mod tests { } } + #[test] + fn default_and_capacity() { + let mut pipeline = CommandPipeline::default(); + assert!(pipeline.ongoing_command().is_none()); + assert!(pipeline.ongoing_response().is_none()); + + let mut pipeline = CommandPipeline::with_capacity(4); + assert!(pipeline.take_ongoing_command().is_none()); + assert!(pipeline.take_ongoing_response().is_none()); + } + #[test] fn insert_and_remove_completed() { let mut pipeline = CommandPipeline::new(); @@ -108,6 +119,44 @@ mod tests { assert!(pipeline.remove(&SmolStr::from("A001")).is_none()); } + #[test] + fn insert_completed_replaces_existing() { + let mut pipeline = CommandPipeline::new(); + pipeline.insert_completed(Command { + tag: SmolStr::from("A001"), + parsed: ParsedCommand::NoOperation, + literal_arg: None, + }); + let old = pipeline.insert_completed(Command { + tag: SmolStr::from("A001"), + parsed: ParsedCommand::Logout, + literal_arg: None, + }); + assert_eq!(old.unwrap().parsed, ParsedCommand::NoOperation); + assert_eq!( + pipeline.remove(&SmolStr::from("A001")).unwrap().parsed, + ParsedCommand::Logout + ); + } + + #[test] + fn remove_prefers_completed_over_ongoing() { + let mut pipeline = CommandPipeline::new(); + pipeline.insert_completed(tagged_command("A001")); + pipeline.set_ongoing_command(Command { + tag: SmolStr::from("A001"), + parsed: ParsedCommand::Logout, + literal_arg: None, + }); + + let removed = pipeline.remove(&SmolStr::from("A001")).unwrap(); + assert_eq!(removed.parsed, ParsedCommand::NoOperation); + assert_eq!( + pipeline.ongoing_command().unwrap().parsed, + ParsedCommand::Logout + ); + } + #[test] fn remove_ongoing_command() { let mut pipeline = CommandPipeline::new(); @@ -119,6 +168,20 @@ mod tests { assert!(pipeline.ongoing_command().is_some()); } + #[test] + fn ongoing_command_lifecycle() { + let mut pipeline = CommandPipeline::new(); + assert!(pipeline.ongoing_command().is_none()); + + pipeline.set_ongoing_command(tagged_command("C001")); + assert_eq!(pipeline.ongoing_command().unwrap().tag.as_str(), "C001"); + assert_eq!( + pipeline.take_ongoing_command().unwrap().tag.as_str(), + "C001" + ); + assert!(pipeline.take_ongoing_command().is_none()); + } + #[test] fn ongoing_response_lifecycle() { let mut pipeline = CommandPipeline::new(); diff --git a/lib/vey-imap-proto/src/response/bad.rs b/lib/vey-imap-proto/src/response/bad.rs index 565482e29..4a1d02dde 100644 --- a/lib/vey-imap-proto/src/response/bad.rs +++ b/lib/vey-imap-proto/src/response/bad.rs @@ -28,3 +28,29 @@ impl BadResponse { writer.write_all_flush(message.as_bytes()).await } } + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn bad_replies_include_tag() { + let mut buf = Vec::new(); + BadResponse::reply_invalid_command(&mut buf, "A001") + .await + .unwrap(); + assert_eq!( + std::str::from_utf8(&buf).unwrap(), + "A001 BAD invalid command\r\n" + ); + + buf.clear(); + BadResponse::reply_append_blocked(&mut buf, "A002") + .await + .unwrap(); + assert_eq!( + std::str::from_utf8(&buf).unwrap(), + "A002 BAD the message is blocked\r\n" + ); + } +} diff --git a/lib/vey-imap-proto/src/response/bye.rs b/lib/vey-imap-proto/src/response/bye.rs index c5a7df720..e56151a63 100644 --- a/lib/vey-imap-proto/src/response/bye.rs +++ b/lib/vey-imap-proto/src/response/bye.rs @@ -41,3 +41,28 @@ impl ByeResponse { impl_method!(reply_upstream_io_error, BYE_UPSTREAM_IO_ERROR); impl_method!(reply_client_protocol_error, BYE_CLIENT_PROTOCOL_ERROR); } + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn bye_replies() { + macro_rules! assert_reply { + ($method:ident, $expected:expr) => {{ + let mut buf = Vec::new(); + ByeResponse::$method(&mut buf).await.unwrap(); + assert_eq!(std::str::from_utf8(&buf).unwrap(), $expected); + }}; + } + + assert_reply!(reply_blocked, BYE_BLOCKED); + assert_reply!(reply_idle_logout, BYE_AUTO_LOGOUT); + assert_reply!(reply_server_quit, BYE_SERVER_QUIT); + assert_reply!(reply_internal_error, BYE_INTERNAL_ERROR); + assert_reply!(reply_upstream_timeout, BYE_UPSTREAM_TIMEOUT); + assert_reply!(reply_upstream_protocol_error, BYE_UPSTREAM_PROTOCOL_ERROR); + assert_reply!(reply_upstream_io_error, BYE_UPSTREAM_IO_ERROR); + assert_reply!(reply_client_protocol_error, BYE_CLIENT_PROTOCOL_ERROR); + } +} diff --git a/lib/vey-imap-proto/src/response/mod.rs b/lib/vey-imap-proto/src/response/mod.rs index a119d8080..cf1b2abac 100644 --- a/lib/vey-imap-proto/src/response/mod.rs +++ b/lib/vey-imap-proto/src/response/mod.rs @@ -317,52 +317,309 @@ fn check_literal_size(left: &[u8]) -> Result, ResponseLineError> { mod tests { use super::*; + fn parse(line: &[u8]) -> Response { + Response::parse_line(line).unwrap() + } + + fn command_data(line: &[u8]) -> UntaggedResponse { + let Response::CommandData(r) = parse(line) else { + panic!("expected command data for {line:?}") + }; + r + } + + fn server_status(line: &[u8]) -> ServerStatus { + let Response::ServerStatus(status) = parse(line) else { + panic!("expected server status for {line:?}") + }; + status + } + + fn tagged(line: &[u8]) -> TaggedResponse { + let Response::CommandResult(r) = parse(line) else { + panic!("expected tagged response for {line:?}") + }; + r + } + #[test] fn bye() { - let rsp = Response::parse_line(b"* BYE Autologout; idle for too long\r\n").unwrap(); - let Response::ServerStatus(status) = rsp else { - panic!("parse failed") - }; - assert_eq!(status, ServerStatus::Close); + assert_eq!( + server_status(b"* BYE Autologout; idle for too long\r\n"), + ServerStatus::Close + ); } #[test] fn capability() { - let rsp = Response::parse_line( + let r = command_data( b"* CAPABILITY STARTTLS AUTH=GSSAPI IMAP4rev2 LOGINDISABLED XPIG-LATIN\r\n", - ) - .unwrap(); - let Response::CommandData(r) = rsp else { - panic!("parse failed") - }; + ); assert_eq!(r.command_data, CommandData::Capability); assert!(r.literal_data.is_none()); } #[test] fn exists() { - let rsp = Response::parse_line(b"* 23 EXISTS\r\n").unwrap(); - let Response::CommandData(r) = rsp else { - panic!("parse failed") - }; + let r = command_data(b"* 23 EXISTS\r\n"); assert_eq!(r.command_data, CommandData::Other); assert!(r.literal_data.is_none()); } #[test] fn fetch() { - let rsp = Response::parse_line(b"* 12 FETCH (BODY[HEADER] {342}\r\n").unwrap(); - let Response::CommandData(r) = rsp else { - panic!("parse failed") - }; + let r = command_data(b"* 12 FETCH (BODY[HEADER] {342}\r\n"); assert_eq!(r.command_data, CommandData::Fetch); assert_eq!(r.literal_data, Some(342)); - let rsp = Response::parse_line(b"* 12 FETCH (BODY[HEADER] {342+}\r\n").unwrap(); - let Response::CommandData(r) = rsp else { - panic!("parse failed") - }; + let r = command_data(b"* 12 FETCH (BODY[HEADER] {342+}\r\n"); assert_eq!(r.command_data, CommandData::Fetch); assert_eq!(r.literal_data, Some(342)); + + let r = command_data(b"* 12 FETCH (FLAGS (\\Seen))\r\n"); + assert_eq!(r.command_data, CommandData::Fetch); + assert!(r.literal_data.is_none()); + } + + #[test] + fn tagged_results() { + let r = tagged(b"A001 OK CAPABILITY completed\r\n"); + assert_eq!(r.tag.as_str(), "A001"); + assert_eq!(r.result, CommandResult::Success); + + let r = tagged(b"A002 NO login failed\r\n"); + assert_eq!(r.result, CommandResult::Fail); + + let r = tagged(b"A003 BAD Command unknown\r\n"); + assert_eq!(r.result, CommandResult::ProtocolError); + + let r = tagged(b"a004 ok done\r\n"); + assert_eq!(r.tag.as_str(), "a004"); + assert_eq!(r.result, CommandResult::Success); + } + + #[test] + fn continuation_request() { + assert!(matches!( + parse(b"+ Ready for literal data\r\n"), + Response::ContinuationRequest + )); + assert!(matches!(parse(b"+ \r\n"), Response::ContinuationRequest)); + } + + #[test] + fn server_status_variants() { + assert_eq!( + server_status(b"* OK IMAP4rev1 Service Ready\r\n"), + ServerStatus::Information + ); + assert_eq!( + server_status(b"* NO [ALERT] System too busy\r\n"), + ServerStatus::Warning + ); + assert_eq!( + server_status(b"* BAD Command line too long\r\n"), + ServerStatus::Error + ); + assert_eq!( + server_status(b"* PREAUTH welcome\r\n"), + ServerStatus::Authenticated + ); + assert_eq!( + server_status(b"* bye shutting down\r\n"), + ServerStatus::Close + ); + } + + #[test] + fn untagged_command_data() { + assert_eq!( + command_data(b"* ENABLED CONDSTORE\r\n").command_data, + CommandData::Enabled + ); + assert_eq!( + command_data(b"* ID (\"name\" \"demo\")\r\n").command_data, + CommandData::Id + ); + assert_eq!( + command_data(b"* LIST (\\Noselect) \"/\" \"\"\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* LSUB () \"/\" INBOX\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* NAMESPACE NIL NIL NIL\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* STATUS INBOX (MESSAGES 231)\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* SEARCH 2 84 882\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* SEARCH\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* ESEARCH (TAG \"A001\") ALL 1:3\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* FLAGS (\\Answered \\Flagged \\Deleted)\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* SORT 2 84 882\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* THREAD (2)(3 6 (4 23)(44 7 96))\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* LANGUAGE (EN)\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* COMPARATOR \"i;unicode-casemap\"\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* VANISHED 300:399\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* QUOTA \"\" (STORAGE 10 512)\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* QUOTAROOT INBOX \"\"\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* ACL INBOX user lrswipkxtecdan\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* LISTRIGHTS INBOX user l r s\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* MYRIGHTS INBOX lrswipkxtecdan\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* CONVERSION image/jpeg image/png\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* CONVERTED image/jpeg image/png\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* METADATA INBOX (/shared/comment \"Hi\")\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* GENURLAUTH imap://x\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* URLFETCH imap://x\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* 1 EXPUNGE\r\n").command_data, + CommandData::Other + ); + assert_eq!( + command_data(b"* 23 RECENT\r\n").command_data, + CommandData::Other + ); + } + + #[test] + fn parse_continue_line_updates_literal() { + let mut r = command_data(b"* 12 FETCH (BODY[HEADER] {10}\r\n"); + assert_eq!(r.literal_data, Some(10)); + + r.parse_continue_line(b"{20+}\r\n").unwrap(); + assert_eq!(r.literal_data, Some(20)); + + r.parse_continue_line(b"\r\n").unwrap(); + assert!(r.literal_data.is_none()); + } + + #[test] + fn malformed_lines_rejected() { + assert!(matches!( + Response::parse_line(b"* BYE done"), + Err(ResponseLineError::NoTrailingSequence) + )); + assert!(matches!( + Response::parse_line(b"OK done\r\n"), + Err(ResponseLineError::NoResultField) + )); + assert!(matches!( + Response::parse_line(b" A001 OK done\r\n"), + Err(ResponseLineError::NotTagPrefixed) + )); + assert!(matches!( + Response::parse_line(b"A001\r\n"), + Err(ResponseLineError::NotTagPrefixed) + )); + assert!(matches!( + Response::parse_line(b"A001 DONE completed\r\n"), + Err(ResponseLineError::InvalidTaggedResult) + )); + assert!(matches!( + Response::parse_line(b"A001 OK\r\n"), + Err(ResponseLineError::NoResultField) + )); + assert!(matches!( + Response::parse_line(b"* UNKNOWN stuff\r\n"), + Err(ResponseLineError::UnknownUntaggedResult) + )); + assert!(matches!( + Response::parse_line(b"* 12 UNKNOWN\r\n"), + Err(ResponseLineError::UnknownUntaggedResult) + )); + assert!(matches!( + Response::parse_line(b"* FOO\r\n"), + Err(ResponseLineError::UnknownUntaggedResult) + )); + } + + #[test] + fn invalid_literal_in_fetch_rejected() { + assert!(matches!( + Response::parse_line(b"* 12 FETCH (BODY[HEADER] {}\r\n"), + Err(ResponseLineError::InvalidLiteralSize) + )); + assert!(matches!( + Response::parse_line(b"* 12 FETCH (BODY[HEADER] {abc}\r\n"), + Err(ResponseLineError::InvalidLiteralSize) + )); + assert!(matches!( + Response::parse_line(b"* 12 FETCH (BODY[HEADER] {12x}\r\n"), + Err(ResponseLineError::InvalidLiteralSize) + )); + } + + #[test] + fn invalid_utf8_rejected() { + assert!(matches!( + Response::parse_line(b"A\xff01 OK done\r\n"), + Err(ResponseLineError::InvalidUtf8Response(_)) + )); + assert!(matches!( + Response::parse_line(b"* CAPABILIT\xff X\r\n"), + Err(ResponseLineError::InvalidUtf8Response(_)) + )); } } diff --git a/lib/vey-io-ext/src/cache/mod.rs b/lib/vey-io-ext/src/cache/mod.rs index 1a6a4af3d..7985044f0 100644 --- a/lib/vey-io-ext/src/cache/mod.rs +++ b/lib/vey-io-ext/src/cache/mod.rs @@ -83,3 +83,49 @@ pub fn create_effective_cache( let query_handle = EffectiveQueryHandle::new(query_receiver, rsp_sender); (cache_runtime, cache_handle, query_handle) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn cache_data_new_and_empty() { + let data = EffectiveCacheData::new("ok".to_string(), 30, Duration::from_secs(5)); + assert_eq!(data.inner().map(String::as_str), Some("ok")); + assert!(data.expire_at <= data.vanish_at); + + let empty = EffectiveCacheData::::empty(10, Duration::from_secs(1)); + assert!(empty.inner().is_none()); + assert!(empty.expire_at <= empty.vanish_at); + } + + #[tokio::test] + async fn query_handle_dedups_in_flight_keys() { + let (_runtime, _cache, mut query) = + create_effective_cache::(NonZeroUsize::MIN); + let key = Arc::new("k".to_string()); + assert!(query.should_send_raw_query(key.clone(), Duration::from_secs(1))); + assert!(!query.should_send_raw_query(key.clone(), Duration::from_secs(1))); + + query.send_rsp_data( + key.clone(), + EffectiveCacheData::new("v".to_string(), 1, Duration::from_secs(1)), + false, + ); + // After response, a new query for the same key may be sent again. + assert!(query.should_send_raw_query(key, Duration::from_secs(1))); + } + + #[tokio::test] + async fn fetch_cache_only_miss_returns_none() { + let (runtime, handle, query) = create_effective_cache::(NonZeroUsize::MIN); + let runtime_task = tokio::spawn(runtime); + let miss = handle + .fetch_cache_only(Arc::new("missing".into()), Duration::from_millis(200)) + .await; + assert!(miss.is_none()); + drop(handle); + drop(query); + let _ = runtime_task.await; + } +} diff --git a/lib/vey-io-ext/src/haproxy/v1.rs b/lib/vey-io-ext/src/haproxy/v1.rs index bc3ca144b..b21c670b4 100644 --- a/lib/vey-io-ext/src/haproxy/v1.rs +++ b/lib/vey-io-ext/src/haproxy/v1.rs @@ -256,4 +256,51 @@ mod tests { Err(ProxyProtocolReadError::InvalidSrcAddr) )); } + + #[test] + fn parse_tcp6_line() { + let reader = ProxyProtocolV1Reader::new(Duration::from_secs(1)); + let result = reader + .parse_buf(b"PROXY TCP6 2001:db8::1 2001:db8::2 12345 443\r\n") + .unwrap() + .unwrap(); + assert_eq!( + result.src_addr, + SocketAddr::from_str("[2001:db8::1]:12345").unwrap() + ); + assert_eq!( + result.dst_addr, + SocketAddr::from_str("[2001:db8::2]:443").unwrap() + ); + } + + #[test] + fn parse_rejects_bad_dst_port() { + let reader = ProxyProtocolV1Reader::new(Duration::from_secs(1)); + let result = reader.parse_buf(b"PROXY TCP4 192.168.1.1 10.0.0.1 80 xyz\r\n"); + assert!(matches!( + result, + Err(ProxyProtocolReadError::InvalidDstAddr) + )); + } + + #[test] + fn parse_rejects_invalid_src_ip() { + let reader = ProxyProtocolV1Reader::new(Duration::from_secs(1)); + let result = reader.parse_buf(b"PROXY TCP4 not-an-ip 10.0.0.1 80 443\r\n"); + assert!(matches!( + result, + Err(ProxyProtocolReadError::InvalidSrcAddr) + )); + } + + #[test] + fn parse_rejects_truncated_fields() { + let reader = ProxyProtocolV1Reader::new(Duration::from_secs(1)); + let result = reader.parse_buf(b"PROXY TCP4 192.168.1.1\r\n"); + assert!(matches!( + result, + Err(ProxyProtocolReadError::InvalidDstAddr) + )); + } } diff --git a/lib/vey-io-ext/src/limit/fixed_window/stream.rs b/lib/vey-io-ext/src/limit/fixed_window/stream.rs index a8d90bc47..f597756e4 100644 --- a/lib/vey-io-ext/src/limit/fixed_window/stream.rs +++ b/lib/vey-io-ext/src/limit/fixed_window/stream.rs @@ -91,5 +91,23 @@ mod tests { limit.set_advance(900); } - // TODO add reset test case + #[test] + fn reset_starts_new_window_budget() { + let mut limit = LocalStreamLimiter::new(10, 100); + assert!(limit.is_set()); + assert_eq!(limit.check(0, 100), StreamLimitAction::AdvanceBy(100)); + limit.set_advance(100); + assert!(matches!(limit.check(1, 1), StreamLimitAction::DelayFor(_))); + + limit.reset(10, 50, 2048); + assert_eq!(limit.check(2048, 40), StreamLimitAction::AdvanceBy(40)); + limit.set_advance(40); + assert_eq!(limit.check(2050, 20), StreamLimitAction::AdvanceBy(10)); + } + + #[test] + fn shift_zero_is_disabled() { + let limit = LocalStreamLimiter::new(0, 1000); + assert!(!limit.is_set()); + } } diff --git a/lib/vey-io-ext/src/limit/token_bucket/stream.rs b/lib/vey-io-ext/src/limit/token_bucket/stream.rs index c6285410c..4ec098d21 100644 --- a/lib/vey-io-ext/src/limit/token_bucket/stream.rs +++ b/lib/vey-io-ext/src/limit/token_bucket/stream.rs @@ -132,4 +132,16 @@ mod tests { limiter.release(100); assert_eq!(limiter.check(1000), StreamLimitAction::AdvanceBy(100)); } + + #[test] + fn group_and_delay_when_tokens_exhausted() { + let config = GlobalStreamSpeedLimitConfig::per_second(10); + let limiter = GlobalStreamLimiter::new(GlobalLimitGroup::User, config); + assert!(matches!(limiter.group(), GlobalLimitGroup::User)); + assert_eq!(limiter.check(10), StreamLimitAction::AdvanceBy(10)); + assert!(matches!( + limiter.check(1), + StreamLimitAction::DelayUntil(_) + )); + } } diff --git a/lib/vey-io-ext/src/stream/ext/limited_write_ext.rs b/lib/vey-io-ext/src/stream/ext/limited_write_ext.rs index 07cbc81ae..0f8916a48 100644 --- a/lib/vey-io-ext/src/stream/ext/limited_write_ext.rs +++ b/lib/vey-io-ext/src/stream/ext/limited_write_ext.rs @@ -30,3 +30,26 @@ pub trait LimitedWriteExt: AsyncWrite { } impl LimitedWriteExt for W {} + +#[cfg(test)] +mod tests { + use super::*; + use tokio::io::AsyncWriteExt; + + #[tokio::test] + async fn write_all_flush_writes_and_flushes() { + let mut writer = Vec::new(); + writer.write_all_flush(b"abcdef").await.unwrap(); + assert_eq!(writer, b"abcdef"); + } + + #[tokio::test] + async fn write_all_vectored_writes_all_slices() { + let mut writer = Vec::new(); + let bufs = [IoSlice::new(b"ab"), IoSlice::new(b"cd"), IoSlice::new(b"ef")]; + writer.write_all_vectored(bufs).await.unwrap(); + // Vec's AsyncWrite may not implement vectored specially; still should complete. + writer.flush().await.unwrap(); + assert_eq!(writer, b"abcdef"); + } +} diff --git a/lib/vey-io-ext/src/udp/mod.rs b/lib/vey-io-ext/src/udp/mod.rs index a639736ea..ead8f6b1c 100644 --- a/lib/vey-io-ext/src/udp/mod.rs +++ b/lib/vey-io-ext/src/udp/mod.rs @@ -97,3 +97,48 @@ impl LimitedUdpRelayConfig { .max(self.packet_size as usize * self.batch_count.min(DEFAULT_UDP_RELAY_BATCH_COUNT)) } } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn default_config_values() { + let cfg = LimitedUdpRelayConfig::default(); + assert_eq!(cfg.packet_size(), DEFAULT_UDP_PACKET_SIZE); + assert_eq!( + cfg.underlying_buffer_size(), + (DEFAULT_UDP_PACKET_SIZE as usize * DEFAULT_UDP_RELAY_BATCH_COUNT) + .max(DEFAULT_UDP_UNDERLYING_BUFFER_SIZE) + ); + } + + #[test] + fn set_packet_size_clamps_to_valid_range() { + let mut cfg = LimitedUdpRelayConfig::default(); + cfg.set_packet_size(1); + assert_eq!(cfg.packet_size(), MINIMUM_UDP_PACKET_SIZE); + cfg.set_packet_size(u16::MAX); + assert_eq!(cfg.packet_size(), MAXIMUM_UDP_PACKET_SIZE); + cfg.set_packet_size(1500); + assert_eq!(cfg.packet_size(), 1500); + } + + #[test] + fn set_yield_count_enforces_minimum() { + let mut cfg = LimitedUdpRelayConfig::default(); + cfg.set_yield_count(1); + assert_eq!(cfg.yield_count, MINIMUM_UDP_RELAY_YIELD_COUNT); + cfg.set_yield_count(4096); + assert_eq!(cfg.yield_count, 4096); + } + + #[test] + fn underlying_buffer_size_grows_with_packet_and_batch() { + let mut cfg = LimitedUdpRelayConfig::default(); + cfg.set_underlying_buffer_size(0); + cfg.set_packet_size(1024); + cfg.set_batch_count(4); + assert_eq!(cfg.underlying_buffer_size(), 1024 * 4); + } +} diff --git a/lib/vey-io-ext/src/udp/split.rs b/lib/vey-io-ext/src/udp/split.rs index 6876f1b2f..9f9acbb7c 100644 --- a/lib/vey-io-ext/src/udp/split.rs +++ b/lib/vey-io-ext/src/udp/split.rs @@ -159,3 +159,31 @@ impl AsyncUdpRecv for RecvHalf { self.0.poll_batch_recvmsg(cx, hdr_v) } } + +#[cfg(test)] +mod tests { + use super::*; + use tokio::net::UdpSocket; + + #[tokio::test] + async fn reunite_same_socket_succeeds() { + let socket = UdpSocket::bind("127.0.0.1:0").await.unwrap(); + let addr = socket.local_addr().unwrap(); + let (recv, send) = split(socket); + let reunited = send.reunite(recv).unwrap(); + assert_eq!(reunited.local_addr().unwrap(), addr); + } + + #[tokio::test] + async fn reunite_different_sockets_fails() { + let a = UdpSocket::bind("127.0.0.1:0").await.unwrap(); + let b = UdpSocket::bind("127.0.0.1:0").await.unwrap(); + let (_ra, sa) = split(a); + let (rb, _sb) = split(b); + let err = sa.reunite(rb).unwrap_err(); + assert_eq!( + err.to_string(), + "tried to reunite halves that are not from the same socket" + ); + } +} diff --git a/lib/vey-syslog/src/backend/mod.rs b/lib/vey-syslog/src/backend/mod.rs index 021ff7e20..bb6f8622d 100644 --- a/lib/vey-syslog/src/backend/mod.rs +++ b/lib/vey-syslog/src/backend/mod.rs @@ -150,3 +150,48 @@ impl SyslogBackendBuilder { } } } + +#[cfg(test)] +mod tests { + use super::*; + use std::net::{Ipv4Addr, Ipv6Addr}; + + #[test] + fn udp_backend_builds_and_sends_to_connected_peer() { + let server = UdpSocket::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0)).unwrap(); + let addr = server.local_addr().unwrap(); + + let backend = SyslogBackendBuilder::Udp(Some(IpAddr::V4(Ipv4Addr::LOCALHOST)), addr) + .build() + .unwrap(); + + let n = backend + .write_many(&[Bytes::from_static(b"hello-syslog")]) + .unwrap(); + assert_eq!(n, 1); + + let mut buf = [0u8; 32]; + let (nr, _) = server.recv_from(&mut buf).unwrap(); + assert_eq!(&buf[..nr], b"hello-syslog"); + } + + #[test] + fn udp_backend_default_bind_accepts_ipv6_unspecified() { + // Bind an IPv6 localhost socket when available; otherwise skip via dual-stack IPv4. + let server = match UdpSocket::bind(SocketAddr::new(IpAddr::V6(Ipv6Addr::LOCALHOST), 0)) { + Ok(s) => s, + Err(_) => return, + }; + let addr = server.local_addr().unwrap(); + assert!(SyslogBackendBuilder::Udp(None, addr).build().is_ok()); + } + + #[cfg(unix)] + #[test] + fn default_backend_is_unix() { + assert!(matches!( + SyslogBackendBuilder::default(), + SyslogBackendBuilder::Unix(None) + )); + } +} diff --git a/lib/vey-syslog/src/format/cee.rs b/lib/vey-syslog/src/format/cee.rs index c451a94bd..4ac573e06 100644 --- a/lib/vey-syslog/src/format/cee.rs +++ b/lib/vey-syslog/src/format/cee.rs @@ -1,6 +1,7 @@ /* * SPDX-License-Identifier: Apache-2.0 * SPDX-FileCopyrightText: 2023-2025 ByteDance and/or its affiliates. + * SPDX-FileCopyrightText: 2026 VEY-OSS Developers. */ use std::io; @@ -125,8 +126,17 @@ fn format_content_as_json( #[cfg(test)] mod tests { use super::*; - - use super::SyslogFormatter; + use crate::Facility; + use slog::{Level, OwnedKVList, Record, RecordLocation, RecordStatic, o}; + + fn sample_header() -> SyslogHeader { + SyslogHeader { + facility: Facility::Daemon, + hostname: None, + process: "test".into(), + pid: 42, + } + } #[test] fn cee_event_flag_constant() { @@ -145,4 +155,71 @@ mod tests { let mut formatter = FormatterRfc5424Cee::new(Some("MID-1".into()), "@cee:".into()); formatter.append_report_ts(true); } + + #[test] + fn rfc3164_cee_format_slog_emits_json_after_flag() { + static LOC: RecordLocation = RecordLocation { + file: file!(), + line: line!(), + column: 0, + module: module_path!(), + function: "", + }; + static RS: RecordStatic = RecordStatic { + location: &LOC, + tag: "", + level: Level::Info, + }; + + let formatter = FormatterRfc3164Cee::new("@cee:".into()); + let msg = format_args!("hello"); + let kv = slog::b!("n" => 1u8); + let record = Record::new(&RS, &msg, kv); + let owned: OwnedKVList = o!("src" => "cee").into(); + let mut buf = Vec::new(); + formatter + .format_slog(&mut buf, &sample_header(), &record, &owned) + .unwrap(); + + let text = String::from_utf8(buf).unwrap(); + assert!(text.contains("@cee:")); + let json = text.rsplit_once("@cee:").unwrap().1; + assert!(json.contains("\"msg\":\"hello\"")); + assert!(json.contains("\"src\":\"cee\"")); + assert!(json.contains("\"n\":1")); + } + + #[test] + fn rfc5424_cee_format_slog_emits_json_with_report_ts() { + static LOC: RecordLocation = RecordLocation { + file: file!(), + line: line!(), + column: 0, + module: module_path!(), + function: "", + }; + static RS: RecordStatic = RecordStatic { + location: &LOC, + tag: "", + level: Level::Error, + }; + + let mut formatter = FormatterRfc5424Cee::new(Some("MID".into()), "@cee:".into()); + formatter.append_report_ts(true); + let msg = format_args!("err"); + let kv = slog::b!(); + let record = Record::new(&RS, &msg, kv); + let owned: OwnedKVList = o!().into(); + let mut buf = Vec::new(); + formatter + .format_slog(&mut buf, &sample_header(), &record, &owned) + .unwrap(); + + let text = String::from_utf8(buf).unwrap(); + assert!(text.contains(" MID ")); + assert!(text.contains("@cee:")); + let json = text.rsplit_once("@cee:").unwrap().1; + assert!(json.contains("\"msg\":\"err\"")); + assert!(json.contains("\"report_ts\":")); + } } diff --git a/lib/vey-syslog/src/format/rfc3164.rs b/lib/vey-syslog/src/format/rfc3164.rs index 587ee15ff..995cfb792 100644 --- a/lib/vey-syslog/src/format/rfc3164.rs +++ b/lib/vey-syslog/src/format/rfc3164.rs @@ -346,4 +346,84 @@ mod tests { .unwrap(); assert_eq!(std::str::from_utf8(&vec).unwrap(), ", a-key=\"a-value\""); } + + #[test] + fn format_header_includes_hostname() { + let lh = SyslogHeader { + facility: Facility::User, + hostname: Some("host1".into()), + process: "app".into(), + pid: 7, + }; + + let mut buffer: Vec = Vec::new(); + let datetime = NaiveDateTime::new( + NaiveDate::from_ymd_opt(2021, 1, 2).unwrap(), + NaiveTime::from_hms_opt(3, 4, 5).unwrap(), + ); + let dt = datetime + .and_local_timezone(Local::from_offset(&FixedOffset::east_opt(0).unwrap())) + .unwrap(); + format_rfc3164_header(&mut buffer, &lh, Level::Error, &dt).unwrap(); + + let s = String::from_utf8(buffer).unwrap(); + assert_eq!(s, "<11>Jan 2 03:04:05 host1 app[7]: "); + } + + #[test] + fn format_str_char_none_and_false() { + let mut vec = Vec::new(); + let mut kv_formatter = FormatterKv(&mut vec); + kv_formatter.emit_str("s".into(), "x\"y").unwrap(); + kv_formatter.emit_char("c".into(), 'Z').unwrap(); + kv_formatter.emit_bool("b".into(), false).unwrap(); + kv_formatter.emit_none("n".into()).unwrap(); + assert_eq!( + std::str::from_utf8(&vec).unwrap(), + ", s=\"x\\\"y\", c=\"Z\", b=false" + ); + } + + #[test] + fn format_slog_appends_msg_kv_and_report_ts() { + use slog::{OwnedKVList, Record, RecordLocation, RecordStatic, o}; + + static LOC: RecordLocation = RecordLocation { + file: file!(), + line: line!(), + column: 0, + module: module_path!(), + function: "", + }; + static RS: RecordStatic = RecordStatic { + location: &LOC, + tag: "", + level: Level::Warning, + }; + + let header = SyslogHeader { + facility: Facility::Daemon, + hostname: None, + process: "test".into(), + pid: 1, + }; + let mut formatter = FormatterRfc3164::new(); + formatter.append_report_ts(true); + + let msg = format_args!("hello"); + let kv = slog::b!("k" => 2u8); + let record = Record::new(&RS, &msg, kv); + let owned: OwnedKVList = o!("host" => "local").into(); + let mut buf = Vec::new(); + formatter + .format_slog(&mut buf, &header, &record, &owned) + .unwrap(); + + let text = String::from_utf8(buf).unwrap(); + assert!(text.starts_with("<28>")); + assert!(text.contains("test[1]: hello")); + assert!(text.contains(", host=\"local\"")); + assert!(text.contains(", k=2")); + assert!(text.contains(", report_ts=")); + } } diff --git a/lib/vey-syslog/src/format/rfc5424.rs b/lib/vey-syslog/src/format/rfc5424.rs index 4f4f9c7cc..945a8a13d 100644 --- a/lib/vey-syslog/src/format/rfc5424.rs +++ b/lib/vey-syslog/src/format/rfc5424.rs @@ -392,4 +392,89 @@ mod tests { .unwrap(); assert_eq!(std::str::from_utf8(&vec).unwrap(), " a-key=\"a-value\""); } + + #[test] + fn format_header_with_hostname_and_message_id() { + let lh = SyslogHeader { + facility: Facility::Local0, + hostname: Some("edge".into()), + process: "vey".into(), + pid: 9, + }; + let datetime = DateTime::parse_from_rfc3339("2021-12-01T10:20:30Z").unwrap(); + let dt = datetime.with_timezone(&Utc); + + let mut buffer = Vec::new(); + format_rfc5424_header( + &mut buffer, + &lh, + Level::Error, + &dt, + &Some("MSGID".into()), + ) + .unwrap(); + + assert_eq!( + String::from_utf8(buffer).unwrap(), + "<131>1 2021-12-01T10:20:30.000000Z edge vey 9 MSGID " + ); + } + + #[test] + fn format_str_escapes_bracket_and_skips_none() { + let mut vec = Vec::new(); + let mut kv_formatter = FormatterKv(&mut vec); + kv_formatter.emit_str("peer".into(), "a]b").unwrap(); + kv_formatter.emit_char("c".into(), ']').unwrap(); + kv_formatter.emit_none("skip".into()).unwrap(); + assert_eq!( + std::str::from_utf8(&vec).unwrap(), + " peer=\"a\\]b\" c=\"\\]\"" + ); + } + + #[test] + fn format_slog_writes_structured_data_and_msg() { + use slog::{OwnedKVList, Record, RecordLocation, RecordStatic, o}; + + static LOC: RecordLocation = RecordLocation { + file: file!(), + line: line!(), + column: 0, + module: module_path!(), + function: "", + }; + static RS: RecordStatic = RecordStatic { + location: &LOC, + tag: "", + level: Level::Info, + }; + + let header = SyslogHeader { + facility: Facility::User, + hostname: None, + process: "test".into(), + pid: 3, + }; + let mut formatter = FormatterRfc5424::new(32473, None); + formatter.append_report_ts(true); + + let msg = format_args!("payload"); + let kv = slog::b!("code" => 1u8); + let record = Record::new(&RS, &msg, kv); + let owned: OwnedKVList = o!("src" => "unit").into(); + let mut buf = Vec::new(); + formatter + .format_slog(&mut buf, &header, &record, &owned) + .unwrap(); + + let text = String::from_utf8(buf).unwrap(); + assert!(text.starts_with("<13>1 ")); + assert!(text.contains(" - test 3 - ")); + assert!(text.contains("[vey-proxy@32473")); + assert!(text.contains(" src=\"unit\"")); + assert!(text.contains(" code=\"1\"")); + assert!(text.contains(" report_ts=\"")); + assert!(text.contains("] payload")); + } } diff --git a/lib/vey-syslog/src/format/serde.rs b/lib/vey-syslog/src/format/serde.rs index 97eb38a9b..5cb160669 100644 --- a/lib/vey-syslog/src/format/serde.rs +++ b/lib/vey-syslog/src/format/serde.rs @@ -109,3 +109,39 @@ impl Serializer for SerdeFormatterKV { Ok(()) } } + +#[cfg(test)] +mod tests { + use super::*; + use slog::Serializer; + + #[test] + fn serde_formatter_kv_emits_json_object() { + let mut buf = Vec::new(); + let mut ser = serde_json::Serializer::new(&mut buf); + let mut kv = SerdeFormatterKV::start(&mut ser, Some(3)).unwrap(); + kv.emit_str("msg".into(), "hi").unwrap(); + kv.emit_u8("n".into(), 7).unwrap(); + kv.emit_bool("ok".into(), true).unwrap(); + kv.end().unwrap(); + + let value: serde_json::Value = serde_json::from_slice(&buf).unwrap(); + assert_eq!(value["msg"], "hi"); + assert_eq!(value["n"], 7); + assert_eq!(value["ok"], true); + } + + #[test] + fn serde_formatter_kv_skips_none() { + let mut buf = Vec::new(); + let mut ser = serde_json::Serializer::new(&mut buf); + let mut kv = SerdeFormatterKV::start(&mut ser, None).unwrap(); + kv.emit_none("missing".into()).unwrap(); + kv.emit_i64("report_ts".into(), 123).unwrap(); + kv.end().unwrap(); + + let value: serde_json::Value = serde_json::from_slice(&buf).unwrap(); + assert!(value.get("missing").is_none()); + assert_eq!(value["report_ts"], 123); + } +} diff --git a/lib/vey-syslog/src/lib.rs b/lib/vey-syslog/src/lib.rs index bc8e78ade..262c57494 100644 --- a/lib/vey-syslog/src/lib.rs +++ b/lib/vey-syslog/src/lib.rs @@ -122,3 +122,60 @@ impl SyslogBuilder { AsyncSyslogStreamer::new(async_conf, header, formatter, &self.backend) } } + +#[cfg(test)] +mod tests { + use super::*; + use std::net::{IpAddr, Ipv4Addr, SocketAddr}; + + #[test] + fn builder_defaults_and_setters() { + let mut builder = SyslogBuilder::with_ident("vey"); + assert_eq!(builder.ident, "vey"); + assert!(matches!(builder.facility, Facility::User)); + assert!(matches!(builder.format, SyslogFormatterKind::Rfc3164)); + assert!(!builder.emit_hostname); + assert!(!builder.append_report_ts); + + builder.set_facility(Facility::Local0); + builder.set_format(SyslogFormatterKind::Rfc5424(32473, Some("MID".into()))); + builder.set_emit_hostname(true); + builder.append_report_ts(true); + builder.set_backend(SyslogBackendBuilder::Udp( + None, + SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 514), + )); + + assert!(matches!(builder.facility, Facility::Local0)); + assert!(builder.emit_hostname); + assert!(builder.append_report_ts); + match &builder.format { + SyslogFormatterKind::Rfc5424(eid, mid) => { + assert_eq!(*eid, 32473); + assert_eq!(mid.as_deref(), Some("MID")); + } + _ => panic!("expected Rfc5424"), + } + assert!(matches!(builder.backend, SyslogBackendBuilder::Udp(_, _))); + } + + #[test] + fn enable_cee_log_syntax_from_rfc3164_and_rfc5424() { + let mut builder = SyslogBuilder::with_ident("vey"); + builder.enable_cee_log_syntax(None); + match &builder.format { + SyslogFormatterKind::Rfc3164Cee(flag) => assert_eq!(flag, "@cee:"), + _ => panic!("expected Rfc3164Cee"), + } + + builder.set_format(SyslogFormatterKind::Rfc5424(1, Some("ID".into()))); + builder.enable_cee_log_syntax(Some("FLAG:".into())); + match &builder.format { + SyslogFormatterKind::Rfc5424Cee(mid, flag) => { + assert_eq!(mid.as_deref(), Some("ID")); + assert_eq!(flag, "FLAG:"); + } + _ => panic!("expected Rfc5424Cee"), + } + } +} diff --git a/lib/vey-syslog/src/types.rs b/lib/vey-syslog/src/types.rs index 1ad114d71..f5caff255 100644 --- a/lib/vey-syslog/src/types.rs +++ b/lib/vey-syslog/src/types.rs @@ -81,4 +81,15 @@ mod tests { assert_eq!(Severity::Alert as u8, 1); assert_eq!(Severity::Debug as u8, 7); } + + #[test] + fn facility_covers_standard_and_local_range() { + assert_eq!(Facility::Auth as u8, 4 << 3); + assert_eq!(Facility::Syslog as u8, 5 << 3); + assert_eq!(Facility::Cron as u8, 9 << 3); + assert_eq!(Facility::Ftp as u8, 11 << 3); + assert_eq!(Facility::Local0 as u8, 16 << 3); + assert_eq!(Facility::Local3 as u8, 19 << 3); + assert_eq!(Facility::Local6 as u8, 22 << 3); + } } diff --git a/sphinx/vey-values/configuration/values/audit.rst b/sphinx/vey-values/configuration/values/audit.rst index e17055cd6..c89ec0e37 100644 --- a/sphinx/vey-values/configuration/values/audit.rst +++ b/sphinx/vey-values/configuration/values/audit.rst @@ -34,6 +34,8 @@ If the value is a map, the following keys are supported: ICAP service URL. The scheme must be either ``icap`` or ``icaps``. When the scheme is ``icaps``, a default TLS client configuration is used. + The default port is ``1344`` for ``icap`` and ``11344`` for ``icaps`` and + may be omitted. * use_unix_socket diff --git a/vey-bench/CHANGELOG b/vey-bench/CHANGELOG index 6dba6ee8d..007bd1d6f 100644 --- a/vey-bench/CHANGELOG +++ b/vey-bench/CHANGELOG @@ -1,4 +1,7 @@ +v0.9.9: + - Feature: h1/h2/h3 targets support the HTTP `QUERY` method (RFC 10008), including `--payload` + v0.9.8: - BUG FIX: SOCKS4 proxy replies parse the bound port correctly again (broken SOCKS4 via proxy no longer fails handshake) - BUG FIX: h1 target records send-header time when the request has a body (header latency stats were missing before) @@ -8,7 +11,6 @@ v0.9.8: - Feature: websocket target reports Req/Conn (conn used times) like other multiplexed targets - Feature: HTTP/3 GREASE is off by default; pass `--grease` when the peer needs it - Feature: optional jemalloc or mimalloc global allocator at compile time; allocator name/version appear in `--version` - - Feature: h1/h2/h3 targets support the HTTP `QUERY` method (RFC 10008), including `--payload` - Optimization: skip noisy percentile / request-distribution summary lines when concurrency is 1 or the sample set is too small - Compatibility: OpenSSL CLI `--tls-supported-groups` / `--proxy-tls-supported-groups` renamed to `--tls-key-exchange-groups` / `--proxy-tls-key-exchange-groups` - Compatibility: bump MSRV to 1.91 diff --git a/vey-proxy/CHANGELOG b/vey-proxy/CHANGELOG index 4bac4602f..d008b220c 100644 --- a/vey-proxy/CHANGELOG +++ b/vey-proxy/CHANGELOG @@ -3,6 +3,9 @@ CHANGELOG Release notes for `vey-proxy`, ordered from newest to oldest. v1.13.10: + - BUG FIX: ICAP `Host` includes the non-default port (and brackets IPv6 literals) + - BUG FIX: ICAP OPTIONS requests include `Encapsulated: null-body=0` (RFC 3507) + - BUG FIX: ICAP/ICAPS service URLs without an explicit port use defaults 1344 / 11344 instead of failing to parse - Feature: enable transparent listen (`tcp_tproxy` / `udp_tproxy` / `sni_proxy`) and foreign bind on FreeBSD (`IP_BINDANY`) and OpenBSD (`SO_BINDANY`) - Feature: enable `direct_fixed` foreign bind on NetBSD (`IP_BINDANY`); the intercepting servers stay unavailable there as NPF gives no original destination - Feature: add FreeBSD `user_cookie` (`SO_USER_COOKIE`) and OpenBSD `rtable` (`SO_RTABLE`) listen / misc socket options alongside Linux `mark` / `netfilter_mark`