Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
package com.jorisjonkers.personalstack.agents.application.sessionbinding

import com.jorisjonkers.personalstack.agents.application.exception.AgentRunnerUnavailableException
import com.jorisjonkers.personalstack.agents.domain.model.Workspace
import com.jorisjonkers.personalstack.agents.domain.model.WorkspaceAgentSession
import com.jorisjonkers.personalstack.agents.domain.port.AgentGatewayClient
import org.slf4j.LoggerFactory
import org.springframework.web.client.ResourceAccessException

/**
* Handles the spawn-with-retry loop for agent gateway sessions. Extracted from
* RunnerSessionBinder to keep that class below the TooManyFunctions threshold.
*/
internal class RunnerAgentSpawner(
private val gateway: AgentGatewayClient,
private val backoffInitialMs: Long,
) {
private val log = LoggerFactory.getLogger(RunnerAgentSpawner::class.java)

fun spawnWithRetry(
workspace: Workspace,
session: WorkspaceAgentSession,
continuation: AgentGatewayClient.ContinuationMetadata?,
): AgentGatewayClient.GatewayAgent {
var lastFailure: ResourceAccessException? = null
repeat(MAX_SPAWN_ATTEMPTS) { attempt ->
try {
return gateway.spawnAgent(buildSpawnRequest(workspace, session, continuation))
} catch (ex: ResourceAccessException) {
lastFailure = ex
logRetry(workspace, ex, attempt)
}
}
throw AgentRunnerUnavailableException(
workspaceId = workspace.id,
runnerStatus = "ConnectionRefused",
retryAfterSeconds = AgentRunnerUnavailableException.DEFAULT_RETRY_AFTER_SECONDS,
cause = lastFailure,
)
}

private fun buildSpawnRequest(
workspace: Workspace,
session: WorkspaceAgentSession,
continuation: AgentGatewayClient.ContinuationMetadata?,
) = AgentGatewayClient.SpawnAgentRequest(
workspace = workspace,
kind = session.kind,
stableSessionId = session.id,
epoch = session.epoch,
continuation = continuation,
resumeCliSessionId = session.cliSessionId,
)

private fun logRetry(
workspace: Workspace,
ex: ResourceAccessException,
attempt: Int,
) {
val sleepMs = backoffInitialMs * (attempt + 1)
log.warn(
"agent spawn attempt {} for workspace {} failed: {} - retrying in {}ms",
attempt + 1,
workspace.id.value,
ex.message,
sleepMs,
)
if (sleepMs > 0) Thread.sleep(sleepMs)
}

companion object {
const val MAX_SPAWN_ATTEMPTS: Int = 3
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
package com.jorisjonkers.personalstack.agents.application.sessionbinding

import com.jorisjonkers.personalstack.agents.application.workspacerunner.RunnerSetupTarget
import com.jorisjonkers.personalstack.agents.application.workspacerunner.RunnerUnavailableReason
import com.jorisjonkers.personalstack.agents.domain.model.RunnerSetupOperation
import com.jorisjonkers.personalstack.agents.domain.model.Workspace
import com.jorisjonkers.personalstack.agents.domain.model.WorkspaceAgentSession

/**
* Guard checks shared by the runner session binding flows. Extracted from
* RunnerSessionBinder to keep that class below the TooManyFunctions threshold.
*/
internal class RunnerBindingGuards(
private val provisioning: RunnerProvisioningCoordinator,
) {
fun hasSetupOperationInProgress(
session: WorkspaceAgentSession,
workspace: Workspace,
): Boolean = sessionHasPendingSetup(session) || workspaceSetupInProgress(workspace)

// Returns a terminal guard result when the session or workspace is blocked, null otherwise.
fun ensureBoundGuard(
session: WorkspaceAgentSession,
workspace: Workspace,
): RunnerSessionBindingResult? =
when {
sessionHasPendingSetup(session) -> RunnerSessionBindingResult.Conflict(current = session)
workspaceSetupInProgress(workspace) ->
RunnerSessionBindingResult.Unavailable(
workspaceId = workspace.id,
runnerStatus = RunnerUnavailableReason.SETUP_OPERATION_IN_PROGRESS.label,
)
else -> null
}

fun checkBindingReadiness(
workspace: Workspace,
target: RunnerSetupTarget,
): RunnerSessionBindingResult.Unavailable? {
val unavailableReason =
when {
workspaceSetupInProgress(workspace) -> RunnerUnavailableReason.SETUP_OPERATION_IN_PROGRESS.label
workspace.runnerBootLeaseId != null -> RunnerUnavailableReason.BOOT_LEASE_HELD.label
!provisioning.isRunnerReadyFor(workspace, target) ->
RunnerUnavailableReason.NOT_READY_AFTER_PROVISION.label
else -> null
}
return unavailableReason?.let {
RunnerSessionBindingResult.Unavailable(workspaceId = workspace.id, runnerStatus = it)
}
}

private fun sessionHasPendingSetup(session: WorkspaceAgentSession): Boolean =
session.pendingSetupId != null || session.pendingSetupVersion != null

private fun workspaceSetupInProgress(workspace: Workspace): Boolean =
workspace.pendingRunnerSetupId != null ||
workspace.pendingRunnerSetupVersion != null ||
workspace.runnerSetupOperation != RunnerSetupOperation.IDLE
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,104 @@
package com.jorisjonkers.personalstack.agents.application.sessionbinding

import com.jorisjonkers.personalstack.agents.application.exception.AgentRunnerUnavailableException
import com.jorisjonkers.personalstack.agents.application.exception.AgentSetupValidationException
import com.jorisjonkers.personalstack.agents.application.observability.AgentsApiTelemetry
import com.jorisjonkers.personalstack.agents.application.observability.FailureReasonLabel
import com.jorisjonkers.personalstack.agents.application.observability.ModeLabel
import com.jorisjonkers.personalstack.agents.application.observability.OperationLabel
import com.jorisjonkers.personalstack.agents.application.observability.OperationTelemetry
import com.jorisjonkers.personalstack.agents.application.observability.OutcomeLabel
import com.jorisjonkers.personalstack.agents.application.observability.RunnerReprovisionTelemetry
import java.io.IOException
import java.time.Duration

/**
* Records telemetry for runner session binding operations. Extracted from
* RunnerSessionBinder to keep that class below the TooManyFunctions threshold.
*/
internal class RunnerBindingMetrics(
private val telemetry: AgentsApiTelemetry,
) {
fun observeBinding(
operation: OperationLabel,
mode: ModeLabel,
block: () -> RunnerSessionBindingResult,
): RunnerSessionBindingResult {
val startedAt = System.nanoTime()
return runCatching(block)
.onSuccess { result ->
val (outcome, reason) =
when (result) {
is RunnerSessionBindingResult.Bound -> OutcomeLabel.SUCCESS to FailureReasonLabel.NONE
is RunnerSessionBindingResult.Conflict -> bindingConflict()
is RunnerSessionBindingResult.Unavailable ->
OutcomeLabel.FAILURE to FailureReasonLabel.UPSTREAM_UNAVAILABLE
}
record(operation, mode, outcome, reason, startedAt)
}.onFailure { ex ->
record(operation, mode, OutcomeLabel.FAILURE, reasonClass(ex), startedAt)
}.getOrThrow()
}

fun <T> observeStage(
operation: OperationLabel,
mode: ModeLabel,
outcome: (T) -> Pair<OutcomeLabel, FailureReasonLabel> = { OutcomeLabel.SUCCESS to FailureReasonLabel.NONE },
block: () -> T,
): T {
val startedAt = System.nanoTime()
return runCatching(block)
.onSuccess { result ->
val (resultOutcome, reason) = outcome(result)
record(operation, mode, resultOutcome, reason, startedAt)
}.onFailure { ex ->
record(operation, mode, OutcomeLabel.FAILURE, reasonClass(ex), startedAt)
}.getOrThrow()
}

fun recordReprovision(
outcome: OutcomeLabel,
reason: FailureReasonLabel,
startedAt: Long,
) {
telemetry.recordRunnerReprovision(RunnerReprovisionTelemetry(outcome = outcome, reason = reason))
record(
operation = OperationLabel.REPROVISION_RUNNER,
mode = ModeLabel.DURABLE,
outcome = outcome,
reason = reason,
startedAt = startedAt,
)
}

fun record(
operation: OperationLabel,
mode: ModeLabel,
outcome: OutcomeLabel,
reason: FailureReasonLabel,
startedAt: Long,
) {
telemetry.recordOperation(
OperationTelemetry(
operation = operation,
mode = mode,
outcome = outcome,
reason = reason,
duration = Duration.ofNanos((System.nanoTime() - startedAt).coerceAtLeast(0)),
),
)
}

fun bindingConflict(): Pair<OutcomeLabel, FailureReasonLabel> = OutcomeLabel.FAILURE to FailureReasonLabel.CAPACITY

fun reasonClass(ex: Throwable): FailureReasonLabel =
when (ex) {
is AgentSetupValidationException,
is IllegalArgumentException,
-> FailureReasonLabel.INVALID_REQUEST

is AgentRunnerUnavailableException -> FailureReasonLabel.UPSTREAM_UNAVAILABLE
is IOException -> FailureReasonLabel.IO_ERROR
else -> FailureReasonLabel.UNKNOWN
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
package com.jorisjonkers.personalstack.agents.application.sessionbinding

import com.jorisjonkers.personalstack.agents.application.exception.AgentRunnerUnavailableException
import com.jorisjonkers.personalstack.agents.application.observability.FailureReasonLabel
import com.jorisjonkers.personalstack.agents.application.observability.OutcomeLabel
import com.jorisjonkers.personalstack.agents.application.workspacerunner.RunnerSetupTarget
import com.jorisjonkers.personalstack.agents.domain.model.Workspace
import com.jorisjonkers.personalstack.agents.domain.port.AgentGatewayClient
import com.jorisjonkers.personalstack.agents.domain.port.AgentRunnerOrchestrator
import com.jorisjonkers.personalstack.agents.domain.port.WorkspaceRepository
import org.slf4j.LoggerFactory

/**
* Coordinates runner provisioning: scale-down, provision, await readiness,
* and record the reprovision telemetry. Extracted from RunnerSessionBinder to
* keep that class below the TooManyFunctions threshold.
*/
internal class RunnerProvisioningCoordinator(
private val orchestrator: AgentRunnerOrchestrator,
private val gateway: AgentGatewayClient,
private val workspaces: WorkspaceRepository,
private val metrics: RunnerBindingMetrics,
private val backoffInitialMs: Long,
) {
private val log = LoggerFactory.getLogger(RunnerProvisioningCoordinator::class.java)

fun forceProvisionAndWait(
workspace: Workspace,
target: RunnerSetupTarget,
runnerGeneration: Long,
): RunnerReady {
val startedAt = System.nanoTime()
return runCatching { provisionAndBuildReady(workspace, target, runnerGeneration) }
.onSuccess {
metrics.recordReprovision(OutcomeLabel.SUCCESS, FailureReasonLabel.NONE, startedAt)
}.onFailure { ex ->
val reason = metrics.reasonClass(ex)
metrics.recordReprovision(OutcomeLabel.FAILURE, reason, startedAt)
}.getOrThrow()
}

fun isRunnerReadyFor(
workspace: Workspace,
target: RunnerSetupTarget,
): Boolean {
val identity = target.spec.identity(workspace.runnerSetupGeneration)
return orchestrator.isReady(workspace, identity) && gateway.isReady(workspace)
}

private fun provisionAndBuildReady(
workspace: Workspace,
target: RunnerSetupTarget,
runnerGeneration: Long,
): RunnerReady {
val handle =
runCatching {
orchestrator.scaleDown(workspace)
orchestrator.provision(workspace, target.spec, runnerGeneration)
}.getOrElse { ex ->
log.warn("reprovision of workspace {} failed during scaleDown/provision", workspace.id.value, ex)
throw AgentRunnerUnavailableException(
workspaceId = workspace.id,
runnerStatus = "ReprovisionFailed",
cause = ex,
)
}
val repointed =
workspace.withPodInfo(
podName = handle.podName,
pvcName = handle.pvcName,
gatewayEndpoint = handle.gatewayEndpoint,
)
val saved = workspaces.save(repointed)
log.info("re-provisioned runner for workspace {} as pod {}", workspace.id.value, handle.podName)
return RunnerReady(
workspace = saved,
ready = awaitRunnerReady(saved, target, runnerGeneration),
provisioning =
RunnerProvisioningResult.Provisioned(
podName = handle.podName,
pvcName = handle.pvcName,
gatewayEndpoint = handle.gatewayEndpoint,
),
)
}

private fun awaitRunnerReady(
workspace: Workspace,
target: RunnerSetupTarget,
runnerGeneration: Long,
): Boolean {
val identity = target.spec.identity(runnerGeneration)
repeat(RunnerAgentSpawner.MAX_SPAWN_ATTEMPTS) { attempt ->
if (orchestrator.isReady(workspace, identity) && gateway.isReady(workspace)) return true
if (attempt < RunnerAgentSpawner.MAX_SPAWN_ATTEMPTS - 1) {
val sleepMs = backoffInitialMs * (attempt + 1)
if (sleepMs > 0) Thread.sleep(sleepMs)
}
}
return false
}
}

internal data class RunnerReady(
val workspace: Workspace,
// The runner was reprovisioned; `ready` is whether it became ready within
// the synchronous window. When false the session is left awaiting rebind.
val ready: Boolean,
val provisioning: RunnerProvisioningResult,
)
Loading