[core] me for stopped sessions
This commit is contained in:
1 parent
95f58df3c2
commit
cbc2d1ef9b
18 files changed
+176
-48
No files matched your search
@@ -1,3 +1,5 @@
|
||||
export abstract class DataStore {
|
||||
abstract init(sessionName?: string): Promise<void>;
|
||||
|
||||
abstract close(): Promise<any>;
|
||||
}
|
||||
@@ -2,6 +2,7 @@ import {
|
||||
BeforeApplicationShutdown,
|
||||
UnprocessableEntityException,
|
||||
} from '@nestjs/common';
|
||||
import { ISessionMeRepository } from '@waha/core/storage/ISessionMeRepository';
|
||||
import { WAHAWebhook } from '@waha/structures/webhooks.dto';
|
||||
import { waitUntil } from '@waha/utils/promiseTimeout';
|
||||
import { VERSION } from '@waha/version';
|
||||
@@ -27,6 +28,7 @@ export abstract class SessionManager implements BeforeApplicationShutdown {
|
||||
public store: any;
|
||||
public sessionAuthRepository: ISessionAuthRepository;
|
||||
public sessionConfigRepository: ISessionConfigRepository;
|
||||
protected sessionMeRepository: ISessionMeRepository;
|
||||
public events: EventEmitter;
|
||||
WAIT_STATUS_INTERVAL = 500;
|
||||
WAIT_STATUS_TIMEOUT = 5_000;
|
||||
@@ -37,8 +39,6 @@ export abstract class SessionManager implements BeforeApplicationShutdown {
|
||||
|
||||
protected abstract get WebhookConductorClass(): typeof WebhookConductor;
|
||||
|
||||
abstract beforeApplicationShutdown(signal?: string);
|
||||
|
||||
//
|
||||
// API Methods
|
||||
//
|
||||
@@ -110,4 +110,8 @@ export abstract class SessionManager implements BeforeApplicationShutdown {
|
||||
}
|
||||
return session;
|
||||
}
|
||||
|
||||
async beforeApplicationShutdown(signal?: string) {
|
||||
this.events.removeAllListeners();
|
||||
}
|
||||
}
|
||||
@@ -1,24 +1,4 @@
|
||||
export class Field {
|
||||
constructor(
|
||||
public fieldName: string,
|
||||
public type: string,
|
||||
) {}
|
||||
}
|
||||
|
||||
export class Index {
|
||||
constructor(
|
||||
public name: string,
|
||||
public columns: string[],
|
||||
) {}
|
||||
}
|
||||
|
||||
export class Schema {
|
||||
constructor(
|
||||
public name: string,
|
||||
public columns: Field[],
|
||||
public indexes: Index[],
|
||||
) {}
|
||||
}
|
||||
import { Field, Index, Schema } from '@waha/core/storage/sqlite3/Schema';
|
||||
|
||||
export const NOWEB_STORE_SCHEMA = [
|
||||
new Schema(
|
||||
|
||||
@@ -0,0 +1,28 @@
|
||||
import { BufferJSON } from '@adiwajshing/baileys/lib/Utils';
|
||||
import {
|
||||
convertProtobufToPlainObject,
|
||||
replaceLongsWithNumber,
|
||||
} from '@waha/core/engines/noweb/utils';
|
||||
import { Sqlite3KVRepository } from '@waha/core/storage/sqlite3/Sqlite3KVRepository';
|
||||
|
||||
/**
|
||||
* Key value repository with extra metadata
|
||||
* Add support for converting protobuf to plain object
|
||||
*/
|
||||
export class NOWEBSqlite3KVRepository<
|
||||
Entity,
|
||||
> extends Sqlite3KVRepository<Entity> {
|
||||
protected stringify(data: any): string {
|
||||
return JSON.stringify(data, BufferJSON.replacer);
|
||||
}
|
||||
|
||||
protected parse(row: any): any {
|
||||
return JSON.parse(row.data, BufferJSON.reviver);
|
||||
}
|
||||
|
||||
protected dump(entity: Entity) {
|
||||
const raw = convertProtobufToPlainObject(entity);
|
||||
replaceLongsWithNumber(raw);
|
||||
return super.dump(raw);
|
||||
}
|
||||
}
|
||||
@@ -1,10 +1,10 @@
|
||||
import { Chat } from '@adiwajshing/baileys';
|
||||
|
||||
import { IChatRepository } from '../IChatRepository';
|
||||
import { Sqlite3KVRepository } from './Sqlite3KVRepository';
|
||||
import { NOWEBSqlite3KVRepository } from './NOWEBSqlite3KVRepository';
|
||||
|
||||
export class Sqlite3ChatRepository
|
||||
extends Sqlite3KVRepository<Chat>
|
||||
extends NOWEBSqlite3KVRepository<Chat>
|
||||
implements IChatRepository
|
||||
{
|
||||
async getAllWithMessages(limit?: number, offset?: number): Promise<Chat[]> {
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
import { Contact } from '@adiwajshing/baileys';
|
||||
|
||||
import { IContactRepository } from '../IContactRepository';
|
||||
import { Sqlite3KVRepository } from './Sqlite3KVRepository';
|
||||
import { NOWEBSqlite3KVRepository } from './NOWEBSqlite3KVRepository';
|
||||
|
||||
export class Sqlite3ContactRepository
|
||||
extends Sqlite3KVRepository<Contact>
|
||||
extends NOWEBSqlite3KVRepository<Contact>
|
||||
implements IContactRepository {}
|
||||
@@ -4,10 +4,10 @@ import {
|
||||
} from '@adiwajshing/baileys/lib/Types/LabelAssociation';
|
||||
import { ILabelAssociationRepository } from '@waha/core/engines/noweb/store/ILabelAssociationsRepository';
|
||||
|
||||
import { Sqlite3KVRepository } from './Sqlite3KVRepository';
|
||||
import { NOWEBSqlite3KVRepository } from './NOWEBSqlite3KVRepository';
|
||||
|
||||
export class Sqlite3LabelAssociationsRepository
|
||||
extends Sqlite3KVRepository<LabelAssociation>
|
||||
extends NOWEBSqlite3KVRepository<LabelAssociation>
|
||||
implements ILabelAssociationRepository
|
||||
{
|
||||
async deleteOne(association: LabelAssociation): Promise<void> {
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
import { Label } from '@adiwajshing/baileys/lib/Types/Label';
|
||||
import { ILabelsRepository } from '@waha/core/engines/noweb/store/ILabelsRepository';
|
||||
|
||||
import { Sqlite3KVRepository } from './Sqlite3KVRepository';
|
||||
import { NOWEBSqlite3KVRepository } from './NOWEBSqlite3KVRepository';
|
||||
|
||||
export class Sqlite3LabelsRepository
|
||||
extends Sqlite3KVRepository<Label>
|
||||
extends NOWEBSqlite3KVRepository<Label>
|
||||
implements ILabelsRepository {}
|
||||
@@ -1,8 +1,8 @@
|
||||
import { IMessagesRepository } from '../IMessagesRepository';
|
||||
import { Sqlite3KVRepository } from './Sqlite3KVRepository';
|
||||
import { NOWEBSqlite3KVRepository } from './NOWEBSqlite3KVRepository';
|
||||
|
||||
export class Sqlite3MessagesRepository
|
||||
extends Sqlite3KVRepository<any>
|
||||
extends NOWEBSqlite3KVRepository<any>
|
||||
implements IMessagesRepository
|
||||
{
|
||||
upsert(messages: any[]): Promise<void> {
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Schema } from '../Schema';
|
||||
import { Schema } from '@waha/core/storage/sqlite3/Schema';
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-var-requires
|
||||
const Database = require('better-sqlite3');
|
||||
|
||||
@@ -4,9 +4,10 @@ import { ILabelAssociationRepository } from '@waha/core/engines/noweb/store/ILab
|
||||
import { ILabelsRepository } from '@waha/core/engines/noweb/store/ILabelsRepository';
|
||||
import { Sqlite3LabelAssociationsRepository } from '@waha/core/engines/noweb/store/sqlite3/Sqlite3LabelAssociationsRepository';
|
||||
import { Sqlite3LabelsRepository } from '@waha/core/engines/noweb/store/sqlite3/Sqlite3LabelsRepository';
|
||||
import { Field, Index, Schema } from '@waha/core/storage/sqlite3/Schema';
|
||||
|
||||
import { INowebStorage } from '../INowebStorage';
|
||||
import { Field, Index, NOWEB_STORE_SCHEMA, Schema } from '../Schema';
|
||||
import { NOWEB_STORE_SCHEMA } from '../Schema';
|
||||
import { Sqlite3ChatRepository } from './Sqlite3ChatRepository';
|
||||
import { Sqlite3ContactRepository } from './Sqlite3ContactRepository';
|
||||
import { Sqlite3MessagesRepository } from './Sqlite3MessagesRepository';
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
import { MeInfo } from '../../structures/sessions.dto';
|
||||
|
||||
export abstract class ISessionMeRepository {
|
||||
abstract upsertMe(sessionName: string, me: MeInfo | null): Promise<void>;
|
||||
|
||||
abstract getMe(sessionName: string): Promise<MeInfo | null>;
|
||||
|
||||
abstract removeMe(sessionName: string): Promise<void>;
|
||||
|
||||
abstract init(): Promise<void>;
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
import { NOWEBSqlite3KVRepository } from '@waha/core/engines/noweb/store/sqlite3/NOWEBSqlite3KVRepository';
|
||||
import { Sqlite3SchemaValidation } from '@waha/core/engines/noweb/store/sqlite3/Sqlite3SchemaValidation';
|
||||
import { ISessionMeRepository } from '@waha/core/storage/ISessionMeRepository';
|
||||
import { LocalStore } from '@waha/core/storage/LocalStore';
|
||||
import { Field, Index, Schema } from '@waha/core/storage/sqlite3/Schema';
|
||||
import { MeInfo } from '@waha/structures/sessions.dto';
|
||||
import * as path from 'path';
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-var-requires
|
||||
const Database = require('better-sqlite3');
|
||||
|
||||
const SCHEMA = new Schema(
|
||||
'me',
|
||||
[new Field('id', 'TEXT'), new Field('data', 'TEXT')],
|
||||
[new Index('me_id_index', ['id'])],
|
||||
);
|
||||
|
||||
class SessionMeInfo {
|
||||
id: string;
|
||||
me?: MeInfo;
|
||||
}
|
||||
|
||||
export class LocalSessionMeRepository
|
||||
extends NOWEBSqlite3KVRepository<SessionMeInfo>
|
||||
implements ISessionMeRepository
|
||||
{
|
||||
constructor(store: LocalStore) {
|
||||
const db = store.getWAHADatabase();
|
||||
super(db, SCHEMA);
|
||||
}
|
||||
|
||||
upsertMe(sessionName: string, me: MeInfo): Promise<void> {
|
||||
return this.upsertOne({ id: sessionName, me: me });
|
||||
}
|
||||
|
||||
async getMe(sessionName: string): Promise<MeInfo | null> {
|
||||
const data = await this.getById(sessionName);
|
||||
return data?.me;
|
||||
}
|
||||
|
||||
removeMe(sessionName: string): Promise<void> {
|
||||
return this.deleteById(sessionName);
|
||||
}
|
||||
|
||||
async init(): Promise<void> {
|
||||
this.migrations();
|
||||
this.validateSchema();
|
||||
}
|
||||
|
||||
private migrations() {
|
||||
this.db.exec(
|
||||
'CREATE TABLE IF NOT EXISTS me (id TEXT PRIMARY KEY, data TEXT)',
|
||||
);
|
||||
this.db.exec('CREATE UNIQUE INDEX IF NOT EXISTS me_id_index ON me (id)');
|
||||
}
|
||||
|
||||
private validateSchema() {
|
||||
new Sqlite3SchemaValidation(SCHEMA, this.db).validate();
|
||||
}
|
||||
}
|
||||
@@ -22,4 +22,6 @@ export abstract class LocalStore extends DataStore {
|
||||
* Get the file path for a session
|
||||
*/
|
||||
abstract getFilePath(session: string, file: string): string;
|
||||
|
||||
abstract getWAHADatabase(): any;
|
||||
}
|
||||
@@ -5,6 +5,9 @@ import * as path from 'path';
|
||||
|
||||
import { LocalStore } from './LocalStore';
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-var-requires
|
||||
const Database = require('better-sqlite3');
|
||||
|
||||
export class LocalStoreCore extends LocalStore {
|
||||
protected readonly baseDirectory: string = path.join(
|
||||
os.tmpdir(),
|
||||
@@ -12,6 +15,7 @@ export class LocalStoreCore extends LocalStore {
|
||||
);
|
||||
|
||||
private readonly engine: string;
|
||||
private db: any;
|
||||
|
||||
constructor(engine: string) {
|
||||
super();
|
||||
@@ -53,4 +57,17 @@ export class LocalStoreCore extends LocalStore {
|
||||
const suffix = crypto.createHash('md5').update(name).digest('hex');
|
||||
return path.join(this.getEngineDirectory(), `${name}-${suffix}`);
|
||||
}
|
||||
|
||||
getWAHADatabase(): any {
|
||||
if (!this.db) {
|
||||
const engineDir = this.getEngineDirectory();
|
||||
const database = path.join(engineDir, 'waha.sqlite3');
|
||||
this.db = new Database(database);
|
||||
}
|
||||
return this.db;
|
||||
}
|
||||
|
||||
async close() {
|
||||
this.db?.close();
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,21 @@
|
||||
export class Field {
|
||||
constructor(
|
||||
public fieldName: string,
|
||||
public type: string,
|
||||
) {}
|
||||
}
|
||||
|
||||
export class Index {
|
||||
constructor(
|
||||
public name: string,
|
||||
public columns: string[],
|
||||
) {}
|
||||
}
|
||||
|
||||
export class Schema {
|
||||
constructor(
|
||||
public name: string,
|
||||
public columns: Field[],
|
||||
public indexes: Index[],
|
||||
) {}
|
||||
}
|
||||
+14
-12
@@ -1,12 +1,7 @@
|
||||
import { BufferJSON } from '@adiwajshing/baileys/lib/Utils';
|
||||
import {
|
||||
convertProtobufToPlainObject,
|
||||
replaceLongsWithNumber,
|
||||
} from '@waha/core/engines/noweb/utils';
|
||||
import { Database } from 'better-sqlite3';
|
||||
import Knex from 'knex';
|
||||
|
||||
import { Field, Schema } from '../Schema';
|
||||
import { Field, Schema } from './Schema';
|
||||
|
||||
/**
|
||||
* Key value repository with extra metadata
|
||||
@@ -52,16 +47,15 @@ export class Sqlite3KVRepository<Entity> {
|
||||
return this.get(query);
|
||||
}
|
||||
|
||||
private dump(entity: Entity) {
|
||||
protected dump(entity: Entity) {
|
||||
const data = {};
|
||||
const raw = convertProtobufToPlainObject(entity);
|
||||
replaceLongsWithNumber(raw);
|
||||
const raw = entity;
|
||||
for (const field of this.columns) {
|
||||
const fn = this.metadata.get(field.fieldName);
|
||||
if (fn) {
|
||||
data[field.fieldName] = fn(raw);
|
||||
} else if (field.fieldName == 'data') {
|
||||
data['data'] = JSON.stringify(raw, BufferJSON.replacer);
|
||||
data['data'] = this.stringify(raw);
|
||||
} else {
|
||||
data[field.fieldName] = raw[field.fieldName];
|
||||
}
|
||||
@@ -143,7 +137,7 @@ export class Sqlite3KVRepository<Entity> {
|
||||
const sql = query.toSQL().sql;
|
||||
const bind = query.toSQL().bindings;
|
||||
const rows: any[] = this.db.prepare(sql).all(bind);
|
||||
return rows.map((row) => JSON.parse(row.data, BufferJSON.reviver));
|
||||
return rows.map((row) => this.parse(row));
|
||||
}
|
||||
|
||||
protected async get(query: any) {
|
||||
@@ -153,7 +147,7 @@ export class Sqlite3KVRepository<Entity> {
|
||||
if (!row) {
|
||||
return null;
|
||||
}
|
||||
return JSON.parse(row.data, BufferJSON.reviver);
|
||||
return this.parse(row);
|
||||
}
|
||||
|
||||
protected async run(query: any) {
|
||||
@@ -161,4 +155,12 @@ export class Sqlite3KVRepository<Entity> {
|
||||
const bind = query.toSQL().bindings;
|
||||
return this.db.prepare(sql).run(bind);
|
||||
}
|
||||
|
||||
protected stringify(data: any): string {
|
||||
return JSON.stringify(data);
|
||||
}
|
||||
|
||||
protected parse(row: any) {
|
||||
return JSON.parse(row.data);
|
||||
}
|
||||
}
|
||||
@@ -132,7 +132,7 @@ export class WAHAWebhook {
|
||||
| object;
|
||||
}
|
||||
|
||||
class WAHAWebhookSessionStatus extends WAHAWebhook {
|
||||
export class WAHAWebhookSessionStatus extends WAHAWebhook {
|
||||
@ApiProperty({
|
||||
description: 'The event is triggered when the session status changes.',
|
||||
})
|
||||
|
||||
Reference in new issue
Block a user