import snowflake.
connector
from flask import Flask, request, jsonify
# Snowflake account information
sf_conn_args = {
'account': '[Link]-east-1',
'user': '[Link]@[Link]',
'authenticator': 'externalbrowser',
'region': 'us-east-1',
'warehouse': 'PILOT_WH',
'database': 'DB_PILOT_DEFAULT',
'schema': 'DATA_COLLABORATION_SCHEMA'
}
conn = [Link](**sf_conn_args)
cursor = [Link]()
default_role_sql = "USE ROLE PILOT_SB_SYSADMIN"
[Link](default_role_sql)
[Link]("COMMIT")
[Link]()
app = Flask(__name__)
def create_user(data):
# Implement the code to create a role in Snowflake
# Create a Snowflake cursor
role_created_by = data['role_created_by']
user_name = data['user_name']
existing_user = data['existing_user']
default_password = data['default_password']
change_password = data['change_password']
cursor = [Link]()
use_role_sql = f"USE ROLE {role_created_by}"
[Link](use_role_sql)
if existing_user == "No" or existing_user == "NO" or existing_user == "no":
create_user_sql = f"CREATE USER IF NOT EXISTS {user_name}
PASSWORD='{default_password}' MUST_CHANGE_PASSWORD = {change_password}"
[Link](create_user_sql)
print(f"{user_name} USER created.")
elif not user_name.strip() or not existing_user.strip() or
existing_user.lower() == "yes":
pass
[Link]("COMMIT")
[Link]()
return "Success!!"
def create_role(data):
role_to_create = data['role_to_create']
warehouse = data['warehouse']
database = data['database']
schema = data['schema']
object_name = data['object_name']
object_type = data['object_type']
privilege = data['privilege']
shared_database = data['shared_database']
target_db_to_create_shared_db_object =
data['target_db_to_create_shared_db_object']
target_schema_to_create_shared_db_object =
data['target_schema_to_create_shared_db_object']
object_name_of_shared_db_object = data['object_name_of_shared_db_object']
cursor = [Link]()
use_wh_sql = f"USE WAREHOUSE {warehouse}"
[Link](use_wh_sql)
print(use_wh_sql)
[Link]("COMMIT")
use_db_sql = f"USE DATABASE {database}"
[Link](use_db_sql)
print(use_db_sql)
[Link]("COMMIT")
# Implement the code to create a role in Snowflake
create_role_sql = f"CREATE ROLE IF NOT EXISTS {role_to_create}"
[Link](create_role_sql)
[Link]("COMMIT")
print(f"{role_to_create} ROLE created/assigned !")
grant_wh_usage_sql = f"grant usage on warehouse {warehouse} to role
{role_to_create}"
[Link](grant_wh_usage_sql)
print(grant_wh_usage_sql)
[Link]("COMMIT")
if shared_database.upper() == "YES":
create_shared_db_object_sql = f"create {object_type} if not exists
{target_db_to_create_shared_db_object}.{target_schema_to_create_shared_db_object}.
{object_name_of_shared_db_object} as select * from {database}.{schema}.
{object_name}"
[Link](create_shared_db_object_sql)
print(create_shared_db_object_sql)
[Link]("COMMIT")
grant_db_usage_sql = f"grant usage on database
{target_db_to_create_shared_db_object} to role {role_to_create}"
[Link](grant_db_usage_sql)
print(grant_db_usage_sql)
[Link]("COMMIT")
grant_schema_usage_sql = f"grant usage on schema
{target_db_to_create_shared_db_object}.{target_schema_to_create_shared_db_object}
to role {role_to_create}"
[Link](grant_schema_usage_sql)
print(grant_schema_usage_sql)
[Link]("COMMIT")
grant_object_privilege_sql = f"grant {privilege} on {object_type}
{target_db_to_create_shared_db_object}.{target_schema_to_create_shared_db_object}.
{object_name_of_shared_db_object} to role {role_to_create}"
print(grant_object_privilege_sql)
[Link](grant_object_privilege_sql)
[Link]("COMMIT")
else:
grant_db_usage_sql = f"grant usage on database {database} to role
{role_to_create}"
[Link](grant_db_usage_sql)
print(grant_db_usage_sql)
[Link]("COMMIT")
grant_schema_usage_sql = f"grant usage on schema {database}.{schema} to
role {role_to_create}"
[Link](grant_schema_usage_sql)
print(grant_schema_usage_sql)
[Link]("COMMIT")
grant_object_privilege_sql = f"grant {privilege} on {object_type}
{database}.{schema}.{object_name} to role {role_to_create}"
print(grant_object_privilege_sql)
[Link](grant_object_privilege_sql)
[Link]("COMMIT")
[Link]()
return "Success!!"
def assign_role_to_user(data):
# implement the code to assign the role to the user
# Create a Snowflake cursor
role_to_create = data['role_to_create']
user_name = data['user_name']
grant_role_to_user = data['grant_role_to_user']
cursor = [Link]()
if grant_role_to_user.lower() == "yes":
assign_role_sql = f"GRANT ROLE {role_to_create} TO USER {user_name}"
print(assign_role_sql)
[Link](assign_role_sql)
[Link]("COMMIT")
print(f"Role {role_to_create} assigned to user {user_name}")
elif not user_name.strip() or grant_role_to_user.lower() == "no" or not
grant_role_to_user.strip():
pass
[Link]()
return "Success!!"
def assign_role_to_role(data):
# implement the code to assign the role to the user
# Create a Snowflake cursor
role_to_create = data['role_to_create']
role_to_assign = data['role_to_assign']
grant_role_to_role = data['grant_role_to_role']
cursor = [Link]()
if grant_role_to_role.lower() == "yes":
assign_role_sql = f"GRANT ROLE {role_to_create} TO ROLE {role_to_assign}"
[Link](assign_role_sql)
[Link]("COMMIT")
print(f"Role assigned to role {role_to_assign}")
elif not role_to_assign.strip() or grant_role_to_role.lower() == "no" or not
grant_role_to_role.strip():
pass
[Link]()
return "Success!!"
def create_dynamic_data_masking_object(data):
# Implement the code to create a dynamic data masking object
policy_created_by = data['policy_created_by']
table_database = data['table_database']
table_schema = data['table_schema']
table_name = data['table_name']
table_column = data['table_column']
policy_name = data['policy_name']
policy_database = data['policy_database']
policy_schema = data['policy_schema']
assigned_to_role = data['assigned_to_role']
policy_action_flag = data['policy_action_flag']
cursor = [Link]()
use_role_sql = f"USE ROLE {policy_created_by}"
[Link](use_role_sql)
use_database_sql = f"USE DATABASE {policy_database}"
[Link](use_database_sql)
use_schema_sql = f"USE SCHEMA {policy_schema}"
[Link](use_schema_sql)
if policy_action_flag.upper() == "SET":
# SQL query to create the masking policy
create_policy_sql = f"""
CREATE MASKING POLICY IF NOT EXISTS {policy_database}.
{policy_schema}.{policy_name} AS (val string) RETURNS string -> CASE
WHEN CURRENT_ROLE() IN ('{assigned_to_role}') THEN val
ELSE '******'
END
"""
# Execute the SQL query
[Link](create_policy_sql)
set_policy_sql = f""" ALTER TABLE IF EXISTS {table_database}.
{table_schema}.{table_name} MODIFY COLUMN {table_column} SET MASKING POLICY
{policy_name} """
[Link](set_policy_sql)
print(f"Policy {policy_name} added to column {table_column} of table
{table_name}")
elif policy_action_flag.upper() == "UNSET":
unset_policy_sql = f""" ALTER TABLE IF EXISTS {table_database}.
{table_schema}.{table_name} MODIFY COLUMN {table_column} UNSET MASKING POLICY """
[Link](unset_policy_sql)
print(f"Policy {policy_name} removed from column {table_column} of table
{table_name}")
# Commit the changes
[Link]("COMMIT")
# Close the cursor and connection
[Link]()
def create_row_access_policy(data):
# Implement the code to create a row access object
# Extract values from the jinja template input
policy_created_by = data['policy_created_by']
policy_action_flag = data['policy_action_flag']
policy_database_name = data['policy_database_name']
policy_schema_name = data['policy_schema_name']
policy_name = data['policy_name']
mapping_table_database_name = data['mapping_table_database_name']
mapping_table_schema_name = data['mapping_table_schema_name']
mapping_table_name = data['mapping_table_name']
mapping_table_parameter_field_name = data['mapping_table_parameter_field_name']
mapping_table_parameter_field_datatype =
data['mapping_table_parameter_field_datatype']
mapping_table_row_access_condn_col_name =
data['mapping_table_row_access_condn_col_name']
mapping_table_target_table_common_column =
data['mapping_table_target_table_common_column']
policy_target_table_database_name = data['policy_target_table_database_name']
policy_target_table_schema_name = data['policy_target_table_schema_name']
policy_target_table_name = data['policy_target_table_name']
policy_target_table_column = data['policy_target_table_column']
default_role_to_access_all_rows = data['default_role_to_access_all_rows']
cursor = [Link]()
use_role_sql = f"USE ROLE {policy_created_by}"
[Link](use_role_sql)
if policy_action_flag.upper() == "ADD":
# Construct the SQL statement with conditional formatting and default
values
policy_sql_statement = f"""
create OR REPLACE row access policy {policy_database_name}.
{policy_schema_name}.{policy_name} as ({mapping_table_parameter_field_name}
{mapping_table_parameter_field_datatype}) returns boolean ->
'{default_role_to_access_all_rows}' = current_role()
or exists (
select 1 from {mapping_table_database_name}.
{mapping_table_schema_name}.{mapping_table_name}
where {mapping_table_row_access_condn_col_name} =
current_role()
and {mapping_table_target_table_common_column} =
{mapping_table_parameter_field_name}
);
"""
# Print the SQL statement
print(policy_sql_statement)
[Link](policy_sql_statement)
apply_policy_sql_statement = f"""
alter table {policy_target_table_database_name}.
{policy_target_table_schema_name}.{policy_target_table_name}
add row access policy {policy_database_name}.{policy_schema_name}.
{policy_name} on ({policy_target_table_column});
"""
# Print the SQL statement
print(apply_policy_sql_statement)
[Link](apply_policy_sql_statement)
elif policy_action_flag.upper() == "REMOVE":
remove_policy_sql_statement = f"""
alter table {policy_target_table_database_name}.
{policy_target_table_schema_name}.{policy_target_table_name}
drop row access policy {policy_database_name}.{policy_schema_name}.
{policy_name};
"""
# Print the SQL statement
print(remove_policy_sql_statement)
[Link](remove_policy_sql_statement)
# Commit the changes
[Link]("COMMIT")
# Close the cursor and connection
[Link]()
def create_share(data):
# Implement the code to create a data share
share_name = data['share_name']
role_created_by = data['role_created_by']
share_description = data['share_description']
action_flag = data['ACTION_FLAG(CREATE/DROP)']
cursor = [Link]()
use_role_sql = f"USE ROLE {role_created_by}"
[Link](use_role_sql)
# SQL query to create the data share
if action_flag.upper() == "CREATE":
create_share_sql = f"""
CREATE SHARE IF NOT EXISTS {share_name}
COMMENT = '{share_description}'
"""
# Execute the SQL query
[Link](create_share_sql)
print(create_share_sql)
elif action_flag.upper() == "DROP":
drop_share_sql = f"DROP SHARE {share_name}"
[Link](drop_share_sql)
print(drop_share_sql)
# Close the cursor and connection
[Link]()
def update_datashare_objects(data):
# Implement the code to add data objects to a data share
action_flag = data['ACTION_FLAG(ADD/REMOVE)']
role_created_by = data['role_created_by']
share_name = data['share_name']
database_name = data['database_name']
schema_name = data['schema_name']
object_name = data['object_name']
object_type = data['object_type']
dbs = data['reference_dbs_incaseof_secure_views']
cursor = [Link]()
use_role_sql = f"USE ROLE {role_created_by}"
[Link](use_role_sql)
print(use_role_sql)
cursor = [Link]()
use_db_sql = f"USE DATABASE {database_name}"
[Link](use_db_sql)
print(use_db_sql)
if action_flag == "ADD":
# Grant USAGE on the database to the share
grant_usage_db_sql = f"GRANT USAGE ON DATABASE {database_name} TO SHARE
{share_name}"
[Link](grant_usage_db_sql)
print(grant_usage_db_sql)
# Grant USAGE on the database to the share
grant_usage_schema_sql = f"GRANT USAGE ON SCHEMA {schema_name} TO SHARE
{share_name}"
[Link](grant_usage_schema_sql)
print(grant_usage_schema_sql)
# Grant SELECT on each data object to the share
if object_type == "SECURE_VIEW":
# Get the values from the 'ReferenceDBs' column as a list
reference_dbs_list = [Link](',').explode().[Link]().tolist()
# Loop through the list of values using a for loop
for dbs_list in reference_dbs_list:
grant_ref_usage_sql = f"grant reference_usage on database
{dbs_list} to share {share_name}"
[Link](grant_ref_usage_sql)
print(grant_ref_usage_sql)
grant_select_sql = f"GRANT SELECT ON {database_name}.{schema_name}.
{object_name} TO SHARE {share_name}"
[Link](grant_select_sql)
print(grant_select_sql)
else:
grant_select_sql = f"GRANT SELECT ON {database_name}.{schema_name}.
{object_name} TO SHARE {share_name}"
[Link](grant_select_sql)
print(grant_select_sql)
elif action_flag == "REMOVE":
revoke_select_sql = f"REVOKE SELECT ON {database_name}.{schema_name}.
{object_name} FROM SHARE {share_name}"
[Link](revoke_select_sql)
# Commit the changes
[Link]("COMMIT")
# Close the cursor and connection
[Link]()
def create_external_tables(data):
# Implement the code to create external stages
role_created_by = data['role_created_by']
table_name = data['table_name']
database_name = data['database_name']
schema_name = data['schema_name']
file_format = data['file_format']
location = data['location']
columns = data['columns']
data_types = data['data_types']
cursor = [Link]()
use_role_sql = f"USE ROLE {role_created_by}"
[Link](use_role_sql)
# Split the comma-separated columns and data types into lists
column_list = [[Link]() for col in [Link](',')]
data_type_list = [[Link]() for dtype in data_types.split(',')]
# Create the external table SQL query
create_table_sql = f"""
CREATE OR REPLACE EXTERNAL TABLE IF NOT EXISTS {database_name}.
{schema_name}.{table_name}
({", ".join([f"{col} {dtype}" for col, dtype in zip(column_list,
data_type_list)])})
LOCATION = '{location}'
FILE_FORMAT = {file_format}
AUTO_REFRESH = FALSE
"""
[Link](create_table_sql)
print(create_table_sql)
print(f"External table '{table_name}' created by '{role_created_by}'")
# Commit the changes
[Link]("COMMIT")
[Link]()
def create_external_stages(data):
# Implement the code to create external stages
role_created_by = data['role_created_by']
stage_name = data['stage_name']
database_name = data['database_name']
schema_name = data['schema_name']
storage_integration_name = data['storage_integration_name']
cursor = [Link]()
use_role_sql = f"USE ROLE {role_created_by}"
[Link](use_role_sql)
create_stage_sql = f"""
CREATE OR REPLACE STAGE IF NOT EXISTS {database_name}.{schema_name}.
{stage_name}
STORAGE_INTEGRATION = {storage_integration_name}
"""
[Link](create_stage_sql)
print(f"External Stage '{stage_name}' created in '{database_name}.
{schema_name}'")
# Commit the changes
[Link]("COMMIT")
[Link]()
@[Link]('/process_request_create_role', methods=['POST'])
def process_request_create_role():
try:
data = request.get_json()
create_user(data)
create_role(data)
assign_role_to_user(data)
assign_role_to_role(data)
return jsonify({"result": "Success!!"})
except Exception as e:
return jsonify({"status": "error", "message": str(e)})
@[Link]('/process_request_create_dynamic_data_masking', methods=['POST'])
def process_request_create_dynamic_data_masking():
try:
data = request.get_json()
create_dynamic_data_masking_object(data)
return jsonify({"status": "success"})
except Exception as e:
return jsonify({"status": "error", "message": str(e)})
@[Link]('/process_request_create_row_access_policy', methods=['POST'])
def process_request_create_row_access_policy():
try:
data = request.get_json()
create_row_access_policy(data)
return jsonify({"status": "success"})
except Exception as e:
return jsonify({"status": "error", "message": str(e)})
@[Link]('/process_request_create_share', methods=['POST'])
def process_request_create_share():
try:
data = request.get_json()
create_share(data)
return jsonify({"status": "success"})
except Exception as e:
return jsonify({"status": "error", "message": str(e)})
@[Link]('/process_request_update_datashare_objects', methods=['POST'])
def process_request_update_datashare_objects():
try:
data = request.get_json()
update_datashare_objects(data)
return jsonify({"status": "success"})
except Exception as e:
return jsonify({"status": "error", "message": str(e)})
@[Link]('/process_request_create_external_table', methods=['POST'])
def process_request_create_external_table():
try:
data = request.get_json()
create_external_tables(data)
return jsonify({"status": "success"})
except Exception as e:
return jsonify({"status": "error", "message": str(e)})
@[Link]('/process_request_create_external_stage', methods=['POST'])
def process_request_create_external_stage():
try:
data = request.get_json()
create_external_stages(data)
return jsonify({"status": "success"})
except Exception as e:
return jsonify({"status": "error", "message": str(e)})
if __name__ == "__main__":
[Link](debug=True)