share.ts 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191
  1. import { FileDiff, Message, Model, Part, Session } from "@opencode-ai/sdk/v2"
  2. import { fn } from "@opencode-ai/util/fn"
  3. import { iife } from "@opencode-ai/util/iife"
  4. import { Identifier } from "@opencode-ai/util/identifier"
  5. import z from "zod"
  6. import { Storage } from "./storage"
  7. import { Binary } from "@opencode-ai/util/binary"
  8. export namespace Share {
  9. export const Info = z.object({
  10. id: z.string(),
  11. secret: z.string(),
  12. sessionID: z.string(),
  13. })
  14. export type Info = z.infer<typeof Info>
  15. export const Data = z.discriminatedUnion("type", [
  16. z.object({
  17. type: z.literal("session"),
  18. data: z.custom<Session>(),
  19. }),
  20. z.object({
  21. type: z.literal("message"),
  22. data: z.custom<Message>(),
  23. }),
  24. z.object({
  25. type: z.literal("part"),
  26. data: z.custom<Part>(),
  27. }),
  28. z.object({
  29. type: z.literal("session_diff"),
  30. data: z.custom<FileDiff[]>(),
  31. }),
  32. z.object({
  33. type: z.literal("model"),
  34. data: z.custom<Model[]>(),
  35. }),
  36. ])
  37. export type Data = z.infer<typeof Data>
  38. export const create = fn(z.object({ sessionID: z.string() }), async (body) => {
  39. const isTest = process.env.NODE_ENV === "test" || body.sessionID.startsWith("test_")
  40. const info: Info = {
  41. id: (isTest ? "test_" : "") + body.sessionID.slice(-8),
  42. sessionID: body.sessionID,
  43. secret: crypto.randomUUID(),
  44. }
  45. const exists = await get(info.id)
  46. if (exists) throw new Errors.AlreadyExists(info.id)
  47. await Storage.write(["share", info.id], info)
  48. return info
  49. })
  50. export async function get(id: string) {
  51. return Storage.read<Info>(["share", id])
  52. }
  53. export const remove = fn(Info.pick({ id: true, secret: true }), async (body) => {
  54. const share = await get(body.id)
  55. if (!share) throw new Errors.NotFound(body.id)
  56. if (share.secret !== body.secret) throw new Errors.InvalidSecret(body.id)
  57. await Storage.remove(["share", body.id])
  58. const list = await Storage.list({ prefix: ["share_data", body.id] })
  59. for (const item of list) {
  60. await Storage.remove(item)
  61. }
  62. })
  63. export const sync = fn(
  64. z.object({
  65. share: Info.pick({ id: true, secret: true }),
  66. data: Data.array(),
  67. }),
  68. async (input) => {
  69. const share = await get(input.share.id)
  70. if (!share) throw new Errors.NotFound(input.share.id)
  71. if (share.secret !== input.share.secret) throw new Errors.InvalidSecret(input.share.id)
  72. await Storage.write(["share_event", input.share.id, Identifier.descending()], input.data)
  73. },
  74. )
  75. type Compaction = {
  76. event?: string
  77. data: Data[]
  78. }
  79. export async function data(shareID: string) {
  80. console.log("reading compaction")
  81. const compaction: Compaction = (await Storage.read<Compaction>(["share_compaction", shareID])) ?? {
  82. data: [],
  83. event: undefined,
  84. }
  85. console.log("reading pending events")
  86. const list = await Storage.list({
  87. prefix: ["share_event", shareID],
  88. before: compaction.event,
  89. }).then((x) => x.toReversed())
  90. console.log("compacting", list.length)
  91. if (list.length > 0) {
  92. const data = await Promise.all(list.map(async (event) => await Storage.read<Data[]>(event))).then((x) => x.flat())
  93. for (const item of data) {
  94. if (!item) continue
  95. const key = (item: Data) => {
  96. switch (item.type) {
  97. case "session":
  98. return "session"
  99. case "message":
  100. return `message/${item.data.id}`
  101. case "part":
  102. return `${item.data.messageID}/${item.data.id}`
  103. case "session_diff":
  104. return "session_diff"
  105. case "model":
  106. return "model"
  107. }
  108. }
  109. const id = key(item)
  110. const result = Binary.search(compaction.data, id, key)
  111. if (result.found) {
  112. compaction.data[result.index] = item
  113. } else {
  114. compaction.data.splice(result.index, 0, item)
  115. }
  116. }
  117. compaction.event = list.at(-1)?.at(-1)
  118. await Storage.write(["share_compaction", shareID], compaction)
  119. }
  120. return compaction.data
  121. }
  122. export const syncOld = fn(
  123. z.object({
  124. share: Info.pick({ id: true, secret: true }),
  125. data: Data.array(),
  126. }),
  127. async (input) => {
  128. const share = await get(input.share.id)
  129. if (!share) throw new Errors.NotFound(input.share.id)
  130. if (share.secret !== input.share.secret) throw new Errors.InvalidSecret(input.share.id)
  131. const promises = []
  132. for (const item of input.data) {
  133. promises.push(
  134. iife(async () => {
  135. switch (item.type) {
  136. case "session":
  137. await Storage.write(["share_data", input.share.id, "session"], item.data)
  138. break
  139. case "message": {
  140. const data = item.data as Message
  141. await Storage.write(["share_data", input.share.id, "message", data.id], item.data)
  142. break
  143. }
  144. case "part": {
  145. const data = item.data as Part
  146. await Storage.write(["share_data", input.share.id, "part", data.messageID, data.id], item.data)
  147. break
  148. }
  149. case "session_diff":
  150. await Storage.write(["share_data", input.share.id, "session_diff"], item.data)
  151. break
  152. case "model":
  153. await Storage.write(["share_data", input.share.id, "model"], item.data)
  154. break
  155. }
  156. }),
  157. )
  158. }
  159. await Promise.all(promises)
  160. },
  161. )
  162. export const Errors = {
  163. NotFound: class extends Error {
  164. constructor(public id: string) {
  165. super(`Share not found: ${id}`)
  166. }
  167. },
  168. InvalidSecret: class extends Error {
  169. constructor(public id: string) {
  170. super(`Share secret invalid: ${id}`)
  171. }
  172. },
  173. AlreadyExists: class extends Error {
  174. constructor(public id: string) {
  175. super(`Share already exists: ${id}`)
  176. }
  177. },
  178. }
  179. }