(input: any)
| 318 | ) |
| 319 | |
| 320 | const buildStream = (input: any) => { |
| 321 | const stream = Stream.unwrap( |
| 322 | mr.contextEffect.pipe( |
| 323 | Effect.flatMap((svcs) => |
| 324 | TheClient |
| 325 | .useSync((client) => { |
| 326 | const rpcStream = (client as any)[requestAttr]!( |
| 327 | Request.make(input) |
| 328 | ) as Stream.Stream<any, any, any> |
| 329 | return rpcStream.pipe( |
| 330 | // Collect server invalidation keys from the "done" chunk, then discard it. |
| 331 | Stream.tap((item: any) => |
| 332 | item._tag === "done" || item._tag === "metadata" |
| 333 | ? Effect.gen(function*() { |
| 334 | const metadata = item.metadata as Invalidation.CommandMetaData |
| 335 | const invalidationKeys = yield* InvalidationKeysFromServer |
| 336 | yield* Effect.forEach(metadata.invalidateQueries, invalidationKeys.add, { |
| 337 | discard: true |
| 338 | }) |
| 339 | const dependencyRecorder = yield* DataDependencies.getDataDependencyRecorder |
| 340 | yield* Effect.forEach( |
| 341 | metadata.dataDependencies.reads, |
| 342 | dependencyRecorder.read, |
| 343 | { discard: true } |
| 344 | ) |
| 345 | yield* Effect.forEach( |
| 346 | metadata.dataDependencies.writes, |
| 347 | dependencyRecorder.write, |
| 348 | { discard: true } |
| 349 | ) |
| 350 | }) |
| 351 | : Effect.void |
| 352 | ), |
| 353 | Stream.filter((item: any) => item._tag === "value"), |
| 354 | Stream.map((item: any) => item.value), |
| 355 | // V2: unwrap StreamFailureChunk — forward keys from failures too, |
| 356 | // then re-fail with the original error so callers see the unmodified |
| 357 | // error type. |
| 358 | Stream.catch((err: any) => |
| 359 | err?._tag === "error" && err?.metadata |
| 360 | ? Stream.fromEffect( |
| 361 | Effect.gen(function*() { |
| 362 | const metadata = err.metadata as Invalidation.CommandMetaData |
| 363 | const invalidationKeys = yield* InvalidationKeysFromServer |
| 364 | yield* Effect.forEach(metadata.invalidateQueries, invalidationKeys.add, { |
| 365 | discard: true |
| 366 | }) |
| 367 | const dependencyRecorder = yield* DataDependencies.getDataDependencyRecorder |
| 368 | yield* Effect.forEach( |
| 369 | metadata.dataDependencies.writes, |
| 370 | dependencyRecorder.write, |
| 371 | { discard: true } |
| 372 | ) |
| 373 | return yield* Effect.fail(err.error) |
| 374 | }) |
| 375 | ) |
| 376 | : Stream.fail(err) |
| 377 | ), |
no test coverage detected
searching dependent graphs…