Access Hub Collector
Đồng bộ từ mã nguồn lúc 10:57, 03/10/2026
Skip to content

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.

Trạng thái
Bản nháp
Phiên bản
0.1, ngày 30/09/2026
Tài liệu cha
L2 SAD Collector
Cấp độ hệ thống
Cấp 3, đề xuất

Ghi chú: khi tài liệu và mã khác nhau, mã thắng

Đã làm nghĩa là có mã và kiểm thử trong repo. Chỉ thiết kế nghĩa là có trong tài liệu nhưng chưa có mã. Chỗ lệch được ghi ở mục nợ kỹ thuật.

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áiBẢ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 L2L2 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 chaL2 SAD Collector
Tài liệu anh emIngest 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ệuNhạy cảm vừa: số liệu theo tenant. Không chứa bí mật
Blast radiusGhi 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 gateChưa chỉ định. Không có ai đã sign-off

Truy vết mục L3 đến thành phần L2:

Mục L3Thành phần hoặc mục L2
1 Phạm vi2.3 Các thành phần chính
2 Yêu cầu3 Functional, 4 NFR
3 Kiến trúc6.1 Component table
4, 6 Domain, dữ liệu7.1 Data Model, 7.2 Data ownership
5 Hợp đồng6.2 Integration (Worker đến VictoriaMetrics)
7 Thuật toán8.1, 8.2 (luồng ingest)
8, 9 Lỗi, suy thoái6.3 Resilience, 12.2 Reliability
11 Bảo mật9 Security (nhãn tenant)
13 Telemetry13 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 ​

mermaid
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ềuBênNội dung
VàoIngest APIPublish(batch) (không chặn, ErrFull khi đầy)
VàoAdmin APIRange, Latest, DeleteServer, DeleteCompany, Ping
RaVictoriaMetricsPOST /api/v1/import, /api/v1/query_range, /api/v1/query, /api/v1/admin/tsdb/delete_series, /health
RaRedisStream 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 dungThuộc về
Xác thực, kiểm giới hạn, giải mã bodyL3 Ingest API
Định tuyến vai trò, khởi động, tắt êminternal/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ùngAccess Hub gọi qua Admin API

2. Requirements ​

2.1 Functional Requirements ​

#Trách nhiệmGiải thíchHiện thực ở
FR-BUS-01Shard ổn địnhShardFor = FNV-32a(server_id) mod 64, luôn trong khoảngTestShardForIsStableAndInRange
FR-BUS-02Publish không chặnĐầy thì ErrFull, đóng thì ErrClosedTestPublishFullAndClosed, TestRedisBusFullAndClosed
FR-BUS-03Giữ thứ tự theo shardLô cùng shard đến worker theo thứ tựTestSubscribeOrderPerShardAndDrainOnCancel, TestRedisBusRoundTripKeepsOrderAndContent
FR-BUS-04Chỉ nhận shard được giaoWorker chỉ đăng ký tập shard sở hữuTestSubscribeOnlyOwnedShards, TestShardSubsetOnlyConsumesItsShards
FR-BUS-05Chia việc giữa nhiều consumerHai tiến trình chia shard, không xử lý đôiTestRedisBusTwoConsumersSplitTheWork
FR-BUS-06Nhận lại việc của consumer chếtXAUTOCLAIM mục treo quá claim_idleTestRedisBusClaimsEntriesOfDeadConsumer
FR-BUS-07Bỏ mục không giải mã đượcKhông kẹt cả shard vì một mục hỏngTestRedisBusSkipsUndecodableEntries
FR-BUS-08Xả khi hủyĐọc hết phần còn lại rồi mới dừngTestSubscribeOrderPerShardAndDrainOnCancel, TestRedisBusDrainsOnCancel
FR-WRK-01Gom lôXả khi đủ 5000 mẫu hoặc sau 500 msTestFlushesWhenFull, TestFlushesAfterMaxWait
FR-WRK-02Nhãn tenant tới Writercompany_id và server_id đi kèm mọi lôTestCompanyAndServerReachWriter
FR-WRK-03Thử lại lỗi tạm thờiGiữ lô, lùi 200 ms nhân đôi đến 30 giâyTestTemporaryErrorRetriesAndKeepsBatch
FR-WRK-04Bỏ lô lỗi vĩnh viễnKhông thử lạiTestPermanentErrorDrops
FR-WRK-05Áp lực ngượcNgừng đọc bus khi mẫu chờ đạt 4 lần trần lôTestBackpressureHoldsBusQueue
FR-WRK-06Xả khi tắtRút bus, ghi nốt trong shutdown_flush (10 giây)TestShutdownDrainsBusAndFlushes, TestShutdownWhileTSDBDownRetriesInFinalFlush, TestFinalFlushGivesUpAtDeadline
FR-TSD-01Chuẩn hóaTên và nhãn về dạng hợp lệ, tiền tố ah_TestNormalize, TestVMWriteImportPayload
FR-TSD-02Ép nhãn tenantLoại nhãn dành riêng của agent; company_id, server_id đặt sau cùngTestCleanLabelsDropsReserved, TestCheckTenantAndPermanent
FR-TSD-03Từ chối thiếu tenant trước khi gọi mạngKhông gửi request khi thiếu company_id hoặc server_idTestVMWriteRejectsMissingTenantWithoutRequest
FR-TSD-04Loại giá trị không hữu hạnNaN, Inf bị bỏTestVMWriteDropsNonFinite
FR-TSD-05Phân loại lỗi5xx, 408, 429 tạm thời; còn lại vĩnh viễnTestVMWriteErrorClassification
FR-TSD-06Truy vấn khoảng chọn rollup theo độ dàiThô, 5 phút, 1 giờTestVMRangePicksRollupBySpan, TestVMRangeSelectorsAndParsing, TestVMRangeErrors
FR-TSD-07Giá trị mới nhấtlast_over_time trong 10 phútTestVMLatest, TestMemoryLatestAndDelete
FR-TSD-08Xóa theo tenantdelete_series với bộ chọn có company_idTestVMDelete, TestMemoryLatestAndDelete
FR-TSD-09Kiểm tra sốngPing gọi /healthTestVMPing

2.2 Non-Functional Requirements ​

Allocated:

NFR L2TargetCách đáp ứng ở L3
NFR-01 (thông lượng)Gom lô lớn ghi ít lần5000 mẫu hoặc 500 ms, NDJSON một request
NFR-04 (tách nhận và ghi)Ingest không chờ TSDBBus 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:

IDTargetParent L2-NFRSatisfied-by
NFR-BUS-01Không mất lô chưa ACK khi consumer chếtNFR-13XAUTOCLAIM mỗi claim_idle/3 (10 giây với mặc định 30 giây), TestRedisBusClaimsEntriesOfDeadConsumer
NFR-BUS-02Không tăng vô hạnNFR-04Trần mềm theo XLEN lấy mẫu mỗi 200 ms, queue_per_shard
NFR-WRK-01Bộ nhớ worker có trầnNFR-01Áp lực ngược ở 4 lần max_batch_samples
NFR-TSD-01Dải truy vấn tối đa 11.000 điểm, bước tối thiểu 30 giâyNFR-01Bộ 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 ​

ACGiven / When / ThenTest ID
AC-BWT-01Given hai consumer cùng nhóm, When có lô, Then mỗi lô xử lý đúng một lầnTestRedisBusTwoConsumersSplitTheWork
AC-BWT-02Given consumer chết giữa chừng, When qua claim_idle, Then consumer khác nhận lạiTestRedisBusClaimsEntriesOfDeadConsumer
AC-BWT-03Given TSDB lỗi tạm thời, When ghi, Then giữ lô và thử lạiTestTemporaryErrorRetriesAndKeepsBatch
AC-BWT-04Given TSDB trả 4xx (trừ 408, 429), When ghi, Then bỏ lôTestPermanentErrorDrops, TestVMWriteErrorClassification
AC-BWT-05Given tắt dịch vụ, When còn mẫu, Then ghi nốt trước hạn 10 giâyTestShutdownDrainsBusAndFlushes
AC-BWT-06Given TSDB còn lỗi khi tắt, When quá hạn, Then bỏ cuộcTestFinalFlushGivesUpAtDeadline
AC-BWT-07Given hai tenant, When ghi và truy vấn, Then không thấy dữ liệu của nhauTestIngestToTSDBPipelineKeepsTenantsApart, TestMemoryDedupAndIsolation
AC-BWT-08Given lô thiếu company_id, When ghi, Then từ chối mà không gọi mạngTestVMWriteRejectsMissingTenantWithoutRequest

2.4 Quality Attribute Scenarios ​

NguồnKích thíchMôi trườngPhản hồiThước đo
AgentGửi liên tụcVictoriaMetrics chậmWorker đạt ngưỡng áp lực ngược, bus đầy dần, Ingest trả 503 kèm Retry-AfterKhông tăng bộ nhớ vô hạn
Vận hànhDừng một node workerBình thườngNode khác nhận lại lô treoSau khoảng claim_idle cộng chu kỳ quét (đề xuất kiểm chứng)
Vận hànhKhởi động lại collectorCòn mẫu chưa ghiXả nốt, hạn 10 giâyMẫu còn lại sau hạn bị mất (đo bằng ahc_tsdb_batches_dropped_total)
Access HubTruy vấn khoảng 60 ngàyBình thườngDùng rollup 1 giờ từ long_urlKhông quá 11.000 điểm

3. Application Architecture ​

3.1 Kiến trúc thời chạy (C&C) ​

mermaid
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"| VM

Mũ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ầnTrách nhiệmVòng đời
Bus inprocKênh Go theo shard, dùng khi một tiến trìnhTheo tiến trình, mất khi restart
Bus RedisStream 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 workerMột goroutine mỗi shard sở hữu: đọc, gom, ghiTheo tiến trình
BatcherGom mẫu theo (company_id, server_id), xả theo kích thước hoặc thời gianTheo shard worker
WriterHiện thực tsdb.Writer (VictoriaMetrics hoặc bộ nhớ)Theo tiến trình
Kết nốiKiểuChi tiết
Ingest đến BusBấ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 ​

mermaid
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 --> RX

Tệ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 ​

mermaid
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 --> Batch

Bất biến:

  • company_id và server_id củ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_.
  • ShardFor chỉ phụ thuộc server_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 ​

InterfacePhương thứcHành vi
BusPublish(batch)Không chặn, ErrFull, ErrClosed
BusSubscribe(shards, handler)Gọi handler theo thứ tự trong shard, chạy đến khi hủy rồi xả
BusDepth()Độ sâu tổng (nguồn của ahc_bus_depth và ahc_ingest_bus_depth)
WriterWrite(batch)Lỗi tạm thời hoặc vĩnh viễn (IsPermanent)
ReaderRange, LatestTruy vấn có company_id bắt buộc
DeleterDeleteServer, DeleteCompanyXóa chuỗi theo bộ chọn
PingKiểm tra sống (dùng cho /readyz)

VictoriaMetrics (gọi ra):

MethodPathDùng cho
POST/api/v1/importGhi NDJSON ({"metric":{...},"values":[...],"timestamps":[...]})
GET, POST/api/v1/query_rangeRange
GET, POST/api/v1/queryLatest (last_over_time trong 10 phút)
POST/api/v1/admin/tsdb/delete_seriesXóa dữ liệu
GET/healthPing

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ỗng

Bộ chọn truy vấn luôn chứa company_id="..." (không có đường bỏ qua).

5.3 Error codes ​

Tình huốngKết quả nội bộGhi chú
VM trả 5xx, 408, 429Lỗi tạm thờiWorker thử lại
VM trả mã khác (4xx)Lỗi vĩnh viễnWorker bỏ lô
Thiếu tenantLỗ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ỗiAdmin 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óaKiểuNội dung
bus:s:<n> (n từ 0 đến 63)streamMỗ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ầnQuy ước
Tên chuỗiah_<tên chuẩn hóa>
Nhãn bắt buộccompany_id, server_id
Nhãn khácNhãn từ agent đã lọc (loại nhãn dành riêng)
RollupChuỗi <tên>:5m_<agg> và <tên>:1h_<agg> do vmalert tạo (đường deploy/, ngoài phạm vi mã Go)
mermaid
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"| R1

Kế 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ệuPhân loạiLưu giữ
Stream busTrung bìnhXóa sau khi ACK; trần mềm theo queue_per_shard
Chuỗi thôTrung bìnhTheo cấu hình VictoriaMetrics (-retentionPeriod, xem deploy/), chưa ghi thành số trong bản nháp này
RollupTrung bìnhTheo 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) ​

mermaid
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à XDEL

7.2 Sequences (đường lỗi) ​

mermaid
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
    end

7.3 State machines ​

Worker có bốn trạng thái:

mermaid
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áiSequence
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ắtTắt êm (L2 mục 8)

7.4 Core algorithms ​

Vấn đềGiải phápTrade-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ểmMáy chủ nóng làm nóng một shard (không cân bằng theo tải)
Không lẫn tenantAdapter ép company_id, server_id cuối cùng, ghi đè bất kỳ nhãn cùng tênNhãn agent trùng tên bị mất (đã chủ đích)
Trần bus với RedisLấy mẫu XLEN mỗi 200 ms so queue_per_shardCó thể vượt trần nhẹ giữa hai lần lấy mẫu
Bus Redis chia việcNhó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 ghiPhân loại lỗi theo mã HTTPLỗi vĩnh viễn làm mất lô (không có dead-letter cho mẫu)
Chọn rollupChọn chuỗi theo độ dài khoảngRollup cần vmalert chạy; không có rollup thì truy vấn dài rỗng
Giá trị mới nhấtlast_over_time 10 phútAgent 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ỗiNguyên nhânCơ chế xử lýTrạng thái cuối
PublishBus đầyErrFull, Ingest trả 503Agent thử lại từ WAL
PublishBus đóngErrClosedIngest trả lỗi khi tắt
Đọc busRedis mất kết nốiThử lại vòng đọcLô vẫn trong stream
Giải mã mụcMục hỏngBỏ qua và ACK để không kẹt shardMất mục đó
Ghi VMLỗi tạm thờiGiữ lô, lùi 200 ms đến 30 sThử đến khi thành công hoặc tắt
Ghi VMLỗi vĩnh viễnBỏ lôahc_tsdb_batches_dropped_total tăng
TắtTSDB vẫn lỗiThử trong hạn 10 s rồi bỏMất mẫu còn lại
Truy vấnVM lỗiTrả lỗi cho AdminAdmin trả 503

8.2 Fail-fast ​

  • tsdb.url rỗng khi vai trò cần TSDB: cảnh báo (TestTSDBWarningOnlyForRolesThatUseTheStore); readiness kiểm tra "tsdb".
  • Tách ingest và worker bắt buộc Redis (TestSplitIngestWorkerRejectsLocalBackends, TestSplitRolesNeedRedis).
  • TestNewRedisRequiresClients: bus Redis cần client.

8.3 Race conditions ​

Tình huốngXử lý
Hai consumer nhận cùng mục sau XAUTOCLAIMAt-least-once; ghi lặp cùng mẫu nhận diện được ở VictoriaMetrics (khử trùng)
Publish khi đóngErrClosed
Đóng khi còn mụcXả 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ỗiHành viKiểu suy thoáiHệ quả
VictoriaMetricsWorker thử lại, áp lực ngược, bus đầy, Ingest 503Giảm chức năngAgent đệm WAL đến 50 MiB, 24 giờ (phía agent)
Redis (bus)Ingest không publish đượcTừ chốiIngest 503; readiness "redis" đỏ
Một worker chếtNhận lại sau claim_idleTrễGhi trễ vài chục giây (đề xuất kiểm chứng)
Cả haiAgent đệm WALMất kết nốiPhụ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 groupChia shard giữa consumerKhô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)
XAUTOCLAIMNhận việc treoGhi lặp
Goroutine mỗi shardThứ tự trong shardShard nóng
Áp lực ngượcNgừng đọcTăng độ trễ ghi
Trần mềm XLENChặn publishVượt nhẹ giữa hai lần lấy mẫu

11. Security ​

11.1 Ba lớp ​

LớpBiện pháp
Truyền thôngRedis 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ệuSố 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 L2Biệ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ừ registryTenant lấy từ bản ghi registry, không từ payload
Xóa dữ liệuDeleteServer, 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 độngMục liên quan
bus.backendautoauto: Redis nếu có redis.addr, ngược lại inproc3.1
bus.queue_per_shard256Trần mềm mỗi shard (dự phòng của NewRedis là 1000, lệch cấu hình)7.4
bus.readers8Số reader3.1
bus.block500 msThời gian chờ mỗi lần XREADGROUP3.1
bus.claim_idle30 giâyNgưỡ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.shardstất cả 64Tập shard sở hữu của tiến trình3.1
worker.shutdown_flush10 giâyHạn xả khi tắt7.3
tsdb.urlrỗngĐịa chỉ VictoriaMetrics ghi và truy vấn gần5.1
tsdb.long_urlrỗngNguồn rollup 1 giờ6.1
tsdb.timeout10 giâyTimeout HTTP8.1
tsdb.max_batch_samples5000Trần mẫu mỗi lô7.1
tsdb.max_batch_wait500 msThời gian chờ gom tối đa7.1
tsdb.retry_max30 giâyTrầ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ênKiểuNhãnNgữ nghĩa
ahc_bus_depthgaugeĐộ sâu bus tổng (không nhãn shard)
ahc_ingest_bus_depthgaugeĐộ sâu nhìn từ Ingest
ahc_tsdb_samples_written_totalcounterMẫu ghi thành công
ahc_tsdb_write_errors_totalcounterkindLỗi ghi theo loại
ahc_tsdb_batches_dropped_totalcounterLô bị bỏ (lỗi vĩnh viễn hoặc hết hạn tắt)
ahc_worker_batches_totalcounterSố lô đã xử lý
ahc_tsdb_flush_secondshistogramThời gian xả
ahc_worker_pending_samplesgaugeMẫ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/healthzTiế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ạiKiểm thửVị trí
Đơn vị BusTestShardForIsStableAndInRange, TestPublishFullAndClosed, TestSubscribeOrderPerShardAndDrainOnCancel, TestSubscribeOnlyOwnedShards, TestBatchSamplesinternal/bus
Bus RedisTestRedisBusRoundTripKeepsOrderAndContent, TestRedisBusFullAndClosed, TestRedisBusDrainsOnCancel, TestRedisBusTwoConsumersSplitTheWork, TestRedisBusClaimsEntriesOfDeadConsumer, TestRedisBusSkipsUndecodableEntries, TestNewRedisRequiresClientsinternal/bus
WorkerTestFlushesAfterMaxWait, TestFlushesWhenFull, TestCompanyAndServerReachWriter, TestTemporaryErrorRetriesAndKeepsBatch, TestPermanentErrorDrops, TestShutdownDrainsBusAndFlushes, TestShutdownWhileTSDBDownRetriesInFinalFlush, TestFinalFlushGivesUpAtDeadline, TestBackpressureHoldsBusQueue, TestOnBatchHookAndEmptyBatch, TestShardSubsetOnlyConsumesItsShardsinternal/worker
TSDB chung và bộ nhớTestNormalize, TestCheckTenantAndPermanent, TestCleanLabelsDropsReserved, TestMemoryDedupAndIsolation, TestMemoryRangeAggregates, TestMemoryLatestAndDeleteinternal/tsdb
VictoriaMetrics (giả lập HTTP)TestVMWriteImportPayload, TestVMWriteRejectsMissingTenantWithoutRequest, TestVMWriteErrorClassification, TestVMWriteDropsNonFinite, TestVMRangeSelectorsAndParsing, TestVMRangePicksRollupBySpan, TestVMRangeErrors, TestVMLatest, TestVMDelete, TestVMPinginternal/tsdb
Tích hợpTestIngestToTSDBPipelineKeepsTenantsApart, TestWorkerFlushesOnShutdown, TestSplitProcessesShareStateThroughRedisinternal/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 ​

MilestoneNội dungPhụ thuộcĐóng góp nghiệm thu
COL-1Khung ứng dụng, cấu hìnhNền
COL-5Bus inproc, worker, ghi TSDB theo lô, adapter VictoriaMetrics (/api/v1/import), kho bộ nhớCOL-1AC-BWT-03..06
COL-7Truy vấn Latest, Range qua Admin API (kế hoạch rollup trong mã)COL-5AC-BWT-04, 07, 08
COL-RBus Redis Streams, tách vai tròCOL-5AC-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-RAC-BWT-08
COL-13 (chưa làm)Bus JetStream (bộ đệm 2 giờ, NFR-07), lease theo shard, phát lạiCOL-RNFR-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 ​

mermaid
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ỏiHành vi tạm thờiOwnerMã theo dõi
1Kết quả đo trên VictoriaMetrics thật (JSON import, khử trùng, thông lượng)Chưa đo; vmcheck opt-inChưa chỉ địnhOQ-BWT-1
2Xác thực tới VictoriaMetrics (nay không có)Mạng nội bộChưa chỉ địnhOQ-BWT-2
3RTO và RPO chính thức của đường ghiChưa định nghĩaChưa chỉ địnhOQ-BWT-3
4bus.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ỉ địnhOQ-BWT-4
5Thay JSON import bằng remote-write (ADR 0005 và 0015)Giữ JSON importChưa chỉ địnhOQ-BWT-5
6Bổ sung lease theo shard (COL-13, ADR 0009 và 0012)Chỉ dùng consumer groupChưa chỉ địnhOQ-BWT-6
7Nhãn shard cho ahc_bus_depth để thấy shard nóngKhông cóChưa chỉ địnhOQ-BWT-7
8Có cần dead-letter cho lô bị bỏ vì lỗi vĩnh viễnBỏ và đếmChưa chỉ địnhOQ-BWT-8

Appendix B. ADR nội bộ ​

Mã ADRQuyết địnhTrạng tháiDriver
ADR 0005Nhập liệu vào VictoriaMetrics bằng remote-writeLệch với mã (mã dùng /api/v1/import)Xem đề xuất ADR 0015
ADR 0009Bus theo shard và leaseLệ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ấtCOL-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ụcHồ sơTrạng thái điềnGiải trình
0 đến 4Bắt buộcĐã điền
5 API contractBắt buộcĐiền theo interface nội bộ và VM APIKhông có API HTTP mở ra ngoài
6 Physical dataBắt buộcĐiềnĐịnh dạng mục stream chưa trích
7.1 đến 7.4Bắt buộcĐã điền
8, 9, 10, 11Bắt buộcĐã điền
12Bắt buộcĐã điềnMặc định lấy từ Default()
13, 14, 15Bắt buộcĐã điềnKiểm thử chưa chạy trong bản nháp
Trang này có giúp được bạn không?
Sửa trang này

Nội dung đồng bộ từ kho mã access-hub-collector lúc 10:57, 03/10/2026. Khi tài liệu và mã khác nhau, mã thắng.