-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathsequence.ts
More file actions
335 lines (312 loc) · 12.5 KB
/
Copy pathsequence.ts
File metadata and controls
335 lines (312 loc) · 12.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
/**
* The sequence engine: runs a validated scene.
*
* `showcontrol`'s `CueRunner` is the prior art, and two of its decisions carry
* over unchanged. Steps run **sequentially** unless a `parallel` says otherwise,
* because a scene that switches an input and then fades a light depends on that
* order. And a step that dispatched without throwing is reported as *dispatched*
* — not as *worked*. The device may not have complied, may not even be there;
* `waitFor` is the only step that claims otherwise, and it is the only one that
* checks.
*
* Effects are injected, so the engine still runs under test with no ioBroker and
* no real clock. What it does not do is decide *when* to run: that is a trigger,
* and ioBroker already has schedules, scripts and Blockly for it.
*/
import type { FailurePolicy, SequenceStep, StateValue } from "../model";
import type { ObjectSource } from "./resolver";
import type { Refusal, StateWrite } from "./actions";
import { plan } from "./actions";
import { read } from "./feedback";
import type { Registry } from "./registry";
import type { SceneBook } from "./scenes";
/**
* How often `waitFor` re-reads while waiting.
*
* Fast enough that a projector reporting ready is not sat on for a visible
* beat, slow enough to be free. The read is against a cache the adapter keeps
* from its subscriptions, not a device round trip.
*/
const POLL_MS = 50;
/** Everything the engine needs from the world outside it. */
export interface Effects {
/** Carry out planned writes. Throwing means they did not land. */
write(writes: ReadonlyArray<StateWrite>): Promise<void>;
/** Wait. */
sleep(ms: number): Promise<void>;
/**
* The object tree as it stands now.
*
* A function rather than a value because a scene's whole purpose is to
* outlive a moment: between a `do` and the `waitFor` that follows it, the
* thing being waited on is expected to change.
*/
source(): ObjectSource;
}
/** Why one step did not do what it said. */
export type FailureReason =
/** The action engine would not plan it. */
| { readonly kind: "refused"; readonly refusal: Refusal }
/** The writes were planned but did not land. */
| { readonly kind: "write-failed"; readonly message: string }
/** `waitFor` gave up. */
| { readonly kind: "timeout"; readonly waited: number; readonly last: StateValue | null }
/** A nested scene was dropped at load, so it is not in the book. */
| { readonly kind: "unknown-scene"; readonly scene: string };
export interface StepFailure {
/** Which scene, since a failure may be inside a nested one. */
readonly scene: string;
/** 1-based index of the step within that scene. */
readonly step: number;
readonly reason: FailureReason;
/** What the policy did about it. */
readonly handled: FailurePolicy["kind"];
}
export interface RunReport {
readonly scene: string;
/** True when every step ran, including ones whose failure was continued past. */
readonly completed: boolean;
/**
* Every failure, including those a `continue` policy passed over.
*
* Recording those is the point: `continue` means "the show goes on", never
* "nothing happened". An operator still needs to know the backup projector
* never came up.
*/
readonly failures: ReadonlyArray<StepFailure>;
}
/** Abort is the default because half a scene leaves equipment as nobody designed it. */
const DEFAULT_POLICY: FailurePolicy = { kind: "abort" };
/**
* What a step leaves the scene doing.
*
* `done` exists for `fallback` alone, and is why this is not a boolean: a
* fallback that succeeded means the scene should stop *without* the run being a
* failure. Switching to the backup projector is the designed outcome, not a
* degraded one — but the steps after it were aimed at the projector that died,
* so they must not run.
*/
type StepOutcome = "next" | "done" | "failed";
/**
* Runs a scene.
*
* @param id - Scene to run; must be in the book
* @param book - Validated scenes
* @param registry - The declared resources
* @param effects - Writes, waits and the object tree
* @returns What happened
*/
export async function run(id: string, book: SceneBook, registry: Registry, effects: Effects): Promise<RunReport> {
const failures: StepFailure[] = [];
const completed = await runScene(id, book, registry, effects, failures);
return { scene: id, completed, failures };
}
/**
* Runs one scene's steps in order.
*
* @param id - Scene id
* @param book - Validated scenes
* @param registry - The declared resources
* @param effects - Writes, waits and the object tree
* @param failures - Collected across the whole run, including nested scenes
* @returns False when a step aborted the scene
*/
async function runScene(
id: string,
book: SceneBook,
registry: Registry,
effects: Effects,
failures: StepFailure[],
): Promise<boolean> {
const scene = book.get(id);
if (!scene) {
// Load-time validation rejects unknown references, so reaching here
// means the scene itself was dropped for a fault of its own.
failures.push({
scene: id,
step: 0,
reason: { kind: "unknown-scene", scene: id },
handled: "abort",
});
return false;
}
for (const [index, step] of scene.steps.entries()) {
const outcome = await runStep(step, scene.id, index + 1, scene.onFailure, book, registry, effects, failures);
if (outcome === "done") {
return true;
}
if (outcome === "failed") {
return false;
}
}
return true;
}
/**
* Runs one step, applying its failure policy.
*
* @param step - The step
* @param scene - Scene it belongs to, for the report
* @param position - 1-based index within that scene
* @param sceneDefault - The scene's own policy, used when the step declares none
* @param book - Validated scenes
* @param registry - The declared resources
* @param effects - Writes, waits and the object tree
* @param failures - Collected across the whole run
* @returns What the scene should do next
*/
async function runStep(
step: SequenceStep,
scene: string,
position: number,
sceneDefault: FailurePolicy | undefined,
book: SceneBook,
registry: Registry,
effects: Effects,
failures: StepFailure[],
): Promise<StepOutcome> {
const policy = ("onFailure" in step ? step.onFailure : undefined) ?? sceneDefault ?? DEFAULT_POLICY;
const attempts = policy.kind === "retry" ? policy.times + 1 : 1;
let reason: FailureReason | undefined;
for (let attempt = 0; attempt < attempts; attempt++) {
if (attempt > 0 && policy.kind === "retry") {
await effects.sleep(policy.delayMs);
}
reason = await attemptStep(step, book, registry, effects, failures);
if (!reason) {
return "next";
}
}
// `reason` is set: the loop runs at least once and only exits early on success.
failures.push({ scene, step: position, reason: reason!, handled: policy.kind });
switch (policy.kind) {
case "continue":
return "next";
case "fallback":
// The fallback replaces the rest of this scene rather than resuming
// it — "switch to the backup projector" does not then want the
// remaining steps aimed at the dead one.
return (await runScene(policy.scene, book, registry, effects, failures)) ? "done" : "failed";
case "abort":
case "retry":
return "failed";
default:
// Unreachable: `SceneBook.load` rejects an unknown policy kind. Kept
// because the alternative to a case matching is falling off the end
// returning undefined, which `runScene` reads as "carry on" — so the
// failure mode of forgetting a case here is a scene that ignores its
// own abort and reports success. Fail closed instead.
return "failed";
}
}
/**
* Carries out one step once.
*
* @param step - The step
* @param book - Validated scenes
* @param registry - The declared resources
* @param effects - Writes, waits and the object tree
* @param failures - Collected across the whole run, for nested scenes
* @returns The reason it failed, or undefined on success
*/
async function attemptStep(
step: SequenceStep,
book: SceneBook,
registry: Registry,
effects: Effects,
failures: StepFailure[],
): Promise<FailureReason | undefined> {
switch (step.kind) {
case "delay":
await effects.sleep(step.ms);
return undefined;
case "do": {
const planned = plan(step.invoke, registry, effects.source());
if (!planned.ok) {
const { ok: _ok, ...refusal } = planned;
return { kind: "refused", refusal };
}
try {
await effects.write(planned.writes);
} catch (error) {
return { kind: "write-failed", message: (error as Error).message };
}
return undefined;
}
case "waitFor":
return waitFor(step, registry, effects);
case "scene":
return (await runScene(step.scene, book, registry, effects, failures))
? undefined
: { kind: "unknown-scene", scene: step.scene };
case "parallel": {
// Children run concurrently but each keeps its own policy, so one
// continuing does not drag the others down with it.
const results = await Promise.all(
step.steps.map((child, index) =>
runStep(child, "parallel", index + 1, undefined, book, registry, effects, failures),
),
);
return results.every(outcome => outcome !== "failed")
? undefined
: { kind: "write-failed", message: "a parallel step failed" };
}
}
}
/**
* Blocks until a feedback matches, or gives up.
*
* Requires the reading to be **healthy**. An unacknowledged or unresolved value
* that happens to equal the target is not the device confirming anything, and a
* scene that carried on from one would be acting on an assumption at the exact
* moment it asked not to.
*
* Honours the feedback's `settleMs`, so a value the owning adapter wrote
* optimistically and a poll later corrects will not satisfy the step. Nothing
* here detects an echo — that is not possible — it waits for the answer to stop
* changing, which is the observable weaker thing.
*
* Matches the semantic name *or* the device's own value, because `toDevice`
* already accepts either when a step writes one. Without that, a scene that
* routes by `1` has to wait on `"Camera 1"`, and the first real scene written
* against hardware fell into exactly that trap: step one switched the mixer and
* step two timed out waiting for a value it had just set.
*
* @param step - The wait
* @param registry - The declared resources
* @param effects - Waits and the object tree
* @returns A timeout reason, or undefined once matched
*/
async function waitFor(
step: Extract<SequenceStep, { kind: "waitFor" }>,
registry: Registry,
effects: Effects,
): Promise<FailureReason | undefined> {
const wanted = String(step.equals);
const settleMs = registry.getFeedback(step.resource, step.capability, step.feedback)?.settleMs ?? 0;
let waited = 0;
let last: StateValue | null = null;
let matchedAt: number | undefined;
for (;;) {
const reading = read(step.resource, step.capability, step.feedback, registry, effects.source());
last = reading?.value ?? null;
const matched = reading !== undefined && (String(reading.value) === wanted || String(reading.raw) === wanted);
if (reading?.healthy && matched) {
// The clock starts at the first match, not at the step, so a value
// arriving late still gets its full settling window.
matchedAt ??= waited;
if (waited - matchedAt >= settleMs) {
return undefined;
}
} else {
// Changed away, so what we saw was not the answer. An optimistic
// echo that a poll corrects lands here, which is the whole point.
matchedAt = undefined;
}
if (waited >= step.timeoutMs) {
return { kind: "timeout", waited, last };
}
const next = Math.min(POLL_MS, step.timeoutMs - waited);
await effects.sleep(next);
waited += next;
}
}