[core] NOWEB - merge @lid and @c.us messages in /overview and /messages
fix #1444, fix #1419, fix #1683, fix #1432
This commit is contained in:
1 parent
f37bf62a30
commit
d1e3b209e7
7 files changed
+170
-24
No files matched your search
@@ -2,6 +2,7 @@ import type { Chat } from '@adiwajshing/baileys';
|
||||
import { SqlKVRepository } from '@waha/core/storage/sql/SqlKVRepository';
|
||||
import { OverviewFilter } from '@waha/structures/chats.dto';
|
||||
import { PaginationParams } from '@waha/structures/pagination.dto';
|
||||
import { Knex } from 'knex';
|
||||
|
||||
export class SqlChatMethods {
|
||||
constructor(private repository: SqlKVRepository<any>) {}
|
||||
@@ -11,22 +12,83 @@ export class SqlChatMethods {
|
||||
broadcast: boolean,
|
||||
filter?: OverviewFilter,
|
||||
): Promise<Chat[]> {
|
||||
// Get chats with conversationTimestamp is not Null
|
||||
let query = this.repository.select().whereNotNull('conversationTimestamp');
|
||||
const knex = this.repository.getKnex();
|
||||
const tableName = this.repository.table;
|
||||
const baseQuery = this.repository
|
||||
.select()
|
||||
.whereNotNull(`${tableName}.conversationTimestamp`);
|
||||
|
||||
if (!broadcast) {
|
||||
// filter out chat by id if it ends at @newsletter or @broadcast
|
||||
query = query
|
||||
.andWhereNot('id', 'like', '%@broadcast')
|
||||
.andWhereNot('id', 'like', '%@newsletter');
|
||||
const annotatedQuery = this.annotateWithPnJid(knex, baseQuery, {
|
||||
tableName,
|
||||
broadcast,
|
||||
filter,
|
||||
}).as('annotated_chats');
|
||||
|
||||
const dedupedQuery = knex
|
||||
.select('*')
|
||||
.from(
|
||||
knex
|
||||
.select(
|
||||
'annotated_chats.*',
|
||||
knex.raw(
|
||||
'ROW_NUMBER() OVER (PARTITION BY annotated_chats.primary_jid ORDER BY annotated_chats.primary_priority ASC, annotated_chats."conversationTimestamp" DESC) as __rownum',
|
||||
),
|
||||
)
|
||||
.from(annotatedQuery)
|
||||
.as('ranked_chats'),
|
||||
)
|
||||
.where('__rownum', 1);
|
||||
|
||||
const pagedQuery = this.repository.pagination(dedupedQuery, pagination);
|
||||
const rows = await pagedQuery;
|
||||
|
||||
return rows.map((row) => {
|
||||
const chat = this.repository.parse(row);
|
||||
if (row.primary_jid) {
|
||||
chat.id = row.primary_jid;
|
||||
}
|
||||
return chat;
|
||||
});
|
||||
}
|
||||
|
||||
private annotateWithPnJid(
|
||||
knex: Knex,
|
||||
query: Knex.QueryBuilder,
|
||||
opts: {
|
||||
tableName: string;
|
||||
broadcast: boolean;
|
||||
filter?: OverviewFilter;
|
||||
},
|
||||
) {
|
||||
const pnExpr = this.buildPnJidExpr(opts.tableName);
|
||||
let annotated = query
|
||||
.leftJoin('lid_map', 'lid_map.id', `${opts.tableName}.id`)
|
||||
.select(
|
||||
knex.raw(`${pnExpr} as primary_jid`),
|
||||
knex.raw(
|
||||
`CASE WHEN ${opts.tableName}.id LIKE '%@lid' THEN 1 ELSE 0 END as primary_priority`,
|
||||
),
|
||||
);
|
||||
|
||||
if (!opts.broadcast) {
|
||||
annotated = annotated
|
||||
.andWhereNot(`${opts.tableName}.id`, 'like', '%@broadcast')
|
||||
.andWhereNot(`${opts.tableName}.id`, 'like', '%@newsletter');
|
||||
}
|
||||
|
||||
// Filter by IDs if provided
|
||||
if (filter?.ids && filter.ids.length > 0) {
|
||||
query = query.whereIn('id', filter.ids);
|
||||
if (opts.filter?.ids && opts.filter.ids.length > 0) {
|
||||
annotated = annotated.andWhere((builder) => {
|
||||
builder
|
||||
.whereIn(`${opts.tableName}.id`, opts.filter.ids)
|
||||
.orWhereIn('lid_map.pn', opts.filter.ids);
|
||||
});
|
||||
}
|
||||
|
||||
query = this.repository.pagination(query, pagination);
|
||||
return await this.repository.all(query);
|
||||
return annotated.select(`${opts.tableName}.*`);
|
||||
}
|
||||
|
||||
private buildPnJidExpr(tableName: string) {
|
||||
const column = `"${tableName}"."id"`;
|
||||
return `CASE WHEN ${column} LIKE '%@lid' THEN COALESCE(lid_map.pn, ${column}) ELSE ${column} END`;
|
||||
}
|
||||
}
|
||||
@@ -1,11 +1,17 @@
|
||||
import { ALL_JID } from '@waha/core/engines/noweb/session.noweb.core';
|
||||
import { SqlKVRepository } from '@waha/core/storage/sql/SqlKVRepository';
|
||||
import { AckToStatus } from '@waha/core/utils/acks';
|
||||
import { isLidUser } from '@waha/core/utils/jids';
|
||||
import { GetChatMessagesFilter } from '@waha/structures/chats.dto';
|
||||
import { PaginationParams } from '@waha/structures/pagination.dto';
|
||||
import { Knex } from 'knex';
|
||||
import { INowebLidPNRepository } from '../INowebLidPNRepository';
|
||||
|
||||
export class SqlMessagesMethods {
|
||||
constructor(private repository: SqlKVRepository<any>) {}
|
||||
constructor(
|
||||
private repository: SqlKVRepository<any>,
|
||||
private lidRepository?: INowebLidPNRepository,
|
||||
) {}
|
||||
|
||||
upsert(messages: any[]): Promise<void> {
|
||||
return this.repository.upsertMany(messages);
|
||||
@@ -18,7 +24,8 @@ export class SqlMessagesMethods {
|
||||
): Promise<any[]> {
|
||||
let query = this.repository.select();
|
||||
if (jid !== ALL_JID) {
|
||||
query = this.repository.select().where({ jid: jid });
|
||||
const pnJid = await this.resolvePnJid(jid);
|
||||
query = this.applyPnJidFilter(query, pnJid);
|
||||
}
|
||||
if (filter['filter.timestamp.lte'] != null) {
|
||||
query = query.where(
|
||||
@@ -34,10 +41,11 @@ export class SqlMessagesMethods {
|
||||
filter['filter.timestamp.gte'],
|
||||
);
|
||||
}
|
||||
const dataColumn = `${this.repository.table}.data`;
|
||||
if (filter['filter.fromMe'] != null) {
|
||||
// filter by data json inside
|
||||
const [sql, value] = this.repository.filterJson(
|
||||
'data',
|
||||
dataColumn,
|
||||
'key.fromMe',
|
||||
filter['filter.fromMe'],
|
||||
);
|
||||
@@ -45,7 +53,11 @@ export class SqlMessagesMethods {
|
||||
}
|
||||
if (filter['filter.ack'] != null) {
|
||||
const status = AckToStatus(filter['filter.ack']);
|
||||
const [sql, value] = this.repository.filterJson('data', 'status', status);
|
||||
const [sql, value] = this.repository.filterJson(
|
||||
dataColumn,
|
||||
'status',
|
||||
status,
|
||||
);
|
||||
query = query.whereRaw(sql, [value]);
|
||||
}
|
||||
query = this.repository.pagination(query, pagination);
|
||||
@@ -56,7 +68,15 @@ export class SqlMessagesMethods {
|
||||
if (jid === ALL_JID) {
|
||||
return this.repository.getBy({ id: id });
|
||||
}
|
||||
return this.repository.getBy({ jid: jid, id: id });
|
||||
const tableName = this.repository.table;
|
||||
const pnJid = await this.resolvePnJid(jid);
|
||||
const baseQuery = this.repository.select().where(`${tableName}.id`, id);
|
||||
const query = this.applyPnJidFilter(baseQuery, pnJid);
|
||||
const rows = await query.limit(1);
|
||||
if (!rows.length) {
|
||||
return null;
|
||||
}
|
||||
return this.repository.parse(rows[0]);
|
||||
}
|
||||
|
||||
async updateByJidAndId(
|
||||
@@ -83,4 +103,39 @@ export class SqlMessagesMethods {
|
||||
deleteAllByJid(jid: string): Promise<void> {
|
||||
return this.repository.deleteBy({ jid: jid });
|
||||
}
|
||||
|
||||
private applyPnJidFilter(
|
||||
query: Knex.QueryBuilder,
|
||||
pnJid: string,
|
||||
): Knex.QueryBuilder {
|
||||
const tableName = this.repository.table;
|
||||
const pnExpr = this.buildPnJidExpr(tableName, 'jid');
|
||||
return query
|
||||
.leftJoin('lid_map', 'lid_map.id', `${tableName}.jid`)
|
||||
.whereRaw(`${pnExpr} = ?`, [pnJid]);
|
||||
}
|
||||
|
||||
private async resolvePnJid(jid: string): Promise<string> {
|
||||
if (!isLidUser(jid)) {
|
||||
return jid;
|
||||
}
|
||||
if (this.lidRepository) {
|
||||
const mapped = await this.lidRepository.findPNByLid(jid);
|
||||
if (mapped) {
|
||||
return mapped;
|
||||
}
|
||||
}
|
||||
const row = await this.repository
|
||||
.getKnex()
|
||||
.select('pn')
|
||||
.from('lid_map')
|
||||
.where('id', jid)
|
||||
.first();
|
||||
return row?.pn || jid;
|
||||
}
|
||||
|
||||
private buildPnJidExpr(tableName: string, column: string) {
|
||||
const colRef = `"${tableName}"."${column}"`;
|
||||
return `CASE WHEN ${colRef} LIKE '%@lid' THEN COALESCE(lid_map.pn, ${colRef}) ELSE ${colRef} END`;
|
||||
}
|
||||
}
|
||||
@@ -16,7 +16,7 @@ export class NOWEBSqlite3KVRepository<
|
||||
return JSON.stringify(data, esm.b.BufferJSON.replacer);
|
||||
}
|
||||
|
||||
protected parse(row: any): any {
|
||||
public parse(row: any): any {
|
||||
return JSON.parse(row.data, esm.b.BufferJSON.reviver);
|
||||
}
|
||||
|
||||
|
||||
@@ -6,17 +6,26 @@ import { PaginationParams } from '@waha/structures/pagination.dto';
|
||||
|
||||
import { IMessagesRepository } from '../IMessagesRepository';
|
||||
import { NOWEBSqlite3KVRepository } from './NOWEBSqlite3KVRepository';
|
||||
import { INowebLidPNRepository } from '../INowebLidPNRepository';
|
||||
import Knex from 'knex';
|
||||
|
||||
export class Sqlite3MessagesRepository
|
||||
extends NOWEBSqlite3KVRepository<any>
|
||||
implements IMessagesRepository
|
||||
{
|
||||
constructor(
|
||||
knex: Knex.Knex,
|
||||
private readonly lidRepository: INowebLidPNRepository,
|
||||
) {
|
||||
super(knex);
|
||||
}
|
||||
|
||||
get schema() {
|
||||
return NowebMessagesSchema;
|
||||
}
|
||||
|
||||
get methods() {
|
||||
return new SqlMessagesMethods(this);
|
||||
return new SqlMessagesMethods(this, this.lidRepository);
|
||||
}
|
||||
|
||||
get metadata() {
|
||||
|
||||
@@ -19,6 +19,7 @@ import { KNEX_SQLITE_CLIENT } from '@waha/core/env';
|
||||
export class Sqlite3Storage extends INowebStorage {
|
||||
private readonly tables: Schema[];
|
||||
private readonly knex: Knex.Knex;
|
||||
private lidRepository: INowebLidPNRepository | null = null;
|
||||
|
||||
constructor(filePath: string) {
|
||||
super();
|
||||
@@ -84,10 +85,17 @@ export class Sqlite3Storage extends INowebStorage {
|
||||
}
|
||||
|
||||
getMessagesRepository() {
|
||||
return new Sqlite3MessagesRepository(this.knex);
|
||||
return new Sqlite3MessagesRepository(this.knex, this.getLidRepository());
|
||||
}
|
||||
|
||||
getLidPNRepository(): INowebLidPNRepository {
|
||||
return new Sqlite3LidPNRepository(this.knex);
|
||||
return this.getLidRepository();
|
||||
}
|
||||
|
||||
private getLidRepository(): INowebLidPNRepository {
|
||||
if (!this.lidRepository) {
|
||||
this.lidRepository = new Sqlite3LidPNRepository(this.knex);
|
||||
}
|
||||
return this.lidRepository;
|
||||
}
|
||||
}
|
||||
@@ -210,8 +210,12 @@ export class SqlKVRepository<Entity> {
|
||||
/**
|
||||
* SQL helpers
|
||||
*/
|
||||
public getKnex(): Knex {
|
||||
return this.knex;
|
||||
}
|
||||
|
||||
public select() {
|
||||
return this.knex.select().from(this.table);
|
||||
return this.knex.select(`${this.table}.*`).from(this.table);
|
||||
}
|
||||
|
||||
protected delete() {
|
||||
@@ -219,7 +223,11 @@ export class SqlKVRepository<Entity> {
|
||||
}
|
||||
|
||||
public pagination(query: any, pagination?: PaginationParams) {
|
||||
const paginator = new this.Paginator(pagination, this.jsonQuery);
|
||||
const paginator = new this.Paginator(
|
||||
pagination,
|
||||
this.jsonQuery,
|
||||
this.table,
|
||||
);
|
||||
return paginator.apply(query);
|
||||
}
|
||||
|
||||
@@ -234,7 +242,7 @@ export class SqlKVRepository<Entity> {
|
||||
return JSON.stringify(data);
|
||||
}
|
||||
|
||||
protected parse(row: any) {
|
||||
public parse(row: any) {
|
||||
return JSON.parse(row.data);
|
||||
}
|
||||
|
||||
|
||||
@@ -56,8 +56,12 @@ export class KnexPaginator extends Paginator {
|
||||
constructor(
|
||||
pagination: PaginationParams,
|
||||
protected jsonQuery: IJsonQuery,
|
||||
protected tableName?: string,
|
||||
) {
|
||||
super(pagination);
|
||||
if (tableName) {
|
||||
this.dataField = `${tableName}.data`;
|
||||
}
|
||||
}
|
||||
|
||||
protected sort(query: any) {
|
||||
|
||||
Reference in new issue
Block a user