use std::collections::{BTreeMap, HashMap};
use std::sync::Arc;
use async_graphql::dynamic::indexmap::IndexMap;
use async_graphql::dynamic::{
Field, FieldFuture, FieldValue, InputObject, InputValue, Object, Type, TypeRef,
};
use async_graphql::{Name, Value as GraphqlValue};
use surrealdb_strand::Strand;
use super::error::{GraphqlError, resolver_error};
use super::schema::{
SchemaContext, graphql_to_sql_kind_with_scope, kind_to_type_with_enum_prefix, unwrap_type,
};
use super::tables::{CachedRecord, cond_from_filter, idiom_to_graphql_name};
use super::utils::{GraphqlValueUtils, execute_plan};
use crate::catalog::providers::TableProvider;
use crate::catalog::{FieldDefinition, TableDefinition, TableType};
use crate::dbs::Session;
use crate::expr::part::Part;
use crate::expr::statements::{
CreateStatement, DeleteStatement, RelateStatement, UpdateStatement, UpsertStatement,
};
use crate::expr::{Cond, Data, Expr, Kind, Literal, LogicalPlan, Output, TopLevelExpr};
use crate::kvs::Datastore;
use crate::val::{Object as SurObject, RecordId, TableName, Value};
fn capitalize_first(s: &str) -> String {
let mut chars = s.chars();
match chars.next() {
None => String::new(),
Some(c) => c.to_uppercase().to_string() + chars.as_str(),
}
}
fn parse_record_id(table_name: &TableName, id: &str) -> Result<RecordId, GraphqlError> {
let rid_str = format!("{table_name}:{id}");
match crate::syn::record_id(&rid_str) {
Ok(x) => Ok(x.into()),
Err(_) => Ok(RecordId::new(table_name.clone(), id.to_string())),
}
}
fn parse_full_record_id(id_str: &str) -> Result<RecordId, GraphqlError> {
crate::syn::record_id(id_str)
.map(|x| x.into())
.map_err(|e| resolver_error(format!("Invalid record ID: {id_str}: {e}")))
}
fn graphql_input_to_sql_object(
input: &IndexMap<Name, GraphqlValue>,
fds: &[FieldDefinition],
skip_fields: &[&str],
tb_name: &str,
) -> Result<SurObject, GraphqlError> {
let mut map: BTreeMap<Strand, Value> = BTreeMap::new();
for (key, val) in input {
let key_str = key.as_str();
if skip_fields.contains(&key_str) {
continue;
}
if matches!(val, GraphqlValue::Null) {
continue;
}
let matched = fds.iter().find(|fd| {
if fd.name.0.len() != 1 {
return false;
}
if super::naming::field_graphql_name(fd) == key_str {
return true;
}
matches!(&fd.name.0[0], Part::Field(n) if n == key_str)
});
let storage_key: String = matched
.and_then(|fd| match &fd.name.0[0] {
Part::Field(n) => Some(n.as_str().to_owned()),
_ => None,
})
.unwrap_or_else(|| key_str.to_owned());
let kind = matched.and_then(|fd| fd.field_kind.clone()).unwrap_or(Kind::Any);
let enum_scope = format!("{tb_name}_{key_str}");
let sql_val = graphql_to_sql_kind_with_scope(val, kind, Some(&enum_scope))?;
map.insert(storage_key.into(), sql_val);
}
Ok(SurObject::from(map))
}
fn kind_to_input_type_ref(
kind: Kind,
types: &mut Vec<Type>,
enum_scope: Option<&str>,
) -> Result<TypeRef, GraphqlError> {
let optional = kind.can_be_none();
match &kind {
Kind::Record(_) => {
let ty = TypeRef::named(TypeRef::ID);
return Ok(if optional {
ty
} else {
TypeRef::NonNull(Box::new(ty))
});
}
Kind::Array(inner, _) if matches!(**inner, Kind::Record(_)) => {
let inner_ty = TypeRef::NonNull(Box::new(TypeRef::named(TypeRef::ID)));
let list_ty = TypeRef::List(Box::new(inner_ty));
return Ok(if optional {
list_ty
} else {
TypeRef::NonNull(Box::new(list_ty))
});
}
Kind::Either(ks) => {
let non_none: Vec<&Kind> =
ks.iter().filter(|k| !matches!(k, Kind::None | Kind::Null)).collect();
if non_none.len() == 1 {
match non_none[0] {
Kind::Record(_) => return Ok(TypeRef::named(TypeRef::ID)),
Kind::Array(inner, _) if matches!(**inner, Kind::Record(_)) => {
let inner_ty = TypeRef::NonNull(Box::new(TypeRef::named(TypeRef::ID)));
return Ok(TypeRef::List(Box::new(inner_ty)));
}
_ => {}
}
}
}
_ => {}
}
kind_to_type_with_enum_prefix(kind, types, true, enum_scope)
}
fn generate_input_types(
tb_name: &str,
fds: &[FieldDefinition],
is_relation: bool,
types: &mut Vec<Type>,
) -> Result<(String, String, String), GraphqlError> {
let cap_name = capitalize_first(tb_name);
let create_name = format!("Create{cap_name}Input");
let update_name = format!("Update{cap_name}Input");
let upsert_name = format!("Upsert{cap_name}Input");
let mut create_input = InputObject::new(&create_name)
.description(format!("Input for creating a `{tb_name}` record"));
let mut update_input = InputObject::new(&update_name)
.description(format!("Input for updating a `{tb_name}` record"));
let mut upsert_input = InputObject::new(&upsert_name)
.description(format!("Input for upserting a `{tb_name}` record"));
create_input = create_input.field(InputValue::new("id", TypeRef::named(TypeRef::ID)));
upsert_input = upsert_input.field(InputValue::new("id", TypeRef::named(TypeRef::ID)));
if is_relation {
create_input = create_input
.field(InputValue::new("in", TypeRef::named_nn(TypeRef::ID)))
.field(InputValue::new("out", TypeRef::named_nn(TypeRef::ID)));
upsert_input = upsert_input
.field(InputValue::new("in", TypeRef::named_nn(TypeRef::ID)))
.field(InputValue::new("out", TypeRef::named_nn(TypeRef::ID)));
}
for fd in fds.iter() {
let Some(ref kind) = fd.field_kind else {
continue;
};
if fd.name.is_id() {
continue;
}
if fd.name.0.len() > 1 {
continue;
}
let raw_name = idiom_to_graphql_name(&fd.name);
if is_relation && (raw_name == "in" || raw_name == "out") {
continue;
}
if fd.computed.is_some() || fd.readonly {
continue;
}
let fd_name = super::naming::field_graphql_name(fd);
let fd_desc = super::naming::description_with_deprecation(
fd.comment.as_deref(),
fd.graphql_deprecated.as_deref(),
);
let enum_scope = format!("{}_{}", tb_name, fd_name);
let create_type = kind_to_input_type_ref(kind.clone(), types, Some(&enum_scope))?;
let mut create_iv = InputValue::new(&fd_name, create_type.clone());
let mut upsert_iv = InputValue::new(&fd_name, create_type);
if let Some(ref d) = fd_desc {
create_iv = create_iv.description(d);
upsert_iv = upsert_iv.description(d);
}
create_input = create_input.field(create_iv);
upsert_input = upsert_input.field(upsert_iv);
let update_type =
unwrap_type(kind_to_input_type_ref(kind.clone(), types, Some(&enum_scope))?);
let mut update_iv = InputValue::new(&fd_name, update_type);
if let Some(ref d) = fd_desc {
update_iv = update_iv.description(d);
}
update_input = update_input.field(update_iv);
}
types.push(Type::InputObject(create_input));
types.push(Type::InputObject(update_input));
types.push(Type::InputObject(upsert_input));
Ok((create_name, update_name, upsert_name))
}
struct MutationTableContext {
cap_name: String,
cap_plural_name: String,
tb_name_str: String,
tb_name: TableName,
is_relation: bool,
fds: Arc<[FieldDefinition]>,
kvs: Arc<Datastore>,
table_filter_name: String,
graphql_deprecated: Option<String>,
}
impl MutationTableContext {
fn describe(&self, base: String) -> String {
match self.graphql_deprecated.as_deref() {
Some(reason) => format!("[Deprecated: {reason}]\n\n{base}"),
None => base,
}
}
}
pub async fn process_mutations(
tbs: Arc<[TableDefinition]>,
types: &mut Vec<Type>,
schema_ctx: &SchemaContext<'_>,
) -> Result<Object, GraphqlError> {
let mut mutation = Object::new("Mutation");
{
let mut seen: HashMap<String, String> = HashMap::new();
for tb in tbs.iter() {
if tb.view.is_some() {
continue;
}
let cap = super::naming::mutation_cap_name(tb);
let cap_plural = super::naming::mutation_cap_plural_name(tb);
let mut names = Vec::with_capacity(8);
for verb in ["create", "update", "upsert", "delete"] {
names.push(format!("{verb}{cap}"));
names.push(format!("{verb}{cap_plural}"));
}
for name in names {
if let Some(prior) = seen.get(&name) {
return Err(super::error::schema_error(format!(
"GraphQL naming collision on mutation `{}` — both `{}` and `{}` \
produce the same mutation field. Set an explicit `GRAPHQL_ALIAS` \
on one of them.",
name, prior, tb.name
)));
}
seen.insert(name, tb.name.as_str().to_owned());
}
}
}
for tb in tbs.iter() {
if tb.view.is_some() {
continue;
}
let tb_name = tb.name.clone();
let tb_name_str = tb_name.as_str().to_string();
let is_relation = matches!(tb.table_type, TableType::Relation(_));
let fds = schema_ctx.tx.all_tb_fields(schema_ctx.ns, schema_ctx.db, &tb.name, None).await?;
let (create_input_name, update_input_name, upsert_input_name) =
generate_input_types(&tb_name_str, &fds, is_relation, types)?;
let ctx = MutationTableContext {
cap_name: super::naming::mutation_cap_name(tb),
cap_plural_name: super::naming::mutation_cap_plural_name(tb),
table_filter_name: format!("_filter_{tb_name_str}"),
tb_name_str,
tb_name,
is_relation,
fds,
kvs: Arc::clone(schema_ctx.datastore),
graphql_deprecated: tb.graphql_deprecated.clone(),
};
mutation = add_create_field(mutation, &ctx, &create_input_name);
mutation = add_update_field(mutation, &ctx, &update_input_name);
mutation = add_upsert_field(mutation, &ctx, &upsert_input_name);
mutation = add_delete_field(mutation, &ctx);
mutation = add_create_many_field(mutation, &ctx, &create_input_name);
mutation = add_update_many_field(mutation, &ctx, &update_input_name);
mutation = add_upsert_many_field(mutation, &ctx, &upsert_input_name);
mutation = add_delete_many_field(mutation, &ctx);
}
Ok(mutation)
}
fn add_create_field(mutation: Object, tc: &MutationTableContext, input_name: &str) -> Object {
let fds = Arc::clone(&tc.fds);
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
let is_relation = tc.is_relation;
mutation.field(
Field::new(
format!("create{}", tc.cap_name),
TypeRef::named(tc.tb_name_str.as_str()),
move |ctx| {
let fds = Arc::clone(&fds);
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let data_obj = get_data_object(args)?;
let id_opt = data_obj.get("id").and_then(GraphqlValueUtils::as_string);
if is_relation {
execute_relate_create(&kvs, sess, &tb_name, data_obj, &fds, id_opt).await
} else {
execute_normal_create(&kvs, sess, &tb_name, data_obj, &fds, id_opt).await
}
})
},
)
.description(tc.describe(format!("Create a new `{}` record", tc.tb_name_str)))
.argument(InputValue::new("data", TypeRef::named_nn(input_name))),
)
}
fn add_update_field(mutation: Object, tc: &MutationTableContext, input_name: &str) -> Object {
let fds = Arc::clone(&tc.fds);
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
mutation.field(
Field::new(
format!("update{}", tc.cap_name),
TypeRef::named(tc.tb_name_str.as_str()),
move |ctx| {
let fds = Arc::clone(&fds);
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let id_str = get_required_id(args)?;
let data_obj = get_data_object(args)?;
let rid = parse_record_id(&tb_name, &id_str)?;
let content =
graphql_input_to_sql_object(data_obj, &fds, &["id"], tb_name.as_str())?;
let data = if content.0.is_empty() {
None
} else {
Some(Data::MergeExpression(Value::Object(content).into_literal()))
};
let stmt = UpdateStatement {
only: true,
what: vec![Value::RecordId(rid).into_literal()],
data,
cond: None,
output: None,
timeout: Expr::Literal(Literal::None),
..Default::default()
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Update(Box::new(stmt)))],
};
let res = execute_plan(&kvs, sess, plan).await?;
extract_single_record(res)
})
},
)
.description(tc.describe(format!("Update an existing `{}` record", tc.tb_name_str)))
.argument(InputValue::new("id", TypeRef::named_nn(TypeRef::ID)))
.argument(InputValue::new("data", TypeRef::named_nn(input_name))),
)
}
fn add_upsert_field(mutation: Object, tc: &MutationTableContext, input_name: &str) -> Object {
let fds = Arc::clone(&tc.fds);
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
mutation.field(
Field::new(
format!("upsert{}", tc.cap_name),
TypeRef::named(tc.tb_name_str.as_str()),
move |ctx| {
let fds = Arc::clone(&fds);
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let id_str = get_required_id(args)?;
let data_obj = get_data_object(args)?;
let rid = parse_record_id(&tb_name, &id_str)?;
let content =
graphql_input_to_sql_object(data_obj, &fds, &["id"], tb_name.as_str())?;
let data = if content.0.is_empty() {
None
} else {
Some(Data::ContentExpression(Value::Object(content).into_literal()))
};
let stmt = UpsertStatement {
only: true,
what: vec![Value::RecordId(rid).into_literal()],
data,
cond: None,
output: None,
timeout: Expr::Literal(Literal::None),
..Default::default()
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Upsert(Box::new(stmt)))],
};
let res = execute_plan(&kvs, sess, plan).await?;
extract_single_record(res)
})
},
)
.description(
tc.describe(format!("Upsert a `{}` record (create or update)", tc.tb_name_str)),
)
.argument(InputValue::new("id", TypeRef::named_nn(TypeRef::ID)))
.argument(InputValue::new("data", TypeRef::named_nn(input_name))),
)
}
fn add_delete_field(mutation: Object, tc: &MutationTableContext) -> Object {
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
mutation.field(
Field::new(
format!("delete{}", tc.cap_name),
TypeRef::named_nn(TypeRef::BOOLEAN),
move |ctx| {
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let id_str = get_required_id(args)?;
let rid = parse_record_id(&tb_name, &id_str)?;
let stmt = DeleteStatement {
only: false,
what: vec![Value::RecordId(rid).into_literal()],
cond: None,
output: None,
timeout: Expr::Literal(Literal::None),
..Default::default()
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Delete(Box::new(stmt)))],
};
let _res = execute_plan(&kvs, sess, plan).await?;
Ok(Some(FieldValue::value(GraphqlValue::Boolean(true))))
})
},
)
.description(tc.describe(format!("Delete a `{}` record by ID", tc.tb_name_str)))
.argument(InputValue::new("id", TypeRef::named_nn(TypeRef::ID))),
)
}
fn add_create_many_field(mutation: Object, tc: &MutationTableContext, input_name: &str) -> Object {
let fds = Arc::clone(&tc.fds);
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
let is_relation = tc.is_relation;
mutation.field(
Field::new(
format!("create{}", tc.cap_plural_name),
TypeRef::named_nn_list_nn(tc.tb_name_str.as_str()),
move |ctx| {
let fds = Arc::clone(&fds);
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let data_list =
args.get("data").and_then(GraphqlValueUtils::as_list).ok_or_else(|| {
resolver_error("Missing required 'data' argument (must be a list)")
})?;
let mut results = Vec::with_capacity(data_list.len());
for item in data_list {
let data_obj = item.as_object().ok_or_else(|| {
resolver_error("Each item in 'data' must be an object")
})?;
let id_opt = data_obj.get("id").and_then(GraphqlValueUtils::as_string);
let res = if is_relation {
execute_relate_create(&kvs, sess, &tb_name, data_obj, &fds, id_opt)
.await
} else {
execute_normal_create(&kvs, sess, &tb_name, data_obj, &fds, id_opt)
.await
};
match res? {
Some(fv) => results.push(fv),
None => {
return Err(resolver_error(
"Create returned no result for bulk create",
)
.into());
}
}
}
Ok(Some(FieldValue::list(results)))
})
},
)
.description(tc.describe(format!("Create multiple `{}` records", tc.tb_name_str)))
.argument(InputValue::new("data", TypeRef::named_nn_list_nn(input_name))),
)
}
fn add_update_many_field(mutation: Object, tc: &MutationTableContext, input_name: &str) -> Object {
let fds = Arc::clone(&tc.fds);
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
mutation.field(
Field::new(
format!("update{}", tc.cap_plural_name),
TypeRef::named_nn_list_nn(tc.tb_name_str.as_str()),
move |ctx| {
let fds = Arc::clone(&fds);
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let data_obj = get_data_object(args)?;
let content =
graphql_input_to_sql_object(data_obj, &fds, &["id"], tb_name.as_str())?;
let data = if content.0.is_empty() {
None
} else {
Some(Data::MergeExpression(Value::Object(content).into_literal()))
};
let cond = parse_where_arg(args, &fds, tb_name.as_str())?;
let stmt = UpdateStatement {
only: false,
what: vec![Expr::Table(tb_name)],
data,
cond,
output: None,
timeout: Expr::Literal(Literal::None),
..Default::default()
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Update(Box::new(stmt)))],
};
let res = execute_plan(&kvs, sess, plan).await?;
extract_record_list(res)
})
},
)
.description(
tc.describe(format!("Update multiple `{}` records matching a filter", tc.tb_name_str)),
)
.argument(InputValue::new("where", TypeRef::named(tc.table_filter_name.as_str())))
.argument(InputValue::new("data", TypeRef::named_nn(input_name))),
)
}
fn add_upsert_many_field(mutation: Object, tc: &MutationTableContext, input_name: &str) -> Object {
let fds = Arc::clone(&tc.fds);
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
mutation.field(
Field::new(
format!("upsert{}", tc.cap_plural_name),
TypeRef::named_nn_list_nn(tc.tb_name_str.as_str()),
move |ctx| {
let fds = Arc::clone(&fds);
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let data_obj = get_data_object(args)?;
let content =
graphql_input_to_sql_object(data_obj, &fds, &["id"], tb_name.as_str())?;
let data = if content.0.is_empty() {
None
} else {
Some(Data::ContentExpression(Value::Object(content).into_literal()))
};
let cond = parse_where_arg(args, &fds, tb_name.as_str())?;
let stmt = UpsertStatement {
only: false,
what: vec![Expr::Table(tb_name)],
data,
cond,
output: None,
timeout: Expr::Literal(Literal::None),
..Default::default()
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Upsert(Box::new(stmt)))],
};
let res = execute_plan(&kvs, sess, plan).await?;
extract_record_list(res)
})
},
)
.description(
tc.describe(format!("Upsert multiple `{}` records matching a filter", tc.tb_name_str)),
)
.argument(InputValue::new("where", TypeRef::named(tc.table_filter_name.as_str())))
.argument(InputValue::new("data", TypeRef::named_nn(input_name))),
)
}
fn add_delete_many_field(mutation: Object, tc: &MutationTableContext) -> Object {
let fds = Arc::clone(&tc.fds);
let kvs = Arc::clone(&tc.kvs);
let tb_name = tc.tb_name.clone();
mutation.field(
Field::new(
format!("delete{}", tc.cap_plural_name),
TypeRef::named_nn(TypeRef::INT),
move |ctx| {
let fds = Arc::clone(&fds);
let kvs = Arc::clone(&kvs);
let tb_name = tb_name.clone();
FieldFuture::new(async move {
let sess = ctx.data::<Arc<Session>>()?;
let args = ctx.args.as_index_map();
let cond = parse_where_arg(args, &fds, tb_name.as_str())?;
let stmt = DeleteStatement {
only: false,
what: vec![Expr::Table(tb_name)],
cond,
output: Some(Output::Before),
timeout: Expr::Literal(Literal::None),
..Default::default()
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Delete(Box::new(stmt)))],
};
let res = execute_plan(&kvs, sess, plan).await?;
let count = match res {
Value::Array(a) => a.len() as i64,
_ => 0,
};
Ok(Some(FieldValue::value(GraphqlValue::Number(count.into()))))
})
},
)
.description(tc.describe(format!(
"Delete multiple `{}` records matching a filter, returns count",
tc.tb_name_str
)))
.argument(InputValue::new("where", TypeRef::named(tc.table_filter_name.as_str()))),
)
}
fn get_data_object(
args: &IndexMap<Name, GraphqlValue>,
) -> Result<&IndexMap<Name, GraphqlValue>, GraphqlError> {
args.get("data")
.ok_or_else(|| resolver_error("Missing required 'data' argument"))
.and_then(|v| v.as_object().ok_or_else(|| resolver_error("'data' must be an object")))
}
fn get_required_id(args: &IndexMap<Name, GraphqlValue>) -> Result<String, GraphqlError> {
args.get("id")
.and_then(GraphqlValueUtils::as_string)
.ok_or_else(|| resolver_error("Missing required 'id' argument"))
}
fn parse_where_arg(
args: &IndexMap<Name, GraphqlValue>,
fds: &[FieldDefinition],
tb_name: &str,
) -> Result<Option<Cond>, GraphqlError> {
match args.get("where") {
Some(GraphqlValue::Object(o)) if !o.is_empty() => {
Ok(Some(cond_from_filter(o, fds, tb_name, &[])?))
}
_ => Ok(None),
}
}
async fn execute_normal_create(
kvs: &Arc<Datastore>,
sess: &Arc<Session>,
tb_name: &TableName,
data_obj: &IndexMap<Name, GraphqlValue>,
fds: &[FieldDefinition],
id_opt: Option<String>,
) -> Result<Option<FieldValue<'static>>, async_graphql::Error> {
let content = graphql_input_to_sql_object(data_obj, fds, &["id"], tb_name.as_str())?;
let what = match id_opt {
Some(id_str) => {
let rid = parse_record_id(tb_name, &id_str)?;
vec![Value::RecordId(rid).into_literal()]
}
None => vec![Expr::Table(tb_name.clone())],
};
let data = if content.0.is_empty() {
None
} else {
Some(Data::ContentExpression(Value::Object(content).into_literal()))
};
let stmt = CreateStatement {
only: true,
what,
data,
output: None,
timeout: Expr::Literal(Literal::None),
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Create(Box::new(stmt)))],
};
let res = execute_plan(kvs, sess, plan).await?;
extract_single_record(res)
}
async fn execute_relate_create(
kvs: &Arc<Datastore>,
sess: &Arc<Session>,
tb_name: &TableName,
data_obj: &IndexMap<Name, GraphqlValue>,
fds: &[FieldDefinition],
id_opt: Option<String>,
) -> Result<Option<FieldValue<'static>>, async_graphql::Error> {
let in_str = data_obj
.get("in")
.and_then(GraphqlValueUtils::as_string)
.ok_or_else(|| resolver_error("Relation create requires 'in' field"))?;
let out_str = data_obj
.get("out")
.and_then(GraphqlValueUtils::as_string)
.ok_or_else(|| resolver_error("Relation create requires 'out' field"))?;
let from_rid = parse_full_record_id(&in_str)?;
let to_rid = parse_full_record_id(&out_str)?;
let content =
graphql_input_to_sql_object(data_obj, fds, &["id", "in", "out"], tb_name.as_str())?;
let through = match id_opt {
Some(id_str) => {
let rid = parse_record_id(tb_name, &id_str)?;
Value::RecordId(rid).into_literal()
}
None => Expr::Table(tb_name.clone()),
};
let data = if content.0.is_empty() {
None
} else {
Some(Data::ContentExpression(Value::Object(content).into_literal()))
};
let stmt = RelateStatement {
only: true,
or_update: false,
through,
from: Value::RecordId(from_rid).into_literal(),
to: Value::RecordId(to_rid).into_literal(),
data,
output: None,
timeout: Expr::Literal(Literal::None),
};
let plan = LogicalPlan {
expressions: vec![TopLevelExpr::Expr(Expr::Relate(Box::new(stmt)))],
};
let res = execute_plan(kvs, sess, plan).await?;
extract_single_record(res)
}
fn extract_single_record(val: Value) -> Result<Option<FieldValue<'static>>, async_graphql::Error> {
match val {
Value::Object(obj) => {
let rid = match obj.get("id") {
Some(Value::RecordId(rid)) => rid.clone(),
_ => return Err(resolver_error("Mutation result missing 'id' field").into()),
};
Ok(Some(FieldValue::owned_any(CachedRecord {
rid,
version: None,
data: obj,
})))
}
Value::None | Value::Null => Ok(None),
_ => {
error!("Unexpected mutation result type: {val:?}");
Err(resolver_error("Unexpected mutation result").into())
}
}
}
fn extract_record_list(val: Value) -> Result<Option<FieldValue<'static>>, async_graphql::Error> {
match val {
Value::Array(arr) => {
let items: Result<Vec<FieldValue>, GraphqlError> = arr
.0
.into_iter()
.map(|v| match v {
Value::Object(obj) => {
let rid = match obj.get("id") {
Some(Value::RecordId(rid)) => rid.clone(),
_ => return Err(resolver_error("Mutation result missing 'id' field")),
};
Ok(FieldValue::owned_any(CachedRecord {
rid,
version: None,
data: obj,
}))
}
_ => {
error!("Expected object in mutation result, found: {v:?}");
Err(resolver_error("Unexpected mutation result format"))
}
})
.collect();
Ok(Some(FieldValue::list(items?)))
}
_ => {
error!("Expected array result for bulk mutation, found: {val:?}");
Err(resolver_error("Unexpected bulk mutation result format").into())
}
}
}