Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line numberDiff line numberDiff line change
Expand Up@@ -54,3 +54,6 @@ dmypy.json
# OS
.DS_Store
Thumbs.db

# Secrets
.env
2 changes: 1 addition & 1 deletion README.md
Original file line numberDiff line numberDiff line change
Expand Up@@ -25,7 +25,7 @@ import asyncio, os
import llpsdk as llp

# Define a callback handler for processing messages
async def on_message(msg):
async def on_message(annotater, msg):
# Process the prompt with your agent.
# Replace this with your own processing logic.
response = msg.prompt
Expand Down
1 change: 1 addition & 0 deletions examples/.env
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
LLP_URL="ws://localhost:4000/agent/websocket"
LLP_API_KEY=
20 changes: 11 additions & 9 deletions examples/simple_agent.py
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
"""Agent example demonstrating basic LLP SDK usage."""
import asyncio
from datetime import timedelta
import llpsdk as llp
import os
from dotenv import load_dotenv
Expand All@@ -9,16 +10,22 @@ async def main() -> None:
"""Run a simple agent that connects, sends presence, and sends a message."""
load_dotenv()
platform_url = os.getenv("LLP_URL")
api_key = os.getenv("LLP_API_KEY")

if platform_url is None:
raise Exception("LLP_URL env var is not defined")

cfg = llp.Config()
if api_key is None:
raise Exception("LLP_API_KEY env var is not defined")

cfg = llp.Config(platform_url=platform_url)
cfg.platform_url = platform_url
client = llp.Client("simple-agent", "testkey", cfg)
client = llp.Client("simple-agent", api_key, cfg)

# Set up handlers
async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
print("Feed msg.prompt into your agent and return the response.")
async def on_message(annotater: llp.Annotater, msg: llp.TextMessage) -> llp.TextMessage:
tc = msg.tool_call("get_weather", '{"city":"Seattle"}', "rainy", timedelta(seconds=1))
await annotater.annotate_tool_call(tc)
return msg.reply("this is my response")

# Register handlers
Expand All@@ -30,11 +37,6 @@ async def on_message(msg: llp.TextMessage) -> llp.TextMessage:
await client.connect()
print(f"Connected! Session ID: {client.session_id}")

# Send a message
msg = llp.TextMessage(recipient="echo-agent", prompt="Hello from Python!")
print(f"Sending message to {msg.recipient}...")
await client.send_message(msg)

# Keep running
print("Agent running. Press Ctrl+C to exit...")
await asyncio.Event().wait()
Expand Down
4 changes: 4 additions & 0 deletions src/llpsdk/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -11,12 +11,14 @@
PlatformError,
TimeoutError,
)
from .handler import Annotater
from .message import (
AuthenticatedResponse,
PresenceMessage,
TextMessage,
)
from .presence import ConnectionStatus, PresenceStatus
from .tool_call import ToolCall

__all__ = [
"Client",
Expand All@@ -33,6 +35,8 @@
"AlreadyClosedError",
"TimeoutError",
"InvalidStatusError",
"Annotater",
"ToolCall",
]

__version__ = "0.1.0"
72 changes: 12 additions & 60 deletions src/llpsdk/client.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -23,6 +23,7 @@
PresenceMessage,
TextMessage,
)
from .tool_call import ToolCall
from .presence import ConnectionStatus, PresenceStatus


Expand DownExpand Up@@ -64,10 +65,6 @@ def __init__(self, name: str, api_key: str, config: Optional[Config] = None) ->
self._write_task: Optional[asyncio.Task[None]] = None
self._stop_event = asyncio.Event()

# Pending messages (for request/response)
self._pending_lock = asyncio.Lock()
self._pending: Dict[str, asyncio.Future[TextMessage]] = {}

# Auth future (for waiting on authentication)
self._auth_future: Optional[asyncio.Future[AuthenticatedResponse]] = None

Expand DownExpand Up@@ -170,7 +167,7 @@ async def close(self) -> None:
async with self._presence_lock:
self._presence = PresenceStatus.unavailable

async def send_async_message(self, message: TextMessage) -> None:
async def _send_async_message(self, message: TextMessage) -> None:
"""
Send a message asynchronously (fire-and-forget).

Expand All@@ -187,50 +184,22 @@ async def send_async_message(self, message: TextMessage) -> None:

await self._send(message.encode())

async def send_message(self, message: TextMessage, timeout: float = 10.0) -> TextMessage:
async def annotate_tool_call(self, tool_call: ToolCall) -> None:
"""
Send a message and wait for response.
Send a tool call annotation to the platform for telemetry.

Args:
message: Message to send
timeout: Response timeout in seconds

Returns:
Response message
tool_call: The tool call to annotate (created via TextMessage.tool_call() or
TextMessage.tool_call_exception())

Raises:
ValueError: If message ID is empty
NotAuthenticatedError: If not authenticated
TimeoutError: If no response within timeout
"""
# ID is REQUIRED for synchronous send
if not message._id:
raise ValueError("Message ID is required for send_message()")

async with self._status_lock:
if self._status != ConnectionStatus.AUTHENTICATED:
raise NotAuthenticatedError("Must connect before sending messages")
raise NotAuthenticatedError("Must connect before annotating tool calls")

# Create future for this message
response_future: asyncio.Future[TextMessage] = asyncio.get_event_loop().create_future()

async with self._pending_lock:
self._pending[message._id] = response_future

try:
# Send message asynchronously
await self.send_async_message(message)

# Wait for response with timeout
response = await asyncio.wait_for(response_future, timeout=timeout)
return response

except asyncio.TimeoutError:
raise TimeoutError(f"No response within {timeout}s")
finally:
# Clean up
async with self._pending_lock:
self._pending.pop(message._id, None)
await self._send(tool_call.encode())

# Properties

Expand DownExpand Up@@ -415,15 +384,9 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
self._auth_future.set_exception(error)
return

# Check if this error is for a pending message
if error.id:
async with self._pending_lock:
if error.id in self._pending:
future = self._pending.pop(error.id)
if not future.done():
future.set_exception(error)
return
return

if msg_type == "ack":
return

if msg_type == "authenticated":
Expand All@@ -438,21 +401,10 @@ async def _handle_message(self, msg_dict: Dict[str, Any]) -> None:
return

if msg_type == "message":
msg_id = msg_dict.get("id", "")

async with self._pending_lock:
if msg_id in self._pending:
future = self._pending[msg_id]
if not future.done():
tm = TextMessage.decode(msg_dict)
future.set_result(tm)
return

# Not a response, call message handler
tm = TextMessage.decode(msg_dict)
reply = await self._handlers.call_message(tm)
reply = await self._handlers.call_message(self, tm)
if reply is not None:
await self.send_async_message(reply)
await self._send_async_message(reply)
return

async def _handle_disconnect(self) -> None:
Expand Down
21 changes: 17 additions & 4 deletions src/llpsdk/handler.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -3,13 +3,24 @@
import asyncio
from typing import Awaitable, Callable, Optional, Union

from typing import Protocol, runtime_checkable

from .message import PresenceMessage, TextMessage
from .tool_call import ToolCall


@runtime_checkable
class Annotater(Protocol):
"""Protocol for annotating tool calls for telemetry."""

async def annotate_tool_call(self, tool_call: ToolCall) -> None: ...


# Handler type signatures (supports both sync and async)
PresenceHandler = Union[
Callable[[PresenceMessage], None], Callable[[PresenceMessage], Awaitable[None]]
]
MessageHandler = Callable[[TextMessage], Awaitable[TextMessage]]
MessageHandler = Callable[["Annotater", TextMessage], Awaitable[TextMessage]]


class HandlerRegistry:
Expand All@@ -36,9 +47,11 @@ async def call_presence(self, update: PresenceMessage) -> None:
else:
self._on_presence(update)

async def call_message(self, message: TextMessage) -> Optional[TextMessage]:
"""Call the message handler if set."""
async def call_message(
self, annotater: Annotater, message: TextMessage
) -> Optional[TextMessage]:
"""Call the message handler if set, passing annotater for tool call telemetry."""
if self._on_message is not None:
result = await self._on_message(message)
result = await self._on_message(annotater, message)
return result
return None
28 changes: 28 additions & 0 deletions src/llpsdk/message.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -4,10 +4,12 @@
import uuid
import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Any, Dict, Optional

from llpsdk.errors import TextMessageEmptyError
from llpsdk.presence import PresenceStatus
from llpsdk.tool_call import ToolCall


class TextMessage:
Expand DownExpand Up@@ -57,6 +59,32 @@ def encode(self) -> str:
def has_attachment(self) -> bool:
return self.attachment != ""

def tool_call(self, name: str, parameters: str, result: str, duration: timedelta) -> ToolCall:
"""Create a successful ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=result,
threw_exception=False,
duration=duration,
)

def tool_call_exception(
self, name: str, parameters: str, error: Exception, duration: timedelta
) -> ToolCall:
"""Create a failed ToolCall annotation from this message."""
return ToolCall(
id=self._id,
recipient=self.sender,
name=name,
parameters=parameters,
result=str(error),
threw_exception=True,
duration=duration,
)

@staticmethod
def decode(msg: Dict[str, Any]) -> "TextMessage":
"""Decodes a JSON dict into a TextMessage object"""
Expand Down
35 changes: 35 additions & 0 deletions src/llpsdk/tool_call.py
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
"""ToolCall message type for telemetry annotation."""

import json
from dataclasses import dataclass
from datetime import timedelta
from typing import Optional


@dataclass
class ToolCall:
"""Tool call annotation sent to the platform for telemetry."""

id: Optional[str]
recipient: str
name: str
parameters: str
result: str
threw_exception: bool
duration: timedelta

def encode(self) -> str:
"""Encode a ToolCall into serialized JSON."""
data = {
"type": "tool_call",
"id": self.id,
"data": {
"to": self.recipient,
"name": self.name,
"parameters": self.parameters,
"result": self.result,
"threw_exception": self.threw_exception,
"duration_ms": int(self.duration.total_seconds() * 1000),
},
}
return json.dumps(data)
Loading