From 20884c1773e213afa83b614e1cde3a1a002441d3 Mon Sep 17 00:00:00 2001 From: matevip Date: Fri, 8 May 2026 15:05:28 +0800 Subject: [PATCH] feat(workflow): payload fs fallback for medium-size payloads --- .../mate/workflow/runtime/PayloadStore.java | 78 ++++++++++++++++--- 1 file changed, 69 insertions(+), 9 deletions(-) diff --git a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/PayloadStore.java b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/PayloadStore.java index a1786a5c..31f6ed6e 100644 --- a/mateclaw-server/src/main/java/vip/mate/workflow/runtime/PayloadStore.java +++ b/mateclaw-server/src/main/java/vip/mate/workflow/runtime/PayloadStore.java @@ -2,11 +2,15 @@ package vip.mate.workflow.runtime; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import vip.mate.workflow.model.WorkflowPayloadEntity; import vip.mate.workflow.repository.WorkflowPayloadMapper; +import java.io.IOException; import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; import java.security.MessageDigest; import java.security.NoSuchAlgorithmException; import java.time.LocalDateTime; @@ -15,24 +19,48 @@ import java.util.Objects; import java.util.UUID; /** - * Write-through facade over {@code mate_workflow_payload}. v0 stores every - * payload inline (storage_kind = "inline"); the table schema reserves room for - * a fs / s3 / oss spill-over flavour but no caller wires that yet. Callers - * receive a stable URI of the form {@code mwf://{workspaceId}/{uuid}} and - * resolve it back via {@link #readString(String)} or {@link #readBytes(String)}. + * Write-through facade over {@code mate_workflow_payload}. Three-tier storage: + * + * + * + *

{@code s3} / {@code oss} columns exist in the schema but the v0 ship + * only writes {@code inline} or {@code fs}; provider configuration ships in + * v1. The fs tier is what unblocks local dev / docker / private deploys + * that don't have an object store configured. */ @Service public class PayloadStore { private static final String SCHEME = "mwf://"; private static final String STORAGE_KIND_INLINE = "inline"; + private static final String STORAGE_KIND_FS = "fs"; private final WorkflowPayloadMapper payloadMapper; private final ObjectMapper objectMapper; + private final long inlineMaxBytes; + private final long hardCapBytes; + private final Path fsRoot; - public PayloadStore(WorkflowPayloadMapper payloadMapper, ObjectMapper objectMapper) { + public PayloadStore(WorkflowPayloadMapper payloadMapper, + ObjectMapper objectMapper, + @Value("${mateclaw.workflow.payload.inline-max-bytes:262144}") long inlineMaxBytes, + @Value("${mateclaw.workflow.payload.hard-cap-bytes:52428800}") long hardCapBytes, + @Value("${mateclaw.workflow.payload.fs.root:./data/workflow-payload}") String fsRoot) { this.payloadMapper = payloadMapper; this.objectMapper = objectMapper; + this.inlineMaxBytes = inlineMaxBytes; + this.hardCapBytes = hardCapBytes; + this.fsRoot = Path.of(fsRoot).toAbsolutePath(); } /** Store a UTF-8 string payload and return its stable URI. */ @@ -51,20 +79,44 @@ public class PayloadStore { } } - /** Store raw bytes and return the URI. */ + /** Store raw bytes and return the URI. Routes by size: inline → fs → reject. */ public String storeBytes(long workspaceId, byte[] bytes, String contentType) { Objects.requireNonNull(bytes, "bytes"); + if (bytes.length > hardCapBytes) { + throw new PayloadStoreException("payload exceeds hard cap of " + + hardCapBytes + " bytes (got " + bytes.length + ")"); + } String uri = SCHEME + workspaceId + "/" + UUID.randomUUID(); WorkflowPayloadEntity row = new WorkflowPayloadEntity(); row.setPayloadUri(uri); row.setWorkspaceId(workspaceId); - row.setContentBytes(bytes); - row.setStorageKind(STORAGE_KIND_INLINE); row.setContentType(contentType); row.setSha256(sha256Hex(bytes)); row.setSizeBytes((long) bytes.length); row.setCreatedAt(LocalDateTime.now()); + + if (bytes.length <= inlineMaxBytes) { + row.setContentBytes(bytes); + row.setStorageKind(STORAGE_KIND_INLINE); + } else { + // Spill to filesystem so we don't bloat the DB row. Path layout + // is {fsRoot}/{workspaceId}/{first2chars}/{uuid} so a single + // workspace can't pile millions of files into one directory. + String relative = workspaceId + "/" + uri.substring(uri.length() - 2) + + "/" + uri.substring(uri.length() - Math.min(36, uri.length())); + Path target = fsRoot.resolve(relative); + try { + Files.createDirectories(target.getParent()); + Files.write(target, bytes); + } catch (IOException e) { + throw new PayloadStoreException("failed to write fs payload " + uri + + ": " + e.getMessage(), e); + } + row.setStorageKind(STORAGE_KIND_FS); + row.setStorageRef(relative); + } + payloadMapper.insert(row); return uri; } @@ -72,6 +124,14 @@ public class PayloadStore { /** Resolve a payload URI to its raw bytes; throws when the URI is unknown. */ public byte[] readBytes(String payloadUri) { WorkflowPayloadEntity row = lookup(payloadUri); + if (STORAGE_KIND_FS.equals(row.getStorageKind())) { + try { + return Files.readAllBytes(fsRoot.resolve(row.getStorageRef())); + } catch (IOException e) { + throw new PayloadStoreException("failed to read fs payload " + payloadUri + + ": " + e.getMessage(), e); + } + } return row.getContentBytes() == null ? new byte[0] : row.getContentBytes(); }