use crate::Role;
use crate::Workspace;
use crate::global_store;
use crate::turso::{self, Connection, TxGuard};
use anyhow::Result;
use serde::Serialize;
use std::path::{Path, PathBuf};
use tracing::warn;
global_store! {
pub static USER_STORE: UserStorage,
constructor = UserStorage::open,
}
#[derive(Clone, Debug)]
pub struct UserStorage {
pub(crate) conn: Connection,
}
const SCHEMA: &str = "\
CREATE TABLE IF NOT EXISTS users (
name TEXT PRIMARY KEY,
permissions TEXT,
selected_workspace TEXT,
selected_role TEXT
);
CREATE TABLE IF NOT EXISTS user_channels (
user_name TEXT NOT NULL REFERENCES users(name),
channel TEXT NOT NULL,
identifier TEXT NOT NULL,
reply_target TEXT,
UNIQUE(channel, identifier)
);";
impl UserStorage {
pub async fn open(root: &Path) -> Result<Self> {
let db_path = root.join("db/users.db");
let conn = turso::open_with_schema(&db_path, SCHEMA).await?;
let this = Self { conn };
this.ensure_admin_user().await?;
Ok(this)
}
async fn ensure_admin_user(&self) -> Result<()> {
let rows = self
.conn
.query("SELECT 1 FROM users WHERE name = 'admin'", turso::params![])
.await?;
if rows.is_empty() {
self.conn
.execute(
"INSERT INTO users (name, permissions) VALUES ('admin', 'full')",
turso::params![],
)
.await?;
init_personal_workspace_dir("admin").await;
}
Ok(())
}
pub async fn add_user(&self, name: &str, permissions: Option<&str>) -> Result<()> {
self.conn
.execute(
"INSERT OR IGNORE INTO users (name, permissions) VALUES (?1, ?2)",
turso::params![name, permissions],
)
.await?;
init_personal_workspace_dir(name).await;
Ok(())
}
pub async fn delete_user(&self, name: &str) -> Result<()> {
let tx = self.conn.begin_tx().await?;
tx.execute(
"DELETE FROM user_channels WHERE user_name = ?1",
turso::params![name],
)
.await?;
tx.execute("DELETE FROM users WHERE name = ?1", turso::params![name])
.await?;
tx.commit().await?;
Ok(())
}
pub async fn get_permissions(&self, name: &str) -> Result<Option<String>> {
let rows = self
.conn
.query(
"SELECT permissions FROM users WHERE name = ?1",
turso::params![name],
)
.await?;
match rows.into_iter().next() {
Some(row) => Ok(row.get::<Option<String>>(0)?),
None => Ok(None),
}
}
pub async fn set_workspace(&self, user_name: &str, workspace_name: Option<&str>) -> Result<()> {
self.conn
.execute(
"INSERT INTO users (name, selected_workspace) VALUES (?1, ?2) \
ON CONFLICT(name) DO UPDATE SET selected_workspace = excluded.selected_workspace",
turso::params![user_name, workspace_name],
)
.await?;
Ok(())
}
pub async fn get_selected_workspace_name(&self, user_name: &str) -> Result<Option<String>> {
let rows = self
.conn
.query(
"SELECT selected_workspace FROM users WHERE name = ?1",
turso::params![user_name],
)
.await?;
match rows.into_iter().next() {
Some(row) => Ok(row.get::<Option<String>>(0)?),
None => Ok(None),
}
}
pub async fn set_active_role(&self, user_name: &str, role_name: &str) -> Result<()> {
self.conn
.execute(
"INSERT INTO users (name, selected_role) VALUES (?1, ?2) \
ON CONFLICT(name) DO UPDATE SET selected_role = excluded.selected_role",
turso::params![user_name, role_name],
)
.await?;
Ok(())
}
pub async fn get_active_role(&self, user_name: &str) -> Result<Option<String>> {
let rows = self
.conn
.query(
"SELECT selected_role FROM users WHERE name = ?1",
turso::params![user_name],
)
.await?;
match rows.into_iter().next() {
Some(row) => Ok(row.get::<Option<String>>(0)?),
None => Ok(None),
}
}
pub async fn bind_channel(
&self,
user_name: &str,
channel: &str,
identifier: &str,
) -> Result<()> {
self.conn
.execute(
"INSERT OR REPLACE INTO user_channels (user_name, channel, identifier) \
VALUES (?1, ?2, ?3)",
turso::params![user_name, channel, identifier],
)
.await?;
Ok(())
}
pub async fn unbind_channel(
&self,
user_name: &str,
channel: &str,
identifier: &str,
) -> Result<()> {
self.conn
.execute(
"DELETE FROM user_channels WHERE user_name = ?1 AND channel = ?2 AND identifier = ?3",
turso::params![user_name, channel, identifier],
)
.await?;
Ok(())
}
pub async fn update_channel_contact(
&self,
channel: &str,
identifier: &str,
reply_target: &str,
) -> Result<()> {
self.conn
.execute(
"UPDATE user_channels SET reply_target = ?1 \
WHERE channel = ?2 AND identifier = ?3",
turso::params![reply_target, channel, identifier],
)
.await?;
Ok(())
}
pub async fn resolve_user_by_channel(
&self,
channel: &str,
identifier: &str,
) -> Result<Option<String>> {
let rows = self
.conn
.query(
"SELECT user_name FROM user_channels WHERE channel = ?1 AND identifier = ?2",
turso::params![channel, identifier],
)
.await?;
match rows.into_iter().next() {
Some(row) => Ok(Some(row.get::<String>(0)?)),
None => Ok(None),
}
}
pub async fn get_user_channels(&self, user_name: &str) -> Result<Vec<ChannelBinding>> {
let rows = self
.conn
.query(
"SELECT channel, identifier, reply_target FROM user_channels WHERE user_name = ?1",
turso::params![user_name],
)
.await?;
let mut bindings = Vec::new();
for row in rows {
bindings.push(ChannelBinding {
channel: row.get::<String>(0)?,
identifier: row.get::<String>(1)?,
reply_target: row.get::<Option<String>>(2)?,
});
}
Ok(bindings)
}
pub async fn find_by_workspace(&self, workspace_name: &str) -> Result<Vec<UserRecord>> {
let rows = self
.conn
.query(
"SELECT name, permissions, selected_workspace, selected_role \
FROM users WHERE selected_workspace = ?1",
turso::params![workspace_name],
)
.await?;
let mut users = Vec::new();
for row in rows {
let name: String = row.get(0)?;
users.push(UserRecord {
name: name.clone(),
permissions: row.get::<Option<String>>(1)?,
selected_workspace: row.get::<Option<String>>(2)?,
selected_role: row.get::<Option<String>>(3)?,
channels: self.get_user_channels(&name).await.unwrap_or_default(),
});
}
Ok(users)
}
pub async fn list_users(&self) -> Result<Vec<UserRecord>> {
let rows = self
.conn
.query(
"SELECT name, permissions, selected_workspace, selected_role FROM users",
turso::params![],
)
.await?;
let mut users = Vec::new();
for row in rows {
let name: String = row.get(0)?;
users.push(UserRecord {
name: name.clone(),
permissions: row.get::<Option<String>>(1)?,
selected_workspace: row.get::<Option<String>>(2)?,
selected_role: row.get::<Option<String>>(3)?,
channels: self.get_user_channels(&name).await.unwrap_or_default(),
});
}
Ok(users)
}
pub async fn update_user(
&self,
name: &str,
role_name: FieldUpdate<'_>,
workspace_name: FieldUpdate<'_>,
permissions: FieldUpdate<'_>,
) -> Result<()> {
let tx = self.conn.begin_tx().await?;
upsert_user_field(&tx, name, "selected_role", role_name).await?;
upsert_user_field(&tx, name, "selected_workspace", workspace_name).await?;
upsert_user_field(&tx, name, "permissions", permissions).await?;
tx.commit().await?;
Ok(())
}
}
#[derive(Debug, Clone, Copy)]
pub enum FieldUpdate<'a> {
Unchanged,
Clear,
Set(&'a str),
}
async fn upsert_user_field(
tx: &TxGuard<'_>,
name: &str,
field: &str,
value: FieldUpdate<'_>,
) -> Result<()> {
let val: Option<&str> = match value {
FieldUpdate::Unchanged => return Ok(()),
FieldUpdate::Clear => None,
FieldUpdate::Set(v) => Some(v),
};
let sql = format!(
"INSERT INTO users (name, {field}) VALUES (?1, ?2) \
ON CONFLICT(name) DO UPDATE SET {field} = excluded.{field}"
);
tx.execute(&sql, turso::params![name, val]).await?;
Ok(())
}
#[derive(Debug, Clone, Serialize)]
pub struct UserRecord {
pub name: String,
pub permissions: Option<String>,
pub selected_workspace: Option<String>,
pub selected_role: Option<String>,
pub channels: Vec<ChannelBinding>,
}
#[derive(Debug, Clone, Serialize)]
pub struct ChannelBinding {
pub channel: String,
pub identifier: String,
pub reply_target: Option<String>,
}
#[must_use]
pub fn personal_workspace_path(user_name: &str) -> PathBuf {
let storage_root = crate::config::default_config_dir()
.unwrap_or_else(|_| std::env::temp_dir().join("mahbot_test_userspaces"));
storage_root.join("userspaces").join(user_name)
}
async fn init_personal_workspace_dir(name: &str) {
let path = personal_workspace_path(name);
if let Err(e) = tokio::fs::create_dir_all(&path).await {
warn!(
path = %path.display(),
error = %e,
"Failed to create personal workspace directory"
);
}
match tokio::process::Command::new("git")
.arg("init")
.arg("-q")
.current_dir(&path)
.output()
.await
{
Ok(o) if o.status.success() => {}
Ok(_) => warn!(
path = %path.display(),
"git init failed for personal workspace (git may not be installed)"
),
Err(e) => warn!(
path = %path.display(),
error = %e,
"git init failed for personal workspace"
),
}
}
pub async fn set_workspace(user_name: &str, name: &str) -> Result<Workspace> {
let ws = match crate::workspace::get_by_name(name).await {
Ok(Some(ws)) => ws,
Ok(None) => anyhow::bail!("Workspace '{name}' not found"),
Err(e) => anyhow::bail!("Database error looking up workspace '{name}': {e}"),
};
let user_store = store();
user_store.set_workspace(user_name, Some(&ws.name)).await?;
Ok(ws)
}
pub async fn get_raw_selected_workspace(user_name: &str) -> Result<Option<String>> {
store().get_selected_workspace_name(user_name).await
}
pub async fn get_workspace(user_name: &str) -> Result<Option<Workspace>> {
let s = store();
let selected = s.get_selected_workspace_name(user_name).await?;
if let Some(ws_name) = selected {
crate::workspace::get_by_name(&ws_name).await
} else {
let path = personal_workspace_path(user_name);
Ok(Some(personal_workspace_struct(user_name, &path)))
}
}
#[must_use]
pub fn personal_workspace_struct(user_name: &str, path: &Path) -> Workspace {
let mut ws = Workspace::from_path(path);
ws.name = format!("personal:{user_name}");
ws.status = "ready".to_string();
ws.maintainer_debounce_mins = 240;
let now = turso::now();
ws.created_at.clone_from(&now);
ws.updated_at = now;
ws
}
pub async fn set_active_role(user_name: &str, role_name: &str) -> Result<()> {
store().set_active_role(user_name, role_name).await
}
pub async fn get_active_role(user_name: &str) -> Result<Option<String>> {
store().get_active_role(user_name).await
}
pub async fn resolve_active_role(user_name: &str) -> Role {
match get_active_role(user_name).await {
Ok(Some(name)) => name.parse::<Role>().unwrap_or(Role::Analyst),
_ => Role::Analyst,
}
}
pub async fn is_full_user(user_name: &str) -> bool {
let Some(store) = USER_STORE.get() else {
return false;
};
matches!(
store.get_permissions(user_name).await,
Ok(Some(p)) if p == "full"
)
}
pub async fn resolve_user_by_channel(channel: &str, identifier: &str) -> Option<String> {
let store = USER_STORE.get()?;
store
.resolve_user_by_channel(channel, identifier)
.await
.unwrap_or(None)
}
pub async fn update_channel_contact(
channel: &str,
identifier: &str,
reply_target: &str,
) -> Result<()> {
store()
.update_channel_contact(channel, identifier, reply_target)
.await
}
#[must_use]
pub fn is_personal_workspace(workspace_name: &str) -> bool {
workspace_name.starts_with("personal:")
}
#[cfg(test)]
pub(crate) mod test_util {
use super::*;
pub(crate) async fn init_test_store() {
if USER_STORE.get().is_some() {
return;
}
let dir = tempfile::TempDir::new().expect("Failed to create temp dir for test user store");
let store = UserStorage::open(dir.path())
.await
.expect("Failed to open test user store");
store
.add_user("alice", Some("full"))
.await
.expect("Failed to add alice");
store
.add_user("bob", None)
.await
.expect("Failed to add bob");
store
.bind_channel("alice", "telegram", "alice")
.await
.expect("Failed to bind alice telegram");
store
.bind_channel("bob", "telegram", "bob")
.await
.expect("Failed to bind bob telegram");
let _ = Box::leak(Box::new(dir));
let _ = USER_STORE.set(store);
}
}