🗿
[capnp][kj] Reader のライフタイムの検証
概要 & 結論
capnp のサブスクリプションにおいて Reader のライフタイムが気になったのでサンプルコードを実装し確認を行った. 特に気にしている点は保管しておいたReaderが新しいサブスクライブ通知によって上書きされないかという点を確認したかった.
結果としては上書きされないということがわかった.
コードは以下の通り
コードの説明
コートの一部を抜粋して説明
スキーマ
サブスクライブ登録を送るとidなどのデータを送信するためのシンプルなスキーマにしている.
# 通知のデータ構造
struct Notification {
id @0 :UInt64;
timestamp @1 :Int64;
kind @2 :Text;
payload @3 :Data;
}
# ポーリング用の通知受信インターフェース
interface PollingNotificationReceiver {
# クライアントがこのメソッドを実装し、サーバーがコンテキスト経由で通知を送信
onNotification @0 (notification :Notification) -> ();
}
# ポーリング購読セッション
interface PollingSubscription {
cancel @0 () -> ();
}
# ポーリング用のNotifier
interface PollingNotifier {
# クライアントがreceiverを渡し、サーバーがそのreceiverに通知を送信
subscribe @0 (filter :Text, receiver :PollingNotificationReceiver)
-> (subscription :PollingSubscription);
}
サーバー側のコード
サーバーからの通知を送信する.
1秒ごとにNotificationを送信. sendNotifications 内で Notification のデータを設定しており, id は一回送信する事に1づつ値が増えていく.
void startNotificationLoop() {
if (!timer_ptr_) return;
// 定期的に通知を送信するループを開始
auto promise =
sendNotifications()
.then([this]() { return timer_ptr_->afterDelay(1 * kj::SECONDS); })
.then([this]() {
startNotificationLoop(); // 再帰的に継続
})
.catch_([](kj::Exception&& e) {
LOG_COUT << "Notification loop error: "
<< e.getDescription().cStr() << std::endl;
});
task_set_->add(kj::mv(promise));
}
クライアント側のコード
通知を受け取ると1回目だけ, 受け取ったデータをもとにReader の内容を表示続けるものである.
kj::Promise<void> onNotification(OnNotificationContext context) override {
const auto notification = context.getParams().getNotification();
LOG_COUT << "[Context Notification] id=" << notification.getId()
<< ", kind=" << notification.getKind().cStr()
<< ", timestamp=" << notification.getTimestamp() << std::endl;
if (!is_start_) {
is_start_ = true;
recursivePrint(notification);
}
return kj::READY_NOW;
}
/**
* @brief 通知データを再帰的に処理する
* @details 100ミリ秒間隔で同じ通知データを継続的に出力し続ける
* @param reader 処理対象の通知データリーダー
*/
void recursivePrint(const ::Notification::Reader& reader) {
LOG_COUT << "[recursivePrint] id=" << reader.getId() << std::endl;
auto promise =
timer->afterDelay(500 * kj::MILLISECONDS).then([this, reader]() {
recursivePrint(reader);
});
taskSet->add(kj::mv(promise));
}
出力データ
クライアント側の出力では, 新たに新しい通知を受け取ったとしても
[08:56:46.345]polling_client.cpp:97]Starting Polling Notifier client...
[08:56:46.348]polling_client.cpp:116]Sending Polling Subscribe request...
[08:56:46.349]polling_client.cpp:122]Polling Subscribe request sent.
[08:56:46.349]polling_client.cpp:125]Polling Subscribe response received.
[08:56:46.349]polling_client.cpp:136][Client] Polling client finished.
[08:56:47.146]polling_client.cpp:48][Context Notification] id=0, kind=polling_demo, timestamp=1753487807146
[08:56:47.147]polling_client.cpp:65][recursivePrint] id=0
[08:56:47.647]polling_client.cpp:65][recursivePrint] id=0
[08:56:48.148]polling_client.cpp:65][recursivePrint] id=0
[08:56:48.148]polling_client.cpp:48][Context Notification] id=1, kind=polling_demo, timestamp=1753487808148
[08:56:48.649]polling_client.cpp:65][recursivePrint] id=0
[08:56:49.150]polling_client.cpp:65][recursivePrint] id=0
[08:56:49.150]polling_client.cpp:48][Context Notification] id=2, kind=polling_demo, timestamp=1753487809150
[08:56:49.651]polling_client.cpp:65][recursivePrint] id=0
[08:56:50.152]polling_client.cpp:65][recursivePrint] id=0
[08:56:50.152]polling_client.cpp:48][Context Notification] id=3, kind=polling_demo, timestamp=1753487810152
[08:56:50.653]polling_client.cpp:65][recursivePrint] id=0
[08:56:51.154]polling_client.cpp:65][recursivePrint] id=0
[08:56:51.154]polling_client.cpp:48][Context Notification] id=4, kind=polling_demo, timestamp=1753487811154
[08:56:51.350]polling_client.cpp:130][Client] Cancelling polling subscription...
[08:56:51.655]polling_client.cpp:65][recursivePrint] id=0
[08:56:52.155]polling_client.cpp:65][recursivePrint] id=0
[08:56:52.656]polling_client.cpp:65][recursivePrint] id=0
[08:56:53.156]polling_client.cpp:65][recursivePrint] id=0
最後に
TODO : Reader 内部の実装がわかっていないため確認する必要がある. メモリ確保や,開放のタイミングを確認しReader のライフタイムをはっきりとさせる.
Discussion