question.ts 4.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153
  1. export * as QuestionV2 from "./question"
  2. import { makeLocationNode } from "./effect/app-node"
  3. import { Context, Deferred, Effect, Layer, Schema } from "effect"
  4. import { Question } from "@kirincode-ai/schema/question"
  5. import { EventV2 } from "./event"
  6. import { SessionSchema } from "./session/schema"
  7. export const ID = Question.ID
  8. export type ID = typeof ID.Type
  9. export const Option = Question.Option
  10. export type Option = typeof Option.Type
  11. export const Info = Question.Info
  12. export type Info = typeof Info.Type
  13. export const Prompt = Question.Prompt
  14. export type Prompt = typeof Prompt.Type
  15. export const Tool = Question.Tool
  16. export type Tool = typeof Tool.Type
  17. export const Request = Question.Request
  18. export type Request = typeof Request.Type
  19. export const Answer = Question.Answer
  20. export type Answer = typeof Answer.Type
  21. export const Reply = Question.Reply
  22. export type Reply = typeof Reply.Type
  23. export const Event = Question.Event
  24. export class RejectedError extends Schema.TaggedErrorClass<RejectedError>()("QuestionV2.RejectedError", {}) {
  25. override get message() {
  26. return "The user dismissed this question"
  27. }
  28. }
  29. export class NotFoundError extends Schema.TaggedErrorClass<NotFoundError>()("QuestionV2.NotFoundError", {
  30. requestID: ID,
  31. }) {}
  32. export interface AskInput {
  33. readonly sessionID: SessionSchema.ID
  34. readonly questions: ReadonlyArray<Info>
  35. readonly tool?: Tool
  36. }
  37. export interface ReplyInput {
  38. readonly requestID: ID
  39. readonly answers: ReadonlyArray<Answer>
  40. }
  41. export interface Interface {
  42. readonly ask: (input: AskInput) => Effect.Effect<ReadonlyArray<Answer>, RejectedError>
  43. readonly reply: (input: ReplyInput) => Effect.Effect<void, NotFoundError>
  44. readonly reject: (requestID: ID) => Effect.Effect<void, NotFoundError>
  45. readonly list: () => Effect.Effect<ReadonlyArray<Request>>
  46. }
  47. export class Service extends Context.Service<Service, Interface>()("@kirincode/v2/Question") {}
  48. interface Pending {
  49. readonly request: Request
  50. readonly deferred: Deferred.Deferred<ReadonlyArray<Answer>, RejectedError>
  51. }
  52. /**
  53. * Location-owned pending prompts. The Location layer map must materialize this
  54. * layer once per embedded Location so replies cannot settle another Location's
  55. * deferred request.
  56. */
  57. const layer = Layer.effect(
  58. Service,
  59. Effect.gen(function* () {
  60. const events = yield* EventV2.Service
  61. const pending = new Map<ID, Pending>()
  62. yield* Effect.addFinalizer(() =>
  63. Effect.forEach(pending.values(), (item) => Deferred.fail(item.deferred, new RejectedError()), {
  64. discard: true,
  65. }).pipe(
  66. Effect.ensuring(
  67. Effect.sync(() => {
  68. pending.clear()
  69. }),
  70. ),
  71. ),
  72. )
  73. const ask = Effect.fn("QuestionV2.ask")((input: AskInput) =>
  74. Effect.uninterruptibleMask((restore) =>
  75. Effect.gen(function* () {
  76. const id = ID.ascending()
  77. const deferred = yield* Deferred.make<ReadonlyArray<Answer>, RejectedError>()
  78. const request: Request = { id, ...input }
  79. pending.set(id, { request, deferred })
  80. return yield* events.publish(Event.Asked, request).pipe(
  81. Effect.andThen(restore(Deferred.await(deferred))),
  82. Effect.ensuring(
  83. Effect.sync(() => {
  84. pending.delete(id)
  85. }),
  86. ),
  87. )
  88. }),
  89. ),
  90. )
  91. const reply = Effect.fn("QuestionV2.reply")((input: ReplyInput) =>
  92. Effect.uninterruptible(
  93. Effect.gen(function* () {
  94. const existing = pending.get(input.requestID)
  95. if (!existing) return yield* new NotFoundError({ requestID: input.requestID })
  96. yield* events.publish(Event.Replied, {
  97. sessionID: existing.request.sessionID,
  98. requestID: existing.request.id,
  99. answers: input.answers.map((answer) => [...answer]),
  100. })
  101. yield* Deferred.succeed(existing.deferred, input.answers)
  102. pending.delete(input.requestID)
  103. }),
  104. ),
  105. )
  106. const reject = Effect.fn("QuestionV2.reject")((requestID: ID) =>
  107. Effect.uninterruptible(
  108. Effect.gen(function* () {
  109. const existing = pending.get(requestID)
  110. if (!existing) return yield* new NotFoundError({ requestID })
  111. yield* events.publish(Event.Rejected, {
  112. sessionID: existing.request.sessionID,
  113. requestID: existing.request.id,
  114. })
  115. yield* Deferred.fail(existing.deferred, new RejectedError())
  116. pending.delete(requestID)
  117. }),
  118. ),
  119. )
  120. const list = Effect.fn("QuestionV2.list")(function* () {
  121. return Array.from(pending.values(), (item) => item.request)
  122. })
  123. return Service.of({ ask, reply, reject, list })
  124. }),
  125. )
  126. export const locationLayer = layer
  127. export const node = makeLocationNode({ service: Service, layer, deps: [EventV2.node] })