use std::{collections::HashSet, path::Path, sync::Mutex};
use chrono::Utc;
use kcode_kweb_db::NodeId;
use kcode_tg_kennedy_bot::{AddUserOutcome, IdentityObservation, IdentitySink, WhitelistSnapshot};
use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
use serde::{Deserialize, Serialize};
const IDENTITY_MIGRATION: &str = include_str!("../migrations/001_initial.sql");
const ANONYMOUS_GROUP_USER_ID: i64 = 1_087_968_824;
const ANONYMOUS_GROUP_HANDLE: &str = "groupanonymousbot";
pub struct Directory {
database: Mutex<Connection>,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct User {
pub handle: String,
pub telegram_user_id: Option<i64>,
pub current_username: Option<String>,
pub display_name: Option<String>,
pub root_node_id: Option<String>,
pub root_ready: bool,
pub can_add_users: bool,
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct Group {
pub group_id: String,
pub root_node_id: Option<String>,
pub root_ready: bool,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ErrorKind {
InvalidInput,
NotFound,
Conflict,
Storage,
}
#[derive(Debug, thiserror::Error)]
#[error("{message}")]
pub struct Error {
kind: ErrorKind,
message: String,
}
impl Error {
pub fn kind(&self) -> ErrorKind {
self.kind
}
pub fn message(&self) -> &str {
&self.message
}
fn invalid(message: impl Into<String>) -> Self {
Self {
kind: ErrorKind::InvalidInput,
message: message.into(),
}
}
fn not_found() -> Self {
Self {
kind: ErrorKind::NotFound,
message: "Telegram directory entry not found.".into(),
}
}
fn conflict(message: impl Into<String>) -> Self {
Self {
kind: ErrorKind::Conflict,
message: message.into(),
}
}
fn storage(error: impl std::fmt::Display) -> Self {
Self {
kind: ErrorKind::Storage,
message: error.to_string(),
}
}
}
pub type Result<T> = std::result::Result<T, Error>;
impl Directory {
pub fn open(path: &Path, bootstrap_handle: &str) -> Result<Self> {
let connection = Connection::open(path)
.map_err(|error| Error::storage(format!("opening {}: {error}", path.display())))?;
connection
.execute_batch(
"PRAGMA journal_mode=WAL; PRAGMA busy_timeout=15000; PRAGMA foreign_keys=ON;",
)
.map_err(Error::storage)?;
connection
.execute_batch(IDENTITY_MIGRATION)
.map_err(Error::storage)?;
let directory = Self {
database: Mutex::new(connection),
};
directory.seed_bootstrap_user(bootstrap_handle)?;
Ok(directory)
}
fn seed_bootstrap_user(&self, handle: &str) -> Result<()> {
let handle = normalize_username(handle);
validate_authorizable_handle(&handle, "Telegram bootstrap handle")?;
let database = self.lock()?;
let now = Utc::now().to_rfc3339();
database
.execute(
"INSERT INTO whitelist_entries(handle,can_add_users,whitelisted_at,updated_at)
VALUES(?1,1,?2,?2)
ON CONFLICT(handle) DO UPDATE SET can_add_users=1,updated_at=excluded.updated_at",
params![handle, now],
)
.map_err(Error::storage)?;
Ok(())
}
fn lock(&self) -> Result<std::sync::MutexGuard<'_, Connection>> {
self.database
.lock()
.map_err(|_| Error::storage("locking Telegram identity directory"))
}
pub fn provisioning_users(&self) -> Result<Vec<User>> {
let database = self.lock()?;
let mut statement = database
.prepare(
"SELECT handle,telegram_user_id,current_username,display_name,root_node_id,root_ready,can_add_users
FROM whitelist_entries WHERE root_ready=0 ORDER BY whitelisted_at,handle",
)
.map_err(Error::storage)?;
statement
.query_map([], row_user)
.map_err(Error::storage)?
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::storage)
}
pub fn provisioning_groups(&self) -> Result<Vec<Group>> {
let database = self.lock()?;
let mut statement = database
.prepare(
"SELECT group_id,root_node_id,root_ready FROM telegram_group_roots
WHERE root_ready=0 ORDER BY datetime(created_at),group_id",
)
.map_err(Error::storage)?;
statement
.query_map([], row_group)
.map_err(Error::storage)?
.collect::<std::result::Result<Vec<_>, _>>()
.map_err(Error::storage)
}
pub fn user(&self, telegram_user_id: i64) -> Result<User> {
if telegram_user_id == ANONYMOUS_GROUP_USER_ID {
return Err(Error::not_found());
}
let database = self.lock()?;
directory_user_by_id(&database, telegram_user_id)?.ok_or_else(Error::not_found)
}
pub fn group(&self, group_id: &str) -> Result<Group> {
let database = self.lock()?;
directory_group_by_id(&database, group_id)?.ok_or_else(Error::not_found)
}
pub fn group_for_root(&self, root_node_id: NodeId) -> Result<Group> {
let database = self.lock()?;
database
.query_row(
"SELECT group_id,root_node_id,root_ready FROM telegram_group_roots
WHERE root_node_id=?1 AND root_ready=1",
[root_node_id.to_string()],
row_group,
)
.optional()
.map_err(Error::storage)?
.ok_or_else(Error::not_found)
}
pub fn authorize_user_id(&self, handle: &str, telegram_user_id: i64) -> Result<User> {
let handle = normalize_username(handle);
validate_authorizable_handle(&handle, "Telegram handle")?;
validate_authorizable_user_id(telegram_user_id)?;
let mut database = self.lock()?;
let transaction = database
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(Error::storage)?;
let current =
directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
if let Some(current_id) = current.telegram_user_id {
if current_id != telegram_user_id {
return Err(Error::conflict(
"This Telegram handle already has a different numeric identity.",
));
}
transaction.commit().map_err(Error::storage)?;
return Ok(current);
}
if directory_user_by_id(&transaction, telegram_user_id)?.is_some() {
return Err(Error::conflict(
"This numeric Telegram identity is already authorized for another handle.",
));
}
let observed = observed_identity_by_id(&transaction, telegram_user_id)?;
let (current_username, display_name) = observed.unwrap_or((None, None));
let changed = transaction
.execute(
"UPDATE whitelist_entries
SET telegram_user_id=?1,current_username=?2,display_name=?3,
resolved_at=?4,updated_at=?4
WHERE handle=?5 AND telegram_user_id IS NULL",
params![
telegram_user_id,
current_username,
display_name,
Utc::now().to_rfc3339(),
handle
],
)
.map_err(Error::storage)?;
if changed != 1 {
return Err(Error::conflict(
"The Telegram handle changed while its numeric identity was being authorized.",
));
}
let user = directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
transaction.commit().map_err(Error::storage)?;
Ok(user)
}
pub fn complete_handle_root(&self, handle: &str, root_node_id: NodeId) -> Result<User> {
let handle = normalize_username(handle);
let root_node_id = root_node_id.to_string();
let mut database = self.lock()?;
let transaction = database
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(Error::storage)?;
let current =
directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
ensure_user_root_compatible(¤t, &root_node_id, "whitelisted handle")?;
transaction
.execute(
"UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
WHERE handle=?3",
params![root_node_id, Utc::now().to_rfc3339(), handle],
)
.map_err(Error::storage)?;
let user = directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
transaction.commit().map_err(Error::storage)?;
Ok(user)
}
pub fn complete_user_root(&self, telegram_user_id: i64, root_node_id: NodeId) -> Result<User> {
validate_authorizable_user_id(telegram_user_id)?;
let root_node_id = root_node_id.to_string();
let mut database = self.lock()?;
let transaction = database
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(Error::storage)?;
let current =
directory_user_by_id(&transaction, telegram_user_id)?.ok_or_else(Error::not_found)?;
ensure_user_root_compatible(¤t, &root_node_id, "Telegram identity")?;
transaction
.execute(
"UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
WHERE telegram_user_id=?3",
params![root_node_id, Utc::now().to_rfc3339(), telegram_user_id],
)
.map_err(Error::storage)?;
let user =
directory_user_by_id(&transaction, telegram_user_id)?.ok_or_else(Error::not_found)?;
transaction.commit().map_err(Error::storage)?;
Ok(user)
}
pub fn complete_group_root(&self, group_id: &str, root_node_id: NodeId) -> Result<Group> {
let root_node_id = root_node_id.to_string();
let mut database = self.lock()?;
let transaction = database
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(Error::storage)?;
let current =
directory_group_by_id(&transaction, group_id)?.ok_or_else(Error::not_found)?;
if current.root_ready && current.root_node_id.as_deref() != Some(&root_node_id) {
return Err(Error::conflict(
"This Telegram group already has a different root node.",
));
}
transaction
.execute(
"UPDATE telegram_group_roots SET root_node_id=?1,root_ready=1,updated_at=?2
WHERE group_id=?3",
params![root_node_id, Utc::now().to_rfc3339(), group_id],
)
.map_err(Error::storage)?;
let group = directory_group_by_id(&transaction, group_id)?.ok_or_else(Error::not_found)?;
transaction.commit().map_err(Error::storage)?;
Ok(group)
}
}
impl IdentitySink for Directory {
fn observe_identity(&self, observation: &IdentityObservation) -> anyhow::Result<()> {
let mut database = self.lock()?;
observe_identity(&mut database, observation)?;
Ok(())
}
fn whitelist(&self) -> anyhow::Result<WhitelistSnapshot> {
let database = self.lock()?;
let telegram_user_ids = database
.prepare(
"SELECT telegram_user_id FROM whitelist_entries
WHERE telegram_user_id IS NOT NULL AND telegram_user_id != ?1
ORDER BY telegram_user_id",
)?
.query_map([ANONYMOUS_GROUP_USER_ID], |row| row.get::<_, i64>(0))?
.collect::<std::result::Result<HashSet<_>, _>>()?;
Ok(WhitelistSnapshot { telegram_user_ids })
}
fn request_add_user(
&self,
requested_by_telegram_user_id: i64,
handle: &str,
) -> anyhow::Result<AddUserOutcome> {
if requested_by_telegram_user_id == ANONYMOUS_GROUP_USER_ID {
return Ok(AddUserOutcome::Forbidden);
}
let database = self.lock()?;
let can_add = directory_user_by_id(&database, requested_by_telegram_user_id)?
.is_some_and(|user| user.can_add_users);
if !can_add {
return Ok(AddUserOutcome::Forbidden);
}
let user = whitelist_handle(&database, handle, requested_by_telegram_user_id)?;
Ok(AddUserOutcome::Whitelisted {
handle: user.handle,
telegram_user_id: user.telegram_user_id,
})
}
fn observe_group(&self, group_id: &str) -> anyhow::Result<()> {
let database = self.lock()?;
let now = Utc::now().to_rfc3339();
database.execute(
"INSERT INTO telegram_group_roots(group_id,created_at,updated_at)
VALUES(?1,?2,?2) ON CONFLICT(group_id) DO NOTHING",
params![group_id, now],
)?;
Ok(())
}
}
fn normalize_username(value: &str) -> String {
value.trim().trim_start_matches('@').to_ascii_lowercase()
}
fn validate_authorizable_handle(handle: &str, label: &str) -> Result<()> {
if handle.is_empty() {
return Err(Error::invalid(format!("{label} must not be empty")));
}
if handle == ANONYMOUS_GROUP_HANDLE {
return Err(Error::invalid(
"Telegram's anonymous-group pseudo-user cannot be authorized.",
));
}
Ok(())
}
fn validate_authorizable_user_id(telegram_user_id: i64) -> Result<()> {
if telegram_user_id <= 0 {
return Err(Error::invalid(
"Telegram user ID must be a positive numeric identity.",
));
}
if telegram_user_id == ANONYMOUS_GROUP_USER_ID {
return Err(Error::invalid(
"Telegram's anonymous-group pseudo-user cannot be authorized.",
));
}
Ok(())
}
fn is_anonymous_group_observation(
telegram_user_id: i64,
normalized_username: Option<&str>,
) -> bool {
telegram_user_id == ANONYMOUS_GROUP_USER_ID
|| normalized_username == Some(ANONYMOUS_GROUP_HANDLE)
}
fn directory_user_by_clause(
database: &Connection,
clause: &str,
value: &dyn rusqlite::ToSql,
) -> Result<Option<User>> {
database
.query_row(
&format!(
"SELECT handle,telegram_user_id,current_username,display_name,root_node_id,root_ready,can_add_users
FROM whitelist_entries WHERE {clause}"
),
[value],
row_user,
)
.optional()
.map_err(Error::storage)
}
fn directory_user_by_id(database: &Connection, telegram_user_id: i64) -> Result<Option<User>> {
directory_user_by_clause(database, "telegram_user_id=?1", &telegram_user_id)
}
fn directory_user_by_handle(database: &Connection, handle: &str) -> Result<Option<User>> {
directory_user_by_clause(database, "handle=?1", &handle)
}
fn observed_identity_by_id(
database: &Connection,
telegram_user_id: i64,
) -> Result<Option<(Option<String>, Option<String>)>> {
database
.query_row(
"SELECT current_username,display_name FROM observed_identities
WHERE telegram_user_id=?1",
[telegram_user_id],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.optional()
.map_err(Error::storage)
}
fn directory_group_by_id(database: &Connection, group_id: &str) -> Result<Option<Group>> {
database
.query_row(
"SELECT group_id,root_node_id,root_ready FROM telegram_group_roots WHERE group_id=?1",
[group_id],
row_group,
)
.optional()
.map_err(Error::storage)
}
fn row_user(row: &rusqlite::Row<'_>) -> rusqlite::Result<User> {
let root_node_id = canonical_root(row.get(4)?, 4)?;
Ok(User {
handle: row.get(0)?,
telegram_user_id: row.get(1)?,
current_username: row.get(2)?,
display_name: row.get(3)?,
root_node_id,
root_ready: row.get::<_, i64>(5)? != 0,
can_add_users: row.get::<_, i64>(6)? != 0,
})
}
fn row_group(row: &rusqlite::Row<'_>) -> rusqlite::Result<Group> {
let root_node_id = canonical_root(row.get(1)?, 1)?;
Ok(Group {
group_id: row.get(0)?,
root_node_id,
root_ready: row.get::<_, i64>(2)? != 0,
})
}
fn canonical_root(value: Option<String>, column: usize) -> rusqlite::Result<Option<String>> {
value
.map(|value| {
value
.parse::<NodeId>()
.map(|id| id.to_string())
.map_err(|error| {
rusqlite::Error::FromSqlConversionFailure(
column,
rusqlite::types::Type::Text,
Box::new(error),
)
})
})
.transpose()
}
fn observe_identity(database: &mut Connection, observation: &IdentityObservation) -> Result<()> {
let normalized = observation
.username
.as_deref()
.map(normalize_username)
.filter(|value| !value.is_empty());
if is_anonymous_group_observation(observation.telegram_user_id, normalized.as_deref()) {
return Ok(());
}
let transaction = database
.transaction_with_behavior(TransactionBehavior::Immediate)
.map_err(Error::storage)?;
let now = Utc::now().to_rfc3339();
transaction
.execute(
"INSERT INTO observed_identities(telegram_user_id,current_username,display_name,first_seen_at,last_seen_at)
VALUES(?1,?2,?3,?4,?4)
ON CONFLICT(telegram_user_id) DO UPDATE SET
current_username=excluded.current_username,
display_name=excluded.display_name,last_seen_at=excluded.last_seen_at",
params![
observation.telegram_user_id,
normalized,
observation.display_name,
now
],
)
.map_err(Error::storage)?;
if directory_user_by_id(&transaction, observation.telegram_user_id)?.is_some() {
transaction
.execute(
"UPDATE whitelist_entries SET current_username=?1,display_name=?2,updated_at=?3
WHERE telegram_user_id=?4",
params![
normalized,
observation.display_name,
now,
observation.telegram_user_id
],
)
.map_err(Error::storage)?;
} else if observation.telegram_user_id > 0
&& let Some(handle) = normalized.as_deref()
&& directory_user_by_handle(&transaction, handle)?
.is_some_and(|user| user.telegram_user_id.is_none())
{
let changed = transaction
.execute(
"UPDATE whitelist_entries
SET telegram_user_id=?1,current_username=?2,display_name=?3,
resolved_at=?4,updated_at=?4
WHERE handle=?5 AND telegram_user_id IS NULL",
params![
observation.telegram_user_id,
normalized,
observation.display_name,
now,
handle
],
)
.map_err(Error::storage)?;
if changed != 1 {
return Err(Error::conflict(
"The Telegram handle changed while its first identity observation was binding.",
));
}
}
transaction.commit().map_err(Error::storage)
}
fn whitelist_handle(database: &Connection, handle: &str, added_by: i64) -> Result<User> {
let handle = normalize_username(handle.trim_matches(['\'', '"']));
validate_authorizable_handle(&handle, "Telegram handle")?;
let now = Utc::now().to_rfc3339();
database
.execute(
"INSERT INTO whitelist_entries(handle,added_by_telegram_user_id,whitelisted_at,updated_at)
VALUES(?1,?2,?3,?3)
ON CONFLICT(handle) DO UPDATE SET updated_at=excluded.updated_at",
params![handle, added_by, now],
)
.map_err(Error::storage)?;
directory_user_by_handle(database, &handle)?.ok_or_else(Error::not_found)
}
fn ensure_user_root_compatible(user: &User, root: &str, label: &str) -> Result<()> {
if user.root_ready && user.root_node_id.as_deref() != Some(root) {
return Err(Error::conflict(format!(
"This {label} already has a different root node."
)));
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use kcode_tg_kennedy_bot::IdentitySink;
use std::sync::{Arc, Barrier};
use std::thread;
fn directory() -> Directory {
let database = Connection::open_in_memory().unwrap();
database
.execute_batch(
"CREATE TABLE kmap_system_roots(
role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
created_at TEXT NOT NULL
);
INSERT INTO kmap_system_roots VALUES(
'user','AAAAAAAB','2026-01-01T00:00:00Z'
);",
)
.unwrap();
database.execute_batch(IDENTITY_MIGRATION).unwrap();
let directory = Directory {
database: Mutex::new(database),
};
directory.seed_bootstrap_user("@taek42").unwrap();
directory
}
fn preauthorize(directory: &Directory, handle: &str) {
directory.authorize_user_id("taek42", 42).unwrap();
assert!(matches!(
directory.request_add_user(42, handle).unwrap(),
AddUserOutcome::Whitelisted {
telegram_user_id: None,
..
}
));
}
#[test]
fn opens_against_the_identity_schema_created_by_kmap_startup() {
let directory = std::env::temp_dir().join(format!(
"kennedy-telegram-identity-startup-test-{}",
uuid::Uuid::new_v4()
));
std::fs::create_dir_all(&directory).unwrap();
let user_database = directory.join("users.sqlite3");
let kmap_schema = Connection::open(&user_database).unwrap();
kmap_schema
.execute_batch(
"CREATE TABLE kmap_system_roots(
role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
created_at TEXT NOT NULL
);
INSERT INTO kmap_system_roots VALUES(
'user','AAAAAAAB','2026-01-01T00:00:00Z'
);",
)
.unwrap();
drop(kmap_schema);
let identity = Directory::open(&user_database, "@taek42").unwrap();
assert!(
identity
.lock()
.unwrap()
.query_row(
"SELECT EXISTS(SELECT 1 FROM whitelist_entries WHERE handle='taek42')",
[],
|row| row.get::<_, i64>(0),
)
.unwrap()
!= 0
);
drop(identity);
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn matching_first_observation_binds_preauthorized_handle() {
let directory = directory();
preauthorize(&directory, "@friend");
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("@FrIeNd".into()),
display_name: "Friend".into(),
})
.unwrap();
let friend = directory.user(77).unwrap();
assert_eq!(friend.handle, "friend");
assert_eq!(friend.current_username.as_deref(), Some("friend"));
assert!(
directory
.whitelist()
.unwrap()
.telegram_user_ids
.contains(&77)
);
}
#[test]
fn nonmatching_observation_remains_metadata_only() {
let directory = directory();
preauthorize(&directory, "friend");
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("stranger".into()),
display_name: "Stranger".into(),
})
.unwrap();
assert_eq!(directory.user(77).unwrap_err().kind(), ErrorKind::NotFound);
assert!(
!directory
.whitelist()
.unwrap()
.telegram_user_ids
.contains(&77)
);
assert_eq!(
observed_identity_by_id(&directory.lock().unwrap(), 77).unwrap(),
Some((Some("stranger".into()), Some("Stranger".into())))
);
}
#[test]
fn existing_numeric_binding_survives_username_change() {
let directory = directory();
preauthorize(&directory, "friend");
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("friend".into()),
display_name: "Friend".into(),
})
.unwrap();
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("renamed".into()),
display_name: "Renamed Friend".into(),
})
.unwrap();
let friend = directory.user(77).unwrap();
assert_eq!(friend.handle, "friend");
assert_eq!(friend.current_username.as_deref(), Some("renamed"));
assert!(
directory
.whitelist()
.unwrap()
.telegram_user_ids
.contains(&77)
);
}
#[test]
fn conflicting_observations_do_not_reassign_id_or_handle() {
let directory = directory();
preauthorize(&directory, "friend");
assert!(matches!(
directory.request_add_user(42, "colleague").unwrap(),
AddUserOutcome::Whitelisted {
telegram_user_id: None,
..
}
));
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("friend".into()),
display_name: "Friend".into(),
})
.unwrap();
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 88,
username: Some("friend".into()),
display_name: "Other".into(),
})
.unwrap();
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("colleague".into()),
display_name: "Renamed Friend".into(),
})
.unwrap();
assert_eq!(directory.user(88).unwrap_err().kind(), ErrorKind::NotFound);
let database = directory.lock().unwrap();
let friend = directory_user_by_handle(&database, "friend")
.unwrap()
.unwrap();
let colleague = directory_user_by_handle(&database, "colleague")
.unwrap()
.unwrap();
assert_eq!(friend.telegram_user_id, Some(77));
assert_eq!(friend.current_username.as_deref(), Some("colleague"));
assert_eq!(colleague.telegram_user_id, None);
}
#[test]
fn explicit_numeric_authorization_replays_and_conflicts() {
let directory = directory();
let first = directory.authorize_user_id("taek42", 42).unwrap();
let replay = directory.authorize_user_id("@TaEk42", 42).unwrap();
assert_eq!(replay, first);
let mismatch = directory.authorize_user_id("taek42", 43).unwrap_err();
assert_eq!(mismatch.kind(), ErrorKind::Conflict);
{
let database = directory.lock().unwrap();
whitelist_handle(&database, "friend", 42).unwrap();
}
let duplicate = directory.authorize_user_id("friend", 42).unwrap_err();
assert_eq!(duplicate.kind(), ErrorKind::Conflict);
}
#[test]
fn anonymous_group_pseudo_user_is_rejected_at_runtime() {
let directory = directory();
directory
.observe_identity(&IdentityObservation {
telegram_user_id: ANONYMOUS_GROUP_USER_ID,
username: Some("GroupAnonymousBot".into()),
display_name: "Anonymous".into(),
})
.unwrap();
directory
.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("@GroupAnonymousBot".into()),
display_name: "Spoof".into(),
})
.unwrap();
let database = directory.lock().unwrap();
assert_eq!(
database
.query_row("SELECT COUNT(*) FROM observed_identities", [], |row| {
row.get::<_, i64>(0)
})
.unwrap(),
0
);
drop(database);
assert_eq!(
directory
.authorize_user_id("taek42", ANONYMOUS_GROUP_USER_ID)
.unwrap_err()
.kind(),
ErrorKind::InvalidInput
);
assert_eq!(
directory
.authorize_user_id(ANONYMOUS_GROUP_HANDLE, 77)
.unwrap_err()
.kind(),
ErrorKind::InvalidInput
);
directory.authorize_user_id("taek42", 42).unwrap();
assert!(
directory
.request_add_user(42, ANONYMOUS_GROUP_HANDLE)
.is_err()
);
assert!(
!directory
.whitelist()
.unwrap()
.telegram_user_ids
.contains(&ANONYMOUS_GROUP_USER_ID)
);
}
#[test]
fn identity_migration_removes_legacy_anonymous_group_pseudo_user() {
let directory = directory();
let database = directory.lock().unwrap();
database
.execute(
"INSERT INTO observed_identities(
telegram_user_id,current_username,display_name,first_seen_at,last_seen_at
) VALUES(1087968824,'GroupAnonymousBot','Group',?1,?1)",
[Utc::now().to_rfc3339()],
)
.unwrap();
database.execute_batch(IDENTITY_MIGRATION).unwrap();
assert_eq!(
database
.query_row(
"SELECT COUNT(*) FROM observed_identities WHERE telegram_user_id=1087968824",
[],
|row| row.get::<_, i64>(0),
)
.unwrap(),
0
);
}
#[test]
fn add_user_capability_and_group_roots_stay_in_kennedy() {
let directory = directory();
directory.authorize_user_id("taek42", 42).unwrap();
assert!(matches!(
directory.request_add_user(77, "@friend").unwrap(),
AddUserOutcome::Forbidden
));
assert!(matches!(
directory.request_add_user(42, "@friend").unwrap(),
AddUserOutcome::Whitelisted {
telegram_user_id: None,
..
}
));
directory.observe_group("opaque-group").unwrap();
let database = directory.lock().unwrap();
let group = directory_group_by_id(&database, "opaque-group")
.unwrap()
.unwrap();
assert_eq!(group.root_node_id, None);
assert!(!group.root_ready);
}
#[test]
fn root_completion_replays_and_conflicts() {
let directory = directory();
directory.authorize_user_id("taek42", 42).unwrap();
directory.observe_group("opaque-group").unwrap();
let user_root = "AAAAAAAC".parse().unwrap();
let user = directory.complete_user_root(42, user_root).unwrap();
assert!(user.root_ready);
assert_eq!(user.root_node_id.as_deref(), Some("AAAAAAAC"));
assert_eq!(
directory
.complete_user_root(42, user_root)
.unwrap()
.root_node_id
.as_deref(),
Some("AAAAAAAC")
);
let group_root = "AAAAAAAD".parse().unwrap();
let group = directory
.complete_group_root("opaque-group", group_root)
.unwrap();
assert!(group.root_ready);
assert_eq!(group.root_node_id.as_deref(), Some("AAAAAAAD"));
assert_eq!(
directory
.complete_group_root("opaque-group", group_root)
.unwrap()
.root_node_id
.as_deref(),
Some("AAAAAAAD")
);
assert_eq!(
directory.group_for_root(group_root).unwrap().group_id,
"opaque-group"
);
assert_eq!(
directory
.group_for_root("AAAAAAAE".parse().unwrap())
.unwrap_err()
.kind(),
ErrorKind::NotFound
);
let mismatch = directory
.complete_group_root("opaque-group", "AAAAAAAE".parse().unwrap())
.unwrap_err();
assert_eq!(mismatch.kind(), ErrorKind::Conflict);
}
#[test]
fn independent_directories_bind_handle_atomically() {
let directory = std::env::temp_dir().join(format!(
"kennedy-telegram-identity-binding-concurrency-test-{}",
uuid::Uuid::new_v4()
));
std::fs::create_dir_all(&directory).unwrap();
let database_path = directory.join("users.sqlite3");
let first = Directory::open(&database_path, "taek42").unwrap();
preauthorize(&first, "friend");
let first = Arc::new(first);
let second = Arc::new(Directory::open(&database_path, "taek42").unwrap());
let barrier = Arc::new(Barrier::new(2));
let first_worker = {
let directory = Arc::clone(&first);
let barrier = Arc::clone(&barrier);
thread::spawn(move || {
barrier.wait();
directory.observe_identity(&IdentityObservation {
telegram_user_id: 77,
username: Some("friend".into()),
display_name: "First".into(),
})
})
};
let second_worker = {
let directory = Arc::clone(&second);
let barrier = Arc::clone(&barrier);
thread::spawn(move || {
barrier.wait();
directory.observe_identity(&IdentityObservation {
telegram_user_id: 88,
username: Some("friend".into()),
display_name: "Second".into(),
})
})
};
first_worker.join().unwrap().unwrap();
second_worker.join().unwrap().unwrap();
let snapshot = first.whitelist().unwrap();
assert_eq!(
[77, 88]
.into_iter()
.filter(|id| snapshot.telegram_user_ids.contains(id))
.count(),
1
);
let database = first.lock().unwrap();
let friend = directory_user_by_handle(&database, "friend")
.unwrap()
.unwrap();
assert!(matches!(friend.telegram_user_id, Some(77 | 88)));
assert_eq!(
database
.query_row("SELECT COUNT(*) FROM observed_identities", [], |row| {
row.get::<_, i64>(0)
})
.unwrap(),
2
);
drop(database);
drop(first);
drop(second);
std::fs::remove_dir_all(directory).unwrap();
}
#[test]
fn independent_directories_complete_roots_atomically() {
let directory = std::env::temp_dir().join(format!(
"kennedy-telegram-identity-concurrency-test-{}",
uuid::Uuid::new_v4()
));
std::fs::create_dir_all(&directory).unwrap();
let database_path = directory.join("users.sqlite3");
let first = Directory::open(&database_path, "taek42").unwrap();
first.authorize_user_id("taek42", 42).unwrap();
let first = Arc::new(first);
let second = Arc::new(Directory::open(&database_path, "taek42").unwrap());
let barrier = Arc::new(Barrier::new(2));
let first_worker = {
let directory = Arc::clone(&first);
let barrier = Arc::clone(&barrier);
thread::spawn(move || {
barrier.wait();
directory
.complete_user_root(42, "AAAAAAAC".parse().unwrap())
.map(|user| user.root_node_id.unwrap())
.map_err(|error| error.kind())
})
};
let second_worker = {
let directory = Arc::clone(&second);
let barrier = Arc::clone(&barrier);
thread::spawn(move || {
barrier.wait();
directory
.complete_user_root(42, "AAAAAAAD".parse().unwrap())
.map(|user| user.root_node_id.unwrap())
.map_err(|error| error.kind())
})
};
let outcomes = [first_worker.join().unwrap(), second_worker.join().unwrap()];
assert_eq!(outcomes.iter().filter(|outcome| outcome.is_ok()).count(), 1);
assert_eq!(
outcomes
.iter()
.filter(|outcome| matches!(outcome, Err(ErrorKind::Conflict)))
.count(),
1
);
drop(first);
drop(second);
std::fs::remove_dir_all(directory).unwrap();
}
}