From 4cf81721f2ea28f662fd5ef47b14a2686d0b73a2 Mon Sep 17 00:00:00 2001 From: Daniel Bruce Date: Mon, 19 Aug 2024 17:30:21 -0400 Subject: [PATCH 1/6] Create the new tabular api service file --- src/tabular-api-service.ts | 0 1 file changed, 0 insertions(+), 0 deletions(-) create mode 100644 src/tabular-api-service.ts diff --git a/src/tabular-api-service.ts b/src/tabular-api-service.ts new file mode 100644 index 000000000..e69de29bb From 932ab9e4af255e3346eb1700dade918bc3151439 Mon Sep 17 00:00:00 2001 From: danieljbruce Date: Wed, 28 Aug 2024 13:43:08 -0400 Subject: [PATCH 2/6] feat: Bigtable authorized views - move the code over to the TabularApiSurface class (#1463) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Move the constructor over to TabularApiService * Move sampleRowKeys over * Move sampleRowKeys functions over and use promisify * Adjust the proxyquire to work with TabularAPIserv * Move all the ReadRows functionality over * Solve the issue with the is dependency * 🦉 Updates from OwlBot post-processor See https://github.com/googleapis/repo-automation-bots/blob/main/packages/owl-bot/README.md * Add header for new class * Only include Table mocks that are necessary * Remove TODO * Rename TabularApiService to TabularApiSurface * Change all imports to tabular-api-surface * surface. not service --------- Co-authored-by: Owl Bot --- protos/protos.json | 3 + src/row.ts | 5 +- src/table.ts | 860 ++---------------------------------- src/tabular-api-service.ts | 0 src/tabular-api-surface.ts | 880 +++++++++++++++++++++++++++++++++++++ src/utils/table.ts | 2 +- test/table.ts | 9 +- 7 files changed, 920 insertions(+), 839 deletions(-) delete mode 100644 src/tabular-api-service.ts create mode 100644 src/tabular-api-surface.ts 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": { diff --git a/src/row.ts b/src/row.ts index a2a9368f5..d7651b423 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 {TabularApiSurface} from './tabular-api-surface'; export interface Rule { column: string; @@ -156,13 +157,13 @@ export class RowError extends Error { */ export class Row { bigtable: Bigtable; - table: Table; + table: TabularApiSurface; id: string; // eslint-disable-next-line @typescript-eslint/no-explicit-any data: any; key?: string; metadata?: {}; - constructor(table: Table, 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 e7c286c0f..ce4d9c2d9 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,29 +23,37 @@ import { CreateFamilyResponse, IColumnFamily, } from './family'; -import {Filter, 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'; - -// 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 * as is from 'is'; +import { + TabularApiSurface, + InsertRowsCallback, + InsertRowsResponse, + MutateCallback, + MutateResponse, + PartialFailureError, + PrefixRange, + GetRowsOptions, + GetRowsCallback, + GetRowsResponse, +} from './tabular-api-surface'; + +export { + InsertRowsCallback, + InsertRowsResponse, + MutateCallback, + MutateResponse, + PartialFailureError, + PrefixRange, + GetRowsOptions, + GetRowsCallback, + GetRowsResponse, }; /** @@ -193,63 +193,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. @@ -355,27 +298,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 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; - end?: BoundData | string; -} export interface CreateBackupConfig extends ModifiableBackupFields { gaxOptions?: CallOptions; @@ -396,34 +318,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 TabularApiSurface { /** * Formats the decodes policy etag value to string. * @@ -692,317 +587,6 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); ); } - /** - * 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; @@ -1413,360 +997,6 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); ); } - 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 - ): 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. - * - * @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; - /** - * 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 @@ -2071,47 +1301,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 deleted file mode 100644 index e69de29bb..000000000 diff --git a/src/tabular-api-surface.ts b/src/tabular-api-surface.ts new file mode 100644 index 000000000..318f32ce7 --- /dev/null +++ b/src/tabular-api-surface.ts @@ -0,0 +1,880 @@ +// 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'; +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 {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) +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]; + +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 +const pumpify = require('pumpify'); + +export class TabularApiSurface { + 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()!; + } + + /** + * 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 + ): 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. + * + * @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; + /** + * 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, + ]); + } +} + +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 + * that a callback is omitted. + */ +promisifyAll(TabularApiSurface, { + 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/src/utils/table.ts b/src/utils/table.ts index 8785ad516..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 '../table'; +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 f0833ef77..1c921a099 100644 --- a/test/table.ts +++ b/test/table.ts @@ -110,7 +110,7 @@ describe('Bigtable/Table', () => { let table: any; before(() => { - Table = proxyquire('../src/table.js', { + const FakeTabularApiSurface = proxyquire('../src/tabular-api-surface.js', { '@google-cloud/promisify': fakePromisify, './family.js': {Family: FakeFamily}, './mutation.js': {Mutation: FakeMutation}, @@ -118,6 +118,13 @@ describe('Bigtable/Table', () => { pumpify, './row.js': {Row: FakeRow}, './chunktransformer.js': {ChunkTransformer: FakeChunkTransformer}, + }).TabularApiSurface; + Table = proxyquire('../src/table.js', { + '@google-cloud/promisify': fakePromisify, + './family.js': {Family: FakeFamily}, + './mutation.js': {Mutation: FakeMutation}, + './row.js': {Row: FakeRow}, + './tabular-api-surface': {TabularApiSurface: FakeTabularApiSurface}, }).Table; }); From cb1ebd4b773805fbdc9c9614525365ef47eeac4b Mon Sep 17 00:00:00 2001 From: danieljbruce Date: Mon, 23 Sep 2024 16:13:37 -0400 Subject: [PATCH 3/6] feat: Bigtable Authorized Views - Allow checkAndMutate and ReadModifyWriteRow calls on Authorized Views (#1464) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Move the constructor over to TabularApiService * Move sampleRowKeys over * Move sampleRowKeys functions over and use promisify * Adjust the proxyquire to work with TabularAPIserv * Move all the ReadRows functionality over * Solve the issue with the is dependency * 🦉 Updates from OwlBot post-processor See https://github.com/googleapis/repo-automation-bots/blob/main/packages/owl-bot/README.md * Add header for new class * Trying out the FilterInformation class * Add a DataUtils module for the shared row function * Outsource functionality of filter to helper * debugging * Adjust proxyquire to include mocks for moved fn * Remove imports * Outsource code to a createRulesUtil function * Change proxyquire to include createRulesUtil * Move increment over to the utils folder * Move the functions into a static class for mocks * Remove mockCreateRules and shorten mock * Stub out FakeRowDataUtil * Remove unused dependencies * Add the functions to work with checkAndMutate and readWriteModifyRow * Fix regressions from the merge * Add documentation for the class * Document the new methods of table * Change the interface of the rowUtils * Pull the generateProperties code out avoid duplica * Move duplicate code out into a getProperties object * Remove console.log * Add method for creating views * Object for making grpc calls for authorized views * Remove import * Update the documentation for the Table * More specific type * Add documentation for each of the functions * Remove TODO * Add headers * Add documentation for new class * Remove imports * Reintroduce before * Remove unused import --------- Co-authored-by: Owl Bot --- src/authorized-view.ts | 255 ++++++++++++++++++++++++++++++++++++ src/chunktransformer.ts | 2 +- src/row-data-utils.ts | 258 +++++++++++++++++++++++++++++++++++++ src/row.ts | 164 ++++------------------- src/table.ts | 14 ++ src/tabular-api-surface.ts | 19 ++- test/row.ts | 74 +++++++---- 7 files changed, 618 insertions(+), 168 deletions(-) create mode 100644 src/authorized-view.ts create mode 100644 src/row-data-utils.ts diff --git a/src/authorized-view.ts b/src/authorized-view.ts new file mode 100644 index 000000000..b2dce5209 --- /dev/null +++ b/src/authorized-view.ts @@ -0,0 +1,255 @@ +// 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 {TabularApiSurface} from './tabular-api-surface'; +import {CallOptions} from 'google-gax'; +import { + CreateRulesCallback, + CreateRulesResponse, + FilterCallback, + FilterConfig, + FilterResponse, + IncrementCallback, + IncrementResponse, + Rule, +} from './row'; +import {RowDataUtils, RowProperties} from './row-data-utils'; +import {RawFilter} from './filter'; +import {Table} from './table'; +import {Family} from './chunktransformer'; + +interface FilterInformation { + filter: RawFilter; + rowId: string; +} + +interface CreateRulesInformation { + rules: Rule | Rule[]; + rowId: string; +} + +interface IncrementInformation { + column: string; + rowId: string; +} + +/** + * The AuthorizedView class is a class that is available to the user that + * contains methods the user can call to work with authorized views. + * + * @class + * @param {Table} table The table that the authorized view exists on. + * @param {string} id Unique identifier of the authorized view. + * + */ +export class AuthorizedView extends TabularApiSurface { + private readonly rowData: {[id: string]: {[index: string]: Family}}; + + constructor(table: Table, viewName: string) { + super(table.instance, table.id, viewName); + this.rowData = {}; + } + + createRules( + createRulesInfo: CreateRulesInformation, + options?: CallOptions + ): Promise; + createRules( + createRulesInfo: CreateRulesInformation, + options: CallOptions, + callback: CreateRulesCallback + ): void; + createRules( + createRulesInfo: CreateRulesInformation, + callback: CreateRulesCallback + ): void; + /** + * Update a row with rules specifying how the row's contents are to be + * transformed into writes. Rules are applied in order, meaning that earlier + * rules will affect the results of later ones. + * + * @throws {error} If no rules are provided. + * + * @param {CreateRulesInformation} createRulesInfo The rules to apply to a row + * along with the row id of the row to update. + * @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.apiResponse The full API response. + * + * @example include:samples/api-reference-doc-snippets/row.js + * region_tag:bigtable_api_create_rules + */ + createRules( + createRulesInfo: CreateRulesInformation, + optionsOrCallback?: CallOptions | CreateRulesCallback, + cb?: CreateRulesCallback + ): void | Promise { + this.initializeRow(createRulesInfo.rowId); + RowDataUtils.createRulesUtil( + createRulesInfo.rules, + this.generateProperties(createRulesInfo.rowId), + optionsOrCallback, + cb + ); + } + + /** + * Mutates a row atomically based on the output of a filter. Depending on + * whether or not any results are yielded, either the `onMatch` or `onNoMatch` + * callback will be executed. + * + * @param {FilterInformation} filter Filter to be applied to the contents of + * the row along with the row id of the affected row. + * @param {object} config Configuration object. + * @param {?object[]} config.onMatch A list of entries to be ran if a match is + * found. + * @param {object[]} [config.onNoMatch] A list of entries to be ran if no + * matches are found. + * @param {object} [config.gaxOptions] Request configuration options, outlined + * here: https://googleapis.github.io/gax-nodejs/global.html#CallOptions. + * @param {function} callback The callback function. + * @param {?error} callback.err An error returned while making this + * request. + * @param {boolean} callback.matched Whether a match was found or not. + * + * @example include:samples/api-reference-doc-snippets/row.js + * region_tag:bigtable_api_row_filter + */ + filter( + filterInfo: FilterInformation, + config?: FilterConfig + ): Promise; + filter( + filterInfo: FilterInformation, + config: FilterConfig, + callback: FilterCallback + ): void; + filter(filterInfo: FilterInformation, callback: FilterCallback): void; + filter( + filterInfo: FilterInformation, + configOrCallback?: FilterConfig | FilterCallback, + cb?: FilterCallback + ): void | Promise { + this.initializeRow(filterInfo.rowId); + RowDataUtils.filterUtil( + filterInfo.filter, + this.generateProperties(filterInfo.rowId), + configOrCallback, + cb + ); + } + + /** + * Generates request properties necessary for making an rpc call for an + * authorized view. + * + * @param {string} id The row id to generate the properties for. + * @private + */ + private generateProperties(id: string): RowProperties { + return { + requestData: { + data: this.rowData[id], + id, + table: this, + bigtable: this.bigtable, + }, + reqOpts: { + authorizedViewName: this.name + '/authorizedViews/' + this.viewName, + }, + }; + } + + increment( + columnInfo: IncrementInformation, + value?: number + ): Promise; + increment( + columnInfo: IncrementInformation, + value: number, + options?: CallOptions + ): Promise; + increment( + columnInfo: IncrementInformation, + options?: CallOptions + ): Promise; + increment( + columnInfo: IncrementInformation, + value: number, + options: CallOptions, + callback: IncrementCallback + ): void; + increment( + columnInfo: IncrementInformation, + value: number, + callback: IncrementCallback + ): void; + increment( + columnInfo: IncrementInformation, + options: CallOptions, + callback: IncrementCallback + ): void; + increment( + columnInfo: IncrementInformation, + callback: IncrementCallback + ): void; + /** + * Increment a specific column within the row. If the column does not + * exist, it is automatically initialized to 0 before being incremented. + * + * @param {IncrementInformation} columnInfo The column we are incrementing a + * value in along with the row id of the affected row. + * @param {number} [value] The amount to increment by, defaults to 1. + * @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 {number} callback.value The updated value of the column. + * @param {object} callback.apiResponse The full API response. + * + * @example include:samples/api-reference-doc-snippets/row.js + * region_tag:bigtable_api_row_increment + */ + increment( + columnInfo: IncrementInformation, + valueOrOptionsOrCallback?: number | CallOptions | IncrementCallback, + optionsOrCallback?: CallOptions | IncrementCallback, + cb?: IncrementCallback + ): void | Promise { + this.initializeRow(columnInfo.rowId); + RowDataUtils.incrementUtils( + columnInfo.column, + this.generateProperties(columnInfo.rowId), + valueOrOptionsOrCallback, + optionsOrCallback, + cb + ); + } + + /** + * Sets the row data for a particular row to an empty object + * + * @param {string} id An string with the key of the row to initialize. + * @private + */ + private initializeRow(id: string) { + if (!this.rowData[id]) { + this.rowData[id] = {}; + } + } +} diff --git a/src/chunktransformer.ts b/src/chunktransformer.ts index 4aec3d9ba..77ec3af10 100644 --- a/src/chunktransformer.ts +++ b/src/chunktransformer.ts @@ -34,7 +34,7 @@ export interface Data { chunks: Chunk[]; lastScannedRowKey?: Buffer; } -interface Family { +export interface Family { [qualifier: string]: Qualifier[]; } export interface Qualifier { diff --git a/src/row-data-utils.ts b/src/row-data-utils.ts new file mode 100644 index 000000000..d2e74a57b --- /dev/null +++ b/src/row-data-utils.ts @@ -0,0 +1,258 @@ +// 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. + +const dotProp = require('dot-prop'); +import {Filter, RawFilter} from './filter'; +import { + CreateRulesCallback, + FilterCallback, + FilterConfig, + FilterConfigOption, + FormatFamiliesOptions, + IncrementCallback, + Rule, +} from './row'; +import {Family} from './chunktransformer'; +import {Bytes, Mutation} from './mutation'; +import {google} from '../protos/protos'; +import {TabularApiSurface} from './tabular-api-surface'; +import arrify = require('arrify'); +import {Bigtable} from './index'; +import {CallOptions} from 'google-gax'; + +interface TabularApiSurfaceRequest { + tableName?: string; + authorizedViewName?: string; +} + +export interface RowProperties { + requestData: { + data?: {[index: string]: Family}; + id: string; + table: TabularApiSurface; + bigtable: Bigtable; + }; + reqOpts: TabularApiSurfaceRequest; +} + +/** + * RowDataUtils is a class containing functionality needed by the Row and + * AuthorizedView classes. Its static methods need to be contained in a class + * so that they can be mocked out using the sinon library as is conventional + * throughout the rest of the client library. + */ +class RowDataUtils { + /** + * Called by `filter` methods for fulfilling table and authorized view requests. + * + * @param {Filter} filter Filter to be applied to the contents of the row. + * @param {RowProperties} properties Properties containing data for the request. + * @param {object} configOrCallback Configuration object. + * @param {function} cb The callback function. + * + */ + static filterUtil( + filter: RawFilter, + properties: RowProperties, + configOrCallback?: FilterConfig | FilterCallback, + cb?: FilterCallback + ) { + const config = typeof configOrCallback === 'object' ? configOrCallback : {}; + const callback = + typeof configOrCallback === 'function' ? configOrCallback : cb!; + const reqOpts = Object.assign( + { + appProfileId: properties.requestData.bigtable.appProfileId, + rowKey: Mutation.convertToBytes(properties.requestData.id), + predicateFilter: Filter.parse(filter), + trueMutations: createFlatMutationsList(config.onMatch!), + falseMutations: createFlatMutationsList(config.onNoMatch!), + }, + properties.reqOpts + ); + properties.requestData.data = {}; + properties.requestData.bigtable.request( + { + client: 'BigtableClient', + method: 'checkAndMutateRow', + reqOpts, + gaxOpts: config.gaxOptions, + }, + (err, apiResponse) => { + if (err) { + callback(err, null, apiResponse); + return; + } + + callback(null, apiResponse!.predicateMatched, apiResponse); + } + ); + + function createFlatMutationsList(entries: FilterConfigOption[]) { + const e2 = arrify(entries).map( + entry => Mutation.parse(entry as Mutation).mutations! + ); + return e2.reduce((a, b) => a.concat(b), []); + } + } + + static formatFamilies_Util( + families: google.bigtable.v2.IFamily[], + options?: FormatFamiliesOptions + ) { + const data = {} as {[index: string]: {}}; + options = options || {}; + families.forEach(family => { + const familyData = (data[family.name!] = {}) as { + [index: string]: {}; + }; + family.columns!.forEach(column => { + const qualifier = Mutation.convertFromBytes( + column.qualifier as string + ) as string; + familyData[qualifier] = column.cells!.map(cell => { + let value = cell.value; + if (options!.decode !== false) { + value = Mutation.convertFromBytes(value as Bytes, { + isPossibleNumber: true, + }) as string; + } + return { + value, + timestamp: cell.timestampMicros, + labels: cell.labels, + }; + }); + }); + }); + return data; + } + + /** + * Called by `createRules` methods for fulfilling table and authorized + * view requests. + * + * @param {object|object[]} rules The rules to apply to this row. + * @param {RowProperties} properties Properties containing data for the request. + * @param {object} [gaxOptions] Request configuration options, outlined here: + * https://googleapis.github.io/gax-nodejs/CallSettings.html. + * @param {function} callback The callback function. + * + */ + static createRulesUtil( + rules: Rule | Rule[], + properties: RowProperties, + optionsOrCallback?: CallOptions | CreateRulesCallback, + cb?: CreateRulesCallback + ) { + const gaxOptions = + typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; + const callback = + typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; + + if (!rules || (rules as Rule[]).length === 0) { + throw new Error('At least one rule must be provided.'); + } + + rules = arrify(rules).map(rule => { + const column = Mutation.parseColumnName(rule.column); + // eslint-disable-next-line @typescript-eslint/no-explicit-any + const ruleData: any = { + familyName: column.family, + columnQualifier: Mutation.convertToBytes(column.qualifier!), + }; + + if (rule.append) { + ruleData.appendValue = Mutation.convertToBytes(rule.append); + } + + if (rule.increment) { + ruleData.incrementAmount = rule.increment; + } + + return ruleData; + }); + + const reqOpts = Object.assign( + { + appProfileId: properties.requestData.bigtable.appProfileId, + rowKey: Mutation.convertToBytes(properties.requestData.id), + rules, + }, + properties.reqOpts + ); + properties.requestData.data = {}; + properties.requestData.bigtable.request( + { + client: 'BigtableClient', + method: 'readModifyWriteRow', + reqOpts, + gaxOpts: gaxOptions, + }, + callback + ); + } + + /** + * @param {string} column The column we are incrementing a value in. + * @param {RowProperties} properties Properties containing data for the request. + * @param {number} [valueOrOptionsOrCallback] The amount to increment by, defaults to 1. + * @param {object} [optionsOrCallback] Request configuration options, outlined here: + * https://googleapis.github.io/gax-nodejs/CallSettings.html. + * @param {function} cb The callback function. + */ + static incrementUtils( + column: string, + properties: RowProperties, + valueOrOptionsOrCallback?: number | CallOptions | IncrementCallback, + optionsOrCallback?: CallOptions | IncrementCallback, + cb?: IncrementCallback + ) { + const value = + typeof valueOrOptionsOrCallback === 'number' + ? valueOrOptionsOrCallback + : 1; + const gaxOptions = + typeof valueOrOptionsOrCallback === 'object' + ? valueOrOptionsOrCallback + : typeof optionsOrCallback === 'object' + ? optionsOrCallback + : {}; + const callback = + typeof valueOrOptionsOrCallback === 'function' + ? valueOrOptionsOrCallback + : typeof optionsOrCallback === 'function' + ? optionsOrCallback + : cb!; + + const reqOpts = { + column, + increment: value, + } as Rule; + + this.createRulesUtil(reqOpts, properties, gaxOptions, (err, resp) => { + if (err) { + callback(err, null, resp); + return; + } + + const data = this.formatFamilies_Util(resp!.row!.families!); + const value = dotProp.get(data, column.replace(':', '.'))[0].value; + + callback(null, value, resp); + }); + } +} + +export {RowDataUtils}; diff --git a/src/row.ts b/src/row.ts index d7651b423..de8c1e4bb 100644 --- a/src/row.ts +++ b/src/row.ts @@ -14,9 +14,7 @@ import {promisifyAll} from '@google-cloud/promisify'; import arrify = require('arrify'); -// eslint-disable-next-line @typescript-eslint/no-var-requires -const dotProp = require('dot-prop'); -import {Filter, RawFilter} from './filter'; +import {RawFilter} from './filter'; import {Mutation, ConvertFromBytesUserOptions, Bytes, Data} from './mutation'; import {Bigtable} from '.'; import { @@ -31,6 +29,7 @@ import {Chunk} from './chunktransformer'; import {CallOptions} from 'google-gax'; import {ServiceError} from 'google-gax'; import {google} from '../protos/protos'; +import {RowDataUtils, RowProperties} from './row-data-utils'; import {TabularApiSurface} from './tabular-api-surface'; export interface Rule { @@ -139,6 +138,18 @@ export class RowError extends Error { } } +/** + * getProperties returns the properties needed to make a request for a table. + * + * @param {Row} row The row to make a request for. + */ +function getProperties(row: Row): RowProperties { + return { + reqOpts: {tableName: row.table.name}, + requestData: row, + }; +} + /** * Create a Row object to interact with your table rows. * @@ -308,32 +319,7 @@ export class Row { families: google.bigtable.v2.IFamily[], options?: FormatFamiliesOptions ) { - const data = {} as {[index: string]: {}}; - options = options || {}; - families.forEach(family => { - const familyData = (data[family.name!] = {}) as { - [index: string]: {}; - }; - family.columns!.forEach(column => { - const qualifier = Mutation.convertFromBytes( - column.qualifier as string - ) as string; - familyData[qualifier] = column.cells!.map(cell => { - let value = cell.value; - if (options!.decode !== false) { - value = Mutation.convertFromBytes(value as Bytes, { - isPossibleNumber: true, - }) as string; - } - return { - value, - timestamp: cell.timestampMicros, - labels: cell.labels, - }; - }); - }); - }); - return data; + return RowDataUtils.formatFamilies_Util(families, options); } create(options?: CreateRowOptions): Promise; @@ -415,49 +401,11 @@ export class Row { optionsOrCallback?: CallOptions | CreateRulesCallback, cb?: CreateRulesCallback ): void | Promise { - const gaxOptions = - typeof optionsOrCallback === 'object' ? optionsOrCallback : {}; - const callback = - typeof optionsOrCallback === 'function' ? optionsOrCallback : cb!; - - if (!rules || (rules as Rule[]).length === 0) { - throw new Error('At least one rule must be provided.'); - } - - rules = arrify(rules).map(rule => { - const column = Mutation.parseColumnName(rule.column); - // eslint-disable-next-line @typescript-eslint/no-explicit-any - const ruleData: any = { - familyName: column.family, - columnQualifier: Mutation.convertToBytes(column.qualifier!), - }; - - if (rule.append) { - ruleData.appendValue = Mutation.convertToBytes(rule.append); - } - - if (rule.increment) { - ruleData.incrementAmount = rule.increment; - } - - return ruleData; - }); - - const reqOpts = { - tableName: this.table.name, - appProfileId: this.bigtable.appProfileId, - rowKey: Mutation.convertToBytes(this.id), + RowDataUtils.createRulesUtil( rules, - }; - this.data = {}; - this.bigtable.request( - { - client: 'BigtableClient', - method: 'readModifyWriteRow', - reqOpts, - gaxOpts: gaxOptions, - }, - callback + getProperties(this), + optionsOrCallback, + cb ); } @@ -622,41 +570,7 @@ export class Row { configOrCallback?: FilterConfig | FilterCallback, cb?: FilterCallback ): void | Promise { - const config = typeof configOrCallback === 'object' ? configOrCallback : {}; - const callback = - typeof configOrCallback === 'function' ? configOrCallback : cb!; - const reqOpts = { - tableName: this.table.name, - appProfileId: this.bigtable.appProfileId, - rowKey: Mutation.convertToBytes(this.id), - predicateFilter: Filter.parse(filter), - trueMutations: createFlatMutationsList(config.onMatch!), - falseMutations: createFlatMutationsList(config.onNoMatch!), - }; - this.data = {}; - this.bigtable.request( - { - client: 'BigtableClient', - method: 'checkAndMutateRow', - reqOpts, - gaxOpts: config.gaxOptions, - }, - (err, apiResponse) => { - if (err) { - callback(err, null, apiResponse); - return; - } - - callback(null, apiResponse!.predicateMatched, apiResponse); - } - ); - - function createFlatMutationsList(entries: FilterConfigOption[]) { - const e2 = arrify(entries).map( - entry => Mutation.parse(entry as Mutation).mutations! - ); - return e2.reduce((a, b) => a.concat(b), []); - } + RowDataUtils.filterUtil(filter, getProperties(this), configOrCallback, cb); } get(options?: GetRowOptions): Promise>; @@ -857,39 +771,13 @@ export class Row { optionsOrCallback?: CallOptions | IncrementCallback, cb?: IncrementCallback ): void | Promise { - const value = - typeof valueOrOptionsOrCallback === 'number' - ? valueOrOptionsOrCallback - : 1; - const gaxOptions = - typeof valueOrOptionsOrCallback === 'object' - ? valueOrOptionsOrCallback - : typeof optionsOrCallback === 'object' - ? optionsOrCallback - : {}; - const callback = - typeof valueOrOptionsOrCallback === 'function' - ? valueOrOptionsOrCallback - : typeof optionsOrCallback === 'function' - ? optionsOrCallback - : cb!; - - const reqOpts = { + RowDataUtils.incrementUtils( column, - increment: value, - } as Rule; - - this.createRules(reqOpts, gaxOptions, (err, resp) => { - if (err) { - callback(err, null, resp); - return; - } - - const data = Row.formatFamilies_(resp!.row!.families!); - const value = dotProp.get(data, column.replace(':', '.'))[0].value; - - callback(null, value, resp); - }); + getProperties(this), + valueOrOptionsOrCallback, + optionsOrCallback, + cb + ); } save(entry: Entry, options?: CallOptions): Promise; diff --git a/src/table.ts b/src/table.ts index ce4d9c2d9..16ac805a5 100644 --- a/src/table.ts +++ b/src/table.ts @@ -43,6 +43,7 @@ import { GetRowsCallback, GetRowsResponse, } from './tabular-api-surface'; +import {AuthorizedView} from './authorized-view'; export { InsertRowsCallback, @@ -319,6 +320,10 @@ export interface CreateBackupConfig extends ModifiableBackupFields { * ``` */ export class Table extends TabularApiSurface { + constructor(instance: Instance, id: string) { + super(instance, id); + } + /** * Formats the decodes policy etag value to string. * @@ -1159,6 +1164,15 @@ export class Table extends TabularApiSurface { ); } + /** + * Gets an Authorized View object for making authorized view grpc calls. + * + * @param {string} viewName The name for the Authorized view + */ + view(viewName: string): AuthorizedView { + return new AuthorizedView(this, viewName); + } + waitForReplication(): Promise; waitForReplication(callback: WaitForReplicationCallback): void; /** diff --git a/src/tabular-api-surface.ts b/src/tabular-api-surface.ts index 318f32ce7..b76782705 100644 --- a/src/tabular-api-surface.ts +++ b/src/tabular-api-surface.ts @@ -127,11 +127,26 @@ export interface PrefixRange { end?: BoundData | string; } +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 const pumpify = require('pumpify'); +/** + * The TabularApiSurface class is a class that contains methods we want to + * expose on both tables and authorized views. It also contains data that will + * be used by both Authorized views and Tables for these shared methods. + * + * @class + * @param {Instance} instance Instance Object. + * @param {string} id Unique identifier of the table. + * + */ export class TabularApiSurface { bigtable: Bigtable; instance: Instance; @@ -139,10 +154,12 @@ export class TabularApiSurface { id: string; metadata?: google.bigtable.admin.v2.ITable; maxRetries?: number; + protected viewName?: string; - constructor(instance: Instance, id: string) { + protected constructor(instance: Instance, id: string, viewName?: string) { this.bigtable = instance.bigtable; this.instance = instance; + this.viewName = viewName; let name; diff --git a/test/row.ts b/test/row.ts index d15004c48..24a22f915 100644 --- a/test/row.ts +++ b/test/row.ts @@ -14,7 +14,7 @@ import * as promisify from '@google-cloud/promisify'; import * as assert from 'assert'; -import {afterEach, before, beforeEach, describe, it} from 'mocha'; +import {afterEach, beforeEach, describe, it} from 'mocha'; import * as proxyquire from 'proxyquire'; import * as sinon from 'sinon'; import {Mutation} from '../src/mutation.js'; @@ -68,6 +68,11 @@ const FakeFilter = { }), }; +const FakeRowDataUtil = proxyquire('../src/row-data-utils.js', { + './mutation.js': {Mutation: FakeMutation}, + './filter.js': {Filter: FakeFilter}, +}).RowDataUtils; + describe('Bigtable/Row', () => { let Row: typeof rw.Row; let RowError: typeof rw.RowError; @@ -78,6 +83,7 @@ describe('Bigtable/Row', () => { '@google-cloud/promisify': fakePromisify, './mutation.js': {Mutation: FakeMutation}, './filter.js': {Filter: FakeFilter}, + './row-data-utils.js': {RowDataUtils: FakeRowDataUtil}, }); Row = Fake.Row; RowError = Fake.RowError; @@ -1318,15 +1324,17 @@ describe('Bigtable/Row', () => { let formatFamiliesSpy: sinon.SinonSpy; beforeEach(() => { - formatFamiliesSpy = sandbox.stub(Row, 'formatFamilies_').returns({ - a: { - b: [ - { - value: 10, - }, - ], - }, - }); + formatFamiliesSpy = sandbox + .stub(FakeRowDataUtil, 'formatFamilies_Util') + .returns({ + a: { + b: [ + { + value: 10, + }, + ], + }, + }); }); afterEach(() => { @@ -1334,18 +1342,20 @@ describe('Bigtable/Row', () => { }); it('should provide the proper request options', done => { - sandbox.stub(row, 'createRules').callsFake((reqOpts, gaxOptions) => { - assert.strictEqual((reqOpts as rw.Rule).column, COLUMN_NAME); - assert.strictEqual((reqOpts as rw.Rule).increment, 1); - assert.deepStrictEqual(gaxOptions, {}); - done(); - }); + sandbox + .stub(FakeRowDataUtil, 'createRulesUtil') + .callsFake((reqOpts, properties, gaxOptions, cb) => { + assert.strictEqual((reqOpts as rw.Rule).column, COLUMN_NAME); + assert.strictEqual((reqOpts as rw.Rule).increment, 1); + assert.deepStrictEqual(gaxOptions, {}); + done(); + }); row.increment(COLUMN_NAME, assert.ifError); }); it('should optionally accept an increment amount', done => { const increment = 10; - sandbox.stub(row, 'createRules').callsFake(reqOpts => { + sandbox.stub(FakeRowDataUtil, 'createRulesUtil').callsFake(reqOpts => { assert.strictEqual((reqOpts as rw.Rule).increment, increment); done(); }); @@ -1354,28 +1364,34 @@ describe('Bigtable/Row', () => { it('should accept gaxOptions', done => { const gaxOptions = {}; - sandbox.stub(row, 'createRules').callsFake((reqOpts, gaxOptions_) => { - assert.strictEqual(gaxOptions_, gaxOptions); - done(); - }); + sandbox + .stub(FakeRowDataUtil, 'createRulesUtil') + .callsFake((reqOpts, properties, gaxOptions_) => { + assert.strictEqual(gaxOptions_, gaxOptions); + done(); + }); row.increment(COLUMN_NAME, gaxOptions, assert.ifError); }); it('should accept increment amount and gaxOptions', done => { const increment = 10; const gaxOptions = {}; - sandbox.stub(row, 'createRules').callsFake((reqOpts, gaxOptions_) => { - assert.strictEqual((reqOpts as rw.Rule).increment, increment); - assert.strictEqual(gaxOptions_, gaxOptions); - done(); - }); + sandbox + .stub(FakeRowDataUtil, 'createRulesUtil') + .callsFake((reqOpts, properties, gaxOptions_) => { + assert.strictEqual((reqOpts as rw.Rule).increment, increment); + assert.strictEqual(gaxOptions_, gaxOptions); + done(); + }); row.increment(COLUMN_NAME, increment, gaxOptions, assert.ifError); }); it('should return an error to the callback', done => { const error = new Error('err'); const response = {}; - sandbox.stub(row, 'createRules').callsArgWith(2, error, response); + sandbox + .stub(FakeRowDataUtil, 'createRulesUtil') + .callsArgWith(3, error, response); row.increment(COLUMN_NAME, (err, value, apiResponse) => { assert.strictEqual(err, error); assert.strictEqual(value, null); @@ -1408,7 +1424,9 @@ describe('Bigtable/Row', () => { }, }; - sandbox.stub(row, 'createRules').callsArgWith(2, null, response); + sandbox + .stub(FakeRowDataUtil, 'createRulesUtil') + .callsArgWith(3, null, response); row.increment(COLUMN_NAME, (err, value, apiResponse) => { assert.ifError(err); assert.strictEqual(value, fakeValue); From d293abd7eaafdaef22bdc54497148cb1a52569b7 Mon Sep 17 00:00:00 2001 From: danieljbruce Date: Mon, 30 Sep 2024 13:34:21 -0400 Subject: [PATCH 4/6] feat: Create a view on the instance, not on the table (#1476) * Move view creation from the table class to instanc * feat: Create a view on the instance, not on the table * Add table name parameter to the documentation * Remove Table --- src/authorized-view.ts | 6 +++--- src/instance.ts | 11 +++++++++++ src/table.ts | 9 --------- 3 files changed, 14 insertions(+), 12 deletions(-) diff --git a/src/authorized-view.ts b/src/authorized-view.ts index b2dce5209..5d8496ad4 100644 --- a/src/authorized-view.ts +++ b/src/authorized-view.ts @@ -26,8 +26,8 @@ import { } from './row'; import {RowDataUtils, RowProperties} from './row-data-utils'; import {RawFilter} from './filter'; -import {Table} from './table'; import {Family} from './chunktransformer'; +import {Instance} from './instance'; interface FilterInformation { filter: RawFilter; @@ -56,8 +56,8 @@ interface IncrementInformation { export class AuthorizedView extends TabularApiSurface { private readonly rowData: {[id: string]: {[index: string]: Family}}; - constructor(table: Table, viewName: string) { - super(table.instance, table.id, viewName); + constructor(instance: Instance, tableName: string, viewName: string) { + super(instance, tableName, viewName); this.rowData = {}; } diff --git a/src/instance.ts b/src/instance.ts index 596b551a9..4e86e368a 100644 --- a/src/instance.ts +++ b/src/instance.ts @@ -68,6 +68,7 @@ import {Bigtable} from '.'; import {google} from '../protos/protos'; import {Backup, RestoreTableCallback, RestoreTableResponse} from './backup'; import {ClusterUtils} from './utils/cluster'; +import {AuthorizedView} from './authorized-view'; export interface ClusterInfo extends BasicClusterConfig { id: string; @@ -1475,6 +1476,16 @@ Please use the format 'my-instance' or '${bigtable.projectName}/instances/my-ins } ); } + + /** + * Gets an Authorized View object for making authorized view grpc calls. + * + * @param {string} tableName The name for the Table + * @param {string} viewName The name for the Authorized view + */ + view(tableName: string, viewName: string): AuthorizedView { + return new AuthorizedView(this, tableName, viewName); + } } /*! Developer Documentation diff --git a/src/table.ts b/src/table.ts index 16ac805a5..0adeb1aa6 100644 --- a/src/table.ts +++ b/src/table.ts @@ -1164,15 +1164,6 @@ export class Table extends TabularApiSurface { ); } - /** - * Gets an Authorized View object for making authorized view grpc calls. - * - * @param {string} viewName The name for the Authorized view - */ - view(viewName: string): AuthorizedView { - return new AuthorizedView(this, viewName); - } - waitForReplication(): Promise; waitForReplication(callback: WaitForReplicationCallback): void; /** From f099da573a1fbfe2e699c822e4f9408a802bad1c Mon Sep 17 00:00:00 2001 From: danieljbruce Date: Tue, 15 Oct 2024 13:31:22 -0400 Subject: [PATCH 5/6] feat: Create the unit tests for Authorized Views ensuring that requests are consistent (#1501) * Create a test for authorized views * Set the ReadRows options * Add authorized view test for createreadstream * Pass authorizedViewName along for createReadStream * Add the unit test for getRows * Add the unit test for getRows * Take out the test mock for the request function * Mutate rows request - test setup * Made corrections and src code changes so unit test passes * Fix the readrows test to end the stream Also add a test for insert and make the test easier to debug. * Finish the sampleRowKeys tests * Refactor part of the test for getting request opts * Authorized views for readRows test should be one assert * Add the mock function for request * Make mutate rows use the mockRequest function * Setup the sampleRowKeys setup fn to use mockReq * Add comments to the mocking functions * Fix comments. Move functions to right place * Setup the readModifyWriteRow tests * For createRules test use new mockCallbackRequest You need to call the callback in order to end the operation to createRules. * Use the new mockCallbackRequest function * Get rid of the unused mock function * Create a test file for createRules not got * increment * Add the test for increment * Change the response value for authorized view * Add the console logs and change return type * Flesh out response so that view call works * Exclude appropriate properties using promisify. * Removed console logs * Finished the filter test * Add the view call for filter * Remove only * Change to more readable values * Remove only * Add header to the test * For tests, expect view to be excluded from promis * run linter * Add a few @params and @returns * Add a comment about making requests for auth views * Added more comments - auth views vs table Also added a comment for calling increment --- src/authorized-view.ts | 5 + src/instance.ts | 1 + src/row-data-utils.ts | 1 - src/row.ts | 2 +- src/table.ts | 1 - src/tabular-api-surface.ts | 50 +++- test/authorized-views.ts | 489 +++++++++++++++++++++++++++++++++++++ test/instance.ts | 1 + 8 files changed, 536 insertions(+), 14 deletions(-) create mode 100644 test/authorized-views.ts diff --git a/src/authorized-view.ts b/src/authorized-view.ts index 5d8496ad4..d0faef7b4 100644 --- a/src/authorized-view.ts +++ b/src/authorized-view.ts @@ -28,6 +28,7 @@ import {RowDataUtils, RowProperties} from './row-data-utils'; import {RawFilter} from './filter'; import {Family} from './chunktransformer'; import {Instance} from './instance'; +import {promisifyAll} from '@google-cloud/promisify'; interface FilterInformation { filter: RawFilter; @@ -253,3 +254,7 @@ export class AuthorizedView extends TabularApiSurface { } } } + +promisifyAll(AuthorizedView, { + exclude: ['initializeRow', 'generateProperties'], +}); diff --git a/src/instance.ts b/src/instance.ts index 4e86e368a..7d5569bd5 100644 --- a/src/instance.ts +++ b/src/instance.ts @@ -1501,6 +1501,7 @@ promisifyAll(Instance, { 'getBackupsStream', 'getTablesStream', 'getAppProfilesStream', + 'view', ], }); diff --git a/src/row-data-utils.ts b/src/row-data-utils.ts index d2e74a57b..d3c3e3a39 100644 --- a/src/row-data-utils.ts +++ b/src/row-data-utils.ts @@ -246,7 +246,6 @@ class RowDataUtils { callback(err, null, resp); return; } - const data = this.formatFamilies_Util(resp!.row!.families!); const value = dotProp.get(data, column.replace(':', '.'))[0].value; diff --git a/src/row.ts b/src/row.ts index de8c1e4bb..f45c26054 100644 --- a/src/row.ts +++ b/src/row.ts @@ -34,7 +34,7 @@ import {TabularApiSurface} from './tabular-api-surface'; export interface Rule { column: string; - append: string; + append?: string; // append is optional since `increment` doesn't supply it for instance. increment?: number; } export interface CreateRowOptions { diff --git a/src/table.ts b/src/table.ts index 0adeb1aa6..d1b94a231 100644 --- a/src/table.ts +++ b/src/table.ts @@ -43,7 +43,6 @@ import { GetRowsCallback, GetRowsResponse, } from './tabular-api-surface'; -import {AuthorizedView} from './authorized-view'; export { InsertRowsCallback, diff --git a/src/tabular-api-surface.ts b/src/tabular-api-surface.ts index b76782705..3e66e42e1 100644 --- a/src/tabular-api-surface.ts +++ b/src/tabular-api-surface.ts @@ -306,10 +306,19 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); // 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; + // If the viewName is provided then request will be made for an + // authorized view. Otherwise, the request is made for a table. + const reqOpts = ( + this.viewName + ? { + authorizedViewName: `${this.name}/authorizedViews/${this.viewName}`, + appProfileId: this.bigtable.appProfileId, + } + : { + tableName: this.name, + appProfileId: this.bigtable.appProfileId, + } + ) as google.bigtable.v2.IReadRowsRequest; const retryOpts = { currentRetryAttempt: 0, // was numConsecutiveErrors @@ -674,13 +683,23 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); return pendingEntryIndices.has(index); }); - const reqOpts = { - tableName: this.name, + // If the viewName is provided then request will be made for an + // authorized view. Otherwise, the request is made for a table. + const baseReqOpts = ( + this.viewName + ? { + authorizedViewName: `${this.name}/authorizedViews/${this.viewName}`, + } + : { + tableName: this.name, + } + ) as google.bigtable.v2.IReadRowsRequest; + const reqOpts = Object.assign(baseReqOpts, { appProfileId: this.bigtable.appProfileId, entries: options.rawMutation ? entryBatch : entryBatch.map(Mutation.parse), - }; + }); const retryOpts = { currentRetryAttempt: numRequestsMade, @@ -817,10 +836,19 @@ Please use the format 'prezzy' or '${instance.name}/tables/prezzy'.`); * ``` */ sampleRowKeysStream(gaxOptions?: CallOptions) { - const reqOpts = { - tableName: this.name, - appProfileId: this.bigtable.appProfileId, - }; + // If the viewName is provided then request will be made for an + // authorized view. Otherwise, the request is made for a table. + const reqOpts = ( + this.viewName + ? { + authorizedViewName: `${this.name}/authorizedViews/${this.viewName}`, + appProfileId: this.bigtable.appProfileId, + } + : { + tableName: this.name, + appProfileId: this.bigtable.appProfileId, + } + ) as google.bigtable.v2.IReadRowsRequest; const rowKeysStream = new Transform({ transform(key, enc, next) { diff --git a/test/authorized-views.ts b/test/authorized-views.ts new file mode 100644 index 000000000..dbff89a8c --- /dev/null +++ b/test/authorized-views.ts @@ -0,0 +1,489 @@ +// 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 {beforeEach, describe} from 'mocha'; +import {AbortableDuplex, Bigtable, RawFilter, RequestCallback} from '../src'; +import {PassThrough} from 'stream'; +import * as assert from 'assert'; +import {Mutation} from '../src/mutation'; +import * as mocha from 'mocha'; +import {Row} from '../src'; + +describe('Bigtable/AuthorizedViews', () => { + describe('Authorized View methods should have requests that match Table and Row requests', () => { + const bigtable = new Bigtable({}); + const fakeTableName = 'fake-table'; + const fakeInstanceName = 'fake-instance'; + const fakeViewName = 'fake-view'; + const instance = bigtable.instance(fakeInstanceName); + const table = instance.table(fakeTableName); + const view = instance.view(fakeTableName, fakeViewName); + + /** This function mocks out the request function and compares the request + * passed into it to ensure it is correct. + * + * @param done The function to call when ending the mocha test + * @param compareFn The function that maps the requestCount to the + * request that we would expect to be passed into `request`. + */ + function mockCallbackRequest( + done: mocha.Done, + compareFn: (requestCount: number) => unknown, + resp?: {} + ) { + let requestCount = 0; + table.bigtable.request = ( + config?: any, + callback?: RequestCallback + ) => { + try { + requestCount++; + delete config['retryOpts']; + assert.deepStrictEqual(config, compareFn(requestCount)); + } catch (err: unknown) { + done(err); + } + if (callback) { + callback(null, resp); + } + const stream = new PassThrough({ + objectMode: true, + }); + setImmediate(() => { + stream.end(); + }); + return stream as {} as AbortableDuplex; + }; + } + + /** + * This function gets the basic structure of the requests we would + * expect when first making a request for a table and then for an + * authorized view. + * + * @param requestCount The number of calls that have been made to request() + * + * @return The expected table or authorized view in the request that will + * reach the request options. + */ + function getBaseRequestOptions(requestCount: number) { + const requestForTable = { + tableName: `projects/{{projectId}}/instances/${fakeInstanceName}/tables/${fakeTableName}`, + }; + const requestForAuthorizedView = { + authorizedViewName: `projects/{{projectId}}/instances/${fakeInstanceName}/tables/${fakeTableName}/authorizedViews/${fakeViewName}`, + }; + return requestCount === 1 ? requestForTable : requestForAuthorizedView; + } + + describe('Table', () => { + describe('should make ReadRows grpc requests', () => { + /** + * This function mocks out the request function to expect a readRows + * request when the tests are run. + * + * @param done The function to call when ending the mocha test + */ + function setupReadRows(done: mocha.Done) { + mockCallbackRequest(done, requestCount => { + return { + client: 'BigtableClient', + method: 'readRows', + gaxOpts: { + maxRetries: 4, + otherArgs: { + headers: { + 'bigtable-attempt': 0, + }, + }, + }, + reqOpts: Object.assign( + { + appProfileId: undefined, + rows: { + rowKeys: [], + rowRanges: [ + { + startKeyClosed: Buffer.from('7'), + endKeyClosed: Buffer.from('9'), + }, + ], + }, + filter: { + columnQualifierRegexFilter: Buffer.from('abc'), + }, + rowsLimit: 5, + }, + getBaseRequestOptions(requestCount) + ), + }; + }); + } + + it('requests for createReadStream should match', done => { + setupReadRows(done); + (async () => { + const opts = { + decode: true, + end: '9', + filter: [{column: 'abc'}], + gaxOptions: { + maxRetries: 4, + }, + limit: 5, + start: '7', + }; + await table.createReadStream(opts); + await view.createReadStream(opts); + done(); + })(); + }); + it('requests for getRows should match', done => { + setupReadRows(done); + (async () => { + const opts = { + decode: true, + end: '9', + filter: [{column: 'abc'}], + gaxOptions: { + maxRetries: 4, + }, + limit: 5, + start: '7', + }; + await table.getRows(opts); + await view.getRows(opts); + done(); + })(); + }); + }); + describe('should make MutateRows grpc requests', () => { + /** + * This function mocks out the request function to expect a mutateRows + * request when the tests are run. + * + * @param done The function to call when ending the mocha test + */ + function setupMutateRows(done: mocha.Done) { + mockCallbackRequest(done, requestCount => { + return { + client: 'BigtableClient', + method: 'mutateRows', + gaxOpts: { + maxRetries: 0, + otherArgs: { + headers: { + ['bigtable-attempt']: 0, + }, + }, + }, + reqOpts: Object.assign(getBaseRequestOptions(requestCount), { + appProfileId: undefined, + entries: [ + { + rowKey: Buffer.from('some-id'), + mutations: [ + { + setCell: { + familyName: 'follows', + columnQualifier: + Mutation.convertToBytes('tjefferson'), + timestampMicros: 2, + value: Mutation.convertToBytes(1), + }, + }, + ], + }, + ], + }), + }; + }); + } + it('requests for mutate should match', done => { + (async () => { + setupMutateRows(done); + const mutation = { + key: 'some-id', + data: { + follows: { + tjefferson: { + value: 1, + timestamp: 2, + }, + }, + }, + method: Mutation.methods.INSERT, + }; + // Currently the client retries on an end event if all the data + // hasn't been sent back so maxRetries needs to be 0. + const gaxOptions = {maxRetries: 0}; + table.maxRetries = 0; + await table.mutate(mutation, {gaxOptions}); + view.maxRetries = 0; + await view.mutate(mutation, {gaxOptions}); + done(); + })(); + }); + it('requests for insert should match', done => { + (async () => { + setupMutateRows(done); + const mutation = { + key: 'some-id', + data: { + follows: { + tjefferson: { + value: 1, + timestamp: 2, + }, + }, + }, + }; + // Currently the client retries on an end event if all the data + // hasn't been sent back so maxRetries needs to be 0. + const gaxOptions = {maxRetries: 0}; + table.maxRetries = 0; + await table.insert(mutation, gaxOptions); + view.maxRetries = 0; + await view.insert(mutation, gaxOptions); + done(); + })(); + }); + }); + describe('should make SampleRowKeys grpc requests', () => { + /** + * This function mocks out the request function to expect a sampleRowKeys + * request when the tests are run. + * + * @param done The function to call when ending the mocha test + */ + function setupSampleRowKeys(done: mocha.Done) { + mockCallbackRequest(done, requestCount => { + return { + client: 'BigtableClient', + method: 'sampleRowKeys', + gaxOpts: { + maxRetries: 4, + }, + reqOpts: Object.assign(getBaseRequestOptions(requestCount), { + appProfileId: undefined, + }), + }; + }); + } + it('requests for sampleRowKeys should match', done => { + setupSampleRowKeys(done); + (async () => { + const opts = { + maxRetries: 4, + }; + await table.sampleRowKeys(opts); + await view.sampleRowKeys(opts); + done(); + })(); + }); + it('requests for sampleRowKeysStream should match', done => { + setupSampleRowKeys(done); + (async () => { + const gaxOptions = {maxRetries: 4}; + await table.sampleRowKeysStream(gaxOptions); + await view.sampleRowKeysStream(gaxOptions); + done(); + })(); + }); + }); + }); + describe('Row', () => { + const rowId = 'row-id'; + let row: Row; + + beforeEach(() => { + row = table.row(rowId); + }); + + describe('should make readModifyWriteRow grpc requests', () => { + /** + * This function mocks out the request function to expect a readRows + * request when the tests are run. + * + * @param done The function to call when ending the mocha test + */ + function setupReadModifyWriteRow(done: mocha.Done) { + mockCallbackRequest( + done, + requestCount => { + return { + client: 'BigtableClient', + method: 'readModifyWriteRow', + gaxOpts: { + maxRetries: 4, + }, + reqOpts: Object.assign( + { + appProfileId: undefined, + rowKey: Buffer.from(rowId), + rules: [ + { + familyName: 'columnFamilyName', + columnQualifier: Buffer.from('columnName'), + incrementAmount: 7, + }, + ], + }, + getBaseRequestOptions(requestCount) + ), + }; + }, + { + row: { + families: [ + { + name: 'columnFamilyName', + columns: [ + { + qualifier: Buffer.from('columnName'), + cells: [ + { + labels: [], + timestampMicros: '4', + value: Mutation.convertToBytes(7), + }, + ], + }, + ], + }, + ], + }, + } + ); + } + + it('requests for createRules should match', done => { + setupReadModifyWriteRow(done); + (async () => { + const rule = { + column: 'columnFamilyName:columnName', + increment: 7, + }; + const gaxOpts = {maxRetries: 4}; + await row.createRules(rule, gaxOpts); + await view.createRules({rules: rule, rowId: rowId}, gaxOpts); + done(); + })(); + }); + it('requests for increment should match', done => { + setupReadModifyWriteRow(done); + (async () => { + // Change the response so that format families can run. + const column = 'columnFamilyName:columnName'; + const gaxOpts = {maxRetries: 4}; + await row.increment(column, 7, gaxOpts); + await view.increment({column, rowId}, 7, gaxOpts); + done(); + })(); + }); + }); + describe('should make checkAndMutateRequest grpc requests', () => { + /** + * This function mocks out the request function to expect a readRows + * request when the tests are run. + * + * @param done The function to call when ending the mocha test + */ + function setupCheckAndMutateRow(done: mocha.Done) { + mockCallbackRequest( + done, + requestCount => { + return { + client: 'BigtableClient', + method: 'checkAndMutateRow', + gaxOpts: { + maxRetries: 4, + }, + reqOpts: Object.assign( + { + appProfileId: undefined, + rowKey: Buffer.from(rowId), + predicateFilter: { + familyNameRegexFilter: 'columnFamilyName', + }, + trueMutations: [ + { + deleteFromColumn: { + familyName: 'columnFamilyName', + columnQualifier: Buffer.from('columnName'), + timeRange: undefined, + }, + }, + ], + falseMutations: [], + }, + getBaseRequestOptions(requestCount) + ), + }; + }, + { + row: { + families: [ + { + name: 'columnFamilyName', + columns: [ + { + qualifier: Buffer.from('columnName'), + cells: [ + { + labels: [], + timestampMicros: '4', + value: Mutation.convertToBytes(7), + }, + ], + }, + ], + }, + ], + }, + } + ); + } + + it('requests for filter should match', done => { + setupCheckAndMutateRow(done); + (async () => { + const filter: RawFilter = { + family: 'columnFamilyName', + value: 'columnName', + }; + const mutations = [ + { + method: 'delete', + data: ['columnFamilyName:columnName'], + }, + ]; + await row.filter(filter, { + onMatch: mutations, + gaxOptions: {maxRetries: 4}, + }); + await view.filter( + {filter, rowId}, + { + onMatch: mutations, + gaxOptions: {maxRetries: 4}, + } + ); + done(); + })(); + }); + }); + }); + }); +}); diff --git a/test/instance.ts b/test/instance.ts index 32c10c31c..66190c55e 100644 --- a/test/instance.ts +++ b/test/instance.ts @@ -55,6 +55,7 @@ const fakePromisify = Object.assign({}, promisify, { 'getBackupsStream', 'getTablesStream', 'getAppProfilesStream', + 'view', ]); }, }); From 530f6da7915cc76099586555e1f369d69b73e9d8 Mon Sep 17 00:00:00 2001 From: danieljbruce Date: Wed, 23 Oct 2024 09:41:12 -0400 Subject: [PATCH 6/6] feat: Bigtable authorized views integration tests (#1504) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * Create a test for authorized views * Set the ReadRows options * Add authorized view test for createreadstream * Pass authorizedViewName along for createReadStream * Add the unit test for getRows * Add the unit test for getRows * Take out the test mock for the request function * Mutate rows request - test setup * Made corrections and src code changes so unit test passes * Fix the readrows test to end the stream Also add a test for insert and make the test easier to debug. * Finish the sampleRowKeys tests * Refactor part of the test for getting request opts * Authorized views for readRows test should be one assert * Add the mock function for request * Make mutate rows use the mockRequest function * Setup the sampleRowKeys setup fn to use mockReq * Add comments to the mocking functions * Fix comments. Move functions to right place * Setup the readModifyWriteRow tests * For createRules test use new mockCallbackRequest You need to call the callback in order to end the operation to createRules. * Use the new mockCallbackRequest function * Get rid of the unused mock function * Create a test file for createRules not got * increment * Add the test for increment * Change the response value for authorized view * Add the console logs and change return type * Flesh out response so that view call works * Exclude appropriate properties using promisify. * Removed console logs * Finished the filter test * Add the view call for filter * Remove only * Change to more readable values * Remove only * Add header to the test * For tests, expect view to be excluded from promis * run linter * Create the before hook for auth views * Note to self on auth views - TODO * Add an insert statement * Make a test for getRows * Add the test for the getRows function. * Finish ‘should fail when writing to a row not in v * Add another mutate test, preserve table after test Add after hook to make sure table values stay the same. Add test for modifying from different column * Remove TODO * refactor the error message in the test * Add mutate, insert and sampleRowKeys tests * Fix sampleRowKeys test to require a fix upstream * Create samplerowkeys and createreadstream tests * Surround all tests in a try block This allows for better error reporting * Fixed the column filter so it works * Reduce verbosity in tests In many cases, 3 lines can be reduced to one variable. This is crucial for making the integration tests shorter. * Replace verbose data with data variable * Add an insert for multiple rows * Fix the check and mutate test to ignore corrupt d * Add the readModifyWriteRow tests * Add awaits * wrap columnFamily and columnIdInView in brackets * Eliminate the TODO * Remove only * Add key back to filter config option --- src/row.ts | 1 + src/table.ts | 2 +- system-test/bigtable.ts | 556 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 558 insertions(+), 1 deletion(-) diff --git a/src/row.ts b/src/row.ts index f45c26054..07aac89ca 100644 --- a/src/row.ts +++ b/src/row.ts @@ -58,6 +58,7 @@ export interface Family { export interface FilterConfigOption { method?: string; data?: Data; + key?: string; } export interface FilterConfig { gaxOptions?: CallOptions; diff --git a/src/table.ts b/src/table.ts index d1b94a231..84989206a 100644 --- a/src/table.ts +++ b/src/table.ts @@ -260,7 +260,7 @@ export type SampleRowKeysCallback = ( err: ServiceError | null, keys?: string[] ) => void; -export type SampleRowsKeysResponse = [string[]]; +export type SampleRowsKeysResponse = [{key: Uint8Array; offset: string}[]]; export type DeleteRowsCallback = ( err: ServiceError | null, apiResponse?: google.protobuf.Empty diff --git a/system-test/bigtable.ts b/system-test/bigtable.ts index 10e70e6bc..0985be7bf 100644 --- a/system-test/bigtable.ts +++ b/system-test/bigtable.ts @@ -22,9 +22,12 @@ import { Backup, BackupTimestamp, Bigtable, + Entry, Instance, InstanceOptions, + MutateOptions, } from '../src'; +import {Mutation} from '../src/mutation'; import {AppProfile} from '../src/app-profile.js'; import {CopyBackupConfig} from '../src/backup.js'; import {Cluster} from '../src/cluster.js'; @@ -33,6 +36,8 @@ import {Row} from '../src/row.js'; import {Table} from '../src/table.js'; import {RawFilter} from '../src/filter'; import {generateId, PREFIX} from './common'; +import {BigtableTableAdminClient} from '../src/v2'; +import {ServiceError} from 'google-gax'; describe('Bigtable', () => { const bigtable = new Bigtable(); @@ -1712,6 +1717,557 @@ describe('Bigtable', () => { }); }); }); + describe('AuthorizedViews', () => { + const tableId = generateId('table'); + const familyName = generateId('column-family-name'); + const rowId = generateId('row-id'); + const otherRowId = generateId('row-id'); + const authorizedViewId = generateId('authorized-view-id'); + const columnIdInView = generateId('column-id'); + const columnIdNotInView = generateId('column-id'); + const cellValueInView = generateId('cell-value'); + const cellValueInView2 = generateId('cell-value'); + const cellValueNotInView = generateId('cell-value'); + const newCellValue = generateId('cell-value'); + const authorizedViewTable = INSTANCE.table(tableId); + const authorizedView = INSTANCE.view(tableId, authorizedViewId); + const columnIdInViewData = { + value: cellValueInView, + labels: [], + timestamp: '77000', + }; + const columnIdInViewData2 = { + value: cellValueInView2, + labels: [], + timestamp: '77000', + }; + const columnIdNotInViewData = { + value: cellValueNotInView, + labels: [], + timestamp: '77000', + }; + const newCellValueData = { + value: newCellValue, + labels: [], + timestamp: '77000', + }; + let authorizedViewTableFullName: string; + let authorizedViewFullName: string; + + /** + * getErrorMessage gets the error message that the test would expect + * for a particular authorizedViewFullName. + * + * @param authorizedViewFullName The full name of the authorized view. + * This method should only be called after authorizedViewFullName has + * been initialized. + */ + function getErrorMessage(authorizedViewFullName: string) { + return `Cannot mutate from ${authorizedViewFullName} because the mutation contains cells outside the Authorized View.`; + } + + /** + * Resets the table to a stable state after running a test. + * + */ + async function resetTable() { + // Change the cell value back to what it was: + await authorizedViewTable.deleteRows(rowId); + const firstMutation = { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: columnIdInViewData, + [columnIdNotInView]: columnIdNotInViewData, + }, + }, + } as {} as Entry; + await authorizedViewTable.insert(firstMutation, {}); + } + + /** + * Creates the table used for the tests. + */ + async function createTable() { + { + // Create a table with just one row. + await authorizedViewTable.create({}); + await authorizedViewTable.createFamily(familyName); + await authorizedViewTable.insert([ + { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: columnIdInViewData, + [columnIdNotInView]: columnIdNotInViewData, + }, + }, + }, + { + key: otherRowId, + data: { + [familyName]: { + [columnIdInView]: columnIdInViewData, + }, + }, + }, + ]); + // The following operations must be performed after table.insert because bigtable.projectId needs to be assigned. + authorizedViewTableFullName = authorizedViewTable.name.replace( + '{{projectId}}', + bigtable.projectId + ); + authorizedViewFullName = `${authorizedViewTableFullName}/authorizedViews/${authorizedViewId}`; + } + { + // Create an authorized view that the integration tests can use. + // The view should only see the columnIdInView column. + const bigtableClient = bigtable.api[ + 'BigtableTableAdminClient' + ] as BigtableTableAdminClient; + await bigtableClient.createAuthorizedView({ + parent: authorizedViewTableFullName, + authorizedViewId, + authorizedView: { + etag: `${authorizedViewId}-etag`, + deletionProtection: false, + subsetView: { + rowPrefixes: [Buffer.from(rowId)], + familySubsets: { + [familyName]: { + qualifiers: [Buffer.from(columnIdInView)], + }, + }, + }, + }, + }); + } + } + + before(async () => { + await createTable(); + }); + + afterEach(async () => { + // Add an after hook to ensure that none of the tests change the table. + const rows = (await authorizedViewTable.getRows())[0]; + assert.strictEqual(rows.length, 2); + assert.strictEqual(rows[0].id, rowId); + assert.deepStrictEqual(rows[0].data, { + [familyName]: { + [columnIdInView]: [columnIdInViewData], + [columnIdNotInView]: [columnIdNotInViewData], + }, + }); + assert.strictEqual(rows[1].id, otherRowId); + assert.deepStrictEqual(rows[1].data, { + [familyName]: { + [columnIdInView]: [columnIdInViewData], + }, + }); + }); + + describe('ReadRows grpc calls', () => { + it('should call getRows for the authorized view', async () => { + const rows = (await authorizedView.getRows())[0]; + // The getRows call will only get one of the rows and only display + // one of the columns visible in the view. + assert.strictEqual(rows[0].id, rowId); + assert.deepStrictEqual(rows[0].data, { + [familyName]: { + [columnIdInView]: [columnIdInViewData], + }, + }); + }); + it('should call createReadStream for the authorized view', done => { + (async () => { + try { + const stream = await authorizedView.createReadStream(); + let receivedDataCount = 0; + stream.on('data', row => { + assert.strictEqual(row.id, rowId); + assert.deepStrictEqual(row.data, { + [familyName]: { + [columnIdInView]: [columnIdInViewData], + }, + }); + receivedDataCount = receivedDataCount + 1; + }); + stream.on('error', () => { + done('An error should not have occurred'); + }); + stream.on('end', () => { + assert.strictEqual(receivedDataCount, 1); + done(); + }); + } catch (e: unknown) { + done(e); + } + })(); + }); + }); + describe('MutateRows grpc calls', () => { + describe('For erroneous calls', () => { + it('should fail when writing to a row not in the authorized view', async () => { + const mutation = { + key: otherRowId, + data: { + [familyName]: { + [columnIdInView]: newCellValue, + }, + }, + method: Mutation.methods.INSERT, + } as {} as Entry; + try { + await authorizedView.mutate(mutation, {} as MutateOptions); + assert.fail('The mutate call should have failed'); + } catch (e: unknown) { + assert.strictEqual( + (e as ServiceError).message, + getErrorMessage(authorizedViewFullName) + ); + } + }); + it('should fail when writing to a column not in the authorized view', async () => { + const mutation = { + key: rowId, + data: { + [familyName]: { + [columnIdNotInView]: newCellValue, + }, + }, + method: Mutation.methods.INSERT, + } as {} as Entry; + try { + await authorizedView.mutate(mutation, {} as MutateOptions); + assert.fail('The mutate call should have failed'); + } catch (e: unknown) { + assert.strictEqual( + (e as ServiceError).message, + getErrorMessage(authorizedViewFullName) + ); + } + }); + }); + it('should mutate a row for a row/column in view', async () => { + // Change the cell in view to a new value. + const firstMutation = { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: newCellValueData, + }, + }, + method: Mutation.methods.INSERT, + } as {} as Entry; + await authorizedView.mutate(firstMutation, {} as MutateOptions); + // Ensure the new cell value change took place + const rows = (await authorizedView.getRows())[0]; + assert.strictEqual(rows[0].id, rowId); + assert.deepStrictEqual(rows[0].data, { + [familyName]: { + [columnIdInView]: [newCellValueData], + }, + }); + // Change the cell value back to what it was before + const secondMutation = { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: columnIdInViewData, + }, + }, + method: Mutation.methods.INSERT, + } as {} as Entry; + await authorizedView.mutate(secondMutation, {} as MutateOptions); + }); + it('should insert a row for a row/column in view', async () => { + // Change the cell in view to a new value. + const firstMutation = { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: newCellValueData, + }, + }, + } as {} as Entry; + await authorizedView.insert(firstMutation, {}); + // Ensure the new cell value change took place + const rows = (await authorizedView.getRows())[0]; + assert.strictEqual(rows[0].id, rowId); + assert.deepStrictEqual(rows[0].data, { + [familyName]: { + [columnIdInView]: [newCellValueData], + }, + }); + // Change the cell value back to what it was before + const secondMutation = { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: columnIdInViewData, + }, + }, + method: Mutation.methods.INSERT, + } as {} as Entry; + await authorizedView.insert(secondMutation, {}); + }); + }); + describe('SampleRowKeys grpc calls', () => { + /** + * This function is for converting a buffer to an integer equal to the + * total value of the buffer. It is useful for comparing the sampleRowKeys + * return value to the row identifier since the difference between the + * total value of these buffers is expected to be exactly 1. + * + * @param buffer The buffer being mapped to be used for comparisons. + */ + function convertBufferToInt(buffer: Uint8Array) { + return buffer + .reverse() + .reduce( + (accumulator, currentValue, index) => + accumulator + Math.pow(currentValue, index), + 0 + ); + } + + it('should get a sample of row keys', async () => { + const rowKeys = await authorizedView.sampleRowKeys(); + assert.strictEqual(rowKeys.length, 1); + assert.strictEqual(rowKeys[0].length, 1); + assert.deepStrictEqual( + convertBufferToInt(rowKeys[0][0].key), + convertBufferToInt(Buffer.from(rowId)) + 1 + ); + }); + it('should call sampleRowKeysStream for the authorized view', done => { + (async () => { + try { + const stream = await authorizedView.sampleRowKeysStream(); + let receivedDataCount = 0; + stream.on('data', (row: {key: Uint8Array; offset: string}) => { + assert.deepStrictEqual( + convertBufferToInt(row.key), + convertBufferToInt(Buffer.from(rowId)) + 1 + ); + receivedDataCount = receivedDataCount + 1; + }); + stream.on('error', () => { + done('An error should not have occurred'); + }); + stream.on('end', () => { + assert.strictEqual(receivedDataCount, 1); + done(); + }); + } catch (e: unknown) { + done(e); + } + })(); + }); + }); + describe('CheckAndMutate grpc calls', () => { + it('should error when the request is made for the row key not in a view', done => { + (async () => { + try { + try { + await authorizedView.filter( + { + rowId: 'some-row-key', + filter: { + column: columnIdInView, + }, + }, + { + onMatch: [ + { + key: rowId, + method: 'delete', + }, + ], + } + ); + done('The call to filter should have failed.'); + } catch (e: unknown) { + assert.strictEqual( + (e as ServiceError).details, + getErrorMessage(authorizedViewFullName) + ); + done(); + } + } catch (e: unknown) { + // Will reach this point if there is an assertion error. + done(e); + } + })(); + }); + it('should call filter for the authorized view', done => { + (async () => { + try { + // Add the row so that the cell offset filter takes effect: + await authorizedViewTable.insert([ + { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: [columnIdInViewData, columnIdInViewData2], + }, + }, + }, + ]); + const rowsAfterAddition = (await authorizedViewTable.getRows())[0]; + assert.strictEqual(rowsAfterAddition.length, 2); + assert.strictEqual(rowsAfterAddition[0].id, rowId); + assert.deepStrictEqual( + rowsAfterAddition[0].data[familyName][columnIdNotInView].length, + 1 + ); + assert.deepStrictEqual( + rowsAfterAddition[0].data[familyName][columnIdInView].length, + 2 + ); + assert.strictEqual(rowsAfterAddition[1].id, otherRowId); + assert.deepStrictEqual(rowsAfterAddition[1].data, { + [familyName]: { + [columnIdInView]: [columnIdInViewData], + }, + }); + // Call filter to conduct a checkAndMutate operation. + const mutations = [ + { + method: 'delete', + data: [`${familyName}:${columnIdInView}`], + }, + ]; + await authorizedView.filter( + { + rowId: rowId, + filter: { + row: { + cellOffset: 1, + }, + }, + }, + { + onMatch: mutations, + } + ); + // Check the rows to ensure the row was deleted by calling `filter`. + const rows = (await authorizedViewTable.getRows())[0]; + assert.strictEqual(rows.length, 2); + assert.strictEqual(rows[0].id, rowId); + assert.deepStrictEqual(rows[0].data, { + [familyName]: { + [columnIdNotInView]: [columnIdNotInViewData], + // [columnIdInView] is deleted by checkAndMutate + }, + }); + assert.strictEqual(rows[1].id, otherRowId); + assert.deepStrictEqual(rows[1].data, { + [familyName]: { + [columnIdInView]: [columnIdInViewData], + }, + }); + // Add the row that was deleted back: + await authorizedViewTable.insert([ + { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: { + value: cellValueInView, + labels: [], + timestamp: 77000, + }, + }, + }, + }, + ]); + done(); + } catch (e: unknown) { + done(e); + } + })(); + }); + }); + describe('ReadModifyWriteRow grpc calls', () => { + it('should apply read/modify/write rules to a row', async () => { + // Append a value to the table: + const rule = { + column: `${familyName}:${columnIdInView}`, + append: '-appended-value', + }; + await authorizedView.createRules({ + rowId, + rules: rule, + }); + // Check that the operation was performed correctly: + const rows = (await authorizedView.getRows())[0]; + rows[0].data[familyName][columnIdInView][0].timestamp = '77000'; + assert.strictEqual(rows.length, 1); + assert.strictEqual(rows[0].id, rowId); + assert.deepStrictEqual(rows[0].data, { + [familyName]: { + [columnIdInView]: [ + { + value: `${cellValueInView}-appended-value`, + labels: [], + timestamp: '77000', + }, + columnIdInViewData, + ], + }, + }); + await resetTable(); + }); + it('should apply increment to a row', async () => { + // First set the row in view cell value to a numeric value: + const originalValue = Math.floor(Math.random() * 1000000000); + await authorizedViewTable.deleteRows(rowId); + const firstMutation = { + key: rowId, + data: { + [familyName]: { + [columnIdInView]: { + value: originalValue, + labels: [], + timestamp: '77000', + }, + [columnIdNotInView]: columnIdNotInViewData, + }, + }, + } as {} as Entry; + await authorizedViewTable.insert(firstMutation, {}); + // Next, increment the value: + await authorizedView.increment( + { + rowId, + column: `${familyName}:${columnIdInView}`, + }, + 1 + ); + const rows = (await authorizedView.getRows())[0]; + rows[0].data[familyName][columnIdInView][0].timestamp = '77000'; + assert.deepStrictEqual(rows[0].data, { + [familyName]: { + [columnIdInView]: [ + { + value: originalValue + 1, + labels: [], + timestamp: '77000', + }, + { + value: originalValue, + labels: [], + timestamp: '77000', + }, + ], + }, + }); + await resetTable(); + }); + }); + }); }); function createInstanceConfig(