File
Blob: src/workerd/api/container.h
| 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 | |
| 15 | namespace workerd::api { |
| 16 | |
| 17 | class Fetcher; |
| 18 | class 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 | |
| 51 | struct 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 | |
| 75 | class 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. |
| 145 | class 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 |