Files
dify-plugin-daemon/cmd/commandline/plugin/templates/python/trigger_provider.py
非法操作 70ba852ecc feat: cli tool support create trigger plugins (#485)
* feat: cli tool support create trigger plugins

* fix(trigger): update placeholder comment in SubscriptionConstructor for webhook registration

---------

Co-authored-by: Harry <xh001x@hotmail.com>
2025-11-05 15:45:58 +08:00

101 lines
3.7 KiB
Python

import time
from collections.abc import Mapping
from typing import Any
from werkzeug import Request, Response
from dify_plugin.entities.oauth import TriggerOAuthCredentials
from dify_plugin.entities.provider_config import CredentialType
from dify_plugin.entities.trigger import EventDispatch, Subscription, UnsubscribeResult
from dify_plugin.errors.trigger import (
SubscriptionError,
TriggerDispatchError,
TriggerProviderCredentialValidationError,
TriggerProviderOAuthError,
UnsubscribeError,
)
from dify_plugin.interfaces.trigger import Trigger, TriggerSubscriptionConstructor
class {{ .PluginName | SnakeToCamel }}Trigger(Trigger):
"""
Handle the webhook event dispatch.
"""
def _dispatch_event(self, subscription: Subscription, request: Request) -> EventDispatch:
payload: Mapping[str, Any] = self._validate_payload(request)
response = Response(response='{"status": "ok"}', status=200, mimetype="application/json")
events: list[str] = self._dispatch_trigger_events(payload=payload)
return EventDispatch(events=events, response=response)
def _dispatch_trigger_events(self, payload: Mapping[str, Any]) -> list[str]:
"""Dispatch events based on webhook payload."""
events = []
# Get the event type from the payload
event_type = payload.get("type", "")
if event_type.startswith("my-event-type"):
events.append("{{ .PluginName }}_event")
return events
def _validate_payload(self, request: Request) -> Mapping[str, Any]:
try:
payload = request.get_json(force=True)
if not payload:
raise TriggerDispatchError("Empty request body")
return payload
except TriggerDispatchError:
raise
except Exception as exc:
raise TriggerDispatchError(f"Failed to parse payload: {exc}") from exc
class {{ .PluginName | SnakeToCamel }}SubscriptionConstructor(TriggerSubscriptionConstructor):
"""Manage {{ .PluginName }} trigger subscriptions."""
def _validate_api_key(self, credentials: dict[str, Any]) -> None:
api_key = credentials.get("api_key")
if not api_key:
raise TriggerProviderCredentialValidationError("API key is required to validate credentials.")
def _create_subscription(
self,
endpoint: str,
parameters: Mapping[str, Any],
credentials: Mapping[str, Any],
credential_type: CredentialType,
) -> Subscription:
events: list[str] = parameters.get("events", [])
# Replace this placeholder with API calls to register a webhook
return Subscription(
expires_at=int(time.time()) + 7 * 24 * 60 * 60,
endpoint=endpoint,
properties={
"external_id": "example-subscription",
"events": events,
},
)
def _delete_subscription(
self,
subscription: Subscription,
credentials: Mapping[str, Any],
credential_type: CredentialType,
) -> UnsubscribeResult:
# Tear down any remote subscription that was created in `_subscribe`.
return UnsubscribeResult(success=True, message="Subscription removed.")
def _refresh_subscription(
self,
subscription: Subscription,
credentials: Mapping[str, Any],
credential_type: CredentialType,
) -> Subscription:
# Extend the subscription lifetime or renew tokens with your upstream service.
return Subscription(
expires_at=int(time.time()) + 7 * 24 * 60 * 60,
endpoint=subscription.endpoint,
properties=subscription.properties,
)