use crate::Role;
use crate::Workspace;
use crate::git_commands::run_git_output;
use crate::turso::{self, TxGuard};
use anyhow::Result;
use serde::Serialize;
use std::path::{Path, PathBuf};
use tracing::warn;
crate::define_store! {
pub static USER_STORE: UserStore,
db_name = "users",
schema = SCHEMA,
post_open = ensure_admin_user,
expect = "USER_STORE not initialized — call init_global() first",
}
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)
);";
crate::columns! {
USERS_COLUMNS [USERS] {
NAME => "name",
PERMISSIONS => "permissions",
SELECTED_WORKSPACE => "selected_workspace",
SELECTED_ROLE => "selected_role",
}
}
crate::columns! {
USER_CHANNEL_COLUMNS [UC] {
CHANNEL => "channel",
IDENTIFIER => "identifier",
REPLY_TARGET => "reply_target",
}
}
impl UserStore {
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_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 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(
&format!("SELECT {USER_CHANNEL_COLUMNS} 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>(COL_UC_CHANNEL)?,
identifier: row.get::<String>(COL_UC_IDENTIFIER)?,
reply_target: row.get::<Option<String>>(COL_UC_REPLY_TARGET)?,
});
}
Ok(bindings)
}
async fn user_record_from_row(&self, row: &turso::Row) -> Result<UserRecord> {
let name: String = row.get(COL_USERS_NAME)?;
Ok(UserRecord {
name: name.clone(),
permissions: row.get::<Option<String>>(COL_USERS_PERMISSIONS)?,
selected_workspace: row.get::<Option<String>>(COL_USERS_SELECTED_WORKSPACE)?,
selected_role: row.get::<Option<String>>(COL_USERS_SELECTED_ROLE)?,
channels: self.get_user_channels(&name).await.unwrap_or_default(),
})
}
pub async fn find_by_workspace(&self, workspace_name: &str) -> Result<Vec<UserRecord>> {
let rows = self
.conn
.query(
&format!(
"SELECT {USERS_COLUMNS} \
FROM users WHERE selected_workspace = ?1"
),
turso::params![workspace_name],
)
.await?;
let mut users = Vec::new();
for row in rows {
users.push(self.user_record_from_row(&row).await?);
}
Ok(users)
}
pub async fn list_users(&self) -> Result<Vec<UserRecord>> {
let rows = self
.conn
.query(
&format!("SELECT {USERS_COLUMNS} FROM users"),
turso::params![],
)
.await?;
let mut users = Vec::new();
for row in rows {
users.push(self.user_record_from_row(&row).await?);
}
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_column(&tx, name, "selected_role", role_name).await?;
upsert_user_column(&tx, name, "selected_workspace", workspace_name).await?;
upsert_user_column(&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_column(
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 run_git_output(&path, &["init", "-q"]).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 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 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 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_else(|e| {
tracing::warn!(error = %e, ?channel, ?identifier, "Failed to resolve user by channel");
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() {
crate::util::test::init_test_stores().await;
if let Some(store) = USER_STORE.get() {
store
.add_user("alice", Some("full"))
.await
.expect("failed to add alice to test USER_STORE");
store
.add_user("bob", None)
.await
.expect("failed to add bob to test USER_STORE");
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");
}
}
}