MCPcopy Create free account
hub / github.com/chirpstack/chirpstack / get_async_receiver

Function get_async_receiver

chirpstack/src/api/backend/mod.rs:627–685  ·  view source on GitHub ↗
(
    transaction_id: u32,
    timeout: Duration,
)

Source from the content-addressed store, hash-verified

625}
626
627pub 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)

Callers 7

test_async_responseFunction · 0.85
start_pr_sessionMethod · 0.85
get_home_net_idMethod · 0.85
start_roamingMethod · 0.85

Calls 4

get_async_redis_connFunction · 0.85
redis_keyFunction · 0.85
as_refMethod · 0.80
to_vecMethod · 0.45

Tested by 1

test_async_responseFunction · 0.68