forked from StellarSplit/StellarSplit
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathcompliance.processor.ts
More file actions
162 lines (147 loc) · 5.57 KB
/
Copy pathcompliance.processor.ts
File metadata and controls
162 lines (147 loc) · 5.57 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
import { Process, Processor, OnQueueFailed } from "@nestjs/bull";
import { Job } from "bull";
import { InjectRepository } from "@nestjs/typeorm";
import { Repository, Between } from "typeorm";
import {
TaxExportRequest,
ExportStatus,
ExportFormat,
} from "./entities/tax-export-request.entity";
import { Split } from "../entities/split.entity";
import { CSVExporterService } from "./exporters/csv-exporter.service";
import { PDFExporterService } from "./exporters/pdf-exporter.service";
import { QBOExporterService } from "./exporters/qbo-exporter.service";
import { JSONExporterService } from "./exporters/json-exporter.service";
import { OFXExporterService } from "./exporters/ofx-exporter.service";
import { ProfileService } from "../profile/profile.service";
import { EmailService } from "../email/email.service";
import { Logger } from "@nestjs/common";
import * as fs from "fs";
import * as path from "path";
import { logJobFailure } from "../common/queue-job-policy";
@Processor("compliance-export")
export class ComplianceProcessor {
private readonly logger = new Logger(ComplianceProcessor.name);
private readonly exportDir = path.join(process.cwd(), "exports");
constructor(
@InjectRepository(TaxExportRequest)
private exportRepo: Repository<TaxExportRequest>,
@InjectRepository(Split)
private splitRepo: Repository<Split>,
private csvExporter: CSVExporterService,
private pdfExporter: PDFExporterService,
private qboExporter: QBOExporterService,
private jsonExporter: JSONExporterService,
private ofxExporter: OFXExporterService,
private emailService: EmailService,
private profileService: ProfileService,
) {
if (!fs.existsSync(this.exportDir)) {
fs.mkdirSync(this.exportDir);
}
}
@Process("generate-export")
async handleExport(job: Job<{ requestId: string }>) {
const { requestId } = job.data;
const request = await this.exportRepo.findOne({ where: { id: requestId } });
if (!request) {
this.logger.error(`Export request ${requestId} not found`);
return;
}
try {
await this.exportRepo.update(requestId, {
status: ExportStatus.PROCESSING,
});
const splits = await this.splitRepo.find({
where: {
creatorWalletAddress: request.userId,
createdAt: Between(request.periodStart, request.periodEnd),
},
relations: ["category"],
});
let content: string | Buffer;
let filename = `tax-export-${requestId}`;
switch (request.exportFormat) {
case ExportFormat.CSV:
content = await this.csvExporter.generate(splits);
filename += ".csv";
break;
case ExportFormat.PDF:
content = await this.pdfExporter.generate(splits);
filename += ".pdf";
break;
case ExportFormat.QBO:
content = await this.qboExporter.generate(splits);
filename += ".qbo";
break;
case ExportFormat.JSON:
content = await this.jsonExporter.generate(splits);
filename += ".json";
break;
case ExportFormat.OFX:
content = await this.ofxExporter.generate(splits);
filename += ".ofx";
break;
default:
throw new Error(`Unsupported export format: ${request.exportFormat}`);
}
const filePath = path.join(this.exportDir, filename);
fs.writeFileSync(filePath, content);
const expiresAt = new Date();
expiresAt.setHours(expiresAt.getHours() + 48);
await this.exportRepo.update(requestId, {
status: ExportStatus.READY,
fileUrl: `http://localhost:3000/api/compliance/export/${requestId}/download`, // Secure download URL
fileSize: fs.statSync(filePath).size,
recordCount: splits.length,
completedAt: new Date(),
expiresAt,
});
// Send email notification
try {
const profile = await this.profileService.getByWalletAddress(request.userId);
const userEmail = profile.email || 'user@example.com'; // fallback
await this.emailService["emailQueue"].add("sendEmail", {
to: userEmail,
type: "export_ready",
context: {
requestId,
format: request.exportFormat,
downloadUrl: `http://localhost:3000/api/compliance/export/${requestId}/download`,
},
});
} catch (error) {
this.logger.error(`Failed to get user email for export ${requestId}:`, error);
// Fallback to placeholder
await this.emailService["emailQueue"].add("sendEmail", {
to: "user@example.com",
type: "export_ready",
context: {
requestId,
format: request.exportFormat,
downloadUrl: `http://localhost:3000/api/compliance/export/${requestId}/download`,
},
});
}
this.logger.log(`Export ${requestId} completed successfully`);
} catch (error) {
logJobFailure(this.logger, job, error, { context: 'compliance-export' });
await this.exportRepo.update(requestId, { status: ExportStatus.FAILED });
}
}
/**
* Dead-letter handler: fires when the compliance export job has exhausted all retries.
*/
@OnQueueFailed()
async onFailed(job: Job<{ requestId: string }>, err: Error) {
logJobFailure(this.logger, job, err, { context: 'compliance-export-dead-letter' });
const { requestId } = job.data;
if (requestId) {
try {
await this.exportRepo.update(requestId, { status: ExportStatus.FAILED });
} catch {
// best effort – the DB may be the cause of the failure
}
}
}
}