OurNight; engineering — POS

0x10. requirement

QR 주문 → 입금 확인 → 조리 → 서빙 → 정산 팝업 주점의 주문·입금 확인·조리·서빙 정보를 실시간으로 연결해 주문 누락과 현장 소통 비용을 줄인다.

0x11. 홀 POS

1. 주문·입금 관리
  • 신규 주문과 입금 확인 요청을 대기열에 표시한다.
  • 테이블 번호, 주문 내역, 주문 차수, 총금액을 보여준다.
  • 스태프가 입금을 확인하고 주문을 확정한다.
  • 확정된 주문을 주방으로 전달한다.
  • 홀에서 주문을 직접 추가하거나 취소·수량 변경할 수 있다.
2. 서빙 관리
  • 조리 완료된 메뉴를 서빙 대기열에 표시한다.
  • 테이블 번호, 메뉴명, 수량을 보여준다.
  • 서빙 후 메뉴 단위로 완료 처리한다.
  • 모든 메뉴가 서빙되면 주문을 완료 처리한다.
3. 호출 관리
  • 물, 티슈, 수저, 컵 등의 고객 요청을 표시한다.
  • 요청을 처리 중 또는 완료 상태로 변경한다.
  • 여러 스태프가 사용할 경우 담당자를 지정해 중복 응대를 막는다.
4. 테이블 관리
  • 테이블별 사용 여부와 진행 중 주문을 표시한다.
  • 손님 퇴장 후 테이블을 정리하고 새 세션으로 초기화한다.
  • 손님이 자리를 옮기면 주문을 새 테이블로 이전한다.
  • 합석 시 두 테이블의 주문과 세션을 하나로 병합한다.
  • 테이블 이용시간과 정리 필요 여부를 확인한다.

0x12. 주방 POS

1. 주문서 관리
  • 입금이 확인된 주문만 주방에 표시한다.
  • 테이블 번호, 메뉴명, 수량, 옵션, 주문 경과시간을 제공한다.
  • 주문을 들어온 순서대로 표시한다.
  • 메뉴 단위로 다음 상태를 관리한다.
조리 대기 → 조리 중 → 조리 완료 → 서빙 완료
  • 조리 완료 시 홀 POS에 알림을 보낸다.
  • 진행 중 주문과 완료 주문을 구분한다.
2. 메뉴·재고 관리
  • 메뉴를 즉시 품절 처리할 수 있다.
  • 품절 상태를 고객 메뉴판에 실시간 반영한다.
  • 재고 추가 또는 품절 해제가 가능해야 한다.

0x13. 도어 POS

  • 남은 테이블 수와 현재 혼잡도를 표시한다.
  • 대기 발생 여부를 설정한다.
  • 입장권을 QR 또는 인증 코드로 검증한다.
  • 유효한 입장권만 승인한다.

0x14. Admin

  • 실시간 매출과 순매출을 조회한다.
  • 주문 수와 메뉴별 판매량을 조회한다.
  • 진행 중·완료된 주문 내역을 조회한다.
  • 취소·환불 요청과 처리 이력을 관리한다.

5. 공통 요구사항

  • 홀·주방·도어 화면의 상태를 실시간 동기화한다.
  • 여러 스태프의 중복 처리를 방지한다.
  • 모든 상태 변경에 처리자와 처리 시각을 기록한다.
  • 네트워크 재연결 후 누락된 변경 사항을 복구한다.
  • 중복 클릭이나 재요청으로 주문이 여러 번 처리되지 않게 한다.

0x20. DESIGN

0x21. happy path

flowchart TD A["REST command로 server state 확정"] B["같은 transaction 안에서 outbox 기록"] C["outbox relay가 Redis Stream에 publish"] D["매장별 SSE reader가 stream event를 읽음"] E["SseEmitter가 event id와 payload를 브라우저에 전달"] F["front가 Redux에 patch하고 lastAppliedEventId 저장"] G{"다음 event가 있는가?"} A --> B     B --> C     C --> D     D --> E     E --> F     F --> G     G -->|"Yes"| D     G -->|"No · 새 event 대기"| H["Redis Stream에서 새 event 대기"]     H -->|"새 event publish"| D

0x22. first connection or server restart

상황
  • POS 처음 엶
  • server restart
문제
  • Redis reader는 현재 stream tail 근처에서 시작한다.
  • SSE connection 이전에 발생한 event는 reader가 놓칠 수 있다.
  • SSE event만으로 현재 POS state를 만들 수 없다.
design
SSE 연결
→ full baseline snapshot load
→ snapshot watermark 모음
→ 가장 작은 watermark까지 covered 처리
→ 이후 SSE event patch

0x23. temporary network loss

상황
  • Wi-Fi가 잠깐 끊겼다.
  • browser가 sleep 상태에 들어갔다.
  • SSE connection만 잠시 사라졌다.
  • Redis Stream에는 누락 event가 아직 남아 있다.
design
감지
EventSource onerror
또는
45초 watchdog timeout
또는
online / focus / visibilitychange
recovery
기존 EventSource cleanup
→ backoff + jitter
→ 새 ticket 발급
→ lastAppliedEventId 포함해 reconnect
→ Redis Stream에서 (lastEventId, +) replay
→ dedupe
→ Redux patch

0x24. incremental recovery invalidated

0x23를 더 이상 믿을 수 없는 상태

상황
상황incremental recovery가 깨지는 이유
장기 유실빠진 event 범위가 너무 크거나 알 수 없다.
stream trim필요한 event가 이미 없어졌다.
replay 과다기술적으로 가능해도 replay 비용과 queue risk가 너무 크다.
unknown eventevent를 어떻게 state에 반영해야 할지 해석할 수 없다.
payload parse 실패event meaning 자체를 만들 수 없다.
queue overflow중간 event가 모두 유지됐다는 보장이 없다.
문제

Level 1: event replay Level 2: delta calculation + partial reload Level 3: full baseline reload

1단계: event replay

다음 조건이면 replay를 믿을 수 없다.

  • stream trim으로 gap expired
  • replay가 300개 초과
  • invalid stream ID
  • unknown event type
  • invalid payload
  • replay exception
2단계: delta calculation + partial reload
sinceEventId 전달
→ Redis Stream scan
→ changed order/table/menu/call ID 계산
3단계: full baseline reload

다음 상황이면 delta recovery도 포기

  • /changes event가 500개 초과
  • changed entity가 종류별 100개 초과
  • replay gap expired
  • payload/type를 믿을 수 없음
  • sinceEventId가 없음
  • stream 내부에 RELOAD_REQUIRED 존재
design flow
flowchart TD A[Incremental recovery 불신 상태 감지] A --> B{감지 위치} B -->|Server| C[RELOAD_REQUIRED control event 전송] B -->|Client| D[parse 실패<br/>unknown mapping<br/>watchdog / queue 불확실성 감지] C --> E[Recovery 시작] D --> E E --> F{유효한 lastAppliedEventId가 있는가?} F -->|없음 또는 0-0| M[Full baseline reload] F -->|있음| G["GET /notifications/changes<br/>since=lastAppliedEventId"] G --> H{changes API가 incremental recovery를<br/>안전하게 계산할 수 있는가?} H -->|아니오| M H -->|예| I{변경 entity 수가<br/>partial limit 안쪽인가?} I -->|아니오| M I -->|예| J[변경된 order / table / menu / call만 조회] J --> K[Partial snapshot을 Redux에 반영] K --> L["changes.currentWatermark까지<br/>lastAppliedEventId 전진"] M --> N[주요 POS snapshot 병렬 조회] N --> O[각 snapshot의<br/>X-Pos-Stream-Watermark 수집] O --> P[필수 scope의 최소 watermark 계산] P --> Q[Baseline state 교체] Q --> R[최소 watermark까지<br/>lastAppliedEventId 전진] L --> S[Incremental event apply 재개] R --> S

0x25. high traffic

design flow
flowchart LR A["Event traffic 증가"] --> B["Bounded queue로 짧은 burst buffer"] B --> C["Worker shard로 전송 작업 분산"] C --> D{"처리 한계 안쪽인가?"} D -->|Yes| E["SSE event 전달 계속"] D -->|No| F["Event 기반 incremental apply 포기"] F --> G["RELOAD_REQUIRED"] G --> H["Partial reload 또는 full baseline"] H --> I["현재 server state와 다시 sync"]

0x26. slow consumer (single device failure)

design flow
flowchart LR A["특정 client의 event 소비 지연"] --> B["Client별 bounded queue가 지연을 격리"] B --> C{"지연이 queue limit 안쪽인가?"} C -->|Yes| D["SSE event 전달 계속"] C -->|No| E["Slow consumer 연결 종료"] E --> F["FE가 backoff와 jitter로 재연결"] F --> G["Partial reload 또는 full baseline"] G --> H["현재 server state와 다시 sync"]

0x27. redis publish failure

상황
  • 주문 DB transaction은 성공했다.
  • outbox row도 저장됐다.
  • Redis가 잠시 unavailable이다.
design
REST command transaction commit
→ outbox row remains unpublished
→ relay publish 실패
→ attempts 증가
→ nextAttemptAt 설정
→ 최대 60초 backoff
→ 이후 retry

0x2F. WHY SSE?



0x20. IMPL — core flow (happy)

대표 사례로 “결제 확인 및 주문 확정” 케이스를 설명하겠다.

0x21. server state 확정

POST /api/pos/{storeId}/hall/orders/{orderId}/confirm
  → HallOrderController.confirmOrder()
  → OrderService.confirmOrder()
  → OrderCommandService.confirmOrder()

0x22. outbox pattern

// PosNotificationPublisher
@Component
class PosNotificationPublisher(... 생략 ...) {
	@TransactionalEventListener(phase = TransactionPhase.BEFORE_COMMIT)
	fun handleCookingRequest(event: CookingRequestedEvent) {
	    if (!isCurrentSlotOrderEvent(event.storeId, event.orderId)) return
	    enqueueOutbox(event.storeId, "KITCHEN_COOK", event)
	}
 
	private fun enqueueOutbox(storeId: UUID, type: String, payload: Any) {
	    outboxRepository.save(
	        PosEventOutbox(
	            storeId = storeId,
	            eventType = type,
	            payload = objectMapper.writeValueAsString(payload),
	        )
	    )
	}
}
  • TransactionPhase.BEFORE_COMMIT
  • domain state change + outbox insert -> 같은 Transaction에 참여

0x23. outbox relay가 Redis Stream에 publish

@Component
class PosOutboxRelay(...) {
	@Scheduled(fixedDelay = 500)
    @Transactional
    fun publishPendingEvents() {
        val events = outboxRepository.lockPublishable(batchSize)
        if (events.isEmpty()) return
 
        events.forEach { event ->
            try {
                val streamId = publishToRedis(event)
                event.markPublished(streamId)
            } catch (e: Exception) {
                event.markFailed(e)
                // log ...
            }
        }
    }
}
 
interface PosEventOutboxRepository : JpaRepository<PosEventOutbox, UUID> {
    @Query(
        value = """
            select *
            from pos_event_outbox
            where published_at is null
              and next_attempt_at <= now()
            order by created_at
            limit :limit
            for update skip locked
        """,
        nativeQuery = true,
    )
    fun lockPublishable(@Param("limit") limit: Int): List<PosEventOutbox>
}

0x24. 매장별 SSE reader가 stream event를 읽음

outline
GET /api/pos/{storeId}/notifications/subscribe
  → SseNotificationService.subscribe()
  → registerEmitter()
  → startRedisReader()
  → readStreamLoop()
code
    fun subscribe(
        expectedStoreId: UUID,
        ticket: String,
        connectionId: String?,
        lastEventId: String?
    ): SseEmitter {
        val storeIdStr = redisTemplate.opsForValue().getAndDelete("ticket:$ticket")
            ?: throw BusinessException(ErrorCode.INVALID_TICKET)
 
        if (storeIdStr != expectedStoreId.toString()) {
            throw BusinessException(ErrorCode.INVALID_TICKET)
        }
 
        val emitter = SseEmitter(60 * 60 * 1000L)
        val chanKey = channelKey(storeIdStr)
        val sKey = streamKey(storeIdStr)
        val safeConnectionId = connectionId?.takeIf { it.isNotBlank() }?.take(128) ?: UUID.randomUUID().toString()
 
        registerEmitter(chanKey, emitter, sKey, safeConnectionId)
        markActive(chanKey)
        logger.info("[SSE-CONNECT] storeId={} connectionId={} lastEventId={}", storeIdStr, safeConnectionId, lastEventId)
 
        enqueueOutbound(chanKey, emitter, SseMsg("connect", "connected", "0-0"))
 
        if (!lastEventId.isNullOrBlank()) {
            replayMissingEvents(sKey, lastEventId, emitter)
            markActive(chanKey)
        }
 
        return emitter
    }
 
    private fun registerEmitter(chanKey: String, emitter: SseEmitter, streamKey: String, connectionId: String) {
        connectionIds[emitter] = connectionId
        outboundQueueOf(emitter)
        startEmitterSender(chanKey, emitter, connectionId)
        emitter.onCompletion { unregisterEmitter(chanKey, emitter, "completion") }
        emitter.onTimeout { emitter.complete(); unregisterEmitter(chanKey, emitter, "timeout") }
        emitter.onError {
            logger.warn("[SSE-ERROR] storeId={} connectionId={}", chanKey, connectionId)
            unregisterEmitter(chanKey, emitter, "error")
        }
 
        emitters.compute(chanKey) { _, list ->
            val newList = list ?: CopyOnWriteArrayList()
            newList.add(emitter)
            newList
        }
 
        startRedisReader(chanKey, streamKey)
    }
 
    private fun startRedisReader(chanKey: String, streamKey: String) {
        streamReaders.compute(chanKey) { _, existing ->
            if (existing != null) return@compute existing
 
            try {
                val startId = currentStreamTailId(streamKey)
                val task = streamReaderExec.submit { readStreamLoop(chanKey, streamKey, startId) }
                logger.info("Redis Stream Reader Started: chanKey={} streamKey={} startId={}", chanKey, streamKey, startId)
                task
            } catch (e: Exception) {
                logger.error("Failed to start redis stream reader", e)
                null
            }
        }
    }
 
    private fun readStreamLoop(chanKey: String, streamKey: String, initialLastSeenId: String) {
        var lastSeenId = initialLastSeenId
        val options = StreamReadOptions.empty()
            .count(100)
            .block(Duration.ofSeconds(2))
 
        while (!Thread.currentThread().isInterrupted && emitters.containsKey(chanKey)) {
            try {
                val records = redisTemplate.opsForStream<String, String>().read(
                    options,
                    StreamOffset.create(streamKey, ReadOffset.from(lastSeenId))
                ) ?: emptyList()
 
                for (record in records) {
                    lastSeenId = record.id.toString()
                    enqueueRecord(chanKey, streamKey, record)
                }
            } catch (e: InterruptedException) {
                Thread.currentThread().interrupt()
            } catch (e: Exception) {
                logger.warn("Redis stream read failed: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
                enqueue(chanKey, SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
                Thread.sleep(1000)
            }
        }
 
        logger.info("Redis Stream Reader Stopped: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
    }

0x25. SseEmitter가 event id와 payload를 브라우저에 전달

outline
매장 queue
  → broadcastNow()
  → emitter별 outbound queue
  → startEmitterSender()
  → sendToEmitter()
  → emitter.send()
code
    private fun broadcastNow(chanKey: String, msg: SseMsg) {
        emitters[chanKey]?.forEach { emitter ->
            enqueueOutbound(chanKey, emitter, msg)
        }
    }
    
    private fun startEmitterSender(chanKey: String, emitter: SseEmitter, connectionId: String) {
        val task = emitterSenderExec.submit {
            while (!Thread.currentThread().isInterrupted && outboundQueues.containsKey(emitter)) {
                val msg = try {
                    outboundQueues[emitter]?.poll(30, TimeUnit.SECONDS)
                } catch (e: InterruptedException) {
                    Thread.currentThread().interrupt()
                    null
                } ?: continue
 
                if (!sendToEmitter(emitter, msg.name, msg.data, msg.id)) {
                    unregisterEmitter(chanKey, emitter, "send-failed")
                    break
                }
            }
            logger.info("[SSE-SENDER-STOP] storeId={} connectionId={}", chanKey, connectionId)
        }
        emitterSenderTasks[emitter] = task
    }
 
    private fun sendToEmitter(emitter: SseEmitter, name: String, data: String, id: String): Boolean {
        try {
            emitter.send(
                SseEmitter.event()
                    .id(id)
                    .name(name)
                    .data(data)
                    .reconnectTime(2000)
            )
            return true
        } catch (e: Exception) {
            logger.debug("Failed to send SSE event: ${e.message}")
            emitter.completeWithError(e)
            return false
        }
    }

0x26. FE: Redux에 patch하고 lastAppliedEventId 저장

0. outline
PosRoot
  → PosRealtimeBridge
  → ticket 발급
  → EventSource 생성
  → MessageEvent 수신
  → eventAdapter
  → Redux realtime action dispatch
  → listener middleware에서 entity patch
  → markStreamEventApplied
  → lastAppliedEventId 갱신
1. bridge mount
// PosRoot.tsx
export default function PosRoot() {
    const { storeId = "" } = useParams();
 
    return (
        <StoreProvider>
            <PosRealtimeBridge storeId={storeId} />
            <div className="contents font-sans">
                <Outlet />
            </div>
        </StoreProvider>
    );
}
2. EventSource 생성과 event id 수신
// sseClient.ts
export function subscribeSse(params: {
    url: string;
    withCredentials?: boolean;
    handlers: Record<string, (raw: string, meta: SseEventMeta) => void>; // eventName -> handler(rawJson)
    onError?: (ev?: Event) => void;
    onOpen?: () => void;
}): SseSubscription {
    const es = new EventSource(params.url, { withCredentials: params.withCredentials ?? true });
 
    for (const [eventName, handler] of Object.entries(params.handlers)) {
        es.addEventListener(eventName, (ev) => {
            const message = ev as MessageEvent;
            handler(message.data, { lastEventId: message.lastEventId });
        });
    }
 
    es.onopen = () => params.onOpen?.();
    es.onerror = (ev) => params.onError?.(ev);
 
    return {
        close: () => es.close(),
    };
}
3. lastAppliedEventId를 reconnect URL에 포함
// PosRealtimeBridge.tsx
export function PosRealtimeBridge({ storeId }: Props) {
    const dispatch = useAppDispatch();
    const lastAppliedEventId = useAppSelector((state) => state.pos.stream.lastAppliedEventId);
    // ....
// PosRealtimeBridge.tsx
export function PosRealtimeBridge({ storeId }: Props) {
    const dispatch = useAppDispatch();
    const lastAppliedEventId = useAppSelector((state) => state.pos.stream.lastAppliedEventId);
    // ....
	// ...
	useEffect(() => {
	//. ...
		const reconcileAfterResume = () => {
		//. ...
			const connectionId = createRequestId();
			const url = getSseUrl(storeId, ticket, connectionId, lastAppliedEventIdRef.current);
4. SSE 이벤트를 Redux action으로 변환
// listener.ts
export const dispatchPosRealtimeEventWithId = (
    event: PosRealtimeEvent,
    streamEventId: string,
    dedupeEventId?: string,
): PosRealtimeAction & { meta: { streamEventId: string; dedupeEventId?: string } } =>
    ({
        type: `pos/realtime/${event.type}`,
        payload: event.payload,
        meta: { streamEventId, dedupeEventId },
    }) as PosRealtimeAction & { meta: { streamEventId: string; dedupeEventId?: string } };
5. Redux entity patch
// listener.ts
posListenerMiddleware.startListening({
	effect: async (action, api) => {
	// ...
		switch (action.type) {
		// ...
            case "pos/realtime/payment/resolved": {
                const payload = action.payload;
                const pos = getPosState(api);
                const orderExists = !!pos.entities.ordersById[payload.orderId];
 
                if (orderExists) {
                    api.dispatch(
                        posActions.patchOrderStatus({
                            orderId: payload.orderId,
                            status: toPosOrderStatus(payload.currentOrderStatus),
                        }),
                    );
                } else {
                    invalidatePayments(api);
                }
 
                api.dispatch(posActions.removePaymentQueue(payload.orderId));
                invalidateHall(api);
                invalidateKitchen(api);
                invalidateTables(api);
                break;
            }
6. lastAppliedEventId 저장
//listener.ts
posListenerMiddleware.startListening({
	// .....
	effect: async (action, api) => {
	//......
        if (isReplayableStreamEvent) {
            api.dispatch(
                posActions.markStreamEventApplied({
                    eventId: streamEventId,
                    dedupeId,
                    snapshotScopes: getSnapshotScopesForRealtimeAction(action.type),
                }),
            );
        }
// slice.ts 
export const posSlice = createSlice({
    name: "pos",
    initialState: initialPosDomainState,
    reducers: {
        markStreamEventApplied(
            state,
            action: PayloadAction<{ eventId: string; dedupeId?: string; snapshotScopes: StreamSnapshotScope[] }>
        ) {
            const eventId = normalizeRedisStreamId(action.payload.eventId);
            rememberAppliedEventId(state, eventId);
            rememberAppliedDedupeId(state, action.payload.dedupeId);
 
            if (isRedisStreamIdAfter(eventId, state.stream.lastAppliedEventId)) {
                state.stream.lastAppliedEventId = eventId;
            }
 
            for (const scope of action.payload.snapshotScopes) {
                const prev = state.stream.lastAppliedEventIdBySnapshot[scope];
                if (isRedisStreamIdAfter(eventId, prev)) {
                    state.stream.lastAppliedEventIdBySnapshot[scope] = eventId;
                }
            }
        },

0x27. 다음 이벤트 또는 Redis Stream에서 대기

XREAD 결과 있음
  → 각 record 처리
  → lastSeenId 갱신
  → 다시 XREAD
 
XREAD 결과 없음
  → 최대 2초 blocking
  → 결과가 없으면 빈 목록
  → while loop를 돌며 다시 XREAD
 
새 event publish
  → blocking read 반환
  → record 처리
// SseNotificationService.kt 
    private fun readStreamLoop(chanKey: String, streamKey: String, initialLastSeenId: String) {
        var lastSeenId = initialLastSeenId
        val options = StreamReadOptions.empty()
            .count(100)
            .block(Duration.ofSeconds(2))
 
        while (!Thread.currentThread().isInterrupted && emitters.containsKey(chanKey)) {
            try {
                val records = redisTemplate.opsForStream<String, String>().read(
                    options,
                    StreamOffset.create(streamKey, ReadOffset.from(lastSeenId))
                ) ?: emptyList()
 
                for (record in records) {
                    lastSeenId = record.id.toString()
                    enqueueRecord(chanKey, streamKey, record)
                }
            } catch (e: InterruptedException) {
                Thread.currentThread().interrupt()
            } catch (e: Exception) {
                logger.warn("Redis stream read failed: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
                enqueue(chanKey, SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
                Thread.sleep(1000)
            }
        }
 
        logger.info("Redis Stream Reader Stopped: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
    }
 

0x2F. 정리

flowchart TD Controller["HallOrderController.confirmOrder"] Command["OrderCommandService.confirmOrder<br/>@Transactional"] Event["ApplicationEventPublisher<br/>HallPaymentResolvedEvent"] Publisher["PosNotificationPublisher<br/>BEFORE_COMMIT"] Outbox[("pos_event_outbox")] Relay["PosOutboxRelay<br/>500ms scheduled"] Stream[("notification:store:{storeId}:pos")] Reader["SseNotificationService<br/>매장별 readStreamLoop"] Emitter["SseEmitter.send<br/>id=Redis Stream ID"] EventSource["sseClient EventSource<br/>MessageEvent.lastEventId"] Bridge["PosRealtimeBridge<br/>eventAdapter + dispatch"] Listener["Redux listener middleware<br/>dedupe + patch"] Slice["posSlice<br/>lastAppliedEventId"] Controller --> Command Command --> Event Event --> Publisher Publisher --> Outbox Outbox --> Relay Relay --> Stream Stream --> Reader Reader --> Emitter Emitter --> EventSource EventSource --> Bridge Bridge --> Listener Listener --> Slice Slice -. "재연결 시 lastEventId" .-> Reader


0x30. IMPL — first connection or server restart

outline
1. SSE 연결
2. 연결 성공 시 baseline 복구 실행
3. 최초 연결(`lastAppliedEventId = 0-0`)이면 full baseline 실행
4. 각 snapshot 응답의 watermark를 scope별로 저장
5. 모든 snapshot 중 가장 작은 watermark까지 stream cursor를 전진
6. SSE 이벤트를 Redux 상태에 patch

0x31. Redis reader가 현재 stream tail에서 시작

//SseNotificationService.kt
val startId = currentStreamTailId(streamKey)
val task = streamReaderExec.submit {
    readStreamLoop(chanKey, streamKey, startId)
}

tail 조회

//SseNotificationService.kt
private fun currentStreamTailId(streamKey: String): String =
    try {
        val info = redisTemplate.opsForStream<String, String>().info(streamKey)
        info.lastGeneratedId() ?: info.lastEntryId() ?: "0-0"
    } catch (_: Exception) {
        "0-0"
    }

startId 이후 이벤트만 읽는다.

// SseNotificationService.kt -> readStreamLoop
StreamOffset.create(streamKey, ReadOffset.from(lastSeenId))

0x32. SSE 연결 성공 후 baseline 복구

// PosRealtimeBridge.tsx
const sub = subscribeSse({
    url,
    handlers,
    onOpen: () => {
        touchSseActivity();
        reconnectDelayMsRef.current = INITIAL_RECONNECT_DELAY_MS;
        startWatchdog();
        reloadBaseline();
    },
});

onOpen에서 reloadBaseline()을 호출하므로 순서는

EventSource 연결 성공
→ baseline 복구 시작

0x33. 최초 POS 연결이면 full baseline으로 전환

초기 stream cursor는 [initialState.ts (line 36)]에서 0-0

initialState.ts
stream: {
    lastAppliedEventId: "0-0",
}

연결 후 실제 호출되는 것은 reloadChangedPosBaseline이지만, cursor가 0-0이면 full baseline으로 전환

//reloadChangedPosBaseline.ts
if (!sinceEventId || sinceEventId === "0-0") {
    await dispatch(reloadPosBaseline(storeId)).unwrap();
    return { storeId, mode: "full" as const };
}

즉 최초 POS 오픈 흐름은

SSE open
→ reloadChangedPosBaseline(sinceEventId = "0-0")
→ reloadPosBaseline()
→ full snapshot load

재연결이고 유효한 cursor가 있으면 전체 baseline 대신 변경분 복구를 먼저 시도

0x34. full baseline snapshot 병렬 조회

//reloadPosBaseline.ts
await Promise.all([
    dispatch(loadPaymentQueue({ storeId, force: true })).unwrap(),
    dispatch(loadServingQueue({ storeId, force: true })).unwrap(),
    dispatch(loadCallQueue({ storeId, force: true })).unwrap(),
    dispatch(loadHallTableZones(storeId)).unwrap(),
    dispatch(loadHallTables({ storeId, force: true })).unwrap(),
    dispatch(loadMenuCatalog(storeId)).unwrap(),
    dispatch(loadHallOrders(...)).unwrap(),
    dispatch(loadKitchenOrders(...)).unwrap(),
]);

이 부분이 full baseline snapshot load에 해당

0x35. 서버가 snapshot watermark를 응답 헤더에 추가

// PosStreamWatermarkService.kt
fun currentWatermark(storeId: UUID): String {
    val key = PosStreamKey.store(storeId)
    val info = redisTemplate.opsForStream<String, String>().info(key)
    return info.lastGeneratedId() ?: info.lastEntryId() ?: "0-0"
}

응답에는 X-Pos-Stream-Watermark 헤더로 포함

ResponseEntity.ok()
    .header(POS_STREAM_WATERMARK_HEADER, watermark)
    .body(body)

프론트는 [streamSnapshot.ts (line 11)]에서 이 헤더를 읽음

0x36. snapshot별 watermark 저장

예를 들어 결제 snapshot은 loadPaymentQueue.ts (line 45)에서 저장

dispatch(posActions.markStreamSnapshotLoaded({
    eventId: snapshot.streamWatermark,
    snapshotScopes: ["payments"],
}));

Reducer는

// slice.ts
for (const scope of action.payload.snapshotScopes) {
    const prev = state.stream.lastAppliedEventIdBySnapshot[scope];
 
    if (isRedisStreamIdAfter(eventId, prev)) {
        state.stream.lastAppliedEventIdBySnapshot[scope] = eventId;
    }
}

payments, servings, calls, hallOrders, kitchenOrders, tables별 watermark가 각각 관리됩니다.

0x37. 가장 작은 watermark까지 covered 처리

//reloadPosBaseline.ts
const coveredWatermark = BASELINE_SNAPSHOT_SCOPES
    .map((scope) => stream.lastAppliedEventIdBySnapshot[scope])
    .reduce((min, eventId) =>
        compareRedisStreamIds(eventId, min) < 0 ? eventId : min
    );
 
dispatch(posActions.markStreamWatermarkCovered({
    eventId: coveredWatermark,
}));

Reducer는

// reducer.ts
if (isRedisStreamIdAfter(eventId, state.stream.lastAppliedEventId)) {
    state.stream.lastAppliedEventId = eventId;
}

가장 작은 값을 사용하는 이유는 일부 snapshot만 최신일 때 stream cursor를 너무 앞으로 보내서 이벤트를 건너뛰지 않도록 하기 위해서

0x38. 이후 SSE 이벤트 patch

SSE 이벤트를 Redux action으로 변환하고 dispatch하는 부분

PosRealtimeBridge.tsx
const event = adaptRealtimeEvent(eventName, parsed);
const reduxAction = dispatchPosRealtimeEventWithId(
    event,
    meta.lastEventId,
    dedupeEventId,
);
 
dispatch(reduxAction);
  • 실제 entity patch와 stream cursor 갱신은 listener.ts (posListenerMiddleware.startListening)에서 처리

  • 이벤트 적용이 끝나면 해당 Redis Stream ID까지 적용됐다고 기록

0x3F. 정리

flowchart TD A["PosRoot에서 PosRealtimeBridge 마운트"] B["Redux lastAppliedEventId 조회"] C["ticket + lastEventId로 EventSource 연결"] D["서버 registerEmitter"] E["매장 reader를 현재 Stream tail에서 시작"] F["EventSource onOpen"] G["reloadChangedPosBaseline"] H{"증분 복구 가능한가?"} I["reloadPosBaseline"] J["6개 scope snapshot 병렬 조회"] K["서버가 X-Pos-Stream-Watermark 응답"] L["scope별 watermark 저장"] M["가장 작은 watermark 계산"] N["lastAppliedEventId covered 처리"] O["SSE event Redux patch"] P["markStreamEventApplied로 cursor 전진"] A --> B B --> C C --> D D --> E C --> F F --> G G --> H H -->|"No · 최초 연결/gap 만료"| I H -->|"Yes"| O I --> J J --> K K --> L L --> M M --> N E --> O N --> O O --> P P --> O


0x40. IMPL — temporary network loss

  • 학교 와이파이는 터질 가능성이 늘 존재한다
  • 이에 다음 전략을 사용
단계전략역할
1Redis Stream cursorRedis Stream에서 읽은 도메인 이벤트의 SSE event id에 Redis Stream ID를 넣는다.
2Dedupe/idempotent apply프론트가 이벤트를 적용하면서 lastAppliedEventId를 갱신한다.
3Client watchdog/reconnect연결 오류, idle timeout, online/focus 복귀 시 재연결한다.
4Last-Event-ID replay재연결 URL에 lastEventId를 실어 누락 구간 replay를 요청한다.
5Snapshot reconciliationreplay가 불가능하면 변경분 reload 또는 full baseline으로 전환한다.

0x41. outline

flowchart TD %% Nodes Normal[Normal SSE Connection] --> CheckID[Check MessageEvent.lastEventId] CheckID --> Redux[Redux listener patches state & stores lastAppliedEventId] Redux --> Disconnect[Network Disconnection] Disconnect --> Schedule[Watchdog / onerror / online / focus schedules reconnect] Schedule --> Ticket[Issue new ticket] Ticket --> Subscribe[Include lastEventId in subscribe URL] Subscribe --> ServerReplay[Server reads Redis Stream range: lastEventId to +] %% Decision ServerReplay --> ReplayDecision{Is replay possible?} %% Branches ReplayDecision -- Yes --> Dispatch[Replay missing events & Dispatch] ReplayDecision -- No --> Reload[Return RELOAD_REQUIRED] Reload --> Fetch[Frontend executes /notifications/changes or Full Baseline]

0x42. 설명

이 구조에서 중요한 점은 네트워크 유실을 “예외”로 보지 않는다는 것이다. POS 단말은 와이파이, 태블릿 sleep, 브라우저 background 상태의 영향을 자주 받기 때문에, 유실 후 재동기화가 기본 경로로 들어가 있다.

서버 전송부는 SSE event id를 명시적으로 넣는다. Redis Stream에서 읽은 도메인 이벤트는 id가 Redis Stream ID이고, PINGconnect, 일부 즉시 reload 신호처럼 stream record가 아닌 제어 이벤트는 0-0을 쓴다.

private fun sendToEmitter(emitter: SseEmitter, name: String, data: String, id: String): Boolean {
    emitter.send(
        SseEmitter.event()
            .id(id)
            .name(name)
            .data(data)
            .reconnectTime(2000)
    )
    return true
}

프론트 EventSource 래퍼는 브라우저의 MessageEvent.lastEventId를 handler로 넘긴다.

es.addEventListener(eventName, (ev) => {
    const message = ev as MessageEvent;
    handler(message.data, { lastEventId: message.lastEventId });
});

프론트가 연결을 다시 만들 때는 마지막 적용 지점을 URL에 싣는다.

const url = getSseUrl(storeId, ticket, connectionId, lastAppliedEventIdRef.current);

그래도 복구가 불가능한 상황은 full baseline으로 떨어진다.

full baseline으로 전환되는 조건이유
Redis Stream이 trim되어 gap이 사라진 경우lastEventId 이후 이벤트를 완전히 재생할 수 없다.
replay 이벤트가 300개를 넘는 경우단일 연결 replay가 너무 커져 SSE 송신 큐를 압박할 수 있다.
/notifications/changes 결과가 500개를 넘는 경우변경분 계산 자체가 너무 커져 partial reload의 장점이 줄어든다.
payload parse 실패이벤트 의미를 신뢰할 수 없다.
unknown event type프론트/백엔드 이벤트 스키마가 어긋났을 가능성이 있다.
SSE queue overflow중간 이벤트를 모두 유지하지 못했으므로 cursor만 믿을 수 없다.
slow consumer해당 단말이 이벤트 소비 속도를 따라오지 못한다.
baseline partial reload 대상 entity가 너무 많은 경우개별 재조회보다 전체 snapshot이 더 안정적이다.

이 점에서 현재 구현은 “유실을 절대 막는다”가 아니라 “유실을 감지하면 빠르게 기준 상태를 다시 잡는다”는 설계다.



0x50. IMPL — incremental recovery invalidated

SSE replay 시도
→ replay 불신 조건이면 RELOAD_REQUIRED
→ 클라이언트 reloadChangedPosBaseline 실행
→ 유효한 lastAppliedEventId가 있으면 /changes 호출
→ delta 계산 가능 + entity 100개 이하이면 partial reload
→ 아니면 full baseline reload
→ watermark까지 cursor 전진
→ SSE event 적용 재개

1단계: Event replay

0x31. replay 시작

클라이언트가 lastEventId를 전달하면 서버가 누락 이벤트를 replay

// SseNotificationService.kt
if (!lastEventId.isNullOrBlank()) {
    replayMissingEvents(sKey, lastEventId, emitter)
    markActive(chanKey)
}

실제 replay 구현은 다음 코드

//SseNotificationService.kt
private fun replayMissingEvents(
    streamKey: String,
    lastEventId: String,
    emitter: SseEmitter,
) {
    try {
        if (isReplayGapExpired(streamKey, lastEventId)) {
            sendReloadSignal(emitter)
            return
        }
 
        val range = Range.rightOpen(lastEventId, "+")
        val events = redisTemplate.opsForStream<String, String>()
            .range(streamKey, range)
            ?: return
 
        var count = 0
        for (record in events) {
            if (count++ >= 300) {
                enqueueReplay(
                    emitter,
                    SseMsg("RELOAD_REQUIRED", "{}", "0-0"),
                )
                return
            }
 
            val data = record.value
            val type = data["type"] ?: "UNKNOWN"
 
            if (type !in knownEventNames) {
                enqueueReplay(
                    emitter,
                    SseMsg(
                        "RELOAD_REQUIRED",
                        "{}",
                        record.id.toString(),
                    ),
                )
                return
            }
 
            enqueueReplay(/* event */)
        }
    } catch (e: Exception) {
        sendReloadSignal(emitter)
    }
}
replay 불신 조건 대응
조건구현 위치
stream trim / gap expiredisReplayGapExpired()
invalid stream IDparseStreamId() 실패 시 gap expired로 처리
replay 300개 초과count++ >= 300
unknown event typetype !in knownEventNames
replay exceptioncatch에서 sendReloadSignal()
invalid payload서버 replay가 아닌 클라이언트 safeJson()에서 감지
Stream trim과 invalid ID 검사
// SseNotificationService.kt
private fun isReplayGapExpired(
    streamKey: String,
    lastEventId: String,
): Boolean {
    val lastId = parseStreamId(lastEventId) ?: return true
 
    val info = try {
        redisTemplate.opsForStream<String, String>().info(streamKey)
    } catch (_: Exception) {
        return false
    }
 
    val firstId = info.firstEntryId() ?: return false
 
    return compareStreamIds(
        parseStreamId(firstId) ?: return false,
        lastId,
    ) > 0
}
  • lastEventId 형식이 잘못되면 parseStreamId()null을 반환하므로 replay를 신뢰하지 ㄴㄴ
  • Redis Stream 첫 ID가 lastEventId보다 뒤라면 필요한 이벤트가 trim된 것으로 판단
Unknown event 감지

서버가 신뢰하는 event 목록은 다음 코드에 정의되어 있음.

// SseNotificationService.kt
private val knownEventNames = setOf(
    "RELOAD_REQUIRED",
    "HALL_PAYMENT",
    "HALL_PAYMENT_RESOLVED",
    // ...
    "KITCHEN_MENU_CHANGED",
)

stream reader와 replay 양쪽에서 이 목록에 없는 이벤트를 발견하면 RELOAD_REQUIRED로 전환

0x32. Payload parse 실패 감지

서버 replay는 payload JSON을 검증하지 않고 클라이언트로 전달. 클라이언트가 JSON 파싱 실패를 감지

//sseClient.ts
export function safeJson<T>(raw: string): T | null {
    try {
        return JSON.parse(raw) as T;
    } catch {
        return null;
    }
}
// PosRealtimeBridge.tsx 
const parsed = safeJson<unknown>(raw);
 
if (!parsed) {
    console.error("[SSE PARSE ERROR]", {
        storeId,
        eventName,
        raw,
    });
 
    reloadBaseline();
    return;
}

event name을 프론트 도메인 event로 매핑하지 못하는 경우도 recovery를 실행

const event = adaptRealtimeEvent(eventName, parsed);
 
if (!event) {
    reloadBaseline();
    return;
}

매핑 구현은 eventAdapter.ts (line 13)에 있습니다.


0x33. Queue overflow

서버에는 세 종류의 queue가 있다.

// 매장 broadcast queue
ArrayBlockingQueue<SseMsg>(2000)
 
// emitter별 replay queue
ArrayBlockingQueue<SseMsg>(500)
 
// emitter별 outbound queue
ArrayBlockingQueue<SseMsg>(500)

Broadcast queue overflow

if (q.offer(msg)) {
    markActive(chanKey)
    return
}
 
logger.warn("queue overflow: chanKey=$chanKey")
q.clear()
q.offer(SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
markActive(chanKey)

중간 이벤트를 모두 버리고 RELOAD_REQUIRED만 남김

Replay queue overflow

if (!q.offer(msg)) {
    q.clear()
    q.offer(SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
}

Outbound queue overflow

logger.warn(
    "[SSE-SLOW-CONSUMER] outbound queue overflow"
)
 
q.clear()
q.offer(SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
emitter.complete()
unregisterEmitter(chanKey, emitter, "slow-consumer")

이 경우에는 RELOAD_REQUIRED를 넣은 직후 연결을 finalize. 따라서 control event 전달보다 연결 종료 → 재연결 → baseline recovery가 실질적인 복구 경로


2단계: Delta calculation + partial reload

0x34. /notifications/changes 호출

프론트 API 코드

// posStreamChangesApi.ts 
 
const sp = new URLSearchParams();
sp.set("since", sinceEventId);
sp.set("limit", "500");
 
const res = await apiRequest(
    `/api/pos/${storeId}/notifications/changes?${sp.toString()}`,
);

백엔드 endpoint는

// PosNotificationController.kt 
@GetMapping("/changes")
fun getChangesSince(
    @RequestParam since: String?,
    @RequestParam(
        required = false,
        defaultValue = "500",
    ) limit: Int,
) = ResponseEntity.ok(
    posStreamChangeService.collectChanges(
        storeId = storeId,
        sinceEventId = since,
        maxEvents = limit.coerceIn(1, 1000),
    ),
)

현재 프론트가 전달하는 기준은 500개

0x35. Redis Stream scan

// PosStreamChangeService.kt
val records = redisTemplate.opsForStream<String, String>()
    .range(
        streamKey,
        Range.rightOpen(normalizedSince, "+"),
    )
    ?: emptyList()

이후 event payload에서 변경된 entity ID를 collect

val orderIds = linkedSetOf<UUID>()
val tableIds = linkedSetOf<UUID>()
val menuIds = linkedSetOf<UUID>()
val callRequestIds = linkedSetOf<UUID>()

대표적인 계산

when (type) {
    "HALL_PAYMENT",
    "HALL_SERVE",
    "KITCHEN_COOK" -> {
        payload.uuid("orderId")?.let(orderIds::add)
        payload.uuid("tableId")?.let(tableIds::add)
        payload.uuid("menuId")?.let(menuIds::add)
    }
 
    "CUSTOMER_CALL" -> {
        payload.uuid("requestId")?.let(callRequestIds::add)
        payload.uuid("tableId")?.let(tableIds::add)
    }
 
    "KITCHEN_MENU_CHANGED" -> {
        payload.uuid("menuId")?.let(menuIds::add)
    }
}

0x36. Delta 계산을 포기하는 조건

조건반환 reason
invalid sinceEventIdINVALID_SINCE_EVENT_ID
sinceEventId == 0-0NO_BASELINE_EVENT_ID
stream trimREPLAY_GAP_EXPIRED
event 500개 초과TOO_MANY_EVENTS
type field 없음MISSING_EVENT_TYPE
stream 안에 RELOAD_REQUIREDSTREAM_RELOAD_REQUIRED
payload parse 실패INVALID_EVENT_PAYLOAD
unknown eventUNKNOWN_EVENT_TYPE:{type}

예를 들어 event 개수 제한은 다음

if (records.size > maxEvents) {
    return reloadRequired(
        normalizedSince,
        currentWatermark,
        "TOO_MANY_EVENTS",
    )
}

payload 검증은 다음 코드

val payload = parsePayload(record.value["payload"])
    ?: return reloadRequired(
        normalizedSince,
        currentWatermark,
        "INVALID_EVENT_PAYLOAD",
    )

0x37. Partial reload 조건

프론트 복구 진입점은 다음

/changes가 delta 계산을 포기했다면 full baseline으로 전환

// reloadChangedPosBaseline.ts
if (changes.requiresReload) {
    await dispatch(reloadPosBaseline(storeId)).unwrap();
 
    return {
        storeId,
        mode: "full",
        reason: changes.reason,
    };
}

변경 entity가 종류별 100개를 초과해도 full baseline으로 전환

const PARTIAL_RELOAD_ENTITY_LIMIT = 100;
 
if (
    changes.changedOrderIds.length > PARTIAL_RELOAD_ENTITY_LIMIT ||
    changes.changedTableIds.length > PARTIAL_RELOAD_ENTITY_LIMIT ||
    changes.changedMenuIds.length > PARTIAL_RELOAD_ENTITY_LIMIT ||
    changes.changedCallRequestIds.length > PARTIAL_RELOAD_ENTITY_LIMIT
) {
    await dispatch(reloadPosBaseline(storeId)).unwrap();
 
    return {
        storeId,
        mode: "full",
        reason: "TOO_MANY_CHANGED_ENTITIES",
    };
}

0x38. 변경 entity만 다시 조회

// reloadChangedPosBaseline.ts
await Promise.all([
    orderIds.length > 0
        ? fetchHallOrdersByIds(storeId, orderIds)
        : Promise.resolve(),
 
    orderIds.length > 0
        ? fetchKitchenOrdersByIds(storeId, orderIds)
        : Promise.resolve(),
 
    orderIds.length > 0
        ? dispatch(loadPaymentQueue({ storeId, force: true })).unwrap()
        : Promise.resolve(),
 
    orderIds.length > 0
        ? dispatch(loadServingQueue({ storeId, force: true })).unwrap()
        : Promise.resolve(),
 
    changes.changedCallRequestIds.length > 0
        ? dispatch(loadCallQueue({ storeId, force: true })).unwrap()
        : Promise.resolve(),
 
    changes.changedMenuIds.length > 0
        ? dispatch(loadMenuCatalog(storeId)).unwrap()
        : Promise.resolve(),
 
    tableIds.length > 0
        ? reloadChangedTables(storeId, tableIds, dispatch, getState)
        : Promise.resolve(),
]);

Partial snapshot 반영 후 /changes가 반환한 현재 watermark까지 cursor를 전진

dispatch(
    posActions.markStreamWatermarkCovered({
        eventId: changes.currentWatermark,
    }),
);

3단계: Full baseline reload

0x39. Full baseline 전환 조건

//reloadChangedPosBaseline.ts
if (!sinceEventId || sinceEventId === "0-0") {
    await dispatch(reloadPosBaseline(storeId)).unwrap();
}

그 외에도 다음 조건에서 full baseline을 execute

  • /changes 요청 실패
  • changes.requiresReload
  • 변경 entity 종류별 100개 초과

0x3A. 주요 snapshot 병렬 조회

// reloadPosBaseline.ts
await Promise.all([
    dispatch(loadPaymentQueue({ storeId, force: true })).unwrap(),
    dispatch(loadServingQueue({ storeId, force: true })).unwrap(),
    dispatch(loadCallQueue({ storeId, force: true })).unwrap(),
    dispatch(loadHallTableZones(storeId)).unwrap(),
    dispatch(loadHallTables({ storeId, force: true })).unwrap(),
    dispatch(loadMenuCatalog(storeId)).unwrap(),
    dispatch(loadHallOrders(/* ... */)).unwrap(),
    dispatch(loadKitchenOrders(/* ... */)).unwrap(),
]);

0x3B. 최소 watermark 계산

const coveredWatermark = BASELINE_SNAPSHOT_SCOPES
    .map(
        (scope) =>
            stream.lastAppliedEventIdBySnapshot[scope],
    )
    .reduce(
        (min, eventId) =>
            compareRedisStreamIds(eventId, min) < 0
                ? eventId
                : min,
    );
 
dispatch(
    posActions.markStreamWatermarkCovered({
        eventId: coveredWatermark,
    }),
);

필수 scope는 다음 여섯 개

payments
servings
calls
hallOrders
kitchenOrders
tables

0x3C. Recovery가 시작되는 클라이언트 감지 지점

PosRealtimeBridge.tsx (line 323)에서 다음 상황을 감지

  • SSE payload JSON parse 실패
  • event mapping 실패
  • RELOAD_REQUIRED 수신
  • SSE open
  • resume/focus/online 주기 reconcile
if (!parsed) {
    reloadBaseline();
    return;
}
 
if (!event) {
    reloadBaseline();
    return;
}
 
if (eventName === "RELOAD_REQUIRED") {
    reloadBaseline();
}

watchdog timeout은 즉시 baseline을 호출하지 않고 먼저 연결을 끊고 재연결

// PosRealtimeBridge.tsx
 
cleanupSubscription();
scheduleReconnect();

재연결 성공 후 onOpen에서 recovery가 실행

onOpen: () => {
    startWatchdog();
    reloadBaseline();
}


0x60. IMPL — high traffic

Event traffic 증가
→ Redis Stream batch read
→ bounded queue에 임시 저장
→ 16개 shard worker가 매장별 queue 처리
→ emitter별 outbound queue로 전달
→ 한계 초과 시 RELOAD_REQUIRED 또는 연결 종료
→ 클라이언트 partial/full baseline 복구

0x61. Event traffic 증가 → batch 처리

Redis Stream reader는 이벤트를 하나씩 요청하지 않고 한 번에 최대 100개씩 읽.

val options = StreamReadOptions.empty()
    .count(100)
    .block(Duration.ofSeconds(2))
val records = redisTemplate.opsForStream<String, String>().read(
    options,
    StreamOffset.create(streamKey, ReadOffset.from(lastSeenId)),
) ?: emptyList()
 
for (record in records) {
    lastSeenId = record.id.toString()
    enqueueRecord(chanKey, streamKey, record)
}

짧은 burst가 들어오면 Redis에서 batch로 읽어 서버 queue에 넣음

Outbox → Redis Stream 발행도 batch

//PosOutboxRelay.kt
private val streamMaxLen = 10_000L
private val batchSize = 100
 
@Scheduled(fixedDelay = 500)
@Transactional
fun publishPendingEvents() {
    val events = outboxRepository.lockPublishable(batchSize)
}

즉 500ms마다 최대 100개의 outbox event를 Redis Stream으로 발행


0x62. Bounded queue로 짧은 burst buffer

// 매장별 broadcast queue
private fun queueOf(chanKey: String) =
    queues.computeIfAbsent(chanKey) {
        ArrayBlockingQueue(2000)
    }
 
// SSE 재연결 replay용 emitter별 queue
private fun replayQueueOf(emitter: SseEmitter) =
    replayQueues.computeIfAbsent(emitter) {
        ArrayBlockingQueue(500)
    }
 
// 실제 브라우저 송신용 emitter별 queue
private fun outboundQueueOf(emitter: SseEmitter) =
    outboundQueues.computeIfAbsent(emitter) {
        ArrayBlockingQueue(500)
    }
Queue범위크기
queues매장별 broadcast2,000
replayQueuesSSE emitter별 replay500
outboundQueuesSSE emitter별 실제 송신500

모두 ArrayBlockingQueue를 써, 메모리가 무한정 증가하지 않게 설계함.


0x63. Worker shard로 전송 작업 분산

private val senderWorkers = 16
private val senderExec = Executors.newFixedThreadPool(senderWorkers)
 
private fun shard(chanKey: String) =
    (chanKey.hashCode() and Int.MAX_VALUE) % senderWorkers
 
private val activeKeys =
    Array(senderWorkers) {
        ConcurrentLinkedQueue<String>()
    }

매장 ID인 chanKey의 hash를 이용해 16개 worker 중 하나에 할당

private fun markActive(chanKey: String) {
    if (activeFlag.putIfAbsent(chanKey, true) == null) {
        activeKeys[shard(chanKey)].offer(chanKey)
    }
}

activeFlag는 같은 매장 key가 worker queue에 반복적으로 쌓이는 것을 막음


0x64. Worker가 queue drain

private fun startSenders() {
    repeat(senderWorkers) { idx ->
        senderExec.submit {
            val myKeys = mutableListOf<String>()
 
            while (!Thread.currentThread().isInterrupted) {
                while (true) {
                    val key = activeKeys[idx].poll() ?: break
                    myKeys.add(key)
                }
 
                // 매장별 queue 처리
            }
        }
    }
}

한 worker가 특정 매장의 이벤트를 무제한으로 처리하지 않도록 한 번에 최대 50개만 꺼냅니다.

var processed = 0
 
while (processed < 50) {
    val msg = q.poll() ?: break
    broadcastNow(chanKey, msg)
    processed++
}

replay도 emitter별로 한 번에 최대 20개만 처리

drainReplay(
    chanKey,
    limitPerEmitter = 20,
)

이 제한은 하나의 매장이나 emitter가 worker를 장시간 독점하는 것을 완화


0x65. 처리 한계 안쪽 → SSE 전달 계속

Broadcast 단계는 매장에 연결된 모든 emitter의 outbound queue에 이벤트를 넣음

private fun broadcastNow(
    chanKey: String,
    msg: SseMsg,
) {
    emitters[chanKey]?.forEach { emitter ->
        enqueueOutbound(chanKey, emitter, msg)
    }
}

각 emitter는 별도 sender task가 outbound queue를 소비해 브라우저로 전송

while (
    !Thread.currentThread().isInterrupted &&
    outboundQueues.containsKey(emitter)
) {
    val msg = outboundQueues[emitter]
        ?.poll(30, TimeUnit.SECONDS)
        ?: continue
 
    if (!sendToEmitter(
            emitter,
            msg.name,
            msg.data,
            msg.id,
        )
    ) {
        unregisterEmitter(
            chanKey,
            emitter,
            "send-failed",
        )
        break
    }
}

따라서 느린 브라우저 하나의 실제 SSE 송신이 매장 broadcast worker를 직접 막지 않음

0x66. 처리 한계 초과 → incremental apply 포기

매장 broadcast queue overflow
private fun enqueue(
    chanKey: String,
    msg: SseMsg,
) {
    val q = queueOf(chanKey)
 
    if (q.offer(msg)) {
        markActive(chanKey)
        return
    }
 
    logger.warn("queue overflow: chanKey=$chanKey")
 
    q.clear()
    q.offer(
        SseMsg(
            "RELOAD_REQUIRED",
            "{}",
            "0-0",
        ),
    )
    markActive(chanKey)
}

overflow가 발생하면 밀려 있던 이벤트를 모두 버리고 RELOAD_REQUIRED로 교체

중간 event를 모두 전달할 수 없음
→ event cursor 기반 incremental apply를 더 이상 신뢰하지 않음
→ snapshot 기반 복구 요청
Replay queue overflow
private fun enqueueReplay(
    emitter: SseEmitter,
    msg: SseMsg,
) {
    val q = replayQueueOf(emitter)
 
    if (!q.offer(msg)) {
        q.clear()
        q.offer(
            SseMsg(
                "RELOAD_REQUIRED",
                "{}",
                "0-0",
            ),
        )
    }
}

누락 이벤트 replay가 queue 처리량을 초과해도 replay를 포기

Outbound queue overflow
if (q.offer(msg)) return
 
logger.warn(
    "[SSE-SLOW-CONSUMER] " +
        "storeId={} connectionId={} " +
        "outbound queue overflow",
    chanKey,
    connectionIds[emitter],
)
 
q.clear()
q.offer(
    SseMsg(
        "RELOAD_REQUIRED",
        "{}",
        "0-0",
    ),
)
 
emitter.complete()
unregisterEmitter(
    chanKey,
    emitter,
    "slow-consumer",
)

특정 브라우저가 SSE 소비 속도를 따라가지 못하면 slow consumer로 판단하고 그 연결만 종료. 다른 POS 연결은 계속 처리


0x67. RELOAD_REQUIRED → 복구 시작

클라이언트는 RELOAD_REQUIRED를 받으면 recovery를 실행

dispatch(reduxAction);
 
if (eventName === "RELOAD_REQUIRED") {
    reloadBaseline();
}

reloadBaseline()은 실제로 reloadChangedPosBaseline()을 실행

await dispatch(
    reloadChangedPosBaseline({
        storeId,
        sinceEventId:
            lastAppliedEventIdRef.current,
    }),
).unwrap();

0x68. Partial reload 또는 full baseline

유효한 cursor가 있으면 /changes를 이용해 partial reload를 시도

lastAppliedEventId
→ /notifications/changes
→ 변경된 order/table/menu/call ID 계산
→ 해당 entity만 재조회

변경량이 많거나 delta 계산을 신뢰할 수 없으면 full baseline으로 전환

if (changes.requiresReload) {
    await dispatch(
        reloadPosBaseline(storeId),
    ).unwrap();
}
if (
    changes.changedOrderIds.length > 100 ||
    changes.changedTableIds.length > 100 ||
    changes.changedMenuIds.length > 100 ||
    changes.changedCallRequestIds.length > 100
) {
    await dispatch(
        reloadPosBaseline(storeId),
    ).unwrap();
}


0x70. IMPL — slow consumer (single device failure)

특정 client 송신 지연
→ emitter별 outbound queue에 이벤트 적재
→ queue limit 500 안쪽이면 계속 송신
→ queue가 가득 차면 해당 emitter만 종료
→ FE EventSource error 감지
→ backoff + jitter 재연결
→ 연결 성공 후 partial/full baseline 복구

0x71. Client별 bounded queue로 지연 격리

각 SSE emitter마다 별도의 outbound queue와 sender task를 생성했음.

// SseNotificationService.kt
private val replayQueues =
    ConcurrentHashMap<SseEmitter, ArrayBlockingQueue<SseMsg>>()
 
private fun replayQueueOf(emitter: SseEmitter) =
    replayQueues.computeIfAbsent(emitter) {
        ArrayBlockingQueue(500)
    }
 
private val outboundQueues =
    ConcurrentHashMap<SseEmitter, ArrayBlockingQueue<SseMsg>>()
 
private fun outboundQueueOf(emitter: SseEmitter) =
    outboundQueues.computeIfAbsent(emitter) {
        ArrayBlockingQueue(500)
    }

client별 실제 송신 queue 크기를 500으로 제한했음.

매장 broadcast 이벤트는 연결된 emitter 각각의 outbound queue에 복사했음.

private fun broadcastNow(
    chanKey: String,
    msg: SseMsg,
) {
    emitters[chanKey]?.forEach { emitter ->
        enqueueOutbound(chanKey, emitter, msg)
    }
}

따라서 특정 client의 queue가 밀려도 다른 client의 outbound queue와 직접 섞이지 않았음.

0x72. Client별 sender task로 송신 분리

emitter가 등록될 때 해당 emitter 전용 queue와 sender task를 생성했음.

private fun registerEmitter(
    chanKey: String,
    emitter: SseEmitter,
    streamKey: String,
    connectionId: String,
) {
    connectionIds[emitter] = connectionId
    outboundQueueOf(emitter)
    startEmitterSender(chanKey, emitter, connectionId)
 
    // ...
}

전용 sender task가 해당 emitter의 queue만 소비했음.

private fun startEmitterSender(
    chanKey: String,
    emitter: SseEmitter,
    connectionId: String,
) {
    val task = emitterSenderExec.submit {
        while (
            !Thread.currentThread().isInterrupted &&
            outboundQueues.containsKey(emitter)
        ) {
            val msg = try {
                outboundQueues[emitter]
                    ?.poll(30, TimeUnit.SECONDS)
            } catch (e: InterruptedException) {
                Thread.currentThread().interrupt()
                null
            } ?: continue
 
            if (!sendToEmitter(
                    emitter,
                    msg.name,
                    msg.data,
                    msg.id,
                )
            ) {
                unregisterEmitter(
                    chanKey,
                    emitter,
                    "send-failed",
                )
                break
            }
        }
    }
 
    emitterSenderTasks[emitter] = task
}

특정 client의 네트워크 송신이 늦어도 매장 broadcast worker가 그 client의 emitter.send()를 직접 기다리지 않도록 구성했음.

0x73. Queue limit 안쪽이면 이벤트 전달 계속

outbound queue에 공간이 있으면 offer()가 성공하고 별도 조치 없이 반환했음.

private fun enqueueOutbound(
    chanKey: String,
    emitter: SseEmitter,
    msg: SseMsg,
) {
    val q = outboundQueues[emitter] ?: return
 
    if (q.offer(msg)) return
 
    // overflow 처리
}

queue가 500개 미만인 동안 sender task가 순차적으로 이벤트를 계속 전송했음.

0x74. Queue limit 초과 시 slow consumer 연결 종료

500개 queue가 가득 차서 offer()가 실패하면 해당 emitter를 slow consumer로 판단했음.

private fun enqueueOutbound(
    chanKey: String,
    emitter: SseEmitter,
    msg: SseMsg,
) {
    val q = outboundQueues[emitter] ?: return
    if (q.offer(msg)) return
 
    logger.warn(
        "[SSE-SLOW-CONSUMER] " +
            "storeId={} connectionId={} " +
            "outbound queue overflow",
        chanKey,
        connectionIds[emitter],
    )
 
    q.clear()
    q.offer(
        SseMsg(
            "RELOAD_REQUIRED",
            "{}",
            "0-0",
        ),
    )
 
    emitter.complete()
    unregisterEmitter(
        chanKey,
        emitter,
        "slow-consumer",
    )
}

처리 순서는 다음과 같았음.

outbound queue overflow
→ 기존 적체 이벤트 제거
→ RELOAD_REQUIRED 삽입
→ emitter 종료
→ 해당 client 연결 제거

0x75. 해당 client만 제거

unregisterEmitter()는 문제가 발생한 emitter의 queue와 sender task만 먼저 제거했음.

private fun unregisterEmitter(
    chanKey: String,
    emitter: SseEmitter,
    reason: String,
) {
    val connectionId = connectionIds.remove(emitter)
 
    replayQueues.remove(emitter)?.clear()
    outboundQueues.remove(emitter)?.clear()
    emitterSenderTasks.remove(emitter)?.cancel(true)
 
    emitters.computeIfPresent(chanKey) { _, list ->
        list.remove(emitter)
 
        if (list.isEmpty()) {
            cancelStreamReader(chanKey)
            queues.remove(chanKey)?.clear()
            activeFlag.remove(chanKey)
            null
        } else {
            list
        }
    }
}

같은 매장에 다른 emitter가 남아 있으면 매장 Redis reader와 broadcast queue를 유지했음. 따라서 단일 slow consumer 장애가 다른 POS 단말의 연결 종료로 전파되지 않았음.

0x76. FE가 연결 종료 감지

서버가 emitter를 종료하면 브라우저의 EventSource error handler가 실행됐음.

const sub = subscribeSse({
    url,
    withCredentials: true,
    handlers,
 
    onError: (error) => {
        console.error("[SSE ERROR]", {
            storeId,
            connectionId,
            error,
            mySeq,
        });
 
        if (
            closedRef.current ||
            connectSeqRef.current !== mySeq
        ) {
            return;
        }
 
        if (subRef.current === sub) {
            cleanupSubscription();
        }
 
        scheduleReconnect();
    },
});

기존 EventSource를 정리한 다음 재연결을 예약했음.

0x77. Backoff와 jitter로 재연결

재연결 지연은 2초부터 시작해 최대 30초까지 증가하도록 설정했음.

const INITIAL_RECONNECT_DELAY_MS = 2_000;
const MAX_RECONNECT_DELAY_MS = 30_000;
const scheduleReconnect = () => {
    if (closedRef.current) return;
    if (connectingRef.current) return;
    if (reconnectTimerRef.current != null) return;
 
    const baseDelayMs =
        reconnectDelayMsRef.current;
 
    const jitterMs = Math.floor(
        Math.random() *
        Math.min(2_000, baseDelayMs * 0.3),
    );
 
    const delayMs = baseDelayMs + jitterMs;
 
    reconnectDelayMsRef.current = Math.min(
        baseDelayMs * 2,
        MAX_RECONNECT_DELAY_MS,
    );
 
    reconnectTimerRef.current =
        window.setTimeout(() => {
            reconnectTimerRef.current = null;
            void connect();
        }, delayMs);
};

동시에 여러 POS가 재연결되더라도 정확히 같은 시점에 요청하지 않도록 jitter를 추가했음.

연결에 성공하면 backoff 값을 다시 2초로 초기화했음.

onOpen: () => {
    touchSseActivity();
 
    reconnectDelayMsRef.current =
        INITIAL_RECONNECT_DELAY_MS;
 
    startWatchdog();
    reloadBaseline();
},

0x78. 재연결 성공 후 recovery 실행

SSE 연결이 다시 열리면 reloadBaseline()을 호출했음.

실제 함수명과 달리 항상 full baseline만 실행하는 것은 아니며, lastAppliedEventId를 기준으로 partial recovery를 먼저 시도했음.

const reloadBaseline = () => {
    if (baselineReloadRef.current) return;
 
    const p = (async () => {
        await dispatch(
            reloadChangedPosBaseline({
                storeId,
                sinceEventId:
                    lastAppliedEventIdRef.current,
            }),
        ).unwrap();
    })();
 
    baselineReloadRef.current = p;
};

0x79. Partial reload 또는 full baseline

유효한 cursor가 없으면 바로 full baseline을 실행했음.

if (!sinceEventId || sinceEventId === "0-0") {
    await dispatch(
        reloadPosBaseline(storeId),
    ).unwrap();
 
    return {
        storeId,
        mode: "full" as const,
    };
}

유효한 cursor가 있으면 /changes를 통해 변경 entity만 조회했음.

await Promise.all([
    orderIds.length > 0
        ? fetchHallOrdersByIds(storeId, orderIds)
        : Promise.resolve(),
 
    orderIds.length > 0
        ? fetchKitchenOrdersByIds(storeId, orderIds)
        : Promise.resolve(),
 
    changes.changedCallRequestIds.length > 0
        ? dispatch(
            loadCallQueue({
                storeId,
                force: true,
            }),
        ).unwrap()
        : Promise.resolve(),
 
    changes.changedMenuIds.length > 0
        ? dispatch(loadMenuCatalog(storeId)).unwrap()
        : Promise.resolve(),
]);

partial recovery를 신뢰할 수 없으면 full baseline으로 전환했음.

if (changes.requiresReload) {
    await dispatch(
        reloadPosBaseline(storeId),
    ).unwrap();
 
    return {
        storeId,
        mode: "full" as const,
        reason: changes.reason,
    };
}


0x80. IMPL — redis publish failure

[Redis 장애에도 이벤트가 유실되지 않고 재시도된다] Redis Stream publish는 DB transaction과 분리되어 있으며, Redis 장애 시 Outbox row가 unpublished 상태로 남아 backoff 후 재시도된다.

정상 Flow
Outbox pending
→ Relay가 조회
→ Redis Stream publish
→ stream ID 저장
→ published 처리
실패 flow
Redis publish 실패
→ attempts 증가
→ nextAttemptAt 설정
→ 이후 relay에서 재조회
model
@Entity
@Table(name = "pos_event_outbox")
class PosEventOutbox(
    // 생략 ...
 
    @Column(name = "published_at")
    var publishedAt: Instant? = null,
 
    @Column(name = "redis_stream_id", length = 80)
    var redisStreamId: String? = null,
 
    @Column(name = "attempts", nullable = false)
    var attempts: Int = 0,
 
    @Column(name = "next_attempt_at", nullable = false)
    var nextAttemptAt: Instant = Instant.now(),
 
    @Column(name = "last_error", columnDefinition = "TEXT")
    var lastError: String? = null,
) {
    fun markPublished(streamId: String, now: Instant = Instant.now()) {
        redisStreamId = streamId
        publishedAt = now
        lastError = null
    }
 
    fun markFailed(error: Throwable) {
        attempts += 1
        lastError = (error.message ?: error::class.java.simpleName).take(1000)
        val delaySeconds = minOf(60L, 1L shl attempts.coerceAtMost(6))
        nextAttemptAt = Instant.now().plus(Duration.ofSeconds(delaySeconds))
    }
}
outbox relay
@Component
class PosOutboxRelay(...) {
	@Scheduled(fixedDelay = 500)
    @Transactional
    fun publishPendingEvents() {
        val events = outboxRepository.lockPublishable(batchSize)
        if (events.isEmpty()) return
 
        events.forEach { event ->
            try {
                val streamId = publishToRedis(event)
                event.markPublished(streamId)
            } catch (e: Exception) {
                event.markFailed(e)
                // log ...
            }
        }
    }
}
 
interface PosEventOutboxRepository : JpaRepository<PosEventOutbox, UUID> {
    @Query(
        value = """
            select *
            from pos_event_outbox
            where published_at is null
              and next_attempt_at <= now()
            order by created_at
            limit :limit
            for update skip locked
        """,
        nativeQuery = true,
    )
    fun lockPublishable(@Param("limit") limit: Int): List<PosEventOutbox>
}
 
consideration
  • lockPublishable()이 어떤 row를 선택하는지
  • 여러 relay instance가 같은 row를 동시에 처리할 수 있는지
  • Redis publish 성공 후 markPublished()가 같은 DB transaction에서 commit되는지
  • Redis publish는 성공했는데 DB commit이 실패할 수 있는지
  • 실패 row가 영구적으로 retry되는지, 최대 횟수가 있는지
  • poison event가 계속 retry되면 어떻게 되는지
한계
  • 현재 구조는 유실을 줄이지만 exactly-once 발행은 아니다
  • 아래 경우에 중복 발행 가능
Redis XADD 성공
→ 애플리케이션 종료 또는 DB commit 실패
→ Outbox에는 published_at이 여전히 null
→ 재시작 후 같은 Outbox 이벤트를 다시 XADD
대처법: FE deduplication
0. BE: SSE에 Outbox id inject
    private fun enqueueRecord(chanKey: String, streamKey: String, message: MapRecord<String, String, String>) {
        try {
            val type = message.value["type"] ?: "UNKNOWN"
            val payload = message.value["payload"] ?: "{}"
            val id = message.id.toString()
            if (type !in knownEventNames) {
                logger.warn("Unknown SSE event type: type={} id={} streamKey={}", type, id, streamKey)
                enqueue(chanKey, SseMsg("RELOAD_REQUIRED", "{}", id))
                return
            }
 
            enqueue(chanKey, SseMsg(type, payloadWithOutboxEventId(payload, message.value["eventId"]), id))
        } catch (e: Exception) {
            logger.error("Error processing redis message", e)
        }
    }
 
    private fun payloadWithOutboxEventId(payload: String, outboxEventId: String?): String {
        if (outboxEventId.isNullOrBlank()) return payload
 
        return try {
            val node = objectMapper.readTree(payload)
            if (!node.isObject) return payload
 
            (node as com.fasterxml.jackson.databind.node.ObjectNode)
                .put("__outboxEventId", outboxEventId)
            objectMapper.writeValueAsString(node)
        } catch (_: Exception) {
            payload
        }
    }
1. FE: 추출 및 dedupe
function readOutboxDedupeEventId(value: unknown): string | undefined {
    if (!value || typeof value !== "object") return undefined;
    const outboxEventId = (value as { __outboxEventId?: unknown }).__outboxEventId;
    return typeof outboxEventId === "string" && outboxEventId.trim()
        ? `outbox:${outboxEventId}`
        : undefined;
}
posListenerMiddleware.startListening({
    predicate: (action): action is PosRealtimeAction =>
        typeof action.type === "string" && action.type.startsWith("pos/realtime/"),
 
    effect: async (action, api) => {
        const meta = (action as { meta?: { streamEventId?: string; dedupeEventId?: string } }).meta;
        const streamEventId = normalizeRedisStreamId(meta?.streamEventId);
        const dedupeEventId = meta?.dedupeEventId;
        const dedupeId = dedupeEventId ?? (streamEventId !== "0-0" ? `stream:${streamEventId}` : undefined);
        const isReplayableStreamEvent = streamEventId !== "0-0";
 
        if (dedupeId) {
            const alreadyApplied = getPosState(api).stream.recentAppliedDedupeIds.includes(dedupeId);
            if (alreadyApplied) {
                if (isReplayableStreamEvent) {
                    api.dispatch(
                        posActions.markStreamEventApplied({
                            eventId: streamEventId,
                            dedupeId,
                            snapshotScopes: getSnapshotScopesForRealtimeAction(action.type),
                        }),
                    );
                }
                return;
            }
        }
 
        switch (action.type) {
            case "pos/realtime/payment/open": {
        // ... 이하 생략 ...


0x90. PS

0x91. SSE authentication and connection

SSE는 요청 헤더에 커스텀 Authorization을 싣기 어렵다. [(draft yet): SSE; Server-Sent Event]

1. 현재 [(draft yet): Authentication & Authorization]
  • kakao Login + Standard Token-based Authorization Scheme을 사용 중
2. ticket 기반 연결
  1. FE가 /api/pos/{storeId}/notifications/ticket을 먼저 호출한다.
  2. BE는 로그인 사용자와 매장 권한을 storeService.assertHost(storeId, user.id)로 확인한다.
  3. 권한이 맞으면 Redis에 30초 TTL의 1회용 티켓을 저장한다.
  4. FE는 /api/pos/{storeId}/notifications/subscribe?ticket=...&connectionId=...&lastEventId=...로 EventSource를 연다.
  5. subscribe 엔드포인트는 Spring Security에서 permitAll이지만, 실제 접근 제어는 ticket 검증으로 처리한다.
  6. BE는 티켓을 getAndDelete로 꺼내므로 같은 티켓은 재사용할 수 없다.
3. 의도
  • 이 구조의 의도는 EventSource의 헤더 제약을 우회하면서도, SSE URL만 알아서는 접속할 수 없게 만드는 것이다.
  • 주의할 점은 티켓 TTL이 30초라서, 티켓 발급 후 구독까지의 시퀀스가 빠르게 이어져야 한다는 것이다. 프론트는 연결 시도마다 새 티켓을 발급받는다.
  • 실제 서버 코드는 티켓을 조회하는 동시에 삭제한다. 이 때문에 같은 티켓으로 두 번 연결할 수 없다.
fun createTicket(storeId: UUID): String {
    val ticket = UUID.randomUUID().toString().replace("-", "")
    redisTemplate.opsForValue().set("ticket:$ticket", storeId.toString(), Duration.ofSeconds(30))
    return ticket
}
 
fun subscribe(
    expectedStoreId: UUID,
    ticket: String,
    connectionId: String?,
    lastEventId: String?
): SseEmitter {
    val storeIdStr = redisTemplate.opsForValue().getAndDelete("ticket:$ticket")
        ?: throw BusinessException(ErrorCode.INVALID_TICKET)
 
    if (storeIdStr != expectedStoreId.toString()) {
        throw BusinessException(ErrorCode.INVALID_TICKET)
    }
 
    val emitter = SseEmitter(60 * 60 * 1000L)
    ...
}
  • 프론트는 EventSource 연결 전에 항상 ticket API를 호출하고, lastEventId가 0-0이 아닐 때만 구독 URL에 붙인다.
async function getPosTicket(storeId: string): Promise<string> {
    const res = await apiRequest(`/api/pos/${storeId}/notifications/ticket`, {
        method: "GET",
        credentials: "include",
    });
    ...
}
 
function getSseUrl(
    storeId: string,
    ticket: string,
    connectionId: string,
    lastEventId: string,
): string {
    const params = new URLSearchParams({ ticket, connectionId });
    if (lastEventId !== "0-0") params.set("lastEventId", lastEventId);
    return `/api/pos/${storeId}/notifications/subscribe?${params.toString()}`;
}