Skip to main content

실시간 파이프라인 아키텍처 (Real-time Pipeline)

본 문서는 vivid-ai 프로젝트의 실시간 데이터 처리 파이프라인 아키텍처를 설명합니다. 이 파이프라인은 ComfyUI와 ai-node-agent에서 발생하는 이벤트를 relay를 거쳐 프론트엔드까지 실시간으로 전달하여 사용자에게 작업 진행 상황을 투명하게 보여주는 것을 목표로 합니다.

1. 개요

실시간 파이프라인은 public relay + shared Redis를 핵심 브리지로 사용하는 구조를 기반으로 합니다. ai-node-agent는 AWS private Redis에 직접 붙지 않고, 단일 public relay에 이벤트를 전달합니다. relay는 이벤트를 정규화한 뒤 shared Redis에 반영하고, backend는 기존처럼 Redis를 구독하여 프론트엔드에 fan-out 합니다.

  • 주요 기술 스택:
    • Publisher: websocket-client / requests (Python)
    • Public Relay: agent 전용 WebSocket + HTTP ingress
    • Internal Broker: Redis Pub/Sub
    • Subscriber: NestJS WebSocket Gateway (socket.io, redis)
    • Real-time Client: Socket.IO Client (Frontend)

2. 아키텍처 다이어그램

3. 구성 요소 (Components)

3.1. AiNodeAgent (Publisher)

  • 위치: ai-server generation runtime
  • 역할: ComfyUI 실행 이벤트와 job 상태를 수집하여 public relay로 전달하는 에이전트입니다.
  • 주요 기능:
    1. ComfyUI WebSocket 연결: 로컬 ComfyUI의 WebSocket(ws://comfyui:8188/ws)에 연결하여 모든 실시간 이벤트를 수신합니다.
    2. 이벤트 전달: 수신한 ComfyUI 이벤트를 relay의 WebSocket/HTTP ingress로 전달합니다.
    3. 메타데이터 추가: 전달 전, 이벤트 페이로드에 nodeId, env, userId 등 메타데이터를 추가하여 어떤 서버의 어떤 작업인지 식별할 수 있도록 합니다.
      • userId는 SQS 메시지 처리 시 prompt_iduserId를 Redis에 매핑해 둔 값을 참조하여 가져옵니다.
    4. 자동 재연결: relay 연결이 끊어지면 자동 재연결을 시도하여 안정성을 보장합니다.

3.2. Agent Relay

  • 위치: AWS public ingress 뒤 relay 서비스
  • 역할: 외부 AiNodeAgent와 shared Redis 사이를 중개하는 정규화 계층입니다.
  • 주요 기능:
    1. Agent 전용 ingress 제공: 외부 노드는 relay 주소 하나만 알면 됩니다.
    2. Transport 분리: heartbeat/progress/node.monitor는 WebSocket, 완료/실패/결과 등록은 HTTP로 분리할 수 있습니다.
    3. 메시지 정규화: 수신 이벤트를 Redis 키/채널 모델에 맞게 변환합니다.
    4. 보안/제어: agent token 검증, rate limit, idempotency, sequence 제어를 수행합니다.

3.3. Redis (Internal Broker)

  • 위치: AWS ElastiCache for Redis
  • 역할: relay와 backend 사이를 중개하는 내부 메시지 브로커입니다.
  • 채널: ai-comfyui-events 채널을 사용하여 모든 ComfyUI 관련 이벤트를 전달합니다.

3.4. AiRealtimeGateway (Subscriber & Broadcaster)

  • 위치: 백엔드 서버 (backend-vivid-ai)
  • 역할: Redis로부터 이벤트를 구독하고, 적절한 프론트엔드 클라이언트에게 정보를 전달(Broadcast)하는 WebSocket 게이트웨이입니다.
  • 주요 기능:
    1. Redis 구독: ai-comfyui-events 채널을 구독하여 AiNodeAgent가 발행한 모든 이벤트를 수신합니다.
    2. 이벤트 필터링 및 전달: 수신한 이벤트의 userId 유무에 따라 적절한 Socket.IO room으로 메시지를 전달합니다.
      • admin room: userId가 없는 전역 상태 이벤트(서버 ONLINE/OFFLINE 등)는 admin room에 있는 모든 클라이언트에게 전달됩니다.
      • user-{userId} room: userId가 있는 특정 사용자의 작업 진행 이벤트는 해당 user-{userId} room에 있는 클라이언트에게만 전달됩니다.
    3. 클라이언트 관리: 프론트엔드 클라이언트의 WebSocket 연결/해제 및 room 참여/이탈을 관리합니다.

3.5. RedisIoAdapter

  • 위치: 백엔드 서버 (backend-vivid-ai)
  • 역할: NestJS의 기본 WebSocket 어댑터를 대체하여 Redis 기반의 어댑터를 사용하게 합니다.
  • 주요 기능:
    • 백엔드 서버가 여러 인스턴스(예: EKS의 여러 Pod)로 수평 확장되더라도, Redis를 통해 Socket.IO 메시지를 모든 인스턴스에 걸쳐 브로드캐스팅할 수 있도록 보장합니다. 이를 통해 어떤 인스턴스에 클라이언트가 연결되어 있든 상관없이 모든 메시지를 정상적으로 수신할 수 있습니다.

4. 데이터 흐름 (End-to-End Flow)

  1. 이벤트 발생: ComfyUI에서 작업 상태 변경(시작, 진행, 완료 등) 이벤트가 발생합니다.
  2. 이벤트 수집: AiNodeAgent가 로컬 WebSocket을 통해 이벤트를 수신합니다.
  3. 메타데이터 추가: AiNodeAgent는 이벤트에 nodeId, env, userId를 추가합니다.
  4. Relay 전달: AiNodeAgent가 이벤트를 relay의 WS/HTTP ingress로 보냅니다.
  5. Redis 반영: relay가 정규화된 이벤트를 Redis 키/채널에 반영합니다.
  6. Redis 구독: AiRealtimeGateway가 Redis 채널로부터 이벤트를 수신합니다.
  7. Room 기반 전달: AiRealtimeGateway는 이벤트의 userId를 확인하고, 해당하는 Socket.IO room으로 메시지를 전달합니다.
  8. 실시간 업데이트: 프론트엔드 클라이언트는 Socket.IO를 통해 메시지를 수신하고, UI를 실시간으로 업데이트합니다.