MCPcopy Create free account
hub / github.com/effect-app/libs / buildStream

Function buildStream

packages/effect-app/src/client/apiClientFactory.ts:320–398  ·  view source on GitHub ↗
(input: any)

Source from the content-addressed store, hash-verified

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 ),

Callers 1

Calls 9

filterMethod · 0.80
catchMethod · 0.80
pipeMethod · 0.65
useSyncMethod · 0.65
makeMethod · 0.65
forEachMethod · 0.65
mapMethod · 0.65
failMethod · 0.45
provideMethod · 0.45

Tested by

no test coverage detected

Used in the wild real call sites across dependent graphs

searching dependent graphs…