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, params};
9use serde::{Deserialize, Serialize};
10
11const IDENTITY_MIGRATION: &str = include_str!("../migrations/001_initial.sql");
12
13pub struct Directory {
15 database: Mutex<Connection>,
16}
17
18#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
20#[serde(rename_all = "camelCase")]
21pub struct User {
22 pub handle: String,
23 pub telegram_user_id: Option<i64>,
24 pub current_username: Option<String>,
25 pub display_name: Option<String>,
26 pub root_node_id: Option<String>,
27 pub root_ready: bool,
28 pub can_add_users: bool,
29}
30
31#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
33#[serde(rename_all = "camelCase")]
34pub struct Group {
35 pub group_id: String,
36 pub root_node_id: Option<String>,
37 pub root_ready: bool,
38}
39
40#[derive(Clone, Copy, Debug, Eq, PartialEq)]
42pub enum ErrorKind {
43 InvalidInput,
44 NotFound,
45 Conflict,
46 Storage,
47}
48
49#[derive(Debug, thiserror::Error)]
51#[error("{message}")]
52pub struct Error {
53 kind: ErrorKind,
54 message: String,
55}
56
57impl Error {
58 pub fn kind(&self) -> ErrorKind {
59 self.kind
60 }
61
62 pub fn message(&self) -> &str {
63 &self.message
64 }
65
66 fn invalid(message: impl Into<String>) -> Self {
67 Self {
68 kind: ErrorKind::InvalidInput,
69 message: message.into(),
70 }
71 }
72
73 fn not_found() -> Self {
74 Self {
75 kind: ErrorKind::NotFound,
76 message: "Telegram directory entry not found.".into(),
77 }
78 }
79
80 fn conflict(message: impl Into<String>) -> Self {
81 Self {
82 kind: ErrorKind::Conflict,
83 message: message.into(),
84 }
85 }
86
87 fn storage(error: impl std::fmt::Display) -> Self {
88 Self {
89 kind: ErrorKind::Storage,
90 message: error.to_string(),
91 }
92 }
93}
94
95pub type Result<T> = std::result::Result<T, Error>;
96
97impl Directory {
98 pub fn open(path: &Path, bootstrap_handle: &str) -> Result<Self> {
100 let connection = Connection::open(path)
101 .map_err(|error| Error::storage(format!("opening {}: {error}", path.display())))?;
102 connection
103 .execute_batch(
104 "PRAGMA journal_mode=WAL; PRAGMA busy_timeout=15000; PRAGMA foreign_keys=ON;",
105 )
106 .map_err(Error::storage)?;
107 connection
108 .execute_batch(IDENTITY_MIGRATION)
109 .map_err(Error::storage)?;
110 let directory = Self {
111 database: Mutex::new(connection),
112 };
113 directory.seed_bootstrap_user(bootstrap_handle)?;
114 Ok(directory)
115 }
116
117 fn seed_bootstrap_user(&self, handle: &str) -> Result<()> {
118 let handle = normalize_username(handle);
119 if handle.is_empty() {
120 return Err(Error::invalid(
121 "Telegram bootstrap handle must not be empty",
122 ));
123 }
124 let database = self.lock()?;
125 let now = Utc::now().to_rfc3339();
126 database
127 .execute(
128 "INSERT INTO whitelist_entries(handle,can_add_users,whitelisted_at,updated_at)
129 VALUES(?1,1,?2,?2)
130 ON CONFLICT(handle) DO UPDATE SET can_add_users=1,updated_at=excluded.updated_at",
131 params![handle, now],
132 )
133 .map_err(Error::storage)?;
134 Ok(())
135 }
136
137 fn lock(&self) -> Result<std::sync::MutexGuard<'_, Connection>> {
138 self.database
139 .lock()
140 .map_err(|_| Error::storage("locking Telegram identity directory"))
141 }
142
143 pub fn provisioning_users(&self) -> Result<Vec<User>> {
145 let database = self.lock()?;
146 let mut statement = database
147 .prepare(
148 "SELECT handle,telegram_user_id,current_username,display_name,root_node_id,root_ready,can_add_users
149 FROM whitelist_entries WHERE root_ready=0 ORDER BY whitelisted_at,handle",
150 )
151 .map_err(Error::storage)?;
152 statement
153 .query_map([], row_user)
154 .map_err(Error::storage)?
155 .collect::<std::result::Result<Vec<_>, _>>()
156 .map_err(Error::storage)
157 }
158
159 pub fn provisioning_groups(&self) -> Result<Vec<Group>> {
161 let database = self.lock()?;
162 let mut statement = database
163 .prepare(
164 "SELECT group_id,root_node_id,root_ready FROM telegram_group_roots
165 WHERE root_ready=0 ORDER BY datetime(created_at),group_id",
166 )
167 .map_err(Error::storage)?;
168 statement
169 .query_map([], row_group)
170 .map_err(Error::storage)?
171 .collect::<std::result::Result<Vec<_>, _>>()
172 .map_err(Error::storage)
173 }
174
175 pub fn user(&self, telegram_user_id: i64) -> Result<User> {
177 let database = self.lock()?;
178 directory_user_by_id(&database, telegram_user_id)?.ok_or_else(Error::not_found)
179 }
180
181 pub fn group(&self, group_id: &str) -> Result<Group> {
183 let database = self.lock()?;
184 directory_group_by_id(&database, group_id)?.ok_or_else(Error::not_found)
185 }
186
187 pub fn group_for_root(&self, root_node_id: NodeId) -> Result<Group> {
189 let database = self.lock()?;
190 database
191 .query_row(
192 "SELECT group_id,root_node_id,root_ready FROM telegram_group_roots
193 WHERE root_node_id=?1 AND root_ready=1",
194 [root_node_id.to_string()],
195 row_group,
196 )
197 .optional()
198 .map_err(Error::storage)?
199 .ok_or_else(Error::not_found)
200 }
201
202 pub fn complete_handle_root(&self, handle: &str, root_node_id: NodeId) -> Result<User> {
204 let handle = normalize_username(handle);
205 let database = self.lock()?;
206 let current = directory_user_by_handle(&database, &handle)?.ok_or_else(Error::not_found)?;
207 ensure_user_root_compatible(¤t, root_node_id, "whitelisted handle")?;
208 database
209 .execute(
210 "UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
211 WHERE handle=?3",
212 params![root_node_id.to_string(), Utc::now().to_rfc3339(), handle],
213 )
214 .map_err(Error::storage)?;
215 directory_user_by_handle(&database, &handle)?.ok_or_else(Error::not_found)
216 }
217
218 pub fn complete_user_root(&self, telegram_user_id: i64, root_node_id: NodeId) -> Result<User> {
220 let database = self.lock()?;
221 let current =
222 directory_user_by_id(&database, telegram_user_id)?.ok_or_else(Error::not_found)?;
223 ensure_user_root_compatible(¤t, root_node_id, "Telegram identity")?;
224 database
225 .execute(
226 "UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
227 WHERE telegram_user_id=?3",
228 params![
229 root_node_id.to_string(),
230 Utc::now().to_rfc3339(),
231 telegram_user_id
232 ],
233 )
234 .map_err(Error::storage)?;
235 directory_user_by_id(&database, telegram_user_id)?.ok_or_else(Error::not_found)
236 }
237
238 pub fn complete_group_root(&self, group_id: &str, root_node_id: NodeId) -> Result<Group> {
240 let database = self.lock()?;
241 let current = directory_group_by_id(&database, group_id)?.ok_or_else(Error::not_found)?;
242 let root_node_id = root_node_id.to_string();
243 if current.root_ready && current.root_node_id.as_deref() != Some(&root_node_id) {
244 return Err(Error::conflict(
245 "This Telegram group already has a different root node.",
246 ));
247 }
248 database
249 .execute(
250 "UPDATE telegram_group_roots SET root_node_id=?1,root_ready=1,updated_at=?2
251 WHERE group_id=?3",
252 params![root_node_id, Utc::now().to_rfc3339(), group_id],
253 )
254 .map_err(Error::storage)?;
255 directory_group_by_id(&database, group_id)?.ok_or_else(Error::not_found)
256 }
257}
258
259impl IdentitySink for Directory {
260 fn observe_identity(&self, observation: &IdentityObservation) -> anyhow::Result<()> {
261 let database = self.lock()?;
262 observe_identity(&database, observation)?;
263 Ok(())
264 }
265
266 fn whitelist(&self) -> anyhow::Result<WhitelistSnapshot> {
267 let database = self.lock()?;
268 let telegram_user_ids = database
269 .prepare(
270 "SELECT telegram_user_id FROM whitelist_entries
271 WHERE telegram_user_id IS NOT NULL ORDER BY telegram_user_id",
272 )?
273 .query_map([], |row| row.get::<_, i64>(0))?
274 .collect::<std::result::Result<HashSet<_>, _>>()?;
275 Ok(WhitelistSnapshot { telegram_user_ids })
276 }
277
278 fn request_add_user(
279 &self,
280 requested_by_telegram_user_id: i64,
281 handle: &str,
282 ) -> anyhow::Result<AddUserOutcome> {
283 let database = self.lock()?;
284 let can_add = directory_user_by_id(&database, requested_by_telegram_user_id)?
285 .is_some_and(|user| user.can_add_users);
286 if !can_add {
287 return Ok(AddUserOutcome::Forbidden);
288 }
289 let user = whitelist_handle(&database, handle, requested_by_telegram_user_id)?;
290 Ok(AddUserOutcome::Whitelisted {
291 handle: user.handle,
292 telegram_user_id: user.telegram_user_id,
293 })
294 }
295
296 fn observe_group(&self, group_id: &str) -> anyhow::Result<()> {
297 let database = self.lock()?;
298 let now = Utc::now().to_rfc3339();
299 database.execute(
300 "INSERT INTO telegram_group_roots(group_id,created_at,updated_at)
301 VALUES(?1,?2,?2) ON CONFLICT(group_id) DO NOTHING",
302 params![group_id, now],
303 )?;
304 Ok(())
305 }
306}
307
308fn normalize_username(value: &str) -> String {
309 value.trim().trim_start_matches('@').to_ascii_lowercase()
310}
311
312fn directory_user_by_clause(
313 database: &Connection,
314 clause: &str,
315 value: &dyn rusqlite::ToSql,
316) -> Result<Option<User>> {
317 database
318 .query_row(
319 &format!(
320 "SELECT handle,telegram_user_id,current_username,display_name,root_node_id,root_ready,can_add_users
321 FROM whitelist_entries WHERE {clause}"
322 ),
323 [value],
324 row_user,
325 )
326 .optional()
327 .map_err(Error::storage)
328}
329
330fn directory_user_by_id(database: &Connection, telegram_user_id: i64) -> Result<Option<User>> {
331 directory_user_by_clause(database, "telegram_user_id=?1", &telegram_user_id)
332}
333
334fn directory_user_by_handle(database: &Connection, handle: &str) -> Result<Option<User>> {
335 directory_user_by_clause(database, "handle=?1", &handle)
336}
337
338fn directory_group_by_id(database: &Connection, group_id: &str) -> Result<Option<Group>> {
339 database
340 .query_row(
341 "SELECT group_id,root_node_id,root_ready FROM telegram_group_roots WHERE group_id=?1",
342 [group_id],
343 row_group,
344 )
345 .optional()
346 .map_err(Error::storage)
347}
348
349fn row_user(row: &rusqlite::Row<'_>) -> rusqlite::Result<User> {
350 let root_node_id = canonical_root(row.get(4)?, 4)?;
351 Ok(User {
352 handle: row.get(0)?,
353 telegram_user_id: row.get(1)?,
354 current_username: row.get(2)?,
355 display_name: row.get(3)?,
356 root_node_id,
357 root_ready: row.get::<_, i64>(5)? != 0,
358 can_add_users: row.get::<_, i64>(6)? != 0,
359 })
360}
361
362fn row_group(row: &rusqlite::Row<'_>) -> rusqlite::Result<Group> {
363 let root_node_id = canonical_root(row.get(1)?, 1)?;
364 Ok(Group {
365 group_id: row.get(0)?,
366 root_node_id,
367 root_ready: row.get::<_, i64>(2)? != 0,
368 })
369}
370
371fn canonical_root(value: Option<String>, column: usize) -> rusqlite::Result<Option<String>> {
372 value
373 .map(|value| {
374 value
375 .parse::<NodeId>()
376 .map(|id| id.to_string())
377 .map_err(|error| {
378 rusqlite::Error::FromSqlConversionFailure(
379 column,
380 rusqlite::types::Type::Text,
381 Box::new(error),
382 )
383 })
384 })
385 .transpose()
386}
387
388fn observe_identity(database: &Connection, observation: &IdentityObservation) -> Result<()> {
389 let now = Utc::now().to_rfc3339();
390 let normalized = observation
391 .username
392 .as_deref()
393 .map(normalize_username)
394 .filter(|value| !value.is_empty());
395 database
396 .execute(
397 "INSERT INTO observed_identities(telegram_user_id,current_username,display_name,first_seen_at,last_seen_at)
398 VALUES(?1,?2,?3,?4,?4)
399 ON CONFLICT(telegram_user_id) DO UPDATE SET
400 current_username=excluded.current_username,
401 display_name=excluded.display_name,last_seen_at=excluded.last_seen_at",
402 params![
403 observation.telegram_user_id,
404 normalized,
405 observation.display_name,
406 now
407 ],
408 )
409 .map_err(Error::storage)?;
410 if directory_user_by_id(database, observation.telegram_user_id)?.is_some() {
411 database
412 .execute(
413 "UPDATE whitelist_entries SET current_username=?1,display_name=?2,updated_at=?3
414 WHERE telegram_user_id=?4",
415 params![
416 normalized,
417 observation.display_name,
418 now,
419 observation.telegram_user_id
420 ],
421 )
422 .map_err(Error::storage)?;
423 return Ok(());
424 }
425 let Some(handle) = normalized else {
426 return Ok(());
427 };
428 let Some(entry) = directory_user_by_handle(database, &handle)? else {
429 return Ok(());
430 };
431 if entry.telegram_user_id.is_some() {
432 return Ok(());
433 }
434 database
435 .execute(
436 "UPDATE whitelist_entries SET telegram_user_id=?1,current_username=?2,display_name=?3,
437 resolved_at=?4,updated_at=?4 WHERE handle=?2 AND telegram_user_id IS NULL",
438 params![
439 observation.telegram_user_id,
440 handle,
441 observation.display_name,
442 now
443 ],
444 )
445 .map_err(Error::storage)?;
446 Ok(())
447}
448
449fn whitelist_handle(database: &Connection, handle: &str, added_by: i64) -> Result<User> {
450 let handle = normalize_username(handle.trim_matches(['\'', '"']));
451 if handle.is_empty() {
452 return Err(Error::invalid("the Telegram handle must not be empty"));
453 }
454 let now = Utc::now().to_rfc3339();
455 database
456 .execute(
457 "INSERT INTO whitelist_entries(handle,current_username,added_by_telegram_user_id,whitelisted_at,updated_at)
458 VALUES(?1,?1,?2,?3,?3)
459 ON CONFLICT(handle) DO UPDATE SET updated_at=excluded.updated_at",
460 params![handle, added_by, now],
461 )
462 .map_err(Error::storage)?;
463 directory_user_by_handle(database, &handle)?.ok_or_else(Error::not_found)
464}
465
466fn ensure_user_root_compatible(user: &User, root: NodeId, label: &str) -> Result<()> {
467 let root = root.to_string();
468 if user.root_ready && user.root_node_id.as_deref() != Some(&root) {
469 return Err(Error::conflict(format!(
470 "This {label} already has a different root node."
471 )));
472 }
473 Ok(())
474}
475
476#[cfg(test)]
477mod tests {
478 use super::*;
479 use kcode_tg_kennedy_bot::IdentitySink;
480
481 fn directory() -> Directory {
482 let database = Connection::open_in_memory().unwrap();
483 database
484 .execute_batch(
485 "CREATE TABLE kmap_system_roots(
486 role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
487 root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
488 created_at TEXT NOT NULL
489 );
490 INSERT INTO kmap_system_roots VALUES(
491 'user','AAAAAAAB','2026-01-01T00:00:00Z'
492 );",
493 )
494 .unwrap();
495 database.execute_batch(IDENTITY_MIGRATION).unwrap();
496 let directory = Directory {
497 database: Mutex::new(database),
498 };
499 directory.seed_bootstrap_user("@taek42").unwrap();
500 directory
501 }
502
503 #[test]
504 fn opens_against_the_identity_schema_created_by_kmap_startup() {
505 let directory = std::env::temp_dir().join(format!(
506 "kennedy-telegram-identity-startup-test-{}",
507 uuid::Uuid::new_v4()
508 ));
509 std::fs::create_dir_all(&directory).unwrap();
510 let user_database = directory.join("users.sqlite3");
511 let kmap_schema = Connection::open(&user_database).unwrap();
512 kmap_schema
513 .execute_batch(
514 "CREATE TABLE kmap_system_roots(
515 role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
516 root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
517 created_at TEXT NOT NULL
518 );
519 INSERT INTO kmap_system_roots VALUES(
520 'user','AAAAAAAB','2026-01-01T00:00:00Z'
521 );",
522 )
523 .unwrap();
524 drop(kmap_schema);
525
526 let identity = Directory::open(&user_database, "@taek42").unwrap();
527 assert!(
528 identity
529 .lock()
530 .unwrap()
531 .query_row(
532 "SELECT EXISTS(SELECT 1 FROM whitelist_entries WHERE handle='taek42')",
533 [],
534 |row| row.get::<_, i64>(0),
535 )
536 .unwrap()
537 != 0
538 );
539 drop(identity);
540 std::fs::remove_dir_all(directory).unwrap();
541 }
542
543 #[test]
544 fn tofu_is_owned_by_kennedy_and_numeric_ids_remain_authoritative() {
545 let directory = directory();
546 directory
547 .observe_identity(&IdentityObservation {
548 telegram_user_id: 42,
549 username: Some("TaEk42".into()),
550 display_name: "David".into(),
551 })
552 .unwrap();
553 assert!(
554 directory
555 .whitelist()
556 .unwrap()
557 .telegram_user_ids
558 .contains(&42)
559 );
560 directory
561 .observe_identity(&IdentityObservation {
562 telegram_user_id: 43,
563 username: Some("taek42".into()),
564 display_name: "Other".into(),
565 })
566 .unwrap();
567 assert!(
568 !directory
569 .whitelist()
570 .unwrap()
571 .telegram_user_ids
572 .contains(&43)
573 );
574 }
575
576 #[test]
577 fn identity_migration_removes_legacy_anonymous_group_pseudo_user() {
578 let directory = directory();
579 let database = directory.lock().unwrap();
580 database
581 .execute(
582 "INSERT INTO observed_identities(
583 telegram_user_id,current_username,display_name,first_seen_at,last_seen_at
584 ) VALUES(1087968824,'GroupAnonymousBot','Group',?1,?1)",
585 [Utc::now().to_rfc3339()],
586 )
587 .unwrap();
588 database.execute_batch(IDENTITY_MIGRATION).unwrap();
589 assert_eq!(
590 database
591 .query_row(
592 "SELECT COUNT(*) FROM observed_identities WHERE telegram_user_id=1087968824",
593 [],
594 |row| row.get::<_, i64>(0),
595 )
596 .unwrap(),
597 0
598 );
599 }
600
601 #[test]
602 fn add_user_capability_and_group_roots_stay_in_kennedy() {
603 let directory = directory();
604 directory
605 .observe_identity(&IdentityObservation {
606 telegram_user_id: 42,
607 username: Some("taek42".into()),
608 display_name: "David".into(),
609 })
610 .unwrap();
611 assert!(matches!(
612 directory.request_add_user(77, "@friend").unwrap(),
613 AddUserOutcome::Forbidden
614 ));
615 assert!(matches!(
616 directory.request_add_user(42, "@friend").unwrap(),
617 AddUserOutcome::Whitelisted { .. }
618 ));
619 directory.observe_group("opaque-group").unwrap();
620 let database = directory.lock().unwrap();
621 let group = directory_group_by_id(&database, "opaque-group")
622 .unwrap()
623 .unwrap();
624 assert_eq!(group.root_node_id, None);
625 assert!(!group.root_ready);
626 }
627
628 #[test]
629 fn root_completion_is_owned_and_conflict_checked_by_kennedy() {
630 let directory = directory();
631 directory
632 .observe_identity(&IdentityObservation {
633 telegram_user_id: 42,
634 username: Some("taek42".into()),
635 display_name: "David".into(),
636 })
637 .unwrap();
638 directory.observe_group("opaque-group").unwrap();
639
640 let user_root = "AAAAAAAC".parse().unwrap();
641 let user = directory.complete_user_root(42, user_root).unwrap();
642 assert!(user.root_ready);
643 assert_eq!(user.root_node_id.as_deref(), Some("AAAAAAAC"));
644
645 let group_root = "AAAAAAAD".parse().unwrap();
646 let group = directory
647 .complete_group_root("opaque-group", group_root)
648 .unwrap();
649 assert!(group.root_ready);
650 assert_eq!(group.root_node_id.as_deref(), Some("AAAAAAAD"));
651 assert_eq!(
652 directory.group_for_root(group_root).unwrap().group_id,
653 "opaque-group"
654 );
655 assert_eq!(
656 directory
657 .group_for_root("AAAAAAAE".parse().unwrap())
658 .unwrap_err()
659 .kind(),
660 ErrorKind::NotFound
661 );
662
663 let mismatch = directory
664 .complete_group_root("opaque-group", "AAAAAAAE".parse().unwrap())
665 .unwrap_err();
666 assert_eq!(mismatch.kind(), ErrorKind::Conflict);
667 }
668}