Coverage for paperless_ai/tables.py: 0%

87 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 09:07 +0000

1"""Thin gateways over the plain relational side tables that sit alongside the 

2vec0 table. Each method takes the sqlite3.Connection to operate on 

3explicitly, rather than owning one -- the store swaps connections during 

4compact()/migration, and migrations always work across two connections 

5(src_conn, dst_conn) at once. 

6 

7PRECONDITION: Callers must set conn.row_factory = sqlite3.Row before passing a 

8connection to any of these gateways' read methods. The read methods across all 

9three classes (DocumentChunksTable.chunk_ids_for_document, IndexMetaTable._get, 

10DocumentMetaTable.all_modified_times, DocumentMetaTable.copy_all) use 

11row["column_name"] dictionary-style indexing, which requires sqlite3.Row as the 

12row factory -- without it, sqlite3.Row is not set, rows are returned as plain 

13tuples, and tuple indices must be integers, raising TypeError. 

14""" 

15 

16import sqlite3 

17from collections.abc import Iterable 

18from typing import NamedTuple 

19 

20 

21class ChunkRow(NamedTuple): 

22 chunk_id: str 

23 document_id: int 

24 

25 

26class DocumentMetaRow(NamedTuple): 

27 document_id: int 

28 modified: str 

29 

30 

31class DocumentChunksTable: 

32 """chunk_id -> document_id, indexed by document_id. Gives O(1) 

33 per-document chunk lookup that vec0's own document_id metadata column 

34 cannot (see PaperlessSqliteVecVectorStore._delete_chunks_by_document_id). 

35 """ 

36 

37 @staticmethod 

38 def create(conn: sqlite3.Connection) -> None: 

39 conn.execute( 

40 "CREATE TABLE IF NOT EXISTS document_chunks " 

41 "(chunk_id TEXT PRIMARY KEY, document_id INTEGER NOT NULL)", 

42 ) 

43 conn.execute( 

44 "CREATE INDEX IF NOT EXISTS idx_document_chunks_document_id " 

45 "ON document_chunks (document_id)", 

46 ) 

47 

48 @staticmethod 

49 def insert_many(conn: sqlite3.Connection, rows: Iterable[ChunkRow]) -> None: 

50 """rows must already be batch-bounded by the caller (e.g. vec0's own 

51 fetchmany() loop) -- this never reads, so it can't itself introduce 

52 an unbounded scan, but a whole-table iterable defeats the point.""" 

53 conn.executemany( 

54 "INSERT INTO document_chunks (chunk_id, document_id) VALUES (?, ?)", 

55 rows, 

56 ) 

57 

58 @staticmethod 

59 def chunk_ids_for_document( 

60 conn: sqlite3.Connection, 

61 document_id: int, 

62 ) -> list[str]: 

63 return [ 

64 row["chunk_id"] 

65 for row in conn.execute( 

66 "SELECT chunk_id FROM document_chunks WHERE document_id = ?", 

67 (document_id,), 

68 ).fetchall() 

69 ] 

70 

71 @staticmethod 

72 def delete_for_document(conn: sqlite3.Connection, document_id: int) -> None: 

73 conn.execute( 

74 "DELETE FROM document_chunks WHERE document_id = ?", 

75 (document_id,), 

76 ) 

77 

78 @staticmethod 

79 def delete_all(conn: sqlite3.Connection) -> None: 

80 conn.execute("DELETE FROM document_chunks") 

81 

82 @staticmethod 

83 def count(conn: sqlite3.Connection) -> int: 

84 """Cheap stand-in for vec0's own row count -- see compact().""" 

85 return conn.execute("SELECT count(*) FROM document_chunks").fetchone()[0] 

86 

87 

88class DocumentMetaTable: 

89 """document_id -> modified, one row per document. Lives outside vec0 

90 because vec0 only inlines TEXT metadata up to 12 bytes and `modified` 

91 (an ISO timestamp) is always longer. 

92 """ 

93 

94 @staticmethod 

95 def create(conn: sqlite3.Connection) -> None: 

96 conn.execute( 

97 "CREATE TABLE IF NOT EXISTS document_meta " 

98 "(document_id INTEGER PRIMARY KEY, modified TEXT NOT NULL)", 

99 ) 

100 

101 @staticmethod 

102 def upsert_many( 

103 conn: sqlite3.Connection, 

104 rows: Iterable[DocumentMetaRow], 

105 ) -> None: 

106 conn.executemany( 

107 "INSERT INTO document_meta (document_id, modified) VALUES (?, ?) " 

108 "ON CONFLICT(document_id) DO UPDATE SET modified = excluded.modified", 

109 rows, 

110 ) 

111 

112 @staticmethod 

113 def delete_for_document(conn: sqlite3.Connection, document_id: int) -> None: 

114 conn.execute( 

115 "DELETE FROM document_meta WHERE document_id = ?", 

116 (document_id,), 

117 ) 

118 

119 @staticmethod 

120 def delete_all(conn: sqlite3.Connection) -> None: 

121 conn.execute("DELETE FROM document_meta") 

122 

123 @staticmethod 

124 def copy_all( 

125 src_conn: sqlite3.Connection, 

126 dst_conn: sqlite3.Connection, 

127 batch_size: int, 

128 ) -> None: 

129 """Stream document_meta from src_conn into dst_conn in bounded 

130 batches. The *only* sanctioned way to move this table across 

131 connections (compact()/migrations) -- an unbounded fetchall here 

132 would defeat the same OOM-avoidance the vec0 row copy already relies 

133 on. batch_size has no default: forces the call site to think about 

134 it (pass BATCH_SIZE).""" 

135 cursor = src_conn.execute( 

136 "SELECT document_id, modified FROM document_meta", 

137 ) 

138 while batch := cursor.fetchmany(batch_size): 

139 DocumentMetaTable.upsert_many( 

140 dst_conn, 

141 (DocumentMetaRow(r["document_id"], r["modified"]) for r in batch), 

142 ) 

143 

144 @staticmethod 

145 def all_modified_times(conn: sqlite3.Connection) -> dict[str, str]: 

146 """Full document_id -> modified map, for get_modified_times()'s 

147 public API only. One unbounded read by design (existing behavior). 

148 Never use this for cross-connection copying; see copy_all().""" 

149 return { 

150 str(row["document_id"]): str(row["modified"] or "") 

151 for row in conn.execute( 

152 "SELECT document_id, modified FROM document_meta", 

153 ) 

154 } 

155 

156 

157class IndexMetaTable: 

158 """Typed accessors over index_meta's key/value rows -- replaces 

159 PaperlessSqliteVecVectorStore._meta_get_on/_meta_set_on, which returned 

160 untyped str | None regardless of whether the key held an int (dim, 

161 schema_version, total_inserts) or a string (embed_model). 

162 """ 

163 

164 @staticmethod 

165 def create(conn: sqlite3.Connection) -> None: 

166 conn.execute( 

167 "CREATE TABLE IF NOT EXISTS index_meta (key TEXT PRIMARY KEY, value TEXT)", 

168 ) 

169 

170 @staticmethod 

171 def _get(conn: sqlite3.Connection, key: str) -> str | None: 

172 row = conn.execute( 

173 "SELECT value FROM index_meta WHERE key = ?", 

174 (key,), 

175 ).fetchone() 

176 return row["value"] if row else None 

177 

178 @staticmethod 

179 def _set(conn: sqlite3.Connection, key: str, value: str) -> None: 

180 conn.execute( 

181 "INSERT INTO index_meta (key, value) VALUES (?, ?) " 

182 "ON CONFLICT(key) DO UPDATE SET value = excluded.value", 

183 (key, value), 

184 ) 

185 

186 @staticmethod 

187 def get_dim(conn: sqlite3.Connection) -> int | None: 

188 value = IndexMetaTable._get(conn, "dim") 

189 return int(value) if value is not None else None 

190 

191 @staticmethod 

192 def set_dim(conn: sqlite3.Connection, dim: int) -> None: 

193 IndexMetaTable._set(conn, "dim", str(dim)) 

194 

195 @staticmethod 

196 def get_embed_model(conn: sqlite3.Connection) -> str | None: 

197 return IndexMetaTable._get(conn, "embed_model") 

198 

199 @staticmethod 

200 def set_embed_model(conn: sqlite3.Connection, name: str) -> None: 

201 IndexMetaTable._set(conn, "embed_model", name) 

202 

203 @staticmethod 

204 def get_schema_version(conn: sqlite3.Connection) -> int | None: 

205 value = IndexMetaTable._get(conn, "schema_version") 

206 return int(value) if value is not None else None 

207 

208 @staticmethod 

209 def set_schema_version(conn: sqlite3.Connection, version: int) -> None: 

210 IndexMetaTable._set(conn, "schema_version", str(version)) 

211 

212 @staticmethod 

213 def get_total_inserts(conn: sqlite3.Connection) -> int: 

214 value = IndexMetaTable._get(conn, "total_inserts") 

215 return int(value) if value is not None else 0 

216 

217 @staticmethod 

218 def increment_total_inserts(conn: sqlite3.Connection, count: int) -> None: 

219 """Add ``count`` to the stored counter in one SQL statement (INSERT 

220 .. ON CONFLICT DO UPDATE with arithmetic), instead of a separate 

221 read-then-write -- called once per add()/upsert_document(), so 

222 halving the statement count here is a real, if small, per-call 

223 saving. This only avoids a read-then-write race within this single 

224 statement; it does not make the counter safe against concurrent 

225 writers in general (callers still rely on the write FileLock for 

226 that). index_meta.value has TEXT affinity, so the incremented 

227 result is stored as its text representation -- get_total_inserts() 

228 already expects that (int(value)), so this is not a behavior 

229 change, only fewer statements. 

230 """ 

231 conn.execute( 

232 "INSERT INTO index_meta (key, value) VALUES ('total_inserts', ?) " 

233 "ON CONFLICT(key) DO UPDATE SET value = " 

234 "CAST(index_meta.value AS INTEGER) + CAST(excluded.value AS INTEGER)", 

235 (str(count),), 

236 ) 

237 

238 @staticmethod 

239 def reset_total_inserts(conn: sqlite3.Connection, count: int) -> None: 

240 """Set total_inserts to an absolute value -- distinct from 

241 increment_total_inserts(): used by compact()'s rebuild and by 

242 m0001_v1_to_v2 after copying live rows into a fresh file, where 

243 total_inserts must become exactly the live row count, not add to 

244 whatever the source file's counter held.""" 

245 IndexMetaTable._set(conn, "total_inserts", str(count))