use tracing::warn;
use nodedb_physical::physical_plan::StorageMode;
use nodedb_types::columnar::StrictSchema;
use crate::bridge::envelope::{ErrorCode, Response};
use crate::data::executor::core_loop::CoreLoop;
use crate::data::executor::handlers::document::sort;
use crate::data::executor::response_codec;
use crate::data::executor::task::ExecutionTask;
use crate::data::executor::{doc_format, strict_format};
fn decode_body(body: &[u8], schema: Option<&StrictSchema>) -> Option<serde_json::Value> {
match schema {
Some(schema) => strict_format::binary_tuple_to_json(body, schema),
None => doc_format::decode_document(body),
}
}
fn value_in_bounds(value: &str, lower: Option<&[u8]>, upper: Option<&[u8]>) -> bool {
let v = value.as_bytes();
if let Some(l) = lower
&& v < l
{
return false;
}
if let Some(u) = upper
&& v >= u
{
return false;
}
true
}
impl CoreLoop {
pub(in crate::data::executor) fn execute_range_scan_bitemporal(
&self,
task: &ExecutionTask,
args: super::snapshot::RangeScanArgs<'_>,
) -> Response {
let super::snapshot::RangeScanArgs {
tid,
collection,
field,
lower,
upper,
limit,
} = args;
let config_key = (
task.request.database_id,
crate::types::TenantId::new(tid),
collection.to_string(),
);
let strict_schema = self.doc_configs.get(&config_key).and_then(|c| {
if let StorageMode::Strict { ref schema } = c.storage_mode {
Some(schema.clone())
} else {
None
}
});
let predicate = |body: &[u8]| match decode_body(body, strict_schema.as_ref()) {
Some(doc) => {
let values =
crate::engine::document::store::extract_index_values(&doc, field, false);
values.iter().any(|v| value_in_bounds(v, lower, upper))
}
None => false,
};
let scan_limit = limit.max(1000);
let mut scanned = match self.sparse.versioned_scan_as_of(
crate::engine::sparse::btree_versioned::VersionedScanParams {
database_id: task.request.database_id.as_u64(),
tenant: tid,
coll: collection,
sys_cutoff_ms: None,
valid_at_ms: None,
limit: scan_limit,
},
&predicate,
) {
Ok(rows) => rows,
Err(e) => {
warn!(core = self.core_id, error = %e, "versioned range scan failed");
return self.response_error(
task,
ErrorCode::Internal {
detail: e.to_string(),
},
);
}
};
if let Some(txn_id) = task.request.txn_id {
let coll_key = (
task.request.database_id,
crate::types::TenantId::new(tid),
collection.to_string(),
);
self.merge_overlay_into_scan(txn_id, &coll_key, &mut scanned, &predicate);
}
let mut rows: Vec<(String, Vec<u8>)> = Vec::with_capacity(scanned.len());
for (id, body) in scanned {
let mp = match &strict_schema {
Some(schema) => match strict_format::binary_tuple_to_msgpack(&body, schema) {
Some(mp) => mp,
None => continue,
},
None => body,
};
rows.push((id, mp));
}
if let Err(e) = sort::sort_rows(&mut rows, &[(field.to_string(), true)]) {
return self.response_error(
task,
ErrorCode::Internal {
detail: format!("in-memory sort failed: {e}"),
},
);
}
rows.truncate(limit);
match response_codec::encode_raw_document_rows(&rows) {
Ok(payload) => self.response_with_payload(task, payload),
Err(e) => self.response_error(
task,
ErrorCode::Internal {
detail: e.to_string(),
},
),
}
}
}