use crate::{
durable::{state_get, state_get_or_default, DurableOrdered, BINDING_DAP_LEADER_COL_JOB_QUEUE},
int_err,
};
use daphne::{
messages::{CollectReq, CollectResp, Id},
DapCollectJob,
};
use prio::{
codec::{Decode, Encode},
vdaf::prg::{Prg, PrgAes128, Seed, SeedStream},
};
use worker::*;
const PENDING_PREFIX: &str = "pending";
const PROCESSED_PREFIX: &str = "processed";
pub(crate) const DURABLE_LEADER_COL_JOB_QUEUE_PUT: &str = "/internal/do/leader_col_job_queue/put";
pub(crate) const DURABLE_LEADER_COL_JOB_QUEUE_GET: &str = "/internal/do/leader_col_job_queue/get";
pub(crate) const DURABLE_LEADER_COL_JOB_QUEUE_FINISH: &str =
"/internal/do/leader_col_job_queue/finish";
pub(crate) const DURABLE_LEADER_COL_JOB_QUEUE_GET_RESULT: &str =
"/internal/do/leader_col_job_queue/get_result";
#[durable_object]
pub struct LeaderCollectionJobQueue {
#[allow(dead_code)]
state: State,
env: Env,
touched: bool,
}
#[durable_object]
impl DurableObject for LeaderCollectionJobQueue {
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_LEADER_COL_JOB_QUEUE);
match (req.path().as_ref(), req.method()) {
(DURABLE_LEADER_COL_JOB_QUEUE_PUT, Method::Post) => {
let collect_req: CollectReq = req.json().await?;
let collect_req_bytes = collect_req.get_encoded();
let collect_id_key_hex = self
.env
.secret("DAP_COLLECT_ID_KEY")
.map_err(|e| format!("failed to load DAP_COLLECT_ID_KEY: {}", e))?
.to_string();
let collect_id_key =
Seed::get_decoded(&hex::decode(&collect_id_key_hex).map_err(int_err)?)
.map_err(int_err)?;
let mut collect_id_bytes = [0; 32];
PrgAes128::seed_stream(&collect_id_key, &collect_req_bytes)
.fill(&mut collect_id_bytes);
let collect_id = Id(collect_id_bytes);
let collect_id_hex = collect_id.to_hex();
let pending_key = format!("pending/id/{}", collect_id_hex);
let pending: bool = state_get_or_default(&self.state, &pending_key).await?;
let processed: Option<CollectResp> = state_get(
&self.state,
&format!("{}/{}", PROCESSED_PREFIX, collect_id_hex),
)
.await?;
if processed.is_none() && !pending {
let queued = DurableOrdered::new_strictly_ordered(
&self.state,
(collect_id, collect_req),
PENDING_PREFIX,
)
.await?;
queued.put(&self.state).await?;
self.state
.storage()
.put(&lookup_key(&collect_id_hex), &queued.key())
.await?;
}
Response::from_json(&collect_id_hex)
}
(DURABLE_LEADER_COL_JOB_QUEUE_GET, Method::Get) => {
let queue: Vec<(Id, CollectReq)> =
DurableOrdered::get_all(&self.state, PENDING_PREFIX)
.await?
.into_iter()
.map(|queued| queued.into_item())
.collect();
Response::from_json(&queue)
}
(DURABLE_LEADER_COL_JOB_QUEUE_FINISH, Method::Post) => {
let (collect_id, collect_resp): (Id, CollectResp) = req.json().await?;
let collect_id_hex = collect_id.to_hex();
let processed_key = format!("{}/{}", PROCESSED_PREFIX, collect_id_hex);
let processed: Option<CollectResp> = state_get(&self.state, &processed_key).await?;
if processed.is_some() {
return Err(int_err(
"LeaderCollectionJobQueue: tried to overwrite collect response",
));
}
let pending_lookup_key = lookup_key(&collect_id_hex);
if let Some(lookup_val) =
state_get::<String>(&self.state, &pending_lookup_key).await?
{
self.state.storage().delete(&lookup_val).await?;
}
let mut storage = self.state.storage();
let f = storage.delete(&pending_lookup_key);
self.state
.storage()
.put(&processed_key, collect_resp)
.await?;
f.await?;
Response::from_json(&())
}
(DURABLE_LEADER_COL_JOB_QUEUE_GET_RESULT, Method::Post) => {
let collect_id: Id = req.json().await?;
let collect_id_hex = collect_id.to_hex();
let pending_lookup_key = lookup_key(&collect_id_hex);
let pending = state_get::<String>(&self.state, &pending_lookup_key)
.await?
.is_some();
let processed_key = format!("{}/{}", PROCESSED_PREFIX, collect_id_hex);
let processed: Option<CollectResp> = state_get(&self.state, &processed_key).await?;
if let Some(collect_resp) = processed {
if pending {
self.state.storage().delete(&pending_lookup_key).await?;
}
Response::from_json(&DapCollectJob::Done(collect_resp))
} else if pending {
Response::from_json(&DapCollectJob::Pending)
} else {
Response::from_json(&DapCollectJob::Unknown)
}
}
_ => Err(int_err(format!(
"LeaderCollectionJobQueue: unexpected request: method={:?}; path={:?}",
req.method(),
req.path()
))),
}
}
}
fn lookup_key(collect_id_hex: &str) -> String {
format!("{}/id/{}", PENDING_PREFIX, collect_id_hex)
}