642 lines
29 KiB
Python
642 lines
29 KiB
Python
"""Automation & Agents plugin — Agent Builder, Automation Builder, Cron-Scheduler, Agent Runner.
|
|
|
|
Houses the core automation engine: AI agent definitions, event/schedule/manual
|
|
workflow automations, cron job scheduling, and execution logging.
|
|
|
|
Supports plugin contributions: when other plugins are activated, their
|
|
agent_definitions, automation_templates, cron_jobs, and heartbeat_configs
|
|
from the manifest are registered. On deactivation, they are removed.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
from typing import Any
|
|
|
|
from app.plugins.base import BasePlugin
|
|
from app.plugins.manifest import (
|
|
CronJobContribution,
|
|
FrontendMenuItem,
|
|
FrontendPageRoute,
|
|
FrontendSettingsPage,
|
|
PluginManifest,
|
|
PluginRouteDef,
|
|
)
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class AutomationPlugin(BasePlugin):
|
|
"""Automation & Agents plugin — Agent Builder, Automation Builder, Cron-Scheduler, Agent Runner."""
|
|
|
|
manifest = PluginManifest(
|
|
name="automation",
|
|
version="1.0.0",
|
|
display_name="Automation & Agents",
|
|
description=(
|
|
"Agent Builder, Automation Builder, Cron-Scheduler, and Agent Runner. "
|
|
"Define AI agents with LLM models and tools, create event/schedule/manual "
|
|
"automations with conditions and actions, schedule cron jobs, and track execution logs."
|
|
),
|
|
dependencies=[],
|
|
routes=[
|
|
PluginRouteDef(
|
|
path="/api/v1/automation",
|
|
module="app.plugins.builtins.automation.routes",
|
|
router_attr="router",
|
|
),
|
|
PluginRouteDef(
|
|
path="/api/v1/agents",
|
|
module="app.plugins.builtins.automation.agent_routes",
|
|
router_attr="router",
|
|
),
|
|
PluginRouteDef(
|
|
path="/api/v1/skills",
|
|
module="app.plugins.builtins.automation.skill_routes",
|
|
router_attr="router",
|
|
),
|
|
],
|
|
events=[
|
|
"contact.created",
|
|
"contact.updated",
|
|
"mail.received",
|
|
"workflow.timeout",
|
|
],
|
|
migrations=["0001_initial.sql", "0002_agent_subtasks.sql", "0003_skill_definitions.sql", "0004_run_steps_phase_f.sql"],
|
|
permissions=[
|
|
"automation:read",
|
|
"automation:write",
|
|
"automation:delete",
|
|
"automation:execute",
|
|
"automation:configure",
|
|
"agents:read",
|
|
"agents:write",
|
|
"agents:delete",
|
|
"agents:execute",
|
|
],
|
|
is_core=True,
|
|
menu_items=[
|
|
FrontendMenuItem(
|
|
label_key="nav.workflows",
|
|
label="Workflows",
|
|
path="/workflows",
|
|
icon="Workflow",
|
|
order=52,
|
|
),
|
|
FrontendMenuItem(
|
|
label_key="nav.importExport",
|
|
label="Import / Export",
|
|
path="/import-export",
|
|
icon="ArrowUpDown",
|
|
order=53,
|
|
),
|
|
FrontendMenuItem(
|
|
label_key="nav.dedupMerge",
|
|
label="Duplikate",
|
|
path="/contacts/dedup",
|
|
icon="Copy",
|
|
order=54,
|
|
),
|
|
FrontendMenuItem(
|
|
label_key="nav.tags",
|
|
label="Tags",
|
|
path="/tags",
|
|
icon="Tag",
|
|
order=55,
|
|
),
|
|
FrontendMenuItem(
|
|
label_key="nav.activity",
|
|
label="Aktivitäten",
|
|
path="/activity",
|
|
icon="Activity",
|
|
order=56,
|
|
),
|
|
],
|
|
page_routes=[
|
|
FrontendPageRoute(
|
|
path="/automation",
|
|
component="@/pages/AutomationDashboard",
|
|
order=50,
|
|
),
|
|
FrontendPageRoute(
|
|
path="/agents",
|
|
component="@/pages/AgentDashboard",
|
|
order=51,
|
|
),
|
|
FrontendPageRoute(
|
|
path="/workflows",
|
|
component="@/pages/Workflows",
|
|
order=52,
|
|
),
|
|
FrontendPageRoute(
|
|
path="/import-export",
|
|
component="@/pages/ImportExport",
|
|
order=53,
|
|
),
|
|
],
|
|
settings_pages=[
|
|
FrontendSettingsPage(
|
|
path="automation",
|
|
label_key="settings.automation",
|
|
label="Automation",
|
|
component="@/pages/AutomationSettings",
|
|
icon="Settings",
|
|
order=60,
|
|
),
|
|
],
|
|
cron_jobs=[
|
|
CronJobContribution(
|
|
name="backup_check",
|
|
cron_expression="0 2 * * *",
|
|
job_type="custom",
|
|
target_name="backup_check",
|
|
plugin_name="automation",
|
|
),
|
|
CronJobContribution(
|
|
name="search_index_check",
|
|
cron_expression="0 3 * * *",
|
|
job_type="custom",
|
|
target_name="search_index_check",
|
|
plugin_name="automation",
|
|
),
|
|
CronJobContribution(
|
|
name="check_workflow_timeouts",
|
|
cron_expression="*/5 * * * *",
|
|
job_type="custom",
|
|
target_name="check_workflow_timeouts",
|
|
plugin_name="automation",
|
|
),
|
|
],
|
|
author="LeoCRM Team",
|
|
min_app_version="1.0.0",
|
|
hooks=["contact.before_create", "contact.after_create"],
|
|
contract_version="1.0.0",
|
|
)
|
|
|
|
def __init__(self) -> None:
|
|
super().__init__()
|
|
# Track contributed definitions by source plugin name
|
|
self._contributed_agents: dict[str, list[str]] = {} # plugin_name -> [agent_name, ...]
|
|
self._contributed_automations: dict[str, list[str]] = {} # plugin_name -> [automation_name, ...]
|
|
self._contributed_cron_jobs: dict[str, list[str]] = {} # plugin_name -> [cron_job_name, ...]
|
|
self._contributed_heartbeats: dict[str, list[str]] = {} # plugin_name -> [agent_name, ...]
|
|
|
|
def get_entity_models(self) -> dict[str, type]:
|
|
from app.plugins.builtins.automation.models import AgentDefinition, AutomationDefinition
|
|
return {"agent_definition": AgentDefinition, "automation_definition": AutomationDefinition}
|
|
|
|
def get_job_modules(self) -> list[str]:
|
|
return [
|
|
"app.plugins.builtins.automation.scheduler",
|
|
"app.plugins.builtins.automation.workflow_timeout",
|
|
"app.plugins.builtins.automation.agent_runner",
|
|
"app.plugins.builtins.automation.execution_engine",
|
|
]
|
|
|
|
async def on_activate(self, db, service_container, event_bus) -> None:
|
|
"""Register event listeners on activation."""
|
|
await super().on_activate(db, service_container, event_bus)
|
|
# Register agent communication tool
|
|
try:
|
|
from app.plugins.builtins.automation.agent_comm import register_agent_comm_tool
|
|
register_agent_comm_tool()
|
|
except Exception:
|
|
logger.exception("Failed to register agent communication tool")
|
|
# Register agent coordinator tools
|
|
try:
|
|
from app.plugins.builtins.automation.agent_coordinator import (
|
|
register_agent_coordinator_tools,
|
|
)
|
|
register_agent_coordinator_tools()
|
|
except Exception:
|
|
logger.exception("Failed to register agent coordinator tools")
|
|
# Register workflow agent tools (I-AW: Agent→Workflow)
|
|
try:
|
|
self._register_workflow_agent_tools()
|
|
except Exception:
|
|
logger.exception("Failed to register workflow agent tools")
|
|
# Register MiniApps from manifest
|
|
try:
|
|
from app.plugins.builtins.kommunikation.contracts import get_miniapp_registry
|
|
registry = get_miniapp_registry()
|
|
for miniapp in self.manifest.miniapps:
|
|
registry.register(
|
|
app_id=miniapp.app_id,
|
|
name=miniapp.name,
|
|
icon=miniapp.icon,
|
|
description=miniapp.description,
|
|
plugin_name=self.manifest.name,
|
|
render_schema=miniapp.render_schema,
|
|
)
|
|
logger.info("Registered MiniApp '%s' from manifest", miniapp.app_id)
|
|
except Exception:
|
|
logger.exception("Failed to register MiniApps from manifest")
|
|
# Register own cron jobs from manifest
|
|
try:
|
|
await self.register_plugin_contributions(db, self.manifest.name, self.manifest)
|
|
logger.info("Registered own cron jobs from manifest")
|
|
except Exception:
|
|
logger.exception("Failed to register own cron jobs")
|
|
# Register pre-built agents in DB (if not already present)
|
|
try:
|
|
from app.plugins.builtins.automation.models import AgentDefinition
|
|
from app.plugins.builtins.automation.prebuilt.email_triage_agent import create_email_triage_agent
|
|
from app.plugins.builtins.automation.prebuilt.contact_enrichment_agent import create_contact_enrichment_agent
|
|
from app.plugins.builtins.automation.prebuilt.follow_up_agent import create_follow_up_agent
|
|
from app.plugins.builtins.automation.prebuilt.report_agent import create_report_agent
|
|
from sqlalchemy import select as sa_select
|
|
|
|
# Get system tenant + admin user for seeding (ARCH-043:
|
|
# deterministic slug lookup instead of arbitrary first row)
|
|
from app.core.db import get_system_tenant
|
|
from app.models.user import User
|
|
|
|
tenant = await get_system_tenant(db)
|
|
if tenant:
|
|
user_result = await db.execute(
|
|
sa_select(User)
|
|
.join(UserTenant, UserTenant.user_id == User.id)
|
|
.where(UserTenant.tenant_id == tenant.id)
|
|
.limit(1)
|
|
)
|
|
user = user_result.scalar_one_or_none()
|
|
if user:
|
|
prebuilt_factories = [
|
|
("E-Mail-Triage-Agent", create_email_triage_agent),
|
|
("Kontakt-Anreicherungs-Agent", create_contact_enrichment_agent),
|
|
("Follow-Up-Agent", create_follow_up_agent),
|
|
("Berichts-Agent", create_report_agent),
|
|
]
|
|
for agent_name, factory in prebuilt_factories:
|
|
# Check if agent already exists
|
|
existing = await db.execute(
|
|
sa_select(AgentDefinition).where(
|
|
AgentDefinition.tenant_id == tenant.id,
|
|
AgentDefinition.name == agent_name,
|
|
)
|
|
)
|
|
if not existing.scalar_one_or_none():
|
|
agent = factory(tenant_id=tenant.id, user_id=user.id)
|
|
db.add(agent)
|
|
logger.info("Registered pre-built agent '%s'", agent_name)
|
|
await db.commit()
|
|
logger.info("Pre-built agents registration complete")
|
|
except Exception:
|
|
logger.exception("Failed to register pre-built agents")
|
|
|
|
logger.info("Automation plugin activated")
|
|
|
|
def _register_workflow_agent_tools(self) -> None:
|
|
"""Register I-AW agent tools for starting and inspecting workflows."""
|
|
import uuid
|
|
from typing import Any
|
|
|
|
from app.ai.tool_registry import get_tool_registry
|
|
registry = get_tool_registry()
|
|
|
|
async def _start_workflow_handler(arguments: dict[str, Any], context: dict[str, Any]) -> dict[str, Any]:
|
|
"""Start a workflow by ID."""
|
|
from app.services.workflow_service import create_instance
|
|
from app.core.db import get_worker_session_factory
|
|
workflow_id = arguments.get("workflow_id", "")
|
|
tenant_id = context.get("tenant_id")
|
|
user_id = context.get("user_id")
|
|
if not workflow_id or not tenant_id:
|
|
return {"error": "workflow_id and tenant_id required"}
|
|
factory = get_worker_session_factory()
|
|
async with factory() as db:
|
|
instance = await create_instance(
|
|
db=db,
|
|
tenant_id=uuid.UUID(str(tenant_id)),
|
|
workflow_id=uuid.UUID(workflow_id),
|
|
initiated_by=uuid.UUID(str(user_id)) if user_id else None,
|
|
)
|
|
await db.commit()
|
|
return {"instance_id": str(instance.get("id", "")), "status": instance.get("status", "created")}
|
|
|
|
registry.register(
|
|
name="start_workflow",
|
|
description="Start a workflow by its ID. Returns the instance ID and status.",
|
|
parameters={
|
|
"type": "object",
|
|
"properties": {
|
|
"workflow_id": {"type": "string", "description": "UUID of the workflow to start"},
|
|
},
|
|
"required": ["workflow_id"],
|
|
},
|
|
handler=_start_workflow_handler,
|
|
plugin_name=self.manifest.name,
|
|
required_permission="workflows:read",
|
|
category="workflow",
|
|
)
|
|
|
|
async def _check_workflow_status_handler(arguments: dict[str, Any], context: dict[str, Any]) -> dict[str, Any]:
|
|
"""Check the status of a workflow instance."""
|
|
from sqlalchemy import select
|
|
from app.models.workflow import WorkflowInstance
|
|
from app.core.db import get_worker_session_factory
|
|
instance_id = arguments.get("instance_id", "")
|
|
tenant_id = context.get("tenant_id")
|
|
if not instance_id or not tenant_id:
|
|
return {"error": "instance_id and tenant_id required"}
|
|
factory = get_worker_session_factory()
|
|
async with factory() as db:
|
|
result = await db.execute(
|
|
select(WorkflowInstance).where(
|
|
WorkflowInstance.id == uuid.UUID(instance_id),
|
|
WorkflowInstance.tenant_id == uuid.UUID(str(tenant_id)),
|
|
)
|
|
)
|
|
inst = result.scalar_one_or_none()
|
|
if not inst:
|
|
return {"error": "Instance not found"}
|
|
return {
|
|
"instance_id": str(inst.id),
|
|
"status": inst.status,
|
|
"current_step": inst.current_step_index,
|
|
"completed_at": inst.completed_at.isoformat() if inst.completed_at else None,
|
|
}
|
|
|
|
registry.register(
|
|
name="check_workflow_status",
|
|
description="Check the status of a workflow instance by its ID.",
|
|
parameters={
|
|
"type": "object",
|
|
"properties": {
|
|
"instance_id": {"type": "string", "description": "UUID of the workflow instance"},
|
|
},
|
|
"required": ["instance_id"],
|
|
},
|
|
handler=_check_workflow_status_handler,
|
|
plugin_name=self.manifest.name,
|
|
required_permission="workflows:read",
|
|
category="workflow",
|
|
)
|
|
logger.info("Registered workflow agent tools: start_workflow, check_workflow_status")
|
|
|
|
async def on_deactivate(self, db, service_container, event_bus) -> None:
|
|
"""Clean up on deactivation."""
|
|
# Contract abmelden
|
|
from app.plugins.builtins.contracts import get_contract_registry
|
|
get_contract_registry().unregister(self.manifest.name)
|
|
|
|
await super().on_deactivate(db, service_container, event_bus)
|
|
# Unregister agent communication tool
|
|
try:
|
|
from app.plugins.builtins.automation.agent_comm import unregister_agent_comm_tool
|
|
unregister_agent_comm_tool()
|
|
except Exception:
|
|
logger.exception("Failed to unregister agent communication tool")
|
|
# Unregister agent coordinator tools
|
|
try:
|
|
from app.plugins.builtins.automation.agent_coordinator import (
|
|
unregister_agent_coordinator_tools,
|
|
)
|
|
unregister_agent_coordinator_tools()
|
|
except Exception:
|
|
logger.exception("Failed to unregister agent coordinator tools")
|
|
# Unregister workflow agent tools from the core AI tool registry
|
|
try:
|
|
from app.ai.tool_registry import get_tool_registry
|
|
get_tool_registry().unregister_plugin(self.manifest.name)
|
|
logger.info("Unregistered AI agent tools for plugin '%s'", self.manifest.name)
|
|
except Exception:
|
|
logger.exception("Failed to unregister AI agent tools")
|
|
# Unregister MiniApps
|
|
try:
|
|
from app.plugins.builtins.kommunikation.contracts import get_miniapp_registry
|
|
registry = get_miniapp_registry()
|
|
registry.unregister_plugin(self.manifest.name)
|
|
logger.info("Unregistered MiniApps for plugin '%s'", self.manifest.name)
|
|
except Exception:
|
|
logger.exception("Failed to unregister MiniApps")
|
|
logger.info("Automation plugin deactivated")
|
|
|
|
# ─── Plugin Contribution Registration ───
|
|
|
|
async def register_plugin_contributions(self, db, plugin_name: str, manifest) -> None:
|
|
"""Register agent definitions, automation templates, cron jobs, and heartbeat configs
|
|
from another plugin's manifest. Uses plugin name prefixing for conflict resolution."""
|
|
from sqlalchemy import select
|
|
|
|
# Get system tenant for contributions (ARCH-043: deterministic slug
|
|
# lookup instead of arbitrary first row)
|
|
from app.core.db import get_system_tenant
|
|
from app.plugins.builtins.automation.models import AutomationCronJob
|
|
from app.plugins.builtins.automation.services import (
|
|
AgentService,
|
|
AutomationService,
|
|
CronJobService,
|
|
)
|
|
tenant = await get_system_tenant(db)
|
|
default_tenant_id = tenant.id if tenant else None
|
|
if default_tenant_id is None:
|
|
logger.warning("No tenant found — skipping plugin contributions registration")
|
|
return
|
|
|
|
# Register agent definitions
|
|
agent_names: list[str] = []
|
|
for agent_def in manifest.agent_definitions:
|
|
prefixed_name = f"{plugin_name}.{agent_def.name}"
|
|
agent_names.append(prefixed_name)
|
|
# Check if already exists (idempotent)
|
|
existing = await AgentService.get_by_name(db, default_tenant_id, prefixed_name) if hasattr(AgentService, 'get_by_name') else None
|
|
if existing is None:
|
|
try:
|
|
await AgentService.create(db, default_tenant_id, {
|
|
"name": prefixed_name,
|
|
"description": agent_def.description,
|
|
"llm_model": agent_def.llm_model,
|
|
"system_prompt": agent_def.system_prompt,
|
|
"tool_ids": agent_def.tool_ids,
|
|
"heartbeat_interval_seconds": agent_def.heartbeat_interval_seconds,
|
|
"mode": agent_def.mode,
|
|
"max_executions_per_hour": agent_def.max_executions_per_hour,
|
|
"max_duration_seconds": agent_def.max_duration_seconds,
|
|
"budget_limit_usd": agent_def.budget_limit_usd,
|
|
"is_active": True,
|
|
})
|
|
logger.info("Registered contributed agent '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to register contributed agent '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
self._contributed_agents[plugin_name] = agent_names
|
|
|
|
# Register automation templates
|
|
automation_names: list[str] = []
|
|
for auto_def in manifest.automation_templates:
|
|
prefixed_name = f"{plugin_name}.{auto_def.name}"
|
|
automation_names.append(prefixed_name)
|
|
existing = await AutomationService.get_by_name(db, default_tenant_id, prefixed_name) if hasattr(AutomationService, 'get_by_name') else None
|
|
if existing is None:
|
|
try:
|
|
await AutomationService.create(db, default_tenant_id, {
|
|
"name": prefixed_name,
|
|
"description": auto_def.description,
|
|
"trigger_type": auto_def.trigger_type,
|
|
"trigger_config": auto_def.trigger_config,
|
|
"conditions": auto_def.conditions,
|
|
"actions": auto_def.actions,
|
|
"is_active": True,
|
|
})
|
|
logger.info("Registered contributed automation '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to register contributed automation '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
self._contributed_automations[plugin_name] = automation_names
|
|
|
|
# Register cron jobs (skip if no tenant exists yet)
|
|
cron_job_names: list[str] = []
|
|
if default_tenant_id is not None:
|
|
for cron_def in manifest.cron_jobs:
|
|
prefixed_name = f"{plugin_name}.{cron_def.name}"
|
|
cron_job_names.append(prefixed_name)
|
|
try:
|
|
# Check if cron job already exists
|
|
result = await db.execute(
|
|
select(AutomationCronJob).where(AutomationCronJob.name == prefixed_name).limit(1)
|
|
)
|
|
existing = result.scalar_one_or_none()
|
|
if existing is None:
|
|
await CronJobService.create(db, default_tenant_id, {
|
|
"name": prefixed_name,
|
|
"cron_expression": cron_def.cron_expression,
|
|
"job_type": cron_def.job_type,
|
|
"target_name": cron_def.target_name,
|
|
"plugin_name": plugin_name,
|
|
"is_active": True,
|
|
})
|
|
logger.info("Registered contributed cron job '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to register contributed cron job '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
self._contributed_cron_jobs[plugin_name] = cron_job_names
|
|
|
|
# Register heartbeat configs
|
|
heartbeat_names: list[str] = []
|
|
for hb_def in manifest.heartbeat_configs:
|
|
prefixed_name = f"{plugin_name}.{hb_def.agent_name}"
|
|
heartbeat_names.append(prefixed_name)
|
|
try:
|
|
# Create or update agent with heartbeat settings
|
|
agent = await AgentService.get_by_name(db, hb_def.tenant_id, prefixed_name) if hasattr(AgentService, 'get_by_name') else None
|
|
if agent is None:
|
|
await AgentService.create(db, hb_def.tenant_id, {
|
|
"name": prefixed_name,
|
|
"description": f"Heartbeat agent contributed by {plugin_name}",
|
|
"heartbeat_interval_seconds": hb_def.interval_seconds,
|
|
"mode": "proactive",
|
|
"is_active": True,
|
|
})
|
|
logger.info("Registered heartbeat agent '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to register heartbeat agent '%s' from plugin '%s'", prefixed_name, plugin_name)
|
|
self._contributed_heartbeats[plugin_name] = heartbeat_names
|
|
|
|
async def unregister_plugin_contributions(self, db, plugin_name: str) -> None:
|
|
"""Remove all contributed definitions from a plugin that is being deactivated."""
|
|
from app.plugins.builtins.automation.services import (
|
|
AgentService,
|
|
AutomationService,
|
|
CronJobService,
|
|
)
|
|
|
|
# Remove contributed agents
|
|
agent_names = self._contributed_agents.pop(plugin_name, [])
|
|
for agent_name in agent_names:
|
|
try:
|
|
agent = await AgentService.get_by_name(db, agent_name)
|
|
if agent:
|
|
await AgentService.delete(db, agent.tenant_id, agent.id)
|
|
logger.info("Unregistered contributed agent '%s' from plugin '%s'", agent_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to unregister contributed agent '%s'", agent_name)
|
|
|
|
# Remove contributed automations
|
|
automation_names = self._contributed_automations.pop(plugin_name, [])
|
|
for auto_name in automation_names:
|
|
try:
|
|
auto = await AutomationService.get_by_name(db, auto_name)
|
|
if auto:
|
|
await AutomationService.delete(db, auto.tenant_id, auto.id)
|
|
logger.info("Unregistered contributed automation '%s' from plugin '%s'", auto_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to unregister contributed automation '%s'", auto_name)
|
|
|
|
# Remove contributed cron jobs
|
|
cron_job_names = self._contributed_cron_jobs.pop(plugin_name, [])
|
|
for cron_name in cron_job_names:
|
|
try:
|
|
from sqlalchemy import select
|
|
|
|
from app.plugins.builtins.automation.models import AutomationCronJob
|
|
result = await db.execute(
|
|
select(AutomationCronJob).where(AutomationCronJob.name == cron_name).limit(1)
|
|
)
|
|
job = result.scalar_one_or_none()
|
|
if job:
|
|
await CronJobService.delete(db, job.tenant_id, job.id)
|
|
logger.info("Unregistered contributed cron job '%s' from plugin '%s'", cron_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to unregister contributed cron job '%s'", cron_name)
|
|
|
|
# Remove contributed heartbeats
|
|
heartbeat_names = self._contributed_heartbeats.pop(plugin_name, [])
|
|
for hb_name in heartbeat_names:
|
|
try:
|
|
agent = await AgentService.get_by_name(db, hb_name)
|
|
if agent:
|
|
await AgentService.delete(db, agent.tenant_id, agent.id)
|
|
logger.info("Unregistered heartbeat agent '%s' from plugin '%s'", hb_name, plugin_name)
|
|
except Exception:
|
|
logger.exception("Failed to unregister heartbeat agent '%s'", hb_name)
|
|
|
|
# ─── Heartbeat Migration ───
|
|
|
|
async def ensure_ai_proactive_heartbeat(self, db) -> None:
|
|
"""Migrate the hardcoded ai_proactive heartbeat to a configurable cron job."""
|
|
from sqlalchemy import select
|
|
|
|
from app.plugins.builtins.automation.models import AutomationCronJob
|
|
from app.plugins.builtins.automation.services import CronJobService
|
|
|
|
# Check if ai_proactive heartbeat cron job already exists
|
|
result = await db.execute(
|
|
select(AutomationCronJob).where(
|
|
AutomationCronJob.name == "ai_proactive.heartbeat"
|
|
).limit(1)
|
|
)
|
|
existing = result.scalar_one_or_none()
|
|
if existing is not None:
|
|
logger.info("ai_proactive heartbeat cron job already exists, skipping migration")
|
|
return
|
|
|
|
# Create the heartbeat cron job
|
|
try:
|
|
await CronJobService.create(db, None, {
|
|
"name": "ai_proactive.heartbeat",
|
|
"cron_expression": "*/5 * * * *", # Every 5 minutes
|
|
"job_type": "agent_heartbeat",
|
|
"target_name": "ai_proactive.heartbeat_agent",
|
|
"plugin_name": "ai_proactive",
|
|
"is_active": True,
|
|
})
|
|
logger.info("Created ai_proactive heartbeat cron job (every 5 minutes)")
|
|
except Exception:
|
|
logger.exception("Failed to create ai_proactive heartbeat cron job")
|
|
|
|
# ─── Event Handlers ───
|
|
|
|
async def on_contact_created(self, payload: dict[str, Any]) -> None:
|
|
"""Handle contact.created event — trigger matching automations."""
|
|
logger.debug("contact.created event received: %s", payload)
|
|
|
|
async def on_contact_updated(self, payload: dict[str, Any]) -> None:
|
|
"""Handle contact.updated event — trigger matching automations."""
|
|
logger.debug("contact.updated event received: %s", payload)
|
|
|
|
async def on_mail_received(self, payload: dict[str, Any]) -> None:
|
|
"""Handle mail.received event — trigger matching automations."""
|
|
logger.debug("mail.received event received: %s", payload)
|
|
|
|
async def on_workflow_timeout(self, payload: dict[str, Any]) -> None:
|
|
"""Handle workflow.timeout event — trigger matching automations."""
|
|
logger.debug("workflow.timeout event received: %s", payload)
|