From 25fe505f418d745739b77af38bc1e2ebf5be7e1d Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Wed, 1 Apr 2026 23:31:36 -0400 Subject: [PATCH] refactor(share): effectify share next --- packages/opencode/src/share/share-next.ts | 545 +++++++++++++--------- 1 file changed, 318 insertions(+), 227 deletions(-) diff --git a/packages/opencode/src/share/share-next.ts b/packages/opencode/src/share/share-next.ts index 2a11094f80..d2519b3bbf 100644 --- a/packages/opencode/src/share/share-next.ts +++ b/packages/opencode/src/share/share-next.ts @@ -1,152 +1,44 @@ -import { Bus } from "@/bus" -import { Account } from "@/account" -import { Config } from "@/config/config" -import { Provider } from "@/provider/provider" -import { ProviderID, ModelID } from "@/provider/schema" -import { Session } from "@/session" -import type { SessionID } from "@/session/schema" -import { MessageV2 } from "@/session/message-v2" -import { Database, eq } from "@/storage/db" -import { SessionShareTable } from "./share.sql" -import { Log } from "@/util/log" import type * as SDK from "@opencode-ai/sdk/v2" +import { Effect, Layer, Option, Scope, ServiceMap, Stream } from "effect" +import { Account } from "@/account" +import { Bus } from "@/bus" +import { InstanceState } from "@/effect/instance-state" +import { makeRuntime } from "@/effect/run-service" +import { Provider } from "@/provider/provider" +import { ModelID, ProviderID } from "@/provider/schema" +import { Session } from "@/session" +import { MessageV2 } from "@/session/message-v2" +import type { SessionID } from "@/session/schema" +import { Database, eq } from "@/storage/db" +import { Config } from "@/config/config" +import { Log } from "@/util/log" +import { SessionShareTable } from "./share.sql" export namespace ShareNext { const log = Log.create({ service: "share-next" }) - - type ApiEndpoints = { - create: string - sync: (shareId: string) => string - remove: (shareId: string) => string - data: (shareId: string) => string - } - - function apiEndpoints(resource: string): ApiEndpoints { - return { - create: `/api/${resource}`, - sync: (shareId) => `/api/${resource}/${shareId}/sync`, - remove: (shareId) => `/api/${resource}/${shareId}`, - data: (shareId) => `/api/${resource}/${shareId}/data`, - } - } - - const legacyApi = apiEndpoints("share") - const consoleApi = apiEndpoints("shares") - - export async function url() { - const req = await request() - return req.baseUrl - } - - export async function request(): Promise<{ - headers: Record - api: ApiEndpoints - baseUrl: string - }> { - const headers: Record = {} - - const active = await Account.active() - if (!active?.active_org_id) { - const baseUrl = await Config.get().then((x) => x.enterprise?.url ?? "https://opncd.ai") - return { headers, api: legacyApi, baseUrl } - } - - const token = await Account.token(active.id) - if (!token) { - throw new Error("No active account token available for sharing") - } - - headers["authorization"] = `Bearer ${token}` - headers["x-org-id"] = active.active_org_id - return { headers, api: consoleApi, baseUrl: active.url } - } - const disabled = process.env["OPENCODE_DISABLE_SHARE"] === "true" || process.env["OPENCODE_DISABLE_SHARE"] === "1" - export async function init() { - if (disabled) return - Bus.subscribe(Session.Event.Updated, async (evt) => { - const session = await Session.get(evt.properties.sessionID) - - await sync(session.id, [ - { - type: "session", - data: session, - }, - ]) - }) - Bus.subscribe(MessageV2.Event.Updated, async (evt) => { - const info = evt.properties.info - await sync(info.sessionID, [ - { - type: "message", - data: evt.properties.info, - }, - ]) - if (info.role === "user") { - await sync(info.sessionID, [ - { - type: "model", - data: [await Provider.getModel(info.model.providerID, info.model.modelID).then((m) => m)], - }, - ]) - } - }) - Bus.subscribe(MessageV2.Event.PartUpdated, async (evt) => { - await sync(evt.properties.part.sessionID, [ - { - type: "part", - data: evt.properties.part, - }, - ]) - }) - Bus.subscribe(Session.Event.Diff, async (evt) => { - await sync(evt.properties.sessionID, [ - { - type: "session_diff", - data: evt.properties.diff, - }, - ]) - }) + type Api = { + create: string + sync: (shareID: string) => string + remove: (shareID: string) => string + data: (shareID: string) => string } - export async function create(sessionID: SessionID) { - if (disabled) return { id: "", url: "", secret: "" } - log.info("creating share", { sessionID }) - const req = await request() - const response = await fetch(`${req.baseUrl}${req.api.create}`, { - method: "POST", - headers: { ...req.headers, "Content-Type": "application/json" }, - body: JSON.stringify({ sessionID: sessionID }), - }) - - if (!response.ok) { - const message = await response.text().catch(() => response.statusText) - throw new Error(`Failed to create share (${response.status}): ${message || response.statusText}`) - } - - const result = (await response.json()) as { id: string; url: string; secret: string } - - Database.use((db) => - db - .insert(SessionShareTable) - .values({ session_id: sessionID, id: result.id, secret: result.secret, url: result.url }) - .onConflictDoUpdate({ - target: SessionShareTable.session_id, - set: { id: result.id, secret: result.secret, url: result.url }, - }) - .run(), - ) - fullSync(sessionID) - return result + type Req = { + headers: Record + api: Api + baseUrl: string } - function get(sessionID: SessionID) { - const row = Database.use((db) => - db.select().from(SessionShareTable).where(eq(SessionShareTable.session_id, sessionID)).get(), - ) - if (!row) return - return { id: row.id, secret: row.secret, url: row.url } + type Share = { + id: string + url: string + secret: string + } + + type State = { + queue: Map; data: Map }> } type Data = @@ -171,6 +63,31 @@ export namespace ShareNext { data: SDK.Model[] } + export interface Interface { + readonly init: () => Effect.Effect + readonly url: () => Effect.Effect + readonly request: () => Effect.Effect + readonly create: (sessionID: SessionID) => Effect.Effect + readonly remove: (sessionID: SessionID) => Effect.Effect + } + + export class Service extends ServiceMap.Service()("@opencode/ShareNext") {} + + const db = (fn: (d: Parameters[0] extends (trx: infer D) => any ? D : never) => T) => + Effect.sync(() => Database.use(fn)) + + function api(resource: string): Api { + return { + create: `/api/${resource}`, + sync: (shareID) => `/api/${resource}/${shareID}/sync`, + remove: (shareID) => `/api/${resource}/${shareID}`, + data: (shareID) => `/api/${resource}/${shareID}/data`, + } + } + + const legacyApi = api("share") + const consoleApi = api("shares") + function key(item: Data) { switch (item.type) { case "session": @@ -186,102 +103,276 @@ export namespace ShareNext { } } - const queue = new Map }>() - async function sync(sessionID: SessionID, data: Data[]) { - if (disabled) return - const existing = queue.get(sessionID) - if (existing) { - for (const item of data) { - existing.data.set(key(item), item) - } - return - } + export const layer: Layer.Layer< + Service, + never, + Account.Service | Bus.Service | Config.Service | Provider.Service | Session.Service + > = Layer.effect( + Service, + Effect.gen(function* () { + const account = yield* Account.Service + const bus = yield* Bus.Service + const cfg = yield* Config.Service + const provider = yield* Provider.Service + const session = yield* Session.Service + const scope = yield* Scope.Scope - const dataMap = new Map() - for (const item of data) { - dataMap.set(key(item), item) - } + const state = yield* InstanceState.make( + Effect.fn("ShareNext.state")(function* () { + const state: State = { queue: new Map() } - const timeout = setTimeout(async () => { - const queued = queue.get(sessionID) - if (!queued) return - queue.delete(sessionID) - const share = get(sessionID) - if (!share) return + yield* Effect.addFinalizer( + Effect.sync(() => { + for (const item of state.queue.values()) { + clearTimeout(item.timeout) + } + state.queue.clear() + }), + ) - const req = await request() - const response = await fetch(`${req.baseUrl}${req.api.sync(share.id)}`, { - method: "POST", - headers: { ...req.headers, "Content-Type": "application/json" }, - body: JSON.stringify({ - secret: share.secret, - data: Array.from(queued.data.values()), + if (disabled) return state + + const watch = (def: D, fn: (evt: { properties: any }) => Effect.Effect) => + bus.subscribe(def as never).pipe( + Stream.runForEach((evt) => + fn(evt).pipe( + Effect.catchAllCause((cause) => + Effect.sync(() => { + log.error("share subscriber failed", { type: def.type, cause }) + }), + ), + ), + ), + Effect.forkScoped, + ) + + yield* watch(Session.Event.Updated, (evt) => + Effect.gen(function* () { + const info = yield* session.get(evt.properties.sessionID) + yield* sync(info.id, [{ type: "session", data: info }]) + }), + ) + yield* watch(MessageV2.Event.Updated, (evt) => + Effect.gen(function* () { + const info = evt.properties.info + yield* sync(info.sessionID, [{ type: "message", data: info }]) + if (info.role !== "user") return + const model = yield* provider.getModel(info.model.providerID, info.model.modelID) + yield* sync(info.sessionID, [{ type: "model", data: [model] }]) + }), + ) + yield* watch(MessageV2.Event.PartUpdated, (evt) => + sync(evt.properties.part.sessionID, [{ type: "part", data: evt.properties.part }]), + ) + yield* watch(Session.Event.Diff, (evt) => + sync(evt.properties.sessionID, [{ type: "session_diff", data: evt.properties.diff }]), + ) + + return state }), + ) + + const request = Effect.fn("ShareNext.request")(function* () { + const headers: Record = {} + const active = yield* account.active() + if (Option.isNone(active) || !active.value.active_org_id) { + const baseUrl = (yield* cfg.get()).enterprise?.url ?? "https://opncd.ai" + return { headers, api: legacyApi, baseUrl } satisfies Req + } + + const token = yield* account.token(active.value.id) + if (Option.isNone(token)) { + throw new Error("No active account token available for sharing") + } + + headers.authorization = `Bearer ${token.value}` + headers["x-org-id"] = active.value.active_org_id + return { headers, api: consoleApi, baseUrl: active.value.url } satisfies Req }) - if (!response.ok) { - log.warn("failed to sync share", { sessionID, shareID: share.id, status: response.status }) - } - }, 1000) - queue.set(sessionID, { timeout, data: dataMap }) + const get = Effect.fnUntraced(function* (sessionID: SessionID) { + const row = yield* db((db) => + db.select().from(SessionShareTable).where(eq(SessionShareTable.session_id, sessionID)).get(), + ) + if (!row) return + return { id: row.id, secret: row.secret, url: row.url } satisfies Share + }) + + const text = Effect.fnUntraced(function* (res: Response) { + return yield* Effect.promise(() => res.text().catch(() => res.statusText)) + }) + + const flush = Effect.fn("ShareNext.flush")(function* (sessionID: SessionID) { + if (disabled) return + const s = yield* InstanceState.get(state) + const queued = s.queue.get(sessionID) + if (!queued) return + + s.queue.delete(sessionID) + + const share = yield* get(sessionID) + if (!share) return + + const req = yield* request() + const res = yield* Effect.promise(() => + fetch(`${req.baseUrl}${req.api.sync(share.id)}`, { + method: "POST", + headers: { ...req.headers, "Content-Type": "application/json" }, + body: JSON.stringify({ + secret: share.secret, + data: Array.from(queued.data.values()), + }), + }), + ) + + if (!res.ok) { + log.warn("failed to sync share", { sessionID, shareID: share.id, status: res.status }) + } + }) + + const sync = Effect.fn("ShareNext.sync")(function* (sessionID: SessionID, data: Data[]) { + if (disabled) return + const s = yield* InstanceState.get(state) + const existing = s.queue.get(sessionID) + if (existing) { + for (const item of data) { + existing.data.set(key(item), item) + } + return + } + + const next = new Map(data.map((item) => [key(item), item])) + const timeout = setTimeout( + InstanceState.bind(() => { + void runPromise(() => flush(sessionID)).catch(() => {}) + }), + 1000, + ) + s.queue.set(sessionID, { timeout, data: next }) + }) + + const full = Effect.fn("ShareNext.full")(function* (sessionID: SessionID) { + log.info("full sync", { sessionID }) + const info = yield* session.get(sessionID) + const diffs = yield* session.diff(sessionID) + const messages = yield* Effect.sync(() => Array.from(MessageV2.stream(sessionID))) + const models = yield* Effect.forEach( + Array.from( + new Map( + messages + .filter((msg) => msg.info.role === "user") + .map((msg) => (msg.info as SDK.UserMessage).model) + .map((item) => [`${item.providerID}/${item.modelID}`, item] as const), + ).values(), + ), + (item) => provider.getModel(ProviderID.make(item.providerID), ModelID.make(item.modelID)), + { concurrency: 8 }, + ) + + yield* sync(sessionID, [ + { type: "session", data: info }, + ...messages.map((item) => ({ type: "message" as const, data: item.info })), + ...messages.flatMap((item) => item.parts.map((part) => ({ type: "part" as const, data: part }))), + { type: "session_diff", data: diffs }, + { type: "model", data: models }, + ]) + }) + + const init = Effect.fn("ShareNext.init")(function* () { + if (disabled) return + yield* InstanceState.get(state) + }) + + const url = Effect.fn("ShareNext.url")(function* () { + return (yield* request()).baseUrl + }) + + const create = Effect.fn("ShareNext.create")(function* (sessionID: SessionID) { + if (disabled) return { id: "", url: "", secret: "" } + log.info("creating share", { sessionID }) + const req = yield* request() + const res = yield* Effect.promise(() => + fetch(`${req.baseUrl}${req.api.create}`, { + method: "POST", + headers: { ...req.headers, "Content-Type": "application/json" }, + body: JSON.stringify({ sessionID }), + }), + ) + + if (!res.ok) { + const msg = yield* text(res) + throw new Error(`Failed to create share (${res.status}): ${msg || res.statusText}`) + } + + const result = (yield* Effect.promise(() => res.json())) as Share + yield* db((db) => + db + .insert(SessionShareTable) + .values({ session_id: sessionID, id: result.id, secret: result.secret, url: result.url }) + .onConflictDoUpdate({ + target: SessionShareTable.session_id, + set: { id: result.id, secret: result.secret, url: result.url }, + }) + .run(), + ) + yield* full(sessionID).pipe(Effect.ignore, Effect.forkIn(scope)) + return result + }) + + const remove = Effect.fn("ShareNext.remove")(function* (sessionID: SessionID) { + if (disabled) return + log.info("removing share", { sessionID }) + const share = yield* get(sessionID) + if (!share) return + + const req = yield* request() + const res = yield* Effect.promise(() => + fetch(`${req.baseUrl}${req.api.remove(share.id)}`, { + method: "DELETE", + headers: { ...req.headers, "Content-Type": "application/json" }, + body: JSON.stringify({ secret: share.secret }), + }), + ) + + if (!res.ok) { + const msg = yield* text(res) + throw new Error(`Failed to remove share (${res.status}): ${msg || res.statusText}`) + } + + yield* db((db) => db.delete(SessionShareTable).where(eq(SessionShareTable.session_id, sessionID)).run()) + }) + + return Service.of({ init, url, request, create, remove }) + }), + ) + + export const defaultLayer = layer.pipe( + Layer.provide(Bus.layer), + Layer.provide(Account.defaultLayer), + Layer.provide(Config.defaultLayer), + Layer.provide(Provider.defaultLayer), + Layer.provide(Session.defaultLayer), + ) + + const { runPromise } = makeRuntime(Service, defaultLayer) + + export async function init() { + return runPromise((svc) => svc.init()) + } + + export async function url() { + return runPromise((svc) => svc.url()) + } + + export async function request(): Promise { + return runPromise((svc) => svc.request()) + } + + export async function create(sessionID: SessionID) { + return runPromise((svc) => svc.create(sessionID)) } export async function remove(sessionID: SessionID) { - if (disabled) return - log.info("removing share", { sessionID }) - const share = get(sessionID) - if (!share) return - - const req = await request() - const response = await fetch(`${req.baseUrl}${req.api.remove(share.id)}`, { - method: "DELETE", - headers: { ...req.headers, "Content-Type": "application/json" }, - body: JSON.stringify({ - secret: share.secret, - }), - }) - - if (!response.ok) { - const message = await response.text().catch(() => response.statusText) - throw new Error(`Failed to remove share (${response.status}): ${message || response.statusText}`) - } - - Database.use((db) => db.delete(SessionShareTable).where(eq(SessionShareTable.session_id, sessionID)).run()) - } - - async function fullSync(sessionID: SessionID) { - log.info("full sync", { sessionID }) - const session = await Session.get(sessionID) - const diffs = await Session.diff(sessionID) - const messages = await Array.fromAsync(MessageV2.stream(sessionID)) - const models = await Promise.all( - Array.from( - new Map( - messages - .filter((m) => m.info.role === "user") - .map((m) => (m.info as SDK.UserMessage).model) - .map((m) => [`${m.providerID}/${m.modelID}`, m] as const), - ).values(), - ).map((m) => Provider.getModel(ProviderID.make(m.providerID), ModelID.make(m.modelID)).then((item) => item)), - ) - await sync(sessionID, [ - { - type: "session", - data: session, - }, - ...messages.map((x) => ({ - type: "message" as const, - data: x.info, - })), - ...messages.flatMap((x) => x.parts.map((y) => ({ type: "part" as const, data: y }))), - { - type: "session_diff", - data: diffs, - }, - { - type: "model", - data: models, - }, - ]) + return runPromise((svc) => svc.remove(sessionID)) } }