use std::collections::hash_map::DefaultHasher;
use std::hash::{Hash, Hasher};
use crate::control::gateway::RouteDecision;
use crate::control::server::graph_dispatch::cluster_resolve::resolve_for_vshard;
use crate::control::state::SharedState;
use crate::engine::graph::pattern::executor::VarLenResume;
use crate::types::VShardId;
pub(super) fn resume_seed_key(resume: &VarLenResume) -> u64 {
let mut frontier_nodes: Vec<&str> = resume.frontier.iter().map(|(n, _)| n.as_str()).collect();
frontier_nodes.sort_unstable();
let mut source: Vec<(&str, &str)> = resume
.source_row
.iter()
.map(|(k, v)| (k.as_str(), v.as_str()))
.collect();
source.sort_unstable();
let mut hasher = DefaultHasher::new();
resume.triple_idx.hash(&mut hasher);
resume.depth.hash(&mut hasher);
frontier_nodes.hash(&mut hasher);
source.hash(&mut hasher);
hasher.finish()
}
pub(super) struct PendingResume {
pub(super) remote_coords: Option<(u64, u64)>,
pub(super) resume: VarLenResume,
}
pub(super) fn resume_to_pending(
state: &SharedState,
resume: VarLenResume,
) -> crate::Result<Option<PendingResume>> {
let Some((node_name, _path)) = resume.frontier.first() else {
return Ok(None);
};
let target_vshard = VShardId::from_key(node_name.as_bytes()).as_u32();
let remote_coords = match resolve_for_vshard(state, target_vshard) {
RouteDecision::Local => None,
RouteDecision::Remote { node_id, vshard_id } => Some((node_id, vshard_id)),
RouteDecision::LeaderUnknown { vshard_id } => {
return Err(crate::Error::NotLeader {
vshard_id: VShardId::new((vshard_id % VShardId::COUNT as u64) as u32),
leader_node: 0,
leader_addr: String::new(),
});
}
RouteDecision::Broadcast { .. } => {
return Err(crate::Error::Internal {
detail: "match scatter: resolve_for_vshard returned Broadcast for a \
single vShard"
.into(),
});
}
};
Ok(Some(PendingResume {
remote_coords,
resume,
}))
}