(
injector: Injector,
request: (ctx: ResourceParamsContext) => HttpRequest<T> | undefined,
defaultValue: T,
debugName?: string,
parse?: (value: unknown) => T,
equal?: ValueEqualityFn<unknown>,
getInitialStream?: (
request: HttpRequest<unknown> | undefined,
) => Signal<ResourceStreamItem<T>> | undefined,
)
| 357 | readonly statusCode = this._statusCode.asReadonly(); |
| 358 | |
| 359 | constructor( |
| 360 | injector: Injector, |
| 361 | request: (ctx: ResourceParamsContext) => HttpRequest<T> | undefined, |
| 362 | defaultValue: T, |
| 363 | debugName?: string, |
| 364 | parse?: (value: unknown) => T, |
| 365 | equal?: ValueEqualityFn<unknown>, |
| 366 | getInitialStream?: ( |
| 367 | request: HttpRequest<unknown> | undefined, |
| 368 | ) => Signal<ResourceStreamItem<T>> | undefined, |
| 369 | ) { |
| 370 | super( |
| 371 | request, |
| 372 | ({params: request, abortSignal}) => { |
| 373 | let sub: Subscription | undefined; |
| 374 | // In the unlikely case the request returns synchronously we want to make sure the observable |
| 375 | // is subscribe even if it isn't initialized yet. |
| 376 | let aborted = false; |
| 377 | |
| 378 | // `once: true`: calling unsubscribe() here (on abort) doesn't trigger the error/complete |
| 379 | // callbacks below, so without it this listener wouldn't get removed on the abort path. |
| 380 | const onAbort = () => { |
| 381 | aborted = true; |
| 382 | sub?.unsubscribe(); |
| 383 | }; |
| 384 | abortSignal.addEventListener('abort', onAbort, {once: true}); |
| 385 | |
| 386 | // Start off stream as undefined. |
| 387 | const stream = signal<ResourceStreamItem<T>>({value: undefined as T}); |
| 388 | let resolve: ((value: Signal<ResourceStreamItem<T>>) => void) | undefined; |
| 389 | const promise = new Promise<Signal<ResourceStreamItem<T>>>((r) => (resolve = r)); |
| 390 | |
| 391 | const send = (value: ResourceStreamItem<T>): void => { |
| 392 | stream.set(value); |
| 393 | resolve?.(stream); |
| 394 | resolve = undefined; |
| 395 | }; |
| 396 | |
| 397 | sub = this.client.request(request!).subscribe({ |
| 398 | next: (event) => { |
| 399 | switch (event.type) { |
| 400 | case HttpEventType.Response: |
| 401 | this._headers.set(event.headers); |
| 402 | this._statusCode.set(event.status); |
| 403 | try { |
| 404 | send({value: parse ? parse(event.body) : (event.body as T)}); |
| 405 | } catch (error) { |
| 406 | send({error: encapsulateResourceError(error)}); |
| 407 | } |
| 408 | break; |
| 409 | case HttpEventType.DownloadProgress: |
| 410 | this._progress.set(event); |
| 411 | break; |
| 412 | } |
| 413 | }, |
| 414 | error: (error) => { |
| 415 | if (error instanceof HttpErrorResponse) { |
| 416 | this._headers.set(error.headers); |
nothing calls this directly
no test coverage detected