mirror of
https://github.com/neondatabase/neon.git
synced 2026-08-15 18:48:38 +00:00
Compare commits
33 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 781897b983 | |||
| e520293090 | |||
| 241e549757 | |||
| 34bea270f0 | |||
| 13f0e7a5b4 | |||
| 3e35f10adc | |||
| 3be3bb7730 | |||
| 01d2c52c82 | |||
| 9f79e7edea | |||
| a22165d41e | |||
| 725be60bb7 | |||
| e516c376d6 | |||
| 8e51c27e1a | |||
| 9e1eb69d55 | |||
| 687ba81366 | |||
| 47bae68a2e | |||
| e8b195acb7 | |||
| 254cb7dc4f | |||
| ed85d97f17 | |||
| 4a216c5f7f | |||
| c5a428a61a | |||
| ff8c481777 | |||
| f25dd75be9 | |||
| b99bed510d | |||
| 580584c8fc | |||
| d823e84ed5 | |||
| 231dfbaed6 | |||
| 5cf53786f9 | |||
| 9b9bbad462 | |||
| 537b2c1ae6 | |||
| 31123d1fa8 | |||
| 4f2ac51bdd | |||
| 7b2f9dc908 |
@@ -23,6 +23,7 @@ docker cp ${ID}:/data/postgres_install.tar.gz .
|
|||||||
tar -xzf postgres_install.tar.gz -C neon_install
|
tar -xzf postgres_install.tar.gz -C neon_install
|
||||||
mkdir neon_install/bin/
|
mkdir neon_install/bin/
|
||||||
docker cp ${ID}:/usr/local/bin/pageserver neon_install/bin/
|
docker cp ${ID}:/usr/local/bin/pageserver neon_install/bin/
|
||||||
|
docker cp ${ID}:/usr/local/bin/pageserver_binutils neon_install/bin/
|
||||||
docker cp ${ID}:/usr/local/bin/safekeeper neon_install/bin/
|
docker cp ${ID}:/usr/local/bin/safekeeper neon_install/bin/
|
||||||
docker cp ${ID}:/usr/local/bin/proxy neon_install/bin/
|
docker cp ${ID}:/usr/local/bin/proxy neon_install/bin/
|
||||||
docker cp ${ID}:/usr/local/v14/bin/ neon_install/v14/bin/
|
docker cp ${ID}:/usr/local/v14/bin/ neon_install/v14/bin/
|
||||||
|
|||||||
@@ -16,5 +16,5 @@ env_name = neon-stress
|
|||||||
console_mgmt_base_url = http://neon-stress-console.local
|
console_mgmt_base_url = http://neon-stress-console.local
|
||||||
bucket_name = neon-storage-ireland
|
bucket_name = neon-storage-ireland
|
||||||
bucket_region = eu-west-1
|
bucket_region = eu-west-1
|
||||||
etcd_endpoints = etcd-stress.local:2379
|
etcd_endpoints = neon-stress-etcd.local:2379
|
||||||
safekeeper_enable_s3_offload = false
|
safekeeper_enable_s3_offload = false
|
||||||
|
|||||||
@@ -46,7 +46,7 @@ jobs:
|
|||||||
runs-on: [self-hosted, zenith-benchmarker]
|
runs-on: [self-hosted, zenith-benchmarker]
|
||||||
|
|
||||||
env:
|
env:
|
||||||
POSTGRES_DISTRIB_DIR: /tmp/pg_install
|
POSTGRES_DISTRIB_DIR: /usr/pgsql
|
||||||
DEFAULT_PG_VERSION: 14
|
DEFAULT_PG_VERSION: 14
|
||||||
|
|
||||||
steps:
|
steps:
|
||||||
@@ -149,6 +149,7 @@ jobs:
|
|||||||
|
|
||||||
strategy:
|
strategy:
|
||||||
fail-fast: false
|
fail-fast: false
|
||||||
|
max-parallel: 1
|
||||||
matrix:
|
matrix:
|
||||||
# neon-captest-new: Run pgbench in a freshly created project
|
# neon-captest-new: Run pgbench in a freshly created project
|
||||||
# neon-captest-reuse: Same, but reusing existing project
|
# neon-captest-reuse: Same, but reusing existing project
|
||||||
|
|||||||
@@ -127,8 +127,8 @@ jobs:
|
|||||||
target/
|
target/
|
||||||
# Fall back to older versions of the key, if no cache for current Cargo.lock was found
|
# Fall back to older versions of the key, if no cache for current Cargo.lock was found
|
||||||
key: |
|
key: |
|
||||||
v8-${{ runner.os }}-${{ matrix.build_type }}-cargo-${{ hashFiles('Cargo.lock') }}
|
v9-${{ runner.os }}-${{ matrix.build_type }}-cargo-${{ hashFiles('Cargo.lock') }}
|
||||||
v8-${{ runner.os }}-${{ matrix.build_type }}-cargo-
|
v9-${{ runner.os }}-${{ matrix.build_type }}-cargo-
|
||||||
|
|
||||||
- name: Cache postgres v14 build
|
- name: Cache postgres v14 build
|
||||||
id: cache_pg_14
|
id: cache_pg_14
|
||||||
@@ -389,7 +389,7 @@ jobs:
|
|||||||
!~/.cargo/registry/src
|
!~/.cargo/registry/src
|
||||||
~/.cargo/git/
|
~/.cargo/git/
|
||||||
target/
|
target/
|
||||||
key: v8-${{ runner.os }}-${{ matrix.build_type }}-cargo-${{ hashFiles('Cargo.lock') }}
|
key: v9-${{ runner.os }}-${{ matrix.build_type }}-cargo-${{ hashFiles('Cargo.lock') }}
|
||||||
|
|
||||||
- name: Get Neon artifact
|
- name: Get Neon artifact
|
||||||
uses: ./.github/actions/download
|
uses: ./.github/actions/download
|
||||||
@@ -494,7 +494,7 @@ jobs:
|
|||||||
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
||||||
|
|
||||||
- name: Kaniko build neon
|
- name: Kaniko build neon
|
||||||
run: /kaniko/executor --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --snapshotMode=redo --context . --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/neon:$GITHUB_RUN_ID
|
run: /kaniko/executor --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --snapshotMode=redo --context . --build-arg GIT_VERSION=${{ github.sha }} --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/neon:$GITHUB_RUN_ID
|
||||||
|
|
||||||
compute-tools-image:
|
compute-tools-image:
|
||||||
runs-on: dev
|
runs-on: dev
|
||||||
@@ -508,7 +508,7 @@ jobs:
|
|||||||
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
||||||
|
|
||||||
- name: Kaniko build compute tools
|
- name: Kaniko build compute tools
|
||||||
run: /kaniko/executor --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --snapshotMode=redo --context . --dockerfile Dockerfile.compute-tools --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-tools:$GITHUB_RUN_ID
|
run: /kaniko/executor --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --snapshotMode=redo --context . --build-arg GIT_VERSION=${{ github.sha }} --dockerfile Dockerfile.compute-tools --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-tools:$GITHUB_RUN_ID
|
||||||
|
|
||||||
compute-node-image:
|
compute-node-image:
|
||||||
runs-on: dev
|
runs-on: dev
|
||||||
@@ -527,7 +527,7 @@ jobs:
|
|||||||
# cloud repo depends on this image name, thus duplicating it
|
# cloud repo depends on this image name, thus duplicating it
|
||||||
# remove compute-node when cloud repo is updated
|
# remove compute-node when cloud repo is updated
|
||||||
- name: Kaniko build compute node with extensions v14 (compatibility)
|
- name: Kaniko build compute node with extensions v14 (compatibility)
|
||||||
run: /kaniko/executor --skip-unused-stages --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --snapshotMode=redo --context . --dockerfile Dockerfile.compute-node-v14 --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node:$GITHUB_RUN_ID
|
run: /kaniko/executor --skip-unused-stages --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --snapshotMode=redo --context . --build-arg GIT_VERSION=${{ github.sha }} --dockerfile Dockerfile.compute-node-v14 --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node:$GITHUB_RUN_ID
|
||||||
|
|
||||||
compute-node-image-v14:
|
compute-node-image-v14:
|
||||||
runs-on: dev
|
runs-on: dev
|
||||||
@@ -543,7 +543,7 @@ jobs:
|
|||||||
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
||||||
|
|
||||||
- name: Kaniko build compute node with extensions v14
|
- name: Kaniko build compute node with extensions v14
|
||||||
run: /kaniko/executor --skip-unused-stages --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --context . --dockerfile Dockerfile.compute-node-v14 --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node-v14:$GITHUB_RUN_ID
|
run: /kaniko/executor --skip-unused-stages --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --context . --build-arg GIT_VERSION=${{ github.sha }} --dockerfile Dockerfile.compute-node-v14 --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node-v14:$GITHUB_RUN_ID
|
||||||
|
|
||||||
|
|
||||||
compute-node-image-v15:
|
compute-node-image-v15:
|
||||||
@@ -560,11 +560,11 @@ jobs:
|
|||||||
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
run: echo "{\"credsStore\":\"ecr-login\"}" > /kaniko/.docker/config.json
|
||||||
|
|
||||||
- name: Kaniko build compute node with extensions v15
|
- name: Kaniko build compute node with extensions v15
|
||||||
run: /kaniko/executor --skip-unused-stages --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --context . --dockerfile Dockerfile.compute-node-v15 --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node-v15:$GITHUB_RUN_ID
|
run: /kaniko/executor --skip-unused-stages --snapshotMode=redo --cache=true --cache-repo 369495373322.dkr.ecr.eu-central-1.amazonaws.com/cache --context . --build-arg GIT_VERSION=${{ github.sha }} --dockerfile Dockerfile.compute-node-v15 --destination 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node-v15:$GITHUB_RUN_ID
|
||||||
|
|
||||||
promote-images:
|
promote-images:
|
||||||
runs-on: dev
|
runs-on: dev
|
||||||
needs: [ neon-image, compute-node-image, compute-node-image-v14, compute-tools-image ]
|
needs: [ neon-image, compute-node-image, compute-node-image-v14, compute-node-image-v15, compute-tools-image ]
|
||||||
if: github.event_name != 'workflow_dispatch'
|
if: github.event_name != 'workflow_dispatch'
|
||||||
container: amazon/aws-cli
|
container: amazon/aws-cli
|
||||||
strategy:
|
strategy:
|
||||||
@@ -573,7 +573,7 @@ jobs:
|
|||||||
# compute-node uses postgres 14, which is default now
|
# compute-node uses postgres 14, which is default now
|
||||||
# cloud repo depends on this image name, thus duplicating it
|
# cloud repo depends on this image name, thus duplicating it
|
||||||
# remove compute-node when cloud repo is updated
|
# remove compute-node when cloud repo is updated
|
||||||
name: [ neon, compute-node, compute-node-v14, compute-tools ]
|
name: [ neon, compute-node, compute-node-v14, compute-node-v15, compute-tools ]
|
||||||
|
|
||||||
steps:
|
steps:
|
||||||
- name: Promote image to latest
|
- name: Promote image to latest
|
||||||
@@ -608,6 +608,9 @@ jobs:
|
|||||||
- name: Pull compute node v14 image from ECR
|
- name: Pull compute node v14 image from ECR
|
||||||
run: crane pull 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node-v14:latest compute-node-v14
|
run: crane pull 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node-v14:latest compute-node-v14
|
||||||
|
|
||||||
|
- name: Pull compute node v15 image from ECR
|
||||||
|
run: crane pull 369495373322.dkr.ecr.eu-central-1.amazonaws.com/compute-node-v15:latest compute-node-v15
|
||||||
|
|
||||||
- name: Pull rust image from ECR
|
- name: Pull rust image from ECR
|
||||||
run: crane pull 369495373322.dkr.ecr.eu-central-1.amazonaws.com/rust:pinned rust
|
run: crane pull 369495373322.dkr.ecr.eu-central-1.amazonaws.com/rust:pinned rust
|
||||||
|
|
||||||
@@ -638,6 +641,9 @@ jobs:
|
|||||||
- name: Push compute node v14 image to Docker Hub
|
- name: Push compute node v14 image to Docker Hub
|
||||||
run: crane push compute-node-v14 neondatabase/compute-node-v14:${{needs.tag.outputs.build-tag}}
|
run: crane push compute-node-v14 neondatabase/compute-node-v14:${{needs.tag.outputs.build-tag}}
|
||||||
|
|
||||||
|
- name: Push compute node v15 image to Docker Hub
|
||||||
|
run: crane push compute-node-v15 neondatabase/compute-node-v15:${{needs.tag.outputs.build-tag}}
|
||||||
|
|
||||||
- name: Push rust image to Docker Hub
|
- name: Push rust image to Docker Hub
|
||||||
run: crane push rust neondatabase/rust:pinned
|
run: crane push rust neondatabase/rust:pinned
|
||||||
|
|
||||||
@@ -650,6 +656,7 @@ jobs:
|
|||||||
crane tag neondatabase/compute-tools:${{needs.tag.outputs.build-tag}} latest
|
crane tag neondatabase/compute-tools:${{needs.tag.outputs.build-tag}} latest
|
||||||
crane tag neondatabase/compute-node:${{needs.tag.outputs.build-tag}} latest
|
crane tag neondatabase/compute-node:${{needs.tag.outputs.build-tag}} latest
|
||||||
crane tag neondatabase/compute-node-v14:${{needs.tag.outputs.build-tag}} latest
|
crane tag neondatabase/compute-node-v14:${{needs.tag.outputs.build-tag}} latest
|
||||||
|
crane tag neondatabase/compute-node-v15:${{needs.tag.outputs.build-tag}} latest
|
||||||
|
|
||||||
calculate-deploy-targets:
|
calculate-deploy-targets:
|
||||||
runs-on: [ self-hosted, Linux, k8s-runner ]
|
runs-on: [ self-hosted, Linux, k8s-runner ]
|
||||||
@@ -768,5 +775,5 @@ jobs:
|
|||||||
- name: Re-deploy proxy
|
- name: Re-deploy proxy
|
||||||
run: |
|
run: |
|
||||||
DOCKER_TAG=${{needs.tag.outputs.build-tag}}
|
DOCKER_TAG=${{needs.tag.outputs.build-tag}}
|
||||||
helm upgrade ${{ matrix.proxy_job }} neondatabase/neon-proxy --namespace default --install -f .github/helm-values/${{ matrix.proxy_config }}.yaml --set image.tag=${DOCKER_TAG} --wait --timeout 15m0s
|
helm upgrade ${{ matrix.proxy_job }} neondatabase/neon-proxy --namespace neon-proxy --install -f .github/helm-values/${{ matrix.proxy_config }}.yaml --set image.tag=${DOCKER_TAG} --wait --timeout 15m0s
|
||||||
helm upgrade ${{ matrix.proxy_job }}-scram neondatabase/neon-proxy --namespace default --install -f .github/helm-values/${{ matrix.proxy_config }}-scram.yaml --set image.tag=${DOCKER_TAG} --wait --timeout 15m0s
|
helm upgrade ${{ matrix.proxy_job }}-scram neondatabase/neon-proxy --namespace neon-proxy --install -f .github/helm-values/${{ matrix.proxy_config }}-scram.yaml --set image.tag=${DOCKER_TAG} --wait --timeout 15m0s
|
||||||
|
|||||||
@@ -106,7 +106,7 @@ jobs:
|
|||||||
!~/.cargo/registry/src
|
!~/.cargo/registry/src
|
||||||
~/.cargo/git
|
~/.cargo/git
|
||||||
target
|
target
|
||||||
key: v4-${{ runner.os }}-cargo-${{ hashFiles('./Cargo.lock') }}-rust
|
key: v5-${{ runner.os }}-cargo-${{ hashFiles('./Cargo.lock') }}-rust
|
||||||
|
|
||||||
- name: Run cargo clippy
|
- name: Run cargo clippy
|
||||||
run: ./run_clippy.sh
|
run: ./run_clippy.sh
|
||||||
|
|||||||
Generated
+99
-2
@@ -497,8 +497,10 @@ dependencies = [
|
|||||||
"chrono",
|
"chrono",
|
||||||
"clap 3.2.16",
|
"clap 3.2.16",
|
||||||
"env_logger",
|
"env_logger",
|
||||||
|
"futures",
|
||||||
"hyper",
|
"hyper",
|
||||||
"log",
|
"log",
|
||||||
|
"notify",
|
||||||
"postgres",
|
"postgres",
|
||||||
"regex",
|
"regex",
|
||||||
"serde",
|
"serde",
|
||||||
@@ -540,11 +542,11 @@ dependencies = [
|
|||||||
"git-version",
|
"git-version",
|
||||||
"nix",
|
"nix",
|
||||||
"once_cell",
|
"once_cell",
|
||||||
"pageserver",
|
"pageserver_api",
|
||||||
"postgres",
|
"postgres",
|
||||||
"regex",
|
"regex",
|
||||||
"reqwest",
|
"reqwest",
|
||||||
"safekeeper",
|
"safekeeper_api",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_with",
|
"serde_with",
|
||||||
"tar",
|
"tar",
|
||||||
@@ -1072,6 +1074,15 @@ dependencies = [
|
|||||||
"winapi",
|
"winapi",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "fsevent-sys"
|
||||||
|
version = "4.1.0"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "76ee7a02da4d231650c7cea31349b889be2f45ddb3ef3032d2ec8185f6313fd2"
|
||||||
|
dependencies = [
|
||||||
|
"libc",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "futures"
|
name = "futures"
|
||||||
version = "0.3.21"
|
version = "0.3.21"
|
||||||
@@ -1493,6 +1504,26 @@ dependencies = [
|
|||||||
"str_stack",
|
"str_stack",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "inotify"
|
||||||
|
version = "0.9.6"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "f8069d3ec154eb856955c1c0fbffefbf5f3c40a104ec912d4797314c1801abff"
|
||||||
|
dependencies = [
|
||||||
|
"bitflags",
|
||||||
|
"inotify-sys",
|
||||||
|
"libc",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "inotify-sys"
|
||||||
|
version = "0.1.5"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "e05c02b5e89bff3b946cedeca278abc628fe811e604f027c45a8aa3cf793d0eb"
|
||||||
|
dependencies = [
|
||||||
|
"libc",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "instant"
|
name = "instant"
|
||||||
version = "0.1.12"
|
version = "0.1.12"
|
||||||
@@ -1552,6 +1583,26 @@ dependencies = [
|
|||||||
"simple_asn1",
|
"simple_asn1",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "kqueue"
|
||||||
|
version = "1.0.6"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "4d6112e8f37b59803ac47a42d14f1f3a59bbf72fc6857ffc5be455e28a691f8e"
|
||||||
|
dependencies = [
|
||||||
|
"kqueue-sys",
|
||||||
|
"libc",
|
||||||
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "kqueue-sys"
|
||||||
|
version = "1.0.3"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "8367585489f01bc55dd27404dcf56b95e6da061a256a666ab23be9ba96a2e587"
|
||||||
|
dependencies = [
|
||||||
|
"bitflags",
|
||||||
|
"libc",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "kstring"
|
name = "kstring"
|
||||||
version = "1.0.6"
|
version = "1.0.6"
|
||||||
@@ -1797,6 +1848,24 @@ dependencies = [
|
|||||||
"minimal-lexical",
|
"minimal-lexical",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "notify"
|
||||||
|
version = "5.0.0"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "ed2c66da08abae1c024c01d635253e402341b4060a12e99b31c7594063bf490a"
|
||||||
|
dependencies = [
|
||||||
|
"bitflags",
|
||||||
|
"crossbeam-channel",
|
||||||
|
"filetime",
|
||||||
|
"fsevent-sys",
|
||||||
|
"inotify",
|
||||||
|
"kqueue",
|
||||||
|
"libc",
|
||||||
|
"mio",
|
||||||
|
"walkdir",
|
||||||
|
"winapi",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "num-bigint"
|
name = "num-bigint"
|
||||||
version = "0.4.3"
|
version = "0.4.3"
|
||||||
@@ -1975,6 +2044,7 @@ dependencies = [
|
|||||||
"nix",
|
"nix",
|
||||||
"num-traits",
|
"num-traits",
|
||||||
"once_cell",
|
"once_cell",
|
||||||
|
"pageserver_api",
|
||||||
"postgres",
|
"postgres",
|
||||||
"postgres-protocol",
|
"postgres-protocol",
|
||||||
"postgres-types",
|
"postgres-types",
|
||||||
@@ -2003,6 +2073,17 @@ dependencies = [
|
|||||||
"workspace_hack",
|
"workspace_hack",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "pageserver_api"
|
||||||
|
version = "0.1.0"
|
||||||
|
dependencies = [
|
||||||
|
"const_format",
|
||||||
|
"serde",
|
||||||
|
"serde_with",
|
||||||
|
"utils",
|
||||||
|
"workspace_hack",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "parking_lot"
|
name = "parking_lot"
|
||||||
version = "0.11.2"
|
version = "0.11.2"
|
||||||
@@ -2371,6 +2452,7 @@ version = "0.1.0"
|
|||||||
dependencies = [
|
dependencies = [
|
||||||
"anyhow",
|
"anyhow",
|
||||||
"async-trait",
|
"async-trait",
|
||||||
|
"atty",
|
||||||
"base64",
|
"base64",
|
||||||
"bstr",
|
"bstr",
|
||||||
"bytes",
|
"bytes",
|
||||||
@@ -2404,6 +2486,8 @@ dependencies = [
|
|||||||
"tokio-postgres",
|
"tokio-postgres",
|
||||||
"tokio-postgres-rustls",
|
"tokio-postgres-rustls",
|
||||||
"tokio-rustls",
|
"tokio-rustls",
|
||||||
|
"tracing",
|
||||||
|
"tracing-subscriber",
|
||||||
"url",
|
"url",
|
||||||
"utils",
|
"utils",
|
||||||
"uuid",
|
"uuid",
|
||||||
@@ -2891,6 +2975,7 @@ dependencies = [
|
|||||||
"postgres_ffi",
|
"postgres_ffi",
|
||||||
"regex",
|
"regex",
|
||||||
"remote_storage",
|
"remote_storage",
|
||||||
|
"safekeeper_api",
|
||||||
"serde",
|
"serde",
|
||||||
"serde_json",
|
"serde_json",
|
||||||
"serde_with",
|
"serde_with",
|
||||||
@@ -2906,6 +2991,17 @@ dependencies = [
|
|||||||
"workspace_hack",
|
"workspace_hack",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "safekeeper_api"
|
||||||
|
version = "0.1.0"
|
||||||
|
dependencies = [
|
||||||
|
"const_format",
|
||||||
|
"serde",
|
||||||
|
"serde_with",
|
||||||
|
"utils",
|
||||||
|
"workspace_hack",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "same-file"
|
name = "same-file"
|
||||||
version = "1.0.6"
|
version = "1.0.6"
|
||||||
@@ -4142,6 +4238,7 @@ dependencies = [
|
|||||||
"bstr",
|
"bstr",
|
||||||
"bytes",
|
"bytes",
|
||||||
"chrono",
|
"chrono",
|
||||||
|
"crossbeam-utils",
|
||||||
"either",
|
"either",
|
||||||
"fail",
|
"fail",
|
||||||
"hashbrown",
|
"hashbrown",
|
||||||
|
|||||||
+5
-5
@@ -44,7 +44,7 @@ COPY . .
|
|||||||
# Show build caching stats to check if it was used in the end.
|
# Show build caching stats to check if it was used in the end.
|
||||||
# Has to be the part of the same RUN since cachepot daemon is killed in the end of this RUN, losing the compilation stats.
|
# Has to be the part of the same RUN since cachepot daemon is killed in the end of this RUN, losing the compilation stats.
|
||||||
RUN set -e \
|
RUN set -e \
|
||||||
&& mold -run cargo build --bin pageserver --bin safekeeper --bin proxy --locked --release \
|
&& mold -run cargo build --bin pageserver --bin pageserver_binutils --bin safekeeper --bin proxy --locked --release \
|
||||||
&& cachepot -s
|
&& cachepot -s
|
||||||
|
|
||||||
# Build final image
|
# Build final image
|
||||||
@@ -63,9 +63,10 @@ RUN set -e \
|
|||||||
&& useradd -d /data neon \
|
&& useradd -d /data neon \
|
||||||
&& chown -R neon:neon /data
|
&& chown -R neon:neon /data
|
||||||
|
|
||||||
COPY --from=build --chown=neon:neon /home/nonroot/target/release/pageserver /usr/local/bin
|
COPY --from=build --chown=neon:neon /home/nonroot/target/release/pageserver /usr/local/bin
|
||||||
COPY --from=build --chown=neon:neon /home/nonroot/target/release/safekeeper /usr/local/bin
|
COPY --from=build --chown=neon:neon /home/nonroot/target/release/pageserver_binutils /usr/local/bin
|
||||||
COPY --from=build --chown=neon:neon /home/nonroot/target/release/proxy /usr/local/bin
|
COPY --from=build --chown=neon:neon /home/nonroot/target/release/safekeeper /usr/local/bin
|
||||||
|
COPY --from=build --chown=neon:neon /home/nonroot/target/release/proxy /usr/local/bin
|
||||||
|
|
||||||
COPY --from=pg-build /home/nonroot/pg_install/v14 /usr/local/v14/
|
COPY --from=pg-build /home/nonroot/pg_install/v14 /usr/local/v14/
|
||||||
COPY --from=pg-build /home/nonroot/pg_install/v15 /usr/local/v15/
|
COPY --from=pg-build /home/nonroot/pg_install/v15 /usr/local/v15/
|
||||||
@@ -85,4 +86,3 @@ VOLUME ["/data"]
|
|||||||
USER neon
|
USER neon
|
||||||
EXPOSE 6400
|
EXPOSE 6400
|
||||||
EXPOSE 9898
|
EXPOSE 9898
|
||||||
CMD ["/bin/bash"]
|
|
||||||
|
|||||||
+46
-13
@@ -5,7 +5,7 @@
|
|||||||
|
|
||||||
ARG TAG=pinned
|
ARG TAG=pinned
|
||||||
# apparently, ARGs don't get replaced in RUN commands in kaniko
|
# apparently, ARGs don't get replaced in RUN commands in kaniko
|
||||||
# ARG POSTGIS_VERSION=3.3.0
|
# ARG POSTGIS_VERSION=3.3.1
|
||||||
# ARG PLV8_VERSION=3.1.4
|
# ARG PLV8_VERSION=3.1.4
|
||||||
# ARG PG_VERSION=v15
|
# ARG PG_VERSION=v15
|
||||||
|
|
||||||
@@ -13,9 +13,12 @@ ARG TAG=pinned
|
|||||||
# Layer "build-deps"
|
# Layer "build-deps"
|
||||||
#
|
#
|
||||||
FROM debian:bullseye-slim AS build-deps
|
FROM debian:bullseye-slim AS build-deps
|
||||||
|
RUN echo "deb http://ftp.debian.org/debian testing main" >> /etc/apt/sources.list && \
|
||||||
|
echo "APT::Default-Release \"stable\";" > /etc/apt/apt.conf.d/default-release && \
|
||||||
|
apt update
|
||||||
RUN apt update && \
|
RUN apt update && \
|
||||||
apt install -y git autoconf automake libtool build-essential bison flex libreadline-dev zlib1g-dev libxml2-dev \
|
apt install -y git autoconf automake libtool build-essential bison flex libreadline-dev zlib1g-dev libxml2-dev \
|
||||||
libcurl4-openssl-dev libossp-uuid-dev
|
libcurl4-openssl-dev libossp-uuid-dev wget pkg-config libglib2.0-dev
|
||||||
|
|
||||||
#
|
#
|
||||||
# Layer "pg-build"
|
# Layer "pg-build"
|
||||||
@@ -42,11 +45,11 @@ RUN cd postgres && \
|
|||||||
FROM build-deps AS postgis-build
|
FROM build-deps AS postgis-build
|
||||||
COPY --from=pg-build /usr/local/pgsql/ /usr/local/pgsql/
|
COPY --from=pg-build /usr/local/pgsql/ /usr/local/pgsql/
|
||||||
RUN apt update && \
|
RUN apt update && \
|
||||||
apt install -y gdal-bin libgdal-dev libprotobuf-c-dev protobuf-c-compiler xsltproc wget
|
apt install -y gdal-bin libgdal-dev libprotobuf-c-dev protobuf-c-compiler xsltproc
|
||||||
|
|
||||||
RUN wget https://download.osgeo.org/postgis/source/postgis-3.3.0.tar.gz && \
|
RUN wget https://download.osgeo.org/postgis/source/postgis-3.3.1.tar.gz && \
|
||||||
tar xvzf postgis-3.3.0.tar.gz && \
|
tar xvzf postgis-3.3.1.tar.gz && \
|
||||||
cd postgis-3.3.0 && \
|
cd postgis-3.3.1 && \
|
||||||
./autogen.sh && \
|
./autogen.sh && \
|
||||||
export PATH="/usr/local/pgsql/bin:$PATH" && \
|
export PATH="/usr/local/pgsql/bin:$PATH" && \
|
||||||
./configure && \
|
./configure && \
|
||||||
@@ -64,15 +67,13 @@ RUN wget https://download.osgeo.org/postgis/source/postgis-3.3.0.tar.gz && \
|
|||||||
# Build plv8
|
# Build plv8
|
||||||
#
|
#
|
||||||
FROM build-deps AS plv8-build
|
FROM build-deps AS plv8-build
|
||||||
COPY --from=postgis-build /usr/local/pgsql/ /usr/local/pgsql/
|
COPY --from=pg-build /usr/local/pgsql/ /usr/local/pgsql/
|
||||||
RUN apt update && \
|
RUN apt update && \
|
||||||
apt install -y git curl wget make ninja-build build-essential libncurses5 python3-dev pkg-config libc++-dev libc++abi-dev libglib2.0-dev
|
apt install -y ninja-build python3-dev libc++-dev libc++abi-dev libncurses5
|
||||||
|
|
||||||
# https://github.com/plv8/plv8/issues/475
|
# https://github.com/plv8/plv8/issues/475
|
||||||
# Debian bullseye provides binutils 2.35 when >= 2.38 is necessary
|
# Debian bullseye provides binutils 2.35 when >= 2.38 is necessary
|
||||||
RUN echo "deb http://ftp.debian.org/debian testing main" >> /etc/apt/sources.list && \
|
RUN apt update && \
|
||||||
echo "APT::Default-Release \"stable\";" > /etc/apt/apt.conf.d/default-release && \
|
|
||||||
apt update && \
|
|
||||||
apt install -y --no-install-recommends -t testing binutils
|
apt install -y --no-install-recommends -t testing binutils
|
||||||
|
|
||||||
RUN wget https://github.com/plv8/plv8/archive/refs/tags/v3.1.4.tar.gz && \
|
RUN wget https://github.com/plv8/plv8/archive/refs/tags/v3.1.4.tar.gz && \
|
||||||
@@ -84,12 +85,46 @@ RUN wget https://github.com/plv8/plv8/archive/refs/tags/v3.1.4.tar.gz && \
|
|||||||
rm -rf /plv8-* && \
|
rm -rf /plv8-* && \
|
||||||
echo 'trusted = true' >> /usr/local/pgsql/share/extension/plv8.control
|
echo 'trusted = true' >> /usr/local/pgsql/share/extension/plv8.control
|
||||||
|
|
||||||
|
#
|
||||||
|
# Layer "h3-pg-build"
|
||||||
|
# Build h3_pg
|
||||||
|
#
|
||||||
|
FROM build-deps AS h3-pg-build
|
||||||
|
COPY --from=pg-build /usr/local/pgsql/ /usr/local/pgsql/
|
||||||
|
|
||||||
|
# packaged cmake is too old
|
||||||
|
RUN apt update && \
|
||||||
|
apt install -y --no-install-recommends -t testing cmake
|
||||||
|
|
||||||
|
RUN wget https://github.com/uber/h3/archive/refs/tags/v4.0.1.tar.gz -O h3.tgz && \
|
||||||
|
tar xvzf h3.tgz && \
|
||||||
|
cd h3-4.0.1 && \
|
||||||
|
mkdir build && \
|
||||||
|
cd build && \
|
||||||
|
cmake .. -DCMAKE_BUILD_TYPE=Release && \
|
||||||
|
make -j $(getconf _NPROCESSORS_ONLN) && \
|
||||||
|
DESTDIR=/h3 make install && \
|
||||||
|
cp -R /h3/usr / && \
|
||||||
|
rm -rf build
|
||||||
|
|
||||||
|
RUN wget https://github.com/zachasme/h3-pg/archive/refs/tags/v4.0.1.tar.gz -O h3-pg.tgz && \
|
||||||
|
tar xvzf h3-pg.tgz && \
|
||||||
|
cd h3-pg-4.0.1 && \
|
||||||
|
export PATH="/usr/local/pgsql/bin:$PATH" && \
|
||||||
|
make -j $(getconf _NPROCESSORS_ONLN) && \
|
||||||
|
make -j $(getconf _NPROCESSORS_ONLN) install && \
|
||||||
|
echo 'trusted = true' >> /usr/local/pgsql/share/extension/h3.control
|
||||||
|
|
||||||
#
|
#
|
||||||
# Layer "neon-pg-ext-build"
|
# Layer "neon-pg-ext-build"
|
||||||
# compile neon extensions
|
# compile neon extensions
|
||||||
#
|
#
|
||||||
FROM build-deps AS neon-pg-ext-build
|
FROM build-deps AS neon-pg-ext-build
|
||||||
COPY --from=postgis-build /usr/local/pgsql/ /usr/local/pgsql/
|
COPY --from=postgis-build /usr/local/pgsql/ /usr/local/pgsql/
|
||||||
|
# plv8 still sometimes crashes during the creation
|
||||||
|
# COPY --from=plv8-build /usr/local/pgsql/ /usr/local/pgsql/
|
||||||
|
COPY --from=h3-pg-build /usr/local/pgsql/ /usr/local/pgsql/
|
||||||
|
COPY --from=h3-pg-build /h3/usr /
|
||||||
COPY pgxn/ pgxn/
|
COPY pgxn/ pgxn/
|
||||||
|
|
||||||
RUN make -j $(getconf _NPROCESSORS_ONLN) \
|
RUN make -j $(getconf _NPROCESSORS_ONLN) \
|
||||||
@@ -137,8 +172,6 @@ RUN mkdir /var/db && useradd -m -d /var/db/postgres postgres && \
|
|||||||
chmod 0750 /var/db/postgres/compute && \
|
chmod 0750 /var/db/postgres/compute && \
|
||||||
echo '/usr/local/lib' >> /etc/ld.so.conf && /sbin/ldconfig
|
echo '/usr/local/lib' >> /etc/ld.so.conf && /sbin/ldconfig
|
||||||
|
|
||||||
# TODO: Check if we can make the extension setup more modular versus a linear build
|
|
||||||
# currently plv8-build copies the output /usr/local/pgsql from postgis-build, etc#
|
|
||||||
COPY --from=postgres-cleanup-layer --chown=postgres /usr/local/pgsql /usr/local
|
COPY --from=postgres-cleanup-layer --chown=postgres /usr/local/pgsql /usr/local
|
||||||
COPY --from=compute-tools --chown=postgres /home/nonroot/target/release-line-debug-size-lto/compute_ctl /usr/local/bin/compute_ctl
|
COPY --from=compute-tools --chown=postgres /home/nonroot/target/release-line-debug-size-lto/compute_ctl /usr/local/bin/compute_ctl
|
||||||
|
|
||||||
|
|||||||
@@ -8,8 +8,10 @@ anyhow = "1.0"
|
|||||||
chrono = "0.4"
|
chrono = "0.4"
|
||||||
clap = "3.0"
|
clap = "3.0"
|
||||||
env_logger = "0.9"
|
env_logger = "0.9"
|
||||||
|
futures = "0.3.13"
|
||||||
hyper = { version = "0.14", features = ["full"] }
|
hyper = { version = "0.14", features = ["full"] }
|
||||||
log = { version = "0.4", features = ["std", "serde"] }
|
log = { version = "0.4", features = ["std", "serde"] }
|
||||||
|
notify = "5.0.0"
|
||||||
postgres = { git = "https://github.com/neondatabase/rust-postgres.git", rev="d052ee8b86fff9897c77b0fe89ea9daba0e1fa38" }
|
postgres = { git = "https://github.com/neondatabase/rust-postgres.git", rev="d052ee8b86fff9897c77b0fe89ea9daba0e1fa38" }
|
||||||
regex = "1"
|
regex = "1"
|
||||||
serde = { version = "1.0", features = ["derive"] }
|
serde = { version = "1.0", features = ["derive"] }
|
||||||
|
|||||||
@@ -178,7 +178,6 @@ impl ComputeNode {
|
|||||||
.args(&["--sync-safekeepers"])
|
.args(&["--sync-safekeepers"])
|
||||||
.env("PGDATA", &self.pgdata) // we cannot use -D in this mode
|
.env("PGDATA", &self.pgdata) // we cannot use -D in this mode
|
||||||
.stdout(Stdio::piped())
|
.stdout(Stdio::piped())
|
||||||
.stderr(Stdio::piped())
|
|
||||||
.spawn()
|
.spawn()
|
||||||
.expect("postgres --sync-safekeepers failed to start");
|
.expect("postgres --sync-safekeepers failed to start");
|
||||||
|
|
||||||
@@ -191,10 +190,10 @@ impl ComputeNode {
|
|||||||
|
|
||||||
if !sync_output.status.success() {
|
if !sync_output.status.success() {
|
||||||
anyhow::bail!(
|
anyhow::bail!(
|
||||||
"postgres --sync-safekeepers exited with non-zero status: {}. stdout: {}, stderr: {}",
|
"postgres --sync-safekeepers exited with non-zero status: {}. stdout: {}",
|
||||||
sync_output.status,
|
sync_output.status,
|
||||||
String::from_utf8(sync_output.stdout).expect("postgres --sync-safekeepers exited, and stdout is not utf-8"),
|
String::from_utf8(sync_output.stdout)
|
||||||
String::from_utf8(sync_output.stderr).expect("postgres --sync-safekeepers exited, and stderr is not utf-8"),
|
.expect("postgres --sync-safekeepers exited, and stdout is not utf-8"),
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -258,14 +257,7 @@ impl ComputeNode {
|
|||||||
.spawn()
|
.spawn()
|
||||||
.expect("cannot start postgres process");
|
.expect("cannot start postgres process");
|
||||||
|
|
||||||
// Try default Postgres port if it is not provided
|
wait_for_postgres(&mut pg, pgdata_path)?;
|
||||||
let port = self
|
|
||||||
.spec
|
|
||||||
.cluster
|
|
||||||
.settings
|
|
||||||
.find("port")
|
|
||||||
.unwrap_or_else(|| "5432".to_string());
|
|
||||||
wait_for_postgres(&mut pg, &port, pgdata_path)?;
|
|
||||||
|
|
||||||
// If connection fails,
|
// If connection fails,
|
||||||
// it may be the old node with `zenith_admin` superuser.
|
// it may be the old node with `zenith_admin` superuser.
|
||||||
|
|||||||
@@ -1,18 +1,19 @@
|
|||||||
use std::fmt::Write;
|
use std::fmt::Write;
|
||||||
|
use std::fs;
|
||||||
use std::fs::File;
|
use std::fs::File;
|
||||||
use std::io::{BufRead, BufReader};
|
use std::io::{BufRead, BufReader};
|
||||||
use std::net::{SocketAddr, TcpStream};
|
|
||||||
use std::os::unix::fs::PermissionsExt;
|
use std::os::unix::fs::PermissionsExt;
|
||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
use std::process::Child;
|
use std::process::Child;
|
||||||
use std::str::FromStr;
|
use std::time::{Duration, Instant};
|
||||||
use std::{fs, thread, time};
|
|
||||||
|
|
||||||
use anyhow::{bail, Result};
|
use anyhow::{bail, Result};
|
||||||
use postgres::{Client, Transaction};
|
use postgres::{Client, Transaction};
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
|
|
||||||
const POSTGRES_WAIT_TIMEOUT: u64 = 60 * 1000; // milliseconds
|
use notify::{RecursiveMode, Watcher};
|
||||||
|
|
||||||
|
const POSTGRES_WAIT_TIMEOUT: Duration = Duration::from_millis(60 * 1000); // milliseconds
|
||||||
|
|
||||||
/// Rust representation of Postgres role info with only those fields
|
/// Rust representation of Postgres role info with only those fields
|
||||||
/// that matter for us.
|
/// that matter for us.
|
||||||
@@ -230,52 +231,112 @@ pub fn get_existing_dbs(client: &mut Client) -> Result<Vec<Database>> {
|
|||||||
Ok(postgres_dbs)
|
Ok(postgres_dbs)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Wait for Postgres to become ready to accept connections:
|
/// Wait for Postgres to become ready to accept connections. It's ready to
|
||||||
/// - state should be `ready` in the `pgdata/postmaster.pid`
|
/// accept connections when the state-field in `pgdata/postmaster.pid` says
|
||||||
/// - and we should be able to connect to 127.0.0.1:5432
|
/// 'ready'.
|
||||||
pub fn wait_for_postgres(pg: &mut Child, port: &str, pgdata: &Path) -> Result<()> {
|
pub fn wait_for_postgres(pg: &mut Child, pgdata: &Path) -> Result<()> {
|
||||||
let pid_path = pgdata.join("postmaster.pid");
|
let pid_path = pgdata.join("postmaster.pid");
|
||||||
let mut slept: u64 = 0; // ms
|
|
||||||
let pause = time::Duration::from_millis(100);
|
|
||||||
|
|
||||||
let timeout = time::Duration::from_millis(10);
|
// PostgreSQL writes line "ready" to the postmaster.pid file, when it has
|
||||||
let addr = SocketAddr::from_str(&format!("127.0.0.1:{}", port)).unwrap();
|
// completed initialization and is ready to accept connections. We want to
|
||||||
|
// react quickly and perform the rest of our initialization as soon as
|
||||||
|
// PostgreSQL starts accepting connections. Use 'notify' to be notified
|
||||||
|
// whenever the PID file is changed, and whenever it changes, read it to
|
||||||
|
// check if it's now "ready".
|
||||||
|
//
|
||||||
|
// You cannot actually watch a file before it exists, so we first watch the
|
||||||
|
// data directory, and once the postmaster.pid file appears, we switch to
|
||||||
|
// watch the file instead. We also wake up every 100 ms to poll, just in
|
||||||
|
// case we miss some events for some reason. Not strictly necessary, but
|
||||||
|
// better safe than sorry.
|
||||||
|
let (tx, rx) = std::sync::mpsc::channel();
|
||||||
|
let (mut watcher, rx): (Box<dyn Watcher>, _) = match notify::recommended_watcher(move |res| {
|
||||||
|
let _ = tx.send(res);
|
||||||
|
}) {
|
||||||
|
Ok(watcher) => (Box::new(watcher), rx),
|
||||||
|
Err(e) => {
|
||||||
|
match e.kind {
|
||||||
|
notify::ErrorKind::Io(os) if os.raw_os_error() == Some(38) => {
|
||||||
|
// docker on m1 macs does not support recommended_watcher
|
||||||
|
// but return "Function not implemented (os error 38)"
|
||||||
|
// see https://github.com/notify-rs/notify/issues/423
|
||||||
|
let (tx, rx) = std::sync::mpsc::channel();
|
||||||
|
|
||||||
loop {
|
// let's poll it faster than what we check the results for (100ms)
|
||||||
// Sleep POSTGRES_WAIT_TIMEOUT at max (a bit longer actually if consider a TCP timeout,
|
let config =
|
||||||
// but postgres starts listening almost immediately, even if it is not really
|
notify::Config::default().with_poll_interval(Duration::from_millis(50));
|
||||||
// ready to accept connections).
|
|
||||||
if slept >= POSTGRES_WAIT_TIMEOUT {
|
let watcher = notify::PollWatcher::new(
|
||||||
bail!("timed out while waiting for Postgres to start");
|
move |res| {
|
||||||
|
let _ = tx.send(res);
|
||||||
|
},
|
||||||
|
config,
|
||||||
|
)?;
|
||||||
|
|
||||||
|
(Box::new(watcher), rx)
|
||||||
|
}
|
||||||
|
_ => return Err(e.into()),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
watcher.watch(pgdata, RecursiveMode::NonRecursive)?;
|
||||||
|
|
||||||
|
let started_at = Instant::now();
|
||||||
|
let mut postmaster_pid_seen = false;
|
||||||
|
loop {
|
||||||
if let Ok(Some(status)) = pg.try_wait() {
|
if let Ok(Some(status)) = pg.try_wait() {
|
||||||
// Postgres exited, that is not what we expected, bail out earlier.
|
// Postgres exited, that is not what we expected, bail out earlier.
|
||||||
let code = status.code().unwrap_or(-1);
|
let code = status.code().unwrap_or(-1);
|
||||||
bail!("Postgres exited unexpectedly with code {}", code);
|
bail!("Postgres exited unexpectedly with code {}", code);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
let res = rx.recv_timeout(Duration::from_millis(100));
|
||||||
|
log::debug!("woken up by notify: {res:?}");
|
||||||
|
// If there are multiple events in the channel already, we only need to be
|
||||||
|
// check once. Swallow the extra events before we go ahead to check the
|
||||||
|
// pid file.
|
||||||
|
while let Ok(res) = rx.try_recv() {
|
||||||
|
log::debug!("swallowing extra event: {res:?}");
|
||||||
|
}
|
||||||
|
|
||||||
// Check that we can open pid file first.
|
// Check that we can open pid file first.
|
||||||
if let Ok(file) = File::open(&pid_path) {
|
if let Ok(file) = File::open(&pid_path) {
|
||||||
|
if !postmaster_pid_seen {
|
||||||
|
log::debug!("postmaster.pid appeared");
|
||||||
|
watcher
|
||||||
|
.unwatch(pgdata)
|
||||||
|
.expect("Failed to remove pgdata dir watch");
|
||||||
|
watcher
|
||||||
|
.watch(&pid_path, RecursiveMode::NonRecursive)
|
||||||
|
.expect("Failed to add postmaster.pid file watch");
|
||||||
|
postmaster_pid_seen = true;
|
||||||
|
}
|
||||||
|
|
||||||
let file = BufReader::new(file);
|
let file = BufReader::new(file);
|
||||||
let last_line = file.lines().last();
|
let last_line = file.lines().last();
|
||||||
|
|
||||||
// Pid file could be there and we could read it, but it could be empty, for example.
|
// Pid file could be there and we could read it, but it could be empty, for example.
|
||||||
if let Some(Ok(line)) = last_line {
|
if let Some(Ok(line)) = last_line {
|
||||||
let status = line.trim();
|
let status = line.trim();
|
||||||
let can_connect = TcpStream::connect_timeout(&addr, timeout).is_ok();
|
log::debug!("last line of postmaster.pid: {status:?}");
|
||||||
|
|
||||||
// Now Postgres is ready to accept connections
|
// Now Postgres is ready to accept connections
|
||||||
if status == "ready" && can_connect {
|
if status == "ready" {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
thread::sleep(pause);
|
// Give up after POSTGRES_WAIT_TIMEOUT.
|
||||||
slept += 100;
|
let duration = started_at.elapsed();
|
||||||
|
if duration >= POSTGRES_WAIT_TIMEOUT {
|
||||||
|
bail!("timed out while waiting for Postgres to start");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
log::info!("PostgreSQL is now running, continuing to configure it");
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -19,7 +19,9 @@ thiserror = "1"
|
|||||||
nix = "0.23"
|
nix = "0.23"
|
||||||
reqwest = { version = "0.11", default-features = false, features = ["blocking", "json", "rustls-tls"] }
|
reqwest = { version = "0.11", default-features = false, features = ["blocking", "json", "rustls-tls"] }
|
||||||
|
|
||||||
pageserver = { path = "../pageserver" }
|
# Note: Do not directly depend on pageserver or safekeeper; use pageserver_api or safekeeper_api
|
||||||
safekeeper = { path = "../safekeeper" }
|
# instead, so that recompile times are better.
|
||||||
|
pageserver_api = { path = "../libs/pageserver_api" }
|
||||||
|
safekeeper_api = { path = "../libs/safekeeper_api" }
|
||||||
utils = { path = "../libs/utils" }
|
utils = { path = "../libs/utils" }
|
||||||
workspace_hack = { version = "0.1", path = "../workspace_hack" }
|
workspace_hack = { version = "0.1", path = "../workspace_hack" }
|
||||||
|
|||||||
@@ -12,12 +12,12 @@ use control_plane::local_env::{EtcdBroker, LocalEnv};
|
|||||||
use control_plane::safekeeper::SafekeeperNode;
|
use control_plane::safekeeper::SafekeeperNode;
|
||||||
use control_plane::storage::PageServerNode;
|
use control_plane::storage::PageServerNode;
|
||||||
use control_plane::{etcd, local_env};
|
use control_plane::{etcd, local_env};
|
||||||
use pageserver::config::defaults::{
|
use pageserver_api::models::TimelineInfo;
|
||||||
|
use pageserver_api::{
|
||||||
DEFAULT_HTTP_LISTEN_ADDR as DEFAULT_PAGESERVER_HTTP_ADDR,
|
DEFAULT_HTTP_LISTEN_ADDR as DEFAULT_PAGESERVER_HTTP_ADDR,
|
||||||
DEFAULT_PG_LISTEN_ADDR as DEFAULT_PAGESERVER_PG_ADDR,
|
DEFAULT_PG_LISTEN_ADDR as DEFAULT_PAGESERVER_PG_ADDR,
|
||||||
};
|
};
|
||||||
use pageserver::http::models::TimelineInfo;
|
use safekeeper_api::{
|
||||||
use safekeeper::defaults::{
|
|
||||||
DEFAULT_HTTP_LISTEN_PORT as DEFAULT_SAFEKEEPER_HTTP_PORT,
|
DEFAULT_HTTP_LISTEN_PORT as DEFAULT_SAFEKEEPER_HTTP_PORT,
|
||||||
DEFAULT_PG_LISTEN_PORT as DEFAULT_SAFEKEEPER_PG_PORT,
|
DEFAULT_PG_LISTEN_PORT as DEFAULT_SAFEKEEPER_PG_PORT,
|
||||||
};
|
};
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ use nix::unistd::Pid;
|
|||||||
use postgres::Config;
|
use postgres::Config;
|
||||||
use reqwest::blocking::{Client, RequestBuilder, Response};
|
use reqwest::blocking::{Client, RequestBuilder, Response};
|
||||||
use reqwest::{IntoUrl, Method};
|
use reqwest::{IntoUrl, Method};
|
||||||
use safekeeper::http::models::TimelineCreateRequest;
|
use safekeeper_api::models::TimelineCreateRequest;
|
||||||
use thiserror::Error;
|
use thiserror::Error;
|
||||||
use utils::{
|
use utils::{
|
||||||
connstring::connection_address,
|
connstring::connection_address,
|
||||||
|
|||||||
@@ -11,7 +11,7 @@ use anyhow::{bail, Context};
|
|||||||
use nix::errno::Errno;
|
use nix::errno::Errno;
|
||||||
use nix::sys::signal::{kill, Signal};
|
use nix::sys::signal::{kill, Signal};
|
||||||
use nix::unistd::Pid;
|
use nix::unistd::Pid;
|
||||||
use pageserver::http::models::{
|
use pageserver_api::models::{
|
||||||
TenantConfigRequest, TenantCreateRequest, TenantInfo, TimelineCreateRequest, TimelineInfo,
|
TenantConfigRequest, TenantCreateRequest, TenantInfo, TimelineCreateRequest, TimelineInfo,
|
||||||
};
|
};
|
||||||
use postgres::{Config, NoTls};
|
use postgres::{Config, NoTls};
|
||||||
|
|||||||
@@ -0,0 +1,163 @@
|
|||||||
|
# Storage messaging
|
||||||
|
|
||||||
|
Safekeepers need to communicate to each other to
|
||||||
|
* Trim WAL on safekeepers;
|
||||||
|
* Decide on which SK should push WAL to the S3;
|
||||||
|
* Decide on when to shut down SK<->pageserver connection;
|
||||||
|
* Understand state of each other to perform peer recovery;
|
||||||
|
|
||||||
|
Pageservers need to communicate to safekeepers to decide which SK should provide
|
||||||
|
WAL to the pageserver.
|
||||||
|
|
||||||
|
This is an iteration on [015-storage-messaging](https://github.com/neondatabase/neon/blob/main/docs/rfcs/015-storage-messaging.md) describing current situation,
|
||||||
|
potential performance issue and ways to address it.
|
||||||
|
|
||||||
|
## Background
|
||||||
|
|
||||||
|
What we have currently is very close to etcd variant described in
|
||||||
|
015-storage-messaging. Basically, we have single `SkTimelineInfo` message
|
||||||
|
periodically sent by all safekeepers to etcd for each timeline.
|
||||||
|
* Safekeepers subscribe to it to learn status of peers (currently they subscribe to
|
||||||
|
'everything', but they can and should fetch data only for timelines they hold).
|
||||||
|
* Pageserver subscribes to it (separate watch per timeline) to learn safekeepers
|
||||||
|
positions; based on that, it decides from which safekeepers to pull WAL.
|
||||||
|
|
||||||
|
Also, safekeepers use etcd elections API to make sure only single safekeeper
|
||||||
|
offloads WAL.
|
||||||
|
|
||||||
|
It works, and callmemaybe is gone. However, this has a performance
|
||||||
|
hazard. Currently deployed etcd can do about 6k puts per second (using its own
|
||||||
|
`benchmark` tool); on my 6 core laptop, while running on tmpfs, this gets to
|
||||||
|
35k. Making benchmark closer to our usage [etcd watch bench](https://github.com/arssher/etcd-client/blob/watch-bench/examples/watch_bench.rs),
|
||||||
|
I get ~10k received messages per second with various number of publisher-subscribers
|
||||||
|
(laptop, tmpfs). Diving this by 12 (3 sks generate msg, 1 ps + 3 sk consume them) we
|
||||||
|
get about 800 active timelines, if message is sent each second. Not extremely
|
||||||
|
low, but quite reachable.
|
||||||
|
|
||||||
|
A lot of idle watches seem to be ok though -- which is good, as pageserver
|
||||||
|
subscribes to all its timelines regardless of their activity.
|
||||||
|
|
||||||
|
Also, running etcd with fsyncs disabled is messy -- data dir must be wiped on
|
||||||
|
each restart or there is a risk of corruption errors.
|
||||||
|
|
||||||
|
The reason is etcd making much more than what we need; it is a fault tolerant
|
||||||
|
store with strong consistency, but I claim all we need here is just simplest pub
|
||||||
|
sub with best effort delivery, because
|
||||||
|
* We already have centralized source of truth for long running data, like which
|
||||||
|
tlis are on which nodes -- the console.
|
||||||
|
* Momentary data (safekeeper/pageserver progress) doesn't make sense to persist.
|
||||||
|
Instead of putting each change to broker, expecting it to reliably deliver it
|
||||||
|
is better to just have constant flow of data for active timelines: 1) they
|
||||||
|
serve as natural heartbeats -- if node can't send, we shouldn't pull WAL from
|
||||||
|
it 2) it is simpler -- no need to track delivery to/from the broker.
|
||||||
|
Moreover, latency here is important: the faster we obtain fresh data, the
|
||||||
|
faster we can switch to proper safekeeper after failure.
|
||||||
|
* As for WAL offloading leader election, it is trivial to achieve through these
|
||||||
|
heartbeats -- just take suitable node through deterministic rule (min node
|
||||||
|
id). Once network is stable, this is a converging process (well, except
|
||||||
|
complicated failure topology, but even then making it converge is not
|
||||||
|
hard). Such elections bear some risk of several offloaders running
|
||||||
|
concurrently for a short period of time, but that's harmless.
|
||||||
|
|
||||||
|
Generally, if one needs strong consistency, electing leader per se is not
|
||||||
|
enough; it must be accompanied with number (logical clock ts), checked at
|
||||||
|
every action to track causality. s3 doesn't provide CAS, so it can't
|
||||||
|
differentiate old/new leader, this must be solved differently.
|
||||||
|
|
||||||
|
We could use etcd CAS (its most powerful/useful primitive actually) to issue
|
||||||
|
these leader numbers (and e.g. prefix files in s3), but currently I don't see
|
||||||
|
need for that.
|
||||||
|
|
||||||
|
|
||||||
|
Obviously best effort pub sub is much more simpler and performant; the one proposed is
|
||||||
|
|
||||||
|
## gRPC broker
|
||||||
|
|
||||||
|
I took tonic and [prototyped](https://github.com/neondatabase/neon/blob/asher/neon-broker/broker/src/broker.rs) the replacement of functionality we currently use
|
||||||
|
with grpc streams and tokio mpsc channels. The implementation description is at the file header.
|
||||||
|
|
||||||
|
It is just 500 lines of code and core functionality is complete. 1-1 pub sub
|
||||||
|
gives about 120k received messages per second; having multiple subscribers in
|
||||||
|
different connecitons quickly scales to 1 million received messages per second.
|
||||||
|
I had concerns about many concurrent streams in singe connection, but 2^20
|
||||||
|
subscribers still work (though eat memory, with 10 publishers 20GB are consumed;
|
||||||
|
in this implementation each publisher holds full copy of all subscribers). There
|
||||||
|
is `bench.rs` nearby which I used for testing.
|
||||||
|
|
||||||
|
`SkTimelineInfo` is wired here, but another message can be added (e.g. if
|
||||||
|
pageservers want to communicate with each other) with templating.
|
||||||
|
|
||||||
|
### Fault tolerance
|
||||||
|
|
||||||
|
Since such broker is stateless, we can run it under k8s. Or add proxying to
|
||||||
|
other members, with best-effort this is simple.
|
||||||
|
|
||||||
|
### Security implications
|
||||||
|
|
||||||
|
Communication happens in a private network that is not exposed to users;
|
||||||
|
additionaly we can add auth to the broker.
|
||||||
|
|
||||||
|
## Alternative: get existing pub-sub
|
||||||
|
|
||||||
|
We could take some existing pub sub solution, e.g. RabbitMQ, Redis. But in this
|
||||||
|
case IMV simplicity of our own outweights external dependency costs (RabbitMQ is
|
||||||
|
much more complicated and needs VM; Redis Rust client maintenance is not
|
||||||
|
ideal...). Also note that projects like CockroachDB and TiDB are based on gRPC
|
||||||
|
as well.
|
||||||
|
|
||||||
|
## Alternative: direct communication
|
||||||
|
|
||||||
|
Apart from being transport, broker solves one more task: discovery, i.e. letting
|
||||||
|
safekeepers and pageservers find each other. We can let safekeepers know, for
|
||||||
|
each timeline, both other safekeepers for this timeline and pageservers serving
|
||||||
|
it. In this case direct communication is possible:
|
||||||
|
- each safekeeper pushes to each other safekeeper status of timelines residing
|
||||||
|
on both of them, letting remove WAL, decide who offloads, decide on peer
|
||||||
|
recovery;
|
||||||
|
- each safekeeper pushes to each pageserver status of timelines residing on
|
||||||
|
both of them, letting pageserver choose from which sk to pull WAL;
|
||||||
|
|
||||||
|
It was mostly described in [014-safekeeper-gossip](https://github.com/neondatabase/neon/blob/main/docs/rfcs/014-safekeepers-gossip.md), but I want to recap on that.
|
||||||
|
|
||||||
|
The main pro is less one dependency: less moving parts, easier to run Neon
|
||||||
|
locally/manually, less places to monitor. Fault tolerance for broker disappears,
|
||||||
|
no kuber or something. To me this is a big thing.
|
||||||
|
|
||||||
|
Also (though not a big thing) idle watches for inactive timelines disappear:
|
||||||
|
naturally safekeepers learn about compute connection first and start pushing
|
||||||
|
status to pageserver(s), notifying it should pull.
|
||||||
|
|
||||||
|
Importantly, I think that eventually knowing and persisting peers and
|
||||||
|
pageservers on safekeepers is inevitable:
|
||||||
|
- Knowing peer safekeepers for the timeline is required for correct
|
||||||
|
automatic membership change -- new member set must be hardened on old
|
||||||
|
majority before proceeding. It is required to get rid of sync-safekeepers
|
||||||
|
as well (peer recovery up to flush_lsn).
|
||||||
|
- Knowing pageservers where the timeline is attached is needed to
|
||||||
|
1. Understand when to shut down activity on the timeline, i.e. push data to
|
||||||
|
the broker. We can have a lot of timelines sleeping quietly which
|
||||||
|
shouldn't occupy resources.
|
||||||
|
2. Preserve WAL for these (currently we offload to s3 and take it from there,
|
||||||
|
but serving locally is better, and we get one less condition on which WAL
|
||||||
|
can be removed from s3).
|
||||||
|
|
||||||
|
I suppose this membership data should be passed to safekeepers directly from the
|
||||||
|
console because
|
||||||
|
1. Console is the original source of this data, conceptually this is the
|
||||||
|
simplest way (rather than passing it through compute or something).
|
||||||
|
2. We already have similar code for deleting timeline on safekeepers
|
||||||
|
(and attaching/detaching timeline on pageserver), this is a typical
|
||||||
|
action -- queue operation against storage node and execute it until it
|
||||||
|
completes (or timeline is dropped).
|
||||||
|
|
||||||
|
Cons of direct communication are
|
||||||
|
- It is more complicated: each safekeeper should maintain set of peers it talks
|
||||||
|
to, and set of timelines for each such peer -- they ought to be multiplexed
|
||||||
|
into single connection.
|
||||||
|
- Totally, we have O(n^2) connections instead of O(n) with broker schema
|
||||||
|
(still O(n) on each node). However, these are relatively stable, async and
|
||||||
|
thus not very expensive, I don't think this is a big problem. Up to 10k
|
||||||
|
storage nodes I doubt connection overhead would be noticeable.
|
||||||
|
|
||||||
|
I'd use gRPC for direct communication, and in this sense gRPC based broker is a
|
||||||
|
step towards it.
|
||||||
+1
-1
@@ -96,7 +96,7 @@ A single virtual environment with all dependencies is described in the single `P
|
|||||||
sudo apt install python3.9
|
sudo apt install python3.9
|
||||||
```
|
```
|
||||||
- Install `poetry`
|
- Install `poetry`
|
||||||
- Exact version of `poetry` is not important, see installation instructions available at poetry's [website](https://python-poetry.org/docs/#installation)`.
|
- Exact version of `poetry` is not important, see installation instructions available at poetry's [website](https://python-poetry.org/docs/#installation).
|
||||||
- Install dependencies via `./scripts/pysync`.
|
- Install dependencies via `./scripts/pysync`.
|
||||||
- Note that CI uses specific Python version (look for `PYTHON_VERSION` [here](https://github.com/neondatabase/docker-images/blob/main/rust/Dockerfile))
|
- Note that CI uses specific Python version (look for `PYTHON_VERSION` [here](https://github.com/neondatabase/docker-images/blob/main/rust/Dockerfile))
|
||||||
so if you have different version some linting tools can yield different result locally vs in the CI.
|
so if you have different version some linting tools can yield different result locally vs in the CI.
|
||||||
|
|||||||
+27
-2
@@ -3,7 +3,7 @@
|
|||||||
//! Otherwise, we might not see all metrics registered via
|
//! Otherwise, we might not see all metrics registered via
|
||||||
//! a default registry.
|
//! a default registry.
|
||||||
use once_cell::sync::Lazy;
|
use once_cell::sync::Lazy;
|
||||||
use prometheus::core::{AtomicU64, GenericGauge, GenericGaugeVec};
|
use prometheus::core::{AtomicU64, Collector, GenericGauge, GenericGaugeVec};
|
||||||
pub use prometheus::opts;
|
pub use prometheus::opts;
|
||||||
pub use prometheus::register;
|
pub use prometheus::register;
|
||||||
pub use prometheus::{core, default_registry, proto};
|
pub use prometheus::{core, default_registry, proto};
|
||||||
@@ -17,6 +17,7 @@ pub use prometheus::{register_int_counter_vec, IntCounterVec};
|
|||||||
pub use prometheus::{register_int_gauge, IntGauge};
|
pub use prometheus::{register_int_gauge, IntGauge};
|
||||||
pub use prometheus::{register_int_gauge_vec, IntGaugeVec};
|
pub use prometheus::{register_int_gauge_vec, IntGaugeVec};
|
||||||
pub use prometheus::{Encoder, TextEncoder};
|
pub use prometheus::{Encoder, TextEncoder};
|
||||||
|
use prometheus::{Registry, Result};
|
||||||
|
|
||||||
mod wrappers;
|
mod wrappers;
|
||||||
pub use wrappers::{CountedReader, CountedWriter};
|
pub use wrappers::{CountedReader, CountedWriter};
|
||||||
@@ -32,13 +33,27 @@ macro_rules! register_uint_gauge_vec {
|
|||||||
}};
|
}};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Special internal registry, to collect metrics independently from the default registry.
|
||||||
|
/// Was introduced to fix deadlock with lazy registration of metrics in the default registry.
|
||||||
|
static INTERNAL_REGISTRY: Lazy<Registry> = Lazy::new(Registry::new);
|
||||||
|
|
||||||
|
/// Register a collector in the internal registry. MUST be called before the first call to `gather()`.
|
||||||
|
/// Otherwise, we can have a deadlock in the `gather()` call, trying to register a new collector
|
||||||
|
/// while holding the lock.
|
||||||
|
pub fn register_internal(c: Box<dyn Collector>) -> Result<()> {
|
||||||
|
INTERNAL_REGISTRY.register(c)
|
||||||
|
}
|
||||||
|
|
||||||
/// Gathers all Prometheus metrics and records the I/O stats just before that.
|
/// Gathers all Prometheus metrics and records the I/O stats just before that.
|
||||||
///
|
///
|
||||||
/// Metrics gathering is a relatively simple and standalone operation, so
|
/// Metrics gathering is a relatively simple and standalone operation, so
|
||||||
/// it might be fine to do it this way to keep things simple.
|
/// it might be fine to do it this way to keep things simple.
|
||||||
pub fn gather() -> Vec<prometheus::proto::MetricFamily> {
|
pub fn gather() -> Vec<prometheus::proto::MetricFamily> {
|
||||||
update_rusage_metrics();
|
update_rusage_metrics();
|
||||||
prometheus::gather()
|
let mut mfs = prometheus::gather();
|
||||||
|
let mut internal_mfs = INTERNAL_REGISTRY.gather();
|
||||||
|
mfs.append(&mut internal_mfs);
|
||||||
|
mfs
|
||||||
}
|
}
|
||||||
|
|
||||||
static DISK_IO_BYTES: Lazy<IntGaugeVec> = Lazy::new(|| {
|
static DISK_IO_BYTES: Lazy<IntGaugeVec> = Lazy::new(|| {
|
||||||
@@ -62,6 +77,16 @@ pub const DISK_WRITE_SECONDS_BUCKETS: &[f64] = &[
|
|||||||
0.000_050, 0.000_100, 0.000_500, 0.001, 0.003, 0.005, 0.01, 0.05, 0.1, 0.3, 0.5,
|
0.000_050, 0.000_100, 0.000_500, 0.001, 0.003, 0.005, 0.01, 0.05, 0.1, 0.3, 0.5,
|
||||||
];
|
];
|
||||||
|
|
||||||
|
pub fn set_build_info_metric(revision: &str) {
|
||||||
|
let metric = register_int_gauge_vec!(
|
||||||
|
"libmetrics_build_info",
|
||||||
|
"Build/version information",
|
||||||
|
&["revision"]
|
||||||
|
)
|
||||||
|
.expect("Failed to register build info metric");
|
||||||
|
metric.with_label_values(&[revision]).set(1);
|
||||||
|
}
|
||||||
|
|
||||||
// Records I/O stats in a "cross-platform" way.
|
// Records I/O stats in a "cross-platform" way.
|
||||||
// Compiles both on macOS and Linux, but current macOS implementation always returns 0 as values for I/O stats.
|
// Compiles both on macOS and Linux, but current macOS implementation always returns 0 as values for I/O stats.
|
||||||
// An alternative is to read procfs (`/proc/[pid]/io`) which does not work under macOS at all, hence abandoned.
|
// An alternative is to read procfs (`/proc/[pid]/io`) which does not work under macOS at all, hence abandoned.
|
||||||
|
|||||||
@@ -0,0 +1,12 @@
|
|||||||
|
[package]
|
||||||
|
name = "pageserver_api"
|
||||||
|
version = "0.1.0"
|
||||||
|
edition = "2021"
|
||||||
|
|
||||||
|
[dependencies]
|
||||||
|
serde = { version = "1.0", features = ["derive"] }
|
||||||
|
serde_with = "1.12.0"
|
||||||
|
const_format = "0.2.21"
|
||||||
|
|
||||||
|
utils = { path = "../utils" }
|
||||||
|
workspace_hack = { version = "0.1", path = "../../workspace_hack" }
|
||||||
@@ -0,0 +1,9 @@
|
|||||||
|
use const_format::formatcp;
|
||||||
|
|
||||||
|
/// Public API types
|
||||||
|
pub mod models;
|
||||||
|
|
||||||
|
pub const DEFAULT_PG_LISTEN_PORT: u16 = 64000;
|
||||||
|
pub const DEFAULT_PG_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_PG_LISTEN_PORT}");
|
||||||
|
pub const DEFAULT_HTTP_LISTEN_PORT: u16 = 9898;
|
||||||
|
pub const DEFAULT_HTTP_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_HTTP_LISTEN_PORT}");
|
||||||
@@ -7,7 +7,17 @@ use utils::{
|
|||||||
lsn::Lsn,
|
lsn::Lsn,
|
||||||
};
|
};
|
||||||
|
|
||||||
use crate::tenant::TenantState;
|
/// A state of a tenant in pageserver's memory.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
|
||||||
|
pub enum TenantState {
|
||||||
|
/// Tenant is fully operational, its background jobs might be running or not.
|
||||||
|
Active { background_jobs_running: bool },
|
||||||
|
/// A tenant is recognized by pageserver, but not yet ready to operate:
|
||||||
|
/// e.g. not present locally and being downloaded or being read into memory from the file system.
|
||||||
|
Paused,
|
||||||
|
/// A tenant is recognized by the pageserver, but no longer used for any operations, as failed to get activated.
|
||||||
|
Broken,
|
||||||
|
}
|
||||||
|
|
||||||
#[serde_as]
|
#[serde_as]
|
||||||
#[derive(Serialize, Deserialize)]
|
#[derive(Serialize, Deserialize)]
|
||||||
@@ -3,9 +3,11 @@
|
|||||||
#![allow(non_snake_case)]
|
#![allow(non_snake_case)]
|
||||||
// bindgen creates some unsafe code with no doc comments.
|
// bindgen creates some unsafe code with no doc comments.
|
||||||
#![allow(clippy::missing_safety_doc)]
|
#![allow(clippy::missing_safety_doc)]
|
||||||
// suppress warnings on rust 1.53 due to bindgen unit tests.
|
// noted at 1.63 that in many cases there's a u32 -> u32 transmutes in bindgen code.
|
||||||
// https://github.com/rust-lang/rust-bindgen/issues/1651
|
#![allow(clippy::useless_transmute)]
|
||||||
#![allow(deref_nullptr)]
|
// modules included with the postgres_ffi macro depend on the types of the specific version's
|
||||||
|
// types, and trigger a too eager lint.
|
||||||
|
#![allow(clippy::duplicate_mod)]
|
||||||
|
|
||||||
use bytes::Bytes;
|
use bytes::Bytes;
|
||||||
use utils::bin_ser::SerializeError;
|
use utils::bin_ser::SerializeError;
|
||||||
|
|||||||
@@ -57,12 +57,10 @@ pub const SIZE_OF_XLOG_RECORD_DATA_HEADER_SHORT: usize = 1 * 2;
|
|||||||
/// in order to let CLOG_TRUNCATE mechanism correctly extend CLOG.
|
/// in order to let CLOG_TRUNCATE mechanism correctly extend CLOG.
|
||||||
const XID_CHECKPOINT_INTERVAL: u32 = 1024;
|
const XID_CHECKPOINT_INTERVAL: u32 = 1024;
|
||||||
|
|
||||||
#[allow(non_snake_case)]
|
|
||||||
pub fn XLogSegmentsPerXLogId(wal_segsz_bytes: usize) -> XLogSegNo {
|
pub fn XLogSegmentsPerXLogId(wal_segsz_bytes: usize) -> XLogSegNo {
|
||||||
(0x100000000u64 / wal_segsz_bytes as u64) as XLogSegNo
|
(0x100000000u64 / wal_segsz_bytes as u64) as XLogSegNo
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(non_snake_case)]
|
|
||||||
pub fn XLogSegNoOffsetToRecPtr(
|
pub fn XLogSegNoOffsetToRecPtr(
|
||||||
segno: XLogSegNo,
|
segno: XLogSegNo,
|
||||||
offset: u32,
|
offset: u32,
|
||||||
@@ -71,7 +69,6 @@ pub fn XLogSegNoOffsetToRecPtr(
|
|||||||
segno * (wal_segsz_bytes as u64) + (offset as u64)
|
segno * (wal_segsz_bytes as u64) + (offset as u64)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(non_snake_case)]
|
|
||||||
pub fn XLogFileName(tli: TimeLineID, logSegNo: XLogSegNo, wal_segsz_bytes: usize) -> String {
|
pub fn XLogFileName(tli: TimeLineID, logSegNo: XLogSegNo, wal_segsz_bytes: usize) -> String {
|
||||||
format!(
|
format!(
|
||||||
"{:>08X}{:>08X}{:>08X}",
|
"{:>08X}{:>08X}{:>08X}",
|
||||||
@@ -81,7 +78,6 @@ pub fn XLogFileName(tli: TimeLineID, logSegNo: XLogSegNo, wal_segsz_bytes: usize
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(non_snake_case)]
|
|
||||||
pub fn XLogFromFileName(fname: &str, wal_seg_size: usize) -> (XLogSegNo, TimeLineID) {
|
pub fn XLogFromFileName(fname: &str, wal_seg_size: usize) -> (XLogSegNo, TimeLineID) {
|
||||||
let tli = u32::from_str_radix(&fname[0..8], 16).unwrap();
|
let tli = u32::from_str_radix(&fname[0..8], 16).unwrap();
|
||||||
let log = u32::from_str_radix(&fname[8..16], 16).unwrap() as XLogSegNo;
|
let log = u32::from_str_radix(&fname[8..16], 16).unwrap() as XLogSegNo;
|
||||||
@@ -89,12 +85,10 @@ pub fn XLogFromFileName(fname: &str, wal_seg_size: usize) -> (XLogSegNo, TimeLin
|
|||||||
(log * XLogSegmentsPerXLogId(wal_seg_size) + seg, tli)
|
(log * XLogSegmentsPerXLogId(wal_seg_size) + seg, tli)
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(non_snake_case)]
|
|
||||||
pub fn IsXLogFileName(fname: &str) -> bool {
|
pub fn IsXLogFileName(fname: &str) -> bool {
|
||||||
return fname.len() == XLOG_FNAME_LEN && fname.chars().all(|c| c.is_ascii_hexdigit());
|
return fname.len() == XLOG_FNAME_LEN && fname.chars().all(|c| c.is_ascii_hexdigit());
|
||||||
}
|
}
|
||||||
|
|
||||||
#[allow(non_snake_case)]
|
|
||||||
pub fn IsPartialXLogFileName(fname: &str) -> bool {
|
pub fn IsPartialXLogFileName(fname: &str) -> bool {
|
||||||
fname.ends_with(".partial") && IsXLogFileName(&fname[0..fname.len() - 8])
|
fname.ends_with(".partial") && IsXLogFileName(&fname[0..fname.len() - 8])
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,12 @@
|
|||||||
|
[package]
|
||||||
|
name = "safekeeper_api"
|
||||||
|
version = "0.1.0"
|
||||||
|
edition = "2021"
|
||||||
|
|
||||||
|
[dependencies]
|
||||||
|
serde = { version = "1.0", features = ["derive"] }
|
||||||
|
serde_with = "1.12.0"
|
||||||
|
const_format = "0.2.21"
|
||||||
|
|
||||||
|
utils = { path = "../utils" }
|
||||||
|
workspace_hack = { version = "0.1", path = "../../workspace_hack" }
|
||||||
@@ -0,0 +1,10 @@
|
|||||||
|
use const_format::formatcp;
|
||||||
|
|
||||||
|
/// Public API types
|
||||||
|
pub mod models;
|
||||||
|
|
||||||
|
pub const DEFAULT_PG_LISTEN_PORT: u16 = 5454;
|
||||||
|
pub const DEFAULT_PG_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_PG_LISTEN_PORT}");
|
||||||
|
|
||||||
|
pub const DEFAULT_HTTP_LISTEN_PORT: u16 = 7676;
|
||||||
|
pub const DEFAULT_HTTP_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_HTTP_LISTEN_PORT}");
|
||||||
@@ -9,6 +9,7 @@ use once_cell::sync::Lazy;
|
|||||||
use routerify::ext::RequestExt;
|
use routerify::ext::RequestExt;
|
||||||
use routerify::RequestInfo;
|
use routerify::RequestInfo;
|
||||||
use routerify::{Middleware, Router, RouterBuilder, RouterService};
|
use routerify::{Middleware, Router, RouterBuilder, RouterService};
|
||||||
|
use tokio::task::JoinError;
|
||||||
use tracing::info;
|
use tracing::info;
|
||||||
|
|
||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
@@ -35,7 +36,13 @@ async fn prometheus_metrics_handler(_req: Request<Body>) -> Result<Response<Body
|
|||||||
let mut buffer = vec![];
|
let mut buffer = vec![];
|
||||||
let encoder = TextEncoder::new();
|
let encoder = TextEncoder::new();
|
||||||
|
|
||||||
let metrics = metrics::gather();
|
let metrics = tokio::task::spawn_blocking(move || {
|
||||||
|
// Currently we take a lot of mutexes while collecting metrics, so it's
|
||||||
|
// better to spawn a blocking task to avoid blocking the event loop.
|
||||||
|
metrics::gather()
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.map_err(|e: JoinError| ApiError::InternalServerError(e.into()))?;
|
||||||
encoder.encode(&metrics, &mut buffer).unwrap();
|
encoder.encode(&metrics, &mut buffer).unwrap();
|
||||||
|
|
||||||
let response = Response::builder()
|
let response = Response::builder()
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ use serde::{Deserialize, Serialize};
|
|||||||
use std::{
|
use std::{
|
||||||
borrow::Cow,
|
borrow::Cow,
|
||||||
collections::HashMap,
|
collections::HashMap,
|
||||||
|
fmt,
|
||||||
future::Future,
|
future::Future,
|
||||||
io::{self, Cursor},
|
io::{self, Cursor},
|
||||||
str,
|
str,
|
||||||
@@ -124,6 +125,19 @@ pub struct CancelKeyData {
|
|||||||
pub cancel_key: i32,
|
pub cancel_key: i32,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl fmt::Display for CancelKeyData {
|
||||||
|
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||||
|
let hi = (self.backend_pid as u64) << 32;
|
||||||
|
let lo = self.cancel_key as u64;
|
||||||
|
let id = hi | lo;
|
||||||
|
|
||||||
|
// This format is more compact and might work better for logs.
|
||||||
|
f.debug_tuple("CancelKeyData")
|
||||||
|
.field(&format_args!("{:x}", id))
|
||||||
|
.finish()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
use rand::distributions::{Distribution, Standard};
|
use rand::distributions::{Distribution, Standard};
|
||||||
impl Distribution<CancelKeyData> for Standard {
|
impl Distribution<CancelKeyData> for Standard {
|
||||||
fn sample<R: rand::Rng + ?Sized>(&self, rng: &mut R) -> CancelKeyData {
|
fn sample<R: rand::Rng + ?Sized>(&self, rng: &mut R) -> CancelKeyData {
|
||||||
|
|||||||
@@ -58,6 +58,7 @@ rstar = "0.9.3"
|
|||||||
num-traits = "0.2.15"
|
num-traits = "0.2.15"
|
||||||
amplify_num = "0.4.1"
|
amplify_num = "0.4.1"
|
||||||
|
|
||||||
|
pageserver_api = { path = "../libs/pageserver_api" }
|
||||||
postgres_ffi = { path = "../libs/postgres_ffi" }
|
postgres_ffi = { path = "../libs/postgres_ffi" }
|
||||||
etcd_broker = { path = "../libs/etcd_broker" }
|
etcd_broker = { path = "../libs/etcd_broker" }
|
||||||
metrics = { path = "../libs/metrics" }
|
metrics = { path = "../libs/metrics" }
|
||||||
|
|||||||
@@ -1,35 +0,0 @@
|
|||||||
//! Main entry point for the dump_layerfile executable
|
|
||||||
//!
|
|
||||||
//! A handy tool for debugging, that's all.
|
|
||||||
use anyhow::Result;
|
|
||||||
use clap::{App, Arg};
|
|
||||||
use pageserver::page_cache;
|
|
||||||
use pageserver::tenant::dump_layerfile_from_path;
|
|
||||||
use pageserver::virtual_file;
|
|
||||||
use std::path::PathBuf;
|
|
||||||
use utils::project_git_version;
|
|
||||||
|
|
||||||
project_git_version!(GIT_VERSION);
|
|
||||||
|
|
||||||
fn main() -> Result<()> {
|
|
||||||
let arg_matches = App::new("Neon dump_layerfile utility")
|
|
||||||
.about("Dump contents of one layer file, for debugging")
|
|
||||||
.version(GIT_VERSION)
|
|
||||||
.arg(
|
|
||||||
Arg::new("path")
|
|
||||||
.help("Path to file to dump")
|
|
||||||
.required(true)
|
|
||||||
.index(1),
|
|
||||||
)
|
|
||||||
.get_matches();
|
|
||||||
|
|
||||||
let path = PathBuf::from(arg_matches.value_of("path").unwrap());
|
|
||||||
|
|
||||||
// Basic initialization of things that don't change after startup
|
|
||||||
virtual_file::init(10);
|
|
||||||
page_cache::init(100);
|
|
||||||
|
|
||||||
dump_layerfile_from_path(&path, true)?;
|
|
||||||
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
@@ -10,6 +10,8 @@ use clap::{App, Arg};
|
|||||||
use daemonize::Daemonize;
|
use daemonize::Daemonize;
|
||||||
|
|
||||||
use fail::FailScenario;
|
use fail::FailScenario;
|
||||||
|
use metrics::set_build_info_metric;
|
||||||
|
|
||||||
use pageserver::{
|
use pageserver::{
|
||||||
config::{defaults::*, PageServerConf},
|
config::{defaults::*, PageServerConf},
|
||||||
http, page_cache, page_service, profiling, task_mgr,
|
http, page_cache, page_service, profiling, task_mgr,
|
||||||
@@ -31,11 +33,20 @@ use utils::{
|
|||||||
|
|
||||||
project_git_version!(GIT_VERSION);
|
project_git_version!(GIT_VERSION);
|
||||||
|
|
||||||
|
const FEATURES: &[&str] = &[
|
||||||
|
#[cfg(feature = "testing")]
|
||||||
|
"testing",
|
||||||
|
#[cfg(feature = "fail/failpoints")]
|
||||||
|
"fail/failpoints",
|
||||||
|
#[cfg(feature = "profiling")]
|
||||||
|
"profiling",
|
||||||
|
];
|
||||||
|
|
||||||
fn version() -> String {
|
fn version() -> String {
|
||||||
format!(
|
format!(
|
||||||
"{GIT_VERSION} profiling:{} failpoints:{}",
|
"{GIT_VERSION} failpoints: {}, features: {:?}",
|
||||||
cfg!(feature = "profiling"),
|
fail::has_failpoints(),
|
||||||
fail::has_failpoints()
|
FEATURES,
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -86,13 +97,7 @@ fn main() -> anyhow::Result<()> {
|
|||||||
.get_matches();
|
.get_matches();
|
||||||
|
|
||||||
if arg_matches.is_present("enabled-features") {
|
if arg_matches.is_present("enabled-features") {
|
||||||
let features: &[&str] = &[
|
println!("{{\"features\": {FEATURES:?} }}");
|
||||||
#[cfg(feature = "testing")]
|
|
||||||
"testing",
|
|
||||||
#[cfg(feature = "profiling")]
|
|
||||||
"profiling",
|
|
||||||
];
|
|
||||||
println!("{{\"features\": {features:?} }}");
|
|
||||||
return Ok(());
|
return Ok(());
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -356,6 +361,8 @@ fn start_pageserver(conf: &'static PageServerConf, daemonize: bool) -> Result<()
|
|||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
|
||||||
|
set_build_info_metric(GIT_VERSION);
|
||||||
|
|
||||||
// All started up! Now just sit and wait for shutdown signal.
|
// All started up! Now just sit and wait for shutdown signal.
|
||||||
signals.handle(|signal| match signal {
|
signals.handle(|signal| match signal {
|
||||||
Signal::Quit => {
|
Signal::Quit => {
|
||||||
|
|||||||
@@ -0,0 +1,144 @@
|
|||||||
|
//! A helper tool to manage pageserver binary files.
|
||||||
|
//! Accepts a file as an argument, attempts to parse it with all ways possible
|
||||||
|
//! and prints its interpreted context.
|
||||||
|
//!
|
||||||
|
//! Separate, `metadata` subcommand allows to print and update pageserver's metadata file.
|
||||||
|
use std::{
|
||||||
|
path::{Path, PathBuf},
|
||||||
|
str::FromStr,
|
||||||
|
};
|
||||||
|
|
||||||
|
use anyhow::Context;
|
||||||
|
use clap::{App, Arg};
|
||||||
|
|
||||||
|
use pageserver::{
|
||||||
|
page_cache,
|
||||||
|
tenant::{dump_layerfile_from_path, metadata::TimelineMetadata},
|
||||||
|
virtual_file,
|
||||||
|
};
|
||||||
|
use postgres_ffi::ControlFileData;
|
||||||
|
use utils::{lsn::Lsn, project_git_version};
|
||||||
|
|
||||||
|
project_git_version!(GIT_VERSION);
|
||||||
|
|
||||||
|
const METADATA_SUBCOMMAND: &str = "metadata";
|
||||||
|
|
||||||
|
fn main() -> anyhow::Result<()> {
|
||||||
|
let arg_matches = App::new("Neon Pageserver binutils")
|
||||||
|
.about("Reads pageserver (and related) binary files management utility")
|
||||||
|
.version(GIT_VERSION)
|
||||||
|
.arg(Arg::new("path").help("Input file path").required(false))
|
||||||
|
.subcommand(
|
||||||
|
App::new(METADATA_SUBCOMMAND)
|
||||||
|
.about("Read and update pageserver metadata file")
|
||||||
|
.arg(
|
||||||
|
Arg::new("metadata_path")
|
||||||
|
.help("Input metadata file path")
|
||||||
|
.required(false),
|
||||||
|
)
|
||||||
|
.arg(
|
||||||
|
Arg::new("disk_consistent_lsn")
|
||||||
|
.long("disk_consistent_lsn")
|
||||||
|
.takes_value(true)
|
||||||
|
.help("Replace disk consistent Lsn"),
|
||||||
|
)
|
||||||
|
.arg(
|
||||||
|
Arg::new("prev_record_lsn")
|
||||||
|
.long("prev_record_lsn")
|
||||||
|
.takes_value(true)
|
||||||
|
.help("Replace previous record Lsn"),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
.get_matches();
|
||||||
|
|
||||||
|
match arg_matches.subcommand() {
|
||||||
|
Some((subcommand_name, subcommand_matches)) => {
|
||||||
|
let path = PathBuf::from(
|
||||||
|
subcommand_matches
|
||||||
|
.value_of("metadata_path")
|
||||||
|
.context("'metadata_path' argument is missing")?,
|
||||||
|
);
|
||||||
|
anyhow::ensure!(
|
||||||
|
subcommand_name == METADATA_SUBCOMMAND,
|
||||||
|
"Unknown subcommand {subcommand_name}"
|
||||||
|
);
|
||||||
|
handle_metadata(&path, subcommand_matches)?;
|
||||||
|
}
|
||||||
|
None => {
|
||||||
|
let path = PathBuf::from(
|
||||||
|
arg_matches
|
||||||
|
.value_of("path")
|
||||||
|
.context("'path' argument is missing")?,
|
||||||
|
);
|
||||||
|
println!(
|
||||||
|
"No subcommand specified, attempting to guess the format for file {}",
|
||||||
|
path.display()
|
||||||
|
);
|
||||||
|
if let Err(e) = read_pg_control_file(&path) {
|
||||||
|
println!(
|
||||||
|
"Failed to read input file as a pg control one: {e:#}\n\
|
||||||
|
Attempting to read it as layer file"
|
||||||
|
);
|
||||||
|
print_layerfile(&path)?;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
};
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn read_pg_control_file(control_file_path: &Path) -> anyhow::Result<()> {
|
||||||
|
let control_file = ControlFileData::decode(&std::fs::read(&control_file_path)?)?;
|
||||||
|
println!("{control_file:?}");
|
||||||
|
let control_file_initdb = Lsn(control_file.checkPoint);
|
||||||
|
println!(
|
||||||
|
"pg_initdb_lsn: {}, aligned: {}",
|
||||||
|
control_file_initdb,
|
||||||
|
control_file_initdb.align()
|
||||||
|
);
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn print_layerfile(path: &Path) -> anyhow::Result<()> {
|
||||||
|
// Basic initialization of things that don't change after startup
|
||||||
|
virtual_file::init(10);
|
||||||
|
page_cache::init(100);
|
||||||
|
dump_layerfile_from_path(path, true)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn handle_metadata(path: &Path, arg_matches: &clap::ArgMatches) -> Result<(), anyhow::Error> {
|
||||||
|
let metadata_bytes = std::fs::read(&path)?;
|
||||||
|
let mut meta = TimelineMetadata::from_bytes(&metadata_bytes)?;
|
||||||
|
println!("Current metadata:\n{meta:?}");
|
||||||
|
let mut update_meta = false;
|
||||||
|
if let Some(disk_consistent_lsn) = arg_matches.value_of("disk_consistent_lsn") {
|
||||||
|
meta = TimelineMetadata::new(
|
||||||
|
Lsn::from_str(disk_consistent_lsn)?,
|
||||||
|
meta.prev_record_lsn(),
|
||||||
|
meta.ancestor_timeline(),
|
||||||
|
meta.ancestor_lsn(),
|
||||||
|
meta.latest_gc_cutoff_lsn(),
|
||||||
|
meta.initdb_lsn(),
|
||||||
|
meta.pg_version(),
|
||||||
|
);
|
||||||
|
update_meta = true;
|
||||||
|
}
|
||||||
|
if let Some(prev_record_lsn) = arg_matches.value_of("prev_record_lsn") {
|
||||||
|
meta = TimelineMetadata::new(
|
||||||
|
meta.disk_consistent_lsn(),
|
||||||
|
Some(Lsn::from_str(prev_record_lsn)?),
|
||||||
|
meta.ancestor_timeline(),
|
||||||
|
meta.ancestor_lsn(),
|
||||||
|
meta.latest_gc_cutoff_lsn(),
|
||||||
|
meta.initdb_lsn(),
|
||||||
|
meta.pg_version(),
|
||||||
|
);
|
||||||
|
update_meta = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
if update_meta {
|
||||||
|
let metadata_bytes = meta.to_bytes()?;
|
||||||
|
std::fs::write(&path, &metadata_bytes)?;
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
@@ -1,75 +0,0 @@
|
|||||||
//! Main entry point for the edit_metadata executable
|
|
||||||
//!
|
|
||||||
//! A handy tool for debugging, that's all.
|
|
||||||
use anyhow::Result;
|
|
||||||
use clap::{App, Arg};
|
|
||||||
use pageserver::tenant::metadata::TimelineMetadata;
|
|
||||||
use std::path::PathBuf;
|
|
||||||
use std::str::FromStr;
|
|
||||||
use utils::{lsn::Lsn, project_git_version};
|
|
||||||
|
|
||||||
project_git_version!(GIT_VERSION);
|
|
||||||
|
|
||||||
fn main() -> Result<()> {
|
|
||||||
let arg_matches = App::new("Neon update metadata utility")
|
|
||||||
.about("Dump or update metadata file")
|
|
||||||
.version(GIT_VERSION)
|
|
||||||
.arg(
|
|
||||||
Arg::new("path")
|
|
||||||
.help("Path to metadata file")
|
|
||||||
.required(true),
|
|
||||||
)
|
|
||||||
.arg(
|
|
||||||
Arg::new("disk_lsn")
|
|
||||||
.short('d')
|
|
||||||
.long("disk_lsn")
|
|
||||||
.takes_value(true)
|
|
||||||
.help("Replace disk constistent lsn"),
|
|
||||||
)
|
|
||||||
.arg(
|
|
||||||
Arg::new("prev_lsn")
|
|
||||||
.short('p')
|
|
||||||
.long("prev_lsn")
|
|
||||||
.takes_value(true)
|
|
||||||
.help("Previous record LSN"),
|
|
||||||
)
|
|
||||||
.get_matches();
|
|
||||||
|
|
||||||
let path = PathBuf::from(arg_matches.value_of("path").unwrap());
|
|
||||||
let metadata_bytes = std::fs::read(&path)?;
|
|
||||||
let mut meta = TimelineMetadata::from_bytes(&metadata_bytes)?;
|
|
||||||
println!("Current metadata:\n{:?}", &meta);
|
|
||||||
|
|
||||||
let mut update_meta = false;
|
|
||||||
|
|
||||||
if let Some(disk_lsn) = arg_matches.value_of("disk_lsn") {
|
|
||||||
meta = TimelineMetadata::new(
|
|
||||||
Lsn::from_str(disk_lsn)?,
|
|
||||||
meta.prev_record_lsn(),
|
|
||||||
meta.ancestor_timeline(),
|
|
||||||
meta.ancestor_lsn(),
|
|
||||||
meta.latest_gc_cutoff_lsn(),
|
|
||||||
meta.initdb_lsn(),
|
|
||||||
meta.pg_version(),
|
|
||||||
);
|
|
||||||
update_meta = true;
|
|
||||||
}
|
|
||||||
|
|
||||||
if let Some(prev_lsn) = arg_matches.value_of("prev_lsn") {
|
|
||||||
meta = TimelineMetadata::new(
|
|
||||||
meta.disk_consistent_lsn(),
|
|
||||||
Some(Lsn::from_str(prev_lsn)?),
|
|
||||||
meta.ancestor_timeline(),
|
|
||||||
meta.ancestor_lsn(),
|
|
||||||
meta.latest_gc_cutoff_lsn(),
|
|
||||||
meta.initdb_lsn(),
|
|
||||||
meta.pg_version(),
|
|
||||||
);
|
|
||||||
update_meta = true;
|
|
||||||
}
|
|
||||||
if update_meta {
|
|
||||||
let metadata_bytes = meta.to_bytes()?;
|
|
||||||
std::fs::write(&path, &metadata_bytes)?;
|
|
||||||
}
|
|
||||||
Ok(())
|
|
||||||
}
|
|
||||||
@@ -30,10 +30,10 @@ pub mod defaults {
|
|||||||
use crate::tenant_config::defaults::*;
|
use crate::tenant_config::defaults::*;
|
||||||
use const_format::formatcp;
|
use const_format::formatcp;
|
||||||
|
|
||||||
pub const DEFAULT_PG_LISTEN_PORT: u16 = 64000;
|
pub use pageserver_api::{
|
||||||
pub const DEFAULT_PG_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_PG_LISTEN_PORT}");
|
DEFAULT_HTTP_LISTEN_ADDR, DEFAULT_HTTP_LISTEN_PORT, DEFAULT_PG_LISTEN_ADDR,
|
||||||
pub const DEFAULT_HTTP_LISTEN_PORT: u16 = 9898;
|
DEFAULT_PG_LISTEN_PORT,
|
||||||
pub const DEFAULT_HTTP_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_HTTP_LISTEN_PORT}");
|
};
|
||||||
|
|
||||||
pub const DEFAULT_WAIT_LSN_TIMEOUT: &str = "60 s";
|
pub const DEFAULT_WAIT_LSN_TIMEOUT: &str = "60 s";
|
||||||
pub const DEFAULT_WAL_REDO_TIMEOUT: &str = "60 s";
|
pub const DEFAULT_WAL_REDO_TIMEOUT: &str = "60 s";
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
pub mod models;
|
|
||||||
pub mod routes;
|
pub mod routes;
|
||||||
pub use routes::make_router;
|
pub use routes::make_router;
|
||||||
|
|
||||||
|
pub use pageserver_api::models;
|
||||||
|
|||||||
@@ -207,6 +207,62 @@ paths:
|
|||||||
schema:
|
schema:
|
||||||
$ref: "#/components/schemas/Error"
|
$ref: "#/components/schemas/Error"
|
||||||
|
|
||||||
|
|
||||||
|
/v1/tenant/{tenant_id}/timeline/{timeline_id}/get_lsn_by_timestamp:
|
||||||
|
parameters:
|
||||||
|
- name: tenant_id
|
||||||
|
in: path
|
||||||
|
required: true
|
||||||
|
schema:
|
||||||
|
type: string
|
||||||
|
format: hex
|
||||||
|
- name: timeline_id
|
||||||
|
in: path
|
||||||
|
required: true
|
||||||
|
schema:
|
||||||
|
type: string
|
||||||
|
format: hex
|
||||||
|
get:
|
||||||
|
description: Get LSN by a timestamp
|
||||||
|
parameters:
|
||||||
|
- name: timestamp
|
||||||
|
in: query
|
||||||
|
required: true
|
||||||
|
schema:
|
||||||
|
type: string
|
||||||
|
format: date-time
|
||||||
|
description: A timestamp to get the LSN
|
||||||
|
responses:
|
||||||
|
"200":
|
||||||
|
description: OK
|
||||||
|
content:
|
||||||
|
application/json:
|
||||||
|
schema:
|
||||||
|
type: string
|
||||||
|
"400":
|
||||||
|
description: Error when no tenant id found in path, no timeline id or invalid timestamp
|
||||||
|
content:
|
||||||
|
application/json:
|
||||||
|
schema:
|
||||||
|
$ref: "#/components/schemas/Error"
|
||||||
|
"401":
|
||||||
|
description: Unauthorized Error
|
||||||
|
content:
|
||||||
|
application/json:
|
||||||
|
schema:
|
||||||
|
$ref: "#/components/schemas/UnauthorizedError"
|
||||||
|
"403":
|
||||||
|
description: Forbidden Error
|
||||||
|
content:
|
||||||
|
application/json:
|
||||||
|
schema:
|
||||||
|
$ref: "#/components/schemas/ForbiddenError"
|
||||||
|
"500":
|
||||||
|
description: Generic operation error
|
||||||
|
content:
|
||||||
|
application/json:
|
||||||
|
schema:
|
||||||
|
$ref: "#/components/schemas/Error"
|
||||||
/v1/tenant/{tenant_id}/attach:
|
/v1/tenant/{tenant_id}/attach:
|
||||||
parameters:
|
parameters:
|
||||||
- name: tenant_id
|
- name: tenant_id
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ use super::models::{
|
|||||||
StatusResponse, TenantConfigRequest, TenantCreateRequest, TenantCreateResponse, TenantInfo,
|
StatusResponse, TenantConfigRequest, TenantCreateRequest, TenantCreateResponse, TenantInfo,
|
||||||
TimelineCreateRequest,
|
TimelineCreateRequest,
|
||||||
};
|
};
|
||||||
|
use crate::pgdatadir_mapping::LsnForTimestamp;
|
||||||
use crate::storage_sync;
|
use crate::storage_sync;
|
||||||
use crate::storage_sync::index::{RemoteIndex, RemoteTimeline};
|
use crate::storage_sync::index::{RemoteIndex, RemoteTimeline};
|
||||||
use crate::tenant::{TenantState, Timeline};
|
use crate::tenant::{TenantState, Timeline};
|
||||||
@@ -265,6 +266,23 @@ fn query_param_present(request: &Request<Body>, param: &str) -> bool {
|
|||||||
.unwrap_or(false)
|
.unwrap_or(false)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn get_query_param(request: &Request<Body>, param_name: &str) -> Result<String, ApiError> {
|
||||||
|
request.uri().query().map_or(
|
||||||
|
Err(ApiError::BadRequest(anyhow!("empty query in request"))),
|
||||||
|
|v| {
|
||||||
|
url::form_urlencoded::parse(v.as_bytes())
|
||||||
|
.into_owned()
|
||||||
|
.find(|(k, _)| k == param_name)
|
||||||
|
.map_or(
|
||||||
|
Err(ApiError::BadRequest(anyhow!(
|
||||||
|
"no {param_name} specified in query parameters"
|
||||||
|
))),
|
||||||
|
|(_, v)| Ok(v),
|
||||||
|
)
|
||||||
|
},
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
async fn timeline_detail_handler(request: Request<Body>) -> Result<Response<Body>, ApiError> {
|
async fn timeline_detail_handler(request: Request<Body>) -> Result<Response<Body>, ApiError> {
|
||||||
let tenant_id: TenantId = parse_request_param(&request, "tenant_id")?;
|
let tenant_id: TenantId = parse_request_param(&request, "tenant_id")?;
|
||||||
let timeline_id: TimelineId = parse_request_param(&request, "timeline_id")?;
|
let timeline_id: TimelineId = parse_request_param(&request, "timeline_id")?;
|
||||||
@@ -329,6 +347,33 @@ async fn timeline_detail_handler(request: Request<Body>) -> Result<Response<Body
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async fn get_lsn_by_timestamp_handler(request: Request<Body>) -> Result<Response<Body>, ApiError> {
|
||||||
|
let tenant_id: TenantId = parse_request_param(&request, "tenant_id")?;
|
||||||
|
check_permission(&request, Some(tenant_id))?;
|
||||||
|
|
||||||
|
let timeline_id: TimelineId = parse_request_param(&request, "timeline_id")?;
|
||||||
|
let timestamp_raw = get_query_param(&request, "timestamp")?;
|
||||||
|
let timestamp = humantime::parse_rfc3339(timestamp_raw.as_str())
|
||||||
|
.with_context(|| format!("Invalid time: {:?}", timestamp_raw))
|
||||||
|
.map_err(ApiError::BadRequest)?;
|
||||||
|
let timestamp_pg = postgres_ffi::to_pg_timestamp(timestamp);
|
||||||
|
|
||||||
|
let timeline = tenant_mgr::get_tenant(tenant_id, true)
|
||||||
|
.and_then(|tenant| tenant.get_timeline(timeline_id))
|
||||||
|
.with_context(|| format!("No timeline {timeline_id} in repository for tenant {tenant_id}"))
|
||||||
|
.map_err(ApiError::NotFound)?;
|
||||||
|
let result = match timeline
|
||||||
|
.find_lsn_for_timestamp(timestamp_pg)
|
||||||
|
.map_err(ApiError::InternalServerError)?
|
||||||
|
{
|
||||||
|
LsnForTimestamp::Present(lsn) => format!("{}", lsn),
|
||||||
|
LsnForTimestamp::Future(_lsn) => "future".into(),
|
||||||
|
LsnForTimestamp::Past(_lsn) => "past".into(),
|
||||||
|
LsnForTimestamp::NoData(_lsn) => "nodata".into(),
|
||||||
|
};
|
||||||
|
json_response(StatusCode::OK, result)
|
||||||
|
}
|
||||||
|
|
||||||
// TODO makes sense to provide tenant config right away the same way as it handled in tenant_create
|
// TODO makes sense to provide tenant config right away the same way as it handled in tenant_create
|
||||||
async fn tenant_attach_handler(request: Request<Body>) -> Result<Response<Body>, ApiError> {
|
async fn tenant_attach_handler(request: Request<Body>) -> Result<Response<Body>, ApiError> {
|
||||||
let tenant_id: TenantId = parse_request_param(&request, "tenant_id")?;
|
let tenant_id: TenantId = parse_request_param(&request, "tenant_id")?;
|
||||||
@@ -337,9 +382,16 @@ async fn tenant_attach_handler(request: Request<Body>) -> Result<Response<Body>,
|
|||||||
info!("Handling tenant attach {tenant_id}");
|
info!("Handling tenant attach {tenant_id}");
|
||||||
|
|
||||||
tokio::task::spawn_blocking(move || match tenant_mgr::get_tenant(tenant_id, false) {
|
tokio::task::spawn_blocking(move || match tenant_mgr::get_tenant(tenant_id, false) {
|
||||||
Ok(_) => Err(ApiError::Conflict(
|
Ok(tenant) => {
|
||||||
"Tenant is already present locally".to_owned(),
|
if tenant.list_timelines().is_empty() {
|
||||||
)),
|
info!("Attaching to tenant {tenant_id} with zero timelines");
|
||||||
|
Ok(())
|
||||||
|
} else {
|
||||||
|
Err(ApiError::Conflict(
|
||||||
|
"Tenant is already present locally".to_owned(),
|
||||||
|
))
|
||||||
|
}
|
||||||
|
}
|
||||||
Err(_) => Ok(()),
|
Err(_) => Ok(()),
|
||||||
})
|
})
|
||||||
.await
|
.await
|
||||||
@@ -732,7 +784,7 @@ async fn tenant_config_handler(mut request: Request<Body>) -> Result<Response<Bo
|
|||||||
json_response(StatusCode::OK, ())
|
json_response(StatusCode::OK, ())
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(any(feature = "testing", feature = "failpoints"))]
|
#[cfg(feature = "testing")]
|
||||||
async fn failpoints_handler(mut request: Request<Body>) -> Result<Response<Body>, ApiError> {
|
async fn failpoints_handler(mut request: Request<Body>) -> Result<Response<Body>, ApiError> {
|
||||||
if !fail::has_failpoints() {
|
if !fail::has_failpoints() {
|
||||||
return Err(ApiError::BadRequest(anyhow!(
|
return Err(ApiError::BadRequest(anyhow!(
|
||||||
@@ -901,6 +953,10 @@ pub fn make_router(
|
|||||||
"/v1/tenant/:tenant_id/timeline/:timeline_id",
|
"/v1/tenant/:tenant_id/timeline/:timeline_id",
|
||||||
timeline_detail_handler,
|
timeline_detail_handler,
|
||||||
)
|
)
|
||||||
|
.get(
|
||||||
|
"/v1/tenant/:tenant_id/timeline/:timeline_id/get_lsn_by_timestamp",
|
||||||
|
get_lsn_by_timestamp_handler,
|
||||||
|
)
|
||||||
.put(
|
.put(
|
||||||
"/v1/tenant/:tenant_id/timeline/:timeline_id/do_gc",
|
"/v1/tenant/:tenant_id/timeline/:timeline_id/do_gc",
|
||||||
testing_api!("run timeline GC", timeline_gc_handler),
|
testing_api!("run timeline GC", timeline_gc_handler),
|
||||||
|
|||||||
@@ -12,7 +12,6 @@
|
|||||||
use anyhow::{bail, ensure, Context, Result};
|
use anyhow::{bail, ensure, Context, Result};
|
||||||
use bytes::{Buf, BufMut, Bytes, BytesMut};
|
use bytes::{Buf, BufMut, Bytes, BytesMut};
|
||||||
use futures::{Stream, StreamExt};
|
use futures::{Stream, StreamExt};
|
||||||
use regex::Regex;
|
|
||||||
use std::io;
|
use std::io;
|
||||||
use std::net::TcpListener;
|
use std::net::TcpListener;
|
||||||
use std::str;
|
use std::str;
|
||||||
@@ -35,7 +34,6 @@ use crate::basebackup;
|
|||||||
use crate::config::{PageServerConf, ProfilingConfig};
|
use crate::config::{PageServerConf, ProfilingConfig};
|
||||||
use crate::import_datadir::{import_basebackup_from_tar, import_wal_from_tar};
|
use crate::import_datadir::{import_basebackup_from_tar, import_wal_from_tar};
|
||||||
use crate::metrics::{LIVE_CONNECTIONS_COUNT, SMGR_QUERY_TIME};
|
use crate::metrics::{LIVE_CONNECTIONS_COUNT, SMGR_QUERY_TIME};
|
||||||
use crate::pgdatadir_mapping::LsnForTimestamp;
|
|
||||||
use crate::profiling::profpoint_start;
|
use crate::profiling::profpoint_start;
|
||||||
use crate::reltag::RelTag;
|
use crate::reltag::RelTag;
|
||||||
use crate::task_mgr;
|
use crate::task_mgr;
|
||||||
@@ -45,7 +43,6 @@ use crate::tenant_mgr;
|
|||||||
use crate::CheckpointConfig;
|
use crate::CheckpointConfig;
|
||||||
|
|
||||||
use postgres_ffi::pg_constants::DEFAULTTABLESPACE_OID;
|
use postgres_ffi::pg_constants::DEFAULTTABLESPACE_OID;
|
||||||
use postgres_ffi::to_pg_timestamp;
|
|
||||||
use postgres_ffi::BLCKSZ;
|
use postgres_ffi::BLCKSZ;
|
||||||
|
|
||||||
// Wrapped in libpq CopyData
|
// Wrapped in libpq CopyData
|
||||||
@@ -1062,33 +1059,6 @@ impl postgres_backend_async::Handler for PageServerHandler {
|
|||||||
Some(tenant.get_pitr_interval().as_secs().to_string().as_bytes()),
|
Some(tenant.get_pitr_interval().as_secs().to_string().as_bytes()),
|
||||||
]))?
|
]))?
|
||||||
.write_message(&BeMessage::CommandComplete(b"SELECT 1"))?;
|
.write_message(&BeMessage::CommandComplete(b"SELECT 1"))?;
|
||||||
} else if query_string.starts_with("get_lsn_by_timestamp ") {
|
|
||||||
// Locate LSN of last transaction with timestamp less or equal than sppecified
|
|
||||||
// TODO lazy static
|
|
||||||
let re = Regex::new(r"^get_lsn_by_timestamp ([[:xdigit:]]+) ([[:xdigit:]]+) '(.*)'$")
|
|
||||||
.unwrap();
|
|
||||||
let caps = re
|
|
||||||
.captures(query_string)
|
|
||||||
.with_context(|| format!("invalid get_lsn_by_timestamp: '{}'", query_string))?;
|
|
||||||
let tenant_id = TenantId::from_str(caps.get(1).unwrap().as_str())?;
|
|
||||||
let timeline_id = TimelineId::from_str(caps.get(2).unwrap().as_str())?;
|
|
||||||
let timestamp = humantime::parse_rfc3339(caps.get(3).unwrap().as_str())?;
|
|
||||||
let timestamp_pg = to_pg_timestamp(timestamp);
|
|
||||||
|
|
||||||
self.check_permission(Some(tenant_id))?;
|
|
||||||
|
|
||||||
let timeline = get_local_timeline(tenant_id, timeline_id)?;
|
|
||||||
pgb.write_message(&BeMessage::RowDescription(&[RowDescriptor::text_col(
|
|
||||||
b"lsn",
|
|
||||||
)]))?;
|
|
||||||
let result = match timeline.find_lsn_for_timestamp(timestamp_pg)? {
|
|
||||||
LsnForTimestamp::Present(lsn) => format!("{}", lsn),
|
|
||||||
LsnForTimestamp::Future(_lsn) => "future".into(),
|
|
||||||
LsnForTimestamp::Past(_lsn) => "past".into(),
|
|
||||||
LsnForTimestamp::NoData(_lsn) => "nodata".into(),
|
|
||||||
};
|
|
||||||
pgb.write_message(&BeMessage::DataRow(&[Some(result.as_bytes())]))?;
|
|
||||||
pgb.write_message(&BeMessage::CommandComplete(b"SELECT 1"))?;
|
|
||||||
} else {
|
} else {
|
||||||
bail!("unknown command");
|
bail!("unknown command");
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -639,6 +639,7 @@ pub fn spawn_storage_sync_task(
|
|||||||
(storage, remote_index_clone, sync_queue),
|
(storage, remote_index_clone, sync_queue),
|
||||||
max_sync_errors,
|
max_sync_errors,
|
||||||
)
|
)
|
||||||
|
.instrument(info_span!("storage_sync_loop"))
|
||||||
.await;
|
.await;
|
||||||
Ok(())
|
Ok(())
|
||||||
},
|
},
|
||||||
|
|||||||
+41
-42
@@ -45,6 +45,7 @@ use crate::tenant_config::TenantConfOpt;
|
|||||||
use crate::virtual_file::VirtualFile;
|
use crate::virtual_file::VirtualFile;
|
||||||
use crate::walredo::WalRedoManager;
|
use crate::walredo::WalRedoManager;
|
||||||
use crate::{CheckpointConfig, TEMP_FILE_SUFFIX};
|
use crate::{CheckpointConfig, TEMP_FILE_SUFFIX};
|
||||||
|
pub use pageserver_api::models::TenantState;
|
||||||
|
|
||||||
use toml_edit;
|
use toml_edit;
|
||||||
use utils::{
|
use utils::{
|
||||||
@@ -118,18 +119,6 @@ pub struct Tenant {
|
|||||||
upload_layers: bool,
|
upload_layers: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A state of a tenant in pageserver's memory.
|
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
|
|
||||||
pub enum TenantState {
|
|
||||||
/// Tenant is fully operational, its background jobs might be running or not.
|
|
||||||
Active { background_jobs_running: bool },
|
|
||||||
/// A tenant is recognized by pageserver, but not yet ready to operate:
|
|
||||||
/// e.g. not present locally and being downloaded or being read into memory from the file system.
|
|
||||||
Paused,
|
|
||||||
/// A tenant is recognized by the pageserver, but no longer used for any operations, as failed to get activated.
|
|
||||||
Broken,
|
|
||||||
}
|
|
||||||
|
|
||||||
/// A repository corresponds to one .neon directory. One repository holds multiple
|
/// A repository corresponds to one .neon directory. One repository holds multiple
|
||||||
/// timelines, forked off from the same initial call to 'initdb'.
|
/// timelines, forked off from the same initial call to 'initdb'.
|
||||||
impl Tenant {
|
impl Tenant {
|
||||||
@@ -400,16 +389,19 @@ impl Tenant {
|
|||||||
timeline_id,
|
timeline_id,
|
||||||
metadata.pg_version()
|
metadata.pg_version()
|
||||||
);
|
);
|
||||||
let timeline = self
|
let ancestor = metadata
|
||||||
.initialize_new_timeline(timeline_id, metadata, &mut timelines_accessor)
|
.ancestor_timeline()
|
||||||
.with_context(|| format!("Failed to initialize timeline {timeline_id}"))?;
|
.and_then(|ancestor_timeline_id| timelines_accessor.get(&ancestor_timeline_id))
|
||||||
|
.cloned();
|
||||||
match timelines_accessor.entry(timeline.timeline_id) {
|
match timelines_accessor.entry(timeline_id) {
|
||||||
Entry::Occupied(_) => bail!(
|
Entry::Occupied(_) => warn!(
|
||||||
"Found freshly initialized timeline {} in the tenant map",
|
"Timeline {}/{} already exists in the tenant map, skipping its initialization",
|
||||||
timeline.timeline_id
|
self.tenant_id, timeline_id
|
||||||
),
|
),
|
||||||
Entry::Vacant(v) => {
|
Entry::Vacant(v) => {
|
||||||
|
let timeline = self
|
||||||
|
.initialize_new_timeline(timeline_id, metadata, ancestor)
|
||||||
|
.with_context(|| format!("Failed to initialize timeline {timeline_id}"))?;
|
||||||
v.insert(timeline);
|
v.insert(timeline);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -609,21 +601,14 @@ impl Tenant {
|
|||||||
&self,
|
&self,
|
||||||
new_timeline_id: TimelineId,
|
new_timeline_id: TimelineId,
|
||||||
new_metadata: TimelineMetadata,
|
new_metadata: TimelineMetadata,
|
||||||
timelines: &mut MutexGuard<HashMap<TimelineId, Arc<Timeline>>>,
|
ancestor: Option<Arc<Timeline>>,
|
||||||
) -> anyhow::Result<Arc<Timeline>> {
|
) -> anyhow::Result<Arc<Timeline>> {
|
||||||
let ancestor = match new_metadata.ancestor_timeline() {
|
if let Some(ancestor_timeline_id) = new_metadata.ancestor_timeline() {
|
||||||
Some(ancestor_timeline_id) => Some(
|
anyhow::ensure!(
|
||||||
timelines
|
ancestor.is_some(),
|
||||||
.get(&ancestor_timeline_id)
|
"Timeline's {new_timeline_id} ancestor {ancestor_timeline_id} was not found"
|
||||||
.cloned()
|
)
|
||||||
.with_context(|| {
|
}
|
||||||
format!(
|
|
||||||
"Timeline's {new_timeline_id} ancestor {ancestor_timeline_id} was not found"
|
|
||||||
)
|
|
||||||
})?,
|
|
||||||
),
|
|
||||||
None => None,
|
|
||||||
};
|
|
||||||
|
|
||||||
let new_disk_consistent_lsn = new_metadata.disk_consistent_lsn();
|
let new_disk_consistent_lsn = new_metadata.disk_consistent_lsn();
|
||||||
let pg_version = new_metadata.pg_version();
|
let pg_version = new_metadata.pg_version();
|
||||||
@@ -1080,8 +1065,12 @@ impl Tenant {
|
|||||||
)
|
)
|
||||||
})?;
|
})?;
|
||||||
|
|
||||||
|
let ancestor = new_metadata
|
||||||
|
.ancestor_timeline()
|
||||||
|
.and_then(|ancestor_timeline_id| timelines.get(&ancestor_timeline_id))
|
||||||
|
.cloned();
|
||||||
let new_timeline = self
|
let new_timeline = self
|
||||||
.initialize_new_timeline(new_timeline_id, new_metadata, timelines)
|
.initialize_new_timeline(new_timeline_id, new_metadata, ancestor)
|
||||||
.with_context(|| {
|
.with_context(|| {
|
||||||
format!(
|
format!(
|
||||||
"Failed to initialize timeline {}/{}",
|
"Failed to initialize timeline {}/{}",
|
||||||
@@ -1105,12 +1094,22 @@ impl Tenant {
|
|||||||
|
|
||||||
/// Create the cluster temporarily in 'initdbpath' directory inside the repository
|
/// Create the cluster temporarily in 'initdbpath' directory inside the repository
|
||||||
/// to get bootstrap data for timeline initialization.
|
/// to get bootstrap data for timeline initialization.
|
||||||
fn run_initdb(conf: &'static PageServerConf, initdbpath: &Path, pg_version: u32) -> Result<()> {
|
fn run_initdb(
|
||||||
info!("running initdb in {}... ", initdbpath.display());
|
conf: &'static PageServerConf,
|
||||||
|
initdb_target_dir: &Path,
|
||||||
|
pg_version: u32,
|
||||||
|
) -> Result<()> {
|
||||||
|
let initdb_bin_path = conf.pg_bin_dir(pg_version).join("initdb");
|
||||||
|
let initdb_lib_dir = conf.pg_lib_dir(pg_version);
|
||||||
|
info!(
|
||||||
|
"running {} in {}, libdir: {}",
|
||||||
|
initdb_bin_path.display(),
|
||||||
|
initdb_target_dir.display(),
|
||||||
|
initdb_lib_dir.display(),
|
||||||
|
);
|
||||||
|
|
||||||
let initdb_path = conf.pg_bin_dir(pg_version).join("initdb");
|
let initdb_output = Command::new(initdb_bin_path)
|
||||||
let initdb_output = Command::new(initdb_path)
|
.args(&["-D", &initdb_target_dir.to_string_lossy()])
|
||||||
.args(&["-D", &initdbpath.to_string_lossy()])
|
|
||||||
.args(&["-U", &conf.superuser])
|
.args(&["-U", &conf.superuser])
|
||||||
.args(&["-E", "utf8"])
|
.args(&["-E", "utf8"])
|
||||||
.arg("--no-instructions")
|
.arg("--no-instructions")
|
||||||
@@ -1118,8 +1117,8 @@ fn run_initdb(conf: &'static PageServerConf, initdbpath: &Path, pg_version: u32)
|
|||||||
// so no need to fsync it
|
// so no need to fsync it
|
||||||
.arg("--no-sync")
|
.arg("--no-sync")
|
||||||
.env_clear()
|
.env_clear()
|
||||||
.env("LD_LIBRARY_PATH", conf.pg_lib_dir(pg_version))
|
.env("LD_LIBRARY_PATH", &initdb_lib_dir)
|
||||||
.env("DYLD_LIBRARY_PATH", conf.pg_lib_dir(pg_version))
|
.env("DYLD_LIBRARY_PATH", &initdb_lib_dir)
|
||||||
.stdout(Stdio::null())
|
.stdout(Stdio::null())
|
||||||
.output()
|
.output()
|
||||||
.context("failed to execute initdb")?;
|
.context("failed to execute initdb")?;
|
||||||
|
|||||||
@@ -556,7 +556,7 @@ impl DeltaLayer {
|
|||||||
|
|
||||||
/// Create a DeltaLayer struct representing an existing file on disk.
|
/// Create a DeltaLayer struct representing an existing file on disk.
|
||||||
///
|
///
|
||||||
/// This variant is only used for debugging purposes, by the 'dump_layerfile' binary.
|
/// This variant is only used for debugging purposes, by the 'pageserver_binutils' binary.
|
||||||
pub fn new_for_path<F>(path: &Path, file: F) -> Result<Self>
|
pub fn new_for_path<F>(path: &Path, file: F) -> Result<Self>
|
||||||
where
|
where
|
||||||
F: FileExt,
|
F: FileExt,
|
||||||
|
|||||||
@@ -177,7 +177,7 @@ impl fmt::Display for ImageFileName {
|
|||||||
///
|
///
|
||||||
/// This is used by DeltaLayer and ImageLayer. Normally, this holds a reference to the
|
/// This is used by DeltaLayer and ImageLayer. Normally, this holds a reference to the
|
||||||
/// global config, and paths to layer files are constructed using the tenant/timeline
|
/// global config, and paths to layer files are constructed using the tenant/timeline
|
||||||
/// path from the config. But in the 'dump_layerfile' binary, we need to construct a Layer
|
/// path from the config. But in the 'pageserver_binutils' binary, we need to construct a Layer
|
||||||
/// struct for a file on disk, without having a page server running, so that we have no
|
/// struct for a file on disk, without having a page server running, so that we have no
|
||||||
/// config. In that case, we use the Path variant to hold the full path to the file on
|
/// config. In that case, we use the Path variant to hold the full path to the file on
|
||||||
/// disk.
|
/// disk.
|
||||||
|
|||||||
@@ -357,7 +357,7 @@ impl ImageLayer {
|
|||||||
|
|
||||||
/// Create an ImageLayer struct representing an existing file on disk.
|
/// Create an ImageLayer struct representing an existing file on disk.
|
||||||
///
|
///
|
||||||
/// This variant is only used for debugging purposes, by the 'dump_layerfile' binary.
|
/// This variant is only used for debugging purposes, by the 'pageserver_binutils' binary.
|
||||||
pub fn new_for_path<F>(path: &Path, file: F) -> Result<ImageLayer>
|
pub fn new_for_path<F>(path: &Path, file: F) -> Result<ImageLayer>
|
||||||
where
|
where
|
||||||
F: std::os::unix::prelude::FileExt,
|
F: std::os::unix::prelude::FileExt,
|
||||||
|
|||||||
@@ -627,7 +627,7 @@ impl Timeline {
|
|||||||
.unwrap_or(self.conf.default_tenant_conf.max_lsn_wal_lag);
|
.unwrap_or(self.conf.default_tenant_conf.max_lsn_wal_lag);
|
||||||
drop(tenant_conf_guard);
|
drop(tenant_conf_guard);
|
||||||
let self_clone = Arc::clone(self);
|
let self_clone = Arc::clone(self);
|
||||||
let _ = spawn_connection_manager_task(
|
spawn_connection_manager_task(
|
||||||
self.conf.broker_etcd_prefix.clone(),
|
self.conf.broker_etcd_prefix.clone(),
|
||||||
self_clone,
|
self_clone,
|
||||||
walreceiver_connect_timeout,
|
walreceiver_connect_timeout,
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ pub mod defaults {
|
|||||||
// This parameter determines L1 layer file size.
|
// This parameter determines L1 layer file size.
|
||||||
pub const DEFAULT_COMPACTION_TARGET_SIZE: u64 = 128 * 1024 * 1024;
|
pub const DEFAULT_COMPACTION_TARGET_SIZE: u64 = 128 * 1024 * 1024;
|
||||||
|
|
||||||
pub const DEFAULT_COMPACTION_PERIOD: &str = "1 s";
|
pub const DEFAULT_COMPACTION_PERIOD: &str = "20 s";
|
||||||
pub const DEFAULT_COMPACTION_THRESHOLD: usize = 10;
|
pub const DEFAULT_COMPACTION_THRESHOLD: usize = 10;
|
||||||
|
|
||||||
pub const DEFAULT_GC_HORIZON: u64 = 64 * 1024 * 1024;
|
pub const DEFAULT_GC_HORIZON: u64 = 64 * 1024 * 1024;
|
||||||
|
|||||||
@@ -107,6 +107,13 @@ pub fn init_tenant_mgr(
|
|||||||
/// Ignores other timelines that might be present for tenant, but were not passed as a parameter.
|
/// Ignores other timelines that might be present for tenant, but were not passed as a parameter.
|
||||||
/// Attempts to load as many entites as possible: if a certain timeline fails during the load, the tenant is marked as "Broken",
|
/// Attempts to load as many entites as possible: if a certain timeline fails during the load, the tenant is marked as "Broken",
|
||||||
/// and the load continues.
|
/// and the load continues.
|
||||||
|
///
|
||||||
|
/// For successful tenant attach, it first has to have a `timelines/` subdirectory and a tenant config file that's loaded into memory successfully.
|
||||||
|
/// If either of the conditions fails, the tenant will be added to memory with [`TenantState::Broken`] state, otherwise we start to load its timelines.
|
||||||
|
/// Alternatively, tenant is considered loaded successfully, if it's already in pageserver's memory (i.e. was loaded already before).
|
||||||
|
///
|
||||||
|
/// Attach happens on startup and sucessful timeline downloads
|
||||||
|
/// (some subset of timeline files, always including its metadata, after which the new one needs to be registered).
|
||||||
pub fn attach_local_tenants(
|
pub fn attach_local_tenants(
|
||||||
conf: &'static PageServerConf,
|
conf: &'static PageServerConf,
|
||||||
remote_index: &RemoteIndex,
|
remote_index: &RemoteIndex,
|
||||||
@@ -122,18 +129,20 @@ pub fn attach_local_tenants(
|
|||||||
);
|
);
|
||||||
debug!("Timelines to attach: {local_timelines:?}");
|
debug!("Timelines to attach: {local_timelines:?}");
|
||||||
|
|
||||||
let tenant = load_local_tenant(conf, tenant_id, remote_index);
|
let mut tenants_accessor = tenants_state::write_tenants();
|
||||||
{
|
let tenant = match tenants_accessor.entry(tenant_id) {
|
||||||
match tenants_state::write_tenants().entry(tenant_id) {
|
hash_map::Entry::Occupied(o) => {
|
||||||
hash_map::Entry::Occupied(_) => {
|
info!("Tenant {tenant_id} was found in pageserver's memory");
|
||||||
error!("Cannot attach tenant {tenant_id}: there's already an entry in the tenant state");
|
Arc::clone(o.get())
|
||||||
continue;
|
|
||||||
}
|
|
||||||
hash_map::Entry::Vacant(v) => {
|
|
||||||
v.insert(Arc::clone(&tenant));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
hash_map::Entry::Vacant(v) => {
|
||||||
|
info!("Tenant {tenant_id} was not found in pageserver's memory, loading it");
|
||||||
|
let tenant = load_local_tenant(conf, tenant_id, remote_index);
|
||||||
|
v.insert(Arc::clone(&tenant));
|
||||||
|
tenant
|
||||||
|
}
|
||||||
|
};
|
||||||
|
drop(tenants_accessor);
|
||||||
|
|
||||||
if tenant.current_state() == TenantState::Broken {
|
if tenant.current_state() == TenantState::Broken {
|
||||||
warn!("Skipping timeline load for broken tenant {tenant_id}")
|
warn!("Skipping timeline load for broken tenant {tenant_id}")
|
||||||
@@ -168,16 +177,28 @@ fn load_local_tenant(
|
|||||||
remote_index.clone(),
|
remote_index.clone(),
|
||||||
conf.remote_storage_config.is_some(),
|
conf.remote_storage_config.is_some(),
|
||||||
));
|
));
|
||||||
match Tenant::load_tenant_config(conf, tenant_id) {
|
|
||||||
Ok(tenant_conf) => {
|
let tenant_timelines_dir = conf.timelines_path(&tenant_id);
|
||||||
tenant.update_tenant_config(tenant_conf);
|
if !tenant_timelines_dir.is_dir() {
|
||||||
tenant.activate(false);
|
error!(
|
||||||
}
|
"Tenant {} has no timelines directory at {}",
|
||||||
Err(e) => {
|
tenant_id,
|
||||||
error!("Failed to read config for tenant {tenant_id}, disabling tenant: {e:?}");
|
tenant_timelines_dir.display()
|
||||||
tenant.set_state(TenantState::Broken);
|
);
|
||||||
|
tenant.set_state(TenantState::Broken);
|
||||||
|
} else {
|
||||||
|
match Tenant::load_tenant_config(conf, tenant_id) {
|
||||||
|
Ok(tenant_conf) => {
|
||||||
|
tenant.update_tenant_config(tenant_conf);
|
||||||
|
tenant.activate(false);
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
error!("Failed to read config for tenant {tenant_id}, disabling tenant: {e:?}");
|
||||||
|
tenant.set_state(TenantState::Broken);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
tenant
|
tenant
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -625,14 +646,10 @@ fn collect_timelines_for_tenant(
|
|||||||
}
|
}
|
||||||
|
|
||||||
if tenant_timelines.is_empty() {
|
if tenant_timelines.is_empty() {
|
||||||
match remove_if_empty(&timelines_dir) {
|
// this is normal, we've removed all broken, empty and temporary timeline dirs
|
||||||
Ok(true) => info!(
|
// but should allow the tenant to stay functional and allow creating new timelines
|
||||||
"Removed empty tenant timelines directory {}",
|
// on a restart, we require tenants to have the timelines dir, so leave it on disk
|
||||||
timelines_dir.display()
|
debug!("Tenant {tenant_id} has no timelines loaded");
|
||||||
),
|
|
||||||
Ok(false) => (),
|
|
||||||
Err(e) => error!("Failed to remove empty tenant timelines directory: {e:?}"),
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
Ok((tenant_id, tenant_timelines))
|
Ok((tenant_id, tenant_timelines))
|
||||||
|
|||||||
@@ -12,6 +12,8 @@ use chrono::{NaiveDateTime, Utc};
|
|||||||
use fail::fail_point;
|
use fail::fail_point;
|
||||||
use futures::StreamExt;
|
use futures::StreamExt;
|
||||||
use postgres::{SimpleQueryMessage, SimpleQueryRow};
|
use postgres::{SimpleQueryMessage, SimpleQueryRow};
|
||||||
|
use postgres_ffi::v14::xlog_utils::normalize_lsn;
|
||||||
|
use postgres_ffi::WAL_SEGMENT_SIZE;
|
||||||
use postgres_protocol::message::backend::ReplicationMessage;
|
use postgres_protocol::message::backend::ReplicationMessage;
|
||||||
use postgres_types::PgLsn;
|
use postgres_types::PgLsn;
|
||||||
use tokio::{pin, select, sync::watch, time};
|
use tokio::{pin, select, sync::watch, time};
|
||||||
@@ -156,6 +158,14 @@ pub async fn handle_walreceiver_connection(
|
|||||||
// There might be some padding after the last full record, skip it.
|
// There might be some padding after the last full record, skip it.
|
||||||
startpoint += startpoint.calc_padding(8u32);
|
startpoint += startpoint.calc_padding(8u32);
|
||||||
|
|
||||||
|
// If the starting point is at a WAL page boundary, skip past the page header. We don't need the page headers
|
||||||
|
// for anything, and in some corner cases, the compute node might have never generated the WAL for page headers
|
||||||
|
//. That happens if you create a branch at page boundary: the start point of the branch is at the page boundary,
|
||||||
|
// but when the compute node first starts on the branch, we normalize the first REDO position to just after the page
|
||||||
|
// header (see generate_pg_control()), so the WAL for the page header is never streamed from the compute node
|
||||||
|
// to the safekeepers.
|
||||||
|
startpoint = normalize_lsn(startpoint, WAL_SEGMENT_SIZE);
|
||||||
|
|
||||||
info!("last_record_lsn {last_rec_lsn} starting replication from {startpoint}, safekeeper is at {end_of_wal}...");
|
info!("last_record_lsn {last_rec_lsn} starting replication from {startpoint}, safekeeper is at {end_of_wal}...");
|
||||||
|
|
||||||
let query = format!("START_REPLICATION PHYSICAL {startpoint}");
|
let query = format!("START_REPLICATION PHYSICAL {startpoint}");
|
||||||
|
|||||||
+4
-1
@@ -5,7 +5,7 @@ edition = "2021"
|
|||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
anyhow = "1.0"
|
anyhow = "1.0"
|
||||||
async-trait = "0.1"
|
atty = "0.2.14"
|
||||||
base64 = "0.13.0"
|
base64 = "0.13.0"
|
||||||
bstr = "0.2.17"
|
bstr = "0.2.17"
|
||||||
bytes = { version = "1.0.1", features = ['serde'] }
|
bytes = { version = "1.0.1", features = ['serde'] }
|
||||||
@@ -35,6 +35,8 @@ thiserror = "1.0.30"
|
|||||||
tokio = { version = "1.17", features = ["macros"] }
|
tokio = { version = "1.17", features = ["macros"] }
|
||||||
tokio-postgres = { git = "https://github.com/neondatabase/rust-postgres.git", rev="d052ee8b86fff9897c77b0fe89ea9daba0e1fa38" }
|
tokio-postgres = { git = "https://github.com/neondatabase/rust-postgres.git", rev="d052ee8b86fff9897c77b0fe89ea9daba0e1fa38" }
|
||||||
tokio-rustls = "0.23.0"
|
tokio-rustls = "0.23.0"
|
||||||
|
tracing = "0.1.36"
|
||||||
|
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
||||||
url = "2.2.2"
|
url = "2.2.2"
|
||||||
uuid = { version = "0.8.2", features = ["v4", "serde"]}
|
uuid = { version = "0.8.2", features = ["v4", "serde"]}
|
||||||
x509-parser = "0.13.2"
|
x509-parser = "0.13.2"
|
||||||
@@ -44,6 +46,7 @@ metrics = { path = "../libs/metrics" }
|
|||||||
workspace_hack = { version = "0.1", path = "../workspace_hack" }
|
workspace_hack = { version = "0.1", path = "../workspace_hack" }
|
||||||
|
|
||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
|
async-trait = "0.1"
|
||||||
rcgen = "0.8.14"
|
rcgen = "0.8.14"
|
||||||
rstest = "0.12"
|
rstest = "0.12"
|
||||||
tokio-postgres-rustls = "0.9.0"
|
tokio-postgres-rustls = "0.9.0"
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ use once_cell::sync::Lazy;
|
|||||||
use serde::{Deserialize, Serialize};
|
use serde::{Deserialize, Serialize};
|
||||||
use std::borrow::Cow;
|
use std::borrow::Cow;
|
||||||
use tokio::io::{AsyncRead, AsyncWrite};
|
use tokio::io::{AsyncRead, AsyncWrite};
|
||||||
|
use tracing::{info, warn};
|
||||||
|
|
||||||
static CPLANE_WAITERS: Lazy<Waiters<mgmt::ComputeReady>> = Lazy::new(Default::default);
|
static CPLANE_WAITERS: Lazy<Waiters<mgmt::ComputeReady>> = Lazy::new(Default::default);
|
||||||
|
|
||||||
@@ -171,6 +172,8 @@ impl BackendType<'_, ClientCredentials<'_>> {
|
|||||||
// support SNI or other means of passing the project name.
|
// support SNI or other means of passing the project name.
|
||||||
// We now expect to see a very specific payload in the place of password.
|
// We now expect to see a very specific payload in the place of password.
|
||||||
if creds.project().is_none() {
|
if creds.project().is_none() {
|
||||||
|
warn!("project name not specified, resorting to the password hack auth flow");
|
||||||
|
|
||||||
let payload = AuthFlow::new(client)
|
let payload = AuthFlow::new(client)
|
||||||
.begin(auth::PasswordHack)
|
.begin(auth::PasswordHack)
|
||||||
.await?
|
.await?
|
||||||
@@ -179,6 +182,7 @@ impl BackendType<'_, ClientCredentials<'_>> {
|
|||||||
|
|
||||||
// Finally we may finish the initialization of `creds`.
|
// Finally we may finish the initialization of `creds`.
|
||||||
// TODO: add missing type safety to ClientCredentials.
|
// TODO: add missing type safety to ClientCredentials.
|
||||||
|
info!(project = &payload.project, "received missing parameter");
|
||||||
creds.project = Some(payload.project.into());
|
creds.project = Some(payload.project.into());
|
||||||
|
|
||||||
let mut config = match &self {
|
let mut config = match &self {
|
||||||
@@ -196,6 +200,7 @@ impl BackendType<'_, ClientCredentials<'_>> {
|
|||||||
// We should use a password from payload as well.
|
// We should use a password from payload as well.
|
||||||
config.password(payload.password);
|
config.password(payload.password);
|
||||||
|
|
||||||
|
info!("user successfully authenticated (using the password hack)");
|
||||||
return Ok(compute::NodeInfo {
|
return Ok(compute::NodeInfo {
|
||||||
reported_auth_ok: false,
|
reported_auth_ok: false,
|
||||||
config,
|
config,
|
||||||
@@ -203,19 +208,31 @@ impl BackendType<'_, ClientCredentials<'_>> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
match self {
|
let res = match self {
|
||||||
Console(endpoint, creds) => {
|
Console(endpoint, creds) => {
|
||||||
|
info!(
|
||||||
|
user = creds.user,
|
||||||
|
project = creds.project(),
|
||||||
|
"performing authentication using the console"
|
||||||
|
);
|
||||||
console::Api::new(&endpoint, extra, &creds)
|
console::Api::new(&endpoint, extra, &creds)
|
||||||
.handle_user(client)
|
.handle_user(client)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
Postgres(endpoint, creds) => {
|
Postgres(endpoint, creds) => {
|
||||||
|
info!("performing mock authentication using a local postgres instance");
|
||||||
postgres::Api::new(&endpoint, &creds)
|
postgres::Api::new(&endpoint, &creds)
|
||||||
.handle_user(client)
|
.handle_user(client)
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
// NOTE: this auth backend doesn't use client credentials.
|
// NOTE: this auth backend doesn't use client credentials.
|
||||||
Link(url) => link::handle_user(&url, client).await,
|
Link(url) => {
|
||||||
}
|
info!("performing link authentication");
|
||||||
|
link::handle_user(&url, client).await
|
||||||
|
}
|
||||||
|
}?;
|
||||||
|
|
||||||
|
info!("user successfully authenticated");
|
||||||
|
Ok(res)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ use serde::{Deserialize, Serialize};
|
|||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
use thiserror::Error;
|
use thiserror::Error;
|
||||||
use tokio::io::{AsyncRead, AsyncWrite};
|
use tokio::io::{AsyncRead, AsyncWrite};
|
||||||
|
use tracing::info;
|
||||||
|
|
||||||
const REQUEST_FAILED: &str = "Console request failed";
|
const REQUEST_FAILED: &str = "Console request failed";
|
||||||
|
|
||||||
@@ -148,10 +149,11 @@ impl<'a> Api<'a> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn get_auth_info(&self) -> Result<AuthInfo, GetAuthInfoError> {
|
async fn get_auth_info(&self) -> Result<AuthInfo, GetAuthInfoError> {
|
||||||
|
let request_id = uuid::Uuid::new_v4().to_string();
|
||||||
let req = self
|
let req = self
|
||||||
.endpoint
|
.endpoint
|
||||||
.get("proxy_get_role_secret")
|
.get("proxy_get_role_secret")
|
||||||
.header("X-Request-ID", uuid::Uuid::new_v4().to_string())
|
.header("X-Request-ID", &request_id)
|
||||||
.query(&[("session_id", self.extra.session_id)])
|
.query(&[("session_id", self.extra.session_id)])
|
||||||
.query(&[
|
.query(&[
|
||||||
("application_name", self.extra.application_name),
|
("application_name", self.extra.application_name),
|
||||||
@@ -160,9 +162,7 @@ impl<'a> Api<'a> {
|
|||||||
])
|
])
|
||||||
.build()?;
|
.build()?;
|
||||||
|
|
||||||
// TODO: use a proper logger
|
info!(id = request_id, url = req.url().as_str(), "request");
|
||||||
println!("cplane request: {}", req.url());
|
|
||||||
|
|
||||||
let resp = self.endpoint.execute(req).await?;
|
let resp = self.endpoint.execute(req).await?;
|
||||||
if !resp.status().is_success() {
|
if !resp.status().is_success() {
|
||||||
return Err(TransportError::HttpStatus(resp.status()).into());
|
return Err(TransportError::HttpStatus(resp.status()).into());
|
||||||
@@ -177,10 +177,11 @@ impl<'a> Api<'a> {
|
|||||||
|
|
||||||
/// Wake up the compute node and return the corresponding connection info.
|
/// Wake up the compute node and return the corresponding connection info.
|
||||||
pub(super) async fn wake_compute(&self) -> Result<ComputeConnCfg, WakeComputeError> {
|
pub(super) async fn wake_compute(&self) -> Result<ComputeConnCfg, WakeComputeError> {
|
||||||
|
let request_id = uuid::Uuid::new_v4().to_string();
|
||||||
let req = self
|
let req = self
|
||||||
.endpoint
|
.endpoint
|
||||||
.get("proxy_wake_compute")
|
.get("proxy_wake_compute")
|
||||||
.header("X-Request-ID", uuid::Uuid::new_v4().to_string())
|
.header("X-Request-ID", &request_id)
|
||||||
.query(&[("session_id", self.extra.session_id)])
|
.query(&[("session_id", self.extra.session_id)])
|
||||||
.query(&[
|
.query(&[
|
||||||
("application_name", self.extra.application_name),
|
("application_name", self.extra.application_name),
|
||||||
@@ -188,9 +189,7 @@ impl<'a> Api<'a> {
|
|||||||
])
|
])
|
||||||
.build()?;
|
.build()?;
|
||||||
|
|
||||||
// TODO: use a proper logger
|
info!(id = request_id, url = req.url().as_str(), "request");
|
||||||
println!("cplane request: {}", req.url());
|
|
||||||
|
|
||||||
let resp = self.endpoint.execute(req).await?;
|
let resp = self.endpoint.execute(req).await?;
|
||||||
if !resp.status().is_success() {
|
if !resp.status().is_success() {
|
||||||
return Err(TransportError::HttpStatus(resp.status()).into());
|
return Err(TransportError::HttpStatus(resp.status()).into());
|
||||||
@@ -227,15 +226,18 @@ where
|
|||||||
GetAuthInfo: Future<Output = Result<AuthInfo, GetAuthInfoError>>,
|
GetAuthInfo: Future<Output = Result<AuthInfo, GetAuthInfoError>>,
|
||||||
WakeCompute: Future<Output = Result<ComputeConnCfg, WakeComputeError>>,
|
WakeCompute: Future<Output = Result<ComputeConnCfg, WakeComputeError>>,
|
||||||
{
|
{
|
||||||
|
info!("fetching user's authentication info");
|
||||||
let auth_info = get_auth_info(endpoint).await?;
|
let auth_info = get_auth_info(endpoint).await?;
|
||||||
|
|
||||||
let flow = AuthFlow::new(client);
|
let flow = AuthFlow::new(client);
|
||||||
let scram_keys = match auth_info {
|
let scram_keys = match auth_info {
|
||||||
AuthInfo::Md5(_) => {
|
AuthInfo::Md5(_) => {
|
||||||
// TODO: decide if we should support MD5 in api v2
|
// TODO: decide if we should support MD5 in api v2
|
||||||
|
info!("auth endpoint chooses MD5");
|
||||||
return Err(auth::AuthError::bad_auth_method("MD5"));
|
return Err(auth::AuthError::bad_auth_method("MD5"));
|
||||||
}
|
}
|
||||||
AuthInfo::Scram(secret) => {
|
AuthInfo::Scram(secret) => {
|
||||||
|
info!("auth endpoint chooses SCRAM");
|
||||||
let scram = auth::Scram(&secret);
|
let scram = auth::Scram(&secret);
|
||||||
Some(compute::ScramKeys {
|
Some(compute::ScramKeys {
|
||||||
client_key: flow.begin(scram).await?.authenticate().await?.as_bytes(),
|
client_key: flow.begin(scram).await?.authenticate().await?.as_bytes(),
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
use crate::{auth, compute, error::UserFacingError, stream::PqStream, waiters};
|
use crate::{auth, compute, error::UserFacingError, stream::PqStream, waiters};
|
||||||
use thiserror::Error;
|
use thiserror::Error;
|
||||||
use tokio::io::{AsyncRead, AsyncWrite};
|
use tokio::io::{AsyncRead, AsyncWrite};
|
||||||
|
use tracing::info;
|
||||||
use utils::pq_proto::{BeMessage as Be, BeParameterStatusMessage};
|
use utils::pq_proto::{BeMessage as Be, BeParameterStatusMessage};
|
||||||
|
|
||||||
#[derive(Debug, Error)]
|
#[derive(Debug, Error)]
|
||||||
@@ -53,14 +54,16 @@ pub async fn handle_user(
|
|||||||
let greeting = hello_message(link_uri, &psql_session_id);
|
let greeting = hello_message(link_uri, &psql_session_id);
|
||||||
|
|
||||||
let db_info = super::with_waiter(psql_session_id, |waiter| async {
|
let db_info = super::with_waiter(psql_session_id, |waiter| async {
|
||||||
// Give user a URL to spawn a new database
|
// Give user a URL to spawn a new database.
|
||||||
|
info!("sending the auth URL to the user");
|
||||||
client
|
client
|
||||||
.write_message_noflush(&Be::AuthenticationOk)?
|
.write_message_noflush(&Be::AuthenticationOk)?
|
||||||
.write_message_noflush(&BeParameterStatusMessage::encoding())?
|
.write_message_noflush(&BeParameterStatusMessage::encoding())?
|
||||||
.write_message(&Be::NoticeResponse(&greeting))
|
.write_message(&Be::NoticeResponse(&greeting))
|
||||||
.await?;
|
.await?;
|
||||||
|
|
||||||
// Wait for web console response (see `mgmt`)
|
// Wait for web console response (see `mgmt`).
|
||||||
|
info!("waiting for console's reply...");
|
||||||
waiter.await?.map_err(LinkAuthError::AuthFailed)
|
waiter.await?.map_err(LinkAuthError::AuthFailed)
|
||||||
})
|
})
|
||||||
.await?;
|
.await?;
|
||||||
|
|||||||
@@ -3,6 +3,7 @@
|
|||||||
use crate::error::UserFacingError;
|
use crate::error::UserFacingError;
|
||||||
use std::borrow::Cow;
|
use std::borrow::Cow;
|
||||||
use thiserror::Error;
|
use thiserror::Error;
|
||||||
|
use tracing::info;
|
||||||
use utils::pq_proto::StartupMessageParams;
|
use utils::pq_proto::StartupMessageParams;
|
||||||
|
|
||||||
#[derive(Debug, Error, PartialEq, Eq, Clone)]
|
#[derive(Debug, Error, PartialEq, Eq, Clone)]
|
||||||
@@ -82,6 +83,13 @@ impl<'a> ClientCredentials<'a> {
|
|||||||
}
|
}
|
||||||
.transpose()?;
|
.transpose()?;
|
||||||
|
|
||||||
|
info!(
|
||||||
|
user = user,
|
||||||
|
dbname = dbname,
|
||||||
|
project = project.as_deref(),
|
||||||
|
"credentials"
|
||||||
|
);
|
||||||
|
|
||||||
Ok(Self {
|
Ok(Self {
|
||||||
user,
|
user,
|
||||||
dbname,
|
dbname,
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ use parking_lot::Mutex;
|
|||||||
use std::net::SocketAddr;
|
use std::net::SocketAddr;
|
||||||
use tokio::net::TcpStream;
|
use tokio::net::TcpStream;
|
||||||
use tokio_postgres::{CancelToken, NoTls};
|
use tokio_postgres::{CancelToken, NoTls};
|
||||||
|
use tracing::info;
|
||||||
use utils::pq_proto::CancelKeyData;
|
use utils::pq_proto::CancelKeyData;
|
||||||
|
|
||||||
/// Enables serving `CancelRequest`s.
|
/// Enables serving `CancelRequest`s.
|
||||||
@@ -18,8 +19,9 @@ impl CancelMap {
|
|||||||
.lock()
|
.lock()
|
||||||
.get(&key)
|
.get(&key)
|
||||||
.and_then(|x| x.clone())
|
.and_then(|x| x.clone())
|
||||||
.with_context(|| format!("unknown session: {:?}", key))?;
|
.with_context(|| format!("query cancellation key not found: {key}"))?;
|
||||||
|
|
||||||
|
info!("cancelling query per user's request using key {key}");
|
||||||
cancel_closure.try_cancel_query().await
|
cancel_closure.try_cancel_query().await
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -41,14 +43,16 @@ impl CancelMap {
|
|||||||
self.0
|
self.0
|
||||||
.lock()
|
.lock()
|
||||||
.try_insert(key, None)
|
.try_insert(key, None)
|
||||||
.map_err(|_| anyhow!("session already exists: {:?}", key))?;
|
.map_err(|_| anyhow!("query cancellation key already exists: {key}"))?;
|
||||||
|
|
||||||
// This will guarantee that the session gets dropped
|
// This will guarantee that the session gets dropped
|
||||||
// as soon as the future is finished.
|
// as soon as the future is finished.
|
||||||
scopeguard::defer! {
|
scopeguard::defer! {
|
||||||
self.0.lock().remove(&key);
|
self.0.lock().remove(&key);
|
||||||
|
info!("dropped query cancellation key {key}");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
info!("registered new query cancellation key {key}");
|
||||||
let session = Session::new(key, self);
|
let session = Session::new(key, self);
|
||||||
f(session).await
|
f(session).await
|
||||||
}
|
}
|
||||||
@@ -102,10 +106,13 @@ impl<'a> Session<'a> {
|
|||||||
fn new(key: CancelKeyData, cancel_map: &'a CancelMap) -> Self {
|
fn new(key: CancelKeyData, cancel_map: &'a CancelMap) -> Self {
|
||||||
Self { key, cancel_map }
|
Self { key, cancel_map }
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Session<'_> {
|
||||||
/// Store the cancel token for the given session.
|
/// Store the cancel token for the given session.
|
||||||
/// This enables query cancellation in [`crate::proxy::handshake`].
|
/// This enables query cancellation in [`crate::proxy::handshake`].
|
||||||
pub fn enable_query_cancellation(self, cancel_closure: CancelClosure) -> CancelKeyData {
|
pub fn enable_query_cancellation(self, cancel_closure: CancelClosure) -> CancelKeyData {
|
||||||
|
info!("enabling query cancellation for this session");
|
||||||
self.cancel_map
|
self.cancel_map
|
||||||
.0
|
.0
|
||||||
.lock()
|
.lock()
|
||||||
@@ -129,13 +136,13 @@ mod tests {
|
|||||||
assert!(CANCEL_MAP.contains(&session));
|
assert!(CANCEL_MAP.contains(&session));
|
||||||
|
|
||||||
tx.send(()).expect("failed to send");
|
tx.send(()).expect("failed to send");
|
||||||
let () = futures::future::pending().await; // sleep forever
|
futures::future::pending::<()>().await; // sleep forever
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}));
|
}));
|
||||||
|
|
||||||
// Wait until the task has been spawned.
|
// Wait until the task has been spawned.
|
||||||
let () = rx.await.context("failed to hear from the task")?;
|
rx.await.context("failed to hear from the task")?;
|
||||||
|
|
||||||
// Drop the session's entry by cancelling the task.
|
// Drop the session's entry by cancelling the task.
|
||||||
task.abort();
|
task.abort();
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ use std::{io, net::SocketAddr};
|
|||||||
use thiserror::Error;
|
use thiserror::Error;
|
||||||
use tokio::net::TcpStream;
|
use tokio::net::TcpStream;
|
||||||
use tokio_postgres::NoTls;
|
use tokio_postgres::NoTls;
|
||||||
|
use tracing::{error, info};
|
||||||
use utils::pq_proto::StartupMessageParams;
|
use utils::pq_proto::StartupMessageParams;
|
||||||
|
|
||||||
#[derive(Debug, Error)]
|
#[derive(Debug, Error)]
|
||||||
@@ -54,6 +55,7 @@ impl NodeInfo {
|
|||||||
use tokio_postgres::config::Host;
|
use tokio_postgres::config::Host;
|
||||||
|
|
||||||
let connect_once = |host, port| {
|
let connect_once = |host, port| {
|
||||||
|
info!("trying to connect to a compute node at {host}:{port}");
|
||||||
TcpStream::connect((host, port)).and_then(|socket| async {
|
TcpStream::connect((host, port)).and_then(|socket| async {
|
||||||
let socket_addr = socket.peer_addr()?;
|
let socket_addr = socket.peer_addr()?;
|
||||||
// This prevents load balancer from severing the connection.
|
// This prevents load balancer from severing the connection.
|
||||||
@@ -72,7 +74,11 @@ impl NodeInfo {
|
|||||||
if ports.len() > 1 && ports.len() != hosts.len() {
|
if ports.len() > 1 && ports.len() != hosts.len() {
|
||||||
return Err(io::Error::new(
|
return Err(io::Error::new(
|
||||||
io::ErrorKind::Other,
|
io::ErrorKind::Other,
|
||||||
format!("couldn't connect: bad compute config, ports and hosts entries' count does not match: {:?}", self.config),
|
format!(
|
||||||
|
"couldn't connect: bad compute config, \
|
||||||
|
ports and hosts entries' count does not match: {:?}",
|
||||||
|
self.config
|
||||||
|
),
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -88,7 +94,7 @@ impl NodeInfo {
|
|||||||
Ok(socket) => return Ok(socket),
|
Ok(socket) => return Ok(socket),
|
||||||
Err(err) => {
|
Err(err) => {
|
||||||
// We can't throw an error here, as there might be more hosts to try.
|
// We can't throw an error here, as there might be more hosts to try.
|
||||||
println!("failed to connect to compute `{host}:{port}`: {err}");
|
error!("failed to connect to a compute node at {host}:{port}: {err}");
|
||||||
connection_error = Some(err);
|
connection_error = Some(err);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -160,8 +166,8 @@ impl NodeInfo {
|
|||||||
.ok_or(ConnectionError::FailedToFetchPgVersion)?
|
.ok_or(ConnectionError::FailedToFetchPgVersion)?
|
||||||
.into();
|
.into();
|
||||||
|
|
||||||
|
info!("connected to user's compute node at {socket_addr}");
|
||||||
let cancel_closure = CancelClosure::new(socket_addr, client.cancel_token());
|
let cancel_closure = CancelClosure::new(socket_addr, client.cancel_token());
|
||||||
|
|
||||||
let db = PostgresConnection { stream, version };
|
let db = PostgresConnection { stream, version };
|
||||||
|
|
||||||
Ok((db, cancel_closure))
|
Ok((db, cancel_closure))
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
use anyhow::anyhow;
|
use anyhow::anyhow;
|
||||||
use hyper::{Body, Request, Response, StatusCode};
|
use hyper::{Body, Request, Response, StatusCode};
|
||||||
use std::net::TcpListener;
|
use std::net::TcpListener;
|
||||||
|
use tracing::info;
|
||||||
use utils::http::{endpoint, error::ApiError, json::json_response, RouterBuilder, RouterService};
|
use utils::http::{endpoint, error::ApiError, json::json_response, RouterBuilder, RouterService};
|
||||||
|
|
||||||
async fn status_handler(_: Request<Body>) -> Result<Response<Body>, ApiError> {
|
async fn status_handler(_: Request<Body>) -> Result<Response<Body>, ApiError> {
|
||||||
@@ -12,9 +13,9 @@ fn make_router() -> RouterBuilder<hyper::Body, ApiError> {
|
|||||||
router.get("/v1/status", status_handler)
|
router.get("/v1/status", status_handler)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn thread_main(http_listener: TcpListener) -> anyhow::Result<()> {
|
pub async fn task_main(http_listener: TcpListener) -> anyhow::Result<()> {
|
||||||
scopeguard::defer! {
|
scopeguard::defer! {
|
||||||
println!("http has shut down");
|
info!("http has shut down");
|
||||||
}
|
}
|
||||||
|
|
||||||
let service = || RouterService::new(make_router().build()?);
|
let service = || RouterService::new(make_router().build()?);
|
||||||
|
|||||||
+15
-7
@@ -23,8 +23,10 @@ use anyhow::{bail, Context};
|
|||||||
use clap::{self, Arg};
|
use clap::{self, Arg};
|
||||||
use config::ProxyConfig;
|
use config::ProxyConfig;
|
||||||
use futures::FutureExt;
|
use futures::FutureExt;
|
||||||
|
use metrics::set_build_info_metric;
|
||||||
use std::{borrow::Cow, future::Future, net::SocketAddr};
|
use std::{borrow::Cow, future::Future, net::SocketAddr};
|
||||||
use tokio::{net::TcpListener, task::JoinError};
|
use tokio::{net::TcpListener, task::JoinError};
|
||||||
|
use tracing::info;
|
||||||
use utils::project_git_version;
|
use utils::project_git_version;
|
||||||
|
|
||||||
project_git_version!(GIT_VERSION);
|
project_git_version!(GIT_VERSION);
|
||||||
@@ -38,6 +40,11 @@ async fn flatten_err(
|
|||||||
|
|
||||||
#[tokio::main]
|
#[tokio::main]
|
||||||
async fn main() -> anyhow::Result<()> {
|
async fn main() -> anyhow::Result<()> {
|
||||||
|
tracing_subscriber::fmt()
|
||||||
|
.with_ansi(atty::is(atty::Stream::Stdout))
|
||||||
|
.with_target(false)
|
||||||
|
.init();
|
||||||
|
|
||||||
let arg_matches = clap::App::new("Neon proxy/router")
|
let arg_matches = clap::App::new("Neon proxy/router")
|
||||||
.version(GIT_VERSION)
|
.version(GIT_VERSION)
|
||||||
.arg(
|
.arg(
|
||||||
@@ -140,26 +147,27 @@ async fn main() -> anyhow::Result<()> {
|
|||||||
auth_backend,
|
auth_backend,
|
||||||
}));
|
}));
|
||||||
|
|
||||||
println!("Version: {GIT_VERSION}");
|
info!("Version: {GIT_VERSION}");
|
||||||
println!("Authentication backend: {}", config.auth_backend);
|
info!("Authentication backend: {}", config.auth_backend);
|
||||||
|
|
||||||
// Check that we can bind to address before further initialization
|
// Check that we can bind to address before further initialization
|
||||||
println!("Starting http on {}", http_address);
|
info!("Starting http on {http_address}");
|
||||||
let http_listener = TcpListener::bind(http_address).await?.into_std()?;
|
let http_listener = TcpListener::bind(http_address).await?.into_std()?;
|
||||||
|
|
||||||
println!("Starting mgmt on {}", mgmt_address);
|
info!("Starting mgmt on {mgmt_address}");
|
||||||
let mgmt_listener = TcpListener::bind(mgmt_address).await?.into_std()?;
|
let mgmt_listener = TcpListener::bind(mgmt_address).await?.into_std()?;
|
||||||
|
|
||||||
println!("Starting proxy on {}", proxy_address);
|
info!("Starting proxy on {proxy_address}");
|
||||||
let proxy_listener = TcpListener::bind(proxy_address).await?;
|
let proxy_listener = TcpListener::bind(proxy_address).await?;
|
||||||
|
|
||||||
let tasks = [
|
let tasks = [
|
||||||
tokio::spawn(http::server::thread_main(http_listener)),
|
tokio::spawn(http::server::task_main(http_listener)),
|
||||||
tokio::spawn(proxy::thread_main(config, proxy_listener)),
|
tokio::spawn(proxy::task_main(config, proxy_listener)),
|
||||||
tokio::task::spawn_blocking(move || mgmt::thread_main(mgmt_listener)),
|
tokio::task::spawn_blocking(move || mgmt::thread_main(mgmt_listener)),
|
||||||
]
|
]
|
||||||
.map(flatten_err);
|
.map(flatten_err);
|
||||||
|
|
||||||
|
set_build_info_metric(GIT_VERSION);
|
||||||
// This will block until all tasks have completed.
|
// This will block until all tasks have completed.
|
||||||
// Furthermore, the first one to fail will cancel the rest.
|
// Furthermore, the first one to fail will cancel the rest.
|
||||||
let _: Vec<()> = futures::future::try_join_all(tasks).await?;
|
let _: Vec<()> = futures::future::try_join_all(tasks).await?;
|
||||||
|
|||||||
+6
-5
@@ -5,6 +5,7 @@ use std::{
|
|||||||
net::{TcpListener, TcpStream},
|
net::{TcpListener, TcpStream},
|
||||||
thread,
|
thread,
|
||||||
};
|
};
|
||||||
|
use tracing::{error, info};
|
||||||
use utils::{
|
use utils::{
|
||||||
postgres_backend::{self, AuthType, PostgresBackend},
|
postgres_backend::{self, AuthType, PostgresBackend},
|
||||||
pq_proto::{BeMessage, SINGLE_COL_ROWDESC},
|
pq_proto::{BeMessage, SINGLE_COL_ROWDESC},
|
||||||
@@ -19,7 +20,7 @@ use utils::{
|
|||||||
///
|
///
|
||||||
pub fn thread_main(listener: TcpListener) -> anyhow::Result<()> {
|
pub fn thread_main(listener: TcpListener) -> anyhow::Result<()> {
|
||||||
scopeguard::defer! {
|
scopeguard::defer! {
|
||||||
println!("mgmt has shut down");
|
info!("mgmt has shut down");
|
||||||
}
|
}
|
||||||
|
|
||||||
listener
|
listener
|
||||||
@@ -27,14 +28,14 @@ pub fn thread_main(listener: TcpListener) -> anyhow::Result<()> {
|
|||||||
.context("failed to set listener to blocking")?;
|
.context("failed to set listener to blocking")?;
|
||||||
loop {
|
loop {
|
||||||
let (socket, peer_addr) = listener.accept().context("failed to accept a new client")?;
|
let (socket, peer_addr) = listener.accept().context("failed to accept a new client")?;
|
||||||
println!("accepted connection from {}", peer_addr);
|
info!("accepted connection from {peer_addr}");
|
||||||
socket
|
socket
|
||||||
.set_nodelay(true)
|
.set_nodelay(true)
|
||||||
.context("failed to set client socket option")?;
|
.context("failed to set client socket option")?;
|
||||||
|
|
||||||
thread::spawn(move || {
|
thread::spawn(move || {
|
||||||
if let Err(err) = handle_connection(socket) {
|
if let Err(err) = handle_connection(socket) {
|
||||||
println!("error: {}", err);
|
error!("{err}");
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
@@ -102,14 +103,14 @@ impl postgres_backend::Handler for MgmtHandler {
|
|||||||
let res = try_process_query(pgb, query_string);
|
let res = try_process_query(pgb, query_string);
|
||||||
// intercept and log error message
|
// intercept and log error message
|
||||||
if res.is_err() {
|
if res.is_err() {
|
||||||
println!("Mgmt query failed: #{:?}", res);
|
error!("mgmt query failed: {res:?}");
|
||||||
}
|
}
|
||||||
res
|
res
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn try_process_query(pgb: &mut PostgresBackend, query_string: &str) -> anyhow::Result<()> {
|
fn try_process_query(pgb: &mut PostgresBackend, query_string: &str) -> anyhow::Result<()> {
|
||||||
println!("Got mgmt query [redacted]"); // Content contains password, don't print it
|
info!("got mgmt query [redacted]"); // Content contains password, don't print it
|
||||||
|
|
||||||
let resp: PsqlSessionResponse = serde_json::from_str(query_string)?;
|
let resp: PsqlSessionResponse = serde_json::from_str(query_string)?;
|
||||||
|
|
||||||
|
|||||||
+36
-17
@@ -8,6 +8,7 @@ use metrics::{register_int_counter, IntCounter};
|
|||||||
use once_cell::sync::Lazy;
|
use once_cell::sync::Lazy;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use tokio::io::{AsyncRead, AsyncWrite};
|
use tokio::io::{AsyncRead, AsyncWrite};
|
||||||
|
use tracing::{error, info, info_span, Instrument};
|
||||||
use utils::pq_proto::{BeMessage as Be, *};
|
use utils::pq_proto::{BeMessage as Be, *};
|
||||||
|
|
||||||
const ERR_INSECURE_CONNECTION: &str = "connection is insecure (try using `sslmode=require`)";
|
const ERR_INSECURE_CONNECTION: &str = "connection is insecure (try using `sslmode=require`)";
|
||||||
@@ -43,17 +44,17 @@ where
|
|||||||
F: std::future::Future<Output = anyhow::Result<R>>,
|
F: std::future::Future<Output = anyhow::Result<R>>,
|
||||||
{
|
{
|
||||||
future.await.map_err(|err| {
|
future.await.map_err(|err| {
|
||||||
println!("error: {}", err);
|
error!("{err}");
|
||||||
err
|
err
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
pub async fn thread_main(
|
pub async fn task_main(
|
||||||
config: &'static ProxyConfig,
|
config: &'static ProxyConfig,
|
||||||
listener: tokio::net::TcpListener,
|
listener: tokio::net::TcpListener,
|
||||||
) -> anyhow::Result<()> {
|
) -> anyhow::Result<()> {
|
||||||
scopeguard::defer! {
|
scopeguard::defer! {
|
||||||
println!("proxy has shut down");
|
info!("proxy has shut down");
|
||||||
}
|
}
|
||||||
|
|
||||||
// When set for the server socket, the keepalive setting
|
// When set for the server socket, the keepalive setting
|
||||||
@@ -63,22 +64,29 @@ pub async fn thread_main(
|
|||||||
let cancel_map = Arc::new(CancelMap::default());
|
let cancel_map = Arc::new(CancelMap::default());
|
||||||
loop {
|
loop {
|
||||||
let (socket, peer_addr) = listener.accept().await?;
|
let (socket, peer_addr) = listener.accept().await?;
|
||||||
println!("accepted connection from {}", peer_addr);
|
info!("accepted connection from {peer_addr}");
|
||||||
|
|
||||||
|
let session_id = uuid::Uuid::new_v4();
|
||||||
let cancel_map = Arc::clone(&cancel_map);
|
let cancel_map = Arc::clone(&cancel_map);
|
||||||
tokio::spawn(log_error(async move {
|
tokio::spawn(
|
||||||
socket
|
log_error(async move {
|
||||||
.set_nodelay(true)
|
info!("spawned a task for {peer_addr}");
|
||||||
.context("failed to set socket option")?;
|
|
||||||
|
|
||||||
handle_client(config, &cancel_map, socket).await
|
socket
|
||||||
}));
|
.set_nodelay(true)
|
||||||
|
.context("failed to set socket option")?;
|
||||||
|
|
||||||
|
handle_client(config, &cancel_map, session_id, socket).await
|
||||||
|
})
|
||||||
|
.instrument(info_span!("client", session = format_args!("{session_id}"))),
|
||||||
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn handle_client(
|
async fn handle_client(
|
||||||
config: &ProxyConfig,
|
config: &ProxyConfig,
|
||||||
cancel_map: &CancelMap,
|
cancel_map: &CancelMap,
|
||||||
|
session_id: uuid::Uuid,
|
||||||
stream: impl AsyncRead + AsyncWrite + Unpin + Send,
|
stream: impl AsyncRead + AsyncWrite + Unpin + Send,
|
||||||
) -> anyhow::Result<()> {
|
) -> anyhow::Result<()> {
|
||||||
// The `closed` counter will increase when this future is destroyed.
|
// The `closed` counter will increase when this future is destroyed.
|
||||||
@@ -88,7 +96,8 @@ async fn handle_client(
|
|||||||
}
|
}
|
||||||
|
|
||||||
let tls = config.tls_config.as_ref();
|
let tls = config.tls_config.as_ref();
|
||||||
let (mut stream, params) = match handshake(stream, tls, cancel_map).await? {
|
let do_handshake = handshake(stream, tls, cancel_map).instrument(info_span!("handshake"));
|
||||||
|
let (mut stream, params) = match do_handshake.await? {
|
||||||
Some(x) => x,
|
Some(x) => x,
|
||||||
None => return Ok(()), // it's a cancellation request
|
None => return Ok(()), // it's a cancellation request
|
||||||
};
|
};
|
||||||
@@ -106,7 +115,7 @@ async fn handle_client(
|
|||||||
async { result }.or_else(|e| stream.throw_error(e)).await?
|
async { result }.or_else(|e| stream.throw_error(e)).await?
|
||||||
};
|
};
|
||||||
|
|
||||||
let client = Client::new(stream, creds, ¶ms);
|
let client = Client::new(stream, creds, ¶ms, session_id);
|
||||||
cancel_map
|
cancel_map
|
||||||
.with_session(|session| client.connect_to_db(session))
|
.with_session(|session| client.connect_to_db(session))
|
||||||
.await
|
.await
|
||||||
@@ -127,7 +136,7 @@ async fn handshake<S: AsyncRead + AsyncWrite + Unpin>(
|
|||||||
let mut stream = PqStream::new(Stream::from_raw(stream));
|
let mut stream = PqStream::new(Stream::from_raw(stream));
|
||||||
loop {
|
loop {
|
||||||
let msg = stream.read_startup_packet().await?;
|
let msg = stream.read_startup_packet().await?;
|
||||||
println!("got message: {:?}", msg);
|
info!("received {msg:?}");
|
||||||
|
|
||||||
use FeStartupPacket::*;
|
use FeStartupPacket::*;
|
||||||
match msg {
|
match msg {
|
||||||
@@ -164,11 +173,13 @@ async fn handshake<S: AsyncRead + AsyncWrite + Unpin>(
|
|||||||
stream.throw_error_str(ERR_INSECURE_CONNECTION).await?;
|
stream.throw_error_str(ERR_INSECURE_CONNECTION).await?;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
info!(session_type = "normal", "successful handshake");
|
||||||
break Ok(Some((stream, params)));
|
break Ok(Some((stream, params)));
|
||||||
}
|
}
|
||||||
CancelRequest(cancel_key_data) => {
|
CancelRequest(cancel_key_data) => {
|
||||||
cancel_map.cancel_session(cancel_key_data).await?;
|
cancel_map.cancel_session(cancel_key_data).await?;
|
||||||
|
|
||||||
|
info!(session_type = "cancellation", "successful handshake");
|
||||||
break Ok(None);
|
break Ok(None);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -183,6 +194,8 @@ struct Client<'a, S> {
|
|||||||
creds: auth::BackendType<'a, auth::ClientCredentials<'a>>,
|
creds: auth::BackendType<'a, auth::ClientCredentials<'a>>,
|
||||||
/// KV-dictionary with PostgreSQL connection params.
|
/// KV-dictionary with PostgreSQL connection params.
|
||||||
params: &'a StartupMessageParams,
|
params: &'a StartupMessageParams,
|
||||||
|
/// Unique connection ID.
|
||||||
|
session_id: uuid::Uuid,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<'a, S> Client<'a, S> {
|
impl<'a, S> Client<'a, S> {
|
||||||
@@ -191,11 +204,13 @@ impl<'a, S> Client<'a, S> {
|
|||||||
stream: PqStream<S>,
|
stream: PqStream<S>,
|
||||||
creds: auth::BackendType<'a, auth::ClientCredentials<'a>>,
|
creds: auth::BackendType<'a, auth::ClientCredentials<'a>>,
|
||||||
params: &'a StartupMessageParams,
|
params: &'a StartupMessageParams,
|
||||||
|
session_id: uuid::Uuid,
|
||||||
) -> Self {
|
) -> Self {
|
||||||
Self {
|
Self {
|
||||||
stream,
|
stream,
|
||||||
creds,
|
creds,
|
||||||
params,
|
params,
|
||||||
|
session_id,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -207,17 +222,20 @@ impl<S: AsyncRead + AsyncWrite + Unpin + Send> Client<'_, S> {
|
|||||||
mut stream,
|
mut stream,
|
||||||
creds,
|
creds,
|
||||||
params,
|
params,
|
||||||
|
session_id,
|
||||||
} = self;
|
} = self;
|
||||||
|
|
||||||
let extra = auth::ConsoleReqExtra {
|
let extra = auth::ConsoleReqExtra {
|
||||||
// Currently it's OK to generate a new UUID **here**, but
|
session_id, // aka this connection's id
|
||||||
// it might be better to move this to `cancellation::Session`.
|
|
||||||
session_id: uuid::Uuid::new_v4(),
|
|
||||||
application_name: params.get("application_name"),
|
application_name: params.get("application_name"),
|
||||||
};
|
};
|
||||||
|
|
||||||
// Authenticate and connect to a compute node.
|
// Authenticate and connect to a compute node.
|
||||||
let auth = creds.authenticate(&extra, &mut stream).await;
|
let auth = creds
|
||||||
|
.authenticate(&extra, &mut stream)
|
||||||
|
.instrument(info_span!("auth"))
|
||||||
|
.await;
|
||||||
|
|
||||||
let node = async { auth }.or_else(|e| stream.throw_error(e)).await?;
|
let node = async { auth }.or_else(|e| stream.throw_error(e)).await?;
|
||||||
let reported_auth_ok = node.reported_auth_ok;
|
let reported_auth_ok = node.reported_auth_ok;
|
||||||
|
|
||||||
@@ -251,6 +269,7 @@ impl<S: AsyncRead + AsyncWrite + Unpin + Send> Client<'_, S> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Starting from here we only proxy the client's traffic.
|
// Starting from here we only proxy the client's traffic.
|
||||||
|
info!("performing the proxy pass...");
|
||||||
let mut db = MetricsStream::new(db.stream, inc_proxied);
|
let mut db = MetricsStream::new(db.stream, inc_proxied);
|
||||||
let mut client = MetricsStream::new(stream.into_inner(), inc_proxied);
|
let mut client = MetricsStream::new(stream.into_inner(), inc_proxied);
|
||||||
let _ = tokio::io::copy_bidirectional(&mut client, &mut db).await?;
|
let _ = tokio::io::copy_bidirectional(&mut client, &mut db).await?;
|
||||||
|
|||||||
@@ -5,6 +5,7 @@ filterwarnings =
|
|||||||
ignore:record_property is incompatible with junit_family:pytest.PytestWarning
|
ignore:record_property is incompatible with junit_family:pytest.PytestWarning
|
||||||
addopts =
|
addopts =
|
||||||
-m 'not remote_cluster'
|
-m 'not remote_cluster'
|
||||||
|
--ignore=test_runner/performance
|
||||||
markers =
|
markers =
|
||||||
remote_cluster
|
remote_cluster
|
||||||
testpaths =
|
testpaths =
|
||||||
|
|||||||
+1
-1
@@ -13,7 +13,7 @@
|
|||||||
# avoid running regular linting script that checks every feature.
|
# avoid running regular linting script that checks every feature.
|
||||||
if [[ "$OSTYPE" == "darwin"* ]]; then
|
if [[ "$OSTYPE" == "darwin"* ]]; then
|
||||||
# no extra features to test currently, add more here when needed
|
# no extra features to test currently, add more here when needed
|
||||||
cargo clippy --locked --all --all-targets -- -A unknown_lints -D warnings
|
cargo clippy --locked --all --all-targets --features testing -- -A unknown_lints -D warnings
|
||||||
else
|
else
|
||||||
# * `-A unknown_lints` – do not warn about unknown lint suppressions
|
# * `-A unknown_lints` – do not warn about unknown lint suppressions
|
||||||
# that people with newer toolchains might use
|
# that people with newer toolchains might use
|
||||||
|
|||||||
+4
-5
@@ -1,11 +1,10 @@
|
|||||||
[toolchain]
|
[toolchain]
|
||||||
# We try to stick to a toolchain version that is widely available on popular distributions, so that most people
|
# We try to stick to a toolchain version that is widely available on popular distributions, so that most people
|
||||||
# can use the toolchain that comes with their operating system. But if there's a feature we miss badly from a later
|
# can use the toolchain that comes with their operating system. But if there's a feature we miss badly from a later
|
||||||
# version, we can consider updating. As of this writing, 1.60 is available on Debian 'experimental' but not yet on
|
# version, we can consider updating.
|
||||||
# 'testing' or even 'unstable', which is a bit more cutting-edge than we'd like. Hopefully the 1.60 packages reach
|
# See https://tracker.debian.org/pkg/rustc for more details on Debian rustc package,
|
||||||
# 'testing' soon (and similarly for the other distributions).
|
# we use "unstable" version number as the highest version used in the project by default.
|
||||||
# See https://tracker.debian.org/pkg/rustc for more details on Debian rustc package.
|
channel = "1.61" # do update GitHub CI cache values for rust builds, when changing this value
|
||||||
channel = "1.60" # do update GitHub CI cache values for rust builds, when changing this value
|
|
||||||
profile = "default"
|
profile = "default"
|
||||||
# The default profile includes rustc, rust-std, cargo, rust-docs, rustfmt and clippy.
|
# The default profile includes rustc, rust-std, cargo, rust-docs, rustfmt and clippy.
|
||||||
# https://rust-lang.github.io/rustup/concepts/profiles.html
|
# https://rust-lang.github.io/rustup/concepts/profiles.html
|
||||||
|
|||||||
@@ -33,6 +33,7 @@ toml_edit = { version = "0.13", features = ["easy"] }
|
|||||||
thiserror = "1"
|
thiserror = "1"
|
||||||
parking_lot = "0.12.1"
|
parking_lot = "0.12.1"
|
||||||
|
|
||||||
|
safekeeper_api = { path = "../libs/safekeeper_api" }
|
||||||
postgres_ffi = { path = "../libs/postgres_ffi" }
|
postgres_ffi = { path = "../libs/postgres_ffi" }
|
||||||
metrics = { path = "../libs/metrics" }
|
metrics = { path = "../libs/metrics" }
|
||||||
utils = { path = "../libs/utils" }
|
utils = { path = "../libs/utils" }
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ use toml_edit::Document;
|
|||||||
use tracing::*;
|
use tracing::*;
|
||||||
use url::{ParseError, Url};
|
use url::{ParseError, Url};
|
||||||
|
|
||||||
|
use metrics::set_build_info_metric;
|
||||||
use safekeeper::broker;
|
use safekeeper::broker;
|
||||||
use safekeeper::control_file;
|
use safekeeper::control_file;
|
||||||
use safekeeper::defaults::{
|
use safekeeper::defaults::{
|
||||||
@@ -291,9 +292,8 @@ fn start_safekeeper(mut conf: SafeKeeperConf, given_id: Option<NodeId>, init: bo
|
|||||||
|
|
||||||
// Register metrics collector for active timelines. It's important to do this
|
// Register metrics collector for active timelines. It's important to do this
|
||||||
// after daemonizing, otherwise process collector will be upset.
|
// after daemonizing, otherwise process collector will be upset.
|
||||||
let registry = metrics::default_registry();
|
|
||||||
let timeline_collector = safekeeper::metrics::TimelineCollector::new();
|
let timeline_collector = safekeeper::metrics::TimelineCollector::new();
|
||||||
registry.register(Box::new(timeline_collector))?;
|
metrics::register_internal(Box::new(timeline_collector))?;
|
||||||
|
|
||||||
let signals = signals::install_shutdown_handlers()?;
|
let signals = signals::install_shutdown_handlers()?;
|
||||||
let mut threads = vec![];
|
let mut threads = vec![];
|
||||||
@@ -364,6 +364,7 @@ fn start_safekeeper(mut conf: SafeKeeperConf, given_id: Option<NodeId>, init: bo
|
|||||||
})?,
|
})?,
|
||||||
);
|
);
|
||||||
|
|
||||||
|
set_build_info_metric(GIT_VERSION);
|
||||||
// TODO: put more thoughts into handling of failed threads
|
// TODO: put more thoughts into handling of failed threads
|
||||||
// We probably should restart them.
|
// We probably should restart them.
|
||||||
|
|
||||||
|
|||||||
@@ -1,3 +1,4 @@
|
|||||||
pub mod models;
|
|
||||||
pub mod routes;
|
pub mod routes;
|
||||||
pub use routes::make_router;
|
pub use routes::make_router;
|
||||||
|
|
||||||
|
pub use safekeeper_api::models;
|
||||||
|
|||||||
@@ -27,14 +27,13 @@ mod timelines_global_map;
|
|||||||
pub use timelines_global_map::GlobalTimelines;
|
pub use timelines_global_map::GlobalTimelines;
|
||||||
|
|
||||||
pub mod defaults {
|
pub mod defaults {
|
||||||
use const_format::formatcp;
|
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
pub const DEFAULT_PG_LISTEN_PORT: u16 = 5454;
|
pub use safekeeper_api::{
|
||||||
pub const DEFAULT_PG_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_PG_LISTEN_PORT}");
|
DEFAULT_HTTP_LISTEN_ADDR, DEFAULT_HTTP_LISTEN_PORT, DEFAULT_PG_LISTEN_ADDR,
|
||||||
|
DEFAULT_PG_LISTEN_PORT,
|
||||||
|
};
|
||||||
|
|
||||||
pub const DEFAULT_HTTP_LISTEN_PORT: u16 = 7676;
|
|
||||||
pub const DEFAULT_HTTP_LISTEN_ADDR: &str = formatcp!("127.0.0.1:{DEFAULT_HTTP_LISTEN_PORT}");
|
|
||||||
pub const DEFAULT_RECALL_PERIOD: Duration = Duration::from_secs(10);
|
pub const DEFAULT_RECALL_PERIOD: Duration = Duration::from_secs(10);
|
||||||
pub const DEFAULT_WAL_BACKUP_RUNTIME_THREADS: usize = 8;
|
pub const DEFAULT_WAL_BACKUP_RUNTIME_THREADS: usize = 8;
|
||||||
}
|
}
|
||||||
|
|||||||
Executable
+11
@@ -0,0 +1,11 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
|
||||||
|
set -euox pipefail
|
||||||
|
|
||||||
|
# Runs all formatting tools to ensure the project is up to date
|
||||||
|
echo 'Reformatting Rust code'
|
||||||
|
cargo fmt
|
||||||
|
echo 'Reformatting Python code'
|
||||||
|
poetry run isort test_runner
|
||||||
|
poetry run flake8 test_runner
|
||||||
|
poetry run black test_runner
|
||||||
@@ -56,6 +56,14 @@ If you want to run all tests that have the string "bench" in their names:
|
|||||||
|
|
||||||
`./scripts/pytest -k bench`
|
`./scripts/pytest -k bench`
|
||||||
|
|
||||||
|
To run tests in parellel we utilize `pytest-xdist` plugin. By default everything runs single threaded. Number of workers can be specified with `-n` argument:
|
||||||
|
|
||||||
|
`./scripts/pytest -n4`
|
||||||
|
|
||||||
|
By default performance tests are excluded. To run them explicitly pass performance tests selection to the script:
|
||||||
|
|
||||||
|
`./scripts/pytest test_runner/performance`
|
||||||
|
|
||||||
Useful environment variables:
|
Useful environment variables:
|
||||||
|
|
||||||
`NEON_BIN`: The directory where neon binaries can be found.
|
`NEON_BIN`: The directory where neon binaries can be found.
|
||||||
|
|||||||
@@ -455,6 +455,9 @@ class RemoteStorageKind(enum.Enum):
|
|||||||
LOCAL_FS = "local_fs"
|
LOCAL_FS = "local_fs"
|
||||||
MOCK_S3 = "mock_s3"
|
MOCK_S3 = "mock_s3"
|
||||||
REAL_S3 = "real_s3"
|
REAL_S3 = "real_s3"
|
||||||
|
# Pass to tests that are generic to remote storage
|
||||||
|
# to ensure the test pass with or without the remote storage
|
||||||
|
NOOP = "noop"
|
||||||
|
|
||||||
|
|
||||||
def available_remote_storages() -> List[RemoteStorageKind]:
|
def available_remote_storages() -> List[RemoteStorageKind]:
|
||||||
@@ -583,7 +586,9 @@ class NeonEnvBuilder:
|
|||||||
test_name: str,
|
test_name: str,
|
||||||
force_enable: bool = True,
|
force_enable: bool = True,
|
||||||
):
|
):
|
||||||
if remote_storage_kind == RemoteStorageKind.LOCAL_FS:
|
if remote_storage_kind == RemoteStorageKind.NOOP:
|
||||||
|
return
|
||||||
|
elif remote_storage_kind == RemoteStorageKind.LOCAL_FS:
|
||||||
self.enable_local_fs_remote_storage(force_enable=force_enable)
|
self.enable_local_fs_remote_storage(force_enable=force_enable)
|
||||||
elif remote_storage_kind == RemoteStorageKind.MOCK_S3:
|
elif remote_storage_kind == RemoteStorageKind.MOCK_S3:
|
||||||
self.enable_mock_s3_remote_storage(bucket_name=test_name, force_enable=force_enable)
|
self.enable_mock_s3_remote_storage(bucket_name=test_name, force_enable=force_enable)
|
||||||
@@ -1131,6 +1136,19 @@ class NeonPageserverHttpClient(requests.Session):
|
|||||||
assert res_json is None
|
assert res_json is None
|
||||||
return res_json
|
return res_json
|
||||||
|
|
||||||
|
def timeline_get_lsn_by_timestamp(
|
||||||
|
self, tenant_id: TenantId, timeline_id: TimelineId, timestamp
|
||||||
|
):
|
||||||
|
log.info(
|
||||||
|
f"Requesting lsn by timestamp {timestamp}, tenant {tenant_id}, timeline {timeline_id}"
|
||||||
|
)
|
||||||
|
res = self.get(
|
||||||
|
f"http://localhost:{self.port}/v1/tenant/{tenant_id}/timeline/{timeline_id}/get_lsn_by_timestamp?timestamp={timestamp}",
|
||||||
|
)
|
||||||
|
self.verbose_error(res)
|
||||||
|
res_json = res.json()
|
||||||
|
return res_json
|
||||||
|
|
||||||
def timeline_checkpoint(self, tenant_id: TenantId, timeline_id: TimelineId):
|
def timeline_checkpoint(self, tenant_id: TenantId, timeline_id: TimelineId):
|
||||||
log.info(f"Requesting checkpoint: tenant {tenant_id}, timeline {timeline_id}")
|
log.info(f"Requesting checkpoint: tenant {tenant_id}, timeline {timeline_id}")
|
||||||
res = self.put(
|
res = self.put(
|
||||||
@@ -1182,6 +1200,7 @@ class AbstractNeonCli(abc.ABC):
|
|||||||
arguments: List[str],
|
arguments: List[str],
|
||||||
extra_env_vars: Optional[Dict[str, str]] = None,
|
extra_env_vars: Optional[Dict[str, str]] = None,
|
||||||
check_return_code=True,
|
check_return_code=True,
|
||||||
|
timeout=None,
|
||||||
) -> "subprocess.CompletedProcess[str]":
|
) -> "subprocess.CompletedProcess[str]":
|
||||||
"""
|
"""
|
||||||
Run the command with the specified arguments.
|
Run the command with the specified arguments.
|
||||||
@@ -1228,6 +1247,7 @@ class AbstractNeonCli(abc.ABC):
|
|||||||
universal_newlines=True,
|
universal_newlines=True,
|
||||||
stdout=subprocess.PIPE,
|
stdout=subprocess.PIPE,
|
||||||
stderr=subprocess.PIPE,
|
stderr=subprocess.PIPE,
|
||||||
|
timeout=timeout,
|
||||||
)
|
)
|
||||||
if not res.returncode:
|
if not res.returncode:
|
||||||
log.info(f"Run success: {res.stdout}")
|
log.info(f"Run success: {res.stdout}")
|
||||||
@@ -1601,6 +1621,14 @@ class WalCraft(AbstractNeonCli):
|
|||||||
res.check_returncode()
|
res.check_returncode()
|
||||||
|
|
||||||
|
|
||||||
|
class ComputeCtl(AbstractNeonCli):
|
||||||
|
"""
|
||||||
|
A typed wrapper around the `compute_ctl` CLI tool.
|
||||||
|
"""
|
||||||
|
|
||||||
|
COMMAND = "compute_ctl"
|
||||||
|
|
||||||
|
|
||||||
class NeonPageserver(PgProtocol):
|
class NeonPageserver(PgProtocol):
|
||||||
"""
|
"""
|
||||||
An object representing a running pageserver.
|
An object representing a running pageserver.
|
||||||
@@ -1933,6 +1961,11 @@ class NeonProxy(PgProtocol):
|
|||||||
def _wait_until_ready(self):
|
def _wait_until_ready(self):
|
||||||
requests.get(f"http://{self.host}:{self.http_port}/v1/status")
|
requests.get(f"http://{self.host}:{self.http_port}/v1/status")
|
||||||
|
|
||||||
|
def get_metrics(self) -> str:
|
||||||
|
request_result = requests.get(f"http://{self.host}:{self.http_port}/metrics")
|
||||||
|
request_result.raise_for_status()
|
||||||
|
return request_result.text
|
||||||
|
|
||||||
def __enter__(self):
|
def __enter__(self):
|
||||||
return self
|
return self
|
||||||
|
|
||||||
|
|||||||
@@ -3,12 +3,15 @@ import statistics
|
|||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
import timeit
|
import timeit
|
||||||
|
from contextlib import closing
|
||||||
from typing import List
|
from typing import List
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
from fixtures.benchmark_fixture import MetricReport
|
from fixtures.benchmark_fixture import MetricReport
|
||||||
from fixtures.compare_fixtures import NeonCompare
|
from fixtures.compare_fixtures import NeonCompare
|
||||||
from fixtures.log_helper import log
|
from fixtures.log_helper import log
|
||||||
|
from fixtures.neon_fixtures import wait_for_last_record_lsn
|
||||||
|
from fixtures.types import Lsn
|
||||||
|
|
||||||
|
|
||||||
def _record_branch_creation_durations(neon_compare: NeonCompare, durs: List[float]):
|
def _record_branch_creation_durations(neon_compare: NeonCompare, durs: List[float]):
|
||||||
@@ -107,3 +110,43 @@ def test_branch_creation_many(neon_compare: NeonCompare, n_branches: int):
|
|||||||
branch_creation_durations.append(dur)
|
branch_creation_durations.append(dur)
|
||||||
|
|
||||||
_record_branch_creation_durations(neon_compare, branch_creation_durations)
|
_record_branch_creation_durations(neon_compare, branch_creation_durations)
|
||||||
|
|
||||||
|
|
||||||
|
# Test measures the branch creation time when branching from a timeline with a lot of relations.
|
||||||
|
#
|
||||||
|
# This test measures the latency of branch creation under two scenarios
|
||||||
|
# 1. The ancestor branch is not under any workloads
|
||||||
|
# 2. The ancestor branch is under a workload (busy)
|
||||||
|
#
|
||||||
|
# To simulate the workload, the test runs a concurrent insertion on the ancestor branch right before branching.
|
||||||
|
def test_branch_creation_many_relations(neon_compare: NeonCompare):
|
||||||
|
env = neon_compare.env
|
||||||
|
|
||||||
|
timeline_id = env.neon_cli.create_branch("root")
|
||||||
|
|
||||||
|
pg = env.postgres.create_start("root")
|
||||||
|
with closing(pg.connect()) as conn:
|
||||||
|
with conn.cursor() as cur:
|
||||||
|
for i in range(10000):
|
||||||
|
cur.execute(f"CREATE TABLE t{i} as SELECT g FROM generate_series(1, 1000) g")
|
||||||
|
|
||||||
|
# Wait for the pageserver to finish processing all the pending WALs,
|
||||||
|
# as we don't want the LSN wait time to be included during the branch creation
|
||||||
|
flush_lsn = Lsn(pg.safe_psql("SELECT pg_current_wal_flush_lsn()")[0][0])
|
||||||
|
wait_for_last_record_lsn(
|
||||||
|
env.pageserver.http_client(), env.initial_tenant, timeline_id, flush_lsn
|
||||||
|
)
|
||||||
|
|
||||||
|
with neon_compare.record_duration("create_branch_time_not_busy_root"):
|
||||||
|
env.neon_cli.create_branch("child_not_busy", "root")
|
||||||
|
|
||||||
|
# run a concurrent insertion to make the ancestor "busy" during the branch creation
|
||||||
|
thread = threading.Thread(
|
||||||
|
target=pg.safe_psql, args=("INSERT INTO t0 VALUES (generate_series(1, 100000))",)
|
||||||
|
)
|
||||||
|
thread.start()
|
||||||
|
|
||||||
|
with neon_compare.record_duration("create_branch_time_busy_root"):
|
||||||
|
env.neon_cli.create_branch("child_busy", "root")
|
||||||
|
|
||||||
|
thread.join()
|
||||||
|
|||||||
@@ -0,0 +1,94 @@
|
|||||||
|
import timeit
|
||||||
|
from pathlib import Path
|
||||||
|
from typing import List
|
||||||
|
|
||||||
|
from fixtures.benchmark_fixture import PgBenchRunResult
|
||||||
|
from fixtures.compare_fixtures import NeonCompare
|
||||||
|
from performance.test_perf_pgbench import utc_now_timestamp
|
||||||
|
|
||||||
|
# -----------------------------------------------------------------------
|
||||||
|
# Start of `test_compare_child_and_root_*` tests
|
||||||
|
# -----------------------------------------------------------------------
|
||||||
|
|
||||||
|
# `test_compare_child_and_root_*` tests compare the performance of a branch and its child branch(s).
|
||||||
|
# A common pattern in those tests is initializing a root branch then creating a child branch(s) from the root.
|
||||||
|
# Each test then runs a similar workload for both child branch and root branch. Each measures and reports
|
||||||
|
# some latencies/metrics during the workload for performance comparison between a branch and its ancestor.
|
||||||
|
|
||||||
|
|
||||||
|
def test_compare_child_and_root_pgbench_perf(neon_compare: NeonCompare):
|
||||||
|
env = neon_compare.env
|
||||||
|
pg_bin = neon_compare.pg_bin
|
||||||
|
|
||||||
|
def run_pgbench_on_branch(branch: str, cmd: List[str]):
|
||||||
|
run_start_timestamp = utc_now_timestamp()
|
||||||
|
t0 = timeit.default_timer()
|
||||||
|
out = pg_bin.run_capture(
|
||||||
|
cmd,
|
||||||
|
)
|
||||||
|
run_duration = timeit.default_timer() - t0
|
||||||
|
run_end_timestamp = utc_now_timestamp()
|
||||||
|
|
||||||
|
stdout = Path(f"{out}.stdout").read_text()
|
||||||
|
|
||||||
|
res = PgBenchRunResult.parse_from_stdout(
|
||||||
|
stdout=stdout,
|
||||||
|
run_duration=run_duration,
|
||||||
|
run_start_timestamp=run_start_timestamp,
|
||||||
|
run_end_timestamp=run_end_timestamp,
|
||||||
|
)
|
||||||
|
neon_compare.zenbenchmark.record_pg_bench_result(branch, res)
|
||||||
|
|
||||||
|
env.neon_cli.create_branch("root")
|
||||||
|
pg_root = env.postgres.create_start("root")
|
||||||
|
pg_bin.run_capture(["pgbench", "-i", pg_root.connstr(), "-s10"])
|
||||||
|
|
||||||
|
env.neon_cli.create_branch("child", "root")
|
||||||
|
pg_child = env.postgres.create_start("child")
|
||||||
|
|
||||||
|
run_pgbench_on_branch("root", ["pgbench", "-c10", "-T10", pg_root.connstr()])
|
||||||
|
run_pgbench_on_branch("child", ["pgbench", "-c10", "-T10", pg_child.connstr()])
|
||||||
|
|
||||||
|
|
||||||
|
def test_compare_child_and_root_write_perf(neon_compare: NeonCompare):
|
||||||
|
env = neon_compare.env
|
||||||
|
env.neon_cli.create_branch("root")
|
||||||
|
pg_root = env.postgres.create_start("root")
|
||||||
|
|
||||||
|
pg_root.safe_psql(
|
||||||
|
"CREATE TABLE foo(key serial primary key, t text default 'foooooooooooooooooooooooooooooooooooooooooooooooooooo')",
|
||||||
|
)
|
||||||
|
|
||||||
|
env.neon_cli.create_branch("child", "root")
|
||||||
|
pg_child = env.postgres.create_start("child")
|
||||||
|
|
||||||
|
with neon_compare.record_duration("root_run_duration"):
|
||||||
|
pg_root.safe_psql("INSERT INTO foo SELECT FROM generate_series(1,1000000)")
|
||||||
|
with neon_compare.record_duration("child_run_duration"):
|
||||||
|
pg_child.safe_psql("INSERT INTO foo SELECT FROM generate_series(1,1000000)")
|
||||||
|
|
||||||
|
|
||||||
|
def test_compare_child_and_root_read_perf(neon_compare: NeonCompare):
|
||||||
|
env = neon_compare.env
|
||||||
|
env.neon_cli.create_branch("root")
|
||||||
|
pg_root = env.postgres.create_start("root")
|
||||||
|
|
||||||
|
pg_root.safe_psql_many(
|
||||||
|
[
|
||||||
|
"CREATE TABLE foo(key serial primary key, t text default 'foooooooooooooooooooooooooooooooooooooooooooooooooooo')",
|
||||||
|
"INSERT INTO foo SELECT FROM generate_series(1,1000000)",
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
|
env.neon_cli.create_branch("child", "root")
|
||||||
|
pg_child = env.postgres.create_start("child")
|
||||||
|
|
||||||
|
with neon_compare.record_duration("root_run_duration"):
|
||||||
|
pg_root.safe_psql("SELECT count(*) from foo")
|
||||||
|
with neon_compare.record_duration("child_run_duration"):
|
||||||
|
pg_child.safe_psql("SELECT count(*) from foo")
|
||||||
|
|
||||||
|
|
||||||
|
# -----------------------------------------------------------------------
|
||||||
|
# End of `test_compare_child_and_root_*` tests
|
||||||
|
# -----------------------------------------------------------------------
|
||||||
@@ -6,6 +6,8 @@ from typing import List
|
|||||||
import pytest
|
import pytest
|
||||||
from fixtures.log_helper import log
|
from fixtures.log_helper import log
|
||||||
from fixtures.neon_fixtures import NeonEnv, PgBin, Postgres
|
from fixtures.neon_fixtures import NeonEnv, PgBin, Postgres
|
||||||
|
from fixtures.types import Lsn
|
||||||
|
from fixtures.utils import query_scalar
|
||||||
from performance.test_perf_pgbench import get_scales_matrix
|
from performance.test_perf_pgbench import get_scales_matrix
|
||||||
|
|
||||||
|
|
||||||
@@ -88,3 +90,39 @@ def test_branching_with_pgbench(
|
|||||||
for pg in pgs:
|
for pg in pgs:
|
||||||
res = pg.safe_psql("SELECT count(*) from pgbench_accounts")
|
res = pg.safe_psql("SELECT count(*) from pgbench_accounts")
|
||||||
assert res[0] == (100000 * scale,)
|
assert res[0] == (100000 * scale,)
|
||||||
|
|
||||||
|
|
||||||
|
# Test branching from an "unnormalized" LSN.
|
||||||
|
#
|
||||||
|
# Context:
|
||||||
|
# When doing basebackup for a newly created branch, pageserver generates
|
||||||
|
# 'pg_control' file to bootstrap WAL segment by specifying the redo position
|
||||||
|
# a "normalized" LSN based on the timeline's starting LSN:
|
||||||
|
#
|
||||||
|
# checkpoint.redo = normalize_lsn(self.lsn, pg_constants::WAL_SEGMENT_SIZE).0;
|
||||||
|
#
|
||||||
|
# This test checks if the pageserver is able to handle a "unnormalized" starting LSN.
|
||||||
|
#
|
||||||
|
# Related: see discussion in https://github.com/neondatabase/neon/pull/2143#issuecomment-1209092186
|
||||||
|
def test_branching_unnormalized_start_lsn(neon_simple_env: NeonEnv, pg_bin: PgBin):
|
||||||
|
XLOG_BLCKSZ = 8192
|
||||||
|
|
||||||
|
env = neon_simple_env
|
||||||
|
|
||||||
|
env.neon_cli.create_branch("b0")
|
||||||
|
pg0 = env.postgres.create_start("b0")
|
||||||
|
|
||||||
|
pg_bin.run_capture(["pgbench", "-i", pg0.connstr()])
|
||||||
|
|
||||||
|
with pg0.cursor() as cur:
|
||||||
|
curr_lsn = Lsn(query_scalar(cur, "SELECT pg_current_wal_flush_lsn()"))
|
||||||
|
|
||||||
|
# Specify the `start_lsn` as a number that is divided by `XLOG_BLCKSZ`
|
||||||
|
# and is smaller than `curr_lsn`.
|
||||||
|
start_lsn = Lsn((int(curr_lsn) - XLOG_BLCKSZ) // XLOG_BLCKSZ * XLOG_BLCKSZ)
|
||||||
|
|
||||||
|
log.info(f"Branching b1 from b0 starting at lsn {start_lsn}...")
|
||||||
|
env.neon_cli.create_branch("b1", "b0", ancestor_start_lsn=start_lsn)
|
||||||
|
pg1 = env.postgres.create_start("b1")
|
||||||
|
|
||||||
|
pg_bin.run_capture(["pgbench", "-i", pg1.connstr()])
|
||||||
|
|||||||
@@ -0,0 +1,19 @@
|
|||||||
|
from fixtures.metrics import parse_metrics
|
||||||
|
from fixtures.neon_fixtures import NeonEnvBuilder, NeonProxy
|
||||||
|
|
||||||
|
|
||||||
|
def test_build_info_metric(neon_env_builder: NeonEnvBuilder, link_proxy: NeonProxy):
|
||||||
|
neon_env_builder.num_safekeepers = 1
|
||||||
|
env = neon_env_builder.init_start()
|
||||||
|
|
||||||
|
parsed_metrics = {}
|
||||||
|
|
||||||
|
parsed_metrics["pageserver"] = parse_metrics(env.pageserver.http_client().get_metrics())
|
||||||
|
parsed_metrics["safekeeper"] = parse_metrics(env.safekeepers[0].http_client().get_metrics_str())
|
||||||
|
parsed_metrics["proxy"] = parse_metrics(link_proxy.get_metrics())
|
||||||
|
|
||||||
|
for component, metrics in parsed_metrics.items():
|
||||||
|
sample = metrics.query_one("libmetrics_build_info")
|
||||||
|
|
||||||
|
assert "revision" in sample.labels
|
||||||
|
assert len(sample.labels["revision"]) > 0
|
||||||
@@ -0,0 +1,203 @@
|
|||||||
|
import os
|
||||||
|
from subprocess import TimeoutExpired
|
||||||
|
|
||||||
|
from fixtures.log_helper import log
|
||||||
|
from fixtures.neon_fixtures import ComputeCtl, NeonEnvBuilder, PgBin
|
||||||
|
|
||||||
|
|
||||||
|
# Test that compute_ctl works and prints "--sync-safekeepers" logs.
|
||||||
|
def test_sync_safekeepers_logs(neon_env_builder: NeonEnvBuilder, pg_bin: PgBin):
|
||||||
|
neon_env_builder.num_safekeepers = 3
|
||||||
|
env = neon_env_builder.init_start()
|
||||||
|
ctl = ComputeCtl(env)
|
||||||
|
|
||||||
|
env.neon_cli.create_branch("test_compute_ctl", "main")
|
||||||
|
pg = env.postgres.create_start("test_compute_ctl")
|
||||||
|
pg.safe_psql("CREATE TABLE t(key int primary key, value text)")
|
||||||
|
|
||||||
|
with open(pg.config_file_path(), "r") as f:
|
||||||
|
cfg_lines = f.readlines()
|
||||||
|
cfg_map = {}
|
||||||
|
for line in cfg_lines:
|
||||||
|
if "=" in line:
|
||||||
|
k, v = line.split("=")
|
||||||
|
cfg_map[k] = v.strip("\n '\"")
|
||||||
|
log.info(f"postgres config: {cfg_map}")
|
||||||
|
pgdata = pg.pg_data_dir_path()
|
||||||
|
pg_bin_path = os.path.join(pg_bin.pg_bin_path, "postgres")
|
||||||
|
|
||||||
|
pg.stop_and_destroy()
|
||||||
|
|
||||||
|
spec = (
|
||||||
|
"""
|
||||||
|
{
|
||||||
|
"format_version": 1.0,
|
||||||
|
|
||||||
|
"timestamp": "2021-05-23T18:25:43.511Z",
|
||||||
|
"operation_uuid": "0f657b36-4b0f-4a2d-9c2e-1dcd615e7d8b",
|
||||||
|
|
||||||
|
"cluster": {
|
||||||
|
"cluster_id": "test-cluster-42",
|
||||||
|
"name": "Neon Test",
|
||||||
|
"state": "restarted",
|
||||||
|
"roles": [
|
||||||
|
],
|
||||||
|
"databases": [
|
||||||
|
],
|
||||||
|
"settings": [
|
||||||
|
{
|
||||||
|
"name": "fsync",
|
||||||
|
"value": "off",
|
||||||
|
"vartype": "bool"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "wal_level",
|
||||||
|
"value": "replica",
|
||||||
|
"vartype": "enum"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "hot_standby",
|
||||||
|
"value": "on",
|
||||||
|
"vartype": "bool"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "neon.safekeepers",
|
||||||
|
"value": """
|
||||||
|
+ f'"{cfg_map["neon.safekeepers"]}"'
|
||||||
|
+ """,
|
||||||
|
"vartype": "string"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "wal_log_hints",
|
||||||
|
"value": "on",
|
||||||
|
"vartype": "bool"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "log_connections",
|
||||||
|
"value": "on",
|
||||||
|
"vartype": "bool"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "shared_buffers",
|
||||||
|
"value": "32768",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "port",
|
||||||
|
"value": """
|
||||||
|
+ f'"{cfg_map["port"]}"'
|
||||||
|
+ """,
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "max_connections",
|
||||||
|
"value": "100",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "max_wal_senders",
|
||||||
|
"value": "10",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "listen_addresses",
|
||||||
|
"value": "0.0.0.0",
|
||||||
|
"vartype": "string"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "wal_sender_timeout",
|
||||||
|
"value": "0",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "password_encryption",
|
||||||
|
"value": "md5",
|
||||||
|
"vartype": "enum"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "maintenance_work_mem",
|
||||||
|
"value": "65536",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "max_parallel_workers",
|
||||||
|
"value": "8",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "max_worker_processes",
|
||||||
|
"value": "8",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "neon.tenant_id",
|
||||||
|
"value": """
|
||||||
|
+ f'"{cfg_map["neon.tenant_id"]}"'
|
||||||
|
+ """,
|
||||||
|
"vartype": "string"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "max_replication_slots",
|
||||||
|
"value": "10",
|
||||||
|
"vartype": "integer"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "neon.timeline_id",
|
||||||
|
"value": """
|
||||||
|
+ f'"{cfg_map["neon.timeline_id"]}"'
|
||||||
|
+ """,
|
||||||
|
"vartype": "string"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "shared_preload_libraries",
|
||||||
|
"value": "neon",
|
||||||
|
"vartype": "string"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "synchronous_standby_names",
|
||||||
|
"value": "walproposer",
|
||||||
|
"vartype": "string"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "neon.pageserver_connstring",
|
||||||
|
"value": """
|
||||||
|
+ f'"{cfg_map["neon.pageserver_connstring"]}"'
|
||||||
|
+ """,
|
||||||
|
"vartype": "string"
|
||||||
|
}
|
||||||
|
]
|
||||||
|
},
|
||||||
|
"delta_operations": [
|
||||||
|
]
|
||||||
|
}
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
|
||||||
|
ps_connstr = cfg_map["neon.pageserver_connstring"]
|
||||||
|
log.info(f"ps_connstr: {ps_connstr}, pgdata: {pgdata}")
|
||||||
|
|
||||||
|
# run compute_ctl and wait for 10s
|
||||||
|
try:
|
||||||
|
ctl.raw_cli(
|
||||||
|
["--connstr", ps_connstr, "--pgdata", pgdata, "--spec", spec, "--pgbin", pg_bin_path],
|
||||||
|
timeout=10,
|
||||||
|
)
|
||||||
|
except TimeoutExpired as exc:
|
||||||
|
ctl_logs = exc.stderr.decode("utf-8")
|
||||||
|
log.info("compute_ctl output:\n" + ctl_logs)
|
||||||
|
|
||||||
|
start = "starting safekeepers syncing"
|
||||||
|
end = "safekeepers synced at LSN"
|
||||||
|
start_pos = ctl_logs.index(start)
|
||||||
|
assert start_pos != -1
|
||||||
|
end_pos = ctl_logs.index(end, start_pos)
|
||||||
|
assert end_pos != -1
|
||||||
|
sync_safekeepers_logs = ctl_logs[start_pos : end_pos + len(end)]
|
||||||
|
log.info("sync_safekeepers_logs:\n" + sync_safekeepers_logs)
|
||||||
|
|
||||||
|
# assert that --sync-safekeepers logs are present in the output
|
||||||
|
assert "connecting with node" in sync_safekeepers_logs
|
||||||
|
assert "connected with node" in sync_safekeepers_logs
|
||||||
|
assert "proposer connected to quorum (2)" in sync_safekeepers_logs
|
||||||
|
assert "got votes from majority (2)" in sync_safekeepers_logs
|
||||||
|
assert "sending elected msg to node" in sync_safekeepers_logs
|
||||||
@@ -15,7 +15,6 @@ def test_lsn_mapping(neon_env_builder: NeonEnvBuilder):
|
|||||||
pgmain = env.postgres.create_start("test_lsn_mapping")
|
pgmain = env.postgres.create_start("test_lsn_mapping")
|
||||||
log.info("postgres is running on 'test_lsn_mapping' branch")
|
log.info("postgres is running on 'test_lsn_mapping' branch")
|
||||||
|
|
||||||
ps_cur = env.pageserver.connect().cursor()
|
|
||||||
cur = pgmain.connect().cursor()
|
cur = pgmain.connect().cursor()
|
||||||
# Create table, and insert rows, each in a separate transaction
|
# Create table, and insert rows, each in a separate transaction
|
||||||
# Disable synchronous_commit to make this initialization go faster.
|
# Disable synchronous_commit to make this initialization go faster.
|
||||||
@@ -38,37 +37,33 @@ def test_lsn_mapping(neon_env_builder: NeonEnvBuilder):
|
|||||||
# Wait until WAL is received by pageserver
|
# Wait until WAL is received by pageserver
|
||||||
wait_for_last_flush_lsn(env, pgmain, env.initial_tenant, new_timeline_id)
|
wait_for_last_flush_lsn(env, pgmain, env.initial_tenant, new_timeline_id)
|
||||||
|
|
||||||
# Check edge cases: timestamp in the future
|
with env.pageserver.http_client() as client:
|
||||||
probe_timestamp = tbl[-1][1] + timedelta(hours=1)
|
# Check edge cases: timestamp in the future
|
||||||
result = query_scalar(
|
probe_timestamp = tbl[-1][1] + timedelta(hours=1)
|
||||||
ps_cur,
|
result = client.timeline_get_lsn_by_timestamp(
|
||||||
f"get_lsn_by_timestamp {env.initial_tenant} {new_timeline_id} '{probe_timestamp.isoformat()}Z'",
|
env.initial_tenant, new_timeline_id, f"{probe_timestamp.isoformat()}Z"
|
||||||
)
|
|
||||||
assert result == "future"
|
|
||||||
|
|
||||||
# timestamp too the far history
|
|
||||||
probe_timestamp = tbl[0][1] - timedelta(hours=10)
|
|
||||||
result = query_scalar(
|
|
||||||
ps_cur,
|
|
||||||
f"get_lsn_by_timestamp {env.initial_tenant} {new_timeline_id} '{probe_timestamp.isoformat()}Z'",
|
|
||||||
)
|
|
||||||
assert result == "past"
|
|
||||||
|
|
||||||
# Probe a bunch of timestamps in the valid range
|
|
||||||
for i in range(1, len(tbl), 100):
|
|
||||||
probe_timestamp = tbl[i][1]
|
|
||||||
|
|
||||||
# Call get_lsn_by_timestamp to get the LSN
|
|
||||||
lsn = query_scalar(
|
|
||||||
ps_cur,
|
|
||||||
f"get_lsn_by_timestamp {env.initial_tenant} {new_timeline_id} '{probe_timestamp.isoformat()}Z'",
|
|
||||||
)
|
)
|
||||||
|
assert result == "future"
|
||||||
|
|
||||||
# Launch a new read-only node at that LSN, and check that only the rows
|
# timestamp too the far history
|
||||||
# that were supposed to be committed at that point in time are visible.
|
probe_timestamp = tbl[0][1] - timedelta(hours=10)
|
||||||
pg_here = env.postgres.create_start(
|
result = client.timeline_get_lsn_by_timestamp(
|
||||||
branch_name="test_lsn_mapping", node_name="test_lsn_mapping_read", lsn=lsn
|
env.initial_tenant, new_timeline_id, f"{probe_timestamp.isoformat()}Z"
|
||||||
)
|
)
|
||||||
assert pg_here.safe_psql("SELECT max(x) FROM foo")[0][0] == i
|
assert result == "past"
|
||||||
|
|
||||||
pg_here.stop_and_destroy()
|
# Probe a bunch of timestamps in the valid range
|
||||||
|
for i in range(1, len(tbl), 100):
|
||||||
|
probe_timestamp = tbl[i][1]
|
||||||
|
lsn = client.timeline_get_lsn_by_timestamp(
|
||||||
|
env.initial_tenant, new_timeline_id, f"{probe_timestamp.isoformat()}Z"
|
||||||
|
)
|
||||||
|
# Call get_lsn_by_timestamp to get the LSN
|
||||||
|
# Launch a new read-only node at that LSN, and check that only the rows
|
||||||
|
# that were supposed to be committed at that point in time are visible.
|
||||||
|
pg_here = env.postgres.create_start(
|
||||||
|
branch_name="test_lsn_mapping", node_name="test_lsn_mapping_read", lsn=lsn
|
||||||
|
)
|
||||||
|
assert pg_here.safe_psql("SELECT max(x) FROM foo")[0][0] == i
|
||||||
|
|
||||||
|
pg_here.stop_and_destroy()
|
||||||
|
|||||||
@@ -54,7 +54,7 @@ tenant_config={checkpoint_distance = 10000, compaction_target_size = 1048576}"""
|
|||||||
for i in {
|
for i in {
|
||||||
"checkpoint_distance": 10000,
|
"checkpoint_distance": 10000,
|
||||||
"compaction_target_size": 1048576,
|
"compaction_target_size": 1048576,
|
||||||
"compaction_period": 1,
|
"compaction_period": 20,
|
||||||
"compaction_threshold": 10,
|
"compaction_threshold": 10,
|
||||||
"gc_horizon": 67108864,
|
"gc_horizon": 67108864,
|
||||||
"gc_period": 100,
|
"gc_period": 100,
|
||||||
@@ -74,7 +74,7 @@ tenant_config={checkpoint_distance = 10000, compaction_target_size = 1048576}"""
|
|||||||
for i in {
|
for i in {
|
||||||
"checkpoint_distance": 20000,
|
"checkpoint_distance": 20000,
|
||||||
"compaction_target_size": 1048576,
|
"compaction_target_size": 1048576,
|
||||||
"compaction_period": 1,
|
"compaction_period": 20,
|
||||||
"compaction_threshold": 10,
|
"compaction_threshold": 10,
|
||||||
"gc_horizon": 67108864,
|
"gc_horizon": 67108864,
|
||||||
"gc_period": 30,
|
"gc_period": 30,
|
||||||
@@ -102,7 +102,7 @@ tenant_config={checkpoint_distance = 10000, compaction_target_size = 1048576}"""
|
|||||||
for i in {
|
for i in {
|
||||||
"checkpoint_distance": 15000,
|
"checkpoint_distance": 15000,
|
||||||
"compaction_target_size": 1048576,
|
"compaction_target_size": 1048576,
|
||||||
"compaction_period": 1,
|
"compaction_period": 20,
|
||||||
"compaction_threshold": 10,
|
"compaction_threshold": 10,
|
||||||
"gc_horizon": 67108864,
|
"gc_horizon": 67108864,
|
||||||
"gc_period": 80,
|
"gc_period": 80,
|
||||||
@@ -125,7 +125,7 @@ tenant_config={checkpoint_distance = 10000, compaction_target_size = 1048576}"""
|
|||||||
for i in {
|
for i in {
|
||||||
"checkpoint_distance": 15000,
|
"checkpoint_distance": 15000,
|
||||||
"compaction_target_size": 1048576,
|
"compaction_target_size": 1048576,
|
||||||
"compaction_period": 1,
|
"compaction_period": 20,
|
||||||
"compaction_threshold": 10,
|
"compaction_threshold": 10,
|
||||||
"gc_horizon": 67108864,
|
"gc_horizon": 67108864,
|
||||||
"gc_period": 80,
|
"gc_period": 80,
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import os
|
import os
|
||||||
|
import shutil
|
||||||
from contextlib import closing
|
from contextlib import closing
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
@@ -7,8 +8,13 @@ from typing import List
|
|||||||
import pytest
|
import pytest
|
||||||
from fixtures.log_helper import log
|
from fixtures.log_helper import log
|
||||||
from fixtures.metrics import PAGESERVER_PER_TENANT_METRICS, parse_metrics
|
from fixtures.metrics import PAGESERVER_PER_TENANT_METRICS, parse_metrics
|
||||||
from fixtures.neon_fixtures import NeonEnv, NeonEnvBuilder
|
from fixtures.neon_fixtures import (
|
||||||
from fixtures.types import Lsn, TenantId
|
NeonEnv,
|
||||||
|
NeonEnvBuilder,
|
||||||
|
RemoteStorageKind,
|
||||||
|
available_remote_storages,
|
||||||
|
)
|
||||||
|
from fixtures.types import Lsn, TenantId, TimelineId
|
||||||
from prometheus_client.samples import Sample
|
from prometheus_client.samples import Sample
|
||||||
|
|
||||||
|
|
||||||
@@ -201,3 +207,63 @@ def test_pageserver_metrics_removed_after_detach(neon_env_builder: NeonEnvBuilde
|
|||||||
|
|
||||||
post_detach_samples = set([x.name for x in get_ps_metric_samples_for_tenant(tenant)])
|
post_detach_samples = set([x.name for x in get_ps_metric_samples_for_tenant(tenant)])
|
||||||
assert post_detach_samples == set()
|
assert post_detach_samples == set()
|
||||||
|
|
||||||
|
|
||||||
|
# Check that empty tenants work with or without the remote storage
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
"remote_storage_kind", available_remote_storages() + [RemoteStorageKind.NOOP]
|
||||||
|
)
|
||||||
|
def test_pageserver_with_empty_tenants(
|
||||||
|
neon_env_builder: NeonEnvBuilder, remote_storage_kind: RemoteStorageKind
|
||||||
|
):
|
||||||
|
neon_env_builder.enable_remote_storage(
|
||||||
|
remote_storage_kind=remote_storage_kind,
|
||||||
|
test_name="test_pageserver_with_empty_tenants",
|
||||||
|
)
|
||||||
|
|
||||||
|
env = neon_env_builder.init_start()
|
||||||
|
client = env.pageserver.http_client()
|
||||||
|
|
||||||
|
tenant_without_timelines_dir = env.initial_tenant
|
||||||
|
log.info(
|
||||||
|
f"Tenant {tenant_without_timelines_dir} becomes broken: it abnormally looses tenants/ directory and is expected to be completely ignored when pageserver restarts"
|
||||||
|
)
|
||||||
|
shutil.rmtree(Path(env.repo_dir) / "tenants" / str(tenant_without_timelines_dir) / "timelines")
|
||||||
|
|
||||||
|
tenant_with_empty_timelines_dir = client.tenant_create()
|
||||||
|
log.info(
|
||||||
|
f"Tenant {tenant_with_empty_timelines_dir} gets all of its timelines deleted: still should be functional"
|
||||||
|
)
|
||||||
|
temp_timelines = client.timeline_list(tenant_with_empty_timelines_dir)
|
||||||
|
for temp_timeline in temp_timelines:
|
||||||
|
client.timeline_delete(
|
||||||
|
tenant_with_empty_timelines_dir, TimelineId(temp_timeline["timeline_id"])
|
||||||
|
)
|
||||||
|
files_in_timelines_dir = sum(
|
||||||
|
1
|
||||||
|
for _p in Path.iterdir(
|
||||||
|
Path(env.repo_dir) / "tenants" / str(tenant_with_empty_timelines_dir) / "timelines"
|
||||||
|
)
|
||||||
|
)
|
||||||
|
assert (
|
||||||
|
files_in_timelines_dir == 0
|
||||||
|
), f"Tenant {tenant_with_empty_timelines_dir} should have an empty timelines/ directory"
|
||||||
|
|
||||||
|
# Trigger timeline reinitialization after pageserver restart
|
||||||
|
env.postgres.stop_all()
|
||||||
|
env.pageserver.stop()
|
||||||
|
env.pageserver.start()
|
||||||
|
|
||||||
|
client = env.pageserver.http_client()
|
||||||
|
tenants = client.tenant_list()
|
||||||
|
|
||||||
|
assert (
|
||||||
|
len(tenants) == 1
|
||||||
|
), "Pageserver should attach only tenants with empty timelines/ dir on restart"
|
||||||
|
loaded_tenant = tenants[0]
|
||||||
|
assert loaded_tenant["id"] == str(
|
||||||
|
tenant_with_empty_timelines_dir
|
||||||
|
), f"Tenant {tenant_with_empty_timelines_dir} should be loaded as the only one with tenants/ directory"
|
||||||
|
assert loaded_tenant["state"] == {
|
||||||
|
"Active": {"background_jobs_running": False}
|
||||||
|
}, "Empty tenant should be loaded and ready for timeline creation"
|
||||||
|
|||||||
@@ -7,19 +7,25 @@
|
|||||||
#
|
#
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import os
|
||||||
|
from pathlib import Path
|
||||||
from typing import List, Tuple
|
from typing import List, Tuple
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
from fixtures.log_helper import log
|
||||||
from fixtures.neon_fixtures import (
|
from fixtures.neon_fixtures import (
|
||||||
NeonEnv,
|
NeonEnv,
|
||||||
NeonEnvBuilder,
|
NeonEnvBuilder,
|
||||||
|
NeonPageserverHttpClient,
|
||||||
Postgres,
|
Postgres,
|
||||||
RemoteStorageKind,
|
RemoteStorageKind,
|
||||||
available_remote_storages,
|
available_remote_storages,
|
||||||
wait_for_last_record_lsn,
|
wait_for_last_record_lsn,
|
||||||
wait_for_upload,
|
wait_for_upload,
|
||||||
|
wait_until,
|
||||||
)
|
)
|
||||||
from fixtures.types import Lsn, TenantId, TimelineId
|
from fixtures.types import Lsn, TenantId, TimelineId
|
||||||
|
from fixtures.utils import query_scalar
|
||||||
|
|
||||||
|
|
||||||
async def tenant_workload(env: NeonEnv, pg: Postgres):
|
async def tenant_workload(env: NeonEnv, pg: Postgres):
|
||||||
@@ -93,3 +99,93 @@ def test_tenants_many(neon_env_builder: NeonEnvBuilder, remote_storage_kind: Rem
|
|||||||
# run final checkpoint manually to flush all the data to remote storage
|
# run final checkpoint manually to flush all the data to remote storage
|
||||||
pageserver_http.timeline_checkpoint(tenant_id, timeline_id)
|
pageserver_http.timeline_checkpoint(tenant_id, timeline_id)
|
||||||
wait_for_upload(pageserver_http, tenant_id, timeline_id, current_lsn)
|
wait_for_upload(pageserver_http, tenant_id, timeline_id, current_lsn)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("remote_storage_kind", [RemoteStorageKind.LOCAL_FS])
|
||||||
|
def test_tenants_attached_after_download(
|
||||||
|
neon_env_builder: NeonEnvBuilder, remote_storage_kind: RemoteStorageKind
|
||||||
|
):
|
||||||
|
neon_env_builder.enable_remote_storage(
|
||||||
|
remote_storage_kind=remote_storage_kind,
|
||||||
|
test_name="remote_storage_kind",
|
||||||
|
)
|
||||||
|
|
||||||
|
data_id = 1
|
||||||
|
data_secret = "very secret secret"
|
||||||
|
|
||||||
|
##### First start, insert secret data and upload it to the remote storage
|
||||||
|
env = neon_env_builder.init_start()
|
||||||
|
pageserver_http = env.pageserver.http_client()
|
||||||
|
pg = env.postgres.create_start("main")
|
||||||
|
|
||||||
|
client = env.pageserver.http_client()
|
||||||
|
|
||||||
|
tenant_id = TenantId(pg.safe_psql("show neon.tenant_id")[0][0])
|
||||||
|
timeline_id = TimelineId(pg.safe_psql("show neon.timeline_id")[0][0])
|
||||||
|
|
||||||
|
for checkpoint_number in range(1, 3):
|
||||||
|
with pg.cursor() as cur:
|
||||||
|
cur.execute(
|
||||||
|
f"""
|
||||||
|
CREATE TABLE t{checkpoint_number}(id int primary key, secret text);
|
||||||
|
INSERT INTO t{checkpoint_number} VALUES ({data_id}, '{data_secret}|{checkpoint_number}');
|
||||||
|
"""
|
||||||
|
)
|
||||||
|
current_lsn = Lsn(query_scalar(cur, "SELECT pg_current_wal_flush_lsn()"))
|
||||||
|
|
||||||
|
# wait until pageserver receives that data
|
||||||
|
wait_for_last_record_lsn(client, tenant_id, timeline_id, current_lsn)
|
||||||
|
|
||||||
|
# run checkpoint manually to be sure that data landed in remote storage
|
||||||
|
pageserver_http.timeline_checkpoint(tenant_id, timeline_id)
|
||||||
|
|
||||||
|
log.info(f"waiting for checkpoint {checkpoint_number} upload")
|
||||||
|
# wait until pageserver successfully uploaded a checkpoint to remote storage
|
||||||
|
wait_for_upload(client, tenant_id, timeline_id, current_lsn)
|
||||||
|
log.info(f"upload of checkpoint {checkpoint_number} is done")
|
||||||
|
|
||||||
|
##### Stop the pageserver, erase its layer file to force it being downloaded from S3
|
||||||
|
env.postgres.stop_all()
|
||||||
|
env.pageserver.stop()
|
||||||
|
|
||||||
|
timeline_dir = Path(env.repo_dir) / "tenants" / str(tenant_id) / "timelines" / str(timeline_id)
|
||||||
|
local_layer_deleted = False
|
||||||
|
for path in Path.iterdir(timeline_dir):
|
||||||
|
if path.name.startswith("00000"):
|
||||||
|
# Looks like a layer file. Remove it
|
||||||
|
os.remove(path)
|
||||||
|
local_layer_deleted = True
|
||||||
|
break
|
||||||
|
assert local_layer_deleted, f"Found no local layer files to delete in directory {timeline_dir}"
|
||||||
|
|
||||||
|
##### Start the pageserver, forcing it to download the layer file and load the timeline into memory
|
||||||
|
env.pageserver.start()
|
||||||
|
client = env.pageserver.http_client()
|
||||||
|
|
||||||
|
wait_until(
|
||||||
|
number_of_iterations=5,
|
||||||
|
interval=1,
|
||||||
|
func=lambda: expect_tenant_to_download_timeline(client, tenant_id),
|
||||||
|
)
|
||||||
|
|
||||||
|
restored_timelines = client.timeline_list(tenant_id)
|
||||||
|
assert (
|
||||||
|
len(restored_timelines) == 1
|
||||||
|
), f"Tenant {tenant_id} should have its timeline reattached after its layer is downloaded from the remote storage"
|
||||||
|
retored_timeline = restored_timelines[0]
|
||||||
|
assert retored_timeline["timeline_id"] == str(
|
||||||
|
timeline_id
|
||||||
|
), f"Tenant {tenant_id} should have its old timeline {timeline_id} restored from the remote storage"
|
||||||
|
|
||||||
|
|
||||||
|
def expect_tenant_to_download_timeline(
|
||||||
|
client: NeonPageserverHttpClient,
|
||||||
|
tenant_id: TenantId,
|
||||||
|
):
|
||||||
|
for tenant in client.tenant_list():
|
||||||
|
if tenant["id"] == str(tenant_id):
|
||||||
|
assert not tenant.get(
|
||||||
|
"has_in_progress_downloads", True
|
||||||
|
), f"Tenant {tenant_id} should have no downloads in progress"
|
||||||
|
return
|
||||||
|
assert False, f"Tenant {tenant_id} is missing on pageserver"
|
||||||
|
|||||||
Vendored
+1
-1
Submodule vendor/postgres-v15 updated: 9383aaa9c2...ff18cec1ee
@@ -19,6 +19,7 @@ anyhow = { version = "1", features = ["backtrace", "std"] }
|
|||||||
bstr = { version = "0.2", features = ["lazy_static", "regex-automata", "serde", "serde1", "serde1-nostd", "std", "unicode"] }
|
bstr = { version = "0.2", features = ["lazy_static", "regex-automata", "serde", "serde1", "serde1-nostd", "std", "unicode"] }
|
||||||
bytes = { version = "1", features = ["serde", "std"] }
|
bytes = { version = "1", features = ["serde", "std"] }
|
||||||
chrono = { version = "0.4", features = ["clock", "libc", "oldtime", "serde", "std", "time", "winapi"] }
|
chrono = { version = "0.4", features = ["clock", "libc", "oldtime", "serde", "std", "time", "winapi"] }
|
||||||
|
crossbeam-utils = { version = "0.8", features = ["once_cell", "std"] }
|
||||||
either = { version = "1", features = ["use_std"] }
|
either = { version = "1", features = ["use_std"] }
|
||||||
fail = { version = "0.5", default-features = false, features = ["failpoints"] }
|
fail = { version = "0.5", default-features = false, features = ["failpoints"] }
|
||||||
hashbrown = { version = "0.12", features = ["ahash", "inline-more", "raw"] }
|
hashbrown = { version = "0.12", features = ["ahash", "inline-more", "raw"] }
|
||||||
|
|||||||
Reference in New Issue
Block a user