use nodedb_sql::ddl_ast::GraphProperties;
use crate::bridge::envelope::PhysicalPlan;
use crate::control::planner::calvin::{build_static_tx_class, submit_calvin_routed};
use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::server::shared::session::{DmlTxnCtx, TransactionState};
use crate::control::server::surrogate_exchange::assign_surrogate_routed;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TraceId, VShardId};
use nodedb_physical::physical_plan::GraphOp;
use nodedb_physical::physical_task::{PhysicalTask, PostSetOp};
use super::super::super::result::{DdlError, DdlResult};
use super::support::ddl_err;
const MAX_EDGE_LABEL_BYTES: usize = 256;
fn validate_edge_label(label: &str) -> Result<(), DdlError> {
if label.is_empty() {
return Err(ddl_err("42601", "edge TYPE label must not be empty"));
}
if label.len() > MAX_EDGE_LABEL_BYTES {
return Err(ddl_err(
"42601",
format!(
"edge TYPE label is {} bytes; maximum is {MAX_EDGE_LABEL_BYTES}",
label.len()
),
));
}
if label.chars().any(|c| c.is_control() || c == '\u{007F}') {
return Err(ddl_err(
"42601",
"edge TYPE label must not contain control characters",
));
}
Ok(())
}
pub async fn insert_edge(
state: &SharedState,
identity: &AuthenticatedIdentity,
database_id: DatabaseId,
edge: EdgeRef,
properties: GraphProperties,
txn_ctx: &DmlTxnCtx<'_>,
) -> Result<Vec<DdlResult>, DdlError> {
let EdgeRef {
collection,
src,
dst,
label,
} = edge;
if collection.is_empty() {
return Err(ddl_err(
"42601",
"GRAPH INSERT EDGE requires IN <collection>",
));
}
if src.is_empty() || dst.is_empty() {
return Err(ddl_err("42601", "GRAPH INSERT EDGE requires FROM and TO"));
}
validate_edge_label(&label)?;
let properties_json = properties_to_json(properties)?;
let tenant_id = identity.tenant_id;
crate::control::planner::implicit_edges::mark_collection_edge_bearing(
state,
database_id,
tenant_id,
&collection,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
let vsrc = VShardId::from_key(src.as_bytes());
let vdst = VShardId::from_key(dst.as_bytes());
let src_surrogate = assign_surrogate_routed(
state,
vsrc,
database_id,
tenant_id,
&collection,
src.as_bytes(),
TraceId::ZERO,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
let dst_surrogate = assign_surrogate_routed(
state,
vdst,
database_id,
tenant_id,
&collection,
dst.as_bytes(),
TraceId::ZERO,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
let edge_put = GraphOp::EdgePut {
collection,
src_id: src,
label,
dst_id: dst,
properties: properties_json.into_bytes(),
src_surrogate,
dst_surrogate,
};
let calvin_available =
state.cluster_transport.is_some() && state.sequencer_inbox.get().is_some();
let single_home = vsrc == vdst || !calvin_available;
if txn_ctx.sessions.transaction_state(txn_ctx.addr) == TransactionState::InBlock {
super::edge_stage::stage_edge_dual_home(
state,
tenant_id,
database_id,
EdgeHomes {
vsrc,
vdst,
single_home,
},
edge_put,
txn_ctx,
)
.await?;
return Ok(vec![DdlResult::Status {
command: "INSERT EDGE".to_string(),
rows_affected: None,
}]);
}
if single_home {
let plan = PhysicalPlan::Graph(edge_put);
crate::control::server::sync::raft_dispatch::dispatch_sync_response(
state,
tenant_id,
vsrc,
plan,
TraceId::ZERO,
crate::event::EventSource::User,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
} else {
let task = PhysicalTask {
tenant_id,
vshard_id: vsrc,
database_id,
plan: PhysicalPlan::Graph(edge_put),
post_set_op: PostSetOp::None,
txn_id: None,
};
let tx_class = build_static_tx_class(&[task], tenant_id, &[])
.map_err(|e| ddl_err("XX000", e.to_string()))?;
submit_calvin_routed(state, tx_class)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
}
Ok(vec![DdlResult::Status {
command: "INSERT EDGE".to_string(),
rows_affected: None,
}])
}
pub struct EdgeRef {
pub collection: String,
pub src: String,
pub dst: String,
pub label: String,
}
pub struct EdgeHomes {
pub vsrc: VShardId,
pub vdst: VShardId,
pub single_home: bool,
}
pub async fn delete_edge(
state: &SharedState,
identity: &AuthenticatedIdentity,
database_id: DatabaseId,
edge: EdgeRef,
txn_ctx: &DmlTxnCtx<'_>,
) -> Result<Vec<DdlResult>, DdlError> {
let EdgeRef {
collection,
src,
dst,
label,
} = edge;
if collection.is_empty() {
return Err(ddl_err(
"42601",
"GRAPH DELETE EDGE requires IN <collection>",
));
}
if src.is_empty() || dst.is_empty() {
return Err(ddl_err("42601", "GRAPH DELETE EDGE requires FROM and TO"));
}
validate_edge_label(&label)?;
let tenant_id = identity.tenant_id;
let vsrc = VShardId::from_key(src.as_bytes());
let vdst = VShardId::from_key(dst.as_bytes());
let src_surrogate = assign_surrogate_routed(
state,
vsrc,
database_id,
tenant_id,
&collection,
src.as_bytes(),
TraceId::ZERO,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
let dst_surrogate = assign_surrogate_routed(
state,
vdst,
database_id,
tenant_id,
&collection,
dst.as_bytes(),
TraceId::ZERO,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
let edge_delete = GraphOp::EdgeDelete {
collection,
src_id: src,
label,
dst_id: dst,
src_surrogate,
dst_surrogate,
};
let calvin_available =
state.cluster_transport.is_some() && state.sequencer_inbox.get().is_some();
let single_home = vsrc == vdst || !calvin_available;
if txn_ctx.sessions.transaction_state(txn_ctx.addr) == TransactionState::InBlock {
super::edge_stage::stage_edge_dual_home(
state,
tenant_id,
database_id,
EdgeHomes {
vsrc,
vdst,
single_home,
},
edge_delete,
txn_ctx,
)
.await?;
return Ok(vec![DdlResult::Status {
command: "DELETE EDGE".to_string(),
rows_affected: None,
}]);
}
if single_home {
let plan = PhysicalPlan::Graph(edge_delete);
crate::control::server::sync::raft_dispatch::dispatch_sync_response(
state,
tenant_id,
vsrc,
plan,
TraceId::ZERO,
crate::event::EventSource::User,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
} else {
let task = PhysicalTask {
tenant_id,
vshard_id: vsrc,
database_id,
plan: PhysicalPlan::Graph(edge_delete),
post_set_op: PostSetOp::None,
txn_id: None,
};
let tx_class = build_static_tx_class(&[task], tenant_id, &[])
.map_err(|e| ddl_err("XX000", e.to_string()))?;
submit_calvin_routed(state, tx_class)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
}
Ok(vec![DdlResult::Status {
command: "DELETE EDGE".to_string(),
rows_affected: None,
}])
}
pub async fn set_node_labels(
state: &SharedState,
identity: &AuthenticatedIdentity,
node_id: String,
labels: Vec<String>,
remove: bool,
) -> Result<Vec<DdlResult>, DdlError> {
if node_id.is_empty() {
return Err(ddl_err(
"42601",
"GRAPH LABEL/UNLABEL requires a quoted node id",
));
}
if labels.is_empty() {
return Err(ddl_err("42601", "missing AS '<label>' [, '<label2>']"));
}
let tenant_id = identity.tenant_id;
let vshard_id = VShardId::from_key(node_id.as_bytes());
let plan = if remove {
PhysicalPlan::Graph(GraphOp::RemoveNodeLabels { node_id, labels })
} else {
PhysicalPlan::Graph(GraphOp::SetNodeLabels { node_id, labels })
};
crate::control::server::wal_dispatch::wal_append_if_write(
&state.wal,
tenant_id,
vshard_id,
DatabaseId::DEFAULT,
&plan,
)
.map_err(|e| ddl_err("XX000", e.to_string()))?;
crate::control::server::sync::raft_dispatch::dispatch_sync_response(
state,
tenant_id,
vshard_id,
plan,
TraceId::ZERO,
crate::event::EventSource::User,
)
.await
.map_err(|e| ddl_err("XX000", e.to_string()))?;
let tag = if remove { "UNLABEL" } else { "LABEL" };
Ok(vec![DdlResult::Status {
command: tag.to_string(),
rows_affected: None,
}])
}
fn properties_to_json(properties: GraphProperties) -> Result<String, DdlError> {
match properties {
GraphProperties::None => Ok(String::new()),
GraphProperties::Quoted(s) => Ok(s),
GraphProperties::Object(obj_str) => {
match nodedb_sql::parser::object_literal::parse_object_literal(&obj_str) {
Some(Ok(fields)) => sonic_rs::to_string(&nodedb_types::Value::Object(fields))
.map_err(|e| ddl_err("XX000", format!("PROPERTIES serialize error: {e}"))),
Some(Err(msg)) => Err(ddl_err(
"42601",
format!("PROPERTIES object literal error: {msg}"),
)),
None => Ok(String::new()),
}
}
}
}