use anyhow::{Result, bail};
use reblessive::tree::Stk;
use super::{CursorDoc, DefineKind};
use crate::catalog::providers::BucketProvider;
use crate::catalog::{BucketDefinition, Permission};
use crate::ctx::FrozenContext;
use crate::dbs::Options;
use crate::err::Error;
use crate::expr::parameterize::expr_to_ident;
use crate::expr::{Base, Expr, FlowResultExt, Literal};
use crate::iam::{Action, ResourceKind};
use crate::val::Value;
#[derive(Clone, Debug, Eq, PartialEq, Hash)]
pub(crate) struct DefineBucketStatement {
pub kind: DefineKind,
pub name: Expr,
pub backend: Option<Expr>,
pub permissions: Permission,
pub readonly: bool,
pub comment: Expr,
}
impl Default for DefineBucketStatement {
fn default() -> Self {
Self {
kind: DefineKind::Default,
name: Expr::Literal(Literal::None),
backend: None,
permissions: Permission::default(),
readonly: false,
comment: Expr::Literal(Literal::None),
}
}
}
impl DefineBucketStatement {
#[instrument(level = "trace", name = "DefineBucketStatement::compute", skip_all)]
pub(crate) async fn compute(
&self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
doc: Option<&CursorDoc>,
) -> Result<Value> {
ctx.is_allowed(opt, Action::Edit, ResourceKind::Bucket, Base::Db)?;
let name = expr_to_ident(stk, ctx, opt, doc, &self.name, "bucket name").await?;
if self.permissions.has_direct_write() {
bail!(Error::PermissionClauseNotReadonly {
kind: "bucket",
name: name.clone(),
});
}
let txn = ctx.tx();
let (ns, db) = ctx.get_ns_db_ids(opt).await?;
if let Some(bucket) = txn.get_db_bucket(ns, db, &name, None).await? {
match self.kind {
DefineKind::Default => {
if !opt.import {
bail!(Error::BuAlreadyExists {
value: bucket.name.to_string(),
});
}
}
DefineKind::Overwrite => {}
DefineKind::IfNotExists => {
return Ok(Value::None);
}
}
}
let backend = if let Some(ref url) = self.backend {
Some(
stk.run(|stk| url.compute(stk, ctx, opt, doc))
.await
.catch_return()?
.coerce_to::<String>()?,
)
} else {
None
};
if let Some(buckets) = ctx.get_buckets() {
buckets.new_backend(ns, db, &name, self.readonly, backend.as_deref()).await?;
} else {
bail!(Error::BucketUnavailable(name));
}
let key = crate::key::database::bu::new(ns, db, &name);
let comment = stk
.run(|stk| self.comment.compute(stk, ctx, opt, doc))
.await
.catch_return()?
.cast_to()?;
let ap = BucketDefinition {
id: None,
name: name.clone().into(),
backend: backend.map(|s| s.into()),
permissions: self.permissions.clone(),
readonly: self.readonly,
comment,
};
txn.set(&key, &ap).await?;
txn.clear_cache();
Ok(Value::None)
}
}