Isolate project device sync transactions

This commit is contained in:
2026-07-23 18:57:18 +02:00
parent cf356bbedc
commit 91baa6c76c
10 changed files with 414 additions and 126 deletions
@@ -1,5 +1,5 @@
import crypto from "node:crypto";
import { and, asc, eq, inArray, isNull } from "drizzle-orm";
import { and, asc, eq, inArray } from "drizzle-orm";
import { db } from "../client.js";
import { circuitDeviceRows } from "../schema/circuit-device-rows.js";
import { circuitLists } from "../schema/circuit-lists.js";
@@ -125,64 +125,6 @@ export class CircuitDeviceRowRepository {
await db.update(circuitDeviceRows).set(values).where(eq(circuitDeviceRows.id, rowId));
}
updateLinkedRowsTransactional(
projectDeviceId: string,
changes: Array<{ rowId: string; input: CircuitDeviceRowUpdateInput }>
) {
db.transaction((tx) => {
for (const change of changes) {
const result = tx
.update(circuitDeviceRows)
.set(toCircuitDeviceRowUpdateValues(change.input))
.where(
and(
eq(circuitDeviceRows.id, change.rowId),
eq(circuitDeviceRows.linkedProjectDeviceId, projectDeviceId)
)
)
.run();
if (result.changes !== 1) {
throw new Error("A linked row changed before the operation could be completed.");
}
}
});
}
disconnectLinkedRowsTransactional(projectDeviceId: string, rowIds: string[]) {
db.transaction((tx) => {
for (const rowId of rowIds) {
const result = tx
.update(circuitDeviceRows)
.set({ linkedProjectDeviceId: null })
.where(
and(
eq(circuitDeviceRows.id, rowId),
eq(circuitDeviceRows.linkedProjectDeviceId, projectDeviceId)
)
)
.run();
if (result.changes !== 1) {
throw new Error("A linked row changed before the operation could be completed.");
}
}
});
}
reconnectRowsTransactional(projectDeviceId: string, rowIds: string[]) {
db.transaction((tx) => {
for (const rowId of rowIds) {
const result = tx
.update(circuitDeviceRows)
.set({ linkedProjectDeviceId: projectDeviceId })
.where(and(eq(circuitDeviceRows.id, rowId), isNull(circuitDeviceRows.linkedProjectDeviceId)))
.run();
if (result.changes !== 1) {
throw new Error("A disconnected row changed before the operation could be undone.");
}
}
});
}
async delete(rowId: string) {
await db.delete(circuitDeviceRows).where(eq(circuitDeviceRows.id, rowId));
}
@@ -0,0 +1,81 @@
import { and, eq, isNull } from "drizzle-orm";
import type {
ProjectDeviceLinkedRowValues,
ProjectDeviceRowSyncStore,
} from "../../domain/ports/project-device-row-sync.store.js";
import type { AppDatabase } from "../database-context.js";
import { circuitDeviceRows } from "../schema/circuit-device-rows.js";
import { toCircuitDeviceRowUpdateValues } from "./circuit-device-row.persistence.js";
export class ProjectDeviceRowSyncRepository implements ProjectDeviceRowSyncStore {
constructor(private readonly database: AppDatabase) {}
updateLinkedRows(
projectDeviceId: string,
changes: Array<{ rowId: string; input: ProjectDeviceLinkedRowValues }>
) {
this.database.transaction((tx) => {
for (const change of changes) {
const result = tx
.update(circuitDeviceRows)
.set(toCircuitDeviceRowUpdateValues(change.input))
.where(
and(
eq(circuitDeviceRows.id, change.rowId),
eq(circuitDeviceRows.linkedProjectDeviceId, projectDeviceId)
)
)
.run();
if (result.changes !== 1) {
throw new Error(
"A linked row changed before the operation could be completed."
);
}
}
});
}
disconnectLinkedRows(projectDeviceId: string, rowIds: string[]) {
this.database.transaction((tx) => {
for (const rowId of rowIds) {
const result = tx
.update(circuitDeviceRows)
.set({ linkedProjectDeviceId: null })
.where(
and(
eq(circuitDeviceRows.id, rowId),
eq(circuitDeviceRows.linkedProjectDeviceId, projectDeviceId)
)
)
.run();
if (result.changes !== 1) {
throw new Error(
"A linked row changed before the operation could be completed."
);
}
}
});
}
reconnectRows(projectDeviceId: string, rowIds: string[]) {
this.database.transaction((tx) => {
for (const rowId of rowIds) {
const result = tx
.update(circuitDeviceRows)
.set({ linkedProjectDeviceId: projectDeviceId })
.where(
and(
eq(circuitDeviceRows.id, rowId),
isNull(circuitDeviceRows.linkedProjectDeviceId)
)
)
.run();
if (result.changes !== 1) {
throw new Error(
"A disconnected row changed before the operation could be undone."
);
}
}
});
}
}
@@ -0,0 +1,28 @@
export interface ProjectDeviceLinkedRowValues {
linkedProjectDeviceId?: string;
name: string;
displayName: string;
phaseType?: string;
connectionKind?: string;
costGroup?: string;
category?: string;
level?: string;
roomId?: string;
roomNumberSnapshot?: string;
roomNameSnapshot?: string;
quantity: number;
powerPerUnit: number;
simultaneityFactor: number;
cosPhi?: number;
remark?: string;
overriddenFields?: string;
}
export interface ProjectDeviceRowSyncStore {
updateLinkedRows(
projectDeviceId: string,
changes: Array<{ rowId: string; input: ProjectDeviceLinkedRowValues }>
): void;
disconnectLinkedRows(projectDeviceId: string, rowIds: string[]): void;
reconnectRows(projectDeviceId: string, rowIds: string[]): void;
}
@@ -1,5 +1,6 @@
import { CircuitDeviceRowRepository } from "../../db/repositories/circuit-device-row.repository.js";
import { ProjectDeviceRepository } from "../../db/repositories/project-device.repository.js";
import type { ProjectDeviceRowSyncStore } from "../ports/project-device-row-sync.store.js";
import {
projectDeviceSyncFields,
type ProjectDeviceSyncField,
@@ -14,10 +15,8 @@ type SyncDependencies = {
CircuitDeviceRowRepository,
| "listLinkedByProjectDevice"
| "listLinkStatesByIds"
| "updateLinkedRowsTransactional"
| "disconnectLinkedRowsTransactional"
| "reconnectRowsTransactional"
>;
deviceRowSyncStore: ProjectDeviceRowSyncStore;
};
export interface ProjectDeviceSyncRestoreRow {
@@ -71,10 +70,19 @@ function buildDifferences(projectDevice: ProjectDevice, row: LinkedRow) {
export class ProjectDeviceSyncService {
private readonly projectDeviceRepository: SyncDependencies["projectDeviceRepository"];
private readonly deviceRowRepository: SyncDependencies["deviceRowRepository"];
private readonly deviceRowSyncStore?: SyncDependencies["deviceRowSyncStore"];
constructor(deps?: Partial<SyncDependencies>) {
this.projectDeviceRepository = deps?.projectDeviceRepository ?? new ProjectDeviceRepository();
this.deviceRowRepository = deps?.deviceRowRepository ?? new CircuitDeviceRowRepository();
this.deviceRowSyncStore = deps?.deviceRowSyncStore;
}
private getDeviceRowSyncStore() {
if (!this.deviceRowSyncStore) {
throw new Error("Project-device row synchronization is not configured.");
}
return this.deviceRowSyncStore;
}
async getPreview(projectId: string, projectDeviceId: string) {
@@ -137,7 +145,7 @@ export class ProjectDeviceSyncService {
next.overriddenFields = serializeOverriddenFields(remainingOverrides);
changes.push({ rowId: row.id, input: next });
}
this.deviceRowRepository.updateLinkedRowsTransactional(projectDeviceId, changes);
this.getDeviceRowSyncStore().updateLinkedRows(projectDeviceId, changes);
return {
preview: await this.getPreview(projectId, projectDeviceId),
@@ -150,7 +158,10 @@ export class ProjectDeviceSyncService {
const selectedRows = this.resolveSelectedRows(linkedRows, rowIds);
const disconnectedRowIds = selectedRows.map((row) => row.id);
this.deviceRowRepository.disconnectLinkedRowsTransactional(projectDeviceId, disconnectedRowIds);
this.getDeviceRowSyncStore().disconnectLinkedRows(
projectDeviceId,
disconnectedRowIds
);
return { disconnectedRowIds, undo: { rowIds: disconnectedRowIds } };
}
@@ -175,7 +186,7 @@ export class ProjectDeviceSyncService {
return { rowId: row.id, input: next };
});
this.deviceRowRepository.updateLinkedRowsTransactional(projectDeviceId, changes);
this.getDeviceRowSyncStore().updateLinkedRows(projectDeviceId, changes);
return this.getPreview(projectId, projectDeviceId);
}
@@ -189,7 +200,7 @@ export class ProjectDeviceSyncService {
if (linkStates.length !== uniqueRowIds.length || linkStates.some((row) => row.linkedProjectDeviceId !== null)) {
throw new Error("One or more rows cannot be reconnected.");
}
this.deviceRowRepository.reconnectRowsTransactional(projectDeviceId, uniqueRowIds);
this.getDeviceRowSyncStore().reconnectRows(projectDeviceId, uniqueRowIds);
return this.getPreview(projectId, projectDeviceId);
}
@@ -0,0 +1,7 @@
import { db } from "../../db/client.js";
import { ProjectDeviceRowSyncRepository } from "../../db/repositories/project-device-row-sync.repository.js";
import { ProjectDeviceSyncService } from "../../domain/services/project-device-sync.service.js";
export const projectDeviceSyncService = new ProjectDeviceSyncService({
deviceRowSyncStore: new ProjectDeviceRowSyncRepository(db),
});
@@ -1,7 +1,7 @@
import type { Request, Response } from "express";
import { GlobalDeviceRepository } from "../../db/repositories/global-device.repository.js";
import { ProjectDeviceRepository } from "../../db/repositories/project-device.repository.js";
import { ProjectDeviceSyncService } from "../../domain/services/project-device-sync.service.js";
import { projectDeviceSyncService } from "../composition/project-device-sync-service.js";
import {
createProjectDeviceSchema,
disconnectProjectDeviceRowsSchema,
@@ -13,8 +13,6 @@ import {
const globalDeviceRepository = new GlobalDeviceRepository();
const projectDeviceRepository = new ProjectDeviceRepository();
const projectDeviceSyncService = new ProjectDeviceSyncService();
export async function listProjectDevicesByProject(req: Request, res: Response) {
const { projectId } = req.params;
if (typeof projectId !== "string") {