-
Notifications
You must be signed in to change notification settings - Fork 198
Expand file tree
/
Copy pathsqlite.js
More file actions
131 lines (116 loc) · 4.34 KB
/
Copy pathsqlite.js
File metadata and controls
131 lines (116 loc) · 4.34 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
/**
* Copyright (c) Forward Email LLC
* SPDX-License-Identifier: BUSL-1.1
*/
// eslint-disable-next-line import/no-unassigned-import
require('#helpers/polyfill-towellformed');
// eslint-disable-next-line import/no-unassigned-import
require('#config/env');
// eslint-disable-next-line import/no-unassigned-import
require('#config/mongoose');
const process = require('node:process');
const { promisify } = require('node:util');
const { setTimeout } = require('node:timers/promises');
const Graceful = require('@ladjs/graceful');
const Redis = require('@ladjs/redis');
const ip = require('ip');
const mongoose = require('mongoose');
const ms = require('ms');
const sharedConfig = require('@ladjs/shared-config');
const SQLite = require('./sqlite-server');
const closeDatabase = require('#helpers/close-database');
const logger = require('#helpers/logger');
const setupMongoose = require('#helpers/setup-mongoose');
const imapSharedConfig = sharedConfig('IMAP');
const client = new Redis(imapSharedConfig.redis, logger);
const subscriber = new Redis(imapSharedConfig.redis, logger);
client.setMaxListeners(0);
subscriber.setMaxListeners(0);
const sqlite = new SQLite({ client, subscriber });
const graceful = new Graceful({
mongooses: [mongoose],
servers: [sqlite.server],
redisClients: [client, subscriber],
logger,
timeoutMs: ms('1m'),
customHandlers: [
//
// Single sequential handler to enforce strict shutdown ordering.
// @ladjs/graceful runs customHandlers in parallel via Promise.all(),
// so we must consolidate into one async function to guarantee:
// 1. Stop accepting new work (isClosing)
// 2. Wait for in-flight requests to drain (refcount polling)
// 3. Close WebSocket server (no new connections)
// 4. Close all database handles (safe after drain)
//
async () => {
// 1. Signal all request handlers to reject new work
sqlite.isClosing = true;
// 2. Wait for in-flight requests to complete (poll refcounts)
// parsePayload checks isClosing and will reject new work,
// so only existing requests need to finish.
const drainStart = Date.now();
const drainTimeout = ms('30s');
if (sqlite.databaseMap && sqlite.databaseMap.size > 0) {
while (Date.now() - drainStart < drainTimeout) {
let activeRefs = 0;
for (const key of sqlite.databaseMap.keys()) {
const entry = sqlite.databaseMap._map.get(key);
if (entry && entry.refcount > 0) activeRefs += entry.refcount;
}
if (activeRefs === 0) break;
await setTimeout(250);
}
}
// 3. Close the WebSocket server (stops accepting new connections,
// terminates existing ones after in-flight work has drained)
try {
await promisify(sqlite.wss.close).bind(sqlite.wss)();
} catch (err) {
logger.error(err);
}
// 4. Close all normal databases (checkpoint WAL first for durability)
if (sqlite.databaseMap && sqlite.databaseMap.size > 0) {
await Promise.allSettled(
[...sqlite.databaseMap.keys()].map(async (key) => {
const db = sqlite.databaseMap.get(key);
if (db) {
sqlite.databaseMap.evict(key);
// Checkpoint WAL to main DB file before closing to prevent
// corruption if pm2 SIGKILLs before OS flushes WAL pages.
try {
if (db.open && !db.readonly) {
db.pragma('wal_checkpoint(PASSIVE)');
}
} catch (err) {
logger.error(err);
}
await closeDatabase(db);
}
})
);
}
// 5. Close all temporary databases
if (sqlite.temporaryDatabaseMap && sqlite.temporaryDatabaseMap.size > 0) {
await sqlite.temporaryDatabaseMap.closeAll();
}
}
]
});
graceful.listen();
(async () => {
try {
await sqlite.listen();
if (process.send) process.send('ready');
const { port } = sqlite.server.address();
logger.info(
`SQLite WebSocket server listening on ${port} (LAN: ${ip.address()}:${port})`,
{ hide_meta: true }
);
await setupMongoose(logger);
} catch (err) {
// Use timeout to prevent hanging if MongoDB pool is exhausted
await Promise.race([logger.error(err), setTimeout(5000)]);
process.exit(1);
}
})();