use ig_client::application::config::{
Config, Credentials, DatabaseConfig, RateLimiterConfig, RestApiConfig, WebSocketConfig,
};
use ig_client::application::http::HttpClient;
use ig_client::application::rate_limiter::{RateLimitClass, RateLimiter};
use ig_client::error::AppError;
use ig_client::model::retry::RetryConfig;
use serde::Deserialize;
use std::time::{Duration, Instant};
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
#[derive(Debug, Deserialize)]
struct Dummy {
#[allow(dead_code)]
ok: bool,
}
fn pool_config(base_url: &str, keys: &str, max_requests: u32, period_seconds: u64) -> Config {
pool_config_burst(base_url, keys, max_requests, period_seconds, max_requests)
}
fn pool_config_burst(
base_url: &str,
keys: &str,
max_requests: u32,
period_seconds: u64,
burst_size: u32,
) -> Config {
Config {
credentials: Credentials {
username: "fake-user".to_string(),
password: "fake-pass".to_string(),
account_id: "ABC12".to_string(),
api_key: keys.to_string(),
client_token: None,
account_token: None,
},
rest_api: RestApiConfig {
base_url: base_url.to_string(),
timeout: 30,
},
websocket: WebSocketConfig {
url: "wss://example.invalid".to_string(),
reconnect_interval: 5,
},
database: DatabaseConfig {
url: "postgres://localhost/none".to_string(),
max_connections: 1,
},
rate_limiter: RateLimiterConfig {
max_requests,
period_seconds,
burst_size,
},
sleep_hours: 1,
page_size: 20,
days_to_look_back: 7,
api_version: Some(3),
}
}
async fn mount_login_ok(server: &MockServer) {
Mock::given(method("POST"))
.and(path("/session"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({
"clientId": "FAKE-CLIENT-1",
"accountId": "ABC12",
"timezoneOffset": 1,
"lightstreamerEndpoint": "https://example.invalid",
"oauthToken": {
"access_token": "FAKE-ACCESS",
"refresh_token": "FAKE-REFRESH",
"scope": "profile",
"token_type": "Bearer",
"expires_in": "600"
}
})))
.mount(server)
.await;
}
async fn mount_data_ok(server: &MockServer) {
Mock::given(method("GET"))
.and(path("/data"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"ok": true})))
.mount(server)
.await;
}
async fn keys_used(server: &MockServer) -> Vec<String> {
server
.received_requests()
.await
.unwrap_or_default()
.iter()
.filter(|r| r.url.path() == "/data")
.map(|r| {
r.headers
.get("X-IG-API-KEY")
.and_then(|v| v.to_str().ok())
.unwrap_or_default()
.to_string()
})
.collect()
}
#[tokio::test]
async fn test_pool_two_keys_first_two_requests_use_distinct_keys_immediately() {
let server = MockServer::start().await;
mount_login_ok(&server).await;
mount_data_ok(&server).await;
let client = HttpClient::new_lazy(pool_config_burst(&server.uri(), "key-a,key-b", 2, 60, 2))
.expect("client builds");
let started = Instant::now();
client.get::<Dummy>("/data", Some(1)).await.expect("first");
client.get::<Dummy>("/data", Some(1)).await.expect("second");
let elapsed = started.elapsed();
let used = keys_used(&server).await;
assert_eq!(used.len(), 2, "both requests were sent");
assert_ne!(used[0], used[1], "the two requests used different keys");
assert!(
elapsed < Duration::from_secs(5),
"neither request waited on a refill, took {elapsed:?}"
);
}
#[tokio::test]
async fn test_pool_third_request_waits_for_the_first_refill() {
let server = MockServer::start().await;
mount_login_ok(&server).await;
mount_data_ok(&server).await;
let client = HttpClient::new_lazy(pool_config_burst(&server.uri(), "key-a,key-b", 2, 60, 2))
.expect("client builds");
client.get::<Dummy>("/data", Some(1)).await.expect("first");
client.get::<Dummy>("/data", Some(1)).await.expect("second");
let third = tokio::time::timeout(
Duration::from_millis(300),
client.get::<Dummy>("/data", Some(1)),
)
.await;
assert!(third.is_err(), "the third request waited for a refill");
assert_eq!(
keys_used(&server).await.len(),
2,
"no third request reached the server while waiting"
);
}
#[tokio::test]
async fn test_send_with_reservation_consumes_exactly_one_token() {
let server = MockServer::start().await;
mount_data_ok(&server).await;
let limiter = RateLimiter::new(&RateLimiterConfig {
max_requests: 3,
period_seconds: 600,
burst_size: 3,
});
assert!(limiter.try_reserve(RateLimitClass::NonTrading), "token 1");
let response = ig_client::application::http::make_http_request_reserved(
&reqwest::Client::new(),
&limiter,
reqwest::Method::GET,
&format!("{}/data", server.uri()),
vec![],
&None::<()>,
RetryConfig {
max_retry_count: Some(0),
retry_delay_secs: Some(0),
},
true,
)
.await;
assert!(response.is_ok(), "the request was sent");
assert!(
limiter.try_reserve(RateLimitClass::NonTrading),
"second token still available"
);
assert!(
limiter.try_reserve(RateLimitClass::NonTrading),
"third token still available - the send did not take a second one"
);
assert!(
!limiter.try_reserve(RateLimitClass::NonTrading),
"and the budget is now genuinely spent"
);
}
#[tokio::test]
async fn test_pool_failed_login_does_not_consume_a_data_token() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/session"))
.respond_with(ResponseTemplate::new(401).set_body_json(serde_json::json!({
"errorCode": "error.security.invalid-details"
})))
.mount(&server)
.await;
mount_data_ok(&server).await;
let client = HttpClient::new_lazy(pool_config_burst(&server.uri(), "key-a", 3, 60, 3))
.expect("client builds");
let failed = client.get::<Dummy>("/data", Some(1)).await;
assert!(failed.is_err(), "the request fails because login fails");
assert!(
keys_used(&server).await.is_empty(),
"no data request was sent"
);
server.reset().await;
mount_login_ok(&server).await;
mount_data_ok(&server).await;
let started = Instant::now();
let second = tokio::time::timeout(
Duration::from_secs(5),
client.get::<Dummy>("/data", Some(1)),
)
.await;
assert!(
second.is_ok(),
"the data token was not spent by the failed login, waited {:?}",
started.elapsed()
);
}
#[tokio::test]
async fn test_pool_concurrent_requests_spread_across_keys() {
let server = MockServer::start().await;
mount_login_ok(&server).await;
mount_data_ok(&server).await;
let client = std::sync::Arc::new(
HttpClient::new_lazy(pool_config_burst(
&server.uri(),
"key-a,key-b,key-c,key-d",
2,
60,
2,
))
.expect("client builds"),
);
let mut tasks = Vec::new();
for _ in 0..4 {
let client = client.clone();
tasks.push(tokio::spawn(async move {
client
.get::<Dummy>("/data", Some(1))
.await
.map(|_: Dummy| ())
}));
}
for task in tasks {
let _ = task.await;
}
let mut used = keys_used(&server).await;
used.sort();
used.dedup();
assert_eq!(
used.len(),
4,
"four concurrent requests used four distinct keys, got {used:?}"
);
}
#[tokio::test]
async fn test_rate_limiter_waits_for_the_key_that_refills_first() {
let slow = RateLimiter::new(&RateLimiterConfig {
max_requests: 1,
period_seconds: 60,
burst_size: 1,
});
let fast = RateLimiter::new(&RateLimiterConfig {
max_requests: 1,
period_seconds: 1,
burst_size: 1,
});
assert!(slow.try_reserve(RateLimitClass::NonTrading));
assert!(fast.try_reserve(RateLimitClass::NonTrading));
let started = Instant::now();
let waits: Vec<std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send>>> = vec![
Box::pin(async move { slow.reserve(RateLimitClass::NonTrading).await }),
Box::pin(async move { fast.reserve(RateLimitClass::NonTrading).await }),
];
let (_, _, _) = futures::future::select_all(waits).await;
assert!(
started.elapsed() < Duration::from_secs(10),
"the wait ended on the fast bucket, took {:?}",
started.elapsed()
);
}
#[tokio::test]
async fn test_pool_api_key_allowance_rotates_immediately() {
let server = MockServer::start().await;
mount_login_ok(&server).await;
Mock::given(method("GET"))
.and(path("/data"))
.respond_with(ResponseTemplate::new(403).set_body_json(serde_json::json!({
"errorCode": "error.public-api.exceeded-api-key-allowance"
})))
.up_to_n_times(1)
.mount(&server)
.await;
mount_data_ok(&server).await;
let client = HttpClient::new_lazy(pool_config(&server.uri(), "key-a,key-b", 5, 1))
.expect("client builds");
let result = client.get::<Dummy>("/data", Some(1)).await;
assert!(result.is_ok(), "the request succeeded on another key");
let used = keys_used(&server).await;
assert_eq!(used.len(), 2, "one rejection plus one success");
assert_ne!(used[0], used[1], "the retry went to a different key");
}
#[tokio::test]
async fn test_pool_account_allowance_does_not_rotate() {
let server = MockServer::start().await;
mount_login_ok(&server).await;
Mock::given(method("GET"))
.and(path("/data"))
.respond_with(ResponseTemplate::new(403).set_body_json(serde_json::json!({
"errorCode": "error.public-api.exceeded-account-allowance"
})))
.mount(&server)
.await;
let client = HttpClient::new_lazy(pool_config(&server.uri(), "key-a,key-b", 5, 1))
.expect("client builds");
let err = client
.get::<Dummy>("/data", Some(1))
.await
.expect_err("account allowance is an error");
assert!(
matches!(err, AppError::AccountAllowanceExceeded),
"got {err:?}"
);
let used = keys_used(&server).await;
assert!(
used.iter().collect::<std::collections::HashSet<_>>().len() <= 1,
"the account allowance did not burn a second key, used {used:?}"
);
}
#[tokio::test]
async fn test_pool_trading_allowance_does_not_rotate() {
let server = MockServer::start().await;
mount_login_ok(&server).await;
Mock::given(method("POST"))
.and(path("/positions/otc"))
.respond_with(ResponseTemplate::new(403).set_body_json(serde_json::json!({
"errorCode": "error.public-api.exceeded-account-trading-allowance"
})))
.mount(&server)
.await;
let client = HttpClient::new_lazy(pool_config(&server.uri(), "key-a,key-b", 5, 1))
.expect("client builds");
let err = client
.post::<_, Dummy>("/positions/otc", serde_json::json!({}), Some(2))
.await
.expect_err("trading allowance is an error");
assert!(
matches!(err, AppError::TradingAllowanceExceeded),
"got {err:?}"
);
let orders: Vec<String> = server
.received_requests()
.await
.unwrap_or_default()
.iter()
.filter(|r| r.url.path() == "/positions/otc")
.map(|r| {
r.headers
.get("X-IG-API-KEY")
.and_then(|v| v.to_str().ok())
.unwrap_or_default()
.to_string()
})
.collect();
assert_eq!(orders.len(), 1, "the order was sent once, got {orders:?}");
}
#[tokio::test]
async fn test_pool_single_key_keeps_previous_behaviour() {
let server = MockServer::start().await;
mount_login_ok(&server).await;
mount_data_ok(&server).await;
let client =
HttpClient::new_lazy(pool_config(&server.uri(), "only-key", 5, 1)).expect("client builds");
client.get::<Dummy>("/data", Some(1)).await.expect("first");
client.get::<Dummy>("/data", Some(1)).await.expect("second");
let used = keys_used(&server).await;
assert_eq!(used, vec!["only-key".to_string(), "only-key".to_string()]);
}