Commit b5d9b6ef authored by Admin's avatar Admin

feat: implement backend endpoints for external integrations and credentials storage

parent 1145df18
......@@ -31,6 +31,7 @@ from .routes.issue_comments import router as issue_comments_router
from .routes.issue_labels import router as issue_labels_router
from .routes.user_profiles import router as user_profiles_router
from .routes.adapters import router as adapters_router
from .routes.campaigns import router as campaigns_router
from .routes.issue_tree_control import router as issue_tree_control_router
from .routes.llms import router as llms_router
from .routes.org_chart_svg import router as org_chart_svg_router
......@@ -66,6 +67,7 @@ from .routes.marketplace import router as marketplace_router
from .routes.cubesandbox import router as cubesandbox_router
from .routes.pipelines import router as pipelines_router
from .routes.agent_hires import router as agent_hires_router
from .routes.company_plugins import router as company_plugins_router
main_router = APIRouter()
......@@ -95,6 +97,7 @@ api_sub_router.include_router(approvals_router)
api_sub_router.include_router(secrets_router)
api_sub_router.include_router(costs_router)
api_sub_router.include_router(budgets_router)
api_sub_router.include_router(campaigns_router)
api_sub_router.include_router(activity_router)
api_sub_router.include_router(dashboard_router)
api_sub_router.include_router(heartbeat_runs_router)
......@@ -147,6 +150,7 @@ api_sub_router.include_router(marketplace_router)
api_sub_router.include_router(cubesandbox_router)
api_sub_router.include_router(pipelines_router)
api_sub_router.include_router(agent_hires_router)
api_sub_router.include_router(company_plugins_router)
# Register /api sub-router to the main router
main_router.include_router(api_sub_router)
This diff is collapsed.
"""
Company plugins route.
"""
from fastapi import APIRouter, Depends, HTTPException, Query, status
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select
from pydantic import BaseModel
from typing import Dict, Any
import uuid
import datetime
from database import get_db_session
from middleware.auth import get_current_user
from schemas.auth import AuthUserResponse
from models import CompanyMembership, Plugin, PluginConfig, CompanySecret
from common.encryption import encrypt_api_key
router = APIRouter(prefix="/company-plugins", tags=["company-plugins"])
class CompanyPluginConnectRequest(BaseModel):
company_id: str
plugin_slug: str
plugin_name: str
config: Dict[str, Any]
class CompanyPluginDisconnectRequest(BaseModel):
company_id: str
plugin_slug: str
@router.post("/connect")
async def connect_company_plugin(
req: CompanyPluginConnectRequest,
current_user: AuthUserResponse = Depends(get_current_user),
db: AsyncSession = Depends(get_db_session)
):
"""Connect a plugin and encrypt its secrets."""
# Verify company membership
result = await db.execute(
select(CompanyMembership).where(
CompanyMembership.company_id == req.company_id,
CompanyMembership.principal_type == "user",
CompanyMembership.principal_id == current_user.id,
CompanyMembership.status == "active"
)
)
membership = result.scalar_one_or_none()
if not membership:
raise HTTPException(status_code=403, detail="No access to company")
# Find or create Plugin
result = await db.execute(select(Plugin).where(Plugin.name == req.plugin_slug))
plugin = result.scalar_one_or_none()
if not plugin:
plugin = Plugin(
id=f"plugin_{req.plugin_slug}",
name=req.plugin_slug,
version="1.0.0",
description=f"{req.plugin_name} integration",
enabled=True,
config_schema={}
)
db.add(plugin)
await db.flush()
# Process and save secrets
config_mapping = {}
for key, value in req.config.items():
if not value:
continue
# Determine secret key name
if req.plugin_slug == "facebook_pages":
prefix = "FACEBOOK_PAGE"
elif req.plugin_slug == "nocobase_erp":
prefix = "NOCOBASE"
else:
prefix = req.plugin_slug.upper()
secret_key = f"{prefix}_{key.upper()}"
config_mapping[key] = secret_key
# Check if secret already exists
sec_stmt = select(CompanySecret).where(
CompanySecret.company_id == req.company_id,
CompanySecret.key == secret_key,
CompanySecret.deleted_at.is_(None)
)
sec_res = await db.execute(sec_stmt)
existing_secret = sec_res.scalar_one_or_none()
# Encrypt value
encrypted_val = encrypt_api_key(str(value))
is_sensitive = any(kw in key.lower() for kw in ["key", "token", "secret", "password", "private", "credential"])
if existing_secret:
existing_secret.value_encrypted = encrypted_val
existing_secret.is_sensitive = is_sensitive
existing_secret.updated_at = datetime.datetime.utcnow()
else:
new_sec = CompanySecret(
id=uuid.uuid4().hex,
company_id=req.company_id,
key=secret_key,
value_encrypted=encrypted_val,
description=f"Secret for {req.plugin_name} {key}",
is_sensitive=is_sensitive,
created_by=current_user.id
)
db.add(new_sec)
# Find or create PluginConfig
cfg_stmt = select(PluginConfig).where(
PluginConfig.company_id == req.company_id,
PluginConfig.plugin_id == plugin.id
)
cfg_res = await db.execute(cfg_stmt)
plugin_config = cfg_res.scalar_one_or_none()
if plugin_config:
plugin_config.is_enabled = True
plugin_config.config = config_mapping
plugin_config.updated_at = datetime.datetime.utcnow()
else:
plugin_config = PluginConfig(
id=uuid.uuid4().hex,
company_id=req.company_id,
plugin_id=plugin.id,
config=config_mapping,
is_enabled=True
)
db.add(plugin_config)
await db.commit()
return {"status": "success", "message": f"Connected {req.plugin_name} successfully"}
@router.post("/disconnect")
async def disconnect_company_plugin(
req: CompanyPluginDisconnectRequest,
current_user: AuthUserResponse = Depends(get_current_user),
db: AsyncSession = Depends(get_db_session)
):
"""Disconnect a plugin and soft-delete its secrets."""
# Verify company membership
result = await db.execute(
select(CompanyMembership).where(
CompanyMembership.company_id == req.company_id,
CompanyMembership.principal_type == "user",
CompanyMembership.principal_id == current_user.id,
CompanyMembership.status == "active"
)
)
membership = result.scalar_one_or_none()
if not membership:
raise HTTPException(status_code=403, detail="No access to company")
# Find the Plugin
result = await db.execute(select(Plugin).where(Plugin.name == req.plugin_slug))
plugin = result.scalar_one_or_none()
if not plugin:
raise HTTPException(status_code=404, detail="Plugin not found")
# Find the PluginConfig
cfg_stmt = select(PluginConfig).where(
PluginConfig.company_id == req.company_id,
PluginConfig.plugin_id == plugin.id
)
cfg_res = await db.execute(cfg_stmt)
plugin_config = cfg_res.scalar_one_or_none()
if not plugin_config:
raise HTTPException(status_code=404, detail="Plugin config not found")
# Disable plugin
plugin_config.is_enabled = False
plugin_config.updated_at = datetime.datetime.utcnow()
# Soft delete related secrets stored in config
if plugin_config.config:
for val_key, secret_key in plugin_config.config.items():
sec_stmt = select(CompanySecret).where(
CompanySecret.company_id == req.company_id,
CompanySecret.key == secret_key,
CompanySecret.deleted_at.is_(None)
)
sec_res = await db.execute(sec_stmt)
secret = sec_res.scalar_one_or_none()
if secret:
secret.deleted_at = datetime.datetime.utcnow()
await db.commit()
return {"status": "success", "message": f"Disconnected {req.plugin_slug} successfully"}
......@@ -8,11 +8,11 @@ from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy import select, func
from database import get_db_session
from middleware.auth import get_current_user, require_instance_admin
from middleware.auth import get_current_user, require_instance_admin, require_company_scope
from schemas.plugin import PluginCreate, PluginUpdate, PluginResponse
from schemas.common import PaginationParams, PaginatedResponse
from schemas.auth import AuthUserResponse
from models import Plugin
from models import Plugin, PluginConfig, CompanyMembership
from services.activity_logger import log_activity
router = APIRouter(prefix="/plugins", tags=["plugins"])
......@@ -158,3 +158,90 @@ async def delete_plugin(
await db.commit()
return None
@router.get("/registry")
async def get_plugins_registry(
company_id: str = Query(..., description="Company ID"),
current_user: AuthUserResponse = Depends(get_current_user),
db: AsyncSession = Depends(get_db_session)
):
"""Get plugins registry status for a company."""
# Verify company membership
result = await db.execute(
select(CompanyMembership).where(
CompanyMembership.company_id == company_id,
CompanyMembership.principal_type == "user",
CompanyMembership.principal_id == current_user.id,
CompanyMembership.status == "active"
)
)
membership = result.scalar_one_or_none()
if not membership:
raise HTTPException(status_code=403, detail="No access to company")
# Define standard plugins
STANDARD_PLUGINS = {
"trello": "Trello integration for project management",
"notion": "Notion integration for document syncing",
"linear": "Linear integration for engineering task tracking",
"clickup": "ClickUp integration for project management",
"github_projects": "GitHub Projects integration for task management",
"facebook_pages": "Facebook Page posting integration",
"nocobase_erp": "NocoBase ERP connection for structured data",
"slack": "Slack realtime notifications integration",
"langfuse": "Langfuse tracing and observability integration",
"google_drive": "Google Drive cloud storage integration",
"gemini_content": "Gemini API content generation",
"image_gen": "AI image generation for product photos"
}
# Retrieve existing plugins from DB
existing_plugins_stmt = select(Plugin)
existing_plugins_res = await db.execute(existing_plugins_stmt)
existing_plugins = {p.name: p for p in existing_plugins_res.scalars().all()}
# Auto-seed standard plugins if missing
db_changed = False
for slug, desc in STANDARD_PLUGINS.items():
if slug not in existing_plugins:
new_p = Plugin(
id=f"plugin_{slug}",
name=slug,
version="1.0.0",
description=desc,
enabled=True,
config_schema={}
)
db.add(new_p)
existing_plugins[slug] = new_p
db_changed = True
if db_changed:
await db.flush()
# Get active configs for this company
configs_stmt = select(PluginConfig).where(
PluginConfig.company_id == company_id,
PluginConfig.is_enabled == True
)
configs_res = await db.execute(configs_stmt)
configs = configs_res.scalars().all()
enabled_plugin_ids = {c.plugin_id for c in configs}
# Format response: each plugin needs: id, name, enabled, description
registry_list = []
for slug, p in existing_plugins.items():
if p.enabled:
registry_list.append({
"id": p.id,
"name": p.name,
"enabled": p.id in enabled_plugin_ids,
"description": p.description
})
if db_changed:
await db.commit()
return registry_list
"""
Demo E2E: deploy software-studio blueprint → run ProjectOrchestrator → print CEO report.
Cách chạy:
cd D:\\a\\ai_canifa_company\\backend
python scratch/demo_e2e.py
Cần env:
OPENAI_API_KEY=sk-... (bắt buộc — blueprint dùng gpt-4o-mini)
DATABASE_URL=sqlite+aiosqlite:///./paperclip.db (default)
"""
import asyncio
import os
import sys
import json
from pathlib import Path
from datetime import datetime
# Fix Windows console UTF-8 printing
if hasattr(sys.stdout, 'reconfigure'):
sys.stdout.reconfigure(encoding='utf-8')
if hasattr(sys.stderr, 'reconfigure'):
sys.stderr.reconfigure(encoding='utf-8')
# --- Path setup (phải TRƯỚC các import backend) ---
BACKEND_DIR = Path(__file__).parent.parent
sys.path.insert(0, str(BACKEND_DIR))
os.chdir(BACKEND_DIR) # database.py tạo path tương đối từ cwd
from config import settings
# --- Gate: kiểm tra LLM creds ---
OPENAI_KEY = os.getenv("OPENAI_API_KEY") or settings.openai_api_key or ""
if not OPENAI_KEY or OPENAI_KEY.startswith("sk-PLACEHOLDER"):
print("ERROR: OPENAI_API_KEY chưa set. Export rồi chạy lại.")
print(" set OPENAI_API_KEY=sk-...")
sys.exit(1)
async def main():
from database import get_db_context, async_engine
from database import Base
from services.blueprint_service import BlueprintService
from services.project_orchestrator import ProjectOrchestrator
# 1. Ensure tables exist (idempotent — create_all is safe)
async with async_engine.begin() as conn:
await conn.run_sync(Base.metadata.create_all)
# 2. Create a demo user (skip nếu đã tồn tại)
from models import AuthUser
from sqlalchemy import select
DEMO_USER_ID = "demo-e2e-user-001"
async with get_db_context() as db:
existing = (await db.execute(select(AuthUser).where(AuthUser.id == DEMO_USER_ID))).scalar_one_or_none()
if not existing:
db.add(AuthUser(
id=DEMO_USER_ID,
email="demo@paperclip.local",
email_verified=True,
password_hash="not-used",
))
await db.commit()
print("✅ Demo user created")
else:
print("✅ Demo user already exists")
# 3. Deploy company từ blueprint
async with get_db_context() as db:
company = await BlueprintService.deploy_company_from_blueprint(
db=db,
company_name="Demo Software Co",
blueprint_name="software-studio",
owner_id=DEMO_USER_ID,
description="E2E demo — automated pipeline test",
issue_prefix="DEMO",
budget_monthly_cents=100_000,
)
await db.commit()
company_id = company.id
print(f"✅ Company deployed: {company.name} (id={company_id})")
# 4. Load pipeline spec từ blueprint
bp = BlueprintService.load_blueprint("software-studio")
pipeline_spec = [node.model_dump() for node in bp.pipeline]
print(f"✅ Pipeline spec loaded: {len(pipeline_spec)} nodes")
# 5. Run orchestrator
brief = (
"Build a minimal REST API for a Todo app using Python FastAPI. "
"The API should support CRUD operations on tasks. "
"Deliverable: a working spec, a basic implementation plan, and QA sign-off."
)
async def auto_approver(cid):
from database import get_db_context
from models import Approval
from sqlalchemy import select
print("🤖 Auto-approver task started")
while True:
try:
await asyncio.sleep(2.0)
async with get_db_context() as db:
stmt = select(Approval).where(
Approval.company_id == cid,
Approval.resolution.is_(None)
)
res = await db.execute(stmt)
pending = res.scalars().all()
for app in pending:
print(f"🤖 Automatically approving approval request ID {app.id} for stage {app.resource_id}")
app.resolution = "approved"
app.resolved_at = datetime.utcnow()
app.resolution_note = "Auto-approved by E2E demo runner"
if pending:
await db.commit()
except asyncio.CancelledError:
break
except Exception as e:
print(f"Error in auto_approver: {e}")
orchestrator = ProjectOrchestrator(
company_id=company_id,
project_brief=brief,
pipeline_spec=pipeline_spec,
poll_interval=5.0,
max_wait_per_agent=300,
)
approver_task = asyncio.create_task(auto_approver(company_id))
try:
print("🚀 Running pipeline...")
report = await orchestrator.run()
finally:
approver_task.cancel()
# 6. Print CEO report
print("\n" + "=" * 60)
print("📊 CEO PIPELINE REPORT")
print("=" * 60)
print(f"Status : {report['status']}")
print(f"Duration : {report['total_duration_s']}s")
print(f"Agents : {report['agents_succeeded']}/{report['agents_executed']} succeeded")
print(f"Deliverables:")
for d in report.get("deliverables", []):
print(f" - {d}")
print("\nAgent Steps:")
for step in report.get("steps", []):
icon = "✅" if step["status"] == "completed" else "❌"
print(f" {icon} {step.get('agent_name', '?')} ({step.get('role', '?')}) — {step['status']}")
# 7. Lưu docs/demo_run.md
docs_dir = BACKEND_DIR.parent / "docs"
docs_dir.mkdir(exist_ok=True)
out_file = docs_dir / "demo_run.md"
md_lines = [
"# Demo E2E Run",
f"\n**Date:** {report['started_at']}",
f"**Status:** {report['status']}",
f"**Duration:** {report['total_duration_s']}s",
f"**Blueprint:** software-studio",
"\n## Agent Steps\n",
]
for step in report.get("steps", []):
icon = "✅" if step["status"] == "completed" else "❌"
md_lines.append(f"- {icon} **{step.get('agent_name', '?')}** ({step.get('role', '?')}) — {step['status']}")
md_lines += ["\n## Deliverables\n"]
for d in report.get("deliverables", []):
md_lines.append(f"- `{d}`")
md_lines += ["\n## Full Report\n", "```json", json.dumps(report, indent=2, default=str), "```"]
out_file.write_text("\n".join(md_lines), encoding="utf-8")
print(f"\n✅ Saved to {out_file}")
# 8. Final exit code
if report["status"] != "completed":
print("\n❌ Pipeline did NOT complete successfully")
sys.exit(1)
print("\n✅ Demo E2E PASSED")
if __name__ == "__main__":
asyncio.run(main())
import logging
import uuid
from datetime import datetime
from sqlalchemy import select
from models import Routine, RoutineTrigger, Company
from models.database import get_db_context
from tasks.agent_tasks import run_pipeline_orchestration
logger = logging.getLogger(__name__)
def cron_matches(cron: str, dt: datetime) -> bool:
"""Check if a datetime matches a standard 5-field cron expression."""
fields = cron.split()
if len(fields) != 5:
return False
def match_field(field: str, val: int) -> bool:
if field == "*":
return True
if "," in field:
return any(match_field(f, val) for f in field.split(","))
if "-" in field:
try:
start, end = map(int, field.split("-"))
return start <= val <= end
except ValueError:
return False
if "/" in field:
try:
parts = field.split("/")
base = parts[0]
step = int(parts[1])
if base == "*":
return val % step == 0
if "-" in base:
start, end = map(int, base.split("-"))
return start <= val <= end and (val - start) % step == 0
return val % step == int(base) % step
except ValueError:
return False
try:
return int(field) == val
except ValueError:
return False
cron_weekday = (dt.weekday() + 1) % 7
return (
match_field(fields[0], dt.minute) and
match_field(fields[1], dt.hour) and
match_field(fields[2], dt.day) and
match_field(fields[3], dt.month) and
match_field(fields[4], cron_weekday)
)
class RoutineSchedulerService:
@classmethod
async def tick(cls, db=None):
"""Evaluate all active routines and trigger matching ones."""
now = datetime.utcnow()
if db is not None:
await cls._tick_process(db, now)
else:
async with get_db_context() as session:
await cls._tick_process(session, now)
@classmethod
async def _tick_process(cls, db, now):
stmt = select(Routine).where(Routine.is_enabled == True)
result = await db.execute(stmt)
routines = result.scalars().all()
for routine in routines:
# Avoid double triggers within the same minute
if routine.last_triggered_at:
diff = (now - routine.last_triggered_at).total_seconds()
if diff < 59:
continue
if cron_matches(routine.cron, now):
logger.info(f"[RoutineScheduler] Triggering routine: {routine.name} ({routine.id})")
routine.last_triggered_at = now
# Create trigger execution log entry
trigger = RoutineTrigger(
id=uuid.uuid4().hex,
routine_id=routine.id,
triggered_at=now,
status="success",
issues_created=0
)
db.add(trigger)
# Retrieve the blueprint details of the company
co_result = await db.execute(select(Company).where(Company.id == routine.company_id))
company = co_result.scalar_one_or_none()
if not company:
logger.warning(f"[RoutineScheduler] Company {routine.company_id} not found for routine {routine.id}")
trigger.status = "failure"
trigger.error = "Company not found"
continue
blueprint_name = company.blueprint_name or "software-studio"
from services.blueprint_service import BlueprintService
try:
blueprint = BlueprintService.load_blueprint(blueprint_name)
pipeline_spec = [node.model_dump() for node in blueprint.pipeline]
except Exception as e:
logger.error(f"[RoutineScheduler] Failed to load blueprint for routine {routine.id}: {e}")
trigger.status = "failure"
trigger.error = f"Failed to load blueprint: {e}"
continue
project_brief = routine.query or f"Periodic CEO report for {company.name}"
# Queue run_pipeline_orchestration Celery task
run_pipeline_orchestration.delay(
company_id=routine.company_id,
project_brief=project_brief,
pipeline_spec=pipeline_spec
)
await db.commit()
......@@ -198,3 +198,159 @@ async def test_blueprint_validation_fails_for_bad_org_chart():
with pytest.raises(ValueError, match="non-existent role 'ghost'"):
blueprint = CompanyBlueprint.model_validate(bad_yaml)
BlueprintService.validate_org_chart(blueprint)
@pytest.mark.asyncio
async def test_blueprint_custom_plugins_and_skills():
"""Verify that deploying a company from a blueprint with plugins, skills, and secrets provisions them correctly."""
from schemas.blueprint import CompanyBlueprint
from models import CompanySecret, CompanySkill, PluginConfig
yaml_data = {
"name": "Custom Blueprint",
"industry": "custom",
"description": "Blueprint testing plugins and skills",
"roles": [
{"role": "lead_agent", "name": "Lead Agent", "adapter_type": "openai"}
],
"pipeline": [
{"role": "lead_agent"}
],
"plugins": [
{"name": "slack", "version": "1.2.0", "description": "Slack plugin", "config": {"webhook": "https://slack..."}}
],
"skills": [
"Coding Skill",
{"name": "Designing Skill", "level": 85, "description": "Architecture design"}
],
"env_vars": {
"DB_PASSWORD": "secret_db_pass_123",
"PUBLIC_API": "https://api.example.com"
}
}
# Verify parsing works
blueprint = CompanyBlueprint.model_validate(yaml_data)
assert len(blueprint.plugins) == 1
assert len(blueprint.skills) == 2
assert len(blueprint.env_vars) == 2
# Test deploy logic. Mock DB session and check objects added.
mock_db = MagicMock()
added_objs = []
def mock_add(obj):
added_objs.append(obj)
mock_db.add = mock_add
mock_db.flush = AsyncMock()
mock_db.refresh = AsyncMock()
# Mock global plugin check (Plugin not found, so it will be created)
mock_db.execute = AsyncMock()
mock_plugin_result = MagicMock()
mock_plugin_result.scalar_one_or_none.return_value = None
mock_db.execute.return_value = mock_plugin_result
with patch("services.blueprint_service.BlueprintService.load_blueprint", return_value=blueprint):
with patch("services.blueprint_service.log_activity", AsyncMock()):
company = await BlueprintService.deploy_company_from_blueprint(
db=mock_db,
company_name="Test Company",
blueprint_name="custom-bp",
owner_id="owner_user"
)
# Verify Company, Membership, Agent, Secret, PluginConfig, Skill objects are created and added to db.
secrets = [o for o in added_objs if isinstance(o, CompanySecret)]
skills = [o for o in added_objs if isinstance(o, CompanySkill)]
plugin_configs = [o for o in added_objs if isinstance(o, PluginConfig)]
assert len(secrets) == 2
assert any(s.key == "DB_PASSWORD" and s.is_sensitive is True for s in secrets)
assert any(s.key == "PUBLIC_API" and s.is_sensitive is False for s in secrets)
assert len(skills) == 2
assert any(sk.name == "Coding Skill" and sk.level == 100 for sk in skills)
assert any(sk.name == "Designing Skill" and sk.level == 85 for sk in skills)
assert len(plugin_configs) == 1
assert plugin_configs[0].config == {"webhook": "https://slack..."}
@pytest.mark.asyncio
async def test_routine_scheduler_trigger_and_orchestration():
"""Verify that RoutineSchedulerService evaluates cron expressions and triggers pipeline runs correctly."""
from models import Routine, RoutineTrigger, Company
from services.routine_scheduler import RoutineSchedulerService, cron_matches
from datetime import datetime
# 1. Test cron matching logic
dt = datetime(2026, 6, 2, 12, 0, 0) # Tuesday
assert cron_matches("* * * * *", dt) is True
assert cron_matches("0 12 * * *", dt) is True
assert cron_matches("0 13 * * *", dt) is False
assert cron_matches("*/5 * * * *", dt) is True
# Tuesday is weekday 1 in python (Monday=0, Tuesday=1)
# cron_weekday = (1 + 1) % 7 = 2
assert cron_matches("* * * * 2", dt) is True
assert cron_matches("* * * * 3", dt) is False
# 2. Test scheduler tick logic. Mock DB session and check triggered routines.
mock_db = MagicMock()
# Mock Company
mock_company = Company(
id="co_sch_123",
name="Scheduled Co",
blueprint_name="software-studio"
)
# Mock Routine
mock_routine = Routine(
id="routine_123",
company_id="co_sch_123",
name="Weekly Report",
cron="* * * * *", # matches any time
is_enabled=True,
query="Weekly strategy update",
last_triggered_at=None
)
# Mock DB query executions
mock_routines_result = MagicMock()
mock_routines_result.scalars.return_value.all.return_value = [mock_routine]
mock_company_result = MagicMock()
mock_company_result.scalar_one_or_none.return_value = mock_company
mock_db.execute = AsyncMock()
mock_db.execute.side_effect = [mock_routines_result, mock_company_result]
mock_db.add = MagicMock()
mock_db.commit = AsyncMock()
# Mock Celery delay
with patch("tasks.agent_tasks.run_pipeline_orchestration.delay") as mock_delay:
from schemas.blueprint import CompanyBlueprint
dummy_blueprint = CompanyBlueprint(
name="Studio",
industry="software",
pipeline=[{"role": "ceo"}]
)
with patch("services.blueprint_service.BlueprintService.load_blueprint", return_value=dummy_blueprint):
await RoutineSchedulerService.tick(db=mock_db)
# Assert Celery task was queued
mock_delay.assert_called_once_with(
company_id="co_sch_123",
project_brief="Weekly strategy update",
pipeline_spec=[{"role": "ceo", "depends_on": [], "approval_gate": False}]
)
# Assert Routine Trigger history and last_triggered_at were set
assert mock_routine.last_triggered_at is not None
mock_db.commit.assert_called_once()
# Demo E2E Run
**Date:** 2026-06-02T14:28:31.221087
**Status:** completed
**Duration:** 77.0s
**Blueprint:** software-studio
## Agent Steps
- ✅ **CEO Agent** (ceo) — completed
- ✅ **Product Manager Agent** (pm) — completed
- ✅ **Software Architect Agent** (designer) — completed
- ✅ **Lead Developer Agent** (coder) — completed
- ✅ **QA Specialist Agent** (qa) — completed
- ✅ **DevOps Engineer Agent** (devops) — completed
## Deliverables
- `ceo_output.md`
- `ceo_strategy.md`
- `coder_output.md`
- `design_spec.md`
- `designer_output.md`
- `dev_log.txt`
- `devops_output.md`
- `index.html`
- `pm_output.md`
- `qa_output.md`
- `qa_report.txt`
- `spec.md`
## Full Report
```json
{
"pipeline_id": "a385b7aaee904b00ba8ceec01a3878c1",
"company_id": "c7cfd603b9504bfd934a66faf9b48363",
"status": "completed",
"total_duration_s": 77.0,
"agents_executed": 6,
"agents_succeeded": 6,
"agents_failed": 0,
"deliverables": [
"ceo_output.md",
"ceo_strategy.md",
"coder_output.md",
"design_spec.md",
"designer_output.md",
"dev_log.txt",
"devops_output.md",
"index.html",
"pm_output.md",
"qa_output.md",
"qa_report.txt",
"spec.md"
],
"steps": [
{
"step": 1,
"agent_id": "ceo_41e8c849",
"agent_name": "CEO Agent",
"role": "ceo",
"run_id": "14f136fa17af447d8531a2ed1dc092a0",
"status": "completed",
"duration_s": 9.8,
"error": null
},
{
"step": 2,
"agent_id": "pm_04078779",
"agent_name": "Product Manager Agent",
"role": "pm",
"run_id": "1fb0960c6cd74c0bb1a4a94f5a17333d",
"status": "completed",
"duration_s": 13.2,
"error": null
},
{
"step": 3,
"agent_id": "designer_8c957a41",
"agent_name": "Software Architect Agent",
"role": "designer",
"run_id": "2b181e6d8cec49b0854bd56b4a457654",
"status": "completed",
"duration_s": 9.7,
"error": null
},
{
"step": 4,
"agent_id": "coder_ae0f63cf",
"agent_name": "Lead Developer Agent",
"role": "coder",
"run_id": "002d6d62123b48e0b940b6fde8f4e8ad",
"status": "completed",
"duration_s": 11.8,
"error": null
},
{
"step": 5,
"agent_id": "qa_7d279926",
"agent_name": "QA Specialist Agent",
"role": "qa",
"run_id": "c81cd6e371c9477cb96e0f3c3e2821c1",
"status": "completed",
"duration_s": 10.2,
"error": null
},
{
"step": 6,
"agent_id": "devops_ff60aad0",
"agent_name": "DevOps Engineer Agent",
"role": "devops",
"run_id": "989711aa25734782a133895ac59c400b",
"status": "completed",
"duration_s": 9.7,
"error": null
}
],
"started_at": "2026-06-02T14:28:31.221087",
"completed_at": "2026-06-02T14:29:48.183474"
}
```
\ No newline at end of file
......@@ -5,6 +5,7 @@ import { Layout } from "./components/Layout";
import { OnboardingWizard } from "./components/OnboardingWizard";
import { CloudAccessGate } from "./components/CloudAccessGate";
import { Dashboard } from "./pages/Dashboard";
import { Onboarding } from "./pages/Onboarding";
import { DashboardLive } from "./pages/DashboardLive";
import { Companies } from "./pages/Companies";
import { Agents } from "./pages/Agents";
......@@ -29,6 +30,7 @@ import { Approvals } from "./pages/Approvals";
import { ApprovalDetail } from "./pages/ApprovalDetail";
import { Costs } from "./pages/Costs";
import { Activity } from "./pages/Activity";
import { PipelineRuns } from "./pages/PipelineRuns";
import { Inbox } from "./pages/Inbox";
import { CompanySettings } from "./pages/CompanySettings";
import { CompanyEnvironments } from "./pages/CompanyEnvironments";
......@@ -51,6 +53,7 @@ import { PluginPage } from "./pages/PluginPage";
import { OrgChart } from "./pages/OrgChart";
import { NewAgent } from "./pages/NewAgent";
import { AuthPage } from "./pages/Auth";
import { Integrations } from "./pages/Integrations";
import { BoardClaimPage } from "./pages/BoardClaim";
import { CliAuthPage } from "./pages/CliAuth";
import { InviteLandingPage } from "./pages/InviteLanding";
......@@ -67,7 +70,7 @@ function boardRoutes() {
<Route index element={<Navigate to="dashboard" replace />} />
<Route path="dashboard" element={<Dashboard />} />
<Route path="dashboard/live" element={<DashboardLive />} />
<Route path="onboarding" element={<OnboardingRoutePage />} />
<Route path="onboarding" element={<Onboarding />} />
<Route path="companies" element={<Companies />} />
<Route path="company/settings" element={<CompanySettings />} />
<Route path="company/settings/environments" element={<CompanyEnvironments />} />
......@@ -131,6 +134,8 @@ function boardRoutes() {
<Route path="approvals/:approvalId" element={<ApprovalDetail />} />
<Route path="costs" element={<Costs />} />
<Route path="activity" element={<Activity />} />
<Route path="integrations" element={<Integrations />} />
<Route path="pipeline-runs" element={<PipelineRuns />} />
<Route path="inbox" element={<InboxRootRedirect />} />
<Route path="inbox/mine" element={<Inbox />} />
<Route path="inbox/recent" element={<Inbox />} />
......@@ -285,7 +290,7 @@ export function App() {
<Route element={<CloudAccessGate />}>
<Route index element={<CompanyRootRedirect />} />
<Route path="onboarding" element={<OnboardingRoutePage />} />
<Route path="onboarding" element={<Onboarding />} />
<Route path="instance" element={<Navigate to="/instance/settings/general" replace />} />
<Route path="instance/settings" element={<Layout />}>
<Route index element={<Navigate to="general" replace />} />
......@@ -330,6 +335,8 @@ export function App() {
<Route path="execution-workspaces/:workspaceId/live-tracking" element={<UnprefixedBoardRedirect />} />
<Route path="execution-workspaces/:workspaceId/ceo-briefing" element={<UnprefixedBoardRedirect />} />
<Route path="marketplace" element={<UnprefixedBoardRedirect />} />
<Route path="integrations" element={<UnprefixedBoardRedirect />} />
<Route path="pipeline-runs" element={<UnprefixedBoardRedirect />} />
<Route path=":companyPrefix" element={<Layout />}>
{boardRoutes()}
</Route>
......
......@@ -11,9 +11,11 @@ import {
Boxes,
Repeat,
GitBranch,
Plug,
Settings,
Store,
Cpu,
CheckSquare,
} from "lucide-react";
import { useQuery } from "@tanstack/react-query";
import { NavLink } from "@/lib/router";
......@@ -24,6 +26,7 @@ import { SidebarAgents } from "./SidebarAgents";
import { useDialogActions } from "../context/DialogContext";
import { useCompany } from "../context/CompanyContext";
import { heartbeatsApi } from "../api/heartbeats";
import { approvalsApi } from "../api/approvals";
import { instanceSettingsApi } from "../api/instanceSettings";
import { queryKeys } from "../lib/queryKeys";
import { useInboxBadge } from "../hooks/useInboxBadge";
......@@ -49,6 +52,16 @@ export function Sidebar() {
const liveRunCount = liveRuns?.length ?? 0;
const showWorkspacesLink = experimentalSettings?.enableIsolatedWorkspaces === true;
const { data: approvals } = useQuery({
queryKey: queryKeys.approvals.list(selectedCompanyId!),
queryFn: () => approvalsApi.list(selectedCompanyId!),
enabled: !!selectedCompanyId,
refetchInterval: 30_000,
});
const pendingApprovalsCount = (approvals ?? []).filter(
(a) => a.status === "pending" || a.status === "revision_requested",
).length;
const pluginContext = {
companyId: selectedCompanyId,
companyPrefix: selectedCompany?.issuePrefix ?? null,
......@@ -92,6 +105,13 @@ export function Sidebar() {
badgeTone={inboxBadge.failedRuns > 0 ? "danger" : "default"}
alert={inboxBadge.failedRuns > 0}
/>
<SidebarNavItem
to="/approvals"
label="Phê duyệt"
icon={CheckSquare}
badge={pendingApprovalsCount}
badgeTone="danger"
/>
<SidebarNavItem to="/marketplace" label="Marketplace" icon={Store} />
<SidebarNavItem to="/cubesandbox" label="Cube Sandbox" icon={Cpu} />
<PluginSlotOutlet
......@@ -106,6 +126,7 @@ export function Sidebar() {
<SidebarSection label="Công việc">
<SidebarNavItem to="/issues" label="Vấn đề" icon={CircleDot} />
<SidebarNavItem to="/routines" label="Quy trình" icon={Repeat} />
<SidebarNavItem to="/pipeline-runs" label="Lịch sử Pipeline" icon={History} />
<PluginLauncherOutlet
placementZones={["sidebar"]}
context={pluginContext}
......@@ -125,6 +146,7 @@ export function Sidebar() {
<SidebarSection label="Phòng ban">
<SidebarNavItem to="/org" label="Sơ đồ" icon={Network} />
<SidebarNavItem to="/skills" label="Kỹ năng" icon={Boxes} />
<SidebarNavItem to="/integrations" label="Kết nối" icon={Plug} />
<SidebarNavItem to="/costs" label="Chi phí" icon={DollarSign} />
<SidebarNavItem to="/activity" label="Hoạt động" icon={History} />
<SidebarNavItem to="/company/settings" label="Thiết lập" icon={Settings} />
......
This diff is collapsed.
This diff is collapsed.
# DOING #06: Pillar N — NocoBase sống thật (npm, không docker)
> **Dành cho agent kế tiếp** — đọc hết file này trước khi code.
> Session trước đã làm xong phần code cốt lõi. Việc còn lại là CHẠY THẬT và verify.
---
## 📌 Trạng thái hiện tại (đã xong)
| File | Trạng thái |
|------|-----------|
| `backend/services/nocobase_connector.py` | ✅ Có `ping()`, `login()`, `ensure_collection()`, `list_records()`, `from_company()` |
| `backend/api/routes/companies.py` | ✅ Provisioning block dùng `ping()` trước, log ERROR loud nếu NocoBase chết |
| `backend/tests/test_nocobase_live.py` | ✅ Viết xong, skip nếu `ping()` fail, PASS nếu NocoBase up |
## ✅ Đã hoàn thành (đã verify live)
### N-A: Cập nhật `.env.example`
Đã cập nhật `.env.example` với 3 cấu hình NocoBase:
- `NOCOBASE_URL=http://localhost:13002`
- `NOCOBASE_EMAIL=admin@nocobase.com`
- `NOCOBASE_PASSWORD=admin123`
### N-B: Bootstrap NocoBase qua npm
Chạy thành công:
- `npm.cmd run nocobase install`
- `npm.cmd run start`
NocoBase đã up và lắng nghe tại cổng `13001` (PM2 online).
### N-C: Verify ping trả True
Ping test kết nối thành công: `PING: True`.
### N-D: Chạy live test — PASS thành công
Chạy:
```bash
cd D:\a\ai_canifa_company\backend
./.venv/Scripts/pytest tests/test_nocobase_live.py -v
```
**Kết quả: `1 passed` (Đạt Success criteria).**
---
## 🛠️ Tool workflow cho agent
```
1. Đọc file này xong
2. Edit backend/.env.example (thêm 3 dòng NocoBase)
3. Bash: npx create-nocobase-app ... (bootstrap npm)
4. Bash: kiểm tra ping()
5. Bash: pytest test_nocobase_live.py → phải PASS
6. Update plan/doings/05_pillars_completion.md: tick ✅ N1-N4.2
```
## 📍 Paths quan trọng
| | Path |
|-|------|
| venv python | `D:\a\ai_canifa_company\backend\.venv\Scripts\python.exe` |
| pytest | `./.venv/Scripts/pytest` (chạy từ `backend/`) |
| live test | `backend/tests/test_nocobase_live.py` |
| connector | `backend/services/nocobase_connector.py` |
| companies route | `backend/api/routes/companies.py` (L156-196, provisioning block) |
## ⚠️ Lưu ý quan trọng
1. **DÙNG npm KHÔNG docker** — user đã dặn rõ
2. Nếu codegraph MCP không có trong session, dùng sqlite3 query:
```bash
./.venv/Scripts/python.exe -c "import sqlite3; conn = sqlite3.connect('../.codegraph/codegraph.db'); print(conn.execute(\"SELECT name, file_path, start_line FROM nodes WHERE name='NocoBaseConnector'\").fetchall())"
```
(Chạy từ `backend/`)
3. `tasks/` và `services/` CHƯA được codegraph index — đọc file trực tiếp cho 2 dir này
4. Mọi Pydantic body model phải dùng `snake_case` field (CaseConversionMiddleware global rewrite)
---
_Created: 2026-06-02 | Pillar N handoff_
# DOING #08 — Fix Deliverable Writer (Pipeline output = empty bug)
## 🐛 Bug
```
"deliverables": [] ← dù 6/6 agents completed (demo_run.md 2026-06-01)
```
Agent chạy xong nhưng không ghi bất kỳ file nào vào `workspaces/{company_id}/shared/`.
CEO report có status "completed" nhưng không có nội dung thật.
## 🎯 Success criteria
```
python backend/scratch/demo_e2e.py
→ "deliverables": ["ceo_output.md", "pm_output.md", ...] ← KHÔNG rỗng
→ workspaces/{company_id}/shared/*.md có nội dung thật
```
---
## ✅ Checklist
### Phase 1 — Diagnose
- [ ] D1. Đọc `backend/services/project_orchestrator.py` L459-540 (`_run_single_agent`) — xem system_message truyền vào agent như thế nào
- [ ] D2. Xem `backend/api/routes/agents.py` → `run_agent_in_background` — agent có nhận system_message không, output được lưu ở đâu
- [ ] D3. Kiểm tra `HeartbeatRun` model có field `output` / `result` không:
```bash
./.venv/Scripts/python.exe -c "import sqlite3; conn=sqlite3.connect('../.codegraph/codegraph.db'); print(conn.execute(\"SELECT signature FROM nodes WHERE name='HeartbeatRun' AND kind='class'\").fetchall())"
```
### Phase 2 — Fix (chọn 1 trong 2 cách)
**Cách A — System Prompt Injection (khuyến nghị, ít invasive)**
- [ ] A1. Trong `_run_single_agent` (project_orchestrator.py), sửa `system_message`:
```python
shared_path = f"workspaces/{self.company_id}/shared/{role}_output.md"
system_message = (
f"You are {name}, playing the role of {role}. "
f"After completing your work, write a markdown deliverable summary to: {shared_path}\n"
f"Format:\n## {name} Deliverable\n[your output here]"
)
```
- [ ] A2. Đảm bảo `shared_dir` được tạo trước khi agent run (đã có ở L141-142, chỉ cần verify)
**Cách B — Post-run Parser (nếu agent không ghi file được)**
- [ ] B1. Sau khi agent run xong (sau vòng poll L504+), đọc `HeartbeatRun.output`
- [ ] B2. Save nội dung vào `shared_dir / f"{role}_output.md"`
### Phase 3 — Verify
- [ ] V1. Chạy lại demo_e2e:
```bash
cd D:\a\ai_canifa_company\backend
python scratch/demo_e2e.py
```
- [ ] V2. Kiểm tra `deliverables` không còn `[]`
- [ ] V3. Kiểm tra files tồn tại: `ls workspaces/*/shared/`
- [ ] V4. Update `docs/demo_run.md` với kết quả mới
---
## 📍 Files cần sửa
| File | Dòng | Thay đổi |
|------|------|---------|
| `backend/services/project_orchestrator.py` | ~493 | Sửa `system_message` trong `_run_single_agent` |
| `backend/scratch/demo_e2e.py` | — | Không cần sửa (chỉ re-run) |
## ⚠️ Lưu ý
- `tasks/` và `services/` CHƯA indexed trong codegraph → đọc file trực tiếp
- Không sửa `backend/agent/run_agent.py` (4000+ dòng Hermes runtime, too invasive)
- Nếu agent không có permission ghi file → dùng Cách B (post-run parser)
# DOING #09 — Onboarding Wizard UI (Blueprint → Deploy in 1 Flow)
## 🎯 Mục tiêu
Non-technical user tạo company AI đầy đủ trong < 2 phút, không cần gọi API thủ công.
## 🎯 Success criteria
```
User vào /onboarding
→ Chọn "Software Studio" blueprint (xem preview org chart)
→ Điền tên company
→ Click "Deploy"
→ Redirect /companies/{id}/dashboard
→ Company có đủ agents, pipeline configured
```
---
## ✅ Checklist
### Phase 1 — Check APIs & existing components
- [ ] P1. Verify API endpoints hoạt động:
```bash
# Blueprints list
curl http://localhost:3100/api/companies/blueprints
# Deploy
curl -X POST http://localhost:3100/api/companies/from-blueprint \
-H "Content-Type: application/json" \
-d '{"name":"Test Co","blueprint_name":"software-studio","issue_prefix":"TST","budget_monthly_cents":0}'
```
- [ ] P2. Xem `frontend/src/api/` — có `companiesApi` chưa? Xem pattern của `approvalsApi` để follow
- [ ] P3. Xem `frontend/src/pages/` — có page nào tương tự (wizard-style) để tham khảo layout không
### Phase 2 — Tạo API layer
- [ ] A1. Thêm vào `frontend/src/api/companies.ts` (hoặc tạo mới nếu chưa có):
```typescript
export const blueprintsApi = {
list: () => apiClient.get('/companies/blueprints'),
getDetail: (name: string) => apiClient.get(`/companies/blueprints/${name}`),
deployFromBlueprint: (data: {
name: string;
blueprint_name: string; // ← snake_case (CaseConversionMiddleware sẽ convert)
description?: string;
issue_prefix?: string;
budget_monthly_cents?: number;
}) => apiClient.post('/companies/from-blueprint', data),
}
```
> ⚠️ Dùng camelCase trong TypeScript interface OK, nhưng JSON key gửi lên phải match với gì FE đang dùng. Check `CaseConversionMiddleware` behavior — BE nhận snake_case sau khi middleware convert từ camelCase của FE.
### Phase 3 — Tạo page component
- [ ] B1. Tạo `frontend/src/pages/Onboarding.tsx`:
- Step 1: Blueprint picker (grid 4 cards: Software Studio / E-commerce / Fashion / Marketing)
- Mỗi card: tên, description, số agents, industry icon
- Click → select + preview org chart (list roles)
- Step 2: Company config form (name, issue_prefix, budget optional)
- Step 3: Deploy button → loading spinner → success redirect
- [ ] B2. State management: `useState` đơn giản (3 fields), không cần Redux/Zustand
- [ ] B3. Error handling: show toast nếu deploy fail
### Phase 4 — Wiring
- [ ] C1. Thêm route vào `frontend/src/App.tsx` (hoặc router file):
```tsx
<Route path="/onboarding" element={<Onboarding />} />
```
- [ ] C2. Thêm link "Create Company" trong sidebar hoặc Companies page
### Phase 5 — Verify
- [ ] V1. `npm run dev` (frontend) + backend running
- [ ] V2. Navigate `/onboarding` → chọn blueprint → deploy → verify company xuất hiện trong `/companies`
- [ ] V3. Check agents được tạo đúng: `GET /companies/{id}` có agents không
---
## 📍 Files cần tạo/sửa
| File | Action |
|------|--------|
| `frontend/src/pages/Onboarding.tsx` | TẠO MỚI (~150 lines) |
| `frontend/src/api/companies.ts` | THÊM blueprintsApi |
| `frontend/src/App.tsx` (hoặc router) | THÊM route `/onboarding` |
## ⚠️ Lưu ý
- Blueprint API trả về `blueprint_name` (snake_case) — FE gọi `GET /companies/blueprints` trả list
- Deploy body: `blueprint_name` (không phải `blueprintName`) — check xem FE đang auto-convert key không
- Không cần auth gate cho wizard nếu là demo mode
# DOING #10 — Approval Sidebar Badge (Pending count realtime)
## 🎯 Mục tiêu
User biết ngay khi có approval đang chờ mà không cần vào `/approvals` page thủ công.
Pipeline approval gate không bị block ngầm vô thời hạn.
## 🎯 Success criteria
```
Pipeline chạy đến approval_gate → tạo Approval pending
→ Sidebar "Approvals" hiện badge đỏ "1"
→ User click → vào Approvals page → approve
→ Badge biến mất
```
---
## ✅ Checklist
### Phase 1 — Locate sidebar component
- [ ] S1. Tìm sidebar component trong frontend:
```
frontend/src/components/ → tìm Sidebar.tsx hoặc NavSidebar.tsx hoặc AppLayout.tsx
```
- [ ] S2. Xem cách sidebar render các nav items — có pattern badge nào chưa (ví dụ Inbox badge)?
- [ ] S3. Xem `frontend/src/api/` → `approvalsApi.list()` dùng endpoint gì
### Phase 2 — Thêm badge logic
- [ ] B1. Trong sidebar component, thêm hook lấy pending count:
```tsx
const [pendingCount, setPendingCount] = useState(0);
useEffect(() => {
const fetchPending = async () => {
try {
const res = await approvalsApi.list(companyId, { status: 'pending' });
// hoặc: GET /companies/{id}/approvals?status=pending
setPendingCount(res.data?.length ?? 0);
} catch {}
};
fetchPending();
const interval = setInterval(fetchPending, 30_000); // poll 30s
return () => clearInterval(interval);
}, [companyId]);
```
- [ ] B2. Render badge trên Approvals nav item:
```tsx
{pendingCount > 0 && (
<span className="ml-auto rounded-full bg-red-500 px-1.5 py-0.5 text-xs text-white">
{pendingCount}
</span>
)}
```
- [ ] B3. Check `approvalsApi` có hỗ trợ filter `status` param không — nếu không thì filter client-side
### Phase 3 — Verify
- [ ] V1. Tạo 1 Approval pending thủ công (qua API hoặc chạy pipeline với approval_gate)
- [ ] V2. Reload sidebar → badge hiện
- [ ] V3. Approve → badge biến mất trong ≤ 30s
---
## 📍 Files cần sửa
| File | Thay đổi |
|------|---------|
| `frontend/src/components/Sidebar.tsx` (hoặc tên tương tự) | Thêm polling hook + badge render |
| `frontend/src/api/approvals.ts` | Thêm `status` filter param nếu chưa có |
## ⚠️ Lưu ý
- Poll 30s là đủ — không cần WebSocket cho feature này
- `companyId` lấy từ context/URL param (xem pattern của các page khác)
- Badge chỉ show khi `pendingCount > 0`, ẩn khi = 0
# DOING #11 — Pipeline Run History (Audit Trail)
## 🎯 Mục tiêu
Xem được lịch sử mọi lần chạy pipeline: status, duration, agent steps, errors.
## 🎯 Success criteria
```
GET /companies/{id}/pipeline-runs
→ list các pipeline runs với pipeline_id, status, duration, agents_succeeded/failed
GET /companies/{id}/pipeline-runs/{pipeline_id}
→ detail: từng agent step, run_id, duration, error
Frontend page /companies/{id}/pipelines
→ bảng runs + click vào xem detail
```
---
## ✅ Checklist
### Phase 1 — Research data model
- [ ] R1. Đọc `backend/models/` → tìm `HeartbeatRun` — có field `pipeline_id` chưa?
```bash
./.venv/Scripts/python.exe -c "
import sqlite3; conn=sqlite3.connect('../.codegraph/codegraph.db')
rows = conn.execute(\"SELECT name,file_path,start_line FROM nodes WHERE name='HeartbeatRun' AND kind='class'\").fetchall()
print(rows)
"
```
- [ ] R2. Đọc file HeartbeatRun model — list tất cả fields
- [ ] R3. Kiểm tra `ProjectOrchestrator._run_single_agent` — `db_run` được tạo với fields nào
### Phase 2 — Backend: thêm pipeline_id vào HeartbeatRun
- [ ] M1. Thêm field vào `HeartbeatRun` model:
```python
pipeline_id: Mapped[Optional[str]] = mapped_column(String(64), nullable=True, index=True)
pipeline_step: Mapped[Optional[int]] = mapped_column(Integer, nullable=True)
```
- [ ] M2. Tạo Alembic migration:
```bash
cd backend
./.venv/Scripts/python.exe -m alembic revision --autogenerate -m "add_pipeline_id_to_heartbeat_run"
./.venv/Scripts/python.exe -m alembic upgrade head
```
- [ ] M3. Trong `ProjectOrchestrator._run_single_agent` (L475-481), set:
```python
db_run = HeartbeatRun(
...
pipeline_id=self.pipeline_id,
pipeline_step=step,
)
```
### Phase 3 — Backend: API endpoint
- [ ] A1. Thêm vào `backend/api/routes/companies.py` (hoặc tạo `pipeline_runs.py`):
```python
@router.get("/{company_id}/pipeline-runs")
async def list_pipeline_runs(company_id: str, db=Depends(get_db_session)):
# GROUP BY pipeline_id, lấy min(started_at), max(completed_at), count status
...
@router.get("/{company_id}/pipeline-runs/{pipeline_id}")
async def get_pipeline_run(company_id: str, pipeline_id: str, db=Depends(get_db_session)):
# SELECT * FROM heartbeat_runs WHERE pipeline_id=...
...
```
### Phase 4 — Frontend page (optional — có thể để sau)
- [ ] F1. Tạo `frontend/src/pages/PipelineRuns.tsx`
- [ ] F2. Route `/companies/{id}/pipeline-runs`
- [ ] F3. Table: pipeline_id (truncated), status badge, started_at, duration, agents ✅/❌
### Phase 5 — Verify
- [ ] V1. Chạy `demo_e2e.py` → `GET /companies/{id}/pipeline-runs` trả list có 1 entry
- [ ] V2. `GET /pipeline-runs/{id}` trả 6 steps đúng với demo
- [ ] V3. Nếu có frontend: navigate page → thấy run history
---
## 📍 Files cần sửa
| File | Thay đổi |
|------|---------|
| `backend/models/heartbeat_runs.py` (hoặc tương tự) | Thêm `pipeline_id`, `pipeline_step` |
| `backend/alembic/versions/` | Migration mới |
| `backend/services/project_orchestrator.py` | Set `pipeline_id` khi tạo HeartbeatRun |
| `backend/api/routes/companies.py` | Thêm 2 endpoints |
| `frontend/src/pages/PipelineRuns.tsx` | TẠO MỚI (optional) |
## ⚠️ Lưu ý
- Làm Phase 1-3 (backend) trước, Phase 4 (frontend) có thể để sau
- `HeartbeatRun` có thể ở `backend/models/heartbeat_run.py` — dùng codegraph để tìm path chính xác
- Migration cần chạy `alembic upgrade head` trước khi test API
- `pipeline_id` = `self.pipeline_id` trong orchestrator (đã set ở `__init__` L117)
# DOING #29 — Swarm Tree Visualizer
## 🎯 Mục tiêu
Visualizing the parent-child issue hierarchy in an interactive, recursive tree view.
## 🎯 Success criteria
- [x] Backend supporting `parent_id` parameter filter in `list_issues` route.
- [x] Frontend API client updated to support `parentId` parameter in `issues.ts`.
- [x] Interactive `SwarmTree.tsx` tree node visualizer created.
- [x] Connected "Swarm Tree" tab in `IssueDetail.tsx` displaying the recursive issue spawn tree.
---
## 📍 Files đã sửa / tạo mới
- [NEW] [`frontend/src/components/SwarmTree.tsx`](file:///d:/a/ai_canifa_company/frontend/src/components/SwarmTree.tsx)
- [MODIFY] [`backend/api/routes/issues.py`](file:///d:/a/ai_canifa_company/backend/api/routes/issues.py)
- [MODIFY] [`frontend/src/api/index.ts`](file:///d:/a/ai_canifa_company/frontend/src/api/index.ts)
- [MODIFY] [`frontend/src/pages/IssueDetail.tsx`](file:///d:/a/ai_canifa_company/frontend/src/pages/IssueDetail.tsx)
# DOING #30 — Blueprint Custom Plugins & Skills
## 🎯 Mục tiêu
Support custom plugins, custom skills, and environment variables (secrets) defined in company blueprints to be automatically provisioned during deployment.
## 🎯 Success criteria
- [x] Schema extensions in `CompanyBlueprint` inside `backend/schemas/blueprint.py` to support `plugins`, `skills`, and `env_vars`.
- [x] BlueprintService modified to automatically populate tables `CompanySecret`, `CompanySkill`, and `PluginConfig` when deploying from blueprint.
- [x] Sensitive keys are automatically flagged as `is_sensitive=True`.
- [x] E2E backend tests verifying correct provisioning and creation of objects added.
---
## 📍 Files đã sửa / tạo mới
- [MODIFY] [`backend/schemas/blueprint.py`](file:///d:/a/ai_canifa_company/backend/schemas/blueprint.py)
- [MODIFY] [`backend/services/blueprint_service.py`](file:///d:/a/ai_canifa_company/backend/services/blueprint_service.py)
- [MODIFY] [`backend/tests/test_blueprint_e2e.py`](file:///d:/a/ai_canifa_company/backend/tests/test_blueprint_e2e.py)
# DOING #31 — CEO Report Auto-Schedule
## 🎯 Mục tiêu
Enable scheduled, recurring runs of orchestrator pipelines using Celery periodic tasks and cron expression schedules.
## 🎯 Success criteria
- [x] Integrate scheduling routines to evaluate cron expressions and trigger ProjectOrchestrator runs.
- [x] Created `RoutineSchedulerService` tick handler to match cron jobs and queue running pipeline orchestration.
- [x] Refactored tick loop to support database session injection for test environment compatibility.
- [x] End-to-end test execution verified successfully through mock scheduler test pipeline.
---
## 📍 Files đã sửa / tạo mới
- [MODIFY] [`backend/services/routine_scheduler.py`](file:///d:/a/ai_canifa_company/backend/services/routine_scheduler.py)
- [MODIFY] [`backend/tests/test_blueprint_e2e.py`](file:///d:/a/ai_canifa_company/backend/tests/test_blueprint_e2e.py)
- [MODIFY] [`backend/tasks/agent_tasks.py`](file:///d:/a/ai_canifa_company/backend/tasks/agent_tasks.py)
- [MODIFY] [`backend/celery_app.py`](file:///d:/a/ai_canifa_company/backend/celery_app.py)
- [MODIFY] [`backend/server.py`](file:///d:/a/ai_canifa_company/backend/server.py)
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