use nodedb_physical::physical_plan::VectorOp;
use crate::types::{DatabaseId, Lsn, TenantId, VShardId};
use crate::wal::manager::WalManager;
pub(crate) fn encode_vector_put_payload(
collection: &str,
vector: &[f32],
dim: usize,
field_name: &str,
surrogate: nodedb_types::Surrogate,
provenance: Option<&nodedb_types::sync::wire::SyncProvenance>,
) -> crate::Result<Vec<u8>> {
let doc_id_compat: Option<String> = None;
zerompk::to_msgpack_vec(&(
collection,
vector,
dim,
field_name,
doc_id_compat,
surrogate.as_u32(),
provenance,
))
.map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal vector insert: {e}"),
})
}
pub(crate) fn encode_vector_batch_put_payload(
collection: &str,
vectors: &[Vec<f32>],
dim: usize,
) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&(collection, vectors, dim)).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal vector batch insert: {e}"),
})
}
pub(crate) fn encode_vector_delete_payload(
collection: &str,
vector_id: u32,
) -> crate::Result<Vec<u8>> {
let prov: Option<nodedb_types::sync::wire::SyncProvenance> = None;
zerompk::to_msgpack_vec(&(collection, vector_id, prov)).map_err(|e| {
crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal vector delete: {e}"),
}
})
}
pub(crate) fn encode_vector_delete_by_surrogate_payload(
collection: &str,
surrogate: nodedb_types::Surrogate,
field_name: &str,
provenance: Option<&nodedb_types::sync::wire::SyncProvenance>,
) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&(collection, surrogate.as_u32(), field_name, provenance)).map_err(
|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal vector delete by surrogate: {e}"),
},
)
}
pub struct VectorPutWalArgs<'a> {
pub collection: &'a str,
pub vector: &'a [f32],
pub dim: usize,
pub field_name: &'a str,
pub surrogate: nodedb_types::Surrogate,
pub provenance: Option<&'a nodedb_types::sync::wire::SyncProvenance>,
}
pub struct VectorDeleteWalArgs<'a> {
pub collection: &'a str,
pub surrogate: nodedb_types::Surrogate,
pub field_name: &'a str,
pub provenance: Option<&'a nodedb_types::sync::wire::SyncProvenance>,
}
pub fn wal_append_vector_put(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
args: VectorPutWalArgs<'_>,
) -> crate::Result<nodedb_types::Lsn> {
let VectorPutWalArgs {
collection,
vector,
dim,
field_name,
surrogate,
provenance,
} = args;
let entry =
encode_vector_put_payload(collection, vector, dim, field_name, surrogate, provenance)?;
let lsn = wal.append_vector_put(tenant_id, vshard_id, database_id, &entry)?;
Ok(lsn)
}
pub(crate) struct VectorDirectUpsertPayload<'a> {
pub collection: &'a str,
pub field: &'a str,
pub surrogate: nodedb_types::Surrogate,
pub vector: &'a [f32],
pub payload: &'a [u8],
pub quantization: nodedb_types::VectorQuantization,
pub storage_dtype: nodedb_types::VectorStorageDtype,
pub payload_indexes: &'a [(String, nodedb_types::PayloadIndexKind)],
}
pub(crate) fn encode_vector_direct_upsert_payload(
args: VectorDirectUpsertPayload<'_>,
) -> crate::Result<Vec<u8>> {
let VectorDirectUpsertPayload {
collection,
field,
surrogate,
vector,
payload,
quantization,
storage_dtype,
payload_indexes,
} = args;
zerompk::to_msgpack_vec(&(
collection,
field,
surrogate.as_u32(),
vector,
payload,
quantization,
storage_dtype,
payload_indexes,
))
.map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal vector direct upsert: {e}"),
})
}
pub(crate) fn encode_sparse_vector_put_payload(
collection: &str,
field_name: &str,
doc_id: &str,
entries: &[(u32, f32)],
) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&(collection, field_name, doc_id, entries)).map_err(|e| {
crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal sparse vector put: {e}"),
}
})
}
pub(crate) fn encode_sparse_vector_delete_payload(
collection: &str,
field_name: &str,
doc_id: &str,
) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&(collection, field_name, doc_id)).map_err(|e| {
crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal sparse vector delete: {e}"),
}
})
}
pub(crate) fn encode_multi_vector_put_payload(
collection: &str,
field_name: &str,
document_surrogate: nodedb_types::Surrogate,
vectors_flat: &[f32],
count: usize,
dim: usize,
) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&(
collection,
field_name,
document_surrogate.as_u32(),
vectors_flat,
count,
dim,
))
.map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal multi-vector put: {e}"),
})
}
pub(crate) fn encode_multi_vector_delete_payload(
collection: &str,
field_name: &str,
document_surrogate: nodedb_types::Surrogate,
) -> crate::Result<Vec<u8>> {
zerompk::to_msgpack_vec(&(collection, field_name, document_surrogate.as_u32())).map_err(|e| {
crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal multi-vector delete: {e}"),
}
})
}
pub(crate) fn wal_append_vector_op(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
op: &VectorOp,
) -> crate::Result<Option<Lsn>> {
let appended = match op {
VectorOp::Insert {
collection,
vector,
dim,
field_name,
surrogate,
pk_bytes: _,
provenance,
} => {
let entry = encode_vector_put_payload(
collection,
vector,
*dim,
field_name,
*surrogate,
provenance.as_ref(),
)?;
Some(wal.append_vector_put(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::BatchInsert {
collection,
vectors,
dim,
surrogates: _,
} => {
let entry = encode_vector_batch_put_payload(collection, vectors, *dim)?;
Some(wal.append_vector_put(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::Delete {
collection,
vector_id,
} => {
let entry = encode_vector_delete_payload(collection, *vector_id)?;
Some(wal.append_vector_delete(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::DeleteBySurrogate {
collection,
surrogate,
field_name,
provenance,
} => {
let entry = encode_vector_delete_by_surrogate_payload(
collection,
*surrogate,
field_name,
provenance.as_ref(),
)?;
Some(wal.append_vector_delete(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::SetParams {
collection,
field_name,
m,
ef_construction,
metric,
index_type,
pq_m,
ivf_cells,
ivf_nprobe,
} => {
let entry = zerompk::to_msgpack_vec(&(
collection,
m,
ef_construction,
metric,
index_type,
pq_m,
ivf_cells,
ivf_nprobe,
field_name,
))
.map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("wal set vector params: {e}"),
})?;
Some(wal.append_vector_params(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::DirectUpsert {
collection,
field,
surrogate,
vector,
payload,
quantization,
storage_dtype,
payload_indexes,
} => {
let entry = encode_vector_direct_upsert_payload(VectorDirectUpsertPayload {
collection,
field,
surrogate: *surrogate,
vector,
payload,
quantization: *quantization,
storage_dtype: *storage_dtype,
payload_indexes,
})?;
Some(wal.append_vector_direct_upsert(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::SparseInsert {
collection,
field_name,
doc_id,
entries,
} => {
let entry = encode_sparse_vector_put_payload(collection, field_name, doc_id, entries)?;
Some(wal.append_sparse_vector_put(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::SparseDelete {
collection,
field_name,
doc_id,
} => {
let entry = encode_sparse_vector_delete_payload(collection, field_name, doc_id)?;
Some(wal.append_sparse_vector_delete(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::MultiVectorInsert {
collection,
field_name,
document_surrogate,
vectors,
count,
dim,
} => {
let entry = encode_multi_vector_put_payload(
collection,
field_name,
*document_surrogate,
vectors,
*count,
*dim,
)?;
Some(wal.append_multi_vector_put(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::MultiVectorDelete {
collection,
field_name,
document_surrogate,
} => {
let entry =
encode_multi_vector_delete_payload(collection, field_name, *document_surrogate)?;
Some(wal.append_multi_vector_delete(tenant_id, vshard_id, database_id, &entry)?)
}
VectorOp::Search { .. }
| VectorOp::MultiSearch { .. }
| VectorOp::SparseSearch { .. }
| VectorOp::MultiVectorScoreSearch { .. }
| VectorOp::QueryStats { .. } => None,
VectorOp::Seal { .. } | VectorOp::CompactIndex { .. } | VectorOp::Rebuild { .. } => None,
};
Ok(appended)
}
pub fn wal_append_vector_delete_by_surrogate(
wal: &WalManager,
tenant_id: TenantId,
vshard_id: VShardId,
database_id: DatabaseId,
args: VectorDeleteWalArgs<'_>,
) -> crate::Result<nodedb_types::Lsn> {
let VectorDeleteWalArgs {
collection,
surrogate,
field_name,
provenance,
} = args;
let entry =
encode_vector_delete_by_surrogate_payload(collection, surrogate, field_name, provenance)?;
let lsn = wal.append_vector_delete(tenant_id, vshard_id, database_id, &entry)?;
Ok(lsn)
}