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, 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 for name in &dropped {
64 debug!("Deleting dropped column {name:?} ...");
65 db.drop_cf(name).or_else(or_else)?;
66 }
67 }
68
69 info!(
70 columns = num_cfds,
71 sequence = %db.latest_sequence_number(),
72 time = ?load_time.elapsed(),
73 "Opened database."
74 );
75
76 Ok(Arc::new(Self {
77 db,
78 pool: ctx.pool.clone(),
79 ctx: ctx.clone(),
80 read_only: config.rocksdb_read_only,
81 secondary: config.rocksdb_secondary,
82 checksums: config.rocksdb_checksums,
83 write_options: WriteOptions::default(),
84 cf_index: OnceLock::new(),
85 corks: AtomicU32::new(0),
86 }))
87}
88
89#[implement(Engine)]
90#[tracing::instrument(name = "configure", skip_all)]
91fn configure_cfds(
92 ctx: &Arc<Context>,
93 db_opts: &Options,
94 desc: &[Descriptor],
95) -> Result<(Vec<ColumnFamilyDescriptor>, Vec<String>)> {
96 let server = &ctx.server;
97 let config = &server.config;
98 let path = &config.database_path;
99 let existing = Self::discover_cfs(path, db_opts)?;
100
101 let missing = existing
103 .iter()
104 .map(String::as_str)
105 .filter(|&name| name != "default")
106 .filter(|&name| !desc.iter().any(|desc| desc.name == name));
107
108 let creating = desc
110 .iter()
111 .filter(|desc| !desc.dropped)
112 .filter(|desc| !existing.contains(desc.name));
113
114 let dropping = desc
116 .iter()
117 .filter(|desc| desc.dropped)
118 .filter(|desc| existing.contains(desc.name))
119 .filter(|_| !config.rocksdb_never_drop_columns);
120
121 let dropped = desc
123 .iter()
124 .filter(|desc| desc.dropped)
125 .filter(|desc| !existing.contains(desc.name));
126
127 debug!(
128 existing = existing.len(),
129 described = desc.len(),
130 missing = missing.clone().count(),
131 dropped = dropped.clone().count(),
132 creating = creating.clone().count(),
133 dropping = dropping.clone().count(),
134 "Discovered database columns"
135 );
136
137 missing.clone().for_each(|name| {
138 debug_warn!("Found undescribed column {name:?} in existing database.");
139 });
140
141 dropped
142 .clone()
143 .map(|desc| desc.name)
144 .for_each(|name| {
145 debug!("Previously dropped column {name:?} no longer found in database.");
146 });
147
148 creating
149 .clone()
150 .map(|desc| desc.name)
151 .for_each(|name| {
152 debug!("Creating new column {name:?} not previously found in existing database.");
153 });
154
155 dropping
156 .clone()
157 .map(|desc| desc.name)
158 .for_each(|name| {
159 warn!(
160 "Column {name:?} has been scheduled for deletion. Storage may not appear \
161 reclaimed until further restart or compaction."
162 );
163 });
164
165 let not_dropped = |desc: &&Descriptor| {
166 !dropped
167 .clone()
168 .any(|dropped| desc.name == dropped.name)
169 };
170
171 let read_only = config.rocksdb_read_only || config.rocksdb_secondary;
173 let openable = |desc: &&Descriptor| !read_only || existing.contains(desc.name);
174
175 let dropping = dropping
176 .map(|desc| desc.name)
177 .map(ToOwned::to_owned)
178 .collect();
179
180 let cfnames = desc
181 .iter()
182 .filter(not_dropped)
183 .filter(openable)
184 .map(|desc| desc.name)
185 .chain(missing.clone());
186
187 let cfds: Vec<_> = desc
188 .iter()
189 .filter(not_dropped)
190 .filter(openable)
191 .copied()
192 .chain(missing.map(|_| descriptor::IGNORED))
193 .zip(cfnames)
194 .inspect(|&(desc, name)| {
195 assert!(
196 desc.ignored || desc.name == name,
197 "{name:?} does not match descriptor {:?}",
198 desc.name
199 );
200 })
201 .inspect(|&(_, name)| debug!(name, "Described column"))
202 .map(|(desc, name)| (desc, name.to_owned()))
203 .map(|(desc, name)| Ok((name, cf_options(ctx, db_opts.clone(), &desc)?)))
204 .map_ok(|(name, opts)| ColumnFamilyDescriptor::new(name, opts))
205 .collect::<Result<_>>()?;
206
207 trace!(?dropping);
208 Ok((cfds, dropping))
209}
210
211#[implement(Engine)]
212#[tracing::instrument(name = "discover", skip_all)]
213fn discover_cfs(path: &Path, opts: &Options) -> Result<BTreeSet<String>> {
214 Db::list_cf(opts, path)
215 .map(|cfs| cfs.into_iter().collect())
216 .or_else(|e| {
217 let remnants = count_remnants(path);
218
219 remnants
220 .eq(&0)
221 .then(BTreeSet::new)
222 .ok_or_else(|| {
223 err!(Database(
224 "Found {remnants} database files in {path:?} but no readable manifest: \
225 {e}. Refusing to initialize a new database over existing data; restore \
226 a complete database copy including CURRENT and MANIFEST, or remove the \
227 remnants to start fresh."
228 ))
229 })
230 })
231}
232
233fn count_remnants(path: &Path) -> usize {
234 read_dir(path)
235 .into_iter()
236 .flatten()
237 .filter_map(Result::ok)
238 .filter(|entry| entry.file_name().to_str().is_some_and(is_remnant))
239 .count()
240}
241
242pub(super) fn is_remnant(name: &str) -> bool {
244 let numbered = |suffix| {
245 name.strip_suffix(suffix)
246 .is_some_and(|stem| !stem.is_empty() && stem.bytes().all(|b| b.is_ascii_digit()))
247 };
248
249 name == "CURRENT" || name.starts_with("MANIFEST-") || numbered(".sst") || numbered(".log")
250}