Skip to content
File

Blob: src/cloudflare/internal/d1-api.ts

typescript836 lines
1// Copyright (c) 2023 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5import { withSpan } from 'cloudflare-internal:tracing-helpers';
6import type { Span } from './tracing';
7 
8interface D1Meta {
9 duration: number;
10 size_after: number;
11 rows_read: number;
12 rows_written: number;
13 last_row_id: number;
14 changed_db: boolean;
15 changes: number;
16 
17 /**
18 * The region of the database instance that executed the query.
19 */
20 served_by_region?: string;
21 
22 /**
23 * The three-letter airport code of the colo that executed the query.
24 */
25 served_by_colo?: string;
26 
27 /**
28 * True if-and-only-if the database instance that executed the query was the primary.
29 */
30 served_by_primary?: boolean;
31 
32 timings?: {
33 /**
34 * The duration of the SQL query execution by the database instance. It doesn't include any network time.
35 */
36 sql_duration_ms: number;
37 };
38 
39 /**
40 * Number of total attempts to execute the query, due to automatic retries.
41 * Note: All other fields in the response like `timings` only apply to the last attempt.
42 */
43 total_attempts?: number;
44}
45 
46interface Fetcher {
47 fetch: typeof fetch;
48}
49 
50type D1Response = {
51 success: true;
52 meta: D1Meta & Record<string, unknown>;
53 error?: never;
54};
55 
56type D1Result<T = unknown> = D1Response & {
57 results: T[];
58};
59 
60type D1RawOptions = {
61 columnNames?: boolean;
62};
63 
64type D1UpstreamFailure = {
65 results?: never;
66 error: string;
67 success: false;
68 meta?: never;
69};
70 
71type D1RowsColumns<T = unknown> = D1Response & {
72 results: {
73 columns: string[];
74 rows: T[][];
75 };
76};
77 
78type D1UpstreamSuccess<T = unknown> =
79 | D1Result<T>
80 | D1Response
81 | D1RowsColumns<T>;
82 
83type D1UpstreamResponse<T = unknown> = D1UpstreamSuccess<T> | D1UpstreamFailure;
84 
85type D1ExecResult = {
86 count: number;
87 duration: number;
88};
89 
90type SQLError = {
91 error: string;
92};
93 
94type ResultsFormat = 'ARRAY_OF_OBJECTS' | 'ROWS_AND_COLUMNS' | 'NONE';
95 
96type D1SessionBookmarkOrConstraint = string;
97type D1SessionBookmark = string;
98// Indicates that the first query should go to the primary, and the rest queries
99// using the same D1DatabaseSession will go to any replica that is consistent with
100// the bookmark maintained by the session (returned by the first query).
101const D1_SESSION_CONSTRAINT_FIRST_PRIMARY = 'first-primary';
102// Indicates that the first query can go anywhere (primary or replica), and the rest queries
103// using the same D1DatabaseSession will go to any replica that is consistent with
104// the bookmark maintained by the session (returned by the first query).
105const D1_SESSION_CONSTRAINT_FIRST_UNCONSTRAINED = 'first-unconstrained';
106 
107// Parsed by the D1 eyeball worker.
108// This header is internal only for our D1 binding, not part of the public D1 REST API.
109// Customers should not use this header otherwise their applications can break when we change this.
110// TODO Rename this to `x-cf-d1-session-bookmark`, with coordination with the D1 internal API.
111const D1_SESSION_COMMIT_TOKEN_HTTP_HEADER = 'x-cf-d1-session-commit-token';
112 
113class D1Database {
114 // TODO(soon): Can we use the # syntax here?
115 // eslint-disable-next-line no-restricted-syntax
116 private readonly alwaysPrimarySession: D1DatabaseSessionAlwaysPrimary;
117 protected readonly fetcher: Fetcher;
118 
119 constructor(fetcher: Fetcher) {
120 this.fetcher = fetcher;
121 this.alwaysPrimarySession = new D1DatabaseSessionAlwaysPrimary(
122 this.fetcher
123 );
124 }
125 
126 prepare(query: string): D1PreparedStatement {
127 return new D1PreparedStatement(this.alwaysPrimarySession, query);
128 }
129 
130 async batch<T = unknown>(
131 statements: D1PreparedStatement[]
132 ): Promise<D1Result<T>[]> {
133 return this.alwaysPrimarySession.batch(statements);
134 }
135 
136 async exec(query: string): Promise<D1ExecResult> {
137 return this.alwaysPrimarySession.exec(query);
138 }
139 
140 withSession(
141 constraintOrBookmark?: D1SessionBookmarkOrConstraint
142 ): D1DatabaseSession {
143 constraintOrBookmark = constraintOrBookmark?.trim();
144 if (constraintOrBookmark == null || constraintOrBookmark === '') {
145 constraintOrBookmark = D1_SESSION_CONSTRAINT_FIRST_UNCONSTRAINED;
146 }
147 return new D1DatabaseSession(this.fetcher, constraintOrBookmark);
148 }
149 
150 /**
151 * @deprecated
152 */
153 async dump(): Promise<ArrayBuffer> {
154 return this.alwaysPrimarySession.dump();
155 }
156}
157 
158class D1DatabaseSession {
159 protected fetcher: Fetcher;
160 protected bookmarkOrConstraint: D1SessionBookmarkOrConstraint;
161 
162 constructor(
163 fetcher: Fetcher,
164 bookmarkOrConstraint: D1SessionBookmarkOrConstraint
165 ) {
166 this.fetcher = fetcher;
167 this.bookmarkOrConstraint = bookmarkOrConstraint;
168 
169 if (!this.bookmarkOrConstraint) {
170 throw new Error('D1_SESSION_ERROR: invalid bookmark or constraint');
171 }
172 }
173 
174 // Update the bookmark IFF the given newBookmark is more recent.
175 // The bookmark held in the session should always be the latest value we
176 // have observed in the responses to our API. There can be cases where we have concurrent
177 // queries running within the same session, and therefore here we ensure we only
178 // retain the latest bookmark received.
179 // @returns the final bookmark after the update.
180 protected _updateBookmark(
181 newBookmark: D1SessionBookmark
182 ): D1SessionBookmark | null {
183 newBookmark = newBookmark.trim();
184 if (!newBookmark) {
185 // We should not be receiving invalid bookmarks, but just be defensive.
186 return this.getBookmark();
187 }
188 const currentBookmark = this.getBookmark();
189 if (
190 currentBookmark === null ||
191 currentBookmark.localeCompare(newBookmark) < 0
192 ) {
193 this.bookmarkOrConstraint = newBookmark;
194 }
195 return this.getBookmark();
196 }
197 
198 prepare(sql: string): D1PreparedStatement {
199 return new D1PreparedStatement(this, sql);
200 }
201 
202 async batch<T = unknown>(
203 statements: D1PreparedStatement[]
204 ): Promise<D1Result<T>[]> {
205 return withSpan('d1_batch', async (span) => {
206 span.setAttribute('db.system.name', 'cloudflare-d1');
207 span.setAttribute('db.operation.name', 'batch');
208 span.setAttribute(
209 'db.query.text',
210 statements.map((s: D1PreparedStatement) => s.statement).join('\n')
211 );
212 span.setAttribute('db.operation.batch.size', statements.length);
213 span.setAttribute('cloudflare.binding.type', 'D1');
214 span.setAttribute(
215 'cloudflare.d1.query.bookmark',
216 this.getBookmark() ?? undefined
217 );
218 
219 const exec = (await this._sendOrThrow(
220 '/query',
221 statements.map((s: D1PreparedStatement) => s.statement),
222 statements.map((s: D1PreparedStatement) => s.params),
223 'ROWS_AND_COLUMNS',
224 span
225 )) as D1UpstreamSuccess<T>[];
226 
227 span.setAttribute(
228 'cloudflare.d1.response.bookmark',
229 this.getBookmark() ?? undefined
230 );
231 addAggregatedD1MetaToSpan(
232 span,
233 exec.map((e) => e.meta)
234 );
235 
236 return exec.map(toArrayOfObjects);
237 });
238 }
239 
240 // Returns the latest bookmark we received from all responses processed so far.
241 // It does not return constraints that might have be passed during the session creation.
242 getBookmark(): D1SessionBookmark | null {
243 switch (this.bookmarkOrConstraint) {
244 // First to any replica, and then anywhere that satisfies the bookmark.
245 case D1_SESSION_CONSTRAINT_FIRST_UNCONSTRAINED:
246 return null;
247 // First to primary, and then anywhere that satisfies the bookmark.
248 case D1_SESSION_CONSTRAINT_FIRST_PRIMARY:
249 return null;
250 default:
251 return this.bookmarkOrConstraint;
252 }
253 }
254 
255 // fetch will append the bookmark header to all outgoing fetch calls.
256 // The response headers are parsed automatically, extracting the bookmark
257 // from the response headers and updating it through `_updateBookmark(token)`.
258 protected async _wrappedFetch(
259 input: RequestInfo | URL,
260 init?: RequestInit
261 ): Promise<Response> {
262 const h = new Headers(init?.headers);
263 
264 // We send either a constraint, or a bookmark, and the eyeball worker will figure out
265 // what to do based on the value. This simulates the same flow as the REST API would behave too.
266 if (this.bookmarkOrConstraint) {
267 h.set(D1_SESSION_COMMIT_TOKEN_HTTP_HEADER, this.bookmarkOrConstraint);
268 }
269 
270 if (!init) {
271 init = { headers: h };
272 } else {
273 init.headers = h;
274 }
275 return this.fetcher.fetch(input, init).then((resp) => {
276 const newBookmark = resp.headers.get(D1_SESSION_COMMIT_TOKEN_HTTP_HEADER);
277 if (newBookmark) {
278 this._updateBookmark(newBookmark);
279 }
280 return resp;
281 });
282 }
283 
284 async _sendOrThrow<T = unknown>(
285 endpoint: string,
286 query: string | string[],
287 params: unknown[],
288 resultsFormat: ResultsFormat,
289 span: Span
290 ): Promise<D1UpstreamSuccess<T>[] | D1UpstreamSuccess<T>> {
291 const results = await this._send(
292 endpoint,
293 query,
294 params,
295 resultsFormat,
296 span
297 );
298 const firstResult = firstIfArray(results);
299 if (!firstResult.success) {
300 span.setAttribute('error.type', firstResult.error);
301 throw new Error(`D1_ERROR: ${firstResult.error}`, {
302 cause: new Error(firstResult.error),
303 });
304 } else {
305 return results as D1UpstreamSuccess<T>[] | D1UpstreamSuccess<T>;
306 }
307 }
308 
309 async _send<T = unknown>(
310 endpoint: string,
311 query: string | string[],
312 params: unknown[],
313 resultsFormat: ResultsFormat,
314 span: Span
315 ): Promise<D1UpstreamResponse<T>[] | D1UpstreamResponse<T>> {
316 /* this needs work - we currently only support ordered ?n params */
317 const body = JSON.stringify(
318 Array.isArray(query)
319 ? query.map((s: string, index: number) => {
320 return { sql: s, params: params[index] };
321 })
322 : {
323 sql: query,
324 params: params,
325 }
326 );
327 
328 const url = new URL(endpoint, 'http://d1');
329 url.searchParams.set('resultsFormat', resultsFormat);
330 const response = await this._wrappedFetch(url.href, {
331 method: 'POST',
332 headers: {
333 'content-type': 'application/json',
334 },
335 body,
336 });
337 
338 try {
339 const answer = await toJson<
340 D1UpstreamResponse<T>[] | D1UpstreamResponse<T>
341 >(response);
342 
343 if (Array.isArray(answer)) {
344 return answer.map((r: D1UpstreamResponse<T>) => mapD1Result<T>(r));
345 } else {
346 return mapD1Result<T>(answer);
347 }
348 } catch (_e: unknown) {
349 const e = _e as Error;
350 const message =
351 (e.cause as Error | undefined)?.message ||
352 e.message ||
353 'Something went wrong';
354 span.setAttribute('error.type', message);
355 throw new Error(`D1_ERROR: ${message}`, {
356 cause: new Error(message),
357 });
358 }
359 }
360}
361 
362class D1DatabaseSessionAlwaysPrimary extends D1DatabaseSession {
363 constructor(fetcher: Fetcher) {
364 // Will always go to primary, since we won't be ever updating this constraint.
365 super(fetcher, D1_SESSION_CONSTRAINT_FIRST_PRIMARY);
366 }
367 
368 // We ignore bookmarks for this special type of session,
369 // since all queries are sent to the primary.
370 override _updateBookmark(
371 _newBookmark: D1SessionBookmark
372 ): D1SessionBookmark | null {
373 return null;
374 }
375 
376 // There is no bookmark returned ever by this special type of session,
377 // since all queries are sent to the primary.
378 override getBookmark(): D1SessionBookmark | null {
379 return null;
380 }
381 
382 //////////////////////////////////////////////////////////////////////////////////////////////
383 // These are only used by the D1Database which is our existing API pre-Sessions API.
384 // For backwards compatibility they always go to the primary database.
385 //
386 
387 async exec(query: string): Promise<D1ExecResult> {
388 return withSpan('d1_exec', async (span) => {
389 span.setAttribute('db.system.name', 'cloudflare-d1');
390 span.setAttribute('db.operation.name', 'exec');
391 span.setAttribute('db.query.text', query);
392 span.setAttribute('cloudflare.binding.type', 'D1');
393 
394 // TODO: splitting by lines is overly simplification because a single line
395 // can contain multiple statements (ex: `select 1; select 2;`).
396 // Also, a statement can span multiple lines...
397 // Either, we should do a more reasonable job to split the query into multiple statements
398 // like we do in the D1 codebase or we report a simpler error without the line number.
399 const lines = query.trim().split('\n');
400 const _exec = await this._send('/execute', lines, [], 'NONE', span);
401 const exec = Array.isArray(_exec) ? _exec : [_exec];
402 
403 let duration = 0;
404 const metas: D1Meta[] = [];
405 for (let i = 0; i < exec.length; i++) {
406 const res = exec[i];
407 if (!res?.success) {
408 span.setAttribute('error.type', `Error in line ${i + 1}`);
409 throw new Error(
410 `D1_EXEC_ERROR: Error in line ${i + 1}: ${lines[i]}${res?.error ? `: ${res.error}` : ''}`,
411 {
412 cause: new Error(
413 `Error in line ${i + 1}: ${lines[i]}${res?.error ? `: ${res.error}` : ''}`
414 ),
415 }
416 );
417 }
418 
419 duration += res.meta.duration;
420 metas.push(res.meta);
421 }
422 
423 if (metas.length) {
424 addAggregatedD1MetaToSpan(span, metas);
425 }
426 return {
427 count: exec.length,
428 duration,
429 };
430 });
431 }
432 
433 /**
434 * DEPRECATED, TO BE REMOVED WITH NEXT BREAKING CHANGE
435 * Only applies to the deprecated v1 alpha databases.
436 */
437 async dump(): Promise<ArrayBuffer> {
438 const response = await this._wrappedFetch('http://d1/dump', {
439 method: 'POST',
440 headers: {
441 'content-type': 'application/json',
442 },
443 });
444 if (response.status !== 200) {
445 try {
446 const err = (await response.json()) as SQLError;
447 throw new Error(`D1_DUMP_ERROR: ${err.error}`, {
448 cause: new Error(err.error),
449 });
450 } catch {
451 throw new Error(`D1_DUMP_ERROR: Status + ${response.status}`, {
452 cause: new Error(`Status ${response.status}`),
453 });
454 }
455 }
456 return await response.arrayBuffer();
457 }
458}
459 
460class D1PreparedStatement {
461 // TODO(soon): Can we use the # syntax here?
462 // eslint-disable-next-line no-restricted-syntax
463 private readonly dbSession: D1DatabaseSession;
464 readonly statement: string;
465 readonly params: unknown[];
466 
467 constructor(
468 dbSession: D1DatabaseSession,
469 statement: string,
470 values?: unknown[]
471 ) {
472 this.dbSession = dbSession;
473 this.statement = statement;
474 this.params = values || [];
475 }
476 
477 bind(...values: unknown[]): D1PreparedStatement {
478 // Validate value types
479 const transformedValues = values.map((r: unknown): unknown => {
480 const rType = typeof r;
481 if (rType === 'number' || rType === 'string') {
482 return r;
483 } else if (rType === 'boolean') {
484 return r ? 1 : 0;
485 } else if (rType === 'object') {
486 // nulls are objects in javascript
487 if (r == null) return r;
488 // arrays with uint8's are good
489 if (
490 Array.isArray(r) &&
491 r.every((b: unknown) => {
492 return typeof b == 'number' && b >= 0 && b < 256;
493 })
494 )
495 return r as unknown[];
496 // convert ArrayBuffer to array
497 if (r instanceof ArrayBuffer) {
498 return Array.from(new Uint8Array(r));
499 }
500 // convert view to array
501 if (ArrayBuffer.isView(r)) {
502 // For some reason TS doesn't think this is valid, but it is!
503 return Array.from(r as unknown as ArrayLike<unknown>);
504 }
505 }
506 
507 throw new Error(
508 `D1_TYPE_ERROR: Type '${rType}' not supported for value '${r}'`,
509 {
510 cause: new Error(`Type '${rType}' not supported for value '${r}'`),
511 }
512 );
513 });
514 return new D1PreparedStatement(
515 this.dbSession,
516 this.statement,
517 transformedValues
518 );
519 }
520 
521 async first<T = unknown>(colName: string): Promise<T | null>;
522 async first<T = Record<string, unknown>>(): Promise<T | null>;
523 async first<T = unknown>(
524 colName?: string
525 ): Promise<Record<string, T> | T | null> {
526 return withSpan('d1_first', async (span) => {
527 span.setAttribute('db.system.name', 'cloudflare-d1');
528 span.setAttribute('db.operation.name', 'first');
529 span.setAttribute('db.query.text', this.statement);
530 span.setAttribute('cloudflare.binding.type', 'D1');
531 span.setAttribute(
532 'cloudflare.d1.query.bookmark',
533 this.dbSession.getBookmark() ?? undefined
534 );
535 
536 const info = firstIfArray(
537 await this.dbSession._sendOrThrow<Record<string, T>>(
538 '/query',
539 this.statement,
540 this.params,
541 'ROWS_AND_COLUMNS',
542 span
543 )
544 );
545 
546 span.setAttribute(
547 'cloudflare.d1.response.bookmark',
548 this.dbSession.getBookmark() ?? undefined
549 );
550 addD1MetaToSpan(span, info.meta);
551 
552 const results = toArrayOfObjects(info).results;
553 const hasResults = results.length > 0;
554 if (!hasResults) return null;
555 
556 const firstResult = results.at(0);
557 if (colName !== undefined) {
558 if (firstResult?.[colName] === undefined) {
559 span.setAttribute('error.type', 'Column not found');
560 throw new Error(`D1_COLUMN_NOTFOUND: Column not found (${colName})`, {
561 cause: new Error('Column not found'),
562 });
563 }
564 return firstResult[colName];
565 } else {
566 return firstResult as Record<string, T>;
567 }
568 });
569 }
570 
571 /* eslint-disable-next-line @typescript-eslint/no-unnecessary-type-parameters */
572 async run<T = Record<string, unknown>>(): Promise<D1Response> {
573 return withSpan('d1_run', async (span) => {
574 span.setAttribute('db.system.name', 'cloudflare-d1');
575 span.setAttribute('db.operation.name', 'run');
576 span.setAttribute('db.query.text', this.statement);
577 span.setAttribute('cloudflare.binding.type', 'D1');
578 span.setAttribute(
579 'cloudflare.d1.query.bookmark',
580 this.dbSession.getBookmark() ?? undefined
581 );
582 
583 const result = firstIfArray(
584 await this.dbSession._sendOrThrow<T>(
585 '/execute',
586 this.statement,
587 this.params,
588 'NONE',
589 span
590 )
591 );
592 
593 span.setAttribute(
594 'cloudflare.d1.response.bookmark',
595 this.dbSession.getBookmark() ?? undefined
596 );
597 addD1MetaToSpan(span, result.meta);
598 return result;
599 });
600 }
601 
602 async all<T = Record<string, unknown>>(): Promise<D1Result<T[]>> {
603 return withSpan('d1_all', async (span) => {
604 span.setAttribute('db.system.name', 'cloudflare-d1');
605 span.setAttribute('db.operation.name', 'all');
606 span.setAttribute('db.query.text', this.statement);
607 span.setAttribute('cloudflare.binding.type', 'D1');
608 span.setAttribute(
609 'cloudflare.d1.query.bookmark',
610 this.dbSession.getBookmark() ?? undefined
611 );
612 
613 const result = firstIfArray(
614 await this.dbSession._sendOrThrow<T[]>(
615 '/query',
616 this.statement,
617 this.params,
618 'ROWS_AND_COLUMNS',
619 span
620 )
621 );
622 
623 span.setAttribute(
624 'cloudflare.d1.response.bookmark',
625 this.dbSession.getBookmark() ?? undefined
626 );
627 addD1MetaToSpan(span, result.meta);
628 
629 return toArrayOfObjects(result);
630 });
631 }
632 
633 async raw<T = unknown[]>(options?: D1RawOptions): Promise<T[]> {
634 return withSpan('d1_all', async (span) => {
635 span.setAttribute('db.system.name', 'cloudflare-d1');
636 span.setAttribute('db.operation.name', 'raw');
637 span.setAttribute('db.query.text', this.statement);
638 span.setAttribute('cloudflare.binding.type', 'D1');
639 span.setAttribute(
640 'cloudflare.d1.query.bookmark',
641 this.dbSession.getBookmark() ?? undefined
642 );
643 
644 const s = firstIfArray(
645 await this.dbSession._sendOrThrow<Record<string, unknown>>(
646 '/query',
647 this.statement,
648 this.params,
649 'ROWS_AND_COLUMNS',
650 span
651 )
652 );
653 
654 span.setAttribute(
655 'cloudflare.d1.response.bookmark',
656 this.dbSession.getBookmark() ?? undefined
657 );
658 addD1MetaToSpan(span, s.meta);
659 
660 // If no results returned, return empty array
661 if (!('results' in s)) return [];
662 
663 // If ARRAY_OF_OBJECTS returned, extract cells
664 if (Array.isArray(s.results)) {
665 const raw: T[] = [];
666 for (const row of s.results) {
667 if (options?.columnNames && raw.length === 0) {
668 raw.push(Array.from(Object.keys(row)) as T);
669 }
670 const entry = Object.keys(row).map((k) => {
671 return row[k];
672 });
673 raw.push(entry as T);
674 }
675 return raw;
676 } else {
677 // Otherwise, data is already in the correct format
678 return [
679 ...(options?.columnNames ? [s.results.columns as T] : []),
680 ...(s.results.rows as T[]),
681 ];
682 }
683 });
684 }
685}
686 
687function firstIfArray<T>(results: T | T[]): T {
688 return Array.isArray(results) ? (results.at(0) as T) : results;
689}
690 
691// This shim may be used against an older version of D1 that doesn't support
692// the ROWS_AND_COLUMNS/NONE interchange format, so be permissive here
693function toArrayOfObjects<T>(response: D1UpstreamSuccess<T>): D1Result<T> {
694 // If 'results' is missing from upstream, add an empty array
695 if (!('results' in response))
696 return {
697 ...response,
698 results: [],
699 };
700 
701 const results = response.results;
702 if (Array.isArray(results)) {
703 return { ...response, results };
704 } else {
705 const { rows, columns } = results;
706 return {
707 ...response,
708 results: rows.map(
709 (row) =>
710 Object.fromEntries(row.map((cell, i) => [columns[i], cell])) as T
711 ),
712 };
713 }
714}
715 
716function mapD1Result<T>(result: D1UpstreamResponse<T>): D1UpstreamResponse<T> {
717 // The rest of the app can guarantee that success is true/false, but from the API
718 // we only guarantee that error is present/absent.
719 return result.error
720 ? {
721 success: false,
722 error: result.error,
723 }
724 : {
725 success: true,
726 meta: (result as D1UpstreamSuccess).meta,
727 ...('results' in result ? { results: result.results } : {}),
728 };
729}
730 
731async function toJson<T = unknown>(response: Response): Promise<T> {
732 const body = await response.text();
733 try {
734 return JSON.parse(body) as T;
735 } catch {
736 throw new Error(`Failed to parse body as JSON, got: ${body}`);
737 }
738}
739 
740type PartialD1Meta = Partial<D1Meta> | undefined;
741 
742function addAggregatedD1MetaToSpan(span: Span, metas: PartialD1Meta[]): void {
743 if (!metas.length) {
744 return;
745 }
746 const aggregatedMeta = aggregateD1Meta(metas);
747 addD1MetaToSpan(span, aggregatedMeta);
748}
749 
750function addD1MetaToSpan(span: Span, meta: D1Meta): void {
751 span.setAttribute('cloudflare.d1.response.size_after', meta.size_after);
752 span.setAttribute('cloudflare.d1.response.rows_read', meta.rows_read);
753 span.setAttribute('cloudflare.d1.response.rows_written', meta.rows_written);
754 span.setAttribute('cloudflare.d1.response.last_row_id', meta.last_row_id);
755 span.setAttribute('cloudflare.d1.response.changed_db', meta.changed_db);
756 span.setAttribute('cloudflare.d1.response.changes', meta.changes);
757 span.setAttribute(
758 'cloudflare.d1.response.served_by_region',
759 meta.served_by_region
760 );
761 span.setAttribute(
762 'cloudflare.d1.response.served_by_colo',
763 meta.served_by_colo
764 );
765 span.setAttribute(
766 'cloudflare.d1.response.served_by_primary',
767 meta.served_by_primary
768 );
769 span.setAttribute(
770 'cloudflare.d1.response.sql_duration_ms',
771 meta.timings?.sql_duration_ms ?? undefined
772 );
773 span.setAttribute(
774 'cloudflare.d1.response.total_attempts',
775 meta.total_attempts
776 );
777}
778 
779// When a query is executing multiple statements, and we receive a D1Meta
780// for each statement, we need to aggregate the meta data before we annotate
781// the telemetry, with different rules for each field.
782function aggregateD1Meta(metas: PartialD1Meta[]): D1Meta {
783 const aggregatedMeta: D1Meta = {
784 duration: 0,
785 size_after: 0,
786 rows_read: 0,
787 rows_written: 0,
788 last_row_id: 0,
789 changed_db: false,
790 changes: 0,
791 };
792 
793 for (const meta of metas) {
794 if (!meta) {
795 continue;
796 }
797 
798 aggregatedMeta.duration += meta.duration ?? 0;
799 // for size_after, we only want the last value
800 aggregatedMeta.size_after = meta.size_after ?? 0;
801 aggregatedMeta.rows_read += meta.rows_read ?? 0;
802 aggregatedMeta.rows_written += meta.rows_written ?? 0;
803 aggregatedMeta.last_row_id = meta.last_row_id ?? 0;
804 if (meta.served_by_region) {
805 aggregatedMeta.served_by_region = meta.served_by_region;
806 }
807 if (meta.served_by_colo) {
808 aggregatedMeta.served_by_colo = meta.served_by_colo;
809 }
810 if (meta.served_by_primary) {
811 aggregatedMeta.served_by_primary = meta.served_by_primary;
812 }
813 if (meta.timings?.sql_duration_ms) {
814 aggregatedMeta.timings = {
815 sql_duration_ms:
816 (aggregatedMeta.timings?.sql_duration_ms ?? 0) +
817 meta.timings.sql_duration_ms,
818 };
819 }
820 if (meta.total_attempts) {
821 aggregatedMeta.total_attempts =
822 (aggregatedMeta.total_attempts ?? 0) + meta.total_attempts;
823 }
824 aggregatedMeta.changes += meta.changes ?? 0;
825 if (meta.changed_db) {
826 aggregatedMeta.changed_db = true;
827 }
828 }
829 
830 return aggregatedMeta;
831}
832 
833export default function makeBinding(env: { fetcher: Fetcher }): D1Database {
834 return new D1Database(env.fetcher);
835}