feat: auto committed
This commit is contained in:
186
services/content/src/shared/sync/neo4j-sync.worker.ts
Normal file
186
services/content/src/shared/sync/neo4j-sync.worker.ts
Normal file
@@ -0,0 +1,186 @@
|
||||
import { logger } from "../observability/logger.js";
|
||||
import { neo4jSyncConsumer } from "../../config/kafka.js";
|
||||
import { getNeo4jSession } from "../../config/neo4j.js";
|
||||
import { CONTENT_TOPICS, EVENT_TYPES } from "../outbox/events.js";
|
||||
|
||||
interface EventEnvelope {
|
||||
event_id: string;
|
||||
aggregate_id: string;
|
||||
event_type: string;
|
||||
occurred_at: number;
|
||||
action: string;
|
||||
kp_id?: string;
|
||||
chapter_id?: string;
|
||||
title?: string;
|
||||
difficulty?: number;
|
||||
prerequisite_id?: string;
|
||||
[key: string]: unknown;
|
||||
}
|
||||
|
||||
const processedEvents = new Set<string>();
|
||||
const DEDUP_MAX_SIZE = 10000;
|
||||
|
||||
export class Neo4jSyncWorker {
|
||||
private running = false;
|
||||
|
||||
async start(): Promise<void> {
|
||||
try {
|
||||
await neo4jSyncConsumer.connect();
|
||||
await neo4jSyncConsumer.subscribe({
|
||||
topic: CONTENT_TOPICS.KNOWLEDGE_POINT,
|
||||
fromBeginning: false,
|
||||
});
|
||||
|
||||
this.running = true;
|
||||
await neo4jSyncConsumer.run({
|
||||
eachMessage: async ({ message }) => {
|
||||
await this.handleMessage(message);
|
||||
},
|
||||
});
|
||||
logger.info("Neo4jSyncWorker started");
|
||||
} catch (err) {
|
||||
logger.warn(
|
||||
{
|
||||
err: err instanceof Error ? err.message : String(err),
|
||||
},
|
||||
"Neo4jSyncWorker failed to start (Neo4j/Kafka may be unavailable)",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async stop(): Promise<void> {
|
||||
this.running = false;
|
||||
try {
|
||||
await neo4jSyncConsumer.stop();
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
logger.info("Neo4jSyncWorker stopped");
|
||||
}
|
||||
|
||||
isRunning(): boolean {
|
||||
return this.running;
|
||||
}
|
||||
|
||||
private async handleMessage(
|
||||
message: { key: Buffer | null; value: Buffer | null },
|
||||
): Promise<void> {
|
||||
const value = message.value?.toString();
|
||||
if (!value) return;
|
||||
|
||||
let event: EventEnvelope;
|
||||
try {
|
||||
event = JSON.parse(value) as EventEnvelope;
|
||||
} catch {
|
||||
logger.error({ value }, "Failed to parse event payload");
|
||||
return;
|
||||
}
|
||||
|
||||
if (processedEvents.has(event.event_id)) {
|
||||
return;
|
||||
}
|
||||
processedEvents.add(event.event_id);
|
||||
if (processedEvents.size > DEDUP_MAX_SIZE) {
|
||||
const first = processedEvents.values().next().value;
|
||||
if (first) processedEvents.delete(first);
|
||||
}
|
||||
|
||||
const session = getNeo4jSession();
|
||||
if (!session) {
|
||||
logger.debug(
|
||||
{ eventType: event.event_type },
|
||||
"Neo4j not available, skipping sync",
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
await this.syncEvent(session, event);
|
||||
} catch (err) {
|
||||
logger.error(
|
||||
{
|
||||
err: err instanceof Error ? err.message : String(err),
|
||||
eventId: event.event_id,
|
||||
eventType: event.event_type,
|
||||
},
|
||||
"Neo4j sync failed",
|
||||
);
|
||||
} finally {
|
||||
await session.close();
|
||||
}
|
||||
}
|
||||
|
||||
private async syncEvent(
|
||||
session: NonNullable<ReturnType<typeof getNeo4jSession>>,
|
||||
event: EventEnvelope,
|
||||
): Promise<void> {
|
||||
switch (event.event_type) {
|
||||
case `edu.content.${EVENT_TYPES.KP_CREATED}`:
|
||||
case `edu.content.${EVENT_TYPES.KP_UPDATED}`:
|
||||
await this.upsertNode(session, event);
|
||||
break;
|
||||
case `edu.content.${EVENT_TYPES.KP_PREREQUISITE_ADDED}`:
|
||||
await this.addPrerequisite(session, event);
|
||||
break;
|
||||
case `edu.content.${EVENT_TYPES.KP_PREREQUISITE_REMOVED}`:
|
||||
await this.removePrerequisite(session, event);
|
||||
break;
|
||||
default:
|
||||
logger.debug(
|
||||
{ eventType: event.event_type },
|
||||
"Unhandled event type for Neo4j sync",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private async upsertNode(
|
||||
session: NonNullable<ReturnType<typeof getNeo4jSession>>,
|
||||
event: EventEnvelope,
|
||||
): Promise<void> {
|
||||
await session.executeWrite((tx) =>
|
||||
tx.run(
|
||||
`MERGE (kp:KnowledgePoint {id: $id})
|
||||
SET kp.title = $title, kp.difficulty = $difficulty, kp.updatedAt = datetime()`,
|
||||
{
|
||||
id: event.kp_id ?? event.aggregate_id,
|
||||
title: event.title ?? "",
|
||||
difficulty: event.difficulty ?? 3,
|
||||
},
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
private async addPrerequisite(
|
||||
session: NonNullable<ReturnType<typeof getNeo4jSession>>,
|
||||
event: EventEnvelope,
|
||||
): Promise<void> {
|
||||
await session.executeWrite((tx) =>
|
||||
tx.run(
|
||||
`MATCH (kp:KnowledgePoint {id: $kpId}), (prereq:KnowledgePoint {id: $prereqId})
|
||||
MERGE (prereq)-[:PREREQUISITE_OF]->(kp)`,
|
||||
{
|
||||
kpId: event.aggregate_id,
|
||||
prereqId: event.prerequisite_id ?? "",
|
||||
},
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
private async removePrerequisite(
|
||||
session: NonNullable<ReturnType<typeof getNeo4jSession>>,
|
||||
event: EventEnvelope,
|
||||
): Promise<void> {
|
||||
await session.executeWrite((tx) =>
|
||||
tx.run(
|
||||
`MATCH (prereq:KnowledgePoint {id: $prereqId})-[r:PREREQUISITE_OF]->(kp:KnowledgePoint {id: $kpId})
|
||||
DELETE r`,
|
||||
{
|
||||
kpId: event.aggregate_id,
|
||||
prereqId: event.prerequisite_id ?? "",
|
||||
},
|
||||
),
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
export const neo4jSyncWorker = new Neo4jSyncWorker();
|
||||
Reference in New Issue
Block a user