From 5e31d2f252ae574c36cd4e9d032264473c655078 Mon Sep 17 00:00:00 2001 From: Cody Lee Date: Thu, 27 Aug 2026 15:49:59 -0500 Subject: [PATCH 1/2] Add InfluxDB v3 support and fix tag/field schema collisions. Introduce explicit version selection with the influxdb3-go client alongside existing v1 and v2 paths, and resolve overlapping tag/field keys required for InfluxDB 3 write validation. Co-authored-by: Cursor --- go.mod | 12 +- go.sum | 43 ++++++ pkg/influxunifi/README.md | 27 +++- pkg/influxunifi/clients.go | 2 +- pkg/influxunifi/influxdb.go | 126 +++++++++++++++--- pkg/influxunifi/integration_test.go | 2 +- .../integration_test_expectations.yaml | 12 +- pkg/influxunifi/report.go | 24 +++- pkg/influxunifi/schema_test.go | 34 +++++ pkg/influxunifi/site.go | 1 - pkg/influxunifi/uap.go | 3 +- pkg/influxunifi/ubb.go | 2 - pkg/influxunifi/uci.go | 2 - pkg/influxunifi/udm.go | 2 - pkg/influxunifi/usg.go | 1 - pkg/influxunifi/uxg.go | 2 - 16 files changed, 242 insertions(+), 53 deletions(-) create mode 100644 pkg/influxunifi/schema_test.go diff --git a/go.mod b/go.mod index 976a61e0..3cd3ca58 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,7 @@ go 1.25.5 require ( github.com/DataDog/datadog-go/v5 v5.9.1 + github.com/InfluxCommunity/influxdb3-go/v2 v2.17.0 github.com/flaticols/countrycodes v0.0.2 github.com/gorilla/mux v1.8.1 github.com/influxdata/influxdb1-client v0.0.0-20220302092344-a9ab5670611c @@ -12,7 +13,7 @@ require ( github.com/prometheus/common v0.70.1 github.com/spf13/pflag v1.0.10 github.com/stretchr/testify v1.12.1 - github.com/unpoller/unifi/v6 v6.0.0 + github.com/unpoller/unifi/v6 v6.0.0 go.opentelemetry.io/otel v1.45.0 go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetricgrpc v1.45.0 go.opentelemetry.io/otel/exporters/otlp/otlpmetric/otlpmetrichttp v1.45.0 @@ -29,14 +30,23 @@ require ( ) require ( + github.com/apache/arrow-go/v18 v18.6.0 // indirect github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/go-logr/logr v1.4.4 // indirect github.com/go-logr/stdr v1.2.2 // indirect + github.com/goccy/go-json v0.10.6 // indirect + github.com/google/flatbuffers v25.12.19+incompatible // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect + github.com/influxdata/line-protocol/v2 v2.2.1 // indirect + github.com/klauspost/compress v1.19.1 // indirect + github.com/klauspost/cpuid/v2 v2.3.0 // indirect + github.com/pierrec/lz4/v4 v4.1.26 // indirect + github.com/zeebo/xxh3 v1.1.0 // indirect go.opentelemetry.io/auto/sdk v1.2.1 // indirect go.opentelemetry.io/otel/trace v1.45.0 // indirect go.opentelemetry.io/proto/otlp v1.11.0 // indirect go.yaml.in/yaml/v3 v3.0.5 // indirect + golang.org/x/exp v0.0.0-20260527015227-08cc5374adb3 // indirect golang.org/x/net v0.58.0 // indirect golang.org/x/text v0.41.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260803160001-6ac0973c030d // indirect diff --git a/go.sum b/go.sum index c5c31e55..5954281e 100644 --- a/go.sum +++ b/go.sum @@ -2,10 +2,18 @@ github.com/BurntSushi/toml v1.5.0 h1:W5quZX/G/csjUnuI8SUYlsHs9M38FC7znL0lIO+DvMg github.com/BurntSushi/toml v1.5.0/go.mod h1:ukJfTF/6rtPPRCnwkur4qwRxa8vTRFBF0uk2lLoLwho= github.com/DataDog/datadog-go/v5 v5.9.1 h1:jOxw/TaxGWok8RIxbpqn2p3RzSnQr/m3Q6TgaHqqOU0= github.com/DataDog/datadog-go/v5 v5.9.1/go.mod h1:2SBt8zJu6r7sRQHZFMQ8oCukWTKj0ymwulmNgQzJ1JM= +github.com/InfluxCommunity/influxdb3-go/v2 v2.17.0 h1:mmaojtPoRZ6QmOXEHt5mMmfYupJ2VPUUYnUEyZf9XX0= +github.com/InfluxCommunity/influxdb3-go/v2 v2.17.0/go.mod h1:SefM3sI/OT1Awe8f0iHc+bpEOEOV/ogCABi74HdtSGc= github.com/Microsoft/go-winio v0.5.0/go.mod h1:JPGBdM1cNvN/6ISo+n8V5iA4v8pBzdOpzfwIujj1a84= github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= github.com/RaveNoX/go-jsoncommentstrip v1.0.0/go.mod h1:78ihd09MekBnJnxpICcwzCMzGrKSKYe4AqU6PDYYpjk= +github.com/andybalholm/brotli v1.2.1 h1:R+f5xP285VArJDRgowrfb9DqL18yVK0gKAW/F+eTWro= +github.com/andybalholm/brotli v1.2.1/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY= +github.com/apache/arrow-go/v18 v18.6.0 h1:GX/Jyd3R7mCLiECAwY9FWbbaYblie2WXBSz4Sw8fNpM= +github.com/apache/arrow-go/v18 v18.6.0/go.mod h1:gm3MiPpY82fLYK5VKPB3WoJbsiLVDfT7flD5/vHReKw= +github.com/apache/thrift v0.22.0 h1:r7mTJdj51TMDe6RtcmNdQxgn9XcyfGDOzegMDRg47uc= +github.com/apache/thrift v0.22.0/go.mod h1:1e7J/O1Ae6ZQMTYdy9xa3w9k+XHWPfRvdPyJeynQ+/g= github.com/apapsch/go-jsonmerge/v2 v2.0.0 h1:axGnT1gRIfimI7gJifB699GoE/oq+F2MU7Dml6nw9rQ= github.com/apapsch/go-jsonmerge/v2 v2.0.0/go.mod h1:lvDnEdqiQrp0O42VQGgmlKpxL1AP2+08jFMw88y4klk= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= @@ -17,18 +25,29 @@ github.com/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1x github.com/cenkalti/backoff/v5 v5.0.3/go.mod h1:rkhZdG3JZukswDf7f0cwqPNk4K0sa+F97BxZthm/crw= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/flaticols/countrycodes v0.0.2 h1:vedxSqHwG3r7lwUK2bfGFWkVcFv7QuSCKFMkywI/rIE= github.com/flaticols/countrycodes v0.0.2/go.mod h1:HCwEez5Z+nf062EOWMPqEh1uLb5QdZSTQmTrq4avBOA= +github.com/frankban/quicktest v1.11.0/go.mod h1:K+q6oSqb0W0Ininfk863uOk1lMy69l/P6txr3mVT54s= +github.com/frankban/quicktest v1.11.2/go.mod h1:K+q6oSqb0W0Ininfk863uOk1lMy69l/P6txr3mVT54s= +github.com/frankban/quicktest v1.13.0 h1:yNZif1OkDfNoDfb9zZa9aXIpejNR4F23Wely0c+Qdqk= +github.com/frankban/quicktest v1.13.0/go.mod h1:qLE0fzW0VuyUAJgPU19zByoIr0HtCHN/r/VLSOOIySU= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= github.com/go-logr/logr v1.4.4 h1:tG4xh9yMsRCAiodLVTxyrkzSZ9+o0L1Kg/+cPVcbP/8= github.com/go-logr/logr v1.4.4/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/goccy/go-json v0.10.6 h1:p8HrPJzOakx/mn/bQtjgNjdTcN+/S6FcG2CTtQOrHVU= +github.com/goccy/go-json v0.10.6/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/google/flatbuffers v25.12.19+incompatible h1:haMV2JRRJCe1998HeW/p0X9UaMTK6SDo0ffLn2+DbLs= +github.com/google/flatbuffers v25.12.19+incompatible/go.mod h1:1AeVuKshWv4vARoZatz6mlQ0JxURH0Kv5+zNeJKJCa8= +github.com/google/go-cmp v0.5.2/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= +github.com/google/go-cmp v0.5.5/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= @@ -43,19 +62,34 @@ github.com/influxdata/influxdb1-client v0.0.0-20220302092344-a9ab5670611c h1:qSH github.com/influxdata/influxdb1-client v0.0.0-20220302092344-a9ab5670611c/go.mod h1:qj24IKcXYK6Iy9ceXlo3Tc+vtHo9lIhSX5JddghvEPo= github.com/influxdata/line-protocol v0.0.0-20210922203350-b1ad95c89adf h1:7JTmneyiNEwVBOHSjoMxiWAqB992atOeepeFYegn5RU= github.com/influxdata/line-protocol v0.0.0-20210922203350-b1ad95c89adf/go.mod h1:xaLFMmpvUxqXtVkUJfg9QmT88cDaCJ3ZKgdZ78oO8Qo= +github.com/influxdata/line-protocol-corpus v0.0.0-20210519164801-ca6fa5da0184/go.mod h1:03nmhxzZ7Xk2pdG+lmMd7mHDfeVOYFyhOgwO61qWU98= +github.com/influxdata/line-protocol-corpus v0.0.0-20210922080147-aa28ccfb8937 h1:MHJNQ+p99hFATQm6ORoLmpUCF7ovjwEFshs/NHzAbig= +github.com/influxdata/line-protocol-corpus v0.0.0-20210922080147-aa28ccfb8937/go.mod h1:BKR9c0uHSmRgM/se9JhFHtTT7JTO67X23MtKMHtZcpo= +github.com/influxdata/line-protocol/v2 v2.0.0-20210312151457-c52fdecb625a/go.mod h1:6+9Xt5Sq1rWx+glMgxhcg2c0DUaehK+5TDcPZ76GypY= +github.com/influxdata/line-protocol/v2 v2.1.0/go.mod h1:QKw43hdUBg3GTk2iC3iyCxksNj7PX9aUSeYOYE/ceHY= +github.com/influxdata/line-protocol/v2 v2.2.1 h1:EAPkqJ9Km4uAxtMRgUubJyqAr6zgWM0dznKMLRauQRE= +github.com/influxdata/line-protocol/v2 v2.2.1/go.mod h1:DmB3Cnh+3oxmG6LOBIxce4oaL4CPj3OmMPgvauXh+tM= github.com/juju/gnuflag v0.0.0-20171113085948-2ce1bb71843d/go.mod h1:2PavIy+JPciBPrBUjwbNvtwB6RQlve+hkpll6QSNmOE= github.com/klauspost/compress v1.19.1 h1:VsB4HPswih7mmZ8WleSFQ75c/Ui1M4trX5oAsJnhSlk= github.com/klauspost/compress v1.19.1/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= +github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= +github.com/kr/pretty v0.2.1/go.mod h1:ipq/a2n7PKx3OHsz4KJII5eveXtPO4qwEXGdVfWzfnI= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= +github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= github.com/kylelemons/godebug v1.1.0/go.mod h1:9/0rRGxNHcop5bhtWyNeEfOS8JIWk580+fNqagV/RAw= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= +github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno= github.com/oapi-codegen/runtime v1.1.1 h1:EXLHh0DXIJnWhdRPN2w4MXAzFyE4CskzhNLUmtpMYro= github.com/oapi-codegen/runtime v1.1.1/go.mod h1:SK9X900oXmPWilYR5/WKPzt3Kqxn/uS/+lbpREv+eCg= +github.com/pierrec/lz4/v4 v4.1.26 h1:GrpZw1gZttORinvzBdXPUXATeqlJjqUG/D87TKMnhjY= +github.com/pierrec/lz4/v4 v4.1.26/go.mod h1:EoQMVJgeeEOMsCqCzqFm2O0cJvljX2nGZjcRIPL34O4= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= @@ -88,6 +122,10 @@ github.com/stretchr/testify v1.12.1/go.mod h1:MDEgiDPPsNp5cuIrHPPCyornHKgEVbtFUm github.com/unpoller/unifi/v6 v6.0.0 h1:NVlAWgHpW1HsS9gRgu/RVHasPfoyrzAl/TB18qMu6wc= github.com/unpoller/unifi/v6 v6.0.0/go.mod h1:ad72+qBh3C4avgn2ZCwfMJlLuEKOMqQzKE9RvzE0Nfg= github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k= +github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= +github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0= +github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= +github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= go.opentelemetry.io/otel v1.45.0 h1:pdrWmLHofpubmArBv1LgFSv1Z0Ie/ppdZzu+kUN5EeU= @@ -118,6 +156,8 @@ golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACk golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.55.0 h1:+KWHjbgOaAQ66dh/YlkZKHlz9ZUlq61AFirAR9ntP8M= golang.org/x/crypto v0.55.0/go.mod h1:uq0V9dE/fzQuJtbnL+2EhWOE63vo164FY8xqEnV9xis= +golang.org/x/exp v0.0.0-20260527015227-08cc5374adb3 h1:VHEvKbpgPXcPXn40t9cDTGK3JZwMikIEyF/CTrFfu7k= +golang.org/x/exp v0.0.0-20260527015227-08cc5374adb3/go.mod h1:d2fgXJLVs4dYDHUk5lwMIfzRzSrWCfGZb0ZqeLa/Vcw= golang.org/x/mod v0.4.2/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= @@ -149,6 +189,7 @@ golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtn golang.org/x/tools v0.1.1/go.mod h1:o0xws9oXOQQZyjljx8fwUC0k7L1pTE6eaCbjGeHmOkk= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golift.io/cnfg v0.2.5 h1:NwhQ+REL9BSTiHYU4MKMawCEzvtjmhE8RlNiE7XroqE= golift.io/cnfg v0.2.5/go.mod h1:iMzXYjvZI7iZphzY75hkFR/VShYeuZznXiQtFsBOSCU= @@ -167,8 +208,10 @@ google.golang.org/grpc v1.83.0/go.mod h1:kDyl6SKsiHKt0uylY5gtn5cEjkrIOhQOGDgIc4J google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.0-20200615113413-eeeca48fe776/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/pkg/influxunifi/README.md b/pkg/influxunifi/README.md index afbfc9af..512906a0 100644 --- a/pkg/influxunifi/README.md +++ b/pkg/influxunifi/README.md @@ -2,18 +2,35 @@ Collects UniFi data from a UniFi controller using the API. -This is meant for InfluxDB users 1.8+ and 2.x series. +This supports InfluxDB 1.x, 2.x, and 3.x. ## Configuration -### InfluxDB 1.8+, 2.x +### InfluxDB 3.x -Note the use of `auth_token` to enable this mode. +Set `version: 3` and provide a token plus database name. For InfluxDB 3 Core/Enterprise, +leave `use_v2_api` unset (defaults to the native v3 write API). For InfluxDB Cloud +Serverless or Clustered, set `use_v2_api: true`. ```yaml influxdb: disable: false - # How often to poll UniFi and report to Datadog. + version: 3 + interval: "30s" + url: http://influxdb3:8181 + auth_token: somesecret + database: unifi + verify_ssl: false +``` + +### InfluxDB 1.8+, 2.x + +Note the use of `auth_token` to enable v2 mode when `version` is omitted. + +```yaml +influxdb: + disable: false + # How often to poll UniFi and report to InfluxDB. interval: "2m" # the influxdb url to post data url: http://somehost:1234 @@ -34,7 +51,7 @@ Note the lack of `auth_token` to enable this mode. ```yaml influxdb: disable: false - # How often to poll UniFi and report to Datadog. + # How often to poll UniFi and report to InfluxDB. interval: "2m" # the influxdb url to post data url: http://somehost:1234 diff --git a/pkg/influxunifi/clients.go b/pkg/influxunifi/clients.go index 807494b8..117125f5 100644 --- a/pkg/influxunifi/clients.go +++ b/pkg/influxunifi/clients.go @@ -30,7 +30,7 @@ func (u *InfluxUnifi) batchClient(r report, s *unifi.Client) { // nolint: funlen "is_wired": s.IsWired.Txt, "is_guest": s.IsGuest.Txt, "use_fixedip": s.UseFixedIP.Txt, - "channel": s.Channel.Txt, + "channel_name": s.Channel.Txt, "vlan": s.Vlan.Txt, "1x_identity": s.Identity1x, } diff --git a/pkg/influxunifi/influxdb.go b/pkg/influxunifi/influxdb.go index bf0a01b3..6017e1a8 100644 --- a/pkg/influxunifi/influxdb.go +++ b/pkg/influxunifi/influxdb.go @@ -6,12 +6,14 @@ import ( "context" "crypto/tls" "fmt" + "net/http" "net/url" "os" "strconv" "strings" "time" + influxdb3 "github.com/InfluxCommunity/influxdb3-go/v2/influxdb3" influx "github.com/influxdata/influxdb-client-go/v2" influxV1 "github.com/influxdata/influxdb1-client/v2" "github.com/unpoller/unifi/v6" @@ -31,6 +33,7 @@ const ( defaultInfluxOrg = "unifi" defaultInfluxBucket = "unifi" defaultInfluxURL = "http://127.0.0.1:8086" + defaultInfluxDBV3 = "unifi" ) // Config defines the data needed to store metrics in InfluxDB. @@ -53,6 +56,13 @@ type Config struct { // BatchSize controls the async batch size for v2 influxdb client mode BatchSize uint `json:"batch_size,omitempty" toml:"batch_size,omitempty" xml:"batch_size" yaml:"batch_size"` + // Version selects the InfluxDB API: 1, 2, or 3. When unset, auth_token enables v2, otherwise v1. + Version uint `json:"version,omitempty" toml:"version,omitempty" xml:"version" yaml:"version"` + // Database is the InfluxDB v3 database to write metrics to. + Database string `json:"database,omitempty" toml:"database,omitempty" xml:"database" yaml:"database"` + // UseV2API when true sends v3 writes through the v2-compatible API (InfluxDB Cloud). + UseV2API bool `json:"use_v2_api,omitempty" toml:"use_v2_api,omitempty" xml:"use_v2_api" yaml:"use_v2_api"` + // URL details which influxdb url to use to report metrics to. URL string `json:"url,omitempty" toml:"url,omitempty" xml:"url" yaml:"url"` // Disable when true will disable the influxdb output. @@ -76,8 +86,9 @@ type InfluxUnifi struct { Collector poller.Collect InfluxV1Client influxV1.Client InfluxV2Client influx.Client + InfluxV3Client *influxdb3.Client LastCheck time.Time - IsVersion2 bool + Version InfluxVersion *InfluxDB } @@ -105,14 +116,10 @@ func init() { // nolint: gochecknoinits func (u *InfluxUnifi) PollController() { interval := u.Interval.Round(time.Second) ticker := time.NewTicker(interval) - version := "1" + version := strconv.FormatUint(uint64(u.Version), 10) - if u.IsVersion2 { - version = "2" - } - - u.Logf("Poller->InfluxDB started, version: %s, interval: %v, dp: %v, db: %s, url: %s, bucket: %s, org: %s", - version, interval, u.DeadPorts, u.DB, u.URL, u.Bucket, u.Org) + u.Logf("Poller->InfluxDB started, version: %s, interval: %v, dp: %v, db: %s, url: %s, bucket: %s, org: %s, database: %s", + version, interval, u.DeadPorts, u.DB, u.URL, u.Bucket, u.Org, u.Database) for u.LastCheck = range ticker.C { u.Poll(interval) @@ -173,7 +180,15 @@ func (u *InfluxUnifi) DebugOutput() (bool, error) { return false, fmt.Errorf("invalid influx URL: %v", err) } - if u.IsVersion2 { + switch u.Version { + case InfluxV3: + client, err := u.newInfluxV3Client() + if err != nil { + return false, err + } + + u.InfluxV3Client = client + case InfluxV2: // we're a version 2 tlsConfig := &tls.Config{InsecureSkipVerify: !u.VerifySSL} // nolint: gosec serverOptions := influx.DefaultOptions().SetTLSConfig(tlsConfig).SetBatchSize(u.BatchSize) @@ -190,7 +205,7 @@ func (u *InfluxUnifi) DebugOutput() (bool, error) { if !ok { return false, fmt.Errorf("unsuccessful ping to influxdb2") } - } else { + default: u.InfluxV1Client, err = influxV1.NewHTTPClient(influxV1.HTTPConfig{ Addr: u.URL, Username: u.User, @@ -233,12 +248,20 @@ func (u *InfluxUnifi) Run(c poller.Collect) error { return err } - if u.IsVersion2 { + switch u.Version { + case InfluxV3: + client, err := u.newInfluxV3Client() + if err != nil { + return err + } + + u.InfluxV3Client = client + case InfluxV2: // we're a version 2 tlsConfig := &tls.Config{InsecureSkipVerify: !u.VerifySSL} // nolint: gosec serverOptions := influx.DefaultOptions().SetTLSConfig(tlsConfig).SetBatchSize(u.BatchSize) u.InfluxV2Client = influx.NewClientWithOptions(u.URL, u.AuthToken, serverOptions) - } else { + default: u.InfluxV1Client, err = influxV1.NewHTTPClient(influxV1.HTTPConfig{ Addr: u.URL, Username: u.User, @@ -268,9 +291,14 @@ func (u *InfluxUnifi) setConfigDefaults() { u.AuthToken = u.getPassFromFile(strings.TrimPrefix(u.AuthToken, "file://")) } - if u.AuthToken != "" { - // Version >= 1.8 influx - u.IsVersion2 = true + u.Version = u.resolveInfluxVersion() + + switch u.Version { + case InfluxV3: + if u.Database == "" { + u.Database = defaultInfluxDBV3 + } + case InfluxV2: if u.Org == "" { u.Org = defaultInfluxOrg } @@ -282,8 +310,7 @@ func (u *InfluxUnifi) setConfigDefaults() { if u.BatchSize == 0 { u.BatchSize = 20 } - } else { - // Version < 1.8 influx + default: if u.User == "" { u.User = defaultInfluxUser } @@ -310,6 +337,45 @@ func (u *InfluxUnifi) setConfigDefaults() { u.Interval = cnfg.Duration{Duration: u.Interval.Round(time.Second)} } +func (u *InfluxUnifi) resolveInfluxVersion() InfluxVersion { + switch u.Config.Version { + case 3: + return InfluxV3 + case 2: + return InfluxV2 + case 1: + return InfluxV1 + default: + if u.AuthToken != "" { + return InfluxV2 + } + + return InfluxV1 + } +} + +func (u *InfluxUnifi) newInfluxV3Client() (*influxdb3.Client, error) { + if u.AuthToken == "" { + return nil, fmt.Errorf("influxdb v3 requires auth_token") + } + + httpClient := &http.Client{ + Transport: &http.Transport{ + TLSClientConfig: &tls.Config{InsecureSkipVerify: !u.VerifySSL}, //nolint:gosec + }, + } + + return influxdb3.New(influxdb3.ClientConfig{ + Host: u.URL, + Token: u.AuthToken, + Database: u.Database, + HTTPClient: httpClient, + WriteOptions: &influxdb3.WriteOptions{ + UseV2Api: u.UseV2API, + }, + }) +} + func (u *InfluxUnifi) getPassFromFile(filename string) string { b, err := os.ReadFile(filename) if err != nil { @@ -324,7 +390,7 @@ func (u *InfluxUnifi) getPassFromFile(filename string) string { // Returns an error if influxdb calls fail, otherwise returns a report. func (u *InfluxUnifi) ReportMetrics(m *poller.Metrics, e *poller.Events) (*Report, error) { r := &Report{ - UseV2: u.IsVersion2, + Version: u.Version, Metrics: m, Events: e, ch: make(chan *metric), @@ -333,7 +399,19 @@ func (u *InfluxUnifi) ReportMetrics(m *poller.Metrics, e *poller.Events) (*Repor } defer close(r.ch) - if u.IsVersion2 { + switch u.Version { + case InfluxV3: + go u.collect(r, r.ch) + u.loopPoints(r) + r.wg.Wait() + + ctx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + + if err := u.InfluxV3Client.WritePoints(ctx, r.v3); err != nil { + return nil, fmt.Errorf("influxdb3.WritePoints: %w", err) + } + case InfluxV2: // Make a new Influx Points Batcher. r.writer = u.InfluxV2Client.WriteAPI(u.Org, u.Bucket) @@ -344,7 +422,7 @@ func (u *InfluxUnifi) ReportMetrics(m *poller.Metrics, e *poller.Events) (*Repor // Flush all the points. r.writer.Flush() - } else { + default: var err error // Make a new Influx Points Batcher. @@ -378,10 +456,14 @@ func (u *InfluxUnifi) collect(r report, ch chan *metric) { tags := u.mergeGlobalTags(m.Tags) - if u.IsVersion2 { + switch u.Version { + case InfluxV3: + pt := influxdb3.NewPoint(m.Table, tags, m.Fields, m.TS) + r.batchV3(m, pt) + case InfluxV2: pt := influx.NewPoint(m.Table, tags, m.Fields, m.TS) r.batchV2(m, pt) - } else { + default: pt, err := influxV1.NewPoint(m.Table, tags, m.Fields, m.TS) if err == nil { r.batchV1(m, pt) diff --git a/pkg/influxunifi/integration_test.go b/pkg/influxunifi/integration_test.go index 0a36c844..71d3d2e2 100644 --- a/pkg/influxunifi/integration_test.go +++ b/pkg/influxunifi/integration_test.go @@ -172,7 +172,7 @@ func TestInfluxV1Integration(t *testing.T) { u := influxunifi.InfluxUnifi{ Collector: testRig.Collector, - IsVersion2: false, + Version: influxunifi.InfluxV1, InfluxV1Client: mockCapture, InfluxDB: &influxunifi.InfluxDB{ Config: &influxunifi.Config{ diff --git a/pkg/influxunifi/integration_test_expectations.yaml b/pkg/influxunifi/integration_test_expectations.yaml index 3ade5f20..a9b5039d 100644 --- a/pkg/influxunifi/integration_test_expectations.yaml +++ b/pkg/influxunifi/integration_test_expectations.yaml @@ -17,7 +17,7 @@ points: clients: tags: - 1x_identity - - channel + - channel_name - dev_cat - dev_family - dev_id @@ -219,7 +219,6 @@ points: site_to_site_enabled: bool tx_bytes-r: int uptime: int - wan_ip: string xput_down: float xput_up: float uap: @@ -307,7 +306,7 @@ points: - source fields: ast_be_xmit: float - channel: float + channel_num: float cu_self_rx: float cu_self_tx: float cu_total: float @@ -320,7 +319,6 @@ points: min_txpower: float nss: float num_sta: float - radio: string radio_caps: float tx_packets: float tx_power: float @@ -434,7 +432,6 @@ points: p2p_throughput: float p2p_tx_rate: float rx_bytes: float - source: string stat_bytes: float stat_duration: float stat_mac_filter_rejections: float @@ -547,7 +544,6 @@ points: upgradeable: bool uptime: float user-num_sta: float - version: string uci: tags: - mac @@ -572,7 +568,6 @@ points: mem_total: float mem_used: float rx_bytes: float - source: string stat_bytes: float stat_rx_bytes: float stat_rx_crypts: float @@ -594,7 +589,6 @@ points: temp_sys: float tx_bytes: float uptime: float - version: string udb: tags: - mac @@ -826,7 +820,6 @@ points: num_handheld: float num_mobile: float rx_bytes: float - source: string speedtest-status_latency: float speedtest-status_ping: float speedtest-status_rundate: float @@ -856,7 +849,6 @@ points: uplink_uptime: float uptime: float user-num_sta: float - version: string usg_networks: tags: - device_name diff --git a/pkg/influxunifi/report.go b/pkg/influxunifi/report.go index 893c5309..f0ee5628 100644 --- a/pkg/influxunifi/report.go +++ b/pkg/influxunifi/report.go @@ -5,15 +5,25 @@ import ( "sync" "time" + influxdb3 "github.com/InfluxCommunity/influxdb3-go/v2/influxdb3" influxV2API "github.com/influxdata/influxdb-client-go/v2/api" influxV2Write "github.com/influxdata/influxdb-client-go/v2/api/write" influxV1 "github.com/influxdata/influxdb1-client/v2" "github.com/unpoller/unpoller/pkg/poller" ) +// InfluxVersion selects the InfluxDB client and write API. +type InfluxVersion uint8 + +const ( + InfluxV1 InfluxVersion = 1 + InfluxV2 InfluxVersion = 2 + InfluxV3 InfluxVersion = 3 +) + // Report is returned to the calling procedure after everything is processed. type Report struct { - UseV2 bool + Version InfluxVersion Metrics *poller.Metrics Events *poller.Events Errors []error @@ -24,6 +34,8 @@ type Report struct { wg sync.WaitGroup bp influxV1.BatchPoints writer influxV2API.WriteAPI + v3 []*influxdb3.Point + v3Mu sync.Mutex } // Counts holds counters and has a lock to deal with routines. @@ -40,6 +52,7 @@ type report interface { error(err error) batchV1(m *metric, pt *influxV1.Point) batchV2(m *metric, pt *influxV2Write.Point) + batchV3(m *metric, pt *influxdb3.Point) metrics() *poller.Metrics events() *poller.Events addCount(item, ...int) @@ -131,6 +144,15 @@ func (r *Report) batchV2(m *metric, p *influxV2Write.Point) { r.writer.WritePoint(p) } +func (r *Report) batchV3(m *metric, p *influxdb3.Point) { + r.addCount(pointT) + r.addCount(fieldT, len(m.Fields)) + r.addCount(bytesT, calculateMetricBytes(m)) + r.v3Mu.Lock() + r.v3 = append(r.v3, p) + r.v3Mu.Unlock() +} + func (r *Report) String() string { r.Counts.RLock() defer r.Counts.RUnlock() diff --git a/pkg/influxunifi/schema_test.go b/pkg/influxunifi/schema_test.go new file mode 100644 index 00000000..a0747e49 --- /dev/null +++ b/pkg/influxunifi/schema_test.go @@ -0,0 +1,34 @@ +package influxunifi_test + +import ( + "os" + "testing" + + "github.com/stretchr/testify/require" + "gopkg.in/yaml.v3" +) + +func TestInfluxSchemaNoTagFieldOverlap(t *testing.T) { + t.Helper() + + yamlFile, err := os.ReadFile("integration_test_expectations.yaml") + require.NoError(t, err) + + var data testExpectations + + err = yaml.Unmarshal(yamlFile, &data) + require.NoError(t, err) + + for measurement, spec := range data.Points { + tags := make(map[string]struct{}, len(spec.Tags)) + for _, tag := range spec.Tags { + tags[tag] = struct{}{} + } + + for field := range spec.Fields { + if _, ok := tags[field]; ok { + t.Errorf("measurement %q has overlapping tag/field key %q", measurement, field) + } + } + } +} diff --git a/pkg/influxunifi/site.go b/pkg/influxunifi/site.go index 796274ec..8bf1cbe0 100644 --- a/pkg/influxunifi/site.go +++ b/pkg/influxunifi/site.go @@ -38,7 +38,6 @@ func (u *InfluxUnifi) batchSite(r report, s *unifi.Site) { "num_disconnected": h.NumDisconnected.Val, "num_pending": h.NumPending.Val, "num_gw": h.NumGw.Val, - "wan_ip": h.WanIP, "num_sta": h.NumSta.Val, "gw_cpu": h.GwSystemStats.CPU.Val, "gw_mem": h.GwSystemStats.Mem.Val, diff --git a/pkg/influxunifi/uap.go b/pkg/influxunifi/uap.go index 08d5afdc..029941bf 100644 --- a/pkg/influxunifi/uap.go +++ b/pkg/influxunifi/uap.go @@ -211,7 +211,7 @@ func (u *InfluxUnifi) processRadTable(r report, t map[string]string, rt unifi.Ra for _, t := range rts { if strings.EqualFold(t.Name, p.Name) { fields["ast_be_xmit"] = t.AstBeXmit.Val - fields["channel"] = t.Channel.Val + fields["channel_num"] = t.Channel.Val fields["cu_self_rx"] = t.CuSelfRx.Val fields["cu_self_tx"] = t.CuSelfTx.Val fields["cu_total"] = t.CuTotal.Val @@ -219,7 +219,6 @@ func (u *InfluxUnifi) processRadTable(r report, t map[string]string, rt unifi.Ra fields["gain"] = t.Gain.Val fields["guest-num_sta"] = t.GuestNumSta.Val fields["num_sta"] = t.NumSta.Val - fields["radio"] = t.Radio fields["tx_packets"] = t.TxPackets.Val fields["tx_power"] = t.TxPower.Val fields["tx_retries"] = t.TxRetries.Val diff --git a/pkg/influxunifi/ubb.go b/pkg/influxunifi/ubb.go index de69c7d2..d6674b99 100644 --- a/pkg/influxunifi/ubb.go +++ b/pkg/influxunifi/ubb.go @@ -41,7 +41,6 @@ func (u *InfluxUnifi) batchUBB(r report, s *unifi.UBB) { // nolint: funlen u.batchSysStats(sysStats, systemStats), u.batchUBBstats(s.Stat), map[string]any{ - "source": s.SourceName, "ip": s.IP, "bytes": s.Bytes.Val, "last_seen": s.LastSeen.Val, @@ -51,7 +50,6 @@ func (u *InfluxUnifi) batchUBB(r report, s *unifi.UBB) { // nolint: funlen "uptime": s.Uptime.Val, "state": s.State.Val, "user-num_sta": s.UserNumSta.Val, - "version": s.Version, "uplink_speed": s.Uplink.Speed.Val, "uplink_max_speed": s.Uplink.MaxSpeed.Val, "uplink_latency": s.Uplink.Latency.Val, diff --git a/pkg/influxunifi/uci.go b/pkg/influxunifi/uci.go index 70560a04..cee6abe6 100644 --- a/pkg/influxunifi/uci.go +++ b/pkg/influxunifi/uci.go @@ -43,7 +43,6 @@ func (u *InfluxUnifi) batchUCI(r report, s *unifi.UCI) { // nolint: funlen fields := Combine( u.batchSysStats(sysStats, systemStats), map[string]any{ - "source": s.SourceName, "ip": s.IP, "bytes": s.Bytes.Val, "last_seen": s.LastSeen.Val, @@ -52,7 +51,6 @@ func (u *InfluxUnifi) batchUCI(r report, s *unifi.UCI) { // nolint: funlen "tx_bytes": s.TxBytes.Val, "uptime": s.Uptime.Val, "state": s.State.Val, - "version": s.Version, }, ) diff --git a/pkg/influxunifi/udm.go b/pkg/influxunifi/udm.go index 764d7bda..26d29cc0 100644 --- a/pkg/influxunifi/udm.go +++ b/pkg/influxunifi/udm.go @@ -101,7 +101,6 @@ func (u *InfluxUnifi) batchUDM(r report, s *unifi.UDM) { // nolint: funlen u.batchUSGstats(s.SpeedtestStatus, s.Stat.Gw, s.Uplink), u.batchSysStats(s.SysStats, s.SystemStats), map[string]any{ - "source": s.SourceName, "ip": s.IP, "bytes": s.Bytes.Val, "last_seen": s.LastSeen.Val, @@ -112,7 +111,6 @@ func (u *InfluxUnifi) batchUDM(r report, s *unifi.UDM) { // nolint: funlen "uptime": s.Uptime.Val, "state": s.State.Val, "user-num_sta": s.UserNumSta.Val, - "version": s.Version, "num_desktop": s.NumDesktop.Val, "num_handheld": s.NumHandheld.Val, "num_mobile": s.NumMobile.Val, diff --git a/pkg/influxunifi/usg.go b/pkg/influxunifi/usg.go index 9905ce6c..9f655ad5 100644 --- a/pkg/influxunifi/usg.go +++ b/pkg/influxunifi/usg.go @@ -39,7 +39,6 @@ func (u *InfluxUnifi) batchUSG(r report, s *unifi.USG) { "uptime": s.Uptime.Val, "state": s.State.Val, "user-num_sta": s.UserNumSta.Val, - "version": s.Version, "num_desktop": s.NumDesktop.Val, "num_handheld": s.NumHandheld.Val, "num_mobile": s.NumMobile.Val, diff --git a/pkg/influxunifi/uxg.go b/pkg/influxunifi/uxg.go index c4d81274..f9640e63 100644 --- a/pkg/influxunifi/uxg.go +++ b/pkg/influxunifi/uxg.go @@ -41,7 +41,6 @@ func (u *InfluxUnifi) batchUXG(r report, s *unifi.UXG) { // nolint: funlen u.batchUSGstats(s.SpeedtestStatus, gw, s.Uplink), u.batchSysStats(s.SysStats, s.SystemStats), map[string]any{ - "source": s.SourceName, "ip": s.IP, "bytes": s.Bytes.Val, "last_seen": s.LastSeen.Val, @@ -52,7 +51,6 @@ func (u *InfluxUnifi) batchUXG(r report, s *unifi.UXG) { // nolint: funlen "uptime": s.Uptime.Val, "state": s.State.Val, "user-num_sta": s.UserNumSta.Val, - "version": s.Version, "num_desktop": s.NumDesktop.Val, "num_handheld": s.NumHandheld.Val, "num_mobile": s.NumMobile.Val, From 9bc7f2c5bf884da3f74fec1102b0b52ed48b0dc4 Mon Sep 17 00:00:00 2001 From: Cody Lee Date: Thu, 27 Aug 2026 16:24:09 -0500 Subject: [PATCH 2/2] Complete InfluxDB v3 rollout: tests, docs, and docker example. Add v3 integration and version tests, migration notes, InfluxDB 3 docker-compose stack, and README updates to finish the remaining plan phases. Co-authored-by: Cursor --- init/docker/README.md | 9 ++ .../docker-compose-influxdb3.env.example | 16 +++ init/docker/docker-compose-influxdb3.yml | 52 +++++++++ pkg/influxunifi/MIGRATION.md | 42 +++++++ pkg/influxunifi/README.md | 2 + pkg/influxunifi/integration_v3_test.go | 106 ++++++++++++++++++ pkg/influxunifi/version_test.go | 84 ++++++++++++++ 7 files changed, 311 insertions(+) create mode 100644 init/docker/docker-compose-influxdb3.env.example create mode 100644 init/docker/docker-compose-influxdb3.yml create mode 100644 pkg/influxunifi/MIGRATION.md create mode 100644 pkg/influxunifi/integration_v3_test.go create mode 100644 pkg/influxunifi/version_test.go diff --git a/init/docker/README.md b/init/docker/README.md index c2c6ec86..9f534212 100644 --- a/init/docker/README.md +++ b/init/docker/README.md @@ -24,3 +24,12 @@ docker exec /usr/bin/unpoller --health The health check is automatically used by Docker and container orchestration platforms (Kubernetes, Docker Swarm, etc.) to determine container health status. + +## InfluxDB 3 Core + +Use `docker-compose-influxdb3.yml` with `docker-compose-influxdb3.env.example` to run UniFi Poller against InfluxDB 3 Core on port `8181`. Copy the env example, set `INFLUXDB_ADMIN_TOKEN` to a token created in InfluxDB 3, then start the stack: + +```bash +cp docker-compose-influxdb3.env.example docker-compose-influxdb3.env +docker compose -f docker-compose-influxdb3.yml --env-file docker-compose-influxdb3.env up -d +``` diff --git a/init/docker/docker-compose-influxdb3.env.example b/init/docker/docker-compose-influxdb3.env.example new file mode 100644 index 00000000..bfe75b0f --- /dev/null +++ b/init/docker/docker-compose-influxdb3.env.example @@ -0,0 +1,16 @@ +# InfluxDB 3 Core +# Create an admin token in InfluxDB 3 and set it here. +INFLUXDB_ADMIN_TOKEN=CHANGEME +INFLUXDB_DATABASE=unpoller + +# Grafana +GRAFANA_USERNAME=admin +GRAFANA_PASSWORD=grafanaadmin + +# UniFi Poller +POLLER_TAG=latest +POLLER_DEBUG=false +POLLER_SAVE_DPI=false +UNIFI_USER=unpoller +UNIFI_PASS=set_this_on_your_controller +UNIFI_URL=https://127.0.0.1:8443 diff --git a/init/docker/docker-compose-influxdb3.yml b/init/docker/docker-compose-influxdb3.yml new file mode 100644 index 00000000..aca12de9 --- /dev/null +++ b/init/docker/docker-compose-influxdb3.yml @@ -0,0 +1,52 @@ +# UniFi Poller with InfluxDB 3 Core. +# See README.md in this directory and pkg/influxunifi/README.md. +version: '3' +services: + influxdb3: + restart: always + image: influxdb:3-core + ports: + - '8181:8181' + command: + - influxdb3 + - serve + - --node-id=node0 + - --object-store=file + - --data-dir=/var/lib/influxdb3/data + - --plugin-dir=/var/lib/influxdb3/plugins + volumes: + - influxdb3-data:/var/lib/influxdb3/data + - influxdb3-plugins:/var/lib/influxdb3/plugins + grafana: + image: grafana/grafana:latest + restart: always + ports: + - '3000:3000' + volumes: + - grafana-storage:/var/lib/grafana + depends_on: + - influxdb3 + environment: + - GF_SECURITY_ADMIN_USER=${GRAFANA_USERNAME} + - GF_SECURITY_ADMIN_PASSWORD=${GRAFANA_PASSWORD} + - GF_INSTALL_PLUGINS=grafana-clock-panel,natel-discrete-panel,grafana-piechart-panel + unifi-poller: + restart: always + image: ghcr.io/unpoller/unpoller:${POLLER_TAG} + depends_on: + - grafana + - influxdb3 + environment: + - UP_INFLUXDB_VERSION=3 + - UP_INFLUXDB_URL=http://influxdb3:8181 + - UP_INFLUXDB_AUTH_TOKEN=${INFLUXDB_ADMIN_TOKEN} + - UP_INFLUXDB_DATABASE=${INFLUXDB_DATABASE} + - UP_UNIFI_DEFAULT_USER=${UNIFI_USER} + - UP_UNIFI_DEFAULT_PASS=${UNIFI_PASS} + - UP_UNIFI_DEFAULT_URL=${UNIFI_URL} + - UP_POLLER_DEBUG=${POLLER_DEBUG} + - UP_UNIFI_DEFAULT_SAVE_DPI=${POLLER_SAVE_DPI} +volumes: + influxdb3-data: + influxdb3-plugins: + grafana-storage: diff --git a/pkg/influxunifi/MIGRATION.md b/pkg/influxunifi/MIGRATION.md new file mode 100644 index 00000000..c4230b33 --- /dev/null +++ b/pkg/influxunifi/MIGRATION.md @@ -0,0 +1,42 @@ +# InfluxDB Output Migration Notes + +## InfluxDB 3 support + +UnPoller now supports InfluxDB 1.x, 2.x, and 3.x from the same `influxdb` output plugin. + +Enable v3 explicitly: + +```yaml +influxdb: + version: 3 + url: http://influxdb3:8181 + auth_token: your-token + database: unifi +``` + +For InfluxDB Cloud Serverless or Clustered, set `use_v2_api: true` so writes use the v2-compatible endpoint. + +See [README.md](README.md) and `init/docker/docker-compose-influxdb3.yml` for examples. + +## Schema changes (v1/v2/v3) + +InfluxDB 3 rejects line protocol where the same key is used as both a tag and a field on one point. The following keys were adjusted for compatibility: + +| Measurement | Change | +|-------------|--------| +| `subsystems` | Removed duplicate field `wan_ip` (remains a tag) | +| `clients` | Renamed tag `channel` to `channel_name` (numeric field `channel` unchanged) | +| `uap_radios` | Renamed field `channel` to `channel_num`; removed duplicate field `radio` (remains a tag) | +| `usg`, `ubb`, `uci`, `udm`, `uxg` | Removed duplicate fields `source` and/or `version` (remain tags) | + +### Grafana dashboards + +Import updated dashboards from the [unpoller/dashboards](https://github.com/unpoller/dashboards) repository (`v2.0.0` InfluxDB JSON files) if panels reference the old tag/field names. + +### Existing InfluxDB 3 databases + +If writes previously failed with errors such as `invalid column type for column 'wan_ip'`, drop and recreate the affected database or bucket, then allow UnPoller to recreate the schema with the corrected layout. + +## Issue reference + +Fixes [unpoller#1061](https://github.com/unpoller/unpoller/issues/1061). diff --git a/pkg/influxunifi/README.md b/pkg/influxunifi/README.md index 512906a0..4e50bf2a 100644 --- a/pkg/influxunifi/README.md +++ b/pkg/influxunifi/README.md @@ -23,6 +23,8 @@ influxdb: verify_ssl: false ``` +See [MIGRATION.md](MIGRATION.md) for schema changes affecting Grafana dashboards. + ### InfluxDB 1.8+, 2.x Note the use of `auth_token` to enable v2 mode when `version` is omitted. diff --git a/pkg/influxunifi/integration_v3_test.go b/pkg/influxunifi/integration_v3_test.go new file mode 100644 index 00000000..a99e8c93 --- /dev/null +++ b/pkg/influxunifi/integration_v3_test.go @@ -0,0 +1,106 @@ +package influxunifi_test + +import ( + "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + "time" + + influxdb3 "github.com/InfluxCommunity/influxdb3-go/v2/influxdb3" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/unpoller/unpoller/pkg/influxunifi" + "github.com/unpoller/unpoller/pkg/unittest" + "golift.io/cnfg" +) + +type v3WriteCapture struct { + mu sync.Mutex + batch []string +} + +func (c *v3WriteCapture) handler(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + http.NotFound(w, r) + + return + } + + if !strings.Contains(r.URL.Path, "write") { + http.NotFound(w, r) + + return + } + + body, err := io.ReadAll(r.Body) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + + return + } + + c.mu.Lock() + c.batch = append(c.batch, string(body)) + c.mu.Unlock() + + w.WriteHeader(http.StatusNoContent) +} + +func (c *v3WriteCapture) lines() []string { + c.mu.Lock() + defer c.mu.Unlock() + + raw := strings.Join(c.batch, "\n") + + return strings.Split(strings.TrimSpace(raw), "\n") +} + +func TestInfluxV3Integration(t *testing.T) { + capture := &v3WriteCapture{} + server := httptest.NewServer(http.HandlerFunc(capture.handler)) + t.Cleanup(server.Close) + + client, err := influxdb3.New(influxdb3.ClientConfig{ + Host: server.URL, + Token: "test-token", + Database: "unpoller", + WriteOptions: &influxdb3.WriteOptions{ + UseV2Api: false, + }, + }) + require.NoError(t, err) + t.Cleanup(func() { _ = client.Close() }) + + testRig := unittest.NewTestSetup(t) + t.Cleanup(testRig.Close) + + u := influxunifi.InfluxUnifi{ + Collector: testRig.Collector, + Version: influxunifi.InfluxV3, + InfluxV3Client: client, + InfluxDB: &influxunifi.InfluxDB{ + Config: &influxunifi.Config{ + Version: 3, + Database: "unpoller", + AuthToken: "test-token", + URL: server.URL, + Interval: cnfg.Duration{Duration: time.Hour}, + }, + }, + } + + testRig.Initialize() + u.Poll(time.Minute) + + lines := capture.lines() + require.NotEmpty(t, lines) + + body := strings.Join(lines, "\n") + assert.Contains(t, body, "subsystems,") + assert.Contains(t, body, "clients,") + assert.Contains(t, body, "channel_name=") + assert.NotRegexp(t, `(?m)^clients,[^ ]* channel=`, body) +} diff --git a/pkg/influxunifi/version_test.go b/pkg/influxunifi/version_test.go new file mode 100644 index 00000000..3e8d5045 --- /dev/null +++ b/pkg/influxunifi/version_test.go @@ -0,0 +1,84 @@ +package influxunifi_test + +import ( + "net/http" + "net/http/httptest" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/unpoller/unpoller/pkg/influxunifi" + "golift.io/cnfg" +) + +func TestInfluxVersionDefaults(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + config influxunifi.Config + version influxunifi.InfluxVersion + }{ + { + name: "v1 default without token", + config: influxunifi.Config{DB: "unifi"}, + version: influxunifi.InfluxV1, + }, + { + name: "v2 default with token", + config: influxunifi.Config{AuthToken: "secret"}, + version: influxunifi.InfluxV2, + }, + { + name: "explicit v3", + config: influxunifi.Config{Version: 3, AuthToken: "secret", Database: "unifi"}, + version: influxunifi.InfluxV3, + }, + { + name: "explicit v1 with token present", + config: influxunifi.Config{Version: 1, AuthToken: "ignored", DB: "unifi"}, + version: influxunifi.InfluxV1, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusNoContent) + })) + t.Cleanup(server.Close) + + u := &influxunifi.InfluxUnifi{ + InfluxDB: &influxunifi.InfluxDB{Config: &tc.config}, + } + u.Config.URL = server.URL + + ok, err := u.DebugOutput() + require.NoError(t, err) + assert.True(t, ok) + assert.Equal(t, tc.version, u.Version) + }) + } +} + +func TestInfluxV3RequiresAuthToken(t *testing.T) { + t.Parallel() + + u := &influxunifi.InfluxUnifi{ + InfluxDB: &influxunifi.InfluxDB{ + Config: &influxunifi.Config{ + Version: 3, + URL: "http://127.0.0.1:8181", + Database: "unifi", + Interval: cnfg.Duration{}, + }, + }, + } + + ok, err := u.DebugOutput() + require.Error(t, err) + assert.False(t, ok) + assert.Contains(t, err.Error(), "auth_token") +}