LavelとWorkermanを統合したチャットルームの実装

テストツール: http://www.blue-zero.com/WebSocket/

<?php

namespace App\Console\Commands;

use Illuminate\Console\Command;
use Workerman\Worker;
use App\Services\ChatService;

class ChatServer extends Command
{
    protected $chatWorker;
    
    // サポートされる操作パラメータ
    // start: サーバー起動
    // stop: サーバー停止
    // reload: ロジックコードの再起動(core workerman_initは再起動不可)
    // status: ステータス確認
    // connections: 接続状態確認(Workermanバージョン>=3.5.0が必要)
    protected $validActions = ['start', 'stop', 'reload', 'status', 'connections'];

    /**
     * コンソールコマンドのシグネチャと説明
     */
    protected $signature = 'chat:server {action}';
    protected $description = 'チャットサーバーを管理します';

    public function __construct()
    {
        parent::__construct();
    }

    /**
     * コンソールコマンドの実行
     */
    public function handle()
    {
        $action = $this->argument('action');
        
        if (!in_array($action, $this->validActions)) {
            $this->error('無効な操作です');
            return;
        }
        
        // Workermanの初期化
        ChatService::initialize($action);
    }
}
<?php

namespace App\Services;

use App\Services\BaseService as Base;
use Illuminate\Support\Facades\DB;
use App\Services\CommonService;
use Workerman\Worker;
use Workerman\Lib\Timer;
use App\Models\OperationLog;
use App\Models\User;

class ChatService extends Base
{
    // グローバル接続数
    static $totalConnections = 0;
    // ルームあたりの最大接続数
    static $maxConnectionsPerRoom;
    // ルーム内のユーザー接続IDマッピング
    static $roomConnectionsMap = [];

    public static function initialize($action = null)
    {
        global $argv;
        
        $argv[0] = 'workerman:websocket';
        $argv[1] = $action;
        
        // ハートビート間隔(秒)
        define('HEARTBEAT_INTERVAL', 30);
        
        // ワーカーの初期化
        $worker = new Worker("websocket://172.17.1.247:9090");
        $worker->name = 'ChatServer';
        // ワーカープロセス数(テスト環境では4)
        $worker->count = 4;
        
        // 接続時の処理
        $worker->onConnect = function($connection) {
            self::$totalConnections++;
            self::handleConnection($connection);
        };
        
        // メッセージ受信時の処理
        $worker->onMessage = function($connection, $data) {
            self::handleMessage($connection, $data);
        };
        
        // 接続終了時の処理
        $worker->onClose = function($connection) {
            self::$totalConnections--;
            self::handleDisconnection($connection);
        };
        
        // ワーカースタート時の処理
        $worker->onWorkerStart = function($worker) {
            // 必要に応じてタイマーを設定
        };
        
        // ワーカーの実行
        Worker::runAll();
    }

    // 接続時の処理
    public static function handleConnection($connection)
    {
        // 必要に応じて接続時の初期処理を行う
    }

    // メッセージ受信時の処理
    public static function handleMessage($connection, $data)
    {
        if (!empty($data)) {
            if (self::isValidJson($data)) {
                $messageData = json_decode($data, true);
                
                switch ($messageData['action_type']) {
                    case 'ping':
                        // ハートビート応答
                        $connection->send(json_encode([
                            'code' => 200, 
                            'message' => 'サーバーは稼働中です',
                            'data' => [],
                            'connections' => self::$totalConnections
                        ]));
                        break;
                        
                    case 'login':
                        // ユーザーログイン処理
                        if ($messageData['is_login'] == 1) {
                            // 匿名ログイン
                            self::$roomConnectionsMap[$messageData['room_id']][$connection->id]['username'] = '匿名ユーザー' . $connection->id;
                            
                            $connection->send(json_encode([
                                'code' => 200, 
                                'message' => '匿名ログイン成功',
                                'data' => self::$roomConnectionsMap,
                                'connections' => self::$totalConnections
                            ]));
                        } elseif ($messageData['is_login'] == 2) {
                            // 認証済みユーザーログイン
                            $user = User::find($messageData['user_id']);
                            if (!$user) {
                                $connection->send(json_encode([
                                    'code' => 201, 
                                    'message' => '無効または不正なユーザーIDです',
                                    'data' => [],
                                    'connections' => self::$totalConnections
                                ]));
                            } else {
                                $userData = $user->toArray();
                                self::$roomConnectionsMap[$messageData['room_id']][$connection->id]['username'] = 
                                    empty($userData['realname']) ? $messageData['user_id'] : $userData['realname'];
                                
                                $connection->send(json_encode([
                                    'code' => 200, 
                                    'message' => 'ログイン成功',
                                    'data' => self::$roomConnectionsMap,
                                    'connections' => self::$totalConnections
                                ]));
                            }
                        } else {
                            $connection->send(json_encode([
                                'code' => 201, 
                                'message' => 'ログインタイプが不正です',
                                'data' => [],
                                'connections' => self::$totalConnections
                            ]));
                        }
                        break;
                        
                    case 'broadcast_to_all':
                        // ルーム内の全ユーザーにブロードキャスト(送信者を除く)
                        foreach ($connection->worker->connections as $client) {
                            if (isset(self::$roomConnectionsMap[$messageData['room_id']][$client->id]) && 
                                $client->id != $connection->id) {
                                $client->send($messageData['message']);
                            }
                        }
                        break;
                        
                    case 'broadcast_to_one':
                        // 特定ユーザーへの送信(実装例では未完成)
                        break;
                        
                    default:
                        $connection->send(json_encode([
                            'code' => 201, 
                            'message' => '不明なアクションタイプです',
                            'data' => [],
                            'connections' => self::$totalConnections
                        ]));
                }
            }
        } else {
            $connection->send(json_encode([
                'code' => 201, 
                'message' => '空のデータです',
                'data' => [],
                'connections' => self::$totalConnections
            ]));
        }
    }

    // 接続終了時の処理
    public static function handleDisconnection($connection)
    {
        // 接続IDに対応するルーム情報を削除
        foreach (self::$roomConnectionsMap as $roomId => $connections) {
            if (isset($connections[$connection->id])) {
                unset(self::$roomConnectionsMap[$roomId][$connection->id]);
            }
        }
    }
    
    // JSON検証ヘルパーメソッド
    private static function isValidJson($string)
    {
        json_decode($string);
        return (json_last_error() === JSON_ERROR_NONE);
    }
}
{"action_type":"login","is_login":1,"room_id":101}  // 匿名ログイン


{"action_type":"broadcast_to_all","room_id":100,"message":"こんにちは、チャットルーム!"} // ルーム内全員にメッセージ送信
注意点:<br></br>1. $connectionは現在の接続情報を保持します<br></br>2. ワーカーIDに基づいてルームを分割していないため、単一ワーカーの接続数を超えた場合に問題が発生する可能性があります<br></br>3. 実際の負荷に応じて各ワーカーの最大接続数を設定できます。これはシンプルなテストデモなので詳細な実装は省略しています<br></br>4. セッションと連携してデータを処理する必要がありますが、この例では送信データに基づいてユーザーを識別しています。ログイン時にセッションデータを処理し、ユーザーリストが必要な場合はルームIDに対応する全ユーザー名を返すようにします<br></br>5. Workermanの使用は比較的簡単ですが、Lavelとの統合には問題があります。例えば、メッセージ通知とチャットルームを一緒に使用できない場合、別の方法が必要になります。単機能の場合は直接統合できます。代替案として、artisanエントリーファイルを複製して新しいエントリポイント(例: artisan1)を追加する方法があります。テストの結果、問題なく動作しますが、最適な解決策ではありません。使用する場合は一時的な対策として利用してください<br></br><br></br>Gatewayを使用するとより簡単に実装できます

2019年7月12日09:43:56

注: 上記は一時的なテストコードです。本番のコードではtry-catchで例外とエラーを処理してください

タグ: laravel workerman websocket チャットルーム

7月26日 18:04 投稿