vector_core/db/
wrappers.rs1use nostr_sdk::prelude::{EventId, Timestamp};
4
5pub const TRANSPORT_NIP17: i64 = 0;
8pub const TRANSPORT_CONCORD: i64 = 1;
9
10pub fn save_processed_wrapper(wrapper_id_bytes: &[u8; 32], wrapper_created_at: u64, transport: i64) -> Result<(), String> {
13 let conn = super::get_write_connection_guard_static()?;
14 conn.execute(
15 "INSERT OR IGNORE INTO processed_wrappers (wrapper_id, wrapper_created_at, transport) VALUES (?1, ?2, ?3)",
16 rusqlite::params![&wrapper_id_bytes[..], wrapper_created_at as i64, transport],
17 ).map_err(|e| format!("Failed to save processed wrapper: {}", e))?;
18 Ok(())
19}
20
21pub fn processed_wrapper_exists(wrapper_id_bytes: &[u8; 32]) -> bool {
25 let conn = match super::get_db_connection_guard_static() {
26 Ok(c) => c,
27 Err(_) => return false,
28 };
29 conn.query_row(
30 "SELECT EXISTS(SELECT 1 FROM processed_wrappers WHERE wrapper_id = ?1)",
31 rusqlite::params![&wrapper_id_bytes[..]],
32 |row| row.get(0),
33 ).unwrap_or(false)
34}
35
36pub fn update_wrapper_timestamp(wrapper_id_bytes: &[u8; 32], wrapper_created_at: u64) -> Result<(), String> {
41 let conn = super::get_write_connection_guard_static()?;
42 conn.execute(
43 "UPDATE processed_wrappers SET wrapper_created_at = ?2 \
44 WHERE wrapper_id = ?1 AND wrapper_created_at = 0",
45 rusqlite::params![&wrapper_id_bytes[..], wrapper_created_at as i64],
46 ).map_err(|e| format!("Failed to backfill wrapper timestamp: {}", e))?;
47 Ok(())
48}
49
50pub fn load_processed_wrappers() -> Result<Vec<[u8; 32]>, String> {
52 let conn = match super::get_db_connection_guard_static() {
53 Ok(c) => c,
54 Err(_) => return Ok(Vec::new()),
55 };
56 let mut stmt = conn.prepare("SELECT wrapper_id FROM processed_wrappers WHERE transport = 0")
59 .map_err(|e| format!("Failed to prepare processed_wrappers query: {}", e))?;
60 let rows = stmt.query_map([], |row| {
61 let blob: Vec<u8> = row.get(0)?;
62 if blob.len() == 32 {
63 let mut arr = [0u8; 32];
64 arr.copy_from_slice(&blob);
65 Ok(arr)
66 } else {
67 Err(rusqlite::Error::InvalidParameterCount(blob.len(), 32))
68 }
69 }).map_err(|e| format!("Failed to query processed_wrappers: {}", e))?;
70
71 Ok(rows.flatten().collect())
72}
73
74pub fn load_processed_wrappers_since(since_secs: u64) -> Result<Vec<[u8; 32]>, String> {
78 let conn = match super::get_db_connection_guard_static() {
79 Ok(c) => c,
80 Err(_) => return Ok(Vec::new()),
81 };
82 let mut stmt = conn.prepare(
83 "SELECT wrapper_id FROM processed_wrappers WHERE transport = 0 AND wrapper_created_at >= ?1",
84 ).map_err(|e| format!("Failed to prepare processed_wrappers query: {}", e))?;
85 let rows = stmt.query_map(rusqlite::params![since_secs as i64], |row| {
86 let blob: Vec<u8> = row.get(0)?;
87 if blob.len() == 32 {
88 let mut arr = [0u8; 32];
89 arr.copy_from_slice(&blob);
90 Ok(arr)
91 } else {
92 Err(rusqlite::Error::InvalidParameterCount(blob.len(), 32))
93 }
94 }).map_err(|e| format!("Failed to query processed_wrappers: {}", e))?;
95
96 Ok(rows.flatten().collect())
97}
98
99pub fn load_recent_wrapper_ids(days: u64) -> Result<Vec<[u8; 32]>, String> {
101 let conn = match super::get_db_connection_guard_static() {
102 Ok(c) => c,
103 Err(_) => return Ok(Vec::new()),
104 };
105
106 let cutoff_secs = std::time::SystemTime::now()
107 .duration_since(std::time::UNIX_EPOCH).unwrap()
108 .as_secs()
109 .saturating_sub(days * 24 * 60 * 60);
110
111 let mut stmt = conn.prepare(
116 "SELECT e.wrapper_event_id FROM events e \
117 JOIN chats c ON e.chat_id = c.id \
118 WHERE e.wrapper_event_id IS NOT NULL AND e.wrapper_event_id != '' \
119 AND e.created_at >= ?1 AND c.chat_type != 2"
120 ).map_err(|e| format!("Failed to prepare wrapper_id query: {}", e))?;
121
122 let hex_ids: Vec<String> = stmt.query_map(rusqlite::params![cutoff_secs as i64], |row| {
123 row.get::<_, String>(0)
124 }).map_err(|e| format!("Failed to query wrapper_ids: {}", e))?
125 .flatten().collect();
126
127 let mut result = Vec::with_capacity(hex_ids.len());
128 for hex in hex_ids {
129 if hex.len() == 64 {
130 result.push(crate::simd::hex::hex_to_bytes_32(&hex));
131 }
132 }
133 Ok(result)
134}
135
136pub fn load_negentropy_items() -> Result<Vec<(EventId, Timestamp)>, String> {
138 let conn = super::get_db_connection_guard_static()
139 .map_err(|_| "No DB connection".to_string())?;
140
141 let mut stmt = conn.prepare(
144 "SELECT wrapper_id, wrapper_created_at FROM processed_wrappers WHERE transport = 0"
145 ).map_err(|e| format!("Failed to prepare negentropy query: {}", e))?;
146
147 let items: Vec<_> = stmt.query_map([], |row| {
148 let blob: Vec<u8> = row.get(0)?;
149 let created_at: i64 = row.get(1)?;
150 Ok((blob, created_at))
151 }).map_err(|e| format!("Failed to query processed_wrappers: {}", e))?
152 .flatten()
153 .filter_map(|(blob, ts)| {
154 if blob.len() == 32 {
155 let mut arr = [0u8; 32];
156 arr.copy_from_slice(&blob);
157 Some((
158 EventId::from_byte_array(arr),
159 Timestamp::from_secs(ts as u64),
160 ))
161 } else {
162 None
163 }
164 })
165 .collect();
166
167 Ok(items)
168}