import { Duration, Effect, HashMap, Layer, Option, PubSub, Queue, Schedule, Scope, Stream, SynchronizedRef, pipe, } from 'effect'; import { EventStreamId, EventStoreError, EventStoreResourceError, eventStoreError, } from '@codeforbreakfast/eventsourcing-store'; interface SubscriptionData { readonly pubsub: PubSub.PubSub; readonly subscribers: number; } export interface InMemorySubscriptionManagerService { readonly subscribeToStream: ( streamId: EventStreamId ) => Effect.Effect, EventStoreError, Scope.Scope>; readonly unsubscribeFromStream: ( streamId: EventStreamId ) => Effect.Effect; readonly getSubscriptionMetrics: () => Effect.Effect< { readonly activeStreams: number; readonly totalSubscribers: number; }, never, never >; } export class InMemorySubscriptionManager extends Effect.Tag('InMemorySubscriptionManager')< InMemorySubscriptionManager, InMemorySubscriptionManagerService >() {} const addSubscriptionDataToMap = (streamId: EventStreamId, pubsub: PubSub.PubSub) => (subs: HashMap.HashMap>) => HashMap.set(subs, streamId, { pubsub, subscribers: 0 }); const addPubsubToSubs = (subs: HashMap.HashMap>, streamId: EventStreamId) => (pubsub: PubSub.PubSub) => pipe(subs, addSubscriptionDataToMap(streamId, pubsub)); const createPubSubAndAddToMap = ( subs: HashMap.HashMap>, streamId: EventStreamId ) => pipe(512, PubSub.bounded, Effect.map(addPubsubToSubs(subs, streamId)), Effect.runSync); const addSubscriptionIfMissing = (streamId: EventStreamId) => (subs: HashMap.HashMap>) => pipe( subs, HashMap.get(streamId), Option.match({ onNone: () => createPubSubAndAddToMap(subs, streamId), onSome: () => subs, }) ); const extractSubscriptionData = (streamId: EventStreamId) => (subscriptions: HashMap.HashMap>) => pipe( subscriptions, HashMap.get(streamId), Option.match({ onNone: () => Effect.fail( new EventStoreResourceError({ resource: `subscription for stream ${streamId}`, operation: 'create', cause: 'Failed to create subscription data', }) ), onSome: Effect.succeed, }) ); const getOrCreateSubscription = ( ref: SynchronizedRef.SynchronizedRef>>, streamId: EventStreamId ): Effect.Effect, EventStoreResourceError, never> => pipe( SynchronizedRef.updateAndGet(ref, addSubscriptionIfMissing(streamId)), Effect.flatMap(extractSubscriptionData(streamId)) ); const updateSubscribersCount = (streamId: EventStreamId, delta: number) => (subscriptions: HashMap.HashMap>) => HashMap.modify(subscriptions, streamId, (data) => ({ ...data, subscribers: Math.max(0, data.subscribers + delta), })); const incrementSubscribers = ( ref: SynchronizedRef.SynchronizedRef>>, streamId: EventStreamId ): Effect.Effect => pipe(SynchronizedRef.update(ref, updateSubscribersCount(streamId, 1))); const decrementSubscribers = ( ref: SynchronizedRef.SynchronizedRef>>, streamId: EventStreamId ): Effect.Effect => pipe(SynchronizedRef.update(ref, updateSubscribersCount(streamId, -1))); const filterActiveSubscriptions = ( subscriptions: HashMap.HashMap> ) => HashMap.filter(subscriptions, (data) => data.subscribers > 0); const cleanupUnusedSubscriptions = ( ref: SynchronizedRef.SynchronizedRef>> ): Effect.Effect => pipe(SynchronizedRef.update(ref, filterActiveSubscriptions)); const createRetrySchedule = () => pipe( Schedule.exponential(Duration.millis(100), 1.5), Schedule.whileOutput((d) => Duration.toMillis(d) < 30000) ); const createStreamFromQueue = (queue: Queue.Dequeue) => pipe( Stream.fromQueue(queue, { shutdown: true }), Stream.map(String), Stream.retry(createRetrySchedule()) ); const decrementAndCleanup = ( ref: SynchronizedRef.SynchronizedRef>>, streamId: EventStreamId ) => pipe( decrementSubscribers(ref, streamId), Effect.tap(() => cleanupUnusedSubscriptions(ref)) ); const subscribeToQueue = (subData: SubscriptionData) => pipe(subData.pubsub, PubSub.subscribe, Effect.map(createStreamFromQueue)); const createStreamWithCleanup = ( ref: SynchronizedRef.SynchronizedRef>>, streamId: EventStreamId ) => (subData: SubscriptionData) => pipe(subData, subscribeToQueue, Effect.ensuring(decrementAndCleanup(ref, streamId))); const createSubscribeError = (streamId: EventStreamId) => eventStoreError.subscribe(streamId, `Failed to subscribe to stream: ${String(streamId)}`); const createUnsubscribeError = (streamId: EventStreamId) => eventStoreError.subscribe(streamId, `Failed to unsubscribe from stream: ${String(streamId)}`); const subscribeToStreamEffect = (streamId: EventStreamId) => (ref: SynchronizedRef.SynchronizedRef>>) => pipe( getOrCreateSubscription(ref, streamId), Effect.tap(() => incrementSubscribers(ref, streamId)), Effect.flatMap(createStreamWithCleanup(ref, streamId)), Effect.mapError(createSubscribeError(streamId)) ); const unsubscribeFromStreamEffect = (streamId: EventStreamId) => (ref: SynchronizedRef.SynchronizedRef>>) => pipe(decrementAndCleanup(ref, streamId), Effect.mapError(createUnsubscribeError(streamId))); const calculateTotalSubscribers = ( subscriptions: HashMap.HashMap> ) => pipe(subscriptions, HashMap.values, (values) => Array.from(values).reduce((sum, data) => sum + data.subscribers, 0) ); const getMetricsEffect = ( ref: SynchronizedRef.SynchronizedRef>> ) => pipe( ref, SynchronizedRef.get, Effect.map((subscriptions) => { const activeStreams = HashMap.size(subscriptions); const totalSubscribers = calculateTotalSubscribers(subscriptions); return { activeStreams, totalSubscribers }; }) ); const subscribeForManager = (ref: SynchronizedRef.SynchronizedRef>>) => (streamId: EventStreamId) => pipe(ref, subscribeToStreamEffect(streamId)); const unsubscribeForManager = (ref: SynchronizedRef.SynchronizedRef>>) => (streamId: EventStreamId) => pipe(ref, unsubscribeFromStreamEffect(streamId)); export const makeInMemorySubscriptionManager = (): Effect.Effect< InMemorySubscriptionManagerService, never, never > => pipe( HashMap.empty(), SynchronizedRef.make>>, Effect.map((ref) => ({ subscribeToStream: subscribeForManager(ref), unsubscribeFromStream: unsubscribeForManager(ref), getSubscriptionMetrics: () => getMetricsEffect(ref), })) ); export const InMemorySubscriptionManagerLive = () => Layer.effect(InMemorySubscriptionManager, makeInMemorySubscriptionManager());