fix: bound formatter memory and improve operations workflows

This commit is contained in:
hectorzhao
2026-09-08 12:46:15 +08:00
parent 50ae37242b
commit 2c228a94e1
27 changed files with 4013 additions and 1230 deletions
File diff suppressed because it is too large Load Diff
+158 -123
View File
@@ -1,16 +1,14 @@
import { BadRequestException, NotFoundException } from '@nestjs/common';
import { Prisma } from '@prisma/client';
import { randomUUID } from 'node:crypto';
import { moneyToNumber } from '../../common/money';
import { PrismaService } from '../../prisma/prisma.service';
import type { MessageQuery, TraceQuery, OperationLogQuery, GatewaySubmitDeadLetterQuery, DownstreamDeliveryQuery, DownstreamDeliveryDashboardQuery, DownstreamRecoveryStatusQuery, MessageSegmentAuditQuery, SignatureQualityQuery } from '../operations.contracts';
import { messageWhere, recognizedCarrierValues, carrierWhere, startOfShanghaiDay, endOfShanghaiDay, qualityBusinessDay, shanghaiDateKey, normalizeGroupBy, returnedTransactionWhere, createdAtRange, downstreamAlertPendingMinutes, downstreamAlertRecentFailedHours, downstreamAlertWindows, downstreamAlertWhere, stalledPendingWhere, downstreamDeliveryScopedWhere, parseDateBoundary, downstreamRecoveryStatusWhere, escapeCsvCell, formatCsvDate, formatExportTimestamp, clientApplicationView, clientReceiptView, clientMessageView, clientBatchTaskView, clientUplinkView, clientAccountView, clientRechargeView, summarizeMessageGroups, groupDownstreamByType, groupDownstreamByApplication, positiveInteger, operationLogLevelWhere, normalizeOperationLog, sanitizeGatewaySubmitException, redactGatewayCommandValue } from '../operations.helpers';
import type { SignatureQualityQuery } from '../operations.contracts';
import { messageWhere, qualityBusinessDay, normalizeGroupBy, positiveInteger } from '../operations.helpers';
// R2 quality query domain. Method bodies are preserved byte-for-byte from the facade baseline.
export class OperationsQualityQueries {
constructor(private readonly prisma: PrismaService) {}
async statistics(query: { tenantId?: string; groupBy?: string }) {
async statistics(query: { tenantId?: string; groupBy?: string }) {
const groupBy = normalizeGroupBy(query.groupBy);
if (groupBy === 'tenantId') {
return this.prisma.smsMessageRecord.groupBy({
@@ -35,24 +33,26 @@ async statistics(query: { tenantId?: string; groupBy?: string }) {
_sum: { amountCents: true, billingUnits: true },
});
}
async sendQuality(date?: string) {
async sendQuality(date?: string) {
const day = qualityBusinessDay(date);
const [channels, signatureSplits, summaryRows, applications] = await Promise.all([
this.prisma.$queryRaw<Array<{
channelId: string;
channelName: string;
total: number;
acceptedCount: number;
submitFailureCount: number;
submitFailureRate: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
unknownRate: number;
failureRate: number;
averageArrivalMs: number | null;
}>>(Prisma.sql`
this.prisma.$queryRaw<
Array<{
channelId: string;
channelName: string;
total: number;
acceptedCount: number;
submitFailureCount: number;
submitFailureRate: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
unknownRate: number;
failureRate: number;
averageArrivalMs: number | null;
}>
>(Prisma.sql`
WITH base AS (
SELECT
submit."channelId" AS channel_id,
@@ -131,22 +131,24 @@ async sendQuality(date?: string) {
GROUP BY channel_id
ORDER BY COUNT(*) DESC, channel_id
`),
this.prisma.$queryRaw<Array<{
id: string;
signatureId: string;
signatureName: string;
tenantId: string;
tenantName: string;
hasDrainage: boolean;
total: number;
acceptedCount: number;
submitFailureCount: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
averageArrivalMs: number | null;
}>>(Prisma.sql`
this.prisma.$queryRaw<
Array<{
id: string;
signatureId: string;
signatureName: string;
tenantId: string;
tenantName: string;
hasDrainage: boolean;
total: number;
acceptedCount: number;
submitFailureCount: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
averageArrivalMs: number | null;
}>
>(Prisma.sql`
WITH base AS (
SELECT
message."signatureId" AS signature_id,
@@ -216,13 +218,15 @@ async sendQuality(date?: string) {
GROUP BY signature.id, signature.name, tenant.id, tenant.name, base.has_drainage
ORDER BY "successCount" DESC, total DESC, signature.name
`),
this.prisma.$queryRaw<Array<{
total: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
}>>(Prisma.sql`
this.prisma.$queryRaw<
Array<{
total: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
}>
>(Prisma.sql`
WITH base AS (
SELECT message.status, message."receiptStatus" AS receipt_status
FROM "SmsMessageRecord" message
@@ -251,17 +255,19 @@ async sendQuality(date?: string) {
END AS "successRate"
FROM base
`),
this.prisma.$queryRaw<Array<{
applicationId: string;
applicationName: string;
tenantId: string;
tenantName: string;
total: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
}>>(Prisma.sql`
this.prisma.$queryRaw<
Array<{
applicationId: string;
applicationName: string;
tenantId: string;
tenantName: string;
total: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
}>
>(Prisma.sql`
WITH base AS (
SELECT
message."applicationId" AS application_id,
@@ -314,28 +320,30 @@ async sendQuality(date?: string) {
};
return { date: day.key, summary, channels, signatures, drainageSignatures, applications };
}
async signatureQuality(query: SignatureQualityQuery) {
async signatureQuality(query: SignatureQualityQuery) {
const day = qualityBusinessDay(query.date);
const page = positiveInteger(query.page, 1);
const pageSize = Math.min(50, positiveInteger(query.pageSize, 10));
const pageSize = Math.min(100, positiveInteger(query.pageSize, 25));
const keyword = query.keyword?.trim() || null;
const keywordPattern = keyword ? `%${keyword}%` : null;
const summaries = await this.prisma.$queryRaw<Array<{
signatureId: string;
signatureName: string;
tenantId: string;
tenantName: string;
applicationNames: string | null;
total: number;
acceptedCount: number;
submitFailureCount: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
averageArrivalMs: number | null;
rowCount: number;
}>>(Prisma.sql`
const summaries = await this.prisma.$queryRaw<
Array<{
signatureId: string;
signatureName: string;
tenantId: string;
tenantName: string;
applicationNames: string | null;
total: number;
acceptedCount: number;
submitFailureCount: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
averageArrivalMs: number | null;
rowCount: number;
}>
>(Prisma.sql`
WITH base AS (
SELECT
message."signatureId" AS signature_id,
@@ -415,23 +423,26 @@ async signatureQuality(query: SignatureQualityQuery) {
OFFSET ${(page - 1) * pageSize}
`);
const signatureIds = summaries.map((item) => item.signatureId);
const drainageBreakdowns = signatureIds.length === 0
? []
: await this.prisma.$queryRaw<Array<{
signatureId: string;
channelId: string;
channelName: string;
carrier: string;
drainageState: 'with' | 'without' | 'unknown';
total: number;
acceptedCount: number;
submitFailureCount: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
averageArrivalMs: number | null;
}>>(Prisma.sql`
const drainageBreakdowns =
signatureIds.length === 0
? []
: await this.prisma.$queryRaw<
Array<{
signatureId: string;
channelId: string;
channelName: string;
carrier: string;
drainageState: 'with' | 'without' | 'unknown';
total: number;
acceptedCount: number;
submitFailureCount: number;
successCount: number;
unknownCount: number;
failureCount: number;
successRate: number;
averageArrivalMs: number | null;
}>
>(Prisma.sql`
WITH base AS (
SELECT
message."signatureId" AS signature_id,
@@ -526,16 +537,19 @@ async signatureQuality(query: SignatureQualityQuery) {
GROUP BY signature_id, channel_id, carrier, drainage_state
ORDER BY signature_id, COUNT(*) DESC, channel_id, carrier, drainage_state
`);
const carrierOverview = signatureIds.length === 0
? []
: await this.prisma.$queryRaw<Array<{
signatureId: string;
carrier: string;
businessMessageCount: number;
finalSuccessCount: number;
finalSuccessRate: number;
averageArrivalMs: number | null;
}>>(Prisma.sql`
const carrierOverview =
signatureIds.length === 0
? []
: await this.prisma.$queryRaw<
Array<{
signatureId: string;
carrier: string;
businessMessageCount: number;
finalSuccessCount: number;
finalSuccessRate: number;
averageArrivalMs: number | null;
}>
>(Prisma.sql`
SELECT
message."signatureId" AS "signatureId",
COALESCE(NULLIF(message.carrier, ''), 'unknown') AS carrier,
@@ -570,6 +584,7 @@ async signatureQuality(query: SignatureQualityQuery) {
ORDER BY message."signatureId", COUNT(*) DESC, carrier
`);
const items = summaries.map(({ rowCount: _rowCount, ...summary }) => {
void _rowCount;
const signatureDrainageBreakdowns = drainageBreakdowns.filter((item) => item.signatureId === summary.signatureId);
const signatureBreakdowns = aggregateChannelCarrierRows(signatureDrainageBreakdowns);
return {
@@ -610,26 +625,41 @@ type SignatureSplitRow = {
function aggregateSignatureRows(rows: SignatureSplitRow[]) {
const grouped = new Map<string, SignatureSplitRow[]>();
rows.forEach((row) => grouped.set(row.signatureId, [...(grouped.get(row.signatureId) ?? []), row]));
return [...grouped.values()].map((parts) => {
const first = parts[0];
const total = parts.reduce((sum, item) => sum + item.total, 0);
const acceptedCount = parts.reduce((sum, item) => sum + item.acceptedCount, 0);
const successCount = parts.reduce((sum, item) => sum + item.successCount, 0);
const arrivalWeight = parts.reduce((sum, item) => sum + (item.averageArrivalMs == null ? 0 : item.successCount), 0);
return {
...first,
id: first.signatureId,
hasDrainage: false,
total,
acceptedCount,
submitFailureCount: parts.reduce((sum, item) => sum + item.submitFailureCount, 0),
successCount,
unknownCount: parts.reduce((sum, item) => sum + item.unknownCount, 0),
failureCount: parts.reduce((sum, item) => sum + item.failureCount, 0),
successRate: acceptedCount === 0 ? 0 : Math.round(successCount * 1000 / acceptedCount) / 10,
averageArrivalMs: arrivalWeight === 0 ? null : Math.round(parts.reduce((sum, item) => sum + (item.averageArrivalMs ?? 0) * item.successCount, 0) / arrivalWeight),
};
}).sort((left, right) => right.successCount - left.successCount || right.total - left.total || left.signatureName.localeCompare(right.signatureName));
return [...grouped.values()]
.map((parts) => {
const first = parts[0];
const total = parts.reduce((sum, item) => sum + item.total, 0);
const acceptedCount = parts.reduce((sum, item) => sum + item.acceptedCount, 0);
const successCount = parts.reduce((sum, item) => sum + item.successCount, 0);
const arrivalWeight = parts.reduce(
(sum, item) => sum + (item.averageArrivalMs == null ? 0 : item.successCount),
0,
);
return {
...first,
id: first.signatureId,
hasDrainage: false,
total,
acceptedCount,
submitFailureCount: parts.reduce((sum, item) => sum + item.submitFailureCount, 0),
successCount,
unknownCount: parts.reduce((sum, item) => sum + item.unknownCount, 0),
failureCount: parts.reduce((sum, item) => sum + item.failureCount, 0),
successRate: acceptedCount === 0 ? 0 : Math.round((successCount * 1000) / acceptedCount) / 10,
averageArrivalMs:
arrivalWeight === 0
? null
: Math.round(
parts.reduce((sum, item) => sum + (item.averageArrivalMs ?? 0) * item.successCount, 0) / arrivalWeight,
),
};
})
.sort(
(left, right) =>
right.successCount - left.successCount ||
right.total - left.total ||
left.signatureName.localeCompare(right.signatureName),
);
}
type DrainageBreakdownRow = {
@@ -670,8 +700,13 @@ function aggregateChannelCarrierRows(rows: DrainageBreakdownRow[]) {
successCount,
unknownCount: parts.reduce((sum, item) => sum + item.unknownCount, 0),
failureCount: parts.reduce((sum, item) => sum + item.failureCount, 0),
successRate: acceptedCount === 0 ? 0 : Math.round(successCount * 1000 / acceptedCount) / 10,
averageArrivalMs: arrivalWeight === 0 ? null : Math.round(parts.reduce((sum, item) => sum + (item.averageArrivalMs ?? 0) * item.successCount, 0) / arrivalWeight),
successRate: acceptedCount === 0 ? 0 : Math.round((successCount * 1000) / acceptedCount) / 10,
averageArrivalMs:
arrivalWeight === 0
? null
: Math.round(
parts.reduce((sum, item) => sum + (item.averageArrivalMs ?? 0) * item.successCount, 0) / arrivalWeight,
),
};
});
}