This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12299-309b15effe99423f73b6007cbe61b97bb6af7b24 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit af80a6704be4734b20a913f46f3838c3770688eb Author: SEZ <[email protected]> AuthorDate: Sun Oct 4 04:23:02 2026 +0000 [Improve][Zeta] Bound the size of REST log-content responses (#12299) --- docs/en/engines/zeta/rest-api-v1.md | 10 ++ docs/en/engines/zeta/rest-api-v2.md | 36 ++++- .../introduction/concepts/incompatible-changes.md | 20 +++ docs/zh/engines/zeta/rest-api-v1.md | 10 +- docs/zh/engines/zeta/rest-api-v2.md | 30 +++- .../introduction/concepts/incompatible-changes.md | 16 ++ .../apache/seatunnel/common/utils/FileUtils.java | 168 +++++++++++++++++++++ .../seatunnel/common/utils/FileUtilsTest.java | 158 +++++++++++++++++++ .../org/apache/seatunnel/engine/e2e/RestApiIT.java | 46 ++++++ .../config/YamlSeaTunnelDomConfigProcessor.java | 9 ++ .../engine/common/config/server/HttpConfig.java | 15 ++ .../common/config/server/ServerConfigOptions.java | 7 + .../config/YamlSeaTunnelConfigParserTest.java | 23 +++ .../src/test/resources/seatunnel.yaml | 3 +- .../engine/server/rest/LogContentReader.java | 80 ++++++++++ .../server/rest/RestHttpGetCommandProcessor.java | 8 +- .../engine/server/rest/service/LogService.java | 12 ++ .../engine/server/rest/servlet/LogBaseServlet.java | 17 ++- .../engine/server/rest/LogContentReaderTest.java | 124 +++++++++++++++ 19 files changed, 783 insertions(+), 9 deletions(-) diff --git a/docs/en/engines/zeta/rest-api-v1.md b/docs/en/engines/zeta/rest-api-v1.md index fff61c8f71..0ae25725b7 100644 --- a/docs/en/engines/zeta/rest-api-v1.md +++ b/docs/en/engines/zeta/rest-api-v1.md @@ -1007,6 +1007,13 @@ If you'd like to view the log list first, you can use a `GET` request to retriev The supported formats are `json` and `html`, with `html` as the default. +#### Response Size Limit + +Reading a log file returns at most `seatunnel.engine.http.log-response-max-size-mb` of content +(64 MB by default), exactly as the [v2 endpoint](rest-api-v2.md#log-response-size-limit) does. A +larger log file is represented by its tail and the response opens with a line saying so. Set the +option to `0` to restore unlimited reads. + #### Examples Retrieve logs for all nodes with the `jobId` of `733584788375666689`: `http://localhost:5801/hazelcast/rest/maps/logs/733584788375666689` @@ -1030,4 +1037,7 @@ Returns a list of logs from the requested node. To get a list of logs from the current node: `http://localhost:5801/hazelcast/rest/maps/log` To get the content of a log file: `http://localhost:5801/hazelcast/rest/maps/log/job-898380162133917698.log` +Log content is limited by `seatunnel.engine.http.log-response-max-size-mb` in the same way as the +all-node endpoint above. + </details> diff --git a/docs/en/engines/zeta/rest-api-v2.md b/docs/en/engines/zeta/rest-api-v2.md index 26382006fa..7541112e35 100644 --- a/docs/en/engines/zeta/rest-api-v2.md +++ b/docs/en/engines/zeta/rest-api-v2.md @@ -16,7 +16,7 @@ The v2 API and the Web UI are both served by the embedded Jetty server. Jetty st There are two different "default" sources that are easy to mix up: -- Code defaults: `enable-http = false`, `enable-https = false`, `port = 8080`, `context-path = ""`, `enable-dynamic-port = false`, `port-range = 100`, `upload-max-file-size-mb = 10`, `upload-max-request-size-mb = 10` +- Code defaults: `enable-http = false`, `enable-https = false`, `port = 8080`, `context-path = ""`, `enable-dynamic-port = false`, `port-range = 100`, `upload-max-file-size-mb = 10`, `upload-max-request-size-mb = 10`, `log-response-max-size-mb = 64` - The packaged `seatunnel.yaml` example: it already sets `enable-http: true` and `port: 8080` As a result, if you start SeaTunnel with the packaged configuration, the Web UI and REST API usually @@ -73,6 +73,7 @@ seatunnel: port: 8080 upload-max-file-size-mb: 10 upload-max-request-size-mb: 10 + log-response-max-size-mb: 64 ``` ## Web UI and Port 8080 Troubleshooting @@ -1450,6 +1451,36 @@ If you want to view the log list first, you can retrieve it via a `GET` request: Supported formats are `json` and `html`, with `html` as the default. +<a id="log-response-size-limit"></a> + +#### Response Size Limit + +The limit applies only to file-content responses, not log listings. Files are decoded as UTF-8, +including active logs and rotated files such as `seatunnel.log.*`. If your logging layout uses a +different platform charset, configure its `layout.charset` as `UTF-8`. The shipped Log4j2 example +rolls files at 100 MB, so the default 64 MB response cap can truncate even a rotated file. +The notice is additional to the capped file content. + +Reading a log file returns at most `seatunnel.engine.http.log-response-max-size-mb` of content +(64 MB by default). A log file larger than that is represented by its last +`log-response-max-size-mb` of content, because for a job that has been running for a long time the +end of the log is the part that explains what happened. + +A truncated response opens with a line naming the actual retained bytes and the file size captured +at the start of the same read, so that a partial log +is not mistaken for a complete one: + +``` +[SeaTunnel] Log truncated: returning 67108792 bytes from the tail of 3435973836 bytes (file size at read start). A partial first line is omitted when possible; an oversized single line returns a UTF-8-safe partial tail. Raise seatunnel.engine.http.log-response-max-size-mb, or set it to 0 for no limit, to return more. +``` + +The content itself starts at the first complete line after the cut, so the response is slightly +smaller than the limit. When a single line is longer than the limit there is no line boundary to +align to and the content starts at the first whole character instead. + +Set the option to `0` to restore unlimited reads - be aware that a single request for a +multi-gigabyte log file then has to fit in the node's heap. + #### Examples Retrieve logs for `jobId` `733584788375666689` across all nodes: `http://localhost:8080/logs/733584788375666689` @@ -1473,6 +1504,9 @@ Returns a list of logs from the requested node. To get a list of logs from the current node: `http://localhost:5801/log` To get the content of a log file: `http://localhost:5801/log/job-898380162133917698.log` +Log content is limited by `seatunnel.engine.http.log-response-max-size-mb` in the same way as the +all-node endpoint above. + </details> ------------------------------------------------------------------------------------------ diff --git a/docs/en/introduction/concepts/incompatible-changes.md b/docs/en/introduction/concepts/incompatible-changes.md index 5719f6fa40..ceae4b0e74 100644 --- a/docs/en/introduction/concepts/incompatible-changes.md +++ b/docs/en/introduction/concepts/incompatible-changes.md @@ -386,4 +386,24 @@ You need to check this document before you upgrade to related version. ### Engine Behavior Changes +- **Behavior change: the REST log-content endpoints return at most 64 MB by default** + - **Affected component**: `seatunnel-engine-server`, REST v2 endpoints `GET /logs/:file` and + `GET /log/:file` and their REST v1 equivalents `GET /hazelcast/rest/maps/logs/:file` and + `GET /hazelcast/rest/maps/log/:file`. + - **Description**: These endpoints read the requested log file whole, which materialises it on the + heap twice, so a single request for the log of a long-running streaming job could exhaust a + node's memory. The new `seatunnel.engine.http.log-response-max-size-mb` option caps how much is + read and defaults to `64`. A larger file is represented by its last `log-response-max-size-mb` + of UTF-8 content, aligned to a complete line when possible (or a partial tail of an oversized + line). The response notice names the actual retained bytes and the file-size snapshot. + - **Impact**: A cluster upgraded without editing `seatunnel.yaml` starts receiving the tail rather + than the whole of any log file above 64 MB, with status `200` as before. Anything that archives + logs through these endpoints - `curl .../logs/<job-id> > job.log`, or the log-analysis flow in + `docs/en/engines/zeta/log-analysis-with-ai.md` - keeps a partial file unless the limit is + raised. The truncation notice on the first line makes a partial response recognisable. + - **Migration Guide**: Set `log-response-max-size-mb: 0` under + `seatunnel.engine.http` to restore the previous unlimited reads, or raise it to a value that + covers the log sizes you collect. Leaving it at the default is recommended, since an unlimited + read of a multi-gigabyte log has to fit in the node's heap. + ### Dependency Upgrades diff --git a/docs/zh/engines/zeta/rest-api-v1.md b/docs/zh/engines/zeta/rest-api-v1.md index 207a478a47..93e50349f3 100644 --- a/docs/zh/engines/zeta/rest-api-v1.md +++ b/docs/zh/engines/zeta/rest-api-v1.md @@ -1009,12 +1009,18 @@ network: 当前支持的格式有`json`和`html`,默认为`html`。 +#### 响应大小限制 + +读取日志文件时最多返回 `seatunnel.engine.http.log-response-max-size-mb` 大小的内容(默认 64 MB), +规则与 [v2 接口](rest-api-v2.md#log-response-size-limit) 完全一致:超过限制的日志文件只返回末尾内容, +并在响应开头附上一行截断提示。把该项设为 `0` 可恢复不限制读取。 + #### 例子 获取所有节点jobId为`733584788375666689`的日志信息:`http://localhost:5801/hazelcast/rest/maps/logs/733584788375666689` 获取所有节点日志列表:`http://localhost:5801/hazelcast/rest/maps/logs` 获取所有节点日志列表以JSON格式返回:`http://localhost:5801/hazelcast/rest/maps/logs?format=json` -获取日志文件内容:`http://localhost:5801/hazelcast/rest/maps/logs/job-898380162133917698.log`` +获取日志文件内容:`http://localhost:5801/hazelcast/rest/maps/logs/job-898380162133917698.log` </details> @@ -1034,4 +1040,6 @@ network: 获取当前节点的日志列表:`http://localhost:5801/hazelcast/rest/maps/log` 获取日志文件内容:`http://localhost:5801/hazelcast/rest/maps/log/job-898380162133917698.log` +日志内容同样受 `seatunnel.engine.http.log-response-max-size-mb` 限制,规则与上面的全节点接口一致。 + </details> diff --git a/docs/zh/engines/zeta/rest-api-v2.md b/docs/zh/engines/zeta/rest-api-v2.md index 39d39928b9..0044f3c8e0 100644 --- a/docs/zh/engines/zeta/rest-api-v2.md +++ b/docs/zh/engines/zeta/rest-api-v2.md @@ -13,7 +13,7 @@ v2 版本的 API 和 Web UI 都由内嵌 Jetty 提供,与 v1 版本保持相 这里需要区分两个容易混淆的“默认值”来源: -- 代码默认值:`enable-http = false`、`enable-https = false`、`port = 8080`、`context-path = ""`、`enable-dynamic-port = false`、`port-range = 100`、`upload-max-file-size-mb = 10`、`upload-max-request-size-mb = 10` +- 代码默认值:`enable-http = false`、`enable-https = false`、`port = 8080`、`context-path = ""`、`enable-dynamic-port = false`、`port-range = 100`、`upload-max-file-size-mb = 10`、`upload-max-request-size-mb = 10`、`log-response-max-size-mb = 64` - 发行包自带的 `seatunnel.yaml` 示例:默认写入了 `enable-http: true` 和 `port: 8080` 因此,直接使用发行包自带配置启动时,Web UI 和 REST API 通常会监听 @@ -68,6 +68,7 @@ seatunnel: port: 8080 upload-max-file-size-mb: 10 upload-max-request-size-mb: 10 + log-response-max-size-mb: 64 ``` ## Web UI 与 8080 排查 @@ -1427,6 +1428,29 @@ curl --location 'http://127.0.0.1:8080/submit-job/upload?restoreMode=CHECKPOINT& 当前支持的格式有`json`和`html`,默认为`html`。 +<a id="log-response-size-limit"></a> + +#### 响应大小限制 + +该限制只作用于文件内容响应,不影响日志列表。活动日志和 `seatunnel.log.*` 等滚动日志均按 UTF-8 解码。 +如果日志布局使用其它平台编码,请将其 `layout.charset` 配置为 `UTF-8`。自带的 Log4j2 示例按 100 MB 滚动文件, +因此默认 64 MB 响应上限也可能截断滚动后的日志文件。截断提示本身不计入文件内容的大小上限。 + +读取日志文件时最多返回 `seatunnel.engine.http.log-response-max-size-mb` 大小的内容(默认 64 MB)。 +超过该限制的日志文件只返回末尾 `log-response-max-size-mb` 的内容——对长时间运行的作业来说,日志末尾 +才是解释问题的部分。 + +被截断的响应会以一行提示开头,写明实际保留的字节数和同一次读取开始时记录的文件大小,避免把不完整的日志当成完整日志: + +``` +[SeaTunnel] Log truncated: returning 67108792 bytes from the tail of 3435973836 bytes (file size at read start). A partial first line is omitted when possible; an oversized single line returns a UTF-8-safe partial tail. Raise seatunnel.engine.http.log-response-max-size-mb, or set it to 0 for no limit, to return more. +``` + +正文从截断点之后的第一个完整行开始,因此实际返回会略小于该限制;如果单行长度本身就超过限制,则没有可 +对齐的换行,正文从第一个完整字符开始。 + +把该项设为 `0` 可恢复不限制读取,但要注意此时单个请求需要把整个多 GB 的日志文件放进节点堆内存。 + #### 例子 @@ -1451,7 +1475,9 @@ curl --location 'http://127.0.0.1:8080/submit-job/upload?restoreMode=CHECKPOINT& #### 例子 获取当前节点的日志列表:`http://localhost:5801/log` -获取日志文件内容:`http://localhost:5801/log/job-898380162133917698.log`` +获取日志文件内容:`http://localhost:5801/log/job-898380162133917698.log` + +日志内容同样受 `seatunnel.engine.http.log-response-max-size-mb` 限制,规则与上面的全节点接口一致。 </details> diff --git a/docs/zh/introduction/concepts/incompatible-changes.md b/docs/zh/introduction/concepts/incompatible-changes.md index 17912d0650..55c85f1360 100644 --- a/docs/zh/introduction/concepts/incompatible-changes.md +++ b/docs/zh/introduction/concepts/incompatible-changes.md @@ -322,4 +322,20 @@ ### 引擎行为变更 +- **行为变更:REST 日志内容接口默认最多返回 64 MB** + - **受影响组件**:`seatunnel-engine-server`,REST v2 接口 `GET /logs/:file`、`GET /log/:file`, + 以及对应的 REST v1 接口 `GET /hazelcast/rest/maps/logs/:file`、`GET /hazelcast/rest/maps/log/:file`。 + - **说明**:这些接口原本会把整个日志文件读入内存,且会在堆上生成两份副本,因此对长时间运行的流作业 + 发起一次日志请求就可能耗尽节点内存。新增的 `seatunnel.engine.http.log-response-max-size-mb` + 选项限制单次读取的大小,默认值为 `64`。超过该限制的文件只返回末尾 `log-response-max-size-mb` + 的 UTF-8 内容,尽量从完整行开始;超长单行则保留部分末尾内容。响应开头的提示写明实际保留的字节数和文件大小快照。 + - **影响**:升级后未修改 `seatunnel.yaml` 的集群,对超过 64 MB 的日志文件将只得到末尾内容, + HTTP 状态码仍为 `200`。所有通过这些接口归档日志的用法——例如 + `curl .../logs/<job-id> > job.log`,或 `docs/zh/engines/zeta/log-analysis-with-ai.md` + 中的日志分析流程——在不调高限制的情况下都只会保存到部分内容。第一行的截断提示可以用来识别 + 响应是否完整。 + - **迁移指南**:在 `seatunnel.engine.http` 下设置 + `log-response-max-size-mb: 0` 可恢复此前的不限制读取,也可以把它调高到 + 足以覆盖需要收集的日志大小。建议保留默认值,因为不限制读取意味着多 GB 的日志需要完整放进节点堆内存。 + ### 依赖升级 diff --git a/seatunnel-common/src/main/java/org/apache/seatunnel/common/utils/FileUtils.java b/seatunnel-common/src/main/java/org/apache/seatunnel/common/utils/FileUtils.java index b228d1be1a..d5cb1a750b 100644 --- a/seatunnel-common/src/main/java/org/apache/seatunnel/common/utils/FileUtils.java +++ b/seatunnel-common/src/main/java/org/apache/seatunnel/common/utils/FileUtils.java @@ -24,17 +24,24 @@ import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException; import lombok.NonNull; import lombok.extern.slf4j.Slf4j; +import java.io.ByteArrayInputStream; import java.io.File; import java.io.FileNotFoundException; import java.io.FileOutputStream; import java.io.IOException; +import java.io.InputStreamReader; import java.io.PrintStream; +import java.io.Reader; import java.net.MalformedURLException; import java.net.URL; +import java.nio.ByteBuffer; +import java.nio.channels.SeekableByteChannel; +import java.nio.charset.StandardCharsets; import java.nio.file.FileVisitOption; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; +import java.nio.file.StandardOpenOption; import java.util.ArrayList; import java.util.Arrays; import java.util.List; @@ -45,6 +52,13 @@ import java.util.stream.Stream; @Slf4j public class FileUtils { + /** + * The largest tail {@link #readFileTailToStr(Path, long)} can return. A byte array cannot hold + * more than {@link Integer#MAX_VALUE} entries and some JVMs reserve a few of those for the + * array header, so a limit above this one cannot be honoured however much heap is available. + */ + public static final long MAX_TAIL_BYTES = Integer.MAX_VALUE - 8; + public static List<URL> searchJarFiles(@NonNull Path directory) throws IOException { if (!directory.toFile().exists()) { return new ArrayList<>(); @@ -75,6 +89,160 @@ public class FileUtils { } } + /** + * Reads a file, keeping at most {@code maxBytes} bytes from the end of it. + * + * <p>Reading a file whole materialises it twice on the heap, once as a byte array and once as a + * string. For files that can grow without bound - engine log files being the case this was + * written for - that turns a single read into a node-wide memory problem. When the file is + * larger than the limit its tail is returned instead, the tail being the part that matters when + * diagnosing a failure. + * + * <p>The tail starts at the first line break after the cut point. If a single line exceeds the + * limit, a partial tail is returned without splitting its first multi-byte UTF-8 character. + * Content is decoded as UTF-8 rather than with the platform default charset used by {@link + * #readFileToStr(Path)} - aligning on character boundaries is only meaningful against a known + * encoding, and a file whose encoding changed as it grew past the limit would be worse than one + * that is consistently wrong. + * + * <p>The limit that actually applies is {@link #effectiveTailLimit(long)} rather than {@code + * maxBytes} itself, so a caller asking for more than a byte array can hold still gets a bounded + * read instead of an {@link OutOfMemoryError}. + * + * @param path file to read + * @param maxBytes maximum number of bytes to keep from the end; a value <= 0 means unlimited + * @return the whole file, or its tail when the file is larger than the effective limit + */ + public static String readFileTailToStr(Path path, long maxBytes) { + return readFileTail(path, maxBytes).getContent(); + } + + /** + * Reads UTF-8 content and truncation metadata from one open file and one size snapshot. + * Positive limits also bound reads of files that grow or are replaced while being read. + */ + public static FileTail readFileTail(Path path, long maxBytes) { + try { + if (maxBytes <= 0) { + byte[] bytes = Files.readAllBytes(path); + return new FileTail(bytes, 0, bytes.length, bytes.length, false); + } + long keep = effectiveTailLimit(maxBytes); + try (SeekableByteChannel channel = + Files.newByteChannel(path, StandardOpenOption.READ)) { + long size = channel.size(); + boolean truncated = size > keep; + ByteBuffer buffer = ByteBuffer.allocate((int) Math.min(size, keep)); + channel.position(truncated ? size - keep : 0); + while (buffer.hasRemaining() && channel.read(buffer) > 0) { + // Never read beyond the snapshot, even if the file grows during this read. + } + byte[] tail = buffer.array(); + int length = buffer.position(); + int start = truncated ? lineStartOffset(tail, length) : 0; + return new FileTail(tail, start, length - start, size, truncated); + } + } catch (IOException e) { + throw CommonError.fileOperationFailed("SeaTunnel", "read", path.toString(), e); + } + } + + /** + * Returns the number of bytes {@link #readFileTailToStr(Path, long)} keeps for the given + * positive limit, which is {@code maxBytes} clamped to {@link #MAX_TAIL_BYTES}. + * + * <p>Comparing the file size against this instead of against {@code maxBytes} is what keeps the + * read bounded. A limit above {@link #MAX_TAIL_BYTES} cannot be honoured, so a file sized + * between the two has to be read as a tail rather than whole - reading it whole would fail with + * {@code OutOfMemoryError: Required array size too large}, which is the very failure the tail + * read exists to prevent. + */ + public static long effectiveTailLimit(long maxBytes) { + return Math.min(maxBytes, MAX_TAIL_BYTES); + } + + /** An immutable read result whose metadata stays valid after the file is rotated or removed. */ + public static final class FileTail { + private final byte[] bytes; + private final int offset; + private final int length; + private final long fileSize; + private final boolean truncated; + + private FileTail(byte[] bytes, int offset, int length, long fileSize, boolean truncated) { + this.bytes = bytes; + this.offset = offset; + this.length = length; + this.fileSize = fileSize; + this.truncated = truncated; + } + + /** Returns the size observed at the start of the read, before any later rotation. */ + public long getFileSize() { + return fileSize; + } + + /** Returns retained file bytes after line/character alignment, excluding any prefix. */ + public int getReturnedBytes() { + return length; + } + + /** Reports whether the file-size snapshot exceeded the effective positive read limit. */ + public boolean isTruncated() { + return truncated; + } + + /** Decodes only the retained window, with no intermediate byte-array copy. */ + public String getContent() { + return new String(bytes, offset, length, StandardCharsets.UTF_8); + } + + /** + * Decodes into the final response builder after its prefix, avoiding a full intermediate + * content String when a caller adds a truncation notice. + */ + public String getContentWithPrefix(String prefix) throws IOException { + StringBuilder response = + new StringBuilder( + (int) Math.min(MAX_TAIL_BYTES, (long) prefix.length() + length)); + response.append(prefix); + try (Reader reader = + new InputStreamReader( + new ByteArrayInputStream(bytes, offset, length), + StandardCharsets.UTF_8)) { + char[] buffer = new char[8192]; + int count; + while ((count = reader.read(buffer)) != -1) { + response.append(buffer, 0, count); + } + } + return response.toString(); + } + } + + /** + * Returns the offset of the first byte to keep in a retained tail window, which is the start of + * the first complete line in it. + */ + private static int lineStartOffset(byte[] bytes, int length) { + // A '\n' on the final byte is the terminator of the line before it, not the start of + // another line, so it is not a boundary to align to. Stopping short of it is what keeps a + // window holding a single newline-terminated line - the shape produced by a log whose last + // entry is a large stack trace - from being reported as empty. + for (int i = 0; i < length - 1; i++) { + if (bytes[i] == '\n') { + return i + 1; + } + } + // A single line longer than the limit leaves no boundary to align to, so drop just the + // leading UTF-8 continuation bytes to avoid starting in the middle of a character. + int start = 0; + while (start < length && (bytes[start] & 0xC0) == 0x80) { + start++; + } + return start; + } + public static void writeStringToFile(String filePath, String str) { PrintStream ps = null; try { diff --git a/seatunnel-common/src/test/java/org/apache/seatunnel/common/utils/FileUtilsTest.java b/seatunnel-common/src/test/java/org/apache/seatunnel/common/utils/FileUtilsTest.java index cfc8eb6595..0ac8b9629d 100644 --- a/seatunnel-common/src/test/java/org/apache/seatunnel/common/utils/FileUtilsTest.java +++ b/seatunnel-common/src/test/java/org/apache/seatunnel/common/utils/FileUtilsTest.java @@ -21,6 +21,7 @@ import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; import lombok.NonNull; @@ -28,11 +29,17 @@ import java.io.BufferedWriter; import java.io.File; import java.io.FileWriter; import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; import java.nio.file.NoSuchFileException; import java.nio.file.Path; import java.nio.file.Paths; public class FileUtilsTest { + + /** The character a UTF-8 decoder emits when it is handed a partial character. */ + private static final String REPLACEMENT_CHAR = "\uFFFD"; + @Test public void testGetFileLineNumber() throws Exception { String filePath = "/tmp/test/file_utils/file1.txt"; @@ -145,4 +152,155 @@ public class FileUtilsTest { "", FileUtils.readFileToStr(Paths.get("/tmp/newfolder/newfolder2/newfolde3/test.txt"))); } + + @Test + void readFileTailToStrReturnsWholeFileWhenItFitsTheLimit(@TempDir Path tempDir) + throws IOException { + Path file = tempDir.resolve("small.log"); + String content = "first line\nsecond line\n"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + long size = Files.size(file); + + Assertions.assertEquals(content, FileUtils.readFileTailToStr(file, size + 1)); + // The limit is inclusive, so a file of exactly the limit is still returned whole. + Assertions.assertEquals(content, FileUtils.readFileTailToStr(file, size)); + } + + @Test + void readFileTailToStrTreatsNonPositiveLimitAsUnlimited(@TempDir Path tempDir) + throws IOException { + Path file = tempDir.resolve("unlimited.log"); + String content = "first line\nsecond line\n"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + + Assertions.assertEquals(content, FileUtils.readFileTailToStr(file, 0)); + Assertions.assertEquals(content, FileUtils.readFileTailToStr(file, -1)); + } + + @Test + void readFileTailToStrHandlesEmptyFile(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("empty.log"); + Files.write(file, new byte[0]); + + Assertions.assertEquals("", FileUtils.readFileTailToStr(file, 16)); + } + + @Test + void readFileTailToStrKeepsTheTailAlignedToALineBoundary(@TempDir Path tempDir) + throws IOException { + Path file = tempDir.resolve("aligned.log"); + StringBuilder content = new StringBuilder(); + for (int i = 0; i < 100; i++) { + content.append("line ").append(i).append('\n'); + } + Files.write(file, content.toString().getBytes(StandardCharsets.UTF_8)); + + String tail = FileUtils.readFileTailToStr(file, 40); + + Assertions.assertTrue(content.toString().endsWith(tail), "tail must be a suffix"); + Assertions.assertTrue(tail.length() < 40, "tail must be smaller than the limit"); + Assertions.assertTrue(tail.startsWith("line "), "tail must start at a line boundary"); + Assertions.assertTrue(tail.endsWith("line 99\n"), "tail must reach the end of the file"); + } + + @Test + void readFileTailToStrNeverSplitsAMultiByteCharacter(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("cjk.log"); + StringBuilder content = new StringBuilder(); + for (int i = 0; i < 200; i++) { + content.append("第").append(i).append("行日志内容").append('\n'); + } + Files.write(file, content.toString().getBytes(StandardCharsets.UTF_8)); + Assertions.assertTrue(Files.size(file) > 400, "test data must exceed the limits swept"); + + // Every character here is three bytes wide, so sweeping the limit guarantees that some of + // these cut points land inside a character. + for (long maxBytes = 100; maxBytes <= 400; maxBytes++) { + String tail = FileUtils.readFileTailToStr(file, maxBytes); + Assertions.assertFalse( + tail.contains(REPLACEMENT_CHAR), "character split at maxBytes=" + maxBytes); + Assertions.assertTrue( + content.toString().endsWith(tail), "not a suffix at maxBytes=" + maxBytes); + Assertions.assertTrue( + tail.startsWith("第"), "not aligned to a line at maxBytes=" + maxBytes); + } + } + + @Test + void readFileTailToStrHandlesALineLongerThanTheLimit(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("single-line.log"); + StringBuilder content = new StringBuilder(); + for (int i = 0; i < 500; i++) { + content.append("啊"); + } + Files.write(file, content.toString().getBytes(StandardCharsets.UTF_8)); + + // There is no line break to align to. 100 is not a multiple of the 3-byte character width, + // so the raw slice would begin on a continuation byte. + String tail = FileUtils.readFileTailToStr(file, 100); + + Assertions.assertFalse( + tail.contains(REPLACEMENT_CHAR), "character split without a line boundary"); + Assertions.assertTrue(content.toString().endsWith(tail), "tail must be a suffix"); + Assertions.assertEquals(33, tail.length()); + } + + @Test + void readFileTailToStrKeepsALineTerminatedAtTheEndOfTheWindow(@TempDir Path tempDir) + throws IOException { + Path file = tempDir.resolve("trailing-newline.log"); + // A log whose last entry is longer than the limit - a stack trace, say - leaves a window + // whose only line break is the file's own trailing one. That break terminates the retained + // line rather than starting a new one, so aligning to it would leave nothing to return. + String content = "earlier line\nabcdefghijklmnopqrstuvwxyz\n"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + + String tail = FileUtils.readFileTailToStr(file, 10); + + Assertions.assertEquals("rstuvwxyz\n", tail); + } + + @Test + void readFileTailToStrAlignsOnABreakThatIsNotTheLastByte(@TempDir Path tempDir) + throws IOException { + Path file = tempDir.resolve("no-trailing-newline.log"); + String content = "earlier line\nlast line without a terminator"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + + // The window is "line\nlast line without a terminator", so there is a real boundary to + // align to and the partial word before it has to be dropped. + String tail = FileUtils.readFileTailToStr(file, 35); + + Assertions.assertEquals("last line without a terminator", tail); + } + + @Test + void readFileTailToStrDecodesUtf8WhateverThePlatformCharsetIs(@TempDir Path tempDir) + throws IOException { + Path file = tempDir.resolve("charset.log"); + String content = "日志内容\n"; + byte[] utf8 = content.getBytes(StandardCharsets.UTF_8); + Files.write(file, utf8); + + // Under the limit, at the limit and unlimited all take a different branch, and all three + // have to agree with the over-limit branch about the encoding - otherwise the same file + // would change encoding as it grew. + Assertions.assertEquals(content, FileUtils.readFileTailToStr(file, utf8.length + 1)); + Assertions.assertEquals(content, FileUtils.readFileTailToStr(file, utf8.length)); + Assertions.assertEquals(content, FileUtils.readFileTailToStr(file, 0)); + } + + @Test + void effectiveTailLimitClampsALimitTooLargeToBeRead() { + Assertions.assertEquals(64L * 1024 * 1024, FileUtils.effectiveTailLimit(64L * 1024 * 1024)); + Assertions.assertEquals( + FileUtils.MAX_TAIL_BYTES, FileUtils.effectiveTailLimit(FileUtils.MAX_TAIL_BYTES)); + // A limit a byte array cannot hold has to be clamped rather than honoured: a file sized + // between the clamp and the configured limit must still be read as a tail, because reading + // it whole would fail with "Required array size too large" however much heap is available. + Assertions.assertEquals( + FileUtils.MAX_TAIL_BYTES, FileUtils.effectiveTailLimit(8L * 1024 * 1024 * 1024)); + Assertions.assertEquals( + FileUtils.MAX_TAIL_BYTES, FileUtils.effectiveTailLimit(Long.MAX_VALUE)); + } } diff --git a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java index 4e67a4f633..381fe8cabf 100644 --- a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java +++ b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java @@ -51,6 +51,8 @@ import io.restassured.common.mapper.TypeRef; import lombok.extern.slf4j.Slf4j; import java.lang.reflect.Field; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.Paths; import java.util.ArrayList; @@ -115,6 +117,7 @@ public class RestApiIT { node1Config.getEngineConfig().getHttpConfig().setPort(8080); node1Config.getEngineConfig().getHttpConfig().setEnabled(true); node1Config.getEngineConfig().getHttpConfig().setEnableDynamicPort(true); + node1Config.getEngineConfig().getHttpConfig().setLogResponseMaxSizeMb(1); node1Config.getHazelcastConfig().setClusterName(testClusterName); node1Config.getEngineConfig().getSlotServiceConfig().setDynamicSlot(false); node1Config.getEngineConfig().getSlotServiceConfig().setSlotNum(20); @@ -173,6 +176,49 @@ public class RestApiIT { ports.put(node2.getCluster().getLocalMember().getAddress().getPort(), httpPort(node2)); } + @Test + public void testLogResponseLimitIsAppliedToV1AndV2() throws Exception { + Path directory = Paths.get(new LogService(node1.node.getNodeEngine()).getLogPath()); + Path file = Files.createTempFile(directory, "rest-response-limit-", ".log"); + StringBuilder content = new StringBuilder("FIRST-RECORD-MUST-BE-TRUNCATED\n"); + for (int i = 0; i < 40000; i++) { + content.append("日志0123456789abcdefghijklmnopqrstuvwxyz\n"); + } + content.append("END-OF-LOG\n"); + String original = content.toString(); + int limit = 1024 * 1024; + Files.write(file, original.getBytes(StandardCharsets.UTF_8)); + Assertions.assertTrue(Files.size(file) > limit); + String v1 = + HOST + node1.getCluster().getLocalMember().getAddress().getPort() + CONTEXT_PATH; + String v2 = buildHttpBaseUrl(httpPort(node1)); + try { + for (String base : Arrays.asList(v1, v2)) { + for (String endpoint : + Arrays.asList(RestConstant.REST_URL_LOG, RestConstant.REST_URL_LOGS)) { + String response = + given().get(base + endpoint + "/" + file.getFileName()) + .then() + .statusCode(200) + .extract() + .asString(); + String[] parts = response.split("\n", 2); + Assertions.assertTrue(parts[0].startsWith("[SeaTunnel] Log truncated:")); + Assertions.assertEquals(2, parts.length); + String tail = parts[1]; + Assertions.assertTrue(tail.getBytes(StandardCharsets.UTF_8).length <= limit); + Assertions.assertTrue(original.endsWith(tail)); + Assertions.assertTrue(tail.endsWith("END-OF-LOG\n")); + Assertions.assertTrue(tail.contains("日志")); + Assertions.assertFalse(tail.contains("\uFFFD")); + Assertions.assertFalse(tail.contains("FIRST-RECORD-MUST-BE-TRUNCATED")); + } + } + } finally { + Files.deleteIfExists(file); + } + } + @Test public void testGetLog() { Arrays.asList(node2, node1) diff --git a/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelDomConfigProcessor.java b/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelDomConfigProcessor.java index 1d3218bb0d..94f34f6dc3 100644 --- a/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelDomConfigProcessor.java +++ b/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelDomConfigProcessor.java @@ -685,6 +685,15 @@ public class YamlSeaTunnelDomConfigProcessor extends AbstractDomConfigProcessor .UPLOAD_MAX_REQUEST_SIZE_MB .key(), getTextContent(node))); + } else if (ServerConfigOptions.MasterServerConfigOptions.LOG_RESPONSE_MAX_SIZE_MB + .key() + .equals(name)) { + httpConfig.setLogResponseMaxSizeMb( + getIntegerValue( + ServerConfigOptions.MasterServerConfigOptions + .LOG_RESPONSE_MAX_SIZE_MB + .key(), + getTextContent(node))); } else { LOGGER.warning("Unrecognized element: " + name); } diff --git a/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/HttpConfig.java b/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/HttpConfig.java index 36e860dfac..6971c465e3 100644 --- a/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/HttpConfig.java +++ b/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/HttpConfig.java @@ -86,6 +86,21 @@ public class HttpConfig implements Serializable { private int uploadMaxRequestSizeMb = ServerConfigOptions.MasterServerConfigOptions.UPLOAD_MAX_REQUEST_SIZE_MB.defaultValue(); + /** The maximum size in MB of a log file returned by the log endpoints. */ + private int logResponseMaxSizeMb = + ServerConfigOptions.MasterServerConfigOptions.LOG_RESPONSE_MAX_SIZE_MB.defaultValue(); + + /** + * Returns {@link #logResponseMaxSizeMb} as a byte count, or -1 when the log endpoints are + * configured to return content of unlimited size. + * + * <p>Both the v1 and the v2 log endpoint read their cap from here, so the two cannot end up + * disagreeing about what the option means. + */ + public long getLogResponseMaxSizeBytes() { + return logResponseMaxSizeMb <= 0 ? -1L : logResponseMaxSizeMb * 1024L * 1024L; + } + public void setPort(int port) { checkPositive(port, ServerConfigOptions.MasterServerConfigOptions.HTTP + " must be > 0"); this.port = port; diff --git a/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/ServerConfigOptions.java b/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/ServerConfigOptions.java index d3a23c5c22..16971f3e85 100644 --- a/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/ServerConfigOptions.java +++ b/seatunnel-engine/seatunnel-engine-common/src/main/java/org/apache/seatunnel/engine/common/config/server/ServerConfigOptions.java @@ -397,6 +397,13 @@ public class ServerConfigOptions { .withDescription( "The maximum total size in MB of a multipart request sent to the http server. A value <= 0 means unlimited."); + public static final Option<Integer> LOG_RESPONSE_MAX_SIZE_MB = + Options.key("log-response-max-size-mb") + .intType() + .defaultValue(64) + .withDescription( + "The maximum size in MB of a log file returned by the log endpoints. Larger files are truncated to their last log-response-max-size-mb of content. A value <= 0 means unlimited."); + public static final Option<HttpConfig> HTTP = Options.key("http") .type(new TypeReference<HttpConfig>() {}) diff --git a/seatunnel-engine/seatunnel-engine-common/src/test/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelConfigParserTest.java b/seatunnel-engine/seatunnel-engine-common/src/test/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelConfigParserTest.java index 0c8bd19c4a..828d2d3879 100644 --- a/seatunnel-engine/seatunnel-engine-common/src/test/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelConfigParserTest.java +++ b/seatunnel-engine/seatunnel-engine-common/src/test/java/org/apache/seatunnel/engine/common/config/YamlSeaTunnelConfigParserTest.java @@ -18,6 +18,7 @@ package org.apache.seatunnel.engine.common.config; import org.apache.seatunnel.common.utils.ReflectionUtils; +import org.apache.seatunnel.engine.common.config.server.HttpConfig; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -81,12 +82,34 @@ public class YamlSeaTunnelConfigParserTest { Assertions.assertEquals(8080, config.getEngineConfig().getHttpConfig().getPort()); Assertions.assertEquals(200, config.getEngineConfig().getHttpConfig().getPortRange()); Assertions.assertEquals(8443, config.getEngineConfig().getHttpConfig().getHttpsPort()); + // An http option without a matching branch in parseHttpConfig is silently dropped with an + // "Unrecognized element" warning, so parsing it is worth asserting explicitly. + Assertions.assertEquals( + 32, config.getEngineConfig().getHttpConfig().getLogResponseMaxSizeMb()); + Assertions.assertEquals( + 32L * 1024 * 1024, + config.getEngineConfig().getHttpConfig().getLogResponseMaxSizeBytes()); Assertions.assertEquals( 30, config.getEngineConfig().getCoordinatorServiceConfig().getCoreThreadNum()); Assertions.assertEquals( 1000, config.getEngineConfig().getCoordinatorServiceConfig().getMaxThreadNum()); } + @Test + public void testLogResponseLimitByteConversion() { + HttpConfig httpConfig = new HttpConfig(); + Assertions.assertEquals(64L * 1024 * 1024, httpConfig.getLogResponseMaxSizeBytes()); + + for (int unlimited : new int[] {0, -1, Integer.MIN_VALUE}) { + httpConfig.setLogResponseMaxSizeMb(unlimited); + Assertions.assertEquals(-1L, httpConfig.getLogResponseMaxSizeBytes()); + } + + httpConfig.setLogResponseMaxSizeMb(Integer.MAX_VALUE); + Assertions.assertEquals( + Integer.MAX_VALUE * 1024L * 1024L, httpConfig.getLogResponseMaxSizeBytes()); + } + @Test public void testCustomizeClientConfig() throws IOException { YamlClientConfigBuilder yamlClientConfigBuilder = diff --git a/seatunnel-engine/seatunnel-engine-common/src/test/resources/seatunnel.yaml b/seatunnel-engine/seatunnel-engine-common/src/test/resources/seatunnel.yaml index da6b831f38..3775207470 100644 --- a/seatunnel-engine/seatunnel-engine-common/src/test/resources/seatunnel.yaml +++ b/seatunnel-engine/seatunnel-engine-common/src/test/resources/seatunnel.yaml @@ -42,4 +42,5 @@ seatunnel: enable-http: true port: 8080 enable-dynamic-port: true - port-range: 200 \ No newline at end of file + port-range: 200 + log-response-max-size-mb: 32 \ No newline at end of file diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/LogContentReader.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/LogContentReader.java new file mode 100644 index 0000000000..54a0b1efb9 --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/LogContentReader.java @@ -0,0 +1,80 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.engine.server.rest; + +import org.apache.seatunnel.common.utils.FileUtils; + +import java.io.IOException; +import java.nio.file.Path; +import java.util.Locale; + +/** + * Reads the log-file content the REST log endpoints return, bounded by {@code + * log-response-max-size-mb}. + * + * <p>Shared by the v1 handler and the v2 servlets so that the two describe a truncated log the same + * way. + */ +public final class LogContentReader { + + /** + * Opening line of a truncated response. Written as a complete line so that a consumer reading + * the body line by line can recognise and drop it. + */ + private static final String TRUNCATION_NOTICE_FORMAT = + "[SeaTunnel] Log truncated: returning %d bytes from the tail of %d bytes" + + " (file size at read start). A partial first line is omitted when possible;" + + " an oversized single line returns a UTF-8-safe partial tail." + + " Raise seatunnel.engine.http.log-response-max-size-mb, or set" + + " it to 0 for no limit, to return more.\n"; + + private LogContentReader() {} + + /** + * Reads the content to return for a log file. + * + * <p>A positive {@code maxBytes} bounds the file content read into memory. A larger file is + * represented by its tail. + * + * <p>A truncated response opens with a notice naming the retained bytes and the file size + * captured by the same read. The tail starts at a complete line when possible, or at a UTF-8 + * character boundary for an oversized single line. The notice lets readers and archiving + * scripts distinguish the retained tail from a complete log. + * + * @param path canonical path of the log file, already resolved against the log directory + * @param maxBytes cap from {@code HttpConfig#getLogResponseMaxSizeBytes()}; <= 0 means no limit + * @return the log content, opening with a truncation notice when the file exceeded the cap + * @throws IOException if the response content cannot be decoded + */ + public static String read(Path path, long maxBytes) throws IOException { + return read(FileUtils.readFileTail(path, maxBytes)); + } + + static String read(FileUtils.FileTail tail) throws IOException { + if (!tail.isTruncated()) { + return tail.getContent(); + } + String notice = + String.format( + Locale.ROOT, + TRUNCATION_NOTICE_FORMAT, + tail.getReturnedBytes(), + tail.getFileSize()); + return tail.getContentWithPrefix(notice); + } +} diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java index 40dce8bdd2..f98705014d 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/RestHttpGetCommandProcessor.java @@ -20,7 +20,6 @@ package org.apache.seatunnel.engine.server.rest; import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils; import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException; -import org.apache.seatunnel.common.utils.FileUtils; import org.apache.seatunnel.common.utils.JsonUtils; import org.apache.seatunnel.engine.server.NodeExtension; import org.apache.seatunnel.engine.server.log.FormatType; @@ -395,6 +394,9 @@ public class RestHttpGetCommandProcessor extends HttpCommandProcessor<HttpGetCom * <p>The requested log file is resolved to its canonical path before reading so that relative * segments and symbolic links cannot escape the canonical log directory. * + * <p>A positive {@code log-response-max-size-mb} bounds the file content read into memory. A + * larger file is represented by its tail, and the response opens with a truncation notice. + * * @param httpGetCommand command used to send the HTTP response * @param logPath configured log directory * @param logName requested log file name from the request URI @@ -413,7 +415,9 @@ public class RestHttpGetCommandProcessor extends HttpCommandProcessor<HttpGetCom logName, canonicalFilePath, canonicalLogDir)); return; } - String logContent = FileUtils.readFileToStr(new File(canonicalFilePath).toPath()); + String logContent = + LogContentReader.read( + new File(canonicalFilePath).toPath(), logService.maxLogResponseBytes()); this.prepareResponse(httpGetCommand, logContent); } catch (IOException e) { httpGetCommand.send400(); diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LogService.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LogService.java index 784f23f845..c03111c2c0 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LogService.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LogService.java @@ -48,6 +48,18 @@ public class LogService extends BaseLogService { super(nodeEngine); } + /** + * Returns the cap in bytes on the content of a single log response, translated from {@code + * log-response-max-size-mb}. A value <= 0 means unlimited. + */ + public long maxLogResponseBytes() { + return getSeaTunnelServer(false) + .getSeaTunnelConfig() + .getEngineConfig() + .getHttpConfig() + .getLogResponseMaxSizeBytes(); + } + public List<String> allLogName() { String logPath = getLogPath(); List<File> logFileList = FileUtils.listFile(logPath); diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/LogBaseServlet.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/LogBaseServlet.java index f79ca21224..a926421599 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/LogBaseServlet.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/LogBaseServlet.java @@ -20,7 +20,7 @@ package org.apache.seatunnel.engine.server.rest.servlet; import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils; import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException; -import org.apache.seatunnel.common.utils.FileUtils; +import org.apache.seatunnel.engine.server.rest.LogContentReader; import com.hazelcast.spi.impl.NodeEngineImpl; import lombok.extern.slf4j.Slf4j; @@ -43,6 +43,9 @@ public class LogBaseServlet extends BaseServlet { * <p>The requested log file is resolved to its canonical path before reading so that relative * segments and symbolic links cannot escape the canonical log directory. * + * <p>A positive {@code log-response-max-size-mb} bounds the file content read into memory. A + * larger file is represented by its tail, and the response opens with a truncation notice. + * * @param resp response used to return status and log content * @param logPath configured log directory * @param logName requested log file name from the request URI @@ -68,7 +71,9 @@ public class LogBaseServlet extends BaseServlet { canonicalLogDir); return; } - String logContent = FileUtils.readFileToStr(new File(canonicalFilePath).toPath()); + String logContent = + LogContentReader.read( + new File(canonicalFilePath).toPath(), maxLogResponseBytes()); write(resp, logContent); } catch (IOException e) { resp.setStatus(HttpServletResponse.SC_BAD_REQUEST); @@ -78,4 +83,12 @@ public class LogBaseServlet extends BaseServlet { log.warn(String.format("Log file content is empty, get log path : %s", logFilePath)); } } + + private long maxLogResponseBytes() { + return getSeaTunnelServer(false) + .getSeaTunnelConfig() + .getEngineConfig() + .getHttpConfig() + .getLogResponseMaxSizeBytes(); + } } diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/LogContentReaderTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/LogContentReaderTest.java new file mode 100644 index 0000000000..79a1283f39 --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/LogContentReaderTest.java @@ -0,0 +1,124 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.engine.server.rest; + +import org.apache.seatunnel.common.utils.FileUtils; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.StandardOpenOption; + +class LogContentReaderTest { + + @Test + void readReturnsAFileWithinTheLimitUnchanged(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("small.log"); + String content = "first line\nsecond line\n"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + + Assertions.assertEquals(content, LogContentReader.read(file, 1024)); + } + + @Test + void readReturnsAFileUnchangedWhenThereIsNoLimit(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("unlimited.log"); + String content = "first line\nsecond line\n"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + + Assertions.assertEquals(content, LogContentReader.read(file, -1)); + Assertions.assertEquals(content, LogContentReader.read(file, 0)); + } + + @Test + void readAnnouncesThatALargeFileWasTruncated(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("large.log"); + StringBuilder content = new StringBuilder(); + for (int i = 0; i < 100; i++) { + content.append("line ").append(i).append('\n'); + } + Files.write(file, content.toString().getBytes(StandardCharsets.UTF_8)); + long size = Files.size(file); + + String response = LogContentReader.read(file, 40); + + String[] lines = response.split("\n", 2); + Assertions.assertTrue( + lines[0].startsWith("[SeaTunnel] Log truncated:"), + "a truncated response has to say so on its first line, got: " + lines[0]); + Assertions.assertTrue( + lines[0].contains( + "returning " + + lines[1].getBytes(StandardCharsets.UTF_8).length + + " bytes from the tail of " + + size), + "the notice must report actual retained bytes, got: " + lines[0]); + Assertions.assertTrue( + lines[0].contains("log-response-max-size-mb"), + "the notice has to name the option that controls the limit, got: " + lines[0]); + // Everything after the notice is the log itself, unchanged. + Assertions.assertTrue( + content.toString().endsWith(lines[1]), "the body has to be the file's own tail"); + Assertions.assertTrue(lines[1].startsWith("line "), "the body starts at a line boundary"); + } + + @Test + void readReturnsAnExactLimitFileWithoutANotice(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("exact.log"); + String content = "第一行\nsecond line\n"; + byte[] bytes = content.getBytes(StandardCharsets.UTF_8); + Files.write(file, bytes); + Assertions.assertEquals(content, LogContentReader.read(file, bytes.length)); + } + + @Test + void readKeepsTheOriginalSnapshotAfterGrowthAndRotation(@TempDir Path tempDir) + throws IOException { + Path file = tempDir.resolve("rotating.log"); + String content = "first line\nsecond line\n"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + FileUtils.FileTail whole = FileUtils.readFileTail(file, 1024); + FileUtils.FileTail tail = FileUtils.readFileTail(file, 15); + Files.write(file, "appended\n".getBytes(StandardCharsets.UTF_8), StandardOpenOption.APPEND); + Files.delete(file); + + Assertions.assertEquals(content, LogContentReader.read(whole)); + String response = LogContentReader.read(tail); + Assertions.assertTrue(response.contains("from the tail of " + content.length() + " bytes")); + Assertions.assertTrue(response.endsWith("second line\n")); + Assertions.assertFalse(response.contains("appended")); + } + + @Test + void readDescribesAnOversizedUtf8LineAccurately(@TempDir Path tempDir) throws IOException { + Path file = tempDir.resolve("oversized.log"); + String content = "日志日志日志日志日志\n"; + Files.write(file, content.getBytes(StandardCharsets.UTF_8)); + String[] response = LogContentReader.read(file, 10).split("\n", 2); + + Assertions.assertTrue(response[0].contains("UTF-8-safe partial tail")); + Assertions.assertTrue(response[0].contains("returning 10 bytes")); + Assertions.assertTrue(content.endsWith(response[1])); + Assertions.assertFalse(response[1].contains("\uFFFD")); + } +}
