fix copy race

This commit is contained in:
mertalev
2026-02-14 03:32:59 -05:00
parent 6b04fa3f94
commit 839fb61340
+17 -13
View File
@@ -394,20 +394,24 @@ export class WriteBuffer {
} }
private async copyInsert(tableName: string, rows: InsertRow[]) { private async copyInsert(tableName: string, rows: InsertRow[]) {
const writable = await this const conn = await this.pgPool.reserve();
.pgPool`COPY ${this.pgPool(tableName)} (code, data, priority, "runAfter") FROM STDIN WITH (FORMAT csv)`.writable(); try {
const now = new Date().toISOString(); const writable = await conn`COPY ${conn(tableName)} (code, data, priority, "runAfter") FROM STDIN WITH (FORMAT csv)`.writable();
for (const row of rows) { const now = new Date().toISOString();
const data = row.data != null ? csvEscape(JSON.stringify(row.data)) : ''; for (const row of rows) {
const priority = row.priority ?? 0; const data = row.data != null ? csvEscape(JSON.stringify(row.data)) : '';
const runAfter = row.runAfter ? row.runAfter.toISOString() : now; const priority = row.priority ?? 0;
writable.write(`${row.code},${data},${priority},${runAfter}\n`); const runAfter = row.runAfter ? row.runAfter.toISOString() : now;
writable.write(`${row.code},${data},${priority},${runAfter}\n`);
}
writable.end();
await new Promise<void>((resolve, reject) => {
writable.on('finish', resolve);
writable.on('error', reject);
});
} finally {
conn.release();
} }
writable.end();
await new Promise<void>((resolve, reject) => {
writable.on('finish', resolve);
writable.on('error', reject);
});
} }
private insertChunk(tableName: string, rows: InsertRow[]) { private insertChunk(tableName: string, rows: InsertRow[]) {