issue/43 Small fixes.

This commit is contained in:
Paul Masurel
2017-01-23 09:37:31 +09:00
parent 926e71a573
commit 6530d43d6a
7 changed files with 39 additions and 55 deletions
+4 -3
View File
@@ -176,11 +176,12 @@ impl Index {
&mut *self.directory
}
/// Reads the meta.json and returns the list of
/// committed segments.
pub fn committed_segments(&self) -> Result<Vec<SegmentMeta>> {
Ok(load_metas(self.directory())?
.committed_segments)
Ok(load_metas(self.directory())?.committed_segments)
}
/// Returns the list of segment ids that are searchable.
pub fn searchable_segment_ids(&self) -> Result<Vec<SegmentId>> {
self.committed_segments()
+10 -10
View File
@@ -20,16 +20,6 @@ pub struct Searcher {
segment_readers: Vec<SegmentReader>,
}
impl fmt::Debug for Searcher {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
let segment_ids = self.segment_readers
.iter()
.map(|segment_reader| segment_reader.segment_id())
.collect::<Vec<_>>();
write!(f, "Searcher({:?})", segment_ids);
Ok(())
}
}
impl Searcher {
@@ -94,4 +84,14 @@ impl From<Vec<SegmentReader>> for Searcher {
segment_readers: segment_readers,
}
}
}
impl fmt::Debug for Searcher {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
let segment_ids = self.segment_readers
.iter()
.map(|segment_reader| segment_reader.segment_id())
.collect::<Vec<_>>();
write!(f, "Searcher({:?})", segment_ids)
}
}
+1
View File
@@ -238,6 +238,7 @@ impl SegmentReader {
self.term_infos.get(term.as_slice())
}
/// Returns the segment id
pub fn segment_id(&self) -> SegmentId {
self.segment_id
}
+1 -1
View File
@@ -134,7 +134,7 @@ fn index_documents(heap: &mut Heap,
let segment_entry = SegmentEntry::new(segment_meta, delete_cursor.clone());
try!(segment_writer.finalize());
let future = segment_update_manager.send(SegmentUpdate::AddSegment(segment_entry));
segment_update_manager.send(SegmentUpdate::AddSegment(segment_entry));
Ok(())
}
+1 -11
View File
@@ -82,11 +82,6 @@ impl SegmentManager {
}
}
pub fn docstamp(&self,) -> u64 {
self.read().docstamp
}
pub fn from_segments(segment_metas: Vec<SegmentMeta>, delete_cursor: DeleteQueueCursor) -> SegmentManager {
SegmentManager {
registers: RwLock::new( SegmentRegisters {
@@ -169,12 +164,7 @@ impl SegmentManager {
warn!("couldn't find segment in SegmentManager");
}
}
pub fn committed_segments(&self,) -> Vec<SegmentId> {
let registers_lock = self.read();
registers_lock.committed.segment_ids()
}
pub fn segment_metas(&self,) -> (Vec<SegmentMeta>, Vec<SegmentMeta>) {
let registers_lock = self.read();
(registers_lock.committed.segment_metas(), registers_lock.uncommitted.segment_metas())
+22 -24
View File
@@ -79,10 +79,9 @@ pub fn save_metas(segment_manager: &SegmentManager,
let metas = create_metas(segment_manager, schema, docstamp);
let mut w = Vec::new();
try!(write!(&mut w, "{}\n", json::as_pretty_json(&metas)));
directory
.atomic_write(&META_FILEPATH, &w[..])
.map_err(From::from)
Ok(directory
.atomic_write(&META_FILEPATH, &w[..])?)
}
#[derive(Clone, Debug)]
@@ -139,17 +138,16 @@ impl SegmentUpdateManager {
let segment_update_manager = SegmentUpdateManager {
channel: segment_update_sender,
};
let segment_updater = SegmentUpdater::new(
SegmentUpdateRunner::new(
index,
segment_manager,
merge_policy,
segment_update_manager.clone(),
segment_update_receiver);
segment_updater.start();
segment_update_receiver).start();
segment_update_manager
}
pub fn send(&self, segment_update: SegmentUpdate) -> Future<(), &'static str> {
pub fn send(&self, segment_update: SegmentUpdate) -> impl Async<Value=(), Error=&'static str> {
let (fullfiller, future) = Future::<(), &'static str>::pair();
self.channel.send((fullfiller, segment_update));
future
@@ -158,8 +156,8 @@ impl SegmentUpdateManager {
}
/// The segment updater is in charge of processing all of the
/// `SegmentUpdate`s.
/// The segment update runner is in charge of processing all
/// of the `SegmentUpdate`s.
///
/// All this processing happens on a single thread
/// consuming a common queue.
@@ -168,7 +166,7 @@ impl SegmentUpdateManager {
/// - indexing threads are sending new segments
/// - merging threads are sending merge operations
/// - the index writer sends "terminate"
pub struct SegmentUpdater {
pub struct SegmentUpdateRunner {
index: Index,
is_cancelled_generation: bool,
segment_update_receiver: SegmentUpdateReceiver,
@@ -179,14 +177,14 @@ pub struct SegmentUpdater {
merging_threads: HashMap<usize, JoinHandle<(Vec<SegmentId>, SegmentEntry)> >,
}
impl SegmentUpdater {
impl SegmentUpdateRunner {
fn new(index: Index,
segment_manager: Arc<SegmentManager>,
merge_policy: Arc<Mutex<Box<MergePolicy>>>,
segment_update_manager: SegmentUpdateManager,
segment_update_receiver: SegmentUpdateReceiver) -> SegmentUpdater {
SegmentUpdater {
segment_update_receiver: SegmentUpdateReceiver) -> SegmentUpdateRunner {
SegmentUpdateRunner {
index: index,
is_cancelled_generation: false,
segment_update_manager: segment_update_manager,
@@ -259,7 +257,6 @@ impl SegmentUpdater {
num_deleted_docs: 0u32,
};
// TODO fix delete cursor
let delete_queue = DeleteQueue::default();
@@ -322,7 +319,7 @@ impl SegmentUpdater {
let generation_before_update = segment_manager.generation();
self.process_one(segment_update);
if generation_before_update != segment_manager.generation() {
// The segment manager has changed, we need to
// - save meta.json
@@ -382,11 +379,12 @@ impl SegmentUpdater {
pub fn process_one(
&mut self,
segment_update: SegmentUpdate) {
info!("Segment update: {:?}", segment_update);
info!("Segment update: {:?}", segment_update);
use self::SegmentUpdate::*;
match segment_update {
SegmentUpdate::AddSegment(segment_entry) => {
AddSegment(segment_entry) => {
if !self.is_cancelled_generation {
self.segment_manager().add_segment(segment_entry);
}
@@ -399,7 +397,7 @@ impl SegmentUpdater {
self.index.delete_segment(segment_entry.segment_id());
}
}
SegmentUpdate::EndMerge(merging_thread_id_opt, segment_ids, segment_entry) => {
EndMerge(merging_thread_id_opt, segment_ids, segment_entry) => {
self.end_merge(
segment_ids,
segment_entry);
@@ -407,21 +405,21 @@ impl SegmentUpdater {
self.merging_threads.remove(&merging_thread_id);
}
}
SegmentUpdate::CancelGeneration => {
CancelGeneration => {
// Called during rollback. The segment
// that will arrive will be ignored
// until a NewGeneration is update arrives.
self.is_cancelled_generation = true;
}
SegmentUpdate::NewGeneration => {
NewGeneration => {
// After rollback, we can resume
// indexing new documents.
self.is_cancelled_generation = false;
}
SegmentUpdate::Commit(docstamp) => {
Commit(docstamp) => {
self.segment_manager().commit(docstamp);
}
SegmentUpdate::Terminate => {
Terminate => {
panic!("We should have left the loop before processing it.");
}
}
-6
View File
@@ -174,12 +174,6 @@ impl<'a> SegmentWriter<'a> {
// .collect::<Vec<_>>()
// }
fn first_opstamp(&self) -> u64 {
*(self.doc_opstamps
.first()
.expect("Last doc opstamp called on an empty segment writer"))
}
pub fn last_opstamp(&self) -> u64 {
*(self.doc_opstamps
.last()