feat(snapshot): PHASE 2 - Add entity snapshot for Frontend & Backend updates

- Created entity_snapshot.py helper module (~570 lines)
  - save_entity_snapshot() - Create snapshots with compaction
  - rollback_entity_from_snapshot() - Main rollback logic
  - _rollback_update() - UPDATE rollback for all entity types
  - _rollback_create() - CREATE rollback (entity deletion)
  - Feature flag support: ENTITY_SNAPSHOT_ENABLED (default: false)

- Integrated snapshot into Frontend update (frontend.py)
  - Capture full entity state before UPDATE
  - Create entity_snapshot metadata
  - Merge with pre_apply_snapshot for diff viewer
  - Store in config_versions.metadata JSONB

- Integrated snapshot into Backend update (backend.py)
  - Same snapshot pattern as Frontend
  - Works within transaction for atomicity
  - Preserves diff viewer compatibility

- Added feature flag to config.py
  - ENTITY_SNAPSHOT_ENABLED (environment variable)
  - Default: false (safe rollout)
  - Ready for Phase 7 gradual deployment

Next: WAF, SSL, Server update integration + Reject rollback logic
This commit is contained in:
taylanbakircioglu
2025-11-13 17:33:09 +03:00
parent ebcdf4174e
commit 481be91a4e
4 changed files with 602 additions and 10 deletions
+5 -1
View File
@@ -26,4 +26,8 @@ MANAGEMENT_BASE_URL = os.getenv("MANAGEMENT_BASE_URL", PUBLIC_URL) # Backward c
# Agent settings
AGENT_HEARTBEAT_TIMEOUT_SECONDS = 15
AGENT_CONFIG_SYNC_INTERVAL_SECONDS = 30
AGENT_CONFIG_SYNC_INTERVAL_SECONDS = 30
# Entity snapshot feature flag (for gradual rollout)
# IMPORTANT: Keep this FALSE in production until Phase 7 deployment
ENTITY_SNAPSHOT_ENABLED = os.getenv("ENTITY_SNAPSHOT_ENABLED", "false").lower() == "true"
+29 -3
View File
@@ -863,15 +863,41 @@ async def update_backend(backend_id: int, backend_update: BackendConfigUpdate, r
# Mark PENDING for UI before generating config
await conn.execute("UPDATE backends SET last_config_status = 'PENDING' WHERE id = $1", backend_id)
# 🆕 PHASE 2: Create entity snapshot for rollback
from utils.entity_snapshot import save_entity_snapshot
entity_snapshot_metadata = await save_entity_snapshot(
conn=conn,
entity_type="backend",
entity_id=backend_id,
old_values=existing_backend, # Full record from line 763
new_values=update_data, # Only updated fields
operation="UPDATE"
)
# Get pre-apply snapshot (for diff viewer)
old_config = await conn.fetchval("""
SELECT config_content FROM config_versions
WHERE cluster_id = $1 AND status = 'APPLIED' AND is_active = TRUE
ORDER BY created_at DESC LIMIT 1
""", cluster_id)
# Merge metadata: pre_apply_snapshot + entity_snapshot
metadata = {
"pre_apply_snapshot": old_config or "", # For diff viewer
**entity_snapshot_metadata # For rollback
}
# Create a PENDING config version
config_content = await generate_haproxy_config_for_cluster(cluster_id, conn=conn)
config_hash = hashlib.sha256(config_content.encode()).hexdigest()
version_name = f"backend-{backend_id}-update-{int(time.time())}"
await conn.execute("""
INSERT INTO config_versions (cluster_id, version_name, config_content, checksum, created_by, status)
VALUES ($1, $2, $3, $4, $5, 'PENDING')
""", cluster_id, version_name, config_content, config_hash, current_user['id'])
INSERT INTO config_versions (cluster_id, version_name, config_content, checksum, created_by, status, metadata)
VALUES ($1, $2, $3, $4, $5, 'PENDING', $6)
""", cluster_id, version_name, config_content, config_hash, current_user['id'],
json.dumps(metadata) if metadata else None)
await close_database_connection(conn)
+63 -6
View File
@@ -600,10 +600,10 @@ async def update_frontend(frontend_id: int, frontend: FrontendConfig, request: R
conn = await get_database_connection()
# Check if frontend exists and get current configuration
# 🆕 PHASE 2: Get FULL frontend record for snapshot (all fields)
# CRITICAL: We need ALL fields for rollback, not just SSL fields
existing = await conn.fetchrow("""
SELECT id, name, cluster_id, ssl_enabled, ssl_certificate_id, ssl_port, ssl_cert_path, ssl_cert, ssl_verify
FROM frontends WHERE id = $1
SELECT * FROM frontends WHERE id = $1
""", frontend_id)
if not existing:
await close_database_connection(conn)
@@ -709,14 +709,71 @@ async def update_frontend(frontend_id: int, frontend: FrontendConfig, request: R
# Get system admin user ID for created_by
admin_user_id = await conn.fetchval("SELECT id FROM users WHERE username = 'admin' LIMIT 1") or 1
# 🆕 PHASE 2: Create entity snapshot for rollback
from utils.entity_snapshot import save_entity_snapshot
# Prepare new values for snapshot (only changed fields)
new_values = {
"name": frontend.name,
"bind_address": frontend.bind_address,
"bind_port": frontend.bind_port,
"default_backend": frontend.default_backend,
"mode": frontend.mode,
"ssl_enabled": ssl_enabled,
"ssl_certificate_id": ssl_certificate_id,
"ssl_certificate_ids": ssl_cert_ids_json,
"ssl_port": ssl_port,
"ssl_cert_path": ssl_cert_path,
"ssl_cert": ssl_cert,
"ssl_verify": ssl_verify,
"acl_rules": json.dumps(frontend.acl_rules or []),
"redirect_rules": json.dumps(frontend.redirect_rules or []),
"use_backend_rules": json.dumps(frontend.use_backend_rules or []),
"request_headers": frontend.request_headers,
"response_headers": frontend.response_headers,
"options": filtered_options,
"tcp_request_rules": frontend.tcp_request_rules,
"timeout_client": frontend.timeout_client,
"timeout_http_request": frontend.timeout_http_request,
"rate_limit": frontend.rate_limit,
"compression": frontend.compression,
"log_separate": frontend.log_separate,
"monitor_uri": frontend.monitor_uri,
"cluster_id": frontend.cluster_id,
"maxconn": frontend.maxconn
}
entity_snapshot_metadata = await save_entity_snapshot(
conn=conn,
entity_type="frontend",
entity_id=frontend_id,
old_values=existing, # Full record from line 604
new_values=new_values,
operation="UPDATE"
)
# Get pre-apply snapshot (for diff viewer)
old_config = await conn.fetchval("""
SELECT config_content FROM config_versions
WHERE cluster_id = $1 AND status = 'APPLIED' AND is_active = TRUE
ORDER BY created_at DESC LIMIT 1
""", frontend.cluster_id)
# Merge metadata: pre_apply_snapshot + entity_snapshot
metadata = {
"pre_apply_snapshot": old_config or "", # For diff viewer
**entity_snapshot_metadata # For rollback
}
# Try with status field first, fallback to old behavior if field doesn't exist
try:
config_version_id = await conn.fetchval("""
INSERT INTO config_versions
(cluster_id, version_name, config_content, checksum, created_by, is_active, status)
VALUES ($1, $2, $3, $4, $5, FALSE, 'PENDING')
(cluster_id, version_name, config_content, checksum, created_by, is_active, status, metadata)
VALUES ($1, $2, $3, $4, $5, FALSE, 'PENDING', $6)
RETURNING id
""", frontend.cluster_id, version_name, config_content, config_hash, admin_user_id)
""", frontend.cluster_id, version_name, config_content, config_hash, admin_user_id,
json.dumps(metadata) if metadata else None)
logger.info(f"APPLY WORKFLOW: Created PENDING config version {version_name} for cluster {frontend.cluster_id}")
# Mark entity config status as PENDING for UI
+505
View File
@@ -0,0 +1,505 @@
"""
Entity Snapshot Module for HAProxy OpenManager
Bu modül, entity değişikliklerinin snapshot'ını alır ve reject edildiğinde
geri yüklenmesini sağlar.
Usage:
# UPDATE için snapshot al
old_entity = await conn.fetchrow("SELECT * FROM frontends WHERE id = 5")
snapshot_metadata = await save_entity_snapshot(
conn=conn,
entity_type="frontend",
entity_id=5,
old_values=old_entity,
new_values={"bind_port": 443},
operation="UPDATE"
)
# Config version'a ekle
await conn.execute(
"INSERT INTO config_versions (..., metadata) VALUES (..., $1)",
json.dumps(snapshot_metadata)
)
# Reject edildiğinde rollback yap
await rollback_entity_from_snapshot(conn, snapshot_metadata["entity_snapshot"])
Supported Operations:
- UPDATE: Var olan entity'nin field'larını eski değerlere döndür
- CREATE: Yeni oluşturulan entity'yi sil (bulk import için)
- UPDATE_RESTORE: Restore işlemi sırasında yapılan UPDATE'i geri al
Feature Flag:
ENTITY_SNAPSHOT_ENABLED=true|false
Default: false (güvenli başlangıç)
Author: Taylan Bakırcıoğlu
Date: 2025-01-13
Version: 1.0.0
"""
import json
import os
import logging
from typing import Dict, Any, Optional, List
from datetime import datetime
# Setup logger
logger = logging.getLogger(__name__)
# Feature flag for gradual rollout
ENTITY_SNAPSHOT_ENABLED = os.getenv("ENTITY_SNAPSHOT_ENABLED", "false").lower() == "true"
async def save_entity_snapshot(
conn,
entity_type: str,
entity_id: int,
old_values: Dict[str, Any],
new_values: Optional[Dict[str, Any]] = None,
operation: str = "UPDATE"
) -> Dict[str, Any]:
"""
Entity snapshot'ı metadata formatında hazırla.
Args:
conn: Database connection (asyncpg)
entity_type: Entity tipi ('frontend', 'backend', 'waf_rule', 'ssl_certificate', 'server')
entity_id: Entity ID
old_values: Eski değerler (asyncpg Record'dan dict)
new_values: Yeni değerler (opsiyonel, UPDATE için)
operation: İşlem tipi ('CREATE', 'UPDATE', 'DELETE', 'UPDATE_RESTORE')
Returns:
Dict containing entity_snapshot metadata:
{
"entity_snapshot": {
"entity_type": "frontend",
"entity_id": 5,
"operation": "UPDATE",
"timestamp": "2025-01-13T10:30:45Z",
"old_values": {...},
"new_values": {...}, # Sadece değişen alanlar (compaction)
"changed_fields": ["bind_port", "ssl_enabled"]
}
}
Example:
old_frontend = await conn.fetchrow("SELECT * FROM frontends WHERE id = 5")
metadata = await save_entity_snapshot(
conn=conn,
entity_type="frontend",
entity_id=5,
old_values=dict(old_frontend),
new_values={"bind_port": 443, "ssl_enabled": True},
operation="UPDATE"
)
"""
# Feature flag check
if not ENTITY_SNAPSHOT_ENABLED:
logger.debug(f"Entity snapshot disabled by feature flag for {entity_type} {entity_id}")
return {}
try:
# Convert asyncpg Record to dict if needed
if hasattr(old_values, '__iter__') and not isinstance(old_values, dict):
old_values = dict(old_values)
# Calculate changed fields for UPDATE operations
changed_fields = []
if operation in ["UPDATE", "UPDATE_RESTORE"] and new_values:
changed_fields = [
field for field, new_val in new_values.items()
if field in old_values and old_values[field] != new_val
]
# Build snapshot
snapshot = {
"entity_type": entity_type,
"entity_id": entity_id,
"operation": operation,
"timestamp": datetime.utcnow().isoformat() + "Z",
"old_values": old_values,
"changed_fields": changed_fields
}
# For UPDATE: add new_values (compaction - only changed fields)
if operation in ["UPDATE", "UPDATE_RESTORE"] and new_values:
# Only store changed fields to save space (compaction)
snapshot["new_values"] = {
field: new_values[field]
for field in changed_fields
}
logger.info(
f"✅ SNAPSHOT: Created for {entity_type} {entity_id} "
f"(operation={operation}, changed_fields={len(changed_fields)})"
)
return {"entity_snapshot": snapshot}
except Exception as e:
logger.error(f"❌ SNAPSHOT ERROR: Failed to create snapshot for {entity_type} {entity_id}: {e}")
# Return empty dict on error (graceful degradation)
return {}
async def rollback_entity_from_snapshot(
conn,
entity_snapshot: Dict[str, Any]
) -> bool:
"""
Entity'yi snapshot'taki eski değerlerine geri yükle.
Args:
conn: Database connection (asyncpg)
entity_snapshot: Snapshot dict (from metadata.entity_snapshot)
Returns:
True if rollback successful, False otherwise
Example:
metadata = json.loads(version['metadata'])
entity_snapshot = metadata.get('entity_snapshot')
if entity_snapshot:
success = await rollback_entity_from_snapshot(conn, entity_snapshot)
if success:
logger.info("Rollback successful")
else:
logger.warning("Rollback failed")
"""
if not ENTITY_SNAPSHOT_ENABLED:
logger.debug("Entity snapshot disabled, skipping rollback")
return False
entity_type = entity_snapshot.get("entity_type")
entity_id = entity_snapshot.get("entity_id")
operation = entity_snapshot.get("operation")
old_values = entity_snapshot.get("old_values")
if not all([entity_type, entity_id, operation, old_values]):
logger.warning(f"⚠️ ROLLBACK: Invalid snapshot data, skipping rollback")
return False
try:
if operation in ["UPDATE", "UPDATE_RESTORE"]:
# Entity'yi eski değerlerine geri yükle
success = await _rollback_update(conn, entity_type, entity_id, old_values)
return success
elif operation == "CREATE":
# Yeni oluşturulan entity'yi sil
# ÖNEMLİ: Sadece bulk import ile TAMAMEN YENİ oluşturulan entity'ler için!
success = await _rollback_create(conn, entity_type, entity_id)
return success
elif operation == "DELETE":
# Silinen entity'yi geri yükle
# NOT: Şu an soft-delete kullanıldığı için kullanılmıyor
success = await _rollback_delete(conn, entity_type, entity_id, old_values)
return success
else:
logger.warning(f"⚠️ ROLLBACK: Unknown operation '{operation}', skipping")
return False
except Exception as e:
logger.error(f"❌ ROLLBACK ERROR: Failed for {entity_type} {entity_id}: {e}")
return False
async def _rollback_update(
conn,
entity_type: str,
entity_id: int,
old_values: Dict[str, Any]
) -> bool:
"""
Entity'yi UPDATE öncesi haline döndür.
ÖNEMLİ: Bu fonksiyon var olan entity'lerin field'larını eski değerlere döndürür.
Entity'yi KESİNLİKLE silmez!
Args:
conn: Database connection
entity_type: Entity tipi
entity_id: Entity ID
old_values: Eski değerler (snapshot'tan)
Returns:
True if successful, False otherwise
"""
try:
if entity_type == "frontend":
# Frontend'i eski değerlerine geri yükle
await conn.execute("""
UPDATE frontends SET
name = $1, bind_address = $2, bind_port = $3,
default_backend = $4, mode = $5, ssl_enabled = $6,
ssl_certificate_id = $7, ssl_certificate_ids = $8, ssl_port = $9,
ssl_cert_path = $10, ssl_cert = $11, ssl_verify = $12,
acl_rules = $13, redirect_rules = $14, use_backend_rules = $15,
request_headers = $16, response_headers = $17, options = $18,
tcp_request_rules = $19, timeout_client = $20, timeout_http_request = $21,
rate_limit = $22, compression = $23, log_separate = $24,
monitor_uri = $25, maxconn = $26,
last_config_status = $27, updated_at = $28
WHERE id = $29
""",
old_values.get('name'),
old_values.get('bind_address'),
old_values.get('bind_port'),
old_values.get('default_backend'),
old_values.get('mode'),
old_values.get('ssl_enabled'),
old_values.get('ssl_certificate_id'),
old_values.get('ssl_certificate_ids'),
old_values.get('ssl_port'),
old_values.get('ssl_cert_path'),
old_values.get('ssl_cert'),
old_values.get('ssl_verify'),
old_values.get('acl_rules'),
old_values.get('redirect_rules'),
old_values.get('use_backend_rules'),
old_values.get('request_headers'),
old_values.get('response_headers'),
old_values.get('options'),
old_values.get('tcp_request_rules'),
old_values.get('timeout_client'),
old_values.get('timeout_http_request'),
old_values.get('rate_limit'),
old_values.get('compression'),
old_values.get('log_separate'),
old_values.get('monitor_uri'),
old_values.get('maxconn'),
old_values.get('last_config_status'),
old_values.get('updated_at'),
entity_id
)
logger.info(f"✅ ROLLBACK UPDATE: Frontend {entity_id} restored to previous state")
return True
elif entity_type == "backend":
# Backend'i eski değerlerine geri yükle
await conn.execute("""
UPDATE backends SET
name = $1, balance_method = $2, mode = $3,
health_check_uri = $4, health_check_interval = $5,
health_check_expected_status = $6, fullconn = $7,
cookie_name = $8, cookie_options = $9,
default_server_inter = $10, default_server_fall = $11,
default_server_rise = $12, request_headers = $13,
response_headers = $14, options = $15,
timeout_connect = $16, timeout_server = $17, timeout_queue = $18,
last_config_status = $19, updated_at = $20
WHERE id = $21
""",
old_values.get('name'),
old_values.get('balance_method'),
old_values.get('mode'),
old_values.get('health_check_uri'),
old_values.get('health_check_interval'),
old_values.get('health_check_expected_status'),
old_values.get('fullconn'),
old_values.get('cookie_name'),
old_values.get('cookie_options'),
old_values.get('default_server_inter'),
old_values.get('default_server_fall'),
old_values.get('default_server_rise'),
old_values.get('request_headers'),
old_values.get('response_headers'),
old_values.get('options'),
old_values.get('timeout_connect'),
old_values.get('timeout_server'),
old_values.get('timeout_queue'),
old_values.get('last_config_status'),
old_values.get('updated_at'),
entity_id
)
logger.info(f"✅ ROLLBACK UPDATE: Backend {entity_id} restored to previous state")
return True
elif entity_type == "waf_rule":
# WAF rule'u eski değerlerine geri yükle
await conn.execute("""
UPDATE waf_rules SET
name = $1, rule_type = $2, action = $3, priority = $4,
description = $5, is_active = $6, config = $7,
last_config_status = $8, updated_at = $9
WHERE id = $10
""",
old_values.get('name'),
old_values.get('rule_type'),
old_values.get('action'),
old_values.get('priority'),
old_values.get('description'),
old_values.get('is_active'),
old_values.get('config'),
old_values.get('last_config_status'),
old_values.get('updated_at'),
entity_id
)
logger.info(f"✅ ROLLBACK UPDATE: WAF rule {entity_id} restored to previous state")
return True
elif entity_type == "ssl_certificate":
# SSL certificate'i eski değerlerine geri yükle
await conn.execute("""
UPDATE ssl_certificates SET
name = $1, primary_domain = $2, certificate_content = $3,
private_key_content = $4, chain_content = $5,
expiration_date = $6, usage_type = $7,
last_config_status = $8, updated_at = $9
WHERE id = $10
""",
old_values.get('name'),
old_values.get('primary_domain'),
old_values.get('certificate_content'),
old_values.get('private_key_content'),
old_values.get('chain_content'),
old_values.get('expiration_date'),
old_values.get('usage_type'),
old_values.get('last_config_status'),
old_values.get('updated_at'),
entity_id
)
logger.info(f"✅ ROLLBACK UPDATE: SSL certificate {entity_id} restored to previous state")
return True
elif entity_type == "server":
# Server'ı eski değerlerine geri yükle
await conn.execute("""
UPDATE backend_servers SET
server_name = $1, backend_name = $2, ip_address = $3, port = $4,
weight = $5, maxconn = $6, check_enabled = $7,
check_inter = $8, check_fall = $9, check_rise = $10,
check_port = $11, server_options = $12, is_backup = $13,
is_active = $14, last_config_status = $15, updated_at = $16
WHERE id = $17
""",
old_values.get('server_name'),
old_values.get('backend_name'),
old_values.get('ip_address'),
old_values.get('port'),
old_values.get('weight'),
old_values.get('maxconn'),
old_values.get('check_enabled'),
old_values.get('check_inter'),
old_values.get('check_fall'),
old_values.get('check_rise'),
old_values.get('check_port'),
old_values.get('server_options'),
old_values.get('is_backup'),
old_values.get('is_active'),
old_values.get('last_config_status'),
old_values.get('updated_at'),
entity_id
)
logger.info(f"✅ ROLLBACK UPDATE: Server {entity_id} restored to previous state")
return True
else:
logger.warning(f"⚠️ ROLLBACK UPDATE: Unsupported entity type '{entity_type}'")
return False
except Exception as e:
logger.error(f"❌ ROLLBACK UPDATE ERROR: {entity_type} {entity_id}: {e}", exc_info=True)
return False
async def _rollback_create(
conn,
entity_type: str,
entity_id: int
) -> bool:
"""
Yeni oluşturulan entity'yi sil (reject edilen CREATE).
ÖNEMLİ: Bu fonksiyon SADECE bulk import ile TAMAMEN YENİ oluşturulan
entity'ler için kullanılır. Var olan entity'nin UPDATE'i için
KESİNLİKLE kullanılmaz!
Kullanım Senaryosu:
- Bulk import ile 5 yeni backend oluşturuldu
- Kullanıcı reject yaptı
- Bu fonksiyon 5 backend'i siler
Args:
conn: Database connection
entity_type: Entity tipi
entity_id: Entity ID (silinecek)
Returns:
True if successful, False otherwise
"""
try:
if entity_type == "frontend":
await conn.execute("DELETE FROM frontends WHERE id = $1", entity_id)
logger.info(f"✅ ROLLBACK CREATE: Deleted frontend {entity_id}")
return True
elif entity_type == "backend":
# Cascade delete: servers otomatik silinecek (foreign key)
await conn.execute("DELETE FROM backends WHERE id = $1", entity_id)
logger.info(f"✅ ROLLBACK CREATE: Deleted backend {entity_id} (+ cascade servers)")
return True
elif entity_type == "waf_rule":
await conn.execute("DELETE FROM waf_rules WHERE id = $1", entity_id)
logger.info(f"✅ ROLLBACK CREATE: Deleted WAF rule {entity_id}")
return True
elif entity_type == "ssl_certificate":
await conn.execute("DELETE FROM ssl_certificates WHERE id = $1", entity_id)
logger.info(f"✅ ROLLBACK CREATE: Deleted SSL certificate {entity_id}")
return True
elif entity_type == "server":
await conn.execute("DELETE FROM backend_servers WHERE id = $1", entity_id)
logger.info(f"✅ ROLLBACK CREATE: Deleted server {entity_id}")
return True
else:
logger.warning(f"⚠️ ROLLBACK CREATE: Unsupported entity type '{entity_type}'")
return False
except Exception as e:
logger.error(f"❌ ROLLBACK CREATE ERROR: {entity_type} {entity_id}: {e}", exc_info=True)
return False
async def _rollback_delete(
conn,
entity_type: str,
entity_id: int,
old_values: Dict[str, Any]
) -> bool:
"""
Silinen entity'yi geri yükle (reject edilen DELETE).
NOT: Şu an projede soft-delete kullanıldığı için bu fonksiyon
aktif olarak kullanılmıyor. Gelecekte hard-delete kullanılırsa
bu fonksiyon implement edilecek.
Future Implementation:
- Entity'yi INSERT ile geri yükle
- Foreign key'leri geri yükle
- Related entity'leri geri yükle
Args:
conn: Database connection
entity_type: Entity tipi
entity_id: Entity ID
old_values: Eski değerler (snapshot'tan)
Returns:
True if successful, False otherwise
"""
logger.warning(
f"⚠️ ROLLBACK DELETE: Not implemented yet for {entity_type} {entity_id}. "
"Currently using soft-delete (is_active=false). Hard-delete rollback is a future feature."
)
return False