Reuse command transaction boundary
This commit is contained in:
@@ -15,8 +15,7 @@ import { circuitDeviceRows } from "../schema/circuit-device-rows.js";
|
||||
import { circuitLists } from "../schema/circuit-lists.js";
|
||||
import { circuits } from "../schema/circuits.js";
|
||||
import { projectDevices } from "../schema/project-devices.js";
|
||||
import { applyProjectHistoryTransition } from "./project-history.persistence.js";
|
||||
import { appendProjectRevision } from "./project-revision.persistence.js";
|
||||
import { executeProjectCommandTransaction } from "./project-command-transaction.persistence.js";
|
||||
|
||||
type CircuitDeviceRow = typeof circuitDeviceRows.$inferSelect;
|
||||
|
||||
@@ -28,116 +27,110 @@ export class ProjectDeviceRowSyncProjectCommandRepository
|
||||
execute(input: ExecuteProjectDeviceRowSyncCommandInput) {
|
||||
assertProjectDeviceRowSyncProjectCommand(input.command);
|
||||
|
||||
return this.database.transaction((tx) => {
|
||||
const projectDevice = tx
|
||||
.select({ id: projectDevices.id })
|
||||
.from(projectDevices)
|
||||
.where(
|
||||
and(
|
||||
eq(
|
||||
projectDevices.id,
|
||||
input.command.payload.projectDeviceId
|
||||
),
|
||||
eq(projectDevices.projectId, input.projectId)
|
||||
)
|
||||
)
|
||||
.get();
|
||||
if (!projectDevice) {
|
||||
throw new Error(
|
||||
"Project device does not belong to project."
|
||||
);
|
||||
}
|
||||
return executeProjectCommandTransaction(
|
||||
this.database,
|
||||
input,
|
||||
(tx) => this.applyCommand(tx, input)
|
||||
);
|
||||
}
|
||||
|
||||
const rowIds = input.command.payload.rows.map(
|
||||
(row) => row.rowId
|
||||
private applyCommand(
|
||||
tx: AppDatabase,
|
||||
input: ExecuteProjectDeviceRowSyncCommandInput
|
||||
) {
|
||||
const projectDevice = tx
|
||||
.select({ id: projectDevices.id })
|
||||
.from(projectDevices)
|
||||
.where(
|
||||
and(
|
||||
eq(
|
||||
projectDevices.id,
|
||||
input.command.payload.projectDeviceId
|
||||
),
|
||||
eq(projectDevices.projectId, input.projectId)
|
||||
)
|
||||
)
|
||||
.get();
|
||||
if (!projectDevice) {
|
||||
throw new Error(
|
||||
"Project device does not belong to project."
|
||||
);
|
||||
const persistedRows = tx
|
||||
.select({
|
||||
row: circuitDeviceRows,
|
||||
projectId: circuitLists.projectId,
|
||||
})
|
||||
.from(circuitDeviceRows)
|
||||
.innerJoin(
|
||||
circuits,
|
||||
eq(circuits.id, circuitDeviceRows.circuitId)
|
||||
)
|
||||
.innerJoin(
|
||||
circuitLists,
|
||||
eq(circuitLists.id, circuits.circuitListId)
|
||||
)
|
||||
.where(inArray(circuitDeviceRows.id, rowIds))
|
||||
.all();
|
||||
}
|
||||
|
||||
const rowIds = input.command.payload.rows.map(
|
||||
(row) => row.rowId
|
||||
);
|
||||
const persistedRows = tx
|
||||
.select({
|
||||
row: circuitDeviceRows,
|
||||
projectId: circuitLists.projectId,
|
||||
})
|
||||
.from(circuitDeviceRows)
|
||||
.innerJoin(
|
||||
circuits,
|
||||
eq(circuits.id, circuitDeviceRows.circuitId)
|
||||
)
|
||||
.innerJoin(
|
||||
circuitLists,
|
||||
eq(circuitLists.id, circuits.circuitListId)
|
||||
)
|
||||
.where(inArray(circuitDeviceRows.id, rowIds))
|
||||
.all();
|
||||
if (
|
||||
persistedRows.length !== rowIds.length ||
|
||||
persistedRows.some(
|
||||
(persisted) => persisted.projectId !== input.projectId
|
||||
)
|
||||
) {
|
||||
throw new Error(
|
||||
"One or more sync rows do not belong to project."
|
||||
);
|
||||
}
|
||||
|
||||
const persistedById = new Map(
|
||||
persistedRows.map((persisted) => [
|
||||
persisted.row.id,
|
||||
persisted.row,
|
||||
])
|
||||
);
|
||||
for (const assignment of input.command.payload.rows) {
|
||||
const persisted = persistedById.get(assignment.rowId);
|
||||
if (
|
||||
persistedRows.length !== rowIds.length ||
|
||||
persistedRows.some(
|
||||
(persisted) => persisted.projectId !== input.projectId
|
||||
)
|
||||
!persisted ||
|
||||
!snapshotMatchesRow(assignment.expected, persisted)
|
||||
) {
|
||||
throw new Error(
|
||||
"One or more sync rows do not belong to project."
|
||||
"Project-device row changed before sync execution."
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const persistedById = new Map(
|
||||
persistedRows.map((persisted) => [
|
||||
persisted.row.id,
|
||||
persisted.row,
|
||||
])
|
||||
);
|
||||
for (const assignment of input.command.payload.rows) {
|
||||
const persisted = persistedById.get(assignment.rowId);
|
||||
if (
|
||||
!persisted ||
|
||||
!snapshotMatchesRow(assignment.expected, persisted)
|
||||
) {
|
||||
throw new Error(
|
||||
"Project-device row changed before sync execution."
|
||||
);
|
||||
}
|
||||
const inverse = createProjectDeviceRowSyncProjectCommand(
|
||||
projectDevice.id,
|
||||
invertProjectDeviceRowSyncOperation(
|
||||
input.command.payload.operation
|
||||
),
|
||||
input.command.payload.rows.map((row) => ({
|
||||
rowId: row.rowId,
|
||||
expected: row.target,
|
||||
target: row.expected,
|
||||
}))
|
||||
);
|
||||
|
||||
for (const assignment of input.command.payload.rows) {
|
||||
const updated = tx
|
||||
.update(circuitDeviceRows)
|
||||
.set(assignment.target)
|
||||
.where(eq(circuitDeviceRows.id, assignment.rowId))
|
||||
.run();
|
||||
if (updated.changes !== 1) {
|
||||
throw new Error(
|
||||
"Project-device row changed during sync execution."
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const inverse = createProjectDeviceRowSyncProjectCommand(
|
||||
projectDevice.id,
|
||||
invertProjectDeviceRowSyncOperation(
|
||||
input.command.payload.operation
|
||||
),
|
||||
input.command.payload.rows.map((row) => ({
|
||||
rowId: row.rowId,
|
||||
expected: row.target,
|
||||
target: row.expected,
|
||||
}))
|
||||
);
|
||||
|
||||
for (const assignment of input.command.payload.rows) {
|
||||
const updated = tx
|
||||
.update(circuitDeviceRows)
|
||||
.set(assignment.target)
|
||||
.where(eq(circuitDeviceRows.id, assignment.rowId))
|
||||
.run();
|
||||
if (updated.changes !== 1) {
|
||||
throw new Error(
|
||||
"Project-device row changed during sync execution."
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
const revision = appendProjectRevision(tx, {
|
||||
projectId: input.projectId,
|
||||
expectedRevision: input.expectedRevision,
|
||||
source: input.source,
|
||||
description: input.description,
|
||||
actorId: input.actorId,
|
||||
forward: input.command,
|
||||
inverse,
|
||||
});
|
||||
applyProjectHistoryTransition(tx, {
|
||||
projectId: input.projectId,
|
||||
source: input.source,
|
||||
recordedChangeSetId: revision.changeSetId,
|
||||
targetChangeSetId: input.historyTargetChangeSetId,
|
||||
});
|
||||
return { revision, inverse };
|
||||
});
|
||||
return inverse;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user