Skip to main content

kcode_telegram_identity/
lib.rs

1//! Persistent Telegram identity authorization and Kennedy Kmap-root directory.
2
3use std::{collections::HashSet, path::Path, sync::Mutex};
4
5use chrono::Utc;
6use kcode_kweb_db::NodeId;
7use kcode_tg_kennedy_bot::{AddUserOutcome, IdentityObservation, IdentitySink, WhitelistSnapshot};
8use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
9use serde::{Deserialize, Serialize};
10
11const IDENTITY_MIGRATION: &str = include_str!("../migrations/001_initial.sql");
12const ANONYMOUS_GROUP_USER_ID: i64 = 1_087_968_824;
13const ANONYMOUS_GROUP_HANDLE: &str = "groupanonymousbot";
14
15/// Durable identity authorization and user/group root mappings.
16pub struct Directory {
17    database: Mutex<Connection>,
18}
19
20/// One whitelisted Telegram identity and its Kennedy root assignment.
21#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
22#[serde(rename_all = "camelCase")]
23pub struct User {
24    pub handle: String,
25    pub telegram_user_id: Option<i64>,
26    pub current_username: Option<String>,
27    pub display_name: Option<String>,
28    pub root_node_id: Option<String>,
29    pub root_ready: bool,
30    pub can_add_users: bool,
31}
32
33/// One opaque Telegram group and its Kennedy root assignment.
34#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
35#[serde(rename_all = "camelCase")]
36pub struct Group {
37    pub group_id: String,
38    pub root_node_id: Option<String>,
39    pub root_ready: bool,
40}
41
42/// Stable categories callers can map to their own transport errors.
43#[derive(Clone, Copy, Debug, Eq, PartialEq)]
44pub enum ErrorKind {
45    InvalidInput,
46    NotFound,
47    Conflict,
48    Storage,
49}
50
51/// A typed identity-directory failure.
52#[derive(Debug, thiserror::Error)]
53#[error("{message}")]
54pub struct Error {
55    kind: ErrorKind,
56    message: String,
57}
58
59impl Error {
60    pub fn kind(&self) -> ErrorKind {
61        self.kind
62    }
63
64    pub fn message(&self) -> &str {
65        &self.message
66    }
67
68    fn invalid(message: impl Into<String>) -> Self {
69        Self {
70            kind: ErrorKind::InvalidInput,
71            message: message.into(),
72        }
73    }
74
75    fn not_found() -> Self {
76        Self {
77            kind: ErrorKind::NotFound,
78            message: "Telegram directory entry not found.".into(),
79        }
80    }
81
82    fn conflict(message: impl Into<String>) -> Self {
83        Self {
84            kind: ErrorKind::Conflict,
85            message: message.into(),
86        }
87    }
88
89    fn storage(error: impl std::fmt::Display) -> Self {
90        Self {
91            kind: ErrorKind::Storage,
92            message: error.to_string(),
93        }
94    }
95}
96
97pub type Result<T> = std::result::Result<T, Error>;
98
99impl Directory {
100    /// Open the directory and ensure the normalized bootstrap handle exists.
101    pub fn open(path: &Path, bootstrap_handle: &str) -> Result<Self> {
102        let connection = Connection::open(path)
103            .map_err(|error| Error::storage(format!("opening {}: {error}", path.display())))?;
104        connection
105            .execute_batch(
106                "PRAGMA journal_mode=WAL; PRAGMA busy_timeout=15000; PRAGMA foreign_keys=ON;",
107            )
108            .map_err(Error::storage)?;
109        connection
110            .execute_batch(IDENTITY_MIGRATION)
111            .map_err(Error::storage)?;
112        let directory = Self {
113            database: Mutex::new(connection),
114        };
115        directory.seed_bootstrap_user(bootstrap_handle)?;
116        Ok(directory)
117    }
118
119    fn seed_bootstrap_user(&self, handle: &str) -> Result<()> {
120        let handle = normalize_username(handle);
121        validate_authorizable_handle(&handle, "Telegram bootstrap handle")?;
122        let database = self.lock()?;
123        let now = Utc::now().to_rfc3339();
124        database
125            .execute(
126                "INSERT INTO whitelist_entries(handle,can_add_users,whitelisted_at,updated_at)
127                 VALUES(?1,1,?2,?2)
128                 ON CONFLICT(handle) DO UPDATE SET can_add_users=1,updated_at=excluded.updated_at",
129                params![handle, now],
130            )
131            .map_err(Error::storage)?;
132        Ok(())
133    }
134
135    fn lock(&self) -> Result<std::sync::MutexGuard<'_, Connection>> {
136        self.database
137            .lock()
138            .map_err(|_| Error::storage("locking Telegram identity directory"))
139    }
140
141    /// List preauthorized handles whose Kweb root assignment is incomplete.
142    pub fn provisioning_users(&self) -> Result<Vec<User>> {
143        let database = self.lock()?;
144        let mut statement = database
145            .prepare(
146                "SELECT handle,telegram_user_id,current_username,display_name,root_node_id,root_ready,can_add_users
147                 FROM whitelist_entries WHERE root_ready=0 ORDER BY whitelisted_at,handle",
148            )
149            .map_err(Error::storage)?;
150        statement
151            .query_map([], row_user)
152            .map_err(Error::storage)?
153            .collect::<std::result::Result<Vec<_>, _>>()
154            .map_err(Error::storage)
155    }
156
157    /// List observed groups whose Kweb root assignment is incomplete.
158    pub fn provisioning_groups(&self) -> Result<Vec<Group>> {
159        let database = self.lock()?;
160        let mut statement = database
161            .prepare(
162                "SELECT group_id,root_node_id,root_ready FROM telegram_group_roots
163                 WHERE root_ready=0 ORDER BY datetime(created_at),group_id",
164            )
165            .map_err(Error::storage)?;
166        statement
167            .query_map([], row_group)
168            .map_err(Error::storage)?
169            .collect::<std::result::Result<Vec<_>, _>>()
170            .map_err(Error::storage)
171    }
172
173    /// Look up one explicitly authorized identity by stable numeric Telegram ID.
174    pub fn user(&self, telegram_user_id: i64) -> Result<User> {
175        if telegram_user_id == ANONYMOUS_GROUP_USER_ID {
176            return Err(Error::not_found());
177        }
178        let database = self.lock()?;
179        directory_user_by_id(&database, telegram_user_id)?.ok_or_else(Error::not_found)
180    }
181
182    /// Look up one previously observed opaque Telegram group.
183    pub fn group(&self, group_id: &str) -> Result<Group> {
184        let database = self.lock()?;
185        directory_group_by_id(&database, group_id)?.ok_or_else(Error::not_found)
186    }
187
188    /// Look up one ready Telegram group by its canonical Kweb root.
189    pub fn group_for_root(&self, root_node_id: NodeId) -> Result<Group> {
190        let database = self.lock()?;
191        database
192            .query_row(
193                "SELECT group_id,root_node_id,root_ready FROM telegram_group_roots
194                 WHERE root_node_id=?1 AND root_ready=1",
195                [root_node_id.to_string()],
196                row_group,
197            )
198            .optional()
199            .map_err(Error::storage)?
200            .ok_or_else(Error::not_found)
201    }
202
203    /// Explicitly authorize one stable numeric Telegram ID for a preauthorized handle.
204    pub fn authorize_user_id(&self, handle: &str, telegram_user_id: i64) -> Result<User> {
205        let handle = normalize_username(handle);
206        validate_authorizable_handle(&handle, "Telegram handle")?;
207        validate_authorizable_user_id(telegram_user_id)?;
208
209        let mut database = self.lock()?;
210        let transaction = database
211            .transaction_with_behavior(TransactionBehavior::Immediate)
212            .map_err(Error::storage)?;
213        let current =
214            directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
215
216        if let Some(current_id) = current.telegram_user_id {
217            if current_id != telegram_user_id {
218                return Err(Error::conflict(
219                    "This Telegram handle already has a different numeric identity.",
220                ));
221            }
222            transaction.commit().map_err(Error::storage)?;
223            return Ok(current);
224        }
225
226        if directory_user_by_id(&transaction, telegram_user_id)?.is_some() {
227            return Err(Error::conflict(
228                "This numeric Telegram identity is already authorized for another handle.",
229            ));
230        }
231
232        let observed = observed_identity_by_id(&transaction, telegram_user_id)?;
233        let (current_username, display_name) = observed.unwrap_or((None, None));
234        let changed = transaction
235            .execute(
236                "UPDATE whitelist_entries
237                 SET telegram_user_id=?1,current_username=?2,display_name=?3,
238                     resolved_at=?4,updated_at=?4
239                 WHERE handle=?5 AND telegram_user_id IS NULL",
240                params![
241                    telegram_user_id,
242                    current_username,
243                    display_name,
244                    Utc::now().to_rfc3339(),
245                    handle
246                ],
247            )
248            .map_err(Error::storage)?;
249        if changed != 1 {
250            return Err(Error::conflict(
251                "The Telegram handle changed while its numeric identity was being authorized.",
252            ));
253        }
254
255        let user = directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
256        transaction.commit().map_err(Error::storage)?;
257        Ok(user)
258    }
259
260    /// Complete the root assignment for a preauthorized handle.
261    pub fn complete_handle_root(&self, handle: &str, root_node_id: NodeId) -> Result<User> {
262        let handle = normalize_username(handle);
263        let root_node_id = root_node_id.to_string();
264        let mut database = self.lock()?;
265        let transaction = database
266            .transaction_with_behavior(TransactionBehavior::Immediate)
267            .map_err(Error::storage)?;
268        let current =
269            directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
270        ensure_user_root_compatible(&current, &root_node_id, "whitelisted handle")?;
271        transaction
272            .execute(
273                "UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
274                 WHERE handle=?3",
275                params![root_node_id, Utc::now().to_rfc3339(), handle],
276            )
277            .map_err(Error::storage)?;
278        let user = directory_user_by_handle(&transaction, &handle)?.ok_or_else(Error::not_found)?;
279        transaction.commit().map_err(Error::storage)?;
280        Ok(user)
281    }
282
283    /// Complete the root assignment for an explicitly authorized Telegram identity.
284    pub fn complete_user_root(&self, telegram_user_id: i64, root_node_id: NodeId) -> Result<User> {
285        validate_authorizable_user_id(telegram_user_id)?;
286        let root_node_id = root_node_id.to_string();
287        let mut database = self.lock()?;
288        let transaction = database
289            .transaction_with_behavior(TransactionBehavior::Immediate)
290            .map_err(Error::storage)?;
291        let current =
292            directory_user_by_id(&transaction, telegram_user_id)?.ok_or_else(Error::not_found)?;
293        ensure_user_root_compatible(&current, &root_node_id, "Telegram identity")?;
294        transaction
295            .execute(
296                "UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
297                 WHERE telegram_user_id=?3",
298                params![root_node_id, Utc::now().to_rfc3339(), telegram_user_id],
299            )
300            .map_err(Error::storage)?;
301        let user =
302            directory_user_by_id(&transaction, telegram_user_id)?.ok_or_else(Error::not_found)?;
303        transaction.commit().map_err(Error::storage)?;
304        Ok(user)
305    }
306
307    /// Complete the root assignment for an observed opaque Telegram group.
308    pub fn complete_group_root(&self, group_id: &str, root_node_id: NodeId) -> Result<Group> {
309        let root_node_id = root_node_id.to_string();
310        let mut database = self.lock()?;
311        let transaction = database
312            .transaction_with_behavior(TransactionBehavior::Immediate)
313            .map_err(Error::storage)?;
314        let current =
315            directory_group_by_id(&transaction, group_id)?.ok_or_else(Error::not_found)?;
316        if current.root_ready && current.root_node_id.as_deref() != Some(&root_node_id) {
317            return Err(Error::conflict(
318                "This Telegram group already has a different root node.",
319            ));
320        }
321        transaction
322            .execute(
323                "UPDATE telegram_group_roots SET root_node_id=?1,root_ready=1,updated_at=?2
324                 WHERE group_id=?3",
325                params![root_node_id, Utc::now().to_rfc3339(), group_id],
326            )
327            .map_err(Error::storage)?;
328        let group = directory_group_by_id(&transaction, group_id)?.ok_or_else(Error::not_found)?;
329        transaction.commit().map_err(Error::storage)?;
330        Ok(group)
331    }
332}
333
334impl IdentitySink for Directory {
335    fn observe_identity(&self, observation: &IdentityObservation) -> anyhow::Result<()> {
336        let database = self.lock()?;
337        observe_identity(&database, observation)?;
338        Ok(())
339    }
340
341    fn whitelist(&self) -> anyhow::Result<WhitelistSnapshot> {
342        let database = self.lock()?;
343        let telegram_user_ids = database
344            .prepare(
345                "SELECT telegram_user_id FROM whitelist_entries
346                 WHERE telegram_user_id IS NOT NULL AND telegram_user_id != ?1
347                 ORDER BY telegram_user_id",
348            )?
349            .query_map([ANONYMOUS_GROUP_USER_ID], |row| row.get::<_, i64>(0))?
350            .collect::<std::result::Result<HashSet<_>, _>>()?;
351        Ok(WhitelistSnapshot { telegram_user_ids })
352    }
353
354    fn request_add_user(
355        &self,
356        requested_by_telegram_user_id: i64,
357        handle: &str,
358    ) -> anyhow::Result<AddUserOutcome> {
359        if requested_by_telegram_user_id == ANONYMOUS_GROUP_USER_ID {
360            return Ok(AddUserOutcome::Forbidden);
361        }
362        let database = self.lock()?;
363        let can_add = directory_user_by_id(&database, requested_by_telegram_user_id)?
364            .is_some_and(|user| user.can_add_users);
365        if !can_add {
366            return Ok(AddUserOutcome::Forbidden);
367        }
368        let user = whitelist_handle(&database, handle, requested_by_telegram_user_id)?;
369        Ok(AddUserOutcome::Whitelisted {
370            handle: user.handle,
371            telegram_user_id: user.telegram_user_id,
372        })
373    }
374
375    fn observe_group(&self, group_id: &str) -> anyhow::Result<()> {
376        let database = self.lock()?;
377        let now = Utc::now().to_rfc3339();
378        database.execute(
379            "INSERT INTO telegram_group_roots(group_id,created_at,updated_at)
380             VALUES(?1,?2,?2) ON CONFLICT(group_id) DO NOTHING",
381            params![group_id, now],
382        )?;
383        Ok(())
384    }
385}
386
387fn normalize_username(value: &str) -> String {
388    value.trim().trim_start_matches('@').to_ascii_lowercase()
389}
390
391fn validate_authorizable_handle(handle: &str, label: &str) -> Result<()> {
392    if handle.is_empty() {
393        return Err(Error::invalid(format!("{label} must not be empty")));
394    }
395    if handle == ANONYMOUS_GROUP_HANDLE {
396        return Err(Error::invalid(
397            "Telegram's anonymous-group pseudo-user cannot be authorized.",
398        ));
399    }
400    Ok(())
401}
402
403fn validate_authorizable_user_id(telegram_user_id: i64) -> Result<()> {
404    if telegram_user_id <= 0 {
405        return Err(Error::invalid(
406            "Telegram user ID must be a positive numeric identity.",
407        ));
408    }
409    if telegram_user_id == ANONYMOUS_GROUP_USER_ID {
410        return Err(Error::invalid(
411            "Telegram's anonymous-group pseudo-user cannot be authorized.",
412        ));
413    }
414    Ok(())
415}
416
417fn is_anonymous_group_observation(
418    telegram_user_id: i64,
419    normalized_username: Option<&str>,
420) -> bool {
421    telegram_user_id == ANONYMOUS_GROUP_USER_ID
422        || normalized_username == Some(ANONYMOUS_GROUP_HANDLE)
423}
424
425fn directory_user_by_clause(
426    database: &Connection,
427    clause: &str,
428    value: &dyn rusqlite::ToSql,
429) -> Result<Option<User>> {
430    database
431        .query_row(
432            &format!(
433                "SELECT handle,telegram_user_id,current_username,display_name,root_node_id,root_ready,can_add_users
434                 FROM whitelist_entries WHERE {clause}"
435            ),
436            [value],
437            row_user,
438        )
439        .optional()
440        .map_err(Error::storage)
441}
442
443fn directory_user_by_id(database: &Connection, telegram_user_id: i64) -> Result<Option<User>> {
444    directory_user_by_clause(database, "telegram_user_id=?1", &telegram_user_id)
445}
446
447fn directory_user_by_handle(database: &Connection, handle: &str) -> Result<Option<User>> {
448    directory_user_by_clause(database, "handle=?1", &handle)
449}
450
451fn observed_identity_by_id(
452    database: &Connection,
453    telegram_user_id: i64,
454) -> Result<Option<(Option<String>, Option<String>)>> {
455    database
456        .query_row(
457            "SELECT current_username,display_name FROM observed_identities
458             WHERE telegram_user_id=?1",
459            [telegram_user_id],
460            |row| Ok((row.get(0)?, row.get(1)?)),
461        )
462        .optional()
463        .map_err(Error::storage)
464}
465
466fn directory_group_by_id(database: &Connection, group_id: &str) -> Result<Option<Group>> {
467    database
468        .query_row(
469            "SELECT group_id,root_node_id,root_ready FROM telegram_group_roots WHERE group_id=?1",
470            [group_id],
471            row_group,
472        )
473        .optional()
474        .map_err(Error::storage)
475}
476
477fn row_user(row: &rusqlite::Row<'_>) -> rusqlite::Result<User> {
478    let root_node_id = canonical_root(row.get(4)?, 4)?;
479    Ok(User {
480        handle: row.get(0)?,
481        telegram_user_id: row.get(1)?,
482        current_username: row.get(2)?,
483        display_name: row.get(3)?,
484        root_node_id,
485        root_ready: row.get::<_, i64>(5)? != 0,
486        can_add_users: row.get::<_, i64>(6)? != 0,
487    })
488}
489
490fn row_group(row: &rusqlite::Row<'_>) -> rusqlite::Result<Group> {
491    let root_node_id = canonical_root(row.get(1)?, 1)?;
492    Ok(Group {
493        group_id: row.get(0)?,
494        root_node_id,
495        root_ready: row.get::<_, i64>(2)? != 0,
496    })
497}
498
499fn canonical_root(value: Option<String>, column: usize) -> rusqlite::Result<Option<String>> {
500    value
501        .map(|value| {
502            value
503                .parse::<NodeId>()
504                .map(|id| id.to_string())
505                .map_err(|error| {
506                    rusqlite::Error::FromSqlConversionFailure(
507                        column,
508                        rusqlite::types::Type::Text,
509                        Box::new(error),
510                    )
511                })
512        })
513        .transpose()
514}
515
516fn observe_identity(database: &Connection, observation: &IdentityObservation) -> Result<()> {
517    let normalized = observation
518        .username
519        .as_deref()
520        .map(normalize_username)
521        .filter(|value| !value.is_empty());
522    if is_anonymous_group_observation(observation.telegram_user_id, normalized.as_deref()) {
523        return Ok(());
524    }
525
526    let now = Utc::now().to_rfc3339();
527    database
528        .execute(
529            "INSERT INTO observed_identities(telegram_user_id,current_username,display_name,first_seen_at,last_seen_at)
530             VALUES(?1,?2,?3,?4,?4)
531             ON CONFLICT(telegram_user_id) DO UPDATE SET
532                 current_username=excluded.current_username,
533                 display_name=excluded.display_name,last_seen_at=excluded.last_seen_at",
534            params![
535                observation.telegram_user_id,
536                normalized,
537                observation.display_name,
538                now
539            ],
540        )
541        .map_err(Error::storage)?;
542
543    if directory_user_by_id(database, observation.telegram_user_id)?.is_some() {
544        database
545            .execute(
546                "UPDATE whitelist_entries SET current_username=?1,display_name=?2,updated_at=?3
547                 WHERE telegram_user_id=?4",
548                params![
549                    normalized,
550                    observation.display_name,
551                    now,
552                    observation.telegram_user_id
553                ],
554            )
555            .map_err(Error::storage)?;
556    }
557    Ok(())
558}
559
560fn whitelist_handle(database: &Connection, handle: &str, added_by: i64) -> Result<User> {
561    let handle = normalize_username(handle.trim_matches(['\'', '"']));
562    validate_authorizable_handle(&handle, "Telegram handle")?;
563    let now = Utc::now().to_rfc3339();
564    database
565        .execute(
566            "INSERT INTO whitelist_entries(handle,added_by_telegram_user_id,whitelisted_at,updated_at)
567             VALUES(?1,?2,?3,?3)
568             ON CONFLICT(handle) DO UPDATE SET updated_at=excluded.updated_at",
569            params![handle, added_by, now],
570        )
571        .map_err(Error::storage)?;
572    directory_user_by_handle(database, &handle)?.ok_or_else(Error::not_found)
573}
574
575fn ensure_user_root_compatible(user: &User, root: &str, label: &str) -> Result<()> {
576    if user.root_ready && user.root_node_id.as_deref() != Some(root) {
577        return Err(Error::conflict(format!(
578            "This {label} already has a different root node."
579        )));
580    }
581    Ok(())
582}
583
584#[cfg(test)]
585mod tests {
586    use super::*;
587    use kcode_tg_kennedy_bot::IdentitySink;
588    use std::sync::{Arc, Barrier};
589    use std::thread;
590
591    fn directory() -> Directory {
592        let database = Connection::open_in_memory().unwrap();
593        database
594            .execute_batch(
595                "CREATE TABLE kmap_system_roots(
596                     role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
597                     root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
598                     created_at TEXT NOT NULL
599                 );
600                 INSERT INTO kmap_system_roots VALUES(
601                     'user','AAAAAAAB','2026-01-01T00:00:00Z'
602                 );",
603            )
604            .unwrap();
605        database.execute_batch(IDENTITY_MIGRATION).unwrap();
606        let directory = Directory {
607            database: Mutex::new(database),
608        };
609        directory.seed_bootstrap_user("@taek42").unwrap();
610        directory
611    }
612
613    #[test]
614    fn opens_against_the_identity_schema_created_by_kmap_startup() {
615        let directory = std::env::temp_dir().join(format!(
616            "kennedy-telegram-identity-startup-test-{}",
617            uuid::Uuid::new_v4()
618        ));
619        std::fs::create_dir_all(&directory).unwrap();
620        let user_database = directory.join("users.sqlite3");
621        let kmap_schema = Connection::open(&user_database).unwrap();
622        kmap_schema
623            .execute_batch(
624                "CREATE TABLE kmap_system_roots(
625                     role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
626                     root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
627                     created_at TEXT NOT NULL
628                 );
629                 INSERT INTO kmap_system_roots VALUES(
630                     'user','AAAAAAAB','2026-01-01T00:00:00Z'
631                 );",
632            )
633            .unwrap();
634        drop(kmap_schema);
635
636        let identity = Directory::open(&user_database, "@taek42").unwrap();
637        assert!(
638            identity
639                .lock()
640                .unwrap()
641                .query_row(
642                    "SELECT EXISTS(SELECT 1 FROM whitelist_entries WHERE handle='taek42')",
643                    [],
644                    |row| row.get::<_, i64>(0),
645                )
646                .unwrap()
647                != 0
648        );
649        drop(identity);
650        std::fs::remove_dir_all(directory).unwrap();
651    }
652
653    #[test]
654    fn usernames_are_observations_and_numeric_authorization_is_explicit() {
655        let directory = directory();
656        directory
657            .observe_identity(&IdentityObservation {
658                telegram_user_id: 42,
659                username: Some("TaEk42".into()),
660                display_name: "David".into(),
661            })
662            .unwrap();
663        assert!(
664            !directory
665                .whitelist()
666                .unwrap()
667                .telegram_user_ids
668                .contains(&42)
669        );
670        assert_eq!(directory.user(42).unwrap_err().kind(), ErrorKind::NotFound);
671
672        let authorized = directory.authorize_user_id("taek42", 42).unwrap();
673        assert_eq!(authorized.telegram_user_id, Some(42));
674        assert_eq!(authorized.current_username.as_deref(), Some("taek42"));
675        assert!(
676            directory
677                .whitelist()
678                .unwrap()
679                .telegram_user_ids
680                .contains(&42)
681        );
682
683        directory
684            .observe_identity(&IdentityObservation {
685                telegram_user_id: 43,
686                username: Some("taek42".into()),
687                display_name: "Spoof".into(),
688            })
689            .unwrap();
690        assert!(
691            !directory
692                .whitelist()
693                .unwrap()
694                .telegram_user_ids
695                .contains(&43)
696        );
697
698        directory
699            .observe_identity(&IdentityObservation {
700                telegram_user_id: 42,
701                username: Some("renamed".into()),
702                display_name: "David".into(),
703            })
704            .unwrap();
705        assert_eq!(
706            directory.user(42).unwrap().current_username.as_deref(),
707            Some("renamed")
708        );
709    }
710
711    #[test]
712    fn explicit_numeric_authorization_replays_and_conflicts() {
713        let directory = directory();
714        let first = directory.authorize_user_id("taek42", 42).unwrap();
715        let replay = directory.authorize_user_id("@TaEk42", 42).unwrap();
716        assert_eq!(replay, first);
717
718        let mismatch = directory.authorize_user_id("taek42", 43).unwrap_err();
719        assert_eq!(mismatch.kind(), ErrorKind::Conflict);
720
721        {
722            let database = directory.lock().unwrap();
723            whitelist_handle(&database, "friend", 42).unwrap();
724        }
725        let duplicate = directory.authorize_user_id("friend", 42).unwrap_err();
726        assert_eq!(duplicate.kind(), ErrorKind::Conflict);
727    }
728
729    #[test]
730    fn anonymous_group_pseudo_user_is_rejected_at_runtime() {
731        let directory = directory();
732        directory
733            .observe_identity(&IdentityObservation {
734                telegram_user_id: ANONYMOUS_GROUP_USER_ID,
735                username: Some("GroupAnonymousBot".into()),
736                display_name: "Anonymous".into(),
737            })
738            .unwrap();
739        directory
740            .observe_identity(&IdentityObservation {
741                telegram_user_id: 77,
742                username: Some("@GroupAnonymousBot".into()),
743                display_name: "Spoof".into(),
744            })
745            .unwrap();
746
747        let database = directory.lock().unwrap();
748        assert_eq!(
749            database
750                .query_row("SELECT COUNT(*) FROM observed_identities", [], |row| {
751                    row.get::<_, i64>(0)
752                })
753                .unwrap(),
754            0
755        );
756        drop(database);
757
758        assert_eq!(
759            directory
760                .authorize_user_id("taek42", ANONYMOUS_GROUP_USER_ID)
761                .unwrap_err()
762                .kind(),
763            ErrorKind::InvalidInput
764        );
765        assert_eq!(
766            directory
767                .authorize_user_id(ANONYMOUS_GROUP_HANDLE, 77)
768                .unwrap_err()
769                .kind(),
770            ErrorKind::InvalidInput
771        );
772
773        directory.authorize_user_id("taek42", 42).unwrap();
774        assert!(
775            directory
776                .request_add_user(42, ANONYMOUS_GROUP_HANDLE)
777                .is_err()
778        );
779    }
780
781    #[test]
782    fn identity_migration_removes_legacy_anonymous_group_pseudo_user() {
783        let directory = directory();
784        let database = directory.lock().unwrap();
785        database
786            .execute(
787                "INSERT INTO observed_identities(
788                     telegram_user_id,current_username,display_name,first_seen_at,last_seen_at
789                 ) VALUES(1087968824,'GroupAnonymousBot','Group',?1,?1)",
790                [Utc::now().to_rfc3339()],
791            )
792            .unwrap();
793        database.execute_batch(IDENTITY_MIGRATION).unwrap();
794        assert_eq!(
795            database
796                .query_row(
797                    "SELECT COUNT(*) FROM observed_identities WHERE telegram_user_id=1087968824",
798                    [],
799                    |row| row.get::<_, i64>(0),
800                )
801                .unwrap(),
802            0
803        );
804    }
805
806    #[test]
807    fn add_user_capability_and_group_roots_stay_in_kennedy() {
808        let directory = directory();
809        directory.authorize_user_id("taek42", 42).unwrap();
810        assert!(matches!(
811            directory.request_add_user(77, "@friend").unwrap(),
812            AddUserOutcome::Forbidden
813        ));
814        assert!(matches!(
815            directory.request_add_user(42, "@friend").unwrap(),
816            AddUserOutcome::Whitelisted {
817                telegram_user_id: None,
818                ..
819            }
820        ));
821        directory.observe_group("opaque-group").unwrap();
822        let database = directory.lock().unwrap();
823        let group = directory_group_by_id(&database, "opaque-group")
824            .unwrap()
825            .unwrap();
826        assert_eq!(group.root_node_id, None);
827        assert!(!group.root_ready);
828    }
829
830    #[test]
831    fn root_completion_replays_and_conflicts() {
832        let directory = directory();
833        directory.authorize_user_id("taek42", 42).unwrap();
834        directory.observe_group("opaque-group").unwrap();
835
836        let user_root = "AAAAAAAC".parse().unwrap();
837        let user = directory.complete_user_root(42, user_root).unwrap();
838        assert!(user.root_ready);
839        assert_eq!(user.root_node_id.as_deref(), Some("AAAAAAAC"));
840        assert_eq!(
841            directory
842                .complete_user_root(42, user_root)
843                .unwrap()
844                .root_node_id
845                .as_deref(),
846            Some("AAAAAAAC")
847        );
848
849        let group_root = "AAAAAAAD".parse().unwrap();
850        let group = directory
851            .complete_group_root("opaque-group", group_root)
852            .unwrap();
853        assert!(group.root_ready);
854        assert_eq!(group.root_node_id.as_deref(), Some("AAAAAAAD"));
855        assert_eq!(
856            directory
857                .complete_group_root("opaque-group", group_root)
858                .unwrap()
859                .root_node_id
860                .as_deref(),
861            Some("AAAAAAAD")
862        );
863        assert_eq!(
864            directory.group_for_root(group_root).unwrap().group_id,
865            "opaque-group"
866        );
867        assert_eq!(
868            directory
869                .group_for_root("AAAAAAAE".parse().unwrap())
870                .unwrap_err()
871                .kind(),
872            ErrorKind::NotFound
873        );
874
875        let mismatch = directory
876            .complete_group_root("opaque-group", "AAAAAAAE".parse().unwrap())
877            .unwrap_err();
878        assert_eq!(mismatch.kind(), ErrorKind::Conflict);
879    }
880
881    #[test]
882    fn independent_directories_complete_roots_atomically() {
883        let directory = std::env::temp_dir().join(format!(
884            "kennedy-telegram-identity-concurrency-test-{}",
885            uuid::Uuid::new_v4()
886        ));
887        std::fs::create_dir_all(&directory).unwrap();
888        let database_path = directory.join("users.sqlite3");
889
890        let first = Directory::open(&database_path, "taek42").unwrap();
891        first.authorize_user_id("taek42", 42).unwrap();
892        let first = Arc::new(first);
893        let second = Arc::new(Directory::open(&database_path, "taek42").unwrap());
894        let barrier = Arc::new(Barrier::new(2));
895
896        let first_worker = {
897            let directory = Arc::clone(&first);
898            let barrier = Arc::clone(&barrier);
899            thread::spawn(move || {
900                barrier.wait();
901                directory
902                    .complete_user_root(42, "AAAAAAAC".parse().unwrap())
903                    .map(|user| user.root_node_id.unwrap())
904                    .map_err(|error| error.kind())
905            })
906        };
907        let second_worker = {
908            let directory = Arc::clone(&second);
909            let barrier = Arc::clone(&barrier);
910            thread::spawn(move || {
911                barrier.wait();
912                directory
913                    .complete_user_root(42, "AAAAAAAD".parse().unwrap())
914                    .map(|user| user.root_node_id.unwrap())
915                    .map_err(|error| error.kind())
916            })
917        };
918
919        let outcomes = [first_worker.join().unwrap(), second_worker.join().unwrap()];
920        assert_eq!(outcomes.iter().filter(|outcome| outcome.is_ok()).count(), 1);
921        assert_eq!(
922            outcomes
923                .iter()
924                .filter(|outcome| matches!(outcome, Err(ErrorKind::Conflict)))
925                .count(),
926            1
927        );
928
929        drop(first);
930        drop(second);
931        std::fs::remove_dir_all(directory).unwrap();
932    }
933}