| 625 | } |
| 626 | |
| 627 | pub async fn get_async_receiver( |
| 628 | transaction_id: u32, |
| 629 | timeout: Duration, |
| 630 | ) -> Result<oneshot::Receiver<Vec<u8>>> { |
| 631 | let (tx, rx) = oneshot::channel(); |
| 632 | |
| 633 | task::spawn(async move { |
| 634 | let mut c = match get_async_redis_conn().await { |
| 635 | Ok(v) => v, |
| 636 | Err(e) => { |
| 637 | error!(error = %e, "Get Redis connection error"); |
| 638 | return; |
| 639 | } |
| 640 | }; |
| 641 | let key = redis_key(format!("backend:async:{}", transaction_id)); |
| 642 | |
| 643 | let srr: StreamReadReply = match redis::cmd("XREAD") |
| 644 | .arg("BLOCK") |
| 645 | .arg(timeout.as_millis() as u64) |
| 646 | .arg("COUNT") |
| 647 | .arg(1_u64) |
| 648 | .arg("STREAMS") |
| 649 | .arg(&key) |
| 650 | .arg("0") |
| 651 | .query_async(&mut c) |
| 652 | .await |
| 653 | { |
| 654 | Ok(v) => v, |
| 655 | Err(e) => { |
| 656 | error!(error = %e, "Read from Redis Stream error"); |
| 657 | return; |
| 658 | } |
| 659 | }; |
| 660 | |
| 661 | for stream_key in &srr.keys { |
| 662 | for stream_id in &stream_key.ids { |
| 663 | for (k, v) in &stream_id.map { |
| 664 | match k.as_ref() { |
| 665 | "pl" => { |
| 666 | if let redis::Value::BulkString(b) = v { |
| 667 | let _ = tx.send(b.to_vec()); |
| 668 | return; |
| 669 | } |
| 670 | } |
| 671 | _ => { |
| 672 | error!( |
| 673 | transaction_id = transaction_id, |
| 674 | key = %key, |
| 675 | "Unexpected key in async stream" |
| 676 | ); |
| 677 | } |
| 678 | } |
| 679 | } |
| 680 | } |
| 681 | } |
| 682 | }); |
| 683 | |
| 684 | Ok(rx) |