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.
spCode làm phạm vi nghiệp vụ theo source LaoCRM hiện tại; không đưa thêm lớp Company vào task này.worker_log cũ để lưu tất cả event.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.
| 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 |
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.
| 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:
worker_command_state: unique (spCode, effectKey); index phục vụ lấy command đến hạn và lease hết hạn. _id đã unique theo commandId.budget_reservation: một reservation cho một command; index phục vụ đối soát theo campaign và trạng thái.event_outbox: unique eventId; index phục vụ lấy bản ghi đến hạn và phục hồi lease.| 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.
| 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.
| 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.
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.
Không giữ transaction Mongo trong lúc gọi HTTP tới LaoCredit hoặc provider.
WAITING_BUDGET; lưu lý do thiếu maxCost hay thiếu số dư SP.feeStatus=UNKNOWN, giữ reservation, kiểm tra theo orderId cũ; không tạo lệnh thu phí mới.actionStatus=UNKNOWN; đối soát hoặc gửi lại với cùng khóa idempotency nếu provider hỗ trợ. Không mặc định timeout là thất bại rồi gọi lại.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.
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.
cost + reservedCost + khoản cần giữ với maxCost bằng cập nhật nguyên tử.| 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.
worker_log trước đây. Dữ liệu thống kê đi Kafka → ClickHouse theo batch.sourceEventId, commandId, executionId, campaignId, spCode, bước xử lý và mã lỗi; dùng Elasticsearch hiện có để tra cứu. Không ghi token hoặc toàn bộ thông tin cá nhân vào log.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.