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=5000; 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 complete_handle_root(&self, handle: &str, root_node_id: NodeId) -> Result<User> {
189 let handle = normalize_username(handle);
190 let database = self.lock()?;
191 let current = directory_user_by_handle(&database, &handle)?.ok_or_else(Error::not_found)?;
192 ensure_user_root_compatible(¤t, root_node_id, "whitelisted handle")?;
193 database
194 .execute(
195 "UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
196 WHERE handle=?3",
197 params![root_node_id.to_string(), Utc::now().to_rfc3339(), handle],
198 )
199 .map_err(Error::storage)?;
200 directory_user_by_handle(&database, &handle)?.ok_or_else(Error::not_found)
201 }
202
203 pub fn complete_user_root(&self, telegram_user_id: i64, root_node_id: NodeId) -> Result<User> {
205 let database = self.lock()?;
206 let current =
207 directory_user_by_id(&database, telegram_user_id)?.ok_or_else(Error::not_found)?;
208 ensure_user_root_compatible(¤t, root_node_id, "Telegram identity")?;
209 database
210 .execute(
211 "UPDATE whitelist_entries SET root_node_id=?1,root_ready=1,updated_at=?2
212 WHERE telegram_user_id=?3",
213 params![
214 root_node_id.to_string(),
215 Utc::now().to_rfc3339(),
216 telegram_user_id
217 ],
218 )
219 .map_err(Error::storage)?;
220 directory_user_by_id(&database, telegram_user_id)?.ok_or_else(Error::not_found)
221 }
222
223 pub fn complete_group_root(&self, group_id: &str, root_node_id: NodeId) -> Result<Group> {
225 let database = self.lock()?;
226 let current = directory_group_by_id(&database, group_id)?.ok_or_else(Error::not_found)?;
227 let root_node_id = root_node_id.to_string();
228 if current.root_ready && current.root_node_id.as_deref() != Some(&root_node_id) {
229 return Err(Error::conflict(
230 "This Telegram group already has a different root node.",
231 ));
232 }
233 database
234 .execute(
235 "UPDATE telegram_group_roots SET root_node_id=?1,root_ready=1,updated_at=?2
236 WHERE group_id=?3",
237 params![root_node_id, Utc::now().to_rfc3339(), group_id],
238 )
239 .map_err(Error::storage)?;
240 directory_group_by_id(&database, group_id)?.ok_or_else(Error::not_found)
241 }
242}
243
244impl IdentitySink for Directory {
245 fn observe_identity(&self, observation: &IdentityObservation) -> anyhow::Result<()> {
246 let database = self.lock()?;
247 observe_identity(&database, observation)?;
248 Ok(())
249 }
250
251 fn whitelist(&self) -> anyhow::Result<WhitelistSnapshot> {
252 let database = self.lock()?;
253 let telegram_user_ids = database
254 .prepare(
255 "SELECT telegram_user_id FROM whitelist_entries
256 WHERE telegram_user_id IS NOT NULL ORDER BY telegram_user_id",
257 )?
258 .query_map([], |row| row.get::<_, i64>(0))?
259 .collect::<std::result::Result<HashSet<_>, _>>()?;
260 Ok(WhitelistSnapshot { telegram_user_ids })
261 }
262
263 fn request_add_user(
264 &self,
265 requested_by_telegram_user_id: i64,
266 handle: &str,
267 ) -> anyhow::Result<AddUserOutcome> {
268 let database = self.lock()?;
269 let can_add = directory_user_by_id(&database, requested_by_telegram_user_id)?
270 .is_some_and(|user| user.can_add_users);
271 if !can_add {
272 return Ok(AddUserOutcome::Forbidden);
273 }
274 let user = whitelist_handle(&database, handle, requested_by_telegram_user_id)?;
275 Ok(AddUserOutcome::Whitelisted {
276 handle: user.handle,
277 telegram_user_id: user.telegram_user_id,
278 })
279 }
280
281 fn observe_group(&self, group_id: &str) -> anyhow::Result<()> {
282 let database = self.lock()?;
283 let now = Utc::now().to_rfc3339();
284 database.execute(
285 "INSERT INTO telegram_group_roots(group_id,created_at,updated_at)
286 VALUES(?1,?2,?2) ON CONFLICT(group_id) DO NOTHING",
287 params![group_id, now],
288 )?;
289 Ok(())
290 }
291}
292
293fn normalize_username(value: &str) -> String {
294 value.trim().trim_start_matches('@').to_ascii_lowercase()
295}
296
297fn directory_user_by_clause(
298 database: &Connection,
299 clause: &str,
300 value: &dyn rusqlite::ToSql,
301) -> Result<Option<User>> {
302 database
303 .query_row(
304 &format!(
305 "SELECT handle,telegram_user_id,current_username,display_name,root_node_id,root_ready,can_add_users
306 FROM whitelist_entries WHERE {clause}"
307 ),
308 [value],
309 row_user,
310 )
311 .optional()
312 .map_err(Error::storage)
313}
314
315fn directory_user_by_id(database: &Connection, telegram_user_id: i64) -> Result<Option<User>> {
316 directory_user_by_clause(database, "telegram_user_id=?1", &telegram_user_id)
317}
318
319fn directory_user_by_handle(database: &Connection, handle: &str) -> Result<Option<User>> {
320 directory_user_by_clause(database, "handle=?1", &handle)
321}
322
323fn directory_group_by_id(database: &Connection, group_id: &str) -> Result<Option<Group>> {
324 database
325 .query_row(
326 "SELECT group_id,root_node_id,root_ready FROM telegram_group_roots WHERE group_id=?1",
327 [group_id],
328 row_group,
329 )
330 .optional()
331 .map_err(Error::storage)
332}
333
334fn row_user(row: &rusqlite::Row<'_>) -> rusqlite::Result<User> {
335 let root_node_id = canonical_root(row.get(4)?, 4)?;
336 Ok(User {
337 handle: row.get(0)?,
338 telegram_user_id: row.get(1)?,
339 current_username: row.get(2)?,
340 display_name: row.get(3)?,
341 root_node_id,
342 root_ready: row.get::<_, i64>(5)? != 0,
343 can_add_users: row.get::<_, i64>(6)? != 0,
344 })
345}
346
347fn row_group(row: &rusqlite::Row<'_>) -> rusqlite::Result<Group> {
348 let root_node_id = canonical_root(row.get(1)?, 1)?;
349 Ok(Group {
350 group_id: row.get(0)?,
351 root_node_id,
352 root_ready: row.get::<_, i64>(2)? != 0,
353 })
354}
355
356fn canonical_root(value: Option<String>, column: usize) -> rusqlite::Result<Option<String>> {
357 value
358 .map(|value| {
359 value
360 .parse::<NodeId>()
361 .map(|id| id.to_string())
362 .map_err(|error| {
363 rusqlite::Error::FromSqlConversionFailure(
364 column,
365 rusqlite::types::Type::Text,
366 Box::new(error),
367 )
368 })
369 })
370 .transpose()
371}
372
373fn observe_identity(database: &Connection, observation: &IdentityObservation) -> Result<()> {
374 let now = Utc::now().to_rfc3339();
375 let normalized = observation
376 .username
377 .as_deref()
378 .map(normalize_username)
379 .filter(|value| !value.is_empty());
380 database
381 .execute(
382 "INSERT INTO observed_identities(telegram_user_id,current_username,display_name,first_seen_at,last_seen_at)
383 VALUES(?1,?2,?3,?4,?4)
384 ON CONFLICT(telegram_user_id) DO UPDATE SET
385 current_username=excluded.current_username,
386 display_name=excluded.display_name,last_seen_at=excluded.last_seen_at",
387 params![
388 observation.telegram_user_id,
389 normalized,
390 observation.display_name,
391 now
392 ],
393 )
394 .map_err(Error::storage)?;
395 if directory_user_by_id(database, observation.telegram_user_id)?.is_some() {
396 database
397 .execute(
398 "UPDATE whitelist_entries SET current_username=?1,display_name=?2,updated_at=?3
399 WHERE telegram_user_id=?4",
400 params![
401 normalized,
402 observation.display_name,
403 now,
404 observation.telegram_user_id
405 ],
406 )
407 .map_err(Error::storage)?;
408 return Ok(());
409 }
410 let Some(handle) = normalized else {
411 return Ok(());
412 };
413 let Some(entry) = directory_user_by_handle(database, &handle)? else {
414 return Ok(());
415 };
416 if entry.telegram_user_id.is_some() {
417 return Ok(());
418 }
419 database
420 .execute(
421 "UPDATE whitelist_entries SET telegram_user_id=?1,current_username=?2,display_name=?3,
422 resolved_at=?4,updated_at=?4 WHERE handle=?2 AND telegram_user_id IS NULL",
423 params![
424 observation.telegram_user_id,
425 handle,
426 observation.display_name,
427 now
428 ],
429 )
430 .map_err(Error::storage)?;
431 Ok(())
432}
433
434fn whitelist_handle(database: &Connection, handle: &str, added_by: i64) -> Result<User> {
435 let handle = normalize_username(handle.trim_matches(['\'', '"']));
436 if handle.is_empty() {
437 return Err(Error::invalid("the Telegram handle must not be empty"));
438 }
439 let now = Utc::now().to_rfc3339();
440 database
441 .execute(
442 "INSERT INTO whitelist_entries(handle,current_username,added_by_telegram_user_id,whitelisted_at,updated_at)
443 VALUES(?1,?1,?2,?3,?3)
444 ON CONFLICT(handle) DO UPDATE SET updated_at=excluded.updated_at",
445 params![handle, added_by, now],
446 )
447 .map_err(Error::storage)?;
448 directory_user_by_handle(database, &handle)?.ok_or_else(Error::not_found)
449}
450
451fn ensure_user_root_compatible(user: &User, root: NodeId, label: &str) -> Result<()> {
452 let root = root.to_string();
453 if user.root_ready && user.root_node_id.as_deref() != Some(&root) {
454 return Err(Error::conflict(format!(
455 "This {label} already has a different root node."
456 )));
457 }
458 Ok(())
459}
460
461#[cfg(test)]
462mod tests {
463 use super::*;
464 use kcode_tg_kennedy_bot::IdentitySink;
465
466 fn directory() -> Directory {
467 let database = Connection::open_in_memory().unwrap();
468 database
469 .execute_batch(
470 "CREATE TABLE kmap_system_roots(
471 role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
472 root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
473 created_at TEXT NOT NULL
474 );
475 INSERT INTO kmap_system_roots VALUES(
476 'user','AAAAAAAB','2026-01-01T00:00:00Z'
477 );",
478 )
479 .unwrap();
480 database.execute_batch(IDENTITY_MIGRATION).unwrap();
481 let directory = Directory {
482 database: Mutex::new(database),
483 };
484 directory.seed_bootstrap_user("@taek42").unwrap();
485 directory
486 }
487
488 #[test]
489 fn opens_against_the_identity_schema_created_by_kmap_startup() {
490 let directory = std::env::temp_dir().join(format!(
491 "kennedy-telegram-identity-startup-test-{}",
492 uuid::Uuid::new_v4()
493 ));
494 std::fs::create_dir_all(&directory).unwrap();
495 let user_database = directory.join("users.sqlite3");
496 let kmap_schema = Connection::open(&user_database).unwrap();
497 kmap_schema
498 .execute_batch(
499 "CREATE TABLE kmap_system_roots(
500 role TEXT PRIMARY KEY CHECK(role IN ('user','kennedy')),
501 root_node_id TEXT NOT NULL UNIQUE CHECK(length(root_node_id)=8),
502 created_at TEXT NOT NULL
503 );
504 INSERT INTO kmap_system_roots VALUES(
505 'user','AAAAAAAB','2026-01-01T00:00:00Z'
506 );",
507 )
508 .unwrap();
509 drop(kmap_schema);
510
511 let identity = Directory::open(&user_database, "@taek42").unwrap();
512 assert!(
513 identity
514 .lock()
515 .unwrap()
516 .query_row(
517 "SELECT EXISTS(SELECT 1 FROM whitelist_entries WHERE handle='taek42')",
518 [],
519 |row| row.get::<_, i64>(0),
520 )
521 .unwrap()
522 != 0
523 );
524 drop(identity);
525 std::fs::remove_dir_all(directory).unwrap();
526 }
527
528 #[test]
529 fn tofu_is_owned_by_kennedy_and_numeric_ids_remain_authoritative() {
530 let directory = directory();
531 directory
532 .observe_identity(&IdentityObservation {
533 telegram_user_id: 42,
534 username: Some("TaEk42".into()),
535 display_name: "David".into(),
536 })
537 .unwrap();
538 assert!(
539 directory
540 .whitelist()
541 .unwrap()
542 .telegram_user_ids
543 .contains(&42)
544 );
545 directory
546 .observe_identity(&IdentityObservation {
547 telegram_user_id: 43,
548 username: Some("taek42".into()),
549 display_name: "Other".into(),
550 })
551 .unwrap();
552 assert!(
553 !directory
554 .whitelist()
555 .unwrap()
556 .telegram_user_ids
557 .contains(&43)
558 );
559 }
560
561 #[test]
562 fn identity_migration_removes_legacy_anonymous_group_pseudo_user() {
563 let directory = directory();
564 let database = directory.lock().unwrap();
565 database
566 .execute(
567 "INSERT INTO observed_identities(
568 telegram_user_id,current_username,display_name,first_seen_at,last_seen_at
569 ) VALUES(1087968824,'GroupAnonymousBot','Group',?1,?1)",
570 [Utc::now().to_rfc3339()],
571 )
572 .unwrap();
573 database.execute_batch(IDENTITY_MIGRATION).unwrap();
574 assert_eq!(
575 database
576 .query_row(
577 "SELECT COUNT(*) FROM observed_identities WHERE telegram_user_id=1087968824",
578 [],
579 |row| row.get::<_, i64>(0),
580 )
581 .unwrap(),
582 0
583 );
584 }
585
586 #[test]
587 fn add_user_capability_and_group_roots_stay_in_kennedy() {
588 let directory = directory();
589 directory
590 .observe_identity(&IdentityObservation {
591 telegram_user_id: 42,
592 username: Some("taek42".into()),
593 display_name: "David".into(),
594 })
595 .unwrap();
596 assert!(matches!(
597 directory.request_add_user(77, "@friend").unwrap(),
598 AddUserOutcome::Forbidden
599 ));
600 assert!(matches!(
601 directory.request_add_user(42, "@friend").unwrap(),
602 AddUserOutcome::Whitelisted { .. }
603 ));
604 directory.observe_group("opaque-group").unwrap();
605 let database = directory.lock().unwrap();
606 let group = directory_group_by_id(&database, "opaque-group")
607 .unwrap()
608 .unwrap();
609 assert_eq!(group.root_node_id, None);
610 assert!(!group.root_ready);
611 }
612
613 #[test]
614 fn root_completion_is_owned_and_conflict_checked_by_kennedy() {
615 let directory = directory();
616 directory
617 .observe_identity(&IdentityObservation {
618 telegram_user_id: 42,
619 username: Some("taek42".into()),
620 display_name: "David".into(),
621 })
622 .unwrap();
623 directory.observe_group("opaque-group").unwrap();
624
625 let user_root = "AAAAAAAC".parse().unwrap();
626 let user = directory.complete_user_root(42, user_root).unwrap();
627 assert!(user.root_ready);
628 assert_eq!(user.root_node_id.as_deref(), Some("AAAAAAAC"));
629
630 let group_root = "AAAAAAAD".parse().unwrap();
631 let group = directory
632 .complete_group_root("opaque-group", group_root)
633 .unwrap();
634 assert!(group.root_ready);
635 assert_eq!(group.root_node_id.as_deref(), Some("AAAAAAAD"));
636
637 let mismatch = directory
638 .complete_group_root("opaque-group", "AAAAAAAE".parse().unwrap())
639 .unwrap_err();
640 assert_eq!(mismatch.kind(), ErrorKind::Conflict);
641 }
642}