|
| 1 | +import { json } from "@remix-run/server-runtime"; |
1 | 2 | import { z } from "zod"; |
2 | | -import { $replica } from "~/db.server"; |
| 3 | +import { $replica, prisma } from "~/db.server"; |
3 | 4 | import { getRealtimeStreamInstance } from "~/services/realtime/v1StreamsGlobal.server"; |
4 | 5 | import { |
5 | 6 | createActionApiRoute, |
@@ -53,26 +54,58 @@ const { action } = createActionApiRoute( |
53 | 54 | return new Response("Target not found", { status: 404 }); |
54 | 55 | } |
55 | 56 |
|
56 | | - // Extract client ID from header, default to "default" if not provided |
57 | | - const clientId = request.headers.get("X-Client-Id") || "default"; |
58 | | - const streamVersion = request.headers.get("X-Stream-Version") || "v1"; |
59 | | - |
60 | | - if (!request.body) { |
61 | | - return new Response("No body provided", { status: 400 }); |
62 | | - } |
| 57 | + if (request.method === "PUT") { |
| 58 | + // This is the "create" endpoint |
| 59 | + const updatedRun = await prisma.taskRun.update({ |
| 60 | + where: { |
| 61 | + friendlyId: targetId, |
| 62 | + runtimeEnvironmentId: authentication.environment.id, |
| 63 | + }, |
| 64 | + data: { |
| 65 | + realtimeStreams: { |
| 66 | + push: params.streamId, |
| 67 | + }, |
| 68 | + }, |
| 69 | + select: { |
| 70 | + realtimeStreamsVersion: true, |
| 71 | + }, |
| 72 | + }); |
63 | 73 |
|
64 | | - const resumeFromChunk = request.headers.get("X-Resume-From-Chunk"); |
65 | | - const resumeFromChunkNumber = resumeFromChunk ? parseInt(resumeFromChunk, 10) : undefined; |
| 74 | + const realtimeStream = getRealtimeStreamInstance( |
| 75 | + authentication.environment, |
| 76 | + updatedRun.realtimeStreamsVersion |
| 77 | + ); |
66 | 78 |
|
67 | | - const realtimeStream = getRealtimeStreamInstance(authentication.environment, streamVersion); |
| 79 | + const { responseHeaders } = await realtimeStream.initializeStream(targetId, params.streamId); |
68 | 80 |
|
69 | | - return realtimeStream.ingestData( |
70 | | - request.body, |
71 | | - targetId, |
72 | | - params.streamId, |
73 | | - clientId, |
74 | | - resumeFromChunkNumber |
75 | | - ); |
| 81 | + return json( |
| 82 | + { |
| 83 | + version: updatedRun.realtimeStreamsVersion, |
| 84 | + }, |
| 85 | + { status: 202, headers: responseHeaders } |
| 86 | + ); |
| 87 | + } else { |
| 88 | + // Extract client ID from header, default to "default" if not provided |
| 89 | + const clientId = request.headers.get("X-Client-Id") || "default"; |
| 90 | + const streamVersion = request.headers.get("X-Stream-Version") || "v1"; |
| 91 | + |
| 92 | + if (!request.body) { |
| 93 | + return new Response("No body provided", { status: 400 }); |
| 94 | + } |
| 95 | + |
| 96 | + const resumeFromChunk = request.headers.get("X-Resume-From-Chunk"); |
| 97 | + const resumeFromChunkNumber = resumeFromChunk ? parseInt(resumeFromChunk, 10) : undefined; |
| 98 | + |
| 99 | + const realtimeStream = getRealtimeStreamInstance(authentication.environment, streamVersion); |
| 100 | + |
| 101 | + return realtimeStream.ingestData( |
| 102 | + request.body, |
| 103 | + targetId, |
| 104 | + params.streamId, |
| 105 | + clientId, |
| 106 | + resumeFromChunkNumber |
| 107 | + ); |
| 108 | + } |
76 | 109 | } |
77 | 110 | ); |
78 | 111 |
|
|
0 commit comments