use std::fs;
use std::thread::{self, JoinHandle};
#[cfg(feature = "http")]
use anyhow::anyhow;
use crossbeam_channel::{Receiver, Sender, unbounded};
#[cfg(feature = "http")]
use eure::query::fetch_url;
use eure::query::{TextFile, TextFileContent};
pub struct IoRequest {
pub file: TextFile,
}
pub struct IoResponse {
pub file: TextFile,
pub result: Result<TextFileContent, anyhow::Error>,
}
pub struct IoPool {
request_tx: Sender<IoRequest>,
response_rx: Receiver<IoResponse>,
_workers: Vec<JoinHandle<()>>,
}
impl IoPool {
pub fn new(num_workers: usize) -> Self {
let (request_tx, request_rx) = unbounded::<IoRequest>();
let (response_tx, response_rx) = unbounded::<IoResponse>();
let mut workers = Vec::with_capacity(num_workers);
for i in 0..num_workers {
let request_rx = request_rx.clone();
let response_tx = response_tx.clone();
let handle = thread::Builder::new()
.name(format!("eure-ls-io-{}", i))
.spawn(move || {
worker_loop(request_rx, response_tx);
})
.expect("failed to spawn IO worker thread");
workers.push(handle);
}
Self {
request_tx,
response_rx,
_workers: workers,
}
}
pub fn request_file(&self, file: TextFile) {
let _ = self.request_tx.send(IoRequest { file });
}
pub fn receiver(&self) -> &Receiver<IoResponse> {
&self.response_rx
}
}
fn worker_loop(request_rx: Receiver<IoRequest>, response_tx: Sender<IoResponse>) {
for request in request_rx {
let result = read_file(&request.file);
let response = IoResponse {
file: request.file,
result,
};
if response_tx.send(response).is_err() {
break;
}
}
}
fn read_file(file: &TextFile) -> Result<TextFileContent, anyhow::Error> {
match file {
TextFile::Local(path) => Ok(TextFileContent(fs::read_to_string(path.as_ref())?)),
#[cfg(feature = "http")]
TextFile::Remote(url) => match fetch_url(url) {
Ok(content) => Ok(content),
Err(e) => Err(anyhow!("Failed to fetch {}: {}", url, e)),
},
#[cfg(not(feature = "http"))]
TextFile::Remote(url) => Err(anyhow::anyhow!(
"HTTP support not enabled, cannot fetch {}",
url
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
#[test]
fn test_read_nonexistent_file() {
let file = TextFile::from_path(PathBuf::from("/nonexistent/path/to/file.eure"));
let result = read_file(&file);
assert!(
result
.unwrap_err()
.downcast_ref::<std::io::Error>()
.is_some()
);
}
}