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
14 changes: 5 additions & 9 deletions packages/server/src/recall/recall-versions.test.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
import { internalTermsIn } from '@beonauto/api/testing';
import { recallDocument } from '@beonauto/recollection/testing';
import { Schema } from 'effect';
import { afterEach, describe, expect, it, vi } from 'vitest';
import { afterEach, describe, expect, it } from 'vitest';

import { alpha, type ReasoningServer } from '../testing/servers/reasoning-server.ts';
import {
brainWithReviews,
foldsHeld,
foldsHeldUntilOpened,
inState,
liveAt,
Expand Down Expand Up @@ -61,15 +62,13 @@ function updated(server: ReasoningServer, name: string, fold: string) {
}

describe('a recall function whose view is being built', { timeout: recallTestTimeoutMs }, () => {
it('answers at once, unavailable, rebuilding, with Retry-After, until its view has caught up', async () => {
it('answers before its view is built, unavailable, rebuilding, with Retry-After, until its view has caught up', async () => {
const gate = foldsHeldUntilOpened();
const server = await servingWith({}, gate);
await brainWithReviews(server, 1);
await standingUntil(server, 'reviews', inState('rebuilding'));
await foldsHeld(gate, 1);

const askedAt = Date.now();
const meanwhile = await recalled(server, 'reviews', { campaign: 'spring' });
const answeredInMs = Date.now() - askedAt;
gate.open();
await standingUntil(server, 'reviews', liveWith(1));
const caughtUp = await recalled(server, 'reviews', { campaign: 'spring' });
Expand All @@ -86,7 +85,6 @@ describe('a recall function whose view is being built', { timeout: recallTestTim
});
expect(meanwhile.headers.get('retry-after')).toBe('5');
expect(internalTermsIn(decodeDetail(meanwhile.body).detail)).toEqual([]);
expect(answeredInMs).toBeLessThan(2000);
expect(caughtUp).toMatchObject({ status: 200, body: { output: [{ verdict: 'approve' }] } });
});

Expand Down Expand Up @@ -114,9 +112,7 @@ describe('a recall function waiting for its view to be built', { timeout: recall
const gate = foldsHeldUntilOpened();
const server = await servingWith({ RECOLLECTION_MAX_REBUILDS: '1', RECOLLECTION_BRAINS_AT_ONCE: '1' }, gate);
await brainWithReviews(server, 1, 'beta');
await vi.waitFor(() => {
expect(gate.held()).toBe(1);
});
await foldsHeld(gate, 1);
await brainWithReviews(server, 1);
await saved(server, 'count', '. + 1');
gate.letThrough(1);
Expand Down
23 changes: 14 additions & 9 deletions packages/server/src/testing/servers/recall-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@ import { alpha, servingReasoning, type ReasoningServer } from './reasoning-serve

export const recallTestTimeoutMs = 60_000;

const untilNearTheTestTimeout = { timeout: recallTestTimeoutMs - 10_000, interval: 50 };

export const anyReview = [
'---',
'description: Reviews a campaign brief, answering whatever the model writes',
Expand Down Expand Up @@ -73,6 +75,12 @@ export function foldsHeldUntilOpened(): Gate {
};
}

export function foldsHeld(gate: Gate, folds: number): Promise<void> {
return vi.waitFor(() => {
expect(gate.held()).toBe(folds);
}, untilNearTheTestTimeout);
}

export async function brainWithReviews(server: ReasoningServer, briefsFirst = 0, brain = 'alpha'): Promise<void> {
const path = `/v1/orgs/acme/brains/${brain}`;
await server.call('POST', '/v1/orgs/acme/brains', { body: { brain, name: brain } });
Expand Down Expand Up @@ -102,15 +110,12 @@ export function standingUntil(
name: string,
holds: (standing: Standing) => boolean,
): Promise<unknown> {
return vi.waitFor(
async () => {
const { body } = await server.call('GET', `${alpha}/specs/recollection/${name}`);
const standing = Option.getOrUndefined(decodeStanding(body))?.standing;
expect(standing !== undefined && holds(standing)).toBe(true);
return body;
},
{ timeout: recallTestTimeoutMs - 10_000, interval: 50 },
);
return vi.waitFor(async () => {
const { body } = await server.call('GET', `${alpha}/specs/recollection/${name}`);
const standing = Option.getOrUndefined(decodeStanding(body))?.standing;
expect(standing !== undefined && holds(standing)).toBe(true);
return body;
}, untilNearTheTestTimeout);
}

export function liveWith(folded: number): (standing: Standing) => boolean {
Expand Down
2 changes: 1 addition & 1 deletion packages/workflow-host/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -184,7 +184,7 @@ The projector reads the rows of a brain without their views, and reads a view on

The projector (`src/projector`) runs while the host holds the database's workflows, woken by its sweep and by the ledger's `AppendSignal`, raised in the process on every append. Each sweep it asks the ledger for the definition streams of the recall type, `EventStore.definitionStreams`, through the ledger's partial index on their names, and takes every brain whose stream moved or that has rows; an append to a brain it does not know asks for them at once, so a brain's first recall function, and a brain made after the start, are found without a restart. A brain is passed in a fiber of its own, at most `brainsAtOnce` at once. A pass of a brain:

1. reads the brain's definition stream from where it last read it and brings the rows in line: a function without a row gets one, a newer version resets its row, a retired function's row is dropped, and the first `rebuildsAtOnce` rows that wait or rebuild, in the order saved, rebuild while the others wait; a stalled row holds no slot;
1. reads the brain's definition stream from where it last read it and brings the rows in line: a function without a row gets one, a newer version resets its row, a retired function's row is dropped, and the first `rebuildsAtOnce` rows that wait or rebuild, in the order saved, rebuild while the others wait; a stalled row holds no slot. A new or reset row is written in the phase its place gives it, so a view that takes a slot, as a brain's only recall function does, is never seen waiting;
2. reads the brain from the smallest checkpoint of its live views, a page of up to 1,000 records, the most one read of the store examines, of the types the views can match, the fact types their filters name, `event_published` when a filter names any other type, and the three definition types, oldest first behind PostgreSQL's horizon, so an event committed late is folded late, never passed over;
3. folds the page in a worker of the pool, warm when one of its module is idle (`src/pages`), the projector holding at most half the pool's workers and at least one: each view folds the events after its own checkpoint that its filters match, never the runs of its own function, each fold under its own work, deadline and size and the view's schema, and each filter's `data` test under the same limits, with work of its own, sharing the fold's deadline; the worker checks the page's 2 seconds of folding, counted from its first fold, before each fold, and ends the page early, before the next fold, once they are spent, so at most one fold runs past them;
4. writes each view that moved once, its checkpoint alone when it only read past events, with one statement conditional on the version and the checkpoint it read, so a host that lost the claim, or one that read a row a newer version has since reset, writes nothing, and a dropped row is never made again;
Expand Down
95 changes: 66 additions & 29 deletions packages/workflow-host/src/projector/view-reconciling.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,8 @@ import { Effect } from 'effect';

import { rowsOf, type DatabaseFailed, type HostDatabase } from '../database/host-database.ts';
import type { Statement } from '../database/statement.ts';
import { viewDetailsOf } from '../views/view-details.ts';
import type { ViewRow } from '../views/view-rows.ts';
import { viewDetailsOf, type ViewDetails } from '../views/view-details.ts';
import type { ViewPhase, ViewRow } from '../views/view-rows.ts';
import { viewRowOf, ViewRowSchema } from '../views/view-rows.ts';
import { phaseSet, viewAdded, viewDropped, viewRenewed, viewsOfBrain, type NewView } from '../views/view-statements.ts';
import {
Expand Down Expand Up @@ -38,49 +38,86 @@ function definitionsOf({ database, definitionType, definitions }: Reconciling, b
);
}

function newViewOf(brain: string, name: string, kept: KeptFunction): NewView | undefined {
interface Saved {
readonly name: string;
readonly saved: number;
}

interface FreshFunction extends Saved {
readonly version: number;
readonly details: ViewDetails;
}

function isBuilding({ phase }: ViewRow): boolean {
return phase === 'waiting' || phase === 'rebuilding';
}

function slotsOf(building: readonly Saved[], rebuildsAtOnce: number): ReadonlySet<string> {
const inOrderSaved = building.toSorted((one, other) => one.saved - other.saved);
return new Set(inOrderSaved.slice(0, rebuildsAtOnce).map(({ name }) => name));
}

function phaseIn(slots: ReadonlySet<string>, name: string): ViewPhase {
return slots.has(name) ? 'rebuilding' : 'waiting';
}

function freshOf(name: string, kept: KeptFunction, row: ViewRow | undefined): readonly FreshFunction[] {
const details = viewDetailsOf(kept.details);
return details === undefined
? undefined
: {
brain,
name,
version: kept.version,
saved: kept.saved,
details: JSON.stringify(details),
phase: 'waiting',
view: JSON.stringify(details.initial),
};
return details === undefined || (row !== undefined && row.version >= kept.version)
? []
: [{ name, saved: kept.saved, version: kept.version, details }];
}

function newViewOf(brain: string, { name, saved, version, details }: FreshFunction, phase: ViewPhase): NewView {
return {
brain,
name,
version,
saved,
details: JSON.stringify(details),
phase,
view: JSON.stringify(details.initial),
};
}

function changesOf(brain: string, rows: readonly ViewRow[], { functions }: BrainDefinitions): readonly Statement[] {
function changesOf(
brain: string,
rows: readonly ViewRow[],
{ functions }: BrainDefinitions,
rebuildsAtOnce: number,
): readonly Statement[] {
const byName = new Map(rows.map((row) => [row.name, row]));
const kept = [...functions].flatMap(([name, function_]: readonly [string, KeptFunction]): readonly Statement[] => {
const added = newViewOf(brain, name, function_);
const row = byName.get(name);
if (added === undefined || (row !== undefined && row.version >= added.version)) {
return [];
}
return [row === undefined ? viewAdded(added) : viewRenewed(added)];
});
const fresh = [...functions].flatMap(([name, kept]: readonly [string, KeptFunction]) =>
freshOf(name, kept, byName.get(name)),
);
const freshNames = new Set(fresh.map(({ name }) => name));
const stillBuilding = rows.filter((row) => isBuilding(row) && functions.has(row.name) && !freshNames.has(row.name));
const slots = slotsOf([...stillBuilding, ...fresh], rebuildsAtOnce);
const promoted = stillBuilding
.filter((row) => row.phase === 'waiting' && slots.has(row.name))
.map((row) => phaseSet(row, 'rebuilding'));
const dropped = rows.filter(({ name }) => !functions.has(name)).map(({ name }) => viewDropped(brain, name));
return [...kept, ...dropped];
const written = fresh.map((function_) => {
const view = newViewOf(brain, function_, phaseIn(slots, function_.name));
return byName.has(view.name) ? viewRenewed(view) : viewAdded(view);
});
return [...promoted, ...dropped, ...written];
}

function phasesOf(rows: readonly ViewRow[], rebuildsAtOnce: number): readonly ViewRow[] {
const building = rows.filter(({ phase }) => phase === 'waiting' || phase === 'rebuilding');
const slots = new Set(building.slice(0, rebuildsAtOnce).map(({ name }) => name));
return rows.map((row) =>
building.includes(row) ? { ...row, phase: slots.has(row.name) ? 'rebuilding' : 'waiting' } : row,
const slots = slotsOf(
rows.filter((row) => isBuilding(row)),
rebuildsAtOnce,
);
return rows.map((row) => (isBuilding(row) ? { ...row, phase: phaseIn(slots, row.name) } : row));
}

export function reconciled(parts: Reconciling, brain: string): Effect.Effect<readonly ViewRow[], DatabaseFailed> {
const { database } = parts;
return Effect.gen(function* () {
const definitions = yield* definitionsOf(parts, brain);
const before = yield* rowsOfBrain(database, brain);
const changes = changesOf(brain, before, definitions);
const changes = changesOf(brain, before, definitions, parts.rebuildsAtOnce);
yield* Effect.forEach(changes, database.write, { discard: true });
const rows = changes.length === 0 ? before : yield* rowsOfBrain(database, brain);
const phased = phasesOf(rows, parts.rebuildsAtOnce);
Expand Down
52 changes: 35 additions & 17 deletions packages/workflow-host/src/triggers/several-triggers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,19 +3,21 @@ import { describe, expect, it } from 'vitest';

import { statement } from '../database/statement.ts';
import {
alpha,
at,
cronTrigger,
eventRecordOf,
eventTrigger,
everyTrigger,
published,
publishedInTurn,
runRecorded,
specRecordAt,
specRecorded,
} from '../reaction-testing/brain-writes.ts';
import { startsReaching, untilScheduleRuns } from '../reaction-testing/kept-triggers.ts';
import { movedClock } from '../reaction-testing/moved-clock.ts';
import { reactingHost } from '../reaction-testing/reacting-host.ts';
import { reactingHost, type ReactingHost } from '../reaction-testing/reacting-host.ts';
import { until } from '../reaction-testing/until.ts';
import { reactionExecutionIdOf } from '../reactions/reaction-ids.ts';
import { mostStartsAMinute } from '../reactions/start-rates.ts';
Expand All @@ -30,15 +32,33 @@ function isoAt(minutes: number): string {
return new Date(activatedAt + minutes * aMinute).toISOString();
}

async function closingAt(start: number, onEvents = closed) {
async function openedAt(start: number) {
const clock = movedClock(start);
const reacting = await reactingHost({ clock });
await specRecorded(reacting.database.store, {
return { reacting, clock, starts: reacting.reactions.starts };
}

function closingSaved({ database }: ReactingHost, onEvents = closed) {
return specRecorded(database.store, {
name: 'close',
version: 1,
triggers: [onEvents, cronTrigger('30 9 * * *'), everyTrigger(15 * aMinute)],
});
return { reacting, clock, starts: reacting.reactions.starts };
}

async function closingAt(start: number, onEvents = closed) {
const opened = await openedAt(start);
await closingSaved(opened.reacting, onEvents);
return opened;
}

function closingStartedIn({ database }: ReactingHost, minute: number, starts: number) {
return Effect.runPromise(
database.write(
statement`INSERT INTO workflow_reaction_rates (brain_key, workflow, minute, starts)
VALUES (${alpha}, 'close', ${minute}, ${starts})`,
),
);
}

function executionIdOf(start: { readonly executionId: string } | undefined): string {
Expand Down Expand Up @@ -112,26 +132,24 @@ describe('a due time of a schedule asked for again', () => {
});

describe('the start rate of a workflow with several triggers', () => {
it(`counts the ${mostStartsAMinute} starts of its event trigger a minute, and never a start of its schedule`, async () => {
const { reacting } = await closingAt(activatedAt + 15 * aMinute + 1000);
it(`counts the starts of its event trigger to ${mostStartsAMinute} a minute, and never a start of its schedule`, async () => {
const minute = activatedAt + 15 * aMinute;
const { reacting } = await openedAt(minute + 1000);
await closingStartedIn(reacting, minute, mostStartsAMinute - 1);
await closingSaved(reacting);
await startsReaching(reacting.reactions.starts, 1);

await Array.from({ length: mostStartsAMinute + 1 }, (_, index) => index).reduce<Promise<void>>(
(before, index) =>
before.then(() => published(reacting.database.store, { id: `e${index}`, type: 'com.acme.closed' })),
Promise.resolve(),
);
const starts = await startsReaching(reacting.reactions.starts, mostStartsAMinute + 1);
await publishedInTurn(reacting.database.store, [
{ id: 'e1', type: 'com.acme.closed' },
{ id: 'e2', type: 'com.acme.closed' },
]);
const starts = await startsReaching(reacting.reactions.starts, 2);
const waiting = await until(
() => Effect.runPromise(reacting.database.read(statement`SELECT execution_id FROM workflow_reaction_backlog`)),
(rows) => rows.length > 0,
);

expect([starts.filter(({ trigger }) => trigger.kind === 'event').length, waiting.length]).toEqual([
mostStartsAMinute,
1,
]);
expect(starts.filter(({ trigger }) => trigger.kind === 'every')).toHaveLength(1);
expect([starts.map(({ trigger }) => trigger.kind), waiting.length]).toEqual([['every', 'event'], 1]);
});
});

Expand Down
Loading
Loading