import { inject, injectable } from "@codemation/core"; import type { HumanTaskActor, HumanTaskStore, JsonValue } from "@codemation/core"; import { CodemationTelemetryAttributeNames, HumanTaskStoreToken } from "@codemation/core"; import { Engine } from "@codemation/core/bootstrap"; import { ApplicationRequestError } from "../ApplicationRequestError"; import { HitlResumeTokenSigner } from "../../hitl/HitlResumeTokenSigner"; import { HitlTimeoutJobScheduler } from "../../hitl/HitlTimeoutJobScheduler"; import { DecisionSchemaValidator } from "./DecisionSchemaValidator"; import { ResumeTelemetryContextForRun } from "../telemetry/ResumeTelemetryContextForRun"; export interface DecideHumanTaskArgs { taskId: string; decision: JsonValue; decidedBy: HumanTaskActor; } export interface DecideHumanTaskResult { status: "decided"; runStatus: "running" | "halted"; } @injectable() export class DecideHumanTaskCommandHandler { private readonly taskStore: HumanTaskStore | undefined; constructor( @inject(HumanTaskStoreToken) taskStore: HumanTaskStore | undefined, @inject(Engine) private readonly engine: Engine, @inject(HitlResumeTokenSigner) private readonly tokenSigner: HitlResumeTokenSigner, @inject(HitlTimeoutJobScheduler) private readonly timeoutScheduler: HitlTimeoutJobScheduler, @inject(DecisionSchemaValidator) private readonly schemaValidator: DecisionSchemaValidator, @inject(ResumeTelemetryContextForRun) private readonly resumeTelemetry: ResumeTelemetryContextForRun, ) { this.taskStore = taskStore; } async decide(args: DecideHumanTaskArgs): Promise { if (!this.taskStore) { throw new ApplicationRequestError(503, "HITL is not available in this configuration"); } const task = await this.taskStore.findById(args.taskId); if (!task) { throw new ApplicationRequestError(404, "HumanTask not found"); } if (task.status !== "pending") { throw new ApplicationRequestError(409, `HumanTask is not pending (current status: ${task.status})`); } const validationResult = this.schemaValidator.validate({ schemaJson: task.decisionSchemaJson, value: args.decision, }); if (!validationResult.valid) { throw new ApplicationRequestError( 422, `Decision does not match the expected schema: ${validationResult.message}`, ); } const decidedAt = new Date(); await this.taskStore.markDecided({ taskId: args.taskId, decision: args.decision, decidedBy: args.decidedBy, decidedAt, }); await this.timeoutScheduler.cancelTimeoutJob(args.taskId); const telemetry = await this.resumeTelemetry.forTask(args.taskId); const latencyMs = decidedAt.getTime() - task.createdAt.getTime(); const decisionPayload = args.decision as Record | null; const decisionStatus = typeof decisionPayload?.["approved"] === "boolean" ? decisionPayload["approved"] ? "approved" : "rejected" : "decided"; await telemetry?.addSpanEvent({ name: "hitl.task.decided", attributes: { [CodemationTelemetryAttributeNames.hitlTaskId]: args.taskId, [CodemationTelemetryAttributeNames.hitlDecisionStatus]: decisionStatus, actor: args.decidedBy.actorId, latencyMs, }, }); const resumeResult = await this.engine.resumeRun({ runId: task.runId, taskId: task.id, resumeContext: { decision: { kind: "decided", value: args.decision, actor: args.decidedBy, decidedAt, }, delivery: task.deliveryRef ?? null, task: { taskId: task.id, runId: task.runId, nodeId: task.nodeId, expiresAt: task.expiresAt, resumeUrl: "", }, }, }); return { status: "decided", runStatus: resumeResult.status === "failed" || resumeResult.status === "halted" ? "halted" : "running", }; } async validateResumeToken(args: { taskId: string; token: string }): Promise<{ schemaHash: string }> { const result = this.tokenSigner.verify(args.token); if (!result.ok) { if (result.reason === "expired") { throw new ApplicationRequestError(410, "Resume token has expired"); } throw new ApplicationRequestError(401, "Invalid resume token"); } if (result.taskId !== args.taskId) { throw new ApplicationRequestError(401, "Token taskId does not match"); } if (!this.taskStore) { throw new ApplicationRequestError(503, "HITL is not available in this configuration"); } const task = await this.taskStore.findById(args.taskId); if (!task) { throw new ApplicationRequestError(404, "HumanTask not found"); } if (task.decisionSchemaHash.slice(0, 8) !== result.schemaHash) { throw new ApplicationRequestError(410, "Schema has changed since this token was issued"); } return { schemaHash: result.schemaHash }; } }