1use std::path::{Path, PathBuf};
12
13#[cfg(test)]
14use std::cell::Cell;
15
16use readcon_core::types::ConFrame;
17
18use crate::corpus::ConCorpus;
19use crate::error::{Error, Result};
20use crate::keys::{FrameKey, TrajId};
21use crate::select::Select;
22
23pub const DEFAULT_N_SHARDS: u32 = 64;
25
26const MANIFEST: &str = "shards.json";
28
29#[cfg(test)]
30thread_local! {
31 pub(crate) static FAIL_DEST_MAN_COPY: Cell<bool> = const { Cell::new(false) };
32}
33
34#[derive(Clone, Debug, serde::Serialize, serde::Deserialize)]
35pub struct ShardManifest {
36 pub n_shards: u32,
37 pub version: u32,
38}
39
40pub struct ShardedConCorpus {
42 root: PathBuf,
43 n_shards: u32,
44 shards: Vec<Option<ConCorpus>>,
46}
47
48impl ShardedConCorpus {
49 pub fn open_existing(root: impl AsRef<Path>) -> Result<Self> {
51 let root = root.as_ref();
52 if !root.join(MANIFEST).is_file() {
53 return Err(Error::Message(format!(
54 "missing shards.json: {}",
55 root.display()
56 )));
57 }
58 Self::open(root, 1)
59 }
60
61 pub fn open(root: impl AsRef<Path>, n_shards: u32) -> Result<Self> {
63 let root = root.as_ref().to_path_buf();
64 std::fs::create_dir_all(&root)?;
65 let manifest_path = root.join(MANIFEST);
66 let n_shards = if manifest_path.is_file() {
67 let s = std::fs::read_to_string(&manifest_path)?;
68 let m: ShardManifest = serde_json::from_str(&s)?;
69 m.n_shards
70 } else {
71 if n_shards == 0 {
72 return Err(Error::Message("n_shards must be >= 1".into()));
73 }
74 let m = ShardManifest {
75 n_shards,
76 version: 1,
77 };
78 std::fs::write(&manifest_path, serde_json::to_string_pretty(&m)?)?;
79 n_shards
80 };
81 let mut shards = Vec::with_capacity(n_shards as usize);
82 shards.resize_with(n_shards as usize, || None);
83 Ok(Self {
84 root,
85 n_shards,
86 shards,
87 })
88 }
89
90 pub fn n_shards(&self) -> u32 {
91 self.n_shards
92 }
93
94 pub fn root(&self) -> &Path {
95 &self.root
96 }
97
98 #[inline]
99 pub fn shard_for_traj(traj_id: TrajId, n_shards: u32) -> u32 {
100 (traj_id % u64::from(n_shards)) as u32
101 }
102
103 fn shard_path(&self, shard_id: u32) -> PathBuf {
104 self.root.join(format!("shard_{shard_id:04}"))
105 }
106
107 fn shard_has_data(&self, shard_id: u32) -> bool {
108 self.shard_path(shard_id).join("data.mdb").is_file()
109 }
110
111 fn release_shards(&mut self) {
112 for slot in &mut self.shards {
113 if let Some(c) = slot.take() {
114 c.close();
115 }
116 }
117 }
118
119 pub fn shard_mut(&mut self, shard_id: u32) -> Result<&ConCorpus> {
121 if shard_id >= self.n_shards {
122 return Err(Error::Message(format!(
123 "shard_id {shard_id} >= n_shards {}",
124 self.n_shards
125 )));
126 }
127 let i = shard_id as usize;
128 if self.shards[i].is_none() {
129 let p = self.shard_path(shard_id);
130 self.shards[i] = Some(ConCorpus::open(p)?);
131 }
132 Ok(self.shards[i].as_ref().unwrap())
133 }
134
135 pub fn open_shard_for_traj(
137 root: impl AsRef<Path>,
138 traj_id: TrajId,
139 ) -> Result<(u32, ConCorpus)> {
140 let root = root.as_ref();
141 let manifest_path = root.join(MANIFEST);
142 let n_shards = if manifest_path.is_file() {
143 let m: ShardManifest = serde_json::from_str(&std::fs::read_to_string(&manifest_path)?)?;
144 m.n_shards
145 } else {
146 let _ = Self::open(root, DEFAULT_N_SHARDS)?;
147 DEFAULT_N_SHARDS
148 };
149 let sid = Self::shard_for_traj(traj_id, n_shards);
150 let corpus = ConCorpus::open(root.join(format!("shard_{sid:04}")))?;
151 Ok((sid, corpus))
152 }
153
154 pub fn open_shard(root: impl AsRef<Path>, shard_id: u32) -> Result<ConCorpus> {
157 let root = root.as_ref();
158 let manifest_path = root.join(MANIFEST);
159 if !manifest_path.is_file() {
160 return Err(Error::Message(format!(
161 "missing shards.json: {}",
162 root.display()
163 )));
164 }
165 let m: ShardManifest = serde_json::from_str(&std::fs::read_to_string(&manifest_path)?)?;
166 let n_shards = m.n_shards;
167 if shard_id >= n_shards {
168 return Err(Error::Message(format!(
169 "shard_id {shard_id} >= n_shards {n_shards}"
170 )));
171 }
172 ConCorpus::open(root.join(format!("shard_{shard_id:04}")))
173 }
174
175 pub fn append_trajectory_path(
176 &mut self,
177 traj_id: TrajId,
178 file: impl AsRef<Path>,
179 ) -> Result<u32> {
180 self.append_trajectory_path_units(traj_id, file, None)
181 }
182
183 pub fn append_trajectory_path_units(
184 &mut self,
185 traj_id: TrajId,
186 file: impl AsRef<Path>,
187 units: Option<serde_json::Value>,
188 ) -> Result<u32> {
189 let sid = Self::shard_for_traj(traj_id, self.n_shards);
190 let c = self.shard_mut(sid)?;
191 c.append_trajectory_path_units(traj_id, file, units)
192 }
193
194 pub fn append_trajectory_str(
195 &mut self,
196 traj_id: TrajId,
197 contents: &str,
198 source: impl Into<String>,
199 ) -> Result<u32> {
200 let sid = Self::shard_for_traj(traj_id, self.n_shards);
201 let c = self.shard_mut(sid)?;
202 c.append_trajectory_str(traj_id, contents, source)
203 }
204
205 pub fn append_trajectory_frames(
206 &mut self,
207 traj_id: TrajId,
208 frames: &[ConFrame],
209 source: impl Into<String>,
210 ) -> Result<u32> {
211 let sid = Self::shard_for_traj(traj_id, self.n_shards);
212 let c = self.shard_mut(sid)?;
213 c.append_trajectory_frames(traj_id, frames, source)
214 }
215
216 pub fn select(&mut self, sel: &Select) -> Result<Vec<FrameKey>> {
218 self.release_shards();
219 let mut out = Vec::new();
220 for sid in 0..self.n_shards {
221 if !self.shard_has_data(sid) {
222 continue;
223 }
224 let c = ConCorpus::open_readonly(self.shard_path(sid))?;
225 out.extend(c.select(sel)?);
226 c.close();
227 }
228 out.sort();
229 if let Some(lim) = sel.limit {
230 out.truncate(lim);
231 }
232 Ok(out)
233 }
234
235 pub fn get_frame_text(&mut self, key: FrameKey) -> Result<String> {
236 self.release_shards();
237 let sid = Self::shard_for_traj(key.traj_id, self.n_shards);
238 if !self.shard_has_data(sid) {
239 return Err(Error::Message(format!(
240 "shard_{sid:04} is not a corpus directory"
241 )));
242 }
243 let c = ConCorpus::open_readonly(self.shard_path(sid))?;
244 let text = c.get_frame_text(key)?;
245 c.close();
246 Ok(text)
247 }
248
249 pub fn reindex_all(&mut self) -> Result<u32> {
250 let mut n = 0u32;
251 for sid in 0..self.n_shards {
252 if self.shard_has_data(sid) {
253 n += self.shard_mut(sid)?.reindex()?;
254 }
255 }
256 Ok(n)
257 }
258
259 pub fn drain_to(src: impl AsRef<Path>, dst: impl AsRef<Path>) -> Result<u32> {
263 let src = src.as_ref();
264 let dst = dst.as_ref();
265 let man = src.join(MANIFEST);
266 if !man.is_file() {
267 return Err(Error::Message("drain: missing shards.json".into()));
268 }
269 let dest_was_new = !dst.exists();
270 let dest_man = dst.join(MANIFEST);
271 let mut copied_manifest = false;
272 let mut created = Vec::new();
273 let written = (|| -> Result<u32> {
274 std::fs::create_dir_all(dst)?;
275 if dest_man.is_file() {
276 let existing: ShardManifest =
277 serde_json::from_str(&std::fs::read_to_string(&dest_man)?)?;
278 let incoming: ShardManifest =
279 serde_json::from_str(&std::fs::read_to_string(&man)?)?;
280 if existing.n_shards != incoming.n_shards {
281 return Err(Error::Message(
282 "drain: dest shards.json n_shards does not match src".into(),
283 ));
284 }
285 }
286 let m: ShardManifest = serde_json::from_str(&std::fs::read_to_string(&man)?)?;
287 for i in 0..m.n_shards {
288 let name = format!("shard_{i:04}");
289 if src.join(&name).join("data.mdb").is_file()
290 && dst.join(&name).join("data.mdb").is_file()
291 {
292 return Err(Error::Message(format!(
293 "drain: dest {name} exists; refuse overwrite. Drain each node to a unique dest, then join-drained."
294 )));
295 }
296 }
297 copied_manifest = !dest_man.is_file();
298 if copied_manifest {
299 #[cfg(test)]
300 if FAIL_DEST_MAN_COPY.with(|f| f.replace(false)) {
301 return Err(Error::Message("drain: dest_man copy".into()));
302 }
303 std::fs::copy(&man, &dest_man)?;
304 }
305 let mut n = 0u32;
306 for i in 0..m.n_shards {
307 let name = format!("shard_{i:04}");
308 let from = src.join(&name);
309 if !from.join("data.mdb").is_file() {
310 continue;
311 }
312 let to = dst.join(&name);
313 created.push(to.clone());
314 let ro = ConCorpus::open_readonly(&from)?;
315 ro.snapshot_to(&to)?;
316 ro.close();
317 n += 1;
318 }
319 Ok(n)
320 })();
321 if written.is_err() {
322 if dest_was_new {
323 let _ = std::fs::remove_dir_all(dst);
324 } else {
325 for p in &created {
326 let _ = std::fs::remove_dir_all(p);
327 }
328 if copied_manifest {
329 let _ = std::fs::remove_file(&dest_man);
330 }
331 }
332 }
333 written
334 }
335}
336
337#[cfg(test)]
338mod tests {
339 use super::*;
340 use std::sync::Arc;
341 use std::thread;
342
343 fn fixture(name: &str) -> PathBuf {
344 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
345 .join("resources/test")
346 .join(name)
347 }
348
349 #[test]
350 fn parallel_writers_different_shards() {
351 let dir = tempfile::tempdir().unwrap();
352 let root = dir.path().join("hpc");
353 let n_shards = 8u32;
355 ShardedConCorpus::open(&root, n_shards).unwrap();
356 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
357 let root = Arc::new(root);
358 let mut joins = Vec::new();
359 for sid in 0..n_shards {
360 let root = Arc::clone(&root);
361 let text = text.clone();
362 joins.push(thread::spawn(move || {
363 let db = ShardedConCorpus::open_shard(root.as_path(), sid).unwrap();
365 let traj = u64::from(sid); db.append_trajectory_str(traj, &text, format!("shard{sid}"))
367 .unwrap()
368 }));
369 }
370 let mut ns = Vec::new();
371 for j in joins {
372 ns.push(j.join().unwrap());
373 }
374 assert!(ns.iter().all(|&n| n >= 1));
375 let mut fan = ShardedConCorpus::open(root.as_path(), n_shards).unwrap();
376 let keys = fan.select(&Select::new().require_symbol("Cu")).unwrap();
377 drop(fan);
378 let drained = dir.path().join("pfs");
379 let ncopy = ShardedConCorpus::drain_to(root.as_path(), &drained).unwrap();
380 assert_eq!(ncopy, n_shards);
381 assert!(drained.join("shards.json").is_file());
382 assert!(drained.join("shard_0000").join("data.mdb").is_file());
383 assert!(!drained.join("shard_0000").join("lock.mdb").is_file());
384 let dest_sz = std::fs::metadata(drained.join("shard_0000").join("data.mdb"))
385 .unwrap()
386 .len();
387 assert!(
388 dest_sz < 64 * 1024 * 1024,
389 "compact snapshot must not materialize the 2 GiB map, got {dest_sz}"
390 );
391 assert!(ShardedConCorpus::drain_to(root.as_path(), &drained).is_err());
392 assert!(drained.join("shards.json").is_file());
393 assert_eq!(keys.len(), 8);
394 let joined = dir.path().join("joined");
395 let mut drained_root = ShardedConCorpus::open(&drained, n_shards).unwrap();
396 let njoin = drained_root.join_to_single_env(&joined).unwrap();
397 assert!(njoin >= n_shards);
398 let single = ConCorpus::open(&joined).unwrap();
399 let jk = single.select(&Select::new().require_symbol("Cu")).unwrap();
400 assert_eq!(jk.len(), keys.len());
401 }
402
403 #[test]
404 fn open_shard_for_traj_writes_manifest() {
405 let dir = tempfile::tempdir().unwrap();
406 let root = dir.path().join("fresh");
407 let (sid, db) = ShardedConCorpus::open_shard_for_traj(&root, 0).unwrap();
408 assert_eq!(sid, 0);
409 drop(db);
410 assert!(root.join("shards.json").is_file());
411 }
412
413 #[test]
414 fn join_drained_duplicate_traj_does_not_create_dest() {
415 let dir = tempfile::tempdir().unwrap();
416 let a = dir.path().join("a");
417 let b = dir.path().join("b");
418 ShardedConCorpus::open(&a, 1).unwrap();
419 ShardedConCorpus::open(&b, 1).unwrap();
420 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
421 ShardedConCorpus::open_shard(&a, 0)
422 .unwrap()
423 .append_trajectory_str(7, &text, "a")
424 .unwrap();
425 ShardedConCorpus::open_shard(&b, 0)
426 .unwrap()
427 .append_trajectory_str(7, &text, "b")
428 .unwrap();
429 let dest = dir.path().join("out");
430 let err = join_drained_roots(&[a, b], &dest).unwrap_err();
431 assert!(err.to_string().contains("traj_id"), "{err}");
432 assert!(!dest.exists());
433 }
434
435 #[test]
436 fn join_drained_roots_missing_source_errors() {
437 let dir = tempfile::tempdir().unwrap();
438 let missing = dir.path().join("nope");
439 let dest = dir.path().join("out");
440 let err = join_drained_roots(&[missing], &dest).unwrap_err();
441 assert!(err.to_string().contains("missing shards.json"), "{err}");
442 assert!(!dest.exists());
443 }
444
445 #[test]
446 fn drain_refuse_does_not_write_manifest() {
447 let dir = tempfile::tempdir().unwrap();
448 let root = dir.path().join("hpc");
449 ShardedConCorpus::open(&root, 2).unwrap();
450 let db = ShardedConCorpus::open_shard(&root, 0).unwrap();
451 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
452 db.append_trajectory_str(0, &text, "s0").unwrap();
453 drop(db);
454 let dest = dir.path().join("pfs");
455 std::fs::create_dir_all(dest.join("shard_0000")).unwrap();
456 std::fs::write(dest.join("shard_0000").join("data.mdb"), b"x").unwrap();
457 let err = ShardedConCorpus::drain_to(&root, &dest).unwrap_err();
458 assert!(err.to_string().contains("refuse overwrite"), "{err}");
459 assert!(
460 !dest.join("shards.json").is_file(),
461 "refuse must not leave dest shards.json"
462 );
463 }
464
465 #[test]
466 fn traj_routing_stable() {
467 assert_eq!(ShardedConCorpus::shard_for_traj(0, 64), 0);
468 assert_eq!(ShardedConCorpus::shard_for_traj(65, 64), 1);
469 }
470
471 #[test]
475 fn multi_shard_writers_select_matches_single_env_baseline() {
476 let dir = tempfile::tempdir().unwrap();
477 let root = dir.path().join("hpc_scale");
478 let baseline = dir.path().join("single_env");
479 let n_shards = 4u32;
480 ShardedConCorpus::open(&root, n_shards).unwrap();
481 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
482 let root_a = Arc::new(root.clone());
483 let mut joins = Vec::new();
484 for sid in 0..n_shards {
485 let root = Arc::clone(&root_a);
486 let text = text.clone();
487 joins.push(thread::spawn(move || {
488 let db = ShardedConCorpus::open_shard(root.as_path(), sid).unwrap();
489 let traj = u64::from(sid);
490 db.append_trajectory_str(traj, &text, format!("s{sid}"))
491 .unwrap()
492 }));
493 }
494 let mut frames_per_traj = Vec::new();
495 for j in joins {
496 frames_per_traj.push(j.join().unwrap());
497 }
498 assert!(frames_per_traj.iter().all(|&n| n >= 1));
499
500 let single = ConCorpus::open(&baseline).unwrap();
502 for sid in 0..n_shards {
503 let traj = u64::from(sid);
504 let n = single
505 .append_trajectory_str(traj, &text, format!("s{sid}"))
506 .unwrap();
507 assert_eq!(n, frames_per_traj[sid as usize]);
508 }
509
510 let mut fan = ShardedConCorpus::open(&root, n_shards).unwrap();
511 let sharded_keys = fan.select(&Select::new().require_symbol("Cu")).unwrap();
512 let base_keys = single.select(&Select::new().require_symbol("Cu")).unwrap();
513 assert_eq!(sharded_keys.len(), base_keys.len());
514 let mut sk: Vec<_> = sharded_keys
515 .iter()
516 .map(|k| (k.traj_id, k.frame_idx))
517 .collect();
518 let mut bk: Vec<_> = base_keys.iter().map(|k| (k.traj_id, k.frame_idx)).collect();
519 sk.sort_unstable();
520 bk.sort_unstable();
521 assert_eq!(sk, bk, "fan-out select must match single-env membership");
522
523 for (tid, fidx) in &sk {
525 let key = crate::keys::FrameKey {
526 traj_id: *tid,
527 frame_idx: *fidx,
528 };
529 let a = fan.get_frame_text(key).unwrap();
530 let b = single.get_frame_text(key).unwrap();
531 assert_eq!(a, b);
532 }
533 }
534}
535
536#[derive(Clone, Copy, Debug, PartialEq, Eq)]
538pub enum CorpusExportKind {
539 ShardedLmdb,
541 SingleEnvLmdb,
543 ExtXyz,
545}
546
547impl CorpusExportKind {
548 pub fn as_str(self) -> &'static str {
549 match self {
550 Self::ShardedLmdb => "sharded-lmdb",
551 Self::SingleEnvLmdb => "single-env-lmdb",
552 Self::ExtXyz => "extxyz",
553 }
554 }
555}
556
557impl ShardedConCorpus {
558 pub fn join_to_single_env(&mut self, dst: impl AsRef<Path>) -> Result<u32> {
563 let dst = dst.as_ref();
564 if dst.exists() {
565 return Err(Error::Message(format!(
566 "join dest exists: {}",
567 dst.display()
568 )));
569 }
570 self.release_shards();
571 let mut preview = std::collections::BTreeMap::new();
572 for sid in 0..self.n_shards {
573 if !self.shard_has_data(sid) {
574 continue;
575 }
576 let c = ConCorpus::open_readonly(self.shard_path(sid))?;
577 for fk in c.list_frame_keys()? {
578 preview_traj(&mut preview, fk.traj_id, u64::from(sid))?;
579 }
580 c.close();
581 }
582 let n = (|| -> Result<u32> {
583 let out = ConCorpus::open(dst)?;
584 let mut seen_traj = std::collections::BTreeSet::new();
585 let n = append_sharded_into(self, &out, &mut seen_traj)?;
586 out.close();
587 Ok(n)
588 })();
589 rollback_new_dest(dst, n)
590 }
591
592 pub fn export_extxyz(
594 &mut self,
595 sel: &Select,
596 out: impl AsRef<Path>,
597 energy_key: &str,
598 ) -> Result<u32> {
599 let joined = std::env::temp_dir().join(format!(
600 "readcon_db_join_{}_{}",
601 std::process::id(),
602 std::time::SystemTime::now()
603 .duration_since(std::time::UNIX_EPOCH)
604 .map(|d| d.as_nanos())
605 .unwrap_or(0)
606 ));
607 if joined.exists() {
608 return Err(Error::Message(
609 "export_extxyz: temp join dest exists".into(),
610 ));
611 }
612 struct RemoveOnDrop(PathBuf);
613 impl Drop for RemoveOnDrop {
614 fn drop(&mut self) {
615 let _ = std::fs::remove_dir_all(&self.0);
616 }
617 }
618 let _guard = RemoveOnDrop(joined.clone());
619 self.join_to_single_env(&joined)?;
620 let db = ConCorpus::open_readonly(&joined)?;
621 let keys = db.select(sel)?;
622 let n = db.export_extxyz(&keys, out, energy_key)?;
623 db.close();
624 Ok(n as u32)
625 }
626
627 pub fn split_single_to_sharded(
630 single: &ConCorpus,
631 dst_root: impl AsRef<Path>,
632 n_shards: u32,
633 ) -> Result<u32> {
634 if n_shards == 0 {
635 return Err(Error::Message("n_shards must be >= 1".into()));
636 }
637 let dst_root = dst_root.as_ref();
638 if dst_root.exists() {
639 return Err(Error::Message(format!(
640 "split dest exists: {}",
641 dst_root.display()
642 )));
643 }
644 let n = (|| -> Result<u32> {
645 let mut sharded = ShardedConCorpus::open(dst_root, n_shards)?;
646 let keys = single.list_frame_keys()?;
647 let mut by_traj: std::collections::BTreeMap<u64, Vec<FrameKey>> =
648 std::collections::BTreeMap::new();
649 for fk in keys {
650 by_traj.entry(fk.traj_id).or_default().push(fk);
651 }
652 let mut n = 0u32;
653 for (tid, mut fks) in by_traj {
654 fks.sort();
655 let mut concat = String::new();
656 for fk in &fks {
657 concat.push_str(&single.get_frame_text(*fk)?);
658 }
659 let nf = sharded.append_trajectory_str(tid, &concat, "split-from-single")?;
660 n += nf;
661 }
662 Ok(n)
663 })();
664 rollback_new_dest(dst_root, n)
665 }
666}
667
668fn append_sharded_into(
669 sh: &mut ShardedConCorpus,
670 out: &ConCorpus,
671 seen_traj: &mut std::collections::BTreeSet<u64>,
672) -> Result<u32> {
673 let mut n = 0u32;
674 for sid in 0..sh.n_shards {
675 if !sh.shard_has_data(sid) {
676 continue;
677 }
678 let shard = ConCorpus::open_readonly(sh.shard_path(sid))?;
679 let keys = shard.list_frame_keys()?;
680 let mut by_traj: std::collections::BTreeMap<u64, Vec<FrameKey>> =
681 std::collections::BTreeMap::new();
682 for fk in keys {
683 by_traj.entry(fk.traj_id).or_default().push(fk);
684 }
685 for (tid, mut fks) in by_traj {
686 if !seen_traj.insert(tid) {
687 return Err(Error::Message(format!(
688 "traj_id {tid} appears in multiple shards or join sources"
689 )));
690 }
691 fks.sort();
692 let mut concat = String::new();
693 for fk in &fks {
694 concat.push_str(&shard.get_frame_text(*fk)?);
695 }
696 n += out.append_trajectory_str(tid, &concat, format!("join-from-shard-{sid}"))?;
697 }
698 shard.close();
699 }
700 Ok(n)
701}
702
703pub fn join_drained_roots(sources: &[PathBuf], dst: impl AsRef<Path>) -> Result<u32> {
706 if sources.is_empty() {
707 return Err(Error::Message("join-drained: no sources".into()));
708 }
709 let dst = dst.as_ref();
710 if dst.exists() {
711 return Err(Error::Message(format!(
712 "join-drained dest exists: {}",
713 dst.display()
714 )));
715 }
716 for src in sources {
717 if !src.join(MANIFEST).is_file() {
718 return Err(Error::Message(format!(
719 "join-drained: missing shards.json: {}",
720 src.display()
721 )));
722 }
723 }
724 {
725 let mut preview = std::collections::BTreeMap::new();
726 for (si, src) in sources.iter().enumerate() {
727 let sh = ShardedConCorpus::open_existing(src)?;
728 for sid in 0..sh.n_shards {
729 if !sh.shard_has_data(sid) {
730 continue;
731 }
732 let shard = ConCorpus::open_readonly(sh.shard_path(sid))?;
733 let owner = ((si as u64) << 32) | u64::from(sid);
734 for fk in shard.list_frame_keys()? {
735 preview_traj(&mut preview, fk.traj_id, owner)?;
736 }
737 shard.close();
738 }
739 }
740 }
741 let n = (|| -> Result<u32> {
742 let out = ConCorpus::open(dst)?;
743 let mut n = 0u32;
744 let mut seen = std::collections::BTreeSet::new();
745 for src in sources {
746 let mut sh = ShardedConCorpus::open_existing(src)?;
747 n += append_sharded_into(&mut sh, &out, &mut seen)?;
748 }
749 out.close();
750 Ok(n)
751 })();
752 rollback_new_dest(dst, n)
753}
754
755pub fn join_corpus_dirs(sources: &[PathBuf], dst: impl AsRef<Path>) -> Result<u32> {
757 let dst = dst.as_ref();
758 if dst.exists() {
759 return Err(Error::Message(format!(
760 "join dest exists: {}",
761 dst.display()
762 )));
763 }
764 for src in sources {
765 if !src.join("data.mdb").is_file() {
766 return Err(Error::Message(format!(
767 "join: missing data.mdb: {}",
768 src.display()
769 )));
770 }
771 }
772 {
773 let mut preview = std::collections::BTreeMap::new();
774 for (si, src) in sources.iter().enumerate() {
775 let c = ConCorpus::open_readonly(src)?;
776 for fk in c.list_frame_keys()? {
777 preview_traj(&mut preview, fk.traj_id, si as u64)?;
778 }
779 c.close();
780 }
781 }
782 let n = (|| -> Result<u32> {
783 let out = ConCorpus::open(dst)?;
784 let mut n = 0u32;
785 let mut seen = std::collections::BTreeSet::new();
786 for src in sources {
787 let c = ConCorpus::open_readonly(src)?;
788 let keys = c.list_frame_keys()?;
789 let mut by_traj: std::collections::BTreeMap<u64, Vec<FrameKey>> =
790 std::collections::BTreeMap::new();
791 for fk in keys {
792 by_traj.entry(fk.traj_id).or_default().push(fk);
793 }
794 for (tid, mut fks) in by_traj {
795 if !seen.insert(tid) {
796 return Err(Error::Message(format!(
797 "duplicate traj_id {tid} across join sources"
798 )));
799 }
800 fks.sort();
801 let mut concat = String::new();
802 for fk in &fks {
803 concat.push_str(&c.get_frame_text(*fk)?);
804 }
805 n += out.append_trajectory_str(tid, &concat, src.display().to_string())?;
806 }
807 c.close();
808 }
809 out.close();
810 Ok(n)
811 })();
812 rollback_new_dest(dst, n)
813}
814
815fn rollback_new_dest<T>(dst: &Path, r: Result<T>) -> Result<T> {
816 if r.is_err() {
817 let _ = std::fs::remove_dir_all(dst);
818 }
819 r
820}
821
822fn preview_traj(
825 preview: &mut std::collections::BTreeMap<u64, u64>,
826 traj_id: u64,
827 owner: u64,
828) -> Result<()> {
829 match preview.entry(traj_id) {
830 std::collections::btree_map::Entry::Occupied(e) if *e.get() != owner => {
831 Err(Error::Message(format!(
832 "traj_id {traj_id} appears in multiple shards or join sources"
833 )))
834 }
835 std::collections::btree_map::Entry::Occupied(_) => Ok(()),
836 std::collections::btree_map::Entry::Vacant(e) => {
837 e.insert(owner);
838 Ok(())
839 }
840 }
841}
842
843pub fn open_single_env_for_export(src: impl AsRef<Path>) -> Result<ConCorpus> {
846 let src = src.as_ref();
847 if src.join(MANIFEST).is_file() {
848 return Err(Error::Message(
849 "compact-export-extxyz: sharded root needs --sharded".into(),
850 ));
851 }
852 ConCorpus::open_readonly(src)
853}
854
855#[cfg(test)]
856mod compaction_tests {
857 use super::*;
858 use crate::keys::FrameKey;
859 use crate::select::Select;
860
861 fn fixture(name: &str) -> PathBuf {
862 PathBuf::from(env!("CARGO_MANIFEST_DIR"))
863 .join("resources/test")
864 .join(name)
865 }
866
867 #[test]
868 fn join_split_reversible_membership() {
869 let dir = tempfile::tempdir().unwrap();
870 let sharded_root = dir.path().join("sharded");
871 let con_text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
872 {
873 let mut s = ShardedConCorpus::open(&sharded_root, 4).unwrap();
874 for tid in [0u64, 1, 2, 3] {
875 s.append_trajectory_str(tid, &con_text, "t").unwrap();
876 }
877 }
878 let mut s = ShardedConCorpus::open(&sharded_root, 4).unwrap();
879 let before = s.select(&Select::new()).unwrap();
880 assert_eq!(before.len(), 4);
881
882 let joined = dir.path().join("joined");
883 let n = s.join_to_single_env(&joined).unwrap();
884 assert_eq!(n, 4);
885 assert!(s.join_to_single_env(&joined).is_err());
886 let joined_c = ConCorpus::open(&joined).unwrap();
887 let mid = joined_c.select(&Select::new()).unwrap();
888 assert_eq!(mid.len(), 4);
889
890 let split_root = dir.path().join("split_again");
891 let n2 = ShardedConCorpus::split_single_to_sharded(&joined_c, &split_root, 4).unwrap();
892 assert_eq!(n2, 4);
893 let mut s2 = ShardedConCorpus::open(&split_root, 4).unwrap();
894 let after = s2.select(&Select::new()).unwrap();
895 assert_eq!(after.len(), before.len());
896 let mut bt: Vec<_> = before.iter().map(|k| k.traj_id).collect();
898 let mut at: Vec<_> = after.iter().map(|k| k.traj_id).collect();
899 bt.sort();
900 at.sort();
901 assert_eq!(bt, at);
902 }
903
904 #[test]
905 fn join_drained_roots_keeps_disjoint_trajs_same_shard() {
906 let dir = tempfile::tempdir().unwrap();
907 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
908 let n_shards = 2u32;
909 let node_a = dir.path().join("node_a");
910 let node_b = dir.path().join("node_b");
911 ShardedConCorpus::open(&node_a, n_shards).unwrap();
912 ShardedConCorpus::open(&node_b, n_shards).unwrap();
913 ShardedConCorpus::open_shard(&node_a, 0)
915 .unwrap()
916 .append_trajectory_str(0, &text, "a")
917 .unwrap();
918 ShardedConCorpus::open_shard(&node_b, 0)
919 .unwrap()
920 .append_trajectory_str(2, &text, "b")
921 .unwrap();
922 let dest_a = dir.path().join("dest_a");
923 let dest_b = dir.path().join("dest_b");
924 assert_eq!(ShardedConCorpus::drain_to(&node_a, &dest_a).unwrap(), 1);
925 assert_eq!(ShardedConCorpus::drain_to(&node_b, &dest_b).unwrap(), 1);
926 assert!(ShardedConCorpus::drain_to(&node_a, &dest_a).is_err());
927 let joined = dir.path().join("joined");
928 let n = join_drained_roots(&[dest_a.clone(), dest_b.clone()], &joined).unwrap();
929 assert!(n >= 2);
930 let db = ConCorpus::open(&joined).unwrap();
931 let keys = db.select(&Select::new()).unwrap();
932 let mut tids: Vec<u64> = keys.iter().map(|k| k.traj_id).collect();
933 tids.sort_unstable();
934 tids.dedup();
935 assert_eq!(tids, vec![0, 2]);
936 assert!(join_drained_roots(&[dest_a, dest_b], &joined).is_err());
937 }
938
939 #[test]
940 fn join_drained_roots_refuses_existing_dest() {
941 let dir = tempfile::tempdir().unwrap();
942 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
943 let node = dir.path().join("node");
944 ShardedConCorpus::open(&node, 2).unwrap();
945 ShardedConCorpus::open_shard(&node, 0)
946 .unwrap()
947 .append_trajectory_str(0, &text, "a")
948 .unwrap();
949 let dest = dir.path().join("dest");
950 assert_eq!(ShardedConCorpus::drain_to(&node, &dest).unwrap(), 1);
951 let joined = dir.path().join("joined");
952 join_drained_roots(&[dest.clone()], &joined).unwrap();
953 let err = join_drained_roots(&[dest], &joined).unwrap_err();
954 assert!(
955 err.to_string().contains("join-drained dest exists"),
956 "{err}"
957 );
958 }
959
960 #[test]
961 fn export_kinds_documented() {
962 assert_eq!(CorpusExportKind::ShardedLmdb.as_str(), "sharded-lmdb");
963 assert_eq!(CorpusExportKind::SingleEnvLmdb.as_str(), "single-env-lmdb");
964 assert_eq!(CorpusExportKind::ExtXyz.as_str(), "extxyz");
965 }
966
967 #[test]
968 fn select_skips_empty_shard_dir_without_data() {
969 let dir = tempfile::tempdir().unwrap();
970 let root = dir.path().join("sharded");
971 let mut s = ShardedConCorpus::open(&root, 4).unwrap();
972 std::fs::create_dir_all(root.join("shard_0001")).unwrap();
973 let keys = s.select(&Select::new()).unwrap();
974 assert!(keys.is_empty());
975 assert!(!root.join("shard_0001").join("data.mdb").exists());
976 let err = s
977 .get_frame_text(FrameKey {
978 traj_id: 1,
979 frame_idx: 0,
980 })
981 .unwrap_err();
982 assert!(
983 err.to_string().contains("is not a corpus directory"),
984 "{err}"
985 );
986 assert!(!root.join("shard_0001").join("data.mdb").exists());
987 }
988
989 #[test]
990 fn select_skips_missing_shard_dirs() {
991 let dir = tempfile::tempdir().unwrap();
992 let root = dir.path().join("sharded");
993 let mut s = ShardedConCorpus::open(&root, 8).unwrap();
994 s.append_trajectory_path(0, fixture("tiny_cuh2.con"))
995 .unwrap();
996 drop(s);
997 let mut s = ShardedConCorpus::open_existing(&root).unwrap();
998 let keys = s.select(&Select::new()).unwrap();
999 assert_eq!(keys.len(), 1);
1000 let present: Vec<u32> = (0..8)
1001 .filter(|&i| root.join(format!("shard_{i:04}")).is_dir())
1002 .collect();
1003 assert_eq!(present, vec![0]);
1004 }
1005
1006 #[test]
1007 fn get_frame_text_missing_shard_does_not_mint() {
1008 let dir = tempfile::tempdir().unwrap();
1009 let root = dir.path().join("sharded");
1010 let mut s = ShardedConCorpus::open(&root, 4).unwrap();
1011 let err = s
1012 .get_frame_text(FrameKey {
1013 traj_id: 1,
1014 frame_idx: 0,
1015 })
1016 .unwrap_err();
1017 assert!(
1018 err.to_string().contains("is not a corpus directory"),
1019 "{err}"
1020 );
1021 assert!(!root.join("shard_0001").exists());
1022 }
1023
1024 #[test]
1025 fn join_drained_roots_keeps_multi_frame_traj() {
1026 let dir = tempfile::tempdir().unwrap();
1027 let text = std::fs::read_to_string(fixture("tiny_multi_cuh2.con")).unwrap();
1028 let node = dir.path().join("node");
1029 ShardedConCorpus::open(&node, 2).unwrap();
1030 ShardedConCorpus::open_shard(&node, 0)
1031 .unwrap()
1032 .append_trajectory_str(0, &text, "a")
1033 .unwrap();
1034 let dest = dir.path().join("dest");
1035 assert_eq!(ShardedConCorpus::drain_to(&node, &dest).unwrap(), 1);
1036 let joined = dir.path().join("joined");
1037 let n = join_drained_roots(&[dest], &joined).unwrap();
1038 assert!(n >= 2, "joined frames={n}");
1039 let db = ConCorpus::open(&joined).unwrap();
1040 let p0 = db
1041 .get_positions(crate::keys::FrameKey {
1042 traj_id: 0,
1043 frame_idx: 0,
1044 })
1045 .unwrap();
1046 let p1 = db
1047 .get_positions(crate::keys::FrameKey {
1048 traj_id: 0,
1049 frame_idx: 1,
1050 })
1051 .unwrap();
1052 assert!((p0[0][0] - 0.6394).abs() < 1e-4);
1053 assert!((p1[2][0] - 8.8549).abs() < 1e-4);
1054 }
1055
1056 #[test]
1057 fn join_corpus_dirs_keeps_multi_frame_traj() {
1058 let dir = tempfile::tempdir().unwrap();
1059 let a = dir.path().join("a");
1060 ConCorpus::open(&a)
1061 .unwrap()
1062 .append_trajectory_path(1, fixture("tiny_multi_cuh2.con"))
1063 .unwrap();
1064 let dest = dir.path().join("out");
1065 let n = join_corpus_dirs(&[a], &dest).unwrap();
1066 assert!(n >= 2, "joined frames={n}");
1067 let db = ConCorpus::open(&dest).unwrap();
1068 let p0 = db
1069 .get_positions(crate::keys::FrameKey {
1070 traj_id: 1,
1071 frame_idx: 0,
1072 })
1073 .unwrap();
1074 let p1 = db
1075 .get_positions(crate::keys::FrameKey {
1076 traj_id: 1,
1077 frame_idx: 1,
1078 })
1079 .unwrap();
1080 assert!((p0[0][0] - 0.6394).abs() < 1e-4);
1081 assert!((p1[2][0] - 8.8549).abs() < 1e-4);
1082 }
1083
1084 #[test]
1085 fn drain_to_n_shards_mismatch_refuses() {
1086 let dir = tempfile::tempdir().unwrap();
1087 let src = dir.path().join("src");
1088 let dest = dir.path().join("dest");
1089 ShardedConCorpus::open(&src, 2).unwrap();
1090 ShardedConCorpus::open(&dest, 4).unwrap();
1091 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
1092 ShardedConCorpus::open_shard(&src, 0)
1093 .unwrap()
1094 .append_trajectory_str(0, &text, "a")
1095 .unwrap();
1096 let err = ShardedConCorpus::drain_to(&src, &dest).unwrap_err();
1097 assert!(err.to_string().contains("n_shards"), "{err}");
1098 let man: ShardManifest =
1099 serde_json::from_str(&std::fs::read_to_string(dest.join("shards.json")).unwrap())
1100 .unwrap();
1101 assert_eq!(man.n_shards, 4);
1102 assert!(!dest.join("shard_0000").join("data.mdb").is_file());
1103 }
1104
1105 #[test]
1106 fn drain_to_new_dest_rollback_on_bad_src() {
1107 let dir = tempfile::tempdir().unwrap();
1108 let src = dir.path().join("src");
1109 ShardedConCorpus::open(&src, 2).unwrap();
1110 let shard0 = src.join("shard_0000");
1111 std::fs::create_dir_all(&shard0).unwrap();
1112 std::fs::write(shard0.join("data.mdb"), b"not-lmdb").unwrap();
1113 let dest = dir.path().join("dest");
1114 assert!(!dest.exists());
1115 assert!(ShardedConCorpus::drain_to(&src, &dest).is_err());
1116 assert!(!dest.exists(), "dest_was_new rollback must remove dest");
1117 }
1118
1119 #[test]
1120 fn drain_to_new_dest_rollback_on_empty_src_manifest() {
1121 let dir = tempfile::tempdir().unwrap();
1122 let src = dir.path().join("src");
1123 std::fs::create_dir_all(&src).unwrap();
1124 std::fs::write(src.join("shards.json"), b"").unwrap();
1125 let dest = dir.path().join("dest");
1126 assert!(!dest.exists());
1127 assert!(ShardedConCorpus::drain_to(&src, &dest).is_err());
1128 assert!(
1129 !dest.exists(),
1130 "dest leftover empty dest_man parse must remove dest"
1131 );
1132 }
1133
1134 #[test]
1135 fn drain_to_new_dest_rollback_on_dest_man_copy() {
1136 let dir = tempfile::tempdir().unwrap();
1137 let src = dir.path().join("src");
1138 ShardedConCorpus::open(&src, 2).unwrap();
1139 let dest = dir.path().join("dest");
1140 std::fs::create_dir_all(dest.join("shards.json")).unwrap();
1141 assert!(ShardedConCorpus::drain_to(&src, &dest).is_err());
1142 assert!(
1143 dest.exists(),
1144 "existing dest stays; dest_man copy failed because dest_man is a dir"
1145 );
1146 assert!(dest.join("shards.json").is_dir());
1147 assert!(!dest.join("shards.json").is_file());
1148 }
1149
1150 #[test]
1151 fn drain_to_new_dest_rollback_on_dest_man_copy_fail() {
1152 let dir = tempfile::tempdir().unwrap();
1153 let src = dir.path().join("src");
1154 ShardedConCorpus::open(&src, 2).unwrap();
1155 let dest = dir.path().join("dest");
1156 assert!(!dest.exists());
1157 FAIL_DEST_MAN_COPY.with(|f| f.set(true));
1158 assert!(ShardedConCorpus::drain_to(&src, &dest).is_err());
1159 assert!(
1160 !dest.exists(),
1161 "dest leftover dest_man copy dest_was_new must remove dest"
1162 );
1163 }
1164
1165 #[test]
1166 fn drain_to_failure_removes_shards_created_this_call() {
1167 let dir = tempfile::tempdir().unwrap();
1168 let src = dir.path().join("src");
1169 ShardedConCorpus::open(&src, 2).unwrap();
1170 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
1171 ShardedConCorpus::open_shard(&src, 0)
1172 .unwrap()
1173 .append_trajectory_str(0, &text, "a")
1174 .unwrap();
1175 ShardedConCorpus::open_shard(&src, 1)
1176 .unwrap()
1177 .append_trajectory_str(1, &text, "b")
1178 .unwrap();
1179 let dest = dir.path().join("dest");
1180 std::fs::create_dir_all(&dest).unwrap();
1181 std::fs::write(dest.join("shard_0001"), b"not-a-dir").unwrap();
1182 assert!(ShardedConCorpus::drain_to(&src, &dest).is_err());
1183 assert!(!dest.join("shard_0000").exists());
1184 }
1185
1186 #[test]
1187 fn rollback_new_dest_removes_dir() {
1188 let dir = tempfile::tempdir().unwrap();
1189 let dest = dir.path().join("out");
1190 std::fs::create_dir_all(&dest).unwrap();
1191 std::fs::write(dest.join("x"), b"y").unwrap();
1192 let err = rollback_new_dest::<u32>(&dest, Err(Error::Message("boom".into()))).unwrap_err();
1193 assert!(err.to_string().contains("boom"), "{err}");
1194 assert!(!dest.exists());
1195 }
1196
1197 #[test]
1198 fn join_to_single_env_keeps_multi_frame_traj() {
1199 let dir = tempfile::tempdir().unwrap();
1200 let root = dir.path().join("sharded");
1201 let mut s = ShardedConCorpus::open(&root, 2).unwrap();
1202 s.append_trajectory_path(0, fixture("tiny_multi_cuh2.con"))
1203 .unwrap();
1204 let dest = dir.path().join("out");
1205 let n = s.join_to_single_env(&dest).unwrap();
1206 assert!(n >= 2, "joined frames={n}");
1207 let db = ConCorpus::open(&dest).unwrap();
1208 let keys = db.select(&Select::new()).unwrap();
1209 assert!(keys.len() >= 2);
1210 assert!(keys.iter().all(|k| k.traj_id == 0));
1211 let p0 = db
1212 .get_positions(crate::keys::FrameKey {
1213 traj_id: 0,
1214 frame_idx: 0,
1215 })
1216 .unwrap();
1217 let p1 = db
1218 .get_positions(crate::keys::FrameKey {
1219 traj_id: 0,
1220 frame_idx: 1,
1221 })
1222 .unwrap();
1223 assert!((p0[0][0] - 0.6394).abs() < 1e-4);
1224 assert!((p1[2][0] - 8.8549).abs() < 1e-4);
1225 }
1226
1227 #[test]
1228 fn join_to_single_env_duplicate_traj_does_not_create_dest() {
1229 let dir = tempfile::tempdir().unwrap();
1230 let root = dir.path().join("sharded");
1231 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
1232 let mut s = ShardedConCorpus::open(&root, 2).unwrap();
1233 s.shard_mut(0)
1234 .unwrap()
1235 .append_trajectory_str(0, &text, "a")
1236 .unwrap();
1237 s.shard_mut(1)
1238 .unwrap()
1239 .append_trajectory_str(0, &text, "b")
1240 .unwrap();
1241 let dest = dir.path().join("out");
1242 let err = s.join_to_single_env(&dest).unwrap_err();
1243 assert!(err.to_string().contains("traj_id"), "{err}");
1244 assert!(!dest.exists());
1245 }
1246
1247 #[test]
1248 fn join_corpus_dirs_duplicate_does_not_create_dest() {
1249 let dir = tempfile::tempdir().unwrap();
1250 let text = std::fs::read_to_string(fixture("tiny_cuh2.con")).unwrap();
1251 let a = dir.path().join("a");
1252 let b = dir.path().join("b");
1253 ConCorpus::open(&a)
1254 .unwrap()
1255 .append_trajectory_str(1, &text, "a")
1256 .unwrap();
1257 ConCorpus::open(&b)
1258 .unwrap()
1259 .append_trajectory_str(1, &text, "b")
1260 .unwrap();
1261 let dest = dir.path().join("out");
1262 let err = join_corpus_dirs(&[a, b], &dest).unwrap_err();
1263 assert!(
1264 err.to_string().contains("traj_id") || err.to_string().contains("duplicate"),
1265 "{err}"
1266 );
1267 assert!(!dest.exists());
1268 }
1269
1270 #[test]
1271 fn join_corpus_dirs_missing_mdb_does_not_create_dest() {
1272 let dir = tempfile::tempdir().unwrap();
1273 let missing = dir.path().join("nope");
1274 std::fs::create_dir_all(&missing).unwrap();
1275 let dest = dir.path().join("out");
1276 let err = join_corpus_dirs(&[missing], &dest).unwrap_err();
1277 assert!(err.to_string().contains("missing data.mdb"), "{err}");
1278 assert!(!dest.exists());
1279 }
1280
1281 #[test]
1282 fn open_single_env_for_export_refuses_sharded_root() {
1283 let dir = tempfile::tempdir().unwrap();
1284 let root = dir.path().join("sharded");
1285 ShardedConCorpus::open(&root, 2).unwrap();
1286 assert!(!root.join("data.mdb").exists());
1287 match open_single_env_for_export(&root) {
1288 Ok(_) => panic!("expected sharded-root refuse"),
1289 Err(err) => assert!(err.to_string().contains("--sharded"), "{err}"),
1290 }
1291 assert!(!root.join("data.mdb").exists());
1292 }
1293
1294 #[test]
1295 fn export_extxyz_removes_temp_join() {
1296 let dir = tempfile::tempdir().unwrap();
1297 let root = dir.path().join("sharded");
1298 let mut s = ShardedConCorpus::open(&root, 2).unwrap();
1299 s.append_trajectory_path(0, fixture("tiny_cuh2.con"))
1300 .unwrap();
1301 let out = dir.path().join("out.xyz");
1302 let n = s.export_extxyz(&Select::new(), &out, "energy").unwrap();
1303 assert!(n >= 1);
1304 assert!(out.is_file());
1305 let prefix = format!("readcon_db_join_{}_", std::process::id());
1306 let leftovers: Vec<_> = std::env::temp_dir()
1307 .read_dir()
1308 .unwrap()
1309 .filter_map(|e| e.ok())
1310 .filter(|e| e.file_name().to_string_lossy().starts_with(&prefix))
1311 .collect();
1312 assert!(leftovers.is_empty(), "{leftovers:?}");
1313 }
1314
1315 #[test]
1316 fn open_existing_error_is_missing_shards_json() {
1317 let dir = tempfile::tempdir().unwrap();
1318 match ShardedConCorpus::open_existing(dir.path().join("nope")) {
1319 Ok(_) => panic!("expected missing shards.json"),
1320 Err(err) => {
1321 assert!(err.to_string().starts_with("missing shards.json:"), "{err}");
1322 assert!(!err.to_string().contains("join-drained"));
1323 }
1324 }
1325 }
1326}