Skip to content
Merged
Show file tree
Hide file tree
Changes from 4 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .changes/unreleased/new-items-20260513-105250.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
kind: new-items
body: Supports Digital Twin Builder Flow item
time: 2026-05-13T10:52:50.730594+03:00
custom:
Author: v-alexmoraru
AuthorLink: https://github.com/v-alexmoraru
1 change: 1 addition & 0 deletions docs/essentials/resource_types.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ Item types are the primary content resources within Fabric workspaces. Each type
| `.ApacheAirflowJob` | Apache Airflow job definitions |
| `.CosmosDBDatabase` | Cosmos DB databases |
| `.DigitalTwinBuilder` | Digital twin builder resources |
| `.DigitalTwinBuilderFlow` | Digital Twin Builder flows |
| `.GraphQuerySet` | Graph query collections |
| `.UserDataFunction` | User data functions |

Expand Down
5 changes: 3 additions & 2 deletions docs/examples/item_examples.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@ fab cd ../ws1.Workspace/lh1.Lakehouse
- `.Lakehouse`: `enableSchemas`
- `.Warehouse`: `enableCaseInsensitive`
- `.KQLDatabase`: `dbtype`, `eventhouseId`, `clusterUri`, `databaseName`
- `.DigitalTwinBuilderFlow`: `digitalTwinBuilderId`
- `.MirroredDatabase`: `mirrorType`, `connectionId`, `defaultSchema`, `database`, `mountedTables`
- `.Report`: `semanticModelId`

Expand Down Expand Up @@ -249,7 +250,7 @@ fab set ws1.Workspace/rep1.Report -q definition.parts[0].payload.datasetReferenc
- `.KQLDatabase`, `.KQLDashboard`, `.KQLQueryset`
- `.Eventhouse`, `.Eventstream`
- `.MirroredDatabase`, `.Reflex`
- `.DigitalTwinBuilder`, `.Map`, `.MountedDataFactory`, `.CopyJob`, `.VariableLibrary`
- `.DigitalTwinBuilder`, `.DigitalTwinBuilderFlow`, `.Map`, `.MountedDataFactory`, `.CopyJob`, `.VariableLibrary`


#### Copy Item to Workspace
Expand Down Expand Up @@ -324,7 +325,7 @@ fab export ws1.Workspace/nb1.Notebook -o /tmp
- `.Report`, `.SemanticModel`
- `.KQLDatabase`, `.KQLDashboard`, `.KQLQueryset`
- `.Eventhouse`, `.Eventstream`, `.MirroredDatabase`
- `.Reflex`, `.DigitalTwinBuilder`, `.Map`, `.MountedDataFactory`, `.CopyJob`, `.VariableLibrary`
- `.Reflex`, `.DigitalTwinBuilder`, `.DigitalTwinBuilderFlow`, `.Map`, `.MountedDataFactory`, `.CopyJob`, `.VariableLibrary`


#### Export to Lakehouse
Expand Down
6 changes: 2 additions & 4 deletions src/fabric_cli/client/fab_api_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,9 +41,7 @@ def _get_session() -> requests.Session:
global _shared_session
if _shared_session is None:
_shared_session = requests.Session()
retries = Retry(
total=3, backoff_factor=1, status_forcelist=[502, 503, 504]
)
retries = Retry(total=3, backoff_factor=1, status_forcelist=[502, 503, 504])
adapter = HTTPAdapter(max_retries=retries)
_shared_session.mount("https://", adapter)
_shared_session.headers.update({"Accept-Encoding": "gzip, deflate"})
Expand Down Expand Up @@ -103,7 +101,7 @@ def do_request(
request_params["continuationToken"] = continuation_token

# Build url
url = f"https://{url}/{uri}"
url = f"https://{url}/{uri.lstrip('/')}"
if request_params:
url += f"?{requests.compat.urlencode(request_params)}"

Expand Down
8 changes: 6 additions & 2 deletions src/fabric_cli/core/fab_config/command_support.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,8 @@ commands:
# - kql_database
- mirrored_database
- cosmos_db_database
# - digital_twin_builder
- digital_twin_builder
- digital_twin_builder_flow
Comment thread
v-alexmoraru marked this conversation as resolved.
Outdated
- reflex
# - eventstream
- mounted_data_factory
Expand Down Expand Up @@ -191,7 +192,8 @@ commands:
# - kql_database
- mirrored_database
- cosmos_db_database
# - digital_twin_builder
- digital_twin_builder
- digital_twin_builder_flow
Comment thread
v-alexmoraru marked this conversation as resolved.
Outdated
- reflex
# - eventstream
- mounted_data_factory
Expand Down Expand Up @@ -256,6 +258,7 @@ commands:
- sql_database
- cosmos_db_database
- digital_twin_builder
- digital_twin_builder_flow
- user_data_function
- graph_query_set
- map
Expand Down Expand Up @@ -283,6 +286,7 @@ commands:
- sql_database
- cosmos_db_database
- digital_twin_builder
- digital_twin_builder_flow
- user_data_function
- map
- lakehouse
Expand Down
10 changes: 6 additions & 4 deletions src/fabric_cli/core/fab_types.py
Original file line number Diff line number Diff line change
Expand Up @@ -140,8 +140,7 @@ def from_string(cls, vws_type_str):
class _BaseItemType(Enum):
@classmethod
def from_string(cls, item_type_str):
raise NotImplementedError(
"This method must be implemented in the subclass")
raise NotImplementedError("This method must be implemented in the subclass")


##################################
Expand Down Expand Up @@ -200,8 +199,7 @@ def from_string(cls, vws_type_str):
if item.value.lower() == vws_type_str.lower():
return item
raise FabricCLIError(
ErrorMessages.Common.invalid_virtual_item_container_type(
vws_type_str),
ErrorMessages.Common.invalid_virtual_item_container_type(vws_type_str),
fab_constant.ERROR_INVALID_ITEM_TYPE,
)

Expand Down Expand Up @@ -262,6 +260,7 @@ class ItemType(_BaseItemType):
DATAMART = "Datamart"
DATA_PIPELINE = "DataPipeline"
DIGITAL_TWIN_BUILDER = "DigitalTwinBuilder"
DIGITAL_TWIN_BUILDER_FLOW = "DigitalTwinBuilderFlow"
ENVIRONMENT = "Environment"
EVENTHOUSE = "Eventhouse"
EVENTSTREAM = "Eventstream"
Expand Down Expand Up @@ -502,6 +501,7 @@ class MirroredDatabaseFolders(Enum):
ItemType.DATA_PIPELINE: "dataPipelines",
ItemType.DATAMART: "datamarts",
ItemType.DIGITAL_TWIN_BUILDER: "digitalTwinBuilders",
ItemType.DIGITAL_TWIN_BUILDER_FLOW: "digitalTwinBuilderFlows",
ItemType.ENVIRONMENT: "environments",
ItemType.EVENTHOUSE: "eventhouses",
ItemType.EVENTSTREAM: "eventstreams",
Expand Down Expand Up @@ -550,6 +550,7 @@ class MirroredDatabaseFolders(Enum):
ItemType.DATAMART: "datamarts",
ItemType.DATA_PIPELINE: "pipelines",
ItemType.DIGITAL_TWIN_BUILDER: "digital-twin-builder",
ItemType.DIGITAL_TWIN_BUILDER_FLOW: "digital-twin-builder-flow",
ItemType.ENVIRONMENT: "sparkenvironments",
ItemType.EVENTHOUSE: "eventhouses",
ItemType.EVENTSTREAM: "eventstreams",
Expand Down Expand Up @@ -599,6 +600,7 @@ class MirroredDatabaseFolders(Enum):
},
ItemType.COSMOS_DB_DATABASE: {"default": ""},
ItemType.DIGITAL_TWIN_BUILDER: {"default": ""},
ItemType.DIGITAL_TWIN_BUILDER_FLOW: {"default": ""},
ItemType.USER_DATA_FUNCTION: {"default": ""},
ItemType.GRAPH_QUERY_SET: {"default": ""},
ItemType.VARIABLE_LIBRARY: {"default": ""},
Expand Down
112 changes: 86 additions & 26 deletions src/fabric_cli/utils/fab_cmd_mkdir_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,41 @@ def add_type_specific_payload(item: Item, args, payload):
"parentEventhouseItemId": _eventhouse_id,
}

case ItemType.DIGITAL_TWIN_BUILDER_FLOW:
_digital_twin_builder_id = params.get("digitaltwinbuilderid")
_workspace_id = item.workspace.id

if _digital_twin_builder_id:
payload_dict["creationPayload"] = {
Comment thread
v-alexmoraru marked this conversation as resolved.
Outdated
"digitalTwinBuilderItemReference": {
"referenceType": "ById",
"itemId": _digital_twin_builder_id,
"workspaceId": _workspace_id,
}
}
else:
fab_logger.log_warning(
"DigitalTwinBuilder not provided in params. Creating one first"
)

_initialize_batch_collection_for_dependency_creation(args)

_digital_twin_builder = Item(
f"{item.short_name}_auto",
None,
item.parent,
"DigitalTwinBuilder",
)
_digital_twin_builder_id = mkdir_item.exec(_digital_twin_builder, args)

payload_dict["creationPayload"] = {
"digitalTwinBuilderItemReference": {
"referenceType": "ById",
"itemId": _digital_twin_builder_id,
"workspaceId": _workspace_id,
}
}

case ItemType.MIRRORED_DATABASE:
_type = "genericmirror"
payload_folder = "MirroredDatabase.GenericMirror"
Expand Down Expand Up @@ -297,6 +332,8 @@ def get_params_per_item_type(item: Item):
optional_params = ["enableCaseInsensitive"]
case ItemType.KQL_DATABASE:
optional_params = ["dbType", "eventhouseId", "clusterUri", "databaseName"]
case ItemType.DIGITAL_TWIN_BUILDER_FLOW:
optional_params = ["digitalTwinBuilderId"]
case ItemType.MIRRORED_DATABASE:
optional_params = [
"mirrorType",
Expand Down Expand Up @@ -392,9 +429,11 @@ def _get_params_per_cred_type(cred_type, is_on_premises_gateway):
f"Unsupported credential type {cred_type}. Skipping validation"
)
return []


def _validate_credential_params(cred_type, provided_cred_params, is_on_premises_gateway):

def _validate_credential_params(
cred_type, provided_cred_params, is_on_premises_gateway
):
ignored_params = []
params = {}
param_keys = _get_params_per_cred_type(cred_type, is_on_premises_gateway)
Expand All @@ -418,31 +457,39 @@ def _validate_credential_params(cred_type, provided_cred_params, is_on_premises_
f"Ignoring unsupported parameters for credential type {cred_type}: {ignored_params}"
)
if is_on_premises_gateway:
provided_cred_params["values"] = _validate_and_get_on_premises_gateway_credential_values(provided_cred_params.get("values"))
provided_cred_params["values"] = (
_validate_and_get_on_premises_gateway_credential_values(
provided_cred_params.get("values")
)
)

for key in param_keys:
params[key] = provided_cred_params[key.lower()]

return params


def _validate_and_get_on_premises_gateway_credential_values(cred_values):
for item in cred_values:
if not isinstance(item, dict):
raise FabricCLIError(
ErrorMessages.Common.invalid_onpremises_gateway_values(),
fab_constant.ERROR_INVALID_INPUT,
)

param_values_keys = ["gatewayId", "encryptedCredentials"]
missing_params = [
key for key in param_values_keys
if not all(key.lower() in {k.lower() for k in item.keys()} for item in cred_values)
key
for key in param_values_keys
if not all(
key.lower() in {k.lower() for k in item.keys()} for item in cred_values
)
]
if len(missing_params) > 0:
raise FabricCLIError(
ErrorMessages.Common.missing_onpremises_gateway_parameters(missing_params),
fab_constant.ERROR_INVALID_INPUT,
)
)

ignored_params = [
key
Expand All @@ -453,9 +500,12 @@ def _validate_and_get_on_premises_gateway_credential_values(cred_values):
if len(ignored_params) > 0:
utils_ui.print_warning(
f"Ignoring unsupported parameters for on-premises gateway: {ignored_params}"
)
)

return [{key: item[key.lower()] for key in param_values_keys if key.lower() in item} for item in cred_values]
return [
{key: item[key.lower()] for key in param_values_keys if key.lower() in item}
for item in cred_values
]


def get_connection_config_from_params(payload, con_type, con_type_def, params):
Expand Down Expand Up @@ -506,7 +556,12 @@ def get_connection_config_from_params(payload, con_type, con_type_def, params):
for item in con_type_def["creationMethods"]
if all(
(k.get("name") or "").lower()
in [key.lower() for key in params.get("connectiondetails").get("parameters").keys()]
in [
key.lower()
for key in params.get("connectiondetails")
.get("parameters")
.keys()
]
for k in item["parameters"]
)
),
Expand Down Expand Up @@ -546,7 +601,9 @@ def get_connection_config_from_params(payload, con_type, con_type_def, params):
missing_params = []
if not provided_params:
# Check if the creation method actually requires parameters
required_params = [p["name"] for p in creation_method["parameters"] if p["required"]]
required_params = [
p["name"] for p in creation_method["parameters"] if p["required"]
]
if required_params:
# Get required and optional parameters from the creation method
req_params_str = ", ".join(required_params)
Expand Down Expand Up @@ -590,7 +647,7 @@ def get_connection_config_from_params(payload, con_type, con_type_def, params):
"type": con_type,
"creationMethod": creation_method["name"],
}

# Only add parameters if there are any
if parsed_params:
connection_request["connectionDetails"]["parameters"] = parsed_params
Expand Down Expand Up @@ -656,8 +713,12 @@ def get_connection_config_from_params(payload, con_type, con_type_def, params):
if "skiptestconnection" in provided_cred_params:
provided_cred_params.pop("skiptestconnection")

is_on_premises_gateway = connection_request.get("connectivityType").lower() == "onpremisesgateway"
connection_params = _validate_credential_params(cred_type, provided_cred_params, is_on_premises_gateway)
is_on_premises_gateway = (
connection_request.get("connectivityType").lower() == "onpremisesgateway"
)
connection_params = _validate_credential_params(
cred_type, provided_cred_params, is_on_premises_gateway
)

connection_request["credentialDetails"] = {
"singleSignOnType": singleSignOnType,
Expand All @@ -667,10 +728,11 @@ def get_connection_config_from_params(payload, con_type, con_type_def, params):
}

connection_request["credentialDetails"]["credentials"]["credentialType"] = cred_type

if is_on_premises_gateway:
connection_request["credentialDetails"]["credentials"]["values"] = connection_params.get(
"values")
connection_request["credentialDetails"]["credentials"]["values"] = (
connection_params.get("values")
)

return connection_request

Expand Down Expand Up @@ -767,27 +829,25 @@ def find_mpe_connection(managed_private_endpoint, targetprivatelinkresourceid):

return None


def _initialize_batch_collection_for_dependency_creation(args):
"""Initialize batch collection for scenarios where dependent items need to be created automatically.

This method is used when creating items that have dependencies that don't exist yet, such as:
- Creating a KQL Database without an EventHouse (auto-creates EventHouse first)
- Creating a Report without a Semantic Model (auto-creates Semantic Model first)

The batch collection allows multiple related items to be created in sequence and then
display a consolidated output message showing all items that were created together.

Args:
args (Namespace): The command arguments namespace that will be augmented with
'output_batch' attribute containing 'items' and 'names' lists
to collect creation results.

Note:
This method only initializes the batch collection if it doesn't already exist,
ensuring it's safe to call multiple times during a dependency creation chain.
"""
if not hasattr(args, 'output_batch'):
args.output_batch = {
'items': [],
'names': []
}
if not hasattr(args, "output_batch"):
args.output_batch = {"items": [], "names": []}
Loading
Loading