Refactor session accounting in daemon, add REPL to it, write a test

This commit is contained in:
Ilya Chernikov
2016-10-02 16:42:27 +02:00
parent 3a52c68973
commit ea4a1df839
3 changed files with 143 additions and 81 deletions
@@ -62,7 +62,12 @@ open class KotlinRemoteReplClientBase(
init { init {
Disposer.register(disposable, Disposable { Disposer.register(disposable, Disposable {
compileService.releaseReplSession(sessionId) try {
compileService.releaseReplSession(sessionId)
}
catch (ex: java.rmi.RemoteException) {
// assuming that communication failed and daemon most likely is already down
}
}) })
} }
@@ -151,27 +151,33 @@ class CompileServiceImpl(
// RMI-exposed API // RMI-exposed API
override fun getDaemonOptions(): CompileService.CallResult<DaemonOptions> = ifAlive { daemonOptions } override fun getDaemonOptions(): CompileService.CallResult<DaemonOptions> = ifAlive {
CompileService.CallResult.Good(daemonOptions)
}
override fun getDaemonJVMOptions(): CompileService.CallResult<DaemonJVMOptions> = ifAlive { daemonJVMOptions } override fun getDaemonJVMOptions(): CompileService.CallResult<DaemonJVMOptions> = ifAlive {
CompileService.CallResult.Good(daemonJVMOptions)
}
override fun registerClient(aliveFlagPath: String?): CompileService.CallResult<Nothing> = ifAlive_Nothing { override fun registerClient(aliveFlagPath: String?): CompileService.CallResult<Nothing> = ifAlive {
synchronized(state.clientProxies) { synchronized(state.clientProxies) {
state.clientProxies.add(ClientOrSessionProxy(aliveFlagPath)) state.clientProxies.add(ClientOrSessionProxy(aliveFlagPath))
} }
CompileService.CallResult.Ok()
} }
override fun getClients(): CompileService.CallResult<List<String>> = ifAlive { override fun getClients(): CompileService.CallResult<List<String>> = ifAlive {
synchronized(state.clientProxies) { synchronized(state.clientProxies) {
state.clientProxies.mapNotNull { it.aliveFlagPath } CompileService.CallResult.Good(state.clientProxies.mapNotNull { it.aliveFlagPath })
} }
} }
// TODO: consider tying a session to a client and use this info to cleanup // TODO: consider tying a session to a client and use this info to cleanup
override fun leaseCompileSession(aliveFlagPath: String?): CompileService.CallResult<Int> = ifAlive(minAliveness = Aliveness.Alive) { override fun leaseCompileSession(aliveFlagPath: String?): CompileService.CallResult<Int> = ifAlive(minAliveness = Aliveness.Alive) {
leaseSessionImpl(ClientOrSessionProxy<Any>(aliveFlagPath)).apply { CompileService.CallResult.Good(
log.info("leased a new session $this, client alive file: $aliveFlagPath") leaseSessionImpl(ClientOrSessionProxy<Any>(aliveFlagPath)).apply {
} log.info("leased a new session $this, client alive file: $aliveFlagPath")
})
} }
private fun<T: Any> leaseSessionImpl(session: ClientOrSessionProxy<T>): Int { private fun<T: Any> leaseSessionImpl(session: ClientOrSessionProxy<T>): Int {
@@ -192,7 +198,7 @@ class CompileServiceImpl(
throw IllegalStateException("Invalid state or algorithm error") throw IllegalStateException("Invalid state or algorithm error")
} }
override fun releaseCompileSession(sessionId: Int) = ifAlive_Nothing(minAliveness = Aliveness.LastSession) { override fun releaseCompileSession(sessionId: Int) = ifAlive(minAliveness = Aliveness.LastSession) {
synchronized(state.sessions) { synchronized(state.sessions) {
state.sessions[sessionId]?.dispose() state.sessions[sessionId]?.dispose()
state.sessions.remove(sessionId) state.sessions.remove(sessionId)
@@ -205,6 +211,7 @@ class CompileServiceImpl(
timer.schedule(0) { timer.schedule(0) {
periodicAndAfterSessionCheck() periodicAndAfterSessionCheck()
} }
CompileService.CallResult.Ok()
} }
override fun checkCompilerId(expectedCompilerId: CompilerId): Boolean = override fun checkCompilerId(expectedCompilerId: CompilerId): Boolean =
@@ -212,27 +219,31 @@ class CompileServiceImpl(
(compilerId.compilerClasspath.all { expectedCompilerId.compilerClasspath.contains(it) }) && (compilerId.compilerClasspath.all { expectedCompilerId.compilerClasspath.contains(it) }) &&
!classpathWatcher.isChanged !classpathWatcher.isChanged
override fun getUsedMemory(): CompileService.CallResult<Long> = ifAlive { usedMemory(withGC = true) } override fun getUsedMemory(): CompileService.CallResult<Long> =
ifAlive { CompileService.CallResult.Good(usedMemory(withGC = true)) }
override fun shutdown(): CompileService.CallResult<Nothing> = ifAliveExclusive_Nothing(minAliveness = Aliveness.LastSession, ignoreCompilerChanged = true) { override fun shutdown(): CompileService.CallResult<Nothing> = ifAliveExclusive(minAliveness = Aliveness.LastSession, ignoreCompilerChanged = true) {
shutdownImpl() shutdownImpl()
CompileService.CallResult.Ok()
} }
override fun scheduleShutdown(graceful: Boolean): CompileService.CallResult<Boolean> = ifAlive(minAliveness = Aliveness.Alive) { override fun scheduleShutdown(graceful: Boolean): CompileService.CallResult<Boolean> = ifAlive(minAliveness = Aliveness.Alive) {
if (!graceful || state.alive.compareAndSet(Aliveness.Alive.ordinal, Aliveness.LastSession.ordinal)) { CompileService.CallResult.Good(
timer.schedule(0) { if (!graceful || state.alive.compareAndSet(Aliveness.Alive.ordinal, Aliveness.LastSession.ordinal)) {
ifAliveExclusive(minAliveness = Aliveness.LastSession, ignoreCompilerChanged = true) { timer.schedule(0) {
if (!graceful || state.sessions.isEmpty()) { ifAliveExclusive(minAliveness = Aliveness.LastSession, ignoreCompilerChanged = true) {
shutdownImpl() if (!graceful || state.sessions.isEmpty()) {
} shutdownImpl()
else { }
log.info("Some sessions are active, waiting for them to finish") else {
log.info("Some sessions are active, waiting for them to finish")
}
CompileService.CallResult.Ok()
}
} }
true
} }
} else false)
true
}
else false
} }
override fun remoteCompile(sessionId: Int, override fun remoteCompile(sessionId: Int,
@@ -280,7 +291,7 @@ class CompileServiceImpl(
evalErrorStream: RemoteOutputStream?, evalErrorStream: RemoteOutputStream?,
evalInputStream: RemoteInputStream?, evalInputStream: RemoteInputStream?,
operationsTracer: RemoteOperationsTracer? operationsTracer: RemoteOperationsTracer?
): CompileService.CallResult<Int> = ifAlive_Raw(minAliveness = Aliveness.Alive) { ): CompileService.CallResult<Int> = ifAlive(minAliveness = Aliveness.Alive) {
if (targetPlatform != CompileService.TargetPlatform.JVM) if (targetPlatform != CompileService.TargetPlatform.JVM)
CompileService.CallResult.Error("Sorry, only JVM target platform is supported now") CompileService.CallResult.Error("Sorry, only JVM target platform is supported now")
else { else {
@@ -292,25 +303,18 @@ class CompileServiceImpl(
} }
} }
private fun<R> withValidRepl(sessionId: Int, body: KotlinJvmReplService.() -> R): CompileService.CallResult<R> =
state.sessions[sessionId]?.let {
(it.data as? KotlinJvmReplService?)?.let {
CompileService.CallResult.Good(it.body())
}
} ?: CompileService.CallResult.Error("Unknown or invalid session $sessionId")
// TODO: add more checks (e.g. is it a repl session) // TODO: add more checks (e.g. is it a repl session)
override fun releaseReplSession(sessionId: Int): CompileService.CallResult<Nothing> = releaseCompileSession(sessionId) override fun releaseReplSession(sessionId: Int): CompileService.CallResult<Nothing> = releaseCompileSession(sessionId)
override fun remoteReplLineCheck(sessionId: Int, codeLine: ReplCodeLine, history: List<ReplCodeLine>): CompileService.CallResult<ReplCheckResult> = override fun remoteReplLineCheck(sessionId: Int, codeLine: ReplCodeLine, history: List<ReplCodeLine>): CompileService.CallResult<ReplCheckResult> =
ifAlive_Raw(minAliveness = Aliveness.Alive) { ifAlive(minAliveness = Aliveness.Alive) {
withValidRepl(sessionId) { withValidRepl(sessionId) {
check(codeLine, history) check(codeLine, history)
} }
} }
override fun remoteReplLineCompile(sessionId: Int, codeLine: ReplCodeLine, history: List<ReplCodeLine>): CompileService.CallResult<ReplCompileResult> = override fun remoteReplLineCompile(sessionId: Int, codeLine: ReplCodeLine, history: List<ReplCodeLine>): CompileService.CallResult<ReplCompileResult> =
ifAlive_Raw(minAliveness = Aliveness.Alive) { ifAlive(minAliveness = Aliveness.Alive) {
withValidRepl(sessionId) { withValidRepl(sessionId) {
compile(codeLine, history) compile(codeLine, history)
} }
@@ -321,7 +325,7 @@ class CompileServiceImpl(
codeLine: ReplCodeLine, codeLine: ReplCodeLine,
history: List<ReplCodeLine> history: List<ReplCodeLine>
): CompileService.CallResult<ReplEvalResult> = ): CompileService.CallResult<ReplEvalResult> =
ifAlive_Raw(minAliveness = Aliveness.Alive) { ifAlive(minAliveness = Aliveness.Alive) {
withValidRepl(sessionId) { withValidRepl(sessionId) {
eval(codeLine, history) eval(codeLine, history)
} }
@@ -373,7 +377,7 @@ class CompileServiceImpl(
private fun periodicAndAfterSessionCheck() { private fun periodicAndAfterSessionCheck() {
ifAlive_Nothing(minAliveness = Aliveness.LastSession) { ifAlive(minAliveness = Aliveness.LastSession) {
// 1. check if unused for a timeout - shutdown // 1. check if unused for a timeout - shutdown
if (shutdownCondition({ daemonOptions.autoshutdownUnusedSeconds != COMPILE_DAEMON_TIMEOUT_INFINITE_S && compilationsCounter.get() == 0 && nowSeconds() - lastUsedSeconds > daemonOptions.autoshutdownUnusedSeconds }, if (shutdownCondition({ daemonOptions.autoshutdownUnusedSeconds != COMPILE_DAEMON_TIMEOUT_INFINITE_S && compilationsCounter.get() == 0 && nowSeconds() - lastUsedSeconds > daemonOptions.autoshutdownUnusedSeconds },
@@ -427,13 +431,14 @@ class CompileServiceImpl(
clearJarCache() clearJarCache()
} }
} }
CompileService.CallResult.Ok()
} }
} }
private fun initiateElections() { private fun initiateElections() {
ifAlive_Nothing { ifAlive {
val aliveWithOpts = walkDaemons(File(daemonOptions.runFilesPathOrDefault), compilerId, filter = { f, p -> p != port }, report = { lvl, msg -> log.info(msg) }) val aliveWithOpts = walkDaemons(File(daemonOptions.runFilesPathOrDefault), compilerId, filter = { f, p -> p != port }, report = { lvl, msg -> log.info(msg) })
.map { Pair(it, it.getDaemonJVMOptions()) } .map { Pair(it, it.getDaemonJVMOptions()) }
@@ -465,6 +470,7 @@ class CompileServiceImpl(
// - shutdown/takeover smaller daemon // - shutdown/takeover smaller daemon
// - run (or better persuade client to run) a bigger daemon (in fact may be even simple shutdown will do, because of client's daemon choosing logic) // - run (or better persuade client to run) a bigger daemon (in fact may be even simple shutdown will do, because of client's daemon choosing logic)
} }
CompileService.CallResult.Ok()
} }
} }
@@ -507,25 +513,24 @@ class CompileServiceImpl(
operationsTracer: RemoteOperationsTracer?, operationsTracer: RemoteOperationsTracer?,
body: (PrintStream, EventManger, Profiler) -> ExitCode): CompileService.CallResult<Int> = body: (PrintStream, EventManger, Profiler) -> ExitCode): CompileService.CallResult<Int> =
ifAlive { ifAlive {
withValidClientOrSessionProxy(sessionId) { session ->
operationsTracer?.before("compile") operationsTracer?.before("compile")
compilationsCounter.incrementAndGet() val rpcProfiler = if (daemonOptions.reportPerf) WallAndThreadTotalProfiler() else DummyProfiler()
val rpcProfiler = if (daemonOptions.reportPerf) WallAndThreadTotalProfiler() else DummyProfiler() val eventManger = EventMangerImpl()
val eventManger = EventMangerImpl() val compilerMessagesStream = PrintStream(BufferedOutputStream(RemoteOutputStreamClient(compilerMessagesStreamProxy, rpcProfiler), REMOTE_STREAM_BUFFER_SIZE))
val compilerMessagesStream = PrintStream(BufferedOutputStream(RemoteOutputStreamClient(compilerMessagesStreamProxy, rpcProfiler), REMOTE_STREAM_BUFFER_SIZE)) val serviceOutputStream = PrintStream(BufferedOutputStream(RemoteOutputStreamClient(serviceOutputStreamProxy, rpcProfiler), REMOTE_STREAM_BUFFER_SIZE))
val serviceOutputStream = PrintStream(BufferedOutputStream(RemoteOutputStreamClient(serviceOutputStreamProxy, rpcProfiler), REMOTE_STREAM_BUFFER_SIZE)) try {
try { CompileService.CallResult.Good(
checkedCompile(args, serviceOutputStream, rpcProfiler) { checkedCompile(args, serviceOutputStream, rpcProfiler) {
val res = body(compilerMessagesStream, eventManger, rpcProfiler).code body(compilerMessagesStream, eventManger, rpcProfiler).code
_lastUsedSeconds = nowSeconds() })
res }
finally {
serviceOutputStream.flush()
compilerMessagesStream.flush()
eventManger.fireCompilationFinished()
operationsTracer?.after("compile")
} }
}
finally {
serviceOutputStream.flush()
compilerMessagesStream.flush()
eventManger.fireCompilationFinished()
operationsTracer?.after("compile")
} }
} }
@@ -588,34 +593,21 @@ class CompileServiceImpl(
(KotlinCoreEnvironment.applicationEnvironment?.jarFileSystem as? CoreJarFileSystem)?.clearHandlersCache() (KotlinCoreEnvironment.applicationEnvironment?.jarFileSystem as? CoreJarFileSystem)?.clearHandlersCache()
} }
private fun<R> ifAlive(minAliveness: Aliveness = Aliveness.Alive, ignoreCompilerChanged: Boolean = false, body: () -> R): CompileService.CallResult<R> = rwlock.read { private fun<R> ifAlive(minAliveness: Aliveness = Aliveness.Alive,
ifAliveChecksImpl(minAliveness, ignoreCompilerChanged) { CompileService.CallResult.Good(body()) } ignoreCompilerChanged: Boolean = false,
} body: () -> CompileService.CallResult<R>
): CompileService.CallResult<R> =
rwlock.read {
ifAliveChecksImpl(minAliveness, ignoreCompilerChanged, body)
}
// TODO: find how to implement it without using unique name for this variant; making name deliberately ugly meanwhile private fun<R> ifAliveExclusive(minAliveness: Aliveness = Aliveness.Alive,
private fun ifAlive_Nothing(minAliveness: Aliveness = Aliveness.Alive, ignoreCompilerChanged: Boolean = false, body: () -> Unit): CompileService.CallResult<Nothing> = rwlock.read { ignoreCompilerChanged: Boolean = false,
ifAliveChecksImpl(minAliveness, ignoreCompilerChanged) { body: () -> CompileService.CallResult<R>
body() ): CompileService.CallResult<R> =
CompileService.CallResult.Ok() rwlock.write {
} ifAliveChecksImpl(minAliveness, ignoreCompilerChanged, body)
} }
private fun<R> ifAlive_Raw(minAliveness: Aliveness = Aliveness.Alive, ignoreCompilerChanged: Boolean = false, body: () -> CompileService.CallResult<R>): CompileService.CallResult<R> = rwlock.read {
ifAliveChecksImpl(minAliveness, ignoreCompilerChanged) { body() }
}
private fun<R> ifAliveExclusive(minAliveness: Aliveness = Aliveness.Alive, ignoreCompilerChanged: Boolean = false, body: () -> R): CompileService.CallResult<R> = rwlock.write {
ifAliveChecksImpl(minAliveness, ignoreCompilerChanged) { CompileService.CallResult.Good(body()) }
}
// see comment to ifAliveNothing
private fun<R> ifAliveExclusive_Nothing(minAliveness: Aliveness = Aliveness.Alive, ignoreCompilerChanged: Boolean = false, body: () -> Unit): CompileService.CallResult<R> = rwlock.write {
ifAliveChecksImpl(minAliveness, ignoreCompilerChanged) {
body()
CompileService.CallResult.Ok()
}
}
inline private fun<R> ifAliveChecksImpl(minAliveness: Aliveness = Aliveness.Alive, ignoreCompilerChanged: Boolean = false, body: () -> CompileService.CallResult<R>): CompileService.CallResult<R> = inline private fun<R> ifAliveChecksImpl(minAliveness: Aliveness = Aliveness.Alive, ignoreCompilerChanged: Boolean = false, body: () -> CompileService.CallResult<R>): CompileService.CallResult<R> =
when { when {
@@ -635,4 +627,27 @@ class CompileServiceImpl(
} }
} }
} }
private inline fun<R> withValidClientOrSessionProxy(sessionId: Int,
body: (ClientOrSessionProxy<Any>?) -> CompileService.CallResult<R>
): CompileService.CallResult<R> {
val session: ClientOrSessionProxy<Any>? =
if (sessionId == CompileService.NO_SESSION) null
else state.sessions[sessionId] ?: return CompileService.CallResult.Error("Unknown or invalid session $sessionId")
try {
compilationsCounter.incrementAndGet()
return body(session)
}
finally {
_lastUsedSeconds = nowSeconds()
}
}
private inline fun<R> withValidRepl(sessionId: Int, body: KotlinJvmReplService.() -> R): CompileService.CallResult<R> =
withValidClientOrSessionProxy(sessionId) { session ->
(session?.data as? KotlinJvmReplService?)?.let {
CompileService.CallResult.Good(it.body())
} ?: CompileService.CallResult.Error("Not a REPL session $sessionId")
}
} }
@@ -567,6 +567,48 @@ class CompilerDaemonTest : KotlinIntegrationTestBase() {
} }
} }
fun testDaemonReplAutoshutdownOnIdle() {
withFlagFile(getTestName(true), ".alive") { flagFile ->
val daemonOptions = DaemonOptions(autoshutdownIdleSeconds = 1, autoshutdownUnusedSeconds = 1, runFilesPath = File(tmpdir, getTestName(true)).absolutePath)
KotlinCompilerClient.shutdownCompileService(compilerId, daemonOptions)
val logFile = createTempFile("kotlin-daemon-test", ".log")
val daemonJVMOptions =
configureDaemonJVMOptions("D$COMPILE_DAEMON_LOG_PATH_PROPERTY=\"${logFile.loggerCompatiblePath}\"",
inheritMemoryLimits = false, inheritAdditionalProperties = false)
val daemon = KotlinCompilerClient.connectToCompileService(compilerId, flagFile, daemonJVMOptions, daemonOptions, DaemonReportingTargets(out = System.err), autostart = true)
assertNotNull("failed to connect daemon", daemon)
val disposable = Disposer.newDisposable()
val repl = KotlinRemoteReplCompiler(disposable, daemon!!, null, CompileService.TargetPlatform.JVM,
classpathFromClassloader(),
ScriptWithNoParam::class.qualifiedName!!,
System.err)
// use repl compiler for >> 1s, making sure that idle/unused timeouts are not firing
for (attempts in 1..10) {
val codeLine1 = ReplCodeLine(1, "3 + 5")
val res1 = repl.compile(codeLine1, emptyList())
val res1c = res1 as? ReplCompileResult.CompiledClasses
TestCase.assertNotNull("Unexpected compile result: $res1", res1c)
Thread.sleep(200)
}
// wait up to 4s (more than 1s idle timeout)
for (attempts in 1..20) {
if (logFile.isLogContainsSequence("Idle timeout exceeded 1s")) break
Thread.sleep(200)
}
Disposer.dispose(disposable)
Thread.sleep(200)
logFile.assertLogContainsSequence("Idle timeout exceeded 1s",
"Shutdown complete")
logFile.delete()
}
}
internal fun withDaemon(body: (CompileService) -> Unit) { internal fun withDaemon(body: (CompileService) -> Unit) {
withFlagFile(getTestName(true), ".alive") { flagFile -> withFlagFile(getTestName(true), ".alive") { flagFile ->
val daemonOptions = DaemonOptions(runFilesPath = File(tmpdir, getTestName(true)).absolutePath) val daemonOptions = DaemonOptions(runFilesPath = File(tmpdir, getTestName(true)).absolutePath)