import type { OperationOptionsBase, ProcessErrorArgs, ServiceBusMessage, ServiceBusMessageBatch, ServiceBusReceivedMessage, ServiceBusReceiver, ServiceBusSender } from "@azure/service-bus" import { ServiceBusClient } from "@azure/service-bus" function makeClient(url: string) { return Effect(() => new ServiceBusClient(url)).acquireRelease( client => Effect.promise(() => client.close()) ) } const Client = Tag() export const LiveServiceBusClient = (url: string) => makeClient(url).toScopedLayer(Client) function makeSender(queueName: string) { return Effect.gen(function*($) { const serviceBusClient = yield* $(Client.access) return yield* $( Effect(() => serviceBusClient.createSender(queueName)).acquireRelease( subscription => Effect.promise(() => subscription.close()) ) ) }) } export const Sender = Tag() export function LiveSender(queueName: string) { return makeSender(queueName).toScopedLayer(Sender) } function makeReceiver(queueName: string) { return Effect.gen(function*($) { const serviceBusClient = yield* $(Client.access) return yield* $( Effect(() => serviceBusClient.createReceiver(queueName)).acquireRelease( r => Effect.promise(() => r.close()) ) ) }) } export const Receiver = Tag() export function LiveReceiver(queueName: string) { return makeReceiver(queueName).toScopedLayer(Receiver) } export function sendMessages( messages: ServiceBusMessage | ServiceBusMessage[] | ServiceBusMessageBatch, options?: OperationOptionsBase ) { return Effect.gen(function*($) { const s = yield* $(Sender.access) return yield* $(Effect.promise(() => s.sendMessages(messages, options))) }) } export function subscribe(hndlr: MessageHandlers) { return Effect.gen(function*($) { const r = yield* $(Receiver.access) const env = yield* $(Effect.environment()) yield* $( Effect(() => r.subscribe({ processError: err => hndlr.processError(err) .provideEnvironment(env) .unsafeRunPromise .catch(console.error), processMessage: msg => hndlr.processMessage(msg) .provideEnvironment(env) .unsafeRunPromise // DO NOT CATCH ERRORS here as they should return to the queue! }) ).acquireRelease( subscription => Effect.promise(() => subscription.close()) ) ) }) } const SubscribeTag = Tag>>() export function Subscription(hndlr: MessageHandlers) { return subscribe(hndlr).toScopedLayer(SubscribeTag) } export interface MessageHandlers { /** * Handler that processes messages from service bus. * * @param message - A message received from Service Bus. */ processMessage(message: ServiceBusReceivedMessage): Effect /** * Handler that processes errors that occur during receiving. * @param args - The error and additional context to indicate where * the error originated. */ processError(args: ProcessErrorArgs): Effect }