serverTime subscription
This commit is contained in:
@@ -1,5 +1,6 @@
|
||||
import { TRPCError } from "@trpc/server"
|
||||
import { Context, Data, Effect, Layer } from "effect"
|
||||
import { observable } from "@trpc/server/observable"
|
||||
import { Context, Data, Effect, Fiber, Layer, Schedule } from "effect"
|
||||
import { TRPCBuilder } from "../trpc/TRPCBuilder"
|
||||
import { RPCProcedureBuilder } from "./procedures/RPCProcedureBuilder"
|
||||
import { todoRouter } from "./routers/todo"
|
||||
@@ -14,6 +15,22 @@ export const router = Effect.gen(function*() {
|
||||
Effect.succeed("pong")
|
||||
)),
|
||||
|
||||
serverTime: procedure
|
||||
.subscription(({ ctx }) =>
|
||||
observable<string>(emit => {
|
||||
const emitter = ctx.fork(
|
||||
Effect.sync(() => emit.next(new Date().toString())).pipe(
|
||||
Effect.repeat(Schedule.spaced("1 second"))
|
||||
)
|
||||
)
|
||||
|
||||
return () => ctx.fork(
|
||||
Fiber.interrupt(emitter)
|
||||
)
|
||||
})
|
||||
),
|
||||
|
||||
|
||||
fail1: procedure.query(({ ctx }) => ctx.run(
|
||||
Effect.fail(new AnError({ aValue: "A value" }))
|
||||
)),
|
||||
|
||||
Reference in New Issue
Block a user