use tracing::{debug, error, warn};
use nodedb_types::sync::wire::AckStatus;
use super::session::SyncSession;
use super::vector_handler::{VectorDispatcher, VectorInsertParams};
use super::wire::*;
use crate::types::{DatabaseId, TenantId, VShardId};
impl SyncSession {
pub async fn handle_vector_insert<D: VectorDispatcher>(
&mut self,
msg: &VectorInsertMsg,
dispatcher: &D,
) -> Option<SyncFrame> {
self.last_activity = std::time::Instant::now();
if !self.authenticated {
let ack = VectorInsertAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: false,
reject_reason: Some("unauthenticated".to_string()),
applied_seq: 0,
status: AckStatus::Applied,
};
return SyncFrame::try_encode(SyncMessageType::VectorInsertAck, &ack);
}
if msg.vector.len() != msg.dim || msg.dim == 0 {
warn!(
session = %self.session_id,
collection = %msg.collection,
id = %msg.id,
batch_id = msg.batch_id,
stated_dim = msg.dim,
actual_len = msg.vector.len(),
"vector sync: dimension mismatch; rejecting"
);
let ack = VectorInsertAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: false,
reject_reason: Some(format!(
"dimension mismatch: stated {}, actual {}",
msg.dim,
msg.vector.len()
)),
applied_seq: 0,
status: AckStatus::Applied,
};
return SyncFrame::try_encode(SyncMessageType::VectorInsertAck, &ack);
}
let surrogate = match dispatcher.assign_surrogate(
DatabaseId::DEFAULT,
self.tenant_id.unwrap_or(TenantId::new(0)),
&msg.collection,
&msg.id,
) {
Ok(s) => s,
Err(e) => {
error!(
session = %self.session_id,
collection = %msg.collection,
id = %msg.id,
batch_id = msg.batch_id,
error = %e,
"vector sync: surrogate assignment failed"
);
let ack = VectorInsertAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: false,
reject_reason: Some(format!("surrogate assignment failed: {e}")),
applied_seq: 0,
status: AckStatus::Applied,
};
return SyncFrame::try_encode(SyncMessageType::VectorInsertAck, &ack);
}
};
let tenant_id = self.tenant_id.unwrap_or(TenantId::new(0));
let vshard = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &msg.collection);
debug!(
session = %self.session_id,
collection = %msg.collection,
id = %msg.id,
batch_id = msg.batch_id,
dim = msg.dim,
lite_id = %msg.lite_id,
"vector insert: dispatching to Data Plane"
);
match dispatcher
.dispatch_insert(
tenant_id,
vshard,
VectorInsertParams {
collection: msg.collection.clone(),
vector: msg.vector.clone(),
dim: msg.dim,
field_name: msg.field_name.clone(),
surrogate,
pk_bytes: Some(msg.id.as_bytes().to_vec()),
},
nodedb_types::sync::wire::SyncProvenance {
producer_id: self.producer_id,
epoch: self.accepted_epoch,
stream_id: nodedb_types::sync::wire::stream_id_for(
nodedb_types::sync::wire::EngineKind::Vector,
&msg.collection,
),
seq: msg.seq,
},
)
.await
{
Ok(payload_bytes) => {
let gate_result = super::ack_decode::decode_sync_ack(
&payload_bytes,
"vector insert",
&self.session_id,
&msg.collection,
msg.seq,
);
self.mutations_processed += 1;
let ack = VectorInsertAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: true,
reject_reason: None,
applied_seq: gate_result.applied_seq,
status: gate_result.status,
};
SyncFrame::try_encode(SyncMessageType::VectorInsertAck, &ack)
}
Err(e) => {
error!(
session = %self.session_id,
collection = %msg.collection,
id = %msg.id,
batch_id = msg.batch_id,
error = %e,
"vector insert dispatch failed"
);
let ack = VectorInsertAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: false,
reject_reason: Some(e.to_string()),
applied_seq: 0,
status: AckStatus::Applied,
};
SyncFrame::try_encode(SyncMessageType::VectorInsertAck, &ack)
}
}
}
pub async fn handle_vector_delete<D: VectorDispatcher>(
&mut self,
msg: &VectorDeleteMsg,
dispatcher: &D,
) -> Option<SyncFrame> {
self.last_activity = std::time::Instant::now();
if !self.authenticated {
let ack = VectorDeleteAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: false,
reject_reason: Some("unauthenticated".to_string()),
applied_seq: 0,
status: AckStatus::Applied,
};
return SyncFrame::try_encode(SyncMessageType::VectorDeleteAck, &ack);
}
let surrogate = match dispatcher.assign_surrogate(
DatabaseId::DEFAULT,
self.tenant_id.unwrap_or(TenantId::new(0)),
&msg.collection,
&msg.id,
) {
Ok(s) => s,
Err(e) => {
error!(
session = %self.session_id,
collection = %msg.collection,
id = %msg.id,
batch_id = msg.batch_id,
error = %e,
"vector sync: surrogate lookup failed for delete"
);
let ack = VectorDeleteAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: false,
reject_reason: Some(format!("surrogate lookup failed: {e}")),
applied_seq: 0,
status: AckStatus::Applied,
};
return SyncFrame::try_encode(SyncMessageType::VectorDeleteAck, &ack);
}
};
let tenant_id = self.tenant_id.unwrap_or(TenantId::new(0));
let vshard = VShardId::from_collection_in_database(DatabaseId::DEFAULT, &msg.collection);
debug!(
session = %self.session_id,
collection = %msg.collection,
id = %msg.id,
batch_id = msg.batch_id,
lite_id = %msg.lite_id,
"vector delete: dispatching to Data Plane"
);
match dispatcher
.dispatch_delete(
tenant_id,
vshard,
msg.collection.clone(),
surrogate,
msg.field_name.clone(),
nodedb_types::sync::wire::SyncProvenance {
producer_id: self.producer_id,
epoch: self.accepted_epoch,
stream_id: nodedb_types::sync::wire::stream_id_for(
nodedb_types::sync::wire::EngineKind::Vector,
&msg.collection,
),
seq: msg.seq,
},
)
.await
{
Ok(payload_bytes) => {
let gate_result = super::ack_decode::decode_sync_ack(
&payload_bytes,
"vector delete",
&self.session_id,
&msg.collection,
msg.seq,
);
self.mutations_processed += 1;
let ack = VectorDeleteAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: true,
reject_reason: None,
applied_seq: gate_result.applied_seq,
status: gate_result.status,
};
SyncFrame::try_encode(SyncMessageType::VectorDeleteAck, &ack)
}
Err(e) => {
error!(
session = %self.session_id,
collection = %msg.collection,
id = %msg.id,
batch_id = msg.batch_id,
error = %e,
"vector delete dispatch failed"
);
let ack = VectorDeleteAckMsg {
collection: msg.collection.clone(),
id: msg.id.clone(),
batch_id: msg.batch_id,
accepted: false,
reject_reason: Some(e.to_string()),
applied_seq: 0,
status: AckStatus::Applied,
};
SyncFrame::try_encode(SyncMessageType::VectorDeleteAck, &ack)
}
}
}
}