use crate::{ChannelOptions, Engine, Vars, Workflow};
use serde_json::{Value as JsonValue, json};
use std::fmt;
#[derive(Debug)]
pub enum Error {
NotFound(String),
Invalid(String),
Internal(String),
Unauthenticated(String),
Denied(String),
}
impl fmt::Display for Error {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Error::NotFound(msg) => write!(f, "not found action '{msg}'"),
Error::Invalid(msg) => f.write_str(msg),
Error::Internal(msg) => f.write_str(msg),
Error::Unauthenticated(msg) => write!(f, "unauthenticated: {msg}"),
Error::Denied(msg) => write!(f, "permission denied: {msg}"),
}
}
}
impl std::error::Error for Error {}
pub type Ret = std::result::Result<JsonValue, Error>;
fn value<T: serde::Serialize>(r: crate::Result<T>) -> Ret {
match r {
Ok(value) => serde_json::to_value(value).map_err(|e| Error::Internal(e.to_string())),
Err(crate::ActError::Unauthenticated(msg)) => Err(Error::Unauthenticated(msg)),
Err(crate::ActError::Denied(msg)) => Err(Error::Denied(msg)),
Err(e) => Err(Error::Internal(e.to_string())),
}
}
fn unit_ok(r: crate::Result<()>) -> Ret {
value(r.map(|_| json!(true)))
}
fn pop(options: &mut Vars, key: &str) -> std::result::Result<String, Error> {
options
.pop::<String>(key)
.ok_or_else(|| Error::Invalid(format!("{key} is required")))
}
fn deny(err: crate::AclError) -> Error {
match err {
crate::AclError::Unauthenticated(msg) => Error::Unauthenticated(msg),
crate::AclError::Denied(msg) => Error::Denied(msg),
}
}
pub async fn apply(engine: &Engine, name: &str, options: Vars) -> Ret {
let principal = engine.anonymous();
apply_as(engine, &principal, name, options).await
}
pub async fn apply_as(
engine: &Engine,
principal: &crate::Principal,
name: &str,
mut options: Vars,
) -> Ret {
if name == crate::acl::ACTION_WHOAMI {
return Ok(principal.to_value());
}
principal.check(name).map_err(deny)?;
let executor = engine.executor(principal);
match name {
"act:push" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().push(&pid, &tid, options).await)
}
"act:remove" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().remove(&pid, &tid, options).await)
}
"act:submit" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().submit(&pid, &tid, options).await)
}
"act:complete" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().complete(&pid, &tid, options).await)
}
"act:abort" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().abort(&pid, &tid, options).await)
}
"act:cancel" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().cancel(&pid, &tid, options).await)
}
"act:back" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().back(&pid, &tid, options).await)
}
"act:skip" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().skip(&pid, &tid, options).await)
}
"act:error" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.act().fail(&pid, &tid, options).await)
}
"model:ls" => {
let query = options
.get::<crate::query::Query>("query")
.unwrap_or_else(|| crate::query::Query::new().limit(100));
value(executor.model().list(&query).await)
}
"model:rm" => {
let id = pop(&mut options, "id")?;
value(executor.model().rm(&id).await)
}
"model:get" => {
let id = pop(&mut options, "id")?;
let fmt = options.get::<String>("fmt").unwrap_or("text".to_string());
value(executor.model().get(&id, &fmt).await)
}
"model:deploy" => {
let model_text = options
.get::<String>("model")
.ok_or_else(|| Error::Invalid("model is required".to_string()))?;
let mut model =
Workflow::from_yml(&model_text).map_err(|e| Error::Invalid(e.to_string()))?;
if let Some(mid) = options.get::<String>("mid") {
model.set_id(&mid);
}
let view = options.get::<JsonValue>("view");
value(executor.model().deploy(&model, view.as_ref()).await)
}
"pack:ls" => {
let query = options
.get::<crate::query::Query>("query")
.unwrap_or_else(|| crate::query::Query::new().limit(100));
value(executor.pack().list(&query).await)
}
"pack:get" => {
let id = pop(&mut options, "id")?;
value(executor.pack().get(&id).await)
}
"pack:publish" => {
let id = options
.get::<String>("id")
.ok_or_else(|| Error::Invalid("package 'id' is required".to_string()))?;
let pack_name = options.get::<String>("name").unwrap_or_default();
let desc = options.get::<String>("desc").unwrap_or_default();
let icon = options.get::<String>("icon").unwrap_or_default();
let doc = options.get::<String>("doc").unwrap_or_default();
let version = options.get::<String>("version").unwrap_or_default();
let schema = options
.get::<serde_json::Value>("schema")
.unwrap_or_default();
let pack_options = options
.get::<Option<serde_json::Value>>("options")
.unwrap_or_default();
let run_as = options.get::<String>("run_as").unwrap_or_default();
let resources = options
.get::<Vec<crate::ActResource>>("resources")
.unwrap_or_default();
let catalog = options.get::<String>("catalog").unwrap_or_default();
let pack = crate::data::Package {
id,
name: pack_name,
desc,
icon,
doc,
version,
schema: schema.to_string(),
options: pack_options.map(|v| v.to_string()),
run_as: std::str::FromStr::from_str(&run_as)
.map_err(|_err| Error::Invalid("package 'run_as' is invalid".to_string()))?,
resources: serde_json::to_string(&resources)
.map_err(|e| Error::Internal(e.to_string()))?,
catalog: std::str::FromStr::from_str(&catalog)
.map_err(|_err| Error::Invalid("package 'catalog' is invalid".to_string()))?,
..Default::default()
};
value(executor.pack().publish(&pack).await)
}
"pack:rm" => {
let id = pop(&mut options, "id")?;
value(executor.pack().rm(&id).await)
}
"proc:start" => {
let id = pop(&mut options, "id")?;
value(executor.proc().start(&id, options).await)
}
"proc:start_from_model" => {
let fmt = pop(&mut options, "fmt")?;
let model = pop(&mut options, "model")?;
value(
executor
.proc()
.start_from_model(&model, &fmt, options)
.await,
)
}
"proc:ls" => {
let query = options
.get::<crate::query::Query>("query")
.unwrap_or_else(|| crate::query::Query::new().limit(100));
value(executor.proc().list(&query).await)
}
"proc:get" => {
let pid = pop(&mut options, "pid")?;
value(executor.proc().get(&pid).await)
}
"task:ls" => {
let query = options
.get::<crate::query::Query>("query")
.unwrap_or_else(|| crate::query::Query::new().limit(100));
value(executor.task().list(&query).await)
}
"task:get" => {
let pid = pop(&mut options, "pid")?;
let tid = pop(&mut options, "tid")?;
value(executor.task().get(&pid, &tid).await)
}
"msg:ls" => {
let query = options
.get::<crate::query::Query>("query")
.unwrap_or_else(|| crate::query::Query::new().limit(100));
value(executor.msg().list(&query).await)
}
"msg:get" => {
let id = pop(&mut options, "id")?;
value(executor.msg().get(&id).await)
}
crate::acl::ACTION_SUBSCRIBE => {
let client_id = options.get::<String>("client_id").unwrap_or_default();
Ok(JsonValue::String(ChannelOptions::subscription_id(
principal.subject(),
&client_id,
)))
}
"msg:ack" => {
let id = pop(&mut options, "id")?;
value(executor.msg().ack(&id).await)
}
"msg:redo" => match options.get::<String>("id") {
Some(id) => value(executor.msg().redeliver(&id).await),
None => value(executor.msg().redo().await),
},
"msg:clear" => {
if let Some(id) = options.get::<String>("id") {
value(executor.msg().clear_delivery(&id).await)
} else {
let pid = options.get::<String>("pid");
value(executor.msg().clear(pid).await)
}
}
"msg:rm" => {
let id = pop(&mut options, "id")?;
value(executor.msg().rm(&id).await)
}
"msg:unsub" => {
let client_id = pop(&mut options, "client_id")?;
let chan_id = ChannelOptions::subscription_id(principal.subject(), &client_id);
value(executor.msg().unsub(&chan_id).await)
}
"evt:ls" => {
let query = options
.get::<crate::query::Query>("query")
.unwrap_or_else(|| crate::query::Query::new().limit(100));
value(executor.evt().list(&query).await)
}
"evt:get" => {
let id = pop(&mut options, "id")?;
value(executor.evt().get(&id).await)
}
"evt:start" => {
let id = pop(&mut options, "id")?;
let params = options.get::<JsonValue>("params").unwrap_or_default();
value(executor.evt().start(&id, ¶ms).await)
}
"snap:upsert" => {
let target = pop(&mut options, "name")?;
let scope = options.get::<String>("scope").unwrap_or_default();
principal.check_scope(&target, &scope).map_err(deny)?;
let rev = options
.get::<u64>("rev")
.ok_or_else(|| Error::Invalid("rev is required".to_string()))?;
let data = options
.pop::<Vars>("data")
.ok_or_else(|| Error::Invalid("data is required".to_string()))?;
unit_ok(engine.snapshot().upsert(&target, &scope, rev, data))
}
"snap:remove" => {
let target = pop(&mut options, "name")?;
let scope = options.get::<String>("scope").unwrap_or_default();
principal.check_scope(&target, &scope).map_err(deny)?;
unit_ok(engine.snapshot().remove(&target, &scope))
}
"snap:get" => {
let target = pop(&mut options, "name")?;
let scope = options.get::<String>("scope").unwrap_or_default();
principal.check_scope(&target, &scope).map_err(deny)?;
match engine.snapshot().read(&target, &scope) {
Some(entry) => {
let data = serde_json::to_value(entry.data)
.map_err(|e| Error::Internal(e.to_string()))?;
Ok(json!({
"scope": scope,
"rev": entry.rev,
"timestamp": entry.timestamp,
"data": data,
}))
}
None => Ok(JsonValue::Null),
}
}
"snap:ls" => {
let target = pop(&mut options, "name")?;
let rows: Vec<JsonValue> = engine
.snapshot()
.list(&target)
.into_iter()
.filter(|(scope, _)| principal.check_scope(&target, scope).is_ok())
.map(|(scope, entry)| {
let data = serde_json::to_value(entry.data)
.map_err(|e| Error::Internal(e.to_string()))?;
Ok(json!({
"scope": scope,
"rev": entry.rev,
"timestamp": entry.timestamp,
"data": data,
}))
})
.collect::<std::result::Result<_, Error>>()?;
Ok(JsonValue::Array(rows))
}
_ => Err(Error::NotFound(name.to_string())),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::utils::consts;
async fn open_engine() -> crate::Engine {
let config = crate::Config {
data: Default::default(),
table: toml::from_str::<toml::Table>("[acl]\nenabled = false\n").unwrap(),
};
crate::Engine::builder()
.set_config(&config)
.start()
.await
.unwrap()
}
#[tokio::test]
async fn snapshot_upsert_remove_roundtrip() {
let engine = open_engine().await;
let payload = Vars::new()
.with("name", "profile")
.with("scope", "u1")
.with("rev", 7u64)
.with("data", Vars::new().with("val", "x"));
let ret = apply(&engine, "snap:upsert", payload).await.unwrap();
assert_eq!(ret, json!(true));
let entry = engine.snapshot().read("profile", "u1").unwrap();
assert_eq!(entry.rev, 7);
assert_eq!(entry.data.get::<String>("val").unwrap(), "x".to_string());
let ret = apply(
&engine,
"snap:remove",
Vars::new().with("name", "profile").with("scope", "u1"),
)
.await
.unwrap();
assert_eq!(ret, json!(true));
assert!(engine.snapshot().read("profile", "u1").is_none());
}
#[tokio::test]
async fn unknown_action_is_not_found() {
let engine = open_engine().await;
let err = apply(&engine, "no:such", Vars::new()).await.unwrap_err();
assert!(matches!(err, Error::NotFound(_)));
assert_eq!(err.to_string(), "not found action 'no:such'");
}
#[tokio::test]
async fn missing_payload_is_invalid() {
let engine = open_engine().await;
let err = apply(&engine, "snap:upsert", Vars::new())
.await
.unwrap_err();
assert!(matches!(err, Error::Invalid(_)));
assert!(err.to_string().contains("name is required"));
}
#[tokio::test]
async fn snapshot_query_roundtrip() {
let engine = open_engine().await;
for (scope, val) in [("u1", 1), ("u2", 2)] {
let payload = Vars::new()
.with("name", "profile")
.with("scope", scope)
.with("rev", val)
.with("data", Vars::new().with("val", val));
let ret = apply(&engine, "snap:upsert", payload).await.unwrap();
assert_eq!(ret, json!(true));
}
let ret = apply(
&engine,
"snap:get",
Vars::new().with("name", "profile").with("scope", "u1"),
)
.await
.unwrap();
assert_eq!(ret["scope"], "u1");
assert_eq!(ret["rev"], 1);
assert_eq!(ret["data"]["val"], 1);
let ret = apply(
&engine,
"snap:get",
Vars::new().with("name", "profile").with("scope", "nope"),
)
.await
.unwrap();
assert_eq!(ret, JsonValue::Null);
let ret = apply(&engine, "snap:ls", Vars::new().with("name", "profile"))
.await
.unwrap();
let rows = ret.as_array().unwrap();
assert_eq!(rows.len(), 2);
assert!(
rows.iter()
.any(|r| r["scope"] == "u1" && r["data"]["val"] == 1)
);
assert!(
rows.iter()
.any(|r| r["scope"] == "u2" && r["data"]["val"] == 2)
);
}
#[tokio::test]
async fn snapshot_query_unknown_target() {
let engine = open_engine().await;
let ret = apply(&engine, "snap:ls", Vars::new().with("name", "none"))
.await
.unwrap();
assert_eq!(ret, JsonValue::Array(vec![]));
}
#[tokio::test]
async fn snapshot_remove_unknown_target_is_internal() {
let engine = open_engine().await;
let err = apply(
&engine,
"snap:remove",
Vars::new().with("name", "none").with("scope", "u1"),
)
.await
.unwrap_err();
assert!(matches!(err, Error::Internal(_)), "got: {err}");
assert!(err.to_string().contains("none"), "got: {err}");
}
const MULTI_TENANT: &str = r#"
[acl]
[[acl.role]]
name = "u1"
tokens = ["token-u1"]
allow = ["model:deploy", "proc:start", "snap:get", "snap:ls", "snap:upsert"]
snapshot = { secrets = ["u1"], profile = ["u1"] }
[[acl.role]]
name = "u2"
tokens = ["token-u2"]
allow = ["model:deploy", "proc:start", "snap:get", "snap:ls", "snap:upsert"]
snapshot = { secrets = ["u2"], profile = ["u2"] }
"#;
async fn acl_engine(text: &str) -> crate::Engine {
let config = crate::Config {
data: Default::default(),
table: toml::from_str::<toml::Table>(text).unwrap(),
};
crate::Engine::builder()
.set_config(&config)
.start()
.await
.unwrap()
}
#[tokio::test]
async fn an_enabled_acl_refuses_the_anonymous_caller() {
let engine = acl_engine(MULTI_TENANT).await;
let err = apply(&engine, "model:ls", Vars::new()).await.unwrap_err();
assert!(matches!(err, Error::Unauthenticated(_)), "got: {err}");
let err = engine.acl().authenticate(Some("bogus")).unwrap_err();
assert!(matches!(err, crate::AclError::Unauthenticated(_)));
}
#[tokio::test]
async fn a_role_runs_only_the_actions_it_allows() {
let engine = acl_engine(MULTI_TENANT).await;
let principal = engine.acl().authenticate(Some("token-u1")).unwrap();
apply_as(&engine, &principal, "model:ls", Vars::new())
.await
.unwrap_err();
let who = apply_as(&engine, &principal, crate::acl::ACTION_WHOAMI, Vars::new())
.await
.unwrap();
assert_eq!(who["subject"], "u1");
assert_eq!(who["unrestricted"], false);
}
#[tokio::test]
async fn a_snapshot_scope_belongs_to_one_subject_only() {
let engine = acl_engine(MULTI_TENANT).await;
let u1 = engine.acl().authenticate(Some("token-u1")).unwrap();
let u2 = engine.acl().authenticate(Some("token-u2")).unwrap();
for (principal, scope, val) in [(&u1, "u1", 1), (&u2, "u2", 2)] {
let payload = Vars::new()
.with("name", "profile")
.with("scope", scope)
.with("rev", 1u64)
.with("data", Vars::new().with("val", val));
assert_eq!(
apply_as(&engine, principal, "snap:upsert", payload)
.await
.unwrap(),
json!(true)
);
}
let payload = Vars::new()
.with("name", "profile")
.with("scope", "u2")
.with("rev", 9u64)
.with("data", Vars::new().with("val", 9));
let err = apply_as(&engine, &u1, "snap:upsert", payload)
.await
.unwrap_err();
assert!(matches!(err, Error::Denied(_)), "got: {err}");
let err = apply_as(
&engine,
&u1,
"snap:get",
Vars::new().with("name", "profile").with("scope", "u2"),
)
.await
.unwrap_err();
assert!(matches!(err, Error::Denied(_)), "got: {err}");
let rows = apply_as(&engine, &u1, "snap:ls", Vars::new().with("name", "profile"))
.await
.unwrap();
assert_eq!(rows.as_array().unwrap().len(), 1);
assert_eq!(rows[0]["scope"], "u1");
assert_eq!(rows[0]["data"]["val"], 1);
}
const MSG_TENANT: &str = r#"
[acl]
[[acl.role]]
name = "u1"
tokens = ["token-u1"]
allow = ["model:deploy", "proc:start", "msg:ack", "msg:unsub"]
[[acl.role]]
name = "u2"
tokens = ["token-u2"]
allow = ["model:deploy", "proc:start", "msg:ack", "msg:unsub"]
"#;
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn an_ack_follows_the_grant_not_the_process_owner() {
let engine = acl_engine(MSG_TENANT).await;
let u1 = engine.acl().authenticate(Some("token-u1")).unwrap();
let u2 = engine.acl().authenticate(Some("token-u2")).unwrap();
let workflow = crate::Workflow::from_yml(
r#"
id: ack_owner
ver: 0.1.0
steps:
- id: step1
uses: acts.core.irq
"#,
)
.unwrap();
apply_as(
&engine,
&u1,
"model:deploy",
Vars::new().with("model", workflow.to_yml().unwrap()),
)
.await
.unwrap();
let delivery = engine.signal(String::new());
let d = delivery.clone();
let chan = engine.channel_with_options(&ChannelOptions {
id: ChannelOptions::subscription_id("u1", "client-1"),
ack: true,
..Default::default()
});
chan.on_message(move |e| {
let d = d.clone();
async move {
if let Some(id) = &e.delivery_id {
d.update(|data| data.clone_from(id));
d.close();
}
}
});
apply_as(
&engine,
&u1,
"proc:start",
Vars::new().with("id", "ack_owner"),
)
.await
.unwrap();
let delivery_id = delivery.recv().await;
assert!(!delivery_id.is_empty());
apply_as(
&engine,
&u2,
"msg:ack",
Vars::new().with("id", delivery_id.clone()),
)
.await
.unwrap();
assert_eq!(
engine
.runtime()
.cache()
.store()
.deliveries()
.find(&delivery_id)
.await
.unwrap()
.status,
crate::data::DeliveryStatus::Acked
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn an_ack_without_the_grant_is_refused() {
let engine = acl_engine(
r#"
[acl]
[[acl.role]]
name = "starter"
tokens = ["token-starter"]
allow = ["model:deploy", "proc:start", "msg:ls"]
"#,
)
.await;
let starter = engine.acl().authenticate(Some("token-starter")).unwrap();
let delivery = engine.signal(String::new());
let d = delivery.clone();
let chan = engine.channel_with_options(&ChannelOptions {
id: "ack-grant-client".to_string(),
ack: true,
..Default::default()
});
chan.on_message(move |e| {
let d = d.clone();
async move {
if let Some(id) = &e.delivery_id {
d.update(|data| data.clone_from(id));
d.close();
}
}
});
let workflow = crate::Workflow::from_yml(
r#"
id: ack_grant
ver: 0.1.0
steps:
- id: step1
uses: acts.core.irq
"#,
)
.unwrap();
apply_as(
&engine,
&starter,
"model:deploy",
Vars::new().with("model", workflow.to_yml().unwrap()),
)
.await
.unwrap();
apply_as(
&engine,
&starter,
"proc:start",
Vars::new().with("id", "ack_grant"),
)
.await
.unwrap();
let delivery_id = delivery.recv().await;
assert!(!delivery_id.is_empty());
let err = apply_as(
&engine,
&starter,
"msg:ack",
Vars::new().with("id", delivery_id),
)
.await
.unwrap_err();
assert!(matches!(err, Error::Denied(_)), "got: {err}");
}
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn an_unsub_reaches_only_the_callers_own_namespace() {
let engine = acl_engine(MSG_TENANT).await;
let u1 = engine.acl().authenticate(Some("token-u1")).unwrap();
let u2 = engine.acl().authenticate(Some("token-u2")).unwrap();
let workflow = crate::Workflow::from_yml(
r#"
id: unsub_scope
ver: 0.1.0
steps:
- id: step1
uses: acts.core.irq
"#,
)
.unwrap();
apply_as(
&engine,
&u1,
"model:deploy",
Vars::new().with("model", workflow.to_yml().unwrap()),
)
.await
.unwrap();
let received = std::sync::Arc::new(parking_lot::Mutex::new(std::collections::HashSet::<
String,
>::new()));
let seen = received.clone();
let chan = engine.channel_with_options(&ChannelOptions {
id: ChannelOptions::subscription_id("u1", "shared-client"),
ack: true,
..Default::default()
});
chan.on_message(move |e| {
let seen = seen.clone();
async move {
seen.lock().insert(e.id.clone());
}
});
apply_as(
&engine,
&u2,
"msg:unsub",
Vars::new().with("client_id", "shared-client"),
)
.await
.unwrap();
apply_as(
&engine,
&u2,
"msg:unsub",
Vars::new().with("client_id", "shared-client"),
)
.await
.unwrap();
apply_as(
&engine,
&u1,
"proc:start",
Vars::new().with("id", "unsub_scope"),
)
.await
.unwrap();
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
while received.lock().is_empty() && tokio::time::Instant::now() < deadline {
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
let after_foreign_unsub = received.lock().clone();
assert!(
!after_foreign_unsub.is_empty(),
"another subject's unsub silenced u1's channel"
);
apply_as(
&engine,
&u1,
"msg:unsub",
Vars::new().with("client_id", "shared-client"),
)
.await
.unwrap();
apply_as(
&engine,
&u1,
"proc:start",
Vars::new().with("id", "unsub_scope"),
)
.await
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
assert_eq!(
*received.lock(),
after_foreign_unsub,
"the subscriber's own unsub did not reach its channel"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn a_run_cannot_seal_another_subjects_scope() {
let engine = acl_engine(MULTI_TENANT).await;
engine
.add_snapshot(
"secrets",
crate::SnapshotOptions {
scope: vec!["uid".to_string()],
..crate::SnapshotOptions::per_proc()
},
)
.unwrap();
let u1 = engine.acl().authenticate(Some("token-u1")).unwrap();
let u2 = engine.acl().authenticate(Some("token-u2")).unwrap();
apply_as(
&engine,
&u2,
"snap:upsert",
Vars::new()
.with("name", "secrets")
.with("scope", "u2")
.with("rev", 1u64)
.with("data", Vars::new().with("TOKEN", "u2-secret")),
)
.await
.unwrap();
let workflow = crate::Workflow::from_yml(
r#"
id: acl_seal
ver: 0.1.0
exposes:
- name: leaked
steps:
- name: read the secret
uses: acts.transform.code
params: |
return { leaked: secrets.TOKEN };
"#,
)
.unwrap();
apply_as(
&engine,
&u1,
"model:deploy",
Vars::new().with("model", workflow.to_yml().unwrap()),
)
.await
.unwrap();
let sig = engine.signal(String::new());
let s = sig.clone();
engine.channel().on_error(move |e| {
let s = s.clone();
async move {
let err = e
.inputs
.get::<String>(crate::utils::consts::ACT_ERR_MESSAGE)
.unwrap_or_default();
s.update(|data| data.clone_from(&err));
s.close();
}
});
apply_as(
&engine,
&u1,
"proc:start",
Vars::new().with("id", "acl_seal").with("uid", "u2"),
)
.await
.unwrap();
let err = sig.recv().await;
assert!(err.contains("not owned by subject 'u1'"), "got: {err}");
engine
.snapshot()
.upsert("secrets", "u1", 1, Vars::new().with("TOKEN", "u1-secret"))
.unwrap();
let sig = engine.signal(String::new());
let s = sig.clone();
engine.channel().on_complete(move |e| {
let s = s.clone();
async move {
s.update(|data| data.clone_from(&e.pid));
s.close();
}
});
let pid = apply_as(
&engine,
&u1,
"proc:start",
Vars::new()
.with("id", "acl_seal")
.with("pid", "acl_seal_owned")
.with("uid", "u1"),
)
.await
.unwrap()
.as_str()
.unwrap()
.to_string();
assert_eq!(sig.recv().await, pid);
let proc = engine.runtime().proc(&pid).await.unwrap().unwrap();
assert_eq!(
proc.root()
.unwrap()
.sealed("secrets")
.unwrap()
.get::<String>("TOKEN")
.unwrap(),
"u1-secret"
);
}
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn a_run_is_confined_to_its_own_workdir() {
let root = std::env::temp_dir().join(format!("acts_workdir_{}", crate::utils::longid()));
let engine = acl_engine(&format!(
r#"
[acl]
workdir = '{}'
[[acl.role]]
name = "u1"
tokens = ["token-u1"]
allow = ["model:deploy", "proc:start"]
"#,
root.display()
))
.await;
assert!(!root.exists());
let workflow = crate::Workflow::from_yml(
r#"
id: workdir_run
ver: 0.1.0
steps:
- name: finish
uses: acts.transform.code
params: |
return { dir: $env.WORK_DIR };
"#,
)
.unwrap();
let u1 = engine.acl().authenticate(Some("token-u1")).unwrap();
assert_eq!(u1.workdir_root(), Some(root.as_path()));
let sig = engine.signal(());
let done = sig.clone();
engine.channel().on_complete(move |_| {
let done = done.clone();
async move { done.close() }
});
apply_as(
&engine,
&u1,
"model:deploy",
Vars::new().with("model", workflow.to_yml().unwrap()),
)
.await
.unwrap();
let pid = apply_as(
&engine,
&u1,
"proc:start",
Vars::new().with("id", "workdir_run").with("pid", "run1"),
)
.await
.unwrap()
.as_str()
.unwrap()
.to_string();
assert_eq!(pid, "run1");
let dir = root.join("run1");
assert!(dir.is_dir(), "workdir {} was not created", dir.display());
let proc = engine.runtime().proc(&pid).await.unwrap().unwrap();
assert_eq!(proc.workdir(), Some(dir.clone()));
assert!(proc.inputs().get::<String>(consts::PROC_WORKDIR).is_none());
sig.recv().await;
assert_eq!(
proc.task_by_uses(crate::utils::test::USES_CODE)
.first()
.unwrap()
.outputs()
.get::<String>("dir")
.unwrap(),
dir.display().to_string()
);
for _ in 0..150 {
if !dir.exists() {
break;
}
let _ = engine.runtime().cache().sweep_removable().await;
tokio::time::sleep(std::time::Duration::from_millis(20)).await;
}
assert!(
!dir.exists(),
"the workdir {} must be removed with the process's rows",
dir.display()
);
std::fs::remove_dir_all(&root).ok();
}
#[tokio::test(flavor = "multi_thread")]
#[serial_test::serial]
async fn a_workdir_refuses_a_pid_that_is_not_one_directory() {
let root = std::env::temp_dir().join(format!("acts_workdir_{}", crate::utils::longid()));
let engine = acl_engine(&format!(
r#"
[acl]
workdir = '{}'
[[acl.role]]
name = "u1"
tokens = ["token-u1"]
allow = ["proc:start_from_model"]
"#,
root.display()
))
.await;
let u1 = engine.acl().authenticate(Some("token-u1")).unwrap();
let model = Vars::new()
.with("model", "id: w\nver: 0.1.0\nsteps:\n - name: s\n")
.with("fmt", "yml");
for pid in ["..", ".", "a/b", "a\\b", "a:b"] {
let err = apply_as(
&engine,
&u1,
"proc:start_from_model",
model.clone().with("pid", pid),
)
.await
.expect_err(pid);
assert!(
err.to_string().contains("cannot be used as a workdir name"),
"pid {pid:?} must be refused, got: {err}"
);
}
assert!(!root.exists() || std::fs::read_dir(&root).unwrap().next().is_none());
std::fs::remove_dir_all(&root).ok();
}
}