Commit 99ee01c4 authored by Admin's avatar Admin

feat: implement 8 department agents with NocoBase integration

parent ae5bac2b
......@@ -17,6 +17,7 @@ from .routes.routines import router as routines_router
from .routes.approvals import router as approvals_router
from .routes.secrets import router as secrets_router
from .routes.costs import router as costs_router
from .routes.budgets import router as budgets_router
from .routes.activity import router as activity_router
from .routes.dashboard import router as dashboard_router
from .routes.heartbeat_runs import router as heartbeat_runs_router
......@@ -93,6 +94,7 @@ api_sub_router.include_router(routines_router)
api_sub_router.include_router(approvals_router)
api_sub_router.include_router(secrets_router)
api_sub_router.include_router(costs_router)
api_sub_router.include_router(budgets_router)
api_sub_router.include_router(activity_router)
api_sub_router.include_router(dashboard_router)
api_sub_router.include_router(heartbeat_runs_router)
......
This diff is collapsed.
"""
Budget management routes.
"""
from typing import Optional, List, Any, Dict
from fastapi import APIRouter, Depends, HTTPException, Query, status
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func
from datetime import datetime, timedelta
import uuid
from pydantic import BaseModel
from database import get_db_session
from middleware.auth import require_company_scope, get_current_user
from schemas.common import CompanyScope
from schemas.auth import AuthUserResponse
from models import BudgetPolicy, Agent, Project, Approval
router = APIRouter(prefix="/companies", tags=["budgets"])
class BudgetPolicyUpsertInput(BaseModel):
scopeType: str # "company", "project", "agent", "user"
scopeId: str
amount: int # monthly limit in cents
isActive: Optional[bool] = True
class BudgetIncidentResolutionInput(BaseModel):
action: str
amount: Optional[int] = None
decisionNote: Optional[str] = None
def map_policy_to_summary(policy: BudgetPolicy) -> Dict[str, Any]:
"""Map a BudgetPolicy model to BudgetPolicySummary structure."""
return {
"policyId": policy.id,
"companyId": policy.company_id,
"scopeType": policy.scope,
"scopeId": policy.scope_id,
"scopeName": policy.name,
"metric": "spend",
"windowKind": "monthly",
"amount": policy.monthly_limit_cents,
"observedAmount": policy.current_spend_cents,
"remainingAmount": max(0, policy.monthly_limit_cents - policy.current_spend_cents),
"utilizationPercent": (policy.current_spend_cents / policy.monthly_limit_cents * 100.0) if policy.monthly_limit_cents > 0 else 0.0,
"warnPercent": 80.0,
"hardStopEnabled": True,
"notifyEnabled": True,
"isActive": policy.is_enabled,
"status": "ok" if policy.current_spend_cents < policy.monthly_limit_cents else "hard_stop",
"paused": not policy.is_enabled or (policy.current_spend_cents >= policy.monthly_limit_cents),
"pauseReason": "budget" if (policy.current_spend_cents >= policy.monthly_limit_cents) else None,
"windowStart": policy.created_at.isoformat() + "Z",
"windowEnd": (policy.created_at + timedelta(days=30)).isoformat() + "Z"
}
@router.get("/{company_id}/budgets/overview")
async def get_budget_overview(
company_id: str,
scope: CompanyScope = Depends(require_company_scope),
db: AsyncSession = Depends(get_db_session),
):
"""Get budget policies and summary overview for a company."""
# 1. Fetch policies
result = await db.execute(
select(BudgetPolicy).where(BudgetPolicy.company_id == company_id)
)
policies = result.scalars().all()
mapped_policies = [map_policy_to_summary(p) for p in policies]
# 2. Count paused agents
agents_result = await db.execute(
select(func.count()).select_from(Agent).where(
Agent.company_id == company_id,
Agent.status == "paused"
)
)
paused_agents = agents_result.scalar() or 0
# 3. Pending budget approvals (represented by general unresolved approvals)
approvals_result = await db.execute(
select(func.count()).select_from(Approval).where(
Approval.company_id == company_id,
Approval.resolved_at.is_(None)
)
)
pending_approvals = approvals_result.scalar() or 0
return {
"companyId": company_id,
"policies": mapped_policies,
"activeIncidents": [],
"pausedAgentCount": paused_agents,
"pausedProjectCount": 0,
"pendingApprovalCount": pending_approvals,
}
@router.post("/{company_id}/budgets/policies")
async def upsert_budget_policy(
company_id: str,
payload: BudgetPolicyUpsertInput,
scope: CompanyScope = Depends(require_company_scope),
db: AsyncSession = Depends(get_db_session),
current_user: AuthUserResponse = Depends(get_current_user),
):
"""Create or update a budget policy."""
# Check if policy already exists
result = await db.execute(
select(BudgetPolicy).where(
BudgetPolicy.company_id == company_id,
BudgetPolicy.scope == payload.scopeType,
BudgetPolicy.scope_id == payload.scopeId
)
)
policy = result.scalar_one_or_none()
# Determine name dynamically if not specified
name = "Budget Policy"
if payload.scopeType == "company":
name = "Company Monthly Budget"
elif payload.scopeType == "agent":
agent_res = await db.execute(select(Agent).where(Agent.id == payload.scopeId))
agent = agent_res.scalar_one_or_none()
name = f"Agent {agent.name} Budget" if agent else "Agent Budget"
elif payload.scopeType == "project":
project_res = await db.execute(select(Project).where(Project.id == payload.scopeId))
project = project_res.scalar_one_or_none()
name = f"Project {project.name} Budget" if project else "Project Budget"
if policy:
# Update
policy.monthly_limit_cents = payload.amount
if payload.isActive is not None:
policy.is_enabled = payload.isActive
policy.name = name
policy.updated_at = datetime.utcnow()
else:
# Create
policy = BudgetPolicy(
id=uuid.uuid4().hex,
company_id=company_id,
name=name,
scope=payload.scopeType,
scope_id=payload.scopeId,
monthly_limit_cents=payload.amount,
current_spend_cents=0,
is_enabled=payload.isActive if payload.isActive is not None else True,
created_at=datetime.utcnow(),
updated_at=datetime.utcnow()
)
db.add(policy)
await db.commit()
await db.refresh(policy)
return map_policy_to_summary(policy)
@router.post("/{company_id}/budget-incidents/{incident_id}/resolve")
async def resolve_budget_incident(
company_id: str,
incident_id: str,
payload: BudgetIncidentResolutionInput,
scope: CompanyScope = Depends(require_company_scope),
db: AsyncSession = Depends(get_db_session),
current_user: AuthUserResponse = Depends(get_current_user),
):
"""Mock resolving a budget incident."""
# Since we don't store incidents in DB, we just mock success
return {
"id": incident_id,
"companyId": company_id,
"status": "resolved",
"resolvedAt": datetime.utcnow().isoformat() + "Z"
}
......@@ -12,7 +12,7 @@ from schemas.company import CompanyCreate, CompanyUpdate, CompanyResponse
from schemas.common import PaginationParams, PaginatedResponse
from schemas.auth import AuthUserResponse
from schemas.common import CompanyScope
from models import Company, CompanyMembership
from models import Company, CompanyMembership, Agent
from services.activity_logger import log_activity
router = APIRouter(prefix="/companies", tags=["companies"])
......@@ -169,3 +169,41 @@ async def update_company(
return company
@router.get("/{company_id}/org")
async def get_company_org(
company_id: str,
scope: CompanyScope = Depends(require_company_scope),
db: AsyncSession = Depends(get_db_session),
current_user: AuthUserResponse = Depends(get_current_user)
):
"""Get org tree structure of a company."""
# Fetch all agents for the company
result = await db.execute(
select(Agent).where(Agent.company_id == company_id)
)
agents = result.scalars().all()
# Build node dictionaries
nodes = {}
for agent in agents:
nodes[agent.id] = {
"id": agent.id,
"name": agent.name,
"role": agent.role,
"status": agent.status,
"reports": []
}
roots: List[dict] = []
for agent in agents:
node = nodes[agent.id]
if agent.reports_to and agent.reports_to in nodes:
parent_node = nodes[agent.reports_to]
parent_node["reports"].append(node)
else:
roots.append(node)
return roots
......@@ -36,3 +36,14 @@ class Agent(Base):
agent_metadata: Mapped[Optional[dict]] = mapped_column("metadata", JSON, nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow)
updated_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow)
@property
def urlKey(self) -> str:
"""Slugified version of agent name, falls back to agent id."""
import re
if not self.name:
return self.id
val = self.name.strip().lower()
val = re.sub(r'[^a-z0-9]+', '-', val)
val = re.sub(r'^-+|-+$', '', val)
return val if val else self.id
......@@ -45,4 +45,5 @@ class AgentUpdate(BaseModel):
class AgentResponse(AgentBase, TimestampMixin):
id: str
\ No newline at end of file
id: str
urlKey: str
\ No newline at end of file
......@@ -182,32 +182,35 @@ class CompanyPathRewriteMiddleware:
company_id = match.group(2)
resource_path = match.group(3)
if resource_path.startswith("runs"):
await self.app(scope, receive, send)
return
# Map specific resource paths for compatibility
if resource_path.startswith("skills"):
resource_path = re.sub(r"^skills(\b|/)", "company-skills\\1", resource_path)
elif resource_path.startswith("members"):
resource_path = re.sub(r"^members(\b|/)", "company-memberships/members\\1", resource_path)
elif resource_path.startswith("user-directory"):
resource_path = re.sub(r"^user-directory(\b|/)", "company-memberships/user-directory\\1", resource_path)
elif resource_path.startswith("join-requests"):
resource_path = re.sub(r"^join-requests(\b|/)", "company-memberships/join-requests\\1", resource_path)
# Rewrite path
scope["path"] = f"{api_prefix}/{resource_path}"
# Only rewrite if it's one of the legacy compatibility paths
is_legacy = (
resource_path.startswith("skills") or
resource_path.startswith("members") or
resource_path.startswith("user-directory") or
resource_path.startswith("join-requests")
)
# Rewrite query string
query_string = scope.get("query_string", b"").decode("utf-8")
addition = f"company_id={company_id}"
if query_string:
query_string = f"{query_string}&{addition}"
else:
query_string = addition
scope["query_string"] = query_string.encode("utf-8")
if is_legacy:
if resource_path.startswith("skills"):
resource_path = re.sub(r"^skills(\b|/)", "company-skills\\1", resource_path)
elif resource_path.startswith("members"):
resource_path = re.sub(r"^members(\b|/)", "company-memberships/members\\1", resource_path)
elif resource_path.startswith("user-directory"):
resource_path = re.sub(r"^user-directory(\b|/)", "company-memberships/user-directory\\1", resource_path)
elif resource_path.startswith("join-requests"):
resource_path = re.sub(r"^join-requests(\b|/)", "company-memberships/join-requests\\1", resource_path)
# Rewrite path
scope["path"] = f"{api_prefix}/{resource_path}"
# Rewrite query string
query_string = scope.get("query_string", b"").decode("utf-8")
addition = f"company_id={company_id}"
if query_string:
query_string = f"{query_string}&{addition}"
else:
query_string = addition
scope["query_string"] = query_string.encode("utf-8")
await self.app(scope, receive, send)
......
import logging
import os
import httpx
from typing import Dict, Any, List, Optional
logger = logging.getLogger(__name__)
class NocoBaseConnector:
def __init__(self, base_url: str = "http://localhost:13001"):
"""
NocoBase API connector.
By default, the server runs on 13001 in development mode.
"""
self.base_url = base_url.rstrip("/")
self.token: Optional[str] = None
self.headers: Dict[str, str] = {
"Accept": "application/json",
"Content-Type": "application/json"
}
async def login(self, email: str = "admin@nocobase.com", password: str = "admin123") -> bool:
"""
Authenticate with NocoBase to obtain a bearer token.
"""
url = f"{self.base_url}/api/auth:signIn"
try:
async with httpx.AsyncClient() as client:
response = await client.post(
url,
json={"email": email, "password": password},
headers={"Accept": "application/json", "Content-Type": "application/json"},
timeout=10.0
)
if response.status_code == 200:
res_data = response.json()
self.token = res_data.get("data", {}).get("token")
if self.token:
self.headers["Authorization"] = f"Bearer {self.token}"
logger.info("Successfully authenticated with NocoBase")
return True
else:
logger.error(f"Token not found in login response: {res_data}")
else:
logger.error(f"Login failed: status={response.status_code}, response={response.text}")
except Exception as e:
logger.error(f"Exception during NocoBase login: {e}", exc_info=True)
return False
async def ensure_logged_in(self) -> None:
"""
Helper method to ensure we have a token before making requests.
"""
if not self.token:
success = await self.login()
if not success:
raise RuntimeError("Failed to authenticate with NocoBase.")
async def list_records(self, collection: str, params: Optional[Dict[str, Any]] = None) -> List[Dict[str, Any]]:
"""
List records in a NocoBase collection.
"""
await self.ensure_logged_in()
url = f"{self.base_url}/api/{collection}:list"
try:
async with httpx.AsyncClient() as client:
response = await client.get(url, headers=self.headers, params=params, timeout=10.0)
if response.status_code == 200:
return response.json().get("data", [])
else:
logger.error(f"List records failed for {collection}: {response.status_code} - {response.text}")
except Exception as e:
logger.error(f"Exception listing records for {collection}: {e}")
return []
async def get_record(self, collection: str, record_id: Any) -> Optional[Dict[str, Any]]:
"""
Retrieve a single record by its token/ID.
"""
await self.ensure_logged_in()
url = f"{self.base_url}/api/{collection}:get"
params = {"filterByTk": record_id}
try:
async with httpx.AsyncClient() as client:
response = await client.get(url, headers=self.headers, params=params, timeout=10.0)
if response.status_code == 200:
return response.json().get("data")
else:
logger.error(f"Get record failed for {collection}/{record_id}: {response.status_code} - {response.text}")
except Exception as e:
logger.error(f"Exception getting record for {collection}/{record_id}: {e}")
return None
async def create_record(self, collection: str, data: Dict[str, Any]) -> Optional[Dict[str, Any]]:
"""
Create a new record in a NocoBase collection.
"""
await self.ensure_logged_in()
url = f"{self.base_url}/api/{collection}:create"
try:
async with httpx.AsyncClient() as client:
response = await client.post(url, headers=self.headers, json=data, timeout=10.0)
if response.status_code in (200, 201):
return response.json().get("data")
else:
logger.error(f"Create record failed for {collection}: {response.status_code} - {response.text}")
except Exception as e:
logger.error(f"Exception creating record in {collection}: {e}")
return None
async def update_record(self, collection: str, record_id: Any, data: Dict[str, Any]) -> Optional[Dict[str, Any]]:
"""
Update an existing record in a NocoBase collection.
"""
await self.ensure_logged_in()
url = f"{self.base_url}/api/{collection}:update"
params = {"filterByTk": record_id}
try:
async with httpx.AsyncClient() as client:
response = await client.post(url, headers=self.headers, params=params, json=data, timeout=10.0)
if response.status_code == 200:
return response.json().get("data")
else:
logger.error(f"Update record failed for {collection}/{record_id}: {response.status_code} - {response.text}")
except Exception as e:
logger.error(f"Exception updating record in {collection}/{record_id}: {e}")
return None
async def delete_record(self, collection: str, record_id: Any) -> bool:
"""
Delete a record from a NocoBase collection.
"""
await self.ensure_logged_in()
url = f"{self.base_url}/api/{collection}:destroy"
params = {"filterByTk": record_id}
try:
async with httpx.AsyncClient() as client:
response = await client.post(url, headers=self.headers, params=params, timeout=10.0)
if response.status_code == 200:
return True
else:
logger.error(f"Delete record failed for {collection}/{record_id}: {response.status_code} - {response.text}")
except Exception as e:
logger.error(f"Exception deleting record in {collection}/{record_id}: {e}")
return False
async def upload_file(self, file_path: str) -> Optional[Dict[str, Any]]:
"""
Upload a file to NocoBase attachments store.
"""
await self.ensure_logged_in()
url = f"{self.base_url}/api/attachments:create"
if not os.path.exists(file_path):
logger.error(f"File not found for upload: {file_path}")
return None
filename = os.path.basename(file_path)
upload_headers = self.headers.copy()
upload_headers.pop("Content-Type", None) # httpx sets boundary automatically for multipart
try:
async with httpx.AsyncClient() as client:
with open(file_path, "rb") as f:
files = {"file": (filename, f)}
response = await client.post(url, headers=upload_headers, files=files, timeout=30.0)
if response.status_code in (200, 201):
return response.json().get("data")
else:
logger.error(f"Upload file failed: {response.status_code} - {response.text}")
except Exception as e:
logger.error(f"Exception uploading file {file_path}: {e}")
return None
......@@ -26,10 +26,11 @@ logger = logging.getLogger(__name__)
# Agent role execution order — roles at the same index run in parallel
# Format: list of "stages", each stage is a list of roles that can run concurrently
PIPELINE_STAGES = [
["pm"], # Stage 1: PM writes spec (must be first)
["coder", "designer"], # Stage 2: Coder + Designer can run in parallel
["qa"], # Stage 3: QA reviews everything (must be last)
["devops"], # Stage 4: DevOps deploys (optional)
["ceo", "pm"], # Stage 1: CEO strategy + PM spec
["rd", "designer", "finance", "hr"], # Stage 2: Independent departments (parallel)
["marketing", "ecom", "coder"], # Stage 3: Depends on design/spec output
["qa"], # Stage 4: QA reviews everything
["devops"], # Stage 5: DevOps deploys (optional)
]
# Flat role order for sorting (PM < Coder < Designer < QA < DevOps)
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment