Skip to content

异步任务

在讨论 AsyncTask 之前,我们需要先讨论 Task。

Task

附加模块通常需要利用 libuv 中的异步助手作为其实现的一部分, 这样,它们就可以安排工作在异步执行,以便它们的方法可以在工作完成之前返回, 这样就可以避免阻塞 Node.js 应用程序的整体执行。

Task 特征提供了一种定义这样的异步任务的方法,该任务需要在 libuv 线程中运行,您可以实现 compute 方法,该方法将在 libuv 线程中调用。

lib.rs
rust
use napi::bindgen_prelude::*;
use napi_derive::napi;

fn fib(n: u32) -> u32 {
  if n <= 1 {
    return n;
  }
  fib(n - 1) + fib(n - 2)
}

pub struct AsyncFib {
  input: u32,
}

#[napi]
impl Task for AsyncFib {
  type Output = u32;
  type JsValue = u32;

  fn compute(&mut self) -> Result<Self::Output> { 
    Ok(fib(self.input)) 
  } 

  fn resolve(&mut self, _: Env, output: u32) -> Result<Self::JsValue> {
    Ok(output)
  }
}

#[napi]
pub fn async_fib(input: u32) -> AsyncTask<AsyncFib> {
  AsyncTask::new(AsyncFib { input })
}

fn compute 在 libuv 线程中运行,您可以在这里运行一些繁重的计算,这不会阻塞 JavaScript 主线程。

你可能会注意到 Task 特征上有两个关联类型,type Output 和 type JsValue, Output 是 compute 方法的返回类型,JsValue 是 resolve 方法的返回类型。

TIP

我们需要分开 type Output 和 type JsValue,因为我们无法在 fn compute 中回调 JavaScript 函数,它不在主线程上执行, 所以我们需要在主线程上运行的 fn resolve,根据 Output 和 Env 创建 JsValue 并在 JavaScript 中回调它。

你可以使用底层 API Env::spawn 在 libuv 线程池中生成一个定义的 Task,参见引用中的示例。

除了 compute 和 resolve,您还可以提供 reject 方法,当 Task 遇到错误时,可以执行一些清理工作,例如 unref 一些对象:

lib.rs
rust
use napi::bindgen_prelude::*;
use napi_derive::napi;

pub struct CountBufferLength {
  data: Buffer,
}

impl CountBufferLength {
  pub fn new(data: Buffer) -> Self {
    Self { data }
  }
}

impl Task for CountBufferLength {
  type Output = usize;
  type JsValue = u32;

  fn compute(&mut self) -> Result<Self::Output> {
    if self.data.len() == 10 {
      return Err(Error::new(
        Status::GenericFailure,
        "Random fatal error".to_string(),
      ));
    }
    Ok((&self.data).len())
  }

  fn resolve(&mut self, _: Env, output: Self::Output) -> Result<Self::JsValue> {
    u32::try_from(output)
      .map_err(|_| Error::new(Status::InvalidArg, "buffer length exceeds u32"))
  }
  fn reject(&mut self, _: Env, err: Error) -> Result<Self::JsValue> {
    // catch the error
    if err.status == Status::GenericFailure {
      Ok(0)
    } else {
      Ok(1)
    }
  }
}

#[napi]
pub fn async_count_buffer_length(data: Buffer) -> AsyncTask<CountBufferLength> {
  AsyncTask::new(CountBufferLength { data })
}

您还可以提供一个 finally 方法,在 Task 被 resolved 或 rejected 后执行一些操作:

lib.rs
rust
use napi::bindgen_prelude::*;
use napi_derive::napi;

pub struct CountBufferLength {
  data: Buffer,
}

impl CountBufferLength {
  pub fn new(data: Buffer) -> Self {
    Self { data }
  }
}

impl Task for CountBufferLength {
  type Output = usize;
  type JsValue = u32;

  fn compute(&mut self) -> Result<Self::Output> {
    if self.data.len() == 10 {
      return Err(Error::new(
        Status::GenericFailure,
        "Random fatal error".to_string(),
      ));
    }
    Ok((&self.data).len())
  }

  fn resolve(&mut self, _: Env, output: Self::Output) -> Result<Self::JsValue> {
    u32::try_from(output)
      .map_err(|_| Error::new(Status::InvalidArg, "buffer length exceeds u32"))
  }

  fn reject(&mut self, _: Env, err: Error) -> Result<Self::JsValue> {
    // catch the error
    if err.status == Status::GenericFailure {
      Ok(0)
    } else {
      Ok(1)
    }
  }
  fn finally(self, _: Env) -> Result<()> {
    println!("finally");
    drop(self.data);
    Ok(())
  }
}

#[napi]
pub fn async_count_buffer_length(data: Buffer) -> AsyncTask<CountBufferLength> {
  AsyncTask::new(CountBufferLength { data })
}

TIP

impl Task for AsyncFib 上面的 #[napi] 宏只是为了生成 .d.ts 文件, 如果这里没有定义 #[napi],生成的 TypeScript 类型里, AsyncTask 的返回值类型将是 Promise<unknown>。

AsyncTask

你定义的 Task 不能直接返回给 JavaScript,JavaScript 引擎不知道如何运行和解析你的 struct 的值, AsyncTask 是可以返回给 JavaScript 引擎的 Task 的包装, 可以使用 Task 和可选的 AbortSignal 来创建它。

lib.rs
rust
#[napi]
fn async_fib(input: u32) -> AsyncTask<AsyncFib> {
  AsyncTask::new(AsyncFib { input })
}

⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️

index.d.ts
ts
export function asyncFib(input: number): Promise<number>

结合 AbortSignal 创建 AsyncTask

在某些场景下,你可能想要中止队列中的 AsyncTask,例如对某些计算任务使用 debounce。 您可以给 AsyncTask 传入 AbortSignal ,这样如果 AsyncTask 还没有启动,您就可以中止它。

lib.rs
rust
use napi::bindgen_prelude::AbortSignal;

#[napi]
fn async_fib(input: u32, signal: AbortSignal) -> AsyncTask<AsyncFib> { 
  AsyncTask::with_signal(AsyncFib { input }, signal)
}

⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️

index.d.ts
ts
export function asyncFib(input: number, signal: AbortSignal): Promise<number>

如果在 libuv 开始执行 AsyncTask 之前调用 AbortController.abort, Node-API 可以取消队列中的工作,Promise 会以 name 为 AbortError 的错误 reject。

test.mjs
js
import { asyncFib } from './index.js'

const controller = new AbortController()

asyncFib(20, controller.signal).catch((e) => {
  console.error(e) // Error: AbortError
})

controller.abort()

如果您不知道 AsyncTask 是否需要中止, 您还可以给 AsyncTask 传入 Option<AbortSignal> :

lib.rs
rust
use napi::bindgen_prelude::AbortSignal;

#[napi]
fn async_fib(input: u32, signal: Option<AbortSignal>) -> AsyncTask<AsyncFib> {
  AsyncTask::with_optional_signal(AsyncFib { input }, signal)
}

⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️ ⬇️

index.d.ts
ts
export function asyncFib(
  input: number,
  signal?: AbortSignal | undefined | null,
): Promise<number>

TIP

如果 AsyncTask 已经开始,Node-API 无法取消正在运行的 compute; 如果已经完成,Promise 结果也已经确定。即使取消时机太晚, Rust 中注册的 AbortSignal::on_abort 回调仍会在 JavaScript signal 触发时运行。 当前 adapter 在转换参数时安装处理函数,不会检查 signal.aborted, 因此请传入尚未 abort 的 signal。它还会给 signal.onabort 赋值,覆盖该属性中 已有的处理函数;独立的 JavaScript 监听器请使用 signal.addEventListener('abort', ...)。

ScopedTask

ScopedTask 与 Task 基本相等,但它会把 &'env Env 传给 resolve 和 reject 方法,这样你就可以通过 &'env Env 创建带有生命周期的 JsValue。

例如:

lib.rs
rust
use napi::{JsString, ScopedTask, bindgen_prelude::*};
use napi_derive::napi;

pub struct CountBufferLength {
  data: Buffer,
}

impl CountBufferLength {
  pub fn new(data: Buffer) -> Self {
    Self { data }
  }
}

impl<'env> ScopedTask<'env> for CountBufferLength { 
  type Output = usize;
  type JsValue = JsString<'env>; 

  fn compute(&mut self) -> Result<Self::Output> {
    if self.data.len() == 10 {
      return Err(Error::new(
        Status::GenericFailure,
        "Random fatal error".to_string(),
      ));
    }
    Ok((&self.data).len())
  }

  fn resolve(&mut self, env: &'env Env, output: Self::Output) -> Result<Self::JsValue> {
    env.create_string(format!("{output}"))
  }

  fn reject(&mut self, env: &'env Env, err: Error) -> Result<Self::JsValue> {
    // catch the error
    if err.status == Status::GenericFailure {
      env.create_string("Random fatal error".to_string())
    } else {
      env.create_string("Random error".to_string())
    }
  }

  fn finally(self, _: Env) -> Result<()> {
    drop(self.data);
    Ok(())
  }
}

#[napi]
pub fn async_count_buffer_length(data: Buffer) -> AsyncTask<CountBufferLength> {
  AsyncTask::new(CountBufferLength { data })
}
最后更新于
LongYinan