implemented 06-watcher.md: non-blocking poll cycle enqueueing changed/new/auto-build branches via the async executor, startup recovery, retention/worktree pruning, JSON auto-build slot state, watcher.pollInterval config, and watcher health state
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Fable 5
parent
3c9eeda5da
commit
a1db450bdb
@@ -126,6 +126,8 @@ class InitCommand(
|
||||
|
||||
# Controls the branch-polling loop.
|
||||
watcher:
|
||||
# delay between poll cycles (s/m/h/d suffix)
|
||||
pollInterval: 10s
|
||||
# max commit age for new origin branches to be pulled automatically
|
||||
newBranchMaxAge: 5d
|
||||
|
||||
|
||||
@@ -2,18 +2,25 @@ package de.hoennig.gittally.config
|
||||
|
||||
import java.time.Duration
|
||||
|
||||
/** Parses durations in the `watcher.newBranchMaxAge` format: `5d` (days) or `12h` (hours). */
|
||||
/**
|
||||
* Parses durations in the format used by `watcher.newBranchMaxAge` and `watcher.pollInterval`:
|
||||
* `5d` (days), `12h` (hours), `10m` (minutes), or `30s` (seconds).
|
||||
*/
|
||||
object DurationParser {
|
||||
private val pattern = Regex("""(\d+)([dh])""")
|
||||
private val pattern = Regex("""(\d+)([dhms])""")
|
||||
|
||||
fun parse(value: String): Duration {
|
||||
val match =
|
||||
pattern.matchEntire(value.trim())
|
||||
?: throw IllegalArgumentException("invalid duration '$value': expected <amount>d or <amount>h, e.g. 5d or 12h")
|
||||
?: throw IllegalArgumentException(
|
||||
"invalid duration '$value': expected <amount>d, <amount>h, <amount>m, or <amount>s, e.g. 5d or 10s",
|
||||
)
|
||||
val (amount, unit) = match.destructured
|
||||
return when (unit) {
|
||||
"d" -> Duration.ofDays(amount.toLong())
|
||||
else -> Duration.ofHours(amount.toLong())
|
||||
"h" -> Duration.ofHours(amount.toLong())
|
||||
"m" -> Duration.ofMinutes(amount.toLong())
|
||||
else -> Duration.ofSeconds(amount.toLong())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -41,6 +41,8 @@ data class ArtifactsConfig(
|
||||
)
|
||||
|
||||
data class WatcherConfig(
|
||||
/** Delay between poll cycles, e.g. `10s` or `1m`. */
|
||||
val pollInterval: String = "10s",
|
||||
val newBranchMaxAge: String = "5d",
|
||||
)
|
||||
|
||||
|
||||
@@ -151,6 +151,15 @@ class GitService(
|
||||
.trim()
|
||||
.ifEmpty { null }
|
||||
|
||||
/** The commit `refs/remotes/origin/[branch]` points at, or null when the branch is not on origin. */
|
||||
fun originHeadCommit(
|
||||
branch: String,
|
||||
workingDir: Path = Paths.get("."),
|
||||
): String? {
|
||||
val result = runner.run(listOf("git", "rev-parse", "--verify", "refs/remotes/origin/$branch"), workingDir)
|
||||
return if (result.isSuccess) result.stdout.trim() else null
|
||||
}
|
||||
|
||||
fun headCommit(workingDir: Path = Paths.get(".")): String =
|
||||
runner
|
||||
.runOrThrow(listOf("git", "rev-parse", "HEAD"), workingDir)
|
||||
|
||||
@@ -0,0 +1,98 @@
|
||||
package de.hoennig.gittally.watcher
|
||||
|
||||
import com.fasterxml.jackson.databind.DeserializationFeature
|
||||
import com.fasterxml.jackson.databind.ObjectMapper
|
||||
import com.fasterxml.jackson.databind.SerializationFeature
|
||||
import com.fasterxml.jackson.module.kotlin.readValue
|
||||
import com.fasterxml.jackson.module.kotlin.registerKotlinModule
|
||||
import org.slf4j.LoggerFactory
|
||||
import java.nio.file.Files
|
||||
import java.nio.file.Path
|
||||
import java.nio.file.StandardCopyOption
|
||||
import java.time.LocalDate
|
||||
import java.time.LocalTime
|
||||
import java.time.format.DateTimeParseException
|
||||
|
||||
/** One recorded auto-build trigger: [branch] was enqueued for the [slot] (UTC `HH:MM`) of [date] (ISO). */
|
||||
data class AutoBuildTrigger(
|
||||
val branch: String,
|
||||
val date: String,
|
||||
val slot: String,
|
||||
)
|
||||
|
||||
/** Auto-build time slot matching (UTC `HH:MM`), like legacy `auto_build_check`. */
|
||||
object AutoBuildSlots {
|
||||
private val log = LoggerFactory.getLogger(AutoBuildSlots::class.java)
|
||||
|
||||
/** The latest valid slot at or before [now], or null when no slot is due yet today. */
|
||||
fun latestDueSlot(
|
||||
times: List<String>,
|
||||
now: LocalTime,
|
||||
): String? =
|
||||
times
|
||||
.mapNotNull { slot ->
|
||||
try {
|
||||
LocalTime.parse(slot.trim()) to slot
|
||||
} catch (_: DateTimeParseException) {
|
||||
log.warn("skipping invalid auto-build time slot '{}': expected HH:MM", slot)
|
||||
null
|
||||
}
|
||||
}.filter { (parsed, _) -> !parsed.isAfter(now) }
|
||||
.maxByOrNull { (parsed, _) -> parsed }
|
||||
?.second
|
||||
}
|
||||
|
||||
/**
|
||||
* Persists which auto-build slots already triggered as a JSON file,
|
||||
* e.g. `.git/gittally/auto-builds.json` (replaces the legacy `auto-builds.tsv`).
|
||||
* Entries of past days are dropped on write, so the file never grows unbounded.
|
||||
*/
|
||||
class FileAutoBuildState(
|
||||
private val file: Path,
|
||||
) {
|
||||
private val log = LoggerFactory.getLogger(FileAutoBuildState::class.java)
|
||||
|
||||
private val json =
|
||||
ObjectMapper()
|
||||
.registerKotlinModule()
|
||||
.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false)
|
||||
.configure(SerializationFeature.INDENT_OUTPUT, true)
|
||||
|
||||
fun isTriggered(
|
||||
branch: String,
|
||||
date: LocalDate,
|
||||
slot: String,
|
||||
): Boolean = AutoBuildTrigger(branch, date.toString(), slot) in load()
|
||||
|
||||
fun markTriggered(
|
||||
branch: String,
|
||||
date: LocalDate,
|
||||
slot: String,
|
||||
) {
|
||||
val current = load().filter { it.date == date.toString() }
|
||||
save(current + AutoBuildTrigger(branch, date.toString(), slot))
|
||||
}
|
||||
|
||||
private fun load(): List<AutoBuildTrigger> {
|
||||
if (!Files.exists(file)) {
|
||||
return emptyList()
|
||||
}
|
||||
return try {
|
||||
json.readValue<List<AutoBuildTrigger>>(file.toFile())
|
||||
} catch (e: Exception) {
|
||||
log.warn("ignoring unreadable auto-builds file {}: {}", file, e.message)
|
||||
emptyList()
|
||||
}
|
||||
}
|
||||
|
||||
private fun save(triggers: List<AutoBuildTrigger>) {
|
||||
Files.createDirectories(file.parent)
|
||||
val tempFile = Files.createTempFile(file.parent, file.fileName.toString(), ".tmp")
|
||||
try {
|
||||
json.writeValue(tempFile.toFile(), triggers)
|
||||
Files.move(tempFile, file, StandardCopyOption.ATOMIC_MOVE, StandardCopyOption.REPLACE_EXISTING)
|
||||
} finally {
|
||||
Files.deleteIfExists(tempFile)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,270 @@
|
||||
package de.hoennig.gittally.watcher
|
||||
|
||||
import de.hoennig.gittally.build.ArtifactKeys
|
||||
import de.hoennig.gittally.build.ArtifactStore
|
||||
import de.hoennig.gittally.build.BuildExecutor
|
||||
import de.hoennig.gittally.build.BuildResultRepository
|
||||
import de.hoennig.gittally.build.BuildStatus
|
||||
import de.hoennig.gittally.build.GitWorktreeWorkspaces
|
||||
import de.hoennig.gittally.config.ConfigLoader
|
||||
import de.hoennig.gittally.config.DurationParser
|
||||
import de.hoennig.gittally.config.GitTallyConfig
|
||||
import de.hoennig.gittally.git.GitService
|
||||
import org.slf4j.LoggerFactory
|
||||
import org.springframework.stereotype.Service
|
||||
import java.nio.file.Files
|
||||
import java.nio.file.Path
|
||||
import java.nio.file.Paths
|
||||
import java.time.Clock
|
||||
import java.time.LocalDate
|
||||
import java.time.LocalTime
|
||||
import java.time.ZoneOffset
|
||||
import java.util.concurrent.Executors
|
||||
import java.util.concurrent.ScheduledExecutorService
|
||||
import java.util.concurrent.TimeUnit
|
||||
|
||||
/**
|
||||
* Replaces the legacy blocking main loop: a non-blocking fixed-delay poll cycle that
|
||||
* fetches origin, enqueues due branches via the async [BuildExecutor], and prunes
|
||||
* retention — it never waits for a build and never touches the primary checkout.
|
||||
* The loop only runs after an explicit [start] (server/watch mode, step 07);
|
||||
* nothing is scheduled during CLI commands or tests.
|
||||
*/
|
||||
@Service
|
||||
class Watcher(
|
||||
private val gitService: GitService,
|
||||
private val buildExecutor: BuildExecutor,
|
||||
private val repository: BuildResultRepository,
|
||||
private val artifactStore: ArtifactStore,
|
||||
private val configLoader: ConfigLoader,
|
||||
private val clock: Clock,
|
||||
) {
|
||||
private val log = LoggerFactory.getLogger(Watcher::class.java)
|
||||
|
||||
private var scheduler: ScheduledExecutorService? = null
|
||||
|
||||
@Volatile
|
||||
private var state = WatcherState()
|
||||
|
||||
fun state(): WatcherState = state
|
||||
|
||||
/**
|
||||
* Runs the startup recovery and schedules the poll loop with the fixed delay
|
||||
* `watcher.pollInterval`; the first poll runs immediately.
|
||||
*/
|
||||
@Synchronized
|
||||
fun start(workingDir: Path = Paths.get(".")) {
|
||||
check(scheduler == null) { "watcher is already running" }
|
||||
recoverOnStartup(workingDir)
|
||||
val interval = DurationParser.parse(configLoader.load(workingDir).watcher.pollInterval)
|
||||
scheduler =
|
||||
Executors
|
||||
.newSingleThreadScheduledExecutor { runnable ->
|
||||
Thread(runnable, "gittally-watcher").apply { isDaemon = true }
|
||||
}.also {
|
||||
it.scheduleWithFixedDelay({ pollSafely(workingDir) }, 0, interval.toMillis(), TimeUnit.MILLISECONDS)
|
||||
}
|
||||
state = state.copy(running = true)
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun stop() {
|
||||
scheduler?.shutdownNow()
|
||||
scheduler = null
|
||||
state = state.copy(running = false)
|
||||
}
|
||||
|
||||
/**
|
||||
* Port of the legacy startup recovery: best-effort fetch, mark stale RUNNING and
|
||||
* superseded PENDING builds as INTERRUPTED, then re-enqueue every branch whose
|
||||
* latest build never finished and which still exists on origin.
|
||||
*/
|
||||
fun recoverOnStartup(workingDir: Path = Paths.get(".")) {
|
||||
try {
|
||||
gitService.fetchOrigin(workingDir)
|
||||
} catch (e: Exception) {
|
||||
log.warn("startup fetch failed; recovering from the last known origin state: {}", e.message)
|
||||
}
|
||||
repository.markStaleRunningAsInterrupted().forEach {
|
||||
log.info("marked stale build of branch {} as interrupted", it.branch)
|
||||
}
|
||||
val restartable =
|
||||
repository
|
||||
.latestPerBranch()
|
||||
.filter { it.status == BuildStatus.INTERRUPTED || it.status == BuildStatus.PENDING }
|
||||
for (result in restartable) {
|
||||
val commit = gitService.originHeadCommit(result.branch, workingDir)
|
||||
if (commit == null) {
|
||||
log.info("not restarting build of branch {}: branch is gone from origin", result.branch)
|
||||
continue
|
||||
}
|
||||
if (result.status == BuildStatus.PENDING) {
|
||||
// the executor queue did not survive the restart; the re-enqueued build supersedes the stale entry
|
||||
repository.updateByArtifactKey(result.artifactKey) { it.copy(status = BuildStatus.INTERRUPTED) }
|
||||
}
|
||||
log.info("restarting unfinished build of branch {}", result.branch)
|
||||
buildExecutor.startBuild(result.branch, commit, workingDir)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* One poll cycle, never blocking on a build: fetch origin (on failure: log, expose
|
||||
* in [state], retry next cycle), enqueue due branches — changed local branches
|
||||
* first, then recent new origin branches, then due auto-build slots — and finally
|
||||
* prune results, artifacts, and worktrees of branches gone from origin.
|
||||
*/
|
||||
fun poll(workingDir: Path = Paths.get(".")) {
|
||||
val startedAt = clock.instant()
|
||||
try {
|
||||
gitService.fetchOrigin(workingDir)
|
||||
} catch (e: Exception) {
|
||||
log.warn("fetching origin failed; retrying next cycle: {}", e.message)
|
||||
state = state.copy(lastPollAt = startedAt, lastFetchError = e.message ?: e.javaClass.simpleName)
|
||||
return
|
||||
}
|
||||
val config = configLoader.load(workingDir)
|
||||
val originBranches = gitService.originBranches(workingDir)
|
||||
enqueueDueBranches(config, originBranches.toSet(), workingDir)
|
||||
prune(config, originBranches, workingDir)
|
||||
state =
|
||||
state.copy(
|
||||
lastPollAt = startedAt,
|
||||
lastFetchError = null,
|
||||
lastPollError = null,
|
||||
queuedBranches =
|
||||
repository
|
||||
.latestPerBranch()
|
||||
.filter { it.status == BuildStatus.PENDING || it.status == BuildStatus.RUNNING }
|
||||
.map { it.branch },
|
||||
)
|
||||
}
|
||||
|
||||
private fun pollSafely(workingDir: Path) {
|
||||
try {
|
||||
poll(workingDir)
|
||||
} catch (e: Exception) {
|
||||
log.error("poll cycle failed", e)
|
||||
state = state.copy(lastPollAt = clock.instant(), lastPollError = e.message ?: e.javaClass.simpleName)
|
||||
}
|
||||
}
|
||||
|
||||
private fun enqueueDueBranches(
|
||||
config: GitTallyConfig,
|
||||
originBranches: Set<String>,
|
||||
workingDir: Path,
|
||||
) {
|
||||
val changedLocal =
|
||||
gitService
|
||||
.localBranches(workingDir)
|
||||
.filter { it in originBranches && gitService.hasNewCommits(it, workingDir) }
|
||||
val newOrigin =
|
||||
gitService.newOriginBranches(DurationParser.parse(config.watcher.newBranchMaxAge), workingDir)
|
||||
for (branch in (changedLocal + newOrigin).distinct()) {
|
||||
startBuildIfDue(branch, allowSameCommit = false, workingDir = workingDir)
|
||||
}
|
||||
enqueueAutoBuilds(config, originBranches, workingDir)
|
||||
}
|
||||
|
||||
/**
|
||||
* Enqueues a build of the branch's origin head unless one is already pending or
|
||||
* running, or that commit was already built. Builds run detached in worktrees and
|
||||
* never move local branch refs, so "already built" is tracked via the result
|
||||
* repository, not by resetting the local ref like legacy. A new commit for a
|
||||
* branch that is still pending/running waits for a later cycle (queue-behind).
|
||||
*/
|
||||
private fun startBuildIfDue(
|
||||
branch: String,
|
||||
allowSameCommit: Boolean,
|
||||
workingDir: Path,
|
||||
): Boolean {
|
||||
val latest = repository.latestFor(branch)
|
||||
if (latest?.status == BuildStatus.PENDING || latest?.status == BuildStatus.RUNNING) {
|
||||
return false
|
||||
}
|
||||
val commit = gitService.originHeadCommit(branch, workingDir) ?: return false
|
||||
if (!allowSameCommit && latest?.commit == commit) {
|
||||
return false
|
||||
}
|
||||
log.info("enqueueing build of branch {} at commit {}", branch, commit)
|
||||
buildExecutor.startBuild(branch, commit, workingDir)
|
||||
return true
|
||||
}
|
||||
|
||||
private fun enqueueAutoBuilds(
|
||||
config: GitTallyConfig,
|
||||
originBranches: Set<String>,
|
||||
workingDir: Path,
|
||||
) {
|
||||
val autoBuildBranches =
|
||||
config.branches.filter { (branch, branchConfig) ->
|
||||
branch != "default" && branchConfig.autoBuild.enabled
|
||||
}
|
||||
if (autoBuildBranches.isEmpty()) {
|
||||
return
|
||||
}
|
||||
val autoBuildState = FileAutoBuildState(workingDir.resolve(AUTO_BUILDS_FILE))
|
||||
val now = clock.instant()
|
||||
val today = LocalDate.ofInstant(now, ZoneOffset.UTC)
|
||||
val timeOfDay = LocalTime.ofInstant(now, ZoneOffset.UTC)
|
||||
for ((branch, branchConfig) in autoBuildBranches) {
|
||||
val slot = AutoBuildSlots.latestDueSlot(branchConfig.autoBuild.times, timeOfDay) ?: continue
|
||||
if (autoBuildState.isTriggered(branch, today, slot)) {
|
||||
continue
|
||||
}
|
||||
if (branch !in originBranches) {
|
||||
log.warn("skipping auto build of branch {}: branch is not on origin", branch)
|
||||
continue
|
||||
}
|
||||
// rebuilding the already-built commit is the point of an auto build
|
||||
if (startBuildIfDue(branch, allowSameCommit = true, workingDir = workingDir)) {
|
||||
autoBuildState.markTriggered(branch, today, slot)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Results first, then artifacts of dropped results, then worktrees of branches gone from origin. */
|
||||
private fun prune(
|
||||
config: GitTallyConfig,
|
||||
originBranches: List<String>,
|
||||
workingDir: Path,
|
||||
) {
|
||||
repository.prune(originBranches, config.artifacts.retentionPerBranch)
|
||||
artifactStore.prune(repository.history())
|
||||
pruneWorktrees(originBranches, workingDir)
|
||||
}
|
||||
|
||||
private fun pruneWorktrees(
|
||||
originBranches: List<String>,
|
||||
workingDir: Path,
|
||||
) {
|
||||
val worktreesDir = workingDir.resolve(GitWorktreeWorkspaces.WORKTREES_DIR)
|
||||
if (!Files.isDirectory(worktreesDir)) {
|
||||
return
|
||||
}
|
||||
val keep = originBranches.map { ArtifactKeys.branchKey(it) }.toMutableSet()
|
||||
// never delete under a build that is still queued or executing
|
||||
buildExecutor.currentBuilds().forEach { keep += ArtifactKeys.branchKey(it.branch) }
|
||||
repository
|
||||
.latestPerBranch()
|
||||
.filter { it.status == BuildStatus.PENDING || it.status == BuildStatus.RUNNING }
|
||||
.forEach { keep += ArtifactKeys.branchKey(it.branch) }
|
||||
var removed = false
|
||||
Files.list(worktreesDir).use { entries ->
|
||||
entries.forEach { entry ->
|
||||
if (Files.isDirectory(entry) && entry.fileName.toString() !in keep) {
|
||||
log.info("removing worktree of branch gone from origin: {}", entry.fileName)
|
||||
entry.toFile().deleteRecursively()
|
||||
removed = true
|
||||
}
|
||||
}
|
||||
}
|
||||
if (removed) {
|
||||
gitService.worktreePrune(workingDir)
|
||||
}
|
||||
}
|
||||
|
||||
companion object {
|
||||
/** Auto-build trigger state next to the build results (replaces legacy `auto-builds.tsv`). */
|
||||
const val AUTO_BUILDS_FILE = ".git/gittally/auto-builds.json"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,12 @@
|
||||
package de.hoennig.gittally.watcher
|
||||
|
||||
import org.springframework.context.annotation.Bean
|
||||
import org.springframework.context.annotation.Configuration
|
||||
import java.time.Clock
|
||||
|
||||
@Configuration
|
||||
class WatcherConfiguration {
|
||||
/** UTC clock, injectable in tests; auto-build slots are UTC `HH:MM` like legacy. */
|
||||
@Bean
|
||||
fun clock(): Clock = Clock.systemUTC()
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
package de.hoennig.gittally.watcher
|
||||
|
||||
import java.time.Instant
|
||||
|
||||
/** Observable watcher health for status endpoints (step 07) and the UI (step 08). */
|
||||
data class WatcherState(
|
||||
/** Whether the poll loop is scheduled. */
|
||||
val running: Boolean = false,
|
||||
/** When the last poll cycle started, successful or not. */
|
||||
val lastPollAt: Instant? = null,
|
||||
/** Why the last `fetchOrigin` failed; null after a successful fetch. */
|
||||
val lastFetchError: String? = null,
|
||||
/** Why the last poll cycle crashed after a successful fetch; null after a clean cycle. */
|
||||
val lastPollError: String? = null,
|
||||
/** Branches whose latest build was PENDING or RUNNING at the end of the last poll. */
|
||||
val queuedBranches: List<String> = emptyList(),
|
||||
)
|
||||
Reference in New Issue
Block a user