Add SEED_DEMO_DATA environment variable to control demo data population, enhance setup wizard with welcome screen and theme toggle, add smooth animations for wizard transitions, improve dashboard endpoint to return structured objects for last_syncs and faulty_nodes instead of raw models, implement comprehensive error handling in cluster sync with failed status tracking and audit logging, fix Proxmox provider
714 lines
29 KiB
Python
714 lines
29 KiB
Python
from datetime import datetime
|
|
import csv
|
|
import io
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException
|
|
from fastapi.responses import StreamingResponse
|
|
from sqlalchemy import func, select
|
|
from sqlalchemy.exc import IntegrityError
|
|
from sqlalchemy.orm import Session
|
|
|
|
from app.api.deps import CurrentUser
|
|
from app.api.v1 import auth
|
|
from app.core.security import hash_password
|
|
from app.db.session import get_db
|
|
from app.models.domain import (
|
|
AuditLog,
|
|
Cluster,
|
|
IpAddress,
|
|
Job,
|
|
Network,
|
|
Node,
|
|
Policy,
|
|
Project,
|
|
Role,
|
|
SecurityGroup,
|
|
SecurityRule,
|
|
ServiceCatalogItem,
|
|
SystemSetting,
|
|
Subnet,
|
|
Tenant,
|
|
User,
|
|
Workload,
|
|
)
|
|
from app.schemas.domain import (
|
|
AuditLogRead,
|
|
ClusterCreate,
|
|
ClusterRead,
|
|
FirewallApplyRequest,
|
|
FirewallPreview,
|
|
IpAddressRead,
|
|
IpReservationCreate,
|
|
JobRead,
|
|
NetworkCreate,
|
|
NetworkRead,
|
|
NodeRead,
|
|
PolicyCreate,
|
|
PolicyRead,
|
|
ProjectCreate,
|
|
ProjectRead,
|
|
RoleCreate,
|
|
RoleRead,
|
|
SecurityRuleCreate,
|
|
SecurityRuleRead,
|
|
ServiceCatalogCreate,
|
|
ServiceCatalogRead,
|
|
SecurityGroupCreate,
|
|
SecurityGroupRead,
|
|
SetupCompleteRequest,
|
|
SetupStatus,
|
|
SubnetCreate,
|
|
SubnetRead,
|
|
TenantCreate,
|
|
TenantRead,
|
|
UserCreate,
|
|
UserRead,
|
|
WorkloadInsight,
|
|
WorkloadRead,
|
|
)
|
|
from app.services.audit import write_audit
|
|
from app.services.firewall_orchestrator import FirewallOrchestrator
|
|
from app.services.providers.base import ProviderConnection
|
|
from app.services.providers.registry import get_provider
|
|
|
|
api_router = APIRouter()
|
|
api_router.include_router(auth.router)
|
|
|
|
|
|
def commit_or_400(db: Session) -> None:
|
|
try:
|
|
db.commit()
|
|
except IntegrityError as exc:
|
|
db.rollback()
|
|
raise HTTPException(status_code=409, detail="Resource conflicts with an existing record") from exc
|
|
|
|
|
|
def setup_setting(db: Session) -> SystemSetting:
|
|
setting = db.get(SystemSetting, "setup")
|
|
if not setting:
|
|
setting = SystemSetting(key="setup", value={"complete": False})
|
|
db.add(setting)
|
|
db.commit()
|
|
db.refresh(setting)
|
|
return setting
|
|
|
|
|
|
@api_router.get("/setup/status", response_model=SetupStatus)
|
|
def setup_status(db: Session = Depends(get_db)) -> SetupStatus:
|
|
setting = setup_setting(db)
|
|
return SetupStatus(
|
|
complete=bool((setting.value or {}).get("complete")),
|
|
has_users=bool(db.scalar(select(func.count()).select_from(User))),
|
|
has_clusters=bool(db.scalar(select(func.count()).select_from(Cluster))),
|
|
)
|
|
|
|
|
|
@api_router.post("/setup/complete", response_model=SetupStatus)
|
|
def complete_setup(payload: SetupCompleteRequest, db: Session = Depends(get_db)) -> SetupStatus:
|
|
setting = setup_setting(db)
|
|
if bool((setting.value or {}).get("complete")):
|
|
raise HTTPException(status_code=409, detail="Setup has already been completed")
|
|
|
|
super_admin = db.scalar(select(Role).where(Role.name == "Super Admin"))
|
|
if not super_admin:
|
|
super_admin = Role(name="Super Admin", permissions=["*"])
|
|
db.add(super_admin)
|
|
db.flush()
|
|
|
|
email = payload.admin_email.strip().lower()
|
|
admin = db.scalar(select(User).where(User.email == email))
|
|
if not admin:
|
|
admin = User(email=email, display_name=payload.admin_name, password_hash=hash_password(payload.admin_password))
|
|
db.add(admin)
|
|
admin.display_name = payload.admin_name
|
|
admin.password_hash = hash_password(payload.admin_password)
|
|
admin.is_active = True
|
|
if super_admin not in admin.roles:
|
|
admin.roles.append(super_admin)
|
|
|
|
if payload.cluster_name and payload.cluster_api_url and payload.cluster_api_token:
|
|
existing_cluster = db.scalar(select(Cluster).where(Cluster.name == payload.cluster_name))
|
|
if not existing_cluster:
|
|
db.add(
|
|
Cluster(
|
|
name=payload.cluster_name,
|
|
api_url=payload.cluster_api_url,
|
|
token_ref=payload.cluster_api_token,
|
|
provider=payload.cluster_provider,
|
|
mode=payload.cluster_mode,
|
|
verify_tls=payload.verify_tls,
|
|
)
|
|
)
|
|
|
|
setting.value = {"complete": True, "completed_at": datetime.utcnow().isoformat()}
|
|
db.add(AuditLog(user_id=admin.id, action="setup.completed", object_type="system", result="success"))
|
|
commit_or_400(db)
|
|
return setup_status(db)
|
|
|
|
|
|
@api_router.get("/dashboard")
|
|
def dashboard(_: CurrentUser, db: Session = Depends(get_db)) -> dict:
|
|
last_syncs = db.scalars(select(Cluster).order_by(Cluster.updated_at.desc()).limit(5)).all()
|
|
faulty_nodes = db.scalars(select(Node).where(Node.status != "online")).all()
|
|
return {
|
|
"clusters": db.scalar(select(func.count()).select_from(Cluster)),
|
|
"nodes": db.scalar(select(func.count()).select_from(Node)),
|
|
"workloads": db.scalar(select(func.count()).select_from(Workload)),
|
|
"networks": db.scalar(select(func.count()).select_from(Network)),
|
|
"open_policy_violations": 1,
|
|
"last_syncs": [
|
|
{
|
|
"id": cluster.id,
|
|
"name": cluster.name,
|
|
"provider": cluster.provider,
|
|
"status": cluster.last_sync_status,
|
|
"error": cluster.last_sync_error,
|
|
"at": cluster.last_sync_at.isoformat() if cluster.last_sync_at else None,
|
|
}
|
|
for cluster in last_syncs
|
|
],
|
|
"faulty_nodes": [
|
|
{"id": node.id, "name": node.name, "status": node.status, "cluster_id": node.cluster_id}
|
|
for node in faulty_nodes
|
|
],
|
|
"top_talkers": [
|
|
{"name": "finance-app-2", "bytes": 942000000},
|
|
{"name": "core-services-1", "bytes": 512000000},
|
|
],
|
|
}
|
|
|
|
|
|
@api_router.get("/users", response_model=list[UserRead])
|
|
def users(_: CurrentUser, db: Session = Depends(get_db)) -> list[User]:
|
|
return db.scalars(select(User).order_by(User.email)).all()
|
|
|
|
|
|
@api_router.post("/users", response_model=UserRead)
|
|
def create_user(payload: UserCreate, user: CurrentUser, db: Session = Depends(get_db)) -> User:
|
|
roles = db.scalars(select(Role).where(Role.id.in_(payload.role_ids))).all() if payload.role_ids else []
|
|
new_user = User(
|
|
email=payload.email.strip().lower(),
|
|
display_name=payload.display_name,
|
|
password_hash=hash_password(payload.password),
|
|
roles=roles,
|
|
)
|
|
db.add(new_user)
|
|
commit_or_400(db)
|
|
db.refresh(new_user)
|
|
write_audit(db, action="user.created", object_type="user", object_id=new_user.id, user_id=user.id)
|
|
return new_user
|
|
|
|
|
|
@api_router.get("/roles", response_model=list[RoleRead])
|
|
def roles(_: CurrentUser, db: Session = Depends(get_db)) -> list[Role]:
|
|
return db.scalars(select(Role).order_by(Role.name)).all()
|
|
|
|
|
|
@api_router.post("/roles", response_model=RoleRead)
|
|
def create_role(payload: RoleCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Role:
|
|
role = Role(name=payload.name, permissions=payload.permissions)
|
|
db.add(role)
|
|
commit_or_400(db)
|
|
db.refresh(role)
|
|
write_audit(db, action="role.created", object_type="role", object_id=role.id, user_id=user.id)
|
|
return role
|
|
|
|
|
|
@api_router.get("/clusters", response_model=list[ClusterRead])
|
|
def clusters(_: CurrentUser, db: Session = Depends(get_db)) -> list[Cluster]:
|
|
return db.scalars(select(Cluster).order_by(Cluster.name)).all()
|
|
|
|
|
|
@api_router.post("/clusters", response_model=ClusterRead)
|
|
def create_cluster(payload: ClusterCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Cluster:
|
|
cluster = Cluster(
|
|
name=payload.name,
|
|
api_url=payload.api_url,
|
|
provider=payload.provider,
|
|
token_ref=payload.api_token,
|
|
mode=payload.mode,
|
|
verify_tls=payload.verify_tls,
|
|
)
|
|
db.add(cluster)
|
|
commit_or_400(db)
|
|
db.refresh(cluster)
|
|
write_audit(db, action="cluster.created", object_type="cluster", object_id=cluster.id, user_id=user.id)
|
|
return cluster
|
|
|
|
|
|
@api_router.post("/clusters/{cluster_id}/test")
|
|
async def test_cluster(cluster_id: str, _: CurrentUser, db: Session = Depends(get_db)) -> dict:
|
|
cluster = db.get(Cluster, cluster_id)
|
|
if not cluster:
|
|
raise HTTPException(status_code=404, detail="Cluster not found")
|
|
try:
|
|
result = await get_provider(cluster.provider).test_connection(
|
|
ProviderConnection(
|
|
api_url=cluster.api_url,
|
|
token=cluster.token_ref or "",
|
|
verify_tls=cluster.verify_tls,
|
|
read_only=cluster.mode == "read_only",
|
|
)
|
|
)
|
|
except Exception as exc:
|
|
raise HTTPException(status_code=502, detail=f"Provider connection failed: {exc}") from exc
|
|
return {"cluster_id": cluster.id, "status": "ok", "provider": cluster.provider, "result": result}
|
|
|
|
|
|
@api_router.post("/clusters/{cluster_id}/sync")
|
|
async def sync_cluster(cluster_id: str, user: CurrentUser, db: Session = Depends(get_db)) -> dict:
|
|
cluster = db.get(Cluster, cluster_id)
|
|
if not cluster:
|
|
raise HTTPException(status_code=404, detail="Cluster not found")
|
|
provider = get_provider(cluster.provider)
|
|
try:
|
|
inventory = await provider.sync_inventory(
|
|
ProviderConnection(
|
|
api_url=cluster.api_url,
|
|
token=cluster.token_ref or "",
|
|
verify_tls=cluster.verify_tls,
|
|
read_only=cluster.mode == "read_only",
|
|
)
|
|
)
|
|
except Exception as exc:
|
|
cluster.last_sync_at = datetime.utcnow()
|
|
cluster.last_sync_status = "failed"
|
|
cluster.last_sync_error = str(exc)
|
|
db.add(Job(kind="proxmox.sync", status="failed", progress=100, logs=[f"Sync failed for {cluster.name}"], error=str(exc)))
|
|
db.commit()
|
|
write_audit(
|
|
db,
|
|
action="cluster.sync",
|
|
object_type="cluster",
|
|
object_id=cluster.id,
|
|
user_id=user.id,
|
|
result="failed",
|
|
error_text=str(exc),
|
|
)
|
|
raise HTTPException(status_code=502, detail=f"Provider sync failed: {exc}") from exc
|
|
cluster.last_sync_at = datetime.utcnow()
|
|
cluster.last_sync_status = "success"
|
|
cluster.last_sync_error = None
|
|
node_by_name = {node.name: node for node in db.scalars(select(Node).where(Node.cluster_id == cluster.id)).all()}
|
|
for raw_node in inventory.get("nodes", []):
|
|
name = raw_node.get("node") or raw_node.get("name")
|
|
if not name:
|
|
continue
|
|
node = node_by_name.get(name)
|
|
if not node:
|
|
node = Node(cluster_id=cluster.id, name=name)
|
|
db.add(node)
|
|
node_by_name[name] = node
|
|
node.status = raw_node.get("status", node.status)
|
|
node.cpu_count = int(raw_node.get("maxcpu") or raw_node.get("cpu_count") or node.cpu_count or 0)
|
|
maxmem = raw_node.get("maxmem")
|
|
node.memory_mb = int(maxmem / 1024 / 1024) if isinstance(maxmem, int | float) else int(raw_node.get("memory_mb") or node.memory_mb or 0)
|
|
|
|
db.flush()
|
|
workload_by_external_id = {
|
|
workload.external_id: workload
|
|
for workload in db.scalars(select(Workload).where(Workload.cluster_id == cluster.id)).all()
|
|
}
|
|
for raw_workload in inventory.get("workloads", []):
|
|
external_id = str(raw_workload.get("vmid") or raw_workload.get("id") or "")
|
|
if not external_id:
|
|
continue
|
|
node_name = raw_workload.get("node")
|
|
node = node_by_name.get(node_name) or next(iter(node_by_name.values()), None)
|
|
if not node:
|
|
continue
|
|
workload = workload_by_external_id.get(external_id)
|
|
if not workload:
|
|
workload = Workload(cluster_id=cluster.id, node_id=node.id, external_id=external_id, name=external_id, kind="qemu")
|
|
db.add(workload)
|
|
workload_by_external_id[external_id] = workload
|
|
workload.node_id = node.id
|
|
workload.name = raw_workload.get("name") or workload.name
|
|
workload.kind = raw_workload.get("type") or raw_workload.get("kind") or workload.kind
|
|
workload.status = raw_workload.get("status") or workload.status
|
|
|
|
network_by_name = {
|
|
network.name: network
|
|
for network in db.scalars(select(Network).where(Network.cluster_id == cluster.id)).all()
|
|
}
|
|
for raw_network in inventory.get("networks", []):
|
|
name = raw_network.get("name") or raw_network.get("iface") or raw_network.get("id")
|
|
if not name:
|
|
continue
|
|
network = network_by_name.get(name)
|
|
if not network:
|
|
network = Network(cluster_id=cluster.id, name=name, kind=raw_network.get("type") or "network")
|
|
db.add(network)
|
|
network_by_name[name] = network
|
|
network.kind = raw_network.get("type") or raw_network.get("kind") or network.kind
|
|
vlan = raw_network.get("vlan") or raw_network.get("vlan_id")
|
|
network.vlan_id = int(vlan) if vlan not in (None, "") else network.vlan_id
|
|
db.add(Job(kind="proxmox.sync", status="success", progress=100, logs=[f"Synced {cluster.name}"]))
|
|
commit_or_400(db)
|
|
write_audit(db, action="cluster.sync", object_type="cluster", object_id=cluster.id, user_id=user.id)
|
|
return {"cluster_id": cluster.id, "status": "success", "inventory_counts": {key: len(value) for key, value in inventory.items()}}
|
|
|
|
|
|
@api_router.get("/nodes", response_model=list[NodeRead])
|
|
def nodes(_: CurrentUser, db: Session = Depends(get_db)) -> list[Node]:
|
|
return db.scalars(select(Node).order_by(Node.name)).all()
|
|
|
|
|
|
@api_router.get("/vms", response_model=list[WorkloadRead])
|
|
def workloads(_: CurrentUser, db: Session = Depends(get_db)) -> list[Workload]:
|
|
return db.scalars(select(Workload).order_by(Workload.name)).all()
|
|
|
|
|
|
@api_router.get("/vms/{workload_id}/insights", response_model=WorkloadInsight)
|
|
def workload_insights(workload_id: str, _: CurrentUser, db: Session = Depends(get_db)) -> WorkloadInsight:
|
|
workload = db.get(Workload, workload_id)
|
|
if not workload:
|
|
raise HTTPException(status_code=404, detail="Workload not found")
|
|
policies = db.scalars(
|
|
select(Policy).where((Policy.project_id == workload.project_id) | (Policy.project_id.is_(None))).order_by(Policy.name)
|
|
).all()
|
|
traffic = [
|
|
{
|
|
"timestamp": datetime.utcnow().isoformat(),
|
|
"source": workload.name,
|
|
"destination": "finance-db-1" if "web" in workload.tags else "core-services-1",
|
|
"protocol": "tcp",
|
|
"port": 5432 if "web" in workload.tags else 22,
|
|
"bytes": 1489200,
|
|
"decision": "allowed",
|
|
},
|
|
{
|
|
"timestamp": datetime.utcnow().isoformat(),
|
|
"source": "unknown-external",
|
|
"destination": workload.name,
|
|
"protocol": "tcp",
|
|
"port": 3389,
|
|
"bytes": 22140,
|
|
"decision": "would_block",
|
|
},
|
|
]
|
|
audit_mode_notes = [
|
|
f"{policy.name} is in audit mode; matching traffic is logged without enforcement."
|
|
for policy in policies
|
|
if policy.enforcement_mode == "audit"
|
|
]
|
|
decision = "audit" if audit_mode_notes else "allowed"
|
|
return WorkloadInsight(
|
|
workload=workload,
|
|
traffic=traffic,
|
|
matching_policies=policies,
|
|
effective_decision=decision,
|
|
audit_mode_notes=audit_mode_notes,
|
|
)
|
|
|
|
|
|
@api_router.get("/networks", response_model=list[NetworkRead])
|
|
def networks(_: CurrentUser, db: Session = Depends(get_db)) -> list[Network]:
|
|
return db.scalars(select(Network).order_by(Network.name)).all()
|
|
|
|
|
|
@api_router.post("/networks", response_model=NetworkRead)
|
|
def create_network(payload: NetworkCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Network:
|
|
if not db.get(Cluster, payload.cluster_id):
|
|
raise HTTPException(status_code=404, detail="Cluster not found")
|
|
network = Network(**payload.model_dump())
|
|
db.add(network)
|
|
commit_or_400(db)
|
|
db.refresh(network)
|
|
write_audit(db, action="network.created", object_type="network", object_id=network.id, user_id=user.id)
|
|
return network
|
|
|
|
|
|
@api_router.get("/ipam/subnets", response_model=list[SubnetRead])
|
|
def subnets(_: CurrentUser, db: Session = Depends(get_db)) -> list[Subnet]:
|
|
return db.scalars(select(Subnet).order_by(Subnet.cidr)).all()
|
|
|
|
|
|
@api_router.post("/ipam/subnets", response_model=SubnetRead)
|
|
def create_subnet(payload: SubnetCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Subnet:
|
|
if not db.get(Network, payload.network_id):
|
|
raise HTTPException(status_code=404, detail="Network not found")
|
|
subnet = Subnet(**payload.model_dump())
|
|
db.add(subnet)
|
|
commit_or_400(db)
|
|
db.refresh(subnet)
|
|
write_audit(db, action="ipam.subnet.created", object_type="subnet", object_id=subnet.id, user_id=user.id)
|
|
return subnet
|
|
|
|
|
|
@api_router.get("/ipam/addresses", response_model=list[IpAddressRead])
|
|
def ipam_addresses(_: CurrentUser, db: Session = Depends(get_db)) -> list[IpAddress]:
|
|
return db.scalars(select(IpAddress).order_by(IpAddress.address)).all()
|
|
|
|
|
|
@api_router.post("/ipam/addresses", response_model=IpAddressRead)
|
|
def reserve_ip(payload: IpReservationCreate, user: CurrentUser, db: Session = Depends(get_db)) -> IpAddress:
|
|
if not db.get(Subnet, payload.subnet_id):
|
|
raise HTTPException(status_code=404, detail="Subnet not found")
|
|
address = IpAddress(subnet_id=payload.subnet_id, address=payload.address, status=payload.status, note=payload.note)
|
|
db.add(address)
|
|
commit_or_400(db)
|
|
db.refresh(address)
|
|
write_audit(db, action="ipam.address.created", object_type="ip_address", object_id=address.id, user_id=user.id)
|
|
return address
|
|
|
|
|
|
@api_router.patch("/ipam/addresses/{address_id}", response_model=IpAddressRead)
|
|
def update_ip(address_id: str, payload: IpReservationCreate, user: CurrentUser, db: Session = Depends(get_db)) -> IpAddress:
|
|
address = db.get(IpAddress, address_id)
|
|
if not address:
|
|
raise HTTPException(status_code=404, detail="IP address not found")
|
|
old_values = {"address": address.address, "status": address.status, "note": address.note}
|
|
address.subnet_id = payload.subnet_id
|
|
address.address = payload.address
|
|
address.status = payload.status
|
|
address.note = payload.note
|
|
commit_or_400(db)
|
|
db.refresh(address)
|
|
write_audit(
|
|
db,
|
|
action="ipam.address.updated",
|
|
object_type="ip_address",
|
|
object_id=address.id,
|
|
user_id=user.id,
|
|
old_values=old_values,
|
|
new_values={"address": address.address, "status": address.status, "note": address.note},
|
|
)
|
|
return address
|
|
|
|
|
|
@api_router.delete("/ipam/addresses/{address_id}")
|
|
def delete_ip(address_id: str, user: CurrentUser, db: Session = Depends(get_db)) -> dict[str, str]:
|
|
address = db.get(IpAddress, address_id)
|
|
if not address:
|
|
raise HTTPException(status_code=404, detail="IP address not found")
|
|
db.delete(address)
|
|
commit_or_400(db)
|
|
write_audit(db, action="ipam.address.deleted", object_type="ip_address", object_id=address_id, user_id=user.id)
|
|
return {"status": "deleted", "id": address_id}
|
|
|
|
|
|
@api_router.get("/ipam/export.csv")
|
|
def export_ipam(_: CurrentUser, db: Session = Depends(get_db)) -> StreamingResponse:
|
|
buffer = io.StringIO()
|
|
writer = csv.writer(buffer)
|
|
writer.writerow(["subnet_id", "address", "status", "workload_id", "note"])
|
|
for address in db.scalars(select(IpAddress).order_by(IpAddress.address)):
|
|
writer.writerow([address.subnet_id, address.address, address.status, address.workload_id or "", address.note or ""])
|
|
buffer.seek(0)
|
|
return StreamingResponse(
|
|
iter([buffer.getvalue()]),
|
|
media_type="text/csv",
|
|
headers={"Content-Disposition": "attachment; filename=nexafabric-ipam.csv"},
|
|
)
|
|
|
|
|
|
@api_router.get("/tenants", response_model=list[TenantRead])
|
|
def tenants(_: CurrentUser, db: Session = Depends(get_db)) -> list[Tenant]:
|
|
return db.scalars(select(Tenant).order_by(Tenant.name)).all()
|
|
|
|
|
|
@api_router.post("/tenants", response_model=TenantRead)
|
|
def create_tenant(payload: TenantCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Tenant:
|
|
tenant = Tenant(**payload.model_dump())
|
|
db.add(tenant)
|
|
commit_or_400(db)
|
|
db.refresh(tenant)
|
|
write_audit(db, action="tenant.created", object_type="tenant", object_id=tenant.id, user_id=user.id)
|
|
return tenant
|
|
|
|
|
|
@api_router.get("/projects", response_model=list[ProjectRead])
|
|
def projects(_: CurrentUser, db: Session = Depends(get_db)) -> list[Project]:
|
|
return db.scalars(select(Project).order_by(Project.name)).all()
|
|
|
|
|
|
@api_router.post("/projects", response_model=ProjectRead)
|
|
def create_project(payload: ProjectCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Project:
|
|
if not db.get(Tenant, payload.tenant_id):
|
|
raise HTTPException(status_code=404, detail="Tenant not found")
|
|
project = Project(**payload.model_dump())
|
|
db.add(project)
|
|
commit_or_400(db)
|
|
db.refresh(project)
|
|
write_audit(db, action="project.created", object_type="project", object_id=project.id, user_id=user.id)
|
|
return project
|
|
|
|
|
|
@api_router.get("/security-groups", response_model=list[SecurityGroupRead])
|
|
def security_groups(_: CurrentUser, db: Session = Depends(get_db)) -> list[SecurityGroup]:
|
|
return db.scalars(select(SecurityGroup).order_by(SecurityGroup.name)).all()
|
|
|
|
|
|
@api_router.post("/security-groups", response_model=SecurityGroupRead)
|
|
def create_security_group(payload: SecurityGroupCreate, user: CurrentUser, db: Session = Depends(get_db)) -> SecurityGroup:
|
|
group = SecurityGroup(project_id=payload.project_id, name=payload.name, description=payload.description)
|
|
db.add(group)
|
|
commit_or_400(db)
|
|
db.refresh(group)
|
|
write_audit(db, action="security_group.created", object_type="security_group", object_id=group.id, user_id=user.id)
|
|
return group
|
|
|
|
|
|
@api_router.get("/security-groups/{group_id}/rules", response_model=list[SecurityRuleRead])
|
|
def security_group_rules(group_id: str, _: CurrentUser, db: Session = Depends(get_db)) -> list[SecurityRule]:
|
|
return db.scalars(
|
|
select(SecurityRule)
|
|
.where(SecurityRule.security_group_id == group_id)
|
|
.order_by(SecurityRule.priority, SecurityRule.created_at)
|
|
).all()
|
|
|
|
|
|
@api_router.post("/security-rules", response_model=SecurityRuleRead)
|
|
def create_security_rule(payload: SecurityRuleCreate, user: CurrentUser, db: Session = Depends(get_db)) -> SecurityRule:
|
|
if not db.get(SecurityGroup, payload.security_group_id):
|
|
raise HTTPException(status_code=404, detail="Security group not found")
|
|
rule = SecurityRule(**payload.model_dump())
|
|
db.add(rule)
|
|
commit_or_400(db)
|
|
db.refresh(rule)
|
|
write_audit(db, action="security_rule.created", object_type="security_rule", object_id=rule.id, user_id=user.id)
|
|
return rule
|
|
|
|
|
|
@api_router.delete("/security-rules/{rule_id}")
|
|
def delete_security_rule(rule_id: str, user: CurrentUser, db: Session = Depends(get_db)) -> dict[str, str]:
|
|
rule = db.get(SecurityRule, rule_id)
|
|
if not rule:
|
|
raise HTTPException(status_code=404, detail="Security rule not found")
|
|
db.delete(rule)
|
|
commit_or_400(db)
|
|
write_audit(db, action="security_rule.deleted", object_type="security_rule", object_id=rule_id, user_id=user.id)
|
|
return {"status": "deleted", "id": rule_id}
|
|
|
|
|
|
@api_router.get("/policies", response_model=list[PolicyRead])
|
|
def policies(_: CurrentUser, db: Session = Depends(get_db)) -> list[Policy]:
|
|
return db.scalars(select(Policy).order_by(Policy.name)).all()
|
|
|
|
|
|
@api_router.post("/policies", response_model=PolicyRead)
|
|
def create_policy(payload: PolicyCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Policy:
|
|
policy = Policy(project_id=payload.project_id, name=payload.name, enabled=payload.enabled, definition=payload.definition)
|
|
db.add(policy)
|
|
commit_or_400(db)
|
|
db.refresh(policy)
|
|
write_audit(db, action="policy.created", object_type="policy", object_id=policy.id, user_id=user.id)
|
|
return policy
|
|
|
|
|
|
@api_router.patch("/policies/{policy_id}", response_model=PolicyRead)
|
|
def update_policy(policy_id: str, payload: PolicyCreate, user: CurrentUser, db: Session = Depends(get_db)) -> Policy:
|
|
policy = db.get(Policy, policy_id)
|
|
if not policy:
|
|
raise HTTPException(status_code=404, detail="Policy not found")
|
|
old_values = {"name": policy.name, "enabled": policy.enabled, "definition": policy.definition, "version": policy.version}
|
|
policy.project_id = payload.project_id
|
|
policy.name = payload.name
|
|
policy.enabled = payload.enabled
|
|
policy.definition = payload.definition
|
|
policy.version += 1
|
|
commit_or_400(db)
|
|
db.refresh(policy)
|
|
write_audit(db, action="policy.updated", object_type="policy", object_id=policy.id, user_id=user.id, old_values=old_values, new_values=payload.model_dump())
|
|
return policy
|
|
|
|
|
|
@api_router.post("/policies/{policy_id}/compile", response_model=PolicyRead)
|
|
def compile_policy(policy_id: str, user: CurrentUser, db: Session = Depends(get_db)) -> Policy:
|
|
from app.services.policy_engine import PolicyEngine
|
|
|
|
policy = db.get(Policy, policy_id)
|
|
if not policy:
|
|
raise HTTPException(status_code=404, detail="Policy not found")
|
|
policy.last_compiled = PolicyEngine().compile(policy)
|
|
db.add(Job(kind="policy.compile", status="success", progress=100, logs=[f"Compiled policy {policy.name}"]))
|
|
commit_or_400(db)
|
|
db.refresh(policy)
|
|
write_audit(db, action="policy.compiled", object_type="policy", object_id=policy.id, user_id=user.id, new_values=policy.last_compiled)
|
|
return policy
|
|
|
|
|
|
@api_router.post("/firewall/preview/{policy_id}", response_model=FirewallPreview)
|
|
async def firewall_preview(policy_id: str, user: CurrentUser, db: Session = Depends(get_db)) -> FirewallPreview:
|
|
policy = db.get(Policy, policy_id)
|
|
cluster = db.scalar(select(Cluster).order_by(Cluster.name).limit(1))
|
|
if not policy or not cluster:
|
|
raise HTTPException(status_code=404, detail="Policy or cluster not found")
|
|
preview = await FirewallOrchestrator().preview(cluster, policy)
|
|
write_audit(db, action="firewall.preview", object_type="policy", object_id=policy.id, user_id=user.id, new_values=preview.model_dump())
|
|
return preview
|
|
|
|
|
|
@api_router.post("/firewall/apply")
|
|
async def firewall_apply(payload: FirewallApplyRequest, user: CurrentUser, db: Session = Depends(get_db)) -> dict:
|
|
if not payload.confirm:
|
|
raise HTTPException(status_code=400, detail="Firewall apply requires confirm=true after preview review")
|
|
policy = db.get(Policy, payload.policy_id)
|
|
cluster = db.get(Cluster, payload.cluster_id) if payload.cluster_id else db.scalar(select(Cluster).order_by(Cluster.name).limit(1))
|
|
if not policy or not cluster:
|
|
raise HTTPException(status_code=404, detail="Policy or cluster not found")
|
|
preview = await FirewallOrchestrator().preview(cluster, policy)
|
|
provider = get_provider(cluster.provider)
|
|
result = await provider.apply_rules(
|
|
ProviderConnection(
|
|
api_url=cluster.api_url,
|
|
token=cluster.token_ref or "",
|
|
verify_tls=cluster.verify_tls,
|
|
read_only=payload.dry_run or cluster.mode == "read_only",
|
|
),
|
|
preview.generated_rules,
|
|
)
|
|
job = Job(
|
|
kind="firewall.apply",
|
|
status="success" if result.get("applied") else "failed",
|
|
progress=100,
|
|
started_at=datetime.utcnow(),
|
|
finished_at=datetime.utcnow(),
|
|
logs=[f"Policy {policy.name}", f"Dry run: {payload.dry_run}", str(result)],
|
|
error=None if result.get("applied") else result.get("reason", "Provider did not apply rules"),
|
|
)
|
|
db.add(job)
|
|
commit_or_400(db)
|
|
write_audit(
|
|
db,
|
|
action="firewall.apply",
|
|
object_type="policy",
|
|
object_id=policy.id,
|
|
user_id=user.id,
|
|
new_values={"request": payload.model_dump(), "result": result},
|
|
result="success" if result.get("applied") else "blocked",
|
|
error_text=None if result.get("applied") else result.get("reason"),
|
|
)
|
|
return {"job_id": job.id, "preview": preview.model_dump(), "provider_result": result}
|
|
|
|
|
|
@api_router.get("/service-catalog", response_model=list[ServiceCatalogRead])
|
|
def service_catalog(_: CurrentUser, db: Session = Depends(get_db)) -> list[ServiceCatalogItem]:
|
|
return db.scalars(select(ServiceCatalogItem).order_by(ServiceCatalogItem.name)).all()
|
|
|
|
|
|
@api_router.post("/service-catalog", response_model=ServiceCatalogRead)
|
|
def create_service(payload: ServiceCatalogCreate, user: CurrentUser, db: Session = Depends(get_db)) -> ServiceCatalogItem:
|
|
service = ServiceCatalogItem(**payload.model_dump())
|
|
db.add(service)
|
|
commit_or_400(db)
|
|
db.refresh(service)
|
|
write_audit(db, action="service.created", object_type="service_catalog", object_id=service.id, user_id=user.id)
|
|
return service
|
|
|
|
|
|
@api_router.get("/jobs", response_model=list[JobRead])
|
|
def jobs(_: CurrentUser, db: Session = Depends(get_db)) -> list[Job]:
|
|
return db.scalars(select(Job).order_by(Job.created_at.desc())).all()
|
|
|
|
|
|
@api_router.get("/audit", response_model=list[AuditLogRead])
|
|
def audit(_: CurrentUser, db: Session = Depends(get_db)) -> list[AuditLog]:
|
|
return db.scalars(select(AuditLog).order_by(AuditLog.created_at.desc()).limit(200)).all()
|
|
|
|
|
|
@api_router.get("/settings")
|
|
def settings(_: CurrentUser) -> dict:
|
|
return {"product": "NexaFabric", "firewall_apply_requires_preview": True, "agent_optional": True}
|