use tracing::debug;
use super::types::KvWriteParams;
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_put(
&mut self,
task: &ExecutionTask,
params: KvWriteParams<'_>,
) -> Response {
let KvWriteParams {
did,
tid,
collection,
key,
value,
ttl_ms,
surrogate,
} = params;
debug!(core = self.core_id, %collection, "kv put");
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: u64 = self.kv_ttl_now_ms(task);
let old = self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key,
value,
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 old.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(value), old_slice);
self.note_kv_write_lsn(task, did, tid, collection, key);
self.response_ok(task)
}
pub(in crate::data::executor) fn execute_kv_insert(
&mut self,
task: &ExecutionTask,
params: KvWriteParams<'_>,
) -> Response {
let KvWriteParams {
did,
tid,
collection,
key,
value,
ttl_ms,
surrogate,
} = params;
debug!(core = self.core_id, %collection, "kv insert");
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);
if self
.kv_engine
.get(did, tid, collection, key, now_ms)
.is_some()
{
let key_str = String::from_utf8_lossy(key);
return self.response_error(
task,
crate::Error::RejectedConstraint {
collection: collection.to_string(),
constraint: "unique".to_string(),
detail: format!(
"duplicate key value '{key_str}' violates primary-key \
uniqueness on '{collection}'"
),
},
);
}
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key,
value,
ttl_ms,
now_ms,
surrogate,
});
if let Some(ref m) = self.metrics {
m.record_kv_put();
}
let key_str = String::from_utf8_lossy(key);
self.emit_write_event(
task,
collection,
crate::event::WriteOp::Insert,
&key_str,
Some(value),
None,
);
self.note_kv_write_lsn(task, did, tid, collection, key);
self.response_ok(task)
}
pub(in crate::data::executor) fn execute_kv_insert_if_absent(
&mut self,
task: &ExecutionTask,
params: KvWriteParams<'_>,
) -> Response {
let KvWriteParams {
did,
tid,
collection,
key,
value,
ttl_ms,
surrogate,
} = params;
debug!(core = self.core_id, %collection, "kv insert-if-absent");
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);
if self
.kv_engine
.get(did, tid, collection, key, now_ms)
.is_some()
{
return self.response_ok(task);
}
self.kv_engine.put(crate::engine::kv::KvPutParams {
database_id: did,
tenant_id: tid,
collection,
key,
value,
ttl_ms,
now_ms,
surrogate,
});
if let Some(ref m) = self.metrics {
m.record_kv_put();
}
let key_str = String::from_utf8_lossy(key);
self.emit_write_event(
task,
collection,
crate::event::WriteOp::Insert,
&key_str,
Some(value),
None,
);
self.note_kv_write_lsn(task, did, tid, collection, key);
self.response_ok(task)
}
pub(in crate::data::executor) fn note_kv_write_lsn(
&mut self,
task: &ExecutionTask,
did: u64,
tid: u64,
collection: &str,
key: &[u8],
) {
let Some(lsn) = task.wal_lsn() else {
return;
};
self.note_write_lsn(
crate::types::DatabaseId::new(did),
crate::types::TenantId::new(tid),
collection,
Some(crate::data::executor::core_loop::write_index::KeyRepr::KvKey(Box::from(key))),
lsn,
);
}
}