1use std::{
2 collections::BTreeSet,
3 fs::read_dir,
4 path::Path,
5 sync::{Arc, OnceLock, atomic::AtomicU32},
6};
7
8use rocksdb::{ColumnFamilyDescriptor, Options, WriteOptions};
9use tuwunel_core::{
10 Result, debug, debug_warn, err, error, implement, info, itertools::Itertools, trace, warn,
11};
12
13use super::{
14 Db, Engine, backup::restore, cf_opts::cf_options, context, db_opts::db_options, descriptor,
15 descriptor::Descriptor, repair::repair,
16};
17use crate::{Context, or_else};
18
19#[implement(Engine)]
20#[tracing::instrument(skip_all)]
21pub(crate) async fn open(ctx: Arc<Context>, desc: &[Descriptor]) -> Result<Arc<Self>> {
22 let server = &ctx.server;
23 let config = &server.config;
24 let path = &config.database_path;
25
26 context::before_open(&ctx, path)?;
27
28 if let Some(backup_id) = config.database_restore_backup {
29 match server.claim_backup_restore() {
30 | true => restore(&ctx, backup_id)?,
31 | false => {
32 info!(%backup_id, "Restore already claimed by this process; not restoring again");
33 },
34 }
35 }
36
37 let db_opts = db_options(
38 config,
39 &ctx.env.lock().expect("environment locked"),
40 &ctx.row_cache.lock().expect("row cache locked"),
41 )?;
42
43 let (cfds, dropped) = Self::configure_cfds(&ctx, &db_opts, desc)?;
44 let num_cfds = cfds.len();
45 debug!("Configured {num_cfds} column descriptors...");
46
47 let load_time = std::time::Instant::now();
48 if config.rocksdb_repair {
49 repair(&db_opts, &config.database_path)?;
50 }
51
52 debug!("Opening database...");
53 let db = if config.rocksdb_read_only {
54 Db::open_cf_descriptors_read_only(&db_opts, path, cfds, false)
55 } else if config.rocksdb_secondary {
56 Db::open_cf_descriptors_as_secondary(&db_opts, path, path, cfds)
57 } else {
58 Db::open_cf_descriptors(&db_opts, path, cfds)
59 }
60 .or_else(or_else)?;
61
62 if !config.rocksdb_read_only && !config.rocksdb_secondary {
63 drop_columns(dropped.iter().map(String::as_str), |name| {
64 db.drop_cf(name).or_else(or_else)
65 })?;
66 }
67
68 info!(
69 columns = num_cfds,
70 sequence = %db.latest_sequence_number(),
71 time = ?load_time.elapsed(),
72 "Opened database."
73 );
74
75 Ok(Arc::new(Self {
76 db,
77 pool: ctx.pool.clone(),
78 ctx: ctx.clone(),
79 read_only: config.rocksdb_read_only,
80 secondary: config.rocksdb_secondary,
81 checksums: config.rocksdb_checksums,
82 write_options: WriteOptions::default(),
83 cf_index: OnceLock::new(),
84 corks: AtomicU32::new(0),
85 }))
86}
87
88pub(super) fn drop_columns<N>(
89 names: impl IntoIterator<Item = N>,
90 mut drop: impl FnMut(&str) -> Result,
91) -> Result
92where
93 N: AsRef<str>,
94{
95 let failed = names
96 .into_iter()
97 .filter(|name| {
98 let name = name.as_ref();
99 debug!(%name, "Deleting dropped database column");
100
101 drop(name)
102 .inspect_err(|error| {
103 error!(%name, ?error, "Failed to delete dropped database column");
104 })
105 .is_err()
106 })
107 .collect::<Vec<_>>();
108
109 failed.is_empty().then_some(()).ok_or_else(|| {
110 err!(Database(
111 "Failed to delete {} dropped database columns: {}",
112 failed.len(),
113 failed
114 .iter()
115 .map(|name: &N| name.as_ref())
116 .format(", ")
117 ))
118 })
119}
120
121#[implement(Engine)]
122#[tracing::instrument(name = "configure", skip_all)]
123fn configure_cfds(
124 ctx: &Arc<Context>,
125 db_opts: &Options,
126 desc: &[Descriptor],
127) -> Result<(Vec<ColumnFamilyDescriptor>, Vec<String>)> {
128 let server = &ctx.server;
129 let config = &server.config;
130 let path = &config.database_path;
131 let existing = Self::discover_cfs(path, db_opts)?;
132
133 let missing = existing
135 .iter()
136 .map(String::as_str)
137 .filter(|&name| name != "default")
138 .filter(|&name| !desc.iter().any(|desc| desc.name == name));
139
140 let creating = desc
142 .iter()
143 .filter(|desc| !desc.dropped)
144 .filter(|desc| !existing.contains(desc.name));
145
146 let dropping = desc
148 .iter()
149 .filter(|desc| desc.dropped)
150 .filter(|desc| existing.contains(desc.name))
151 .filter(|_| !config.rocksdb_never_drop_columns);
152
153 let dropped = desc
155 .iter()
156 .filter(|desc| desc.dropped)
157 .filter(|desc| !existing.contains(desc.name));
158
159 debug!(
160 existing = existing.len(),
161 described = desc.len(),
162 missing = missing.clone().count(),
163 dropped = dropped.clone().count(),
164 creating = creating.clone().count(),
165 dropping = dropping.clone().count(),
166 "Discovered database columns"
167 );
168
169 missing.clone().for_each(|name| {
170 debug_warn!("Found undescribed column {name:?} in existing database.");
171 });
172
173 dropped
174 .clone()
175 .map(|desc| desc.name)
176 .for_each(|name| {
177 debug!("Previously dropped column {name:?} no longer found in database.");
178 });
179
180 creating
181 .clone()
182 .map(|desc| desc.name)
183 .for_each(|name| {
184 debug!("Creating new column {name:?} not previously found in existing database.");
185 });
186
187 dropping
188 .clone()
189 .map(|desc| desc.name)
190 .for_each(|name| {
191 warn!(
192 "Column {name:?} has been scheduled for deletion. Storage may not appear \
193 reclaimed until further restart or compaction."
194 );
195 });
196
197 let not_dropped = |desc: &&Descriptor| {
198 !dropped
199 .clone()
200 .any(|dropped| desc.name == dropped.name)
201 };
202
203 let read_only = config.rocksdb_read_only || config.rocksdb_secondary;
205 let openable = |desc: &&Descriptor| !read_only || existing.contains(desc.name);
206
207 let dropping = dropping
208 .map(|desc| desc.name)
209 .map(ToOwned::to_owned)
210 .collect();
211
212 let cfnames = desc
213 .iter()
214 .filter(not_dropped)
215 .filter(openable)
216 .map(|desc| desc.name)
217 .chain(missing.clone());
218
219 let cfds: Vec<_> = desc
220 .iter()
221 .filter(not_dropped)
222 .filter(openable)
223 .copied()
224 .chain(missing.map(|_| descriptor::IGNORED))
225 .zip(cfnames)
226 .inspect(|&(desc, name)| {
227 assert!(
228 desc.ignored || desc.name == name,
229 "{name:?} does not match descriptor {:?}",
230 desc.name
231 );
232 })
233 .inspect(|&(_, name)| debug!(name, "Described column"))
234 .map(|(desc, name)| (desc, name.to_owned()))
235 .map(|(desc, name)| Ok((name, cf_options(ctx, db_opts.clone(), &desc)?)))
236 .map_ok(|(name, opts)| ColumnFamilyDescriptor::new(name, opts))
237 .collect::<Result<_>>()?;
238
239 trace!(?dropping);
240 Ok((cfds, dropping))
241}
242
243#[implement(Engine)]
244#[tracing::instrument(name = "discover", skip_all)]
245fn discover_cfs(path: &Path, opts: &Options) -> Result<BTreeSet<String>> {
246 Db::list_cf(opts, path)
247 .map(|cfs| cfs.into_iter().collect())
248 .or_else(|e| {
249 let remnants = count_remnants(path);
250
251 remnants
252 .eq(&0)
253 .then(BTreeSet::new)
254 .ok_or_else(|| {
255 err!(Database(
256 "Found {remnants} database files in {path:?} but no readable manifest: \
257 {e}. Refusing to initialize a new database over existing data; restore \
258 a complete database copy including CURRENT and MANIFEST, or remove the \
259 remnants to start fresh."
260 ))
261 })
262 })
263}
264
265fn count_remnants(path: &Path) -> usize {
266 read_dir(path)
267 .into_iter()
268 .flatten()
269 .filter_map(Result::ok)
270 .filter(|entry| entry.file_name().to_str().is_some_and(is_remnant))
271 .count()
272}
273
274pub(super) fn is_remnant(name: &str) -> bool {
276 let numbered = |suffix| {
277 name.strip_suffix(suffix)
278 .is_some_and(|stem| !stem.is_empty() && stem.bytes().all(|b| b.is_ascii_digit()))
279 };
280
281 name == "CURRENT" || name.starts_with("MANIFEST-") || numbered(".sst") || numbered(".log")
282}