import { mutation, query } from "./_generated/server"; import { v } from "convex/values"; import { requireUserId } from "./_utils/auth"; import { nowIso } from "./_utils/time"; async function requireWorkspaceMember(ctx: any, workspaceId: string, userId: string) { const member = await ctx.db .query("workspace_members") .withIndex("by_workspace_user", (q: any) => q.eq("workspace_id", workspaceId).eq("user_id", userId)) .first(); if (!member) { throw new Error("无权限"); } return member; } export const recordCommandLog = mutation({ args: { workspaceId: v.string(), id: v.string(), requestId: v.string(), traceId: v.string(), commandId: v.string(), commandName: v.string(), actorId: v.string(), actorType: v.string(), sourceChannel: v.string(), sourceClient: v.string(), status: v.union(v.literal("pending"), v.literal("succeeded"), v.literal("failed"), v.literal("rolled_back")), targetPageId: v.optional(v.union(v.string(), v.null())), targetBlockId: v.optional(v.union(v.string(), v.null())), payload: v.any(), payloadSummary: v.string(), refs: v.array(v.string()), idempotencyKey: v.optional(v.union(v.string(), v.null())), error: v.optional(v.union(v.string(), v.null())), createdAt: v.string(), finishedAt: v.optional(v.union(v.string(), v.null())), }, handler: async (ctx, args) => { const userId = await requireUserId(ctx); await requireWorkspaceMember(ctx, args.workspaceId, userId); const existing = await ctx.db .query("command_logs") .withIndex("by_command_log_id", (q) => q.eq("id", args.id)) .first(); if (existing) { return { ok: true, id: existing.id, duplicated: true }; } await ctx.db.insert("command_logs", { id: args.id, workspace_id: args.workspaceId, request_id: args.requestId, trace_id: args.traceId, command_id: args.commandId, command_name: args.commandName, actor_id: args.actorId, actor_type: args.actorType, source_channel: args.sourceChannel, source_client: args.sourceClient, status: args.status, target_page_id: args.targetPageId ?? null, target_block_id: args.targetBlockId ?? null, payload: args.payload, payload_summary: args.payloadSummary, refs: args.refs, idempotency_key: args.idempotencyKey ?? null, error: args.error ?? null, created_at: args.createdAt, finished_at: args.finishedAt ?? null, }); return { ok: true, id: args.id, duplicated: false }; }, }); export const recordDomainEvent = mutation({ args: { workspaceId: v.string(), id: v.string(), requestId: v.string(), traceId: v.string(), commandId: v.string(), commandLogId: v.string(), eventType: v.string(), aggregateType: v.string(), aggregateId: v.string(), eventVersion: v.number(), status: v.union(v.literal("pending"), v.literal("committed"), v.literal("rejected"), v.literal("failed")), actorType: v.string(), payload: v.any(), createdAt: v.string(), }, handler: async (ctx, args) => { const userId = await requireUserId(ctx); await requireWorkspaceMember(ctx, args.workspaceId, userId); const existing = await ctx.db .query("domain_events") .withIndex("by_domain_event_id", (q) => q.eq("id", args.id)) .first(); if (existing) { return { ok: true, id: existing.id, duplicated: true }; } await ctx.db.insert("domain_events", { id: args.id, workspace_id: args.workspaceId, request_id: args.requestId, trace_id: args.traceId, command_id: args.commandId, command_log_id: args.commandLogId, event_type: args.eventType, aggregate_type: args.aggregateType, aggregate_id: args.aggregateId, event_version: args.eventVersion, status: args.status, actor_type: args.actorType, payload: args.payload, created_at: args.createdAt, }); return { ok: true, id: args.id, duplicated: false }; }, }); export const listByTrace = query({ args: { workspaceId: v.string(), traceId: v.string() }, handler: async (ctx, args) => { const userId = await requireUserId(ctx); await requireWorkspaceMember(ctx, args.workspaceId, userId); const commandLogs = await ctx.db .query("command_logs") .withIndex("by_workspace_trace", (q) => q.eq("workspace_id", args.workspaceId).eq("trace_id", args.traceId)) .collect(); const domainEvents = await ctx.db .query("domain_events") .withIndex("by_workspace_trace", (q) => q.eq("workspace_id", args.workspaceId).eq("trace_id", args.traceId)) .collect(); return { trace_id: args.traceId, command_logs: commandLogs, domain_events: domainEvents, generated_at: nowIso(), }; }, }); export const listByRequest = query({ args: { workspaceId: v.string(), requestId: v.string() }, handler: async (ctx, args) => { const userId = await requireUserId(ctx); await requireWorkspaceMember(ctx, args.workspaceId, userId); const commandLogs = await ctx.db .query("command_logs") .withIndex("by_workspace_request", (q) => q.eq("workspace_id", args.workspaceId).eq("request_id", args.requestId)) .collect(); const domainEvents = await ctx.db .query("domain_events") .withIndex("by_workspace_request", (q) => q.eq("workspace_id", args.workspaceId).eq("request_id", args.requestId)) .collect(); return { request_id: args.requestId, command_logs: commandLogs, domain_events: domainEvents, generated_at: nowIso(), }; }, });