feat(sdk): emit runtime bridge dispatch observability events
This commit is contained in:
@@ -20,7 +20,7 @@ import { resolveGsdToolsPath } from './query-gsd-tools-path.js';
|
||||
import { createGSDToolsRuntime } from './query-gsd-tools-runtime.js';
|
||||
import { QueryCommandExecutor } from './query-command-executor.js';
|
||||
import { QueryHotpathMethods } from './query-hotpath-methods.js';
|
||||
import { QueryRuntimeBridge } from './query-runtime-bridge.js';
|
||||
import { QueryRuntimeBridge, type RuntimeBridgeOptions } from './query-runtime-bridge.js';
|
||||
|
||||
export { GSDToolsError } from './gsd-tools-error.js';
|
||||
|
||||
@@ -57,6 +57,8 @@ export class GSDTools {
|
||||
strictSdk?: boolean;
|
||||
/** Explicit subprocess bridge policy. Default false for SDK-native mode. */
|
||||
allowFallbackToSubprocess?: boolean;
|
||||
/** Structured runtime bridge dispatch observability callback. */
|
||||
onDispatchEvent?: RuntimeBridgeOptions['onDispatchEvent'];
|
||||
}) {
|
||||
this.projectDir = opts.projectDir;
|
||||
this.gsdToolsPath =
|
||||
@@ -77,6 +79,7 @@ export class GSDTools {
|
||||
execRawFallback: (legacyCommand, legacyArgs) => this.execRaw(legacyCommand, legacyArgs),
|
||||
strictSdk: opts.strictSdk,
|
||||
allowFallbackToSubprocess: opts.allowFallbackToSubprocess ?? false,
|
||||
onDispatchEvent: opts.onDispatchEvent,
|
||||
});
|
||||
|
||||
this.bridge = runtime.bridge;
|
||||
|
||||
@@ -25,20 +25,34 @@ export interface TransportPolicyLike {
|
||||
allowFallbackToSubprocess: boolean;
|
||||
}
|
||||
|
||||
export interface TransportDecision {
|
||||
dispatchMode: 'native' | 'subprocess';
|
||||
reason?: 'workstream_forced' | 'native_not_preferred' | 'native_unregistered' | 'native_failure_fallback';
|
||||
}
|
||||
|
||||
export class GSDTransport {
|
||||
constructor(
|
||||
private readonly registry: QueryRegistry,
|
||||
private readonly adapters: TransportAdapters,
|
||||
) {}
|
||||
|
||||
async run(request: TransportRequest, policy: TransportPolicyLike): Promise<unknown> {
|
||||
if (this.shouldUseNative(request, policy)) {
|
||||
async run(
|
||||
request: TransportRequest,
|
||||
policy: TransportPolicyLike,
|
||||
onDecision?: (decision: TransportDecision) => void,
|
||||
): Promise<unknown> {
|
||||
const useNative = this.shouldUseNative(request, policy);
|
||||
if (useNative) {
|
||||
try {
|
||||
const native = await this.adapters.dispatchNative(request);
|
||||
onDecision?.({ dispatchMode: 'native' });
|
||||
return this.projectNativeOutput(request, native.data);
|
||||
} catch (error) {
|
||||
if (this.shouldRethrowNativeError(error, policy)) throw error;
|
||||
onDecision?.({ dispatchMode: 'subprocess', reason: 'native_failure_fallback' });
|
||||
}
|
||||
} else {
|
||||
onDecision?.({ dispatchMode: 'subprocess', reason: this.subprocessReason(request, policy) });
|
||||
}
|
||||
|
||||
return this.dispatchSubprocess(request);
|
||||
@@ -49,6 +63,13 @@ export class GSDTransport {
|
||||
return !forceSubprocess && policy.preferNative && this.registry.has(request.registryCommand);
|
||||
}
|
||||
|
||||
private subprocessReason(request: TransportRequest, policy: TransportPolicyLike): TransportDecision['reason'] {
|
||||
if (request.workstream) return 'workstream_forced';
|
||||
if (!policy.preferNative) return 'native_not_preferred';
|
||||
if (!this.registry.has(request.registryCommand)) return 'native_unregistered';
|
||||
return 'native_not_preferred';
|
||||
}
|
||||
|
||||
private shouldRethrowNativeError(error: unknown, policy: TransportPolicyLike): boolean {
|
||||
if (!policy.allowFallbackToSubprocess) return true;
|
||||
// Do not subprocess-fallback after a timed-out native dispatch:
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { resolveTransportPolicy } from './gsd-transport-policy.js';
|
||||
import type { GSDTransport } from './gsd-transport.js';
|
||||
import type { GSDTransport, TransportDecision } from './gsd-transport.js';
|
||||
import type { TransportMode } from './gsd-transport-policy.js';
|
||||
|
||||
export interface QueryExecutionRequest {
|
||||
@@ -12,6 +12,7 @@ export interface QueryExecutionRequest {
|
||||
workstream?: string;
|
||||
preferNativeQuery: boolean;
|
||||
allowFallbackToSubprocess?: boolean;
|
||||
onTransportDecision?: (decision: TransportDecision) => void;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -39,6 +40,7 @@ export class QueryExecutionPolicy {
|
||||
allowFallbackToSubprocess:
|
||||
request.allowFallbackToSubprocess ?? policy.allowFallbackToSubprocess,
|
||||
},
|
||||
request.onTransportDecision,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ import { QueryNativeDirectAdapter } from './query-native-direct-adapter.js';
|
||||
import { QueryNativeHotpathAdapter } from './query-native-hotpath-adapter.js';
|
||||
import { formatQueryRawOutput } from './query-raw-output-projection.js';
|
||||
import { createQueryNativeErrorFactory, createQueryToolsErrorFactory } from './query-tools-error-factory.js';
|
||||
import { QueryRuntimeBridge } from './query-runtime-bridge.js';
|
||||
import { QueryRuntimeBridge, type RuntimeBridgeOptions } from './query-runtime-bridge.js';
|
||||
|
||||
export interface GSDToolsRuntime {
|
||||
bridge: QueryRuntimeBridge;
|
||||
@@ -25,6 +25,7 @@ export function createGSDToolsRuntime(opts: {
|
||||
execRawFallback: (legacyCommand: string, legacyArgs: string[]) => Promise<string>;
|
||||
strictSdk?: boolean;
|
||||
allowFallbackToSubprocess?: boolean;
|
||||
onDispatchEvent?: RuntimeBridgeOptions['onDispatchEvent'];
|
||||
}): GSDToolsRuntime {
|
||||
const registry = createRegistry(opts.eventStream, opts.sessionId);
|
||||
|
||||
@@ -74,6 +75,7 @@ export function createGSDToolsRuntime(opts: {
|
||||
{
|
||||
strictSdk: opts.strictSdk,
|
||||
allowFallbackToSubprocess: opts.allowFallbackToSubprocess,
|
||||
onDispatchEvent: opts.onDispatchEvent,
|
||||
},
|
||||
);
|
||||
|
||||
|
||||
102
sdk/src/query-runtime-bridge.test.ts
Normal file
102
sdk/src/query-runtime-bridge.test.ts
Normal file
@@ -0,0 +1,102 @@
|
||||
import { describe, it, expect, vi } from 'vitest';
|
||||
import { QueryRuntimeBridge } from './query-runtime-bridge.js';
|
||||
import { GSDToolsError } from './gsd-tools-error.js';
|
||||
|
||||
describe('QueryRuntimeBridge observability', () => {
|
||||
it('emits query_dispatch success event with transport decision', async () => {
|
||||
const onDispatchEvent = vi.fn();
|
||||
const executionPolicy = {
|
||||
execute: vi.fn(async (request: { onTransportDecision?: (d: unknown) => void }) => {
|
||||
request.onTransportDecision?.({ dispatchMode: 'subprocess', reason: 'workstream_forced' });
|
||||
return { ok: true };
|
||||
}),
|
||||
};
|
||||
|
||||
const bridge = new QueryRuntimeBridge(
|
||||
{ has: () => true } as never,
|
||||
executionPolicy as never,
|
||||
{ dispatch: vi.fn() } as never,
|
||||
() => true,
|
||||
{ onDispatchEvent },
|
||||
);
|
||||
|
||||
await bridge.execute({
|
||||
legacyCommand: 'state',
|
||||
legacyArgs: ['load'],
|
||||
registryCommand: 'state.load',
|
||||
registryArgs: [],
|
||||
mode: 'json',
|
||||
projectDir: '/tmp',
|
||||
workstream: 'ws-1',
|
||||
});
|
||||
|
||||
expect(onDispatchEvent).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
type: 'query_dispatch',
|
||||
command: 'state.load',
|
||||
dispatchMode: 'subprocess',
|
||||
reason: 'workstream_forced',
|
||||
outcome: 'success',
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it('emits query_dispatch error event with errorKind', async () => {
|
||||
const onDispatchEvent = vi.fn();
|
||||
const executionPolicy = {
|
||||
execute: vi.fn(async () => {
|
||||
throw GSDToolsError.timeout('timeout', 'state', ['load'], '', 500);
|
||||
}),
|
||||
};
|
||||
|
||||
const bridge = new QueryRuntimeBridge(
|
||||
{ has: () => true } as never,
|
||||
executionPolicy as never,
|
||||
{ dispatch: vi.fn() } as never,
|
||||
() => true,
|
||||
{ onDispatchEvent },
|
||||
);
|
||||
|
||||
await expect(
|
||||
bridge.execute({
|
||||
legacyCommand: 'state',
|
||||
legacyArgs: ['load'],
|
||||
registryCommand: 'state.load',
|
||||
registryArgs: [],
|
||||
mode: 'json',
|
||||
projectDir: '/tmp',
|
||||
}),
|
||||
).rejects.toThrow();
|
||||
|
||||
expect(onDispatchEvent).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
type: 'query_dispatch',
|
||||
command: 'state.load',
|
||||
outcome: 'error',
|
||||
errorKind: 'timeout',
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it('emits hotpath event', async () => {
|
||||
const onDispatchEvent = vi.fn();
|
||||
const bridge = new QueryRuntimeBridge(
|
||||
{ has: () => true } as never,
|
||||
{ execute: vi.fn() } as never,
|
||||
{ dispatch: vi.fn(async () => 'ok') } as never,
|
||||
() => true,
|
||||
{ onDispatchEvent },
|
||||
);
|
||||
|
||||
await bridge.dispatchHotpath('commit', ['msg'], 'commit', ['msg'], 'raw');
|
||||
|
||||
expect(onDispatchEvent).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
type: 'query_hotpath_dispatch',
|
||||
command: 'commit',
|
||||
dispatchMode: 'native_hotpath',
|
||||
outcome: 'success',
|
||||
}),
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -5,6 +5,7 @@ import { resolveQueryCommand } from './query/query-command-resolution-strategy.j
|
||||
import { QueryExecutionPolicy } from './query-execution-policy.js';
|
||||
import { QueryNativeHotpathAdapter } from './query-native-hotpath-adapter.js';
|
||||
import { GSDToolsError } from './gsd-tools-error.js';
|
||||
import type { TransportDecision } from './gsd-transport.js';
|
||||
|
||||
export interface RuntimeBridgeExecuteInput {
|
||||
legacyCommand: string;
|
||||
@@ -16,6 +17,39 @@ export interface RuntimeBridgeExecuteInput {
|
||||
workstream?: string;
|
||||
}
|
||||
|
||||
export interface RuntimeBridgeDispatchEvent {
|
||||
type: 'query_dispatch';
|
||||
command: string;
|
||||
legacyCommand: string;
|
||||
mode: TransportMode;
|
||||
dispatchMode: 'native' | 'subprocess' | 'native_hotpath';
|
||||
reason?: TransportDecision['reason'];
|
||||
durationMs: number;
|
||||
outcome: 'success' | 'error';
|
||||
errorKind?: 'timeout' | 'failure';
|
||||
}
|
||||
|
||||
export interface RuntimeBridgeHotpathEvent {
|
||||
type: 'query_hotpath_dispatch';
|
||||
command: string;
|
||||
legacyCommand: string;
|
||||
mode: TransportMode;
|
||||
dispatchMode: 'native_hotpath';
|
||||
durationMs: number;
|
||||
outcome: 'success' | 'error';
|
||||
errorKind?: 'timeout' | 'failure';
|
||||
}
|
||||
|
||||
export interface RuntimeBridgeEvent {
|
||||
type: 'query_dispatch' | 'query_hotpath_dispatch';
|
||||
}
|
||||
|
||||
export interface RuntimeBridgeOptions {
|
||||
strictSdk?: boolean;
|
||||
allowFallbackToSubprocess?: boolean;
|
||||
onDispatchEvent?: (event: RuntimeBridgeDispatchEvent | RuntimeBridgeHotpathEvent) => void;
|
||||
}
|
||||
|
||||
/**
|
||||
* SDK Runtime Bridge Module.
|
||||
* Owns dispatch routing through the execution policy seam and hotpath/native fallback behavior.
|
||||
@@ -26,10 +60,7 @@ export class QueryRuntimeBridge {
|
||||
private readonly executionPolicy: QueryExecutionPolicy,
|
||||
private readonly nativeHotpathAdapter: QueryNativeHotpathAdapter,
|
||||
private readonly shouldUseNativeQuery: () => boolean,
|
||||
private readonly options?: {
|
||||
strictSdk?: boolean;
|
||||
allowFallbackToSubprocess?: boolean;
|
||||
},
|
||||
private readonly options?: RuntimeBridgeOptions,
|
||||
) {}
|
||||
|
||||
getRegistry(): QueryRegistry {
|
||||
@@ -41,26 +72,71 @@ export class QueryRuntimeBridge {
|
||||
}
|
||||
|
||||
async execute(input: RuntimeBridgeExecuteInput): Promise<unknown> {
|
||||
const startedAt = Date.now();
|
||||
if (this.options?.strictSdk && !this.registry.has(input.registryCommand)) {
|
||||
throw GSDToolsError.failure(
|
||||
const error = GSDToolsError.failure(
|
||||
`Strict SDK mode: command '${input.registryCommand}' has no native adapter`,
|
||||
input.legacyCommand,
|
||||
input.legacyArgs,
|
||||
null,
|
||||
);
|
||||
this.options?.onDispatchEvent?.({
|
||||
type: 'query_dispatch',
|
||||
command: input.registryCommand,
|
||||
legacyCommand: input.legacyCommand,
|
||||
mode: input.mode,
|
||||
dispatchMode: 'subprocess',
|
||||
reason: 'native_unregistered',
|
||||
durationMs: Date.now() - startedAt,
|
||||
outcome: 'error',
|
||||
errorKind: 'failure',
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
|
||||
return this.executionPolicy.execute({
|
||||
legacyCommand: input.legacyCommand,
|
||||
legacyArgs: input.legacyArgs,
|
||||
registryCommand: input.registryCommand,
|
||||
registryArgs: input.registryArgs,
|
||||
mode: input.mode,
|
||||
projectDir: input.projectDir,
|
||||
workstream: input.workstream,
|
||||
preferNativeQuery: this.shouldUseNativeQuery(),
|
||||
allowFallbackToSubprocess: this.options?.allowFallbackToSubprocess,
|
||||
});
|
||||
let transportDecision: TransportDecision | undefined;
|
||||
try {
|
||||
const result = await this.executionPolicy.execute({
|
||||
legacyCommand: input.legacyCommand,
|
||||
legacyArgs: input.legacyArgs,
|
||||
registryCommand: input.registryCommand,
|
||||
registryArgs: input.registryArgs,
|
||||
mode: input.mode,
|
||||
projectDir: input.projectDir,
|
||||
workstream: input.workstream,
|
||||
preferNativeQuery: this.shouldUseNativeQuery(),
|
||||
allowFallbackToSubprocess: this.options?.allowFallbackToSubprocess,
|
||||
onTransportDecision: (decision) => {
|
||||
transportDecision = decision;
|
||||
},
|
||||
});
|
||||
|
||||
this.options?.onDispatchEvent?.({
|
||||
type: 'query_dispatch',
|
||||
command: input.registryCommand,
|
||||
legacyCommand: input.legacyCommand,
|
||||
mode: input.mode,
|
||||
dispatchMode: transportDecision?.dispatchMode ?? 'native',
|
||||
reason: transportDecision?.reason,
|
||||
durationMs: Date.now() - startedAt,
|
||||
outcome: 'success',
|
||||
});
|
||||
return result;
|
||||
} catch (error) {
|
||||
const kind = error instanceof GSDToolsError ? error.classification.kind : 'failure';
|
||||
this.options?.onDispatchEvent?.({
|
||||
type: 'query_dispatch',
|
||||
command: input.registryCommand,
|
||||
legacyCommand: input.legacyCommand,
|
||||
mode: input.mode,
|
||||
dispatchMode: transportDecision?.dispatchMode ?? 'native',
|
||||
reason: transportDecision?.reason,
|
||||
durationMs: Date.now() - startedAt,
|
||||
outcome: 'error',
|
||||
errorKind: kind,
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async dispatchHotpath(
|
||||
@@ -70,12 +146,38 @@ export class QueryRuntimeBridge {
|
||||
registryArgs: string[],
|
||||
mode: TransportMode,
|
||||
): Promise<unknown> {
|
||||
return this.nativeHotpathAdapter.dispatch(
|
||||
legacyCommand,
|
||||
legacyArgs,
|
||||
registryCommand,
|
||||
registryArgs,
|
||||
mode,
|
||||
);
|
||||
const startedAt = Date.now();
|
||||
try {
|
||||
const result = await this.nativeHotpathAdapter.dispatch(
|
||||
legacyCommand,
|
||||
legacyArgs,
|
||||
registryCommand,
|
||||
registryArgs,
|
||||
mode,
|
||||
);
|
||||
this.options?.onDispatchEvent?.({
|
||||
type: 'query_hotpath_dispatch',
|
||||
command: registryCommand,
|
||||
legacyCommand,
|
||||
mode,
|
||||
dispatchMode: 'native_hotpath',
|
||||
durationMs: Date.now() - startedAt,
|
||||
outcome: 'success',
|
||||
});
|
||||
return result;
|
||||
} catch (error) {
|
||||
const kind = error instanceof GSDToolsError ? error.classification.kind : 'failure';
|
||||
this.options?.onDispatchEvent?.({
|
||||
type: 'query_hotpath_dispatch',
|
||||
command: registryCommand,
|
||||
legacyCommand,
|
||||
mode,
|
||||
dispatchMode: 'native_hotpath',
|
||||
durationMs: Date.now() - startedAt,
|
||||
outcome: 'error',
|
||||
errorKind: kind,
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user