Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Feature/event listener drivers #735

Merged
merged 10 commits into from
Apr 10, 2024
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
10 changes: 9 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,15 +8,23 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0

### Added
- `list_files_from_disk` activity to `FileManager` Tool.
- Support for Drivers in `EventListener`.
- `AmazonSqsEventListenerDriver` for sending events to an Amazon SQS queue.
- `AwsIotCoreEventListenerDriver` for sending events to a topic on AWS IoT Core.
- `GriptapeCloudEventListenerDriver` for sending events to Griptape Cloud.
- `WebhookEventListenerDriver` for sending events to a webhook.
- `LocalEventListenerDriver` for sending events to a callback function.


### Changed
- Improved RAG performance in `VectorQueryEngine`.
- **BREAKING**: Secret fields (ex: api_key) removed from serialized Drivers.
- **BREAKING**: Remove `FileLoader`.
- **BREAKING**: `CsvLoader` no longer accepts `str` file paths as a source. It will now accept the content of the CSV file as a `str` or `bytes` object.
- **BREAKING**: `PdfLoader` no longer accepts `str` file content, `Path` file paths or `IO` objects as sources. Instead, it will only accept the content of the PDF file as a `bytes` object.
- **BREAKING**: `TextLoader` no longer accepts `Path` file paths as a source. It will now accept the content of the text file as a `str` or `bytes` object.
- **BREAKING**: `FileManager.default_loader` is now `None` by default.
- **BREAKING**: Replaced `EventListener.handler` with `EventListener.driver` and `LocalEventListenerDriver`.
- Improved RAG performance in `VectorQueryEngine`.

## [0.24.2] - 2024-04-04

Expand Down
13 changes: 13 additions & 0 deletions griptape/drivers/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,13 @@
from .web_scraper.trafilatura_web_scraper_driver import TrafilaturaWebScraperDriver
from .web_scraper.markdownify_web_scraper_driver import MarkdownifyWebScraperDriver

from .event_listener.base_event_listener_driver import BaseEventListenerDriver
from .event_listener.amazon_sqs_event_listener_driver import AmazonSqsEventListenerDriver
from .event_listener.webhook_event_listener_driver import WebhookEventListenerDriver
from .event_listener.aws_iot_core_event_listener_driver import AwsIotCoreEventListenerDriver
from .event_listener.griptape_cloud_event_listener_driver import GriptapeCloudEventListenerDriver
from .event_listener.local_event_listener_driver import LocalEventListenerDriver

__all__ = [
"BasePromptDriver",
"OpenAiChatPromptDriver",
Expand Down Expand Up @@ -161,4 +168,10 @@
"BaseWebScraperDriver",
"TrafilaturaWebScraperDriver",
"MarkdownifyWebScraperDriver",
"BaseEventListenerDriver",
"AmazonSqsEventListenerDriver",
"WebhookEventListenerDriver",
"AwsIotCoreEventListenerDriver",
"GriptapeCloudEventListenerDriver",
"LocalEventListenerDriver",
]
Empty file.
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
from __future__ import annotations

from typing import TYPE_CHECKING, Any
import json

from attr import Factory, define, field

from griptape.drivers.event_listener.base_event_listener_driver import BaseEventListenerDriver
from griptape.events.base_event import BaseEvent
from griptape.utils import import_optional_dependency

if TYPE_CHECKING:
import boto3


@define
class AmazonSqsEventListenerDriver(BaseEventListenerDriver):
queue_url: str = field(kw_only=True)
session: boto3.Session = field(default=Factory(lambda: import_optional_dependency("boto3").Session()), kw_only=True)
sqs_client: Any = field(default=Factory(lambda self: self.session.client("sqs"), takes_self=True))

def try_publish_event(self, event: BaseEvent) -> None:
self.sqs_client.send_message(QueueUrl=self.queue_url, MessageBody=json.dumps({"event": event.to_dict()}))
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
from __future__ import annotations

from typing import TYPE_CHECKING, Any

import json
from attr import Factory, define, field

from griptape.drivers.event_listener.base_event_listener_driver import BaseEventListenerDriver
from griptape.events.base_event import BaseEvent
from griptape.utils import import_optional_dependency

if TYPE_CHECKING:
import boto3


@define
class AwsIotCoreEventListenerDriver(BaseEventListenerDriver):
iot_endpoint: str = field(kw_only=True)
topic: str = field(kw_only=True)
session: boto3.Session = field(default=Factory(lambda: import_optional_dependency("boto3").Session()), kw_only=True)
iotdata_client: Any = field(default=Factory(lambda self: self.session.client("iot-data"), takes_self=True))

def try_publish_event(self, event: BaseEvent) -> None:
self.iotdata_client.publish(topic=self.topic, payload=json.dumps({"event": event.to_dict()}))
20 changes: 20 additions & 0 deletions griptape/drivers/event_listener/base_event_listener_driver.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
from __future__ import annotations
from abc import ABC, abstractmethod
from concurrent import futures
from attr import define, field, Factory
from typing import TYPE_CHECKING

if TYPE_CHECKING:
from griptape.events import BaseEvent


@define
class BaseEventListenerDriver(ABC):
futures_executor: futures.Executor = field(default=Factory(lambda: futures.ThreadPoolExecutor()), kw_only=True)

def publish_event(self, event: BaseEvent) -> None:
self.futures_executor.submit(self.try_publish_event, event)

@abstractmethod
def try_publish_event(self, event: BaseEvent) -> None:
...
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
from __future__ import annotations

import os
import requests

from urllib.parse import urljoin
from attr import define, field, Factory

from griptape.drivers.event_listener.base_event_listener_driver import BaseEventListenerDriver
from griptape.events.base_event import BaseEvent


@define
class GriptapeCloudEventListenerDriver(BaseEventListenerDriver):
base_url: str = field(default="https://cloud.griptape.ai", kw_only=True)
api_key: str = field(kw_only=True)
headers: dict = field(
default=Factory(lambda self: {"Authorization": f"Bearer {self.api_key}"}, takes_self=True), kw_only=True
)
run_id: str = field(default=Factory(lambda: os.getenv("GT_CLOUD_RUN_ID")), kw_only=True)

@run_id.validator # pyright: ignore
def validate_run_id(self, _, run_id: str):
if run_id is None:
raise ValueError(
"run_id must be set either in the constructor or as an environment variable (GT_CLOUD_RUN_ID)."
)

def try_publish_event(self, event: BaseEvent) -> None:
url = urljoin(self.base_url.strip("/"), f"/api/runs/{self.run_id}/events/")

requests.post(url=url, json=event.to_dict(), headers=self.headers)
18 changes: 18 additions & 0 deletions griptape/drivers/event_listener/local_event_listener_driver.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
from __future__ import annotations

from typing import Callable, Any
from attr import define, field

from griptape.drivers.event_listener.base_event_listener_driver import BaseEventListenerDriver
from griptape.events.base_event import BaseEvent


@define
class LocalEventListenerDriver(BaseEventListenerDriver):
handler: Callable[[BaseEvent], Any] = field(default=None, kw_only=True)

def publish_event(self, event: BaseEvent) -> None:
self.try_publish_event(event)

def try_publish_event(self, event: BaseEvent) -> None:
self.handler(event)
17 changes: 17 additions & 0 deletions griptape/drivers/event_listener/webhook_event_listener_driver.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
from __future__ import annotations

import requests

from attr import define, field

from griptape.drivers.event_listener.base_event_listener_driver import BaseEventListenerDriver
from griptape.events.base_event import BaseEvent


@define
class WebhookEventListenerDriver(BaseEventListenerDriver):
webhook_url: str = field(kw_only=True)
headers: dict = field(default=None, kw_only=True)

def try_publish_event(self, event: BaseEvent) -> None:
requests.post(url=self.webhook_url, json={"event": event.to_dict()}, headers=self.headers)
17 changes: 14 additions & 3 deletions griptape/events/event_listener.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,20 @@
from typing import Callable, Optional, Type, Any
from __future__ import annotations
from typing import Optional, TYPE_CHECKING
from attrs import define, field
from .base_event import BaseEvent

if TYPE_CHECKING:
from griptape.drivers import BaseEventListenerDriver


@define
class EventListener:
handler: Callable[[BaseEvent], Any] = field()
event_types: Optional[list[Type[BaseEvent]]] = field(default=None, kw_only=True)
event_types: Optional[list[type[BaseEvent]]] = field(default=None, kw_only=True)
driver: Optional[BaseEventListenerDriver] = field(default=None, kw_only=True)

def publish_event(self, event: BaseEvent) -> None:
event_types = self.event_types

if event_types is None or type(event) in event_types:
if self.driver is not None:
self.driver.publish_event(event)
6 changes: 1 addition & 5 deletions griptape/structures/structure.py
Original file line number Diff line number Diff line change
Expand Up @@ -251,11 +251,7 @@ def remove_event_listener(self, event_listener: EventListener) -> None:

def publish_event(self, event: BaseEvent) -> None:
for event_listener in self.event_listeners:
handler = event_listener.handler
event_types = event_listener.event_types

if event_types is None or type(event) in event_types:
handler(event)
event_listener.publish_event(event)

def context(self, task: BaseTask) -> dict[str, Any]:
return {"args": self.execution_args, "structure": self}
Expand Down
5 changes: 4 additions & 1 deletion griptape/utils/stream.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,11 +53,14 @@ def run(self, *args) -> Iterator[TextArtifact]:
t.join()

def _run_structure(self, *args):
from griptape.drivers import LocalEventListenerDriver

def event_handler(event: BaseEvent):
self._event_queue.put(event)

stream_event_listener = EventListener(
event_handler, event_types=[CompletionChunkEvent, FinishPromptEvent, FinishStructureRunEvent]
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Now I see why you overrode publish_event in LocalEventListenerDriver

Copy link
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧵

driver=LocalEventListenerDriver(handler=event_handler),
event_types=[CompletionChunkEvent, FinishPromptEvent, FinishStructureRunEvent],
)
self.structure.add_event_listener(stream_event_listener)

Expand Down
22 changes: 19 additions & 3 deletions poetry.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

5 changes: 4 additions & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,9 @@ drivers-embedding-google = ["google-generativeai"]
drivers-web-scraper-trafilatura = ["trafilatura"]
drivers-web-scraper-markdownify = ["playwright", "beautifulsoup4", "markdownify"]

drivers-event-listener-amazon-sqs = ["boto3"]
drivers-event-listener-amazon-iot = ["boto3"]

loaders-dataframe = ["pandas"]
loaders-pdf = ["pypdf"]
loaders-image = ["pillow"]
Expand Down Expand Up @@ -134,7 +137,7 @@ pytest-mock = "*"
mongomock = "*"

twine = ">=4"
moto = {extras = ["dynamodb"], version = "^4.1.8"}
moto = {extras = ["dynamodb", "iotdata", "sqs"], version = "^4.2.13"}
pytest-xdist = "^3.3.1"
pytest-cov = "^4.1.0"
pytest-env = "^1.1.1"
Expand Down
Empty file.
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
from pytest import fixture
from moto import mock_sqs
import boto3
from tests.mocks.mock_event import MockEvent
from griptape.drivers.event_listener.amazon_sqs_event_listener_driver import AmazonSqsEventListenerDriver
from tests.utils.aws import mock_aws_credentials


class TestAmazonSqsEventListenerDriver:
@fixture()
def run_before_and_after_tests(self):
mock_aws_credentials()

@fixture()
def driver(self):
mock = mock_sqs()
mock.start()

session = boto3.Session(region_name="us-east-1")
response = session.client("sqs").create_queue(QueueName="foo-bar")
queue_url = response["QueueUrl"]

yield AmazonSqsEventListenerDriver(queue_url=queue_url, session=session)

mock.stop()

def test_init(self, driver):
assert driver

def test_try_publish_event(self, driver):
driver.try_publish_event(event=MockEvent())
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
from pytest import fixture
from moto import mock_iotdata
import boto3
from tests.mocks.mock_event import MockEvent
from griptape.drivers.event_listener.aws_iot_core_event_listener_driver import AwsIotCoreEventListenerDriver
from tests.utils.aws import mock_aws_credentials


@mock_iotdata
class TestAwsIotCoreEventListenerDriver:
@fixture()
def run_before_and_after_tests(self):
mock_aws_credentials()

@fixture()
def driver(self):
return AwsIotCoreEventListenerDriver(
iot_endpoint="foo bar", topic="fizz buzz", session=boto3.Session(region_name="us-east-1")
)

def test_init(self, driver):
assert driver

def test_try_publish_event(self, driver):
driver.try_publish_event(event=MockEvent())
Loading
Loading