Compare commits
11
Commits
f433c5bdb2
..
v1.0.1
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d3bce39783 | ||
|
|
7c8d0bdc99 | ||
|
|
b35017b832 | ||
|
|
ebe222784a | ||
|
|
fd9ea30184 | ||
|
|
654504e3fc | ||
|
|
63a7494982 | ||
|
|
21ce966d82 | ||
|
|
567b1d437c | ||
|
|
cca4e1a0ef | ||
|
|
4f45e4c9a6 |
@@ -0,0 +1,74 @@
|
|||||||
|
name: Build & Release Docker Image
|
||||||
|
|
||||||
|
on:
|
||||||
|
push:
|
||||||
|
tags:
|
||||||
|
- 'v*.*.*'
|
||||||
|
workflow_dispatch:
|
||||||
|
inputs:
|
||||||
|
tag:
|
||||||
|
description: 'Existing tag to release (e.g. v0.1.0)'
|
||||||
|
required: true
|
||||||
|
type: string
|
||||||
|
|
||||||
|
jobs:
|
||||||
|
build-and-push:
|
||||||
|
runs-on: ubuntu-latest
|
||||||
|
permissions:
|
||||||
|
contents: read
|
||||||
|
packages: write
|
||||||
|
env:
|
||||||
|
IMAGE: git.warky.dev/wdevs/pgsql-broker
|
||||||
|
TAG: ${{ github.event_name == 'workflow_dispatch' && inputs.tag || github.ref_name }}
|
||||||
|
steps:
|
||||||
|
- name: Validate release tag
|
||||||
|
run: |
|
||||||
|
case "$TAG" in
|
||||||
|
v*) ;;
|
||||||
|
*) echo "Release tags must start with v (received: $TAG)" >&2; exit 1 ;;
|
||||||
|
esac
|
||||||
|
|
||||||
|
- uses: actions/checkout@v4
|
||||||
|
with:
|
||||||
|
ref: ${{ github.event_name == 'workflow_dispatch' && inputs.tag || github.ref }}
|
||||||
|
|
||||||
|
- uses: docker/setup-buildx-action@v3
|
||||||
|
|
||||||
|
- name: Verify package registry credentials
|
||||||
|
env:
|
||||||
|
PACKAGE_REGISTRY_USERNAME: ${{ secrets.PACKAGE_REGISTRY_USERNAME }}
|
||||||
|
PACKAGE_REGISTRY_TOKEN: ${{ secrets.PACKAGE_REGISTRY_TOKEN }}
|
||||||
|
run: |
|
||||||
|
test -n "$PACKAGE_REGISTRY_USERNAME" || {
|
||||||
|
echo 'PACKAGE_REGISTRY_USERNAME is required to publish the Docker image.' >&2
|
||||||
|
exit 1
|
||||||
|
}
|
||||||
|
test -n "$PACKAGE_REGISTRY_TOKEN" || {
|
||||||
|
echo 'PACKAGE_REGISTRY_TOKEN is required to publish the Docker image.' >&2
|
||||||
|
exit 1
|
||||||
|
}
|
||||||
|
|
||||||
|
- name: Log in to the Warky container registry
|
||||||
|
uses: docker/login-action@v3
|
||||||
|
with:
|
||||||
|
registry: git.warky.dev
|
||||||
|
username: ${{ secrets.PACKAGE_REGISTRY_USERNAME }}
|
||||||
|
password: ${{ secrets.PACKAGE_REGISTRY_TOKEN }}
|
||||||
|
|
||||||
|
- name: Build and push image
|
||||||
|
uses: docker/build-push-action@v5
|
||||||
|
with:
|
||||||
|
context: .
|
||||||
|
file: Dockerfile
|
||||||
|
push: true
|
||||||
|
build-args: |
|
||||||
|
VERSION=${{ env.TAG }}
|
||||||
|
COMMIT=${{ github.sha }}
|
||||||
|
BUILD_TIME=${{ github.event.head_commit.timestamp || github.event.repository.updated_at }}
|
||||||
|
tags: |
|
||||||
|
${{ env.IMAGE }}:${{ env.TAG }}
|
||||||
|
${{ env.IMAGE }}:latest
|
||||||
|
labels: |
|
||||||
|
org.opencontainers.image.source=${{ github.server_url }}/${{ github.repository }}
|
||||||
|
org.opencontainers.image.revision=${{ github.sha }}
|
||||||
|
org.opencontainers.image.version=${{ env.TAG }}
|
||||||
@@ -12,6 +12,19 @@ jobs:
|
|||||||
integration-test:
|
integration-test:
|
||||||
runs-on: ubuntu-latest
|
runs-on: ubuntu-latest
|
||||||
|
|
||||||
|
services:
|
||||||
|
postgres:
|
||||||
|
image: postgres:13
|
||||||
|
env:
|
||||||
|
POSTGRES_DB: broker_test
|
||||||
|
POSTGRES_USER: user
|
||||||
|
POSTGRES_PASSWORD: password
|
||||||
|
options: >-
|
||||||
|
--health-cmd="pg_isready -U user"
|
||||||
|
--health-interval=5s
|
||||||
|
--health-timeout=5s
|
||||||
|
--health-retries=10
|
||||||
|
|
||||||
steps:
|
steps:
|
||||||
- name: Checkout code
|
- name: Checkout code
|
||||||
uses: actions/checkout@v4
|
uses: actions/checkout@v4
|
||||||
@@ -19,16 +32,28 @@ jobs:
|
|||||||
- name: Set up Go
|
- name: Set up Go
|
||||||
uses: actions/setup-go@v5
|
uses: actions/setup-go@v5
|
||||||
with:
|
with:
|
||||||
go-version: '1.25'
|
go-version: '1.26'
|
||||||
cache: true
|
cache: true
|
||||||
|
|
||||||
- name: Set up Python
|
- name: Check formatting
|
||||||
uses: actions/setup-python@v5
|
run: |
|
||||||
with:
|
make fmt
|
||||||
python-version: '3.12'
|
git diff --exit-code -- '*.go'
|
||||||
|
|
||||||
- name: Install podman-compose
|
- name: Install golangci-lint
|
||||||
run: pip install podman-compose
|
run: go install github.com/golangci/golangci-lint/v2/cmd/golangci-lint@latest
|
||||||
|
|
||||||
|
- name: Run lint
|
||||||
|
run: make lint
|
||||||
|
|
||||||
- name: Run all tests
|
- name: Run all tests
|
||||||
run: make test-all
|
env:
|
||||||
|
# act_runner runs this job in its own container alongside the
|
||||||
|
# postgres service container, both on the job's Docker network.
|
||||||
|
# "localhost" from inside the job container is the job container
|
||||||
|
# itself, not the runner host, so the service must be reached by
|
||||||
|
# its network alias (the services: key) and container-internal
|
||||||
|
# port -- not a published host port.
|
||||||
|
TEST_DB_HOST: postgres
|
||||||
|
TEST_DB_PORT: 5432
|
||||||
|
run: make test-ci TEST_DB_HOST="$TEST_DB_HOST" TEST_DB_PORT="$TEST_DB_PORT"
|
||||||
|
|||||||
@@ -19,7 +19,7 @@ jobs:
|
|||||||
- name: Set up Go
|
- name: Set up Go
|
||||||
uses: actions/setup-go@v5
|
uses: actions/setup-go@v5
|
||||||
with:
|
with:
|
||||||
go-version: '1.25'
|
go-version: '1.26'
|
||||||
cache: true
|
cache: true
|
||||||
|
|
||||||
- name: Get version from tag
|
- name: Get version from tag
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
.PHONY: all build clean test test-all test-integration-go test-unit-go test-connection schema-install broker-start broker-stop install deps docker-up docker-down help
|
.PHONY: all build clean test test-all test-ci test-integration-go test-unit-go test-connection generate-test-config schema-install broker-start broker-stop install deps docker-up docker-down docker-build release help
|
||||||
|
|
||||||
# Build variables
|
# Build variables
|
||||||
BINARY_NAME=pgsql-broker
|
BINARY_NAME=pgsql-broker
|
||||||
@@ -30,9 +30,16 @@ COMPOSE_CMD := $(shell \
|
|||||||
|
|
||||||
|
|
||||||
|
|
||||||
|
# Test database connection info. Override in CI to point at a dynamically
|
||||||
|
# assigned Postgres (e.g. a services: block port), avoiding a fixed host port
|
||||||
|
# that can collide with other jobs on a shared runner.
|
||||||
|
TEST_DB_HOST ?= 127.0.0.1
|
||||||
|
TEST_DB_PORT ?= 5433
|
||||||
|
TEST_CONFIG := $(BIN_DIR)/broker.test.runtime.yaml
|
||||||
|
|
||||||
# Version information
|
# Version information
|
||||||
VERSION ?= $(shell git describe --tags --always --dirty 2>/dev/null || echo "dev")
|
VERSION ?= $(shell git describe --tags --always --dirty 2>/dev/null || echo "dev")
|
||||||
BUILD_TIME=$(shell date -u '+2026-01-02_19:58:30')
|
BUILD_TIME=$(shell date -u '+%Y-%m-%d_%H:%M:%S')
|
||||||
COMMIT=$(shell git rev-parse --short HEAD 2>/dev/null || echo "unknown")
|
COMMIT=$(shell git rev-parse --short HEAD 2>/dev/null || echo "unknown")
|
||||||
|
|
||||||
# Inject version info
|
# Inject version info
|
||||||
@@ -63,28 +70,52 @@ test-local-unit: deps ## Run local unit tests
|
|||||||
@echo "Running local unit tests..."
|
@echo "Running local unit tests..."
|
||||||
@$(GO) test -v -race -cover $(shell $(GO) list ./... | grep -v /tests/integration)
|
@$(GO) test -v -race -cover $(shell $(GO) list ./... | grep -v /tests/integration)
|
||||||
|
|
||||||
test-all: test-teardown test-setup test-connection schema-install broker-start test-local-unit test-integration-go broker-stop test-teardown ## Run all unit and integration tests
|
test-all: test-teardown test-setup test-connection schema-install test-local-unit test-integration-go test-teardown ## Run all unit and integration tests (starts its own Postgres via docker-compose)
|
||||||
|
|
||||||
|
test-ci: test-connection schema-install test-local-unit test-integration-go ## Run all unit and integration tests against an externally-provided Postgres (CI services: block)
|
||||||
|
|
||||||
test-connection: deps ## Test database connection with retry
|
test-connection: deps ## Test database connection with retry
|
||||||
@echo "Testing database connection..."
|
@echo "Testing database connection (host=$(TEST_DB_HOST) port=$(TEST_DB_PORT))..."
|
||||||
@$(GO) test -v ./tests/integration/connection_test.go
|
@TEST_DB_HOST=$(TEST_DB_HOST) TEST_DB_PORT=$(TEST_DB_PORT) $(GO) test -v -run '^TestConnection$$' ./tests/integration/...
|
||||||
|
|
||||||
schema-install: build ## Install database schema using the broker CLI
|
generate-test-config: ## (internal) render broker.test.yaml with TEST_DB_HOST/TEST_DB_PORT
|
||||||
|
@mkdir -p $(BIN_DIR)
|
||||||
|
@sed -e "s/^ host: .*/ host: $(TEST_DB_HOST)/" -e "s/^ port: .*/ port: $(TEST_DB_PORT)/" broker.test.yaml > $(TEST_CONFIG)
|
||||||
|
|
||||||
|
schema-install: build generate-test-config ## Install database schema using the broker CLI
|
||||||
@echo "Installing database schema..."
|
@echo "Installing database schema..."
|
||||||
@$(BIN_DIR)/$(BINARY_NAME) install --config broker.test.yaml
|
@$(BIN_DIR)/$(BINARY_NAME) install --config $(TEST_CONFIG)
|
||||||
|
|
||||||
test-setup: build ## Start test environment (docker-compose)
|
test-setup: build ## Start test environment (docker-compose/podman-compose)
|
||||||
@echo "Starting test environment..."
|
@echo "Starting test environment (using $(COMPOSE_CMD))..."
|
||||||
@podman-compose -f tests/docker-compose.yml up -d
|
@if [ "$(CONTAINER_RUNTIME)" = "none" ]; then \
|
||||||
|
echo "Error: Neither Docker nor Podman is installed"; \
|
||||||
|
exit 1; \
|
||||||
|
fi
|
||||||
|
@$(COMPOSE_CMD) -f tests/docker-compose.yml up -d
|
||||||
|
@echo "Waiting for PostgreSQL to be ready..."
|
||||||
|
@for i in 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15; do \
|
||||||
|
if $(COMPOSE_CMD) -f tests/docker-compose.yml exec -T postgres pg_isready -U user > /dev/null 2>&1; then \
|
||||||
|
echo "PostgreSQL is ready"; exit 0; \
|
||||||
|
fi; \
|
||||||
|
sleep 2; \
|
||||||
|
done; \
|
||||||
|
echo "ERROR: PostgreSQL did not become ready in time"; \
|
||||||
|
$(COMPOSE_CMD) -f tests/docker-compose.yml logs postgres; \
|
||||||
|
exit 1
|
||||||
|
|
||||||
test-teardown: ## Stop test environment (docker-compose)
|
test-teardown: ## Stop test environment (docker-compose/podman-compose)
|
||||||
@echo "Stopping test environment..."
|
@echo "Stopping test environment (using $(COMPOSE_CMD))..."
|
||||||
@podman-compose -f tests/docker-compose.yml down -v --rmi all
|
@if [ "$(CONTAINER_RUNTIME)" = "none" ]; then \
|
||||||
@sleep 5 # Give Docker time to release resources
|
echo "Neither Docker nor Podman is installed, skipping teardown"; \
|
||||||
|
else \
|
||||||
|
$(COMPOSE_CMD) -f tests/docker-compose.yml down -v --rmi all || true; \
|
||||||
|
fi
|
||||||
|
@sleep 5 # Give the container runtime time to release resources
|
||||||
|
|
||||||
broker-start: build ## Start the broker in the background
|
broker-start: build generate-test-config ## Start the broker in the background
|
||||||
@echo "Starting broker..."
|
@echo "Starting broker..."
|
||||||
@setsid $(BIN_DIR)/$(BINARY_NAME) start --config broker.test.yaml > broker.log 2>&1 < /dev/null & echo $$! > broker.pid
|
@setsid $(BIN_DIR)/$(BINARY_NAME) start --config $(TEST_CONFIG) > broker.log 2>&1 < /dev/null & echo $$! > broker.pid
|
||||||
@sleep 5 # Give the broker a moment to start
|
@sleep 5 # Give the broker a moment to start
|
||||||
|
|
||||||
broker-stop: ## Stop the broker
|
broker-stop: ## Stop the broker
|
||||||
@@ -97,8 +128,8 @@ broker-stop: ## Stop the broker
|
|||||||
fi
|
fi
|
||||||
|
|
||||||
test-integration-go: ## Run Go integration tests
|
test-integration-go: ## Run Go integration tests
|
||||||
@echo "Running Go integration tests..."
|
@echo "Running Go integration tests (host=$(TEST_DB_HOST) port=$(TEST_DB_PORT))..."
|
||||||
@$(GO) test -v ./tests/integration/...
|
@TEST_DB_HOST=$(TEST_DB_HOST) TEST_DB_PORT=$(TEST_DB_PORT) $(GO) test -v ./tests/integration/...
|
||||||
|
|
||||||
install: build ## Install the binary to GOPATH/bin
|
install: build ## Install the binary to GOPATH/bin
|
||||||
@echo "Installing to GOPATH/bin..."
|
@echo "Installing to GOPATH/bin..."
|
||||||
@@ -172,6 +203,15 @@ docker-down: ## Stop PostgreSQL test database
|
|||||||
fi
|
fi
|
||||||
@echo "PostgreSQL stopped"
|
@echo "PostgreSQL stopped"
|
||||||
|
|
||||||
|
docker-build: ## Build the pgsql-broker runtime image
|
||||||
|
@if [ "$(CONTAINER_RUNTIME)" = "none" ]; then echo "Error: Neither Docker nor Podman is installed"; exit 1; fi
|
||||||
|
@$(CONTAINER_RUNTIME) build \
|
||||||
|
--build-arg VERSION=$(VERSION) \
|
||||||
|
--build-arg COMMIT=$(COMMIT) \
|
||||||
|
--build-arg BUILD_TIME=$(BUILD_TIME) \
|
||||||
|
-t pgsql-broker:$(VERSION) -t pgsql-broker:latest \
|
||||||
|
-f Dockerfile .
|
||||||
|
|
||||||
release: ## Create and push a new release tag (auto-increments patch version)
|
release: ## Create and push a new release tag (auto-increments patch version)
|
||||||
@echo "Creating new release..."
|
@echo "Creating new release..."
|
||||||
@latest_tag=$$(git describe --tags --abbrev=0 2>/dev/null || echo ""); \
|
@latest_tag=$$(git describe --tags --abbrev=0 2>/dev/null || echo ""); \
|
||||||
|
|||||||
@@ -2,6 +2,14 @@
|
|||||||
|
|
||||||
A robust, event-driven job processing system for PostgreSQL that uses LISTEN/NOTIFY for real-time job execution. It supports multiple queues, priority-based scheduling, multi-tenant row-level security, and can be used both as a standalone service or as a Go library.
|
A robust, event-driven job processing system for PostgreSQL that uses LISTEN/NOTIFY for real-time job execution. It supports multiple queues, priority-based scheduling, multi-tenant row-level security, and can be used both as a standalone service or as a Go library.
|
||||||
|
|
||||||
|
## Status
|
||||||
|
[](https://git.warky.dev/wdevs/pgsql-broker/actions?workflow=integration.yml)
|
||||||
|
|
||||||
|
[](https://git.warky.dev/wdevs/pgsql-broker/actions?workflow=release.yml)
|
||||||
|
|
||||||
|
[](https://git.warky.dev/wdevs/pgsql-broker/actions?workflow=docker-release.yml)
|
||||||
|
|
||||||
|
|
||||||
## Features
|
## Features
|
||||||
|
|
||||||
- **Multi-Database Support**: Single broker process can manage multiple database connections
|
- **Multi-Database Support**: Single broker process can manage multiple database connections
|
||||||
@@ -264,6 +272,33 @@ See the [examples](./examples/) directory for complete examples.
|
|||||||
|
|
||||||
## Docker
|
## Docker
|
||||||
|
|
||||||
|
### Prebuilt image
|
||||||
|
|
||||||
|
Published to `git.warky.dev/wdevs/pgsql-broker` on every `v*.*.*` tag (`.github/workflows/docker-release.yml`), tagged with the version and `latest`.
|
||||||
|
|
||||||
|
```bash
|
||||||
|
docker pull git.warky.dev/wdevs/pgsql-broker:latest
|
||||||
|
docker run --rm -v $(pwd)/broker.yaml:/etc/pgsql-broker/broker.yaml:ro git.warky.dev/wdevs/pgsql-broker:latest
|
||||||
|
```
|
||||||
|
|
||||||
|
Config must be mounted at `/etc/pgsql-broker/broker.yaml` — no config is baked into the image.
|
||||||
|
|
||||||
|
### Prebuilt image: `docker-compose.prebuilt.yml`
|
||||||
|
|
||||||
|
Same stack as `docker-compose.yml` below, but pulls `git.warky.dev/wdevs/pgsql-broker:latest` instead of building from source — no local Go toolchain or Dockerfile needed.
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cp broker.docker.example.yaml broker.docker.yaml # set broker_runtime password
|
||||||
|
cp .env.example .env # set POSTGRES_PASSWORD + 3 BROKER_*_PASSWORD
|
||||||
|
docker-compose -f docker-compose.prebuilt.yml up -d
|
||||||
|
```
|
||||||
|
|
||||||
|
### Building the image locally
|
||||||
|
|
||||||
|
```bash
|
||||||
|
make docker-build # builds pgsql-broker:<VERSION> and pgsql-broker:latest via Dockerfile
|
||||||
|
```
|
||||||
|
|
||||||
### Production: `Dockerfile` + `docker-compose.yml`
|
### Production: `Dockerfile` + `docker-compose.yml`
|
||||||
|
|
||||||
Multi-stage build (`golang:1.26-alpine` → `alpine:3.22`, static binary, non-root user). The compose stack runs Postgres, a one-shot `migrate` service (`install --with-roles`), then the `broker` service connecting as `broker_runtime`.
|
Multi-stage build (`golang:1.26-alpine` → `alpine:3.22`, static binary, non-root user). The compose stack runs Postgres, a one-shot `migrate` service (`install --with-roles`), then the `broker` service connecting as `broker_runtime`.
|
||||||
@@ -372,7 +407,7 @@ go test -v ./tests/integration/... # integration tests (needs local Postgres
|
|||||||
docker-compose -f docker-compose.test.yml up --build --abort-on-container-exit --exit-code-from tests
|
docker-compose -f docker-compose.test.yml up --build --abort-on-container-exit --exit-code-from tests
|
||||||
```
|
```
|
||||||
|
|
||||||
Integration tests expect Postgres reachable at `localhost:5433` (see `tests/integration/`), including `rls_test.go` (multi-tenant isolation) and `stage5_test.go`.
|
Integration tests expect Postgres reachable at `127.0.0.1:5433` (see `tests/integration/`), including `rls_test.go` (multi-tenant isolation) and `stage5_test.go`. `127.0.0.1` is used instead of `localhost` because some CI Docker hosts publish container ports on IPv4 only, and `localhost` can resolve to `::1` first and fail.
|
||||||
|
|
||||||
### Project Structure
|
### Project Structure
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -1,6 +1,6 @@
|
|||||||
databases:
|
databases:
|
||||||
- name: test
|
- name: test
|
||||||
host: localhost
|
host: 127.0.0.1
|
||||||
port: 5433
|
port: 5433
|
||||||
database: broker_test
|
database: broker_test
|
||||||
user: user
|
user: user
|
||||||
|
|||||||
+2
-1
@@ -227,7 +227,8 @@ func runInstall() error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Install/verify on all configured databases
|
// Install/verify on all configured databases
|
||||||
for i, dbCfg := range cfg.Databases {
|
for i := range cfg.Databases {
|
||||||
|
dbCfg := &cfg.Databases[i]
|
||||||
logger.Info("processing database", "index", i, "name", dbCfg.Name, "host", dbCfg.Host, "database", dbCfg.Database)
|
logger.Info("processing database", "index", i, "name", dbCfg.Name, "host", dbCfg.Host, "database", dbCfg.Database)
|
||||||
|
|
||||||
// Create database adapter. With --with-roles, the config file's own
|
// Create database adapter. With --with-roles, the config file's own
|
||||||
|
|||||||
@@ -0,0 +1,46 @@
|
|||||||
|
services:
|
||||||
|
postgres:
|
||||||
|
image: docker.io/library/postgres:16-alpine
|
||||||
|
environment:
|
||||||
|
POSTGRES_DB: broker
|
||||||
|
POSTGRES_USER: postgres
|
||||||
|
POSTGRES_PASSWORD: ${POSTGRES_PASSWORD:?set POSTGRES_PASSWORD in .env}
|
||||||
|
volumes:
|
||||||
|
- postgres-data:/var/lib/postgresql/data
|
||||||
|
healthcheck:
|
||||||
|
test: ["CMD-SHELL", "pg_isready -U postgres -d broker"]
|
||||||
|
interval: 2s
|
||||||
|
timeout: 3s
|
||||||
|
retries: 30
|
||||||
|
restart: unless-stopped
|
||||||
|
|
||||||
|
# One-shot: applies migrations and creates/rotates the least-privilege
|
||||||
|
# broker_admin/broker_runtime/broker_enqueue roles, then exits. The
|
||||||
|
# broker service below only starts once this completes successfully.
|
||||||
|
migrate:
|
||||||
|
image: git.warky.dev/wdevs/pgsql-broker:latest
|
||||||
|
depends_on:
|
||||||
|
postgres:
|
||||||
|
condition: service_healthy
|
||||||
|
volumes:
|
||||||
|
- ./broker.docker.yaml:/etc/pgsql-broker/broker.yaml:ro
|
||||||
|
environment:
|
||||||
|
PGUSER: postgres
|
||||||
|
PGPASSWORD: ${POSTGRES_PASSWORD:?set POSTGRES_PASSWORD in .env}
|
||||||
|
BROKER_ADMIN_PASSWORD: ${BROKER_ADMIN_PASSWORD:?set BROKER_ADMIN_PASSWORD in .env}
|
||||||
|
BROKER_RUNTIME_PASSWORD: ${BROKER_RUNTIME_PASSWORD:?set BROKER_RUNTIME_PASSWORD in .env}
|
||||||
|
BROKER_ENQUEUE_PASSWORD: ${BROKER_ENQUEUE_PASSWORD:?set BROKER_ENQUEUE_PASSWORD in .env}
|
||||||
|
command: ["install", "--with-roles"]
|
||||||
|
restart: "no"
|
||||||
|
|
||||||
|
broker:
|
||||||
|
image: git.warky.dev/wdevs/pgsql-broker:latest
|
||||||
|
depends_on:
|
||||||
|
migrate:
|
||||||
|
condition: service_completed_successfully
|
||||||
|
volumes:
|
||||||
|
- ./broker.docker.yaml:/etc/pgsql-broker/broker.yaml:ro
|
||||||
|
restart: unless-stopped
|
||||||
|
|
||||||
|
volumes:
|
||||||
|
postgres-data:
|
||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -22,6 +23,10 @@ type PostgresConfig struct {
|
|||||||
MaxIdleConns int
|
MaxIdleConns int
|
||||||
ConnMaxLifetime time.Duration
|
ConnMaxLifetime time.Duration
|
||||||
ConnMaxIdleTime time.Duration
|
ConnMaxIdleTime time.Duration
|
||||||
|
// ApplicationName identifies this instance's pool connections in
|
||||||
|
// pg_stat_activity (e.g. "PGSQL_BROKER_INSTANCE1"). The LISTEN
|
||||||
|
// connection appends "_LISTENER" to this value.
|
||||||
|
ApplicationName string
|
||||||
}
|
}
|
||||||
|
|
||||||
// PostgresAdapter implements DBAdapter for PostgreSQL
|
// PostgresAdapter implements DBAdapter for PostgreSQL
|
||||||
@@ -170,7 +175,7 @@ func (p *PostgresAdapter) Query(ctx context.Context, query string, args ...inter
|
|||||||
|
|
||||||
// Listen starts listening on a PostgreSQL notification channel
|
// Listen starts listening on a PostgreSQL notification channel
|
||||||
func (p *PostgresAdapter) Listen(ctx context.Context, channel string, handler NotificationHandler) error {
|
func (p *PostgresAdapter) Listen(ctx context.Context, channel string, handler NotificationHandler) error {
|
||||||
connStr := p.buildConnectionString()
|
connStr := p.buildConnectionStringWithAppName(p.config.ApplicationName + "_LISTENER")
|
||||||
|
|
||||||
reportProblem := func(ev pq.ListenerEventType, err error) {
|
reportProblem := func(ev pq.ListenerEventType, err error) {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -211,7 +216,11 @@ func (p *PostgresAdapter) Listen(ctx context.Context, channel string, handler No
|
|||||||
p.logger.Info("stopping listener", "channel", channel)
|
p.logger.Info("stopping listener", "channel", channel)
|
||||||
return
|
return
|
||||||
case <-time.After(90 * time.Second):
|
case <-time.After(90 * time.Second):
|
||||||
SafeGo(p.logger, "listener-ping-"+channel, func() { listener.Ping() })
|
SafeGo(p.logger, "listener-ping-"+channel, func() {
|
||||||
|
if err := listener.Ping(); err != nil {
|
||||||
|
p.logger.Error("listener ping failed", "channel", channel, "error", err)
|
||||||
|
}
|
||||||
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
@@ -232,24 +241,42 @@ func (p *PostgresAdapter) Unlisten(ctx context.Context, channel string) error {
|
|||||||
return listener.Unlisten(channel)
|
return listener.Unlisten(channel)
|
||||||
}
|
}
|
||||||
|
|
||||||
// buildConnectionString builds a PostgreSQL connection string
|
// buildConnectionString builds a PostgreSQL connection string for the
|
||||||
|
// pooled connection, using the adapter's own application name.
|
||||||
func (p *PostgresAdapter) buildConnectionString() string {
|
func (p *PostgresAdapter) buildConnectionString() string {
|
||||||
|
return p.buildConnectionStringWithAppName(p.config.ApplicationName)
|
||||||
|
}
|
||||||
|
|
||||||
|
// buildConnectionStringWithAppName builds a PostgreSQL connection string
|
||||||
|
// with the given application_name, so pooled and LISTEN connections can be
|
||||||
|
// told apart in pg_stat_activity.
|
||||||
|
func (p *PostgresAdapter) buildConnectionStringWithAppName(appName string) string {
|
||||||
sslMode := p.config.SSLMode
|
sslMode := p.config.SSLMode
|
||||||
if sslMode == "" {
|
if sslMode == "" {
|
||||||
sslMode = "disable"
|
sslMode = "disable"
|
||||||
}
|
}
|
||||||
|
|
||||||
return fmt.Sprintf(
|
return fmt.Sprintf(
|
||||||
"host=%s port=%d user=%s password=%s dbname=%s sslmode=%s options='-c search_path=broker,public'",
|
"host=%s port=%d user=%s password=%s dbname=%s sslmode=%s application_name=%s options='-c search_path=broker,public'",
|
||||||
p.config.Host,
|
p.config.Host,
|
||||||
p.config.Port,
|
p.config.Port,
|
||||||
p.config.User,
|
p.config.User,
|
||||||
p.config.Password,
|
p.config.Password,
|
||||||
p.config.Database,
|
p.config.Database,
|
||||||
sslMode,
|
sslMode,
|
||||||
|
quoteDSNValue(appName),
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// quoteDSNValue escapes a value for use in a libpq keyword/value connection
|
||||||
|
// string, single-quoting it and backslash-escaping embedded backslashes and
|
||||||
|
// quotes per the libpq connection string format.
|
||||||
|
func quoteDSNValue(v string) string {
|
||||||
|
v = strings.ReplaceAll(v, `\`, `\\`)
|
||||||
|
v = strings.ReplaceAll(v, `'`, `\'`)
|
||||||
|
return "'" + v + "'"
|
||||||
|
}
|
||||||
|
|
||||||
// Conn returns a single physical connection pinned out of the pool, for
|
// Conn returns a single physical connection pinned out of the pool, for
|
||||||
// session-scoped operations (e.g. pg_try_advisory_lock) that must survive
|
// session-scoped operations (e.g. pg_try_advisory_lock) that must survive
|
||||||
// across calls and must not be silently reaped or handed to another caller
|
// across calls and must not be silently reaped or handed to another caller
|
||||||
|
|||||||
@@ -42,14 +42,15 @@ func (b *Broker) Start() error {
|
|||||||
b.logger.Info("starting broker", "database_count", len(b.config.Databases))
|
b.logger.Info("starting broker", "database_count", len(b.config.Databases))
|
||||||
|
|
||||||
// Create and start an instance for each database
|
// Create and start an instance for each database
|
||||||
for i, dbCfg := range b.config.Databases {
|
for i := range b.config.Databases {
|
||||||
|
dbCfg := &b.config.Databases[i]
|
||||||
b.logger.Info("starting database instance", "name", dbCfg.Name, "host", dbCfg.Host, "database", dbCfg.Database)
|
b.logger.Info("starting database instance", "name", dbCfg.Name, "host", dbCfg.Host, "database", dbCfg.Database)
|
||||||
|
|
||||||
// Create database adapter
|
// Create database adapter
|
||||||
dbAdapter := adapter.NewPostgresAdapter(dbCfg.ToPostgresConfig(), b.logger)
|
dbAdapter := adapter.NewPostgresAdapter(dbCfg.ToPostgresConfig(), b.logger)
|
||||||
|
|
||||||
// Create database instance
|
// Create database instance
|
||||||
instance, err := NewDatabaseInstance(b.config, &dbCfg, dbAdapter, b.logger, b.version, b.ctx)
|
instance, err := NewDatabaseInstance(b.config, dbCfg, dbAdapter, b.logger, b.version, b.ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
// Stop any already-started instances
|
// Stop any already-started instances
|
||||||
b.stopInstances()
|
b.stopInstances()
|
||||||
|
|||||||
@@ -2,10 +2,11 @@ package config
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/spf13/viper"
|
|
||||||
"git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter"
|
"git.warky.dev/wdevs/pgsql-broker/pkg/broker/adapter"
|
||||||
|
"github.com/spf13/viper"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Config holds all broker configuration
|
// Config holds all broker configuration
|
||||||
@@ -47,7 +48,7 @@ type BrokerConfig struct {
|
|||||||
QueueTimerSec int `mapstructure:"queue_timer_sec"`
|
QueueTimerSec int `mapstructure:"queue_timer_sec"`
|
||||||
QueueBufferSize int `mapstructure:"queue_buffer_size"`
|
QueueBufferSize int `mapstructure:"queue_buffer_size"`
|
||||||
WorkerIdleTimeoutSec int `mapstructure:"worker_idle_timeout_sec"`
|
WorkerIdleTimeoutSec int `mapstructure:"worker_idle_timeout_sec"`
|
||||||
NotifyRetrySeconds time.Duration `mapstructure:"notify_retry_seconds"`
|
NotifyRetryInterval time.Duration `mapstructure:"notify_retry_seconds"`
|
||||||
EnableDebug bool `mapstructure:"enable_debug"`
|
EnableDebug bool `mapstructure:"enable_debug"`
|
||||||
// LeaseSeconds is how long a claimed job's lease is valid for before
|
// LeaseSeconds is how long a claimed job's lease is valid for before
|
||||||
// broker_recover_stale_jobs considers it abandoned.
|
// broker_recover_stale_jobs considers it abandoned.
|
||||||
@@ -132,7 +133,8 @@ func validateConfig(config *Config) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Validate each database configuration
|
// Validate each database configuration
|
||||||
for i, db := range config.Databases {
|
for i := range config.Databases {
|
||||||
|
db := &config.Databases[i]
|
||||||
if db.Name == "" {
|
if db.Name == "" {
|
||||||
return fmt.Errorf("database[%d]: name is required", i)
|
return fmt.Errorf("database[%d]: name is required", i)
|
||||||
}
|
}
|
||||||
@@ -195,5 +197,6 @@ func (d *DatabaseConfig) ToPostgresConfig() adapter.PostgresConfig {
|
|||||||
MaxIdleConns: d.MaxIdleConns,
|
MaxIdleConns: d.MaxIdleConns,
|
||||||
ConnMaxLifetime: d.ConnMaxLifetime,
|
ConnMaxLifetime: d.ConnMaxLifetime,
|
||||||
ConnMaxIdleTime: d.ConnMaxIdleTime,
|
ConnMaxIdleTime: d.ConnMaxIdleTime,
|
||||||
|
ApplicationName: fmt.Sprintf("PGSQL_BROKER_%s", strings.ToUpper(d.Name)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -205,7 +205,9 @@ func (i *Installer) ApplyMigrations(ctx context.Context) error {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if err := execStatements(ctx, tx, string(content)); err != nil {
|
if err := execStatements(ctx, tx, string(content)); err != nil {
|
||||||
tx.Rollback()
|
if rbErr := tx.Rollback(); rbErr != nil {
|
||||||
|
i.logger.Error("failed to rollback migration transaction", "version", m.version, "error", rbErr)
|
||||||
|
}
|
||||||
return fmt.Errorf("failed to apply migration %s: %w", m.name, err)
|
return fmt.Errorf("failed to apply migration %s: %w", m.name, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -213,7 +215,9 @@ func (i *Installer) ApplyMigrations(ctx context.Context) error {
|
|||||||
"INSERT INTO broker.broker_schema_migrations (version, name) VALUES ($1, $2)",
|
"INSERT INTO broker.broker_schema_migrations (version, name) VALUES ($1, $2)",
|
||||||
m.version, m.name,
|
m.version, m.name,
|
||||||
); err != nil {
|
); err != nil {
|
||||||
tx.Rollback()
|
if rbErr := tx.Rollback(); rbErr != nil {
|
||||||
|
i.logger.Error("failed to rollback migration transaction", "version", m.version, "error", rbErr)
|
||||||
|
}
|
||||||
return fmt.Errorf("failed to record migration %s: %w", m.name, err)
|
return fmt.Errorf("failed to record migration %s: %w", m.name, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -293,7 +297,9 @@ func (i *Installer) InstallRoles(ctx context.Context, passwords RolePasswords) e
|
|||||||
rendered := replacer.Replace(string(content))
|
rendered := replacer.Replace(string(content))
|
||||||
|
|
||||||
if err := execStatements(ctx, tx, rendered); err != nil {
|
if err := execStatements(ctx, tx, rendered); err != nil {
|
||||||
tx.Rollback()
|
if rbErr := tx.Rollback(); rbErr != nil {
|
||||||
|
i.logger.Error("failed to rollback roles script transaction", "name", name, "error", rbErr)
|
||||||
|
}
|
||||||
return fmt.Errorf("failed to apply roles script %s: %w", name, err)
|
return fmt.Errorf("failed to apply roles script %s: %w", name, err)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -326,7 +332,7 @@ func execStatements(ctx context.Context, tx adapter.DBTransaction, sqlText strin
|
|||||||
// $$-quoted function bodies intact.
|
// $$-quoted function bodies intact.
|
||||||
// splitSQLStatements splits a SQL script into individual statements on
|
// splitSQLStatements splits a SQL script into individual statements on
|
||||||
// top-level semicolons, ignoring semicolons that appear inside single-quoted
|
// top-level semicolons, ignoring semicolons that appear inside single-quoted
|
||||||
// strings ('...', with '' as an escaped quote), double-quoted identifiers,
|
// strings ('...', with ” as an escaped quote), double-quoted identifiers,
|
||||||
// line comments (--), and dollar-quoted bodies ($$...$$ or $tag$...$tag$).
|
// line comments (--), and dollar-quoted bodies ($$...$$ or $tag$...$tag$).
|
||||||
func splitSQLStatements(sqlText string) []string {
|
func splitSQLStatements(sqlText string) []string {
|
||||||
var result []string
|
var result []string
|
||||||
|
|||||||
@@ -204,26 +204,36 @@ func (w *Worker) processJobs(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if err := w.setTenantTx(ctx, tx); err != nil {
|
if err := w.setTenantTx(ctx, tx); err != nil {
|
||||||
tx.Rollback()
|
if rbErr := tx.Rollback(); rbErr != nil {
|
||||||
|
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||||
|
}
|
||||||
w.logger.Error("failed to set tenant", "error", err)
|
w.logger.Error("failed to set tenant", "error", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
jobID, leaseToken, err := w.fetchNextJobTx(ctx, tx)
|
jobID, leaseToken, err := w.fetchNextJobTx(ctx, tx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
tx.Rollback()
|
if rbErr := tx.Rollback(); rbErr != nil {
|
||||||
|
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||||
|
}
|
||||||
w.logger.Error("failed to fetch job", "error", err)
|
w.logger.Error("failed to fetch job", "error", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if jobID <= 0 {
|
if jobID <= 0 {
|
||||||
tx.Rollback() // No job found, rollback
|
// No job found, rollback
|
||||||
|
if rbErr := tx.Rollback(); rbErr != nil {
|
||||||
|
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||||
|
}
|
||||||
return // No more jobs
|
return // No more jobs
|
||||||
}
|
}
|
||||||
|
|
||||||
// Run the job
|
// Run the job
|
||||||
if err := w.runJobTx(ctx, tx, jobID, leaseToken); err != nil {
|
if err := w.runJobTx(ctx, tx, jobID, leaseToken); err != nil {
|
||||||
tx.Rollback() // Rollback on genuine infra failure
|
// Rollback on genuine infra failure
|
||||||
|
if rbErr := tx.Rollback(); rbErr != nil {
|
||||||
|
w.logger.Error("failed to rollback transaction", "error", rbErr)
|
||||||
|
}
|
||||||
w.logger.Error("failed to run job", "job_id", jobID, "error", err)
|
w.logger.Error("failed to run job", "job_id", jobID, "error", err)
|
||||||
} else {
|
} else {
|
||||||
if err := tx.Commit(); err != nil {
|
if err := tx.Commit(); err != nil {
|
||||||
@@ -247,13 +257,13 @@ func (w *Worker) setTenantTx(ctx context.Context, tx adapter.DBTransaction) erro
|
|||||||
|
|
||||||
// fetchNextJobTx fetches the next job from the queue within a transaction,
|
// fetchNextJobTx fetches the next job from the queue within a transaction,
|
||||||
// claiming it with a lease that must be presented back to broker_run.
|
// claiming it with a lease that must be presented back to broker_run.
|
||||||
func (w *Worker) fetchNextJobTx(ctx context.Context, tx adapter.DBTransaction) (int64, string, error) {
|
func (w *Worker) fetchNextJobTx(ctx context.Context, tx adapter.DBTransaction) (jobID int64, leaseToken string, err error) {
|
||||||
var retval int
|
var retval int
|
||||||
var errmsg string
|
var errmsg string
|
||||||
var nullableJobID sql.NullInt64
|
var nullableJobID sql.NullInt64
|
||||||
var nullableLeaseToken sql.NullString
|
var nullableLeaseToken sql.NullString
|
||||||
|
|
||||||
err := tx.QueryRow(ctx,
|
err = tx.QueryRow(ctx,
|
||||||
"SELECT p_retval, p_errmsg, p_job_id, p_lease_token FROM broker.broker_get($1, $2, $3)",
|
"SELECT p_retval, p_errmsg, p_job_id, p_lease_token FROM broker.broker_get($1, $2, $3)",
|
||||||
w.QueueNumber, w.InstanceID, w.leaseSeconds,
|
w.QueueNumber, w.InstanceID, w.leaseSeconds,
|
||||||
).Scan(&retval, &errmsg, &nullableJobID, &nullableLeaseToken)
|
).Scan(&retval, &errmsg, &nullableJobID, &nullableLeaseToken)
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package integration
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"database/sql"
|
"database/sql"
|
||||||
|
"fmt"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -10,7 +11,7 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
func TestConnection(t *testing.T) {
|
func TestConnection(t *testing.T) {
|
||||||
connStr := "user=user password=password dbname=broker_test port=5433 sslmode=disable"
|
connStr := fmt.Sprintf("user=user password=password dbname=broker_test host=%s port=%d sslmode=disable", testDBHost(), testDBPort())
|
||||||
var db *sql.DB
|
var db *sql.DB
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package integration
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
|
"fmt"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
_ "github.com/lib/pq"
|
_ "github.com/lib/pq"
|
||||||
@@ -47,7 +48,7 @@ func TestRLSTenantIsolation(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
runtimeDB, err := sql.Open("postgres",
|
runtimeDB, err := sql.Open("postgres",
|
||||||
"user=test_broker_runtime password=test-pass dbname=broker_test host=localhost port=5433 sslmode=disable options='-c search_path=broker,public'")
|
fmt.Sprintf("user=test_broker_runtime password=test-pass dbname=broker_test host=%s port=%d sslmode=disable options='-c search_path=broker,public'", testDBHost(), testDBPort()))
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
defer runtimeDB.Close()
|
defer runtimeDB.Close()
|
||||||
require.NoError(t, runtimeDB.Ping())
|
require.NoError(t, runtimeDB.Ping())
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package integration
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -14,11 +15,13 @@ import (
|
|||||||
"git.warky.dev/wdevs/pgsql-broker/pkg/broker/install"
|
"git.warky.dev/wdevs/pgsql-broker/pkg/broker/install"
|
||||||
)
|
)
|
||||||
|
|
||||||
const stage5ConnStr = "user=user password=password dbname=broker_test host=localhost port=5433 sslmode=disable"
|
func stage5ConnStr() string {
|
||||||
|
return fmt.Sprintf("user=user password=password dbname=broker_test host=%s port=%d sslmode=disable", testDBHost(), testDBPort())
|
||||||
|
}
|
||||||
|
|
||||||
func newStage5Adapter(logger adapter.Logger) *adapter.PostgresAdapter {
|
func newStage5Adapter(logger adapter.Logger) *adapter.PostgresAdapter {
|
||||||
return adapter.NewPostgresAdapter(adapter.PostgresConfig{
|
return adapter.NewPostgresAdapter(adapter.PostgresConfig{
|
||||||
Host: "localhost", Port: 5433, Database: "broker_test",
|
Host: testDBHost(), Port: testDBPort(), Database: "broker_test",
|
||||||
User: "user", Password: "password", SSLMode: "disable",
|
User: "user", Password: "password", SSLMode: "disable",
|
||||||
MaxOpenConns: 10, MaxIdleConns: 2,
|
MaxOpenConns: 10, MaxIdleConns: 2,
|
||||||
ConnMaxLifetime: 5 * time.Minute, ConnMaxIdleTime: 10 * time.Minute,
|
ConnMaxLifetime: 5 * time.Minute, ConnMaxIdleTime: 10 * time.Minute,
|
||||||
@@ -30,7 +33,7 @@ func newStage5Adapter(logger adapter.Logger) *adapter.PostgresAdapter {
|
|||||||
func setupStage5Schema(t *testing.T) *sql.DB {
|
func setupStage5Schema(t *testing.T) *sql.DB {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
db, err := connectWithRetry(stage5ConnStr, 10, 2*time.Second)
|
db, err := connectWithRetry(stage5ConnStr(), 10, 2*time.Second)
|
||||||
require.NoError(t, err)
|
require.NoError(t, err)
|
||||||
|
|
||||||
cleanupSchema(t, db)
|
cleanupSchema(t, db)
|
||||||
|
|||||||
@@ -0,0 +1,26 @@
|
|||||||
|
package integration
|
||||||
|
|
||||||
|
import (
|
||||||
|
"os"
|
||||||
|
"strconv"
|
||||||
|
)
|
||||||
|
|
||||||
|
// testDBHost and testDBPort let CI point the integration suite at a
|
||||||
|
// dynamically-assigned Postgres (TEST_DB_HOST/TEST_DB_PORT), avoiding a fixed
|
||||||
|
// host port that can collide with other jobs on a shared runner. Local dev
|
||||||
|
// keeps working unset, defaulting to the docker-compose test stack.
|
||||||
|
func testDBHost() string {
|
||||||
|
if h := os.Getenv("TEST_DB_HOST"); h != "" {
|
||||||
|
return h
|
||||||
|
}
|
||||||
|
return "127.0.0.1"
|
||||||
|
}
|
||||||
|
|
||||||
|
func testDBPort() int {
|
||||||
|
if p := os.Getenv("TEST_DB_PORT"); p != "" {
|
||||||
|
if n, err := strconv.Atoi(p); err == nil {
|
||||||
|
return n
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return 5433
|
||||||
|
}
|
||||||
@@ -3,6 +3,7 @@ package integration
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
@@ -27,7 +28,7 @@ func TestBrokerWorkflow(t *testing.T) {
|
|||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
|
|
||||||
// Database connection string
|
// Database connection string
|
||||||
connStr := "user=user password=password dbname=broker_test host=localhost port=5433 sslmode=disable"
|
connStr := fmt.Sprintf("user=user password=password dbname=broker_test host=%s port=%d sslmode=disable", testDBHost(), testDBPort())
|
||||||
|
|
||||||
// Connect to database with retry logic
|
// Connect to database with retry logic
|
||||||
db, err := connectWithRetry(connStr, 10, 2*time.Second)
|
db, err := connectWithRetry(connStr, 10, 2*time.Second)
|
||||||
@@ -43,8 +44,8 @@ func TestBrokerWorkflow(t *testing.T) {
|
|||||||
|
|
||||||
// Create database adapter
|
// Create database adapter
|
||||||
postgresConfig := adapter.PostgresConfig{
|
postgresConfig := adapter.PostgresConfig{
|
||||||
Host: "localhost",
|
Host: testDBHost(),
|
||||||
Port: 5433,
|
Port: testDBPort(),
|
||||||
Database: "broker_test",
|
Database: "broker_test",
|
||||||
User: "user",
|
User: "user",
|
||||||
Password: "password",
|
Password: "password",
|
||||||
@@ -78,8 +79,8 @@ func TestBrokerWorkflow(t *testing.T) {
|
|||||||
Databases: []config.DatabaseConfig{
|
Databases: []config.DatabaseConfig{
|
||||||
{
|
{
|
||||||
Name: "test_db",
|
Name: "test_db",
|
||||||
Host: "localhost",
|
Host: testDBHost(),
|
||||||
Port: 5433,
|
Port: testDBPort(),
|
||||||
Database: "broker_test",
|
Database: "broker_test",
|
||||||
User: "user",
|
User: "user",
|
||||||
Password: "password",
|
Password: "password",
|
||||||
@@ -93,7 +94,7 @@ func TestBrokerWorkflow(t *testing.T) {
|
|||||||
QueueTimerSec: 1, // Short interval for testing
|
QueueTimerSec: 1, // Short interval for testing
|
||||||
QueueBufferSize: 10,
|
QueueBufferSize: 10,
|
||||||
WorkerIdleTimeoutSec: 5,
|
WorkerIdleTimeoutSec: 5,
|
||||||
NotifyRetrySeconds: 5 * time.Second,
|
NotifyRetryInterval: 5 * time.Second,
|
||||||
EnableDebug: true,
|
EnableDebug: true,
|
||||||
},
|
},
|
||||||
Logging: config.LoggingConfig{
|
Logging: config.LoggingConfig{
|
||||||
|
|||||||
Reference in New Issue
Block a user