| 398 | } |
| 399 | |
| 400 | async fn enqueue( |
| 401 | &self, |
| 402 | request: Request<api::EnqueueMulticastGroupQueueItemRequest>, |
| 403 | ) -> Result<Response<api::EnqueueMulticastGroupQueueItemResponse>, Status> { |
| 404 | let req_enq = match &request.get_ref().queue_item { |
| 405 | Some(v) => v, |
| 406 | None => { |
| 407 | return Err(Status::invalid_argument("queue_item is missing")); |
| 408 | } |
| 409 | }; |
| 410 | |
| 411 | let mg_id = Uuid::from_str(&req_enq.multicast_group_id).map_err(|e| e.status())?; |
| 412 | |
| 413 | self.validator |
| 414 | .validate( |
| 415 | request.extensions(), |
| 416 | validator::ValidateMulticastGroupQueueAccess::new(validator::Flag::Create, mg_id), |
| 417 | ) |
| 418 | .await?; |
| 419 | |
| 420 | let f_cnt = downlink::multicast::enqueue(multicast::MulticastGroupQueueItem { |
| 421 | multicast_group_id: mg_id.into(), |
| 422 | f_port: req_enq.f_port as i16, |
| 423 | data: req_enq.data.clone(), |
| 424 | expires_at: if let Some(expires_at) = req_enq.expires_at { |
| 425 | let expires_at: std::time::SystemTime = expires_at |
| 426 | .try_into() |
| 427 | .map_err(|e: prost_types::TimestampError| e.status())?; |
| 428 | Some(expires_at.into()) |
| 429 | } else { |
| 430 | None |
| 431 | }, |
| 432 | ..Default::default() |
| 433 | }) |
| 434 | .await |
| 435 | .map_err(|e| e.status())?; |
| 436 | |
| 437 | let mut resp = Response::new(api::EnqueueMulticastGroupQueueItemResponse { f_cnt }); |
| 438 | resp.metadata_mut().insert( |
| 439 | "x-log-multicast_group_id", |
| 440 | req_enq.multicast_group_id.parse().unwrap(), |
| 441 | ); |
| 442 | |
| 443 | Ok(resp) |
| 444 | } |
| 445 | |
| 446 | async fn flush_queue( |
| 447 | &self, |