使用Rust通过gRPC调用GCS的WriteObjectRequest时遇参数错误
我用gRPC和生成的proto文件向Google Cloud Storage执行简单上传,调用WriteObjectRequest时遇到报错:
An x-goog-request-params request metadata property must be provided for this request.
奇怪的是,请求元数据里明明已经包含x-goog-request-params,但服务端检测不到;而用相同方式添加该参数的ReadObjectRequest能正常识别,下载验证也没问题。以下是相关Rust代码及调试输出,求可行的解决方法:
use gcloud_sdk::google::storage::v2::{storage_client::StorageClient, write_object_request, ChecksummedData, Object, ReadObjectRequest, WriteObjectRequest, WriteObjectSpec}; use tonic::{metadata::MetadataValue, transport::Channel, Request}; use futures::stream; use hyper_rustls; use tokio_stream::StreamExt; use tonic::transport::ClientTlsConfig; const BUCKET_NAME: &str = "test-bucket"; const GOOGLE_AUTH_TOKEN: &str = "token-xxxx"; struct GcsClient { client: StorageClient<Channel>, bucket: String, } impl GcsClient { async fn new(bucket: String) -> Result<Self, Box<dyn std::error::Error>> { let channel = Channel::from_static("https://storage.googleapis.com") .connect_timeout(std::time::Duration::from_secs(5)) .timeout(std::time::Duration::from_secs(30)) .tcp_nodelay(true) .http2_adaptive_window(true) .http2_keep_alive_interval(std::time::Duration::from_secs(30)) .tls_config(ClientTlsConfig::new().with_native_roots())? .connect() .await?; Ok(Self { client: StorageClient::new(channel), bucket, }) } fn get_formatted_bucket(&self) -> String { format!("projects/_/buckets/{}", self.bucket) } async fn get_token() -> Result<String, Box<dyn std::error::Error>> { Ok(GOOGLE_AUTH_TOKEN.to_string()) } async fn auth_request<T>(&self, request: T) -> Request<T> { let token = Self::get_token().await.expect("Failed to get authentication token"); let mut request = Request::new(request); // Authorization header request.metadata_mut().insert( "authorization", MetadataValue::try_from(&format!("Bearer {}", token)).unwrap(), ); let formatted_bucket = self.get_formatted_bucket(); let encoded_bucket = urlencoding::encode(&formatted_bucket); // Adding required x-goog-request-params based on request type let params = if std::any::type_name::<T>().contains("WriteObjectRequest") { format!("write_object_spec.resource.bucket={}", encoded_bucket) } else if std::any::type_name::<T>().contains("ReadObjectRequest") { format!("bucket={}", encoded_bucket) } else if std::any::type_name::<T>().contains("StartResumableWriteRequest") { format!("write_object_spec.resource.bucket={}", encoded_bucket) } else if std::any::type_name::<T>().contains("QueryWriteStatusRequest") { format!("upload_id={}", encoded_bucket) } else { format!("bucket={}", encoded_bucket) }; request.metadata_mut().insert( "x-goog-request-params", MetadataValue::try_from(¶ms).unwrap(), ); request } async fn simple_upload(&mut self, object_name: &str, data: Vec<u8>) -> Result<(), Box<dyn std::error::Error>> { // Create write specification let write_spec = WriteObjectSpec { resource: Some(Object { name: object_name.to_string(), bucket: self.get_formatted_bucket(), ..Default::default() }), ..Default::default() }; // Create a single request with both spec and data let write_request = WriteObjectRequest { first_message: Some(write_object_request::FirstMessage::WriteObjectSpec(write_spec)), write_offset: 0, data: Some(write_object_request::Data::ChecksummedData(ChecksummedData { content: data, crc32c: None, })), finish_write: true, ..Default::default() }; // Create authenticated request and convert to stream let auth_request = self.auth_request(write_request).await; // Print the request metadata for debugging println!("Request metadata: {:#?}", auth_request.metadata()); // Output from above: // Request metadata: MetadataMap { // headers: { // "authorization": "Bearer token-xxxx", // "x-goog-request-params": "write_object_spec.resource.bucket=projects%2F_%2Fbuckets%2Ftest-bucket", // }, // } // Send the request match self.client.write_object(stream::iter(vec![auth_request.into_inner()])).await { Ok(response) => { println!("Upload successful! Response: {:?}", response); Ok(()) } Err(e) => { println!("Error details: {:#?}", e); // Output from above: // Error details: Status { // code: InvalidArgument, // message: "An x-goog-request-params request metadata property must be provided for this request.", // details: b"\x08\x03\x12UAn x-goog-request-params request metadata property must be provided for this request.\x1ah\n(type.googleapis.com/google.rpc.ErrorInfo\x12<\n\"GRPC_INVALID_X_GOOG_REQUEST_PARAMS\x12\x16storage.googleapis.com", // metadata: MetadataMap { // headers: { // "grpc-server-stats-bin": "AAAzYToAAAAAAA", // "google.rpc.errorinfo-bin": "CiJHUlBDX0lOVkFMSURfWF9HT09HX1JFUVVFU1RfUEFSQU1TEhZzdG9yYWdlLmdvb2dsZWFwaXMuY29t", // "endpoint-load-metrics-bin": "MbFRfe/oHjJASbtUC3gdlq0/", // "content-type": "application/grpc", // "grpc-accept-encoding": "identity, deflate, gzip", // "content-length": "0", // "date": "Tue, 07 Jan 2025 21:58:42 GMT", // "alt-svc": "h3=\":443\"; ma=2592000,h3-29=\":443\"; ma=2592000", // }, // }, // source: None, // } Err(Box::new(e)) } } } async fn download_object(&mut self, object_name: &str) -> Result<Vec<u8>, Box<dyn std::error::Error>> { let request = ReadObjectRequest { bucket: self.get_formatted_bucket(), object: object_name.to_string(), ..Default::default() }; let request = self.auth_request(request).await; println!("Request metadata: {:#?}", request.metadata()); // Output from above: // Request metadata: MetadataMap { // headers: { // "authorization": "Bearer token-xxxx", // "x-goog-request-params": "bucket=projects%2F_%2Fbuckets%2Ftest-bucket", // }, // } let response = self.client.read_object(request).await?; let mut stream = response.into_inner(); let mut content = Vec::new(); while let Some(chunk) = StreamExt::next(&mut stream).await { let chunk = chunk?; if let Some(data) = chunk.checksummed_data { content.extend(data.content); } } Ok(content) } } #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { let bucket = BUCKET_NAME.to_string(); let mut client = GcsClient::new(bucket).await?; // Test small file upload println!("\n=== Testing Small File Upload ==="); let small_content = b"Hello, GCS! This is a test.".to_vec(); let small_object_name = "test-small-file.txt"; println!("Uploading small file: {}", small_object_name); client.simple_upload(small_object_name, small_content.clone()).await?; println!("Small file upload complete!"); // The above output is failing with the error specified in the function. // Download and verify small file println!("\n=== Downloading Small File ==="); let downloaded = client.download_object(small_object_name).await?; println!("Downloaded content: {}", String::from_utf8_lossy(&downloaded)); println!("Content verification: {}", downloaded == small_content); // Output from above: // Downloaded content: Hello, GCS! This is a test. // Content verification: true Ok(()) }
可行的解决尝试方向
修复流式请求的元数据传递逻辑:
WriteObject是流式gRPC接口,你当前的代码把带元数据的Request<WriteObjectRequest>调用into_inner()取出请求体再包装成流,这会导致元数据丢失。正确做法是把请求流包裹在Request中,再给外层Request添加元数据:
调整simple_upload中的代码:// 先创建请求流 let request_stream = stream::iter(vec![write_request]); // 创建包含流的Request,再添加元数据 let mut auth_request = Request::new(request_stream); // 手动添加授权和x-goog-request-params元数据 let token = Self::get_token().await.expect("Failed to get authentication token"); auth_request.metadata_mut().insert( "authorization", MetadataValue::try_from(&format!("Bearer {}", token)).unwrap(), ); let formatted_bucket = self.get_formatted_bucket(); let encoded_bucket = urlencoding::encode(&formatted_bucket); let params = format!("bucket={}", encoded_bucket); auth_request.metadata_mut().insert( "x-goog-request-params", MetadataValue::try_from(¶ms).unwrap(), ); // 发送请求 match self.client.write_object(auth_request).await { // 后续逻辑不变 }修正x-goog-request-params的参数名:当前给
WriteObjectRequest设置的参数是write_object_spec.resource.bucket=xxx,但GCS gRPC API可能要求该接口的参数名就是bucket(和ReadObject一致)。尝试把Write对应的参数改成format!("bucket={}", encoded_bucket)再测试。检查元数据键名的正确性:确认
x-goog-request-params的键名没有拼写错误,虽然HTTP头大小写不敏感,但tonic的Metadata处理可能严格区分大小写,确保键名完全匹配。开启gRPC日志排查实际请求头:添加环境变量
RUST_LOG=tonic::transport=debug后运行程序,查看实际发送的HTTP请求头,确认x-goog-request-params是否真的被发送到服务端。
内容的提问来源于stack exchange,提问作者Aman Kumar

