( conn: Mysql.PoolConnection, sql: string, params?: ReadonlyArray<any> )
| 473 | return ["", []] |
| 474 | }, |
| 475 | onRecordUpdate() { |
| 476 | return ["", []] |
| 477 | } |
| 478 | }) |
| 479 | |
| 480 | const escape = Statement.defaultEscape("`") |
| 481 | |
| 482 | function queryStream( |
| 483 | conn: Mysql.PoolConnection, |
| 484 | sql: string, |
| 485 | params?: ReadonlyArray<any> |
| 486 | ) { |
| 487 | return asyncPauseResume<any, SqlError>(Effect.fnUntraced(function*(emit) { |
| 488 | const query = (conn as any).query(sql, params).stream() |
| 489 | yield* Effect.addFinalizer(() => Effect.sync(() => query.destroy() as void)) |
| 490 | |
| 491 | let buffer: Array<any> = [] |
| 492 | let taskPending = false |
| 493 | query.on( |
| 494 | "error", |
| 495 | (cause: unknown) => |
| 496 | emit.fail(new SqlError({ reason: classifyError(cause, "Failed to stream statement", "stream") })) |
| 497 | ) |
| 498 | query.on("data", (row: any) => { |
| 499 | buffer.push(row) |
| 500 | if (!taskPending) { |
| 501 | taskPending = true |
| 502 | queueMicrotask(() => { |
| 503 | const items = buffer |
| 504 | buffer = [] |
| 505 | emit.array(items) |
| 506 | taskPending = false |
| 507 | }) |
| 508 | } |
| 509 | }) |
| 510 | query.on("end", () => { |
| 511 | if (buffer.length > 0) { |
| 512 | emit.array(buffer) |
| 513 | buffer = [] |
| 514 | } |
| 515 | emit.end() |
| 516 | }) |
| 517 | |
| 518 | return { |
| 519 | onPause() { |
| 520 | query.pause() |
| 521 | }, |
| 522 | onResume() { |
| 523 | query.resume() |
no test coverage detected
searching dependent graphs…