use crate::{
durable::{state_get_or_default, BINDING_DAP_AGGREGATE_STORE},
int_err,
};
use daphne::DapAggregateShare;
use worker::*;
pub(crate) const DURABLE_AGGREGATE_STORE_GET: &str = "/internal/do/aggregate_store/get";
pub(crate) const DURABLE_AGGREGATE_STORE_MERGE: &str = "/internal/do/aggregate_store/merge";
pub(crate) const DURABLE_AGGREGATE_STORE_MARK_COLLECTED: &str =
"/internal/do/aggregate_store/mark_collected";
pub(crate) const DURABLE_AGGREGATE_STORE_CHECK_COLLECTED: &str =
"/internal/do/aggregate_store/check_collected";
#[durable_object]
pub struct AggregateStore {
#[allow(dead_code)]
state: State,
env: Env,
touched: bool,
}
#[durable_object]
impl DurableObject for AggregateStore {
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, BINDING_DAP_AGGREGATE_STORE);
match (req.path().as_ref(), req.method()) {
(DURABLE_AGGREGATE_STORE_MERGE, Method::Post) => {
let agg_share_delta = req.json().await?;
let mut agg_share: DapAggregateShare =
state_get_or_default(&self.state, "agg_share").await?;
agg_share.merge(agg_share_delta).map_err(int_err)?;
self.state.storage().put("agg_share", agg_share).await?;
Response::from_json(&())
}
(DURABLE_AGGREGATE_STORE_GET, Method::Get) => {
let agg_share: DapAggregateShare =
state_get_or_default(&self.state, "agg_share").await?;
Response::from_json(&agg_share)
}
(DURABLE_AGGREGATE_STORE_MARK_COLLECTED, Method::Post) => {
self.state.storage().put("collected", true).await?;
Response::from_json(&())
}
(DURABLE_AGGREGATE_STORE_CHECK_COLLECTED, Method::Get) => {
let collected: bool = state_get_or_default(&self.state, "collected").await?;
Response::from_json(&collected)
}
_ => Err(int_err(format!(
"AggregatesStore: unexpected request: method={:?}; path={:?}",
req.method(),
req.path()
))),
}
}
}