Skip to content

Commit 071c5bb

Browse files
committed
Prune stale payjoin sessions on DB open
Delete payjoin sessions older than 30 days when the payjoin database is accessed. Remove related event rows in the same cleanup pass.
1 parent f3e688c commit 071c5bb

1 file changed

Lines changed: 81 additions & 1 deletion

File tree

src/payjoin/db.rs

Lines changed: 81 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,14 +60,17 @@ impl From<Error> for payjoin::ImplementationError {
6060

6161
/// Default filename for the payjoin database
6262
pub const DB_FILENAME: &str = "payjoin.sqlite";
63+
const SESSION_RETENTION_SECS: i64 = 30 * 24 * 60 * 60;
6364

6465
pub fn open_payjoin_db(
6566
datadir: Option<PathBuf>,
6667
wallet_name: &str,
6768
) -> std::result::Result<Arc<Database>, BDKCliError> {
6869
let wallet_dir = prepare_home_dir(datadir)?.join(wallet_name);
6970
std::fs::create_dir_all(&wallet_dir).map_err(|e| BDKCliError::Generic(e.to_string()))?;
70-
Ok(Arc::new(Database::create(wallet_dir.join(DB_FILENAME))?))
71+
let db = Arc::new(Database::create(wallet_dir.join(DB_FILENAME))?);
72+
db.prune_expired_sessions()?;
73+
Ok(db)
7174
}
7275

7376
/// Returns the current Unix timestamp in seconds
@@ -162,6 +165,83 @@ impl Database {
162165
Ok(was_seen_before)
163166
}
164167

168+
/// Removes old completed sessions and stale incomplete sessions plus their event logs.
169+
pub fn prune_expired_sessions(&self) -> Result<()> {
170+
let cutoff = now() - SESSION_RETENTION_SECS;
171+
let mut conn = self.conn();
172+
let tx = conn.transaction()?;
173+
let stale_send_session_ids = {
174+
let mut stmt = tx.prepare(
175+
"SELECT session_id FROM send_sessions
176+
WHERE (completed_at IS NOT NULL AND completed_at < ?1)
177+
OR (
178+
completed_at IS NULL
179+
AND session_id IN (
180+
SELECT session_id FROM send_session_events
181+
GROUP BY session_id
182+
HAVING MAX(created_at) < ?1
183+
)
184+
)",
185+
)?;
186+
let rows = stmt.query_map(params![cutoff], |row| row.get::<_, i64>(0))?;
187+
let mut ids = Vec::new();
188+
for row in rows {
189+
ids.push(row?);
190+
}
191+
ids
192+
};
193+
let stale_receive_session_ids = {
194+
let mut stmt = tx.prepare(
195+
"SELECT session_id FROM receive_sessions
196+
WHERE (completed_at IS NOT NULL AND completed_at < ?1)
197+
OR (
198+
completed_at IS NULL
199+
AND session_id IN (
200+
SELECT session_id FROM receive_session_events
201+
GROUP BY session_id
202+
HAVING MAX(created_at) < ?1
203+
)
204+
)",
205+
)?;
206+
let rows = stmt.query_map(params![cutoff], |row| row.get::<_, i64>(0))?;
207+
let mut ids = Vec::new();
208+
for row in rows {
209+
ids.push(row?);
210+
}
211+
ids
212+
};
213+
let deleted_any =
214+
!stale_send_session_ids.is_empty() || !stale_receive_session_ids.is_empty();
215+
216+
for session_id in stale_send_session_ids {
217+
tx.execute(
218+
"DELETE FROM send_session_events WHERE session_id = ?1",
219+
params![session_id],
220+
)?;
221+
tx.execute(
222+
"DELETE FROM send_sessions WHERE session_id = ?1",
223+
params![session_id],
224+
)?;
225+
}
226+
227+
for session_id in stale_receive_session_ids {
228+
tx.execute(
229+
"DELETE FROM receive_session_events WHERE session_id = ?1",
230+
params![session_id],
231+
)?;
232+
tx.execute(
233+
"DELETE FROM receive_sessions WHERE session_id = ?1",
234+
params![session_id],
235+
)?;
236+
}
237+
238+
tx.commit()?;
239+
if deleted_any {
240+
conn.execute("VACUUM", [])?;
241+
}
242+
Ok(())
243+
}
244+
165245
/// Returns IDs of all active (incomplete) receive sessions
166246
pub fn get_recv_session_ids(&self) -> Result<Vec<SessionId>> {
167247
let conn = self.conn();

0 commit comments

Comments
 (0)