Skip to content

Commit c7c1192

Browse files
Claudehotlong
andauthored
Phase 2.2: Add event emission to ObjectQL hooks
- Add IRealtimeService import to ObjectQL engine - Add setRealtimeService() method to ObjectQL - Publish data.record.created events in insert() - Publish data.record.updated events in update() - Publish data.record.deleted events in delete() - Bridge realtime service in ObjectQLPlugin.start() Agent-Logs-Url: https://github.com/objectstack-ai/framework/sessions/6a38da2b-8234-4c22-b456-ec507a8f82f2 Co-authored-by: hotlong <50353452+hotlong@users.noreply.github.com>
1 parent e60cdd1 commit c7c1192

2 files changed

Lines changed: 114 additions & 8 deletions

File tree

packages/objectql/src/engine.ts

Lines changed: 99 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,18 @@
11
// Copyright (c) 2025 ObjectStack. Licensed under the Apache-2.0 license.
22

33
import { QueryAST, HookContext, ServiceObject } from '@objectstack/spec/data';
4-
import {
4+
import {
55
EngineQueryOptions,
6-
DataEngineInsertOptions,
7-
EngineUpdateOptions,
6+
DataEngineInsertOptions,
7+
EngineUpdateOptions,
88
EngineDeleteOptions,
99
EngineAggregateOptions,
10-
EngineCountOptions
10+
EngineCountOptions
1111
} from '@objectstack/spec/data';
1212
import { ExecutionContext, ExecutionContextSchema } from '@objectstack/spec/kernel';
1313
import { DriverInterface, IDataEngine, Logger, createLogger } from '@objectstack/core';
1414
import { CoreServiceName } from '@objectstack/spec/system';
15+
import { IRealtimeService, RealtimeEventPayload } from '@objectstack/spec/contracts';
1516
import { SchemaRegistry } from './registry.js';
1617

1718
export type HookHandler = (context: HookContext) => Promise<void> | void;
@@ -69,7 +70,7 @@ export class ObjectQL implements IDataEngine {
6970
private drivers = new Map<string, DriverInterface>();
7071
private defaultDriver: string | null = null;
7172
private logger: Logger;
72-
73+
7374
// Per-object hooks with priority support
7475
private hooks: Map<string, HookEntry[]> = new Map([
7576
['beforeFind', []], ['afterFind', []],
@@ -86,10 +87,13 @@ export class ObjectQL implements IDataEngine {
8687

8788
// Action registry: key = "objectName:actionName"
8889
private actions = new Map<string, { handler: (ctx: any) => Promise<any> | any; package?: string }>();
89-
90+
9091
// Host provided context additions (e.g. Server router)
9192
private hostContext: Record<string, any> = {};
9293

94+
// Realtime service for event publishing
95+
private realtimeService?: IRealtimeService;
96+
9397
constructor(hostContext: Record<string, any> = {}) {
9498
this.hostContext = hostContext;
9599
// Use provided logger or create a new one
@@ -514,8 +518,8 @@ export class ObjectQL implements IDataEngine {
514518
}
515519

516520
this.drivers.set(driver.name, driver);
517-
this.logger.info('Registered driver', {
518-
driverName: driver.name,
521+
this.logger.info('Registered driver', {
522+
driverName: driver.name,
519523
version: driver.version
520524
});
521525

@@ -525,6 +529,17 @@ export class ObjectQL implements IDataEngine {
525529
}
526530
}
527531

532+
/**
533+
* Set the realtime service for publishing data change events.
534+
* Should be called after kernel resolves the realtime service.
535+
*
536+
* @param service - An IRealtimeService instance for event publishing
537+
*/
538+
setRealtimeService(service: IRealtimeService): void {
539+
this.realtimeService = service;
540+
this.logger.info('RealtimeService configured for data events');
541+
}
542+
528543
/**
529544
* Helper to get object definition
530545
*/
@@ -883,6 +898,42 @@ export class ObjectQL implements IDataEngine {
883898
hookContext.result = result;
884899
await this.triggerHooks('afterInsert', hookContext);
885900

901+
// Publish data.record.created event to realtime service
902+
if (this.realtimeService) {
903+
try {
904+
if (Array.isArray(result)) {
905+
// Bulk insert - publish event for each record
906+
for (const record of result) {
907+
const event: RealtimeEventPayload = {
908+
type: 'data.record.created',
909+
object,
910+
payload: {
911+
recordId: record.id,
912+
after: record,
913+
},
914+
timestamp: new Date().toISOString(),
915+
};
916+
await this.realtimeService.publish(event);
917+
}
918+
this.logger.debug(`Published ${result.length} data.record.created events`, { object });
919+
} else {
920+
const event: RealtimeEventPayload = {
921+
type: 'data.record.created',
922+
object,
923+
payload: {
924+
recordId: result.id,
925+
after: result,
926+
},
927+
timestamp: new Date().toISOString(),
928+
};
929+
await this.realtimeService.publish(event);
930+
this.logger.debug('Published data.record.created event', { object, recordId: result.id });
931+
}
932+
} catch (error) {
933+
this.logger.warn('Failed to publish data event', { object, error });
934+
}
935+
}
936+
886937
return hookContext.result;
887938
} catch (e) {
888939
this.logger.error('Insert operation failed', e as Error, { object });
@@ -937,6 +988,27 @@ export class ObjectQL implements IDataEngine {
937988
hookContext.event = 'afterUpdate';
938989
hookContext.result = result;
939990
await this.triggerHooks('afterUpdate', hookContext);
991+
992+
// Publish data.record.updated event to realtime service
993+
if (this.realtimeService) {
994+
try {
995+
const event: RealtimeEventPayload = {
996+
type: 'data.record.updated',
997+
object,
998+
payload: {
999+
recordId: hookContext.input.id || result?.id,
1000+
changes: hookContext.input.data,
1001+
after: result,
1002+
},
1003+
timestamp: new Date().toISOString(),
1004+
};
1005+
await this.realtimeService.publish(event);
1006+
this.logger.debug('Published data.record.updated event', { object, recordId: hookContext.input.id });
1007+
} catch (error) {
1008+
this.logger.warn('Failed to publish data event', { object, error });
1009+
}
1010+
}
1011+
9401012
return hookContext.result;
9411013
} catch (e) {
9421014
this.logger.error('Update operation failed', e as Error, { object });
@@ -990,6 +1062,25 @@ export class ObjectQL implements IDataEngine {
9901062
hookContext.event = 'afterDelete';
9911063
hookContext.result = result;
9921064
await this.triggerHooks('afterDelete', hookContext);
1065+
1066+
// Publish data.record.deleted event to realtime service
1067+
if (this.realtimeService) {
1068+
try {
1069+
const event: RealtimeEventPayload = {
1070+
type: 'data.record.deleted',
1071+
object,
1072+
payload: {
1073+
recordId: hookContext.input.id || result?.id,
1074+
},
1075+
timestamp: new Date().toISOString(),
1076+
};
1077+
await this.realtimeService.publish(event);
1078+
this.logger.debug('Published data.record.deleted event', { object, recordId: hookContext.input.id });
1079+
} catch (error) {
1080+
this.logger.warn('Failed to publish data event', { object, error });
1081+
}
1082+
}
1083+
9931084
return hookContext.result;
9941085
} catch (e) {
9951086
this.logger.error('Delete operation failed', e as Error, { object });

packages/objectql/src/plugin.ts

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -113,6 +113,21 @@ export class ObjectQLPlugin implements Plugin {
113113
ctx.logger.debug('Discovered and registered app service (legacy)', { serviceName: name });
114114
}
115115
}
116+
117+
// Bridge realtime service from kernel service registry to ObjectQL.
118+
// RealtimeServicePlugin registers as 'realtime' service during init().
119+
// This enables ObjectQL to publish data change events.
120+
try {
121+
const realtimeService = ctx.getService('realtime');
122+
if (realtimeService) {
123+
ctx.logger.info('[ObjectQLPlugin] Bridging realtime service to ObjectQL for event publishing');
124+
this.ql.setRealtimeService(realtimeService);
125+
}
126+
} catch (e: any) {
127+
ctx.logger.debug('[ObjectQLPlugin] No realtime service found — data events will not be published', {
128+
error: e.message,
129+
});
130+
}
116131
}
117132

118133
// Initialize drivers (calls driver.connect() which sets up persistence)

0 commit comments

Comments
 (0)