[ExecuteSalaryCurrentService] ครอบ transaction + respone success/fail count
All checks were successful
Build & Deploy on Dev / build (push) Successful in 1m3s

This commit is contained in:
harid 2026-06-23 13:34:06 +07:00
parent 00c35e8974
commit 64be68d0a3
3 changed files with 366 additions and 291 deletions

View file

@ -106,7 +106,7 @@ import { reOrderCommandRecivesAndDelete } from "../services/CommandService";
import { RetirementService } from "../services/RetirementService"; import { RetirementService } from "../services/RetirementService";
import { ExecuteOfficerProfileService } from "../services/ExecuteOfficerProfileService"; import { ExecuteOfficerProfileService } from "../services/ExecuteOfficerProfileService";
import { ExecuteSalaryService } from "../services/ExecuteSalaryService"; import { ExecuteSalaryService } from "../services/ExecuteSalaryService";
import { ExecuteSalaryCurrentService } from "../services/ExecuteSalaryCurrentService"; import { ExecuteSalaryCurrentService, ExecuteSalaryResult } from "../services/ExecuteSalaryCurrentService";
import { ExecuteSalaryEmployeeCurrentService } from "../services/ExecuteSalaryEmployeeCurrentService"; import { ExecuteSalaryEmployeeCurrentService } from "../services/ExecuteSalaryEmployeeCurrentService";
import { ExecuteSalaryLeaveService } from "../services/ExecuteSalaryLeaveService"; import { ExecuteSalaryLeaveService } from "../services/ExecuteSalaryLeaveService";
import { ExecuteSalaryEmployeeLeaveService } from "../services/ExecuteSalaryEmployeeLeaveService"; import { ExecuteSalaryEmployeeLeaveService } from "../services/ExecuteSalaryEmployeeLeaveService";
@ -3702,11 +3702,14 @@ export class CommandController extends Controller {
}[]; }[];
}, },
) { ) {
await new ExecuteSalaryCurrentService().executeSalaryCurrent(body.data, { const result: ExecuteSalaryResult = await new ExecuteSalaryCurrentService().executeSalaryCurrent(
body.data,
{
user: { sub: req.user.sub, name: req.user.name }, user: { sub: req.user.sub, name: req.user.name },
req, req,
}); },
return new HttpSuccess(); );
return new HttpSuccess(result);
} }
@Post("excexute/salary-employee-current") @Post("excexute/salary-employee-current")

View file

@ -1,4 +1,4 @@
import { Double } from "typeorm"; import { Double, EntityManager } from "typeorm";
import { AppDataSource } from "../database/data-source"; import { AppDataSource } from "../database/data-source";
import HttpError from "../interfaces/http-error"; import HttpError from "../interfaces/http-error";
import HttpStatusCode from "../interfaces/http-status"; import HttpStatusCode from "../interfaces/http-status";
@ -60,6 +60,16 @@ export interface SalaryCurrentExecutionContext {
req?: any; req?: any;
} }
/**
* batch independent
* fail rollback (per-item transaction)
*/
export interface ExecuteSalaryResult {
successCount: number;
failureCount: number;
failures: { profileId: string; reason: string }[];
}
/** /**
* Service ProfileSalary + () * Service ProfileSalary + ()
* *
@ -69,28 +79,29 @@ export interface SalaryCurrentExecutionContext {
* - consumer rabbitmq handler service (Linear Flow) * - consumer rabbitmq handler service (Linear Flow)
* *
* Behavior preserve CommandController.newSalaryAndUpdateCurrent * Behavior preserve CommandController.newSalaryAndUpdateCurrent
*
* Batch semantics: ประมวลผลทุกคนแบบ sequential ()
* transaction race condition batch
* posMaster/position throw rollback
* success/failure count + fail
*/ */
export class ExecuteSalaryCurrentService { export class ExecuteSalaryCurrentService {
private commandRepository = AppDataSource.getRepository(Command); private commandRepository = AppDataSource.getRepository(Command);
private profileRepository = AppDataSource.getRepository(Profile); private profileRepository = AppDataSource.getRepository(Profile);
private salaryRepo = AppDataSource.getRepository(ProfileSalary);
private salaryHistoryRepo = AppDataSource.getRepository(ProfileSalaryHistory);
private posMasterRepository = AppDataSource.getRepository(PosMaster);
private positionRepository = AppDataSource.getRepository(Position);
private orgRootRepository = AppDataSource.getRepository(OrgRoot); private orgRootRepository = AppDataSource.getRepository(OrgRoot);
/** /**
* ProfileSalary + * ProfileSalary + batch
*
* @returns success/failure
*/ */
async executeSalaryCurrent( async executeSalaryCurrent(
data: SalaryCurrentItem[], data: SalaryCurrentItem[],
ctx: SalaryCurrentExecutionContext, ctx: SalaryCurrentExecutionContext,
): Promise<void> { ): Promise<ExecuteSalaryResult> {
console.log("[ExecuteSalaryCurrentService] Starting executeSalaryCurrent"); console.log("[ExecuteSalaryCurrentService] Starting executeSalaryCurrent");
console.log("[ExecuteSalaryCurrentService] Request body count:", data?.length); console.log("[ExecuteSalaryCurrentService] Request body count:", data?.length);
const req = ctx.req;
// ───────────────────────────────────────────────────────────── // ─────────────────────────────────────────────────────────────
// Normalize date fields (ผ่าน handler จะได้ string → ต้องแปลงเป็น Date) // Normalize date fields (ผ่าน handler จะได้ string → ต้องแปลงเป็น Date)
// ───────────────────────────────────────────────────────────── // ─────────────────────────────────────────────────────────────
@ -149,14 +160,69 @@ export class ExecuteSalaryCurrentService {
.orgRootShortName ?? ""; .orgRootShortName ?? "";
} }
} }
await Promise.all(
data.map(async (item) => { // ─────────────────────────────────────────────────────────────
const profile: any = await this.profileRepository.findOneBy({ id: item.profileId }); // Per-item transaction: แต่ละคนมี transaction ของตัวเอง (sequential)
// ประมวลทีละคนเพื่อกัน race condition เมื่อหลายคนใน batch อ้างอิง
// posMaster/position ตัวเดียวกัน คนที่ throw จะ rollback เฉพาะตัว (manager)
// และไม่กระทบคนอื่นใน batch
// ─────────────────────────────────────────────────────────────
const failures: ExecuteSalaryResult["failures"] = [];
let successCount = 0;
for (const item of data ?? []) {
try {
await AppDataSource.transaction(async (manager) => {
await this.processOne(item, ctx, manager, _posNumCodeSit, _posNumCodeSitAbb);
});
successCount++;
} catch (err) {
const reason =
err instanceof HttpError
? err.message
: err instanceof Error
? err.message
: "unexpected error";
console.error(
`[ExecuteSalaryCurrentService] Failed profileId=${item.profileId}: ${reason}`,
err,
);
failures.push({ profileId: item.profileId ?? "unknown", reason });
}
}
console.log(
`[ExecuteSalaryCurrentService] executeSalaryCurrent completed — success: ${successCount}, failure: ${failures.length}`,
);
return { successCount, failureCount: failures.length, failures };
}
/**
* 1 transaction (manager)
* save manager.getRepository(...) transaction
* throw rollback ( partial commit)
*/
private async processOne(
item: SalaryCurrentItem,
ctx: SalaryCurrentExecutionContext,
manager: EntityManager,
_posNumCodeSit: string,
_posNumCodeSitAbb: string,
): Promise<void> {
const req = ctx.req;
const profileRepository = manager.getRepository(Profile);
const salaryRepo = manager.getRepository(ProfileSalary);
const salaryHistoryRepo = manager.getRepository(ProfileSalaryHistory);
const posMasterRepository = manager.getRepository(PosMaster);
const positionRepository = manager.getRepository(Position);
const profile: any = await profileRepository.findOneBy({ id: item.profileId });
if (!profile) { if (!profile) {
throw new HttpError(HttpStatusCode.NOT_FOUND, "ไม่พบข้อมูลทะเบียนประวัตินี้"); throw new HttpError(HttpStatusCode.NOT_FOUND, "ไม่พบข้อมูลทะเบียนประวัตินี้");
} }
let _null: any = null; let _null: any = null;
const dest_item = await this.salaryRepo.findOne({ const dest_item = await salaryRepo.findOne({
where: { profileId: item.profileId }, where: { profileId: item.profileId },
order: { order: "DESC" }, order: { order: "DESC" },
}); });
@ -177,14 +243,14 @@ export class ExecuteSalaryCurrentService {
Object.assign(dataSalary, { ...item, ...meta }); Object.assign(dataSalary, { ...item, ...meta });
const history = new ProfileSalaryHistory(); const history = new ProfileSalaryHistory();
Object.assign(history, { ...dataSalary, id: undefined }); Object.assign(history, { ...dataSalary, id: undefined });
await this.salaryRepo.save(dataSalary, { data: req }); await salaryRepo.save(dataSalary, { data: req });
setLogDataDiff(req, { before, after: dataSalary }); setLogDataDiff(req, { before, after: dataSalary });
history.commandId = item.commandId ?? _null; history.commandId = item.commandId ?? _null;
history.profileSalaryId = dataSalary.id; history.profileSalaryId = dataSalary.id;
await this.salaryHistoryRepo.save(history, { data: req }); await salaryHistoryRepo.save(history, { data: req });
// STEP 1: หา posMaster ที่จะใช้งานตาม id ที่ส่งมา // STEP 1: หา posMaster ที่จะใช้งานตาม id ที่ส่งมา
let posMaster = await this.posMasterRepository.findOne({ let posMaster = await posMasterRepository.findOne({
where: { id: item.posmasterId }, where: { id: item.posmasterId },
relations: { relations: {
orgRevision: true, orgRevision: true,
@ -203,7 +269,7 @@ export class ExecuteSalaryCurrentService {
// ถ้าไม่อยู่ในโครงสร้างปัจจุบัน ให้หาตัวใหม่จาก ancestorDNA // ถ้าไม่อยู่ในโครงสร้างปัจจุบัน ให้หาตัวใหม่จาก ancestorDNA
if (!isCurrent && posMaster?.ancestorDNA) { if (!isCurrent && posMaster?.ancestorDNA) {
posMaster = await this.posMasterRepository.findOne({ posMaster = await posMasterRepository.findOne({
where: { where: {
ancestorDNA: posMaster.ancestorDNA, ancestorDNA: posMaster.ancestorDNA,
orgRevision: { orgRevision: {
@ -229,7 +295,7 @@ export class ExecuteSalaryCurrentService {
throw new HttpError(HttpStatusCode.NOT_FOUND, "ไม่พบข้อมูลตำแหน่งนี้"); throw new HttpError(HttpStatusCode.NOT_FOUND, "ไม่พบข้อมูลตำแหน่งนี้");
} }
const posMasterOld = await this.posMasterRepository.findOne({ const posMasterOld = await posMasterRepository.findOne({
where: { where: {
current_holderId: item.profileId, current_holderId: item.profileId,
orgRevisionId: posMaster.orgRevisionId, orgRevisionId: posMaster.orgRevisionId,
@ -240,7 +306,7 @@ export class ExecuteSalaryCurrentService {
posMasterOld.lastUpdatedAt = new Date(); posMasterOld.lastUpdatedAt = new Date();
} }
const positionOld = await this.positionRepository.findOne({ const positionOld = await positionRepository.findOne({
where: { where: {
posMasterId: posMasterOld?.id, posMasterId: posMasterOld?.id,
positionIsSelected: true, positionIsSelected: true,
@ -255,10 +321,10 @@ export class ExecuteSalaryCurrentService {
}); });
positionOld.positionIsSelected = false; positionOld.positionIsSelected = false;
await this.positionRepository.save(positionOld); await positionRepository.save(positionOld);
} }
const checkPosition = await this.positionRepository.find({ const checkPosition = await positionRepository.find({
where: { where: {
posMasterId: posMaster!.id, // ใช้ posMaster ตัวใหม่ (ที่อาจจะเปลี่ยนจาก ancestorDNA) posMasterId: posMaster!.id, // ใช้ posMaster ตัวใหม่ (ที่อาจจะเปลี่ยนจาก ancestorDNA)
positionIsSelected: true, positionIsSelected: true,
@ -282,7 +348,7 @@ export class ExecuteSalaryCurrentService {
positionIsSelected: false, positionIsSelected: false,
}; };
}); });
await this.positionRepository.save(clearPosition); await positionRepository.save(clearPosition);
} }
posMaster.current_holderId = item.profileId; posMaster.current_holderId = item.profileId;
@ -290,10 +356,11 @@ export class ExecuteSalaryCurrentService {
// posMaster.conditionReason = _null; // posMaster.conditionReason = _null;
// posMaster.isCondition = false; // posMaster.isCondition = false;
if (posMasterOld != null) { if (posMasterOld != null) {
await this.posMasterRepository.save(posMasterOld); await posMasterRepository.save(posMasterOld);
await CreatePosMasterHistoryOfficer(posMasterOld.id, req); // ส่ง manager เข้าไปเพื่อให้อยู่ใน transaction เดียวกัน
await CreatePosMasterHistoryOfficer(posMasterOld.id, req, null, null, manager);
} }
await this.posMasterRepository.save(posMaster); await posMasterRepository.save(posMaster);
// STEP 2: กำหนด position ใหม่ // STEP 2: กำหนด position ใหม่
// Match position ตามลำดับ priority: // Match position ตามลำดับ priority:
@ -312,7 +379,7 @@ export class ExecuteSalaryCurrentService {
// CONDITION 1: เช็คจาก positionId ตรง // CONDITION 1: เช็คจาก positionId ตรง
// ═══════════════════════════════════════════════════════════ // ═══════════════════════════════════════════════════════════
if (item.positionId) { if (item.positionId) {
const positionById = await this.positionRepository.findOne({ const positionById = await positionRepository.findOne({
where: { where: {
id: item.positionId, id: item.positionId,
posMasterId: posMaster.id, // ต้องอยู่ใน posMaster ที่ถูกต้อง posMasterId: posMaster.id, // ต้องอยู่ใน posMaster ที่ถูกต้อง
@ -351,7 +418,7 @@ export class ExecuteSalaryCurrentService {
whereCondition.positionArea = item.positionArea; whereCondition.positionArea = item.positionArea;
} }
const positionBy7Fields = await this.positionRepository.findOne({ const positionBy7Fields = await positionRepository.findOne({
where: whereCondition, where: whereCondition,
relations: ["posExecutive"], relations: ["posExecutive"],
order: { orderNo: "ASC" }, order: { orderNo: "ASC" },
@ -366,7 +433,7 @@ export class ExecuteSalaryCurrentService {
// CONDITION 3: Match 3 ฟิลด์ (ถ้า Condition 2 ไม่ match) // CONDITION 3: Match 3 ฟิลด์ (ถ้า Condition 2 ไม่ match)
// ═══════════════════════════════════════════════════════════ // ═══════════════════════════════════════════════════════════
if (!positionNew && item.positionName && posTypeId && posLevelId) { if (!positionNew && item.positionName && posTypeId && posLevelId) {
const positionBy3Fields = await this.positionRepository.findOne({ const positionBy3Fields = await positionRepository.findOne({
where: { where: {
posMasterId: posMaster.id, posMasterId: posMaster.id,
positionName: item.positionName, positionName: item.positionName,
@ -386,7 +453,7 @@ export class ExecuteSalaryCurrentService {
// // FALLBACK: ถ้าทั้ง 3 ไม่ match ให้เลือก position แรกใน posMaster // // FALLBACK: ถ้าทั้ง 3 ไม่ match ให้เลือก position แรกใน posMaster
// // ═══════════════════════════════════════════════════════════ // // ═══════════════════════════════════════════════════════════
// if (!positionNew) { // if (!positionNew) {
// const fallbackPositions = await this.positionRepository.find({ // const fallbackPositions = await positionRepository.find({
// where: { // where: {
// posMasterId: posMaster.id, // posMasterId: posMaster.id,
// }, // },
@ -419,13 +486,10 @@ export class ExecuteSalaryCurrentService {
} }
profile.amount = item.amount ?? null; profile.amount = item.amount ?? null;
profile.amountSpecial = item.amountSpecial ?? null; profile.amountSpecial = item.amountSpecial ?? null;
await this.profileRepository.save(profile); await profileRepository.save(profile);
await this.positionRepository.save(positionNew); await positionRepository.save(positionNew);
} }
await CreatePosMasterHistoryOfficer(posMaster.id, req); // ส่ง manager เข้าไปเพื่อให้อยู่ใน transaction เดียวกัน
}), await CreatePosMasterHistoryOfficer(posMaster.id, req, null, null, manager);
);
console.log("[ExecuteSalaryCurrentService] executeSalaryCurrent completed successfully");
} }
} }

View file

@ -388,8 +388,16 @@ async function handler(msg: amqp.ConsumeMessage): Promise<boolean> {
await new ExecuteOfficerProfileService().executeCreateOfficerProfile(resultData, ctx); await new ExecuteOfficerProfileService().executeCreateOfficerProfile(resultData, ctx);
console.log(`[AMQ] Processed ${resultData.length} profiles via ExecuteOfficerProfileService`); console.log(`[AMQ] Processed ${resultData.length} profiles via ExecuteOfficerProfileService`);
} else if (isSalaryCurrent) { } else if (isSalaryCurrent) {
await new ExecuteSalaryCurrentService().executeSalaryCurrent(resultData, ctx); const salaryResult = await new ExecuteSalaryCurrentService().executeSalaryCurrent(
console.log(`[AMQ] Processed ${resultData.length} profiles via ExecuteSalaryCurrentService`); resultData,
ctx,
);
console.log(
`[AMQ] Processed via ExecuteSalaryCurrentService — success: ${salaryResult.successCount}, failure: ${salaryResult.failureCount}`,
);
for (const f of salaryResult.failures) {
console.error(`[AMQ] ExecuteSalaryCurrentService failed profileId=${f.profileId}: ${f.reason}`);
}
} else if (isSalaryEmployeeCurrent) { } else if (isSalaryEmployeeCurrent) {
await new ExecuteSalaryEmployeeCurrentService().executeSalaryEmployeeCurrent(resultData, ctx); await new ExecuteSalaryEmployeeCurrentService().executeSalaryEmployeeCurrent(resultData, ctx);
console.log(`[AMQ] Processed ${resultData.length} profiles via ExecuteSalaryEmployeeCurrentService`); console.log(`[AMQ] Processed ${resultData.length} profiles via ExecuteSalaryEmployeeCurrentService`);