diff --git a/src/core/index.rs b/src/core/index.rs index 417bbebd5..05296188a 100644 --- a/src/core/index.rs +++ b/src/core/index.rs @@ -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> { - 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> { self.committed_segments() diff --git a/src/core/searcher.rs b/src/core/searcher.rs index 0ea6cf840..839e00172 100644 --- a/src/core/searcher.rs +++ b/src/core/searcher.rs @@ -20,16 +20,6 @@ pub struct Searcher { segment_readers: Vec, } -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::>(); - write!(f, "Searcher({:?})", segment_ids); - Ok(()) - } -} impl Searcher { @@ -94,4 +84,14 @@ impl From> 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::>(); + write!(f, "Searcher({:?})", segment_ids) + } } \ No newline at end of file diff --git a/src/core/segment_reader.rs b/src/core/segment_reader.rs index 15b873a79..d8518f87f 100644 --- a/src/core/segment_reader.rs +++ b/src/core/segment_reader.rs @@ -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 } diff --git a/src/indexer/index_writer.rs b/src/indexer/index_writer.rs index 567cf5f4b..c9c424255 100644 --- a/src/indexer/index_writer.rs +++ b/src/indexer/index_writer.rs @@ -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(()) } diff --git a/src/indexer/segment_manager.rs b/src/indexer/segment_manager.rs index 3b9ad78af..6fcdb253d 100644 --- a/src/indexer/segment_manager.rs +++ b/src/indexer/segment_manager.rs @@ -82,11 +82,6 @@ impl SegmentManager { } } - pub fn docstamp(&self,) -> u64 { - self.read().docstamp - } - - pub fn from_segments(segment_metas: Vec, 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 { - let registers_lock = self.read(); - registers_lock.committed.segment_ids() - } - + pub fn segment_metas(&self,) -> (Vec, Vec) { let registers_lock = self.read(); (registers_lock.committed.segment_metas(), registers_lock.uncommitted.segment_metas()) diff --git a/src/indexer/segment_updater.rs b/src/indexer/segment_updater.rs index a5b5e4ec3..060044f95 100644 --- a/src/indexer/segment_updater.rs +++ b/src/indexer/segment_updater.rs @@ -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 { 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, SegmentEntry)> >, } -impl SegmentUpdater { +impl SegmentUpdateRunner { fn new(index: Index, segment_manager: Arc, merge_policy: Arc>>, 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."); } } diff --git a/src/indexer/segment_writer.rs b/src/indexer/segment_writer.rs index 074f6e7d1..204eb9a37 100644 --- a/src/indexer/segment_writer.rs +++ b/src/indexer/segment_writer.rs @@ -174,12 +174,6 @@ impl<'a> SegmentWriter<'a> { // .collect::>() // } - 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()