发布 v1.10 #5
@@ -32,6 +32,7 @@ import tech.easyflow.common.constant.Constants;
|
|||||||
import tech.easyflow.common.domain.Result;
|
import tech.easyflow.common.domain.Result;
|
||||||
import tech.easyflow.common.entity.LoginAccount;
|
import tech.easyflow.common.entity.LoginAccount;
|
||||||
import tech.easyflow.common.satoken.util.SaTokenUtil;
|
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.exceptions.BusinessException;
|
||||||
import tech.easyflow.common.web.jsonbody.JsonBody;
|
import tech.easyflow.common.web.jsonbody.JsonBody;
|
||||||
import tech.easyflow.publicapi.dto.PublicWorkflowInfo;
|
import tech.easyflow.publicapi.dto.PublicWorkflowInfo;
|
||||||
@@ -194,6 +195,7 @@ public class PublicWorkflowController {
|
|||||||
multipartRequest.getMultiFileMap());
|
multipartRequest.getMultiFileMap());
|
||||||
WorkflowApiPreparedUpload preparedUpload =
|
WorkflowApiPreparedUpload preparedUpload =
|
||||||
workflowApiUploadLifecycleService.prepare(
|
workflowApiUploadLifecycleService.prepare(
|
||||||
|
RequestIdContext.get(request),
|
||||||
workflow.getContent(),
|
workflow.getContent(),
|
||||||
metadata.getVariables(),
|
metadata.getVariables(),
|
||||||
fileParts);
|
fileParts);
|
||||||
@@ -205,12 +207,12 @@ public class PublicWorkflowController {
|
|||||||
topology,
|
topology,
|
||||||
executeId ->
|
executeId ->
|
||||||
workflowApiUploadLifecycleService.bindExecution(
|
workflowApiUploadLifecycleService.bindExecution(
|
||||||
preparedUpload.getRequestId(),
|
preparedUpload.getUploadId(),
|
||||||
executeId));
|
executeId));
|
||||||
} catch (RuntimeException | Error error) {
|
} catch (RuntimeException | Error error) {
|
||||||
try {
|
try {
|
||||||
workflowApiUploadLifecycleService.abort(
|
workflowApiUploadLifecycleService.abort(
|
||||||
preparedUpload.getRequestId());
|
preparedUpload.getUploadId());
|
||||||
} catch (RuntimeException cleanupError) {
|
} catch (RuntimeException cleanupError) {
|
||||||
error.addSuppressed(cleanupError);
|
error.addSuppressed(cleanupError);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -40,6 +40,12 @@ public class PublicApiInterceptor implements HandlerInterceptor {
|
|||||||
String apiKey = request.getHeader("ApiKey");
|
String apiKey = request.getHeader("ApiKey");
|
||||||
|
|
||||||
if (apiKey == null || apiKey.isBlank()) {
|
if (apiKey == null || apiKey.isBlank()) {
|
||||||
|
if (!isWorkflowApi(requestURI)) {
|
||||||
|
Result<Void> failed = Result.fail(401, "密钥不正确");
|
||||||
|
response.setStatus(HttpServletResponse.SC_UNAUTHORIZED);
|
||||||
|
ResponseUtil.renderJson(response, failed);
|
||||||
|
return false;
|
||||||
|
}
|
||||||
Result<PublicApiErrorDetail> failed = Result.fail(
|
Result<PublicApiErrorDetail> failed = Result.fail(
|
||||||
"缺少 ApiKey 请求头",
|
"缺少 ApiKey 请求头",
|
||||||
new PublicApiErrorDetail(
|
new PublicApiErrorDetail(
|
||||||
@@ -62,4 +68,15 @@ public class PublicApiInterceptor implements HandlerInterceptor {
|
|||||||
);
|
);
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 判断是否为工作流公共 API,避免专项错误契约影响其他公共接口。
|
||||||
|
*
|
||||||
|
* @param requestUri 请求 URI
|
||||||
|
* @return 是否为工作流公共 API
|
||||||
|
*/
|
||||||
|
private boolean isWorkflowApi(String requestUri) {
|
||||||
|
return requestUri != null
|
||||||
|
&& requestUri.contains("/public-api/workflow/");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -112,6 +112,7 @@ public class PublicWorkflowControllerRoutingTest {
|
|||||||
when(multipartMapper.map(any()))
|
when(multipartMapper.map(any()))
|
||||||
.thenReturn(Map.of("file", List.of()));
|
.thenReturn(Map.of("file", List.of()));
|
||||||
when(uploadLifecycleService.prepare(
|
when(uploadLifecycleService.prepare(
|
||||||
|
anyString(),
|
||||||
eq("{}"),
|
eq("{}"),
|
||||||
anyMap(),
|
anyMap(),
|
||||||
anyMap()))
|
anyMap()))
|
||||||
@@ -181,6 +182,7 @@ public class PublicWorkflowControllerRoutingTest {
|
|||||||
eq("{}"),
|
eq("{}"),
|
||||||
anyMap());
|
anyMap());
|
||||||
verify(uploadLifecycleService, never()).prepare(
|
verify(uploadLifecycleService, never()).prepare(
|
||||||
|
anyString(),
|
||||||
anyString(),
|
anyString(),
|
||||||
anyMap(),
|
anyMap(),
|
||||||
anyMap());
|
anyMap());
|
||||||
@@ -209,12 +211,14 @@ public class PublicWorkflowControllerRoutingTest {
|
|||||||
mockMvc.perform(multipart(RUN_PATH)
|
mockMvc.perform(multipart(RUN_PATH)
|
||||||
.file(metadata)
|
.file(metadata)
|
||||||
.file(file)
|
.file(file)
|
||||||
.header("ApiKey", "key"))
|
.header("ApiKey", "key")
|
||||||
|
.header("X-Request-Id", "request-multipart"))
|
||||||
.andExpect(status().isOk())
|
.andExpect(status().isOk())
|
||||||
.andExpect(jsonPath("$.data").value("execute-1"));
|
.andExpect(jsonPath("$.data").value("execute-1"));
|
||||||
|
|
||||||
verify(multipartMapper).map(any());
|
verify(multipartMapper).map(any());
|
||||||
verify(uploadLifecycleService).prepare(
|
verify(uploadLifecycleService).prepare(
|
||||||
|
eq("request-multipart"),
|
||||||
eq("{}"),
|
eq("{}"),
|
||||||
anyMap(),
|
anyMap(),
|
||||||
anyMap());
|
anyMap());
|
||||||
|
|||||||
@@ -74,6 +74,59 @@ public class PublicApiInterceptorTest {
|
|||||||
Assert.assertTrue(body.toString().contains("request-1"));
|
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"));
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 验证通过接口权限校验的访问令牌会写入请求,供资源级鉴权复用。
|
* 验证通过接口权限校验的访问令牌会写入请求,供资源级鉴权复用。
|
||||||
*
|
*
|
||||||
|
|||||||
@@ -122,21 +122,23 @@ public class XFIleStorageServiceImpl implements FileStorageService {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 获取上传文件的 Content-Type,并为文本文件补充 UTF-8 编码。
|
* 获取上传文件的 Content-Type,并在客户端未声明时为文本文件补充 UTF-8 编码。
|
||||||
*
|
*
|
||||||
* @param file 上传文件
|
* @param file 上传文件
|
||||||
* @return 文件媒体类型
|
* @return 文件媒体类型
|
||||||
*/
|
*/
|
||||||
public static String getFileContentType(MultipartFile file) {
|
public static String getFileContentType(MultipartFile file) {
|
||||||
String originalFilename = file.getOriginalFilename();
|
String originalFilename = file.getOriginalFilename();
|
||||||
String contentType = null;
|
String contentType = file.getContentType();
|
||||||
if (originalFilename != null && originalFilename.toLowerCase().endsWith(".txt")) {
|
if (StringUtils.hasText(contentType)) {
|
||||||
contentType = "text/plain; charset=utf-8";
|
return contentType;
|
||||||
} else {
|
|
||||||
// 其他类型文件可以按需设置
|
|
||||||
contentType = file.getContentType();
|
|
||||||
}
|
}
|
||||||
return contentType;
|
if (StringUtils.endsWithIgnoreCase(
|
||||||
|
originalFilename,
|
||||||
|
".txt")) {
|
||||||
|
return "text/plain; charset=utf-8";
|
||||||
|
}
|
||||||
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -109,6 +109,53 @@ public class XFIleStorageServiceImplTest {
|
|||||||
assertTrue(platform.exists);
|
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。
|
* 验证 MinIO 可恢复读取使用已配置客户端和句柄中的精确对象键,不请求公开 URL。
|
||||||
*
|
*
|
||||||
@@ -398,6 +445,8 @@ public class XFIleStorageServiceImplTest {
|
|||||||
private String uploadPath;
|
private String uploadPath;
|
||||||
/** 上传文件名。 */
|
/** 上传文件名。 */
|
||||||
private String uploadFilename;
|
private String uploadFilename;
|
||||||
|
/** 上传媒体类型。 */
|
||||||
|
private String uploadContentType;
|
||||||
/** recorder 删除调用次数。 */
|
/** recorder 删除调用次数。 */
|
||||||
private int recorderDeleteCalls;
|
private int recorderDeleteCalls;
|
||||||
/** recorder 删除是否抛出异常。 */
|
/** recorder 删除是否抛出异常。 */
|
||||||
@@ -536,6 +585,7 @@ public class XFIleStorageServiceImplTest {
|
|||||||
*/
|
*/
|
||||||
@Override
|
@Override
|
||||||
public org.dromara.x.file.storage.core.upload.UploadPretreatment setContentType(String contentType) {
|
public org.dromara.x.file.storage.core.upload.UploadPretreatment setContentType(String contentType) {
|
||||||
|
delegate.uploadContentType = contentType;
|
||||||
return this;
|
return this;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -562,20 +612,42 @@ public class XFIleStorageServiceImplTest {
|
|||||||
private static final class BytesMultipartFile implements MultipartFile {
|
private static final class BytesMultipartFile implements MultipartFile {
|
||||||
/** 文件内容。 */
|
/** 文件内容。 */
|
||||||
private final byte[] bytes;
|
private final byte[] bytes;
|
||||||
|
/** 文件名。 */
|
||||||
|
private final String filename;
|
||||||
|
/** 文件媒体类型。 */
|
||||||
|
private final String contentType;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 创建上传文件替身。
|
* 创建上传文件替身。
|
||||||
*
|
*
|
||||||
* @param bytes 文件内容
|
* @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} */
|
/** {@inheritDoc} */
|
||||||
@Override public String getName() { return "file"; }
|
@Override public String getName() { return "file"; }
|
||||||
/** {@inheritDoc} */
|
/** {@inheritDoc} */
|
||||||
@Override public String getOriginalFilename() { return "content.bin"; }
|
@Override public String getOriginalFilename() { return filename; }
|
||||||
/** {@inheritDoc} */
|
/** {@inheritDoc} */
|
||||||
@Override public String getContentType() { return "application/octet-stream"; }
|
@Override public String getContentType() { return contentType; }
|
||||||
/** {@inheritDoc} */
|
/** {@inheritDoc} */
|
||||||
@Override public boolean isEmpty() { return bytes.length == 0; }
|
@Override public boolean isEmpty() { return bytes.length == 0; }
|
||||||
/** {@inheritDoc} */
|
/** {@inheritDoc} */
|
||||||
|
|||||||
@@ -7,29 +7,29 @@ import java.util.Map;
|
|||||||
*/
|
*/
|
||||||
public class WorkflowApiPreparedUpload {
|
public class WorkflowApiPreparedUpload {
|
||||||
|
|
||||||
private final String requestId;
|
private final String uploadId;
|
||||||
private final Map<String, Object> variables;
|
private final Map<String, Object> variables;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 创建文件准备结果。
|
* 创建文件准备结果。
|
||||||
*
|
*
|
||||||
* @param requestId 临时上传请求 ID
|
* @param uploadId 内部临时上传 ID
|
||||||
* @param variables 已注入文件描述的工作流变量
|
* @param variables 已注入文件描述的工作流变量
|
||||||
*/
|
*/
|
||||||
public WorkflowApiPreparedUpload(
|
public WorkflowApiPreparedUpload(
|
||||||
String requestId,
|
String uploadId,
|
||||||
Map<String, Object> variables) {
|
Map<String, Object> variables) {
|
||||||
this.requestId = requestId;
|
this.uploadId = uploadId;
|
||||||
this.variables = variables;
|
this.variables = variables;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 获取临时上传请求 ID。
|
* 获取内部临时上传 ID。
|
||||||
*
|
*
|
||||||
* @return 临时上传请求 ID
|
* @return 内部临时上传 ID
|
||||||
*/
|
*/
|
||||||
public String getRequestId() {
|
public String getUploadId() {
|
||||||
return requestId;
|
return uploadId;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -17,7 +17,6 @@ import tech.easyflow.common.web.exceptions.BusinessException;
|
|||||||
|
|
||||||
import java.net.ConnectException;
|
import java.net.ConnectException;
|
||||||
import java.net.SocketTimeoutException;
|
import java.net.SocketTimeoutException;
|
||||||
import java.net.UnknownHostException;
|
|
||||||
import java.net.http.HttpTimeoutException;
|
import java.net.http.HttpTimeoutException;
|
||||||
import java.time.Duration;
|
import java.time.Duration;
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
@@ -89,12 +88,14 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
/**
|
/**
|
||||||
* 校验、存储 multipart 文件并注入工作流变量。
|
* 校验、存储 multipart 文件并注入工作流变量。
|
||||||
*
|
*
|
||||||
|
* @param requestId HTTP 请求关联标识
|
||||||
* @param workflowContent 已发布工作流内容
|
* @param workflowContent 已发布工作流内容
|
||||||
* @param variables 普通运行变量
|
* @param variables 普通运行变量
|
||||||
* @param fileParts 以工作流文件参数名分组的 multipart 文件
|
* @param fileParts 以工作流文件参数名分组的 multipart 文件
|
||||||
* @return 临时上传准备结果
|
* @return 临时上传准备结果
|
||||||
*/
|
*/
|
||||||
public WorkflowApiPreparedUpload prepare(
|
public WorkflowApiPreparedUpload prepare(
|
||||||
|
String requestId,
|
||||||
String workflowContent,
|
String workflowContent,
|
||||||
Map<String, Object> variables,
|
Map<String, Object> variables,
|
||||||
Map<String, List<MultipartFile>> fileParts) {
|
Map<String, List<MultipartFile>> fileParts) {
|
||||||
@@ -126,7 +127,10 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
|
|
||||||
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
||||||
long now = System.currentTimeMillis();
|
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.setCreatedAt(now);
|
||||||
record.setCleanupAt(now + STAGED_RETENTION.toMillis());
|
record.setCleanupAt(now + STAGED_RETENTION.toMillis());
|
||||||
|
|
||||||
@@ -142,7 +146,7 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
workflowContent,
|
workflowContent,
|
||||||
resolvedVariables);
|
resolvedVariables);
|
||||||
return new WorkflowApiPreparedUpload(
|
return new WorkflowApiPreparedUpload(
|
||||||
record.getRequestId(),
|
record.getUploadId(),
|
||||||
normalized);
|
normalized);
|
||||||
} catch (RuntimeException | Error error) {
|
} catch (RuntimeException | Error error) {
|
||||||
try {
|
try {
|
||||||
@@ -157,14 +161,14 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
/**
|
/**
|
||||||
* 在工作流首个节点启动前绑定执行 ID。
|
* 在工作流首个节点启动前绑定执行 ID。
|
||||||
*
|
*
|
||||||
* @param requestId 临时上传请求 ID
|
* @param uploadId 内部临时上传 ID
|
||||||
* @param executeId 工作流执行 ID
|
* @param executeId 工作流执行 ID
|
||||||
*/
|
*/
|
||||||
public void bindExecution(String requestId, String executeId) {
|
public void bindExecution(String uploadId, String executeId) {
|
||||||
uploadStore.bindExecution(requestId, executeId);
|
uploadStore.bindExecution(uploadId, executeId);
|
||||||
WorkflowApiUploadRecord record = uploadStore.find(requestId)
|
WorkflowApiUploadRecord record = uploadStore.find(uploadId)
|
||||||
.orElseThrow(() -> new IllegalStateException(
|
.orElseThrow(() -> new IllegalStateException(
|
||||||
"工作流临时上传记录不存在: " + requestId));
|
"工作流临时上传记录不存在: " + uploadId));
|
||||||
uploadStore.schedule(
|
uploadStore.schedule(
|
||||||
record,
|
record,
|
||||||
System.currentTimeMillis() + ACTIVE_RECHECK.toMillis());
|
System.currentTimeMillis() + ACTIVE_RECHECK.toMillis());
|
||||||
@@ -186,10 +190,10 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
/**
|
/**
|
||||||
* 终止尚未成功启动的临时上传并立即清理。
|
* 终止尚未成功启动的临时上传并立即清理。
|
||||||
*
|
*
|
||||||
* @param requestId 临时上传请求 ID
|
* @param uploadId 内部临时上传 ID
|
||||||
*/
|
*/
|
||||||
public void abort(String requestId) {
|
public void abort(String uploadId) {
|
||||||
cleanupRequest(requestId, true);
|
cleanupRequest(uploadId, true);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -209,15 +213,12 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
if (claimed.isEmpty()) {
|
if (claimed.isEmpty()) {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
String requestId = claimed.get();
|
String uploadId = claimed.get();
|
||||||
processed++;
|
processed++;
|
||||||
try {
|
try {
|
||||||
cleanupRequest(requestId, false);
|
cleanupRequest(uploadId, false);
|
||||||
} catch (RuntimeException error) {
|
} catch (RuntimeException error) {
|
||||||
LOG.error(
|
logCleanupFailure(uploadId, error);
|
||||||
"清理工作流 API 临时上传失败,requestId={}",
|
|
||||||
requestId,
|
|
||||||
error);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return processed;
|
return processed;
|
||||||
@@ -341,7 +342,7 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
try {
|
try {
|
||||||
writeHandle = fileStorageService.prepareRecoverableWrite(
|
writeHandle = fileStorageService.prepareRecoverableWrite(
|
||||||
STORAGE_PATH_PREFIX
|
STORAGE_PATH_PREFIX
|
||||||
+ record.getRequestId(),
|
+ record.getUploadId(),
|
||||||
buildStorageFilename(
|
buildStorageFilename(
|
||||||
file,
|
file,
|
||||||
record.getStoredFiles().size()));
|
record.getStoredFiles().size()));
|
||||||
@@ -499,7 +500,6 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
while (current != null) {
|
while (current != null) {
|
||||||
if (current instanceof ConnectException
|
if (current instanceof ConnectException
|
||||||
|| current instanceof SocketTimeoutException
|
|| current instanceof SocketTimeoutException
|
||||||
|| current instanceof UnknownHostException
|
|
||||||
|| current instanceof HttpTimeoutException
|
|| current instanceof HttpTimeoutException
|
||||||
|| current instanceof TimeoutException) {
|
|| current instanceof TimeoutException) {
|
||||||
return true;
|
return true;
|
||||||
@@ -590,14 +590,14 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
/**
|
/**
|
||||||
* 在分布式锁下清理一条上传记录。
|
* 在分布式锁下清理一条上传记录。
|
||||||
*
|
*
|
||||||
* @param requestId 上传请求 ID
|
* @param uploadId 内部上传 ID
|
||||||
* @param force 是否忽略工作流运行状态立即清理
|
* @param force 是否忽略工作流运行状态立即清理
|
||||||
* @return 是否成功删除记录
|
* @return 是否成功删除记录
|
||||||
*/
|
*/
|
||||||
private boolean cleanupRequest(String requestId, boolean force) {
|
private boolean cleanupRequest(String uploadId, boolean force) {
|
||||||
RedisLockExecutor.LockHandle handle =
|
RedisLockExecutor.LockHandle handle =
|
||||||
redisLockExecutor.tryAcquire(
|
redisLockExecutor.tryAcquire(
|
||||||
CLEANUP_LOCK_PREFIX + requestId,
|
CLEANUP_LOCK_PREFIX + uploadId,
|
||||||
CLEANUP_LOCK_WAIT,
|
CLEANUP_LOCK_WAIT,
|
||||||
CLEANUP_LOCK_LEASE);
|
CLEANUP_LOCK_LEASE);
|
||||||
if (handle == null) {
|
if (handle == null) {
|
||||||
@@ -605,9 +605,9 @@ public class WorkflowApiUploadLifecycleService {
|
|||||||
}
|
}
|
||||||
try (handle) {
|
try (handle) {
|
||||||
WorkflowApiUploadRecord record =
|
WorkflowApiUploadRecord record =
|
||||||
uploadStore.find(requestId).orElse(null);
|
uploadStore.find(uploadId).orElse(null);
|
||||||
if (record == null) {
|
if (record == null) {
|
||||||
uploadStore.removeMissingIndex(requestId);
|
uploadStore.removeMissingIndex(uploadId);
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
if (!force && shouldRetain(record)) {
|
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);
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 判断上传记录是否仍被运行中或挂起的工作流使用。
|
* 判断上传记录是否仍被运行中或挂起的工作流使用。
|
||||||
*
|
*
|
||||||
|
|||||||
@@ -8,6 +8,9 @@ import java.util.List;
|
|||||||
*/
|
*/
|
||||||
public class WorkflowApiUploadRecord {
|
public class WorkflowApiUploadRecord {
|
||||||
|
|
||||||
|
/** 服务端生成、用于 Redis 和对象存储定位的内部上传 ID。 */
|
||||||
|
private String uploadId;
|
||||||
|
/** 调用方可见、用于跨线程排障的 HTTP 请求关联标识。 */
|
||||||
private String requestId;
|
private String requestId;
|
||||||
private String executeId;
|
private String executeId;
|
||||||
private List<WorkflowApiStoredFile> storedFiles = new ArrayList<>();
|
private List<WorkflowApiStoredFile> storedFiles = new ArrayList<>();
|
||||||
@@ -19,18 +22,41 @@ public class WorkflowApiUploadRecord {
|
|||||||
private long cleanupAt;
|
private long cleanupAt;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 获取上传请求 ID。
|
* 获取内部上传 ID。
|
||||||
*
|
*
|
||||||
* @return 上传请求 ID
|
* <p>兼容旧版 Redis 记录:旧记录只包含 {@code requestId},
|
||||||
|
* 其值曾作为内部上传主键。</p>
|
||||||
|
*
|
||||||
|
* @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() {
|
public String getRequestId() {
|
||||||
return requestId;
|
return requestId;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 设置上传请求 ID。
|
* 设置 HTTP 请求关联标识。
|
||||||
*
|
*
|
||||||
* @param requestId 上传请求 ID
|
* @param requestId HTTP 请求关联标识
|
||||||
*/
|
*/
|
||||||
public void setRequestId(String requestId) {
|
public void setRequestId(String requestId) {
|
||||||
this.requestId = requestId;
|
this.requestId = requestId;
|
||||||
|
|||||||
@@ -98,50 +98,50 @@ public class WorkflowApiUploadStore {
|
|||||||
public void save(WorkflowApiUploadRecord record) {
|
public void save(WorkflowApiUploadRecord record) {
|
||||||
requireRecord(record);
|
requireRecord(record);
|
||||||
redisTemplate.opsForValue().set(
|
redisTemplate.opsForValue().set(
|
||||||
recordKey(record.getRequestId()),
|
recordKey(record.getUploadId()),
|
||||||
serialize(record));
|
serialize(record));
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 将上传请求绑定到工作流执行实例。
|
* 将上传请求绑定到工作流执行实例。
|
||||||
*
|
*
|
||||||
* @param requestId 上传请求 ID
|
* @param uploadId 内部上传 ID
|
||||||
* @param executeId 工作流执行 ID
|
* @param executeId 工作流执行 ID
|
||||||
*/
|
*/
|
||||||
public void bindExecution(String requestId, String executeId) {
|
public void bindExecution(String uploadId, String executeId) {
|
||||||
if (!StringUtils.hasText(executeId)) {
|
if (!StringUtils.hasText(executeId)) {
|
||||||
throw new IllegalArgumentException("工作流执行 ID 不能为空");
|
throw new IllegalArgumentException("工作流执行 ID 不能为空");
|
||||||
}
|
}
|
||||||
WorkflowApiUploadRecord record = find(requestId)
|
WorkflowApiUploadRecord record = find(uploadId)
|
||||||
.orElseThrow(() -> new IllegalStateException(
|
.orElseThrow(() -> new IllegalStateException(
|
||||||
"工作流临时上传记录不存在: " + requestId));
|
"工作流临时上传记录不存在: " + uploadId));
|
||||||
String previousExecuteId = record.getExecuteId();
|
String previousExecuteId = record.getExecuteId();
|
||||||
record.setExecuteId(executeId);
|
record.setExecuteId(executeId);
|
||||||
redisTemplate.execute(
|
redisTemplate.execute(
|
||||||
BIND_EXECUTION_SCRIPT,
|
BIND_EXECUTION_SCRIPT,
|
||||||
Arrays.asList(
|
Arrays.asList(
|
||||||
recordKey(requestId),
|
recordKey(uploadId),
|
||||||
executionKey(executeId),
|
executionKey(executeId),
|
||||||
StringUtils.hasText(previousExecuteId)
|
StringUtils.hasText(previousExecuteId)
|
||||||
? executionKey(previousExecuteId)
|
? executionKey(previousExecuteId)
|
||||||
: recordKey(requestId)),
|
: recordKey(uploadId)),
|
||||||
serialize(record),
|
serialize(record),
|
||||||
requestId,
|
uploadId,
|
||||||
String.valueOf(EXECUTION_INDEX_TTL.toMillis()),
|
String.valueOf(EXECUTION_INDEX_TTL.toMillis()),
|
||||||
StringUtils.hasText(previousExecuteId) ? "1" : "0");
|
StringUtils.hasText(previousExecuteId) ? "1" : "0");
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 按上传请求 ID 查找记录。
|
* 按内部上传 ID 查找记录。
|
||||||
*
|
*
|
||||||
* @param requestId 上传请求 ID
|
* @param uploadId 内部上传 ID
|
||||||
* @return 上传记录
|
* @return 上传记录
|
||||||
*/
|
*/
|
||||||
public Optional<WorkflowApiUploadRecord> find(String requestId) {
|
public Optional<WorkflowApiUploadRecord> find(String uploadId) {
|
||||||
if (!StringUtils.hasText(requestId)) {
|
if (!StringUtils.hasText(uploadId)) {
|
||||||
return Optional.empty();
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
String value = redisTemplate.opsForValue().get(recordKey(requestId));
|
String value = redisTemplate.opsForValue().get(recordKey(uploadId));
|
||||||
if (!StringUtils.hasText(value)) {
|
if (!StringUtils.hasText(value)) {
|
||||||
return Optional.empty();
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
@@ -151,7 +151,7 @@ public class WorkflowApiUploadStore {
|
|||||||
WorkflowApiUploadRecord.class));
|
WorkflowApiUploadRecord.class));
|
||||||
} catch (JsonProcessingException error) {
|
} catch (JsonProcessingException error) {
|
||||||
throw new IllegalStateException(
|
throw new IllegalStateException(
|
||||||
"读取工作流临时上传记录失败: " + requestId,
|
"读取工作流临时上传记录失败: " + uploadId,
|
||||||
error);
|
error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -167,16 +167,16 @@ public class WorkflowApiUploadStore {
|
|||||||
if (!StringUtils.hasText(executeId)) {
|
if (!StringUtils.hasText(executeId)) {
|
||||||
return Optional.empty();
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
String requestId = redisTemplate.opsForValue().get(
|
String uploadId = redisTemplate.opsForValue().get(
|
||||||
executionKey(executeId));
|
executionKey(executeId));
|
||||||
Optional<WorkflowApiUploadRecord> record = find(requestId);
|
Optional<WorkflowApiUploadRecord> record = find(uploadId);
|
||||||
if (StringUtils.hasText(requestId)
|
if (StringUtils.hasText(uploadId)
|
||||||
&& (record.isEmpty()
|
&& (record.isEmpty()
|
||||||
|| !executeId.equals(record.get().getExecuteId()))) {
|
|| !executeId.equals(record.get().getExecuteId()))) {
|
||||||
redisTemplate.execute(
|
redisTemplate.execute(
|
||||||
REMOVE_EXECUTION_INDEX_SCRIPT,
|
REMOVE_EXECUTION_INDEX_SCRIPT,
|
||||||
List.of(executionKey(executeId)),
|
List.of(executionKey(executeId)),
|
||||||
requestId);
|
uploadId);
|
||||||
return Optional.empty();
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
return record;
|
return record;
|
||||||
@@ -199,17 +199,17 @@ public class WorkflowApiUploadStore {
|
|||||||
*
|
*
|
||||||
* @param now 当前 Unix 毫秒时间戳
|
* @param now 当前 Unix 毫秒时间戳
|
||||||
* @param retryAt 领取后默认重试时间
|
* @param retryAt 领取后默认重试时间
|
||||||
* @return 领取到的上传请求 ID
|
* @return 领取到的内部上传 ID
|
||||||
*/
|
*/
|
||||||
public Optional<String> claimExpired(
|
public Optional<String> claimExpired(
|
||||||
long now,
|
long now,
|
||||||
long retryAt) {
|
long retryAt) {
|
||||||
String requestId = redisTemplate.execute(
|
String uploadId = redisTemplate.execute(
|
||||||
CLAIM_EXPIRED_SCRIPT,
|
CLAIM_EXPIRED_SCRIPT,
|
||||||
List.of(CLEANUP_INDEX),
|
List.of(CLEANUP_INDEX),
|
||||||
String.valueOf(now),
|
String.valueOf(now),
|
||||||
String.valueOf(retryAt));
|
String.valueOf(retryAt));
|
||||||
return Optional.ofNullable(requestId);
|
return Optional.ofNullable(uploadId);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -224,23 +224,23 @@ public class WorkflowApiUploadStore {
|
|||||||
redisTemplate.execute(
|
redisTemplate.execute(
|
||||||
REMOVE_SCRIPT,
|
REMOVE_SCRIPT,
|
||||||
Arrays.asList(
|
Arrays.asList(
|
||||||
recordKey(record.getRequestId()),
|
recordKey(record.getUploadId()),
|
||||||
hasExecuteId
|
hasExecuteId
|
||||||
? executionKey(record.getExecuteId())
|
? executionKey(record.getExecuteId())
|
||||||
: recordKey(record.getRequestId()),
|
: recordKey(record.getUploadId()),
|
||||||
CLEANUP_INDEX),
|
CLEANUP_INDEX),
|
||||||
record.getRequestId(),
|
record.getUploadId(),
|
||||||
hasExecuteId ? "1" : "0");
|
hasExecuteId ? "1" : "0");
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 删除已经缺少详情记录的残留清理索引。
|
* 删除已经缺少详情记录的残留清理索引。
|
||||||
*
|
*
|
||||||
* @param requestId 上传请求 ID
|
* @param uploadId 内部上传 ID
|
||||||
*/
|
*/
|
||||||
public void removeMissingIndex(String requestId) {
|
public void removeMissingIndex(String uploadId) {
|
||||||
if (StringUtils.hasText(requestId)) {
|
if (StringUtils.hasText(uploadId)) {
|
||||||
redisTemplate.opsForZSet().remove(CLEANUP_INDEX, requestId);
|
redisTemplate.opsForZSet().remove(CLEANUP_INDEX, uploadId);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -256,7 +256,7 @@ public class WorkflowApiUploadStore {
|
|||||||
} catch (JsonProcessingException error) {
|
} catch (JsonProcessingException error) {
|
||||||
throw new IllegalStateException(
|
throw new IllegalStateException(
|
||||||
"写入工作流临时上传记录失败: "
|
"写入工作流临时上传记录失败: "
|
||||||
+ record.getRequestId(),
|
+ record.getUploadId(),
|
||||||
error);
|
error);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -272,14 +272,14 @@ public class WorkflowApiUploadStore {
|
|||||||
redisTemplate.execute(
|
redisTemplate.execute(
|
||||||
SAVE_SCHEDULED_SCRIPT,
|
SAVE_SCHEDULED_SCRIPT,
|
||||||
Arrays.asList(
|
Arrays.asList(
|
||||||
recordKey(record.getRequestId()),
|
recordKey(record.getUploadId()),
|
||||||
CLEANUP_INDEX,
|
CLEANUP_INDEX,
|
||||||
hasExecuteId
|
hasExecuteId
|
||||||
? executionKey(record.getExecuteId())
|
? executionKey(record.getExecuteId())
|
||||||
: recordKey(record.getRequestId())),
|
: recordKey(record.getUploadId())),
|
||||||
serialize(record),
|
serialize(record),
|
||||||
String.valueOf(record.getCleanupAt()),
|
String.valueOf(record.getCleanupAt()),
|
||||||
record.getRequestId(),
|
record.getUploadId(),
|
||||||
hasExecuteId ? "1" : "0",
|
hasExecuteId ? "1" : "0",
|
||||||
String.valueOf(EXECUTION_INDEX_TTL.toMillis()));
|
String.valueOf(EXECUTION_INDEX_TTL.toMillis()));
|
||||||
}
|
}
|
||||||
@@ -318,19 +318,19 @@ public class WorkflowApiUploadStore {
|
|||||||
* @param record 上传记录
|
* @param record 上传记录
|
||||||
*/
|
*/
|
||||||
private void requireRecord(WorkflowApiUploadRecord record) {
|
private void requireRecord(WorkflowApiUploadRecord record) {
|
||||||
if (record == null || !StringUtils.hasText(record.getRequestId())) {
|
if (record == null || !StringUtils.hasText(record.getUploadId())) {
|
||||||
throw new IllegalArgumentException("工作流临时上传请求 ID 不能为空");
|
throw new IllegalArgumentException("工作流内部上传 ID 不能为空");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 构建记录 Redis Key。
|
* 构建记录 Redis Key。
|
||||||
*
|
*
|
||||||
* @param requestId 上传请求 ID
|
* @param uploadId 内部上传 ID
|
||||||
* @return Redis Key
|
* @return Redis Key
|
||||||
*/
|
*/
|
||||||
private String recordKey(String requestId) {
|
private String recordKey(String uploadId) {
|
||||||
return RECORD_KEY_PREFIX + requestId;
|
return RECORD_KEY_PREFIX + uploadId;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -1,5 +1,7 @@
|
|||||||
package tech.easyflow.ai.easyagentsflow.upload;
|
package tech.easyflow.ai.easyagentsflow.upload;
|
||||||
|
|
||||||
|
import org.slf4j.Logger;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
import org.springframework.beans.factory.annotation.Qualifier;
|
import org.springframework.beans.factory.annotation.Qualifier;
|
||||||
import org.springframework.stereotype.Component;
|
import org.springframework.stereotype.Component;
|
||||||
import org.springframework.util.StringUtils;
|
import org.springframework.util.StringUtils;
|
||||||
@@ -16,12 +18,14 @@ import java.util.regex.Pattern;
|
|||||||
/**
|
/**
|
||||||
* 读取经过 Public Workflow API 上传记录验证的临时文件。
|
* 读取经过 Public Workflow API 上传记录验证的临时文件。
|
||||||
*
|
*
|
||||||
* <p>外部文件描述只提供公开读取路径。仅当路径中的随机请求 ID、Redis 上传记录、
|
* <p>外部文件描述只提供公开读取路径。仅当路径中的随机上传 ID、Redis 上传记录、
|
||||||
* 完整文件 URL 和可恢复存储句柄全部匹配时,才允许绕过公网 URL 限制并直接读取物理对象。</p>
|
* 完整文件 URL 和可恢复存储句柄全部匹配时,才允许绕过公网 URL 限制并直接读取物理对象。</p>
|
||||||
*/
|
*/
|
||||||
@Component
|
@Component
|
||||||
public class WorkflowApiUploadedFileReader {
|
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 String STORAGE_PATH_PREFIX = "workflow-api-upload/";
|
||||||
private static final Pattern MANAGED_PATH_PATTERN = Pattern.compile(
|
private static final Pattern MANAGED_PATH_PATTERN = Pattern.compile(
|
||||||
"(?:^|/)workflow-api-upload/([0-9a-f]{32})/([^/]+)$");
|
"(?:^|/)workflow-api-upload/([0-9a-f]{32})/([^/]+)$");
|
||||||
@@ -54,8 +58,11 @@ public class WorkflowApiUploadedFileReader {
|
|||||||
if (managedPath == null) {
|
if (managedPath == null) {
|
||||||
return Optional.empty();
|
return Optional.empty();
|
||||||
}
|
}
|
||||||
WorkflowApiUploadRecord record = uploadStore.find(managedPath.requestId())
|
WorkflowApiUploadRecord record = uploadStore.find(managedPath.uploadId())
|
||||||
.orElseThrow(() -> new IOException("工作流上传文件已失效,请重新上传"));
|
.orElseThrow(() -> new IOException("工作流上传文件已失效,请重新上传"));
|
||||||
|
if (!managedPath.uploadId().equals(record.getUploadId())) {
|
||||||
|
throw new IOException("工作流上传文件引用与上传记录不匹配");
|
||||||
|
}
|
||||||
WorkflowApiStoredFile storedFile = record.getStoredFiles().stream()
|
WorkflowApiStoredFile storedFile = record.getStoredFiles().stream()
|
||||||
.filter(file -> file != null && filePath.equals(file.filePath()))
|
.filter(file -> file != null && filePath.equals(file.filePath()))
|
||||||
.findFirst()
|
.findFirst()
|
||||||
@@ -67,7 +74,7 @@ public class WorkflowApiUploadedFileReader {
|
|||||||
} catch (IllegalArgumentException exception) {
|
} catch (IllegalArgumentException exception) {
|
||||||
throw new IOException("工作流上传文件存储定位符无效", exception);
|
throw new IOException("工作流上传文件存储定位符无效", exception);
|
||||||
}
|
}
|
||||||
String expectedStoragePath = STORAGE_PATH_PREFIX + managedPath.requestId() + "/";
|
String expectedStoragePath = STORAGE_PATH_PREFIX + managedPath.uploadId() + "/";
|
||||||
if (!expectedStoragePath.equals(handle.getPath())
|
if (!expectedStoragePath.equals(handle.getPath())
|
||||||
|| !managedPath.filename().equals(handle.getFilename())) {
|
|| !managedPath.filename().equals(handle.getFilename())) {
|
||||||
throw new IOException("工作流上传文件存储定位与上传请求不匹配");
|
throw new IOException("工作流上传文件存储定位与上传请求不匹配");
|
||||||
@@ -75,6 +82,12 @@ public class WorkflowApiUploadedFileReader {
|
|||||||
try {
|
try {
|
||||||
return Optional.of(fileStorageService.readRecoverable(handle));
|
return Optional.of(fileStorageService.readRecoverable(handle));
|
||||||
} catch (RuntimeException exception) {
|
} catch (RuntimeException exception) {
|
||||||
|
LOG.error(
|
||||||
|
"读取工作流 API 上传文件失败, uploadId={}, requestId={}, executeId={}",
|
||||||
|
record.getUploadId(),
|
||||||
|
record.getRequestId(),
|
||||||
|
record.getExecuteId(),
|
||||||
|
exception);
|
||||||
throw new IOException("读取工作流上传文件失败", exception);
|
throw new IOException("读取工作流上传文件失败", exception);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -92,7 +105,7 @@ public class WorkflowApiUploadedFileReader {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 从 URL 或相对路径中解析受管请求 ID 与固定文件名。
|
* 从 URL 或相对路径中解析受管上传 ID 与固定文件名。
|
||||||
*
|
*
|
||||||
* @param filePath 原始文件路径
|
* @param filePath 原始文件路径
|
||||||
* @return 受管路径信息
|
* @return 受管路径信息
|
||||||
@@ -121,9 +134,9 @@ public class WorkflowApiUploadedFileReader {
|
|||||||
/**
|
/**
|
||||||
* 系统受管上传路径中的可信定位片段。
|
* 系统受管上传路径中的可信定位片段。
|
||||||
*
|
*
|
||||||
* @param requestId 随机上传请求 ID
|
* @param uploadId 随机内部上传 ID
|
||||||
* @param filename 系统生成的存储文件名
|
* @param filename 系统生成的存储文件名
|
||||||
*/
|
*/
|
||||||
private record ManagedPath(String requestId, String filename) {
|
private record ManagedPath(String uploadId, String filename) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -25,6 +25,8 @@ import java.util.Set;
|
|||||||
*/
|
*/
|
||||||
public class WorkflowApiUploadLifecycleServiceTest {
|
public class WorkflowApiUploadLifecycleServiceTest {
|
||||||
|
|
||||||
|
private static final String HTTP_REQUEST_ID = "http-request-1";
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 验证同名文件 Part 会按顺序保存并注入文件对象数组。
|
* 验证同名文件 Part 会按顺序保存并注入文件对象数组。
|
||||||
*/
|
*/
|
||||||
@@ -63,11 +65,12 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
secondHandle.encodeLocator()));
|
secondHandle.encodeLocator()));
|
||||||
|
|
||||||
WorkflowApiPreparedUpload prepared = fixture.service.prepare(
|
WorkflowApiPreparedUpload prepared = fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of("user_input", "解析"),
|
Map.of("user_input", "解析"),
|
||||||
Map.of("documents", List.of(first, second)));
|
Map.of("documents", List.of(first, second)));
|
||||||
|
|
||||||
Assert.assertNotNull(prepared.getRequestId());
|
Assert.assertNotNull(prepared.getUploadId());
|
||||||
Assert.assertEquals("解析", prepared.getVariables().get("user_input"));
|
Assert.assertEquals("解析", prepared.getVariables().get("user_input"));
|
||||||
@SuppressWarnings("unchecked")
|
@SuppressWarnings("unchecked")
|
||||||
List<Map<String, Object>> documents =
|
List<Map<String, Object>> documents =
|
||||||
@@ -85,6 +88,12 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
ArgumentCaptor.forClass(WorkflowApiUploadRecord.class);
|
ArgumentCaptor.forClass(WorkflowApiUploadRecord.class);
|
||||||
Mockito.verify(fixture.uploadStore).create(
|
Mockito.verify(fixture.uploadStore).create(
|
||||||
recordCaptor.capture());
|
recordCaptor.capture());
|
||||||
|
Assert.assertEquals(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
|
recordCaptor.getValue().getRequestId());
|
||||||
|
Assert.assertNotEquals(
|
||||||
|
recordCaptor.getValue().getRequestId(),
|
||||||
|
recordCaptor.getValue().getUploadId());
|
||||||
Assert.assertEquals(
|
Assert.assertEquals(
|
||||||
List.of("/files/first.pdf", "/files/second.docx"),
|
List.of("/files/first.pdf", "/files/second.docx"),
|
||||||
recordCaptor.getValue().getStoredFiles().stream()
|
recordCaptor.getValue().getStoredFiles().stream()
|
||||||
@@ -130,6 +139,7 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
handle.encodeLocator()));
|
handle.encodeLocator()));
|
||||||
|
|
||||||
WorkflowApiPreparedUpload prepared = fixture.service.prepare(
|
WorkflowApiPreparedUpload prepared = fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of(),
|
Map.of(),
|
||||||
Map.of("documents", List.of(file)));
|
Map.of("documents", List.of(file)));
|
||||||
@@ -187,6 +197,7 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
BusinessException exception = Assert.assertThrows(
|
BusinessException exception = Assert.assertThrows(
|
||||||
BusinessException.class,
|
BusinessException.class,
|
||||||
() -> fixture.service.prepare(
|
() -> fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of(),
|
Map.of(),
|
||||||
Map.of("documents", List.of(file))));
|
Map.of("documents", List.of(file))));
|
||||||
@@ -232,6 +243,7 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
BusinessException exception = Assert.assertThrows(
|
BusinessException exception = Assert.assertThrows(
|
||||||
BusinessException.class,
|
BusinessException.class,
|
||||||
() -> fixture.service.prepare(
|
() -> fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of(),
|
Map.of(),
|
||||||
Map.of("documents", List.of(file))));
|
Map.of("documents", List.of(file))));
|
||||||
@@ -245,6 +257,53 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
.remove(Mockito.any(WorkflowApiUploadRecord.class));
|
.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 在写入存储前被拒绝。
|
* 验证未知文件 Part 在写入存储前被拒绝。
|
||||||
*/
|
*/
|
||||||
@@ -258,6 +317,7 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
BusinessException exception = Assert.assertThrows(
|
BusinessException exception = Assert.assertThrows(
|
||||||
BusinessException.class,
|
BusinessException.class,
|
||||||
() -> fixture.service.prepare(
|
() -> fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of(),
|
Map.of(),
|
||||||
Map.of("unknown", List.of(file))));
|
Map.of("unknown", List.of(file))));
|
||||||
@@ -284,6 +344,7 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
BusinessException exception = Assert.assertThrows(
|
BusinessException exception = Assert.assertThrows(
|
||||||
BusinessException.class,
|
BusinessException.class,
|
||||||
() -> fixture.service.prepare(
|
() -> fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of(),
|
Map.of(),
|
||||||
Map.of("documents", List.of(file))));
|
Map.of("documents", List.of(file))));
|
||||||
@@ -313,6 +374,7 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
BusinessException exception = Assert.assertThrows(
|
BusinessException exception = Assert.assertThrows(
|
||||||
BusinessException.class,
|
BusinessException.class,
|
||||||
() -> fixture.service.prepare(
|
() -> fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of(),
|
Map.of(),
|
||||||
Map.of("documents", List.of(file))));
|
Map.of("documents", List.of(file))));
|
||||||
@@ -356,6 +418,7 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
IllegalStateException exception = Assert.assertThrows(
|
IllegalStateException exception = Assert.assertThrows(
|
||||||
IllegalStateException.class,
|
IllegalStateException.class,
|
||||||
() -> fixture.service.prepare(
|
() -> fixture.service.prepare(
|
||||||
|
HTTP_REQUEST_ID,
|
||||||
"flow",
|
"flow",
|
||||||
Map.of(),
|
Map.of(),
|
||||||
Map.of("documents", List.of(file))));
|
Map.of("documents", List.of(file))));
|
||||||
@@ -374,7 +437,8 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
public void abortShouldDeleteStoredFilesAndRecord() {
|
public void abortShouldDeleteStoredFilesAndRecord() {
|
||||||
Fixture fixture = fixture();
|
Fixture fixture = fixture();
|
||||||
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
||||||
record.setRequestId("request-1");
|
record.setUploadId("upload-1");
|
||||||
|
record.setRequestId(HTTP_REQUEST_ID);
|
||||||
FileStorageWriteHandle firstHandle = handle("a.pdf");
|
FileStorageWriteHandle firstHandle = handle("a.pdf");
|
||||||
FileStorageWriteHandle secondHandle = handle("b.pdf");
|
FileStorageWriteHandle secondHandle = handle("b.pdf");
|
||||||
record.setStoredFiles(List.of(
|
record.setStoredFiles(List.of(
|
||||||
@@ -391,10 +455,10 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
Mockito.any(),
|
Mockito.any(),
|
||||||
Mockito.any()))
|
Mockito.any()))
|
||||||
.thenReturn(handle);
|
.thenReturn(handle);
|
||||||
Mockito.when(fixture.uploadStore.find("request-1"))
|
Mockito.when(fixture.uploadStore.find("upload-1"))
|
||||||
.thenReturn(Optional.of(record));
|
.thenReturn(Optional.of(record));
|
||||||
|
|
||||||
fixture.service.abort("request-1");
|
fixture.service.abort("upload-1");
|
||||||
|
|
||||||
Mockito.verify(fixture.fileStorageService)
|
Mockito.verify(fixture.fileStorageService)
|
||||||
.deleteRecoverable(firstHandle);
|
.deleteRecoverable(firstHandle);
|
||||||
@@ -486,16 +550,17 @@ public class WorkflowApiUploadLifecycleServiceTest {
|
|||||||
/**
|
/**
|
||||||
* 创建包含单个临时文件的上传记录。
|
* 创建包含单个临时文件的上传记录。
|
||||||
*
|
*
|
||||||
* @param requestId 请求 ID
|
* @param uploadId 内部上传 ID
|
||||||
* @param handle 文件句柄
|
* @param handle 文件句柄
|
||||||
* @return 上传记录
|
* @return 上传记录
|
||||||
*/
|
*/
|
||||||
private WorkflowApiUploadRecord storedRecord(
|
private WorkflowApiUploadRecord storedRecord(
|
||||||
String requestId,
|
String uploadId,
|
||||||
FileStorageWriteHandle handle) {
|
FileStorageWriteHandle handle) {
|
||||||
WorkflowApiUploadRecord record =
|
WorkflowApiUploadRecord record =
|
||||||
new WorkflowApiUploadRecord();
|
new WorkflowApiUploadRecord();
|
||||||
record.setRequestId(requestId);
|
record.setUploadId(uploadId);
|
||||||
|
record.setRequestId(HTTP_REQUEST_ID);
|
||||||
record.setStoredFiles(List.of(
|
record.setStoredFiles(List.of(
|
||||||
new WorkflowApiStoredFile(
|
new WorkflowApiStoredFile(
|
||||||
"/files/" + handle.getFilename(),
|
"/files/" + handle.getFilename(),
|
||||||
|
|||||||
@@ -116,18 +116,34 @@ public class WorkflowApiUploadStoreTest {
|
|||||||
.contains("redis.call('del', KEYS[3])"));
|
.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
|
* @param executeId 执行 ID
|
||||||
* @return 上传记录
|
* @return 上传记录
|
||||||
*/
|
*/
|
||||||
private WorkflowApiUploadRecord record(
|
private WorkflowApiUploadRecord record(
|
||||||
String requestId,
|
String uploadId,
|
||||||
String executeId) {
|
String executeId) {
|
||||||
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
||||||
record.setRequestId(requestId);
|
record.setUploadId(uploadId);
|
||||||
|
record.setRequestId("http-request-1");
|
||||||
record.setExecuteId(executeId);
|
record.setExecuteId(executeId);
|
||||||
record.setCleanupAt(1_000L);
|
record.setCleanupAt(1_000L);
|
||||||
return record;
|
return record;
|
||||||
|
|||||||
@@ -18,13 +18,13 @@ import java.util.Optional;
|
|||||||
*/
|
*/
|
||||||
public class WorkflowApiUploadedFileReaderTest {
|
public class WorkflowApiUploadedFileReaderTest {
|
||||||
|
|
||||||
private static final String REQUEST_ID =
|
private static final String UPLOAD_ID =
|
||||||
"0123456789abcdef0123456789abcdef";
|
"0123456789abcdef0123456789abcdef";
|
||||||
private static final String FILENAME =
|
private static final String FILENAME =
|
||||||
"000-abcdefabcdefabcdefabcdefabcdefab.docx";
|
"000-abcdefabcdefabcdefabcdefabcdefab.docx";
|
||||||
private static final String FILE_URL =
|
private static final String FILE_URL =
|
||||||
"http://127.0.0.1:39000/easyflow/attachment/"
|
"http://127.0.0.1:39000/easyflow/attachment/"
|
||||||
+ "workflow-api-upload/" + REQUEST_ID + "/" + FILENAME;
|
+ "workflow-api-upload/" + UPLOAD_ID + "/" + FILENAME;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 验证 URL、上传记录与恢复句柄完全匹配后按固定后端读取。
|
* 验证 URL、上传记录与恢复句柄完全匹配后按固定后端读取。
|
||||||
@@ -40,7 +40,7 @@ public class WorkflowApiUploadedFileReaderTest {
|
|||||||
FileStorageWriteHandle handle = handle();
|
FileStorageWriteHandle handle = handle();
|
||||||
WorkflowApiUploadRecord record = record(FILE_URL, handle);
|
WorkflowApiUploadRecord record = record(FILE_URL, handle);
|
||||||
byte[] content = "document-content".getBytes(StandardCharsets.UTF_8);
|
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))
|
Mockito.when(fileStorageService.readRecoverable(handle))
|
||||||
.thenReturn(new ByteArrayInputStream(content));
|
.thenReturn(new ByteArrayInputStream(content));
|
||||||
|
|
||||||
@@ -81,7 +81,7 @@ public class WorkflowApiUploadedFileReaderTest {
|
|||||||
FileStorageService fileStorageService = Mockito.mock(FileStorageService.class);
|
FileStorageService fileStorageService = Mockito.mock(FileStorageService.class);
|
||||||
WorkflowApiUploadedFileReader reader =
|
WorkflowApiUploadedFileReader reader =
|
||||||
new WorkflowApiUploadedFileReader(uploadStore, fileStorageService);
|
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 exception = Assert.assertThrows(
|
||||||
IOException.class,
|
IOException.class,
|
||||||
@@ -100,7 +100,7 @@ public class WorkflowApiUploadedFileReaderTest {
|
|||||||
FileStorageService fileStorageService = Mockito.mock(FileStorageService.class);
|
FileStorageService fileStorageService = Mockito.mock(FileStorageService.class);
|
||||||
WorkflowApiUploadedFileReader reader =
|
WorkflowApiUploadedFileReader reader =
|
||||||
new WorkflowApiUploadedFileReader(uploadStore, fileStorageService);
|
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())));
|
record(FILE_URL + "?different=true", handle())));
|
||||||
|
|
||||||
IOException exception = Assert.assertThrows(
|
IOException exception = Assert.assertThrows(
|
||||||
@@ -121,7 +121,7 @@ public class WorkflowApiUploadedFileReaderTest {
|
|||||||
"local",
|
"local",
|
||||||
"",
|
"",
|
||||||
"/tmp/easyflow-test",
|
"/tmp/easyflow-test",
|
||||||
"workflow-api-upload/" + REQUEST_ID,
|
"workflow-api-upload/" + UPLOAD_ID,
|
||||||
FILENAME);
|
FILENAME);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -136,7 +136,8 @@ public class WorkflowApiUploadedFileReaderTest {
|
|||||||
String fileUrl,
|
String fileUrl,
|
||||||
FileStorageWriteHandle handle) {
|
FileStorageWriteHandle handle) {
|
||||||
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
WorkflowApiUploadRecord record = new WorkflowApiUploadRecord();
|
||||||
record.setRequestId(REQUEST_ID);
|
record.setUploadId(UPLOAD_ID);
|
||||||
|
record.setRequestId("http-request-1");
|
||||||
record.setStoredFiles(List.of(new WorkflowApiStoredFile(
|
record.setStoredFiles(List.of(new WorkflowApiStoredFile(
|
||||||
fileUrl,
|
fileUrl,
|
||||||
handle.encodeLocator())));
|
handle.encodeLocator())));
|
||||||
|
|||||||
@@ -6,12 +6,17 @@ import jakarta.servlet.http.HttpServletRequest;
|
|||||||
import jakarta.servlet.http.HttpServletResponse;
|
import jakarta.servlet.http.HttpServletResponse;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
import org.springframework.http.MediaType;
|
||||||
import org.springframework.stereotype.Component;
|
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.method.HandlerMethod;
|
||||||
import org.springframework.web.servlet.HandlerInterceptor;
|
import org.springframework.web.servlet.HandlerInterceptor;
|
||||||
import org.springframework.web.servlet.ModelAndView;
|
import org.springframework.web.servlet.ModelAndView;
|
||||||
import org.springframework.web.util.ContentCachingResponseWrapper;
|
import org.springframework.web.util.ContentCachingResponseWrapper;
|
||||||
import tech.easyflow.common.util.RequestUtil;
|
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 tech.easyflow.log.annotation.LogReporterDisabled;
|
||||||
|
|
||||||
import java.lang.reflect.Method;
|
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("EasyFlow action report -------- ").append(timestamp).append(" -------------------------\n");
|
||||||
sb.append("Request : ").append(request.getMethod())
|
sb.append("Request : ").append(request.getMethod())
|
||||||
.append(" ").append(request.getRequestURI()).append("\n");
|
.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 表单参数),脱敏
|
// 打印参数(GET / POST 表单参数),脱敏
|
||||||
Map<String, String[]> params = request.getParameterMap();
|
Map<String, String[]> params = request.getParameterMap();
|
||||||
if (!params.isEmpty()) {
|
if (!params.isEmpty() && !multipartRequest) {
|
||||||
Map<String, Object> maskedParams = new LinkedHashMap<>();
|
Map<String, Object> maskedParams = new LinkedHashMap<>();
|
||||||
for (Map.Entry<String, String[]> entry : params.entrySet()) {
|
for (Map.Entry<String, String[]> entry : params.entrySet()) {
|
||||||
String key = entry.getKey();
|
String key = entry.getKey();
|
||||||
@@ -127,9 +138,17 @@ public class ActionReportInterceptor implements HandlerInterceptor {
|
|||||||
sb.append("Params : ").append(JSON.toJSONString(maskedParams)).append("\n");
|
sb.append("Params : ").append(JSON.toJSONString(maskedParams)).append("\n");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (multipartRequest) {
|
||||||
|
sb.append("Parts : ")
|
||||||
|
.append(buildMultipartSummary(request))
|
||||||
|
.append("\n");
|
||||||
|
}
|
||||||
|
|
||||||
// ====== 读取 POST Body ======
|
// ====== 读取 POST Body ======
|
||||||
String methodStr = request.getMethod();
|
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);
|
String body = RequestUtil.readBodyString(request);
|
||||||
if (body != null && !body.trim().isEmpty()) {
|
if (body != null && !body.trim().isEmpty()) {
|
||||||
try {
|
try {
|
||||||
@@ -208,8 +227,14 @@ public class ActionReportInterceptor implements HandlerInterceptor {
|
|||||||
if (ex != null) {
|
if (ex != null) {
|
||||||
sb.append('\n')
|
sb.append('\n')
|
||||||
.append("Status : FAILED\n")
|
.append("Status : FAILED\n")
|
||||||
.append("Exception : ").append(ex.getClass().getSimpleName())
|
.append("Exception : ")
|
||||||
.append(": ").append(ex.getMessage() != null ? ex.getMessage().split("\n")[0] : "Unknown");
|
.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;
|
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<Map<String, Object>> 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<String, Object> 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<String> textPartNames = new LinkedHashSet<>(
|
||||||
|
multipart.getParameterMap().keySet());
|
||||||
|
Map<String, Object> summary = new LinkedHashMap<>();
|
||||||
|
summary.put("files", parts);
|
||||||
|
summary.put("textPartNames", textPartNames);
|
||||||
|
return JSON.toJSONString(summary);
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* 构建方法签名:methodName(paramType paramName, ...)
|
* 构建方法签名:methodName(paramType paramName, ...)
|
||||||
*/
|
*/
|
||||||
|
|||||||
@@ -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<String, MultipartFile> 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() {
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user