(uf: stream::UplinkFrameLog)
| 438 | } |
| 439 | |
| 440 | pub fn device_uplink_frame_log(uf: stream::UplinkFrameLog) -> Validator { |
| 441 | Box::new(move || { |
| 442 | let uf = uf.clone(); |
| 443 | Box::pin(async move { |
| 444 | let key = redis_key(format!("device:{{{}}}:stream:frame", uf.dev_eui)); |
| 445 | let srr: StreamReadReply = redis::cmd("XREAD") |
| 446 | .arg("COUNT") |
| 447 | .arg(1_usize) |
| 448 | .arg("STREAMS") |
| 449 | .arg(&key) |
| 450 | .arg("0") |
| 451 | .query_async(&mut get_async_redis_conn().await.unwrap()) |
| 452 | .await |
| 453 | .unwrap(); |
| 454 | |
| 455 | for stream_key in &srr.keys { |
| 456 | for stream_id in &stream_key.ids { |
| 457 | for (k, v) in &stream_id.map { |
| 458 | assert_eq!("up", k); |
| 459 | if let redis::Value::BulkString(b) = v { |
| 460 | let mut pl = |
| 461 | stream::UplinkFrameLog::decode(&mut Cursor::new(b)).unwrap(); |
| 462 | pl.time = None; // we don't have control over this value |
| 463 | assert_eq!(uf, pl); |
| 464 | } else { |
| 465 | panic!("Invalid payload"); |
| 466 | } |
| 467 | |
| 468 | return; |
| 469 | } |
| 470 | } |
| 471 | } |
| 472 | }) |
| 473 | }) |
| 474 | } |
| 475 | |
| 476 | pub fn scheduler_run_after_set(dev_eui: EUI64) -> Validator { |
| 477 | Box::new(move || { |
nothing calls this directly
no test coverage detected