Skip to content

stream_configs

stream_configs

Stream configuration classes organized by platform.

This package provides typed stream configuration classes organized by major platform (Snowflake, AWS, Domo, Google, etc.). These classes follow the AccountConfig-style pattern and provide type-safe attribute access.

AWSAthena_StreamConfig dataclass

AWSAthena_StreamConfig(
    data_provider_type: str = "aws-athena",
    parent: Any = None,
    raw: dict = dict(),
    query: str = None,
    database_name: str = None,
    table_name: str = None,
)

Bases: StreamConfig_Base

AWS Athena stream configuration (typed).

AdobeAnalyticsV2_StreamConfig dataclass

AdobeAnalyticsV2_StreamConfig(
    data_provider_type: str = "adobe-analytics-v2",
    parent: Any = None,
    raw: dict = dict(),
    query: str = None,
    adobe_report_suite_id: str = None,
)

Bases: StreamConfig_Base

Adobe Analytics v2 stream configuration (typed).

AmazonAthenaHighBandwidth_StreamConfig dataclass

AmazonAthenaHighBandwidth_StreamConfig(
    data_provider_type: str = "amazon-athena-high-bandwidth",
    parent: Any = None,
    raw: dict = dict(),
    query: str = None,
    database_name: str = None,
)

Bases: StreamConfig_Base

Amazon Athena high bandwidth stream configuration (typed).

AmazonS3AssumeRole_StreamConfig dataclass

AmazonS3AssumeRole_StreamConfig(
    data_provider_type: str = "amazon_s3_assumerole",
    parent: Any = None,
    raw: dict = dict(),
    files_discovery: str = None,
)

Bases: StreamConfig_Base

Amazon S3 assume role stream configuration (typed).

ConformedProperty dataclass

ConformedProperty(
    name: str,
    mappings: dict[str, str],
    description: str = None,
    is_repr: bool = False,
)

Maps semantic property names to platform-specific config keys.

Defines a "conformed" or "semantic" property that exists conceptually across multiple data providers, even though each provider may use a different parameter name.

Attributes:

Name Type Description
name str

Semantic name of the property (e.g., "query", "database")

mappings dict[str, str]

Dict mapping provider_type to typed_config attribute name

description str

Human-readable description of what this property represents

supported_providers list[str]

List of provider types that support this property

is_repr bool

Whether to include this property in repr output (default: False)

Example

query = ConformedProperty( name="query", description="SQL query or data selection statement", mappings={ "snowflake": "query", "aws-athena": "query", "amazon-athena-high-bandwidth": "query", "postgresql": "query", }, is_repr=True # Show in repr )

supported_providers property

supported_providers: list[str]

List of provider types that support this property.

get_key_for_provider

get_key_for_provider(provider_type: str) -> str | None

Get the typed_config attribute name for a specific provider.

Parameters:

Name Type Description Default
provider_type str

Data provider type (e.g., "snowflake", "aws-athena")

required

Returns:

Type Description
str | None

Attribute name on the typed config class, or None if not supported

Example

query_prop.get_key_for_provider("snowflake") "query" query_prop.get_key_for_provider("unknown-provider") None

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_conformed.py
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
def get_key_for_provider(self, provider_type: str) -> str | None:
    """Get the typed_config attribute name for a specific provider.

    Args:
        provider_type: Data provider type (e.g., "snowflake", "aws-athena")

    Returns:
        Attribute name on the typed config class, or None if not supported

    Example:
        >>> query_prop.get_key_for_provider("snowflake")
        "query"
        >>> query_prop.get_key_for_provider("unknown-provider")
        None
    """
    return self.mappings.get(provider_type)

ConformedPropertyReprMixin

Mixin that adds smart repr for classes with conformed properties.

Usage

class DomoStream(ConformedPropertyReprMixin, DomoEntity): ...

DatasetCopy_StreamConfig dataclass

DatasetCopy_StreamConfig(
    data_provider_type: str = "dataset-copy",
    parent: Any = None,
    raw: dict = dict(),
    dataset_id: str = None,
)

Bases: StreamConfig_Base

Dataset copy stream configuration (typed).

Default_StreamConfig dataclass

Default_StreamConfig(
    data_provider_type: str = "default",
    parent: Any = None,
    raw: dict = dict(),
)

Bases: StreamConfig_Base

Default stream configuration for unknown provider types (typed).

This is a catch-all config that accepts any parameters. Used when the specific provider type is not recognized.

DomoCSV_StreamConfig dataclass

DomoCSV_StreamConfig(
    data_provider_type: str = "domo-csv",
    parent: Any = None,
    raw: dict = dict(),
    url: str = None,
)

Bases: StreamConfig_Base

Domo CSV stream configuration (typed).

GoogleSheets_StreamConfig dataclass

GoogleSheets_StreamConfig(
    data_provider_type: str = "google-sheets",
    parent: Any = None,
    raw: dict = dict(),
    spreadsheet_id_file_name: str = None,
)

Bases: StreamConfig_Base

Google Sheets stream configuration (typed).

GoogleSpreadsheets_StreamConfig dataclass

GoogleSpreadsheets_StreamConfig(
    data_provider_type: str = "google-spreadsheets",
    parent: Any = None,
    raw: dict = dict(),
    spreadsheet_id_file_name: str = None,
)

Bases: StreamConfig_Base

Google Spreadsheets stream configuration (typed).

PostgreSQL_StreamConfig dataclass

PostgreSQL_StreamConfig(
    data_provider_type: str = "postgresql",
    parent: Any = None,
    raw: dict = dict(),
    query: str = None,
    database_name: str = None,
    schema_name: str = None,
    host: str = None,
    port: str = None,
)

Bases: StreamConfig_Base

PostgreSQL stream configuration (typed).

Qualtrics_StreamConfig dataclass

Qualtrics_StreamConfig(
    data_provider_type: str = "qualtrics",
    parent: Any = None,
    raw: dict = dict(),
    qualtrics_survey_id: str = None,
)

Bases: StreamConfig_Base

Qualtrics stream configuration (typed).

SharePointOnline_StreamConfig dataclass

SharePointOnline_StreamConfig(
    data_provider_type: str = "sharepoint-online",
    parent: Any = None,
    raw: dict = dict(),
    site_url: str = None,
)

Bases: StreamConfig_Base

SharePoint Online stream configuration (typed).

SnowflakeKeyPairAuth_StreamConfig dataclass

SnowflakeKeyPairAuth_StreamConfig(
    data_provider_type: str = "snowflakekeypairauthentication",
    parent: Any = None,
    raw: dict = dict(),
    query: str = None,
    database_name: str = None,
    schema_name: str = None,
    warehouse_name: str = None,
    report_type: str = None,
    query_tag: str = None,
    fetch_size: str = None,
    update_mode: str = None,
    convert_timezone: str = None,
    cloud: str = None,
)

Bases: StreamConfig_Base

Snowflake keypair authentication stream configuration.

Includes additional fields for keypair auth (query tags, fetch size, etc.).

Example

config = SnowflakeKeyPairAuth_StreamConfig.from_dict({ ... "query": "SELECT * FROM table", ... "databaseName": "SA_PRD", ... "queryTag": "domoD2C123", ... "fetchSize": "1000" ... }) config.query_tag "domoD2C123"

Snowflake_StreamConfig dataclass

Snowflake_StreamConfig(
    data_provider_type: str = "snowflake",
    parent: Any = None,
    raw: dict = dict(),
    query: str = None,
    database_name: str = None,
    warehouse_name: str = None,
    schema_name: str = None,
)

Bases: StreamConfig_Base

Snowflake stream configuration (typed, follows AccountConfig pattern).

Provides type-safe access to Snowflake stream parameters.

Example

config = Snowflake_StreamConfig.from_dict({ ... "query": "SELECT * FROM table", ... "databaseName": "SA_PRD", ... "warehouseName": "COMPUTE_WH" ... }) config.query # Type-safe attribute access "SELECT * FROM table" config.database_name "SA_PRD"

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,
    )

create_stream_repr

create_stream_repr(
    stream_obj: Any,
    max_value_length: int = 50,
    max_total_length: int = 200,
    priority_props: list[str] | None = None,
    include_missing_mappings: bool = False,
) -> str

Create a custom repr string for DomoStream with conformed properties.

Only includes properties where is_repr=True in CONFORMED_PROPERTIES registry.

Parameters:

Name Type Description Default
stream_obj Any

The DomoStream instance

required
max_value_length int

Maximum length for individual property values

50
max_total_length int

Maximum total length of repr string

200
priority_props list[str] | None

List of properties to show first (defaults to common SQL properties)

None
include_missing_mappings bool

If True, append count of missing mappings

False

Returns:

Type Description
str

Formatted repr string

Example

stream = DomoStream(id="123", data_provider_key="snowflake") print(create_stream_repr(stream)) DomoStream(id='123', provider='snowflake', sql='SELECT...', database='SA_PRD')

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_repr.py
104
105
106
107
108
109
110
111
112
113
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
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
def create_stream_repr(  # noqa: C901
    stream_obj: Any,
    max_value_length: int = 50,
    max_total_length: int = 200,
    priority_props: list[str] | None = None,
    include_missing_mappings: bool = False,
) -> str:
    """Create a custom __repr__ string for DomoStream with conformed properties.

    Only includes properties where is_repr=True in CONFORMED_PROPERTIES registry.

    Args:
        stream_obj: The DomoStream instance
        max_value_length: Maximum length for individual property values
        max_total_length: Maximum total length of repr string
        priority_props: List of properties to show first (defaults to common SQL properties)
        include_missing_mappings: If True, append count of missing mappings

    Returns:
        Formatted repr string

    Example:
        >>> stream = DomoStream(id="123", data_provider_key="snowflake")
        >>> print(create_stream_repr(stream))
        DomoStream(id='123', provider='snowflake', sql='SELECT...', database='SA_PRD')
    """
    from ._conformed import CONFORMED_PROPERTIES

    priority_props = priority_props or ["sql", "database", "warehouse", "schema"]

    parts = []

    # Always include ID and provider
    parts.append(f"DomoStream(id='{stream_obj.id}'")

    if hasattr(stream_obj, "data_provider_name") and stream_obj.data_provider_name:
        parts.append(f"provider='{stream_obj.data_provider_name}'")
    elif hasattr(stream_obj, "data_provider_key") and stream_obj.data_provider_key:
        parts.append(f"provider='{stream_obj.data_provider_key}'")

    # Get property name mapping (registry name -> property name)
    property_map = {
        "query": "sql",
        "database": "database",
        "schema": "schema",
        "warehouse": "warehouse",
        "table": "table",
        "report_id": "report_id",
        "spreadsheet": "spreadsheet",
        "bucket": "bucket",
        "dataset_id": "dataset_id",
        "file_url": "file_url",
        "host": "host",
        "port": "port",
    }

    # Collect property values (ONLY where is_repr=True)
    prop_values = []

    # Check priority properties first
    for prop_name in priority_props:
        # Find the registry name for this property
        registry_name = next(
            (k for k, v in property_map.items() if v == prop_name), None
        )
        if not registry_name:
            continue

        # Check if this property is marked for repr
        conf_prop = CONFORMED_PROPERTIES.get(registry_name)
        if not conf_prop or not conf_prop.is_repr:
            continue

        if hasattr(stream_obj, prop_name):
            try:
                value = getattr(stream_obj, prop_name, None)
                if value is not None:
                    prop_values.append((prop_name, value, True))  # True = priority
            except DomoError:
                continue

    # Check other properties (ONLY where is_repr=True)
    for registry_name, conf_prop in CONFORMED_PROPERTIES.items():
        # Skip if not marked for repr
        if not conf_prop.is_repr:
            continue

        prop_name = property_map.get(registry_name)
        if not prop_name or prop_name in priority_props:
            continue

        if hasattr(stream_obj, prop_name):
            try:
                value = getattr(stream_obj, prop_name, None)
                if value is not None:
                    prop_values.append(
                        (prop_name, value, False)
                    )  # False = not priority
            except DomoError:
                continue

    # Format property values
    for prop_name, value, is_priority in prop_values:
        # Truncate long strings
        if isinstance(value, str) and len(value) > max_value_length:
            value = value[:max_value_length] + "..."

        parts.append(f"{prop_name}='{value}'")

        # Check if we're getting too long
        current_repr = ", ".join(parts) + ")"
        if len(current_repr) > max_total_length and not is_priority:
            parts.append("...")
            break

    # Optionally include missing mappings count
    if include_missing_mappings:
        missing = get_missing_mappings(stream_obj)
        if missing:
            parts.append(f"_missing_mappings={len(missing)}")

    return ", ".join(parts) + ")"

get_available_config_keys

get_available_config_keys(stream_obj: Any) -> list[str]

Get list of typed_config keys that are NOT mapped to conformed properties.

This helps identify gaps in conformed property mappings - keys that exist in the typed_config but don't have a corresponding conformed property.

Parameters:

Name Type Description Default
stream_obj Any

The DomoStream instance

required

Returns:

Type Description
list[str]

List of typed_config attribute names that aren't mapped to conformed properties

Example

stream.data_provider_key = "snowflake" stream._available_config_keys ['role', 'authenticator', 'private_key'] # Keys not yet mapped

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_repr.py
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
def get_available_config_keys(stream_obj: Any) -> list[str]:
    """Get list of typed_config keys that are NOT mapped to conformed properties.

    This helps identify gaps in conformed property mappings - keys that exist
    in the typed_config but don't have a corresponding conformed property.

    Args:
        stream_obj: The DomoStream instance

    Returns:
        List of typed_config attribute names that aren't mapped to conformed properties

    Example:
        >>> stream.data_provider_key = "snowflake"
        >>> stream._available_config_keys
        ['role', 'authenticator', 'private_key']  # Keys not yet mapped
    """
    from ._conformed import CONFORMED_PROPERTIES

    typed_config = getattr(stream_obj, "typed_config", None)
    if not typed_config:
        return []

    provider_key = getattr(stream_obj, "data_provider_key", None)
    if not provider_key:
        return []

    # Get all attribute names from typed_config (excluding private/magic)
    all_keys = [
        attr
        for attr in dir(typed_config)
        if not attr.startswith("_") and not callable(getattr(typed_config, attr))
    ]

    # Get all mapped keys for this provider
    mapped_keys = set()
    for conf_prop in CONFORMED_PROPERTIES.values():
        key = conf_prop.get_key_for_provider(provider_key)
        if key:
            mapped_keys.add(key)

    # Return keys that aren't mapped
    unmapped = [key for key in all_keys if key not in mapped_keys]

    # Filter out common base class attributes
    exclude = {
        "data_provider_type",
        "transport_type",
        "data_provider_key",
        "raw",
        "auth",
    }
    return [key for key in unmapped if key not in exclude]

get_conformed_properties_for_repr

get_conformed_properties_for_repr(
    stream_obj: Any,
) -> dict[str, Any]

Get dictionary of conformed properties and their values for a stream.

Parameters:

Name Type Description Default
stream_obj Any

The DomoStream instance

required

Returns:

Type Description
dict[str, Any]

Dict of {property_name: value} for all non-None conformed properties

Example

props = get_conformed_properties_for_repr(stream) props {'sql': 'SELECT * FROM table', 'database': 'SA_PRD', 'warehouse': 'COMPUTE_WH'}

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_repr.py
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
def get_conformed_properties_for_repr(stream_obj: Any) -> dict[str, Any]:
    """Get dictionary of conformed properties and their values for a stream.

    Args:
        stream_obj: The DomoStream instance

    Returns:
        Dict of {property_name: value} for all non-None conformed properties

    Example:
        >>> props = get_conformed_properties_for_repr(stream)
        >>> props
        {'sql': 'SELECT * FROM table', 'database': 'SA_PRD', 'warehouse': 'COMPUTE_WH'}
    """
    result = {}

    conformed_prop_names = [
        "sql",
        "database",
        "schema",
        "warehouse",
        "table",
        "report_id",
        "spreadsheet",
        "bucket",
        "dataset_id",
        "file_url",
        "host",
        "port",
    ]

    for prop_name in conformed_prop_names:
        if hasattr(stream_obj, prop_name):
            try:
                value = getattr(stream_obj, prop_name, None)
                if value is not None:
                    result[prop_name] = value
            except DomoError:
                continue

    return result

get_missing_mappings

get_missing_mappings(stream_obj: Any) -> list[str]

Get list of conformed properties that don't have mappings for this provider.

Parameters:

Name Type Description Default
stream_obj Any

The DomoStream instance

required

Returns:

Type Description
list[str]

List of property names that exist in registry but don't support this provider

Example

stream.data_provider_key = "google-sheets" get_missing_mappings(stream) ['query', 'database', 'warehouse'] # SQL properties not supported

Source code in src/crew_dcs/classes/DomoDataset/stream_configs/_repr.py
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
def get_missing_mappings(stream_obj: Any) -> list[str]:
    """Get list of conformed properties that don't have mappings for this provider.

    Args:
        stream_obj: The DomoStream instance

    Returns:
        List of property names that exist in registry but don't support this provider

    Example:
        >>> stream.data_provider_key = "google-sheets"
        >>> get_missing_mappings(stream)
        ['query', 'database', 'warehouse']  # SQL properties not supported
    """
    from ._conformed import CONFORMED_PROPERTIES

    missing = []
    provider_key = getattr(stream_obj, "data_provider_key", None)

    if not provider_key:
        return missing

    for prop_name, conf_prop in CONFORMED_PROPERTIES.items():
        # Skip properties not meant for repr
        if not conf_prop.is_repr:
            continue

        # Check if provider supports this property
        if provider_key not in conf_prop.supported_providers:
            missing.append(prop_name)

    return missing

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

Modules