mod search;
mod transactions;
use futures::stream::{self, StreamExt};
use crate::error::Result;
use crate::models::filings::CongressionalTrade;
use super::client::SenateTradesClient;
use search::search_recent_ptrs;
use transactions::matching_transactions;
const MAX_FILINGS_SCANNED: usize = 100;
const CONCURRENT_FETCHES: usize = 8;
pub(crate) async fn fetch_congressional_trades_response(
symbol: &str,
) -> Result<Vec<CongressionalTrade>> {
let symbol = symbol.to_uppercase();
let client = SenateTradesClient::launch().await?;
let result = scan(&client, &symbol).await;
client.close().await;
result
}
async fn scan(client: &SenateTradesClient, symbol: &str) -> Result<Vec<CongressionalTrade>> {
let mut entries = search_recent_ptrs(client).await?;
entries.truncate(MAX_FILINGS_SCANNED);
let trades = stream::iter(entries)
.map(|entry| async move { matching_transactions(client, entry, symbol).await })
.buffer_unordered(CONCURRENT_FETCHES)
.collect::<Vec<_>>()
.await
.into_iter()
.flatten()
.collect();
Ok(trades)
}
#[cfg(test)]
mod tests {
#[tokio::test]
#[ignore = "requires network access and a local chromium binary"]
async fn test_live_congressional_trades() {
let trades = super::fetch_congressional_trades_response("AAPL")
.await
.unwrap();
assert!(!trades.is_empty());
}
}