index.ts 6.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221
  1. import { Deferred, Effect, Layer, Schema, ServiceMap } from "effect"
  2. import { Bus } from "@/bus"
  3. import { BusEvent } from "@/bus/bus-event"
  4. import { InstanceState } from "@/effect/instance-state"
  5. import { makeRunPromise } from "@/effect/run-service"
  6. import { SessionID, MessageID } from "@/session/schema"
  7. import { Log } from "@/util/log"
  8. import z from "zod"
  9. import { QuestionID } from "./schema"
  10. export namespace Question {
  11. const log = Log.create({ service: "question" })
  12. // Schemas
  13. export const Option = z
  14. .object({
  15. label: z.string().describe("Display text (1-5 words, concise)"),
  16. description: z.string().describe("Explanation of choice"),
  17. })
  18. .meta({ ref: "QuestionOption" })
  19. export type Option = z.infer<typeof Option>
  20. export const Info = z
  21. .object({
  22. question: z.string().describe("Complete question"),
  23. header: z.string().describe("Very short label (max 30 chars)"),
  24. options: z.array(Option).describe("Available choices"),
  25. multiple: z.boolean().optional().describe("Allow selecting multiple choices"),
  26. custom: z.boolean().optional().describe("Allow typing a custom answer (default: true)"),
  27. })
  28. .meta({ ref: "QuestionInfo" })
  29. export type Info = z.infer<typeof Info>
  30. export const Request = z
  31. .object({
  32. id: QuestionID.zod,
  33. sessionID: SessionID.zod,
  34. questions: z.array(Info).describe("Questions to ask"),
  35. tool: z
  36. .object({
  37. messageID: MessageID.zod,
  38. callID: z.string(),
  39. })
  40. .optional(),
  41. })
  42. .meta({ ref: "QuestionRequest" })
  43. export type Request = z.infer<typeof Request>
  44. export const Answer = z.array(z.string()).meta({ ref: "QuestionAnswer" })
  45. export type Answer = z.infer<typeof Answer>
  46. export const Reply = z.object({
  47. answers: z
  48. .array(Answer)
  49. .describe("User answers in order of questions (each answer is an array of selected labels)"),
  50. })
  51. export type Reply = z.infer<typeof Reply>
  52. export const Event = {
  53. Asked: BusEvent.define("question.asked", Request),
  54. Replied: BusEvent.define(
  55. "question.replied",
  56. z.object({
  57. sessionID: SessionID.zod,
  58. requestID: QuestionID.zod,
  59. answers: z.array(Answer),
  60. }),
  61. ),
  62. Rejected: BusEvent.define(
  63. "question.rejected",
  64. z.object({
  65. sessionID: SessionID.zod,
  66. requestID: QuestionID.zod,
  67. }),
  68. ),
  69. }
  70. export class RejectedError extends Schema.TaggedErrorClass<RejectedError>()("QuestionRejectedError", {}) {
  71. override get message() {
  72. return "The user dismissed this question"
  73. }
  74. }
  75. interface PendingEntry {
  76. info: Request
  77. deferred: Deferred.Deferred<Answer[], RejectedError>
  78. }
  79. interface State {
  80. pending: Map<QuestionID, PendingEntry>
  81. }
  82. // Service
  83. export interface Interface {
  84. readonly ask: (input: {
  85. sessionID: SessionID
  86. questions: Info[]
  87. tool?: { messageID: MessageID; callID: string }
  88. }) => Effect.Effect<Answer[], RejectedError>
  89. readonly reply: (input: { requestID: QuestionID; answers: Answer[] }) => Effect.Effect<void>
  90. readonly reject: (requestID: QuestionID) => Effect.Effect<void>
  91. readonly list: () => Effect.Effect<Request[]>
  92. }
  93. export class Service extends ServiceMap.Service<Service, Interface>()("@opencode/Question") {}
  94. export const layer = Layer.effect(
  95. Service,
  96. Effect.gen(function* () {
  97. const state = yield* InstanceState.make<State>(
  98. Effect.fn("Question.state")(function* () {
  99. const state = {
  100. pending: new Map<QuestionID, PendingEntry>(),
  101. }
  102. yield* Effect.addFinalizer(() =>
  103. Effect.gen(function* () {
  104. for (const item of state.pending.values()) {
  105. yield* Deferred.fail(item.deferred, new RejectedError())
  106. }
  107. state.pending.clear()
  108. }),
  109. )
  110. return state
  111. }),
  112. )
  113. const ask = Effect.fn("Question.ask")(function* (input: {
  114. sessionID: SessionID
  115. questions: Info[]
  116. tool?: { messageID: MessageID; callID: string }
  117. }) {
  118. const pending = (yield* InstanceState.get(state)).pending
  119. const id = QuestionID.ascending()
  120. log.info("asking", { id, questions: input.questions.length })
  121. const deferred = yield* Deferred.make<Answer[], RejectedError>()
  122. const info: Request = {
  123. id,
  124. sessionID: input.sessionID,
  125. questions: input.questions,
  126. tool: input.tool,
  127. }
  128. pending.set(id, { info, deferred })
  129. Bus.publish(Event.Asked, info)
  130. return yield* Effect.ensuring(
  131. Deferred.await(deferred),
  132. Effect.sync(() => {
  133. pending.delete(id)
  134. }),
  135. )
  136. })
  137. const reply = Effect.fn("Question.reply")(function* (input: { requestID: QuestionID; answers: Answer[] }) {
  138. const pending = (yield* InstanceState.get(state)).pending
  139. const existing = pending.get(input.requestID)
  140. if (!existing) {
  141. log.warn("reply for unknown request", { requestID: input.requestID })
  142. return
  143. }
  144. pending.delete(input.requestID)
  145. log.info("replied", { requestID: input.requestID, answers: input.answers })
  146. Bus.publish(Event.Replied, {
  147. sessionID: existing.info.sessionID,
  148. requestID: existing.info.id,
  149. answers: input.answers,
  150. })
  151. yield* Deferred.succeed(existing.deferred, input.answers)
  152. })
  153. const reject = Effect.fn("Question.reject")(function* (requestID: QuestionID) {
  154. const pending = (yield* InstanceState.get(state)).pending
  155. const existing = pending.get(requestID)
  156. if (!existing) {
  157. log.warn("reject for unknown request", { requestID })
  158. return
  159. }
  160. pending.delete(requestID)
  161. log.info("rejected", { requestID })
  162. Bus.publish(Event.Rejected, {
  163. sessionID: existing.info.sessionID,
  164. requestID: existing.info.id,
  165. })
  166. yield* Deferred.fail(existing.deferred, new RejectedError())
  167. })
  168. const list = Effect.fn("Question.list")(function* () {
  169. const pending = (yield* InstanceState.get(state)).pending
  170. return Array.from(pending.values(), (x) => x.info)
  171. })
  172. return Service.of({ ask, reply, reject, list })
  173. }),
  174. )
  175. const runPromise = makeRunPromise(Service, layer)
  176. export async function ask(input: {
  177. sessionID: SessionID
  178. questions: Info[]
  179. tool?: { messageID: MessageID; callID: string }
  180. }): Promise<Answer[]> {
  181. return runPromise((s) => s.ask(input))
  182. }
  183. export async function reply(input: { requestID: QuestionID; answers: Answer[] }) {
  184. return runPromise((s) => s.reply(input))
  185. }
  186. export async function reject(requestID: QuestionID) {
  187. return runPromise((s) => s.reject(requestID))
  188. }
  189. export async function list() {
  190. return runPromise((s) => s.list())
  191. }
  192. }