File
Blob: src/workerd/api/container.c++
| 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 | #include "container.h" |
| 6 | |
| 7 | #include <workerd/api/http.h> |
| 8 | #include <workerd/api/streams/readable.h> |
| 9 | #include <workerd/api/streams/writable.h> |
| 10 | #include <workerd/api/system-streams.h> |
| 11 | #include <workerd/io/features.h> |
| 12 | #include <workerd/io/io-context.h> |
| 13 | |
| 14 | #include <capnp/compat/byte-stream.h> |
| 15 | #include <kj/filesystem.h> |
| 16 | |
| 17 | namespace workerd::api { |
| 18 | |
| 19 | namespace { |
| 20 | |
| 21 | kj::Maybe<kj::Path> parseRestorePath(kj::StringPtr path) { |
| 22 | JSG_REQUIRE(path.size() > 0 && path[0] == '/', TypeError, |
| 23 | "Directory snapshot restore path must be absolute. Got: ", path); |
| 24 | |
| 25 | try { |
| 26 | auto parsed = kj::Path::parse(path.slice(1)); |
| 27 | if (parsed.size() == 0) { |
| 28 | return kj::none; |
| 29 | } |
| 30 | return kj::mv(parsed); |
| 31 | } catch (kj::Exception&) { |
| 32 | JSG_FAIL_REQUIRE( |
| 33 | TypeError, "Directory snapshot restore path contains invalid components: ", path); |
| 34 | } |
| 35 | } |
| 36 | |
| 37 | void requireValidEnvNameAndValue(kj::StringPtr name, kj::StringPtr value) { |
| 38 | JSG_REQUIRE(name.findFirst('=') == kj::none, Error, |
| 39 | "Environment variable names cannot contain '=': ", name); |
| 40 | JSG_REQUIRE(name.findFirst('\0') == kj::none, Error, |
| 41 | "Environment variable names cannot contain '\\0': ", name); |
| 42 | JSG_REQUIRE(value.findFirst('\0') == kj::none, Error, |
| 43 | "Environment variable values cannot contain '\\0': ", name); |
| 44 | } |
| 45 | |
| 46 | kj::String getExecOutputMode(jsg::Optional<kj::String> maybeMode, kj::StringPtr kind) { |
| 47 | auto mode = kj::mv(maybeMode).orDefault(kj::str("pipe")); |
| 48 | JSG_REQUIRE(mode == "pipe" || mode == "ignore" || (kind == "stderr" && mode == "combined"), |
| 49 | TypeError, "Invalid ", kind, " option: ", mode); |
| 50 | return mode; |
| 51 | } |
| 52 | |
| 53 | kj::Array<kj::byte> emptyByteArray() { |
| 54 | return kj::heapArray<kj::byte>(0); |
| 55 | } |
| 56 | |
| 57 | capnp::ByteStream::Client makeExecPipe( |
| 58 | capnp::ByteStreamFactory& factory, kj::Own<kj::AsyncOutputStream> output) { |
| 59 | return factory.kjToCapnp(capnp::ExplicitEndOutputStream::wrap(kj::mv(output), []() {})); |
| 60 | } |
| 61 | |
| 62 | } // namespace |
| 63 | |
| 64 | // ======================================================================================= |
| 65 | // ExecOutput / ExecProcess |
| 66 | |
| 67 | ExecOutput::ExecOutput( |
| 68 | kj::Array<kj::byte> stdoutBytes, kj::Array<kj::byte> stderrBytes, int exitCode) |
| 69 | : stdoutBytes(kj::mv(stdoutBytes)), |
| 70 | stderrBytes(kj::mv(stderrBytes)), |
| 71 | exitCode(exitCode) {} |
| 72 | |
| 73 | jsg::JsArrayBuffer ExecOutput::getStdout(jsg::Lock& js) { |
| 74 | return jsg::JsArrayBuffer::create(js, stdoutBytes); |
| 75 | } |
| 76 | |
| 77 | jsg::JsArrayBuffer ExecOutput::getStderr(jsg::Lock& js) { |
| 78 | return jsg::JsArrayBuffer::create(js, stderrBytes); |
| 79 | } |
| 80 | |
| 81 | ExecProcess::ExecProcess(jsg::Optional<jsg::Ref<WritableStream>> stdinStream, |
| 82 | jsg::Optional<jsg::Ref<ReadableStream>> stdoutStream, |
| 83 | jsg::Optional<jsg::Ref<ReadableStream>> stderrStream, |
| 84 | int pid, |
| 85 | rpc::Container::ProcessHandle::Client handle) |
| 86 | : stdinStream(kj::mv(stdinStream)), |
| 87 | stdoutStream(kj::mv(stdoutStream)), |
| 88 | stderrStream(kj::mv(stderrStream)), |
| 89 | pid(pid), |
| 90 | handle(IoContext::current().addObject(kj::heap(kj::mv(handle)))) {} |
| 91 | |
| 92 | jsg::Optional<jsg::Ref<WritableStream>> ExecProcess::getStdin() { |
| 93 | return stdinStream.map([](jsg::Ref<WritableStream>& stream) { return stream.addRef(); }); |
| 94 | } |
| 95 | |
| 96 | jsg::Optional<jsg::Ref<ReadableStream>> ExecProcess::getStdout() { |
| 97 | return stdoutStream.map([](jsg::Ref<ReadableStream>& stream) { return stream.addRef(); }); |
| 98 | } |
| 99 | |
| 100 | jsg::Optional<jsg::Ref<ReadableStream>> ExecProcess::getStderr() { |
| 101 | return stderrStream.map([](jsg::Ref<ReadableStream>& stream) { return stream.addRef(); }); |
| 102 | } |
| 103 | |
| 104 | void ExecProcess::ensureExitCodePromise(jsg::Lock& js) { |
| 105 | if (exitCodePromise != kj::none) { |
| 106 | return; |
| 107 | } |
| 108 | |
| 109 | // jsg::Promise is single-use. Keep the original Promise<int> as the public `exitCode` |
| 110 | // property and a separate whenResolved() branch for helpers like output(). |
| 111 | auto self = JSG_THIS; |
| 112 | auto promise = IoContext::current().awaitIo(js, |
| 113 | handle->waitRequest(capnp::MessageSize{4, 0}) |
| 114 | .send() |
| 115 | .then([self = kj::mv(self)]( |
| 116 | capnp::Response<rpc::Container::ProcessHandle::WaitResults>&& results) mutable { |
| 117 | auto exitCode = results.getExitCode(); |
| 118 | self->resolvedExitCode = exitCode; |
| 119 | return exitCode; |
| 120 | })); |
| 121 | |
| 122 | exitCodePromiseCopy = promise.whenResolved(js); |
| 123 | exitCodePromise = jsg::MemoizedIdentity<jsg::Promise<int>>(kj::mv(promise)); |
| 124 | } |
| 125 | |
| 126 | jsg::Promise<int> ExecProcess::getExitCodeForOutput(jsg::Lock& js) { |
| 127 | ensureExitCodePromise(js); |
| 128 | |
| 129 | KJ_IF_SOME(exitCode, resolvedExitCode) { |
| 130 | return js.resolvedPromise(static_cast<int>(exitCode)); |
| 131 | } |
| 132 | |
| 133 | auto self = JSG_THIS; |
| 134 | return KJ_ASSERT_NONNULL(exitCodePromiseCopy) |
| 135 | .whenResolved(js) |
| 136 | .then(js, [self = kj::mv(self)](jsg::Lock&) -> int { |
| 137 | return static_cast<int>(KJ_ASSERT_NONNULL(self->resolvedExitCode)); |
| 138 | }); |
| 139 | } |
| 140 | |
| 141 | jsg::MemoizedIdentity<jsg::Promise<int>>& ExecProcess::getExitCode(jsg::Lock& js) { |
| 142 | ensureExitCodePromise(js); |
| 143 | return KJ_ASSERT_NONNULL(exitCodePromise); |
| 144 | } |
| 145 | |
| 146 | jsg::Promise<jsg::Ref<ExecOutput>> ExecProcess::output(jsg::Lock& js) { |
| 147 | JSG_REQUIRE(!outputCalled, TypeError, "output() can only be called once."); |
| 148 | outputCalled = true; |
| 149 | |
| 150 | auto stdoutPromise = js.resolvedPromise(emptyByteArray()); |
| 151 | KJ_IF_SOME(stream, stdoutStream) { |
| 152 | JSG_REQUIRE(!stream->isDisturbed(), TypeError, |
| 153 | "Cannot call output() after stdout has started being consumed."); |
| 154 | stdoutPromise = |
| 155 | stream->getController() |
| 156 | .readAllBytes(js, IoContext::current().getLimitEnforcer().getBufferingLimit()) |
| 157 | .then(js, [](jsg::Lock&, jsg::BufferSource bytes) { |
| 158 | return kj::heapArray(bytes.asArrayPtr()); |
| 159 | }); |
| 160 | } |
| 161 | |
| 162 | auto stderrPromise = js.resolvedPromise(emptyByteArray()); |
| 163 | KJ_IF_SOME(stream, stderrStream) { |
| 164 | JSG_REQUIRE(!stream->isDisturbed(), TypeError, |
| 165 | "Cannot call output() after stderr has started being consumed."); |
| 166 | stderrPromise = stream->getController() |
| 167 | .readAllBytes(js, kj::maxValue) |
| 168 | .then(js, [](jsg::Lock&, jsg::BufferSource bytes) { |
| 169 | return kj::heapArray(bytes.asArrayPtr()); |
| 170 | }); |
| 171 | } |
| 172 | |
| 173 | auto exitCodePromise = getExitCodeForOutput(js); |
| 174 | |
| 175 | return stdoutPromise.then(js, |
| 176 | [stderrPromise = kj::mv(stderrPromise), exitCodePromise = kj::mv(exitCodePromise)]( |
| 177 | jsg::Lock& js, |
| 178 | kj::Array<kj::byte> stdoutBytes) mutable -> jsg::Promise<jsg::Ref<ExecOutput>> { |
| 179 | return stderrPromise.then(js, |
| 180 | [stdoutBytes = kj::mv(stdoutBytes), exitCodePromise = kj::mv(exitCodePromise)]( |
| 181 | jsg::Lock& js, |
| 182 | kj::Array<kj::byte> stderrBytes) mutable -> jsg::Promise<jsg::Ref<ExecOutput>> { |
| 183 | return exitCodePromise.then(js, |
| 184 | [stdoutBytes = kj::mv(stdoutBytes), stderrBytes = kj::mv(stderrBytes)]( |
| 185 | jsg::Lock& js, int exitCode) mutable -> jsg::Ref<ExecOutput> { |
| 186 | return js.alloc<ExecOutput>(kj::mv(stdoutBytes), kj::mv(stderrBytes), exitCode); |
| 187 | }); |
| 188 | }); |
| 189 | }); |
| 190 | } |
| 191 | |
| 192 | void ExecProcess::kill(jsg::Lock& js, jsg::Optional<int> signal) { |
| 193 | auto signo = signal.orDefault(15); |
| 194 | JSG_REQUIRE(signo > 0 && signo <= 64, RangeError, "Invalid signal number."); |
| 195 | |
| 196 | auto req = handle->killRequest(capnp::MessageSize{4, 0}); |
| 197 | req.setSigno(signo); |
| 198 | IoContext::current().addTask(req.sendIgnoringResult()); |
| 199 | } |
| 200 | // ======================================================================================= |
| 201 | // Basic lifecycle methods |
| 202 | |
| 203 | Container::Container(rpc::Container::Client rpcClient, bool running) |
| 204 | : rpcClient(IoContext::current().addObject(kj::heap(kj::mv(rpcClient)))), |
| 205 | running(running) {} |
| 206 | |
| 207 | void Container::start(jsg::Lock& js, jsg::Optional<StartupOptions> maybeOptions) { |
| 208 | auto flags = FeatureFlags::get(js); |
| 209 | JSG_REQUIRE(!running, Error, "start() cannot be called on a container that is already running."); |
| 210 | |
| 211 | StartupOptions options = kj::mv(maybeOptions).orDefault({}); |
| 212 | |
| 213 | auto req = rpcClient->startRequest(); |
| 214 | KJ_IF_SOME(entrypoint, options.entrypoint) { |
| 215 | auto list = req.initEntrypoint(entrypoint.size()); |
| 216 | for (auto i: kj::indices(entrypoint)) { |
| 217 | list.set(i, entrypoint[i]); |
| 218 | } |
| 219 | } |
| 220 | req.setEnableInternet(options.enableInternet); |
| 221 | |
| 222 | KJ_IF_SOME(env, options.env) { |
| 223 | auto list = req.initEnvironmentVariables(env.fields.size()); |
| 224 | for (auto i: kj::indices(env.fields)) { |
| 225 | auto field = &env.fields[i]; |
| 226 | requireValidEnvNameAndValue(field->name, field->value); |
| 227 | |
| 228 | list.set(i, str(field->name, "=", field->value)); |
| 229 | } |
| 230 | } |
| 231 | |
| 232 | if (flags.getWorkerdExperimental()) { |
| 233 | KJ_IF_SOME(hardTimeoutMs, options.hardTimeout) { |
| 234 | JSG_REQUIRE(hardTimeoutMs > 0, RangeError, "Hard timeout must be greater than 0"); |
| 235 | req.setHardTimeoutMs(hardTimeoutMs); |
| 236 | } |
| 237 | } |
| 238 | |
| 239 | KJ_IF_SOME(labels, options.labels) { |
| 240 | auto list = req.initLabels(labels.fields.size()); |
| 241 | for (auto i: kj::indices(labels.fields)) { |
| 242 | auto& field = labels.fields[i]; |
| 243 | JSG_REQUIRE(field.name.size() > 0, Error, "Label names cannot be empty"); |
| 244 | for (auto c: field.name) { |
| 245 | JSG_REQUIRE(static_cast<kj::byte>(c) >= 0x20, Error, |
| 246 | "Label names cannot contain control characters (index ", i, ")"); |
| 247 | } |
| 248 | for (auto c: field.value) { |
| 249 | JSG_REQUIRE(static_cast<kj::byte>(c) >= 0x20, Error, |
| 250 | "Label values cannot contain control characters (index ", i, ")"); |
| 251 | } |
| 252 | list[i].setName(field.name); |
| 253 | list[i].setValue(field.value); |
| 254 | } |
| 255 | } |
| 256 | |
| 257 | KJ_IF_SOME(directorySnapshots, options.directorySnapshots) { |
| 258 | auto list = req.initDirectorySnapshots(directorySnapshots.size()); |
| 259 | for (auto i: kj::indices(directorySnapshots)) { |
| 260 | auto entry = list[i]; |
| 261 | auto& restore = directorySnapshots[i]; |
| 262 | auto& snap = restore.snapshot; |
| 263 | auto effectiveRestorePath = snap.dir.asPtr(); |
| 264 | KJ_IF_SOME(mp, restore.mountPoint) { |
| 265 | effectiveRestorePath = mp.asPtr(); |
| 266 | } |
| 267 | |
| 268 | JSG_REQUIRE_NONNULL(parseRestorePath(effectiveRestorePath), Error, |
| 269 | "Directory snapshot cannot be restored to root directory."); |
| 270 | |
| 271 | entry.setSnapshotId(snap.id); |
| 272 | entry.setRestorePath(effectiveRestorePath); |
| 273 | } |
| 274 | } |
| 275 | |
| 276 | KJ_IF_SOME(containerSnapshot, options.containerSnapshot) { |
| 277 | req.setContainerSnapshotId(containerSnapshot.id); |
| 278 | } |
| 279 | |
| 280 | req.setCompatibilityFlags(flags); |
| 281 | |
| 282 | IoContext::current().addTask(req.sendIgnoringResult()); |
| 283 | |
| 284 | running = true; |
| 285 | } |
| 286 | |
| 287 | jsg::Promise<kj::Maybe<Container::Info>> Container::inspect(jsg::Lock& js) { |
| 288 | return IoContext::current().awaitIo(js, rpcClient->inspectRequest().send(), |
| 289 | [](jsg::Lock& js, |
| 290 | capnp::Response<rpc::Container::InspectResults> results) -> kj::Maybe<Info> { |
| 291 | auto info = results.getInfo(); |
| 292 | if (info.isNone()) { |
| 293 | return kj::none; |
| 294 | } |
| 295 | return Info{ |
| 296 | .labels = |
| 297 | jsg::Dict<kj::String>{ |
| 298 | .fields = |
| 299 | KJ_MAP(label, info.getStarted().getLabels()) { |
| 300 | return jsg::Dict<kj::String>::Field{ |
| 301 | .name = kj::str(label.getName()), |
| 302 | .value = kj::str(label.getValue()), |
| 303 | }; |
| 304 | }, |
| 305 | }, |
| 306 | }; |
| 307 | }); |
| 308 | } |
| 309 | |
| 310 | jsg::Promise<Container::DirectorySnapshot> Container::snapshotDirectory( |
| 311 | jsg::Lock& js, DirectorySnapshotOptions options) { |
| 312 | JSG_REQUIRE( |
| 313 | running, Error, "snapshotDirectory() cannot be called on a container that is not running."); |
| 314 | JSG_REQUIRE(options.dir.size() > 0 && options.dir.startsWith("/"), TypeError, |
| 315 | "snapshotDirectory() requires an absolute directory path (starting with '/')."); |
| 316 | |
| 317 | auto req = rpcClient->snapshotDirectoryRequest(); |
| 318 | req.setDir(options.dir); |
| 319 | |
| 320 | KJ_IF_SOME(name, options.name) { |
| 321 | req.setName(name); |
| 322 | } |
| 323 | |
| 324 | return IoContext::current() |
| 325 | .awaitIo(js, req.send()) |
| 326 | .then( |
| 327 | js, [](jsg::Lock& js, capnp::Response<rpc::Container::SnapshotDirectoryResults> results) { |
| 328 | auto snapshot = results.getSnapshot(); |
| 329 | JSG_REQUIRE(snapshot.getSize() <= jsg::MAX_SAFE_INTEGER, RangeError, |
| 330 | "Snapshot size exceeds Number.MAX_SAFE_INTEGER"); |
| 331 | |
| 332 | jsg::Optional<kj::String> name = kj::none; |
| 333 | if (snapshot.getName().size() > 0) { |
| 334 | name = kj::str(snapshot.getName()); |
| 335 | } |
| 336 | |
| 337 | return Container::DirectorySnapshot{kj::str(snapshot.getId()), |
| 338 | static_cast<double>(snapshot.getSize()), kj::str(snapshot.getDir()), kj::mv(name)}; |
| 339 | }); |
| 340 | } |
| 341 | |
| 342 | jsg::Promise<Container::Snapshot> Container::snapshotContainer( |
| 343 | jsg::Lock& js, SnapshotOptions options) { |
| 344 | JSG_REQUIRE( |
| 345 | running, Error, "snapshotContainer() cannot be called on a container that is not running."); |
| 346 | |
| 347 | auto req = rpcClient->snapshotContainerRequest(); |
| 348 | |
| 349 | KJ_IF_SOME(name, options.name) { |
| 350 | req.setName(name); |
| 351 | } |
| 352 | |
| 353 | return IoContext::current() |
| 354 | .awaitIo(js, req.send()) |
| 355 | .then( |
| 356 | js, [](jsg::Lock& js, capnp::Response<rpc::Container::SnapshotContainerResults> results) { |
| 357 | auto snapshot = results.getSnapshot(); |
| 358 | JSG_REQUIRE(snapshot.getSize() <= jsg::MAX_SAFE_INTEGER, RangeError, |
| 359 | "Snapshot size exceeds Number.MAX_SAFE_INTEGER"); |
| 360 | |
| 361 | jsg::Optional<kj::String> name = kj::none; |
| 362 | if (snapshot.getName().size() > 0) { |
| 363 | name = kj::str(snapshot.getName()); |
| 364 | } |
| 365 | |
| 366 | return Container::Snapshot{ |
| 367 | kj::str(snapshot.getId()), static_cast<double>(snapshot.getSize()), kj::mv(name)}; |
| 368 | }); |
| 369 | } |
| 370 | |
| 371 | jsg::Promise<void> Container::setInactivityTimeout(jsg::Lock& js, int64_t durationMs) { |
| 372 | JSG_REQUIRE( |
| 373 | durationMs > 0, TypeError, "setInactivityTimeout() cannot be called with a durationMs <= 0"); |
| 374 | |
| 375 | auto req = rpcClient->setInactivityTimeoutRequest(); |
| 376 | |
| 377 | req.setDurationMs(durationMs); |
| 378 | return IoContext::current().awaitIo(js, req.sendIgnoringResult()); |
| 379 | } |
| 380 | |
| 381 | jsg::Promise<void> Container::interceptOutboundHttp( |
| 382 | jsg::Lock& js, kj::String addr, jsg::Ref<Fetcher> binding) { |
| 383 | auto& ioctx = IoContext::current(); |
| 384 | auto channel = binding->getSubrequestChannel(ioctx); |
| 385 | |
| 386 | // Get a channel token for RPC usage, the container runtime can use this |
| 387 | // token later to redeem a Fetcher. |
| 388 | auto token = channel->getToken(IoChannelFactory::ChannelTokenUsage::RPC); |
| 389 | |
| 390 | auto req = rpcClient->setEgressHttpRequest(); |
| 391 | req.setHostPort(addr); |
| 392 | req.setChannelToken(token); |
| 393 | return ioctx.awaitIo(js, req.sendIgnoringResult()); |
| 394 | } |
| 395 | |
| 396 | jsg::Promise<void> Container::interceptAllOutboundHttp(jsg::Lock& js, jsg::Ref<Fetcher> binding) { |
| 397 | auto& ioctx = IoContext::current(); |
| 398 | auto channel = binding->getSubrequestChannel(ioctx); |
| 399 | auto token = channel->getToken(IoChannelFactory::ChannelTokenUsage::RPC); |
| 400 | |
| 401 | // Register for all IPv4 and IPv6 addresses (on port 80) |
| 402 | auto reqV4 = rpcClient->setEgressHttpRequest(); |
| 403 | reqV4.setHostPort("0.0.0.0/0"_kj); |
| 404 | reqV4.setChannelToken(token); |
| 405 | |
| 406 | auto reqV6 = rpcClient->setEgressHttpRequest(); |
| 407 | reqV6.setHostPort("::/0"_kj); |
| 408 | reqV6.setChannelToken(token); |
| 409 | |
| 410 | return ioctx.awaitIo(js, |
| 411 | kj::joinPromisesFailFast(kj::arr(reqV4.sendIgnoringResult(), reqV6.sendIgnoringResult()))); |
| 412 | } |
| 413 | |
| 414 | jsg::Promise<void> Container::interceptOutboundHttps( |
| 415 | jsg::Lock& js, kj::String addr, jsg::Ref<Fetcher> binding) { |
| 416 | auto& ioctx = IoContext::current(); |
| 417 | auto channel = binding->getSubrequestChannel(ioctx); |
| 418 | auto token = channel->getToken(IoChannelFactory::ChannelTokenUsage::RPC); |
| 419 | |
| 420 | auto req = rpcClient->setEgressHttpsRequest(); |
| 421 | req.setHostPort(addr); |
| 422 | req.setChannelToken(token); |
| 423 | |
| 424 | return ioctx.awaitIo(js, req.sendIgnoringResult()); |
| 425 | } |
| 426 | |
| 427 | jsg::Promise<jsg::Ref<ExecProcess>> Container::exec( |
| 428 | jsg::Lock& js, kj::Array<kj::String> cmd, jsg::Optional<ExecOptions> maybeOptions) { |
| 429 | JSG_REQUIRE(running, Error, "exec() cannot be called on a container that is not running."); |
| 430 | JSG_REQUIRE(cmd.size() > 0, TypeError, "exec() requires a non-empty command array."); |
| 431 | |
| 432 | auto options = kj::mv(maybeOptions).orDefault({}); |
| 433 | auto stdoutMode = getExecOutputMode(kj::mv(options.$stdout), "stdout"); |
| 434 | auto stderrMode = getExecOutputMode(kj::mv(options.$stderr), "stderr"); |
| 435 | bool combinedOutput = stderrMode == "combined"; |
| 436 | JSG_REQUIRE(!combinedOutput || stdoutMode == "pipe", TypeError, |
| 437 | "stderr: \"combined\" requires stdout to be \"pipe\"."); |
| 438 | |
| 439 | auto& ioContext = IoContext::current(); |
| 440 | auto& byteStreamFactory = ioContext.getByteStreamFactory(); |
| 441 | |
| 442 | auto req = rpcClient->execRequest(); |
| 443 | auto cmdList = req.initCmd(cmd.size()); |
| 444 | for (auto i: kj::indices(cmd)) { |
| 445 | cmdList.set(i, cmd[i]); |
| 446 | } |
| 447 | |
| 448 | // Init the kj pipes to create the stdout/err bytestreams |
| 449 | kj::Maybe<kj::Own<kj::AsyncInputStream>> stdoutInput; |
| 450 | if (stdoutMode == "pipe") { |
| 451 | auto pipe = kj::newOneWayPipe(); |
| 452 | req.setStdoutWriter(makeExecPipe(byteStreamFactory, kj::mv(pipe.out))); |
| 453 | stdoutInput = kj::mv(pipe.in); |
| 454 | } |
| 455 | |
| 456 | kj::Maybe<kj::Own<kj::AsyncInputStream>> stderrInput; |
| 457 | if (!combinedOutput && stderrMode == "pipe") { |
| 458 | auto pipe = kj::newOneWayPipe(); |
| 459 | req.setStderrWriter(makeExecPipe(byteStreamFactory, kj::mv(pipe.out))); |
| 460 | stderrInput = kj::mv(pipe.in); |
| 461 | } |
| 462 | |
| 463 | auto params = req.initParams(); |
| 464 | params.setCombinedOutput(combinedOutput); |
| 465 | |
| 466 | // Some basic validation... |
| 467 | KJ_IF_SOME(cwd, options.cwd) { |
| 468 | JSG_REQUIRE(cwd.findFirst('\0') == kj::none, TypeError, "cwd cannot contain '\\0' characters."); |
| 469 | params.setWorkingDirectory(cwd); |
| 470 | } |
| 471 | |
| 472 | KJ_IF_SOME(user, options.user) { |
| 473 | JSG_REQUIRE( |
| 474 | user.findFirst('\0') == kj::none, TypeError, "user cannot contain '\\0' characters."); |
| 475 | params.setUser(user); |
| 476 | } |
| 477 | |
| 478 | KJ_IF_SOME(env, options.env) { |
| 479 | auto envList = params.initEnv(env.fields.size()); |
| 480 | for (auto i: kj::indices(env.fields)) { |
| 481 | auto field = &env.fields[i]; |
| 482 | requireValidEnvNameAndValue(field->name, field->value); |
| 483 | envList.set(i, str(field->name, "=", field->value)); |
| 484 | } |
| 485 | } |
| 486 | |
| 487 | // We have to await, because PID won't be available until the response resolves |
| 488 | return ioContext.awaitIo(js, req.send()) |
| 489 | .then(js, |
| 490 | [&ioContext, &byteStreamFactory, options = kj::mv(options), |
| 491 | stdoutInput = kj::mv(stdoutInput), stderrInput = kj::mv(stderrInput)]( |
| 492 | jsg::Lock& js, capnp::Response<rpc::Container::ExecResults> results) mutable |
| 493 | -> jsg::Ref<ExecProcess> { |
| 494 | auto process = results.getProcess(); |
| 495 | auto handle = process.getHandle(); |
| 496 | auto pid = process.getPid(); |
| 497 | |
| 498 | // Init the ReadableStreams (stdout/stderr) |
| 499 | jsg::Optional<jsg::Ref<ReadableStream>> stdoutStream = kj::none; |
| 500 | KJ_IF_SOME(input, stdoutInput) { |
| 501 | auto source = newSystemStream(kj::mv(input), StreamEncoding::IDENTITY, ioContext); |
| 502 | stdoutStream = js.alloc<ReadableStream>(ioContext, kj::mv(source)); |
| 503 | } |
| 504 | |
| 505 | // stderrInput is only set if using "pipe" on stderr and not "combined" |
| 506 | jsg::Optional<jsg::Ref<ReadableStream>> stderrStream = kj::none; |
| 507 | KJ_IF_SOME(input, stderrInput) { |
| 508 | auto source = newSystemStream(kj::mv(input), StreamEncoding::IDENTITY, ioContext); |
| 509 | stderrStream = js.alloc<ReadableStream>(ioContext, kj::mv(source)); |
| 510 | } |
| 511 | |
| 512 | jsg::Optional<jsg::Ref<WritableStream>> stdinStream = kj::none; |
| 513 | |
| 514 | // If stdin is undefined, the JS API promises immediate EOF. We still use the pipelined stdin() |
| 515 | // capability so exec() doesn't wait on an extra round-trip. |
| 516 | KJ_IF_SOME(stdinOption, options.$stdin) { |
| 517 | auto stdinRequest = handle.stdinWriterRequest(capnp::MessageSize{4, 0}); |
| 518 | // Get the stdinWriter() ByteStream, use the pipelined capability |
| 519 | auto stdinPipeline = stdinRequest.send(); |
| 520 | // ... adapt bytestream into a writer |
| 521 | auto stdinWriter = byteStreamFactory.capnpToKjExplicitEnd(stdinPipeline.getWriter()); |
| 522 | |
| 523 | KJ_SWITCH_ONEOF(stdinOption) { |
| 524 | // user sets ReadableStream... |
| 525 | KJ_CASE_ONEOF(readable, jsg::Ref<ReadableStream>) { |
| 526 | auto sink = newSystemStream(kj::mv(stdinWriter), StreamEncoding::IDENTITY, ioContext); |
| 527 | auto pipePromise = |
| 528 | (ioContext.waitForDeferredProxy(readable->pumpTo(js, kj::mv(sink), true))); |
| 529 | ioContext.addTask(pipePromise.attach(readable.addRef())); |
| 530 | } |
| 531 | // user sets "pipe"... they want to consume the API with the stdin WritableStream |
| 532 | KJ_CASE_ONEOF(mode, kj::String) { |
| 533 | JSG_REQUIRE( |
| 534 | mode == "pipe", TypeError, "stdin must be a ReadableStream or the string \"pipe\"."); |
| 535 | auto sink = newSystemStream(kj::mv(stdinWriter), StreamEncoding::IDENTITY, ioContext); |
| 536 | auto writable = js.alloc<WritableStream>(ioContext, kj::mv(sink), |
| 537 | ioContext.getMetrics().tryCreateWritableByteStreamObserver()); |
| 538 | stdinStream = kj::mv(writable); |
| 539 | } |
| 540 | } |
| 541 | |
| 542 | // all good, we have the stdinStream set |
| 543 | } else { |
| 544 | auto stdinRequest = handle.stdinWriterRequest(capnp::MessageSize{4, 0}); |
| 545 | auto stdinPipeline = stdinRequest.send(); |
| 546 | auto stdinWriter = byteStreamFactory.capnpToKjExplicitEnd(stdinPipeline.getWriter()); |
| 547 | ioContext.addTask(stdinWriter->end().attach(kj::mv(stdinWriter))); |
| 548 | } |
| 549 | |
| 550 | // return the instance to the process after getting pipeline of the process handle |
| 551 | return js.alloc<ExecProcess>( |
| 552 | kj::mv(stdinStream), kj::mv(stdoutStream), kj::mv(stderrStream), pid, kj::mv(handle)); |
| 553 | }); |
| 554 | } |
| 555 | |
| 556 | jsg::Promise<void> Container::interceptOutboundTcp( |
| 557 | jsg::Lock& js, kj::String addr, jsg::Ref<Fetcher> binding) { |
| 558 | auto& ioctx = IoContext::current(); |
| 559 | auto channel = binding->getSubrequestChannel(ioctx); |
| 560 | |
| 561 | // Get a channel token for RPC usage, the container runtime can use this |
| 562 | // token later to redeem a Fetcher whose connect() handler processes the TCP stream. |
| 563 | auto token = channel->getToken(IoChannelFactory::ChannelTokenUsage::RPC); |
| 564 | |
| 565 | auto req = rpcClient->setEgressTcpRequest(); |
| 566 | req.setHostPort(addr); |
| 567 | req.setChannelToken(token); |
| 568 | return ioctx.awaitIo(js, req.sendIgnoringResult()); |
| 569 | } |
| 570 | |
| 571 | jsg::Promise<void> Container::monitor(jsg::Lock& js) { |
| 572 | JSG_REQUIRE(running, Error, "monitor() cannot be called on a container that is not running."); |
| 573 | |
| 574 | return IoContext::current() |
| 575 | .awaitIo(js, rpcClient->monitorRequest(capnp::MessageSize{4, 0}).send()) |
| 576 | .then(js, [this](jsg::Lock& js, capnp::Response<rpc::Container::MonitorResults> results) { |
| 577 | running = false; |
| 578 | auto exitCode = results.getExitCode(); |
| 579 | KJ_IF_SOME(d, destroyReason) { |
| 580 | jsg::Value error = kj::mv(d); |
| 581 | destroyReason = kj::none; |
| 582 | js.throwException(kj::mv(error)); |
| 583 | } |
| 584 | |
| 585 | if (exitCode != 0) { |
| 586 | auto err = js.error(kj::str("Container exited with unexpected exit code: ", exitCode)); |
| 587 | KJ_ASSERT_NONNULL(err.tryCast<jsg::JsObject>()).set(js, "exitCode", js.num(exitCode)); |
| 588 | js.throwException(err); |
| 589 | } |
| 590 | }, [this](jsg::Lock& js, jsg::Value&& error) { |
| 591 | running = false; |
| 592 | destroyReason = kj::none; |
| 593 | js.throwException(kj::mv(error)); |
| 594 | }); |
| 595 | } |
| 596 | |
| 597 | jsg::Promise<void> Container::destroy(jsg::Lock& js, jsg::Optional<jsg::Value> error) { |
| 598 | if (!running) return js.resolvedPromise(); |
| 599 | |
| 600 | if (destroyReason == kj::none) { |
| 601 | destroyReason = kj::mv(error); |
| 602 | } |
| 603 | |
| 604 | return IoContext::current().awaitIo( |
| 605 | js, rpcClient->destroyRequest(capnp::MessageSize{4, 0}).sendIgnoringResult()); |
| 606 | } |
| 607 | |
| 608 | void Container::signal(jsg::Lock& js, int signo) { |
| 609 | JSG_REQUIRE(signo > 0 && signo <= 64, RangeError, "Invalid signal number."); |
| 610 | JSG_REQUIRE(running, Error, "signal() cannot be called on a container that is not running."); |
| 611 | |
| 612 | auto req = rpcClient->signalRequest(capnp::MessageSize{4, 0}); |
| 613 | req.setSigno(signo); |
| 614 | IoContext::current().addTask(req.sendIgnoringResult()); |
| 615 | } |
| 616 | |
| 617 | // ======================================================================================= |
| 618 | // getTcpPort() |
| 619 | |
| 620 | // `getTcpPort()` returns a `Fetcher`, on which `fetch()` and `connect()` can be called. `Fetcher` |
| 621 | // is a JavaScript wrapper around `WorkerInterface`, so we need to implement that. |
| 622 | class Container::TcpPortWorkerInterface final: public WorkerInterface { |
| 623 | public: |
| 624 | TcpPortWorkerInterface(capnp::ByteStreamFactory& byteStreamFactory, |
| 625 | kj::EntropySource& entropySource, |
| 626 | const kj::HttpHeaderTable& headerTable, |
| 627 | rpc::Container::Port::Client port) |
| 628 | : byteStreamFactory(byteStreamFactory), |
| 629 | entropySource(entropySource), |
| 630 | headerTable(headerTable), |
| 631 | port(kj::mv(port)) {} |
| 632 | |
| 633 | // Implements fetch(), i.e., HTTP requests. We form a TCP connection, then run HTTP over it |
| 634 | // (as opposed to, say, speaking http-over-capnp to the container service). |
| 635 | kj::Promise<void> request(kj::HttpMethod method, |
| 636 | kj::StringPtr url, |
| 637 | const kj::HttpHeaders& headers, |
| 638 | kj::AsyncInputStream& requestBody, |
| 639 | kj::HttpService::Response& response) override { |
| 640 | // URLs should have been validated earlier in the stack, so parsing the URL should succeed. |
| 641 | auto parsedUrl = KJ_REQUIRE_NONNULL(kj::Url::tryParse(url, kj::Url::Context::HTTP_PROXY_REQUEST, |
| 642 | {.percentDecode = false, .allowEmpty = true}), |
| 643 | "invalid url?", url); |
| 644 | |
| 645 | // We don't support TLS. |
| 646 | JSG_REQUIRE(parsedUrl.scheme != "https", Error, |
| 647 | "Connecting to a container using HTTPS is not currently supported; use HTTP instead. " |
| 648 | "TLS is unnecessary anyway, as the connection is already secure by default."); |
| 649 | |
| 650 | // Schemes other than http: and https: should have been rejected earlier, but let's verify. |
| 651 | KJ_REQUIRE(parsedUrl.scheme == "http"); |
| 652 | |
| 653 | // We need to convert the URL from proxy format (full URL in request line) to host format |
| 654 | // (path in request line, hostname in Host header). |
| 655 | auto newHeaders = headers.cloneShallow(); |
| 656 | newHeaders.setPtr(kj::HttpHeaderId::HOST, parsedUrl.host); |
| 657 | auto noHostUrl = parsedUrl.toString(kj::Url::Context::HTTP_REQUEST); |
| 658 | |
| 659 | // Make a TCP connection... |
| 660 | auto pipe = kj::newTwoWayPipe(); |
| 661 | kj::Maybe<kj::Exception> connectionException = kj::none; |
| 662 | |
| 663 | auto connectionPromise = connectImpl(*pipe.ends[1]); |
| 664 | |
| 665 | // ... and then stack an HttpClient on it ... |
| 666 | auto client = kj::newHttpClient(headerTable, *pipe.ends[0], {.entropySource = entropySource}); |
| 667 | |
| 668 | // ... and then adapt that to an HttpService ... |
| 669 | auto service = kj::newHttpService(*client); |
| 670 | |
| 671 | // ... fork connection promises so we can keep the original exception around ... |
| 672 | auto connectionPromiseForked = connectionPromise.fork(); |
| 673 | auto connectionPromiseBranch = connectionPromiseForked.addBranch(); |
| 674 | auto connectionPromiseToKeepException = connectionPromiseForked.addBranch(); |
| 675 | |
| 676 | // ... and now we can just forward our call to that ... |
| 677 | try { |
| 678 | co_await service->request(method, noHostUrl, newHeaders, requestBody, response) |
| 679 | .exclusiveJoin( |
| 680 | // never done as we do not want a Connection RPC exiting successfully |
| 681 | // affecting the request |
| 682 | connectionPromiseBranch.then([]() -> kj::Promise<void> { return kj::NEVER_DONE; })); |
| 683 | } catch (...) { |
| 684 | auto exception = kj::getCaughtExceptionAsKj(); |
| 685 | connectionException = kj::some(kj::mv(exception)); |
| 686 | } |
| 687 | |
| 688 | // ... and last but not least, if the connect() call succeeded but the connection |
| 689 | // was broken, we throw that exception. |
| 690 | KJ_IF_SOME(exception, connectionException) { |
| 691 | co_await connectionPromiseToKeepException; |
| 692 | kj::throwFatalException(kj::mv(exception)); |
| 693 | } |
| 694 | } |
| 695 | |
| 696 | // Implements connect(), i.e., forms a raw socket. |
| 697 | kj::Promise<void> connect(kj::StringPtr host, |
| 698 | const kj::HttpHeaders& headers, |
| 699 | kj::AsyncIoStream& connection, |
| 700 | ConnectResponse& response, |
| 701 | kj::HttpConnectSettings settings) override { |
| 702 | JSG_REQUIRE(!settings.useTls, Error, |
| 703 | "Connencting to a container using TLS is not currently supported. It is unnecessary " |
| 704 | "anyway, as the connection is already secure by default."); |
| 705 | |
| 706 | auto promise = connectImpl(connection); |
| 707 | |
| 708 | kj::HttpHeaders responseHeaders(headerTable); |
| 709 | response.accept(200, "OK", responseHeaders); |
| 710 | |
| 711 | return promise; |
| 712 | } |
| 713 | |
| 714 | // The only `CustomEvent` that can happen through `Fetcher` is a JSRPC call. Maybe we will |
| 715 | // support this someday? But not today. |
| 716 | kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override { |
| 717 | return event->notSupported(); |
| 718 | } |
| 719 | |
| 720 | // There's no way to invoke the remaining event types via `Fetcher`. |
| 721 | kj::Promise<void> prewarm(kj::StringPtr url) override { |
| 722 | KJ_UNREACHABLE; |
| 723 | } |
| 724 | kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override { |
| 725 | KJ_UNREACHABLE; |
| 726 | } |
| 727 | kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override { |
| 728 | KJ_UNREACHABLE; |
| 729 | } |
| 730 | |
| 731 | private: |
| 732 | capnp::ByteStreamFactory& byteStreamFactory; |
| 733 | kj::EntropySource& entropySource; |
| 734 | const kj::HttpHeaderTable& headerTable; |
| 735 | rpc::Container::Port::Client port; |
| 736 | |
| 737 | // Connect to the port and pump bytes to/from `connection`. Used by both request() and |
| 738 | // connect(). |
| 739 | kj::Promise<void> connectImpl(kj::AsyncIoStream& connection) { |
| 740 | // A lot of the following is copied from |
| 741 | // capnp::HttpOverCapnpFactory::KjToCapnpHttpServiceAdapter::connect(). |
| 742 | auto req = port.connectRequest(capnp::MessageSize{4, 1}); |
| 743 | auto downPipe = kj::newOneWayPipe(); |
| 744 | req.setDown(byteStreamFactory.kjToCapnp(kj::mv(downPipe.out))); |
| 745 | auto pipeline = req.send(); |
| 746 | |
| 747 | // Make sure the request message isn't pinned into memory through the co_await below. |
| 748 | { auto drop = kj::mv(req); } |
| 749 | |
| 750 | auto downPumpTask = |
| 751 | downPipe.in->pumpTo(connection) |
| 752 | .then([&connection, down = kj::mv(downPipe.in)](uint64_t) -> kj::Promise<void> { |
| 753 | connection.shutdownWrite(); |
| 754 | return kj::NEVER_DONE; |
| 755 | }); |
| 756 | auto up = pipeline.getUp(); |
| 757 | |
| 758 | auto upStream = byteStreamFactory.capnpToKjExplicitEnd(up); |
| 759 | auto upPumpTask = connection.pumpTo(*upStream) |
| 760 | .then([&upStream = *upStream](uint64_t) mutable { |
| 761 | return upStream.end(); |
| 762 | }).then([up = kj::mv(up), upStream = kj::mv(upStream)]() mutable -> kj::Promise<void> { |
| 763 | return kj::NEVER_DONE; |
| 764 | }); |
| 765 | |
| 766 | co_await pipeline.ignoreResult(); |
| 767 | co_await kj::joinPromisesFailFast(kj::arr(kj::mv(upPumpTask), kj::mv(downPumpTask))); |
| 768 | } |
| 769 | }; |
| 770 | |
| 771 | // `Fetcher` actually wants us to give it a factory that creates a new `WorkerInterface` for each |
| 772 | // request, so this is that. |
| 773 | class Container::TcpPortOutgoingFactory final: public Fetcher::OutgoingFactory { |
| 774 | public: |
| 775 | TcpPortOutgoingFactory(capnp::ByteStreamFactory& byteStreamFactory, |
| 776 | kj::EntropySource& entropySource, |
| 777 | const kj::HttpHeaderTable& headerTable, |
| 778 | rpc::Container::Port::Client port) |
| 779 | : byteStreamFactory(byteStreamFactory), |
| 780 | entropySource(entropySource), |
| 781 | headerTable(headerTable), |
| 782 | port(kj::mv(port)) {} |
| 783 | |
| 784 | kj::Own<WorkerInterface> newSingleUseClient(kj::Maybe<kj::String> cfStr) override { |
| 785 | // At present we have no use for `cfStr`. |
| 786 | return IoContext::current().getSubrequestNoChecks([&](auto& tracing, auto& channelFactory) { |
| 787 | return kj::heap<TcpPortWorkerInterface>(byteStreamFactory, entropySource, headerTable, port); |
| 788 | }, {.inHouse = false, .wrapMetrics = false}); |
| 789 | } |
| 790 | |
| 791 | private: |
| 792 | capnp::ByteStreamFactory& byteStreamFactory; |
| 793 | kj::EntropySource& entropySource; |
| 794 | const kj::HttpHeaderTable& headerTable; |
| 795 | rpc::Container::Port::Client port; |
| 796 | }; |
| 797 | |
| 798 | jsg::Ref<Fetcher> Container::getTcpPort(jsg::Lock& js, int port) { |
| 799 | JSG_REQUIRE(port > 0 && port < 65536, TypeError, "Invalid port number: ", port); |
| 800 | |
| 801 | auto req = rpcClient->getTcpPortRequest(capnp::MessageSize{4, 0}); |
| 802 | req.setPort(port); |
| 803 | |
| 804 | auto& ioctx = IoContext::current(); |
| 805 | |
| 806 | kj::Own<Fetcher::OutgoingFactory> factory = |
| 807 | kj::heap<TcpPortOutgoingFactory>(ioctx.getByteStreamFactory(), ioctx.getEntropySource(), |
| 808 | ioctx.getHeaderTable(), req.send().getPort()); |
| 809 | |
| 810 | return js.alloc<Fetcher>( |
| 811 | ioctx.addObject(kj::mv(factory)), Fetcher::RequiresHostAndProtocol::YES, true); |
| 812 | } |
| 813 | |
| 814 | } // namespace workerd::api |