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

payjoin / rust-payjoin / 29448737091

15 Jul 2026 08:33PM UTC coverage: 86.187% (+0.2%) from 85.938%
29448737091

Pull #1707

github

web-flow
Merge bb5eb7fc2 into 1da296b10
Pull Request #1707: Prepare release-ready C# NuGet package workflow

13758 of 15963 relevant lines covered (86.19%)

345.48 hits per line

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

82.9
/payjoin-cli/src/db/v2.rs
1
use std::sync::Arc;
2

3
use payjoin::persist::SessionPersister;
4
use payjoin::receive::v2::SessionEvent as ReceiverSessionEvent;
5
use payjoin::send::v2::SessionEvent as SenderSessionEvent;
6
use payjoin::HpkePublicKey;
7
use rusqlite::{params, OptionalExtension};
8

9
use super::*;
10

11
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
12
pub(crate) struct SessionId(pub(crate) uuid::Uuid);
13

14
impl std::fmt::Display for SessionId {
15
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!(f, "{}", self.0) }
60✔
16
}
17

18
#[derive(Clone, Debug)]
19
pub(crate) struct SenderPersister {
20
    db: Arc<Database>,
21
    session_id: SessionId,
22
}
23

24
impl SenderPersister {
25
    pub fn new(
9✔
26
        db: Arc<Database>,
9✔
27
        pj_uri: &str,
9✔
28
        receiver_pubkey: &HpkePublicKey,
9✔
29
    ) -> crate::db::Result<Self> {
9✔
30
        let conn = db.get_connection()?;
9✔
31
        let receiver_pubkey_bytes = receiver_pubkey.to_compressed_bytes();
9✔
32

33
        let (duplicate_uri, duplicate_rk): (bool, bool) = conn.query_row(
9✔
34
            "SELECT \
9✔
35
                EXISTS(SELECT 1 FROM send_sessions WHERE pj_uri = ?1), \
9✔
36
                EXISTS(SELECT 1 FROM send_sessions WHERE receiver_pubkey = ?2)",
9✔
37
            params![pj_uri, &receiver_pubkey_bytes],
9✔
38
            |row| Ok((row.get(0)?, row.get(1)?)),
9✔
39
        )?;
×
40

41
        if duplicate_uri {
9✔
42
            return Err(Error::DuplicateSendSession(DuplicateKind::Uri));
2✔
43
        }
7✔
44
        if duplicate_rk {
7✔
45
            return Err(Error::DuplicateSendSession(DuplicateKind::ReceiverPubkey));
1✔
46
        }
6✔
47

48
        let session_id = uuid::Uuid::now_v7();
6✔
49
        conn.execute(
6✔
50
            "INSERT INTO send_sessions (session_id, pj_uri, receiver_pubkey) VALUES (?1, ?2, ?3)",
6✔
51
            params![session_id.to_string(), pj_uri, &receiver_pubkey_bytes],
6✔
52
        )?;
6✔
53

54
        Ok(Self { db, session_id: SessionId(session_id) })
6✔
55
    }
9✔
56

57
    pub fn from_id(db: Arc<Database>, id: SessionId) -> Self { Self { db, session_id: id } }
4✔
58

59
    pub fn session_id(&self) -> SessionId { self.session_id.clone() }
18✔
60
}
61
impl SessionPersister for SenderPersister {
62
    type SessionEvent = SenderSessionEvent;
63
    type InternalStorageError = crate::db::error::Error;
64

65
    fn save_event(
8✔
66
        &self,
8✔
67
        event: SenderSessionEvent,
8✔
68
    ) -> std::result::Result<(), Self::InternalStorageError> {
8✔
69
        let conn = self.db.get_connection()?;
8✔
70
        let event_data = serde_json::to_string(&event).map_err(Error::Serialize)?;
8✔
71

72
        conn.execute(
8✔
73
            "INSERT INTO send_session_events (session_id, event_data, created_at) VALUES (?1, ?2, ?3)",
8✔
74
            params![self.session_id.0.to_string(), event_data, now()],
8✔
75
        )?;
8✔
76

77
        Ok(())
8✔
78
    }
8✔
79

80
    fn load(
4✔
81
        &self,
4✔
82
    ) -> std::result::Result<Box<dyn Iterator<Item = SenderSessionEvent>>, Self::InternalStorageError>
4✔
83
    {
84
        let conn = self.db.get_connection()?;
4✔
85
        let mut stmt = conn.prepare(
4✔
86
            "SELECT event_data FROM send_session_events WHERE session_id = ?1 ORDER BY id ASC",
4✔
87
        )?;
4✔
88

89
        let event_rows = stmt.query_map(params![self.session_id.0.to_string()], |row| {
8✔
90
            let event_data: String = row.get(0)?;
8✔
91
            Ok(event_data)
8✔
92
        })?;
8✔
93

94
        let events: Vec<SenderSessionEvent> = event_rows
4✔
95
            .map(|row| {
8✔
96
                let event_data = row.expect("Failed to read event data from database");
8✔
97
                serde_json::from_str::<SenderSessionEvent>(&event_data)
8✔
98
                    .expect("Database corruption: failed to deserialize session event")
8✔
99
            })
8✔
100
            .collect();
4✔
101

102
        Ok(Box::new(events.into_iter()))
4✔
103
    }
4✔
104

105
    fn close(&self) -> std::result::Result<(), Self::InternalStorageError> {
4✔
106
        let conn = self.db.get_connection()?;
4✔
107

108
        conn.execute(
4✔
109
            "UPDATE send_sessions SET completed_at = ?1 WHERE session_id = ?2",
4✔
110
            params![now(), self.session_id.0.to_string()],
4✔
111
        )?;
4✔
112

113
        Ok(())
4✔
114
    }
4✔
115
}
116

117
#[derive(Clone)]
118
pub(crate) struct ReceiverPersister {
119
    db: Arc<Database>,
120
    session_id: SessionId,
121
}
122

123
impl ReceiverPersister {
124
    pub fn new(db: Arc<Database>) -> crate::db::Result<Self> {
5✔
125
        let conn = db.get_connection()?;
5✔
126

127
        let session_id = uuid::Uuid::now_v7();
5✔
128
        conn.execute(
5✔
129
            "INSERT INTO receive_sessions (session_id) VALUES (?1)",
5✔
130
            params![session_id.to_string()],
5✔
131
        )?;
5✔
132

133
        Ok(Self { db, session_id: SessionId(session_id) })
5✔
134
    }
5✔
135

136
    pub fn from_id(db: Arc<Database>, id: SessionId) -> Self { Self { db, session_id: id } }
6✔
137

138
    pub fn session_id(&self) -> SessionId { self.session_id.clone() }
34✔
139
}
140

141
impl SessionPersister for ReceiverPersister {
142
    type SessionEvent = ReceiverSessionEvent;
143
    type InternalStorageError = crate::db::error::Error;
144

145
    fn save_event(
17✔
146
        &self,
17✔
147
        event: ReceiverSessionEvent,
17✔
148
    ) -> std::result::Result<(), Self::InternalStorageError> {
17✔
149
        let conn = self.db.get_connection()?;
17✔
150
        let event_data = serde_json::to_string(&event).map_err(Error::Serialize)?;
17✔
151

152
        conn.execute(
17✔
153
            "INSERT INTO receive_session_events (session_id, event_data, created_at) VALUES (?1, ?2, ?3)",
17✔
154
            params![self.session_id.0.to_string(), event_data, now()],
17✔
155
        )?;
17✔
156

157
        Ok(())
17✔
158
    }
17✔
159

160
    fn load(
6✔
161
        &self,
6✔
162
    ) -> std::result::Result<
6✔
163
        Box<dyn Iterator<Item = ReceiverSessionEvent>>,
6✔
164
        Self::InternalStorageError,
6✔
165
    > {
6✔
166
        let conn = self.db.get_connection()?;
6✔
167
        let mut stmt = conn.prepare(
6✔
168
            "SELECT event_data FROM receive_session_events WHERE session_id = ?1 ORDER BY id ASC",
6✔
169
        )?;
6✔
170

171
        let event_rows = stmt.query_map(params![self.session_id.0.to_string()], |row| {
27✔
172
            let event_data: String = row.get(0)?;
27✔
173
            Ok(event_data)
27✔
174
        })?;
27✔
175

176
        let events: Vec<ReceiverSessionEvent> = event_rows
6✔
177
            .map(|row| {
27✔
178
                let event_data = row.expect("Failed to read event data from database");
27✔
179
                serde_json::from_str::<ReceiverSessionEvent>(&event_data)
27✔
180
                    .expect("Database corruption: failed to deserialize session event")
27✔
181
            })
27✔
182
            .collect();
6✔
183

184
        Ok(Box::new(events.into_iter()))
6✔
185
    }
6✔
186

187
    fn close(&self) -> std::result::Result<(), Self::InternalStorageError> {
3✔
188
        let conn = self.db.get_connection()?;
3✔
189

190
        conn.execute(
3✔
191
            "UPDATE receive_sessions SET completed_at = ?1 WHERE session_id = ?2",
3✔
192
            params![now(), self.session_id.0.to_string()],
3✔
193
        )?;
3✔
194

195
        Ok(())
3✔
196
    }
3✔
197
}
198

199
impl Database {
200
    pub(crate) fn get_recv_session_ids(&self) -> Result<Vec<SessionId>> {
5✔
201
        let conn = self.get_connection()?;
5✔
202
        let mut stmt =
5✔
203
            conn.prepare("SELECT session_id FROM receive_sessions WHERE completed_at IS NULL")?;
5✔
204

205
        let session_rows = stmt.query_map([], |row| {
5✔
206
            let session_id: String = row.get(0)?;
3✔
207
            let session_id = uuid::Uuid::parse_str(&session_id)
3✔
208
                .expect("Database corruption: invalid session_id UUID");
3✔
209
            Ok(SessionId(session_id))
3✔
210
        })?;
3✔
211

212
        let mut session_ids = Vec::new();
5✔
213
        for session_row in session_rows {
5✔
214
            let session_id = session_row?;
3✔
215
            session_ids.push(session_id);
3✔
216
        }
217

218
        Ok(session_ids)
5✔
219
    }
5✔
220

221
    pub(crate) fn get_send_session_ids(&self) -> Result<Vec<SessionId>> {
5✔
222
        let conn = self.get_connection()?;
5✔
223
        let mut stmt =
5✔
224
            conn.prepare("SELECT session_id FROM send_sessions WHERE completed_at IS NULL")?;
5✔
225

226
        let session_rows = stmt.query_map([], |row| {
5✔
227
            let session_id: String = row.get(0)?;
×
228
            let session_id = uuid::Uuid::parse_str(&session_id)
×
229
                .expect("Database corruption: invalid session_id UUID");
×
230
            Ok(SessionId(session_id))
×
231
        })?;
×
232

233
        let mut session_ids = Vec::new();
5✔
234
        for session_row in session_rows {
5✔
235
            let session_id = session_row?;
×
236
            session_ids.push(session_id);
×
237
        }
238

239
        Ok(session_ids)
5✔
240
    }
5✔
241

242
    pub(crate) fn get_send_session_id_by_receiver_pk(
4✔
243
        &self,
4✔
244
        receiver_pubkey: &HpkePublicKey,
4✔
245
    ) -> Result<Option<SessionId>> {
4✔
246
        let conn = self.get_connection()?;
4✔
247
        let receiver_pubkey_bytes = receiver_pubkey.to_compressed_bytes();
4✔
248
        let mut stmt = conn.prepare(
4✔
249
            "SELECT session_id FROM send_sessions WHERE receiver_pubkey = ?1 AND completed_at IS NULL",
4✔
250
        )?;
4✔
251
        let result = stmt.query_row(params![&receiver_pubkey_bytes], |row| {
4✔
252
            let session_id: String = row.get(0)?;
1✔
253
            let session_id = uuid::Uuid::parse_str(&session_id)
1✔
254
                .expect("Database corruption: invalid session_id UUID");
1✔
255
            Ok(SessionId(session_id))
1✔
256
        });
1✔
257
        Ok(result.optional()?)
4✔
258
    }
4✔
259

260
    pub(crate) fn get_inactive_send_session_ids(&self) -> Result<Vec<(SessionId, u64)>> {
×
261
        let conn = self.get_connection()?;
×
262
        let mut stmt = conn.prepare(
×
263
            "SELECT session_id, completed_at FROM send_sessions WHERE completed_at IS NOT NULL",
×
264
        )?;
×
265
        let session_rows = stmt.query_map([], |row| {
×
266
            let session_id: String = row.get(0)?;
×
267
            let session_id = uuid::Uuid::parse_str(&session_id)
×
268
                .expect("Database corruption: invalid session_id UUID");
×
269
            let completed_at: u64 = row.get(1)?;
×
270
            Ok((SessionId(session_id), completed_at))
×
271
        })?;
×
272

273
        let mut session_ids = Vec::new();
×
274
        for session_row in session_rows {
×
275
            let (session_id, completed_at) = session_row?;
×
276
            session_ids.push((session_id, completed_at));
×
277
        }
278
        Ok(session_ids)
×
279
    }
×
280

281
    pub(crate) fn get_inactive_recv_session_ids(&self) -> Result<Vec<(SessionId, u64)>> {
×
282
        let conn = self.get_connection()?;
×
283
        let mut stmt = conn.prepare(
×
284
            "SELECT session_id, completed_at FROM receive_sessions WHERE completed_at IS NOT NULL",
×
285
        )?;
×
286
        let session_rows = stmt.query_map([], |row| {
×
287
            let session_id: String = row.get(0)?;
×
288
            let session_id = uuid::Uuid::parse_str(&session_id)
×
289
                .expect("Database corruption: invalid session_id UUID");
×
290
            let completed_at: u64 = row.get(1)?;
×
291
            Ok((SessionId(session_id), completed_at))
×
292
        })?;
×
293

294
        let mut session_ids = Vec::new();
×
295
        for session_row in session_rows {
×
296
            let (session_id, completed_at) = session_row?;
×
297
            session_ids.push((session_id, completed_at));
×
298
        }
299
        Ok(session_ids)
×
300
    }
×
301

302
    /// Look up a sender session by ID regardless of active/inactive state.
303
    pub(crate) fn send_session_exists(&self, session_id: &SessionId) -> Result<bool> {
4✔
304
        let conn = self.get_connection()?;
4✔
305
        let exists: bool = conn.query_row(
4✔
306
            "SELECT EXISTS(SELECT 1 FROM send_sessions WHERE session_id = ?1)",
4✔
307
            params![session_id.0.to_string()],
4✔
308
            |row| row.get(0),
4✔
309
        )?;
×
310
        Ok(exists)
4✔
311
    }
4✔
312

313
    /// Look up a receiver session by ID regardless of active/inactive state.
314
    pub(crate) fn recv_session_exists(&self, session_id: &SessionId) -> Result<bool> {
4✔
315
        let conn = self.get_connection()?;
4✔
316
        let exists: bool = conn.query_row(
4✔
317
            "SELECT EXISTS(SELECT 1 FROM receive_sessions WHERE session_id = ?1)",
4✔
318
            params![session_id.0.to_string()],
4✔
319
            |row| row.get(0),
4✔
320
        )?;
×
321
        Ok(exists)
4✔
322
    }
4✔
323
}
324

325
#[cfg(all(test, feature = "v2"))]
326
mod tests {
327
    use std::sync::Arc;
328

329
    use payjoin::HpkeKeyPair;
330

331
    use super::*;
332

333
    fn create_test_db() -> Arc<Database> {
3✔
334
        // Use an in-memory database for tests
335
        let manager = r2d2_sqlite::SqliteConnectionManager::memory()
3✔
336
            .with_init(|conn| conn.execute_batch("PRAGMA locking_mode = EXCLUSIVE;"));
30✔
337
        let pool = r2d2::Pool::new(manager).expect("pool creation should succeed");
3✔
338
        let conn = pool.get().expect("connection should succeed");
3✔
339
        Database::init_schema(&conn).expect("schema init should succeed");
3✔
340
        Arc::new(Database(pool))
3✔
341
    }
3✔
342

343
    fn make_receiver_pubkey() -> payjoin::HpkePublicKey { HpkeKeyPair::gen_keypair().1 }
5✔
344

345
    // Second call with the same URI (same active session) should return DuplicateSendSession(Uri).
346
    #[test]
347
    fn test_duplicate_uri_returns_error() {
1✔
348
        let db = create_test_db();
1✔
349
        let rk1 = make_receiver_pubkey();
1✔
350
        let rk2 = make_receiver_pubkey();
1✔
351
        let uri = "bitcoin:addr1?pj=https://example.com/BBBBBBBB";
1✔
352

353
        SenderPersister::new(db.clone(), uri, &rk1).expect("first session should succeed");
1✔
354

355
        let err = SenderPersister::new(db, uri, &rk2).expect_err("duplicate URI should fail");
1✔
356
        assert!(
1✔
357
            matches!(err, Error::DuplicateSendSession(DuplicateKind::Uri)),
1✔
358
            "expected DuplicateSendSession(Uri), got: {err:?}"
359
        );
360
    }
1✔
361

362
    // Same receiver pubkey under a different URI should return DuplicateSendSession(ReceiverPubkey).
363
    #[test]
364
    fn test_duplicate_rk_returns_error() {
1✔
365
        let db = create_test_db();
1✔
366
        let rk = make_receiver_pubkey();
1✔
367
        let uri1 = "bitcoin:addr1?pj=https://example.com/CCCCCCCC";
1✔
368
        let uri2 = "bitcoin:addr1?pj=https://example.com/DDDDDDDD";
1✔
369

370
        SenderPersister::new(db.clone(), uri1, &rk).expect("first session should succeed");
1✔
371

372
        let err = SenderPersister::new(db, uri2, &rk).expect_err("duplicate RK should fail");
1✔
373
        assert!(
1✔
374
            matches!(err, Error::DuplicateSendSession(DuplicateKind::ReceiverPubkey)),
1✔
375
            "expected DuplicateSendSession(ReceiverPubkey), got: {err:?}"
376
        );
377
    }
1✔
378

379
    // After a session is marked completed, a new session with the same URI must still be rejected
380
    // to prevent address reuse, HPKE receiver-key reuse
381
    #[test]
382
    fn test_completed_session_blocks_reuse() {
1✔
383
        let db = create_test_db();
1✔
384
        let rk1 = make_receiver_pubkey();
1✔
385
        let rk2 = make_receiver_pubkey();
1✔
386
        let uri = "bitcoin:addr1?pj=https://example.com/EEEEEEEE";
1✔
387

388
        let persister =
1✔
389
            SenderPersister::new(db.clone(), uri, &rk1).expect("first session should succeed");
1✔
390

391
        // Mark the session as completed
392
        use payjoin::persist::SessionPersister;
393
        persister.close().expect("close should succeed");
1✔
394

395
        // A new session with the same URI must be rejected even after completion
396
        let err = SenderPersister::new(db, uri, &rk2)
1✔
397
            .expect_err("reuse of a completed session URI must be rejected");
1✔
398
        assert!(
1✔
399
            matches!(err, Error::DuplicateSendSession(DuplicateKind::Uri)),
1✔
400
            "expected DuplicateSendSession(Uri), got: {err:?}"
401
        );
402
    }
1✔
403
}
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