From 43bcc41981d71f67af9d7d01780ef745a1ebdf14 Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Tue, 20 Aug 2024 10:02:50 -0400 Subject: [PATCH 01/13] Move the constructor over to TabularApiService --- src/table.ts | 30 ++---------------------------- src/tabular-api-service.ts | 33 +++++++++++++++++++++++++++++++++ 2 files changed, 35 insertions(+), 28 deletions(-) diff --git a/src/table.ts b/src/table.ts index e7c286c0f..de3b72ff5 100644 --- a/src/table.ts +++ b/src/table.ts @@ -43,6 +43,7 @@ import {CreateBackupCallback, CreateBackupResponse} from './cluster'; import {google} from '../protos/protos'; import {Duplex} from 'stream'; import {TableUtils} from './utils/table'; +import {TabularApiService} from './tabular-api-service'; // See protos/google/rpc/code.proto // (4=DEADLINE_EXCEEDED, 8=RESOURCE_EXHAUSTED, 10=ABORTED, 14=UNAVAILABLE) @@ -396,34 +397,7 @@ export interface CreateBackupConfig extends ModifiableBackupFields { * const table = instance.table('prezzy'); * ``` */ -export class Table { - bigtable: Bigtable; - instance: Instance; - name: string; - id: string; - metadata?: google.bigtable.admin.v2.ITable; - maxRetries?: number; - constructor(instance: Instance, id: string) { - this.bigtable = instance.bigtable; - this.instance = instance; - - let name; - - if (id.includes('/')) { - if (id.startsWith(`${instance.name}/tables/`)) { - name = id; - } else { - throw new Error(`Table id '${id}' is not formatted correctly. -Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); - } - } else { - name = `${instance.name}/tables/${id}`; - } - - this.name = name; - this.id = name.split('/').pop()!; - } - +export class Table extends TabularApiService { /** * Formats the decodes policy etag value to string. * diff --git a/src/tabular-api-service.ts b/src/tabular-api-service.ts index e69de29bb..b23c03202 100644 --- a/src/tabular-api-service.ts +++ b/src/tabular-api-service.ts @@ -0,0 +1,33 @@ +import {Instance} from './instance'; +import {Bigtable} from './index'; +import {google} from '../protos/protos'; + +export class TabularApiService { + bigtable: Bigtable; + instance: Instance; + name: string; + id: string; + metadata?: google.bigtable.admin.v2.ITable; + maxRetries?: number; + + constructor(instance: Instance, id: string) { + this.bigtable = instance.bigtable; + this.instance = instance; + + let name; + + if (id.includes('/')) { + if (id.startsWith(`${instance.name}/tables/`)) { + name = id; + } else { + throw new Error(`Table id '${id}' is not formatted correctly. +Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); + } + } else { + name = `${instance.name}/tables/${id}`; + } + + this.name = name; + this.id = name.split('/').pop()!; + } +} From a76219e7c7bda771d7a07b2ba52e27e9f8b5bde7 Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Tue, 20 Aug 2024 10:36:53 -0400 Subject: [PATCH 02/13] Move sampleRowKeys over --- src/table.ts | 88 ---------------------------------- src/tabular-api-service.ts | 97 +++++++++++++++++++++++++++++++++++++- 2 files changed, 96 insertions(+), 89 deletions(-) diff --git a/src/table.ts b/src/table.ts index de3b72ff5..5556ff677 100644 --- a/src/table.ts +++ b/src/table.ts @@ -1653,94 +1653,6 @@ export class Table extends TabularApiService { return new Row(this, key); } - sampleRowKeys(gaxOptions?: CallOptions): Promise; - sampleRowKeys(gaxOptions: CallOptions, callback: SampleRowKeysCallback): void; - sampleRowKeys(callback?: SampleRowKeysCallback): void; - /** - * Returns a sample of row keys in the table. The returned row keys will delimit - * contiguous sections of the table of approximately equal size, which can be - * used to break up the data for distributed tasks like mapreduces. - * - * @param {object} [gaxOptions] Request configuration options, outlined here: - * https://googleapis.github.io/gax-nodejs/CallSettings.html. - * @param {function} [callback] The callback function. - * @param {?error} callback.err An error returned while making this request. - * @param {object[]} callback.keys The list of keys. - * - * @example include:samples/api-reference-doc-snippets/table.js - * region_tag:bigtable_api_sample_row_keys - */ - sampleRowKeys( - optionsOrCallback?: CallOptions | SampleRowKeysCallback, - cb?: SampleRowKeysCallback - ): void | Promise { - const callback = - typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; - const gaxOptions = - typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; - this.sampleRowKeysStream(gaxOptions) - .on('error', callback) - .pipe( - concat((keys: string[]) => { - callback(null, keys); - }) - ); - } - - /** - * Returns a sample of row keys in the table as a readable object stream. - * - * See {@link Table#sampleRowKeys} for more details. - * - * @param {object} [gaxOptions] Request configuration options, outlined here: - * https://googleapis.github.io/gax-nodejs/CallSettings.html. - * @returns {stream} - * - * @example - * ``` - * table.sampleRowKeysStream() - * .on('error', console.error) - * .on('data', function(key) { - * // Do something with the `key` object. - * }); - * - * //- - * // If you anticipate many results, you can end a stream early to prevent - * // unnecessary processing. - * //- - * table.sampleRowKeysStream() - * .on('data', function(key) { - * this.end(); - * }); - * ``` - */ - sampleRowKeysStream(gaxOptions?: CallOptions) { - const reqOpts = { - tableName: this.name, - appProfileId: this.bigtable.appProfileId, - }; - - const rowKeysStream = new Transform({ - transform(key, enc, next) { - next(null, { - key: key.rowKey, - offset: key.offsetBytes, - }); - }, - objectMode: true, - }); - - return pumpify.obj([ - this.bigtable.request({ - client: 'BigtableClient', - method: 'sampleRowKeys', - reqOpts, - gaxOpts: Object.assign({}, gaxOptions), - }), - rowKeysStream, - ]); - } - setIamPolicy( policy: Policy, gaxOptions?: CallOptions diff --git a/src/tabular-api-service.ts b/src/tabular-api-service.ts index b23c03202..765bad1cd 100644 --- a/src/tabular-api-service.ts +++ b/src/tabular-api-service.ts @@ -1,6 +1,13 @@ import {Instance} from './instance'; -import {Bigtable} from './index'; +import {Bigtable, SampleRowKeysCallback, SampleRowsKeysResponse} from './index'; import {google} from '../protos/protos'; +import {CallOptions} from 'google-gax'; +import {Transform} from 'stream'; + +// eslint-disable-next-line @typescript-eslint/no-var-requires +const concat = require('concat-stream'); +// eslint-disable-next-line @typescript-eslint/no-var-requires +const pumpify = require('pumpify'); export class TabularApiService { bigtable: Bigtable; @@ -30,4 +37,92 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); this.name = name; this.id = name.split('/').pop()!; } + + sampleRowKeys(gaxOptions?: CallOptions): Promise; + sampleRowKeys(gaxOptions: CallOptions, callback: SampleRowKeysCallback): void; + sampleRowKeys(callback?: SampleRowKeysCallback): void; + /** + * Returns a sample of row keys in the table. The returned row keys will delimit + * contiguous sections of the table of approximately equal size, which can be + * used to break up the data for distributed tasks like mapreduces. + * + * @param {object} [gaxOptions] Request configuration options, outlined here: + * https://googleapis.github.io/gax-nodejs/CallSettings.html. + * @param {function} [callback] The callback function. + * @param {?error} callback.err An error returned while making this request. + * @param {object[]} callback.keys The list of keys. + * + * @example include:samples/api-reference-doc-snippets/table.js + * region_tag:bigtable_api_sample_row_keys + */ + sampleRowKeys( + optionsOrCallback?: CallOptions | SampleRowKeysCallback, + cb?: SampleRowKeysCallback + ): void | Promise { + const callback = + typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; + const gaxOptions = + typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; + this.sampleRowKeysStream(gaxOptions) + .on('error', callback) + .pipe( + concat((keys: string[]) => { + callback(null, keys); + }) + ); + } + + /** + * Returns a sample of row keys in the table as a readable object stream. + * + * See {@link Table#sampleRowKeys} for more details. + * + * @param {object} [gaxOptions] Request configuration options, outlined here: + * https://googleapis.github.io/gax-nodejs/CallSettings.html. + * @returns {stream} + * + * @example + * ``` + * table.sampleRowKeysStream() + * .on('error', console.error) + * .on('data', function(key) { + * // Do something with the `key` object. + * }); + * + * //- + * // If you anticipate many results, you can end a stream early to prevent + * // unnecessary processing. + * //- + * table.sampleRowKeysStream() + * .on('data', function(key) { + * this.end(); + * }); + * ``` + */ + sampleRowKeysStream(gaxOptions?: CallOptions) { + const reqOpts = { + tableName: this.name, + appProfileId: this.bigtable.appProfileId, + }; + + const rowKeysStream = new Transform({ + transform(key, enc, next) { + next(null, { + key: key.rowKey, + offset: key.offsetBytes, + }); + }, + objectMode: true, + }); + + return pumpify.obj([ + this.bigtable.request({ + client: 'BigtableClient', + method: 'sampleRowKeys', + reqOpts, + gaxOpts: Object.assign({}, gaxOptions), + }), + rowKeysStream, + ]); + } } From 0f020ee41a83d9fef1aee18e289e5838d8289a24 Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Tue, 20 Aug 2024 10:57:04 -0400 Subject: [PATCH 03/13] Move sampleRowKeys functions over and use promisify --- src/tabular-api-service.ts | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/src/tabular-api-service.ts b/src/tabular-api-service.ts index 765bad1cd..819f768e5 100644 --- a/src/tabular-api-service.ts +++ b/src/tabular-api-service.ts @@ -1,3 +1,4 @@ +import {promisifyAll} from '@google-cloud/promisify'; import {Instance} from './instance'; import {Bigtable, SampleRowKeysCallback, SampleRowsKeysResponse} from './index'; import {google} from '../protos/protos'; @@ -126,3 +127,12 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); ]); } } + +/*! Developer Documentation + * + * All async methods (except for streams) will return a Promise in the event + * that a callback is omitted. + */ +promisifyAll(TabularApiService, { + exclude: ['family', 'row'], +}); From 4db469b09c758248b264c640304518f4ae00decc Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Tue, 20 Aug 2024 14:33:58 -0400 Subject: [PATCH 04/13] Adjust the proxyquire to work with TabularAPIserv --- src/table.ts | 290 +++---------------------------------- src/tabular-api-service.ts | 289 +++++++++++++++++++++++++++++++++++- test/table.ts | 11 ++ 3 files changed, 318 insertions(+), 272 deletions(-) diff --git a/src/table.ts b/src/table.ts index 5556ff677..fdaf4cbc8 100644 --- a/src/table.ts +++ b/src/table.ts @@ -43,18 +43,26 @@ import {CreateBackupCallback, CreateBackupResponse} from './cluster'; import {google} from '../protos/protos'; import {Duplex} from 'stream'; import {TableUtils} from './utils/table'; -import {TabularApiService} from './tabular-api-service'; - -// See protos/google/rpc/code.proto -// (4=DEADLINE_EXCEEDED, 8=RESOURCE_EXHAUSTED, 10=ABORTED, 14=UNAVAILABLE) -const RETRYABLE_STATUS_CODES = new Set([4, 8, 10, 14]); -// (1=CANCELLED) -const IGNORED_STATUS_CODES = new Set([1]); - -const DEFAULT_BACKOFF_SETTINGS: BackoffSettings = { - initialRetryDelayMillis: 10, - retryDelayMultiplier: 2, - maxRetryDelayMillis: 60000, +import { + DEFAULT_BACKOFF_SETTINGS, + getNextDelay, + IGNORED_STATUS_CODES, + populateAttemptHeader, + RETRYABLE_STATUS_CODES, + TabularApiService, + InsertRowsCallback, + InsertRowsResponse, + MutateCallback, + MutateResponse, + PartialFailureError, +} from './tabular-api-service'; + +export { + InsertRowsCallback, + InsertRowsResponse, + MutateCallback, + MutateResponse, + PartialFailureError, }; /** @@ -362,16 +370,6 @@ export type GetRowsCallback = ( apiResponse?: google.bigtable.v2.ReadRowsResponse ) => void; export type GetRowsResponse = [Row[], google.bigtable.v2.ReadRowsResponse]; -export type InsertRowsCallback = ( - err: ServiceError | PartialFailureError | null, - apiResponse?: google.protobuf.Empty -) => void; -export type InsertRowsResponse = [google.protobuf.Empty]; -export type MutateCallback = ( - err: ServiceError | PartialFailureError | null, - apiResponse?: google.protobuf.Empty -) => void; -export type MutateResponse = [google.protobuf.Empty]; export interface PrefixRange { start?: BoundData | string; @@ -1425,214 +1423,6 @@ export class Table extends TabularApiService { ); } - insert( - entries: Entry | Entry[], - gaxOptions?: CallOptions - ): Promise; - insert( - entries: Entry | Entry[], - gaxOptions: CallOptions, - callback: InsertRowsCallback - ): void; - insert(entries: Entry | Entry[], callback: InsertRowsCallback): void; - /** - * Insert or update rows in your table. It should be noted that gRPC only allows - * you to send payloads that are less than or equal to 4MB. If you're inserting - * more than that you may need to send smaller individual requests. - * - * @param {object|object[]} entries List of entries to be inserted. - * See {@link Table#mutate}. - * @param {object} [gaxOptions] Request configuration options, outlined here: - * https://googleapis.github.io/gax-nodejs/CallSettings.html. - * @param {function} callback The callback function. - * @param {?error} callback.err An error returned while making this request. - * @param {object[]} callback.err.errors If present, these represent partial - * failures. It's possible for part of your request to be completed - * successfully, while the other part was not. - * - * @example include:samples/api-reference-doc-snippets/table.js - * region_tag:bigtable_api_insert_rows - */ - insert( - entries: Entry | Entry[], - optionsOrCallback?: CallOptions | InsertRowsCallback, - cb?: InsertRowsCallback - ): void | Promise { - const callback = - typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; - const gaxOptions = - typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; - entries = arrify(entries).map((entry: Entry) => { - entry.method = Mutation.methods.INSERT; - return entry; - }); - return this.mutate(entries, {gaxOptions}, callback); - } - - mutate( - entries: Entry | Entry[], - options?: MutateOptions - ): Promise; - mutate( - entries: Entry | Entry[], - options: MutateOptions, - callback: MutateCallback - ): void; - mutate(entries: Entry | Entry[], callback: MutateCallback): void; - /** - * Apply a set of changes to be atomically applied to the specified row(s). - * Mutations are applied in order, meaning that earlier mutations can be masked - * by later ones. - * - * @param {object|object[]} entries List of entities to be inserted or - * deleted. - * @param {object} [options] Configuration object. - * @param {object} [options.gaxOptions] Request configuration options, outlined - * here: https://googleapis.github.io/gax-nodejs/global.html#CallOptions. - * @param {boolean} [options.rawMutation] If set to `true` will treat entries - * as a raw Mutation object. See {@link Mutation#parse}. - * @param {function} callback The callback function. - * @param {?error} callback.err An error returned while making this request. - * @param {object[]} callback.err.errors If present, these represent partial - * failures. It's possible for part of your request to be completed - * successfully, while the other part was not. - * - * @example include:samples/api-reference-doc-snippets/table.js - * region_tag:bigtable_api_mutate_rows - */ - mutate( - entriesRaw: Entry | Entry[], - optionsOrCallback?: MutateOptions | MutateCallback, - cb?: MutateCallback - ): void | Promise { - const callback = - typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; - const options = - typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; - const entries: Entry[] = (arrify(entriesRaw) as Entry[]).reduce( - (a, b) => a.concat(b), - [] - ); - - let numRequestsMade = 0; - - const maxRetries = is.number(this.maxRetries) ? this.maxRetries! : 3; - const pendingEntryIndices = new Set( - entries.map((entry: Entry, index: number) => index) - ); - const entryToIndex = new Map( - entries.map((entry: Entry, index: number) => [entry, index]) - ); - const mutationErrorsByEntryIndex = new Map(); - - const isRetryable = (err: ServiceError | null) => { - // Don't retry if there are no more entries or retry attempts - if (pendingEntryIndices.size === 0 || numRequestsMade >= maxRetries + 1) { - return false; - } - // If the error is empty but there are still outstanding mutations, - // it means that there are retryable errors in the mutate response - // even when the RPC succeeded - return !err || RETRYABLE_STATUS_CODES.has(err.code); - }; - - const onBatchResponse = (err: ServiceError | null) => { - // Return if the error happened before a request was made - if (numRequestsMade === 0) { - callback(err); - return; - } - - if (isRetryable(err)) { - const backOffSettings = - options.gaxOptions?.retry?.backoffSettings || - DEFAULT_BACKOFF_SETTINGS; - const nextDelay = getNextDelay(numRequestsMade, backOffSettings); - setTimeout(makeNextBatchRequest, nextDelay); - return; - } - - // If there's no more pending mutations, set the error - // to null - if (pendingEntryIndices.size === 0) { - err = null; - } - - if (mutationErrorsByEntryIndex.size !== 0) { - const mutationErrors = Array.from(mutationErrorsByEntryIndex.values()); - callback(new PartialFailureError(mutationErrors, err)); - return; - } - - callback(err); - }; - - const makeNextBatchRequest = () => { - const entryBatch = entries.filter((entry: Entry, index: number) => { - return pendingEntryIndices.has(index); - }); - - const reqOpts = { - tableName: this.name, - appProfileId: this.bigtable.appProfileId, - entries: options.rawMutation - ? entryBatch - : entryBatch.map(Mutation.parse), - }; - - const retryOpts = { - currentRetryAttempt: numRequestsMade, - // Handling retries in this client. Specify the retry options to - // make sure nothing is retried in retry-request. - noResponseRetries: 0, - shouldRetryFn: (_: any) => { - return false; - }, - }; - - options.gaxOptions = populateAttemptHeader( - numRequestsMade, - options.gaxOptions - ); - - this.bigtable - .request({ - client: 'BigtableClient', - method: 'mutateRows', - reqOpts, - gaxOpts: options.gaxOptions, - retryOpts, - }) - .on('error', (err: ServiceError) => { - onBatchResponse(err); - }) - .on('data', (obj: google.bigtable.v2.IMutateRowsResponse) => { - obj.entries!.forEach(entry => { - const originalEntry = entryBatch[entry.index as number]; - const originalEntriesIndex = entryToIndex.get(originalEntry)!; - - // Mutation was successful. - if (entry.status!.code === 0) { - pendingEntryIndices.delete(originalEntriesIndex); - mutationErrorsByEntryIndex.delete(originalEntriesIndex); - return; - } - if (!RETRYABLE_STATUS_CODES.has(entry.status!.code!)) { - pendingEntryIndices.delete(originalEntriesIndex); - } - const errorDetails = entry.status; - // eslint-disable-next-line @typescript-eslint/no-explicit-any - (errorDetails as any).entry = originalEntry; - mutationErrorsByEntryIndex.set(originalEntriesIndex, errorDetails); - }); - }) - .on('end', onBatchResponse); - numRequestsMade++; - }; - - makeNextBatchRequest(); - } - /** * Get a reference to a table row. * @@ -1957,47 +1747,7 @@ promisifyAll(Table, { exclude: ['family', 'row'], }); -function getNextDelay(numConsecutiveErrors: number, config: BackoffSettings) { - // 0 - 100 ms jitter - const jitter = Math.floor(Math.random() * 100); - const calculatedNextRetryDelay = - config.initialRetryDelayMillis * - Math.pow(config.retryDelayMultiplier, numConsecutiveErrors) + - jitter; - - return Math.min(calculatedNextRetryDelay, config.maxRetryDelayMillis); -} - -function populateAttemptHeader(attempt: number, gaxOpts?: CallOptions) { - gaxOpts = gaxOpts || {}; - gaxOpts.otherArgs = gaxOpts.otherArgs || {}; - gaxOpts.otherArgs.headers = gaxOpts.otherArgs.headers || {}; - gaxOpts.otherArgs.headers['bigtable-attempt'] = attempt; - return gaxOpts; -} - export interface GoogleInnerError { reason?: string; message?: string; } - -export class PartialFailureError extends Error { - errors?: GoogleInnerError[]; - constructor(errors: GoogleInnerError[], rpcError?: ServiceError | null) { - super(); - this.errors = errors; - this.name = 'PartialFailureError'; - let messages = errors.map(e => e.message); - if (messages.length > 1) { - messages = messages.map((message, i) => ` ${i + 1}. ${message}`); - messages.unshift( - 'Multiple errors occurred during the request. Please see the `errors` array for complete details.\n' - ); - messages.push('\n'); - } - this.message = messages.join('\n'); - if (rpcError) { - this.message += 'Request failed with: ' + rpcError.message; - } - } -} diff --git a/src/tabular-api-service.ts b/src/tabular-api-service.ts index 819f768e5..7df826b90 100644 --- a/src/tabular-api-service.ts +++ b/src/tabular-api-service.ts @@ -1,9 +1,43 @@ import {promisifyAll} from '@google-cloud/promisify'; +import arrify = require('arrify'); import {Instance} from './instance'; -import {Bigtable, SampleRowKeysCallback, SampleRowsKeysResponse} from './index'; +import {Mutation} from './mutation'; +import { + Bigtable, + Entry, + MutateOptions, + SampleRowKeysCallback, + SampleRowsKeysResponse, +} from './index'; +import {BackoffSettings} from 'google-gax/build/src/gax'; import {google} from '../protos/protos'; -import {CallOptions} from 'google-gax'; +import {CallOptions, ServiceError} from 'google-gax'; import {Transform} from 'stream'; +import * as is from 'is'; +import {GoogleInnerError} from './table'; + +// See protos/google/rpc/code.proto +// (4=DEADLINE_EXCEEDED, 8=RESOURCE_EXHAUSTED, 10=ABORTED, 14=UNAVAILABLE) +export const RETRYABLE_STATUS_CODES = new Set([4, 8, 10, 14]); +// (1=CANCELLED) +export const IGNORED_STATUS_CODES = new Set([1]); + +export const DEFAULT_BACKOFF_SETTINGS: BackoffSettings = { + initialRetryDelayMillis: 10, + retryDelayMultiplier: 2, + maxRetryDelayMillis: 60000, +}; + +export type InsertRowsCallback = ( + err: ServiceError | PartialFailureError | null, + apiResponse?: google.protobuf.Empty +) => void; +export type InsertRowsResponse = [google.protobuf.Empty]; +export type MutateCallback = ( + err: ServiceError | PartialFailureError | null, + apiResponse?: google.protobuf.Empty +) => void; +export type MutateResponse = [google.protobuf.Empty]; // eslint-disable-next-line @typescript-eslint/no-var-requires const concat = require('concat-stream'); @@ -39,6 +73,214 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); this.id = name.split('/').pop()!; } + insert( + entries: Entry | Entry[], + gaxOptions?: CallOptions + ): Promise; + insert( + entries: Entry | Entry[], + gaxOptions: CallOptions, + callback: InsertRowsCallback + ): void; + insert(entries: Entry | Entry[], callback: InsertRowsCallback): void; + /** + * Insert or update rows in your table. It should be noted that gRPC only allows + * you to send payloads that are less than or equal to 4MB. If you're inserting + * more than that you may need to send smaller individual requests. + * + * @param {object|object[]} entries List of entries to be inserted. + * See {@link Table#mutate}. + * @param {object} [gaxOptions] Request configuration options, outlined here: + * https://googleapis.github.io/gax-nodejs/CallSettings.html. + * @param {function} callback The callback function. + * @param {?error} callback.err An error returned while making this request. + * @param {object[]} callback.err.errors If present, these represent partial + * failures. It's possible for part of your request to be completed + * successfully, while the other part was not. + * + * @example include:samples/api-reference-doc-snippets/table.js + * region_tag:bigtable_api_insert_rows + */ + insert( + entries: Entry | Entry[], + optionsOrCallback?: CallOptions | InsertRowsCallback, + cb?: InsertRowsCallback + ): void | Promise { + const callback = + typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; + const gaxOptions = + typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; + entries = arrify(entries).map((entry: Entry) => { + entry.method = Mutation.methods.INSERT; + return entry; + }); + return this.mutate(entries, {gaxOptions}, callback); + } + + mutate( + entries: Entry | Entry[], + options?: MutateOptions + ): Promise; + mutate( + entries: Entry | Entry[], + options: MutateOptions, + callback: MutateCallback + ): void; + mutate(entries: Entry | Entry[], callback: MutateCallback): void; + /** + * Apply a set of changes to be atomically applied to the specified row(s). + * Mutations are applied in order, meaning that earlier mutations can be masked + * by later ones. + * + * @param {object|object[]} entries List of entities to be inserted or + * deleted. + * @param {object} [options] Configuration object. + * @param {object} [options.gaxOptions] Request configuration options, outlined + * here: https://googleapis.github.io/gax-nodejs/global.html#CallOptions. + * @param {boolean} [options.rawMutation] If set to `true` will treat entries + * as a raw Mutation object. See {@link Mutation#parse}. + * @param {function} callback The callback function. + * @param {?error} callback.err An error returned while making this request. + * @param {object[]} callback.err.errors If present, these represent partial + * failures. It's possible for part of your request to be completed + * successfully, while the other part was not. + * + * @example include:samples/api-reference-doc-snippets/table.js + * region_tag:bigtable_api_mutate_rows + */ + mutate( + entriesRaw: Entry | Entry[], + optionsOrCallback?: MutateOptions | MutateCallback, + cb?: MutateCallback + ): void | Promise { + const callback = + typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; + const options = + typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; + const entries: Entry[] = (arrify(entriesRaw) as Entry[]).reduce( + (a, b) => a.concat(b), + [] + ); + + let numRequestsMade = 0; + + const maxRetries = is.number(this.maxRetries) ? this.maxRetries! : 3; + const pendingEntryIndices = new Set( + entries.map((entry: Entry, index: number) => index) + ); + const entryToIndex = new Map( + entries.map((entry: Entry, index: number) => [entry, index]) + ); + const mutationErrorsByEntryIndex = new Map(); + + const isRetryable = (err: ServiceError | null) => { + // Don't retry if there are no more entries or retry attempts + if (pendingEntryIndices.size === 0 || numRequestsMade >= maxRetries + 1) { + return false; + } + // If the error is empty but there are still outstanding mutations, + // it means that there are retryable errors in the mutate response + // even when the RPC succeeded + return !err || RETRYABLE_STATUS_CODES.has(err.code); + }; + + const onBatchResponse = (err: ServiceError | null) => { + // Return if the error happened before a request was made + if (numRequestsMade === 0) { + callback(err); + return; + } + + if (isRetryable(err)) { + const backOffSettings = + options.gaxOptions?.retry?.backoffSettings || + DEFAULT_BACKOFF_SETTINGS; + const nextDelay = getNextDelay(numRequestsMade, backOffSettings); + setTimeout(makeNextBatchRequest, nextDelay); + return; + } + + // If there's no more pending mutations, set the error + // to null + if (pendingEntryIndices.size === 0) { + err = null; + } + + if (mutationErrorsByEntryIndex.size !== 0) { + const mutationErrors = Array.from(mutationErrorsByEntryIndex.values()); + callback(new PartialFailureError(mutationErrors, err)); + return; + } + + callback(err); + }; + + const makeNextBatchRequest = () => { + const entryBatch = entries.filter((entry: Entry, index: number) => { + return pendingEntryIndices.has(index); + }); + + const reqOpts = { + tableName: this.name, + appProfileId: this.bigtable.appProfileId, + entries: options.rawMutation + ? entryBatch + : entryBatch.map(Mutation.parse), + }; + + const retryOpts = { + currentRetryAttempt: numRequestsMade, + // Handling retries in this client. Specify the retry options to + // make sure nothing is retried in retry-request. + noResponseRetries: 0, + shouldRetryFn: (_: any) => { + return false; + }, + }; + + options.gaxOptions = populateAttemptHeader( + numRequestsMade, + options.gaxOptions + ); + + this.bigtable + .request({ + client: 'BigtableClient', + method: 'mutateRows', + reqOpts, + gaxOpts: options.gaxOptions, + retryOpts, + }) + .on('error', (err: ServiceError) => { + onBatchResponse(err); + }) + .on('data', (obj: google.bigtable.v2.IMutateRowsResponse) => { + obj.entries!.forEach(entry => { + const originalEntry = entryBatch[entry.index as number]; + const originalEntriesIndex = entryToIndex.get(originalEntry)!; + + // Mutation was successful. + if (entry.status!.code === 0) { + pendingEntryIndices.delete(originalEntriesIndex); + mutationErrorsByEntryIndex.delete(originalEntriesIndex); + return; + } + if (!RETRYABLE_STATUS_CODES.has(entry.status!.code!)) { + pendingEntryIndices.delete(originalEntriesIndex); + } + const errorDetails = entry.status; + // eslint-disable-next-line @typescript-eslint/no-explicit-any + (errorDetails as any).entry = originalEntry; + mutationErrorsByEntryIndex.set(originalEntriesIndex, errorDetails); + }); + }) + .on('end', onBatchResponse); + numRequestsMade++; + }; + + makeNextBatchRequest(); + } + sampleRowKeys(gaxOptions?: CallOptions): Promise; sampleRowKeys(gaxOptions: CallOptions, callback: SampleRowKeysCallback): void; sampleRowKeys(callback?: SampleRowKeysCallback): void; @@ -128,6 +370,28 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); } } +export function getNextDelay( + numConsecutiveErrors: number, + config: BackoffSettings +) { + // 0 - 100 ms jitter + const jitter = Math.floor(Math.random() * 100); + const calculatedNextRetryDelay = + config.initialRetryDelayMillis * + Math.pow(config.retryDelayMultiplier, numConsecutiveErrors) + + jitter; + + return Math.min(calculatedNextRetryDelay, config.maxRetryDelayMillis); +} + +export function populateAttemptHeader(attempt: number, gaxOpts?: CallOptions) { + gaxOpts = gaxOpts || {}; + gaxOpts.otherArgs = gaxOpts.otherArgs || {}; + gaxOpts.otherArgs.headers = gaxOpts.otherArgs.headers || {}; + gaxOpts.otherArgs.headers['bigtable-attempt'] = attempt; + return gaxOpts; +} + /*! Developer Documentation * * All async methods (except for streams) will return a Promise in the event @@ -136,3 +400,24 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); promisifyAll(TabularApiService, { exclude: ['family', 'row'], }); + +export class PartialFailureError extends Error { + errors?: GoogleInnerError[]; + constructor(errors: GoogleInnerError[], rpcError?: ServiceError | null) { + super(); + this.errors = errors; + this.name = 'PartialFailureError'; + let messages = errors.map(e => e.message); + if (messages.length > 1) { + messages = messages.map((message, i) => ` ${i + 1}. ${message}`); + messages.unshift( + 'Multiple errors occurred during the request. Please see the `errors` array for complete details.\n' + ); + messages.push('\n'); + } + this.message = messages.join('\n'); + if (rpcError) { + this.message += 'Request failed with: ' + rpcError.message; + } + } +} diff --git a/test/table.ts b/test/table.ts index f0833ef77..881951eea 100644 --- a/test/table.ts +++ b/test/table.ts @@ -110,6 +110,16 @@ describe('Bigtable/Table', () => { let table: any; before(() => { + const FakeTabularApiService = proxyquire('../src/tabular-api-service.js', { + '@google-cloud/promisify': fakePromisify, + './family.js': {Family: FakeFamily}, + './mutation.js': {Mutation: FakeMutation}, + './filter.js': {Filter: FakeFilter}, + pumpify, + './row.js': {Row: FakeRow}, + './chunktransformer.js': {ChunkTransformer: FakeChunkTransformer}, + }).TabularApiService; + // TODO: Consider removing this proxyquire for Table Table = proxyquire('../src/table.js', { '@google-cloud/promisify': fakePromisify, './family.js': {Family: FakeFamily}, @@ -117,6 +127,7 @@ describe('Bigtable/Table', () => { './filter.js': {Filter: FakeFilter}, pumpify, './row.js': {Row: FakeRow}, + './tabular-api-service': {TabularApiService: FakeTabularApiService}, './chunktransformer.js': {ChunkTransformer: FakeChunkTransformer}, }).Table; }); From 7491a3d1077a9095e094b32bd9f9d6fdedfb0adc Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Thu, 22 Aug 2024 11:06:21 -0400 Subject: [PATCH 05/13] Move all the ReadRows functionality over --- src/row.ts | 5 +- src/table.ts | 463 +------------------------------------ src/tabular-api-service.ts | 445 ++++++++++++++++++++++++++++++++++- src/utils/table.ts | 2 +- 4 files changed, 457 insertions(+), 458 deletions(-) diff --git a/src/row.ts b/src/row.ts index a2a9368f5..ccbd65986 100644 --- a/src/row.ts +++ b/src/row.ts @@ -31,6 +31,7 @@ import {Chunk} from './chunktransformer'; import {CallOptions} from 'google-gax'; import {ServiceError} from 'google-gax'; import {google} from '../protos/protos'; +import {TabularApiService} from './tabular-api-service'; export interface Rule { column: string; @@ -156,13 +157,13 @@ export class RowError extends Error { */ export class Row { bigtable: Bigtable; - table: Table; + table: TabularApiService; id: string; // eslint-disable-next-line @typescript-eslint/no-explicit-any data: any; key?: string; metadata?: {}; - constructor(table: Table, key: string) { + constructor(table: TabularApiService, key: string) { this.bigtable = table.bigtable; this.table = table; this.id = key; diff --git a/src/table.ts b/src/table.ts index fdaf4cbc8..a26e420bc 100644 --- a/src/table.ts +++ b/src/table.ts @@ -15,14 +15,6 @@ import {promisifyAll} from '@google-cloud/promisify'; import arrify = require('arrify'); import {ServiceError} from 'google-gax'; -import {BackoffSettings} from 'google-gax/build/src/gax'; -import {PassThrough, Transform} from 'stream'; - -// eslint-disable-next-line @typescript-eslint/no-var-requires -const concat = require('concat-stream'); -import * as is from 'is'; -// eslint-disable-next-line @typescript-eslint/no-var-requires -const pumpify = require('pumpify'); import { Family, @@ -31,30 +23,26 @@ import { CreateFamilyResponse, IColumnFamily, } from './family'; -import {Filter, BoundData, RawFilter} from './filter'; +import {BoundData, RawFilter} from './filter'; import {Mutation} from './mutation'; import {Row} from './row'; -import {ChunkTransformer} from './chunktransformer'; import {CallOptions} from 'google-gax'; -import {Bigtable, AbortableDuplex} from '.'; import {Instance} from './instance'; import {ModifiableBackupFields} from './backup'; import {CreateBackupCallback, CreateBackupResponse} from './cluster'; import {google} from '../protos/protos'; -import {Duplex} from 'stream'; import {TableUtils} from './utils/table'; import { - DEFAULT_BACKOFF_SETTINGS, - getNextDelay, - IGNORED_STATUS_CODES, - populateAttemptHeader, - RETRYABLE_STATUS_CODES, TabularApiService, InsertRowsCallback, InsertRowsResponse, MutateCallback, MutateResponse, PartialFailureError, + PrefixRange, + GetRowsOptions, + GetRowsCallback, + GetRowsResponse, } from './tabular-api-service'; export { @@ -63,6 +51,10 @@ export { MutateCallback, MutateResponse, PartialFailureError, + PrefixRange, + GetRowsOptions, + GetRowsCallback, + GetRowsResponse, }; /** @@ -202,63 +194,6 @@ export interface GetTablesOptions { pageToken?: string; } -export interface GetRowsOptions { - /** - * If set to `false` it will not decode Buffer values returned from Bigtable. - */ - decode?: boolean; - - /** - * The encoding to use when converting Buffer values to a string. - */ - encoding?: string; - - /** - * End value for key range. - */ - end?: string; - - /** - * Row filters allow you to both make advanced queries and format how the data is returned. - */ - filter?: RawFilter; - - /** - * Request configuration options, outlined here: https://googleapis.github.io/gax-nodejs/CallSettings.html. - */ - gaxOptions?: CallOptions; - - /** - * A list of row keys. - */ - keys?: string[]; - - /** - * Maximum number of rows to be returned. - */ - limit?: number; - - /** - * Prefix that the row key must match. - */ - prefix?: string; - - /** - * List of prefixes that a row key must match. - */ - prefixes?: string[]; - - /** - * A list of key ranges. - */ - ranges?: PrefixRange[]; - - /** - * Start value for key range. - */ - start?: string; -} - export interface GetMetadataOptions { /** * Request configuration options, outlined here: https://googleapis.github.io/gax-nodejs/CallSettings.html. @@ -364,17 +299,6 @@ export type GetReplicationStatesResponse = [ Map, google.bigtable.admin.v2.ITable, ]; -export type GetRowsCallback = ( - err: ServiceError | null, - rows?: Row[], - apiResponse?: google.bigtable.v2.ReadRowsResponse -) => void; -export type GetRowsResponse = [Row[], google.bigtable.v2.ReadRowsResponse]; - -export interface PrefixRange { - start?: BoundData | string; - end?: BoundData | string; -} export interface CreateBackupConfig extends ModifiableBackupFields { gaxOptions?: CallOptions; @@ -664,317 +588,6 @@ export class Table extends TabularApiService { ); } - /** - * Get {@link Row} objects for the rows currently in your table as a - * readable object stream. - * - * @param {object} [options] Configuration object. - * @param {boolean} [options.decode=true] If set to `false` it will not decode - * Buffer values returned from Bigtable. - * @param {boolean} [options.encoding] The encoding to use when converting - * Buffer values to a string. - * @param {string} [options.end] End value for key range. - * @param {Filter} [options.filter] Row filters allow you to - * both make advanced queries and format how the data is returned. - * @param {object} [options.gaxOptions] Request configuration options, outlined - * here: https://googleapis.github.io/gax-nodejs/CallSettings.html. - * @param {string[]} [options.keys] A list of row keys. - * @param {number} [options.limit] Maximum number of rows to be returned. - * @param {string} [options.prefix] Prefix that the row key must match. - * @param {string[]} [options.prefixes] List of prefixes that a row key must - * match. - * @param {object[]} [options.ranges] A list of key ranges. - * @param {string} [options.start] Start value for key range. - * @returns {stream} - * - * @example include:samples/api-reference-doc-snippets/table.js - * region_tag:bigtable_api_table_readstream - */ - createReadStream(opts?: GetRowsOptions) { - const options = opts || {}; - const maxRetries = is.number(this.maxRetries) ? this.maxRetries! : 10; - let activeRequestStream: AbortableDuplex | null; - let rowKeys: string[]; - let filter: {} | null; - const rowsLimit = options.limit || 0; - const hasLimit = rowsLimit !== 0; - - let numConsecutiveErrors = 0; - let numRequestsMade = 0; - let retryTimer: NodeJS.Timeout | null; - - rowKeys = options.keys || []; - - const ranges = TableUtils.getRanges(options); - - // If rowKeys and ranges are both empty, the request is a full table scan. - // Add an empty range to simplify the resumption logic. - if (rowKeys.length === 0 && ranges.length === 0) { - ranges.push({}); - } - - if (options.filter) { - filter = Filter.parse(options.filter); - } - - let chunkTransformer: ChunkTransformer; - let rowStream: Duplex; - - let userCanceled = false; - // The key of the last row that was emitted by the per attempt pipeline - // Note: this must be updated from the operation level userStream to avoid referencing buffered rows that will be - // discarded in the per attempt subpipeline (rowStream) - let lastRowKey = ''; - let rowsRead = 0; - const userStream = new PassThrough({ - objectMode: true, - readableHighWaterMark: 0, // We need to disable readside buffering to allow for acceptable behavior when the end user cancels the stream early. - writableHighWaterMark: 0, // We need to disable writeside buffering because in nodejs 14 the call to _transform happens after write buffering. This creates problems for tracking the last seen row key. - transform(row, _encoding, callback) { - if (userCanceled) { - callback(); - return; - } - if (TableUtils.lessThanOrEqualTo(row.id, lastRowKey)) { - /* - Sometimes duplicate rows reach this point. To avoid delivering - duplicate rows to the user, rows are thrown away if they don't exceed - the last row key. We can expect each row to reach this point and rows - are delivered in order so if the last row key equals or exceeds the - row id then we know data for this row has already reached this point - and been delivered to the user. In this case we want to throw the row - away and we do not want to deliver this row to the user again. - */ - callback(); - return; - } - lastRowKey = row.id; - rowsRead++; - callback(null, row); - }, - }); - - // The caller should be able to call userStream.end() to stop receiving - // more rows and cancel the stream prematurely. But also, the 'end' event - // will be emitted if the stream ended normally. To tell these two - // situations apart, we'll save the "original" end() function, and - // will call it on rowStream.on('end'). - const originalEnd = userStream.end.bind(userStream); - - // Taking care of this extra listener when piping and unpiping userStream: - const rowStreamPipe = (rowStream: Duplex, userStream: PassThrough) => { - rowStream.pipe(userStream, {end: false}); - rowStream.on('end', originalEnd); - }; - const rowStreamUnpipe = (rowStream: Duplex, userStream: PassThrough) => { - rowStream?.unpipe(userStream); - rowStream?.removeListener('end', originalEnd); - }; - - // eslint-disable-next-line @typescript-eslint/no-explicit-any - userStream.end = (chunk?: any, encoding?: any, cb?: () => void) => { - rowStreamUnpipe(rowStream, userStream); - userCanceled = true; - if (activeRequestStream) { - activeRequestStream.abort(); - } - if (retryTimer) { - clearTimeout(retryTimer); - } - return originalEnd(chunk, encoding, cb); - }; - - const makeNewRequest = () => { - // Avoid cancelling an expired timer if user - // cancelled the stream in the middle of a retry - retryTimer = null; - - // eslint-disable-next-line @typescript-eslint/no-explicit-any - chunkTransformer = new ChunkTransformer({decode: options.decode} as any); - - const reqOpts = { - tableName: this.name, - appProfileId: this.bigtable.appProfileId, - } as google.bigtable.v2.IReadRowsRequest; - - const retryOpts = { - currentRetryAttempt: 0, // was numConsecutiveErrors - // Handling retries in this client. Specify the retry options to - // make sure nothing is retried in retry-request. - noResponseRetries: 0, - shouldRetryFn: (_: any) => { - return false; - }, - }; - - if (lastRowKey) { - // Readjust and/or remove ranges based on previous valid row reads. - // Iterate backward since items may need to be removed. - for (let index = ranges.length - 1; index >= 0; index--) { - const range = ranges[index]; - const startValue = is.object(range.start) - ? (range.start as BoundData).value - : range.start; - const endValue = is.object(range.end) - ? (range.end as BoundData).value - : range.end; - const startKeyIsRead = - !startValue || - TableUtils.lessThanOrEqualTo( - startValue as string, - lastRowKey as string - ); - const endKeyIsNotRead = - !endValue || - (endValue as Buffer).length === 0 || - TableUtils.lessThan(lastRowKey as string, endValue as string); - if (startKeyIsRead) { - if (endKeyIsNotRead) { - // EndKey is not read, reset the range to start from lastRowKey open - range.start = { - value: lastRowKey, - inclusive: false, - }; - } else { - // EndKey is read, remove this range - ranges.splice(index, 1); - } - } - } - - // Remove rowKeys already read. - rowKeys = rowKeys.filter(rowKey => - TableUtils.greaterThan(rowKey, lastRowKey as string) - ); - - // If there was a row limit in the original request and - // we've already read all the rows, end the stream and - // do not retry. - if (hasLimit && rowsLimit === rowsRead) { - userStream.end(); - return; - } - // If all the row keys and ranges are read, end the stream - // and do not retry. - if (rowKeys.length === 0 && ranges.length === 0) { - userStream.end(); - return; - } - } - - // Create the new reqOpts - reqOpts.rows = {}; - - // TODO: preprocess all the keys and ranges to Bytes - reqOpts.rows.rowKeys = rowKeys.map( - Mutation.convertToBytes - ) as {} as Uint8Array[]; - - reqOpts.rows.rowRanges = ranges.map(range => - Filter.createRange( - range.start as BoundData, - range.end as BoundData, - 'Key' - ) - ); - - if (filter) { - reqOpts.filter = filter; - } - - if (hasLimit) { - reqOpts.rowsLimit = rowsLimit - rowsRead; - } - - const gaxOpts = populateAttemptHeader( - numRequestsMade, - options.gaxOptions - ); - - const requestStream = this.bigtable.request({ - client: 'BigtableClient', - method: 'readRows', - reqOpts, - gaxOpts, - retryOpts, - }); - - activeRequestStream = requestStream!; - - const toRowStream = new Transform({ - transform: (rowData, _, next) => { - if ( - userCanceled || - // eslint-disable-next-line @typescript-eslint/no-explicit-any - (userStream as any)._writableState.ended - ) { - return next(); - } - const row = this.row(rowData.key); - row.data = rowData.data; - next(null, row); - }, - objectMode: true, - }); - - rowStream = pumpify.obj([requestStream, chunkTransformer, toRowStream]); - - // Retry on "received rst stream" errors - const isRstStreamError = (error: ServiceError): boolean => { - if (error.code === 13 && error.message) { - const error_message = (error.message || '').toLowerCase(); - return ( - error.code === 13 && - (error_message.includes('rst_stream') || - error_message.includes('rst stream')) - ); - } - return false; - }; - - rowStream - .on('error', (error: ServiceError) => { - rowStreamUnpipe(rowStream, userStream); - activeRequestStream = null; - if (IGNORED_STATUS_CODES.has(error.code)) { - // We ignore the `cancelled` "error", since we are the ones who cause - // it when the user calls `.abort()`. - userStream.end(); - return; - } - numConsecutiveErrors++; - numRequestsMade++; - if ( - numConsecutiveErrors <= maxRetries && - (RETRYABLE_STATUS_CODES.has(error.code) || isRstStreamError(error)) - ) { - const backOffSettings = - options.gaxOptions?.retry?.backoffSettings || - DEFAULT_BACKOFF_SETTINGS; - const nextRetryDelay = getNextDelay( - numConsecutiveErrors, - backOffSettings - ); - retryTimer = setTimeout(makeNewRequest, nextRetryDelay); - } else { - userStream.emit('error', error); - } - }) - .on('data', _ => { - // Reset error count after a successful read so the backoff - // time won't keep increasing when as stream had multiple errors - numConsecutiveErrors = 0; - }) - .on('end', () => { - activeRequestStream = null; - }); - rowStreamPipe(rowStream, userStream); - }; - - makeNewRequest(); - return userStream; - } - delete(gaxOptions?: CallOptions): Promise; delete(gaxOptions: CallOptions, callback: DeleteTableCallback): void; delete(callback: DeleteTableCallback): void; @@ -1385,64 +998,6 @@ export class Table extends TabularApiService { ); } - getRows(options?: GetRowsOptions): Promise; - getRows(options: GetRowsOptions, callback: GetRowsCallback): void; - getRows(callback: GetRowsCallback): void; - /** - * Get {@link Row} objects for the rows currently in your table. - * - * This method is not recommended for large datasets as it will buffer all rows - * before returning the results. Instead we recommend using the streaming API - * via {@link Table#createReadStream}. - * - * @param {object} [options] Configuration object. See - * {@link Table#createReadStream} for a complete list of options. - * @param {object} [options.gaxOptions] Request configuration options, outlined - * here: https://googleapis.github.io/gax-nodejs/CallSettings.html. - * @param {function} callback The callback function. - * @param {?error} callback.err An error returned while making this request. - * @param {Row[]} callback.rows List of Row objects. - * - * @example include:samples/api-reference-doc-snippets/table.js - * region_tag:bigtable_api_get_rows - */ - getRows( - optionsOrCallback?: GetRowsOptions | GetRowsCallback, - cb?: GetRowsCallback - ): void | Promise { - const callback = - typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; - const options = - typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; - this.createReadStream(options) - .on('error', callback) - .pipe( - concat((rows: Row[]) => { - callback(null, rows); - }) - ); - } - - /** - * Get a reference to a table row. - * - * @throws {error} If a key is not provided. - * - * @param {string} key The row key. - * @returns {Row} - * - * @example - * ``` - * const row = table.row('lincoln'); - * ``` - */ - row(key: string): Row { - if (!key) { - throw new Error('A row key must be provided.'); - } - return new Row(this, key); - } - setIamPolicy( policy: Policy, gaxOptions?: CallOptions diff --git a/src/tabular-api-service.ts b/src/tabular-api-service.ts index 7df826b90..8206b809e 100644 --- a/src/tabular-api-service.ts +++ b/src/tabular-api-service.ts @@ -3,18 +3,23 @@ import arrify = require('arrify'); import {Instance} from './instance'; import {Mutation} from './mutation'; import { + AbortableDuplex, Bigtable, Entry, MutateOptions, SampleRowKeysCallback, SampleRowsKeysResponse, } from './index'; +import {Filter, BoundData, RawFilter} from './filter'; +import {Row} from './row'; +import {ChunkTransformer} from './chunktransformer'; import {BackoffSettings} from 'google-gax/build/src/gax'; import {google} from '../protos/protos'; import {CallOptions, ServiceError} from 'google-gax'; -import {Transform} from 'stream'; +import {Duplex, PassThrough, Transform} from 'stream'; import * as is from 'is'; import {GoogleInnerError} from './table'; +import {TableUtils} from './utils/table'; // See protos/google/rpc/code.proto // (4=DEADLINE_EXCEEDED, 8=RESOURCE_EXHAUSTED, 10=ABORTED, 14=UNAVAILABLE) @@ -39,6 +44,75 @@ export type MutateCallback = ( ) => void; export type MutateResponse = [google.protobuf.Empty]; +export interface GetRowsOptions { + /** + * If set to `false` it will not decode Buffer values returned from Bigtable. + */ + decode?: boolean; + + /** + * The encoding to use when converting Buffer values to a string. + */ + encoding?: string; + + /** + * End value for key range. + */ + end?: string; + + /** + * Row filters allow you to both make advanced queries and format how the data is returned. + */ + filter?: RawFilter; + + /** + * Request configuration options, outlined here: https://googleapis.github.io/gax-nodejs/CallSettings.html. + */ + gaxOptions?: CallOptions; + + /** + * A list of row keys. + */ + keys?: string[]; + + /** + * Maximum number of rows to be returned. + */ + limit?: number; + + /** + * Prefix that the row key must match. + */ + prefix?: string; + + /** + * List of prefixes that a row key must match. + */ + prefixes?: string[]; + + /** + * A list of key ranges. + */ + ranges?: PrefixRange[]; + + /** + * Start value for key range. + */ + start?: string; +} + +export type GetRowsCallback = ( + err: ServiceError | null, + rows?: Row[], + apiResponse?: google.bigtable.v2.ReadRowsResponse +) => void; +export type GetRowsResponse = [Row[], google.bigtable.v2.ReadRowsResponse]; + +export interface PrefixRange { + start?: BoundData | string; + end?: BoundData | string; +} + // eslint-disable-next-line @typescript-eslint/no-var-requires const concat = require('concat-stream'); // eslint-disable-next-line @typescript-eslint/no-var-requires @@ -73,6 +147,355 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); this.id = name.split('/').pop()!; } + /** + * Get {@link Row} objects for the rows currently in your table as a + * readable object stream. + * + * @param {object} [options] Configuration object. + * @param {boolean} [options.decode=true] If set to `false` it will not decode + * Buffer values returned from Bigtable. + * @param {boolean} [options.encoding] The encoding to use when converting + * Buffer values to a string. + * @param {string} [options.end] End value for key range. + * @param {Filter} [options.filter] Row filters allow you to + * both make advanced queries and format how the data is returned. + * @param {object} [options.gaxOptions] Request configuration options, outlined + * here: https://googleapis.github.io/gax-nodejs/CallSettings.html. + * @param {string[]} [options.keys] A list of row keys. + * @param {number} [options.limit] Maximum number of rows to be returned. + * @param {string} [options.prefix] Prefix that the row key must match. + * @param {string[]} [options.prefixes] List of prefixes that a row key must + * match. + * @param {object[]} [options.ranges] A list of key ranges. + * @param {string} [options.start] Start value for key range. + * @returns {stream} + * + * @example include:samples/api-reference-doc-snippets/table.js + * region_tag:bigtable_api_table_readstream + */ + createReadStream(opts?: GetRowsOptions) { + const options = opts || {}; + const maxRetries = is.number(this.maxRetries) ? this.maxRetries! : 10; + let activeRequestStream: AbortableDuplex | null; + let rowKeys: string[]; + let filter: {} | null; + const rowsLimit = options.limit || 0; + const hasLimit = rowsLimit !== 0; + + let numConsecutiveErrors = 0; + let numRequestsMade = 0; + let retryTimer: NodeJS.Timeout | null; + + rowKeys = options.keys || []; + + const ranges = TableUtils.getRanges(options); + + // If rowKeys and ranges are both empty, the request is a full table scan. + // Add an empty range to simplify the resumption logic. + if (rowKeys.length === 0 && ranges.length === 0) { + ranges.push({}); + } + + if (options.filter) { + filter = Filter.parse(options.filter); + } + + let chunkTransformer: ChunkTransformer; + let rowStream: Duplex; + + let userCanceled = false; + // The key of the last row that was emitted by the per attempt pipeline + // Note: this must be updated from the operation level userStream to avoid referencing buffered rows that will be + // discarded in the per attempt subpipeline (rowStream) + let lastRowKey = ''; + let rowsRead = 0; + const userStream = new PassThrough({ + objectMode: true, + readableHighWaterMark: 0, // We need to disable readside buffering to allow for acceptable behavior when the end user cancels the stream early. + writableHighWaterMark: 0, // We need to disable writeside buffering because in nodejs 14 the call to _transform happens after write buffering. This creates problems for tracking the last seen row key. + transform(row, _encoding, callback) { + if (userCanceled) { + callback(); + return; + } + if (TableUtils.lessThanOrEqualTo(row.id, lastRowKey)) { + /* + Sometimes duplicate rows reach this point. To avoid delivering + duplicate rows to the user, rows are thrown away if they don't exceed + the last row key. We can expect each row to reach this point and rows + are delivered in order so if the last row key equals or exceeds the + row id then we know data for this row has already reached this point + and been delivered to the user. In this case we want to throw the row + away and we do not want to deliver this row to the user again. + */ + callback(); + return; + } + lastRowKey = row.id; + rowsRead++; + callback(null, row); + }, + }); + + // The caller should be able to call userStream.end() to stop receiving + // more rows and cancel the stream prematurely. But also, the 'end' event + // will be emitted if the stream ended normally. To tell these two + // situations apart, we'll save the "original" end() function, and + // will call it on rowStream.on('end'). + const originalEnd = userStream.end.bind(userStream); + + // Taking care of this extra listener when piping and unpiping userStream: + const rowStreamPipe = (rowStream: Duplex, userStream: PassThrough) => { + rowStream.pipe(userStream, {end: false}); + rowStream.on('end', originalEnd); + }; + const rowStreamUnpipe = (rowStream: Duplex, userStream: PassThrough) => { + rowStream?.unpipe(userStream); + rowStream?.removeListener('end', originalEnd); + }; + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + userStream.end = (chunk?: any, encoding?: any, cb?: () => void) => { + rowStreamUnpipe(rowStream, userStream); + userCanceled = true; + if (activeRequestStream) { + activeRequestStream.abort(); + } + if (retryTimer) { + clearTimeout(retryTimer); + } + return originalEnd(chunk, encoding, cb); + }; + + const makeNewRequest = () => { + // Avoid cancelling an expired timer if user + // cancelled the stream in the middle of a retry + retryTimer = null; + + // eslint-disable-next-line @typescript-eslint/no-explicit-any + chunkTransformer = new ChunkTransformer({decode: options.decode} as any); + + const reqOpts = { + tableName: this.name, + appProfileId: this.bigtable.appProfileId, + } as google.bigtable.v2.IReadRowsRequest; + + const retryOpts = { + currentRetryAttempt: 0, // was numConsecutiveErrors + // Handling retries in this client. Specify the retry options to + // make sure nothing is retried in retry-request. + noResponseRetries: 0, + shouldRetryFn: (_: any) => { + return false; + }, + }; + + if (lastRowKey) { + // Readjust and/or remove ranges based on previous valid row reads. + // Iterate backward since items may need to be removed. + for (let index = ranges.length - 1; index >= 0; index--) { + const range = ranges[index]; + const startValue = is.object(range.start) + ? (range.start as BoundData).value + : range.start; + const endValue = is.object(range.end) + ? (range.end as BoundData).value + : range.end; + const startKeyIsRead = + !startValue || + TableUtils.lessThanOrEqualTo( + startValue as string, + lastRowKey as string + ); + const endKeyIsNotRead = + !endValue || + (endValue as Buffer).length === 0 || + TableUtils.lessThan(lastRowKey as string, endValue as string); + if (startKeyIsRead) { + if (endKeyIsNotRead) { + // EndKey is not read, reset the range to start from lastRowKey open + range.start = { + value: lastRowKey, + inclusive: false, + }; + } else { + // EndKey is read, remove this range + ranges.splice(index, 1); + } + } + } + + // Remove rowKeys already read. + rowKeys = rowKeys.filter(rowKey => + TableUtils.greaterThan(rowKey, lastRowKey as string) + ); + + // If there was a row limit in the original request and + // we've already read all the rows, end the stream and + // do not retry. + if (hasLimit && rowsLimit === rowsRead) { + userStream.end(); + return; + } + // If all the row keys and ranges are read, end the stream + // and do not retry. + if (rowKeys.length === 0 && ranges.length === 0) { + userStream.end(); + return; + } + } + + // Create the new reqOpts + reqOpts.rows = {}; + + // TODO: preprocess all the keys and ranges to Bytes + reqOpts.rows.rowKeys = rowKeys.map( + Mutation.convertToBytes + ) as {} as Uint8Array[]; + + reqOpts.rows.rowRanges = ranges.map(range => + Filter.createRange( + range.start as BoundData, + range.end as BoundData, + 'Key' + ) + ); + + if (filter) { + reqOpts.filter = filter; + } + + if (hasLimit) { + reqOpts.rowsLimit = rowsLimit - rowsRead; + } + + const gaxOpts = populateAttemptHeader( + numRequestsMade, + options.gaxOptions + ); + + const requestStream = this.bigtable.request({ + client: 'BigtableClient', + method: 'readRows', + reqOpts, + gaxOpts, + retryOpts, + }); + + activeRequestStream = requestStream!; + + const toRowStream = new Transform({ + transform: (rowData, _, next) => { + if ( + userCanceled || + // eslint-disable-next-line @typescript-eslint/no-explicit-any + (userStream as any)._writableState.ended + ) { + return next(); + } + const row = this.row(rowData.key); + row.data = rowData.data; + next(null, row); + }, + objectMode: true, + }); + + rowStream = pumpify.obj([requestStream, chunkTransformer, toRowStream]); + + // Retry on "received rst stream" errors + const isRstStreamError = (error: ServiceError): boolean => { + if (error.code === 13 && error.message) { + const error_message = (error.message || '').toLowerCase(); + return ( + error.code === 13 && + (error_message.includes('rst_stream') || + error_message.includes('rst stream')) + ); + } + return false; + }; + + rowStream + .on('error', (error: ServiceError) => { + rowStreamUnpipe(rowStream, userStream); + activeRequestStream = null; + if (IGNORED_STATUS_CODES.has(error.code)) { + // We ignore the `cancelled` "error", since we are the ones who cause + // it when the user calls `.abort()`. + userStream.end(); + return; + } + numConsecutiveErrors++; + numRequestsMade++; + if ( + numConsecutiveErrors <= maxRetries && + (RETRYABLE_STATUS_CODES.has(error.code) || isRstStreamError(error)) + ) { + const backOffSettings = + options.gaxOptions?.retry?.backoffSettings || + DEFAULT_BACKOFF_SETTINGS; + const nextRetryDelay = getNextDelay( + numConsecutiveErrors, + backOffSettings + ); + retryTimer = setTimeout(makeNewRequest, nextRetryDelay); + } else { + userStream.emit('error', error); + } + }) + .on('data', _ => { + // Reset error count after a successful read so the backoff + // time won't keep increasing when as stream had multiple errors + numConsecutiveErrors = 0; + }) + .on('end', () => { + activeRequestStream = null; + }); + rowStreamPipe(rowStream, userStream); + }; + + makeNewRequest(); + return userStream; + } + + getRows(options?: GetRowsOptions): Promise; + getRows(options: GetRowsOptions, callback: GetRowsCallback): void; + getRows(callback: GetRowsCallback): void; + /** + * Get {@link Row} objects for the rows currently in your table. + * + * This method is not recommended for large datasets as it will buffer all rows + * before returning the results. Instead we recommend using the streaming API + * via {@link Table#createReadStream}. + * + * @param {object} [options] Configuration object. See + * {@link Table#createReadStream} for a complete list of options. + * @param {object} [options.gaxOptions] Request configuration options, outlined + * here: https://googleapis.github.io/gax-nodejs/CallSettings.html. + * @param {function} callback The callback function. + * @param {?error} callback.err An error returned while making this request. + * @param {Row[]} callback.rows List of Row objects. + * + * @example include:samples/api-reference-doc-snippets/table.js + * region_tag:bigtable_api_get_rows + */ + getRows( + optionsOrCallback?: GetRowsOptions | GetRowsCallback, + cb?: GetRowsCallback + ): void | Promise { + const callback = + typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; + const options = + typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; + this.createReadStream(options) + .on('error', callback) + .pipe( + concat((rows: Row[]) => { + callback(null, rows); + }) + ); + } + insert( entries: Entry | Entry[], gaxOptions?: CallOptions @@ -281,6 +704,26 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); makeNextBatchRequest(); } + /** + * Get a reference to a table row. + * + * @throws {error} If a key is not provided. + * + * @param {string} key The row key. + * @returns {Row} + * + * @example + * ``` + * const row = table.row('lincoln'); + * ``` + */ + row(key: string): Row { + if (!key) { + throw new Error('A row key must be provided.'); + } + return new Row(this, key); + } + sampleRowKeys(gaxOptions?: CallOptions): Promise; sampleRowKeys(gaxOptions: CallOptions, callback: SampleRowKeysCallback): void; sampleRowKeys(callback?: SampleRowKeysCallback): void; diff --git a/src/utils/table.ts b/src/utils/table.ts index 8785ad516..6c2d1ea3f 100644 --- a/src/utils/table.ts +++ b/src/utils/table.ts @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -import {GetRowsOptions, PrefixRange} from '../table'; +import {GetRowsOptions, PrefixRange} from '../tabular-api-service'; import {Mutation} from '../mutation'; export class TableUtils { From 111a9edb800b8b1c0fe534eff71cc9dc24b1e5ec Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Thu, 22 Aug 2024 11:10:15 -0400 Subject: [PATCH 06/13] Solve the issue with the is dependency --- src/table.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/table.ts b/src/table.ts index a26e420bc..381e90722 100644 --- a/src/table.ts +++ b/src/table.ts @@ -23,15 +23,14 @@ import { CreateFamilyResponse, IColumnFamily, } from './family'; -import {BoundData, RawFilter} from './filter'; import {Mutation} from './mutation'; -import {Row} from './row'; import {CallOptions} from 'google-gax'; import {Instance} from './instance'; import {ModifiableBackupFields} from './backup'; import {CreateBackupCallback, CreateBackupResponse} from './cluster'; import {google} from '../protos/protos'; import {TableUtils} from './utils/table'; +import * as is from 'is'; import { TabularApiService, InsertRowsCallback, From 4e82e17040e59c68a0d11d7c027423e1cf4a7817 Mon Sep 17 00:00:00 2001 From: Owl Bot Date: Thu, 22 Aug 2024 17:30:41 +0000 Subject: [PATCH 07/13] =?UTF-8?q?=F0=9F=A6=89=20Updates=20from=20OwlBot=20?= =?UTF-8?q?post-processor?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit See https://github.com/googleapis/repo-automation-bots/blob/main/packages/owl-bot/README.md --- protos/protos.json | 3 +++ 1 file changed, 3 insertions(+) diff --git a/protos/protos.json b/protos/protos.json index 13a104bf7..edeeecd74 100644 --- a/protos/protos.json +++ b/protos/protos.json @@ -1,4 +1,7 @@ { + "options": { + "syntax": "proto3" + }, "nested": { "google": { "nested": { From f5364a0b5c25fc43588db8f1af5aee80edcd17b5 Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Thu, 22 Aug 2024 14:35:01 -0400 Subject: [PATCH 08/13] Add header for new class --- src/tabular-api-service.ts | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/src/tabular-api-service.ts b/src/tabular-api-service.ts index 8206b809e..72eff1317 100644 --- a/src/tabular-api-service.ts +++ b/src/tabular-api-service.ts @@ -1,3 +1,17 @@ +// Copyright 2024 Google LLC +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// https://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + import {promisifyAll} from '@google-cloud/promisify'; import arrify = require('arrify'); import {Instance} from './instance'; From e06689c3437a146ea697829022e9bee11b3e4725 Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Fri, 23 Aug 2024 14:23:02 -0400 Subject: [PATCH 09/13] Only include Table mocks that are necessary --- test/table.ts | 3 --- 1 file changed, 3 deletions(-) diff --git a/test/table.ts b/test/table.ts index 881951eea..28702f630 100644 --- a/test/table.ts +++ b/test/table.ts @@ -124,11 +124,8 @@ describe('Bigtable/Table', () => { '@google-cloud/promisify': fakePromisify, './family.js': {Family: FakeFamily}, './mutation.js': {Mutation: FakeMutation}, - './filter.js': {Filter: FakeFilter}, - pumpify, './row.js': {Row: FakeRow}, './tabular-api-service': {TabularApiService: FakeTabularApiService}, - './chunktransformer.js': {ChunkTransformer: FakeChunkTransformer}, }).Table; }); From bf27e6bc2d5e40cab32fc8627d44a78159571eb2 Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Fri, 23 Aug 2024 14:23:44 -0400 Subject: [PATCH 10/13] Remove TODO --- test/table.ts | 1 - 1 file changed, 1 deletion(-) diff --git a/test/table.ts b/test/table.ts index 28702f630..013a3a3b5 100644 --- a/test/table.ts +++ b/test/table.ts @@ -119,7 +119,6 @@ describe('Bigtable/Table', () => { './row.js': {Row: FakeRow}, './chunktransformer.js': {ChunkTransformer: FakeChunkTransformer}, }).TabularApiService; - // TODO: Consider removing this proxyquire for Table Table = proxyquire('../src/table.js', { '@google-cloud/promisify': fakePromisify, './family.js': {Family: FakeFamily}, From 7aa12fb451191897ee3522881946e3e6a791e89f Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Wed, 28 Aug 2024 13:22:54 -0400 Subject: [PATCH 11/13] Rename TabularApiService to TabularApiSurface --- src/row.ts | 6 +++--- src/table.ts | 6 +++--- src/{tabular-api-service.ts => tabular-api-surface.ts} | 4 ++-- src/utils/table.ts | 2 +- test/table.ts | 6 +++--- 5 files changed, 12 insertions(+), 12 deletions(-) rename src/{tabular-api-service.ts => tabular-api-surface.ts} (99%) diff --git a/src/row.ts b/src/row.ts index ccbd65986..d7651b423 100644 --- a/src/row.ts +++ b/src/row.ts @@ -31,7 +31,7 @@ import {Chunk} from './chunktransformer'; import {CallOptions} from 'google-gax'; import {ServiceError} from 'google-gax'; import {google} from '../protos/protos'; -import {TabularApiService} from './tabular-api-service'; +import {TabularApiSurface} from './tabular-api-surface'; export interface Rule { column: string; @@ -157,13 +157,13 @@ export class RowError extends Error { */ export class Row { bigtable: Bigtable; - table: TabularApiService; + table: TabularApiSurface; id: string; // eslint-disable-next-line @typescript-eslint/no-explicit-any data: any; key?: string; metadata?: {}; - constructor(table: TabularApiService, key: string) { + constructor(table: TabularApiSurface, key: string) { this.bigtable = table.bigtable; this.table = table; this.id = key; diff --git a/src/table.ts b/src/table.ts index 381e90722..ce4d9c2d9 100644 --- a/src/table.ts +++ b/src/table.ts @@ -32,7 +32,7 @@ import {google} from '../protos/protos'; import {TableUtils} from './utils/table'; import * as is from 'is'; import { - TabularApiService, + TabularApiSurface, InsertRowsCallback, InsertRowsResponse, MutateCallback, @@ -42,7 +42,7 @@ import { GetRowsOptions, GetRowsCallback, GetRowsResponse, -} from './tabular-api-service'; +} from './tabular-api-surface'; export { InsertRowsCallback, @@ -318,7 +318,7 @@ export interface CreateBackupConfig extends ModifiableBackupFields { * const table = instance.table('prezzy'); * ``` */ -export class Table extends TabularApiService { +export class Table extends TabularApiSurface { /** * Formats the decodes policy etag value to string. * diff --git a/src/tabular-api-service.ts b/src/tabular-api-surface.ts similarity index 99% rename from src/tabular-api-service.ts rename to src/tabular-api-surface.ts index 72eff1317..318f32ce7 100644 --- a/src/tabular-api-service.ts +++ b/src/tabular-api-surface.ts @@ -132,7 +132,7 @@ const concat = require('concat-stream'); // eslint-disable-next-line @typescript-eslint/no-var-requires const pumpify = require('pumpify'); -export class TabularApiService { +export class TabularApiSurface { bigtable: Bigtable; instance: Instance; name: string; @@ -854,7 +854,7 @@ export function populateAttemptHeader(attempt: number, gaxOpts?: CallOptions) { * All async methods (except for streams) will return a Promise in the event * that a callback is omitted. */ -promisifyAll(TabularApiService, { +promisifyAll(TabularApiSurface, { exclude: ['family', 'row'], }); diff --git a/src/utils/table.ts b/src/utils/table.ts index 6c2d1ea3f..775542c9a 100644 --- a/src/utils/table.ts +++ b/src/utils/table.ts @@ -12,7 +12,7 @@ // See the License for the specific language governing permissions and // limitations under the License. -import {GetRowsOptions, PrefixRange} from '../tabular-api-service'; +import {GetRowsOptions, PrefixRange} from '../tabular-api-surface'; import {Mutation} from '../mutation'; export class TableUtils { diff --git a/test/table.ts b/test/table.ts index 013a3a3b5..3aa85f2fa 100644 --- a/test/table.ts +++ b/test/table.ts @@ -110,7 +110,7 @@ describe('Bigtable/Table', () => { let table: any; before(() => { - const FakeTabularApiService = proxyquire('../src/tabular-api-service.js', { + const FakeTabularApiSurface = proxyquire('../src/tabular-api-service.js', { '@google-cloud/promisify': fakePromisify, './family.js': {Family: FakeFamily}, './mutation.js': {Mutation: FakeMutation}, @@ -118,13 +118,13 @@ describe('Bigtable/Table', () => { pumpify, './row.js': {Row: FakeRow}, './chunktransformer.js': {ChunkTransformer: FakeChunkTransformer}, - }).TabularApiService; + }).TabularApiSurface; Table = proxyquire('../src/table.js', { '@google-cloud/promisify': fakePromisify, './family.js': {Family: FakeFamily}, './mutation.js': {Mutation: FakeMutation}, './row.js': {Row: FakeRow}, - './tabular-api-service': {TabularApiService: FakeTabularApiService}, + './tabular-api-service': {TabularApiService: FakeTabularApiSurface}, }).Table; }); From 6f3c1bdfeae3b1491e37e0706ae556de9937cb6b Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Wed, 28 Aug 2024 13:29:54 -0400 Subject: [PATCH 12/13] Change all imports to tabular-api-surface --- test/table.ts | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/test/table.ts b/test/table.ts index 3aa85f2fa..488aae379 100644 --- a/test/table.ts +++ b/test/table.ts @@ -110,7 +110,7 @@ describe('Bigtable/Table', () => { let table: any; before(() => { - const FakeTabularApiSurface = proxyquire('../src/tabular-api-service.js', { + const FakeTabularApiSurface = proxyquire('../src/tabular-api-surface.js', { '@google-cloud/promisify': fakePromisify, './family.js': {Family: FakeFamily}, './mutation.js': {Mutation: FakeMutation}, @@ -124,7 +124,7 @@ describe('Bigtable/Table', () => { './family.js': {Family: FakeFamily}, './mutation.js': {Mutation: FakeMutation}, './row.js': {Row: FakeRow}, - './tabular-api-service': {TabularApiService: FakeTabularApiSurface}, + './tabular-api-surface': {TabularApiService: FakeTabularApiSurface}, }).Table; }); From f8baecc75813be49ab7c06ea61ed55d0029902dd Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Wed, 28 Aug 2024 13:34:57 -0400 Subject: [PATCH 13/13] surface. not service --- test/table.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/test/table.ts b/test/table.ts index 488aae379..1c921a099 100644 --- a/test/table.ts +++ b/test/table.ts @@ -124,7 +124,7 @@ describe('Bigtable/Table', () => { './family.js': {Family: FakeFamily}, './mutation.js': {Mutation: FakeMutation}, './row.js': {Row: FakeRow}, - './tabular-api-surface': {TabularApiService: FakeTabularApiSurface}, + './tabular-api-surface': {TabularApiSurface: FakeTabularApiSurface}, }).Table; });