Skip to content
File

Blob: src/cloudflare/internal/workflows-api.ts

typescript144 lines
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 
5export class NonRetryableError extends Error {
6 constructor(message: string, name = 'NonRetryableError') {
7 super(message);
8 this.name = name;
9 }
10}
11 
12interface Fetcher {
13 fetch: typeof fetch;
14}
15 
16async 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 
42class 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 
99class 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 
139export function makeBinding(env: { fetcher: Fetcher }): Workflow {
140 return new WorkflowImpl(env.fetcher);
141}
142 
143export default makeBinding;