Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
feat: Introduce CubeStoreCacheDriver (#5511)
- Loading branch information
Showing
30 changed files
with
369 additions
and
92 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
File renamed without changes.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,3 +1,5 @@ | ||
export * from './BaseDriver'; | ||
export * from './utils'; | ||
export * from './driver.interface'; | ||
export * from './queue-driver.interface'; | ||
export * from './cache-driver.interface'; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,23 @@ | ||
export interface QueueDriverConnectionInterface { | ||
getResultBlocking(queryKey: string): Promise<unknown>; | ||
getResult(queryKey: string): Promise<unknown>; | ||
addToQueue(queryKey: string): Promise<unknown>; | ||
getToProcessQueries(): Promise<unknown>; | ||
getActiveQueries(): Promise<unknown>; | ||
getOrphanedQueries(): Promise<unknown>; | ||
getStalledQueries(): Promise<unknown>; | ||
getQueryStageState(onlyKeys: any): Promise<unknown>; | ||
updateHeartBeat(queryKey: string): Promise<void>; | ||
getNextProcessingId(): Promise<string>; | ||
retrieveForProcessing(queryKey: string, processingId: string): Promise<unknown>; | ||
freeProcessingLock(queryKe: string, processingId: string, activated: unknown): Promise<unknown>; | ||
optimisticQueryUpdate(queryKey, toUpdate, processingId): Promise<unknown>; | ||
cancelQuery(queryKey: string): Promise<unknown>; | ||
setResultAndRemoveQuery(queryKey: string, executionResult: any, processingId: any): Promise<unknown>; | ||
release(): Promise<void>; | ||
} | ||
|
||
export interface QueueDriverInterface { | ||
createConnection(): Promise<QueueDriverConnectionInterface>; | ||
release(connection: QueueDriverConnectionInterface): Promise<void>; | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,17 +1,15 @@ | ||
const fromExports = require('./dist/src'); | ||
const { CubeStoreDriver } = require('./dist/src/CubeStoreDriver'); | ||
const { CubeStoreDevDriver } = require('./dist/src/CubeStoreDevDriver'); | ||
const { isCubeStoreSupported, CubeStoreHandler } = require('./dist/src/rexport'); | ||
|
||
/** | ||
* After 5 years working with TypeScript, now I know | ||
* that commonjs and nodejs require is not compatibility with using export default | ||
*/ | ||
module.exports = CubeStoreDriver; | ||
const toExport = CubeStoreDriver; | ||
|
||
/** | ||
* It's needed to move our CLI to destructing style on import | ||
* Please sync this file with src/index.ts | ||
*/ | ||
module.exports.CubeStoreDevDriver = CubeStoreDevDriver; | ||
module.exports.isCubeStoreSupported = isCubeStoreSupported; | ||
module.exports.CubeStoreHandler = CubeStoreHandler; | ||
// eslint-disable-next-line no-restricted-syntax | ||
for (const [key, module] of Object.entries(fromExports)) { | ||
toExport[key] = module; | ||
} | ||
|
||
module.exports = toExport; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
90 changes: 90 additions & 0 deletions
90
packages/cubejs-cubestore-driver/src/CubeStoreCacheDriver.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,90 @@ | ||
import { createCancelablePromise, MaybeCancelablePromise } from '@cubejs-backend/shared'; | ||
import { CacheDriverInterface } from '@cubejs-backend/base-driver'; | ||
|
||
import { CubeStoreDriver } from './CubeStoreDriver'; | ||
|
||
export class CubeStoreCacheDriver implements CacheDriverInterface { | ||
public constructor( | ||
protected readonly connection: CubeStoreDriver | ||
) {} | ||
|
||
public withLock = ( | ||
key: string, | ||
cb: () => MaybeCancelablePromise<any>, | ||
expiration: number = 60, | ||
freeAfter: boolean = true, | ||
) => createCancelablePromise(async (tkn) => { | ||
if (tkn.isCanceled()) { | ||
return false; | ||
} | ||
|
||
const rows = await this.connection.query('CACHE SET NX TTL ? ? ?', [expiration, key, '1']); | ||
if (rows && rows.length === 1 && rows[0]?.success === 'true') { | ||
if (tkn.isCanceled()) { | ||
if (freeAfter) { | ||
await this.connection.query('CACHE REMOVE ?', [ | ||
key | ||
]); | ||
} | ||
|
||
return false; | ||
} | ||
|
||
try { | ||
await tkn.with(cb()); | ||
} finally { | ||
if (freeAfter) { | ||
await this.connection.query('CACHE REMOVE ?', [ | ||
key | ||
]); | ||
} | ||
} | ||
|
||
return true; | ||
} | ||
|
||
return false; | ||
}); | ||
|
||
public async get(key: string) { | ||
const rows = await this.connection.query('CACHE GET ?', [ | ||
key | ||
]); | ||
if (rows && rows.length === 1) { | ||
return JSON.parse(rows[0].value); | ||
} | ||
|
||
return null; | ||
} | ||
|
||
public async set(key: string, value, expiration) { | ||
const strValue = JSON.stringify(value); | ||
await this.connection.query('CACHE SET TTL ? ? ?', [expiration, key, strValue]); | ||
|
||
return { | ||
key, | ||
bytes: Buffer.byteLength(strValue), | ||
}; | ||
} | ||
|
||
public async remove(key: string) { | ||
await this.connection.query('CACHE REMOVE ?', [ | ||
key | ||
]); | ||
} | ||
|
||
public async keysStartingWith(prefix: string) { | ||
const rows = await this.connection.query('CACHE KEYS ?', [ | ||
prefix | ||
]); | ||
return rows.map((row) => row.key); | ||
} | ||
|
||
public async cleanup(): Promise<void> { | ||
// | ||
} | ||
|
||
public async testConnection(): Promise<void> { | ||
return this.connection.testConnection(); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,4 +1,4 @@ | ||
export * from './CubeStoreQuery'; | ||
export * from './CubeStoreCacheDriver'; | ||
export * from './CubeStoreDriver'; | ||
export * from './CubeStoreDevDriver'; | ||
export * from './rexport'; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1 +1,2 @@ | ||
dist | ||
.cubestore |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
7 changes: 6 additions & 1 deletion
7
packages/cubejs-query-orchestrator/src/orchestrator/BaseQueueDriver.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,7 +1,12 @@ | ||
import { QueueDriverConnectionInterface, QueueDriverInterface } from '@cubejs-backend/base-driver'; | ||
import { getCacheHash } from './utils'; | ||
|
||
export abstract class BaseQueueDriver { | ||
export abstract class BaseQueueDriver implements QueueDriverInterface { | ||
public redisHash(queryKey) { | ||
return getCacheHash(queryKey); | ||
} | ||
|
||
abstract createConnection(): Promise<QueueDriverConnectionInterface>; | ||
|
||
abstract release(connection: QueueDriverConnectionInterface): Promise<void>; | ||
} |
3 changes: 1 addition & 2 deletions
3
packages/cubejs-query-orchestrator/src/orchestrator/LocalCacheDriver.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.