diff --git a/packages/opencode/src/share/share-next.ts b/packages/opencode/src/share/share-next.ts index 605fb01cc3..b2d771f7cd 100644 --- a/packages/opencode/src/share/share-next.ts +++ b/packages/opencode/src/share/share-next.ts @@ -64,11 +64,11 @@ export namespace ShareNext { } 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 + 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") {} @@ -103,11 +103,7 @@ export namespace ShareNext { } } - export const layer: Layer.Layer< - Service, - never, - Account.Service | Bus.Service | Config.Service | Provider.Service | Session.Service - > = Layer.effect( + export const layer = Layer.effect( Service, Effect.gen(function* () { const account = yield* Account.Service @@ -117,26 +113,57 @@ export namespace ShareNext { const session = yield* Session.Service const scope = yield* Scope.Scope - const state = yield* InstanceState.make( - Effect.fn("ShareNext.state")(function* () { - const state: State = { queue: new Map() } + function sync(sessionID: SessionID, data: Data[]): Effect.Effect { + return Effect.gen(function* () { + 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 + } - yield* Effect.addFinalizer( + const next = new Map(data.map((item) => [key(item), item])) + const timeout = setTimeout( + InstanceState.bind(() => { + void runPromise(() => + flush(sessionID).pipe( + Effect.catchCause((cause) => + Effect.sync(() => { + log.error("share flush failed", { sessionID, cause }) + }), + ), + ), + ) + }), + 1000, + ) + s.queue.set(sessionID, { timeout, data: next }) + }) + } + + const state: InstanceState = yield* InstanceState.make( + Effect.fn("ShareNext.state")(function* (_ctx): Effect.Effect { + const cache: State = { queue: new Map() } + + yield* Effect.addFinalizer(() => Effect.sync(() => { - for (const item of state.queue.values()) { + for (const item of cache.queue.values()) { clearTimeout(item.timeout) } - state.queue.clear() + cache.queue.clear() }), ) - if (disabled) return state + if (disabled) return cache 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.catchCause((cause) => Effect.sync(() => { log.error("share subscriber failed", { type: def.type, cause }) }), @@ -168,7 +195,7 @@ export namespace ShareNext { sync(evt.properties.sessionID, [{ type: "session_diff", data: evt.properties.diff }]), ) - return state + return cache }), ) @@ -230,35 +257,6 @@ export namespace ShareNext { } }) - 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).pipe( - Effect.catchAllCause((cause) => - Effect.sync(() => { - log.error("share flush failed", { sessionID, cause }) - }), - ), - ), - ) - }), - 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) @@ -324,7 +322,7 @@ export namespace ShareNext { .run(), ) yield* full(sessionID).pipe( - Effect.catchAllCause((cause) => + Effect.catchCause((cause) => Effect.sync(() => { log.error("share full sync failed", { sessionID, cause }) }),