Skip to content

DomoDataflow

DomoDataflow

DomoDataflow dataclass

DomoDataflow(
    auth: DomoAuth,
    id: str,
    raw: dict,
    name: str | None = None,
    owner: str | None = None,
    description: str | None = None,
    tags: list[str] | None = None,
    version_id: int | None = None,
    version_number: int | None = None,
    versions: list[dict[str, Any]] | None = None,
    jupyter_workspace_config: dict | None = None,
    Definition: DomoDataflow_Definition | None = None,
    History: DomoDataflow_History | None = None,
    JupyterWorkspace: DomoJupyterWorkspace | None = None,
)

Bases: DomoEntity_w_Lineage

Actions property

Actions

Shorthand for self.Definition.Actions.

Canvas property

Canvas

Shorthand for self.Definition.Canvas.

entity_name property

entity_name: str

Get the display name for this dataflow.

Dataflows use the 'name' field as their display name.

Returns:

Type Description
str

Dataflow name, or dataflow ID as fallback

optimize_canvas async

optimize_canvas(
    save: bool = True,
    add_sections: bool = True,
    debug_api: bool = False,
    session: AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs
)

Fetch definition, optimize canvas layout, and optionally save.

Convenience method that ensures the definition is loaded, runs the layout optimizer, and optionally saves the result back to the Domo API.

Parameters:

Name Type Description Default
save bool

Whether to save after optimizing (default True).

True
add_sections bool

Whether to create auto-sections (default True).

True
debug_api bool

Enable debug logging for API calls.

False
session AsyncClient | None

Optional httpx client session.

None
context RouteContext | None

Optional RouteContext for the request.

None
**context_kwargs

Additional context parameters.

{}

Returns:

Type Description

The modified CanvasElements instance.

Example::

canvas = await dataflow.optimize_canvas(save=True, add_sections=True)
Source code in src/crew_dcs/classes/DomoDataflow/core.py
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
async def optimize_canvas(
    self,
    save: bool = True,
    add_sections: bool = True,
    debug_api: bool = False,
    session: httpx.AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
):
    """Fetch definition, optimize canvas layout, and optionally save.

    Convenience method that ensures the definition is loaded, runs
    the layout optimizer, and optionally saves the result back to
    the Domo API.

    Args:
        save: Whether to save after optimizing (default True).
        add_sections: Whether to create auto-sections (default True).
        debug_api: Enable debug logging for API calls.
        session: Optional httpx client session.
        context: Optional RouteContext for the request.
        **context_kwargs: Additional context parameters.

    Returns:
        The modified CanvasElements instance.

    Example::

        canvas = await dataflow.optimize_canvas(save=True, add_sections=True)
    """
    if self.Definition is None:
        raise RuntimeError("Cannot optimize canvas: Definition not initialized")

    # Ensure definition is loaded
    await self.Definition.get(
        debug_api=debug_api,
        session=session,
        context=context,
        **context_kwargs,
    )

    # Run optimizer
    from .layout_optimizer import optimize_canvas as _optimize

    canvas = _optimize(self.Definition, add_sections=add_sections)

    # Optionally save
    if save:
        await self.Definition.save(
            debug_api=debug_api,
            session=session,
            version_note="Optimized canvas layout",
            context=context,
            **context_kwargs,
        )

    return canvas

DomoDataflowNotFoundError

DomoDataflowNotFoundError(cls_instance, search_name: str)

Bases: ClassError

Exception raised when dataflow search operations return no results.

Source code in src/crew_dcs/classes/DomoDataflow/exceptions.py
17
18
19
def __init__(self, cls_instance, search_name: str):
    message = f"No dataflow found matching '{search_name}'"
    super().__init__(cls_instance=cls_instance, message=message)

DomoDataflow_ActionResult dataclass

DomoDataflow_ActionResult(
    id: str,
    type: str = None,
    name: str = None,
    is_success: bool = None,
    rows_processed: int = None,
    begin_time: datetime = None,
    end_time: datetime = None,
    duration_in_sec: int = None,
)

Result of an action execution from dataflow history.

DomoDataflow_Actions dataclass

DomoDataflow_Actions(
    auth: DomoAuth,
    dataflow: DomoDataflowProtocol,
    dataflow_id: str,
    actions: list[DomoDataflow_Action_Base] = list(),
)

Manager class for dataflow actions.

This class wraps the actions list from a dataflow and provides computed properties for filtering and analysis.

Attributes:

Name Type Description
auth DomoAuth

DomoAuth instance from parent dataflow

dataflow DomoDataflowProtocol

Reference to parent DomoDataflow

dataflow_id str

ID of the parent dataflow

actions list[DomoDataflow_Action_Base]

List of action objects

Example

dataflow = await DomoDataflow.get_by_id(auth=auth, dataflow_id=123) await dataflow.Actions.get() print(f"Has data science tiles: {dataflow.Actions.has_datascience_tiles}") print(f"Input datasets: {dataflow.Actions.input_datasets}")

action_type_counts property

action_type_counts: dict[str, int]

Get count of actions by action type.

Returns:

Type Description
dict[str, int]

Dictionary mapping action_type to count

Example

counts = dataflow.Actions.action_type_counts for action_type, count in counts.items(): ... print(f"{action_type}: {count}")

column_relationships property

column_relationships: set[ColumnRelationship]

Aggregate column-level lineage across all actions.

Collects ColumnRelationships from every action's column_relationships property and returns the union.

Example

rels = dataflow.Actions.column_relationships print(f"Found {len(rels)} column relationships")

datascience_tiles property

datascience_tiles: list[DomoDataflow_Action_Base]

Get list of data science actions.

Returns:

Type Description
list[DomoDataflow_Action_Base]

Filtered list of data science actions

Example

ds_tiles = dataflow.Actions.datascience_tiles print(f"Found {len(ds_tiles)} data science tiles")

disabled_actions property

disabled_actions: list[DomoDataflow_Action_Base]

Get list of disabled actions.

Returns:

Type Description
list[DomoDataflow_Action_Base]

Filtered list of disabled actions

Example

disabled = dataflow.Actions.disabled_actions print(f"Found {len(disabled)} disabled actions")

has_datascience_tiles property

has_datascience_tiles: bool

Check if the dataflow has any data science tiles.

Returns:

Type Description
bool

True if any action is a data science tile, False otherwise

Example

if dataflow.Actions.has_datascience_tiles: ... print("This dataflow uses data science operations")

input_datasets property

input_datasets: list[dict[str, Any]]

Get list of input datasets (LoadFromVault actions).

Returns:

Type Description
list[dict[str, Any]]

List of dicts with 'id' and 'name' keys for each input dataset

Example

for ds in dataflow.Actions.input_datasets: ... print(f"Input: {ds['name']} ({ds['id']})")

output_datasets property

output_datasets: list[dict[str, Any]]

Get list of output datasets (PublishToVault/WriteToVault actions).

Returns:

Type Description
list[dict[str, Any]]

List of dicts with 'id' and 'name' keys for each output dataset

Example

for ds in dataflow.Actions.output_datasets: ... print(f"Output: {ds['name']} ({ds['id']})")

tile_type_counts property

tile_type_counts: dict[str, int]

Get count of actions by tile type.

Returns:

Type Description
dict[str, int]

Dictionary mapping tile_type to count

Example

counts = dataflow.Actions.tile_type_counts for tile_type, count in counts.items(): ... print(f"{tile_type}: {count}")

topo_sorted property

topo_sorted: list[DomoDataflow_Action_Base]

Return actions in topological (dependency) order.

Uses the same Kahn's-algorithm sort as convert_entire_workflow_to_sql_tiles so callers get a stable, reproducible ordering that respects dependsOn edges.

Example

for action in dataflow.Actions.topo_sorted: ... print(action.name)

add_action

add_action(obj: dict) -> DomoDataflow_Action_Base

Create a typed action from a raw dict and append it to this manager.

Uses polymorphic dispatch on DomoDataflow_Action_Base.from_dict() so the returned object is the correct concrete subclass.

Parameters:

Name Type Description Default
obj dict

Raw action dict (e.g. from to_action_dict())

required

Returns:

Type Description
DomoDataflow_Action_Base

The newly created typed action (already appended to self.actions)

Example

new_dict = existing_pub.to_action_dict(upstream_id=..., output_name=...) new_action = defn.Actions.add_action(new_dict) tile = new_action.to_canvas_tile()

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
def add_action(self, obj: dict) -> DomoDataflow_Action_Base:
    """Create a typed action from a raw dict and append it to this manager.

    Uses polymorphic dispatch on DomoDataflow_Action_Base.from_dict() so the
    returned object is the correct concrete subclass.

    Args:
        obj: Raw action dict (e.g. from to_action_dict())

    Returns:
        The newly created typed action (already appended to self.actions)

    Example:
        >>> new_dict = existing_pub.to_action_dict(upstream_id=..., output_name=...)
        >>> new_action = defn.Actions.add_action(new_dict)
        >>> tile = new_action.to_canvas_tile()
    """
    action = DomoDataflow_Action_Base.from_dict(obj, all_actions=self.actions)
    self.actions.append(action)
    return action

convert_entire_workflow_to_sql

convert_entire_workflow_to_sql() -> list[dict[str, Any]]

Alias of convert_entire_workflow_to_sql_tiles().

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
540
541
542
def convert_entire_workflow_to_sql(self) -> list[dict[str, Any]]:
    """Alias of convert_entire_workflow_to_sql_tiles()."""
    return self.convert_entire_workflow_to_sql_tiles()

convert_entire_workflow_to_sql_tiles

convert_entire_workflow_to_sql_tiles() -> (
    list[dict[str, Any]]
)

Convert the full workflow into ordered SQL tile approximations.

Returns:

Type Description
list[dict[str, Any]]

List of tile conversion payloads with SQL, status, and dependencies.

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
530
531
532
533
534
535
536
537
538
def convert_entire_workflow_to_sql_tiles(self) -> list[dict[str, Any]]:
    """Convert the full workflow into ordered SQL tile approximations.

    Returns:
        List of tile conversion payloads with SQL, status, and dependencies.
    """
    if not self.actions:
        return []
    return convert_entire_workflow_to_sql_tiles(self.actions)

convert_tile_to_sql

convert_tile_to_sql(
    action_id: str,
) -> dict[str, Any] | None

Convert a single tile/action to SQL with dependency context.

Parameters:

Name Type Description Default
action_id str

The tile/action ID to convert.

required

Returns:

Type Description
dict[str, Any] | None

A conversion dict for the matching tile, or None if not found.

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
515
516
517
518
519
520
521
522
523
524
525
526
527
528
def convert_tile_to_sql(self, action_id: str) -> dict[str, Any] | None:
    """Convert a single tile/action to SQL with dependency context.

    Args:
        action_id: The tile/action ID to convert.

    Returns:
        A conversion dict for the matching tile, or None if not found.
    """
    conversions = self.convert_entire_workflow_to_sql_tiles()
    for row in conversions:
        if row.get("tile_id") == action_id:
            return row
    return None

from_parent classmethod

from_parent(
    parent: DomoDataflowProtocol,
    actions: list[DomoDataflow_Action_Base] | None = None,
) -> DomoDataflow_Actions

Create an Actions manager from a parent dataflow.

Parameters:

Name Type Description Default
parent DomoDataflowProtocol

Parent DomoDataflow instance

required
actions list[DomoDataflow_Action_Base] | None

Optional initial list of actions

None

Returns:

Type Description
DomoDataflow_Actions

DomoDataflow_Actions instance

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
@classmethod
def from_parent(
    cls,
    parent: DomoDataflowProtocol,
    actions: list[DomoDataflow_Action_Base] | None = None,
) -> DomoDataflow_Actions:
    """Create an Actions manager from a parent dataflow.

    Args:
        parent: Parent DomoDataflow instance
        actions: Optional initial list of actions

    Returns:
        DomoDataflow_Actions instance
    """
    return cls(
        auth=parent.auth,
        dataflow=parent,
        dataflow_id=parent.id,
        actions=actions or [],
    )

get async

get(
    session: AsyncClient | None = None,
    debug_api: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoDataflow_Action_Base]

Get actions from the dataflow definition.

This method fetches the latest dataflow definition and populates the actions list.

Parameters:

Name Type Description Default
session AsyncClient | None

Optional httpx client session

None
debug_api bool

Enable debug logging for API calls

False
context RouteContext | None

Optional RouteContext for the request

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description
list[DomoDataflow_Action_Base]

List of action objects

Example

actions = await dataflow.Actions.get() for action in actions: ... print(f"{action.name}: {action.action_type}")

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
 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
102
103
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
async def get(
    self,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoDataflow_Action_Base]:
    """Get actions from the dataflow definition.

    This method fetches the latest dataflow definition and populates
    the actions list.

    Args:
        session: Optional httpx client session
        debug_api: Enable debug logging for API calls
        context: Optional RouteContext for the request
        **context_kwargs: Additional context parameters

    Returns:
        List of action objects

    Example:
        >>> actions = await dataflow.Actions.get()
        >>> for action in actions:
        ...     print(f"{action.name}: {action.action_type}")
    """
    # Get the dataflow definition (this populates dataflow.raw)
    # Use the parent dataflow's Definition manager get method

    # If this Actions manager is part of a Definition, use its get method
    # Otherwise, call the definition's get directly
    if hasattr(self.dataflow, "Definition") and self.dataflow.Definition:
        await self.dataflow.Definition.get(
            context=context,
            session=session,  # type: ignore
            debug_api=debug_api,
        )
    else:
        # Fallback: directly fetch and update raw
        from ....routes import dataflow as dataflow_routes

        res = await dataflow_routes.get_dataflow_by_id(
            auth=self.auth,
            dataflow_id=self.dataflow_id,
            context=context,
        )
        if res.is_success:
            self.dataflow.raw = res.response

    # Parse actions from the raw definition
    if self.dataflow.raw and self.dataflow.raw.get("actions"):
        self.actions = [
            DomoDataflow_Action_Base.from_dict(
                action_dict, all_actions=self.actions
            )
            for action_dict in self.dataflow.raw["actions"]
        ]

    return self.actions

get_action_by_id

get_action_by_id(
    action_id: str,
) -> DomoDataflow_Action_Base | None

Get an action by its ID.

Parameters:

Name Type Description Default
action_id str

ID of the action to find

required

Returns:

Type Description
DomoDataflow_Action_Base | None

The action object, or None if not found

Example

action = dataflow.Actions.get_action_by_id("abc123") if action: ... print(f"Found action: {action.name}")

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
def get_action_by_id(self, action_id: str) -> DomoDataflow_Action_Base | None:
    """Get an action by its ID.

    Args:
        action_id: ID of the action to find

    Returns:
        The action object, or None if not found

    Example:
        >>> action = dataflow.Actions.get_action_by_id("abc123")
        >>> if action:
        ...     print(f"Found action: {action.name}")
    """
    for action in self.actions:
        if action.id == action_id:
            return action
    return None

get_actions_by_type

get_actions_by_type(
    action_type: str,
) -> list[DomoDataflow_Action_Base]

Get all actions of a specific type.

Parameters:

Name Type Description Default
action_type str

The action type to filter by (e.g., "Filter", "GroupBy")

required

Returns:

Type Description
list[DomoDataflow_Action_Base]

List of actions matching the type

Example

filters = dataflow.Actions.get_actions_by_type("Filter") print(f"Found {len(filters)} filter actions")

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
333
334
335
336
337
338
339
340
341
342
343
344
345
346
def get_actions_by_type(self, action_type: str) -> list[DomoDataflow_Action_Base]:
    """Get all actions of a specific type.

    Args:
        action_type: The action type to filter by (e.g., "Filter", "GroupBy")

    Returns:
        List of actions matching the type

    Example:
        >>> filters = dataflow.Actions.get_actions_by_type("Filter")
        >>> print(f"Found {len(filters)} filter actions")
    """
    return [action for action in self.actions if action.action_type == action_type]

get_script_content

get_script_content(action_id: str) -> str | None

Extract script content from Python, R, or SQL script tiles.

Parameters:

Name Type Description Default
action_id str

ID of the action to extract script from

required

Returns:

Type Description
str | None

Script content as string, or None if not a script tile or script not found

Example

script = dataflow.Actions.get_script_content("abc123") if script: ... print(script)

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
def get_script_content(self, action_id: str) -> str | None:
    """Extract script content from Python, R, or SQL script tiles.

    Args:
        action_id: ID of the action to extract script from

    Returns:
        Script content as string, or None if not a script tile or script not found

    Example:
        >>> script = dataflow.Actions.get_script_content("abc123")
        >>> if script:
        ...     print(script)
    """
    action = self.get_action_by_id(action_id)
    if not action:
        return None

    # Check for direct script field
    if hasattr(action, "script") and action.script:  # type: ignore
        return action.script  # type: ignore

    # Check settings
    if action.settings:
        script = (
            action.settings.get("script")
            or action.settings.get("code")
            or action.settings.get("pythonScript")
            or action.settings.get("rScript")
            or action.settings.get("sql")
            or action.settings.get("query")
            or action.settings.get("sqlScript")
        )
        if script:
            return script

    # Check raw data
    if action.raw:
        return action.raw.get("script")

    return None

remove_action

remove_action(
    action_id: str,
) -> DomoDataflow_Action_Base | None

Remove an action by its ID.

Parameters:

Name Type Description Default
action_id str

ID of the action to remove.

required

Returns:

Type Description
DomoDataflow_Action_Base | None

The removed action, or None if not found.

Example::

removed = defn.Actions.remove_action("MergeJoin-abc123")
Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
def remove_action(self, action_id: str) -> DomoDataflow_Action_Base | None:
    """Remove an action by its ID.

    Args:
        action_id: ID of the action to remove.

    Returns:
        The removed action, or None if not found.

    Example::

        removed = defn.Actions.remove_action("MergeJoin-abc123")
    """
    for i, action in enumerate(self.actions):
        if action.id == action_id:
            return self.actions.pop(i)
    return None

remove_actions_by_type

remove_actions_by_type(
    action_type: str,
) -> list[DomoDataflow_Action_Base]

Remove all actions of a given type.

Idempotent — safe to call repeatedly (returns empty list on second call).

Parameters:

Name Type Description Default
action_type str

The action type string to remove (e.g., "SplitJoin").

required

Returns:

Type Description
list[DomoDataflow_Action_Base]

List of removed actions.

Example::

removed = defn.Actions.remove_actions_by_type("SplitJoin")
Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
def remove_actions_by_type(
    self, action_type: str
) -> list[DomoDataflow_Action_Base]:
    """Remove all actions of a given type.

    Idempotent — safe to call repeatedly (returns empty list on second call).

    Args:
        action_type: The action type string to remove (e.g., "SplitJoin").

    Returns:
        List of removed actions.

    Example::

        removed = defn.Actions.remove_actions_by_type("SplitJoin")
    """
    keep, removed = [], []
    for action in self.actions:
        (removed if action.action_type == action_type else keep).append(action)
    self.actions = keep
    return removed

remove_actions_where

remove_actions_where(
    predicate,
) -> list[DomoDataflow_Action_Base]

Remove all actions matching a predicate.

Idempotent — safe to call repeatedly.

Parameters:

Name Type Description Default
predicate

Callable accepting a DomoDataflow_Action_Base and returning True for actions to remove.

required

Returns:

Type Description
list[DomoDataflow_Action_Base]

List of removed actions.

Example::

removed = defn.Actions.remove_actions_where(
    lambda a: a.action_type == "PublishToVault" and a.name == "MetaData_Pages_NoStats"
)
Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
def remove_actions_where(self, predicate) -> list[DomoDataflow_Action_Base]:
    """Remove all actions matching a predicate.

    Idempotent — safe to call repeatedly.

    Args:
        predicate: Callable accepting a DomoDataflow_Action_Base and
            returning True for actions to remove.

    Returns:
        List of removed actions.

    Example::

        removed = defn.Actions.remove_actions_where(
            lambda a: a.action_type == "PublishToVault" and a.name == "MetaData_Pages_NoStats"
        )
    """
    keep, removed = [], []
    for action in self.actions:
        (removed if predicate(action) else keep).append(action)
    self.actions = keep
    return removed

to_erd

to_erd(
    *,
    title: str | None = None,
    trace_from_entity: str | None = None,
    trace_from_column: str | None = None
) -> MermaidERDiagram

Generate a Mermaid ER diagram from column relationships.

Collects column relationships from all actions and converts them to an ER diagram showing column-level lineage between tiles.

Parameters:

Name Type Description Default
title str | None

Optional diagram title (defaults to dataflow name)

None
trace_from_entity str | None

If provided, trace only the upstream chain from this entity (e.g., a tile ID)

None
trace_from_column str | None

If provided, trace only the upstream chain for this specific column (e.g., "Revenue")

None

Returns:

Type Description
MermaidERDiagram

MermaidERDiagram with entities and relationships

Example
Full dataflow ERD

diagram = dataflow.Actions.to_erd()

Trace a specific column upstream

diagram = dataflow.Actions.to_erd( ... trace_from_entity="publish-1", ... trace_from_column="Revenue", ... )

Source code in src/crew_dcs/classes/DomoDataflow/action/manager.py
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
def to_erd(
    self,
    *,
    title: str | None = None,
    trace_from_entity: str | None = None,
    trace_from_column: str | None = None,
) -> MermaidERDiagram:
    """Generate a Mermaid ER diagram from column relationships.

    Collects column relationships from all actions and converts them
    to an ER diagram showing column-level lineage between tiles.

    Args:
        title: Optional diagram title (defaults to dataflow name)
        trace_from_entity: If provided, trace only the upstream chain
                          from this entity (e.g., a tile ID)
        trace_from_column: If provided, trace only the upstream chain
                          for this specific column (e.g., "Revenue")

    Returns:
        MermaidERDiagram with entities and relationships

    Example:
        >>> # Full dataflow ERD
        >>> diagram = dataflow.Actions.to_erd()
        >>> # Trace a specific column upstream
        >>> diagram = dataflow.Actions.to_erd(
        ...     trace_from_entity="publish-1",
        ...     trace_from_column="Revenue",
        ... )
    """
    from ....integrations.graphs.mermaid import ColumnRelationshipERConverter

    rels = self.column_relationships
    entity_names = {a.id: a.name or a.id for a in self.actions}
    diagram_title = title or getattr(self, "_dataflow_name", None)
    return ColumnRelationshipERConverter.convert(
        rels,
        entity_names=entity_names,
        title=diagram_title,
        trace_from_entity=trace_from_entity,
        trace_from_column=trace_from_column,
    )

DomoDataflow_Definition dataclass

DomoDataflow_Definition(
    auth: DomoAuth,
    dataflow: DomoDataflowProtocol,
    dataflow_id: str,
    Actions: DomoDataflow_Actions,
    TriggerSettings: DomoTriggerSettings | None = None,
    version_id: int | None = None,
    version_number: int | None = None,
    Canvas: CanvasElements | None = None,
    procedures: list | None = None,
    raw: dict | None = None,
)

Manager class for dataflow definition.

This class wraps the complete dataflow definition including actions, triggers, GUI layout, and version information.

Attributes:

Name Type Description
auth DomoAuth

DomoAuth instance from parent dataflow

dataflow DomoDataflowProtocol

Reference to parent DomoDataflow

dataflow_id str

ID of the parent dataflow

Actions DomoDataflow_Actions

Manager for dataflow actions

TriggerSettings DomoTriggerSettings | None

Manager for trigger configuration

version_id int | None

Current version ID

version_number int | None

Current version number

gui dict | None

GUI canvas layout configuration

procedures list | None

List of procedures (legacy)

raw dict | None

Raw definition data from API

Example

dataflow = await DomoDataflow.get_by_id(auth=auth, dataflow_id=123) await dataflow.Definition.get() print(f"Actions: {len(dataflow.Definition.Actions.actions)}") print(f"Version: {dataflow.Definition.version_number}")

gui property

gui: dict | None

Raw gui dict for backwards compatibility.

canvas_edges

canvas_edges() -> list[tuple[str, str]]

Return (source_id, target_id) pairs derived from action dependsOn.

Cross-domain join: needs tile IDs from Canvas and dependsOn from Actions. Owned here — not delegated to either sub-manager.

Source code in src/crew_dcs/classes/DomoDataflow/definition.py
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
def canvas_edges(self) -> list[tuple[str, str]]:
    """Return (source_id, target_id) pairs derived from action dependsOn.

    Cross-domain join: needs tile IDs from Canvas and dependsOn from Actions.
    Owned here — not delegated to either sub-manager.
    """
    if not self.Canvas:
        return []
    tile_ids = {t.id for t in self.Canvas.tiles}
    return [
        (dep_id, action.id)
        for action in self.Actions.actions
        for dep_id in (action.depends_on or [])
        if dep_id in tile_ids and action.id in tile_ids
    ]

canvas_layout

canvas_layout() -> list[dict]

Return all canvas elements as normalised dicts.

Delegates to CanvasElements.layout() but enriches Tile labels with action names from the Actions manager.

Source code in src/crew_dcs/classes/DomoDataflow/definition.py
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
def canvas_layout(self) -> list[dict]:
    """Return all canvas elements as normalised dicts.

    Delegates to CanvasElements.layout() but enriches Tile labels with
    action names from the Actions manager.
    """
    if not self.Canvas:
        return []
    action_names: dict[str, str] = {
        a.id: a.name for a in self.Actions.actions if a.id and a.name
    }
    layout = self.Canvas.layout()
    for el in layout:
        if el["type"] == "Tile" and el["id"] in action_names:
            el["label"] = action_names[el["id"]]
    return layout

from_parent classmethod

from_parent(
    parent: DomoDataflowProtocol,
    actions: list | None = None,
    trigger_settings: dict | None = None,
) -> DomoDataflow_Definition

Create a Definition manager from a parent dataflow.

Parameters:

Name Type Description Default
parent DomoDataflowProtocol

Parent DomoDataflow instance

required
actions list | None

Optional initial list of actions

None
trigger_settings dict | None

Optional trigger settings dict

None

Returns:

Type Description
DomoDataflow_Definition

DomoDataflow_Definition instance

Source code in src/crew_dcs/classes/DomoDataflow/definition.py
 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
102
103
104
105
106
107
108
109
110
111
112
113
114
@classmethod
def from_parent(
    cls,
    parent: DomoDataflowProtocol,
    actions: list | None = None,
    trigger_settings: dict | None = None,
) -> DomoDataflow_Definition:
    """Create a Definition manager from a parent dataflow.

    Args:
        parent: Parent DomoDataflow instance
        actions: Optional initial list of actions
        trigger_settings: Optional trigger settings dict

    Returns:
        DomoDataflow_Definition instance
    """
    # Create Actions manager
    actions_manager = DomoDataflow_Actions.from_parent(
        parent=parent, actions=actions
    )

    # Create TriggerSettings if provided
    trigger_mgr = None
    if trigger_settings:
        trigger_mgr = DomoTriggerSettings.from_parent(
            parent=parent, obj=trigger_settings
        )

    return cls(
        auth=parent.auth,
        dataflow=parent,
        dataflow_id=parent.id,
        Actions=actions_manager,
        TriggerSettings=trigger_mgr,
        version_id=parent.version_id,
        version_number=parent.version_number,
        Canvas=CanvasElements.from_dict(
            parent.raw.get("gui") if parent.raw else None
        ),
        procedures=parent.raw.get("procedures") if parent.raw else None,
        raw=parent.raw,
    )

get async

get(
    debug_api: bool = False,
    return_raw: bool = False,
    session: AsyncClient | None = None,
    context: RouteContext | None = None,
    **context_kwargs
)

Fetch and refresh the dataflow definition.

This method fetches the latest dataflow definition from the API and updates all definition components (actions, triggers, GUI, etc.).

Parameters:

Name Type Description Default
debug_api bool

Enable debug logging for API calls

False
return_raw bool

Return raw API response instead of updating

False
session AsyncClient | None

Optional httpx client session

None
context RouteContext | None

Optional RouteContext for the request

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description

Self if return_raw=False, ResponseGetData if return_raw=True

Example

await dataflow.Definition.get() print(f"Loaded {len(dataflow.Definition.Actions.actions)} actions")

Source code in src/crew_dcs/classes/DomoDataflow/definition.py
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
async def get(
    self,
    debug_api: bool = False,
    return_raw: bool = False,
    session: httpx.AsyncClient | None = None,
    context: RouteContext | None = None,
    **context_kwargs,
):
    """Fetch and refresh the dataflow definition.

    This method fetches the latest dataflow definition from the API
    and updates all definition components (actions, triggers, GUI, etc.).

    Args:
        debug_api: Enable debug logging for API calls
        return_raw: Return raw API response instead of updating
        session: Optional httpx client session
        context: Optional RouteContext for the request
        **context_kwargs: Additional context parameters

    Returns:
        Self if return_raw=False, ResponseGetData if return_raw=True

    Example:
        >>> await dataflow.Definition.get()
        >>> print(f"Loaded {len(dataflow.Definition.Actions.actions)} actions")
    """
    context = RouteContext.build_context(
        session=session, debug_api=debug_api, **context_kwargs
    )

    res = await dataflow_routes.get_dataflow_by_id(
        auth=self.auth,
        dataflow_id=self.dataflow_id,
        context=context,
    )

    if return_raw:
        return res

    if not res.is_success:
        return self

    # Update raw data
    self.raw = res.response
    self.dataflow.raw = res.response

    # Update version info
    self.version_id = res.response.get("versionId")
    self.version_number = res.response.get("versionNumber")
    self.dataflow.version_id = self.version_id
    self.dataflow.version_number = self.version_number

    # Update Canvas and procedures
    self.Canvas = CanvasElements.from_dict(res.response.get("gui"))
    self.procedures = res.response.get("procedures")

    # Update parent dataflow attributes
    self.dataflow.name = res.response.get("name")
    self.dataflow.description = res.response.get("description")
    self.dataflow.owner = res.response.get("owner")

    # Re-parse actions from the new definition
    if self.raw and self.raw.get("actions"):
        self.Actions.actions = [
            DomoDataflow_Action_Base.from_dict(action_dict, all_actions=[])
            for action_dict in self.raw["actions"]
        ]

    # Update trigger settings if present
    if self.raw.get("triggerSettings"):
        self.TriggerSettings = DomoTriggerSettings.from_parent(
            parent=self.dataflow, obj=self.raw["triggerSettings"]
        )

    return self

next_free_position

next_free_position(
    width: int = 224, height: int = 96, padding: int = 32
) -> tuple[int, int]

Return (x, y) for a new element that does not overlap any existing element.

Source code in src/crew_dcs/classes/DomoDataflow/definition.py
308
309
310
311
312
313
314
315
316
def next_free_position(
    self, width: int = 224, height: int = 96, padding: int = 32
) -> tuple[int, int]:
    """Return (x, y) for a new element that does not overlap any existing element."""
    if not self.Canvas:
        return (32, 32)
    return self.Canvas.next_free_position(
        width=width, height=height, padding=padding
    )

optimize_canvas

optimize_canvas(
    add_sections: bool = True,
    save: bool = False,
    debug_api: bool = False,
    session: AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> CanvasElements

Optimize the canvas layout using topological positioning.

Computes a layered layout from the action dependency graph, applies standard tile colours by action type, and optionally groups tiles into sections by upstream root.

Parameters:

Name Type Description Default
add_sections bool

Whether to create auto-sections (default True).

True
save bool

Whether to save the definition after optimizing (default False).

False
debug_api bool

Enable debug logging for API calls (only used if save=True).

False
session AsyncClient | None

Optional httpx client session (only used if save=True).

None
context RouteContext | None

Optional RouteContext for the request (only used if save=True).

None
**context_kwargs

Additional context parameters (only used if save=True).

{}

Returns:

Type Description
CanvasElements

The modified :class:CanvasElements instance.

Example::

defn.optimize_canvas(add_sections=True)
# or with auto-save:
defn.optimize_canvas(save=True, version_note="Optimized layout")
Source code in src/crew_dcs/classes/DomoDataflow/definition.py
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
def optimize_canvas(
    self,
    add_sections: bool = True,
    save: bool = False,
    debug_api: bool = False,
    session: httpx.AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
) -> CanvasElements:
    """Optimize the canvas layout using topological positioning.

    Computes a layered layout from the action dependency graph,
    applies standard tile colours by action type, and optionally
    groups tiles into sections by upstream root.

    Args:
        add_sections: Whether to create auto-sections (default True).
        save: Whether to save the definition after optimizing (default False).
        debug_api: Enable debug logging for API calls (only used if save=True).
        session: Optional httpx client session (only used if save=True).
        context: Optional RouteContext for the request (only used if save=True).
        **context_kwargs: Additional context parameters (only used if save=True).

    Returns:
        The modified :class:`CanvasElements` instance.

    Example::

        defn.optimize_canvas(add_sections=True)
        # or with auto-save:
        defn.optimize_canvas(save=True, version_note="Optimized layout")
    """
    from .layout_optimizer import optimize_canvas as _optimize

    canvas = _optimize(self, add_sections=add_sections)

    if save:
        # save() is async so we return a coroutine that the caller awaits
        return self.save(
            debug_api=debug_api,
            session=session,
            version_note="Optimized canvas layout",
            context=context,
            **context_kwargs,
        )

    return canvas

save async

save(
    debug_api: bool = False,
    session: AsyncClient | None = None,
    *,
    version_note: str | None = None,
    context: RouteContext | None = None,
    **context_kwargs
)

Write current Canvas + Actions state back to the API.

Builds the full definition payload from the typed managers so callers don't have to hand-craft raw dicts. Equivalent to calling update(new_definition) with a payload assembled from self.Canvas and self.Actions.

Parameters:

Name Type Description Default
debug_api bool

Enable debug logging for API calls.

False
session AsyncClient | None

Optional httpx client session.

None
version_note str | None

Optional version description shown in the Domo UI history (onboardFlowVersion.description). If omitted the existing value (or server default) is preserved.

None
context RouteContext | None

Optional RouteContext for the request.

None
**context_kwargs

Additional context parameters.

{}

Example::

await defn.save(version_note="Replaced MergeJoin with SplitJoin")
Source code in src/crew_dcs/classes/DomoDataflow/definition.py
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
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
async def save(
    self,
    debug_api: bool = False,
    session: httpx.AsyncClient | None = None,
    *,
    version_note: str | None = None,
    context: RouteContext | None = None,
    **context_kwargs,
):
    """Write current Canvas + Actions state back to the API.

    Builds the full definition payload from the typed managers so callers
    don't have to hand-craft raw dicts. Equivalent to calling
    update(new_definition) with a payload assembled from self.Canvas and
    self.Actions.

    Args:
        debug_api: Enable debug logging for API calls.
        session: Optional httpx client session.
        version_note: Optional version description shown in the Domo UI
            history (``onboardFlowVersion.description``). If omitted the
            existing value (or server default) is preserved.
        context: Optional RouteContext for the request.
        **context_kwargs: Additional context parameters.

    Example::

        await defn.save(version_note="Replaced MergeJoin with SplitJoin")
    """
    if not self.raw:
        raise RuntimeError(
            "Cannot save: Definition has not been fetched yet (call get() first)"
        )

    new_def = dict(self.raw)
    new_def["actions"] = [a.raw for a in self.Actions.actions if a.raw]
    if self.Canvas:
        new_def["gui"] = self.Canvas.to_dict()

    if version_note is not None:
        new_def.setdefault("onboardFlowVersion", {})
        new_def["onboardFlowVersion"]["description"] = version_note

    return await self.update(
        new_def,
        debug_api=debug_api,
        session=session,
        context=context,
        **context_kwargs,
    )

save_and_export_svg async

save_and_export_svg(
    svg_path: str,
    *,
    title: str | None = None,
    version_note: str | None = None,
    debug_api: bool = False,
    session: AsyncClient | None = None,
    context: RouteContext | None = None,
    **context_kwargs
) -> str

Save the definition and export an SVG rendering in one step.

Convenience method that calls :meth:save then :meth:to_svg and writes the SVG to svg_path.

Parameters:

Name Type Description Default
svg_path str

File path to write the SVG to.

required
title str | None

Optional title for the SVG (defaults to dataflow name).

None
version_note str | None

Optional version description for the Domo UI history.

None
debug_api bool

Enable debug logging for API calls.

False
session AsyncClient | None

Optional httpx client session.

None
context RouteContext | None

Optional RouteContext for the request.

None
**context_kwargs

Additional context parameters.

{}

Returns:

Type Description
str

The absolute path to the written SVG file.

Example::

await defn.save_and_export_svg("/tmp/dataflow-509.svg",
    version_note="Replaced MergeJoin with SplitJoin")
Source code in src/crew_dcs/classes/DomoDataflow/definition.py
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
async def save_and_export_svg(
    self,
    svg_path: str,
    *,
    title: str | None = None,
    version_note: str | None = None,
    debug_api: bool = False,
    session: httpx.AsyncClient | None = None,
    context: RouteContext | None = None,
    **context_kwargs,
) -> str:
    """Save the definition and export an SVG rendering in one step.

    Convenience method that calls :meth:`save` then :meth:`to_svg`
    and writes the SVG to *svg_path*.

    Args:
        svg_path: File path to write the SVG to.
        title: Optional title for the SVG (defaults to dataflow name).
        version_note: Optional version description for the Domo UI history.
        debug_api: Enable debug logging for API calls.
        session: Optional httpx client session.
        context: Optional RouteContext for the request.
        **context_kwargs: Additional context parameters.

    Returns:
        The absolute path to the written SVG file.

    Example::

        await defn.save_and_export_svg("/tmp/dataflow-509.svg",
            version_note="Replaced MergeJoin with SplitJoin")
    """
    await self.save(
        debug_api=debug_api,
        session=session,
        version_note=version_note,
        context=context,
        **context_kwargs,
    )
    svg = self.to_svg(title=title)
    import os

    from crew_dcs.utils.files import upsert_folder

    svg_path = os.path.normpath(svg_path)
    upsert_folder(os.path.dirname(svg_path))
    with open(svg_path, "w") as f:
        f.write(svg)
    return svg_path

to_svg

to_svg(title: str | None = None) -> str

Render the canvas as an SVG string (call get() first).

Source code in src/crew_dcs/classes/DomoDataflow/definition.py
383
384
385
386
387
388
389
390
391
def to_svg(self, title: str | None = None) -> str:
    """Render the canvas as an SVG string (call get() first)."""
    from .canvas import render_svg

    return render_svg(
        layout=self.canvas_layout(),
        edges=self.canvas_edges(),
        title=title or (self.dataflow.name if self.dataflow else "Dataflow Canvas"),
    )

update async

update(
    new_definition: dict,
    debug_api: bool = False,
    debug_num_stacks_to_drop: int = 2,
    session: AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs
)

Update the dataflow definition.

This method updates the dataflow definition on the server and then refreshes the local state.

Parameters:

Name Type Description Default
new_definition dict

New definition dictionary

required
debug_api bool

Enable debug logging for API calls

False
debug_num_stacks_to_drop int

Stack frames to drop in debug logging

2
session AsyncClient | None

Optional httpx client session

None
context RouteContext | None

Optional RouteContext for the request

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description

Self with updated definition

Example

new_def = {"name": "Updated Name", "actions": [...]} await dataflow.Definition.update(new_def)

Source code in src/crew_dcs/classes/DomoDataflow/definition.py
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
226
227
228
229
230
231
232
233
234
235
236
237
238
async def update(
    self,
    new_definition: dict,
    debug_api: bool = False,
    debug_num_stacks_to_drop: int = 2,
    session: httpx.AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
):
    """Update the dataflow definition.

    This method updates the dataflow definition on the server and
    then refreshes the local state.

    Args:
        new_definition: New definition dictionary
        debug_api: Enable debug logging for API calls
        debug_num_stacks_to_drop: Stack frames to drop in debug logging
        session: Optional httpx client session
        context: Optional RouteContext for the request
        **context_kwargs: Additional context parameters

    Returns:
        Self with updated definition

    Example:
        >>> new_def = {"name": "Updated Name", "actions": [...]}
        >>> await dataflow.Definition.update(new_def)
    """
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        debug_num_stacks_to_drop=debug_num_stacks_to_drop,
        **context_kwargs,
    )

    await dataflow_routes.update_dataflow_definition(
        auth=self.auth,
        dataflow_id=self.dataflow_id,
        dataflow_definition=new_definition,
        context=context,
    )

    return await self.get(return_raw=False, session=session, debug_api=debug_api)

DomoDataflow_History dataclass

DomoDataflow_History(
    auth: DomoAuth,
    dataflow_id: str,
    dataflow: Any = None,
    execution_history: (
        list[DomoDataflow_History_Execution] | None
    ) = None,
)

get_execution_history async

get_execution_history(
    auth: DomoAuth | None = None,
    maximum: int = 10,
    return_raw: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoDataflow_History_Execution]

retrieves metadata about execution history. includes details like execution status.

Source code in src/crew_dcs/classes/DomoDataflow/history.py
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
async def get_execution_history(
    self,
    auth: DomoAuth | None = None,
    maximum: int = 10,  # maximum number of execution histories to retrieve
    return_raw: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoDataflow_History_Execution]:
    """retrieves metadata about execution history.
    includes details like execution status.
    """

    auth = auth or self.auth or self.dataflow.auth

    base_context = RouteContext.build_context(
        parent_class=self.__class__.__name__,
        **context_kwargs,
    )
    context = RouteContext.build_context(context=context or base_context)

    res = await dataflow_routes.get_dataflow_execution_history(
        auth=auth,
        dataflow_id=self.dataflow_id,
        maximum=maximum,
        context=context,
    )

    if return_raw:
        return res

    execution_history = [
        DomoDataflow_History_Execution.from_dict(df_obj, auth)
        for df_obj in res.response
    ]

    await ce.gather_with_concurrency(
        *[
            domo_execution.get_actions(context=context)
            for domo_execution in execution_history
        ],
        n=20,
    )

    self.execution_history = execution_history

    return self.execution_history

DomoLineage_Dataflow dataclass

DomoLineage_Dataflow(
    auth: DomoAuth,
    parent: Any = None,
    lineage: list[DomoLineage_Link] = list(),
    immediate_dependencies: list[DomoLineage_Link] = list(),
    immediate_dependents: list[DomoLineage_Link] = list(),
    downstream_lineage: list[DomoLineage_Link] = list(),
)

Bases: DomoLineage

Lineage handler for dataflow entities.

MagicETL_Dataflow dataclass

MagicETL_Dataflow(
    auth: DomoAuth,
    id: str,
    raw: dict,
    name: str | None = None,
    owner: str | None = None,
    description: str | None = None,
    tags: list[str] | None = None,
    version_id: int | None = None,
    version_number: int | None = None,
    versions: list[dict[str, Any]] | None = None,
    jupyter_workspace_config: dict | None = None,
    Definition: DomoDataflow_Definition | None = None,
    History: DomoDataflow_History | None = None,
    JupyterWorkspace: DomoJupyterWorkspace | None = None,
)

Bases: DomoDataflow

DomoDataflow subclass for Magic ETL (databaseType == "MAGIC") dataflows.

Adds Magic ETL-specific functionality, most notably instructions which fetches the full tile/function reference from the Domo expression-docs API.

instructions async

instructions(
    return_raw: bool = False,
    debug_api: bool = False,
    debug_num_stacks_to_drop: int = 2,
    session: AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    pin_version: str | None = None,
    **context_kwargs
)

Retrieve the Magic ETL tile/function reference for this instance's Domo instance.

Calls v2 endpoint by default, falls back to v1 if unavailable (non-auth errors). Logs a warning if fallback is used. Raises if both fail.

Parameters:

Name Type Description Default
return_raw bool

Return the raw ResponseGetData object instead of the parsed response body.

False
debug_api bool

Enable API debugging output.

False
debug_num_stacks_to_drop int

Stack frames to drop in debug logging.

2
session AsyncClient | None

Optional httpx.AsyncClient to reuse.

None
context RouteContext | None

RouteContext for request configuration.

None
pin_version str | None

'v2', 'v1', or None for auto-fallback.

None
**context_kwargs

Additional context parameters.

{}

Returns:

Type Description

Parsed tile documentation (list / dict) when return_raw=False,

or the ResponseGetData object when return_raw=True.

Source code in src/crew_dcs/classes/DomoDataflow/magic_etl.py
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 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
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
async def instructions(
    self,
    return_raw: bool = False,
    debug_api: bool = False,
    debug_num_stacks_to_drop: int = 2,
    session: httpx.AsyncClient | None = None,
    *,
    context: RouteContext | None = None,
    pin_version: str | None = None,  # 'v2', 'v1', or None for auto
    **context_kwargs,
):
    """Retrieve the Magic ETL tile/function reference for this instance's Domo instance.

    Calls v2 endpoint by default, falls back to v1 if unavailable (non-auth errors).
    Logs a warning if fallback is used. Raises if both fail.

    Args:
        return_raw: Return the raw ``ResponseGetData`` object instead of the
            parsed response body.
        debug_api: Enable API debugging output.
        debug_num_stacks_to_drop: Stack frames to drop in debug logging.
        session: Optional ``httpx.AsyncClient`` to reuse.
        context: ``RouteContext`` for request configuration.
        pin_version: 'v2', 'v1', or None for auto-fallback.
        **context_kwargs: Additional context parameters.

    Returns:
        Parsed tile documentation (list / dict) when ``return_raw=False``,
        or the ``ResponseGetData`` object when ``return_raw=True``.
    """
    import logging

    logger = logging.getLogger("crew_dcs.magic_etl.instructions")

    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        debug_num_stacks_to_drop=debug_num_stacks_to_drop,
        **context_kwargs,
    )

    # Allow pinning version for explicit calls
    if pin_version == "v1":
        res = await dataflow_routes.get_magic_etl_tile_documentation_v1(
            auth=self.auth,
            context=context,
        )
        if return_raw:
            return res
        return res.response
    if pin_version == "v2":
        res = await dataflow_routes.get_magic_etl_tile_documentation(
            auth=self.auth,
            context=context,
        )
        if return_raw:
            return res
        return res.response

    # Default: try v2, fallback to v1 if non-auth error
    try:
        res = await dataflow_routes.get_magic_etl_tile_documentation(
            auth=self.auth,
            context=context,
        )
        if return_raw:
            return res
        return res.response
    except Exception as e:
        # Only fallback if not auth error
        if hasattr(e, "status") and e.status in {401, 403}:
            raise
        logger.warning(
            "Falling back to v1 expression-docs endpoint due to: %s", str(e)
        )
        try:
            res = await dataflow_routes.get_magic_etl_tile_documentation_v1(
                auth=self.auth,
                context=context,
            )
            if return_raw:
                return res
            return res.response
        except Exception as e2:  # noqa: BLE001
            raise RuntimeError(  # noqa: B904
                f"Both v2 and v1 expression-docs endpoints failed: v2 error: {e}, v1 error: {e2}"
            )

Modules