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
4 changes: 2 additions & 2 deletions pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,11 +1,11 @@
[project]
name = "uipath-runtime"
version = "0.12.7"
version = "0.13.0"
description = "Runtime abstractions and interfaces for building agents and automation scripts in the UiPath ecosystem"
readme = { file = "README.md", content-type = "text/markdown" }
requires-python = ">=3.11"
dependencies = [
"uipath-core>=0.5.28, <0.6.0",
"uipath-core>=0.5.31, <0.6.0",
"vaderSentiment>=3.3.2, <4.0",
"chardet>=5.2.0, <8.0",
]
Expand Down
5 changes: 4 additions & 1 deletion src/uipath/runtime/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
UiPathStreamNotSupportedError,
UiPathStreamOptions,
)
from uipath.runtime.chat.protocol import UiPathChatProtocol
from uipath.runtime.chat.protocol import UiPathChatMetaEventProtocol, UiPathChatProtocol
from uipath.runtime.chat.runtime import UiPathChatRuntime
from uipath.runtime.context import UiPathRuntimeContext
from uipath.runtime.debug.breakpoint import UiPathBreakpointResult
Expand Down Expand Up @@ -45,6 +45,7 @@
from uipath.runtime.storage import UiPathRuntimeStorageProtocol
from uipath.runtime.workspace import (
AttachmentRegistryEntry,
ConversationalWorkspaceRuntime,
HydrationPolicy,
HydrationRuntime,
Workspace,
Expand Down Expand Up @@ -80,9 +81,11 @@
"UiPathBreakpointResult",
"UiPathStreamNotSupportedError",
"UiPathResumeTriggerName",
"UiPathChatMetaEventProtocol",
"UiPathChatProtocol",
"UiPathChatRuntime",
"AttachmentRegistryEntry",
"ConversationalWorkspaceRuntime",
"HydrationPolicy",
"HydrationRuntime",
"get_workspace_path",
Expand Down
4 changes: 2 additions & 2 deletions src/uipath/runtime/chat/__init__.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
"""Chat bridge protocol and runtime for conversational agents."""

from uipath.runtime.chat.protocol import UiPathChatProtocol
from uipath.runtime.chat.protocol import UiPathChatMetaEventProtocol, UiPathChatProtocol
from uipath.runtime.chat.runtime import UiPathChatRuntime

__all__ = ["UiPathChatProtocol", "UiPathChatRuntime"]
__all__ = ["UiPathChatMetaEventProtocol", "UiPathChatProtocol", "UiPathChatRuntime"]
11 changes: 10 additions & 1 deletion src/uipath/runtime/chat/protocol.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
"""Abstract conversation bridge interface."""

from typing import Any, Protocol
from typing import Any, Protocol, runtime_checkable

from uipath.core.chat import (
UiPathConversationMessageEvent,
Expand Down Expand Up @@ -70,3 +70,12 @@ async def emit_exchange_error_event(self, error: Exception) -> None:
async def wait_for_resume(self) -> dict[str, Any]:
"""Wait for the interrupt_end event to be received."""
...


@runtime_checkable
class UiPathChatMetaEventProtocol(Protocol):
"""Optional chat-bridge capability for conversation metadata events."""

async def emit_meta_event(self, meta_event: dict[str, Any]) -> None:
"""Send an exchange-scoped conversation metadata event."""
...
10 changes: 9 additions & 1 deletion src/uipath/runtime/chat/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,11 @@
UiPathRuntimeProtocol,
UiPathStreamOptions,
)
from uipath.runtime.chat.protocol import UiPathChatProtocol
from uipath.runtime.chat.protocol import UiPathChatMetaEventProtocol, UiPathChatProtocol
from uipath.runtime.errors import UiPathBaseRuntimeError
from uipath.runtime.errors.contract import UiPathErrorCategory
from uipath.runtime.events import (
UiPathRuntimeConversationMetaEvent,
UiPathRuntimeEvent,
UiPathRuntimeMessageEvent,
)
Expand Down Expand Up @@ -86,6 +87,13 @@ async def stream(
if isinstance(event, UiPathRuntimeMessageEvent):
if event.payload:
await self.chat_bridge.emit_message_event(event.payload)
elif isinstance(event, UiPathRuntimeConversationMetaEvent):
if isinstance(self.chat_bridge, UiPathChatMetaEventProtocol):
await self.chat_bridge.emit_meta_event(event.payload)
else:
logger.warning(
"Chat bridge does not support conversation metadata events"
)

if isinstance(event, UiPathRuntimeResult):
runtime_result = event
Expand Down
2 changes: 2 additions & 0 deletions src/uipath/runtime/events/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
UiPathRuntimeStatePhase,
)
from uipath.runtime.events.state import (
UiPathRuntimeConversationMetaEvent,
UiPathRuntimeMessageEvent,
UiPathRuntimeStateEvent,
)
Expand All @@ -26,4 +27,5 @@
# Runtime events
"UiPathRuntimeStateEvent",
"UiPathRuntimeMessageEvent",
"UiPathRuntimeConversationMetaEvent",
]
1 change: 1 addition & 0 deletions src/uipath/runtime/events/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ class UiPathRuntimeEventType(str, Enum):
"""Types of events that can be emitted during execution."""

RUNTIME_MESSAGE = "runtime_message"
CONVERSATION_META = "conversation_meta"
RUNTIME_STATE = "runtime_state"
RUNTIME_ERROR = "runtime_error"
RUNTIME_RESULT = "runtime_result"
Expand Down
15 changes: 14 additions & 1 deletion src/uipath/runtime/events/state.py
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,15 @@ class UiPathRuntimeMessageEvent(UiPathRuntimeEvent):
)


class UiPathRuntimeConversationMetaEvent(UiPathRuntimeEvent):
"""Conversation metadata emitted by a runtime."""

payload: dict[str, Any]
event_type: UiPathRuntimeEventType = Field(
default=UiPathRuntimeEventType.CONVERSATION_META, frozen=True
)


class UiPathRuntimeStateEvent(UiPathRuntimeEvent):
"""Event emitted when agent state is updated.

Expand Down Expand Up @@ -82,4 +91,8 @@ class UiPathRuntimeStateEvent(UiPathRuntimeEvent):
)


__all__ = ["UiPathRuntimeMessageEvent", "UiPathRuntimeStateEvent"]
__all__ = [
"UiPathRuntimeConversationMetaEvent",
"UiPathRuntimeMessageEvent",
"UiPathRuntimeStateEvent",
]
2 changes: 2 additions & 0 deletions src/uipath/runtime/workspace/__init__.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
"""Workspace persistence primitives for runtime implementations."""

from uipath.runtime.workspace.context import get_workspace_path
from uipath.runtime.workspace.conversational import ConversationalWorkspaceRuntime
from uipath.runtime.workspace.hydration import (
HydrationPolicy,
HydrationRuntime,
Expand All @@ -14,6 +15,7 @@

__all__ = [
"AttachmentRegistryEntry",
"ConversationalWorkspaceRuntime",
"HydrationPolicy",
"HydrationRuntime",
"get_workspace_path",
Expand Down
211 changes: 211 additions & 0 deletions src/uipath/runtime/workspace/conversational.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,211 @@
"""Attachment-backed workspace persistence for conversational agents."""

import logging
from collections.abc import Mapping
from typing import Any, AsyncGenerator

from uipath.runtime.base import (
UiPathExecuteOptions,
UiPathRuntimeProtocol,
UiPathStreamOptions,
)
from uipath.runtime.events import (
UiPathRuntimeConversationMetaEvent,
UiPathRuntimeEvent,
)
from uipath.runtime.result import UiPathRuntimeResult, UiPathRuntimeStatus
from uipath.runtime.schema import UiPathRuntimeSchema
from uipath.runtime.workspace.hydrator import WorkspaceHydrator
from uipath.runtime.workspace.registry_store import WorkspaceRegistryStore

logger = logging.getLogger(__name__)

CONVERSATION_META_EVENTS_INPUT_KEY = "uipath__conversation_meta_events"
WORKSPACE_FILES_META_KEY = "workspaceFiles"
WORKSPACE_FILE_PATH_KEY = "path"
WORKSPACE_FILE_ATTACHMENT_KEY = "attachmentKey"


class _InvalidWorkspaceSnapshot(ValueError):
pass


def _meta_event_payload(
event: Mapping[object, object],
) -> Mapping[object, object] | None:
exchange = event.get("exchange")
if isinstance(exchange, Mapping):
exchange_meta_event = exchange.get("metaEvent")
if isinstance(exchange_meta_event, Mapping):
return exchange_meta_event

meta_event = event.get("metaEvent")
return meta_event if isinstance(meta_event, Mapping) else None


def _attachment_keys_from_meta_events(

Check failure on line 46 in src/uipath/runtime/workspace/conversational.py

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Refactor this function to reduce its Cognitive Complexity from 19 to the 15 allowed.

See more on https://sonarcloud.io/project/issues?id=UiPath_uipath-runtime-python&issues=AZ_H8ShdRt9J36t0Aqf9&open=AZ_H8ShdRt9J36t0Aqf9&pullRequest=145
input: Mapping[str, object] | None,
) -> dict[str, str] | None:
events = (input or {}).get(CONVERSATION_META_EVENTS_INPUT_KEY)
if not isinstance(events, list):
return None

for event in reversed(events):
if not isinstance(event, Mapping):
continue
meta_event = _meta_event_payload(event)
if meta_event is None or WORKSPACE_FILES_META_KEY not in meta_event:
continue

workspace_files = meta_event[WORKSPACE_FILES_META_KEY]
if not isinstance(workspace_files, list):
raise _InvalidWorkspaceSnapshot

attachment_keys_by_path: dict[str, str] = {}
for workspace_file in workspace_files:
if not isinstance(workspace_file, Mapping):
raise _InvalidWorkspaceSnapshot
path = workspace_file.get(WORKSPACE_FILE_PATH_KEY)
attachment_key = workspace_file.get(WORKSPACE_FILE_ATTACHMENT_KEY)
if not isinstance(path, str) or not isinstance(attachment_key, str):
raise _InvalidWorkspaceSnapshot
attachment_keys_by_path[path] = attachment_key
return attachment_keys_by_path

return None


def _without_conversation_meta_events(
input: dict[str, Any] | None,
) -> dict[str, Any] | None:
if input is None or CONVERSATION_META_EVENTS_INPUT_KEY not in input:
return input
return {
key: value
for key, value in input.items()
if key != CONVERSATION_META_EVENTS_INPUT_KEY
}


class ConversationalWorkspaceRuntime:
"""Persists workspace attachments between conversational jobs."""

def __init__(
self,
delegate: UiPathRuntimeProtocol,
*,
hydrator: WorkspaceHydrator,
registry_store: WorkspaceRegistryStore | None = None,
):
"""Initialize the wrapper with its delegate and hydrator."""
self.delegate = delegate
self.hydrator = hydrator
self.registry_store = registry_store
self._registry: dict[str, dict[str, Any]] = {}
self._hydrated = False

async def execute(
self,
input: dict[str, Any] | None = None,
options: UiPathExecuteOptions | None = None,
) -> UiPathRuntimeResult:
"""Execute by draining the stream."""
result: UiPathRuntimeResult | None = None
stream_options = (
UiPathStreamOptions.model_validate(options.model_dump())
if options is not None
else None
)
async for event in self.stream(input, options=stream_options):
if isinstance(event, UiPathRuntimeResult):
result = event
if result is None:
raise RuntimeError("Delegate stream completed without a runtime result")
return result

async def stream(
self,
input: dict[str, Any] | None = None,
options: UiPathStreamOptions | None = None,
) -> AsyncGenerator[UiPathRuntimeEvent, None]:
"""Hydrate, stream the delegate, then emit the workspace snapshot."""
await self._hydrate(input)
delegate_input = _without_conversation_meta_events(input)
final_result: UiPathRuntimeResult | None = None

async for event in self.delegate.stream(delegate_input, options=options):
if isinstance(event, UiPathRuntimeResult):
final_result = event
else:
yield event

if final_result is None:
return
if final_result.status == UiPathRuntimeStatus.SUCCESSFUL:
yield await self._dehydrate()
yield final_result

async def get_schema(self) -> UiPathRuntimeSchema:
"""Passthrough schema from delegate runtime."""
return await self.delegate.get_schema()

async def dispose(self) -> None:
"""Release resources owned by this wrapper."""

async def _hydrate(self, input: Mapping[str, object] | None) -> None:
if self._hydrated:
return

persisted_registry = (
await self.registry_store.try_load() if self.registry_store else None
)
if persisted_registry is not None:
self._registry = persisted_registry
source = "suspended job state"
else:
try:
attachment_keys_by_path = _attachment_keys_from_meta_events(input)
except _InvalidWorkspaceSnapshot:
logger.warning("Ignoring malformed conversational workspace snapshot")
source = "existing workspace"
else:
if attachment_keys_by_path is None:
source = "existing workspace"
else:
source = "conversation metadata"
self._registry = await self.hydrator.hydrate_from_attachments(
attachment_keys_by_path
)

logger.info(
"Conversational workspace initialized: %d file(s) from %s",
len(self._registry),
source,
)
self._hydrated = True

async def _dehydrate(self) -> UiPathRuntimeConversationMetaEvent:
registry = self._registry
if self.registry_store is not None:
persisted_registry = await self.registry_store.try_load()
if persisted_registry is not None:
registry = persisted_registry

self._registry = await self.hydrator.dehydrate(registry)
if self.registry_store is not None:
await self.registry_store.save(self._registry)

workspace_files = [
{
WORKSPACE_FILE_PATH_KEY: virtual_path,
WORKSPACE_FILE_ATTACHMENT_KEY: entry["attachment_key"],
}
for virtual_path, entry in sorted(self._registry.items())
]
logger.info(
"Conversational workspace dehydrate: emitting %d file(s)",
len(workspace_files),
)
return UiPathRuntimeConversationMetaEvent(
payload={WORKSPACE_FILES_META_KEY: workspace_files}
)
2 changes: 1 addition & 1 deletion src/uipath/runtime/workspace/hydration.py
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,7 @@ async def dispose(self) -> None:

async def _hydrate(self) -> None:
registry = await self.registry_store.load()
hydrated = await self._get_hydrator().hydrate(registry)
hydrated = await self._get_hydrator().hydrate_from_registry(registry)
if hydrated != registry:
await self.registry_store.save(hydrated)

Expand Down
Loading
Loading