From aaa6e5e2f998461709e2cef4849335d88385e9e1 Mon Sep 17 00:00:00 2001 From: Grigorii Kirgizov Date: Thu, 14 Jan 2021 21:07:40 +0300 Subject: [PATCH] Refactoring: simplify RewindStage & related logic (MPSCR-21). Add docs, some renames. ProcessingSession -> SessionManager RewindVolatileOccurrencesStage -> RewindStage --- .../core/internal/ConstraintsProcessing.kt | 3 +- .../core/internal/EvaluationSessionImpl.kt | 55 +++++- .../reactor/core/internal/IncrementalStage.kt | 146 ++++++--------- .../core/internal/ProcessingStrategy.kt | 172 +++++++++--------- 4 files changed, 184 insertions(+), 192 deletions(-) diff --git a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ConstraintsProcessing.kt b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ConstraintsProcessing.kt index b865a8bc..d47e14b3 100644 --- a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ConstraintsProcessing.kt +++ b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ConstraintsProcessing.kt @@ -44,9 +44,10 @@ internal class ConstraintsProcessing( val trace: EvaluationTrace = EvaluationTrace.NULL, val profiler: Profiler? = null + // fixme: can get rid of inheritance from journal, use composition instead ) : StoreAwareJournalImpl(journal, logicalState), IncrSpecHolder { - private var incrementalProcessing: ProcessingStrategy = DefaultProcessing() + private var incrementalProcessing: ProcessingStrategy = EmptyProcessing() fun setStrategy(strategy: ProcessingStrategy) { this.incrementalProcessing = strategy } diff --git a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/EvaluationSessionImpl.kt b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/EvaluationSessionImpl.kt index 16cf984d..df3bdcb7 100644 --- a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/EvaluationSessionImpl.kt +++ b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/EvaluationSessionImpl.kt @@ -33,21 +33,37 @@ internal typealias OccurrenceStore = Collection internal fun emptyStore(): OccurrenceStore = emptyList() -internal interface ProcessingSession { +/** + * Handles creation of the first and following sessions, + * properly ending sessions and getting their results. + */ +internal interface SessionManager { + /** + * Creates first session when there's no [SessionToken] available. + */ fun firstSession(): SessionParts + /** + * Creates following session when previous [token] is available. + */ fun nextSession(token: SessionToken): SessionParts /** - * Clears state unneeded between sessions and - * returns [SessionToken] with session results. + * Clears state unneeded between sessions and return + * next [SessionToken] required for following session. */ fun endSession(session: SessionParts): SessionToken + /** + * Starts [session] and returns overall [EvaluationResult] of the program. + */ fun runSession(session: SessionParts, main: Constraint): EvaluationResult } +/** + * Bundle of all entities involved in a session. + */ internal data class SessionParts( val preambleInfo: PreambleInfo, val ruleIndex: RuleIndex, @@ -77,7 +93,7 @@ internal class EvaluationSessionImpl private constructor ( (this as? SessionTokenImpl)?.principalObservers?.isNotEmpty() ?: false private fun launch(token: SessionToken?, store: OccurrenceStore, main: Constraint): EvaluationResult { - val sessionProcessing: ProcessingSession = + val sessionProcessing: SessionManager = with(incrementality) { when { ability().allowed() -> when { @@ -114,7 +130,7 @@ internal class EvaluationSessionImpl private constructor ( } } - open inner class DefaultProcessingSession: ProcessingSession { + open inner class DefaultProcessingSession: SessionManager { override fun firstSession(): SessionParts = getSession(null) @@ -130,15 +146,13 @@ internal class EvaluationSessionImpl private constructor ( val logicalState = LogicalState() val dispatchingFront = Dispatcher(ruleIndex).front() - val principalObservers = LogicalBindObserverDispatcher() - val processingStrategy = GroundProcessing(incrementality, principalObservers) - + val processingStrategy = EmptyProcessing() val processing = ConstraintsProcessing(dispatchingFront, journal, logicalState, incrementality, trace, profiler) processing.setStrategy(processingStrategy) val controller = ControllerImpl(supervisor, processing, incrementality, trace, profiler) - return SessionParts(program.preambleInfo(), ruleIndex, journal, logicalState, controller, processing, processingStrategy, principalObservers) + return SessionParts(program.preambleInfo(), ruleIndex, journal, logicalState, controller, processing, processingStrategy, PrincipalObserverDispatcher.EMPTY) } override fun endSession(session: SessionParts): SessionToken = with(session) { @@ -208,6 +222,28 @@ internal class EvaluationSessionImpl private constructor ( } open inner class IncrementalProcessingSession(): DefaultProcessingSession() { + + /** + * Same as [DefaultProcessingSession.firstSession], + * but uses [GroundProcessing] instead of [EmptyProcessing] + */ + override fun firstSession(): SessionParts { + val ruleIndex = RuleIndex(program.rules()) + val journal = MatchJournalImpl(incrementality) + val logicalState = LogicalState() + val dispatchingFront = Dispatcher(ruleIndex).front() + + val principalObservers = LogicalBindObserverDispatcher() + val processingStrategy = GroundProcessing(incrementality, principalObservers) + + val processing = ConstraintsProcessing(dispatchingFront, journal, logicalState, incrementality, trace, profiler) + processing.setStrategy(processingStrategy) + + val controller = ControllerImpl(supervisor, processing, incrementality, trace, profiler) + + return SessionParts(program.preambleInfo(), ruleIndex, journal, logicalState, controller, processing, processingStrategy, principalObservers) + } + override fun nextSession(token: SessionToken): SessionParts { val tkn = token as SessionTokenImpl val logicalState = tkn.logicalState @@ -273,7 +309,6 @@ internal class EvaluationSessionImpl private constructor ( override fun endSession(session: SessionParts): SessionToken = with(session) { val histView = journal.view() val outputOccurrences = histView.filterOccurrences(inputStore) - // todo: make output store in other processing strategies? outputOccurrences.forEach{ it.terminate(logicalState) } processing.resetStore() // clear observers diff --git a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/IncrementalStage.kt b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/IncrementalStage.kt index bf8709c8..ff7276d5 100644 --- a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/IncrementalStage.kt +++ b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/IncrementalStage.kt @@ -31,20 +31,12 @@ import java.util.PriorityQueue */ internal interface IncrementalStage: IncrSpecHolder { -// fun onNext(reader: MatchJournal.ChunkReader) +// fun receive(data: Iterable): Boolean + +// fun onNext(reader: MatchJournal.ChunkReader): Collection } -/** - * Used for pruning invalid internal state in impls of [IncrementalStage]. - */ -internal interface InvalidatingInfo { - fun isInvalid(entity: Justified): Boolean - - companion object Empty: InvalidatingInfo { - override fun isInvalid(entity: Justified): Boolean = false - } -} /** * Invalidation stage includes several activities: @@ -61,7 +53,7 @@ internal class InvalidationStage( private val activationSink: ContinuedActivationSink, private val stateCleaner: ConstraintsProcessing.ProgramStateCleaner, private val trace: EvaluationTrace -): IncrementalStage, InvalidatingInfo { +): IncrementalStage { private val invalidJustifications = mutableListOf() @@ -74,9 +66,6 @@ internal class InvalidationStage( fun invalidatedRules(): List = invalidRuleIdsAll.toList() - override fun isInvalid(entity: Justified): Boolean = - entity.justifiedByAny(invalidJustifications) - fun receive(invalid: Iterable) { invalidJustifications.addAll(invalid) } fun receive(invalid: Justified) { invalidJustifications.add(invalid) } @@ -153,9 +142,9 @@ internal class InvalidationStage( */ internal class AdditionStage( override val ispec: IncrementalSpec, + private val posTracker: PosTracking, private val addedRules: Iterable, private val activationSink: ContinuedActivationSink, - private val invalidatingInfo: InvalidatingInfo, private val ruleOrdering: RuleOrdering, private val ruleIndex: RuleIndex, private val trace: EvaluationTrace @@ -187,6 +176,12 @@ internal class AdditionStage( offerCandidates(reader) } + fun onRewind(reader: ChunkReader) = with(posTracker) { + activationCandidates.removeIf { + isNew(it.occChunk) || isFuture(it.occChunk) + } + } + fun receive(candidates: Iterable) = activationCandidates.addAll(candidates) @@ -216,12 +211,6 @@ internal class AdditionStage( val candidateRule = candidate.rule val occChunk = candidate.occChunk - // Relevant for when rewind happenned - if (invalidatingInfo.isInvalid(occChunk)) { - aIt.remove() - continue - } - val pos = if (ruleOrdering.canBeInserted(candidateRule, occChunk, reader.next) || reader.atEnd()) reader.current.toPos() @@ -240,7 +229,7 @@ internal class AdditionStage( // todo; extract more restricted interface (w/o rewind) for use in stages internal class PosTracking( - private val journalIndex: MatchJournal.Index, + val index: MatchJournal.Index, private val journal: MatchJournal, initPos: MatchJournal.Pos = journal.initialChunk().toPos() ) { @@ -260,12 +249,12 @@ internal class PosTracking( private fun updateReferencePos(current: MatchJournal.Chunk) { - if (journalIndex.isKnown(current)) { + if (index.isKnown(current)) { lastVisited = current.toPos() } } - fun rewind(pos: MatchJournal.Pos) = with(journalIndex) { + fun rewind(pos: MatchJournal.Pos) = with(index) { assert(isKnown(pos.chunk)) assert(pos before lastVisited) @@ -285,31 +274,32 @@ internal class PosTracking( fun isOld(chunk: MatchJournal.Chunk): Boolean = - journalIndex.isKnown(chunk) + index.isKnown(chunk) fun isOld(occ: Occurrence): Boolean = - journalIndex.isKnown(occ) + index.isKnown(occ) fun isNew(chunk: MatchJournal.Chunk): Boolean = - !journalIndex.isKnown(chunk) + !index.isKnown(chunk) fun isNew(occ: Occurrence): Boolean = - !journalIndex.isKnown(occ) + !index.isKnown(occ) fun isFront(): Boolean = - with(journalIndex) { lastVisited afterOrEq front } + with(index) { lastVisited afterOrEq front } fun isFuture(pos: MatchJournal.Pos): Boolean = - with(journalIndex) { pos after lastVisited } + with(index) { pos after lastVisited } + fun isFuture(chunk: MatchJournal.Chunk) = isFuture(chunk.toPos()) fun isPast(pos: MatchJournal.Pos): Boolean = - with(journalIndex) { pos before lastVisited } + with(index) { pos before lastVisited } + fun isPast(chunk: MatchJournal.Chunk) = isPast(chunk.toPos()) } internal class PostponeMatchesStage( override val ispec: IncrementalSpec, private val posTracker: PosTracking, - private val journalIndex: MatchJournal.Index, private val ruleOrdering: RuleOrdering ): IncrementalStage { @@ -353,7 +343,7 @@ internal class PostponeMatchesStage( for (m in matches) { // Returns null for matches with occurrences only from this session // because journalIndex indexes only previous session. - val occChunk = journalIndex.activationPos(m) + val occChunk = posTracker.index.activationPos(m) if (occChunk != null && posTracker.isFuture(occChunk.toPos())) { postponedMatches.getOrPut(occChunk.identity, ::mutableListOf).add(m) @@ -380,7 +370,6 @@ internal class PostponeMatchesStage( internal class ContinueOccurrencesStage( override val ispec: IncrementalSpec, - private val invalidatingInfo: InvalidatingInfo, private val journalIndex: MatchJournal.Index ): IncrementalStage, ContinuedActivationSink { @@ -416,28 +405,23 @@ internal class ContinueOccurrencesStage( private val seen: MutableSet = HashSet() - private fun ExecPos.isInvalid() = with(invalidatingInfo) { - isInvalid(reactivated) - } - - fun onNext(reader: ChunkReader): Collection { - val continued = mutableListOf() - while (queue.isNotEmpty()) { - if (queue.top().isInvalid()) { - queue.pop() - continue - } - if (!(reader at queue.top())) break - - continued.add( queue.pop().reactivatedOcc ) + fun onNext(reader: ChunkReader) = mutableListOf().apply { + while (queue.isNotEmpty() && reader at queue.top()) { + add(queue.pop().reactivatedOcc) } - return continued } - fun onRewind(reader: ChunkReader) = with(journalIndex) { - val nextPos = reader.next.toPos() - seen.removeIf { seenPos -> - isNew(seenPos.continueFrom.chunk) || seenPos.continueFrom afterOrEq nextPos + fun onRewind(reader: ChunkReader): Unit = with(journalIndex) { + val rewindPos = reader.next.toPos() + + // clear queue + queue.removeIf { + it.reactivated.toPos() afterOrEq rewindPos + } + + // clear seen + seen.removeIf { + isNew(it.continueFrom.chunk) || it.continueFrom afterOrEq rewindPos } } @@ -465,63 +449,41 @@ internal class ContinueOccurrencesStage( } -// fixme: docs -internal class RewindVolatileOccurrencesStage( +internal class RewindStage( override val ispec: IncrementalSpec, private val posTracker: PosTracking, - private val journalIndex: MatchJournal.Index, - private val invalidated: InvalidatingInfo, private val principalObserver: PrincipalObserverDispatcher ): IncrementalStage { - private val toRewind: PriorityQueue = PriorityQueue(journalIndex.chunkComparator) + private val toRewind: PriorityQueue = PriorityQueue(posTracker.index.chunkComparator) private val seen = hashSetOf() + // precondition: all received occurrences are principal and valid (i.e. not invalidated) fun receive(occs: Sequence): Boolean = occs.filter(::takeVolatile) - .mapNotNull(journalIndex::activatingChunkOf) - .mapNotNull(journalIndex::matchChunkOf) - .filter(seen::add) + .mapNotNull(posTracker.index::activatingChunkOf) // get chunk corresponding to occurrence activation + .mapNotNull(posTracker.index::matchChunkOf) // get chunk that activated it + .filter(seen::add) // filter already seen .toList().let(toRewind::addAll) - fun maybeReset(reader: ChunkReader): MatchJournal.Pos? = + /** + * Returns collection of past chunks for rewind + * in sorted order from earliest to latest. + */ + fun onNext(reader: ChunkReader): Collection = if (needRewind()) { - toRewind.peek()!!.toPos() - } else null + toRewind.toList().also { toRewind.clear() } + } else emptyList() - fun reevaluate(reader: ChunkReader): Iterable { - val reevaluated = mutableListOf() - while (toRewind.isNotEmpty()) { - if (reader atNext toRewind.peek()!!) { - reevaluated.add(toRewind.remove()!!) - } else break - } - return reevaluated - } - - fun needRewind(): Boolean = peek()?.let { - posTracker.isPast(it.toPos()) + fun needRewind(): Boolean = toRewind.peek()?.let { + posTracker.isPast(it) } ?: false - /** - * Checked peek() for [toRewind] which prunes stale positions. - */ - private fun peek(): MatchJournal.MatchChunk? { - while (toRewind.isNotEmpty()) { - val top = toRewind.peek()!! - if (invalidated.isInvalid(top)) { - toRewind.remove() - continue - } - return top - } - return null - } private fun takeVolatile(occ: Occurrence) = // NB: only occurrences from previous session are considered - journalIndex.isKnown(occ) && + posTracker.isOld(occ) && with(principalObserver) { removeTriggered(occ) } } diff --git a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ProcessingStrategy.kt b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ProcessingStrategy.kt index 1a869477..97996142 100644 --- a/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ProcessingStrategy.kt +++ b/reactor/Core/src/jetbrains/mps/logic/reactor/core/internal/ProcessingStrategy.kt @@ -25,18 +25,37 @@ import jetbrains.mps.logic.reactor.program.Rule /** * Facade interface for incremental processing. - * Designates points in program evaluation where incremental processing - * must be injected. It also provides an entry point for evaluation. + * + * Specifies points in program evaluation where + * incremental processing must be must be injected: + * [offerMatch], [processMatch], [processOccurrenceMatches], + * [processActivated], [processInvalidated]. + * + * Provides an entry point for evaluation: [run]. + * + * Provides additional session output through + * [invalidatedFeedback] & [invalidatedRules]. */ internal interface ProcessingStrategy { /** - * Information returned back from incremental processing session. + * Output of incremental processing session. + * + * Specifies keys for invalid program feedback + * so that caller could clear it. + * + * @return collection of invalid program feedback keys. */ fun invalidatedFeedback(): FeedbackKeySet /** - * Unique tags of principal rules which matches where invalidated. + * Output of incremental processing session. + * + * Specifies unique tags of principal rules which + * matches were invalidated. From this rule information + * it can be inferred which rule origins were affected. + * + * @return collection of unique tags of affected rules. */ fun invalidatedRules(): List @@ -46,51 +65,75 @@ internal interface ProcessingStrategy { fun run(processing: ConstraintsProcessing, controller: Controller, main: Constraint): FeedbackStatus /** + * Injection point into program evaluation process. + * * Called before [Controller.offerMatch]. - * Can omit match by returning false. + * Allows to omit match (by returning `false`) + * from usual evaluation process and handle it + * in some special manner according to strategy. + * + * @return `true` if match isn't affected by the strategy, + * `false` if it was and must be omitted. */ fun offerMatch(match: RuleMatchEx): Boolean /** - * Pre-process accepted match (i.e. after [Controller.offerMatch]) + * Injection point into program evaluation process. + * + * Pre-process [match] accepted by [Controller.offerMatch] * before general processing in [Controller.processBody]. + * Aimed at processing heads of the [match], while [match] + * itself is evaluated in a usual manner by [Controller]. */ fun processMatch(match: RuleMatchEx) /** - * Processes list of matches on active [Occurrence]. - * Should return a filtered [matches] list. + * Injection point into program evaluation process. + * + * Processes [matches] of activated [Occurrence] + * returned by [Dispatcher.DispatchingFront.matches]. + * Allows to filter [matches] according to the strategy + * e.g. postpone them or drop entirely. + * + * In contrast to [offerMatch], happens earlier and allows + * to process all new matches of [active] at once. + * + * @return list of filtered [matches]. */ fun processOccurrenceMatches(active: Occurrence, matches: List): List /** + * Injection point into program evaluation process. + * * Called on new activated [Occurrence]. * Allows to handle program's logical state, * e.g. add processing-specific observers with [observable]. */ - fun processActivated(active: Occurrence, observable: LogicalStateObservable): Unit + fun processActivated(active: Occurrence, observable: LogicalStateObservable) // fixme: essentially InvalidationStage.invalidateChunk redirects back here -- unnecessary circle /** + * Injection point into program evaluation process. + * * Called for invalidated [Occurrence]s. * It's a pair method for [processActivated] to discharge * its effects, e.g. to clear program logical state. */ - fun processInvalidated(occ: Occurrence, observable: LogicalStateObservable): Unit + fun processInvalidated(occ: Occurrence, observable: LogicalStateObservable) } /** - * Default non-incremental processing with stubs. + * Default non-incremental processing with stubs. Does nothing. */ -internal open class DefaultProcessing: ProcessingStrategy { +internal open class EmptyProcessing: ProcessingStrategy { override fun invalidatedFeedback(): FeedbackKeySet = emptySet() override fun invalidatedRules(): List = emptyList() /** - * Simply redirects evaluation to [Controller]. + * Entry point. Simply redirects evaluation to [Controller]. */ override fun run(processing: ConstraintsProcessing, controller: Controller, main: Constraint) = controller.activate(main) @@ -108,12 +151,15 @@ internal open class DefaultProcessing: ProcessingStrategy { /** - * Processing strategy that observes logical vars and ensures basic incremental contract. + * Strategy that observes logical vars used in principal constraints + * with a help of [PrincipalObserverDispatcher]. + * + * Required for [IncrementalProcessing] for an initial session run. */ internal open class GroundProcessing( override val ispec: IncrementalSpec, private val principalObservers: PrincipalObserverDispatcher = PrincipalObserverDispatcher.EMPTY -): DefaultProcessing(), IncrSpecHolder { +): EmptyProcessing(), IncrSpecHolder { override fun processActivated(active: Occurrence, observable: LogicalStateObservable) { if (active.isPrincipal) { @@ -133,8 +179,9 @@ internal open class GroundProcessing( /** * Facade implementation for incremental processing algorithm. * - * It includes 4 stages that operate on [MatchJournalImpl.Cursor]. + * It includes 5 stages that operate on [MatchJournalImpl.Cursor]. * Stages are: + * - [RewindStage] * - [InvalidationStage] * - [AdditionStage] * - [PostponeMatchesStage] @@ -145,8 +192,7 @@ internal open class GroundProcessing( * general processing in [Controller] & [ConstraintsProcessing]. * * Methods [processMatch] & [processOccurrenceMatches] serve as - * a bridge back from [ConstraintsProcessing] to specific stages - * in [IncrementalProcessing]. + * a bridge back from [ConstraintsProcessing] to specific stages. */ internal class IncrementalProcessing( ispec: IncrementalSpec, @@ -161,28 +207,18 @@ internal class IncrementalProcessing( private val journalIndex = journal.index() private val ruleOrdering = RuleOrdering(ruleIndex) - private val invalidatingInfo = InvalidatingInfoDispatcher() private val posTracker = PosTracking(journalIndex, journal) - private val continuator = ContinueOccurrencesStage(ispec, invalidatingInfo, journalIndex) + private val continuator = ContinueOccurrencesStage(ispec, journalIndex) private val invalidator = InvalidationStage(ispec, posTracker, droppedRules.toSet(), continuator, stateCleaner, trace) - private val adder = AdditionStage(ispec, newRules, continuator, invalidatingInfo, ruleOrdering, ruleIndex, trace) - private val postponer = PostponeMatchesStage(ispec, posTracker, journalIndex, ruleOrdering) - private val rewinder = RewindVolatileOccurrencesStage(ispec, posTracker, journalIndex, invalidatingInfo, principalObservers) + private val adder = AdditionStage(ispec, posTracker, newRules, continuator, ruleOrdering, ruleIndex, trace) + private val postponer = PostponeMatchesStage(ispec, posTracker, ruleOrdering) + private val rewinder = RewindStage(ispec, posTracker, principalObservers) init { - // lateinit to break constructor cycle - invalidatingInfo.src = invalidator - principalObservers.setTriggerReceiver(this::receiveBindTriggered) } - private class InvalidatingInfoDispatcher : InvalidatingInfo { - lateinit var src: InvalidatingInfo - - override fun isInvalid(entity: Justified): Boolean = src.isInvalid(entity) - } - override fun invalidatedFeedback(): FeedbackKeySet = invalidator.invalidatedFeedback() @@ -210,7 +246,7 @@ internal class IncrementalProcessing( while (true) { posTracker.onNext(cursor) - rewinder.rewind(cursor) + rewind(cursor) invalidate(cursor) val postponedMatches = postponer.onNext(cursor) @@ -251,7 +287,7 @@ internal class IncrementalProcessing( postponer.postponeFutureMatches(matches) } else matches - private fun RewindVolatileOccurrencesStage.receiveStrict(occs: Sequence, match: RuleMatchEx) = + private fun RewindStage.receiveStrict(occs: Sequence, match: RuleMatchEx) = rewinder.receive(occs).also { if (ispec.assertLevel() == IncrementalSpec.AssertLevel.AssertContracts) { throw IncrementalContractViolationException( @@ -260,29 +296,29 @@ internal class IncrementalProcessing( } } - private fun RewindVolatileOccurrencesStage.rewind(cursor: ChunkReader): Boolean = - maybeReset(cursor)?.let { rewindPos -> + private fun rewind(cursor: ChunkReader) { + val toRewind = rewinder.onNext(cursor) + + toRewind.firstOrNull()?.let { earliest -> // This is the only place where journal can be reset to past // fixme: modifying not through cursor; not the cleanest way - posTracker.rewind(rewindPos) + posTracker.rewind(earliest.toPos()) - // Some stages on rewind need to take actions (e.g. clear internal state) + // Will invalidate these chunks and reevaluate them + invalidator.receive(toRewind) + + // Some stages on rewind need to clear internal state + adder.onRewind(cursor) continuator.onRewind(cursor) - - reevaluate(cursor).let { - // will invalidate these chunks and reevaluate them - invalidator.receive(it) - } - true - } ?: false - + } + } private fun invalidate(cursor: RemovingJournalIterator): Boolean { var haveChanges = false do { // handle all new chunks that were invalidated by rewind if (posTracker.isNew(cursor.next)) { - invalidator.receive(listOf(cursor.next)) + invalidator.receive(cursor.next) } if (invalidator.onNext(cursor)) { @@ -300,48 +336,6 @@ internal class IncrementalProcessing( } -/** - * Very restricted incremental strategy for processing program preamble. - * - * Includes 2 stages: - * - [AdditionStage] which adds potential matches given occurrences from preamble - * - [ContinueOccurrencesStage] which actually evaluates them - * - * Journal invalidation and injected intermediate processing with - * [processMatch] & [processOccurrenceMatches] are not needed for this. - * So this strategy is very close to the default [DefaultProcessing]. - */ -internal class PreambleProcessing( - ispec: IncrementalSpec, - val journal: MatchJournal, - newRules: Iterable, - ruleIndex: RuleIndex, - trace: EvaluationTrace -): GroundProcessing(ispec) { - - private val journalIndex = journal.index() - private val ruleOrdering = RuleOrdering(ruleIndex) - - private val continuator = ContinueOccurrencesStage(ispec, InvalidatingInfo.Empty, journalIndex) - private val adder = AdditionStage(ispec, newRules, continuator, InvalidatingInfo.Empty, ruleOrdering, ruleIndex, trace) - - override fun run(processing: ConstraintsProcessing, controller: Controller, main: Constraint): FeedbackStatus { - var status: FeedbackStatus = FeedbackStatus.NORMAL() - val cursor = journal.cursor - while (true) { - adder.onNext(cursor) - - status = continuator.runContinued(processing, controller, cursor) - if (!status.operational) break - if (cursor.atEnd()) break - - cursor.next() - } - return status - } -} - - /** * Strategy that works with occurrence store. * Doesn't directly work with [MatchJournal] except for default logging.