Skip to content

Commit 73b2246

Browse files
fix(service-automation): a flow saved through the metadata API is armed on the running engine at once (#21746)
Fixes #21725 Clause-②: no A flow saved active through `PUT /api/v1/meta/flow/:name` is now armed on the running engine at once. A save that sets `status: 'obsolete'` or `'invalid'` disarms it, and a delete unregisters it. Hooks and actions saved through the same door already behaved this way. The work follows triage's direction (one door, one arming rule): `service-automation` subscribes to the metadata protocol's existing mutation signal, and no new event is added. ## Reproduced first, at `316be321ef` This was measured at the producer behind the REST route (real kernel, ObjectQLPlugin's metadata protocol, driver-sql on better-sqlite3 `:memory:`), before the change: - `saveMetaItem({ type: 'flow', name: 'qa_meta_flow', item })` answered `"Saved flow 'qa_meta_flow' (env-wide, state=active) [seq=1]"`, the card's step 2. - `getMetaItemsForExecution({ type: 'flow' })` listed the flow. - `engine.getFlow('qa_meta_flow')` was `null`, both at once and 50 ms later. - `engine.execute('qa_meta_flow')` returned `Flow 'qa_meta_flow' not found`, the card's step 3. Where each side lives: - **The engine's flow set** is filled only by `registerFlow`, at `service-automation/src/engine.ts:4415` (`this.flows.set`). Its callers are the boot pull (`plugin.ts:1280`), the `kernel:ready` sync (`:2434`) and the `metadata:reloaded` re-sync (`:2383`, hook at `:1371`). - **The direct save's signal** fires at `metadata-protocol/src/protocol.ts:19404` (`saveMetaItem` → `emitMetadataMutation`). ObjectQL's listener (`objectql/src/plugin.ts:328-352`) re-binds hooks and actions on it, and nothing in `service-automation` listened. ## The signal is reachable without an `objectql` edit `onMetadataMutation` is a public method of the `protocol` **kernel service** (`ObjectStackProtocolImplementation`, `@objectstack/metadata-protocol`). It is a server-side extension, not part of the wire contract. `plugin-email`, `plugin-security`, `service-i18n` and core's authored-translation sync already subscribe through `ctx.getService('protocol')`. The `MetadataMutationEvent` type is imported type-only. No file in `objectql`, `spec` or `metadata-protocol` changes. ## The change (`packages/services/service-automation/src/plugin.ts`) - **Subscription.** `subscribeFlowMutations` runs at `start()`, after the inert-mode return, beside the other runtime re-bind hooks, so `armRuntime: false` arms nothing. It handles `type === 'flow'` and ignores `state: 'draft'`, because a draft leaves the live row unchanged. The listener only records the name, and one queued sync drains every name reported before it starts. - **The sync (`syncMutatedFlows`)** re-reads each named flow through the execution view, with the same read, precedence (`resolveFlowContenders`) and env-wide scope (no organization) as the boot: - **In the view:** `registerFlow`. For `status: 'obsolete' | 'invalid'` the flow stays registered and unbound, which is exactly what a boot gives such a flow. - **Gone from the view:** `withdrawFlow`, the re-sync's own teardown. This applies only to a flow armed from the view, so a flow registered straight into the engine, with no stored row, is left alone. - **Read failed:** nothing changes. - **One flow sync queue.** The `kernel:ready` bind, the `metadata:reloaded` re-sync and the mutation sync now run serially, so no read and its registrations interleave with another's. - **Idempotency with the existing paths.** - A publish raises the mutation signal and then `metadata:reloaded`. The re-sync skips a flow the mutation sync just armed from the identical stored body, so a per-item publish registers the published flow **once**. - `PUT /api/v1/automation/:name` registers before it saves (`runtime/src/domains/automation.ts`, `registerAndSaveFlow`). Its save no longer causes a second registration, because the sync skips a definition whose canonical parse the engine already holds. - **Shutdown.** `destroy()` unsubscribes and waits for any queued sync. ## Measured beside the change - **Timing.** The signal is synchronous and fire-and-forget, so the save's answer comes before the arming. Measured: the flow is not yet registered when `await saveMetaItem` resolves, is registered after one `setImmediate` turn, and is armed 4.9 ms after the save on sqlite `:memory:`. The window is one read of the protocol's view. Hooks and actions on this door have the same shape. - **A run in flight when its flow is replaced** finishes on the definition it started with: the `tail_v1` node ran. - **A suspended run** resumes against the definition registered when it resumes: the `tail_v2` node ran. Once the flow is deleted, resume answers `RUN_NOT_FOUND` with `Flow 'held' not found for run …`. A re-registration through a publish or `PUT /automation/:name` already behaved this way (`resumeInternal` reads the current `this.flows` entry). - **Tenancy.** `flow` has no per-org channel. An organization-scoped flow save is refused at the door with `403 NOT_OVERRIDABLE` before any signal is raised (pinned). The sync never passes an event's organization to the read (pinned with a protocol stand-in), so it cannot arm beyond the boot's reach. - **PM assumption 3, refined.** "Unregister when inactive" became "registered and disarmed", which is what boot, kernel:ready and the re-sync all do with an `obsolete` flow. A delete unregisters. ## Pins `flow-metadata-save-arming.integration.test.ts` runs the real write path. `saveMetaItem`, `publishMetaItem` and `deleteMetaItem` are the producers behind the REST doors. The record-change trigger is a stand-in bound through ObjectQL's `registerHook` / `unregisterHooksByPackage`, because `@objectstack/trigger-record-change` depends on this package. A real insert fires it. 1. An active save is triggerable with no restart and no publish, and `execute` succeeds. 2. A record-triggered flow saved that way fires on the next matching insert. 3. A save to `status: 'obsolete'` leaves the flow registered with `enabled: false, bound: false`, and an insert does not fire it. 4. A delete unregisters the flow (`getFlow` returns `null`), and an insert does not fire it. 5. A draft save registers nothing. 6. The per-item publish door registers the published flow exactly once, and one insert launches it once. 7. `PUT /automation/:name`'s register-then-save order registers once. 8. An organization-scoped save is refused with `403 NOT_OVERRIDABLE` and arms nothing. `flow-metadata-mutation-reach.test.ts` covers two cases: the event's organization never reaches the read, and the env-wide control is armed. `inert-mode.test.ts` checks that inert mode subscribes nothing and that the default mode subscribes exactly once. ## Ablations (each through `scripts/ablation-replace.mjs`, run at `b8d422a2f4`, restore proven blob == HEAD with an empty `git diff HEAD`) Tests import `./plugin.js` directly, so `src` is the subject and no `dist` is involved. | Mutation | Predicted | Observed | |:---|:---|:---| | A1: remove the `subscribeFlowMutations(ctx)` call | the four direct-door rows (1–4), both reach rows and "subscribes once" go red; publish-once, draft, PUT-once, org-refused and inert-none stay green | 7 failed / 12 passed, exactly that split | | A2: the sync skips a body whose status is `obsolete`/`invalid` (handler ignores inactive state) | only row 3 (disarm) goes red | 1 failed / 18 passed: row 3 | | A3: drop the re-sync's skip for a mutation-armed body | only row 6 goes red, with 2 registrations | 1 failed / 22 passed: row 6, "expected length 1 but got 2" | | A4: drop the already-held skip | only row 7 goes red | 1 failed / 18 passed: row 7 | The first attempt at A2 never ran. `ablation-replace` refused it because the replacement still contained the anchor, so the count did not drop, and it restored the file. It was re-anchored and run as recorded above. ## Verification (code at `b8d422a2f4`; HEAD `ae04563d15` adds only docs and the changeset, with an empty `git diff b8d422a ae04563 -- packages/`) - `pnpm --filter @objectstack/service-automation test`: 172 files, **2110 passed**. - `pnpm --filter @objectstack/service-automation typecheck`: `tsc --noEmit` and `check:test-typecheck` OK, 0 debt. - **Derived gates at `ae04563d15`.** `dispatch-gates --commands` derived 92 commands. All 92 ran; 90 exited 0 on the first run. - `check:skill-examples` and `check:dual-build-cjs-loads` first exited 3 (PREREQUISITE NOT MET: no dist). After a full `turbo run build`, both exited 0. Readings: 262 examples; 106 entry points across 66 packages load. - `dispatch-gates --ran` reconciled 92 derived, 92 run, 0 NOT-MEASURED, all with exit codes. - **Narrowed lint** at `ae04563d15`: `eslint --no-inline-config --format json` on the four changed TS files read 4 files with 0 errors and 0 warnings. `--print-config` shows no `parserOptions.project` or `projectService`, so type-aware linting is off and this diff cannot move any untouched file's verdict. The repo-wide `pnpm lint` is left to CI. ## Docs - `content/docs/automation/flows.mdx` → "Discovery & Registration" gains one paragraph on how a runtime-authored flow follows its stored row. No existing sentence there was false. - `content/docs/references/api/automation-api.mdx:31-32` is auto-generated from `packages/spec`. It names `PUT /automation/:name` as one door that arms, which stays true. - `skills/**`: no sentence was made false (not edited; governed). ## Acceptance notes - **The answer-before-arm window**, one view read, is the signal's documented fire-and-forget shape and is shared with hooks and actions. The awaited ADR-0094 projector seam would close it, but it is a single slot per type that this package already fills for credential pruning, and it is not replayed to peer replicas, which `onMetadataMutation` is. Not taken: triage named the mutation signal. - `skills/objectstack-platform/references/plugin-hooks.md:52` lists the dev reload and publish-drafts as `metadata:reloaded` announcers and omits the per-item publish door. This predates this PR, is incomplete rather than false, and is not changed by it. Noted, not filed (carrier: none). --- _Generated by [Claude Code](https://claude.ai/code/session_01DiCSbmJrkzNhuEAier4VoJ)_ --------- Co-authored-by: Claude <noreply@anthropic.com>
1 parent 309224d commit 73b2246

6 files changed

Lines changed: 623 additions & 9 deletions

File tree

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
---
2+
'@objectstack/service-automation': patch
3+
---
4+
5+
A flow saved through the metadata API is armed on the running engine at once, as hooks and actions saved through the same door already are
6+
7+
Clause-②: no
8+
9+
`PUT /api/v1/meta/flow/:name` answered `200 "Saved flow … (env-wide, state=active)"` and `GET /api/v1/meta/flow/:name` served the row, but the automation engine registered nothing until the process restarted: `GET /api/v1/automation/:name` and `POST /api/v1/automation/:name/trigger` answered `404 Flow not found`, and a record-triggered flow never fired. The engine armed flows only at boot, at `kernel:ready` and on `metadata:reloaded`, and only the publish doors announce that event.
10+
11+
The automation service now listens to the metadata protocol's post-write signal (`onMetadataMutation`), the one ObjectQL already re-binds authored hooks and actions on, and makes the engine follow the stored row of each flow it names:
12+
13+
- an active save or a publish registers the flow, or re-registers it over the definition the engine held;
14+
- a save whose `status` is `'obsolete'` or `'invalid'` keeps it registered and unbound, as a boot does;
15+
- a delete unregisters it;
16+
- a draft save changes nothing until it is published.
17+
18+
The flow is re-read through the same execution view, precedence and env-wide scope the boot reads, so a save never arms a flow beyond the reach a restart would give it. The save answers first, and the registration follows it by one read of the stored metadata.
19+
20+
A publish raises this signal and `metadata:reloaded` together, and the published flow is still registered once, not twice. `PUT /api/v1/automation/:name`, which registers the flow before it saves it, is not registered a second time either.
21+
22+
A run already executing keeps the definition it started with. A suspended run resumes against the definition registered when it resumes, and answers `RUN_NOT_FOUND` once its flow is deleted. Both already held for a re-registration through a publish or `PUT /api/v1/automation/:name`.

‎content/docs/automation/flows.mdx‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1845,6 +1845,15 @@ export default defineStack({
18451845
});
18461846
```
18471847

1848+
A flow authored at runtime through the metadata API follows its stored row
1849+
on the running engine, with no restart. An active save
1850+
(`PUT /api/v1/meta/flow/:name`) or a publish registers the flow, or
1851+
re-registers it over the definition the engine held. A save
1852+
whose `status` is `'obsolete'` or `'invalid'` keeps it registered and
1853+
unbound, as a boot does. A delete unregisters it. A draft save changes
1854+
nothing until it is published. The write answers first, and the
1855+
registration follows it by one read of the stored metadata.
1856+
18481857
Registration is not arming. A `record_change`, `schedule`, `time_relative` or
18491858
`api` flow fires only when its trigger is installed, and that takes **both**
18501859
capability tokens on the same stack: `requires: ['automation', 'triggers']`.
Lines changed: 104 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,104 @@
1+
// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license.
2+
3+
/**
4+
* #21725 — the metadata-mutation sync arms what the boot arms, and no more.
5+
*
6+
* The boot binds flows from the protocol's execution view read with no
7+
* organization: the env-wide set. A mutation event carries the written row's
8+
* organization, and the real door refuses an org-scoped flow write
9+
* (`flow-metadata-save-arming.integration.test.ts` measures that), but a
10+
* signal can still name one — a delete of a pre-#6190 phantom row, or a peer
11+
* replica's replay. Handing that organization to the read would arm a row the
12+
* boot never arms. So the sync reads exactly what the boot reads, whatever the
13+
* event names.
14+
*
15+
* Harness: the automation plugin over a protocol stand-in that records every
16+
* read and lets the test raise the mutation signal, the shape
17+
* `flow-publish-rebind.test.ts` uses for `metadata:reloaded`.
18+
*/
19+
20+
import { describe, it, expect } from 'vitest';
21+
import { LiteKernel } from '@objectstack/core';
22+
import { AutomationServicePlugin } from './plugin.js';
23+
import type { AutomationEngine } from './engine.js';
24+
25+
type MutationListener = (evt: { type: string; name: string; state: string; organizationId?: string | null }) => void;
26+
27+
const flow = (name: string) => ({
28+
name,
29+
label: name,
30+
type: 'autolaunched',
31+
status: 'active',
32+
nodes: [
33+
{ id: 'start', type: 'start', label: 'Start' },
34+
{ id: 'end', type: 'end', label: 'End' },
35+
],
36+
edges: [{ id: 'e1', source: 'start', target: 'end' }],
37+
});
38+
39+
/** A protocol whose env-wide view and org view differ, recording each read's request. */
40+
function protocolStandIn() {
41+
const listeners: MutationListener[] = [];
42+
const reads: Array<Record<string, unknown>> = [];
43+
const envWide: unknown[] = [];
44+
const orgScoped: unknown[] = [flow('org_only')];
45+
return {
46+
service: {
47+
async getMetaItemsForExecution(request: { type: string; organizationId?: string }) {
48+
reads.push({ ...request });
49+
if (request.type !== 'flow') return { items: [] };
50+
return { items: request.organizationId ? [...envWide, ...orgScoped] : envWide };
51+
},
52+
onMetadataMutation(listener: MutationListener) {
53+
listeners.push(listener);
54+
return () => listeners.splice(listeners.indexOf(listener), 1);
55+
},
56+
},
57+
envWide,
58+
reads,
59+
emit: (evt: Parameters<MutationListener>[0]) => listeners.forEach((l) => l(evt)),
60+
};
61+
}
62+
63+
async function boot(proto: ReturnType<typeof protocolStandIn>) {
64+
const kernel = new LiteKernel({ logger: { level: 'silent' } } as never);
65+
kernel.use(new AutomationServicePlugin({ suspendedRunStore: 'memory' }));
66+
kernel.use({
67+
name: 'test.harness',
68+
type: 'standard',
69+
version: '1.0.0',
70+
dependencies: [],
71+
async init(ctx: { registerService(name: string, svc: unknown): void }) {
72+
ctx.registerService('protocol', proto.service);
73+
},
74+
async start() {},
75+
} as never);
76+
await kernel.bootstrap();
77+
return { kernel, engine: kernel.getService<AutomationEngine>('automation') };
78+
}
79+
80+
describe('the metadata-mutation sync reads what the boot reads (#21725)', () => {
81+
it("never hands the event's organization to the read, so an org-only row is not armed", async () => {
82+
const proto = protocolStandIn();
83+
const { kernel, engine } = await boot(proto);
84+
proto.reads.length = 0;
85+
86+
proto.emit({ type: 'flow', name: 'org_only', state: 'active', organizationId: 'org_a' });
87+
// `destroy()` waits for every queued flow sync.
88+
await kernel.shutdown();
89+
90+
expect(proto.reads, 'one read, of the env-wide view').toEqual([{ type: 'flow' }]);
91+
expect(await engine.getFlow('org_only')).toBeNull();
92+
});
93+
94+
it('arms an env-wide flow the same signal names — the control for the row above', async () => {
95+
const proto = protocolStandIn();
96+
const { kernel, engine } = await boot(proto);
97+
proto.envWide.push(flow('env_flow'));
98+
99+
proto.emit({ type: 'flow', name: 'env_flow', state: 'active', organizationId: null });
100+
await kernel.shutdown();
101+
102+
expect(await engine.getFlow('env_flow')).not.toBeNull();
103+
});
104+
});
Lines changed: 254 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,254 @@
1+
// Copyright (c) 2026 ObjectStack. Licensed under the Apache-2.0 license.
2+
3+
/**
4+
* #21725 — a flow saved through the metadata API is armed on the running
5+
* engine, on the same signal that already re-binds hooks and actions.
6+
*
7+
* The defect: `PUT /api/v1/meta/flow/:name` answered `200 "Saved flow … (env-wide,
8+
* state=active)"`, `GET /meta/flow/:name` served the row, and the engine held
9+
* nothing — `GET /api/v1/automation/:name` and its trigger door answered 404
10+
* and a record-triggered flow never fired, until a restart. The engine armed
11+
* flows only at boot, at `kernel:ready` and on `metadata:reloaded`, and only
12+
* the publish doors announce that event. A direct save raises the protocol's
13+
* post-persistence signal (`onMetadataMutation`); ObjectQL re-binds authored
14+
* hooks and actions on it, and `service-automation` did not listen.
15+
*
16+
* These tests run the REAL write path — ObjectKernel, ObjectQLPlugin's
17+
* metadata protocol (the producer behind `PUT /meta/flow/:name`, the per-item
18+
* publish door and `DELETE /meta/flow/:name`), driver-sql on better-sqlite3
19+
* `:memory:` and the real engine. `saveMetaItem` IS what the REST route calls;
20+
* the HTTP layer adds nothing on this path. The one stand-in is the
21+
* record-change trigger: `@objectstack/trigger-record-change` depends on this
22+
* package, so it binds through the same ObjectQL `registerHook` /
23+
* `unregisterHooksByPackage` pair the real trigger uses, and a real insert
24+
* fires it.
25+
*
26+
* The signal is fire-and-forget, so a save's answer precedes the arming by one
27+
* read of the protocol's view: the positive rows wait for it (`vi.waitFor`),
28+
* and the rows that assert something did NOT happen shut the kernel down
29+
* first — the plugin's `destroy()` waits for every flow sync already queued.
30+
*/
31+
32+
import { describe, it, expect, afterEach, vi } from 'vitest';
33+
import { ObjectKernel } from '@objectstack/core';
34+
import { ObjectQLPlugin, type ObjectQL } from '@objectstack/objectql';
35+
import { SqlDriver } from '@objectstack/driver-sql';
36+
import type { AutomationContext } from '@objectstack/spec/contracts';
37+
import { AutomationServicePlugin } from './plugin.js';
38+
import type { AutomationEngine, FlowTrigger, FlowTriggerBinding } from './engine.js';
39+
40+
/** The saved surface these tests drive: `ObjectStackProtocolImplementation`'s write doors. */
41+
interface MetadataWriteDoors {
42+
saveMetaItem(request: { type: string; name: string; item: unknown; mode?: 'draft' | 'publish' }): Promise<unknown>;
43+
publishMetaItem(request: { type: string; name: string }): Promise<unknown>;
44+
deleteMetaItem(request: { type: string; name: string }): Promise<unknown>;
45+
}
46+
47+
const TICKET = 'qa_ticket';
48+
const SYS = { isSystem: true } as const;
49+
50+
/** The repro's flow: a minimal start → end, saved active. */
51+
const plainFlow = (name: string, status: 'active' | 'obsolete' = 'active') => ({
52+
name,
53+
label: name,
54+
type: 'autolaunched',
55+
status,
56+
nodes: [
57+
{ id: 'start', type: 'start', label: 'Start' },
58+
{ id: 'end', type: 'end', label: 'End' },
59+
],
60+
edges: [{ id: 'e1', source: 'start', target: 'end' }],
61+
});
62+
63+
/** A flow launched by every insert into `qa_ticket`. */
64+
const ticketCreatedFlow = (name: string, status: 'active' | 'obsolete' = 'active') => ({
65+
...plainFlow(name, status),
66+
nodes: [
67+
{
68+
id: 'start',
69+
type: 'start',
70+
label: 'Start',
71+
config: { objectName: TICKET, triggerType: 'record-after-create' },
72+
},
73+
{ id: 'end', type: 'end', label: 'End' },
74+
],
75+
});
76+
77+
/**
78+
* A `record_change` trigger bound through ObjectQL's own hook registry, the
79+
* mechanism `RecordChangeTrigger` uses: `start` registers an `afterInsert` hook
80+
* under a per-flow package id, `stop` drops everything under it.
81+
*/
82+
function objectqlRecordTrigger(ql: ObjectQL) {
83+
const fired: string[] = [];
84+
const packageOf = (flowName: string) => `qa-record-trigger:${flowName}`;
85+
const trigger: FlowTrigger = {
86+
type: 'record_change',
87+
start(binding: FlowTriggerBinding, callback: (ctx: AutomationContext) => Promise<void>) {
88+
ql.unregisterHooksByPackage(packageOf(binding.flowName));
89+
ql.registerHook(
90+
'afterInsert',
91+
async (hook: { result?: unknown }) => {
92+
fired.push(binding.flowName);
93+
await callback({ record: hook.result, object: binding.object, event: 'afterInsert' } as never);
94+
},
95+
{ object: binding.object, packageId: packageOf(binding.flowName) },
96+
);
97+
},
98+
stop(flowName: string) {
99+
ql.unregisterHooksByPackage(packageOf(flowName));
100+
},
101+
};
102+
return { trigger, fired };
103+
}
104+
105+
describe('a flow saved through the metadata API is armed on the running engine (#21725)', () => {
106+
let kernel: ObjectKernel;
107+
let ql: ObjectQL;
108+
let protocol: MetadataWriteDoors;
109+
let automation: AutomationEngine;
110+
let fired: string[];
111+
112+
afterEach(async () => {
113+
vi.restoreAllMocks();
114+
if (kernel?.getState() === 'running') await kernel.shutdown();
115+
});
116+
117+
async function boot() {
118+
kernel = new ObjectKernel({ logger: { level: 'fatal' }, gracefulShutdown: false } as never);
119+
await kernel.use(new ObjectQLPlugin());
120+
await kernel.use(new AutomationServicePlugin({ suspendedRunStore: 'memory' }));
121+
await kernel.bootstrap();
122+
123+
ql = kernel.getService<ObjectQL>('objectql');
124+
automation = kernel.getService<AutomationEngine>('automation');
125+
protocol = kernel.getService<MetadataWriteDoors>('protocol');
126+
127+
const driver = new SqlDriver({ client: 'better-sqlite3', connection: { filename: ':memory:' }, useNullAsDefault: true });
128+
await driver.connect();
129+
ql.registerDriver(driver as never, true);
130+
ql.registry.registerObject(
131+
{ name: TICKET, label: 'Ticket', fields: { title: { name: 'title', label: 'Title', type: 'text' } } } as never,
132+
'qa-21725',
133+
'qa-21725',
134+
);
135+
await ql.syncSchemas();
136+
137+
const rec = objectqlRecordTrigger(ql);
138+
automation.registerTrigger(rec.trigger);
139+
fired = rec.fired;
140+
}
141+
142+
const saveActive = (name: string, item: unknown) => protocol.saveMetaItem({ type: 'flow', name, item });
143+
const insertTicket = (title: string) => ql.insert(TICKET, { title }, { context: SYS } as never);
144+
const runtimeState = (name: string) => automation.getFlowRuntimeStates().find((s) => s.name === name);
145+
146+
/** Every flow sync already queued has run: `destroy()` waits for them. */
147+
const settleFlowSync = () => kernel.shutdown();
148+
149+
it('an active save is triggerable without a restart or a publish', async () => {
150+
await boot();
151+
const saved = await saveActive('qa_meta_flow', plainFlow('qa_meta_flow'));
152+
// The card's step 2, as the door answers it.
153+
expect((saved as { message?: string }).message).toContain("Saved flow 'qa_meta_flow' (env-wide, state=active)");
154+
155+
await vi.waitFor(async () => {
156+
expect(await automation.getFlow('qa_meta_flow'), 'GET /automation/:name reads this').not.toBeNull();
157+
});
158+
const run = await automation.execute('qa_meta_flow', {} as never);
159+
expect(run.success, `the trigger door's run: ${JSON.stringify(run)}`).toBe(true);
160+
});
161+
162+
it('a record-triggered flow saved that way fires on the next matching write', async () => {
163+
await boot();
164+
await saveActive('ticket_created', ticketCreatedFlow('ticket_created'));
165+
await vi.waitFor(() => {
166+
expect(runtimeState('ticket_created')?.bound, 'the record trigger is bound').toBe(true);
167+
});
168+
169+
await insertTicket('first');
170+
expect(fired).toEqual(['ticket_created']);
171+
});
172+
173+
it("a save that leaves the flow `status: 'obsolete'` disarms it, as the boot registers such a flow", async () => {
174+
await boot();
175+
await saveActive('ticket_created', ticketCreatedFlow('ticket_created'));
176+
await vi.waitFor(() => expect(runtimeState('ticket_created')?.bound).toBe(true));
177+
178+
await saveActive('ticket_created', ticketCreatedFlow('ticket_created', 'obsolete'));
179+
await vi.waitFor(() => {
180+
expect(runtimeState('ticket_created'), 'still registered, disabled and unbound').toMatchObject({
181+
status: 'obsolete',
182+
enabled: false,
183+
bound: false,
184+
});
185+
});
186+
await insertTicket('after deactivation');
187+
expect(fired, 'a disarmed flow does not fire').toEqual([]);
188+
});
189+
190+
it('a delete through the same door unregisters it', async () => {
191+
await boot();
192+
await saveActive('ticket_created', ticketCreatedFlow('ticket_created'));
193+
await vi.waitFor(() => expect(runtimeState('ticket_created')?.bound).toBe(true));
194+
195+
await protocol.deleteMetaItem({ type: 'flow', name: 'ticket_created' });
196+
await vi.waitFor(async () => {
197+
expect(await automation.getFlow('ticket_created'), 'GET /automation/:name answers 404 again').toBeNull();
198+
});
199+
await insertTicket('after delete');
200+
expect(fired, 'an unregistered flow does not fire').toEqual([]);
201+
});
202+
203+
it('a draft save arms nothing: the live row is unchanged', async () => {
204+
await boot();
205+
const register = vi.spyOn(automation, 'registerFlow');
206+
await protocol.saveMetaItem({ type: 'flow', name: 'drafted', item: plainFlow('drafted'), mode: 'draft' });
207+
await settleFlowSync();
208+
209+
expect(register.mock.calls.filter(([name]) => name === 'drafted')).toEqual([]);
210+
expect(await automation.getFlow('drafted')).toBeNull();
211+
});
212+
213+
it('the per-item publish door (#10219) still arms the published flow exactly once', async () => {
214+
await boot();
215+
const register = vi.spyOn(automation, 'registerFlow');
216+
await protocol.saveMetaItem({ type: 'flow', name: 'ticket_created', item: ticketCreatedFlow('ticket_created'), mode: 'draft' });
217+
// The door awaits its `metadata:reloaded` announce, which runs behind the
218+
// mutation sync the same publish raised — both are done when it answers.
219+
await protocol.publishMetaItem({ type: 'flow', name: 'ticket_created' });
220+
221+
expect(register.mock.calls.filter(([name]) => name === 'ticket_created')).toHaveLength(1);
222+
await insertTicket('after publish');
223+
expect(fired, 'one binding, so one launch per write').toEqual(['ticket_created']);
224+
});
225+
226+
it('PUT /automation/:name, which registers before it saves, is not registered a second time', async () => {
227+
await boot();
228+
const register = vi.spyOn(automation, 'registerFlow');
229+
// `registerAndSaveFlow`'s order (runtime `domains/automation.ts`): the
230+
// engine first, then the metadata door's own save.
231+
automation.registerFlow('door_flow', plainFlow('door_flow'));
232+
await saveActive('door_flow', plainFlow('door_flow'));
233+
await settleFlowSync();
234+
235+
expect(register.mock.calls.filter(([name]) => name === 'door_flow')).toHaveLength(1);
236+
expect(await automation.getFlow('door_flow')).not.toBeNull();
237+
});
238+
239+
it('a flow is env-wide only: an organization-scoped save is refused and arms nothing', async () => {
240+
// The tenancy half of the boot's reach. `flow` declares no per-org
241+
// channel, so the door refuses the write before it lands — no
242+
// mutation is raised, and a flow write's signal never names an
243+
// organization the boot's env-wide read would not.
244+
await boot();
245+
const register = vi.spyOn(automation, 'registerFlow');
246+
await expect(
247+
protocol.saveMetaItem({ type: 'flow', name: 'org_flow', item: plainFlow('org_flow'), organizationId: 'org_a' } as never),
248+
).rejects.toMatchObject({ code: 'NOT_OVERRIDABLE', status: 403 });
249+
await settleFlowSync();
250+
251+
expect(register.mock.calls.filter(([name]) => name === 'org_flow')).toEqual([]);
252+
expect(await automation.getFlow('org_flow')).toBeNull();
253+
});
254+
});

0 commit comments

Comments
 (0)