pub mod counters;
pub mod fetch;
pub mod storage;
pub use counters::PlatformCounters;
pub use fetch::{FetchState, RequestId};
pub use storage::MemoryStorage;
use std::sync::{Arc, Mutex, Weak};
use blitz_traits::platform::{
FetchError, FetchHandler, FetchProvider, FetchRequest, FetchResponse, OriginKey, StorageError,
StorageProvider,
};
use crate::fetch::{InFlight, Slot};
pub type ReadyWaker = Arc<dyn Fn() + Send + Sync + 'static>;
pub struct PlatformHost {
origin: OriginKey,
fetch_provider: Arc<dyn FetchProvider>,
storage_provider: Arc<dyn StorageProvider>,
inflight: Arc<Mutex<InFlight>>,
waker: Option<ReadyWaker>,
counters: Mutex<PlatformCounters>,
}
impl PlatformHost {
pub fn new(
origin: OriginKey,
fetch_provider: Arc<dyn FetchProvider>,
storage_provider: Arc<dyn StorageProvider>,
) -> Self {
Self {
origin,
fetch_provider,
storage_provider,
inflight: Arc::new(Mutex::new(InFlight::default())),
waker: None,
counters: Mutex::new(PlatformCounters::default()),
}
}
pub fn with_waker(mut self, waker: ReadyWaker) -> Self {
self.waker = Some(waker);
self
}
pub fn origin(&self) -> &OriginKey {
&self.origin
}
pub fn counters(&self) -> PlatformCounters {
*self.counters.lock().unwrap()
}
pub fn start_fetch(&self, request: FetchRequest) -> RequestId {
let sent = request.body.as_ref().map(|body| body.len()).unwrap_or(0);
let id = {
let mut inflight = self.inflight.lock().unwrap();
inflight.begin()
};
{
let mut counters = self.counters.lock().unwrap();
counters.fetches_started += 1;
counters.fetch_bytes_sent += sent as u64;
}
self.fetch_provider.fetch(
request,
Box::new(Completion {
id,
inflight: Arc::downgrade(&self.inflight),
waker: self.waker.clone(),
}),
);
id
}
pub fn take_ready(&self) -> Vec<RequestId> {
let ready = {
let mut inflight = self.inflight.lock().unwrap();
std::mem::take(&mut inflight.ready)
};
if !ready.is_empty() {
let received: u64 = ready
.iter()
.filter_map(|id| self.with_response(*id, |response| response.body.len() as u64))
.sum();
let mut counters = self.counters.lock().unwrap();
counters.fetches_completed += ready.len() as u64;
counters.fetch_bytes_received += received;
}
ready
}
pub fn state(&self, id: RequestId) -> FetchState {
let inflight = self.inflight.lock().unwrap();
match inflight.slots.get(&id) {
None => FetchState::Unknown,
Some(Slot::Pending) => FetchState::Pending,
Some(Slot::Done(answer)) => match answer.as_ref() {
Ok(_) => FetchState::Response,
Err(error) => FetchState::Failed(error.clone()),
},
}
}
pub fn with_response<T>(
&self,
id: RequestId,
read: impl FnOnce(&FetchResponse) -> T,
) -> Option<T> {
let inflight = self.inflight.lock().unwrap();
match inflight.slots.get(&id) {
Some(Slot::Done(answer)) => answer.as_ref().as_ref().ok().map(read),
_ => None,
}
}
pub fn release(&self, id: RequestId) -> bool {
let mut inflight = self.inflight.lock().unwrap();
inflight.ready.retain(|ready| *ready != id);
inflight.slots.remove(&id).is_some()
}
pub fn tracked_requests(&self) -> usize {
self.inflight.lock().unwrap().slots.len()
}
pub fn storage_get(&self, key: &str) -> Option<String> {
let value = self.storage_provider.get(&self.origin, key);
let mut counters = self.counters.lock().unwrap();
counters.storage_reads += 1;
counters.storage_bytes_read += value.as_ref().map(|v| v.len()).unwrap_or(0) as u64;
value
}
pub fn storage_set(&self, key: &str, value: &str) -> Result<(), StorageError> {
let result = self.storage_provider.set(&self.origin, key, value);
if result.is_ok() {
let mut counters = self.counters.lock().unwrap();
counters.storage_writes += 1;
counters.storage_bytes_written += (key.len() + value.len()) as u64;
}
result
}
pub fn storage_remove(&self, key: &str) {
self.storage_provider.remove(&self.origin, key);
self.counters.lock().unwrap().storage_writes += 1;
}
pub fn storage_clear(&self) {
self.storage_provider.clear(&self.origin);
self.counters.lock().unwrap().storage_writes += 1;
}
}
struct Completion {
id: RequestId,
inflight: Weak<Mutex<InFlight>>,
waker: Option<ReadyWaker>,
}
impl FetchHandler for Completion {
fn complete(self: Box<Self>, result: Result<FetchResponse, FetchError>) {
let Some(inflight) = self.inflight.upgrade() else {
return;
};
{
let mut inflight = inflight.lock().unwrap();
let Some(slot) = inflight.slots.get_mut(&self.id) else {
return;
};
*slot = Slot::Done(Box::new(result));
inflight.ready.push(self.id);
}
if let Some(waker) = &self.waker {
waker();
}
}
}
#[cfg(test)]
mod tests;