chore: rebrand project to crank

This commit is contained in:
a.tolmachev
2026-03-28 00:58:56 +03:00
parent 6821d0c64a
commit 26335e8d9b
101 changed files with 550 additions and 538 deletions
+168
View File
@@ -0,0 +1,168 @@
use std::{collections::BTreeMap, str::FromStr, time::Duration};
use base64::{Engine as _, engine::general_purpose::STANDARD};
use crank_core::GrpcTarget;
use prost::Message;
use prost_reflect::{DescriptorPool, MethodDescriptor, prost_types::FileDescriptorSet};
use tonic::{
Request,
client::Grpc,
metadata::{KeyAndValueRef, MetadataKey, MetadataValue},
transport::Endpoint,
};
use crate::{GrpcAdapterError, GrpcRequest, GrpcResponse, codec::JsonCodec};
#[derive(Clone, Debug, Default)]
pub struct GrpcAdapter;
impl GrpcAdapter {
pub fn new() -> Self {
Self
}
pub async fn execute(
&self,
target: &GrpcTarget,
request: &GrpcRequest,
) -> Result<GrpcResponse, GrpcAdapterError> {
let method = resolve_method(target)?;
if method.is_client_streaming() || method.is_server_streaming() {
return Err(GrpcAdapterError::UnsupportedMethodKind {
service: format!("{}.{}", target.package, target.service),
method: target.method.clone(),
});
}
let endpoint = Endpoint::from_shared(target.server_addr.clone())?
.timeout(Duration::from_millis(request.timeout_ms));
let channel = endpoint.connect().await?;
let mut grpc = Grpc::new(channel);
grpc.ready()
.await
.map_err(|error| GrpcAdapterError::Status {
code: tonic::Code::Unavailable,
message: error.to_string(),
})?;
let service_name = method.parent_service().full_name().to_owned();
let method_name = method.name().to_owned();
let path = tonic::codegen::http::uri::PathAndQuery::from_str(&format!(
"/{service_name}/{method_name}"
))
.map_err(|_| GrpcAdapterError::InvalidMethodPath {
service: service_name,
method: method_name,
})?;
let codec = JsonCodec::new(method.input(), method.output());
let mut tonic_request = Request::new(request.body.clone());
apply_headers(&mut tonic_request, &request.headers)?;
let response = grpc.unary(tonic_request, path, codec).await?;
let headers = normalize_headers(response.metadata());
let body = response.into_inner();
Ok(GrpcResponse {
status_code: 200,
headers,
body,
})
}
}
fn resolve_method(target: &GrpcTarget) -> Result<MethodDescriptor, GrpcAdapterError> {
let bytes = STANDARD
.decode(&target.descriptor_set_b64)
.map_err(|_| GrpcAdapterError::InvalidDescriptorEncoding)?;
let descriptor_set = FileDescriptorSet::decode(bytes.as_slice())
.map_err(|_| GrpcAdapterError::InvalidDescriptorSet)?;
let pool = DescriptorPool::from_file_descriptor_set(descriptor_set)
.map_err(|_| GrpcAdapterError::InvalidDescriptorPool)?;
let full_service_name = if target.package.is_empty() {
target.service.clone()
} else {
format!("{}.{}", target.package, target.service)
};
let service = pool
.get_service_by_name(&full_service_name)
.ok_or_else(|| GrpcAdapterError::ServiceNotFound {
service: full_service_name.clone(),
})?;
service
.methods()
.find(|method| method.name() == target.method)
.ok_or_else(|| GrpcAdapterError::MethodNotFound {
service: full_service_name,
method: target.method.clone(),
})
}
fn apply_headers(
request: &mut Request<serde_json::Value>,
headers: &BTreeMap<String, String>,
) -> Result<(), GrpcAdapterError> {
for (key, value) in headers {
let metadata_key =
MetadataKey::from_str(key).map_err(|source| GrpcAdapterError::InvalidMetadataKey {
key: key.clone(),
source,
})?;
let metadata_value = MetadataValue::from_str(value).map_err(|source| {
GrpcAdapterError::InvalidMetadataValue {
key: key.clone(),
source,
}
})?;
request.metadata_mut().insert(metadata_key, metadata_value);
}
Ok(())
}
fn normalize_headers(metadata: &tonic::metadata::MetadataMap) -> BTreeMap<String, String> {
metadata
.iter()
.filter_map(|entry| match entry {
KeyAndValueRef::Ascii(key, value) => {
Some((key.as_str().to_owned(), value.to_str().ok()?.to_owned()))
}
KeyAndValueRef::Binary(_, _) => None,
})
.collect()
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use crank_core::{DescriptorId, GrpcTarget};
use serde_json::json;
use crate::{GrpcAdapter, GrpcRequest, test_support};
#[tokio::test]
async fn executes_unary_grpc_request() {
let server_addr = test_support::spawn_unary_echo_server().await;
let adapter = GrpcAdapter::new();
let target = GrpcTarget {
server_addr,
package: "echo".to_owned(),
service: "EchoService".to_owned(),
method: "UnaryEcho".to_owned(),
descriptor_ref: DescriptorId::new("desc_echo"),
descriptor_set_b64: test_support::descriptor_set_b64(),
};
let request = GrpcRequest {
headers: BTreeMap::new(),
body: json!({ "message": "hello" }),
timeout_ms: 1_000,
};
let response = adapter.execute(&target, &request).await.unwrap();
assert_eq!(response.body, json!({ "message": "hello" }));
}
}