api.ts 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174
  1. import { DurableObject } from "cloudflare:workers"
  2. import { randomUUID } from "node:crypto"
  3. import { Resource } from "sst"
  4. type Bindings = {
  5. SYNC_SERVER: DurableObjectNamespace<SyncServer>
  6. }
  7. export class SyncServer extends DurableObject {
  8. async fetch(req: Request) {
  9. console.log("SyncServer subscribe")
  10. const webSocketPair = new WebSocketPair()
  11. const [client, server] = Object.values(webSocketPair)
  12. this.ctx.acceptWebSocket(server)
  13. setTimeout(async () => {
  14. const data = await this.ctx.storage.list()
  15. data.forEach((content: any, key) => {
  16. if (key === "shareID") return
  17. server.send(JSON.stringify({ key, content: content }))
  18. })
  19. }, 0)
  20. return new Response(null, {
  21. status: 101,
  22. webSocket: client,
  23. })
  24. }
  25. async webSocketMessage(ws, message) {}
  26. async webSocketClose(ws, code, reason, wasClean) {
  27. ws.close(code, "Durable Object is closing WebSocket")
  28. }
  29. async publish(key: string, content: any) {
  30. await this.ctx.storage.put(key, content)
  31. const clients = this.ctx.getWebSockets()
  32. console.log("SyncServer publish", key, "to", clients.length, "subscribers")
  33. clients.forEach((client) => client.send(JSON.stringify({ key, content })))
  34. }
  35. async setShareID(shareID: string) {
  36. await this.ctx.storage.put("shareID", shareID)
  37. }
  38. async getShareID() {
  39. return this.ctx.storage.get<string>("shareID")
  40. }
  41. async clear() {
  42. await this.ctx.storage.deleteAll()
  43. }
  44. }
  45. export default {
  46. async fetch(request: Request, env: Bindings, ctx: ExecutionContext) {
  47. const url = new URL(request.url)
  48. const splits = url.pathname.split("/")
  49. const method = splits[1]
  50. if (request.method === "GET" && method === "") {
  51. return new Response("Hello, world!", {
  52. headers: { "Content-Type": "text/plain" },
  53. })
  54. }
  55. if (request.method === "POST" && method === "share_create") {
  56. const body = await request.json<any>()
  57. const sessionID = body.sessionID
  58. // Get existing shareID
  59. const id = env.SYNC_SERVER.idFromName(sessionID)
  60. const stub = env.SYNC_SERVER.get(id)
  61. if (await stub.getShareID())
  62. return new Response("Error: Session already shared", { status: 400 })
  63. const shareID = randomUUID()
  64. await stub.setShareID(shareID)
  65. return new Response(JSON.stringify({ shareID }), {
  66. headers: { "Content-Type": "application/json" },
  67. })
  68. }
  69. if (request.method === "POST" && method === "share_delete") {
  70. const body = await request.json<any>()
  71. const sessionID = body.sessionID
  72. const shareID = body.shareID
  73. // validate shareID
  74. if (!shareID)
  75. return new Response("Error: Share ID is required", { status: 400 })
  76. // Delete from durable object
  77. const id = env.SYNC_SERVER.idFromName(sessionID)
  78. const stub = env.SYNC_SERVER.get(id)
  79. if ((await stub.getShareID()) !== shareID)
  80. return new Response("Error: Share ID does not match", { status: 400 })
  81. await stub.clear()
  82. return new Response(JSON.stringify({}), {
  83. headers: { "Content-Type": "application/json" },
  84. })
  85. }
  86. if (request.method === "POST" && method === "share_sync") {
  87. const body = await request.json<any>()
  88. const sessionID = body.sessionID
  89. const shareID = body.shareID
  90. const key = body.key
  91. const content = body.content
  92. console.log("share_sync", sessionID, shareID, key, content)
  93. // validate key
  94. if (
  95. !key.startsWith(`session/info/${sessionID}`) &&
  96. !key.startsWith(`session/message/${sessionID}/`)
  97. )
  98. return new Response("Error: Invalid key", { status: 400 })
  99. // validate shareID
  100. if (!shareID)
  101. return new Response("Error: Share ID is required", { status: 400 })
  102. // send message to server
  103. const id = env.SYNC_SERVER.idFromName(sessionID)
  104. const stub = env.SYNC_SERVER.get(id)
  105. if ((await stub.getShareID()) !== shareID)
  106. return new Response("Error: Share ID does not match", { status: 400 })
  107. await stub.publish(key, content)
  108. // store message
  109. await Resource.Bucket.put(
  110. `${shareID}/${key}.json`,
  111. JSON.stringify(content),
  112. )
  113. return new Response(JSON.stringify({}), {
  114. headers: { "Content-Type": "application/json" },
  115. })
  116. }
  117. if (request.method === "GET" && method === "share_poll") {
  118. // Expect to receive a WebSocket Upgrade request.
  119. // If there is one, accept the request and return a WebSocket Response.
  120. const upgradeHeader = request.headers.get("Upgrade")
  121. if (!upgradeHeader || upgradeHeader !== "websocket") {
  122. return new Response("Error: Upgrade header is required", {
  123. status: 426,
  124. })
  125. }
  126. // get query parameters
  127. const sessionID = url.searchParams.get("id")
  128. if (!sessionID)
  129. return new Response("Error: Share ID is required", { status: 400 })
  130. // subscribe to server
  131. const id = env.SYNC_SERVER.idFromName(sessionID)
  132. const stub = env.SYNC_SERVER.get(id)
  133. if (!(await stub.getShareID()))
  134. return new Response("Error: Session not shared", { status: 400 })
  135. return stub.fetch(request)
  136. }
  137. },
  138. }