import { v } from "convex/values"; import { mutation, query, internalMutation } from "./_generated/server"; import { inputRequestStatusValidator, inputTypeValidator, inputConfigValidator, } from "./schema"; // Return validator for input requests const inputRequestReturnValidator = v.object({ _id: v.id("feedbackInputRequests"), _creationTime: v.number(), threadId: v.string(), toolCallId: v.string(), status: inputRequestStatusValidator, inputType: inputTypeValidator, prompt: v.string(), config: v.optional(inputConfigValidator), response: v.optional(v.string()), respondedAt: v.optional(v.number()), agentName: v.string(), userId: v.optional(v.string()), createdAt: v.number(), expiresAt: v.optional(v.number()), }); /** * Internal mutation to create an input request (called by agent tools) */ export const createRequest = internalMutation({ args: { threadId: v.string(), toolCallId: v.string(), inputType: inputTypeValidator, prompt: v.string(), config: v.optional(inputConfigValidator), agentName: v.string(), userId: v.optional(v.string()), expiresInMs: v.optional(v.number()), }, returns: v.id("feedbackInputRequests"), handler: async (ctx, args) => { const now = Date.now(); const expiresAt = args.expiresInMs ? now + args.expiresInMs : undefined; const requestId = await ctx.db.insert("feedbackInputRequests", { threadId: args.threadId, toolCallId: args.toolCallId, status: "pending", inputType: args.inputType, prompt: args.prompt, config: args.config, agentName: args.agentName, userId: args.userId, createdAt: now, expiresAt, }); return requestId; }, }); /** * Get a specific input request by ID */ export const get = query({ args: { requestId: v.id("feedbackInputRequests"), }, returns: v.union(inputRequestReturnValidator, v.null()), handler: async (ctx, args) => { return await ctx.db.get(args.requestId); }, }); /** * Get the pending input request for a thread */ export const getPendingForThread = query({ args: { threadId: v.string(), }, returns: v.union(inputRequestReturnValidator, v.null()), handler: async (ctx, args) => { const request = await ctx.db .query("feedbackInputRequests") .withIndex("by_thread_status", (q) => q.eq("threadId", args.threadId).eq("status", "pending") ) .first(); return request; }, }); /** * Submit a response to an input request */ export const submitResponse = mutation({ args: { requestId: v.id("feedbackInputRequests"), response: v.string(), // JSON stringified response }, returns: v.object({ success: v.boolean(), threadId: v.string(), toolCallId: v.string(), }), handler: async (ctx, args) => { const request = await ctx.db.get(args.requestId); if (!request) { throw new Error("Input request not found"); } if (request.status !== "pending") { throw new Error(`Input request is already ${request.status}`); } // Check expiration if (request.expiresAt && Date.now() > request.expiresAt) { await ctx.db.patch(args.requestId, { status: "expired" }); throw new Error("Input request has expired"); } // Update the request with the response await ctx.db.patch(args.requestId, { status: "answered", response: args.response, respondedAt: Date.now(), }); return { success: true, threadId: request.threadId, toolCallId: request.toolCallId, }; }, }); /** * Cancel a pending input request */ export const cancel = mutation({ args: { requestId: v.id("feedbackInputRequests"), }, returns: v.null(), handler: async (ctx, args) => { const request = await ctx.db.get(args.requestId); if (!request) { throw new Error("Input request not found"); } if (request.status !== "pending") { throw new Error(`Cannot cancel request with status: ${request.status}`); } await ctx.db.patch(args.requestId, { status: "cancelled" }); return null; }, });