From 9d2fa39a2d53f97e28436532d04d6630e2b71661 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=99=88=E5=AD=90=E9=BB=98?= <925456043@qq.com> Date: Mon, 10 Aug 2026 11:41:09 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E5=AE=8C=E5=96=84=E5=B7=A5=E4=BD=9C?= =?UTF-8?q?=E6=B5=81=E5=85=AC=E5=85=B1=E6=8E=A5=E5=8F=A3=E4=B8=8A=E4=BC=A0?= =?UTF-8?q?=E4=B8=8E=E9=94=99=E8=AF=AF=E5=A5=91=E7=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 区分 HTTP 请求标识与内部上传标识,补齐 Redis 旧记录兼容和关联日志 - 保持归一化 MIME 一致,并隔离工作流鉴权错误契约对其他公共接口的影响 - 收口 Multipart 操作日志与对象存储故障分类,归档范围:S05 --- .../controller/PublicWorkflowController.java | 6 +- .../interceptor/PublicApiInterceptor.java | 17 ++++ .../PublicWorkflowControllerRoutingTest.java | 6 +- .../interceptor/PublicApiInterceptorTest.java | 53 ++++++++++ .../impl/XFIleStorageServiceImpl.java | 18 ++-- .../impl/XFIleStorageServiceImplTest.java | 78 ++++++++++++++- .../upload/WorkflowApiPreparedUpload.java | 16 ++-- .../WorkflowApiUploadLifecycleService.java | 73 +++++++++----- .../upload/WorkflowApiUploadRecord.java | 34 ++++++- .../upload/WorkflowApiUploadStore.java | 74 +++++++------- .../upload/WorkflowApiUploadedFileReader.java | 25 +++-- ...WorkflowApiUploadLifecycleServiceTest.java | 79 +++++++++++++-- .../upload/WorkflowApiUploadStoreTest.java | 22 ++++- .../WorkflowApiUploadedFileReaderTest.java | 15 +-- .../log/reporter/ActionReportInterceptor.java | 96 ++++++++++++++++++- .../reporter/ActionReportInterceptorTest.java | 76 +++++++++++++++ 16 files changed, 573 insertions(+), 115 deletions(-) create mode 100644 easyflow-modules/easyflow-module-log/src/test/java/tech/easyflow/log/reporter/ActionReportInterceptorTest.java diff --git a/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/controller/PublicWorkflowController.java b/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/controller/PublicWorkflowController.java index 8e8c5193..b6b20f01 100644 --- a/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/controller/PublicWorkflowController.java +++ b/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/controller/PublicWorkflowController.java @@ -32,6 +32,7 @@ import tech.easyflow.common.constant.Constants; import tech.easyflow.common.domain.Result; import tech.easyflow.common.entity.LoginAccount; import tech.easyflow.common.satoken.util.SaTokenUtil; +import tech.easyflow.common.web.error.RequestIdContext; import tech.easyflow.common.web.exceptions.BusinessException; import tech.easyflow.common.web.jsonbody.JsonBody; import tech.easyflow.publicapi.dto.PublicWorkflowInfo; @@ -194,6 +195,7 @@ public class PublicWorkflowController { multipartRequest.getMultiFileMap()); WorkflowApiPreparedUpload preparedUpload = workflowApiUploadLifecycleService.prepare( + RequestIdContext.get(request), workflow.getContent(), metadata.getVariables(), fileParts); @@ -205,12 +207,12 @@ public class PublicWorkflowController { topology, executeId -> workflowApiUploadLifecycleService.bindExecution( - preparedUpload.getRequestId(), + preparedUpload.getUploadId(), executeId)); } catch (RuntimeException | Error error) { try { workflowApiUploadLifecycleService.abort( - preparedUpload.getRequestId()); + preparedUpload.getUploadId()); } catch (RuntimeException cleanupError) { error.addSuppressed(cleanupError); } diff --git a/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptor.java b/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptor.java index 4fff8ee4..21e5eec2 100644 --- a/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptor.java +++ b/easyflow-api/easyflow-api-public/src/main/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptor.java @@ -40,6 +40,12 @@ public class PublicApiInterceptor implements HandlerInterceptor { String apiKey = request.getHeader("ApiKey"); if (apiKey == null || apiKey.isBlank()) { + if (!isWorkflowApi(requestURI)) { + Result failed = Result.fail(401, "密钥不正确"); + response.setStatus(HttpServletResponse.SC_UNAUTHORIZED); + ResponseUtil.renderJson(response, failed); + return false; + } Result failed = Result.fail( "缺少 ApiKey 请求头", new PublicApiErrorDetail( @@ -62,4 +68,15 @@ public class PublicApiInterceptor implements HandlerInterceptor { ); return true; } + + /** + * 判断是否为工作流公共 API,避免专项错误契约影响其他公共接口。 + * + * @param requestUri 请求 URI + * @return 是否为工作流公共 API + */ + private boolean isWorkflowApi(String requestUri) { + return requestUri != null + && requestUri.contains("/public-api/workflow/"); + } } diff --git a/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/controller/PublicWorkflowControllerRoutingTest.java b/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/controller/PublicWorkflowControllerRoutingTest.java index 5ca73408..e574cf3e 100644 --- a/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/controller/PublicWorkflowControllerRoutingTest.java +++ b/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/controller/PublicWorkflowControllerRoutingTest.java @@ -112,6 +112,7 @@ public class PublicWorkflowControllerRoutingTest { when(multipartMapper.map(any())) .thenReturn(Map.of("file", List.of())); when(uploadLifecycleService.prepare( + anyString(), eq("{}"), anyMap(), anyMap())) @@ -181,6 +182,7 @@ public class PublicWorkflowControllerRoutingTest { eq("{}"), anyMap()); verify(uploadLifecycleService, never()).prepare( + anyString(), anyString(), anyMap(), anyMap()); @@ -209,12 +211,14 @@ public class PublicWorkflowControllerRoutingTest { mockMvc.perform(multipart(RUN_PATH) .file(metadata) .file(file) - .header("ApiKey", "key")) + .header("ApiKey", "key") + .header("X-Request-Id", "request-multipart")) .andExpect(status().isOk()) .andExpect(jsonPath("$.data").value("execute-1")); verify(multipartMapper).map(any()); verify(uploadLifecycleService).prepare( + eq("request-multipart"), eq("{}"), anyMap(), anyMap()); diff --git a/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptorTest.java b/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptorTest.java index ae2b6a3d..922d6e4a 100644 --- a/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptorTest.java +++ b/easyflow-api/easyflow-api-public/src/test/java/tech/easyflow/publicapi/interceptor/PublicApiInterceptorTest.java @@ -74,6 +74,59 @@ public class PublicApiInterceptorTest { Assert.assertTrue(body.toString().contains("request-1")); } + /** + * 验证非工作流公共 API 缺少令牌时继续保留既有通用错误契约。 + * + * @throws Exception 拦截器处理失败时抛出 + */ + @Test + public void shouldKeepGenericMissingKeyContractForOtherPublicApis() + throws Exception { + StringWriter body = new StringWriter(); + AtomicInteger status = new AtomicInteger(); + HttpServletRequest request = proxy( + HttpServletRequest.class, + (instance, method, args) -> { + if ("getRequestURI".equals(method.getName())) { + return "/public-api/knowledge-share/detail"; + } + if ("getHeader".equals(method.getName())) { + return null; + } + throw new AssertionError( + "测试路径不应调用 HttpServletRequest." + + method.getName()); + }); + HttpServletResponse response = proxy( + HttpServletResponse.class, + (instance, method, args) -> { + if ("setStatus".equals(method.getName())) { + status.set((Integer) args[0]); + return null; + } + if ("setContentType".equals(method.getName())) { + return null; + } + if ("getWriter".equals(method.getName())) { + return new PrintWriter(body); + } + throw new AssertionError( + "测试路径不应调用 HttpServletResponse." + + method.getName()); + }); + + boolean allowed = new PublicApiInterceptor() + .preHandle(request, response, new Object()); + + Assert.assertFalse(allowed); + Assert.assertEquals( + HttpServletResponse.SC_UNAUTHORIZED, + status.get()); + Assert.assertTrue(body.toString().contains("\"errorCode\":401")); + Assert.assertTrue(body.toString().contains("密钥不正确")); + Assert.assertFalse(body.toString().contains("工作流 Public API Key")); + } + /** * 验证通过接口权限校验的访问令牌会写入请求,供资源级鉴权复用。 * diff --git a/easyflow-commons/easyflow-common-file-storage/src/main/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImpl.java b/easyflow-commons/easyflow-common-file-storage/src/main/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImpl.java index 06ff94da..f1386325 100644 --- a/easyflow-commons/easyflow-common-file-storage/src/main/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImpl.java +++ b/easyflow-commons/easyflow-common-file-storage/src/main/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImpl.java @@ -122,21 +122,23 @@ public class XFIleStorageServiceImpl implements FileStorageService { } /** - * 获取上传文件的 Content-Type,并为文本文件补充 UTF-8 编码。 + * 获取上传文件的 Content-Type,并在客户端未声明时为文本文件补充 UTF-8 编码。 * * @param file 上传文件 * @return 文件媒体类型 */ public static String getFileContentType(MultipartFile file) { String originalFilename = file.getOriginalFilename(); - String contentType = null; - if (originalFilename != null && originalFilename.toLowerCase().endsWith(".txt")) { - contentType = "text/plain; charset=utf-8"; - } else { - // 其他类型文件可以按需设置 - contentType = file.getContentType(); + String contentType = file.getContentType(); + if (StringUtils.hasText(contentType)) { + return contentType; } - return contentType; + if (StringUtils.endsWithIgnoreCase( + originalFilename, + ".txt")) { + return "text/plain; charset=utf-8"; + } + return null; } /** diff --git a/easyflow-commons/easyflow-common-file-storage/src/test/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImplTest.java b/easyflow-commons/easyflow-common-file-storage/src/test/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImplTest.java index c9a271d8..4a844058 100644 --- a/easyflow-commons/easyflow-common-file-storage/src/test/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImplTest.java +++ b/easyflow-commons/easyflow-common-file-storage/src/test/java/tech/easyflow/common/filestorage/impl/XFIleStorageServiceImplTest.java @@ -109,6 +109,53 @@ public class XFIleStorageServiceImplTest { assertTrue(platform.exists); } + /** + * 验证可恢复上传保留上游已经归一化的 MIME,不再按 txt 扩展名二次覆盖。 + * + * @throws Exception 注入测试替身失败 + */ + @Test + public void recoverableSavePreservesNormalizedTextContentType() + throws Exception { + RecoverablePlatform platform = new RecoverablePlatform( + "minio-main", + "attachment", + "https://files/"); + RecoverableStorageService delegate = + new RecoverableStorageService(platform); + XFIleStorageServiceImpl service = createService(delegate); + FileStorageWriteHandle handle = service.prepareRecoverableWrite( + "workflow-api-upload/upload-1", + "content.txt"); + + service.saveRecoverable( + new BytesMultipartFile( + "content".getBytes( + java.nio.charset.StandardCharsets.UTF_8), + "content.txt", + "application/octet-stream"), + handle); + + assertEquals( + "application/octet-stream", + delegate.uploadContentType); + } + + /** + * 验证客户端未声明 MIME 时继续为 txt 文件补充 UTF-8 文本类型。 + */ + @Test + public void textFileWithoutContentTypeUsesUtf8Fallback() { + BytesMultipartFile file = new BytesMultipartFile( + new byte[]{1}, + "content.TXT", + null); + + assertEquals( + "text/plain; charset=utf-8", + XFIleStorageServiceImpl.getFileContentType(file)); + } + /** * 验证 MinIO 可恢复读取使用已配置客户端和句柄中的精确对象键,不请求公开 URL。 * @@ -398,6 +445,8 @@ public class XFIleStorageServiceImplTest { private String uploadPath; /** 上传文件名。 */ private String uploadFilename; + /** 上传媒体类型。 */ + private String uploadContentType; /** recorder 删除调用次数。 */ private int recorderDeleteCalls; /** recorder 删除是否抛出异常。 */ @@ -536,6 +585,7 @@ public class XFIleStorageServiceImplTest { */ @Override public org.dromara.x.file.storage.core.upload.UploadPretreatment setContentType(String contentType) { + delegate.uploadContentType = contentType; return this; } @@ -562,20 +612,42 @@ public class XFIleStorageServiceImplTest { private static final class BytesMultipartFile implements MultipartFile { /** 文件内容。 */ private final byte[] bytes; + /** 文件名。 */ + private final String filename; + /** 文件媒体类型。 */ + private final String contentType; /** * 创建上传文件替身。 * * @param bytes 文件内容 */ - private BytesMultipartFile(byte[] bytes) { this.bytes = bytes.clone(); } + private BytesMultipartFile(byte[] bytes) { + this(bytes, "content.bin", "application/octet-stream"); + } + + /** + * 创建指定文件名和媒体类型的上传文件替身。 + * + * @param bytes 文件内容 + * @param filename 文件名 + * @param contentType 文件媒体类型 + */ + private BytesMultipartFile( + byte[] bytes, + String filename, + String contentType) { + this.bytes = bytes.clone(); + this.filename = filename; + this.contentType = contentType; + } /** {@inheritDoc} */ @Override public String getName() { return "file"; } /** {@inheritDoc} */ - @Override public String getOriginalFilename() { return "content.bin"; } + @Override public String getOriginalFilename() { return filename; } /** {@inheritDoc} */ - @Override public String getContentType() { return "application/octet-stream"; } + @Override public String getContentType() { return contentType; } /** {@inheritDoc} */ @Override public boolean isEmpty() { return bytes.length == 0; } /** {@inheritDoc} */ diff --git a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiPreparedUpload.java b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiPreparedUpload.java index d010c58d..c9a8a348 100644 --- a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiPreparedUpload.java +++ b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiPreparedUpload.java @@ -7,29 +7,29 @@ import java.util.Map; */ public class WorkflowApiPreparedUpload { - private final String requestId; + private final String uploadId; private final Map variables; /** * 创建文件准备结果。 * - * @param requestId 临时上传请求 ID + * @param uploadId 内部临时上传 ID * @param variables 已注入文件描述的工作流变量 */ public WorkflowApiPreparedUpload( - String requestId, + String uploadId, Map variables) { - this.requestId = requestId; + this.uploadId = uploadId; this.variables = variables; } /** - * 获取临时上传请求 ID。 + * 获取内部临时上传 ID。 * - * @return 临时上传请求 ID + * @return 内部临时上传 ID */ - public String getRequestId() { - return requestId; + public String getUploadId() { + return uploadId; } /** diff --git a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleService.java b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleService.java index 239fd5b7..cb8bd3f2 100644 --- a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleService.java +++ b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleService.java @@ -17,7 +17,6 @@ import tech.easyflow.common.web.exceptions.BusinessException; import java.net.ConnectException; import java.net.SocketTimeoutException; -import java.net.UnknownHostException; import java.net.http.HttpTimeoutException; import java.time.Duration; import java.util.ArrayList; @@ -89,12 +88,14 @@ public class WorkflowApiUploadLifecycleService { /** * 校验、存储 multipart 文件并注入工作流变量。 * + * @param requestId HTTP 请求关联标识 * @param workflowContent 已发布工作流内容 * @param variables 普通运行变量 * @param fileParts 以工作流文件参数名分组的 multipart 文件 * @return 临时上传准备结果 */ public WorkflowApiPreparedUpload prepare( + String requestId, String workflowContent, Map variables, Map> fileParts) { @@ -126,7 +127,10 @@ public class WorkflowApiUploadLifecycleService { WorkflowApiUploadRecord record = new WorkflowApiUploadRecord(); long now = System.currentTimeMillis(); - record.setRequestId(UUID.randomUUID().toString().replace("-", "")); + record.setUploadId(UUID.randomUUID().toString().replace("-", "")); + record.setRequestId(StringUtils.hasText(requestId) + ? requestId + : record.getUploadId()); record.setCreatedAt(now); record.setCleanupAt(now + STAGED_RETENTION.toMillis()); @@ -142,7 +146,7 @@ public class WorkflowApiUploadLifecycleService { workflowContent, resolvedVariables); return new WorkflowApiPreparedUpload( - record.getRequestId(), + record.getUploadId(), normalized); } catch (RuntimeException | Error error) { try { @@ -157,14 +161,14 @@ public class WorkflowApiUploadLifecycleService { /** * 在工作流首个节点启动前绑定执行 ID。 * - * @param requestId 临时上传请求 ID + * @param uploadId 内部临时上传 ID * @param executeId 工作流执行 ID */ - public void bindExecution(String requestId, String executeId) { - uploadStore.bindExecution(requestId, executeId); - WorkflowApiUploadRecord record = uploadStore.find(requestId) + public void bindExecution(String uploadId, String executeId) { + uploadStore.bindExecution(uploadId, executeId); + WorkflowApiUploadRecord record = uploadStore.find(uploadId) .orElseThrow(() -> new IllegalStateException( - "工作流临时上传记录不存在: " + requestId)); + "工作流临时上传记录不存在: " + uploadId)); uploadStore.schedule( record, System.currentTimeMillis() + ACTIVE_RECHECK.toMillis()); @@ -186,10 +190,10 @@ public class WorkflowApiUploadLifecycleService { /** * 终止尚未成功启动的临时上传并立即清理。 * - * @param requestId 临时上传请求 ID + * @param uploadId 内部临时上传 ID */ - public void abort(String requestId) { - cleanupRequest(requestId, true); + public void abort(String uploadId) { + cleanupRequest(uploadId, true); } /** @@ -209,15 +213,12 @@ public class WorkflowApiUploadLifecycleService { if (claimed.isEmpty()) { break; } - String requestId = claimed.get(); + String uploadId = claimed.get(); processed++; try { - cleanupRequest(requestId, false); + cleanupRequest(uploadId, false); } catch (RuntimeException error) { - LOG.error( - "清理工作流 API 临时上传失败,requestId={}", - requestId, - error); + logCleanupFailure(uploadId, error); } } return processed; @@ -341,7 +342,7 @@ public class WorkflowApiUploadLifecycleService { try { writeHandle = fileStorageService.prepareRecoverableWrite( STORAGE_PATH_PREFIX - + record.getRequestId(), + + record.getUploadId(), buildStorageFilename( file, record.getStoredFiles().size())); @@ -499,7 +500,6 @@ public class WorkflowApiUploadLifecycleService { while (current != null) { if (current instanceof ConnectException || current instanceof SocketTimeoutException - || current instanceof UnknownHostException || current instanceof HttpTimeoutException || current instanceof TimeoutException) { return true; @@ -590,14 +590,14 @@ public class WorkflowApiUploadLifecycleService { /** * 在分布式锁下清理一条上传记录。 * - * @param requestId 上传请求 ID + * @param uploadId 内部上传 ID * @param force 是否忽略工作流运行状态立即清理 * @return 是否成功删除记录 */ - private boolean cleanupRequest(String requestId, boolean force) { + private boolean cleanupRequest(String uploadId, boolean force) { RedisLockExecutor.LockHandle handle = redisLockExecutor.tryAcquire( - CLEANUP_LOCK_PREFIX + requestId, + CLEANUP_LOCK_PREFIX + uploadId, CLEANUP_LOCK_WAIT, CLEANUP_LOCK_LEASE); if (handle == null) { @@ -605,9 +605,9 @@ public class WorkflowApiUploadLifecycleService { } try (handle) { WorkflowApiUploadRecord record = - uploadStore.find(requestId).orElse(null); + uploadStore.find(uploadId).orElse(null); if (record == null) { - uploadStore.removeMissingIndex(requestId); + uploadStore.removeMissingIndex(uploadId); return false; } if (!force && shouldRetain(record)) { @@ -623,6 +623,31 @@ public class WorkflowApiUploadLifecycleService { } } + /** + * 记录带 HTTP 请求关联标识的异步清理异常。 + * + * @param uploadId 内部上传 ID + * @param error 清理异常 + */ + private void logCleanupFailure( + String uploadId, + RuntimeException error) { + WorkflowApiUploadRecord record = null; + try { + record = uploadStore.find(uploadId).orElse(null); + } catch (RuntimeException lookupError) { + if (lookupError != error) { + error.addSuppressed(lookupError); + } + } + LOG.error( + "清理工作流 API 临时上传失败, uploadId={}, requestId={}, executeId={}", + uploadId, + record == null ? null : record.getRequestId(), + record == null ? null : record.getExecuteId(), + error); + } + /** * 判断上传记录是否仍被运行中或挂起的工作流使用。 * diff --git a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadRecord.java b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadRecord.java index d2e2436b..60899c4c 100644 --- a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadRecord.java +++ b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadRecord.java @@ -8,6 +8,9 @@ import java.util.List; */ public class WorkflowApiUploadRecord { + /** 服务端生成、用于 Redis 和对象存储定位的内部上传 ID。 */ + private String uploadId; + /** 调用方可见、用于跨线程排障的 HTTP 请求关联标识。 */ private String requestId; private String executeId; private List storedFiles = new ArrayList<>(); @@ -19,18 +22,41 @@ public class WorkflowApiUploadRecord { private long cleanupAt; /** - * 获取上传请求 ID。 + * 获取内部上传 ID。 * - * @return 上传请求 ID + *

兼容旧版 Redis 记录:旧记录只包含 {@code requestId}, + * 其值曾作为内部上传主键。

+ * + * @return 内部上传 ID + */ + public String getUploadId() { + return uploadId == null || uploadId.isBlank() + ? requestId + : uploadId; + } + + /** + * 设置内部上传 ID。 + * + * @param uploadId 内部上传 ID + */ + public void setUploadId(String uploadId) { + this.uploadId = uploadId; + } + + /** + * 获取 HTTP 请求关联标识。 + * + * @return HTTP 请求关联标识 */ public String getRequestId() { return requestId; } /** - * 设置上传请求 ID。 + * 设置 HTTP 请求关联标识。 * - * @param requestId 上传请求 ID + * @param requestId HTTP 请求关联标识 */ public void setRequestId(String requestId) { this.requestId = requestId; diff --git a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStore.java b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStore.java index c6d96f3e..6adc7cdb 100644 --- a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStore.java +++ b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStore.java @@ -98,50 +98,50 @@ public class WorkflowApiUploadStore { public void save(WorkflowApiUploadRecord record) { requireRecord(record); redisTemplate.opsForValue().set( - recordKey(record.getRequestId()), + recordKey(record.getUploadId()), serialize(record)); } /** * 将上传请求绑定到工作流执行实例。 * - * @param requestId 上传请求 ID + * @param uploadId 内部上传 ID * @param executeId 工作流执行 ID */ - public void bindExecution(String requestId, String executeId) { + public void bindExecution(String uploadId, String executeId) { if (!StringUtils.hasText(executeId)) { throw new IllegalArgumentException("工作流执行 ID 不能为空"); } - WorkflowApiUploadRecord record = find(requestId) + WorkflowApiUploadRecord record = find(uploadId) .orElseThrow(() -> new IllegalStateException( - "工作流临时上传记录不存在: " + requestId)); + "工作流临时上传记录不存在: " + uploadId)); String previousExecuteId = record.getExecuteId(); record.setExecuteId(executeId); redisTemplate.execute( BIND_EXECUTION_SCRIPT, Arrays.asList( - recordKey(requestId), + recordKey(uploadId), executionKey(executeId), StringUtils.hasText(previousExecuteId) ? executionKey(previousExecuteId) - : recordKey(requestId)), + : recordKey(uploadId)), serialize(record), - requestId, + uploadId, String.valueOf(EXECUTION_INDEX_TTL.toMillis()), StringUtils.hasText(previousExecuteId) ? "1" : "0"); } /** - * 按上传请求 ID 查找记录。 + * 按内部上传 ID 查找记录。 * - * @param requestId 上传请求 ID + * @param uploadId 内部上传 ID * @return 上传记录 */ - public Optional find(String requestId) { - if (!StringUtils.hasText(requestId)) { + public Optional find(String uploadId) { + if (!StringUtils.hasText(uploadId)) { return Optional.empty(); } - String value = redisTemplate.opsForValue().get(recordKey(requestId)); + String value = redisTemplate.opsForValue().get(recordKey(uploadId)); if (!StringUtils.hasText(value)) { return Optional.empty(); } @@ -151,7 +151,7 @@ public class WorkflowApiUploadStore { WorkflowApiUploadRecord.class)); } catch (JsonProcessingException error) { throw new IllegalStateException( - "读取工作流临时上传记录失败: " + requestId, + "读取工作流临时上传记录失败: " + uploadId, error); } } @@ -167,16 +167,16 @@ public class WorkflowApiUploadStore { if (!StringUtils.hasText(executeId)) { return Optional.empty(); } - String requestId = redisTemplate.opsForValue().get( + String uploadId = redisTemplate.opsForValue().get( executionKey(executeId)); - Optional record = find(requestId); - if (StringUtils.hasText(requestId) + Optional record = find(uploadId); + if (StringUtils.hasText(uploadId) && (record.isEmpty() || !executeId.equals(record.get().getExecuteId()))) { redisTemplate.execute( REMOVE_EXECUTION_INDEX_SCRIPT, List.of(executionKey(executeId)), - requestId); + uploadId); return Optional.empty(); } return record; @@ -199,17 +199,17 @@ public class WorkflowApiUploadStore { * * @param now 当前 Unix 毫秒时间戳 * @param retryAt 领取后默认重试时间 - * @return 领取到的上传请求 ID + * @return 领取到的内部上传 ID */ public Optional claimExpired( long now, long retryAt) { - String requestId = redisTemplate.execute( + String uploadId = redisTemplate.execute( CLAIM_EXPIRED_SCRIPT, List.of(CLEANUP_INDEX), String.valueOf(now), String.valueOf(retryAt)); - return Optional.ofNullable(requestId); + return Optional.ofNullable(uploadId); } /** @@ -224,23 +224,23 @@ public class WorkflowApiUploadStore { redisTemplate.execute( REMOVE_SCRIPT, Arrays.asList( - recordKey(record.getRequestId()), + recordKey(record.getUploadId()), hasExecuteId ? executionKey(record.getExecuteId()) - : recordKey(record.getRequestId()), + : recordKey(record.getUploadId()), CLEANUP_INDEX), - record.getRequestId(), + record.getUploadId(), hasExecuteId ? "1" : "0"); } /** * 删除已经缺少详情记录的残留清理索引。 * - * @param requestId 上传请求 ID + * @param uploadId 内部上传 ID */ - public void removeMissingIndex(String requestId) { - if (StringUtils.hasText(requestId)) { - redisTemplate.opsForZSet().remove(CLEANUP_INDEX, requestId); + public void removeMissingIndex(String uploadId) { + if (StringUtils.hasText(uploadId)) { + redisTemplate.opsForZSet().remove(CLEANUP_INDEX, uploadId); } } @@ -256,7 +256,7 @@ public class WorkflowApiUploadStore { } catch (JsonProcessingException error) { throw new IllegalStateException( "写入工作流临时上传记录失败: " - + record.getRequestId(), + + record.getUploadId(), error); } } @@ -272,14 +272,14 @@ public class WorkflowApiUploadStore { redisTemplate.execute( SAVE_SCHEDULED_SCRIPT, Arrays.asList( - recordKey(record.getRequestId()), + recordKey(record.getUploadId()), CLEANUP_INDEX, hasExecuteId ? executionKey(record.getExecuteId()) - : recordKey(record.getRequestId())), + : recordKey(record.getUploadId())), serialize(record), String.valueOf(record.getCleanupAt()), - record.getRequestId(), + record.getUploadId(), hasExecuteId ? "1" : "0", String.valueOf(EXECUTION_INDEX_TTL.toMillis())); } @@ -318,19 +318,19 @@ public class WorkflowApiUploadStore { * @param record 上传记录 */ private void requireRecord(WorkflowApiUploadRecord record) { - if (record == null || !StringUtils.hasText(record.getRequestId())) { - throw new IllegalArgumentException("工作流临时上传请求 ID 不能为空"); + if (record == null || !StringUtils.hasText(record.getUploadId())) { + throw new IllegalArgumentException("工作流内部上传 ID 不能为空"); } } /** * 构建记录 Redis Key。 * - * @param requestId 上传请求 ID + * @param uploadId 内部上传 ID * @return Redis Key */ - private String recordKey(String requestId) { - return RECORD_KEY_PREFIX + requestId; + private String recordKey(String uploadId) { + return RECORD_KEY_PREFIX + uploadId; } /** diff --git a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReader.java b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReader.java index 9b7802f5..c6a8173e 100644 --- a/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReader.java +++ b/easyflow-modules/easyflow-module-ai/src/main/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReader.java @@ -1,5 +1,7 @@ package tech.easyflow.ai.easyagentsflow.upload; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.stereotype.Component; import org.springframework.util.StringUtils; @@ -16,12 +18,14 @@ import java.util.regex.Pattern; /** * 读取经过 Public Workflow API 上传记录验证的临时文件。 * - *

外部文件描述只提供公开读取路径。仅当路径中的随机请求 ID、Redis 上传记录、 + *

外部文件描述只提供公开读取路径。仅当路径中的随机上传 ID、Redis 上传记录、 * 完整文件 URL 和可恢复存储句柄全部匹配时,才允许绕过公网 URL 限制并直接读取物理对象。

*/ @Component public class WorkflowApiUploadedFileReader { + private static final Logger LOG = LoggerFactory.getLogger( + WorkflowApiUploadedFileReader.class); private static final String STORAGE_PATH_PREFIX = "workflow-api-upload/"; private static final Pattern MANAGED_PATH_PATTERN = Pattern.compile( "(?:^|/)workflow-api-upload/([0-9a-f]{32})/([^/]+)$"); @@ -54,8 +58,11 @@ public class WorkflowApiUploadedFileReader { if (managedPath == null) { return Optional.empty(); } - WorkflowApiUploadRecord record = uploadStore.find(managedPath.requestId()) + WorkflowApiUploadRecord record = uploadStore.find(managedPath.uploadId()) .orElseThrow(() -> new IOException("工作流上传文件已失效,请重新上传")); + if (!managedPath.uploadId().equals(record.getUploadId())) { + throw new IOException("工作流上传文件引用与上传记录不匹配"); + } WorkflowApiStoredFile storedFile = record.getStoredFiles().stream() .filter(file -> file != null && filePath.equals(file.filePath())) .findFirst() @@ -67,7 +74,7 @@ public class WorkflowApiUploadedFileReader { } catch (IllegalArgumentException exception) { throw new IOException("工作流上传文件存储定位符无效", exception); } - String expectedStoragePath = STORAGE_PATH_PREFIX + managedPath.requestId() + "/"; + String expectedStoragePath = STORAGE_PATH_PREFIX + managedPath.uploadId() + "/"; if (!expectedStoragePath.equals(handle.getPath()) || !managedPath.filename().equals(handle.getFilename())) { throw new IOException("工作流上传文件存储定位与上传请求不匹配"); @@ -75,6 +82,12 @@ public class WorkflowApiUploadedFileReader { try { return Optional.of(fileStorageService.readRecoverable(handle)); } catch (RuntimeException exception) { + LOG.error( + "读取工作流 API 上传文件失败, uploadId={}, requestId={}, executeId={}", + record.getUploadId(), + record.getRequestId(), + record.getExecuteId(), + exception); throw new IOException("读取工作流上传文件失败", exception); } } @@ -92,7 +105,7 @@ public class WorkflowApiUploadedFileReader { } /** - * 从 URL 或相对路径中解析受管请求 ID 与固定文件名。 + * 从 URL 或相对路径中解析受管上传 ID 与固定文件名。 * * @param filePath 原始文件路径 * @return 受管路径信息 @@ -121,9 +134,9 @@ public class WorkflowApiUploadedFileReader { /** * 系统受管上传路径中的可信定位片段。 * - * @param requestId 随机上传请求 ID + * @param uploadId 随机内部上传 ID * @param filename 系统生成的存储文件名 */ - private record ManagedPath(String requestId, String filename) { + private record ManagedPath(String uploadId, String filename) { } } diff --git a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleServiceTest.java b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleServiceTest.java index 0dc0e41e..947feb28 100644 --- a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleServiceTest.java +++ b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadLifecycleServiceTest.java @@ -25,6 +25,8 @@ import java.util.Set; */ public class WorkflowApiUploadLifecycleServiceTest { + private static final String HTTP_REQUEST_ID = "http-request-1"; + /** * 验证同名文件 Part 会按顺序保存并注入文件对象数组。 */ @@ -63,11 +65,12 @@ public class WorkflowApiUploadLifecycleServiceTest { secondHandle.encodeLocator())); WorkflowApiPreparedUpload prepared = fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of("user_input", "解析"), Map.of("documents", List.of(first, second))); - Assert.assertNotNull(prepared.getRequestId()); + Assert.assertNotNull(prepared.getUploadId()); Assert.assertEquals("解析", prepared.getVariables().get("user_input")); @SuppressWarnings("unchecked") List> documents = @@ -85,6 +88,12 @@ public class WorkflowApiUploadLifecycleServiceTest { ArgumentCaptor.forClass(WorkflowApiUploadRecord.class); Mockito.verify(fixture.uploadStore).create( recordCaptor.capture()); + Assert.assertEquals( + HTTP_REQUEST_ID, + recordCaptor.getValue().getRequestId()); + Assert.assertNotEquals( + recordCaptor.getValue().getRequestId(), + recordCaptor.getValue().getUploadId()); Assert.assertEquals( List.of("/files/first.pdf", "/files/second.docx"), recordCaptor.getValue().getStoredFiles().stream() @@ -130,6 +139,7 @@ public class WorkflowApiUploadLifecycleServiceTest { handle.encodeLocator())); WorkflowApiPreparedUpload prepared = fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of(), Map.of("documents", List.of(file))); @@ -187,6 +197,7 @@ public class WorkflowApiUploadLifecycleServiceTest { BusinessException exception = Assert.assertThrows( BusinessException.class, () -> fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of(), Map.of("documents", List.of(file)))); @@ -232,6 +243,7 @@ public class WorkflowApiUploadLifecycleServiceTest { BusinessException exception = Assert.assertThrows( BusinessException.class, () -> fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of(), Map.of("documents", List.of(file)))); @@ -245,6 +257,53 @@ public class WorkflowApiUploadLifecycleServiceTest { .remove(Mockito.any(WorkflowApiUploadRecord.class)); } + /** + * 验证无法解析的存储主机按配置故障返回不可重试的 50001。 + */ + @Test + public void prepareShouldTreatUnknownStorageHostAsPermanentFailure() { + Fixture fixture = fixture(); + MultipartFile file = file( + "report.pdf", + "application/pdf", + 10L); + Mockito.when(fixture.parameterResolver + .resolveFileParameterNames("flow")) + .thenReturn(Set.of("documents")); + Mockito.when(fixture.parameterResolver.normalizeRuntimeVariables( + Mockito.eq("flow"), + Mockito.anyMap())) + .thenAnswer(invocation -> new LinkedHashMap<>( + invocation.getArgument(1))); + FileStorageWriteHandle handle = handle("report.pdf"); + Mockito.when(fixture.fileStorageService.prepareRecoverableWrite( + Mockito.anyString(), + Mockito.anyString())) + .thenReturn(handle); + Mockito.when(fixture.fileStorageService.saveRecoverable( + file, + handle)) + .thenThrow(new IllegalStateException( + "storage endpoint unavailable", + new java.net.UnknownHostException( + "invalid-storage-host"))); + + BusinessException exception = Assert.assertThrows( + BusinessException.class, + () -> fixture.service.prepare( + HTTP_REQUEST_ID, + "flow", + Map.of(), + Map.of("documents", List.of(file)))); + + Assert.assertEquals(500, exception.getHttpStatus()); + Assert.assertEquals(50001, exception.getErrorCode()); + Assert.assertFalse(exception.getMessage().contains( + "invalid-storage-host")); + Mockito.verify(fixture.fileStorageService) + .deleteRecoverable(handle); + } + /** * 验证未知文件 Part 在写入存储前被拒绝。 */ @@ -258,6 +317,7 @@ public class WorkflowApiUploadLifecycleServiceTest { BusinessException exception = Assert.assertThrows( BusinessException.class, () -> fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of(), Map.of("unknown", List.of(file)))); @@ -284,6 +344,7 @@ public class WorkflowApiUploadLifecycleServiceTest { BusinessException exception = Assert.assertThrows( BusinessException.class, () -> fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of(), Map.of("documents", List.of(file)))); @@ -313,6 +374,7 @@ public class WorkflowApiUploadLifecycleServiceTest { BusinessException exception = Assert.assertThrows( BusinessException.class, () -> fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of(), Map.of("documents", List.of(file)))); @@ -356,6 +418,7 @@ public class WorkflowApiUploadLifecycleServiceTest { IllegalStateException exception = Assert.assertThrows( IllegalStateException.class, () -> fixture.service.prepare( + HTTP_REQUEST_ID, "flow", Map.of(), Map.of("documents", List.of(file)))); @@ -374,7 +437,8 @@ public class WorkflowApiUploadLifecycleServiceTest { public void abortShouldDeleteStoredFilesAndRecord() { Fixture fixture = fixture(); WorkflowApiUploadRecord record = new WorkflowApiUploadRecord(); - record.setRequestId("request-1"); + record.setUploadId("upload-1"); + record.setRequestId(HTTP_REQUEST_ID); FileStorageWriteHandle firstHandle = handle("a.pdf"); FileStorageWriteHandle secondHandle = handle("b.pdf"); record.setStoredFiles(List.of( @@ -391,10 +455,10 @@ public class WorkflowApiUploadLifecycleServiceTest { Mockito.any(), Mockito.any())) .thenReturn(handle); - Mockito.when(fixture.uploadStore.find("request-1")) + Mockito.when(fixture.uploadStore.find("upload-1")) .thenReturn(Optional.of(record)); - fixture.service.abort("request-1"); + fixture.service.abort("upload-1"); Mockito.verify(fixture.fileStorageService) .deleteRecoverable(firstHandle); @@ -486,16 +550,17 @@ public class WorkflowApiUploadLifecycleServiceTest { /** * 创建包含单个临时文件的上传记录。 * - * @param requestId 请求 ID + * @param uploadId 内部上传 ID * @param handle 文件句柄 * @return 上传记录 */ private WorkflowApiUploadRecord storedRecord( - String requestId, + String uploadId, FileStorageWriteHandle handle) { WorkflowApiUploadRecord record = new WorkflowApiUploadRecord(); - record.setRequestId(requestId); + record.setUploadId(uploadId); + record.setRequestId(HTTP_REQUEST_ID); record.setStoredFiles(List.of( new WorkflowApiStoredFile( "/files/" + handle.getFilename(), diff --git a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStoreTest.java b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStoreTest.java index 73ed519a..2ccba25a 100644 --- a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStoreTest.java +++ b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadStoreTest.java @@ -116,18 +116,34 @@ public class WorkflowApiUploadStoreTest { .contains("redis.call('del', KEYS[3])")); } + /** + * 验证旧版 Redis 记录仍可把历史 requestId 作为内部上传 ID 读取。 + * + * @throws Exception JSON 反序列化失败 + */ + @Test + public void legacyRecordShouldResolveHistoricalUploadId() + throws Exception { + WorkflowApiUploadRecord record = new ObjectMapper().readValue( + "{\"requestId\":\"legacy-upload-1\"}", + WorkflowApiUploadRecord.class); + + Assert.assertEquals("legacy-upload-1", record.getUploadId()); + } + /** * 创建测试上传记录。 * - * @param requestId 上传请求 ID + * @param uploadId 内部上传 ID * @param executeId 执行 ID * @return 上传记录 */ private WorkflowApiUploadRecord record( - String requestId, + String uploadId, String executeId) { WorkflowApiUploadRecord record = new WorkflowApiUploadRecord(); - record.setRequestId(requestId); + record.setUploadId(uploadId); + record.setRequestId("http-request-1"); record.setExecuteId(executeId); record.setCleanupAt(1_000L); return record; diff --git a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReaderTest.java b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReaderTest.java index 39ffdb26..6e678e55 100644 --- a/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReaderTest.java +++ b/easyflow-modules/easyflow-module-ai/src/test/java/tech/easyflow/ai/easyagentsflow/upload/WorkflowApiUploadedFileReaderTest.java @@ -18,13 +18,13 @@ import java.util.Optional; */ public class WorkflowApiUploadedFileReaderTest { - private static final String REQUEST_ID = + private static final String UPLOAD_ID = "0123456789abcdef0123456789abcdef"; private static final String FILENAME = "000-abcdefabcdefabcdefabcdefabcdefab.docx"; private static final String FILE_URL = "http://127.0.0.1:39000/easyflow/attachment/" - + "workflow-api-upload/" + REQUEST_ID + "/" + FILENAME; + + "workflow-api-upload/" + UPLOAD_ID + "/" + FILENAME; /** * 验证 URL、上传记录与恢复句柄完全匹配后按固定后端读取。 @@ -40,7 +40,7 @@ public class WorkflowApiUploadedFileReaderTest { FileStorageWriteHandle handle = handle(); WorkflowApiUploadRecord record = record(FILE_URL, handle); byte[] content = "document-content".getBytes(StandardCharsets.UTF_8); - Mockito.when(uploadStore.find(REQUEST_ID)).thenReturn(Optional.of(record)); + Mockito.when(uploadStore.find(UPLOAD_ID)).thenReturn(Optional.of(record)); Mockito.when(fileStorageService.readRecoverable(handle)) .thenReturn(new ByteArrayInputStream(content)); @@ -81,7 +81,7 @@ public class WorkflowApiUploadedFileReaderTest { FileStorageService fileStorageService = Mockito.mock(FileStorageService.class); WorkflowApiUploadedFileReader reader = new WorkflowApiUploadedFileReader(uploadStore, fileStorageService); - Mockito.when(uploadStore.find(REQUEST_ID)).thenReturn(Optional.empty()); + Mockito.when(uploadStore.find(UPLOAD_ID)).thenReturn(Optional.empty()); IOException exception = Assert.assertThrows( IOException.class, @@ -100,7 +100,7 @@ public class WorkflowApiUploadedFileReaderTest { FileStorageService fileStorageService = Mockito.mock(FileStorageService.class); WorkflowApiUploadedFileReader reader = new WorkflowApiUploadedFileReader(uploadStore, fileStorageService); - Mockito.when(uploadStore.find(REQUEST_ID)).thenReturn(Optional.of( + Mockito.when(uploadStore.find(UPLOAD_ID)).thenReturn(Optional.of( record(FILE_URL + "?different=true", handle()))); IOException exception = Assert.assertThrows( @@ -121,7 +121,7 @@ public class WorkflowApiUploadedFileReaderTest { "local", "", "/tmp/easyflow-test", - "workflow-api-upload/" + REQUEST_ID, + "workflow-api-upload/" + UPLOAD_ID, FILENAME); } @@ -136,7 +136,8 @@ public class WorkflowApiUploadedFileReaderTest { String fileUrl, FileStorageWriteHandle handle) { WorkflowApiUploadRecord record = new WorkflowApiUploadRecord(); - record.setRequestId(REQUEST_ID); + record.setUploadId(UPLOAD_ID); + record.setRequestId("http-request-1"); record.setStoredFiles(List.of(new WorkflowApiStoredFile( fileUrl, handle.encodeLocator()))); diff --git a/easyflow-modules/easyflow-module-log/src/main/java/tech/easyflow/log/reporter/ActionReportInterceptor.java b/easyflow-modules/easyflow-module-log/src/main/java/tech/easyflow/log/reporter/ActionReportInterceptor.java index b9a118ed..71fa5907 100644 --- a/easyflow-modules/easyflow-module-log/src/main/java/tech/easyflow/log/reporter/ActionReportInterceptor.java +++ b/easyflow-modules/easyflow-module-log/src/main/java/tech/easyflow/log/reporter/ActionReportInterceptor.java @@ -6,12 +6,17 @@ import jakarta.servlet.http.HttpServletRequest; import jakarta.servlet.http.HttpServletResponse; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.http.MediaType; import org.springframework.stereotype.Component; +import org.springframework.web.multipart.MultipartFile; +import org.springframework.web.multipart.MultipartHttpServletRequest; import org.springframework.web.method.HandlerMethod; import org.springframework.web.servlet.HandlerInterceptor; import org.springframework.web.servlet.ModelAndView; import org.springframework.web.util.ContentCachingResponseWrapper; import tech.easyflow.common.util.RequestUtil; +import tech.easyflow.common.web.error.RequestIdContext; +import tech.easyflow.common.web.multipart.MultipartFileMetadataNormalizer; import tech.easyflow.log.annotation.LogReporterDisabled; import java.lang.reflect.Method; @@ -110,10 +115,16 @@ public class ActionReportInterceptor implements HandlerInterceptor { sb.append("EasyFlow action report -------- ").append(timestamp).append(" -------------------------\n"); sb.append("Request : ").append(request.getMethod()) .append(" ").append(request.getRequestURI()).append("\n"); + String requestId = RequestIdContext.get(request); + if (requestId != null) { + sb.append("RequestId : ").append(requestId).append("\n"); + } + + boolean multipartRequest = isMultipartRequest(request); // 打印参数(GET / POST 表单参数),脱敏 Map params = request.getParameterMap(); - if (!params.isEmpty()) { + if (!params.isEmpty() && !multipartRequest) { Map maskedParams = new LinkedHashMap<>(); for (Map.Entry entry : params.entrySet()) { String key = entry.getKey(); @@ -127,9 +138,17 @@ public class ActionReportInterceptor implements HandlerInterceptor { sb.append("Params : ").append(JSON.toJSONString(maskedParams)).append("\n"); } + if (multipartRequest) { + sb.append("Parts : ") + .append(buildMultipartSummary(request)) + .append("\n"); + } + // ====== 读取 POST Body ====== String methodStr = request.getMethod(); - if ("POST".equalsIgnoreCase(methodStr) || "PUT".equalsIgnoreCase(methodStr) || "PATCH".equalsIgnoreCase(methodStr)) { + if (!multipartRequest && ("POST".equalsIgnoreCase(methodStr) + || "PUT".equalsIgnoreCase(methodStr) + || "PATCH".equalsIgnoreCase(methodStr))) { String body = RequestUtil.readBodyString(request); if (body != null && !body.trim().isEmpty()) { try { @@ -208,8 +227,14 @@ public class ActionReportInterceptor implements HandlerInterceptor { if (ex != null) { sb.append('\n') .append("Status : FAILED\n") - .append("Exception : ").append(ex.getClass().getSimpleName()) - .append(": ").append(ex.getMessage() != null ? ex.getMessage().split("\n")[0] : "Unknown"); + .append("Exception : ") + .append(ex.getClass().getSimpleName()); + if (requestId == null) { + sb.append(": ") + .append(ex.getMessage() != null + ? ex.getMessage().split("\n")[0] + : "Unknown"); + } } // ====== 耗时 ====== @@ -258,6 +283,67 @@ public class ActionReportInterceptor implements HandlerInterceptor { return result; } + /** + * 判断请求是否为 Multipart 表单。 + * + * @param request 当前请求 + * @return 是否为 Multipart 请求 + */ + private boolean isMultipartRequest(HttpServletRequest request) { + String contentType = request.getContentType(); + return contentType != null + && contentType.regionMatches( + true, + 0, + MediaType.MULTIPART_FORM_DATA_VALUE, + 0, + MediaType.MULTIPART_FORM_DATA_VALUE.length()); + } + + /** + * 构建不读取文件正文的 Multipart Part 摘要。 + * + * @param request 当前请求 + * @return JSON 摘要 + */ + String buildMultipartSummary(HttpServletRequest request) { + if (!(request instanceof MultipartHttpServletRequest multipart)) { + return "{\"available\":false}"; + } + List> parts = new ArrayList<>(); + multipart.getMultiFileMap().forEach((partName, files) -> { + if (files == null || files.isEmpty()) { + parts.add(Map.of("partName", partName, "fileCount", 0)); + return; + } + for (MultipartFile file : files) { + if (file == null) { + parts.add(Map.of("partName", partName, "file", "null")); + continue; + } + String filename = + MultipartFileMetadataNormalizer.sanitizeFilename( + file.getOriginalFilename()); + Map summary = new LinkedHashMap<>(); + summary.put("partName", partName); + summary.put("fileName", filename); + summary.put("size", file.getSize()); + summary.put( + "contentType", + MultipartFileMetadataNormalizer.normalizeContentType( + filename, + file.getContentType())); + parts.add(summary); + } + }); + Set textPartNames = new LinkedHashSet<>( + multipart.getParameterMap().keySet()); + Map summary = new LinkedHashMap<>(); + summary.put("files", parts); + summary.put("textPartNames", textPartNames); + return JSON.toJSONString(summary); + } + /** * 构建方法签名:methodName(paramType paramName, ...) */ @@ -279,4 +365,4 @@ public class ActionReportInterceptor implements HandlerInterceptor { sig.append(")"); return sig.toString(); } -} \ No newline at end of file +} diff --git a/easyflow-modules/easyflow-module-log/src/test/java/tech/easyflow/log/reporter/ActionReportInterceptorTest.java b/easyflow-modules/easyflow-module-log/src/test/java/tech/easyflow/log/reporter/ActionReportInterceptorTest.java new file mode 100644 index 00000000..1239c518 --- /dev/null +++ b/easyflow-modules/easyflow-module-log/src/test/java/tech/easyflow/log/reporter/ActionReportInterceptorTest.java @@ -0,0 +1,76 @@ +package tech.easyflow.log.reporter; + +import jakarta.servlet.http.HttpServletResponse; +import org.junit.Test; +import org.springframework.util.LinkedMultiValueMap; +import org.springframework.web.method.HandlerMethod; +import org.springframework.web.multipart.MultipartFile; +import org.springframework.web.multipart.MultipartHttpServletRequest; + +import java.util.Map; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * {@link ActionReportInterceptor} Multipart 日志安全测试。 + */ +public class ActionReportInterceptorTest { + + /** + * 验证 Multipart 日志只读取 Part 元数据,不读取原始请求体。 + * + * @throws Exception 构造处理器或日志回调失败 + */ + @Test + public void multipartReportShouldNotReadRawRequestBody() + throws Exception { + ActionReportInterceptor interceptor = + new ActionReportInterceptor( + new ActionLogReporterProperties()); + MultipartHttpServletRequest request = + mock(MultipartHttpServletRequest.class); + HttpServletResponse response = mock(HttpServletResponse.class); + MultipartFile file = mock(MultipartFile.class); + LinkedMultiValueMap files = + new LinkedMultiValueMap<>(); + files.add("files.document", file); + when(request.getMethod()).thenReturn("POST"); + when(request.getRequestURI()) + .thenReturn("/public-api/workflow/runAsync"); + when(request.getContentType()) + .thenReturn("multipart/form-data; boundary=test"); + when(request.getParameterMap()).thenReturn(Map.of()); + when(request.getMultiFileMap()).thenReturn(files); + when(file.getOriginalFilename()) + .thenReturn("C:\\fakepath\\report.docx"); + when(file.getContentType()).thenReturn("Other"); + when(file.getSize()).thenReturn(128L); + HandlerMethod handler = new HandlerMethod( + new SampleController(), + SampleController.class.getDeclaredMethod("run")); + + interceptor.preHandle(request, response, handler); + interceptor.afterCompletion( + request, + response, + handler, + null); + + verify(request, never()).getInputStream(); + } + + /** + * 提供测试用处理方法。 + */ + private static final class SampleController { + + /** + * 测试处理方法。 + */ + private void run() { + } + } +}