Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 70 additions & 0 deletions sdk/python/examples/09e_signals.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
# Copyright (c) 2025 Agentspan
# Licensed under the MIT License. See LICENSE file in the project root for details.

"""Signals — Send context to a running agent mid-execution.

Demonstrates:
- Starting a long-running agent
- Sending a normal signal to redirect it
- Sending an urgent signal for faster pickup
- Polling for signal disposition

Requirements:
- AGENTSPAN_SERVER_URL=http://localhost:6767/api
- AGENTSPAN_LLM_MODEL=openai/gpt-4o-mini
- OPENAI_API_KEY set
"""

import time
from agentspan.agents import Agent, AgentRuntime
from settings import settings

agent = Agent(
name="researcher",
model=settings.llm_model,
instructions="You are a research assistant. Research topics thoroughly before answering.",
signal_mode="evaluate", # LLM will accept or reject signals
)


if __name__ == "__main__":
with AgentRuntime() as runtime:
# Deploy to server. CLI alternative (recommended for CI/CD):
# agentspan deploy examples.09e_signals
runtime.deploy(agent)

# Start a long-running research task
handle = runtime.start(agent, "Research the history of quantum computing in detail.")
print(f"Started: {handle.execution_id}")

# Wait a moment for the agent to start working
time.sleep(3)

# Send a normal signal to redirect focus
receipt = handle.signal(
message="Focus only on developments after 2015. Skip early history.",
priority="normal",
sender="user",
)
print(f"Signal queued: {receipt.signal_id}")

# Send an urgent signal (picked up sooner — after current task, not full loop)
urgent_receipt = handle.signal(
message="Budget constraint: wrap up in the next 2 paragraphs.",
priority="urgent",
sender="manager",
)
print(f"Urgent signal queued: {urgent_receipt.signal_id}")

# Stream events to see signals being accepted/rejected
for event in handle.stream():
print(f" [{event.type}] {event.data}")
if event.type in ("signal_accepted", "signal_rejected"):
print(f" → Signal {(event.data or {}).get('signalId')} was {event.type}")

result = runtime.get_result(handle.execution_id)
result.print_result()

# Check final signal disposition
status = runtime.get_signal_status(receipt.signal_id)
print(f"Signal disposition: {status.disposition}")
9 changes: 9 additions & 0 deletions sdk/python/src/agentspan/agents/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,9 @@ def get_weather(city: str) -> str:
# Memory
from agentspan.agents.memory import ConversationMemory

# Signal types
from agentspan.agents.signal import SignalReceipt, SignalStatus

# Result types
from agentspan.agents.result import (
AgentEvent,
Expand Down Expand Up @@ -191,6 +194,7 @@ def resolve_credentials(input_data: dict, names: list) -> dict:
mcp_tool,
pdf_tool,
search_tool,
signal_tool,
tool,
video_tool,
)
Expand Down Expand Up @@ -221,6 +225,8 @@ def resolve_credentials(input_data: dict, names: list) -> dict:
"http_tool",
"human_tool",
"mcp_tool",
"signal_tool",
"wait_for_message_tool",
"image_tool",
"audio_tool",
"video_tool",
Expand Down Expand Up @@ -312,4 +318,7 @@ def resolve_credentials(input_data: dict, names: list) -> dict:
"skill",
"load_skills",
"SkillLoadError",
# Signal types
"SignalReceipt",
"SignalStatus",
]
2 changes: 2 additions & 0 deletions sdk/python/src/agentspan/agents/agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -360,6 +360,7 @@ def __init__(
required_tools: Optional[List[str]] = None,
gate: Optional[Any] = None,
credentials: Optional[List[Any]] = None,
signal_mode: str = "evaluate",
) -> None:
if not name or not isinstance(name, str):
raise ValueError("Agent name must be a non-empty string")
Expand Down Expand Up @@ -455,6 +456,7 @@ def __init__(
self.thinking_budget_tokens = thinking_budget_tokens
self.required_tools: List[str] = list(required_tools) if required_tools else []
self.gate = gate
self.signal_mode = signal_mode
# ── Code execution setup ─────────────────────────────────────
self.code_execution_config: Optional[Any] = None
if code_execution is not None:
Expand Down
1 change: 1 addition & 0 deletions sdk/python/src/agentspan/agents/config_serializer.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ def _serialize_agent(self, agent: "Agent") -> dict:
"maxTurns": agent.max_turns,
"timeoutSeconds": agent.timeout_seconds,
"external": agent.external,
"signalMode": getattr(agent, "signal_mode", "evaluate"),
}

# Instructions
Expand Down
27 changes: 27 additions & 0 deletions sdk/python/src/agentspan/agents/result.py
Original file line number Diff line number Diff line change
Expand Up @@ -263,6 +263,25 @@ def send(self, message: str) -> None:
"""Send a message to a waiting agent (multi-turn conversation)."""
self.respond({"message": message})

def signal(self, message: str, *, priority: str = "normal",
data: dict | None = None, sender: str | None = None,
propagate: bool = True) -> "SignalReceipt":
"""Send a signal to this running execution."""
from agentspan.agents.signal import SignalReceipt # noqa: F401
return self._runtime.signal(
execution_id=self.execution_id,
message=message,
priority=priority,
data=data,
sender=sender,
propagate=propagate,
)

async def signal_async(self, message: str, **kwargs) -> "SignalReceipt":
"""Async version of :meth:`signal`."""
return await self._runtime.signal_async(
execution_id=self.execution_id, message=message, **kwargs)

# ── Execution control ───────────────────────────────────────────

def pause(self) -> None:
Expand Down Expand Up @@ -350,6 +369,9 @@ class EventType(str, Enum):
DONE = "done"
GUARDRAIL_PASS = "guardrail_pass"
GUARDRAIL_FAIL = "guardrail_fail"
SIGNAL_RECEIVED = "signal_received"
SIGNAL_ACCEPTED = "signal_accepted"
SIGNAL_REJECTED = "signal_rejected"


@dataclass
Expand Down Expand Up @@ -381,6 +403,7 @@ class AgentEvent:
output: Any = None
execution_id: str = ""
guardrail_name: Optional[str] = None
data: Optional[Dict[str, Any]] = None

def __post_init__(self):
if self.args and isinstance(self.args, dict):
Expand Down Expand Up @@ -519,6 +542,10 @@ def send(self, message: str) -> None:
"""Send a message to a waiting agent (multi-turn conversation)."""
self.handle.send(message)

def signal(self, message: str, **kwargs) -> "SignalReceipt":
"""Send a signal to the running agent this stream is attached to."""
return self.handle.signal(message, **kwargs)

@property
def execution_id(self) -> str:
"""The Conductor execution ID."""
Expand Down
152 changes: 152 additions & 0 deletions sdk/python/src/agentspan/agents/runtime/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -3209,6 +3209,27 @@ def _sse_to_agent_event(sse_event: Dict[str, Any], execution_id: str) -> Optiona
if event_type is None:
return None

# Signal events carry structured data in the ``data`` field.
signal_data: Optional[Dict[str, Any]] = None
if event_type == "signal_received":
signal_data = {
"signalId": data.get("signalId"),
"message": data.get("content"),
"sender": data.get("sender"),
"priority": data.get("priority"),
}
elif event_type == "signal_accepted":
signal_data = {
"signalId": data.get("signalId"),
"sender": data.get("sender"),
}
elif event_type == "signal_rejected":
signal_data = {
"signalId": data.get("signalId"),
"sender": data.get("sender"),
"rejectionReason": data.get("rejectionReason"),
}

return AgentEvent(
type=event_type,
content=data.get("content"),
Expand All @@ -3219,6 +3240,7 @@ def _sse_to_agent_event(sse_event: Dict[str, Any], execution_id: str) -> Optiona
output=data.get("output"),
execution_id=data.get("executionId", execution_id),
guardrail_name=data.get("guardrailName"),
data=signal_data,
)

# ── Fire-and-forget execution ───────────────────────────────────
Expand Down Expand Up @@ -4368,6 +4390,136 @@ async def cancel_async(self, execution_id: str, reason: str = "") -> None:
),
)

# ── Signal methods ───────────────────────────────────────────────

def signal(
self,
*,
execution_id: Optional[str] = None,
agent_name: Optional[str] = None,
message: str,
priority: str = "normal",
data: Optional[Dict[str, Any]] = None,
sender: Optional[str] = None,
propagate: bool = True,
correlation_id: Optional[str] = None,
) -> "SignalReceipt":
"""Send a signal to a running agent execution.

Args:
execution_id: Target execution ID (UUID). Use this or agent_name.
agent_name: Target agent name — signals all matching RUNNING executions.
message: Natural language message (max 4096 chars).
priority: "normal" (default) or "urgent".
data: Optional structured payload dict.
sender: Optional attribution string (e.g. your agent name).
propagate: Whether to propagate to active sub-workflows (default True).
correlation_id: Optional correlation ID for filtering by agent_name.

Returns:
:class:`SignalReceipt` confirming the signal was queued.
"""
import requests as req_lib

from agentspan.agents.signal import SignalReceipt

payload: Dict[str, Any] = {
"message": message,
"priority": priority,
"propagate": propagate,
}
if data is not None:
payload["data"] = data
if sender is not None:
payload["sender"] = sender

if execution_id is not None:
url = self._agent_api_url(f"/{execution_id}/signal")
elif agent_name is not None:
params = f"?agentName={agent_name}"
if correlation_id:
params += f"&correlationId={correlation_id}"
url = self._agent_api_url(f"/signal{params}")
else:
raise ValueError("Either execution_id or agent_name must be provided")

resp = req_lib.post(url, json=payload, headers=self._agent_api_headers(), timeout=30)
try:
resp.raise_for_status()
except req_lib.exceptions.HTTPError as exc:
_raise_api_error(exc, url=url)
resp_data = resp.json()
return SignalReceipt(
signal_id=resp_data["signalId"],
execution_id=resp_data.get("executionId", execution_id or ""),
status=resp_data.get("status", "queued"),
)

def broadcast(
self,
*,
execution_ids: List[str],
message: str,
priority: str = "normal",
**kwargs: Any,
) -> "List[SignalReceipt]":
"""Send the same signal to multiple executions.

Returns a list of :class:`SignalReceipt` objects (one per execution).
"""
return [
self.signal(execution_id=wf_id, message=message, priority=priority, **kwargs)
for wf_id in execution_ids
]

def get_signal_status(self, signal_id: str) -> "SignalStatus":
"""Poll for signal disposition.

Args:
signal_id: The signal ID from :class:`SignalReceipt`.

Returns:
:class:`SignalStatus` reflecting current disposition.
"""
import requests as req_lib

from agentspan.agents.signal import SignalStatus

url = self._agent_api_url(f"/signal/{signal_id}/status")
resp = req_lib.get(url, headers=self._agent_api_headers(content_type=""), timeout=30)
try:
resp.raise_for_status()
except req_lib.exceptions.HTTPError as exc:
_raise_api_error(exc, url=url)
resp_data = resp.json()
return SignalStatus(
signal_id=resp_data["signalId"],
execution_id=resp_data["executionId"],
delivered=resp_data.get("delivered", False),
disposition=resp_data.get("disposition", "pending"),
rejection_reason=resp_data.get("rejectionReason"),
)

async def signal_async(
self,
*,
execution_id: Optional[str] = None,
agent_name: Optional[str] = None,
message: str,
**kwargs: Any,
) -> "SignalReceipt":
"""Async version of :meth:`signal`."""
loop = asyncio.get_event_loop()
return await loop.run_in_executor(
None,
lambda: self.signal(
execution_id=execution_id,
agent_name=agent_name,
message=message,
**kwargs,
),
)

# ── Session continuity helpers ────────────────────────────────────

def _get_session_messages(self, session_id: str, agent_name: str) -> List[Dict[str, Any]]:
Expand Down
24 changes: 24 additions & 0 deletions sdk/python/src/agentspan/agents/signal.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
# Copyright (c) 2025 Agentspan
# Licensed under the MIT License. See LICENSE file in the project root for details.

from __future__ import annotations
from dataclasses import dataclass
from typing import Optional


@dataclass
class SignalReceipt:
"""Returned immediately when a signal is sent — confirms it was queued."""
signal_id: str
execution_id: str
status: str # Always "queued" at send time


@dataclass
class SignalStatus:
"""Current disposition of a signal — returned by get_signal_status()."""
signal_id: str
execution_id: str
delivered: bool
disposition: str # "pending" | "accepted" | "rejected" | "accepted_implicit"
rejection_reason: Optional[str] = None
Loading
Loading