import { Transform } from 'node:stream'; import type { FastifyInstance, FastifyReply, FastifyRequest, preParsingHookHandler } from 'fastify'; import { z } from 'zod'; import type { AuthRepository } from '../auth/auth-repository.js'; import { authenticateAccessToken } from '../auth/authenticate.js'; import type { AccessProfile } from '../auth/rbac-repository.js'; import { MediaStorage, MediaValidationError } from '../content/media-storage.js'; import { CleaningTaskError, cleaningTaskStatuses, type CleaningActor, type CleaningTaskRepository } from '../cleaning/cleaning-task-repository.js'; import { CleaningPayoutError, type CleaningPayoutService } from '../cleaning/cleaning-payout-service.js'; import { WechatPayError, type WechatNotificationHeaders } from '../payments/wechat-pay-client.js'; const listSchema = z.object({ page: z.coerce.number().int().min(1).default(1), pageSize: z.coerce.number().int().min(1).max(50).default(20), status: z.enum(cleaningTaskStatuses).optional() }); const paramsSchema = z.object({ taskId: z.string().regex(/^[1-9]\d{0,19}$/) }); const submitSchema = z.object({ photoUrls: z.array(z.string().url()).max(9).default([]), note: z.string().trim().max(512).optional() }); const assignSchema = z.object({ cleanerUserId: z.string().regex(/^[1-9]\d{0,19}$/), note: z.string().trim().max(512).optional() }).strict(); const memberParamsSchema = z.object({ taskId: z.string().regex(/^[1-9]\d{0,19}$/), cleanerUserId: z.string().regex(/^[1-9]\d{0,19}$/) }); const memberSchema = z.object({ cleanerUserId: z.string().regex(/^[1-9]\d{0,19}$/), rewardCents: z.coerce.number().int().min(0).max(1000000), note: z.string().trim().max(512).optional() }).strict(); const memberRemoveSchema = z.object({ note: z.string().trim().max(512).optional() }).strict(); const completeSchema = z.object({ note: z.string().trim().max(512).optional() }).strict(); const rejectSchema = z.object({ reason: z.string().trim().min(1).max(512) }).strict(); const settlementStatusSchema = z.enum(['DRAFT', 'CONFIRMED', 'PAID', 'CANCELLED']); const settlementListSchema = z.object({ page: z.coerce.number().int().min(1).default(1), pageSize: z.coerce.number().int().min(1).max(50).default(20), status: settlementStatusSchema.optional() }); const settlementParamsSchema = z.object({ settlementId: z.string().regex(/^[1-9]\d{0,19}$/) }); const settlementGenerateSchema = z.object({ cleanerUserId: z.string().regex(/^[1-9]\d{0,19}$/), storeId: z.string().regex(/^[1-9]\d{0,19}$/).optional(), note: z.string().trim().max(512).optional() }).strict(); const settlementConfirmSchema = z.object({ note: z.string().trim().max(512).optional() }).strict(); const settlementPaidSchema = z.object({ payoutChannel: z.string().trim().min(1).max(32), payoutReference: z.string().trim().min(1).max(128), note: z.string().trim().max(512).optional() }).strict(); const settlementPayoutFailureSchema = z.object({ payoutChannel: z.string().trim().max(32).optional(), payoutReference: z.string().trim().max(128).optional(), error: z.string().trim().min(1).max(512), note: z.string().trim().max(512).optional() }).strict(); const settlementWechatTransferSchema = z.object({ mode: z.enum(['API', 'MOCK']).default('API'), note: z.string().trim().max(512).optional() }).strict(); const settlementWechatSyncSchema = z.object({ note: z.string().trim().max(512).optional() }).strict(); const statisticsSchema = z.object({ from: z.string().regex(/^\d{4}-\d{2}-\d{2}(?:[ T]\d{2}:\d{2}:\d{2})?$/).optional(), to: z.string().regex(/^\d{4}-\d{2}-\d{2}(?:[ T]\d{2}:\d{2}:\d{2})?$/).optional(), storeId: z.string().regex(/^[1-9]\d{0,19}$/).optional() }).strict(); const reclaimSchema = z.object({ olderThanMinutes: z.coerce.number().int().min(5).max(1440).default(60), limit: z.coerce.number().int().min(1).max(100).default(20) }); export interface CleaningRouteOptions { repository: Pick; mediaStorage?: MediaStorage; payoutService?: Pick; authRepository: Pick; accessControl: { getAccessProfile(tenantId: string, userId: string): Promise }; jwtSecret: string; } export async function registerCleaningRoutes( app: FastifyInstance, options: CleaningRouteOptions ): Promise { if (options.mediaStorage && !app.hasContentTypeParser('application/octet-stream')) { app.addContentTypeParser( 'application/octet-stream', { parseAs: 'buffer', bodyLimit: 8 * 1024 * 1024 }, (_request, body, done) => done(null, body) ); } app.get('/app-api/cleaning/tasks/hall', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const query = listSchema.safeParse(request.query); if (!query.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.listHall({ ...actor, ...query.data }), traceId: request.traceId })); }); app.get('/app-api/cleaning/tasks/mine', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const query = listSchema.safeParse(request.query); if (!query.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.listMine({ ...actor, ...query.data }), traceId: request.traceId })); }); app.get('/app-api/cleaning/stats', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.stats(actor), traceId: request.traceId })); }); app.get('/admin-api/cleaning/tasks', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const query = listSchema.safeParse(request.query); if (!query.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.listManage({ ...actor, ...query.data }), traceId: request.traceId })); }); app.get('/admin-api/cleaning/statistics', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const query = statisticsSchema.safeParse(request.query); if (!query.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.managerStatistics({ ...actor, ...query.data }), traceId: request.traceId })); }); app.get('/admin-api/cleaning/settlement-candidates', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const query = listSchema.omit({ status: true }).safeParse(request.query); if (!query.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.settlementCandidates({ ...actor, ...query.data }), traceId: request.traceId })); }); app.get('/admin-api/cleaning/settlements', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const query = settlementListSchema.safeParse(request.query); if (!query.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.listSettlements({ ...actor, ...query.data }), traceId: request.traceId })); }); app.get('/admin-api/cleaning/settlements/:settlementId', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const params = settlementParamsSchema.safeParse(request.params); if (!params.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.getSettlementDetail({ ...actor, settlementId: params.data.settlementId }), traceId: request.traceId })); }); app.post('/app-api/cleaning/tasks/:taskId/claim', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); if (!params.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.claim({ ...actor, taskId: params.data.taskId }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/tasks/:taskId/assign', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); const body = assignSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.assign({ ...actor, taskId: params.data.taskId, cleanerUserId: body.data.cleanerUserId, note: body.data.note }), traceId: request.traceId })); }); app.get('/admin-api/cleaning/tasks/:taskId/members', async (request, reply) => { const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const params = paramsSchema.safeParse(request.params); if (!params.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.listMembers({ ...actor, taskId: params.data.taskId }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/tasks/:taskId/members', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); const body = memberSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.addMember({ ...actor, taskId: params.data.taskId, cleanerUserId: body.data.cleanerUserId, rewardCents: body.data.rewardCents, note: body.data.note }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/tasks/:taskId/members/:cleanerUserId/remove', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = memberParamsSchema.safeParse(request.params); const body = memberRemoveSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.removeMember({ ...actor, taskId: params.data.taskId, cleanerUserId: params.data.cleanerUserId, note: body.data.note }), traceId: request.traceId })); }); app.post('/app-api/cleaning/tasks/:taskId/start', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); if (!params.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.start({ ...actor, taskId: params.data.taskId }), traceId: request.traceId })); }); app.post('/app-api/cleaning/tasks/:taskId/rework', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); if (!params.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.rework({ ...actor, taskId: params.data.taskId }), traceId: request.traceId })); }); app.post('/app-api/cleaning/tasks/:taskId/photos', async (request, reply) => { if (!options.mediaStorage) return reply.status(501).send({ code: 'CLEANING_PHOTO_UPLOAD_UNAVAILABLE', message: 'Cleaning photo upload is not configured.', traceId: request.traceId }); const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); const originalName = singleHeader(request.headers['x-file-name']); if (!params.success || !originalName || !Buffer.isBuffer(request.body)) { return invalid(reply, request.traceId); } return handle(reply, request.traceId, async () => { await options.repository.assertCanUploadPhoto({ ...actor, taskId: params.data.taskId }); const image = await options.mediaStorage!.storeImage({ tenantId: actor.tenantId, originalName, contentType: singleHeader(request.headers['x-image-content-type']) ?? '', body: request.body as Buffer }); return reply.status(201).send({ code: 0, data: image, traceId: request.traceId }); }); }); app.post('/app-api/cleaning/tasks/:taskId/submit', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); const body = submitSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.submit({ ...actor, taskId: params.data.taskId, photoUrls: body.data.photoUrls, note: body.data.note }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/tasks/:taskId/complete', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); const body = completeSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.complete({ ...actor, taskId: params.data.taskId, note: body.data.note }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/tasks/:taskId/reject', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = paramsSchema.safeParse(request.params); const body = rejectSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.reject({ ...actor, taskId: params.data.taskId, reason: body.data.reason }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/reclaim-timeouts', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const body = reclaimSchema.safeParse(request.body ?? {}); if (!body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.reclaimTimeouts({ ...actor, ...body.data }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/settlements', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const body = settlementGenerateSchema.safeParse(request.body ?? {}); if (!body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => reply.status(201).send({ code: 0, data: await options.repository.generateSettlement({ ...actor, ...body.data }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/settlements/:settlementId/confirm', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = settlementParamsSchema.safeParse(request.params); const body = settlementConfirmSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.confirmSettlement({ ...actor, settlementId: params.data.settlementId, note: body.data.note }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/settlements/:settlementId/paid', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = settlementParamsSchema.safeParse(request.params); const body = settlementPaidSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.markSettlementPaid({ ...actor, settlementId: params.data.settlementId, payoutChannel: body.data.payoutChannel, payoutReference: body.data.payoutReference, note: body.data.note }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/settlements/:settlementId/payout-failure', async (request, reply) => { const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = settlementParamsSchema.safeParse(request.params); const body = settlementPayoutFailureSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.repository.recordSettlementPayoutFailure({ ...actor, settlementId: params.data.settlementId, payoutChannel: body.data.payoutChannel, payoutReference: body.data.payoutReference, error: body.data.error, note: body.data.note }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/settlements/:settlementId/wechat-transfer', async (request, reply) => { if (!options.payoutService) return reply.status(501).send({ code: 'CLEANING_PAYOUT_UNAVAILABLE', message: 'Cleaning payout service is not configured.', traceId: request.traceId }); const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = settlementParamsSchema.safeParse(request.params); const body = settlementWechatTransferSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.payoutService!.executeWechatTransfer({ ...actor, settlementId: params.data.settlementId, mode: body.data.mode, note: body.data.note }), traceId: request.traceId })); }); app.get('/admin-api/cleaning/settlements/:settlementId/wechat-transfer/preflight', async (request, reply) => { if (!options.payoutService) return reply.status(501).send({ code: 'CLEANING_PAYOUT_UNAVAILABLE', message: 'Cleaning payout service is not configured.', traceId: request.traceId }); const actor = await requireActor(request, reply, options, 'read'); if (!actor) return; const params = settlementParamsSchema.safeParse(request.params); if (!params.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.payoutService!.preflightWechatTransfer({ ...actor, settlementId: params.data.settlementId }), traceId: request.traceId })); }); app.post('/admin-api/cleaning/settlements/:settlementId/wechat-transfer/sync', async (request, reply) => { if (!options.payoutService) return reply.status(501).send({ code: 'CLEANING_PAYOUT_UNAVAILABLE', message: 'Cleaning payout service is not configured.', traceId: request.traceId }); const actor = await requireActor(request, reply, options, 'write'); if (!actor) return; const params = settlementParamsSchema.safeParse(request.params); const body = settlementWechatSyncSchema.safeParse(request.body ?? {}); if (!params.success || !body.success) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => ({ code: 0, data: await options.payoutService!.syncWechatTransfer({ ...actor, settlementId: params.data.settlementId, note: body.data.note }), traceId: request.traceId })); }); app.post('/app-api/cleaning/wechat-transfer/notify', { preParsing: captureRawBody }, async (request, reply) => { if (!options.payoutService) return reply.status(501).send({ code: 'CLEANING_PAYOUT_UNAVAILABLE', message: 'Cleaning payout service is not configured.', traceId: request.traceId }); const headers = wechatHeaders(request); if (!headers) return invalid(reply, request.traceId); return handle(reply, request.traceId, async () => { await options.payoutService!.processWechatTransferNotification( headers, request.rawBody, request.traceId ); return { code: 'SUCCESS', message: '成功', traceId: request.traceId }; }); }); } async function requireActor( request: FastifyRequest, reply: FastifyReply, options: CleaningRouteOptions, mode: 'read' | 'write' ): Promise { const auth = await authenticateAccessToken( request.headers.authorization, options.authRepository, options.jwtSecret ); if (!auth) { reply.status(401).send({ code: 'AUTH_SESSION_INVALID', message: 'Authentication required.', traceId: request.traceId }); return null; } const access = await options.accessControl.getAccessProfile(auth.session.tenantId, auth.session.user.id); const permission = mode === 'read' ? 'cleaning.task.read' : 'cleaning.task.write'; if (!access.capabilities.includes(permission) && !access.capabilities.includes('tenant.manage') && !access.roles.includes('PLATFORM_ADMIN')) { reply.status(403).send({ code: 'CLEANING_TASK_FORBIDDEN', message: 'Cleaning task permission is required.', traceId: request.traceId }); return null; } return { tenantId: auth.session.tenantId, userId: auth.session.user.id, access, traceId: request.traceId }; } async function handle(reply: FastifyReply, traceId: string, work: () => Promise) { try { return await work(); } catch (error) { if (error instanceof MediaValidationError) { return reply.status(400).send({ code: error.code, message: 'The cleaning task request cannot be completed.', traceId }); } if (!(error instanceof CleaningTaskError) && !(error instanceof CleaningPayoutError) && !(error instanceof WechatPayError)) throw error; const code = error.code; const statusCode = code === 'CLEANING_TASK_FORBIDDEN' || code === 'CLEANING_SETTLEMENT_FORBIDDEN' ? 403 : 409; return reply.status(statusCode).send({ code, message: 'The cleaning task request cannot be completed.', traceId }); } } function singleHeader(value: string | string[] | undefined) { return Array.isArray(value) ? value[0] : value; } function invalid(reply: FastifyReply, traceId: string) { return reply.status(400).send({ code: 'INVALID_CLEANING_TASK_REQUEST', message: 'The cleaning task request is invalid.', traceId }); } function wechatHeaders(request: FastifyRequest): WechatNotificationHeaders | null { const read = (key: string) => singleHeader(request.headers[key]); const headers = { timestamp: read('wechatpay-timestamp'), nonce: read('wechatpay-nonce'), serial: read('wechatpay-serial'), signature: read('wechatpay-signature') }; return headers.timestamp && headers.nonce && headers.serial && headers.signature ? headers as WechatNotificationHeaders : null; } const captureRawBody: preParsingHookHandler = (request, _reply, payload, done) => { const chunks: Buffer[] = []; const capture = new Transform({ transform(chunk, _encoding, callback) { chunks.push(Buffer.from(chunk)); callback(null, chunk); }, flush(callback) { request.rawBody = Buffer.concat(chunks).toString('utf8'); callback(); } }); const transformed = payload.pipe(capture) as typeof payload; transformed.receivedEncodedLength = payload.receivedEncodedLength; done(null, transformed); };