26 lines
1.8 KiB
JavaScript
26 lines
1.8 KiB
JavaScript
// Maintenance only: stop caller-analytics worker first; does not touch calls or billing.
|
|
import { PrismaClient } from '@prisma/client';
|
|
import { analyticsHash, COUNT_KEYS } from '../packages/database/dist/index.js';
|
|
import assert from 'node:assert/strict';
|
|
assert(process.argv.includes('--worker-stopped'), 'Stop the analytics worker before rebuilding derived buckets');
|
|
const db = new PrismaClient();
|
|
try {
|
|
await db.$transaction(async tx => {
|
|
const rows = await tx.$queryRaw`SELECT * FROM caller_analysis_states`;
|
|
assert(rows.length <= 100000, 'Use a batched maintenance job for larger data sets');
|
|
const buckets = new Map();
|
|
for (const row of rows) {
|
|
const minute = new Date(Math.floor(+row.started_at / 60000) * 60000);
|
|
const id = analyticsHash([minute.toISOString(), row.view, row.customer_id, row.customer_gateway_id, row.vendor_id, row.vendor_gateway_id, row.caller, row.city, row.carrier]);
|
|
const bucket = buckets.get(id) || {...row, id, started_at:minute, counts:Object.fromEntries(COUNT_KEYS.map(k=>[k,0]))};
|
|
for (const key of COUNT_KEYS) bucket.counts[key] += Number(row.counts[key]);
|
|
buckets.set(id,bucket);
|
|
}
|
|
await tx.$executeRaw`DELETE FROM caller_analysis_minute_buckets`;
|
|
for (const r of buckets.values()) await tx.$executeRaw`INSERT INTO caller_analysis_minute_buckets
|
|
(id,started_at,view,customer_id,customer_gateway_id,vendor_id,vendor_gateway_id,caller,city,carrier,counts)
|
|
VALUES (${r.id},${r.started_at},${r.view},${r.customer_id},${r.customer_gateway_id},${r.vendor_id},${r.vendor_gateway_id},${r.caller},${r.city},${r.carrier},${JSON.stringify(r.counts)})`;
|
|
console.log(JSON.stringify({states:rows.length,buckets:buckets.size,rebuilt:true}));
|
|
},{timeout:60000});
|
|
} finally {await db.$disconnect();}
|