Loopinメッセンジャー開発記

Loopinメッセンジャー開発記

1. はじめに(Introduction)

1.1 プロジェクトの背景

従来の古いメッセージングサービスの構造は、変化するウェブ環境に対応することが難しく、保守の面でも多くの課題がありました。本プロジェクト「Loopin」は、このようなSmalltalkシステムを廃止し、最新のVizendプラットフォームへ再構築(Refactoring)することを目標として始動しました。

メッセンジャーサービスの核心的な価値は、「ユーザー間の断絶のないリアルタイムコミュニケーション体験」にあります。これを実現するために直面した技術的課題は、大きく2つありました。1つ目は、従来のHTTPプロトコルによる一方向のリクエスト・レスポンスモデルから脱却し、クライアントとサーバーが永続的なコネクションを確立して、遅延なくメッセージを送受信できるメッセージングチャネルを構築することです。2つ目は、コラボレーションとコミュニケーションの効率を最大化するため、ユーザーが現在接続中かどうかをリアルタイムに判定して画面に表示する接続状態管理(Presence)エンジンを実装することでした。

1.2 プロジェクトの目標

技術を選択する際に最も警戒すべき点は、まさに「オーバーエンジニアリング(Overengineering)」です。大規模なグローバルトラフィックに耐える必要があるアーキテクチャではないにもかかわらず、流行に流されてRedis、Kafkaなどの複雑な分散ソリューションを無条件にインフラへ追加すると、運用負担やデプロイの難易度、さらにはクラウド費用の不必要な増加を招きます。Loopinプロジェクトは、「シンプルさと効率性(KISS - Keep It Simple, Stupid)」を中核的な価値としました。外部インフラ要素を可能な限り排除し、サービスのメインストレージであるリレーショナルデータベース(RDB)だけを多角的に活用することで、リアルタイムPresenceとメッセージのライフサイクルを完全に制御することを最優先の目標として定めました。

2. システムアーキテクチャの概要(System Architecture)

Loopinメッセンジャーのデータフローは、不要なレイヤーを省略し、明確な責任分担に基づいて構成されています。リアルタイム性を保証する通信層と、データの完全性を担保する永続化層が完全に同期して動作します。

image1.png

[図1]LoopinシステムアーキテクチャおよびWebSocketベースのデータフロー図

  • React Frontend: STOMPプロトコルを使用してサーバーとのコネクションを一度だけ確立した後、ブラウザがバックグラウンドで常時待機します。ユーザーの入力が行われると非同期でメッセージを送信し、受信したイベントをコンポーネント単位で再レンダリングして、ユーザーに即座にフィードバックを提供します。

  • Spring Boot Backend: WebSocketエンドポイントの受け入れ、JWTトークンによる接続時のセキュリティ検証(Interceptor)、STOMPアドレス体系を利用したルーティング制御、内部ライフサイクルイベント(Event Listener)の処理を担当します。

  • Relational Database: 永続的なチャットメッセージだけでなく、リアルタイム接続制御スキーマも内部に保持し、セッション状態とセッション履歴情報を安定的に保存します。

3. WebSocketとSTOMPを利用したリアルタイムメッセージング

3.1 STOMPサブプロトコル導入の背景

HTML5標準仕様であるRaw WebSocketは、TCPソケットと同様に双方向通信のトンネルを開くだけで、データの具体的な形式や宛先を指定する規格は存在しません。つまり、クライアントとサーバーが送受信する文字列がチャットデータなのか、入室を宣言するメッセージなのか、エラーメッセージなのかを区別するために、開発者が独自のテキスト解析ルールを定義し、カスタムハンドラーを実装する必要があります。これは不具合が発生しやすく、アーキテクチャを複雑にします。
Loopinはこのような無駄を削減するため、WebSocketフレーム上で動作するSTOMP(Simple Text Oriented Messaging Protocol)サブプロトコルを全面的に導入しました。STOMPは、コマンド(COMMAND)、ヘッダー(Headers)、ボディ(Body)という形式で構造化されたフレームを備えているため、メッセージベースシステムのルーティングマッピングを画期的に直感化します。

3.2 Spring Boot WebSocketインフラの構成ソースコード

Spring BootでSTOMPベースの内蔵メッセージブローカーを有効化し、エンドポイントをマッピングするための実際の実装コードは次のとおりです。

@Configuration
@EnableWebSocketMessageBroker
public class WebSocketConfig implements WebSocketMessageBrokerConfigurer {

    @Override
    public void registerStompEndpoints(StompEndpointRegistry registry) {
        // 프론트엔드가 핸드셰이크를 요청할 종단점 주소 매핑
        // SockJS 폴백을 적용하여 웹소켓이 차단된 프록시나 구형 브라우저 환경 지원
        registry.addEndpoint("/ws")
                .setAllowedOriginPatterns("*")
                .withSockJS()
                .setSessionCookieNeeded(false);
    }

    @Override
    public void configureMessageBroker(MessageBrokerRegistry registry) {
        registry.enableSimpleBroker("/user");
        registry.setUserDestinationPrefix("/user");
        registry.setApplicationDestinationPrefixes("/app");
    }
}

4. データベース(DB)ベースのPresence(接続状態)機能設計

4.1 追加インフラ(Redis)を排除する妥当性とRDBを選択したアーキテクチャ上の意図

多くの開発ガイドでは、リアルタイムセッションデータは揮発性が高いという理由から、RedisのようなインメモリKey-Value NoSQLデータベースに保存することを推奨しています。しかし、Loopinプロジェクトの規模とドメイン特性を詳細に把握すると、これは明らかなリソースの浪費でした。リレーショナルデータベース(RDB)だけを活用して構造を設計した具体的な理由は、次のとおりです。

1つ目は、シンプルなデプロイアーキテクチャを維持できる保守上の利便性です。単一インスタンス、またはプライマリ・レプリカ(Primary-Replica)構成の既存DB構成をそのまま活用すれば、監視アラート体系やバックアップポリシーなどを一元化できます。2つ目は、強固なデータ完全性と統計履歴の追跡性です。ユーザーがいつログインし、ログアウトしたのかという履歴は、単なる状態値以上の意味を持ちます。システム障害やセキュリティ監査の際にはセッションログの履歴が不可欠ですが、RDBはリレーショナルスキーマとインデックスを活用して、精密な時系列ログ分析クエリに対応できます。3つ目は、ドメインデータとの結合が容易であることです。「自分の友達リストのうち、現在接続中のメンバーだけを上位に並べる」といった要件を実装する際、RedisとRDBにデータが分散していると、アプリケーション層で2回のIOを発生させ、手動でメモリ結合を行う必要があります。しかし、RDB単一環境であれば、複雑な構造を必要とせず、単純なJOINまたはINサブクエリ演算によって、ミリ秒以内に処理できます。

4.2 データベーススキーマのモデリング

-- Presence 테이블
CREATE TABLE IF NOT EXISTS active_participant
(
    id             VARCHAR(255) NOT NULL,
    actor_id       VARCHAR(255),
    stage_id       VARCHAR(255),
    pavilion_id    VARCHAR(255),
    entity_version BIGINT       NOT NULL,
    registered_by  VARCHAR(255),
    registered_on  BIGINT       NOT NULL,
    modified_by    VARCHAR(255),
    modified_on    BIGINT       NOT NULL,
    session_id     VARCHAR(255),
    login_time     BIGINT       NOT NULL,
    CONSTRAINT pk_active_participant PRIMARY KEY (id)
);

4.3 Presenceイベントハンドラーの実装

Loopin Presenceテーブルのアクティブカウントをリアルタイムに検証する方式です。

@Component
@RequiredArgsConstructor
public class PresenceEventListener {
    //
    private static final String SIMP_CONNECT_MESSAGE_HEADER = "simpConnectMessage";
    private static final String NATIVE_HEADERS = "nativeHeaders";
    private static final String SIMP_USER_HEADER = "simpUser";
    private final PresenceTask presenceTask;
    
    @EventListener
    public void handleSessionConnect(SessionConnectEvent event) {
        SimpMessageHeaderAccessor headers = SimpMessageHeaderAccessor.wrap(event.getMessage());
        log.info("### STOMP CONNECT Attempt ### SessionId: {}", headers.getSessionId());
    }

    @EventListener
    public void handleSessionConnected(SessionConnectedEvent event) {
        //
        SimpMessageHeaderAccessor headers = SimpMessageHeaderAccessor.wrap(event.getMessage());
        log.info("WebSocket Session Connected! SessionId: {}", headers.getSessionId());
        Map<String, ArrayList<String>> nativeHeaders = this.extractNativeHeaders(headers);
        
        if (nativeHeaders == null || nativeHeaders.isEmpty()) {
            throw new IllegalArgumentException("Native headers cannot be found in Socket header.");
        }
        List<String> actorId = nativeHeaders.get(SIMP_USER_HEADER);
        log.info("WebSocket Session Connected! SessionId: {}, ActorId: {}, Headers: {}", headers.getSessionId(), actorId, nativeHeaders);
        LoginEvent loginEvent = new LoginEvent(actorId.getFirst());
        presenceTask.addParticipant(headers.getSessionId(), loginEvent);
    }

    @EventListener
    public void handleSessionDisconnect(SessionDisconnectEvent event) {
        log.info("WebSocket Session Disconnect! sessionId : {}", event.getSessionId());
        presenceTask.removeParticipant(event.getSessionId());
    }

    @SuppressWarnings("unchecked")
    private Map<String, ArrayList<String>> extractNativeHeaders(SimpMessageHeaderAccessor headers) {
        //
        GenericMessage<?> gm = (GenericMessage<?>) headers.getMessageHeaders().get(SIMP_CONNECT_MESSAGE_HEADER);
        return (Map<String, ArrayList<String>>) gm.getHeaders().get(NATIVE_HEADERS);
    }
}

@Service
@Transactional
@RequiredArgsConstructor
public class PresenceTask {
    //
    private final ActiveParticipantLogic activeParticipantLogic;

    public void addParticipant(String sessionId, LoginEvent event) {
        // 1인 1세션 정책: 동일 사용자의 기존 세션 정보를 모두 정리한 후 새 세션 등록
        activeParticipantLogic.removeByActorId(event.getActorId());

        ActiveParticipantCdo cdo = new ActiveParticipantCdo();
        cdo.setActorId(event.getActorId());
        cdo.setSessionId(sessionId);
        cdo.setLoginTime(System.currentTimeMillis());
        activeParticipantLogic.registerActiveParticipant(cdo);
    }

    public LoginEvent getParticipant(String sessionId) {
        ActiveParticipant activeParticipant = activeParticipantLogic.findBySessionId(sessionId);
        if (activeParticipant == null) {
            return null;
        }
        return new LoginEvent(activeParticipant.getActorId());
    }

    public void removeParticipant(String sessionId) {
        ActiveParticipant activeParticipant = activeParticipantLogic.findBySessionId(sessionId);
        if (activeParticipant != null) {
            activeParticipantLogic.removeActiveParticipant(activeParticipant.getId());
        }
    }

    public Map<String, LoginEvent> getActiveSessions() {
        return activeParticipantLogic.findActiveParticipants(null).stream()
                .collect(Collectors.toMap(
                        ActiveParticipant::getSessionId,
                        ap -> new LoginEvent(ap.getActorId())
                ));
    }

    public boolean isCitizenOnline(String citizenId) {
        return activeParticipantLogic.existsByActorId(citizenId);
    }
}

5. 結論と振り返り(Conclusion & Retrospective)

従来のSmalltalkメッセンジャーシステムを、最新のVizendプラットフォーム上のLoopinメッセンジャーサービスへ再構築する作業では、大規模で高価な外部分散キャッシュレイヤーを使わなくても、複合インデックス戦略と巧みに設計された相互検証用セッションログテーブルのスキーマ設計だけで、安定性とデータ可視性に優れたPresenceエンジンを見事に完成させることができました。

今回のプロジェクトでは、WebSocketとSTOMPを活用してリアルタイムメッセージングの基本的な流れを自ら実装し、フレームワークの背後でセッションがどのように制御されているのか、WebSocketのセッション切断例外を実際にハンドリングしながら、リアルタイム通信の流れなどを段階的に理解できる貴重な機会となりました。華やかな最新技術を追い求めるのではなく、現在の要件とコストに合わせてデータベースだけで問題を解決することで、シンプルなアーキテクチャがもたらすプロジェクト保守上のメリットを改めて学ぶことができました。

BigJumbo

Site footer