diff --git a/cmd/trading/main.go b/cmd/trading/main.go index 8fa4f46..e714798 100644 --- a/cmd/trading/main.go +++ b/cmd/trading/main.go @@ -10,6 +10,7 @@ import ( "sig-pub/pkg/grpc/discovery" "sig-pub/pkg/grpc/interceptor" "sig-pub/pkg/mq" + "sig-pub/pkg/storage/ck" "sig-pub/pkg/storage/persist" "sig-pub/pkg/utils/exit" "sig-pub/pkg/zlog" @@ -55,7 +56,19 @@ func main() { if err := rdb.Init(); err != nil { panic(err) } - tradingDataPersist := trading.NewTradingDataPersist(rdb) + // clickhouse + ckBatchWriter := ck.NewClickhouseBatchWriter(conf.Database.Clickhouse) + if err := ckBatchWriter.Init(); err != nil { + panic(err) + } + ckDB := ck.NewClickhouseDB(conf.Database.Clickhouse) + if err := ckDB.Init(); err != nil { + panic(err) + } + tradingDataPersist := trading.NewTradingDataPersist(rdb, ckDB, ckBatchWriter) + if err := tradingDataPersist.Init(); err != nil { + panic(err) + } // new market grpc client marketClient, err := client.NewMarketClient( diff --git a/config/config.toml b/config/config.toml index 70ebead..62c0d18 100644 --- a/config/config.toml +++ b/config/config.toml @@ -34,13 +34,14 @@ postgres = { DSN = "host=127.0.0.1 port=5432 user=postgres password=123456 dbnam [database.clickhouse] # dsn = "clickhouse://user:password@192.168.0.137:9000/game?dial_timeout=10s&read_timeout=20s" -addr = "127.0.0.1:7900" +addr = "127.0.0.1:9000" database = "sig" username = "root" password = "123456" dial_timeout = 10 read_timeout = 20 -logsql = false +logsql = true +logsqlCk = false [database.kvrocks] addr = "127.0.0.1:7666" diff --git a/docker-compose.yml b/docker-compose.yml index 085638c..9ba09b0 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,11 +1,14 @@ +# docker compose up -d sig-consul sig-kvrocks sig-postgres sig-vm sig-nats sig-clickhouse + services: sig-clickhouse: - image: 'clickhouse/clickhouse-server:24.12' + image: 'clickhouse/clickhouse-server:25.10' user: 'root' container_name: sig-clickhouse hostname: sig-clickhouse environment: - # - CLICKHOUSE_DB=sig + - CLICKHOUSE_RUN_AS_ROOT=1 + - CLICKHOUSE_DB=sig - CLICKHOUSE_USER=root - CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT=1 - CLICKHOUSE_PASSWORD=123456 @@ -16,8 +19,8 @@ services: # - ./fs/clickhouse/config.d/config.xml:/etc/clickhouse-server/config.d/config.xml # - ./fs/clickhouse/users.d/users.xml:/etc/clickhouse-server/users.d/users.xml ports: - - '7123:8123' - - '7900:9000' + - '8123:8123' + - '9000:9000' ulimits: nofile: soft: "262144" diff --git a/go.mod b/go.mod index 7916719..0a5e64f 100644 --- a/go.mod +++ b/go.mod @@ -1,6 +1,6 @@ module sig-pub -go 1.23.8 +go 1.24.0 toolchain go1.24.7 @@ -29,8 +29,8 @@ require ( go.etcd.io/etcd/api/v3 v3.6.1 go.etcd.io/etcd/client/v3 v3.6.1 go.uber.org/zap v1.27.0 - golang.org/x/net v0.38.0 - golang.org/x/sync v0.13.0 + golang.org/x/net v0.44.0 + golang.org/x/sync v0.17.0 golang.org/x/time v0.8.0 gonum.org/v1/gonum v0.16.0 google.golang.org/grpc v1.71.1 @@ -38,10 +38,13 @@ require ( gopkg.in/natefinch/lumberjack.v2 v2.2.1 gorm.io/driver/mysql v1.5.7 gorm.io/driver/postgres v1.6.0 - gorm.io/gorm v1.25.12 + gorm.io/gorm v1.30.0 ) require ( + github.com/ClickHouse/ch-go v0.68.0 // indirect + github.com/ClickHouse/clickhouse-go/v2 v2.40.3 // indirect + github.com/andybalholm/brotli v1.2.0 // indirect github.com/apapsch/go-jsonmerge/v2 v2.0.0 // indirect github.com/armon/go-metrics v0.4.1 // indirect github.com/bytedance/sonic/loader v0.2.4 // indirect @@ -54,6 +57,8 @@ require ( github.com/fsnotify/fsnotify v1.8.0 // indirect github.com/gabriel-vasile/mimetype v1.4.8 // indirect github.com/gin-contrib/sse v1.0.0 // indirect + github.com/go-faster/city v1.0.1 // indirect + github.com/go-faster/errors v0.7.1 // indirect github.com/go-playground/locales v0.14.1 // indirect github.com/go-playground/universal-translator v0.18.1 // indirect github.com/go-playground/validator/v10 v10.26.0 // indirect @@ -70,6 +75,7 @@ require ( github.com/hashicorp/go-immutable-radix v1.3.1 // indirect github.com/hashicorp/go-multierror v1.1.1 // indirect github.com/hashicorp/go-rootcerts v1.0.2 // indirect + github.com/hashicorp/go-version v1.7.0 // indirect github.com/hashicorp/golang-lru v0.5.4 // indirect github.com/hashicorp/serf v0.10.1 // indirect github.com/influxdata/line-protocol v0.0.0-20200327222509-2487e7298839 // indirect @@ -91,8 +97,12 @@ require ( github.com/nats-io/nkeys v0.4.11 // indirect github.com/nats-io/nuid v1.0.1 // indirect github.com/oapi-codegen/runtime v1.0.0 // indirect + github.com/paulmach/orb v0.11.1 // indirect github.com/pelletier/go-toml/v2 v2.2.3 // indirect + github.com/pierrec/lz4/v4 v4.1.22 // indirect github.com/sagikazarmark/locafero v0.7.0 // indirect + github.com/segmentio/asm v1.2.0 // indirect + github.com/shopspring/decimal v1.4.0 // indirect github.com/sourcegraph/conc v0.3.0 // indirect github.com/spf13/afero v1.12.0 // indirect github.com/spf13/pflag v1.0.6 // indirect @@ -102,13 +112,17 @@ require ( github.com/valyala/fastrand v1.1.0 // indirect github.com/valyala/histogram v1.2.0 // indirect go.etcd.io/etcd/client/pkg/v3 v3.6.1 // indirect + go.opentelemetry.io/otel v1.38.0 // indirect + go.opentelemetry.io/otel/trace v1.38.0 // indirect go.uber.org/multierr v1.11.0 // indirect + go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/arch v0.15.0 // indirect - golang.org/x/crypto v0.37.0 // indirect + golang.org/x/crypto v0.42.0 // indirect golang.org/x/exp v0.0.0-20250305212735-054e65f0b394 // indirect - golang.org/x/sys v0.32.0 // indirect - golang.org/x/text v0.24.0 // indirect + golang.org/x/sys v0.36.0 // indirect + golang.org/x/text v0.29.0 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20250303144028-a0af3efb3deb // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20250303144028-a0af3efb3deb // indirect gopkg.in/yaml.v3 v3.0.1 // indirect + gorm.io/driver/clickhouse v0.7.0 // indirect ) diff --git a/go.sum b/go.sum index eca0de2..15eb908 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,7 @@ +github.com/ClickHouse/ch-go v0.68.0 h1:zd2VD8l2aVYnXFRyhTyKCrxvhSz1AaY4wBUXu/f0GiU= +github.com/ClickHouse/ch-go v0.68.0/go.mod h1:C89Fsm7oyck9hr6rRo5gqqiVtaIY6AjdD0WFMyNRQ5s= +github.com/ClickHouse/clickhouse-go/v2 v2.40.3 h1:46jB4kKwVDUOnECpStKMVXxvR0Cg9zeV9vdbPjtn6po= +github.com/ClickHouse/clickhouse-go/v2 v2.40.3/go.mod h1:qO0HwvjCnTB4BPL/k6EE3l4d9f/uF+aoimAhJX70eKA= github.com/DataDog/datadog-go v3.2.0+incompatible/go.mod h1:LButxg5PwREeZtORoXG3tL4fMGNddJ+vMq1mwgfaqoQ= github.com/RaveNoX/go-jsoncommentstrip v1.0.0/go.mod h1:78ihd09MekBnJnxpICcwzCMzGrKSKYe4AqU6PDYYpjk= github.com/VictoriaMetrics/metrics v1.36.0 h1:f3SZMpLgIG4hJm2zfDs6wicxQ/QNWBZekY5rEGgbHKs= @@ -6,6 +10,8 @@ github.com/alecthomas/template v0.0.0-20160405071501-a0175ee3bccc/go.mod h1:LOuy github.com/alecthomas/template v0.0.0-20190718012654-fb15b899a751/go.mod h1:LOuyumcjzFXgccqObfd/Ljyb9UuFJ6TxHnclSeseNhc= github.com/alecthomas/units v0.0.0-20151022065526-2efee857e7cf/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= github.com/alecthomas/units v0.0.0-20190717042225-c3de453c63f4/go.mod h1:ybxpYRFXyAe+OPACYpWeL0wqObRcbAqCMya13uyzqw0= +github.com/andybalholm/brotli v1.2.0 h1:ukwgCxwYrmACq68yiUqwIWnGY0cTPox/M94sVwToPjQ= +github.com/andybalholm/brotli v1.2.0/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY= 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/armon/circbuf v0.0.0-20150827004946-bbbad097214e/go.mod h1:3U/XgcO3hCbHZ8TKRvWD2dDTCfh9M9ya+I9JpbB7O8o= @@ -69,12 +75,17 @@ github.com/gin-contrib/sse v1.0.0 h1:y3bT1mUWUxDpW4JLQg/HnTqV4rozuW4tC9eFKTxYI9E github.com/gin-contrib/sse v1.0.0/go.mod h1:zNuFdwarAygJBht0NTKiSi3jRf6RbqeILZ9Sp6Slhe0= github.com/gin-gonic/gin v1.10.0 h1:nTuyha1TYqgedzytsKYqna+DfLos46nTv2ygFy86HFU= github.com/gin-gonic/gin v1.10.0/go.mod h1:4PMNQiOhvDRa013RKVbsiNwoyezlm2rm0uX/T7kzp5Y= +github.com/go-faster/city v1.0.1 h1:4WAxSZ3V2Ws4QRDrscLEDcibJY8uf41H6AhXDrNDcGw= +github.com/go-faster/city v1.0.1/go.mod h1:jKcUJId49qdW3L1qKHH/3wPeUstCVpVSXTM6vO3VcTw= +github.com/go-faster/errors v0.7.1 h1:MkJTnDoEdi9pDabt1dpWf7AA8/BaSYZqibYyhZ20AYg= +github.com/go-faster/errors v0.7.1/go.mod h1:5ySTjWFiphBs07IKuiL69nxdfd5+fzh1u7FPGZP2quo= github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-logfmt/logfmt v0.3.0/go.mod h1:Qt1PoO58o5twSAckw1HlFXLmHsOX5/0LbT9GBnD5lWE= github.com/go-logfmt/logfmt v0.4.0/go.mod h1:3RMwSq7FuexP4Kalkev3ejPJsZTpXXBr9+V4qmtdjCk= github.com/go-logr/logr v1.4.2 h1:6pFjapn8bFcIbiKo3XT4j/BhANplGihG6tvd+8rYgrY= github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= 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/go-playground/assert/v2 v2.2.0 h1:JvknZsQTYeFEAhQwI4qEt9cyV5ONwRHC+lYKSsYSR8s= @@ -101,8 +112,10 @@ github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69 github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= github.com/golang/protobuf v1.3.1/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= github.com/golang/protobuf v1.3.2/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= +github.com/golang/protobuf v1.5.0/go.mod h1:FsONVRAS9T7sI+LIUmWTfcYkHO4aIWwzhcaSAoJOfIk= github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/golang/snappy v0.0.1/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= github.com/golang/snappy v0.0.4 h1:yAGX7huGHXlcLOEtBnF4w7FQwA26wojNCwOYAEhLjQM= github.com/golang/snappy v0.0.4/go.mod h1:/XxbfmMg8lxefKM7IXC3fBNl/7bRcc72aCRzEWrmP2Q= github.com/google/btree v0.0.0-20180813153112-4030bb1f1f0c/go.mod h1:lNA+9X1NB3Zf8V7Ke586lFgjr2dZNuvo3lPJSGZ5JPQ= @@ -111,6 +124,8 @@ github.com/google/btree v1.0.1/go.mod h1:xXMiIv4Fb/0kKde4SpL7qlzvu5cMJDRkFDxJfI9 github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= +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/gofuzz v1.0.0/go.mod h1:dBl0BpW6vV/+mYPU4Po3pmUjxk6FQPldtuIdl/M65Eg= @@ -157,6 +172,8 @@ github.com/hashicorp/go-uuid v1.0.3 h1:2gKiV6YVmrJ1i2CKKa9obLvRieoRGviZFL26PcT/C github.com/hashicorp/go-uuid v1.0.3/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro= github.com/hashicorp/go-version v1.2.1 h1:zEfKbn2+PDgroKdiOzqiE8rsmLqU2uwi5PB5pBJ3TkI= github.com/hashicorp/go-version v1.2.1/go.mod h1:fltr4n8CU8Ke44wwGCBoEymUuxUHl09ZGVZPK5anwXA= +github.com/hashicorp/go-version v1.7.0 h1:5tqGy27NaOTB8yJKUZELlFAS/LTKJkrmONwQKeRZfjY= +github.com/hashicorp/go-version v1.7.0/go.mod h1:fltr4n8CU8Ke44wwGCBoEymUuxUHl09ZGVZPK5anwXA= github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= github.com/hashicorp/golang-lru v0.5.4 h1:YDjusn29QI/Das2iO9M0BHnIbxPeyuCHsjMW+lJfyTc= github.com/hashicorp/golang-lru v0.5.4/go.mod h1:iADmTwqILo4mZ8BN3D2Q6+9jd8WM5uGBxy+E8yxSoD4= @@ -192,6 +209,7 @@ github.com/juju/gnuflag v0.0.0-20171113085948-2ce1bb71843d/go.mod h1:2PavIy+JPci github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7VTCxuUUipMqKk8s4w= github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= +github.com/klauspost/compress v1.13.6/go.mod h1:/3/Vjq9QcHkK5uEr5lBEmyoZ1iFhe47etQ6QUkpK6sk= github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= github.com/klauspost/cpuid/v2 v2.0.9/go.mod h1:FInQzS24/EEf25PyTYn52gqo7WaD8xa0213Md/qVLRg= @@ -245,6 +263,7 @@ github.com/modern-go/reflect2 v0.0.0-20180701023420-4b7aa43c6742/go.mod h1:bx2lN github.com/modern-go/reflect2 v1.0.1/go.mod h1:bx2lNnkwVCuqBIxFjflWJWanXIb3RllmbCylyMrvgv0= github.com/modern-go/reflect2 v1.0.2 h1:xBagoLtFs94CBntxluKeaWgTMpvLxC4ur3nMaC9Gz0M= github.com/modern-go/reflect2 v1.0.2/go.mod h1:yWuevngMOJpCy52FWWMvUC8ws7m/LJsjYzDa0/r8luk= +github.com/montanaflynn/stats v0.0.0-20171201202039-1bf9dbcd8cbe/go.mod h1:wL8QJuTMNUDYhXwkmfOly8iTdp5TEcJFWZD2D7SIkUc= github.com/mostynb/go-grpc-compression v1.2.3 h1:42/BKWMy0KEJGSdWvzqIyOZ95YcR9mLPqKctH7Uo//I= github.com/mostynb/go-grpc-compression v1.2.3/go.mod h1:AghIxF3P57umzqM9yz795+y1Vjs47Km/Y2FE6ouQ7Lg= github.com/mwitkow/go-conntrack v0.0.0-20161129095857-cc309e4a2223/go.mod h1:qRWi+5nqEBWmkhHvq77mSJWrCKwh8bxhgT7d/eI7P4U= @@ -259,8 +278,13 @@ github.com/oapi-codegen/runtime v1.0.0/go.mod h1:LmCUMQuPB4M/nLXilQXhHw+BLZdDb18 github.com/pascaldekloe/goe v0.0.0-20180627143212-57f6aae5913c/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= github.com/pascaldekloe/goe v0.1.0 h1:cBOtyMzM9HTpWjXfbbunk26uA6nG3a8n06Wieeh0MwY= github.com/pascaldekloe/goe v0.1.0/go.mod h1:lzWF7FIEvWOWxwDKqyGYQf6ZUaNfKdP144TG7ZOy1lc= +github.com/paulmach/orb v0.11.1 h1:3koVegMC4X/WeiXYz9iswopaTwMem53NzTJuTF20JzU= +github.com/paulmach/orb v0.11.1/go.mod h1:5mULz1xQfs3bmQm63QEJA6lNGujuRafwA5S/EnuLaLU= +github.com/paulmach/protoscan v0.2.1/go.mod h1:SpcSwydNLrxUGSDvXvO0P7g7AuhJ7lcKfDlhJCDw2gY= github.com/pelletier/go-toml/v2 v2.2.3 h1:YmeHyLY8mFWbdkNWwpr+qIL2bEqT0o95WSdkNHvL12M= github.com/pelletier/go-toml/v2 v2.2.3/go.mod h1:MfCQTFTvCcUyyvvwm1+G6H/jORL20Xlb6rzQu9GuUkc= +github.com/pierrec/lz4/v4 v4.1.22 h1:cKFw6uJDK+/gfw5BcDL0JL5aBsAFdsIT18eRtLj7VIU= +github.com/pierrec/lz4/v4 v4.1.22/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= github.com/pkg/errors v0.8.0/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= @@ -290,6 +314,10 @@ github.com/sagikazarmark/locafero v0.7.0 h1:5MqpDsTGNDhY8sGp0Aowyf0qKsPrhewaLSsF github.com/sagikazarmark/locafero v0.7.0/go.mod h1:2za3Cg5rMaTMoG/2Ulr9AwtFaIppKXTRYnozin4aB5k= github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529 h1:nn5Wsu0esKSJiIVhscUtVbo7ada43DJhG55ua/hjS5I= github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529/go.mod h1:DxrIzT+xaE7yg65j358z/aeFdxmN0P9QXhEzd20vsDc= +github.com/segmentio/asm v1.2.0 h1:9BQrFxC+YOHJlTlHGkTrFWf59nbL3XnCoFLTwDCI7ys= +github.com/segmentio/asm v1.2.0/go.mod h1:BqMnlJP91P8d+4ibuonYZw9mfnzI9HfxselHZr5aAcs= +github.com/shopspring/decimal v1.4.0 h1:bxl37RwXBklmTi0C79JfXCEBD1cqqHt0bbgBAGFp81k= +github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+DMd9qYNcwME= github.com/sirupsen/logrus v1.2.0/go.mod h1:LxeOpSwHxABJmUn/MG1IvRgCAasNZTLOkJPxbbu5VWo= github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE= github.com/sourcegraph/conc v0.3.0 h1:OQTbbt6P72L20UqAkXXuLOj79LfEanQ+YQFNpLA9ySo= @@ -312,6 +340,7 @@ github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/ github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= +github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.2/go.mod h1:R6va5+xMeoiuVRoj+gSkQ7d3FALtqAAGI1FQKckRals= @@ -320,8 +349,10 @@ github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo= github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA= github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/subosito/gotenv v1.6.0 h1:9NlTDc1FTs4qu0DDq7AEtTPNw6SVm7uBMsUCUjABIf8= github.com/subosito/gotenv v1.6.0/go.mod h1:Dk4QP5c2W3ibzajGcXpNraDfq2IrhjMIvMSWPKKo0FU= +github.com/tidwall/pretty v1.0.0/go.mod h1:XNkn88O1ChpSDQmQeStsy+sBenx6DDtFZJxhVysOjyk= github.com/tv42/httpunix v0.0.0-20150427012821-b75d8614f926/go.mod h1:9ESjWnEqriFuLhtthL60Sar/7RFoluCcXsuvEwTV5KM= github.com/twitchyliquid64/golang-asm v0.15.1 h1:SU5vSMR7hnwNxj24w34ZyCi/FmDZTkS4MhqMhdFk5YI= github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2+aY1QWCk3Cedj/Gdt08= @@ -331,6 +362,10 @@ github.com/valyala/fastrand v1.1.0 h1:f+5HkLW4rsgzdNoleUOB69hyT9IlD2ZQh9GyDMfb5G github.com/valyala/fastrand v1.1.0/go.mod h1:HWqCzkrkg6QXT8V2EXWvXCoow7vLwOFN002oeRzjapQ= github.com/valyala/histogram v1.2.0 h1:wyYGAZZt3CpwUiIb9AU/Zbllg1llXyrtApRS815OLoQ= github.com/valyala/histogram v1.2.0/go.mod h1:Hb4kBwb4UxsaNbbbh+RRz8ZR6pdodR57tzWUS3BUzXY= +github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= +github.com/xdg-go/scram v1.1.1/go.mod h1:RaEWvsqvNKKvBPvcKeFjrG2cJqOkHTiyTpzz23ni57g= +github.com/xdg-go/stringprep v1.0.3/go.mod h1:W3f5j4i+9rC0kuIEJL0ky1VpHXQU3ocBgklLGvcBnW8= +github.com/youmark/pkcs8 v0.0.0-20181117223130-1be2e3e5546d/go.mod h1:rHwXgn7JulP+udvsHwJoVG1YGAP6VLg4y9I5dyZdqmA= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= go.etcd.io/etcd/api/v3 v3.6.1 h1:yJ9WlDih9HT457QPuHt/TH/XtsdN2tubyxyQHSHPsEo= @@ -339,24 +374,32 @@ go.etcd.io/etcd/client/pkg/v3 v3.6.1 h1:CxDVv8ggphmamrXM4Of8aCC8QHzDM4tGcVr9p2BS go.etcd.io/etcd/client/pkg/v3 v3.6.1/go.mod h1:aTkCp+6ixcVTZmrJGa7/Mc5nMNs59PEgBbq+HCmWyMc= go.etcd.io/etcd/client/v3 v3.6.1 h1:KelkcizJGsskUXlsxjVrSmINvMMga0VWwFF0tSPGEP0= go.etcd.io/etcd/client/v3 v3.6.1/go.mod h1:fCbPUdjWNLfx1A6ATo9syUmFVxqHH9bCnPLBZmnLmMY= +go.mongodb.org/mongo-driver v1.11.4/go.mod h1:PTSz5yu21bkT/wXpkS7WR5f0ddqw5quethTUn9WM+2g= go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A= go.opentelemetry.io/otel v1.34.0 h1:zRLXxLCgL1WyKsPVrgbSdMN4c0FMkDAskSTQP+0hdUY= go.opentelemetry.io/otel v1.34.0/go.mod h1:OWFPOQ+h4G8xpyjgqo4SxJYdDQ/qmRH+wivy7zzx9oI= +go.opentelemetry.io/otel v1.38.0 h1:RkfdswUDRimDg0m2Az18RKOsnI8UDzppJAtj01/Ymk8= +go.opentelemetry.io/otel v1.38.0/go.mod h1:zcmtmQ1+YmQM9wrNsTGV/q/uyusom3P8RxwExxkZhjM= go.opentelemetry.io/otel/metric v1.34.0 h1:+eTR3U0MyfWjRDhmFMxe2SsW64QrZ84AOhvqS7Y+PoQ= go.opentelemetry.io/otel/metric v1.34.0/go.mod h1:CEDrp0fy2D0MvkXE+dPV7cMi8tWZwX3dmaIhwPOaqHE= +go.opentelemetry.io/otel/metric v1.38.0 h1:Kl6lzIYGAh5M159u9NgiRkmoMKjvbsKtYRwgfrA6WpA= go.opentelemetry.io/otel/sdk v1.34.0 h1:95zS4k/2GOy069d321O8jWgYsW3MzVV+KuSPKp7Wr1A= go.opentelemetry.io/otel/sdk v1.34.0/go.mod h1:0e/pNiaMAqaykJGKbi+tSjWfNNHMTxoC9qANsCzbyxU= go.opentelemetry.io/otel/sdk/metric v1.34.0 h1:5CeK9ujjbFVL5c1PhLuStg1wxA7vQv7ce1EK0Gyvahk= go.opentelemetry.io/otel/sdk/metric v1.34.0/go.mod h1:jQ/r8Ze28zRKoNRdkjCZxfs6YvBTG1+YIqyFVFYec5w= go.opentelemetry.io/otel/trace v1.34.0 h1:+ouXS2V8Rd4hp4580a8q23bg0azF2nI8cqLYnC8mh/k= go.opentelemetry.io/otel/trace v1.34.0/go.mod h1:Svm7lSjQD7kG7KJ/MUHPVXSDGz2OX4h0M2jHBhmSfRE= +go.opentelemetry.io/otel/trace v1.38.0 h1:Fxk5bKrDZJUH+AMyyIXGcFAPah0oRcT+LuNtJrmcNLE= +go.opentelemetry.io/otel/trace v1.38.0/go.mod h1:j1P9ivuFsTceSWe1oY+EeW3sc+Pp42sO++GHkg4wwhs= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= +go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= golang.org/x/arch v0.15.0 h1:QtOrQd0bTUnhNVNndMpLHNWrDmYzZ2KDqSrEymqInZw= golang.org/x/arch v0.15.0/go.mod h1:JmwW7aLIoRUKgaTzhkiEFxvcEiQGyOg9BMonBJUS7EE= golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4= @@ -364,8 +407,11 @@ golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACk golang.org/x/crypto v0.0.0-20190923035154-9ee001bba392/go.mod h1:/lpIB1dKB+9EgE3H3cr1v9wB50oz8l4C4h62xy7jSTY= golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= +golang.org/x/crypto v0.0.0-20220622213112-05595931fe9d/go.mod h1:IxCIyHEi3zRg3s0A5j5BB6A9Jmi73HwBIUl50j+osU4= golang.org/x/crypto v0.37.0 h1:kJNSjF/Xp7kU0iB2Z+9viTPMW4EqqsrywMXLJOOsXSE= golang.org/x/crypto v0.37.0/go.mod h1:vg+k43peMZ0pUMhYmVAWysMK35e6ioLh3wB8ZCAfbVc= +golang.org/x/crypto v0.42.0 h1:chiH31gIWm57EkTXpwnqf8qeuMUi0yekh6mT2AvFlqI= +golang.org/x/crypto v0.42.0/go.mod h1:4+rDnOTJhQCx2q7/j6rAN5XDw8kPjeaXEUR2eL94ix8= golang.org/x/exp v0.0.0-20250305212735-054e65f0b394 h1:nDVHiLt8aIbd/VzvPWN6kSOPE7+F/fNFDSXLVYkE/Iw= golang.org/x/exp v0.0.0-20250305212735-054e65f0b394/go.mod h1:sIifuuw/Yco/y6yb6+bDNfyeQ/MdPUy/hKEMYQV17cM= golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= @@ -379,8 +425,11 @@ golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLL golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20210410081132-afb366fc7cd1/go.mod h1:9tjilg8BloeKEkVJvy7fQ90B1CfIiPueXVOjqfkSzI8= +golang.org/x/net v0.0.0-20211112202133-69e39bad7dc2/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y= golang.org/x/net v0.38.0 h1:vRMAPTMaeGqVhG5QyLJHqNDwecKTomGeqbnfZyKlBI8= golang.org/x/net v0.38.0/go.mod h1:ivrbrMbzFq5J41QOQh0siUuly180yBYtLp+CKbEaFx8= +golang.org/x/net v0.44.0 h1:evd8IRDyfNBMBTTY5XRF1vaZlD+EmWx6x8PkhR04H/I= +golang.org/x/net v0.44.0/go.mod h1:ECOoLqd5U3Lhyeyo/QDCEVQ4sNgYsqvCZ722XogGieY= golang.org/x/sync v0.0.0-20181108010431-42b317875d0f/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20181221193216-37e7f081c4d4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= @@ -389,6 +438,8 @@ golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.13.0 h1:AauUjRAJ9OSnvULf/ARrrVywoJDy0YS2AwQ98I37610= golang.org/x/sync v0.13.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/sync v0.17.0 h1:l60nONMj9l5drqw6jlhIELNv9I0A4OFgRsG9k2oT9Ug= +golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= golang.org/x/sys v0.0.0-20180823144017-11551d06cbcc/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20181116152217-5ac8a444bdc5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -406,6 +457,8 @@ golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210303074136-134d130e1a04/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210630005230-0f9fa26af87c/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210927094055-39ccf1dd6fa6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220503163025-988cb79eb6c6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= @@ -414,13 +467,18 @@ golang.org/x/sys v0.0.0-20220811171246-fbc7d0a398ab/go.mod h1:oPkhp1MJrh7nUepCBc golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.32.0 h1:s77OFDvIQeibCmezSnk/q6iAfkdiQaJi4VzroCFrN20= golang.org/x/sys v0.32.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.36.0 h1:KVRy2GtZBrk1cBYA7MKu5bEZFxQk4NIDV6RLVcC8o0k= +golang.org/x/sys v0.36.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= +golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= golang.org/x/text v0.24.0 h1:dd5Bzh4yt5KYA8f9CJHCP4FB4D51c2c6JvN37xJJkJ0= golang.org/x/text v0.24.0/go.mod h1:L8rBsPeo2pSS+xqN0d5u2ikmjtmoJbDBT1b7nHvFCdU= +golang.org/x/text v0.29.0 h1:1neNs90w9YzJ9BocxfsQNHKuAT4pkghyXc4nhZ6sJvk= +golang.org/x/text v0.29.0/go.mod h1:7MhJOA9CD2qZyOKYazxdYMF85OwPdEr9jTtBpO7ydH4= golang.org/x/time v0.8.0 h1:9i3RxcPv3PZnitoVGMPDKZSq1xW1gK1Xy3ArNOGZfEg= golang.org/x/time v0.8.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= @@ -440,10 +498,13 @@ google.golang.org/genproto/googleapis/rpc v0.0.0-20250303144028-a0af3efb3deb h1: google.golang.org/genproto/googleapis/rpc v0.0.0-20250303144028-a0af3efb3deb/go.mod h1:LuRYeWDFV6WOn90g357N17oMCaxpgCnbi/44qJvDn2I= google.golang.org/grpc v1.71.1 h1:ffsFWr7ygTUscGPI0KKK6TLrGz0476KUvvsbqWK0rPI= google.golang.org/grpc v1.71.1/go.mod h1:H0GRtasmQOh9LkFoCPDu3ZrwUtD1YGE+b2vYBYd/8Ec= +google.golang.org/protobuf v1.26.0-rc.1/go.mod h1:jlhhOSvTdKEhbULTjvd4ARK9grFBp09yW+WbY/TyQbw= +google.golang.org/protobuf v1.27.1/go.mod h1:9q0QmTI4eRPtz6boOQmLYwt+qCgq0jsYwAQnmE0givc= google.golang.org/protobuf v1.36.6 h1:z1NpPI8ku2WgiWnf+t9wTPsn6eP1L7ksHUlkfLvd9xY= google.golang.org/protobuf v1.36.6/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY= gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20180628173108-788fd7840127/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/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= @@ -456,6 +517,8 @@ gopkg.in/yaml.v2 v2.2.5/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/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= +gorm.io/driver/clickhouse v0.7.0 h1:BCrqvgONayvZRgtuA6hdya+eAW5P2QVagV3OlEp1vtA= +gorm.io/driver/clickhouse v0.7.0/go.mod h1:TmNo0wcVTsD4BBObiRnCahUgHJHjBIwuRejHwYt3JRs= gorm.io/driver/mysql v1.5.7 h1:MndhOPYOfEp2rHKgkZIhJ16eVUIRf2HmzgoPmh7FCWo= gorm.io/driver/mysql v1.5.7/go.mod h1:sEtPWMiqiN1N1cMXoXmBbd8C6/l+TESwriotuRRpkDM= gorm.io/driver/postgres v1.6.0 h1:2dxzU8xJ+ivvqTRph34QX+WrRaJlmfyPqXmoGVjMBa4= @@ -463,4 +526,6 @@ gorm.io/driver/postgres v1.6.0/go.mod h1:vUw0mrGgrTK+uPHEhAdV4sfFELrByKVGnaVRkXD gorm.io/gorm v1.25.7/go.mod h1:hbnx/Oo0ChWMn1BIhpy1oYozzpM15i4YPuHDmfYtwg8= gorm.io/gorm v1.25.12 h1:I0u8i2hWQItBq1WfE0o2+WuL9+8L21K9e2HHSTE/0f8= gorm.io/gorm v1.25.12/go.mod h1:xh7N7RHfYlNc5EmcI/El95gXusucDrQnHXe0+CgWcLQ= +gorm.io/gorm v1.30.0 h1:qbT5aPv1UH8gI99OsRlvDToLxW5zR7FzS9acZDOZcgs= +gorm.io/gorm v1.30.0/go.mod h1:8Z33v652h4//uMA76KjeDH8mJXPm1QNCYrMeatR0DOE= nullprogram.com/x/optparse v1.0.0/go.mod h1:KdyPE+Igbe0jQUrVfMqDMeJQIJZEuyV7pjYmp6pbG50= diff --git a/internal/trading/backtest/account.go b/internal/trading/backtest/account.go index 821b429..be40f80 100644 --- a/internal/trading/backtest/account.go +++ b/internal/trading/backtest/account.go @@ -2,47 +2,155 @@ package backtest import ( "fmt" + "math" "sig-pub/pkg/trade" "sig-pub/pkg/types" "sig-pub/pkg/types/decimals" "sig-pub/pkg/utils/conver" "sig-pub/pkg/utils/lang" + "sort" "time" ) type BacktestAccount struct { trade.ITradeAccount cash float64 + initialCash float64 positions map[int64]*trade.Position - trades map[int64]*trade.Trade - closeTrades map[int64]*trade.Trade + trades map[int64]*Trade + closeTrades map[int64]*Trade simulator *TradeSimulator profit float64 winningTrades int losingTrades int fee float64 + maxDrawdown [4]float64 // 最大回撤 } func NewBacktestAccount(cash float64, simulator *TradeSimulator) *BacktestAccount { return &BacktestAccount{ + // keep initial cash for return calculations + initialCash: cash, cash: cash, positions: make(map[int64]*trade.Position), - trades: make(map[int64]*trade.Trade), - closeTrades: make(map[int64]*trade.Trade), + trades: make(map[int64]*Trade), + closeTrades: make(map[int64]*Trade), simulator: simulator, + maxDrawdown: [4]float64{cash, cash, cash, cash}, } } +// initialCash stores the starting capital for return calculations +// (placed here to avoid changing exported API) +func (a *BacktestAccount) InitialCash() float64 { return a.initialCash } + +// SharpeRatio computes an annualized Sharpe ratio based on closed trades. +// rfAnnual is the annual risk-free rate expressed as a decimal (e.g. 0.01 for 1%). +// Method: +// - For each closed trade, we compute a period return = trade.Pnl / initialCash. +// - Period lengths are derived from successive trade close timestamps (ms). +// - Excess returns = periodReturn - rfAnnual * periodYears. +// - Sharpe = mean(excess) / stddev(excess) * sqrt(periodsPerYear) +// This provides a reasonable approximation when equity snapshots are not available. +func (a *BacktestAccount) SharpeRatio(rfAnnual float64) float64 { + if a.initialCash <= 0 { + return 0 + } + n := len(a.closeTrades) + if n < 2 { + return 0 + } + + trades := make([]*Trade, 0, n) + for _, t := range a.closeTrades { + trades = append(trades, t) + } + sort.Slice(trades, func(i, j int) bool { return trades[i].CloseTime < trades[j].CloseTime }) + + // returns per closed trade (relative to initial capital) + returns := make([]float64, 0, n) + // periods in seconds between closes; length will be n-1 initially + periodsSec := make([]float64, 0, n-1) + for i, t := range trades { + returns = append(returns, t.Pnl/a.initialCash) + if i > 0 { + // CloseTs is in milliseconds in this codebase + dtSec := float64(t.CloseTime-trades[i-1].CloseTime) / 1000.0 + if dtSec <= 0 { + dtSec = 1.0 + } + periodsSec = append(periodsSec, dtSec) + } + } + if len(periodsSec) == 0 { + return 0 + } + // average period (seconds) used to approximate period length for the first return + sumDt := 0.0 + for _, d := range periodsSec { + sumDt += d + } + avgDt := sumDt / float64(len(periodsSec)) + + // Build final periods slice aligned with returns length + periods := make([]float64, 0, n) + periods = append(periods, avgDt) + periods = append(periods, periodsSec...) + + const secsYear = 365.0 * 24.0 * 3600.0 + excess := make([]float64, len(returns)) + for i := range returns { + years := periods[i] / secsYear + excess[i] = returns[i] - rfAnnual*years + } + + meanEx := mean(excess) + sd := stddev(excess) + if sd == 0 { + return 0 + } + // approximate number of periods per year + periodsPerYear := secsYear / avgDt + return meanEx / sd * math.Sqrt(periodsPerYear) +} + +func mean(x []float64) float64 { + if len(x) == 0 { + return 0 + } + s := 0.0 + for _, v := range x { + s += v + } + return s / float64(len(x)) +} + +func stddev(x []float64) float64 { + if len(x) <= 1 { + return 0 + } + m := mean(x) + s := 0.0 + for _, v := range x { + d := v - m + s += d * d + } + // population or sample? use sample (n-1) + return math.Sqrt(s / float64(len(x)-1)) +} + // CurrentEquity 根据当前价格对仓位进行 mark-to-market,返回账户净值 func (a *BacktestAccount) CurrentEquity(price float64) float64 { equity := a.cash for _, p := range a.positions { switch p.Side { case types.SideLong: - equity += (price - p.EntryPx) * p.Qty + // equity += (price - p.EntryPx) * p.Qty + equity += price * p.Qty case types.SideShort: - equity += (p.EntryPx - price) * p.Qty + // equity += (p.EntryPx - price) * p.Qty + equity += (p.EntryPx - price + p.EntryPx) * p.Qty } } return equity @@ -152,12 +260,26 @@ func (a *BacktestAccount) ClosePosition(pos *trade.Position, kline types.Kline, a.closeTrades[t.Id] = t if trade, ok := a.trades[pos.TradeId]; ok { - trade.Pnl = profit trade.ClosePrice = closePrice trade.CloseFee = t.Fee - trade.CloseTs = kline.Ts + trade.CloseTime = kline.Ts trade.CloseCause = cause - trade.HoldTime = conver.TimeDurationFormat(time.Duration(trade.CloseTs-trade.Time)*time.Millisecond, ".") + trade.Pnl = profit + trade.HoldTime = conver.TimeDurationFormat(time.Duration(trade.CloseTime-trade.Time)*time.Millisecond, ".") + trade.PeakPx = pos.PeakPx + } + + // 记录最大回撤 [high, low, high, low] + currentEquity := a.CurrentEquity(closePrice) + if currentEquity < a.maxDrawdown[1] { + a.maxDrawdown[1] = currentEquity + } + if currentEquity > a.maxDrawdown[0] { + if a.maxDrawdown[0]-a.maxDrawdown[1] > a.maxDrawdown[2]-a.maxDrawdown[3] { + a.maxDrawdown[2], a.maxDrawdown[3] = a.maxDrawdown[0], a.maxDrawdown[1] + } + a.maxDrawdown[0] = currentEquity + a.maxDrawdown[1] = currentEquity } return } diff --git a/internal/trading/backtest/trade_simulator.go b/internal/trading/backtest/trade_simulator.go index 4d2cfb4..2c9f11d 100644 --- a/internal/trading/backtest/trade_simulator.go +++ b/internal/trading/backtest/trade_simulator.go @@ -2,7 +2,6 @@ package backtest import ( "math" - "sig-pub/pkg/trade" "sig-pub/pkg/types" "sig-pub/pkg/types/decimals" ) @@ -19,7 +18,7 @@ func NewTradeSimulator(feePct, slippagePct float64) *TradeSimulator { } // ExecuteMarket 执行市价单,使用kline信息决定成交价(使用close以及滑点) -func (s *TradeSimulator) ExecuteMarket(side types.Side, qty float64, closePrice float64, ts int64) (trd *trade.Trade, ok bool) { +func (s *TradeSimulator) ExecuteMarket(side types.Side, qty float64, closePrice float64, ts int64) (trd *Trade, ok bool) { // base price use close base := closePrice slippage := s.SlippagePct @@ -34,7 +33,7 @@ func (s *TradeSimulator) ExecuteMarket(side types.Side, qty float64, closePrice return } fee := math.Abs(base*qty) * s.FeePct - trd = &trade.Trade{Side: side, Qty: qty, Price: base, Fee: fee, Time: ts} + trd = &Trade{Side: side, Qty: qty, Price: base, Fee: fee, Time: ts} s.TradeId++ trd.Id = s.TradeId ok = true @@ -42,7 +41,7 @@ func (s *TradeSimulator) ExecuteMarket(side types.Side, qty float64, closePrice } // ExecuteLimit 简单实现: 如果limit价格被kline的high/low包含则成交 -func (s *TradeSimulator) ExecuteLimit(side types.Side, qty float64, limitPx float64, k types.Kline, ts int64) (trd *trade.Trade, filled bool) { +func (s *TradeSimulator) ExecuteLimit(side types.Side, qty float64, limitPx float64, k types.Kline, ts int64) (trd *Trade, filled bool) { h := decimals.MustToFloat64(k.High) l := decimals.MustToFloat64(k.Low) switch side { @@ -52,14 +51,14 @@ func (s *TradeSimulator) ExecuteLimit(side types.Side, qty float64, limitPx floa // assume filled at min(limitPx, open) px := math.Min(limitPx, decimals.MustToFloat64(k.Open)) fee := math.Abs(px*qty) * s.FeePct - trd = &trade.Trade{Side: side, Qty: qty, Price: px * (1 + s.SlippagePct), Fee: fee, Time: ts} + trd = &Trade{Side: side, Qty: qty, Price: px * (1 + s.SlippagePct), Fee: fee, Time: ts} return trd, true } case types.SideShort: if h >= limitPx { px := math.Max(limitPx, decimals.MustToFloat64(k.Open)) fee := math.Abs(px*qty) * s.FeePct - trd = &trade.Trade{Side: side, Qty: qty, Price: px * (1 - s.SlippagePct), Fee: fee, Time: ts} + trd = &Trade{Side: side, Qty: qty, Price: px * (1 - s.SlippagePct), Fee: fee, Time: ts} return trd, true } } diff --git a/internal/trading/backtest/trading_plan_backtester.go b/internal/trading/backtest/trading_plan_backtester.go index d109920..503f155 100644 --- a/internal/trading/backtest/trading_plan_backtester.go +++ b/internal/trading/backtest/trading_plan_backtester.go @@ -5,6 +5,7 @@ import ( "errors" "fmt" "sig-pub/api/pb" + "sig-pub/internal/trading/sig" "sig-pub/pkg/data/entity" "sig-pub/pkg/indicator" "sig-pub/pkg/strategy" @@ -12,15 +13,18 @@ import ( "sig-pub/pkg/types" "sig-pub/pkg/utils/collect" "sig-pub/pkg/zlog" + "time" "github.com/bytedance/sonic" ) +// todo 将backtest独立成单独服务横向扩展 type TradingPlanBacktester struct { indicatorReg *indicator.IndicatorRegistry sigStrategyReg *strategy.SigStrategyRegistry exchangeClient pb.ExchangeServiceClient + plan entity.TradePlan account *BacktestAccount sigStrategyType strategy.SigStrategyType sigStrategy strategy.ISigStrategy @@ -38,6 +42,7 @@ func NewTradingPlanBacktester(indicatorReg *indicator.IndicatorRegistry, sigStra } func (b *TradingPlanBacktester) Init(cash float64, plan entity.TradePlan) (err error) { + b.plan = plan // trade account simulator := NewTradeSimulator(0.0005, 0.0008) b.account = NewBacktestAccount(cash, simulator) @@ -85,19 +90,30 @@ func (b *TradingPlanBacktester) Init(cash float64, plan entity.TradePlan) (err e } // 核心引擎,模拟交易、持仓跟踪、费用计算 -func (b *TradingPlanBacktester) Backtest(ctx context.Context, sr *pb.SeriesRange) (err error) { +func (b *TradingPlanBacktester) Backtest(ctx context.Context, sr *pb.SeriesRange) (test *BacktestTradingPlan, err error) { + test = &BacktestTradingPlan{ + Id: time.Now().Unix(), + UserId: 10001, + PlanId: b.plan.Id, + InstId: b.plan.InstId, + Exchange: pb.ExchangeType(b.plan.Exchange), + Interval: b.plan.Interval, + SeriesBefore: sr.Before, + SeriesAfter: sr.After, + Ctime: time.Now().UnixMilli(), + Cash: b.account.cash, + } // 交易信号回测器 sigStrategyBacktester := NewSigStrategyBacktester(b.sigStrategyType, b.sigStrategy, b.indicatorReg, b.exchangeClient) - onInterval1m, sigTimes := 0, 0 sigStrategyBacktester.SubKline(types.Interval1m, func(interval types.Interval, k types.Kline) (err error) { - onInterval1m++ // 检查仓位平仓 return b.closeByKlineInterval1m(k) }) - err = sigStrategyBacktester.Backtest(ctx, sr, nil, func(sigSide types.Side, k types.Kline) (err error) { - sigTimes++ + intervalSeries := types.NewIntervalState[*sig.KlineSeries]() + err = sigStrategyBacktester.Backtest(ctx, sr, intervalSeries, func(sigSide types.Side, k types.Kline) (err error) { + test.Singals++ // 根据交易信号检查仓位平仓 if err = b.closeBySigSingal(sigSide, k); err != nil { return @@ -109,16 +125,47 @@ func (b *TradingPlanBacktester) Backtest(ctx context.Context, sr *pb.SeriesRange return } - // todo 回测报告 - var trades []*trade.Trade + // 读最新的k线 + kSeries := intervalSeries.Get(types.Interval(sr.Interval)) + lastCandle, ok := kSeries.Get(0) + if !ok { + err = fmt.Errorf("get series last candle error") + return + } + // 关闭所有未平仓仓位 + b.forceCloseAllHoldingPosition(lastCandle) + + // 回测结果 + test.Etime = time.Now().UnixMilli() + test.EndCash = b.account.cash + test.Profit = b.account.profit + test.TotalTrades = len(b.account.trades) + test.WinningTrades = b.account.winningTrades + test.LosingTrades = b.account.losingTrades + test.Fee = b.account.fee + // trades + var trades []*Trade for _, trade := range b.account.trades { - if trade.ClosePrice > 0 { - trades = append(trades, trade) + trade.BacktestId = test.Id + trade.Ctime = test.Ctime + trades = append(trades, trade) + } + collect.SortAsc(trades, func(t *Trade) int64 { return t.Id }) + test.Trades = trades + // 最大回撤 + drawdown := b.account.maxDrawdown + test.MaxDrawdown = max((drawdown[0]-drawdown[1])/drawdown[0], (drawdown[2]-drawdown[3])/drawdown[2]) + return +} + +// forceCloseAllHoldingPosition 关闭所有未平仓仓位 +func (b *TradingPlanBacktester) forceCloseAllHoldingPosition(k types.Kline) (err error) { + for _, pos := range b.account.positions { + err = b.account.ClosePosition(pos, k, trade.CauseCloseForced) + if err != nil { + return } } - collect.SortDesc(trades, func(t *trade.Trade) float64 { return t.Pnl }) - exposure := b.account.cash + b.account.CurrentExposure() - _ = exposure return } diff --git a/internal/trading/backtest/types.go b/internal/trading/backtest/types.go index e1da4f0..e8dd789 100644 --- a/internal/trading/backtest/types.go +++ b/internal/trading/backtest/types.go @@ -1,13 +1,15 @@ package backtest import ( + "sig-pub/api/pb" "sig-pub/pkg/trade" + "sig-pub/pkg/types" ) type BacktestResult struct { StartTs int64 EndTs int64 - Trades []*trade.Trade + Trades []*Trade Positions []*trade.Position Cash float64 Equity float64 @@ -17,6 +19,8 @@ type BacktestResult struct { // max_drawdown 最大回撤 // num_trades 单数 // win_rate 胜率 + // 资金利用率 150%/天 + // 资金曲线,收益分布,持仓分析 } type TradeStat struct { @@ -70,3 +74,57 @@ type TradeStat struct { // LongestDaysWithoutNewHigh int64 `json:"longest_days_without_new_high"` // 权益/盈亏最长未创新高天数 // LongestPeriodWithoutNewHigh string `json:"longest_period_without_new_high"` // 权益/盈亏最长未创新高时间段 } + +// BacktestTradingPlan 交易计划回测结果 +type BacktestTradingPlan struct { + Id int64 `json:"id" gorm:"column:id"` // 测试id + UserId int64 `json:"userId" gorm:"column:user_id"` // 用户id + PlanId int64 `json:"planId" gorm:"column:plan_id"` // 交易计划id + InstId string `json:"instId" gorm:"column:inst_id"` // 交易产品id + Exchange pb.ExchangeType `json:"exchange" gorm:"column:exchange"` // 交易所 + Interval string `json:"interval" gorm:"column:interval"` // 交易周期 + SeriesBefore int64 `json:"seriesBefore" gorm:"column:series_before"` // 回测周期开始时间 + SeriesAfter int64 `json:"seriesAfter" gorm:"column:series_after"` // 回测周期结束时间 + Ctime int64 `json:"ctime" gorm:"column:ctime"` // 创建时间(测试时间) + Etime int64 `json:"etime" gorm:"column:etime"` // 测试结束时间 + Cash float64 `json:"cash" gorm:"column:cash"` // 起始金额 + EndCash float64 `json:"endCash" gorm:"column:end_cash"` // 结束金额 + Profit float64 `json:"profit" gorm:"column:profit"` // 利润 + Singals int `json:"singals" gorm:"column:singals"` // 交易信号数 + TotalTrades int `json:"totalTrades" gorm:"column:total_trades"` // 总单数 + WinningTrades int `json:"winningTrades" gorm:"column:winning_trades"` // 盈利单数 + LosingTrades int `json:"losingTrades" gorm:"column:losing_trades"` // 亏损单数 + Fee float64 `json:"fee" gorm:"column:fee"` // 总手续费 + MaxDrawdown float64 `json:"maxDrawdown" gorm:"column:max_drawdown"` // 最大回撤 + Trades []*Trade `json:"-" gorm:"-"` // 回测交易单 +} + +func (BacktestTradingPlan) TableName() string { + return "backtest_trading_plan" +} + +// Trade 交易计划回测单 +type Trade struct { + Id int64 `json:"id" gorm:"column:id"` // 交易id + BacktestId int64 `json:"backtestId" gorm:"column:backtest_id"` // 回测id + Ctime int64 `json:"ctime" gorm:"column:ctime"` // 创建时间 + Side types.Side `json:"side" gorm:"column:side"` // 交易方向 + Qty float64 `json:"qty" gorm:"column:qty"` // 交易量 + Price float64 `json:"price" gorm:"column:price"` // 开仓价格 + Fee float64 `json:"fee" gorm:"column:fee"` // 开仓手续费 + Leverage int32 `json:"leverage" gorm:"column:leverage"` // 杠杆倍数 + Time int64 `json:"time" gorm:"column:time"` // 开仓时间 + TimeK int64 `json:"timeK" gorm:"column:time_k"` // 开仓K线时间 + ClosePrice float64 `json:"closePrice" gorm:"column:close_price"` // 平仓价格 + CloseFee float64 `json:"closeFee" gorm:"column:close_fee"` // 平仓手续费 + CloseTime int64 `json:"closeTime" gorm:"column:close_time"` // 平仓时间 + CloseTimeK int64 `json:"closeTimeK" gorm:"column:close_time_k"` // 平仓K线时间 + CloseCause trade.Cause `json:"closeCause" gorm:"column:close_cause"` // 平仓原因 ["stoploss", "takeprofit", "trailing", "retrace", "signal"](“止损”、“止盈”、“动态跟踪”、“回撤”、“信号”) + Pnl float64 `json:"pnl" gorm:"column:pnl"` // 盈利/亏损 pnl = (t.ClosePrice-t.Price)*t.Qty - t.Fee - t.CloseFee + HoldTime string `json:"holdTime" gorm:"column:hold_time"` // 持仓时间 + PeakPx float64 `json:"peakPx" gorm:"column:peakPx"` // highest (for long) or lowest (for short) observed price since entry +} + +func (Trade) TableName() string { + return "backtest_trading_trade" +} diff --git a/internal/trading/trading_data_persist.go b/internal/trading/trading_data_persist.go index 9e17f7b..a846df1 100644 --- a/internal/trading/trading_data_persist.go +++ b/internal/trading/trading_data_persist.go @@ -1,22 +1,33 @@ package trading import ( + "context" + "sig-pub/internal/trading/backtest" "sig-pub/pkg/data" "sig-pub/pkg/data/entity" + "sig-pub/pkg/storage/ck" "sig-pub/pkg/storage/persist" ) type TradingDataPersist struct { - db *persist.RDB + db *persist.RDB + ckDB *ck.ClickhouseDB + ckBatchWriter *ck.ClickhouseBatchWriter } -func NewTradingDataPersist(db *persist.RDB) *TradingDataPersist { +func NewTradingDataPersist(db *persist.RDB, ckDB *ck.ClickhouseDB, ckBatchWriter *ck.ClickhouseBatchWriter) *TradingDataPersist { return &TradingDataPersist{ - db: db, + db: db, + ckDB: ckDB, + ckBatchWriter: ckBatchWriter, } } func (p *TradingDataPersist) Init() (err error) { + err = p.ckDB.AutoMigrateTables( + &backtest.BacktestTradingPlan{}, + &backtest.Trade{}, + ) return } @@ -32,3 +43,13 @@ func (p *TradingDataPersist) GetTradePlanById(planId int64) (plan *entity.TradeP } return } + +// SaveBacktestTradingPlan 保存交易计划回测结果 +func (p *TradingDataPersist) SaveBacktestTradingPlan(ctx context.Context, backtestTradingPlan *backtest.BacktestTradingPlan) (err error) { + err = ck.InsertBatch(p.ckBatchWriter, ctx, []*backtest.BacktestTradingPlan{backtestTradingPlan}) + if err != nil { + return + } + err = ck.InsertBatch(p.ckBatchWriter, ctx, backtestTradingPlan.Trades) + return +} diff --git a/internal/trading/trading_service.go b/internal/trading/trading_service.go index 344bb1d..2ddd4cd 100644 --- a/internal/trading/trading_service.go +++ b/internal/trading/trading_service.go @@ -297,6 +297,10 @@ func (svc *TradingService) Backtest(ctx context.Context, planId, stime, etime in if err = tester.Init(10000, *plan); err != nil { return } - err = tester.Backtest(ctx, sr) + backtestTradingPlan, err := tester.Backtest(ctx, sr) + if err != nil { + return + } + err = svc.tradingDataPersist.SaveBacktestTradingPlan(ctx, backtestTradingPlan) return } diff --git a/pkg/config/config.go b/pkg/config/config.go index e70bf1c..defa845 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -1,14 +1,21 @@ package config import ( + "context" + "net" + "os" "sig-pub/pkg/zlog" + "time" "github.com/redis/go-redis/v9" clientv3 "go.etcd.io/etcd/client/v3" + "golang.org/x/crypto/ssh" "gorm.io/driver/mysql" "gorm.io/driver/postgres" "gorm.io/gorm" "gorm.io/gorm/logger" + + clickhousev2 "github.com/ClickHouse/clickhouse-go/v2" ) type Configuration struct { @@ -116,15 +123,99 @@ func (conf PostgresConfig) NewGormDB() (db *gorm.DB, err error) { } type ClickhouseConfig struct { - Logsql bool - Addr string - Database string - Username string - Password string - DialTimeout int // second - ReadTimeout int // second - NotAutoMigrateTable bool - DsnParams map[string]any + Logsql bool + LogsqlCk bool + Addr string + Database string + Username string + Password string + DialTimeout int // second + ReadTimeout int // second + DsnParams map[string]any + SSH bool + SSHUser string + SSHAddr string + SSHPassword string + SSHKey string +} + +func (cfg ClickhouseConfig) Clickhousev2Options() (options *clickhousev2.Options, err error) { + options = &clickhousev2.Options{ + Addr: []string{cfg.Addr}, + DialTimeout: time.Duration(cfg.DialTimeout) * time.Second, + ReadTimeout: time.Duration(cfg.ReadTimeout) * time.Second, + Auth: clickhousev2.Auth{ + Database: cfg.Database, + Username: cfg.Username, + Password: cfg.Password, + }, + ClientInfo: clickhousev2.ClientInfo{ + Products: []struct { + Name string + Version string + }{ + {Name: "sig", Version: "0.1"}, + }, + }, + TLS: nil, + // TLS: &tls.Config{ + // InsecureSkipVerify: true, + // }, + // Settings: clickhouse0.Settings{ + // "max_execution_time": 60, + // "max_query_size": 10000000, // sql最大长度单位byte + // }, + Settings: clickhousev2.Settings(cfg.DsnParams), + Compression: &clickhousev2.Compression{ + Method: clickhousev2.CompressionLZ4, + }, + Debug: cfg.LogsqlCk, + Debugf: zlog.Debugf, + } + + if cfg.SSH { + var auth []ssh.AuthMethod + if cfg.SSHPassword != "" { + // SSH 密码认证 + auth = append(auth, ssh.Password(cfg.SSHPassword)) + } + if cfg.SSHKey != "" { + // SSH 私钥认证 + k, e := os.ReadFile(cfg.SSHKey) + if e != nil { + return nil, e + } + signer, e := ssh.ParsePrivateKey(k) + if e != nil { + return nil, e + } + auth = append(auth, ssh.PublicKeysCallback(func() ([]ssh.Signer, error) { + return []ssh.Signer{signer}, nil + })) + } + sshConfig := &ssh.ClientConfig{ + User: cfg.SSHUser, + Auth: auth, + HostKeyCallback: ssh.InsecureIgnoreHostKey(), // 忽略主机密钥检查(仅供开发使用) + Timeout: 5 * time.Second, // SSH 连接超时 + } + + // 连接 SSH 服务器 + sshClient, e := ssh.Dial("tcp", cfg.SSHAddr, sshConfig) + if e != nil { + return nil, e + } + // 创建本地监听器 (本地端口) + localListener, e := sshClient.Dial("tcp", cfg.Addr) // 远程 ClickHouse 地址及端口 + if e != nil { + return nil, e + } + // 设置本地 dail 代理 + options.DialContext = func(ctx context.Context, addr string) (net.Conn, error) { + return localListener, nil + } + } + return } // ============== tsdb ============== diff --git a/pkg/data/sdata/trade.go b/pkg/data/sdata/trade.go new file mode 100644 index 0000000..dd5b11e --- /dev/null +++ b/pkg/data/sdata/trade.go @@ -0,0 +1 @@ +package sdata diff --git a/pkg/grpc/interceptor/recover_interceptor.go b/pkg/grpc/interceptor/recover_interceptor.go index 5eb3779..3db7302 100755 --- a/pkg/grpc/interceptor/recover_interceptor.go +++ b/pkg/grpc/interceptor/recover_interceptor.go @@ -29,7 +29,9 @@ func RecoverInterceptor(ctx context.Context, req any, server *grpc.UnaryServerIn // zlog.Debugf("call method: %s", server.FullMethod) // -> /TradingService/IndicatorSeries resp, err = handler(ctx, req) - if err == nil { + if err != nil { + zlog.Errorf("grpc request error: method=%s, %v", server.FullMethod, err) + } else { if empty, ok := resp.(*emptypb.Empty); ok && empty == nil { resp = &emptypb.Empty{} // grpc: error while marshaling: proto: Marshal called with nil } else if resp == nil { diff --git a/pkg/storage/ck/clickhouse.go b/pkg/storage/ck/clickhouse.go new file mode 100644 index 0000000..693678b --- /dev/null +++ b/pkg/storage/ck/clickhouse.go @@ -0,0 +1,73 @@ +package ck + +import ( + "sig-pub/pkg/config" + "sig-pub/pkg/zlog" + + clickhousev2 "github.com/ClickHouse/clickhouse-go/v2" + "gorm.io/driver/clickhouse" + "gorm.io/gorm" + "gorm.io/gorm/logger" +) + +type ClickhouseDB struct { + cfg config.ClickhouseConfig + db *gorm.DB +} + +func NewClickhouseDB(cfg config.ClickhouseConfig) *ClickhouseDB { + return &ClickhouseDB{ + cfg: cfg, + } +} + +// InitClickhouse 初始化ck +// https://github.com/go-gorm/clickhouse +// dsn: clickhouse://gorm:gorm@localhost:9942/gorm?dial_timeout=10s&read_timeout=20s +func (c *ClickhouseDB) Init() (err error) { + gormConfig := &gorm.Config{} + if c.cfg.Logsql { + gormConfig.Logger = logger.Default.LogMode(logger.Info) // 打印sql + } + + // initial db + options, err := c.cfg.Clickhousev2Options() + if err != nil { + return + } + ckDB := clickhousev2.OpenDB(options) + c.db, err = gorm.Open(clickhouse.New(clickhouse.Config{Conn: ckDB}), gormConfig) + if err != nil { + return + } + return +} + +func (c *ClickhouseDB) DB() *gorm.DB { + return c.db +} + +// sql查询数据 +func (c *ClickhouseDB) Select(r interface{}, sql string, args ...any) (err error) { + return c.db.Raw(sql, args...).Scan(r).Error +} + +// 批量插入 +func (c *ClickhouseDB) Insert(data any) (err error) { + err = c.db.Create(data).Error + return +} + +// AutoMigrateTables 自动对齐表结构,自动根据字段修改数据库表结构,只会加改不会删字段 +func (c *ClickhouseDB) AutoMigrateTables(gormStructs ...any) (err error) { + if len(gormStructs) == 0 { + return + } + err = c.db.AutoMigrate(gormStructs...) + if err != nil { + zlog.Errorf("type auto migrate clickhouse table error: %v", err) + return + } + // db.Set("gorm:table_options", "ENGINE=Distributed(cluster, default, hits)").AutoMigrate(&entity.TradeRecord{}) + return +} diff --git a/pkg/storage/ck/clickhouse.sql b/pkg/storage/ck/clickhouse.sql new file mode 100644 index 0000000..aa7903e --- /dev/null +++ b/pkg/storage/ck/clickhouse.sql @@ -0,0 +1,27 @@ +-- optimize table sig.col_log_vbdiamond final; + +create database sig; + +create table sig.backtest_trading_plan ( + id Int64 comment 'id', + user_id Int64 comment '用户id', + plan_id Int64 comment '交易计划id', + ctime Int64 comment '创建时间' +) ENGINE = MergeTree() +PARTITION BY toYYYYMM(toDateTime(ctime/1000)) +primary key (plan_id, user_id, id) +order by (plan_id, user_id, id) +SETTINGS index_granularity = 8192 +comment '交易计划回测'; + +create table sig.backtest_trading_trade ( + id Int64 comment 'id', + backtest_id Int64 comment '回测id', + ctime Int64 comment '创建时间' +) ENGINE = MergeTree() +PARTITION BY toYYYYMM(toDateTime(ctime/1000)) +primary key (backtest_id, id) +order by (backtest_id, id) +SETTINGS index_granularity = 8192 +comment '交易计划回测订单'; + diff --git a/pkg/storage/ck/clickhouse_batch_writer.go b/pkg/storage/ck/clickhouse_batch_writer.go new file mode 100644 index 0000000..e77f326 --- /dev/null +++ b/pkg/storage/ck/clickhouse_batch_writer.go @@ -0,0 +1,160 @@ +package ck + +import ( + "context" + "fmt" + "reflect" + "sig-pub/pkg/config" + "strings" + "time" + + "github.com/ClickHouse/clickhouse-go/v2/lib/driver" + + clickhousev2 "github.com/ClickHouse/clickhouse-go/v2" +) + +type ClickhouseBatchWriter struct { + cfg config.ClickhouseConfig + conn driver.Conn +} + +func NewClickhouseBatchWriter(cfg config.ClickhouseConfig) *ClickhouseBatchWriter { + return &ClickhouseBatchWriter{ + cfg: cfg, + } +} + +func (c *ClickhouseBatchWriter) Init() (err error) { + options, err := c.cfg.Clickhousev2Options() + if err != nil { + return + } + // initial clickhouse conn + c.conn, err = clickhousev2.Open(options) + if err != nil { + return + } + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err = c.conn.Ping(ctx); err != nil { + if exception, ok := err.(*clickhousev2.Exception); ok { + err = fmt.Errorf("exception [%d] %s %s", exception.Code, exception.Message, exception.StackTrace) + } + return + } + return +} + +func reflectGormData(data any) (table string, columns []string, values map[string]any, err error) { + // value := reflect.ValueOf(data) + // tableNameM := value.MethodByName("TableName") + // _ = tableNameM.Call(nil) + defer func() { + if r := recover(); r != nil { + err = fmt.Errorf("%T parse gorm data error %v", data, r) + } + }() + + values = make(map[string]any) + refValue := reflect.ValueOf(data) + refType := reflect.TypeOf(data) + + tableNameM := refValue.MethodByName("TableName") + if !tableNameM.IsValid() { + err = fmt.Errorf("type %T not has method TableName", data) + return + } + rsp := tableNameM.Call(nil) + table = rsp[0].Interface().(string) + + if refValue.Kind() == reflect.Ptr { + refValue = refValue.Elem() + } + if refType.Kind() == reflect.Ptr { + refType = refType.Elem() + } + + for i := 0; i < refType.NumField(); i++ { + field := refType.Field(i) + value := refValue.Field(i).Interface() + // value + tag := field.Tag.Get("gorm") + column := getGormTagColumnName(tag) + if column == "" || strings.HasPrefix(column, "-") { // 忽略字段 + continue + } + columns = append(columns, column) + values[column] = value + } + return +} + +func getGormTagColumnName(tag string) (column string) { + i1 := strings.Index(tag, "column:") + if i1 < 0 { + return + } + i2 := strings.Index(tag[i1+7:], ";") + if i2 < 0 { + return tag[i1+7:] + } + return tag[i1+7 : i2+7] +} + +// InsertBatch 批量插入gorm标签的结构体, 结构体用指针! +func InsertBatch[T any](c *ClickhouseBatchWriter, ctx context.Context, datas []T) (err error) { + if len(datas) == 0 { + return + } + + var tColumns = make(map[string][]string, 3) + var tBatchs = make(map[string]driver.Batch, 3) + + for _, data := range datas { + // 优化原生批量插入: https://clickhouse.com/docs/en/integrations/go#batch-insert + table, columns, values, e := reflectGormData(data) + if e != nil { + err = e + return + } + // 初始化该表 prepare + if _, ok := tColumns[table]; !ok { + tColumns[table] = columns + + prepareSql := fmt.Sprintf("INSERT INTO %s(", table) + for i, column := range columns { + if i == 0 { + prepareSql += column + } else { + prepareSql += ("," + column) + } + } + prepareSql += ") SETTINGS async_insert=1, wait_for_async_insert=0" + batch, e := c.conn.PrepareBatch(ctx, prepareSql) + if e != nil { + err = e + return + } + tBatchs[table] = batch + } + columns = tColumns[table] + batch := tBatchs[table] + var args = make([]any, 0, len(columns)) + for _, column := range columns { + args = append(args, values[column]) + } + err = batch.Append(args...) + if err != nil { + return + } + } + // 批量插入 + for _, batch := range tBatchs { + err = batch.Send() + if err != nil { + return + } + } + return +} diff --git a/pkg/trade/close_strategy.go b/pkg/trade/close_strategy.go index 2413743..1b6c4f3 100644 --- a/pkg/trade/close_strategy.go +++ b/pkg/trade/close_strategy.go @@ -48,21 +48,11 @@ func (s *CloseStrategy) OnPrice(price float64, pos *Position) (closePos bool, ca return } // update peak px - peakPx, ok := pos.GetStateF64("peakPx") - updatePeakPx := false - if !ok { - updatePeakPx = true - } else { - if pos.Side == types.SideLong && price > peakPx { - updatePeakPx = true - } - if pos.Side == types.SideShort && price < peakPx { - updatePeakPx = true - } + if pos.Side == types.SideLong && (price > pos.PeakPx) { + pos.PeakPx = price } - if updatePeakPx { - peakPx = price - pos.SetState("peakPx", price) + if pos.Side == types.SideShort && (price < pos.PeakPx) { + pos.PeakPx = price } entry := pos.EntryPx @@ -79,7 +69,7 @@ func (s *CloseStrategy) OnPrice(price float64, pos *Position) (closePos bool, ca // 基于最高利润动态止盈 if len(s.ProfitRetracePcts) > 0 { // peak profit fraction - peakProfit := (peakPx - entry) / entry + peakProfit := (pos.PeakPx - entry) / entry minProfitToTrail, trailingPct := float64(0), float64(0) for _, profit := range s.ProfitRetracePcts { if len(profit) != 2 { @@ -87,15 +77,13 @@ func (s *CloseStrategy) OnPrice(price float64, pos *Position) (closePos bool, ca } _minProfitToTrail := profit[0] // 启动最高利润回撤的最小盈利阈值 _trailingPct := profit[1] // 基于最高利润回撤触发平仓 - if peakProfit >= _minProfitToTrail { - if _minProfitToTrail > minProfitToTrail { - minProfitToTrail = _minProfitToTrail - trailingPct = _trailingPct - } + if peakProfit >= _minProfitToTrail && _minProfitToTrail > minProfitToTrail { + minProfitToTrail = _minProfitToTrail + trailingPct = _trailingPct } } if minProfitToTrail > 0 && trailingPct > 0 { - trail := peakPx * (1 - trailingPct) + trail := entry + (pos.PeakPx-entry)*(1-trailingPct) if price <= trail { return true, CauseCloseTrailing } @@ -114,7 +102,7 @@ func (s *CloseStrategy) OnPrice(price float64, pos *Position) (closePos bool, ca // 基于最高利润动态止盈 if len(s.ProfitRetracePcts) > 0 { // peak profit fraction - peakProfit := (entry - peakPx) / entry + peakProfit := (entry - pos.PeakPx) / entry minProfitToTrail, trailingPct := float64(0), float64(0) for _, profit := range s.ProfitRetracePcts { if len(profit) != 2 { @@ -128,7 +116,7 @@ func (s *CloseStrategy) OnPrice(price float64, pos *Position) (closePos bool, ca } } if minProfitToTrail > 0 && trailingPct > 0 { - trail := peakPx * (1 + trailingPct) + trail := entry - (entry-pos.PeakPx)*(1+trailingPct) if price >= trail { return true, CauseCloseTrailing } diff --git a/pkg/trade/risk/manager.go b/pkg/trade/risk/manager.go.txt similarity index 100% rename from pkg/trade/risk/manager.go rename to pkg/trade/risk/manager.go.txt diff --git a/pkg/trade/types.go b/pkg/trade/types.go index 984406a..eee7cdc 100644 --- a/pkg/trade/types.go +++ b/pkg/trade/types.go @@ -2,18 +2,17 @@ package trade import ( "sig-pub/pkg/types" - - "github.com/spf13/cast" ) type Cause int32 const ( // _ Cause = iota - CauseCloseStoploss Cause = 1001 // 平仓:固定止损 - CauseCloseTakeprofit Cause = 1002 // 平仓:固定止盈 - CauseCloseReverseSingal Cause = 1003 // 平仓:策略反向信号 - CauseCloseTrailing Cause = 1004 // 平仓:移动止损基于最高利润点回撤百分比 + CauseCloseForced Cause = 1001 // 平仓:强制平仓 + CauseCloseStoploss Cause = 1002 // 平仓:固定止损 + CauseCloseTakeprofit Cause = 1003 // 平仓:固定止盈 + CauseCloseReverseSingal Cause = 1004 // 平仓:策略反向信号 + CauseCloseTrailing Cause = 1005 // 平仓:移动止损基于最高利润点回撤百分比 CauseRiskAlreadyTrade Cause = 2001 // 风控:已有持仓 CauseRiskSideAlreadyTrade Cause = 2002 // 风控:相同方向已有持仓 ) @@ -22,66 +21,32 @@ func (c Cause) String() string { switch c { default: return "" + case CauseCloseForced: + return "closeForced" case CauseCloseStoploss: - return "stoploss" + return "closeStoploss" case CauseCloseTakeprofit: - return "takeprofit" + return "closeTakeprofit" + case CauseCloseReverseSingal: + return "closeReverseSingal" + case CauseCloseTrailing: + return "closeTrailing" + case CauseRiskAlreadyTrade: + return "riskAlreadyTrade" + case CauseRiskSideAlreadyTrade: + return "riskSideAlreadyTrade" } } -type Trade struct { - Id int64 // 交易id - Side types.Side // 交易方向 - Qty float64 // 交易量 - Price float64 // 开仓价格 - Fee float64 // 开仓手续费 - Time int64 // 开仓时间 - ClosePrice float64 // 平仓价格 - CloseFee float64 // 平仓手续费 - CloseTs int64 // 平仓时间 - CloseCause Cause // 平仓原因 ["stoploss", "takeprofit", "trailing", "retrace", "signal"](“止损”、“止盈”、“动态跟踪”、“回撤”、“信号”) - Pnl float64 // 盈利/亏损 pnl = (t.ClosePrice-t.Price)*t.Qty - t.Fee - t.CloseFee - HoldTime string // 持仓时间 -} - // Position 持仓仓位 type Position struct { - TradeId int64 // 交易订单id - Status int32 // 1.交易中 2.持仓中 3.已平仓 - Side types.Side // 交易方向 - Qty float64 // 交易量 - EntryPx float64 // 入场价格 - EntryTs int64 // 入场时间 - PeakPx float64 // highest (for long) or lowest (for short) observed price since entry - Fee float64 // 手续费 - FeeRate float64 // 手续费率 - State map[string]any // 持仓持久化状态 -} - -func (pos *Position) SetState(k string, v any) { - if pos.State == nil { - pos.State = make(map[string]any) - } - pos.State[k] = v -} - -func (pos *Position) GetState(k string) (v any, ok bool) { - if pos.State == nil { - return - } - v, ok = pos.State[k] - return -} - -func (pos *Position) GetStateF64(k string) (f float64, ok bool) { - if pos.State == nil { - return - } - v, ok := pos.State[k] - if !ok { - return - } - f, err := cast.ToFloat64E(v) - ok = err == nil - return + TradeId int64 // 交易订单id + Status int32 // 1.交易中 2.持仓中 3.已平仓 + Side types.Side // 交易方向 + Qty float64 // 交易量 + EntryPx float64 // 入场价格 + EntryTs int64 // 入场时间 + PeakPx float64 // highest (for long) or lowest (for short) observed price since entry + Fee float64 // 手续费 + FeeRate float64 // 手续费率 }