L3 - Nền tảng giám sát máy chủ - Collector - Bus, Worker, TSDB
Quy ước tên: L3 - <P&L> - <Hệ thống L2> - <Thành phần>. P&L: chưa chỉ định. Hệ thống L2: Collector. Thành phần: Bus (hàng đợi theo shard), Worker (gom lô và ghi), TSDB adapter (VictoriaMetrics và bộ nhớ). Mã: internal/bus, internal/worker, internal/tsdb, phần nối dây trong internal/app/wiring.go.
Quy ước đánh dấu: "Đã làm" là có mã và kiểm thử. "Chỉ thiết kế" là chỉ có trong docs/. "Đề xuất" là ý của tác giả bản nháp. Mã thắng khi lệch tài liệu. Khóa Redis mang tiền tố redis.prefix (mặc định ah:).
0. Front matter
Thông tin tài liệu đầy đủ
| Trạng thái | BẢN NHÁP, phiên bản 0.1, 2026-09-30. Người phê duyệt: chưa chỉ định |
| Thành phần (Component) | Bus (inproc, Redis Streams), Worker (shard worker, batcher), TSDB adapter (VictoriaMetrics, bộ nhớ), bộ lập kế hoạch truy vấn |
| Truy vết L2 | L2 SAD Collector: thành phần "Bus", "Worker", "TSDB adapter" (mục 2.3, 6.1), FR ghi và truy vấn chuỗi thời gian (mục 3), NFR thông lượng, cô lập tenant, đệm khi TSDB lỗi (mục 4), mục 7.1 (Batch, Sample), rủi ro AR-002, AR-003, AR-004. Mục tiêu L1: G1, G3 (qua L2) |
| Tài liệu cha | L2 SAD Collector |
| Tài liệu anh em | Ingest API, Registry, Enroll, Presence, Alerting, Outbox, Admin API |
| Tier | Đề xuất Cấp 3 (kế thừa L2, chưa xác nhận) |
| Phân loại dữ liệu | Nhạy cảm vừa: số liệu theo tenant. Không chứa bí mật |
| Blast radius | Ghi sai nhãn company_id làm lẫn tenant. Worker chậm làm đầy bus và ingest trả 503. Mất lô khi lỗi vĩnh viễn |
| Sign-off gate | Chưa chỉ định. Không có ai đã sign-off |
Truy vết mục L3 đến thành phần L2:
| Mục L3 | Thành phần hoặc mục L2 |
|---|---|
| 1 Phạm vi | 2.3 Các thành phần chính |
| 2 Yêu cầu | 3 Functional, 4 NFR |
| 3 Kiến trúc | 6.1 Component table |
| 4, 6 Domain, dữ liệu | 7.1 Data Model, 7.2 Data ownership |
| 5 Hợp đồng | 6.2 Integration (Worker đến VictoriaMetrics) |
| 7 Thuật toán | 8.1, 8.2 (luồng ingest) |
| 8, 9 Lỗi, suy thoái | 6.3 Resilience, 12.2 Reliability |
| 11 Bảo mật | 9 Security (nhãn tenant) |
| 13 Telemetry | 13 Observability |
1. Scope and Non-Goals
1.1 Vai trò
Tách phần nhận (Ingest) khỏi phần ghi (Worker) để agent không phải chờ TSDB. Bus giữ lô đã nhận theo shard (64 shard logic, khóa server_id), Worker gom mẫu thành lô lớn rồi ghi vào VictoriaMetrics. TSDB adapter là ranh giới duy nhất chạm kho chuỗi thời gian, chịu trách nhiệm ép nhãn company_id và server_id.
1.2 Context diagram
flowchart LR
classDef bc fill:#1f3a5f,stroke:#4a90d9,color:#fff;
classDef owned fill:#2d4a3e,stroke:#5fb37a,color:#fff;
classDef datastore fill:#3a2d4a,stroke:#a06fd9,color:#fff;
ING["Ingest API"]:::bc
ADM["Admin API"]:::bc
BUS[["Bus (shard)"]]:::datastore
WK["Worker"]:::owned
TS["TSDB adapter"]:::owned
VM[("VictoriaMetrics")]:::datastore
ING -.->|"publish lô"| BUS
WK -->|"đọc theo shard"| BUS
WK -->|"ghi lô"| TS
ADM -->|"truy vấn, xóa"| TS
TS -->|"import, query, delete"| VM| Chiều | Bên | Nội dung |
|---|---|---|
| Vào | Ingest API | Publish(batch) (không chặn, ErrFull khi đầy) |
| Vào | Admin API | Range, Latest, DeleteServer, DeleteCompany, Ping |
| Ra | VictoriaMetrics | POST /api/v1/import, /api/v1/query_range, /api/v1/query, /api/v1/admin/tsdb/delete_series, /health |
| Ra | Redis | Stream ah:bus:s:<n>, nhóm consumer workers |
1.3 Trong phạm vi
- Interface Bus, hiện thực inproc và Redis Streams.
- Ánh xạ khóa shard
ShardFor(server_id). - Worker: gom lô, ghi, thử lại, áp lực ngược, xả khi tắt.
- TSDB: interface
Writer,Reader,Deleter, hiện thực bộ nhớ và VictoriaMetrics. - Ép nhãn tenant, chuẩn hóa tên và nhãn, loại giá trị không hữu hạn.
- Kế hoạch truy vấn theo khoảng (thô, rollup 5 phút, rollup 1 giờ).
1.4 Ngoài phạm vi
| Nội dung | Thuộc về |
|---|---|
| Xác thực, kiểm giới hạn, giải mã body | L3 Ingest API |
| Định tuyến vai trò, khởi động, tắt êm | internal/app (L2 mục 6.1) |
| Vận hành VictoriaMetrics (lưu giữ, rollup bằng vmalert, sao lưu) | deploy/, docs/09 |
| Bộ đệm JetStream 2 giờ (NFR-07) | Chưa làm (COL-13 chỉ thiết kế) |
| Endpoint HTTP truy vấn cho người dùng | Access Hub gọi qua Admin API |
2. Requirements
2.1 Functional Requirements
| # | Trách nhiệm | Giải thích | Hiện thực ở |
|---|---|---|---|
| FR-BUS-01 | Shard ổn định | ShardFor = FNV-32a(server_id) mod 64, luôn trong khoảng | TestShardForIsStableAndInRange |
| FR-BUS-02 | Publish không chặn | Đầy thì ErrFull, đóng thì ErrClosed | TestPublishFullAndClosed, TestRedisBusFullAndClosed |
| FR-BUS-03 | Giữ thứ tự theo shard | Lô cùng shard đến worker theo thứ tự | TestSubscribeOrderPerShardAndDrainOnCancel, TestRedisBusRoundTripKeepsOrderAndContent |
| FR-BUS-04 | Chỉ nhận shard được giao | Worker chỉ đăng ký tập shard sở hữu | TestSubscribeOnlyOwnedShards, TestShardSubsetOnlyConsumesItsShards |
| FR-BUS-05 | Chia việc giữa nhiều consumer | Hai tiến trình chia shard, không xử lý đôi | TestRedisBusTwoConsumersSplitTheWork |
| FR-BUS-06 | Nhận lại việc của consumer chết | XAUTOCLAIM mục treo quá claim_idle | TestRedisBusClaimsEntriesOfDeadConsumer |
| FR-BUS-07 | Bỏ mục không giải mã được | Không kẹt cả shard vì một mục hỏng | TestRedisBusSkipsUndecodableEntries |
| FR-BUS-08 | Xả khi hủy | Đọc hết phần còn lại rồi mới dừng | TestSubscribeOrderPerShardAndDrainOnCancel, TestRedisBusDrainsOnCancel |
| FR-WRK-01 | Gom lô | Xả khi đủ 5000 mẫu hoặc sau 500 ms | TestFlushesWhenFull, TestFlushesAfterMaxWait |
| FR-WRK-02 | Nhãn tenant tới Writer | company_id và server_id đi kèm mọi lô | TestCompanyAndServerReachWriter |
| FR-WRK-03 | Thử lại lỗi tạm thời | Giữ lô, lùi 200 ms nhân đôi đến 30 giây | TestTemporaryErrorRetriesAndKeepsBatch |
| FR-WRK-04 | Bỏ lô lỗi vĩnh viễn | Không thử lại | TestPermanentErrorDrops |
| FR-WRK-05 | Áp lực ngược | Ngừng đọc bus khi mẫu chờ đạt 4 lần trần lô | TestBackpressureHoldsBusQueue |
| FR-WRK-06 | Xả khi tắt | Rút bus, ghi nốt trong shutdown_flush (10 giây) | TestShutdownDrainsBusAndFlushes, TestShutdownWhileTSDBDownRetriesInFinalFlush, TestFinalFlushGivesUpAtDeadline |
| FR-TSD-01 | Chuẩn hóa | Tên và nhãn về dạng hợp lệ, tiền tố ah_ | TestNormalize, TestVMWriteImportPayload |
| FR-TSD-02 | Ép nhãn tenant | Loại nhãn dành riêng của agent; company_id, server_id đặt sau cùng | TestCleanLabelsDropsReserved, TestCheckTenantAndPermanent |
| FR-TSD-03 | Từ chối thiếu tenant trước khi gọi mạng | Không gửi request khi thiếu company_id hoặc server_id | TestVMWriteRejectsMissingTenantWithoutRequest |
| FR-TSD-04 | Loại giá trị không hữu hạn | NaN, Inf bị bỏ | TestVMWriteDropsNonFinite |
| FR-TSD-05 | Phân loại lỗi | 5xx, 408, 429 tạm thời; còn lại vĩnh viễn | TestVMWriteErrorClassification |
| FR-TSD-06 | Truy vấn khoảng chọn rollup theo độ dài | Thô, 5 phút, 1 giờ | TestVMRangePicksRollupBySpan, TestVMRangeSelectorsAndParsing, TestVMRangeErrors |
| FR-TSD-07 | Giá trị mới nhất | last_over_time trong 10 phút | TestVMLatest, TestMemoryLatestAndDelete |
| FR-TSD-08 | Xóa theo tenant | delete_series với bộ chọn có company_id | TestVMDelete, TestMemoryLatestAndDelete |
| FR-TSD-09 | Kiểm tra sống | Ping gọi /health | TestVMPing |
2.2 Non-Functional Requirements
Allocated:
| NFR L2 | Target | Cách đáp ứng ở L3 |
|---|---|---|
| NFR-01 (thông lượng) | Gom lô lớn ghi ít lần | 5000 mẫu hoặc 500 ms, NDJSON một request |
| NFR-04 (tách nhận và ghi) | Ingest không chờ TSDB | Bus với Publish không chặn |
| NFR-12 (cô lập tenant) | Không lẫn tenant | Ép nhãn ở adapter, bộ chọn truy vấn luôn có company_id |
| NFR-13 (chịu lỗi worker) | Mất worker không mất lô | Redis Streams, consumer group, XAUTOCLAIM |
Inherited: client Redis từ internal/redisx (timeout, danh sách lệnh bị chặn), HTTP server của internal/app.
Owned:
| ID | Target | Parent L2-NFR | Satisfied-by |
|---|---|---|---|
| NFR-BUS-01 | Không mất lô chưa ACK khi consumer chết | NFR-13 | XAUTOCLAIM mỗi claim_idle/3 (10 giây với mặc định 30 giây), TestRedisBusClaimsEntriesOfDeadConsumer |
| NFR-BUS-02 | Không tăng vô hạn | NFR-04 | Trần mềm theo XLEN lấy mẫu mỗi 200 ms, queue_per_shard |
| NFR-WRK-01 | Bộ nhớ worker có trần | NFR-01 | Áp lực ngược ở 4 lần max_batch_samples |
| NFR-TSD-01 | Dải truy vấn tối đa 11.000 điểm, bước tối thiểu 30 giây | NFR-01 | Bộ lập kế hoạch truy vấn |
Không đặt số đo hiệu năng cụ thể vì chưa có kết quả đo trên VictoriaMetrics thật (đề xuất: đo bằng vmcheck và tải giả, xem OQ-BWT-1).
2.3 Acceptance Criteria
| AC | Given / When / Then | Test ID |
|---|---|---|
| AC-BWT-01 | Given hai consumer cùng nhóm, When có lô, Then mỗi lô xử lý đúng một lần | TestRedisBusTwoConsumersSplitTheWork |
| AC-BWT-02 | Given consumer chết giữa chừng, When qua claim_idle, Then consumer khác nhận lại | TestRedisBusClaimsEntriesOfDeadConsumer |
| AC-BWT-03 | Given TSDB lỗi tạm thời, When ghi, Then giữ lô và thử lại | TestTemporaryErrorRetriesAndKeepsBatch |
| AC-BWT-04 | Given TSDB trả 4xx (trừ 408, 429), When ghi, Then bỏ lô | TestPermanentErrorDrops, TestVMWriteErrorClassification |
| AC-BWT-05 | Given tắt dịch vụ, When còn mẫu, Then ghi nốt trước hạn 10 giây | TestShutdownDrainsBusAndFlushes |
| AC-BWT-06 | Given TSDB còn lỗi khi tắt, When quá hạn, Then bỏ cuộc | TestFinalFlushGivesUpAtDeadline |
| AC-BWT-07 | Given hai tenant, When ghi và truy vấn, Then không thấy dữ liệu của nhau | TestIngestToTSDBPipelineKeepsTenantsApart, TestMemoryDedupAndIsolation |
| AC-BWT-08 | Given lô thiếu company_id, When ghi, Then từ chối mà không gọi mạng | TestVMWriteRejectsMissingTenantWithoutRequest |
2.4 Quality Attribute Scenarios
| Nguồn | Kích thích | Môi trường | Phản hồi | Thước đo |
|---|---|---|---|---|
| Agent | Gửi liên tục | VictoriaMetrics chậm | Worker đạt ngưỡng áp lực ngược, bus đầy dần, Ingest trả 503 kèm Retry-After | Không tăng bộ nhớ vô hạn |
| Vận hành | Dừng một node worker | Bình thường | Node khác nhận lại lô treo | Sau khoảng claim_idle cộng chu kỳ quét (đề xuất kiểm chứng) |
| Vận hành | Khởi động lại collector | Còn mẫu chưa ghi | Xả nốt, hạn 10 giây | Mẫu còn lại sau hạn bị mất (đo bằng ahc_tsdb_batches_dropped_total) |
| Access Hub | Truy vấn khoảng 60 ngày | Bình thường | Dùng rollup 1 giờ từ long_url | Không quá 11.000 điểm |
3. Application Architecture
3.1 Kiến trúc thời chạy (C&C)
flowchart LR
classDef bc fill:#1f3a5f,stroke:#4a90d9,color:#fff;
classDef owned fill:#2d4a3e,stroke:#5fb37a,color:#fff;
classDef datastore fill:#3a2d4a,stroke:#a06fd9,color:#fff;
PUB["Ingest publisher"]:::bc
BUS[["Redis Streams"]]:::datastore
SW["Shard worker x N"]:::owned
BT["Batcher"]:::owned
WR["Writer (VM)"]:::owned
VM[("VictoriaMetrics")]:::datastore
PUB -.->|"XADD"| BUS
SW -->|"XREADGROUP"| BUS
SW --> BT
BT --> WR
WR -->|"POST import"| VMMũi tên chỉ chiều khởi tạo lời gọi. Worker là consumer nên trỏ tới broker. Publisher ghi vào broker (nét đứt vì bất đồng bộ).
| Thành phần | Trách nhiệm | Vòng đời |
|---|---|---|
| Bus inproc | Kênh Go theo shard, dùng khi một tiến trình | Theo tiến trình, mất khi restart |
| Bus Redis | Stream ah:bus:s:<n> mỗi shard, nhóm workers, XREADGROUP COUNT 32, XACK và XDEL sau khi xử lý | Bền trong Redis |
| Shard worker | Một goroutine mỗi shard sở hữu: đọc, gom, ghi | Theo tiến trình |
| Batcher | Gom mẫu theo (company_id, server_id), xả theo kích thước hoặc thời gian | Theo shard worker |
| Writer | Hiện thực tsdb.Writer (VictoriaMetrics hoặc bộ nhớ) | Theo tiến trình |
| Kết nối | Kiểu | Chi tiết |
|---|---|---|
| Ingest đến Bus | Bất đồng bộ | Publish không chặn, ErrFull khi vượt trần mềm |
| Worker đến Bus | Đồng bộ (kéo) | Khối bus.block, tối đa bus.readers reader (chi tiết khóa cấu hình xem 12.1) |
| Worker đến VictoriaMetrics | Đồng bộ | tsdb.timeout |
3.2 Module view
flowchart TB
classDef bc fill:#1f3a5f,stroke:#4a90d9,color:#fff;
classDef owned fill:#2d4a3e,stroke:#5fb37a,color:#fff;
APP["internal/app"]:::bc
ING["internal/ingest"]:::bc
ADM["internal/admin"]:::bc
BUSP["internal/bus"]:::owned
WKP["internal/worker"]:::owned
TSP["internal/tsdb"]:::owned
RX["internal/redisx"]:::owned
APP --> BUSP
APP --> WKP
APP --> TSP
ING --> BUSP
ADM --> TSP
WKP --> BUSP
WKP --> TSP
BUSP --> RXTệp: internal/bus/bus.go (interface, ShardFor, inproc), redis.go (Redis Streams). internal/worker/worker.go. internal/tsdb/tsdb.go (interface, chuẩn hóa), memory.go, vm.go (VictoriaMetrics).
4. Domain Model
classDiagram
class Batch {
<<bus>>
company_id
server_id
agent_id
samples[]
}
class Sample {
name
labels
value
timestamp_ms
}
class Bus {
<<interface>>
Publish
Subscribe
Depth
Close
}
class Writer {
<<interface>>
Write
}
class Reader {
<<interface>>
Range
Latest
}
class Deleter {
<<interface>>
DeleteServer
DeleteCompany
}
Batch "1" --> "*" Sample
Bus --> Batch
Writer --> BatchBất biến:
company_idvàserver_idcủa lô lấy từ danh tính đã xác thực (registry), không từ payload agent.- Chỉ chấp nhận giá trị hữu hạn; nhãn dành riêng của agent bị loại; tên chuẩn hóa và thêm tiền tố
ah_. ShardForchỉ phụ thuộcserver_id, nên mọi mẫu của một máy chủ đi qua cùng shard (giữ thứ tự).- Redis là nguồn thứ tự cho từng shard. Ghi nhãn tenant thực hiện ở Writer, không ở Bus.
5. API Contract
Không có API HTTP phơi bày ra ngoài. Hợp đồng gồm interface Go nội bộ và các lời gọi HTTP tới VictoriaMetrics.
5.1 Operations
| Interface | Phương thức | Hành vi |
|---|---|---|
Bus | Publish(batch) | Không chặn, ErrFull, ErrClosed |
Bus | Subscribe(shards, handler) | Gọi handler theo thứ tự trong shard, chạy đến khi hủy rồi xả |
Bus | Depth() | Độ sâu tổng (nguồn của ahc_bus_depth và ahc_ingest_bus_depth) |
Writer | Write(batch) | Lỗi tạm thời hoặc vĩnh viễn (IsPermanent) |
Reader | Range, Latest | Truy vấn có company_id bắt buộc |
Deleter | DeleteServer, DeleteCompany | Xóa chuỗi theo bộ chọn |
Ping | Kiểm tra sống (dùng cho /readyz) |
VictoriaMetrics (gọi ra):
| Method | Path | Dùng cho |
|---|---|---|
| POST | /api/v1/import | Ghi NDJSON ({"metric":{...},"values":[...],"timestamps":[...]}) |
| GET, POST | /api/v1/query_range | Range |
| GET, POST | /api/v1/query | Latest (last_over_time trong 10 phút) |
| POST | /api/v1/admin/tsdb/delete_series | Xóa dữ liệu |
| GET | /health | Ping |
Ghi chú lệch tài liệu: ADR 0005 và docs/03 nói remote-write; mã dùng JSON import (/api/v1/import). Ghi nhận thành đề xuất ADR 0015 (xem L2 TD list). Chưa kiểm chứng trên VictoriaMetrics thật trong phiên soạn (vmcheck chỉ chạy khi bật).
5.2 Request và Response schema
Ghi (import): mỗi dòng NDJSON
metric!: { __name__!: "ah_<tên>", company_id!, server_id!, <nhãn khác>? }
values!: number[] (hữu hạn)
timestamps!: int[] (ms)
Range(company!, server?, metric!, from!, to!, agg?) -> series[]: { labels, points[]: [ts, value] }
Latest(company!, server!, metric!) -> { value, ts } hoặc rỗngBộ chọn truy vấn luôn chứa company_id="..." (không có đường bỏ qua).
5.3 Error codes
| Tình huống | Kết quả nội bộ | Ghi chú |
|---|---|---|
| VM trả 5xx, 408, 429 | Lỗi tạm thời | Worker thử lại |
| VM trả mã khác (4xx) | Lỗi vĩnh viễn | Worker bỏ lô |
| Thiếu tenant | Lỗi vĩnh viễn, không gọi mạng | |
| Khoảng truy vấn không hợp lệ | Lỗi tham số | Admin trả 400 (xem L3 Admin) |
| Backend truy vấn lỗi | Admin trả 503 unavailable |
5.4 Versioning
Interface nội bộ đổi cùng mã. Với VictoriaMetrics: gắn chặt tiền tố /api/v1/... của VictoriaMetrics; nâng phiên bản VictoriaMetrics phải kiểm thử lại định dạng import (đề xuất ghi vào runbook nâng cấp).
5.5 Authz
Không có xác thực người dùng ở tầng này; tin cậy tuyệt đối vào lời gọi nội bộ. Kết nối tới VictoriaMetrics trong dev là HTTP thuần nội bộ; sản xuất nên đặt VictoriaMetrics chỉ nghe trên mạng nội bộ (đề xuất; nội dung docs/09). Không có cơ chế xác thực tới VictoriaMetrics trong mã (đề xuất ghi nhận rủi ro).
6. Physical Data Schema
6.1 Mapping
Redis Streams:
| Khóa | Kiểu | Nội dung |
|---|---|---|
bus:s:<n> (n từ 0 đến 63) | stream | Mỗi mục là một lô đã mã hóa; nhóm consumer workers |
Mã hóa mục: lô tuần tự hóa (định dạng cụ thể xem internal/bus/redis.go, chưa trích ở bản nháp này). XACK rồi XDEL sau khi xử lý xong.
VictoriaMetrics (chuỗi thời gian):
| Thành phần | Quy ước |
|---|---|
| Tên chuỗi | ah_<tên chuẩn hóa> |
| Nhãn bắt buộc | company_id, server_id |
| Nhãn khác | Nhãn từ agent đã lọc (loại nhãn dành riêng) |
| Rollup | Chuỗi <tên>:5m_<agg> và <tên>:1h_<agg> do vmalert tạo (đường deploy/, ngoài phạm vi mã Go) |
flowchart LR
classDef datastore fill:#3a2d4a,stroke:#a06fd9,color:#fff;
RAW[("ah_name thô")]:::datastore
R5[("ah_name:5m_agg")]:::datastore
R1[("ah_name:1h_agg long_url")]:::datastore
RAW -.->|"vmalert rollup"| R5
R5 -.->|"vmalert rollup"| R1Kế hoạch truy vấn: khoảng đến 7 ngày dùng chuỗi thô; đến 30 ngày dùng :5m_<agg>; xa hơn dùng :1h_<agg> từ tsdb.long_url. Bước tối thiểu 30 giây, tối đa 11.000 điểm.
6.2 Phân loại và lưu giữ
| Dữ liệu | Phân loại | Lưu giữ |
|---|---|---|
| Stream bus | Trung bình | Xóa sau khi ACK; trần mềm theo queue_per_shard |
| Chuỗi thô | Trung bình | Theo cấu hình VictoriaMetrics (-retentionPeriod, xem deploy/), chưa ghi thành số trong bản nháp này |
| Rollup | Trung bình | Theo cấu hình long_url |
Mục tiêu lưu giữ chính thức: L2 mục 7.3 và docs/09. Bản nháp này không nhắc lại số để tránh lệch.
7. Algorithms
7.1 Sequences (đường thành công)
sequenceDiagram
participant IN as Ingest
participant BS as Bus (Redis)
participant WK as Worker
participant TS as TSDB adapter
participant VM as VictoriaMetrics (ext)
IN->>BS: Publish(lô, shard = ShardFor(server_id))
BS-->>IN: OK
WK->>BS: XREADGROUP (COUNT 32)
BS-->>WK: lô
WK->>WK: gom đến 5000 mẫu hoặc 500 ms
WK->>TS: Write(company, server, mẫu)
TS->>VM: POST /api/v1/import (NDJSON)
VM-->>TS: 204
WK->>BS: XACK và XDEL7.2 Sequences (đường lỗi)
sequenceDiagram
participant WK as Worker
participant TS as TSDB adapter
participant VM as VictoriaMetrics (ext)
WK->>TS: Write(lô)
TS->>VM: POST /api/v1/import
alt lỗi tạm thời (5xx, 408, 429, mạng)
VM--)TS: lỗi
TS-->>WK: lỗi tạm thời
WK->>WK: lùi 200 ms nhân đôi đến 30 s, giữ lô
WK->>TS: Write lại
else lỗi vĩnh viễn (4xx khác)
VM-->>TS: 4xx
TS-->>WK: lỗi vĩnh viễn
WK->>WK: bỏ lô, ghi ahc_tsdb_batches_dropped_total
end7.3 State machines
Worker có bốn trạng thái:
stateDiagram-v2
direction LR
state "Đang thu (collecting)" as COL
state "Đang xả (flushing)" as FL
state "Đang chờ thử lại (backoff)" as BO
state "Đang tắt (draining)" as DR
[*] --> COL
COL --> FL: đủ mẫu hoặc hết thời gian
FL --> COL: ghi thành công
FL --> BO: lỗi tạm thời
BO --> FL: hết thời gian chờ
FL --> COL: lỗi vĩnh viễn (bỏ lô)
COL --> DR: hủy ngữ cảnh
BO --> DR: hủy ngữ cảnh
DR --> [*]: xả xong hoặc hết hạn 10 sÁp lực ngược: khi số mẫu chờ đạt 4 lần max_batch_samples, worker ngừng đọc bus, bus đầy dần đến ErrFull và Ingest trả 503.
| Chuyển trạng thái | Sequence |
|---|---|
| Thu đến xả | 7.1 |
| Xả đến chờ | 7.2 nhánh tạm thời |
| Xả đến thu (bỏ lô) | 7.2 nhánh vĩnh viễn |
| Thu đến tắt | Tắt êm (L2 mục 8) |
7.4 Core algorithms
| Vấn đề | Giải pháp | Trade-off |
|---|---|---|
| Thứ tự theo máy chủ | ShardFor = FNV-32a mod 64 và một reader mỗi shard tại một thời điểm | Máy chủ nóng làm nóng một shard (không cân bằng theo tải) |
| Không lẫn tenant | Adapter ép company_id, server_id cuối cùng, ghi đè bất kỳ nhãn cùng tên | Nhãn agent trùng tên bị mất (đã chủ đích) |
| Trần bus với Redis | Lấy mẫu XLEN mỗi 200 ms so queue_per_shard | Có thể vượt trần nhẹ giữa hai lần lấy mẫu |
| Bus Redis chia việc | Nhóm workers, XAUTOCLAIM mỗi claim_idle/3 (10 giây mặc định) cho mục treo quá claim_idle (30 giây mặc định) | Lô có thể ghi hai lần nếu consumer chậm bị nhận lại (at-least-once, VictoriaMetrics khử trùng theo dấu thời gian và nhãn, chưa kiểm chứng ở mọi trường hợp) |
| Lỗi mạng ghi | Phân loại lỗi theo mã HTTP | Lỗi vĩnh viễn làm mất lô (không có dead-letter cho mẫu) |
| Chọn rollup | Chọn chuỗi theo độ dài khoảng | Rollup cần vmalert chạy; không có rollup thì truy vấn dài rỗng |
| Giá trị mới nhất | last_over_time 10 phút | Agent im hơn 10 phút thì không có giá trị mới nhất |
Không có bộ ngắt mạch (circuit breaker) quanh VictoriaMetrics; chỉ có lùi mũ và áp lực ngược (TD của L2).
8. Error Handling
8.1 Bảng lỗi theo bước
| Bước lỗi | Nguyên nhân | Cơ chế xử lý | Trạng thái cuối |
|---|---|---|---|
| Publish | Bus đầy | ErrFull, Ingest trả 503 | Agent thử lại từ WAL |
| Publish | Bus đóng | ErrClosed | Ingest trả lỗi khi tắt |
| Đọc bus | Redis mất kết nối | Thử lại vòng đọc | Lô vẫn trong stream |
| Giải mã mục | Mục hỏng | Bỏ qua và ACK để không kẹt shard | Mất mục đó |
| Ghi VM | Lỗi tạm thời | Giữ lô, lùi 200 ms đến 30 s | Thử đến khi thành công hoặc tắt |
| Ghi VM | Lỗi vĩnh viễn | Bỏ lô | ahc_tsdb_batches_dropped_total tăng |
| Tắt | TSDB vẫn lỗi | Thử trong hạn 10 s rồi bỏ | Mất mẫu còn lại |
| Truy vấn | VM lỗi | Trả lỗi cho Admin | Admin trả 503 |
8.2 Fail-fast
tsdb.urlrỗng khi vai trò cần TSDB: cảnh báo (TestTSDBWarningOnlyForRolesThatUseTheStore); readiness kiểm tra "tsdb".- Tách
ingestvàworkerbắt buộc Redis (TestSplitIngestWorkerRejectsLocalBackends,TestSplitRolesNeedRedis). TestNewRedisRequiresClients: bus Redis cần client.
8.3 Race conditions
| Tình huống | Xử lý |
|---|---|
Hai consumer nhận cùng mục sau XAUTOCLAIM | At-least-once; ghi lặp cùng mẫu nhận diện được ở VictoriaMetrics (khử trùng) |
| Publish khi đóng | ErrClosed |
| Đóng khi còn mục | Xả trước khi thoát (TestRedisBusDrainsOnCancel) |
9. Degradation
Nguyên tắc: ưu tiên giữ ingest chạy và giữ dữ liệu trong bus càng lâu càng tốt, chỉ từ chối khi bus đầy.
9.1 Dependency matrix
| Dependency lỗi | Hành vi | Kiểu suy thoái | Hệ quả |
|---|---|---|---|
| VictoriaMetrics | Worker thử lại, áp lực ngược, bus đầy, Ingest 503 | Giảm chức năng | Agent đệm WAL đến 50 MiB, 24 giờ (phía agent) |
| Redis (bus) | Ingest không publish được | Từ chối | Ingest 503; readiness "redis" đỏ |
| Một worker chết | Nhận lại sau claim_idle | Trễ | Ghi trễ vài chục giây (đề xuất kiểm chứng) |
| Cả hai | Agent đệm WAL | Mất kết nối | Phục hồi theo WAL |
9.2 Backup và recovery
Bus không phải nơi lưu lâu; mất Redis (khi AOF không bảo vệ) làm mất lô chưa ghi. Redis dev có AOF everysec (mất tối đa khoảng 1 giây). Dữ liệu chuỗi thời gian thuộc VictoriaMetrics: sao lưu bằng snapshot của VictoriaMetrics (thuộc docs/09, chưa có trong mã). RTO và RPO chính thức: chưa định nghĩa (đề xuất xác định ở cấp L1, OQ-BWT-3). Bộ đệm JetStream 2 giờ (NFR-07) chưa làm.
10. Concurrency
10.1 Ranh giới giao dịch
Không có giao dịch xuyên hệ thống. Mỗi lô là đơn vị at-least-once: XADD (Ingest), XREADGROUP rồi ghi VM rồi XACK (Worker). Nếu tiến trình chết giữa ghi VM và ACK thì lô ghi lại lần nữa.
10.2 Cơ chế
| Cơ chế | Mô tả | Mối nguy |
|---|---|---|
| Consumer group | Chia shard giữa consumer | Không có khóa sở hữu shard trong mã (chỉ dựa vào nhóm); ADR 0009 nói lease theo shard, mã chưa có (COL-13) |
XAUTOCLAIM | Nhận việc treo | Ghi lặp |
| Goroutine mỗi shard | Thứ tự trong shard | Shard nóng |
| Áp lực ngược | Ngừng đọc | Tăng độ trễ ghi |
Trần mềm XLEN | Chặn publish | Vượt nhẹ giữa hai lần lấy mẫu |
11. Security
11.1 Ba lớp
| Lớp | Biện pháp |
|---|---|
| Truyền thông | Redis có mật khẩu, TLS tùy chọn; VictoriaMetrics chỉ trong mạng nội bộ (đề xuất; mã không xác thực tới VM) |
| Mã | Ép nhãn tenant ở một nơi; loại nhãn dành riêng; xác thực định dạng tên |
| Dữ liệu | Số liệu không chứa bí mật; xóa theo company_id qua Admin |
11.2 Pipeline phân quyền năm bước
Thành phần này không tự phân quyền: nhận company_id, server_id đã được Ingest xác thực (registry). Điều kiện tiên quyết: company_id và server_id khác rỗng, đã kiểm tra ở Write (không gọi mạng nếu thiếu).
11.3 Chỉ mục neo
| Nguyên tắc L2 | Biện pháp nội bộ |
|---|---|
| NFR-12 cô lập tenant | Ép nhãn ở adapter, bộ chọn truy vấn luôn có company_id |
| P1 danh tính từ registry | Tenant lấy từ bản ghi registry, không từ payload |
| Xóa dữ liệu | DeleteServer, DeleteCompany tìm theo nhãn |
Rủi ro còn mở: không có xác thực tới VictoriaMetrics (OQ-BWT-2); giá trị nhãn company_id nhiễm ký tự đặc biệt vào bộ chọn truy vấn được chặn bởi tsdb.ValidID ở Admin, đề xuất rà soát thêm đường Ingest.
12. Configuration
12.1 Tunables
Biến môi trường: AHC_TSDB_URL, AHC_TSDB_LONG_URL, AHC_REDIS_*. Còn lại chỉ YAML.
| Tham số | Default | Ý nghĩa và tác động | Mục liên quan |
|---|---|---|---|
bus.backend | auto | auto: Redis nếu có redis.addr, ngược lại inproc | 3.1 |
bus.queue_per_shard | 256 | Trần mềm mỗi shard (dự phòng của NewRedis là 1000, lệch cấu hình) | 7.4 |
bus.readers | 8 | Số reader | 3.1 |
bus.block | 500 ms | Thời gian chờ mỗi lần XREADGROUP | 3.1 |
bus.claim_idle | 30 giây | Ngưỡng nhận lại mục treo (quét mỗi một phần ba, tức 10 giây) | 7.4 |
worker.shards | tất cả 64 | Tập shard sở hữu của tiến trình | 3.1 |
worker.shutdown_flush | 10 giây | Hạn xả khi tắt | 7.3 |
tsdb.url | rỗng | Địa chỉ VictoriaMetrics ghi và truy vấn gần | 5.1 |
tsdb.long_url | rỗng | Nguồn rollup 1 giờ | 6.1 |
tsdb.timeout | 10 giây | Timeout HTTP | 8.1 |
tsdb.max_batch_samples | 5000 | Trần mẫu mỗi lô | 7.1 |
tsdb.max_batch_wait | 500 ms | Thời gian chờ gom tối đa | 7.1 |
tsdb.retry_max | 30 giây | Trần lùi mũ | 7.2 |
Mặc định lấy từ Default() trong internal/config/config.go. Hằng số cố định trong mã: bước truy vấn tối thiểu 30 giây, tối đa 11.000 điểm, cửa sổ Latest 10 phút, lô đọc 32, chu kỳ quét XAUTOCLAIM bằng claim_idle chia 3 (10 giây mặc định).
12.2 Feature flags
Không có. bus.backend và redis.addr chọn hiện thực bus. tsdb.url rỗng chọn ghi vào bộ nhớ (chỉ dev và kiểm thử).
13. Telemetry
13.1 Metrics
| Tên | Kiểu | Nhãn | Ngữ nghĩa |
|---|---|---|---|
ahc_bus_depth | gauge | Độ sâu bus tổng (không nhãn shard) | |
ahc_ingest_bus_depth | gauge | Độ sâu nhìn từ Ingest | |
ahc_tsdb_samples_written_total | counter | Mẫu ghi thành công | |
ahc_tsdb_write_errors_total | counter | kind | Lỗi ghi theo loại |
ahc_tsdb_batches_dropped_total | counter | Lô bị bỏ (lỗi vĩnh viễn hoặc hết hạn tắt) | |
ahc_worker_batches_total | counter | Số lô đã xử lý | |
ahc_tsdb_flush_seconds | histogram | Thời gian xả | |
ahc_worker_pending_samples | gauge | Mẫu chờ trong worker |
Lệch tài liệu: docs/10 dùng tên khác (ví dụ độ sâu theo shard, ahc_worker_lag). Mã thắng, xem L2 TD-004. Không có nhãn shard nên khó phát hiện shard nóng (đề xuất bổ sung).
13.2 Log schema
JSON một dòng: time, level, msg, collector_id. Sự kiện: khởi động và dừng worker, lỗi ghi VM (mã, loại), bỏ lô, mục bus hỏng, nhận lại mục treo. Không ghi giá trị mẫu.
13.3 Cảnh báo đến runbook
Chưa có luật cảnh báo trong mã. Đề xuất: ahc_bus_depth cao kéo dài, ahc_tsdb_write_errors_total tăng, ahc_tsdb_batches_dropped_total tăng bất kỳ, ahc_worker_pending_samples chạm trần. Runbook: docs/09 mục 6.
13.4 Probes
| Probe | Điểm cuối | Ý nghĩa |
|---|---|---|
| Liveness | /healthz | Tiến trình sống |
| Readiness | /readyz ("redis", "tsdb") | Redis và VictoriaMetrics trả lời |
13.5 Trace propagation
Chưa có trace qua bus (lô không mang X-Request-Id; đề xuất bổ sung nếu cần điều tra).
14. Test Plan
| Loại | Kiểm thử | Vị trí |
|---|---|---|
| Đơn vị Bus | TestShardForIsStableAndInRange, TestPublishFullAndClosed, TestSubscribeOrderPerShardAndDrainOnCancel, TestSubscribeOnlyOwnedShards, TestBatchSamples | internal/bus |
| Bus Redis | TestRedisBusRoundTripKeepsOrderAndContent, TestRedisBusFullAndClosed, TestRedisBusDrainsOnCancel, TestRedisBusTwoConsumersSplitTheWork, TestRedisBusClaimsEntriesOfDeadConsumer, TestRedisBusSkipsUndecodableEntries, TestNewRedisRequiresClients | internal/bus |
| Worker | TestFlushesAfterMaxWait, TestFlushesWhenFull, TestCompanyAndServerReachWriter, TestTemporaryErrorRetriesAndKeepsBatch, TestPermanentErrorDrops, TestShutdownDrainsBusAndFlushes, TestShutdownWhileTSDBDownRetriesInFinalFlush, TestFinalFlushGivesUpAtDeadline, TestBackpressureHoldsBusQueue, TestOnBatchHookAndEmptyBatch, TestShardSubsetOnlyConsumesItsShards | internal/worker |
| TSDB chung và bộ nhớ | TestNormalize, TestCheckTenantAndPermanent, TestCleanLabelsDropsReserved, TestMemoryDedupAndIsolation, TestMemoryRangeAggregates, TestMemoryLatestAndDelete | internal/tsdb |
| VictoriaMetrics (giả lập HTTP) | TestVMWriteImportPayload, TestVMWriteRejectsMissingTenantWithoutRequest, TestVMWriteErrorClassification, TestVMWriteDropsNonFinite, TestVMRangeSelectorsAndParsing, TestVMRangePicksRollupBySpan, TestVMRangeErrors, TestVMLatest, TestVMDelete, TestVMPing | internal/tsdb |
| Tích hợp | TestIngestToTSDBPipelineKeepsTenantsApart, TestWorkerFlushesOnShutdown, TestSplitProcessesShareStateThroughRedis | internal/app |
| Chưa có | Kiểm thử với VictoriaMetrics thật (vmcheck opt-in), tải, đo độ trễ nhận lại XAUTOCLAIM, kiểm tra khử trùng khi ghi lặp |
Bản nháp này không chạy kiểm thử (tránh Redis 6380 và VictoriaMetrics sống).
15. Implementation Sequence
15.1 Milestone matrix
| Milestone | Nội dung | Phụ thuộc | Đóng góp nghiệm thu |
|---|---|---|---|
| COL-1 | Khung ứng dụng, cấu hình | Nền | |
| COL-5 | Bus inproc, worker, ghi TSDB theo lô, adapter VictoriaMetrics (/api/v1/import), kho bộ nhớ | COL-1 | AC-BWT-03..06 |
| COL-7 | Truy vấn Latest, Range qua Admin API (kế hoạch rollup trong mã) | COL-5 | AC-BWT-04, 07, 08 |
| COL-R | Bus Redis Streams, tách vai trò | COL-5 | AC-BWT-01, 02 |
| COL-10 (xong 2026-10-01) | Rollup (vmalert, vm-long), luật sinh bằng collector -rollup-rules, API series chọn nguồn theo tuổi của from và theo step, gộp lại ô rollup, đọc lại vm-raw khi rollup trống (internal/tsdb/rollup.go, vm.go) | COL-R | AC-BWT-08 |
| COL-13 (chưa làm) | Bus JetStream (bộ đệm 2 giờ, NFR-07), lease theo shard, phát lại | COL-R | NFR-07, khớp ADR 0009 |
Tên và thứ tự milestone lấy theo docs/11-roadmap.md; nếu có chênh, roadmap thắng.
15.2 Dependency flowchart
flowchart LR
C1["COL-1 khung"] --> C5["COL-5 bus, worker"]
C5 --> C7["COL-7 truy vấn qua admin"]
C5 --> C10["COL-10 rollup (xong)"]
C5 --> CR["COL-R Redis Streams"]
CR --> C13["COL-13 JetStream, lease shard"]Appendix A. Open Questions
| # | Câu hỏi | Hành vi tạm thời | Owner | Mã theo dõi |
|---|---|---|---|---|
| 1 | Kết quả đo trên VictoriaMetrics thật (JSON import, khử trùng, thông lượng) | Chưa đo; vmcheck opt-in | Chưa chỉ định | OQ-BWT-1 |
| 2 | Xác thực tới VictoriaMetrics (nay không có) | Mạng nội bộ | Chưa chỉ định | OQ-BWT-2 |
| 3 | RTO và RPO chính thức của đường ghi | Chưa định nghĩa | Chưa chỉ định | OQ-BWT-3 |
| 4 | bus.queue_per_shard mặc định 256 nhưng dự phòng NewRedis là 1000: thống nhất một giá trị | Giữ nguyên hai giá trị | Chưa chỉ định | OQ-BWT-4 |
| 5 | Thay JSON import bằng remote-write (ADR 0005 và 0015) | Giữ JSON import | Chưa chỉ định | OQ-BWT-5 |
| 6 | Bổ sung lease theo shard (COL-13, ADR 0009 và 0012) | Chỉ dùng consumer group | Chưa chỉ định | OQ-BWT-6 |
| 7 | Nhãn shard cho ahc_bus_depth để thấy shard nóng | Không có | Chưa chỉ định | OQ-BWT-7 |
| 8 | Có cần dead-letter cho lô bị bỏ vì lỗi vĩnh viễn | Bỏ và đếm | Chưa chỉ định | OQ-BWT-8 |
Appendix B. ADR nội bộ
| Mã ADR | Quyết định | Trạng thái | Driver |
|---|---|---|---|
| ADR 0005 | Nhập liệu vào VictoriaMetrics bằng remote-write | Lệch với mã (mã dùng /api/v1/import) | Xem đề xuất ADR 0015 |
| ADR 0009 | Bus theo shard và lease | Lệch với mã (Redis Streams, chưa lease) | Xem đề xuất ADR 0012 |
| ADR 0012 (đề xuất) | Redis Streams làm bus, trạng thái dùng chung | Đề xuất | COL-R |
| ADR 0015 (đề xuất) | JSON import thay remote-write | Đề xuất | Đơn giản, dễ kiểm thử |
| ADR nội bộ (đề xuất) | Lô lỗi vĩnh viễn bị bỏ, không dead-letter | Đề xuất, là hành vi mã | Tránh kẹt shard |
| ADR nội bộ (đề xuất) | Nhãn tenant ép ở adapter, không ở Bus | Đề xuất, là hành vi mã | Một điểm kiểm soát |
Appendix C. Section Profile
| Phân mục | Hồ sơ | Trạng thái điền | Giải trình |
|---|---|---|---|
| 0 đến 4 | Bắt buộc | Đã điền | |
| 5 API contract | Bắt buộc | Điền theo interface nội bộ và VM API | Không có API HTTP mở ra ngoài |
| 6 Physical data | Bắt buộc | Điền | Định dạng mục stream chưa trích |
| 7.1 đến 7.4 | Bắt buộc | Đã điền | |
| 8, 9, 10, 11 | Bắt buộc | Đã điền | |
| 12 | Bắt buộc | Đã điền | Mặc định lấy từ Default() |
| 13, 14, 15 | Bắt buộc | Đã điền | Kiểm thử chưa chạy trong bản nháp |