Skip to content
File

Blob: src/workerd/api/container.c++

30.8 KB
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 
17namespace workerd::api {
18 
19namespace {
20 
21kj::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 
37void 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 
46kj::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 
53kj::Array<kj::byte> emptyByteArray() {
54 return kj::heapArray<kj::byte>(0);
55}
56 
57capnp::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 
67ExecOutput::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 
73jsg::JsArrayBuffer ExecOutput::getStdout(jsg::Lock& js) {
74 return jsg::JsArrayBuffer::create(js, stdoutBytes);
75}
76 
77jsg::JsArrayBuffer ExecOutput::getStderr(jsg::Lock& js) {
78 return jsg::JsArrayBuffer::create(js, stderrBytes);
79}
80 
81ExecProcess::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 
92jsg::Optional<jsg::Ref<WritableStream>> ExecProcess::getStdin() {
93 return stdinStream.map([](jsg::Ref<WritableStream>& stream) { return stream.addRef(); });
94}
95 
96jsg::Optional<jsg::Ref<ReadableStream>> ExecProcess::getStdout() {
97 return stdoutStream.map([](jsg::Ref<ReadableStream>& stream) { return stream.addRef(); });
98}
99 
100jsg::Optional<jsg::Ref<ReadableStream>> ExecProcess::getStderr() {
101 return stderrStream.map([](jsg::Ref<ReadableStream>& stream) { return stream.addRef(); });
102}
103 
104void 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 
126jsg::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 
141jsg::MemoizedIdentity<jsg::Promise<int>>& ExecProcess::getExitCode(jsg::Lock& js) {
142 ensureExitCodePromise(js);
143 return KJ_ASSERT_NONNULL(exitCodePromise);
144}
145 
146jsg::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 
192void 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 
203Container::Container(rpc::Container::Client rpcClient, bool running)
204 : rpcClient(IoContext::current().addObject(kj::heap(kj::mv(rpcClient)))),
205 running(running) {}
206 
207void 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 
287jsg::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 
310jsg::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 
342jsg::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 
371jsg::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 
381jsg::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 
396jsg::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 
414jsg::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 
427jsg::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 
556jsg::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 
571jsg::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 
597jsg::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 
608void 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.
622class 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.
773class 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 
798jsg::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