Skip to content

lineage

lineage

Lineage package for Domo entity dependency tracking.

This package provides lineage handling for tracking dependencies between Domo entities (datasets, cards, pages, dataflows, publications, etc.).

Modules:

Name Description
base

Core DomoLineage class (fields, properties, get(), get_datacenter_lineage)

registry

Decorators, registries, and import maps for lineage type registration

federation

Federation connection mixin (subscriber→publisher linking)

publish_lineage

Publish lineage mixin (subscription/publication chain tracing)

column_lineage

Column-level lineage mixin (get_column_lineage, to_erd)

lineage_link

Link registry and handlers

federation_context

Federation/publish support

publish_resolver

Subscription resolution utilities

protocols

Structural types to avoid circular imports

Note

Entity-specific lineage handlers and lineage link classes live alongside their entities (e.g., classes/DomoCard/lineage.py). They are imported here only for registration side effects.

DomoLineage dataclass

DomoLineage(
    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: DomoBase

Lineage handler for Domo entities.

The parent attribute refers to the entity that this lineage is based off of (the entity we're getting lineage for). This matches the from_parent() method.

Note: parent is NOT a dependency. Dependencies are what the parent entity depends on (upstream entities), and are returned by the get() method.

is_federated property

is_federated: bool

Check if the parent entity is federated.

Returns:

Name Type Description
bool bool

True if parent entity is federated, False otherwise

is_published property

is_published: bool

Check if entity has been identified as published.

Returns:

Name Type Description
bool bool

True if Federation attribute exists and indicates published state

parent_type property

parent_type: str

Return the lineage type string for this lineage's parent entity.

check_is_federated async

check_is_federated(
    *,
    check_is_published: bool = False,
    parent_auth_retrieval_fn: (
        Callable[[str], Any | Awaitable[Any]] | None
    ) = None,
    session: AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None
) -> bool

Check if this entity is federated, optionally check if published.

Parameters:

Name Type Description Default
check_is_published bool

If True, also check publish state and populate Federation context

False
parent_auth_retrieval_fn Callable[[str], Any | Awaitable[Any]] | None

Function to retrieve parent instance auth (required if check_is_published=True)

None
session AsyncClient | None

HTTP client session

None
debug_api bool

Enable API debug logging

False
max_subscriptions_to_check int | None

Limit subscription search scope

None

Returns:

Name Type Description
bool bool

True if entity is federated, False otherwise

Source code in src/crew_dcs/classes/subentity/lineage/base.py
1786
1787
1788
1789
1790
1791
1792
1793
1794
1795
1796
1797
1798
1799
1800
1801
1802
1803
1804
1805
1806
1807
1808
1809
1810
1811
1812
1813
1814
1815
1816
1817
1818
1819
1820
1821
1822
1823
1824
1825
1826
1827
1828
1829
1830
1831
async def check_is_federated(
    self,
    *,
    check_is_published: bool = False,
    parent_auth_retrieval_fn: Callable[[str], Any | Awaitable[Any]] | None = None,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
) -> bool:
    """Check if this entity is federated, optionally check if published.

    Args:
        check_is_published: If True, also check publish state and populate Federation context
        parent_auth_retrieval_fn: Function to retrieve parent instance auth (required if check_is_published=True)
        session: HTTP client session
        debug_api: Enable API debug logging
        max_subscriptions_to_check: Limit subscription search scope

    Returns:
        bool: True if entity is federated, False otherwise
    """
    if not self.parent:
        raise ValueError(
            "Parent must be set. Use get_lineage_from_entity() to create lineage with a parent."
        )

    # Check federation state
    is_federated_entity = (
        hasattr(self.parent, "is_federated") and self.parent.is_federated
    )

    # Optionally check publish state
    if check_is_published and is_federated_entity:
        if not parent_auth_retrieval_fn:
            raise ValueError(
                "parent_auth_retrieval_fn is required when check_is_published=True"
            )

        await self.check_is_published(
            parent_auth_retrieval_fn=parent_auth_retrieval_fn,
            session=session,
            debug_api=debug_api,
            max_subscriptions_to_check=max_subscriptions_to_check,
        )

    return is_federated_entity

check_is_published async

check_is_published(
    *,
    parent_auth_retrieval_fn: Callable[
        [str], Any | Awaitable[Any]
    ],
    entity_type: str | None = None,
    entity_id: str | None = None,
    session: AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None
) -> bool

Check if this entity is published and populate Federation context.

Parameters:

Name Type Description Default
parent_auth_retrieval_fn Callable[[str], Any | Awaitable[Any]]

Function to retrieve parent instance auth

required
entity_type str | None

Override entity type (defaults to parent.entity_type)

None
entity_id str | None

Override entity ID (defaults to parent.id)

None
session AsyncClient | None

HTTP client session

None
debug_api bool

Enable API debug logging

False
max_subscriptions_to_check int | None

Limit subscription search scope

None

Returns:

Name Type Description
bool bool

True if entity is published, False otherwise

Source code in src/crew_dcs/classes/subentity/lineage/base.py
1732
1733
1734
1735
1736
1737
1738
1739
1740
1741
1742
1743
1744
1745
1746
1747
1748
1749
1750
1751
1752
1753
1754
1755
1756
1757
1758
1759
1760
1761
1762
1763
1764
1765
1766
1767
1768
1769
1770
1771
1772
1773
1774
1775
1776
1777
1778
1779
1780
1781
1782
1783
1784
async def check_is_published(
    self,
    *,
    parent_auth_retrieval_fn: Callable[[str], Any | Awaitable[Any]],
    entity_type: str | None = None,
    entity_id: str | None = None,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
) -> bool:
    """Check if this entity is published and populate Federation context.

    Args:
        parent_auth_retrieval_fn: Function to retrieve parent instance auth
        entity_type: Override entity type (defaults to parent.entity_type)
        entity_id: Override entity ID (defaults to parent.id)
        session: HTTP client session
        debug_api: Enable API debug logging
        max_subscriptions_to_check: Limit subscription search scope

    Returns:
        bool: True if entity is published, False otherwise
    """
    from .federation_context import FederationContext

    if not self.parent:
        raise ValueError(
            "Parent must be set. Use get_lineage_from_entity() to create lineage with a parent."
        )

    # Create or reuse Federation helper
    if self.Federation is None:
        self.Federation = FederationContext(parent=self.parent)

    # Check if already published
    if self.Federation.is_published:
        return True

    # Determine entity type and ID
    content_type = entity_type or self.parent.entity_type
    target_id = entity_id or str(self.parent.id)

    # Perform publish check
    is_published = await self.Federation.check_if_published(
        retrieve_parent_auth_fn=parent_auth_retrieval_fn,
        entity_type=content_type,
        entity_id=target_id,
        session=session,
        debug_api=debug_api,
        max_subscriptions_to_check=max_subscriptions_to_check,
    )

    return is_published  # noqa: RET504

from_parent classmethod

from_parent(parent, auth: DomoAuth = None)

Create a DomoLineage instance from a parent entity synchronously.

This is a simple factory for synchronous initialization (e.g., post_init). For async usage with publish checking, use get_lineage_from_entity() instead.

Source code in src/crew_dcs/classes/subentity/lineage/base.py
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
@classmethod
def from_parent(cls, parent, auth: DomoAuth = None):
    """Create a DomoLineage instance from a parent entity synchronously.

    This is a simple factory for synchronous initialization (e.g., __post_init__).
    For async usage with publish checking, use get_lineage_from_entity() instead.
    """
    parent_class_name = parent.__class__.__name__
    lineage_class = _LINEAGE_REGISTRY.get(parent_class_name)

    if lineage_class is None:
        _try_register_lineage(parent_class_name)
        lineage_class = _LINEAGE_REGISTRY.get(parent_class_name)

    if lineage_class is None:
        raise ValueError(
            f"No lineage class registered for parent type: {parent_class_name}. "
            f"Known types: {sorted(_LINEAGE_REGISTRY.keys())}"
        )

    return lineage_class(
        auth=auth or parent.auth,
        parent=parent,
    )

get async

get(
    session: AsyncClient = None,
    debug_api: bool = False,
    return_raw: bool = False,
    parent_auth: DomoAuth = None,
    parent_auth_retrieval_fn: Callable | None = None,
    max_depth: int = 100,
    lineage_cache: (
        dict[tuple[Any, ...], list[DomoLineage_Link]] | None
    ) = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoLineage_Link]

Get lineage for this entity.

Parameters:

Name Type Description Default
session AsyncClient

HTTP session for reuse

None
debug_api bool

Enable API debug logging

False
return_raw bool

Return raw response

False
parent_auth DomoAuth

Authentication for publisher instance

None
parent_auth_retrieval_fn Callable | None

Callable to retrieve publisher auth by domain

None
max_depth int

Maximum lineage depth (default: 100 for complete upstream lineage)

100
context RouteContext | None

Route context for API configuration

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description
list[DomoLineage_Link]

List of lineage links (default) or LineageGraph (if return_graph=True)

Raises:

Type Description
ClassError

If parent is not set

Source code in src/crew_dcs/classes/subentity/lineage/base.py
1833
1834
1835
1836
1837
1838
1839
1840
1841
1842
1843
1844
1845
1846
1847
1848
1849
1850
1851
1852
1853
1854
1855
1856
1857
1858
1859
1860
1861
1862
1863
1864
1865
1866
1867
1868
1869
1870
1871
1872
1873
1874
1875
1876
1877
1878
1879
1880
1881
1882
1883
1884
1885
1886
1887
1888
1889
1890
1891
1892
1893
1894
1895
1896
1897
1898
1899
1900
1901
1902
1903
1904
1905
1906
1907
1908
1909
1910
1911
1912
1913
1914
1915
1916
1917
1918
1919
1920
1921
1922
1923
1924
1925
1926
1927
1928
1929
1930
1931
1932
1933
1934
1935
1936
1937
1938
1939
1940
1941
1942
1943
1944
1945
1946
1947
1948
1949
async def get(
    self,
    session: httpx.AsyncClient = None,
    debug_api: bool = False,
    return_raw: bool = False,
    parent_auth: DomoAuth = None,
    parent_auth_retrieval_fn: Callable | None = None,
    max_depth: int = 100,
    lineage_cache: dict[tuple[Any, ...], list[DomoLineage_Link]] | None = None,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoLineage_Link]:
    """Get lineage for this entity.

    Args:
        session: HTTP session for reuse
        debug_api: Enable API debug logging
        return_raw: Return raw response
        parent_auth: Authentication for publisher instance
        parent_auth_retrieval_fn: Callable to retrieve publisher auth by domain
        max_depth: Maximum lineage depth (default: 100 for complete upstream lineage)
        context: Route context for API configuration
        **context_kwargs: Additional context parameters

    Returns:
        List of lineage links (default) or LineageGraph (if return_graph=True)

    Raises:
        ClassError: If parent is not set
    """

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

    if not self.parent:
        raise ClassError(
            message="Parent must be set. Use from_parent() to create lineage with a parent.",
            cls_instance=self,
        )

    cache = lineage_cache if lineage_cache is not None else self._lineage_cache
    cache_key = self._build_lineage_cache_key(
        max_depth=max_depth,
        parent_auth=parent_auth,
        parent_auth_retrieval_fn=parent_auth_retrieval_fn,
        return_raw=return_raw,
    )

    if cache_key in cache:
        await logger.debug(
            "DomoLineage.get: returning cached lineage "
            f"(len={len(cache[cache_key])})"
        )
        self.lineage = list(cache[cache_key])
        self._cached_lineage_params = self._build_cache_key(
            parent_auth=parent_auth,
            parent_auth_retrieval_fn=parent_auth_retrieval_fn,
            return_raw=return_raw,
        )
        return list(self.lineage)

    self._seen_link_keys.clear()

    entity_key = (str(self.parent.id), str(self.parent.entity_type))
    self._seen_link_keys.add(entity_key)

    # Per documented design (lineage/AGENTS.md): get() ALWAYS calls
    # get_datacenter_lineage() first — this is correct for all entity types.
    # If the entity is federated AND auth is provided, additionally traverse
    # the subscription/publication/federation chain via get_federated_lineage().
    await self.get_datacenter_lineage(
        max_depth=max_depth,
        context=context,
        **context_kwargs,
    )

    if self.is_federated and (parent_auth or parent_auth_retrieval_fn):
        federated_links = await self.get_federated_lineage(
            parent_auth=parent_auth,
            parent_auth_retrieval_fn=parent_auth_retrieval_fn,
            session=session,
            debug_api=debug_api,
            return_raw=return_raw,
            context=context,
            **context_kwargs,
        )
        if federated_links:
            self.lineage = self._merge_lineage_links(
                list(self.lineage), federated_links
            )

    # Connect federated entities to their publisher counterparts
    if parent_auth or parent_auth_retrieval_fn:
        self.lineage = await self._connect_federated_entities(
            self.lineage,
            parent_auth=parent_auth,
            parent_auth_retrieval_fn=parent_auth_retrieval_fn,
            session=session,
            debug_api=debug_api,
            context=context,
            **context_kwargs,
        )

    self._cached_lineage_params = self._build_cache_key(
        parent_auth=parent_auth,
        parent_auth_retrieval_fn=parent_auth_retrieval_fn,
        return_raw=return_raw,
    )

    cache[cache_key] = list(self.lineage)

    return self.lineage

get_datacenter_lineage async

get_datacenter_lineage(
    session: AsyncClient = None,
    debug_api: bool = False,
    return_raw: bool = False,
    max_depth: int | None = None,
    traverse_up: bool = True,
    traverse_down: bool = False,
    populate_self: bool = True,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoLineage_Link]

Query the datacenter lineage API and return ALL nodes.

Returns all nodes from the lineage API response as DomoLineage_Link objects. Each node includes its nested parent/child relationships.

Parameters:

Name Type Description Default
session AsyncClient

HTTP session for reuse

None
debug_api bool

Enable API debug logging

False
return_raw bool

Return raw response

False
max_depth int | None

Maximum lineage depth

None
traverse_up bool

Whether to traverse upstream (dependencies). Default True.

True
traverse_down bool

Whether to traverse downstream (dependents). Default False for backward compatibility.

False
populate_self bool

If True (default), set self.lineage = result. If False, return the result without mutating self.lineage. Use False when calling from get_downstream() to avoid overwriting upstream lineage.

True
context RouteContext | None

Route context for API configuration

None
**context_kwargs

Additional context parameters

{}
Source code in src/crew_dcs/classes/subentity/lineage/base.py
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
async def get_datacenter_lineage(
    self,
    session: httpx.AsyncClient = None,
    debug_api: bool = False,
    return_raw: bool = False,
    max_depth: int | None = None,
    traverse_up: bool = True,
    traverse_down: bool = False,
    populate_self: bool = True,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoLineage_Link]:
    """Query the datacenter lineage API and return ALL nodes.

    Returns all nodes from the lineage API response as DomoLineage_Link objects.
    Each node includes its nested parent/child relationships.

    Args:
        session: HTTP session for reuse
        debug_api: Enable API debug logging
        return_raw: Return raw response
        max_depth: Maximum lineage depth
        traverse_up: Whether to traverse upstream (dependencies). Default True.
        traverse_down: Whether to traverse downstream (dependents). Default False
            for backward compatibility.
        populate_self: If True (default), set self.lineage = result. If False,
            return the result without mutating self.lineage. Use False when
            calling from get_downstream() to avoid overwriting upstream lineage.
        context: Route context for API configuration
        **context_kwargs: Additional context parameters
    """
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        **context_kwargs,
    )

    res = await datacenter_routes.get_lineage_upstream(
        auth=self.parent.auth,
        entity_type=self.parent_type,
        entity_id=self.parent.id,
        max_depth=max_depth,
        traverse_up=traverse_up,
        traverse_down=traverse_down,
        context=context,
    )

    if return_raw:
        return res

    await logger.debug(
        f"get_datacenter_lineage: {len(res.response)} nodes in response"
    )

    result = [
        DomoLineage_Link.from_dict(
            lineage_api_obj=lineage_api_obj,
            auth=self.auth,
        )
        for lineage_api_obj in res.response.values()
    ]

    # Load entities for all links
    for link in result:
        try:
            link.entity = await link.get_entity(
                context=context, debug_api=debug_api
            )
        except DomoError as e:
            await logger.warning(
                f"Failed to load entity for {link.type} {link.id}: {e}"
            )

    await logger.debug(
        f"get_datacenter_lineage: {len(result)} links created, "
        f"{sum(1 for link in result if link.entity)} entities loaded"
    )

    if populate_self:
        self.lineage = result

    return result

get_downstream async

get_downstream(
    session: AsyncClient = None,
    debug_api: bool = False,
    return_raw: bool = False,
    max_depth: int = 1,
    *,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoLineage_Link]

Get downstream lineage (dependents) for this entity.

Returns entities that depend on this entity — the "impact" direction. Default max_depth=1 for safety (popular datasets may have 1000+ dependents).

When the parent entity is published (has Federation.is_published), federated dependents on subscriber instances are also discovered and appended to the root link's dependents.

Parameters:

Name Type Description Default
session AsyncClient

HTTP session for reuse

None
debug_api bool

Enable API debug logging

False
return_raw bool

Return raw response

False
max_depth int

Maximum downstream depth (default 1 for safety)

1
parent_auth DomoAuth | None

Authentication for publisher instance (for federation)

None
parent_auth_retrieval_fn Callable | None

Callable to retrieve publisher auth by domain

None
context RouteContext | None

Route context for API configuration

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description
list[DomoLineage_Link]

List of downstream lineage links

Source code in src/crew_dcs/classes/subentity/lineage/base.py
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
async def get_downstream(
    self,
    session: httpx.AsyncClient = None,
    debug_api: bool = False,
    return_raw: bool = False,
    max_depth: int = 1,
    *,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoLineage_Link]:
    """Get downstream lineage (dependents) for this entity.

    Returns entities that depend on this entity — the "impact" direction.
    Default max_depth=1 for safety (popular datasets may have 1000+ dependents).

    When the parent entity is published (has Federation.is_published),
    federated dependents on subscriber instances are also discovered and
    appended to the root link's dependents.

    Args:
        session: HTTP session for reuse
        debug_api: Enable API debug logging
        return_raw: Return raw response
        max_depth: Maximum downstream depth (default 1 for safety)
        parent_auth: Authentication for publisher instance (for federation)
        parent_auth_retrieval_fn: Callable to retrieve publisher auth by domain
        context: Route context for API configuration
        **context_kwargs: Additional context parameters

    Returns:
        List of downstream lineage links
    """
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        **context_kwargs,
    )

    result = await self.get_datacenter_lineage(
        session=session,
        debug_api=debug_api,
        return_raw=return_raw,
        traverse_down=True,
        traverse_up=False,
        populate_self=False,
        max_depth=max_depth,
        context=context,
    )

    # Find the root link (matching self.parent.id and self.parent_type)
    root_link = next(
        (
            link
            for link in result
            if str(link.id) == str(self.parent.id) and link.type == self.parent_type
        ),
        None,
    )

    if root_link is not None:
        self.immediate_dependents = root_link.dependents

    # Augment with federated dependents if the entity is published
    if parent_auth or parent_auth_retrieval_fn:
        result = await self._connect_federated_dependents(
            result,
            parent_auth=parent_auth,
            parent_auth_retrieval_fn=parent_auth_retrieval_fn,
            session=session,
            debug_api=debug_api,
            context=context,
            **context_kwargs,
        )

        # Re-derive root_link and immediate_dependents after federation
        # augmentation may have added dependents to the root link.
        root_link = next(
            (
                link
                for link in result
                if str(link.id) == str(self.parent.id)
                and link.type == self.parent_type
            ),
            None,
        )
        if root_link is not None:
            self.immediate_dependents = root_link.dependents

    self.downstream_lineage = result

    return self.downstream_lineage

get_federated_lineage async

get_federated_lineage(
    session: AsyncClient | None = None,
    debug_api: bool = False,
    return_raw: bool = False,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    _debug_num_stacks_to_drop: int = 3,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoLineage_Link]

Get lineage for a federated entity.

Always delegates to trace_publish_lineage() which builds the full chain: subscriber → subscription → publication → publisher → publisher lineage.

If the subscription hasn't been discovered yet, ensure_subscription() is called first. If no subscription is found the entity is not published and an empty list is returned.

Source code in src/crew_dcs/classes/subentity/lineage/base.py
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
async def get_federated_lineage(
    self,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    return_raw: bool = False,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    _debug_num_stacks_to_drop: int = 3,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoLineage_Link]:
    """Get lineage for a federated entity.

    Always delegates to trace_publish_lineage() which builds the full chain:
    subscriber → subscription → publication → publisher → publisher lineage.

    If the subscription hasn't been discovered yet, ensure_subscription() is
    called first.  If no subscription is found the entity is not published
    and an empty list is returned.
    """
    if not self.parent and (not self.parent_type):
        raise ValueError(
            "Parent must be set. Use from_parent() to create lineage with a parent."
        )

    # Ensure federation context exists
    if not getattr(self.parent, "Federation", None):
        self.parent.enable_federation_support()

    federation = self.parent.Federation

    # Discover subscription if not already known
    if not federation.is_published:
        await federation.ensure_subscription(
            retrieve_parent_auth_fn=parent_auth_retrieval_fn,
            parent_auth=parent_auth,
            entity_type=self.parent.entity_type,
            entity_id=str(self.parent.id),
            session=session,
            debug_api=debug_api,
        )

    if not federation.subscription:
        await logger.warning(
            f"get_federated_lineage: entity {self.parent.id} is federated but no "
            f"subscription found — skipping publisher lineage traversal.",
            entity_id=str(self.parent.id),
            entity_type=self.parent.entity_type,
        )
        return []

    # Delegate to trace_publish_lineage which builds the full chain
    # (subscriber → subscription → publication → publisher → publisher lineage)
    return await self.trace_publish_lineage(
        publish_helper=federation,
        parent_auth=parent_auth,
        parent_auth_retrieval_fn=parent_auth_retrieval_fn,
        context=context,
        **context_kwargs,
    )

get_impact async

get_impact(
    max_depth: int = 1,
    entity_types: list[str] | None = None,
    session: AsyncClient = None,
    debug_api: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[Any]

Get flat list of all downstream entities (blast radius).

Returns deduplicated entities, optionally filtered by type. Default max_depth=1 for safety.

Parameters:

Name Type Description Default
max_depth int

Maximum downstream depth (default 1 for safety)

1
entity_types list[str] | None

Optional list of entity type strings to filter by (e.g. ["DATA_SOURCE", "DATAFLOW"]). If None, all types are included.

None
session AsyncClient

HTTP session for reuse

None
debug_api bool

Enable API debug logging

False
context RouteContext | None

Route context for API configuration

None
**context_kwargs

Additional context parameters

{}

Returns:

Type Description
list[Any]

Deduplicated list of downstream entities, optionally filtered by type

Source code in src/crew_dcs/classes/subentity/lineage/base.py
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
async def get_impact(
    self,
    max_depth: int = 1,
    entity_types: list[str] | None = None,
    session: httpx.AsyncClient = None,
    debug_api: bool = False,
    *,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[Any]:
    """Get flat list of all downstream entities (blast radius).

    Returns deduplicated entities, optionally filtered by type.
    Default max_depth=1 for safety.

    Args:
        max_depth: Maximum downstream depth (default 1 for safety)
        entity_types: Optional list of entity type strings to filter by
            (e.g. ["DATA_SOURCE", "DATAFLOW"]). If None, all types are included.
        session: HTTP session for reuse
        debug_api: Enable API debug logging
        context: Route context for API configuration
        **context_kwargs: Additional context parameters

    Returns:
        Deduplicated list of downstream entities, optionally filtered by type
    """
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        **context_kwargs,
    )

    await self.get_downstream(
        max_depth=max_depth,
        session=session,
        debug_api=debug_api,
        context=context,
        **context_kwargs,
    )

    # Extract entities from all links, skipping links where entity is None
    seen: set[tuple[str, str]] = set()
    entities: list[Any] = []

    for link in self.downstream_lineage:
        if link.entity is None:
            continue

        # Exclude the parent entity itself from the blast radius
        if str(link.entity.id) == str(self.parent.id):
            continue

        entity_key = (
            str(link.entity.id),
            getattr(link.entity, "entity_type", "") or "",
        )

        if entity_key in seen:
            continue

        seen.add(entity_key)

        if entity_types is not None:  # noqa: SIM102
            if getattr(link.entity, "entity_type", "") not in entity_types:
                continue

        entities.append(link.entity)

    return entities

get_lineage_from_entity async classmethod

get_lineage_from_entity(
    entity,
    auth: DomoAuth = None,
    check_is_published: bool = False,
    parent_auth_retrieval_fn: Callable | None = None,
    session: AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
)

Create a DomoLineage instance from an entity.

Parameters:

Name Type Description Default
entity

The entity to create lineage for

required
auth DomoAuth

Authentication object (defaults to entity.auth)

None
check_is_published bool

If True, check if entity is published and populate Federation context

False
parent_auth_retrieval_fn Callable | None

Function to retrieve parent instance auth for federated entities

None
session AsyncClient | None

HTTP client session

None
debug_api bool

Enable API debug logging

False
max_subscriptions_to_check int | None

Limit subscription search scope

None

Returns:

Type Description

DomoLineage instance with optional Federation context populated

Source code in src/crew_dcs/classes/subentity/lineage/base.py
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
@classmethod
async def get_lineage_from_entity(
    cls,
    entity,
    auth: DomoAuth = None,
    check_is_published: bool = False,
    parent_auth_retrieval_fn: Callable | None = None,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
):
    """Create a DomoLineage instance from an entity.

    Args:
        entity: The entity to create lineage for
        auth: Authentication object (defaults to entity.auth)
        check_is_published: If True, check if entity is published and populate Federation context
        parent_auth_retrieval_fn: Function to retrieve parent instance auth for federated entities
        session: HTTP client session
        debug_api: Enable API debug logging
        max_subscriptions_to_check: Limit subscription search scope

    Returns:
        DomoLineage instance with optional Federation context populated
    """
    parent_class_name = entity.__class__.__name__
    lineage_class = _LINEAGE_REGISTRY.get(parent_class_name)

    if lineage_class is None:
        _try_register_lineage(parent_class_name)
        lineage_class = _LINEAGE_REGISTRY.get(parent_class_name)

    if lineage_class is None:
        raise ValueError(
            f"No lineage class registered for parent type: {parent_class_name}. "
            f"Known types: {sorted(_LINEAGE_REGISTRY.keys())}"
        )

    lineage = lineage_class(
        auth=auth or entity.auth,
        parent=entity,
    )

    if check_is_published:
        await lineage.check_is_published(
            parent_auth_retrieval_fn=parent_auth_retrieval_fn,
            session=session,
            debug_api=debug_api,
            max_subscriptions_to_check=max_subscriptions_to_check,
        )

    return lineage

get_publish_lineage async

get_publish_lineage(
    *,
    publish_helper=None,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoLineage_Link]

Build lineage chain for published entities.

.. deprecated:: Use :meth:trace_publish_lineage instead. This thin wrapper exists only for backward compatibility.

Source code in src/crew_dcs/classes/subentity/lineage/base.py
1709
1710
1711
1712
1713
1714
1715
1716
1717
1718
1719
1720
1721
1722
1723
1724
1725
1726
1727
1728
1729
1730
async def get_publish_lineage(
    self,
    *,
    publish_helper=None,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoLineage_Link]:
    """Build lineage chain for published entities.

    .. deprecated::
        Use :meth:`trace_publish_lineage` instead.  This thin wrapper
        exists only for backward compatibility.
    """
    return await self.trace_publish_lineage(
        publish_helper=publish_helper,
        parent_auth=parent_auth,
        parent_auth_retrieval_fn=parent_auth_retrieval_fn,
        context=context,
        **context_kwargs,
    )

to_dict

to_dict(
    override_fn: Callable | None = None,
    return_snake_case: bool = False,
) -> dict

Serialize lineage without circular dependency expansion.

Source code in src/crew_dcs/classes/subentity/lineage/base.py
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
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
382
383
384
385
386
387
388
389
390
391
392
def to_dict(  # noqa: C901
    self, override_fn: Callable | None = None, return_snake_case: bool = False
) -> dict:
    """Serialize lineage without circular dependency expansion."""
    if override_fn:
        return override_fn(self)

    def _format_key(key: str) -> str:
        return key if return_snake_case else convert_snake_to_pascal(key)

    def _serialize_entity(entity: Any) -> Any:
        if entity is None:
            return None
        if hasattr(entity, "to_dict") and callable(entity.to_dict):
            return entity.to_dict(return_snake_case=return_snake_case)
        return entity

    def _serialize_link(link: DomoLineage_Link) -> dict:
        payload = {
            _format_key("id"): link.id,
            _format_key("type"): link.type,
        }

        if link.entity:
            payload[_format_key("entity")] = _serialize_entity(link.entity)

        if link.dependencies:
            payload[_format_key("dependencies")] = [
                {
                    _format_key("id"): dep.id,
                    _format_key("type"): dep.type,
                }
                for dep in link.dependencies
            ]

        if link.dependents:
            payload[_format_key("dependents")] = [
                {
                    _format_key("id"): dep.id,
                    _format_key("type"): dep.type,
                }
                for dep in link.dependents
            ]

        return payload

    parent_payload = None
    if self.parent is not None:
        parent_payload = {
            _format_key("id"): getattr(self.parent, "id", None),
            _format_key("type"): getattr(self.parent, "entity_type", None),
            _format_key("name"): getattr(
                self.parent,
                "name",
                getattr(self.parent, "title", None),
            ),
        }

    federation_payload = None
    if self.Federation is not None:
        if hasattr(self.Federation, "to_dict") and callable(
            self.Federation.to_dict
        ):
            federation_payload = self.Federation.to_dict(
                return_snake_case=return_snake_case
            )
        else:
            federation_payload = self.Federation

    result = {
        _format_key("parent"): parent_payload,
        _format_key("parent_type"): self.parent_type if self.parent else None,
        _format_key("lineage"): [_serialize_link(link) for link in self.lineage],
        _format_key("immediate_dependencies"): [
            _serialize_link(link) for link in self.immediate_dependencies
        ],
        _format_key("immediate_dependents"): [
            _serialize_link(link) for link in self.immediate_dependents
        ],
        _format_key("federation"): federation_payload,
        _format_key("is_federated"): self.is_federated,
        _format_key("is_published"): self.is_published,
    }

    if self.downstream_lineage:
        result[_format_key("downstream_lineage")] = [
            _serialize_link(link) for link in self.downstream_lineage
        ]

    return result

trace_publish_lineage async

trace_publish_lineage(
    *,
    publish_helper=None,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    max_depth: int = 100,
    context: RouteContext | None = None,
    **context_kwargs
) -> list[DomoLineage_Link]

Build the full publish-lineage chain for a federated published entity.

Orchestrates: 1. Ensure subscription is resolved (via FederationContext). 2. Resolve publisher auth. 3. Fetch parent publication. 4. Call _resolve_publisher_entity (overridable hook) to find the publisher-side entity that corresponds to self.parent. 5. Fetch the publisher entity's own lineage. 6. Assemble the chain: subscriber entity → subscription → publication → publisher entity → publisher entity lineage.

Subclasses override _resolve_publisher_entity to customise step 4. For example DomoLineage_Card handles cards that are indirectly published via pages.

Parameters:

Name Type Description Default
publish_helper

FederationContext (defaults to self.Federation).

None
parent_auth DomoAuth | None

Explicit publisher auth (alternative to retrieval fn).

None
parent_auth_retrieval_fn Callable | None

Callable(domain) → DomoAuth.

None
max_depth int

Depth limit for publisher entity lineage.

100
context RouteContext | None

Route context for API calls.

None

Returns:

Type Description
list[DomoLineage_Link]

Flat list of lineage links representing the full chain.

Source code in src/crew_dcs/classes/subentity/lineage/base.py
1508
1509
1510
1511
1512
1513
1514
1515
1516
1517
1518
1519
1520
1521
1522
1523
1524
1525
1526
1527
1528
1529
1530
1531
1532
1533
1534
1535
1536
1537
1538
1539
1540
1541
1542
1543
1544
1545
1546
1547
1548
1549
1550
1551
1552
1553
1554
1555
1556
1557
1558
1559
1560
1561
1562
1563
1564
1565
1566
1567
1568
1569
1570
1571
1572
1573
1574
1575
1576
1577
1578
1579
1580
1581
1582
1583
1584
1585
1586
1587
1588
1589
1590
1591
1592
1593
1594
1595
1596
1597
1598
1599
1600
1601
1602
1603
1604
1605
1606
1607
1608
1609
1610
1611
1612
1613
1614
1615
1616
1617
1618
1619
1620
1621
1622
1623
1624
1625
1626
1627
1628
1629
1630
1631
1632
1633
1634
1635
1636
1637
1638
1639
1640
1641
1642
1643
1644
1645
1646
1647
1648
1649
1650
1651
1652
1653
1654
1655
1656
1657
1658
1659
1660
1661
1662
1663
1664
1665
1666
1667
1668
1669
1670
1671
1672
1673
1674
1675
1676
1677
1678
1679
1680
1681
1682
1683
1684
1685
1686
1687
1688
1689
1690
1691
1692
1693
1694
1695
1696
1697
1698
1699
1700
1701
1702
1703
1704
1705
1706
1707
async def trace_publish_lineage(
    self,
    *,
    publish_helper=None,
    parent_auth: DomoAuth | None = None,
    parent_auth_retrieval_fn: Callable | None = None,
    max_depth: int = 100,
    context: RouteContext | None = None,
    **context_kwargs,
) -> list[DomoLineage_Link]:
    """Build the full publish-lineage chain for a federated published entity.

    Orchestrates:
    1. Ensure subscription is resolved (via FederationContext).
    2. Resolve publisher auth.
    3. Fetch parent publication.
    4. Call ``_resolve_publisher_entity`` (overridable hook) to find
       the publisher-side entity that corresponds to ``self.parent``.
    5. Fetch the publisher entity's own lineage.
    6. Assemble the chain:
       subscriber entity → subscription → publication → publisher entity
       → publisher entity lineage.

    Subclasses override ``_resolve_publisher_entity`` to customise step 4.
    For example ``DomoLineage_Card`` handles cards that are indirectly
    published via pages.

    Args:
        publish_helper: FederationContext (defaults to ``self.Federation``).
        parent_auth: Explicit publisher auth (alternative to retrieval fn).
        parent_auth_retrieval_fn: Callable(domain) → DomoAuth.
        max_depth: Depth limit for publisher entity lineage.
        context: Route context for API calls.

    Returns:
        Flat list of lineage links representing the full chain.
    """
    from ...DomoEverywhere.lineage import (
        DomoLineageLink_Publication,
        DomoLineageLink_Subscription,
    )

    # --- 1. Ensure subscription ---
    helper = publish_helper or self.Federation
    if not helper:
        raise ValueError(
            "No federation helper available. "
            "Call check_is_published() first or provide publish_helper."
        )

    subscription = await helper.ensure_subscription(
        retrieve_parent_auth_fn=parent_auth_retrieval_fn,
        parent_auth=parent_auth,
        entity_type=self.parent.entity_type,
        entity_id=str(self.parent.id),
        context=context,
        **context_kwargs,
    )
    if not subscription:
        raise DomoError(
            f"Failed to resolve subscription for published entity "
            f"{self.parent.entity_type}:{self.parent.id}"
        )

    # --- 2. Resolve publisher auth ---
    publisher_auth = await helper.get_publisher_auth(
        parent_auth=parent_auth,
        context=context,
        **context_kwargs,
    )

    # --- 3. Fetch parent publication ---
    publication = await helper.get_parent_publication(
        parent_auth=publisher_auth,
        context=context,
        **context_kwargs,
    )

    # --- 4. Resolve publisher entity (overridable hook) ---
    subscriber_domain = subscription.subscriber_domain
    resolution = await self._resolve_publisher_entity(
        publication=publication,
        publisher_auth=publisher_auth,
        subscriber_domain=subscriber_domain,
        context=context,
        **context_kwargs,
    )
    publisher_entity = resolution.entity
    intermediate_entities = resolution.intermediate_entities

    # --- 5. Fetch publisher entity lineage ---
    publisher_entity_lineage: list[DomoLineage_Link] = []
    publisher_entity_link = None

    if publisher_entity:
        publisher_entity_lineage = await publisher_entity.Lineage.get(
            return_raw=False,
            max_depth=max_depth,
            context=context,
            **context_kwargs,
        )

        # Find the publisher entity's own link within its lineage — it
        # already has the correct direct dependencies wired by
        # get_datacenter_lineage().  Only create a new link as a fallback.
        publisher_entity_link = next(
            (
                link
                for link in publisher_entity_lineage
                if str(link.id) == str(publisher_entity.id)
            ),
            None,
        )

        if publisher_entity_link is None:
            link_type = _map_entity_type(publisher_entity.entity_type)
            link_cls = _get_lineage_link_class(link_type)
            publisher_entity_link = link_cls(
                auth=publisher_auth,
                id=str(publisher_entity.id),
                entity=publisher_entity,
                _type=None,
                dependencies=[],
                dependents=[],
            )

    # --- 5b. Build intermediate links (e.g., page for indirect card publication) ---
    # Intermediate entities sit between publication and publisher_entity.
    # Build in reverse so each link's dependencies point to the next in the chain.
    intermediate_links: list[DomoLineage_Link] = []
    next_dependency = publisher_entity_link

    for intermediate_entity in reversed(intermediate_entities):
        inter_link_type = _map_entity_type(intermediate_entity.entity_type)
        inter_link_cls = _get_lineage_link_class(inter_link_type)
        inter_link = inter_link_cls(
            auth=publisher_auth,
            id=str(intermediate_entity.id),
            entity=intermediate_entity,
            _type=None,
            dependencies=[next_dependency] if next_dependency else [],
            dependents=[],
        )
        if next_dependency:
            next_dependency.dependents = [inter_link]
        intermediate_links.insert(0, inter_link)
        next_dependency = inter_link

    # Determine what the publication link should point to
    publication_first_dep = (
        intermediate_links[0] if intermediate_links else publisher_entity_link
    )

    # --- 6. Assemble the chain ---
    publication_link = DomoLineageLink_Publication(
        auth=publisher_auth,
        id=str(publication.id),
        entity=publication,
        _type="PUBLICATION",
        dependents=[],
        dependencies=([publication_first_dep] if publication_first_dep else []),
    )

    subscription_link = DomoLineageLink_Subscription(
        auth=self.parent.auth,
        id=str(subscription.id),
        entity=subscription,
        _type="SUBSCRIPTION",
        dependents=[],
        dependencies=[publication_link],
    )

    subscriber_link = _get_lineage_link_class(
        _map_entity_type(self.parent.entity_type)
    )(
        auth=self.parent.auth,
        id=str(self.parent.id),
        entity=self.parent,
        _type=None,
        dependencies=[subscription_link],
        dependents=[],
    )

    # Wire back-pointers
    subscription_link.dependents = [subscriber_link]
    publication_link.dependents = [subscription_link]
    if publication_first_dep:
        publication_first_dep.dependents = [publication_link]

    chain = [
        subscriber_link,
        subscription_link,
        publication_link,
    ]
    chain.extend(intermediate_links)
    # publisher_entity_link is already within publisher_entity_lineage
    # (found by id match), so only extend — don't append separately.
    chain.extend(publisher_entity_lineage)

    return chain
DomoLineage_Link(
    auth: DomoAuth,
    id: str,
    entity: Any,
    _type: str | None = None,
    dependencies: list[DomoLineage_Link] = list(),
    dependents: list[DomoLineage_Link] = list(),
)

Bases: DomoBase

Represents a link in the dependency lineage graph.

Terminology: Uses dependency graph terminology:

  • dependencies: Upstream dependencies (what this entity depends on) Example: A card's dependencies are the datasets it uses for data Example: A view's dependencies are the datasets it depends on

  • dependents: Downstream dependents (what depends on this entity) Example: A dataset's dependents are the cards/pages/views that use it

This matches standard dependency graph terminology and the Domo API structure. Note: This is different from composition (e.g., a page "contains" cards).

The datacenter lineage payload frequently references entities only by ID/type; we therefore construct placeholder links with entity=None and materialize the real entity later via :meth:get_entity. The _type cache allows those placeholder links to behave consistently (hashing, comparisons, rendering) until their backing entity is fetched.

type property

type: str

Get the lineage type, derived from entity if available, otherwise from cached _type.

from_dict classmethod

from_dict(
    lineage_api_obj: dict[str, Any], auth: DomoAuth
) -> DomoLineage_Link

Create a DomoLineage_Link from a lineage API object dict. Gets the appropriate link class based on the 'type' field.

Source code in src/crew_dcs/classes/subentity/lineage/lineage_link.py
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
@classmethod
def from_dict(
    cls,
    lineage_api_obj: dict[str, Any],
    auth: DomoAuth,
) -> DomoLineage_Link:
    """
    Create a DomoLineage_Link from a lineage API object dict.
    Gets the appropriate link class based on the 'type' field.
    """
    # Get the specific link class for this type
    link_cls = _get_lineage_link_class(lineage_api_obj["type"])

    return link_cls(
        auth=auth,
        id=lineage_api_obj["id"],
        entity=None,  # Placeholder - entity not loaded yet
        _type=lineage_api_obj["type"],
        dependents=cls._create_lineage_links_from_dicts(
            lineage_api_obj.get("children", []), auth=auth
        ),
        dependencies=cls._create_lineage_links_from_dicts(
            lineage_api_obj.get("parents", []), auth=auth
        ),
    )

get_entity abstractmethod async

get_entity(
    debug_api: bool = False,
    session: AsyncClient | None = None,
    *,
    context=None
)

Get the entity associated with this lineage link. Uses self.id and self.auth to fetch the entity.

This method should be implemented by subclasses to return the appropriate entity.

Source code in src/crew_dcs/classes/subentity/lineage/lineage_link.py
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
@abstractmethod
async def get_entity(
    self,
    debug_api: bool = False,
    session: httpx.AsyncClient | None = None,
    *,
    context=None,
):
    """
    Get the entity associated with this lineage link.
    Uses self.id and self.auth to fetch the entity.

    This method should be implemented by subclasses to return the appropriate entity.
    """
    raise NotImplementedError("Subclasses must implement this method.")

FederatedLineageAuthRequiredError

FederatedLineageAuthRequiredError(
    *,
    entity_id: str,
    entity_type: str,
    domo_instance: str | None,
    partial_lineage: list[DomoLineage_Link] | None = None
)

Bases: ClassError

Raised when federated lineage is requested without publisher auth.

.. deprecated:: This exception is no longer raised internally. Partial-lineage situations are now handled by logging a warning instead. Kept for backward-compatible except clauses in downstream code.

Source code in src/crew_dcs/classes/subentity/lineage/base.py
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
def __init__(
    self,
    *,
    entity_id: str,
    entity_type: str,
    domo_instance: str | None,
    partial_lineage: list[DomoLineage_Link] | None = None,
):
    self.partial_lineage = partial_lineage or []
    super().__init__(
        message=(
            "Federated lineage requires parent_auth or parent_auth_retrieval_fn to"
            " retrieve publisher-side dependencies."
        ),
        entity_id=entity_id,
        parent_class=entity_type,
        domo_instance=domo_instance,
    )

FederationContext dataclass

FederationContext(
    parent: DomoEntityLineageProtocol | None = None,
    subscription: DomoSubscriptionProtocol | None = None,
    parent_publication: (
        DomoPublicationProtocol | None
    ) = None,
    publisher_entity: (
        DomoEntityLineageProtocol | None
    ) = None,
    _parent_auth_fn: (
        Callable[[str], DomoAuth | Awaitable[DomoAuth]]
        | None
    ) = None,
    _parent_auth: DomoAuth | None = None,
    _content_type: str | None = None,
    _entity_id: str | None = None,
    _auth: DomoAuth | None = None,
)

Bases: DomoBase

Holds resolved federation state for federated entities.

Attached to subscriber-side entities that are copies from another instance. Manages subscription/publication discovery and caching.

Can be attached to a parent entity via constructor, or created standalone using from_entity_id() for probing publish state without a full entity.

auth property

auth: DomoAuth

Get auth from parent entity or standalone _auth.

entity_id property

entity_id: str

Get entity ID from stored value or parent.

check_if_published async

check_if_published(
    *,
    retrieve_parent_auth_fn: Callable[
        [str], DomoAuth | Awaitable[DomoAuth]
    ],
    entity_type: str,
    entity_id: str | None = None,
    context: RouteContext | None = None,
    session: AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
    **context_kwargs
) -> bool

Discover whether the entity participates in a subscription.

Source code in src/crew_dcs/classes/subentity/lineage/federation_context.py
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
async def check_if_published(
    self,
    *,
    retrieve_parent_auth_fn: Callable[[str], DomoAuth | Awaitable[DomoAuth]],
    entity_type: str,
    entity_id: str | None = None,
    context: RouteContext | None = None,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
    **context_kwargs,
) -> bool:
    """Discover whether the entity participates in a subscription."""
    from .publish_resolver import PublishResolver

    if not retrieve_parent_auth_fn:
        raise ValueError(
            "retrieve_parent_auth_fn is required to resolve published entities."
        )

    # Build context from provided parameters
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        **context_kwargs,
    )

    self._parent_auth_fn = retrieve_parent_auth_fn
    self._content_type = entity_type
    if entity_id:
        self._entity_id = entity_id

    target_id = self.entity_id

    await logger.debug(
        f"Checking if {entity_type} {target_id} is published",
        extra={
            "entity_type": entity_type,
            "entity_id": target_id,
            "subscriber_instance": self.auth.domo_instance,
        },
    )

    resolver = PublishResolver(
        subscriber_auth=self.auth,
        parent_auth_retrieval_fn=retrieve_parent_auth_fn,
        session=context.session,
        debug_api=context.debug_api,
        max_subscriptions_to_check=max_subscriptions_to_check,
    )

    try:
        # PublishResolver returns raw dict - hydrate into DomoSubscription
        from ...DomoEverywhere.core import DomoSubscription

        subscription_data = await resolver.get_subscription_for_entity(
            entity_type=entity_type,
            subscriber_entity_id=target_id,
        )

        self.subscription = DomoSubscription.from_dict(
            auth=self.auth,
            parent_publication=None,
            obj=subscription_data,
        )
        await logger.info(
            f"✅ {entity_type} {target_id} is published",
            extra={
                "entity_type": entity_type,
                "entity_id": target_id,
                "subscription_id": (
                    self.subscription.id if self.subscription else None
                ),
            },
        )
    except ValueError as exc:
        await logger.debug(
            f"❌ {entity_type} {target_id} is not published: {exc}",
            extra={
                "entity_type": entity_type,
                "entity_id": target_id,
                "error": str(exc),
            },
        )
        self.subscription = None

    return self.is_published

ensure_subscription async

ensure_subscription(
    *,
    retrieve_parent_auth_fn: (
        Callable[[str], DomoAuth | Awaitable[DomoAuth]]
        | None
    ) = None,
    parent_auth: DomoAuth | None = None,
    entity_type: str | None = None,
    entity_id: str | None = None,
    session: AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
    context: RouteContext | None = None,
    **context_kwargs
) -> DomoSubscriptionProtocol | None

Ensure subscription data is loaded, running discovery if needed.

Parameters:

Name Type Description Default
retrieve_parent_auth_fn Callable[[str], DomoAuth | Awaitable[DomoAuth]] | None

Callable to retrieve parent auth by domain

None
parent_auth DomoAuth | None

Pre-existing parent auth (alternative to retrieval function)

None
entity_type str | None

Entity type (DATA_SOURCE, CARD, PAGE, etc.)

None
entity_id str | None

The entity identifier

None
session AsyncClient | None

HTTP session for API calls

None
debug_api bool

Enable debug logging

False
max_subscriptions_to_check int | None

Limit subscription search

None

Returns:

Type Description
DomoSubscriptionProtocol | None

The subscription object if found

Raises:

Type Description
ValueError

If no auth method available

Source code in src/crew_dcs/classes/subentity/lineage/federation_context.py
219
220
221
222
223
224
225
226
227
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
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
async def ensure_subscription(
    self,
    *,
    retrieve_parent_auth_fn: (
        Callable[[str], DomoAuth | Awaitable[DomoAuth]] | None
    ) = None,
    parent_auth: DomoAuth | None = None,
    entity_type: str | None = None,
    entity_id: str | None = None,
    session: httpx.AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
    context: RouteContext | None = None,
    **context_kwargs,
) -> DomoSubscriptionProtocol | None:
    """Ensure subscription data is loaded, running discovery if needed.

    Args:
        retrieve_parent_auth_fn: Callable to retrieve parent auth by domain
        parent_auth: Pre-existing parent auth (alternative to retrieval function)
        entity_type: Entity type (DATA_SOURCE, CARD, PAGE, etc.)
        entity_id: The entity identifier
        session: HTTP session for API calls
        debug_api: Enable debug logging
        max_subscriptions_to_check: Limit subscription search

    Returns:
        The subscription object if found

    Raises:
        ValueError: If no auth method available
    """

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

    if self.subscription:
        return self.subscription

    # Store parent_auth if provided
    if parent_auth:
        self._parent_auth = parent_auth

    fn = retrieve_parent_auth_fn or self._parent_auth_fn
    # Allow proceeding if we have either a retrieval function OR direct parent_auth
    if not fn and not self._parent_auth:
        raise ValueError(
            "retrieve_parent_auth_fn or parent_auth must be provided to resolve subscriptions."
        )

    # Determine content type - try stored, then parent if available
    content_type = entity_type or self._content_type
    if not content_type and self.parent is not None:
        content_type = self.parent.entity_type

    if not content_type:
        raise ValueError("entity_type must be provided or set via _content_type")

    # Use entity_id property which handles both standalone and parent modes
    target_id = entity_id or self.entity_id

    # If we have parent_auth but no fn, create a simple lambda that returns the auth
    effective_fn = fn
    if not effective_fn and self._parent_auth:
        effective_fn = lambda _domain, **_kwargs: self._parent_auth  # noqa: E731

    await self.check_if_published(
        retrieve_parent_auth_fn=effective_fn,
        entity_type=content_type,
        entity_id=target_id,
        context=context,
        session=session,
        debug_api=debug_api,
        max_subscriptions_to_check=max_subscriptions_to_check,
    )
    return self.subscription

from_entity_id classmethod

from_entity_id(
    *, auth: DomoAuth, entity_id: str, entity_type: str
) -> FederationContext

Create a federation context from entity identifiers (no parent entity required).

Use this factory when you need to check publish state without constructing a full entity object.

Parameters:

Name Type Description Default
auth DomoAuth

Authentication for the subscriber instance

required
entity_id str

The entity identifier to check

required
entity_type str

The entity type (DATA_SOURCE, CARD, PAGE, etc.)

required

Returns:

Type Description
FederationContext

FederationContext instance ready for publish state checks

Source code in src/crew_dcs/classes/subentity/lineage/federation_context.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
@classmethod
def from_entity_id(
    cls,
    *,
    auth: DomoAuth,
    entity_id: str,
    entity_type: str,
) -> FederationContext:
    """Create a federation context from entity identifiers (no parent entity required).

    Use this factory when you need to check publish state without constructing
    a full entity object.

    Args:
        auth: Authentication for the subscriber instance
        entity_id: The entity identifier to check
        entity_type: The entity type (DATA_SOURCE, CARD, PAGE, etc.)

    Returns:
        FederationContext instance ready for publish state checks
    """
    return cls(
        parent=None,
        _auth=auth,
        _entity_id=str(entity_id),
        _content_type=entity_type,
    )

get_parent_publication async

get_parent_publication(
    *,
    parent_auth: DomoAuth | None = None,
    is_fetch_content_details: bool = True,
    context: RouteContext | None = None,
    **context_kwargs
) -> DomoPublicationProtocol

Fetch and cache the parent publication for this entity.

Parameters:

Name Type Description Default
parent_auth DomoAuth | None

Pre-existing publisher auth

None
is_fetch_content_details bool

If True, automatically call get_content_details() on the publication before returning. This ensures content mapping is available for indirect resolution. Default: True.

True
context RouteContext | None

Route context

None
**context_kwargs

Additional context args

{}

Returns:

Type Description
DomoPublicationProtocol

DomoPublication with content details loaded (if is_fetch_content_details=True)

Source code in src/crew_dcs/classes/subentity/lineage/federation_context.py
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
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
async def get_parent_publication(
    self,
    *,
    parent_auth: DomoAuth | None = None,
    is_fetch_content_details: bool = True,
    context: RouteContext | None = None,
    **context_kwargs,
) -> DomoPublicationProtocol:
    """Fetch and cache the parent publication for this entity.

    Args:
        parent_auth: Pre-existing publisher auth
        is_fetch_content_details: If True, automatically call get_content_details()
            on the publication before returning. This ensures content mapping is
            available for indirect resolution. Default: True.
        context: Route context
        **context_kwargs: Additional context args

    Returns:
        DomoPublication with content details loaded (if is_fetch_content_details=True)
    """
    context = RouteContext.build_context(
        context=context,
        **context_kwargs,
    )

    from ...DomoEverywhere.core import DomoPublication

    if not self.subscription:
        raise ValueError(
            "Subscription must be loaded before fetching parent publication."
        )

    publisher_auth = await self._resolve_parent_auth(parent_auth, context=context)

    if not self.parent_publication:
        self.parent_publication = await DomoPublication.get_by_id(
            publication_id=self.subscription.publication_id,
            auth=publisher_auth,
            context=context,
        )

    # Automatically fetch content details if requested
    if is_fetch_content_details and self.parent_publication:
        await self.parent_publication.get_content_details(
            subscriber_domain=self.auth.domo_instance,
            context=context,
        )

    return self.parent_publication

get_publisher_auth async

get_publisher_auth(
    parent_auth: Any = None,
    context: RouteContext | None = None,
    **context_kwargs
) -> DomoAuth

Public helper to resolve publisher authentication.

Source code in src/crew_dcs/classes/subentity/lineage/federation_context.py
343
344
345
346
347
348
349
350
351
352
353
354
355
async def get_publisher_auth(
    self,
    parent_auth: Any = None,
    context: RouteContext | None = None,
    **context_kwargs,
) -> DomoAuth:
    """Public helper to resolve publisher authentication."""

    context = RouteContext.build_context(
        context=context,
        **context_kwargs,
    )
    return await self._resolve_parent_auth(parent_auth, context=context)

hydrate_from_existing

hydrate_from_existing(
    *,
    subscription: DomoSubscriptionProtocol,
    parent_auth_retrieval_fn: (
        Callable[[str], DomoAuth | Awaitable[DomoAuth]]
        | None
    ) = None,
    parent_auth: DomoAuth | None = None,
    content_type: str,
    entity_id: str
)

Attach an already-discovered subscription to this helper.

Parameters:

Name Type Description Default
subscription DomoSubscriptionProtocol

The subscription object

required
parent_auth_retrieval_fn Callable[[str], DomoAuth | Awaitable[DomoAuth]] | None

Callable to retrieve parent auth by domain

None
parent_auth DomoAuth | None

Pre-existing parent auth (alternative to retrieval function)

None
content_type str

Entity type (DATA_SOURCE, CARD, PAGE, etc.)

required
entity_id str

The entity identifier

required
Source code in src/crew_dcs/classes/subentity/lineage/federation_context.py
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
def hydrate_from_existing(
    self,
    *,
    subscription: DomoSubscriptionProtocol,
    parent_auth_retrieval_fn: (
        Callable[[str], DomoAuth | Awaitable[DomoAuth]] | None
    ) = None,
    parent_auth: DomoAuth | None = None,
    content_type: str,
    entity_id: str,
):
    """Attach an already-discovered subscription to this helper.

    Args:
        subscription: The subscription object
        parent_auth_retrieval_fn: Callable to retrieve parent auth by domain
        parent_auth: Pre-existing parent auth (alternative to retrieval function)
        content_type: Entity type (DATA_SOURCE, CARD, PAGE, etc.)
        entity_id: The entity identifier
    """
    self.subscription = subscription
    self._parent_auth_fn = parent_auth_retrieval_fn
    self._parent_auth = parent_auth
    self._content_type = content_type
    self._entity_id = entity_id

PublishResolver dataclass

PublishResolver(
    subscriber_auth: DomoAuth,
    parent_auth_retrieval_fn: Callable[
        [str, RouteContext], Any | Awaitable[Any]
    ],
    session: AsyncClient | None = None,
    debug_api: bool = False,
    max_subscriptions_to_check: int | None = None,
)

Resolve subscriptions for subscriber-side entities.

get_subscription_for_entity async

get_subscription_for_entity(
    *,
    entity_type: ContentType,
    subscriber_entity_id: str,
    context: RouteContext | None = None,
    **context_kwargs: Any
) -> dict[str, Any]

Find the subscription that contains the given subscriber entity.

Returns:

Type Description
dict[str, Any]

Raw subscription summary dict from the API. Callers should hydrate

dict[str, Any]

this into a DomoSubscription instance in the classes layer.

Source code in src/crew_dcs/classes/subentity/lineage/publish_resolver.py
 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
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
226
227
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
269
270
271
272
273
274
275
276
277
278
async def get_subscription_for_entity(  # noqa: C901
    self,
    *,
    entity_type: ContentType,
    subscriber_entity_id: str,
    context: RouteContext | None = None,
    **context_kwargs: Any,
) -> dict[str, Any]:
    """Find the subscription that contains the given subscriber entity.

    Returns:
        Raw subscription summary dict from the API. Callers should hydrate
        this into a DomoSubscription instance in the classes layer.
    """
    if not self.parent_auth_retrieval_fn:
        raise ValueError(
            "parent_auth_retrieval_fn is required to resolve subscriptions. "
            "This function should accept a publisher_domain and return DomoAuth "
            "for that instance (sync or async)."
        )

    context = RouteContext.build_context(
        context=context,
        session=self.session,
        debug_api=self.debug_api,
        debug_num_stacks_to_drop=2,
        parent_class=self.__class__.__name__,
    )

    summaries_res = await publish_routes.get_subscription_summaries(
        auth=self.subscriber_auth,
        context=context,
    )

    if not summaries_res.is_success or not summaries_res.response:
        raise ValueError(
            f"Failed to retrieve subscriptions for instance "
            f"{self.subscriber_auth.domo_instance}"
        )

    subscriptions_checked = 0
    total_subscriptions = len(summaries_res.response)
    target_id = str(subscriber_entity_id)
    normalized_entity_type = (entity_type or "").upper()
    acceptable_content_types = {normalized_entity_type}
    # entity_type is now consistently "DATA_SOURCE" (not "DATASET"),
    # but the API may return either value, so accept both aliases
    if normalized_entity_type in ("DATA_SOURCE", "DATASET"):
        acceptable_content_types.update({"DATA_SOURCE", "DATASET"})

    await logger.debug(
        f"Found {total_subscriptions} subscriptions",
        extra={
            "entity_type": entity_type,
            "subscriber_entity_id": subscriber_entity_id,
            "subscriber_instance": self.subscriber_auth.domo_instance,
            "total_subscriptions": total_subscriptions,
        },
    )

    for summary in summaries_res.response:
        if (
            self.max_subscriptions_to_check is not None
            and subscriptions_checked >= self.max_subscriptions_to_check
        ):
            await logger.debug(
                f"Reached max subscriptions limit ({self.max_subscriptions_to_check}), stopping search",
                extra={"max_subscriptions": self.max_subscriptions_to_check},
            )
            break

        subscription_id = summary.get("subscriptionId")
        publication_id = summary.get("publicationId")
        subscriber_domain = summary.get("subscriberDomain")
        publisher_domain = summary.get("publisherDomain")

        if (
            not subscription_id
            or not publication_id
            or not subscriber_domain
            or not publisher_domain
        ):
            continue

        subscriptions_checked += 1

        sub_idx = (
            self.max_subscriptions_to_check
            if self.max_subscriptions_to_check
            else total_subscriptions
        )
        await logger.debug(
            f"Checking subscription {subscriptions_checked}/{sub_idx}: {subscription_id} - {publisher_domain}",
            extra={
                "subscription_id": subscription_id,
                "publisher_domain": publisher_domain,
                "publication_id": publication_id,
                "subscription_number": subscriptions_checked,
            },
        )

        publisher_auth = None

        try:
            publisher_auth = await self._get_publisher_auth(
                publisher_domain, context=context
            )
        except DomoError as exc:  # pragma: no cover - best effort logging
            await logger.warning(
                f"Unable to fetch auth for publisher {publisher_domain}: {exc}",
                extra={
                    "publisher_domain": publisher_domain,
                    "error": str(exc),
                    "error_type": type(exc).__name__,
                },
            )
            continue
        except (
            ValueError,
            KeyError,
            LookupError,
        ) as exc:  # pragma: no cover - auth lookup failure
            await logger.warning(
                f"Auth credentials not found for publisher {publisher_domain}: {exc}",
                extra={
                    "publisher_domain": publisher_domain,
                    "error": str(exc),
                    "error_type": type(exc).__name__,
                },
            )
            continue
        except (
            Exception  # noqa: BLE001
        ) as exc:  # pragma: no cover - unexpected error
            await logger.error(
                f"❌ Unexpected error fetching auth for publisher {publisher_domain}: {exc}",
                extra={
                    "publisher_domain": publisher_domain,
                    "error": str(exc),
                    "error_type": type(exc).__name__,
                },
                exc_info=True,
            )
            continue

        if not publisher_auth:
            continue

        try:
            content_res = await publish_routes.get_subscriber_content_details(
                auth=publisher_auth,
                publication_id=publication_id,
                subscriber_instance=subscriber_domain,
                context=context,
            )

        except DomoError as exc:  # pragma: no cover - best effort logging
            await logger.warning(
                f"Unable to fetch subscriber content details for {subscription_id}: {exc}",
                extra={
                    "subscription_id": subscription_id,
                    "error": str(exc),
                    "error_type": type(exc).__name__,
                },
            )
            continue
        except (
            Exception  # noqa: BLE001
        ) as exc:  # pragma: no cover - unexpected error
            await logger.error(
                f"❌ Unexpected error fetching subscriber content for {subscription_id}: {exc}",
                extra={
                    "subscription_id": subscription_id,
                    "error": str(exc),
                    "error_type": type(exc).__name__,
                },
                exc_info=True,
            )
            continue

        if not (content_res.is_success and content_res.response):
            await logger.debug(
                f"No subscriber content details for subscription {subscription_id}",
                extra={"subscription_id": subscription_id},
            )
            continue

        for item in content_res.response:
            content_type = (item.get("contentType") or "").upper()
            if (
                content_type in acceptable_content_types
                and str(item.get("subscriberObjectId")) == target_id
            ):
                await logger.info(
                    f"✅ Found {normalized_entity_type} {target_id} in subscription {subscription_id}",
                    extra={
                        "entity_type": normalized_entity_type,
                        "entity_id": target_id,
                        "subscription_id": subscription_id,
                        "publisher_domain": publisher_domain,
                    },
                )
                # Return raw dict - caller hydrates into DomoSubscription
                return summary

    # No subscription found
    if self.max_subscriptions_to_check:
        raise ValueError(
            f"Entity {subscriber_entity_id} (type {entity_type}) is not part of "
            f"any subscription after checking {subscriptions_checked} "
            f"subscriptions. Try increasing max_subscriptions_to_check "
            f"(currently {self.max_subscriptions_to_check})."
        )

    raise ValueError(
        f"Entity {subscriber_entity_id} (type {entity_type}) is not part of any "
        f"subscription after checking all {subscriptions_checked} subscriptions."
    )

PublisherResolution dataclass

PublisherResolution(
    entity: Any | None = None,
    intermediate_entities: list[Any] = list(),
)

Result of resolving a publisher entity, with optional intermediate entities.

Attributes:

Name Type Description
entity Any | None

The resolved publisher-side entity (e.g., a card or dataset), or None.

intermediate_entities list[Any]

Entities between the publication and the resolved entity in the lineage chain (e.g., a page when a card is published indirectly as part of a page).

get_lineage_type

get_lineage_type(class_name: str) -> str

Get the lineage type string for a given class name.

Parameters:

Name Type Description Default
class_name str

The class name to look up

required

Returns:

Type Description
str

The lineage type string (e.g., "DATAFLOW", "DATA_SOURCE")

Raises:

Type Description
ValueError

If the class name is not registered

Source code in src/crew_dcs/classes/subentity/lineage/base.py
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
def get_lineage_type(class_name: str) -> str:
    """Get the lineage type string for a given class name.

    Args:
        class_name: The class name to look up

    Returns:
        The lineage type string (e.g., "DATAFLOW", "DATA_SOURCE")

    Raises:
        ValueError: If the class name is not registered
    """
    lineage_type = _LINEAGE_TYPE_REGISTRY.get(class_name)
    if lineage_type is None:
        raise ValueError(
            f"No lineage type registered for class '{class_name}'. "
            f"Known classes: {sorted(_LINEAGE_TYPE_REGISTRY.keys())}"
        )
    return lineage_type

register_lineage

register_lineage(*parent_class_names: str)

Decorator to register a DomoLineage subclass for parent class names.

Parameters:

Name Type Description Default
*parent_class_names str

One or more parent class names that should use this lineage class

()
Example

@register_lineage("DomoPage", "DomoPage_Default") @dataclass class DomoLineage_Page(DomoLineage): ...

Source code in src/crew_dcs/classes/subentity/lineage/base.py
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
def register_lineage(*parent_class_names: str):
    """Decorator to register a DomoLineage subclass for parent class names.

    Args:
        *parent_class_names: One or more parent class names that should use this lineage class

    Example:
        @register_lineage("DomoPage", "DomoPage_Default")
        @dataclass
        class DomoLineage_Page(DomoLineage):
            ...
    """

    def decorator(cls: type[DomoLineage]) -> type[DomoLineage]:
        for parent_class_name in parent_class_names:
            _LINEAGE_REGISTRY[parent_class_name] = cls
        return cls

    return decorator
register_lineage_link(link_type: str)

Decorator to register a DomoLineage_Link subclass by datacenter type string.

Source code in src/crew_dcs/classes/subentity/lineage/lineage_link.py
273
274
275
276
277
278
279
280
def register_lineage_link(link_type: str):
    """Decorator to register a DomoLineage_Link subclass by datacenter type string."""

    def decorator(cls: type[DomoLineage_Link]) -> type[DomoLineage_Link]:
        _LINEAGE_LINK_REGISTRY[link_type] = cls
        return cls

    return decorator

register_lineage_type

register_lineage_type(*class_names: str, lineage_type: str)

Decorator to register lineage type for entity classes.

This decorator registers one or more entity class names with their corresponding lineage type string (e.g., "DATAFLOW", "DATA_SOURCE", "PAGE").

Parameters:

Name Type Description Default
*class_names str

One or more class names that should use this lineage type

()
lineage_type str

The lineage type string (e.g., "DATAFLOW", "DATA_SOURCE")

required
Example

@register_lineage_type("DomoDataflow", lineage_type="DATAFLOW") @dataclass class DomoDataflow(DomoEntity_w_Lineage): ...

Source code in src/crew_dcs/classes/subentity/lineage/base.py
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
def register_lineage_type(*class_names: str, lineage_type: str):
    """Decorator to register lineage type for entity classes.

    This decorator registers one or more entity class names with their
    corresponding lineage type string (e.g., "DATAFLOW", "DATA_SOURCE", "PAGE").

    Args:
        *class_names: One or more class names that should use this lineage type
        lineage_type: The lineage type string (e.g., "DATAFLOW", "DATA_SOURCE")

    Example:
        @register_lineage_type("DomoDataflow", lineage_type="DATAFLOW")
        @dataclass
        class DomoDataflow(DomoEntity_w_Lineage):
            ...
    """

    def decorator(cls):
        for class_name in class_names:
            _LINEAGE_TYPE_REGISTRY[class_name] = lineage_type
        return cls

    return decorator

Modules