205 lines
5.3 KiB
TypeScript
205 lines
5.3 KiB
TypeScript
import { v } from "convex/values";
|
|
|
|
import { mutation, query } from "./_generated/server";
|
|
|
|
const commandStatus = v.union(
|
|
v.literal("queued"),
|
|
v.literal("claimed"),
|
|
v.literal("running"),
|
|
v.literal("succeeded"),
|
|
v.literal("failed"),
|
|
v.literal("cancelled")
|
|
);
|
|
|
|
export const enqueue = mutation({
|
|
args: {
|
|
daemonId: v.string(),
|
|
method: v.string(),
|
|
args: v.any(),
|
|
actorKey: v.array(v.string()),
|
|
},
|
|
handler: async (ctx, args) => {
|
|
const timestamp = Date.now();
|
|
return await ctx.db.insert("daemonCommands", {
|
|
...args,
|
|
status: "queued",
|
|
attempts: 0,
|
|
createdAt: timestamp,
|
|
updatedAt: timestamp,
|
|
});
|
|
},
|
|
});
|
|
|
|
export const list = query({
|
|
args: {
|
|
daemonId: v.string(),
|
|
status: v.optional(commandStatus),
|
|
},
|
|
handler: async (ctx, args) => {
|
|
if (args.status) {
|
|
return await ctx.db
|
|
.query("daemonCommands")
|
|
.withIndex("by_daemonId_and_status", (q) =>
|
|
q.eq("daemonId", args.daemonId).eq("status", args.status!)
|
|
)
|
|
.order("desc")
|
|
.take(100);
|
|
}
|
|
return await ctx.db
|
|
.query("daemonCommands")
|
|
.filter((q) => q.eq(q.field("daemonId"), args.daemonId))
|
|
.order("desc")
|
|
.take(100);
|
|
},
|
|
});
|
|
export const get = query({
|
|
args: { commandId: v.id("daemonCommands") },
|
|
handler: async (ctx, args) =>
|
|
await ctx.db.get("daemonCommands", args.commandId),
|
|
});
|
|
|
|
export const available = query({
|
|
args: { daemonId: v.string() },
|
|
handler: async (ctx, args) => {
|
|
const queued = await ctx.db
|
|
.query("daemonCommands")
|
|
.withIndex("by_daemonId_and_status", (q) =>
|
|
q.eq("daemonId", args.daemonId).eq("status", "queued")
|
|
)
|
|
.order("asc")
|
|
.take(25);
|
|
const expiredClaims = await ctx.db
|
|
.query("daemonCommands")
|
|
.withIndex("by_daemonId_and_status", (q) =>
|
|
q.eq("daemonId", args.daemonId).eq("status", "claimed")
|
|
)
|
|
.filter((q) => q.lt(q.field("leaseExpiresAt"), Date.now()))
|
|
.order("asc")
|
|
.take(25);
|
|
return [...expiredClaims, ...queued].slice(0, 25);
|
|
},
|
|
});
|
|
|
|
export const claim = mutation({
|
|
args: {
|
|
commandId: v.id("daemonCommands"),
|
|
daemonId: v.string(),
|
|
sessionId: v.string(),
|
|
leaseDurationMs: v.number(),
|
|
},
|
|
handler: async (ctx, args) => {
|
|
const command = await ctx.db.get("daemonCommands", args.commandId);
|
|
if (!command || command.daemonId !== args.daemonId) {
|
|
return null;
|
|
}
|
|
const timestamp = Date.now();
|
|
const canClaim =
|
|
command.status === "queued" ||
|
|
(command.status === "claimed" &&
|
|
command.leaseExpiresAt !== undefined &&
|
|
command.leaseExpiresAt <= timestamp);
|
|
if (!canClaim) {
|
|
return null;
|
|
}
|
|
await ctx.db.patch("daemonCommands", command._id, {
|
|
status: "claimed",
|
|
claimedBySessionId: args.sessionId,
|
|
leaseExpiresAt: timestamp + args.leaseDurationMs,
|
|
attempts: command.attempts + 1,
|
|
updatedAt: timestamp,
|
|
});
|
|
return { ...command, status: "claimed" as const };
|
|
},
|
|
});
|
|
|
|
export const start = mutation({
|
|
args: {
|
|
commandId: v.id("daemonCommands"),
|
|
sessionId: v.string(),
|
|
},
|
|
handler: async (ctx, args) => {
|
|
const command = await ctx.db.get("daemonCommands", args.commandId);
|
|
if (
|
|
!command ||
|
|
command.status !== "claimed" ||
|
|
command.claimedBySessionId !== args.sessionId
|
|
) {
|
|
return false;
|
|
}
|
|
const timestamp = Date.now();
|
|
await ctx.db.patch("daemonCommands", command._id, {
|
|
status: "running",
|
|
startedAt: timestamp,
|
|
updatedAt: timestamp,
|
|
});
|
|
return true;
|
|
},
|
|
});
|
|
|
|
export const succeed = mutation({
|
|
args: {
|
|
commandId: v.id("daemonCommands"),
|
|
sessionId: v.string(),
|
|
result: v.any(),
|
|
},
|
|
handler: async (ctx, args) => {
|
|
const command = await ctx.db.get("daemonCommands", args.commandId);
|
|
if (!command || command.claimedBySessionId !== args.sessionId) {
|
|
return false;
|
|
}
|
|
const timestamp = Date.now();
|
|
await ctx.db.patch("daemonCommands", command._id, {
|
|
status: "succeeded",
|
|
result: args.result,
|
|
completedAt: timestamp,
|
|
updatedAt: timestamp,
|
|
leaseExpiresAt: undefined,
|
|
});
|
|
return true;
|
|
},
|
|
});
|
|
|
|
export const fail = mutation({
|
|
args: {
|
|
commandId: v.id("daemonCommands"),
|
|
sessionId: v.string(),
|
|
error: v.string(),
|
|
},
|
|
handler: async (ctx, args) => {
|
|
const command = await ctx.db.get("daemonCommands", args.commandId);
|
|
if (!command || command.claimedBySessionId !== args.sessionId) {
|
|
return false;
|
|
}
|
|
const timestamp = Date.now();
|
|
await ctx.db.patch("daemonCommands", command._id, {
|
|
status: "failed",
|
|
error: args.error,
|
|
completedAt: timestamp,
|
|
updatedAt: timestamp,
|
|
leaseExpiresAt: undefined,
|
|
});
|
|
return true;
|
|
},
|
|
});
|
|
|
|
export const cancel = mutation({
|
|
args: { commandId: v.id("daemonCommands") },
|
|
handler: async (ctx, args) => {
|
|
const command = await ctx.db.get("daemonCommands", args.commandId);
|
|
if (
|
|
!command ||
|
|
["succeeded", "failed", "cancelled"].includes(command.status)
|
|
) {
|
|
return false;
|
|
}
|
|
const timestamp = Date.now();
|
|
await ctx.db.patch("daemonCommands", command._id, {
|
|
status: "cancelled",
|
|
completedAt: timestamp,
|
|
updatedAt: timestamp,
|
|
leaseExpiresAt: undefined,
|
|
});
|
|
return true;
|
|
},
|
|
});
|