MCPcopy Create free account
hub / github.com/anomalyco/opencode / subscribeDurable

Function subscribeDurable

packages/core/src/event.ts:565–583  ·  view source on GitHub ↗
(aggregateID: string)

Source from the content-addressed store, hash-verified

563 )
564
565 const subscribeDurable = (aggregateID: string) =>
566 Effect.gen(function* () {
567 const wake = yield* PubSub.sliding<void>(1)
568 const subscription = yield* PubSub.subscribe(wake)
569 yield* Effect.acquireRelease(
570 Effect.sync(() => {
571 const wakes = pubsub.durable.get(aggregateID) ?? new Set()
572 wakes.add(wake)
573 pubsub.durable.set(aggregateID, wakes)
574 }),
575 () =>
576 Effect.sync(() => {
577 const wakes = pubsub.durable.get(aggregateID)
578 wakes?.delete(wake)
579 if (wakes?.size === 0) pubsub.durable.delete(aggregateID)
580 }).pipe(Effect.andThen(PubSub.shutdown(wake))),
581 )
582 return subscription
583 })
584
585 const durable = (input: { readonly aggregateID: string; readonly after?: number }): Stream.Stream<Payload> =>
586 Stream.unwrap(

Callers 1

durableFunction · 0.85

Calls 6

syncMethod · 0.80
subscribeMethod · 0.65
getMethod · 0.65
addMethod · 0.65
setMethod · 0.45
deleteMethod · 0.45

Tested by

no test coverage detected