Implement audited findings
Some checks failed
Dependency security audit / rustsec (push) Failing after 3s
Some checks failed
Dependency security audit / rustsec (push) Failing after 3s
This commit is contained in:
@@ -210,6 +210,7 @@ impl AsyncExecutor {
|
||||
where
|
||||
F: FnOnce() -> Result<AsyncPayload, String> + Send + 'static,
|
||||
{
|
||||
self.reap_finished();
|
||||
let sender = self.sender.clone();
|
||||
let task = thread::spawn(move || {
|
||||
let _ignored = sender.send(AsyncResult {
|
||||
@@ -223,6 +224,34 @@ impl AsyncExecutor {
|
||||
.push(task);
|
||||
}
|
||||
|
||||
fn reap_finished(&self) {
|
||||
let mut tasks = self
|
||||
.tasks
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner);
|
||||
let mut completed = Vec::new();
|
||||
let mut index = 0;
|
||||
while index < tasks.len() {
|
||||
if tasks[index].is_finished() {
|
||||
completed.push(tasks.swap_remove(index));
|
||||
} else {
|
||||
index += 1;
|
||||
}
|
||||
}
|
||||
drop(tasks);
|
||||
for task in completed {
|
||||
let _ignored = task.join();
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn retained_task_count(&self) -> usize {
|
||||
self.tasks
|
||||
.lock()
|
||||
.unwrap_or_else(std::sync::PoisonError::into_inner)
|
||||
.len()
|
||||
}
|
||||
|
||||
pub fn progress_reporter(
|
||||
&self,
|
||||
token: RequestToken,
|
||||
@@ -253,6 +282,7 @@ impl AsyncExecutor {
|
||||
}
|
||||
|
||||
pub fn drain(&self) -> impl Iterator<Item = AsyncResult> + '_ {
|
||||
self.reap_finished();
|
||||
self.receiver.try_iter()
|
||||
}
|
||||
}
|
||||
@@ -313,4 +343,28 @@ mod tests {
|
||||
.expect("worker result");
|
||||
assert_eq!(app.apply_result(result), ResultDisposition::Applied);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn completed_task_handles_are_reaped_during_normal_activity() {
|
||||
let executor = AsyncExecutor::new();
|
||||
let mut app = App::new();
|
||||
let submitted = 256;
|
||||
for _ in 0..submitted {
|
||||
let token = app.begin_request();
|
||||
executor.submit(token, || Err("complete".to_owned()));
|
||||
}
|
||||
for _ in 0..submitted {
|
||||
executor
|
||||
.receiver
|
||||
.recv_timeout(Duration::from_secs(2))
|
||||
.expect("worker result");
|
||||
}
|
||||
|
||||
let deadline = std::time::Instant::now() + Duration::from_secs(2);
|
||||
while executor.retained_task_count() != 0 && std::time::Instant::now() < deadline {
|
||||
let _ = executor.drain().count();
|
||||
thread::yield_now();
|
||||
}
|
||||
assert_eq!(executor.retained_task_count(), 0);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user