Tech Wiki

[Rust 실전 로드맵 25] Tokio 채널·select·취소로 종료 가능한 작업자 만들기

bounded mpsc에 작업을 넣고 select!로 종료 신호를 받는 것만으로는 작업자가 안전하게 멈추지 않습니다. 큐가 가득 찼을 때 무엇을 반환할지, 이미 받은 작업을 끝낼지 취소할지, 종료 요청이 겹치면 어떤 상태가 이길지까지 계약으로 정해야 합니다. 마지막에는 spawn한 task를 반드시 join해야 합니다.

이 글은 Rust 2024 edition, rustc·Cargo 1.98.1, Tokio 1.53.1을 기준으로 합니다. 검사 작업자 하나를 만들어 bounded backpressure, Running < Drain < Cancel 상태 전이, select!의 취소 안전성, paused time 테스트를 확인합니다. 네트워크 I/O는 사용하지 않습니다.

1. bounded mpsc로 유입량에 상한을 둡니다

Tokio의 bounded mpsc는 지정한 수만큼 메시지를 버퍼링합니다. 버퍼가 가득 차면 send().await는 수신자가 자리를 만들 때까지 기다립니다. 생산자가 소비자보다 계속 빠를 때 대기와 메모리 사용의 경계를 큐 용량으로 드러내는 backpressure입니다.

대기하지 않을 호출자에게는 try_send가 더 분명합니다. 예제의 try_submit은 빈자리가 없으면 TrySendError::Full(Check)를 그대로 돌려줍니다. 테스트는 작업자가 첫 작업을 시작했다는 event를 받은 다음 용량 1인 큐에 두 번째 작업을 넣고, 세 번째 작업이 Full인지 확인합니다. scheduler가 어느 순간에 작업자를 poll할지 추측하지 않습니다.

monitor
    .try_submit(Check::new("buffered", Duration::from_secs(60)))
    .unwrap();
let error = monitor
    .try_submit(Check::new("rejected", Duration::from_secs(60)))
    .unwrap_err();
assert!(matches!(error, TrySendError::Full(check) if check.name() == "rejected"));

1.1. 잘못된 capacity는 channel 생성 전에 typed error로 막습니다

Tokio 1.53.1의 mpsc::channel은 capacity가 0이거나 너무 크면 panic합니다. 라이브러리 경계에서 그대로 panic하게 두지 않고 설정 오류로 바꾸는 편이 호출자에게 낫습니다. start_monitor는 channel을 만들기 전에 두 경계를 검사합니다.

if queue_capacity == 0 {
    return Err(StartError::ZeroCapacity);
}
if queue_capacity > Semaphore::MAX_PERMITS {
    return Err(StartError::CapacityTooLarge {
        requested: queue_capacity,
        maximum: Semaphore::MAX_PERMITS,
    });
}

let (jobs, receiver) = mpsc::channel(queue_capacity);

상한 초과 테스트는 Semaphore::MAX_PERMITS.checked_add(1)로 잘못된 값을 만듭니다. 테스트 코드 자체에서 정수 overflow가 생기지 않도록 한 선택입니다. 결과는 StartError::CapacityTooLarge { requested, maximum }이며 capacity 0StartError::ZeroCapacity입니다.

1.2. 입력을 닫은 뒤 새 submitter를 만들 수 없습니다

close_input은 monitor가 가진 sender를 버리고 공유 admission을 동기적으로 닫습니다. 이후 submitter()는 unwrap이나 panic 대신 Result를 반환합니다.

pub fn submitter(&self) -> Result<Submitter, SubmitterError> {
    self.jobs
        .as_ref()
        .map(|jobs| Submitter {
            jobs: jobs.clone(),
            admission: Arc::clone(&self.admission),
        })
        .ok_or(SubmitterError::InputClosed)
}

이미 복제된 Submitter도 같은 admission state를 봅니다. stop이나 close_input이 이 state를 닫은 뒤 시작한 submit은 sender를 가지고 있어도 실패합니다. 종료 전에 admission을 획득한 호출만 in-flight counter에 남고, RAII guard가 success, send failure, future drop 어느 경로에서도 counter를 줄입니다.

2. watch 알림과 단조 상태를 분리합니다

watch channel은 마지막으로 보낸 값 하나만 보관합니다. 종료 상태 전달에는 잘 맞지만 순서대로 쌓이는 event log는 아닙니다. Cancel을 보낸 뒤 Drain을 보내면 worker가 보기 전에 최신 값이 Drain으로 바뀔 수 있습니다. 취소를 되돌릴 수 있게 됩니다.

예제는 watch를 깨우기 위한 신호로 쓰고 실제 우선순위는 AtomicU8에 둡니다. 값 0, 1, 2가 각각 Running, Drain, Cancel이며 허용되는 방향은 Running < Drain < Cancel뿐입니다.

2.1. drain은 올릴 수만 있고 cancel은 terminal입니다

정상 종료 요청은 fetch_max(1)을 사용합니다. 현재 상태가 이미 2라면 낮아지지 않습니다. 취소 요청은 2를 저장합니다. 두 요청 모두 짧은 mutex 임계 구역에서 admission을 먼저 닫고 상태를 갱신한 뒤 watch를 보냅니다. mutex는 await를 가로질러 잡지 않습니다.

let mut admission = self.admission.state.lock().expect("admission mutex poisoned");
if admission.worker_stopped {
    return Err(RequestError::WorkerStopped);
}
admission.open = false;
self.stop_state.store(2, Ordering::Release);
self.stop
    .send(StopSignal::from_state(&self.stop_state))
    .map_err(|_| RequestError::WorkerStopped)

그래서 worker가 관찰하기 전에 Cancel 다음 Drain이 들어와도 최종 상태는 Cancel입니다. Drain 다음 Cancel도 취소로 승격됩니다. 더 까다로운 경우도 있습니다. worker가 이미 drain을 관찰해 queue를 비우기 시작한 뒤라도 나중에 들어온 cancel은 terminal 상태로 승격되어야 합니다.

2.2. 알림 수신과 상태 판정을 같은 것으로 보지 않습니다

watch::Receiver::changed는 값이 바뀌었다는 사실을 기다립니다. 어떤 종료 모드를 적용할지는 borrow_and_update()만 보고 결정하지 않고 atomic state를 다시 읽어 판정합니다. 최신 값 하나만 남기는 watch의 성질과 되돌아가지 않는 정책을 분리한 셈입니다.

이 구조는 모든 상태 기계를 watch + atomic으로 만들라는 일반 해법은 아닙니다. 여기서는 상태가 세 단계이고 병합 규칙이 max로 표현되기 때문에 맞습니다. 요청 순서를 모두 보존해야 한다면 별도 queue가 필요합니다.

3. select!에는 취소 안전한 branch만 놓습니다

tokio::select!는 여러 branch를 기다리다가 먼저 완료된 branch의 handler를 실행하고 나머지 future를 취소합니다. 여기서 취소는 잃은 future를 drop한다는 뜻입니다. 따라서 select!가 임의의 future를 안전하게 취소해 주는 것은 아닙니다. branch로 넣을 operation 자체가 중간 drop 뒤에도 계약을 지켜야 합니다.

3.1. recv와 changed의 문서화된 보장을 사용합니다

Tokio 1.53.1의 mpsc::Receiver::recv는 cancel safe입니다. 다른 branch가 먼저 끝나면 해당 호출이 메시지를 받지 않았다고 보장합니다. watch::Receiver::changed도 cancel safe이며, 잃은 호출 때문에 값이 seen으로 표시되지 않습니다. 작업자 loop가 이 두 API를 함께 기다릴 수 있는 근거입니다.

let next = tokio::select! {
    biased;
    changed = stop.changed() => {
        let _ = stop.borrow_and_update();
        if StopSignal::from_state(stop_state) == StopSignal::Cancel {
            cancel_queue(jobs, admission, events, cancelled).await;
            return StopReason::Cancelled;
        }
        if changed.is_err() {
            watch_open = false;
        }
        continue;
    }
    next = jobs.recv() => next,
};

반대로 Sender::sendselect! branch에 넣었다가 다른 branch가 이기면 메시지는 전송되지 않지만 값도 drop되어 잃을 수 있습니다. 메시지 손실을 피하려면 reserve로 capacity를 먼저 확보한 뒤 Permit으로 보냅니다. 이 예제는 sendselect! 안에 넣지 않습니다. submit이 성공하면 작업 소유권은 worker queue로 넘어갑니다.

3.2. biased는 우선순위를 주지만 공정성 책임도 넘깁니다

기본 select!는 먼저 검사할 branch를 pseudo-random하게 고릅니다. 어느 정도 공정성을 제공하지만 특정 순서를 계약하지 않습니다. 현재 예제의 네 select! site는 모두 biased;를 명시해 위에서 아래로 poll합니다. 종료 branch가 작업 수신이나 timer보다 앞섭니다.

이 우선순위는 shutdown latency를 제한하려는 의도입니다. 다만 biased mode에서는 공정성을 호출자가 책임집니다. 앞 branch가 계속 ready인데 loop가 계속 돈다면 뒤 branch가 굶을 수 있습니다. 여기서는 종료를 관찰하면 loop를 끝내거나 drain mode로 전환하므로 종료 branch가 무한히 ready인 채 정상 작업을 가리는 구조는 아닙니다. 다른 loop에 그대로 복사할 때는 이 조건을 다시 검토해야 합니다.

4. drain과 cancel의 결과를 다르게 기록합니다

종료 요청은 입력 admission부터 동기적으로 닫습니다. Cancel은 receiver도 닫아 capacity를 기다리는 reservation을 깨운 뒤 active admission이 모두 success, failure 또는 future drop으로 정리될 때까지 기다립니다. 그 다음 queue를 비우므로 submitOk를 반환한 작업은 최종 report의 completed 또는 cancelled 중 하나에 반드시 들어갑니다. DrainCancel의 결과는 StopReason::CleanShutdownStopReason::Cancelled로 구분합니다.

4.1. clean drain은 현재 timer와 accepted queue를 끝냅니다

작업 중 Drain을 관찰하면 receiver를 close해 새 전송을 막고 현재 Checkdrain_queue로 넘깁니다. 함수는 현재 check의 Sleep을 새로 만들어 pin하고 완료까지 기다린 뒤 queue에 이미 들어온 다음 작업을 처리합니다. drain 요청 자체가 현재 작업의 완료를 취소하지 않습니다.

let timer = sleep(check.duration);
tokio::pin!(timer);
loop {
    if StopSignal::from_state(stop_state) == StopSignal::Cancel {
        let _ = events.send(WorkerEvent::Cancelled(check.name.clone()));
        cancelled.push(check.name);
        cancel_queue(jobs, admission, events, cancelled).await;
        return StopReason::Cancelled;
    }
    if watch_open {
        tokio::select! {
            biased;
            changed = stop.changed() => {
                let _ = stop.borrow_and_update();
                if StopSignal::from_state(stop_state) == StopSignal::Cancel {
                    let _ = events.send(WorkerEvent::Cancelled(check.name.clone()));
                    cancelled.push(check.name);
                    cancel_queue(jobs, admission, events, cancelled).await;
                    return StopReason::Cancelled;
                }
                if changed.is_err() {
                    watch_open = false;
                }
            }
            () = &mut timer => {
                let _ = events.send(WorkerEvent::Completed(check.name.clone()));
                completed.push(check.name);
                break;
            }
        }
    } else {
        (&mut timer).await;
        let _ = events.send(WorkerEvent::Completed(check.name.clone()));
        completed.push(check.name);
        break;
    }
}

여기서 Sleep을 pin한 채 같은 future를 계속 poll한다는 점이 중요합니다. select! loop를 돌 때마다 새 sleep을 만들면 deadline이 계속 뒤로 밀릴 수 있습니다. drain 도중 watch sender가 닫히면 watch_openfalse로 바꾸고 같은 pinned timer를 직접 await합니다.

4.2. cancel은 in-flight sleep을 drop하고 queue까지 분류합니다

취소를 관찰한 branch에서 scope를 빠져나오면 in-flight Sleep future가 drop됩니다. Drop하면 timer가 취소되며 추가 cleanup은 필요 없습니다. 이 말은 timer future에 한정됩니다. 임의의 future나 외부 I/O에도 같은 결론을 적용할 수는 없습니다.

현재 check는 WorkerEvent::Cancelled와 report의 cancelled에 추가합니다. cancel_queue는 receiver를 닫고 wait_until_settled로 in-flight admission이 끝나기를 기다린 뒤 try_recv로 accepted check를 모두 분류합니다. Receiver::close 전에 얻은 outstanding permit이 뒤늦게 성공할 수 있기 때문에 이 대기 순서가 필요합니다. 이미 완료한 check는 completed에 남습니다. 모든 terminal return은 같은 admission mutex 아래에서 최신 atomic state를 다시 읽고 worker_stopped를 기록합니다. 따라서 그 잠금보다 먼저 성공한 request_cancelCleanShutdown으로 보고될 수 없습니다.

5. paused time과 event로 종료 테스트를 결정적으로 만듭니다

실제 시간에 맞춰 sleep한 뒤 "충분히 기다렸을 것"이라고 가정하면 느리고 흔들리는 테스트가 됩니다. Tokio test runtime의 start_paused = truetime::advance를 쓰면 wall clock을 기다리지 않고 timer 경계를 진행시킬 수 있습니다. paused runtime은 할 일이 없을 때 다음 pending timer로 시간을 자동 전진할 수도 있으므로 독립적으로 ready인 task 사이의 세부 poll 순서를 주장해서는 안 됩니다.

5.1. 먼저 관찰하고 필요한 만큼만 시간을 전진시킵니다

clean shutdown 테스트는 current_thread runtime을 paused 상태로 시작합니다. Started("alpha") event를 받아 첫 작업이 실제로 in flight임을 확인한 뒤 drain을 요청하고 virtual time을 10초 전진합니다. 완료 event를 확인한 다음 두 번째 timer에 5초를 더 전진합니다.

monitor.request_clean_shutdown().unwrap();
advance(Duration::from_secs(10)).await;
assert_eq!(
    monitor.next_event().await,
    Some(WorkerEvent::Completed("alpha".into()))
);
assert_eq!(
    monitor.next_event().await,
    Some(WorkerEvent::Started("beta".into()))
);
advance(Duration::from_secs(5)).await;
assert_eq!(
    monitor.next_event().await,
    Some(WorkerEvent::Completed("beta".into()))
);

backpressure 테스트도 먼저 Started event를 기다립니다. 취소 테스트는 완료된 작업 하나, in-flight 작업 하나, queued 작업 하나를 구분한 뒤 취소합니다. 별도로 cancel→drain, drain→cancel, drain 관찰 후 cancel을 검사해 단조 상태와 escalation을 고정합니다.

5.2. 모든 종료 경로는 JoinHandle을 기다립니다

Monitor::join(self)은 자신이 가진 sender와 stop sender를 drop한 다음 handle.await를 호출합니다. cancel(self)도 취소 요청을 보낸 뒤 join으로 이어집니다. clean drain, explicit cancel, 모든 submitter drop 경로에서 detached task를 cleanup으로 인정하지 않습니다.

WorkerReport::joined_count는 이 예제에서 1입니다. worker task 하나를 실제로 join했다는 계약을 출력과 테스트가 함께 확인합니다. JoinHandle await가 실패하면 JoinWorkerError::TaskFailed로 반환합니다.

6. 실행 결과와 적용 범위를 함께 확인합니다

프로젝트 디렉터리에서 아래 명령을 실행합니다.

cd examples/article-25-tokio-channels-select-cancellation
cargo fmt --all -- --check
cargo check --all-targets --all-features
cargo clippy --all-targets --all-features -- -D warnings
cargo test --all-features
cargo run --quiet

Rust 1.98.1, Cargo 1.98.1, Tokio 1.53.1에서 다섯 명령은 종료 코드 0을 반환해야 합니다. 테스트 모음에는 통합 test 14개와 unit test 1개가 있습니다. 실행 출력은 다음과 같습니다.

stop=clean
completed=alpha,beta
cancelled=0
joined=1

6.1. 테스트가 고정하는 계약

15개 테스트는 capacity 경계, bounded queue의 Full, clean drain, cancel/drain의 단조 전이, terminal race, stop 이후 admission 거부, concurrent accepted-send 분류, close 전 outstanding permit, submitter cleanup, 정확한 binary 출력을 다룹니다. multi-thread terminal stress는 보조 테스트이고, permit test와 admission-close test가 결정적 회귀 테스트입니다. 명령 순서에는 format, check, clippy도 포함됩니다.

출력의 stop=clean은 이 binary run이 clean drain 경로를 택했다는 뜻입니다. 모든 실행에서 항상 clean하다는 주장은 아닙니다. cancellation 테스트에서는 StopReason::Cancelled와 취소된 작업 이름을 따로 검증합니다.

6.2. 운영 코드로 옮길 때 정할 것

이 설계를 적용하려면 먼저 accepted의 경계를 정해야 합니다. 이 예제에서는 submit 성공 시점부터 queue가 작업을 소유합니다. 이어서 drain이 accepted 작업을 모두 마쳐야 하는지, cancel이 in-flight operation을 drop해도 되는지, shutdown 우선순위가 처리량보다 앞서는지 결정합니다. 각 branch API의 취소 안전성도 문서에서 확인해야 합니다.

여기서 확인한 것은 timer로 표현한 검사 작업 하나입니다. 실제 database write나 socket protocol은 future를 drop했을 때 외부 효과가 어디까지 진행됐는지 별도 계약이 필요합니다. select! 자체를 안전 보증으로 삼지 말고 recvchanged처럼 문서화된 operation 단위로 판단하십시오. 마지막 cleanup 경계에서는 worker를 join하고 report를 확인해야 합니다.

전체 소스 코드

이 글의 전체 실행 가능한 소스는 GitHub의 Chapter 25 프로젝트에서 확인할 수 있습니다.

출처


2개 응답

  1. […] 이전 글Tokio 채널·select·취소로 종료 가능한 작업자 만들기 […]

  2. […] 다음 글Tokio 채널·select·취소로 종료 가능한 작업자 만들기 […]

답글 남기기

이메일 주소는 공개되지 않습니다. 필수 필드는 *로 표시됩니다

Tech Wiki

Built with WordPress · Learn in public.