Skip to content

get_data

get_data

GetDataError

GetDataError(
    message: str,
    url: str,
    response: ResponseGetData | None = None,
)

Bases: DomoError

Raised for Domo-specific get_data failures (e.g. VPN / allowlist HTML).

Source code in src/crew_dcs/client/get_data.py
48
49
50
51
52
53
54
55
def __init__(
    self,
    message: str,
    url: str,
    response: rgd.ResponseGetData | None = None,
) -> None:
    super().__init__(message=message, domo_instance=url)
    self.response = response

create_headers

create_headers(
    auth: DomoAuth = None,
    content_type: str | None = None,
    headers: dict | None = None,
) -> dict

Creates default headers for interacting with Domo APIs.

Source code in src/crew_dcs/client/get_data.py
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
def create_headers(
    auth: "dmda.DomoAuth" = None,  # The authentication object containing the Domo API token.
    content_type: (
        str | None
    ) = None,  # The content type for the request. Defaults to None.
    headers: dict
    | None = None,  # Any additional headers for the request. Defaults to None.
) -> dict:  # The headers for the request.
    """
    Creates default headers for interacting with Domo APIs.
    """

    if headers is None:
        headers = {}

    headers = {
        "Content-Type": content_type or "application/json",
        "Connection": "keep-alive",
        "accept": "application/json, text/plain",
        **headers,
    }
    if auth:
        headers.update(**auth.auth_header)
    return headers

get_data async

get_data(
    url: str,
    method: str,
    auth: DomoAuth = None,
    content_type: str | None = None,
    headers: dict | None = None,
    body: dict | list | str | None = None,
    params: dict | None = None,
    context: RouteContext | None = None,
    debug_api: bool | None = None,
    session: AsyncClient | None = None,
    return_raw: bool | None = None,
    is_follow_redirects: bool | None = None,
    timeout: int | None = None,
    parent_class: str | None = None,
    debug_num_stacks_to_drop: int | None = None,
    is_verify: bool | None = None,
    dry_run: bool | None = None,
    return_on_domo_block: bool | None = None,
    **kwargs
) -> ResponseGetData

Asynchronously performs an HTTP request to retrieve data from a Domo API endpoint.

Parameters:

Name Type Description Default
url str

API endpoint URL

required
method str

HTTP method (GET, POST, PUT, DELETE, etc.)

required
auth DomoAuth

Authentication object containing credentials

None
content_type str | None

Optional content type header

None
headers dict | None

Additional HTTP headers

None
body dict | list | str | None

Request body (dict, list, or string)

None
params dict | None

Query parameters

None
context RouteContext | None

Optional RouteContext with debug/session settings (takes precedence over individual params)

None
debug_api bool | None

Enable API debugging (overridden by context if provided)

None
session AsyncClient | None

Optional httpx client session (overridden by context if provided)

None
return_raw bool | None

Attach the httpx.Response on ResponseGetData.raw_response (does not suppress errors).

None
is_follow_redirects bool | None

Follow HTTP redirects

None
timeout int | None

Request timeout in seconds

None
parent_class str | None

Optional parent class name for debugging (overridden by context if provided)

None
debug_num_stacks_to_drop int | None

Number of stack frames to drop in debug output (overridden by context if provided)

None
is_verify bool | None

SSL verification flag

None
dry_run bool | None

If True, return request parameters without executing

None
return_on_domo_block bool | None

If True, return a ResponseGetData for VPN/allowlist HTML instead of raising GetDataError (still logged). Transport/httpx errors are always re-raised as-is after logging.

None

Returns:

Type Description
ResponseGetData

ResponseGetData object containing the response

Source code in src/crew_dcs/client/get_data.py
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
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
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
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
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
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
@dmce.run_with_retry()
@log_call(
    action_name="get_data",
    level_name="client",
    log_level="DEBUG",
    config=LogDecoratorConfig(result_processor=ResponseGetDataProcessor()),
    color="light_blue",
)
async def get_data(  # noqa: C901
    url: str,
    method: str,
    auth: "dmda.DomoAuth" = None,
    content_type: str | None = None,
    headers: dict | None = None,
    body: dict | list | str | None = None,
    params: dict | None = None,
    context: RouteContext | None = None,
    debug_api: bool | None = None,
    session: httpx.AsyncClient | None = None,
    return_raw: bool | None = None,
    is_follow_redirects: bool | None = None,
    timeout: int | None = None,
    parent_class: str | None = None,
    debug_num_stacks_to_drop: int | None = None,
    is_verify: bool | None = None,
    dry_run: bool | None = None,
    return_on_domo_block: bool | None = None,
    **kwargs,
) -> rgd.ResponseGetData:
    """Asynchronously performs an HTTP request to retrieve data from a Domo API endpoint.

    Args:
        url: API endpoint URL
        method: HTTP method (GET, POST, PUT, DELETE, etc.)
        auth: Authentication object containing credentials
        content_type: Optional content type header
        headers: Additional HTTP headers
        body: Request body (dict, list, or string)
        params: Query parameters
        context: Optional RouteContext with debug/session settings (takes precedence over individual params)
        debug_api: Enable API debugging (overridden by context if provided)
        session: Optional httpx client session (overridden by context if provided)
        return_raw: Attach the ``httpx.Response`` on ``ResponseGetData.raw_response`` (does not suppress errors).
        is_follow_redirects: Follow HTTP redirects
        timeout: Request timeout in seconds
        parent_class: Optional parent class name for debugging (overridden by context if provided)
        debug_num_stacks_to_drop: Number of stack frames to drop in debug output (overridden by context if provided)
        is_verify: SSL verification flag
        dry_run: If True, return request parameters without executing
        return_on_domo_block: If True, return a ``ResponseGetData`` for VPN/allowlist HTML instead of raising
            ``GetDataError`` (still logged). Transport/httpx errors are always re-raised as-is after logging.

    Returns:
        ResponseGetData object containing the response
    """

    # Build context from individual parameters (manual params override context)
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        debug_num_stacks_to_drop=debug_num_stacks_to_drop,
        parent_class=parent_class,
        is_follow_redirects=is_follow_redirects,
        is_verify=is_verify,
        dry_run=dry_run,
    )

    # Set defaults for params not in RouteContext
    timeout = timeout or DEFAULT_TIMEOUT
    return_raw = return_raw if return_raw is not None else False
    return_on_domo_block = (
        return_on_domo_block if return_on_domo_block is not None else False
    )

    if context.debug_api:
        print(f"[DEBUG] get_data: {method} {url}", flush=True)
        await logger.debug(f"🐛 Debugging get_data: {method} {url}")

    # Create headers and session
    headers = create_headers(
        auth=auth, content_type=content_type, headers=headers or {}
    )

    # Reuse auth's pooled session when connection reuse is enabled. This
    # amortizes TCP/TLS handshakes across many calls (the bulk-operation win).
    # Independent of response caching (context.cache_responses).
    session_from_auth = None
    if (
        context.session is None
        and context.reuse_session
        and auth
        and hasattr(auth, "_cached_session")
        and _is_session_usable(auth._cached_session)
    ):
        session_from_auth = auth._cached_session

    # Create the session (with response-caching transport iff cache_responses).
    session, is_close_session = _create_httpx_session(
        session=session_from_auth or context.session,
        is_verify=context.is_verify,
        cache_responses=context.cache_responses,
        cache_config=context.cache_config,
    )

    # Cache a freshly-created session on auth for reuse. Don't close it — auth
    # owns its lifecycle (see auth.close_session()). Recreate if the previously
    # cached session is closed or bound to a dead/other event loop.
    if (
        context.reuse_session
        and auth
        and hasattr(auth, "_cached_session")
        and not _is_session_usable(auth._cached_session)
        and is_close_session  # Session was just created
    ):
        auth._cached_session = session
        context.session = session  # Update context to reference cached session
        is_close_session = False  # auth owns lifecycle; don't close here
    elif (
        context.reuse_session
        and auth
        and hasattr(auth, "_cached_session")
        and _is_session_usable(auth._cached_session)
        and context.session != auth._cached_session
    ):
        # auth has a usable cached session but context doesn't reference it
        context.session = auth._cached_session

    # Build metadata and additional information
    request_metadata = rgd.RequestMetadata(
        url=url, headers=headers, body=body, params=params
    )

    additional_information = {}
    if context.parent_class:
        additional_information["parent_class"] = context.parent_class

    if context.debug_api:
        message = f"[DEBUG] Request Metadata: {request_metadata.to_dict()}"
        print(message, flush=True)
        await logger.debug(message)

    # Handle dry run mode
    if context.dry_run:
        message = "[DEBUG] Dry run mode enabled. Request not sent."
        print(message, flush=True)
        await logger.debug(message)

        additional_information["dry_run"] = True
        return rgd.ResponseGetData(
            status=200,
            response={
                "dry_run": True,
                "method": method,
                "url": url,
                "headers": headers,
                "body": body,
                "params": params,
                "auth": {
                    "domo_instance": auth.domo_instance if auth else None,
                    "auth_type": type(auth).__name__ if auth else None,
                },
            },
            is_success=True,
            request_metadata=request_metadata,
            additional_information=additional_information,
        )

    try:
        # Build and execute request
        request_kwargs = {
            "method": method,
            "url": url,
            "headers": headers,
            "params": params,
            "follow_redirects": context.is_follow_redirects,
            "timeout": timeout,
        }

        if isinstance(body, dict | list):
            request_kwargs["json"] = body
        elif isinstance(body, str):
            request_kwargs["content"] = body

        try:
            response = await session.request(**request_kwargs)
        except httpx.HTTPError as e:
            await _log_transport_failure(method, request_metadata, e)
            raise TransportError(
                message=f"{method} request transport failed",
                exception=e,
                domo_instance=auth.domo_instance if auth else None,
                additional_context={"url": url, "method": method},
            ) from e

        if context.debug_api:
            print(f"[DEBUG] Response Status: {response.status_code}", flush=True)
            try:
                response_json = response.json()
            except (ValueError, TypeError):
                print(f"[DEBUG] Response Text: {response.text[:500]}", flush=True)
            else:
                pretty_json = json.dumps(response_json, indent=2, ensure_ascii=False)
                print(f"[DEBUG] Response JSON:\n{pretty_json}", flush=True)
            await logger.debug(f"Response Status: {response.status_code}")

        # Domo VPN / allowlist HTML (skip when body is clearly JSON — see _body_may_contain_domo_html_block_page).
        domo_kind: str | None = None
        if _response_may_contain_domo_html_block_page(response):
            domo_kind = _detect_domo_block_kind(response.text, response.status_code)
        if domo_kind:
            result = rgd.ResponseGetData.from_httpx_response(
                res=response,
                request_metadata=request_metadata,
                additional_information=dict(additional_information),
                raw_response=response if return_raw else None,
            )
            _apply_expected_json_mismatch_failure(result, headers, domo_block=domo_kind)
            ip_address = rgd.find_ip(response.text)
            msg = (
                f"Blocked by VPN: {ip_address}"
                if domo_kind == "vpn"
                else f"Blocked by Allowlist: {ip_address}"
            )
            return await _handle_domo_block_response(
                result=result,
                method=method,
                url=url,
                request_metadata=request_metadata,
                message=msg,
                return_on_domo_block=return_on_domo_block,
            )

        if return_raw:
            result = rgd.ResponseGetData.from_httpx_response(
                res=response,
                request_metadata=request_metadata,
                additional_information=additional_information,
                raw_response=response,
            )
        else:
            result = rgd.ResponseGetData.from_httpx_response(
                res=response,
                request_metadata=request_metadata,
                additional_information=additional_information,
            )

        _apply_expected_json_mismatch_failure(result, headers, domo_block=None)

        if not result.is_success:
            await _log_get_data_unsuccessful(method, result, request_metadata)

        return result

    finally:
        if is_close_session:
            await session.aclose()

get_data_stream async

get_data_stream(
    url: str,
    auth: DomoAuth,
    method: str = "GET",
    content_type: str | None = "application/json",
    headers: dict | None = None,
    params: dict | None = None,
    context: RouteContext | None = None,
    debug_api: bool | None = None,
    timeout: int = DEFAULT_STREAM_TIMEOUT,
    parent_class: str | None = None,
    session: AsyncClient | None = None,
    is_verify: bool | None = None,
    is_follow_redirects: bool | None = None,
    return_on_domo_block: bool | None = None,
) -> ResponseGetData

Asynchronously streams data from a Domo API endpoint.

Parameters:

Name Type Description Default
url str

API endpoint URL.

required
auth DomoAuth

Authentication object for Domo APIs.

required
method str

HTTP method to use, default is GET.

'GET'
content_type str | None

Optional content type header.

'application/json'
headers dict | None

Additional HTTP headers.

None
params dict | None

Query parameters for the request.

None
context RouteContext | None

Optional RouteContext with debug/session settings (takes precedence over individual params)

None
debug_api bool | None

Enable debugging information (overridden by context if provided).

None
timeout int

Maximum time to wait for a response (in seconds).

DEFAULT_STREAM_TIMEOUT
parent_class str | None

(Optional) Name of the calling class (overridden by context if provided).

None
session AsyncClient | None

Optional HTTPX session to be used (overridden by context if provided).

None
is_verify bool | None

SSL verification flag.

None
is_follow_redirects bool | None

Follow HTTP redirects if True.

None
return_on_domo_block bool | None

If True, return ResponseGetData for VPN/allowlist HTML instead of raising GetDataError. Transport/httpx errors are logged and re-raised unchanged.

None

Returns:

Type Description
ResponseGetData

An instance of ResponseGetData containing the streamed response data.

Source code in src/crew_dcs/client/get_data.py
694
695
696
697
698
699
700
701
702
703
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
798
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
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
@dmce.run_with_retry()
@log_call(
    action_name="get_data_stream",
    level_name="client",
    log_level="DEBUG",
    config=LogDecoratorConfig(result_processor=ResponseGetDataProcessor()),
)
async def get_data_stream(  # noqa: C901
    url: str,
    auth: dmda.DomoAuth,
    method: str = "GET",
    content_type: str | None = "application/json",
    headers: dict | None = None,
    params: dict | None = None,
    context: RouteContext | None = None,
    debug_api: bool | None = None,
    timeout: int = DEFAULT_STREAM_TIMEOUT,
    parent_class: str | None = None,
    session: httpx.AsyncClient | None = None,
    is_verify: bool | None = None,
    is_follow_redirects: bool | None = None,
    return_on_domo_block: bool | None = None,
) -> rgd.ResponseGetData:
    """Asynchronously streams data from a Domo API endpoint.

    Args:
        url: API endpoint URL.
        auth: Authentication object for Domo APIs.
        method: HTTP method to use, default is GET.
        content_type: Optional content type header.
        headers: Additional HTTP headers.
        params: Query parameters for the request.
        context: Optional RouteContext with debug/session settings (takes precedence over individual params)
        debug_api: Enable debugging information (overridden by context if provided).
        timeout: Maximum time to wait for a response (in seconds).
        parent_class: (Optional) Name of the calling class (overridden by context if provided).
        session: Optional HTTPX session to be used (overridden by context if provided).
        is_verify: SSL verification flag.
        is_follow_redirects: Follow HTTP redirects if True.
        return_on_domo_block: If True, return ``ResponseGetData`` for VPN/allowlist HTML instead of raising
            ``GetDataError``. Transport/httpx errors are logged and re-raised unchanged.

    Returns:
        An instance of ResponseGetData containing the streamed response data.
    """
    # Build context from individual parameters (same pattern as get_data)
    context = RouteContext.build_context(
        context=context,
        session=session,
        debug_api=debug_api,
        parent_class=parent_class,
        is_verify=is_verify,
        is_follow_redirects=is_follow_redirects,
    )

    if context.debug_api:
        message = f"[DEBUG] get_data_stream: {method} {url}"
        print(message, flush=True)
        await logger.debug(message)

    if auth and not auth.token:
        await auth.get_auth_token()

    headers = headers or {}
    headers.update({"Connection": "keep-alive"})
    headers = create_headers(headers=headers, content_type=content_type, auth=auth)

    # Create request metadata
    request_metadata = rgd.RequestMetadata(
        url=url,
        headers=headers,
        body=None,
        params=params,
    )

    # Create additional information with parent_class
    additional_information = {}
    if context.parent_class:
        additional_information["parent_class"] = context.parent_class

    if context.debug_api:
        pprint(
            {
                "method": method,
                "url": url,
                "headers": headers,
                "params": params,
            }
        )

    session, is_close_session = _create_httpx_session(
        session=context.session,
        is_verify=context.is_verify,
    )

    return_on_domo_block = (
        return_on_domo_block if return_on_domo_block is not None else False
    )

    try:
        try:
            async with session.stream(
                method,
                url=url,
                headers=headers,
                follow_redirects=context.is_follow_redirects,
                timeout=timeout,
            ) as res:
                resp_headers = dict(res.headers) if res.headers else {}
                elapsed = None
                try:
                    if hasattr(res, "elapsed") and res.elapsed:
                        elapsed = res.elapsed.total_seconds()
                except RuntimeError:
                    elapsed = None

                if res.status_code != 200:
                    response_text = (
                        res.text if hasattr(res, "text") else str(await res.aread())
                    )
                    domo_kind = None
                    if _body_may_contain_domo_html_block_page(
                        resp_headers.get("content-type") or "", response_text
                    ):
                        domo_kind = _detect_domo_block_kind(
                            response_text, res.status_code
                        )
                    if domo_kind:
                        result = rgd.ResponseGetData(
                            status=res.status_code,
                            response=response_text,
                            is_success=True,
                            request_metadata=request_metadata,
                            additional_information=dict(additional_information),
                            response_headers=resp_headers,
                            elapsed=elapsed,
                        )
                        _apply_expected_json_mismatch_failure(
                            result, headers, domo_block=domo_kind
                        )
                        ip_address = rgd.find_ip(response_text)
                        msg = (
                            f"Blocked by VPN: {ip_address}"
                            if domo_kind == "vpn"
                            else f"Blocked by Allowlist: {ip_address}"
                        )
                        return await _handle_domo_block_response(
                            result=result,
                            method=method,
                            url=url,
                            request_metadata=request_metadata,
                            message=msg,
                            return_on_domo_block=return_on_domo_block,
                        )
                    res_obj = rgd.ResponseGetData(
                        status=res.status_code,
                        response=response_text,
                        is_success=False,
                        request_metadata=request_metadata,
                        additional_information=dict(additional_information),
                        response_headers=resp_headers,
                        elapsed=elapsed,
                    )
                    _apply_expected_json_mismatch_failure(res_obj, headers)
                    await _log_get_data_unsuccessful(method, res_obj, request_metadata)
                    return res_obj

                resp_ct = resp_headers.get("content-type") or ""
                expect_json = _request_expects_json_response(headers)
                if expect_json and not _response_content_type_indicates_json(resp_ct):
                    peek = bytearray()
                    async for chunk in res.aiter_bytes():
                        peek.extend(chunk)
                        if len(peek) >= STREAM_BODY_PEEK_MAX:
                            break
                    text_peek = bytes(peek).decode("utf-8", errors="replace")
                    domo_kind = None
                    if _body_may_contain_domo_html_block_page(resp_ct, text_peek):
                        domo_kind = _detect_domo_block_kind(text_peek, res.status_code)
                    async for _chunk in res.aiter_bytes():
                        pass
                    if domo_kind:
                        result = rgd.ResponseGetData(
                            status=200,
                            response=text_peek,
                            is_success=True,
                            request_metadata=request_metadata,
                            additional_information=dict(additional_information),
                            response_headers=resp_headers,
                            elapsed=elapsed,
                        )
                        _apply_expected_json_mismatch_failure(
                            result, headers, domo_block=domo_kind
                        )
                        ip_address = rgd.find_ip(text_peek)
                        msg = (
                            f"Blocked by VPN: {ip_address}"
                            if domo_kind == "vpn"
                            else f"Blocked by Allowlist: {ip_address}"
                        )
                        return await _handle_domo_block_response(
                            result=result,
                            method=method,
                            url=url,
                            request_metadata=request_metadata,
                            message=msg,
                            return_on_domo_block=return_on_domo_block,
                        )
                    res_obj = rgd.ResponseGetData(
                        status=200,
                        response=text_peek,
                        is_success=False,
                        request_metadata=request_metadata,
                        additional_information={
                            **additional_information,
                            "api_contract": "expected_json_non_json_stream",
                        },
                        response_headers=resp_headers,
                        elapsed=elapsed,
                    )
                    _apply_expected_json_mismatch_failure(res_obj, headers)
                    await _log_get_data_unsuccessful(method, res_obj, request_metadata)
                    return res_obj

                content = bytearray()
                async for chunk in res.aiter_bytes():
                    content += chunk

                res_obj = rgd.ResponseGetData(
                    status=res.status_code,
                    response=bytes(content),
                    is_success=True,
                    request_metadata=request_metadata,
                    additional_information=additional_information,
                    response_headers=resp_headers,
                    elapsed=elapsed,
                )
                _apply_expected_json_mismatch_failure(res_obj, headers)
                if not res_obj.is_success:
                    await _log_get_data_unsuccessful(method, res_obj, request_metadata)
                return res_obj
        except httpx.HTTPError as e:
            await _log_transport_failure(method, request_metadata, e)
            raise TransportError(
                message=f"{method} stream transport failed",
                exception=e,
                domo_instance=auth.domo_instance if auth else None,
                additional_context={"url": url, "method": method},
            ) from e

    finally:
        if is_close_session:
            await session.aclose()

httpx_session_context async

httpx_session_context(
    session: AsyncClient | None = None,
    is_verify: bool = False,
    cache_responses: bool = True,
    cache_config: dict | None = None,
)

Context manager for httpx session lifecycle with optional caching.

Provides automatic session cleanup and supports HTTP response caching via CachedAsyncHTTPTransport. If a session is provided, it's reused and the caller manages its lifecycle. Otherwise, a new session is created and automatically closed when exiting the context.

Parameters:

Name Type Description Default
session AsyncClient | None

Optional existing session to reuse (caller manages lifecycle)

None
is_verify bool

SSL verification flag

False
cache_responses bool

Enable HTTP response caching (default: True)

True
cache_config dict | None

Optional cache configuration dict. Valid keys: - cache_size: Maximum number of cached responses - default_ttl: Default TTL in seconds for cached responses - invalidation_strategy: Strategy for cache invalidation - custom_invalidation_rules: Optional custom rules for CUSTOM strategy - collection_cache_size: Maximum number of cached collections - collection_cache_max_records: Maximum records per collection cache entry - ttl_config: Optional TTL configuration by URL pattern - collection_ttl_config: Optional collection TTL configuration

None

Yields:

Type Description

httpx.AsyncClient session

Example

Basic usage with caching enabled

async with httpx_session_context(cache_responses=True) as session: ... response = await session.get("https://api.example.com") ...

Reuse existing session

async with httpx_session_context(session=existing_session) as session: ... response = await session.get("https://api.example.com") ...

Custom cache configuration

cache_config = {"cache_size": 2000, "default_ttl": 600} async with httpx_session_context(cache_responses=True, cache_config=cache_config) as session: ... response = await session.get("https://api.example.com")

Source code in src/crew_dcs/client/get_data.py
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
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
@asynccontextmanager
async def httpx_session_context(
    session: httpx.AsyncClient | None = None,
    is_verify: bool = False,
    cache_responses: bool = True,
    cache_config: dict | None = None,
):
    """Context manager for httpx session lifecycle with optional caching.

    Provides automatic session cleanup and supports HTTP response caching
    via CachedAsyncHTTPTransport. If a session is provided, it's reused
    and the caller manages its lifecycle. Otherwise, a new session is
    created and automatically closed when exiting the context.

    Args:
        session: Optional existing session to reuse (caller manages lifecycle)
        is_verify: SSL verification flag
        cache_responses: Enable HTTP response caching (default: True)
        cache_config: Optional cache configuration dict. Valid keys:
            - cache_size: Maximum number of cached responses
            - default_ttl: Default TTL in seconds for cached responses
            - invalidation_strategy: Strategy for cache invalidation
            - custom_invalidation_rules: Optional custom rules for CUSTOM strategy
            - collection_cache_size: Maximum number of cached collections
            - collection_cache_max_records: Maximum records per collection cache entry
            - ttl_config: Optional TTL configuration by URL pattern
            - collection_ttl_config: Optional collection TTL configuration

    Yields:
        httpx.AsyncClient session

    Example:
        >>> # Basic usage with caching enabled
        >>> async with httpx_session_context(cache_responses=True) as session:
        ...     response = await session.get("https://api.example.com")
        ...
        >>> # Reuse existing session
        >>> async with httpx_session_context(session=existing_session) as session:
        ...     response = await session.get("https://api.example.com")
        ...
        >>> # Custom cache configuration
        >>> cache_config = {"cache_size": 2000, "default_ttl": 600}
        >>> async with httpx_session_context(cache_responses=True, cache_config=cache_config) as session:
        ...     response = await session.get("https://api.example.com")
    """
    if session is not None:
        # Reuse provided session (caller manages lifecycle)
        yield session
    else:
        # Create new session with optional caching
        if cache_responses:
            from .cached_transport import CachedAsyncHTTPTransport

            # Filter cache_config to only include valid CachedAsyncHTTPTransport parameters
            valid_transport_params = {
                "cache_size",
                "default_ttl",
                "invalidation_strategy",
                "custom_invalidation_rules",
                "collection_cache_size",
                "collection_cache_max_records",
                "ttl_config",
                "collection_ttl_config",
            }

            filtered_config = (
                {k: v for k, v in cache_config.items() if k in valid_transport_params}
                if cache_config
                else {}
            )

            transport = CachedAsyncHTTPTransport(verify=is_verify, **filtered_config)
            session = httpx.AsyncClient(transport=transport)
        else:
            session = httpx.AsyncClient(verify=is_verify)

        try:
            yield session
        finally:
            await session.aclose()

looper async

looper(
    auth: DomoAuth,
    session: AsyncClient | None = None,
    url: str | None = None,
    offset_params: dict | None = None,
    arr_fn: Callable | None = None,
    loop_until_end: bool = False,
    method="POST",
    body: dict | None = None,
    fixed_params: dict | None = None,
    offset_params_in_body: bool = False,
    body_fn=None,
    limit=1000,
    skip=0,
    maximum=0,
    context: RouteContext | None = None,
    debug_api: bool = False,
    debug_loop: bool = False,
    debug_num_stacks_to_drop: int = 1,
    parent_class: str | None = None,
    timeout: int = 10,
    wait_sleep: int = 0,
    is_verify: bool = False,
    return_raw: bool = False,
    return_on_domo_block: bool = False,
    use_collection_cache: bool = True,
    collection_cache_ttl: int | None = None,
    invalidate_collection_cache: bool = False,
) -> ResponseGetData

Iteratively retrieves paginated data from a Domo API endpoint with optional collection caching.

Parameters:

Name Type Description Default
auth DomoAuth

Authentication object for Domo APIs.

required
session AsyncClient | None

HTTPX AsyncClient session used for making requests (overridden by context if provided).

None
url str | None

API endpoint URL for data retrieval.

None
offset_params dict | None

Dictionary specifying the pagination keys (e.g., 'offset', 'limit').

None
arr_fn Callable | None

Function to extract records from the API response.

None
loop_until_end bool

If True, continues fetching until no new records are returned.

False
method

HTTP method to use (default is POST).

'POST'
body dict | None

Request payload (if required).

None
fixed_params dict | None

Fixed query parameters to include in every request.

None
offset_params_in_body bool

Whether to include pagination parameters inside the request body.

False
body_fn

Function to modify the request body before each request.

None
limit

Number of records to retrieve per request.

1000
skip

Initial offset value.

0
maximum

Maximum number of records to retrieve.

0
context RouteContext | None

Optional RouteContext with debug/session settings (takes precedence over individual params)

None
debug_api bool

Enable debugging output for API calls.

False
debug_loop bool

Enable debugging output for the looping process.

False
debug_num_stacks_to_drop int

Number of stack frames to drop in traceback for debugging (overridden by context if provided).

1
parent_class str | None

(Optional) Name of the calling class (overridden by context if provided).

None
timeout int

Request timeout value.

10
wait_sleep int

Time to wait between consecutive requests (in seconds).

0
is_verify bool

SSL verification flag.

False
return_raw bool

Flag to return the raw response instead of processed data.

False
return_on_domo_block bool

Forwarded to get_data — return ResponseGetData on VPN/allowlist HTML instead of raising GetDataError on the first failed page.

False
use_collection_cache bool

Enable collection-level caching (caches complete paginated results).

True
collection_cache_ttl int | None

TTL in seconds for collection cache (None = use transport default).

None
invalidate_collection_cache bool

Force invalidation of collection cache before fetching.

False

Returns:

Type Description
ResponseGetData

An instance of ResponseGetData containing the aggregated data and pagination metadata.

Source code in src/crew_dcs/client/get_data.py
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
@log_call(action_name="looper", level_name="client", log_level="DEBUG")
async def looper(  # noqa: C901
    auth: dmda.DomoAuth,
    session: httpx.AsyncClient | None = None,
    url: str | None = None,
    offset_params: dict | None = None,
    arr_fn: Callable | None = None,
    loop_until_end: bool = False,  # usually you'll set this to true.  it will override maximum
    method="POST",
    body: dict | None = None,
    fixed_params: dict | None = None,
    offset_params_in_body: bool = False,
    body_fn=None,
    limit=1000,
    skip=0,
    maximum=0,
    context: RouteContext | None = None,
    debug_api: bool = False,
    debug_loop: bool = False,
    debug_num_stacks_to_drop: int = 1,
    parent_class: str | None = None,
    timeout: int = 10,
    wait_sleep: int = 0,
    is_verify: bool = False,
    return_raw: bool = False,
    return_on_domo_block: bool = False,
    # NEW: Collection cache parameters
    use_collection_cache: bool = True,
    collection_cache_ttl: int | None = None,
    invalidate_collection_cache: bool = False,
) -> rgd.ResponseGetData:
    """Iteratively retrieves paginated data from a Domo API endpoint with optional collection caching.

    Args:
        auth: Authentication object for Domo APIs.
        session: HTTPX AsyncClient session used for making requests (overridden by context if provided).
        url: API endpoint URL for data retrieval.
        offset_params: Dictionary specifying the pagination keys (e.g., 'offset', 'limit').
        arr_fn: Function to extract records from the API response.
        loop_until_end: If True, continues fetching until no new records are returned.
        method: HTTP method to use (default is POST).
        body: Request payload (if required).
        fixed_params: Fixed query parameters to include in every request.
        offset_params_in_body: Whether to include pagination parameters inside the request body.
        body_fn: Function to modify the request body before each request.
        limit: Number of records to retrieve per request.
        skip: Initial offset value.
        maximum: Maximum number of records to retrieve.
        context: Optional RouteContext with debug/session settings (takes precedence over individual params)
        debug_api: Enable debugging output for API calls.
        debug_loop: Enable debugging output for the looping process.
        debug_num_stacks_to_drop: Number of stack frames to drop in traceback for debugging (overridden by context if provided).
        parent_class: (Optional) Name of the calling class (overridden by context if provided).
        timeout: Request timeout value.
        wait_sleep: Time to wait between consecutive requests (in seconds).
        is_verify: SSL verification flag.
        return_raw: Flag to return the raw response instead of processed data.
        return_on_domo_block: Forwarded to ``get_data`` — return ``ResponseGetData`` on VPN/allowlist HTML
            instead of raising ``GetDataError`` on the first failed page.
        use_collection_cache: Enable collection-level caching (caches complete paginated results).
        collection_cache_ttl: TTL in seconds for collection cache (None = use transport default).
        invalidate_collection_cache: Force invalidation of collection cache before fetching.

    Returns:
        An instance of ResponseGetData containing the aggregated data and pagination metadata.
    """
    # Extract parameters from context if provided (manual params take precedence)
    if isinstance(context, RouteContext):
        session = session or context.session
        debug_num_stacks_to_drop = (
            context.debug_num_stacks_to_drop
            if context.debug_num_stacks_to_drop is not None
            else debug_num_stacks_to_drop
        )
        parent_class = context.parent_class or parent_class
        # Only use context.debug_api if manual param is False (default)
        debug_api = debug_api or (context.debug_api if context.debug_api else False)

    is_close_session = False

    session, is_close_session = _create_httpx_session(session, is_verify=is_verify)

    # Phase 2: Collection Cache Check
    if use_collection_cache and hasattr(session, "_transport"):
        transport = session._transport

        if hasattr(transport, "get_collection_cache"):
            auth_instance = auth.domo_instance

            # Manual invalidation if requested
            if invalidate_collection_cache:  # noqa: SIM102
                if hasattr(transport, "invalidate_collection_by_url"):
                    await transport.invalidate_collection_by_url(url)
                    if debug_loop:
                        print(
                            f"[DEBUG] Collection cache invalidated for {url}",
                            flush=True,
                        )
                        await logger.debug(f"Collection cache invalidated for {url}")

            # Try collection cache
            cached = await transport.get_collection_cache(
                base_url=url,
                params=fixed_params or {},
                auth_instance=auth_instance,
            )

            if cached:
                if debug_loop:
                    print(
                        f"[CACHE HIT] Retrieved {cached.total_records} records from collection cache",
                        flush=True,
                    )
                    await logger.info(
                        f"Collection cache hit: {cached.total_records} records, "
                        f"cached_at={cached.cached_at.isoformat()}, "
                        f"saved {cached.request_count} API calls"
                    )

                # Return cached data
                return rgd.ResponseGetData(
                    status=200,
                    response=cached.data,
                    is_success=True,
                    additional_information={
                        "cache_hit": True,
                        "from_collection_cache": True,
                        "cached_at": cached.cached_at.isoformat(),
                        "total_records": cached.total_records,
                        "original_request_count": cached.request_count,
                    },
                )

    all_rows = []
    is_loop = True
    request_count = 0  # Track number of API calls made

    res: rgd.ResponseGetData | None = None

    if maximum and maximum <= limit and not loop_until_end:
        limit = maximum

    while is_loop:
        params = fixed_params or {}

        if offset_params_in_body:
            if body is None:
                body = {}
            body.update(
                {offset_params.get("offset"): skip, offset_params.get("limit"): limit}
            )

        else:
            params.update(
                {offset_params.get("offset"): skip, offset_params.get("limit"): limit}
            )

        if body_fn:
            body = body_fn(skip, limit, body)

        if debug_loop:
            print(f"\n🚀 Retrieving records {skip} through {skip + limit} via {url}")
            await logger.debug(
                f"\n🚀 Retrieving records {skip} through {skip + limit} via {url}"
            )
            # pprint(params)

            message = {
                "action": "looper_request",
                "params": params,
                "body": body,
                "skip": skip,
                "limit": limit,
            }

        res = await get_data(
            auth=auth,
            url=url,
            method=method,
            params=params,
            body=body,
            timeout=timeout,
            debug_api=debug_api,
            debug_num_stacks_to_drop=debug_num_stacks_to_drop,
            session=session,
            context=context,
            return_raw=return_raw,
            return_on_domo_block=return_on_domo_block,
        )
        request_count += 1  # Track number of API calls

        if not res or not res.is_success:
            if is_close_session:
                await session.aclose()

            fallback_metadata = rgd.RequestMetadata(
                url=url or "",
                headers={},
                body=None,
                params=params or fixed_params or {},
            )
            fallback_info = {}
            if parent_class:
                fallback_info["parent_class"] = parent_class
            return res or rgd.ResponseGetData(
                status=500,
                response="No response",
                is_success=False,
                request_metadata=fallback_metadata,
                additional_information=fallback_info,
            )

        if return_raw:
            return res

        new_records = arr_fn(res)

        all_rows += new_records

        if len(new_records) == 0:
            is_loop = False

        if maximum and len(all_rows) >= maximum and not loop_until_end:
            is_loop = False

        message = f"🐛 Looper iteration complete: {{'all_rows': {len(all_rows)}, 'new_records': {len(new_records)}, 'skip': {skip}, 'limit': {limit}}}"

        if debug_loop:
            print(message, flush=True)
            await logger.debug(message)

        if maximum and skip + limit > maximum and not loop_until_end:
            limit = maximum - len(all_rows)

        skip += len(new_records)
        await asyncio.sleep(wait_sleep)

    if debug_loop:
        message = f"\n🎉 Success - {len(all_rows)} records retrieved from {url} in query looper (made {request_count} API calls)\n"
        print(message, flush=True)
        await logger.info(message)

    # Phase 2: Store Collection Cache
    if (
        use_collection_cache
        and hasattr(session, "_transport")
        and all_rows
        and request_count > 0
    ):
        transport = session._transport

        if hasattr(transport, "store_collection_cache"):
            auth_instance = auth.domo_instance

            await transport.store_collection_cache(
                base_url=url,
                params=fixed_params or {},
                auth_instance=auth_instance,
                data=all_rows,
                request_count=request_count,
                ttl=collection_cache_ttl,
            )

            if debug_loop:
                print(
                    f"[CACHE STORE] Cached {len(all_rows)} records ({request_count} API calls saved for future)",
                    flush=True,
                )
                await logger.info(
                    f"Collection cached: {len(all_rows)} records, {request_count} requests"
                )

    if is_close_session:
        await session.aclose()

    if not res:
        fallback_metadata = rgd.RequestMetadata(
            url=url or "",
            headers={},
            body=None,
            params=fixed_params or {},
        )
        fallback_info = {}
        if parent_class:
            fallback_info["parent_class"] = parent_class
        return rgd.ResponseGetData(
            status=500,
            response="No response received",
            is_success=False,
            request_metadata=fallback_metadata,
            additional_information=fallback_info,
        )

    return rgd.ResponseGetData.from_looper(res=res, array=all_rows)

managed_context async

managed_context(
    context: RouteContext | None = None, **context_kwargs
)

Yield a :class:RouteContext with a session that is CLOSED on block exit.

Opt-in utility for explicitly scoping a short-lived session that must not outlive the block — e.g. a one-off script that shouldn't leave a pooled session cached on auth.

NOTE: this is NOT the default reuse mechanism. Normal calls should just pass auth and rely on RouteContext.reuse_session (auth-level connection pooling, default on), which reuses one session across many calls and is the optimization for bulk operations. Use managed_context only when you want the session torn down deterministically instead of pooled on the auth.

  • If the incoming context already has an open session, it is reused and NOT closed (the caller owns its lifecycle).
  • Otherwise a session is created for the block and closed on exit.

The incoming context is never mutated (build_context copies it).

Source code in src/crew_dcs/client/get_data.py
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
@asynccontextmanager
async def managed_context(
    context: RouteContext | None = None,
    **context_kwargs,
):
    """Yield a :class:`RouteContext` with a session that is CLOSED on block exit.

    Opt-in utility for explicitly scoping a short-lived session that must not
    outlive the block — e.g. a one-off script that shouldn't leave a pooled
    session cached on ``auth``.

    NOTE: this is NOT the default reuse mechanism. Normal calls should just pass
    ``auth`` and rely on ``RouteContext.reuse_session`` (auth-level connection
    pooling, default on), which reuses one session across many calls and is the
    optimization for bulk operations. Use ``managed_context`` only when you want
    the session torn down deterministically instead of pooled on the auth.

    - If the incoming context already has an *open* session, it is reused and NOT
      closed (the caller owns its lifecycle).
    - Otherwise a session is created for the block and closed on exit.

    The incoming context is never mutated (``build_context`` copies it).
    """
    ctx = RouteContext.build_context(context=context, **context_kwargs)

    if ctx.session is not None and not getattr(ctx.session, "is_closed", False):
        yield ctx
        return

    async with httpx_session_context(
        cache_responses=ctx.cache_responses,
        is_verify=ctx.is_verify,
        cache_config=ctx.cache_config,
    ) as session:
        ctx.session = session
        yield ctx

route_function

route_function(
    func: Callable[..., Any],
) -> Callable[..., Any]

Decorator for route functions to ensure they receive a built RouteContext.

This decorator automatically builds a RouteContext from individual parameters or a pre-built context, eliminating the need for route functions to manually call RouteContext.build_context().

⚠️ IMPORTANT: Route functions decorated with @route_function MUST NOT call RouteContext.build_context() inside their function body. The decorator handles context building automatically. Route functions should only accept context: RouteContext | None = None as a parameter and pass it through to get_data/looper.

Parameters:

Name Type Description Default
func Callable[..., Any]

The function to decorate.

required

Returns:

Type Description
Callable[..., Any]

Callable[..., Any]: The decorated function.

The decorated function takes the following arguments

args (Any): Positional arguments for the decorated function. parent_class (str, optional): The parent class. Defaults to None. debug_num_stacks_to_drop (int, optional): The number of stacks to drop for debugging. Defaults to 1. debug_api (bool, optional): Whether to debug the API. Defaults to False. session (httpx.AsyncClient, optional): The HTTPX client session. Defaults to None. log_level (LogLevel | str, optional): The log level. Defaults to LogLevel.INFO. dry_run (bool, optional): If True, return request parameters without executing. Defaults to False. context (RouteContext, optional): A pre-built route context object. Defaults to None. *kwargs (Any): Additional keyword arguments for the decorated function.

Source code in src/crew_dcs/client/get_data.py
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
def route_function(func: Callable[..., Any]) -> Callable[..., Any]:  # noqa: C901
    """
    Decorator for route functions to ensure they receive a built RouteContext.

    This decorator automatically builds a RouteContext from individual parameters
    or a pre-built context, eliminating the need for route functions to manually
    call RouteContext.build_context().

    ⚠️ IMPORTANT: Route functions decorated with @route_function MUST NOT call
    RouteContext.build_context() inside their function body. The decorator handles
    context building automatically. Route functions should only accept
    `context: RouteContext | None = None` as a parameter and pass it through to
    get_data/looper.

    Args:
        func (Callable[..., Any]): The function to decorate.

    Returns:
        Callable[..., Any]: The decorated function.

    The decorated function takes the following arguments:
        *args (Any): Positional arguments for the decorated function.
        parent_class (str, optional): The parent class. Defaults to None.
        debug_num_stacks_to_drop (int, optional): The number of stacks to drop for debugging. Defaults to 1.
        debug_api (bool, optional): Whether to debug the API. Defaults to False.
        session (httpx.AsyncClient, optional): The HTTPX client session. Defaults to None.
        log_level (LogLevel | str, optional): The log level. Defaults to LogLevel.INFO.
        dry_run (bool, optional): If True, return request parameters without executing. Defaults to False.
        context (RouteContext, optional): A pre-built route context object. Defaults to None.
        **kwargs (Any): Additional keyword arguments for the decorated function.
    """
    import inspect

    # Check if the function accepts 'context' parameter
    sig = inspect.signature(func)
    is_accepts_context = "context" in sig.parameters or any(
        param.kind == inspect.Parameter.VAR_KEYWORD for param in sig.parameters.values()
    )

    @wraps(func)
    async def wrapper(  # noqa: C901
        *args: Any,
        parent_class: str | None = None,
        debug_num_stacks_to_drop: int = 1,
        debug_api: bool = False,
        session: httpx.AsyncClient | None = None,
        log_level: str | None = None,
        dry_run: bool = False,
        context: RouteContext | None = None,
        **kwargs: Any,
    ) -> Any:
        # Build context from parameters using RouteContext.build_context()
        # Only pass parameters that are not None to avoid overwriting with defaults
        context_params = {}
        if session is not None:
            context_params["session"] = session
        if debug_api is not False:  # Only pass if explicitly True
            context_params["debug_api"] = debug_api
        if debug_num_stacks_to_drop != 1:  # Only pass if not default
            context_params["debug_num_stacks_to_drop"] = debug_num_stacks_to_drop
        if parent_class is not None:
            context_params["parent_class"] = parent_class
        if log_level is not None:
            context_params["log_level"] = log_level
        if dry_run is not False:  # Only pass if explicitly True
            context_params["dry_run"] = dry_run

        context = RouteContext.build_context(context, **context_params)

        # Build kwargs for the function call
        call_kwargs = {**kwargs}

        # Only pass context if the function accepts it
        if is_accepts_context:
            call_kwargs["context"] = context
        else:
            # Pass individual parameters if function doesn't accept context
            if "debug_api" in sig.parameters:
                call_kwargs["debug_api"] = debug_api
            if "session" in sig.parameters:
                call_kwargs["session"] = session
            if "parent_class" in sig.parameters:
                call_kwargs["parent_class"] = parent_class
            if "debug_num_stacks_to_drop" in sig.parameters:
                call_kwargs["debug_num_stacks_to_drop"] = debug_num_stacks_to_drop
            if "dry_run" in sig.parameters:
                call_kwargs["dry_run"] = dry_run

        result = await func(*args, **call_kwargs)

        if not isinstance(result, rgd.ResponseGetData):
            raise RouteFunctionResponseTypeError(result)

        return result

    return wrapper