import sys
from common.utils import *

@flow(
    task_runner=DaskTaskRunner(cluster_kwargs={"n_workers": 2, "processes": True}),
    persist_result=True, 
    result_storage=s3_prefect_result_block,
)
async def main_flow():
    try:
        logger = get_run_logger()
        test_task = databricks_run_now.with_options(name="prefect-test").submit()
        test_task2 = databricks_run_now.with_options(name="prefect-test2").submit()

        # all LAST tasks with separate dependency, to identify flow complete
        futures = [test_task, test_task2]
        results = await asyncio.gather(*[future.result() for future in futures])

    except Exception as e:
        logger.error(f"Flow run failed, {e}")

if __name__ == "__main__":
    local_file_name = "template_flow.py" # current file name
    current_folder_name = "template" # current folder path

    s3_key = f"{deploy_s3_folder}{current_folder_name}/{local_file_name}" # s3 deployment path    
    s3_prefect_bucket_block.upload_from_path("common/utils.py", f"{deploy_s3_folder}common/utils.py")

    # upload flow code to s3, for prefect worker to pull for execution
    s3_prefect_bucket_block.upload_from_path(f"{current_folder_name}/{local_file_name}", s3_key)

    # Specific s3 as the source of deployment code to pull from
    flow.from_source(
        source=f"s3://{deploy_s3_bucket}/{deploy_s3_folder}",
        entrypoint=f"{current_folder_name}/{local_file_name}:main_flow"
    ).deploy(
        name=f"{flow_name}-{deployment_suffix}",
        work_pool_name="default-worker-pool",
        job_variables={
            "env": {
                "EXTRA_PIP_PACKAGES": "prefect-aws"
            }
        }
    )