use crate::config;
use chrono::{NaiveDateTime, Utc};
use rocksdb::{BlockBasedOptions, Cache, ColumnFamily, Options, DB};
use serde::{Deserialize, Serialize};
use std::path::Path;
use tracing::info;
const DEFAULT_DIR_NAME: &str = "metadata";
const DEFAULT_MEMTABLE_MEMORY_BUDGET: usize = 32 * 1024 * 1024;
const DEFAULT_MAX_OPEN_FILES: i32 = 10_000;
const DEFAULT_BLOCK_SIZE: usize = 64 * 1024;
const DEFAULT_CACHE_SIZE: usize = 16 * 1024 * 1024;
const TASK_CF_NAME: &str = "task";
const PIECE_CF_NAME: &str = "piece";
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct Task {
pub id: String,
pub piece_length: u64,
pub uploaded_count: u64,
pub updated_at: NaiveDateTime,
pub created_at: NaiveDateTime,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct Piece {
pub number: u32,
pub offset: u64,
pub length: u64,
pub digest: String,
pub uploaded_count: u64,
pub updated_at: NaiveDateTime,
pub created_at: NaiveDateTime,
}
pub struct Metadata {
db: DB,
}
impl Metadata {
pub fn new(data_dir: &Path) -> super::Result<Metadata> {
let mut options = Options::default();
options.create_if_missing(true);
options.create_missing_column_families(true);
options.optimize_level_style_compaction(DEFAULT_MEMTABLE_MEMORY_BUDGET);
options.increase_parallelism(num_cpus::get() as i32);
options.set_max_open_files(DEFAULT_MAX_OPEN_FILES);
let mut block_options = BlockBasedOptions::default();
block_options.set_block_cache(&Cache::new_lru_cache(DEFAULT_CACHE_SIZE));
block_options.set_block_size(DEFAULT_BLOCK_SIZE);
block_options.set_cache_index_and_filter_blocks(true);
block_options.set_pin_l0_filter_and_index_blocks_in_cache(true);
block_options.set_bloom_filter(10.0, false);
options.set_block_based_table_factory(&block_options);
let dir = data_dir.join(config::NAME).join(DEFAULT_DIR_NAME);
let db = DB::open_cf(&options, &dir, [TASK_CF_NAME, PIECE_CF_NAME])?;
info!("create metadata directory: {:?}", dir);
Ok(Metadata { db })
}
pub fn download_task_started(&self, id: &str, piece_length: u64) -> super::Result<()> {
let task = match self.get_task(id)? {
Some(mut task) => {
task.updated_at = Utc::now().naive_utc();
task
}
None => Task {
id: id.to_string(),
piece_length,
updated_at: Utc::now().naive_utc(),
created_at: Utc::now().naive_utc(),
..Default::default()
},
};
self.put_task(id, &task)
}
pub fn upload_task_finished(&self, id: &str) -> super::Result<()> {
match self.get_task(id)? {
Some(mut task) => {
task.uploaded_count += 1;
task.updated_at = Utc::now().naive_utc();
self.put_task(id, &task)
}
None => Err(super::Error::TaskNotFound(id.to_string())),
}
}
pub fn get_task(&self, id: &str) -> super::Result<Option<Task>> {
let handle = self.cf_handle(TASK_CF_NAME)?;
match self.db.get_cf(handle, id)? {
Some(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
None => Ok(None),
}
}
pub fn download_piece_started(&self, id: &str, number: u32) -> super::Result<()> {
self.put_piece(
id,
&Piece {
number,
updated_at: Utc::now().naive_utc(),
created_at: Utc::now().naive_utc(),
..Default::default()
},
)
}
pub fn download_piece_finished(
&self,
id: &str,
offset: u64,
length: u64,
digest: &str,
) -> super::Result<()> {
match self.get_piece(id)? {
Some(mut piece) => {
piece.offset = offset;
piece.length = length;
piece.digest = digest.to_string();
piece.updated_at = Utc::now().naive_utc();
self.put_piece(id, &piece)
}
None => Err(super::Error::PieceNotFound(id.to_string())),
}
}
pub fn upload_piece_finished(&self, id: &str) -> super::Result<()> {
match self.get_piece(id)? {
Some(mut piece) => {
piece.uploaded_count += 1;
piece.updated_at = Utc::now().naive_utc();
self.put_piece(id, &piece)
}
None => Err(super::Error::PieceNotFound(id.to_string())),
}
}
pub fn get_piece(&self, id: &str) -> super::Result<Option<Piece>> {
let handle = self.cf_handle(PIECE_CF_NAME)?;
match self.db.get_cf(handle, id.as_bytes())? {
Some(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
None => Ok(None),
}
}
pub fn piece_id(&self, task_id: &str, number: u32) -> String {
format!("{}-{}", task_id, number)
}
fn put_task(&self, id: &str, task: &Task) -> super::Result<()> {
let handle = self.cf_handle(TASK_CF_NAME)?;
let json = serde_json::to_string(&task)?;
self.db.put_cf(handle, id.as_bytes(), json.as_bytes())?;
Ok(())
}
fn put_piece(&self, id: &str, piece: &Piece) -> super::Result<()> {
let handle = self.cf_handle(PIECE_CF_NAME)?;
let json = serde_json::to_string(&piece)?;
self.db.put_cf(handle, id.as_bytes(), json.as_bytes())?;
Ok(())
}
fn cf_handle(&self, cf_name: &str) -> super::Result<&ColumnFamily> {
self.db
.cf_handle(cf_name)
.ok_or_else(|| super::Error::ColumnFamilyNotFound(cf_name.to_string()))
}
}