| 1103 | } |
| 1104 | |
| 1105 | async fn enqueue( |
| 1106 | &self, |
| 1107 | request: Request<api::EnqueueDeviceQueueItemRequest>, |
| 1108 | ) -> Result<Response<api::EnqueueDeviceQueueItemResponse>, Status> { |
| 1109 | let req = request.get_ref(); |
| 1110 | |
| 1111 | let req_qi = match &req.queue_item { |
| 1112 | Some(v) => v, |
| 1113 | None => { |
| 1114 | return Err(Status::invalid_argument("queue_item is missing")); |
| 1115 | } |
| 1116 | }; |
| 1117 | let dev_eui = EUI64::from_str(&req_qi.dev_eui).map_err(|e| e.status())?; |
| 1118 | |
| 1119 | if req.flush_queue { |
| 1120 | self.validator |
| 1121 | .validate( |
| 1122 | request.extensions(), |
| 1123 | validator::ValidateDeviceQueueAccess::new(validator::Flag::Delete, dev_eui), |
| 1124 | ) |
| 1125 | .await?; |
| 1126 | |
| 1127 | device_queue::flush_non_pending_for_dev_eui(&dev_eui) |
| 1128 | .await |
| 1129 | .map_err(|e| e.status())?; |
| 1130 | } |
| 1131 | |
| 1132 | self.validator |
| 1133 | .validate( |
| 1134 | request.extensions(), |
| 1135 | validator::ValidateDeviceQueueAccess::new(validator::Flag::Create, dev_eui), |
| 1136 | ) |
| 1137 | .await?; |
| 1138 | |
| 1139 | let mut data = req_qi.data.clone(); |
| 1140 | let mut f_port = req_qi.f_port as u8; |
| 1141 | |
| 1142 | if let Some(obj) = &req_qi.object { |
| 1143 | let dev = device::get(&dev_eui).await.map_err(|e| e.status())?; |
| 1144 | let dp = device_profile::get(&dev.device_profile_id) |
| 1145 | .await |
| 1146 | .map_err(|e| e.status())?; |
| 1147 | |
| 1148 | (f_port, data) = codec::struct_to_binary( |
| 1149 | dp.payload_codec_runtime, |
| 1150 | req_qi.f_port as u8, |
| 1151 | &dev.variables, |
| 1152 | &dp.payload_codec_script, |
| 1153 | obj, |
| 1154 | ) |
| 1155 | .await |
| 1156 | .map_err(|e| e.status())?; |
| 1157 | } |
| 1158 | |
| 1159 | let qi = device_queue::DeviceQueueItem { |
| 1160 | id: Uuid::new_v4().into(), |
| 1161 | dev_eui, |
| 1162 | f_port: f_port as i16, |