68e88d0d71
This allows seamless migration to different object storage. Existing runs that have offloaded payloads/outputs will continue to use the default object store (configured using `OBJECT_STORE_*` env vars). You can add additional stores by setting new env vars: - `OBJECT_STORE_DEFAULT_PROTOCOL` this determines where new run large payloads will get stored. - If you set that you need to set new env vars for that protocol. Example: ``` OBJECT_STORE_DEFAULT_PROTOCOL=“s3" OBJECT_STORE_S3_BASE_URL=https://s3.us-east-1.amazonaws.com OBJECT_STORE_S3_ACCESS_KEY_ID=<val> OBJECT_STORE_S3_SECRET_ACCESS_KEY=<val> OBJECT_STORE_S3_REGION=us-east-1 OBJECT_STORE_S3_SERVICE=s3 ``` --------- Co-authored-by: nicktrn <55853254+nicktrn@users.noreply.github.com>
61 lines
2.0 KiB
TypeScript
61 lines
2.0 KiB
TypeScript
import { IOPacket, packetRequiresOffloading, tryCatch } from "@trigger.dev/core/v3";
|
|
import { PayloadProcessor, TriggerTaskRequest } from "../types";
|
|
import { env } from "~/env.server";
|
|
import { startActiveSpan } from "~/v3/tracer.server";
|
|
import { uploadPacketToObjectStore } from "~/v3/objectStore.server";
|
|
import { ServiceValidationError } from "~/v3/services/common.server";
|
|
|
|
export class DefaultPayloadProcessor implements PayloadProcessor {
|
|
async process(request: TriggerTaskRequest): Promise<IOPacket> {
|
|
return await startActiveSpan("handlePayloadPacket()", async (span) => {
|
|
const payload = request.body.payload;
|
|
const payloadType = request.body.options?.payloadType ?? "application/json";
|
|
|
|
const packet = this.#createPayloadPacket(payload, payloadType);
|
|
|
|
if (!packet.data) {
|
|
return packet;
|
|
}
|
|
|
|
const { needsOffloading, size } = packetRequiresOffloading(
|
|
packet,
|
|
env.TASK_PAYLOAD_OFFLOAD_THRESHOLD
|
|
);
|
|
|
|
span.setAttribute("needsOffloading", needsOffloading);
|
|
span.setAttribute("size", size);
|
|
|
|
if (!needsOffloading) {
|
|
return packet;
|
|
}
|
|
|
|
const filename = `${request.friendlyId}/payload.json`;
|
|
|
|
const [uploadError, uploadedFilename] = await tryCatch(
|
|
uploadPacketToObjectStore(filename, packet.data, packet.dataType, request.environment, env.OBJECT_STORE_DEFAULT_PROTOCOL)
|
|
);
|
|
|
|
if (uploadError) {
|
|
throw new ServiceValidationError("Failed to upload large payload to object store", 500); // This is retryable
|
|
}
|
|
|
|
return {
|
|
data: uploadedFilename!,
|
|
dataType: "application/store",
|
|
};
|
|
});
|
|
}
|
|
|
|
#createPayloadPacket(payload: any, payloadType: string): IOPacket {
|
|
if (payloadType === "application/json") {
|
|
return { data: JSON.stringify(payload), dataType: "application/json" };
|
|
}
|
|
|
|
if (typeof payload === "string") {
|
|
return { data: payload, dataType: payloadType };
|
|
}
|
|
|
|
return { dataType: payloadType };
|
|
}
|
|
}
|