use tracing::debug;
use super::transfer_compute::{TransferError, compute_transfer};
use crate::bridge::envelope::{ErrorCode, Response};
use crate::data::executor::core_loop::CoreLoop;
use crate::data::executor::response_codec;
use crate::data::executor::task::ExecutionTask;
use crate::engine::kv::current_ms;
pub(in crate::data::executor) struct TransferParams<'a> {
pub did: u64,
pub tid: u64,
pub collection: &'a str,
pub source_key: &'a [u8],
pub dest_key: &'a [u8],
pub field: &'a str,
pub amount: f64,
pub debit_surrogate: nodedb_types::Surrogate,
pub credit_surrogate: nodedb_types::Surrogate,
}
pub(in crate::data::executor) struct TransferItemParams<'a> {
pub did: u64,
pub tid: u64,
pub source_collection: &'a str,
pub dest_collection: &'a str,
pub item_key: &'a [u8],
pub dest_key: &'a [u8],
pub surrogate: nodedb_types::Surrogate,
}
impl CoreLoop {
pub(in crate::data::executor) fn execute_kv_transfer(
&mut self,
task: &ExecutionTask,
params: TransferParams<'_>,
) -> Response {
let TransferParams {
did,
tid,
collection,
source_key,
dest_key,
field,
amount,
debit_surrogate,
credit_surrogate,
} = params;
debug!(core = self.core_id, %collection, %field, amount, "kv transfer");
if self.kv_engine.is_over_budget() {
return self.response_error(task, ErrorCode::ResourcesExhausted);
}
let now_ms = current_ms();
let source_val = self.kv_engine.get(did, tid, collection, source_key, now_ms);
let dest_val = self.kv_engine.get(did, tid, collection, dest_key, now_ms);
let Some(source_bytes) = source_val else {
return self.response_error(task, ErrorCode::NotFound);
};
let dest_bytes = dest_val.unwrap_or_default();
let dest_ref = if dest_bytes.is_empty() {
None
} else {
Some(dest_bytes.as_slice())
};
let computed = match compute_transfer(&source_bytes, dest_ref, field, amount) {
Ok(c) => c,
Err(TransferError::TypeMismatch(detail)) => {
return self.response_error(
task,
ErrorCode::TypeMismatch {
collection: collection.to_string(),
detail,
},
);
}
Err(TransferError::InsufficientBalance { have, need }) => {
return self.response_error(
task,
ErrorCode::InsufficientBalance {
collection: collection.to_string(),
detail: format!("source has {have}, need {need}"),
},
);
}
};
let new_source = computed.new_source;
let new_dest = computed.new_dest;
let source_balance_after = computed.source_balance_after;
let dest_balance_after = computed.dest_balance_after;
if source_key <= dest_key {
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key: source_key,
value: &new_source,
ttl_ms: 0,
now_ms,
surrogate: debit_surrogate,
});
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key: dest_key,
value: &new_dest,
ttl_ms: 0,
now_ms,
surrogate: credit_surrogate,
});
} else {
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key: dest_key,
value: &new_dest,
ttl_ms: 0,
now_ms,
surrogate: credit_surrogate,
});
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key: source_key,
value: &new_source,
ttl_ms: 0,
now_ms,
surrogate: debit_surrogate,
});
}
if let Some(ref m) = self.metrics {
m.record_kv_put();
m.record_kv_put();
}
let src_str = String::from_utf8_lossy(source_key);
let dst_str = String::from_utf8_lossy(dest_key);
self.emit_write_event(
task,
collection,
crate::event::WriteOp::Update,
&src_str,
Some(&new_source),
Some(&source_bytes),
);
self.emit_write_event(
task,
collection,
crate::event::WriteOp::Update,
&dst_str,
Some(&new_dest),
if dest_bytes.is_empty() {
None
} else {
Some(&dest_bytes)
},
);
match response_codec::encode_json(&serde_json::json!({
"source_key": src_str,
"dest_key": dst_str,
"field": field,
"amount": amount,
"source_balance": source_balance_after,
"dest_balance": dest_balance_after,
})) {
Ok(payload) => self.response_with_payload(task, payload),
Err(e) => self.response_error(
task,
ErrorCode::Internal {
detail: e.to_string(),
},
),
}
}
pub(in crate::data::executor) fn execute_kv_transfer_item(
&mut self,
task: &ExecutionTask,
params: TransferItemParams<'_>,
) -> Response {
let TransferItemParams {
did,
tid,
source_collection,
dest_collection,
item_key,
dest_key,
surrogate,
} = params;
debug!(core = self.core_id, %source_collection, %dest_collection, "kv transfer item");
if self.kv_engine.is_over_budget() {
return self.response_error(task, ErrorCode::ResourcesExhausted);
}
let now_ms = current_ms();
let Some(item_data) = self
.kv_engine
.get(did, tid, source_collection, item_key, now_ms)
else {
return self.response_error(task, ErrorCode::NotFound);
};
self.kv_engine
.delete(did, tid, source_collection, &[item_key.to_vec()], now_ms);
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection: dest_collection,
key: dest_key,
value: &item_data,
ttl_ms: 0,
now_ms,
surrogate,
});
if let Some(ref m) = self.metrics {
m.record_kv_delete();
m.record_kv_put();
}
let item_str = String::from_utf8_lossy(item_key);
let dest_str = String::from_utf8_lossy(dest_key);
self.emit_write_event(
task,
source_collection,
crate::event::WriteOp::Delete,
&item_str,
None,
Some(&item_data),
);
self.emit_write_event(
task,
dest_collection,
crate::event::WriteOp::Insert,
&dest_str,
Some(&item_data),
None,
);
match response_codec::encode_json(&serde_json::json!({
"item_key": item_str,
"dest_key": dest_str,
"source_collection": source_collection,
"dest_collection": dest_collection,
})) {
Ok(payload) => self.response_with_payload(task, payload),
Err(e) => self.response_error(
task,
ErrorCode::Internal {
detail: e.to_string(),
},
),
}
}
}