From 3c3645f8d55c01aeec95036521891c1d692450d4 Mon Sep 17 00:00:00 2001 From: Hungerdream <1710233908@qq.com> Date: Mon, 27 Jul 2026 15:16:57 +0800 Subject: [PATCH] feat(delivery): extend MySQL delivery parameters and tighten resource bounds - Add topology/port/data_disk, DB options and 8 advanced parameters with whitelist validation, rendered via extra_vars into the playbook - Hybrid port allocation over pool 13306-13999 - Tighten bounds: memory 2048-65536 MiB, storage 20-2000 GiB --- ansible/mysql-deploy.yml | 160 +++++++++++++--- ansible/templates/mysql-instance.cnf.j2 | 66 +++++-- server/.env.example | 2 + server/internal/config/config.go | 10 + server/internal/service/delivery.go | 227 +++++++++++++++++++++-- server/internal/service/delivery_test.go | 112 +++++++---- 6 files changed, 492 insertions(+), 85 deletions(-) diff --git a/ansible/mysql-deploy.yml b/ansible/mysql-deploy.yml index 24d2cf2..fb07795 100644 --- a/ansible/mysql-deploy.yml +++ b/ansible/mysql-deploy.yml @@ -1,73 +1,134 @@ --- -- name: Preflight native MySQL delivery +- name: Native MySQL delivery hosts: "{{ target_hosts }}" become: true gather_facts: true any_errors_fatal: true vars: + # --- identity --- mysql_instance: "{{ instance_name }}" mysql_version_value: "{{ mysql_version | default('8.0') }}" mysql_package_map: "8.0": "mysql-server=8.0.46-0ubuntu0.24.04.3" mysql_package_name: "{{ mysql_package_map[mysql_version_value] }}" - mysql_port_value: "{{ mysql_port | default(3307) | int }}" - mysql_memory_mb_value: "{{ memory_mb | default(2048) | int }}" - mysql_storage_gb_value: "{{ storage_gb | default(20) | int }}" + + # --- platform-allocated inputs (safe fallbacks when not passed in) --- + # 端口池 13306–13999:平台从池分配并传入;未传时兜底池首端口,占用探测在目标机执行。 + mysql_port_value: "{{ mysql_port | default(13306) | int }}" + # GR 组通信端口:仅 mgr_3 使用,平台成对分配;兜底为 SQL 端口 +10000。 + mysql_gr_port_value: "{{ gr_port | default((mysql_port | default(13306) | int) + 10000) | int }}" + # 数据盘挂载点:平台按白名单选定并传入;兜底 /data。 + mysql_data_disk: "{{ data_disk | default('/data') }}" + + # --- resource quota (Go→AWX 契约单位保持 MB/GB,不擅改) --- + mysql_memory_mb_value: "{{ memory_mb | default(4096) | int }}" + mysql_storage_gb_value: "{{ storage_gb | default(50) | int }}" + + # --- mount-point-agnostic layout: {data_disk}/mysql-delivery/{instance_id}/... --- + mysql_base_dir: "{{ mysql_data_disk }}/mysql-delivery/{{ mysql_instance }}" mysql_install_dir: "/opt/mysql-delivery/{{ mysql_instance }}" - mysql_data_dir: "/data/mysql-delivery/{{ mysql_instance }}/data" - mysql_log_dir: "/data/mysql-delivery/{{ mysql_instance }}/log" + mysql_data_dir: "{{ mysql_base_dir }}/data" + mysql_log_dir: "{{ mysql_base_dir }}/logs" + mysql_binlog_dir: "{{ mysql_base_dir }}/logs/binlog" + mysql_redo_dir: "{{ mysql_base_dir }}/logs/redo" + mysql_tmp_dir: "{{ mysql_base_dir }}/tmp" + # socket/pid 走 /run tmpfs(重启自动清理),由 systemd RuntimeDirectory 创建。 mysql_run_dir: "/run/mysql-delivery-{{ mysql_instance }}" mysql_config_file: "/etc/mysql/mysql-delivery/{{ mysql_instance }}.cnf" + + # --- database config (user form, with doc default baselines) --- + mysql_timezone: "{{ timezone | default('+08:00') }}" + mysql_lower_case_table_names: "{{ lower_case_table_names | default(1) | int }}" + mysql_character_set: "{{ character_set | default('utf8mb4') }}" + mysql_collation: "{{ collation | default('utf8mb4_general_ci') }}" + + # --- advanced params (overridable; default baseline assumes SSD for io_capacity) --- + mysql_flush_log_at_trx_commit: "{{ innodb_flush_log_at_trx_commit | default(1) | int }}" + mysql_sync_binlog: "{{ sync_binlog | default(1) | int }}" + mysql_io_capacity: "{{ innodb_io_capacity | default(2000) | int }}" + mysql_long_query_time: "{{ long_query_time | default(1) }}" + mysql_binlog_expire_logs_seconds: "{{ binlog_expire_logs_seconds | default(604800) | int }}" + mysql_max_binlog_size: "{{ max_binlog_size | default('256M') }}" + mgr_group_name: "{{ group_replication_group_name | default('aaaaaaaa-aaaa-aaaa-aaaa-aaaaaaaaaaaa') }}" + + # --- expected node count per topology --- + mysql_expected_hosts: "{{ {'standalone': 1, 'primary_replica': 2, 'mgr_3': 3}[topology | default('standalone')] }}" + + # --- secrets (injected via environment) --- mysql_root_password_value: "{{ lookup('ansible.builtin.env', 'XINFRA_MYSQL_ROOT_PASSWORD') }}" mysql_admin_password_value: "{{ lookup('ansible.builtin.env', 'XINFRA_MYSQL_ADMIN_PASSWORD') }}" + pre_tasks: - - name: Validate prototype parameters + - name: Validate delivery parameters ansible.builtin.assert: that: - - topology == 'standalone' + - topology in ['standalone', 'primary_replica', 'mgr_3'] - mysql_instance is match('^[a-z0-9][a-z0-9-]{0,62}$') - mysql_version_value in mysql_package_map - - (mysql_port_value | int) >= 1024 - - (mysql_port_value | int) <= 65535 - - (mysql_memory_mb_value | int) >= 1024 - - (mysql_memory_mb_value | int) <= 4096 - - (mysql_storage_gb_value | int) >= 10 - - (mysql_storage_gb_value | int) <= 100 + - mysql_data_disk is match('^/') + - (mysql_port_value | int) >= 13306 + - (mysql_port_value | int) <= 13999 + - (mysql_memory_mb_value | int) >= 2048 + - (mysql_memory_mb_value | int) <= 65536 + - (mysql_storage_gb_value | int) >= 20 + - (mysql_storage_gb_value | int) <= 2000 + - (mysql_lower_case_table_names | int) in [0, 1] - mysql_root_password_value | length >= 16 - mysql_admin_password_value | length >= 16 - fail_msg: The first machine prototype only supports safe standalone parameters + fail_msg: >- + Delivery parameters out of the supported target-state whitelist + (topology / version-package-map / port pool 13306-13999 / memory 2-64G / storage 20-2000G). + quiet: true no_log: true - name: Validate topology host count ansible.builtin.assert: that: - - ansible_play_hosts_all | length == 1 - fail_msg: The selected topology does not match the target host count + - (ansible_play_hosts_all | length | int) == (mysql_expected_hosts | int) + fail_msg: "topology={{ topology }} expects {{ mysql_expected_hosts }} host(s), got {{ ansible_play_hosts_all | length }}" run_once: true - - name: Check port ownership + - name: Probe target-host port occupancy (SQL + GR) ansible.builtin.shell: | set -o pipefail - if ss -lntH "sport = :{{ mysql_port_value }}" | grep -q .; then - systemctl is-active --quiet "mysql-delivery@{{ mysql_instance }}.service" - fi + for p in {{ mysql_probe_ports | join(' ') }}; do + if ss -lntH "sport = :${p}" | grep -q .; then + # already listening: only tolerated when owned by this instance service + if ! systemctl is-active --quiet "mysql-delivery@{{ mysql_instance }}.service"; then + echo "port ${p} already in use on target host" >&2 + exit 3 + fi + fi + done args: executable: /bin/bash + vars: + mysql_probe_ports: "{{ [mysql_port_value, mysql_gr_port_value] if topology == 'mgr_3' else [mysql_port_value] }}" changed_when: false - - name: Read currently available memory + - name: Read currently available memory (point-in-time guard) ansible.builtin.shell: awk '/^MemAvailable:/ { print int($2 / 1024) }' /proc/meminfo args: executable: /bin/bash register: mysql_available_memory changed_when: false - - name: Check available memory and disk + - name: Resolve data-disk mountpoint available space + ansible.builtin.set_fact: + mysql_mount_avail: >- + {{ (ansible_mounts | selectattr('mount', 'equalto', mysql_data_disk) + | map(attribute='size_available') | first) + | default(ansible_mounts | selectattr('mount', 'equalto', '/') + | map(attribute='size_available') | first) }} + + - name: Check available memory and data-disk space ansible.builtin.assert: that: - (mysql_available_memory.stdout | int) >= (mysql_memory_mb_value | int) - - (ansible_mounts | selectattr('mount', 'equalto', '/') | map(attribute='size_available') | first | int) >= (mysql_storage_gb_value | int) * 1073741824 - fail_msg: Target host does not have enough available memory or disk + - (mysql_mount_avail | int) >= (mysql_storage_gb_value | int) * 1073741824 + fail_msg: >- + Target host lacks memory or free space on {{ mysql_data_disk }}; + host-wide Σ-quota budgeting is the platform's responsibility, this is a last-resort guard. tasks: - name: Detect an existing Ubuntu MySQL package @@ -95,7 +156,7 @@ path: /etc/apparmor.d/usr.sbin.mysqld register: mysql_apparmor_profile - - name: Allow the delivery data and run paths in the MySQL AppArmor profile + - name: Authorize the delivery paths in the MySQL AppArmor profile ansible.builtin.copy: dest: /etc/apparmor.d/local/usr.sbin.mysqld owner: root @@ -103,8 +164,8 @@ mode: '0644' content: | # Managed by XINFRA MySQL delivery - grant per-instance native paths - /data/mysql-delivery/ r, - /data/mysql-delivery/** rwk, + {{ mysql_data_disk }}/mysql-delivery/ r, + {{ mysql_data_disk }}/mysql-delivery/** rwk, /run/mysql-delivery-*/ rw, /run/mysql-delivery-*/** rwk, when: mysql_apparmor_profile.stat.exists @@ -115,7 +176,7 @@ when: mysql_apparmor_profile.stat.exists and mysql_apparmor_local.changed changed_when: true - - name: Create instance directories + - name: Create instance directories (mount-point-agnostic layout) ansible.builtin.file: path: "{{ item.path }}" state: directory @@ -125,9 +186,12 @@ loop: - { path: /etc/mysql/mysql-delivery, owner: root, group: mysql, mode: '0750' } - { path: "{{ mysql_install_dir }}", owner: root, group: root, mode: '0755' } + - { path: "{{ mysql_base_dir }}", owner: mysql, group: mysql, mode: '0750' } - { path: "{{ mysql_data_dir }}", owner: mysql, group: mysql, mode: '0750' } - { path: "{{ mysql_log_dir }}", owner: mysql, group: mysql, mode: '0750' } - - { path: "{{ mysql_run_dir }}", owner: mysql, group: mysql, mode: '0755' } + - { path: "{{ mysql_binlog_dir }}", owner: mysql, group: mysql, mode: '0750' } + - { path: "{{ mysql_redo_dir }}", owner: mysql, group: mysql, mode: '0750' } + - { path: "{{ mysql_tmp_dir }}", owner: mysql, group: mysql, mode: '0750' } - name: Link native binaries into the instance directory ansible.builtin.file: @@ -146,6 +210,33 @@ {{ 'standalone' if topology == 'standalone' else ('primary' if ansible_play_hosts_all.index(inventory_hostname) == 0 else ('replica' if topology == 'primary_replica' else 'mgr')) }} + # host-wide unique: platform may pass explicit mysql_server_id; otherwise derive + # port + node index (ports are unique per host in the pool model). + mysql_server_id_value: >- + {{ mysql_server_id | default((mysql_port_value | int) + + (ansible_play_hosts_all.index(inventory_hostname))) | int }} + + - name: Resolve memory tier (GB) + ansible.builtin.set_fact: + mysql_memory_gb: "{{ ((mysql_memory_mb_value | int) // 1024) | int }}" + + - name: Resolve linkage-derived defaults (illustrative tiers, pending hardware calibration) + ansible.builtin.set_fact: + mysql_max_connections: >- + {{ (max_connections | int) if (max_connections is defined and (max_connections | string) != 'auto') + else (200 if (mysql_memory_gb | int) <= 2 + else 500 if (mysql_memory_gb | int) <= 4 + else 1000 if (mysql_memory_gb | int) <= 8 + else 2000 if (mysql_memory_gb | int) <= 16 + else 4000 if (mysql_memory_gb | int) <= 32 + else 8000 if (mysql_memory_gb | int) <= 64 + else 16000) }} + mysql_redo_capacity: >- + {{ innodb_redo_log_capacity if (innodb_redo_log_capacity is defined) + else ('128M' if (mysql_memory_gb | int) <= 4 + else '256M' if (mysql_memory_gb | int) <= 16 + else '512M' if (mysql_memory_gb | int) <= 32 + else '1G') }} - name: Write instance configuration ansible.builtin.template: @@ -229,6 +320,15 @@ port: "{{ mysql_port_value }}" timeout: 30 + - name: Note pending HA runtime orchestration + ansible.builtin.debug: + msg: >- + topology={{ topology }} deployed with HA-ready config (GTID/binlog/relay/GR settings in place), + but replication wiring (CHANGE REPLICATION SOURCE) and Group Replication bootstrap + (START GROUP_REPLICATION) are not yet automated — nodes start config-ready only. + when: topology != 'standalone' + run_once: true + handlers: - name: Restart MySQL delivery instance ansible.builtin.systemd_service: diff --git a/ansible/templates/mysql-instance.cnf.j2 b/ansible/templates/mysql-instance.cnf.j2 index a46dc54..c295279 100644 --- a/ansible/templates/mysql-instance.cnf.j2 +++ b/ansible/templates/mysql-instance.cnf.j2 @@ -1,35 +1,79 @@ +# Managed by XINFRA MySQL delivery - generated, do not edit by hand +# instance={{ mysql_instance }} version={{ mysql_version_value }} topology={{ topology }} +[client] +socket={{ mysql_run_dir }}/mysql.sock +port={{ mysql_port_value }} + [mysqld] user=mysql basedir=/usr datadir={{ mysql_data_dir }} +tmpdir={{ mysql_tmp_dir }} socket={{ mysql_run_dir }}/mysql.sock pid-file={{ mysql_run_dir }}/mysql.pid port={{ mysql_port_value }} bind-address=0.0.0.0 -mysqlx=0 -server-id={{ 101 + mysql_node_index | int }} +skip-name-resolve=ON +server-id={{ mysql_server_id_value }} +local-infile=OFF + +# --- charset / collation / timezone / identifier case --- +character-set-server={{ mysql_character_set }} +collation-server={{ mysql_collation }} +default-time-zone={{ mysql_timezone }} +lower-case-table-names={{ mysql_lower_case_table_names }} + +# --- error / slow log --- log-error={{ mysql_log_dir }}/error.log slow-query-log=ON slow-query-log-file={{ mysql_log_dir }}/slow.log -skip-name-resolve=ON -max-connections=100 -innodb-buffer-pool-size={{ ((mysql_memory_mb_value | int) * 55 / 100) | int }}M -innodb-log-file-size=128M -log-bin={{ mysql_log_dir }}/mysql-bin +long-query-time={{ mysql_long_query_time }} + +# --- binlog (PITR + replication base) --- +log-bin={{ mysql_binlog_dir }}/binlog binlog-format=ROW +sync-binlog={{ mysql_sync_binlog }} +max-binlog-size={{ mysql_max_binlog_size }} +binlog-expire-logs-seconds={{ mysql_binlog_expire_logs_seconds }} +{% if topology != 'standalone' %} gtid-mode=ON enforce-gtid-consistency=ON -relay-log={{ mysql_log_dir }}/relay-bin +relay-log={{ mysql_binlog_dir }}/relay-bin +{% endif %} + +# --- InnoDB core --- +innodb-buffer-pool-size={{ ((mysql_memory_mb_value | int) * 55 / 100) | int }}M +innodb-flush-method=O_DIRECT +innodb-flush-log-at-trx-commit={{ mysql_flush_log_at_trx_commit }} +innodb-io-capacity={{ mysql_io_capacity }} +max-connections={{ mysql_max_connections }} + +# --- redo log (version-sensitive) --- +innodb-log-group-home-dir={{ mysql_redo_dir }} +{% if mysql_version_value in ['8.0', '8.4'] %} +innodb-redo-log-capacity={{ mysql_redo_capacity }} +{% else %} +innodb-log-file-size={{ mysql_redo_capacity }} +{% endif %} + +# --- X Protocol disabled (version-sensitive) --- +{% if mysql_version_value in ['8.0', '8.4'] %} +mysqlx=0 +{% elif mysql_version_value == '5.7' %} +loose-mysqlx=0 +{% endif %} {% if topology == 'mgr_3' %} +# --- Group Replication (config-ready; runtime bootstrap handled out-of-band) --- plugin-load-add=group_replication.so transaction-write-set-extraction=XXHASH64 loose-group-replication-group-name={{ mgr_group_name }} loose-group-replication-start-on-boot=OFF -loose-group-replication-local-address={{ ansible_host | default(inventory_hostname) }}:{{ mysql_port_value | int + 10000 }} -loose-group-replication-group-seeds={% for host in ansible_play_hosts_all %}{{ hostvars[host].ansible_host | default(host) }}:{{ mysql_port_value | int + 10000 }}{% if not loop.last %},{% endif %}{% endfor %} +loose-group-replication-local-address={{ ansible_host | default(inventory_hostname) }}:{{ mysql_gr_port_value }} +loose-group-replication-group-seeds={% for host in ansible_play_hosts_all %}{{ hostvars[host].ansible_host | default(host) }}:{{ mysql_gr_port_value }}{% if not loop.last %},{% endif %}{% endfor %} + loose-group-replication-ip-allowlist={% for host in ansible_play_hosts_all %}{{ hostvars[host].ansible_host | default(host) }}{% if not loop.last %},{% endif %}{% endfor %} + loose-group-replication-single-primary-mode=ON loose-group-replication-enforce-update-everywhere-checks=OFF {% endif %} - diff --git a/server/.env.example b/server/.env.example index 7ce4835..5e6d1a5 100644 --- a/server/.env.example +++ b/server/.env.example @@ -22,6 +22,8 @@ DELIVERY_RESERVATION_TTL_MINUTES=120 DELIVERY_GLOBAL_LIMIT=2 DELIVERY_TARGET_LIMIT=2 DELIVERY_BUSINESS_LIMIT=1 +# 数据盘挂载点白名单(逗号分隔,第一项为默认值) +DELIVERY_DATA_DISKS=/data AWX_BASE_URL= AWX_TOKEN= AWX_USERNAME= diff --git a/server/internal/config/config.go b/server/internal/config/config.go index 98a10e7..30f46d6 100644 --- a/server/internal/config/config.go +++ b/server/internal/config/config.go @@ -70,6 +70,7 @@ type Config struct { DeliveryGlobalLimit int DeliveryTargetLimit int DeliveryBusinessLimit int + DeliveryDataDisks []string } func Load() Config { @@ -138,6 +139,7 @@ func Load() Config { DeliveryGlobalLimit: envInt("DELIVERY_GLOBAL_LIMIT", 2), DeliveryTargetLimit: envInt("DELIVERY_TARGET_LIMIT", 2), DeliveryBusinessLimit: envInt("DELIVERY_BUSINESS_LIMIT", 1), + DeliveryDataDisks: envCSV("DELIVERY_DATA_DISKS", "/data"), } } @@ -241,6 +243,14 @@ func envInt(key string, fallback int) int { return parsed } +func envCSV(key, fallback string) []string { + items := splitCSV(env(key, fallback)) + if len(items) == 0 { + items = splitCSV(fallback) + } + return items +} + func defaultPublicBaseURL(httpAddr string) string { addr := strings.TrimSpace(httpAddr) if addr == "" { diff --git a/server/internal/service/delivery.go b/server/internal/service/delivery.go index 1186212..1c6911c 100644 --- a/server/internal/service/delivery.go +++ b/server/internal/service/delivery.go @@ -30,9 +30,26 @@ type MySQLDeliveryInput struct { Namespace string `json:"namespace" binding:"required"` InstanceName string `json:"instance_name" binding:"required"` MySQLVersion string `json:"mysql_version"` + Topology string `json:"topology"` + MySQLPort int `json:"mysql_port"` + DataDisk string `json:"data_disk"` CPUMilli int64 `json:"cpu_milli" binding:"required"` MemoryMi int64 `json:"memory_mi" binding:"required"` StorageGi int64 `json:"storage_gi" binding:"required"` + // 数据库配置(选填,缺省由 playbook 基线兜底) + Timezone string `json:"timezone"` + LowerCaseTableNames *int `json:"lower_case_table_names"` + CharacterSet string `json:"character_set"` + Collation string `json:"collation"` + // 高级参数(选填,零值视为未设置) + MaxConnections string `json:"max_connections"` + InnodbRedoLogCapacity string `json:"innodb_redo_log_capacity"` + InnodbFlushLogAtTrxCommit *int `json:"innodb_flush_log_at_trx_commit"` + SyncBinlog *int `json:"sync_binlog"` + InnodbIOCapacity int `json:"innodb_io_capacity"` + LongQueryTime float64 `json:"long_query_time"` + BinlogExpireLogsSeconds int64 `json:"binlog_expire_logs_seconds"` + MaxBinlogSize string `json:"max_binlog_size"` } type deliveryPayload struct { @@ -90,6 +107,32 @@ func firstFreeHost(hosts []targetHost, occupied []string) *targetHost { return nil } +// 端口池 13306-13999:混合分配模型,用户留空时自动分配,可覆盖为池内指定端口。 +const ( + mysqlPortPoolStart = 13306 + mysqlPortPoolEnd = 13999 +) + +// allocatePort 在目标主机已占用端口集上做混合分配:指定端口验冲突,未指定则取池内首个空闲端口。 +func allocatePort(requested int, used []int) (int, error) { + taken := make(map[int]bool, len(used)) + for _, p := range used { + taken[p] = true + } + if requested != 0 { + if taken[requested] { + return 0, fmt.Errorf("mysql_port %d is already allocated on the target host", requested) + } + return requested, nil + } + for p := mysqlPortPoolStart; p <= mysqlPortPoolEnd; p++ { + if !taken[p] { + return p, nil + } + } + return 0, fmt.Errorf("mysql port pool %d-%d is exhausted on the target host", mysqlPortPoolStart, mysqlPortPoolEnd) +} + type DeliveryService struct { db *gorm.DB cfg config.Config @@ -170,7 +213,7 @@ func (s *DeliveryService) CreateTask(ctx context.Context, userID uint64, isAdmin if idempotencyKey == "" || len(idempotencyKey) > 128 { return nil, false, fmt.Errorf("Idempotency-Key header is required and must not exceed 128 characters") } - if err := validateDeliveryInput(input); err != nil { + if err := validateDeliveryInput(input, s.cfg.DeliveryDataDisks); err != nil { return nil, false, err } @@ -204,6 +247,16 @@ func (s *DeliveryService) CreateTask(ctx context.Context, userID uint64, isAdmin if input.MySQLVersion == "" { input.MySQLVersion = "8.0" } + if input.Topology == "" { + input.Topology = "standalone" + } + if input.DataDisk == "" { + if len(s.cfg.DeliveryDataDisks) > 0 { + input.DataDisk = s.cfg.DeliveryDataDisks[0] + } else { + input.DataDisk = "/data" + } + } payload := deliveryPayload{MySQLDeliveryInput: input, TargetType: target.TargetType} raw, err := json.Marshal(payload) if err != nil { @@ -233,21 +286,105 @@ func (s *DeliveryService) CreateTask(ctx context.Context, userID uint64, isAdmin return &task, false, nil } +// 版本白名单与 playbook 的 mysql_package_map 保持同步。 var supportedMySQLVersions = map[string]bool{"8.0": true} -func validateDeliveryInput(input MySQLDeliveryInput) error { +// 拓扑白名单:playbook 已支持 primary_replica/mgr_3 的配置渲染, +// 但调度器仍是单主机模型且复制编排未自动化,本期仅放开 standalone。 +var supportedTopologies = map[string]bool{"standalone": true} + +var supportedCharsets = map[string]bool{"utf8mb4": true, "utf8": true, "gbk": true, "latin1": true} + +// 高级参数档位白名单(与 docs/mysql-parameter-selection.md 保持一致) +var ( + supportedMaxConnections = map[string]bool{"auto": true, "200": true, "500": true, "1000": true, "2000": true, "4000": true, "8000": true, "16000": true} + supportedLogSizes = map[string]bool{"128M": true, "256M": true, "512M": true, "1G": true} + supportedIOCapacities = map[int]bool{200: true, 2000: true, 5000: true} + supportedLongQueryTimes = map[float64]bool{0.5: true, 1: true, 2: true, 5: true, 10: true} + supportedBinlogExpireSecs = map[int64]bool{86400: true, 259200: true, 604800: true, 1209600: true} +) + +// timezone 仅接受偏移量(±HH:MM)、SYSTEM 或命名时区(如 Asia/Shanghai)。 +var timezonePattern = regexp.MustCompile(`^([+-](0\d|1[0-4]):[0-5]\d|SYSTEM|[A-Za-z]+(?:/[A-Za-z0-9_+-]+)+)$`) + +func validateDeliveryInput(input MySQLDeliveryInput, dataDisks []string) error { if len(input.Namespace) > 63 || !dnsLabelPattern.MatchString(input.Namespace) { return fmt.Errorf("namespace must be a valid Kubernetes DNS label") } if len(input.InstanceName) > 63 || !dnsLabelPattern.MatchString(input.InstanceName) { return fmt.Errorf("instance_name must be a valid Kubernetes DNS label") } - if input.CPUMilli < 100 || input.CPUMilli > 64000 || input.MemoryMi < 1024 || input.MemoryMi > 4096 || input.StorageGi < 10 || input.StorageGi > 100 { - return fmt.Errorf("requested resources are outside the supported range (memory: 1024-4096 MiB, storage: 10-100 GiB)") + // 与文档目标态一致(memory 2048-65536 MiB / storage 20-2000 GiB),playbook assert 同步。 + 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)") } if input.MySQLVersion != "" && !supportedMySQLVersions[input.MySQLVersion] { return fmt.Errorf("unsupported mysql_version %q, supported: 8.0", input.MySQLVersion) } + if input.Topology != "" && !supportedTopologies[input.Topology] { + return fmt.Errorf("unsupported topology %q, supported: standalone (primary_replica/mgr_3 pending scheduler support)", input.Topology) + } + if input.MySQLPort != 0 && (input.MySQLPort < mysqlPortPoolStart || input.MySQLPort > mysqlPortPoolEnd) { + return fmt.Errorf("mysql_port must be left empty for auto allocation or within the pool %d-%d", mysqlPortPoolStart, mysqlPortPoolEnd) + } + if input.DataDisk != "" { + allowed := dataDisks + if len(allowed) == 0 { + allowed = []string{"/data"} + } + found := false + for _, disk := range allowed { + if input.DataDisk == disk { + found = true + break + } + } + if !found { + return fmt.Errorf("data_disk %q is not in the allowed mount point list %v", input.DataDisk, allowed) + } + } + if input.Timezone != "" && !timezonePattern.MatchString(input.Timezone) { + return fmt.Errorf("timezone must be an offset like +08:00, SYSTEM, or a named zone like Asia/Shanghai") + } + if input.LowerCaseTableNames != nil && *input.LowerCaseTableNames != 0 && *input.LowerCaseTableNames != 1 { + return fmt.Errorf("lower_case_table_names must be 0 or 1") + } + if input.CharacterSet != "" && !supportedCharsets[input.CharacterSet] { + return fmt.Errorf("unsupported character_set %q, supported: utf8mb4, utf8, gbk, latin1", input.CharacterSet) + } + if input.Collation != "" { + charset := input.CharacterSet + if charset == "" { + charset = "utf8mb4" + } + if !strings.HasPrefix(input.Collation, charset+"_") { + return fmt.Errorf("collation %q does not match character_set %q", input.Collation, charset) + } + } + if input.MaxConnections != "" && !supportedMaxConnections[input.MaxConnections] { + return fmt.Errorf("max_connections must be one of auto, 200, 500, 1000, 2000, 4000, 8000, 16000") + } + if input.InnodbRedoLogCapacity != "" && !supportedLogSizes[input.InnodbRedoLogCapacity] { + return fmt.Errorf("innodb_redo_log_capacity must be one of 128M, 256M, 512M, 1G") + } + if input.InnodbFlushLogAtTrxCommit != nil && (*input.InnodbFlushLogAtTrxCommit < 0 || *input.InnodbFlushLogAtTrxCommit > 2) { + return fmt.Errorf("innodb_flush_log_at_trx_commit must be 0, 1 or 2") + } + if input.SyncBinlog != nil && *input.SyncBinlog != 0 && *input.SyncBinlog != 1 { + return fmt.Errorf("sync_binlog must be 0 or 1") + } + if input.InnodbIOCapacity != 0 && !supportedIOCapacities[input.InnodbIOCapacity] { + return fmt.Errorf("innodb_io_capacity must be one of 200, 2000, 5000") + } + if input.LongQueryTime != 0 && !supportedLongQueryTimes[input.LongQueryTime] { + return fmt.Errorf("long_query_time must be one of 0.5, 1, 2, 5, 10") + } + if input.BinlogExpireLogsSeconds != 0 && !supportedBinlogExpireSecs[input.BinlogExpireLogsSeconds] { + return fmt.Errorf("binlog_expire_logs_seconds must be one of 86400, 259200, 604800, 1209600") + } + if input.MaxBinlogSize != "" && !supportedLogSizes[input.MaxBinlogSize] { + return fmt.Errorf("max_binlog_size must be one of 128M, 256M, 512M, 1G") + } return nil } @@ -364,16 +501,29 @@ func (s *DeliveryService) claimAndReserve(ctx context.Context) (*model.DeliveryT if host == nil { return fmt.Errorf("defer: no free host available on target") } + // 端口池混合分配:同主机已占端口 = 非终态任务分配端口 ∪ 存量 active 实例端口。 + var usedPorts []int + if err := tx.Model(&model.DeliveryTask{}).Where("target_host = ? AND status NOT IN ?", host.Name, occupiedExclude).Pluck("mysql_port", &usedPorts).Error; err != nil { + return err + } + var instancePorts []int + if err := tx.Model(&model.MySQLInstance{}).Where("node_name = ? AND status = ?", host.Name, "active").Pluck("port", &instancePorts).Error; err != nil { + return err + } + port, portErr := allocatePort(payload.MySQLPort, append(usedPorts, instancePorts...)) + if portErr != nil { + return s.failInTransaction(tx, &task, model.TaskValidationFailed, portErr.Error()) + } reservation := model.ResourceReservation{TaskID: task.ID, BusinessLineID: task.BusinessLineID, TargetID: task.TargetID, CPUMilli: payload.CPUMilli, MemoryMi: payload.MemoryMi, StorageGi: payload.StorageGi, InstanceCount: 1, Status: "reserved", ExpiresAt: time.Now().Add(time.Duration(s.cfg.ReservationTTLMinutes) * time.Minute)} if err := tx.Create(&reservation).Error; err != nil { return err } - if err := tx.Model(&model.DeliveryTask{}).Where("id = ?", task.ID).Updates(map[string]any{"target_host": host.Name, "target_host_ip": host.IP, "mysql_port": meta.MySQLPort}).Error; err != nil { + if err := tx.Model(&model.DeliveryTask{}).Where("id = ?", task.ID).Updates(map[string]any{"target_host": host.Name, "target_host_ip": host.IP, "mysql_port": port}).Error; err != nil { return err } task.TargetHost = host.Name task.TargetHostIP = host.IP - task.MySQLPort = meta.MySQLPort + task.MySQLPort = port if err := s.transitionTx(tx, &task, model.TaskDispatching, "resources reserved", ""); err != nil { return err } @@ -516,13 +666,7 @@ func (s *DeliveryService) CreateExecution(ctx context.Context, taskID, payloadHa if err := s.db.WithContext(ctx).Create(&execution).Error; err != nil { return nil, false, err } - job, err := s.awx.Launch(ctx, target.AWXTemplateID, AWXLaunchRequest{InventoryID: target.AWXInventoryID, Limit: task.TargetHost, ExtraVars: map[string]any{ - "task_id": task.ID, "payload_hash": task.PayloadHash, - "target_hosts": task.TargetHost, "topology": meta.Topology, - "instance_name": payload.InstanceName, "mysql_port": task.MySQLPort, - "memory_mb": payload.MemoryMi, "storage_gb": payload.StorageGi, - "mysql_version": payload.MySQLVersion, - }}) + job, err := s.awx.Launch(ctx, target.AWXTemplateID, AWXLaunchRequest{InventoryID: target.AWXInventoryID, Limit: task.TargetHost, ExtraVars: deliveryExtraVars(&task, payload, meta)}) if err != nil { _ = s.db.WithContext(ctx).Model(&execution).Updates(map[string]any{"status": "failed", "finished_at": time.Now()}) return nil, false, err @@ -538,6 +682,63 @@ func (s *DeliveryService) CreateExecution(ctx context.Context, taskID, payloadHa return &execution, false, nil } +// deliveryExtraVars 组装传给 playbook 的变量:必传项固定注入, +// 选填项仅在用户显式设置时下发,未设置时由 playbook 默认基线兜底。 +func deliveryExtraVars(task *model.DeliveryTask, payload deliveryPayload, meta targetMetadata) map[string]any { + topology := payload.Topology + if topology == "" { + // 存量任务的 ImmutablePayload 无 topology 字段,回退到 target metadata。 + topology = meta.Topology + } + vars := map[string]any{ + "task_id": task.ID, "payload_hash": task.PayloadHash, + "target_hosts": task.TargetHost, "topology": topology, + "instance_name": payload.InstanceName, "mysql_port": task.MySQLPort, + "memory_mb": payload.MemoryMi, "storage_gb": payload.StorageGi, + "mysql_version": payload.MySQLVersion, + } + if payload.DataDisk != "" { + vars["data_disk"] = payload.DataDisk + } + if payload.Timezone != "" { + vars["timezone"] = payload.Timezone + } + if payload.LowerCaseTableNames != nil { + vars["lower_case_table_names"] = *payload.LowerCaseTableNames + } + if payload.CharacterSet != "" { + vars["character_set"] = payload.CharacterSet + } + if payload.Collation != "" { + vars["collation"] = payload.Collation + } + if payload.MaxConnections != "" { + vars["max_connections"] = payload.MaxConnections + } + if payload.InnodbRedoLogCapacity != "" { + vars["innodb_redo_log_capacity"] = payload.InnodbRedoLogCapacity + } + if payload.InnodbFlushLogAtTrxCommit != nil { + vars["innodb_flush_log_at_trx_commit"] = *payload.InnodbFlushLogAtTrxCommit + } + if payload.SyncBinlog != nil { + vars["sync_binlog"] = *payload.SyncBinlog + } + if payload.InnodbIOCapacity != 0 { + vars["innodb_io_capacity"] = payload.InnodbIOCapacity + } + if payload.LongQueryTime != 0 { + vars["long_query_time"] = payload.LongQueryTime + } + if payload.BinlogExpireLogsSeconds != 0 { + vars["binlog_expire_logs_seconds"] = payload.BinlogExpireLogsSeconds + } + if payload.MaxBinlogSize != "" { + vars["max_binlog_size"] = payload.MaxBinlogSize + } + return vars +} + 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 { diff --git a/server/internal/service/delivery_test.go b/server/internal/service/delivery_test.go index bc18b82..508b46f 100644 --- a/server/internal/service/delivery_test.go +++ b/server/internal/service/delivery_test.go @@ -2,43 +2,93 @@ package service import "testing" +func intPtr(v int) *int { return &v } + func TestValidateDeliveryInput(t *testing.T) { - valid := MySQLDeliveryInput{BusinessLineID: 1, TargetID: 1, Namespace: "team-a", InstanceName: "mysql-01", CPUMilli: 500, MemoryMi: 1024, StorageGi: 10} - if err := validateDeliveryInput(valid); err != nil { + dataDisks := []string{"/data", "/disk1"} + valid := MySQLDeliveryInput{BusinessLineID: 1, TargetID: 1, Namespace: "team-a", InstanceName: "mysql-01", CPUMilli: 500, MemoryMi: 2048, StorageGi: 20} + if err := validateDeliveryInput(valid, dataDisks); err != nil { t.Fatalf("valid input rejected: %v", err) } - withVersion := valid - withVersion.MySQLVersion = "8.0" - if err := validateDeliveryInput(withVersion); err != nil { - t.Fatalf("valid input with version 8.0 rejected: %v", err) + full := valid + full.MySQLVersion = "8.0" + full.Topology = "standalone" + full.MySQLPort = 13306 + full.DataDisk = "/disk1" + full.MemoryMi = 8192 + full.StorageGi = 2000 + full.Timezone = "+08:00" + full.LowerCaseTableNames = intPtr(0) + full.CharacterSet = "utf8mb4" + full.Collation = "utf8mb4_general_ci" + full.MaxConnections = "auto" + full.InnodbRedoLogCapacity = "256M" + full.InnodbFlushLogAtTrxCommit = intPtr(2) + full.SyncBinlog = intPtr(0) + full.InnodbIOCapacity = 2000 + full.LongQueryTime = 0.5 + full.BinlogExpireLogsSeconds = 604800 + full.MaxBinlogSize = "512M" + if err := validateDeliveryInput(full, dataDisks); err != nil { + t.Fatalf("valid full input rejected: %v", err) } - for name, input := range map[string]MySQLDeliveryInput{ - "uppercase namespace": valid, - "bad instance": valid, - "too little memory": valid, - "too much memory": valid, - "too little storage": valid, - "too much storage": valid, - "unsupported version": valid, + namedZone := valid + namedZone.Timezone = "Asia/Shanghai" + if err := validateDeliveryInput(namedZone, dataDisks); err != nil { + t.Fatalf("named timezone rejected: %v", err) + } + for name, mutate := range map[string]func(*MySQLDeliveryInput){ + "uppercase namespace": func(in *MySQLDeliveryInput) { in.Namespace = "Team-A" }, + "bad instance": func(in *MySQLDeliveryInput) { in.InstanceName = "mysql_01" }, + "too little memory": func(in *MySQLDeliveryInput) { in.MemoryMi = 1024 }, + "too much memory": func(in *MySQLDeliveryInput) { in.MemoryMi = 131072 }, + "too little storage": func(in *MySQLDeliveryInput) { in.StorageGi = 10 }, + "too much storage": func(in *MySQLDeliveryInput) { in.StorageGi = 4000 }, + "unsupported version": func(in *MySQLDeliveryInput) { in.MySQLVersion = "5.7" }, + "unsupported topology": func(in *MySQLDeliveryInput) { in.Topology = "mgr_3" }, + "port below pool": func(in *MySQLDeliveryInput) { in.MySQLPort = 3307 }, + "port above pool": func(in *MySQLDeliveryInput) { in.MySQLPort = 14000 }, + "data disk not in list": func(in *MySQLDeliveryInput) { in.DataDisk = "/mnt/other" }, + "bad timezone": func(in *MySQLDeliveryInput) { in.Timezone = "UTC+8" }, + "bad lower case": func(in *MySQLDeliveryInput) { in.LowerCaseTableNames = intPtr(2) }, + "bad charset": func(in *MySQLDeliveryInput) { in.CharacterSet = "big5" }, + "collation mismatch": func(in *MySQLDeliveryInput) { in.CharacterSet = "gbk"; in.Collation = "utf8mb4_general_ci" }, + "collation vs default": func(in *MySQLDeliveryInput) { in.Collation = "gbk_chinese_ci" }, + "bad max connections": func(in *MySQLDeliveryInput) { in.MaxConnections = "300" }, + "bad redo capacity": func(in *MySQLDeliveryInput) { in.InnodbRedoLogCapacity = "2G" }, + "bad flush log": func(in *MySQLDeliveryInput) { in.InnodbFlushLogAtTrxCommit = intPtr(3) }, + "bad sync binlog": func(in *MySQLDeliveryInput) { in.SyncBinlog = intPtr(2) }, + "bad io capacity": func(in *MySQLDeliveryInput) { in.InnodbIOCapacity = 500 }, + "bad long query time": func(in *MySQLDeliveryInput) { in.LongQueryTime = 3 }, + "bad binlog expire": func(in *MySQLDeliveryInput) { in.BinlogExpireLogsSeconds = 3600 }, + "bad max binlog size": func(in *MySQLDeliveryInput) { in.MaxBinlogSize = "64M" }, } { - switch name { - case "uppercase namespace": - input.Namespace = "Team-A" - case "bad instance": - input.InstanceName = "mysql_01" - case "too little memory": - input.MemoryMi = 512 - case "too much memory": - input.MemoryMi = 8192 - case "too little storage": - input.StorageGi = 5 - case "too much storage": - input.StorageGi = 200 - case "unsupported version": - input.MySQLVersion = "5.7" - } - if err := validateDeliveryInput(input); err == nil { + input := valid + mutate(&input) + if err := validateDeliveryInput(input, dataDisks); err == nil { t.Errorf("%s was accepted", name) } } } + +func TestAllocatePort(t *testing.T) { + if port, err := allocatePort(0, nil); err != nil || port != mysqlPortPoolStart { + t.Fatalf("expected first pool port %d, got %d err=%v", mysqlPortPoolStart, port, err) + } + if port, err := allocatePort(0, []int{13306, 13307}); err != nil || port != 13308 { + t.Fatalf("expected 13308 skipping occupied, got %d err=%v", port, err) + } + if port, err := allocatePort(13400, []int{13306}); err != nil || port != 13400 { + t.Fatalf("expected requested port 13400, got %d err=%v", port, err) + } + if _, err := allocatePort(13306, []int{13306}); err == nil { + t.Fatal("requested occupied port was accepted") + } + used := make([]int, 0, mysqlPortPoolEnd-mysqlPortPoolStart+1) + for p := mysqlPortPoolStart; p <= mysqlPortPoolEnd; p++ { + used = append(used, p) + } + if _, err := allocatePort(0, used); err == nil { + t.Fatal("exhausted pool still allocated a port") + } +}