1use crate::wire::{VShardEnvelope, VShardMessageType};
15
16pub enum HandleResult {
18 Response(VShardEnvelope),
20 NoResponse,
22 Error(String),
24}
25
26pub trait VShardHandler: Send + Sync + 'static {
31 fn handle_vshard_envelope(
33 &self,
34 envelope: VShardEnvelope,
35 ) -> impl std::future::Future<Output = HandleResult> + Send;
36}
37
38pub fn dispatch_by_type(envelope: &VShardEnvelope) -> DispatchTarget {
43 match envelope.msg_type {
44 VShardMessageType::TsScatterRequest => DispatchTarget::TimeseriesScan,
46 VShardMessageType::TsScatterResponse => DispatchTarget::TimeseriesCoordinator,
47 VShardMessageType::TsRetentionCommand => DispatchTarget::TimeseriesRetention,
48 VShardMessageType::TsRetentionAck => DispatchTarget::TimeseriesCoordinator,
49 VShardMessageType::TsArchiveCommand => DispatchTarget::TimeseriesArchive,
50 VShardMessageType::TsArchiveAck => DispatchTarget::TimeseriesCoordinator,
51
52 VShardMessageType::SegmentChunk
54 | VShardMessageType::SegmentComplete
55 | VShardMessageType::WalTail
56 | VShardMessageType::RoutingUpdate
57 | VShardMessageType::RoutingAck => DispatchTarget::Migration,
58
59 VShardMessageType::GhostCreate
60 | VShardMessageType::GhostDelete
61 | VShardMessageType::GhostVerifyRequest
62 | VShardMessageType::GhostVerifyResponse => DispatchTarget::Ghost,
63
64 VShardMessageType::MigrationBaseCopy => DispatchTarget::Migration,
65 VShardMessageType::GsiForward => DispatchTarget::Forward,
66 VShardMessageType::EdgeValidation => DispatchTarget::GraphValidation,
67
68 VShardMessageType::VectorScatterRequest => DispatchTarget::VectorSearch,
70 VShardMessageType::VectorScatterResponse => DispatchTarget::VectorCoordinator,
71
72 VShardMessageType::VectorCoarseRouteRequest => DispatchTarget::VectorCoarseRoute,
74 VShardMessageType::VectorCoarseRouteResponse => DispatchTarget::VectorCoordinator,
75
76 VShardMessageType::VectorBuildExchangeRequest => DispatchTarget::VectorBuildExchange,
78 VShardMessageType::VectorBuildExchangeResponse => DispatchTarget::VectorBuildExchange,
79
80 VShardMessageType::VectorMemRegionRequest => DispatchTarget::VectorMemRegion,
82 VShardMessageType::VectorMemRegionResponse => DispatchTarget::VectorMemRegion,
83
84 VShardMessageType::SpatialScatterRequest => DispatchTarget::SpatialSearch,
86 VShardMessageType::SpatialScatterResponse => DispatchTarget::SpatialCoordinator,
87
88 VShardMessageType::CrossShardEvent => DispatchTarget::EventPlane,
90 VShardMessageType::CrossShardEventAck => DispatchTarget::EventPlane,
91 VShardMessageType::NotifyBroadcast => DispatchTarget::EventPlane,
92 VShardMessageType::NotifyBroadcastAck => DispatchTarget::EventPlane,
93
94 VShardMessageType::ArrayShardSliceReq => DispatchTarget::ArrayShard,
96 VShardMessageType::ArrayShardSliceResp => DispatchTarget::ArrayCoordinator,
97 VShardMessageType::ArrayShardAggReq => DispatchTarget::ArrayShard,
98 VShardMessageType::ArrayShardAggResp => DispatchTarget::ArrayCoordinator,
99 VShardMessageType::ArrayShardPutReq => DispatchTarget::ArrayShard,
100 VShardMessageType::ArrayShardPutResp => DispatchTarget::ArrayCoordinator,
101 VShardMessageType::ArrayShardDeleteReq => DispatchTarget::ArrayShard,
102 VShardMessageType::ArrayShardDeleteResp => DispatchTarget::ArrayCoordinator,
103 VShardMessageType::ArrayShardSurrogateBitmapReq => DispatchTarget::ArrayShard,
104 VShardMessageType::ArrayShardSurrogateBitmapResp => DispatchTarget::ArrayCoordinator,
105 }
106}
107
108#[derive(Debug, Clone, Copy, PartialEq, Eq)]
110pub enum DispatchTarget {
111 GraphValidation,
113 TimeseriesScan,
115 TimeseriesCoordinator,
117 TimeseriesRetention,
119 TimeseriesArchive,
121 VectorSearch,
123 VectorCoordinator,
125 VectorCoarseRoute,
127 VectorBuildExchange,
129 VectorMemRegion,
132 Migration,
134 Ghost,
136 Forward,
138 SpatialSearch,
140 SpatialCoordinator,
142 EventPlane,
144 ArrayShard,
146 ArrayCoordinator,
148}
149
150pub fn build_ts_scatter_response(
152 source_node: u64,
153 target_node: u64,
154 vshard_id: u32,
155 partials_json: &[u8],
156) -> VShardEnvelope {
157 VShardEnvelope::new(
158 VShardMessageType::TsScatterResponse,
159 source_node,
160 target_node,
161 vshard_id,
162 partials_json.to_vec(),
163 )
164}
165
166pub fn build_ts_retention_ack(
168 source_node: u64,
169 target_node: u64,
170 vshard_id: u32,
171 result_json: &[u8],
172) -> VShardEnvelope {
173 VShardEnvelope::new(
174 VShardMessageType::TsRetentionAck,
175 source_node,
176 target_node,
177 vshard_id,
178 result_json.to_vec(),
179 )
180}
181
182#[cfg(test)]
183mod tests {
184 use super::*;
185
186 #[test]
187 fn dispatch_ts_scatter() {
188 let env = VShardEnvelope::new(VShardMessageType::TsScatterRequest, 1, 2, 42, vec![]);
189 assert_eq!(dispatch_by_type(&env), DispatchTarget::TimeseriesScan);
190 }
191
192 #[test]
193 fn dispatch_ts_retention() {
194 let env = VShardEnvelope::new(VShardMessageType::TsRetentionCommand, 1, 2, 42, vec![]);
195 assert_eq!(dispatch_by_type(&env), DispatchTarget::TimeseriesRetention);
196 }
197
198 #[test]
199 fn dispatch_migration() {
200 let env = VShardEnvelope::new(VShardMessageType::SegmentChunk, 1, 2, 42, vec![]);
201 assert_eq!(dispatch_by_type(&env), DispatchTarget::Migration);
202 }
203
204 #[test]
205 fn all_message_types_dispatched() {
206 let all_types = [
209 VShardMessageType::SegmentChunk,
210 VShardMessageType::SegmentComplete,
211 VShardMessageType::WalTail,
212 VShardMessageType::RoutingUpdate,
213 VShardMessageType::RoutingAck,
214 VShardMessageType::GhostCreate,
215 VShardMessageType::GhostDelete,
216 VShardMessageType::GhostVerifyRequest,
217 VShardMessageType::GhostVerifyResponse,
218 VShardMessageType::MigrationBaseCopy,
219 VShardMessageType::GsiForward,
220 VShardMessageType::EdgeValidation,
221 VShardMessageType::TsScatterRequest,
222 VShardMessageType::TsScatterResponse,
223 VShardMessageType::TsRetentionCommand,
224 VShardMessageType::TsRetentionAck,
225 VShardMessageType::TsArchiveCommand,
226 VShardMessageType::TsArchiveAck,
227 VShardMessageType::VectorScatterRequest,
228 VShardMessageType::VectorScatterResponse,
229 VShardMessageType::VectorCoarseRouteRequest,
230 VShardMessageType::VectorCoarseRouteResponse,
231 VShardMessageType::VectorBuildExchangeRequest,
232 VShardMessageType::VectorBuildExchangeResponse,
233 VShardMessageType::VectorMemRegionRequest,
234 VShardMessageType::VectorMemRegionResponse,
235 VShardMessageType::SpatialScatterRequest,
236 VShardMessageType::SpatialScatterResponse,
237 VShardMessageType::CrossShardEvent,
238 VShardMessageType::CrossShardEventAck,
239 VShardMessageType::NotifyBroadcast,
240 VShardMessageType::NotifyBroadcastAck,
241 VShardMessageType::ArrayShardSliceReq,
242 VShardMessageType::ArrayShardSliceResp,
243 VShardMessageType::ArrayShardAggReq,
244 VShardMessageType::ArrayShardAggResp,
245 VShardMessageType::ArrayShardPutReq,
246 VShardMessageType::ArrayShardPutResp,
247 VShardMessageType::ArrayShardDeleteReq,
248 VShardMessageType::ArrayShardDeleteResp,
249 VShardMessageType::ArrayShardSurrogateBitmapReq,
250 VShardMessageType::ArrayShardSurrogateBitmapResp,
251 ];
252 for msg_type in all_types {
253 let env = VShardEnvelope::new(msg_type, 1, 2, 0, vec![]);
254 let _ = dispatch_by_type(&env); }
256 }
257}