Files
dht/src/search/src/storage/rocks.rs
T

1080 lines
38 KiB
Rust

// 负责实现 RocksDB 打开配置批量写入精确查询和关闭流程
use std::sync::{
Arc, Mutex, RwLock,
atomic::{AtomicU64, Ordering},
};
use rocksdb::{DB, Direction, IteratorMode, WriteBatch};
use crate::domain::{ContentFilter, InfoHash, RejectedMetadata, TorrentRecord, VerificationResult};
use super::{
keys::{
CONTENT_GROUP_COUNT_KEY, PENDING_INDEX_COUNT_KEY, TORRENT_COUNT_KEY,
VERIFICATION_QUEUE_COUNT_KEY, content_group_key, content_group_prefix, content_member_key,
content_members_prefix, decode_content_member_info_hash, decode_pending_content_key,
decode_verification_lease, decode_verification_task, filter_pending_key,
filter_projection_key, pending_index_key, pending_index_prefix, rejected_metadata_key,
torrent_key, torrent_prefix, verification_lease_key, verification_lease_prefix,
verification_locator_key, verification_task_key, verification_task_prefix,
},
repository::{
ContentGroupTask, ContentVariants, IndexDocument, IndexInventory, IndexRefreshDiagnostics,
StorageError, TorrentRepository, UpsertOutcome, VerificationEnqueueOutcome,
VerificationPriority, VerificationRequest, dynamic_projection_hash, projection_hash,
refresh_bucket,
},
};
mod filter_migration;
mod lifecycle;
const DEFAULT_BLOCK_CACHE_BYTES: usize = 64 * 1024 * 1024;
const UNKNOWN_DYNAMIC_PROJECTION: [u8; 32] = [0; 32];
const BASELINED_DYNAMIC_PROJECTION: [u8; 32] = [u8::MAX; 32];
type VerificationTaskEntry = (Box<[u8]>, InfoHash);
#[derive(Debug, Clone, Copy, serde::Serialize, serde::Deserialize)]
struct ContentGroupState {
revision: u64,
member_count: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct ProjectionState {
hash: [u8; 32],
dynamic_hash: [u8; 32],
refresh_bucket: u64,
visible: bool,
}
pub struct RocksTorrentRepository {
db: DB,
write_lock: Mutex<()>,
rejection_rule_id: [u8; 32],
canonical_filter: Arc<ContentFilter>,
content_filter: RwLock<Arc<ContentFilter>>,
inventory: InventoryState,
index_refresh_scheduled: AtomicU64,
index_refresh_suppressed: AtomicU64,
}
#[derive(Default)]
struct InventoryState {
stored_torrents: AtomicU64,
searchable_groups: AtomicU64,
pending_documents: AtomicU64,
}
#[derive(Debug, Clone, Copy, Default)]
struct InventoryDelta {
stored_torrents: u64,
searchable_groups: u64,
pending_documents: i64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum RefreshDecision {
Scheduled,
Suppressed,
}
impl RocksTorrentRepository {
fn current_filter(&self) -> Arc<ContentFilter> {
self.content_filter
.read()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone()
}
pub(crate) fn set_content_filter(&self, filter: Arc<ContentFilter>) {
*self
.content_filter
.write()
.unwrap_or_else(std::sync::PoisonError::into_inner) = filter;
}
fn encode(record: &TorrentRecord) -> Result<Vec<u8>, StorageError> {
rmp_serde::to_vec_named(record).map_err(Into::into)
}
fn decode(bytes: &[u8]) -> Result<TorrentRecord, StorageError> {
rmp_serde::from_slice(bytes).map_err(Into::into)
}
fn encode_rejection(rejection: &RejectedMetadata) -> Result<Vec<u8>, StorageError> {
rmp_serde::to_vec_named(rejection).map_err(Into::into)
}
fn decode_rejection(bytes: &[u8]) -> Result<RejectedMetadata, StorageError> {
rmp_serde::from_slice(bytes).map_err(Into::into)
}
fn projection_state(
&self,
content_key: &[u8; 32],
) -> Result<Option<ProjectionState>, StorageError> {
self.db
.get(filter_projection_key(content_key))?
.map(|bytes| decode_projection(&bytes))
.transpose()
}
fn encode_group(state: ContentGroupState) -> Result<Vec<u8>, StorageError> {
rmp_serde::to_vec_named(&state).map_err(Into::into)
}
fn group_state(
&self,
content_key: &[u8; 32],
) -> Result<Option<ContentGroupState>, StorageError> {
self.db
.get(content_group_key(content_key))?
.map(|bytes| rmp_serde::from_slice(&bytes).map_err(StorageError::from))
.transpose()
}
fn dirty_group(
&self,
batch: &mut WriteBatch,
content_key: &[u8; 32],
inserted: bool,
) -> Result<InventoryDelta, StorageError> {
let current = self.group_state(content_key)?;
let pending_value = self.db.get(pending_index_key(content_key))?;
let pending = pending_value.is_some();
let force_write = pending_value
.as_deref()
.map(decode_pending_revision)
.transpose()?
.is_some_and(|(_, force_write)| force_write);
let mut state = current.unwrap_or(ContentGroupState {
revision: 0,
member_count: 0,
});
state.revision = state.revision.saturating_add(1);
if inserted {
state.member_count = state.member_count.saturating_add(1);
}
batch.put(content_group_key(content_key), Self::encode_group(state)?);
if force_write {
batch.put(
pending_index_key(content_key),
encode_forced_pending_revision(state.revision),
);
} else {
batch.put(pending_index_key(content_key), state.revision.to_be_bytes());
}
Ok(InventoryDelta {
searchable_groups: u64::from(current.is_none()),
pending_documents: i64::from(!pending),
..InventoryDelta::default()
})
}
fn schedule_observation_refresh(
&self,
batch: &mut WriteBatch,
content_key: &[u8; 32],
observed_at: u64,
) -> Result<RefreshDecision, StorageError> {
let Some(mut state) = self.projection_state(content_key)? else {
return Ok(RefreshDecision::Scheduled);
};
let bucket = refresh_bucket(observed_at);
if state.dynamic_hash == UNKNOWN_DYNAMIC_PROJECTION {
state.dynamic_hash = BASELINED_DYNAMIC_PROJECTION;
state.refresh_bucket = bucket;
batch.put(filter_projection_key(content_key), encode_projection(state));
return Ok(RefreshDecision::Suppressed);
}
if bucket <= state.refresh_bucket {
return Ok(RefreshDecision::Suppressed);
}
state.refresh_bucket = bucket;
batch.put(filter_projection_key(content_key), encode_projection(state));
Ok(RefreshDecision::Scheduled)
}
fn record_refresh_decision(&self, decision: RefreshDecision) {
match decision {
RefreshDecision::Scheduled => {
self.index_refresh_scheduled.fetch_add(1, Ordering::Relaxed);
}
RefreshDecision::Suppressed => {
self.index_refresh_suppressed
.fetch_add(1, Ordering::Relaxed);
}
}
}
fn inventory_snapshot(&self) -> IndexInventory {
IndexInventory {
stored_torrents: self.inventory.stored_torrents.load(Ordering::Relaxed),
searchable_groups: self.inventory.searchable_groups.load(Ordering::Relaxed),
pending_documents: self.inventory.pending_documents.load(Ordering::Relaxed),
}
}
fn commit_inventory_batch(
&self,
mut batch: WriteBatch,
delta: InventoryDelta,
) -> Result<(), StorageError> {
let current = self.inventory_snapshot();
let next = IndexInventory {
stored_torrents: current
.stored_torrents
.saturating_add(delta.stored_torrents),
searchable_groups: current
.searchable_groups
.saturating_add(delta.searchable_groups),
pending_documents: if delta.pending_documents >= 0 {
current
.pending_documents
.saturating_add(delta.pending_documents as u64)
} else {
current
.pending_documents
.saturating_sub(delta.pending_documents.unsigned_abs())
},
};
batch.put(TORRENT_COUNT_KEY, next.stored_torrents.to_be_bytes());
batch.put(
CONTENT_GROUP_COUNT_KEY,
next.searchable_groups.to_be_bytes(),
);
batch.put(
PENDING_INDEX_COUNT_KEY,
next.pending_documents.to_be_bytes(),
);
self.db.write(batch)?;
self.set_inventory(next);
Ok(())
}
fn set_inventory(&self, inventory: IndexInventory) {
self.inventory
.stored_torrents
.store(inventory.stored_torrents, Ordering::Relaxed);
self.inventory
.searchable_groups
.store(inventory.searchable_groups, Ordering::Relaxed);
self.inventory
.pending_documents
.store(inventory.pending_documents, Ordering::Relaxed);
}
fn initialize_inventory(&self, force_scan: bool) -> Result<(), StorageError> {
if !force_scan
&& let (Some(stored), Some(groups), Some(pending)) = (
read_counter(&self.db, TORRENT_COUNT_KEY)?,
read_counter(&self.db, CONTENT_GROUP_COUNT_KEY)?,
read_counter(&self.db, PENDING_INDEX_COUNT_KEY)?,
)
{
self.set_inventory(IndexInventory {
stored_torrents: stored,
searchable_groups: groups,
pending_documents: pending,
});
return Ok(());
}
let inventory = IndexInventory {
stored_torrents: count_prefix(&self.db, torrent_prefix())?,
searchable_groups: count_prefix(&self.db, content_group_prefix())?,
pending_documents: count_prefix(&self.db, pending_index_prefix())?,
};
let mut batch = WriteBatch::default();
batch.put(TORRENT_COUNT_KEY, inventory.stored_torrents.to_be_bytes());
batch.put(
CONTENT_GROUP_COUNT_KEY,
inventory.searchable_groups.to_be_bytes(),
);
batch.put(
PENDING_INDEX_COUNT_KEY,
inventory.pending_documents.to_be_bytes(),
);
self.db.write(batch)?;
self.set_inventory(inventory);
Ok(())
}
fn verification_queue_len_inner(&self) -> Result<usize, StorageError> {
let Some(value) = self.db.get(VERIFICATION_QUEUE_COUNT_KEY)? else {
return Ok(0);
};
let bytes: [u8; 8] = value
.as_slice()
.try_into()
.map_err(|_| StorageError::CorruptVerificationQueueCount)?;
Ok(u64::from_be_bytes(bytes).min(usize::MAX as u64) as usize)
}
fn first_verification_task(
&self,
high_priority: bool,
) -> Result<Option<VerificationTaskEntry>, StorageError> {
let prefix = verification_task_prefix(high_priority);
let mut iterator = self
.db
.iterator(IteratorMode::From(&prefix, Direction::Forward));
let Some(entry) = iterator.next() else {
return Ok(None);
};
let (key, _) = entry?;
let Some((_, _, info_hash)) = decode_verification_task(&key) else {
return Ok(None);
};
Ok(Some((key, info_hash)))
}
}
impl TorrentRepository for RocksTorrentRepository {
fn get(&self, info_hash: InfoHash) -> Result<Option<TorrentRecord>, StorageError> {
self.db
.get(torrent_key(info_hash))?
.map(|bytes| Self::decode(&bytes))
.transpose()
}
fn get_visible(&self, info_hash: InfoHash) -> Result<Option<TorrentRecord>, StorageError> {
let filter = self.current_filter();
Ok(self
.get(info_hash)?
.and_then(|record| filter.public_record(&record)))
}
fn upsert(&self, mut observation: TorrentRecord) -> Result<UpsertOutcome, StorageError> {
let visible: Vec<_> = observation
.files
.iter()
.filter(|file| !self.canonical_filter.is_file_hidden(file))
.cloned()
.collect();
observation.searchable = !visible.is_empty();
observation.content_key = if observation.searchable {
crate::domain::content_key(&visible)?
} else {
[0; 32]
};
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(mut current) = self.get(observation.info_hash)? {
current.observe_again(observation.last_seen, &observation.source_peers);
let mut batch = WriteBatch::default();
batch.put(torrent_key(current.info_hash), Self::encode(&current)?);
let (delta, refresh) = if current.searchable {
let refresh = self.schedule_observation_refresh(
&mut batch,
&current.content_key,
current.last_seen,
)?;
let delta = if refresh == RefreshDecision::Scheduled {
self.dirty_group(&mut batch, &current.content_key, false)?
} else {
InventoryDelta::default()
};
(delta, Some(refresh))
} else {
(InventoryDelta::default(), None)
};
self.commit_inventory_batch(batch, delta)?;
if let Some(refresh) = refresh {
self.record_refresh_decision(refresh);
}
return Ok(UpsertOutcome::Updated {
seen_count: current.seen_count,
});
}
let mut batch = WriteBatch::default();
batch.put(
torrent_key(observation.info_hash),
Self::encode(&observation)?,
);
batch.delete(rejected_metadata_key(observation.info_hash));
let mut delta = InventoryDelta {
stored_torrents: 1,
..InventoryDelta::default()
};
if observation.searchable {
batch.put(
content_member_key(&observation.content_key, observation.info_hash),
[],
);
let group_delta = self.dirty_group(&mut batch, &observation.content_key, true)?;
delta.searchable_groups = group_delta.searchable_groups;
delta.pending_documents = group_delta.pending_documents;
}
self.commit_inventory_batch(batch, delta)?;
Ok(UpsertOutcome::Inserted)
}
fn insert_if_absent(&self, mut observation: TorrentRecord) -> Result<bool, StorageError> {
let visible: Vec<_> = observation
.files
.iter()
.filter(|file| !self.canonical_filter.is_file_hidden(file))
.cloned()
.collect();
observation.searchable = !visible.is_empty();
observation.content_key = if observation.searchable {
crate::domain::content_key(&visible)?
} else {
[0; 32]
};
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.get(observation.info_hash)?.is_some() {
return Ok(false);
}
let mut batch = WriteBatch::default();
batch.put(
torrent_key(observation.info_hash),
Self::encode(&observation)?,
);
batch.delete(rejected_metadata_key(observation.info_hash));
let mut delta = InventoryDelta {
stored_torrents: 1,
..InventoryDelta::default()
};
if observation.searchable {
batch.put(
content_member_key(&observation.content_key, observation.info_hash),
[],
);
let group_delta = self.dirty_group(&mut batch, &observation.content_key, true)?;
delta.searchable_groups = group_delta.searchable_groups;
delta.pending_documents = group_delta.pending_documents;
}
self.commit_inventory_batch(batch, delta)?;
Ok(true)
}
fn rejection(&self, info_hash: InfoHash) -> Result<Option<RejectedMetadata>, StorageError> {
self.db
.get(rejected_metadata_key(info_hash))?
.map(|bytes| Self::decode_rejection(&bytes))
.transpose()
}
fn record_rejection(&self, mut rejection: RejectedMetadata) -> Result<(), StorageError> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.get(rejection.info_hash)?.is_some() {
return Ok(());
}
if let Some(mut current) = self.rejection(rejection.info_hash)?
&& current.rule_id == rejection.rule_id
{
current.observe_again(rejection.last_seen);
rejection = current;
}
self.db.put(
rejected_metadata_key(rejection.info_hash),
Self::encode_rejection(&rejection)?,
)?;
Ok(())
}
fn observe_existing(
&self,
info_hash: InfoHash,
observed_at: u64,
) -> Result<bool, StorageError> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(mut record) = self.get(info_hash)? else {
let Some(mut rejection) = self.rejection(info_hash)? else {
return Ok(false);
};
if rejection.rule_id != self.rejection_rule_id {
return Ok(false);
}
rejection.observe_again(observed_at);
self.db.put(
rejected_metadata_key(info_hash),
Self::encode_rejection(&rejection)?,
)?;
return Ok(true);
};
record.observe_again(observed_at, &[]);
let mut batch = WriteBatch::default();
batch.put(torrent_key(info_hash), Self::encode(&record)?);
let (delta, refresh) = if record.searchable {
let refresh = self.schedule_observation_refresh(
&mut batch,
&record.content_key,
record.last_seen,
)?;
let delta = if refresh == RefreshDecision::Scheduled {
self.dirty_group(&mut batch, &record.content_key, false)?
} else {
InventoryDelta::default()
};
(delta, Some(refresh))
} else {
(InventoryDelta::default(), None)
};
self.commit_inventory_batch(batch, delta)?;
if let Some(refresh) = refresh {
self.record_refresh_decision(refresh);
}
Ok(true)
}
fn filter_unknown_and_observe(
&self,
info_hashes: &[InfoHash],
observed_at: u64,
) -> Result<Vec<InfoHash>, StorageError> {
if info_hashes.is_empty() {
return Ok(Vec::new());
}
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let keys: Vec<_> = info_hashes.iter().copied().map(torrent_key).collect();
let rejected_keys: Vec<_> = info_hashes
.iter()
.copied()
.map(rejected_metadata_key)
.collect();
let records = self.db.multi_get(keys.iter());
let rejections = self.db.multi_get(rejected_keys.iter());
let mut unknown = Vec::with_capacity(info_hashes.len());
let mut batch = WriteBatch::default();
let mut updated = 0_usize;
let mut observed_groups = std::collections::BTreeMap::new();
for ((((info_hash, key), rejected_key), record), rejection) in info_hashes
.iter()
.zip(&keys)
.zip(&rejected_keys)
.zip(records)
.zip(rejections)
{
if let Some(bytes) = record? {
let mut record = Self::decode(&bytes)?;
record.observe_again(observed_at, &[]);
batch.put(key, Self::encode(&record)?);
if record.searchable {
observed_groups
.entry(record.content_key)
.and_modify(|last_seen: &mut u64| {
*last_seen = (*last_seen).max(record.last_seen);
})
.or_insert(record.last_seen);
}
updated += 1;
continue;
}
if let Some(bytes) = rejection? {
let mut rejection = Self::decode_rejection(&bytes)?;
if rejection.rule_id == self.rejection_rule_id {
rejection.observe_again(observed_at);
batch.put(rejected_key, Self::encode_rejection(&rejection)?);
updated += 1;
continue;
}
}
unknown.push(*info_hash);
}
if updated > 0 {
let mut delta = InventoryDelta::default();
let mut refreshes = Vec::with_capacity(observed_groups.len());
for (content_key, last_seen) in observed_groups {
let refresh =
self.schedule_observation_refresh(&mut batch, &content_key, last_seen)?;
if refresh == RefreshDecision::Scheduled {
let group_delta = self.dirty_group(&mut batch, &content_key, false)?;
delta.searchable_groups = delta
.searchable_groups
.saturating_add(group_delta.searchable_groups);
delta.pending_documents = delta
.pending_documents
.saturating_add(group_delta.pending_documents);
}
refreshes.push(refresh);
}
self.commit_inventory_batch(batch, delta)?;
for refresh in refreshes {
self.record_refresh_decision(refresh);
}
}
Ok(unknown)
}
fn pending_index(&self, limit: usize) -> Result<Vec<ContentGroupTask>, StorageError> {
if limit == 0 {
return Ok(Vec::new());
}
let prefix = pending_index_prefix();
let iterator = self
.db
.iterator(IteratorMode::From(&prefix, Direction::Forward));
let mut tasks = Vec::with_capacity(limit.min(1024));
for entry in iterator {
let (key, value) = entry?;
let Some(content_key) = decode_pending_content_key(&key) else {
break;
};
let (revision, force_write) = decode_pending_revision(&value)?;
tasks.push(ContentGroupTask {
content_key,
revision,
force_write,
});
if tasks.len() == limit {
break;
}
}
Ok(tasks)
}
fn content_group(
&self,
content_key: &[u8; 32],
now: u64,
) -> Result<Option<crate::domain::ContentGroup>, StorageError> {
let filter = self.current_filter();
self.content_group_with_filter(content_key, now, &filter)
}
fn index_document(
&self,
content_key: &[u8; 32],
now: u64,
) -> Result<IndexDocument, StorageError> {
let group = self.content_group(content_key, now)?;
let projection_hash = projection_hash(group.as_ref());
let dynamic_projection_hash = dynamic_projection_hash(group.as_ref());
let bucket = group.as_ref().map_or_else(
|| refresh_bucket(now),
|group| refresh_bucket(group.last_seen),
);
let previous = self.projection_state(content_key)?;
let requires_write = previous.is_none_or(|state| {
state.hash != projection_hash
|| state.dynamic_hash != dynamic_projection_hash
|| state.visible != group.is_some()
});
Ok(IndexDocument {
content_key: *content_key,
projection_hash,
dynamic_projection_hash,
refresh_bucket: bucket,
requires_write,
group,
})
}
fn mark_indexed(
&self,
content_key: &[u8; 32],
revision: u64,
projection_hash: [u8; 32],
dynamic_projection_hash: [u8; 32],
refresh_bucket: u64,
visible: bool,
) -> Result<bool, StorageError> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(state) = self.group_state(content_key)? else {
return Err(StorageError::MissingContentGroup(hex::encode(content_key)));
};
if state.revision != revision {
return Ok(false);
}
let mut batch = WriteBatch::default();
batch.put(content_group_key(content_key), Self::encode_group(state)?);
batch.delete(pending_index_key(content_key));
batch.delete(filter_pending_key(content_key));
batch.put(
filter_projection_key(content_key),
encode_projection(ProjectionState {
hash: projection_hash,
dynamic_hash: dynamic_projection_hash,
refresh_bucket,
visible,
}),
);
self.commit_inventory_batch(
batch,
InventoryDelta {
pending_documents: -1,
..InventoryDelta::default()
},
)?;
Ok(true)
}
fn prepare_full_reindex(&self) -> Result<u64, StorageError> {
const BATCH_SIZE: usize = 1_000;
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let prefix = content_group_prefix();
let iterator = self
.db
.iterator(IteratorMode::From(&prefix, Direction::Forward));
let mut batch = WriteBatch::default();
let mut batch_len = 0_usize;
let mut total = 0_u64;
for entry in iterator {
let (key, value) = entry?;
if !key.starts_with(&prefix) {
break;
}
let content_key: [u8; 32] = key[1..]
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?;
let state: ContentGroupState = rmp_serde::from_slice(&value)?;
batch.put(
pending_index_key(&content_key),
encode_forced_pending_revision(state.revision),
);
batch_len += 1;
total += 1;
if batch_len == BATCH_SIZE {
self.db.write(batch)?;
batch = WriteBatch::default();
batch_len = 0;
}
}
if batch_len > 0 {
self.db.write(batch)?;
}
let inventory = self.inventory_snapshot();
let mut batch = WriteBatch::default();
batch.put(
PENDING_INDEX_COUNT_KEY,
inventory.searchable_groups.to_be_bytes(),
);
self.db.write(batch)?;
self.inventory
.pending_documents
.store(inventory.searchable_groups, Ordering::Relaxed);
Ok(total)
}
fn content_variants(
&self,
content_key: &[u8; 32],
offset: usize,
limit: usize,
) -> Result<ContentVariants, StorageError> {
let filter = self.current_filter();
let prefix = content_members_prefix(content_key);
let iterator = self
.db
.iterator(IteratorMode::From(&prefix, Direction::Forward));
let mut records = Vec::with_capacity(limit);
let mut total = 0_u64;
for entry in iterator {
let (key, _) = entry?;
let Some(info_hash) = decode_content_member_info_hash(&key, content_key) else {
break;
};
if let Some(record) = self
.get(info_hash)?
.and_then(|record| filter.public_record(&record))
{
if total >= offset as u64 && records.len() < limit {
records.push(record);
}
total = total.saturating_add(1);
}
}
Ok(ContentVariants { total, records })
}
fn enqueue_verification(
&self,
info_hashes: &[InfoHash],
priority: VerificationPriority,
requested_at: u64,
capacity: usize,
) -> Result<VerificationEnqueueOutcome, StorageError> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut count = self.verification_queue_len_inner()?;
let mut outcome = VerificationEnqueueOutcome::default();
let mut seen = std::collections::HashSet::with_capacity(info_hashes.len());
for info_hash in info_hashes {
if !seen.insert(*info_hash) {
outcome.deduplicated += 1;
continue;
}
let Some(record) = self.get(*info_hash)? else {
outcome.deduplicated += 1;
continue;
};
if record.availability.next_check_at > requested_at {
outcome.deduplicated += 1;
continue;
}
let locator_key = verification_locator_key(*info_hash);
if let Some(existing_key) = self.db.get(locator_key)? {
let can_promote = priority == VerificationPriority::High
&& decode_verification_task(&existing_key).is_some_and(|(high, _, _)| !high);
if can_promote {
let new_key = verification_task_key(true, requested_at, *info_hash);
let mut batch = WriteBatch::default();
batch.delete(existing_key);
batch.put(new_key, []);
batch.put(locator_key, new_key);
self.db.write(batch)?;
}
outcome.deduplicated += 1;
continue;
}
if count >= capacity {
if priority == VerificationPriority::High {
if let Some((evicted_key, evicted_hash)) =
self.first_verification_task(false)?
{
let mut batch = WriteBatch::default();
batch.delete(evicted_key);
batch.delete(verification_locator_key(evicted_hash));
count = count.saturating_sub(1);
batch.put(VERIFICATION_QUEUE_COUNT_KEY, (count as u64).to_be_bytes());
self.db.write(batch)?;
} else {
outcome.rejected_full += 1;
continue;
}
} else {
outcome.rejected_full += 1;
continue;
}
}
let task_key = verification_task_key(
priority == VerificationPriority::High,
requested_at,
*info_hash,
);
let mut batch = WriteBatch::default();
batch.put(task_key, []);
batch.put(locator_key, task_key);
count += 1;
batch.put(VERIFICATION_QUEUE_COUNT_KEY, (count as u64).to_be_bytes());
self.db.write(batch)?;
outcome.accepted += 1;
}
Ok(outcome)
}
fn claim_verification(
&self,
now: u64,
lease_secs: u64,
) -> Result<Option<VerificationRequest>, StorageError> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let lease_prefix = verification_lease_prefix();
let iterator = self
.db
.iterator(IteratorMode::From(&lease_prefix, Direction::Forward));
let mut recovery = WriteBatch::default();
let mut recovered = false;
for entry in iterator {
let (key, value) = entry?;
let Some((lease_until, info_hash)) = decode_verification_lease(&key) else {
break;
};
if lease_until > now {
break;
}
let high = value.first().copied() == Some(1);
let task_key = verification_task_key(high, now, info_hash);
recovery.delete(&key);
recovery.put(task_key, []);
recovery.put(verification_locator_key(info_hash), task_key);
recovered = true;
}
if recovered {
self.db.write(recovery)?;
}
let selected = match self.first_verification_task(true)? {
Some((key, hash)) => Some((VerificationPriority::High, key, hash)),
None => self
.first_verification_task(false)?
.map(|(key, hash)| (VerificationPriority::Normal, key, hash)),
};
let Some((priority, task_key, info_hash)) = selected else {
return Ok(None);
};
let lease_key = verification_lease_key(now.saturating_add(lease_secs), info_hash);
let mut batch = WriteBatch::default();
batch.delete(task_key);
batch.put(
lease_key,
[u8::from(priority == VerificationPriority::High)],
);
batch.put(verification_locator_key(info_hash), lease_key);
self.db.write(batch)?;
Ok(Some(VerificationRequest { info_hash }))
}
fn finish_verification(
&self,
info_hash: InfoHash,
result: VerificationResult,
) -> Result<(), StorageError> {
let _guard = self
.write_lock
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let Some(mut record) = self.get(info_hash)? else {
return Err(StorageError::MissingRecord(info_hash));
};
record.apply_verification(result);
let locator_key = verification_locator_key(info_hash);
let queued_key = self.db.get(locator_key)?;
let count = self.verification_queue_len_inner()?.saturating_sub(1);
let mut batch = WriteBatch::default();
batch.put(torrent_key(info_hash), Self::encode(&record)?);
let delta = if record.searchable {
self.dirty_group(&mut batch, &record.content_key, false)?
} else {
InventoryDelta::default()
};
if let Some(queued_key) = queued_key {
batch.delete(queued_key);
}
batch.delete(locator_key);
batch.put(VERIFICATION_QUEUE_COUNT_KEY, (count as u64).to_be_bytes());
self.commit_inventory_batch(batch, delta)?;
Ok(())
}
fn verification_queue_len(&self) -> Result<usize, StorageError> {
self.verification_queue_len_inner()
}
fn index_inventory(&self) -> IndexInventory {
self.inventory_snapshot()
}
fn index_refresh_diagnostics(&self) -> IndexRefreshDiagnostics {
IndexRefreshDiagnostics {
scheduled: self.index_refresh_scheduled.load(Ordering::Relaxed),
suppressed: self.index_refresh_suppressed.load(Ordering::Relaxed),
}
}
}
fn encode_projection(state: ProjectionState) -> [u8; 73] {
let mut bytes = [0_u8; 73];
bytes[0] = u8::from(state.visible);
bytes[1..33].copy_from_slice(&state.hash);
bytes[33..65].copy_from_slice(&state.dynamic_hash);
bytes[65..].copy_from_slice(&state.refresh_bucket.to_be_bytes());
bytes
}
fn decode_projection(bytes: &[u8]) -> Result<ProjectionState, StorageError> {
if !matches!(bytes.len(), 33 | 73) || bytes[0] > 1 {
return Err(StorageError::CorruptContentGroup);
}
Ok(ProjectionState {
visible: bytes[0] == 1,
hash: bytes[1..33]
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
dynamic_hash: if bytes.len() == 73 {
bytes[33..65]
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?
} else {
[0; 32]
},
refresh_bucket: if bytes.len() == 73 {
u64::from_be_bytes(
bytes[65..]
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
)
} else {
0
},
})
}
fn read_counter(db: &DB, key: &[u8]) -> Result<Option<u64>, StorageError> {
db.get(key)?
.map(|value| {
value
.as_slice()
.try_into()
.map(u64::from_be_bytes)
.map_err(|_| StorageError::CorruptContentGroup)
})
.transpose()
}
fn encode_forced_pending_revision(revision: u64) -> [u8; 9] {
let mut bytes = [0_u8; 9];
bytes[0] = 1;
bytes[1..].copy_from_slice(&revision.to_be_bytes());
bytes
}
fn decode_pending_revision(bytes: &[u8]) -> Result<(u64, bool), StorageError> {
match bytes {
bytes if bytes.len() == 8 => Ok((
u64::from_be_bytes(
bytes
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
),
false,
)),
[1, revision @ ..] if revision.len() == 8 => Ok((
u64::from_be_bytes(
revision
.try_into()
.map_err(|_| StorageError::CorruptContentGroup)?,
),
true,
)),
_ => Err(StorageError::CorruptContentGroup),
}
}
fn count_prefix<const N: usize>(db: &DB, prefix: [u8; N]) -> Result<u64, StorageError> {
let iterator = db.iterator(IteratorMode::From(&prefix, Direction::Forward));
let mut count = 0_u64;
for entry in iterator {
let (key, _) = entry?;
if !key.starts_with(&prefix) {
break;
}
count = count.saturating_add(1);
}
Ok(count)
}
#[cfg(test)]
mod tests;