Ruby
 Computer >> コンピューター >  >> プログラミング >> Ruby

Rubyで学ぶリアルタイム通信:プッシュ、Pub/Sub、WebSocketの仕組み

Railsでリアルタイム機能を実装するのは、Action Cableのようなライブラリのおかげで格段に簡単になりました。このAppSignal Academyの記事では、リアルタイム更新の実装方法を掘り下げながら、最小構成のWebSocketサーバーを実際に構築し、その内部動作(アンダー・ザ・フード)を探っていきます。

今回構築するアプリケーションは、データをプッシュで送信し、Pub/Sub(パブリッシュ/サブスクライブ)モデルをWebSocket上で利用します。コードに入る前に、まずこれら3つの概念について整理しておきましょう。

  • プッシュ(Push):受信側がデータをポーリング(定期的に問い合わせ)するのではなく、送信側から能動的にデータを届ける方式を指します。株価表示、チャットアプリケーション、運用監視コンソールなど、リアルタイム性が求められる更新には必須の仕組みです。

  • Pub/Sub(パブリッシュ/サブスクライブ):発行(publish)と購読(subscribe)によるデータ配信の相互作用モデルです。1990年代にTIBCOがウォール街で広めたことで有名になりました。受信者は特定のサブジェクト(トピック)を購読し、パブリッシャーがそのサブジェクトへデータを送信するのを待ちます。ワイルドカードによるパターンマッチングで、発行されたメッセージとリスナーを照合するのが一般的ですが、シンプルな実装では名前付きチャネルのみを使う場合もあります。筆者自身もTIBCOの初期の頃から関わってきたため、ワイルドカードパターンマッチングの柔軟性には思い入れがあります。

  • WebSocket:主にWebブラウザとアプリケーション間でデータをやり取りするためのプロトコルです。HTTP接続をWebSocket接続にアップグレードすると、両エンドポイント間で双方向のデータ送受信が可能になります。WebSocketを使えば、アプリケーションからブラウザへデータをプッシュできるだけでなく、ブラウザ上のJavaScriptコードからアプリケーションへ、POSTやPUT以外の手段でデータを送り返すことも可能です。なかなか便利ですよね?

WebSocketの内部動作

ここで、WebSocketサーバーがどのように動作するのか、具体例を見てみましょう。ブラウザ上のクライアントは、JavaScriptコードを使ってサーバーへのWebSocket接続を試みます。

var sock = new WebSocket("ws://" + document.URL.split("/")[2] + "/upgrade");

サーバーは、アップグレード要求の指示を含むHTTPリクエストを受け取ります。通常、アップグレードするかどうかの判断はアプリケーション側に委ねられます。その方法は、アプリに提供されるAPIによって異なります。Rack対応のサーバーであれば、ソケットをハイジャックして開発者がプロトコルの詳細をすべて処理できるオプションが用意されています。また、ある提案されたPRによれば、アップグレードに対するレスポンスを返すだけで十分とされています。

アップグレードは、サーバーとクライアント間の一連のやり取り(ハンドシェイク)です。すべてのブラウザや一部のサーバー用gemは、こうした詳細を隠蔽しています。接続が確立されると、WebSocketプロトコルに従ってメッセージを交換できるようになります。

内部では、エンコード、デコード、そしてメッセージ交換プロトコルという「魔法」が処理を担っています。メッセージはSHA1で暗号化された固定幅のバイナリ構造体に、末尾にペイロードを付加した形式で構成されます。WebSocketプロトコルには、ping/pongによるハートビートや、接続の開始・終了時のメッセージ交換など、複数のメッセージタイプとやり取りが定義されています。これこそが、接続ハイジャック方式を使わないサーバーが裏で自動的に行っている処理なのです。

実装してみよう

ここでは、起動すると現在時刻を待ち受けている全クライアントに配信し続けるクロックスレッドの例を使います。サーバーにはAgooを採用します。高速でありながら、複雑さを最小限に抑えられるためです。

まずはJavaScriptでクライアントを作成し、HTMLページ上に現在時刻を表示します。新しいWebSocketを作成したら、onopenコールバックを設定してステータス表示用のHTML要素を書き換えます。onmessageコールバックはmessage要素を更新します。コールバックは、パブリッシュ/サブスクライブのような非同期処理を扱う際によく使われるデザインパターンです。

<!-- websocket.html -->
<html>
  <body>
    <p id="status">...</p>
    <p id="message">... waiting ...</p>
 
    <script type="text/javascript">
      var sock = new WebSocket(
        "ws://" + document.URL.split("/")[2] + "/upgrade"
      );
      sock.onopen = function () {
        document.getElementById("status").textContent = "connected";
      };
      sock.onmessage = function (msg) {
        document.getElementById("message").textContent = msg.data;
      };
    </script>
  </body>
</html>

クライアントが完成したので、次はサーバーを実装しましょう。Rack APIを使用するRubyアプリケーションです。Clockクラス自体が、/upgradeパスへのすべてのHTTPリクエストを処理するハンドラになります。リクエストがアップグレード要求であれば、HTTPステータスコード200(Success)を返し、そうでなければ404(Page Not Found)を返します。#callメソッドのもう一つの役割は、WebSocketハンドラの割り当てです。

class Clock
  def self.call(env)
    unless env['rack.upgrade?'].nil?
      env['rack.upgrade'] = Clock
      [ 200, { }, [ ] ]
    else
      [ 404, { }, [ ] ]
    end
  end
end

このAPIはコールバックベースです。今回のサーバーで必要なのは#on_openコールバックだけです。これにより、「time」というサブジェクトへの購読を作成できます。メッセージは、サブジェクト(トピック)で識別されるチャネルを通じて交換されます。#on_openは、WebSocket接続が確立されたときに呼び出されます。

class Clock
  # ...
 
  def self.on_open(client)
    client.subscribe('time')
  end
end

続いて、1秒ごとに現在時刻をパブリッシュするスレッドで配信を開始します。Agoo.publishの呼び出しにより「time」サブジェクト上でメッセージが送信され、すべての購読者がそのメッセージを受け取ります。サーバーは購読と接続を管理し、メッセージをJavaScriptクライアントへ配信します。クライアント側ではHTML要素が更新されます。

require 'agoo'
 
Thread.new {
  loop do
    now = Time.now
    Agoo.publish('time', "%02d:%02d:%02d" % [now.hour, now.min, now.sec])
    sleep(1)
  end
}

あとは、サーバーを初期化して起動するコードを追加するだけです。Agoo::Server.handle(:GET, '/upgrade', Clock)の呼び出しにより、サーバーは/upgradeURLパスへのHTTP GETリクエストを待ち受け、それらのリクエストをClockクラスに渡すようになります。これにより、ルーティングをRubyの外側で処理できるため、パフォーマンスと柔軟性が向上します。

Agoo::Server.init(6464, '.', thread_count: 0)
Agoo::Server.handle(:GET, '/upgrade', Clock)
Agoo::Server.start

もう少しで完成です。次のコマンドでサーバーを起動してください。

$ ruby pubsub.rb

以下のようなログが出力されれば、サーバーが正常に起動し、ポート6464で待ち受けていることを示しています。

I 2018/08/14 19:49:45.170618000 INFO: Agoo 2.5.0 with pid 40366 is listening on https://:6464.

動作確認をしてみよう

それでは、https://localhost:6464/websocket.html を開いてみましょう。接続確立時に一瞬ちらついた後、接続ステータスと現在時刻が表示されるはずです。時刻はクロックが進むにつれて1秒ごとに更新されていきます。

connected

19:50:12

おめでとうございます!これでパブリッシュ&サブスクライブ型のWebアプリケーションが完成しました ;-)

今回はWebSocketの使い方を見てきました。同じ目的を達成する選択肢として、Server-Sent Events(SSE)もあります。完全なソースコードのサンプルにはSSEの実装も含まれています。さらに詳しく知りたい方は、今回使用したAgooサーバーや、Iodine WebSocket Serverをチェックしてみてください。

ご質問やコメントがあれば、お気軽に@AppSignalまでご連絡ください。

  1. Rubyでのログ出力をマスターする:LoggerとLogrageの使い方徹底解説

    Rubyでのログ出力入門:LoggerとLogrageの使い方 ロギングは、アプリケーション開発において最も重要なタスクの一つです。ログは以下のような場面で活用されます。 アプリ内部で何が起きているかを把握したいとき アプリケーションを監視したいとき 特定のデータに関するメトリクスを収集したいとき 新しいプログラミング言語を学ぶ際、最初に選ばれるのはその言語がネイティブに備えているロギング機構でしょう。標準機能は通常、扱いやすく、ドキュメントも充実しており、コミュニティでも広く使われています。 ただし、ログデータの内容や扱い方は、企業の方針、ビジネスの性質、アプリケーションの種類に

  2. JedisライブラリでRedisのPub/Subシステムを実装する方法【Javaサンプルコード付き】

    このチュートリアルでは、Java向けクライアントライブラリ「Jedis」を使用して、RedisのPub/Sub(パブリッシュ/サブスクライブ)システムを実装する方法を解説します。 Jedisライブラリとは Jedisは、Redisデータストア用のJavaクライアントライブラリです。軽量で非常に扱いやすく、Redis 2.8.x、3.x.x以降のバージョンと完全な互換性があります。シンプルなAPI设计で、Redisの各種コマンドを直感的に操作できるのが特徴です。 RedisのPub/Subシステムとは Redisは、Publish/Subscribe(パブリッシュ/サブスクライブ)型のメッセージ