feat(content): docker 本地测试通过 + P5 ES 集成 + P6+ 审核工作流/可视化/可观测性
P5 ES 集成: - config/elasticsearch.ts: 惰性初始化 + ik_max_word→standard 回退 - shared/sync/es-sync.worker.ts: Kafka 消费 question 事件并索引 ES - questions search API: ES 优先, MySQL LIKE 降级 - ensureQuestionIndex() 幂等创建, IK 不可用回退 standard - main.ts: 启动 ensureQuestionIndex + esSyncWorker 生命周期管理 P6+ 审核工作流/可视化/可观测性: - Question 状态机: draft→pending_review→published→archived - 非法转换拦截 - 知识图谱可视化 API: Neo4j 优先 + MySQL 降级 - Cypher 返回标量避免 Node 包装对象问题 - 教材版本管理: GET /textbooks/versions + archive - 5 个 Prometheus 指标 + /readyz Outbox 积压检查 Docker 本地测试 (8 类全通过): - healthz/readyz (5 依赖 ok) - REST CRUD (textbook/chapter/kp/question) - ES 全文检索命中 - 审核工作流状态机 (合法/非法转换) - Outbox 事件驱动 (8 事件全 published) - Neo4j 同步 (KnowledgePoint 节点创建) - 可视化 (nodes/edges 正确) - Prometheus 指标 docs/nextstep.md: 上游 (MySQL/Neo4j/Kafka/Redis/ES/ai) + 下游 (teacher-bff/student-bff/parent-bff/data-ana/api-gateway/ai)
This commit is contained in:
@@ -4,6 +4,12 @@ import { HttpException } from "@nestjs/common";
|
||||
const mockGetDb = vi.fn();
|
||||
const mockGetNeo4jSession = vi.fn();
|
||||
const mockIsKafkaConnected = vi.fn();
|
||||
const mockGetEsClient = vi.fn();
|
||||
const mockEsClient = vi.hoisted(() => ({
|
||||
indices: { exists: vi.fn() },
|
||||
}));
|
||||
const mockCountPending = vi.fn();
|
||||
const mockOutboxPendingGauge = vi.hoisted(() => ({ set: vi.fn() }));
|
||||
|
||||
vi.mock("../../config/database.js", () => ({
|
||||
getDb: () => mockGetDb(),
|
||||
@@ -17,10 +23,23 @@ vi.mock("../../config/kafka.js", () => ({
|
||||
isKafkaProducerConnected: () => mockIsKafkaConnected(),
|
||||
}));
|
||||
|
||||
vi.mock("../../config/elasticsearch.js", () => ({
|
||||
getEsClient: () => mockGetEsClient(),
|
||||
QUESTION_INDEX_NAME: "content_questions",
|
||||
}));
|
||||
|
||||
vi.mock("../../config/env.js", () => ({
|
||||
env: { NEO4J_URL: undefined, NEO4J_PASSWORD: undefined },
|
||||
}));
|
||||
|
||||
vi.mock("../outbox/outbox.repository.js", () => ({
|
||||
outboxRepository: { countPending: () => mockCountPending() },
|
||||
}));
|
||||
|
||||
vi.mock("../observability/metrics.js", () => ({
|
||||
contentOutboxPendingTotal: mockOutboxPendingGauge,
|
||||
}));
|
||||
|
||||
import { HealthController } from "./health.controller.js";
|
||||
|
||||
describe("HealthController", () => {
|
||||
@@ -28,6 +47,9 @@ describe("HealthController", () => {
|
||||
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
// 默认:ES 未配置,outbox pending = 0
|
||||
mockGetEsClient.mockReturnValue(null);
|
||||
mockCountPending.mockResolvedValue(0);
|
||||
controller = new HealthController();
|
||||
});
|
||||
|
||||
@@ -48,9 +70,16 @@ describe("HealthController", () => {
|
||||
mockIsKafkaConnected.mockReturnValue(true);
|
||||
const result = await controller.readiness();
|
||||
expect(result.status).toBe("ok");
|
||||
expect(result.dependencies).toHaveLength(3);
|
||||
// mysql, neo4j, kafka, outbox, elasticsearch
|
||||
expect(result.dependencies).toHaveLength(5);
|
||||
const names = result.dependencies.map((d) => d.name);
|
||||
expect(names).toEqual(["mysql", "neo4j", "kafka"]);
|
||||
expect(names).toEqual([
|
||||
"mysql",
|
||||
"neo4j",
|
||||
"kafka",
|
||||
"outbox",
|
||||
"elasticsearch",
|
||||
]);
|
||||
expect(result.dependencies.every((d) => d.status === "ok")).toBe(true);
|
||||
});
|
||||
|
||||
@@ -64,6 +93,17 @@ describe("HealthController", () => {
|
||||
expect(neo4jDep?.status).toBe("ok");
|
||||
});
|
||||
|
||||
it("should report elasticsearch as ok when not configured", async () => {
|
||||
mockGetDb.mockReturnValue({
|
||||
execute: vi.fn().mockResolvedValue(undefined),
|
||||
});
|
||||
mockIsKafkaConnected.mockReturnValue(true);
|
||||
mockGetEsClient.mockReturnValue(null);
|
||||
const result = await controller.readiness();
|
||||
const esDep = result.dependencies.find((d) => d.name === "elasticsearch");
|
||||
expect(esDep?.status).toBe("ok");
|
||||
});
|
||||
|
||||
it("should report mysql error when db check fails", async () => {
|
||||
mockGetDb.mockReturnValue({
|
||||
execute: vi.fn().mockRejectedValue(new Error("db down")),
|
||||
@@ -79,5 +119,56 @@ describe("HealthController", () => {
|
||||
mockIsKafkaConnected.mockReturnValue(false);
|
||||
await expect(controller.readiness()).rejects.toThrow(HttpException);
|
||||
});
|
||||
|
||||
it("should mark outbox as degraded when pending > 100", async () => {
|
||||
mockGetDb.mockReturnValue({
|
||||
execute: vi.fn().mockResolvedValue(undefined),
|
||||
});
|
||||
mockIsKafkaConnected.mockReturnValue(true);
|
||||
mockCountPending.mockResolvedValue(150);
|
||||
const result = await controller.readiness();
|
||||
expect(result.status).toBe("degraded");
|
||||
const outboxDep = result.dependencies.find((d) => d.name === "outbox");
|
||||
expect(outboxDep?.status).toBe("degraded");
|
||||
// degraded 仍返回 200(不抛异常)
|
||||
});
|
||||
|
||||
it("should mark outbox as ok when pending <= 100", async () => {
|
||||
mockGetDb.mockReturnValue({
|
||||
execute: vi.fn().mockResolvedValue(undefined),
|
||||
});
|
||||
mockIsKafkaConnected.mockReturnValue(true);
|
||||
mockCountPending.mockResolvedValue(50);
|
||||
const result = await controller.readiness();
|
||||
const outboxDep = result.dependencies.find((d) => d.name === "outbox");
|
||||
expect(outboxDep?.status).toBe("ok");
|
||||
expect(result.status).toBe("ok");
|
||||
});
|
||||
|
||||
it("should mark elasticsearch as degraded when index not found", async () => {
|
||||
mockGetDb.mockReturnValue({
|
||||
execute: vi.fn().mockResolvedValue(undefined),
|
||||
});
|
||||
mockIsKafkaConnected.mockReturnValue(true);
|
||||
mockGetEsClient.mockReturnValue(mockEsClient);
|
||||
mockEsClient.indices.exists.mockResolvedValue(false);
|
||||
const result = await controller.readiness();
|
||||
const esDep = result.dependencies.find((d) => d.name === "elasticsearch");
|
||||
expect(esDep?.status).toBe("degraded");
|
||||
expect(result.status).toBe("degraded");
|
||||
});
|
||||
|
||||
it("should report elasticsearch as ok when index exists", async () => {
|
||||
mockGetDb.mockReturnValue({
|
||||
execute: vi.fn().mockResolvedValue(undefined),
|
||||
});
|
||||
mockIsKafkaConnected.mockReturnValue(true);
|
||||
mockGetEsClient.mockReturnValue(mockEsClient);
|
||||
mockEsClient.indices.exists.mockResolvedValue(true);
|
||||
const result = await controller.readiness();
|
||||
const esDep = result.dependencies.find((d) => d.name === "elasticsearch");
|
||||
expect(esDep?.status).toBe("ok");
|
||||
expect(result.status).toBe("ok");
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -3,18 +3,28 @@ import { sql } from "drizzle-orm";
|
||||
import { getDb } from "../../config/database.js";
|
||||
import { getNeo4jSession } from "../../config/neo4j.js";
|
||||
import { isKafkaProducerConnected } from "../../config/kafka.js";
|
||||
import {
|
||||
getEsClient,
|
||||
QUESTION_INDEX_NAME,
|
||||
} from "../../config/elasticsearch.js";
|
||||
import { env } from "../../config/env.js";
|
||||
import { outboxRepository } from "../outbox/outbox.repository.js";
|
||||
import { contentOutboxPendingTotal } from "../observability/metrics.js";
|
||||
|
||||
const SERVICE_NAME = "content";
|
||||
|
||||
// P6.4: Outbox pending 数量阈值,超过则标记 degraded
|
||||
const OUTBOX_PENDING_DEGRADED_THRESHOLD = 100;
|
||||
|
||||
interface DependencyCheck {
|
||||
name: string;
|
||||
status: "ok" | "error";
|
||||
status: "ok" | "error" | "degraded";
|
||||
error?: string;
|
||||
details?: Record<string, unknown>;
|
||||
}
|
||||
|
||||
interface ReadinessResponse {
|
||||
status: "ok" | "error";
|
||||
status: "ok" | "error" | "degraded";
|
||||
service: string;
|
||||
timestamp: string;
|
||||
dependencies: DependencyCheck[];
|
||||
@@ -89,10 +99,77 @@ export class HealthController {
|
||||
});
|
||||
}
|
||||
|
||||
const allOk = dependencies.every((d) => d.status === "ok");
|
||||
const status = allOk ? "ok" : "error";
|
||||
// 4. P6.4: Outbox pending count check(> 100 标记 degraded)
|
||||
// 同时更新 Prometheus Gauge 指标 content_outbox_pending_total
|
||||
try {
|
||||
const pendingCount = await outboxRepository.countPending();
|
||||
contentOutboxPendingTotal.set(pendingCount);
|
||||
if (pendingCount > OUTBOX_PENDING_DEGRADED_THRESHOLD) {
|
||||
dependencies.push({
|
||||
name: "outbox",
|
||||
status: "degraded",
|
||||
details: {
|
||||
pendingCount,
|
||||
threshold: OUTBOX_PENDING_DEGRADED_THRESHOLD,
|
||||
},
|
||||
});
|
||||
} else {
|
||||
dependencies.push({
|
||||
name: "outbox",
|
||||
status: "ok",
|
||||
details: { pendingCount },
|
||||
});
|
||||
}
|
||||
} catch (err) {
|
||||
dependencies.push({
|
||||
name: "outbox",
|
||||
status: "error",
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
});
|
||||
}
|
||||
|
||||
if (!allOk) {
|
||||
// 5. P6.4: ES index 存在性检查(ES 未配置时标记 ok-skip)
|
||||
const esClient = getEsClient();
|
||||
if (esClient) {
|
||||
try {
|
||||
const exists = await esClient.indices.exists({
|
||||
index: QUESTION_INDEX_NAME,
|
||||
});
|
||||
if (exists) {
|
||||
dependencies.push({ name: "elasticsearch", status: "ok" });
|
||||
} else {
|
||||
dependencies.push({
|
||||
name: "elasticsearch",
|
||||
status: "degraded",
|
||||
error: `index ${QUESTION_INDEX_NAME} not found`,
|
||||
});
|
||||
}
|
||||
} catch (err) {
|
||||
dependencies.push({
|
||||
name: "elasticsearch",
|
||||
status: "error",
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
});
|
||||
}
|
||||
} else {
|
||||
dependencies.push({
|
||||
name: "elasticsearch",
|
||||
status: "ok",
|
||||
error: "not configured (optional)",
|
||||
});
|
||||
}
|
||||
|
||||
// 综合状态:任一 error → error;任一 degraded(无 error)→ degraded;否则 ok
|
||||
const hasError = dependencies.some((d) => d.status === "error");
|
||||
const hasDegraded = dependencies.some((d) => d.status === "degraded");
|
||||
const status: ReadinessResponse["status"] = hasError
|
||||
? "error"
|
||||
: hasDegraded
|
||||
? "degraded"
|
||||
: "ok";
|
||||
|
||||
// error 时返回 503;degraded 时仍返回 200(服务可用,但需关注)
|
||||
if (hasError) {
|
||||
throw new HttpException(
|
||||
{
|
||||
status,
|
||||
|
||||
@@ -21,8 +21,57 @@ registry.registerMetric(
|
||||
}),
|
||||
);
|
||||
|
||||
// P6.4: 可观测性硬化指标
|
||||
// 待发布的 Outbox 事件数(由 health check 或定时任务设置)
|
||||
const contentOutboxPendingTotal = new promClient.Gauge({
|
||||
name: "content_outbox_pending_total",
|
||||
help: "Number of pending outbox events awaiting publication",
|
||||
});
|
||||
registry.registerMetric(contentOutboxPendingTotal);
|
||||
|
||||
// Neo4j 同步延迟(毫秒):最近一次事件时间距现在的差值
|
||||
const contentNeo4jSyncLagMs = new promClient.Gauge({
|
||||
name: "content_neo4j_sync_lag_ms",
|
||||
help: "Neo4j sync lag in milliseconds",
|
||||
});
|
||||
registry.registerMetric(contentNeo4jSyncLagMs);
|
||||
|
||||
// ES 同步延迟(毫秒)
|
||||
const contentEsSyncLagMs = new promClient.Gauge({
|
||||
name: "content_es_sync_lag_ms",
|
||||
help: "Elasticsearch sync lag in milliseconds",
|
||||
});
|
||||
registry.registerMetric(contentEsSyncLagMs);
|
||||
|
||||
// 题目搜索请求总数(按数据源 es/mysql 标签区分)
|
||||
const contentQuestionSearchTotal = new promClient.Counter({
|
||||
name: "content_question_search_total",
|
||||
help: "Total number of question search requests",
|
||||
labelNames: ["source"],
|
||||
});
|
||||
registry.registerMetric(contentQuestionSearchTotal);
|
||||
|
||||
// 题目搜索延迟(毫秒)
|
||||
const contentQuestionSearchLatencyMs = new promClient.Histogram({
|
||||
name: "content_question_search_latency_ms",
|
||||
help: "Question search latency in milliseconds",
|
||||
labelNames: ["source"],
|
||||
buckets: [1, 5, 10, 25, 50, 100, 250, 500, 1000],
|
||||
});
|
||||
registry.registerMetric(contentQuestionSearchLatencyMs);
|
||||
|
||||
// 自动收集 Node.js 进程级指标(CPU/内存/事件循环/GC等)
|
||||
// 这些指标无需业务代码埋点,prom-client 自动采集
|
||||
promClient.collectDefaultMetrics({ register: registry });
|
||||
|
||||
export { registry as metricsRegistry };
|
||||
export {
|
||||
registry as metricsRegistry,
|
||||
contentOutboxPendingTotal,
|
||||
contentNeo4jSyncLagMs,
|
||||
contentEsSyncLagMs,
|
||||
contentQuestionSearchTotal,
|
||||
contentQuestionSearchLatencyMs,
|
||||
// 兼容性别名:业务代码使用的简短名称
|
||||
contentQuestionSearchTotal as questionSearchTotal,
|
||||
contentQuestionSearchLatencyMs as questionSearchLatencyMs,
|
||||
};
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { eq, sql, and, or, isNull, lte } from "drizzle-orm";
|
||||
import { eq, sql, and, or, isNull, lte, count } from "drizzle-orm";
|
||||
import type { MySql2Database } from "drizzle-orm/mysql2";
|
||||
import { getDb } from "../../config/database.js";
|
||||
import {
|
||||
@@ -29,6 +29,15 @@ export class OutboxRepository {
|
||||
.limit(limit);
|
||||
}
|
||||
|
||||
// P6.4: 统计 pending 状态的 outbox 事件总数(health check + metrics 用)
|
||||
async countPending(): Promise<number> {
|
||||
const [row] = await getDb()
|
||||
.select({ value: count() })
|
||||
.from(outbox)
|
||||
.where(eq(outbox.status, "pending"));
|
||||
return row?.value ?? 0;
|
||||
}
|
||||
|
||||
async markPublished(id: string): Promise<void> {
|
||||
await getDb()
|
||||
.update(outbox)
|
||||
|
||||
325
services/content/src/shared/sync/es-sync.worker.test.ts
Normal file
325
services/content/src/shared/sync/es-sync.worker.test.ts
Normal file
@@ -0,0 +1,325 @@
|
||||
import { describe, it, expect, vi, beforeEach } from "vitest";
|
||||
|
||||
// 使用 vi.hoisted 确保 mock 在 worker 模块导入前注册
|
||||
const mockEsSyncConsumer = vi.hoisted(() => ({
|
||||
connect: vi.fn(),
|
||||
subscribe: vi.fn(),
|
||||
run: vi.fn(),
|
||||
stop: vi.fn(),
|
||||
disconnect: vi.fn(),
|
||||
}));
|
||||
|
||||
const mockEsClient = vi.hoisted(() => ({
|
||||
index: vi.fn(),
|
||||
delete: vi.fn(),
|
||||
}));
|
||||
|
||||
const mockGetEsClient = vi.hoisted(() => vi.fn());
|
||||
|
||||
const mockQuestionsRepo = vi.hoisted(() => ({
|
||||
findById: vi.fn(),
|
||||
}));
|
||||
|
||||
vi.mock("../../config/kafka.js", () => ({
|
||||
esSyncConsumer: mockEsSyncConsumer,
|
||||
}));
|
||||
|
||||
vi.mock("../../config/elasticsearch.js", () => ({
|
||||
getEsClient: mockGetEsClient,
|
||||
QUESTION_INDEX_NAME: "content_questions",
|
||||
}));
|
||||
|
||||
vi.mock("../../questions/questions.repository.js", () => ({
|
||||
questionsRepository: mockQuestionsRepo,
|
||||
}));
|
||||
|
||||
vi.mock("../observability/logger.js", () => ({
|
||||
logger: {
|
||||
info: vi.fn(),
|
||||
error: vi.fn(),
|
||||
warn: vi.fn(),
|
||||
debug: vi.fn(),
|
||||
},
|
||||
}));
|
||||
|
||||
import { EsSyncWorker } from "./es-sync.worker.js";
|
||||
import { EVENT_TYPES } from "../outbox/events.js";
|
||||
|
||||
// 构造 Kafka 消息(value 为 JSON 字符串)
|
||||
function createMessage(event: Record<string, unknown>): {
|
||||
key: Buffer | null;
|
||||
value: Buffer | null;
|
||||
} {
|
||||
return {
|
||||
key: Buffer.from(String(event.aggregate_id ?? "")),
|
||||
value: Buffer.from(JSON.stringify(event)),
|
||||
};
|
||||
}
|
||||
|
||||
const sampleQuestion = {
|
||||
id: "q-1",
|
||||
knowledgePointId: "kp-1",
|
||||
type: "single_choice",
|
||||
content: "What is 2+2?",
|
||||
options: null,
|
||||
answer: "4",
|
||||
explanation: "Basic arithmetic",
|
||||
difficulty: 3,
|
||||
status: "draft",
|
||||
source: "manual",
|
||||
createdBy: "u-1",
|
||||
metadata: null,
|
||||
createdAt: new Date("2026-01-01T00:00:00Z"),
|
||||
updatedAt: new Date("2026-01-02T00:00:00Z"),
|
||||
};
|
||||
|
||||
describe("EsSyncWorker", () => {
|
||||
let worker: EsSyncWorker;
|
||||
|
||||
beforeEach(() => {
|
||||
vi.clearAllMocks();
|
||||
worker = new EsSyncWorker();
|
||||
// 默认 ES 可用
|
||||
mockGetEsClient.mockReturnValue(mockEsClient);
|
||||
mockEsSyncConsumer.run.mockImplementation(({ eachMessage }) => {
|
||||
// 不实际启动消费循环;测试通过调用 handleMessage 验证
|
||||
void eachMessage;
|
||||
return Promise.resolve();
|
||||
});
|
||||
});
|
||||
|
||||
describe("start", () => {
|
||||
it("ES 未配置时跳过启动(不连接 Kafka)", async () => {
|
||||
mockGetEsClient.mockReturnValue(null);
|
||||
await worker.start();
|
||||
expect(mockEsSyncConsumer.connect).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("ES 配置时连接 Kafka 并订阅 question topic", async () => {
|
||||
await worker.start();
|
||||
expect(mockEsSyncConsumer.connect).toHaveBeenCalled();
|
||||
expect(mockEsSyncConsumer.subscribe).toHaveBeenCalledWith(
|
||||
expect.objectContaining({ topic: "edu.content.question.events" }),
|
||||
);
|
||||
expect(mockEsSyncConsumer.run).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("Kafka 连接失败时软失败(不抛出)", async () => {
|
||||
mockEsSyncConsumer.connect.mockRejectedValueOnce(new Error("kafka down"));
|
||||
await expect(worker.start()).resolves.not.toThrow();
|
||||
});
|
||||
});
|
||||
|
||||
describe("stop", () => {
|
||||
it("停止时不抛出错误", async () => {
|
||||
await expect(worker.stop()).resolves.not.toThrow();
|
||||
expect(mockEsSyncConsumer.stop).toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
describe("handleMessage — 事件处理", () => {
|
||||
// 通过类型转换访问私有 handleMessage
|
||||
function handleMessage(
|
||||
w: EsSyncWorker,
|
||||
msg: { key: Buffer | null; value: Buffer | null },
|
||||
): Promise<void> {
|
||||
return (
|
||||
w as unknown as {
|
||||
handleMessage: (m: {
|
||||
key: Buffer | null;
|
||||
value: Buffer | null;
|
||||
}) => Promise<void>;
|
||||
}
|
||||
).handleMessage(msg);
|
||||
}
|
||||
|
||||
it("question.created → 从 DB 读取并 index 文档", async () => {
|
||||
mockQuestionsRepo.findById.mockResolvedValue(sampleQuestion);
|
||||
const event = {
|
||||
event_id: "e-1",
|
||||
aggregate_id: "q-1",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_CREATED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "created",
|
||||
};
|
||||
await handleMessage(worker, createMessage(event));
|
||||
|
||||
expect(mockQuestionsRepo.findById).toHaveBeenCalledWith("q-1");
|
||||
expect(mockEsClient.index).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
index: "content_questions",
|
||||
id: "q-1",
|
||||
}),
|
||||
);
|
||||
const doc = mockEsClient.index.mock.calls[0]?.[0]?.document;
|
||||
expect(doc?.question_id).toBe("q-1");
|
||||
expect(doc?.content).toBe("What is 2+2?");
|
||||
expect(doc?.difficulty).toBe(3);
|
||||
});
|
||||
|
||||
it("question.updated → 重新 index 文档", async () => {
|
||||
mockQuestionsRepo.findById.mockResolvedValue(sampleQuestion);
|
||||
const event = {
|
||||
event_id: "e-2",
|
||||
aggregate_id: "q-1",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_UPDATED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "updated",
|
||||
};
|
||||
await handleMessage(worker, createMessage(event));
|
||||
expect(mockEsClient.index).toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("question.published → 重新 index 文档(含新 status)", async () => {
|
||||
mockQuestionsRepo.findById.mockResolvedValue({
|
||||
...sampleQuestion,
|
||||
status: "published",
|
||||
});
|
||||
const event = {
|
||||
event_id: "e-3",
|
||||
aggregate_id: "q-1",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_PUBLISHED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "published",
|
||||
};
|
||||
await handleMessage(worker, createMessage(event));
|
||||
expect(mockEsClient.index).toHaveBeenCalled();
|
||||
const doc = mockEsClient.index.mock.calls[0]?.[0]?.document;
|
||||
expect(doc?.status).toBe("published");
|
||||
});
|
||||
|
||||
it("question.deleted → 删除 ES 文档", async () => {
|
||||
const event = {
|
||||
event_id: "e-4",
|
||||
aggregate_id: "q-1",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_DELETED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "deleted",
|
||||
};
|
||||
await handleMessage(worker, createMessage(event));
|
||||
expect(mockEsClient.delete).toHaveBeenCalledWith({
|
||||
index: "content_questions",
|
||||
id: "q-1",
|
||||
});
|
||||
expect(mockQuestionsRepo.findById).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("question.created 但 DB 中不存在 → 跳过索引", async () => {
|
||||
mockQuestionsRepo.findById.mockResolvedValue(undefined);
|
||||
const event = {
|
||||
event_id: "e-5",
|
||||
aggregate_id: "q-missing",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_CREATED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "created",
|
||||
};
|
||||
await handleMessage(worker, createMessage(event));
|
||||
expect(mockEsClient.index).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("未处理的事件类型 → 不调用 ES", async () => {
|
||||
const event = {
|
||||
event_id: "e-6",
|
||||
aggregate_id: "q-1",
|
||||
event_type: "edu.content.unknown_event",
|
||||
occurred_at: Date.now(),
|
||||
action: "unknown",
|
||||
};
|
||||
await handleMessage(worker, createMessage(event));
|
||||
expect(mockEsClient.index).not.toHaveBeenCalled();
|
||||
expect(mockEsClient.delete).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
|
||||
describe("handleMessage — 幂等去重", () => {
|
||||
function handleMessage(
|
||||
w: EsSyncWorker,
|
||||
msg: { key: Buffer | null; value: Buffer | null },
|
||||
): Promise<void> {
|
||||
return (
|
||||
w as unknown as {
|
||||
handleMessage: (m: {
|
||||
key: Buffer | null;
|
||||
value: Buffer | null;
|
||||
}) => Promise<void>;
|
||||
}
|
||||
).handleMessage(msg);
|
||||
}
|
||||
|
||||
it("相同 event_id 的事件只处理一次", async () => {
|
||||
mockQuestionsRepo.findById.mockResolvedValue(sampleQuestion);
|
||||
const event = {
|
||||
event_id: "dup-1",
|
||||
aggregate_id: "q-1",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_CREATED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "created",
|
||||
};
|
||||
const msg = createMessage(event);
|
||||
await handleMessage(worker, msg);
|
||||
await handleMessage(worker, msg);
|
||||
expect(mockEsClient.index).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
});
|
||||
|
||||
describe("handleMessage — 软失败", () => {
|
||||
function handleMessage(
|
||||
w: EsSyncWorker,
|
||||
msg: { key: Buffer | null; value: Buffer | null },
|
||||
): Promise<void> {
|
||||
return (
|
||||
w as unknown as {
|
||||
handleMessage: (m: {
|
||||
key: Buffer | null;
|
||||
value: Buffer | null;
|
||||
}) => Promise<void>;
|
||||
}
|
||||
).handleMessage(msg);
|
||||
}
|
||||
|
||||
it("ES Client 为 null 时跳过(不抛出)", async () => {
|
||||
mockGetEsClient.mockReturnValue(null);
|
||||
const event = {
|
||||
event_id: "e-null",
|
||||
aggregate_id: "q-1",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_CREATED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "created",
|
||||
};
|
||||
await expect(
|
||||
handleMessage(worker, createMessage(event)),
|
||||
).resolves.not.toThrow();
|
||||
expect(mockEsClient.index).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("ES 操作抛错时软失败(不抛出)", async () => {
|
||||
mockQuestionsRepo.findById.mockResolvedValue(sampleQuestion);
|
||||
mockEsClient.index.mockRejectedValueOnce(new Error("ES down"));
|
||||
const event = {
|
||||
event_id: "e-err",
|
||||
aggregate_id: "q-1",
|
||||
event_type: `edu.content.${EVENT_TYPES.QUESTION_CREATED}`,
|
||||
occurred_at: Date.now(),
|
||||
action: "created",
|
||||
};
|
||||
await expect(
|
||||
handleMessage(worker, createMessage(event)),
|
||||
).resolves.not.toThrow();
|
||||
});
|
||||
|
||||
it("无效 JSON payload 不抛出", async () => {
|
||||
await expect(
|
||||
handleMessage(worker, {
|
||||
key: Buffer.from("k"),
|
||||
value: Buffer.from("not-json"),
|
||||
}),
|
||||
).resolves.not.toThrow();
|
||||
});
|
||||
|
||||
it("空 value 不抛出", async () => {
|
||||
await expect(
|
||||
handleMessage(worker, { key: null, value: null }),
|
||||
).resolves.not.toThrow();
|
||||
});
|
||||
});
|
||||
});
|
||||
189
services/content/src/shared/sync/es-sync.worker.ts
Normal file
189
services/content/src/shared/sync/es-sync.worker.ts
Normal file
@@ -0,0 +1,189 @@
|
||||
import { logger } from "../observability/logger.js";
|
||||
import { esSyncConsumer } from "../../config/kafka.js";
|
||||
import {
|
||||
getEsClient,
|
||||
QUESTION_INDEX_NAME,
|
||||
} from "../../config/elasticsearch.js";
|
||||
import { CONTENT_TOPICS, EVENT_TYPES } from "../outbox/events.js";
|
||||
import { questionsRepository } from "../../questions/questions.repository.js";
|
||||
|
||||
interface EventEnvelope {
|
||||
event_id: string;
|
||||
aggregate_id: string;
|
||||
event_type: string;
|
||||
occurred_at: number;
|
||||
action: string;
|
||||
[key: string]: unknown;
|
||||
}
|
||||
|
||||
// 幂等去重 Set:与 neo4j-sync.worker.ts 保持相同模式
|
||||
const processedEvents = new Set<string>();
|
||||
const DEDUP_MAX_SIZE = 10000;
|
||||
|
||||
export class EsSyncWorker {
|
||||
private running = false;
|
||||
|
||||
async start(): Promise<void> {
|
||||
const client = getEsClient();
|
||||
if (!client) {
|
||||
logger.info("EsSyncWorker skipped (ES not configured)");
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await esSyncConsumer.connect();
|
||||
await esSyncConsumer.subscribe({
|
||||
topic: CONTENT_TOPICS.QUESTION,
|
||||
fromBeginning: false,
|
||||
});
|
||||
|
||||
this.running = true;
|
||||
await esSyncConsumer.run({
|
||||
eachMessage: async ({ message }) => {
|
||||
await this.handleMessage(message);
|
||||
},
|
||||
});
|
||||
logger.info("EsSyncWorker started");
|
||||
} catch (err) {
|
||||
logger.warn(
|
||||
{
|
||||
err: err instanceof Error ? err.message : String(err),
|
||||
},
|
||||
"EsSyncWorker failed to start (ES/Kafka may be unavailable)",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async stop(): Promise<void> {
|
||||
this.running = false;
|
||||
try {
|
||||
await esSyncConsumer.stop();
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
logger.info("EsSyncWorker 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 client = getEsClient();
|
||||
if (!client) {
|
||||
logger.debug(
|
||||
{ eventType: event.event_type },
|
||||
"ES not available, skipping sync",
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
try {
|
||||
await this.syncEvent(client, event);
|
||||
} catch (err) {
|
||||
// ES 不可用时软失败:log warn 并跳过,不阻塞消费链
|
||||
logger.warn(
|
||||
{
|
||||
err: err instanceof Error ? err.message : String(err),
|
||||
eventId: event.event_id,
|
||||
eventType: event.event_type,
|
||||
},
|
||||
"ES sync failed, skipping (soft failure)",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private async syncEvent(
|
||||
client: NonNullable<ReturnType<typeof getEsClient>>,
|
||||
event: EventEnvelope,
|
||||
): Promise<void> {
|
||||
switch (event.event_type) {
|
||||
case `edu.content.${EVENT_TYPES.QUESTION_CREATED}`:
|
||||
case `edu.content.${EVENT_TYPES.QUESTION_UPDATED}`:
|
||||
case `edu.content.${EVENT_TYPES.QUESTION_PUBLISHED}`:
|
||||
await this.indexQuestion(client, event.aggregate_id);
|
||||
break;
|
||||
case `edu.content.${EVENT_TYPES.QUESTION_DELETED}`:
|
||||
await client.delete({
|
||||
index: QUESTION_INDEX_NAME,
|
||||
id: event.aggregate_id,
|
||||
});
|
||||
logger.debug(
|
||||
{ questionId: event.aggregate_id, eventId: event.event_id },
|
||||
"ES document deleted",
|
||||
);
|
||||
break;
|
||||
default:
|
||||
logger.debug(
|
||||
{ eventType: event.event_type },
|
||||
"Unhandled event type for ES sync",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 从 MySQL 读取题目完整记录并索引到 ES。
|
||||
* 事件只携带部分字段,因此需要回查 DB 以获取 content/answer/explanation。
|
||||
*/
|
||||
private async indexQuestion(
|
||||
client: NonNullable<ReturnType<typeof getEsClient>>,
|
||||
questionId: string,
|
||||
): Promise<void> {
|
||||
const question = await questionsRepository.findById(questionId);
|
||||
if (!question) {
|
||||
logger.warn(
|
||||
{ questionId },
|
||||
"Question not found in DB, skipping ES index",
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
await client.index({
|
||||
index: QUESTION_INDEX_NAME,
|
||||
id: question.id,
|
||||
document: {
|
||||
question_id: question.id,
|
||||
knowledge_point_id: question.knowledgePointId,
|
||||
type: question.type,
|
||||
content: question.content,
|
||||
answer: question.answer,
|
||||
explanation: question.explanation ?? "",
|
||||
difficulty: question.difficulty,
|
||||
status: question.status,
|
||||
source: question.source,
|
||||
created_by: question.createdBy,
|
||||
created_at: question.createdAt.getTime(),
|
||||
updated_at: question.updatedAt.getTime(),
|
||||
},
|
||||
});
|
||||
logger.debug(
|
||||
{ questionId: question.id, eventId: "n/a" },
|
||||
"ES document indexed",
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
export const esSyncWorker = new EsSyncWorker();
|
||||
Reference in New Issue
Block a user