fix(share): correct effect typing
This commit is contained in:
@@ -64,11 +64,11 @@ export namespace ShareNext {
|
||||
}
|
||||
|
||||
export interface Interface {
|
||||
readonly init: () => Effect.Effect<void>
|
||||
readonly url: () => Effect.Effect<string>
|
||||
readonly request: () => Effect.Effect<Req>
|
||||
readonly create: (sessionID: SessionID) => Effect.Effect<Share>
|
||||
readonly remove: (sessionID: SessionID) => Effect.Effect<void>
|
||||
readonly init: () => Effect.Effect<void, unknown>
|
||||
readonly url: () => Effect.Effect<string, unknown>
|
||||
readonly request: () => Effect.Effect<Req, unknown>
|
||||
readonly create: (sessionID: SessionID) => Effect.Effect<Share, unknown>
|
||||
readonly remove: (sessionID: SessionID) => Effect.Effect<void, unknown>
|
||||
}
|
||||
|
||||
export class Service extends ServiceMap.Service<Service, Interface>()("@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<State>(
|
||||
Effect.fn("ShareNext.state")(function* () {
|
||||
const state: State = { queue: new Map() }
|
||||
function sync(sessionID: SessionID, data: Data[]): Effect.Effect<void, unknown> {
|
||||
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<State> = yield* InstanceState.make<State>(
|
||||
Effect.fn("ShareNext.state")(function* (_ctx): Effect.Effect<State, never, Scope.Scope> {
|
||||
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 = <D extends { type: string }>(def: D, fn: (evt: { properties: any }) => Effect.Effect<void>) =>
|
||||
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 })
|
||||
}),
|
||||
|
||||
Reference in New Issue
Block a user