diff --git a/core/exports/RACommons.exports b/core/exports/RACommons.exports index 4ce71a0d70..107dfe2745 100644 --- a/core/exports/RACommons.exports +++ b/core/exports/RACommons.exports @@ -86,6 +86,10 @@ _rac_http_download_cancel # Connect — protocol/session policy plus the platform transport adapter ABI. _rac_connect_client_create_hello_proto _rac_connect_client_validate_host_proto +_rac_connect_cluster_join_proto +_rac_connect_cluster_start_proto +_rac_connect_cluster_stop_proto +_rac_connect_cluster_validate_activation_proto _rac_connect_get_platform_policy_proto _rac_connect_host_accept_client_proto _rac_connect_host_close_session_proto diff --git a/core/include/rac/connect/rac_connect.h b/core/include/rac/connect/rac_connect.h index 4ba5a26af1..a7464aafe7 100644 --- a/core/include/rac/connect/rac_connect.h +++ b/core/include/rac/connect/rac_connect.h @@ -102,6 +102,43 @@ RAC_API rac_result_t rac_connect_host_validate_cancel_proto(const uint8_t* reque size_t request_size, rac_proto_buffer_t* out_validation); +/* =========================================================================== + * Cluster Orchestration (Issue #541: Local AI Cluster / Pipeline Sharding) + * =========================================================================== */ + +/** + * Start a local AI cluster coordinator session from a serialized + * runanywhere.v1.ClusterStartRequest. Validates peer assignments and returns + * the active ClusterState. + */ +RAC_API rac_result_t rac_connect_cluster_start_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_cluster_state); + +/** + * Join a cluster as a peer from a serialized runanywhere.v1.ClusterJoinRequest. + * Validates the peer's assigned layer range in the cluster and returns + * a ClusterStartResponse echoing the peer's capability. + */ +RAC_API rac_result_t rac_connect_cluster_join_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_response); + +/** + * Stop an active cluster coordinator session and clear peer states. + */ +RAC_API rac_result_t rac_connect_cluster_stop_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_cluster_state); + +/** + * Validate and route an intermediate activation tensor payload between + * pipeline stages. Accepts ClusterActivation and verifies cluster/session binding. + */ +RAC_API rac_result_t rac_connect_cluster_validate_activation_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_validation); + #ifdef __cplusplus } #endif diff --git a/core/src/connect/rac_connect.cpp b/core/src/connect/rac_connect.cpp index 1bae89e59b..2fa7a21acf 100644 --- a/core/src/connect/rac_connect.cpp +++ b/core/src/connect/rac_connect.cpp @@ -5,11 +5,14 @@ #include "rac/connect/rac_connect.h" +#include #include #include #include #include +#include #include +#include #include "rac/core/rac_uuid.h" @@ -630,4 +633,331 @@ rac_result_t rac_connect_host_validate_cancel_proto(const uint8_t* request_bytes #endif } +/* =========================================================================== + * Cluster Orchestration Implementations (Issue #541) + * =========================================================================== */ + +#if defined(RAC_HAVE_PROTOBUF) +struct ClusterRuntime { + bool is_active = false; + std::string cluster_id; + v1::ConnectModelDescriptor model; + std::vector assignments; + std::unordered_map peer_assignments; +}; + +ClusterRuntime& cluster_runtime() { + static ClusterRuntime instance; + return instance; +} + +bool validate_layer_assignments(const google::protobuf::RepeatedPtrField& assignments, + std::string* out_rejection_reason) { + if (assignments.empty()) { + if (out_rejection_reason != nullptr) { + *out_rejection_reason = "Cluster must have at least one layer assignment"; + } + return false; + } + + std::unordered_set seen_instance_ids; + std::vector sorted_assignments; + sorted_assignments.reserve(assignments.size()); + + for (int i = 0; i < assignments.size(); ++i) { + const v1::ClusterLayerAssignment& assignment = assignments[i]; + if (assignment.instance_id().empty()) { + if (out_rejection_reason != nullptr) { + *out_rejection_reason = "Peer assignment has empty instance_id"; + } + return false; + } + if (!seen_instance_ids.insert(assignment.instance_id()).second) { + if (out_rejection_reason != nullptr) { + *out_rejection_reason = "Duplicate instance_id in cluster assignment"; + } + return false; + } + if (assignment.layer_end() <= assignment.layer_start()) { + if (out_rejection_reason != nullptr) { + *out_rejection_reason = "Invalid layer range: layer_end must be greater than layer_start"; + } + return false; + } + sorted_assignments.push_back(assignment); + } + + std::sort(sorted_assignments.begin(), sorted_assignments.end(), + [](const v1::ClusterLayerAssignment& a, const v1::ClusterLayerAssignment& b) { + return a.layer_start() < b.layer_start(); + }); + + uint32_t expected_start = 0; + for (const auto& assignment : sorted_assignments) { + if (assignment.layer_start() != expected_start) { + if (out_rejection_reason != nullptr) { + *out_rejection_reason = "Layer range gap or mismatch in cluster assignment"; + } + return false; + } + expected_start = assignment.layer_end(); + } + return true; +} +#endif + +rac_result_t rac_connect_cluster_start_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_cluster_state) { + if (out_cluster_state == nullptr) { + return RAC_ERROR_NULL_POINTER; + } + rac_proto_buffer_init(out_cluster_state); + +#if !defined(RAC_HAVE_PROTOBUF) + (void)request_bytes; + (void)request_size; + rac_proto_buffer_set_error(out_cluster_state, RAC_ERROR_FEATURE_NOT_AVAILABLE, + "Cluster orchestration requires protobuf support"); + return RAC_ERROR_FEATURE_NOT_AVAILABLE; +#else + v1::ClusterStartRequest request; + const rac_result_t parse_result = + parse_message(request_bytes, request_size, &request, out_cluster_state, + "Invalid ClusterStartRequest protobuf payload"); + if (parse_result != RAC_SUCCESS) { + return parse_result; + } + + std::lock_guard lock(runtime_mutex()); + ClusterRuntime& cluster = cluster_runtime(); + + if (request.cluster_id().empty()) { + rac_proto_buffer_set_error(out_cluster_state, RAC_ERROR_INVALID_ARGUMENT, + "ClusterStartRequest missing cluster_id"); + return RAC_ERROR_INVALID_ARGUMENT; + } + + if (cluster.is_active) { + if (cluster.cluster_id == request.cluster_id()) { + // Idempotent start: return existing state + v1::ClusterState state; + state.set_is_active(true); + state.set_cluster_id(cluster.cluster_id); + *state.mutable_model() = cluster.model; + for (const auto& item : cluster.assignments) { + *state.add_assignments() = item; + } + state.set_peer_count(static_cast(cluster.assignments.size())); + return serialize_message(state, out_cluster_state); + } + rac_proto_buffer_set_error(out_cluster_state, RAC_ERROR_INVALID_ARGUMENT, + "Another cluster is already active"); + return RAC_ERROR_INVALID_ARGUMENT; + } + + if (!is_valid_model(request.model())) { + rac_proto_buffer_set_error(out_cluster_state, RAC_ERROR_INVALID_ARGUMENT, + "ClusterStartRequest missing or invalid model descriptor"); + return RAC_ERROR_INVALID_ARGUMENT; + } + + std::string rejection; + if (!validate_layer_assignments(request.assignments(), &rejection)) { + rac_proto_buffer_set_error(out_cluster_state, RAC_ERROR_INVALID_ARGUMENT, + rejection.c_str()); + return RAC_ERROR_INVALID_ARGUMENT; + } + + cluster.is_active = true; + cluster.cluster_id = request.cluster_id(); + cluster.model = request.model(); + cluster.assignments.clear(); + cluster.peer_assignments.clear(); + + v1::ClusterState state; + state.set_is_active(true); + state.set_cluster_id(cluster.cluster_id); + *state.mutable_model() = cluster.model; + + for (const v1::ClusterLayerAssignment& item : request.assignments()) { + cluster.assignments.push_back(item); + cluster.peer_assignments[item.instance_id()] = item; + *state.add_assignments() = item; + } + state.set_peer_count(static_cast(request.assignments_size())); + + return serialize_message(state, out_cluster_state); +#endif +} + +rac_result_t rac_connect_cluster_join_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_response) { + if (out_response == nullptr) { + return RAC_ERROR_NULL_POINTER; + } + rac_proto_buffer_init(out_response); + +#if !defined(RAC_HAVE_PROTOBUF) + (void)request_bytes; + (void)request_size; + rac_proto_buffer_set_error(out_response, RAC_ERROR_FEATURE_NOT_AVAILABLE, + "Cluster orchestration requires protobuf support"); + return RAC_ERROR_FEATURE_NOT_AVAILABLE; +#else + v1::ClusterJoinRequest request; + const rac_result_t parse_result = + parse_message(request_bytes, request_size, &request, out_response, + "Invalid ClusterJoinRequest protobuf payload"); + if (parse_result != RAC_SUCCESS) { + return parse_result; + } + + std::lock_guard lock(runtime_mutex()); + const ClusterRuntime& cluster = cluster_runtime(); + + v1::ClusterStartResponse response; + if (request.cluster_id().empty()) { + response.set_accepted(false); + response.set_rejection_reason("Cluster ID is empty"); + } else if (request.instance_id().empty()) { + response.set_accepted(false); + response.set_rejection_reason("Peer instance_id is empty"); + } else if (!cluster.is_active) { + response.set_accepted(false); + response.set_rejection_reason("Cluster is not active"); + } else if (request.cluster_id() != cluster.cluster_id) { + response.set_accepted(false); + response.set_rejection_reason("Cluster ID does not match active cluster"); + } else if (!request.peer_capability().instance_id().empty() && + request.peer_capability().instance_id() != request.instance_id()) { + response.set_accepted(false); + response.set_rejection_reason("Peer capability instance_id does not match request identity"); + } else { + auto it = cluster.peer_assignments.find(request.instance_id()); + if (it == cluster.peer_assignments.end()) { + response.set_accepted(false); + response.set_rejection_reason("Peer instance_id not found in cluster layer assignments"); + } else { + response.set_accepted(true); + *response.mutable_peer_capability() = request.peer_capability(); + if (response.peer_capability().instance_id().empty()) { + response.mutable_peer_capability()->set_instance_id(request.instance_id()); + } + } + } + + return serialize_message(response, out_response); +#endif +} + +rac_result_t rac_connect_cluster_stop_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_cluster_state) { + if (out_cluster_state == nullptr) { + return RAC_ERROR_NULL_POINTER; + } + rac_proto_buffer_init(out_cluster_state); + +#if !defined(RAC_HAVE_PROTOBUF) + (void)request_bytes; + (void)request_size; + rac_proto_buffer_set_error(out_cluster_state, RAC_ERROR_FEATURE_NOT_AVAILABLE, + "Cluster orchestration requires protobuf support"); + return RAC_ERROR_FEATURE_NOT_AVAILABLE; +#else + v1::ClusterStopRequest request; + if (request_bytes != nullptr && request_size > 0) { + const rac_result_t parse_result = + parse_message(request_bytes, request_size, &request, out_cluster_state, + "Invalid ClusterStopRequest protobuf payload"); + if (parse_result != RAC_SUCCESS) { + return parse_result; + } + } + + std::lock_guard lock(runtime_mutex()); + ClusterRuntime& cluster = cluster_runtime(); + + if (!request.cluster_id().empty() && cluster.is_active && + request.cluster_id() != cluster.cluster_id) { + rac_proto_buffer_set_error(out_cluster_state, RAC_ERROR_INVALID_ARGUMENT, + "ClusterStopRequest cluster_id does not match the active cluster"); + return RAC_ERROR_INVALID_ARGUMENT; + } + + const std::string stopped_cluster_id = cluster.cluster_id; + cluster.is_active = false; + cluster.cluster_id.clear(); + cluster.model.Clear(); + cluster.assignments.clear(); + cluster.peer_assignments.clear(); + + v1::ClusterState state; + state.set_is_active(false); + state.set_cluster_id(stopped_cluster_id); + state.set_peer_count(0); + return serialize_message(state, out_cluster_state); +#endif +} + +rac_result_t rac_connect_cluster_validate_activation_proto(const uint8_t* request_bytes, + size_t request_size, + rac_proto_buffer_t* out_validation) { + if (out_validation == nullptr) { + return RAC_ERROR_NULL_POINTER; + } + rac_proto_buffer_init(out_validation); + +#if !defined(RAC_HAVE_PROTOBUF) + (void)request_bytes; + (void)request_size; + rac_proto_buffer_set_error(out_validation, RAC_ERROR_FEATURE_NOT_AVAILABLE, + "Cluster orchestration requires protobuf support"); + return RAC_ERROR_FEATURE_NOT_AVAILABLE; +#else + v1::ClusterActivation request; + const rac_result_t parse_result = + parse_message(request_bytes, request_size, &request, out_validation, + "Invalid ClusterActivation protobuf payload"); + if (parse_result != RAC_SUCCESS) { + return parse_result; + } + + v1::ConnectInvocationValidation validation; + std::lock_guard lock(runtime_mutex()); + const ClusterRuntime& cluster = cluster_runtime(); + + if (!cluster.is_active) { + validation.set_rejection_reason("Cluster is not active"); + } else if (request.cluster_id() != cluster.cluster_id) { + validation.set_rejection_reason("Cluster ID mismatch on activation tensor"); + } else if (request.request_id().empty()) { + validation.set_rejection_reason("Activation request_id is empty"); + } else if (request.tensor_data().empty()) { + validation.set_rejection_reason("Activation tensor_data is empty"); + } else if (request.seq_len() == 0 || request.hidden_size() == 0) { + validation.set_rejection_reason("Invalid activation tensor dimensions (seq_len/hidden_size zero)"); + } else { + bool valid_boundary = (request.from_layer() == 0); + for (const auto& a : cluster.assignments) { + if (request.from_layer() == a.layer_start() || request.from_layer() == a.layer_end()) { + valid_boundary = true; + break; + } + } + const uint64_t min_elements = static_cast(request.seq_len()) * request.hidden_size(); + if (!valid_boundary) { + validation.set_rejection_reason("Activation from_layer does not match any cluster layer boundary"); + } else if (request.tensor_data().size() < min_elements) { + validation.set_rejection_reason("Activation tensor_data size is smaller than tensor dimensions"); + } else { + validation.set_accepted(true); + } + } + return serialize_message(validation, out_validation); +#endif +} + } // extern "C" diff --git a/core/tests/test_connect_proto_abi.cpp b/core/tests/test_connect_proto_abi.cpp index d47ef56a49..b369413fe2 100644 --- a/core/tests/test_connect_proto_abi.cpp +++ b/core/tests/test_connect_proto_abi.cpp @@ -274,6 +274,170 @@ void test_session_cap_and_cancel_validation() { stop_host(); } +void test_cluster_orchestration() { + // Reset cluster runtime state to ensure clean test isolation + { + v1::ClusterStopRequest reset_request; + v1::ClusterState reset_state; + call_proto(rac_connect_cluster_stop_proto, reset_request, &reset_state); + } + + // 1. Invalid cluster start (empty cluster id) + v1::ClusterStartRequest invalid_request; + v1::ClusterState invalid_state; + CHECK(call_proto(rac_connect_cluster_start_proto, invalid_request, &invalid_state) == RAC_ERROR_INVALID_ARGUMENT, + "Cluster start rejects empty request"); + + // 2. Invalid cluster start (duplicate instance_id) + v1::ClusterStartRequest dup_request; + dup_request.set_cluster_id("test-cluster-dup"); + dup_request.mutable_model()->set_model_id("llama-3-70b"); + dup_request.mutable_model()->set_display_name("Llama 3 70B"); + auto* d1 = dup_request.add_assignments(); + d1->set_instance_id("peer-1"); + d1->set_layer_start(0); + d1->set_layer_end(16); + auto* d2 = dup_request.add_assignments(); + d2->set_instance_id("peer-1"); + d2->set_layer_start(16); + d2->set_layer_end(32); + + v1::ClusterState dup_state; + CHECK(call_proto(rac_connect_cluster_start_proto, dup_request, &dup_state) == RAC_ERROR_INVALID_ARGUMENT, + "Cluster start rejects duplicate peer instance_id"); + + // 3. Invalid cluster start (layer gap: [0, 10) and [15, 30)) + v1::ClusterStartRequest gap_request; + gap_request.set_cluster_id("test-cluster-gap"); + gap_request.mutable_model()->set_model_id("llama-3-70b"); + gap_request.mutable_model()->set_display_name("Llama 3 70B"); + auto* a1 = gap_request.add_assignments(); + a1->set_instance_id("mac-coord"); + a1->set_layer_start(0); + a1->set_layer_end(10); + a1->set_is_coordinator(true); + auto* a2 = gap_request.add_assignments(); + a2->set_instance_id("android-peer"); + a2->set_layer_start(15); + a2->set_layer_end(30); + + v1::ClusterState gap_state; + CHECK(call_proto(rac_connect_cluster_start_proto, gap_request, &gap_state) == RAC_ERROR_INVALID_ARGUMENT, + "Cluster start rejects layer assignment gap"); + + // 4. Valid cluster start with unordered assignments: [16, 32) then [0, 16) + v1::ClusterStartRequest valid_request; + valid_request.set_cluster_id("cluster-alpha"); + valid_request.mutable_model()->set_model_id("llama-3-70b"); + valid_request.mutable_model()->set_display_name("Llama 3 70B"); + auto* v2_assign = valid_request.add_assignments(); + v2_assign->set_instance_id("android-snapdragon"); + v2_assign->set_layer_start(16); + v2_assign->set_layer_end(32); + + auto* v1_assign = valid_request.add_assignments(); + v1_assign->set_instance_id("mac-coordinator"); + v1_assign->set_layer_start(0); + v1_assign->set_layer_end(16); + v1_assign->set_is_coordinator(true); + + v1::ClusterState active_state; + CHECK(call_proto(rac_connect_cluster_start_proto, valid_request, &active_state) == RAC_SUCCESS, + "Cluster coordinator starts with valid layer sharding"); + CHECK(active_state.is_active(), "Cluster state is marked active"); + CHECK(active_state.peer_count() == 2, "Cluster state tracks 2 participating devices"); + CHECK(active_state.cluster_id() == "cluster-alpha", "Cluster ID matches request"); + + // 5. Idempotent cluster start with same cluster_id + v1::ClusterState idempotent_state; + CHECK(call_proto(rac_connect_cluster_start_proto, valid_request, &idempotent_state) == RAC_SUCCESS, + "Repeated cluster start with same ID is idempotent"); + CHECK(idempotent_state.is_active(), "Idempotent start returns active state"); + + // 6. Conflicting cluster start with different cluster_id while active is rejected + v1::ClusterStartRequest conflict_request = valid_request; + conflict_request.set_cluster_id("cluster-beta"); + v1::ClusterState conflict_state; + CHECK(call_proto(rac_connect_cluster_start_proto, conflict_request, &conflict_state) == RAC_ERROR_INVALID_ARGUMENT, + "Starting a different cluster while one is active is rejected"); + + // 7. Peer joins cluster with ClusterJoinRequest + v1::ClusterJoinRequest join_request; + join_request.set_cluster_id("cluster-alpha"); + join_request.set_instance_id("android-snapdragon"); + join_request.mutable_peer_capability()->set_instance_id("android-snapdragon"); + join_request.mutable_peer_capability()->set_available_memory_bytes(8589934592ULL); + join_request.mutable_peer_capability()->set_acceleration(v1::ACCELERATION_PREFERENCE_NPU); + + v1::ClusterStartResponse join_response; + CHECK(call_proto(rac_connect_cluster_join_proto, join_request, &join_response) == RAC_SUCCESS, + "Peer join request is validated successfully against active cluster"); + CHECK(join_response.accepted(), "Peer join response is accepted"); + CHECK(join_response.peer_capability().available_memory_bytes() == 8589934592ULL, + "Peer capability is populated in join response"); + + // 8. Peer join with unknown instance_id is rejected + v1::ClusterJoinRequest unknown_peer_join = join_request; + unknown_peer_join.set_instance_id("unknown-device"); + unknown_peer_join.mutable_peer_capability()->set_instance_id("unknown-device"); + v1::ClusterStartResponse unknown_peer_response; + CHECK(call_proto(rac_connect_cluster_join_proto, unknown_peer_join, &unknown_peer_response) == RAC_SUCCESS, + "Unknown peer join returns typed validation response"); + CHECK(!unknown_peer_response.accepted(), "Unknown peer instance_id is rejected"); + + // 9. Activation validation between pipeline stages + v1::ClusterActivation valid_activation; + valid_activation.set_cluster_id("cluster-alpha"); + valid_activation.set_request_id("req-token-101"); + valid_activation.set_from_layer(16); + valid_activation.set_tensor_data(std::string(4096 * 2, 'A')); + valid_activation.set_seq_len(1); + valid_activation.set_hidden_size(4096); + + v1::ConnectInvocationValidation activation_validation; + CHECK(call_proto(rac_connect_cluster_validate_activation_proto, valid_activation, &activation_validation) == RAC_SUCCESS, + "Valid intermediate activation tensor is validated"); + CHECK(activation_validation.accepted(), "Activation tensor passes validation"); + + // 10. Activation validation rejection on mismatched cluster ID + v1::ClusterActivation wrong_cluster_activation = valid_activation; + wrong_cluster_activation.set_cluster_id("cluster-wrong"); + v1::ConnectInvocationValidation wrong_cluster_validation; + CHECK(call_proto(rac_connect_cluster_validate_activation_proto, wrong_cluster_activation, &wrong_cluster_validation) == RAC_SUCCESS, + "Mismatched cluster activation returns typed validation"); + CHECK(!wrong_cluster_validation.accepted(), "Mismatched cluster ID activation is rejected"); + + // 11. Activation validation rejection on invalid layer boundary + v1::ClusterActivation invalid_boundary_activation = valid_activation; + invalid_boundary_activation.set_from_layer(7); // not a boundary in [0, 16) or [16, 32) + v1::ConnectInvocationValidation invalid_boundary_validation; + CHECK(call_proto(rac_connect_cluster_validate_activation_proto, invalid_boundary_activation, &invalid_boundary_validation) == RAC_SUCCESS, + "Invalid layer boundary activation returns typed validation"); + CHECK(!invalid_boundary_validation.accepted(), "Invalid layer boundary is rejected"); + + // 12. Stop cluster with mismatched cluster ID is rejected + v1::ClusterStopRequest wrong_stop_request; + wrong_stop_request.set_cluster_id("cluster-different"); + v1::ClusterState wrong_stop_state; + CHECK(call_proto(rac_connect_cluster_stop_proto, wrong_stop_request, &wrong_stop_state) == RAC_ERROR_INVALID_ARGUMENT, + "Stopping with mismatched cluster_id is rejected"); + + // 13. Stop active cluster cleanly + v1::ClusterStopRequest stop_request; + stop_request.set_cluster_id("cluster-alpha"); + v1::ClusterState stopped_state; + CHECK(call_proto(rac_connect_cluster_stop_proto, stop_request, &stopped_state) == RAC_SUCCESS, + "Cluster stop terminates active cluster"); + CHECK(!stopped_state.is_active(), "Cluster state is marked inactive after stop"); + CHECK(stopped_state.cluster_id() == "cluster-alpha", "Stopped cluster ID is preserved in response"); + + // 14. Activation rejected after stop + v1::ConnectInvocationValidation stopped_activation_validation; + CHECK(call_proto(rac_connect_cluster_validate_activation_proto, valid_activation, &stopped_activation_validation) == RAC_SUCCESS, + "Activation check after stop returns typed validation"); + CHECK(!stopped_activation_validation.accepted(), "Activation tensor is rejected when cluster is inactive"); +} + #endif } // namespace @@ -288,6 +452,7 @@ int main() { test_client_admission(); test_host_handshake_and_reconnect_deduplication(); test_session_cap_and_cancel_validation(); + test_cluster_orchestration(); std::fprintf(stdout, " %d checks, %d failures\n", g_checks, g_failures); return g_failures == 0 ? 0 : 1; #endif diff --git a/idl/SCHEMA_LOCK b/idl/SCHEMA_LOCK index 4928cd35b9..ea6eae0449 100644 --- a/idl/SCHEMA_LOCK +++ b/idl/SCHEMA_LOCK @@ -12,7 +12,7 @@ # # Verify: idl/codegen/schema_lock.sh --check # ============================================================================= -IDL_VERSION=1.1.0 +IDL_VERSION=1.2.0 IDL_PROTO_COUNT=39 -IDL_SCHEMA_SHA256=377868bf8e12b327c32978215cafb36c2b1b3157a9be05c7960ac474d55231d2 +IDL_SCHEMA_SHA256=86468efcd7638fd690aaca7256baf66d60813b94a7d6639f143ccd49d0cfe61f IDL_PROTOC_VERSION=35.1 diff --git a/idl/VERSION b/idl/VERSION index 9084fa2f71..26aaba0e86 100644 --- a/idl/VERSION +++ b/idl/VERSION @@ -1 +1 @@ -1.1.0 +1.2.0 diff --git a/idl/connect.proto b/idl/connect.proto index d1b9eea36f..aa5d5fa1b6 100644 --- a/idl/connect.proto +++ b/idl/connect.proto @@ -9,6 +9,7 @@ syntax = "proto3"; package runanywhere.v1; +import "hardware_profile.proto"; import "llm_service.proto"; option cc_enable_arenas = true; @@ -223,3 +224,69 @@ message ConnectHostFrame { ConnectHeartbeatResponse heartbeat = 2; } } + +// ----------------------------------------------------------------------------- +// Cluster-Specific Schemas (Issue #541: Local AI Cluster / Pipeline Sharding) +// ----------------------------------------------------------------------------- + +// Compute and memory capacity advertised by a peer participating in a cluster. +message ClusterPeerCapability { + string instance_id = 1; + ConnectPlatform platform = 2; + uint64 available_memory_bytes = 3; + AccelerationPreference acceleration = 4; + NpuCapability npu = 5; + uint32 max_layers = 6; +} + +// Layer range assigned to a peer by the coordinator. +message ClusterLayerAssignment { + string instance_id = 1; + // Half-open range: this peer owns layers [layer_start, layer_end), with layer_end excluded. + uint32 layer_start = 2; + uint32 layer_end = 3; + bool is_coordinator = 4; +} + +// Coordinator starts cluster and sends assignments to all peers. +message ClusterStartRequest { + string cluster_id = 1; + ConnectModelDescriptor model = 2; + repeated ClusterLayerAssignment assignments = 3; +} + +// Sent by a peer device when joining a cluster to validate its assignment and register capabilities. +message ClusterJoinRequest { + string cluster_id = 1; + string instance_id = 2; + ClusterPeerCapability peer_capability = 3; +} + +message ClusterStartResponse { + bool accepted = 1; + string rejection_reason = 2; + ClusterPeerCapability peer_capability = 3; +} + +// Intermediate activation tensor passed between pipeline stages. +message ClusterActivation { + string cluster_id = 1; + string request_id = 2; + uint32 from_layer = 3; + bytes tensor_data = 4; + uint32 seq_len = 5; + uint32 hidden_size = 6; +} + +message ClusterStopRequest { + string cluster_id = 1; +} + +message ClusterState { + bool is_active = 1; + string cluster_id = 2; + ConnectModelDescriptor model = 3; + repeated ClusterLayerAssignment assignments = 4; + uint32 peer_count = 5; +} +