-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdb.py
More file actions
2333 lines (2111 loc) · 96.9 KB
/
Copy pathdb.py
File metadata and controls
2333 lines (2111 loc) · 96.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
# db.py
"""
Kontext Database -- SQLite backend for memory storage.
Single source of truth for all memory entries, relations, conflicts, and session state.
Flat markdown files are generated FROM this database, not the other way around.
"""
import json
import os
import re
import hashlib
import sqlite3
import uuid
from datetime import datetime, timezone
from collections import deque
from pathlib import Path
import struct
import shutil
# Minimal English stopword list for FTS5 query tokenization. Kept tight —
# over-filtering hurts short queries like "who is X". Only remove words that
# appear in almost every natural-language question and carry no signal.
_FTS_STOPWORDS = frozenset({
"a", "an", "the", "is", "are", "was", "were", "be", "been", "being",
"do", "does", "did", "have", "has", "had",
"i", "me", "my", "mine", "you", "your", "we", "our", "it", "its",
"and", "or", "but", "if", "so",
"of", "to", "in", "on", "for", "at", "by", "with", "from", "as",
"what", "who", "when", "where", "why", "how", "which",
"this", "that", "these", "those",
"s", "t", "m", "re", "ve", "ll", "d", # leftover from contraction splits
})
def _tokenize_fts_query(query: str) -> list[str]:
"""Extract FTS5-usable tokens from a natural-language query.
Splits on non-word chars, drops tokens <2 chars, strips stopwords.
Preserves the original if no usable tokens remain so exact-phrase
matches (e.g. "GTX 1080") still work via the phrase fallback.
"""
raw = re.findall(r"[A-Za-z0-9_]+", query)
return [t for t in raw if len(t) >= 2 and t.lower() not in _FTS_STOPWORDS]
def _default_db_path() -> str:
"""Return the runtime DB path outside the source tree unless explicitly overridden."""
override = os.environ.get("KONTEXT_DB_PATH")
if override:
return str(Path(override).expanduser())
data_root = os.environ.get("KONTEXT_DATA_DIR")
if data_root:
base = Path(data_root).expanduser()
elif os.name == "nt" and os.environ.get("APPDATA"):
base = Path(os.environ["APPDATA"]) / "Kontext"
else:
base = Path(os.environ.get("XDG_DATA_HOME", Path.home() / ".local" / "share")) / "kontext"
base.mkdir(parents=True, exist_ok=True)
target = base / "kontext.db"
legacy = Path(__file__).parent / "kontext.db"
if not target.exists() and legacy.exists():
try:
source = sqlite3.connect(str(legacy))
try:
dest = sqlite3.connect(str(target))
try:
source.backup(dest)
finally:
dest.close()
finally:
source.close()
except sqlite3.Error:
try:
shutil.copy2(legacy, target)
except OSError:
pass
return str(target)
def _normalize_workspace(workspace: str = "") -> str:
"""Canonicalize workspace identifiers so slash style changes still match."""
if not workspace:
return ""
try:
raw = str(Path(workspace).expanduser())
except Exception:
raw = str(workspace)
return os.path.normcase(os.path.normpath(raw)).replace("\\", "/")
def _workspace_session_path(workspace: str = "") -> Path:
"""Return the fallback session file path for a specific workspace."""
if not workspace:
return Path.home() / ".claude" / "_last_session.md"
try:
base = Path(workspace).expanduser().name or "workspace"
except Exception:
base = "workspace"
safe_base = "".join(ch if ch.isalnum() or ch in ("-", "_") else "-" for ch in base).strip("-")
if not safe_base:
safe_base = "workspace"
digest = hashlib.sha1(workspace.encode("utf-8")).hexdigest()[:12]
return Path.home() / ".claude" / "sessions" / f"{safe_base}-{digest}.md"
def _cloud_state_path(db_path: str) -> Path:
"""Return the adjacent cloud-link state path for a DB file."""
return Path(db_path).with_suffix(".cloud.json")
def _utc_timestamp() -> str:
"""Return a stable UTC timestamp for sync envelopes."""
return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
def _escape_like(value: str) -> str:
"""Escape LIKE wildcards so user searches stay literal."""
return value.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_")
# --- Schema migrations ---
#
# Each migration is (version, callable(conn)). Migrations are run in order,
# only if their version > current schema_version. Bodies must be idempotent
# (safe to re-run on a fresh DB) — the version gate is the primary guard but
# IF NOT EXISTS / IF EXISTS clauses are belt-and-suspenders.
def _migration_1_session_summary_cols(conn):
"""Add summary + files_touched columns to sessions (older schemas lacked them)."""
cols = {row[1] for row in conn.execute("PRAGMA table_info(sessions)").fetchall()}
for col in ("summary", "files_touched"):
if col not in cols:
conn.execute(f"ALTER TABLE sessions ADD COLUMN {col} TEXT DEFAULT ''")
def _migration_2_dedup_and_unique_indexes(conn):
"""Dedup duplicate rows then install UNIQUE indexes for race-safe inserts."""
conn.executescript("""
DELETE FROM entries WHERE id NOT IN (
SELECT MIN(id) FROM entries GROUP BY file, fact
);
DELETE FROM relations WHERE id NOT IN (
SELECT MIN(id) FROM relations GROUP BY entity_a, relation, entity_b
);
DELETE FROM conflicts WHERE id NOT IN (
SELECT MIN(id) FROM conflicts GROUP BY file, entry_a, entry_b
);
CREATE UNIQUE INDEX IF NOT EXISTS idx_entries_unique
ON entries(file, fact);
CREATE UNIQUE INDEX IF NOT EXISTS idx_relations_unique
ON relations(entity_a, relation, entity_b);
CREATE UNIQUE INDEX IF NOT EXISTS idx_conflicts_unique
ON conflicts(file, entry_a, entry_b);
""")
def _has_fts5(conn) -> bool:
"""Detect FTS5 support in the running SQLite build."""
try:
conn.execute("CREATE VIRTUAL TABLE IF NOT EXISTS _fts5_probe USING fts5(x, tokenize='trigram')")
conn.execute("DROP TABLE IF EXISTS _fts5_probe")
return True
except Exception:
return False
def _migration_3_fts5_entries(conn):
"""Build a trigram-tokenized FTS5 index over entries.fact + sync triggers.
Skipped silently if the SQLite build lacks FTS5 — search_entries falls back
to LIKE in that case.
"""
if not _has_fts5(conn):
return
conn.executescript("""
CREATE VIRTUAL TABLE IF NOT EXISTS entries_fts USING fts5(
fact,
content='entries',
content_rowid='id',
tokenize='trigram'
);
CREATE TRIGGER IF NOT EXISTS entries_ai AFTER INSERT ON entries BEGIN
INSERT INTO entries_fts(rowid, fact) VALUES (new.id, new.fact);
END;
CREATE TRIGGER IF NOT EXISTS entries_ad AFTER DELETE ON entries BEGIN
INSERT INTO entries_fts(entries_fts, rowid, fact) VALUES('delete', old.id, old.fact);
END;
CREATE TRIGGER IF NOT EXISTS entries_au AFTER UPDATE ON entries BEGIN
INSERT INTO entries_fts(entries_fts, rowid, fact) VALUES('delete', old.id, old.fact);
INSERT INTO entries_fts(rowid, fact) VALUES (new.id, new.fact);
END;
""")
# Backfill from existing rows (the contentless 'rebuild' command handles this)
conn.execute("INSERT INTO entries_fts(entries_fts) VALUES('rebuild')")
def _migration_4_access_count(conn):
"""Add access_count to entries for usage-driven ranking."""
cols = {row[1] for row in conn.execute("PRAGMA table_info(entries)").fetchall()}
if "access_count" not in cols:
conn.execute("ALTER TABLE entries ADD COLUMN access_count INTEGER DEFAULT 0")
def _migration_5_tool_events(conn):
"""Capture PostToolUse events for session intelligence."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS tool_events (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT DEFAULT '',
tool_name TEXT NOT NULL,
summary TEXT NOT NULL,
file_path TEXT DEFAULT NULL,
grade REAL DEFAULT 5.0,
promoted INTEGER DEFAULT 0,
created_at TEXT DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_tool_events_session ON tool_events(session_id);
CREATE INDEX IF NOT EXISTS idx_tool_events_created ON tool_events(created_at);
""")
def _migration_6_user_prompts(conn):
"""User prompt history with FTS5 for searchable session context."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS user_prompts (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT DEFAULT '',
content TEXT NOT NULL,
created_at TEXT DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_user_prompts_session ON user_prompts(session_id);
CREATE INDEX IF NOT EXISTS idx_user_prompts_created ON user_prompts(created_at);
""")
if not _has_fts5(conn):
return
conn.executescript("""
CREATE VIRTUAL TABLE IF NOT EXISTS user_prompts_fts USING fts5(
content,
content='user_prompts',
content_rowid='id',
tokenize='trigram'
);
CREATE TRIGGER IF NOT EXISTS prompts_ai AFTER INSERT ON user_prompts BEGIN
INSERT INTO user_prompts_fts(rowid, content) VALUES (new.id, new.content);
END;
CREATE TRIGGER IF NOT EXISTS prompts_ad AFTER DELETE ON user_prompts BEGIN
INSERT INTO user_prompts_fts(user_prompts_fts, rowid, content)
VALUES('delete', old.id, old.content);
END;
CREATE TRIGGER IF NOT EXISTS prompts_au AFTER UPDATE ON user_prompts BEGIN
INSERT INTO user_prompts_fts(user_prompts_fts, rowid, content)
VALUES('delete', old.id, old.content);
INSERT INTO user_prompts_fts(rowid, content) VALUES(new.id, new.content);
END;
""")
conn.execute("INSERT INTO user_prompts_fts(user_prompts_fts) VALUES('rebuild')")
def _migration_7_session_intelligence(conn):
"""Add investigated + learned columns to sessions for richer auto-summaries."""
cols = {row[1] for row in conn.execute("PRAGMA table_info(sessions)").fetchall()}
for col in ("investigated", "learned"):
if col not in cols:
conn.execute(f"ALTER TABLE sessions ADD COLUMN {col} TEXT DEFAULT ''")
def _migration_8_hook_session_id(conn):
"""Associate hook-generated session summaries with Claude hook session IDs."""
cols = {row[1] for row in conn.execute("PRAGMA table_info(sessions)").fetchall()}
if "hook_session_id" not in cols:
conn.execute("ALTER TABLE sessions ADD COLUMN hook_session_id TEXT DEFAULT ''")
conn.execute("""
CREATE UNIQUE INDEX IF NOT EXISTS idx_sessions_hook_session_id
ON sessions(hook_session_id)
WHERE hook_session_id != ''
""")
def _migration_9_session_workspace(conn):
"""Store workspace identity so task restore is scoped instead of global."""
cols = {row[1] for row in conn.execute("PRAGMA table_info(sessions)").fetchall()}
if "workspace" not in cols:
conn.execute("ALTER TABLE sessions ADD COLUMN workspace TEXT DEFAULT ''")
def _migration_10_cloud_identity(conn):
"""Add workspace, device, and manifest tables for cloud sync."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS workspaces (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
recovery_key_id TEXT NOT NULL,
created_at TEXT DEFAULT (datetime('now'))
);
CREATE TABLE IF NOT EXISTS devices (
id TEXT PRIMARY KEY,
workspace_id TEXT NOT NULL,
label TEXT NOT NULL,
device_class TEXT NOT NULL CHECK(device_class IN ('interactive', 'server')),
public_key BLOB NOT NULL,
enrolled_at TEXT DEFAULT (datetime('now')),
revoked_at TEXT DEFAULT NULL,
FOREIGN KEY(workspace_id) REFERENCES workspaces(id)
);
CREATE TABLE IF NOT EXISTS sync_manifests (
workspace_id TEXT PRIMARY KEY,
schema_version INTEGER NOT NULL,
embedding_model TEXT NOT NULL,
ranking_version TEXT NOT NULL,
prompt_routing_version TEXT NOT NULL,
FOREIGN KEY(workspace_id) REFERENCES workspaces(id)
);
""")
def _migration_11_cloud_ops(conn):
"""Add append-only history ops and per-device sync cursors."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS history_ops (
id TEXT PRIMARY KEY,
workspace_id TEXT NOT NULL,
device_id TEXT NOT NULL,
op_kind TEXT NOT NULL,
entity_type TEXT NOT NULL,
entity_id TEXT NOT NULL,
payload BLOB NOT NULL,
created_at TEXT NOT NULL,
applied_at TEXT DEFAULT NULL,
FOREIGN KEY(workspace_id) REFERENCES workspaces(id),
FOREIGN KEY(device_id) REFERENCES devices(id)
);
CREATE INDEX IF NOT EXISTS idx_history_ops_workspace_created
ON history_ops(workspace_id, created_at, id);
CREATE TABLE IF NOT EXISTS sync_cursors (
workspace_id TEXT NOT NULL,
device_id TEXT NOT NULL,
lane TEXT NOT NULL CHECK(lane IN ('history', 'canonical', 'history_push', 'canonical_push')),
cursor TEXT NOT NULL,
PRIMARY KEY (workspace_id, device_id, lane),
FOREIGN KEY(workspace_id) REFERENCES workspaces(id),
FOREIGN KEY(device_id) REFERENCES devices(id)
);
""")
def _migration_13_workspace_auth(conn):
"""Add bearer-token auth columns to workspaces and harden revocation semantics."""
cols = {row[1] for row in conn.execute("PRAGMA table_info(workspaces)").fetchall()}
if "api_token_hash" not in cols:
conn.execute("ALTER TABLE workspaces ADD COLUMN api_token_hash TEXT DEFAULT NULL")
if "api_token_salt" not in cols:
conn.execute("ALTER TABLE workspaces ADD COLUMN api_token_salt TEXT DEFAULT NULL")
def _migration_14_push_cursors(conn):
"""Widen sync_cursors.lane CHECK to allow history_push / canonical_push."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS sync_cursors__new (
workspace_id TEXT NOT NULL,
device_id TEXT NOT NULL,
lane TEXT NOT NULL CHECK(lane IN ('history', 'canonical', 'history_push', 'canonical_push')),
cursor TEXT NOT NULL,
PRIMARY KEY (workspace_id, device_id, lane),
FOREIGN KEY(workspace_id) REFERENCES workspaces(id),
FOREIGN KEY(device_id) REFERENCES devices(id)
);
INSERT INTO sync_cursors__new (workspace_id, device_id, lane, cursor)
SELECT workspace_id, device_id, lane, cursor FROM sync_cursors;
DROP TABLE sync_cursors;
ALTER TABLE sync_cursors__new RENAME TO sync_cursors;
""")
def _migration_17_sessions_files_loaded(conn):
"""Track memory files pulled during a session via kontext_query/search."""
cols = {row[1] for row in conn.execute("PRAGMA table_info(sessions)").fetchall()}
if "files_loaded" not in cols:
conn.execute("ALTER TABLE sessions ADD COLUMN files_loaded TEXT DEFAULT ''")
def _migration_18_retrieval_evals(conn):
"""Store per-query retrieval eval results for regression tracking."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS retrieval_evals (
id INTEGER PRIMARY KEY AUTOINCREMENT,
run_id TEXT NOT NULL,
query_text TEXT NOT NULL,
category TEXT NOT NULL DEFAULT '',
held_out INTEGER NOT NULL DEFAULT 0,
recall_at_3 REAL NOT NULL DEFAULT 0.0,
recall_at_5 REAL NOT NULL DEFAULT 0.0,
mrr REAL NOT NULL DEFAULT 0.0,
latency_ms INTEGER NOT NULL DEFAULT 0,
mode TEXT NOT NULL DEFAULT 'rrf',
rank TEXT NOT NULL DEFAULT 'first_seen',
top_files TEXT DEFAULT '',
git_sha TEXT DEFAULT '',
schema_version INTEGER DEFAULT 0,
ts TEXT DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ','now'))
);
CREATE INDEX IF NOT EXISTS idx_re_run_id ON retrieval_evals(run_id);
CREATE INDEX IF NOT EXISTS idx_re_ts ON retrieval_evals(ts);
""")
def _migration_15_retrieval_queries(conn):
"""Log every kontext_query / kontext_search call for regression tracking."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS retrieval_queries (
id INTEGER PRIMARY KEY AUTOINCREMENT,
query_text TEXT NOT NULL,
tool_name TEXT NOT NULL,
semantic_flag INTEGER NOT NULL DEFAULT 0,
expanded_query TEXT DEFAULT '',
results_json TEXT DEFAULT '',
result_count INTEGER NOT NULL DEFAULT 0,
latency_ms INTEGER NOT NULL DEFAULT 0,
session_id TEXT DEFAULT '',
ts TEXT DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ','now'))
);
CREATE INDEX IF NOT EXISTS idx_rq_ts ON retrieval_queries(ts);
CREATE INDEX IF NOT EXISTS idx_rq_session ON retrieval_queries(session_id);
CREATE INDEX IF NOT EXISTS idx_rq_tool ON retrieval_queries(tool_name);
""")
def _migration_12_canonical_objects(conn):
"""Add canonical object and revision tables for curated memory sync."""
conn.executescript("""
CREATE TABLE IF NOT EXISTS canonical_objects (
id TEXT PRIMARY KEY,
workspace_id TEXT NOT NULL,
object_type TEXT NOT NULL,
head_revision TEXT NOT NULL DEFAULT '',
tombstoned INTEGER DEFAULT 0,
FOREIGN KEY(workspace_id) REFERENCES workspaces(id)
);
CREATE TABLE IF NOT EXISTS canonical_revisions (
id TEXT PRIMARY KEY,
object_id TEXT NOT NULL,
parent_revision TEXT DEFAULT NULL,
device_id TEXT NOT NULL,
payload BLOB NOT NULL,
created_at TEXT NOT NULL,
accepted INTEGER NOT NULL DEFAULT 0,
FOREIGN KEY(object_id) REFERENCES canonical_objects(id),
FOREIGN KEY(device_id) REFERENCES devices(id)
);
CREATE INDEX IF NOT EXISTS idx_canonical_revisions_object_created
ON canonical_revisions(object_id, created_at, id);
CREATE TABLE IF NOT EXISTS snapshots (
id TEXT PRIMARY KEY,
workspace_id TEXT NOT NULL,
manifest_hash TEXT NOT NULL,
blob_path TEXT NOT NULL,
created_at TEXT DEFAULT (datetime('now')),
FOREIGN KEY(workspace_id) REFERENCES workspaces(id)
);
""")
def _migration_19_memory_type(conn):
"""Partition entries into episodic / semantic / procedural.
Enables partition-aware retrieval: callback queries filter by
memory_type='episodic', stable-attribute queries filter by
memory_type='semantic', habit queries by 'procedural'. Schema is
permissive (DEFAULT 'semantic', no CHECK constraint to keep ALTER
cheap on large tables); app code enforces the enum at write time.
Backfill heuristic — zero LLM, single SQL pass:
- `episodic` = entries whose source contains a dated AI tag
([Claude 2026-04], [ChatGPT 2025-08], etc.) or whose fact
starts with a bracketed date.
- `procedural` = entries describing rules/preferences/workflows
("always X", "never X", "prefer X", "How to X", "Protocol:").
- `semantic` = everything else (DEFAULT, no UPDATE needed).
"""
cols = {row[1] for row in conn.execute("PRAGMA table_info(entries)").fetchall()}
if "memory_type" in cols:
return
conn.execute("ALTER TABLE entries ADD COLUMN memory_type TEXT DEFAULT 'semantic'")
conn.execute("CREATE INDEX IF NOT EXISTS idx_entries_memory_type ON entries(memory_type)")
conn.execute("""
UPDATE entries SET memory_type = 'episodic'
WHERE memory_type = 'semantic'
AND (source LIKE '%Claude 20%'
OR source LIKE '%Codex 20%'
OR source LIKE '%Mastermind 20%'
OR source LIKE '%ChatGPT 20%'
OR source LIKE '%Gemini 20%'
OR source LIKE '%WhatsApp 20%'
OR fact LIKE '[20__-__%')
""")
conn.execute("""
UPDATE entries SET memory_type = 'procedural'
WHERE memory_type = 'semantic'
AND (fact LIKE 'always %' OR fact LIKE 'Always %'
OR fact LIKE 'never %' OR fact LIKE 'Never %'
OR fact LIKE 'prefer%' OR fact LIKE 'Prefer%'
OR fact LIKE 'When %' OR fact LIKE 'when %'
OR fact LIKE 'How to %' OR fact LIKE 'how to %'
OR fact LIKE 'Protocol:%' OR fact LIKE 'Rule:%'
OR fact LIKE '% workflow%' OR fact LIKE '% rule:%')
""")
def _migration_21_score_history(conn):
"""Daily Kontext Score snapshots — one row per calendar day.
Replaces the dashboard's 14-day history chart which was projecting
today's values backward via Math.sin (synthetic). Populated by
dream.phase_history_snapshot on each dream cycle; dashboard reads
the trailing 14 rows. Date is PRIMARY KEY so repeated same-day
snapshots upsert rather than duplicating.
"""
conn.executescript("""
CREATE TABLE IF NOT EXISTS score_history (
snapshot_date TEXT PRIMARY KEY,
score INTEGER NOT NULL DEFAULT 0,
breadth INTEGER NOT NULL DEFAULT 0,
depth INTEGER NOT NULL DEFAULT 0,
recency INTEGER NOT NULL DEFAULT 0,
longevity INTEGER NOT NULL DEFAULT 0,
linkage INTEGER NOT NULL DEFAULT 0,
captures INTEGER NOT NULL DEFAULT 0,
prompts INTEGER NOT NULL DEFAULT 0,
entries_active INTEGER NOT NULL DEFAULT 0,
recorded_at TEXT DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ','now'))
);
CREATE INDEX IF NOT EXISTS idx_score_history_date ON score_history(snapshot_date);
""")
def _migration_20_relation_typed_edges(conn):
"""Typed edges on `relations`: supersedes / contradicts / caused_by /
instance_of / temporal_before / co_occurs.
Before: free-text `relation` column (usually 'co_occurs_with'). After:
add `rel_type` as a closed vocabulary enforced by app code. Dream's
resolve phase can write supersedes edges directly rather than only
tombstoning losers — giving the graph a notion of "this fact replaced
that one" rather than just "these co-occur."
Backfill: map existing `relation` strings to the enum via LIKE.
Anything ambiguous stays 'co_occurs'.
"""
cols = {row[1] for row in conn.execute("PRAGMA table_info(relations)").fetchall()}
if "rel_type" in cols:
return
conn.execute("ALTER TABLE relations ADD COLUMN rel_type TEXT DEFAULT 'co_occurs'")
conn.execute("CREATE INDEX IF NOT EXISTS idx_relations_rel_type ON relations(rel_type)")
conn.execute("""
UPDATE relations SET rel_type = CASE
WHEN relation LIKE '%supersede%' THEN 'supersedes'
WHEN relation LIKE '%contradict%' THEN 'contradicts'
WHEN relation LIKE '%caused%' OR relation LIKE '%causes%' THEN 'caused_by'
WHEN relation LIKE '%instance%' OR relation LIKE '%is_a%' THEN 'instance_of'
WHEN relation LIKE '%temporal%' OR relation LIKE '%before%' THEN 'temporal_before'
ELSE 'co_occurs'
END
WHERE rel_type = 'co_occurs'
""")
MIGRATIONS = [
(1, _migration_1_session_summary_cols),
(2, _migration_2_dedup_and_unique_indexes),
(3, _migration_3_fts5_entries),
(4, _migration_4_access_count),
(5, _migration_5_tool_events),
(6, _migration_6_user_prompts),
(7, _migration_7_session_intelligence),
(8, _migration_8_hook_session_id),
(9, _migration_9_session_workspace),
(10, _migration_10_cloud_identity),
(11, _migration_11_cloud_ops),
(12, _migration_12_canonical_objects),
(13, _migration_13_workspace_auth),
(14, _migration_14_push_cursors),
(15, _migration_15_retrieval_queries),
(17, _migration_17_sessions_files_loaded),
(18, _migration_18_retrieval_evals),
(19, _migration_19_memory_type),
(20, _migration_20_relation_typed_edges),
(21, _migration_21_score_history),
]
LATEST_SCHEMA_VERSION = max(v for v, _ in MIGRATIONS)
class _Transaction:
"""Context manager for SQLite transactions with rollback on failure.
Also sets db._in_batch so nested _execute() calls skip their auto-commit —
this lets phases/methods that call _execute() compose into one atomic unit.
"""
def __init__(self, db):
self.db = db
self.conn = db.conn
self._prev_batch = False
def __enter__(self):
self._prev_batch = self.db._in_batch
self.db._in_batch = True
# Only open a real transaction at the outermost level.
if not self._prev_batch:
self.conn.execute("BEGIN IMMEDIATE")
return self.conn
def __exit__(self, exc_type, exc_val, exc_tb):
self.db._in_batch = self._prev_batch
if self._prev_batch:
# Nested — let the outer scope decide commit/rollback.
return False
if exc_type is None:
self.conn.commit()
else:
self.conn.rollback()
return False # re-raise exceptions
class KontextDB:
"""SQLite-backed memory database."""
def __init__(self, db_path: str = None):
if db_path is None:
db_path = _default_db_path()
self.db_path = db_path
self.conn = sqlite3.connect(db_path)
self.conn.row_factory = sqlite3.Row
self.conn.execute("PRAGMA journal_mode=WAL")
self.conn.execute("PRAGMA synchronous=NORMAL") # 2-10x faster writes; WAL durability retained
self.conn.execute("PRAGMA cache_size=-16000") # 16 MB page cache
self.conn.execute("PRAGMA foreign_keys=ON")
self._in_batch = False # True inside a transaction() context — suppresses _execute auto-commits
self._embed_cache = None # Lazy-built dict: {id: (file, fact, source, grade, tier, vec_tuple)}
self._fts_enabled = False # set after _create_tables → _migrate runs
self._create_tables()
self._fts_enabled = self.conn.execute(
"SELECT 1 FROM sqlite_master WHERE type='table' AND name='entries_fts'"
).fetchone() is not None
def _create_tables(self):
self.conn.executescript("""
CREATE TABLE IF NOT EXISTS schema_version (
id INTEGER PRIMARY KEY CHECK (id = 1),
version INTEGER NOT NULL DEFAULT 0
);
INSERT OR IGNORE INTO schema_version (id, version) VALUES (1, 0);
CREATE TABLE IF NOT EXISTS entries (
id INTEGER PRIMARY KEY AUTOINCREMENT,
file TEXT NOT NULL,
fact TEXT NOT NULL,
source TEXT DEFAULT '',
grade REAL DEFAULT 5,
tier TEXT DEFAULT 'active' CHECK(tier IN ('active', 'historical', 'cold')),
created_at TEXT DEFAULT (datetime('now')),
updated_at TEXT DEFAULT (datetime('now')),
last_accessed TEXT DEFAULT (datetime('now')),
embedding BLOB DEFAULT NULL
);
CREATE TABLE IF NOT EXISTS relations (
id INTEGER PRIMARY KEY AUTOINCREMENT,
entity_a TEXT NOT NULL,
relation TEXT NOT NULL,
entity_b TEXT NOT NULL,
confidence REAL DEFAULT 1.0,
source TEXT DEFAULT '',
created_at TEXT DEFAULT (datetime('now'))
);
CREATE TABLE IF NOT EXISTS conflicts (
id INTEGER PRIMARY KEY AUTOINCREMENT,
file TEXT NOT NULL,
entry_a TEXT NOT NULL,
entry_b TEXT NOT NULL,
status TEXT DEFAULT 'pending' CHECK(status IN ('pending', 'resolved')),
resolution TEXT DEFAULT NULL,
created_at TEXT DEFAULT (datetime('now')),
resolved_at TEXT DEFAULT NULL
);
CREATE TABLE IF NOT EXISTS sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
hook_session_id TEXT DEFAULT '',
workspace TEXT DEFAULT '',
project TEXT DEFAULT '',
status TEXT DEFAULT '',
next_step TEXT DEFAULT '',
key_decisions TEXT DEFAULT '',
summary TEXT DEFAULT '',
files_touched TEXT DEFAULT '',
created_at TEXT DEFAULT (datetime('now'))
);
CREATE TABLE IF NOT EXISTS file_meta (
filename TEXT PRIMARY KEY,
file_type TEXT DEFAULT 'user',
description TEXT DEFAULT '',
updated_at TEXT DEFAULT (datetime('now'))
);
CREATE INDEX IF NOT EXISTS idx_entries_file ON entries(file);
CREATE INDEX IF NOT EXISTS idx_entries_tier ON entries(tier);
CREATE INDEX IF NOT EXISTS idx_entries_grade ON entries(grade);
CREATE INDEX IF NOT EXISTS idx_relations_entity_a ON relations(entity_a);
CREATE INDEX IF NOT EXISTS idx_relations_entity_b ON relations(entity_b);
""")
# Drop the useless leading-wildcard-LIKE index if it exists from older schemas.
self.conn.execute("DROP INDEX IF EXISTS idx_entries_fact")
self.conn.commit()
self._migrate()
def _migrate(self):
"""Run any pending migrations from MIGRATIONS, in order, idempotently.
Each migration is wrapped in a transaction and bumps schema_version on
success. Already-applied migrations are skipped via the version check.
"""
# Legacy DBs (pre-framework) won't have schema_version — create it.
self.conn.execute("""
CREATE TABLE IF NOT EXISTS schema_version (
id INTEGER PRIMARY KEY CHECK (id = 1),
version INTEGER NOT NULL DEFAULT 0
)
""")
self.conn.execute("INSERT OR IGNORE INTO schema_version (id, version) VALUES (1, 0)")
self.conn.commit()
current = self.conn.execute("SELECT version FROM schema_version WHERE id = 1").fetchone()[0]
for version, fn in MIGRATIONS:
if version <= current:
continue
with self.transaction() as conn:
fn(conn)
conn.execute("UPDATE schema_version SET version = ? WHERE id = 1", (version,))
current = version
def _execute(self, sql, params=()):
cursor = self.conn.execute(sql, params)
if not self._in_batch:
self.conn.commit()
return cursor
def transaction(self):
"""Context manager for explicit transactions. Use for multi-step operations.
Inside the block, any call to self._execute() defers its commit — the
whole block commits atomically on success and rolls back on exception.
Nested transactions compose (inner ones defer to the outermost).
Usage:
with db.transaction():
db.add_entry(...) # normally auto-commits; now deferred
db.add_entry(...)
# single commit on success
"""
return _Transaction(self)
def list_tables(self):
cursor = self.conn.execute("SELECT name FROM sqlite_master WHERE type='table'")
return [row[0] for row in cursor.fetchall()]
def _load_cloud_link_state(self) -> dict:
"""Load adjacent cloud-link state when this DB is enrolled in sync."""
path = _cloud_state_path(self.db_path)
if not path.exists():
return {}
try:
state = json.loads(path.read_text(encoding="utf-8"))
except (json.JSONDecodeError, OSError):
return {}
workspace_id = str(state.get("workspace_id", "")).strip()
device_id = str(state.get("device_id", "")).strip()
if not workspace_id or not device_id:
return {}
workspace = self.conn.execute(
"SELECT 1 FROM workspaces WHERE id = ?",
(workspace_id,),
).fetchone()
device = self.conn.execute(
"SELECT 1 FROM devices WHERE id = ? AND workspace_id = ?",
(device_id, workspace_id),
).fetchone()
if workspace is None or device is None:
return {}
return {
"workspace_id": workspace_id,
"device_id": device_id,
}
def _maybe_append_local_history_op(self, op_kind: str, entity_type: str,
entity_id: str, payload: dict) -> bool:
"""Append a linked local write to the history lane when sync is active."""
state = self._load_cloud_link_state()
if not state:
return False
from cloud.codec import pack_payload
return self.append_history_op(
op_id=f"op-{uuid.uuid4().hex}",
workspace_id=state["workspace_id"],
device_id=state["device_id"],
op_kind=op_kind,
entity_type=entity_type,
entity_id=str(entity_id),
payload=pack_payload(payload),
created_at=_utc_timestamp(),
)
# --- Entries ---
def add_entry(self, file: str, fact: str, source: str = "", grade: float = 5,
tier: str = "active", emit_cloud: bool = True,
created_at: str | None = None) -> int:
"""Add an entry. Race-safe via UNIQUE(file, fact) index + INSERT OR IGNORE."""
with self.transaction():
if created_at:
cursor = self._execute(
"INSERT OR IGNORE INTO entries (file, fact, source, grade, tier, created_at, updated_at) "
"VALUES (?, ?, ?, ?, ?, ?, ?)",
(file, fact, source, grade, tier, created_at, created_at)
)
else:
cursor = self._execute(
"INSERT OR IGNORE INTO entries (file, fact, source, grade, tier) VALUES (?, ?, ?, ?, ?)",
(file, fact, source, grade, tier)
)
row = self.conn.execute(
"SELECT id FROM entries WHERE file = ? AND fact = ?", (file, fact)
).fetchone()
entry_id = row[0] if row else 0
if cursor.rowcount > 0 and emit_cloud and entry_id:
self._maybe_append_local_history_op(
op_kind="entry.written",
entity_type="entry",
entity_id=str(entry_id),
payload={
"file": file,
"fact": fact,
"source": source,
"grade": grade,
"tier": tier,
},
)
return entry_id
def update_entry(self, entry_id: int, **kwargs):
"""Update specific fields of an entry."""
allowed = {"fact", "source", "grade", "tier", "file", "embedding"}
updates = {k: v for k, v in kwargs.items() if k in allowed}
if not updates:
return
updates["updated_at"] = datetime.now(timezone.utc).isoformat()
set_clause = ", ".join(f"{k} = ?" for k in updates)
values = list(updates.values()) + [entry_id]
self._execute(f"UPDATE entries SET {set_clause} WHERE id = ?", values)
if "fact" in updates or "file" in updates:
self._embed_cache = None # fact text changed → cached copy stale
def get_entry(self, entry_id: int) -> dict | None:
row = self.conn.execute("SELECT * FROM entries WHERE id = ?", (entry_id,)).fetchone()
if row:
# Write amplification fix: only bump last_accessed if stale >1 hour.
# last_accessed is only consumed by decay/purge, so sub-hour precision is worthless.
prev = row["last_accessed"]
stale = (
not prev
or self.conn.execute(
"SELECT datetime('now', '-1 hour') > ?", (prev,)
).fetchone()[0] == 1
)
if stale:
self._execute(
"UPDATE entries SET last_accessed = datetime('now') WHERE id = ?",
(entry_id,),
)
return dict(row)
return None
def get_entries(self, file: str = None, tier: str = None, min_grade: float = None) -> list[dict]:
sql = "SELECT * FROM entries WHERE 1=1"
params = []
if file:
sql += " AND file = ?"
params.append(file)
if tier:
sql += " AND tier = ?"
params.append(tier)
if min_grade is not None:
sql += " AND grade >= ?"
params.append(min_grade)
sql += " ORDER BY grade DESC, updated_at DESC"
return [dict(r) for r in self.conn.execute(sql, params).fetchall()]
def search_entries(self, query: str, limit: int = 20,
file: str = None, tier: str = None,
min_grade: float = None) -> list[dict]:
"""Substring search on fact content with optional filters.
Uses FTS5 trigram index when available (O(log n) on indexed prefixes
of 3+ chars), falling back to escaped LIKE on older SQLite builds.
"""
if self._fts_enabled and query and query.strip():
# Build a tokenized OR query first so natural-language questions
# like "What's my PFA tax regime?" match any entry containing
# "PFA", "tax", or "regime" (ranked by grade). Fall back to the
# full quoted phrase if tokenization produces nothing useful, so
# exact multi-word matches ("GTX 1080") still work.
tokens = _tokenize_fts_query(query)
candidate_match_clauses: list[str] = []
if tokens:
candidate_match_clauses.append(
" OR ".join('"' + t.replace('"', '""') + '"' for t in tokens)
)
# Always include the literal phrase as a second attempt.
candidate_match_clauses.append('"' + query.replace('"', '""') + '"')
base_sql = (
"SELECT e.* FROM entries e "
"JOIN entries_fts f ON e.id = f.rowid "
"WHERE f.fact MATCH ?"
)
suffix = ""
extra_params: list = []
if file:
suffix += " AND e.file = ?"
extra_params.append(file)
if tier:
suffix += " AND e.tier = ?"
extra_params.append(tier)
if min_grade is not None:
suffix += " AND e.grade >= ?"
extra_params.append(min_grade)
suffix += " ORDER BY e.grade DESC LIMIT ?"
for match_clause in candidate_match_clauses:
params: list = [match_clause] + extra_params + [limit]
try:
rows = self.conn.execute(base_sql + suffix, params).fetchall()
except sqlite3.OperationalError:
continue # Trigram <3 chars or FTS syntax edge → try next clause
if rows:
return [dict(r) for r in rows]
# All FTS attempts returned empty or errored → fall through to LIKE
safe_query = _escape_like(query)
sql = "SELECT * FROM entries WHERE fact LIKE ? ESCAPE '\\'"
params = [f"%{safe_query}%"]
if file:
sql += " AND file = ?"
params.append(file)
if tier:
sql += " AND tier = ?"
params.append(tier)
if min_grade is not None:
sql += " AND grade >= ?"
params.append(min_grade)
sql += " ORDER BY grade DESC LIMIT ?"
params.append(limit)
return [dict(r) for r in self.conn.execute(sql, params).fetchall()]
def delete_entry(self, entry_id: int):
self._execute("DELETE FROM entries WHERE id = ?", (entry_id,))
self._embed_cache = None # invalidate
def list_files(self) -> dict:
"""Return dict of {filename: entry_count}."""
rows = self.conn.execute(
"SELECT file, COUNT(*) as cnt FROM entries GROUP BY file ORDER BY cnt DESC"
).fetchall()
return {r[0]: r[1] for r in rows}
def get_recent_changes(self, hours: int = 24) -> list[dict]:
return [dict(r) for r in self.conn.execute(
"SELECT * FROM entries WHERE updated_at >= datetime('now', ? || ' hours') ORDER BY updated_at DESC",
(f"-{int(hours)}",)
).fetchall()]
def decay_scores(self, days_threshold: int = 60, decay_amount: float = 0.5):
"""Reduce grade of entries not accessed in days_threshold days. Minimum grade: 1."""
self._execute("""
UPDATE entries SET
grade = MAX(1, grade - ?),
tier = CASE
WHEN MAX(1, grade - ?) < 5 THEN 'cold'
WHEN MAX(1, grade - ?) < 8 THEN 'historical'
ELSE tier
END