실시간 파이프라인 아키텍처 (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)
- Publisher:
2. 아키텍처 다이어그램
3. 구성 요소 (Components)
3.1. AiNodeAgent (Publisher)
- 위치:
ai-servergeneration runtime - 역할: ComfyUI 실행 이벤트와 job 상태를 수집하여 public relay로 전달하는 에이전트입니다.
- 주요 기능:
- ComfyUI WebSocket 연결: 로컬 ComfyUI의 WebSocket(
ws://comfyui:8188/ws)에 연결하여 모든 실시간 이벤트를 수신합니다. - 이벤트 전달: 수신한 ComfyUI 이벤트를 relay의 WebSocket/HTTP ingress로 전달합니다.
- 메타데이터 추가: 전달 전, 이벤트 페이로드에
nodeId,env,userId등 메타데이터를 추가하여 어떤 서버의 어떤 작업인지 식별할 수 있도록 합니다.userId는 SQS 메시지 처리 시prompt_id와userId를 Redis에 매핑해 둔 값을 참조하여 가져옵니다.
- 자동 재연결: relay 연결이 끊어지면 자동 재연결을 시도하여 안정성을 보장합니다.
- ComfyUI WebSocket 연결: 로컬 ComfyUI의 WebSocket(
3.2. Agent Relay
- 위치: AWS public ingress 뒤 relay 서비스
- 역할: 외부
AiNodeAgent와 shared Redis 사이를 중개하는 정규화 계층입니다. - 주요 기능:
- Agent 전용 ingress 제공: 외부 노드는 relay 주소 하나만 알면 됩니다.
- Transport 분리: heartbeat/progress/node.monitor는 WebSocket, 완료/실패/결과 등록은 HTTP로 분리할 수 있습니다.
- 메시지 정규화: 수신 이벤트를 Redis 키/채널 모델에 맞게 변환합니다.
- 보안/제어: 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 게이트웨이입니다.
- 주요 기능:
- Redis 구독:
ai-comfyui-events채널을 구독하여AiNodeAgent가 발행한 모든 이벤트를 수신합니다. - 이벤트 필터링 및 전달: 수신한 이벤트의
userId유무에 따라 적절한 Socket.IOroom으로 메시지를 전달합니다.adminroom:userId가 없는 전역 상태 이벤트(서버 ONLINE/OFFLINE 등)는adminroom에 있는 모든 클라이언트에게 전달됩니다.user-{userId}room:userId가 있는 특정 사용자의 작업 진행 이벤트는 해당user-{userId}room에 있는 클라이언트에게만 전달됩니다.
- 클라이언트 관리: 프론트엔드 클라이언트의 WebSocket 연결/해제 및
room참여/이탈을 관리합니다.
- Redis 구독:
3.5. RedisIoAdapter
- 위치: 백엔드 서버 (
backend-vivid-ai) - 역할: NestJS의 기본 WebSocket 어댑터를 대체하여 Redis 기반의 어댑터를 사용하게 합니다.
- 주요 기능:
- 백엔드 서버가 여러 인스턴스(예: EKS의 여러 Pod)로 수평 확장되더라도, Redis를 통해 Socket.IO 메시지를 모든 인스턴스에 걸쳐 브로드캐스팅할 수 있도록 보장합니다. 이를 통해 어떤 인스턴스에 클라이언트가 연결되어 있든 상관없이 모든 메시지를 정상적으로 수신할 수 있습니다.
4. 데이터 흐름 (End-to-End Flow)
- 이벤트 발생: ComfyUI에서 작업 상태 변경(시작, 진행, 완료 등) 이벤트가 발생합니다.
- 이벤트 수집:
AiNodeAgent가 로컬 WebSocket을 통해 이벤트를 수신합니다. - 메타데이터 추가:
AiNodeAgent는 이벤트에nodeId,env,userId를 추가합니다. - Relay 전달:
AiNodeAgent가 이벤트를 relay의 WS/HTTP ingress로 보냅니다. - Redis 반영: relay가 정규화된 이벤트를 Redis 키/채널에 반영합니다.
- Redis 구독:
AiRealtimeGateway가 Redis 채널로부터 이벤트를 수신합니다. - Room 기반 전달:
AiRealtimeGateway는 이벤트의userId를 확인하고, 해당하는 Socket.IOroom으로 메시지를 전달합니다. - 실시간 업데이트: 프론트엔드 클라이언트는
Socket.IO를 통해 메시지를 수신하고, UI를 실시간으로 업데이트합니다.