520 lines
18 KiB
TypeScript
520 lines
18 KiB
TypeScript
import { Router } from "express";
|
|
import type { Request } from "express";
|
|
import {
|
|
issueRecoveryActions,
|
|
issues as issueRows,
|
|
type Db,
|
|
} from "@paperclipai/db";
|
|
import { and, eq, inArray, isNotNull, or, sql } from "drizzle-orm";
|
|
import { conflict } from "../errors.js";
|
|
import {
|
|
EXECUTION_RECONCILIATION_CAUSES,
|
|
createIssueTreeHoldSchema,
|
|
isUuidLike,
|
|
previewIssueTreeControlSchema,
|
|
releaseIssueTreeHoldSchema,
|
|
} from "@paperclipai/shared";
|
|
import { validate } from "../middleware/validate.js";
|
|
import {
|
|
heartbeatService,
|
|
issueService,
|
|
issueTreeControlService,
|
|
logActivity,
|
|
} from "../services/index.js";
|
|
import { assertBoard, getAccessibleResource, getActorInfo } from "./authz.js";
|
|
|
|
const TREE_RUN_CANCELLATION_RESPONSE_WAIT_MS = 1_000;
|
|
const RESUME_EXECUTABLE_STATUSES = ["todo", "in_progress", "in_review"];
|
|
|
|
function errorToMessage(error: unknown) {
|
|
return error instanceof Error ? error.message : String(error);
|
|
}
|
|
|
|
async function waitForRunCancellationTasks(tasks: Promise<void>[]) {
|
|
let timeout: ReturnType<typeof setTimeout> | null = null;
|
|
try {
|
|
await Promise.race([
|
|
Promise.all(tasks),
|
|
new Promise((resolve) => {
|
|
timeout = setTimeout(resolve, TREE_RUN_CANCELLATION_RESPONSE_WAIT_MS);
|
|
}),
|
|
]);
|
|
} finally {
|
|
if (timeout) clearTimeout(timeout);
|
|
}
|
|
}
|
|
|
|
export function issueTreeControlRoutes(db: Db) {
|
|
const router = Router();
|
|
const issuesSvc = issueService(db);
|
|
const treeControlSvc = issueTreeControlService(db);
|
|
const heartbeat = heartbeatService(db);
|
|
|
|
async function resolveRootIssue(req: Request) {
|
|
const rootIssueId = req.params.id as string;
|
|
const root = await issuesSvc.getById(rootIssueId);
|
|
return root;
|
|
}
|
|
|
|
router.post("/issues/:id/tree-control/preview", validate(previewIssueTreeControlSchema), async (req, res) => {
|
|
assertBoard(req);
|
|
const root = await getAccessibleResource(req, res, resolveRootIssue(req), "Root issue not found");
|
|
if (!root) return;
|
|
|
|
const preview = await treeControlSvc.preview(root.companyId, root.id, req.body);
|
|
const actor = getActorInfo(req);
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_control_previewed",
|
|
entityType: "issue",
|
|
entityId: root.id,
|
|
details: {
|
|
mode: preview.mode,
|
|
totals: preview.totals,
|
|
warningCodes: preview.warnings.map((warning) => warning.code),
|
|
},
|
|
});
|
|
|
|
res.json(preview);
|
|
});
|
|
|
|
router.post("/issues/:id/tree-holds", validate(createIssueTreeHoldSchema), async (req, res) => {
|
|
assertBoard(req);
|
|
const root = await getAccessibleResource(req, res, resolveRootIssue(req), "Root issue not found");
|
|
if (!root) return;
|
|
|
|
const actor = getActorInfo(req);
|
|
const actorInput = {
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
userId: actor.actorType === "user" ? actor.actorId : null,
|
|
runId: actor.runId,
|
|
};
|
|
let result = await treeControlSvc.createHold(root.companyId, root.id, {
|
|
...req.body,
|
|
actor: actorInput,
|
|
});
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_hold_created",
|
|
entityType: "issue",
|
|
entityId: root.id,
|
|
details: {
|
|
holdId: result.hold.id,
|
|
mode: result.hold.mode,
|
|
reason: result.hold.reason,
|
|
totals: result.preview.totals,
|
|
warningCodes: result.preview.warnings.map((warning) => warning.code),
|
|
},
|
|
});
|
|
|
|
const runCancellationTasks: Promise<void>[] = [];
|
|
if (result.hold.mode === "pause" || result.hold.mode === "cancel") {
|
|
const interruptedRunIds = [...new Set(result.preview.activeRuns.map((run) => run.id))];
|
|
for (const heartbeatRunId of interruptedRunIds) {
|
|
const cancellationTask = (async () => {
|
|
try {
|
|
await heartbeat.cancelRun(heartbeatRunId);
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_hold_run_interrupted",
|
|
entityType: "heartbeat_run",
|
|
entityId: heartbeatRunId,
|
|
details: {
|
|
holdId: result.hold.id,
|
|
rootIssueId: root.id,
|
|
reason: result.hold.mode === "pause" ? "active_subtree_pause_hold" : "subtree_cancel_operation",
|
|
},
|
|
});
|
|
} catch (error) {
|
|
await Promise.resolve(logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_hold_run_interrupt_failed",
|
|
entityType: "heartbeat_run",
|
|
entityId: heartbeatRunId,
|
|
details: {
|
|
holdId: result.hold.id,
|
|
rootIssueId: root.id,
|
|
reason: result.hold.mode === "pause" ? "active_subtree_pause_hold" : "subtree_cancel_operation",
|
|
error: errorToMessage(error),
|
|
},
|
|
})).catch(() => null);
|
|
}
|
|
})();
|
|
runCancellationTasks.push(cancellationTask);
|
|
}
|
|
|
|
const cancelledWakeups = await treeControlSvc.cancelUnclaimedWakeupsForTree(
|
|
root.companyId,
|
|
root.id,
|
|
result.hold.mode === "pause"
|
|
? "Cancelled because an active subtree pause hold was created"
|
|
: "Cancelled because a subtree cancel operation was applied",
|
|
);
|
|
for (const wakeup of cancelledWakeups) {
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_hold_wakeup_deferred",
|
|
entityType: "agent_wakeup_request",
|
|
entityId: wakeup.id,
|
|
details: {
|
|
holdId: result.hold.id,
|
|
rootIssueId: root.id,
|
|
agentId: wakeup.agentId,
|
|
previousReason: wakeup.reason,
|
|
},
|
|
});
|
|
}
|
|
}
|
|
|
|
if (result.hold.mode === "cancel") {
|
|
const statusUpdate = await treeControlSvc.cancelIssueStatusesForHold(root.companyId, root.id, result.hold.id);
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_cancel_status_updated",
|
|
entityType: "issue",
|
|
entityId: root.id,
|
|
details: {
|
|
holdId: result.hold.id,
|
|
cancelledIssueIds: statusUpdate.updatedIssueIds,
|
|
cancelledIssueCount: statusUpdate.updatedIssueIds.length,
|
|
},
|
|
});
|
|
}
|
|
|
|
if (runCancellationTasks.length > 0) {
|
|
await waitForRunCancellationTasks(runCancellationTasks);
|
|
}
|
|
|
|
if (result.hold.mode === "restore") {
|
|
let statusUpdate;
|
|
try {
|
|
statusUpdate = await treeControlSvc.restoreIssueStatusesForHold(root.companyId, root.id, result.hold.id, {
|
|
reason: result.hold.reason,
|
|
actor: actorInput,
|
|
});
|
|
} catch (error) {
|
|
await treeControlSvc.releaseHold(root.companyId, root.id, result.hold.id, {
|
|
reason: "Restore operation failed before subtree status updates completed",
|
|
metadata: {
|
|
cleanup: "restore_failed_before_apply",
|
|
},
|
|
actor: actorInput,
|
|
}).catch(() => null);
|
|
throw error;
|
|
}
|
|
if (statusUpdate.restoreHold) {
|
|
result = { ...result, hold: statusUpdate.restoreHold };
|
|
}
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_restore_status_updated",
|
|
entityType: "issue",
|
|
entityId: root.id,
|
|
details: {
|
|
holdId: result.hold.id,
|
|
restoredIssueIds: statusUpdate.updatedIssueIds,
|
|
restoredIssueCount: statusUpdate.updatedIssueIds.length,
|
|
releasedCancelHoldIds: statusUpdate.releasedCancelHoldIds,
|
|
},
|
|
});
|
|
|
|
const wakeAgents = typeof req.body.metadata === "object"
|
|
&& req.body.metadata !== null
|
|
&& (req.body.metadata as Record<string, unknown>).wakeAgents === true;
|
|
if (wakeAgents) {
|
|
for (const restoredIssue of statusUpdate.updatedIssues) {
|
|
if (!restoredIssue.assigneeAgentId) continue;
|
|
const wakeRun = await heartbeat
|
|
.wakeup(restoredIssue.assigneeAgentId, {
|
|
source: "assignment",
|
|
triggerDetail: "system",
|
|
reason: "issue_tree_restored",
|
|
payload: {
|
|
issueId: restoredIssue.id,
|
|
rootIssueId: root.id,
|
|
restoreHoldId: result.hold.id,
|
|
},
|
|
requestedByActorType: actor.actorType,
|
|
requestedByActorId: actor.actorId,
|
|
contextSnapshot: {
|
|
issueId: restoredIssue.id,
|
|
taskId: restoredIssue.id,
|
|
wakeReason: "issue_tree_restored",
|
|
source: "issue.tree_restore",
|
|
rootIssueId: root.id,
|
|
restoreHoldId: result.hold.id,
|
|
},
|
|
})
|
|
.catch(() => null);
|
|
if (!wakeRun) continue;
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_restore_wakeup_requested",
|
|
entityType: "heartbeat_run",
|
|
entityId: wakeRun.id,
|
|
issueId: restoredIssue.id,
|
|
details: {
|
|
holdId: result.hold.id,
|
|
rootIssueId: root.id,
|
|
issueId: restoredIssue.id,
|
|
agentId: restoredIssue.assigneeAgentId,
|
|
},
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
res
|
|
.status(result.hold.mode === "restore" || result.hold.mode === "resume" ? 200 : 201)
|
|
.json(result);
|
|
});
|
|
|
|
router.get("/issues/:id/tree-control/state", async (req, res) => {
|
|
assertBoard(req);
|
|
const issueId = req.params.id as string;
|
|
const issue = await getAccessibleResource(req, res, issuesSvc.getById(issueId), "Issue not found");
|
|
if (!issue) return;
|
|
const activePauseHold = await treeControlSvc.getActivePauseHoldGate(issue.companyId, issue.id);
|
|
res.json({ activePauseHold });
|
|
});
|
|
|
|
router.get("/issues/:id/tree-holds", async (req, res) => {
|
|
assertBoard(req);
|
|
const root = await getAccessibleResource(req, res, resolveRootIssue(req), "Root issue not found");
|
|
if (!root) return;
|
|
const statusParam = typeof req.query.status === "string" ? req.query.status : null;
|
|
const modeParam = typeof req.query.mode === "string" ? req.query.mode : null;
|
|
const includeMembers = req.query.includeMembers === "true";
|
|
const holds = await treeControlSvc.listHolds(root.companyId, root.id, {
|
|
status: statusParam === "active" || statusParam === "released" ? statusParam : undefined,
|
|
mode:
|
|
modeParam === "pause" || modeParam === "resume" || modeParam === "cancel" || modeParam === "restore"
|
|
? modeParam
|
|
: undefined,
|
|
includeMembers,
|
|
});
|
|
res.json(holds);
|
|
});
|
|
|
|
router.get("/issues/:id/tree-holds/:holdId", async (req, res) => {
|
|
assertBoard(req);
|
|
const root = await getAccessibleResource(req, res, resolveRootIssue(req), "Root issue not found");
|
|
if (!root) return;
|
|
|
|
const holdId = req.params.holdId as string;
|
|
if (!isUuidLike(holdId)) {
|
|
res.status(400).json({ error: "Invalid hold ID" });
|
|
return;
|
|
}
|
|
|
|
const hold = await treeControlSvc.getHold(root.companyId, holdId);
|
|
if (!hold || hold.rootIssueId !== root.id) {
|
|
res.status(404).json({ error: "Issue tree hold not found" });
|
|
return;
|
|
}
|
|
res.json(hold);
|
|
});
|
|
|
|
router.post(
|
|
"/issues/:id/tree-holds/:holdId/release",
|
|
validate(releaseIssueTreeHoldSchema),
|
|
async (req, res) => {
|
|
assertBoard(req);
|
|
const root = await getAccessibleResource(
|
|
req,
|
|
res,
|
|
resolveRootIssue(req),
|
|
"Root issue not found",
|
|
);
|
|
if (!root) return;
|
|
|
|
const holdId = req.params.holdId as string;
|
|
if (!isUuidLike(holdId)) {
|
|
res.status(400).json({ error: "Invalid hold ID" });
|
|
return;
|
|
}
|
|
|
|
// Releasing a pause does not authorize replay of uncertain provider actions.
|
|
// Check before release so a rejected wake leaves the subtree paused.
|
|
if (req.body.metadata?.wakeAgents === true) {
|
|
const activeHold = await treeControlSvc.getHold(root.companyId, holdId);
|
|
const issueIds =
|
|
activeHold?.mode === "pause" && activeHold.rootIssueId === root.id
|
|
? (activeHold.members ?? [])
|
|
.filter((member) => !member.skipped)
|
|
.map((member) => member.issueId)
|
|
: [];
|
|
if (issueIds.length > 0) {
|
|
const [blocked] = await db
|
|
.select({ identifier: issueRows.identifier })
|
|
.from(issueRecoveryActions)
|
|
.innerJoin(
|
|
issueRows,
|
|
and(
|
|
eq(issueRows.id, issueRecoveryActions.sourceIssueId),
|
|
eq(issueRows.companyId, root.companyId),
|
|
),
|
|
)
|
|
.where(
|
|
and(
|
|
eq(issueRecoveryActions.companyId, root.companyId),
|
|
inArray(issueRecoveryActions.sourceIssueId, issueIds),
|
|
inArray(issueRows.status, RESUME_EXECUTABLE_STATUSES),
|
|
isNotNull(issueRows.assigneeAgentId),
|
|
inArray(issueRecoveryActions.cause, [
|
|
...EXECUTION_RECONCILIATION_CAUSES,
|
|
]),
|
|
or(
|
|
inArray(issueRecoveryActions.status, ["active", "escalated"]),
|
|
sql`${issueRecoveryActions.evidence}->'automaticRecovery'->>'replay' = 'blocked'`,
|
|
),
|
|
),
|
|
)
|
|
.limit(1);
|
|
if (blocked)
|
|
throw conflict(
|
|
`Cannot wake ${blocked.identifier ?? "this task"} until its stopped execution is reconciled. Resume without waking agents, or review the stopped run first.`,
|
|
);
|
|
}
|
|
}
|
|
const actor = getActorInfo(req);
|
|
const hold = await treeControlSvc.releaseHold(
|
|
root.companyId,
|
|
root.id,
|
|
holdId,
|
|
{
|
|
...req.body,
|
|
actor: {
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
userId: actor.actorType === "user" ? actor.actorId : null,
|
|
runId: actor.runId,
|
|
},
|
|
},
|
|
);
|
|
await logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_hold_released",
|
|
entityType: "issue",
|
|
entityId: root.id,
|
|
details: {
|
|
holdId: hold.id,
|
|
mode: hold.mode,
|
|
reason: hold.releaseReason,
|
|
memberCount: hold.members?.length ?? 0,
|
|
},
|
|
});
|
|
|
|
const wakeFailures: Array<{ issueId: string; message: string }> = [];
|
|
if (hold.mode === "pause" && req.body.metadata?.wakeAgents === true) {
|
|
for (const member of hold.members ?? []) {
|
|
if (member.skipped) continue;
|
|
try {
|
|
const issue = await issuesSvc.getById(member.issueId);
|
|
if (
|
|
!issue ||
|
|
issue.companyId !== root.companyId ||
|
|
!issue.assigneeAgentId ||
|
|
!RESUME_EXECUTABLE_STATUSES.includes(issue.status)
|
|
)
|
|
continue;
|
|
await heartbeat.wakeup(issue.assigneeAgentId, {
|
|
source: "assignment",
|
|
triggerDetail: "system",
|
|
reason: "issue_tree_resumed",
|
|
payload: {
|
|
issueId: issue.id,
|
|
rootIssueId: root.id,
|
|
holdId: hold.id,
|
|
},
|
|
requestedByActorType: actor.actorType,
|
|
requestedByActorId: actor.actorId,
|
|
contextSnapshot: {
|
|
issueId: issue.id,
|
|
taskId: issue.id,
|
|
wakeReason: "issue_tree_resumed",
|
|
source: "issue.tree_resume",
|
|
rootIssueId: root.id,
|
|
holdId: hold.id,
|
|
},
|
|
});
|
|
} catch (error) {
|
|
const message = errorToMessage(error);
|
|
wakeFailures.push({ issueId: member.issueId, message });
|
|
await Promise.resolve(
|
|
logActivity(db, {
|
|
companyId: root.companyId,
|
|
actorType: actor.actorType,
|
|
actorId: actor.actorId,
|
|
agentId: actor.agentId,
|
|
runId: actor.runId,
|
|
agentApiKeyId: actor.agentApiKeyId,
|
|
action: "issue.tree_resume_wake_failed",
|
|
entityType: "issue",
|
|
entityId: root.id,
|
|
details: {
|
|
holdId: hold.id,
|
|
issueId: member.issueId,
|
|
error: message,
|
|
},
|
|
}),
|
|
).catch(() => null);
|
|
}
|
|
}
|
|
}
|
|
|
|
res.json({ ...hold, ...(wakeFailures.length ? { wakeFailures } : {}) });
|
|
},
|
|
);
|
|
|
|
return router;
|
|
}
|