Task: Nền xử lý event/action chung theo kiến trúc v2

1. Mục tiêu và phạm vi

Hiện tại mới triển khai LaoCRM và các phần liên quan. Các bảng worker_command_state, event_outbox, budget_reservation, campaign_action_event, campaign_fee_event mới nằm trong thiết kế, chưa được triển khai trong source.

Task này xây dựng nền chung để nhận event, tạo lệnh thực hiện action, kiểm soát ngân sách, xử lý lỗi/retry và đưa kết quả sang báo cáo. Sau đó các action mới như UPoint, gia hạn LaoTV sẽ tích hợp vào nền này.

Không phải mọi event đều tạo command. Event nghiệp vụ có action và người nhận từ một lần chạy campaign mới đi qua nền thực thi này. Tracking click/open, PostHog và dữ liệu chỉ phục vụ thống kê đi luồng Kafka → analytics riêng.

2. Các thành phần trong luồng

Thành phần Hiểu đơn giản Hiện trạng
Event Dữ liệu đầu vào kích hoạt campaign/action Có; cần bổ sung định danh ổn định
Execution Tiến độ một chuỗi action của một người nhận Có campaign_action_execution; cần sửa
Command Lệnh thực hiện một action cụ thể Có message; chưa có bảng trạng thái riêng
Outbox Hàng đợi bền vững trong Mongo để gửi Kafka Chưa có
Reservation Khoản ngân sách đang giữ cho một command Chưa có
Action event Kết quả nghiệp vụ để thống kê Cần bổ sung luồng và bảng mới
Fee event Kết quả thu phí để thống kê Cần bổ sung luồng và bảng mới

3. Sơ đồ luồng chung

[Diagram]

Các mũi tên thể hiện dữ liệu đi đâu, không có nghĩa Mongo, Kafka và ClickHouse cùng tham gia một transaction. Outbox giúp phục hồi việc gửi; consumer phải xử lý được message trùng.

4. Các bảng cần thêm và sửa

4.1. Bảng MongoDB mới

Bảng mới Dữ liệu chính Nơi ghi / mục đích
worker_command_state _id=commandId, spCode, executionId, campaignId, effectKey, action và tham số đã chốt, giá đã chốt, actionStatus, feeStatus, attempt, lease, providerReference, lỗi gần nhất, revision Worker ghi; biết command đang làm gì, chống xử lý trùng, phục hồi khi worker chết
budget_reservation _id=commandId, campaignId, spCode, số tiền, tiền tệ, trạng thái HELD/CONSUMED/RELEASED, transferRecordId, version và thời gian Worker ghi; kiểm soát tổng ngân sách khi nhiều command chạy đồng thời
event_outbox _id/eventId, đối tượng liên quan, revision, loại event, topic, partition key, payload hoặc tham chiếu, trạng thái, lease, lần gửi, retryAt, publishedAt Dispatcher/worker ghi cùng thay đổi Mongo; relay gửi Kafka

Yêu cầu index chính:

4.2. Bảng MongoDB hiện có cần sửa

Bảng Thay đổi cần làm
campaign_action_execution Bổ sung định danh execution ổn định trong phạm vi SP; unique theo định danh đó. Giữ chuỗi action/tiến độ hiện có, cập nhật có điều kiện theo command và trạng thái để kết quả trùng không chạy action tiếp hai lần
campaign Bổ sung reservedCost; cập nhật ngân sách nguyên tử. Rà soát cost/maxCost/expectedCost hiện dùng Double, chuyển xử lý tiền sang kiểu chính xác với kế hoạch chuyển dữ liệu phù hợp
laocredit_transfer_record Liên kết commandId, reservation và mục đích thu phí; định danh orderId ổn định, unique theo hợp đồng LaoCredit; hỗ trợ trạng thái chưa xác định và dấu đã cộng chi phí để không cộng lại

Không cần thêm lại User, SP, role, permission hoặc Company trong task này. Các thay đổi đó thuộc phần LaoCRM đã triển khai.

4.3. Hai bảng ClickHouse mới cho báo cáo

Bảng mới Dữ liệu chính Ý nghĩa
campaign_action_event eventId, commandId, spCode, campaign/run/snapshot, action, người nhận, trạng thái, sourceRevision, thời điểm, tham chiếu provider Thống kê action thành công/thất bại/chưa rõ kết quả
campaign_fee_event eventId, transferId, commandId, spCode, campaign, số tiền/tiền tệ, trạng thái, sourceRevision, thời điểm thu phí Thống kê phí thực tế đã được xác nhận

ClickHouse là nơi đọc báo cáo; Mongo giữ trạng thái thực thi và sổ thu phí gốc. Không suy ra phí thực tế bằng số tin nhắn nhân đơn giá hiện tại.

commandId/transferId là khóa nghiệp vụ để nhận diện dữ liệu, không phải unique constraint của ClickHouse. Thiết kế ingestion/query phải lấy revision hợp lệ và chịu được replay; không cộng trực tiếp mọi dòng nhận được vì sẽ đếm trùng.

Các bảng tổng hợp và report_partition_manifest / analytics_projection_state trong thiết kế báo cáo sẽ triển khai ở task báo cáo tiếp theo. Task này phải cung cấp đủ định danh, revision và thời gian để task đó tổng hợp chính xác.

5. Các sự kiện cần dùng và bổ sung

Message / event Đường đi Việc cần làm
EventMessage hiện có Nguồn event → campaign-dispatcher Bổ sung sourceEventId ổn định từ nguồn; truyền spCode và các định danh campaign/run hiện có
CampaignActionCommandMessage hiện có Dispatcher → topic command theo worker → worker Chốt commandId, effectKey, execution và action identity; gửi lại vẫn giữ định danh cũ
CampaignActionResultMessage hiện có Worker → campaign-action-result → dispatcher Chỉ kết quả kết thúc mới chuyển tiến độ chuỗi; kết quả trùng không chuyển lần hai
CampaignActionOutcomeEvent mới, tên đề xuất Worker → campaign-action-event → analytics Gửi trạng thái nghiệp vụ cùng eventId, command, SP, campaign, revision, thời gian
CampaignFeeEvent mới, tên đề xuất Worker/đối soát → campaign-fee-event → analytics Gửi trạng thái thu phí cùng transfer, command, số tiền, tiền tệ, revision

Mỗi thay đổi trạng thái có eventId riêng. Relay gửi lại một thay đổi thì giữ nguyên eventId; thay đổi tiếp theo tăng revision.

campaign-action-event và campaign-fee-event là tên logic đề xuất cho hai topic mới. Bổ sung key cấu hình/consumer group theo quy ước topic hiện có của repo; không hard-code tên topic. Các topic hiện có tiếp tục dùng cấu hình đã đồng bộ.

Kết quả gửi tin và action nghiệp vụ phải có nghĩa rõ ràng: provider nhận yêu cầu chưa chắc là tin đã được giao. Không thống kê cùng một kết quả subscriber action từ cả delivery stream cũ và action stream mới.

6. Đầu việc để giao triển khai

Tên service mới dưới đây là đề xuất để giao việc, chưa có trong source.

Task Module / service Nội dung và dữ liệu đi qua Phụ thuộc
WF01 — Chốt định danh và hợp đồng event common, campaign-dispatcher, nguồn phát event Bổ sung source identity; execution/command/effect identity ổn định; quy định revision và các trạng thái. Không tạo request identity mới mỗi lần consume cùng event Làm đầu tiên
WF02 — Thêm bảng trạng thái và index common entity/repository và migrations Thêm ba bảng Mongo; sửa execution, campaign và transfer record như mục 4; chốt cách lưu tiền chính xác WF01
WF03 — Dispatcher ghi command qua outbox CampaignEventConsumer, CampaignDispatcherService, CampaignActionCommandPublisher Nhận event → tạo/lấy execution → chốt command → ghi execution + outbox trong transaction Mongo. Thay đường publish command trực tiếp WF02
WF04 — Gửi outbox và phục hồi Service relay trong campaign-dispatcher, dùng chung model/repository trong common Claim batch → gửi đúng topic/key → chờ broker ACK → đánh dấu published. Phục hồi lease và retry với cùng eventId. Dùng một cơ chế relay, không dựng thêm microservice WF02–03
WF05 — Nền nhận command chung cho worker common: service xử lý trạng thái dùng chung; consumer các worker Nhận command → tạo/lấy state → claim có điều kiện → chốt tham số → gọi bước phí/handler → lưu kết quả. Giữ handler gửi tin của v2 WF01–02
WF06 — Ngân sách và thu phí v1 common: service ngân sách/thu phí dùng chung; LaoCreditApiClient; các worker gọi Campaign → reservation → LaoCredit → transfer record → cập nhật cost đúng một lần. Thiếu ngân sách thì chờ/pause; kết quả phí chưa rõ thì đối soát WF02, WF05
WF07 — Kết quả, retry và tiến độ chuỗi Worker + CampaignDispatcherService State + outbox → result/facts. Dispatcher chỉ chuyển action khi có kết quả kết thúc hợp lệ. Bổ sung phục hồi command, resume và xử lý dừng campaign WF03–06
WF08 — Dữ liệu đầu vào báo cáo ads-analytics-flink-job, schema ClickHouse Consume hai topic mới → ghi batch hai bảng mới; xử lý revision/replay và thống nhất số liệu với delivery stream. Công bố thời điểm dữ liệu đã cập nhật WF01, WF07
WF09 — Áp dụng nền chung cho từng worker subscriber-action-worker, sms-worker, mail-worker, whatsapp-worker, push-worker Gắn consumer/handler hiện có vào nền chung; cấu hình giới hạn đồng thời, retry và timeout theo provider. Làm lần lượt để tránh thay tất cả cùng lúc WF05–07

WF01–02 cần chốt trước. Sau đó dispatcher/outbox và worker/budget có thể chia cho hai người làm theo cùng hợp đồng. WF08 có thể triển khai song song khi hợp đồng fact event đã ổn định.

7. Luồng xử lý các tình huống chính

7.1. Chạy bình thường

  1. Dispatcher nhận event, xác định execution và command bằng định danh ổn định.
  2. Ghi execution và command outbox trong cùng transaction Mongo.
  3. Relay gửi command; worker claim trạng thái để một command chỉ có một tiến trình đang sở hữu lease hợp lệ.
  4. Nếu action thuộc diện thu phí: chốt giá → giữ ngân sách → gọi LaoCredit bằng orderId cố định → ghi phí thành công → cộng cost và tiêu thụ reservation đúng một lần.
  5. Worker gọi handler/provider; action không thu phí bỏ qua bước 4.
  6. Ghi kết quả command cùng outbox trong transaction Mongo. Relay gửi result và các fact event.
  7. Dispatcher chuyển sang action kế tiếp; analytics cập nhật dữ liệu báo cáo độc lập.

Không giữ transaction Mongo trong lúc gọi HTTP tới LaoCredit hoặc provider.

7.2. Thiếu ngân sách hoặc số dư

7.3. Timeout, worker chết hoặc nhận lại message

Lease/Redis lock không đảm bảo provider chỉ thực hiện một lần: HTTP cũ có thể đã đến provider. Cần xác minh hợp đồng idempotency hoặc API tra cứu kết quả của từng provider; trường hợp không có phải có luồng xử lý thủ công.

7.4. Pause/stop và chạy lại chủ động

8. Chính sách phí phải mang từ v1

Danh sách action có thu phí trong luồng transfer v1 hiện tại: TOP_UP, ADD_MONEY, ADD_DATA, ADD_PROMOTION, ADD_SMS, ADD_CALL_MINUTES, NOTIFICATION_SMS, NOTIFICATION_MAIL, ADD_UPOINT.

9. Task action bổ sung sau nền chung

Task riêng Module Phần hiện tại còn thiếu
ADD_UPOINT subscriber-action-worker, phần giá trong API/common DTO/action type đã có; provider hiện chưa triển khai. Bổ sung handler/client, route qua executor hiện có; bổ sung giá UPoint vào quote và phép tính
Gia hạn LaoTV subscriber-action-worker Luồng gọi đã có một phần; bổ sung/thống nhất tham số v1 như phần trăm gia hạn theo hợp đồng nghiệp vụ thực tế

Hai task này dùng nền claim command, ngân sách, outbox và báo cáo chung. Không tạo thêm worker type/microservice chỉ cho từng action.

10. Hiệu suất, báo cáo và kiểm tra lỗi

11. Tiêu chí hoàn thành để review

Trước khi viết code cần chốt thêm khả năng idempotency/tra cứu của từng provider, hợp đồng thu phí LaoCredit, sự kiện kích hoạt resume và thời gian giữ dữ liệu. Đây là các chi tiết triển khai còn thiếu, không phải các bảng đã tồn tại.