Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
37 commits
Select commit Hold shift + click to select a range
27c0701
Migrate storage gRPC error codes to tonic::Status
sitaowang1998 Jun 26, 2026
00c3196
Implement the missing storage gRPC services
sitaowang1998 Jun 26, 2026
b41f226
Refactor ready-task builder to take a task-id closure instead of a trait
sitaowang1998 Jun 26, 2026
1e35e9c
Add ResendReadyTasks and DeleteResourceGroup services with job-gone h…
sitaowang1998 Jun 26, 2026
a1ee868
Merge branch 'main' into storage-grpc-migration
sitaowang1998 Jun 26, 2026
b4fece6
Merge storage-grpc-migration (with main #342) into storage-grpc-services
sitaowang1998 Jun 26, 2026
9c601f1
Merge branch 'storage-grpc-services' into add-missing-services
sitaowang1998 Jun 26, 2026
b98b085
Rename map_inbound_status to inbound_status_to_error to match main's …
sitaowang1998 Jun 26, 2026
b241c34
Rename map_inbound_status and map_liveness_status to match main's sta…
sitaowang1998 Jun 26, 2026
ff31980
Merge storage-grpc-migration (status_to_error rename) into storage-gr…
sitaowang1998 Jun 26, 2026
22dbe6a
Merge storage-grpc-services (status_to_error rename) into add-missing…
sitaowang1998 Jun 26, 2026
429fdbb
Merge branch 'storage-grpc-migration' of github.com:sitaowang1998/spi…
sitaowang1998 Jun 26, 2026
25fc42d
Merge origin/main into storage-grpc-migration
sitaowang1998 Jun 27, 2026
72b1021
Merge storage-grpc-migration (with main #358) into storage-grpc-services
sitaowang1998 Jun 27, 2026
f4790c9
Merge storage-grpc-services (with main #358/#360) into add-missing-se…
sitaowang1998 Jun 27, 2026
c84d324
Update proto.
LinZhihao-723 Jun 28, 2026
9c2f88f
Merge origin/main (tonic 0.14 bump) into storage-grpc-migration
sitaowang1998 Jun 28, 2026
2cf1432
Merge storage-grpc-migration (tonic 0.14 + common.Void responses) int…
sitaowang1998 Jun 28, 2026
7849212
Merge storage-grpc-services (tonic 0.14 + common.Void responses) into…
sitaowang1998 Jun 28, 2026
896eef6
Remove session id 0 check
sitaowang1998 Jun 28, 2026
3c0d0e8
Use void for empty resposne
sitaowang1998 Jun 28, 2026
c137a3b
Map inbound queue closed to internal error
sitaowang1998 Jun 28, 2026
0bba7ae
Merge branch 'storage-grpc-migration' into storage-grpc-services
sitaowang1998 Jun 28, 2026
2aa8d8e
Map inbound queue closed to internal error on the server
sitaowang1998 Jun 28, 2026
d809a62
Merge branch 'storage-grpc-services' into add-missing-services
sitaowang1998 Jun 28, 2026
3e8b1d7
Update resend_ready_tasks doc to drop InboundClosed reference
sitaowang1998 Jun 28, 2026
925219a
Merge origin/main into storage-grpc-services
sitaowang1998 Jun 29, 2026
9b7f0b6
Merge branch 'storage-grpc-services' into add-missing-services
sitaowang1998 Jun 29, 2026
3ec8e32
Fix build_ready_tasks docstring to reference actual lane marker types
sitaowang1998 Jun 29, 2026
c998dac
Merge storage-grpc-services into add-missing-services
sitaowang1998 Jun 29, 2026
a0c8513
Merge branch 'main' into storage-grpc-services
sitaowang1998 Jun 29, 2026
dbf7c7d
Merge branch 'main' into storage-grpc-services
sitaowang1998 Jun 30, 2026
ad69055
Merge branch 'main' into add-missing-services
sitaowang1998 Jun 30, 2026
2fee410
Merge branch 'main' into storage-grpc-services
sitaowang1998 Jul 1, 2026
bf84e09
Merge branch 'main' into storage-grpc-services
sitaowang1998 Jul 2, 2026
e612765
Merge remote-tracking branch 'origin/storage-grpc-services' into add-…
sitaowang1998 Jul 2, 2026
e4944ab
Add DeleteResourceGroup in client
sitaowang1998 Jul 2, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions components/spider-client/src/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -224,4 +224,29 @@ impl SpiderClient {
.verify_resource_group(resource_group_id, password)
.await
}

/// Deletes a resource group after verifying its password.
///
/// # Returns
///
/// `Ok(())` on success.
///
/// # Errors
///
/// Returns an error if:
///
/// * [`ClientError::InvalidArgument`] if the storage server rejects the request as invalid.
/// * [`ClientError::Unauthenticated`] if the resource group is unknown or the password is
/// invalid.
/// * [`ClientError::Transport`] if the gRPC transport fails or the connection is lost.
/// * [`ClientError::Server`] for any other server-reported error.
pub async fn delete_resource_group(
&self,
resource_group_id: ResourceGroupId,
password: Vec<u8>,
) -> Result<(), ClientError> {
self.resource_group
.delete_resource_group(resource_group_id, password)
.await
}
}
29 changes: 29 additions & 0 deletions components/spider-client/src/grpc/resource_group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,35 @@ impl ResourceGroupManagementClient {
.map_err(|status| resource_group_status_to_error(&status))?;
Ok(())
}

/// Deletes a resource group after verifying its password.
///
/// # Returns
///
/// `Ok(())` on success — the storage server's response is empty, so success is implicit.
///
/// # Errors
///
/// Returns an error if:
///
/// * Forwards [`ResourceGroupManagementServiceClient::delete_resource_group`]'s status on
/// failure.
pub async fn delete_resource_group(
&self,
resource_group_id: ResourceGroupId,
password: Vec<u8>,
) -> Result<(), ClientError> {
let request = storage::DeleteResourceGroupRequest {
resource_group_id: resource_group_id.get(),
password,
};
self.connection_pool
.get_client()
.delete_resource_group(request)
.await
.map_err(|status| resource_group_status_to_error(&status))?;
Ok(())
}
}

/// Maps a resource-group-management gRPC [`Status`] to a [`ClientError`].
Expand Down
176 changes: 176 additions & 0 deletions components/spider-proto-rust/src/generated/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,13 @@ pub struct VerifyResourceGroupRequest {
pub password: ::prost::alloc::vec::Vec<u8>,
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct DeleteResourceGroupRequest {
#[prost(uint64, tag = "1")]
pub resource_group_id: u64,
#[prost(bytes = "vec", tag = "2")]
pub password: ::prost::alloc::vec::Vec<u8>,
}
#[derive(Clone, PartialEq, Eq, Hash, ::prost::Message)]
pub struct RegisterExecutionManagerRequest {
#[prost(string, tag = "1")]
pub ip_address: ::prost::alloc::string::String,
Expand Down Expand Up @@ -1635,6 +1642,32 @@ pub mod inbound_queue_service_client {
);
self.inner.unary(req, path, codec).await
}
pub async fn resend_ready_tasks(
&mut self,
request: impl tonic::IntoRequest<super::super::common::Void>,
) -> std::result::Result<
tonic::Response<super::super::common::Void>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic_prost::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/storage.InboundQueueService/ResendReadyTasks",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new("storage.InboundQueueService", "ResendReadyTasks"),
);
self.inner.unary(req, path, codec).await
}
}
}
/// Generated server implementations.
Expand Down Expand Up @@ -1671,6 +1704,13 @@ pub mod inbound_queue_service_server {
tonic::Response<super::PollReadyTasksResponse>,
tonic::Status,
>;
async fn resend_ready_tasks(
&self,
request: tonic::Request<super::super::common::Void>,
) -> std::result::Result<
tonic::Response<super::super::common::Void>,
tonic::Status,
>;
}
#[derive(Debug)]
pub struct InboundQueueServiceServer<T> {
Expand Down Expand Up @@ -1895,6 +1935,55 @@ pub mod inbound_queue_service_server {
};
Box::pin(fut)
}
"/storage.InboundQueueService/ResendReadyTasks" => {
#[allow(non_camel_case_types)]
struct ResendReadyTasksSvc<T: InboundQueueService>(pub Arc<T>);
impl<
T: InboundQueueService,
> tonic::server::UnaryService<super::super::common::Void>
for ResendReadyTasksSvc<T> {
type Response = super::super::common::Void;
type Future = BoxFuture<
tonic::Response<Self::Response>,
tonic::Status,
>;
fn call(
&mut self,
request: tonic::Request<super::super::common::Void>,
) -> Self::Future {
let inner = Arc::clone(&self.0);
let fut = async move {
<T as InboundQueueService>::resend_ready_tasks(
&inner,
request,
)
.await
};
Box::pin(fut)
}
}
let accept_compression_encodings = self.accept_compression_encodings;
let send_compression_encodings = self.send_compression_encodings;
let max_decoding_message_size = self.max_decoding_message_size;
let max_encoding_message_size = self.max_encoding_message_size;
let inner = self.inner.clone();
let fut = async move {
let method = ResendReadyTasksSvc(inner);
let codec = tonic_prost::ProstCodec::default();
let mut grpc = tonic::server::Grpc::new(codec)
.apply_compression_config(
accept_compression_encodings,
send_compression_encodings,
)
.apply_max_message_size_config(
max_decoding_message_size,
max_encoding_message_size,
);
let res = grpc.unary(method, req).await;
Ok(res)
};
Box::pin(fut)
}
_ => {
Box::pin(async move {
let mut response = http::Response::new(
Expand Down Expand Up @@ -2086,6 +2175,35 @@ pub mod resource_group_management_service_client {
);
self.inner.unary(req, path, codec).await
}
pub async fn delete_resource_group(
&mut self,
request: impl tonic::IntoRequest<super::DeleteResourceGroupRequest>,
) -> std::result::Result<
tonic::Response<super::super::common::Void>,
tonic::Status,
> {
self.inner
.ready()
.await
.map_err(|e| {
tonic::Status::unknown(
format!("Service was not ready: {}", e.into()),
)
})?;
let codec = tonic_prost::ProstCodec::default();
let path = http::uri::PathAndQuery::from_static(
"/storage.ResourceGroupManagementService/DeleteResourceGroup",
);
let mut req = request.into_request();
req.extensions_mut()
.insert(
GrpcMethod::new(
"storage.ResourceGroupManagementService",
"DeleteResourceGroup",
),
);
self.inner.unary(req, path, codec).await
}
}
}
/// Generated server implementations.
Expand Down Expand Up @@ -2115,6 +2233,13 @@ pub mod resource_group_management_service_server {
tonic::Response<super::super::common::Void>,
tonic::Status,
>;
async fn delete_resource_group(
&self,
request: tonic::Request<super::DeleteResourceGroupRequest>,
) -> std::result::Result<
tonic::Response<super::super::common::Void>,
tonic::Status,
>;
}
#[derive(Debug)]
pub struct ResourceGroupManagementServiceServer<T> {
Expand Down Expand Up @@ -2295,6 +2420,57 @@ pub mod resource_group_management_service_server {
};
Box::pin(fut)
}
"/storage.ResourceGroupManagementService/DeleteResourceGroup" => {
#[allow(non_camel_case_types)]
struct DeleteResourceGroupSvc<T: ResourceGroupManagementService>(
pub Arc<T>,
);
impl<
T: ResourceGroupManagementService,
> tonic::server::UnaryService<super::DeleteResourceGroupRequest>
for DeleteResourceGroupSvc<T> {
type Response = super::super::common::Void;
type Future = BoxFuture<
tonic::Response<Self::Response>,
tonic::Status,
>;
fn call(
&mut self,
request: tonic::Request<super::DeleteResourceGroupRequest>,
) -> Self::Future {
let inner = Arc::clone(&self.0);
let fut = async move {
<T as ResourceGroupManagementService>::delete_resource_group(
&inner,
request,
)
.await
};
Box::pin(fut)
}
}
let accept_compression_encodings = self.accept_compression_encodings;
let send_compression_encodings = self.send_compression_encodings;
let max_decoding_message_size = self.max_decoding_message_size;
let max_encoding_message_size = self.max_encoding_message_size;
let inner = self.inner.clone();
let fut = async move {
let method = DeleteResourceGroupSvc(inner);
let codec = tonic_prost::ProstCodec::default();
let mut grpc = tonic::server::Grpc::new(codec)
.apply_compression_config(
accept_compression_encodings,
send_compression_encodings,
)
.apply_max_message_size_config(
max_decoding_message_size,
max_encoding_message_size,
);
let res = grpc.unary(method, req).await;
Ok(res)
};
Box::pin(fut)
}
_ => {
Box::pin(async move {
let mut response = http::Response::new(
Expand Down
1 change: 1 addition & 0 deletions components/spider-proto-rust/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ pub mod error;
pub mod id;
pub mod io;
pub mod job;
pub mod scheduler_registration;
pub mod unpack;

#[allow(clippy::all, clippy::nursery, clippy::pedantic)]
Expand Down
40 changes: 40 additions & 0 deletions components/spider-proto-rust/src/scheduler_registration.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
//! Conversions between protobuf scheduler messages and their Spider core representations.

use spider_core::types::scheduler::RegisteredScheduler;

use crate::storage;

impl From<RegisteredScheduler> for storage::Scheduler {
fn from(scheduler: RegisteredScheduler) -> Self {
Self {
scheduler_id: scheduler.id.get(),
ip_address: scheduler.ip_address.to_string(),
port: u32::from(scheduler.port),
}
}
}

#[cfg(test)]
mod tests {
use std::net::IpAddr;

use spider_core::types::{id::SchedulerId, scheduler::RegisteredScheduler};

use crate::storage;

#[test]
fn registered_scheduler_to_protocol_carries_id_ip_and_port() {
const SCHEDULER_ID: SchedulerId = SchedulerId::from(42);
const PORT: u16 = 5678;

let scheduler = storage::Scheduler::from(RegisteredScheduler {
id: SCHEDULER_ID,
ip_address: IpAddr::V4("127.0.0.1".parse().expect("valid IP")),
port: PORT,
});

assert_eq!(scheduler.scheduler_id, SCHEDULER_ID.get());
assert_eq!(scheduler.ip_address, "127.0.0.1");
assert_eq!(scheduler.port, u32::from(PORT));
}
}
Loading
Loading