태그 보관물: rust

gRPC Streaming으로 상태 전송 받기 예제

gRPC에서 Observer pattern을 구현하는 예제를 고민했던 적이 있었는데(gRPC를 이용한 Observer Pattern) 알고보니 gRPC에서는 ReadReactor / WriteReactor 보다 간단하게 object의 상태 변경을 감지할 수 있는 메커니즘을 제공하고 있었다. gRPC를 이용한 Observer pattern 구현 예제를 찾기 어려웠던 이유가 어쩌면 그런 복잡한 패턴 없이도 구현이 가능하기 때문이었는지도 모르겠다.

이 포스팅에서는 Stream을 통해 변경되는 값들을 subscribe하는 간단한 예제를 Rust로 구현해 본다. gRPC의 Rust 구현체인 tonic과 proto compiler인 prost, 비동기 처리를 위해 tokio를 사용했다.

Proto file – stockprice.proto

랜덤하게 생성된 주가 정보를 gRPC로 구독하는 간단한 서비스 예제를 가정해 보자. Proto file은 UpdateStockPrice()를 정의하고 호출되면 StockPriceResponse의 stream을 받아오도록 다음과 같이 구성한다.

syntax = "proto3";
package stockservice;

import "google/protobuf/empty.proto";

message StockPriceResponse {
    string symbol = 1;
    double price = 2;
    uint64 sequence = 3;
}

service StockService {
    rpc UpdateStockPrice(google.protobuf.Empty) returns (stream StockPriceResponse);
}

Server

StockServiceImpl 구조체를 선언하고 서비스를 제공하기 위한 StockService를 다음과 같이 구현한다.

// Server 구현체
struct StockServiceImpl;

// tokio_stream을 이용해서 gRPC stream을 구현.
type StockStream = Pin<Box<dyn Stream<Item = Result<StockPriceResponse, Status>> + Send>>;

#[tonic::async_trait]
impl StockService for StockServiceImpl {
    type UpdateStockPriceStream = StockStream;

    async fn update_stock_price(
        &self,
        _request: Request<()>,
    ) -> Result<Response<Self::UpdateStockPriceStream>, Status> {

        // 1초 간격으로 stock price 생성
        let mut interval = tokio_stream::wrappers::IntervalStream::new(tokio::time::interval(
            Duration::from_secs(1),
        ));

        // Stream으로 전송할 데이터를 생성하고 yield로 반환.
        let output = async_stream::stream! {
            let mut rng = SmallRng::from_os_rng();
            let mut sequence: u64 = 0;

            // 매 interval마다 순번을 증가시키고 랜덤 주가를 생성해 클라이언트로 전송
            while interval.next().await.is_some() {
                sequence += 1;
                let p = generate_stock_price(&mut rng, sequence);
                yield Ok(p);
            }
        };

        Ok(Response::new(
            Box::pin(output) as Self::UpdateStockPriceStream
        ))
    }
}

이 코드에서 핵심은 async_stream::stream!() 매크로로 데이터를 전송할 데이터임을 정의하고, while / yeild로 data를 전송하는 다음의 부분이다. 이 부분에서 데이터를 갱신하고 stream에 지속적으로 update 시킨다.

        let output = async_stream::stream! {
           ...

            // 매 interval마다 순번을 증가시키고 랜덤 주가를 생성해 클라이언트로 전송
            while interval.next().await.is_some() {
                // Data 생성 및 전송
                ...
                yield Ok(p);
            }
        };

Client

Server에 비해 Client 쪽은 비교적 구현이 간단한데, prost가 proto file로 부터 생성해 준 StockServiceClient를 이용해서 stream을 생성하고, 이 때 반환받은 stream으로 부터 계속 데이터를 읽어 주면된다.

#[tokio::main]
async fn main() -> Result<(), Box<dyn Error>> {
    let mut client = StockServiceClient::connect("http://[::1]:50051").await?;

    // grpc 서비스 stream 생성.
    let mut stream = client
        .update_stock_price(Request::new(()))
        .await?
        .into_inner();

    // stream에서 데이터를 받아서 화면에 출력.
    while let Some(update) = stream.message().await? {
        println!(
            "[{}] {}: {:.2}",
            update.sequence, update.symbol, update.price
        );
    }

    Ok(())
}

코드 위치 및 수행결과

전체 구현 코드는 다음 위치에 있다. https://github.com/litcoder/grpcobsr/tree/stream_rust

Rust와 Python unit test의 공존

예전 글에서 vscode의 test explorer에 Rust의 unittest가 보이도록 설정하는 방법을 다룬 적이 있었는데, 여기에 Python test case (여기서는 pytest)도 함께 표시되도록 하려면 .vscode/settings.json을 편집해 주어야 한다.

해당 프로젝트의 경우 가상환경을 source root에 두지 않고 서브 디렉토리인 service/.venv 안에 넣어 두었기 때문에 python.defaultInterpreterPath 값을 이곳으로 직접 설정해 주고 pytest 사용을 위한 설정도 해 두었다.

{
    // Rust unit test
    "rust-analyzer.testExplorer": true,

    // Python virtual environment용 인터프리터 설정
    "python.defaultInterpreterPath": "${workspaceFolder}/service/.venv/bin/python", 
    
    // pytest 사용
    "python.testing.pytestEnabled": true,
    "python.testing.unittestEnabled": false,

    // Python unit test가 있는 디렉토리 경로
    "python.testing.pytestArgs": [
        "service/test"
    ],
    "python-envs.defaultEnvManager": "ms-python.python:venv"
}

그리고 나서 vscode GUI에서도 다시 한번 venv를 설정해 준다. 그냥 프로그램을 재 실행 했으면 이 부분은 건너 뛰어도 되었을 것 같긴 한데, 오류가 계속 뜨길래 수동으로 설정해 주었다.

이렇게 하고 나면 test explorer에 두 언어의 test case들이 모두 표시되는 평화로운 공존상태가 된다.

Rust 프로그래밍, 코드 가독성을 높이기 위한 and_then() 활용

Rust는 선언적 프로그래밍을 적극적으로 사용하기에 method chaining을 많이 사용하는데 이는 가독성을 높이고 컴파일러가 최적화를 수행하는데 적합한 형태이기 때문이다. 이 포스팅에서는 method chaining을 유지하면서도 중첩된 Result 혹은 Option 형식을 피할 수 있는 and_then()에 대해 알아본다.

match 문을 이용한 처리

세개의 함수 각각이 다음과 같이 Result를 반환한다고 해보자.

use std::io::{Error, ErrorKind};

/// 1. 입력된 사용자 JSON 검증, 빈 문자열이 아니면 Ok 반환.
fn validate_json(raw: &str) -> Result<&str, Error> {
  if raw.is_empty() {
    Err(Error::new(ErrorKind::InvalidData, "empty user input"))
  } else {
    Ok(raw)
  }
}

/// 2. 사용자 id 반환
fn get_user_id(username: &str) -> Result<u32, Error> {
  if username == "admin" {
    Ok(1)
  } else {
    Err(Error::new(ErrorKind::NotFound, "user not found"))
  }
}

/// 3. 토큰 생성
fn generate_token(user_id: u32) -> Result<String, Error> {
  if user_id == 1 {
    Ok(format!("token_for_user_{}", user_id))
  } else {
    Err(Error::new(ErrorKind::InvalidInput, "user id unavailable"))
  }
}

각 함수들은 1. 유효한 JSON인지 판별하고, 2. 사용자를 검색해서 id를 반환한 후, 3. 해당 사용자에 대한 토큰을 생성해서 반환한다.

이를 match 문을 이용해서 처리 하면 다음과 같이 될 것이다.

fn main() {
  let raw_input = "admin";

  // 1. JSON 유효성 판별
  match validate_json(raw_input) {
    Ok(user) => {
      // 2. 유효하면 사용자 id 검색
      match get_user_id(user) {
        Ok(id) => {
          // 3. 사용자가 유효하면 토큰 생성
          match generate_token(id) {
            Ok(tok) => {
              // 토큰 출력
              println!("{:?}", tok);
            },
            // 토큰 생성시 에러 처리
            Err(e) => eprint!("Error: {e}")
          }
        },
        // 유효하지 않은 사용자
        Err(e) => eprint!("Error: {e}")
      }
    },
    // JSON이 유효하지 않음
    Err(e) => eprint!("Error: {e}")
    }
  }
}

모든 match arm들을 명시해야 하는 match 구문의 제약 때문에 ResultOkErr에 대한 경우를 모두 적어 주어야 하고 이 때문에 코드가 장황하고 가독성이 떨어진다.

map() 활용

이전의 “Rust 프로그래밍에서 map() 활용” 에서 언급 했듯이 map()은 container뿐만 아니라 ResultOptionOk / Some인 경우 동작을 정의하는데 사용할 수 있다. 따라서 map()을 이용하면 다음과 같이 작성할 수 있다.

fn main() {
  let raw_input = "admin";

  let tok = validate_json(raw_input)
      .map(|user| get_user_id(user)
      .map(|id| generate_token(id)));

  // Ok(Ok(Ok("token_for_user_1")))
  print!("{:?}", tok);
}

보기에는 훨씬 깔끔해 졌지만, 각 함수들이 Result를 반환하기 때문에 tok변수의 type이 Result<Result<Result<String, ...>, ...>, ...> 같은 여러번 중첩된 형식이 되어 여러번 unwrap 해주어야 하는 보기 싫은 형태가 된다.

and_then() 활용

이렇게 method chaining으로 중첩된 ResultOption이 반환되는 경우에는 and_then()을 사용해보자.

fn main() {
  let raw_input = "admin";

  let tok = validate_json(raw_input)
      .and_then(|user| get_user_id(user))
      .and_then(|id| generate_token(id));
 
  // Ok("token_for_user_1")  
  print!("{:?}", tok);
}

이렇게 method chaining을 하면 중복된 반환형식 들이 모두 제거되고 Result<String, Error> 형식으로 tok 변수가 반환되어 훨씬 가독성 높은 코드를 유지할 수 있다.

결론

map()and_then()은 모두 Rust의 선언형 프로그래밍을 가능하게 하는 중요한 요소로 비슷해 보이지만, closure가 무엇을 반환하는가에 따라 그 활용이 달라진다.

  • map(): 내부 값의 형태만 바꿀때(실패 가능성이 없는 단순 변환)
  • and_then(): 변환 과정에서 또 다른 Option이나 Result가 발생할 때(실패 가능성이 있는 연쇄 작업)

동일한 코드를 matchif let으로 작성할 수도 있지만, 이른바 “match hell”의 조짐이 보이고 그로 인해 코드의 가독성이 떨어진다면 and_then()을 활용해 보자.