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

kimata / my-py-lib / 21092969449

17 Jan 2026 10:39AM UTC coverage: 62.974% (+2.0%) from 60.961%
21092969449

push

github

kimata
refactor: 型安全性向上とコード品質改善(第4弾)

## 主な変更

### センサー
- ADS1015/ADS1115 の重複コードを ads_base.py に統合
- センサー系の型注釈改善

### json_util.py
- iso_pattern の重複定義をモジュールレベル定数に統一
- DateTimeJSONEncoder.default() の type: ignore を削除

### store/mercari
- MercariItem dataclass を導入し、item 辞書を型安全に
- ProgressObserver Protocol の型定義を MercariItem に変更

### store/amazon
- AmazonItem.to_dict() を dataclasses.asdict() で簡潔化

### その他
- flask_util.py の type: ignore を cast() に置換
- lifecycle_manager.py を削除(lifecycle/manager.py に移行済み)
- 各種型注釈・docstring の改善

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

191 of 246 new or added lines in 29 files covered. (77.64%)

112 existing lines in 12 files now uncovered.

3439 of 5461 relevant lines covered (62.97%)

0.63 hits per line

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

96.71
/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 dataclasses import dataclass
1✔
25
from typing import Any, Literal
1✔
26

27
# sqlite3.connect の isolation_level パラメータの型
28
IsolationLevel = Literal["DEFERRED", "EXCLUSIVE", "IMMEDIATE"] | None
1✔
29

30

31
@dataclass(frozen=True)
1✔
32
class SQLiteConnectionParams:
1✔
33
    """SQLite接続パラメータを保持するデータクラス"""
34

35
    timeout: float
1✔
36
    check_same_thread: bool
1✔
37
    isolation_level: IsolationLevel
1✔
38

39

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

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

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

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

57

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

62
    Args:
63
        conn: SQLiteデータベース接続
64
        timeout: データベース接続のタイムアウト時間(秒)
65
    """
66
    # 同期モード: NFSやローカルストレージではFULLが最も安全
67
    conn.execute("PRAGMA synchronous=FULL")
1✔
68

69
    # WALモード使用時の設定
70
    journal_mode = os.environ.get("SQLITE_JOURNAL_MODE", "WAL")
1✔
71
    if journal_mode == "WAL":
1✔
72
        conn.execute("PRAGMA wal_autocheckpoint=1000")
1✔
73
        conn.execute("PRAGMA journal_size_limit=67108864")  # 64MB
1✔
74

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

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

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

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

88
    # 外部キー制約を有効化
89
    conn.execute("PRAGMA foreign_keys=ON")
1✔
90

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

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

98

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

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

106
    Args:
107
        db_path: データベースファイルのパス
108
    """
109
    wal_path = db_path.with_suffix(db_path.suffix + "-wal")
1✔
110
    shm_path = db_path.with_suffix(db_path.suffix + "-shm")
1✔
111

112
    if wal_path.exists() or shm_path.exists():
1✔
UNCOV
113
        logging.warning("残存するWAL/SHMファイルを削除します: %s", db_path)
×
UNCOV
114
        with contextlib.suppress(Exception):
×
UNCOV
115
            wal_path.unlink(missing_ok=True)
×
UNCOV
116
        with contextlib.suppress(Exception):
×
UNCOV
117
            shm_path.unlink(missing_ok=True)
×
118

119

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

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

127
        Args:
128
            db_path: データベースファイルのパス
129
            timeout: データベース接続のタイムアウト時間(秒)
130
        """
131
        self.db_path = pathlib.Path(db_path)
1✔
132
        self.timeout = timeout
1✔
133
        self.conn: sqlite3.Connection | None = None
1✔
134

135
    def _acquire_lock(self, lock_file: Any) -> bool:
1✔
136
        """ロックの取得を試みる"""
137
        if os.environ.get("SQLITE_LOCK_MODE") == "NONBLOCK":
1✔
138
            try:
1✔
139
                fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
1✔
140
                return True
1✔
141
            except BlockingIOError:
1✔
142
                return False
1✔
143
        else:
144
            fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
1✔
145
            return True
1✔
146

147
    def _get_connection_params(self) -> SQLiteConnectionParams:
1✔
148
        """SQLite接続パラメータを取得"""
149
        checkpoint_dir = os.environ.get("SQLITE_CHECKPOINT_DIR")
1✔
150
        if checkpoint_dir:
1✔
151
            pathlib.Path(checkpoint_dir).mkdir(parents=True, exist_ok=True)
1✔
152

153
        return SQLiteConnectionParams(
1✔
154
            timeout=self.timeout,
155
            check_same_thread=False,
156
            isolation_level="DEFERRED",
157
        )
158

159
    def _create_connection(self) -> sqlite3.Connection:
1✔
160
        """実際の接続処理"""
161
        self.db_path.parent.mkdir(parents=True, exist_ok=True)
1✔
162

163
        is_new_db = not self.db_path.exists()
1✔
164

165
        if is_new_db:
1✔
166
            # 新規作成時は排他制御を行う
167
            lock_path = self.db_path.with_suffix(".lock")
1✔
168
            max_retries = 5
1✔
169
            retry_count = 0
1✔
170

171
            while retry_count < max_retries:
1✔
172
                try:
1✔
173
                    with lock_path.open("w") as lock_file:
1✔
174
                        if not self._acquire_lock(lock_file):
1✔
175
                            retry_count += 1
1✔
176
                            time.sleep(0.1 * retry_count)
1✔
177
                            continue
1✔
178

179
                        try:
1✔
180
                            is_new_db = not self.db_path.exists()
1✔
181
                            params = self._get_connection_params()
1✔
182
                            self.conn = sqlite3.connect(
1✔
183
                                self.db_path,
184
                                timeout=params.timeout,
185
                                check_same_thread=params.check_same_thread,
186
                                isolation_level=params.isolation_level,
187
                            )
188

189
                            if is_new_db:
1✔
190
                                init_persistent(self.conn)
1✔
191
                                logging.info("新規SQLiteデータベースを作成しました: %s", self.db_path)
1✔
192
                            init_connection(self.conn, timeout=self.timeout)
1✔
193
                        finally:
194
                            fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
1✔
195

196
                    with contextlib.suppress(Exception):
1✔
197
                        lock_path.unlink()
1✔
198
                    break
1✔
199

200
                except Exception:
1✔
201
                    retry_count += 1
1✔
202
                    if retry_count >= max_retries:
1✔
203
                        logging.exception("データベース作成時のロック取得に失敗しました")
1✔
204
                        raise
1✔
205
                    time.sleep(0.1 * retry_count)
1✔
206
        else:
207
            # 既存のデータベースへの接続
208
            cleanup_stale_files(self.db_path)
1✔
209
            params = self._get_connection_params()
1✔
210
            self.conn = sqlite3.connect(
1✔
211
                self.db_path,
212
                timeout=params.timeout,
213
                check_same_thread=params.check_same_thread,
214
                isolation_level=params.isolation_level,
215
            )
216
            init_connection(self.conn, timeout=self.timeout)
1✔
217
            logging.debug("既存のSQLiteデータベースに接続しました: %s", self.db_path)
1✔
218

219
        assert self.conn is not None  # noqa: S101
1✔
220
        return self.conn
1✔
221

222
    def __enter__(self) -> sqlite3.Connection:
1✔
223
        """Context Manager として使用する場合のenter"""
224
        return self._create_connection()
1✔
225

226
    def __exit__(
1✔
227
        self,
228
        exc_type: type[BaseException] | None,
229
        _exc_val: BaseException | None,
230
        _exc_tb: Any,
231
    ) -> None:
232
        """Context Manager として使用する場合のexit"""
233
        if self.conn is not None:
1✔
234
            if exc_type is None:
1✔
235
                self.conn.commit()
1✔
236
            else:
237
                self.conn.rollback()
1✔
238
            self.conn.close()
1✔
239

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

244

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

249
    Context Managerとしても通常の関数としても使用可能
250

251
    Args:
252
        db_path: データベースファイルのパス
253
        timeout: データベース接続のタイムアウト時間(秒)
254

255
    Returns:
256
        DatabaseConnection: Context Managerとしても通常の接続取得にも使用可能
257

258
    Usage:
259
        # Context Manager として使用
260
        with my_lib.sqlite_util.connect(db_path) as conn:
261
            conn.execute("SELECT * FROM table")
262

263
        # 通常の関数として使用
264
        db_conn = my_lib.sqlite_util.connect(db_path)
265
        conn = db_conn.get()
266
        try:
267
            conn.execute("SELECT * FROM table")
268
        finally:
269
            conn.close()
270
    """
271
    return DatabaseConnection(db_path, timeout=timeout)
1✔
272

273

274
def recover(db_path: str | pathlib.Path) -> None:
1✔
275
    """データベースの復旧を試みる"""
276
    try:
1✔
277
        db_path = pathlib.Path(db_path)
1✔
278

279
        journal_mode = os.environ.get("SQLITE_JOURNAL_MODE", "WAL")
1✔
280

281
        if journal_mode == "WAL":
1✔
282
            wal_path = db_path.with_suffix(db_path.suffix + "-wal")
1✔
283
            shm_path = db_path.with_suffix(db_path.suffix + "-shm")
1✔
284

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

289
            if shm_path.exists():
1✔
290
                logging.warning("共有メモリファイル %s を削除します", shm_path)
1✔
291
                shm_path.unlink()
1✔
292
        else:
293
            journal_path = db_path.with_suffix(db_path.suffix + "-journal")
1✔
294
            if journal_path.exists():
1✔
295
                logging.warning("ジャーナルファイル %s を削除してデータベースを復旧します", journal_path)
1✔
296
                journal_path.unlink()
1✔
297

298
        try:
1✔
299
            conn = sqlite3.connect(db_path, timeout=5.0)
1✔
300
            result = conn.execute("PRAGMA integrity_check").fetchone()
1✔
301
            if result[0] != "ok":
1✔
302
                raise sqlite3.DatabaseError(f"整合性チェック失敗: {result[0]}")
1✔
303

304
            conn.execute("PRAGMA quick_check")
1✔
305

306
            if os.environ.get("SQLITE_AUTO_VACUUM") == "1":
1✔
307
                logging.info("データベースのVACUUMを実行します")
1✔
308
                conn.execute("VACUUM")
1✔
309

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

313
        except sqlite3.Error:
1✔
314
            logging.exception("データベースの整合性チェックに失敗")
1✔
315
            backup_path = db_path.with_suffix(f".backup.{int(time.time())}")
1✔
316
            db_path.rename(backup_path)
1✔
317
            logging.warning("破損したデータベースを %s にバックアップし、新規作成します", backup_path)
1✔
318

319
    except OSError:
1✔
320
        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