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 bound 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 a bound 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 mut database = self.lock()?;
337        observe_identity(&mut 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: &mut 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 transaction = database
527        .transaction_with_behavior(TransactionBehavior::Immediate)
528        .map_err(Error::storage)?;
529    let now = Utc::now().to_rfc3339();
530    transaction
531        .execute(
532            "INSERT INTO observed_identities(telegram_user_id,current_username,display_name,first_seen_at,last_seen_at)
533             VALUES(?1,?2,?3,?4,?4)
534             ON CONFLICT(telegram_user_id) DO UPDATE SET
535                 current_username=excluded.current_username,
536                 display_name=excluded.display_name,last_seen_at=excluded.last_seen_at",
537            params![
538                observation.telegram_user_id,
539                normalized,
540                observation.display_name,
541                now
542            ],
543        )
544        .map_err(Error::storage)?;
545
546    if directory_user_by_id(&transaction, observation.telegram_user_id)?.is_some() {
547        transaction
548            .execute(
549                "UPDATE whitelist_entries SET current_username=?1,display_name=?2,updated_at=?3
550                 WHERE telegram_user_id=?4",
551                params![
552                    normalized,
553                    observation.display_name,
554                    now,
555                    observation.telegram_user_id
556                ],
557            )
558            .map_err(Error::storage)?;
559    } else if observation.telegram_user_id > 0
560        && let Some(handle) = normalized.as_deref()
561        && directory_user_by_handle(&transaction, handle)?
562            .is_some_and(|user| user.telegram_user_id.is_none())
563    {
564        let changed = transaction
565            .execute(
566                "UPDATE whitelist_entries
567                 SET telegram_user_id=?1,current_username=?2,display_name=?3,
568                     resolved_at=?4,updated_at=?4
569                 WHERE handle=?5 AND telegram_user_id IS NULL",
570                params![
571                    observation.telegram_user_id,
572                    normalized,
573                    observation.display_name,
574                    now,
575                    handle
576                ],
577            )
578            .map_err(Error::storage)?;
579        if changed != 1 {
580            return Err(Error::conflict(
581                "The Telegram handle changed while its first identity observation was binding.",
582            ));
583        }
584    }
585
586    transaction.commit().map_err(Error::storage)
587}
588
589fn whitelist_handle(database: &Connection, handle: &str, added_by: i64) -> Result<User> {
590    let handle = normalize_username(handle.trim_matches(['\'', '"']));
591    validate_authorizable_handle(&handle, "Telegram handle")?;
592    let now = Utc::now().to_rfc3339();
593    database
594        .execute(
595            "INSERT INTO whitelist_entries(handle,added_by_telegram_user_id,whitelisted_at,updated_at)
596             VALUES(?1,?2,?3,?3)
597             ON CONFLICT(handle) DO UPDATE SET updated_at=excluded.updated_at",
598            params![handle, added_by, now],
599        )
600        .map_err(Error::storage)?;
601    directory_user_by_handle(database, &handle)?.ok_or_else(Error::not_found)
602}
603
604fn ensure_user_root_compatible(user: &User, root: &str, label: &str) -> Result<()> {
605    if user.root_ready && user.root_node_id.as_deref() != Some(root) {
606        return Err(Error::conflict(format!(
607            "This {label} already has a different root node."
608        )));
609    }
610    Ok(())
611}
612
613#[cfg(test)]
614mod tests {
615    use super::*;
616    use kcode_tg_kennedy_bot::IdentitySink;
617    use std::sync::{Arc, Barrier};
618    use std::thread;
619
620    fn directory() -> Directory {
621        let database = Connection::open_in_memory().unwrap();
622        database
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        database.execute_batch(IDENTITY_MIGRATION).unwrap();
635        let directory = Directory {
636            database: Mutex::new(database),
637        };
638        directory.seed_bootstrap_user("@taek42").unwrap();
639        directory
640    }
641
642    fn preauthorize(directory: &Directory, handle: &str) {
643        directory.authorize_user_id("taek42", 42).unwrap();
644        assert!(matches!(
645            directory.request_add_user(42, handle).unwrap(),
646            AddUserOutcome::Whitelisted {
647                telegram_user_id: None,
648                ..
649            }
650        ));
651    }
652
653    #[test]
654    fn opens_against_the_identity_schema_created_by_kmap_startup() {
655        let directory = std::env::temp_dir().join(format!(
656            "kennedy-telegram-identity-startup-test-{}",
657            uuid::Uuid::new_v4()
658        ));
659        std::fs::create_dir_all(&directory).unwrap();
660        let user_database = directory.join("users.sqlite3");
661        let kmap_schema = Connection::open(&user_database).unwrap();
662        kmap_schema
663            .execute_batch(
664                "CREATE TABLE kmap_system_roots(
665                     role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
666                     root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
667                     created_at TEXT NOT NULL
668                 );
669                 INSERT INTO kmap_system_roots VALUES(
670                     'user','AAAAAAAB','2026-01-01T00:00:00Z'
671                 );",
672            )
673            .unwrap();
674        drop(kmap_schema);
675
676        let identity = Directory::open(&user_database, "@taek42").unwrap();
677        assert!(
678            identity
679                .lock()
680                .unwrap()
681                .query_row(
682                    "SELECT EXISTS(SELECT 1 FROM whitelist_entries WHERE handle='taek42')",
683                    [],
684                    |row| row.get::<_, i64>(0),
685                )
686                .unwrap()
687                != 0
688        );
689        drop(identity);
690        std::fs::remove_dir_all(directory).unwrap();
691    }
692
693    #[test]
694    fn matching_first_observation_binds_preauthorized_handle() {
695        let directory = directory();
696        preauthorize(&directory, "@friend");
697
698        directory
699            .observe_identity(&IdentityObservation {
700                telegram_user_id: 77,
701                username: Some("@FrIeNd".into()),
702                display_name: "Friend".into(),
703            })
704            .unwrap();
705
706        let friend = directory.user(77).unwrap();
707        assert_eq!(friend.handle, "friend");
708        assert_eq!(friend.current_username.as_deref(), Some("friend"));
709        assert!(
710            directory
711                .whitelist()
712                .unwrap()
713                .telegram_user_ids
714                .contains(&77)
715        );
716    }
717
718    #[test]
719    fn nonmatching_observation_remains_metadata_only() {
720        let directory = directory();
721        preauthorize(&directory, "friend");
722
723        directory
724            .observe_identity(&IdentityObservation {
725                telegram_user_id: 77,
726                username: Some("stranger".into()),
727                display_name: "Stranger".into(),
728            })
729            .unwrap();
730
731        assert_eq!(directory.user(77).unwrap_err().kind(), ErrorKind::NotFound);
732        assert!(
733            !directory
734                .whitelist()
735                .unwrap()
736                .telegram_user_ids
737                .contains(&77)
738        );
739        assert_eq!(
740            observed_identity_by_id(&directory.lock().unwrap(), 77).unwrap(),
741            Some((Some("stranger".into()), Some("Stranger".into())))
742        );
743    }
744
745    #[test]
746    fn existing_numeric_binding_survives_username_change() {
747        let directory = directory();
748        preauthorize(&directory, "friend");
749        directory
750            .observe_identity(&IdentityObservation {
751                telegram_user_id: 77,
752                username: Some("friend".into()),
753                display_name: "Friend".into(),
754            })
755            .unwrap();
756
757        directory
758            .observe_identity(&IdentityObservation {
759                telegram_user_id: 77,
760                username: Some("renamed".into()),
761                display_name: "Renamed Friend".into(),
762            })
763            .unwrap();
764
765        let friend = directory.user(77).unwrap();
766        assert_eq!(friend.handle, "friend");
767        assert_eq!(friend.current_username.as_deref(), Some("renamed"));
768        assert!(
769            directory
770                .whitelist()
771                .unwrap()
772                .telegram_user_ids
773                .contains(&77)
774        );
775    }
776
777    #[test]
778    fn conflicting_observations_do_not_reassign_id_or_handle() {
779        let directory = directory();
780        preauthorize(&directory, "friend");
781        assert!(matches!(
782            directory.request_add_user(42, "colleague").unwrap(),
783            AddUserOutcome::Whitelisted {
784                telegram_user_id: None,
785                ..
786            }
787        ));
788        directory
789            .observe_identity(&IdentityObservation {
790                telegram_user_id: 77,
791                username: Some("friend".into()),
792                display_name: "Friend".into(),
793            })
794            .unwrap();
795
796        directory
797            .observe_identity(&IdentityObservation {
798                telegram_user_id: 88,
799                username: Some("friend".into()),
800                display_name: "Other".into(),
801            })
802            .unwrap();
803        directory
804            .observe_identity(&IdentityObservation {
805                telegram_user_id: 77,
806                username: Some("colleague".into()),
807                display_name: "Renamed Friend".into(),
808            })
809            .unwrap();
810
811        assert_eq!(directory.user(88).unwrap_err().kind(), ErrorKind::NotFound);
812        let database = directory.lock().unwrap();
813        let friend = directory_user_by_handle(&database, "friend")
814            .unwrap()
815            .unwrap();
816        let colleague = directory_user_by_handle(&database, "colleague")
817            .unwrap()
818            .unwrap();
819        assert_eq!(friend.telegram_user_id, Some(77));
820        assert_eq!(friend.current_username.as_deref(), Some("colleague"));
821        assert_eq!(colleague.telegram_user_id, None);
822    }
823
824    #[test]
825    fn explicit_numeric_authorization_replays_and_conflicts() {
826        let directory = directory();
827        let first = directory.authorize_user_id("taek42", 42).unwrap();
828        let replay = directory.authorize_user_id("@TaEk42", 42).unwrap();
829        assert_eq!(replay, first);
830
831        let mismatch = directory.authorize_user_id("taek42", 43).unwrap_err();
832        assert_eq!(mismatch.kind(), ErrorKind::Conflict);
833
834        {
835            let database = directory.lock().unwrap();
836            whitelist_handle(&database, "friend", 42).unwrap();
837        }
838        let duplicate = directory.authorize_user_id("friend", 42).unwrap_err();
839        assert_eq!(duplicate.kind(), ErrorKind::Conflict);
840    }
841
842    #[test]
843    fn anonymous_group_pseudo_user_is_rejected_at_runtime() {
844        let directory = directory();
845        directory
846            .observe_identity(&IdentityObservation {
847                telegram_user_id: ANONYMOUS_GROUP_USER_ID,
848                username: Some("GroupAnonymousBot".into()),
849                display_name: "Anonymous".into(),
850            })
851            .unwrap();
852        directory
853            .observe_identity(&IdentityObservation {
854                telegram_user_id: 77,
855                username: Some("@GroupAnonymousBot".into()),
856                display_name: "Spoof".into(),
857            })
858            .unwrap();
859
860        let database = directory.lock().unwrap();
861        assert_eq!(
862            database
863                .query_row("SELECT COUNT(*) FROM observed_identities", [], |row| {
864                    row.get::<_, i64>(0)
865                })
866                .unwrap(),
867            0
868        );
869        drop(database);
870
871        assert_eq!(
872            directory
873                .authorize_user_id("taek42", ANONYMOUS_GROUP_USER_ID)
874                .unwrap_err()
875                .kind(),
876            ErrorKind::InvalidInput
877        );
878        assert_eq!(
879            directory
880                .authorize_user_id(ANONYMOUS_GROUP_HANDLE, 77)
881                .unwrap_err()
882                .kind(),
883            ErrorKind::InvalidInput
884        );
885
886        directory.authorize_user_id("taek42", 42).unwrap();
887        assert!(
888            directory
889                .request_add_user(42, ANONYMOUS_GROUP_HANDLE)
890                .is_err()
891        );
892        assert!(
893            !directory
894                .whitelist()
895                .unwrap()
896                .telegram_user_ids
897                .contains(&ANONYMOUS_GROUP_USER_ID)
898        );
899    }
900
901    #[test]
902    fn identity_migration_removes_legacy_anonymous_group_pseudo_user() {
903        let directory = directory();
904        let database = directory.lock().unwrap();
905        database
906            .execute(
907                "INSERT INTO observed_identities(
908                     telegram_user_id,current_username,display_name,first_seen_at,last_seen_at
909                 ) VALUES(1087968824,'GroupAnonymousBot','Group',?1,?1)",
910                [Utc::now().to_rfc3339()],
911            )
912            .unwrap();
913        database.execute_batch(IDENTITY_MIGRATION).unwrap();
914        assert_eq!(
915            database
916                .query_row(
917                    "SELECT COUNT(*) FROM observed_identities WHERE telegram_user_id=1087968824",
918                    [],
919                    |row| row.get::<_, i64>(0),
920                )
921                .unwrap(),
922            0
923        );
924    }
925
926    #[test]
927    fn add_user_capability_and_group_roots_stay_in_kennedy() {
928        let directory = directory();
929        directory.authorize_user_id("taek42", 42).unwrap();
930        assert!(matches!(
931            directory.request_add_user(77, "@friend").unwrap(),
932            AddUserOutcome::Forbidden
933        ));
934        assert!(matches!(
935            directory.request_add_user(42, "@friend").unwrap(),
936            AddUserOutcome::Whitelisted {
937                telegram_user_id: None,
938                ..
939            }
940        ));
941        directory.observe_group("opaque-group").unwrap();
942        let database = directory.lock().unwrap();
943        let group = directory_group_by_id(&database, "opaque-group")
944            .unwrap()
945            .unwrap();
946        assert_eq!(group.root_node_id, None);
947        assert!(!group.root_ready);
948    }
949
950    #[test]
951    fn root_completion_replays_and_conflicts() {
952        let directory = directory();
953        directory.authorize_user_id("taek42", 42).unwrap();
954        directory.observe_group("opaque-group").unwrap();
955
956        let user_root = "AAAAAAAC".parse().unwrap();
957        let user = directory.complete_user_root(42, user_root).unwrap();
958        assert!(user.root_ready);
959        assert_eq!(user.root_node_id.as_deref(), Some("AAAAAAAC"));
960        assert_eq!(
961            directory
962                .complete_user_root(42, user_root)
963                .unwrap()
964                .root_node_id
965                .as_deref(),
966            Some("AAAAAAAC")
967        );
968
969        let group_root = "AAAAAAAD".parse().unwrap();
970        let group = directory
971            .complete_group_root("opaque-group", group_root)
972            .unwrap();
973        assert!(group.root_ready);
974        assert_eq!(group.root_node_id.as_deref(), Some("AAAAAAAD"));
975        assert_eq!(
976            directory
977                .complete_group_root("opaque-group", group_root)
978                .unwrap()
979                .root_node_id
980                .as_deref(),
981            Some("AAAAAAAD")
982        );
983        assert_eq!(
984            directory.group_for_root(group_root).unwrap().group_id,
985            "opaque-group"
986        );
987        assert_eq!(
988            directory
989                .group_for_root("AAAAAAAE".parse().unwrap())
990                .unwrap_err()
991                .kind(),
992            ErrorKind::NotFound
993        );
994
995        let mismatch = directory
996            .complete_group_root("opaque-group", "AAAAAAAE".parse().unwrap())
997            .unwrap_err();
998        assert_eq!(mismatch.kind(), ErrorKind::Conflict);
999    }
1000
1001    #[test]
1002    fn independent_directories_bind_handle_atomically() {
1003        let directory = std::env::temp_dir().join(format!(
1004            "kennedy-telegram-identity-binding-concurrency-test-{}",
1005            uuid::Uuid::new_v4()
1006        ));
1007        std::fs::create_dir_all(&directory).unwrap();
1008        let database_path = directory.join("users.sqlite3");
1009
1010        let first = Directory::open(&database_path, "taek42").unwrap();
1011        preauthorize(&first, "friend");
1012        let first = Arc::new(first);
1013        let second = Arc::new(Directory::open(&database_path, "taek42").unwrap());
1014        let barrier = Arc::new(Barrier::new(2));
1015
1016        let first_worker = {
1017            let directory = Arc::clone(&first);
1018            let barrier = Arc::clone(&barrier);
1019            thread::spawn(move || {
1020                barrier.wait();
1021                directory.observe_identity(&IdentityObservation {
1022                    telegram_user_id: 77,
1023                    username: Some("friend".into()),
1024                    display_name: "First".into(),
1025                })
1026            })
1027        };
1028        let second_worker = {
1029            let directory = Arc::clone(&second);
1030            let barrier = Arc::clone(&barrier);
1031            thread::spawn(move || {
1032                barrier.wait();
1033                directory.observe_identity(&IdentityObservation {
1034                    telegram_user_id: 88,
1035                    username: Some("friend".into()),
1036                    display_name: "Second".into(),
1037                })
1038            })
1039        };
1040
1041        first_worker.join().unwrap().unwrap();
1042        second_worker.join().unwrap().unwrap();
1043
1044        let snapshot = first.whitelist().unwrap();
1045        assert_eq!(
1046            [77, 88]
1047                .into_iter()
1048                .filter(|id| snapshot.telegram_user_ids.contains(id))
1049                .count(),
1050            1
1051        );
1052        let database = first.lock().unwrap();
1053        let friend = directory_user_by_handle(&database, "friend")
1054            .unwrap()
1055            .unwrap();
1056        assert!(matches!(friend.telegram_user_id, Some(77 | 88)));
1057        assert_eq!(
1058            database
1059                .query_row("SELECT COUNT(*) FROM observed_identities", [], |row| {
1060                    row.get::<_, i64>(0)
1061                })
1062                .unwrap(),
1063            2
1064        );
1065        drop(database);
1066
1067        drop(first);
1068        drop(second);
1069        std::fs::remove_dir_all(directory).unwrap();
1070    }
1071
1072    #[test]
1073    fn independent_directories_complete_roots_atomically() {
1074        let directory = std::env::temp_dir().join(format!(
1075            "kennedy-telegram-identity-concurrency-test-{}",
1076            uuid::Uuid::new_v4()
1077        ));
1078        std::fs::create_dir_all(&directory).unwrap();
1079        let database_path = directory.join("users.sqlite3");
1080
1081        let first = Directory::open(&database_path, "taek42").unwrap();
1082        first.authorize_user_id("taek42", 42).unwrap();
1083        let first = Arc::new(first);
1084        let second = Arc::new(Directory::open(&database_path, "taek42").unwrap());
1085        let barrier = Arc::new(Barrier::new(2));
1086
1087        let first_worker = {
1088            let directory = Arc::clone(&first);
1089            let barrier = Arc::clone(&barrier);
1090            thread::spawn(move || {
1091                barrier.wait();
1092                directory
1093                    .complete_user_root(42, "AAAAAAAC".parse().unwrap())
1094                    .map(|user| user.root_node_id.unwrap())
1095                    .map_err(|error| error.kind())
1096            })
1097        };
1098        let second_worker = {
1099            let directory = Arc::clone(&second);
1100            let barrier = Arc::clone(&barrier);
1101            thread::spawn(move || {
1102                barrier.wait();
1103                directory
1104                    .complete_user_root(42, "AAAAAAAD".parse().unwrap())
1105                    .map(|user| user.root_node_id.unwrap())
1106                    .map_err(|error| error.kind())
1107            })
1108        };
1109
1110        let outcomes = [first_worker.join().unwrap(), second_worker.join().unwrap()];
1111        assert_eq!(outcomes.iter().filter(|outcome| outcome.is_ok()).count(), 1);
1112        assert_eq!(
1113            outcomes
1114                .iter()
1115                .filter(|outcome| matches!(outcome, Err(ErrorKind::Conflict)))
1116                .count(),
1117            1
1118        );
1119
1120        drop(first);
1121        drop(second);
1122        std::fs::remove_dir_all(directory).unwrap();
1123    }
1124}