use crate::{
durable::{state_set_if_not_exists, BINDING_DAP_REPORTS_PROCESSED},
int_err,
};
use futures::future::try_join_all;
use worker::*;
pub(crate) const DURABLE_REPORTS_PROCESSED_MARK_AGGREGATED: &str =
"/internal/do/report_store/mark_aggregated";
#[durable_object]
pub struct ReportsProcessed {
#[allow(dead_code)]
state: State,
env: Env,
touched: bool,
}
impl ReportsProcessed {
async fn to_checked(&self, report_id_hex: String) -> Result<Option<String>> {
let key = format!("processed/{}", report_id_hex);
let processed: bool = state_set_if_not_exists(&self.state, &key, &true)
.await?
.unwrap_or(false);
if processed {
Ok(Some(report_id_hex))
} else {
Ok(None)
}
}
}
#[durable_object]
impl DurableObject for ReportsProcessed {
fn new(state: State, env: Env) -> Self {
Self {
state,
env,
touched: false,
}
}
async fn fetch(&mut self, mut req: Request) -> Result<Response> {
let id_hex = self.state.id().to_string();
ensure_garbage_collected!(req, self, id_hex.clone(), BINDING_DAP_REPORTS_PROCESSED);
match (req.path().as_ref(), req.method()) {
(DURABLE_REPORTS_PROCESSED_MARK_AGGREGATED, Method::Post) => {
let report_id_hex_set: Vec<String> = req.json().await?;
let mut requests = Vec::new();
for report_id_hex in report_id_hex_set.into_iter() {
requests.push(self.to_checked(report_id_hex));
}
let responses: Vec<Option<String>> = try_join_all(requests).await?;
let res: Vec<String> = responses.into_iter().flatten().collect();
Response::from_json(&res)
}
_ => Err(int_err(format!(
"ReportsProcessed: unexpected request: method={:?}; path={:?}",
req.method(),
req.path()
))),
}
}
}