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
15 changes: 11 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,9 +1,16 @@
# Convert Python SDK

This package contains the first parity-focused slice of the Python Convert SDK:
This package is the Python port of the Convert Fullstack SDK.

Current functionality includes:

- deterministic MurmurHash3-based bucketing
- rule comparisons
- rule traversal
- JS-equivalent rule comparisons and traversal
- static-config and `sdkKey` initialization
- experience and feature evaluation through visitor contexts
- conversion tracking and API queue release
- in-memory and datastore-backed visitor persistence
- helper `DataStore` and `FileLogger` utilities
- configurable logger clients and log levels

Networking, tracking, config fetching, and the public SDK surface are intentionally left for follow-up PRs.
The SDK entrypoint is `ConvertSDK`, which can create visitor contexts and evaluate the same CDN config shape used by the JavaScript SDK.
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "convertcom-python-sdk"
version = "0.1.0"
description = "Python port of the Convert Fullstack SDK deterministic core"
description = "Python port of the Convert Fullstack SDK"
readme = "README.md"
requires-python = ">=3.10"
dependencies = []
Expand Down
17 changes: 16 additions & 1 deletion src/convertcom_sdk/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,21 @@
from .events.event_manager import EventManager
from .experience.experience_manager import ExperienceManager
from .features.feature_manager import FeatureManager
from .logger.log_manager import LogManager
from .segments.segments_manager import SegmentsManager
from .bucketing.bucketing_manager import BucketingAllocation, BucketingManager
from .enums import BucketingError, EntityType, FeatureStatus, RuleError, SystemEvents
from .enums import (
BucketingError,
EntityType,
FeatureStatus,
LogLevel,
LogMethod,
RuleError,
SystemEvents,
)
from .rules.rule_manager import RuleManager
from .sdk import ConvertSDK
from .utils import DataStore, FileLogger

__all__ = [
"ApiManager",
Expand All @@ -23,13 +33,18 @@
"ConvertSDK",
"Core",
"DataManager",
"DataStore",
"DataStoreManager",
"DEFAULT_CONFIG",
"EntityType",
"EventManager",
"ExperienceManager",
"FeatureManager",
"FeatureStatus",
"FileLogger",
"LogLevel",
"LogManager",
"LogMethod",
"RuleError",
"RuleManager",
"SegmentsManager",
Expand Down
4 changes: 4 additions & 0 deletions src/convertcom_sdk/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -144,5 +144,9 @@ def onReady(self) -> None:

def close(self) -> None:
self._stop_refresh_timer()
if self._data_manager.data_store_manager and hasattr(
self._data_manager.data_store_manager, "close"
):
self._data_manager.data_store_manager.close()
if hasattr(self._api_manager, "close"):
self._api_manager.close()
23 changes: 22 additions & 1 deletion src/convertcom_sdk/data/data_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ def __init__(
rule_manager: RuleManager,
data_store_manager: DataStoreManager | None = None,
api_manager: Any | None = None,
async_storage: bool = True,
) -> None:
self._config = dict(config or {})
self._lock = threading.RLock()
Expand All @@ -29,6 +30,7 @@ def __init__(
self._rule_manager = rule_manager
self._data_store_manager = data_store_manager
self._api_manager = api_manager
self._async_storage = async_storage
self._environment = self._config.get("environment")
self._account_id = self._data.get("account_id")
self._project_id = (self._data.get("project") or {}).get("id")
Expand Down Expand Up @@ -86,13 +88,32 @@ def put_data(self, visitor_id: str, new_data: Mapping[str, Any] | None = None) -
store_key = self.get_store_key(visitor_id)
current = self.get_data(visitor_id) or {}
updated = object_deep_merge(current, new_data)
if updated == current:
return
with self._lock:
self._bucketed_visitors[store_key] = updated
self._bucketed_visitors.move_to_end(store_key)
while len(self._bucketed_visitors) > self._local_store_limit:
self._bucketed_visitors.popitem(last=False)
if self._data_store_manager:
self._data_store_manager.set(store_key, updated)
stored_segments = dict(current.get("segments") or {})
data = {key: value for key, value in current.items() if key != "segments"}
report_segments = self.filter_report_segments(stored_segments).get("segments") or {}
new_segments = self.filter_report_segments(new_data.get("segments") or {}).get(
"segments"
)
stored_value = (
object_deep_merge(
data,
{"segments": {**report_segments, **new_segments}},
)
if new_segments
else updated
)
Comment on lines +105 to +112

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

high

There's a potential bug here. If new_data contains only non-reportable segments, new_segments will be empty, and the else branch will be taken. This sets stored_value to updated, which includes the non-reportable segments. This seems to contradict the goal of persisting only reportable segments. The stored_value should always be constructed using only filtered, reportable segments.

            stored_value = object_deep_merge(
                data,
                {"segments": {**report_segments, **new_segments}},
            )

if self._async_storage and hasattr(self._data_store_manager, "enqueue"):
self._data_store_manager.enqueue(store_key, stored_value)
else:
self._data_store_manager.set(store_key, stored_value)

def filter_report_segments(
self, visitor_properties: Mapping[str, Any] | None
Expand Down
110 changes: 98 additions & 12 deletions src/convertcom_sdk/data/data_store_manager.py
Original file line number Diff line number Diff line change
@@ -1,27 +1,122 @@
from __future__ import annotations

import threading
from collections.abc import Mapping
from typing import Any

from convertcom_sdk.enums import SystemEvents
from convertcom_sdk.utils.object_utils import object_deep_merge

DEFAULT_BATCH_SIZE = 1
DEFAULT_RELEASE_INTERVAL = 5000


class DataStoreManager:
def __init__(
self,
config: Mapping[str, Any] | None = None,
*,
data_store: Any = None,
event_manager: Any | None = None,
logger_manager: Any | None = None,
) -> None:
del config
config = dict(config or {})
events = config.get("events") or {}
self._lock = threading.RLock()
self._event_manager = event_manager
self._logger_manager = logger_manager
self._mapper = config.get("mapper") or (lambda value: value)
self.batch_size = int(events.get("batch_size") or DEFAULT_BATCH_SIZE)
self.release_interval = int(
events.get("release_interval") or DEFAULT_RELEASE_INTERVAL
)
self._requests_queue: dict[str, Any] = {}
self._requests_queue_timer: threading.Timer | None = None
self.data_store = data_store if self.is_valid_data_store(data_store) else None

def set(self, key: str, data: Any) -> None:
if self.data_store is not None:
if self.data_store is None:
return
try:
self.data_store.set(key, data)
except Exception as error:
if self._logger_manager:
self._logger_manager.error(
"DataStoreManager.set()",
{"error": str(error)},
)

def get(self, key: str) -> Any:
if self.data_store is None:
return None
return self.data_store.get(key)
try:
return self.data_store.get(key)
except Exception as error:
if self._logger_manager:
self._logger_manager.error(
"DataStoreManager.get()",
{"error": str(error)},
)
return None

def enqueue(self, key: str, data: Any) -> None:
if self._logger_manager:
self._logger_manager.trace(
"DataStoreManager.enqueue()",
self._mapper({"key": key, "data": data}),
)
with self._lock:
self._requests_queue = object_deep_merge(
self._requests_queue,
{key: data},
)
queue_length = len(self._requests_queue)
if queue_length >= self.batch_size:
self.release_queue("size")
elif queue_length == 1:
self.start_queue()

def release_queue(self, reason: str | None = None) -> None:
if self._logger_manager:
self._logger_manager.info(
"DataStoreManager.release_queue()",
{"reason": reason or ""},
)
with self._lock:
queued_items = dict(self._requests_queue)
self._requests_queue = {}
self.stop_queue()
for key, value in queued_items.items():
self.set(key, value)
Comment on lines +89 to +90

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

high

Calling set in a loop for each item can be inefficient, causing repeated file reads and writes if using a file-based datastore like the DataStore utility. To optimize this, you can perform a bulk update if the underlying data_store supports it (e.g., with a set_many method).

Suggested change
for key, value in queued_items.items():
self.set(key, value)
if hasattr(self.data_store, "set_many") and callable(
getattr(self.data_store, "set_many")
):
try:
self.data_store.set_many(queued_items)
except Exception as error:
if self._logger_manager:
self._logger_manager.error(
"DataStoreManager.release_queue() bulk set failed",
{"error": str(error)},
)
else:
for key, value in queued_items.items():
self.set(key, value)

if self._event_manager:
self._event_manager.fire(
SystemEvents.DATA_STORE_QUEUE_RELEASED,
{"reason": reason or ""},
)

def releaseQueue(self, reason: str | None = None) -> None:
self.release_queue(reason)
Comment on lines +97 to +98

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

medium

While I understand the goal of achieving parity with the JS SDK, adding camelCase aliases for Python methods goes against PEP 8, which recommends snake_case for functions and methods. This can be confusing for Python developers and harm maintainability. It's best to stick to idiomatic Python and remove these aliases.


def stop_queue(self) -> None:
with self._lock:
timer = self._requests_queue_timer
self._requests_queue_timer = None
if timer:
timer.cancel()

def start_queue(self) -> None:
self.stop_queue()
with self._lock:
timer = threading.Timer(
self.release_interval / 1000.0,
lambda: self.release_queue("timeout"),
)
timer.daemon = True
self._requests_queue_timer = timer
timer.start()

def close(self) -> None:
self.stop_queue()

def is_valid_data_store(self, data_store: Any) -> bool:
return bool(
Expand All @@ -31,12 +126,3 @@ def is_valid_data_store(self, data_store: Any) -> bool:
and hasattr(data_store, "set")
and callable(data_store.set)
)

def release_queue(self, reason: str | None = None) -> Any:
if self.data_store is None:
return None
if hasattr(self.data_store, "release_queue") and callable(self.data_store.release_queue):
return self.data_store.release_queue(reason)
if hasattr(self.data_store, "releaseQueue") and callable(self.data_store.releaseQueue):
return self.data_store.releaseQueue(reason)
return None
20 changes: 19 additions & 1 deletion src/convertcom_sdk/enums.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
from enum import Enum
from enum import Enum, IntEnum


class RuleError(str, Enum):
Expand Down Expand Up @@ -50,3 +50,21 @@ class EntityType(str, Enum):
GOAL = "goal"
EXPERIENCE = "experience"
VARIATION = "variation"


class LogLevel(IntEnum):
TRACE = 0
DEBUG = 1
INFO = 2
WARN = 3
ERROR = 4
SILENT = 5


class LogMethod(str, Enum):
LOG = "log"
TRACE = "trace"
DEBUG = "debug"
INFO = "info"
WARN = "warn"
ERROR = "error"
14 changes: 12 additions & 2 deletions src/convertcom_sdk/events/event_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,11 +5,17 @@


class EventManager:
def __init__(self, config: Mapping[str, Any] | None = None) -> None:
def __init__(
self,
config: Mapping[str, Any] | None = None,
*,
logger_manager: Any | None = None,
) -> None:
config = dict(config or {})
self._listeners: dict[str, list[Callable[[Any, Any], None]]] = {}
self._deferred: dict[str, dict[str, Any]] = {}
self._mapper = config.get("mapper") or (lambda value: value)
self._logger_manager = logger_manager

def _event_key(self, event: Any) -> str:
return getattr(event, "value", event)
Expand All @@ -36,6 +42,10 @@ def fire(
event_key = self._event_key(event)
listeners = list(self._listeners.get(event_key, []))
for fn in listeners:
fn(self._mapper(args), err)
try:
fn(self._mapper(args), err)
except Exception as error:
if self._logger_manager:
self._logger_manager.error("EventManager.fire()", error)
if deferred and event_key not in self._deferred:
self._deferred[event_key] = {"args": args, "err": err}
3 changes: 3 additions & 0 deletions src/convertcom_sdk/logger/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
from .log_manager import LogManager

__all__ = ["LogManager"]
Loading