この記事で「はじめに」以下は全部 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 するだけのループを書きます。
defp deps do
[
{:zenohex, "~> 0.9.0"}
]
end
{: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_get は schedule = "DirtyIo" が付いていて、タイムアウト付きでちゃんとリモートからの返信を待つ、正真正銘のリクエスト/レスポンスです。
Elixir の GenServer に例えるなら、put は cast、get は call くらいの非対称性がある、と考えるとしっくりきます。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 して、書き込んだ値が読み返せることを確認してから返る」ラッパーを書きました。上記のエラーケースも「まだ確認できていない」として扱うようにしてあります。
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 に追加して
ください。
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 の状態の保存・復元のような「即座に読めてほしい」用途で使う場合は、この非対称性を頭の片隅に置いておくと良さそうです。