diff --git a/docs/decisions.md b/docs/decisions.md index 50b383b..793031a 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -568,11 +568,13 @@ and replays as before: routing, the join edge and `iteration_max` are unchanged. Only the new `launch-many` step knows about branches: it starts every branch at once, and starts the ones a crash left unstarted. -**The list is a value; branch results are values.** A source's list rides in -its report like any produced value, so the store needs no new channel to count -branches, and a body reaches heavy content through the bag by entry. A join -collects values only, in branch order. A file output of a branch would need a -blob assembled at open; nothing asked for it yet. +**The list is a value; branch results are values or files.** A source's list +rides in its report like any produced value, so the store needs no new channel +to count branches, and a body reaches heavy content through the bag by entry. +A join collects what the branches produced in branch order: a value as it is, +a file as its text, read from the blob store when the join opens. An agent +branch hands back a file far more easily than a value, which has to ride inside +its one-line result marker. **A failed branch fails the region after the rest have finished.** Cancelling siblings loses work that already cost; collecting partial results lets a join diff --git a/docs/entities/assembly-line.md b/docs/entities/assembly-line.md index 0c7872a..c5dbac3 100644 --- a/docs/entities/assembly-line.md +++ b/docs/entities/assembly-line.md @@ -100,8 +100,9 @@ cancel its siblings. An empty list passes through `to` as a success at once. The edges out of `to`, and their `iteration_max`, are the ordinary ones. A need with `collect: ` is filled, at open, with a JSON array of what -the latest round of branches produced under `` (a **value**), in -branch order, however they reported. The fan-out region does not nest, a +the latest round of branches produced under ``, in branch order, +however they reported: a produced value as it is, a produced file as its +text (UTF-8). The fan-out region does not nest, a body belongs to one fan-out, and a body cannot itself fan out. ### Edges @@ -142,7 +143,7 @@ edges: - a fan-out names a real body, lists a value its own station produces, has a success edge to its body, and its body has an edge for `failed`; a body belongs to one fan-out and does not fan out itself; a collected name is a - value that exactly one body produces (two fan-outs whose bodies produce the + name that exactly one body produces (two fan-outs whose bodies produce the same name would mix their branches) - every required need of every node is seeded at start or produced on every path into it; a node with a custom start event must have its required diff --git a/packages/store/src/assembly-run-store-fanout.test.ts b/packages/store/src/assembly-run-store-fanout.test.ts index 1445253..24e4284 100644 --- a/packages/store/src/assembly-run-store-fanout.test.ts +++ b/packages/store/src/assembly-run-store-fanout.test.ts @@ -1,8 +1,9 @@ import { describe, expect, it } from "vitest"; +import { BlobsStore } from "./blobs.js"; import { FAN_STATIONS, fanLine } from "./fanout.fixtures.js"; import { setupStoreFixture } from "./assembly-run-store.fixtures.js"; -const { store, events, definitions } = setupStoreFixture(); +const { store, events, definitions, pool } = setupStoreFixture(); async function runAfterSplit(listed: string[]): Promise { await seedFanLine(); @@ -174,4 +175,30 @@ describe("AssemblyRunStore across a fan-out", () => { merge: [{ runId, nodeId: "merge", iteration: 1 }], }); }); + + it("hands the join the text of each branch's file output, in branch order", async () => { + await definitions().put("line", "fan", fanLine()); + const fileWorker = { ...FAN_STATIONS.worker!, produces: [{ name: "result", kind: "file" as const }] }; + + await definitions().put("station", "splitter", FAN_STATIONS.splitter!); + await definitions().put("station", "worker", fileWorker); + await definitions().put("station", "merger", FAN_STATIONS.merger!); + const { run } = await store().start({ lineId: "fan", repo: null, startItems: {} }); + const { visit: entry } = await store().openVisit(run.id, "split", 1); + + await store().report(entry.id, { outcome: "success", produced: { items: JSON.stringify(["a", "b"]) } }); + const blobs = new BlobsStore({ connection: pool() }); + const branches = await openBranches(run.id, 2); + + for (const [index, visitId] of [...branches.entries()].reverse().map(([index, id]) => [index, id] as const)) { + const { hash } = await blobs.put(Buffer.from(`patch ${index}`)); + + await store().report(visitId, { outcome: "success", produced: { result: hash } }); + } + + const { visit } = await store().openVisit(run.id, "merge", 1); + const { needs } = visit.brief; + + expect(JSON.parse(needs.results!)).toEqual(["patch 0", "patch 1"]); + }); }); diff --git a/packages/store/src/fan-needs.ts b/packages/store/src/fan-needs.ts index a66673e..7484227 100644 --- a/packages/store/src/fan-needs.ts +++ b/packages/store/src/fan-needs.ts @@ -29,20 +29,38 @@ function newestListing(visits: readonly Visit[], sourceId: string, over: string) return listings.filter((listed) => listed !== undefined).at(-1); } -/** Each `collect` need of the station, filled with a JSON array of what the newest round of branches produced under that name, in branch order. */ -export function collectNeeds(node: LineNode, station: StationBody, visits: readonly Visit[]): Record { +/** How a collected output is read back: what kind a station produces a name as, and the text of a file's blob. */ +export interface CollectReaders { + produceKind(stationHash: string | null, name: string): Promise<"value" | "file">; + blobText(hash: string): Promise; +} + +/** Each `collect` need of the station, filled with a JSON array of what the newest round of branches produced under that name, in branch order: a value as it is, a file as its text. */ +export async function collectNeeds(node: LineNode, station: StationBody, visits: readonly Visit[], readers: CollectReaders): Promise> { const collecting = station.needs.filter((need) => need.collect !== undefined); + const filled = await Promise.all( + collecting.map(async (need): Promise<[string, Item]> => { + const texts = await collected(visits, need.collect!, readers); - return Object.fromEntries( - collecting.map((need): [string, Item] => [node.bind?.[need.name] ?? need.name, { kind: "value", ref: JSON.stringify(collected(visits, need.collect!)), by: "fanout" }]), + return [node.bind?.[need.name] ?? need.name, { kind: "value", ref: JSON.stringify(texts), by: "fanout" }]; + }), ); + + return Object.fromEntries(filled); } -function collected(visits: readonly Visit[], name: string): string[] { +async function collected(visits: readonly Visit[], name: string, readers: CollectReaders): Promise { const branches = visits.filter((visit) => visit.branch !== null && producedValue(visit, name) !== undefined); const newest = Math.max(0, ...branches.map((visit) => visit.iteration)); const round = branches.filter((visit) => visit.iteration === newest); const ordered = round.sort((first, second) => first.branch! - second.branch!); - return ordered.map((visit) => producedValue(visit, name)!); + return Promise.all(ordered.map((visit) => textOf(visit, name, readers))); +} + +async function textOf(visit: Visit, name: string, readers: CollectReaders): Promise { + const produced = producedValue(visit, name)!; + const kind = await readers.produceKind(visit.stationHash, name); + + return kind === "file" ? readers.blobText(produced) : produced; } diff --git a/packages/store/src/line-fanout.test.ts b/packages/store/src/line-fanout.test.ts index ef19a70..607c013 100644 --- a/packages/store/src/line-fanout.test.ts +++ b/packages/store/src/line-fanout.test.ts @@ -104,7 +104,7 @@ describe("validateLine on a fan-out", () => { }), }), ).toContain( - 'node "merge" collects "nothing", which no fan-out body produces as a value', + 'node "merge" collects "nothing", which no fan-out body produces', ); }); diff --git a/packages/store/src/line-fanout.ts b/packages/store/src/line-fanout.ts index 98bf575..6c399b1 100644 --- a/packages/store/src/line-fanout.ts +++ b/packages/store/src/line-fanout.ts @@ -123,7 +123,7 @@ function collectProblem( if (count === 0) { return [ - `node "${nodeId}" collects "${collect}", which no fan-out body produces as a value`, + `node "${nodeId}" collects "${collect}", which no fan-out body produces`, ]; } @@ -134,7 +134,7 @@ function collectProblem( : []; } -// How many fan-out bodies produce each value name: a collected name must come from exactly one, or the branches of two fan-outs would mix. +// How many fan-out bodies produce each name: a collected name must come from exactly one, or the branches of two fan-outs would mix. function producersByName( line: LineBody, sources: FanSource[], @@ -144,7 +144,7 @@ function producersByName( const stationNodes = line.nodes.filter(hasStation); const names = stationNodes .filter((node) => bodyIds.has(node.id)) - .flatMap((node) => valuesProducedBy(node, bodies)); + .flatMap((node) => namesProducedBy(node, bodies)); return names.reduce( (counts, name) => counts.set(name, (counts.get(name) ?? 0) + 1), @@ -152,13 +152,11 @@ function producersByName( ); } -function valuesProducedBy( +function namesProducedBy( node: LineNode & { station: string }, bodies: Bodies, ): string[] { const produces = resolvedStationOf(node, bodies)?.produces ?? []; - return produces - .filter((produce) => produce.kind === "value") - .map((produce) => produce.name); + return produces.map((produce) => produce.name); } diff --git a/packages/store/src/open-visit.ts b/packages/store/src/open-visit.ts index 0a48cf9..e205b54 100644 --- a/packages/store/src/open-visit.ts +++ b/packages/store/src/open-visit.ts @@ -1,6 +1,7 @@ // Resolution at open (docs/assembly_run_storage.md, "Resolution at open"): station, agent definition and variant, needs from the bag, conversation, failure context, deadline. -import { collectNeeds, itemNeed } from "./fan-needs.js"; +import { BlobsStore } from "./blobs.js"; +import { collectNeeds, itemNeed, type CollectReaders } from "./fan-needs.js"; import { forRepo } from "./repo-name.js"; import type { Pool } from "pg"; import { DefinitionsStore } from "./definitions.js"; @@ -76,7 +77,7 @@ export class OpenVisitResolver { const station = await this.stationRef(node.station!); const bag = await this.deps.bag(request.run.id); const visits = await visitsWith(this.deps.pool, request.run.id); - const fanned = { ...itemNeed(request.line, node, request.branch, visits), ...collectNeeds(node, station.body, visits) }; + const fanned = { ...itemNeed(request.line, node, request.branch, visits), ...(await collectNeeds(node, station.body, visits, this.collectReaders())) }; const { needs: resolvedNeeds, missing } = resolveNeeds(station.body.needs, node.bind, { ...bag, ...fanned }); if (missing.length > 0) { @@ -113,6 +114,31 @@ export class OpenVisitResolver { }; } + private collectReaders(): CollectReaders { + const blobs = new BlobsStore({ connection: this.deps.pool }); + // Every branch of a fan-out ran the same station, so one open looks it up once. + const stations = new Map>(); + + const stationOf = (hash: string): Promise => { + const known = stations.get(hash) ?? this.deps.definitions.byHashOnly("station", hash).then((row) => row?.body); + + stations.set(hash, known); + + return known; + }; + + return { + produceKind: async (stationHash, name) => producedKind(stationHash ? await stationOf(stationHash) : undefined, name), + blobText: async (hash) => { + const blob = await blobs.get(hash); + + enforce(blob, `a collected file is gone: ${hash}`); + + return blob.bytes.toString("utf8"); + }, + }; + } + private async stationRef(ref: string): Promise { const [id, hash] = ref.split("@"); const row = hash @@ -271,3 +297,9 @@ function dispatchTagsFor(station: ResolvedStation, settings: { body: AgentSettin return []; } +/** What a station produces `name` as; a value when it declares no such produce. */ +function producedKind(station: StationBody | undefined, name: string): "value" | "file" { + const produces = station?.produces ?? []; + + return produces.find((produce) => produce.name === name)?.kind ?? "value"; +}