느린 구독자가 생중계를 멈췄다 · 확인·기한·회수의 주인을 정한다 · 실습
느린 구독자를 격리해 생중계를 구한다
목표
큐 정책과 ACK 기한, 연결 수명을 구현해 실제 TCP에서 느린 구독자를 격리합니다.
왜 중요한가
구독자 한 명의 지연이 전체 발행을 막거나 메모리를 계속 사용하게 할 수 있습니다. 데이터의 손실 허용 범위와 업무 확인을 구분하고, 취소 때도 자원을 회수하는 프로그램을 만듭니다. Python 함수·예외·async/await 기초가 필요합니다. 표준 라이브러리만 사용하며 설치나 인터넷은 필요 없습니다.
단계
1. 끝없이 쌓이는 큐에 상한을 건다 — delivery.py에 make_queue(capacity)를 구현하세요. 비어 있는 asyncio.Queue를 반환하고 maxsize는 입력과 같아야 합니다. 양의 int만 허용하며 bool·실수·문자열·0·음수는 ValueError입니다.
2. 기다리되 취소된 발행은 남기지 않는다 — async offer_wait(queue, item)를 추가하세요. 공간이 생길 때까지 기다려 넣고 None을 반환합니다. 대기 중 취소는 CancelledError로 전달하며 기존 항목과 순서를 보존합니다. 취소된 항목을 나중에 삽입하지 않습니다.
3. 옛 좌표를 버리고 현재 위치를 남긴다 — 동기 함수 offer_latest(queue, item)를 추가하세요. 공간이 있으면 삽입 후 None, 가득 차면 가장 오래 기다린 항목 하나를 제거하고 새 항목을 넣은 뒤 제거한 항목을 반환합니다. 버린 항목의 task_done도 한 번 대응시킵니다.
4. 멈춘 구독자를 정책 예외로 알린다 — Exception 하위 클래스 SlowConsumer와 동기 offer_disconnect(queue, item)를 구현하세요. 공간이 있으면 삽입 후 None을 반환합니다. 가득 차면 기존 큐를 바꾸지 않고 SlowConsumer를 냅니다. 연결을 닫는 일은 이 함수의 책임이 아닙니다.
5. 중복은 건너뛰고 누락은 숨기지 않는다 — Exception 하위 클래스 Gap과 Cursor(last=-1)를 추가하세요. last는 인스턴스별 공개 속성이고 -1~2147483647 범위 int입니다. accept(seq)는 0~2147483647 int만 받습니다. seq<=last면 False와 상태 유지, seq==last+1이면 last 갱신 후 True, 더 큰 순번이면 Gap과 상태 유지입니다. bool 등 잘못된 타입·범위는 ValueError입니다. 메모리 커서이며 영속 업무 처리는 구현하지 않습니다.
6. 송신과 업무 확인에 한 번만 시간을 준다 — async send_one(reader, writer, seq, timeout, clock=None)을 만드세요. seq는 5단계와 같은 유효 순번이며 timeout은 bool을 제외한 양의 유한 int/float이고 잘못되면 ValueError입니다. 기본 시계는 time.monotonic, 주입 시계는 인자 없이 초를 반환합니다. ASCII EVENT 순번 뒤 개행을 writer.write로 한 번 쓰고 drain을 기다린 다음 reader.readline으로 정확히 같은 순번의 ACK와 개행을 받아야 None입니다. EOF·잘못된 ACK는 ConnectionError, 전체 기한 초과는 TimeoutError입니다. drain과 ACK는 같은 마감 시각을 공유하며 완료 후에도 기한을 확인합니다. 이 함수는 writer를 닫지 않고 오류·취소를 전달합니다.
7. 비어 있거나 실패해도 연결을 회수한다 — async serve_queue(reader, writer, queue, timeout)를 추가하세요. 양의 유한 timeout이 전달됩니다. 반복해서 get한 seq를 send_one으로 처리하며, 그 호출이 끝나거나 실패·취소된 뒤에만 인수한 항목에 task_done을 정확히 한 번 호출합니다. 빈 큐 대기를 포함해 모든 종료 경로에서 writer.close 후 wait_closed를 timeout 안에 기다립니다. 종료 대기의 시간 초과·ConnectionError에는 writer.transport.abort로 회수하고 원래 송신 오류·취소를 숨기지 않습니다. 아직 큐에 남은 항목 정리는 호출자의 책임입니다.
8. 느린 한 명을 분리하고 생중계를 계속한다 — 동기 broadcast(queues, item, disconnect)를 완성하세요. queues는 이름→큐 딕셔너리이고 순회 스냅샷의 각 큐에 offer_disconnect를 적용합니다. SlowConsumer인 대상만 딕셔너리에서 제거하고 disconnect(name)를 한 번 호출한 뒤 다른 구독자에게 계속 전달합니다. 다른 예외는 전달하고 정상 반환은 None입니다. 콜백은 예외 없이 해당 작업자 취소·연결 회수·잔여 큐 정리를 담당합니다. 실제 TCP 검사에서 slow는 첫 ACK를 보류하고 fast는 매번 확인합니다. 0~8 이벤트를 fast가 모두 받고 slow만 분리되어야 합니다.
참고
모든 함수는 /root/realtime/delivery.py 하나에 둡니다. mkdir -p /root/realtime로 작업 폴더를 만드세요. 예시의 아직 구현하지 않은 함수는 틀로 두되 완료한 함수를 덮어쓰지 마세요. 채점은 이전 단계 계약도 검사합니다. 한 단계 실행은 5초 제한이며 학습자가 작성할 시간의 제한이 아닙니다. TCP 서버·임시 포트·구독자는 검사기가 준비하고 종료합니다. 이 실습은 커널 버퍼 포화·인터넷 성능·영속 전달·실제 재접속 복구를 보장하지 않습니다. 실습 세션이 끝나면 파일은 유지되지 않으니 필요한 코드는 종료 전에 별도로 보관하세요.
단계 8개
- 끝없이 쌓이는 큐에 상한을 건다
- 기다리되 취소된 발행은 남기지 않는다
- 옛 좌표를 버리고 현재 위치를 남긴다
- 멈춘 구독자를 정책 예외로 알린다
- 중복은 건너뛰고 누락은 숨기지 않는다
- 송신과 업무 확인에 한 번만 시간을 준다
- 비어 있거나 실패해도 연결을 회수한다
- 느린 한 명을 분리하고 생중계를 계속한다