Commit d64a3f7b authored by Admin's avatar Admin

feat: add core company features, onboarding wizard, pipeline history, swarm trees, and blueprints

parent b5d9b6ef
# Canifa AI Platform — Backend & Tools # Paperclip AI Platform — Backend & Tools
A unified AI ecosystem powering Canifa's fashion retail intelligence, content automation, and inventory monitoring. A unified, multi-tenant AI ecosystem powering fashion retail intelligence, content automation, software engineering, and industry-specific business processes.
This platform serves as the central hub for all AI-driven workflows at Canifa, integrating Large Language Models (Gemini, GPT) with real-time business data (PostgreSQL, StarRocks, SQLite) to provide actionable insights and automated content. This platform serves as the central hub for all AI-driven workflows, integrating Large Language Models (Gemini, GPT, Claude) with business data and dynamic NocoBase collection engines.
--- ---
...@@ -13,38 +13,28 @@ This platform serves as the central hub for all AI-driven workflows at Canifa, i ...@@ -13,38 +13,28 @@ This platform serves as the central hub for all AI-driven workflows at Canifa, i
<img src="public/screenshots/03_chatbot_ui.png" width="31%" alt="Chatbot UI" /> <img src="public/screenshots/03_chatbot_ui.png" width="31%" alt="Chatbot UI" />
<img src="public/screenshots/06_product_desc.png" width="31%" alt="Product Description" /> <img src="public/screenshots/06_product_desc.png" width="31%" alt="Product Description" />
</p> </p>
<p align="center">
<img src="public/screenshots/04_stock_cache.png" width="45%" alt="Stock Cache" />
<img src="public/screenshots/05_fashion_matches.png" width="45%" alt="Fashion Matches" />
</p>
--- ---
## 🧠 Core Modules ## 🧠 Core Modules
### 1. 👗 AI Stylist & Fashion Matches ### 1. 🏢 Industry Blueprints & Spawning
- **Deterministic Styling**: Maps raw Magento categories into a clean 4-slot framework (Top, Bottom, Set, Accessories). - **Dynamic Provisioning**: Spawns tailored agent departments and hierarchies (e.g., Software Studio, E-commerce Brand, Marketing Agency) from single-file YAML blueprints.
- **Style Pairing**: Generates context-aware outfit combinations based on product attributes and seasonal trends. - **NocoBase Connector**: Per-company database schema isolation, dynamically provisioning collections and fields at company launch.
- **Grounded Retrieval**: Prevents AI hallucinations by anchoring recommendations in real product metadata. - **Topological Run Pipeline**: Parallel stage-based execution DAG with approval gates and manual review checkpoints.
### 2. ✍️ Ultra Product Descriptions
- **Automated Copywriting**: Generates web-ready, SEO-optimized descriptions from raw product data and images.
- **Human-in-the-loop**: Integrated approval workflow for AI-generated content before publishing to Magento.
- **N8n Integration**: Seamless data flow between Google Sheets, n8n, and the Backend API.
### 3. 📦 Stock & Inventory Intelligence ### 2. 👗 AI Stylist & Catalog Intelligence
- **Real-time Monitoring**: Dashboard for tracking stock levels across multiple warehouses. - **Styling pairing**: Standardizes product listings into category slots and generates styling recommendations based on product attributes and current trends.
- **Cache Management**: High-performance SQLite-based caching for instant stock lookups. - **Human-in-the-loop**: Integrated approval flow for AI-generated product copywriting and catalog updates.
- **StarRocks Sync**: Enterprise-grade data warehousing for historical analysis and forecasting.
### 4. 📅 Content Automation ### 3. 📦 Inventory & Stock Intelligence
- **Composer**: AI-assisted creation of marketing materials and social media posts. - **Real-time Monitoring**: Dashboard for tracking stock levels across warehouses.
- **Calendar & Approval**: Full lifecycle management for promotional content. - **StarRocks / BI sync**: Historical analysis and forecasting capabilities.
- **Social Inbox**: Centralized monitoring of customer interactions across social channels.
--- ---
## 🛠 Tech Stack ## 🛠 Tech Stack
- **Backend**: FastAPI (Async Python 3.10+) - **Backend**: FastAPI (Async Python 3.10+)
- **Frontend**: React + TypeScript + Vite + Tailwind CSS
- **AI Framework**: LangChain & LangGraph - **AI Framework**: LangChain & LangGraph
...@@ -22,6 +22,11 @@ BUDGET_ENFORCEMENT_ENABLED=True ...@@ -22,6 +22,11 @@ BUDGET_ENFORCEMENT_ENABLED=True
ANTHROPIC_API_KEY= ANTHROPIC_API_KEY=
OPENAI_API_KEY= OPENAI_API_KEY=
# NocoBase ERP integration
NOCOBASE_URL=http://localhost:13002
NOCOBASE_EMAIL=admin@nocobase.com
NOCOBASE_PASSWORD=admin123
# Logging # Logging
LOG_LEVEL=INFO LOG_LEVEL=INFO
......
...@@ -202,6 +202,7 @@ def init_agent( ...@@ -202,6 +202,7 @@ def init_agent(
checkpoint_max_total_size_mb: int = 500, checkpoint_max_total_size_mb: int = 500,
checkpoint_max_file_size_mb: int = 10, checkpoint_max_file_size_mb: int = 10,
pass_session_id: bool = False, pass_session_id: bool = False,
team_id: str = None,
): ):
""" """
Initialize the AI Agent. Initialize the AI Agent.
...@@ -453,6 +454,8 @@ def init_agent( ...@@ -453,6 +454,8 @@ def init_agent(
# Store toolset filtering options # Store toolset filtering options
agent.enabled_toolsets = enabled_toolsets agent.enabled_toolsets = enabled_toolsets
agent.disabled_toolsets = disabled_toolsets agent.disabled_toolsets = disabled_toolsets
if not hasattr(agent, "team_id"):
agent.team_id = team_id
# Model response configuration # Model response configuration
agent.max_tokens = max_tokens # None = use model default agent.max_tokens = max_tokens # None = use model default
......
...@@ -414,11 +414,13 @@ class AIAgent: ...@@ -414,11 +414,13 @@ class AIAgent:
checkpoint_max_file_size_mb: int = 10, checkpoint_max_file_size_mb: int = 10,
pass_session_id: bool = False, pass_session_id: bool = False,
company_id: str = None, company_id: str = None,
team_id: str = None,
model_name: str = None, model_name: str = None,
**kwargs **kwargs
): ):
"""Forwarder — see ``agent.agent_init.init_agent``.""" """Forwarder — see ``agent.agent_init.init_agent``."""
self.company_id = company_id self.company_id = company_id
self.team_id = team_id
self.model_name = model_name self.model_name = model_name
if model_name and not model: if model_name and not model:
model = model_name model = model_name
...@@ -489,6 +491,7 @@ class AIAgent: ...@@ -489,6 +491,7 @@ class AIAgent:
checkpoint_max_total_size_mb=checkpoint_max_total_size_mb, checkpoint_max_total_size_mb=checkpoint_max_total_size_mb,
checkpoint_max_file_size_mb=checkpoint_max_file_size_mb, checkpoint_max_file_size_mb=checkpoint_max_file_size_mb,
pass_session_id=pass_session_id, pass_session_id=pass_session_id,
team_id=team_id,
) )
def _get_session_db_for_recall(self): def _get_session_db_for_recall(self):
...@@ -4171,7 +4174,7 @@ class AIAgent: ...@@ -4171,7 +4174,7 @@ class AIAgent:
if hasattr(self, "company_id") and self.company_id: if hasattr(self, "company_id") and self.company_id:
try: try:
from tools.get_tools import register_company_tools from tools.get_tools import register_company_tools
register_company_tools(self.company_id) register_company_tools(self.company_id, getattr(self, "team_id", None), getattr(self, "enabled_toolsets", None))
# Re-load tools list from the registry so they are visible to the agent loop # Re-load tools list from the registry so they are visible to the agent loop
from model_tools import get_tool_definitions from model_tools import get_tool_definitions
self.tools = get_tool_definitions(quiet_mode=self.quiet_mode) self.tools = get_tool_definitions(quiet_mode=self.quiet_mode)
......
"""Facebook Page posting tool for LLM agents."""
import logging
import os
import uuid
import httpx
logger = logging.getLogger(__name__)
GRAPH_API_URL = "https://graph.facebook.com/v19.0"
async def post_to_facebook(message: str) -> dict:
"""Post a message to the configured Facebook Page.
If FACEBOOK_PAGE_ACCESS_TOKEN and FACEBOOK_PAGE_ID are set, makes a real
POST to the Graph API. Otherwise runs in dry-run mode.
"""
access_token = os.environ.get("FACEBOOK_PAGE_ACCESS_TOKEN", "")
page_id = os.environ.get("FACEBOOK_PAGE_ID", "")
if access_token and page_id:
return await _post_real(page_id, access_token, message)
logger.info("[facebook_tool] DRY-RUN: would post: %s", message[:120])
return {
"post_id": f"dry_run_{uuid.uuid4().hex[:12]}",
"status": "ok",
"mode": "dry-run",
"message_preview": message[:200],
}
async def _post_real(page_id: str, token: str, message: str) -> dict:
url = f"{GRAPH_API_URL}/{page_id}/feed"
async with httpx.AsyncClient(timeout=30) as client:
resp = await client.post(url, data={"message": message, "access_token": token})
resp.raise_for_status()
data = resp.json()
return {"post_id": data.get("id", ""), "status": "ok", "mode": "real"}
def get_facebook_tools() -> list[dict]:
"""Return tool definitions compatible with OpenAI function-calling format."""
return [
{
"type": "function",
"function": {
"name": "post_to_facebook",
"description": "Post a message to the company Facebook Page. Runs in dry-run mode when credentials are not configured.",
"parameters": {
"type": "object",
"properties": {
"message": {
"type": "string",
"description": "The text content to publish on the Facebook Page.",
},
},
"required": ["message"],
},
},
}
]
...@@ -3,7 +3,7 @@ import logging ...@@ -3,7 +3,7 @@ import logging
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
def get_all_tools(company_id: str | None = None) -> list[BaseTool]: def get_all_tools(company_id: str | None = None, team_id: str | None = None, enabled_toolsets: list[str] | None = None) -> list[BaseTool]:
"""Return all tools for the Agent, including Workspace tools, MCP tools, and dynamic Skill tools.""" """Return all tools for the Agent, including Workspace tools, MCP tools, and dynamic Skill tools."""
tools = [] tools = []
...@@ -20,29 +20,38 @@ def get_all_tools(company_id: str | None = None) -> list[BaseTool]: ...@@ -20,29 +20,38 @@ def get_all_tools(company_id: str | None = None) -> list[BaseTool]:
try: try:
from .mcp_client import MCPClientManager from .mcp_client import MCPClientManager
mcp_tools = MCPClientManager.get_instance().get_tools() mcp_tools = MCPClientManager.get_instance().get_tools()
logger.info(f"Loaded {len(mcp_tools)} MCP tools") if enabled_toolsets is not None:
tools.extend(mcp_tools) filtered_mcp = []
for t in mcp_tools:
# t.name format: mcp_{server_name}_{tool_name}
if any(t.name.startswith(f"mcp_{ts}_") for ts in enabled_toolsets):
filtered_mcp.append(t)
logger.info(f"Loaded {len(filtered_mcp)} MCP tools (filtered from {len(mcp_tools)})")
tools.extend(filtered_mcp)
else:
logger.info(f"Loaded {len(mcp_tools)} MCP tools")
tools.extend(mcp_tools)
except Exception as e: except Exception as e:
logger.error(f"Error loading MCP tools: {e}", exc_info=True) logger.error(f"Error loading MCP tools: {e}", exc_info=True)
# 3. Load Dynamic Skills # 3. Load Dynamic Skills
try: try:
from .skills_loader import load_skills_as_tools from .skills_loader import load_skills_as_tools
skill_tools = load_skills_as_tools(company_id) skill_tools = load_skills_as_tools(company_id, team_id)
logger.info(f"Loaded {len(skill_tools)} skill tools for company_id: {company_id}") logger.info(f"Loaded {len(skill_tools)} skill tools for company_id: {company_id}, team_id: {team_id}")
tools.extend(skill_tools) tools.extend(skill_tools)
except Exception as e: except Exception as e:
logger.error(f"Error loading skill tools: {e}", exc_info=True) logger.error(f"Error loading skill tools: {e}", exc_info=True)
return tools return tools
def register_company_tools(company_id: str | None = None) -> None: def register_company_tools(company_id: str | None = None, team_id: str | None = None, enabled_toolsets: list[str] | None = None) -> None:
"""Load all workspace, MCP, and skill tools for the company, and register them in the Hermes ToolRegistry.""" """Load all workspace, MCP, and skill tools for the company, and register them in the Hermes ToolRegistry."""
from .registry import registry from .registry import registry
from langchain_core.utils.function_calling import convert_to_openai_function from langchain_core.utils.function_calling import convert_to_openai_function
import inspect import inspect
tools = get_all_tools(company_id) tools = get_all_tools(company_id, team_id, enabled_toolsets)
for tool in tools: for tool in tools:
try: try:
openai_schema = convert_to_openai_function(tool) openai_schema = convert_to_openai_function(tool)
......
...@@ -30,7 +30,7 @@ def parse_skill_md(file_path: str) -> Optional[Dict[str, Any]]: ...@@ -30,7 +30,7 @@ def parse_skill_md(file_path: str) -> Optional[Dict[str, Any]]:
logger.error(f"Error parsing skill markdown at {file_path}: {e}") logger.error(f"Error parsing skill markdown at {file_path}: {e}")
return None return None
def get_company_skill_names(company_id: str) -> List[str]: def get_company_skill_names(company_id: str, team_id: str | None = None) -> List[str]:
root_dir = os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))) root_dir = os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
db_path = os.path.join(root_dir, "backend", "paperclip.db") db_path = os.path.join(root_dir, "backend", "paperclip.db")
if not os.path.exists(db_path): if not os.path.exists(db_path):
...@@ -40,7 +40,10 @@ def get_company_skill_names(company_id: str) -> List[str]: ...@@ -40,7 +40,10 @@ def get_company_skill_names(company_id: str) -> List[str]:
try: try:
conn = sqlite3.connect(db_path) conn = sqlite3.connect(db_path)
cursor = conn.cursor() cursor = conn.cursor()
cursor.execute("SELECT name FROM company_skills WHERE company_id = ?", (company_id,)) if team_id:
cursor.execute("SELECT name FROM company_skills WHERE company_id = ? AND (team_id = ? OR team_id IS NULL)", (company_id, team_id))
else:
cursor.execute("SELECT name FROM company_skills WHERE company_id = ? AND team_id IS NULL", (company_id,))
rows = cursor.fetchall() rows = cursor.fetchall()
conn.close() conn.close()
return [row[0] for row in rows] return [row[0] for row in rows]
...@@ -48,7 +51,7 @@ def get_company_skill_names(company_id: str) -> List[str]: ...@@ -48,7 +51,7 @@ def get_company_skill_names(company_id: str) -> List[str]:
logger.error(f"Error querying company skills from SQLite: {e}") logger.error(f"Error querying company skills from SQLite: {e}")
return [] return []
def load_skills_as_tools(company_id: str | None = None) -> List[StructuredTool]: def load_skills_as_tools(company_id: str | None = None, team_id: str | None = None) -> List[StructuredTool]:
root_dir = os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))) root_dir = os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
skills_dir = os.path.join(root_dir, ".agent", "skills") skills_dir = os.path.join(root_dir, ".agent", "skills")
...@@ -59,8 +62,8 @@ def load_skills_as_tools(company_id: str | None = None) -> List[StructuredTool]: ...@@ -59,8 +62,8 @@ def load_skills_as_tools(company_id: str | None = None) -> List[StructuredTool]:
# Get the names of skills we should load # Get the names of skills we should load
allowed_skill_names = None allowed_skill_names = None
if company_id: if company_id:
allowed_skill_names = get_company_skill_names(company_id) allowed_skill_names = get_company_skill_names(company_id, team_id)
logger.info(f"Loaded allowed skill names for company {company_id}: {allowed_skill_names}") logger.info(f"Loaded allowed skill names for company {company_id}, team {team_id}: {allowed_skill_names}")
tools = [] tools = []
......
...@@ -8,6 +8,7 @@ import os ...@@ -8,6 +8,7 @@ import os
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), '..'))) sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), '..')))
from models.database import async_engine, Base from models.database import async_engine, Base
import models # Register all models on Base.metadata
from config import settings from config import settings
config = context.config config = context.config
...@@ -29,24 +30,28 @@ def run_migrations_offline(): ...@@ -29,24 +30,28 @@ def run_migrations_offline():
target_metadata=target_metadata, target_metadata=target_metadata,
literal_binds=True, literal_binds=True,
dialect_opts={"paramstyle": "named"}, dialect_opts={"paramstyle": "named"},
render_as_batch=True,
) )
with context.begin_transaction(): with context.begin_transaction():
context.run_migrations() context.run_migrations()
def run_migrations_online(): import asyncio
"""Run migrations in 'online' mode."""
connectable = async_engine
with connectable.connect() as connection: def do_run_migrations(connection):
context.configure( context.configure(connection=connection, target_metadata=target_metadata, render_as_batch=True)
connection=connection, with context.begin_transaction():
target_metadata=target_metadata, context.run_migrations()
)
with context.begin_transaction(): async def run_async_migrations():
context.run_migrations() connectable = async_engine
async with connectable.connect() as connection:
await connection.run_sync(do_run_migrations)
def run_migrations_online():
"""Run migrations in 'online' mode."""
asyncio.run(run_async_migrations())
if context.is_offline_mode(): if context.is_offline_mode():
......
"""${message}
Revision ID: ${up_revision}
Revises: ${down_revision | comma,n}
Create Date: ${create_date}
"""
from typing import Sequence, Union
from alembic import op
import sqlalchemy as sa
${imports if imports else ""}
# revision identifiers, used by Alembic.
revision: str = ${repr(up_revision)}
down_revision: Union[str, None] = ${repr(down_revision)}
branch_labels: Union[str, Sequence[str], None] = ${repr(branch_labels)}
depends_on: Union[str, Sequence[str], None] = ${repr(depends_on)}
def upgrade() -> None:
${upgrades if upgrades else "pass"}
def downgrade() -> None:
${downgrades if downgrades else "pass"}
...@@ -217,11 +217,8 @@ async def run_agent_in_background( ...@@ -217,11 +217,8 @@ async def run_agent_in_background(
from config import settings from config import settings
if settings.openai_api_key == "wrapper-test-key" and os.environ.get("REAL_AGENT_RUN") != "true": if os.environ.get("STUB_AGENTS") == "true":
import asyncio
import json
from pathlib import Path from pathlib import Path
from datetime import datetime
# Ensure shared workspaces folder exists # Ensure shared workspaces folder exists
shared_dir = Path("workspaces") / company_id / "shared" shared_dir = Path("workspaces") / company_id / "shared"
...@@ -452,15 +449,35 @@ async def run_agent_in_background( ...@@ -452,15 +449,35 @@ async def run_agent_in_background(
except Exception as e: except Exception as e:
logging.getLogger(__name__).error(f"Failed to spawn CubeSandbox instance: {e}", exc_info=True) logging.getLogger(__name__).error(f"Failed to spawn CubeSandbox instance: {e}", exc_info=True)
from services.runtime_env import resolve_runtime_env, resolve_enabled_toolsets
async with get_db_context() as db:
result = await db.execute(select(Agent).where(Agent.id == agent_id))
db_agent = result.scalar_one_or_none()
team_id = db_agent.team_id if db_agent else None
env_vars = await resolve_runtime_env(db, company_id, team_id)
enabled_toolsets = await resolve_enabled_toolsets(db, company_id, team_id)
original_env = {}
injected_keys = []
for k, v in env_vars.items():
if k in os.environ:
original_env[k] = os.environ[k]
else:
injected_keys.append(k)
os.environ[k] = str(v)
ai_agent = AIAgent( ai_agent = AIAgent(
model_name=model_name, model_name=model_name,
company_id=company_id, company_id=company_id,
team_id=team_id,
tool_progress_callback=log_writer.on_subagent_progress, tool_progress_callback=log_writer.on_subagent_progress,
step_callback=log_writer.on_step, step_callback=log_writer.on_step,
status_callback=log_writer.on_status, status_callback=log_writer.on_status,
api_key=settings.openai_api_key, api_key=settings.openai_api_key,
base_url=base_url, base_url=base_url,
api_mode=api_mode, api_mode=api_mode,
enabled_toolsets=enabled_toolsets,
) )
...@@ -507,6 +524,16 @@ async def run_agent_in_background( ...@@ -507,6 +524,16 @@ async def run_agent_in_background(
del os.environ["CUBE_SANDBOX_ID"] del os.environ["CUBE_SANDBOX_ID"]
except Exception as e: except Exception as e:
logging.getLogger(__name__).error(f"Failed to delete CubeSandbox instance: {e}", exc_info=True) logging.getLogger(__name__).error(f"Failed to delete CubeSandbox instance: {e}", exc_info=True)
# Teardown runtime env vars
try:
for k in injected_keys:
if k in os.environ:
del os.environ[k]
for k, v in original_env.items():
os.environ[k] = v
except Exception:
pass
if workspace: if workspace:
try: try:
import shutil import shutil
...@@ -526,7 +553,6 @@ async def run_agent_in_background( ...@@ -526,7 +553,6 @@ async def run_agent_in_background(
# Look up agent role to name the output file # Look up agent role to name the output file
async with get_db_context() as db: async with get_db_context() as db:
from models import Agent
agent_res = await db.execute(select(Agent).where(Agent.id == agent_id)) agent_res = await db.execute(select(Agent).where(Agent.id == agent_id))
agent_obj = agent_res.scalar_one_or_none() agent_obj = agent_res.scalar_one_or_none()
role_name = (agent_obj.role or "agent").lower().replace(" ", "_") if agent_obj else "agent" role_name = (agent_obj.role or "agent").lower().replace(" ", "_") if agent_obj else "agent"
......
...@@ -3,7 +3,9 @@ Approval routes. ...@@ -3,7 +3,9 @@ Approval routes.
""" """
from typing import Optional, List from typing import Optional, List
from datetime import datetime
from fastapi import APIRouter, Depends, HTTPException, Query, status, Body from fastapi import APIRouter, Depends, HTTPException, Query, status, Body
from pydantic import BaseModel
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func from sqlalchemy import select, func
...@@ -18,6 +20,92 @@ from services.activity_logger import log_activity ...@@ -18,6 +20,92 @@ from services.activity_logger import log_activity
router = APIRouter(prefix="/approvals", tags=["approvals"]) router = APIRouter(prefix="/approvals", tags=["approvals"])
class ApprovalDecisionBody(BaseModel):
# NOTE: snake_case on purpose — CaseConversionMiddleware rewrites the
# incoming camelCase body (`decisionNote`) to snake_case before it reaches
# the route, so the model field must be snake_case to bind.
decision_note: Optional[str] = None
async def _apply_resolution(
db: AsyncSession,
current_user: AuthUserResponse,
approval_id: str,
resolution: str,
note: Optional[str],
) -> Approval:
"""Resolve an approval, deriving the company from the approval record itself.
Used by the FE-facing /approve and /reject endpoints which pass only the
approval id (no company_id). Sets `resolution` which the pipeline
orchestrator polls to resume/abort a gated stage.
"""
result = await db.execute(select(Approval).where(Approval.id == approval_id))
approval = result.scalar_one_or_none()
if not approval:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Approval not found")
# Verify membership against the approval's own company
membership_result = await db.execute(
select(CompanyMembership).where(
CompanyMembership.company_id == approval.company_id,
CompanyMembership.principal_type == "user",
CompanyMembership.principal_id == current_user.id,
CompanyMembership.status == "active",
)
)
if not membership_result.scalar_one_or_none():
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="No access to company")
if approval.resolution is not None:
raise HTTPException(status_code=400, detail="Approval already resolved")
approval.resolution = resolution
approval.resolved_by = current_user.id
approval.resolved_at = datetime.utcnow()
if note:
approval.resolution_note = note
await log_activity(
db=db,
company_id=approval.company_id,
actor_type="user",
actor_id=current_user.id,
action="approval.resolved",
resource_type="approval",
resource_id=approval_id,
changes={"resolution": (None, resolution)},
)
await db.commit()
await db.refresh(approval)
return approval
@router.post("/{approval_id}/approve", response_model=ApprovalResponse)
async def approve_approval(
approval_id: str,
db: AsyncSession = Depends(get_db_session),
current_user: AuthUserResponse = Depends(get_current_user),
body: Optional[ApprovalDecisionBody] = Body(default=None),
):
"""Approve an approval request (resumes a gated pipeline stage)."""
note = body.decision_note if body else None
return await _apply_resolution(db, current_user, approval_id, "approved", note)
@router.post("/{approval_id}/reject", response_model=ApprovalResponse)
async def reject_approval(
approval_id: str,
db: AsyncSession = Depends(get_db_session),
current_user: AuthUserResponse = Depends(get_current_user),
body: Optional[ApprovalDecisionBody] = Body(default=None),
):
"""Reject an approval request (fails a gated pipeline stage)."""
note = body.decision_note if body else None
return await _apply_resolution(db, current_user, approval_id, "rejected", note)
@router.post("", response_model=ApprovalResponse, status_code=status.HTTP_201_CREATED) @router.post("", response_model=ApprovalResponse, status_code=status.HTTP_201_CREATED)
async def create_approval( async def create_approval(
approval: ApprovalCreate, approval: ApprovalCreate,
...@@ -95,19 +183,34 @@ async def list_approvals( ...@@ -95,19 +183,34 @@ async def list_approvals(
@router.get("/{approval_id}", response_model=ApprovalResponse) @router.get("/{approval_id}", response_model=ApprovalResponse)
async def get_approval( async def get_approval(
approval_id: str, approval_id: str,
company_id: str = Query(..., description="Company ID"), company_id: Optional[str] = Query(None, description="Company ID (optional; derived from approval if omitted)"),
scope: CompanyScope = Depends(require_company_scope), scope: Optional[CompanyScope] = Depends(require_company_scope),
db: AsyncSession = Depends(get_db_session) db: AsyncSession = Depends(get_db_session),
current_user: AuthUserResponse = Depends(get_current_user),
): ):
"""Get a specific approval.""" """Get a specific approval. Accepts either an explicit company_id or derives it."""
result = await db.execute( result = await db.execute(select(Approval).where(Approval.id == approval_id))
select(Approval).where(Approval.id == approval_id, Approval.company_id == company_id)
)
approval = result.scalar_one_or_none() approval = result.scalar_one_or_none()
if not approval: if not approval:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Approval not found") raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Approval not found")
if company_id and approval.company_id != company_id:
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Approval not found")
# When no company scope was supplied, verify membership against the approval's company
if not company_id:
membership_result = await db.execute(
select(CompanyMembership).where(
CompanyMembership.company_id == approval.company_id,
CompanyMembership.principal_type == "user",
CompanyMembership.principal_id == current_user.id,
CompanyMembership.status == "active",
)
)
if not membership_result.scalar_one_or_none():
raise HTTPException(status_code=status.HTTP_403_FORBIDDEN, detail="No access to company")
return approval return approval
......
This diff is collapsed.
This diff is collapsed.
...@@ -84,6 +84,7 @@ async def list_issues( ...@@ -84,6 +84,7 @@ async def list_issues(
status: Optional[str] = Query(None, description="Filter by status"), status: Optional[str] = Query(None, description="Filter by status"),
priority: Optional[str] = Query(None, description="Filter by priority"), priority: Optional[str] = Query(None, description="Filter by priority"),
assignee_agent_id: Optional[str] = Query(None, description="Filter by assignee"), assignee_agent_id: Optional[str] = Query(None, description="Filter by assignee"),
parent_id: Optional[str] = Query(None, description="Filter by parent issue"),
pagination: PaginationParams = Depends(), pagination: PaginationParams = Depends(),
db: AsyncSession = Depends(get_db_session), db: AsyncSession = Depends(get_db_session),
current_user: AuthUserResponse = Depends(get_current_user) current_user: AuthUserResponse = Depends(get_current_user)
...@@ -104,6 +105,8 @@ async def list_issues( ...@@ -104,6 +105,8 @@ async def list_issues(
stmt = stmt.where(Issue.priority == priority) stmt = stmt.where(Issue.priority == priority)
if assignee_agent_id: if assignee_agent_id:
stmt = stmt.where(Issue.assignee_agent_id == assignee_agent_id) stmt = stmt.where(Issue.assignee_agent_id == assignee_agent_id)
if parent_id:
stmt = stmt.where(Issue.parent_id == parent_id)
stmt = stmt.order_by(Issue.created_at.desc()).offset(pagination.offset).limit(pagination.limit) stmt = stmt.order_by(Issue.created_at.desc()).offset(pagination.offset).limit(pagination.limit)
result = await db.execute(stmt) result = await db.execute(stmt)
......
name: "Canifa Fashion Company"
industry: "fashion"
description: "Giả lập công ty thời trang Canifa hoàn chỉnh với 7 phòng ban, 13 agents, pipeline DAG 7 stages. Từ CEO → Publisher đẩy bài lên Facebook."
roles:
- role: "ceo"
name: "CEO Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "strategic direction, weekly reporting, approval, corporate oversight"
reports_to: null
- role: "pm"
name: "Product Manager Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "specifications writing, project tracking, requirements gathering, PRD creation"
reports_to: "ceo"
- role: "rd"
name: "R&D Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "trend analysis, fashion market research, competitor analysis, fabric database"
reports_to: "pm"
- role: "designer"
name: "Designer Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "mockups, collection concept art, banner design, poster generation"
reports_to: "pm"
- role: "finance"
name: "Finance Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "bookkeeping, profit and loss reports, invoice verification, budgeting"
reports_to: "pm"
- role: "hr"
name: "HR Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "employee onboarding, recruiting, payroll coordination, time-off requests"
reports_to: "pm"
- role: "content"
name: "Content Creator Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "facebook post writing, ad copy, social media content, brand voice"
reports_to: "pm"
- role: "ads"
name: "Ads Manager Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "ad budget allocation, campaign planning, performance tracking"
reports_to: "pm"
- role: "ecom"
name: "E-commerce Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "product listing description, SEO keywords, review sentiment monitoring"
reports_to: "pm"
- role: "coder"
name: "IT Dev Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "full-stack development, bug fixing, test scripts, automated deployments"
reports_to: "pm"
- role: "qa"
name: "QA Specialist Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "review specifications, run test cases, report bugs, content review"
reports_to: "pm"
- role: "publisher"
name: "Publisher Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "facebook posting, content publishing, deployment"
reports_to: "pm"
pipeline:
- role: "ceo"
depends_on: []
approval_gate: false
- role: "pm"
depends_on: ["ceo"]
approval_gate: false
- role: "rd"
depends_on: ["pm"]
approval_gate: false
- role: "designer"
depends_on: ["pm"]
approval_gate: false
- role: "finance"
depends_on: ["pm"]
approval_gate: false
- role: "hr"
depends_on: ["pm"]
approval_gate: false
- role: "content"
depends_on: ["designer", "rd"]
approval_gate: false
- role: "ads"
depends_on: ["finance", "rd"]
approval_gate: false
- role: "ecom"
depends_on: ["designer"]
approval_gate: false
- role: "coder"
depends_on: ["designer", "content"]
approval_gate: false
- role: "qa"
depends_on: ["coder", "content", "ecom"]
approval_gate: false
- role: "publisher"
depends_on: ["qa"]
approval_gate: true
teams:
- id: executive
name: "Ban Giám Đốc"
roles: [ceo]
- id: product
name: "Phòng Sản Phẩm"
roles: [pm, designer, rd]
- id: marketing
name: "Phòng Marketing"
roles: [content, ads, publisher]
- id: sales
name: "Phòng Kinh Doanh"
roles: [ecom]
- id: finance_team
name: "Phòng Tài Chính"
roles: [finance]
- id: hr_team
name: "Phòng Nhân Sự"
roles: [hr]
- id: engineering
name: "Phòng IT"
roles: [coder, qa]
plugins:
- plugin: nocobase_erp
teams: [all]
secrets:
- key: NOCOBASE_URL
- key: NOCOBASE_EMAIL
- key: NOCOBASE_PASSWORD
- plugin: facebook_pages
teams: [marketing]
secrets:
- key: FACEBOOK_PAGE_ACCESS_TOKEN
- key: FACEBOOK_PAGE_ID
nocobase_collections:
- name: "products"
title: "📦 Products"
fields:
- name: "title"
type: "string"
uiSchema:
title: "Product Title"
x-component: "Input"
- name: "outfits"
title: "👗 Outfits"
fields:
- name: "title"
type: "string"
uiSchema:
title: "Outfit Title"
x-component: "Input"
- name: "campaigns"
title: "📣 Campaigns"
fields:
- name: "title"
type: "string"
uiSchema:
title: "Campaign Title"
x-component: "Input"
...@@ -34,3 +34,12 @@ celery_app.conf.update( ...@@ -34,3 +34,12 @@ celery_app.conf.update(
# Auto-discover tasks in tasks/ directory # Auto-discover tasks in tasks/ directory
celery_app.autodiscover_tasks(["tasks"]) celery_app.autodiscover_tasks(["tasks"])
# Periodic task beat schedules
celery_app.conf.beat_schedule = {
"check-routines-every-minute": {
"task": "tasks.check_scheduled_routines",
"schedule": 60.0,
}
}
...@@ -67,6 +67,12 @@ class Settings(BaseSettings): ...@@ -67,6 +67,12 @@ class Settings(BaseSettings):
openai_api_key: Optional[str] = Field(default=None, env="OPENAI_API_KEY") openai_api_key: Optional[str] = Field(default=None, env="OPENAI_API_KEY")
openai_base_url: Optional[str] = Field(default=None, env="OPENAI_BASE_URL") openai_base_url: Optional[str] = Field(default=None, env="OPENAI_BASE_URL")
# NocoBase ERP
nocobase_url: Optional[str] = Field(default=None, env="NOCOBASE_URL")
nocobase_email: Optional[str] = Field(default=None, env="NOCOBASE_EMAIL")
nocobase_password: Optional[str] = Field(default=None, env="NOCOBASE_PASSWORD")
# Redis configuration # Redis configuration
redis_cache_url: str = Field(default="localhost", env="REDIS_CACHE_URL") redis_cache_url: str = Field(default="localhost", env="REDIS_CACHE_URL")
......
...@@ -54,8 +54,17 @@ from .plugin_config import PluginConfig ...@@ -54,8 +54,17 @@ from .plugin_config import PluginConfig
from .plugin_jobs import PluginJob from .plugin_jobs import PluginJob
from .sidebar_preferences import SidebarPreference from .sidebar_preferences import SidebarPreference
from .cli_auth_challenges import CliAuthChallenge from .cli_auth_challenges import CliAuthChallenge
from .teams import Team
# ── Summer Camp 2026 Campaign models ────────────────────────────────────────
from .campaigns import Campaign
from .product_specs import ProductSpec
from .materials import Material
from .ad_plans import AdPlan
from .budget_items import BudgetItem
__all__ = [ __all__ = [
"Team",
"CliAuthChallenge", "CliAuthChallenge",
"Base", "Base",
"async_engine", "async_engine",
...@@ -111,4 +120,10 @@ __all__ = [ ...@@ -111,4 +120,10 @@ __all__ = [
"PluginConfig", "PluginConfig",
"PluginJob", "PluginJob",
"SidebarPreference", "SidebarPreference",
# Campaign models
"Campaign",
"ProductSpec",
"Material",
"AdPlan",
"BudgetItem",
] ]
"""
Ad Plan / Campaign Marketing Plan — drafted by Marketing, executed by Ads Manager.
"""
from sqlalchemy import String, Text, DateTime, ForeignKey, JSON, Integer, Numeric
from sqlalchemy.orm import Mapped, mapped_column
from datetime import datetime
from typing import Optional
from .database import Base
class AdPlan(Base):
__tablename__ = "ad_plans"
id: Mapped[str] = mapped_column(
String(36), primary_key=True, default=lambda: __import__("uuid").uuid4().hex
)
company_id: Mapped[str] = mapped_column(
String(36), ForeignKey("companies.id"), nullable=False, index=True
)
campaign_id: Mapped[str] = mapped_column(
String(36), ForeignKey("campaigns.id"), nullable=False, index=True
)
created_by_agent_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("agents.id"), nullable=True
)
# identity
name: Mapped[str] = mapped_column(String(255), nullable=False)
platform: Mapped[str] = mapped_column(
String(40), nullable=False
) # facebook | instagram | tiktok | google
objective: Mapped[Optional[str]] = mapped_column(String(80), nullable=True)
status: Mapped[str] = mapped_column(
String(30), nullable=False, default="draft"
) # draft | scheduled | running | paused | completed
# content brief
ad_copy: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
creative_brief: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
target_audience: Mapped[Optional[dict]] = mapped_column(JSON, nullable=True)
hashtags: Mapped[list] = mapped_column(JSON, nullable=False, default=list)
# schedule
scheduled_start: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
scheduled_end: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
# budget (synced from budget_items)
allocated_budget_vnd: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
spent_vnd: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
# performance
impressions: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
clicks: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
conversions: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
revenue_vnd: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
extra: Mapped[dict] = mapped_column("extra", JSON, nullable=False, default=dict)
created_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
updated_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
...@@ -16,6 +16,7 @@ class Agent(Base): ...@@ -16,6 +16,7 @@ class Agent(Base):
id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex) id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex)
company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False) company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False)
team_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("teams.id"), nullable=True)
name: Mapped[str] = mapped_column(String(255), nullable=False) name: Mapped[str] = mapped_column(String(255), nullable=False)
role: Mapped[str] = mapped_column(String(50), nullable=False, default="general") role: Mapped[str] = mapped_column(String(50), nullable=False, default="general")
title: Mapped[Optional[str]] = mapped_column(String(255), nullable=True) title: Mapped[Optional[str]] = mapped_column(String(255), nullable=True)
......
"""
Budget Item — detailed cost breakdown per campaign (Finance / HR).
"""
from sqlalchemy import String, Text, DateTime, ForeignKey, JSON, Integer
from sqlalchemy.orm import Mapped, mapped_column
from datetime import datetime
from typing import Optional
from .database import Base
class BudgetItem(Base):
__tablename__ = "budget_items"
id: Mapped[str] = mapped_column(
String(36), primary_key=True, default=lambda: __import__("uuid").uuid4().hex
)
company_id: Mapped[str] = mapped_column(
String(36), ForeignKey("companies.id"), nullable=False, index=True
)
campaign_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("campaigns.id"), nullable=True, index=True
)
product_spec_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("product_specs.id"), nullable=True
)
ad_plan_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("ad_plans.id"), nullable=True
)
created_by_agent_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("agents.id"), nullable=True
)
# identity
category: Mapped[str] = mapped_column(String(60), nullable=False)
# categories: production | material | ad_spend | design | logistics | hr | other
label: Mapped[str] = mapped_column(String(255), nullable=False)
# amounts (VND, integer — no float precision loss)
estimated_vnd: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
actual_vnd: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
# tracking
status: Mapped[str] = mapped_column(
String(30), nullable=False, default="planned"
) # planned | approved | committed | paid | cancelled
approved_by_agent_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("agents.id"), nullable=True
)
approved_at: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
notes: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
attachments: Mapped[list] = mapped_column(JSON, nullable=False, default=list)
extra: Mapped[dict] = mapped_column("extra", JSON, nullable=False, default=dict)
created_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
updated_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
"""
Campaign model — Summer Camp 2026.
"""
from sqlalchemy import String, Text, DateTime, ForeignKey, JSON, Integer
from sqlalchemy.orm import Mapped, mapped_column
from datetime import datetime
from typing import Optional
from .database import Base
class Campaign(Base):
__tablename__ = "campaigns"
id: Mapped[str] = mapped_column(
String(36), primary_key=True, default=lambda: __import__("uuid").uuid4().hex
)
company_id: Mapped[str] = mapped_column(
String(36), ForeignKey("companies.id"), nullable=False
)
created_by_agent_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("agents.id"), nullable=True
)
# identity
name: Mapped[str] = mapped_column(String(255), nullable=False)
slug: Mapped[str] = mapped_column(String(100), nullable=False, index=True)
season: Mapped[str] = mapped_column(String(50), nullable=False)
year: Mapped[int] = mapped_column(Integer, nullable=False)
status: Mapped[str] = mapped_column(
String(30), nullable=False, default="draft", index=True
) # draft | approved | active | paused | completed | archived
# branding
tagline: Mapped[Optional[str]] = mapped_column(String(500), nullable=True)
theme_color: Mapped[Optional[str]] = mapped_column(String(20), nullable=True)
description: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
visual_brief: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
# schedule
launch_date: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
end_date: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
publish_date: Mapped[Optional[datetime]] = mapped_column(DateTime, nullable=True)
# budget envelope
total_budget_vnd: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
currency: Mapped[str] = mapped_column(String(10), nullable=False, default="VND")
# targets
target_revenue_vnd: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
target_units: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
target_roi: Mapped[Optional[float]] = mapped_column(Integer, nullable=True) # x100
# rich metadata
extra: Mapped[dict] = mapped_column("extra", JSON, nullable=False, default=dict)
created_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
updated_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
...@@ -31,5 +31,7 @@ class Company(Base): ...@@ -31,5 +31,7 @@ class Company(Base):
feedback_data_sharing_consent_by_user_id: Mapped[Optional[str]] = mapped_column(String(255), nullable=True) feedback_data_sharing_consent_by_user_id: Mapped[Optional[str]] = mapped_column(String(255), nullable=True)
feedback_data_sharing_terms_version: Mapped[Optional[str]] = mapped_column(String(50), nullable=True) feedback_data_sharing_terms_version: Mapped[Optional[str]] = mapped_column(String(50), nullable=True)
brand_color: Mapped[Optional[str]] = mapped_column(String(20), nullable=True) brand_color: Mapped[Optional[str]] = mapped_column(String(20), nullable=True)
blueprint_name: Mapped[Optional[str]] = mapped_column(String(100), nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow) created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow)
updated_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow) updated_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow)
...@@ -13,11 +13,12 @@ from typing import Optional ...@@ -13,11 +13,12 @@ from typing import Optional
class CompanySecret(Base): class CompanySecret(Base):
__tablename__ = "company_secrets" __tablename__ = "company_secrets"
__table_args__ = ( __table_args__ = (
UniqueConstraint('company_id', 'key', name='uq_company_secrets_company_key'), UniqueConstraint('company_id', 'team_id', 'key', name='uq_company_secrets_company_key'),
) )
id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex) id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex)
company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False) company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False)
team_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("teams.id"), nullable=True)
key: Mapped[str] = mapped_column(String(255), nullable=False) key: Mapped[str] = mapped_column(String(255), nullable=False)
value_encrypted: Mapped[str] = mapped_column(Text, nullable=False) # encrypted value value_encrypted: Mapped[str] = mapped_column(Text, nullable=False) # encrypted value
description: Mapped[Optional[str]] = mapped_column(Text, nullable=True) description: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
......
...@@ -15,6 +15,7 @@ class CompanySkill(Base): ...@@ -15,6 +15,7 @@ class CompanySkill(Base):
id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex) id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex)
company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False) company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False)
team_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("teams.id"), nullable=True)
name: Mapped[str] = mapped_column(String(100), nullable=False) name: Mapped[str] = mapped_column(String(100), nullable=False)
description: Mapped[Optional[str]] = mapped_column(Text, nullable=True) description: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
level: Mapped[int] = mapped_column(Integer, nullable=False, default=1) # 1-100 scale level: Mapped[int] = mapped_column(Integer, nullable=False, default=1) # 1-100 scale
......
...@@ -26,5 +26,7 @@ class HeartbeatRun(Base): ...@@ -26,5 +26,7 @@ class HeartbeatRun(Base):
error_message: Mapped[Optional[str]] = mapped_column(Text, nullable=True) error_message: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
execution_cost_cents: Mapped[int] = mapped_column(Integer, nullable=False, default=0) execution_cost_cents: Mapped[int] = mapped_column(Integer, nullable=False, default=0)
execution_cost_currency: Mapped[str] = mapped_column(String(10), nullable=False, default="USD") execution_cost_currency: Mapped[str] = mapped_column(String(10), nullable=False, default="USD")
pipeline_id: Mapped[Optional[str]] = mapped_column(String(64), nullable=True, index=True)
pipeline_step: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow) created_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow)
updated_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow) updated_at: Mapped[datetime] = mapped_column(DateTime, nullable=False, default=datetime.utcnow)
"""
Material / Fabric record — used by R&D to specify fabric for each product spec.
"""
from sqlalchemy import String, Text, DateTime, ForeignKey, JSON, Integer, Numeric
from sqlalchemy.orm import Mapped, mapped_column
from datetime import datetime
from typing import Optional
from .database import Base
class Material(Base):
__tablename__ = "materials"
id: Mapped[str] = mapped_column(
String(36), primary_key=True, default=lambda: __import__("uuid").uuid4().hex
)
company_id: Mapped[str] = mapped_column(
String(36), ForeignKey("companies.id"), nullable=False, index=True
)
campaign_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("campaigns.id"), nullable=True, index=True
)
product_spec_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("product_specs.id"), nullable=True, index=True
)
created_by_agent_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("agents.id"), nullable=True
)
# identity
name: Mapped[str] = mapped_column(String(255), nullable=False)
material_type: Mapped[str] = mapped_column(
String(60), nullable=False
) # fabric | button | zipper | label | thread | packaging
sub_type: Mapped[Optional[str]] = mapped_column(String(80), nullable=True)
sku: Mapped[Optional[str]] = mapped_column(String(60), nullable=True)
# properties
composition: Mapped[Optional[str]] = mapped_column(String(200), nullable=True)
weight_gsm: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
origin: Mapped[Optional[str]] = mapped_column(String(100), nullable=True)
color: Mapped[Optional[str]] = mapped_column(String(40), nullable=True)
finish: Mapped[Optional[dict]] = mapped_column(JSON, nullable=True)
# supply
supplier: Mapped[Optional[str]] = mapped_column(String(200), nullable=True)
moq: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
unit_cost_vnd: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
lead_time_days: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
stock_qty: Mapped[Optional[int]] = mapped_column(Integer, nullable=False, default=0)
# approval
rd_approved: Mapped[bool] = mapped_column(Integer, nullable=False, default=0)
notes: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
extra: Mapped[dict] = mapped_column("extra", JSON, nullable=False, default=dict)
created_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
updated_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
...@@ -15,6 +15,7 @@ class Memory(Base): ...@@ -15,6 +15,7 @@ class Memory(Base):
id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex) id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex)
company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False) company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False)
team_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("teams.id"), nullable=True)
agent_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("agents.id"), nullable=True) agent_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("agents.id"), nullable=True)
issue_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("issues.id"), nullable=True) issue_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("issues.id"), nullable=True)
namespace: Mapped[str] = mapped_column(String(100), nullable=False) namespace: Mapped[str] = mapped_column(String(100), nullable=False)
......
...@@ -15,6 +15,7 @@ class PluginConfig(Base): ...@@ -15,6 +15,7 @@ class PluginConfig(Base):
id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex) id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex)
company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False) company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False)
team_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("teams.id"), nullable=True)
plugin_id: Mapped[str] = mapped_column(String(36), ForeignKey("plugins.id"), nullable=False) plugin_id: Mapped[str] = mapped_column(String(36), ForeignKey("plugins.id"), nullable=False)
config: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict) config: Mapped[dict] = mapped_column(JSON, nullable=False, default=dict)
is_enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True) is_enabled: Mapped[bool] = mapped_column(Boolean, nullable=False, default=True)
......
"""
Product Specification model — per-product detail for campaigns.
"""
from sqlalchemy import String, Text, DateTime, ForeignKey, JSON, Integer, Numeric
from sqlalchemy.orm import Mapped, mapped_column, relationship
from datetime import datetime
from typing import Optional
from .database import Base
class ProductSpec(Base):
__tablename__ = "product_specs"
id: Mapped[str] = mapped_column(
String(36), primary_key=True, default=lambda: __import__("uuid").uuid4().hex
)
company_id: Mapped[str] = mapped_column(
String(36), ForeignKey("companies.id"), nullable=False, index=True
)
campaign_id: Mapped[str] = mapped_column(
String(36), ForeignKey("campaigns.id"), nullable=False, index=True
)
created_by_agent_id: Mapped[Optional[str]] = mapped_column(
String(36), ForeignKey("agents.id"), nullable=True
)
# identity
name: Mapped[str] = mapped_column(String(255), nullable=False)
sku: Mapped[str] = mapped_column(String(60), nullable=False, index=True)
product_type: Mapped[str] = mapped_column(
String(40), nullable=False
) # graphic_tshirt | linen_shorts
status: Mapped[str] = mapped_column(
String(30), nullable=False, default="draft"
) # draft | approved | in_production | ready | discontinued
# designer spec
size_chart: Mapped[Optional[dict]] = mapped_column(JSON, nullable=True)
colorways: Mapped[list] = mapped_column(JSON, nullable=False, default=list)
style_notes: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
design_files: Mapped[list] = mapped_column(JSON, nullable=False, default=list)
# pricing
cost_price_vnd: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
retail_price_vnd: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
# inventory targets
production_qty: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
stock_safety_qty: Mapped[Optional[int]] = mapped_column(Integer, nullable=True, default=0)
extra: Mapped[dict] = mapped_column("extra", JSON, nullable=False, default=dict)
created_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
updated_at: Mapped[datetime] = mapped_column(
DateTime, nullable=False, default=datetime.utcnow
)
"""
Teams model.
"""
from sqlalchemy import String, ForeignKey, DateTime
from sqlalchemy.orm import Mapped, mapped_column
from datetime import datetime
from typing import Optional
from .database import Base
class Team(Base):
__tablename__ = "teams"
id: Mapped[str] = mapped_column(String(36), primary_key=True, default=lambda: __import__('uuid').uuid4().hex)
company_id: Mapped[str] = mapped_column(String(36), ForeignKey("companies.id"), nullable=False)
name: Mapped[str] = mapped_column(String(255), nullable=False)
slug: Mapped[Optional[str]] = mapped_column(String(100), nullable=True)
parent_team_id: Mapped[Optional[str]] = mapped_column(String(36), ForeignKey("teams.id"), 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)
...@@ -41,6 +41,12 @@ from .heartbeat_run_event import * ...@@ -41,6 +41,12 @@ from .heartbeat_run_event import *
from .vote import * from .vote import *
from .blueprint import * from .blueprint import *
from .common import TimestampMixin from .common import TimestampMixin
from .campaign import *
from .product_spec import *
from .material import *
from .ad_plan import *
from .budget_item import *
__all__ = [ __all__ = [
...@@ -180,4 +186,24 @@ __all__ = [ ...@@ -180,4 +186,24 @@ __all__ = [
"BlueprintRole", "BlueprintRole",
"BlueprintPipelineNode", "BlueprintPipelineNode",
"CompanyBlueprint", "CompanyBlueprint",
"CampaignBase",
"CampaignCreate",
"CampaignUpdate",
"CampaignResponse",
"ProductSpecBase",
"ProductSpecCreate",
"ProductSpecUpdate",
"ProductSpecResponse",
"MaterialBase",
"MaterialCreate",
"MaterialUpdate",
"MaterialResponse",
"AdPlanBase",
"AdPlanCreate",
"AdPlanUpdate",
"AdPlanResponse",
"BudgetItemBase",
"BudgetItemCreate",
"BudgetItemUpdate",
"BudgetItemResponse",
] ]
\ No newline at end of file
"""
Ad Plan schemas.
"""
from datetime import datetime
from typing import Optional, Any, List
from pydantic import BaseModel, ConfigDict, Field
from .common import TimestampMixin
class AdPlanBase(BaseModel):
company_id: str
campaign_id: str
created_by_agent_id: Optional[str] = None
name: str = Field(..., min_length=1, max_length=255)
platform: str = Field(..., pattern="^(facebook|instagram|tiktok|google|other)$")
objective: Optional[str] = Field(None, max_length=80)
status: str = Field("draft", pattern="^(draft|scheduled|running|paused|completed)$")
ad_copy: Optional[str] = None
creative_brief: Optional[str] = None
target_audience: Optional[dict] = None
hashtags: list = Field(default_factory=list)
scheduled_start: Optional[datetime] = None
scheduled_end: Optional[datetime] = None
allocated_budget_vnd: Optional[int] = Field(None, ge=0)
spent_vnd: int = 0
impressions: Optional[int] = Field(None, ge=0)
clicks: Optional[int] = Field(None, ge=0)
conversions: Optional[int] = Field(None, ge=0)
revenue_vnd: Optional[int] = Field(None, ge=0)
extra: dict = Field(default_factory=dict)
model_config = ConfigDict(from_attributes=True)
class AdPlanCreate(AdPlanBase):
model_config = ConfigDict(from_attributes=True)
class AdPlanUpdate(BaseModel):
model_config = ConfigDict(from_attributes=True)
name: Optional[str] = Field(None, min_length=1, max_length=255)
platform: Optional[str] = Field(None, pattern="^(facebook|instagram|tiktok|google|other)$")
objective: Optional[str] = Field(None, max_length=80)
status: Optional[str] = Field(None, pattern="^(draft|scheduled|running|paused|completed)$")
ad_copy: Optional[str] = None
creative_brief: Optional[str] = None
target_audience: Optional[dict] = None
hashtags: Optional[list] = None
scheduled_start: Optional[datetime] = None
scheduled_end: Optional[datetime] = None
allocated_budget_vnd: Optional[int] = Field(None, ge=0)
spent_vnd: Optional[int] = Field(None, ge=0)
impressions: Optional[int] = Field(None, ge=0)
clicks: Optional[int] = Field(None, ge=0)
conversions: Optional[int] = Field(None, ge=0)
revenue_vnd: Optional[int] = Field(None, ge=0)
extra: Optional[dict] = None
class AdPlanResponse(AdPlanBase, TimestampMixin):
id: str
...@@ -10,7 +10,7 @@ from .common import TimestampMixin ...@@ -10,7 +10,7 @@ from .common import TimestampMixin
class ApprovalBase(BaseModel): class ApprovalBase(BaseModel):
company_id: str company_id: str
resource_type: str = Field(..., pattern="^(agent|issue|budget)$") resource_type: str = Field(..., pattern="^(agent|issue|budget|pipeline_stage)$")
resource_id: str resource_id: str
required_approvers: int = Field(1, ge=1) required_approvers: int = Field(1, ge=1)
created_by: Optional[str] = None created_by: Optional[str] = None
...@@ -29,6 +29,9 @@ class ApprovalUpdate(BaseModel): ...@@ -29,6 +29,9 @@ class ApprovalUpdate(BaseModel):
class ApprovalResponse(ApprovalBase, TimestampMixin): class ApprovalResponse(ApprovalBase, TimestampMixin):
id: str id: str
# Approval model has no updated_at column — override the mixin's required field
updated_at: Optional[datetime] = None
resolved_at: Optional[datetime] = None resolved_at: Optional[datetime] = None
resolution: Optional[str] = None resolution: Optional[str] = None
resolved_by: Optional[str] = None resolved_by: Optional[str] = None
\ No newline at end of file resolution_note: Optional[str] = None
\ No newline at end of file
...@@ -33,6 +33,9 @@ class CompanyBlueprint(BaseModel): ...@@ -33,6 +33,9 @@ class CompanyBlueprint(BaseModel):
roles: List[BlueprintRole] = Field(default_factory=list) roles: List[BlueprintRole] = Field(default_factory=list)
pipeline: List[BlueprintPipelineNode] = Field(default_factory=list) pipeline: List[BlueprintPipelineNode] = Field(default_factory=list)
nocobase_collections: List[Dict] = Field(default_factory=list) nocobase_collections: List[Dict] = Field(default_factory=list)
plugins: Optional[List[Dict[str, Any]]] = Field(default_factory=list)
skills: Optional[List[Any]] = Field(default_factory=list)
env_vars: Optional[Dict[str, str]] = Field(default_factory=dict)
model_config = ConfigDict(from_attributes=True) model_config = ConfigDict(from_attributes=True)
......
"""
Budget Item schemas.
"""
from datetime import datetime
from typing import Optional, Any, List
from pydantic import BaseModel, ConfigDict, Field
from .common import TimestampMixin
class BudgetItemBase(BaseModel):
company_id: str
campaign_id: Optional[str] = None
product_spec_id: Optional[str] = None
ad_plan_id: Optional[str] = None
created_by_agent_id: Optional[str] = None
category: str = Field(..., pattern="^(production|material|ad_spend|design|logistics|hr|other)$")
label: str = Field(..., min_length=1, max_length=255)
estimated_vnd: int = Field(..., ge=0)
actual_vnd: int = Field(0, ge=0)
status: str = Field("planned", pattern="^(planned|approved|committed|paid|cancelled)$")
approved_by_agent_id: Optional[str] = None
approved_at: Optional[datetime] = None
notes: Optional[str] = None
attachments: list = Field(default_factory=list)
extra: dict = Field(default_factory=dict)
model_config = ConfigDict(from_attributes=True)
class BudgetItemCreate(BudgetItemBase):
model_config = ConfigDict(from_attributes=True)
class BudgetItemUpdate(BaseModel):
model_config = ConfigDict(from_attributes=True)
campaign_id: Optional[str] = None
product_spec_id: Optional[str] = None
ad_plan_id: Optional[str] = None
category: Optional[str] = Field(None, pattern="^(production|material|ad_spend|design|logistics|hr|other)$")
label: Optional[str] = Field(None, min_length=1, max_length=255)
estimated_vnd: Optional[int] = Field(None, ge=0)
actual_vnd: Optional[int] = Field(None, ge=0)
status: Optional[str] = Field(None, pattern="^(planned|approved|committed|paid|cancelled)$")
approved_by_agent_id: Optional[str] = None
approved_at: Optional[datetime] = None
notes: Optional[str] = None
attachments: Optional[list] = None
extra: Optional[dict] = None
class BudgetItemResponse(BudgetItemBase, TimestampMixin):
id: str
"""
Campaign schemas.
"""
from datetime import datetime
from typing import Optional, Any
from pydantic import BaseModel, ConfigDict, Field
from .common import TimestampMixin
class CampaignBase(BaseModel):
company_id: str
created_by_agent_id: Optional[str] = None
name: str = Field(..., min_length=1, max_length=255)
slug: str = Field(..., min_length=1, max_length=100, pattern=r"^[a-z0-9-]+$")
season: str = Field(..., min_length=1, max_length=50)
year: int = Field(..., ge=2000, le=2100)
status: str = Field("draft", pattern="^(draft|approved|active|paused|completed|archived)$")
tagline: Optional[str] = Field(None, max_length=500)
theme_color: Optional[str] = Field(None, max_length=20)
description: Optional[str] = None
visual_brief: Optional[str] = None
launch_date: Optional[datetime] = None
end_date: Optional[datetime] = None
publish_date: Optional[datetime] = None
total_budget_vnd: Optional[int] = Field(None, ge=0)
currency: str = "VND"
target_revenue_vnd: Optional[int] = Field(None, ge=0)
target_units: Optional[int] = Field(None, ge=0)
target_roi: Optional[float] = Field(None, ge=0)
extra: dict = Field(default_factory=dict)
model_config = ConfigDict(from_attributes=True)
class CampaignCreate(CampaignBase):
model_config = ConfigDict(from_attributes=True)
class CampaignUpdate(BaseModel):
model_config = ConfigDict(from_attributes=True)
name: Optional[str] = Field(None, min_length=1, max_length=255)
slug: Optional[str] = Field(None, min_length=1, max_length=100, pattern=r"^[a-z0-9-]+$")
season: Optional[str] = Field(None, min_length=1, max_length=50)
year: Optional[int] = Field(None, ge=2000, le=2100)
status: Optional[str] = Field(None, pattern="^(draft|approved|active|paused|completed|archived)$")
tagline: Optional[str] = Field(None, max_length=500)
theme_color: Optional[str] = Field(None, max_length=20)
description: Optional[str] = None
visual_brief: Optional[str] = None
launch_date: Optional[datetime] = None
end_date: Optional[datetime] = None
publish_date: Optional[datetime] = None
total_budget_vnd: Optional[int] = Field(None, ge=0)
target_revenue_vnd: Optional[int] = Field(None, ge=0)
target_units: Optional[int] = Field(None, ge=0)
target_roi: Optional[float] = Field(None, ge=0)
extra: Optional[dict] = None
class CampaignResponse(CampaignBase, TimestampMixin):
id: str
...@@ -14,6 +14,8 @@ class CompanyBase(BaseModel): ...@@ -14,6 +14,8 @@ class CompanyBase(BaseModel):
issue_prefix: Optional[str] = Field(None, min_length=2, max_length=10) issue_prefix: Optional[str] = Field(None, min_length=2, max_length=10)
budget_monthly_cents: int = Field(0, ge=0) budget_monthly_cents: int = Field(0, ge=0)
require_board_approval_for_new_agents: bool = False require_board_approval_for_new_agents: bool = False
blueprint_name: Optional[str] = None
class CompanyCreate(CompanyBase): class CompanyCreate(CompanyBase):
......
"""
Material schemas.
"""
from datetime import datetime
from typing import Optional, Any
from pydantic import BaseModel, ConfigDict, Field
from .common import TimestampMixin
class MaterialBase(BaseModel):
company_id: str
campaign_id: Optional[str] = None
product_spec_id: Optional[str] = None
created_by_agent_id: Optional[str] = None
name: str = Field(..., min_length=1, max_length=255)
material_type: str = Field(..., pattern="^(fabric|button|zipper|label|thread|packaging|other)$")
sub_type: Optional[str] = Field(None, max_length=80)
sku: Optional[str] = Field(None, max_length=60)
composition: Optional[str] = Field(None, max_length=200)
weight_gsm: Optional[int] = Field(None, ge=0)
origin: Optional[str] = Field(None, max_length=100)
color: Optional[str] = Field(None, max_length=40)
finish: Optional[dict] = None
supplier: Optional[str] = Field(None, max_length=200)
moq: Optional[int] = Field(None, ge=0)
unit_cost_vnd: Optional[int] = Field(None, ge=0)
lead_time_days: Optional[int] = Field(None, ge=0)
stock_qty: int = 0
rd_approved: bool = False
notes: Optional[str] = None
extra: dict = Field(default_factory=dict)
model_config = ConfigDict(from_attributes=True)
class MaterialCreate(MaterialBase):
model_config = ConfigDict(from_attributes=True)
class MaterialUpdate(BaseModel):
model_config = ConfigDict(from_attributes=True)
name: Optional[str] = Field(None, min_length=1, max_length=255)
material_type: Optional[str] = Field(None, pattern="^(fabric|button|zipper|label|thread|packaging|other)$")
sub_type: Optional[str] = Field(None, max_length=80)
sku: Optional[str] = Field(None, max_length=60)
composition: Optional[str] = Field(None, max_length=200)
weight_gsm: Optional[int] = Field(None, ge=0)
origin: Optional[str] = Field(None, max_length=100)
color: Optional[str] = Field(None, max_length=40)
finish: Optional[dict] = None
supplier: Optional[str] = Field(None, max_length=200)
moq: Optional[int] = Field(None, ge=0)
unit_cost_vnd: Optional[int] = Field(None, ge=0)
lead_time_days: Optional[int] = Field(None, ge=0)
stock_qty: Optional[int] = None
rd_approved: Optional[bool] = None
notes: Optional[str] = None
extra: Optional[dict] = None
class MaterialResponse(MaterialBase, TimestampMixin):
id: str
"""
Product Specification schemas.
"""
from datetime import datetime
from typing import Optional, Any, List
from pydantic import BaseModel, ConfigDict, Field
from .common import TimestampMixin
class ProductSpecBase(BaseModel):
company_id: str
campaign_id: str
created_by_agent_id: Optional[str] = None
name: str = Field(..., min_length=1, max_length=255)
sku: str = Field(..., min_length=1, max_length=60)
product_type: str = Field(..., pattern="^(graphic_tshirt|linen_shorts)$")
status: str = Field("draft", pattern="^(draft|approved|in_production|ready|discontinued)$")
size_chart: Optional[dict] = None
colorways: list = Field(default_factory=list)
style_notes: Optional[str] = None
design_files: list = Field(default_factory=list)
cost_price_vnd: Optional[int] = Field(None, ge=0)
retail_price_vnd: Optional[int] = Field(None, ge=0)
production_qty: Optional[int] = Field(None, ge=0)
stock_safety_qty: int = 0
extra: dict = Field(default_factory=dict)
model_config = ConfigDict(from_attributes=True)
class ProductSpecCreate(ProductSpecBase):
model_config = ConfigDict(from_attributes=True)
class ProductSpecUpdate(BaseModel):
model_config = ConfigDict(from_attributes=True)
name: Optional[str] = Field(None, min_length=1, max_length=255)
sku: Optional[str] = Field(None, min_length=1, max_length=60)
product_type: Optional[str] = Field(None, pattern="^(graphic_tshirt|linen_shorts)$")
status: Optional[str] = Field(None, pattern="^(draft|approved|in_production|ready|discontinued)$")
size_chart: Optional[dict] = None
colorways: Optional[list] = None
style_notes: Optional[str] = None
design_files: Optional[list] = None
cost_price_vnd: Optional[int] = Field(None, ge=0)
retail_price_vnd: Optional[int] = Field(None, ge=0)
production_qty: Optional[int] = Field(None, ge=0)
stock_safety_qty: Optional[int] = None
extra: Optional[dict] = None
class ProductSpecResponse(ProductSpecBase, TimestampMixin):
id: str
...@@ -64,6 +64,22 @@ async def lifespan(app: FastAPI): ...@@ -64,6 +64,22 @@ async def lifespan(app: FastAPI):
except Exception as e: except Exception as e:
logger.error(f"❌ Failed to start MCP clients: {e}", exc_info=True) logger.error(f"❌ Failed to start MCP clients: {e}", exc_info=True)
# Start Heartbeat Scheduler Loop if enabled (fallback/primary scheduler when Celery Beat is not active)
if settings.heartbeat_scheduler_enabled:
async def scheduler_loop():
await asyncio.sleep(5)
from services.routine_scheduler import RoutineSchedulerService
logger.info("⏰ Heartbeat Routine Scheduler Loop started")
while True:
try:
await RoutineSchedulerService.tick()
except Exception as e:
logger.error(f"❌ Error in routine scheduler tick: {e}", exc_info=True)
interval_s = getattr(settings, 'heartbeat_scheduler_interval_ms', 1000) / 1000.0
await asyncio.sleep(max(interval_s, 1.0))
asyncio.create_task(scheduler_loop())
yield yield
# ── Shutdown ── # ── Shutdown ──
...@@ -177,15 +193,22 @@ class CompanyPathRewriteMiddleware: ...@@ -177,15 +193,22 @@ class CompanyPathRewriteMiddleware:
if scope["type"] in ("http", "websocket"): if scope["type"] in ("http", "websocket"):
path = scope.get("path", "") path = scope.get("path", "")
match = re.match(r"^(/api)?/companies/([^/]+)/(.+)$", path) match = re.match(r"^(/api)?/companies/([^/]+)/(.+)$", path)
if match: if match and match.group(2) != "blueprints":
api_prefix = match.group(1) or "" api_prefix = match.group(1) or ""
company_id = match.group(2) company_id = match.group(2)
resource_path = match.group(3) resource_path = match.group(3)
# Check if this is a path we should NOT rewrite (runs API expects prefix) # Check if this is a path we should NOT rewrite (runs, JSON org, budgets expect prefix)
is_runs = resource_path.startswith("runs") should_not_rewrite = (
resource_path.startswith("runs")
or resource_path == "org"
or resource_path.startswith("org/")
or resource_path.startswith("budgets")
or resource_path.startswith("budget-incidents")
)
print(f"DEBUG MIDDLEWARE match: path={path}, resource_path={resource_path}, should_not_rewrite={should_not_rewrite}", flush=True)
if not is_runs: if not should_not_rewrite:
if resource_path.startswith("skills"): if resource_path.startswith("skills"):
resource_path = re.sub(r"^skills(\b|/)", "company-skills\\1", resource_path) resource_path = re.sub(r"^skills(\b|/)", "company-skills\\1", resource_path)
elif resource_path.startswith("members"): elif resource_path.startswith("members"):
...@@ -194,8 +217,9 @@ class CompanyPathRewriteMiddleware: ...@@ -194,8 +217,9 @@ class CompanyPathRewriteMiddleware:
resource_path = re.sub(r"^user-directory(\b|/)", "company-memberships/user-directory\\1", resource_path) resource_path = re.sub(r"^user-directory(\b|/)", "company-memberships/user-directory\\1", resource_path)
elif resource_path.startswith("join-requests"): elif resource_path.startswith("join-requests"):
resource_path = re.sub(r"^join-requests(\b|/)", "company-memberships/join-requests\\1", resource_path) resource_path = re.sub(r"^join-requests(\b|/)", "company-memberships/join-requests\\1", resource_path)
elif resource_path.startswith("org"): elif resource_path.startswith("org-svg"):
resource_path = re.sub(r"^org(\b|/)", "org-chart-svg\\1", resource_path) resource_path = re.sub(r"^org-svg(\b|/)", "org-chart-svg\\1", resource_path)
# Rewrite path # Rewrite path
scope["path"] = f"{api_prefix}/{resource_path}" scope["path"] = f"{api_prefix}/{resource_path}"
......
...@@ -11,7 +11,7 @@ from typing import List, Dict, Any, Optional ...@@ -11,7 +11,7 @@ from typing import List, Dict, Any, Optional
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select from sqlalchemy import select
from models import Company, Agent, CompanyMembership from models import Company, Agent, CompanyMembership, CompanySecret, CompanySkill, Plugin, PluginConfig
from schemas.blueprint import CompanyBlueprint, BlueprintRole, BlueprintPipelineNode from schemas.blueprint import CompanyBlueprint, BlueprintRole, BlueprintPipelineNode
from services.activity_logger import log_activity from services.activity_logger import log_activity
...@@ -210,6 +210,84 @@ class BlueprintService: ...@@ -210,6 +210,84 @@ class BlueprintService:
else: else:
raise ValueError(f"reports_to role '{parent_role}' not found") raise ValueError(f"reports_to role '{parent_role}' not found")
# 4. Provision environment secrets
if blueprint.env_vars:
for key, value in blueprint.env_vars.items():
is_sensitive = any(kw in key.lower() for kw in ["key", "token", "secret", "password", "private", "credential"])
secret = CompanySecret(
id=uuid.uuid4().hex,
company_id=db_company.id,
key=key,
value_encrypted=value,
description=f"Environment variable {key} provisioned from blueprint",
is_sensitive=is_sensitive,
created_by=owner_id
)
db.add(secret)
# 5. Provision custom plugins
if blueprint.plugins:
for plugin_spec in blueprint.plugins:
name = plugin_spec.get("name")
if not name:
continue
version = plugin_spec.get("version", "1.0.0")
description = plugin_spec.get("description", f"{name} plugin")
config = plugin_spec.get("config", {})
# Check if global plugin exists
plugin_result = await db.execute(select(Plugin).where(Plugin.name == name))
db_plugin = plugin_result.scalar_one_or_none()
if not db_plugin:
db_plugin = Plugin(
id=uuid.uuid4().hex,
name=name,
version=version,
description=description,
enabled=True,
config_schema={}
)
db.add(db_plugin)
await db.flush()
# Configure for this company
plugin_config = PluginConfig(
id=uuid.uuid4().hex,
company_id=db_company.id,
plugin_id=db_plugin.id,
config=config,
is_enabled=True
)
db.add(plugin_config)
# 6. Provision company skills
if blueprint.skills:
for skill_spec in blueprint.skills:
if isinstance(skill_spec, str):
name = skill_spec
level = 100
description = f"Capability: {name}"
is_custom = False
elif isinstance(skill_spec, dict):
name = skill_spec.get("name")
if not name:
continue
level = skill_spec.get("level", 100)
description = skill_spec.get("description", f"Capability: {name}")
is_custom = skill_spec.get("is_custom", True)
else:
continue
skill = CompanySkill(
id=uuid.uuid4().hex,
company_id=db_company.id,
name=name,
description=description,
level=level,
is_custom=is_custom
)
db.add(skill)
await db.flush() await db.flush()
# Log activity # Log activity
......
This diff is collapsed.
...@@ -50,6 +50,24 @@ class NocoBaseConnector: ...@@ -50,6 +50,24 @@ class NocoBaseConnector:
return connector return connector
async def ping(self) -> bool:
"""Health-check NocoBase reachability (no auth required).
Hits the public app-info endpoint. Returns True only if the server
answers with a 2xx. Used to gate provisioning + live integration tests
so a dead NocoBase is detected loudly instead of silently swallowed.
"""
url = f"{self.base_url}/api/app:getInfo"
try:
async with httpx.AsyncClient() as client:
response = await client.get(url, timeout=5.0)
if 200 <= response.status_code < 300:
return True
logger.warning(f"NocoBase ping non-2xx: {response.status_code} at {url}")
except Exception as e:
logger.warning(f"NocoBase ping failed at {url}: {e}")
return False
async def ensure_collection(self, name: str, fields: List[Dict[str, Any]]) -> bool: async def ensure_collection(self, name: str, fields: List[Dict[str, Any]]) -> bool:
""" """
Ensure that a NocoBase collection exists, creating it if not. Ensure that a NocoBase collection exists, creating it if not.
......
...@@ -478,19 +478,25 @@ class ProjectOrchestrator: ...@@ -478,19 +478,25 @@ class ProjectOrchestrator:
agent_id=agent_id, agent_id=agent_id,
status="running", status="running",
started_at=start, started_at=start,
pipeline_id=self.pipeline_id,
pipeline_step=step,
) )
db.add(db_run) db.add(db_run)
logger.info(f"[Pipeline {self.pipeline_id}] Step {step}/{total}: Waking {name} (run_id={run_id})") logger.info(f"[Pipeline {self.pipeline_id}] Step {step}/{total}: Waking {name} (run_id={run_id})")
try: try:
system_message = (
f"You are {name}, playing the role of {role}.\n"
f"After completing your work, you MUST write a detailed markdown summary of your findings, code, decisions, or results to a file named 'deliverable.md' in your workspace root directory. Keep the file name exactly 'deliverable.md' so it can be collected as your deliverable."
)
await run_agent_in_background( await run_agent_in_background(
run_id=run_id, run_id=run_id,
agent_id=agent_id, agent_id=agent_id,
company_id=self.company_id, company_id=self.company_id,
model_name=model_name, model_name=model_name,
user_message=prompt, user_message=prompt,
system_message=f"You are {name}, playing the role of {role}.", system_message=system_message,
conversation_history=None, conversation_history=None,
) )
except Exception as e: except Exception as e:
......
"""
Runtime Environment Resolver.
Resolves environment variables and enabled toolsets for a given agent execution.
"""
import logging
from typing import Optional, Dict, List
from sqlalchemy.future import select
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.orm import selectinload
from models.company_secrets import CompanySecret
from models.plugin_config import PluginConfig
from models.plugins import Plugin
from common.encryption import decrypt_api_key
logger = logging.getLogger(__name__)
async def resolve_runtime_env(db: AsyncSession, company_id: str, team_id: Optional[str] = None) -> Dict[str, str]:
"""
Resolve environment variables for an agent run.
Team secrets override company secrets if they share the same key.
"""
env_vars = {}
# Query company secrets
query = select(CompanySecret).where(CompanySecret.company_id == company_id)
if team_id:
query = query.where(
(CompanySecret.team_id == team_id) | (CompanySecret.team_id.is_(None))
)
else:
query = query.where(CompanySecret.team_id.is_(None))
result = await db.execute(query)
secrets = result.scalars().all()
# Process secrets. Team specific secrets (team_id is not None) override company ones (team_id is None).
# Sorting by team_id ensures that None comes first, then specific team_id overrides.
secrets = sorted(secrets, key=lambda x: 1 if x.team_id else 0)
for secret in secrets:
val = secret.value_encrypted
try:
val = decrypt_api_key(val)
except Exception:
# Fallback to raw if not actually encrypted (e.g., from blueprint provisioning without encryption)
pass
env_vars[secret.key] = val
return env_vars
async def resolve_enabled_toolsets(db: AsyncSession, company_id: str, team_id: Optional[str] = None) -> Optional[List[str]]:
"""
Resolve enabled toolsets (plugins) for an agent run.
If there are no plugin configs for this company/team, returns None (no filtering).
If there are plugin configs, returns a list of enabled plugin names.
Team configs override company configs.
"""
# Check if there are any plugin configs at all for this company
check_query = select(PluginConfig).where(PluginConfig.company_id == company_id)
check_result = await db.execute(check_query.limit(1))
if not check_result.scalar_one_or_none():
return None # No configs -> no filtering (backward compatibility)
query = select(PluginConfig, Plugin.name).join(
Plugin, Plugin.id == PluginConfig.plugin_id
).where(
PluginConfig.company_id == company_id
)
if team_id:
query = query.where(
(PluginConfig.team_id == team_id) | (PluginConfig.team_id.is_(None))
)
else:
query = query.where(PluginConfig.team_id.is_(None))
result = await db.execute(query)
rows = result.all()
# Deduplicate by plugin_id, prioritizing team_id configs
resolved_configs = {}
for cfg, plugin_name in sorted(rows, key=lambda x: 1 if x[0].team_id else 0):
resolved_configs[cfg.plugin_id] = (cfg, plugin_name)
# Gather enabled plugin names
enabled_toolsets = []
for cfg, plugin_name in resolved_configs.values():
if cfg.is_enabled:
enabled_toolsets.append(plugin_name)
return enabled_toolsets
"""
Issue-Driven Swarm service.
When an agent finishes a parent issue it may emit follow-up work ("child specs").
This module turns those specs into real child Issue rows (linked via parent_id /
request_depth / origin_kind="agent_spawn") and dispatches a queued HeartbeatRun
for each child so the swarm propagates through the existing checkout/wake path.
A hard depth limit (MAX_SWARM_DEPTH) prevents runaway recursion: an agent that
keeps spawning children eventually hits the cap and no further children are made.
"""
import logging
import uuid
from datetime import datetime
from typing import Optional
from sqlalchemy.ext.asyncio import AsyncSession
from models import Issue, HeartbeatRun
logger = logging.getLogger(__name__)
# Maximum spawn depth. A root issue is depth 0; its children depth 1, etc.
# At request_depth >= MAX_SWARM_DEPTH we refuse to spawn further children.
MAX_SWARM_DEPTH = 3
async def spawn_child_issues(
db: AsyncSession,
parent_issue: Issue,
child_specs: list[dict],
created_by_agent_id: Optional[str] = None,
origin_run_id: Optional[str] = None,
) -> list[Issue]:
"""Create child issues under ``parent_issue`` from agent-emitted specs.
Each spec is a dict with at least ``title`` (required). Optional keys:
``description``, ``priority``, ``assignee_agent_id``.
Returns the created Issue rows (flushed, not committed — caller commits).
Returns an empty list if the parent is already at/over MAX_SWARM_DEPTH,
or if no spec carries a title.
"""
if parent_issue.request_depth >= MAX_SWARM_DEPTH:
logger.warning(
"[Swarm] Refusing to spawn: parent issue %s at depth %d >= MAX_SWARM_DEPTH %d",
parent_issue.id, parent_issue.request_depth, MAX_SWARM_DEPTH,
)
return []
children: list[Issue] = []
for spec in child_specs:
title = (spec or {}).get("title")
if not title:
logger.warning("[Swarm] Skipping child spec without title under %s", parent_issue.id)
continue
child = Issue(
company_id=parent_issue.company_id,
project_id=parent_issue.project_id,
goal_id=parent_issue.goal_id,
parent_id=parent_issue.id,
title=title,
description=spec.get("description"),
status="open",
priority=spec.get("priority", parent_issue.priority),
assignee_agent_id=spec.get("assignee_agent_id") or parent_issue.assignee_agent_id,
created_by_agent_id=created_by_agent_id,
origin_kind="agent_spawn",
origin_id=parent_issue.id,
origin_run_id=origin_run_id,
request_depth=parent_issue.request_depth + 1,
)
db.add(child)
children.append(child)
if children:
await db.flush()
for c in children:
await db.refresh(c)
logger.info(
"[Swarm] Spawned %d child issue(s) under %s (depth %d)",
len(children), parent_issue.id, parent_issue.request_depth + 1,
)
return children
async def dispatch_child_issues(
db: AsyncSession,
children: list[Issue],
) -> list[HeartbeatRun]:
"""Queue a HeartbeatRun for each child that has an assignee agent.
Mirrors the wakeup_checkout path so spawned issues actually get executed.
Children without an assignee are left for a router/human to pick up.
Returns the created (queued) runs, flushed but not committed.
"""
runs: list[HeartbeatRun] = []
for child in children:
if not child.assignee_agent_id:
continue
run = HeartbeatRun(
id=uuid.uuid4().hex,
company_id=child.company_id,
agent_id=child.assignee_agent_id,
issue_id=child.id,
status="queued",
execution_cost_cents=0,
execution_cost_currency="USD",
)
db.add(run)
child.checkout_run_id = run.id
child.monitor_wake_requested_at = datetime.utcnow()
runs.append(run)
if runs:
await db.flush()
logger.info("[Swarm] Dispatched %d queued run(s) for spawned children", len(runs))
return runs
...@@ -21,6 +21,47 @@ def _run_async(coro): ...@@ -21,6 +21,47 @@ def _run_async(coro):
loop.close() loop.close()
async def _maybe_spawn_children(run_id: str, agent_id: str, workspace):
"""Issue-Driven Swarm hook.
After an issue's run completes, if the agent wrote a `child_issues.json`
array into its workspace, turn each entry into a child Issue (linked via
parent_id / request_depth) and queue a run for it. Depth is capped inside
swarm_service so this cannot recurse forever.
"""
import json
from pathlib import Path
from sqlalchemy import select
from models.database import get_db_context
from models.heartbeat_runs import HeartbeatRun
from models.issues import Issue
from services.swarm_service import spawn_child_issues, dispatch_child_issues
specs_file = Path(workspace) / "child_issues.json"
if not specs_file.exists():
return
try:
specs = json.loads(specs_file.read_text(encoding="utf-8"))
except Exception:
logger.warning(f"[Swarm] child_issues.json invalid for run {run_id}")
return
if not isinstance(specs, list) or not specs:
return
async with get_db_context() as db:
run = (await db.execute(select(HeartbeatRun).where(HeartbeatRun.id == run_id))).scalar_one_or_none()
if not run or not run.issue_id:
return
parent = (await db.execute(select(Issue).where(Issue.id == run.issue_id))).scalar_one_or_none()
if not parent:
return
children = await spawn_child_issues(
db, parent, specs, created_by_agent_id=agent_id, origin_run_id=run_id
)
await dispatch_child_issues(db, children)
await db.commit()
@shared_task( @shared_task(
bind=True, bind=True,
name="tasks.run_agent", name="tasks.run_agent",
...@@ -83,6 +124,12 @@ def run_agent_task(self, run_id: str, agent_id: str, company_id: str, model_name ...@@ -83,6 +124,12 @@ def run_agent_task(self, run_id: str, agent_id: str, company_id: str, model_name
) )
await db.commit() await db.commit()
# 5b. Issue-Driven Swarm: spawn + dispatch any child issues the agent emitted
try:
await _maybe_spawn_children(run_id, agent_id, workspace)
except Exception as swarm_exc:
logger.error(f"[Swarm] spawn failed for run {run_id}: {swarm_exc}", exc_info=True)
logger.info(f"[Celery] Agent task completed: run_id={run_id}") logger.info(f"[Celery] Agent task completed: run_id={run_id}")
return {"status": "completed", "run_id": run_id} return {"status": "completed", "run_id": run_id}
...@@ -109,3 +156,84 @@ def run_agent_task(self, run_id: str, agent_id: str, company_id: str, model_name ...@@ -109,3 +156,84 @@ def run_agent_task(self, run_id: str, agent_id: str, company_id: str, model_name
os.environ.pop("AGENT_WORKSPACE", None) os.environ.pop("AGENT_WORKSPACE", None)
return _run_async(_execute()) return _run_async(_execute())
@shared_task(
name="tasks.run_pipeline_orchestration",
max_retries=1,
)
def run_pipeline_orchestration(company_id: str, project_brief: str, pipeline_spec: list[dict] | None = None):
"""
Celery task to run a complete topological project pipeline using ProjectOrchestrator.
Writes the final report to workspaces/{company_id}/shared/docs/weekly_report_{date}.md
"""
logger.info(f"[Celery] Starting pipeline orchestration task: company_id={company_id}")
async def _execute():
from services.project_orchestrator import ProjectOrchestrator
from pathlib import Path
from datetime import datetime
orchestrator = ProjectOrchestrator(
company_id=company_id,
project_brief=project_brief,
pipeline_spec=pipeline_spec
)
result = await orchestrator.run()
# Write final markdown report to shared/docs/weekly_report_{date}.md
shared_dir = Path("workspaces") / company_id / "shared"
docs_dir = shared_dir / "docs"
docs_dir.mkdir(parents=True, exist_ok=True)
date_str = datetime.utcnow().strftime("%Y%m%d")
report_file = docs_dir / f"weekly_report_{date_str}.md"
# Generate markdown content for the report
md_lines = [
f"# Weekly CEO Executive Report — {datetime.utcnow().strftime('%Y-%m-%d')}",
"",
f"**Pipeline ID:** `{result.get('pipeline_id')}`",
f"**Company ID:** `{result.get('company_id')}`",
f"**Status:** `{result.get('status')}`",
f"**Duration:** {result.get('total_duration_s')} seconds",
f"**Execution Summary:** {result.get('agents_succeeded')} / {result.get('agents_executed')} agents completed successfully.",
"",
"## Executed Steps",
""
]
for step in result.get("steps", []):
md_lines.append(
f"- **{step.get('agent_name')}** ({step.get('role')}): "
f"Status: `{step.get('status')}` (Duration: {step.get('duration_s')}s)"
)
if step.get("error"):
md_lines.append(f" - *Error:* {step.get('error')}")
md_lines.extend([
"",
"## Produced Deliverables",
""
])
if result.get("deliverables"):
for d in result.get("deliverables"):
md_lines.append(f"- `{d}`")
else:
md_lines.append("*No deliverables registered.*")
report_file.write_text("\n".join(md_lines), encoding="utf-8")
logger.info(f"[Celery] Pipeline report written to: {report_file}")
return result
return _run_async(_execute())
@shared_task(name="tasks.check_scheduled_routines")
def check_scheduled_routines():
"""Celery task to run periodic routine scheduler tick."""
from services.routine_scheduler import RoutineSchedulerService
_run_async(RoutineSchedulerService.tick())
# Campaign API tests
[pytest]
asyncio_mode = auto
"""
Human-in-loop approval flow test.
Reproduces and verifies the FE<->BE contract that the pipeline orchestrator
depends on: POST /approvals/{id}/approve|reject must set `resolution`, which
ProjectOrchestrator polls to resume (approved) or fail (rejected) a gated stage.
Uses a real in-memory SQLite DB (no mocks) so the SQL wiring is actually exercised.
"""
import sys
from pathlib import Path
import pytest
import httpx
from unittest.mock import MagicMock
from sqlalchemy import select
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import sessionmaker
from sqlalchemy.pool import StaticPool
sys.path.insert(0, str(Path(__file__).parent.parent))
from server import app
from database import get_db_session, Base
from middleware.auth import get_current_user
from models import Company, CompanyMembership, Approval
@pytest.mark.asyncio
async def test_approve_and_reject_set_resolution():
engine = create_async_engine(
"sqlite+aiosqlite:///:memory:",
poolclass=StaticPool,
connect_args={"check_same_thread": False},
)
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
TestSession = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
# Seed company + membership + two pending pipeline-stage approvals
async with TestSession() as s:
s.add(Company(id="co1", name="Co", status="active", issue_prefix="PAP", budget_monthly_cents=0))
s.add(CompanyMembership(
id="m1", company_id="co1", principal_type="user",
principal_id="u1", membership_role="instance_admin", status="active",
))
s.add(Approval(id="ap1", company_id="co1", resource_type="pipeline_stage",
resource_id="pipe1", required_approvers=1, created_by="system"))
s.add(Approval(id="ap2", company_id="co1", resource_type="pipeline_stage",
resource_id="pipe2", required_approvers=1, created_by="system"))
await s.commit()
mock_user = MagicMock()
mock_user.id = "u1"
async def _override_user():
return mock_user
async def _override_db():
async with TestSession() as session:
yield session
app.dependency_overrides[get_current_user] = _override_user
app.dependency_overrides[get_db_session] = _override_db
try:
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://testserver") as client:
# Company-scoped list (FE path) returns both pending approvals
r = await client.get("/api/companies/co1/approvals")
assert r.status_code == 200, r.text
assert len(r.json()) == 2
# Approve ap1 (no company_id in body — derived from approval)
r = await client.post("/api/approvals/ap1/approve", json={"decisionNote": "ship it"})
assert r.status_code == 200, r.text
assert r.json()["resolution"] == "approved"
# Reject ap2
r = await client.post("/api/approvals/ap2/reject", json={})
assert r.status_code == 200, r.text
assert r.json()["resolution"] == "rejected"
# Detail GET without company_id query (FE calls get(id))
r = await client.get("/api/approvals/ap1")
assert r.status_code == 200, r.text
assert r.json()["resolution"] == "approved"
# Double-resolve is rejected
r = await client.post("/api/approvals/ap1/approve", json={})
assert r.status_code == 400
finally:
app.dependency_overrides.clear()
# The persisted resolution is exactly what the orchestrator polls
async with TestSession() as s:
a1 = (await s.execute(select(Approval).where(Approval.id == "ap1"))).scalar_one()
a2 = (await s.execute(select(Approval).where(Approval.id == "ap2"))).scalar_one()
assert a1.resolution == "approved"
assert a1.resolved_by == "u1"
assert a1.resolution_note == "ship it"
assert a2.resolution == "rejected"
await engine.dispose()
@pytest.mark.asyncio
async def test_approve_forbidden_without_membership():
engine = create_async_engine(
"sqlite+aiosqlite:///:memory:",
poolclass=StaticPool,
connect_args={"check_same_thread": False},
)
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
TestSession = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
async with TestSession() as s:
s.add(Company(id="co1", name="Co", status="active", issue_prefix="PAP", budget_monthly_cents=0))
# NOTE: no membership for u_stranger
s.add(Approval(id="ap1", company_id="co1", resource_type="pipeline_stage",
resource_id="pipe1", required_approvers=1, created_by="system"))
await s.commit()
mock_user = MagicMock()
mock_user.id = "u_stranger"
async def _override_user():
return mock_user
async def _override_db():
async with TestSession() as session:
yield session
app.dependency_overrides[get_current_user] = _override_user
app.dependency_overrides[get_db_session] = _override_db
try:
transport = httpx.ASGITransport(app=app)
async with httpx.AsyncClient(transport=transport, base_url="http://testserver") as client:
r = await client.post("/api/approvals/ap1/approve", json={})
assert r.status_code == 403, r.text
finally:
app.dependency_overrides.clear()
async with TestSession() as s:
a1 = (await s.execute(select(Approval).where(Approval.id == "ap1"))).scalar_one()
assert a1.resolution is None # not resolved by a non-member
await engine.dispose()
This diff is collapsed.
"""
NocoBase LIVE integration test (no mocks).
This is the proof that NocoBase is "alive", not "wired but dead":
- It talks to a REAL NocoBase over HTTP.
- It is skipped (not failed) when NocoBase is unreachable, so CI stays green
without the server — but the moment NocoBase is up it actually runs and
creates a real collection.
Bring NocoBase up via npm (NOT docker), then run this test:
npx create-nocobase-app nocobase-app -d sqlite
cd nocobase-app && npm run nocobase install && npm run start
# then, with NOCOBASE_URL/EMAIL/PASSWORD set:
pytest tests/test_nocobase_live.py -v
Env (with sensible defaults):
NOCOBASE_URL=http://localhost:13001
NOCOBASE_EMAIL=admin@nocobase.com
NOCOBASE_PASSWORD=admin123
"""
import os
import sys
from pathlib import Path
import pytest
sys.path.insert(0, str(Path(__file__).parent.parent))
from services.nocobase_connector import NocoBaseConnector
def _connector() -> NocoBaseConnector:
return NocoBaseConnector(
base_url=os.getenv("NOCOBASE_URL", "http://localhost:13001"),
email=os.getenv("NOCOBASE_EMAIL", "admin@nocobase.com"),
password=os.getenv("NOCOBASE_PASSWORD", "admin123"),
)
@pytest.mark.asyncio
async def test_nocobase_create_and_list_collection_live():
connector = _connector()
if not await connector.ping():
pytest.skip(
f"NocoBase not reachable at {connector.base_url} — bring it up via npm "
"(see module docstring) to run this live test."
)
# 1. Authenticate against the real server
assert await connector.login(), "NocoBase login failed — check NOCOBASE_EMAIL/PASSWORD"
# 2. Create a collection (idempotent)
created = await connector.ensure_collection(
"test_products",
[
{"name": "title", "type": "string", "interface": "input"},
{"name": "price", "type": "double", "interface": "number"},
],
)
assert created, "ensure_collection returned False against live NocoBase"
# 3. The collection now shows up in the collections registry
collections = await connector.list_records("collections")
names = {c.get("name") for c in collections}
assert "test_products" in names, f"test_products not found in live collections: {sorted(names)}"
"""
Issue-Driven Swarm tests.
Verifies the core spawn contract (success criterion S2):
- A parent issue + N child specs -> N child issues with parent_id set and
request_depth = parent + 1, origin_kind="agent_spawn".
- A parent already at MAX_SWARM_DEPTH spawns nothing (no runaway recursion).
- Dispatch queues a HeartbeatRun per child that has an assignee.
Uses a real in-memory SQLite DB (no mocks).
"""
import sys
from pathlib import Path
import pytest
from sqlalchemy import select
from sqlalchemy.ext.asyncio import create_async_engine, AsyncSession
from sqlalchemy.orm import sessionmaker
from sqlalchemy.pool import StaticPool
sys.path.insert(0, str(Path(__file__).parent.parent))
from database import Base
from models import Company, Issue, HeartbeatRun
from services.swarm_service import (
spawn_child_issues,
dispatch_child_issues,
MAX_SWARM_DEPTH,
)
async def _make_session():
engine = create_async_engine(
"sqlite+aiosqlite:///:memory:",
poolclass=StaticPool,
connect_args={"check_same_thread": False},
)
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
TestSession = sessionmaker(engine, class_=AsyncSession, expire_on_commit=False)
return engine, TestSession
@pytest.mark.asyncio
async def test_spawn_creates_children_with_parent_link_and_depth():
engine, TestSession = await _make_session()
async with TestSession() as s:
s.add(Company(id="co1", name="Co", status="active", issue_prefix="PAP", budget_monthly_cents=0))
s.add(Issue(id="root", company_id="co1", title="Build feature",
status="in_progress", priority="high", request_depth=0))
await s.commit()
async with TestSession() as s:
parent = (await s.execute(select(Issue).where(Issue.id == "root"))).scalar_one()
specs = [
{"title": "Write backend"},
{"title": "Write frontend", "priority": "medium", "description": "the UI"},
]
children = await spawn_child_issues(s, parent, specs, created_by_agent_id="agentX", origin_run_id="run1")
await s.commit()
assert len(children) == 2
for c in children:
assert c.parent_id == "root"
assert c.request_depth == 1
assert c.origin_kind == "agent_spawn"
assert c.origin_id == "root"
assert c.origin_run_id == "run1"
assert c.created_by_agent_id == "agentX"
assert c.company_id == "co1"
# priority inherited when not specified; explicit kept
by_title = {c.title: c for c in children}
assert by_title["Write backend"].priority == "high" # inherited from parent
assert by_title["Write frontend"].priority == "medium"
assert by_title["Write frontend"].description == "the UI"
# persisted
async with TestSession() as s:
rows = (await s.execute(select(Issue).where(Issue.parent_id == "root"))).scalars().all()
assert len(rows) == 2
await engine.dispose()
@pytest.mark.asyncio
async def test_spawn_blocked_at_max_depth():
engine, TestSession = await _make_session()
async with TestSession() as s:
s.add(Company(id="co1", name="Co", status="active", issue_prefix="PAP", budget_monthly_cents=0))
s.add(Issue(id="deep", company_id="co1", title="Deep issue",
status="in_progress", priority="medium", request_depth=MAX_SWARM_DEPTH))
await s.commit()
async with TestSession() as s:
parent = (await s.execute(select(Issue).where(Issue.id == "deep"))).scalar_one()
children = await spawn_child_issues(s, parent, [{"title": "Should not exist"}])
await s.commit()
assert children == []
async with TestSession() as s:
rows = (await s.execute(select(Issue).where(Issue.parent_id == "deep"))).scalars().all()
assert len(rows) == 0 # depth guard held — no runaway recursion
await engine.dispose()
@pytest.mark.asyncio
async def test_specs_without_title_are_skipped():
engine, TestSession = await _make_session()
async with TestSession() as s:
s.add(Company(id="co1", name="Co", status="active", issue_prefix="PAP", budget_monthly_cents=0))
s.add(Issue(id="root", company_id="co1", title="Parent",
status="open", priority="low", request_depth=0))
await s.commit()
async with TestSession() as s:
parent = (await s.execute(select(Issue).where(Issue.id == "root"))).scalar_one()
children = await spawn_child_issues(s, parent, [{"title": "ok"}, {"description": "no title"}, {}])
await s.commit()
assert len(children) == 1
assert children[0].title == "ok"
await engine.dispose()
@pytest.mark.asyncio
async def test_dispatch_queues_runs_only_for_assigned_children():
engine, TestSession = await _make_session()
async with TestSession() as s:
s.add(Company(id="co1", name="Co", status="active", issue_prefix="PAP", budget_monthly_cents=0))
s.add(Issue(id="root", company_id="co1", title="Parent",
status="in_progress", priority="medium", request_depth=0))
await s.commit()
async with TestSession() as s:
parent = (await s.execute(select(Issue).where(Issue.id == "root"))).scalar_one()
# one child with explicit assignee, one without (parent has none to inherit)
children = await spawn_child_issues(s, parent, [
{"title": "assigned", "assignee_agent_id": "agent1"},
{"title": "unassigned"},
])
runs = await dispatch_child_issues(s, children)
await s.commit()
assert len(runs) == 1
assert runs[0].agent_id == "agent1"
assert runs[0].status == "queued"
# the assigned child got its checkout_run_id wired
assigned = next(c for c in children if c.title == "assigned")
assert assigned.checkout_run_id == runs[0].id
async with TestSession() as s:
all_runs = (await s.execute(select(HeartbeatRun))).scalars().all()
assert len(all_runs) == 1
await engine.dispose()
import pytest
import os
import uuid
import datetime
import sys
from unittest.mock import patch, AsyncMock
# Add backend to path
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from sqlalchemy.ext.asyncio import AsyncSession
from models.companies import Company
from models.agents import Agent
from models.company_secrets import CompanySecret
from models.plugins import Plugin
from models.plugin_config import PluginConfig
from models.heartbeat_runs import HeartbeatRun
from api.routes.agents import run_agent_in_background
from services.runtime_env import resolve_runtime_env, resolve_enabled_toolsets
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker
import pytest_asyncio
from database import Base
engine = create_async_engine("sqlite+aiosqlite:///:memory:")
TestingSessionLocal = async_sessionmaker(autocommit=False, autoflush=False, bind=engine)
@pytest_asyncio.fixture
async def db():
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
async with TestingSessionLocal() as session:
yield session
async with engine.begin() as conn:
await conn.run_sync(Base.metadata.drop_all)
@pytest.mark.asyncio
async def test_team_env_injection_and_cleanup(db: AsyncSession):
from services.runtime_env import resolve_runtime_env
# Setup Company
company_id = uuid.uuid4().hex
company = Company(
id=company_id,
name="Test Env Company",
description="test",
status="active"
)
db.add(company)
# Setup Agent with team_id
team_id = "team_alpha"
agent_id = uuid.uuid4().hex
agent = Agent(
id=agent_id,
company_id=company_id,
team_id=team_id,
name="Alpha Agent",
role="tester",
title="Tester",
status="active"
)
db.add(agent)
# Setup Company-level Secret
company_secret_key = f"COMP_TOKEN_{uuid.uuid4().hex[:8]}"
comp_secret = CompanySecret(
id=uuid.uuid4().hex,
company_id=company_id,
team_id=None,
key=company_secret_key,
value_encrypted="comp_secret_value",
is_sensitive=True,
)
db.add(comp_secret)
# Setup Team-level Secret
team_secret_key = f"TEAM_TOKEN_{uuid.uuid4().hex[:8]}"
team_secret = CompanySecret(
id=uuid.uuid4().hex,
company_id=company_id,
team_id=team_id,
key=team_secret_key,
value_encrypted="team_secret_value",
is_sensitive=True,
)
db.add(team_secret)
# Setup Override Secret (Company level has one value, Team level overrides it)
override_key = f"OVERRIDE_TOKEN_{uuid.uuid4().hex[:8]}"
override_comp = CompanySecret(
id=uuid.uuid4().hex,
company_id=company_id,
team_id=None,
key=override_key,
value_encrypted="base_value",
is_sensitive=True,
)
db.add(override_comp)
override_team = CompanySecret(
id=uuid.uuid4().hex,
company_id=company_id,
team_id=team_id,
key=override_key,
value_encrypted="overridden_value",
is_sensitive=True,
)
db.add(override_team)
await db.commit()
# Test resolving env for the specific team
env_vars = await resolve_runtime_env(db, company_id, team_id)
# Assertions
assert company_secret_key in env_vars
assert env_vars[company_secret_key] == "comp_secret_value"
assert team_secret_key in env_vars
assert env_vars[team_secret_key] == "team_secret_value"
assert override_key in env_vars
assert env_vars[override_key] == "overridden_value", "Team secret should override company secret"
# Test resolving env for no team (company level only)
env_vars_comp_only = await resolve_runtime_env(db, company_id, None)
assert company_secret_key in env_vars_comp_only
assert team_secret_key not in env_vars_comp_only, "Team secret should not leak to company-level context"
assert env_vars_comp_only[override_key] == "base_value", "Company level should see the base value"
@pytest.mark.asyncio
async def test_resolve_enabled_toolsets(db: AsyncSession):
company_id = uuid.uuid4().hex
team_id = "team_beta"
# Create plugin
plugin = Plugin(
id=uuid.uuid4().hex,
name="test_plugin",
version="1.0",
enabled=True
)
db.add(plugin)
# Create plugin config for team
config = PluginConfig(
id=uuid.uuid4().hex,
company_id=company_id,
team_id=team_id,
plugin_id=plugin.id,
is_enabled=True,
config={}
)
db.add(config)
await db.commit()
# Resolve
toolsets = await resolve_enabled_toolsets(db, company_id, team_id)
assert toolsets == ["test_plugin"]
# Resolve for another team should be None or empty
toolsets_other = await resolve_enabled_toolsets(db, company_id, "other_team")
assert toolsets_other == []
# Canifa · Brand Spec # Paperclip · Brand Spec
> 采集日期:2026-04-30 > 收集日期:2026-05-30
> 资产来源:canifa.com 官网提取 > 资产来源:paperclip.ai 提取
> 资产完整度:完整
## 🎯 核心资产(一等公民)
## 🎯 核心资产(一等公民)
### Logo
### Logo - 主版本:`backend/static/paperclip-brand/logo.svg`
- 主版本:`backend/static/canifa-brand/logo.svg` - 使用场景:Chatbot Header, Watermark, App Launcher
- 使用场景:Chatbot Header, Watermark, App Launcher - 视觉特征:现代字标 "Paperclip" 无衬线设计,结合圆角几何形状。
- 视觉特征:经典的 "CANIFA" 无衬线字标,置于红色圆角矩形块中。
### 品牌气质
### 品牌气质 - 核心关键词:Efficiency, Automation, Reliability, Modular.
- 核心关键词:Fashion for All, Gia đình, Năng động, Tin cậy.
- 2026 夏季主题:Trạm Hè Đa Sắc (Vibrant Summer Station) ## 🎨 辅助资产
## 🎨 辅助资产 ### 色板
- **Paperclip Blue (Primary)**: `#0F62FE` (Logo & CTA)
### 色板 - **Deep Navy/Black**: `#161616` (Typography & Secondary elements)
- **Canifa Red (Primary)**: `#E2231A` (Logo & CTA) - **Tech Teal/Accent**: `#00F0FF` (Highlighter / Accent elements)
- **Deep Navy/Black**: `#333F48` (Typography & Secondary elements) - **White**: `#FFFFFF` (Background)
- **Summer Yellow/Lime**: `#C4FF1C` (Highlighter / New Campaign) - **Light Gray**: `#F4F4F4` (Container background)
- **White**: `#FFFFFF` (Background)
- **Light Gray**: `#F5F5F5` (Container background) ### 字型
- **Display**: Sans-serif (Montserrat, Outfit, or Inter)
### 字型 - **Body**: Sans-serif (Highly readable)
- **Display**: Sans-serif (Clean, modern like Montserrat or Inter)
- **Body**: Sans-serif (Highly readable) ### 交互签名
- 极简、流畅、圆角适中(8px - 12px)。
### 交互签名 - 强调效率与高响应性。
- 极简、流畅、圆角适中(8px - 12px)。
- 强调 "Joy" (Niềm vui) thông qua các hiệu ứng chuyển cảnh nhẹ nhàng.
# Paperclip OS — Industry Blueprints Guide
Paperclip OS is a generic, blueprint-driven AI Company Operating System. Instead of being tied to a single industry or organization structure, Paperclip can bootstrap entirely different organizations (e.g., Software Studio, Fashion Retail, E-commerce, Marketing Agency) from a single YAML blueprint definition.
This guide explains how to write, customize, and deploy blueprints in Paperclip.
---
## 1. Blueprint Structure
A blueprint is a YAML file stored in the `backend/blueprints/` folder. It has three main sections:
1. **Metadata**: Define the template name, industry, and description.
2. **Roles (Organizational Chart)**: Specify the AI agents, their credentials, adapters, capabilities, and hierarchy (who reports to whom).
3. **Pipeline (Task Execution DAG)**: Graph definition detailing the sequence of task executions and approval gates.
Here is a complete example of a blueprint structure:
```yaml
name: "Software Studio"
industry: "software"
description: "A blueprint for custom software development projects, featuring PMs, system architects, coders, and QA engineers."
roles:
- role: "ceo"
name: "CEO Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "strategic direction, weekly reporting, approval, corporate oversight"
reports_to: null
- role: "pm"
name: "Product Manager Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "product specifications writing, feature scoping, task prioritization"
reports_to: "ceo"
- role: "designer"
name: "Software Architect Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "system architecture, database schema design, class diagrams, API interfaces"
reports_to: "pm"
- role: "coder"
name: "Lead Developer Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "full-stack coding, writing algorithms, API implementation, code refactoring"
reports_to: "pm"
- role: "qa"
name: "QA Specialist Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "writing unit tests, bug hunting, API response verification"
reports_to: "pm"
- role: "devops"
name: "DevOps Engineer Agent"
adapter_type: "openai"
adapter_config:
model_name: "gpt-4o-mini"
runtime_config:
max_turns_per_run: 10
capabilities: "CI/CD pipeline configuration, Docker container build, deployment scripts, monitoring"
reports_to: "ceo"
pipeline:
- role: "ceo"
depends_on: []
approval_gate: false
- role: "pm"
depends_on: ["ceo"]
approval_gate: true
- role: "designer"
depends_on: ["pm"]
approval_gate: false
- role: "coder"
depends_on: ["designer"]
approval_gate: false
- role: "qa"
depends_on: ["coder"]
approval_gate: true
- role: "devops"
depends_on: ["qa"]
approval_gate: false
```
---
## 2. Roles Configuration
The `roles` array specifies the AI agents created for the company:
- `role`: The unique lowercase key for the role (e.g. `ceo`, `designer`, `coder`).
- `name`: The display title/name of the agent.
- `adapter_type`: The runtime executor adapter (e.g. `openai`, `claude_local`, `gemini_local`, `cursor`, `opencode_local`).
- `adapter_config`: Arguments passed to the adapter. E.g. `model_name` or custom environment variables.
- `runtime_config`: Platform-level configurations, such as `max_turns_per_run`.
- `capabilities`: Plain-text description of the agent's skillset, injected into its system prompts.
- `reports_to`: Optional parent role identifier establishing the reporting hierarchy (used for task routing and escalation).
---
## 3. Pipeline / DAG Configuration
The `pipeline` list defines the dependency graph (DAG) for project execution runs:
- `role`: The agent responsible for this pipeline stage.
- `depends_on`: A list of parent roles that must complete execution before this stage can start. Parallel branches run concurrently automatically.
- `approval_gate`: A boolean flag. If `true`, the pipeline pauses after this agent finishes its work, demanding user/board review and approval before triggering dependent stages.
---
## 4. Cycle Validation
Paperclip OS validates all custom blueprints upon loading. The service uses Kahn's algorithm to perform cycle checks and validates that all role references in both `reports_to` and `depends_on` are correctly defined in the blueprint's roles array. If a cycle or an invalid reference is detected, the API will refuse to load the blueprint and raise a `400 Bad Request`.
# Demo E2E Run # Demo E2E Run
**Date:** 2026-06-02T14:28:31.221087 **Date:** 2026-06-02T16:31:09.656355
**Status:** completed **Status:** completed
**Duration:** 77.0s **Duration:** 77.0s
**Blueprint:** software-studio **Blueprint:** software-studio
...@@ -33,8 +33,8 @@ ...@@ -33,8 +33,8 @@
```json ```json
{ {
"pipeline_id": "a385b7aaee904b00ba8ceec01a3878c1", "pipeline_id": "34fc39fd679443cfbf2b5b2062855374",
"company_id": "c7cfd603b9504bfd934a66faf9b48363", "company_id": "d187083b62ab4ff0ab7674383e7b3e49",
"status": "completed", "status": "completed",
"total_duration_s": 77.0, "total_duration_s": 77.0,
"agents_executed": 6, "agents_executed": 6,
...@@ -57,66 +57,66 @@ ...@@ -57,66 +57,66 @@
"steps": [ "steps": [
{ {
"step": 1, "step": 1,
"agent_id": "ceo_41e8c849", "agent_id": "ceo_cf76fda2",
"agent_name": "CEO Agent", "agent_name": "CEO Agent",
"role": "ceo", "role": "ceo",
"run_id": "14f136fa17af447d8531a2ed1dc092a0", "run_id": "7d68ebab076841d592e6d92ace061a09",
"status": "completed", "status": "completed",
"duration_s": 9.8, "duration_s": 10.0,
"error": null "error": null
}, },
{ {
"step": 2, "step": 2,
"agent_id": "pm_04078779", "agent_id": "pm_fae7b171",
"agent_name": "Product Manager Agent", "agent_name": "Product Manager Agent",
"role": "pm", "role": "pm",
"run_id": "1fb0960c6cd74c0bb1a4a94f5a17333d", "run_id": "95f1d638f3994eb3992648345b1deddd",
"status": "completed", "status": "completed",
"duration_s": 13.2, "duration_s": 13.3,
"error": null "error": null
}, },
{ {
"step": 3, "step": 3,
"agent_id": "designer_8c957a41", "agent_id": "designer_5f34c481",
"agent_name": "Software Architect Agent", "agent_name": "Software Architect Agent",
"role": "designer", "role": "designer",
"run_id": "2b181e6d8cec49b0854bd56b4a457654", "run_id": "6be1be0b76d04463b9b9e8448c7a1119",
"status": "completed", "status": "completed",
"duration_s": 9.7, "duration_s": 9.6,
"error": null "error": null
}, },
{ {
"step": 4, "step": 4,
"agent_id": "coder_ae0f63cf", "agent_id": "coder_5eebeb72",
"agent_name": "Lead Developer Agent", "agent_name": "Lead Developer Agent",
"role": "coder", "role": "coder",
"run_id": "002d6d62123b48e0b940b6fde8f4e8ad", "run_id": "0dc497854f284291b4529fc648e7884a",
"status": "completed", "status": "completed",
"duration_s": 11.8, "duration_s": 11.7,
"error": null "error": null
}, },
{ {
"step": 5, "step": 5,
"agent_id": "qa_7d279926", "agent_id": "qa_43ef2a4f",
"agent_name": "QA Specialist Agent", "agent_name": "QA Specialist Agent",
"role": "qa", "role": "qa",
"run_id": "c81cd6e371c9477cb96e0f3c3e2821c1", "run_id": "23c544f066ae45048ee65ef14579fda0",
"status": "completed", "status": "completed",
"duration_s": 10.2, "duration_s": 10.2,
"error": null "error": null
}, },
{ {
"step": 6, "step": 6,
"agent_id": "devops_ff60aad0", "agent_id": "devops_ce827a83",
"agent_name": "DevOps Engineer Agent", "agent_name": "DevOps Engineer Agent",
"role": "devops", "role": "devops",
"run_id": "989711aa25734782a133895ac59c400b", "run_id": "452891ae28bc4fc0bf6a7234fc081021",
"status": "completed", "status": "completed",
"duration_s": 9.7, "duration_s": 9.6,
"error": null "error": null
} }
], ],
"started_at": "2026-06-02T14:28:31.221087", "started_at": "2026-06-02T16:31:09.656355",
"completed_at": "2026-06-02T14:29:48.183474" "completed_at": "2026-06-02T16:32:26.635376"
} }
``` ```
\ No newline at end of file
...@@ -19,3 +19,5 @@ export { sidebarPreferencesApi } from "./sidebarPreferences"; ...@@ -19,3 +19,5 @@ export { sidebarPreferencesApi } from "./sidebarPreferences";
export { inboxDismissalsApi } from "./inboxDismissals"; export { inboxDismissalsApi } from "./inboxDismissals";
export { companySkillsApi } from "./companySkills"; export { companySkillsApi } from "./companySkills";
export { blueprintsApi } from "./blueprints"; export { blueprintsApi } from "./blueprints";
export { pipelineRunsApi } from "./pipelineRuns";
import { api } from "./client";
export interface PipelineRunStep {
id: string;
agentName: string;
role: string;
status: string;
step: number;
startedAt: string;
completedAt?: string;
error?: string;
executionCostCents: number;
}
export interface PipelineRun {
pipelineId: string;
status: string;
startedAt: string;
completedAt?: string;
durationS: number;
steps: PipelineRunStep[];
agentsExecuted: number;
agentsSucceeded: number;
agentsFailed: number;
}
export const pipelineRunsApi = {
list: (companyId: string) =>
api.get<PipelineRun[]>(`/companies/${companyId}/pipeline-runs`),
get: (companyId: string, pipelineId: string) =>
api.get<PipelineRun>(`/companies/${companyId}/pipeline-runs/${pipelineId}`),
};
This diff is collapsed.
...@@ -76,6 +76,7 @@ import { IssueScheduledRetryCard } from "../components/IssueScheduledRetryCard"; ...@@ -76,6 +76,7 @@ import { IssueScheduledRetryCard } from "../components/IssueScheduledRetryCard";
import { IssueProperties } from "../components/IssueProperties"; import { IssueProperties } from "../components/IssueProperties";
import { IssueRunLedger } from "../components/IssueRunLedger"; import { IssueRunLedger } from "../components/IssueRunLedger";
import { IssueWorkspaceCard } from "../components/IssueWorkspaceCard"; import { IssueWorkspaceCard } from "../components/IssueWorkspaceCard";
import { SwarmTree } from "../components/SwarmTree";
import type { MentionOption } from "../components/MarkdownEditor"; import type { MentionOption } from "../components/MarkdownEditor";
import { ImageGalleryModal } from "../components/ImageGalleryModal"; import { ImageGalleryModal } from "../components/ImageGalleryModal";
import { ScrollToBottom } from "../components/ScrollToBottom"; import { ScrollToBottom } from "../components/ScrollToBottom";
...@@ -123,6 +124,7 @@ import { ...@@ -123,6 +124,7 @@ import {
Eye, Eye,
EyeOff, EyeOff,
Flag, Flag,
GitBranch,
Hexagon, Hexagon,
ListTree, ListTree,
MessageSquare, MessageSquare,
...@@ -3895,6 +3897,10 @@ export function IssueDetail() { ...@@ -3895,6 +3897,10 @@ export function IssueDetail() {
<ListTree className="h-3.5 w-3.5" /> <ListTree className="h-3.5 w-3.5" />
Related work Related work
</TabsTrigger> </TabsTrigger>
<TabsTrigger value="swarm-tree" className="gap-1.5">
<GitBranch className="h-3.5 w-3.5" />
Swarm Tree
</TabsTrigger>
{issuePluginTabItems.map((item) => ( {issuePluginTabItems.map((item) => (
<TabsTrigger key={item.value} value={item.value}> <TabsTrigger key={item.value} value={item.value}>
{item.label} {item.label}
...@@ -4006,6 +4012,16 @@ export function IssueDetail() { ...@@ -4006,6 +4012,16 @@ export function IssueDetail() {
<IssueRelatedWorkPanel relatedWork={issue.relatedWork} /> <IssueRelatedWorkPanel relatedWork={issue.relatedWork} />
</TabsContent> </TabsContent>
<TabsContent value="swarm-tree">
{detailTab === "swarm-tree" ? (
<SwarmTree
rootIssue={issue}
allIssues={childIssues}
agents={agents ?? []}
/>
) : null}
</TabsContent>
{activePluginTab && ( {activePluginTab && (
<TabsContent value={activePluginTab.value}> <TabsContent value={activePluginTab.value}>
<PluginSlotMount <PluginSlotMount
......
This diff is collapsed.
# 🔨 Doing #12 — Team Schema + Migration (Chunk 1, prereq)
**Goal:** Đẻ entity `Team` + cột `team_id` nullable trên 5 bảng tài nguyên, đổi unique constraint của `company_secrets`. Một migration round-trip được. KHÔNG đụng runtime/logic — chỉ schema.
**Dependency:** Không (đây là prereq của tất cả các chunk sau).
**Idea nguồn:** `plan/ideas/33_multi_team_company.md`
---
## Research-first (đọc trước khi sửa)
- [x] `backend/models/agents.py` — xem cách 1 model khai báo (Base, `company_id` FK, `__tablename__`, mixin timestamp nếu có)
- [x] `backend/models/company_secrets.py` — tìm `UniqueConstraint(company_id, key)` hiện tại (chỗ phải đổi)
- [x] `backend/models/memory.py`, `company_skills.py`, `plugin_config.py` — xác nhận đều scope `company_id`, chưa có `team_id`
- [x] `backend/alembic/versions/` (hoặc thư mục migration) — xem migration gần nhất + cách dùng `batch_alter_table` (SQLite KHÔNG ALTER constraint in-place)
- [x] Xác nhận file DB: `backend/paperclip.db` (hay đường dẫn khác trong `alembic.ini`/`env.py`)
## Phases
### Phase 1 — Model
- [x] Tạo `backend/models/teams.py`: `Team` (id `String(36)` PK, `company_id` FK→companies.id NOT NULL, `name` String(255) NOT NULL, `slug` String(100), `parent_team_id` FK→teams.id **nullable**, `created_at`/`updated_at`)
- [x] Đăng ký model vào nơi import chung (giống các model khác — kiểm `models/__init__.py`)
- [x] Thêm `team_id String(36) FK→teams.id nullable` vào: `agents`, `memory`, `company_skills`, `company_secrets`, `plugin_config`
- [x] Sửa `company_secrets`: `UniqueConstraint(company_id, key)` → `UniqueConstraint(company_id, team_id, key)`
### Phase 2 — Migration (batch mode)
- [x] 1 file Alembic: tạo bảng `teams` **TRƯỚC** → `batch_alter_table` add cột `team_id` cho từng bảng → trong batch của `company_secrets` swap unique constraint (drop cũ, tạo mới có team_id)
- [x] Viết `downgrade()` đối xứng (drop cột, drop bảng teams, khôi phục constraint cũ)
## Verify
- [x] Copy DB ra bản tạm rồi chạy: `alembic upgrade head` → `alembic downgrade -1` → `alembic upgrade head` round-trip KHÔNG lỗi
- [x] `rtk python -c "import sqlite3; c=sqlite3.connect('paperclip.db'); print(c.execute('SELECT count(*) FROM agents WHERE team_id IS NOT NULL').fetchone())"` → `(0,)` (mọi row cũ null = backward-compat)
- [x] `rtk python -c "import sqlite3; c=sqlite3.connect('paperclip.db'); print([r[1] for r in c.execute('PRAGMA table_info(company_secrets)')])"` → có cột `team_id`
- [x] Backend import sạch: `rtk python -c "from models.teams import Team; print('ok')"` (chạy từ `backend/` với venv)
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
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