🐰

【MoonBit】oneshotでイベントをasyncに変換する

に公開

MoonBitの非同期周りの練習がてら、oneshot.mbtというライブラリを作りました。既にmooncaskes.ioにも公開しているので、moon add nuskey8/oneshotで簡単に追加できます。

https://github.com/nuskey8/oneshot.mbt

これは何かというと、要するにJavaScriptのPromise.withResolvers()やC#のTaskCompletionSourceなどの要領で、ワンショットのチャネルを通してイベントをasyncに対応させるためのユーティリティです。Rustのtokioにも似た機能があるほか、oneshotというクレートも存在しています。

使い方はこんな感じ。

async fn main {
  let (put, get) = @oneshot.new()

  // なんか適当なコールバックを受け取る関数
  setTimeout(() => {
    put("Hello!")
  }, 1000)

  let value = get() // このgetが非同期呼び出し
  println(value) // Hello!
}

もう少し複雑な例として、MoonBitの@asyncライブラリを使ったコードも。

async fn main {
  @async.with_task_group(group => {
    // Create a oneshot channel
    let (put, get) = @oneshot.new()

    // Spawn a sender and a receiver
    group.spawn(() => {
      @async.sleep(1000)
      put("Hello!")
    })
    |> ignore
    group.spawn(() => {
      let msg = get()
      println(msg) // Hello!
    })
    |> ignore
  })
}

MoonBitの非同期はいわゆるStructured Concurrencyを採用していて、木構造のグループで非同期のタスクを管理します。ここでは1000ms後に値をセットするタスクと、その値を受け取って出力するタスクの2つを並列に走らせています。

見慣れないコードばかりで複雑に見えるかもしれませんが、ライブラリ本体は以下の一行だけです。

let (put, get) = @oneshot.new()

MoonBitは非同期メソッドの呼び出しにawaitが不要(暗黙的)であるため分かりづらいですが、ここのget()が非同期関数で、put()で値が飛んでくるまで待機するようになっています。

一応同じようなことは標準ライブラリであるmoonbitlang/asyncQueueを用いることでも実現可能ですが、こちらはより高機能である分、単一の値を待ち受けるにはややオーバースペック感があります。また、複数の書き込みが発生する可能性があるため、put()が非同期になるのも少々扱いづらいです。

async fn main {
  let queue = @async.Queue::new(kind=Unbounded)
  defer queue.close()

  setTimeout(() => {
    // putが非同期なので呼び出せない、代わりにtry_putを使う
    queue.try_put("Hello!") |> ignore
  }, 1000)

  let value = queue.get()
  println(value) // Hello!
}

oneshot.mbtはメッセージが1回のみであることを保証するため、put()が他の書き込みを待機する必要がありません。また、内部実装も@async.Queueに比べて非常にシンプルです。

内部実装

MoonBitの非同期周りはまだ安定していないため変更の可能性がありますが、現時点ではプリミティブな非同期処理のためのディレクティブとして%async.run%async/suspendが提供されており、moonbitlang/asyncはこれをベースに実装されています。

///| `run_async` spawn a new coroutine and execute an async function in it
fn run_async(f : async () -> Unit noraise) -> Unit = "%async.run"

///| `suspend` will suspend the execution of the current coroutine.
/// The suspension will be handled by a callback passed to `suspend`
async fn[T, E : Error] suspend(
  // `f` is a callback for handling suspension
  f : (
    // the first parameter of `f` is used to resume the execution of the coroutine normally
    (T) -> Unit,
    // the second parameter of `f` is used to cancel the execution of the current coroutine
    // by throwing an error at suspension point
    (E) -> Unit,
  ) -> Unit,
) -> T raise E = "%async.suspend"

%async.runはそのままですが、%async.suspendがやや複雑に見えるかもしれません。これは要するにJavaScriptのnew Promise((resolve, reject) => {...})だと思えばわかりやすいでしょうか。

ただし、これは基本的に言語が内部で利用するための機能である(ユーザーの利用を想定していない)ため、%async.runで非同期処理を複数同時に走らせたりなどすると容易にクラッシュします。その場合は特にエラーなども吐かないので、これに気づかないと結構ハマります。

oneshot.mbtはこのsuspendを利用し、put()で値が書き込まれたら継続処理をトリガーするような実装になっています。

///|
pub suberror OneShotChannelClosed

///|
/// OneShotChannel internal state
priv struct Channel[T] {
  mut value : T?
  mut resume_fn : ((T) -> Unit)?
}

///|
/// Creates a new oneshot channel.
pub fn[T] new() -> (
  (T) -> Unit raise OneShotChannelClosed,
  async () -> T raise OneShotChannelClosed,
) {
  let channel : Channel[T] = { value: None, resume_fn: None }
  fn writer(v : T) raise OneShotChannelClosed {
    if channel.value is Some(_) {
      raise OneShotChannelClosed
    }
    channel.value = Some(v)
    if channel.resume_fn is Some(resume_fn) {
      resume_fn(v)
      channel.resume_fn = None
    }
  }

  async fn reader() -> T raise OneShotChannelClosed {
    if channel.value is Some(v) {
      return v
    }
    suspend(fn(ok, _ : (OneShotChannelClosed) -> Unit) -> Unit {
      if channel.value is Some(v) {
        ok(v)
      } else {
        channel.resume_fn = Some(ok)
      }
    })
  }

  (writer, reader)
}

中身はこれだけです。今のところMoonBitの非同期処理はシングルスレッドであることを前提としているので、特に排他処理などは行なっていません。

まとめ

元々は一つの記事だけの予定でしたが、なんか筆が乗ってしまったのでもう一つ作ってみました。書いてて楽しい言語は良いですね。

Discussion