【问题标题】:How to send a message to another actor from actix SyncContext?如何从 actix SyncContext 向另一个参与者发送消息?
【发布时间】:2021-02-23 10:20:24
【问题描述】:

我想实现一个长时间运行的后台任务,它可以向其他Actors 报告进度。我已经做到了。 但我也希望能够再次取消长时间运行的后台任务。

到目前为止,我得到的是:

use actix::prelude::*;

struct Worker {}

impl Actor for Worker {
    type Context = SyncContext<Self>;
}

struct Manager {
    worker: Addr<Worker>,
}

impl Actor for Manager {
    type Context = Context<Self>;
}

impl Supervised for Manager {}

impl SystemService for Manager {
    fn service_started(&mut self, _ctx: &mut Context<Self>) {}
}

struct Work {}

#[derive(Message)]
#[rtype(result = "()")]
struct PerformWork(Work);

#[derive(Message)]
#[rtype(result = "()")]
pub struct ReportProgress(i32);

impl Handler<PerformWork> for Worker {
    type Result = ();

    fn handle(&mut self, msg: PerformWork, ctx: &mut Self::Context) -> Self::Result {
        for i in 0..10000000 {
            // Report progress
            Manager::from_registry().do_send(ReportProgress(i));
            // Do some very slow I/O.
            thread::sleep(time::Duration::from_millis(1));
        }
    }
}

impl Handler<ReportProgress> for Manager {
    type Result = ();

    fn handle(&mut self, msg: ReportProgress, ctx: &mut Self::Context) -> Self::Result {
        // Do something with the progress here
    }
}

Manager 还处理一个Message,它将PerformWork Message 发送到Worker

我想给ReportProgress Message 一个bool 返回类型,这将允许Worker 决定它是否应该跳出它的循环。但是,我无法将带有返回结果的Message 发送到Manager。 使用 send() 而不是 do_send() 返回一个 Future,我无法在 SyncContext 内解决。

非常感谢任何想法。

更多背景知识:

  • 真正慢的 I/O 是串行通信。
  • actix 是 0.10 版

【问题讨论】:

    标签: rust rust-actix


    【解决方案1】:

    我找到了一个解决方案,但我不相信它是一个好的解决方案。

    我添加了一个Arc&lt;AtomicBool&gt;&gt;,它传递给WorkerManager 保留对AtomicBool 的引用并且可以修改它。如果AtomicBoolManager 修改,Worker 将跳出循环。

    use actix::prelude::*;
    use std::sync::atomic::{AtomicBool, Ordering};
    
    struct Worker {}
    
    impl Actor for Worker {
        type Context = SyncContext<Self>;
    }
    
    struct Manager {
        worker: Addr<Worker>,
    }
    
    impl Actor for Manager {
        type Context = Context<Self>;
    }
    
    impl Supervised for Manager {}
    
    impl SystemService for Manager {
        fn service_started(&mut self, _ctx: &mut Context<Self>) {}
    }
    
    struct Work {}
    
    #[derive(Message)]
    #[rtype(result = "()")]
    struct PerformWork(Work, Arc<AtomicBool>>);
    
    #[derive(Message)]
    #[rtype(result = "()")]
    pub struct ReportProgress(i32);
    
    impl Handler<PerformWork> for Worker {
        type Result = ();
    
        fn handle(&mut self, msg: PerformWork, ctx: &mut Self::Context) -> Self::Result {
            for i in 0..10000000 {
                // Report progress
                Manager::from_registry().do_send(ReportProgress(i));
                if msg.1.load(Ordering::Relaxed) {
                    break;
                }
                // Do some very slow I/O.
                thread::sleep(time::Duration::from_millis(1));
            }
        }
    }
    
    impl Handler<ReportProgress> for Manager {
        type Result = ();
    
        fn handle(&mut self, msg: ReportProgress, ctx: &mut Self::Context) -> Self::Result {
            // Do something with the progress here
        }
    }
    

    【讨论】:

      猜你喜欢
      • 2021-09-11
      • 1970-01-01
      • 1970-01-01
      • 2021-12-31
      • 1970-01-01
      • 2015-07-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多