Devin.KR

메시지 경계와 역압

90분 안팎

학습 목표

순수 Python 메시지 큐에서 timestamp·단위·용량과 DROP 정책을 검사합니다.

개념

큐는 속도 차이를 보관하지만 없애지는 않습니다

수집은 100ms마다 결과를 만들고 전송은 일시적으로 느려질 수 있습니다. 큐가 있으면 잠깐의 차이를 흡수하지만 전송이 계속 느리면 결국 가득 찹니다. 큐 크기를 늘리는 것만으로 평균 처리율 부족을 해결할 수 없습니다. 이번 레슨은 메시지의 시각·단위·용량과 폐기 이유를 검사하는 Python 모델을 만들고, 앞 모듈 상태 머신에 C 메시지 큐를 붙이는 미션의 경계를 정합니다.

역압은 뒤 단계의 처리 능력이 앞 단계의 생산에 영향을 주도록 전달하는 생각입니다. 기다리는 push로 생산자를 멈추거나 샘플 주기를 낮추는 선택이 있습니다. 여기의 결정적 메시지 모델은 생산자를 기다리게 하지 않고 DROP_NEW로 손실을 명시합니다. 이는 과부하 처리 정책이며 자동으로 생산 속도를 줄이는 역압 구현은 아닙니다. 센서 주기를 바꿀지, 오래된 값 대신 최신 값을 남길지는 장치 요구사항에 따라 따로 승인할 결정입니다.

메시지를 의미 있는 값 묶음으로 정의합니다

Python 입력 메시지는 seq·timestamp_ms·temp_mc·unit 네 필드를 갖습니다. seq는 생산 시도 순번, timestamp_ms는 성공한 측정 시각, temp_mc는 밀리 섭씨도 정수, unit은 mC 문자열입니다. timestamp를 큐에 넣은 시각이나 전송 시각으로 다시 찍지 않습니다. 대기 때문에 값이 오래되어도 원래 측정 시각이 남아 있어야 소비자가 신선도를 판단할 수 있습니다. 전송 시각이 필요하면 별도 필드를 추가합니다.

단위 문자열을 검사하지 않고 temp_mc 숫자만 읽으면 25를 25밀리도로 볼지 25도로 볼지 모호해집니다. C 미션의 BusRecord는 필드명과 명세로 단위를 고정하므로 구조체에 매번 문자열을 저장하지 않습니다. Python 연습은 경계를 눈에 보이게 하기 위해 unit을 검사합니다. 실제 네트워크와 장치 사이에서 포맷을 바꾸면 버전·바이트 순서·범위 검증도 계약에 넣어야 합니다. 이번 연습의 JSON을 UART의 기존 13바이트 프레임과 동일한 포맷이라고 부르지 않습니다.

입력은 capacity·now_ms·messages를 가진 JSON 객체입니다. capacity는 1부터 16까지, now_ms는 0부터 1000000까지의 정수입니다. messages는 최대 1000건입니다. seq는 0부터 1000000까지의 정수이며 입력 순서대로 엄격히 증가합니다. 순번 간격은 허용하되 중복과 역행은 INVALID입니다. timestamp_ms는 0부터 now_ms까지이고 temp_mc는 -40000부터 125000까지입니다. 범위는 이 연습의 센서 계약이며 모든 온도 센서의 측정 범위를 뜻하지 않습니다.

모든 필수 필드와 타입을 먼저 검사합니다. 정수처럼 보이는 소수나 true를 허용하지 않으며 unit은 mC만 받습니다. 미래 시각은 시계 동기화가 어긋난 입력으로 보고 INVALID입니다. 오류 하나가 있으면 전체 입력을 거절하고 부분 전송 결과를 출력하지 않습니다. 실제 C 미션은 uint32_t의 차이로 시계 순환을 처리하지만 Python 연습은 시계 순환을 포함하지 않는 제한된 관찰 구간입니다. 두 모델의 시간 계약을 섞지 않습니다.

넣을 때의 손실과 꺼낼 때의 손실을 셉니다

메시지를 입력 순서대로 전부 넣은 다음 now_ms에서 큐를 순서대로 비웁니다. 넣는 동안 소비자는 실행하지 않습니다. 자리가 있으면 메시지 사본을 넣고, 가득 차면 새 메시지를 버리며 DROP_NEW 수를 올립니다. 기존 메시지를 몰래 지우는 DROP_OLD는 사용하지 않습니다. 이 일괄 입력 모델은 생산·소비 교대 실행을 흉내 내지 않으므로 capacity를 초과한 입력 수가 곧 DROP_NEW 수입니다.

꺼낸 메시지의 나이는 now_ms에서 timestamp_ms를 뺀 값입니다. 나이 300ms 이상이면 STALE로 버리고, 299ms까지는 sent에 seq를 넣습니다. 따라서 capacity 2·now_ms 400에 시각 101, 100, 400인 메시지 세 개를 넣으면 첫 두 개만 보관됩니다. 첫 메시지는 나이 299로 전송, 두 번째는 300으로 폐기, 세 번째는 신선해도 넣는 순간 공간이 없어 DROP_NEW입니다. 먼저 신선도 검사를 하고 오래된 칸을 미리 치우면 이 계약과 다른 정책이 됩니다.

출력은 sent=순번목록 DROP_NEW=N STALE=M입니다. 전송 목록은 쉼표로 연결하고 없으면 sent=-로 적습니다. 빈 messages는 정상 입력이며 두 카운터가 0입니다. DROP_NEW와 STALE을 합쳐 오류 한 개로 기록하지 않습니다. 같은 손실 수라도 생산 과부하와 소비 지연은 처방이 다릅니다. 정상 입력에서 보관된 수는 전송 수와 STALE 수의 합이며 전체 입력 수는 여기에 DROP_NEW를 더한 값과 같습니다.

앞 모듈의 수집과 전송을 큐로 분리합니다

미션 압축 embedded-concurrency-queues는 m07-state-recovery의 solution 소스와 테스트를 포함합니다. recovery_tick은 기존 회귀를 위해 유지하고 새 수집 경로는 pipeline_collect입니다. 단일 가상 루프에서 timed_button → timer_poll → pipeline_collect → pipeline_dispatch 순으로 호출합니다. 새 경로와 옛 recovery_tick 또는 bus_consume을 함께 호출하면 동일 이벤트의 소유권이 겹치므로 피합니다. README와 CONCURRENCY-SPEC.md에서 각 함수의 소유 데이터를 확인합니다.

pipeline_collect는 MEASURE에서 성공한 pending을 확정하고 TRANSMIT 분기에서 용량16 MessageQueue로 값 복사합니다. 수집은 Buffer·SPI·UART에 쓰지 않고 다음 상태를 IDLE로 돌립니다. pipeline_send가 꺼낸 메시지를 신선도 검사한 뒤 기존 프레임 인코딩·SPI 기록·Buffer 저장·UART 공개에 사용합니다. seq와 ADC 원시값 raw는 내부 메시지에 추가하지만 UART 프레임 형식은 앞 모듈의 계약을 유지합니다. 내부 순번을 전송 프레임의 새 필드로 착각하지 않습니다.

분리·정지·ERROR·명시 재연결 시 대기 큐를 비우고 제거 수를 CANCEL에 누적합니다. 저장된 Buffer·SPI와 마지막 UART 이력은 유지합니다. 재연결 뒤 새 측정만 큐에 들어오며 seq는 다시 0으로 돌리지 않습니다. CANCEL은 이미 보관했던 메시지의 무효화이고 DROP_NEW는 넣지 못한 새 메시지의 손실입니다. 분리 이후 큐에 남은 오래된 온도가 먼저 전송되는지 검사해야 복구가 실제 공개 경계까지 이어졌다고 말할 수 있습니다.

C 수집 상태와 HAL은 단일 루프가 소유합니다. 실제 pthread 검사는 공유 큐에 대해 따로 실행합니다. 전체 Recovery 구조체와 장치 HAL을 두 스레드가 동시에 호출하도록 만든 프로젝트가 아닙니다. 같은 큐 규칙을 검증해도 RTOS 태스크의 실행 기한·ISR 안전성·전송 장치 직렬화는 추가 설계가 필요합니다. task_model.py는 비선점 마감과 결정적 증가 유실을 검사하며 pthread 시간 측정값으로 실시간 보장을 주장하지 않습니다.

통합 증거와 실패 코드의 책임을 남깁니다

미션은 make test로 이전 회귀, 큐 10000건, 과부하 20건 중 4건 DROP_NEW, 나이 299/300 경계, 분리·정지 CANCEL과 재연결 새 값 전송을 검사합니다. starter는 큐의 seq 보존 TODO가 실패합니다. 그 결함을 수정해 기존 테스트까지 유지하고 CONCURRENCY-EVIDENCE.md에 실제 요약과 시각 조건을 기록합니다. 저장소 원고의 정상 출력을 보고 시작 코드 결과도 같을 것이라고 적지 않습니다.

전송 함수의 BUS_TIMEOUT·BUS_FULL·BUS_LENGTH 같은 반환은 소비자 실패를 뜻합니다. SPI 실패 시 CS를 해제하고 이전 공개 기록은 보존합니다. pipeline_dispatch가 전송 실패를 ERROR·valid 초기화·오류 LED와 남은 큐 CANCEL에 연결합니다. STALE도 자동 재시도하지 않고 명시 재연결을 기다립니다. 제출 문서에는 전송 실패를 상태·LED 정책에 어떻게 연결했는지와 재시도하지 않는 메시지의 처리 이유를 적습니다. JSON 필드가 없어서 INVALID가 나오면 키를, DROP_NEW가 예상보다 많으면 용량과 생산·소비 순서를, STALE만 틀리면 측정 시각과 등호 경계를 점검합니다.

ROS 메시지의 필드 이름·타입 원리는 더 읽기의 장에서 확인하되 ROS 2 설치나 메시지 생성은 하지 않습니다. 이 레슨의 제출물은 단위와 시간의 계약, 세 손실 이유의 구분, 실제 실행 증거입니다. 수집률을 낮추거나 최신 값을 남기는 다른 정책으로 바꾸고 싶다면 먼저 새로운 acceptance와 경계 입력을 정의하여 기존 FIFO 증거와 비교합니다.

따라하기

넣기와 비우기의 경계

모든 메시지를 먼저 넣습니다. 신선한 셋째 메시지도 포화 시 거절됩니다. 두 번째는 나이300 경계입니다. 표준 입력은 다음과 같습니다.

{"capacity": 2, "now_ms": 400, "messages": [{"seq": 0, "timestamp_ms": 101, "temp_mc": 1000, "unit": "mC"}, {"seq": 1, "timestamp_ms": 100, "temp_mc": 1000, "unit": "mC"}, {"seq": 2, "timestamp_ms": 400, "temp_mc": 1000, "unit": "mC"}]}

모든 메시지를 먼저 넣습니다. 신선한 셋째 메시지도 포화 시 거절됩니다. 두 번째는 나이300 경계입니다.

import json,sys

def integer(x,low,high):
    return type(x) is int and low<=x<=high
try:
    data=json.load(sys.stdin)
    cap=data['capacity']; now=data['now_ms']; messages=data['messages']
    if not integer(cap,1,16) or not integer(now,0,1000000) or type(messages) is not list or len(messages)>1000: raise ValueError()
    previous=-1
    for m in messages:
        if type(m) is not dict or not integer(m['seq'],0,1000000) or m['seq']<=previous: raise ValueError()
        if not integer(m['timestamp_ms'],0,now) or not integer(m['temp_mc'],-40000,125000) or m['unit']!='mC': raise ValueError()
        previous=m['seq']
    queue=[]; dropped=0; stale=0; sent=[]
    for m in messages:
        if len(queue)==cap: dropped+=1
        else: queue.append(dict(m))
    for m in queue:
        if now-m['timestamp_ms']>=300: stale+=1
        else: sent.append(str(m['seq']))
    print('sent='+(','.join(sent) or '-'),f'DROP_NEW={dropped}',f'STALE={stale}')
except (ValueError,KeyError,TypeError):
    print('INVALID')

실행 결과

sent=0 DROP_NEW=1 STALE=1

측정 시각과 단위 검증

C 단위는 이 경계에서 자동 변환하지 않고 INVALID로 거절합니다. 입력 전체 검사가 전송보다 앞서는지 확인합니다. 표준 입력은 다음과 같습니다.

{"capacity": 2, "now_ms": 100, "messages": [{"seq": 0, "timestamp_ms": 100, "temp_mc": 1000, "unit": "C"}]}

C 단위는 이 경계에서 자동 변환하지 않고 INVALID로 거절합니다. 입력 전체 검사가 전송보다 앞서는지 확인합니다.

import json,sys

def integer(x,low,high):
    return type(x) is int and low<=x<=high
try:
    data=json.load(sys.stdin)
    cap=data['capacity']; now=data['now_ms']; messages=data['messages']
    if not integer(cap,1,16) or not integer(now,0,1000000) or type(messages) is not list or len(messages)>1000: raise ValueError()
    previous=-1
    for m in messages:
        if type(m) is not dict or not integer(m['seq'],0,1000000) or m['seq']<=previous: raise ValueError()
        if not integer(m['timestamp_ms'],0,now) or not integer(m['temp_mc'],-40000,125000) or m['unit']!='mC': raise ValueError()
        previous=m['seq']
    queue=[]; dropped=0; stale=0; sent=[]
    for m in messages:
        if len(queue)==cap: dropped+=1
        else: queue.append(dict(m))
    for m in queue:
        if now-m['timestamp_ms']>=300: stale+=1
        else: sent.append(str(m['seq']))
    print('sent='+(','.join(sent) or '-'),f'DROP_NEW={dropped}',f'STALE={stale}')
except (ValueError,KeyError,TypeError):
    print('INVALID')

실행 결과

INVALID

값 복사로 생산자 변경 격리

다음 측정에서 생산자 변수를 바꿔도 대기 메시지에는 이전 seq와 측정 시각이 남습니다. 내부 필드가 정수라 얕은 사본으로 충분합니다.

m={'seq':0,'timestamp_ms':100,'temp_mc':1000,'unit':'mC'}
queue=[dict(m)]
m['seq']=1; m['timestamp_ms']=200
print(queue[0]['seq'],queue[0]['timestamp_ms'])

실행 결과

0 100

누적 C 미션 실행

embedded-concurrency-queues 압축 폴더에서 실행합니다. 아래 요약은 완성본의 실제 출력입니다. 이어 make test로 앞 모듈 회귀도 검사합니다. pipeline_dispatch가 전송 오류를 ERROR·LED로 연결하는 코드와 test_pipeline.c의 보존 조건을 함께 읽습니다.

실행 명령: make -s queue-test pipeline-test model-test

실행 결과

queue: 135 checks, 0 failures; delivered=10000
pipeline: 55 checks, 0 failures; DROP_NEW=4 STALE=1 CANCEL=17 sent=2
task-model: 5 checks, 0 failures; missed=1 unprotected=1 protected=2

확인 문제

실습

capacity·now_ms·messages JSON을 읽어 메시지 전체의 정수 타입·범위·필드·mC 단위·엄격 증가 seq를 먼저 검사합니다. 순서대로 모두 넣고 DROP_NEW를 센 뒤 now_ms에서 비웁니다. 나이300ms 이상은 STALE, 그 외는 sent에 seq를 넣습니다. 출력은 sent=목록 DROP_NEW=N STALE=M이며 빈 목록은 sent=-입니다. 오류 입력은 INVALID 한 줄입니다. 본문의 범위와 미래 시각 금지 규칙을 적용합니다. C 미션은 별도 다운로드에서 이어갑니다.

모범 답안
import json,sys

def integer(x,low,high):
    return type(x) is int and low<=x<=high
try:
    data=json.load(sys.stdin)
    cap=data['capacity']; now=data['now_ms']; messages=data['messages']
    if not integer(cap,1,16) or not integer(now,0,1000000) or type(messages) is not list or len(messages)>1000: raise ValueError()
    previous=-1
    for m in messages:
        if type(m) is not dict or not integer(m['seq'],0,1000000) or m['seq']<=previous: raise ValueError()
        if not integer(m['timestamp_ms'],0,now) or not integer(m['temp_mc'],-40000,125000) or m['unit']!='mC': raise ValueError()
        previous=m['seq']
    queue=[]; dropped=0; stale=0; sent=[]
    for m in messages:
        if len(queue)==cap: dropped+=1
        else: queue.append(dict(m))
    for m in queue:
        if now-m['timestamp_ms']>=300: stale+=1
        else: sent.append(str(m['seq']))
    print('sent='+(','.join(sent) or '-'),f'DROP_NEW={dropped}',f'STALE={stale}')
except (ValueError,KeyError,TypeError):
    print('INVALID')

더 읽기

면접 질문

  • 센서가 응답하지 않을 때 점검할 순서를 설명해 주시면 됩니다.
  • 과부하·오래된 값·복구 취소 손실을 어떤 근거로 구분하는지 설명해 주세요.