• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

kimata / my-py-lib / 21089817782

17 Jan 2026 06:17AM UTC coverage: 60.961% (+0.02%) from 60.946%
21089817782

push

github

kimata
refactor: SQLite接続設定を永続設定と接続設定に分離

- init() を init_persistent() と init_connection() に分離
  - init_persistent(): DBファイルに永続化される設定(新規作成時のみ)
  - init_connection(): 接続ごとに必要な設定(毎回の接続時)
- cleanup_stale_files() を追加: CephFS/NFS環境で残存するWAL/SHMファイルを削除
- 既存DB接続時にも init_connection() を呼ぶように修正
- テストを新しい関数名に更新

これによりCephFS/NFS環境でのSQLite「locking protocol」エラーを防止

Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>

21 of 26 new or added lines in 1 file covered. (80.77%)

3387 of 5556 relevant lines covered (60.96%)

0.61 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

96.53
/src/my_lib/sqlite_util.py
1
#!/usr/bin/env python3
2
"""SQLiteデータベースのユーティリティ関数
3

4
CephFS/NFS/Kubernetes manual storage class環境に最適化。
5
WALモード + 排他ロック + mmap無効化により、分散ファイルシステムでも安全に動作。
6

7
環境変数による設定:
8
    SQLITE_JOURNAL_MODE: ジャーナルモード (WAL/DELETE/TRUNCATE/PERSIST) デフォルト: WAL
9
    SQLITE_MMAP_SIZE: mmapサイズ(バイト単位、0で無効化)デフォルト: 0
10
    SQLITE_LOCKING_MODE: ロックモード (NORMAL/EXCLUSIVE) デフォルト: EXCLUSIVE
11
    SQLITE_LOCK_MODE: fcntlロックモード (BLOCK/NONBLOCK) デフォルト: BLOCK
12
    SQLITE_CHECKPOINT_DIR: WALチェックポイント用ディレクトリ
13
"""
14

15
from __future__ import annotations
1✔
16

17
import contextlib
1✔
18
import fcntl
1✔
19
import logging
1✔
20
import os
1✔
21
import pathlib
1✔
22
import sqlite3
1✔
23
import time
1✔
24
from typing import Any
1✔
25

26

27
def init_persistent(conn: sqlite3.Connection) -> None:
1✔
28
    """
29
    DBファイルに永続化されるPRAGMA設定を行う(新規作成時のみ)
30

31
    Args:
32
        conn: SQLiteデータベース接続
33
    """
34
    # ジャーナルモード: WALモードはNFSでも動作するが、DELETEモードの方が安全な場合もある
35
    journal_mode = os.environ.get("SQLITE_JOURNAL_MODE", "WAL")
1✔
36
    conn.execute(f"PRAGMA journal_mode={journal_mode}")
1✔
37

38
    # ページサイズ: 標準的な4096バイトを使用
39
    conn.execute("PRAGMA page_size=4096")
1✔
40

41
    conn.commit()
1✔
42
    logging.info("SQLiteデータベースの永続設定を初期化しました")
1✔
43

44

45
def init_connection(conn: sqlite3.Connection, *, timeout: float = 60.0) -> None:
1✔
46
    """
47
    接続ごとに必要なPRAGMA設定を行う(毎回の接続時に呼ぶ)
48

49
    Args:
50
        conn: SQLiteデータベース接続
51
        timeout: データベース接続のタイムアウト時間(秒)
52
    """
53
    # 同期モード: NFSやローカルストレージではFULLが最も安全
54
    conn.execute("PRAGMA synchronous=FULL")
1✔
55

56
    # WALモード使用時の設定
57
    journal_mode = os.environ.get("SQLITE_JOURNAL_MODE", "WAL")
1✔
58
    if journal_mode == "WAL":
1✔
59
        conn.execute("PRAGMA wal_autocheckpoint=1000")
1✔
60
        conn.execute("PRAGMA journal_size_limit=67108864")  # 64MB
1✔
61

62
    # キャッシュサイズ: 控えめに設定(Pod のメモリ制限を考慮)
63
    conn.execute("PRAGMA cache_size=-32000")  # 約32MB
1✔
64

65
    # テンポラリストレージをメモリに設定
66
    conn.execute("PRAGMA temp_store=MEMORY")
1✔
67

68
    # mmapサイズ: NFSでは無効化(0で無効化)
69
    mmap_size = int(os.environ.get("SQLITE_MMAP_SIZE", "0"))
1✔
70
    conn.execute(f"PRAGMA mmap_size={mmap_size}")
1✔
71

72
    # ロックタイムアウト(NFSレイテンシを考慮)
73
    conn.execute(f"PRAGMA busy_timeout={int(timeout * 1000)}")
1✔
74

75
    # 外部キー制約を有効化
76
    conn.execute("PRAGMA foreign_keys=ON")
1✔
77

78
    # ロックモード: CephFS/NFSでは排他ロックモードが必要
79
    locking_mode = os.environ.get("SQLITE_LOCKING_MODE", "EXCLUSIVE")
1✔
80
    conn.execute(f"PRAGMA locking_mode={locking_mode}")
1✔
81

82
    conn.commit()
1✔
83
    logging.debug("SQLiteデータベースの接続設定を適用しました")
1✔
84

85

86
def cleanup_stale_files(db_path: pathlib.Path) -> None:
1✔
87
    """
88
    CephFS/NFS環境で残存するWAL/SHMファイルを削除する
89

90
    Podクラッシュ後などにこれらのファイルが残っていると
91
    「locking protocol」エラーが発生するため、接続前に削除する。
92

93
    Args:
94
        db_path: データベースファイルのパス
95
    """
96
    wal_path = db_path.with_suffix(db_path.suffix + "-wal")
1✔
97
    shm_path = db_path.with_suffix(db_path.suffix + "-shm")
1✔
98

99
    if wal_path.exists() or shm_path.exists():
1✔
NEW
100
        logging.warning("残存するWAL/SHMファイルを削除します: %s", db_path)
×
NEW
101
        with contextlib.suppress(Exception):
×
NEW
102
            wal_path.unlink(missing_ok=True)
×
NEW
103
        with contextlib.suppress(Exception):
×
NEW
104
            shm_path.unlink(missing_ok=True)
×
105

106

107
class DatabaseConnection:
1✔
108
    """SQLite接続をContext Managerとしても通常の関数としても使用可能にするラッパー"""
109

110
    def __init__(self, db_path: str | pathlib.Path, *, timeout: float = 60.0) -> None:
1✔
111
        """
112
        データベース接続の初期化
113

114
        Args:
115
            db_path: データベースファイルのパス
116
            timeout: データベース接続のタイムアウト時間(秒)
117
        """
118
        self.db_path = pathlib.Path(db_path)
1✔
119
        self.timeout = timeout
1✔
120
        self.conn: sqlite3.Connection | None = None
1✔
121

122
    def _acquire_lock(self, lock_file: Any) -> bool:
1✔
123
        """ロックの取得を試みる"""
124
        if os.environ.get("SQLITE_LOCK_MODE") == "NONBLOCK":
1✔
125
            try:
1✔
126
                fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
1✔
127
                return True
1✔
128
            except BlockingIOError:
1✔
129
                return False
1✔
130
        else:
131
            fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
1✔
132
            return True
1✔
133

134
    def _get_connection_params(self) -> dict[str, Any]:
1✔
135
        """SQLite接続パラメータを取得"""
136
        params: dict[str, Any] = {
1✔
137
            "timeout": self.timeout,
138
            "check_same_thread": False,
139
            "isolation_level": "DEFERRED",
140
        }
141

142
        checkpoint_dir = os.environ.get("SQLITE_CHECKPOINT_DIR")
1✔
143
        if checkpoint_dir:
1✔
144
            pathlib.Path(checkpoint_dir).mkdir(parents=True, exist_ok=True)
1✔
145

146
        return params
1✔
147

148
    def _create_connection(self) -> sqlite3.Connection:
1✔
149
        """実際の接続処理"""
150
        self.db_path.parent.mkdir(parents=True, exist_ok=True)
1✔
151

152
        is_new_db = not self.db_path.exists()
1✔
153

154
        if is_new_db:
1✔
155
            # 新規作成時は排他制御を行う
156
            lock_path = self.db_path.with_suffix(".lock")
1✔
157
            max_retries = 5
1✔
158
            retry_count = 0
1✔
159

160
            while retry_count < max_retries:
1✔
161
                try:
1✔
162
                    with lock_path.open("w") as lock_file:
1✔
163
                        if not self._acquire_lock(lock_file):
1✔
164
                            retry_count += 1
1✔
165
                            time.sleep(0.1 * retry_count)
1✔
166
                            continue
1✔
167

168
                        try:
1✔
169
                            is_new_db = not self.db_path.exists()
1✔
170
                            self.conn = sqlite3.connect(self.db_path, **self._get_connection_params())
1✔
171

172
                            if is_new_db:
1✔
173
                                init_persistent(self.conn)
1✔
174
                                logging.info("新規SQLiteデータベースを作成しました: %s", self.db_path)
1✔
175
                            init_connection(self.conn, timeout=self.timeout)
1✔
176
                        finally:
177
                            fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
1✔
178

179
                    with contextlib.suppress(Exception):
1✔
180
                        lock_path.unlink()
1✔
181
                    break
1✔
182

183
                except Exception:
1✔
184
                    retry_count += 1
1✔
185
                    if retry_count >= max_retries:
1✔
186
                        logging.exception("データベース作成時のロック取得に失敗しました")
1✔
187
                        raise
1✔
188
                    time.sleep(0.1 * retry_count)
1✔
189
        else:
190
            # 既存のデータベースへの接続
191
            cleanup_stale_files(self.db_path)
1✔
192
            self.conn = sqlite3.connect(self.db_path, **self._get_connection_params())
1✔
193
            init_connection(self.conn, timeout=self.timeout)
1✔
194
            logging.debug("既存のSQLiteデータベースに接続しました: %s", self.db_path)
1✔
195

196
        assert self.conn is not None  # noqa: S101
1✔
197
        return self.conn
1✔
198

199
    def __enter__(self) -> sqlite3.Connection:
1✔
200
        """Context Manager として使用する場合のenter"""
201
        return self._create_connection()
1✔
202

203
    def __exit__(
1✔
204
        self,
205
        exc_type: type[BaseException] | None,
206
        _exc_val: BaseException | None,
207
        _exc_tb: Any,
208
    ) -> None:
209
        """Context Manager として使用する場合のexit"""
210
        if self.conn is not None:
1✔
211
            if exc_type is None:
1✔
212
                self.conn.commit()
1✔
213
            else:
214
                self.conn.rollback()
1✔
215
            self.conn.close()
1✔
216

217
    def get(self) -> sqlite3.Connection:
1✔
218
        """通常の関数として使用する場合(使用後は必ずcloseすること)"""
219
        return self._create_connection()
1✔
220

221

222
def connect(db_path: str | pathlib.Path, *, timeout: float = 60.0) -> DatabaseConnection:
1✔
223
    """
224
    Kubernetes manual storage class環境に適したSQLiteデータベースに接続する
225

226
    Context Managerとしても通常の関数としても使用可能
227

228
    Args:
229
        db_path: データベースファイルのパス
230
        timeout: データベース接続のタイムアウト時間(秒)
231

232
    Returns:
233
        DatabaseConnection: Context Managerとしても通常の接続取得にも使用可能
234

235
    Usage:
236
        # Context Manager として使用
237
        with my_lib.sqlite_util.connect(db_path) as conn:
238
            conn.execute("SELECT * FROM table")
239

240
        # 通常の関数として使用
241
        db_conn = my_lib.sqlite_util.connect(db_path)
242
        conn = db_conn.get()
243
        try:
244
            conn.execute("SELECT * FROM table")
245
        finally:
246
            conn.close()
247
    """
248
    return DatabaseConnection(db_path, timeout=timeout)
1✔
249

250

251
def recover(db_path: str | pathlib.Path) -> None:
1✔
252
    """データベースの復旧を試みる"""
253
    try:
1✔
254
        db_path = pathlib.Path(db_path)
1✔
255

256
        journal_mode = os.environ.get("SQLITE_JOURNAL_MODE", "WAL")
1✔
257

258
        if journal_mode == "WAL":
1✔
259
            wal_path = db_path.with_suffix(db_path.suffix + "-wal")
1✔
260
            shm_path = db_path.with_suffix(db_path.suffix + "-shm")
1✔
261

262
            if wal_path.exists():
1✔
263
                logging.warning("WALファイル %s を削除してデータベースを復旧します", wal_path)
1✔
264
                wal_path.unlink()
1✔
265

266
            if shm_path.exists():
1✔
267
                logging.warning("共有メモリファイル %s を削除します", shm_path)
1✔
268
                shm_path.unlink()
1✔
269
        else:
270
            journal_path = db_path.with_suffix(db_path.suffix + "-journal")
1✔
271
            if journal_path.exists():
1✔
272
                logging.warning("ジャーナルファイル %s を削除してデータベースを復旧します", journal_path)
1✔
273
                journal_path.unlink()
1✔
274

275
        try:
1✔
276
            conn = sqlite3.connect(db_path, timeout=5.0)
1✔
277
            result = conn.execute("PRAGMA integrity_check").fetchone()
1✔
278
            if result[0] != "ok":
1✔
279
                raise sqlite3.DatabaseError(f"整合性チェック失敗: {result[0]}")
1✔
280

281
            conn.execute("PRAGMA quick_check")
1✔
282

283
            if os.environ.get("SQLITE_AUTO_VACUUM") == "1":
1✔
284
                logging.info("データベースのVACUUMを実行します")
1✔
285
                conn.execute("VACUUM")
1✔
286

287
            conn.close()
1✔
288
            logging.info("データベースの整合性チェックが成功しました")
1✔
289

290
        except Exception:
1✔
291
            logging.exception("データベースの整合性チェックに失敗")
1✔
292
            backup_path = db_path.with_suffix(f".backup.{int(time.time())}")
1✔
293
            db_path.rename(backup_path)
1✔
294
            logging.warning("破損したデータベースを %s にバックアップし、新規作成します", backup_path)
1✔
295

296
    except Exception:
1✔
297
        logging.exception("データベース復旧中にエラーが発生")
1✔
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc