Skip to content

_base

_base

Base classes and utilities for stream configurations.

This module contains: - StreamConfig_Base: Base class for typed stream configs (AccountConfig-style) - register_stream_config: Decorator for registering typed config classes

StreamConfig_Base dataclass

StreamConfig_Base(
    data_provider_type: str = None,
    parent: Any = None,
    raw: dict = dict(),
)

Bases: DomoBase

Base class for typed stream configurations.

Similar to DomoAccount_Config, this provides a consistent interface for handling stream execution parameters across different data providers.

from_dict classmethod

from_dict(obj: dict[str, Any], parent: Any = None)

Create StreamConfig from dictionary of config parameters.

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_base.py
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
@classmethod
def from_dict(cls, obj: dict[str, Any], parent: Any = None):
    """Create StreamConfig from dictionary of config parameters."""
    field_map = {}
    if "_field_map" in cls.__dataclass_fields__:
        field_map_field = cls.__dataclass_fields__["_field_map"]
        if (
            hasattr(field_map_field, "default_factory")
            and field_map_field.default_factory
        ):
            field_map = field_map_field.default_factory()

    init_kwargs = {}
    for api_key, api_value in obj.items():
        python_attr = field_map.get(api_key)
        if not python_attr:
            python_attr = _camel_to_snake(api_key)

        if python_attr in cls.__dataclass_fields__:
            init_kwargs[python_attr] = api_value

    return cls(parent=parent, raw=obj, **init_kwargs)

to_dict

to_dict(
    return_snake_case: bool = False, *, for_api: bool = True
) -> list[dict] | dict

Convert to list of config dictionaries for API submission or export.

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_base.py
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
def to_dict(
    self,
    return_snake_case: bool = False,
    *,
    for_api: bool = True,
) -> list[dict] | dict:
    """Convert to list of config dictionaries for API submission or export."""
    reverse_map = {v: k for k, v in self._field_map.items()}

    export_fields = self._fields_for_export or [
        f
        for f in self.__dataclass_fields__
        if not f.startswith("_")
        and f not in ["data_provider_type", "parent", "raw"]
    ]

    if for_api:
        result = []
        for attr_name in export_fields:
            value = getattr(self, attr_name, None)
            if value is not None:
                api_key = reverse_map.get(attr_name, _snake_to_camel(attr_name))
                result.append(
                    {"name": api_key, "type": "string", "value": str(value)}
                )
        return result

    result: dict[str, Any] = {}
    for attr_name in export_fields:
        value = getattr(self, attr_name, None)
        if value is None:
            continue
        key = (
            attr_name
            if return_snake_case
            else reverse_map.get(attr_name, _snake_to_camel(attr_name))
        )
        result[key] = value

    return result

Stream_CRUD_Error

Stream_CRUD_Error(
    operation: str,
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs
)

Bases: RouteError

Raised when stream create, update, delete, or execute operations fail.

Source code in src/crew_dcs/routes/dataset/stream.py
61
62
63
64
65
66
67
68
69
70
71
72
73
74
def __init__(
    self,
    operation: str,
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs,
):
    super().__init__(
        message=message or f"Stream {operation} operation failed",
        entity_id=stream_id,
        res=res,
        **kwargs,
    )

Stream_GET_Error

Stream_GET_Error(
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs
)

Bases: RouteError

Raised when stream retrieval operations fail.

Source code in src/crew_dcs/routes/dataset/stream.py
43
44
45
46
47
48
49
50
51
52
53
54
55
def __init__(
    self,
    stream_id: str | None = None,
    message: str | None = None,
    res=None,
    **kwargs,
):
    super().__init__(
        message=message or "Stream retrieval failed",
        entity_id=stream_id,
        res=res,
        **kwargs,
    )

register_stream_config

register_stream_config(data_provider_type: str)

Decorator to register a StreamConfig_Base subclass.

Parameters:

Name Type Description Default
data_provider_type str

The data provider type identifier (e.g., 'snowflake')

required
Example

@register_stream_config('snowflake') @dataclass class Snowflake_StreamConfig(StreamConfig_Base): data_provider_type: str = "snowflake" query: str = None database_name: str = None

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_base.py
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
def register_stream_config(data_provider_type: str):
    """Decorator to register a StreamConfig_Base subclass.

    Args:
        data_provider_type: The data provider type identifier (e.g., 'snowflake')

    Example:
        @register_stream_config('snowflake')
        @dataclass
        class Snowflake_StreamConfig(StreamConfig_Base):
            data_provider_type: str = "snowflake"
            query: str = None
            database_name: str = None
    """

    def decorator(cls: type[StreamConfig_Base]) -> type[StreamConfig_Base]:
        _CONFIG_REGISTRY[data_provider_type] = cls
        return cls

    return decorator