3
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?

はじめてな Elixir(38) Zenohex の put/get で「書いたはずの値が読めない」を潰す

3
Last updated at Posted at 2026-08-15

この記事で「はじめに」以下は全部 Claude Code で生成した記事です。もちろんプロンプトは私が指示したものです。書きっぷりが私っぽく、よく出来てます。中で出てくる Github のリポジトリの中も全部 Claude Code が生成してます。急ぎで投稿するほうが良いだろうし、自分の時間使ってられないのでこのような試みをやってみました。お楽しみください。
なお、タイトルや記事の雰囲気は Zenohex が悪いふうに読めますが、元々の Zenoh が持ってる問題が Zenohex 下でも起こったという話です。Zenohex が悪いのではないので念のため。

はじめに

以前 はじめてな Elixir(36) Zenohex (v0.5.1) で Pub/Sub する で紹介した Zenohex を、今度は Pub/Sub ではなく put/get(Zenoh の Storage 機能)で使う個人プロジェクトをやっていました。GenServer が終了するときに自分の状態を Zenoh 経由の Storage に put しておき、別プロセス(別の BEAM でも構いません)が起動時にその状態を get で拾って引き継ぐ、という「お手軽な Handoff」を作っていたのですが、たまに引き継いだ値が1つ古いことに気づきました。掘っていくとなかなか面白い非対称性に行き当たったので、原因と対策を共有します。

現象を再現する

話を単純にするため、put した直後に同じ key へ get するだけのループを書きます。

mix.exs (抜粋)
defp deps do
  [
    {:zenohex, "~> 0.9.0"}
  ]
end
put_get_race.exs
{:ok, session_id} = Zenohex.Session.open(config)

Enum.each(1..2000, fn i ->
  payload = Integer.to_string(i)
  :ok = Zenohex.Session.put(session_id, key, payload)

  {:ok, replies} = Zenohex.Session.get(session_id, key, 3_000, consolidation: :latest)

  case Enum.find(replies, &match?(%Zenohex.Sample{}, &1)) do
    %Zenohex.Sample{payload: ^payload} -> :ok
    %Zenohex.Sample{payload: other} -> IO.puts("stale! put #{payload} but got #{other}")
    nil -> IO.puts("no reply at all")
  end
end)

2000回のうちのごく一部で stale! が出ます。手元では 78/2000(3.9%)ほどでした。ただし面白いのは、直後にもう一度 get すればほぼ100%(手元の観測では最短1ms後)で正しい値が読めることです。つまり「値が消える」のではなく「反映にわずかなラグがある」だけのようでした。

原因を調べる

Zenohex.Session.put/4 は Rustler 経由で zenoh-rust の put をそのまま呼んでいます。実装(zenohex の NIF 実装)を覗くとこうなっています。

fn session_put(...) -> rustler::NifResult<rustler::Atom> {
    ...
    publication_builder
        .apply_opts(opts)?
        .wait()   // ← ここで待っているのはローカルの送信処理の完了だけ
        ...
    Ok(rustler::types::atom::ok())
}

.wait() が待っているのは「ローカルのセッションが送信キューに積み終えた」ことだけで、相手(storage を持つ zenohd ルーター)が実際に受け取って反映し終えたことまでは保証していません。一方 session_getschedule = "DirtyIo" が付いていて、タイムアウト付きでちゃんとリモートからの返信を待つ、正真正銘のリクエスト/レスポンスです。

Elixir の GenServer に例えるなら、putcastgetcall くらいの非対称性がある、と考えるとしっくりきます。cast を送った直後に「もう反映されてるはず」と決め打ちして call すると、たまに反映前に問い合わせてしまう、というわけです。

これは私だけが気づいた話ではなく、zenoh 本家にもちょうど 同じ課題を扱う Issue が立っていました(本記事執筆時点でまだ Open)。Issue の一文がそのまま今回の現象の説明になっています。

Zenoh's pub/sub path is fire-and-forget — session.put() returns when the message is sent, not when it's stored.

さらに厄介なケース

ついでにもうひとつ見つけたことがあります。一度も put されたことのない key に get すると、{:ok, []} ではなくエラーになることがあります。

iex> Zenohex.Session.get(session_id, "never/published/key", 3_000, consolidation: :latest)
{:error, "native/zenohex_nif/src/session.rs:371: channel is empty and closed"}

これは stale read とは別種の問題で、「storage が空だから空リストが返る」のではなく、内部の応答チャネルまわりでエラーとして表面化するようです。{:ok, replies} = Zenohex.Session.get(...) のように決め打ちでパターンマッチしていると、ここで素直に落ちます。

ZenohAckPut を作った

というわけで、「put した直後に同じ key へ get して、書き込んだ値が読み返せることを確認してから返る」ラッパーを書きました。上記のエラーケースも「まだ確認できていない」として扱うようにしてあります。

zenohackput.ex
defmodule ZenohAckPut do
  @default_confirm_timeout_ms 3_000
  @default_confirm_interval_ms 1
  @default_query_timeout_ms 3_000

  def put(session_id, key_expr, payload, put_opts \\ [], confirm_opts \\ []) do
    confirm_timeout_ms =
      Keyword.get(confirm_opts, :confirm_timeout_ms, @default_confirm_timeout_ms)

    confirm_interval_ms =
      Keyword.get(confirm_opts, :confirm_interval_ms, @default_confirm_interval_ms)

    query_timeout_ms = Keyword.get(confirm_opts, :query_timeout_ms, @default_query_timeout_ms)

    with :ok <- Zenohex.Session.put(session_id, key_expr, payload, put_opts) do
      deadline = System.monotonic_time(:millisecond) + confirm_timeout_ms
      confirm(session_id, key_expr, payload, query_timeout_ms, confirm_interval_ms, deadline)
    end
  end

  defp confirm(session_id, key_expr, payload, query_timeout_ms, confirm_interval_ms, deadline) do
    if fetch(session_id, key_expr, query_timeout_ms) == payload do
      :ok
    else
      if System.monotonic_time(:millisecond) >= deadline do
        {:error, :not_confirmed}
      else
        Process.sleep(confirm_interval_ms)
        confirm(session_id, key_expr, payload, query_timeout_ms, confirm_interval_ms, deadline)
      end
    end
  end

  defp fetch(session_id, key_expr, query_timeout_ms) do
    case Zenohex.Session.get(session_id, key_expr, query_timeout_ms, consolidation: :latest) do
      {:ok, replies} ->
        case Enum.find(replies, &match?(%Zenohex.Sample{}, &1)) do
          %Zenohex.Sample{payload: found_payload} -> found_payload
          nil -> nil
        end

      {:error, _reason} ->
        nil
    end
  end
end

使い方はこんな感じです。

iex> ZenohAckPut.put(session_id, "key/expr", "payload")
:ok

戻り値は3種類です。

  • :ok : put が成功し、read-after-write の確認も取れた
  • {:error, :not_confirmed} : put 自体は成功したが、タイムアウト以内に確認が取れなかった(put が失敗したわけではないので注意)
  • {:error, reason} : put そのものが失敗した

先程のループを ZenohAckPut.put に差し替えて2000回まわしたところ、stale read も未確認タイムアウトも0件になりました。

これは単体のモジュールとして GitHub に切り出してあります。素の put/get で現象を再現するスクリプトや、ZenohAckPut.put がそれを解消することを確認するスクリプトも同梱しています。

まだ Hex には公開していないので、使う場合は git 依存として mix.exs に追加して
ください。

mix.exs
defp deps do
  [
    {:zenohackput, git: "https://github.com/kikuyuta/zenohackput.git"}
  ]
end

まとめ

Zenoh の put は fire-and-forget、get はリクエスト/レスポンスという非対称性があり、put 直後の get はまれに(手元の計測で数%程度)古い値を読んでしまうことがあります。原因は zenoh 本家でも Issue になっている既知のギャップで、今のところプロトコルレベルでの確認手段はありません。アプリケーション層で「putした直後にgetして確認する」だけのシンプルなラッパーで実用上は十分に解消できました。

Zenoh の put/get を GenServer の状態の保存・復元のような「即座に読めてほしい」用途で使う場合は、この非対称性を頭の片隅に置いておくと良さそうです。

参考文献

3
0
0

Register as a new user and use Qiita more conveniently

  1. You get articles that match your needs
  2. You can efficiently read back useful information
  3. You can use dark theme
What you can do with signing up
3
0

Delete article

Deleted articles cannot be recovered.

Draft of this article would be also deleted.

Are you sure you want to delete this article?