Skip to main content

platform/storage/
db.rs

1//! 数据库连接和操作管理
2//!
3//! 提供基于 sqlx 的数据库连接池和基本操作
4
5use anyhow::Result;
6use sqlx::sqlite::{SqliteConnectOptions, SqlitePool, SqlitePoolOptions};
7use std::path::Path;
8use std::str::FromStr;
9use std::time::Duration;
10
11/// 数据库管理器
12#[derive(Clone)]
13pub struct Database {
14    pool: SqlitePool,
15}
16
17impl Database {
18    /// 创建新的数据库实例
19    ///
20    /// # Arguments
21    /// * `path` - 数据库文件存储目录路径,必须已存在
22    ///   主数据库文件将存储为 `{path}/actrix.db`
23    pub async fn new<P: AsRef<Path>>(path: P) -> Result<Self> {
24        let db_file = path.as_ref().join("actrix.db");
25
26        // 创建连接选项并启用 WAL 模式
27        let options = SqliteConnectOptions::from_str(&format!("sqlite:{}", db_file.display()))?
28            .create_if_missing(true)
29            .journal_mode(sqlx::sqlite::SqliteJournalMode::Wal)
30            .synchronous(sqlx::sqlite::SqliteSynchronous::Normal)
31            .busy_timeout(Duration::from_secs(5));
32
33        // 创建连接池
34        let pool = SqlitePoolOptions::new()
35            .max_connections(10)
36            .connect_with(options)
37            .await?;
38
39        let db = Self { pool };
40
41        // 初始化数据库表结构
42        db.initialize_schema().await?;
43
44        Ok(db)
45    }
46
47    /// 初始化数据库表结构
48    async fn initialize_schema(&self) -> Result<()> {
49        // 创建 Realm 表
50        sqlx::query(
51            "CREATE TABLE IF NOT EXISTS realm (
52                id INTEGER PRIMARY KEY AUTOINCREMENT,
53                name TEXT NOT NULL,
54                status TEXT NOT NULL DEFAULT 'Active',
55                enabled INTEGER NOT NULL DEFAULT 1,
56                expires_at INTEGER,
57                created_at INTEGER NOT NULL,
58                updated_at INTEGER,
59                secret_current TEXT NOT NULL DEFAULT '',
60                secret_previous_hash TEXT,
61                secret_previous_valid_until INTEGER
62            )",
63        )
64        .execute(&self.pool)
65        .await?;
66
67        // Set autoincrement start to 2^25 = 33554432
68        // Only insert if not already present (fresh database)
69        sqlx::query("INSERT OR IGNORE INTO sqlite_sequence(name, seq) VALUES('realm', 33554431)")
70            .execute(&self.pool)
71            .await?;
72
73        // 创建访问控制列表表
74        sqlx::query(
75            "CREATE TABLE IF NOT EXISTS actoracl (
76                rowid INTEGER PRIMARY KEY AUTOINCREMENT,
77                realm_id INTEGER NOT NULL,
78                source_realm_id INTEGER,
79                from_type TEXT NOT NULL,
80                to_type TEXT NOT NULL,
81                access INTEGER NOT NULL
82            )",
83        )
84        .execute(&self.pool)
85        .await?;
86
87        // 创建索引
88        sqlx::query(
89            "CREATE INDEX IF NOT EXISTS idx_realm_name
90             ON realm(name)",
91        )
92        .execute(&self.pool)
93        .await?;
94
95        sqlx::query(
96            "CREATE INDEX IF NOT EXISTS idx_actoracl_realm_id
97             ON actoracl(realm_id)",
98        )
99        .execute(&self.pool)
100        .await?;
101
102        sqlx::query(
103            "CREATE INDEX IF NOT EXISTS idx_actoracl_lookup
104             ON actoracl(realm_id, source_realm_id, from_type, to_type)",
105        )
106        .execute(&self.pool)
107        .await?;
108
109        // Pending registration data: AIS writes, signaling reads on WS connect
110        sqlx::query(
111            "CREATE TABLE IF NOT EXISTS pending_registration (
112                serial_number INTEGER PRIMARY KEY,
113                realm_id INTEGER NOT NULL,
114                service_spec_blob BLOB,
115                ws_address TEXT,
116                created_at INTEGER NOT NULL
117            )",
118        )
119        .execute(&self.pool)
120        .await?;
121
122        // Migrate: add ws_address column if it doesn't exist (for existing databases)
123        let _ = sqlx::query("ALTER TABLE pending_registration ADD COLUMN ws_address TEXT")
124            .execute(&self.pool)
125            .await; // intentionally ignore error (column may already exist)
126
127        // MFR (Manufacturer Registry) tables
128        // name = GitHub user/org login (lowercased), serves as manufacturer identity
129        sqlx::query(
130            "CREATE TABLE IF NOT EXISTS mfr (
131                id           INTEGER PRIMARY KEY AUTOINCREMENT,
132                name         TEXT    NOT NULL UNIQUE,
133                public_key   TEXT    NOT NULL DEFAULT '',
134                contact      TEXT,
135                status       TEXT    NOT NULL DEFAULT 'pending',
136                created_at   INTEGER NOT NULL,
137                updated_at   INTEGER,
138                verified_at  INTEGER,
139                suspended_at INTEGER,
140                revoked_at   INTEGER,
141                key_expires_at  INTEGER
142            )",
143        )
144        .execute(&self.pool)
145        .await?;
146
147        // GitHub verification challenge for identity verification
148        sqlx::query(
149            "CREATE TABLE IF NOT EXISTS mfr_challenge (
150                id          INTEGER PRIMARY KEY AUTOINCREMENT,
151                mfr_id      INTEGER NOT NULL REFERENCES mfr(id),
152                token       TEXT    NOT NULL,
153                verify_url  TEXT    NOT NULL DEFAULT '',
154                expires_at  INTEGER NOT NULL,
155                verified_at INTEGER,
156                created_at  INTEGER NOT NULL
157            )",
158        )
159        .execute(&self.pool)
160        .await?;
161
162        sqlx::query(
163            "CREATE TABLE IF NOT EXISTS mfr_package (
164                id           INTEGER PRIMARY KEY AUTOINCREMENT,
165                mfr_id       INTEGER NOT NULL REFERENCES mfr(id),
166                manufacturer TEXT    NOT NULL,
167                name         TEXT    NOT NULL,
168                version      TEXT    NOT NULL,
169                type_str     TEXT    NOT NULL,
170                target       TEXT    NOT NULL,
171                manifest     TEXT    NOT NULL,
172                signature    TEXT    NOT NULL,
173                signing_key_id TEXT,
174                status       TEXT    NOT NULL DEFAULT 'active',
175                published_at INTEGER NOT NULL,
176                revoked_at   INTEGER,
177                UNIQUE(manufacturer, name, version, target)
178            )",
179        )
180        .execute(&self.pool)
181        .await?;
182
183        // Migrate: add target column if it doesn't exist (for existing databases)
184        let _ = sqlx::query(
185            "ALTER TABLE mfr_package ADD COLUMN target TEXT NOT NULL DEFAULT 'wasm32-wasip1'",
186        )
187        .execute(&self.pool)
188        .await; // intentionally ignore error (column may already exist)
189
190        sqlx::query("CREATE INDEX IF NOT EXISTS idx_mfr_package_type ON mfr_package(type_str)")
191            .execute(&self.pool)
192            .await?;
193
194        sqlx::query(
195            "CREATE INDEX IF NOT EXISTS idx_mfr_package_mfr ON mfr_package(mfr_id, status)",
196        )
197        .execute(&self.pool)
198        .await?;
199
200        // Migrate: add proto_files column for proto filing (JSON text, nullable)
201        let _ = sqlx::query("ALTER TABLE mfr_package ADD COLUMN proto_files TEXT")
202            .execute(&self.pool)
203            .await; // intentionally ignore error (column may already exist)
204
205        // Migrate: record the MFR key that authenticated each newly published package.
206        // Existing rows remain nullable and are resolved from their signed manifest.
207        let _ = sqlx::query("ALTER TABLE mfr_package ADD COLUMN signing_key_id TEXT")
208            .execute(&self.pool)
209            .await; // intentionally ignore error (column may already exist)
210
211        // Migrate: add key_id column (auto-assigned on activate/renew)
212        let _ = sqlx::query("ALTER TABLE mfr ADD COLUMN key_id TEXT NOT NULL DEFAULT ''")
213            .execute(&self.pool)
214            .await;
215
216        // MFR key history: stores retired public keys for JWKS-style rotation
217        sqlx::query(
218            "CREATE TABLE IF NOT EXISTS mfr_key_history (
219                id           INTEGER PRIMARY KEY AUTOINCREMENT,
220                mfr_id       INTEGER NOT NULL REFERENCES mfr(id),
221                key_id       TEXT    NOT NULL,
222                public_key   TEXT    NOT NULL,
223                status       TEXT    NOT NULL DEFAULT 'retired',
224                created_at   INTEGER NOT NULL,
225                retired_at   INTEGER NOT NULL
226            )",
227        )
228        .execute(&self.pool)
229        .await?;
230
231        sqlx::query(
232            "CREATE INDEX IF NOT EXISTS idx_mfr_key_history_lookup
233             ON mfr_key_history(mfr_id, key_id)",
234        )
235        .execute(&self.pool)
236        .await?;
237
238        // Publish nonce table for Challenge-Response authentication on /mfr/pkg/publish
239        sqlx::query(
240            "CREATE TABLE IF NOT EXISTS mfr_publish_nonce (
241                id         INTEGER PRIMARY KEY AUTOINCREMENT,
242                mfr_id     INTEGER NOT NULL REFERENCES mfr(id),
243                nonce      BLOB    NOT NULL UNIQUE,
244                status     TEXT    NOT NULL DEFAULT 'pending',
245                created_at INTEGER NOT NULL,
246                expires_at INTEGER NOT NULL
247            )",
248        )
249        .execute(&self.pool)
250        .await?;
251
252        sqlx::query(
253            "CREATE INDEX IF NOT EXISTS idx_mfr_publish_nonce_expires
254             ON mfr_publish_nonce(expires_at)",
255        )
256        .execute(&self.pool)
257        .await?;
258
259        // AIS unpublished package manufacturer-proof nonce table. AIS inserts
260        // after manufacturer_auth_signature verification; the unique index is the
261        // replay guard.
262        sqlx::query(
263            "CREATE TABLE IF NOT EXISTS ais_manufacturer_auth_nonce (
264                id           INTEGER PRIMARY KEY AUTOINCREMENT,
265                manufacturer TEXT    NOT NULL,
266                key_id       TEXT    NOT NULL,
267                nonce        BLOB    NOT NULL,
268                created_at   INTEGER NOT NULL,
269                expires_at   INTEGER NOT NULL,
270                UNIQUE(manufacturer, key_id, nonce)
271            )",
272        )
273        .execute(&self.pool)
274        .await?;
275
276        sqlx::query(
277            "CREATE INDEX IF NOT EXISTS idx_ais_manufacturer_auth_nonce_expires
278             ON ais_manufacturer_auth_nonce(expires_at)",
279        )
280        .execute(&self.pool)
281        .await?;
282
283        // Backfill key_id for existing MFRs that have a public_key but empty key_id.
284        // This runs on every startup but is a no-op when all rows already have a key_id.
285        {
286            use base64::Engine as _;
287            use sha2::{Digest, Sha256};
288
289            let rows: Vec<(i64, String)> = sqlx::query_as(
290                "SELECT id, public_key FROM mfr WHERE key_id = '' AND public_key != ''",
291            )
292            .fetch_all(&self.pool)
293            .await
294            .unwrap_or_default();
295
296            for (id, public_key_b64) in rows {
297                if let Ok(bytes) = base64::engine::general_purpose::STANDARD.decode(&public_key_b64)
298                {
299                    let hash = Sha256::digest(&bytes);
300                    let hex_str: String = hash.iter().map(|b| format!("{b:02x}")).collect();
301                    let key_id = format!("mfr-{}", &hex_str[..16]);
302                    let _ = sqlx::query("UPDATE mfr SET key_id = ? WHERE id = ?")
303                        .bind(&key_id)
304                        .bind(id)
305                        .execute(&self.pool)
306                        .await;
307                    crate::recording::info!(
308                        "backfilled key_id from public_key fingerprint: id={}, key_id={}",
309                        id,
310                        key_id
311                    );
312                }
313            }
314        }
315
316        // AIS renewal token table — only stores SHA-256(token), never the raw token.
317        sqlx::query(
318            "CREATE TABLE IF NOT EXISTS ais_renewal_token (
319                id INTEGER PRIMARY KEY AUTOINCREMENT,
320                actor_id TEXT NOT NULL,
321                token_hash BLOB NOT NULL UNIQUE,
322                expires_at INTEGER NOT NULL,
323                created_at INTEGER NOT NULL
324            )",
325        )
326        .execute(&self.pool)
327        .await?;
328
329        sqlx::query(
330            "CREATE INDEX IF NOT EXISTS idx_ais_renewal_token_actor
331             ON ais_renewal_token(actor_id)",
332        )
333        .execute(&self.pool)
334        .await?;
335
336        sqlx::query(
337            "CREATE INDEX IF NOT EXISTS idx_ais_renewal_token_expires
338             ON ais_renewal_token(expires_at)",
339        )
340        .execute(&self.pool)
341        .await?;
342
343        Ok(())
344    }
345
346    /// 获取数据库连接池
347    pub fn get_pool(&self) -> &SqlitePool {
348        &self.pool
349    }
350
351    /// 执行 SQL 语句并返回影响的行数
352    pub async fn execute(&self, sql: &str) -> Result<u64> {
353        let result = sqlx::query(sql).execute(&self.pool).await?;
354        Ok(result.rows_affected())
355    }
356}
357
358use tokio::sync::OnceCell;
359
360/// 全局数据库实例
361static GLOBAL_DATABASE: OnceCell<Database> = OnceCell::const_new();
362
363/// 设置全局数据库路径
364pub async fn set_db_path(path: &Path) -> Result<()> {
365    let database = Database::new(path).await?;
366    GLOBAL_DATABASE
367        .set(database)
368        .map_err(|_| anyhow::anyhow!("Database already initialized"))?;
369    Ok(())
370}
371
372/// 获取全局数据库实例
373pub fn get_database() -> &'static Database {
374    GLOBAL_DATABASE
375        .get()
376        .expect("Database not initialized. Call set_db_path first.")
377}
378
379/// 检查数据库是否已初始化
380pub fn is_database_initialized() -> bool {
381    GLOBAL_DATABASE.get().is_some()
382}