NLinker icon

Simple async service

NLinker | PRO | 02/02/20 05:57:14 PM UTC | 0 ⭐ | 8088 👁️ | Never ⏰ | []
Rust |

2.39 KB

|

None

|

0 👍

/

0 👎

// https://play.rust-lang.org/?version=stable&mode=debug&edition=2018&gist=fef597285e9b04c52512daf255ec988c
 
use num_cpus::get;
use std::fmt::Display;
use std::net::{TcpStream, ToSocketAddrs};
use std::sync::mpsc::{Receiver, SyncSender};
use std::sync::{mpsc, Arc};
use std::thread::spawn;
use tokio; // 0.2.9
use tokio::runtime::Builder;
use tokio::sync::Semaphore;
 
#[derive(Debug, Clone)]
struct Response {
    result: String,
    hostname: String,
    process_time: String,
    status: bool,
}
async fn process_host<A>(
    hostname: A,
    tx: SyncSender<Response>,
    connection_pool: Arc<Semaphore>,
) -> Response
where
    A: ToSocketAddrs + Display,
{
    let guard = connection_pool.acquire().await;
    let res = Response {
        result: "none".to_string(),
        hostname: hostname.to_string(),
        process_time: "sometime".to_string(),
        status: false,
    };
    let tcp = match TcpStream::connect(&hostname) {
        Ok(a) => a,
        Err(e) => {
            let res = Response {
                result: e.to_string(),
                hostname: hostname.to_string(),
                process_time: "sometime".to_string(),
                status: false,
            };
            tx.send(res.clone()).unwrap();
            return res;
        }
    };
    // chain of async calls.
    res
}
 
fn main() {
    let hosts = vec!["1.1.1.1:22", "8.8.8.8:22" , "8.8.4.4:22"];
    let num_of_threads = Arc::new(Semaphore::new(10));
    let mut reactor = Builder::new()
        .enable_all()
        .threaded_scheduler()
        .core_threads(get())
        .build()
        .unwrap();
    let (tx, rx): (SyncSender<Response>, Receiver<Response>) = mpsc::sync_channel(0);
    let queue_len =100;
    spawn(move || incremental_save(rx,  queue_len));
 
    let tasks: Vec<_> = hosts.into_iter().map(|host| {
        reactor.spawn(process_host(
            host,
            tx.clone(),
            num_of_threads.clone(),
        ))
    }).collect();
    
    reactor.block_on(futures::future::join_all(tasks));
}
fn incremental_save(rx: Receiver<Response>, queue_len: u64) {
    for _ in 0..=queue_len {
        let received = match rx.recv() {
            Ok(a) => a,
            Err(e) => {
                eprintln!("incremental_save: {}", e);
                break;
            }
        };
        println! {"{:?}", received};
    }
}

Comments

  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎

    
        
  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎

    
        
  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎

    
        
  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎

    
        
  •  icon
    01/01/70 12:00:00 AM UTC
    Plain Text |

    0 B

    |

    👍

    /

    👎