diff --git a/.github/workflows/compatibility.yml b/.github/workflows/compatibility.yml new file mode 100644 index 00000000000..f223e0c8df9 --- /dev/null +++ b/.github/workflows/compatibility.yml @@ -0,0 +1,112 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +name: Compatibility Tests + +on: + push: + branches: [unstable] + pull_request: + paths: + - '.github/workflows/compatibility.yml' + - 'src/storage/**' + - 'src/types/**' + - 'src/commands/**' + - 'src/config/**' + - 'tests/gocase/**' + release: + types: [published] + +concurrency: + group: ${{ github.workflow }}-${{ github.event_name }}-${{ github.event.number || github.run_id }} + cancel-in-progress: true + +jobs: + build: + name: Build current kvrocks image + runs-on: ubuntu-22.04 + steps: + - uses: actions/checkout@v6 + + # The old-version images are pulled from Docker Hub by the test; only + # the "new" side is built here, once, and handed to the parallel test + # jobs. Uses the same buildx method as the nightly image build; push + # defaults to false, so the image stays local. + - name: Set up Docker Buildx + uses: docker/setup-buildx-action@4d04d5d9486b7bd6fa91e7baf45bbb4f8b9deedd + + - name: Build current kvrocks image + uses: docker/build-push-action@53b7df96c91f9c12dcc8a07bcb9ccacbed38856a # v7.3.0 + with: + context: . + load: true + tags: apache/kvrocks:current + + - name: Save image + run: docker save apache/kvrocks:current | gzip > /tmp/kvrocks-current.tar.gz + + - uses: actions/upload-artifact@v7 + with: + name: kvrocks-current-image + path: /tmp/kvrocks-current.tar.gz + retention-days: 1 + + compatibility-test: + name: "${{ matrix.version }} -> current" + runs-on: ubuntu-22.04 + needs: build + strategy: + fail-fast: false + matrix: + version: [v2.3.0, v2.6.0, v2.7.0, v2.10.0, v2.14.0, v2.15.0] + env: + KVROCKS_OLD_VERSION: ${{ matrix.version }} + # Test old-version data against the PR's own build, not the released + # :latest image, so an encoding regression in the change under test + # is caught. + KVROCKS_NEW_VERSION: current + steps: + - uses: actions/checkout@v6 + + - uses: actions/download-artifact@v8 + with: + name: kvrocks-current-image + path: /tmp + + - name: Load image + run: gunzip -c /tmp/kvrocks-current.tar.gz | docker load + + - uses: actions/setup-go@v7 + with: + go-version-file: 'tests/gocase/go.mod' + cache: false + + # Pull the old-version image up front so the test does not race a lazy + # pull inside the container lifecycle. + - name: Pull old kvrocks image + run: docker pull docker.io/apache/kvrocks:${KVROCKS_OLD_VERSION#v} + + # Each job runs on a fresh runner, so GOCACHE starts empty and the go + # test result cache cannot replay one version's result for another. + # KVROCKS_OLD_VERSION is read at runtime and is NOT part of the go test + # cache key, so if a shared cache is ever added (actions/cache on + # GOCACHE, or collapsing the matrix into a single looping job), version + # the cache key or pass -count=1 to go test. + - name: Run compatibility test + run: | + cd tests/gocase + go test -tags compat -run TestCompatibilityAllTypes ./integration/compatibility/ diff --git a/tests/gocase/go.mod b/tests/gocase/go.mod index b5fb70e0183..333c95d4fb0 100644 --- a/tests/gocase/go.mod +++ b/tests/gocase/go.mod @@ -1,27 +1,70 @@ module github.com/apache/kvrocks/tests/gocase -go 1.25 +go 1.25.0 require ( + github.com/docker/docker v28.5.2+incompatible github.com/redis/go-redis/v9 v9.17.2 - github.com/shirou/gopsutil/v4 v4.25.12 + github.com/shirou/gopsutil/v4 v4.26.2 github.com/stretchr/testify v1.11.1 + github.com/testcontainers/testcontainers-go v0.41.0 golang.org/x/exp v0.0.0-20251219203646-944ab1f22d93 golang.org/x/sync v0.19.0 ) require ( + dario.cat/mergo v1.0.2 // indirect + github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c // indirect + github.com/Microsoft/go-winio v0.6.2 // indirect + github.com/cenkalti/backoff/v4 v4.3.0 // indirect + github.com/cenkalti/backoff/v5 v5.0.3 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/containerd/errdefs v1.0.0 // indirect + github.com/containerd/errdefs/pkg v0.3.0 // indirect + github.com/containerd/log v0.1.0 // indirect + github.com/containerd/platforms v0.2.1 // indirect + github.com/cpuguy83/dockercfg v0.3.2 // indirect github.com/davecgh/go-spew v1.1.1 // indirect github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect - github.com/ebitengine/purego v0.9.1 // indirect + github.com/distribution/reference v0.6.0 // indirect + github.com/docker/go-connections v0.6.0 // indirect + github.com/docker/go-units v0.5.0 // indirect + github.com/ebitengine/purego v0.10.0 // indirect + github.com/felixge/httpsnoop v1.0.4 // indirect + github.com/go-logr/logr v1.4.3 // indirect + github.com/go-logr/stdr v1.2.2 // indirect github.com/go-ole/go-ole v1.3.0 // indirect + github.com/google/uuid v1.6.0 // indirect + github.com/klauspost/compress v1.18.2 // indirect github.com/lufia/plan9stats v0.0.0-20251013123823-9fd1530e3ec3 // indirect + github.com/magiconair/properties v1.8.10 // indirect + github.com/moby/docker-image-spec v1.3.1 // indirect + github.com/moby/go-archive v0.2.0 // indirect + github.com/moby/patternmatcher v0.6.0 // indirect + github.com/moby/sys/sequential v0.6.0 // indirect + github.com/moby/sys/user v0.4.0 // indirect + github.com/moby/sys/userns v0.1.0 // indirect + github.com/moby/term v0.5.2 // indirect + github.com/morikuni/aec v1.0.0 // indirect + github.com/opencontainers/go-digest v1.0.0 // indirect + github.com/opencontainers/image-spec v1.1.1 // indirect + github.com/pkg/errors v0.9.1 // indirect github.com/pmezard/go-difflib v1.0.0 // indirect github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect + github.com/sirupsen/logrus v1.9.3 // indirect github.com/tklauser/go-sysconf v0.3.16 // indirect github.com/tklauser/numcpus v0.11.0 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect - golang.org/x/sys v0.39.0 // indirect + go.opentelemetry.io/auto/sdk v1.2.1 // indirect + go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.49.0 // indirect + go.opentelemetry.io/otel v1.42.0 // indirect + go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.42.0 // indirect + go.opentelemetry.io/otel/metric v1.42.0 // indirect + go.opentelemetry.io/otel/sdk v1.42.0 // indirect + go.opentelemetry.io/otel/trace v1.42.0 // indirect + go.opentelemetry.io/proto/otlp v1.10.0 // indirect + golang.org/x/crypto v0.48.0 // indirect + golang.org/x/sys v0.41.0 // indirect + google.golang.org/protobuf v1.36.11 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/tests/gocase/go.sum b/tests/gocase/go.sum index f1bab797459..19f6a91c063 100644 --- a/tests/gocase/go.sum +++ b/tests/gocase/go.sum @@ -1,48 +1,176 @@ +dario.cat/mergo v1.0.2 h1:85+piFYR1tMbRrLcDwR18y4UKJ3aH1Tbzi24VRW1TK8= +dario.cat/mergo v1.0.2/go.mod h1:E/hbnu0NxMFBjpMIE34DRGLWqDy0g5FuKDhCb31ngxA= +github.com/AdaLogics/go-fuzz-headers v0.0.0-20240806141605-e8a1dd7889d6 h1:He8afgbRMd7mFxO99hRNu+6tazq8nFF9lIwo9JFroBk= +github.com/AdaLogics/go-fuzz-headers v0.0.0-20240806141605-e8a1dd7889d6/go.mod h1:8o94RPi1/7XTJvwPpRSzSUedZrtlirdB3r9Z20bi2f8= +github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c h1:udKWzYgxTojEKWjV8V+WSxDXJ4NFATAsZjh8iIbsQIg= +github.com/Azure/go-ansiterm v0.0.0-20250102033503-faa5f7b0171c/go.mod h1:xomTg63KZ2rFqZQzSB4Vz2SUXa1BpHTVz9L5PTmPC4E= +github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= +github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= +github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK343L8= +github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= +github.com/cenkalti/backoff/v5 v5.0.3 h1:ZN+IMa753KfX5hd8vVaMixjnqRZ3y8CuJKRKj1xcsSM= +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/containerd/errdefs v1.0.0 h1:tg5yIfIlQIrxYtu9ajqY42W3lpS19XqdxRQeEwYG8PI= +github.com/containerd/errdefs v1.0.0/go.mod h1:+YBYIdtsnF4Iw6nWZhJcqGSg/dwvV7tyJ/kCkyJ2k+M= +github.com/containerd/errdefs/pkg v0.3.0 h1:9IKJ06FvyNlexW690DXuQNx2KA2cUJXx151Xdx3ZPPE= +github.com/containerd/errdefs/pkg v0.3.0/go.mod h1:NJw6s9HwNuRhnjJhM7pylWwMyAkmCQvQ4GpJHEqRLVk= +github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= +github.com/containerd/log v0.1.0/go.mod h1:VRRf09a7mHDIRezVKTRCrOq78v577GXq3bSa3EhrzVo= +github.com/containerd/platforms v0.2.1 h1:zvwtM3rz2YHPQsF2CHYM8+KtB5dvhISiXh5ZpSBQv6A= +github.com/containerd/platforms v0.2.1/go.mod h1:XHCb+2/hzowdiut9rkudds9bE5yJ7npe7dG/wG+uFPw= +github.com/cpuguy83/dockercfg v0.3.2 h1:DlJTyZGBDlXqUZ2Dk2Q3xHs/FtnooJJVaad2S9GKorA= +github.com/cpuguy83/dockercfg v0.3.2/go.mod h1:sugsbF4//dDlL/i+S+rtpIWp+5h0BHJHfjj5/jFyUJc= +github.com/creack/pty v1.1.18 h1:n56/Zwd5o6whRC5PMGretI4IdRLlmBXYNjScPaBgsbY= +github.com/creack/pty v1.1.18/go.mod h1:MOBLtS5ELjhRRrroQr9kyvTxUAFNvYEK993ew/Vr4O4= +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= -github.com/ebitengine/purego v0.9.1 h1:a/k2f2HQU3Pi399RPW1MOaZyhKJL9w/xFpKAg4q1s0A= -github.com/ebitengine/purego v0.9.1/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= +github.com/distribution/reference v0.6.0 h1:0IXCQ5g4/QMHHkarYzh5l+u8T3t73zM5QvfrDyIgxBk= +github.com/distribution/reference v0.6.0/go.mod h1:BbU0aIcezP1/5jX/8MP0YiH4SdvB5Y4f/wlDRiLyi3E= +github.com/docker/docker v28.5.2+incompatible h1:DBX0Y0zAjZbSrm1uzOkdr1onVghKaftjlSWt4AFexzM= +github.com/docker/docker v28.5.2+incompatible/go.mod h1:eEKB0N0r5NX/I1kEveEz05bcu8tLC/8azJZsviup8Sk= +github.com/docker/go-connections v0.6.0 h1:LlMG9azAe1TqfR7sO+NJttz1gy6KO7VJBh+pMmjSD94= +github.com/docker/go-connections v0.6.0/go.mod h1:AahvXYshr6JgfUJGdDCs2b5EZG/vmaMAntpSFH5BFKE= +github.com/docker/go-units v0.5.0 h1:69rxXcBk27SvSaaxTtLh/8llcHD8vYHT7WSdRZ/jvr4= +github.com/docker/go-units v0.5.0/go.mod h1:fgPhTUdO+D/Jk86RDLlptpiXQzgHJF7gydDDbaIK4Dk= +github.com/ebitengine/purego v0.10.0 h1:QIw4xfpWT6GWTzaW5XEKy3HXoqrJGx1ijYHzTF0/ISU= +github.com/ebitengine/purego v0.10.0/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ= +github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2Wg= +github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U= +github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/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/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0= github.com/go-ole/go-ole v1.3.0 h1:Dt6ye7+vXGIKZ7Xtk4s6/xVdGDQynvom7xCFEdWr6uE= github.com/go-ole/go-ole v1.3.0/go.mod h1:5LS6F96DhAwUc7C+1HLexzMXY1xGRSryjyPPKW6zv78= 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= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0 h1:HWRh5R2+9EifMyIHV7ZV+MIZqgz+PMpZ14Jynv3O2Zs= +github.com/grpc-ecosystem/grpc-gateway/v2 v2.28.0/go.mod h1:JfhWUomR1baixubs02l85lZYYOm7LV6om4ceouMv45c= +github.com/klauspost/compress v1.18.2 h1:iiPHWW0YrcFgpBYhsA6D1+fqHssJscY/Tm/y2Uqnapk= +github.com/klauspost/compress v1.18.2/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/lufia/plan9stats v0.0.0-20251013123823-9fd1530e3ec3 h1:PwQumkgq4/acIiZhtifTV5OUqqiP82UAl0h87xj/l9k= github.com/lufia/plan9stats v0.0.0-20251013123823-9fd1530e3ec3/go.mod h1:autxFIvghDt3jPTLoqZ9OZ7s9qTGNAWmYCjVFWPX/zg= +github.com/magiconair/properties v1.8.10 h1:s31yESBquKXCV9a/ScB3ESkOjUYYv+X0rg8SYxI99mE= +github.com/magiconair/properties v1.8.10/go.mod h1:Dhd985XPs7jluiymwWYZ0G4Z61jb3vdS329zhj2hYo0= +github.com/moby/docker-image-spec v1.3.1 h1:jMKff3w6PgbfSa69GfNg+zN/XLhfXJGnEx3Nl2EsFP0= +github.com/moby/docker-image-spec v1.3.1/go.mod h1:eKmb5VW8vQEh/BAr2yvVNvuiJuY6UIocYsFu/DxxRpo= +github.com/moby/go-archive v0.2.0 h1:zg5QDUM2mi0JIM9fdQZWC7U8+2ZfixfTYoHL7rWUcP8= +github.com/moby/go-archive v0.2.0/go.mod h1:mNeivT14o8xU+5q1YnNrkQVpK+dnNe/K6fHqnTg4qPU= +github.com/moby/patternmatcher v0.6.0 h1:GmP9lR19aU5GqSSFko+5pRqHi+Ohk1O69aFiKkVGiPk= +github.com/moby/patternmatcher v0.6.0/go.mod h1:hDPoyOpDY7OrrMDLaYoY3hf52gNCR/YOUYxkhApJIxc= +github.com/moby/sys/atomicwriter v0.1.0 h1:kw5D/EqkBwsBFi0ss9v1VG3wIkVhzGvLklJ+w3A14Sw= +github.com/moby/sys/atomicwriter v0.1.0/go.mod h1:Ul8oqv2ZMNHOceF643P6FKPXeCmYtlQMvpizfsSoaWs= +github.com/moby/sys/sequential v0.6.0 h1:qrx7XFUd/5DxtqcoH1h438hF5TmOvzC/lspjy7zgvCU= +github.com/moby/sys/sequential v0.6.0/go.mod h1:uyv8EUTrca5PnDsdMGXhZe6CCe8U/UiTWd+lL+7b/Ko= +github.com/moby/sys/user v0.4.0 h1:jhcMKit7SA80hivmFJcbB1vqmw//wU61Zdui2eQXuMs= +github.com/moby/sys/user v0.4.0/go.mod h1:bG+tYYYJgaMtRKgEmuueC0hJEAZWwtIbZTB+85uoHjs= +github.com/moby/sys/userns v0.1.0 h1:tVLXkFOxVu9A64/yh59slHVv9ahO9UIev4JZusOLG/g= +github.com/moby/sys/userns v0.1.0/go.mod h1:IHUYgu/kao6N8YZlp9Cf444ySSvCmDlmzUcYfDHOl28= +github.com/moby/term v0.5.2 h1:6qk3FJAFDs6i/q3W/pQ97SX192qKfZgGjCQqfCJkgzQ= +github.com/moby/term v0.5.2/go.mod h1:d3djjFCrjnB+fl8NJux+EJzu0msscUP+f8it8hPkFLc= +github.com/morikuni/aec v1.0.0 h1:nP9CBfwrvYnBRgY6qfDQkygYDmYwOilePFkwzv4dU8A= +github.com/morikuni/aec v1.0.0/go.mod h1:BbKIizmSmc5MMPqRYbxO4ZU0S0+P200+tUnFx7PXmsc= +github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= +github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= +github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= +github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= +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 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 h1:o4JXh1EVt9k/+g42oCprj/FisM4qX9L3sZB3upGN2ZU= github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55/go.mod h1:OmDBASR4679mdNQnz2pUhc2G8CO2JrUAVFDRBDP/hJE= github.com/redis/go-redis/v9 v9.17.2 h1:P2EGsA4qVIM3Pp+aPocCJ7DguDHhqrXNhVcEp4ViluI= github.com/redis/go-redis/v9 v9.17.2/go.mod h1:u410H11HMLoB+TP67dz8rL9s6QW2j76l0//kSOd3370= -github.com/shirou/gopsutil/v4 v4.25.12 h1:e7PvW/0RmJ8p8vPGJH4jvNkOyLmbkXgXW4m6ZPic6CY= -github.com/shirou/gopsutil/v4 v4.25.12/go.mod h1:EivAfP5x2EhLp2ovdpKSozecVXn1TmuG7SMzs/Wh4PU= +github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= +github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/shirou/gopsutil/v4 v4.26.2 h1:X8i6sicvUFih4BmYIGT1m2wwgw2VG9YgrDTi7cIRGUI= +github.com/shirou/gopsutil/v4 v4.26.2/go.mod h1:LZ6ewCSkBqUpvSOf+LsTGnRinC6iaNUNMGBtDkJBaLQ= +github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= +github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/objx v0.5.2 h1:xuMeJ0Sdp5ZMRXx/aWO6RZxdr3beISkG5/G/aIRr3pY= +github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/testcontainers/testcontainers-go v0.41.0 h1:mfpsD0D36YgkxGj2LrIyxuwQ9i2wCKAD+ESsYM1wais= +github.com/testcontainers/testcontainers-go v0.41.0/go.mod h1:pdFrEIfaPl24zmBjerWTTYaY0M6UHsqA1YSvsoU40MI= github.com/tklauser/go-sysconf v0.3.16 h1:frioLaCQSsF5Cy1jgRBrzr6t502KIIwQ0MArYICU0nA= github.com/tklauser/go-sysconf v0.3.16/go.mod h1:/qNL9xxDhc7tx3HSRsLWNnuzbVfh3e7gh/BmM179nYI= github.com/tklauser/numcpus v0.11.0 h1:nSTwhKH5e1dMNsCdVBukSZrURJRoHbSEQjdEbY+9RXw= github.com/tklauser/numcpus v0.11.0/go.mod h1:z+LwcLq54uWZTX0u/bGobaV34u6V7KNlTZejzM6/3MQ= github.com/yusufpapurcu/wmi v1.2.4 h1:zFUKzehAFReQwLys1b/iSMl+JQGSCSjtVqQn9bBrPo0= github.com/yusufpapurcu/wmi v1.2.4/go.mod h1:SBZ9tNy3G9/m5Oi98Zks0QjeHVDvuK0qfxQmPyzfmi0= +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/contrib/instrumentation/net/http/otelhttp v0.49.0 h1:jq9TW8u3so/bN+JPT166wjOI6/vQPF6Xe7nMNIltagk= +go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.49.0/go.mod h1:p8pYQP+m5XfbZm9fxtSKAbM6oIllS7s2AfxrChvc7iw= +go.opentelemetry.io/otel v1.42.0 h1:lSQGzTgVR3+sgJDAU/7/ZMjN9Z+vUip7leaqBKy4sho= +go.opentelemetry.io/otel v1.42.0/go.mod h1:lJNsdRMxCUIWuMlVJWzecSMuNjE7dOYyWlqOXWkdqCc= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.42.0 h1:THuZiwpQZuHPul65w4WcwEnkX2QIuMT+UFoOrygtoJw= +go.opentelemetry.io/otel/exporters/otlp/otlptrace v1.42.0/go.mod h1:J2pvYM5NGHofZ2/Ru6zw/TNWnEQp5crgyDeSrYpXkAw= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.41.0 h1:inYW9ZhgqiDqh6BioM7DVHHzEGVq76Db5897WLGZ5Go= +go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracehttp v1.41.0/go.mod h1:Izur+Wt8gClgMJqO/cZ8wdeeMryJ/xxiOVgFSSfpDTY= +go.opentelemetry.io/otel/metric v1.42.0 h1:2jXG+3oZLNXEPfNmnpxKDeZsFI5o4J+nz6xUlaFdF/4= +go.opentelemetry.io/otel/metric v1.42.0/go.mod h1:RlUN/7vTU7Ao/diDkEpQpnz3/92J9ko05BIwxYa2SSI= +go.opentelemetry.io/otel/sdk v1.42.0 h1:LyC8+jqk6UJwdrI/8VydAq/hvkFKNHZVIWuslJXYsDo= +go.opentelemetry.io/otel/sdk v1.42.0/go.mod h1:rGHCAxd9DAph0joO4W6OPwxjNTYWghRWmkHuGbayMts= +go.opentelemetry.io/otel/trace v1.42.0 h1:OUCgIPt+mzOnaUTpOQcBiM/PLQ/Op7oq6g4LenLmOYY= +go.opentelemetry.io/otel/trace v1.42.0/go.mod h1:f3K9S+IFqnumBkKhRJMeaZeNk9epyhnCmQh/EysQCdc= +go.opentelemetry.io/proto/otlp v1.10.0 h1:IQRWgT5srOCYfiWnpqUYz9CVmbO8bFmKcwYxpuCSL2g= +go.opentelemetry.io/proto/otlp v1.10.0/go.mod h1:/CV4QoCR/S9yaPj8utp3lvQPoqMtxXdzn7ozvvozVqk= +golang.org/x/crypto v0.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts= +golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos= golang.org/x/exp v0.0.0-20251219203646-944ab1f22d93 h1:fQsdNF2N+/YewlRZiricy4P1iimyPKZ/xwniHj8Q2a0= golang.org/x/exp v0.0.0-20251219203646-944ab1f22d93/go.mod h1:EPRbTFwzwjXj9NpYyyrvenVh9Y+GFeEvMNh7Xuz7xgU= +golang.org/x/net v0.50.0 h1:ucWh9eiCGyDR3vtzso0WMQinm2Dnt8cFMuQa9K33J60= +golang.org/x/net v0.50.0/go.mod h1:UgoSli3F/pBgdJBHCTc+tp3gmrU4XswgGRgtnwWTfyM= golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4= golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= golang.org/x/sys v0.0.0-20190916202348-b4ddaad3f8a3/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210616094352-59db8d763f22/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= -golang.org/x/sys v0.39.0 h1:CvCKL8MeisomCi6qNZ+wbb0DN9E5AATixKsvNtMoMFk= -golang.org/x/sys v0.39.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM= +golang.org/x/sys v0.41.0 h1:Ivj+2Cp/ylzLiEU89QhWblYnOE9zerudt9Ftecq2C6k= +golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/term v0.40.0 h1:36e4zGLqU4yhjlmxEaagx2KuYbJq3EwY8K943ZsHcvg= +golang.org/x/term v0.40.0/go.mod h1:w2P8uVp06p2iyKKuvXIm7N/y0UCRt3UfJTfZ7oOpglM= +golang.org/x/text v0.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk= +golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= +golang.org/x/time v0.0.0-20220210224613-90d013bbcef8 h1:vVKdlvoWBphwdxWKrFZEuM0kGgGLxUOYcY4U/2Vjg44= +golang.org/x/time v0.0.0-20220210224613-90d013bbcef8/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= +google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57 h1:JLQynH/LBHfCTSbDWl+py8C+Rg/k1OVH3xfcaiANuF0= +google.golang.org/genproto/googleapis/api v0.0.0-20260209200024-4cfbd4190f57/go.mod h1:kSJwQxqmFXeo79zOmbrALdflXQeAYcUbgS7PbpMknCY= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260209200024-4cfbd4190f57 h1:mWPCjDEyshlQYzBpMNHaEof6UX1PmHcaUODUywQ0uac= +google.golang.org/genproto/googleapis/rpc v0.0.0-20260209200024-4cfbd4190f57/go.mod h1:j9x/tPzZkyxcgEFkiKEEGxfvyumM01BEtsW8xzOahRQ= +google.golang.org/grpc v1.79.2 h1:fRMD94s2tITpyJGtBBn7MkMseNpOZU8ZxgC3MMBaXRU= +google.golang.org/grpc v1.79.2/go.mod h1:KmT0Kjez+0dde/v2j9vzwoAScgEPx/Bw1CYChhHLrHQ= +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-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.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gotest.tools/v3 v3.5.2 h1:7koQfIKdy+I8UTetycgUqXWSDwpgv193Ka+qRsmBY8Q= +gotest.tools/v3 v3.5.2/go.mod h1:LtdLGcnqToBH83WByAAi/wiwSFCArdFIUV/xxN4pcjA= diff --git a/tests/gocase/integration/compatibility/compatibility_test.go b/tests/gocase/integration/compatibility/compatibility_test.go new file mode 100644 index 00000000000..fbe3338b795 --- /dev/null +++ b/tests/gocase/integration/compatibility/compatibility_test.go @@ -0,0 +1,163 @@ +//go:build compat + +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package compatibility + +import ( + "context" + "fmt" + "io" + "os" + "strings" + "testing" + "time" + + dimage "github.com/docker/docker/api/types/image" + + "github.com/redis/go-redis/v9" + "github.com/stretchr/testify/require" + "github.com/testcontainers/testcontainers-go" + "github.com/testcontainers/testcontainers-go/wait" +) + +const ( + kvrocksPort = "6666" + kvrocksDataDir = "/var/lib/kvrocks/db" + defaultVersion = "v2.3.0" +) + +var oldVersion = os.Getenv("KVROCKS_OLD_VERSION") +var newVersion = os.Getenv("KVROCKS_NEW_VERSION") + +func init() { + if oldVersion == "" { + oldVersion = defaultVersion + } + if newVersion == "" { + newVersion = "latest" + } +} + +func startKvrocksContainer(ctx context.Context, t *testing.T, image, volumeName string) (testcontainers.Container, error) { + client, err := testcontainers.NewDockerClientWithOpts(ctx) + if err != nil { + return nil, fmt.Errorf("connect to docker: %w", err) + } + imageInspect, err := client.ImageInspect(ctx, image) + if err != nil { + // Image not present locally (e.g. fresh CI runner): pull it first. + // Container creation would pull anyway, but the entrypoint rewrite + // below needs the image's config. + rc, pullErr := client.ImagePull(ctx, image, dimage.PullOptions{}) + if pullErr != nil { + return nil, fmt.Errorf("pull image %s: %w", image, pullErr) + } + _, _ = io.Copy(io.Discard, rc) // consume until the pull completes + rc.Close() + imageInspect, err = client.ImageInspect(ctx, image) + if err != nil { + return nil, fmt.Errorf("inspect image %s after pull: %w", image, err) + } + } + originalEntrypointCmd := strings.Join(imageInspect.Config.Entrypoint, " ") + workingDir := imageInspect.Config.WorkingDir + + // older images had default user as root; in the entrypoint, we check if we're running as uid 999, then: + // 1. set ownership of the database mounted volume, as older versions will have this directory owned by root, + // 2. change user to 999:999 kvrocks + preambleCmd := fmt.Sprintf("[ $(id -u) -eq 999 ] && chown -R 999:999 /var/lib/kvrocks/db && su kvrocks || cd %s", workingDir) + // Use /bin/sh: alpine-based images (v2.6.0, v2.7.0) have no bash. Bind to + // 0.0.0.0: older entrypoints omit --bind and their conf defaults to + // 127.0.0.1, unreachable via the published port. + entrypoint := []string{"/bin/sh", "-c", preambleCmd + " && " + originalEntrypointCmd + " --bind 0.0.0.0"} + + req := testcontainers.ContainerRequest{ + Image: image, + Entrypoint: entrypoint, + ExposedPorts: []string{fmt.Sprintf("%s/tcp", kvrocksPort)}, + User: "root", + Mounts: []testcontainers.ContainerMount{ + { + Source: testcontainers.GenericVolumeMountSource{Name: volumeName}, + Target: kvrocksDataDir, + }, + }, + WaitingFor: wait.ForListeningPort(kvrocksPort), + } + + return testcontainers.GenericContainer(ctx, testcontainers.GenericContainerRequest{ + ContainerRequest: req, + Started: true, + }) +} + +func terminateContainer(ctx context.Context, c testcontainers.Container, volumeName string) { + err := c.Terminate(ctx, testcontainers.RemoveVolumes(volumeName)) + if err != nil { + fmt.Printf("Warning: failed to terminate container and remove volume: %v\n", err) + } +} + +func getContainerAddr(c testcontainers.Container, ctx context.Context, t *testing.T) string { + hostPort, err := c.MappedPort(ctx, kvrocksPort) + require.NoError(t, err) + // Host() returns "localhost", which can resolve to ::1 and fail on hosts + // without IPv6 docker bindings; dial the IPv4 loopback directly. + return fmt.Sprintf("127.0.0.1:%s", hostPort.Port()) +} + +func TestCompatibilityAllTypes(t *testing.T) { + ctx := context.Background() + + // Unique per run: a crashed container leaks its volume, and a reused name + // would then make every later run fail opening RocksDB (LOCK contention). + volumeName := fmt.Sprintf("kvrocks-compat-%s-%s-%d", oldVersion, t.Name(), time.Now().UnixNano()) + // Docker Hub tags releases without the 'v' prefix (e.g. 2.3.0, not v2.3.0) + oldImage := fmt.Sprintf("docker.io/apache/kvrocks:%s", strings.TrimPrefix(oldVersion, "v")) + // newImage is env-overridable so CI can build the PR's own source and test + // against it instead of the last released image. + newImage := fmt.Sprintf("docker.io/apache/kvrocks:%s", newVersion) + + // Start old kvrocks and populate data + oldC, err := startKvrocksContainer(ctx, t, oldImage, volumeName) + require.NoError(t, err) + + oldClient := redis.NewClient(&redis.Options{Addr: getContainerAddr(oldC, ctx, t)}) + require.NoError(t, oldClient.Ping(ctx).Err()) + + // Populate test data based on old version capabilities + testData, err := PopulateTestData(ctx, oldClient, oldVersion) + require.NoError(t, err) + + require.NoError(t, oldClient.Close()) + require.NoError(t, oldC.Terminate(ctx)) + + // Start new kvrocks with same volume and verify data + newC, err := startKvrocksContainer(ctx, t, newImage, volumeName) + require.NoError(t, err) + defer terminateContainer(ctx, newC, volumeName) + + newClient := redis.NewClient(&redis.Options{Addr: getContainerAddr(newC, ctx, t)}) + defer func() { require.NoError(t, newClient.Close()) }() + + require.NoError(t, newClient.Ping(ctx).Err()) + require.NoError(t, VerifyTestData(ctx, newClient, testData)) +} diff --git a/tests/gocase/integration/compatibility/testdata.go b/tests/gocase/integration/compatibility/testdata.go new file mode 100644 index 00000000000..1cbd15e8351 --- /dev/null +++ b/tests/gocase/integration/compatibility/testdata.go @@ -0,0 +1,1104 @@ +//go:build compat + +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package compatibility + +import ( + "context" + "encoding/json" + "fmt" + "math" + "math/rand" + "strconv" + "strings" + + "github.com/redis/go-redis/v9" +) + +// SeededRandom provides deterministic random generation for reproducible tests +type SeededRandom struct { + rng *rand.Rand +} + +// NewSeededRandom creates a SeededRandom with the given seed +func NewSeededRandom(seed int64) *SeededRandom { + return &SeededRandom{rng: rand.New(rand.NewSource(seed))} +} + +// String generates a random alphanumeric string of the given length +func (r *SeededRandom) String(length int) string { + const charset = "abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789" + b := make([]byte, length) + for i := range b { + b[i] = charset[r.rng.Intn(len(charset))] + } + return string(b) +} + +// Bytes generates random bytes of the given length +func (r *SeededRandom) Bytes(length int) []byte { + b := make([]byte, length) + r.rng.Read(b) + return b +} + +// Int63 returns a random int64 in range [0, upper) +func (r *SeededRandom) Int63(upper int64) int64 { + return r.rng.Int63n(upper) +} + +// Int returns a random int in range [0, upper) +func (r *SeededRandom) Int(upper int) int { + return r.rng.Intn(upper) +} + +// Float64 returns a random float64 in [0, 1) +func (r *SeededRandom) Float64() float64 { + return r.rng.Float64() +} + +// TestData contains all populated test data for verification +type TestData struct { + Strings map[string]string + Hashes map[string]map[string]string + Lists map[string][]string + Sets map[string][]string + ZSets map[string][]ZSetEntry + Streams map[string][]StreamEntry + Bitmaps map[string][]byte + BloomFilters map[string][]string + JSONs map[string]string + HLLs map[string][]string + TDigests map[string][]float64 + TimeSeries map[string][]TSSample +} + +// TSSample represents a time series sample +type TSSample struct { + Timestamp uint64 + Value float64 +} + +// ZSetEntry represents a sorted set member +type ZSetEntry struct { + Score float64 + Member string +} + +// StreamEntry represents a stream entry +type StreamEntry struct { + ID string + Field string + Value string +} + +// Data type identifiers +const ( + TypeString = "string" + TypeHash = "hash" + TypeList = "list" + TypeSet = "set" + TypeZSet = "zset" + TypeStream = "stream" + TypeBitmap = "bitmap" + TypeBloomFilter = "bloomfilter" + TypeJSON = "json" + TypeHyperLogLog = "hyperloglog" + TypeTDigest = "tdigest" + TypeTimeSeries = "timeseries" +) + +// supportedTypesByVersion maps version strings to supported data types. +// Verified against the command registration in each release tag (e.g. +// bf.reserve in v2.6.0, json.set in v2.7.0, pfadd in v2.10.0, +// tdigest.create in v2.14.0; timeseries exists in v2.14.0 source but its +// REDIS_REGISTER_COMMANDS is commented out, so it is not listed there and +// first appears in v2.15.0). +var supportedTypesByVersion = map[string][]string{ + "v2.3.0": {TypeString, TypeHash, TypeList, TypeSet, TypeZSet, TypeStream, TypeBitmap}, + "v2.6.0": {TypeString, TypeHash, TypeList, TypeSet, TypeZSet, TypeStream, TypeBitmap, TypeBloomFilter}, + "v2.7.0": {TypeString, TypeHash, TypeList, TypeSet, TypeZSet, TypeStream, TypeBitmap, TypeBloomFilter, TypeJSON}, + "v2.10.0": {TypeString, TypeHash, TypeList, TypeSet, TypeZSet, TypeStream, TypeBitmap, TypeBloomFilter, TypeJSON, TypeHyperLogLog}, + "v2.14.0": {TypeString, TypeHash, TypeList, TypeSet, TypeZSet, TypeStream, TypeBitmap, TypeBloomFilter, TypeJSON, TypeHyperLogLog, TypeTDigest}, + "v2.15.0": {TypeString, TypeHash, TypeList, TypeSet, TypeZSet, TypeStream, TypeBitmap, TypeBloomFilter, TypeJSON, TypeHyperLogLog, TypeTDigest, TypeTimeSeries}, +} + +// IsTypeSupported checks if a data type is supported by the given version +func IsTypeSupported(version string, dataType string) bool { + types, ok := supportedTypesByVersion[version] + if !ok { + return false + } + for _, t := range types { + if t == dataType { + return true + } + } + return false +} + +// PopulateTestData populates all supported data types for the given version +func PopulateTestData(ctx context.Context, client *redis.Client, version string) (*TestData, error) { + data := &TestData{ + Strings: make(map[string]string), + Hashes: make(map[string]map[string]string), + Lists: make(map[string][]string), + Sets: make(map[string][]string), + ZSets: make(map[string][]ZSetEntry), + Streams: make(map[string][]StreamEntry), + Bitmaps: make(map[string][]byte), + BloomFilters: make(map[string][]string), + JSONs: make(map[string]string), + HLLs: make(map[string][]string), + TDigests: make(map[string][]float64), + TimeSeries: make(map[string][]TSSample), + } + + // Use version-based seed for determinism + seed := int64(12345) + r := NewSeededRandom(seed) + + // Populate strings (seed 0) + if IsTypeSupported(version, TypeString) { + var err error + data.Strings, err = populateStrings(ctx, client, r) + if err != nil { + return nil, fmt.Errorf("strings: %w", err) + } + } + + // Populate hashes (seed 1) + if IsTypeSupported(version, TypeHash) { + r1 := NewSeededRandom(seed + 1) + var err error + data.Hashes, err = populateHashes(ctx, client, r1) + if err != nil { + return nil, fmt.Errorf("hashes: %w", err) + } + } + + // Populate lists (seed 2) + if IsTypeSupported(version, TypeList) { + r2 := NewSeededRandom(seed + 2) + var err error + data.Lists, err = populateLists(ctx, client, r2) + if err != nil { + return nil, fmt.Errorf("lists: %w", err) + } + } + + // Populate sets (seed 3) + if IsTypeSupported(version, TypeSet) { + r3 := NewSeededRandom(seed + 3) + var err error + data.Sets, err = populateSets(ctx, client, r3) + if err != nil { + return nil, fmt.Errorf("sets: %w", err) + } + } + + // Populate zsets (seed 4) + if IsTypeSupported(version, TypeZSet) { + r4 := NewSeededRandom(seed + 4) + var err error + data.ZSets, err = populateZSets(ctx, client, r4) + if err != nil { + return nil, fmt.Errorf("zsets: %w", err) + } + } + + // Populate streams (seed 5) + if IsTypeSupported(version, TypeStream) { + r5 := NewSeededRandom(seed + 5) + var err error + data.Streams, err = populateStreams(ctx, client, r5) + if err != nil { + return nil, fmt.Errorf("streams: %w", err) + } + } + + // Populate bitmaps (seed 6) + if IsTypeSupported(version, TypeBitmap) { + r6 := NewSeededRandom(seed + 6) + var err error + data.Bitmaps, err = populateBitmaps(ctx, client, r6) + if err != nil { + return nil, fmt.Errorf("bitmaps: %w", err) + } + } + + // Populate bloom filters (seed 7) + if IsTypeSupported(version, TypeBloomFilter) { + r7 := NewSeededRandom(seed + 7) + var err error + data.BloomFilters, err = populateBloomFilters(ctx, client, r7) + if err != nil { + return nil, fmt.Errorf("bloom filters: %w", err) + } + } + + // Populate JSON docs (seed 8) + if IsTypeSupported(version, TypeJSON) { + r8 := NewSeededRandom(seed + 8) + var err error + data.JSONs, err = populateJSONs(ctx, client, r8) + if err != nil { + return nil, fmt.Errorf("json: %w", err) + } + } + + // Populate hyperloglogs (seed 9) + if IsTypeSupported(version, TypeHyperLogLog) { + r9 := NewSeededRandom(seed + 9) + var err error + data.HLLs, err = populateHLLs(ctx, client, r9) + if err != nil { + return nil, fmt.Errorf("hyperloglogs: %w", err) + } + } + + // Populate tdigests (seed 10) + if IsTypeSupported(version, TypeTDigest) { + r10 := NewSeededRandom(seed + 10) + var err error + data.TDigests, err = populateTDigests(ctx, client, r10) + if err != nil { + return nil, fmt.Errorf("tdigests: %w", err) + } + } + + // Populate time series (seed 11) + if IsTypeSupported(version, TypeTimeSeries) { + r11 := NewSeededRandom(seed + 11) + var err error + data.TimeSeries, err = populateTimeSeries(ctx, client, r11) + if err != nil { + return nil, fmt.Errorf("time series: %w", err) + } + } + + return data, nil +} + +// VerifyTestData verifies all test data is preserved +func VerifyTestData(ctx context.Context, client *redis.Client, data *TestData) error { + seed := int64(12345) + + if data.Strings != nil { + r := NewSeededRandom(seed) + if err := verifyStrings(ctx, client, data.Strings, r); err != nil { + return fmt.Errorf("strings: %w", err) + } + } + + if data.Hashes != nil { + r := NewSeededRandom(seed + 1) + if err := verifyHashes(ctx, client, data.Hashes, r); err != nil { + return fmt.Errorf("hashes: %w", err) + } + } + + if data.Lists != nil { + r := NewSeededRandom(seed + 2) + if err := verifyLists(ctx, client, data.Lists, r); err != nil { + return fmt.Errorf("lists: %w", err) + } + } + + if data.Sets != nil { + r := NewSeededRandom(seed + 3) + if err := verifySets(ctx, client, data.Sets, r); err != nil { + return fmt.Errorf("sets: %w", err) + } + } + + if data.ZSets != nil { + r := NewSeededRandom(seed + 4) + if err := verifyZSets(ctx, client, data.ZSets, r); err != nil { + return fmt.Errorf("zsets: %w", err) + } + } + + if data.Streams != nil { + r := NewSeededRandom(seed + 5) + if err := verifyStreams(ctx, client, data.Streams, r); err != nil { + return fmt.Errorf("streams: %w", err) + } + } + + if data.Bitmaps != nil { + r := NewSeededRandom(seed + 6) + if err := verifyBitmaps(ctx, client, data.Bitmaps, r); err != nil { + return fmt.Errorf("bitmaps: %w", err) + } + } + + // Maps are always initialized; empty means the old version did not + // support the type, so there is nothing to verify. + if len(data.BloomFilters) > 0 { + r := NewSeededRandom(seed + 7) + if err := verifyBloomFilters(ctx, client, data.BloomFilters, r); err != nil { + return fmt.Errorf("bloom filters: %w", err) + } + } + + if len(data.JSONs) > 0 { + r := NewSeededRandom(seed + 8) + if err := verifyJSONs(ctx, client, data.JSONs, r); err != nil { + return fmt.Errorf("json: %w", err) + } + } + + if len(data.HLLs) > 0 { + r := NewSeededRandom(seed + 9) + if err := verifyHLLs(ctx, client, data.HLLs, r); err != nil { + return fmt.Errorf("hyperloglogs: %w", err) + } + } + + if len(data.TDigests) > 0 { + r := NewSeededRandom(seed + 10) + if err := verifyTDigests(ctx, client, data.TDigests, r); err != nil { + return fmt.Errorf("tdigests: %w", err) + } + } + + if len(data.TimeSeries) > 0 { + r := NewSeededRandom(seed + 11) + if err := verifyTimeSeries(ctx, client, data.TimeSeries, r); err != nil { + return fmt.Errorf("time series: %w", err) + } + } + + return nil +} + +// populateStrings creates random string keys and returns what was stored +func populateStrings(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string]string, error) { + result := make(map[string]string) + + // Generate N random string keys + for i := 0; i < 10; i++ { + key := fmt.Sprintf("str:%d", i) + value := r.String(64) // 64 char random string + if err := client.Set(ctx, key, value, 0).Err(); err != nil { + return nil, fmt.Errorf("set %s: %w", key, err) + } + result[key] = value + } + + return result, nil +} + +// verifyStrings re-generates expected values and compares +func verifyStrings(ctx context.Context, client *redis.Client, stored map[string]string, r *SeededRandom) error { + for i := 0; i < 10; i++ { + key := fmt.Sprintf("str:%d", i) + expected := r.String(64) + actual, err := client.Get(ctx, key).Result() + if err != nil { + return fmt.Errorf("get %s: %w", key, err) + } + if actual != expected { + return fmt.Errorf("key %s: expected %q, got %q", key, expected, actual) + } + } + return nil +} + +// populateHashes creates random hash keys and returns what was stored +func populateHashes(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string]map[string]string, error) { + result := make(map[string]map[string]string) + + for i := 0; i < 5; i++ { + key := fmt.Sprintf("hash:%d", i) + fields := make(map[string]string) + fieldCount := 3 + r.Int(5) + + for j := 0; j < fieldCount; j++ { + field := fmt.Sprintf("f%d", j) + value := r.String(32) + fields[field] = value + } + + // Convert to flat args for HSET + args := make([]interface{}, 0, len(fields)*2) + for k, v := range fields { + args = append(args, k, v) + } + if err := client.HSet(ctx, key, args...).Err(); err != nil { + return nil, fmt.Errorf("hset %s: %w", key, err) + } + result[key] = fields + } + + return result, nil +} + +// verifyHashes re-generates expected values and compares +func verifyHashes(ctx context.Context, client *redis.Client, stored map[string]map[string]string, r *SeededRandom) error { + for i := 0; i < 5; i++ { + key := fmt.Sprintf("hash:%d", i) + fieldCount := 3 + r.Int(5) + + // Re-generate expected fields + fields := make(map[string]string) + for j := 0; j < fieldCount; j++ { + field := fmt.Sprintf("f%d", j) + fields[field] = r.String(32) + } + + // Compare against stored + for f, expVal := range fields { + actVal, err := client.HGet(ctx, key, f).Result() + if err != nil { + return fmt.Errorf("hget %s %s: %w", key, f, err) + } + if actVal != expVal { + return fmt.Errorf("hash %s field %s: expected %q, got %q", key, f, expVal, actVal) + } + } + } + return nil +} + +// populateLists creates random list keys and returns what was stored +func populateLists(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]string, error) { + result := make(map[string][]string) + + for i := 0; i < 5; i++ { + key := fmt.Sprintf("list:%d", i) + elemCount := 5 + r.Int(20) + elems := make([]string, elemCount) + + for j := 0; j < elemCount; j++ { + elems[j] = r.String(32) + } + + args := make([]interface{}, len(elems)) + for j, e := range elems { + args[j] = e + } + if err := client.RPush(ctx, key, args...).Err(); err != nil { + return nil, fmt.Errorf("rpush %s: %w", key, err) + } + result[key] = elems + } + + return result, nil +} + +// verifyLists re-generates expected values and compares +func verifyLists(ctx context.Context, client *redis.Client, stored map[string][]string, r *SeededRandom) error { + for i := 0; i < 5; i++ { + key := fmt.Sprintf("list:%d", i) + elemCount := 5 + r.Int(20) + + actual, err := client.LRange(ctx, key, 0, -1).Result() + if err != nil { + return fmt.Errorf("lrange %s: %w", key, err) + } + + for j := 0; j < elemCount; j++ { + expected := r.String(32) + if actual[j] != expected { + return fmt.Errorf("list %s[%d]: expected %q, got %q", key, j, expected, actual[j]) + } + } + } + return nil +} + +// populateSets creates random set keys and returns what was stored +func populateSets(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]string, error) { + result := make(map[string][]string) + + for i := 0; i < 5; i++ { + key := fmt.Sprintf("set:%d", i) + memberCount := 10 + r.Int(30) + seen := make(map[string]bool) + members := make([]string, 0, memberCount) + + for j := 0; j < memberCount; j++ { + member := r.String(32) + for seen[member] { + member = r.String(32) + } + seen[member] = true + members = append(members, member) + if err := client.SAdd(ctx, key, member).Err(); err != nil { + return nil, fmt.Errorf("sadd %s: %w", key, err) + } + } + result[key] = members + } + + return result, nil +} + +// verifySets re-generates expected values and compares (order-independent) +func verifySets(ctx context.Context, client *redis.Client, stored map[string][]string, r *SeededRandom) error { + for i := 0; i < 5; i++ { + key := fmt.Sprintf("set:%d", i) + memberCount := 10 + r.Int(30) + + actual, err := client.SMembers(ctx, key).Result() + if err != nil { + return fmt.Errorf("smembers %s: %w", key, err) + } + + // Re-generate expected and check membership + seen := make(map[string]bool) + for j := 0; j < memberCount; j++ { + member := r.String(32) + for seen[member] { + member = r.String(32) + } + seen[member] = true + found := false + for _, a := range actual { + if a == member { + found = true + break + } + } + if !found { + return fmt.Errorf("set %s: missing member %q", key, member) + } + } + } + return nil +} + +// populateZSets creates random sorted set keys and returns what was stored +func populateZSets(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]ZSetEntry, error) { + result := make(map[string][]ZSetEntry) + + for i := 0; i < 5; i++ { + key := fmt.Sprintf("zset:%d", i) + entryCount := 10 + r.Int(30) + entries := make([]ZSetEntry, entryCount) + + for j := 0; j < entryCount; j++ { + member := r.String(32) + score := r.Float64() * 1000 + entries[j] = ZSetEntry{Score: score, Member: member} + if err := client.ZAdd(ctx, key, redis.Z{Score: score, Member: member}).Err(); err != nil { + return nil, fmt.Errorf("zadd %s: %w", key, err) + } + } + result[key] = entries + } + + return result, nil +} + +// verifyZSets re-generates expected values and compares +func verifyZSets(ctx context.Context, client *redis.Client, stored map[string][]ZSetEntry, r *SeededRandom) error { + for i := 0; i < 5; i++ { + key := fmt.Sprintf("zset:%d", i) + entryCount := 10 + r.Int(30) + + actual, err := client.ZRangeWithScores(ctx, key, 0, -1).Result() + if err != nil { + return fmt.Errorf("zrange %s: %w", key, err) + } + + if len(actual) != entryCount { + return fmt.Errorf("zset %s: expected %d entries, got %d", key, entryCount, len(actual)) + } + + // Re-generate expected entries and build lookup by member + expected := make(map[string]float64) + for j := 0; j < entryCount; j++ { + member := r.String(32) + score := r.Float64() * 1000 + expected[member] = score + } + + // Verify all actual entries match expected + for _, z := range actual { + member, _ := z.Member.(string) + expScore, ok := expected[member] + if !ok { + return fmt.Errorf("zset %s: unexpected member %q", key, member) + } + if z.Score != expScore { + return fmt.Errorf("zset %s member %q: expected score %.2f, got %.2f", + key, member, expScore, z.Score) + } + } + } + return nil +} + +// populateStreams creates random stream keys and returns what was stored +func populateStreams(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]StreamEntry, error) { + result := make(map[string][]StreamEntry) + + for i := 0; i < 3; i++ { + key := fmt.Sprintf("stream:%d", i) + entryCount := 5 + r.Int(15) + entries := make([]StreamEntry, entryCount) + + for j := 0; j < entryCount; j++ { + field := fmt.Sprintf("field%d", j) + value := r.String(32) + entries[j] = StreamEntry{Field: field, Value: value} + + id, err := client.XAdd(ctx, &redis.XAddArgs{ + Stream: key, + Values: map[string]interface{}{field: value}, + }).Result() + if err != nil { + return nil, fmt.Errorf("xadd %s: %w", key, err) + } + entries[j].ID = id + } + result[key] = entries + } + + return result, nil +} + +// verifyStreams re-generates expected values and compares +func verifyStreams(ctx context.Context, client *redis.Client, stored map[string][]StreamEntry, r *SeededRandom) error { + for i := 0; i < 3; i++ { + key := fmt.Sprintf("stream:%d", i) + entryCount := 5 + r.Int(15) + + actual, err := client.XRange(ctx, key, "-", "+").Result() + if err != nil { + return fmt.Errorf("xrange %s: %w", key, err) + } + + for j := 0; j < entryCount; j++ { + field := fmt.Sprintf("field%d", j) + expected := r.String(32) + if actual[j].Values[field] != expected { + return fmt.Errorf("stream %s[%d]: expected %s=%q, got %v", + key, j, field, expected, actual[j].Values) + } + } + } + return nil +} + +// populateBitmaps creates random bitmap keys and returns what was stored +func populateBitmaps(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]byte, error) { + result := make(map[string][]byte) + + for i := 0; i < 5; i++ { + key := fmt.Sprintf("bitmap:%d", i) + size := 64 + r.Int(256) // 64-320 bytes + bytes := r.Bytes(size) + if err := client.Set(ctx, key, string(bytes), 0).Err(); err != nil { + return nil, fmt.Errorf("set %s: %w", key, err) + } + result[key] = bytes + } + + return result, nil +} + +// verifyBitmaps re-generates expected values and compares +func verifyBitmaps(ctx context.Context, client *redis.Client, stored map[string][]byte, r *SeededRandom) error { + for i := 0; i < 5; i++ { + key := fmt.Sprintf("bitmap:%d", i) + size := 64 + r.Int(256) + expected := r.Bytes(size) + + actual, err := client.Get(ctx, key).Result() + if err != nil { + return fmt.Errorf("get %s: %w", key, err) + } + if actual != string(expected) { + return fmt.Errorf("bitmap %s: mismatch", key) + } + } + return nil +} + +// populateBloomFilters adds random members to bloom filter keys +func populateBloomFilters(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]string, error) { + result := make(map[string][]string) + + for i := 0; i < 3; i++ { + key := fmt.Sprintf("bf:%d", i) + memberCount := 20 + r.Int(50) + members := make([]string, memberCount) + + for j := 0; j < memberCount; j++ { + members[j] = fmt.Sprintf("m%d-%s", j, r.String(16)) + if err := client.BFAdd(ctx, key, members[j]).Err(); err != nil { + return nil, fmt.Errorf("bf.add %s: %w", key, err) + } + } + result[key] = members + } + + return result, nil +} + +// verifyBloomFilters checks every populated member is still reported present +func verifyBloomFilters(ctx context.Context, client *redis.Client, stored map[string][]string, r *SeededRandom) error { + for i := 0; i < 3; i++ { + key := fmt.Sprintf("bf:%d", i) + memberCount := 20 + r.Int(50) + + for j := 0; j < memberCount; j++ { + member := fmt.Sprintf("m%d-%s", j, r.String(16)) + present, err := client.BFExists(ctx, key, member).Result() + if err != nil { + return fmt.Errorf("bf.exists %s %s: %w", key, member, err) + } + if !present { + return fmt.Errorf("bf %s: member %q lost", key, member) + } + } + } + return nil +} + +// populateJSONs stores random nested JSON documents +func populateJSONs(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string]string, error) { + result := make(map[string]string) + + for i := 0; i < 3; i++ { + key := fmt.Sprintf("json:%d", i) + doc := map[string]interface{}{ + "name": r.String(16), + "age": r.Int(100), + "tags": []string{r.String(8), r.String(8)}, + "nested": map[string]interface{}{"flag": r.Int(2) == 0}, + } + raw, err := json.Marshal(doc) + if err != nil { + return nil, fmt.Errorf("marshal json %s: %w", key, err) + } + if err := client.JSONSet(ctx, key, "$", doc).Err(); err != nil { + return nil, fmt.Errorf("json.set %s: %w", key, err) + } + result[key] = string(raw) + } + + return result, nil +} + +// verifyJSONs re-generates expected docs and compares, then checks JSON.DEL +func verifyJSONs(ctx context.Context, client *redis.Client, stored map[string]string, r *SeededRandom) error { + for i := 0; i < 3; i++ { + key := fmt.Sprintf("json:%d", i) + doc := map[string]interface{}{ + "name": r.String(16), + "age": r.Int(100), + "tags": []string{r.String(8), r.String(8)}, + "nested": map[string]interface{}{"flag": r.Int(2) == 0}, + } + expected, err := json.Marshal(doc) + if err != nil { + return fmt.Errorf("marshal expected %s: %w", key, err) + } + actual, err := client.JSONGet(ctx, key).Result() + if err != nil { + return fmt.Errorf("json.get %s: %w", key, err) + } + if actual != string(expected) { + return fmt.Errorf("json %s: expected %s, got %s", key, expected, actual) + } + + deleted, err := client.Do(ctx, "JSON.DEL", key, "$").Int64() + if err != nil { + return fmt.Errorf("json.del %s: %w", key, err) + } + if deleted != 1 { + return fmt.Errorf("json.del %s: expected 1, got %d", key, deleted) + } + } + return nil +} + +// hllCountTolerance is generous: HLL counts are estimates with ~0.81% error. +const hllCountTolerance = 0.05 +const hllMergeTolerance = 0.10 + +// populateHLLs adds many unique members to two keys (disjoint prefixes) +func populateHLLs(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]string, error) { + result := make(map[string][]string) + + for i := 0; i < 2; i++ { + key := fmt.Sprintf("hll:%d", i) + memberCount := 1000 + r.Int(2000) + members := make([]string, memberCount) + + for j := 0; j < memberCount; j++ { + members[j] = fmt.Sprintf("h%d-%s", i, r.String(24)) + if err := client.PFAdd(ctx, key, members[j]).Err(); err != nil { + return nil, fmt.Errorf("pfadd %s: %w", key, err) + } + } + result[key] = members + } + + return result, nil +} + +// verifyHLLs checks cardinality estimates and PFMERGE union +func verifyHLLs(ctx context.Context, client *redis.Client, stored map[string][]string, r *SeededRandom) error { + counts := make([]int64, 2) + for i := 0; i < 2; i++ { + key := fmt.Sprintf("hll:%d", i) + memberCount := int64(1000 + r.Int(2000)) + for j := 0; j < int(memberCount); j++ { + _ = fmt.Sprintf("h%d-%s", i, r.String(24)) // consume the same random sequence + } + + actual, err := client.PFCount(ctx, key).Result() + if err != nil { + return fmt.Errorf("pfcount %s: %w", key, err) + } + if diff := float64(actual-memberCount) / float64(memberCount); diff < -hllCountTolerance || diff > hllCountTolerance { + return fmt.Errorf("hll %s: count %d, expected ~%d", key, actual, memberCount) + } + counts[i] = memberCount + } + + // Merge both keys and check the union cardinality (members are disjoint) + if err := client.PFMerge(ctx, "hll:merged", "hll:0", "hll:1").Err(); err != nil { + return fmt.Errorf("pfmerge: %w", err) + } + expectedUnion := counts[0] + counts[1] + union, err := client.PFCount(ctx, "hll:merged").Result() + if err != nil { + return fmt.Errorf("pfcount merged: %w", err) + } + if diff := float64(union-expectedUnion) / float64(expectedUnion); diff < -hllMergeTolerance || diff > hllMergeTolerance { + return fmt.Errorf("hll merged: count %d, expected ~%d", union, expectedUnion) + } + return nil +} + +// populateTDigests creates t-digests with random values +func populateTDigests(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]float64, error) { + result := make(map[string][]float64) + + for i := 0; i < 3; i++ { + key := fmt.Sprintf("td:%d", i) + if err := client.Do(ctx, "TDIGEST.CREATE", key).Err(); err != nil { + return nil, fmt.Errorf("tdigest.create %s: %w", key, err) + } + valueCount := 100 + r.Int(300) + values := make([]float64, valueCount) + args := make([]interface{}, 0, valueCount+1) + args = append(args, key) + for j := 0; j < valueCount; j++ { + values[j] = r.Float64() * 1000 + args = append(args, values[j]) + } + if err := client.Do(ctx, append([]interface{}{"TDIGEST.ADD"}, args...)...).Err(); err != nil { + return nil, fmt.Errorf("tdigest.add %s: %w", key, err) + } + result[key] = values + } + + return result, nil +} + +// tdigestInfoObservations extracts the exact value count from TDIGEST.INFO. +// kvrocks replies with a RESP2 flat array of field/value pairs. +func tdigestInfoObservations(ctx context.Context, client *redis.Client, key string) (int64, error) { + info, err := client.Do(ctx, "TDIGEST.INFO", key).Result() + if err != nil { + return 0, fmt.Errorf("tdigest.info %s: %w", key, err) + } + switch v := info.(type) { + case []interface{}: + for j := 0; j+1 < len(v); j += 2 { + if strings.ToLower(fmt.Sprint(v[j])) == "observations" { + return strconv.ParseInt(fmt.Sprint(v[j+1]), 10, 64) + } + } + case map[interface{}]interface{}: + for k, obs := range v { + if strings.ToLower(fmt.Sprint(k)) == "observations" { + return strconv.ParseInt(fmt.Sprint(obs), 10, 64) + } + } + } + return 0, fmt.Errorf("tdigest.info %s: observations field not found", key) +} + +// verifyTDigests checks exact value counts, quantile sanity, and TDIGEST.MERGE +func verifyTDigests(ctx context.Context, client *redis.Client, stored map[string][]float64, r *SeededRandom) error { + counts := make([]int64, 3) + for i := 0; i < 3; i++ { + key := fmt.Sprintf("td:%d", i) + valueCount := 100 + r.Int(300) + values := make([]float64, valueCount) + for j := 0; j < valueCount; j++ { + values[j] = r.Float64() * 1000 + } + + obs, err := tdigestInfoObservations(ctx, client, key) + if err != nil { + return err + } + if obs != int64(valueCount) { + return fmt.Errorf("tdigest %s: observations %d, expected %d", key, obs, valueCount) + } + counts[i] = int64(valueCount) + + // Quantile estimates must stay within the data range + quants, err := client.Do(ctx, "TDIGEST.QUANTILE", key, 0, 0.5, 1).Result() + if err != nil { + return fmt.Errorf("tdigest.quantile %s: %w", key, err) + } + qs := make([]float64, 3) + for j, raw := range quants.([]interface{}) { + qs[j], err = strconv.ParseFloat(fmt.Sprint(raw), 64) + if err != nil { + return fmt.Errorf("tdigest.quantile %s: bad value %v: %w", key, raw, err) + } + } + lo, hi := values[0], values[0] + for _, v := range values { + if v < lo { + lo = v + } + if v > hi { + hi = v + } + } + span := hi - lo + if qs[0] < lo-0.15*span || qs[0] > lo+0.15*span || qs[2] < hi-0.15*span || qs[2] > hi+0.15*span { + return fmt.Errorf("tdigest %s: quantiles %v outside expected range [%f, %f]", key, qs, lo, hi) + } + if qs[0] > qs[1] || qs[1] > qs[2] { + return fmt.Errorf("tdigest %s: quantiles not monotonic: %v", key, qs) + } + } + + // Merge two digests and check the total observation count + if err := client.Do(ctx, "TDIGEST.MERGE", "td:merged", 2, "td:0", "td:1").Err(); err != nil { + return fmt.Errorf("tdigest.merge: %w", err) + } + mergedObs, err := tdigestInfoObservations(ctx, client, "td:merged") + if err != nil { + return err + } + if mergedObs != counts[0]+counts[1] { + return fmt.Errorf("tdigest merged: observations %d, expected %d", mergedObs, counts[0]+counts[1]) + } + return nil +} + +// tsValueTolerance allows for float round-trip error in double values +const tsValueTolerance = 1e-9 + +// populateTimeSeries adds random samples to time series keys (TS.ADD +// auto-creates the series, so no TS.CREATE is needed) +func populateTimeSeries(ctx context.Context, client *redis.Client, r *SeededRandom) (map[string][]TSSample, error) { + result := make(map[string][]TSSample) + + for i := 0; i < 3; i++ { + key := fmt.Sprintf("ts:%d", i) + sampleCount := 20 + r.Int(30) + base := uint64(1000000000 + i*1000000) + samples := make([]TSSample, sampleCount) + + for j := 0; j < sampleCount; j++ { + samples[j] = TSSample{Timestamp: base + uint64(j*1000), Value: r.Float64() * 1000} + if err := client.Do(ctx, "TS.ADD", key, samples[j].Timestamp, samples[j].Value).Err(); err != nil { + return nil, fmt.Errorf("ts.add %s: %w", key, err) + } + } + result[key] = samples + } + + return result, nil +} + +// parseTSSample parses a kvrocks TS reply sample: an array of [timestamp, value] +func parseTSSample(raw interface{}) (uint64, float64, error) { + s, ok := raw.([]interface{}) + if !ok || len(s) != 2 { + return 0, 0, fmt.Errorf("unexpected TS sample reply: %v", raw) + } + ts, err := strconv.ParseUint(fmt.Sprint(s[0]), 10, 64) + if err != nil { + return 0, 0, fmt.Errorf("bad TS timestamp %v: %w", s[0], err) + } + val, err := strconv.ParseFloat(fmt.Sprint(s[1]), 64) + if err != nil { + return 0, 0, fmt.Errorf("bad TS value %v: %w", s[1], err) + } + return ts, val, nil +} + +// verifyTimeSeries checks TS.GET (last sample) and TS.RANGE (all samples) +func verifyTimeSeries(ctx context.Context, client *redis.Client, stored map[string][]TSSample, r *SeededRandom) error { + for i := 0; i < 3; i++ { + key := fmt.Sprintf("ts:%d", i) + sampleCount := 20 + r.Int(30) + base := uint64(1000000000 + i*1000000) + expected := make([]TSSample, sampleCount) + for j := 0; j < sampleCount; j++ { + expected[j] = TSSample{Timestamp: base + uint64(j*1000), Value: r.Float64() * 1000} + } + + // TS.GET returns the last sample + got, err := client.Do(ctx, "TS.GET", key).Result() + if err != nil { + return fmt.Errorf("ts.get %s: %w", key, err) + } + lastTs, lastVal, err := parseTSSample(got.([]interface{})[0]) + if err != nil { + return fmt.Errorf("ts.get %s: %w", key, err) + } + last := TSSample{Timestamp: lastTs, Value: lastVal} + expLast := expected[sampleCount-1] + if last.Timestamp != expLast.Timestamp || math.Abs(last.Value-expLast.Value) > tsValueTolerance { + return fmt.Errorf("ts %s: last sample (%d, %f), expected (%d, %f)", + key, last.Timestamp, last.Value, expLast.Timestamp, expLast.Value) + } + + // TS.RANGE returns all samples in order + rangeReply, err := client.Do(ctx, "TS.RANGE", key, 0, uint64(1)<<62).Result() + if err != nil { + return fmt.Errorf("ts.range %s: %w", key, err) + } + samples, ok := rangeReply.([]interface{}) + if !ok || len(samples) != sampleCount { + return fmt.Errorf("ts %s: range returned %v samples, expected %d", key, len(samples), sampleCount) + } + for j := 0; j < sampleCount; j++ { + ts, val, err := parseTSSample(samples[j]) + if err != nil { + return fmt.Errorf("ts.range %s[%d]: %w", key, j, err) + } + if ts != expected[j].Timestamp || math.Abs(val-expected[j].Value) > tsValueTolerance { + return fmt.Errorf("ts %s[%d]: (%d, %f), expected (%d, %f)", + key, j, ts, val, expected[j].Timestamp, expected[j].Value) + } + } + } + return nil +}