File
Blob: src/workerd/server/pyodide.c++
| 1 | // Copyright (c) 2017-2022 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 | #include "pyodide.h" |
| 5 | |
| 6 | #include <workerd/api/pyodide/pyodide.h> |
| 7 | |
| 8 | #include <kj/array.h> |
| 9 | #include <kj/common.h> |
| 10 | #include <kj/compat/gzip.h> |
| 11 | #include <kj/compat/tls.h> |
| 12 | #include <kj/debug.h> |
| 13 | #include <kj/string.h> |
| 14 | |
| 15 | namespace workerd::server { |
| 16 | |
| 17 | // Helper functions for bundle file operations |
| 18 | kj::Path getPyodideBundleFileName(kj::StringPtr version) { |
| 19 | return kj::Path(kj::str("pyodide_", version, ".capnp.bin")); |
| 20 | } |
| 21 | |
| 22 | kj::Maybe<kj::Own<const kj::ReadableFile>> getPyodideBundleFile( |
| 23 | const kj::Maybe<kj::Own<const kj::Directory>>& maybeDir, kj::StringPtr version) { |
| 24 | KJ_IF_SOME(dir, maybeDir) { |
| 25 | kj::Path filename = getPyodideBundleFileName(version); |
| 26 | auto file = dir->tryOpenFile(filename); |
| 27 | |
| 28 | return file; |
| 29 | } |
| 30 | |
| 31 | return kj::none; |
| 32 | } |
| 33 | |
| 34 | void writePyodideBundleFileToDisk(const kj::Maybe<kj::Own<const kj::Directory>>& maybeDir, |
| 35 | kj::StringPtr version, |
| 36 | kj::ArrayPtr<byte> bytes) { |
| 37 | KJ_IF_SOME(dir, maybeDir) { |
| 38 | kj::Path filename = getPyodideBundleFileName(version); |
| 39 | auto replacer = dir->replaceFile(filename, kj::WriteMode::CREATE | kj::WriteMode::MODIFY); |
| 40 | |
| 41 | replacer->get().writeAll(bytes); |
| 42 | replacer->commit(); |
| 43 | } |
| 44 | } |
| 45 | |
| 46 | // Used to preload the Pyodide bundle during workerd startup |
| 47 | kj::Promise<kj::Maybe<jsg::Bundle::Reader>> fetchPyodideBundle( |
| 48 | const api::pyodide::PythonConfig& pyConfig, |
| 49 | kj::String version, |
| 50 | kj::Network& network, |
| 51 | kj::Timer& timer) { |
| 52 | if (pyConfig.pyodideBundleManager.getPyodideBundle(version) != kj::none) { |
| 53 | co_return pyConfig.pyodideBundleManager.getPyodideBundle(version); |
| 54 | } |
| 55 | |
| 56 | auto maybePyodideBundleFile = getPyodideBundleFile(pyConfig.pyodideDiskCacheRoot, version); |
| 57 | KJ_IF_SOME(pyodideBundleFile, maybePyodideBundleFile) { |
| 58 | auto body = pyodideBundleFile->readAllBytes(); |
| 59 | pyConfig.pyodideBundleManager.setPyodideBundleData(kj::str(version), kj::mv(body)); |
| 60 | co_return pyConfig.pyodideBundleManager.getPyodideBundle(version); |
| 61 | } |
| 62 | |
| 63 | if (version == "dev") { |
| 64 | // the "dev" version is special and indicates we're using the tip-of-tree version built for testing |
| 65 | // so we shouldn't fetch it from the internet, only check for its existence in the disk cache |
| 66 | co_return kj::none; |
| 67 | } |
| 68 | |
| 69 | kj::String url = |
| 70 | kj::str("https://pyodide-capnp-bin.edgeworker.net/pyodide_", version, ".capnp.bin"); |
| 71 | KJ_LOG(INFO, "Loading Pyodide bundle from internet", url); |
| 72 | kj::HttpHeaderTable table; |
| 73 | |
| 74 | kj::TlsContext::Options options; |
| 75 | options.useSystemTrustStore = true; |
| 76 | |
| 77 | kj::Own<kj::TlsContext> tls = kj::heap<kj::TlsContext>(kj::mv(options)); |
| 78 | auto tlsNetwork = tls->wrapNetwork(network); |
| 79 | auto client = kj::newHttpClient(timer, table, network, *tlsNetwork); |
| 80 | |
| 81 | kj::HttpHeaders headers(table); |
| 82 | |
| 83 | auto req = client->request(kj::HttpMethod::GET, url.asPtr(), headers); |
| 84 | |
| 85 | auto res = co_await req.response; |
| 86 | KJ_ASSERT(res.statusCode == 200, "Request for Pyodide bundle failed", url); |
| 87 | auto body = co_await res.body->readAllBytes(); |
| 88 | |
| 89 | writePyodideBundleFileToDisk(pyConfig.pyodideDiskCacheRoot, version, body); |
| 90 | |
| 91 | pyConfig.pyodideBundleManager.setPyodideBundleData(kj::str(version), kj::mv(body)); |
| 92 | |
| 93 | co_return pyConfig.pyodideBundleManager.getPyodideBundle(version); |
| 94 | } |
| 95 | |
| 96 | // Downloads a package with retry logic (up to 3 attempts with 5-second delays) |
| 97 | kj::Promise<kj::Maybe<kj::Array<byte>>> downloadPackageWithRetry(kj::HttpClient& client, |
| 98 | kj::Timer& timer, |
| 99 | kj::HttpHeaderTable& headerTable, |
| 100 | kj::StringPtr url, |
| 101 | kj::StringPtr path) { |
| 102 | constexpr uint retryLimit = 3; |
| 103 | kj::HttpHeaders headers(headerTable); |
| 104 | |
| 105 | for (uint retryCount = 0; retryCount < retryLimit; ++retryCount) { |
| 106 | if (retryCount > 0) { |
| 107 | // Sleep for 5 seconds before retrying |
| 108 | co_await timer.afterDelay(5 * kj::SECONDS); |
| 109 | KJ_LOG(INFO, "Retrying package download", path, "attempt", retryCount + 1, "of", retryLimit); |
| 110 | } |
| 111 | |
| 112 | try { |
| 113 | auto req = client.request(kj::HttpMethod::GET, url, headers); |
| 114 | auto res = co_await req.response; |
| 115 | |
| 116 | if (res.statusCode != 200) { |
| 117 | KJ_LOG(WARNING, "Failed to download package", path, res.statusCode, "attempt", |
| 118 | retryCount + 1, "of", retryLimit); |
| 119 | continue; // Try again in the next iteration |
| 120 | } |
| 121 | |
| 122 | // Request succeeded, read the body |
| 123 | co_return co_await res.body->readAllBytes(); |
| 124 | } catch (kj::Exception& e) { |
| 125 | if (retryCount + 1 >= retryLimit) { |
| 126 | // This was our last attempt |
| 127 | KJ_LOG(WARNING, "Failed to download package after all retry attempts", path, e, "attempts", |
| 128 | retryLimit); |
| 129 | } else { |
| 130 | KJ_LOG(WARNING, "Failed to download package", path, e, "attempt", retryCount + 1, "of", |
| 131 | retryLimit, "will retry"); |
| 132 | } |
| 133 | } |
| 134 | } |
| 135 | |
| 136 | co_return kj::none; // All retry attempts failed |
| 137 | } |
| 138 | |
| 139 | // Loads a single Python package, either from disk cache or by downloading it |
| 140 | kj::Promise<void> loadPyodidePackage(const api::pyodide::PythonConfig& pyConfig, |
| 141 | const api::pyodide::PyodidePackageManager& pyodidePackageManager, |
| 142 | kj::StringPtr packagesVersion, |
| 143 | kj::StringPtr filename, |
| 144 | kj::Network& network, |
| 145 | kj::Timer& timer) { |
| 146 | |
| 147 | auto path = kj::str("python-package-bucket/", packagesVersion, "/", filename); |
| 148 | // First check if we already have this package in memory |
| 149 | if (pyodidePackageManager.getPyodidePackage(path) != kj::none) { |
| 150 | co_return; |
| 151 | } |
| 152 | |
| 153 | // Then check disk cache |
| 154 | KJ_IF_SOME(diskCachePath, pyConfig.packageDiskCacheRoot) { |
| 155 | auto parsedPath = kj::Path::parse(filename); |
| 156 | if (diskCachePath->exists(parsedPath)) { |
| 157 | try { |
| 158 | auto file = diskCachePath->openFile(parsedPath); |
| 159 | auto blob = file->readAllBytes(); |
| 160 | |
| 161 | // Decompress the package |
| 162 | kj::ArrayInputStream ais(blob); |
| 163 | kj::GzipInputStream gzip(ais); |
| 164 | auto decompressed = gzip.readAllBytes(); |
| 165 | |
| 166 | // Store in memory |
| 167 | pyodidePackageManager.setPyodidePackageData(kj::str(path), kj::mv(decompressed)); |
| 168 | co_return; |
| 169 | } catch (kj::Exception& e) { |
| 170 | // Something went wrong while reading or processing the file |
| 171 | KJ_LOG(WARNING, "Failed to read or process package from disk cache", path, e); |
| 172 | } |
| 173 | } |
| 174 | } |
| 175 | |
| 176 | // Need to fetch from network |
| 177 | kj::HttpHeaderTable table; |
| 178 | kj::TlsContext::Options tlsOptions; |
| 179 | tlsOptions.useSystemTrustStore = true; |
| 180 | kj::Own<kj::TlsContext> tlsContext = kj::heap<kj::TlsContext>(kj::mv(tlsOptions)); |
| 181 | |
| 182 | auto tlsNetwork = tlsContext->wrapNetwork(network); |
| 183 | auto client = kj::newHttpClient(timer, table, network, *tlsNetwork); |
| 184 | |
| 185 | kj::String url = kj::str(api::pyodide::PYTHON_PACKAGES_URL, path); |
| 186 | |
| 187 | auto maybeBody = co_await downloadPackageWithRetry(*client, timer, table, url, path); |
| 188 | KJ_IF_SOME(body, maybeBody) { |
| 189 | // Successfully downloaded the package |
| 190 | // Save the compressed data to disk cache (if enabled) |
| 191 | KJ_IF_SOME(diskCachePath, pyConfig.packageDiskCacheRoot) { |
| 192 | try { |
| 193 | auto parsedPath = kj::Path::parse(path); |
| 194 | auto file = diskCachePath->openFile(parsedPath, |
| 195 | kj::WriteMode::CREATE | kj::WriteMode::MODIFY | kj::WriteMode::CREATE_PARENT); |
| 196 | file->writeAll(body); |
| 197 | } catch (kj::Exception& e) { |
| 198 | KJ_LOG(WARNING, "Failed to write package to disk cache", e); |
| 199 | } |
| 200 | } |
| 201 | |
| 202 | // Now decompress and store in memory |
| 203 | kj::ArrayInputStream ais(body); |
| 204 | kj::GzipInputStream gzip(ais); |
| 205 | auto decompressed = gzip.readAllBytes(); |
| 206 | |
| 207 | pyodidePackageManager.setPyodidePackageData(kj::str(path), kj::mv(decompressed)); |
| 208 | } else { |
| 209 | KJ_FAIL_ASSERT("Failed to download package after all retry attempts", path); |
| 210 | } |
| 211 | |
| 212 | co_return; |
| 213 | } |
| 214 | |
| 215 | kj::Promise<void> fetchPyodidePackages(const api::pyodide::PythonConfig& pyConfig, |
| 216 | const api::pyodide::PyodidePackageManager& pyodidePackageManager, |
| 217 | kj::ArrayPtr<kj::String> pythonRequirements, |
| 218 | workerd::PythonSnapshotRelease::Reader pythonSnapshotRelease, |
| 219 | kj::Network& network, |
| 220 | kj::Timer& timer) { |
| 221 | auto packagesVersion = pythonSnapshotRelease.getPackages(); |
| 222 | |
| 223 | auto pyodideLock = api::pyodide::getPyodideLock(pythonSnapshotRelease); |
| 224 | if (pyodideLock == kj::none) { |
| 225 | KJ_LOG(WARNING, "No lock file found for Python packages version", packagesVersion); |
| 226 | co_return; |
| 227 | } |
| 228 | |
| 229 | auto filenames = api::pyodide::getPythonPackageFiles( |
| 230 | KJ_ASSERT_NONNULL(pyodideLock), pythonRequirements, packagesVersion); |
| 231 | |
| 232 | kj::Vector<kj::Promise<void>> promises(filenames.size()); |
| 233 | for (const auto& filename: filenames) { |
| 234 | promises.add(loadPyodidePackage( |
| 235 | pyConfig, pyodidePackageManager, packagesVersion, filename, network, timer)); |
| 236 | } |
| 237 | |
| 238 | co_await kj::joinPromisesFailFast(promises.releaseAsArray()); |
| 239 | } |
| 240 | |
| 241 | } // namespace workerd::server |