Skip to content
File

Blob: src/workerd/api/container.h

cpp309 lines
1// Copyright (c) 2025 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#pragma once
6// Container management API for Durable Object-attached containers.
7//
8#include <workerd/api/streams/readable.h>
9#include <workerd/api/streams/writable.h>
10#include <workerd/io/compatibility-date.h>
11#include <workerd/io/container.capnp.h>
12#include <workerd/io/io-own.h>
13#include <workerd/jsg/jsg.h>
14 
15namespace workerd::api {
16 
17class Fetcher;
18class ExecOutput: public jsg::Object {
19 public:
20 ExecOutput(kj::Array<kj::byte> stdoutBytes, kj::Array<kj::byte> stderrBytes, int exitCode);
21 
22 jsg::JsArrayBuffer getStdout(jsg::Lock& js);
23 jsg::JsArrayBuffer getStderr(jsg::Lock& js);
24 int getExitCode() const {
25 return exitCode;
26 }
27 
28 JSG_RESOURCE_TYPE(ExecOutput) {
29 JSG_LAZY_READONLY_INSTANCE_PROPERTY(stdout, getStdout);
30 JSG_LAZY_READONLY_INSTANCE_PROPERTY(stderr, getStderr);
31 JSG_READONLY_PROTOTYPE_PROPERTY(exitCode, getExitCode);
32 
33 JSG_TS_OVERRIDE({
34 readonly stdout: ArrayBuffer;
35 readonly stderr: ArrayBuffer;
36 readonly exitCode: number;
37 });
38 }
39 
40 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
41 tracker.trackField("stdout", stdoutBytes);
42 tracker.trackField("stderr", stderrBytes);
43 }
44 
45 private:
46 kj::Array<kj::byte> stdoutBytes;
47 kj::Array<kj::byte> stderrBytes;
48 int exitCode;
49};
50 
51struct ExecOptions {
52 // $ prefix avoids collision with stdin/stdout/stderr macros from <stdio.h>;
53 // JSG_STRUCT strips the $ when exposing to JS.
54 jsg::Optional<kj::OneOf<jsg::Ref<ReadableStream>, kj::String>> $stdin;
55 jsg::Optional<kj::String> $stdout;
56 jsg::Optional<kj::String> $stderr;
57 jsg::Optional<kj::String> cwd;
58 jsg::Optional<jsg::Dict<kj::String>> env;
59 jsg::Optional<kj::String> user;
60 
61 JSG_STRUCT($stdin, $stdout, $stderr, cwd, env, user);
62 JSG_STRUCT_TS_OVERRIDE(ContainerExecOptions {
63 stdin?: ReadableStream | "pipe";
64 stdout?: "pipe" | "ignore";
65 stderr?: "pipe" | "ignore" | "combined";
66 cwd?: string;
67 env?: Record<string, string>;
68 user?: string;
69 $stdin: never;
70 $stdout: never;
71 $stderr: never;
72 });
73};
74 
75class ExecProcess: public jsg::Object {
76 public:
77 ExecProcess(jsg::Optional<jsg::Ref<WritableStream>> stdinStream,
78 jsg::Optional<jsg::Ref<ReadableStream>> stdoutStream,
79 jsg::Optional<jsg::Ref<ReadableStream>> stderrStream,
80 int pid,
81 rpc::Container::ProcessHandle::Client handle);
82 
83 jsg::Optional<jsg::Ref<WritableStream>> getStdin();
84 jsg::Optional<jsg::Ref<ReadableStream>> getStdout();
85 jsg::Optional<jsg::Ref<ReadableStream>> getStderr();
86 int getPid() const {
87 return pid;
88 }
89 jsg::MemoizedIdentity<jsg::Promise<int>>& getExitCode(jsg::Lock& js);
90 
91 jsg::Promise<jsg::Ref<ExecOutput>> output(jsg::Lock& js);
92 void kill(jsg::Lock& js, jsg::Optional<int> signal);
93 
94 JSG_RESOURCE_TYPE(ExecProcess) {
95 JSG_READONLY_PROTOTYPE_PROPERTY(stdin, getStdin);
96 JSG_READONLY_PROTOTYPE_PROPERTY(stdout, getStdout);
97 JSG_READONLY_PROTOTYPE_PROPERTY(stderr, getStderr);
98 JSG_READONLY_PROTOTYPE_PROPERTY(pid, getPid);
99 JSG_READONLY_PROTOTYPE_PROPERTY(exitCode, getExitCode);
100 JSG_METHOD(output);
101 JSG_METHOD(kill);
102 
103 JSG_TS_OVERRIDE({
104 readonly stdin: WritableStream | null;
105 readonly stdout: ReadableStream | null;
106 readonly stderr: ReadableStream | null;
107 readonly pid: number;
108 readonly exitCode: Promise<number>;
109 output(): Promise<ExecOutput>;
110 kill(signal?: number): void;
111 });
112 }
113 
114 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
115 tracker.trackField("stdin", stdinStream);
116 tracker.trackField("stdout", stdoutStream);
117 tracker.trackField("stderr", stderrStream);
118 tracker.trackField("exitCodePromise", exitCodePromise);
119 tracker.trackField("exitCodePromiseCopy", exitCodePromiseCopy);
120 }
121 
122 private:
123 void ensureExitCodePromise(jsg::Lock& js);
124 jsg::Promise<int> getExitCodeForOutput(jsg::Lock& js);
125 
126 jsg::Optional<jsg::Ref<WritableStream>> stdinStream;
127 jsg::Optional<jsg::Ref<ReadableStream>> stdoutStream;
128 jsg::Optional<jsg::Ref<ReadableStream>> stderrStream;
129 int pid;
130 IoOwn<rpc::Container::ProcessHandle::Client> handle;
131 kj::Maybe<jsg::MemoizedIdentity<jsg::Promise<int>>> exitCodePromise;
132 kj::Maybe<jsg::Promise<void>> exitCodePromiseCopy;
133 kj::Maybe<int> resolvedExitCode;
134 bool outputCalled = false;
135 
136 void visitForGc(jsg::GcVisitor& visitor) {
137 visitor.visit(stdinStream, stdoutStream, stderrStream, exitCodePromise, exitCodePromiseCopy);
138 }
139};
140 
141// Implements the `ctx.container` API for durable-object-attached containers. This API allows
142// the DO to supervise the attached container (lightweight virtual machine), including starting,
143// stopping, monitoring, making requests to the container, intercepting outgoing network requests,
144// etc.
145class Container: public jsg::Object {
146 public:
147 Container(rpc::Container::Client rpcClient, bool running);
148 
149 struct DirectorySnapshot {
150 kj::String id;
151 double size;
152 kj::String dir;
153 jsg::Optional<kj::String> name;
154 
155 JSG_STRUCT(id, size, dir, name);
156 };
157 
158 struct DirectorySnapshotOptions {
159 kj::String dir;
160 jsg::Optional<kj::String> name;
161 
162 JSG_STRUCT(dir, name);
163 };
164 
165 struct DirectorySnapshotRestoreParams {
166 DirectorySnapshot snapshot;
167 jsg::Optional<kj::String> mountPoint;
168 
169 JSG_STRUCT(snapshot, mountPoint);
170 };
171 
172 struct Snapshot {
173 kj::String id;
174 double size;
175 jsg::Optional<kj::String> name;
176 
177 JSG_STRUCT(id, size, name);
178 };
179 
180 struct SnapshotOptions {
181 jsg::Optional<kj::String> name;
182 
183 JSG_STRUCT(name);
184 };
185 
186 struct Info {
187 jsg::Dict<kj::String> labels;
188 
189 JSG_STRUCT(labels);
190 };
191 
192 struct StartupOptions {
193 jsg::Optional<kj::Array<kj::String>> entrypoint;
194 bool enableInternet = false;
195 jsg::Optional<jsg::Dict<kj::String>> env;
196 jsg::Optional<int64_t> hardTimeout;
197 jsg::Optional<jsg::Dict<kj::String>> labels;
198 jsg::Optional<kj::Array<DirectorySnapshotRestoreParams>> directorySnapshots;
199 jsg::Optional<Snapshot> containerSnapshot;
200 
201 // TODO(containers): Allow intercepting stdin/stdout/stderr by specifying streams here.
202 
203 JSG_STRUCT(entrypoint,
204 enableInternet,
205 env,
206 hardTimeout,
207 labels,
208 directorySnapshots,
209 containerSnapshot);
210 JSG_STRUCT_TS_OVERRIDE_DYNAMIC(CompatibilityFlags::Reader flags) {
211 if (flags.getWorkerdExperimental()) {
212 JSG_TS_OVERRIDE(ContainerStartupOptions {
213 entrypoint?: string[];
214 enableInternet: boolean;
215 env?: Record<string, string>;
216 hardTimeout?: number | bigint;
217 labels?: Record<string, string>;
218 directorySnapshots?: ContainerDirectorySnapshotRestoreParams[];
219 containerSnapshot?: ContainerSnapshot;
220 });
221 } else {
222 JSG_TS_OVERRIDE(ContainerStartupOptions {
223 entrypoint?: string[];
224 enableInternet: boolean;
225 env?: Record<string, string>;
226 hardTimeout?: never;
227 labels?: Record<string, string>;
228 directorySnapshots?: ContainerDirectorySnapshotRestoreParams[];
229 containerSnapshot?: ContainerSnapshot;
230 });
231 }
232 }
233 };
234 
235 bool getRunning() const {
236 return running;
237 }
238 
239 // Methods correspond closely to the RPC interface in `container.capnp`.
240 void start(jsg::Lock& js, jsg::Optional<StartupOptions> options);
241 jsg::Promise<void> monitor(jsg::Lock& js);
242 jsg::Promise<void> destroy(jsg::Lock& js, jsg::Optional<jsg::Value> error);
243 void signal(jsg::Lock& js, int signo);
244 jsg::Ref<Fetcher> getTcpPort(jsg::Lock& js, int port);
245 jsg::Promise<void> setInactivityTimeout(jsg::Lock& js, int64_t durationMs);
246 jsg::Promise<void> interceptOutboundHttp(
247 jsg::Lock& js, kj::String addr, jsg::Ref<Fetcher> binding);
248 jsg::Promise<void> interceptAllOutboundHttp(jsg::Lock& js, jsg::Ref<Fetcher> binding);
249 jsg::Promise<void> interceptOutboundHttps(
250 jsg::Lock& js, kj::String addr, jsg::Ref<Fetcher> binding);
251 jsg::Promise<void> interceptOutboundTcp(
252 jsg::Lock& js, kj::String addr, jsg::Ref<Fetcher> binding);
253 jsg::Promise<DirectorySnapshot> snapshotDirectory(
254 jsg::Lock& js, DirectorySnapshotOptions options);
255 jsg::Promise<Snapshot> snapshotContainer(jsg::Lock& js, SnapshotOptions options);
256 jsg::Promise<jsg::Ref<ExecProcess>> exec(
257 jsg::Lock& js, kj::Array<kj::String> cmd, jsg::Optional<ExecOptions> options);
258 
259 jsg::Promise<kj::Maybe<Info>> inspect(jsg::Lock& js);
260 
261 // TODO(containers): listenTcp()
262 
263 JSG_RESOURCE_TYPE(Container, CompatibilityFlags::Reader flags) {
264 JSG_READONLY_PROTOTYPE_PROPERTY(running, getRunning);
265 JSG_METHOD(start);
266 JSG_METHOD(monitor);
267 JSG_METHOD(destroy);
268 JSG_METHOD(signal);
269 JSG_METHOD(getTcpPort);
270 JSG_METHOD(setInactivityTimeout);
271 
272 JSG_METHOD(interceptOutboundHttp);
273 JSG_METHOD(interceptAllOutboundHttp);
274 JSG_METHOD(snapshotDirectory);
275 JSG_METHOD(snapshotContainer);
276 JSG_METHOD(interceptOutboundHttps);
277 if (flags.getWorkerdExperimental()) {
278 JSG_METHOD(exec);
279 JSG_METHOD(interceptOutboundTcp);
280 JSG_METHOD(inspect);
281 }
282 }
283 
284 void visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
285 tracker.trackField("destroyReason", destroyReason);
286 }
287 
288 private:
289 IoOwn<rpc::Container::Client> rpcClient;
290 bool running;
291 
292 kj::Maybe<jsg::Value> destroyReason;
293 
294 void visitForGc(jsg::GcVisitor& visitor) {
295 visitor.visit(destroyReason);
296 }
297 
298 class TcpPortWorkerInterface;
299 class TcpPortOutgoingFactory;
300};
301 
302#define EW_CONTAINER_ISOLATE_TYPES \
303 api::ExecOutput, api::ExecOptions, api::ExecProcess, api::Container, \
304 api::Container::DirectorySnapshot, api::Container::DirectorySnapshotOptions, \
305 api::Container::DirectorySnapshotRestoreParams, api::Container::Snapshot, \
306 api::Container::SnapshotOptions, api::Container::StartupOptions, api::Container::Info
307 
308} // namespace workerd::api