mirror of
https://gitee.com/mateos/mateclaw.git
synced 2026-09-16 20:34:39 +08:00
feat(workflow): payload fs fallback for medium-size payloads
This commit is contained in:
parent
e09c62b529
commit
20884c1773
@ -2,11 +2,15 @@ package vip.mate.workflow.runtime;
|
|||||||
|
|
||||||
import com.fasterxml.jackson.core.JsonProcessingException;
|
import com.fasterxml.jackson.core.JsonProcessingException;
|
||||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||||
|
import org.springframework.beans.factory.annotation.Value;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import vip.mate.workflow.model.WorkflowPayloadEntity;
|
import vip.mate.workflow.model.WorkflowPayloadEntity;
|
||||||
import vip.mate.workflow.repository.WorkflowPayloadMapper;
|
import vip.mate.workflow.repository.WorkflowPayloadMapper;
|
||||||
|
|
||||||
|
import java.io.IOException;
|
||||||
import java.nio.charset.StandardCharsets;
|
import java.nio.charset.StandardCharsets;
|
||||||
|
import java.nio.file.Files;
|
||||||
|
import java.nio.file.Path;
|
||||||
import java.security.MessageDigest;
|
import java.security.MessageDigest;
|
||||||
import java.security.NoSuchAlgorithmException;
|
import java.security.NoSuchAlgorithmException;
|
||||||
import java.time.LocalDateTime;
|
import java.time.LocalDateTime;
|
||||||
@ -15,24 +19,48 @@ import java.util.Objects;
|
|||||||
import java.util.UUID;
|
import java.util.UUID;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Write-through facade over {@code mate_workflow_payload}. v0 stores every
|
* Write-through facade over {@code mate_workflow_payload}. Three-tier storage:
|
||||||
* 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
|
* <ul>
|
||||||
* receive a stable URI of the form {@code mwf://{workspaceId}/{uuid}} and
|
* <li><b>inline</b> (≤ {@code inlineMaxBytes}, default 256KB) — bytes go into
|
||||||
* resolve it back via {@link #readString(String)} or {@link #readBytes(String)}.
|
* {@code content_bytes}. Cheapest and lets one DB query reconstruct the
|
||||||
|
* payload.</li>
|
||||||
|
* <li><b>fs</b> (≤ {@code hardCapBytes}) — bytes go to a workspace-scoped
|
||||||
|
* file under {@code mateclaw.workflow.payload.fs.root}; the row stores
|
||||||
|
* only the relative path in {@code storage_ref}. Default for any
|
||||||
|
* deployment that hasn't enabled a configured object-storage provider.</li>
|
||||||
|
* <li>Anything above the hard cap is rejected at write time so a runaway
|
||||||
|
* fan-out can't fill the disk silently.</li>
|
||||||
|
* </ul>
|
||||||
|
*
|
||||||
|
* <p>{@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
|
@Service
|
||||||
public class PayloadStore {
|
public class PayloadStore {
|
||||||
|
|
||||||
private static final String SCHEME = "mwf://";
|
private static final String SCHEME = "mwf://";
|
||||||
private static final String STORAGE_KIND_INLINE = "inline";
|
private static final String STORAGE_KIND_INLINE = "inline";
|
||||||
|
private static final String STORAGE_KIND_FS = "fs";
|
||||||
|
|
||||||
private final WorkflowPayloadMapper payloadMapper;
|
private final WorkflowPayloadMapper payloadMapper;
|
||||||
private final ObjectMapper objectMapper;
|
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.payloadMapper = payloadMapper;
|
||||||
this.objectMapper = objectMapper;
|
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. */
|
/** 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) {
|
public String storeBytes(long workspaceId, byte[] bytes, String contentType) {
|
||||||
Objects.requireNonNull(bytes, "bytes");
|
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();
|
String uri = SCHEME + workspaceId + "/" + UUID.randomUUID();
|
||||||
|
|
||||||
WorkflowPayloadEntity row = new WorkflowPayloadEntity();
|
WorkflowPayloadEntity row = new WorkflowPayloadEntity();
|
||||||
row.setPayloadUri(uri);
|
row.setPayloadUri(uri);
|
||||||
row.setWorkspaceId(workspaceId);
|
row.setWorkspaceId(workspaceId);
|
||||||
row.setContentBytes(bytes);
|
|
||||||
row.setStorageKind(STORAGE_KIND_INLINE);
|
|
||||||
row.setContentType(contentType);
|
row.setContentType(contentType);
|
||||||
row.setSha256(sha256Hex(bytes));
|
row.setSha256(sha256Hex(bytes));
|
||||||
row.setSizeBytes((long) bytes.length);
|
row.setSizeBytes((long) bytes.length);
|
||||||
row.setCreatedAt(LocalDateTime.now());
|
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);
|
payloadMapper.insert(row);
|
||||||
return uri;
|
return uri;
|
||||||
}
|
}
|
||||||
@ -72,6 +124,14 @@ public class PayloadStore {
|
|||||||
/** Resolve a payload URI to its raw bytes; throws when the URI is unknown. */
|
/** Resolve a payload URI to its raw bytes; throws when the URI is unknown. */
|
||||||
public byte[] readBytes(String payloadUri) {
|
public byte[] readBytes(String payloadUri) {
|
||||||
WorkflowPayloadEntity row = lookup(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();
|
return row.getContentBytes() == null ? new byte[0] : row.getContentBytes();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Loading…
Reference in New Issue
Block a user