From 7a9517a7b8d53efe58dbf86787441782d023e0d4 Mon Sep 17 00:00:00 2001 From: Dani Date: Fri, 31 Jul 2026 03:04:32 -0400 Subject: [PATCH] feat: add ergonomic EverOS client wrapper MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Hand-maintained facade over the generated MemoryApi/StorageApi — plain kwargs/dicts in, response .data out — exported as `from everos_cloud import EverOS`. - everos_cloud/client.py: EverOS with add/search/get/flush/edit/delete + end-to-end upload() (presign + direct-to-S3) and presign(); unified EverOSError hierarchy; default request timeout; flush getattr-fallback so it survives the flush_api_v2_memory_flush_post -> flush_memory rename. - everos_cloud/__init__.py: export EverOS (guarded try/except so a bare generated tree without client.py still imports). - tests/test_client.py: 18 offline unit tests (API layer mocked). - ci.yml: run the wrapper tests + assert the top-level EverOS export. Validated: 18 unit tests green; the direct-to-S3 upload mechanics verified against real S3 (204, incl. the presigned Content-Type condition) via the storage sign endpoint. Co-Authored-By: Claude Opus 4.8 (1M context) --- .github/workflows/ci.yml | 12 +- everos_cloud/__init__.py | 13 ++ everos_cloud/client.py | 361 +++++++++++++++++++++++++++++++++++++++ tests/test_client.py | 224 ++++++++++++++++++++++++ 4 files changed, 609 insertions(+), 1 deletion(-) create mode 100644 everos_cloud/client.py create mode 100644 tests/test_client.py diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 8551157..b8ae973 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -41,13 +41,16 @@ jobs: if: steps.detect.outputs.has_pkg == 'true' run: | python - <<'PY' - from everos_cloud import ApiClient, Configuration, MemoryApi + from everos_cloud import ApiClient, Configuration, MemoryApi, EverOS cfg = Configuration(access_token="sk-test") assert cfg.host == "https://api.evermind.ai", cfg.host assert cfg.auth_settings()["BearerAuth"]["value"] == "Bearer sk-test" MemoryApi(ApiClient(cfg)) for m in ("add_memory", "search_memory", "get_memory", "delete_memory", "edit_profile"): assert hasattr(MemoryApi, m), m + # ergonomic facade is exported at the top level + for m in ("add", "search", "get", "flush", "edit", "delete", "upload", "close"): + assert hasattr(EverOS, m), m print("import + config OK") PY @@ -88,6 +91,13 @@ jobs: print("quickstart drift smoke OK") PY + - name: Wrapper unit tests + # Offline tests for the hand-maintained ergonomic client (everos_cloud/client.py). + if: steps.detect.outputs.has_pkg == 'true' && hashFiles('tests/test_client.py') != '' + run: | + pip install pytest + pytest -q tests/ + - name: Build sdist + wheel if: steps.detect.outputs.has_pkg == 'true' run: | diff --git a/everos_cloud/__init__.py b/everos_cloud/__init__.py index 67859ca..603fd51 100644 --- a/everos_cloud/__init__.py +++ b/everos_cloud/__init__.py @@ -91,3 +91,16 @@ from everos_cloud.models.update_operation import UpdateOperation from everos_cloud.models.validation_error import ValidationError from everos_cloud.models.validation_error_loc_inner import ValidationErrorLocInner + +# Ergonomic high-level client — hand-maintained, preserved across regeneration +# (usage in quickstart.md). Guarded so a bare generated tree without client.py +# (e.g. the factory's build output) still imports cleanly. +try: + from everos_cloud.client import ( + EverOS, + EverOSAPIError, + EverOSError, + EverOSStorageError, + ) +except ImportError: # pragma: no cover + pass diff --git a/everos_cloud/client.py b/everos_cloud/client.py new file mode 100644 index 0000000..43597fb --- /dev/null +++ b/everos_cloud/client.py @@ -0,0 +1,361 @@ +"""High-level ergonomic client for the EverOS Cloud Memory API. + +A thin, hand-maintained facade over the generated ``MemoryApi`` / ``StorageApi``: +plain kwargs / dicts in, response ``.data`` out. The generated typed client stays +available via ``client.memory`` / ``client.storage`` for full control. + + from everos_cloud import EverOS + + client = EverOS(api_key="sk-...") + client.add(session_id="s1", messages=[{"role": "user", "content": "I love hiking"}]) + hits = client.search("outdoor hobbies") + +Errors: every failure raised by this facade derives from :class:`EverOSError` — +``EverOSAPIError`` for memory HTTP errors, ``EverOSStorageError`` for object-upload +failures. +""" +from __future__ import annotations + +import mimetypes +import os +import time +from typing import Any, Mapping, Sequence, Union + +import urllib3 + +from everos_cloud import ApiClient, Configuration, MemoryApi, StorageApi +from everos_cloud.exceptions import ApiException +from everos_cloud.models import ( + AddInput, + AddOperation, + Content, + DeleteInput, + DeleteOperation, + EditInput, + EditInputOperationsInner, + FlushInput, + GetInput, + MessageItem, + SearchInput, + SignObjectItem, + SignRequest, + UpdateOperation, +) + +__all__ = ["EverOS", "EverOSError", "EverOSAPIError", "EverOSStorageError"] + +DEFAULT_TIMEOUT = 60.0 # seconds; agentic search / LLM rerank can be slow + +_OP_CLASSES = {"add": AddOperation, "update": UpdateOperation, "delete": DeleteOperation} + +_IMAGE_EXT = {".jpg", ".jpeg", ".png", ".gif", ".webp", ".bmp", ".svg"} +_VIDEO_EXT = {".mp4", ".mov", ".avi", ".mkv", ".webm", ".m4v"} + +MessageLike = Union[MessageItem, Mapping[str, Any]] + + +def _guess_file_type(name: str) -> str: + ext = os.path.splitext(name)[1].lower() + if ext in _IMAGE_EXT: + return "image" + if ext in _VIDEO_EXT: + return "video" + return "file" + + +def _clean(**kwargs: Any) -> dict: + """Drop ``None`` values so model defaults (e.g. search ``method='hybrid'``) survive.""" + return {k: v for k, v in kwargs.items() if v is not None} + + +class EverOSError(Exception): + """Base class for every error raised by the EverOS facade.""" + + +class EverOSAPIError(EverOSError): + """A memory endpoint returned an HTTP error (wraps the SDK's ``ApiException``).""" + + def __init__(self, status: Any, body: Any = None) -> None: + super().__init__(f"EverOS API error: status={status}") + self.status = status + self.body = body + + +class EverOSStorageError(EverOSError): + """The storage (object-sign) endpoint returned a non-zero business status.""" + + def __init__(self, status: int, error: Any = None) -> None: + super().__init__(f"EverOS storage error: status={status} error={error}") + self.status = status + self.error = error + + +class EverOS: + """Ergonomic wrapper around the EverOS Cloud Memory API. + + ``timeout`` (seconds) applies to every request; override per client. Note the + ergonomic defaults in :meth:`add`: a message without ``timestamp`` is stamped + with the current time, and without ``sender_id`` defaults to its ``role`` — pass + them explicitly when backfilling historical or multi-party conversations. + """ + + def __init__( + self, + api_key: str, + *, + host: str | None = None, + app_id: str = "default", + project_id: str = "default", + timeout: float = DEFAULT_TIMEOUT, + ) -> None: + cfg = Configuration(access_token=api_key) + if host: + cfg.host = host + self._api_client = ApiClient(cfg) + # Escape hatches: the generated low-level clients, for anything the facade omits. + self.memory = MemoryApi(self._api_client) + self.storage = StorageApi(self._api_client) + self._app_id = app_id + self._project_id = project_id + self._timeout = timeout + + # -- lifecycle ----------------------------------------------------------- + def __enter__(self) -> "EverOS": + return self + + def __exit__(self, *exc: Any) -> None: + self.close() + + def close(self) -> None: + """Release pooled HTTP connections. (The generated ApiClient has no close().)""" + pool = getattr(self._api_client.rest_client, "pool_manager", None) + if pool is not None: + pool.clear() + + # -- helpers ------------------------------------------------------------- + def _scope(self, app_id: str | None, project_id: str | None) -> dict: + return {"app_id": app_id or self._app_id, "project_id": project_id or self._project_id} + + def _call(self, fn: Any, *args: Any, **kwargs: Any) -> Any: + """Invoke a generated method with the client timeout, normalizing errors.""" + kwargs.setdefault("_request_timeout", self._timeout) + try: + return fn(*args, **kwargs) + except ApiException as exc: # pragma: no cover - exercised via integration + raise EverOSAPIError(exc.status, getattr(exc, "body", None)) from exc + + @staticmethod + def _to_message(m: MessageLike) -> MessageItem: + if isinstance(m, MessageItem): + return m + m = dict(m) + content = m.get("content") + if not isinstance(content, Content): + content = Content(content) + role = m.get("role", "user") + return MessageItem( + **_clean( + sender_id=m.get("sender_id") or role, + sender_name=m.get("sender_name"), + role=role, + timestamp=m["timestamp"] if m.get("timestamp") is not None else int(time.time()), + content=content, + tool_calls=m.get("tool_calls"), + tool_call_id=m.get("tool_call_id"), + ) + ) + + @staticmethod + def _to_operation(op: Any) -> EditInputOperationsInner: + if isinstance(op, EditInputOperationsInner): + return op + op = dict(op) + cls = _OP_CLASSES.get(op.get("action")) + if cls is None: + raise ValueError(f"unknown edit operation action: {op.get('action')!r}") + return EditInputOperationsInner(cls(**op)) + + # -- memory -------------------------------------------------------------- + def add( + self, + session_id: str, + messages: Sequence[MessageLike], + *, + mode: str | None = None, + async_mode: bool | None = None, + app_id: str | None = None, + project_id: str | None = None, + ) -> Any: + """Add conversation messages to a session. Returns the ``AddData`` result.""" + payload = AddInput( + **_clean( + session_id=session_id, + messages=[self._to_message(m) for m in messages], + mode=mode, + async_mode=async_mode, + **self._scope(app_id, project_id), + ) + ) + return self._call(self.memory.add_memory, payload).data + + def search( + self, + query: str, + *, + method: str | None = None, + top_k: int | None = None, + user_id: str | None = None, + agent_id: str | None = None, + include_profile: bool | None = None, + min_score: float | None = None, + radius: float | None = None, + enable_llm_rerank: bool | None = None, + filters: Any = None, + app_id: str | None = None, + project_id: str | None = None, + ) -> Any: + """Search memories (keyword / vector / hybrid / agentic). Returns ``SearchData``.""" + payload = SearchInput( + **_clean( + query=query, + method=method, + top_k=top_k, + user_id=user_id, + agent_id=agent_id, + include_profile=include_profile, + min_score=min_score, + radius=radius, + enable_llm_rerank=enable_llm_rerank, + filters=filters, + **self._scope(app_id, project_id), + ) + ) + return self._call(self.memory.search_memory, payload).data + + def get( + self, + memory_type: str, + *, + user_id: str | None = None, + agent_id: str | None = None, + page: int | None = None, + page_size: int | None = None, + sort_by: str | None = None, + sort_order: str | None = None, + filters: Any = None, + app_id: str | None = None, + project_id: str | None = None, + ) -> Any: + """Get memories (paginated) by type. Returns ``GetData``.""" + payload = GetInput( + **_clean( + memory_type=memory_type, + user_id=user_id, + agent_id=agent_id, + page=page, + page_size=page_size, + sort_by=sort_by, + sort_order=sort_order, + filters=filters, + **self._scope(app_id, project_id), + ) + ) + return self._call(self.memory.get_memory, payload).data + + def flush(self, session_id: str, *, app_id: str | None = None, project_id: str | None = None) -> Any: + """Force extraction for a session. Returns ``FlushData``.""" + payload = FlushInput(**_clean(session_id=session_id, **self._scope(app_id, project_id))) + # Method name is being cleaned up (flush_api_v2_memory_flush_post -> flush_memory); + # call whichever the installed SDK exposes so the facade survives the rename. + fn = getattr(self.memory, "flush_memory", None) or self.memory.flush_api_v2_memory_flush_post + return self._call(fn, payload).data + + def edit( + self, + user_id: str, + operations: Sequence[Any], + *, + app_id: str | None = None, + project_id: str | None = None, + ) -> Any: + """Bulk-edit a user's profile. ``operations`` are dicts (``action`` add/update/delete).""" + payload = EditInput( + **_clean( + user_id=user_id, + operations=[self._to_operation(op) for op in operations], + **self._scope(app_id, project_id), + ) + ) + return self._call(self.memory.edit_profile, payload).data + + def delete( + self, + *, + user_id: str | None = None, + agent_id: str | None = None, + session_id: str | None = None, + app_id: str | None = None, + project_id: str | None = None, + ) -> Any: + """Scoped soft-delete of memories. Returns ``DeleteData``.""" + payload = DeleteInput( + **_clean( + user_id=user_id, + agent_id=agent_id, + session_id=session_id, + **self._scope(app_id, project_id), + ) + ) + return self._call(self.memory.delete_memory, payload).data + + # -- storage ------------------------------------------------------------- + def presign(self, objects: Sequence[Any]) -> Any: + """Presign objects for direct-to-S3 upload (advanced/batch, up to 50). + + Returns the ``SignResponse`` data; raises ``EverOSStorageError`` on a + non-zero business status. Most callers want :meth:`upload` instead. + """ + items = [o if isinstance(o, SignObjectItem) else SignObjectItem(**o) for o in objects] + envelope = self._call(self.storage.sign_objects, SignRequest(object_list=items)) + if envelope.status != 0: + raise EverOSStorageError(envelope.status, getattr(envelope, "error", None)) + return envelope.result.data + + def upload( + self, + path: str, + *, + file_type: str | None = None, + file_id: str | None = None, + file_name: str | None = None, + ) -> str: + """Upload a local file end to end and return its ``object_key``. + + Presigns the object, POSTs the bytes straight to S3, and returns the + ``object_key`` you then reference in a message's multimodal content. + ``file_type`` (image / file / video) is inferred from the extension if omitted. + The whole file is read into memory — fine for images/docs, mind large videos. + """ + file_name = file_name or os.path.basename(path) + file_id = file_id or file_name + file_type = file_type or _guess_file_type(file_name) + with open(path, "rb") as fh: + body = fh.read() + + data = self.presign([{"file_id": file_id, "file_name": file_name, "file_type": file_type}]) + obj = data.object_list[0] + signed = obj.object_signed_info + + status = self._s3_post(signed.url, dict(signed.fields), file_name, body, self._timeout) + if status not in (200, 201, 204): + raise EverOSStorageError(status, "direct-to-S3 upload failed") + return obj.object_key + + @staticmethod + def _s3_post(url: str, fields: dict, filename: str, body: bytes, timeout: float = DEFAULT_TIMEOUT) -> int: + """POST a presigned multipart form to S3. The ``file`` field must come last.""" + content_type = mimetypes.guess_type(filename)[0] or "application/octet-stream" + form = dict(fields) + form["file"] = (filename, body, content_type) # urllib3: (filename, data, content_type) + resp = urllib3.PoolManager().request("POST", url, fields=form, timeout=timeout) + return resp.status diff --git a/tests/test_client.py b/tests/test_client.py new file mode 100644 index 0000000..8ceb376 --- /dev/null +++ b/tests/test_client.py @@ -0,0 +1,224 @@ +"""Offline unit tests for the EverOS facade — the API layer is mocked, no network.""" +from types import SimpleNamespace +from unittest.mock import MagicMock + +import pytest +from everos_cloud import EverOS, EverOSAPIError, EverOSError, EverOSStorageError +from everos_cloud.exceptions import ApiException +from everos_cloud.models import AddInput, EditInput, MessageItem, SearchInput + + +def _client(memory=None, storage=None): + c = EverOS("sk-test") + if memory is not None: + c.memory = memory + if storage is not None: + c.storage = storage + return c + + +def test_add_builds_payload_and_returns_data(): + mem = MagicMock() + mem.add_memory.return_value = SimpleNamespace(data="ADD_RESULT") + c = _client(memory=mem) + + out = c.add(session_id="s1", messages=[{"role": "user", "content": "I love hiking"}]) + + assert out == "ADD_RESULT" + payload = mem.add_memory.call_args.args[0] + assert isinstance(payload, AddInput) + assert payload.session_id == "s1" + assert payload.app_id == "default" and payload.project_id == "default" + msg = payload.messages[0] + assert isinstance(msg, MessageItem) + assert msg.role == "user" + assert msg.sender_id == "user" # defaults to role when omitted + assert isinstance(msg.timestamp, int) # defaults to now + # string content shorthand became a Content wrapper + assert msg.content is not None + + +def test_message_passthrough_and_explicit_fields(): + mem = MagicMock() + mem.add_memory.return_value = SimpleNamespace(data=None) + c = _client(memory=mem) + c.add(session_id="s1", messages=[ + {"sender_id": "u9", "role": "assistant", "timestamp": 1700000000, "content": "hi"}, + ]) + msg = mem.add_memory.call_args.args[0].messages[0] + assert msg.sender_id == "u9" and msg.role == "assistant" and msg.timestamp == 1700000000 + + +def test_search_drops_none_so_defaults_survive(): + mem = MagicMock() + mem.search_memory.return_value = SimpleNamespace(data="S") + c = _client(memory=mem) + + assert c.search("outdoor hobbies") == "S" + payload = mem.search_memory.call_args.args[0] + assert isinstance(payload, SearchInput) + assert payload.query == "outdoor hobbies" + assert payload.method == "hybrid" # model default preserved (not overwritten with None) + + +def test_get_and_delete(): + mem = MagicMock() + mem.get_memory.return_value = SimpleNamespace(data="G") + mem.delete_memory.return_value = SimpleNamespace(data="D") + c = _client(memory=mem) + + assert c.get("episode", page=2) == "G" + assert mem.get_memory.call_args.args[0].memory_type == "episode" + assert mem.get_memory.call_args.args[0].page == 2 + + assert c.delete(user_id="u1", session_id="s1") == "D" + dp = mem.delete_memory.call_args.args[0] + assert dp.user_id == "u1" and dp.session_id == "s1" + + +def test_edit_wraps_operations(): + mem = MagicMock() + mem.edit_profile.return_value = SimpleNamespace(data="E") + c = _client(memory=mem) + + c.edit("u1", [{"action": "add", "type": "explicit_info", + "data": {"category": "hobby", "description": "hiking"}, "reason": "x"}]) + payload = mem.edit_profile.call_args.args[0] + assert isinstance(payload, EditInput) + assert payload.user_id == "u1" + assert len(payload.operations) == 1 + + +def test_edit_rejects_unknown_action(): + c = _client(memory=MagicMock()) + with pytest.raises(ValueError): + c.edit("u1", [{"action": "nope", "type": "explicit_info", "data": {}}]) + + +def test_flush_prefers_flush_memory_when_present(): + calls = {} + + def new_flush(p, **kw): + calls["new"] = p + return SimpleNamespace(data="NEW") + + def old_flush(p, **kw): + calls["old"] = p + return SimpleNamespace(data="OLD") + + mem = SimpleNamespace(flush_memory=new_flush, flush_api_v2_memory_flush_post=old_flush) + assert _client(memory=mem).flush("s1") == "NEW" + assert "new" in calls and "old" not in calls + + +def test_flush_falls_back_to_generated_name(): + mem = SimpleNamespace( + flush_api_v2_memory_flush_post=lambda p, **kw: SimpleNamespace(data="OLD"), + ) + assert _client(memory=mem).flush("s1") == "OLD" + + +def test_presign_success_returns_data(): + storage = MagicMock() + storage.sign_objects.return_value = SimpleNamespace(status=0, result=SimpleNamespace(data="SIGNED")) + c = _client(storage=storage) + assert c.presign([{"file_id": "f1", "file_name": "a.jpg", "file_type": "image"}]) == "SIGNED" + + +def test_presign_raises_on_business_error(): + storage = MagicMock() + storage.sign_objects.return_value = SimpleNamespace(status=2018, error="validation failed") + c = _client(storage=storage) + with pytest.raises(EverOSStorageError) as ei: + c.presign([{"file_id": "f1", "file_name": "a.jpg", "file_type": "image"}]) + assert ei.value.status == 2018 + + +def test_custom_scope_defaults_propagate(): + mem = MagicMock() + mem.add_memory.return_value = SimpleNamespace(data=None) + c = EverOS("sk-test", app_id="myapp", project_id="proj") + c.memory = mem + c.add(session_id="s1", messages=[{"role": "user", "content": "hi"}]) + payload = mem.add_memory.call_args.args[0] + assert payload.app_id == "myapp" and payload.project_id == "proj" + + +def _signed_envelope(object_key="k1", url="https://s3.example/upload", fields=None): + signed = SimpleNamespace(url=url, fields=fields or {"key": "abc", "policy": "xyz"}, max_size=100) + obj = SimpleNamespace(object_key=object_key, object_signed_info=signed) + return SimpleNamespace(status=0, result=SimpleNamespace(data=SimpleNamespace(object_list=[obj]))) + + +def test_upload_end_to_end(tmp_path): + storage = MagicMock() + storage.sign_objects.return_value = _signed_envelope(object_key="obj-123") + c = _client(storage=storage) + captured = {} + + def fake_post(url, fields, filename, body, timeout=None): + captured.update(url=url, fields=fields, filename=filename, body=body, timeout=timeout) + return 204 + + c._s3_post = fake_post + f = tmp_path / "photo.jpg" + f.write_bytes(b"BYTES") + + key = c.upload(str(f)) + + assert key == "obj-123" + req = storage.sign_objects.call_args.args[0] + assert req.object_list[0].file_type == "image" # inferred from .jpg + assert req.object_list[0].file_name == "photo.jpg" + assert captured["body"] == b"BYTES" + assert captured["filename"] == "photo.jpg" + assert "key" in captured["fields"] # presigned form fields forwarded + + +def test_upload_raises_on_s3_failure(tmp_path): + storage = MagicMock() + storage.sign_objects.return_value = _signed_envelope() + c = _client(storage=storage) + c._s3_post = lambda *a: 403 + f = tmp_path / "doc.pdf" + f.write_bytes(b"x") + with pytest.raises(EverOSStorageError): + c.upload(str(f)) + + +def test_guess_file_type(): + from everos_cloud.client import _guess_file_type + assert _guess_file_type("a.PNG") == "image" + assert _guess_file_type("clip.mp4") == "video" + assert _guess_file_type("report.pdf") == "file" + assert _guess_file_type("noext") == "file" + + +def test_memory_api_error_is_wrapped(): + mem = MagicMock() + mem.search_memory.side_effect = ApiException(status=401, reason="unauthorized") + c = _client(memory=mem) + with pytest.raises(EverOSAPIError) as ei: + c.search("q") + assert ei.value.status == 401 + assert isinstance(ei.value, EverOSError) # unified hierarchy + + +def test_timeout_is_applied_to_calls(): + mem = MagicMock() + mem.get_memory.return_value = SimpleNamespace(data=None) + c = EverOS("sk-test", timeout=5) + c.memory = mem + c.get("episode") + assert mem.get_memory.call_args.kwargs.get("_request_timeout") == 5 + + +def test_storage_error_is_in_hierarchy(): + assert issubclass(EverOSStorageError, EverOSError) + + +def test_close_and_context_manager_do_not_raise(): + # exercises __enter__/__exit__ -> close() against the real generated ApiClient + with EverOS("sk-test"): + pass + EverOS("sk-test").close() # idempotent / direct