Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/fix-subscribe-abort-controller.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
"nostream": patch
---

fix: abort in-flight streaming queries when a subscription is cancelled
12 changes: 6 additions & 6 deletions src/handlers/subscribe-message-handler.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { anyPass, equals, isNil, map, omit, propSatisfies, uniqWith } from 'ramda'
// import { addAbortSignal } from 'stream'
import { addAbortSignal } from 'stream'
import { pipeline } from 'stream/promises'

import {
Expand All @@ -24,18 +24,18 @@ import { WebSocketAdapterEvent } from '../constants/adapter'
const logger = createLogger('subscribe-message-handler')

export class SubscribeMessageHandler implements IMessageHandler, IAbortable {
//private readonly abortController: AbortController
private readonly abortController: AbortController

public constructor(
private readonly webSocket: IWebSocketAdapter,
private readonly eventRepository: IEventRepository,
private readonly settings: () => Settings,
) {
//this.abortController = new AbortController()
this.abortController = new AbortController()
}

public abort(): void {
//this.abortController.abort()
this.abortController.abort()
}

public async handleMessage(message: SubscribeMessage): Promise<void> {
Expand Down Expand Up @@ -90,11 +90,11 @@ export class SubscribeMessageHandler implements IMessageHandler, IAbortable {

const findEvents = this.eventRepository.findByFilters(filters).stream()

// const abortableFindEvents = addAbortSignal(this.abortController.signal, findEvents)
const abortableFindEvents = addAbortSignal(this.abortController.signal, findEvents)

try {
await pipeline(
findEvents,
abortableFindEvents,
streamFilter(propSatisfies(isNil, 'deleted_at')),
streamMap(toNostrEvent),
streamFilter(isTagUnexpired),
Expand Down
13 changes: 13 additions & 0 deletions test/unit/handlers/subscribe-message-handler.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -357,6 +357,19 @@ describe('SubscribeMessageHandler', () => {
await expect(promise).to.eventually.be.rejectedWith(error)
expect(destroySpy).to.have.been.called
})

it('aborts and destroys the event stream when abort() is called', async () => {
isClientSubscribedToEventStub.returns(always(true))

const destroySpy = sandbox.spy(stream, 'destroy')

const promise = (handler as any).fetchAndSend(subscriptionId, filters)

handler.abort()

await expect(promise).to.eventually.be.rejected
expect(destroySpy).to.have.been.called
})
})

describe('.isClientSubscribedToEvent', () => {
Expand Down
Loading