Skip to content

Commit a6661ca

Browse files
kriszypclaude
andcommitted
fix(openapi): query worker thread for REST OpenAPI spec on main thread
When the operations API runs on the main thread (post PR #355), resources are only registered on worker threads, so /api/openapi/rest returned an empty spec. This adds a cross-thread ITC mechanism: the main thread broadcasts RESOURCE_OPENAPI_REQUEST; the first worker with registered resources generates the spec via generateJsonApi and sends it back via RESOURCE_OPENAPI_RESPONSE. Falls back to local resources when running in single-thread mode (resources.size > 0). Resolves #299 (partial — addresses the OpenAPI endpoint concern) Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
1 parent ce2c466 commit a6661ca

5 files changed

Lines changed: 184 additions & 2 deletions

File tree

server/itc/serverHandlers.js

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ const serverItcHandlers = {
2121
[hdbTerms.ITC_EVENT_TYPES.SCHEMA]: schemaHandler,
2222
[hdbTerms.ITC_EVENT_TYPES.USER]: userHandler,
2323
[hdbTerms.ITC_EVENT_TYPES.COMPONENT_STATUS_REQUEST]: componentStatusRequestHandler,
24+
[hdbTerms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_REQUEST]: resourceOpenApiRequestHandler,
2425
};
2526

2627
/**
@@ -155,4 +156,46 @@ async function componentStatusRequestHandler(event) {
155156
}
156157
}
157158

159+
/**
160+
* Handles incoming requests for the REST OpenAPI spec from the main thread.
161+
* Generates the spec from the local resources (which are only registered on worker threads)
162+
* and sends it back to the requesting thread.
163+
*/
164+
async function resourceOpenApiRequestHandler(event) {
165+
try {
166+
const validate = validateEvent(event);
167+
if (validate) {
168+
hdbLogger.error(validate);
169+
return;
170+
}
171+
172+
hdbLogger.trace(`ITC resourceOpenApiRequestHandler received request:`, event);
173+
174+
const { resources } = require('../../resources/Resources.ts');
175+
if (!resources || resources.size === 0) {
176+
// This thread has no registered resources — don't respond so another worker can.
177+
return;
178+
}
179+
const { generateJsonApi } = require('../../resources/openApi.ts');
180+
const openapi = generateJsonApi(resources, event.message.serverHttpURL);
181+
182+
const originatorThreadId = event.message.originator;
183+
const responseMessage = {
184+
type: hdbTerms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_RESPONSE,
185+
message: {
186+
requestId: event.message.requestId,
187+
openapi,
188+
},
189+
};
190+
191+
if (!threads.sendToThread(originatorThreadId, responseMessage)) {
192+
hdbLogger.trace(
193+
`Dropping resource OpenAPI response for request ${event.message.requestId}: originator thread ${originatorThreadId} is unreachable`
194+
);
195+
}
196+
} catch (error) {
197+
hdbLogger.error('Error handling resource OpenAPI request:', error);
198+
}
199+
}
200+
158201
module.exports = serverItcHandlers;

server/operationsServer.ts

Lines changed: 47 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -30,6 +30,8 @@ type ParsedSqlObject = any;
3030
import { generateJsonApi } from '../resources/openApi.ts';
3131
import { Resources } from '../resources/Resources.ts';
3232
import { ServerError } from '../utility/errors/hdbError.ts';
33+
import { sendItcEvent } from './threads/itc.js';
34+
import { onMessageByType } from './threads/manageThreads.js';
3335

3436
const DEFAULT_HEADERS_TIMEOUT = 60000;
3537
const REQ_MAX_BODY_SIZE = env.get(terms.CONFIG_PARAMS.OPERATIONSAPI_NETWORK_MAXREQUESTBODYSIZE) ?? 1024 * 1024 * 1024; //this defaults to 1GB in bytes
@@ -211,12 +213,55 @@ function buildServer(isHttps: boolean, resources: Resources): FastifyInstance {
211213
return app;
212214
}
213215

216+
let nextOpenApiRequestId = 1;
217+
let openApiResponseListenerAttached = false;
218+
const pendingOpenApiRequests = new Map<number, (openapi: unknown) => void>();
219+
220+
function attachOpenApiResponseListener() {
221+
if (openApiResponseListenerAttached) return;
222+
onMessageByType(terms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_RESPONSE, ({ message }: any) => {
223+
const resolve = pendingOpenApiRequests.get(message.requestId);
224+
if (resolve) {
225+
pendingOpenApiRequests.delete(message.requestId);
226+
resolve(message.openapi);
227+
}
228+
});
229+
openApiResponseListenerAttached = true;
230+
}
231+
232+
function queryWorkerForOpenApi(serverHttpURL: string): Promise<unknown> {
233+
attachOpenApiResponseListener();
234+
const requestId = nextOpenApiRequestId++;
235+
return new Promise<unknown>((resolve, reject) => {
236+
const timeoutHandle = setTimeout(() => {
237+
pendingOpenApiRequests.delete(requestId);
238+
reject(new ServerError('Timeout fetching OpenAPI spec from worker thread', 503));
239+
}, 5000);
240+
pendingOpenApiRequests.set(requestId, (openapi) => {
241+
clearTimeout(timeoutHandle);
242+
resolve(openapi);
243+
});
244+
sendItcEvent({
245+
type: terms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_REQUEST,
246+
message: { requestId, serverHttpURL },
247+
}).catch((err: unknown) => {
248+
clearTimeout(timeoutHandle);
249+
pendingOpenApiRequests.delete(requestId);
250+
reject(err);
251+
});
252+
});
253+
}
254+
214255
function restOpenAPIHandler(resources: Resources) {
215256
const httpPort = env.get(terms.CONFIG_PARAMS.HTTP_PORT);
216257
const httpSecurePort = env.get(terms.CONFIG_PARAMS.HTTP_SECUREPORT);
217-
return (req: FastifyRequest & { hdb_user?: { role?: { permission?: { super_user: boolean } } } }) => {
258+
return async (req: FastifyRequest & { hdb_user?: { role?: { permission?: { super_user: boolean } } } }) => {
218259
if (req.hdb_user?.role?.permission?.super_user) {
219-
return generateJsonApi(resources, calculateRestHttpURL(httpPort, httpSecurePort, req));
260+
const serverHttpURL = calculateRestHttpURL(httpPort, httpSecurePort, req);
261+
if (resources.size > 0) {
262+
return generateJsonApi(resources, serverHttpURL);
263+
}
264+
return queryWorkerForOpenApi(serverHttpURL);
220265
} else {
221266
harperLogger.warn(
222267
`{"ip":"${req.socket.remoteAddress}", "error":"attempt to access /api/openapi/rest without being super_user"`

server/threads/manageThreads.js

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -119,6 +119,8 @@ listenersByType.set(hdbTerms.ITC_EVENT_TYPES.CHILD_STARTED, null);
119119
listenersByType.set(hdbTerms.ITC_EVENT_TYPES.SCHEMA, null);
120120
listenersByType.set(hdbTerms.ITC_EVENT_TYPES.USER, null);
121121
listenersByType.set(hdbTerms.ITC_EVENT_TYPES.COMPONENT_STATUS_REQUEST, null);
122+
listenersByType.set(hdbTerms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_REQUEST, null);
123+
listenersByType.set(hdbTerms.ITC_EVENT_TYPES.RESOURCE_OPENAPI_RESPONSE, null);
122124

123125
function startWorker(path, options = {}) {
124126
// Take a percentage of total memory to determine the max memory for each thread. The percentage is based

unitTests/server/itc/serverHandlers.test.js

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -221,4 +221,94 @@ describe('Test hdbChildIpcHandler module', () => {
221221
sendToThreadStub.restore();
222222
});
223223
});
224+
225+
describe('Test resourceOpenApiRequestHandler function', () => {
226+
let resource_openapi_handler;
227+
228+
before(() => {
229+
resource_openapi_handler = server_itc_handlers.__get__('resourceOpenApiRequestHandler');
230+
});
231+
232+
// Tests validation: invalid events should be rejected and logged
233+
it('logs error on invalid event (missing type)', async () => {
234+
const test_event = {
235+
message: { originator: 1, requestId: 42, serverHttpURL: 'http://localhost' },
236+
};
237+
await resource_openapi_handler(test_event);
238+
expect(log_error_stub).to.have.been.called;
239+
});
240+
241+
// Tests validation: invalid events should be rejected and logged
242+
it('logs error on invalid event (missing message)', async () => {
243+
const test_event = {
244+
type: 'resource_openapi_request',
245+
};
246+
await resource_openapi_handler(test_event);
247+
expect(log_error_stub).to.have.been.called;
248+
});
249+
250+
// Tests validation: invalid events should be rejected and logged
251+
it('logs error on invalid event (missing originator)', async () => {
252+
const test_event = {
253+
type: 'resource_openapi_request',
254+
message: { requestId: 42, serverHttpURL: 'http://localhost' },
255+
};
256+
await resource_openapi_handler(test_event);
257+
expect(log_error_stub).to.have.been.called;
258+
});
259+
260+
it('sends OpenAPI response directly when originator is reachable', async () => {
261+
sandbox.resetHistory();
262+
const sendToThreadStub = sandbox.stub(global.threads, 'sendToThread').returns(true);
263+
// Inject a minimal mock for generateJsonApi and a non-empty resources Map
264+
const mockOpenapi = { openapi: '3.0.3', paths: {} };
265+
const mockResources = new Map([['test', { path: 'test', Resource: { isError: false } }]]);
266+
server_itc_handlers.__set__('require', (path) => {
267+
if (path.includes('Resources')) return { resources: mockResources };
268+
if (path.includes('openApi')) return { generateJsonApi: () => mockOpenapi };
269+
return require(path);
270+
});
271+
272+
const test_event = {
273+
type: 'resource_openapi_request',
274+
message: { originator: 5, requestId: 99, serverHttpURL: 'http://localhost:9925' },
275+
};
276+
await resource_openapi_handler(test_event);
277+
278+
expect(sendToThreadStub).to.have.been.calledOnce;
279+
expect(sendToThreadStub.firstCall.args[0]).to.equal(5);
280+
const responseMessage = sendToThreadStub.firstCall.args[1];
281+
expect(responseMessage.type).to.equal('resource_openapi_response');
282+
expect(responseMessage.message.requestId).to.equal(99);
283+
expect(responseMessage.message.openapi).to.deep.equal(mockOpenapi);
284+
expect(log_error_stub).to.not.have.been.called;
285+
sendToThreadStub.restore();
286+
// Restore original require
287+
server_itc_handlers.__set__('require', require);
288+
});
289+
290+
it('drops response silently when originator is unreachable', async () => {
291+
sandbox.resetHistory();
292+
const sendToThreadStub = sandbox.stub(global.threads, 'sendToThread').returns(false);
293+
const mockResources = new Map([['test', { path: 'test', Resource: { isError: false } }]]);
294+
server_itc_handlers.__set__('require', (path) => {
295+
if (path.includes('Resources')) return { resources: mockResources };
296+
if (path.includes('openApi')) return { generateJsonApi: () => ({}) };
297+
return require(path);
298+
});
299+
300+
const test_event = {
301+
type: 'resource_openapi_request',
302+
message: { originator: 99, requestId: 7, serverHttpURL: 'http://localhost:9925' },
303+
};
304+
await resource_openapi_handler(test_event);
305+
306+
expect(sendToThreadStub).to.have.been.calledOnce;
307+
expect(log_error_stub).to.not.have.been.called;
308+
const traceCalls = log_trace_stub.getCalls().map((call) => String(call.args[0]));
309+
expect(traceCalls.some((msg) => msg.includes('Dropping resource OpenAPI response'))).to.be.true;
310+
sendToThreadStub.restore();
311+
server_itc_handlers.__set__('require', require);
312+
});
313+
});
224314
});

utility/hdbTerms.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -802,6 +802,8 @@ export const ITC_EVENT_TYPES = {
802802
START_JOB: 'start_job',
803803
COMPONENT_STATUS_REQUEST: 'component_status_request',
804804
COMPONENT_STATUS_RESPONSE: 'component_status_response',
805+
RESOURCE_OPENAPI_REQUEST: 'resource_openapi_request',
806+
RESOURCE_OPENAPI_RESPONSE: 'resource_openapi_response',
805807
} as const;
806808

807809
/** Supported thread types */

0 commit comments

Comments
 (0)