(
state: &Arc<SharedState>,
node_id: u64,
tenant_id: u64,
plan: &PhysicalPlan,
)
| 199 | } |
| 200 | |
| 201 | async fn snapshot_remote( |
| 202 | state: &Arc<SharedState>, |
| 203 | node_id: u64, |
| 204 | tenant_id: u64, |
| 205 | plan: &PhysicalPlan, |
| 206 | ) -> Result<Vec<u8>, Error> { |
| 207 | let transport = state |
| 208 | .cluster_transport |
| 209 | .as_ref() |
| 210 | .ok_or_else(|| Error::Internal { |
| 211 | detail: format!("backup: cluster_transport unavailable but node {node_id} is remote"), |
| 212 | })?; |
| 213 | |
| 214 | let plan_bytes = plan_wire::encode(plan).map_err(|e| Error::Internal { |
| 215 | detail: format!("backup: plan encode failed: {e}"), |
| 216 | })?; |
| 217 | let req = RaftRpc::ExecuteRequest(ExecuteRequest { |
| 218 | plan_bytes, |
| 219 | tenant_id, |
| 220 | database_id: DatabaseId::DEFAULT.as_u64(), |
| 221 | deadline_remaining_ms: NODE_SNAPSHOT_TIMEOUT.as_millis() as u64, |
| 222 | trace_id: TraceId::generate().0, |
| 223 | descriptor_versions: Vec::new(), |
| 224 | }); |
| 225 | |
| 226 | let resp = transport |
| 227 | .send_rpc(node_id, req) |
| 228 | .await |
| 229 | .map_err(|e| Error::Internal { |
| 230 | detail: format!("backup: snapshot RPC to node {node_id} failed: {e}"), |
| 231 | })?; |
| 232 | match resp { |
| 233 | RaftRpc::ExecuteResponse(ExecuteResponse { |
| 234 | success: true, |
| 235 | mut payloads, |
| 236 | .. |
| 237 | }) => { |
| 238 | // CreateTenantSnapshot returns exactly one payload. |
| 239 | if payloads.len() != 1 { |
| 240 | return Err(Error::Internal { |
| 241 | detail: format!( |
| 242 | "backup: expected 1 payload from node {node_id}, got {}", |
| 243 | payloads.len() |
| 244 | ), |
| 245 | }); |
| 246 | } |
| 247 | Ok(payloads.remove(0)) |
| 248 | } |
| 249 | RaftRpc::ExecuteResponse(ExecuteResponse { |
| 250 | error: Some(err), .. |
| 251 | }) => Err(map_typed_error(err, node_id)), |
| 252 | RaftRpc::ExecuteResponse(_) => Err(Error::Internal { |
| 253 | detail: format!("backup: empty error response from node {node_id}"), |
| 254 | }), |
| 255 | other => Err(Error::Internal { |
| 256 | detail: format!( |
| 257 | "backup: unexpected RPC response variant from node {node_id}: {other:?}" |
| 258 | ), |
no test coverage detected