Press n or j to go to the next uncovered block, b, p or k for the previous block.
| 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 | 28x 28x 28x 12x 12x 12x 6x 6x 2x 1x 1x 6x 6x 4x 2x | /**
* @module
*
* @author Bayu Dwiyan Satria
* @version 1.0.0
* @since 1.3.0
*/
import { Binding } from '@/core/bindings/Binding'
import { lazySettings } from '@/utils/lazySettings'
import { WorkflowSettings } from '@/types/WorkflowSettings'
import type { CloudflareEnv } from '@/types/CloudflareEnv'
/**
* Accessor for a Workflows binding.
*
* @remarks
* A Workflow is durable execution: a class of steps that survives restarts,
* retries a failed step without replaying the ones before it, and can sleep for
* a month between two of them. This accessor is the *producer* side — starting
* instances and reading them back. The steps themselves live in a
* `WorkflowEntrypoint` class the Worker exports, which this package neither
* provides nor wraps: a base class that must be exported from the Worker's own
* entry, and named in its own `wrangler.json`, has nothing an adapter can add.
*
* The method names mirror the platform's — `create`, `createBatch`, `get` —
* because an accessor that renames operations makes Cloudflare's documentation
* stop matching this one. `status` is the single addition, saving the two-step
* that reading an instance's state otherwise takes.
*
* ## Why there is no `WorkflowService`
*
* Every binding here but this one and Browser Rendering has a capability
* service beside it, and the reason differs from Browser's. There the obstacle
* was a dependency; here it is that the kernel has no capability to implement.
* `@bayudwiyansatria/core` names a `CacheStore`, a `DataStore`, a
* `MessageQueue` — abstractions with more than one plausible provider. Durable
* execution has none yet, and inventing one to sit in front of a single
* implementation with no consumers would be ceremony rather than abstraction,
* which the binding guide says in as many words.
*
* That leaves nothing lost, because unlike Browser Rendering the binding's own
* operations cost no dependency: `Workflow`, `WorkflowInstance` and their
* options come from `@cloudflare/workers-types`, already a peer. A consumer
* gets the real surface here rather than a bare handle.
*
* @typeParam Params The event payload instances of this Workflow are started with.
*
* @example
* ```ts
* const screens = new WorkflowBinding<ScreenRequest>()
*
* const run = await screens.create(env, { universe: 'IDX', filters })
*
* return ctx.json({ id: run.id })
* ```
*
* @class
*
* @author Bayu Dwiyan Satria
* @version 1.0.0
* @since 1.3.0
*/
export class WorkflowBinding<Params = unknown, E extends CloudflareEnv = CloudflareEnv> extends Binding<
Workflow<Params>,
E
> {
/**
* This accessor's settings, resolved on first read.
*/
private readonly settings: () => WorkflowSettings
/**
* Constructs a WorkflowBinding instance.
*
* @param settings Resolved settings. Defaults to `resolve('workflow')`, so callers
* normally construct this with no arguments at all.
*/
constructor(settings?: WorkflowSettings) {
const settle = lazySettings<WorkflowSettings>('workflow', settings)
super(() => settle().binding)
this.settings = settle
}
/**
* Starts an instance.
*
* @remarks
* Returns as soon as the instance is accepted, not when it finishes — a
* Workflow that runs for an hour returns a handle in milliseconds. Read the
* outcome later with {@link status}.
*
* Supplying an `id` in `options` makes the start idempotent in the only sense
* Cloudflare offers: a second start under an id that already exists throws
* rather than producing a second run.
*
* @param env The Worker environment.
* @param params The event payload the instance is triggered with.
* @param options Creation options — an explicit id, or a retention policy overriding the configured one.
* @returns A handle to the new instance.
* @throws {@link MissingBindingError} When the Workflow binding is absent.
*/
public async create(
env: E,
params?: Params,
options: WorkflowInstanceCreateOptions<Params> = {}
): Promise<WorkflowInstance> {
const request = params === undefined ? { ...options } : { ...options, params }
return await this.resolve(env).create(this.withRetention(request))
}
/**
* Starts several instances in one call.
*
* @remarks
* Cloudflare caps a batch at 100 instances, or at 1 MiB of payload, whichever
* is reached first — so a caller fanning out over a large universe splits the
* work itself rather than relying on this to do it. Nothing here chunks on
* the caller's behalf, because the right chunk size depends on how big the
* payloads are and only the caller knows that.
*
* @param env The Worker environment.
* @param batch Creation options, one entry per instance.
* @returns Handles to the new instances, in the order they were given.
* @throws {@link MissingBindingError} When the Workflow binding is absent.
*/
public async createBatch(env: E, batch: WorkflowInstanceCreateOptions<Params>[]): Promise<WorkflowInstance[]> {
return await this.resolve(env).createBatch(batch.map(options => this.withRetention(options)))
}
/**
* Reads back an instance started earlier.
*
* @param env The Worker environment.
* @param id The instance id.
* @returns A handle to the instance, which can be paused, resumed, terminated or read.
* @throws {@link MissingBindingError} When the Workflow binding is absent.
*/
public async get(env: E, id: string): Promise<WorkflowInstance> {
return await this.resolve(env).get(id)
}
/**
* Reads an instance's current status.
*
* @remarks
* The one method here without a counterpart on the binding, and the only one
* worth adding: reporting a run's progress is what a status endpoint does on
* every request, and it would otherwise take two awaits every time.
*
* A finished instance carries its result in `output`; a failed one carries
* `error`. Both disappear once retention lapses — see {@link WorkflowSettings}.
*
* @param env The Worker environment.
* @param id The instance id.
* @returns The instance's status, plus its output or error when it has one.
* @throws {@link MissingBindingError} When the Workflow binding is absent.
*/
public async status(env: E, id: string): Promise<InstanceStatus> {
return await (await this.get(env, id)).status()
}
/**
* Applies the configured retention policy to options that do not set one.
*
* @remarks
* All or nothing, deliberately. A call that names `retention` at all owns
* both halves of it, so a caller asking for a longer error retention on one
* run does not silently inherit the configured success retention beside it.
*
* @param options The caller-supplied creation options.
* @returns Creation options with the configured retention filled in.
*/
private withRetention(options: WorkflowInstanceCreateOptions<Params>): WorkflowInstanceCreateOptions<Params> {
const { successRetention, errorRetention } = this.settings()
if (options.retention !== undefined || (successRetention === undefined && errorRetention === undefined)) {
return options
}
return { ...options, retention: { successRetention, errorRetention } }
}
}
|