Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions infrastructure/evault-core/src/config/database.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -21,4 +21,12 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
})
58 changes: 48 additions & 10 deletions platforms/cerberus/src/controllers/WebhookController.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -113,21 +113,59 @@ export class WebhookController {
Array.isArray(local.data.participants)
) {
console.log("Processing participants:", local.data.participants);

// Use Promise.allSettled with timeout to prevent webhook hang
const participantPromises = local.data.participants.map(
async (ref: string) => {
if (ref && typeof ref === "string") {
const userId = ref.split("(")[1].split(")")[0];
console.log("Extracted userId:", userId);
return await this.userService.getUserById(userId);
async (ref: string, index: number) => {
if (!ref || typeof ref !== "string") {
return null;
}

try {
const userId = ref.split("(")[1]?.split(")")[0];
if (!userId) {
console.warn(`⚠️ Could not extract userId from ref: ${ref}`);
return null;
}

console.log(`Extracted userId [${index}]: ${userId}`);

// Add 5-second timeout to prevent indefinite hang
const timeoutPromise = new Promise<null>((_, reject) =>
setTimeout(() => reject(new Error(`Timeout loading user ${userId}`)), 5000)
);

const userPromise = this.userService.userRepository.findOne({
where: { id: userId },
// Skip heavy relations in webhook context - only need basic user data
});

const user = await Promise.race([userPromise, timeoutPromise]);

if (user) {
console.log(`✅ Loaded user [${index}]: ${userId}`);
} else {
console.warn(`⚠️ User not found [${index}]: ${userId}`);
}

return user;
} catch (error) {
console.error(`❌ Error loading participant [${index}]:`, error instanceof Error ? error.message : error);
return null;
}
return null;
}
);

participants = (
await Promise.all(participantPromises)
).filter((user): user is User => user !== null);
console.log("Found participants:", participants.length);
// Use allSettled to handle failures gracefully without blocking
const settledResults = await Promise.allSettled(participantPromises);

participants = settledResults
.filter((result): result is PromiseFulfilledResult<User | null> =>
result.status === 'fulfilled' && result.value !== null
)
.map(result => result.value as User);

console.log(`Found ${participants.length} participants (${settledResults.filter(r => r.status === 'rejected').length} failed)`);
}

// Process admins - filter out nulls and extract IDs
Expand Down
13 changes: 13 additions & 0 deletions platforms/cerberus/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -34,5 +34,18 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
// Maximum number of connections in pool
max: 10,
// Minimum number of connections in pool
min: 2,
// Maximum time (ms) a connection can be idle before being released
idleTimeoutMillis: 30000,
// Maximum time (ms) to wait for a connection from pool
connectionTimeoutMillis: 5000,
// Query timeout (ms) - fail queries that take too long
statement_timeout: 10000,
},
});

169 changes: 155 additions & 14 deletions platforms/cerberus/src/web3adapter/watchers/subscriber.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -109,23 +109,39 @@ export class PostgresSubscriber implements EntitySubscriberInterface {

/**
* Called after entity update.
* NOTE: We pass metadata to handleChangeWithReload so the entity reload happens
* AFTER the transaction commits (inside setTimeout), avoiding stale/partial reads.
*/
async afterUpdate(event: UpdateEvent<any>) {
let entity = event.entity;
if (entity) {
entity = (await this.enrichEntity(
entity,
event.metadata.tableName,
event.metadata.target
)) as ObjectLiteral;
// Try different ways to get the entity ID
let entityId = event.entity?.id || event.databaseEntity?.id;

if (!entityId && event.entity) {
// Look for common ID field names
entityId = event.entity.id || event.entity.Id || event.entity.ID || event.entity._id;
}
this.handleChange(
// @ts-ignore
entity ?? event.entityId,
event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s"
);

if (!entityId) {
console.warn(`⚠️ afterUpdate: Could not determine entity ID for ${event.metadata.tableName}`);
return;
}

const entityName = typeof event.metadata.target === 'function'
? event.metadata.target.name
: event.metadata.target;

const tableName = event.metadata.tableName.endsWith("s")
? event.metadata.tableName
: event.metadata.tableName + "s";

// Pass reload metadata instead of entity - actual DB read happens in setTimeout
this.handleChangeWithReload({
entityId,
tableName,
relations: this.getRelationsForEntity(entityName),
tableTarget: event.metadata.target,
rawTableName: event.metadata.tableName,
});
}

/**
Expand All@@ -152,6 +168,131 @@ export class PostgresSubscriber implements EntitySubscriberInterface {
// This prevents the error when trying to access entity.id
}

/**
* Handle update changes by reloading entity AFTER transaction commits.
* This avoids stale/partial reads that occur when we use event.entity which only contains changed fields.
*/
private async handleChangeWithReload(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

console.log(`🔍 handleChangeWithReload called for: ${tableName}, entityId: ${entityId}`);

// Check if this is a junction table - skip for now
if (tableName === "group_participants") {
return;
}

// @ts-ignore
const junctionInfo = JUNCTION_TABLE_MAP[tableName];
if (junctionInfo) {
// Junction tables handled separately
return;
}

// Small delay to ensure transaction has committed before we read
// Groups and messages sync quickly (50ms), other entities use standard delay
const delayMs = (tableName.toLowerCase() === "groups" || tableName.toLowerCase() === "messages") ? 50 : 3_000;

setTimeout(async () => {
try {
await this.executeReloadAndSend(params);
} catch (error) {
console.error(`❌ Error in handleChangeWithReload setTimeout for ${tableName}:`, error);
}
}, delayMs);
}

/**
* Execute the entity reload and send webhook - called from within setTimeout
* when transaction has definitely committed.
*/
private async executeReloadAndSend(params: {
entityId: string;
tableName: string;
relations: string[];
tableTarget: any;
rawTableName: string;
}): Promise<void> {
const { entityId, tableName, relations, tableTarget, rawTableName } = params;

// NOW reload entity - transaction has committed, data is fresh and complete
const repository = AppDataSource.getRepository(tableTarget);
let entity = await repository.findOne({
where: { id: entityId },
relations: relations.length > 0 ? relations : undefined
});

if (!entity) {
console.warn(`⚠️ executeReloadAndSend: Entity ${entityId} not found after reload`);
return;
}

// Enrich entity with additional data
entity = (await this.enrichEntity(
entity,
rawTableName,
tableTarget
)) as ObjectLiteral;

// Convert to plain data
const data = this.entityToPlain(entity);

if (!data.id) {
return;
}

// For Message entities, only process system messages
if (tableName === "messages") {
const isSystemMessage = data.text && data.text.includes('$$system-message$$');
if (!isSystemMessage) {
return;
}
}

let globalId = await this.adapter.mappingDb.getGlobalId(entityId);
globalId = globalId ?? "";

if (this.adapter.lockedIds.includes(globalId)) {
console.log("Entity already locked, skipping:", globalId, entityId);
return;
}

if (this.adapter.lockedIds.includes(entityId)) {
console.log("Local entity locked (webhook created), skipping:", entityId);
return;
}

console.log(
"sending packet for global Id",
globalId,
entityId,
"table:",
tableName
);

// Log the full data being sent for system messages
if (tableName === "messages") {
console.log("📤 [SUBSCRIBER] Sending message data:");
console.log(" - Data keys:", Object.keys(data));
console.log(" - Data.sender:", data.sender);
console.log(" - Data.group:", data.group ? `Group ID: ${data.group.id}` : "null");
console.log(" - Data.text (first 100):", data.text?.substring(0, 100));
console.log(" - Data.isSystemMessage:", data.isSystemMessage);
}

const envelope = await this.adapter.handleChange({
data,
tableName: tableName.toLowerCase(),
});
console.log("📥 [SUBSCRIBER] Envelope response:", envelope);
}

/**
* Handle entity changes and send to web3adapter
*/
Expand Down
8 changes: 8 additions & 0 deletions platforms/dreamsync-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -30,6 +30,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/eCurrency-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -22,6 +22,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
Expand Down
8 changes: 8 additions & 0 deletions platforms/eReputation-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -26,6 +26,14 @@ export const dataSourceOptions: DataSourceOptions = {
migrations: [path.join(__dirname, "migrations", "*.ts")],
logging: process.env.NODE_ENV === "development",
subscribers: [PostgresSubscriber],
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/emover-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/esigner-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -27,6 +27,14 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});


8 changes: 8 additions & 0 deletions platforms/evoting-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -29,6 +29,14 @@ export const dataSourceOptions: DataSourceOptions = {
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
};

export const AppDataSource = new DataSource(dataSourceOptions);
8 changes: 8 additions & 0 deletions platforms/file-manager-api/src/database/data-source.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -41,5 +41,13 @@ export const AppDataSource = new DataSource({
ca: process.env.DB_CA_CERT,
}
: false,
// Connection pool configuration to prevent exhaustion
extra: {
max: 10,
min: 2,
idleTimeoutMillis: 30000,
connectionTimeoutMillis: 5000,
statement_timeout: 10000,
},
});

Loading