1use 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
15pub struct Directory {
17 database: Mutex<Connection>,
18}
19
20#[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#[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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
44pub enum ErrorKind {
45 InvalidInput,
46 NotFound,
47 Conflict,
48 Storage,
49}
50
51#[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 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 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 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 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 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 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 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 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(¤t, &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 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(¤t, &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 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}