diff --git a/.gitignore b/.gitignore index dd56108..b9303cc 100644 --- a/.gitignore +++ b/.gitignore @@ -17,6 +17,8 @@ dist/ .env .env.local .env.*.local +config/local/*.env +!config/local/*.env.example # IDE .vscode/ diff --git a/ansible/files/postgresql-xinfra@.service b/ansible/files/postgresql-xinfra@.service new file mode 100644 index 0000000..1d45f05 --- /dev/null +++ b/ansible/files/postgresql-xinfra@.service @@ -0,0 +1,21 @@ +[Unit] +Description=XINFRA PostgreSQL delivery instance %i +After=network-online.target +Wants=network-online.target + +[Service] +Type=notify +User=postgres +Group=postgres +RuntimeDirectory=postgresql-xinfra/%i +RuntimeDirectoryMode=0755 +ExecStart=/bin/sh -c 'exec /usr/lib/postgresql/${POSTGRESQL_VERSION}/bin/postgres -D "${POSTGRESQL_DATA_DIR}" -c "config_file=${POSTGRESQL_CONFIG_FILE}" -c "hba_file=${POSTGRESQL_HBA_FILE}"' +ExecReload=/bin/kill -HUP $MAINPID +Restart=on-failure +RestartSec=5s +TimeoutStartSec=120s +TimeoutStopSec=120s +LimitNOFILE=65535 + +[Install] +WantedBy=multi-user.target diff --git a/ansible/postgresql-deploy.yml b/ansible/postgresql-deploy.yml new file mode 100644 index 0000000..d2f0d22 --- /dev/null +++ b/ansible/postgresql-deploy.yml @@ -0,0 +1,331 @@ +--- +- name: Native PostgreSQL delivery + hosts: "{{ target_hosts }}" + become: true + gather_facts: true + any_errors_fatal: true + vars: + postgresql_version_value: "{{ postgresql_version | default('16') }}" + postgresql_instance_plan: "{{ postgresql_instances[inventory_hostname] | default({}) }}" + postgresql_instance_id: "{{ postgresql_instance_plan.instance_id | default(instance_name) }}" + postgresql_role: "{{ postgresql_instance_plan.role | default('standalone') }}" + postgresql_data_root: "{{ data_root | default('/data/postgresql') }}" + postgresql_data_dir: "{{ postgresql_instance_plan.data_dir | default(postgresql_data_root ~ '/' ~ postgresql_instance_id ~ '/data') }}" + postgresql_conf_dir: "{{ postgresql_instance_plan.config_dir | default(postgresql_data_root ~ '/' ~ postgresql_instance_id ~ '/conf') }}" + postgresql_log_dir: "{{ postgresql_instance_plan.log_dir | default(postgresql_data_root ~ '/' ~ postgresql_instance_id ~ '/log') }}" + postgresql_run_dir: "/run/postgresql-xinfra/{{ postgresql_instance_id }}" + postgresql_config_file: "{{ postgresql_conf_dir }}/postgresql.conf" + postgresql_hba_file: "{{ postgresql_conf_dir }}/pg_hba.conf" + postgresql_port_value: "{{ postgresql_instance_plan.port | default(postgresql_port | default(15432)) | int }}" + postgresql_max_connections: "{{ max_connections | default(200, true) | int }}" + postgresql_shared_buffers: "{{ ((memory_mb | default(4096) | int) * 25 / 100) | int }}MB" + postgresql_effective_cache_size: "{{ ((memory_mb | default(4096) | int) * 70 / 100) | int }}MB" + postgresql_max_wal_senders: "{{ (replica_count | default(0) | int) + 4 }}" + postgresql_max_replication_slots: "{{ (replica_count | default(0) | int) + 2 }}" + postgresql_replication_user: "{{ lookup('ansible.builtin.env', 'XINFRA_POSTGRES_REPLICATION_USER') | default('xinfra_replication', true) }}" + postgresql_replication_password: "{{ lookup('ansible.builtin.env', 'XINFRA_POSTGRES_REPLICATION_PASSWORD') }}" + postgresql_admin_password: "{{ lookup('ansible.builtin.env', 'XINFRA_POSTGRES_ADMIN_PASSWORD') }}" + postgresql_primary_host: "{{ hostvars[ansible_play_hosts_all[0]].ansible_host | default(ansible_play_hosts_all[0]) }}" + postgresql_primary_port: "{{ (postgresql_instances[ansible_play_hosts_all[0]].port | default(15432)) | int }}" + postgresql_replication_slot: "{{ postgresql_instance_plan.replication_slot | default('xinfra_' ~ postgresql_instance_id | replace('-', '_')) }}" + + pre_tasks: + - name: Validate PostgreSQL delivery parameters + ansible.builtin.assert: + that: + - postgresql_version_value in ['15', '16'] + - topology | default('standalone') in ['standalone', 'primary_replica'] + - ansible_distribution == 'Ubuntu' + - ansible_distribution_version is version('24.04', '>=') + - ansible_architecture in ['x86_64', 'aarch64'] + - postgresql_instance_id is match('^[a-z0-9][a-z0-9-]{0,62}$') + - (postgresql_port_value | int) >= 15432 and (postgresql_port_value | int) <= 15999 + - (memory_mb | default(4096) | int) >= 2048 + - (storage_gb | default(50) | int) >= 20 + - (postgresql_max_connections | int) >= 1 + - postgresql_admin_password | length >= 16 + - (topology | default('standalone') == 'standalone') or (postgresql_replication_password | length >= 16) + fail_msg: "PostgreSQL parameters are outside the supported target-state whitelist" + no_log: true + + - name: Check PostgreSQL port is free + ansible.builtin.shell: + cmd: | + if ss -lntH 'sport = :{{ postgresql_port_value }}' | grep -q .; then + systemctl is-active --quiet 'postgresql-xinfra@{{ postgresql_instance_id }}.service' + fi + register: postgresql_port_check + failed_when: postgresql_port_check.rc != 0 + changed_when: false + + - name: Check target memory and data-root capacity + ansible.builtin.assert: + that: + - (ansible_memtotal_mb | int) >= (memory_mb | default(4096) | int) + fail_msg: "Target host does not have enough memory or storage for the PostgreSQL instance" + + - name: Ensure PostgreSQL data root exists + ansible.builtin.file: + path: "{{ postgresql_data_root }}" + state: directory + owner: postgres + group: postgres + mode: '0750' + + - name: Check PostgreSQL data filesystem capacity + ansible.builtin.shell: + cmd: "df -P -B1 {{ postgresql_data_root | quote }} | awk 'NR==2 {print $4}'" + register: postgresql_root_available_bytes + changed_when: false + failed_when: >- + postgresql_root_available_bytes.rc != 0 or + (postgresql_root_available_bytes.stdout | trim | int) < + (storage_gb | default(50) | int) * 1073741824 + + - name: Install PostgreSQL packages + ansible.builtin.apt: + name: + - "postgresql-{{ postgresql_version_value }}" + - postgresql-client + state: present + update_cache: true + cache_valid_time: 3600 + + - name: Stop distribution-managed cluster + ansible.builtin.systemd_service: + name: "postgresql@{{ postgresql_version_value }}-main.service" + state: stopped + enabled: false + failed_when: false + + - name: Create isolated PostgreSQL directories + ansible.builtin.file: + path: "{{ item }}" + state: directory + owner: postgres + group: postgres + mode: '0700' + loop: + - "{{ postgresql_data_root }}/{{ postgresql_instance_id }}" + - "{{ postgresql_data_dir }}" + - "{{ postgresql_conf_dir }}" + - "{{ postgresql_log_dir }}" + + - name: Initialize primary or standalone data directory + ansible.builtin.command: + cmd: "/usr/lib/postgresql/{{ postgresql_version_value }}/bin/initdb -D {{ postgresql_data_dir }}" + creates: "{{ postgresql_data_dir }}/PG_VERSION" + become_user: postgres + when: postgresql_role != 'replica' + + - name: Render PostgreSQL configuration + ansible.builtin.template: + src: templates/postgresql-instance.conf.j2 + dest: "{{ postgresql_config_file }}" + owner: postgres + group: postgres + mode: '0600' + no_log: "{{ postgresql_role == 'replica' }}" + + - name: Render pg_hba.conf + ansible.builtin.copy: + dest: "{{ postgresql_hba_file }}" + owner: postgres + group: postgres + mode: '0600' + content: | + local all postgres peer + local all all peer + host all all 127.0.0.1/32 scram-sha-256 + {% for host in ansible_play_hosts_all %} + {% set host_address = hostvars[host].ansible_host | default(host) %} + host all all {{ host_address }}{% if host_address is match('^[0-9.]+$') %}/32{% endif %} scram-sha-256 + host replication {{ postgresql_replication_user }} {{ host_address }}{% if host_address is match('^[0-9.]+$') %}/32{% endif %} scram-sha-256 + {% endfor %} + + - name: Validate primary or standalone PostgreSQL configuration syntax + ansible.builtin.command: + argv: + - "/usr/lib/postgresql/{{ postgresql_version_value }}/bin/postgres" + - "-D" + - "{{ postgresql_data_dir }}" + - "-C" + - "port" + - "-c" + - "config_file={{ postgresql_config_file }}" + - "-c" + - "hba_file={{ postgresql_hba_file }}" + become_user: postgres + changed_when: false + when: postgresql_role != 'replica' + + - name: Install PostgreSQL systemd template + ansible.builtin.copy: + src: files/postgresql-xinfra@.service + dest: /etc/systemd/system/postgresql-xinfra@.service + owner: root + group: root + mode: '0644' + register: postgresql_unit + + - name: Configure runtime unit environment + ansible.builtin.file: + dest: "/etc/systemd/system/postgresql-xinfra@{{ postgresql_instance_id }}.service.d" + state: directory + mode: '0755' + + - name: Write runtime unit environment + ansible.builtin.copy: + dest: "/etc/systemd/system/postgresql-xinfra@{{ postgresql_instance_id }}.service.d/version.conf" + mode: '0644' + content: | + [Service] + Environment=POSTGRESQL_VERSION={{ postgresql_version_value }} + Environment=POSTGRESQL_DATA_DIR={{ postgresql_data_dir }} + Environment=POSTGRESQL_CONFIG_FILE={{ postgresql_config_file }} + Environment=POSTGRESQL_HBA_FILE={{ postgresql_hba_file }} + + - name: Reload systemd + ansible.builtin.systemd_service: + daemon_reload: true + + - name: Start PostgreSQL delivery instance + ansible.builtin.systemd_service: + name: "postgresql-xinfra@{{ postgresql_instance_id }}.service" + state: started + enabled: true + when: postgresql_role != 'replica' + + - name: Wait for primary or standalone PostgreSQL TCP port + ansible.builtin.wait_for: + host: "{{ ansible_host | default(inventory_hostname) }}" + port: "{{ postgresql_port_value }}" + timeout: 60 + when: postgresql_role != 'replica' + + - name: Configure PostgreSQL administrator and replication roles + ansible.builtin.shell: + cmd: | + /usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -v ON_ERROR_STOP=1 <<'SQL' + ALTER ROLE postgres PASSWORD '{{ postgresql_admin_password | replace("'", "''") }}'; + {% if topology | default('standalone') == 'primary_replica' %} + DO $do$BEGIN + IF NOT EXISTS (SELECT FROM pg_roles WHERE rolname = '{{ postgresql_replication_user }}') THEN + CREATE ROLE {{ postgresql_replication_user }} WITH REPLICATION LOGIN PASSWORD '{{ postgresql_replication_password | replace("'", "''") }}'; + ELSE + ALTER ROLE {{ postgresql_replication_user }} WITH REPLICATION LOGIN PASSWORD '{{ postgresql_replication_password | replace("'", "''") }}'; + END IF; + END$do$; + {% endif %} + SQL + become_user: postgres + no_log: true + when: postgresql_role in ['primary', 'standalone'] + + - name: Create physical replication slots on primary + ansible.builtin.shell: + cmd: | + /usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -v ON_ERROR_STOP=1 <<'SQL' + {% for replica_index in range(1, ansible_play_hosts_all | length) %} + SELECT pg_create_physical_replication_slot('xinfra_{{ cluster_name | replace("-", "_") }}_replica_{{ replica_index }}') + WHERE NOT EXISTS (SELECT FROM pg_replication_slots WHERE slot_name = 'xinfra_{{ cluster_name | replace("-", "_") }}_replica_{{ replica_index }}'); + {% endfor %} + SQL + become_user: postgres + changed_when: false + when: postgresql_role == 'primary' and topology | default('standalone') == 'primary_replica' + + - name: Initialize replica from primary + ansible.builtin.shell: + cmd: >- + PGPASSWORD={{ postgresql_replication_password | quote }} + /usr/lib/postgresql/{{ postgresql_version_value }}/bin/pg_basebackup + -h {{ postgresql_primary_host }} -p {{ postgresql_primary_port }} + -U {{ postgresql_replication_user }} -D {{ postgresql_data_dir }} + -Fp -Xs -P -R -S {{ postgresql_replication_slot }} + args: + creates: "{{ postgresql_data_dir }}/PG_VERSION" + become_user: postgres + no_log: true + when: postgresql_role == 'replica' + + - name: Validate replica PostgreSQL configuration syntax + ansible.builtin.command: + argv: + - "/usr/lib/postgresql/{{ postgresql_version_value }}/bin/postgres" + - "-D" + - "{{ postgresql_data_dir }}" + - "-C" + - "port" + - "-c" + - "config_file={{ postgresql_config_file }}" + - "-c" + - "hba_file={{ postgresql_hba_file }}" + become_user: postgres + changed_when: false + no_log: true + when: postgresql_role == 'replica' + + - name: Start PostgreSQL replica instance + ansible.builtin.systemd_service: + name: "postgresql-xinfra@{{ postgresql_instance_id }}.service" + state: started + enabled: true + when: postgresql_role == 'replica' + + - name: Wait for PostgreSQL replica TCP port + ansible.builtin.wait_for: + host: "{{ ansible_host | default(inventory_hostname) }}" + port: "{{ postgresql_port_value }}" + timeout: 60 + when: postgresql_role == 'replica' + + - name: Verify PostgreSQL readiness + ansible.builtin.command: + cmd: "/usr/lib/postgresql/{{ postgresql_version_value }}/bin/pg_isready -h 127.0.0.1 -p {{ postgresql_port_value }}" + changed_when: false + + - name: Verify primary or standalone role + ansible.builtin.shell: + cmd: >- + test "$(/usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -Atc 'SELECT pg_is_in_recovery()')" = "f" + become_user: postgres + changed_when: false + when: postgresql_role in ['primary', 'standalone'] + + - name: Verify replica recovery role + ansible.builtin.shell: + cmd: >- + test "$(/usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -Atc 'SELECT pg_is_in_recovery()')" = "t" + become_user: postgres + changed_when: false + when: postgresql_role == 'replica' + + - name: Verify primary streaming replicas and replication slots + ansible.builtin.shell: + cmd: | + streaming=$(/usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -Atc "SELECT count(*) FROM pg_stat_replication WHERE state = 'streaming'") + slots=$(/usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -Atc "SELECT count(*) FROM pg_replication_slots WHERE slot_type = 'physical'") + test "$streaming" -ge "{{ replica_count | default(0) | int }}" + test "$slots" -ge "{{ replica_count | default(0) | int }}" + become_user: postgres + register: postgresql_replication_health + until: postgresql_replication_health.rc == 0 + retries: 30 + delay: 2 + changed_when: false + when: postgresql_role == 'primary' and topology | default('standalone') == 'primary_replica' + + - name: Verify delivered version, port and data directory + ansible.builtin.shell: + cmd: | + set -euo pipefail + actual_version=$(/usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -Atc "SHOW server_version") + actual_port=$(/usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -Atc "SHOW port") + actual_data=$(/usr/bin/psql -h {{ postgresql_run_dir }} -p {{ postgresql_port_value }} -U postgres -d postgres -Atc "SHOW data_directory") + test "${actual_version%%.*}" = "{{ postgresql_version_value }}" + test "$actual_port" = "{{ postgresql_port_value }}" + test "$actual_data" = "{{ postgresql_data_dir }}" + executable: /bin/bash + become_user: postgres + changed_when: false diff --git a/ansible/templates/postgresql-instance.conf.j2 b/ansible/templates/postgresql-instance.conf.j2 new file mode 100644 index 0000000..c1b3418 --- /dev/null +++ b/ansible/templates/postgresql-instance.conf.j2 @@ -0,0 +1,28 @@ +# Managed by XINFRA PostgreSQL delivery. instance={{ postgresql_instance_id }} +listen_addresses = '*' +port = {{ postgresql_port_value }} +data_directory = '{{ postgresql_data_dir }}' +hba_file = '{{ postgresql_hba_file }}' +unix_socket_directories = '{{ postgresql_run_dir }}' +external_pid_file = '{{ postgresql_run_dir }}/postmaster.pid' +logging_collector = on +log_directory = '{{ postgresql_log_dir }}' +log_filename = 'postgresql-%Y-%m-%d_%H%M%S.log' +log_min_messages = warning +max_connections = {{ postgresql_max_connections }} +shared_buffers = '{{ postgresql_shared_buffers }}' +effective_cache_size = '{{ postgresql_effective_cache_size }}' +work_mem = '4MB' +maintenance_work_mem = '64MB' +wal_level = replica +max_wal_senders = {{ postgresql_max_wal_senders }} +max_replication_slots = {{ postgresql_max_replication_slots }} +wal_keep_size = '256MB' +hot_standby = on +synchronous_standby_names = '' +{% if postgresql_role == 'primary' %} +primary_conninfo = '' +{% elif postgresql_role == 'replica' %} +primary_conninfo = 'host={{ postgresql_primary_host }} port={{ postgresql_primary_port }} user={{ postgresql_replication_user }} password={{ postgresql_replication_password | replace("'", "''") }} application_name={{ postgresql_instance_id }}' +primary_slot_name = '{{ postgresql_replication_slot }}' +{% endif %} diff --git a/config/local/README.md b/config/local/README.md new file mode 100644 index 0000000..0ae7807 --- /dev/null +++ b/config/local/README.md @@ -0,0 +1,49 @@ +# Local PostgreSQL AWX setup + +This directory keeps the local configuration used to expose the PostgreSQL +delivery playbook through AWX. Files ending in `.env` are ignored; their +`.env.example` counterparts document the supported parameters. + +The local object mapping is: + +| AWX object | Local value | +| --- | --- | +| Project | `XINFRA PostgreSQL Project` | +| Project source | local `ansible/`, committed into AWX task pod local Git | +| Inventory | `XINFRA PostgreSQL Inventory` | +| Inventory host | `postgresql-218-11-5-224` -> public `218.11.5.224`, SSH via `192.168.1.5` | +| Machine credential | `XINFRA PostgreSQL SSH`, user `root` | +| Private key | `/Users/lx/xengineer-cs2.pem` | +| Secret credential | `XINFRA PostgreSQL Runtime Secrets` | +| Job Template | `XINFRA PostgreSQL Delivery` | +| Playbook | `postgresql-deploy.yml` | +| Instance Group | `xinfra-e2e-execution` | + +Run the idempotent setup whenever the AWX task pod restarts or the local +playbook changes. This AWX installation has project persistence disabled and +cannot use a Manual Project, so the Project is Git-backed with launch updates +disabled. The script refreshes a local Git source inside the AWX task pod and +runs a project update so AWX rebuilds its playbook index. + +```bash +./scripts/setup-postgresql-awx.sh +``` + +Start the backend with AWX delivery scheduling enabled: + +```bash +./scripts/run-server-postgresql-local.sh +``` + +Then start the frontend and open `/service/catalog/pgsql`. Select +`XINFRA PostgreSQL Delivery` as the target. The inventory currently contains +one host, so only `standalone` can be deployed until more distinct inventory +hosts are added. + +The frontend dev server listens on `http://127.0.0.1:5173`. Port `3000` is +reserved for the local `authserver` MySQL container and must not be used by +Vite. Set `XINFRA_BACKEND_PROXY_TARGET` only when the backend is listening on a +non-default address. + +The AWX HTTP API is `http://127.0.0.1:59123`. Local port `59121` is Minikube +SSH and must not be used as `AWX_BASE_URL`. diff --git a/config/local/postgresql-awx.env.example b/config/local/postgresql-awx.env.example new file mode 100644 index 0000000..84cee5c --- /dev/null +++ b/config/local/postgresql-awx.env.example @@ -0,0 +1,43 @@ +# Local AWX API. Port 59121 is Minikube SSH, not the AWX HTTP API. +AWX_BASE_URL=http://127.0.0.1:59123 +AWX_TOKEN= +AWX_USERNAME=admin +AWX_PASSWORD= + +# When AWX_PASSWORD is empty, the setup script reads this Kubernetes Secret. +AWX_K8S_NAMESPACE=awx +AWX_K8S_ADMIN_PASSWORD_SECRET=awx-demo-admin-password +AWX_K8S_TASK_APP=awx-demo-task +AWX_K8S_TASK_POD= +AWX_K8S_TASK_CONTAINER=awx-demo-task + +AWX_ORGANIZATION=Default +AWX_PROJECT_NAME="XINFRA PostgreSQL Project" +AWX_PROJECT_DESCRIPTION="Local-staged XINFRA PostgreSQL Ansible project" +AWX_PROJECT_SCM_URL=file:///var/lib/awx/projects/_xinfra_postgresql_source +AWX_PROJECT_SOURCE_PATH=_xinfra_postgresql_source +AWX_PROJECT_SOURCE_DIR=ansible + +AWX_INVENTORY_NAME="XINFRA PostgreSQL Inventory" +AWX_INVENTORY_DESCRIPTION="PostgreSQL host pool for local delivery testing" +AWX_HOST_NAME=postgresql-218-11-5-224 +AWX_HOST_ADDRESS=218.11.5.224 +AWX_HOST_ANSIBLE_ADDRESS=192.168.1.5 +AWX_HOST_SSH_PORT=22 +AWX_HOST_SSH_USER=root + +AWX_MACHINE_CREDENTIAL_NAME="XINFRA PostgreSQL SSH" +AWX_SSH_PRIVATE_KEY_FILE=/Users/lx/xengineer-cs2.pem + +AWX_POSTGRES_CREDENTIAL_TYPE_NAME="XINFRA PostgreSQL Secrets" +AWX_POSTGRES_CREDENTIAL_NAME="XINFRA PostgreSQL Runtime Secrets" +# Leave both passwords empty on first setup to generate them inside AWX. +# To rotate them later, set both values together before running setup again. +XINFRA_POSTGRES_ADMIN_PASSWORD= +XINFRA_POSTGRES_REPLICATION_USER=xinfra_replication +XINFRA_POSTGRES_REPLICATION_PASSWORD= + +AWX_JOB_TEMPLATE_NAME="XINFRA PostgreSQL Delivery" +AWX_JOB_TEMPLATE_DESCRIPTION="PostgreSQL host_pool delivery template for XINFRA" +AWX_PLAYBOOK=postgresql-deploy.yml +AWX_INSTANCE_GROUP=xinfra-e2e-execution diff --git a/config/local/server-postgresql.env.example b/config/local/server-postgresql.env.example new file mode 100644 index 0000000..8224347 --- /dev/null +++ b/config/local/server-postgresql.env.example @@ -0,0 +1,16 @@ +# Loaded by scripts/run-server-postgresql-local.sh before server/.env. +AWX_BASE_URL=http://127.0.0.1:59123 +AWX_TOKEN= +AWX_USERNAME=admin +AWX_PASSWORD= + +AWX_K8S_NAMESPACE=awx +AWX_K8S_ADMIN_PASSWORD_SECRET=awx-demo-admin-password + +DELIVERY_SCHEDULER_ENABLED=true +DELIVERY_POLL_SECONDS=3 +DELIVERY_RESERVATION_TTL_MINUTES=120 +DELIVERY_GLOBAL_LIMIT=2 +DELIVERY_TARGET_LIMIT=2 +DELIVERY_BUSINESS_LIMIT=1 +DELIVERY_DATA_DISKS=/data diff --git a/frontend/src/api/delivery.ts b/frontend/src/api/delivery.ts index f1e8a02..5698b03 100644 --- a/frontend/src/api/delivery.ts +++ b/frontend/src/api/delivery.ts @@ -5,6 +5,7 @@ export interface DeliveryTarget { id: number name: string target_type: string + service_type?: string awx_inventory_id: number awx_template_id: number enabled: boolean @@ -48,18 +49,36 @@ export interface CreateMySQLDeliveryPayload { mysql_root_password?: string } +export interface CreatePostgreSQLDeliveryPayload { + business_line_id: number + target_id: number + namespace: string + cluster_name: string + version_major: string + topology: 'standalone' | 'primary_replica' + replica_count?: number + target_hosts?: string[] + cpu_milli: number + memory_mi: number + storage_gi: number + data_root?: string + max_connections?: number +} + export interface DeliveryTask { id: string business_line_id: number requested_by: number component: string target_type: string + service_type?: string target_id: number namespace: string instance_name: string target_host?: string target_host_ip?: string mysql_port?: number + postgresql_port?: number status: string error_message?: string created_at: string @@ -109,8 +128,8 @@ export interface MySQLServiceLedgerItem { } export const deliveryApi = { - async listTargets(): Promise { - const data = await authRequest('/auth/api/v1/delivery/targets?component=mysql') + async listTargets(component = 'mysql'): Promise { + const data = await authRequest(`/auth/api/v1/delivery/targets?component=${encodeURIComponent(component)}`) return Array.isArray(data.items) ? data.items : [] }, @@ -156,6 +175,17 @@ export const deliveryApi = { }) }, + async createPostgreSQL(payload: CreatePostgreSQLDeliveryPayload): Promise { + const data = await authRequest('/auth/api/v1/delivery/postgresql', { + method: 'POST', + headers: { + 'Idempotency-Key': createIdempotencyKey({ ...payload, service: 'postgresql' }), + }, + body: JSON.stringify(payload), + }) + return data.task + }, + async getTask(taskId: string): Promise<{ task: DeliveryTask; events: TaskEvent[] }> { const data = await authRequest(`/auth/api/v1/delivery/tasks/${encodeURIComponent(taskId)}`) return { @@ -204,13 +234,13 @@ export const deliveryApi = { }, } -function createIdempotencyKey(payload: CreateMySQLDeliveryPayload) { +function createIdempotencyKey(payload: { business_line_id: number; target_id: number; namespace: string; instance_name?: string; cluster_name?: string; service?: string }) { return [ - 'mysql', + payload.service || 'mysql', payload.business_line_id, payload.target_id, payload.namespace, - payload.instance_name, + payload.instance_name || payload.cluster_name, Date.now(), ].join(':') } diff --git a/frontend/src/views/service/Catalog.vue b/frontend/src/views/service/Catalog.vue index ebc6b1f..fe065b5 100644 --- a/frontend/src/views/service/Catalog.vue +++ b/frontend/src/views/service/Catalog.vue @@ -82,32 +82,46 @@
-
+ +
数据库管理员密码(root)
{{ mode.label }} @@ -175,8 +189,10 @@
-
-
安装目录{{ mysqlInstallPath }}
-
数据目录{{ mysqlDataPath }}
- 程序安装在系统盘,数据盘只决定实例数据、日志和临时文件的落点。 +
安装目录{{ installPath }}
+
数据目录{{ dataPath }}
+ {{ isPostgreSQL ? '端口由平台从 15432–15999 自动分配。' : '程序安装在系统盘,数据盘只决定实例数据、日志和临时文件的落点。' }}
-
+
数据库配置可选,默认值已按标准基线填充 @@ -254,39 +270,39 @@
-
-
本地 Socket{{ mysqlSocketPath }}
-
配置文件{{ mysqlConfigPath }}
-
systemd 服务{{ mysqlServiceName }}
+
本地 Socket{{ databaseSocketPath }}
+
配置文件{{ databaseConfigPath }}
+
systemd 服务{{ databaseServiceName }}
注册系统{{ activeService.registerTo }}
@@ -618,6 +634,8 @@ interface WorkbenchSnapshot { timezone: string dataDisk: string targetHost: string + targetHosts: string[] + replicaCount: number maxConnections: string redoCapacity: string flushLogAtCommit: number @@ -694,30 +712,30 @@ const basicServices = ref([ }, { key: 'pgsql', - name: 'PgSQL', + name: 'PostgreSQL', icon: 'Pg', bgColor: 'var(--logo-ap-bg)', iconColor: 'var(--tag-purple-text)', - description: '流复制主从架构,规划中 · 接入 CloudDM 统一审核体系。', - playbook: 'roles/pgsql-deploy', - version: '规划中', - status: '规划中', - template: '', - runner: '', - versions: [], - modes: [], - specs: [], - disks: [], - charsets: [], - defaultPort: 5432, - defaultPaths: { install: '', data: '', log: '' }, - registerTo: '', - configId: '', + description: '原生 PostgreSQL 单机或一主多从交付,使用独立数据目录与 systemd 实例。', + playbook: 'postgresql-deploy.yml', + version: 'v15 / v16', + status: '可交付', + template: 'XINFRA PostgreSQL Delivery', + runner: 'xinfra-e2e-execution', + versions: ['PostgreSQL 16', 'PostgreSQL 15'], + modes: [ + { value: 'standalone', label: '单实例' }, + { value: 'primary_replica', label: '一主多从' }, + ], + specs: ['0.5C / 2G', '1C / 2G', '2C / 4G', '4C / 8G', '8C / 16G', '16C / 32G'], + disks: ['20 GB', '50 GB', '100 GB', '500 GB', '1000 GB'], + charsets: ['UTF8'], + defaultPort: 0, + defaultPaths: { install: '/usr/lib/postgresql', data: '/data/postgresql', log: '/data/postgresql' }, + registerTo: 'PostgreSQL 专用资源台账', + configId: 'PG-NATIVE-V1', assetId: '', - healthText: '', - statusColor: 'var(--text-dim)', - disabled: true, - opacity: 0.75, + healthText: 'PostgreSQL 就绪 · TCP 探测通过', }, { key: 'redis', @@ -817,7 +835,7 @@ const rollbackReleasing = ref(false) const cloudDMRetrying = ref(false) const rollbackActionBusy = computed(() => rollbackRetrying.value || rollbackReleasing.value) const canManageRollback = computed(() => authStore.isAdmin && lastDeliveryStatus.value === 'rollback_failed' && Boolean(deploymentId.value)) -const canRetryCloudDM = computed(() => deliveryRegisterFailed.value && Boolean(deploymentId.value)) +const canRetryCloudDM = computed(() => !isPostgreSQL.value && deliveryRegisterFailed.value && Boolean(deploymentId.value)) const deliveryForm = reactive({ instanceName: generateInstanceName(activeServiceKey.value, currentName.value), @@ -833,6 +851,8 @@ const deliveryForm = reactive({ timezone: '+08:00', dataDisk: '/data', targetHost: '', + targetHosts: [] as string[], + replicaCount: 1, maxConnections: 'auto', redoCapacity: 'auto', flushLogAtCommit: 1, @@ -844,6 +864,7 @@ const deliveryForm = reactive({ }) const maxConnectionOptions = ['auto', '200', '500', '1000', '2000', '4000', '8000', '16000'] +const postgresqlMaxConnectionOptions = ['auto', '200', '500', '1000', '2000', '4000', '8000'] const logSizeOptions = ['128M', '256M', '512M', '1G'] const ioCapacityOptions = [200, 2000, 5000] const longQueryOptions = [0.5, 1, 2, 5, 10] @@ -864,24 +885,40 @@ const deliveryStageStepIndex: Record = { } const activeService = computed(() => basicServices.value.find((service) => service.key === activeServiceKey.value && !service.disabled)) +const isPostgreSQL = computed(() => activeServiceKey.value === 'pgsql') const selectedTarget = computed(() => deliveryTargets.value.find((target) => target.id === selectedTargetId.value)) const targetHosts = computed(() => parseTargetHosts(selectedTarget.value?.metadata)) const selectedHost = computed(() => targetHosts.value.find((host) => host.name === deliveryForm.targetHost)) +const activeMaxConnectionOptions = computed(() => isPostgreSQL.value ? postgresqlMaxConnectionOptions : maxConnectionOptions) const collationOptions = computed(() => collationMap[deliveryForm.charset] || []) const storageLabel = computed(() => `${deliveryForm.storageGi} GB`) const instanceNameError = computed(() => instanceNameValidationError(deliveryForm.instanceName)) const dataDiskError = computed(() => dataDiskValidationError(deliveryForm.dataDisk)) -const rootPasswordError = computed(() => rootPasswordValidationError(deliveryForm.rootPassword)) -const mysqlInstallPath = computed(() => instanceNameError.value ? '请先填写有效实例名' : `/opt/mysql-delivery/${deliveryForm.instanceName}`) -const mysqlDataPath = computed(() => { +const rootPasswordError = computed(() => isPostgreSQL.value ? '' : rootPasswordValidationError(deliveryForm.rootPassword)) +const installPath = computed(() => { + if (instanceNameError.value) return '请先填写有效实例名' + return isPostgreSQL.value ? `/usr/lib/postgresql/${postgresqlVersionValue(deliveryForm.version)}` : `/opt/mysql-delivery/${deliveryForm.instanceName}` +}) +const dataPath = computed(() => { if (instanceNameError.value) return '请先填写有效实例名' if (dataDiskError.value) return '请先选择有效挂载点' + if (isPostgreSQL.value) return `${deliveryForm.dataDisk}/${deliveryForm.instanceName}/data` return `${deliveryForm.dataDisk}/mysql-delivery/${deliveryForm.instanceName}/data` }) const taskNo = computed(() => `CMP-20260721-${activeServiceKey.value === 'mysql' ? '0024' : '0023'}`) const currentModeLabel = computed(() => activeService.value?.modes.find((mode) => mode.value === deliveryForm.mode)?.label || '-') const topologySummary = computed(() => `${currentModeLabel.value} · ${deliveryForm.spec} · ${storageLabel.value}`) const topologyNodes = computed(() => { + if (isPostgreSQL.value) { + const count = deliveryForm.mode === 'primary_replica' ? deliveryForm.replicaCount + 1 : 1 + const pinned = deliveryForm.mode === 'primary_replica' ? deliveryForm.targetHosts : deliveryForm.targetHost ? [deliveryForm.targetHost] : [] + const candidates = (pinned.length ? pinned : targetHosts.value.slice(0, count).map((host) => host.name)) + .map((name) => targetHosts.value.find((host) => host.name === name) || { name }) + return Array.from({ length: count }, (_, index) => { + const host = candidates[index] + return { name: host?.name || '等待调度', ip: host?.ip || '由 AWX Inventory 分配', role: index === 0 ? (count === 1 ? '单实例' : '主节点') : `从节点 ${index}` } + }) + } if (deliveryForm.mode === 'single') { const previewHost = selectedHost.value || targetHosts.value[0] return [{ name: previewHost?.name || '等待调度', ip: previewHost?.ip || '由 AWX Inventory 分配', role: '单实例' }] @@ -912,37 +949,53 @@ const runnerPreview = computed(() => { `component: ${activeService.value?.key || '-'}`, `instance: ${deliveryForm.instanceName}`, `topology: ${currentModeLabel.value}`, - `limit: ${deliveryForm.targetHost || 'auto'}`, - `data_disk: ${deliveryForm.dataDisk}`, + `limit: ${selectedTargetHostText.value}`, + `${isPostgreSQL.value ? 'data_root' : 'data_disk'}: ${deliveryForm.dataDisk}`, `resources: ${deliveryForm.spec} / ${storageLabel.value}`, - `charset: ${deliveryForm.charset}`, - `collation: ${deliveryForm.collation}`, - `timezone: ${deliveryForm.timezone}`, + ...(isPostgreSQL.value ? [`replica_count: ${deliveryForm.mode === 'primary_replica' ? deliveryForm.replicaCount : 0}`] : [ + `charset: ${deliveryForm.charset}`, + `collation: ${deliveryForm.collation}`, + `timezone: ${deliveryForm.timezone}`, + ]), `register_to: ${activeService.value?.registerTo || '-'}`, ].join('\n') }) const connectionReady = computed(() => Boolean(deliveredHost.value && deliveredPort.value)) const credentialEligible = computed(() => { - if (!connectionReady.value || !['finished', 'register_failed'].includes(lastDeliveryStatus.value)) return false + if (isPostgreSQL.value || !connectionReady.value || !['finished', 'register_failed'].includes(lastDeliveryStatus.value)) return false return credentialAvailable.value !== false }) const revealedRootCredential = computed(() => revealedCredentials.value.find((item) => item.username === 'root')) const resultReady = computed(() => !taskStateLoading.value && !running.value && isTerminalDeliveryStatus(lastDeliveryStatus.value)) +const selectedTargetHostText = computed(() => { + if (isPostgreSQL.value && deliveryForm.mode === 'primary_replica') { + return deliveryForm.targetHosts.length ? deliveryForm.targetHosts.join(', ') : '由 AWX Inventory 分配' + } + return deliveryForm.targetHost || '由 AWX Inventory 分配' +}) const previewAddress = computed(() => `${selectedHost.value?.ip || 'AWX 自动分配'}:${deliveryForm.port || '自动端口'}`) const resultAddress = computed(() => connectionReady.value ? `${deliveredHost.value}:${deliveredPort.value}` : '交付接口未返回地址') -const connectionUri = computed(() => connectionReady.value ? `mysql://root@${deliveredHost.value}:${deliveredPort.value}` : '') -const mysqlCommand = computed(() => connectionReady.value ? `mysql -h ${deliveredHost.value} -P ${deliveredPort.value} -u root -p` : '等待交付接口返回主机和端口') -const mysqlSocketPath = computed(() => `/run/mysql-delivery-${deliveryForm.instanceName}/mysql.sock`) -const mysqlConfigPath = computed(() => `/etc/mysql/mysql-delivery/${deliveryForm.instanceName}.cnf`) -const mysqlServiceName = computed(() => `mysql-delivery@${deliveryForm.instanceName}.service`) +const connectionUri = computed(() => { + if (!connectionReady.value) return '' + return isPostgreSQL.value ? `postgresql://postgres@${deliveredHost.value}:${deliveredPort.value}/postgres` : `mysql://root@${deliveredHost.value}:${deliveredPort.value}` +}) +const databaseCommand = computed(() => { + if (!connectionReady.value) return '等待交付接口返回主机和端口' + return isPostgreSQL.value + ? `psql -h ${deliveredHost.value} -p ${deliveredPort.value} -U postgres -d postgres` + : `mysql -h ${deliveredHost.value} -P ${deliveredPort.value} -u root -p` +}) +const databaseSocketPath = computed(() => isPostgreSQL.value ? `/run/postgresql-xinfra/${deliveryForm.instanceName}/.s.PGSQL.${deliveredPort.value || ''}` : `/run/mysql-delivery-${deliveryForm.instanceName}/mysql.sock`) +const databaseConfigPath = computed(() => isPostgreSQL.value ? `${deliveryForm.dataDisk}/${deliveryForm.instanceName}/conf/postgresql.conf` : `/etc/mysql/mysql-delivery/${deliveryForm.instanceName}.cnf`) +const databaseServiceName = computed(() => isPostgreSQL.value ? `postgresql-xinfra@${deliveryForm.instanceName}.service` : `mysql-delivery@${deliveryForm.instanceName}.service`) const resultTitle = computed(() => { if (deliveryRolledBack.value) return '交付失败,已完成回退' if (deliveryAcknowledged.value) return '交付失败,已确认清理' - if (deliveryRegisterFailed.value) return 'MySQL 已交付,CloudDM 注册失败' + if (deliveryRegisterFailed.value) return `${activeService.value?.name || '数据库'} 已交付,CloudDM 注册失败` if (lastDeliveryStatus.value === 'rollback_failed') return '交付失败,回退失败' if (lastDeliveryStatus.value === 'canceled') return '交付任务已取消' if (deliveryFailed.value) return '交付失败' - return activeServiceKey.value === 'mysql' ? 'MySQL 实例已交付' : 'OpenResty 集群已交付' + return isPostgreSQL.value ? 'PostgreSQL 集群已交付' : 'MySQL 实例已交付' }) const resultSubtitle = computed(() => { if (deliveryRolledBack.value) return deliveryError.value || '目标机实例文件、配置和资源记录已清理' @@ -982,6 +1035,7 @@ watch(currentName, (name) => { watch(selectedTargetId, () => { deliveryForm.targetHost = '' + deliveryForm.targetHosts = [] precheckPassed.value = false }) @@ -1001,6 +1055,7 @@ watch(deliveryForm, () => { watch(activeServiceKey, () => { hydrateServiceDefaults() + loadDeliveryTargets() if (activeServiceKey.value) { workbenchReady = false restoreWorkbench() @@ -1096,6 +1151,8 @@ function persistWorkbench() { timezone: deliveryForm.timezone, dataDisk: deliveryForm.dataDisk, targetHost: deliveryForm.targetHost, + targetHosts: deliveryForm.targetHosts, + replicaCount: deliveryForm.replicaCount, maxConnections: deliveryForm.maxConnections, redoCapacity: deliveryForm.redoCapacity, flushLogAtCommit: deliveryForm.flushLogAtCommit, @@ -1153,9 +1210,9 @@ async function restoreWorkbench() { try { const data = await deliveryApi.getTask(snapshot.deploymentId) deliveredHost.value = data.task.target_host_ip || deliveredHost.value - deliveredPort.value = data.task.mysql_port || deliveredPort.value + deliveredPort.value = (isPostgreSQL.value ? data.task.postgresql_port : data.task.mysql_port) || deliveredPort.value lastDeliveryStatus.value = data.task.status - credentialAvailable.value = data.task.credential_available + credentialAvailable.value = isPostgreSQL.value ? false : data.task.credential_available deliveryError.value = data.task.error_message || deliveryError.value seenEventIds.value = new Set() applyTaskEvents(data.events, false) @@ -1188,7 +1245,7 @@ function hydrateServiceDefaults() { if (!service) return deliveryForm.instanceName = generateInstanceName(service.key, currentName.value) deliveryForm.version = service.versions[0] || '' - deliveryForm.mode = service.modes[1]?.value || service.modes[0]?.value || 'single' + deliveryForm.mode = service.modes[0]?.value || 'single' deliveryForm.spec = service.key === 'mysql' ? (service.specs[0] || '') : (service.specs[1] || service.specs[0] || '') @@ -1202,6 +1259,8 @@ function hydrateServiceDefaults() { deliveryForm.timezone = '+08:00' deliveryForm.dataDisk = service.defaultPaths.data || '/data' deliveryForm.targetHost = '' + deliveryForm.targetHosts = [] + deliveryForm.replicaCount = 1 deliveryForm.rootPassword = generateRootPassword() deliveryForm.maxConnections = 'auto' deliveryForm.redoCapacity = 'auto' @@ -1252,6 +1311,16 @@ function resetExecutionState() { } function defaultSteps(): DeliveryStep[] { + if (isPostgreSQL.value) { + return [ + ['资源锁定与主机检查', '配额、主机、端口和实例目录'], + ['安装 PostgreSQL 与初始化集群', '软件包、数据目录、运行用户'], + ['应用实例配置', 'postgresql.conf、pg_hba.conf、systemd'], + ['配置复制拓扑', '复制账号、slot、primary_conninfo'], + ['数据库健康检查', '端口、角色、版本和数据目录'], + ['资源入账与交付归档', 'PostgreSQL 集群、实例与资源台账'], + ].map(([name, desc]) => ({ name, desc, state: 'pending' as StepState })) + } const isMysql = activeServiceKey.value === 'mysql' const copy = isMysql ? [ @@ -1283,10 +1352,21 @@ async function precheck() { ElMessage.warning('当前没有可用的 AWX 部署模板,请联系平台管理员配置') return } - if (deliveryForm.port !== undefined && (deliveryForm.port < 13306 || deliveryForm.port > 13999)) { + if (!isPostgreSQL.value && deliveryForm.port !== undefined && (deliveryForm.port < 13306 || deliveryForm.port > 13999)) { ElMessage.warning('手动端口必须在 13306–13999 范围内') return } + if (isPostgreSQL.value && deliveryForm.mode === 'primary_replica') { + const requiredHosts = deliveryForm.replicaCount + 1 + if (targetHosts.value.length < requiredHosts) { + ElMessage.warning(`一主 ${deliveryForm.replicaCount} 从需要 ${requiredHosts} 台不同主机,当前 Inventory 只有 ${targetHosts.value.length} 台`) + return + } + if (deliveryForm.targetHosts.length > 0 && deliveryForm.targetHosts.length !== requiredHosts) { + ElMessage.warning(`请按主节点、从节点顺序选择 ${requiredHosts} 台主机,或清空后由平台自动分配`) + return + } + } try { const conflict = await findInstanceNameConflict(businessLineId, normalizeDNSLabel(deliveryForm.instanceName)) if (conflict) { @@ -1325,9 +1405,10 @@ async function findInstanceNameConflict(businessLineId: number, instanceName: st 'rolling_back', 'rollback_failed', ]) + const component = isPostgreSQL.value ? 'postgresql' : 'mysql' const [tasks, services] = await Promise.all([ - deliveryApi.listTasks({ businessLineId, component: 'mysql' }), - deliveryApi.listMySQLServices(businessLineId), + deliveryApi.listTasks({ businessLineId, component }), + isPostgreSQL.value ? Promise.resolve([]) : deliveryApi.listMySQLServices(businessLineId), ]) if (tasks.some((task) => task.instance_name === instanceName && occupiedStatuses.has(task.status))) { return `实例名称 ${instanceName} 已被运行中或尚未释放的交付任务占用` @@ -1350,7 +1431,7 @@ async function createTask() { return } if (!activeService.value) return - if (activeService.value.key !== 'mysql') { + if (!['mysql', 'pgsql'].includes(activeService.value.key)) { ElMessage.warning(`${activeService.value.name} 的 AWX 交付模板尚未接入`) return } @@ -1364,9 +1445,11 @@ async function createTask() { activeView.value = 'execution' steps.value = defaultSteps().map((step) => ({ ...step, state: 'pending' })) try { - const task = await deliveryApi.createMySQL(mysqlDeliveryPayload(businessLineId)) + const task = isPostgreSQL.value + ? await deliveryApi.createPostgreSQL(postgresqlDeliveryPayload(businessLineId)) + : await deliveryApi.createMySQL(mysqlDeliveryPayload(businessLineId)) deploymentId.value = task.id - credentialAvailable.value = task.credential_available + credentialAvailable.value = isPostgreSQL.value ? false : task.credential_available deliveryForm.rootPassword = '' upsertDeliveryHistory(task.id, task.status, task.error_message || '', task.created_at) deliveryLog.value += `\n[task] ${task.id} created by ${currentName.value}` @@ -1430,6 +1513,28 @@ function mysqlDeliveryPayload(businessLineId: number) { return payload } +function postgresqlDeliveryPayload(businessLineId: number) { + const resources = parseSpec(deliveryForm.spec) + const selectedHosts = deliveryForm.mode === 'primary_replica' + ? deliveryForm.targetHosts + : deliveryForm.targetHost ? [deliveryForm.targetHost] : [] + return { + business_line_id: businessLineId, + target_id: selectedTargetId.value || 0, + namespace: normalizeDNSLabel(currentName.value), + cluster_name: normalizeDNSLabel(deliveryForm.instanceName), + version_major: postgresqlVersionValue(deliveryForm.version), + topology: deliveryForm.mode as 'standalone' | 'primary_replica', + replica_count: deliveryForm.mode === 'primary_replica' ? deliveryForm.replicaCount : 0, + target_hosts: selectedHosts.length ? selectedHosts : undefined, + cpu_milli: resources.cpuMilli, + memory_mi: resources.memoryMi, + storage_gi: deliveryForm.storageGi, + data_root: '/data/postgresql', + max_connections: deliveryForm.maxConnections === 'auto' ? undefined : Number(deliveryForm.maxConnections), + } +} + function parseTargetHosts(metadata?: string): TargetHostOption[] { if (!metadata) return [] try { @@ -1460,8 +1565,8 @@ async function pollTask(id: string) { try { const data = await deliveryApi.getTask(id) deliveredHost.value = data.task.target_host_ip || deliveredHost.value - deliveredPort.value = data.task.mysql_port || deliveredPort.value - credentialAvailable.value = data.task.credential_available + deliveredPort.value = (isPostgreSQL.value ? data.task.postgresql_port : data.task.mysql_port) || deliveredPort.value + credentialAvailable.value = isPostgreSQL.value ? false : data.task.credential_available upsertDeliveryHistory(id, data.task.status, data.task.error_message || '', data.task.created_at) applyTaskEvents(data.events) applyDeliveryStatus(data.task.status, data.task.error_message || '') @@ -1568,10 +1673,17 @@ function applyDeliveryStatus(status: string, message: string) { } async function loadDeliveryTargets() { + if (!activeService.value) { + deliveryTargets.value = [] + selectedTargetId.value = undefined + return + } targetsLoading.value = true try { - deliveryTargets.value = await deliveryApi.listTargets() - const preferred = deliveryTargets.value.find((target) => target.name === mysqlDeliveryTemplateName) + deliveryTargets.value = await deliveryApi.listTargets(isPostgreSQL.value ? 'postgresql' : activeServiceKey.value) + const preferred = isPostgreSQL.value + ? deliveryTargets.value.find((target) => !/callback|rollback/i.test(target.name)) + : deliveryTargets.value.find((target) => target.name === mysqlDeliveryTemplateName) selectedTargetId.value = preferred?.id || deliveryTargets.value[0]?.id } catch (error) { ElMessage.error(error instanceof Error ? error.message : '获取部署目标失败') @@ -1664,6 +1776,17 @@ function defaultMountPathOptions(): DeliveryMountPath[] { return ['/data', '/disk1', '/mnt', '/opt/mysql-delivery'].map((path) => ({ path, available_gi: 0 })) } +function postgresqlVersionValue(version: string) { + return version.match(/\d+/)?.[0] || '16' +} + +function selectDeliveryMode(mode: string) { + deliveryForm.mode = mode + deliveryForm.targetHost = '' + deliveryForm.targetHosts = [] + if (isPostgreSQL.value && mode === 'primary_replica') deliveryForm.replicaCount = 1 +} + function normalizeDNSLabel(value: string) { const normalized = value .toLowerCase() @@ -1676,7 +1799,7 @@ function normalizeDNSLabel(value: string) { } function generateInstanceName(serviceKey: string, businessLineName: string) { - const prefix = serviceKey === 'mysql' ? 'mysql' : 'nginx' + const prefix = serviceKey === 'pgsql' ? 'postgresql' : serviceKey === 'mysql' ? 'mysql' : 'nginx' // 限长业务线段,确保 63 字符截断不会吃掉随机后缀(prefix 5 + 连字符 2 + 后缀 6 = 13) const businessLine = normalizeDNSLabel(businessLineName).slice(0, 50).replace(/-+$/g, '') const alphabet = 'abcdefghijklmnopqrstuvwxyz0123456789' diff --git a/frontend/src/views/service/Management.vue b/frontend/src/views/service/Management.vue index 8be01c1..a56dbfd 100644 --- a/frontend/src/views/service/Management.vue +++ b/frontend/src/views/service/Management.vue @@ -173,6 +173,42 @@
+
+
+

PostgreSQL 交付记录

+ 共 {{ postgresqlDeliveryRecords.length }} 条 · 当前业务线:{{ currentName }} +
+
+ + + + + + + + + + + + + + + + + + + + + + + + + + +
集群名称任务编号连接地址命名空间交付状态创建时间操作
当前业务线暂无 PostgreSQL 交付记录
{{ task.instance_name }}{{ task.id }}{{ deliveryEndpoint(task) }}{{ task.namespace }}{{ deliveryStatusText(task.status) }}{{ formatDeliveryTime(task.created_at) }}查看详情
+
+
+ 查看一次性密码 -
+
一次性凭证已关闭 历史任务仍可查看连接信息,但平台不会再次展示 root 明文密码。
@@ -269,6 +305,7 @@ const servicePage = ref(1) const servicePageSize = ref(10) const deliveryRecords = ref([]) +const postgresqlDeliveryRecords = computed(() => deliveryRecords.value.filter((task) => task.service_type === 'postgresql' || task.component === 'postgresql')) const deliveryLoading = ref(false) const deliveryDetailVisible = ref(false) const selectedDeliveryTask = ref() @@ -348,7 +385,7 @@ const loadDeliveryRecords = async () => { } deliveryLoading.value = true try { - const items = await deliveryApi.listTasks({ businessLineId, component: 'mysql' }) + const items = await deliveryApi.listTasks({ businessLineId }) deliveryRecords.value = items.sort((a, b) => Date.parse(b.created_at) - Date.parse(a.created_at)) } catch (error) { deliveryRecords.value = [] @@ -382,7 +419,7 @@ function openDeliveryDetail(task: DeliveryTask) { } function latestDeliveryTaskForService(serviceName: string) { - return deliveryRecords.value.find((task) => task.instance_name === serviceName) + return deliveryRecords.value.find((task) => isMySQLDeliveryTask(task) && task.instance_name === serviceName) } function latestCredentialTaskForService(serviceName: string) { @@ -391,7 +428,13 @@ function latestCredentialTaskForService(serviceName: string) { } function credentialEligible(task: DeliveryTask) { - return ['finished', 'register_failed'].includes(task.status) && Boolean(task.credential_available) + return isMySQLDeliveryTask(task) + && ['finished', 'register_failed'].includes(task.status) + && Boolean(task.credential_available) +} + +function isMySQLDeliveryTask(task: DeliveryTask) { + return task.service_type === undefined || task.service_type === '' || task.service_type === 'mysql' } function markCredentialUnavailable(taskID: string) { @@ -435,7 +478,8 @@ function deliveryStatusClass(status: string) { function deliveryEndpoint(task: DeliveryTask) { const host = task.target_host_ip || task.target_host if (!host) return '等待主机分配' - return task.mysql_port ? `${host}:${task.mysql_port}` : host + const port = task.service_type === 'postgresql' ? task.postgresql_port : task.mysql_port + return port ? `${host}:${port}` : host } function formatDeliveryTime(value: string, withYear = false) { @@ -478,7 +522,7 @@ function downloadDeliveryCredential(task: DeliveryTask) { item.service || 'mysql', item.instance_name || task.instance_name, item.host || task.target_host_ip || task.target_host || '', - item.port || task.mysql_port || '', + item.port || (task.service_type === 'postgresql' ? task.postgresql_port : task.mysql_port) || '', item.username, item.account_host || '', item.password, diff --git a/frontend/vite.config.ts b/frontend/vite.config.ts index c06670d..87d2e65 100644 --- a/frontend/vite.config.ts +++ b/frontend/vite.config.ts @@ -1,44 +1,51 @@ -import { defineConfig } from 'vite' +import { defineConfig, loadEnv } from 'vite' import vue from '@vitejs/plugin-vue' import { resolve } from 'path' import AutoImport from 'unplugin-auto-import/vite' import Components from 'unplugin-vue-components/vite' import { ElementPlusResolver } from 'unplugin-vue-components/resolvers' -export default defineConfig({ - plugins: [ - vue(), - AutoImport({ - resolvers: [ElementPlusResolver()], - }), - Components({ - resolvers: [ElementPlusResolver()], - }), - ], - resolve: { - alias: { - '@': resolve(__dirname, 'src'), - }, - }, - server: { - port: 3000, - proxy: { - '/auth': { - target: 'http://localhost:8083', - changeOrigin: true, - }, - '/api': { - target: 'http://localhost:8080', - changeOrigin: true, - }, - '/swagger': { - target: 'http://localhost:8080', - changeOrigin: true, +export default defineConfig(({ mode }) => { + const env = loadEnv(mode, __dirname, '') + const backendTarget = env.XINFRA_BACKEND_PROXY_TARGET || 'http://localhost:8083' + + return { + plugins: [ + vue(), + AutoImport({ + resolvers: [ElementPlusResolver()], + }), + Components({ + resolvers: [ElementPlusResolver()], + }), + ], + resolve: { + alias: { + '@': resolve(__dirname, 'src'), }, }, - }, - build: { - outDir: 'dist', - sourcemap: false, - }, + server: { + host: '127.0.0.1', + port: 5173, + strictPort: true, + proxy: { + '/auth': { + target: backendTarget, + changeOrigin: true, + }, + '/api': { + target: backendTarget, + changeOrigin: true, + }, + '/swagger': { + target: backendTarget, + changeOrigin: true, + }, + }, + }, + build: { + outDir: 'dist', + sourcemap: false, + }, + } }) diff --git a/scripts/run-server-postgresql-local.sh b/scripts/run-server-postgresql-local.sh new file mode 100755 index 0000000..864135b --- /dev/null +++ b/scripts/run-server-postgresql-local.sh @@ -0,0 +1,29 @@ +#!/usr/bin/env bash + +set -euo pipefail + +SCRIPT_DIR=$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd) +REPO_ROOT=$(cd -- "${SCRIPT_DIR}/.." && pwd) +CONFIG_FILE=${1:-"${REPO_ROOT}/config/local/server-postgresql.env"} + +[[ -f ${CONFIG_FILE} ]] || { + printf 'error: configuration file not found: %s\n' "${CONFIG_FILE}" >&2 + exit 1 +} + +set -a +# shellcheck disable=SC1090 +source "${CONFIG_FILE}" +set +a + +if [[ -z ${AWX_TOKEN:-} && -z ${AWX_PASSWORD:-} ]]; then + AWX_K8S_NAMESPACE=${AWX_K8S_NAMESPACE:-awx} + AWX_K8S_ADMIN_PASSWORD_SECRET=${AWX_K8S_ADMIN_PASSWORD_SECRET:-awx-demo-admin-password} + AWX_PASSWORD=$(kubectl get secret \ + -n "${AWX_K8S_NAMESPACE}" "${AWX_K8S_ADMIN_PASSWORD_SECRET}" \ + -o jsonpath='{.data.password}' | base64 --decode) + export AWX_PASSWORD +fi + +cd "${REPO_ROOT}/server" +exec go run ./cmd/server diff --git a/scripts/setup-postgresql-awx.sh b/scripts/setup-postgresql-awx.sh new file mode 100755 index 0000000..076e1f1 --- /dev/null +++ b/scripts/setup-postgresql-awx.sh @@ -0,0 +1,343 @@ +#!/usr/bin/env bash + +set -euo pipefail + +SCRIPT_DIR=$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd) +REPO_ROOT=$(cd -- "${SCRIPT_DIR}/.." && pwd) +CONFIG_FILE=${1:-"${REPO_ROOT}/config/local/postgresql-awx.env"} + +die() { + printf 'error: %s\n' "$*" >&2 + exit 1 +} + +require_command() { + command -v "$1" >/dev/null 2>&1 || die "required command not found: $1" +} + +require_value() { + local name=$1 + [[ -n ${!name:-} ]] || die "${name} must be set in ${CONFIG_FILE}" +} + +[[ -f ${CONFIG_FILE} ]] || die "configuration file not found: ${CONFIG_FILE}" + +set -a +# shellcheck disable=SC1090 +source "${CONFIG_FILE}" +set +a + +for command_name in curl jq kubectl base64 openssl; do + require_command "${command_name}" +done + +for variable_name in \ + AWX_BASE_URL AWX_ORGANIZATION AWX_PROJECT_NAME AWX_PROJECT_SCM_URL \ + AWX_PROJECT_SOURCE_PATH AWX_PROJECT_SOURCE_DIR AWX_INVENTORY_NAME \ + AWX_HOST_NAME AWX_HOST_ADDRESS \ + AWX_MACHINE_CREDENTIAL_NAME AWX_SSH_PRIVATE_KEY_FILE \ + AWX_POSTGRES_CREDENTIAL_TYPE_NAME AWX_POSTGRES_CREDENTIAL_NAME \ + AWX_JOB_TEMPLATE_NAME AWX_PLAYBOOK; do + require_value "${variable_name}" +done + +AWX_BASE_URL=${AWX_BASE_URL%/} +AWX_USERNAME=${AWX_USERNAME:-admin} +AWX_K8S_NAMESPACE=${AWX_K8S_NAMESPACE:-awx} +AWX_K8S_ADMIN_PASSWORD_SECRET=${AWX_K8S_ADMIN_PASSWORD_SECRET:-awx-demo-admin-password} +AWX_K8S_TASK_APP=${AWX_K8S_TASK_APP:-awx-demo-task} +AWX_K8S_TASK_CONTAINER=${AWX_K8S_TASK_CONTAINER:-${AWX_K8S_TASK_APP}} +AWX_HOST_SSH_PORT=${AWX_HOST_SSH_PORT:-22} +AWX_HOST_SSH_USER=${AWX_HOST_SSH_USER:-root} +AWX_HOST_ANSIBLE_ADDRESS=${AWX_HOST_ANSIBLE_ADDRESS:-${AWX_HOST_ADDRESS}} +XINFRA_POSTGRES_REPLICATION_USER=${XINFRA_POSTGRES_REPLICATION_USER:-xinfra_replication} + +[[ -r ${AWX_SSH_PRIVATE_KEY_FILE} ]] || die "SSH private key is not readable: ${AWX_SSH_PRIVATE_KEY_FILE}" +case ${AWX_PROJECT_SOURCE_PATH} in + _[A-Za-z0-9_-]*) ;; + *) die "AWX_PROJECT_SOURCE_PATH must be a safe name beginning with '_'" ;; +esac + +PROJECT_SOURCE=$(cd -- "${REPO_ROOT}/${AWX_PROJECT_SOURCE_DIR}" && pwd) +case ${PROJECT_SOURCE}/ in + "${REPO_ROOT}/"*) ;; + *) die "AWX_PROJECT_SOURCE_DIR must resolve inside the repository" ;; +esac +[[ -f ${PROJECT_SOURCE}/${AWX_PLAYBOOK} ]] || die "playbook not found: ${PROJECT_SOURCE}/${AWX_PLAYBOOK}" + +if [[ -z ${AWX_TOKEN:-} && -z ${AWX_PASSWORD:-} ]]; then + AWX_PASSWORD=$(kubectl get secret \ + -n "${AWX_K8S_NAMESPACE}" "${AWX_K8S_ADMIN_PASSWORD_SECRET}" \ + -o jsonpath='{.data.password}' | base64 --decode) +fi +[[ -n ${AWX_TOKEN:-} || -n ${AWX_PASSWORD:-} ]] || die "AWX_TOKEN or AWX_PASSWORD is required" + +AWX_AUTH_ARGS=() +if [[ -n ${AWX_TOKEN:-} ]]; then + AWX_AUTH_ARGS=(-H "Authorization: Bearer ${AWX_TOKEN}") +else + AWX_AUTH_ARGS=(-u "${AWX_USERNAME}:${AWX_PASSWORD}") +fi + +awx_request() { + local method=$1 + local path=$2 + local payload=${3:-} + local args=(-sS --fail-with-body -X "${method}" -H 'Accept: application/json') + args+=("${AWX_AUTH_ARGS[@]}") + if [[ -n ${payload} ]]; then + args+=(-H 'Content-Type: application/json' --data "${payload}") + fi + curl "${args[@]}" "${AWX_BASE_URL}${path}" +} + +find_named_id() { + local endpoint=$1 + local name=$2 + awx_request GET "${endpoint}?page_size=200" | jq -er --arg name "${name}" \ + '.results[] | select(.name == $name) | .id' | head -n 1 +} + +find_inventory_host_id() { + local inventory_id=$1 + awx_request GET "/api/v2/inventories/${inventory_id}/hosts/?page_size=200" | \ + jq -er --arg name "${AWX_HOST_NAME}" '.results[] | select(.name == $name) | .id' | head -n 1 +} + +upsert_named_object() { + local endpoint=$1 + local name=$2 + local payload=$3 + local id + id=$(find_named_id "${endpoint}" "${name}" 2>/dev/null || true) + if [[ -n ${id} ]]; then + awx_request PATCH "${endpoint}${id}/" "${payload}" | jq -er '.id' + else + awx_request POST "${endpoint}" "${payload}" | jq -er '.id' + fi +} + +printf 'Checking AWX API at %s...\n' "${AWX_BASE_URL}" +awx_request GET '/api/v2/ping/' >/dev/null + +if [[ -z ${AWX_K8S_TASK_POD:-} ]]; then + AWX_K8S_TASK_POD=$(kubectl get pods -n "${AWX_K8S_NAMESPACE}" \ + -l "app.kubernetes.io/name=${AWX_K8S_TASK_APP}" \ + -o jsonpath='{.items[0].metadata.name}') +fi +[[ -n ${AWX_K8S_TASK_POD} ]] || die "could not find the AWX task pod" + +PROJECT_SOURCE_REMOTE_PATH=/var/lib/awx/projects/${AWX_PROJECT_SOURCE_PATH} +printf 'Refreshing local SCM source in %s/%s:%s...\n' \ + "${AWX_K8S_NAMESPACE}" "${AWX_K8S_TASK_POD}" "${PROJECT_SOURCE_REMOTE_PATH}" +kubectl exec -n "${AWX_K8S_NAMESPACE}" "${AWX_K8S_TASK_POD}" \ + -c "${AWX_K8S_TASK_CONTAINER}" -- \ + sh -c 'mkdir -p "$1"; find "$1" -mindepth 1 -maxdepth 1 -exec rm -rf -- {} +' sh "${PROJECT_SOURCE_REMOTE_PATH}" +kubectl cp "${PROJECT_SOURCE}/." \ + "${AWX_K8S_NAMESPACE}/${AWX_K8S_TASK_POD}:${PROJECT_SOURCE_REMOTE_PATH}" \ + -c "${AWX_K8S_TASK_CONTAINER}" +kubectl exec -n "${AWX_K8S_NAMESPACE}" "${AWX_K8S_TASK_POD}" \ + -c "${AWX_K8S_TASK_CONTAINER}" -- sh -c \ + 'cd "$1" && git init -q && git config user.email xinfra-local@localhost && git config user.name xinfra-local && git add -A && if ! git diff --cached --quiet; then git commit -qm local; fi' \ + sh "${PROJECT_SOURCE_REMOTE_PATH}" + +organization_id=$(find_named_id '/api/v2/organizations/' "${AWX_ORGANIZATION}" 2>/dev/null || true) +[[ -n ${organization_id} ]] || die "AWX organization not found: ${AWX_ORGANIZATION}" + +project_payload=$(jq -nc \ + --arg name "${AWX_PROJECT_NAME}" \ + --arg description "${AWX_PROJECT_DESCRIPTION:-}" \ + --arg scm_url "${AWX_PROJECT_SCM_URL}" \ + --argjson organization "${organization_id}" \ + '{name:$name, description:$description, organization:$organization, scm_type:"git", scm_url:$scm_url, scm_update_on_launch:false, scm_clean:false, scm_delete_on_update:false}') +project_id=$(upsert_named_object '/api/v2/projects/' "${AWX_PROJECT_NAME}" "${project_payload}") + +project_local_path=$(awx_request GET "/api/v2/projects/${project_id}/" | jq -er '.local_path') +case ${project_local_path} in + _[A-Za-z0-9_-]*) ;; + *) die "AWX returned an unsafe project local_path: ${project_local_path}" ;; +esac + +PROJECT_REMOTE_PATH=/var/lib/awx/projects/${project_local_path} +project_status=$(awx_request GET "/api/v2/projects/${project_id}/" | jq -r '.status') +if [[ ${project_status} == pending || ${project_status} == running ]]; then + for _ in $(seq 1 60); do + sleep 2 + project_status=$(awx_request GET "/api/v2/projects/${project_id}/" | jq -r '.status') + [[ ${project_status} == pending || ${project_status} == running ]] || break + done +fi +if [[ ${project_status} != successful ]]; then + printf 'Refreshing the AWX-managed project directory...\n' + kubectl exec -n "${AWX_K8S_NAMESPACE}" "${AWX_K8S_TASK_POD}" \ + -c "${AWX_K8S_TASK_CONTAINER}" -- \ + sh -c 'mkdir -p "$1"; find "$1" -mindepth 1 -maxdepth 1 -exec rm -rf -- {} +' sh "${PROJECT_REMOTE_PATH}" +fi + +project_update_id=$(awx_request POST "/api/v2/projects/${project_id}/update/" '{}' | jq -er '.id') +for _ in $(seq 1 60); do + project_status=$(awx_request GET "/api/v2/project_updates/${project_update_id}/" | jq -r '.status') + [[ ${project_status} == pending || ${project_status} == waiting || ${project_status} == running ]] || break + sleep 2 +done +[[ ${project_status} == successful ]] || die "AWX project update ${project_update_id} finished with status ${project_status}" + +if ! awx_request GET "/api/v2/projects/${project_id}/playbooks/" | \ + jq -e --arg playbook "${AWX_PLAYBOOK}" 'index($playbook) != null' >/dev/null; then + die "AWX project ${project_id} does not expose playbook ${AWX_PLAYBOOK}" +fi + +inventory_payload=$(jq -nc \ + --arg name "${AWX_INVENTORY_NAME}" \ + --arg description "${AWX_INVENTORY_DESCRIPTION:-}" \ + --argjson organization "${organization_id}" \ + '{name:$name, description:$description, organization:$organization, kind:""}') +inventory_id=$(upsert_named_object '/api/v2/inventories/' "${AWX_INVENTORY_NAME}" "${inventory_payload}") + +host_variables=$(jq -nc \ + --arg ansible_host "${AWX_HOST_ANSIBLE_ADDRESS}" \ + --arg xinfra_public_address "${AWX_HOST_ADDRESS}" \ + --arg ansible_user "${AWX_HOST_SSH_USER}" \ + --argjson ansible_port "${AWX_HOST_SSH_PORT}" \ + '{ansible_host:$ansible_host, xinfra_public_address:$xinfra_public_address, ansible_user:$ansible_user, ansible_port:$ansible_port, ansible_python_interpreter:"/usr/bin/python3"}') +host_payload=$(jq -nc \ + --arg name "${AWX_HOST_NAME}" \ + --arg variables "${host_variables}" \ + --argjson inventory "${inventory_id}" \ + '{name:$name, inventory:$inventory, enabled:true, variables:$variables}') +host_id=$(find_inventory_host_id "${inventory_id}" 2>/dev/null || true) +if [[ -n ${host_id} ]]; then + host_id=$(awx_request PATCH "/api/v2/hosts/${host_id}/" "${host_payload}" | jq -er '.id') +else + host_id=$(awx_request POST '/api/v2/hosts/' "${host_payload}" | jq -er '.id') +fi + +machine_type_id=$(awx_request GET '/api/v2/credential_types/?page_size=200' | \ + jq -er '.results[] | select(.kind == "ssh") | .id' | head -n 1) +ssh_key_data=$(<"${AWX_SSH_PRIVATE_KEY_FILE}") +machine_credential_payload=$(jq -nc \ + --arg name "${AWX_MACHINE_CREDENTIAL_NAME}" \ + --arg description "SSH access for the XINFRA PostgreSQL host pool" \ + --arg username "${AWX_HOST_SSH_USER}" \ + --arg ssh_key_data "${ssh_key_data}" \ + --argjson organization "${organization_id}" \ + --argjson credential_type "${machine_type_id}" \ + '{name:$name, description:$description, organization:$organization, credential_type:$credential_type, inputs:{username:$username, ssh_key_data:$ssh_key_data}}') +machine_credential_id=$(upsert_named_object '/api/v2/credentials/' "${AWX_MACHINE_CREDENTIAL_NAME}" "${machine_credential_payload}") + +postgres_credential_type_payload=$(jq -nc \ + --arg name "${AWX_POSTGRES_CREDENTIAL_TYPE_NAME}" \ + '{ + name:$name, + description:"Injects PostgreSQL delivery secrets as execution environment variables", + kind:"cloud", + inputs:{ + fields:[ + {id:"admin_password", label:"PostgreSQL administrator password", type:"string", secret:true}, + {id:"replication_user", label:"PostgreSQL replication user", type:"string", default:"xinfra_replication"}, + {id:"replication_password", label:"PostgreSQL replication password", type:"string", secret:true} + ], + required:["admin_password", "replication_password"] + }, + injectors:{env:{ + XINFRA_POSTGRES_ADMIN_PASSWORD:"{{ admin_password }}", + XINFRA_POSTGRES_REPLICATION_USER:"{{ replication_user }}", + XINFRA_POSTGRES_REPLICATION_PASSWORD:"{{ replication_password }}" + }} + }') +postgres_credential_type_id=$(upsert_named_object '/api/v2/credential_types/' \ + "${AWX_POSTGRES_CREDENTIAL_TYPE_NAME}" "${postgres_credential_type_payload}") + +postgres_credential_id=$(find_named_id '/api/v2/credentials/' "${AWX_POSTGRES_CREDENTIAL_NAME}" 2>/dev/null || true) +if [[ -z ${postgres_credential_id} ]]; then + admin_password=${XINFRA_POSTGRES_ADMIN_PASSWORD:-$(openssl rand -base64 24 | tr -d '\n')} + replication_password=${XINFRA_POSTGRES_REPLICATION_PASSWORD:-$(openssl rand -base64 24 | tr -d '\n')} + postgres_credential_payload=$(jq -nc \ + --arg name "${AWX_POSTGRES_CREDENTIAL_NAME}" \ + --arg description "Runtime secrets for XINFRA PostgreSQL delivery" \ + --arg admin_password "${admin_password}" \ + --arg replication_user "${XINFRA_POSTGRES_REPLICATION_USER}" \ + --arg replication_password "${replication_password}" \ + --argjson organization "${organization_id}" \ + --argjson credential_type "${postgres_credential_type_id}" \ + '{name:$name, description:$description, organization:$organization, credential_type:$credential_type, inputs:{admin_password:$admin_password, replication_user:$replication_user, replication_password:$replication_password}}') + postgres_credential_id=$(awx_request POST '/api/v2/credentials/' "${postgres_credential_payload}" | jq -er '.id') +elif [[ -n ${XINFRA_POSTGRES_ADMIN_PASSWORD:-} || -n ${XINFRA_POSTGRES_REPLICATION_PASSWORD:-} ]]; then + [[ -n ${XINFRA_POSTGRES_ADMIN_PASSWORD:-} && -n ${XINFRA_POSTGRES_REPLICATION_PASSWORD:-} ]] || \ + die "set both PostgreSQL passwords together when rotating an existing credential" + postgres_credential_payload=$(jq -nc \ + --arg admin_password "${XINFRA_POSTGRES_ADMIN_PASSWORD}" \ + --arg replication_user "${XINFRA_POSTGRES_REPLICATION_USER}" \ + --arg replication_password "${XINFRA_POSTGRES_REPLICATION_PASSWORD}" \ + '{inputs:{admin_password:$admin_password, replication_user:$replication_user, replication_password:$replication_password}}') + postgres_credential_id=$(awx_request PATCH "/api/v2/credentials/${postgres_credential_id}/" \ + "${postgres_credential_payload}" | jq -er '.id') +fi + +prevent_fallback=false +instance_group_id= +if [[ -n ${AWX_INSTANCE_GROUP:-} ]]; then + instance_group_id=$(find_named_id '/api/v2/instance_groups/' "${AWX_INSTANCE_GROUP}" 2>/dev/null || true) + [[ -n ${instance_group_id} ]] || die "AWX instance group not found: ${AWX_INSTANCE_GROUP}" + prevent_fallback=true +fi + +job_template_payload=$(jq -nc \ + --arg name "${AWX_JOB_TEMPLATE_NAME}" \ + --arg description "${AWX_JOB_TEMPLATE_DESCRIPTION:-PostgreSQL delivery template}" \ + --arg playbook "${AWX_PLAYBOOK}" \ + --argjson organization "${organization_id}" \ + --argjson inventory "${inventory_id}" \ + --argjson project "${project_id}" \ + --argjson prevent_fallback "${prevent_fallback}" \ + '{ + name:$name, + description:$description, + organization:$organization, + inventory:$inventory, + project:$project, + playbook:$playbook, + job_type:"run", + ask_inventory_on_launch:true, + ask_variables_on_launch:true, + ask_limit_on_launch:true, + allow_simultaneous:true, + prevent_instance_group_fallback:$prevent_fallback + }') +job_template_id=$(upsert_named_object '/api/v2/job_templates/' \ + "${AWX_JOB_TEMPLATE_NAME}" "${job_template_payload}") + +while IFS= read -r existing_group_id; do + [[ -z ${existing_group_id} || ${existing_group_id} == "${instance_group_id}" ]] && continue + awx_request POST "/api/v2/job_templates/${job_template_id}/instance_groups/" \ + "$(jq -nc --argjson id "${existing_group_id}" '{id:$id,disassociate:true}')" >/dev/null +done < <(awx_request GET "/api/v2/job_templates/${job_template_id}/instance_groups/" | jq -r '.results[].id') + +if ! awx_request GET "/api/v2/job_templates/${job_template_id}/credentials/" | \ + jq -e --argjson id "${machine_credential_id}" '.results | any(.id == $id)' >/dev/null; then + awx_request POST "/api/v2/job_templates/${job_template_id}/credentials/" \ + "$(jq -nc --argjson id "${machine_credential_id}" '{id:$id}')" >/dev/null +fi +if ! awx_request GET "/api/v2/job_templates/${job_template_id}/credentials/" | \ + jq -e --argjson id "${postgres_credential_id}" '.results | any(.id == $id)' >/dev/null; then + awx_request POST "/api/v2/job_templates/${job_template_id}/credentials/" \ + "$(jq -nc --argjson id "${postgres_credential_id}" '{id:$id}')" >/dev/null +fi +if [[ -n ${instance_group_id} ]]; then + if ! awx_request GET "/api/v2/job_templates/${job_template_id}/instance_groups/" | \ + jq -e --argjson id "${instance_group_id}" '.results | any(.id == $id)' >/dev/null; then + awx_request POST "/api/v2/job_templates/${job_template_id}/instance_groups/" \ + "$(jq -nc --argjson id "${instance_group_id}" '{id:$id}')" >/dev/null + fi +fi + +printf '\nPostgreSQL AWX configuration is ready:\n' +printf ' Project: %s (id=%s)\n' "${AWX_PROJECT_NAME}" "${project_id}" +printf ' Inventory: %s (id=%s)\n' "${AWX_INVENTORY_NAME}" "${inventory_id}" +printf ' Host: %s -> %s (id=%s)\n' "${AWX_HOST_NAME}" "${AWX_HOST_ADDRESS}" "${host_id}" +printf ' SSH credential: %s (id=%s)\n' "${AWX_MACHINE_CREDENTIAL_NAME}" "${machine_credential_id}" +printf ' PG credential: %s (id=%s)\n' "${AWX_POSTGRES_CREDENTIAL_NAME}" "${postgres_credential_id}" +printf ' Job Template: %s (id=%s)\n' "${AWX_JOB_TEMPLATE_NAME}" "${job_template_id}" +if [[ -n ${instance_group_id} ]]; then + printf ' Instance Group: %s (id=%s)\n' "${AWX_INSTANCE_GROUP}" "${instance_group_id}" +fi diff --git a/server/.env.example b/server/.env.example index 9dbf569..18a080e 100644 --- a/server/.env.example +++ b/server/.env.example @@ -41,9 +41,11 @@ SSO_ENABLED=true BOOTSTRAP_ADMIN_USERNAME=admin BOOTSTRAP_ADMIN_PASSWORD=change-this-bootstrap-admin-password -# MySQL service delivery (AWX is required when the scheduler is enabled) +# Database delivery (MySQL/PostgreSQL; AWX is required when the scheduler is enabled) DELIVERY_SCHEDULER_ENABLED=false DELIVERY_DISPATCH_SECONDS=5 +DELIVERY_POLL_SECONDS=5 +DELIVERY_BUSINESS_LIMIT=1 DELIVERY_CALLBACK_BASE_URL=http://authserver-backend.authserver.svc.cluster.local:8083 DELIVERY_CREDENTIAL_SECRET= DELIVERY_RESERVATION_TTL_MINUTES=120 @@ -54,6 +56,7 @@ DELIVERY_HOST_INSTANCE_LIMIT=4 # 数据盘挂载点白名单(逗号分隔,第一项为默认值) DELIVERY_DATA_DISKS=/data,/disk1,/mnt,/opt/mysql-delivery AWX_BASE_URL= +# AWX API must point at the HTTP API endpoint, not a Minikube SSH port. AWX_TOKEN= AWX_USERNAME= AWX_PASSWORD= @@ -66,6 +69,12 @@ DELIVERY_MYSQL_INSPECT_TIMEOUT_SECONDS=90 # AWX Job Template ID for ansible/mysql-rollback.yml; required for automatic cleanup DELIVERY_ROLLBACK_TEMPLATE_ID=0 DELIVERY_SERVICE_TOKEN= + +# PostgreSQL runtime secrets are injected by an AWX Credential and are never +# persisted in xinfra task payloads: +# XINFRA_POSTGRES_ADMIN_PASSWORD +# XINFRA_POSTGRES_REPLICATION_PASSWORD +# XINFRA_POSTGRES_REPLICATION_USER (optional, defaults to xinfra_replication) CLOUDDM_REGISTER_URL= CLOUDDM_DELETE_URL= CLOUDDM_API_TOKEN= diff --git a/server/internal/config/config.go b/server/internal/config/config.go index 63b3f60..51c021b 100644 --- a/server/internal/config/config.go +++ b/server/internal/config/config.go @@ -78,12 +78,14 @@ type Config struct { DeliveryServiceToken string DeliverySchedulerEnabled bool DeliveryDispatchSeconds int + DeliveryPollSeconds int DeliveryCallbackBaseURL string DeliveryCredentialSecret string ReservationTTLMinutes int DeliveryGlobalLimit int DeliveryTargetLimit int DeliveryHostInstanceLimit int + DeliveryBusinessLimit int DeliveryDataDisks []string SINABaseURL string SINAUsername string @@ -164,12 +166,14 @@ func Load() Config { DeliveryServiceToken: env("DELIVERY_SERVICE_TOKEN", ""), DeliverySchedulerEnabled: envBool("DELIVERY_SCHEDULER_ENABLED", false), DeliveryDispatchSeconds: envInt("DELIVERY_DISPATCH_SECONDS", 5), + DeliveryPollSeconds: envInt("DELIVERY_POLL_SECONDS", 5), DeliveryCallbackBaseURL: trimURL(env("DELIVERY_CALLBACK_BASE_URL", publicBaseURL)), DeliveryCredentialSecret: env("DELIVERY_CREDENTIAL_SECRET", env("JWT_SECRET", "change-this-secret")), ReservationTTLMinutes: envInt("DELIVERY_RESERVATION_TTL_MINUTES", 120), DeliveryGlobalLimit: envInt("DELIVERY_GLOBAL_LIMIT", 2), DeliveryTargetLimit: envInt("DELIVERY_TARGET_LIMIT", 2), DeliveryHostInstanceLimit: envInt("DELIVERY_HOST_INSTANCE_LIMIT", 4), + DeliveryBusinessLimit: envInt("DELIVERY_BUSINESS_LIMIT", 1), DeliveryDataDisks: splitCSV(env("DELIVERY_DATA_DISKS", "/data,/disk1,/mnt,/opt/mysql-delivery")), SINABaseURL: trimURL(env("SINA_BASE_URL", "https://sinai.qiniu.io:443")), SINAUsername: env("SINA_USERNAME", ""), diff --git a/server/internal/database/database.go b/server/internal/database/database.go index 8b1863e..8653df1 100644 --- a/server/internal/database/database.go +++ b/server/internal/database/database.go @@ -27,6 +27,9 @@ func AutoMigrate(db *gorm.DB) error { &model.ResourceReservation{}, &model.DeploymentResult{}, &model.DeploymentCredential{}, + &model.PostgreSQLCluster{}, + &model.PostgreSQLInstance{}, + &model.PostgreSQLResourceUsage{}, &model.ResourceUsage{}, &model.ExecutionJob{}, &model.RollbackJob{}, diff --git a/server/internal/handler/delivery.go b/server/internal/handler/delivery.go index 9e79e7b..eb4416f 100644 --- a/server/internal/handler/delivery.go +++ b/server/internal/handler/delivery.go @@ -15,12 +15,22 @@ import ( "gorm.io/gorm" ) -type DeliveryHandler struct{ service *service.DeliveryService } +type DeliveryHandler struct { + service *service.DeliveryService +} + +type PostgreSQLDeliveryHandler struct { + service *service.PostgreSQLDeliveryService +} func NewDeliveryHandler(s *service.DeliveryService) *DeliveryHandler { return &DeliveryHandler{service: s} } +func NewPostgreSQLDeliveryHandler(s *service.PostgreSQLDeliveryService) *PostgreSQLDeliveryHandler { + return &PostgreSQLDeliveryHandler{service: s} +} + type DeliveryCallbackHandler struct { service *service.DeliveryService token string @@ -108,6 +118,30 @@ func (h *DeliveryHandler) CreateMySQL(c *gin.Context) { c.JSON(status, gin.H{"task": task, "idempotent_replay": existing}) } +// CreatePostgreSQL submits a standalone or primary/replica PostgreSQL delivery. +func (h *PostgreSQLDeliveryHandler) CreatePostgreSQL(c *gin.Context) { + claims, ok := CurrentClaims(c) + if !ok { + c.JSON(http.StatusUnauthorized, gin.H{"error": "missing current user"}) + return + } + var req service.PostgreSQLDeliveryInput + if err := c.ShouldBindJSON(&req); err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + task, existing, err := h.service.CreateTask(c.Request.Context(), claims.UserID, claims.IsAdmin, c.GetHeader("Idempotency-Key"), req) + if err != nil { + c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) + return + } + status := http.StatusAccepted + if existing { + status = http.StatusOK + } + c.JSON(status, gin.H{"task": task, "idempotent_replay": existing}) +} + // List 获取交付任务列表 // @Summary 获取交付任务列表 // @Description 返回当前用户可见的交付任务列表(管理员可见全部) diff --git a/server/internal/model/delivery.go b/server/internal/model/delivery.go index 681bd01..4722d77 100644 --- a/server/internal/model/delivery.go +++ b/server/internal/model/delivery.go @@ -39,12 +39,14 @@ type DeliveryTask struct { RequestedBy uint64 `gorm:"not null;index" json:"requested_by"` Component string `gorm:"size:32;not null;default:mysql;index" json:"component"` TargetType string `gorm:"size:32;not null" json:"target_type"` + ServiceType string `gorm:"size:32;not null;default:'mysql';index" json:"service_type"` TargetID uint64 `gorm:"not null;index" json:"target_id"` Namespace string `gorm:"size:63;not null;index" json:"namespace"` InstanceName string `gorm:"size:63;not null" json:"instance_name"` - TargetHost string `gorm:"size:128" json:"target_host,omitempty"` + TargetHost string `gorm:"size:1024" json:"target_host,omitempty"` TargetHostIP string `gorm:"size:64" json:"target_host_ip,omitempty"` MySQLPort int `gorm:"column:mysql_port;not null;default:3307" json:"mysql_port"` + PostgreSQLPort int `gorm:"column:postgresql_port;not null;default:15432" json:"postgresql_port,omitempty"` Status string `gorm:"size:32;not null;index" json:"status"` ImmutablePayload string `gorm:"type:json;not null" json:"immutable_payload"` PayloadHash string `gorm:"size:64;not null" json:"payload_hash"` @@ -126,6 +128,69 @@ type DeploymentCredential struct { UpdatedAt time.Time `json:"updated_at"` } +type PostgreSQLCluster struct { + ID uint64 `gorm:"primaryKey" json:"id"` + TaskID string `gorm:"size:36;not null;uniqueIndex" json:"task_id"` + BusinessLineID uint64 `gorm:"not null;index" json:"business_line_id"` + TargetID uint64 `gorm:"not null;index" json:"target_id"` + Name string `gorm:"size:63;not null;uniqueIndex" json:"name"` + VersionMajor string `gorm:"size:8;not null" json:"version_major"` + Topology string `gorm:"size:32;not null" json:"topology"` + ReplicationMode string `gorm:"size:32;not null;default:'async'" json:"replication_mode"` + FailoverMode string `gorm:"size:32;not null;default:'manual'" json:"failover_mode"` + PrimaryInstanceID uint64 `gorm:"not null;default:0" json:"primary_instance_id"` + Status string `gorm:"size:32;not null;index" json:"status"` + BackupStatus string `gorm:"size:32;not null;default:'not_configured'" json:"backup_status"` + MonitoringStatus string `gorm:"size:32;not null;default:'not_configured'" json:"monitoring_status"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +type PostgreSQLInstance struct { + ID uint64 `gorm:"primaryKey" json:"id"` + TaskID string `gorm:"size:36;not null;index" json:"task_id"` + ClusterID uint64 `gorm:"not null;index" json:"cluster_id"` + BusinessLineID uint64 `gorm:"not null;index" json:"business_line_id"` + TargetID uint64 `gorm:"not null;index" json:"target_id"` + InstanceID string `gorm:"size:63;not null;uniqueIndex" json:"instance_id"` + Hostname string `gorm:"size:255;not null;uniqueIndex:idx_postgresql_host_port,priority:1" json:"hostname"` + HostIP string `gorm:"size:64;not null" json:"host_ip"` + Port int `gorm:"not null;index:idx_postgresql_host_port,priority:2" json:"port"` + DataDir string `gorm:"size:512;not null;uniqueIndex" json:"data_dir"` + ConfigDir string `gorm:"size:512;not null" json:"config_dir"` + LogDir string `gorm:"size:512;not null" json:"log_dir"` + SystemdUnit string `gorm:"size:128;not null" json:"systemd_unit"` + VersionMajor string `gorm:"size:8;not null" json:"version_major"` + VersionFull string `gorm:"size:32;not null;default:''" json:"version_full"` + Role string `gorm:"size:16;not null" json:"role"` + HAComponentRole string `gorm:"size:16;not null;default:'database'" json:"ha_component_role"` + UpstreamInstanceID string `gorm:"size:63;not null;default:''" json:"upstream_instance_id"` + ReplicationSlotName string `gorm:"size:63;not null;default:''" json:"replication_slot_name"` + Status string `gorm:"size:32;not null;index" json:"status"` + BackupStatus string `gorm:"size:32;not null;default:'not_configured'" json:"backup_status"` + MonitoringStatus string `gorm:"size:32;not null;default:'not_configured'" json:"monitoring_status"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} + +type PostgreSQLResourceUsage struct { + ID uint64 `gorm:"primaryKey" json:"id"` + TaskID string `gorm:"size:36;not null;index" json:"task_id"` + InstanceID uint64 `gorm:"not null;uniqueIndex" json:"instance_id"` + ClusterID uint64 `gorm:"not null;index" json:"cluster_id"` + BusinessLineID uint64 `gorm:"not null;index" json:"business_line_id"` + TargetID uint64 `gorm:"not null;index" json:"target_id"` + CPUMilli int64 `gorm:"not null" json:"cpu_milli"` + MemoryMi int64 `gorm:"not null" json:"memory_mi"` + StorageGi int64 `gorm:"not null" json:"storage_gi"` + Port int `gorm:"not null" json:"port"` + DataDir string `gorm:"size:512;not null" json:"data_dir"` + Status string `gorm:"size:32;not null;index" json:"status"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` + ReleasedAt *time.Time `json:"released_at,omitempty"` +} + type ExecutionJob struct { ID uint64 `gorm:"primaryKey" json:"id"` TaskID string `gorm:"size:36;not null;uniqueIndex" json:"task_id"` diff --git a/server/internal/router/router.go b/server/internal/router/router.go index 994db1d..34226c1 100644 --- a/server/internal/router/router.go +++ b/server/internal/router/router.go @@ -73,8 +73,10 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { wayneRoleBindingService := service.NewWayneRoleBindingService(deps.Config, deps.DB) deliveryService := service.NewDeliveryService(deps.Config, deps.DB, auditService) machineService := service.NewMachineService(deps.Config, deps.DB) + postgresqlDeliveryService := service.NewPostgreSQLDeliveryService(deps.Config, deps.DB, deliveryService) if deps.Config.DeliverySchedulerEnabled { go deliveryService.Run(context.Background()) + go postgresqlDeliveryService.Run(context.Background()) } machineService.Run(context.Background()) @@ -90,6 +92,7 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { oauthHandler := handler.NewOAuthHandler(deps.Config, deps.DB, auditService) deliveryHandler := handler.NewDeliveryHandler(deliveryService) deliveryCallbackHandler := handler.NewDeliveryCallbackHandler(deliveryService, deps.Config.AWXWebhookToken) + postgresqlDeliveryHandler := handler.NewPostgreSQLDeliveryHandler(postgresqlDeliveryService) containerServiceHandler := handler.NewContainerServiceHandler(deps.DB, wayneRoleBindingService) taskLogHandler := handler.NewTaskLogHandler(deps.DB, deliveryService, wayneRoleBindingService) machineHandler := handler.NewMachineHandler(machineService) @@ -155,6 +158,7 @@ func registerAuthServerRoutes(r *gin.Engine, deps Dependencies) { protected.GET("/delivery/targets/:target_id/hosts/:host/mount-paths", deliveryHandler.TargetHostMountPaths) protected.PUT("/delivery/quotas", deliveryHandler.UpsertQuota) protected.POST("/delivery/mysql", deliveryHandler.CreateMySQL) + protected.POST("/delivery/postgresql", postgresqlDeliveryHandler.CreatePostgreSQL) protected.GET("/delivery/tasks", deliveryHandler.List) protected.GET("/delivery/business-lines/:id/mysql-services", deliveryHandler.MySQLServiceLedger) protected.POST("/delivery/business-lines/:id/mysql-services/sync", deliveryHandler.SyncMySQLServiceLedger) diff --git a/server/internal/service/awx.go b/server/internal/service/awx.go index 8e7b41f..5cef55e 100644 --- a/server/internal/service/awx.go +++ b/server/internal/service/awx.go @@ -363,6 +363,11 @@ var ansibleHostPattern = regexp.MustCompile(`(?m)^\s*ansible_host\s*:\s*"?([^"\s func AWXHostIP(host AWXInventoryHost) string { var parsed map[string]any if err := json.Unmarshal([]byte(host.Variables), &parsed); err == nil { + for _, key := range []string{"xinfra_public_address", "xinfra_host_address"} { + if value, ok := parsed[key].(string); ok && strings.TrimSpace(value) != "" { + return strings.TrimSpace(value) + } + } if value, ok := parsed["ansible_host"].(string); ok && strings.TrimSpace(value) != "" { return strings.TrimSpace(value) } diff --git a/server/internal/service/awx_test.go b/server/internal/service/awx_test.go index 09c8706..5e8a1fe 100644 --- a/server/internal/service/awx_test.go +++ b/server/internal/service/awx_test.go @@ -73,3 +73,13 @@ func TestAWXClientCreatesTokenFromCredentials(t *testing.T) { t.Fatalf("token requests = %d, want 1", tokenRequests) } } + +func TestAWXHostIPPrefersXInfraPublicAddress(t *testing.T) { + host := AWXInventoryHost{ + Name: "postgresql-218-11-5-224", + Variables: `{"ansible_host":"127.0.0.1","xinfra_public_address":"218.11.5.224"}`, + } + if got := AWXHostIP(host); got != "218.11.5.224" { + t.Fatalf("AWXHostIP = %q, want public address", got) + } +} diff --git a/server/internal/service/delivery.go b/server/internal/service/delivery.go index 79d1f5d..931df17 100644 --- a/server/internal/service/delivery.go +++ b/server/internal/service/delivery.go @@ -160,7 +160,7 @@ type DeploymentCredentialView struct { // targetMetadata describes the native VM候选节点池以及部署形态,由 AWX inventory hosts 动态组装。 type targetMetadata struct { Topology string `json:"topology"` - MySQLPort int `json:"mysql_port"` + MySQLPort int `json:"mysql_port,omitempty"` Hosts []targetHost `json:"hosts"` } @@ -410,6 +410,14 @@ func (s *DeliveryService) ListTargets(ctx context.Context, component string) ([] if err != nil { continue } + if component == postgresqlServiceType { + target.TargetType = "host_pool" + meta := parseTargetMetadata(target.Metadata) + meta.MySQLPort = 0 + if raw, marshalErr := json.Marshal(meta); marshalErr == nil { + target.Metadata = string(raw) + } + } targets = append(targets, target) } return targets, nil @@ -692,6 +700,7 @@ func (s *DeliveryService) CreateTask(ctx context.Context, userID uint64, isAdmin RequestedBy: userID, Component: "mysql", TargetType: target.TargetType, + ServiceType: "mysql", TargetID: target.ID, Namespace: input.Namespace, InstanceName: input.InstanceName, @@ -1376,6 +1385,9 @@ func (s *DeliveryService) RevealDeploymentCredentials(ctx context.Context, taskI if err != nil { return nil, err } + if !isMySQLServiceType(task.ServiceType) { + return nil, fmt.Errorf("task service type %q does not provide MySQL credentials", task.ServiceType) + } if task.Status != model.TaskFinished && task.Status != model.TaskRegisterFailed { return nil, fmt.Errorf("task credentials are available only after a successful deployment") } @@ -1593,7 +1605,7 @@ func (s *DeliveryService) claimAndReserve(ctx context.Context) (*model.DeliveryT var target DeliveryTarget dispatchable := false err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}).Where("status = ?", model.TaskPending).Order("created_at ASC").First(&task).Error; err != nil { + if err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}).Where("status = ? AND (service_type = '' OR service_type = 'mysql')", model.TaskPending).Order("created_at ASC").First(&task).Error; err != nil { return err } var targetErr error @@ -1930,6 +1942,9 @@ func (s *DeliveryService) CreateExecution(ctx context.Context, taskID, payloadHa if err := s.db.WithContext(ctx).First(&task, "id = ?", taskID).Error; err != nil { return nil, false, err } + if task.ServiceType != "" && task.ServiceType != "mysql" { + return nil, false, fmt.Errorf("task service type %q is not handled by the MySQL executor", task.ServiceType) + } if task.PayloadHash != payloadHash || task.IdempotencyKey != idempotencyKey { return nil, false, fmt.Errorf("execution request does not match the immutable task payload") } @@ -2063,6 +2078,9 @@ func (s *DeliveryService) HandleStageEvent(ctx context.Context, taskID string, i if err := tx.First(&task, "id = ?", taskID).Error; err != nil { return err } + if !isMySQLServiceType(task.ServiceType) { + return fmt.Errorf("task service type %q is not handled by MySQL stage callbacks", task.ServiceType) + } if input.AWXJobID != "" { var execution model.ExecutionJob if err := tx.Where("task_id = ?", task.ID).First(&execution).Error; err != nil { @@ -2118,6 +2136,14 @@ func (s *DeliveryService) HandleAWXJobNotification(ctx context.Context, input AW if err != nil { return err } + var task model.DeliveryTask + if err := s.db.WithContext(ctx).Select("id", "service_type").First(&task, "id = ?", execution.TaskID).Error; err != nil { + return err + } + if !isMySQLServiceType(task.ServiceType) { + // PostgreSQL jobs are finalized by PostgreSQLDeliveryService.PollOnce. + return nil + } message := awxNotificationMessage(input) defer s.broadcastTask(ctx, execution.TaskID) switch status { @@ -2275,7 +2301,10 @@ func (s *DeliveryService) finishExecution(ctx context.Context, execution *model. func (s *DeliveryService) PollOnce(ctx context.Context) error { var jobs []model.ExecutionJob - if err := s.db.WithContext(ctx).Where("status = ?", "running").Find(&jobs).Error; err != nil { + if err := s.db.WithContext(ctx). + Joins("JOIN delivery_tasks ON delivery_tasks.id = execution_jobs.task_id"). + Where("execution_jobs.status = ? AND (delivery_tasks.service_type = '' OR delivery_tasks.service_type = ?)", "running", "mysql"). + Find(&jobs).Error; err != nil { return err } for _, execution := range jobs { @@ -2309,6 +2338,9 @@ func (s *DeliveryService) completeTask(ctx context.Context, taskID string) error if err := s.db.WithContext(ctx).First(&task, "id = ?", taskID).Error; err != nil { return err } + if !isMySQLServiceType(task.ServiceType) { + return fmt.Errorf("task service type %q cannot be completed by the MySQL delivery service", task.ServiceType) + } var payload deliveryPayload if err := json.Unmarshal([]byte(task.ImmutablePayload), &payload); err != nil { return err @@ -2380,6 +2412,9 @@ func (s *DeliveryService) RetryCloudDMRegistration(ctx context.Context, taskID s if err := s.db.WithContext(ctx).First(&task, "id = ?", taskID).Error; err != nil { return err } + if !isMySQLServiceType(task.ServiceType) { + return fmt.Errorf("task service type %q cannot use the MySQL CloudDM registration flow", task.ServiceType) + } if task.Status != model.TaskRegisterFailed { return fmt.Errorf("task %s is in state %q and cannot retry CloudDM registration", taskID, task.Status) } @@ -2657,6 +2692,9 @@ func (s *DeliveryService) beginRollback(ctx context.Context, taskID, reason stri if err := s.db.WithContext(ctx).First(&task, "id = ?", taskID).Error; err != nil { return err } + if !isMySQLServiceType(task.ServiceType) { + return fmt.Errorf("task service type %q cannot use the MySQL rollback flow", task.ServiceType) + } if rollbackProtectedStatus(task.Status) { return nil } @@ -2743,6 +2781,9 @@ func (s *DeliveryService) RetryRollback(ctx context.Context, taskID string) erro if err := s.db.WithContext(ctx).First(&task, "id = ?", taskID).Error; err != nil { return err } + if !isMySQLServiceType(task.ServiceType) { + return fmt.Errorf("task service type %q cannot use the MySQL rollback flow", task.ServiceType) + } if task.Status != model.TaskRollbackFailed { return fmt.Errorf("task %s is in state %q and cannot retry rollback", taskID, task.Status) } @@ -2770,6 +2811,9 @@ func (s *DeliveryService) AcknowledgeRollbackRelease(ctx context.Context, taskID if err := tx.First(&task, "id = ?", taskID).Error; err != nil { return err } + if !isMySQLServiceType(task.ServiceType) { + return fmt.Errorf("task service type %q cannot use the MySQL rollback flow", task.ServiceType) + } if task.Status != model.TaskRollbackFailed { return fmt.Errorf("task %s is in state %q and cannot acknowledge rollback release", taskID, task.Status) } @@ -2888,6 +2932,9 @@ func (s *DeliveryService) completeRollback(ctx context.Context, taskID string) e if err := tx.First(&task, "id = ?", taskID).Error; err != nil { return err } + if !isMySQLServiceType(task.ServiceType) { + return fmt.Errorf("task service type %q cannot use the MySQL rollback flow", task.ServiceType) + } if task.Status != model.TaskRollingBack { return fmt.Errorf("task %s is in state %q, cannot complete rollback", taskID, task.Status) } @@ -2907,6 +2954,10 @@ func (s *DeliveryService) completeRollback(ctx context.Context, taskID string) e }) } +func isMySQLServiceType(serviceType string) bool { + return serviceType == "" || serviceType == "mysql" +} + func (s *DeliveryService) failTask(ctx context.Context, task *model.DeliveryTask, status, message string) error { return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { var current model.DeliveryTask diff --git a/server/internal/service/delivery_test.go b/server/internal/service/delivery_test.go index 7a68fb3..006e2ce 100644 --- a/server/internal/service/delivery_test.go +++ b/server/internal/service/delivery_test.go @@ -293,6 +293,15 @@ func TestRegisterFailedIsProtectedFromRollback(t *testing.T) { } } +func TestMySQLServiceTypeCompatibility(t *testing.T) { + if !isMySQLServiceType("") || !isMySQLServiceType("mysql") { + t.Fatal("legacy and explicit MySQL tasks must remain supported") + } + if isMySQLServiceType(postgresqlServiceType) { + t.Fatal("PostgreSQL tasks must not enter MySQL execution paths") + } +} + func TestRollbackLaunchExpired(t *testing.T) { now := time.Date(2026, 7, 28, 12, 0, 0, 0, time.UTC) started := now.Add(-rollbackLaunchTimeout - time.Second) diff --git a/server/internal/service/postgresql_delivery.go b/server/internal/service/postgresql_delivery.go new file mode 100644 index 0000000..9f8a14a --- /dev/null +++ b/server/internal/service/postgresql_delivery.go @@ -0,0 +1,607 @@ +package service + +import ( + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "net" + "regexp" + "strconv" + "strings" + "time" + + "github.com/1024XEngineer/xinfra/server/internal/config" + "github.com/1024XEngineer/xinfra/server/internal/model" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +const ( + postgresqlPortPoolStart = 15432 + postgresqlPortPoolEnd = 15999 + postgresqlServiceType = "postgresql" +) + +var supportedPostgreSQLVersions = map[string]bool{"15": true, "16": true} +var supportedPostgreSQLTopologies = map[string]bool{"standalone": true, "primary_replica": true} +var postgresqlNamePattern = regexp.MustCompile(`^[a-z0-9](?:[-a-z0-9]*[a-z0-9])?$`) + +type PostgreSQLDeliveryInput struct { + BusinessLineID uint64 `json:"business_line_id" binding:"required"` + TargetID uint64 `json:"target_id" binding:"required"` + Namespace string `json:"namespace" binding:"required"` + ClusterName string `json:"cluster_name" binding:"required"` + VersionMajor string `json:"version_major"` + PostgreSQLVersion string `json:"postgresql_version,omitempty"` + Topology string `json:"topology"` + ReplicaCount int `json:"replica_count"` + TargetHosts []string `json:"target_hosts,omitempty"` + CPUMilli int64 `json:"cpu_milli" binding:"required"` + MemoryMi int64 `json:"memory_mi" binding:"required"` + StorageGi int64 `json:"storage_gi" binding:"required"` + DataRoot string `json:"data_root"` + MaxConnections int `json:"max_connections"` +} + +type postgresqlDeliveryPayload struct { + PostgreSQLDeliveryInput + TargetType string `json:"target_type"` +} + +func validatePostgreSQLDeliveryInput(input PostgreSQLDeliveryInput, dataDisks []string) error { + _ = dataDisks + if input.VersionMajor == "" { + input.VersionMajor = input.PostgreSQLVersion + } + if input.PostgreSQLVersion != "" && input.VersionMajor != "" && input.PostgreSQLVersion != input.VersionMajor { + return fmt.Errorf("version_major and postgresql_version must match") + } + if len(input.Namespace) > 63 || !dnsLabelPattern.MatchString(input.Namespace) { + return fmt.Errorf("namespace must be a valid Kubernetes DNS label") + } + if len(input.ClusterName) > 63 || !postgresqlNamePattern.MatchString(input.ClusterName) { + return fmt.Errorf("cluster_name must be a valid DNS label") + } + if !supportedPostgreSQLVersions[input.VersionMajor] { + return fmt.Errorf("unsupported version_major %q, supported: 15, 16", input.VersionMajor) + } + topology := input.Topology + if topology == "" { + topology = "standalone" + } + if !supportedPostgreSQLTopologies[topology] { + return fmt.Errorf("unsupported topology %q, supported: standalone, primary_replica", topology) + } + if topology == "standalone" && input.ReplicaCount != 0 { + return fmt.Errorf("standalone topology cannot have replicas") + } + if topology == "primary_replica" && (input.ReplicaCount < 1 || input.ReplicaCount > 7) { + return fmt.Errorf("replica_count must be between 1 and 7 for primary_replica") + } + if input.CPUMilli < 100 || input.CPUMilli > 64000 || input.MemoryMi < 2048 || input.MemoryMi > 65536 || input.StorageGi < 20 || input.StorageGi > 2000 { + return fmt.Errorf("requested resources are outside the supported range (memory: 2048-65536 MiB, storage: 20-2000 GiB)") + } + nodes := 1 + if topology == "primary_replica" { + nodes += input.ReplicaCount + } + if len(input.TargetHosts) > 0 && len(input.TargetHosts) != nodes { + return fmt.Errorf("target_hosts must contain exactly %d distinct hosts", nodes) + } + seen := map[string]bool{} + for _, host := range input.TargetHosts { + if len(host) > 253 || !hostNamePattern.MatchString(host) || seen[host] { + return fmt.Errorf("target_hosts must contain unique valid inventory host names") + } + seen[host] = true + } + if input.DataRoot != "" { + if input.DataRoot != "/data/postgresql" { + return fmt.Errorf("data_root is fixed to /data/postgresql in the first release") + } + } + if input.MaxConnections < 0 || input.MaxConnections > 10000 { + return fmt.Errorf("max_connections must be between 0 and 10000") + } + return nil +} + +func allocatePostgreSQLPort(requested int, used []int) (int, error) { + taken := make(map[int]bool, len(used)) + for _, port := range used { + taken[port] = true + } + if requested != 0 { + if requested < postgresqlPortPoolStart || requested > postgresqlPortPoolEnd { + return 0, fmt.Errorf("postgresql_port must be within %d-%d", postgresqlPortPoolStart, postgresqlPortPoolEnd) + } + if taken[requested] { + return 0, fmt.Errorf("postgresql port %d is already allocated on the target host", requested) + } + return requested, nil + } + for port := postgresqlPortPoolStart; port <= postgresqlPortPoolEnd; port++ { + if !taken[port] { + return port, nil + } + } + return 0, fmt.Errorf("postgresql port pool %d-%d is exhausted on the target host", postgresqlPortPoolStart, postgresqlPortPoolEnd) +} + +type postgresqlPortProbe func(context.Context, string, int) bool + +func allocateReachablePostgreSQLPort(ctx context.Context, host string, used []int, inUse postgresqlPortProbe) (int, error) { + occupied := append([]int(nil), used...) + for { + port, err := allocatePostgreSQLPort(0, occupied) + if err != nil { + return 0, err + } + if !inUse(ctx, host, port) { + return port, nil + } + occupied = append(occupied, port) + } +} + +func postgresqlPortInUse(ctx context.Context, host string, port int) bool { + probeCtx, cancel := context.WithTimeout(ctx, 500*time.Millisecond) + defer cancel() + conn, err := (&net.Dialer{}).DialContext(probeCtx, "tcp", net.JoinHostPort(host, strconv.Itoa(port))) + if err != nil { + return false + } + _ = conn.Close() + return true +} + +type PostgreSQLDeliveryService struct { + db *gorm.DB + cfg config.Config + awx *AWXClient + common *DeliveryService +} + +func NewPostgreSQLDeliveryService(cfg config.Config, db *gorm.DB, common *DeliveryService) *PostgreSQLDeliveryService { + return &PostgreSQLDeliveryService{db: db, cfg: cfg, awx: NewAWXClient(cfg.AWXBaseURL, cfg.AWXToken, cfg.AWXUsername, cfg.AWXPassword), common: common} +} + +func (s *PostgreSQLDeliveryService) CreateTask(ctx context.Context, userID uint64, isAdmin bool, idempotencyKey string, input PostgreSQLDeliveryInput) (*model.DeliveryTask, bool, error) { + idempotencyKey = strings.TrimSpace(idempotencyKey) + if idempotencyKey == "" || len(idempotencyKey) > 128 { + return nil, false, fmt.Errorf("Idempotency-Key header is required and must not exceed 128 characters") + } + if err := validatePostgreSQLDeliveryInput(input, s.cfg.DeliveryDataDisks); err != nil { + return nil, false, err + } + if input.VersionMajor == "" { + input.VersionMajor = input.PostgreSQLVersion + } + var existing model.DeliveryTask + if err := s.db.WithContext(ctx).Where("idempotency_key = ?", idempotencyKey).First(&existing).Error; err == nil { + if existing.RequestedBy != userID || existing.ServiceType != postgresqlServiceType { + return nil, false, fmt.Errorf("idempotency key is already in use") + } + return &existing, true, nil + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return nil, false, err + } + if !isAdmin { + var count int64 + if err := s.db.WithContext(ctx).Model(&model.BusinessLineUser{}).Where("business_line_id = ? AND user_id = ?", input.BusinessLineID, userID).Count(&count).Error; err != nil { + return nil, false, err + } + if count == 0 { + return nil, false, fmt.Errorf("user is not authorized for this business line") + } + } + if input.Topology == "" { + input.Topology = "standalone" + } + if input.DataRoot == "" { + input.DataRoot = "/data/postgresql" + } + target, err := getPostgreSQLTarget(ctx, s.awx, input.TargetID) + if err != nil { + return nil, false, err + } + payload := postgresqlDeliveryPayload{PostgreSQLDeliveryInput: input, TargetType: target.TargetType} + raw, err := json.Marshal(payload) + if err != nil { + return nil, false, err + } + digest := sha256.Sum256(raw) + task := model.DeliveryTask{ID: randomUUID(), BusinessLineID: input.BusinessLineID, RequestedBy: userID, Component: postgresqlServiceType, TargetType: target.TargetType, ServiceType: postgresqlServiceType, TargetID: input.TargetID, Namespace: input.Namespace, InstanceName: input.ClusterName, Status: model.TaskPending, ImmutablePayload: string(raw), PayloadHash: hex.EncodeToString(digest[:]), IdempotencyKey: idempotencyKey} + if err := s.db.WithContext(ctx).Create(&task).Error; err != nil { + if lookupErr := s.db.WithContext(ctx).Where("idempotency_key = ?", idempotencyKey).First(&existing).Error; lookupErr == nil { + return &existing, true, nil + } + return nil, false, err + } + _ = s.db.WithContext(ctx).Create(&model.TaskEvent{TaskID: task.ID, ToState: model.TaskPending, Message: "postgresql delivery task created"}).Error + return &task, false, nil +} + +func getPostgreSQLTarget(ctx context.Context, awx *AWXClient, templateID uint64) (DeliveryTarget, error) { + template, err := awx.GetJobTemplate(ctx, templateID) + if err != nil { + return DeliveryTarget{}, fmt.Errorf("deployment target is unavailable: %w", err) + } + if !strings.Contains(strings.ToLower(template.Name+" "+template.Description), "postgresql") { + return DeliveryTarget{}, fmt.Errorf("AWX job template %d is not a PostgreSQL target", templateID) + } + hosts, err := awx.ListInventoryHosts(ctx, template.Inventory) + if err != nil { + return DeliveryTarget{}, err + } + meta := targetMetadata{Topology: "standalone"} + for _, host := range hosts { + if host.Enabled { + meta.Hosts = append(meta.Hosts, targetHost{Name: host.Name, IP: AWXHostIP(host)}) + } + } + raw, err := json.Marshal(meta) + if err != nil { + return DeliveryTarget{}, err + } + return DeliveryTarget{ID: template.ID, Name: template.Name, TargetType: "host_pool", AWXInventoryID: template.Inventory, AWXTemplateID: template.ID, Enabled: true, Metadata: string(raw)}, nil +} + +func selectPostgreSQLHosts(hosts []targetHost, requested []string, count int) ([]targetHost, error) { + if len(requested) == 0 { + if len(hosts) < count { + return nil, fmt.Errorf("deployment target has %d hosts, but %d PostgreSQL instances are required", len(hosts), count) + } + return append([]targetHost(nil), hosts[:count]...), nil + } + byName := map[string]targetHost{} + for _, host := range hosts { + byName[host.Name] = host + } + selected := make([]targetHost, 0, len(requested)) + for _, name := range requested { + host, ok := byName[name] + if !ok { + return nil, fmt.Errorf("target host %q is not in the candidate host pool", name) + } + selected = append(selected, host) + } + return selected, nil +} + +func (s *PostgreSQLDeliveryService) claimAndReserve(ctx context.Context) (*model.DeliveryTask, error) { + var task model.DeliveryTask + var payload postgresqlDeliveryPayload + err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}).Where("status = ? AND service_type = ?", model.TaskPending, postgresqlServiceType).Order("created_at ASC").First(&task).Error; err != nil { + return err + } + if err := json.Unmarshal([]byte(task.ImmutablePayload), &payload); err != nil { + return s.common.transitionTx(tx, &task, model.TaskValidationFailed, "stored PostgreSQL deployment payload is invalid", "stored PostgreSQL deployment payload is invalid") + } + activeStates := []string{model.TaskValidating, model.TaskDispatching, model.TaskRunning, model.TaskCanceling} + checks := []struct { + query string + args []any + limit int + }{ + {"status IN ?", []any{activeStates}, s.cfg.DeliveryGlobalLimit}, + {"status IN ? AND target_id = ?", []any{activeStates, task.TargetID}, s.cfg.DeliveryTargetLimit}, + {"status IN ? AND business_line_id = ?", []any{activeStates, task.BusinessLineID}, s.cfg.DeliveryBusinessLimit}, + } + for _, check := range checks { + if check.limit <= 0 { + continue + } + var count int64 + if err := tx.Model(&model.DeliveryTask{}).Where(check.query, check.args...).Count(&count).Error; err != nil { + return err + } + if count >= int64(check.limit) { + return fmt.Errorf("defer: delivery concurrency limit reached") + } + } + target, err := getPostgreSQLTarget(ctx, s.awx, task.TargetID) + if err != nil { + return s.common.transitionTx(tx, &task, model.TaskValidationFailed, err.Error(), err.Error()) + } + meta := parseTargetMetadata(target.Metadata) + nodes := 1 + if payload.Topology == "primary_replica" { + nodes += payload.ReplicaCount + } + selected, err := selectPostgreSQLHosts(meta.Hosts, payload.TargetHosts, nodes) + if err != nil { + return s.common.transitionTx(tx, &task, model.TaskValidationFailed, err.Error(), err.Error()) + } + quotaOK, err := checkPostgreSQLResourceQuota(tx, task.BusinessLineID, task.TargetID, payload, int64(nodes)) + if err != nil { + return err + } + if !quotaOK { + return s.common.transitionTx(tx, &task, model.TaskValidationFailed, "resource quota is insufficient", "resource quota is insufficient") + } + var existingCluster model.PostgreSQLCluster + if err := tx.Where("name = ?", payload.ClusterName).First(&existingCluster).Error; err == nil { + return s.common.transitionTx(tx, &task, model.TaskValidationFailed, "PostgreSQL cluster name is already in use", "PostgreSQL cluster name is already in use") + } else if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + cluster := model.PostgreSQLCluster{TaskID: task.ID, BusinessLineID: task.BusinessLineID, TargetID: task.TargetID, Name: payload.ClusterName, VersionMajor: payload.VersionMajor, Topology: payload.Topology, ReplicationMode: "async", FailoverMode: "manual", Status: "provisioning", BackupStatus: "not_configured", MonitoringStatus: "not_configured"} + if err := tx.Create(&cluster).Error; err != nil { + return err + } + primaryPort := 0 + for index, host := range selected { + var usedPorts []int + if err := tx.Model(&model.PostgreSQLInstance{}).Where("hostname = ? AND status IN ?", host.Name, []string{"provisioning", "active", "quarantined"}).Pluck("port", &usedPorts).Error; err != nil { + return err + } + port, err := allocateReachablePostgreSQLPort(ctx, host.IP, usedPorts, postgresqlPortInUse) + if err != nil { + _ = tx.Delete(&cluster).Error + return s.common.transitionTx(tx, &task, model.TaskValidationFailed, err.Error(), err.Error()) + } + role := "replica" + instanceID := fmt.Sprintf("%s-replica-%d", payload.ClusterName, index) + upstream := payload.ClusterName + "-primary" + slot := fmt.Sprintf("xinfra_%s_replica_%d", strings.ReplaceAll(payload.ClusterName, "-", "_"), index) + if index == 0 { + role = "primary" + instanceID = payload.ClusterName + "-primary" + upstream = "" + slot = "" + primaryPort = port + } + if payload.Topology == "standalone" { + role = "standalone" + instanceID = payload.ClusterName + upstream = "" + slot = "" + } + root := strings.TrimRight(payload.DataRoot, "/") + "/" + instanceID + instance := model.PostgreSQLInstance{TaskID: task.ID, ClusterID: cluster.ID, BusinessLineID: task.BusinessLineID, TargetID: task.TargetID, InstanceID: instanceID, Hostname: host.Name, HostIP: host.IP, Port: port, DataDir: root + "/data", ConfigDir: root + "/conf", LogDir: root + "/log", SystemdUnit: "postgresql-xinfra@" + instanceID + ".service", VersionMajor: payload.VersionMajor, Role: role, HAComponentRole: "database", UpstreamInstanceID: upstream, ReplicationSlotName: slot, Status: "provisioning", BackupStatus: "not_configured", MonitoringStatus: "not_configured"} + if err := tx.Create(&instance).Error; err != nil { + return err + } + if index == 0 { + cluster.PrimaryInstanceID = instance.ID + } + } + if err := tx.Save(&cluster).Error; err != nil { + return err + } + reservation := model.ResourceReservation{TaskID: task.ID, BusinessLineID: task.BusinessLineID, TargetID: task.TargetID, CPUMilli: payload.CPUMilli * int64(nodes), MemoryMi: payload.MemoryMi * int64(nodes), StorageGi: payload.StorageGi * int64(nodes), InstanceCount: int64(nodes), Status: "reserved", ExpiresAt: time.Now().Add(time.Duration(s.cfg.ReservationTTLMinutes) * time.Minute)} + if err := tx.Create(&reservation).Error; err != nil { + return err + } + hostNames := make([]string, 0, len(selected)) + for _, host := range selected { + hostNames = append(hostNames, host.Name) + } + task.TargetHost = strings.Join(hostNames, ",") + task.TargetHostIP = selected[0].IP + task.PostgreSQLPort = primaryPort + if err := tx.Model(&model.DeliveryTask{}).Where("id = ?", task.ID).Updates(map[string]any{"target_host": task.TargetHost, "target_host_ip": task.TargetHostIP, "postgresql_port": primaryPort}).Error; err != nil { + return err + } + return s.common.transitionTx(tx, &task, model.TaskDispatching, "PostgreSQL resources, ports and directories reserved", "") + }) + return &task, err +} + +func checkPostgreSQLResourceQuota(tx *gorm.DB, businessLineID, targetID uint64, payload postgresqlDeliveryPayload, nodes int64) (bool, error) { + var quota model.ResourceQuota + if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("business_line_id = ? AND target_id = ?", businessLineID, targetID).First("a).Error; err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return true, nil + } + return false, err + } + type totals struct{ CPU, Memory, Storage, Instances int64 } + var mysqlUsed, postgresUsed, reserved totals + if err := tx.Model(&model.ResourceUsage{}).Select("COALESCE(SUM(cpu_milli),0) cpu, COALESCE(SUM(memory_mi),0) memory, COALESCE(SUM(storage_gi),0) storage, COALESCE(SUM(instance_count),0) instances").Where("business_line_id = ? AND target_id = ? AND status = ?", businessLineID, targetID, "active").Scan(&mysqlUsed).Error; err != nil { + return false, err + } + if err := tx.Model(&model.PostgreSQLResourceUsage{}).Select("COALESCE(SUM(cpu_milli),0) cpu, COALESCE(SUM(memory_mi),0) memory, COALESCE(SUM(storage_gi),0) storage, COUNT(*) instances").Where("business_line_id = ? AND target_id = ? AND status = ?", businessLineID, targetID, "active").Scan(&postgresUsed).Error; err != nil { + return false, err + } + if err := tx.Model(&model.ResourceReservation{}).Select("COALESCE(SUM(cpu_milli),0) cpu, COALESCE(SUM(memory_mi),0) memory, COALESCE(SUM(storage_gi),0) storage, COALESCE(SUM(instance_count),0) instances").Where("business_line_id = ? AND target_id = ? AND status = ? AND expires_at > ?", businessLineID, targetID, "reserved", time.Now()).Scan(&reserved).Error; err != nil { + return false, err + } + requestedCPU := payload.CPUMilli * nodes + requestedMemory := payload.MemoryMi * nodes + requestedStorage := payload.StorageGi * nodes + return mysqlUsed.CPU+postgresUsed.CPU+reserved.CPU+requestedCPU <= quota.CPUMilli && + mysqlUsed.Memory+postgresUsed.Memory+reserved.Memory+requestedMemory <= quota.MemoryMi && + mysqlUsed.Storage+postgresUsed.Storage+reserved.Storage+requestedStorage <= quota.StorageGi && + mysqlUsed.Instances+postgresUsed.Instances+reserved.Instances+nodes <= quota.InstanceLimit, nil +} + +func (s *PostgreSQLDeliveryService) CreateExecution(ctx context.Context, task *model.DeliveryTask) error { + var payload postgresqlDeliveryPayload + if err := json.Unmarshal([]byte(task.ImmutablePayload), &payload); err != nil { + return err + } + var instances []model.PostgreSQLInstance + if err := s.db.WithContext(ctx).Where("task_id = ?", task.ID).Order("id ASC").Find(&instances).Error; err != nil { + return err + } + instanceVars := make(map[string]any, len(instances)) + for _, instance := range instances { + instanceVars[instance.Hostname] = map[string]any{"instance_id": instance.InstanceID, "port": instance.Port, "data_dir": instance.DataDir, "config_dir": instance.ConfigDir, "log_dir": instance.LogDir, "role": instance.Role, "replication_slot": instance.ReplicationSlotName} + } + now := time.Now() + execution := model.ExecutionJob{TaskID: task.ID, IdempotencyKey: task.IdempotencyKey, ExecutorJobID: "pending", Status: "launching", StartedAt: &now} + if err := s.db.WithContext(ctx).Create(&execution).Error; err != nil { + return err + } + target, err := getPostgreSQLTarget(ctx, s.awx, task.TargetID) + if err != nil { + return err + } + extraVars := map[string]any{"task_id": task.ID, "payload_hash": task.PayloadHash, "target_hosts": task.TargetHost, "cluster_name": payload.ClusterName, "postgresql_version": payload.VersionMajor, "postgresql_port": task.PostgreSQLPort, "topology": payload.Topology, "replica_count": payload.ReplicaCount, "postgresql_instances": instanceVars, "memory_mb": payload.MemoryMi, "storage_gb": payload.StorageGi} + if payload.MaxConnections > 0 { + extraVars["max_connections"] = payload.MaxConnections + } + job, err := s.awx.Launch(ctx, target.AWXTemplateID, AWXLaunchRequest{InventoryID: target.AWXInventoryID, Limit: task.TargetHost, ExtraVars: extraVars}) + if err != nil { + _ = s.db.WithContext(ctx).Model(&execution).Updates(map[string]any{"status": "failed", "finished_at": time.Now()}) + return err + } + return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if err := tx.Model(&execution).Updates(map[string]any{"executor_job_id": fmt.Sprint(job.ID), "status": "running"}).Error; err != nil { + return err + } + return s.common.transitionTx(tx, task, model.TaskRunning, "PostgreSQL AWX job started", "") + }) +} + +func (s *PostgreSQLDeliveryService) DispatchOnce(ctx context.Context) error { + task, err := s.claimAndReserve(ctx) + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) || strings.HasPrefix(err.Error(), "defer:") { + return nil + } + return err + } + if err := s.CreateExecution(ctx, task); err != nil { + return s.fail(ctx, task.ID, model.TaskExecutionFailed, err.Error()) + } + return nil +} + +func (s *PostgreSQLDeliveryService) PollOnce(ctx context.Context) error { + var jobs []model.ExecutionJob + if err := s.db.WithContext(ctx).Joins("JOIN delivery_tasks ON delivery_tasks.id = execution_jobs.task_id").Where("execution_jobs.status = ? AND delivery_tasks.service_type = ?", "running", postgresqlServiceType).Find(&jobs).Error; err != nil { + return err + } + for _, execution := range jobs { + job, err := s.awx.GetJob(ctx, execution.ExecutorJobID) + if err != nil { + // A temporary AWX failure does not mean the deployment failed. Retry on + // the next scheduler tick while the execution remains running. + continue + } + switch strings.ToLower(job.Status) { + case "pending", "waiting", "running", "new": + continue + case "canceled": + _ = s.fail(ctx, execution.TaskID, model.TaskCanceled, "AWX job was canceled") + case "successful": + if err := s.complete(ctx, execution.TaskID); err != nil { + _ = s.fail(ctx, execution.TaskID, model.TaskValidationFailed, err.Error()) + } + default: + _ = s.fail(ctx, execution.TaskID, model.TaskExecutionFailed, "AWX job finished with status "+job.Status) + } + } + return nil +} + +func (s *PostgreSQLDeliveryService) complete(ctx context.Context, taskID string) error { + var task model.DeliveryTask + if err := s.db.WithContext(ctx).First(&task, "id = ? AND service_type = ?", taskID, postgresqlServiceType).Error; err != nil { + return err + } + var payload postgresqlDeliveryPayload + if err := json.Unmarshal([]byte(task.ImmutablePayload), &payload); err != nil { + return err + } + var instances []model.PostgreSQLInstance + if err := s.db.WithContext(ctx).Where("task_id = ?", taskID).Order("id ASC").Find(&instances).Error; err != nil { + return err + } + if len(instances) == 0 { + return fmt.Errorf("PostgreSQL task has no planned instances") + } + for _, instance := range instances { + if err := postgresReady(ctx, net.JoinHostPort(instance.HostIP, fmt.Sprint(instance.Port))); err != nil { + return fmt.Errorf("postgresql health check failed for %s: %w", instance.InstanceID, err) + } + } + now := time.Now() + return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + if err := s.common.transitionTx(tx, &task, model.TaskRegistering, "PostgreSQL health checks passed", ""); err != nil { + return err + } + for _, instance := range instances { + if err := tx.Model(&instance).Updates(map[string]any{"status": "active", "version_full": instance.VersionMajor}).Error; err != nil { + return err + } + if err := tx.Create(&model.PostgreSQLResourceUsage{TaskID: task.ID, InstanceID: instance.ID, ClusterID: instance.ClusterID, BusinessLineID: instance.BusinessLineID, TargetID: instance.TargetID, CPUMilli: payload.CPUMilli, MemoryMi: payload.MemoryMi, StorageGi: payload.StorageGi, Port: instance.Port, DataDir: instance.DataDir, Status: "active"}).Error; err != nil { + return err + } + } + if err := tx.Model(&model.ResourceReservation{}).Where("task_id = ? AND status = ?", task.ID, "reserved").Update("status", "consumed").Error; err != nil { + return err + } + if err := tx.Model(&model.PostgreSQLCluster{}).Where("task_id = ?", task.ID).Update("status", "active").Error; err != nil { + return err + } + if err := tx.Model(&model.ExecutionJob{}).Where("task_id = ?", task.ID).Updates(map[string]any{"status": "successful", "finished_at": now}).Error; err != nil { + return err + } + return s.common.transitionTx(tx, &task, model.TaskFinished, "PostgreSQL delivery completed and recorded in the PostgreSQL resource ledger", "") + }) +} + +func postgresReady(ctx context.Context, address string) error { + dialer := net.Dialer{Timeout: 5 * time.Second} + conn, err := dialer.DialContext(ctx, "tcp", address) + if err != nil { + return err + } + return conn.Close() +} + +func (s *PostgreSQLDeliveryService) fail(ctx context.Context, taskID, status, message string) error { + return s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { + var task model.DeliveryTask + if err := tx.First(&task, "id = ? AND service_type = ?", taskID, postgresqlServiceType).Error; err != nil { + return err + } + if err := s.common.transitionTx(tx, &task, status, message, message); err != nil { + return err + } + reservationStatus := "released" + if status == model.TaskExecutionFailed || (status == model.TaskValidationFailed && strings.Contains(strings.ToLower(message), "health")) { + reservationStatus = "quarantined" + } + if err := tx.Model(&model.ResourceReservation{}).Where("task_id = ? AND status = ?", taskID, "reserved").Update("status", reservationStatus).Error; err != nil { + return err + } + instanceStatus := "failed" + if reservationStatus == "quarantined" { + instanceStatus = "quarantined" + } + if err := tx.Model(&model.PostgreSQLInstance{}).Where("task_id = ? AND status = ?", taskID, "provisioning").Update("status", instanceStatus).Error; err != nil { + return err + } + if err := tx.Model(&model.PostgreSQLCluster{}).Where("task_id = ? AND status = ?", taskID, "provisioning").Update("status", instanceStatus).Error; err != nil { + return err + } + return tx.Model(&model.ExecutionJob{}).Where("task_id = ? AND status IN ?", taskID, []string{"launching", "running"}).Updates(map[string]any{"status": "failed", "finished_at": time.Now()}).Error + }) +} + +func (s *PostgreSQLDeliveryService) Run(ctx context.Context) { + interval := time.Duration(s.cfg.DeliveryPollSeconds) * time.Second + if interval < time.Second { + interval = time.Second + } + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + _ = s.DispatchOnce(ctx) + _ = s.PollOnce(ctx) + } + } +} diff --git a/server/internal/service/postgresql_delivery_test.go b/server/internal/service/postgresql_delivery_test.go new file mode 100644 index 0000000..f5c79a0 --- /dev/null +++ b/server/internal/service/postgresql_delivery_test.go @@ -0,0 +1,104 @@ +package service + +import ( + "context" + "testing" +) + +func TestValidatePostgreSQLDeliveryInput(t *testing.T) { + valid := PostgreSQLDeliveryInput{BusinessLineID: 1, TargetID: 2, Namespace: "team-a", ClusterName: "orders-pg", VersionMajor: "16", Topology: "standalone", CPUMilli: 500, MemoryMi: 2048, StorageGi: 20} + if err := validatePostgreSQLDeliveryInput(valid, []string{"/data"}); err != nil { + t.Fatalf("valid standalone input rejected: %v", err) + } + aliasVersion := valid + aliasVersion.VersionMajor = "" + aliasVersion.PostgreSQLVersion = "15" + if err := validatePostgreSQLDeliveryInput(aliasVersion, []string{"/data"}); err != nil { + t.Fatalf("postgresql_version alias rejected: %v", err) + } + replicated := valid + replicated.VersionMajor = "15" + replicated.Topology = "primary_replica" + replicated.ReplicaCount = 2 + replicated.TargetHosts = []string{"pg-a", "pg-b", "pg-c"} + replicated.DataRoot = "/data/postgresql" + replicated.MaxConnections = 500 + if err := validatePostgreSQLDeliveryInput(replicated, []string{"/data"}); err != nil { + t.Fatalf("valid primary_replica input rejected: %v", err) + } + for name, mutate := range map[string]func(*PostgreSQLDeliveryInput){ + "unsupported version": func(in *PostgreSQLDeliveryInput) { in.VersionMajor = "14" }, + "bad topology": func(in *PostgreSQLDeliveryInput) { in.Topology = "patroni" }, + "standalone replicas": func(in *PostgreSQLDeliveryInput) { in.ReplicaCount = 1 }, + "missing replicas": func(in *PostgreSQLDeliveryInput) { in.Topology = "primary_replica" }, + "too many replicas": func(in *PostgreSQLDeliveryInput) { in.Topology = "primary_replica"; in.ReplicaCount = 8 }, + "duplicate hosts": func(in *PostgreSQLDeliveryInput) { + in.Topology = "primary_replica" + in.ReplicaCount = 1 + in.TargetHosts = []string{"pg-a", "pg-a"} + }, + "wrong host count": func(in *PostgreSQLDeliveryInput) { + in.Topology = "primary_replica" + in.ReplicaCount = 2 + in.TargetHosts = []string{"pg-a", "pg-b"} + }, + "bad root": func(in *PostgreSQLDeliveryInput) { in.DataRoot = "/tmp/postgresql" }, + "unsupported root disk": func(in *PostgreSQLDeliveryInput) { in.DataRoot = "/disk1/postgresql" }, + } { + input := valid + mutate(&input) + if err := validatePostgreSQLDeliveryInput(input, []string{"/data"}); err == nil { + t.Errorf("%s was accepted", name) + } + } +} + +func TestAllocatePostgreSQLPort(t *testing.T) { + if port, err := allocatePostgreSQLPort(0, nil); err != nil || port != postgresqlPortPoolStart { + t.Fatalf("expected pool start %d, got %d err=%v", postgresqlPortPoolStart, port, err) + } + if port, err := allocatePostgreSQLPort(0, []int{15432, 15433}); err != nil || port != 15434 { + t.Fatalf("expected 15434, got %d err=%v", port, err) + } + if _, err := allocatePostgreSQLPort(15432, []int{15432}); err == nil { + t.Fatal("occupied requested port was accepted") + } + if _, err := allocatePostgreSQLPort(5432, nil); err == nil { + t.Fatal("port outside the PostgreSQL pool was accepted") + } +} + +func TestAllocateReachablePostgreSQLPortSkipsListeningPorts(t *testing.T) { + probed := []int{} + port, err := allocateReachablePostgreSQLPort(context.Background(), "pg.example", nil, func(_ context.Context, host string, port int) bool { + if host != "pg.example" { + t.Fatalf("unexpected probe host %q", host) + } + probed = append(probed, port) + return port == postgresqlPortPoolStart + }) + if err != nil { + t.Fatalf("allocate reachable port: %v", err) + } + if port != postgresqlPortPoolStart+1 { + t.Fatalf("port = %d, want %d", port, postgresqlPortPoolStart+1) + } + if len(probed) != 2 || probed[0] != postgresqlPortPoolStart || probed[1] != postgresqlPortPoolStart+1 { + t.Fatalf("unexpected probes: %v", probed) + } +} + +func TestSelectPostgreSQLHosts(t *testing.T) { + hosts := []targetHost{{Name: "pg-a"}, {Name: "pg-b"}, {Name: "pg-c"}} + selected, err := selectPostgreSQLHosts(hosts, nil, 2) + if err != nil || len(selected) != 2 || selected[0].Name != "pg-a" || selected[1].Name != "pg-b" { + t.Fatalf("unexpected automatic selection: %+v err=%v", selected, err) + } + selected, err = selectPostgreSQLHosts(hosts, []string{"pg-c", "pg-a"}, 2) + if err != nil || selected[0].Name != "pg-c" || selected[1].Name != "pg-a" { + t.Fatalf("unexpected pinned selection: %+v err=%v", selected, err) + } + if _, err := selectPostgreSQLHosts(hosts, []string{"missing"}, 1); err == nil { + t.Fatal("host outside the pool was accepted") + } +}