use core::marker::PhantomData;
use std::{
fs,
path::{Path, PathBuf},
time::SystemTime,
};
use serde::{Deserialize, Serialize};
use crate::{
bolts::{current_time, shmem::ShMemProvider},
corpus::{Corpus, CorpusId, HasTestcase},
events::{llmp::LlmpEventConverter, Event, EventConfig, EventFirer},
executors::{Executor, ExitKind, HasObservers},
fuzzer::{Evaluator, EvaluatorObservers, ExecutionProcessor},
inputs::{Input, InputConverter, UsesInput},
stages::Stage,
state::{HasClientPerfMonitor, HasCorpus, HasExecutions, HasMetadata, HasRand, UsesState},
Error,
};
#[derive(Serialize, Deserialize, Debug)]
pub struct SyncFromDiskMetadata {
pub last_time: SystemTime,
}
crate::impl_serdeany!(SyncFromDiskMetadata);
impl SyncFromDiskMetadata {
#[must_use]
pub fn new(last_time: SystemTime) -> Self {
Self { last_time }
}
}
#[derive(Debug)]
pub struct SyncFromDiskStage<CB, E, EM, Z> {
sync_dir: PathBuf,
load_callback: CB,
phantom: PhantomData<(E, EM, Z)>,
}
impl<CB, E, EM, Z> UsesState for SyncFromDiskStage<CB, E, EM, Z>
where
E: UsesState,
{
type State = E::State;
}
impl<CB, E, EM, Z> Stage<E, EM, Z> for SyncFromDiskStage<CB, E, EM, Z>
where
CB: FnMut(&mut Z, &mut Z::State, &Path) -> Result<<Z::State as UsesInput>::Input, Error>,
E: UsesState<State = Z::State>,
EM: UsesState<State = Z::State>,
Z: Evaluator<E, EM>,
Z::State: HasClientPerfMonitor + HasCorpus + HasRand + HasMetadata,
{
#[inline]
fn perform(
&mut self,
fuzzer: &mut Z,
executor: &mut E,
state: &mut Z::State,
manager: &mut EM,
_corpus_idx: CorpusId,
) -> Result<(), Error> {
let last = state
.metadata_map()
.get::<SyncFromDiskMetadata>()
.map(|m| m.last_time);
let path = self.sync_dir.clone();
if let Some(max_time) =
self.load_from_directory(&path, &last, fuzzer, executor, state, manager)?
{
if last.is_none() {
state
.metadata_map_mut()
.insert(SyncFromDiskMetadata::new(max_time));
} else {
state
.metadata_map_mut()
.get_mut::<SyncFromDiskMetadata>()
.unwrap()
.last_time = max_time;
}
}
#[cfg(feature = "introspection")]
state.introspection_monitor_mut().finish_stage();
Ok(())
}
}
impl<CB, E, EM, Z> SyncFromDiskStage<CB, E, EM, Z>
where
CB: FnMut(&mut Z, &mut Z::State, &Path) -> Result<<Z::State as UsesInput>::Input, Error>,
E: UsesState<State = Z::State>,
EM: UsesState<State = Z::State>,
Z: Evaluator<E, EM>,
Z::State: HasClientPerfMonitor + HasCorpus + HasRand + HasMetadata,
{
#[must_use]
pub fn new(sync_dir: PathBuf, load_callback: CB) -> Self {
Self {
sync_dir,
load_callback,
phantom: PhantomData,
}
}
fn load_from_directory(
&mut self,
in_dir: &Path,
last: &Option<SystemTime>,
fuzzer: &mut Z,
executor: &mut E,
state: &mut Z::State,
manager: &mut EM,
) -> Result<Option<SystemTime>, Error> {
let mut max_time = None;
for entry in fs::read_dir(in_dir)? {
let entry = entry?;
let path = entry.path();
let attributes = fs::metadata(&path);
if attributes.is_err() {
continue;
}
let attr = attributes?;
if attr.is_file() && attr.len() > 0 {
if let Ok(time) = attr.modified() {
if let Some(l) = last {
if time.duration_since(*l).is_err() {
continue;
}
}
max_time = Some(max_time.map_or(time, |t: SystemTime| t.max(time)));
let input = (self.load_callback)(fuzzer, state, &path)?;
fuzzer.evaluate_input(state, executor, manager, input)?;
}
} else if attr.is_dir() {
let dir_max_time =
self.load_from_directory(&path, last, fuzzer, executor, state, manager)?;
if let Some(time) = dir_max_time {
max_time = Some(max_time.map_or(time, |t: SystemTime| t.max(time)));
}
}
}
Ok(max_time)
}
}
pub type SyncFromDiskFunction<S, Z> =
fn(&mut Z, &mut S, &Path) -> Result<<S as UsesInput>::Input, Error>;
impl<E, EM, Z> SyncFromDiskStage<SyncFromDiskFunction<Z::State, Z>, E, EM, Z>
where
E: UsesState<State = Z::State>,
EM: UsesState<State = Z::State>,
Z: Evaluator<E, EM>,
Z::State: HasClientPerfMonitor + HasCorpus + HasRand + HasMetadata,
{
#[must_use]
pub fn with_from_file(sync_dir: PathBuf) -> Self {
fn load_callback<S: UsesInput, Z>(
_: &mut Z,
_: &mut S,
p: &Path,
) -> Result<S::Input, Error> {
Input::from_file(p)
}
Self {
sync_dir,
load_callback: load_callback::<_, _>,
phantom: PhantomData,
}
}
}
#[derive(Serialize, Deserialize, Debug)]
pub struct SyncFromBrokerMetadata {
pub last_id: Option<CorpusId>,
}
crate::impl_serdeany!(SyncFromBrokerMetadata);
impl SyncFromBrokerMetadata {
#[must_use]
pub fn new(last_id: Option<CorpusId>) -> Self {
Self { last_id }
}
}
#[derive(Debug)]
pub struct SyncFromBrokerStage<IC, ICB, DI, S, SP>
where
SP: ShMemProvider + 'static,
S: UsesInput,
IC: InputConverter<From = S::Input, To = DI>,
ICB: InputConverter<From = DI, To = S::Input>,
DI: Input,
{
client: LlmpEventConverter<IC, ICB, DI, S, SP>,
}
impl<IC, ICB, DI, S, SP> UsesState for SyncFromBrokerStage<IC, ICB, DI, S, SP>
where
SP: ShMemProvider + 'static,
S: UsesInput,
IC: InputConverter<From = S::Input, To = DI>,
ICB: InputConverter<From = DI, To = S::Input>,
DI: Input,
{
type State = S;
}
impl<E, EM, IC, ICB, DI, S, SP, Z> Stage<E, EM, Z> for SyncFromBrokerStage<IC, ICB, DI, S, SP>
where
EM: UsesState<State = S> + EventFirer,
S: UsesInput
+ HasClientPerfMonitor
+ HasExecutions
+ HasCorpus
+ HasRand
+ HasMetadata
+ HasTestcase,
SP: ShMemProvider,
E: HasObservers<State = S> + Executor<EM, Z>,
for<'a> E::Observers: Deserialize<'a>,
Z: EvaluatorObservers<E::Observers, State = S> + ExecutionProcessor<E::Observers, State = S>,
IC: InputConverter<From = S::Input, To = DI>,
ICB: InputConverter<From = DI, To = S::Input>,
DI: Input,
{
#[inline]
fn perform(
&mut self,
fuzzer: &mut Z,
executor: &mut E,
state: &mut Z::State,
manager: &mut EM,
_corpus_idx: CorpusId,
) -> Result<(), Error> {
if self.client.can_convert() {
let last_id = state
.metadata_map()
.get::<SyncFromBrokerMetadata>()
.and_then(|m| m.last_id);
let mut cur_id =
last_id.map_or_else(|| state.corpus().first(), |id| state.corpus().next(id));
while let Some(id) = cur_id {
let input = state.corpus().cloned_input_for_id(id)?;
self.client.fire(
state,
Event::NewTestcase {
input,
observers_buf: None,
exit_kind: ExitKind::Ok,
corpus_size: 0, client_config: EventConfig::AlwaysUnique,
time: current_time(),
executions: 0,
forward_id: None,
},
)?;
cur_id = state.corpus().next(id);
}
let last = state.corpus().last();
if last_id.is_none() {
state
.metadata_map_mut()
.insert(SyncFromBrokerMetadata::new(last));
} else {
state
.metadata_map_mut()
.get_mut::<SyncFromBrokerMetadata>()
.unwrap()
.last_id = last;
}
}
self.client.process(fuzzer, state, executor, manager)?;
#[cfg(feature = "introspection")]
state.introspection_monitor_mut().finish_stage();
Ok(())
}
}
impl<IC, ICB, DI, S, SP> SyncFromBrokerStage<IC, ICB, DI, S, SP>
where
SP: ShMemProvider + 'static,
S: UsesInput,
IC: InputConverter<From = S::Input, To = DI>,
ICB: InputConverter<From = DI, To = S::Input>,
DI: Input,
{
#[must_use]
pub fn new(client: LlmpEventConverter<IC, ICB, DI, S, SP>) -> Self {
Self { client }
}
}