use crate::{
durable::{
durable_name_queue,
leader_agg_job_queue::{
DURABLE_LEADER_AGG_JOB_QUEUE_FINISH, DURABLE_LEADER_AGG_JOB_QUEUE_PUT,
},
report_id_hex_from_report, state_get, state_set_if_not_exists, DurableConnector,
DurableOrdered, BINDING_DAP_LEADER_AGG_JOB_QUEUE, BINDING_DAP_REPORTS_PENDING,
},
int_err,
};
use serde::{Deserialize, Serialize};
use worker::*;
pub(crate) const DURABLE_REPORTS_PENDING_GET: &str = "/internal/do/reports_pending/get";
pub(crate) const DURABLE_REPORTS_PENDING_PUT: &str = "/internal/do/reports_pending/put";
#[derive(Deserialize, Serialize)]
#[serde(rename_all = "snake_case")]
pub(crate) enum ReportsPendingResult {
Ok,
ErrReportExists,
}
#[durable_object]
pub struct ReportsPending {
#[allow(dead_code)]
state: State,
env: Env,
touched: bool,
}
#[durable_object]
impl DurableObject for ReportsPending {
fn new(state: State, env: Env) -> Self {
Self {
state,
env,
touched: false,
}
}
async fn fetch(&mut self, mut req: Request) -> Result<Response> {
let durable = DurableConnector::new(&self.env);
let id_hex = self.state.id().to_string();
ensure_garbage_collected!(req, self, id_hex.clone(), BINDING_DAP_REPORTS_PENDING);
match (req.path().as_ref(), req.method()) {
(DURABLE_REPORTS_PENDING_GET, Method::Post) => {
let reports_requested: usize = req.json().await?;
let opt = ListOptions::new()
.prefix("pending/")
.limit(reports_requested);
let iter = self.state.storage().list_with_options(opt).await?.entries();
let mut item = iter.next()?;
let mut reports = Vec::with_capacity(reports_requested);
let mut keys = Vec::with_capacity(reports_requested);
while !item.done() {
#[allow(deprecated)]
let (key, report_hex): (String, String) = item.value().into_serde()?;
reports.push(report_hex);
keys.push(key);
item = iter.next()?;
}
self.state.storage().delete_multiple(keys).await?;
let empty = self
.state
.storage()
.list_with_options(ListOptions::new().prefix("pending/").limit(1))
.await?
.size()
== 0;
if empty
&& self
.env
.var("DAP_ISSUE73_DISABLE_AGG_JOB_QUEUE_GARBAGE_COLLECTION")?
.to_string()
== "true"
{
let agg_job: Option<DurableOrdered<String>> =
state_get(&self.state, "agg_job").await?;
if let Some(agg_job) = agg_job {
durable
.post(
BINDING_DAP_LEADER_AGG_JOB_QUEUE,
DURABLE_LEADER_AGG_JOB_QUEUE_FINISH,
durable_name_queue(0),
&agg_job,
)
.await?;
self.state.storage().delete("agg_job").await?;
}
}
console_debug!(
"drained {} reports from bucket {}",
reports.len(),
self.state.id().to_string()
);
Response::from_json(&reports)
}
(DURABLE_REPORTS_PENDING_PUT, Method::Post) => {
let report_hex: String = req.json().await?;
let report_id_hex = report_id_hex_from_report(&report_hex)
.ok_or_else(|| int_err("failed to parse report_id from report"))?;
let key = format!("pending/{}", report_id_hex);
let exists = state_set_if_not_exists::<String>(&self.state, &key, &report_hex)
.await?
.is_some();
if exists {
return Response::from_json(&ReportsPendingResult::ErrReportExists);
}
let agg_job: Option<DurableOrdered<String>> =
state_get(&self.state, "agg_job").await?;
if agg_job.is_none() {
let agg_job = DurableOrdered::new_roughly_ordered(id_hex, "agg_job");
durable
.post(
BINDING_DAP_LEADER_AGG_JOB_QUEUE,
DURABLE_LEADER_AGG_JOB_QUEUE_PUT,
durable_name_queue(0),
&agg_job,
)
.await?;
self.state.storage().put("agg_job", agg_job).await?;
}
Response::from_json(&ReportsPendingResult::Ok)
}
_ => Err(int_err(format!(
"ReportsPending: unexpected request: method={:?}; path={:?}",
req.method(),
req.path()
))),
}
}
}