import { db } from "../db";
import { workflows } from "../db/schema";
import { getTree as realGetTree, getBlob as realGetBlob } from "../git/repository";
import { parseWorkflow } from "./workflow-parser";
import { enqueueRun as realEnqueueRun } from "./workflow-runner";
const WORKFLOWS_DIR = ".gluecron/workflows";
export interface PushWorkflowSyncResult {
synced: number;
enqueued: number;
errors: string[];
}
export async function syncAndEnqueuePushWorkflows(
opts: {
owner: string;
repo: string;
repositoryId: string;
branch: string;
commitSha: string;
triggeredBy?: string | null;
},
deps: {
getTree: typeof realGetTree;
getBlob: typeof realGetBlob;
enqueueRun: typeof realEnqueueRun;
} = { getTree: realGetTree, getBlob: realGetBlob, enqueueRun: realEnqueueRun }
): Promise<PushWorkflowSyncResult> {
const result: PushWorkflowSyncResult = { synced: 0, enqueued: 0, errors: [] };
const entries = await deps.getTree(opts.owner, opts.repo, opts.commitSha, WORKFLOWS_DIR);
const ymlFiles = entries.filter(
(e) => e.type === "blob" && /\.ya?ml$/i.test(e.name)
);
for (const file of ymlFiles) {
const path = `${WORKFLOWS_DIR}/${file.name}`;
try {
const blob = await deps.getBlob(opts.owner, opts.repo, opts.commitSha, path);
if (!blob || blob.isBinary) continue;
const parseResult = parseWorkflow(blob.content);
if (!parseResult.ok) {
result.errors.push(`${file.name}: ${parseResult.error}`);
continue;
}
const { workflow } = parseResult;
const [row] = await db
.insert(workflows)
.values({
repositoryId: opts.repositoryId,
name: workflow.name || file.name,
path,
yaml: blob.content,
parsed: JSON.stringify(workflow),
onEvents: JSON.stringify(workflow.on),
})
.onConflictDoUpdate({
target: [workflows.repositoryId, workflows.path],
set: {
name: workflow.name || file.name,
yaml: blob.content,
parsed: JSON.stringify(workflow),
onEvents: JSON.stringify(workflow.on),
updatedAt: new Date(),
},
})
.returning({ id: workflows.id, disabled: workflows.disabled });
if (!row) continue;
result.synced += 1;
if (!row.disabled && workflow.on.includes("push")) {
try {
await deps.enqueueRun({
workflowId: row.id,
repositoryId: opts.repositoryId,
event: "push",
ref: opts.branch,
commitSha: opts.commitSha,
triggeredBy: opts.triggeredBy ?? null,
});
result.enqueued += 1;
} catch (err) {
result.errors.push(
`enqueue ${file.name}: ${err instanceof Error ? err.message : String(err)}`
);
}
}
} catch (err) {
result.errors.push(
`${file.name}: ${err instanceof Error ? err.message : String(err)}`
);
}
}
return result;
}
|