use tracing::debug;
use super::types::KvInsertOnConflictUpdateParams;
use crate::bridge::envelope::{ErrorCode, Response};
use crate::data::executor::core_loop::CoreLoop;
use crate::data::executor::task::ExecutionTask;
impl CoreLoop {
pub(in crate::data::executor) fn execute_kv_insert_on_conflict_update(
&mut self,
task: &ExecutionTask,
params: KvInsertOnConflictUpdateParams<'_>,
) -> Response {
let KvInsertOnConflictUpdateParams {
did,
tid,
collection,
key,
value,
ttl_ms,
updates,
surrogate,
} = params;
debug!(core = self.core_id, %collection, "kv insert-on-conflict-update");
if self.kv_engine.is_over_budget() {
return self.response_error(
task,
ErrorCode::Internal {
detail: "KV memory budget exceeded, retry later".into(),
},
);
}
let now_ms = self.kv_ttl_now_ms(task);
let existing_bytes = self.kv_engine.get(did, tid, collection, key, now_ms);
let stored_bytes: Vec<u8> = match &existing_bytes {
None => value.to_vec(),
Some(existing_raw) => {
let existing_val = match nodedb_types::value_from_msgpack(existing_raw) {
Ok(v) => v,
Err(_) => {
return self.response_error(
task,
ErrorCode::Internal {
detail: "failed to decode existing KV value for ON CONFLICT \
DO UPDATE"
.into(),
},
);
}
};
let excluded_val = match nodedb_types::value_from_msgpack(value) {
Ok(v) => v,
Err(_) => {
return self.response_error(
task,
ErrorCode::Internal {
detail: "failed to decode incoming KV value for ON CONFLICT \
DO UPDATE"
.into(),
},
);
}
};
let merged = crate::data::executor::handlers::upsert::apply_on_conflict_updates(
existing_val,
&excluded_val,
updates,
);
match nodedb_types::value_to_msgpack(&merged) {
Ok(b) => b,
Err(_) => {
return self.response_error(
task,
ErrorCode::Internal {
detail: "failed to encode merged KV value".into(),
},
);
}
}
}
};
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key,
value: &stored_bytes,
ttl_ms,
now_ms,
surrogate,
});
if let Some(ref m) = self.metrics {
m.record_kv_put();
}
let key_str = String::from_utf8_lossy(key);
let (op, old_slice): (_, Option<&[u8]>) = match existing_bytes.as_deref() {
Some(o) => (crate::event::WriteOp::Update, Some(o)),
None => (crate::event::WriteOp::Insert, None),
};
self.emit_write_event(
task,
collection,
op,
&key_str,
Some(&stored_bytes),
old_slice,
);
self.response_ok(task)
}
}