use std::borrow::Cow;
use std::sync::Arc;
use anyhow::{Result, bail, ensure};
use reblessive::tree::Stk;
use surrealdb_types::ToSql;
use uuid::Uuid;
use crate::catalog::providers::TableProvider;
use crate::catalog::{self, DatabaseId, Error as CatalogError, NamespaceId, Relation, TableType};
use crate::ctx::FrozenContext;
use crate::dbs::Options;
use crate::doc::CursorDoc;
use crate::exe::FlowResultExt;
use crate::exec::Error as ExecError;
use crate::expr::statements::define::DefineKind;
use crate::expr::statements::define::field::{
DefineDefault, DefineFieldStatement, kind_contains_object,
};
use crate::expr::{Base, Idiom, Kind, KindLiteral, Part, RecordIdKeyLit};
use crate::iam::{Action, AuthLimit, ResourceKind};
use crate::key::KVKeyDecode;
use crate::key::schema::{FieldKey, RefCachePrefix, ReferenceKey, ReferencePrefix};
use crate::kvs::{Direction, NORMAL_BATCH_SIZE, Transaction};
use crate::legacy::{expr_to_ident, expr_to_idiom};
use crate::val::{TableName, Value};
pub(crate) async fn define_field_statement_to_definition(
this: &DefineFieldStatement,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
doc: Option<&CursorDoc>,
) -> Result<catalog::FieldDefinition> {
let comment = stk
.run(|stk| crate::legacy::expr_compute(&this.comment, stk, ctx, opt, doc))
.await
.catch_return()?
.cast_to()?;
let name: Idiom = expr_to_idiom(stk, ctx, opt, doc, &this.name, "field name").await?;
let table: TableName =
expr_to_ident(stk, ctx, opt, doc, &this.what, "table name").await?.into();
if this.computed.is_some() {
let (ns, db) = ctx.get_ns_db_ids(opt).await?;
for ix in ctx.tx().all_tb_indexes(ns, db, &table, None).await?.iter() {
if ix.cols.iter().any(|col| col.starts_with(&name)) {
bail!(ExecError::ComputedFieldCannotBeIndexed {
index: ix.name.to_string(),
field: name.to_raw_string(),
})
}
}
}
Ok(catalog::FieldDefinition {
name,
table,
field_kind: this.field_kind.clone(),
flexible: this.flexible,
readonly: this.readonly,
inline: this.inline,
value: this.value.clone(),
assert: this.assert.clone(),
computed: this.computed.clone(),
default: match &this.default {
DefineDefault::None => catalog::DefineDefault::None,
DefineDefault::Set(x) => catalog::DefineDefault::Set(x.clone()),
DefineDefault::Always(x) => catalog::DefineDefault::Always(x.clone()),
},
select_permission: this.permissions.select.clone(),
create_permission: this.permissions.create.clone(),
update_permission: this.permissions.update.clone(),
comment,
reference: this.reference.clone(),
auth_limit: AuthLimit::new_from_auth(opt.auth.as_ref()).into(),
graphql_alias: this.graphql_alias.clone(),
graphql_deprecated: this.graphql_deprecated.clone(),
})
}
#[instrument(level = "trace", name = "DefineFieldStatement::compute", skip_all)]
pub(crate) async fn define_field_statement_compute(
this: &DefineFieldStatement,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
doc: Option<&CursorDoc>,
) -> Result<Value> {
let definition =
crate::legacy::define_field_statement_to_definition(this, stk, ctx, opt, doc).await?;
ctx.is_allowed(opt, Action::Edit, ResourceKind::Field, Base::Db)?;
crate::legacy::expr::statements::define::validate_graphql_alias(&this.graphql_alias, "field")?;
let (ns_name, db_name) = opt.ns_db()?;
let (ns, db) = ctx.get_ns_db_ids(opt).await?;
crate::legacy::define_field_statement_validate_computed_options(
this,
ns,
db,
ctx.tx(),
&definition,
)
.await?;
crate::legacy::define_field_statement_validate_computed_cycles(
this,
ns,
db,
ctx.tx(),
&definition,
)
.await?;
crate::legacy::define_field_statement_validate_reference_options(this, &definition)?;
crate::legacy::define_field_statement_disallow_mismatched_types(this, ctx, ns, db, &definition)
.await?;
validate_id_field_restrictions(&definition)?;
crate::legacy::define_field_statement_validate_flexible_restrictions(
this,
ctx,
ns,
db,
&definition,
)
.await?;
let txn = ctx.tx();
let tb = txn.get_or_add_tb(Some(ctx), ns_name, db_name, &definition.table, None).await?;
let tb_name = tb.name.clone();
anyhow::ensure!(
crate::kvs::lightweight::lightweight_relation(&tb.table_type).is_none(),
crate::exec::Error::Thrown(
"a LIGHTWEIGHT relation cannot carry a field definition: its edges store no records"
.to_owned(),
)
);
if definition.inline && !opt.import {
anyhow::ensure!(
existing_table_is_relation(&txn, ns, db, &definition.table).await?,
crate::exec::Error::Thrown(
"INLINE requires an existing TYPE RELATION table: the value is embedded into \
the edges' adjacency payloads"
.to_owned(),
)
);
let name = definition.name.to_raw_string();
anyhow::ensure!(
!matches!(name.as_str(), "id" | "in" | "out"),
crate::exec::Error::Thrown(
"the `id`, `in` and `out` fields already live in the adjacency keys and cannot \
be INLINE"
.to_owned(),
)
);
anyhow::ensure!(
definition.name.len() == 1,
crate::exec::Error::Thrown("only a top-level field can be INLINE".to_owned(),)
);
anyhow::ensure!(
definition.computed.is_none(),
crate::exec::Error::Thrown(
"a COMPUTED field cannot be INLINE: its value is derived at read time, so \
there is nothing stored to embed"
.to_owned(),
)
);
let fields = txn.all_tb_fields(ns, db, &definition.table, None).await?;
let others = fields.iter().filter(|f| f.inline && f.name != definition.name).count();
anyhow::ensure!(
others < u8::MAX as usize,
crate::exec::Error::Thrown(
"a table cannot carry more than 255 INLINE fields: the adjacency payload's \
value count is a single byte"
.to_owned(),
)
);
}
let fd = definition.name.to_raw_string();
let existing = txn.get_tb_field(ns, db, &tb_name, &fd, None).await?;
if let Some(existing) = &existing {
match this.kind {
DefineKind::Default => {
if !opt.import {
bail!(CatalogError::FdAlreadyExists {
name: existing.name.to_sql(),
});
}
}
DefineKind::Overwrite => {}
DefineKind::IfNotExists => {
return Ok(Value::None);
}
}
}
txn.put_tb_field(ns, db, &tb_name, &definition).await?;
if !opt.import
&& let Some(existing) = &existing
{
purge_dropped_reference_keys(&txn, ns, db, &tb_name, existing, Some(&definition)).await?;
}
let mut tb = catalog::TableDefinition {
cache_fields_ts: Uuid::now_v7(),
..(*tb).clone()
};
if fd.as_str() == "in" {
if let TableType::Relation(ref relation) = tb.table_type {
if let Some(kind) = this.field_kind.as_ref() {
let Kind::Record(field_kind) = kind else {
bail!(ExecError::Thrown("in field on a relation must be a record".into(),))
};
if *field_kind != relation.from {
tb.table_type = TableType::Relation(Relation {
from: field_kind.clone(),
..relation.clone()
});
txn.replace_tb(ns_name, db_name, &tb).await?;
txn.clear_cache();
return Ok(Value::None);
}
}
}
}
if fd.as_str() == "out" {
if let TableType::Relation(ref relation) = tb.table_type {
if let Some(kind) = this.field_kind.as_ref() {
let Kind::Record(field_kind) = kind else {
bail!(ExecError::Thrown("out field on a relation must be a record".into(),))
};
if *field_kind != relation.to {
tb.table_type = TableType::Relation(Relation {
to: field_kind.clone(),
..relation.clone()
});
txn.replace_tb(ns_name, db_name, &tb).await?;
txn.clear_cache();
return Ok(Value::None);
}
}
}
}
let was_inline = existing.as_ref().is_some_and(|e| e.inline);
let inline_toggled = definition.inline != was_inline;
let inline_kind_changed = definition.inline
&& was_inline
&& existing.as_ref().is_some_and(|e| e.field_kind != definition.field_kind);
if inline_toggled || inline_kind_changed {
tb.graph_inline_gen = tb.graph_inline_gen.wrapping_add(1);
}
txn.replace_tb(ns_name, db_name, &tb).await?;
crate::legacy::define_field_statement_process_recursive_definitions(
this,
ns,
db,
Arc::clone(&txn),
&definition,
)
.await?;
txn.clear_cache();
Ok(Value::None)
}
pub(crate) async fn define_field_statement_disallow_mismatched_types(
this: &DefineFieldStatement,
ctx: &FrozenContext,
ns: NamespaceId,
db: DatabaseId,
definition: &catalog::FieldDefinition,
) -> Result<()> {
let fds = ctx.tx().all_tb_fields(ns, db, &definition.table, None).await?;
if let Some(self_kind) = &this.field_kind {
for fd in fds.iter() {
if definition.name.starts_with(&fd.name)
&& definition.name != fd.name
&& let Some(fd_kind) = &fd.field_kind
{
let path = definition.name[fd.name.len()..].to_vec();
if !fd_kind.allows_nested_kind(&path, self_kind) {
bail!(ExecError::MismatchedFieldTypes {
name: definition.name.to_sql(),
kind: self_kind.to_sql(),
existing_name: fd.name.to_sql(),
existing_kind: fd_kind.to_sql(),
});
}
}
}
}
Ok(())
}
pub(crate) async fn define_field_statement_validate_flexible_restrictions(
this: &DefineFieldStatement,
ctx: &FrozenContext,
ns: NamespaceId,
db: DatabaseId,
definition: &catalog::FieldDefinition,
) -> Result<()> {
if this.flexible {
ensure!(
this.field_kind.as_ref().is_some_and(kind_contains_object),
ExecError::Thrown("FLEXIBLE can only be used with types containing object".into())
);
let txn = ctx.tx();
let Some(tb) = txn.get_tb(ns, db, &definition.table, None).await? else {
bail!(CatalogError::TbNotFound {
name: definition.table.clone(),
});
};
ensure!(
tb.schemafull,
ExecError::Thrown("FLEXIBLE can only be used in SCHEMAFULL tables".into())
);
}
Ok(())
}
pub(crate) async fn define_field_statement_process_recursive_definitions(
this: &DefineFieldStatement,
ns: NamespaceId,
db: DatabaseId,
txn: Arc<Transaction>,
definition: &catalog::FieldDefinition,
) -> Result<()> {
let fields = txn.all_tb_fields(ns, db, &definition.table, None).await.ok();
if let Some(mut cur_kind) = this.field_kind.as_ref().and_then(|x| x.inner_kind()) {
let mut name = definition.name.clone();
loop {
if let Kind::Any = cur_kind {
break;
}
let new_kind = cur_kind.inner_kind();
name.0.push(Part::All);
let fd = name.to_sql();
let key = FieldKey {
ns,
db,
tb: Cow::Borrowed(&definition.table),
fd: Cow::Borrowed(&fd),
};
let val = if let Some(existing) =
fields.as_ref().and_then(|x| x.iter().find(|x| x.name == name))
{
catalog::FieldDefinition {
field_kind: Some(cur_kind.clone()),
flexible: existing.flexible || definition.flexible,
..existing.clone()
}
} else {
catalog::FieldDefinition {
name: name.clone(),
table: definition.table.clone(),
field_kind: Some(cur_kind.clone()),
flexible: definition.flexible,
..Default::default()
}
};
txn.set_key(&key, &val.to_stored()).await?;
if let Some(new_kind) = new_kind {
cur_kind = new_kind;
} else {
break;
}
}
}
Ok(())
}
pub(crate) async fn define_field_statement_validate_computed_options(
this: &DefineFieldStatement,
ns: NamespaceId,
db: DatabaseId,
txn: Arc<Transaction>,
definition: &catalog::FieldDefinition,
) -> Result<()> {
let fields = txn.all_tb_fields(ns, db, &definition.table, None).await?;
if let Some(computed) = this.computed.as_ref() {
ensure!(!computed.contains_mutation(), ExecError::ComputedWrite(definition.name.to_sql()));
ensure!(!definition.name.is_id(), ExecError::IdFieldKeywordConflict("COMPUTED".into()));
ensure!(
definition.name.len() == 1,
ExecError::ComputedNestedField(definition.name.to_sql())
);
ensure!(this.value.is_none(), ExecError::ComputedKeywordConflict("VALUE".into()));
ensure!(this.assert.is_none(), ExecError::ComputedKeywordConflict("ASSERT".into()));
ensure!(this.reference.is_none(), ExecError::ComputedKeywordConflict("REFERENCE".into()));
ensure!(
matches!(this.default, DefineDefault::None),
ExecError::ComputedKeywordConflict("DEFAULT".into())
);
ensure!(!this.readonly, ExecError::ComputedKeywordConflict("READONLY".into()));
for field in fields.iter() {
if field.name.starts_with(&definition.name) && field.name != definition.name {
bail!(ExecError::ComputedNestedFieldConflict(
definition.name.to_sql(),
field.name.to_sql()
));
}
}
} else {
for field in fields.iter() {
if field.computed.is_some()
&& definition.name.starts_with(&field.name)
&& field.name != definition.name
{
bail!(ExecError::ComputedParentFieldConflict(
definition.name.to_sql(),
field.name.to_sql()
));
}
}
}
Ok(())
}
pub(crate) async fn define_field_statement_validate_computed_cycles(
_this: &DefineFieldStatement,
ns: NamespaceId,
db: DatabaseId,
txn: Arc<Transaction>,
definition: &catalog::FieldDefinition,
) -> Result<()> {
if !definition.has_production_clause() {
return Ok(());
}
let fields = txn.all_tb_fields(ns, db, &definition.table, None).await?;
let field_name = definition.name.to_raw_string();
let mut graph: std::collections::BTreeMap<String, Vec<String>> =
std::collections::BTreeMap::new();
for fd in fields.iter() {
if !fd.has_production_clause() {
continue;
}
let name = fd.name.to_raw_string();
if name == field_name {
continue;
}
graph.insert(name, fd.production_dependencies());
}
graph.insert(field_name, definition.production_dependencies());
let mut state: std::collections::BTreeMap<&str, u8> = std::collections::BTreeMap::new();
for key in graph.keys() {
state.insert(key.as_str(), 0);
}
for start in graph.keys() {
if state.get(start.as_str()) == Some(&2) {
continue;
}
let mut stack: Vec<(&str, usize)> = vec![(start.as_str(), 0)];
let mut path: Vec<&str> = vec![start.as_str()];
state.insert(start.as_str(), 1);
while let Some((node, idx)) = stack.last_mut() {
let neighbors = graph.get(*node).map(|v| v.as_slice()).unwrap_or(&[]);
if *idx < neighbors.len() {
let neighbor = neighbors[*idx].as_str();
*idx += 1;
if !graph.contains_key(neighbor) {
continue;
}
match state.get(neighbor) {
Some(1) => {
let cycle_start = path.iter().position(|&n| n == neighbor).unwrap_or(0);
let cycle: Vec<String> =
path[cycle_start..].iter().map(|s| (*s).to_string()).collect();
let cycle_str = format!("{} -> {}", cycle.join(" -> "), neighbor);
bail!(ExecError::ComputedFieldCycle(cycle_str));
}
Some(0) | None => {
state.insert(neighbor, 1);
path.push(neighbor);
stack.push((neighbor, 0));
}
_ => {
}
}
} else {
state.insert(node, 2);
path.pop();
stack.pop();
}
}
}
Ok(())
}
pub(crate) fn define_field_statement_validate_reference_options(
this: &DefineFieldStatement,
definition: &catalog::FieldDefinition,
) -> Result<()> {
if this.reference.is_some() {
ensure!(
definition.name.len() == 1,
ExecError::ReferenceNestedField(definition.name.to_sql())
);
fn valid(kind: &Kind, outer: bool) -> bool {
match kind {
Kind::None | Kind::Record(_) => true,
Kind::Array(kind, _) | Kind::Set(kind, _) => outer && valid(kind, false),
Kind::Literal(KindLiteral::Array(kinds)) => {
outer && kinds.iter().all(|k| valid(k, false))
}
_ => false,
}
}
let is_record_id = match this.field_kind.as_ref() {
Some(Kind::Either(kinds)) => kinds.iter().all(|k| valid(k, true)),
Some(Kind::Array(kind, _)) | Some(Kind::Set(kind, _)) => match kind.as_ref() {
Kind::Either(kinds) => kinds.iter().all(|k| valid(k, true)),
Kind::Record(_) => true,
_ => false,
},
Some(Kind::Literal(KindLiteral::Array(kinds))) => kinds.iter().all(|k| valid(k, true)),
Some(Kind::Record(_)) => true,
_ => false,
};
ensure!(
is_record_id,
ExecError::ReferenceTypeConflict(
this.field_kind.as_ref().unwrap_or(&Kind::Any).to_sql()
)
);
}
Ok(())
}
pub(crate) fn validate_id_field_restrictions(def: &catalog::FieldDefinition) -> Result<()> {
if !def.name.is_id() {
return Ok(());
}
ensure!(def.value.is_none(), ExecError::IdFieldKeywordConflict("VALUE".into()));
ensure!(def.reference.is_none(), ExecError::IdFieldKeywordConflict("REFERENCE".into()));
ensure!(def.computed.is_none(), ExecError::IdFieldKeywordConflict("COMPUTED".into()));
ensure!(
!matches!(def.default, catalog::DefineDefault::Always(_)),
ExecError::IdFieldKeywordConflict("DEFAULT ALWAYS".into())
);
ensure!(!def.readonly, ExecError::IdFieldKeywordConflict("READONLY".into()));
ensure!(!def.flexible, ExecError::IdFieldKeywordConflict("FLEXIBLE".into()));
if let Some(kind) = &def.field_kind {
ensure!(
RecordIdKeyLit::kind_supported(kind),
ExecError::IdFieldUnsupportedKind(kind.to_sql())
);
}
Ok(())
}
pub(crate) async fn purge_dropped_reference_keys(
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
ft: &TableName,
old: &catalog::FieldDefinition,
new: Option<&catalog::FieldDefinition>,
) -> Result<()> {
if old.reference.is_none() {
return Ok(());
}
let ff = old.name.to_sql();
let old_kind: Option<&Kind> = old.field_kind.as_ref();
let new_kind: Option<&Kind> = new.and_then(|n| n.field_kind.as_ref());
for target in txn.all_tb(ns, db, None).await?.iter() {
let target = target.name.clone();
let target = ⌖
let old_can_target = old_kind.is_none_or(|k| k.reference_can_target(target));
if !old_can_target {
continue;
}
let new_can_target = new.is_some_and(|n| {
n.reference.is_some() && new_kind.is_none_or(|k| k.reference_can_target(target))
});
if new_can_target {
continue;
}
let range = ReferencePrefix {
ns,
db,
tb: Cow::Borrowed(target),
}
.range()?;
let mut orphaned: Vec<Vec<u8>> = Vec::new();
let mut cursor = txn.open_keys_cursor(range, Direction::Forward, 0, None).await?;
loop {
let batch = cursor.next_batch(NORMAL_BATCH_SIZE).await?;
if batch.is_empty() {
break;
}
for raw in batch.iter() {
let key = ReferenceKey::decode_key(raw)?;
if key.foreign_table.as_ref() == ft && key.foreign_field.as_ref() == ff.as_str() {
orphaned.push(raw.to_vec());
}
}
}
drop(cursor);
for raw in &orphaned {
let key = ReferenceKey::decode_key(raw)?;
txn.del_key(&key).await?;
}
if !orphaned.is_empty() {
let key = RefCachePrefix {
ns,
db,
tb: Cow::Borrowed(target),
};
txn.del_prefix_key(&key).await?;
}
}
Ok(())
}
async fn existing_table_is_relation(
txn: &Transaction,
ns: NamespaceId,
db: DatabaseId,
table: &TableName,
) -> Result<bool> {
Ok(txn
.get_tb(ns, db, table, None)
.await?
.is_some_and(|tb| matches!(tb.table_type, TableType::Relation(_))))
}