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
7 changes: 6 additions & 1 deletion src/convertcom_sdk/__init__.py
Original file line number Diff line number Diff line change
@@ -1,22 +1,27 @@
from .api.api_manager import ApiManager
from .data.data_manager import DataManager
from .data.data_store_manager import DataStoreManager
from .events.event_manager import EventManager
from .experience.experience_manager import ExperienceManager
from .features.feature_manager import FeatureManager
from .segments.segments_manager import SegmentsManager
from .bucketing.bucketing_manager import BucketingAllocation, BucketingManager
from .enums import BucketingError, FeatureStatus, RuleError
from .enums import BucketingError, FeatureStatus, RuleError, SystemEvents
from .rules.rule_manager import RuleManager

__all__ = [
"ApiManager",
"BucketingAllocation",
"BucketingError",
"BucketingManager",
"DataManager",
"DataStoreManager",
"EventManager",
"ExperienceManager",
"FeatureManager",
"FeatureStatus",
"RuleError",
"RuleManager",
"SegmentsManager",
"SystemEvents",
]
3 changes: 3 additions & 0 deletions src/convertcom_sdk/api/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
from .api_manager import ApiManager

__all__ = ["ApiManager"]
235 changes: 235 additions & 0 deletions src/convertcom_sdk/api/api_manager.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,235 @@
from __future__ import annotations

import copy
import os
import threading
from collections.abc import Callable, Mapping
from dataclasses import dataclass, field
from typing import Any

from convertcom_sdk.enums import SystemEvents
from convertcom_sdk.utils.http_client import HttpResponse, request as http_request

DEFAULT_HEADERS = {"Content-Type": "application/json"}
DEFAULT_BATCH_SIZE = 10
DEFAULT_RELEASE_INTERVAL = 10000
DEFAULT_CONFIG_ENDPOINT = os.getenv("CONFIG_ENDPOINT", "")
DEFAULT_TRACK_ENDPOINT = os.getenv("TRACK_ENDPOINT", "")
Comment on lines +16 to +17

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

The default values for DEFAULT_CONFIG_ENDPOINT and DEFAULT_TRACK_ENDPOINT are empty strings if the respective environment variables are not set. This can lead to request calls with an invalid base_url, which will likely result in runtime errors when attempting to make HTTP requests. It's safer to either raise an error during initialization if these critical endpoints are not configured or provide a more robust default that prevents silent failures.



@dataclass
class VisitorQueueItem:
visitorId: str

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

The field visitorId uses camelCase which is inconsistent with Python's snake_case convention as recommended by PEP 8 for variable names. Please consider renaming it to visitor_id for better readability and adherence to style guidelines.

Suggested change
visitorId: str
visitor_id: str

events: list[dict[str, Any]] = field(default_factory=list)
segments: dict[str, Any] | None = None


class VisitorsQueue:
def __init__(self) -> None:
self.length = 0
self.items: list[VisitorQueueItem] = []

def push(
self,
visitor_id: str,
event_request: Mapping[str, Any],
segments: Mapping[str, Any] | None = None,
) -> None:
for item in self.items:
if item.visitorId == visitor_id:
item.events.append(dict(event_request))
self.length += 1
return
self.items.append(
VisitorQueueItem(
visitorId=visitor_id,
events=[dict(event_request)],
segments=dict(segments) if segments else None,
)
)
self.length += 1

def reset(self) -> None:
self.items = []
self.length = 0


class ApiManager:
def __init__(
self,
config: Mapping[str, Any] | None = None,
*,
event_manager: Any | None = None,

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

The type hint for event_manager is Any | None. For better type safety and clarity, it would be more precise to use EventManager | None since EventManager is a defined class within the SDK.

Suggested change
event_manager: Any | None = None,
event_manager: EventManager | None = None,

request_sender: Callable[..., HttpResponse] | None = None,
) -> None:
config = dict(config or {})
self._lock = threading.RLock()
self._event_manager = event_manager
self._request_sender = request_sender or http_request
endpoints = ((config.get("api") or {}).get("endpoint")) or {}
network = config.get("network") or {}
events = config.get("events") or {}
self._data = dict(config.get("data") or {})
self._enrich_data = not bool(config.get("dataStore"))
self._environment = config.get("environment")
self._mapper = config.get("mapper") or (lambda value: value)
self._config_endpoint = endpoints.get("config") or DEFAULT_CONFIG_ENDPOINT
self._track_endpoint = endpoints.get("track") or DEFAULT_TRACK_ENDPOINT
self._default_headers = dict(DEFAULT_HEADERS)
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._account_id = self._data.get("account_id")
self._project_id = (self._data.get("project") or {}).get("id")
self._sdk_key = config.get("sdkKey") or self._build_default_sdk_key()
if config.get("sdkKeySecret"):
self._default_headers["Authorization"] = (
f"Bearer {config['sdkKeySecret']}"
)
Comment on lines +87 to +89

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

Accessing config['sdkKeySecret'] directly after checking config.get("sdkKeySecret") can lead to a KeyError if config.get("sdkKeySecret") returns a falsy value (e.g., 0 or False) but the key is not actually present. It's safer to store the result of config.get("sdkKeySecret") in a variable and use that variable.

Suggested change
self._default_headers["Authorization"] = (
f"Bearer {config['sdkKeySecret']}"
)
sdk_key_secret = config.get("sdkKeySecret")
if sdk_key_secret:
self._default_headers["Authorization"] = f"Bearer {sdk_key_secret}"

self._tracking_event = {
"enrichData": self._enrich_data,
"accountId": self._account_id,
"projectId": self._project_id,
"visitors": [],
}
self._tracking_enabled = bool(network.get("tracking"))
self._tracking_source = network.get("source") or "python-sdk"
self._cache_level = network.get("cacheLevel")
self._requests_queue = VisitorsQueue()
self._requests_queue_timer: threading.Timer | None = None

def _build_default_sdk_key(self) -> str:
if self._account_id and self._project_id:
return f"{self._account_id}/{self._project_id}"
return ""

def request(
self,
method: str,
path: Mapping[str, str],
data: Mapping[str, Any] | None = None,
headers: Mapping[str, Any] | None = None,
) -> HttpResponse:
request_headers = {
**self._default_headers,
**dict(headers or {}),
}
return self._request_sender(
method=method,
base_url=path["base"],
route=path["route"],
headers=request_headers,
data=dict(data or {}),
)

def enqueue(
self,
visitor_id: str,
event_request: Mapping[str, Any],
segments: Mapping[str, Any] | None = None,
) -> None:
with self._lock:
self._requests_queue.push(visitor_id, event_request, segments)
queue_length = self._requests_queue.length
tracking_enabled = self._tracking_enabled
if tracking_enabled:
if queue_length == 1:
self.start_queue()
elif queue_length == self.batch_size:
self.release_queue("size")

def release_queue(self, reason: str | None = None) -> HttpResponse | None:
with self._lock:
if not self._requests_queue.length:
return None
self.stop_queue()
payload = dict(self._tracking_event)
payload["visitors"] = [
{
"visitorId": item.visitorId,
"events": copy.deepcopy(item.events),
**({"segments": copy.deepcopy(item.segments)} if item.segments else {}),
}
for item in self._requests_queue.items
]
payload["source"] = self._tracking_source
try:
response = self.request(
"post",
{
"base": self._track_endpoint.replace(
"[project_id]", str(self._project_id or "")
),
"route": f"/track/{self._sdk_key}",
},
self._mapper(payload),
)
except Exception as error:

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

Catching a generic Exception can mask other programming errors and make debugging more difficult. It's generally better to catch more specific exceptions. In this context, HttpError (which is raised by _request_sender) would be a more appropriate exception to catch for network-related issues.

Suggested change
except Exception as error:
except HttpError as error:

self.start_queue()
if self._event_manager:
self._event_manager.fire(
SystemEvents.API_QUEUE_RELEASED,
{"reason": reason},
error,
)
return None

with self._lock:
self._requests_queue.reset()
if self._event_manager:
self._event_manager.fire(
SystemEvents.API_QUEUE_RELEASED,
{"reason": reason, "result": response, "visitors": payload["visitors"]},
None,
)
return response

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 enable_tracking(self) -> None:
self._tracking_enabled = True
self.release_queue("trackingEnabled")

def disable_tracking(self) -> None:
self._tracking_enabled = False

def set_data(self, data: Mapping[str, Any] | None) -> None:
self._data = dict(data or {})
self._account_id = self._data.get("account_id")
self._project_id = (self._data.get("project") or {}).get("id")
self._tracking_event["accountId"] = self._account_id
self._tracking_event["projectId"] = self._project_id
if not self._sdk_key:
self._sdk_key = self._build_default_sdk_key()
Comment on lines +219 to +220

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

The logic for rebuilding _sdk_key in set_data only triggers if _sdk_key is falsy. If _sdk_key was initially provided (and thus truthy) but _account_id or _project_id are subsequently changed via set_data, the _sdk_key might become stale if it was originally derived from these IDs. To ensure consistency, _sdk_key should be rebuilt unconditionally if _account_id or _project_id are updated, or if sdkKey is not explicitly provided in the initial config.

        self._sdk_key = config.get("sdkKey") or self._build_default_sdk_key()


def get_config(self) -> dict[str, Any]:
query = "?" if self._cache_level == "low" or self._environment else ""
if self._environment:
query += f"environment={self._environment}"
if self._cache_level == "low":
query += "_conv_low_cache=1"
response = self.request(
"get",
{
"base": self._config_endpoint,
"route": f"/config/{self._sdk_key}{query}",
},
)
return dict(response.data or {})
13 changes: 13 additions & 0 deletions src/convertcom_sdk/enums.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,3 +27,16 @@ class SegmentsKeys(str, Enum):
CAMPAIGN = "campaign"
VISITOR_TYPE = "visitorType"
CUSTOM_SEGMENTS = "customSegments"


class SystemEvents(str, Enum):
READY = "ready"
CONFIG_UPDATED = "config.updated"
API_QUEUE_RELEASED = "api.queue.released"
BUCKETING = "bucketing"
CONVERSION = "conversion"
SEGMENTS = "segments"
LOCATION_ACTIVATED = "location.activated"
LOCATION_DEACTIVATED = "location.deactivated"
AUDIENCES = "audiences"
DATA_STORE_QUEUE_RELEASED = "datastore.queue.released"
3 changes: 3 additions & 0 deletions src/convertcom_sdk/events/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
from .event_manager import EventManager

__all__ = ["EventManager"]
35 changes: 35 additions & 0 deletions src/convertcom_sdk/events/event_manager.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,35 @@
from __future__ import annotations

from collections.abc import Callable, Mapping
from typing import Any


class EventManager:
def __init__(self, config: Mapping[str, 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)

def on(self, event: str, fn: Callable[[Any, Any], None]) -> None:
self._listeners.setdefault(str(event), []).append(fn)
if str(event) in self._deferred:
deferred = self._deferred[str(event)]
self.fire(str(event), deferred.get("args"), deferred.get("err"))

def remove_listeners(self, event: str) -> None:
self._listeners.pop(str(event), None)
self._deferred.pop(str(event), None)

def fire(
self,
event: str,
args: Any = None,
err: Any = None,
deferred: bool = False,
) -> None:
listeners = list(self._listeners.get(str(event), []))
for fn in listeners:
fn(self._mapper(args), err)
if deferred and str(event) not in self._deferred:
self._deferred[str(event)] = {"args": args, "err": err}
3 changes: 3 additions & 0 deletions src/convertcom_sdk/utils/__init__.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,12 @@
from .comparisons import DEFAULT_COMPARISON_PROCESSOR
from .http_client import HttpError, HttpResponse
from .type_utils import cast_type
from .string_utils import camel_case, generate_hash, is_numeric, to_number

__all__ = [
"DEFAULT_COMPARISON_PROCESSOR",
"HttpError",
"HttpResponse",
"camel_case",
"cast_type",
"generate_hash",
Expand Down
Loading