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 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}