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
1 change: 1 addition & 0 deletions src/utils/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,5 +18,6 @@ export * from './skip-result/skip-result.util.js'
export * from './to-paginated/to-paginated.util.js'
export * from './transform-params/transform-params.util.js'
export * from './sort-query-properties/sort-query-properties.util.js'
export * from './wait-for-service-event/wait-for-service-event.util.js'
export * from './walk-query/walk-query.util.js'
export * from './zip-data-result/zip-data-result.util.js'
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
---
title: waitForServiceEvent
category: utils
---
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
import type { HookContext } from '@feathersjs/feathers'
import type { MemoryService } from '@feathersjs/memory'
import { feathers } from '@feathersjs/feathers'
import { expectTypeOf } from 'vitest'
import { waitForServiceEvent } from './wait-for-service-event.util.js'

type User = {
id: number
name: string
email: string
}

const app = feathers<{ users: MemoryService<User> }>()

it('currying returns a callable bound to the app', () => {
expectTypeOf(waitForServiceEvent(app)).toBeFunction()
// service-agnostic defaults are accepted
waitForServiceEvent(app, { timeout: false })
waitForServiceEvent(app, { timeout: 1000 })
})

it('resolves a [data, { event, context }] tuple typed by the service', async () => {
const waitForEvent = waitForServiceEvent(app)

const [data, meta] = await waitForEvent('users', 'created')
expectTypeOf(data).toEqualTypeOf<User>()
expectTypeOf(meta.context).toEqualTypeOf<HookContext>()
})

it('event is the literal union of the requested events', async () => {
const waitForEvent = waitForServiceEvent(app)

const [, single] = await waitForEvent('users', 'created')
expectTypeOf(single.event).toEqualTypeOf<'created'>()

const [, many] = await waitForEvent('users', ['created', 'patched'])
expectTypeOf(many.event).toEqualTypeOf<'created' | 'patched'>()
})

it('filter receives the record and the hook context', () => {
const waitForEvent = waitForServiceEvent(app)

waitForEvent('users', 'created', {
filter: (data, context) => {
expectTypeOf(data).toEqualTypeOf<User>()
expectTypeOf(context).toEqualTypeOf<HookContext>()
return true
},
})
})

it('rejects unknown service paths', () => {
const waitForEvent = waitForServiceEvent(app)

// @ts-expect-error - 'unknown' is not a registered service
waitForEvent('unknown', 'created')
})
145 changes: 145 additions & 0 deletions src/utils/wait-for-service-event/wait-for-service-event.util.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,145 @@
import { feathers } from '@feathersjs/feathers'
import { MemoryService } from '@feathersjs/memory'
import { waitForServiceEvent } from './wait-for-service-event.util.js'

type User = {
id: number
name: string
email?: string
}

const setup = () => {
const app = feathers<{ users: MemoryService<User, Partial<User>> }>()

app.use('users', new MemoryService({ id: 'id', startId: 1, multi: true }))

const usersService = app.service('users')
const waitForEvent = waitForServiceEvent(app)

return { app, usersService, waitForEvent }
}

describe('waitForServiceEvent', function () {
it('resolves with [data, { event, context }] when the event fires', async function () {
const { usersService, waitForEvent } = setup()

const promise = waitForEvent('users', 'created')
const created = await usersService.create({ name: 'jane' })
const [data, { event, context }] = await promise

expect(event).toBe('created')
expect(data).toEqual(created)
expect(context.method).toBe('create')
// listeners are cleaned up after resolving
expect(usersService.listenerCount('created')).toBe(0)
})

it('only resolves for events whose data passes the filter', async function () {
const { usersService, waitForEvent } = setup()

const promise = waitForEvent('users', 'created', {
filter: (user) => user.name === 'match',
})

await usersService.create({ name: 'nomatch' })
const wanted = await usersService.create({ name: 'match' })
const [data] = await promise

expect(data).toEqual(wanted)
expect(usersService.listenerCount('created')).toBe(0)
})

it('passes the hook context as the second filter argument', async function () {
const { usersService, waitForEvent } = setup()

let seenMethod: string | undefined

const promise = waitForEvent('users', 'created', {
filter: (_user, context) => {
seenMethod = context.method
return true
},
})

await usersService.create({ name: 'jane' })
await promise

expect(seenMethod).toBe('create')
})

it('resolves on the first of multiple events and reports which fired', async function () {
const { usersService, waitForEvent } = setup()

const created = await usersService.create({ name: 'jane' })

const promise = waitForEvent('users', ['created', 'patched'])
const patched = await usersService.patch(created.id, { name: 'jane2' })
const [data, { event }] = await promise

expect(event).toBe('patched')
expect(data).toEqual(patched)
expect(usersService.listenerCount('created')).toBe(0)
expect(usersService.listenerCount('patched')).toBe(0)
})

it('rejects after the timeout and detaches listeners', async function () {
const { usersService, waitForEvent } = setup()

await expect(
waitForEvent('users', 'created', { timeout: 50 }),
).rejects.toThrow(/Timeout waiting for event "created" on service "users"/)

expect(usersService.listenerCount('created')).toBe(0)
})

it('uses the timeout bound as a default when currying', async function () {
const { app } = setup()
const waitForEvent = waitForServiceEvent(app, { timeout: 50 })

await expect(waitForEvent('users', 'created')).rejects.toThrow(/Timeout/)
})

it('lets a per-call timeout override the bound default', async function () {
const { app, usersService } = setup()
const waitForEvent = waitForServiceEvent(app, { timeout: 1 })

const promise = waitForEvent('users', 'created', { timeout: false })
const created = await usersService.create({ name: 'jane' })
const [data] = await promise

expect(data).toEqual(created)
})

it('never times out when timeout is false', async function () {
const { usersService, waitForEvent } = setup()

const promise = waitForEvent('users', 'created', { timeout: false })
const created = await usersService.create({ name: 'jane' })
const [data] = await promise

expect(data).toEqual(created)
})

it('rejects when the AbortSignal aborts while waiting', async function () {
const { usersService, waitForEvent } = setup()

const ac = new AbortController()
const promise = waitForEvent('users', 'created', { signal: ac.signal })
ac.abort(new Error('stop waiting'))

await expect(promise).rejects.toThrow('stop waiting')
expect(usersService.listenerCount('created')).toBe(0)
})

it('rejects immediately when the signal is already aborted', async function () {
const { usersService, waitForEvent } = setup()

const ac = new AbortController()
ac.abort(new Error('already gone'))

await expect(
waitForEvent('users', 'created', { signal: ac.signal }),
).rejects.toThrow('already gone')
expect(usersService.listenerCount('created')).toBe(0)
})
})
177 changes: 177 additions & 0 deletions src/utils/wait-for-service-event/wait-for-service-event.util.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
import type { Application, HookContext } from '@feathersjs/feathers'
import type { KeyOf, NeverFallback } from '../../internal.utils.js'
import type { InferFindResultSingle } from '../../utility-types/infer-service-methods.js'

/**
* The standard Feathers service events, plus any custom event name a service
* might emit.
*/
export type ServiceEventName =
| 'created'
| 'updated'
| 'patched'
| 'removed'
| (string & {})

export type WaitForServiceEventOptions<Result = unknown> = {
/**
* Reject after this many milliseconds. Pass `false` to wait indefinitely.
*
* @default 5000
*/
timeout?: number | false
/**
* Only resolve for events whose data passes this predicate. Receives the
* emitted record and the `HookContext` (the second argument Feathers emits).
*/
filter?: (data: Result, context: HookContext) => boolean
/**
* Abort waiting via an `AbortSignal`. The promise rejects with the signal's
* `reason` (or a generic abort error) and all listeners are detached.
*/
signal?: AbortSignal
}

/**
* Service-agnostic defaults that can be bound once when currying with the app.
* `filter` is intentionally omitted because it depends on the per-service
* record type.
*/
export type WaitForServiceEventDefaults = Pick<
WaitForServiceEventOptions,
'timeout' | 'signal'
>

export type WaitForServiceEventResult<Event extends string, Result> = [
/** The emitted record. */
data: Result,
meta: {
/** The event that fired (one of the requested events). */
event: Event
/** The `HookContext` Feathers emitted alongside the record. */
context: HookContext
},
]

/**
* Wait for a service event to fire and resolve with the emitted record. Useful
* in tests to await the result of an asynchronous service event, a bit like
* `promisify` for Feathers events.
*
* Curried: bind the `app` (and optional defaults) once, then call the returned
* function per service/event. Resolves with a `[data, { event, context }]`
* tuple: `data` is typed as the service's record type, and `event` is the union
* of the requested events.
*
* Feathers emits events as `emit(event, record, context)` and fires one event
* per record, so each resolution carries a single record and its `HookContext`.
*
* @example
* ```ts
* import { waitForServiceEvent } from 'feathers-utils/utils'
*
* const app = feathers()
* const waitForEvent = waitForServiceEvent(app)
*
* // Wait for the next `users` record to be created.
* const [user] = await waitForEvent('users', 'created')
*
* // Wait for a specific record, with a custom timeout and filter.
* const [data, { event }] = await waitForEvent(
* 'users',
* ['created', 'patched'],
* { filter: (user) => user.email === 'jane@example.com', timeout: 1000 },
* )
* ```
*
* @see https://utils.feathersjs.com/utils/wait-for-service-event.html
*/
export function waitForServiceEvent<Services>(
app: Application<Services>,
defaultOptions?: WaitForServiceEventDefaults,
) {
return function waitForEvent<
Path extends KeyOf<Services>,
const Event extends ServiceEventName,
Service extends Services[Path] = Services[Path],
Result = NeverFallback<InferFindResultSingle<Service>, unknown>,
>(
servicePath: Path,
eventOrEvents: Event | Event[],
options?: WaitForServiceEventOptions<Result>,
): Promise<WaitForServiceEventResult<Event, Result>> {
const events = (
Array.isArray(eventOrEvents) ? eventOrEvents : [eventOrEvents]
) as Event[]

const timeout = options?.timeout ?? defaultOptions?.timeout ?? 5000
const filter = options?.filter
const signal = options?.signal ?? defaultOptions?.signal

const service = app.service(servicePath)

return new Promise<WaitForServiceEventResult<Event, Result>>(
(resolve, reject) => {
let timer: ReturnType<typeof setTimeout> | undefined

// [event, listener] pairs so we know which event fired (Node's
// EventEmitter does not pass the event name to the listener) and can
// detach each listener precisely on cleanup.
const listeners = events.map((event) => {
const listener = (data: Result, context: HookContext) => {
if (filter && !filter(data, context)) {
return
}
cleanup()
resolve([data, { event, context }])
}
return [event, listener] as const
})

const abortError = () =>
signal?.reason ??
new Error(
`Aborted waiting for event "${events.join(', ')}" on service "${String(servicePath)}"`,
)

const onAbort = () => {
cleanup()
reject(abortError())
}

function cleanup() {
if (timer) {
clearTimeout(timer)
timer = undefined
}
for (const [event, listener] of listeners) {
;(service as any).off(event, listener)
}
signal?.removeEventListener('abort', onAbort)
}

if (signal?.aborted) {
reject(abortError())
return
}

if (timeout !== false) {
timer = setTimeout(() => {
cleanup()
reject(
new Error(
`Timeout waiting for event "${events.join(', ')}" on service "${String(servicePath)}"`,
),
)
}, timeout)
}

signal?.addEventListener('abort', onAbort, { once: true })

for (const [event, listener] of listeners) {
;(service as any).on(event, listener)
}
},
)
}
}
Loading
Loading