Merge pull request 'mesh-delivery: the owner of deliveries and delivery groups; the forge's note, view and status tools (hq ADR 0239)' (#99) from feat/mesh-delivery into main

This commit was merged in pull request #99.
This commit is contained in:
2026-10-06 23:00:36 +00:00
22 changed files with 4698 additions and 2 deletions
+17
View File
@@ -50,6 +50,8 @@ export interface GiteaPull {
base?: string;
updated_at?: string;
html_url: string;
/** The description: where a pull request says it goes after another (`after: <repository>`, novox/hq ADR 0239). */
body?: string;
}
/** A commit status, as the forge keeps it: what a pull request shows beside its head commit. */
@@ -338,6 +340,20 @@ export class GiteaClient {
await this.request(`/repos/${owner}/${repo}/statuses/${sha}`, { method: "POST", body: JSON.stringify(status) });
}
/** A commit's statuses by context, the newest of each: what a merge's head was checked as (novox/hq ADR 0239). */
async commitStatuses(owner: string, repo: string, sha: string): Promise<Record<string, string>> {
const raw = await this.request<any>(`/repos/${owner}/${repo}/commits/${sha}/status`);
const out: Record<string, string> = {};
for (const s of raw?.statuses ?? []) if (s?.context && !out[s.context]) out[s.context] = String(s.status ?? s.state ?? "");
return out;
}
/** Replace a comment's body: the delivery's view, kept current in place (novox/hq ADR 0239). */
async editComment(owner: string, repo: string, id: number, body: string): Promise<{ id: number; html_url: string }> {
const c = await this.request<any>(`/repos/${owner}/${repo}/issues/comments/${id}`, { method: "PATCH", body: JSON.stringify({ body }) });
return { id: Number(c?.id ?? id), html_url: String(c?.html_url ?? "") };
}
// ---- Branch protection ----
/** A repository's branch protection rules. */
@@ -483,6 +499,7 @@ export class GiteaClient {
base: p.base?.ref,
updated_at: p.updated_at ?? undefined,
html_url: p.html_url,
body: typeof p.body === "string" ? p.body : undefined,
};
}
}
+97
View File
@@ -0,0 +1,97 @@
// What the forge's holder does for a delivery (novox/hq ADR 0239): the commit's note under
// refs/notes/mesh-plan, the delivery's view on its pull request, and its statuses — asked by mesh-delivery,
// the delivery's owner, through this module's tools. The forge is this module's; mesh-delivery never writes
// to it itself.
//
// **The note is written in the forge's own repository, as the forge's own user.** The forge's API reads a
// note and writes none, and a push needs a clone and a credential nobody else should hold; the forge's
// container holds the bare repository and git. So the line is appended there, by `git notes append`, as the
// account the forge runs as — the mesh's own forge, never a person's or an agent's hand. Appending is
// idempotent here: a line the note already has is not added again, so asking twice adds nothing.
//
// Pure functions and an injectable runner, so this is tested without a forge.
import { execFile } from "node:child_process";
import { promisify } from "node:util";
/** How a command is run: docker on the forge's machine, or a test's. */
export type Runner = (file: string, args: string[]) => Promise<{ stdout: string; code: number }>;
const execFileP = promisify(execFile);
export const run: Runner = async (file, args) => {
try {
const { stdout } = await execFileP(file, args, { maxBuffer: 4 << 20 });
return { stdout, code: 0 };
} catch (err) {
const e = err as { code?: number; stdout?: string };
return { stdout: e.stdout ?? "", code: typeof e.code === "number" ? e.code : 1 };
}
};
const name = /^[A-Za-z0-9_.-]+$/;
const sha = /^[0-9a-fA-F]{7,64}$/;
const ref = /^[a-z0-9][a-z0-9-]*$/;
/** A note line as it may be written: one line, bounded, nothing that is not text. */
export function noteLine(line: string): string {
const one = String(line ?? "").replace(/[\r\n\t]+/g, " ").replace(/\s+/g, " ").trim();
if (!one) throw new Error("an empty line is no note");
return one.length > 2000 ? one.slice(0, 1999) + "…" : one;
}
/** Where the forge keeps a repository, inside its container. */
export function gitDir(owner: string, repo: string): string {
if (!name.test(owner) || !name.test(repo)) throw new Error(`${owner}/${repo} is not a repository's name`);
return `/data/git/repositories/${owner.toLowerCase()}/${repo.toLowerCase()}.git`;
}
/** The commands the note is read and appended with, in the forge's container, as its user. */
export function noteArgs(container: string, owner: string, repo: string, commit: string, notesRef: string, line?: string): string[] {
if (!name.test(container)) throw new Error(`${container} is not a container's name`);
if (!sha.test(commit)) throw new Error(`${commit} is not a commit`);
if (!ref.test(notesRef)) throw new Error(`${notesRef} is not a notes ref`);
const base = ["exec", "-u", "git", container, "git", "--git-dir", gitDir(owner, repo),
"-c", "user.name=mesh", "-c", "user.email=mesh@mesh.invalid", "notes", `--ref=${notesRef}`];
return line === undefined ? [...base, "show", commit] : [...base, "append", "-m", noteLine(line), commit];
}
/** Append a line to a commit's note, unless the note already holds it. Answers whether it was added. */
export async function appendNote(runner: Runner, container: string, owner: string, repo: string, commit: string,
notesRef: string, line: string): Promise<{ added: boolean; lines: number }> {
const wanted = noteLine(line);
const shown = await runner("docker", noteArgs(container, owner, repo, commit, notesRef));
// No note yet is git's exit 1 with nothing on stdout; anything else unreadable is said.
const lines = shown.code === 0 ? shown.stdout.split("\n").filter((l) => l.trim() !== "") : [];
if (lines.includes(wanted)) return { added: false, lines: lines.length };
const appended = await runner("docker", noteArgs(container, owner, repo, commit, notesRef, wanted));
if (appended.code !== 0) throw new Error(`the note on ${commit.slice(0, 8)} could not be appended (git exited ${appended.code})`);
return { added: true, lines: lines.length + 1 };
}
/** The marker that makes one comment of a pull request the delivery's view. */
export const VIEW_MARKER = "<!-- mesh-delivery:view -->";
/** The view's body, marked: mesh-delivery writes it, this keeps exactly one of them per pull request. */
export function viewBody(body: string): string {
const b = String(body ?? "");
return b.includes(VIEW_MARKER) ? b : `${VIEW_MARKER}\n${b}`;
}
/** Which comment is the view: the first that carries the marker, or none yet. */
export function viewComment<T extends { id: number; body: string }>(comments: T[]): T | undefined {
return comments.find((c) => c.body.includes(VIEW_MARKER));
}
const states = new Set(["pending", "success", "error", "failure", "warning"]);
/** A status mesh-delivery asks for, checked: one of the forge's states, a context of the mesh's own, a
* bounded description. */
export function deliveryStatus(context: string, state: string, description: string, target?: string) {
if (!/^mesh\/[a-z-]+$/.test(context)) throw new Error(`${context} is not a status of the mesh's own`);
if (!states.has(state)) throw new Error(`${state} is not a status the forge keeps`);
let d = String(description ?? "").replace(/\s+/g, " ").trim();
if (d.length > 140) d = d.slice(0, 139) + "…";
return { state: state as "pending" | "success" | "error" | "failure" | "warning", context, description: d,
...(target && /^https?:\/\//.test(target) ? { target_url: target } : {}) };
}
+41 -1
View File
@@ -7,6 +7,8 @@
// module.gitea.pull.updated — an open pull request's head moved, opened or pushed to: the mesh checks it
// before it merges (novox/hq to-be 45 §9)
//
// module.gitea.pull.closed — a pull request closed unmerged: the delivery it was is stopped (novox/hq ADR 0239)
//
// Consumes mesh-controller.checked — a pull request's merge check, judged — and sets it as the head
// commit's statuses: `mesh/merge-gate`, the modules of the mesh's graph the change touches, and
// `mesh/repo-check`, the repository's own merge-check.sh (novox/hq ADR 0237 as amended); with a comment
@@ -109,6 +111,18 @@ async function pollMerged(client: GiteaClient): Promise<void> {
for (const repo of repos) {
const pulls = await client.listPullRequests(repo.owner, repo.name, { state: "closed", sort: "recentupdate", limit: "20" });
for (const pull of pulls) {
// Closed unmerged (novox/hq ADR 0239): announced once, since the watching began, so its delivery stops.
const closedKey = `closed:${repo.full_name}#${pull.number}`;
if (!pull.merged && pull.state === "closed" && !announced.has(closedKey)) {
if (primedMerges && !!pull.updated_at && !!since && pull.updated_at > since) {
await emit("pull.closed", { owner: repo.owner, repo: repo.name, number: pull.number, title: pull.title,
head: pull.head, head_sha: pull.head_sha, base: pull.base, html_url: pull.html_url });
console.log(`[gitea] announced ${repo.full_name}#${pull.number} closed unmerged`);
}
announced.add(closedKey);
changed = true;
continue;
}
if (!pull.merged || !pull.merge_commit_sha || announced.has(pull.merge_commit_sha)) continue;
// Announced only if merged since the watching began; recorded either way, so it is looked
// at once.
@@ -130,11 +144,24 @@ async function pollMerged(client: GiteaClient): Promise<void> {
`so the mesh reads its files by the old rule — ${err instanceof Error ? err.message : String(err)}`);
}
}
// The head it merged, and what its head was checked as (novox/hq ADR 0239): the delivery it was, made
// from the forge's word when its owner never heard the head.
let headChecks: Record<string, string> | null = null;
if (pull.head_sha) {
try {
headChecks = await client.commitStatuses(repo.owner, repo.name, pull.head_sha);
} catch (err) {
console.error(`[gitea] ${repo.full_name}#${pull.number}: its head's statuses could not be read — ${err instanceof Error ? err.message : err}`);
}
}
await emit("pull.merged", {
owner: repo.owner,
repo: repo.name,
number: pull.number,
title: pull.title,
body: pull.body,
head_sha: pull.head_sha,
...(headChecks ? { head_checks: headChecks } : {}),
head: pull.head,
base: pull.base,
merge_commit_sha: pull.merge_commit_sha,
@@ -215,6 +242,7 @@ async function pollPulls(client: GiteaClient): Promise<void> {
repo: repo.name,
number: pull.number,
title: pull.title,
body: pull.body,
base: pull.base,
head: pull.head,
head_sha: pull.head_sha,
@@ -249,11 +277,23 @@ async function pollPulls(client: GiteaClient): Promise<void> {
// change is reviewed. An error — the check could not run — is the forge's `error`, never a success.
async function setVerdict(client: GiteaClient, event: { body: unknown }): Promise<void> {
const c = (event.body ?? {}) as Checked;
if (c.group) {
// A delivery group's composed check (novox/hq ADR 0239): its verdict is the group owner's to say, on each
// member's head, as mesh/delivery-group — never this head's merge gate.
return;
}
if (!c.owner || !c.repo || !c.commit || !c.verdict) {
console.error("[gitea] a merge check's verdict named no repository, commit or verdict; ignored");
return;
}
for (const status of statusesFor(c)) await client.setCommitStatus(c.owner, c.repo, c.commit, status);
// Each status links to the pull request, where the delivery's view is kept (novox/hq ADR 0239).
let target: string | undefined;
if (c.number) {
target = await client.getPullRequest(c.owner, c.repo, c.number).then((p) => p.html_url || undefined).catch(() => undefined);
}
for (const status of statusesFor(c)) {
await client.setCommitStatus(c.owner, c.repo, c.commit, target ? { ...status, target_url: target } : status);
}
const comment = commentFor(c);
if (comment && c.number) await client.addComment(c.owner, c.repo, c.number, comment);
console.log(`[gitea] ${c.owner}/${c.repo}#${c.number ?? "?"} at ${c.commit.slice(0, 8)}: merge check ${c.verdict}`);
+2 -1
View File
@@ -41,7 +41,8 @@
"repo.created",
"issue.opened",
"pull.merged",
"pull.updated"
"pull.updated",
"pull.closed"
],
"consumes": [
"mesh-controller.checked"
+2
View File
@@ -44,6 +44,8 @@ export interface Checked {
"repo-check"?: Layer;
/** The change plan of the commit checked (novox/hq ADR 0238): what a merge of it would build and send. */
plan?: ChangePlan;
/** Set on a delivery group's composed check (novox/hq ADR 0239): not this head's merge gate. */
group?: string;
}
/** A change plan, as the controller says it (mesh-controller internal/link, ChangePlan). */
+56
View File
@@ -0,0 +1,56 @@
import assert from "node:assert/strict";
import { test } from "node:test";
// What the forge's holder does for a delivery (novox/hq ADR 0239): a note appended once, in the forge's own
// repository as its own user; one view per pull request; only the mesh's statuses.
test("a note line is appended once, as the forge's user, in its own repository", async () => {
const { appendNote, noteArgs } = await import("../delivery.ts");
const notes: Record<string, string[]> = {};
const calls: string[][] = [];
const runner = async (file: string, args: string[]) => {
calls.push([file, ...args]);
const commit = args[args.length - 1];
if (args.includes("show")) {
const lines = notes[commit];
return lines ? { stdout: lines.join("\n") + "\n", code: 0 } : { stdout: "", code: 1 };
}
const line = args[args.indexOf("-m") + 1];
(notes[commit] ??= []).push(line);
return { stdout: "", code: 0 };
};
const sha = "0123456789abcdef0123456789abcdef01234567";
assert.deepEqual(await appendNote(runner, "gitea", "Novox", "Mesh-Catalog", sha, "mesh-plan", "a -> b\n(merged)"),
{ added: true, lines: 1 });
assert.deepEqual(await appendNote(runner, "gitea", "Novox", "Mesh-Catalog", sha, "mesh-plan", "a -> b (merged)"),
{ added: false, lines: 1 }, "the same line twice adds nothing");
assert.deepEqual(await appendNote(runner, "gitea", "Novox", "Mesh-Catalog", sha, "mesh-plan", "b -> c"),
{ added: true, lines: 2 });
const append = calls.find((c) => c.includes("append"))!;
assert.deepEqual(append.slice(0, 7), ["docker", "exec", "-u", "git", "gitea", "git", "--git-dir"]);
assert.equal(append[7], "/data/git/repositories/novox/mesh-catalog.git");
assert.ok(append.includes("--ref=mesh-plan"));
assert.throws(() => noteArgs("gitea", "novox", "x; rm -rf /", sha, "mesh-plan"), /not a repository/);
assert.throws(() => noteArgs("gitea", "novox", "x", "HEAD", "mesh-plan"), /not a commit/);
assert.throws(() => noteArgs("gitea", "novox", "x", sha, "../commits"), /not a notes ref/);
});
test("a note that cannot be written is said, never read as written", async () => {
const { appendNote } = await import("../delivery.ts");
const runner = async (_file: string, args: string[]) => ({ stdout: "", code: args.includes("append") ? 128 : 1 });
await assert.rejects(appendNote(runner, "gitea", "novox", "x", "abcdef1234567", "mesh-plan", "l"), /could not be appended/);
});
test("one view per pull request, found by its marker; only the mesh's statuses", async () => {
const { viewBody, viewComment, deliveryStatus, VIEW_MARKER } = await import("../delivery.ts");
assert.ok(viewBody("**Delivery**").startsWith(VIEW_MARKER));
assert.equal(viewBody(`${VIEW_MARKER}\nx`), `${VIEW_MARKER}\nx`, "a marked body is kept as it is");
const comments = [{ id: 1, body: "a review" }, { id: 2, body: `${VIEW_MARKER}\nold` }, { id: 3, body: `${VIEW_MARKER}\nlater` }];
assert.equal(viewComment(comments)?.id, 2);
assert.equal(viewComment([{ id: 1, body: "x" }]), undefined);
const s = deliveryStatus("mesh/delivery", "pending", "y".repeat(300), "https://forge.invalid/novox/x/pulls/1");
assert.ok(s.description.length <= 140 && s.target_url);
assert.throws(() => deliveryStatus("ci/other", "success", "x"), /mesh's own/);
assert.throws(() => deliveryStatus("mesh/delivery", "green", "x"), /not a status/);
assert.equal(deliveryStatus("mesh/delivery", "success", "x", "javascript:alert(1)").target_url, undefined);
});
+2
View File
@@ -410,6 +410,8 @@ test("the tools register once there is a way to a token, and the first call mint
"gitea_get_file",
"gitea_list_branches",
"gitea_delete_branch",
"gitea_branch_protection_get", "gitea_branch_protection_set",
"gitea_note_append", "gitea_delivery_view", "gitea_commit_status",
"gitea_list_labels", "gitea_create_label",
"gitea_api",
],
+58
View File
@@ -11,6 +11,10 @@
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
import { emit } from "@novox/mesh-sdk/events";
import { GiteaClient, type BranchProtection } from "../client.js";
import { appendNote, deliveryStatus, run, viewBody, viewComment } from "../delivery.js";
/** The forge's container, where its repositories and git are: the note is written there (novox/hq ADR 0239). */
const forgeContainer = process.env.MESH_GITEA_CONTAINER || "gitea";
/** A protection rule as a person reads it: what it guards, not every field the forge keeps. */
export function summarised(p: BranchProtection) {
@@ -412,6 +416,60 @@ export function getGiteaTools(gitea: GiteaClient): ToolDefinition[] {
},
},
// ---- A delivery's note, view and statuses (novox/hq ADR 0239), asked by mesh-delivery ----
{
name: "gitea_note_append",
description: "Append one line to a commit's git note under refs/notes/<ref> (mesh-plan for a delivery), written in the forge's own repository as the forge's own account; a line the note already holds is not added again. `git log --notes=mesh-plan` shows it.",
input: {
owner: { type: "string", description: "the repository owner" },
repo: { type: "string", description: "the repository name" },
sha: { type: "string", description: "the commit" },
ref: { type: "string", description: "the notes ref's name (default mesh-plan)" },
line: { type: "string", description: "the line, one line" },
},
run: async (args) => {
const said = await appendNote(run, forgeContainer, String(args.owner), String(args.repo), String(args.sha),
String(args.ref ?? "mesh-plan") || "mesh-plan", String(args.line ?? ""));
return { commit: String(args.sha), ...said };
},
},
{
name: "gitea_delivery_view",
description: "Keep a delivery's view on its pull request: one comment, marked as the delivery's, created the first time and edited in place after — the page every status of the delivery links to.",
input: {
owner: { type: "string", description: "the repository owner" },
repo: { type: "string", description: "the repository name" },
number: { type: "number", description: "the pull request's number" },
body: { type: "string", description: "the view, markdown" },
},
run: async (args) => {
const owner = String(args.owner), repo = String(args.repo), number = Number(args.number);
const body = viewBody(String(args.body ?? ""));
const existing = viewComment(await gitea.listComments(owner, repo, number));
if (existing) return { edited: await gitea.editComment(owner, repo, existing.id, body) };
return { created: await gitea.addComment(owner, repo, number, body) };
},
},
{
name: "gitea_commit_status",
description: "Set one of the mesh's statuses on a commit (mesh/delivery, mesh/delivery-group): pending, success, error, failure or warning, a short description, and the page it links to.",
input: {
owner: { type: "string", description: "the repository owner" },
repo: { type: "string", description: "the repository name" },
sha: { type: "string", description: "the commit" },
context: { type: "string", description: "the status's name, mesh/…" },
state: { type: "string", description: "pending | success | error | failure | warning" },
description: { type: "string", description: "one line" },
target_url: { type: "string", description: "the page it links to (the pull request)" },
},
run: async (args) => {
const status = deliveryStatus(String(args.context), String(args.state), String(args.description ?? ""),
args.target_url ? String(args.target_url) : undefined);
await gitea.setCommitStatus(String(args.owner), String(args.repo), String(args.sha), status);
return { set: status };
},
},
// ---- Labels ----
{
name: "gitea_list_labels",
@@ -0,0 +1,295 @@
package main
import (
"encoding/json"
"fmt"
"sort"
"strings"
"time"
)
// What a transition owes outside the state (novox/hq ADR 0239 decision 5), each by its owner: the event on
// the bus (said by this module), a line on the commit's note and the pull request's view and statuses (by
// the forge's holder, asked through its tools). Owed, never lost: kept with the delivery before anything is
// tried, and tried again until done.
// transitionEvent is what `mesh-delivery.transition` says: no secret, no address.
func transitionEvent(d *Delivery, t Transition) map[string]any {
e := map[string]any{"id": d.ID, "repository": d.Repository, "commit": d.Commit, "from": t.From, "to": t.To,
"event": t.Event, "why": t.Why, "at": t.At}
if t.By != "" {
e["by"] = t.By
}
if d.Number > 0 {
e["number"] = d.Number
}
if d.Group != "" {
e["group"] = d.Group
}
if d.MergedAs != "" {
e["merged_as"] = d.MergedAs
}
if d.Walk != nil {
e["walk"] = d.Walk.ID
}
if d.Plan != nil {
e["plan"] = d.Plan.Summary
}
return e
}
// noteLine is one transition as the commit's note keeps it.
func noteLine(d *Delivery, t Transition) string {
by := ""
if t.By != "" {
by = " by " + t.By
}
line := fmt.Sprintf("%s delivery %s: %s -> %s (%s%s): %s", t.At.UTC().Format(time.RFC3339), d.ID, orNothing(t.From),
t.To, t.Event, by, t.Why)
if t.To == Proposed && d.Plan != nil {
line += "; predicted: " + d.Plan.Summary
}
if d.Group != "" && (t.To == Proposed || t.To == Published) {
line += "; group " + d.Group
}
return oneLine(line)
}
// executedLine is what a delivery did, written on the commit it landed as when it ends.
func executedLine(d *Delivery) string {
var parts []string
for _, s := range d.Steps {
parts = append(parts, fmt.Sprintf("%s on %s %s", s.Module, s.Machine, s.State))
}
walk := ""
if d.Walk != nil {
walk = " (walk " + d.Walk.ID + ")"
}
return oneLine(fmt.Sprintf("%s delivery %s executed%s: %s", time.Now().UTC().Format(time.RFC3339), d.ID, walk,
strings.Join(parts, "; ")))
}
func oneLine(s string) string { return strings.Join(strings.Fields(s), " ") }
// statusOf is the commit status `mesh/delivery`: where the delivery stands, in the forge's words.
func statusOf(d *Delivery) (string, string) {
why := ""
if len(d.Transitions) > 0 {
why = d.Transitions[len(d.Transitions)-1].Why
}
switch d.State {
case Proposed, Checked:
return "pending", "checking"
case Ready:
return "success", "ready: it delivers once merged"
case Rejected:
return "failure", clip("rejected: "+why, 140)
case Published:
return "pending", "published: its walk waits for its turn"
case Held:
return "pending", clip("held for a person: "+d.HeldWhy, 140)
case Delivering:
passed := 0
for _, s := range d.Steps {
if s.State == StepPassed {
passed++
}
}
return "pending", fmt.Sprintf("delivering: %d machine step(s) passed", passed)
case Delivered:
return "success", "delivered"
case Failed:
return "failure", clip("failed: "+why, 140)
case Superseded:
return "warning", clip("superseded: "+why, 140)
case Stopped:
return "error", clip("stopped: "+why, 140)
}
return "pending", string(d.State)
}
// viewMarker is how the forge's holder finds the one view comment of a pull request.
const viewMarker = "<!-- mesh-delivery:view -->"
// ViewBody is a delivery's view on its pull request, written from the delivery whole.
func ViewBody(d *Delivery, g *Group) string {
var b strings.Builder
fmt.Fprintf(&b, "%s\n**Delivery** `%s` — **%s** since %s\n\n", viewMarker, d.ID, d.State, d.Since.UTC().Format(time.RFC3339))
if d.HeldWhy != "" && d.State == Held {
fmt.Fprintf(&b, "Held for a person: %s — `mesh-delivery.release` with why lets it go on.\n\n", d.HeldWhy)
}
if d.Plan != nil {
fmt.Fprintf(&b, "**Delivery plan** — %s\n", d.Plan.Summary)
for i, tier := range d.Plan.Tiers {
fmt.Fprintf(&b, "- tier %d: %s\n", i, strings.Join(tier, ", "))
}
for _, m := range d.Plan.Machines {
var parts []string
if len(m.Receives) > 0 {
parts = append(parts, "receives "+strings.Join(m.Receives, ", "))
}
if len(m.Waits) > 0 {
parts = append(parts, "waits for a person: "+strings.Join(m.Waits, ", "))
}
fmt.Fprintf(&b, "- %s: %s\n", m.Machine, strings.Join(parts, "; "))
}
for _, s := range d.Plan.Steps {
fmt.Fprintf(&b, "- %s\n", s)
}
b.WriteString("\n")
}
if g != nil {
fmt.Fprintf(&b, "**Group** `%s`, in order: %s\n", g.ID, strings.Join(g.Order, " → "))
for _, p := range g.Pairs {
fmt.Fprintf(&b, "- %s before %s: %s\n", p.Before, p.After, p.Why)
}
if g.Check != nil && g.Check.Verdict != "" {
fmt.Fprintf(&b, "- composed together: %s — %s\n", g.Check.Verdict, g.Check.Summary)
}
b.WriteString("\n")
}
if len(d.Steps) > 0 {
b.WriteString("**Machines**\n")
for _, s := range d.Steps {
fmt.Fprintf(&b, "- %s on %s: %s", s.Module, s.Machine, s.State)
if s.Why != "" {
fmt.Fprintf(&b, " — %s", s.Why)
}
b.WriteString("\n")
}
b.WriteString("\n")
}
b.WriteString("**Transitions**\n")
from := 0
if len(d.Transitions) > 12 {
from = len(d.Transitions) - 12
}
for _, t := range d.Transitions[from:] {
fmt.Fprintf(&b, "- %s %s → %s (%s): %s\n", t.At.UTC().Format("2006-01-02 15:04"), orNothing(t.From), t.To, t.Event, t.Why)
}
b.WriteString("\nThe commit's note under `refs/notes/mesh-plan` keeps every transition: `git log --notes=mesh-plan`.\n")
return b.String()
}
// Flush does what is owed, kept deliveries first: a delivery not yet kept says nothing.
func (h *Holder) Flush() {
type job struct {
id string
group bool
e Effect
body string
repo [2]string
num int
url string
}
h.mu.Lock()
for id := range h.dirty {
if d := h.deliveries[id]; d != nil {
h.keep(d)
}
}
for id := range h.dirtyGroup {
if g := h.groups[id]; g != nil {
h.keepGroup(g)
}
}
var jobs []job
ids := make([]string, 0, len(h.deliveries))
for id := range h.deliveries {
ids = append(ids, id)
}
sort.Strings(ids)
for _, id := range ids {
d := h.deliveries[id]
if h.dirty[id] {
continue
}
for _, e := range d.Owed {
j := job{id: id, e: e, repo: [2]string{d.Owner(), d.Repo()}, num: d.Number, url: d.HTMLURL}
if e.Kind == EffectView {
j.body = ViewBody(d, h.groups[d.Group])
}
jobs = append(jobs, j)
}
}
for id, g := range h.groups {
if h.dirtyGroup[id] {
continue
}
for _, e := range g.Owed {
jobs = append(jobs, job{id: id, group: true, e: e})
}
}
h.mu.Unlock()
for _, j := range jobs {
err := h.do(j.e, j.repo, j.num, j.body, j.url, j.group)
h.mu.Lock()
owed := func(list []Effect) []Effect {
for i, o := range list {
if o.Kind == j.e.Kind && o.Commit == j.e.Commit && o.Context == j.e.Context && o.Line == j.e.Line &&
o.Since.Equal(j.e.Since) {
if err == nil {
return append(list[:i:i], list[i+1:]...)
}
list[i].Tries++
list[i].Last = err.Error()
return list
}
}
return list
}
if j.group {
if g := h.groups[j.id]; g != nil {
g.Owed = owed(g.Owed)
h.keepGroup(g)
}
} else if d := h.deliveries[j.id]; d != nil {
d.Owed = owed(d.Owed)
h.keep(d)
}
h.mu.Unlock()
if err != nil && j.e.Tries%10 == 0 {
h.Logf("[mesh-delivery] %s: %s not done yet (%v); tried again", j.id, j.e.Kind, err)
}
}
}
// do is one owed effect.
func (h *Holder) do(e Effect, repo [2]string, number int, body, url string, group bool) error {
switch e.Kind {
case EffectEmit:
var payload any
if err := json.Unmarshal([]byte(e.Line), &payload); err != nil {
return nil // nothing sayable: dropped
}
event := "transition"
if group {
event = "group"
}
return h.Emit(event, payload)
case EffectNote:
if h.Forge == nil {
return fmt.Errorf("no forge")
}
return h.Forge.Note(repo[0], repo[1], e.Commit, e.Line)
case EffectView:
if h.Forge == nil || number == 0 {
return nil
}
return h.Forge.View(repo[0], repo[1], number, body)
case EffectStatus:
if h.Forge == nil {
return fmt.Errorf("no forge")
}
if group {
// A group's status names its member's repository and page in Last.
where, page, _ := strings.Cut(e.Last, " ")
ownerRepo, _, _ := strings.Cut(where, "#")
owner, name, _ := strings.Cut(ownerRepo, "/")
return h.Forge.Status(owner, name, e.Commit, e.Context, e.State, e.Line, page)
}
return h.Forge.Status(repo[0], repo[1], e.Commit, e.Context, e.State, e.Line, url)
}
return nil
}
@@ -0,0 +1,374 @@
package main
import (
"encoding/json"
"errors"
"fmt"
"sort"
"strings"
"sync"
"testing"
"time"
)
// What a test stands the holder on: a store in memory that can be made to fail, a controller that records
// what it was asked and answers what it is told, a forge that records every note, view and status.
type memStore struct {
mu sync.Mutex
deliveries map[string][]byte
groups map[string][]byte
failing bool
}
func newMemStore() *memStore {
return &memStore{deliveries: map[string][]byte{}, groups: map[string][]byte{}}
}
func (s *memStore) PutDelivery(d *Delivery) error {
s.mu.Lock()
defer s.mu.Unlock()
if s.failing {
return errors.New("the bus is away")
}
s.deliveries[d.ID] = mustJSON(d)
return nil
}
func (s *memStore) DeleteDelivery(id string) error {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.deliveries, id)
return nil
}
func (s *memStore) PutGroup(g *Group) error {
s.mu.Lock()
defer s.mu.Unlock()
if s.failing {
return errors.New("the bus is away")
}
s.groups[g.ID] = mustJSON(g)
return nil
}
func (s *memStore) DeleteGroup(id string) error {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.groups, id)
return nil
}
func (s *memStore) Deliveries() ([]*Delivery, error) {
s.mu.Lock()
defer s.mu.Unlock()
var out []*Delivery
for _, raw := range s.deliveries {
var d Delivery
if err := json.Unmarshal(raw, &d); err != nil {
return nil, err
}
out = append(out, &d)
}
return out, nil
}
func (s *memStore) Groups() ([]*Group, error) {
s.mu.Lock()
defer s.mu.Unlock()
var out []*Group
for _, raw := range s.groups {
var g Group
if err := json.Unmarshal(raw, &g); err != nil {
return nil, err
}
out = append(out, &g)
}
return out, nil
}
type fakeController struct {
mu sync.Mutex
walks map[string]Walk
held bool
order func([]Member) Order
checks []string
asked int
delivers []string
stops []string
plans int
down bool
}
func newFakeController() *fakeController {
return &fakeController{walks: map[string]Walk{}, held: true}
}
func (c *fakeController) Plan(repository, base, head string, paths, dirs, removed []string) (*DeliveryPlan, error) {
c.mu.Lock()
defer c.mu.Unlock()
c.plans++
if c.down {
return nil, errors.New("down")
}
return &DeliveryPlan{Repository: repository, Base: base, Head: head, Moved: []string{"app"}, Summary: "builds app"}, nil
}
func (c *fakeController) Order(ms []Member) (Order, error) {
c.mu.Lock()
defer c.mu.Unlock()
if c.down {
return Order{}, errors.New("down")
}
if c.order != nil {
return c.order(ms), nil
}
var ids []string
for _, m := range ms {
ids = append(ids, m.ID)
}
sort.Strings(ids)
return Order{Order: ids}, nil
}
func (c *fakeController) Check(group string, ms []Member) (string, error) {
c.mu.Lock()
defer c.mu.Unlock()
if c.down {
return "", errors.New("down")
}
c.asked++
id := fmt.Sprintf("check-%d", c.asked)
var heads []string
for _, m := range ms {
heads = append(heads, m.ID)
}
c.checks = append(c.checks, group+":"+strings.Join(heads, ","))
return id, nil
}
func (c *fakeController) Deliver(walk, why string) error {
c.mu.Lock()
defer c.mu.Unlock()
if c.down {
return errors.New("down")
}
w, known := c.walks[walk]
if !known || w.Delivery == nil || w.Delivery.Go != nil {
return fmt.Errorf("%s waits for nobody", walk)
}
now := time.Now()
w.Delivery.Go, w.Delivery.By, w.Delivery.Why = &now, byOwner, why
w.Revision++
c.walks[walk] = w
c.delivers = append(c.delivers, walk)
return nil
}
func (c *fakeController) Stop(walk, why, by string) error {
c.mu.Lock()
defer c.mu.Unlock()
if c.down {
return errors.New("down")
}
c.stops = append(c.stops, walk)
if w, known := c.walks[walk]; known {
w.State = walkFailed
if w.Delivery == nil {
w.Delivery = &WalkDelivery{}
}
w.Delivery.Stopped, w.Delivery.StoppedWhy = by, why
w.Revision++
c.walks[walk] = w
}
return nil
}
func (c *fakeController) Walks(walk string) (bool, []Walk, error) {
c.mu.Lock()
defer c.mu.Unlock()
if c.down {
return false, nil, errors.New("down")
}
var out []Walk
for id, w := range c.walks {
if walk == "" || walk == id {
out = append(out, w)
}
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
return c.held, out, nil
}
// put is the controller keeping a walk; it answers the walk as `plan-moved` would say it.
func (c *fakeController) put(w Walk) Walk {
c.mu.Lock()
defer c.mu.Unlock()
if old, known := c.walks[w.ID]; known {
w.Revision = old.Revision + 1
} else if w.Revision == 0 {
w.Revision = 1
}
c.walks[w.ID] = w
return w
}
func (c *fakeController) walk(id string) Walk {
c.mu.Lock()
defer c.mu.Unlock()
return c.walks[id]
}
type fakeForge struct {
mu sync.Mutex
notes map[string][]string
views map[string]string
statuses map[string]string
down bool
}
func newFakeForge() *fakeForge {
return &fakeForge{notes: map[string][]string{}, views: map[string]string{}, statuses: map[string]string{}}
}
func (f *fakeForge) Note(owner, repo, commit, line string) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.down {
return errors.New("the forge is away")
}
key := owner + "/" + repo + "@" + commit
for _, l := range f.notes[key] {
if l == line {
return nil // appended once
}
}
f.notes[key] = append(f.notes[key], line)
return nil
}
func (f *fakeForge) View(owner, repo string, number int, body string) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.down {
return errors.New("the forge is away")
}
f.views[fmt.Sprintf("%s/%s#%d", owner, repo, number)] = body
return nil
}
func (f *fakeForge) Status(owner, repo, commit, context, state, description, target string) error {
f.mu.Lock()
defer f.mu.Unlock()
if f.down {
return errors.New("the forge is away")
}
f.statuses[owner+"/"+repo+"@"+commit+" "+context] = state + " " + description + " → " + target
return nil
}
func (f *fakeForge) notesOn(commit string) []string {
f.mu.Lock()
defer f.mu.Unlock()
for key, lines := range f.notes {
if strings.HasSuffix(key, "@"+commit) {
return append([]string(nil), lines...)
}
}
return nil
}
// world is a holder over fakes, with a clock a test moves.
type world struct {
t *testing.T
h *Holder
store *memStore
ctl *fakeController
forge *fakeForge
now time.Time
said []map[string]any
mu sync.Mutex
}
func newWorld(t *testing.T) *world {
return newWorldOver(t, newMemStore(), newFakeController(), newFakeForge(), time.Date(2026, 10, 6, 12, 0, 0, 0, time.UTC))
}
func newWorldOver(t *testing.T, store *memStore, ctl *fakeController, forge *fakeForge, now time.Time) *world {
w := &world{t: t, store: store, ctl: ctl, forge: forge, now: now}
w.h = &Holder{Store: store, Controller: ctl, Forge: forge,
Now: func() time.Time { w.mu.Lock(); defer w.mu.Unlock(); return w.now },
Logf: func(format string, a ...any) { t.Logf(format, a...) },
Emit: func(event string, body any) error {
w.mu.Lock()
defer w.mu.Unlock()
var m map[string]any
_ = json.Unmarshal(mustJSON(body), &m)
m["event-kind"] = event
w.said = append(w.said, m)
return nil
}}
w.h.init()
if err := w.h.Load(); err != nil {
t.Fatal(err)
}
return w
}
func (w *world) later(d time.Duration) {
w.mu.Lock()
w.now = w.now.Add(d)
w.mu.Unlock()
}
// settleAll ticks and flushes until nothing changes, a few times.
func (w *world) settleAll() {
for range 4 {
w.h.Tick()
w.h.Flush()
}
}
func (w *world) delivery(id string) *Delivery {
w.t.Helper()
w.h.mu.Lock()
defer w.h.mu.Unlock()
d := w.h.deliveries[id]
if d == nil {
w.t.Fatalf("no delivery %s; there are %v", id, w.ids())
}
c := *d
return &c
}
func (w *world) ids() []string {
var out []string
for id := range w.h.deliveries {
out = append(out, id)
}
sort.Strings(out)
return out
}
func (w *world) state(id string) State { return w.delivery(id).State }
// pr is a pull request's head as the forge announces it.
func pr(repo string, number int, sha, branch string, paths ...string) PullEvent {
owner, name, _ := strings.Cut(repo, "/")
return PullEvent{Owner: owner, Repo: name, Number: number, Base: "main", Head: branch, HeadSHA: sha,
HTMLURL: fmt.Sprintf("https://forge.invalid/%s/pulls/%d", repo, number), Paths: paths}
}
// verdict is the controller's `checked` for a head.
func verdict(repo string, number int, sha, gate string) CheckedEvent {
owner, name, _ := strings.Cut(repo, "/")
c := CheckedEvent{Owner: owner, Repo: name, Number: number, Commit: sha, Verdict: gate, Summary: "the gate said " + gate,
ID: "check-of-" + sha, Plan: &DeliveryPlan{Repository: repo, Moved: []string{"app"}, Summary: "builds app → anchor"}}
return c
}
// merged is the forge's `pull.merged`.
func merged(repo string, number int, head, merge string) PullEvent {
p := pr(repo, number, head, "feature")
p.MergeCommit = merge
p.MergedAt = "2026-10-06T12:30:00Z"
return p
}
// aWalk is the controller's walk of a merge commit, waiting for its delivery's word.
func aWalk(id, repo, commit string, waits bool) Walk {
w := Walk{ID: id, Repository: repo, Branch: "main", Commit: commit, Created: time.Date(2026, 10, 6, 12, 31, 0, 0, time.UTC),
State: walkBuilding, Tiers: [][]string{{"app"}}, Modules: map[string]*WalkModule{"app": {}}}
if waits {
w.Delivery = &WalkDelivery{Awaits: "mesh-delivery"}
}
return w
}
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,515 @@
package main
import (
"strings"
"testing"
"time"
)
// The holder, end to end over fakes (novox/hq ADR 0239 *How it is checked*): a pull request's head proposed,
// checked and ready; merged, published, let go, delivering and delivered; a gate that fails, put back and
// failed; a newer head superseding; a restart resuming; the trunk rule; and every transition kept before
// it is said, noted on the commit and shown on the pull request.
const (
head = "aaaaaaaaaaaa1111"
merge = "mmmmmmmmmmmm9999"
)
func TestAPullRequestsHeadIsProposedCheckedAndReady(t *testing.T) {
w := newWorld(t)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
id := IDOf("novox/app", head)
if w.state(id) != Proposed {
t.Fatalf("a head is %s", w.state(id))
}
w.h.Checked(verdict("novox/app", 7, head, "pass"))
d := w.delivery(id)
if d.State != Ready || d.Plan == nil || d.Plan.Summary != "builds app → anchor" {
t.Fatalf("a passing verdict left it %s with plan %+v", d.State, d.Plan)
}
var path []State
for _, tr := range d.Transitions {
path = append(path, tr.To)
}
if strings.Join(statesOf(path), ",") != "proposed,checked,ready" {
t.Fatalf("its transitions are %v", path)
}
w.settleAll()
notes := w.forge.notesOn(head)
if len(notes) != 3 || !strings.Contains(notes[2], "checked -> ready") {
t.Fatalf("the head's note is %q", notes)
}
if !strings.Contains(w.forge.statuses["novox/app@"+head+" mesh/delivery"], "success ready") {
t.Fatalf("the head's status is %q", w.forge.statuses)
}
if v := w.forge.views["novox/app#7"]; !strings.Contains(v, viewMarker) || !strings.Contains(v, "**ready**") {
t.Fatalf("the pull request's view is %q", v)
}
// A failing verdict for another head rejects it, and a newer head supersedes the older.
w.h.PullUpdated(pr("novox/app", 7, "bbbbbbbbbbbb2222", "feat/x"))
if w.state(id) != Superseded {
t.Fatalf("a newer head left the older %s", w.state(id))
}
w.h.Checked(verdict("novox/app", 7, "bbbbbbbbbbbb2222", "fail"))
if s := w.state(IDOf("novox/app", "bbbbbbbbbbbb2222")); s != Rejected {
t.Fatalf("a failing verdict left it %s", s)
}
}
func statesOf(ss []State) []string {
var out []string
for _, s := range ss {
out = append(out, string(s))
}
return out
}
// A merge publishes; the walk waits for this owner's word; the word is given; the walk's record moves the
// delivery through its machines to delivered — and the note on the merge commit says what was executed.
func TestAMergeIsPublishedLetGoDeliveringAndDelivered(t *testing.T) {
w := newWorld(t)
id := IDOf("novox/app", head)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "pass"))
// The walk is heard before the merge: it waits for its delivery.
walk := w.ctl.put(aWalk("plan-1", "novox/app", merge, true))
w.h.WalkMoved(walk)
if w.state(id) != Ready {
t.Fatalf("a walk for a commit no merge was heard of moved the delivery to %s", w.state(id))
}
w.h.PullMerged(merged("novox/app", 7, head, merge))
if w.state(id) != Published {
t.Fatalf("a merge with its walk open left it %s", w.state(id))
}
w.settleAll()
if len(w.ctl.delivers) != 1 || w.ctl.delivers[0] != "plan-1" {
t.Fatalf("its walk was let go %v", w.ctl.delivers)
}
if w.state(id) != Delivering {
t.Fatalf("let go, it is %s", w.state(id))
}
// The first machine sent, judged, passed; the rest sent; the walk done.
sent := time.Date(2026, 10, 6, 12, 40, 0, 0, time.UTC)
walk = w.ctl.walk("plan-1")
walk.State = walkRolling
walk.Modules["app"] = &WalkModule{State: "built", First: []string{"anchor"}, FirstAt: &sent, Build: "build-1",
Gate: &WalkGate{Machines: []string{"anchor"}, Since: &sent, Passes: 1}}
w.h.WalkMoved(w.ctl.put(walk))
if d := w.delivery(id); len(d.Steps) != 1 || d.Steps[0].State != StepJudging {
t.Fatalf("judging on its first machine reads %+v", d.Steps)
}
judged := sent.Add(3 * time.Minute)
walk.Modules["app"].Gate.Verdict, walk.Modules["app"].Gate.JudgedAt = "passed", &judged
walk.Modules["app"].SentAt = &judged
walk.State = walkDone
w.h.WalkMoved(w.ctl.put(walk))
d := w.delivery(id)
if d.State != Delivered {
t.Fatalf("a walk done left it %s", d.State)
}
if len(d.Steps) != 2 || d.Steps[0].State != StepPassed || d.Steps[1].Machine != "the rest running it" {
t.Fatalf("its machine steps are %+v", d.Steps)
}
w.settleAll()
onMerge := w.forge.notesOn(merge)
if len(onMerge) == 0 || !strings.Contains(strings.Join(onMerge, "\n"), "executed (walk plan-1): app on anchor passed") {
t.Fatalf("the merge commit's note is %q", onMerge)
}
// Said on the bus, kept first: every transition emitted once.
n := 0
for _, s := range w.said {
if s["id"] == id {
n++
}
}
if n != len(d.Transitions) {
t.Fatalf("%d transitions, %d said", len(d.Transitions), n)
}
}
// A first machine's gate fails: what it carried is put back (the controller's), the step is rolled back,
// and the delivery failed — never delivered.
func TestAFailedGateIsRolledBackAndTheDeliveryFails(t *testing.T) {
w := newWorld(t)
id := IDOf("novox/app", head)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "pass"))
w.h.PullMerged(merged("novox/app", 7, head, merge))
w.h.WalkMoved(w.ctl.put(aWalk("plan-1", "novox/app", merge, true)))
w.settleAll()
sent := time.Date(2026, 10, 6, 12, 40, 0, 0, time.UTC)
walk := w.ctl.walk("plan-1")
walk.State = walkRolling
walk.Modules["app"] = &WalkModule{State: "built", First: []string{"anchor"}, FirstAt: &sent,
Gate: &WalkGate{Machines: []string{"anchor"}, Since: &sent}}
w.h.WalkMoved(w.ctl.put(walk))
at := sent.Add(10 * time.Minute)
walk.Modules["app"].Gate.Verdict, walk.Modules["app"].Gate.Why, walk.Modules["app"].Gate.JudgedAt =
"failed", "its tools were not served", &at
walk.Modules["app"].Gate.Rollback = "rolled-back"
walk.State = walkFailed
walk.Note = "app failed its gate on anchor"
w.h.WalkMoved(w.ctl.put(walk))
d := w.delivery(id)
if d.State != Failed || len(d.Steps) != 1 || d.Steps[0].State != StepRolledBack {
t.Fatalf("a failed gate left it %s, steps %+v", d.State, d.Steps)
}
// An older record arriving late changes nothing; a step the table does not hold is refused.
w.h.WalkMoved(aWalk("plan-1", "novox/app", merge, false))
if w.state(id) != Failed {
t.Fatal("an older record of the walk moved a final delivery")
}
}
// A newer merge's walk took over this one's: superseded, from published or delivering.
func TestANewerWalkSupersedesAnOlderDelivery(t *testing.T) {
w := newWorld(t)
id := IDOf("novox/app", head)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "pass"))
w.h.PullMerged(merged("novox/app", 7, head, merge))
walk := aWalk("plan-1", "novox/app", merge, true)
w.h.WalkMoved(w.ctl.put(walk))
walk.State = walkSuperseded
walk.Note = "superseded at tier 0 by plan-2"
w.ctl.down = true // its word could not be given meanwhile
w.settleAll()
w.h.WalkMoved(w.ctl.put(walk))
if s := w.state(id); s != Superseded {
t.Fatalf("its walk superseded left it %s", s)
}
}
// The trunk rule: a ready delivery whose commit never reached the trunk is never published, whatever walks
// are heard of; and a walk for a commit no delivery landed as is not taken as its.
func TestACommitOffTheTrunkIsNeverPublished(t *testing.T) {
w := newWorld(t)
id := IDOf("novox/app", head)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "pass"))
w.h.WalkMoved(w.ctl.put(aWalk("plan-9", "novox/app", head, true))) // a walk of the head itself: no merge
w.settleAll()
if s := w.state(id); s != Ready {
t.Fatalf("an off-trunk head is %s", s)
}
if len(w.ctl.delivers) != 0 {
t.Fatalf("a walk was let go for a commit off the trunk: %v", w.ctl.delivers)
}
}
// A restart resumes: a new holder over the same state stands every delivery where it was, reads the walks it
// missed back from the controller, and goes on.
func TestARestartResumesFromTheState(t *testing.T) {
w := newWorld(t)
id := IDOf("novox/app", head)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "pass"))
w.h.PullMerged(merged("novox/app", 7, head, merge))
w.ctl.put(aWalk("plan-1", "novox/app", merge, true))
if err := w.h.Reconcile(); err != nil {
t.Fatal(err)
}
w.settleAll()
if s := w.state(id); s != Delivering {
t.Fatalf("before the restart it is %s", s)
}
// While it is down, the walk finishes and nobody says so.
walk := w.ctl.walk("plan-1")
walk.State = walkDone
w.ctl.put(walk)
again := newWorldOver(t, w.store, w.ctl, w.forge, w.now.Add(time.Hour))
if s := again.state(id); s != Delivering {
t.Fatalf("read back, it is %s", s)
}
if err := again.h.Reconcile(); err != nil {
t.Fatal(err)
}
if s := again.state(id); s != Delivered {
t.Fatalf("the walk it missed was not read back: %s", s)
}
if len(w.ctl.delivers) != 1 {
t.Fatalf("a restarted holder let the walk go again: %v", w.ctl.delivers)
}
}
// Kept before it is said: with the state unreachable, nothing is emitted or noted; once it is kept, all of it.
func TestATransitionIsKeptBeforeItIsSaid(t *testing.T) {
w := newWorld(t)
w.store.failing = true
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.settleAll()
if len(w.said) != 0 || len(w.forge.notesOn(head)) != 0 {
t.Fatalf("said before it was kept: %v %v", w.said, w.forge.notesOn(head))
}
w.store.failing = false
w.settleAll()
if len(w.said) != 1 || len(w.forge.notesOn(head)) != 1 {
t.Fatalf("once kept, said %d and noted %d", len(w.said), len(w.forge.notesOn(head)))
}
// A forge that is away is asked again; a note is appended once however often it is asked.
w.forge.down = true
w.h.Checked(verdict("novox/app", 7, head, "pass"))
w.settleAll()
if len(w.forge.notesOn(head)) != 1 {
t.Fatal("noted while the forge was away")
}
w.forge.down = false
w.settleAll()
w.settleAll()
if notes := w.forge.notesOn(head); len(notes) != 3 {
t.Fatalf("after the forge came back the note has %d line(s): %q", len(notes), notes)
}
if d := w.delivery(IDOf("novox/app", head)); len(d.Owed) != 0 {
t.Fatalf("still owed: %+v", d.Owed)
}
}
// At the switch: a walk already running is adopted as a delivery that is delivering; a waiting walk whose pull
// request this owner never heard of waits for a person, and a person's release lets it go.
func TestTheSwitchAdoptsRunningWalksAndHoldsUnknownWaitingOnes(t *testing.T) {
w := newWorld(t)
running := aWalk("plan-old", "novox/mesh-catalog", "cccccccccccc3333", false)
running.Created = w.now.Add(-time.Hour)
running.State = walkRolling
w.ctl.put(running)
waiting := aWalk("plan-new", "novox/app", "dddddddddddd4444", true)
waiting.Created = w.now.Add(time.Minute)
w.ctl.put(waiting)
if err := w.h.Reconcile(); err != nil {
t.Fatal(err)
}
w.h.Tick()
if s := w.state(IDOf("novox/mesh-catalog", "cccccccccccc3333")); s != Delivering {
t.Fatalf("a running walk at the switch was adopted as %s", s)
}
if _, made := w.h.deliveries[IDOf("novox/app", "dddddddddddd4444")]; made {
t.Fatal("a new walk was made a delivery before its merge had time to be heard")
}
w.later(3 * time.Minute)
w.h.Tick()
held := IDOf("novox/app", "dddddddddddd4444")
if s := w.state(held); s != Held {
t.Fatalf("a waiting walk with no pull request is %s", s)
}
if _, err := w.h.Release(held, "", "jochen"); err == nil {
t.Fatal("released without why")
}
if _, err := w.h.Release(held, "checked by hand", "jochen"); err != nil {
t.Fatal(err)
}
if s := w.state(held); s != Delivering || len(w.ctl.delivers) != 1 {
t.Fatalf("released, it is %s and let go %v", s, w.ctl.delivers)
}
}
// mesh-delivery's controller down: nothing is lost — its word is asked again on the next tick.
func TestAControllerThatDoesNotAnswerIsAskedAgain(t *testing.T) {
w := newWorld(t)
id := IDOf("novox/app", head)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "pass"))
w.h.PullMerged(merged("novox/app", 7, head, merge))
w.h.WalkMoved(w.ctl.put(aWalk("plan-1", "novox/app", merge, true)))
w.ctl.down = true
w.settleAll()
if s := w.state(id); s != Published {
t.Fatalf("with the controller down it is %s", s)
}
w.ctl.down = false
w.settleAll()
if s := w.state(id); s != Delivering {
t.Fatalf("with the controller back it is %s", s)
}
}
// A merge before the check passed is held for a person; a merge this owner never heard the head of is made
// from the forge's word and published when the forge says its gate passed.
func TestAMergeWithoutAPassingCheckIsHeld(t *testing.T) {
w := newWorld(t)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "fail"))
w.h.PullMerged(merged("novox/app", 7, head, merge))
if s := w.state(IDOf("novox/app", head)); s != Held {
t.Fatalf("a rejected head merged is %s", s)
}
// Never heard: the forge's statuses on its head decide.
m := merged("novox/lab", 3, "eeeeeeeeeeee5555", "ffffffffffff6666")
m.HeadChecks = map[string]string{"mesh/merge-gate": "success", "mesh/repo-check": "warning"}
m.Paths = []string{"README.md"}
w.h.PullMerged(m)
w.settleAll()
d := w.delivery(IDOf("novox/lab", "eeeeeeeeeeee5555"))
// The fake planner says it moves app; with no walk yet it waits published.
if d.State != Ready && d.State != Published {
t.Fatalf("a merge whose head passed by the forge's word is %s", d.State)
}
}
// A delivery group (ADR 0239 decision 4): two heads sharing a branch name across repositories, ordered by the
// controller, checked together, ready only when both are and their composition passed — and delivered in
// order, whatever order they merged in.
func TestAGroupIsOrderedCheckedTogetherAndDeliveredInOrder(t *testing.T) {
w := newWorld(t)
const ctlHead, catHead = "c1c1c1c1c1c1c1c1", "d2d2d2d2d2d2d2d2"
ctl, cat := IDOf("novox/mesh-controller", ctlHead), IDOf("novox/mesh-catalog", catHead)
w.ctl.order = func(ms []Member) Order {
return Order{Order: []string{ctl, cat}, Pairs: []Pair{{Before: ctl, After: cat, Why: "version skew"}}}
}
w.h.PullUpdated(pr("novox/mesh-controller", 1, ctlHead, "feat/x", "internal/a.go"))
w.h.PullUpdated(pr("novox/mesh-catalog", 2, catHead, "feat/x", "modules/app/module.json"))
if w.delivery(ctl).Group != "feat/x" || w.delivery(cat).Group != "feat/x" {
t.Fatalf("one branch name in two repositories made no group: %q %q", w.delivery(ctl).Group, w.delivery(cat).Group)
}
w.h.Checked(verdict("novox/mesh-controller", 1, ctlHead, "pass"))
w.settleAll()
if len(w.ctl.checks) != 0 {
t.Fatal("the group was checked composed before every member was ready")
}
w.h.Checked(verdict("novox/mesh-catalog", 2, catHead, "pass"))
w.settleAll()
if len(w.ctl.checks) != 1 || w.ctl.checks[0] != "feat/x:"+ctl+","+cat {
t.Fatalf("the composed check asked %v, wanted the members in order", w.ctl.checks)
}
if state, _ := GroupState(w.h.groups["feat/x"], w.h.membersOf(w.h.groups["feat/x"])); state != "checking" {
t.Fatalf("before its composed verdict the group is %s", state)
}
w.h.Checked(CheckedEvent{Group: "feat/x", ID: "check-1", Verdict: "pass", Summary: "both compose"})
w.settleAll()
if state, _ := GroupState(w.h.groups["feat/x"], w.h.membersOf(w.h.groups["feat/x"])); state != "ready" {
t.Fatalf("composed and ready, the group is %s", state)
}
if !strings.HasPrefix(w.forge.statuses["novox/mesh-catalog@"+catHead+" mesh/delivery-group"], "success") {
t.Fatalf("the group's status on a member's head: %v", w.forge.statuses)
}
// The catalogue merges first; its walk waits — the controller has not been delivered.
w.h.PullMerged(merged("novox/mesh-catalog", 2, catHead, "m2m2m2m2m2m2"))
w.h.WalkMoved(w.ctl.put(aWalk("plan-cat", "novox/mesh-catalog", "m2m2m2m2m2m2", true)))
w.settleAll()
if w.state(cat) != Published || len(w.ctl.delivers) != 0 {
t.Fatalf("a member merged before the one it goes after was let go: %s %v", w.state(cat), w.ctl.delivers)
}
// The controller merges: its walk is on the controller's own path, started at the merge.
w.h.PullMerged(merged("novox/mesh-controller", 1, ctlHead, "m1m1m1m1m1m1"))
ctlWalk := w.ctl.put(aWalk("plan-ctl", "novox/mesh-controller", "m1m1m1m1m1m1", false))
w.h.WalkMoved(ctlWalk)
w.settleAll()
if w.state(ctl) != Delivering || len(w.ctl.delivers) != 0 {
t.Fatalf("the controller's member is %s, let go %v", w.state(ctl), w.ctl.delivers)
}
ctlWalk.State = walkDone
w.h.WalkMoved(w.ctl.put(ctlWalk))
w.settleAll()
if w.state(ctl) != Delivered || len(w.ctl.delivers) != 1 || w.ctl.delivers[0] != "plan-cat" {
t.Fatalf("once the controller was delivered, the catalogue was not let go: %s %v", w.state(ctl), w.ctl.delivers)
}
if w.state(cat) != Delivering {
t.Fatalf("the catalogue's member is %s", w.state(cat))
}
}
// A member that fails stops every member after it, naming it; what was delivered before it stays.
func TestAFailedMemberStopsTheMembersAfterItAndLeavesThoseBefore(t *testing.T) {
w := newWorld(t)
heads := map[string]string{"novox/a": "a1a1a1a1a1a1a1", "novox/b": "b2b2b2b2b2b2b2", "novox/c": "c3c3c3c3c3c3c3"}
ids := map[string]string{}
for i, repo := range []string{"novox/a", "novox/b", "novox/c"} {
ids[repo] = IDOf(repo, heads[repo])
w.h.PullUpdated(pr(repo, i+1, heads[repo], "feat/y", "x"))
w.h.Checked(verdict(repo, i+1, heads[repo], "pass"))
}
w.settleAll()
w.h.Checked(CheckedEvent{Group: "feat/y", ID: "check-1", Verdict: "pass", Summary: "compose"})
for i, repo := range []string{"novox/a", "novox/b", "novox/c"} {
mergeCommit := strings.Repeat(string(rune('p'+i)), 12)
w.h.PullMerged(merged(repo, i+1, heads[repo], mergeCommit))
w.h.WalkMoved(w.ctl.put(aWalk("plan-"+repo[6:], repo, mergeCommit, true)))
}
w.settleAll()
a := w.ctl.walk("plan-a")
a.State = walkDone
w.h.WalkMoved(w.ctl.put(a))
w.settleAll()
b := w.ctl.walk("plan-b")
if b.Delivery == nil || b.Delivery.Go == nil {
t.Fatalf("b was not let go after a was delivered: %v", w.ctl.delivers)
}
b.State = walkFailed
b.Note = "b failed its gate on anchor"
w.h.WalkMoved(w.ctl.put(b))
w.settleAll()
if w.state(ids["novox/a"]) != Delivered || w.state(ids["novox/b"]) != Failed || w.state(ids["novox/c"]) != Stopped {
t.Fatalf("a %s, b %s, c %s", w.state(ids["novox/a"]), w.state(ids["novox/b"]), w.state(ids["novox/c"]))
}
c := w.delivery(ids["novox/c"])
if why := c.Transitions[len(c.Transitions)-1].Why; !strings.Contains(why, ids["novox/b"]) {
t.Fatalf("c's stop does not name the member that failed: %q", why)
}
if len(w.ctl.stops) != 1 || w.ctl.stops[0] != "plan-c" {
t.Fatalf("c's walk was not ended: %v", w.ctl.stops)
}
if state, _ := GroupState(w.h.groups["feat/y"], w.h.membersOf(w.h.groups["feat/y"])); state != "failed" {
t.Fatalf("the group is %s", state)
}
}
// A group whose declared order contradicts its inferred one is rejected, and checked never.
func TestAGroupWithACycleIsRejected(t *testing.T) {
w := newWorld(t)
w.ctl.order = func(ms []Member) Order { return Order{Cycle: []string{ms[0].ID, ms[1].ID}} }
w.h.PullUpdated(pr("novox/a", 1, "a1a1a1a1a1a1a1", "feat/z"))
w.h.PullUpdated(pr("novox/b", 2, "b2b2b2b2b2b2b2", "feat/z"))
w.h.Checked(verdict("novox/a", 1, "a1a1a1a1a1a1a1", "pass"))
w.h.Checked(verdict("novox/b", 2, "b2b2b2b2b2b2b2", "pass"))
w.settleAll()
g := w.h.groups["feat/z"]
if state, why := GroupState(g, w.h.membersOf(g)); state != "rejected" || !strings.Contains(why, "contradicts") {
t.Fatalf("a cycle left the group %s: %s", state, why)
}
if len(w.ctl.checks) != 0 {
t.Fatal("a group with a cycle was checked")
}
}
// Stalled: a delivery past its state's bound is listed with what H2 may do; close takes only that.
func TestStalledAndCloseWorkFromTheTable(t *testing.T) {
w := newWorld(t)
id := IDOf("novox/app", head)
w.h.PullUpdated(pr("novox/app", 7, head, "feat/x", "src/a.go"))
w.h.Checked(verdict("novox/app", 7, head, "pass"))
w.h.PullMerged(merged("novox/app", 7, head, merge))
walk := w.ctl.put(aWalk("plan-1", "novox/app", merge, false))
w.h.WalkMoved(walk)
if w.state(id) != Delivering {
t.Fatalf("it is %s", w.state(id))
}
if _, err := w.h.Close(id, "too early"); err == nil {
t.Fatal("closed within its bound")
}
w.later(3 * time.Hour)
stalled := w.h.Stalled()
if len(stalled) != 1 || stalled[0].ID != id || !strings.Contains(stalled[0].H2, "close") {
t.Fatalf("stalled %+v", stalled)
}
if _, err := w.h.Close(id, "delivery stalled"); err == nil {
t.Fatal("closed with nothing in the walk's record saying it moved")
}
walk.State = walkDone
w.ctl.put(walk) // finished, and the event was missed
said, err := w.h.Close(id, "delivery stalled")
if err != nil || w.state(id) != Delivered {
t.Fatalf("H2's close left it %s: %v", w.state(id), err)
}
if !strings.Contains(said, "delivered") {
t.Fatalf("close says %q", said)
}
// Held is the operator's: H2 has no transition there.
w.h.PullUpdated(pr("novox/app", 8, "ffffffffffff0000", "feat/q"))
w.h.PullMerged(merged("novox/app", 8, "ffffffffffff0000", "eeeeeeeeeeee0000"))
w.later(48 * time.Hour)
if _, err := w.h.Close(IDOf("novox/app", "ffffffffffff0000"), "x"); err == nil || !strings.Contains(err.Error(), "operator's") {
t.Fatalf("H2 closed a held delivery: %v", err)
}
}
@@ -0,0 +1,253 @@
// mesh-delivery: the holder of the mesh-delivery seat (novox/hq ADR 0239, to-be 47). A Go bundle the node's
// runtime launches. It owns the delivery — one commit in one repository, from its pull request's head to
// every machine — and the delivery group above it: their states in one compiled table, kept in its own state
// on the bus, every transition said as an event, noted on the commit and shown on the pull request through
// the forge's holder. It never sends to a machine: it asks the controller, through its seat's verbs. stdout
// is the MCP channel; what this module says, it says on stderr.
package main
import (
"encoding/json"
"fmt"
"os"
"strings"
"sync"
"time"
stdio "git.novox.be/novox/mesh-sdk/go"
)
func logf(format string, a ...any) { fmt.Fprintf(os.Stderr, format+"\n", a...) }
// listening is whether the events this owner follows reach it, in words.
type listening struct {
mu sync.Mutex
now string
}
func (l *listening) set(s string) { l.mu.Lock(); l.now = s; l.mu.Unlock() }
func (l *listening) get() string { l.mu.Lock(); defer l.mu.Unlock(); return l.now }
// The events this owner follows: the forge's pull requests, and the controller's verdicts and walks.
var followed = []string{"gitea.pull.updated", "gitea.pull.merged", "gitea.pull.closed",
ControllerSeat + ".checked", ControllerSeat + ".plan-moved"}
func main() {
h := &Holder{
Store: kvStore{},
Controller: seatController{ask: stdio.Ask},
Forge: toolForge{ask: stdio.Ask},
Emit: func(event string, body any) error { return stdio.Emit(event, body) },
Logf: logf,
}
h.init()
l := &listening{now: "not yet: starting"}
go run(h, l)
if err := stdio.Serve("", tools(h, l)); err != nil {
logf("%v", err)
os.Exit(1)
}
}
// run reads back the state, then the controller's walks, then takes the events, and keeps time — each
// retried, each failure said, never given up on quietly.
func run(h *Holder, l *listening) {
time.Sleep(500 * time.Millisecond) // Serve first: the state is reached through it
for wait := 2 * time.Second; ; wait = min(wait*2, time.Minute) {
err := h.Load()
if err == nil {
break
}
logf("[mesh-delivery] cannot read back its deliveries yet (%v); asking again in %s", err, wait)
time.Sleep(wait)
}
if err := h.Reconcile(); err != nil {
logf("[mesh-delivery] the controller's walks cannot be read yet (%v); asked again each half minute", err)
}
go func() {
tick := time.NewTicker(10 * time.Second)
defer tick.Stop()
last := time.Now()
for now := range tick.C {
if now.Sub(last) >= 30*time.Second {
last = now
if err := h.Reconcile(); err != nil {
logf("[mesh-delivery] the controller's walks cannot be read: %v", err)
}
}
h.Tick()
h.Flush()
}
}()
for _, pattern := range followed {
pattern := pattern
for wait := 2 * time.Second; ; wait = min(wait*2, time.Minute) {
err := stdio.Subscribe(pattern, func(e stdio.Envelope) error { return handle(h, e) })
if err == nil {
break
}
l.set("not yet: " + err.Error())
logf("[mesh-delivery] not hearing %s yet (%v); asking again in %s", pattern, err, wait)
time.Sleep(wait)
}
}
l.set("listening")
logf("[mesh-delivery] listening: %s", strings.Join(followed, ", "))
}
// handle takes one event. An event that cannot be read is said and acknowledged: read again it would fail
// again, and the controller's walks are read back by comparison anyway.
func handle(h *Holder, e stdio.Envelope) error {
switch e.Key {
case "gitea.pull.updated", "gitea.pull.merged", "gitea.pull.closed":
var p PullEvent
if err := json.Unmarshal(e.Body, &p); err != nil {
h.Refuse(fmt.Errorf("%s is not readable: %v", e.Key, err))
return nil
}
switch e.Key {
case "gitea.pull.updated":
h.PullUpdated(p)
case "gitea.pull.merged":
h.PullMerged(p)
default:
h.PullClosed(p)
}
case ControllerSeat + ".checked":
var c CheckedEvent
if err := json.Unmarshal(e.Body, &c); err != nil {
h.Refuse(fmt.Errorf("a verdict is not readable: %v", err))
return nil
}
h.Checked(c)
case ControllerSeat + ".plan-moved":
var w Walk
if err := json.Unmarshal(e.Body, &w); err != nil {
h.Refuse(fmt.Errorf("a walk is not readable: %v", err))
return nil
}
h.WalkMoved(w)
}
return nil
}
func str(description string) map[string]any {
return map[string]any{"type": "string", "description": description}
}
func strArg(a map[string]any, k string) string { s, _ := a[k].(string); return strings.TrimSpace(s) }
func boolArg(a map[string]any, k string) bool {
switch v := a[k].(type) {
case bool:
return v
case string:
return v == "true"
}
return false
}
// aPerson is who an act through the seat is by: the seat does not name its caller.
const aPerson = "a person, through the mesh-delivery seat"
func need(a map[string]any, keys ...string) error {
for _, k := range keys {
if strArg(a, k) == "" {
return fmt.Errorf("needs %q", k)
}
}
return nil
}
func tools(h *Holder, l *listening) []stdio.Tool {
seat := "mesh-delivery."
return []stdio.Tool{
{Name: seat + "deliveries",
Description: "Every delivery not final, and those that ended in the last day, one line each: its id " +
"(owner/repository@commit), its state, its group, what it waits for and since when. Narrowed by state, " +
"repository or group.",
Input: map[string]any{"type": "object", "properties": map[string]any{"state": str("one state"),
"repository": str("owner/repository"), "group": str("a group's id (its branch name)"),
"all": map[string]any{"type": "string", "enum": []string{"true", "false"}}}},
Run: func(a map[string]any) (any, error) {
return h.Deliveries(strArg(a, "state"), strArg(a, "repository"), strArg(a, "group"), boolArg(a, "all")), nil
}},
{Name: seat + "show",
Description: "One delivery or group whole: its delivery plan, every transition with when and why, the " +
"machine steps of its walk, its group and its order.",
Input: map[string]any{"type": "object", "properties": map[string]any{"id": str("a delivery's id, or a group's")},
"required": []string{"id"}},
Run: func(a map[string]any) (any, error) { return h.Show(strArg(a, "id")) }},
{Name: seat + "groups",
Description: "Every delivery group: its members in order, why each pair is ordered, its composed check " +
"and its state, derived from its members.",
Input: map[string]any{"type": "object", "properties": map[string]any{
"all": map[string]any{"type": "string", "enum": []string{"true", "false"}}}},
Run: func(a map[string]any) (any, error) { return h.Groups(boolArg(a, "all")), nil }},
{Name: seat + "what-if",
Description: "The delivery plan a change would have, asked of the controller's planner and kept nowhere.",
Input: map[string]any{"type": "object", "properties": map[string]any{"repository": str("owner/repository"),
"paths": str("the files it changes, comma-separated"), "base": str("the branch it merges into (default main)")},
"required": []string{"repository", "paths"}},
Run: func(a map[string]any) (any, error) {
if err := need(a, "repository", "paths"); err != nil {
return nil, err
}
return h.WhatIf(strArg(a, "repository"), strArg(a, "base"), strings.Split(strArg(a, "paths"), ","))
}},
{Name: seat + "table",
Description: "The state table every delivery runs by: each transition with its guard, each state's bound " +
"and what healer H2 may do once it has passed; and the machine steps' table.",
Run: func(map[string]any) (any, error) { return TableText(), nil }},
{Name: seat + "stalled",
Description: "Every delivery held past its state's bound, with the transition the table lets healer H2 take.",
Run: func(map[string]any) (any, error) { return h.Stalled(), nil }},
{Name: seat + "recheck",
Description: "Check a rejected or ready delivery again: proposed again, and its check asked. With why.",
Input: map[string]any{"type": "object", "properties": map[string]any{"id": str("the delivery's id"),
"why": str("why, kept with the transition")}, "required": []string{"id", "why"}},
Run: func(a map[string]any) (any, error) {
if err := need(a, "id", "why"); err != nil {
return nil, err
}
return h.Recheck(strArg(a, "id"), strArg(a, "why"), aPerson)
}},
{Name: seat + "release",
Description: "A person's word that a held delivery goes on: it starts delivering. With why.",
Input: map[string]any{"type": "object", "properties": map[string]any{"id": str("the delivery's id"),
"why": str("why, kept with the transition")}, "required": []string{"id", "why"}},
Run: func(a map[string]any) (any, error) {
if err := need(a, "id", "why"); err != nil {
return nil, err
}
return h.Release(strArg(a, "id"), strArg(a, "why"), aPerson)
}},
{Name: seat + "stop",
Description: "Stop a delivery that is not final, with why: its walk is ended through the controller, and " +
"in a delivering group every member after it is stopped too, naming it.",
Input: map[string]any{"type": "object", "properties": map[string]any{"id": str("the delivery's id"),
"why": str("why, kept with the transition")}, "required": []string{"id", "why"}},
Run: func(a map[string]any) (any, error) {
if err := need(a, "id", "why"); err != nil {
return nil, err
}
return h.Stop(strArg(a, "id"), strArg(a, "why"), aPerson)
}},
{Name: seat + "close",
Description: "Healer H2's verb: for a delivery held past its bound, read its walk again and take the " +
"transition the table names for that state. Refused for any other.",
Input: map[string]any{"type": "object", "properties": map[string]any{"id": str("the delivery's id"),
"why": str("the condition it answers")}, "required": []string{"id", "why"}},
Run: func(a map[string]any) (any, error) {
if err := need(a, "id", "why"); err != nil {
return nil, err
}
return h.Close(strArg(a, "id"), strArg(a, "why"))
}},
{Name: "delivery_status",
Description: "Whether this owner works: how many deliveries in each state, the groups, what is owed to " +
"the forge and the bus and not yet done, walks with no delivery yet, whether its events arrive, when " +
"the controller's walks were last read or why they could not be, and the last refusals.",
Run: func(map[string]any) (any, error) { return h.Status(l.get()), nil }},
}
}
@@ -0,0 +1,121 @@
package main
import (
"encoding/json"
"os"
"path/filepath"
"reflect"
"sort"
"strings"
"testing"
)
// The manifest says what the code does: the seat it claims and the verbs it serves, the events it follows
// and says, the tools it calls, its state — and names nothing of one installation.
type manifest struct {
Module string `json:"module"`
Claims []struct {
Name string `json:"name"`
Scope string `json:"scope"`
Serves []string `json:"serves"`
} `json:"claims"`
Seats []any `json:"seats"`
Consumes []string `json:"consumes"`
Emits []string `json:"emits"`
Invokes []string `json:"invokes"`
State []string `json:"state"`
Tools []string `json:"tools"`
}
func readManifest(t *testing.T) (manifest, string) {
t.Helper()
raw, err := os.ReadFile(filepath.Join("..", "..", "module.json"))
if err != nil {
t.Fatal(err)
}
var m manifest
if err := json.Unmarshal(raw, &m); err != nil {
t.Fatal(err)
}
return m, string(raw)
}
func TestItClaimsTheDeliverySeatAndFollowsWhatItHandles(t *testing.T) {
m, raw := readManifest(t)
if m.Module != "mesh-delivery" || len(m.Claims) != 1 || m.Claims[0].Name != "mesh-delivery" || m.Claims[0].Scope != "mesh" {
t.Fatalf("claims %+v", m.Claims)
}
// The seat is the mesh's own (mesh-*): a module claims it and never declares it.
if len(m.Seats) != 0 {
t.Fatal("the manifest declares a seat; mesh-delivery is the mesh's own")
}
if !reflect.DeepEqual(m.Consumes, followed) {
t.Fatalf("consumes %v, follows %v", m.Consumes, followed)
}
if !reflect.DeepEqual(m.Emits, []string{"transition", "group"}) {
t.Fatalf("emits %v", m.Emits)
}
if !reflect.DeepEqual(m.State, []string{"deliveries", "groups"}) {
t.Fatalf("state %v", m.State)
}
// It calls exactly what its ports ask: the controller's six verbs and the forge's three tools.
want := []string{}
for _, v := range []string{"delivery-plan", "delivery-order", "delivery-check", "deliver", "delivery-stop", "delivery-walks"} {
want = append(want, "seat:"+ControllerSeat+"."+v)
}
want = append(want, "gitea.gitea_note_append", "gitea.gitea_delivery_view", "gitea.gitea_commit_status")
if !reflect.DeepEqual(m.Invokes, want) {
t.Fatalf("invokes %v, wanted %v", m.Invokes, want)
}
for _, never := range []string{"/home/", "jochen", "g14", "shanks", "novox", "zurag", "anchor"} {
if strings.Contains(strings.ToLower(raw), never) {
t.Errorf("module.json names %q", never)
}
}
}
func TestTheToolsAgreeWithTheManifest(t *testing.T) {
m, _ := readManifest(t)
var own, verbs []string
for _, tool := range tools(&Holder{}, &listening{}) {
if strings.TrimSpace(tool.Description) == "" {
t.Errorf("%s has no description", tool.Name)
}
if seat, verb, ok := strings.Cut(tool.Name, "."); ok {
if seat != "mesh-delivery" {
t.Errorf("%s is a verb of a seat this does not hold", tool.Name)
}
verbs = append(verbs, verb)
continue
}
own = append(own, tool.Name)
}
sort.Strings(own)
if !reflect.DeepEqual(own, m.Tools) {
t.Fatalf("serves %v, lists %v", own, m.Tools)
}
if !reflect.DeepEqual(verbs, m.Claims[0].Serves) {
t.Fatalf("verbs %v, claim %v", verbs, m.Claims[0].Serves)
}
}
// The forge's and the controller's answers are read whatever wraps them.
func TestAnAnswerIsReadWhateverWrapsIt(t *testing.T) {
for _, c := range []struct {
raw string
ok bool
answer string
}{
{`{"output":"x","ok":true,"answer":{"held":true}}`, true, `{"held":true}`},
{`{"output":"refused: no plan","ok":false}`, false, ``},
{`"{\"output\":\"x\",\"ok\":true,\"answer\":[1]}"`, true, `[1]`},
{`{"content":[{"type":"text","text":"{\"output\":\"x\",\"ok\":true,\"answer\":2}"}]}`, true, `2`},
{`{"content":[{"type":"text","text":"boom"}],"isError":true}`, false, ``},
} {
answer, _, ok := answerOf(json.RawMessage(c.raw))
if ok != c.ok || (c.ok && string(answer) != c.answer) {
t.Errorf("%s → %s %v", c.raw, answer, ok)
}
}
}
@@ -0,0 +1,436 @@
package main
import (
"encoding/json"
"fmt"
"regexp"
"sort"
"strings"
"time"
)
// What mesh-delivery owns (novox/hq ADR 0239): the delivery — one commit in one repository — and the
// delivery group above it. Each is kept whole, one key of the module's own state per delivery and per
// group, written only by this holder.
// Delivery is one commit in one repository, from its pull request's head to every machine.
type Delivery struct {
// ID is owner/repository@<twelve characters of its commit>.
ID string `json:"id"`
Repository string `json:"repository"`
Commit string `json:"commit"`
// The pull request, when there is one: its number, title, base, head branch and page.
Number int `json:"number,omitempty"`
Title string `json:"title,omitempty"`
Base string `json:"base,omitempty"`
Branch string `json:"branch,omitempty"`
HTMLURL string `json:"html_url,omitempty"`
CloneURL string `json:"clone_url,omitempty"`
// The diffset, as the forge said it: what the planner maps onto modules.
Paths []string `json:"paths,omitempty"`
PathsTruncated bool `json:"paths_truncated,omitempty"`
Removed []string `json:"removed,omitempty"`
ModuleDirs []string `json:"module_dirs,omitempty"`
ModuleDirsSaid bool `json:"module_dirs_said,omitempty"`
// After are the repositories its pull request says it goes after (`after: <repository>` lines).
After []string `json:"after,omitempty"`
// Group is the delivery group it is in: its head branch name. Empty for a lone delivery.
Group string `json:"group,omitempty"`
State State `json:"state"`
Since time.Time `json:"since"`
Created time.Time `json:"created"`
// Plan is its delivery plan, as the controller's planner computed it (the build plan, the deploy plan,
// the steps that are not an ordinary send).
Plan *DeliveryPlan `json:"plan,omitempty"`
// Check is its own check's verdict.
Check *Verdict `json:"check,omitempty"`
// The facts the table's guards read. NewerHead is the head that superseded it; ClosedUnmerged the pull
// request closed; MergedAs the commit it landed on the trunk as; HeldWhy what holds it for a person.
NewerHead string `json:"newer_head,omitempty"`
ClosedUnmerged bool `json:"closed_unmerged,omitempty"`
MergedAs string `json:"merged_as,omitempty"`
MergedAt time.Time `json:"merged_at"`
HeldWhy string `json:"held_why,omitempty"`
// Walk is the controller's walk of its trunk commit, as last read; Steps the machines it reached.
Walk *WalkSeen `json:"walk,omitempty"`
Steps []Step `json:"steps,omitempty"`
// LetGo is when this owner asked the controller to start its walk, and why.
LetGo time.Time `json:"let_go"`
LetGoWhy string `json:"let_go_why,omitempty"`
Transitions []Transition `json:"transitions,omitempty"`
// Owed are the forge's and the bus's effects of its transitions not yet done: said only after the
// delivery is kept, retried until done.
Owed []Effect `json:"owed,omitempty"`
}
// Transition is one move of the table, kept.
type Transition struct {
At time.Time `json:"at"`
From State `json:"from"`
To State `json:"to"`
Event Event `json:"event"`
Why string `json:"why"`
By string `json:"by,omitempty"`
}
// Verdict is a check's verdict: the gate's and the repository's own.
type Verdict struct {
ID string `json:"id,omitempty"`
At time.Time `json:"at"`
Gate string `json:"gate"`
Summary string `json:"summary"`
Repo string `json:"repo,omitempty"`
RepoSaid string `json:"repo_summary,omitempty"`
}
// Passes is whether the verdict lets the delivery be ready: the gate passed or warned, and the
// repository's own check did not fail or error.
func (v *Verdict) Passes() bool {
if v == nil {
return false
}
gate := v.Gate == "pass" || v.Gate == "warning"
repo := v.Repo == "" || v.Repo == "pass" || v.Repo == "warning"
return gate && repo
}
// DeliveryPlan is a delivery plan as the controller's planner says it (mesh-controller link.ChangePlan).
type DeliveryPlan struct {
Repository string `json:"repository"`
Base string `json:"base"`
Head string `json:"head"`
Moved []string `json:"moved,omitempty"`
Dependents []string `json:"dependents,omitempty"`
New []string `json:"new,omitempty"`
Unread []string `json:"unread,omitempty"`
Tiers [][]string `json:"tiers,omitempty"`
Machines []MachinePlan `json:"machines,omitempty"`
Steps []string `json:"steps,omitempty"`
Summary string `json:"summary"`
}
// MachinePlan is one machine's part of a plan.
type MachinePlan struct {
Machine string `json:"machine"`
Receives []string `json:"receives,omitempty"`
Waits []string `json:"waits,omitempty"`
}
// MovesNothing is whether the delivery's plan moves no module and adds none: nothing for a walk to do. A
// delivery with no plan known is not said to move nothing.
func (d *Delivery) MovesNothing() bool {
return d.Plan != nil && len(d.Plan.Moved)+len(d.Plan.New) == 0
}
// Step is one machine's part of a delivering walk.
type Step struct {
Machine string `json:"machine"`
Module string `json:"module"`
Build string `json:"build,omitempty"`
State StepState `json:"state"`
Since time.Time `json:"since"`
Why string `json:"why,omitempty"`
}
// Effect is one thing owed outside the state: an event said, a line appended to a commit's note, the view
// on the pull request refreshed, a status set.
type Effect struct {
Kind string `json:"kind"`
Commit string `json:"commit,omitempty"`
Line string `json:"line,omitempty"`
Context string `json:"context,omitempty"`
State string `json:"state,omitempty"`
Tries int `json:"tries,omitempty"`
Last string `json:"last,omitempty"`
Since time.Time `json:"since"`
}
// The kinds of effect.
const (
EffectEmit = "emit"
EffectNote = "note"
EffectView = "view"
EffectStatus = "status"
)
// Owner and Repo split the repository.
func (d *Delivery) Owner() string { o, _, _ := strings.Cut(d.Repository, "/"); return o }
func (d *Delivery) Repo() string { _, r, _ := strings.Cut(d.Repository, "/"); return r }
// IDOf is a delivery's id: the repository and twelve characters of its commit.
func IDOf(repository, commit string) string {
return strings.ToLower(repository) + "@" + shortOf(commit, 12)
}
func shortOf(commit string, n int) string {
if len(commit) > n {
return commit[:n]
}
return commit
}
// sameCommit is whether two commits are one, however either is abbreviated.
func sameCommit(a, b string) bool {
if a == "" || b == "" {
return false
}
if len(a) > len(b) {
a, b = b, a
}
return len(a) >= 7 && strings.HasPrefix(strings.ToLower(b), strings.ToLower(a))
}
// afterLines are the `after: <repository>` lines of a pull request's description.
var afterLines = regexp.MustCompile(`(?im)^\s*after:\s*([A-Za-z0-9_.\-/]+)\s*$`)
// AfterIn reads them.
func AfterIn(body string) []string {
var out []string
for _, m := range afterLines.FindAllStringSubmatch(body, -1) {
out = append(out, strings.TrimSuffix(m[1], ".git"))
}
sort.Strings(out)
return out
}
// The walk: the controller's plan of one trunk commit, as its `plan-moved` and `delivery-walks` say it
// (mesh-controller inventory.Plan). Read, never written here.
// The states of a walk as the controller keeps them.
const (
walkBuilding = "building"
walkRolling = "rolling"
walkDone = "done"
walkFailed = "failed"
walkSuperseded = "superseded"
)
// Walk is a walk whole, as the controller says it.
type Walk struct {
ID string `json:"id"`
Repository string `json:"repository"`
Branch string `json:"branch,omitempty"`
Commit string `json:"commit"`
Created time.Time `json:"created"`
Updated time.Time `json:"updated"`
State string `json:"state"`
Tier int `json:"tier"`
Tiers [][]string `json:"tiers"`
Modules map[string]*WalkModule `json:"modules"`
Note string `json:"note,omitempty"`
Revision int64 `json:"revision"`
Release *WalkRelease `json:"release,omitempty"`
Delivery *WalkDelivery `json:"delivery,omitempty"`
}
// WalkModule is one module of a walk.
type WalkModule struct {
State string `json:"state,omitempty"`
SentAt *time.Time `json:"sent_at,omitempty"`
First []string `json:"first,omitempty"`
FirstAt *time.Time `json:"first_at,omitempty"`
Commit string `json:"commit,omitempty"`
Why string `json:"why,omitempty"`
Build string `json:"build,omitempty"`
Gate *WalkGate `json:"gate,omitempty"`
GatedBy string `json:"gated_by,omitempty"`
}
// WalkGate is a module's judging on its first machine.
type WalkGate struct {
Machines []string `json:"machines"`
From string `json:"from,omitempty"`
To string `json:"to,omitempty"`
Since *time.Time `json:"since,omitempty"`
Passes int `json:"passes,omitempty"`
Verdict string `json:"verdict,omitempty"`
Why string `json:"why,omitempty"`
JudgedAt *time.Time `json:"judged_at,omitempty"`
Rollback string `json:"rollback,omitempty"`
Carried []WalkCarry `json:"carried,omitempty"`
}
// WalkCarry is one build a gated send carried.
type WalkCarry struct {
Module string `json:"module"`
Node string `json:"node"`
From string `json:"from,omitempty"`
To string `json:"to"`
Build string `json:"build,omitempty"`
}
// WalkRelease is a backlog walk's way through the machines.
type WalkRelease struct {
Order []string `json:"order"`
Next int `json:"next"`
Gate *WalkGate `json:"gate,omitempty"`
Done []string `json:"done,omitempty"`
By string `json:"by,omitempty"`
}
// WalkDelivery is what a walk waits for.
type WalkDelivery struct {
Awaits string `json:"awaits"`
Go *time.Time `json:"go,omitempty"`
By string `json:"by,omitempty"`
Why string `json:"why,omitempty"`
Stopped string `json:"stopped,omitempty"`
StoppedWhy string `json:"stopped_why,omitempty"`
}
// WalkSeen is what a delivery keeps of its walk.
type WalkSeen struct {
ID string `json:"id"`
Commit string `json:"commit"`
State string `json:"state"`
Note string `json:"note,omitempty"`
Revision int64 `json:"revision"`
Waits bool `json:"waits,omitempty"`
WaitedFor string `json:"waited_for,omitempty"`
GoBy string `json:"go_by,omitempty"`
StoppedBy string `json:"stopped_by,omitempty"`
Tier int `json:"tier"`
Tiers int `json:"tiers"`
Updated time.Time `json:"updated"`
}
// Seen is a walk as a delivery keeps it.
func (w Walk) Seen() *WalkSeen {
s := &WalkSeen{ID: w.ID, Commit: w.Commit, State: w.State, Note: w.Note, Revision: w.Revision, Tier: w.Tier,
Tiers: len(w.Tiers), Updated: w.Updated}
if d := w.Delivery; d != nil {
s.WaitedFor = d.Awaits
s.Waits = d.Awaits != "" && d.Go == nil
s.GoBy = d.By
s.StoppedBy = d.Stopped
}
return s
}
// Started is whether the walk is past its wait: it waited for nobody, or was let go — whatever came of it
// since. A walk superseded or stopped while it waited never started.
func (s *WalkSeen) Started() bool { return s != nil && !s.Waits && (s.WaitedFor == "" || s.GoBy != "") }
// Waited is whether it waited for a word at all.
func (s *WalkSeen) Waited() bool { return s != nil && s.WaitedFor != "" }
// LetGoByAPerson is whether a person's `plans go` started it.
func (s *WalkSeen) LetGoByAPerson() bool { return s != nil && strings.HasPrefix(s.GoBy, "a person") }
// Open is whether the controller still works it.
func (w Walk) Open() bool { return w.State == walkBuilding || w.State == walkRolling }
// StepsOf are the machine steps a walk's record says, for one delivery: every module it moved, on its first
// machines from its gate, and the rest once sent.
func StepsOf(w Walk) []Step {
var out []Step
names := make([]string, 0, len(w.Modules))
for n := range w.Modules {
names = append(names, n)
}
sort.Strings(names)
for _, name := range names {
m := w.Modules[name]
if m == nil || m.FirstAt == nil {
if m != nil && m.SentAt != nil && m.FirstAt == nil && m.Why == "" {
out = append(out, Step{Machine: "every machine running it", Module: name, Build: m.Build,
State: StepSent, Since: *m.SentAt})
}
continue
}
gate := m.Gate
state, why, since := StepSent, "", *m.FirstAt
if gate != nil {
switch {
case gate.Rollback == "rolled-back":
state, why = StepRolledBack, gate.Why
case gate.Verdict == "failed":
state, why = StepFailed, gate.Why
case gate.Verdict == "passed":
state, why = StepPassed, gate.Why
case gate.Since != nil:
state = StepJudging
why = fmt.Sprintf("%d healthy judging(s) so far", gate.Passes)
}
if gate.JudgedAt != nil {
since = *gate.JudgedAt
}
}
for _, machine := range m.First {
out = append(out, Step{Machine: machine, Module: name, Build: m.Build, State: state, Since: since, Why: why})
}
if m.SentAt != nil && state == StepPassed {
out = append(out, Step{Machine: "the rest running it", Module: name, Build: m.Build, State: StepSent,
Since: *m.SentAt})
}
}
return out
}
// Group is a delivery group: two or more deliveries sharing a head branch name, one level only.
type Group struct {
ID string `json:"id"`
Members []string `json:"members"`
Created time.Time `json:"created"`
Updated time.Time `json:"updated"`
// Order is the members in the order they are delivered, Pairs why, Cycle the members whose order
// contradicts itself; OrderOf the member commits the order was worked out for.
Order []string `json:"order,omitempty"`
Pairs []Pair `json:"pairs,omitempty"`
Cycle []string `json:"cycle,omitempty"`
OrderOf map[string]string `json:"order_of,omitempty"`
// Check is its composed check: asked of the controller for these member commits, and its verdict.
Check *GroupCheck `json:"check,omitempty"`
// Closed is set once a member reached the trunk: a group delivering takes no new member.
Closed bool `json:"closed,omitempty"`
// Said is the derived state last said, so a change of it is said once.
Said string `json:"said,omitempty"`
Owed []Effect `json:"owed,omitempty"`
}
// Pair is one "before" among a group's members.
type Pair struct {
Before string `json:"before"`
After string `json:"after"`
Why string `json:"why"`
}
// GroupCheck is a group's composed check.
type GroupCheck struct {
Asked string `json:"asked"`
At time.Time `json:"at"`
Commits map[string]string `json:"commits"`
Verdict string `json:"verdict,omitempty"`
Summary string `json:"summary,omitempty"`
}
// Commits is each member's commit, by id.
func commitsOf(ds []*Delivery) map[string]string {
out := map[string]string{}
for _, d := range ds {
out[d.ID] = d.Commit
}
return out
}
// sameMembers is whether two member→commit maps are one.
func sameMembers(a, b map[string]string) bool {
if len(a) != len(b) {
return false
}
for k, v := range a {
if b[k] != v {
return false
}
}
return true
}
func mustJSON(v any) json.RawMessage {
raw, _ := json.Marshal(v)
return raw
}
@@ -0,0 +1,316 @@
package main
import (
"encoding/json"
"errors"
"fmt"
"sort"
"strings"
stdio "git.novox.be/novox/mesh-sdk/go"
)
// What mesh-delivery asks of others, and where it keeps what it owns. **It never sends to a machine, writes
// a declaration, a grant or a bus object, and never registers a build** (novox/hq ADR 0239 decision 6): it
// asks the controller, through the controller's seat, and the forge's holder, through its tools.
// ControllerSeat is the seat whose verbs the controller serves.
const ControllerSeat = "mesh-controller"
// Member is one delivery as the controller's planner takes it (mesh-controller orderMember).
type Member struct {
ID string `json:"id"`
Repository string `json:"repository"`
Base string `json:"base,omitempty"`
Head string `json:"head,omitempty"`
Number int `json:"number,omitempty"`
Paths []string `json:"paths,omitempty"`
PathsTruncated bool `json:"paths_truncated,omitempty"`
Removed []string `json:"removed,omitempty"`
ModuleDirs []string `json:"module_dirs,omitempty"`
ModuleDirsSaid bool `json:"module_dirs_said,omitempty"`
CloneURL string `json:"clone_url,omitempty"`
After []string `json:"after,omitempty"`
}
// MemberOf is a delivery as a member.
func MemberOf(d *Delivery) Member {
return Member{ID: d.ID, Repository: d.Repository, Base: d.Base, Head: d.Commit, Number: d.Number, Paths: d.Paths,
PathsTruncated: d.PathsTruncated, Removed: d.Removed, ModuleDirs: d.ModuleDirs, ModuleDirsSaid: d.ModuleDirsSaid,
CloneURL: d.CloneURL, After: d.After}
}
// Order is a group's order as the controller works it out.
type Order struct {
Order []string `json:"order"`
Pairs []Pair `json:"pairs,omitempty"`
Cycle []string `json:"cycle,omitempty"`
}
// Controller is the controller's verbs this owner asks with.
type Controller interface {
// Plan is the delivery plan of a diffset, kept nowhere.
Plan(repository, base, head string, paths, moduleDirs, removed []string) (*DeliveryPlan, error)
// Order is a group's order from the graph and its members' `after:` lines.
Order(members []Member) (Order, error)
// Check asks a group's composed check, or — with no group and one member — that head's own check again.
Check(group string, members []Member) (string, error)
// Deliver lets a waiting walk start.
Deliver(walk, why string) error
// Stop ends a walk, said as stopped by whom.
Stop(walk, why, by string) error
// Walks are the walks the controller keeps — one, given its id — and whether this seat has a holder on record.
Walks(walk string) (bool, []Walk, error)
}
// Forge is what the forge's holder does for a delivery.
type Forge interface {
// Note appends one line to a commit's note under refs/notes/mesh-plan; a line already there is not added.
Note(owner, repo, commit, line string) error
// View keeps the delivery's view current on its pull request: one comment, edited in place.
View(owner, repo string, number int, body string) error
// Status sets one status of a commit, linking to the view.
Status(owner, repo, commit, context, state, description, target string) error
}
// Store keeps deliveries and groups: one key each, in the module's own state (ADR 0201).
type Store interface {
PutDelivery(d *Delivery) error
DeleteDelivery(id string) error
Deliveries() ([]*Delivery, error)
PutGroup(g *Group) error
DeleteGroup(id string) error
Groups() ([]*Group, error)
}
// --- over the runtime -----------------------------------------------------------------------------------
// asking is how a tool is asked: stdio.Ask, or a test's.
type asking func(key string, body any) (json.RawMessage, error)
// seatController asks the controller's seat through the runtime.
type seatController struct{ ask asking }
// answerOf reads a verb's answer whatever wraps it: the controller's {output, ok, answer}, a text the
// runtime handed over, or the protocol's content list.
func answerOf(raw json.RawMessage) (json.RawMessage, string, bool) {
var s string
if json.Unmarshal(raw, &s) == nil {
return answerOf(json.RawMessage(s))
}
var m map[string]json.RawMessage
if json.Unmarshal(raw, &m) != nil {
return raw, string(raw), true
}
if content, has := m["content"]; has {
var items []struct {
Text string `json:"text"`
}
var e struct {
IsError bool `json:"isError"`
}
_ = json.Unmarshal(raw, &e)
if json.Unmarshal(content, &items) == nil && len(items) > 0 {
inner, out, ok := answerOf(json.RawMessage(items[0].Text))
return inner, out, ok && !e.IsError
}
}
var v struct {
Output string `json:"output"`
OK *bool `json:"ok"`
Answer json.RawMessage `json:"answer"`
}
if json.Unmarshal(raw, &v) == nil && v.OK != nil {
return v.Answer, v.Output, *v.OK
}
return raw, string(raw), true
}
func lastLine(s string) string {
lines := strings.Split(strings.TrimSpace(s), "\n")
return strings.TrimSpace(lines[len(lines)-1])
}
func (c seatController) verb(name string, args map[string]any) (json.RawMessage, error) {
raw, err := c.ask("seat:"+ControllerSeat+"."+name, args)
if err != nil {
return nil, err
}
answer, output, ok := answerOf(raw)
if !ok {
return nil, fmt.Errorf("the controller refused %s: %s", name, lastLine(output))
}
return answer, nil
}
func (c seatController) Plan(repository, base, head string, paths, moduleDirs, removed []string) (*DeliveryPlan, error) {
args := map[string]any{"repository": repository, "paths": strings.Join(paths, ",")}
if base != "" {
args["base"] = base
}
if head != "" {
args["head"] = head
}
if len(moduleDirs) > 0 {
args["module-dirs"] = strings.Join(moduleDirs, ",")
}
if len(removed) > 0 {
args["removed"] = strings.Join(removed, ",")
}
raw, err := c.verb("delivery-plan", args)
if err != nil {
return nil, err
}
var a struct {
Plan *DeliveryPlan `json:"plan"`
}
if err := json.Unmarshal(raw, &a); err != nil || a.Plan == nil {
return nil, errors.New("the controller's delivery-plan answered no plan")
}
return a.Plan, nil
}
func (c seatController) Order(members []Member) (Order, error) {
raw, err := c.verb("delivery-order", map[string]any{"members": string(mustJSON(members))})
if err != nil {
return Order{}, err
}
var o Order
if err := json.Unmarshal(raw, &o); err != nil {
return Order{}, fmt.Errorf("the controller's delivery-order is not readable: %w", err)
}
return o, nil
}
func (c seatController) Check(group string, members []Member) (string, error) {
args := map[string]any{"members": string(mustJSON(members))}
if group != "" {
args["group"] = group
}
raw, err := c.verb("delivery-check", args)
if err != nil {
return "", err
}
var a struct {
Asked string `json:"asked"`
Rechecked string `json:"rechecked"`
}
_ = json.Unmarshal(raw, &a)
if a.Asked != "" {
return a.Asked, nil
}
return a.Rechecked, nil
}
func (c seatController) Deliver(walk, why string) error {
_, err := c.verb("deliver", map[string]any{"plan": walk, "why": why})
return err
}
func (c seatController) Stop(walk, why, by string) error {
_, err := c.verb("delivery-stop", map[string]any{"plan": walk, "why": why, "by": by})
return err
}
func (c seatController) Walks(walk string) (bool, []Walk, error) {
args := map[string]any{}
if walk != "" {
args["plan"] = walk
}
raw, err := c.verb("delivery-walks", args)
if err != nil {
return false, nil, err
}
var a struct {
Held bool `json:"held"`
Walks []Walk `json:"walks"`
}
if err := json.Unmarshal(raw, &a); err != nil {
return false, nil, fmt.Errorf("the controller's delivery-walks is not readable: %w", err)
}
return a.Held, a.Walks, nil
}
// toolForge asks the forge's holder through its tools.
type toolForge struct{ ask asking }
func (f toolForge) tool(name string, args map[string]any) error {
raw, err := f.ask("gitea."+name, args)
if err != nil {
return err
}
if _, output, ok := answerOf(raw); !ok {
return fmt.Errorf("the forge refused %s: %s", name, lastLine(output))
}
return nil
}
func (f toolForge) Note(owner, repo, commit, line string) error {
return f.tool("gitea_note_append", map[string]any{"owner": owner, "repo": repo, "sha": commit,
"ref": "mesh-plan", "line": line})
}
func (f toolForge) View(owner, repo string, number int, body string) error {
return f.tool("gitea_delivery_view", map[string]any{"owner": owner, "repo": repo, "number": number, "body": body})
}
func (f toolForge) Status(owner, repo, commit, context, state, description, target string) error {
return f.tool("gitea_commit_status", map[string]any{"owner": owner, "repo": repo, "sha": commit,
"context": context, "state": state, "description": description, "target_url": target})
}
// kvStore is the module's own state: `deliveries` and `groups`, one key each.
type kvStore struct{}
// kvKey is an id as a bucket key: only the characters a key may carry.
func kvKey(id string) string {
var b strings.Builder
for _, r := range id {
switch {
case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9', r == '-', r == '_', r == '=':
b.WriteRune(r)
default:
b.WriteString("_" + fmt.Sprintf("%02x", r) + "_")
}
}
return b.String()
}
func (kvStore) PutDelivery(d *Delivery) error {
_, err := stdio.State("deliveries").Put(kvKey(d.ID), d)
return err
}
func (kvStore) DeleteDelivery(id string) error { return stdio.State("deliveries").Delete(kvKey(id)) }
func (kvStore) PutGroup(g *Group) error {
_, err := stdio.State("groups").Put(kvKey(g.ID), g)
return err
}
func (kvStore) DeleteGroup(id string) error { return stdio.State("groups").Delete(kvKey(id)) }
func readAll[T any](state string) ([]*T, error) {
keys, err := stdio.State(state).Keys()
if err != nil {
return nil, err
}
sort.Strings(keys)
var out []*T
for _, k := range keys {
e, err := stdio.State(state).Get(k)
if err != nil {
return nil, err
}
if e == nil {
continue
}
var v T
if err := json.Unmarshal(e.Value, &v); err != nil {
logf("[mesh-delivery] %s %s cannot be read (%v); left as it is", state, k, err)
continue
}
out = append(out, &v)
}
return out, nil
}
func (kvStore) Deliveries() ([]*Delivery, error) { return readAll[Delivery]("deliveries") }
func (kvStore) Groups() ([]*Group, error) { return readAll[Group]("groups") }
@@ -0,0 +1,391 @@
package main
import (
"fmt"
"strings"
"time"
)
// The state table (novox/hq ADR 0239 decision 2, to-be 47): every transition a delivery may take, each with
// its guard, compiled once. The code moves a delivery only through Apply, which finds the row for the
// delivery's state and the event, runs its guard and refuses anything else, naming what was asked. A test
// walks the table: every row it holds, and every pair it does not.
// State is where a delivery is.
type State string
// The states, in the order a delivery passes through them. None is a delivery not yet made: the rows from
// it are how one comes to exist.
const (
None State = ""
Proposed State = "proposed"
Checked State = "checked"
Ready State = "ready"
Rejected State = "rejected"
Published State = "published"
Delivering State = "delivering"
Held State = "held"
Delivered State = "delivered"
Failed State = "failed"
Superseded State = "superseded"
Stopped State = "stopped"
)
// AllStates are every state a delivery can be in, for the table's own test and its verb.
var AllStates = []State{Proposed, Checked, Ready, Rejected, Published, Delivering, Held, Delivered, Failed,
Superseded, Stopped}
// Final is whether nothing follows a state.
func (s State) Final() bool {
return s == Delivered || s == Failed || s == Superseded || s == Stopped
}
// Event is what moves a delivery.
type Event string
// The events. Observed ones are taken by settle whenever their guard holds; the others are acts — a
// person's, or healer H2's — taken only when asked.
const (
EvAnnounced Event = "announced" // the forge announced a pull request's head
EvAppeared Event = "appeared" // a trunk commit the mesh is delivering, with no pull request known
EvAdopted Event = "adopted" // a walk already running at the switch, or on the controller's own path
EvChecked Event = "checked" // its check's verdict arrived
EvAccepted Event = "accepted" // the verdict passes
EvRefused Event = "refused" // the verdict fails, or the check could not run
EvRecheck Event = "recheck" // a person asked for it again, or its group changed
EvNewHead Event = "new-head" // a newer head of the same pull request
EvClosed Event = "closed" // the pull request closed unmerged
EvMerged Event = "merged" // it reached the trunk
EvMergedUnchecked Event = "merged-unchecked" // it reached the trunk without a passing check
EvGo Event = "go" // its walk started
EvHold Event = "hold" // something holds it for a person
EvRelease Event = "release" // a person's word that it goes on
EvDone Event = "done" // every machine of its deploy plan took it, or nothing was to be sent
EvFailed Event = "failed" // its walk failed: a gate, a build, a machine
EvSuperseded Event = "superseded" // a newer delivery to the same trunk took over its walk
EvStop Event = "stop" // a person stopped it, or its group did
)
// Facts are what a guard reads beyond the delivery itself: the moment, and who asks, for an act.
type Facts struct {
Now time.Time
By string
Why string
}
// Row is one transition the table holds.
type Row struct {
From []State
Event Event
To State
// Guard says, in words, what must be true; Holds checks it, answering why not.
Guard string
Holds func(d *Delivery, f Facts) error
// Act says the row is taken only when asked — a person's or H2's — never by settle.
Act bool
}
func unless(cond bool, why string) error {
if cond {
return nil
}
return fmt.Errorf("%s", why)
}
var notFinal = []State{Proposed, Checked, Ready, Rejected, Published, Delivering, Held}
// Table is every transition there is. The order within one state is the order settle tries them.
var Table = []Row{
{From: []State{None}, Event: EvAnnounced, To: Proposed,
Guard: "a pull request's head the forge announced",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Number > 0 && d.Commit != "", "no pull request's head")
}},
{From: []State{None}, Event: EvAppeared, To: Held,
Guard: "a commit on the trunk, its walk waiting, and no pull request known: it waits for a person",
Holds: func(d *Delivery, f Facts) error {
return unless(d.MergedAs != "" && d.HeldWhy != "", "no trunk commit, or nothing says why it is held")
}},
{From: []State{None}, Event: EvAdopted, To: Delivering,
Guard: "a walk already running: at the switch, or on the controller's own path",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Walk != nil && d.Walk.Started(), "no running walk")
}},
{From: []State{Proposed, Checked, Ready, Rejected}, Event: EvNewHead, To: Superseded,
Guard: "a newer head of the same pull request",
Holds: func(d *Delivery, f Facts) error { return unless(d.NewerHead != "", "no newer head") }},
{From: []State{Proposed, Checked, Ready, Rejected}, Event: EvClosed, To: Stopped,
Guard: "the pull request closed unmerged",
Holds: func(d *Delivery, f Facts) error { return unless(d.ClosedUnmerged, "the pull request is open") }},
{From: []State{Ready}, Event: EvMerged, To: Published,
Guard: "on the trunk its modules follow: merged there, and its walk opened — or nothing for a walk to move",
Holds: func(d *Delivery, f Facts) error {
return unless(d.MergedAs != "" && (d.Walk != nil || d.MovesNothing()), "not on the trunk yet")
}},
{From: []State{Proposed}, Event: EvChecked, To: Checked,
Guard: "the verdict names this commit",
Holds: func(d *Delivery, f Facts) error { return unless(d.Check != nil, "no verdict") }},
{From: []State{Checked}, Event: EvAccepted, To: Ready,
Guard: "the gate passed or warned, and the repository's own check did not fail",
Holds: func(d *Delivery, f Facts) error { return unless(d.Check.Passes(), "the verdict does not pass") }},
{From: []State{Checked}, Event: EvRefused, To: Rejected,
Guard: "the gate failed or could not run, or the repository's own check failed",
Holds: func(d *Delivery, f Facts) error { return unless(!d.Check.Passes(), "the verdict passes") }},
{From: []State{Proposed, Checked, Rejected}, Event: EvMergedUnchecked, To: Held,
Guard: "merged without a passing check — a verdict known with it is taken first: only a person decides that it goes on",
Holds: func(d *Delivery, f Facts) error { return unless(d.MergedAs != "", "not merged") }},
{From: []State{Rejected, Ready}, Event: EvRecheck, To: Proposed, Act: true,
Guard: "a person asked, with why, or its group changed",
Holds: func(d *Delivery, f Facts) error { return unless(f.Why != "", "a recheck says why") }},
{From: []State{Published}, Event: EvSuperseded, To: Superseded,
Guard: "a newer delivery to the same trunk took over its walk",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Walk != nil && d.Walk.State == walkSuperseded, "its walk is not superseded")
}},
{From: []State{Published}, Event: EvDone, To: Delivered,
Guard: "nothing for a walk to move",
Holds: func(d *Delivery, f Facts) error { return unless(d.Walk == nil && d.MovesNothing(), "a walk moves it") }},
{From: []State{Published, Ready}, Event: EvHold, To: Held,
Guard: "something holds it for a person: its group's composed check did not pass for the heads that merged, " +
"or it merged and no walk was opened for it",
Holds: func(d *Delivery, f Facts) error {
return unless(d.HeldWhy != "" && d.MergedAs != "", "nothing holds it, or it is not on the trunk")
}},
{From: []State{Published}, Event: EvGo, To: Delivering,
Guard: "its walk started: let go by this owner in its group's order, by a person, or on the controller's own path",
Holds: func(d *Delivery, f Facts) error { return unless(d.Walk != nil && d.Walk.Started(), "its walk waits") }},
{From: []State{Held}, Event: EvGo, To: Delivering,
Guard: "its walk started on a word that was not this owner's: the controller's own path, which waits for " +
"nobody, or a person's `plans go`",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Walk != nil && d.Walk.Started() && (!d.Walk.Waited() || d.Walk.LetGoByAPerson()),
"its walk was not started by a person or the controller's own path")
}},
{From: []State{Held}, Event: EvRelease, To: Delivering, Act: true,
Guard: "a person's decision, with why",
Holds: func(d *Delivery, f Facts) error {
return unless(f.Why != "" && f.By != "", "a release is a person's word, with why")
}},
{From: []State{Delivering}, Event: EvSuperseded, To: Superseded,
Guard: "a newer delivery to the same trunk took over its walk",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Walk != nil && d.Walk.State == walkSuperseded, "its walk is not superseded")
}},
{From: []State{Delivering}, Event: EvDone, To: Delivered,
Guard: "its walk is done: every machine of its deploy plan passed, or was left as its policy says",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Walk != nil && d.Walk.State == walkDone, "its walk is not done")
}},
{From: []State{Delivering}, Event: EvStop, To: Stopped,
Guard: "its walk was stopped through this owner",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Walk != nil && d.Walk.State == walkFailed && d.Walk.StoppedBy != "", "its walk was not stopped")
}},
{From: []State{Delivering}, Event: EvFailed, To: Failed,
Guard: "its walk failed: a gate on a first machine (what it carried put back), a build, a machine",
Holds: func(d *Delivery, f Facts) error {
return unless(d.Walk != nil && d.Walk.State == walkFailed && d.Walk.StoppedBy == "", "its walk did not fail")
}},
{From: notFinal, Event: EvStop, To: Stopped, Act: true,
Guard: "a person, with why, or its group stopped by a member before it",
Holds: func(d *Delivery, f Facts) error { return unless(f.Why != "" && f.By != "", "a stop says who and why") }},
}
// rowFor is the row a state and event name; nil when the table holds none.
func rowFor(from State, ev Event, act bool) *Row {
for i := range Table {
r := &Table[i]
if r.Event != ev || r.Act != act {
continue
}
for _, s := range r.From {
if s == from {
return r
}
}
}
return nil
}
// ErrRefused is a transition the table does not hold, or whose guard does not.
type ErrRefused struct {
ID string
From State
Event Event
Why string
}
func (e ErrRefused) Error() string {
from := string(e.From)
if from == "" {
from = "nothing"
}
return fmt.Sprintf("%s: %s from %s is refused — %s", e.ID, e.Event, from, e.Why)
}
// Apply takes one transition, or refuses it. The delivery keeps the transition, bounded.
func Apply(d *Delivery, ev Event, act bool, f Facts) (Transition, error) {
r := rowFor(d.State, ev, act)
if r == nil {
kind := "observed"
if act {
kind = "asked"
}
return Transition{}, ErrRefused{ID: d.ID, From: d.State, Event: ev,
Why: "the table holds no such transition (" + kind + ")"}
}
if err := r.Holds(d, f); err != nil {
return Transition{}, ErrRefused{ID: d.ID, From: d.State, Event: ev, Why: "its guard does not hold: " + err.Error()}
}
why := f.Why
if why == "" {
why = r.Guard
}
t := Transition{At: f.Now, From: d.State, To: r.To, Event: ev, Why: why, By: f.By}
d.State, d.Since = r.To, f.Now
d.Transitions = append(d.Transitions, t)
if len(d.Transitions) > keptTransitions {
d.Transitions = d.Transitions[len(d.Transitions)-keptTransitions:]
}
return t, nil
}
// keptTransitions is how many transitions a delivery keeps; its note on the commit keeps every one.
const keptTransitions = 100
// settle takes every observed transition whose guard holds, in the table's order, until none does — so a
// delivery always stands where its facts put it, whatever order they arrived in.
func settle(d *Delivery, f Facts) []Transition {
var taken []Transition
for range len(Table) + 1 {
moved := false
for i := range Table {
r := &Table[i]
if r.Act || !contains(r.From, d.State) {
continue
}
if t, err := Apply(d, r.Event, false, f); err == nil {
taken = append(taken, t)
moved = true
break
}
}
if !moved {
break
}
}
return taken
}
func contains(states []State, s State) bool {
for _, x := range states {
if x == s {
return true
}
}
return false
}
// Bounds are how long a delivery may be in a state before `stalled` lists it, and what healer H2 may do
// then — always a transition the table holds, never one of its own (novox/hq ADR 0239 decision 9).
type Bound struct {
For time.Duration
// H2 is the transitions H2 may take after the bound, by event; empty when the state is the operator's.
H2 []Event
Says string
}
// Bounds are every state's.
var Bounds = map[State]Bound{
Proposed: {For: time.Hour, Says: "the check's own watchdog (S6) speaks for a check that is late"},
Checked: {For: time.Minute, Says: "a verdict is decided at once"},
Published: {For: 30 * time.Minute, H2: []Event{EvSuperseded, EvGo}, Says: "its walk waits for its turn or its word"},
Held: {For: 24 * time.Hour, Says: "it waits for the operator"},
Delivering: {For: 2 * time.Hour, H2: []Event{EvDone, EvSuperseded, EvFailed}, Says: "its walk runs"},
}
// StepState is where one machine stands in a delivering walk.
type StepState string
// The machine steps (ADR 0239 decision 2): sent, judged, passed or failed, and put back.
const (
StepSent StepState = "sent"
StepJudging StepState = "judging"
StepPassed StepState = "passed"
StepFailed StepState = "failed"
StepRolledBack StepState = "rolled-back"
)
// StepTable is every move a machine's step may make. From "" is a step first heard of.
var StepTable = map[StepState][]StepState{
"": {StepSent, StepJudging, StepPassed, StepFailed, StepRolledBack},
StepSent: {StepJudging, StepPassed, StepFailed},
StepJudging: {StepPassed, StepFailed},
StepFailed: {StepRolledBack},
}
// StepMay is whether a step may move from one state to another: directly, or through states the walk's
// record passed between two readings of it (sent, judged and failed and put back, read once).
func StepMay(from, to StepState) bool {
if from == to {
return true
}
seen := map[StepState]bool{}
next := []StepState{from}
for len(next) > 0 {
s := next[0]
next = next[1:]
for _, t := range StepTable[s] {
if t == to {
return true
}
if !seen[t] {
seen[t] = true
next = append(next, t)
}
}
}
return false
}
// TableText is the table as the `table` verb answers it.
func TableText() map[string]any {
var rows []map[string]any
for _, r := range Table {
from := make([]string, 0, len(r.From))
for _, s := range r.From {
if s == None {
from = append(from, "(none)")
continue
}
from = append(from, string(s))
}
kind := "observed"
if r.Act {
kind = "asked"
}
rows = append(rows, map[string]any{"from": strings.Join(from, ", "), "event": r.Event, "to": r.To,
"guard": r.Guard, "taken": kind})
}
bounds := map[string]any{}
for s, b := range Bounds {
var h2 []string
for _, e := range b.H2 {
h2 = append(h2, string(e))
}
bounds[string(s)] = map[string]any{"bound": b.For.String(), "h2": h2, "says": b.Says}
}
steps := map[string][]StepState{}
for from, to := range StepTable {
name := string(from)
if name == "" {
name = "(first heard)"
}
steps[name] = to
}
return map[string]any{"transitions": rows, "bounds": bounds, "machine-steps": steps}
}
@@ -0,0 +1,204 @@
package main
import (
"errors"
"testing"
"time"
)
// The table walked (novox/hq ADR 0239 decision 2): every row it holds can be taken when its guard holds and
// is refused when it does not; every pair it does not hold is refused, by name; and the rules the table
// exists for are true of the table itself, not of the code around it.
// aDeliveryFor is a delivery whose facts make the row's guard hold.
func aDeliveryFor(r Row) (*Delivery, Facts) {
now := time.Date(2026, 10, 6, 12, 0, 0, 0, time.UTC)
d := &Delivery{ID: "novox/app@c0ffee000000", Repository: "novox/app", Commit: "c0ffee000000", Number: 7}
f := Facts{Now: now}
switch r.Event {
case EvAppeared:
d.MergedAs, d.HeldWhy = "c0ffee000000", "no pull request"
case EvAdopted, EvGo:
d.Walk = &WalkSeen{ID: "plan-1", State: walkRolling}
if r.From[0] == Held {
d.Walk.WaitedFor, d.Walk.GoBy = "mesh-delivery", "a person (jochen)"
}
case EvNewHead:
d.NewerHead = "beef"
case EvClosed:
d.ClosedUnmerged = true
case EvMerged:
d.MergedAs, d.Walk = "aaaa1111", &WalkSeen{ID: "plan-1", State: walkBuilding, Waits: true, WaitedFor: "mesh-delivery"}
case EvMergedUnchecked:
d.MergedAs = "aaaa1111"
case EvChecked, EvAccepted:
d.Check = &Verdict{Gate: "pass"}
case EvRefused:
d.Check = &Verdict{Gate: "fail"}
case EvRecheck, EvRelease, EvStop:
f.Why, f.By = "a person's reason", "jochen"
if r.Event == EvStop && !r.Act {
f = Facts{Now: now}
d.Walk = &WalkSeen{ID: "plan-1", State: walkFailed, StoppedBy: "mesh-delivery for jochen"}
}
case EvHold:
d.HeldWhy, d.MergedAs = "its group's check did not pass", "aaaa1111"
case EvDone:
if r.From[0] == Published {
d.Plan = &DeliveryPlan{Summary: "moves nothing"}
} else {
d.Walk = &WalkSeen{ID: "plan-1", State: walkDone}
}
case EvFailed:
d.Walk = &WalkSeen{ID: "plan-1", State: walkFailed}
case EvSuperseded:
d.Walk = &WalkSeen{ID: "plan-1", State: walkSuperseded}
}
return d, f
}
func TestEveryRowOfTheTableIsTakenWhenItsGuardHoldsAndRefusedWhenNot(t *testing.T) {
for _, r := range Table {
for _, from := range r.From {
d, f := aDeliveryFor(r)
d.State = from
got, err := Apply(d, r.Event, r.Act, f)
if err != nil {
t.Errorf("%s —%s→ %s was refused with its guard holding: %v", orNothing(from), r.Event, r.To, err)
continue
}
if d.State != r.To || got.From != from || got.To != r.To || len(d.Transitions) != 1 {
t.Errorf("%s —%s→ left it %s, transition %+v", orNothing(from), r.Event, d.State, got)
}
// The guard not holding: a delivery with no facts.
bare := &Delivery{ID: "novox/app@bare", State: from, Check: &Verdict{Gate: "pass"}}
if r.Event == EvAccepted || r.Event == EvChecked {
bare.Check = nil
if r.Event == EvAccepted {
bare.Check = &Verdict{Gate: "fail"}
}
}
if r.Event == EvRefused {
bare.Check = &Verdict{Gate: "pass"}
}
if _, err := Apply(bare, r.Event, r.Act, Facts{Now: f.Now}); err == nil {
t.Errorf("%s —%s→ %s was taken with no fact for its guard", orNothing(from), r.Event, r.To)
} else if bare.State != from || len(bare.Transitions) != 0 {
t.Errorf("a refused transition moved the delivery to %s", bare.State)
}
}
}
}
func TestEveryPairTheTableDoesNotHoldIsRefusedByName(t *testing.T) {
events := []Event{EvAnnounced, EvAppeared, EvAdopted, EvChecked, EvAccepted, EvRefused, EvRecheck, EvNewHead,
EvClosed, EvMerged, EvMergedUnchecked, EvGo, EvHold, EvRelease, EvDone, EvFailed, EvSuperseded, EvStop}
for _, from := range append([]State{None}, AllStates...) {
for _, ev := range events {
for _, act := range []bool{false, true} {
if rowFor(from, ev, act) != nil {
continue
}
d := &Delivery{ID: "novox/app@c0ffee", State: from, MergedAs: "x", NewerHead: "y", Check: &Verdict{Gate: "pass"},
Walk: &WalkSeen{ID: "p", State: walkDone}, HeldWhy: "h", ClosedUnmerged: true, Number: 1, Commit: "c"}
_, err := Apply(d, ev, act, Facts{Now: time.Now(), Why: "w", By: "b"})
var refused ErrRefused
if !errors.As(err, &refused) || refused.Event != ev || refused.From != from {
t.Errorf("%s —%s (act %v)→ was not refused by name: %v", orNothing(from), ev, act, err)
}
if d.State != from {
t.Errorf("a pair the table does not hold moved %s to %s", orNothing(from), d.State)
}
}
}
}
}
// The rules the table is for, read off the table.
func TestTheTableKeepsItsRules(t *testing.T) {
for _, r := range Table {
for _, from := range r.From {
// Nothing leaves a final state.
if from.Final() {
t.Errorf("%s is final and the table leaves it by %s", from, r.Event)
}
// Nothing goes back before the trunk once a delivery is on it.
if (from == Published || from == Delivering || from == Held) &&
(r.To == Proposed || r.To == Checked || r.To == Ready || r.To == Rejected) {
t.Errorf("%s goes back to %s by %s", from, r.To, r.Event)
}
}
// The trunk rule: only a ready delivery is published, by its merge; nothing off the trunk delivers.
if r.To == Published && (r.Event != EvMerged || len(r.From) != 1 || r.From[0] != Ready) {
t.Errorf("%v —%s→ published: only a ready delivery's merge publishes", r.From, r.Event)
}
if r.To == Delivering {
for _, from := range r.From {
if from == Proposed || from == Checked || from == Ready || from == Rejected {
t.Errorf("%s delivers by %s without reaching the trunk", from, r.Event)
}
}
}
// Held is left only by a person or a walk that started on a word that was not this owner's.
if contains(r.From, Held) && r.To == Delivering && r.Event == EvRelease && !r.Act {
t.Error("a held delivery is released by observation")
}
}
// Every state that is not final has a bound or is the operator's to wait in, and every H2 transition is
// a row the table holds from that state.
for state, b := range Bounds {
for _, ev := range b.H2 {
if rowFor(state, ev, false) == nil {
t.Errorf("H2 may take %s from %s, which the table does not hold", ev, state)
}
}
}
// A ready delivery off the trunk is refused publication, whatever else is true of it.
d := &Delivery{ID: "novox/app@c0", State: Ready, Walk: &WalkSeen{ID: "plan-x", State: walkRolling}}
if _, err := Apply(d, EvMerged, false, Facts{Now: time.Now()}); err == nil {
t.Fatal("a commit not on the trunk was published")
}
}
func TestTheMachineStepTable(t *testing.T) {
for _, c := range []struct {
from, to StepState
may bool
}{
{"", StepSent, true}, {StepSent, StepJudging, true}, {StepJudging, StepPassed, true},
{StepJudging, StepFailed, true}, {StepFailed, StepRolledBack, true}, {StepSent, StepFailed, true},
{StepPassed, StepFailed, false}, {StepRolledBack, StepPassed, false}, {StepPassed, StepSent, false},
{StepJudging, StepSent, false}, {StepFailed, StepPassed, false},
} {
if StepMay(c.from, c.to) != c.may {
t.Errorf("%q → %q: may %v, wanted %v", c.from, c.to, !c.may, c.may)
}
}
}
// settle stands a delivery where its facts put it, whatever order they came in.
func TestSettleIsTheOrderOfTheFactsNotOfTheirArrival(t *testing.T) {
now := time.Now()
d := &Delivery{ID: "novox/app@c0", Number: 1, Commit: "c0"}
if _, err := Apply(d, EvAnnounced, false, Facts{Now: now}); err != nil {
t.Fatal(err)
}
// The verdict, the merge and a finished walk, all known before anything settled.
d.Check = &Verdict{Gate: "pass"}
d.MergedAs = "m1"
d.Walk = &WalkSeen{ID: "plan-1", State: walkDone}
ts := settle(d, Facts{Now: now})
var path []State
for _, t := range ts {
path = append(path, t.To)
}
want := []State{Checked, Ready, Published, Delivering, Delivered}
if len(path) != len(want) {
t.Fatalf("settled through %v, wanted %v", path, want)
}
for i := range want {
if path[i] != want[i] {
t.Fatalf("settled through %v, wanted %v", path, want)
}
}
}
@@ -0,0 +1,381 @@
package main
import (
"errors"
"fmt"
"sort"
"strings"
"time"
)
// The seat's verbs (novox/hq ADR 0239, the mesh-delivery seat in the controller's set): five that read, and
// the acts — a person's recheck, release and stop, and healer H2's close, each only by a transition the table
// holds.
// Line is one delivery as `deliveries` lists it.
type Line struct {
ID string `json:"id"`
State State `json:"state"`
Since string `json:"since"`
For string `json:"for"`
Group string `json:"group,omitempty"`
Number int `json:"number,omitempty"`
Waits string `json:"waits,omitempty"`
Plan string `json:"plan,omitempty"`
Walk string `json:"walk,omitempty"`
Owed int `json:"owed,omitempty"`
Stalled bool `json:"stalled,omitempty"`
}
// waitsFor is what a delivery waits for, in words.
func (h *Holder) waitsFor(d *Delivery) string {
switch d.State {
case Proposed:
return "its check"
case Ready:
return "its merge"
case Rejected:
return "a new head, or a recheck"
case Published:
if d.Group != "" {
return "its turn in group " + d.Group
}
return "its walk's start"
case Held:
return "a person: " + d.HeldWhy
case Delivering:
if d.Walk != nil {
return "its walk: " + clip(d.Walk.Note, 120)
}
}
return ""
}
// Deliveries is the `deliveries` verb.
func (h *Holder) Deliveries(state, repository, group string, all bool) []Line {
h.mu.Lock()
defer h.mu.Unlock()
now := h.Now()
var out []Line
for _, d := range h.deliveries {
if state != "" && string(d.State) != state || repository != "" && !strings.EqualFold(d.Repository, repository) ||
group != "" && d.Group != group {
continue
}
if d.State.Final() && !all && now.Sub(d.Since) > 24*time.Hour {
continue
}
l := Line{ID: d.ID, State: d.State, Since: d.Since.UTC().Format(time.RFC3339), For: now.Sub(d.Since).Round(time.Second).String(),
Group: d.Group, Number: d.Number, Waits: h.waitsFor(d), Owed: len(d.Owed)}
if d.Plan != nil {
l.Plan = d.Plan.Summary
}
if d.Walk != nil {
l.Walk = d.Walk.ID
}
if b, bounded := Bounds[d.State]; bounded && !d.State.Final() && now.Sub(d.Since) > b.For {
l.Stalled = true
}
out = append(out, l)
}
sort.Slice(out, func(i, j int) bool {
if out[i].State.Final() != out[j].State.Final() {
return !out[i].State.Final()
}
return out[i].Since > out[j].Since
})
return out
}
// Show is the `show` verb: a delivery whole, or a group with its members.
func (h *Holder) Show(id string) (any, error) {
h.mu.Lock()
defer h.mu.Unlock()
if d := h.deliveries[id]; d != nil {
out := map[string]any{"delivery": d, "waits": h.waitsFor(d)}
if g := h.groups[d.Group]; g != nil {
state, why := GroupState(g, h.membersOf(g))
out["group"] = map[string]any{"group": g, "state": state, "why": why}
}
return out, nil
}
if g := h.groups[id]; g != nil {
members := h.membersOf(g)
state, why := GroupState(g, members)
return map[string]any{"group": g, "state": state, "why": why, "members": members}, nil
}
return nil, fmt.Errorf("no delivery or group %q: `deliveries` lists them", id)
}
// GroupLine is one group as `groups` lists it.
type GroupLine struct {
ID string `json:"id"`
State string `json:"state"`
Why string `json:"why"`
Order []string `json:"order"`
Pairs []Pair `json:"pairs,omitempty"`
Cycle []string `json:"cycle,omitempty"`
Check string `json:"composed,omitempty"`
Members []string `json:"members"`
}
// Groups is the `groups` verb.
func (h *Holder) Groups(all bool) []GroupLine {
h.mu.Lock()
defer h.mu.Unlock()
var out []GroupLine
for _, g := range h.groups {
members := h.membersOf(g)
if len(members) == 0 {
continue
}
state, why := GroupState(g, members)
if !all && (state == "delivered" || state == "failed" || state == "stopped") {
continue
}
l := GroupLine{ID: g.ID, State: state, Why: why, Order: g.Order, Pairs: g.Pairs, Cycle: g.Cycle}
for _, d := range members {
l.Members = append(l.Members, fmt.Sprintf("%s (%s)", d.ID, d.State))
}
if g.Check != nil {
l.Check = g.Check.Verdict
if l.Check == "" {
l.Check = "asked as " + g.Check.Asked
}
}
out = append(out, l)
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
return out
}
// StalledLine is one delivery past its state's bound.
type StalledLine struct {
ID string `json:"id"`
State State `json:"state"`
For string `json:"for"`
Bound string `json:"bound"`
H2 string `json:"h2"`
Says string `json:"says"`
}
// Stalled is the `stalled` verb: what the controller's self-check raises as `delivery.<id>.stalled`.
func (h *Holder) Stalled() []StalledLine {
h.mu.Lock()
defer h.mu.Unlock()
now := h.Now()
var out []StalledLine
for _, d := range h.deliveries {
b, bounded := Bounds[d.State]
if !bounded || d.State.Final() || now.Sub(d.Since) <= b.For {
continue
}
h2 := "none: the state is the operator's"
if len(b.H2) > 0 {
var evs []string
for _, e := range b.H2 {
evs = append(evs, string(e))
}
h2 = "close, by " + strings.Join(evs, " or ") + ", when the walk's record says so"
}
out = append(out, StalledLine{ID: d.ID, State: d.State, For: now.Sub(d.Since).Round(time.Second).String(),
Bound: b.For.String(), H2: h2, Says: b.Says})
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
return out
}
// Recheck is the `recheck` verb: a rejected or ready delivery proposed again, and its check asked again.
func (h *Holder) Recheck(id, why, by string) (string, error) {
h.mu.Lock()
d := h.deliveries[id]
if d == nil {
h.mu.Unlock()
return "", fmt.Errorf("no delivery %q", id)
}
t, err := Apply(d, EvRecheck, true, Facts{Now: h.Now(), Why: why, By: by})
if err != nil {
h.mu.Unlock()
return "", err
}
d.Check = nil
h.moved(d, []Transition{t})
m := MemberOf(d)
h.mu.Unlock()
asked, err := h.Controller.Check("", []Member{m})
if err != nil {
return "", fmt.Errorf("%s is proposed again and its check could not be asked yet (%v); the forge's next "+
"announcement asks it", id, err)
}
return fmt.Sprintf("%s is proposed again; its check was asked (%s)", id, asked), nil
}
// Release is the `release` verb: a held delivery goes on, on a person's word.
func (h *Holder) Release(id, why, by string) (string, error) {
h.mu.Lock()
d := h.deliveries[id]
if d == nil {
h.mu.Unlock()
return "", fmt.Errorf("no delivery %q", id)
}
t, err := Apply(d, EvRelease, true, Facts{Now: h.Now(), Why: why, By: by})
if err != nil {
h.mu.Unlock()
return "", err
}
d.HeldWhy = ""
h.moved(d, []Transition{t})
var walk string
if d.Walk != nil && d.Walk.Waits {
walk = d.Walk.ID
}
h.mu.Unlock()
if walk == "" {
return fmt.Sprintf("%s is delivering, on %s's word", id, by), nil
}
if err := h.Controller.Deliver(walk, "released by "+by+": "+why); err != nil {
return "", fmt.Errorf("%s is delivering on %s's word, and its walk %s could not be let go yet (%v): the next "+
"tick asks again", id, by, walk, err)
}
h.mu.Lock()
if d := h.deliveries[id]; d != nil {
d.LetGo, d.LetGoWhy = h.Now(), "released by "+by
if d.Walk != nil && d.Walk.ID == walk {
d.Walk.Waits, d.Walk.GoBy = false, byOwner
}
h.keep(d)
}
h.mu.Unlock()
return fmt.Sprintf("%s is delivering on %s's word; its walk %s is let go", id, by, walk), nil
}
// Stop is the `stop` verb: a delivery stopped, its walk ended through the controller, and in a group every
// member after it stopped too.
func (h *Holder) Stop(id, why, by string) (string, error) {
h.mu.Lock()
d := h.deliveries[id]
if d == nil {
h.mu.Unlock()
return "", fmt.Errorf("no delivery %q", id)
}
t, err := Apply(d, EvStop, true, Facts{Now: h.Now(), Why: why, By: by})
if err != nil {
h.mu.Unlock()
return "", err
}
h.moved(d, []Transition{t})
var walk string
if d.Walk != nil && (d.Walk.State == walkBuilding || d.Walk.State == walkRolling) {
walk = d.Walk.ID
}
var asks []func()
if g := h.groups[d.Group]; g != nil && g.Closed {
asks = h.groupTurns(g, h.membersOf(g))
} else if g != nil {
h.leaveGroup(d)
}
h.mu.Unlock()
for _, a := range asks {
a()
}
if walk == "" {
return fmt.Sprintf("%s is stopped, on %s's word", id, by), nil
}
if err := h.Controller.Stop(walk, why, by); err != nil {
return "", fmt.Errorf("%s is stopped, and its walk %s could not be ended yet: %v — `plans stop %s --why …` "+
"ends it by hand", id, walk, err, walk)
}
return fmt.Sprintf("%s is stopped, on %s's word; its walk %s is ended", id, by, walk), nil
}
// Close is healer H2's verb: for a delivery past its bound, the walk's record is read again and the
// delivery settled; only a transition the table names for that state counts, any other answer is refused.
func (h *Holder) Close(id, why string) (string, error) {
h.mu.Lock()
d := h.deliveries[id]
if d == nil {
h.mu.Unlock()
return "", fmt.Errorf("no delivery %q", id)
}
b, bounded := Bounds[d.State]
if !bounded || len(b.H2) == 0 {
h.mu.Unlock()
return "", fmt.Errorf("%s is %s: H2 has no transition there — the state is the operator's", id, d.State)
}
if h.Now().Sub(d.Since) <= b.For {
h.mu.Unlock()
return "", fmt.Errorf("%s has been %s for %s, within its bound of %s: nothing to close", id, d.State,
h.Now().Sub(d.Since).Round(time.Second), b.For)
}
if d.Walk == nil {
h.mu.Unlock()
return "", fmt.Errorf("%s has no walk on record: H2 has nothing to read", id)
}
walk, was := d.Walk.ID, d.State
h.mu.Unlock()
_, walks, err := h.Controller.Walks(walk)
if err != nil {
return "", fmt.Errorf("the walk %s cannot be read: %w", walk, err)
}
for _, w := range walks {
h.WalkMoved(w)
}
h.mu.Lock()
defer h.mu.Unlock()
d = h.deliveries[id]
if d.State == was {
return "", errors.New(id + " is still " + string(was) + ": its walk's record does not say it moved — the condition is the operator's")
}
return fmt.Sprintf("%s: %s → %s, as its walk's record says (%s)", id, was, d.State, why), nil
}
// WhatIf is the `what-if` verb.
func (h *Holder) WhatIf(repository, base string, paths []string) (*DeliveryPlan, error) {
if base == "" {
base = "main"
}
return h.Controller.Plan(repository, base, "", paths, nil, nil)
}
// Status is the module's own `delivery_status`.
func (h *Holder) Status(listening string) map[string]any {
h.mu.Lock()
defer h.mu.Unlock()
count := map[State]int{}
owed := 0
for _, d := range h.deliveries {
count[d.State]++
owed += len(d.Owed)
}
for _, g := range h.groups {
owed += len(g.Owed)
}
byState := map[string]int{}
for s, n := range count {
byState[string(s)] = n
}
out := map[string]any{"deliveries": byState, "groups": len(h.groups), "owed": owed, "unkept": len(h.dirty),
"walks-unmatched": len(h.unmatched), "events": listening, "started": h.started.Format(time.RFC3339),
"seat-on-record": h.held}
if !h.reconciled.IsZero() {
out["walks-read"] = h.reconciled.Format(time.RFC3339)
}
if h.unreached != "" {
out["controller-unreached"] = h.unreached
}
if len(h.refusals) > 0 {
from := 0
if len(h.refusals) > 10 {
from = len(h.refusals) - 10
}
out["refused"] = h.refusals[from:]
}
if len(h.backlog) > 0 {
var ids []string
for _, w := range h.backlog {
ids = append(ids, w.ID+" "+w.State)
}
out["backlog-walks"] = ids
}
return out
}
+5
View File
@@ -0,0 +1,5 @@
module mesh-delivery
go 1.22
require git.novox.be/novox/mesh-sdk/go v0.1.7
+2
View File
@@ -0,0 +1,2 @@
git.novox.be/novox/mesh-sdk/go v0.1.7 h1:C0sTQmtTiyYH7bnqZb7PusXnqA37gKuT7Nqjn9gG47w=
git.novox.be/novox/mesh-sdk/go v0.1.7/go.mod h1:GFuZUElBZ9A++mxgIKo97aXXo+kV0uJ/UkbhQPPIbrY=
+75
View File
@@ -0,0 +1,75 @@
{
"module": "mesh-delivery",
"version": "1",
"slug": "deliver",
"claims": [
{
"name": "mesh-delivery",
"scope": "mesh",
"serves": [
"deliveries",
"show",
"groups",
"what-if",
"table",
"stalled",
"recheck",
"release",
"stop",
"close"
]
}
],
"consumes": [
"gitea.pull.updated",
"gitea.pull.merged",
"gitea.pull.closed",
"mesh-controller.checked",
"mesh-controller.plan-moved"
],
"emits": [
"transition",
"group"
],
"invokes": [
"seat:mesh-controller.delivery-plan",
"seat:mesh-controller.delivery-order",
"seat:mesh-controller.delivery-check",
"seat:mesh-controller.deliver",
"seat:mesh-controller.delivery-stop",
"seat:mesh-controller.delivery-walks",
"gitea.gitea_note_append",
"gitea.gitea_delivery_view",
"gitea.gitea_commit_status"
],
"state": [
"deliveries",
"groups"
],
"tools": [
"delivery_status"
],
"resources": [
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
}
],
"build": {
"artifacts": [
{
"name": "tools",
"kind": "bundle",
"language": "go",
"system": "arch",
"from": "cmd/mesh-delivery",
"binary": "mesh-delivery",
"loads": [
"mesh-delivery"
]
}
]
}
}