-
Notifications
You must be signed in to change notification settings - Fork 10
Expand file tree
/
Copy pathserverHandlers.js
More file actions
202 lines (181 loc) · 7.34 KB
/
Copy pathserverHandlers.js
File metadata and controls
202 lines (181 loc) · 7.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
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
'use strict';
/* global threads */
const hdbLogger = require('../../utility/logging/harper_logger.ts');
const hdbTerms = require('../../utility/hdbTerms.ts');
const cleanLmdbMap =
require('../../utility/lmdb/cleanLMDBMap.ts').default || require('../../utility/lmdb/cleanLMDBMap.ts');
const userSchema = require('../../security/user.ts');
const { validateEvent } = require('../threads/itc.js');
const harperBridge =
require('../../dataLayer/harperBridge/harperBridge.ts').default ||
require('../../dataLayer/harperBridge/harperBridge.ts');
const process = require('process');
const { resetDatabases } = require('../../resources/databases.ts');
/**
* This object/functions are passed to the ITC client instance and dynamically added as event handlers.
* @type {{schema: ((function(*): Promise<void>)|*), job: ((function(*): Promise<void>)|*), user: ((function(): Promise<void>)|*)}}
*/
const serverItcHandlers = {
[hdbTerms.ITC_EVENT_TYPES.SCHEMA]: schemaHandler,
[hdbTerms.ITC_EVENT_TYPES.USER]: userHandler,
[hdbTerms.ITC_EVENT_TYPES.COMPONENT_STATUS_REQUEST]: componentStatusRequestHandler,
[hdbTerms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_REQUEST]: resourceOpenApiRequestHandler,
};
/**
* Updates the global hdbSchema object.
* @param event
* @returns {Promise<void>}
*/
async function schemaHandler(event) {
const validate = validateEvent(event);
if (validate) {
hdbLogger.error(validate);
return;
}
hdbLogger.trace(`ITC schemaHandler received schema event:`, event);
await cleanLmdbMap(event.message);
await syncSchemaMetadata(event.message);
}
/**
* Switch statement to handle schema-related messages from other forked processes - i.e. if another process completes an
* operation that updates schema and, therefore, requires that we update the global schema value for the process
*
* @param msg
* @returns {Promise<void>}
*/
async function syncSchemaMetadata(msg) {
try {
// TODO: Eventually should indicate which database/table changed so we don't have to scan everything
let databases = resetDatabases();
if (msg.table && msg.database)
// wait for a write to finish to ensure all writes have been written
await databases[msg.database][msg.table].put(Symbol.for('write-verify'), null);
} catch (e) {
hdbLogger.error(e);
}
}
const userListeners = [];
/**
* Updates the global hdbUsers object by querying the hdbRole table.
* @param event
* @returns {Promise<void>}
*/
async function userHandler(event) {
try {
try {
harperBridge.resetReadTxn(hdbTerms.SYSTEM_SCHEMA_NAME, hdbTerms.SYSTEM_TABLE_NAMES.USER_TABLE_NAME);
harperBridge.resetReadTxn(hdbTerms.SYSTEM_SCHEMA_NAME, hdbTerms.SYSTEM_TABLE_NAMES.ROLE_TABLE_NAME);
} catch (error) {
// this can happen during tests, best to ignore
hdbLogger.warn(error);
}
const validate = validateEvent(event);
if (validate) {
hdbLogger.error(validate);
return;
}
hdbLogger.trace(`ITC userHandler ${hdbTerms.HDB_ITC_CLIENT_PREFIX}${process.pid} received user event:`, event);
await userSchema.setUsersWithRolesCache();
for (let listener of userListeners) listener();
} catch (err) {
hdbLogger.error(err);
}
}
userHandler.addListener = function (listener) {
userListeners.push(listener);
};
/**
* Handles incoming requests for component status from inter-thread communication (ITC).
* Validates the event, retrieves the current thread's component statuses, and sends a response
* back to the originator thread with the requested information.
*
* @async
* @function componentStatusRequestHandler
* @param {Object} event - The event object containing the request details.
* @param {Object} event.message - The message object within the event.
* @param {string} event.message.originator - The identifier of the thread that originated the request.
* @param {string} event.message.requestId - The unique identifier for the request.
* @returns {Promise<void>} Sends a response back to the originator thread or logs an error if validation fails.
*/
async function componentStatusRequestHandler(event) {
try {
const validate = validateEvent(event);
if (validate) {
hdbLogger.error(validate);
return;
}
hdbLogger.trace(`ITC componentStatusRequestHandler received request:`, event);
// Get current thread's component status
const { internal } = require('../../components/status/index.ts');
const { getWorkerIndex } = require('../threads/manageThreads.js');
const componentStatuses = internal.componentStatusRegistry.getAllStatuses();
// Convert Map to array for serialization
const statusArray = Array.from(componentStatuses.entries());
// Get worker index and determine if this is the main thread
const workerIndex = getWorkerIndex();
const isMainThread = workerIndex === undefined;
// Send response directly back to the originating thread. validateEvent already
// ensures originator is present.
const originatorThreadId = event.message.originator;
const responseMessage = {
type: hdbTerms.ITC_EVENT_TYPES.COMPONENT_STATUS_RESPONSE,
message: {
requestId: event.message.requestId,
statuses: statusArray,
workerIndex: workerIndex,
isMainThread: isMainThread,
},
};
if (threads.sendToThread(originatorThreadId, responseMessage)) {
hdbLogger.trace(`Sent component status response directly to thread ${originatorThreadId}`);
} else {
// Originator's port is no longer in connectedPorts (thread exited / disconnected
// during the request). Dropping the response is correct — the originator is
// unreachable, and the collector's own timeout will handle the missing reply.
hdbLogger.trace(
`Dropping component status response for request ${event.message.requestId}: originator thread ${originatorThreadId} is unreachable`
);
}
} catch (error) {
hdbLogger.error('Error handling component status request:', error);
}
}
/**
* Handles incoming requests for the REST OpenAPI spec from the main thread.
* Generates the spec from the local resources (which are only registered on worker threads)
* and sends it back to the requesting thread.
*/
async function resourceOpenApiRequestHandler(event) {
try {
const validate = validateEvent(event);
if (validate) {
hdbLogger.error(validate);
return;
}
hdbLogger.trace(`ITC resourceOpenApiRequestHandler received request:`, event);
const { resources } = require('../../resources/Resources.ts');
// Only respond if this thread has registered resources. Job-type workers with an empty
// resources map must stay silent so that an app worker with real resources replies first.
// If no worker has resources the main thread gets a 503 after the timeout, which is a
// more honest response than silently returning an empty spec.
if (!resources || resources.size === 0) return;
const { generateJsonApi } = require('../../resources/openApi.ts');
const openapi = generateJsonApi(resources, event.message.serverHttpURL);
const originatorThreadId = event.message.originator;
const responseMessage = {
type: hdbTerms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_RESPONSE,
message: {
requestId: event.message.requestId,
openapi,
},
};
if (!threads.sendToThread(originatorThreadId, responseMessage)) {
hdbLogger.trace(
`Dropping resource OpenAPI response for request ${event.message.requestId}: originator thread ${originatorThreadId} is unreachable`
);
}
} catch (error) {
hdbLogger.error('Error handling resource OpenAPI request:', error);
}
}
module.exports = serverItcHandlers;