From 0b9f9b23358e882e555d88d64d2f5e84b170b076 Mon Sep 17 00:00:00 2001 From: devlikepro Date: Fri, 10 Jan 2025 10:56:56 +0700 Subject: [PATCH] [core] Add GOWS engine --- .dockerignore | 1 + .eslintignore | 1 + .github/workflows/build.yaml | 20 +- .gitignore | 1 + .pre-commit-config.yaml | 1 + Dockerfile | 67 +- Makefile | 19 + package.json | 5 + src/core/abc/EngineBootstrap.ts | 15 + src/core/abc/manager.abc.ts | 28 +- src/core/abc/session.abc.ts | 6 +- src/core/app.module.core.ts | 2 + src/core/config/GowsEngineConfigService.ts | 26 + src/core/engines/gows/EventsFromObservable.ts | 31 + src/core/engines/gows/GowsBootstrap.ts | 72 ++ .../engines/gows/GowsEventStreamObservable.ts | 69 ++ src/core/engines/gows/GowsSubprocess.ts | 98 ++ src/core/engines/gows/session.gows.core.ts | 890 ++++++++++++++++++ src/core/engines/gows/store/GowsAuth.ts | 4 + .../engines/gows/store/GowsAuthFactoryCore.ts | 26 + src/core/engines/gows/store/GowsAuthSimple.ts | 19 + src/core/engines/gows/types.ts | 24 + src/core/engines/noweb/session.noweb.core.ts | 8 +- src/core/manager.core.ts | 19 +- src/structures/enums.dto.ts | 1 + src/structures/webhooks.dto.ts | 2 + src/utils/reactive/ops/onlyEvent.ts | 14 + yarn.lock | 247 ++++- 28 files changed, 1699 insertions(+), 17 deletions(-) create mode 100644 .eslintignore create mode 100644 src/core/abc/EngineBootstrap.ts create mode 100644 src/core/config/GowsEngineConfigService.ts create mode 100644 src/core/engines/gows/EventsFromObservable.ts create mode 100644 src/core/engines/gows/GowsBootstrap.ts create mode 100644 src/core/engines/gows/GowsEventStreamObservable.ts create mode 100644 src/core/engines/gows/GowsSubprocess.ts create mode 100644 src/core/engines/gows/session.gows.core.ts create mode 100644 src/core/engines/gows/store/GowsAuth.ts create mode 100644 src/core/engines/gows/store/GowsAuthFactoryCore.ts create mode 100644 src/core/engines/gows/store/GowsAuthSimple.ts create mode 100644 src/core/engines/gows/types.ts create mode 100644 src/utils/reactive/ops/onlyEvent.ts diff --git a/.dockerignore b/.dockerignore index 6489c7b3..a2fbd4a0 100644 --- a/.dockerignore +++ b/.dockerignore @@ -1,4 +1,5 @@ Dockerfile +src/core/engines/gows/grpc .*sessions sessions files diff --git a/.eslintignore b/.eslintignore new file mode 100644 index 00000000..1ac85cf1 --- /dev/null +++ b/.eslintignore @@ -0,0 +1 @@ +src/core/engines/gows/grpc/* diff --git a/.github/workflows/build.yaml b/.github/workflows/build.yaml index d8da0bd4..1eabcbab 100644 --- a/.github/workflows/build.yaml +++ b/.github/workflows/build.yaml @@ -45,7 +45,7 @@ jobs: # browser: "chrome" # goss: 'goss-linux-arm' - # No browser - x86 + # NOWEB - x86 - runner: 'buildjet-4vcpu-ubuntu-2204' tag: 'noweb' platform: 'amd64' @@ -53,7 +53,7 @@ jobs: engine: 'NOWEB' goss: 'goss-linux-amd64' - # No browser - ARM + # NOWEB - ARM - runner: 'buildjet-4vcpu-ubuntu-2204-arm' tag: 'noweb-arm' platform: 'linux/arm64' @@ -61,6 +61,22 @@ jobs: engine: 'NOWEB' goss: 'goss-linux-arm' + # GOWS - x86 + - runner: 'buildjet-4vcpu-ubuntu-2204' + tag: 'gows' + platform: 'amd64' + browser: 'none' + engine: 'GOWS' + goss: 'goss-linux-amd64' + + # No browser - ARM + - runner: 'buildjet-4vcpu-ubuntu-2204-arm' + tag: 'gows-arm' + platform: 'linux/arm64' + browser: 'none' + engine: 'GOWS' + goss: 'goss-linux-arm' + steps: - name: Checkout uses: actions/checkout@v3 diff --git a/.gitignore b/.gitignore index 648b025b..ace00edc 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,4 @@ +src/core/engines/gows/grpc .env *.heapsnapshot .yarn diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml index 9f182e94..11a820d7 100644 --- a/.pre-commit-config.yaml +++ b/.pre-commit-config.yaml @@ -2,6 +2,7 @@ exclude: | (?x)^( ^src/dashboard/.*| ^src/core/engines/webjs/.*html| + ^src/core/engines/gows/grpc/.*| ^src/plus/engines/webjs/.*html| ^.yarn/.* )$ diff --git a/Dockerfile b/Dockerfile index 0fc7a0b4..2fdd1af7 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,3 +1,9 @@ +# GOWS +ARG GOWS_GITHUB_REPO=devlikeapro/gows +ARG GOWS_SHA=95ca63a1dee5faec75232a9f59a548dfd4e414ff +ARG WAHA_DASHBOARD_GITHUB_REPO=devlikeapro/dashboard +ARG WAHA_DASHBOARD_SHA=a1c839d8da94d1cade182addffdfc42fc0cf7ed3 + # # Build # @@ -5,6 +11,10 @@ ARG NODE_VERSION=22.8-bullseye FROM node:${NODE_VERSION} AS build ENV PUPPETEER_SKIP_DOWNLOAD=True +# install protoc +RUN apt-get update && \ + apt-get install protobuf-compiler -y + # npm packages WORKDIR /src COPY package.json . @@ -18,6 +28,16 @@ RUN yarn install WORKDIR /src ADD . /src RUN yarn install + +# Build node grpc client for GOWS engine +WORKDIR /gows/proto +ARG GOWS_GITHUB_REPO +ARG GOWS_SHA +RUN wget https://raw.githubusercontent.com/${GOWS_GITHUB_REPO}/${GOWS_SHA}/proto/gows.proto +WORKDIR /src + +RUN make proto-gows + RUN yarn build && find ./dist -name "*.d.ts" -delete # @@ -26,15 +46,41 @@ RUN yarn build && find ./dist -name "*.d.ts" -delete FROM node:${NODE_VERSION} AS dashboard # Download WAHA Dashboard -ENV WAHA_DASHBOARD_SHA fed4e50e88e4d26c610e3289fd0d8657cb866543 +ARG WAHA_DASHBOARD_GITHUB_REPO +ARG WAHA_DASHBOARD_SHA RUN \ - wget https://github.com/devlikeapro/dashboard/archive/${WAHA_DASHBOARD_SHA}.zip \ + wget https://github.com/${WAHA_DASHBOARD_GITHUB_REPO}/archive/${WAHA_DASHBOARD_SHA}.zip \ && unzip ${WAHA_DASHBOARD_SHA}.zip -d /tmp/dashboard \ && mkdir -p /dashboard \ && mv /tmp/dashboard/dashboard-${WAHA_DASHBOARD_SHA}/* /dashboard/ \ && rm -rf ${WAHA_DASHBOARD_SHA}.zip \ && rm -rf /tmp/dashboard/dashboard-${WAHA_DASHBOARD_SHA} +# +# GOWS +# +FROM golang:1.23-bullseye AS gows +# install protoc +RUN apt-get update && \ + apt-get install protobuf-compiler -y + +# Image processing for thumbnails +RUN apt-get update \ + && apt-get install -y libvips-dev \ + && rm -rf /var/lib/apt/lists/* + +ARG GOWS_GITHUB_REPO +ARG GOWS_SHA +RUN git clone https://github.com/${GOWS_GITHUB_REPO} gows && \ + cd gows && \ + git checkout ${GOWS_SHA} + +WORKDIR /go/gows +RUN go install google.golang.org/protobuf/cmd/protoc-gen-go@latest +RUN go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest +RUN make all + + # # Final # @@ -51,6 +97,11 @@ RUN echo "USE_BROWSER=$USE_BROWSER" # Install ffmpeg to generate previews for videos RUN apt-get update && apt-get install -y ffmpeg --no-install-recommends && rm -rf /var/lib/apt/lists/* +# Image processing for thumbnails +RUN apt-get update \ + && apt-get install -y libvips \ + && rm -rf /var/lib/apt/lists/* + # Install zip and unzip - either for chromium or chrome RUN if [ "$USE_BROWSER" = "chromium" ] || [ "$USE_BROWSER" = "chrome" ]; then \ apt-get update \ @@ -100,7 +151,13 @@ RUN if [ "$USE_BROWSER" = "chrome" ]; then \ && rm -rf /var/lib/apt/lists/*; \ fi -# Set the ENV for NOWEB docker image +# GOWS requirements +# libc6 +RUN apt-get update \ + && apt-get install -y libc6 \ + && rm -rf /var/lib/apt/lists/* + +# Set the ENV for docker image ENV WHATSAPP_DEFAULT_ENGINE=$WHATSAPP_DEFAULT_ENGINE # Attach sources, install packages @@ -109,6 +166,10 @@ COPY package.json ./ COPY --from=build /src/node_modules ./node_modules COPY --from=build /src/dist ./dist COPY --from=dashboard /dashboard ./dist/dashboard +COPY --from=gows /go/gows/bin/gows /app/gows +ENV WAHA_GOWS_PATH /app/gows +ENV WAHA_GOWS_SOCKET /tmp/gows.sock + COPY entrypoint.sh /entrypoint.sh # Chokidar options to monitor file changes diff --git a/Makefile b/Makefile index b7422bfc..ff216922 100644 --- a/Makefile +++ b/Makefile @@ -7,6 +7,9 @@ build-chrome: build-noweb: docker build . -t devlikeapro/waha:noweb --build-arg USE_BROWSER=none --build-arg WHATSAPP_DEFAULT_ENGINE=NOWEB +build-plus-gows: + docker build . -t devlikeapro/waha-plus:gows --build-arg USE_BROWSER=none --build-arg WHATSAPP_DEFAULT_ENGINE=GOWS + build-all: build build-chrome build-noweb build-plus: @@ -42,3 +45,19 @@ up-webjs: start-proxy: docker run --rm -d --name squid-container -e TZ=UTC -p 3128:3128 ubuntu/squid:5.2-22.04_beta + +proto-gows: + rm -rf src/core/engines/gows/grpc + mkdir -p src/core/engines/gows/grpc + protoc \ + --plugin=protoc-gen-ts=node_modules/.bin/protoc-gen-ts \ + --plugin=protoc-gen-grpc=node_modules/.bin/grpc_tools_node_protoc_plugin \ + --js_out=import_style=commonjs,binary:src/core/engines/gows/grpc \ + --grpc_out=grpc_js:src/core/engines/gows/grpc \ + --ts_out=grpc_js:src/core/engines/gows/grpc \ + -I ../gows/proto ../gows/proto/gows.proto + +gows: + cd ../gows && \ + export PATH=${HOME}/go/bin:${PATH} && \ + make all \ No newline at end of file diff --git a/package.json b/package.json index fc5b3923..266909db 100644 --- a/package.json +++ b/package.json @@ -55,6 +55,7 @@ "express-basic-auth": "^1.2.1", "file-type": "16.5.4", "fs-extra": "^11.2.0", + "google-protobuf": "^3.21.4", "https-proxy-agent": "^7.0.0", "joi": "^17.13.3", "knex": "^3.1.0", @@ -94,6 +95,8 @@ "libsignal": "github:devlikeapro/libsignal-node#fork-master" }, "devDependencies": { + "@grpc/grpc-js": "^1.12.4", + "@grpc/proto-loader": "^0.7.13", "@nestjs/cli": "^9.0.0", "@nestjs/schematics": "^9.0.1", "@nestjs/testing": "^9.0.9", @@ -107,8 +110,10 @@ "eslint-config-prettier": "^6.10.0", "eslint-plugin-import": "^2.20.1", "eslint-plugin-simple-import-sort": "^10.0.0", + "grpc-tools": "^1.12.4", "jest": "^29.7.0", "prettier": "^1.19.1", + "protoc-gen-ts": "^0.8.7", "supertest": "^4.0.2", "ts-jest": "^29.1.3", "ts-loader": "^6.2.1", diff --git a/src/core/abc/EngineBootstrap.ts b/src/core/abc/EngineBootstrap.ts new file mode 100644 index 00000000..d6de138f --- /dev/null +++ b/src/core/abc/EngineBootstrap.ts @@ -0,0 +1,15 @@ +export interface EngineBootstrap { + bootstrap(): Promise; + + shutdown(): Promise; +} + +export class NoopEngineBootstrap implements EngineBootstrap { + async bootstrap(): Promise { + return; + } + + async shutdown(): Promise { + return; + } +} diff --git a/src/core/abc/manager.abc.ts b/src/core/abc/manager.abc.ts index 0238edff..b0a11ee1 100644 --- a/src/core/abc/manager.abc.ts +++ b/src/core/abc/manager.abc.ts @@ -1,8 +1,15 @@ import { BeforeApplicationShutdown, + OnApplicationBootstrap, UnprocessableEntityException, } from '@nestjs/common'; import { WhatsappConfigService } from '@waha/config.service'; +import { + EngineBootstrap, + NoopEngineBootstrap, +} from '@waha/core/abc/EngineBootstrap'; +import { GowsEngineConfigService } from '@waha/core/config/GowsEngineConfigService'; +import { GowsBootstrap } from '@waha/core/engines/gows/GowsBootstrap'; import { ISessionMeRepository } from '@waha/core/storage/ISessionMeRepository'; import { ISessionWorkerRepository } from '@waha/core/storage/ISessionWorkerRepository'; import { WAHAWebhook } from '@waha/structures/webhooks.dto'; @@ -29,7 +36,9 @@ import { WhatsappSession } from './session.abc'; // eslint-disable-next-line @typescript-eslint/no-var-requires const AsyncLock = require('async-lock'); -export abstract class SessionManager implements BeforeApplicationShutdown { +export abstract class SessionManager + implements BeforeApplicationShutdown, OnApplicationBootstrap +{ public store: any; public sessionAuthRepository: ISessionAuthRepository; public sessionConfigRepository: ISessionConfigRepository; @@ -42,8 +51,9 @@ export abstract class SessionManager implements BeforeApplicationShutdown { LOCK_TIMEOUT = 10_000; protected constructor( - protected config: WhatsappConfigService, protected log: PinoLogger, + protected config: WhatsappConfigService, + protected gowsConfigService: GowsEngineConfigService, ) { this.lock = new AsyncLock({ maxPending: Infinity, @@ -160,7 +170,21 @@ export abstract class SessionManager implements BeforeApplicationShutdown { beforeApplicationShutdown(signal?: string) { return; } + + onApplicationBootstrap(): any { + return; + } + + protected getEngineBootstrap(engine: WAHAEngine): EngineBootstrap { + const logger = this.log.logger.child({ engine: engine.toLowerCase() }); + if (engine === WAHAEngine.GOWS) { + const config = this.gowsConfigService.getBootstrapConfig(); + return new GowsBootstrap(logger, config); + } + return new NoopEngineBootstrap(); + } } + export function populateSessionInfo( event: WAHAEvents, session: WhatsappSession, diff --git a/src/core/abc/session.abc.ts b/src/core/abc/session.abc.ts index b50548f3..967bab2c 100644 --- a/src/core/abc/session.abc.ts +++ b/src/core/abc/session.abc.ts @@ -20,6 +20,7 @@ import { complete } from '@waha/utils/reactive/complete'; import { SwitchObservable } from '@waha/utils/reactive/SwitchObservable'; import * as fs from 'fs'; import * as lodash from 'lodash'; +import { PinoLogger } from 'nestjs-pino'; import * as NodeCache from 'node-cache'; import { Logger } from 'pino'; import { @@ -53,6 +54,7 @@ import { MessageTextRequest, MessageVideoRequest, MessageVoiceRequest, + SendSeenRequest, } from '../../structures/chatting.dto'; import { ContactQuery, ContactRequest } from '../../structures/contacts.dto'; import { @@ -380,7 +382,7 @@ export abstract class WhatsappSession { abstract reply(request: MessageReplyRequest); - abstract sendSeen(chat: ChatRequest); + abstract sendSeen(chat: SendSeenRequest); abstract startTyping(chat: ChatRequest); @@ -530,7 +532,7 @@ export abstract class WhatsappSession { if (!has || refresh) { await this.refreshProfilePicture(id); } - return this.profilePictures.get(id); + return this.profilePictures.get(id) || null; } protected async refreshProfilePicture(id: string) { diff --git a/src/core/app.module.core.ts b/src/core/app.module.core.ts index ef79b84a..68e3a2ab 100644 --- a/src/core/app.module.core.ts +++ b/src/core/app.module.core.ts @@ -10,6 +10,7 @@ import { ServerDebugController, } from '@waha/api/server.controller'; import { WebsocketGatewayCore } from '@waha/core/api/websocket.gateway.core'; +import { GowsEngineConfigService } from '@waha/core/config/GowsEngineConfigService'; import { WebJSEngineConfigService } from '@waha/core/config/WebJSEngineConfigService'; import { MediaLocalStorageModule } from '@waha/core/media/local/media.local.storage.module'; import { MediaLocalStorageConfig } from '@waha/core/media/local/MediaLocalStorageConfig'; @@ -148,6 +149,7 @@ const PROVIDERS = [ DashboardConfigServiceCore, SwaggerConfigServiceCore, WebJSEngineConfigService, + GowsEngineConfigService, WhatsappConfigService, EngineConfigService, WebsocketGatewayCore, diff --git a/src/core/config/GowsEngineConfigService.ts b/src/core/config/GowsEngineConfigService.ts new file mode 100644 index 00000000..c16714d5 --- /dev/null +++ b/src/core/config/GowsEngineConfigService.ts @@ -0,0 +1,26 @@ +import { Injectable } from '@nestjs/common'; +import { ConfigService } from '@nestjs/config'; +import { BootstrapConfig } from '@waha/core/engines/gows/GowsBootstrap'; +import { GowsConfig } from '@waha/core/engines/gows/session.gows.core'; + +@Injectable() +export class GowsEngineConfigService { + constructor(protected configService: ConfigService) {} + + getBootstrapConfig(): BootstrapConfig { + return { + path: this.configService.get('WAHA_GOWS_PATH'), + socket: this.getSocket(), + }; + } + + getSocket() { + return this.configService.get('WAHA_GOWS_SOCKET', '/tmp/gows.sock'); + } + + getConfig(): GowsConfig { + return { + connection: 'unix:' + this.getSocket(), + }; + } +} diff --git a/src/core/engines/gows/EventsFromObservable.ts b/src/core/engines/gows/EventsFromObservable.ts new file mode 100644 index 00000000..52a831da --- /dev/null +++ b/src/core/engines/gows/EventsFromObservable.ts @@ -0,0 +1,31 @@ +import { EventEmitter } from 'events'; +import { Observable, Subscription } from 'rxjs'; + +interface EventValue { + event: string; + data: T; +} + +export class EventsFromObservable extends EventEmitter { + private subscription: Subscription; + + constructor(private observable: Observable>) { + super(); + } + + on(event: Events, listener: (...args: any[]) => void): this { + return super.on(event as string, listener); + } + + start() { + this.subscription = this.observable.subscribe({ + next: (data) => this.emit(data.event, data.data), + complete: () => this.removeAllListeners(), + error: (err) => null, + }); + } + + stop() { + this.subscription.unsubscribe(); + } +} diff --git a/src/core/engines/gows/GowsBootstrap.ts b/src/core/engines/gows/GowsBootstrap.ts new file mode 100644 index 00000000..d2653999 --- /dev/null +++ b/src/core/engines/gows/GowsBootstrap.ts @@ -0,0 +1,72 @@ +import { EngineBootstrap } from '@waha/core/abc/EngineBootstrap'; +import { GowsSubprocess } from '@waha/core/engines/gows/GowsSubprocess'; +import { promises as fs } from 'fs'; +import { Logger } from 'pino'; + +export async function isUnixSocket(socketPath: string): Promise { + try { + // Get the file stats + const stats = await fs.lstat(socketPath); + // Check if the path is a socket + return stats.isSocket(); + } catch (error) { + // Handle errors (e.g., file does not exist) + if (error.code === 'ENOENT') { + return false; // Path does not exist + } + // Re-throw other errors for debugging purposes + throw error; + } +} + +export interface BootstrapConfig { + path: string; + socket: string; +} + +export class GowsBootstrap implements EngineBootstrap { + private gows: GowsSubprocess; + + constructor( + private logger: Logger, + private config: BootstrapConfig, + ) {} + + async bootstrap(): Promise { + if (!this.config.path) { + this.logger.warn('GOWS path is not set, skipping GOWS initialization.'); + this.logger.warn('Make sure to run GOWS manually.'); + this.checkSocket(this.config.socket); + return; + } + + this.gows = new GowsSubprocess( + this.logger, + this.config.path, + this.config.socket, + ); + this.gows.start(() => { + this.logger.info(`GOWS stopped, exiting...`); + process.kill(process.pid, 'SIGTERM'); + }); + await this.gows.waitWhenReady(10_000); + await this.checkSocket(this.config.socket); + return; + } + + async shutdown(): Promise { + if (this.gows) { + await this.gows.stop(); + } + } + + /** + * Check the provided path exists and is a socket + * @param path + */ + async checkSocket(path: string) { + if (!(await isUnixSocket(path))) { + throw new Error(`Invalid socket path: ${path}`); + } + } +} diff --git a/src/core/engines/gows/GowsEventStreamObservable.ts b/src/core/engines/gows/GowsEventStreamObservable.ts new file mode 100644 index 00000000..ffb9d4cf --- /dev/null +++ b/src/core/engines/gows/GowsEventStreamObservable.ts @@ -0,0 +1,69 @@ +import * as grpc from '@grpc/grpc-js'; +import { messages } from '@waha/core/engines/gows/grpc/gows'; +import { EnginePayload } from '@waha/structures/webhooks.dto'; +import { sleep } from '@waha/utils/promiseTimeout'; +import { Logger } from 'pino'; +import { Observable } from 'rxjs'; + +/** + * Observable that listens to a gRPC stream and emits EnginePayload objects. + * Pass a factory function that returns a client and a stream. + */ +export class GowsEventStreamObservable extends Observable { + _client: grpc.Client; + CLIENT_CLOSE_TIMEOUT = 1_000; + + constructor( + private logger: Logger, + factory: () => { + client: grpc.Client; + stream: grpc.ClientReadableStream; + }, + ) { + super((subscriber) => { + this.logger.debug('Creating grpc client and stream...'); + const { client, stream } = factory(); + this._client = client; + + stream.on('data', (raw) => { + const obj = raw.toObject(); + obj.data = JSON.parse(obj.data); + subscriber?.next(obj); + }); + + stream.on('end', (...args) => { + this.logger.debug('Stream ended', args); + subscriber?.complete(); + subscriber = null; + }); + + stream.on('error', async (err: any) => { + const CLIENT_CANCELLED_CODE = grpc.status.CANCELLED; + if (err.code === CLIENT_CANCELLED_CODE) { + this.logger.debug('Stream cancelled by client'); + return; + } + this.logger.error(err, 'Stream error'); + // Give some time to node event loop to process the error + await sleep(100); + subscriber?.error(err); + subscriber = null; + }); + + return async () => { + this.logger.debug('Closing stream client...'); + client.close(); + await sleep(this.CLIENT_CLOSE_TIMEOUT); + this.logger.debug('Stream client closed'); + + this.logger.debug('Cancelling stream...'); + stream.cancel(); + this.logger.debug('Stream cancelled'); + }; + }); + } + + get client(): Omit { + return this._client; + } +} diff --git a/src/core/engines/gows/GowsSubprocess.ts b/src/core/engines/gows/GowsSubprocess.ts new file mode 100644 index 00000000..dd4b3976 --- /dev/null +++ b/src/core/engines/gows/GowsSubprocess.ts @@ -0,0 +1,98 @@ +import { sleep, waitUntil } from '@waha/utils/promiseTimeout'; +import { spawn } from 'child_process'; +import { Logger } from 'pino'; + +export class GowsSubprocess { + private checkIntervalMs: number = 100; + private readyDelayMs: number = 1_000; + private readyText = 'gRPC server started!'; + + private child: any; + private ready: boolean = false; + + constructor( + private logger: Logger, + readonly path: string, + readonly socket: string, + ) {} + + start(onExit: (code: number) => void) { + this.logger.info('Starting GOWS subprocess...'); + this.logger.debug(`GOWS path '${this.path}', socket: '${this.socket}'...`); + + this.child = spawn(this.path, ['--socket', this.socket]); + this.logger.debug(`GOWS started with PID: ${this.child.pid}`); + this.child.on('close', async (code, singal) => { + const msg = code + ? `GOWS subprocess closed with code ${code}` + : `GOWS subprocess closed by signal ${singal}`; + this.logger.debug(msg); + onExit(code); + }); + this.child.on('error', (err) => { + this.logger.error(`GOWS subprocess error: ${err}`); + }); + + this.child.stderr.setEncoding('utf8'); + this.child.stderr.on('data', (data) => { + this.logger.error(data); + }); + + this.child.stdout.setEncoding('utf8'); + this.child.stdout.on('data', async (data) => { + // remove empty line at the end, split by \n + const lines = data.trim().split('\n'); + lines.forEach((line) => this.log(line)); + }); + this.listenReady(); + } + + listenReady() { + this.child.stdout.on('data', async (data) => { + if (this.ready) { + return; + } + if (!data.includes(this.readyText)) { + return; + } + await sleep(this.readyDelayMs); + this.ready = true; + this.logger.info('GOWS is ready'); + }); + } + + async waitWhenReady(timeout: number) { + const started = await waitUntil( + async () => this.ready, + this.checkIntervalMs, + timeout, + ); + if (!started) { + const msg = 'GOWS did not start after 10 seconds'; + this.logger.error(msg); + throw new Error(msg); + } + } + + async stop() { + this.logger.info('Stopping GOWS subprocess...'); + this.child?.kill(); + this.logger.info('GOWS subprocess stopped'); + } + + private log(msg) { + if (msg.startsWith('ERROR | ')) { + this.logger.error(msg.slice(8)); + } else if (msg.startsWith('WARN | ')) { + this.logger.warn(msg.slice(7)); + } else if (msg.startsWith('INFO | ')) { + this.logger.info(msg.slice(7)); + } else if (msg.startsWith('DEBUG | ')) { + this.logger.debug(msg.slice(8)); + } else if (msg.startsWith('TRACE | ')) { + this.logger.trace(msg.slice(8)); + } else { + this.logger.info(msg); + } + } +} diff --git a/src/core/engines/gows/session.gows.core.ts b/src/core/engines/gows/session.gows.core.ts new file mode 100644 index 00000000..986aa77a --- /dev/null +++ b/src/core/engines/gows/session.gows.core.ts @@ -0,0 +1,890 @@ +import { isJidGroup, jidNormalizedUser } from '@adiwajshing/baileys'; +import * as grpc from '@grpc/grpc-js'; +import { connectivityState } from '@grpc/grpc-js'; +import { UnprocessableEntityException } from '@nestjs/common'; +import { WhatsappSession } from '@waha/core/abc/session.abc'; +import { EventsFromObservable } from '@waha/core/engines/gows/EventsFromObservable'; +import { GowsEventStreamObservable } from '@waha/core/engines/gows/GowsEventStreamObservable'; +import { messages } from '@waha/core/engines/gows/grpc/gows'; +import { GowsAuthFactoryCore } from '@waha/core/engines/gows/store/GowsAuthFactoryCore'; +import { + parseMessageIdSerialized, + toCusFormat, + toJID, +} from '@waha/core/engines/noweb/session.noweb.core'; +import { AvailableInPlusVersion } from '@waha/core/exceptions'; +import { IMediaEngineProcessor } from '@waha/core/media/IMediaEngineProcessor'; +import { QR } from '@waha/core/QR'; +import { + ChatRequest, + CheckNumberStatusQuery, + MessageFileRequest, + MessageForwardRequest, + MessageImageRequest, + MessageLocationRequest, + MessageReactionRequest, + MessageReplyRequest, + MessageTextRequest, + MessageVoiceRequest, + SendSeenRequest, +} from '@waha/structures/chatting.dto'; +import { + ACK_UNKNOWN, + WAHAEvents, + WAHAPresenceStatus, + WAHASessionStatus, + WAMessageAck, +} from '@waha/structures/enums.dto'; +import { + WAHAChatPresences, + WAHAPresenceData, +} from '@waha/structures/presence.dto'; +import { WAMessage, WAMessageReaction } from '@waha/structures/responses.dto'; +import { MeInfo, ProxyConfig } from '@waha/structures/sessions.dto'; +import { EnginePayload, WAMessageAckBody } from '@waha/structures/webhooks.dto'; +import { sleep, waitUntil } from '@waha/utils/promiseTimeout'; +import { onlyEvent } from '@waha/utils/reactive/ops/onlyEvent'; +import * as NodeCache from 'node-cache'; +import { + debounceTime, + filter, + groupBy, + merge, + mergeMap, + Observable, + partition, + retry, + share, +} from 'rxjs'; +import { map } from 'rxjs/operators'; +import { promisify } from 'util'; + +import * as gows from './types'; +import MessageServiceClient = messages.MessageServiceClient; + +enum WhatsMeowEvent { + CONNECTED = 'gows.ConnectedEventData', + DISCONNECTED = 'events.Disconnected', + KEEP_ALIVE_TIMEOUT = 'events.KeepAliveTimeout', + QR_CHANNEL_ITEM = 'whatsmeow.QRChannelItem', + MESSAGE = 'events.Message', + RECEIPT = 'events.Receipt', + PRESENCE = 'events.Presence', + CHAT_PRESENCE = 'events.ChatPresence', + PUSH_NAME_SETTING = 'events.PushNameSetting', + LOGGED_OUT = 'events.LoggedOut', +} + +export interface GowsConfig { + connection: string; +} + +export class WhatsappSessionGoWSCore extends WhatsappSession { + protected authFactory = new GowsAuthFactoryCore(); + + protected qr: QR; + public client: MessageServiceClient; + protected stream$: GowsEventStreamObservable; + protected all$: Observable; + protected events: EventsFromObservable; + + protected me: MeInfo | null; + public session: messages.Session; + protected presences: any; + + public constructor(config) { + super(config); + this.qr = new QR(); + this.session = new messages.Session({ id: this.name }); + this.presences = new NodeCache({ + stdTTL: 60 * 60, // 1 hour + useClones: false, + }); + } + + async start() { + this.status = WAHASessionStatus.STARTING; + this.buildStreams(); + this.subscribeEvents(); + this.subscribeEngineEvents2(); + if (this.isDebugEnabled()) { + this.listenEngineEventsInDebugMode(); + } + + // start session + const auth = await this.authFactory.buildAuth(this.sessionStore, this.name); + + const level = messages.LogLevel[this.logger.level.toUpperCase()]; + const request = new messages.StartSessionRequest({ + id: this.name, + config: new messages.SessionConfig({ + store: new messages.SessionStoreConfig({ + address: auth.address(), + dialect: auth.dialect(), + }), + log: new messages.SessionLogConfig({ + level: level ?? messages.LogLevel.INFO, + }), + proxy: new messages.SessionProxyConfig({ + url: this.getProxyUrl(this.proxyConfig), + }), + }), + }); + + this.client = new MessageServiceClient( + this.engineConfig.connection, + grpc.credentials.createInsecure(), + ); + + try { + await promisify(this.client.StartSession)(request); + } catch (err) { + this.logger.error('Failed to start the client'); + this.logger.error(err, err.stack); + this.status = WAHASessionStatus.FAILED; + } + } + + protected getProxyUrl(config: ProxyConfig): string { + if (!config || !config.server) { + return ''; + } + const { server, username, password } = config; + const auth = username && password ? `${username}:${password}@` : ''; + const schema = 'http'; + return `${schema}://${auth}${server}`; + } + + buildStreams() { + this.stream$ = new GowsEventStreamObservable( + this.loggerBuilder.child({ grpc: 'stream' }), + () => { + const client = new messages.EventStreamClient( + this.engineConfig.connection, + grpc.credentials.createInsecure(), + ); + const stream = client.StreamEvents(this.session); + return { client, stream }; + }, + ); + + // Retry on error with delay + this.all$ = this.stream$.pipe(retry({ delay: 2_000 }), share()); + } + + subscribeEvents() { + // Handle connection status + this.events = new EventsFromObservable(this.all$); + const events = this.events; + events.on(WhatsMeowEvent.CONNECTED, (data) => { + this.status = WAHASessionStatus.WORKING; + this.me = { + id: toCusFormat(jidNormalizedUser(data.ID)), + pushName: data.PushName, + }; + }); + + events.on(WhatsMeowEvent.DISCONNECTED, () => { + if (this.status != WAHASessionStatus.STARTING) { + this.status = WAHASessionStatus.STARTING; + } + }); + events.on(WhatsMeowEvent.KEEP_ALIVE_TIMEOUT, () => { + if (this.status != WAHASessionStatus.STARTING) { + this.status = WAHASessionStatus.STARTING; + } + }); + events.on(WhatsMeowEvent.QR_CHANNEL_ITEM, async (data) => { + if (!data.Event) { + return; + } + if (data.Event == 'success') { + return; + } + if (data.Event != 'code') { + this.logger.warn(data, 'Failed QR item event'); + this.status = WAHASessionStatus.FAILED; + return; + } + const qr = data.Code; + if (!qr) { + return; + } + this.qr.save(qr); + this.printQR(this.qr); + this.status = WAHASessionStatus.SCAN_QR_CODE; + }); + events.on(WhatsMeowEvent.PUSH_NAME_SETTING, (data) => { + this.me = { ...this.me, pushName: data.Action.name }; + }); + events.on(WhatsMeowEvent.LOGGED_OUT, () => { + this.logger.error('Logged out'); + this.status = WAHASessionStatus.FAILED; + }); + events.on(WhatsMeowEvent.PRESENCE, (event: gows.Presence) => { + if (isJidGroup(event.From)) { + // So group is not "online" + return; + } + const Chat = event.From; + const Sender = event.From; + const stored = this.presences.get(Chat) || []; + // remove values by event.Sender + const filtered = stored.filter( + (p) => p.Sender !== Sender && p.From !== Sender, + ); + // add new value + this.presences.set(Chat, [...filtered, event]); + }); + events.on(WhatsMeowEvent.CHAT_PRESENCE, (event: gows.ChatPresence) => { + const Chat = event.Chat; + const Sender = event.Sender; + const stored = this.presences.get(Chat) || []; + // remove values by event.Sender + const filtered = stored.filter( + (p) => p.Sender !== Sender && p.From !== Sender, + ); + // add new value + this.presences.set(Chat, [...filtered, event]); + }); + events.start(); + } + + subscribeEngineEvents2() { + const all$ = this.all$; + this.events2.get(WAHAEvents.ENGINE_EVENT).switch(all$); + + const messages$ = all$.pipe(onlyEvent(WhatsMeowEvent.MESSAGE)); + + let [messagesFromMe$, messagesFromOthers$] = partition(messages$, isMine); + messagesFromMe$ = messagesFromMe$.pipe( + mergeMap((msg) => this.processIncomingMessage(msg, true)), + share(), // share it so we don't process twice in message.any + ); + messagesFromOthers$ = messagesFromOthers$.pipe( + mergeMap((msg) => this.processIncomingMessage(msg, true)), + share(), // share it so we don't process twice in message.any + ); + const messagesFromAll$ = merge(messagesFromMe$, messagesFromOthers$); + this.events2.get(WAHAEvents.MESSAGE).switch(messagesFromOthers$); + this.events2.get(WAHAEvents.MESSAGE_ANY).switch(messagesFromAll$); + + const receipt$ = all$.pipe(onlyEvent(WhatsMeowEvent.RECEIPT)); + const messageAck$ = receipt$.pipe( + mergeMap(this.receiptToMessageAck.bind(this)), + // Create a composite key for deduplication + groupBy((message) => `${message.id}-${message.ack}`), + mergeMap((group$) => + group$.pipe( + debounceTime(1000), // Wait 1 second for deduplication + map((message) => message), // Pass the latest message after debounce + ), + ), + ); + this.events2.get(WAHAEvents.MESSAGE_ACK).switch(messageAck$); + + const messageReactions$ = messages$.pipe( + filter((msg) => !!msg?.Message?.reactionMessage), + map(this.processMessageReaction.bind(this)), + ); + this.events2.get(WAHAEvents.MESSAGE_REACTION).switch(messageReactions$); + + const presence$ = all$.pipe( + onlyEvent(WhatsMeowEvent.PRESENCE), + filter((event) => !isJidGroup(event.From)), + ); + const chatPresence$ = all$.pipe(onlyEvent(WhatsMeowEvent.CHAT_PRESENCE)); + const presenceUpdates$ = merge(presence$, chatPresence$).pipe( + map((event) => this.toWahaPresences(event.From || event.Chat, [event])), + ); + this.events2.get(WAHAEvents.PRESENCE_UPDATE).switch(presenceUpdates$); + } + + async fetchContactProfilePicture(id: string): Promise { + const jid = toJID(this.ensureSuffix(id)); + const request = new messages.ProfilePictureRequest({ + jid: jid, + session: this.session, + }); + const response = await promisify(this.client.GetProfilePicture)(request); + const url = response.toObject().url; + return url; + } + + protected listenEngineEventsInDebugMode() { + this.events2.get(WAHAEvents.ENGINE_EVENT).subscribe((data) => { + this.logger.debug({ events: data }, `GOWS event`); + }); + } + + async stop(): Promise { + if (this.client) { + const response = await promisify(this.client.StopSession)(this.session); + response.toObject(); + } + this.status = WAHASessionStatus.STOPPED; + this.events?.stop(); + this.stopEvents(); + this.client?.close(); + } + + public async requestCode(phoneNumber: string, method: string, params?: any) { + if (this.status == WAHASessionStatus.STARTING) { + this.logger.debug('Waiting for connection update...'); + await waitUntil( + async () => this.status === WAHASessionStatus.SCAN_QR_CODE, + 100, + 2000, + ); + } + + if (this.status != WAHASessionStatus.SCAN_QR_CODE) { + const err = `Can request code only in SCAN_QR_CODE status. The current status is ${this.status}`; + throw new UnprocessableEntityException(err); + } + + const request = new messages.PairCodeRequest({ + session: this.session, + phone: phoneNumber, + }); + const response = await promisify(this.client.RequestCode)(request); + const code: string = response.toObject().code; + this.logger.info(`Your code: ${code}`); + return { code: code }; + } + + async unpair() { + await promisify(this.client.Logout)(this.session); + } + + public getSessionMeInfo(): MeInfo | null { + return this.me; + } + + /** + * START - Methods for API + */ + + /** + * Auth methods + */ + public getQR(): QR { + return this.qr; + } + + async getScreenshot(): Promise { + if (this.status === WAHASessionStatus.STARTING) { + throw new UnprocessableEntityException( + `The session is starting, please try again after few seconds`, + ); + } else if (this.status === WAHASessionStatus.SCAN_QR_CODE) { + return this.qr.get(); + } else if (this.status === WAHASessionStatus.WORKING) { + throw new UnprocessableEntityException( + `Can not get screenshot for non chrome based engine.`, + ); + } else { + throw new UnprocessableEntityException(`Unknown status - ${this.status}`); + } + } + + async sendText(request: MessageTextRequest) { + const jid = toJID(this.ensureSuffix(request.chatId)); + const message = new messages.MessageRequest({ + jid: jid, + text: request.text, + session: this.session, + }); + const response = await promisify(this.client.SendMessage)(message); + const data = response.toObject(); + return this.messageResponse(jid, data); + } + + protected messageResponse(jid, data) { + const id = buildMessageId({ + ID: data.id, + IsFromMe: true, + IsGroup: isJidGroup(jid), + Chat: jid, + Sender: this.me.id, + }); + return { + id: id, + _data: data, + }; + } + + checkNumberStatus(request: CheckNumberStatusQuery) { + throw new Error('Method not implemented.'); + } + + sendLocation(request: MessageLocationRequest) { + throw new Error('Method not implemented.'); + } + + forwardMessage(request: MessageForwardRequest): Promise { + throw new Error('Method not implemented.'); + } + + sendImage(request: MessageImageRequest) { + throw new AvailableInPlusVersion(); + } + + sendFile(request: MessageFileRequest) { + throw new AvailableInPlusVersion(); + } + + sendVoice(request: MessageVoiceRequest) { + throw new AvailableInPlusVersion(); + } + + reply(request: MessageReplyRequest) { + throw new Error('Method not implemented.'); + } + + async sendSeen(request: SendSeenRequest) { + const key = parseMessageIdSerialized(request.messageId); + const req = new messages.MarkReadRequest({ + session: this.session, + jid: key.remoteJid, + messageId: key.id, + sender: key.fromMe ? this.me.id : key.participant, + }); + const response = await promisify(this.client.MarkRead)(req); + response.toObject(); + return; + } + + startTyping(chat: ChatRequest) { + throw new Error('Method not implemented.'); + } + + stopTyping(chat: ChatRequest) { + throw new Error('Method not implemented.'); + } + + async setReaction(request: MessageReactionRequest) { + const key = parseMessageIdSerialized(request.messageId); + const message = new messages.MessageReaction({ + session: this.session, + jid: key.remoteJid, + messageId: key.id, + reaction: request.reaction, + sender: key.fromMe ? this.me.id : key.participant, + }); + const response = await promisify(this.client.SendReaction)(message); + const data = response.toObject(); + return this.messageResponse(key.remoteJid, data); + } + + public async setPresence(presence: WAHAPresenceStatus, chatId?: string) { + let request: any; + let method: any; + const jid = chatId ? toJID(this.ensureSuffix(chatId)) : null; + switch (presence) { + case WAHAPresenceStatus.ONLINE: + request = new messages.PresenceRequest({ + session: this.session, + status: messages.Presence.AVAILABLE, + }); + method = this.client.SendPresence; + break; + case WAHAPresenceStatus.OFFLINE: + request = new messages.PresenceRequest({ + session: this.session, + status: messages.Presence.UNAVAILABLE, + }); + method = this.client.SendPresence; + break; + case WAHAPresenceStatus.TYPING: + request = new messages.ChatPresenceRequest({ + session: this.session, + jid: jid, + status: messages.ChatPresence.TYPING, + }); + method = this.client.SendChatPresence; + break; + case WAHAPresenceStatus.RECORDING: + request = new messages.ChatPresenceRequest({ + session: this.session, + jid: jid, + status: messages.ChatPresence.RECORDING, + }); + method = this.client.SendChatPresence; + break; + case WAHAPresenceStatus.PAUSED: + request = new messages.ChatPresenceRequest({ + session: this.session, + jid: jid, + status: messages.ChatPresence.PAUSED, + }); + method = this.client.SendChatPresence; + break; + + default: + throw new Error('Invalid presence status'); + } + await promisify(method)(request); + } + + public async getPresences(): Promise { + const result: WAHAChatPresences[] = []; + for (const remoteJid in this.presences.keys()) { + const storedPresences = this.presences.get(remoteJid); + result.push(this.toWahaPresences(remoteJid, storedPresences)); + } + return result; + } + + public async getPresence(chatId: string): Promise { + const remoteJid = toJID(chatId); + if (!(remoteJid in this.presences.keys())) { + await this.subscribePresence(remoteJid); + await sleep(1000); + } + const result = this.presences.get(remoteJid) || []; + return this.toWahaPresences(remoteJid, result); + } + + async subscribePresence(chatId: string) { + const jid = toJID(chatId); + const req = new messages.SubscribePresenceRequest({ + session: this.session, + jid: jid, + }); + const response = await promisify(this.client.SubscribePresence)(req); + return response.toObject(); + } + + protected toWahaPresenceData( + data: gows.Presence | gows.ChatPresence, + ): WAHAPresenceData { + if ('From' in data) { + data = data as gows.Presence; + const lastKnownPresence = data.Unavailable + ? WAHAPresenceStatus.OFFLINE + : WAHAPresenceStatus.ONLINE; + return { + participant: toCusFormat(data.From), + lastKnownPresence: lastKnownPresence, + lastSeen: parseTimestamp(data.LastSeen), + }; + } + + data = data as gows.ChatPresence; + let lastKnownPresence: WAHAPresenceStatus; + if (data.State === gows.ChatPresenceState.PAUSED) { + lastKnownPresence = WAHAPresenceStatus.PAUSED; + } else if ( + data.State === gows.ChatPresenceState.COMPOSING && + data.Media === gows.ChatPresenceMedia.TEXT + ) { + lastKnownPresence = WAHAPresenceStatus.TYPING; + } else if ( + data.State === gows.ChatPresenceState.COMPOSING && + data.Media === gows.ChatPresenceMedia.AUDIO + ) { + lastKnownPresence = WAHAPresenceStatus.RECORDING; + } + return { + participant: toCusFormat(data.Sender), + lastKnownPresence: lastKnownPresence, + lastSeen: null, + }; + } + + protected toWahaPresences( + jid, + result: null | gows.Presence[] | gows.ChatPresence[], + ): WAHAChatPresences { + const chatId = toCusFormat(jid); + return { + id: chatId, + presences: result?.map(this.toWahaPresenceData.bind(this)), + }; + } + + // + // END - Methods for API + // + + private async processIncomingMessage(message, downloadMedia = true) { + // if there is no text or media message + if (!message) return; + if (!message.Message) return; + // Ignore reactions, we have dedicated handler for that + if (message.Message.reactionMessage) return; + // Ignore poll votes, we have dedicated handler for that + if (message.Message.pollUpdateMessage) return; + // Ignore protocol messages + if (message.Message.protocolMessage) return; + + if (downloadMedia) { + try { + message = await this.downloadMedia(message); + } catch (e) { + this.logger.error('Failed when tried to download media for a message'); + this.logger.error(e, e.stack); + } + } + return this.toWAMessage(message); + } + + protected downloadMedia(message) { + const processor = new GOWSEngineMediaProcessor(this); + return this.mediaManager.processMedia(processor, message, this.name); + } + + protected toWAMessage(message): WAMessage { + const fromToParticipant = getFromToParticipant(message); + const id = buildMessageId(message); + const body = this.extractBody(message.Message); + const replyTo = null; // TODO: this.extractReplyTo(message.message); + const ack = null; // TODO: Extract + + return { + id: id, + timestamp: parseTimestamp(message.Info.Timestamp), + from: toCusFormat(fromToParticipant.from), + fromMe: message.Info.IsFromMe, + body: body, + to: toCusFormat(fromToParticipant.to), + participant: toCusFormat(fromToParticipant.participant), + // Media + hasMedia: Boolean(message.media), + media: message.media, + mediaUrl: message.media?.url, + // @ts-ignore + ack: ack, + // @ts-ignore + ackName: WAMessageAck[ack] || ACK_UNKNOWN, + replyTo: replyTo, + _data: message, + }; + } + + private extractBody(message) { + if (!message) { + return null; + } + let body = message.Conversation || message.conversation; + if (!body) { + body = message.extendedTextMessage?.text; + } + if (!body) { + const media = extractMediaContent(message); + body = media?.caption; + } + return body; + } + + public async getEngineInfo() { + const clientState = this.client?.getChannel().getConnectivityState(false); + const streamState = this.stream$?.client + ?.getChannel() + .getConnectivityState(false); + const grpc = { + client: connectivityState[clientState] || 'NO_CLIENT', + stream: connectivityState[streamState] || 'NO_STREAM', + }; + if (!this.client) { + return { + grpc: grpc, + }; + } + + let gows; + try { + const response = await promisify(this.client.GetSessionState)( + this.session, + ); + const info = response.toObject(); + gows = { ...info }; + } catch (err) { + gows = { error: err }; + } + return { + grpc: grpc, + gows: gows, + }; + } + + receiptToMessageAck(receipt: any): WAMessageAckBody[] { + const fromToParticipant = getFromToParticipant(receipt); + + let ack; + switch (receipt.Type) { + case '': + ack = WAMessageAck.DEVICE; + break; + case 'server-error': + ack = WAMessageAck.ERROR; + break; + case 'inactive': + ack = WAMessageAck.DEVICE; + break; + case 'active': + ack = WAMessageAck.DEVICE; + break; + case 'read': + ack = WAMessageAck.READ; + break; + case 'played': + ack = WAMessageAck.PLAYED; + break; + default: + return []; + } + const acks = []; + for (const messageId of receipt.MessageIDs) { + const msg = { + ...receipt, + ID: messageId, + // Reverse the IsFromMe flag + IsFromMe: !receipt.IsFromMe, + Sender: receipt.MessageSender || this.me?.id, + }; + const id = buildMessageId(msg); + const body: WAMessageAckBody = { + id: id, + from: toCusFormat(fromToParticipant.from), + to: toCusFormat(fromToParticipant.to), + participant: toCusFormat(fromToParticipant.participant), + fromMe: msg.IsFromMe, + ack: ack, + ackName: WAMessageAck[ack] || ACK_UNKNOWN, + _data: receipt, + }; + acks.push(body); + } + return acks; + } + + private processMessageReaction(message): WAMessageReaction | null { + if (!message) return null; + if (!message.Message) return null; + if (!message.Message.reactionMessage) return null; + + const id = buildMessageId(message); + const fromToParticipant = getFromToParticipant(message); + const reactionMessage = message.Message.reactionMessage; + const messageId = this.buildMessageIdFromKey(reactionMessage.key); + const reaction: WAMessageReaction = { + id: id, + timestamp: parseTimestamp(message.Info.Timestamp), + from: toCusFormat(fromToParticipant.from), + fromMe: message.Info.IsFromMe, + to: toCusFormat(fromToParticipant.to), + participant: toCusFormat(fromToParticipant.participant), + reaction: { + text: reactionMessage.text, + messageId: messageId, + }, + }; + return reaction; + } + + buildMessageIdFromKey(key: WARawKey) { + const sender = key.fromMe ? this.me.id : key.participant || key.remoteJID; + const info: MessageIdData = { + Chat: key.remoteJID, + Sender: sender, + ID: key.ID, + IsFromMe: key.fromMe, + IsGroup: isJidGroup(key.remoteJID), + }; + return buildMessageId(info); + } +} + +export class GOWSEngineMediaProcessor implements IMediaEngineProcessor { + constructor(public session: WhatsappSessionGoWSCore) {} + + hasMedia(message: any): boolean { + return Boolean(extractMediaContent(message.Message)); + } + + getMessageId(message: any): string { + return message.Info.ID; + } + + getMimetype(message: any): string { + const content = extractMediaContent(message.Message); + return content.mimetype; + } + + async getMediaBuffer(message: any): Promise { + const data = JSON.stringify(message.Message); + const request = new messages.DownloadMediaRequest({ + session: this.session.session, + message: data, + }); + const response = await promisify(this.session.client.DownloadMedia)( + request, + ); + const obj = response.toObject(); + return Buffer.from(obj.content); + } + + getFilename(message: any): string | null { + const content = extractMediaContent(message.Message); + return content?.fileName; + } +} + +export function extractMediaContent(message: any) { + const mediaContent = + message?.documentMessage || + message?.imageMessage || + message?.videoMessage || + message?.audioMessage || + message?.stickerMessage; + return mediaContent; +} + +function isMine(message) { + return !!message.Info.IsFromMe; +} + +function getFromToParticipant(message) { + const info = message.Info || message; + return { + from: info.Chat, + to: info.IsGroup ? info.Sender : null, + participant: info.IsGroup ? info.Sender : null, + }; +} + +interface MessageIdData { + Info?: MessageIdData; + Chat: string; + Sender: string; + ID: string; + IsFromMe: boolean; + IsGroup: boolean; +} + +function buildMessageId(message: MessageIdData) { + const info = message.Info || message; + const chatId = toCusFormat(info.Chat); + const participant = toCusFormat(info.Sender); + if (info.IsGroup) { + return `${info.IsFromMe}_${chatId}_${info.ID}_${participant}`; + } + return `${info.IsFromMe}_${chatId}_${info.ID}`; +} + +interface WARawKey { + remoteJID: string; + fromMe: boolean; + ID: string; + participant?: string; +} + +function parseTimestamp(timestamp: string): number { + if (timestamp.startsWith('0001')) { + return null; + } + // "2024-12-25T14:28:42+03:00" => 1234567890 + return new Date(timestamp).getTime(); +} diff --git a/src/core/engines/gows/store/GowsAuth.ts b/src/core/engines/gows/store/GowsAuth.ts new file mode 100644 index 00000000..ca55b502 --- /dev/null +++ b/src/core/engines/gows/store/GowsAuth.ts @@ -0,0 +1,4 @@ +export interface GowsAuth { + address(): string; + dialect(): string; +} diff --git a/src/core/engines/gows/store/GowsAuthFactoryCore.ts b/src/core/engines/gows/store/GowsAuthFactoryCore.ts new file mode 100644 index 00000000..30f20f2d --- /dev/null +++ b/src/core/engines/gows/store/GowsAuthFactoryCore.ts @@ -0,0 +1,26 @@ +import * as path from 'node:path'; + +import { DataStore } from '@waha/core/abc/DataStore'; +import { GowsAuthSimple } from '@waha/core/engines/gows/store/GowsAuthSimple'; +import { LocalStore } from '@waha/core/storage/LocalStore'; + +import { GowsAuth } from './GowsAuth'; + +export class GowsAuthFactoryCore { + buildAuth(store: DataStore, name: string): Promise { + if (store instanceof LocalStore) return this.buildSqlite3(store, name); + throw new Error(`Unsupported store type '${store.constructor.name}'`); + } + + protected async buildSqlite3( + store: LocalStore, + name: string, + ): Promise { + await store.init(name); + const authFolder = store.getSessionDirectory(name); + // resolve path + const authFolderFullPath = path.resolve(authFolder); + const connection = `file:${authFolderFullPath}/gows.db`; + return new GowsAuthSimple(connection, 'sqlite3'); + } +} diff --git a/src/core/engines/gows/store/GowsAuthSimple.ts b/src/core/engines/gows/store/GowsAuthSimple.ts new file mode 100644 index 00000000..8f14104c --- /dev/null +++ b/src/core/engines/gows/store/GowsAuthSimple.ts @@ -0,0 +1,19 @@ +import { GowsAuth } from '@waha/core/engines/gows/store/GowsAuth'; + +export class GowsAuthSimple implements GowsAuth { + private _address: string; + private _dialect: string; + + constructor(address: string, dialect: string) { + this._address = address; + this._dialect = dialect; + } + + address() { + return this._address; + } + + dialect(): string { + return this._dialect; + } +} diff --git a/src/core/engines/gows/types.ts b/src/core/engines/gows/types.ts new file mode 100644 index 00000000..bb10f9e3 --- /dev/null +++ b/src/core/engines/gows/types.ts @@ -0,0 +1,24 @@ +export interface Presence { + From: string; + Unavailable: boolean; + LastSeen: string; +} + +export enum ChatPresenceState { + COMPOSING = 'composing', + PAUSED = 'paused', +} + +export enum ChatPresenceMedia { + TEXT = '', + AUDIO = 'audio', +} + +export interface ChatPresence { + Chat: string; + Sender: string; + IsFromMe: boolean; + IsGroup: boolean; + State: ChatPresenceState; + Media: ChatPresenceMedia; +} diff --git a/src/core/engines/noweb/session.noweb.core.ts b/src/core/engines/noweb/session.noweb.core.ts index 8d9bc555..cf043e2f 100644 --- a/src/core/engines/noweb/session.noweb.core.ts +++ b/src/core/engines/noweb/session.noweb.core.ts @@ -1986,7 +1986,7 @@ export class NOWEBEngineMediaProcessor implements IMediaEngineProcessor { /** * Convert from 11111111111@s.whatsapp.net to 11111111111@c.us */ -function toCusFormat(remoteJid) { +export function toCusFormat(remoteJid) { if (!remoteJid) { return remoteJid; } @@ -2008,7 +2008,9 @@ function toCusFormat(remoteJid) { if (remoteJid == 'me') { return remoteJid; } - const number = remoteJid.split('@')[0]; + let number = remoteJid.split('@')[0]; + // remove :{device} part + number = number.split(':')[0]; return ensureSuffix(number); } @@ -2051,7 +2053,7 @@ function buildMessageId({ id, remoteJid, fromMe, participant }: WAMessageKey) { * false_11111111111@c.us_AAA * {id: "AAA", remoteJid: "11111111111@s.whatsapp.net", "fromMe": false} */ -function parseMessageIdSerialized( +export function parseMessageIdSerialized( messageId: string, soft: boolean = false, ): WAMessageKey { diff --git a/src/core/manager.core.ts b/src/core/manager.core.ts index ebbce8ce..0f63a5e0 100644 --- a/src/core/manager.core.ts +++ b/src/core/manager.core.ts @@ -3,7 +3,10 @@ import { NotFoundException, UnprocessableEntityException, } from '@nestjs/common'; +import { EngineBootstrap } from '@waha/core/abc/EngineBootstrap'; +import { GowsEngineConfigService } from '@waha/core/config/GowsEngineConfigService'; import { WebJSEngineConfigService } from '@waha/core/config/WebJSEngineConfigService'; +import { WhatsappSessionGoWSCore } from '@waha/core/engines/gows/session.gows.core'; import { WebhookConductor } from '@waha/core/integrations/webhooks/WebhookConductor'; import { MediaStorageFactory } from '@waha/core/media/MediaStorageFactory'; import { DefaultMap } from '@waha/utils/DefaultMap'; @@ -66,19 +69,22 @@ export class SessionManagerCore extends SessionManager { protected readonly EngineClass: typeof WhatsappSession; protected events2: DefaultMap>; + protected readonly engineBootstrap: EngineBootstrap; constructor( config: WhatsappConfigService, private engineConfigService: EngineConfigService, private webjsEngineConfigService: WebJSEngineConfigService, + gowsConfigService: GowsEngineConfigService, log: PinoLogger, private mediaStorageFactory: MediaStorageFactory, ) { - super(config, log); + super(log, config, gowsConfigService); this.session = DefaultSessionStatus.STOPPED; this.sessionConfig = null; const engineName = this.engineConfigService.getDefaultEngineName(); this.EngineClass = this.getEngine(engineName); + this.engineBootstrap = this.getEngineBootstrap(engineName); this.events2 = new DefaultMap>( (key) => @@ -89,7 +95,6 @@ export class SessionManagerCore extends SessionManager { this.store = new LocalStoreCore(engineName.toLowerCase()); this.sessionAuthRepository = new LocalSessionAuthRepository(this.store); - this.startPredefinedSessions(); this.clearStorage().catch((error) => { this.log.error({ error }, 'Error while clearing storage'); }); @@ -100,6 +105,8 @@ export class SessionManagerCore extends SessionManager { return WhatsappSessionWebJSCore; } else if (engine === WAHAEngine.NOWEB) { return WhatsappSessionNoWebCore; + } else if (engine === WAHAEngine.GOWS) { + return WhatsappSessionGoWSCore; } else { throw new NotFoundException(`Unknown whatsapp engine '${engine}'.`); } @@ -117,6 +124,12 @@ export class SessionManagerCore extends SessionManager { return; } await this.stop(this.DEFAULT, true); + await this.engineBootstrap.shutdown(); + } + + async onApplicationBootstrap() { + await this.engineBootstrap.bootstrap(); + this.startPredefinedSessions(); } private async clearStorage() { @@ -179,6 +192,8 @@ export class SessionManagerCore extends SessionManager { }; if (this.EngineClass === WhatsappSessionWebJSCore) { sessionConfig.engineConfig = this.webjsEngineConfigService.getConfig(); + } else if (this.EngineClass === WhatsappSessionGoWSCore) { + sessionConfig.engineConfig = this.gowsConfigService.getConfig(); } await this.sessionAuthRepository.init(name); // @ts-ignore diff --git a/src/structures/enums.dto.ts b/src/structures/enums.dto.ts index 6dd5e89c..d7220ebd 100644 --- a/src/structures/enums.dto.ts +++ b/src/structures/enums.dto.ts @@ -41,6 +41,7 @@ export enum WAHASessionStatus { export enum WAHAEngine { WEBJS = 'WEBJS', NOWEB = 'NOWEB', + GOWS = 'GOWS', } export enum WAHAPresenceStatus { diff --git a/src/structures/webhooks.dto.ts b/src/structures/webhooks.dto.ts index 59f938e6..b199bffe 100644 --- a/src/structures/webhooks.dto.ts +++ b/src/structures/webhooks.dto.ts @@ -32,6 +32,8 @@ export class WAMessageAckBody { fromMe: boolean; ack: WAMessageAck; ackName: string; + + _data?: any; } export class WAGroupPayload { diff --git a/src/utils/reactive/ops/onlyEvent.ts b/src/utils/reactive/ops/onlyEvent.ts new file mode 100644 index 00000000..ba4daf02 --- /dev/null +++ b/src/utils/reactive/ops/onlyEvent.ts @@ -0,0 +1,14 @@ +import { filter, pipe } from 'rxjs'; +import { map } from 'rxjs/operators'; + +interface EventValue { + event: string; + data: T; +} + +export function onlyEvent(event: any) { + return pipe( + filter>((obj) => obj.event === event), + map((event) => event.data), + ); +} diff --git a/yarn.lock b/yarn.lock index 5287c410..442f130e 100644 --- a/yarn.lock +++ b/yarn.lock @@ -1242,6 +1242,30 @@ __metadata: languageName: node linkType: hard +"@grpc/grpc-js@npm:^1.12.4": + version: 1.12.4 + resolution: "@grpc/grpc-js@npm:1.12.4" + dependencies: + "@grpc/proto-loader": ^0.7.13 + "@js-sdsl/ordered-map": ^4.4.2 + checksum: 6dab69c70173735b775717361662ed5a931d2ca1bc9f3ed5f7095a05f26eb59489b1805b0a8f78b77f34d48a896f358ca951f0ce38d8053d18b0d5d0fd92837f + languageName: node + linkType: hard + +"@grpc/proto-loader@npm:^0.7.13": + version: 0.7.13 + resolution: "@grpc/proto-loader@npm:0.7.13" + dependencies: + lodash.camelcase: ^4.3.0 + long: ^5.0.0 + protobufjs: ^7.2.5 + yargs: ^17.7.2 + bin: + proto-loader-gen-types: build/bin/proto-loader-gen-types.js + checksum: 399c1b8a4627f93dc31660d9636ea6bf58be5675cc7581e3df56a249369e5be02c6cd0d642c5332b0d5673bc8621619bc06fb045aa3e8f57383737b5d35930dc + languageName: node + linkType: hard + "@hapi/boom@npm:^9.1.3": version: 9.1.4 resolution: "@hapi/boom@npm:9.1.4" @@ -1780,6 +1804,13 @@ __metadata: languageName: node linkType: hard +"@js-sdsl/ordered-map@npm:^4.4.2": + version: 4.4.2 + resolution: "@js-sdsl/ordered-map@npm:4.4.2" + checksum: a927ae4ff8565ecb75355cc6886a4f8fadbf2af1268143c96c0cce3ba01261d241c3f4ba77f21f3f017a00f91dfe9e0673e95f830255945c80a0e96c6d30508a + languageName: node + linkType: hard + "@lukeed/csprng@npm:^1.0.0": version: 1.1.0 resolution: "@lukeed/csprng@npm:1.1.0" @@ -1787,6 +1818,25 @@ __metadata: languageName: node linkType: hard +"@mapbox/node-pre-gyp@npm:^1.0.5": + version: 1.0.11 + resolution: "@mapbox/node-pre-gyp@npm:1.0.11" + dependencies: + detect-libc: ^2.0.0 + https-proxy-agent: ^5.0.0 + make-dir: ^3.1.0 + node-fetch: ^2.6.7 + nopt: ^5.0.0 + npmlog: ^5.0.1 + rimraf: ^3.0.2 + semver: ^7.3.5 + tar: ^6.1.11 + bin: + node-pre-gyp: bin/node-pre-gyp + checksum: b848f6abc531a11961d780db813cc510ca5a5b6bf3184d72134089c6875a91c44d571ba6c1879470020803f7803609e7b2e6e429651c026fe202facd11d444b8 + languageName: node + linkType: hard + "@microsoft/tsdoc@npm:^0.14.2": version: 0.14.2 resolution: "@microsoft/tsdoc@npm:0.14.2" @@ -3675,6 +3725,13 @@ __metadata: languageName: node linkType: hard +"abbrev@npm:1": + version: 1.1.1 + resolution: "abbrev@npm:1.1.1" + checksum: a4a97ec07d7ea112c517036882b2ac22f3109b7b19077dc656316d07d308438aac28e4d9746dc4d84bf6b1e75b4a7b0a5f3cb30592419f128ca9a8cee3bcfa17 + languageName: node + linkType: hard + "abbrev@npm:^2.0.0": version: 2.0.0 resolution: "abbrev@npm:2.0.0" @@ -3753,6 +3810,15 @@ __metadata: languageName: node linkType: hard +"agent-base@npm:6": + version: 6.0.2 + resolution: "agent-base@npm:6.0.2" + dependencies: + debug: 4 + checksum: f52b6872cc96fd5f622071b71ef200e01c7c4c454ee68bc9accca90c98cfb39f2810e3e9aa330435835eedc8c23f4f8a15267f67c6e245d2b33757575bdac49d + languageName: node + linkType: hard + "agent-base@npm:^7.0.2, agent-base@npm:^7.1.0, agent-base@npm:^7.1.1": version: 7.1.1 resolution: "agent-base@npm:7.1.1" @@ -3935,6 +4001,13 @@ __metadata: languageName: node linkType: hard +"aproba@npm:^1.0.3 || ^2.0.0": + version: 2.0.0 + resolution: "aproba@npm:2.0.0" + checksum: 5615cadcfb45289eea63f8afd064ab656006361020e1735112e346593856f87435e02d8dcc7ff0d11928bc7d425f27bc7c2a84f6c0b35ab0ff659c814c138a24 + languageName: node + linkType: hard + "archiver-utils@npm:^2.1.0": version: 2.1.0 resolution: "archiver-utils@npm:2.1.0" @@ -3986,6 +4059,16 @@ __metadata: languageName: node linkType: hard +"are-we-there-yet@npm:^2.0.0": + version: 2.0.0 + resolution: "are-we-there-yet@npm:2.0.0" + dependencies: + delegates: ^1.0.0 + readable-stream: ^3.6.0 + checksum: 6c80b4fd04ecee6ba6e737e0b72a4b41bdc64b7d279edfc998678567ff583c8df27e27523bc789f2c99be603ffa9eaa612803da1d886962d2086e7ff6fa90c7c + languageName: node + linkType: hard + "arg@npm:^4.1.0": version: 4.1.3 resolution: "arg@npm:4.1.3" @@ -5077,6 +5160,15 @@ __metadata: languageName: node linkType: hard +"color-support@npm:^1.1.2": + version: 1.1.3 + resolution: "color-support@npm:1.1.3" + bin: + color-support: bin.js + checksum: 9b7356817670b9a13a26ca5af1c21615463b500783b739b7634a0c2047c16cef4b2865d7576875c31c3cddf9dd621fa19285e628f20198b233a5cfdda6d0793b + languageName: node + linkType: hard + "color@npm:^4.2.3": version: 4.2.3 resolution: "color@npm:4.2.3" @@ -5176,6 +5268,13 @@ __metadata: languageName: node linkType: hard +"console-control-strings@npm:^1.0.0, console-control-strings@npm:^1.1.0": + version: 1.1.0 + resolution: "console-control-strings@npm:1.1.0" + checksum: 8755d76787f94e6cf79ce4666f0c5519906d7f5b02d4b884cf41e11dcd759ed69c57da0670afd9236d229a46e0f9cf519db0cd829c6dca820bb5a5c3def584ed + languageName: node + linkType: hard + "content-disposition@npm:0.5.4": version: 0.5.4 resolution: "content-disposition@npm:0.5.4" @@ -5572,6 +5671,13 @@ __metadata: languageName: node linkType: hard +"delegates@npm:^1.0.0": + version: 1.0.0 + resolution: "delegates@npm:1.0.0" + checksum: a51744d9b53c164ba9c0492471a1a2ffa0b6727451bdc89e31627fdf4adda9d51277cfcbfb20f0a6f08ccb3c436f341df3e92631a3440226d93a8971724771fd + languageName: node + linkType: hard + "depd@npm:2.0.0": version: 2.0.0 resolution: "depd@npm:2.0.0" @@ -6873,6 +6979,23 @@ __metadata: languageName: node linkType: hard +"gauge@npm:^3.0.0": + version: 3.0.2 + resolution: "gauge@npm:3.0.2" + dependencies: + aproba: ^1.0.3 || ^2.0.0 + color-support: ^1.1.2 + console-control-strings: ^1.0.0 + has-unicode: ^2.0.1 + object-assign: ^4.1.1 + signal-exit: ^3.0.0 + string-width: ^4.2.3 + strip-ansi: ^6.0.1 + wide-align: ^1.1.2 + checksum: 81296c00c7410cdd48f997800155fbead4f32e4f82109be0719c63edc8560e6579946cc8abd04205297640691ec26d21b578837fd13a4e96288ab4b40b1dc3e9 + languageName: node + linkType: hard + "gensync@npm:^1.0.0-beta.2": version: 1.0.0-beta.2 resolution: "gensync@npm:1.0.0-beta.2" @@ -7065,6 +7188,13 @@ __metadata: languageName: node linkType: hard +"google-protobuf@npm:^3.21.4": + version: 3.21.4 + resolution: "google-protobuf@npm:3.21.4" + checksum: 048fa2cb579f5f88c977774b2ae36851807379d9329a6895fe3685df69ba6c927e2ff463d08d5eecd56becd9c65bca406f34f90e27984a3077e8629bb3a2a766 + languageName: node + linkType: hard + "gopd@npm:^1.0.1": version: 1.0.1 resolution: "gopd@npm:1.0.1" @@ -7081,6 +7211,18 @@ __metadata: languageName: node linkType: hard +"grpc-tools@npm:^1.12.4": + version: 1.12.4 + resolution: "grpc-tools@npm:1.12.4" + dependencies: + "@mapbox/node-pre-gyp": ^1.0.5 + bin: + grpc_tools_node_protoc: bin/protoc.js + grpc_tools_node_protoc_plugin: bin/protoc_plugin.js + checksum: 5ea946c2d231e2b58be48b14e12f2982b44845f62d05327cf84f67fae30355c7a781d8443b0988f9db77d4eb219f15e5f98e869f6a14c4276084726e48dc8ec0 + languageName: node + linkType: hard + "has-bigints@npm:^1.0.1, has-bigints@npm:^1.0.2": version: 1.0.2 resolution: "has-bigints@npm:1.0.2" @@ -7134,6 +7276,13 @@ __metadata: languageName: node linkType: hard +"has-unicode@npm:^2.0.1": + version: 2.0.1 + resolution: "has-unicode@npm:2.0.1" + checksum: 1eab07a7436512db0be40a710b29b5dc21fa04880b7f63c9980b706683127e3c1b57cb80ea96d47991bdae2dfe479604f6a1ba410106ee1046a41d1bd0814400 + languageName: node + linkType: hard + "hasown@npm:^2.0.0, hasown@npm:^2.0.1, hasown@npm:^2.0.2": version: 2.0.2 resolution: "hasown@npm:2.0.2" @@ -7199,6 +7348,16 @@ __metadata: languageName: node linkType: hard +"https-proxy-agent@npm:^5.0.0": + version: 5.0.1 + resolution: "https-proxy-agent@npm:5.0.1" + dependencies: + agent-base: 6 + debug: 4 + checksum: 571fccdf38184f05943e12d37d6ce38197becdd69e58d03f43637f7fa1269cf303a7d228aa27e5b27bbd3af8f09fd938e1c91dcfefff2df7ba77c20ed8dfc765 + languageName: node + linkType: hard + "https-proxy-agent@npm:^7.0.0, https-proxy-agent@npm:^7.0.1, https-proxy-agent@npm:^7.0.2": version: 7.0.4 resolution: "https-proxy-agent@npm:7.0.4" @@ -8537,6 +8696,13 @@ __metadata: languageName: node linkType: hard +"lodash.camelcase@npm:^4.3.0": + version: 4.3.0 + resolution: "lodash.camelcase@npm:4.3.0" + checksum: cb9227612f71b83e42de93eccf1232feeb25e705bdb19ba26c04f91e885bfd3dd5c517c4a97137658190581d3493ea3973072ca010aab7e301046d90740393d1 + languageName: node + linkType: hard + "lodash.clonedeep@npm:^4.5.0": version: 4.5.0 resolution: "lodash.clonedeep@npm:4.5.0" @@ -8663,6 +8829,15 @@ __metadata: languageName: node linkType: hard +"make-dir@npm:^3.1.0": + version: 3.1.0 + resolution: "make-dir@npm:3.1.0" + dependencies: + semver: ^6.0.0 + checksum: 484200020ab5a1fdf12f393fe5f385fc8e4378824c940fba1729dcd198ae4ff24867bc7a5646331e50cead8abff5d9270c456314386e629acec6dff4b8016b78 + languageName: node + linkType: hard + "make-dir@npm:^4.0.0": version: 4.0.0 resolution: "make-dir@npm:4.0.0" @@ -9206,7 +9381,7 @@ __metadata: languageName: node linkType: hard -"node-fetch@npm:^2.6.1, node-fetch@npm:^2.6.9": +"node-fetch@npm:^2.6.1, node-fetch@npm:^2.6.7, node-fetch@npm:^2.6.9": version: 2.7.0 resolution: "node-fetch@npm:2.7.0" dependencies: @@ -9279,6 +9454,17 @@ __metadata: languageName: node linkType: hard +"nopt@npm:^5.0.0": + version: 5.0.0 + resolution: "nopt@npm:5.0.0" + dependencies: + abbrev: 1 + bin: + nopt: bin/nopt.js + checksum: d35fdec187269503843924e0114c0c6533fb54bbf1620d0f28b4b60ba01712d6687f62565c55cc20a504eff0fbe5c63e22340c3fad549ad40469ffb611b04f2f + languageName: node + linkType: hard + "nopt@npm:^7.0.0": version: 7.2.1 resolution: "nopt@npm:7.2.1" @@ -9306,6 +9492,18 @@ __metadata: languageName: node linkType: hard +"npmlog@npm:^5.0.1": + version: 5.0.1 + resolution: "npmlog@npm:5.0.1" + dependencies: + are-we-there-yet: ^2.0.0 + console-control-strings: ^1.1.0 + gauge: ^3.0.0 + set-blocking: ^2.0.0 + checksum: 516b2663028761f062d13e8beb3f00069c5664925871a9b57989642ebe09f23ab02145bf3ab88da7866c4e112cafff72401f61a672c7c8a20edc585a7016ef5f + languageName: node + linkType: hard + "nth-check@npm:^2.0.1": version: 2.1.1 resolution: "nth-check@npm:2.1.1" @@ -10087,6 +10285,35 @@ __metadata: languageName: node linkType: hard +"protobufjs@npm:^7.2.5": + version: 7.4.0 + resolution: "protobufjs@npm:7.4.0" + dependencies: + "@protobufjs/aspromise": ^1.1.2 + "@protobufjs/base64": ^1.1.2 + "@protobufjs/codegen": ^2.0.4 + "@protobufjs/eventemitter": ^1.1.0 + "@protobufjs/fetch": ^1.1.0 + "@protobufjs/float": ^1.0.2 + "@protobufjs/inquire": ^1.1.0 + "@protobufjs/path": ^1.1.2 + "@protobufjs/pool": ^1.1.0 + "@protobufjs/utf8": ^1.1.0 + "@types/node": ">=13.7.0" + long: ^5.0.0 + checksum: ba0e6b60541bbf818bb148e90f5eb68bd99004e29a6034ad9895a381cbd352be8dce5376e47ae21b2e05559f2505b4a5f4a3c8fa62402822c6ab4dcdfb89ffb3 + languageName: node + linkType: hard + +"protoc-gen-ts@npm:^0.8.7": + version: 0.8.7 + resolution: "protoc-gen-ts@npm:0.8.7" + bin: + protoc-gen-ts: protoc-gen-ts.js + checksum: 556d89ef43b152b9f457dab79d51f649f5108fb1ef447da41f0c3b6f33df1a5bd8ef3392c0ce9b25730e68ccb0c972a52ca37d50775ce4602a94b502aeab1e88 + languageName: node + linkType: hard + "proxy-addr@npm:~2.0.7": version: 2.0.7 resolution: "proxy-addr@npm:2.0.7" @@ -10923,7 +11150,7 @@ __metadata: languageName: node linkType: hard -"signal-exit@npm:^3.0.2, signal-exit@npm:^3.0.3, signal-exit@npm:^3.0.7": +"signal-exit@npm:^3.0.0, signal-exit@npm:^3.0.2, signal-exit@npm:^3.0.3, signal-exit@npm:^3.0.7": version: 3.0.7 resolution: "signal-exit@npm:3.0.7" checksum: a2f098f247adc367dffc27845853e9959b9e88b01cb301658cfe4194352d8d2bb32e18467c786a7fe15f1d44b233ea35633d076d5e737870b7139949d1ab6318 @@ -11185,7 +11412,7 @@ __metadata: languageName: node linkType: hard -"string-width-cjs@npm:string-width@^4.2.0, string-width@npm:^4.0.0, string-width@npm:^4.1.0, string-width@npm:^4.2.0, string-width@npm:^4.2.2, string-width@npm:^4.2.3": +"string-width-cjs@npm:string-width@^4.2.0, string-width@npm:^1.0.2 || 2 || 3 || 4, string-width@npm:^4.0.0, string-width@npm:^4.1.0, string-width@npm:^4.2.0, string-width@npm:^4.2.2, string-width@npm:^4.2.3": version: 4.2.3 resolution: "string-width@npm:4.2.3" dependencies: @@ -12225,6 +12452,8 @@ __metadata: "@adiwajshing/keyed-db": ^0.2.4 "@aws-sdk/client-s3": ^3.633.0 "@aws-sdk/s3-request-presigner": ^3.633.0 + "@grpc/grpc-js": ^1.12.4 + "@grpc/proto-loader": ^0.7.13 "@nestjs/axios": ^3.0.2 "@nestjs/cli": ^9.0.0 "@nestjs/common": ^9.0.9 @@ -12267,6 +12496,8 @@ __metadata: express-basic-auth: ^1.2.1 file-type: 16.5.4 fs-extra: ^11.2.0 + google-protobuf: ^3.21.4 + grpc-tools: ^1.12.4 https-proxy-agent: ^7.0.0 jest: ^29.7.0 joi: ^17.13.3 @@ -12285,6 +12516,7 @@ __metadata: prettier: ^1.19.1 pretty-bytes: 5.6.0 promise-retry: ^2.0.1 + protoc-gen-ts: ^0.8.7 puppeteer: ^23.6.0 qrcode: ^1.5.1 qrcode-terminal: ^0.12.0 @@ -12494,6 +12726,15 @@ __metadata: languageName: node linkType: hard +"wide-align@npm:^1.1.2": + version: 1.1.5 + resolution: "wide-align@npm:1.1.5" + dependencies: + string-width: ^1.0.2 || 2 || 3 || 4 + checksum: d5fc37cd561f9daee3c80e03b92ed3e84d80dde3365a8767263d03dacfc8fa06b065ffe1df00d8c2a09f731482fcacae745abfbb478d4af36d0a891fad4834d3 + languageName: node + linkType: hard + "widest-line@npm:^3.1.0": version: 3.1.0 resolution: "widest-line@npm:3.1.0"