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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 7 additions & 5 deletions docs/decisions.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 4 additions & 3 deletions docs/entities/assembly-line.md
Original file line number Diff line number Diff line change
Expand Up @@ -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: <name>` is filled, at open, with a JSON array of what
the latest round of branches produced under `<name>` (a **value**), in
branch order, however they reported. The fan-out region does not nest, a
the latest round of branches produced under `<name>`, 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
Expand Down Expand Up @@ -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
Expand Down
29 changes: 28 additions & 1 deletion packages/store/src/assembly-run-store-fanout.test.ts
Original file line number Diff line number Diff line change
@@ -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<string> {
await seedFanLine();
Expand Down Expand Up @@ -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"]);
});
});
30 changes: 24 additions & 6 deletions packages/store/src/fan-needs.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, Item> {
/** 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<string>;
}

/** 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<Record<string, Item>> {
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<string[]> {
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<string> {
const produced = producedValue(visit, name)!;
const kind = await readers.produceKind(visit.stationHash, name);

return kind === "file" ? readers.blobText(produced) : produced;
}
2 changes: 1 addition & 1 deletion packages/store/src/line-fanout.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
);
});

Expand Down
12 changes: 5 additions & 7 deletions packages/store/src/line-fanout.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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`,
];
}

Expand All @@ -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[],
Expand All @@ -144,21 +144,19 @@ 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),
new Map<string, number>(),
);
}

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);
}
36 changes: 34 additions & 2 deletions packages/store/src/open-visit.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -113,6 +114,31 @@ export class OpenVisitResolver {
};
}

private collectReaders(): CollectReaders {
Comment thread
gedaiu marked this conversation as resolved.
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<string, Promise<StationBody | undefined>>();

const stationOf = (hash: string): Promise<StationBody | undefined> => {
const known = stations.get(hash) ?? this.deps.definitions.byHashOnly<StationBody>("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<ResolvedStation> {
const [id, hash] = ref.split("@");
const row = hash
Expand Down Expand Up @@ -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";
}
Loading