Refactor ProcessingStateImpl: move logic related to execution queue to its own class
This commit is contained in:
parent
577662e202
commit
2dba8890f3
|
|
@ -0,0 +1,121 @@
|
||||||
|
/*
|
||||||
|
* Copyright 2014-2019 JetBrains s.r.o.
|
||||||
|
*
|
||||||
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
* you may not use this file except in compliance with the License.
|
||||||
|
* You may obtain a copy of the License at
|
||||||
|
*
|
||||||
|
* http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
*
|
||||||
|
* Unless required by applicable law or agreed to in writing, software
|
||||||
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
* See the License for the specific language governing permissions and
|
||||||
|
* limitations under the License.
|
||||||
|
*/
|
||||||
|
|
||||||
|
package jetbrains.mps.logic.reactor.core.internal
|
||||||
|
|
||||||
|
import jetbrains.mps.logic.reactor.core.Controller
|
||||||
|
import jetbrains.mps.logic.reactor.core.Occurrence
|
||||||
|
import jetbrains.mps.logic.reactor.core.RuleMatchEx
|
||||||
|
import jetbrains.mps.logic.reactor.util.Id
|
||||||
|
import java.util.*
|
||||||
|
|
||||||
|
internal class ExecutionQueue(
|
||||||
|
private val journalIndex: MatchJournal.Index,
|
||||||
|
private val ruleOrdering: RuleOrdering
|
||||||
|
) {
|
||||||
|
|
||||||
|
private data class ExecPos(val pos: MatchJournal.Pos, val activeOcc: Occurrence)
|
||||||
|
|
||||||
|
// It is a position in Journal from previous session,
|
||||||
|
// from which incremental execution continues.
|
||||||
|
// Needed for pos comparison in postponeFutureMatches.
|
||||||
|
private lateinit var lastIncrementalRootPos: MatchJournal.Pos
|
||||||
|
|
||||||
|
private val postponedMatches: MutableMap<Id<Occurrence>, List<RuleMatchEx>> = HashMap()
|
||||||
|
|
||||||
|
private val execQueue: Queue<ExecPos> =
|
||||||
|
PriorityQueue<ExecPos>(1 + journalIndex.size / 2) { // just an estimate
|
||||||
|
lhs, rhs -> journalIndex.compare(lhs.pos, rhs.pos)
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
fun run(controller: Controller, state: ProcessingStateImpl): FeedbackStatus.NORMAL {
|
||||||
|
if (execQueue.isNotEmpty()) {
|
||||||
|
state.resetStore()
|
||||||
|
|
||||||
|
var prevPos: MatchJournal.Pos? = null
|
||||||
|
do {
|
||||||
|
val execPos = execQueue.poll()
|
||||||
|
|
||||||
|
// Handles the case when several matches are added to the same position.
|
||||||
|
// Then shouldn't replay, because currentPos is valid and more recent (!) than execPos.
|
||||||
|
if (execPos.pos != prevPos) {
|
||||||
|
state.replay(controller, execPos.pos)
|
||||||
|
lastIncrementalRootPos = execPos.pos
|
||||||
|
}
|
||||||
|
prevPos = execPos.pos
|
||||||
|
|
||||||
|
state.reactivate(controller, execPos.activeOcc)
|
||||||
|
} while (execQueue.isNotEmpty())
|
||||||
|
}
|
||||||
|
// Also replay to the end after queue is fully executed
|
||||||
|
state.replay(controller, state.last().toPos())
|
||||||
|
// fixme: get FeedbackStatus out of reactivate()
|
||||||
|
return FeedbackStatus.NORMAL()
|
||||||
|
}
|
||||||
|
|
||||||
|
fun withPostponedMatches(active: Occurrence, matches: List<RuleMatchEx>): List<RuleMatchEx> =
|
||||||
|
postponedMatches.remove(Id(active))?.let { postponed ->
|
||||||
|
// Sort matches according to rule priorities
|
||||||
|
(matches + postponed).sortedBy { ruleOrdering.orderOf(it.rule()) }
|
||||||
|
} ?: matches
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Determines, filters out and enqueues future matches.
|
||||||
|
* Returns only current matches.
|
||||||
|
*/
|
||||||
|
fun postponeFutureMatches(matches: List<RuleMatchEx>): List<RuleMatchEx> {
|
||||||
|
val currentMatches = mutableListOf<RuleMatchEx>()
|
||||||
|
for (m in matches) {
|
||||||
|
|
||||||
|
// Returns null for matches with occurrences only from this session
|
||||||
|
// because journalIndex indexes only previous session.
|
||||||
|
val pos = journalIndex.activationPos(m)
|
||||||
|
|
||||||
|
// if it is a future match
|
||||||
|
if (pos != null && journalIndex.compare(lastIncrementalRootPos, pos) < 0) {
|
||||||
|
val idOcc = Id(pos.occ)
|
||||||
|
postponedMatches[idOcc] = (postponedMatches[idOcc] ?: emptyList()) + listOf(m)
|
||||||
|
offer(pos)
|
||||||
|
} else {
|
||||||
|
currentMatches.add(m)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return currentMatches
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
// todo: do need to reactivate only the main, matching~activating match?
|
||||||
|
// (i.e. don't reactivate additional, inactive heads that only completed the match?)
|
||||||
|
fun offerAll(occs: Iterable<Occurrence>): Boolean =
|
||||||
|
execQueue.addAll(occs.mapNotNull { occ ->
|
||||||
|
journalIndex.activatingChunkOf(Id(occ))?.let { chunk ->
|
||||||
|
ExecPos(chunk.toPos(), occ)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
fun offer(posInJournal: MatchJournal.Pos, activeOcc: Occurrence): Boolean =
|
||||||
|
execQueue.offer(ExecPos(posInJournal, activeOcc))
|
||||||
|
|
||||||
|
// special case when position in trace corresponds to position of activated occurrence
|
||||||
|
fun offer(posInJournal: MatchJournal.Pos): Boolean =
|
||||||
|
execQueue.offer(ExecPos(posInJournal, posInJournal.occ))
|
||||||
|
|
||||||
|
fun isEmpty(): Boolean = execQueue.isEmpty()
|
||||||
|
|
||||||
|
fun isNotEmpty(): Boolean = execQueue.isNotEmpty()
|
||||||
|
|
||||||
|
}
|
||||||
|
|
@ -108,6 +108,11 @@ interface MatchJournal : MutableIterable<MatchJournal.Chunk> {
|
||||||
match.signature().mapNotNull { occSig ->
|
match.signature().mapNotNull { occSig ->
|
||||||
occSig?.let { activatingChunkOf(it)?.toPos() }
|
occSig?.let { activatingChunkOf(it)?.toPos() }
|
||||||
}.maxWith(this) // compare positions: find latest
|
}.maxWith(this) // compare positions: find latest
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Length of the indexed [MatchJournal]
|
||||||
|
*/
|
||||||
|
val size: Int
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
|
||||||
|
|
@ -181,6 +181,8 @@ internal open class MatchJournalImpl(
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
override val size: Int = chunks.count()
|
||||||
|
|
||||||
override fun activatingChunkOf(occId: Id<Occurrence>) = occChunks[occId]
|
override fun activatingChunkOf(occId: Id<Occurrence>) = occChunks[occId]
|
||||||
|
|
||||||
// todo: throw for invalid positions?
|
// todo: throw for invalid positions?
|
||||||
|
|
|
||||||
|
|
@ -54,23 +54,24 @@ internal class ProcessingStateImpl(private var dispatchingFront: Dispatcher.Disp
|
||||||
|
|
||||||
private val journalIndex: MatchJournal.Index = journal.index()
|
private val journalIndex: MatchJournal.Index = journal.index()
|
||||||
|
|
||||||
// It is a position in Journal from previous session,
|
private val execQueue: ExecutionQueue = ExecutionQueue(journalIndex, ruleOrdering)
|
||||||
// from which incremental execution continues.
|
|
||||||
// Needed for pos comparison in postponeFutureMatches.
|
|
||||||
private lateinit var lastIncrementalRootPos: MatchJournal.Pos
|
|
||||||
|
|
||||||
private val postponedMatches: MutableMap<Id<Occurrence>, List<RuleMatchEx>> = HashMap()
|
|
||||||
|
|
||||||
private val execQueue: Queue<ExecPos> =
|
|
||||||
PriorityQueue<ExecPos>(1 + this.count() / 2) { // just an estimate
|
|
||||||
lhs, rhs -> journalIndex.compare(lhs.pos, rhs.pos)
|
|
||||||
}
|
|
||||||
|
|
||||||
private data class ExecPos(val pos: MatchJournal.Pos, val activeOcc: Occurrence)
|
|
||||||
|
|
||||||
private data class MatchCandidate(val rule: Rule, val occChunk: MatchJournal.OccChunk)
|
private data class MatchCandidate(val rule: Rule, val occChunk: MatchJournal.OccChunk)
|
||||||
|
|
||||||
|
|
||||||
|
fun reactivate(controller: Controller, activeOcc: Occurrence) {
|
||||||
|
// If the occurrence is still in the store after replay (i.e. if it's valid to activate it)
|
||||||
|
if (activeOcc.stored) {
|
||||||
|
// Forget that occ was seen.
|
||||||
|
// Incremental reactivation isn't like the usual reactivation,
|
||||||
|
// it should proceed more like usual activation.
|
||||||
|
this.dispatchingFront = dispatchingFront.forgetSeen(activeOcc)
|
||||||
|
trace.reactivateIncremental(activeOcc)
|
||||||
|
controller.reactivate(activeOcc)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fun invalidateDependentRules(ruleIds: Set<Any>) {
|
fun invalidateDependentRules(ruleIds: Set<Any>) {
|
||||||
val it = this.iterator()
|
val it = this.iterator()
|
||||||
while (it.hasNext()) {
|
while (it.hasNext()) {
|
||||||
|
|
@ -129,18 +130,12 @@ internal class ProcessingStateImpl(private var dispatchingFront: Dispatcher.Disp
|
||||||
// We removed the match, so need to reactivate all still valid occurrences from the head
|
// We removed the match, so need to reactivate all still valid occurrences from the head
|
||||||
// by definition of Chunk and principal rule, all occurrences from the head are principal
|
// by definition of Chunk and principal rule, all occurrences from the head are principal
|
||||||
val matchedOccs = chunk.match.allHeads() as Iterable<Occurrence>
|
val matchedOccs = chunk.match.allHeads() as Iterable<Occurrence>
|
||||||
val (invalidatedOccs, validOccs) = matchedOccs.partition { occ ->
|
val validOccs = matchedOccs.filter { occ ->
|
||||||
occ.justifications().intersects(justificationRoots)
|
!occ.justifications().intersects(justificationRoots)
|
||||||
}
|
}
|
||||||
assert(matchedOccs.all { it.isPrincipal() })
|
assert(matchedOccs.all { it.isPrincipal() })
|
||||||
|
|
||||||
// todo: do need to reactivate only the main, matching~activating match?
|
execQueue.offerAll(validOccs)
|
||||||
// (i.e. don't reactivate additional, inactive heads that only completed the match?)
|
|
||||||
execQueue.addAll(validOccs.mapNotNull { occ ->
|
|
||||||
journalIndex.activatingChunkOf(Id(occ))?.let { chunk ->
|
|
||||||
ExecPos(chunk.toPos(), occ)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -196,7 +191,7 @@ internal class ProcessingStateImpl(private var dispatchingFront: Dispatcher.Disp
|
||||||
else if (!it.hasNext()) chunk.toPos()
|
else if (!it.hasNext()) chunk.toPos()
|
||||||
else continue
|
else continue
|
||||||
|
|
||||||
execQueue.offer(ExecPos(pos, occChunk.occ))
|
execQueue.offer(pos, occChunk.occ)
|
||||||
trace.potentialMatch(occChunk.occ, candRule)
|
trace.potentialMatch(occChunk.occ, candRule)
|
||||||
// Drop the candidate if appropriate activation place is found.
|
// Drop the candidate if appropriate activation place is found.
|
||||||
aIt.remove()
|
aIt.remove()
|
||||||
|
|
@ -208,38 +203,8 @@ internal class ProcessingStateImpl(private var dispatchingFront: Dispatcher.Disp
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fun launchQueue(controller: Controller): FeedbackStatus.NORMAL {
|
fun launchQueue(controller: Controller): FeedbackStatus.NORMAL =
|
||||||
if (execQueue.isNotEmpty()) {
|
execQueue.run(controller, this)
|
||||||
resetStore()
|
|
||||||
|
|
||||||
var prevPos: MatchJournal.Pos? = null
|
|
||||||
do {
|
|
||||||
val execPos = execQueue.poll()
|
|
||||||
|
|
||||||
// Handles the case when several matches are added to the same position.
|
|
||||||
// Then shouldn't replay, because currentPos is valid and more recent (!) than execPos.
|
|
||||||
if (execPos.pos != prevPos) {
|
|
||||||
replay(controller, execPos.pos)
|
|
||||||
lastIncrementalRootPos = execPos.pos
|
|
||||||
}
|
|
||||||
prevPos = execPos.pos
|
|
||||||
|
|
||||||
// If the occurrence is still in the store after replay (i.e. if it's valid to activate it)
|
|
||||||
if (execPos.activeOcc.stored) {
|
|
||||||
// Forget that occ was seen.
|
|
||||||
// Incremental reactivation isn't like the usual reactivation,
|
|
||||||
// it should proceed more like usual activation.
|
|
||||||
this.dispatchingFront = dispatchingFront.forgetSeen(execPos.activeOcc)
|
|
||||||
trace.reactivateIncremental(execPos.activeOcc)
|
|
||||||
controller.reactivate(execPos.activeOcc)
|
|
||||||
}
|
|
||||||
} while (execQueue.isNotEmpty())
|
|
||||||
}
|
|
||||||
// Also replay to the end after queue is fully executed
|
|
||||||
replay(controller, this.last().toPos())
|
|
||||||
// fixme: get FeedbackStatus out of reactivate()
|
|
||||||
return FeedbackStatus.NORMAL()
|
|
||||||
}
|
|
||||||
|
|
||||||
fun snapshot(): SessionToken =
|
fun snapshot(): SessionToken =
|
||||||
SessionTokenImpl(view(), ruleOrdering.ruleTags, dispatchingFront.state())
|
SessionTokenImpl(view(), ruleOrdering.ruleTags, dispatchingFront.state())
|
||||||
|
|
@ -274,9 +239,10 @@ internal class ProcessingStateImpl(private var dispatchingFront: Dispatcher.Disp
|
||||||
if (isFront() || !active.isPrincipal()) {
|
if (isFront() || !active.isPrincipal()) {
|
||||||
matches
|
matches
|
||||||
} else {
|
} else {
|
||||||
postponeFutureMatches(matches)
|
assert( matches.all { ispec.isPrincipal(it.rule()) } )
|
||||||
|
execQueue.postponeFutureMatches(matches)
|
||||||
}
|
}
|
||||||
val currentMatches = withPostponedMatches(active, newCurrentMatches)
|
val currentMatches = execQueue.withPostponedMatches(active, newCurrentMatches)
|
||||||
|
|
||||||
val outStatus = currentMatches.fold(inStatus) { status, match ->
|
val outStatus = currentMatches.fold(inStatus) { status, match ->
|
||||||
// TODO: paranoid check. should be isAlive() instead
|
// TODO: paranoid check. should be isAlive() instead
|
||||||
|
|
@ -296,42 +262,6 @@ internal class ProcessingStateImpl(private var dispatchingFront: Dispatcher.Disp
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private fun withPostponedMatches(active: Occurrence, matches: List<RuleMatchEx>): List<RuleMatchEx> =
|
|
||||||
postponedMatches.remove(Id(active))?.let { postponed ->
|
|
||||||
// Sort matches according to rule priorities
|
|
||||||
(matches + postponed).sortedBy { ruleOrdering.orderOf(it.rule()) }
|
|
||||||
} ?: matches
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Determines, filters out and enqueues (to execution queue) future matches.
|
|
||||||
* Returns only current matches.
|
|
||||||
*/
|
|
||||||
private fun postponeFutureMatches(matches: List<RuleMatchEx>): List<RuleMatchEx> {
|
|
||||||
assert(
|
|
||||||
matches.all { ispec.isPrincipal(it.rule()) },
|
|
||||||
{ "non-principal ctrs in head of principal rule: ${ matches.filter { !ispec.isPrincipal(it.rule()) } }" }
|
|
||||||
)
|
|
||||||
|
|
||||||
val currentMatches = mutableListOf<RuleMatchEx>()
|
|
||||||
for (m in matches) {
|
|
||||||
|
|
||||||
// Returns null for matches with occurrences only from this session
|
|
||||||
// because journalIndex indexes only previous session.
|
|
||||||
val pos = journalIndex.activationPos(m)
|
|
||||||
|
|
||||||
// if it is a future match
|
|
||||||
if (pos != null && journalIndex.compare(lastIncrementalRootPos, pos) < 0) {
|
|
||||||
val idOcc = Id(pos.occ)
|
|
||||||
postponedMatches[idOcc] = (postponedMatches[idOcc] ?: emptyList()) + listOf(m)
|
|
||||||
execQueue.offer(ExecPos(pos, pos.occ))
|
|
||||||
} else {
|
|
||||||
currentMatches.add(m)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return currentMatches
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
private inline fun FeedbackStatus.then(action: (FeedbackStatus) -> FeedbackStatus) : FeedbackStatus =
|
private inline fun FeedbackStatus.then(action: (FeedbackStatus) -> FeedbackStatus) : FeedbackStatus =
|
||||||
if (operational) action(this) else this
|
if (operational) action(this) else this
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue