From e77b2960f9e5c3d4c201bbe21af54dbf0f7e3920 Mon Sep 17 00:00:00 2001 From: kami Date: Wed, 24 Jun 2026 17:02:03 +0000 Subject: [PATCH] feat(tasks): dependency-aware work graph (blockers, ready, blocking) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The DEPENDS_ON/BLOCKS link types were inert metadata — nothing reasoned about them. TaskGraph turns them into answerable questions: what a task waits on (its DEPENDS_ON targets plus any task that BLOCKS it), whether it is ready (TODO with no unmet blocker — a terminal dependency stops blocking), and what finishing it would unblock. Readiness surfaces workable tasks to claim; it never assigns (honouring the rejected /tasks/next scheduler). Surfaced across the stack: - TaskService.ready / blockers / blocking. - Context bundle: a `blocked` flag + `blocked_by` list, rendered as a prominent "BLOCKED — waiting on ..." line, so an agent calling task_context learns it should wait before charging at the task. - GET /tasks?ready=true, GET /tasks/{id}/blockers; correx task ready / blockers. Cross-project dependencies are scoped to a project for now (documented in TaskGraph). Tests cover both link directions, terminal-dependency clearing, the ready filter, the reverse blocking direction, the bundle flag, and REST. Co-Authored-By: Claude Opus 4.8 --- .../correx/apps/cli/commands/TaskCommand.kt | 26 ++++++++ .../correx/apps/server/routes/TaskRoutes.kt | 12 ++++ .../apps/server/routes/TaskRoutesTest.kt | 29 ++++++++ .../correx/core/tasks/TaskContextAssembler.kt | 26 ++++++++ .../correx/core/tasks/TaskContextBundle.kt | 6 ++ .../kotlin/com/correx/core/tasks/TaskGraph.kt | 54 +++++++++++++++ .../com/correx/core/tasks/TaskService.kt | 20 ++++++ .../core/tasks/TaskContextAssemblerTest.kt | 21 ++++++ .../com/correx/core/tasks/TaskGraphTest.kt | 66 +++++++++++++++++++ .../com/correx/core/tasks/TaskServiceTest.kt | 16 +++++ 10 files changed, 276 insertions(+) create mode 100644 core/tasks/src/main/kotlin/com/correx/core/tasks/TaskGraph.kt create mode 100644 core/tasks/src/test/kotlin/com/correx/core/tasks/TaskGraphTest.kt diff --git a/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/TaskCommand.kt b/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/TaskCommand.kt index 32dd5a21..a29521ce 100644 --- a/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/TaskCommand.kt +++ b/apps/cli/src/main/kotlin/com/correx/apps/cli/commands/TaskCommand.kt @@ -89,8 +89,10 @@ class TaskCommand : CliktCommand(name = "task") { init { subcommands( TaskListCommand(), + TaskReadyCommand(), TaskSearchCommand(), TaskShowCommand(), + TaskBlockersCommand(), TaskHistoryCommand(), TaskCreateCommand(), TaskClaimCommand(), @@ -133,6 +135,30 @@ class TaskListCommand : TaskHttpCommand("list") { } } +class TaskReadyCommand : TaskHttpCommand("ready") { + private val project by option("--project", help = "Limit to a project") + + override fun run() = http { client -> + val query = project?.let { "&project=${enc(it)}" }.orEmpty() + val raw = client.get("${baseUrl()}/tasks?ready=true$query").bodyAsText() + if (jsonOut()) println(raw) else println(renderTaskList(taskJson.decodeFromString(raw))) + } +} + +class TaskBlockersCommand : TaskHttpCommand("blockers") { + private val id by argument("TASK_ID") + + override fun run() = http { client -> + val resp = client.get("${baseUrl()}/tasks/$id/blockers") + if (resp.status == HttpStatusCode.NotFound) { + System.err.println("No such task: $id") + } else { + val raw = resp.bodyAsText() + if (jsonOut()) println(raw) else println(renderTaskList(taskJson.decodeFromString(raw))) + } + } +} + class TaskSearchCommand : TaskHttpCommand("search") { private val query by argument("QUERY", help = "Search text (quote multi-word queries)") private val project by option("--project", help = "Limit to a project") diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/routes/TaskRoutes.kt b/apps/server/src/main/kotlin/com/correx/apps/server/routes/TaskRoutes.kt index 2a6da8c6..81a0d994 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/routes/TaskRoutes.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/routes/TaskRoutes.kt @@ -126,6 +126,12 @@ fun Route.taskRoutes(module: ServerModule) { private fun Route.taskCollectionRoutes(service: TaskService) { get { + // ?ready=true short-circuits to the dependency-aware workable set (TODO, no unmet blockers). + if (call.request.queryParameters["ready"]?.toBoolean() == true) { + val readyProject = call.request.queryParameters["project"] + val workable = service.ready(readyProject?.let { ProjectId(it) }) + return@get call.respond(workable.map { it.toResponse() }) + } val statusRaw = call.request.queryParameters["status"] val statusFilter = statusRaw?.let { parseEnum(it) ?: return@get call.respond(HttpStatusCode.BadRequest, "Invalid status: $it") @@ -209,6 +215,12 @@ private fun Route.taskCrudRoutes(service: TaskService, assembler: TaskContextAss if (events.isEmpty()) return@get call.respond(HttpStatusCode.NotFound, "Task not found") call.respondText(TaskHistory.render(events), ContentType.Text.Plain) } + // Unfinished tasks this one is waiting on (its dependencies + inbound blocks). + get("/blockers") { + val id = taskId() ?: return@get call.respond(HttpStatusCode.BadRequest, "Missing task id") + if (service.getTask(id) == null) return@get call.respond(HttpStatusCode.NotFound, "Task not found") + call.respond(service.blockers(id).map { it.toResponse() }) + } } private fun Route.taskTransitionRoutes(service: TaskService) { diff --git a/apps/server/src/test/kotlin/com/correx/apps/server/routes/TaskRoutesTest.kt b/apps/server/src/test/kotlin/com/correx/apps/server/routes/TaskRoutesTest.kt index 75062afb..a56f713e 100644 --- a/apps/server/src/test/kotlin/com/correx/apps/server/routes/TaskRoutesTest.kt +++ b/apps/server/src/test/kotlin/com/correx/apps/server/routes/TaskRoutesTest.kt @@ -256,6 +256,35 @@ class TaskRoutesTest { } } + @Test + fun `ready returns unblocked TODO tasks and blockers explains a wait`(@TempDir tempDir: Path) { + testApplication { + application { configureServer(buildModule(tempDir)) } + val client = createClient {} + client.createTask() // demo-1, TODO + client.post("/tasks") { + contentType(ContentType.Application.Json) + setBody("""{"project":"demo","title":"second","goal":"g"}""") + } + // demo-2 depends on demo-1 (id infers to a TASK target). + client.post("/tasks/demo-2/links") { + contentType(ContentType.Application.Json) + setBody("""{"targetId":"demo-1","type":"DEPENDS_ON"}""") + } + + val ready = client.get("/tasks?ready=true") + assertEquals(HttpStatusCode.OK, ready.status) + val rows = testJson.parseToJsonElement(ready.bodyAsText()) as JsonArray + assertEquals(1, rows.size) + assertEquals("demo-1", (rows[0] as JsonObject)["id"]!!.jsonPrimitive.content) + + val blockers = client.get("/tasks/demo-2/blockers") + assertEquals(HttpStatusCode.OK, blockers.status) + val blk = testJson.parseToJsonElement(blockers.bodyAsText()) as JsonArray + assertEquals("demo-1", (blk.single() as JsonObject)["id"]!!.jsonPrimitive.content) + } + } + @Test fun `GET history renders the lifecycle timeline`(@TempDir tempDir: Path) { testApplication { diff --git a/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextAssembler.kt b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextAssembler.kt index b7bd3b0b..fd8eaeab 100644 --- a/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextAssembler.kt +++ b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextAssembler.kt @@ -1,6 +1,7 @@ package com.correx.core.tasks import com.correx.core.events.types.TaskId +import com.correx.core.events.types.TaskLinkType import com.correx.core.events.types.TaskTargetKind /** @@ -38,6 +39,8 @@ class TaskContextAssembler( } } + val blockedBy = computeBlockedBy(taskId, state) + return TaskContextBundle( id = taskId.value, key = state.key, @@ -46,6 +49,8 @@ class TaskContextAssembler( goal = state.goal, acceptanceCriteria = state.acceptanceCriteria, relevantFiles = state.affectedPaths, + blocked = blockedBy.isNotEmpty(), + blockedBy = blockedBy, dependencies = resolved.dependencies, documents = resolved.documents, artifacts = resolved.artifacts, @@ -124,6 +129,27 @@ class TaskContextAssembler( } } + /** + * The task's unfinished dependency blockers, resolved over its project board (see [TaskGraph]), + * each labelled by the direction of the edge. Empty when the project board can't be read. + */ + private fun computeBlockedBy(taskId: TaskId, state: TaskState): List { + val board = state.projectId?.let { runCatching { service.list(it) }.getOrNull() } ?: return emptyList() + return TaskGraph.unmetBlockers(Task(taskId, state), board).map { b -> + val viaDependsOn = state.links.any { + it.targetKind == TaskTargetKind.TASK && + it.type == TaskLinkType.DEPENDS_ON && + it.targetId == b.taskId.value + } + RelatedTask( + id = b.taskId.value, + status = b.state.status.name, + title = b.state.title, + link = if (viaDependsOn) "DEPENDS_ON" else "BLOCKS", + ) + } + } + private suspend fun retrieveKnowledge(state: TaskState): List { val active = retriever ?: return emptyList() val query = listOfNotNull(state.title, state.goal).joinToString(" ").trim() diff --git a/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextBundle.kt b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextBundle.kt index db2090e6..9b809bb4 100644 --- a/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextBundle.kt +++ b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskContextBundle.kt @@ -16,6 +16,10 @@ data class TaskContextBundle( val goal: String?, val acceptanceCriteria: List, val relevantFiles: List, + // Dependency-derived readiness (distinct from an explicit BLOCKED status): blockedBy holds the + // unfinished tasks this one is waiting on, so an agent calling task_context knows to wait. + val blocked: Boolean = false, + val blockedBy: List = emptyList(), val dependencies: List, val documents: List, val artifacts: List = emptyList(), @@ -28,6 +32,8 @@ data class TaskContextBundle( fun render(): String = buildString { appendLine("task $id [$status] ${title.orEmpty()}".trimEnd()) goal?.let { appendLine("goal: $it") } + if (blocked) appendLine("BLOCKED — waiting on ${blockedBy.size} unfinished dependency(ies):") + appendList("blocked_by", blockedBy) { " - ${it.id} [${it.status}] ${it.title.orEmpty()} (${it.link})".trimEnd() } appendList("acceptance_criteria", acceptanceCriteria) { " - $it" } if (relevantFiles.isNotEmpty()) appendLine("relevant_files: ${relevantFiles.joinToString(", ")}") appendList("dependencies", dependencies) { " - ${it.id} [${it.status}] ${it.title.orEmpty()} (${it.link})".trimEnd() } diff --git a/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskGraph.kt b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskGraph.kt new file mode 100644 index 00000000..2a2e05b6 --- /dev/null +++ b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskGraph.kt @@ -0,0 +1,54 @@ +package com.correx.core.tasks + +import com.correx.core.events.types.TaskLinkType +import com.correx.core.events.types.TaskTargetKind + +/** + * Dependency reasoning over the task work graph — turns the `DEPENDS_ON`/`BLOCKS` link types from + * inert metadata into answerable questions: what is a task waiting on, what is it ready to be + * worked, and what does finishing it unblock. + * + * Two directions express the same edge, so both count: a `DEPENDS_ON B` link on A means A waits on + * B; a `BLOCKS A` link on B means the same. A blocker is "unmet" until it reaches a terminal + * status ([FINISHED]); a CANCELLED dependency no longer blocks. All reasoning is over the supplied + * board, so callers scope it (one project, or the whole cross-project set); a link to a task absent + * from the board is treated as already-resolved rather than an unknown blocker. + */ +object TaskGraph { + + private val FINISHED = setOf(TaskStatus.DONE, TaskStatus.CANCELLED) + + /** Tasks that must finish before [task] proceeds: its `DEPENDS_ON` targets plus tasks that `BLOCKS` it. */ + fun blockers(task: Task, board: Collection): List { + val byId = board.associateBy { it.taskId.value } + val dependsOn = task.linkTargets(TaskLinkType.DEPENDS_ON).mapNotNull { byId[it] } + val inbound = board.filter { it.hasTaskLink(TaskLinkType.BLOCKS, task.taskId.value) } + return (dependsOn + inbound).distinctBy { it.taskId.value } + } + + /** The blockers that have not finished — the concrete reason [task] cannot start yet. */ + fun unmetBlockers(task: Task, board: Collection): List = + blockers(task, board).filter { it.state.status !in FINISHED } + + /** True when [task] is TODO and has no unmet blockers — workable now (to be claimed, not assigned). */ + fun isReady(task: Task, board: Collection): Boolean = + task.state.status == TaskStatus.TODO && unmetBlockers(task, board).isEmpty() + + /** Every workable task on the board, in board order. */ + fun ready(board: Collection): List = + board.filter { isReady(it, board) } + + /** Tasks [task] is holding up: its own `BLOCKS` targets plus tasks that `DEPENDS_ON` it. */ + fun blocking(task: Task, board: Collection): List { + val byId = board.associateBy { it.taskId.value } + val downstream = task.linkTargets(TaskLinkType.BLOCKS).mapNotNull { byId[it] } + val dependents = board.filter { it.hasTaskLink(TaskLinkType.DEPENDS_ON, task.taskId.value) } + return (downstream + dependents).distinctBy { it.taskId.value } + } + + private fun Task.linkTargets(type: TaskLinkType): List = + state.links.filter { it.targetKind == TaskTargetKind.TASK && it.type == type }.map { it.targetId } + + private fun Task.hasTaskLink(type: TaskLinkType, targetId: String): Boolean = + state.links.any { it.targetKind == TaskTargetKind.TASK && it.type == type && it.targetId == targetId } +} diff --git a/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskService.kt b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskService.kt index 9413b057..2f83006b 100644 --- a/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskService.kt +++ b/core/tasks/src/main/kotlin/com/correx/core/tasks/TaskService.kt @@ -84,6 +84,26 @@ class TaskService( fun search(query: String, projectId: ProjectId? = null): List = TaskSearch.search(projectId?.let { list(it) } ?: listAll(), query) + /** + * Tasks ready to be picked up: TODO with every dependency satisfied (see [TaskGraph]). Surfaces + * workable tasks for an agent or human to claim — it never assigns. Scoped to one project when + * given, else the whole cross-project board. + */ + fun ready(projectId: ProjectId? = null): List = + TaskGraph.ready(projectId?.let { list(it) } ?: listAll()) + + /** Unfinished tasks that must complete before [taskId] can proceed (resolved within its project). */ + fun blockers(taskId: TaskId): List { + val task = getTask(taskId) ?: return emptyList() + return TaskGraph.unmetBlockers(task, list(projectOf(taskId))) + } + + /** Tasks that [taskId] is holding up — what completing it would unblock (within its project). */ + fun blocking(taskId: TaskId): List { + val task = getTask(taskId) ?: return emptyList() + return TaskGraph.blocking(task, list(projectOf(taskId))) + } + suspend fun createTask( projectId: ProjectId, title: String, diff --git a/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskContextAssemblerTest.kt b/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskContextAssemblerTest.kt index 13f623f1..03e45880 100644 --- a/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskContextAssemblerTest.kt +++ b/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskContextAssemblerTest.kt @@ -53,6 +53,27 @@ class TaskContextAssemblerTest { assertEquals("kickoff", bundle.notes.single().body) } + @Test + fun `bundle flags an unfinished dependency as a blocker and clears it when done`() = runBlocking { + val dep = service.createTask(project, "Design auth model", "model") // auth-1 + service.claim(dep.taskId, "claude-opus") // auth-1 IN_PROGRESS + val task = service.createTask(project, "Implement JWT refresh", "auth") // auth-2 + service.link(task.taskId, dep.taskId.value, TaskLinkType.DEPENDS_ON, TaskTargetKind.TASK) + + val blocked = assembler.assemble(task.taskId)!! + assertTrue(blocked.blocked) + assertEquals("auth-1", blocked.blockedBy.single().id) + assertEquals("IN_PROGRESS", blocked.blockedBy.single().status) + assertEquals("DEPENDS_ON", blocked.blockedBy.single().link) + assertTrue(blocked.render().contains("BLOCKED")) + + service.submitForReview(dep.taskId) + service.complete(dep.taskId) + val cleared = assembler.assemble(task.taskId)!! + assertTrue(!cleared.blocked) + assertTrue(cleared.blockedBy.isEmpty()) + } + @Test fun `render produces a compact labelled bundle`() = runBlocking { val task = service.createTask(project, "JWT refresh", "stay authed", listOf("rotates")) diff --git a/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskGraphTest.kt b/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskGraphTest.kt new file mode 100644 index 00000000..320628b2 --- /dev/null +++ b/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskGraphTest.kt @@ -0,0 +1,66 @@ +package com.correx.core.tasks + +import com.correx.core.events.types.ProjectId +import com.correx.core.events.types.TaskId +import com.correx.core.events.types.TaskLinkType +import com.correx.core.events.types.TaskTargetKind +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertFalse +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test + +class TaskGraphTest { + + private fun task(id: String, status: TaskStatus = TaskStatus.TODO, links: List = emptyList()) = + Task(TaskId(id), TaskState(status = status, key = id, projectId = ProjectId("p"), title = id, links = links)) + + private fun dependsOn(target: String) = TaskLink(target, TaskLinkType.DEPENDS_ON, TaskTargetKind.TASK) + private fun blocks(target: String) = TaskLink(target, TaskLinkType.BLOCKS, TaskTargetKind.TASK) + + @Test + fun `an unfinished depends_on target blocks the task`() { + val a = task("p-1", links = listOf(dependsOn("p-2"))) + val b = task("p-2", status = TaskStatus.IN_PROGRESS) + val board = listOf(a, b) + + assertEquals(listOf("p-2"), TaskGraph.unmetBlockers(a, board).map { it.taskId.value }) + assertFalse(TaskGraph.isReady(a, board)) + } + + @Test + fun `a finished dependency no longer blocks`() { + val a = task("p-1", links = listOf(dependsOn("p-2"))) + val done = task("p-2", status = TaskStatus.DONE) + val board = listOf(a, done) + + assertTrue(TaskGraph.unmetBlockers(a, board).isEmpty()) + assertTrue(TaskGraph.isReady(a, board)) + } + + @Test + fun `an inbound BLOCKS edge blocks the target`() { + val a = task("p-1") + val blocker = task("p-2", status = TaskStatus.IN_PROGRESS, links = listOf(blocks("p-1"))) + val board = listOf(a, blocker) + + assertEquals(listOf("p-2"), TaskGraph.unmetBlockers(a, board).map { it.taskId.value }) + } + + @Test + fun `ready lists only unblocked TODO tasks`() { + val a = task("p-1", links = listOf(dependsOn("p-2"))) // blocked + val b = task("p-2", status = TaskStatus.IN_PROGRESS) // not TODO + val c = task("p-3") // unblocked TODO → ready + + assertEquals(listOf("p-3"), TaskGraph.ready(listOf(a, b, c)).map { it.taskId.value }) + } + + @Test + fun `blocking lists what completing a task unblocks`() { + val dependent = task("p-1", links = listOf(dependsOn("p-2"))) + val target = task("p-2") + val board = listOf(dependent, target) + + assertEquals(listOf("p-1"), TaskGraph.blocking(target, board).map { it.taskId.value }) + } +} diff --git a/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskServiceTest.kt b/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskServiceTest.kt index b3a8bb62..f9c19936 100644 --- a/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskServiceTest.kt +++ b/core/tasks/src/test/kotlin/com/correx/core/tasks/TaskServiceTest.kt @@ -105,6 +105,22 @@ class TaskServiceTest { assertTrue(lines.last().contains("completed")) } + @Test + fun `ready surfaces unblocked work and blockers explain the wait`() = runBlocking { + val a = service.createTask(project, "A", "g").taskId // auth-1 + val b = service.createTask(project, "B", "g").taskId // auth-2 depends on auth-1 + service.link(b, a.value, TaskLinkType.DEPENDS_ON, TaskTargetKind.TASK) + + assertEquals(setOf("auth-1"), service.ready(project).map { it.taskId.value }.toSet()) + assertEquals(listOf("auth-1"), service.blockers(b).map { it.taskId.value }) + assertEquals(listOf("auth-2"), service.blocking(a).map { it.taskId.value }) + + // Finishing auth-1 unblocks auth-2. + service.claim(a, "x"); service.submitForReview(a); service.complete(a) + assertEquals(setOf("auth-2"), service.ready(project).map { it.taskId.value }.toSet()) + assertTrue(service.blockers(b).isEmpty()) + } + @Test fun `history survives soft-delete but is empty for an unknown task`() = runBlocking { val id = service.createTask(project, "throwaway", "oops").taskId