File
Blob: src/cloudflare/internal/workflows-api.ts
| 1 | // Copyright (c) 2024 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | export class NonRetryableError extends Error { |
| 6 | constructor(message: string, name = 'NonRetryableError') { |
| 7 | super(message); |
| 8 | this.name = name; |
| 9 | } |
| 10 | } |
| 11 | |
| 12 | interface Fetcher { |
| 13 | fetch: typeof fetch; |
| 14 | } |
| 15 | |
| 16 | async function callFetcher<T>( |
| 17 | fetcher: Fetcher, |
| 18 | path: string, |
| 19 | body: object |
| 20 | ): Promise<T> { |
| 21 | const res = await fetcher.fetch(`http://workflow-binding.local${path}`, { |
| 22 | method: 'POST', |
| 23 | headers: { |
| 24 | 'Content-Type': 'application/json', |
| 25 | 'X-Version': '1', |
| 26 | }, |
| 27 | body: JSON.stringify(body), |
| 28 | }); |
| 29 | |
| 30 | const response = (await res.json()) as { |
| 31 | result: T; |
| 32 | error?: WorkflowError; |
| 33 | }; |
| 34 | |
| 35 | if (res.ok) { |
| 36 | return response.result; |
| 37 | } else { |
| 38 | throw new Error(response.error?.message); |
| 39 | } |
| 40 | } |
| 41 | |
| 42 | class InstanceImpl implements WorkflowInstance { |
| 43 | // TODO(soon): Can we use the # syntax here? |
| 44 | // eslint-disable-next-line no-restricted-syntax |
| 45 | private readonly fetcher: Fetcher; |
| 46 | readonly id: string; |
| 47 | |
| 48 | constructor(id: string, fetcher: Fetcher) { |
| 49 | this.id = id; |
| 50 | this.fetcher = fetcher; |
| 51 | } |
| 52 | |
| 53 | async pause(): Promise<void> { |
| 54 | await callFetcher(this.fetcher, '/pause', { |
| 55 | id: this.id, |
| 56 | }); |
| 57 | } |
| 58 | async resume(): Promise<void> { |
| 59 | await callFetcher(this.fetcher, '/resume', { |
| 60 | id: this.id, |
| 61 | }); |
| 62 | } |
| 63 | |
| 64 | async terminate(): Promise<void> { |
| 65 | await callFetcher(this.fetcher, '/terminate', { |
| 66 | id: this.id, |
| 67 | }); |
| 68 | } |
| 69 | |
| 70 | async restart(options?: WorkflowInstanceRestartOptions): Promise<void> { |
| 71 | await callFetcher(this.fetcher, '/restart', { |
| 72 | ...options, |
| 73 | id: this.id, |
| 74 | }); |
| 75 | } |
| 76 | |
| 77 | async status(): Promise<InstanceStatus> { |
| 78 | const result = await callFetcher<InstanceStatus>(this.fetcher, '/status', { |
| 79 | id: this.id, |
| 80 | }); |
| 81 | return result; |
| 82 | } |
| 83 | |
| 84 | async sendEvent({ |
| 85 | type, |
| 86 | payload, |
| 87 | }: { |
| 88 | type: string; |
| 89 | payload: unknown; |
| 90 | }): Promise<void> { |
| 91 | await callFetcher(this.fetcher, '/send-event', { |
| 92 | type, |
| 93 | payload, |
| 94 | id: this.id, |
| 95 | }); |
| 96 | } |
| 97 | } |
| 98 | |
| 99 | class WorkflowImpl { |
| 100 | // TODO(soon): Can we use the # syntax here? |
| 101 | // eslint-disable-next-line no-restricted-syntax |
| 102 | private readonly fetcher: Fetcher; |
| 103 | |
| 104 | constructor(fetcher: Fetcher) { |
| 105 | this.fetcher = fetcher; |
| 106 | } |
| 107 | |
| 108 | async get(id: string): Promise<WorkflowInstance> { |
| 109 | const result = await callFetcher<{ |
| 110 | id: string; |
| 111 | }>(this.fetcher, '/get', { id }); |
| 112 | |
| 113 | return new InstanceImpl(result.id, this.fetcher); |
| 114 | } |
| 115 | |
| 116 | async create( |
| 117 | options?: WorkflowInstanceCreateOptions |
| 118 | ): Promise<WorkflowInstance> { |
| 119 | const result = await callFetcher<{ |
| 120 | id: string; |
| 121 | }>(this.fetcher, '/create', options ?? {}); |
| 122 | |
| 123 | return new InstanceImpl(result.id, this.fetcher); |
| 124 | } |
| 125 | |
| 126 | async createBatch( |
| 127 | options: WorkflowInstanceCreateOptions[] |
| 128 | ): Promise<WorkflowInstance[]> { |
| 129 | const results = await callFetcher< |
| 130 | { |
| 131 | id: string; |
| 132 | }[] |
| 133 | >(this.fetcher, '/createBatch', options); |
| 134 | |
| 135 | return results.map((result) => new InstanceImpl(result.id, this.fetcher)); |
| 136 | } |
| 137 | } |
| 138 | |
| 139 | export function makeBinding(env: { fetcher: Fetcher }): Workflow { |
| 140 | return new WorkflowImpl(env.fetcher); |
| 141 | } |
| 142 | |
| 143 | export default makeBinding; |