1use byteorder::{ByteOrder, LittleEndian};
6use db::SystemDB;
7use kvdb::DBTransaction;
8use parking_lot::Mutex;
9use rlp::Rlp;
10use rlp_derive::{RlpDecodable, RlpEncodable};
11use std::{
12 collections::HashMap,
13 fs::{File, OpenOptions},
14 io::{Error, Read, Seek, SeekFrom, Write},
15 path::PathBuf,
16 sync::Arc,
17};
18
19const COL_DB: u32 = 0;
22const NUM_COLUMNS: u32 = 1;
24
25const DB_KEY_LOG_DEVICE_NUM: &[u8] = b"log_device_num";
26
27const NUM_OF_STRIPES_PER_SEGMENT: u64 = 2000;
28const META_DATA_DB_DIR: &str = "meta_db";
29const LOG_DEVICE_DIR_PREFIX: &str = "log_device_";
30const SEGMENT_FILE_NAME_PREFIX: &str = "segment_";
31
32#[derive(Clone, Copy, Debug, Default, RlpDecodable, RlpEncodable)]
45pub struct StripeReference {
46 segment_id: u64,
48 offset: u64,
50}
51
52#[derive(Clone, Copy, Debug, Default, RlpDecodable, RlpEncodable)]
53pub struct StripeInfo {
54 stripe_ref: StripeReference,
56 stripe_id: u64,
58}
59
60pub struct LogDeviceManager {
61 path_dir: PathBuf,
62 db: Arc<SystemDB>,
63 devices: Mutex<Vec<Arc<LogDevice>>>,
64}
65
66impl LogDeviceManager {
67 pub fn new(path_dir: PathBuf) -> Self {
68 let mut db_dir_path = path_dir.clone();
69 db_dir_path.push(META_DATA_DB_DIR);
70 let db_config = db::db_config(
71 &db_dir_path,
72 None,
73 db::DatabaseCompactionProfile::default(),
74 NUM_COLUMNS,
75 false, );
77
78 let db = db::open_database(db_dir_path.to_str().unwrap(), &db_config)
79 .unwrap();
80
81 let mut log_device_manager = LogDeviceManager {
82 path_dir,
83 db,
84 devices: Mutex::new(Vec::new()),
85 };
86 log_device_manager.initialize();
87 log_device_manager
88 }
89
90 fn initialize(&mut self) {
91 let device_num = self.get_device_num_from_db();
92 let mut devices = self.devices.lock();
93 for i in 0..device_num {
94 let mut log_device_filename = String::from(LOG_DEVICE_DIR_PREFIX);
95 log_device_filename.push_str(i.to_string().as_str());
96 let mut device_path_dir = self.path_dir.clone();
97 device_path_dir.push(log_device_filename.as_str());
98 let log_device = LogDevice::new(
99 device_path_dir,
100 i,
101 self.db.clone(),
102 true, );
104 devices.push(Arc::new(log_device));
105 }
106 }
107
108 fn get_device_num_from_db(&self) -> usize {
109 let res = self
110 .db
111 .key_value()
112 .get(COL_DB, DB_KEY_LOG_DEVICE_NUM)
113 .expect("Low level database error.");
114
115 match res {
116 Some(value) => LittleEndian::read_u64(&value) as usize,
117 None => 0,
118 }
119 }
120
121 fn set_device_num_to_db(&self, device_num: usize) {
122 let mut tx = DBTransaction::new();
123 let mut value = [0; 8];
124 LittleEndian::write_u64(&mut value[0..8], device_num as u64);
125 tx.put(COL_DB, DB_KEY_LOG_DEVICE_NUM, &value);
126 self.db.key_value().write(tx).expect("DB write failed.");
127 }
128
129 pub fn get_device_num(&self) -> usize { self.devices.lock().len() }
130
131 pub fn get_device(&self, device_id: usize) -> Option<Arc<LogDevice>> {
132 Some(self.devices.lock().get(device_id)?.clone())
133 }
134
135 pub fn create_new_device(&self) -> usize {
136 let new_device_id = self.get_device_num();
137 let mut log_device_filename = String::from(LOG_DEVICE_DIR_PREFIX);
138 log_device_filename.push_str(new_device_id.to_string().as_str());
139 let mut device_path_dir = self.path_dir.clone();
140 device_path_dir.push(log_device_filename.as_str());
141 let log_device = LogDevice::new(
142 device_path_dir,
143 new_device_id,
144 self.db.clone(),
145 false, );
147 self.devices.lock().push(Arc::new(log_device));
148 let new_device_num = new_device_id + 1;
149 self.set_device_num_to_db(new_device_num);
150 self.db.key_value().flush().expect("DB flush failed.");
151 new_device_id
152 }
153}
154
155pub struct LogDevice {
156 device_id: usize,
157 tail_db_key: String,
158 head_db_key: String,
159 db: Arc<SystemDB>,
160 inner: Mutex<LogDeviceInner>,
161}
162
163impl LogDevice {
164 pub fn new(
165 path_dir: PathBuf, device_id: usize, db: Arc<SystemDB>, open: bool,
166 ) -> Self {
167 let mut log_device = LogDevice {
168 device_id,
169 tail_db_key: String::default(),
170 head_db_key: String::default(),
171 db: db.clone(),
172 inner: Mutex::new(LogDeviceInner::new(path_dir)),
173 };
174
175 log_device.tail_db_key = log_device.get_tail_key();
176 log_device.head_db_key = log_device.get_head_key();
177 let (head, tail) = if open {
178 let tail = log_device
179 .get_stripe_info_from_db(log_device.tail_db_key.as_bytes())
180 .unwrap();
181 let head = log_device
182 .get_stripe_info_from_db(log_device.head_db_key.as_bytes())
183 .unwrap();
184 (head, tail)
185 } else {
186 let tail = StripeInfo {
187 stripe_ref: StripeReference {
188 segment_id: 0,
189 offset: 0,
190 },
191 stripe_id: 0,
192 };
193
194 let head = tail;
195 log_device.set_stripe_info_to_db(
196 log_device.tail_db_key.as_bytes(),
197 &tail,
198 );
199 log_device.set_stripe_info_to_db(
200 log_device.head_db_key.as_bytes(),
201 &head,
202 );
203 db.key_value().flush().expect("DB flush failed.");
204 (head, tail)
205 };
206
207 log_device.inner.lock().initialize(head, tail);
208 log_device
209 }
210
211 fn get_tail_key(&self) -> String {
212 let mut tail_key = String::from(LOG_DEVICE_DIR_PREFIX);
213 tail_key.push_str(self.device_id.to_string().as_str());
214 tail_key.push_str("_tail");
215 tail_key
216 }
217
218 fn get_head_key(&self) -> String {
219 let mut head_key = String::from(LOG_DEVICE_DIR_PREFIX);
220 head_key.push_str(self.device_id.to_string().as_str());
221 head_key.push_str("_head");
222 head_key
223 }
224
225 fn get_stripe_info_from_db(&self, key: &[u8]) -> Option<StripeInfo> {
226 let res = self
227 .db
228 .key_value()
229 .get(COL_DB, key)
230 .expect("Low level database error.");
231 match res {
232 Some(value) => {
233 let rlp = Rlp::new(&value);
234 let stripe_info: StripeInfo = rlp.as_val().expect("rlp error");
235 Some(stripe_info)
236 }
237 None => None,
238 }
239 }
240
241 fn set_stripe_info_to_db(&self, key: &[u8], stripe_info: &StripeInfo) {
242 let value = rlp::encode(stripe_info);
243 let mut tx = DBTransaction::new();
244 tx.put(COL_DB, key, &value);
245 self.db.key_value().write(tx).expect("DB write failed.");
246 }
247
248 pub fn append_stripe(&self, stripe: &[u8]) -> Result<StripeInfo, Error> {
249 let (appended_stripe, tail) =
250 self.inner.lock().append_stripe(stripe)?;
251 self.set_stripe_info_to_db(self.tail_db_key.as_bytes(), &tail);
252 self.db.key_value().flush()?;
253 Ok(appended_stripe)
254 }
255
256 pub fn get_stripe(
257 &self, stripe_ref: &StripeReference,
258 ) -> Result<Vec<u8>, Error> {
259 self.inner.lock().get_stripe(stripe_ref)
260 }
261
262 pub fn trim(&self, stripe: &StripeInfo) {
263 let mut inner = self.inner.lock();
264 let new_head = inner.check_trim(stripe.stripe_ref.segment_id);
265 if let Some(new_head) = new_head {
266 self.set_stripe_info_to_db(self.head_db_key.as_bytes(), &new_head);
267 self.db.key_value().flush().expect("DB flush failed.");
268 inner.trim(&new_head);
269 }
270 }
271
272 pub fn segment_to_file_name(segment_id: u64) -> String {
273 let mut filename = String::from(SEGMENT_FILE_NAME_PREFIX);
274 filename.push_str(segment_id.to_string().as_str());
275 filename
276 }
277}
278
279struct LogDeviceInner {
280 tail: StripeInfo,
282 head: StripeInfo,
285 path_dir: PathBuf,
287 file_cache: HashMap<u64, File>,
290}
291
292impl LogDeviceInner {
293 pub fn new(path_dir: PathBuf) -> Self {
294 LogDeviceInner {
295 tail: StripeInfo::default(),
296 head: StripeInfo::default(),
297 path_dir,
298 file_cache: HashMap::new(),
299 }
300 }
301
302 fn initialize(&mut self, head: StripeInfo, tail: StripeInfo) {
303 self.head = head;
304 self.tail = tail;
305
306 let segment_path =
308 self.segment_to_path(self.tail.stripe_ref.segment_id);
309 let create = if segment_path.exists() {
310 false
311 } else {
312 assert_eq!(self.tail.stripe_ref.segment_id, 0);
313 assert_eq!(self.tail.stripe_ref.offset, 0);
314 assert_eq!(self.tail.stripe_id, 0);
315 std::fs::create_dir_all(&self.path_dir)
316 .expect("Failed to create log_device dir.");
317 true
318 };
319 let mut segment_file = OpenOptions::new()
320 .read(true)
321 .write(true)
322 .create_new(create)
323 .open(&segment_path)
324 .expect("Failed to open segment file.");
325 let offset = segment_file
326 .seek(SeekFrom::Start(self.tail.stripe_ref.offset))
327 .expect("Failed to seek segment file.");
328 assert_eq!(offset, self.tail.stripe_ref.offset);
329 self.file_cache
330 .insert(self.tail.stripe_ref.segment_id, segment_file);
331 }
332
333 fn segment_to_path(&self, segment: u64) -> PathBuf {
334 let segment_filename = LogDevice::segment_to_file_name(segment);
335 let mut segment_path = self.path_dir.clone();
336 segment_path.push(segment_filename.as_str());
337 segment_path
338 }
339
340 pub fn append_stripe(
341 &mut self, stripe: &[u8],
342 ) -> Result<(StripeInfo, StripeInfo), Error> {
343 let payload_size = LittleEndian::read_u32(&stripe[0..4]) as usize;
345 assert_eq!(payload_size + 4, stripe.len(), "Incorrect payload size.");
346
347 if self.tail.stripe_id == NUM_OF_STRIPES_PER_SEGMENT {
348 self.tail.stripe_ref.segment_id += 1;
350 self.tail.stripe_ref.offset = 0;
351 self.tail.stripe_id = 0;
352
353 let segment_path =
355 self.segment_to_path(self.tail.stripe_ref.segment_id);
356 let segment_file = OpenOptions::new()
357 .read(true)
358 .write(true)
359 .create_new(true)
360 .open(&segment_path)?;
361 self.file_cache
362 .insert(self.tail.stripe_ref.segment_id, segment_file);
363 }
364
365 let segment_file = self
367 .file_cache
368 .get_mut(&self.tail.stripe_ref.segment_id)
369 .unwrap();
370 let write_size = segment_file.write(stripe)?;
371 assert_eq!(write_size, stripe.len());
373 let offset = self.tail.stripe_ref.offset + write_size as u64;
374 assert_eq!(segment_file.seek(SeekFrom::End(0)).unwrap(), offset);
375 segment_file.flush()?;
376
377 let appended_stripe = self.tail;
378 self.tail.stripe_id += 1;
380 self.tail.stripe_ref.offset = offset;
381 Ok((appended_stripe, self.tail))
382 }
383
384 pub fn get_stripe(
385 &mut self, stripe_ref: &StripeReference,
386 ) -> Result<Vec<u8>, Error> {
387 if !self.file_cache.contains_key(&stripe_ref.segment_id) {
388 let segment_path = self.segment_to_path(stripe_ref.segment_id);
390 let segment_file = OpenOptions::new()
391 .read(true)
392 .write(true)
393 .open(&segment_path)?;
394 self.file_cache.insert(stripe_ref.segment_id, segment_file);
395 }
396
397 let segment_file =
398 self.file_cache.get_mut(&stripe_ref.segment_id).unwrap();
399 let offset = segment_file.seek(SeekFrom::Start(stripe_ref.offset))?;
400 assert_eq!(offset, stripe_ref.offset);
401 let mut stripe: Vec<u8> = vec![0; 4];
402 let read_size = segment_file.read(&mut stripe[0..4])?;
403 assert_eq!(read_size, 4);
404 let payload_size = LittleEndian::read_u32(&stripe[0..4]) as usize;
405 if payload_size != 0 {
406 stripe.resize(payload_size + 4, 0);
407 let read_size =
408 segment_file.read(&mut stripe[4..4 + payload_size])?;
409 assert_eq!(read_size, payload_size);
410 }
411 Ok(stripe)
412 }
413
414 pub fn check_trim(&self, segment_id: u64) -> Option<StripeInfo> {
415 if segment_id >= self.head.stripe_ref.segment_id
416 && segment_id <= self.tail.stripe_ref.segment_id
417 {
418 Some(StripeInfo {
419 stripe_ref: StripeReference {
420 segment_id,
421 offset: 0,
422 },
423 stripe_id: 0,
424 })
425 } else {
426 None
427 }
428 }
429
430 pub fn trim(&mut self, new_head: &StripeInfo) {
431 let old_head = self.head;
432 self.head = *new_head;
433
434 for segment in
435 old_head.stripe_ref.segment_id..self.head.stripe_ref.segment_id
436 {
437 self.file_cache.remove(&segment);
438 let segment_path = self.segment_to_path(segment);
439 if segment_path.exists() {
440 std::fs::remove_file(&segment_path).ok();
441 }
442 }
443 }
444}
445
446#[cfg(test)]
447mod tests {
448 use super::{
449 LogDevice, LogDeviceManager, LOG_DEVICE_DIR_PREFIX,
450 NUM_OF_STRIPES_PER_SEGMENT,
451 };
452 use crate::{StripeInfo, StripeReference};
453 use byteorder::{ByteOrder, LittleEndian};
454 use rand::Rng;
455 use std::{path::PathBuf, sync::Arc};
456
457 fn gen_random_and_append(
458 log_device: Arc<LogDevice>, stripes: &mut Vec<Vec<u8>>,
459 stripe_refs: &mut Vec<StripeReference>, start: usize, end: usize,
460 ) {
461 for i in start..end {
462 let mut stripe: Vec<u8> = Vec::new();
463 let stripe_size = rand::rng().random_range(4..1024 * 64);
464 stripe.resize(stripe_size, i as u8);
465 let payload_size = stripe_size - 4;
466 LittleEndian::write_u32(&mut stripe[0..4], payload_size as u32);
467 let stripe_info = log_device.append_stripe(&stripe).unwrap();
468 stripes.push(stripe);
469 stripe_refs.push(stripe_info.stripe_ref);
470 }
471 }
472
473 fn read_and_check(
474 log_device: Arc<LogDevice>, stripes: &[Vec<u8>],
475 stripe_refs: &[StripeReference], start: usize, end: usize,
476 ) {
477 for i in start..end {
478 let stripe = &stripes[i];
479 let stripe_ref = &stripe_refs[i];
480 let read_stripe = log_device
481 .get_stripe(stripe_ref)
482 .expect("Failed to read stripe");
483 let matching = stripe
484 .iter()
485 .zip(read_stripe.iter())
486 .filter(|&(a, b)| a == b)
487 .count();
488 assert_eq!(matching, stripe.len());
489 assert_eq!(matching, read_stripe.len());
490 }
491 }
492
493 fn create_and_append(
494 stripes: &mut Vec<Vec<u8>>, stripe_refs: &mut Vec<StripeReference>,
495 ) {
496 let path_dir = String::from("./ldm_open");
497 let path_dir = PathBuf::from(path_dir);
498 std::fs::remove_dir_all(&path_dir).ok();
499 std::fs::create_dir_all(&path_dir).ok();
500 let log_device_manager = LogDeviceManager::new(path_dir.clone());
501 assert_eq!(log_device_manager.get_device_num(), 0);
502 let device_id = log_device_manager.create_new_device();
503 assert_eq!(log_device_manager.get_device_num(), 1);
504 assert_eq!(log_device_manager.get_device_num_from_db(), 1);
505 let log_device = log_device_manager.get_device(device_id).unwrap();
506
507 gen_random_and_append(log_device.clone(), stripes, stripe_refs, 0, 10);
508 read_and_check(log_device.clone(), stripes, stripe_refs, 0, 10);
509 }
510
511 fn open_and_append_and_read(
512 stripes: &mut Vec<Vec<u8>>, stripe_refs: &mut Vec<StripeReference>,
513 ) {
514 let path_dir = String::from("./ldm_open");
515 let path_dir = PathBuf::from(path_dir);
516 let log_device_manager = LogDeviceManager::new(path_dir.clone());
517 assert_eq!(log_device_manager.get_device_num(), 1);
518 let log_device = log_device_manager.get_device(0).unwrap();
519
520 gen_random_and_append(log_device.clone(), stripes, stripe_refs, 10, 20);
521 read_and_check(log_device.clone(), stripes, stripe_refs, 0, 20);
522 std::fs::remove_dir_all(&path_dir).ok();
523 }
524
525 #[test]
526 fn test_open_log_device() {
527 let mut stripes = Vec::new();
528 let mut stripe_refs = Vec::new();
529
530 create_and_append(&mut stripes, &mut stripe_refs);
531 open_and_append_and_read(&mut stripes, &mut stripe_refs);
532 }
533
534 #[test]
535 fn test_append_log_device() {
536 let path_dir = String::from("./ldm_append");
537 let path_dir = PathBuf::from(path_dir);
538 std::fs::remove_dir_all(&path_dir).ok();
539 std::fs::create_dir_all(&path_dir).ok();
540 let log_device_manager = LogDeviceManager::new(path_dir.clone());
541 assert_eq!(log_device_manager.get_device_num(), 0);
542 let device_id = log_device_manager.create_new_device();
543 assert_eq!(log_device_manager.get_device_num(), 1);
544 let log_device = log_device_manager.get_device(device_id).unwrap();
545 let mut stripes = Vec::new();
546 let mut stripe_refs = Vec::new();
547
548 gen_random_and_append(
549 log_device.clone(),
550 &mut stripes,
551 &mut stripe_refs,
552 0,
553 10,
554 );
555 read_and_check(log_device.clone(), &stripes, &stripe_refs, 0, 10);
556 std::fs::remove_dir_all(&path_dir).ok();
557 }
558
559 #[test]
560 fn test_trim_log_device() {
561 let path_dir = String::from("./ldm_trim");
562 let path_dir = PathBuf::from(path_dir);
563 std::fs::remove_dir_all(&path_dir).ok();
564 std::fs::create_dir_all(&path_dir).ok();
565 let log_device_manager = LogDeviceManager::new(path_dir.clone());
566 assert_eq!(log_device_manager.get_device_num(), 0);
567 let device_id = log_device_manager.create_new_device();
568 assert_eq!(log_device_manager.get_device_num(), 1);
569 let log_device = log_device_manager.get_device(device_id).unwrap();
570 let mut stripes = Vec::new();
571 let mut stripe_refs = Vec::new();
572
573 gen_random_and_append(
574 log_device.clone(),
575 &mut stripes,
576 &mut stripe_refs,
577 0,
578 4 * NUM_OF_STRIPES_PER_SEGMENT as usize,
579 );
580
581 let mut log_device_path_dir = path_dir.clone();
582 let mut log_device_dir = String::from(LOG_DEVICE_DIR_PREFIX);
583 log_device_dir.push('0');
584 log_device_path_dir.push(log_device_dir.as_str());
585
586 let mut segment_0_path = log_device_path_dir.clone();
587 segment_0_path.push("segment_0");
588 let mut segment_1_path = log_device_path_dir.clone();
589 segment_1_path.push("segment_1");
590 let mut segment_2_path = log_device_path_dir.clone();
591 segment_2_path.push("segment_2");
592 let mut segment_3_path = log_device_path_dir.clone();
593 segment_3_path.push("segment_3");
594
595 assert!(segment_0_path.exists());
596 assert!(segment_1_path.exists());
597 assert!(segment_2_path.exists());
598 assert!(segment_3_path.exists());
599
600 let strip_info = StripeInfo {
601 stripe_ref: StripeReference {
602 segment_id: 2,
603 offset: 0,
604 },
605 stripe_id: 0,
606 };
607 log_device.trim(&strip_info);
608
609 assert!(!segment_0_path.exists());
610 assert!(!segment_1_path.exists());
611 assert!(segment_2_path.exists());
612 assert!(segment_3_path.exists());
613
614 read_and_check(
615 log_device.clone(),
616 &stripes,
617 &stripe_refs,
618 2 * NUM_OF_STRIPES_PER_SEGMENT as usize,
619 4 * NUM_OF_STRIPES_PER_SEGMENT as usize,
620 );
621
622 std::fs::remove_dir_all(&path_dir).ok();
623 }
624
625 #[test]
626 fn test_create_log_device() {
627 let path_dir = String::from("./ldm_create");
628 let path_dir = PathBuf::from(path_dir);
629 std::fs::remove_dir_all(&path_dir).ok();
630 std::fs::create_dir_all(&path_dir).ok();
631 let log_device_manager = LogDeviceManager::new(path_dir.clone());
632 assert_eq!(log_device_manager.get_device_num(), 0);
633 assert_eq!(log_device_manager.get_device_num_from_db(), 0);
634 log_device_manager.create_new_device();
635 assert_eq!(log_device_manager.get_device_num(), 1);
636 assert_eq!(log_device_manager.get_device_num_from_db(), 1);
637 log_device_manager.create_new_device();
638 assert_eq!(log_device_manager.get_device_num(), 2);
639 assert_eq!(log_device_manager.get_device_num_from_db(), 2);
640 log_device_manager.create_new_device();
641 assert_eq!(log_device_manager.get_device_num(), 3);
642 assert_eq!(log_device_manager.get_device_num_from_db(), 3);
643 log_device_manager.create_new_device();
644 assert_eq!(log_device_manager.get_device_num(), 4);
645 assert_eq!(log_device_manager.get_device_num_from_db(), 4);
646
647 let mut log_device_path_dir = path_dir.clone();
648 let mut log_device_dir = String::from(LOG_DEVICE_DIR_PREFIX);
649 log_device_dir.push('0');
650 log_device_path_dir.push(log_device_dir.as_str());
651 let mut segment_0_path = log_device_path_dir.clone();
652 segment_0_path.push("segment_0");
653 assert!(segment_0_path.exists());
654
655 let mut log_device_path_dir = path_dir.clone();
656 let mut log_device_dir = String::from(LOG_DEVICE_DIR_PREFIX);
657 log_device_dir.push('1');
658 log_device_path_dir.push(log_device_dir.as_str());
659 let mut segment_0_path = log_device_path_dir.clone();
660 segment_0_path.push("segment_0");
661 assert!(segment_0_path.exists());
662
663 let mut log_device_path_dir = path_dir.clone();
664 let mut log_device_dir = String::from(LOG_DEVICE_DIR_PREFIX);
665 log_device_dir.push('2');
666 log_device_path_dir.push(log_device_dir.as_str());
667 let mut segment_0_path = log_device_path_dir.clone();
668 segment_0_path.push("segment_0");
669 assert!(segment_0_path.exists());
670
671 let mut log_device_path_dir = path_dir.clone();
672 let mut log_device_dir = String::from(LOG_DEVICE_DIR_PREFIX);
673 log_device_dir.push('3');
674 log_device_path_dir.push(log_device_dir.as_str());
675 let mut segment_0_path = log_device_path_dir.clone();
676 segment_0_path.push("segment_0");
677 assert!(segment_0_path.exists());
678
679 std::fs::remove_dir_all(&path_dir).ok();
680 }
681}