Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 16 additions & 7 deletions src/jobs/garmin-dump-flow.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -216,12 +216,21 @@ describe("attachGarminFitImportFlow", () => {
{ id: "import-job-1", queue: "bull:import" },
);

const flow: FlowJob | undefined = mockAdd.mock.calls[0]?.[0];
const fitJobIds = flow?.children?.map((child) => child.opts?.jobId);
const extractionJobIds = flow?.children?.map((child) => child.children?.[0]?.opts?.jobId);
expect(fitJobIds).toHaveLength(1_294);
expect(extractionJobIds).toHaveLength(1_294);
expect(new Set(fitJobIds).size).toBe(1_294);
expect(new Set(extractionJobIds).size).toBe(1_294);
const flows: FlowJob[] = mockAdd.mock.calls.map(([flow]) => flow);
const allFitJobIds = flows.flatMap(
(flow) => flow.children?.map((child) => child.opts?.jobId) ?? [],
);
const allExtractionJobIds = flows.flatMap(
(flow) =>
flow.children?.flatMap((child) => child.children?.map((c) => c.opts?.jobId) ?? []) ?? [],
);
expect(allFitJobIds).toHaveLength(1_294);
expect(allExtractionJobIds).toHaveLength(1_294);
expect(new Set(allFitJobIds).size).toBe(1_294);
expect(new Set(allExtractionJobIds).size).toBe(1_294);

expect(flows).toHaveLength(Math.ceil(1_294 / 500));
const batchIds = flows.map((flow) => flow.opts?.jobId);
expect(new Set(batchIds).size).toBe(flows.length);
});
});
32 changes: 26 additions & 6 deletions src/jobs/garmin-dump-flow.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ const FIT_FILE_IMPORT_JOB_NAME = "fit-file-import";
const FIT_FILE_IMPORT_BATCH_JOB_NAME = "fit-file-import-batch";
const ZIP_ENTRY_EXTRACT_JOB_NAME = "zip-entry-extract";

export const FLOW_BATCH_SIZE = 500;

export interface GarminImportParent {
id: string;
queue: string;
Expand Down Expand Up @@ -71,7 +73,9 @@ function fitFileImportFlowChild(
};
}

async function addGarminFitImportFlow(
async function addBatchFlow(
batchId: string,
entries: GarminFitJobEntry[],
preparedImport: PreparedGarminDumpImport,
parent?: GarminImportParent,
) {
Expand All @@ -80,21 +84,37 @@ async function addGarminFitImportFlow(
queueName: FIT_FILE_IMPORT_BATCH_QUEUE,
data: { type: "fit-file-import-batch" },
opts: {
jobId: preparedImport.batchId,
jobId: batchId,
...(parent ? { parent } : {}),
ignoreDependencyOnFailure: true,
removeOnComplete: { age: 86_400 },
removeOnFail: { age: 604_800 },
},
children: preparedImport.fitJobEntries.map((fitJobEntry) =>
fitFileImportFlowChild(preparedImport, fitJobEntry),
),
children: entries.map((fitJobEntry) => fitFileImportFlowChild(preparedImport, fitJobEntry)),
});
}

export function createBatchId(
preparedImport: PreparedGarminDumpImport,
batchIndex: number,
): string {
if (batchIndex === 0) {
return preparedImport.batchId;
}
return `${preparedImport.batchId}-batch-${batchIndex}`;
}

export async function attachGarminFitImportFlow(
preparedImport: PreparedGarminDumpImport,
parent: GarminImportParent,
): Promise<void> {
await addGarminFitImportFlow(preparedImport, parent);
const allEntries = preparedImport.fitJobEntries;
const batchCount = Math.ceil(allEntries.length / FLOW_BATCH_SIZE);

for (let batchIndex = 0; batchIndex < batchCount; batchIndex++) {
const start = batchIndex * FLOW_BATCH_SIZE;
const batchEntries = allEntries.slice(start, start + FLOW_BATCH_SIZE);
const batchId = createBatchId(preparedImport, batchIndex);
await addBatchFlow(batchId, batchEntries, preparedImport, parent);
}
}
44 changes: 44 additions & 0 deletions src/jobs/process-garmin-dump-import-job.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@ vi.mock("../providers/garmin-dump.ts", () => ({

vi.mock("./garmin-dump-flow.ts", () => ({
attachGarminFitImportFlow: vi.fn(),
createBatchId: vi.fn((preparedImport: { batchId: string }, batchIndex: number) =>
batchIndex === 0 ? preparedImport.batchId : `${preparedImport.batchId}-batch-${batchIndex}`,
),
FLOW_BATCH_SIZE: 500,
}));

const mockDb: SyncDatabase = {
Expand Down Expand Up @@ -340,6 +344,46 @@ describe("processGarminDumpImportJob", () => {
]);
});

it("finalizes zero FIT files without throwing", async () => {
const job = createJob({
...waitingCheckpoint(),
totalFitFiles: 0,
batchIds: [],
});
job.moveToWaitingChildren.mockResolvedValue(false);

const result = await processGarminDumpImportJob(job, mockDb);

expect(result.recordsSynced).toBe(2);
expect(result.errors).toEqual([{ message: "one summary was invalid" }]);
expect(result.totalFitFiles).toBe(0);
});

it("merges partial batch failures into successful batch results", async () => {
const job = createJob({
...waitingCheckpoint(),
batchIds: ["garmin-fit-batch-1", "garmin-fit-batch-1-batch-1"],
});
job.moveToWaitingChildren.mockResolvedValue(false);
job.getChildrenValues.mockResolvedValue({
"bull:fit-file-import-batch:garmin-fit-batch-1": {
recordsSynced: 3,
errors: [],
},
});
job.getIgnoredChildrenFailures.mockResolvedValue({
"bull:fit-file-import-batch:garmin-fit-batch-1-batch-1": "batch Redis connection failed",
});

const result = await processGarminDumpImportJob(job, mockDb);

expect(result.recordsSynced).toBe(5);
expect(result.errors).toEqual([
{ message: "one summary was invalid" },
{ message: "Garmin FIT batch failed: batch Redis connection failed" },
]);
});

it("bounds combined preparation and FIT error groups", async () => {
const baseErrors = Array.from({ length: 7 }, (_, errorIndex) => ({
message: `summary error ${errorIndex + 1}`,
Expand Down
49 changes: 32 additions & 17 deletions src/jobs/process-garmin-dump-import-job.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ import {
prepareGarminDumpImport,
} from "../providers/garmin-dump.ts";
import type { SyncResult } from "../providers/types.ts";
import { attachGarminFitImportFlow } from "./garmin-dump-flow.ts";
import { attachGarminFitImportFlow, createBatchId, FLOW_BATCH_SIZE } from "./garmin-dump-flow.ts";
import type { ImportJobData } from "./queues.ts";

const importResultSchema = z.object({
Expand All @@ -23,6 +23,7 @@ export const garminImportCheckpointSchema = z.object({
phase: z.enum(["prepared", "waiting-children"]),
startedAtMs: z.number().int().nonnegative(),
batchId: z.string().min(1),
batchIds: z.array(z.string().min(1)).optional(),
baseResult: importResultSchema,
tempDirectories: z.array(z.string()),
totalFitFiles: z.number().int().nonnegative(),
Expand Down Expand Up @@ -70,20 +71,30 @@ function batchResultFromChildren(
checkpoint: GarminImportCheckpoint,
childValues: Record<string, unknown>,
): z.infer<typeof importResultSchema> | null {
const batchEntry = Object.entries(childValues).find(([jobKey]) =>
jobKey.endsWith(`:${checkpoint.batchId}`),
const batchIds = checkpoint.batchIds ?? [checkpoint.batchId];
const batchEntries = Object.entries(childValues).filter(([jobKey]) =>
batchIds.some((id) => jobKey.endsWith(`:${id}`)),
);
return batchEntry ? importResultSchema.parse(batchEntry[1]) : null;
if (batchEntries.length === 0) return null;

let totalRecordsSynced = 0;
const allErrors: Array<{ message: string }> = [];
for (const [, value] of batchEntries) {
const parsed = importResultSchema.parse(value);
totalRecordsSynced += parsed.recordsSynced;
allErrors.push(...parsed.errors);
}
return { recordsSynced: totalRecordsSynced, errors: allErrors };
}

function batchFailureFromChildren(
checkpoint: GarminImportCheckpoint,
ignoredFailures: Record<string, string>,
): string | null {
const batchEntry = Object.entries(ignoredFailures).find(([jobKey]) =>
jobKey.endsWith(`:${checkpoint.batchId}`),
);
return batchEntry?.[1] ?? null;
): string[] {
const batchIds = checkpoint.batchIds ?? [checkpoint.batchId];
return Object.entries(ignoredFailures)
.filter(([jobKey]) => batchIds.some((id) => jobKey.endsWith(`:${id}`)))
.map(([, message]) => message);
}

function uniqueDirectories(...directoryGroups: readonly string[][]): string[] {
Expand Down Expand Up @@ -149,7 +160,9 @@ export async function processGarminDumpImportJob(
{ ...preparedImport, tempDirectories },
{ id: job.id, queue: job.queueQualifiedName },
);
checkpoint = { ...checkpoint, phase: "waiting-children" };
const batchCount = Math.ceil(preparedImport.totalFitFiles / FLOW_BATCH_SIZE);
const batchIds = Array.from({ length: batchCount }, (_, i) => createBatchId(preparedImport, i));
checkpoint = { ...checkpoint, phase: "waiting-children", batchIds };
await persistCheckpoint(job, checkpoint);
await logGarminImportPhase(
job,
Expand All @@ -158,9 +171,11 @@ export async function processGarminDumpImportJob(
}

if (await job.moveToWaitingChildren()) {
const batchCount = checkpoint.batchIds?.length ?? 1;
const batchLabel = batchCount > 1 ? ` across ${batchCount} batches` : "";
await logGarminImportPhase(
job,
`[phase] Waiting for ${checkpoint.totalFitFiles} Garmin FIT ${checkpoint.totalFitFiles === 1 ? "file" : "files"} in batch ${checkpoint.batchId}`,
`[phase] Waiting for ${checkpoint.totalFitFiles} Garmin FIT ${checkpoint.totalFitFiles === 1 ? "file" : "files"} in batch ${checkpoint.batchId}${batchLabel}`,
);
throw new WaitingChildrenError();
}
Expand All @@ -170,15 +185,15 @@ export async function processGarminDumpImportJob(
job.getIgnoredChildrenFailures(),
]);
const batchResult = batchResultFromChildren(checkpoint, childValues);
Comment thread
Asherlc marked this conversation as resolved.
const batchFailure = batchFailureFromChildren(checkpoint, ignoredFailures);
if (!batchResult && !batchFailure) {
const batchFailures = batchFailureFromChildren(checkpoint, ignoredFailures);
if (!batchResult && batchFailures.length === 0 && checkpoint.totalFitFiles > 0) {
throw new Error(`Garmin FIT batch result is missing: ${checkpoint.batchId}`);
}

const fitResult = batchResult ?? {
recordsSynced: 0,
errors: [{ message: `Garmin FIT batch failed: ${batchFailure}` }],
};
const fitResult = batchResult ?? { recordsSynced: 0, errors: [] };
for (const message of batchFailures) {
fitResult.errors.push({ message: `Garmin FIT batch failed: ${message}` });
}
await cleanupPreparedGarminDumpImport(checkpoint.tempDirectories);
const result: GarminDumpImportResult = {
provider: "garmin-dump",
Expand Down
2 changes: 1 addition & 1 deletion src/jobs/queues.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ export const fitFileImportActivitySummarySchema = z.object({
startedAtIso: z.string(),
endedAtIso: z.string(),
name: z.string(),
raw: z.unknown(),
raw: z.unknown().optional(),
});

export const fitFileImportJobDataSchema = z.object({
Expand Down
3 changes: 1 addition & 2 deletions src/providers/garmin-dump.ts
Original file line number Diff line number Diff line change
Expand Up @@ -539,7 +539,6 @@ function garminSummaryToFitJobSummary(
startedAtIso: startedAt.toISOString(),
endedAtIso: endedAt.toISOString(),
name: summary.name ?? `Garmin ${activityType.replace(/_/g, " ")}`,
raw: summary,
};
}

Expand Down Expand Up @@ -683,7 +682,7 @@ export async function prepareGarminDumpImport(
errors: errors.map((error) => ({ message: error.message })),
},
fitJobEntries,
tempDirectories: parsedDump.tempDirectories,
tempDirectories: [...parsedDump.tempDirectories],
totalFitFiles: fitJobEntries.length,
};
preparationComplete = true;
Expand Down