feat: add comprehensive CRUD endpoints, cluster sync improvements, and firewall orchestration
CI / backend (push) Failing after 2s
CI / frontend (push) Failing after 30s

Add create endpoints for users, roles, tenants, projects, networks, subnets, and security rules with audit logging, implement commit_or_400 helper for IntegrityError handling with 409 responses, enhance cluster sync to populate nodes, workloads, and networks from provider inventory with last_sync_at tracking, add update/delete operations for IP addresses and policies with version tracking, implement IP
This commit is contained in:
2026-07-09 12:33:36 +02:00
parent dea27b5459
commit a911d36f34
15 changed files with 1267 additions and 36 deletions
+358 -13
View File
@@ -1,19 +1,30 @@
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,
Subnet,
Tenant,
User,
@@ -23,19 +34,31 @@ 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,
SubnetCreate,
SubnetRead,
TenantCreate,
TenantRead,
UserCreate,
UserRead,
WorkloadRead,
)
@@ -48,6 +71,14 @@ 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
@api_router.get("/dashboard")
def dashboard(_: CurrentUser, db: Session = Depends(get_db)) -> dict:
return {
@@ -70,6 +101,37 @@ 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()
@@ -80,23 +142,35 @@ def create_cluster(payload: ClusterCreate, user: CurrentUser, db: Session = Depe
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)
db.commit()
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")
def test_cluster(cluster_id: str, _: CurrentUser, db: Session = Depends(get_db)) -> dict:
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")
return {"cluster_id": cluster.id, "status": "configured", "message": "Provider connection is ready for live token validation."}
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")
@@ -113,10 +187,65 @@ async def sync_cluster(cluster_id: str, user: CurrentUser, db: Session = Depends
read_only=cluster.mode == "read_only",
)
)
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}"]))
db.commit()
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()}}
@@ -136,40 +265,134 @@ 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.get("/ipam/addresses", response_model=list[IpAddressRead])
def ipam_addresses(_: CurrentUser, db: Session = Depends(get_db)):
from app.models.domain import IpAddress
@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)):
from app.models.domain import IpAddress
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)
db.commit()
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()
@@ -179,12 +402,44 @@ def security_groups(_: CurrentUser, db: Session = Depends(get_db)) -> list[Secur
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)
db.commit()
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()
@@ -194,12 +449,44 @@ def policies(_: CurrentUser, db: Session = Depends(get_db)) -> list[Policy]:
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)
db.commit()
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)
@@ -211,6 +498,64 @@ async def firewall_preview(policy_id: str, user: CurrentUser, db: Session = Depe
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()
+101
View File
@@ -26,20 +26,85 @@ class UserRead(OrmModel):
is_active: bool
class UserCreate(BaseModel):
email: str
display_name: str
password: str = Field(min_length=12)
role_ids: list[str] = []
class RoleCreate(BaseModel):
name: str
permissions: list[str] = []
class RoleRead(OrmModel):
id: str
name: str
permissions: list[str]
class ClusterCreate(BaseModel):
name: str
api_url: str
api_token: str = Field(min_length=8)
provider: str = "proxmox"
mode: str = "read_only"
verify_tls: bool = True
class TenantCreate(BaseModel):
name: str
description: str | None = None
class ProjectCreate(BaseModel):
tenant_id: str
name: str
description: str | None = None
class NetworkCreate(BaseModel):
cluster_id: str
project_id: str | None = None
name: str
kind: str = "bridge"
vlan_id: int | None = None
mtu: int = 1500
gateway: str | None = None
dns: list[str] = []
dhcp_enabled: bool = False
tags: list[str] = []
description: str | None = None
class SubnetCreate(BaseModel):
network_id: str
cidr: str
gateway: str | None = None
dns: list[str] = []
dhcp_enabled: bool = False
class SecurityGroupCreate(BaseModel):
project_id: str | None = None
name: str
description: str | None = None
class SecurityRuleCreate(BaseModel):
security_group_id: str
direction: str = "ingress"
action: str = "allow"
protocol: str = "tcp"
source: str = "any"
destination: str = "any"
port: str | None = None
priority: int = 1000
logging: bool = False
description: str | None = None
class PolicyCreate(BaseModel):
project_id: str | None = None
name: str
@@ -54,6 +119,20 @@ class IpReservationCreate(BaseModel):
note: str | None = None
class ServiceCatalogCreate(BaseModel):
name: str
protocol: str
ports: str
editable: bool = True
class FirewallApplyRequest(BaseModel):
policy_id: str
cluster_id: str | None = None
confirm: bool = False
dry_run: bool = True
class ClusterRead(OrmModel):
id: str
name: str
@@ -139,6 +218,20 @@ class SecurityGroupRead(OrmModel):
description: str | None
class SecurityRuleRead(OrmModel):
id: str
security_group_id: str
direction: str
action: str
protocol: str
source: str
destination: str
port: str | None
priority: int
logging: bool
description: str | None
class PolicyRead(OrmModel):
id: str
project_id: str | None
@@ -149,6 +242,14 @@ class PolicyRead(OrmModel):
last_compiled: dict[str, Any] | None
class ServiceCatalogRead(OrmModel):
id: str
name: str
protocol: str
ports: str
editable: bool
class FirewallPreview(BaseModel):
policy_id: str
dry_run: bool = True