diff --git a/poetry.lock b/poetry.lock index 3d2615d26..889ef581a 100644 --- a/poetry.lock +++ b/poetry.lock @@ -1800,6 +1800,18 @@ requests-oauthlib = ">=0.5.0" [package.extras] async = ["aiodns", "aiohttp (>=3.0)"] +[[package]] +name = "networkx" +version = "2.6.3" +description = "Python package for creating and manipulating graphs and networks" +category = "main" +optional = true +python-versions = ">=3.7" +files = [ + {file = "networkx-2.6.3-py3-none-any.whl", hash = "sha256:80b6b89c77d1dfb64a4c7854981b60aeea6360ac02c6d4e4913319e0a313abef"}, + {file = "networkx-2.6.3.tar.gz", hash = "sha256:c0946ed31d71f1b732b5aaa6da5a0388a345019af232ce2f49c766e2d6795c51"}, +] + [[package]] name = "mypy-extensions" version = "1.0.0" @@ -2903,7 +2915,7 @@ azure = ["azure-identity", "azure-mgmt-authorization", "azure-mgmt-compute", "az gateway = ["flask", "lz4", "pynacl", "pyopenssl", "werkzeug"] gcp = ["google-api-python-client", "google-auth", "google-cloud-compute", "google-cloud-storage"] ibm = ["ibm-cloud-sdk-core", "ibm-cos-sdk", "ibm-vpc"] -solver = ["cvxpy", "graphviz", "matplotlib", "numpy"] +solver = ["networkx", "cvxpy", "graphviz", "matplotlib", "numpy"] [metadata] lock-version = "2.0" diff --git a/pyproject.toml b/pyproject.toml index d59fe74a6..a582eed68 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -54,6 +54,7 @@ cvxpy = { version = ">=1.1.0", optional = true } graphviz = { version = ">=0.15", optional = true } matplotlib = { version = ">=3.0.0", optional = true } numpy = { version = ">=1.19.0", optional = true } +networkx = { version = ">=2.5", optional = true } # gateway dependencies flask = { version = "^2.1.2", optional = true } @@ -71,7 +72,7 @@ ibm = ["ibm-cloud-sdk-core", "ibm-cos-sdk", "ibm-vpc"] scp = ["boto3"] all = ["boto3", "azure-identity", "azure-mgmt-authorization", "azure-mgmt-compute", "azure-mgmt-network", "azure-mgmt-resource", "azure-mgmt-storage", "azure-mgmt-subscription", "azure-storage-blob", "google-api-python-client", "google-auth", "google-cloud-compute", "google-cloud-storage", "ibm-cloud-sdk-core", "ibm-cos-sdk", "ibm-vpc"] gateway = ["flask", "lz4", "pynacl", "pyopenssl", "werkzeug"] -solver = ["cvxpy", "graphviz", "matplotlib", "numpy"] +solver = ["networkx", "cvxpy", "graphviz", "matplotlib", "numpy"] [tool.poetry.dev-dependencies] pytest = ">=6.0.0" diff --git a/scripts/pack_docker.sh b/scripts/pack_docker.sh index 49e940298..c2d92c620 100755 --- a/scripts/pack_docker.sh +++ b/scripts/pack_docker.sh @@ -12,15 +12,15 @@ set +e >&2 echo -e "${BGreen}Building docker image${NC}" set -e ->&2 sudo DOCKER_BUILDKIT=1 docker build -t skyplane --platform linux/x86_64 . +>&2 DOCKER_BUILDKIT=1 docker build -t skyplane --platform linux/x86_64 . set +e DOCKER_URL="ghcr.io/$1/skyplane:local-$(openssl rand -hex 16)" >&2 echo -e "${BGreen}Uploading docker image to $DOCKER_URL${NC}" set -e ->&2 sudo docker tag skyplane $DOCKER_URL ->&2 sudo docker push $DOCKER_URL ->&2 sudo docker system prune -f +>&2 docker tag skyplane $DOCKER_URL +>&2 docker push $DOCKER_URL +>&2 docker system prune -f set +e >&2 echo -e "${BGreen}SKYPLANE_DOCKER_IMAGE=$DOCKER_URL${NC}" diff --git a/scripts/requirements-gateway.txt b/scripts/requirements-gateway.txt index 687f3aeeb..be1e298a1 100644 --- a/scripts/requirements-gateway.txt +++ b/scripts/requirements-gateway.txt @@ -33,3 +33,4 @@ numpy pandas pyarrow typer +networkx \ No newline at end of file diff --git a/skyplane/api/dataplane.py b/skyplane/api/dataplane.py index 69916093d..391f9a758 100644 --- a/skyplane/api/dataplane.py +++ b/skyplane/api/dataplane.py @@ -8,6 +8,7 @@ import nacl.secret import nacl.utils +import typer import urllib3 from pathlib import Path from typing import TYPE_CHECKING, Dict, List, Optional @@ -111,6 +112,7 @@ def _start_gateway( # write gateway programs gateway_program_filename = Path(f"{gateway_log_dir}/gateway_program_{gateway_node.gateway_id}.json") + print(gateway_program_filename) with open(gateway_program_filename, "w") as f: f.write(gateway_node.gateway_program.to_json()) @@ -232,9 +234,11 @@ def provision( def copy_gateway_logs(self): # copy logs from all gateways in parallel def copy_log(instance): - out_file = self.transfer_dir / f"gateway_{instance.uuid()}.stdout" - err_file = self.transfer_dir / f"gateway_{instance.uuid()}.stderr" - logger.fs.info(f"[Dataplane.copy_gateway_logs] Copying logs from {instance.uuid()}: {out_file}") + out_file = f"{self.transfer_dir}/gateway_{instance.uuid()}.stdout" + err_file = f"{self.transfer_dir}/gateway_{instance.uuid()}.stderr" + typer.secho(f"Downloading log: {self.transfer_dir}/gateway_{instance.uuid()}.stdout", fg="bright_black") + typer.secho(f"Downloading log: {self.transfer_dir}/gateway_{instance.uuid()}.stderr", fg="bright_black") + instance.run_command("sudo docker logs -t skyplane_gateway 2> /tmp/gateway.stderr > /tmp/gateway.stdout") instance.download_file("/tmp/gateway.stdout", out_file) instance.download_file("/tmp/gateway.stderr", err_file) diff --git a/skyplane/api/pipeline.py b/skyplane/api/pipeline.py index 6fa3face4..3ffae2b1c 100644 --- a/skyplane/api/pipeline.py +++ b/skyplane/api/pipeline.py @@ -10,7 +10,13 @@ from skyplane.api.transfer_job import CopyJob, SyncJob, TransferJob from skyplane.api.config import TransferConfig -from skyplane.planner.planner import MulticastDirectPlanner, DirectPlannerSourceOneSided, DirectPlannerDestOneSided +from skyplane.planner.planner import ( + MulticastDirectPlanner, + UnicastDirectPlanner, + UnicastILPPlanner, + MulticastILPPlanner, + MulticastMDSTPlanner, +) from skyplane.planner.topology import TopologyPlanGateway from skyplane.utils import logger from skyplane.utils.definitions import tmp_log_dir @@ -31,7 +37,7 @@ def __init__( transfer_config: TransferConfig, # cloud_regions: dict, max_instances: Optional[int] = 1, - n_connections: Optional[int] = 64, + num_connections: Optional[int] = 32, planning_algorithm: Optional[str] = "direct", debug: Optional[bool] = False, ): @@ -47,7 +53,7 @@ def __init__( # self.cloud_regions = cloud_regions # TODO: set max instances with VM CPU limits and/or config self.max_instances = max_instances - self.n_connections = n_connections + self.n_connections = num_connections self.provisioner = provisioner self.transfer_config = transfer_config self.http_pool = urllib3.PoolManager(retries=urllib3.Retry(total=3)) @@ -61,12 +67,24 @@ def __init__( # planner self.planning_algorithm = planning_algorithm + if self.planning_algorithm == "direct": - self.planner = MulticastDirectPlanner(self.max_instances, self.n_connections, self.transfer_config) - elif self.planning_algorithm == "src_one_sided": - self.planner = DirectPlannerSourceOneSided(self.max_instances, self.n_connections, self.transfer_config) - elif self.planning_algorithm == "dst_one_sided": - self.planner = DirectPlannerDestOneSided(self.max_instances, self.n_connections, self.transfer_config) + # TODO: should find some ways to merge direct / Ndirect + #self.planner = UnicastDirectPlanner(self.max_instances, num_connections) + #self.planner = MulticastDirectPlanner(self.max_instances, self.n_connections, self.transfer_config) + self.planner = MulticastDirectPlanner(self.max_instances, self.n_connections) + #elif self.planning_algorithm == "Ndirect": + # self.planner = MulticastDirectPlanner(self.max_instances, num_connections) + elif self.planning_algorithm == "MDST": + self.planner = MulticastMDSTPlanner(self.max_instances, num_connections) + elif self.planning_algorithm == "ILP": + self.planner = MulticastILPPlanner(self.max_instances, num_connections) + elif self.planning_algorithm == "UnicastILP": + self.planner = UnicastILPPlanner(self.max_instances, num_connections) + #elif self.planning_algorithm == "src_one_sided": + # self.planner = DirectPlannerSourceOneSided(self.max_instances, self.n_connections, self.transfer_config) + #elif self.planning_algorithm == "dst_one_sided": + # self.planner = DirectPlannerDestOneSided(self.max_instances, self.n_connections, self.transfer_config) else: raise ValueError(f"No such planning algorithm {planning_algorithm}") diff --git a/skyplane/api/transfer_job.py b/skyplane/api/transfer_job.py index 7bf0f2509..92a48abad 100644 --- a/skyplane/api/transfer_job.py +++ b/skyplane/api/transfer_job.py @@ -25,6 +25,8 @@ from skyplane import exceptions from skyplane.api.config import TransferConfig from skyplane.chunk import Chunk +from skyplane.obj_store.azure_blob_interface import AzureBlobObject +from skyplane.obj_store.gcs_interface import GCSObject from skyplane.obj_store.storage_interface import StorageInterface from skyplane.obj_store.object_store_interface import ObjectStoreObject, ObjectStoreInterface from skyplane.utils import logger @@ -102,6 +104,7 @@ def _run_multipart_chunk_thread( src_object = transfer_pair.src_obj dest_objects = transfer_pair.dst_objs dest_key = transfer_pair.dst_key + print("dest_key: ", dest_key) if isinstance(self.src_iface, ObjectStoreInterface): mime_type = self.src_iface.get_obj_mime_type(src_object.key) # create multipart upload request per destination @@ -130,6 +133,7 @@ def _run_multipart_chunk_thread( for _ in range(num_chunks): file_size_bytes = min(chunk_size_bytes, src_object.size - offset) assert file_size_bytes > 0, f"file size <= 0 {file_size_bytes}" + print("partition", part_num, self.num_partitions) chunk = Chunk( src_key=src_object.key, dest_key=dest_key, # dest_object.key, # TODO: upload basename (no prefix) @@ -283,10 +287,10 @@ def transfer_pair_generator( dest_provider, dest_region = dst_iface.region_tag().split(":") try: dest_key = self.map_object_key_prefix(src_prefix, obj.key, dst_prefix, recursive=recursive) - assert ( - dest_key[: len(dst_prefix)] == dst_prefix - ), f"Destination key {dest_key} does not start with destination prefix {dst_prefix}" - dest_keys.append(dest_key[len(dst_prefix) :]) + # TODO: why is it changed here? + # dest_keys.append(dest_key[len(dst_prefix) :]) + + dest_keys.append(dest_key) except exceptions.MissingObjectException as e: logger.fs.exception(e) raise e from None @@ -347,6 +351,7 @@ def chunk(self, transfer_pair_generator: Generator[TransferPair, None, None]) -> multipart_chunk_threads.append(t) # begin chunking loop + part_num = 0 for transfer_pair in transfer_pair_generator: # print("transfer_pair", transfer_pair.src_obj.key, transfer_pair.dst_objs) src_obj = transfer_pair.src_obj @@ -362,9 +367,11 @@ def chunk(self, transfer_pair_generator: Generator[TransferPair, None, None]) -> dest_key=transfer_pair.dst_key, # TODO: get rid of dest_key, and have write object have info on prefix (or have a map here) chunk_id=uuid.uuid4().hex, chunk_length_bytes=transfer_pair.src_obj.size, - partition_id=str(0), # TODO: fix this to distribute across multiple partitions + #partition_id=str(0), # TODO: fix this to distribute across multiple partitions + partition_id=str(part_num % self.num_partitions), ) ) + part_num += 1 if self.transfer_config.multipart_enabled: # drain multipart chunk queue and yield with updated chunk IDs @@ -512,8 +519,12 @@ def dst_prefixes(self) -> List[str]: if not hasattr(self, "_dst_prefix"): if self.transfer_type == "unicast": self._dst_prefix = [str(parse_path(self.dst_paths[0])[2])] + print("return dst_prefixes for unicast", self._dst_prefix) else: + for path in self.dst_paths: + print("Parsing result for multicast", parse_path(path)) self._dst_prefix = [str(parse_path(path)[2]) for path in self.dst_paths] + print("return dst_prefixes for multicast", self._dst_prefix) return self._dst_prefix @property @@ -681,9 +692,10 @@ def chunk_request(server, chunk_batch, n_added): # send chunk requests to source gateways chunk_batch = [cr.chunk for cr in batch if cr.chunk is not None] # TODO: allow multiple partition ids per chunk - for chunk in chunk_batch: # assign job UUID as partition ID - chunk.partition_id = self.uuid + #for chunk in chunk_batch: # assign job UUID as partition ID + # chunk.partition_id = self.uuid min_idx = queue_size.index(min(queue_size)) + print([b.chunk.partition_id for b in batch if b.chunk]) n_added = 0 while n_added < len(chunk_batch): # TODO: should update every source instance queue size diff --git a/skyplane/gateway/gateway_program.py b/skyplane/gateway/gateway_program.py index c27427fd9..7c4d16f69 100644 --- a/skyplane/gateway/gateway_program.py +++ b/skyplane/gateway/gateway_program.py @@ -117,7 +117,7 @@ def add_operators(self, ops: List[GatewayOperator], parent_handle: Optional[str] parent_op = self._ops[parent_handle] if parent_handle else None ops_handles = [] for op in ops: - ops_handles.append(self.add_operator(op, parent_op, partition_id)) + ops_handles.append(self.add_operator(op, parent_op.handle, partition_id)) return ops_handles @@ -138,6 +138,8 @@ def to_dict(self): """ program_all = [] for partition_id, op_list in self._plan.items(): + partition_id = list(partition_id) # convert tuple to list + # build gateway program representation program = [] for op in op_list: @@ -147,11 +149,12 @@ def to_dict(self): exists = False for p in program_all: if p["value"] == program: # equivalent partition exists - p["partitions"].append(partition_id) + for pid in partition_id: + p["partitions"].append(str(pid)) exists = True break if not exists: - program_all.append({"value": program, "partitions": [partition_id]}) + program_all.append({"value": program, "partitions": str(partition_id)}) return program_all diff --git a/skyplane/obj_store/s3_interface.py b/skyplane/obj_store/s3_interface.py index f502cd77e..3ba6a9ee6 100644 --- a/skyplane/obj_store/s3_interface.py +++ b/skyplane/obj_store/s3_interface.py @@ -199,6 +199,8 @@ def upload_object( ): dst_object_name, src_file_path = str(dst_object_name), str(src_file_path) s3_client = self._s3_client() + print("Destination object name: ", dst_object_name) + print("Source file path: ", src_file_path) assert len(dst_object_name) > 0, f"Destination object name must be non-empty: '{dst_object_name}'" b64_md5sum = base64.b64encode(check_md5).decode("utf-8") if check_md5 else None checksum_args = dict(ContentMD5=b64_md5sum) if b64_md5sum else dict() diff --git a/skyplane/planner/planner.py b/skyplane/planner/planner.py index dde2846ba..8ca42e716 100644 --- a/skyplane/planner/planner.py +++ b/skyplane/planner/planner.py @@ -1,13 +1,12 @@ -from collections import defaultdict +import functools from importlib.resources import path -from typing import Dict, List, Optional, Tuple, Tuple -import os -import csv - +from typing import List, Optional, Tuple +import numpy as np +import collections +import pandas as pd +from skyplane.planner.solver_ilp import ThroughputSolverILP +from skyplane.planner.solver import ThroughputProblem, BroadcastProblem, BroadcastSolution, GBIT_PER_GBYTE from skyplane import compute -from skyplane.api.config import TransferConfig -from skyplane.utils import logger - from skyplane.planner.topology import TopologyPlan from skyplane.gateway.gateway_program import ( GatewayProgram, @@ -18,488 +17,643 @@ GatewayReceive, GatewaySend, ) - +import networkx as nx from skyplane.api.transfer_job import TransferJob -import json - -from skyplane.utils.fn import do_parallel -from skyplane.config_paths import config_path, azure_standardDv5_quota_path, aws_quota_path, gcp_quota_path, scp_quota_path -from skyplane.config import SkyplaneConfig +from pathlib import Path +from importlib.resources import files +from random import sample class Planner: - def __init__(self, transfer_config: TransferConfig, quota_limits_file: Optional[str] = None): - self.transfer_config = transfer_config - self.config = SkyplaneConfig.load_config(config_path) - self.n_instances = self.config.get_flag("max_instances") - - # Loading the quota information, add ibm cloud when it is supported - quota_limits = {} - if quota_limits_file is not None: - with open(quota_limits_file, "r") as f: - quota_limits = json.load(f) - else: - if os.path.exists(aws_quota_path): - with aws_quota_path.open("r") as f: - quota_limits["aws"] = json.load(f) - if os.path.exists(azure_standardDv5_quota_path): - with azure_standardDv5_quota_path.open("r") as f: - quota_limits["azure"] = json.load(f) - if os.path.exists(gcp_quota_path): - with gcp_quota_path.open("r") as f: - quota_limits["gcp"] = json.load(f) - if os.path.exists(scp_quota_path): - with scp_quota_path.open("r") as f: - quota_limits["scp"] = json.load(f) - self.quota_limits = quota_limits - - # Loading the vcpu information - a dictionary of dictionaries - # {"cloud_provider": {"instance_name": vcpu_cost}} - self.vcpu_info = defaultdict(dict) - with path("skyplane.data", "vcpu_info.csv") as file_path: - with open(file_path, "r") as csvfile: - reader = csv.reader(csvfile) - next(reader) # Skip the header row - - for row in reader: - instance_name, cloud_provider, vcpu_cost = row - vcpu_cost = int(vcpu_cost) - self.vcpu_info[cloud_provider][instance_name] = vcpu_cost - - def plan(self) -> TopologyPlan: - raise NotImplementedError - - def _vm_to_vcpus(self, cloud_provider: str, vm: str) -> int: - """Gets the vcpu_cost of the given vm instance (instance_name) - - :param cloud_provider: name of the cloud_provider - :type cloud_provider: str - :param instance_name: name of the vm instance - :type instance_name: str - """ - return self.vcpu_info[cloud_provider][vm] - - def _get_quota_limits_for(self, cloud_provider: str, region: str, spot: bool = False) -> Optional[int]: - """Gets the quota info from the saved files. Returns None if quota_info isn't loaded during `skyplane init` - or if the quota info doesn't include the region. - - :param cloud_provider: name of the cloud provider of the region - :type cloud_provider: str - :param region: name of the region for which to get the quota for - :type region: int - :param spot: whether to use spot specified by the user config (default: False) - :type spot: bool - """ - quota_limits = self.quota_limits.get(cloud_provider, None) - if not quota_limits: - # User needs to reinitialize to save the quota information - return None - if cloud_provider == "gcp": - region_family = "-".join(region.split("-")[:2]) - if region_family in quota_limits: - return quota_limits[region_family] - elif cloud_provider == "azure": - if region in quota_limits: - return quota_limits[region] - elif cloud_provider == "aws": - for quota in quota_limits: - if quota["region_name"] == region: - return quota["spot_standard_vcpus"] if spot else quota["on_demand_standard_vcpus"] - elif cloud_provider == "scp": - for quota in quota_limits: - if quota["service_zone_name"] == region: - return quota["on_demand_standard_vcpus"] - return None - - def _calculate_vm_types(self, region_tag: str) -> Optional[Tuple[str, int]]: - """Calculates the largest allowed vm type according to the regional quota limit as well as - how many of these vm types can we launch to avoid QUOTA_EXCEEDED errors. Returns None if quota - information wasn't properly loaded or allowed vcpu list is wrong. - - :param region_tag: tag of the node we are calculating the above for, example -> "aws:us-east-1" - :type region_tag: str - """ - cloud_provider, region = region_tag.split(":") - - # Get the quota limit - quota_limit = self._get_quota_limits_for( - cloud_provider=cloud_provider, region=region, spot=getattr(self.transfer_config, f"{cloud_provider}_use_spot_instances") - ) - - config_vm_type = getattr(self.transfer_config, f"{cloud_provider}_instance_class") - - # No quota limits (quota limits weren't initialized properly during skyplane init) - if quota_limit is None: - logger.warning( - f"Quota limit file not found for {region_tag}. Try running `skyplane init --reinit-{cloud_provider}` to load the quota information" - ) - # return default instance type and number of instances - return config_vm_type, self.n_instances - - config_vcpus = self._vm_to_vcpus(cloud_provider, config_vm_type) - if config_vcpus <= quota_limit: - return config_vm_type, quota_limit // config_vcpus - - vm_type, vcpus = None, None - for instance_name, vcpu_cost in sorted(self.vcpu_info[cloud_provider].items(), key=lambda x: x[1], reverse=True): - if vcpu_cost <= quota_limit: - vm_type, vcpus = instance_name, vcpu_cost - break - - # shouldn't happen, but just in case we use more complicated vm types in the future - assert vm_type is not None and vcpus is not None - - # number of instances allowed by the quota with the selected vm type - n_instances = quota_limit // vcpus - logger.warning( - f"Falling back to instance class `{vm_type}` at {region_tag} " - f"due to cloud vCPU limit of {quota_limit}. You can visit https://skyplane.org/en/latest/increase_vcpus.html " - "to learn more about how to increase your cloud vCPU limits for any cloud provider." - ) - return (vm_type, n_instances) - - def _get_vm_type_and_instances( - self, src_region_tag: Optional[str] = None, dst_region_tags: Optional[List[str]] = None - ) -> Tuple[Dict[str, str], int]: - """Dynamically calculates the vm type each region can use (both the source region and all destination regions) - based on their quota limits and calculates the number of vms to launch in all regions by conservatively - taking the minimum of all regions to stay consistent. - - :param src_region_tag: the source region tag (default: None) - :type src_region_tag: Optional[str] - :param dst_region_tags: a list of the destination region tags (defualt: None) - :type dst_region_tags: Optional[List[str]] - """ - - # One of them has to provided - # assert src_region_tag is not None or dst_region_tags is not None, "There needs to be at least one source or destination" - src_tags = [src_region_tag] if src_region_tag is not None else [] - dst_tags = dst_region_tags if dst_region_tags is not None else [] - - if src_region_tag: - assert len(src_region_tag.split(":")) == 2, f"Source region tag {src_region_tag} must be in the form of `cloud_provider:region`" - if dst_region_tags: - assert ( - len(dst_region_tags[0].split(":")) == 2 - ), f"Destination region tag {dst_region_tags} must be in the form of `cloud_provider:region`" - - # do_parallel returns tuples of (region_tag, (vm_type, n_instances)) - vm_info = do_parallel(self._calculate_vm_types, src_tags + dst_tags) - # Specifies the vm_type for each region - vm_types = {v[0]: v[1][0] for v in vm_info} # type: ignore - # Taking the minimum so that we can use the same number of instances for both source and destination - n_instances = min([self.n_instances] + [v[1][1] for v in vm_info]) # type: ignore - return vm_types, n_instances - - -class UnicastDirectPlanner(Planner): - # DO NOT USE THIS - broken for single-region transfers - def __init__(self, n_instances: int, n_connections: int, transfer_config: TransferConfig, quota_limits_file: Optional[str] = None): - super().__init__(transfer_config, quota_limits_file) + def __init__(self, n_instances: int, n_connections: int, n_partitions: Optional[int] = 1): self.n_instances = n_instances self.n_connections = n_connections + self.n_partitions = n_partitions + + def logical_plan(self) -> nx.DiGraph: + # create logical plan in nx.DiGraph format + raise NotImplementedError def plan(self, jobs: List[TransferJob]) -> TopologyPlan: - # make sure only single destination - for job in jobs: - assert len(job.dst_ifaces) == 1, f"DirectPlanner only support single destination jobs, got {len(job.dst_ifaces)}" + # create physical plan in TopologyPlan format + raise NotImplementedError + def verify_job_src_dsts(self, jobs: List[TransferJob], multicast=False) -> Tuple[str, List[str]]: src_region_tag = jobs[0].src_iface.region_tag() - dst_region_tag = jobs[0].dst_ifaces[0].region_tag() - - assert len(src_region_tag.split(":")) == 2, f"Source region tag {src_region_tag} must be in the form of `cloud_provider:region`" - assert ( - len(dst_region_tag.split(":")) == 2 - ), f"Destination region tag {dst_region_tag} must be in the form of `cloud_provider:region`" - - # jobs must have same sources and destinations - for job in jobs[1:]: - assert job.src_iface.region_tag() == src_region_tag, "All jobs must have same source region" - assert job.dst_ifaces[0].region_tag() == dst_region_tag, "All jobs must have same destination region" - - plan = TopologyPlan(src_region_tag=src_region_tag, dest_region_tags=[dst_region_tag]) - # Dynammically calculate n_instances based on quota limits - vm_types, n_instances = self._get_vm_type_and_instances(src_region_tag=src_region_tag, dst_region_tags=[dst_region_tag]) + if multicast: + # multicast checking + dst_region_tags = [iface.region_tag() for iface in jobs[0].dst_ifaces] - # TODO: support on-sided transfers but not requiring VMs to be created in source/destination regions - for i in range(n_instances): - plan.add_gateway(src_region_tag, vm_types[src_region_tag]) - plan.add_gateway(dst_region_tag, vm_types[dst_region_tag]) - - # ids of gateways in dst region - dst_gateways = plan.get_region_gateways(dst_region_tag) - - src_program = GatewayProgram() - dst_program = GatewayProgram() - - for job in jobs: - src_bucket = job.src_iface.bucket() - dst_bucket = job.dst_ifaces[0].bucket() - - # give each job a different partition id, so we can read/write to different buckets - partition_id = job.uuid - - # source region gateway program - obj_store_read = src_program.add_operator( - GatewayReadObjectStore(src_bucket, src_region_tag, self.n_connections), partition_id=partition_id + # jobs must have same sources and destinations + for job in jobs[1:]: + assert job.src_iface.region_tag() == src_region_tag, "All jobs must have same source region" + assert [iface.region_tag() for iface in job.dst_ifaces] == dst_region_tags, "Add jobs must have same destination set" + else: + # unicast checking + for job in jobs: + assert len(job.dst_ifaces) == 1, f"DirectPlanner only support single destination jobs, got {len(job.dst_ifaces)}" + + # jobs must have same sources and destinations + dst_region_tag = jobs[0].dst_ifaces[0].region_tag() + for job in jobs[1:]: + assert job.src_iface.region_tag() == src_region_tag, "All jobs must have same source region" + assert job.dst_ifaces[0].region_tag() == dst_region_tag, "All jobs must have same destination region" + dst_region_tags = [dst_region_tag] + + return src_region_tag, dst_region_tags + + @functools.lru_cache(maxsize=None) + def make_nx_graph(self, tp_grid_path: Optional[Path] = files("data") / "throughput.csv") -> nx.DiGraph: + # create throughput / cost graph for all regions for planner + G = nx.DiGraph() + throughput = pd.read_csv(tp_grid_path) + for _, row in throughput.iterrows(): + if row["src_region"] == row["dst_region"]: + continue + G.add_edge(row["src_region"], row["dst_region"], cost=None, throughput=row["throughput_sent"] / 1e9) + + for edge in G.edges.data(): + if edge[-1]["cost"] is None: + edge[-1]["cost"] = compute.CloudProvider.get_transfer_cost(edge[0], edge[1]) + + assert all([edge[-1]["cost"] is not None for edge in G.edges.data()]) + return G + + def add_src_or_overlay_operator( + self, + solution_graph: nx.DiGraph, + gateway_program: GatewayProgram, + region: str, + partition_ids: List[int], + partition_offset: int, + plan: TopologyPlan, + bucket_info: Optional[Tuple[str, str]] = None, + dst_op: Optional[GatewayReceive] = None, + ) -> bool: + """ + :param solution_graph: nx.DiGraph of solution + :param gateway_program: GatewayProgram of region to add operator to + :param region: region to add operator to + :param partition_ids: list of partition ids to add operator to + :param partition_offset: offset of partition ids + :param plan: TopologyPlan of solution [for getting gateway ids] + :param bucket_info: tuple of (bucket_name, bucket_region) for object store + :param dst_op: if None, then this is either the source node or a overlay node; otherwise, this is the destination overlay node + """ + g = solution_graph + # partition_ids are set of ids that follow the same path from the out edges of the region + any_id = partition_ids[0] - partition_offset + next_regions = set([edge[1] for edge in g.out_edges(region, data=True) if str(any_id) in edge[-1]["partitions"]]) + + # if partition_ids does not have a next region, then we cannot add an operator + if len(next_regions) == 0: + print( + f"Region {region}, any id: {any_id}, partition ids: {partition_ids}, has no next region to forward data to: {g.out_edges(region, data=True)}" ) - mux_or = src_program.add_operator(GatewayMuxOr(), parent_handle=obj_store_read, partition_id=partition_id) - for i in range(n_instances): - src_program.add_operator( - GatewaySend( - target_gateway_id=dst_gateways[i].gateway_id, - region=src_region_tag, - num_connections=self.n_connections, - compress=True, - encrypt=True, - ), - parent_handle=mux_or, - partition_id=partition_id, + return + + # identify if this is a destination overlay node or not + if dst_op is None: + # source node or overlay node + # TODO: add generate data locally operator + if bucket_info is None: + receive_op = GatewayReceive() + else: + receive_op = GatewayReadObjectStore( + bucket_name=bucket_info[0], bucket_region=bucket_info[1], num_connections=self.n_connections ) - - # dst region gateway program - recv_op = dst_program.add_operator(GatewayReceive(decompress=True, decrypt=True), partition_id=partition_id) - dst_program.add_operator( - GatewayWriteObjectStore(dst_bucket, dst_region_tag, self.n_connections), parent_handle=recv_op, partition_id=partition_id + else: + # destination overlay node, dst_op is the parent node + receive_op = dst_op + + # find set of regions to send to for all partitions in partition_ids + + region_to_id_map = {} + for next_region in next_regions: + region_to_id_map[next_region] = [] + for i in range(solution_graph.nodes[next_region]["num_vms"]): + region_to_id_map[next_region].append(plan.get_region_gateways(next_region)[i].gateway_id) + + # use muxand or muxor for partition_id + operation = "MUX_AND" if len(next_regions) > 1 else "MUX_OR" + mux_op = GatewayMuxAnd() if len(next_regions) > 1 else GatewayMuxOr() + + # non-dst node: add receive_op into gateway program + if dst_op is None: + gateway_program.add_operator(op=receive_op, partition_id=tuple(partition_ids)) + + # MUX_AND: send this partition to multiple regions + if operation == "MUX_AND": + if dst_op is not None and dst_op.op_type == "mux_and": + mux_op = receive_op + else: # do not add any nested mux_and if dst_op parent is mux_and + gateway_program.add_operator(op=mux_op, parent_handle=receive_op.handle, partition_id=tuple(partition_ids)) + + for next_region, next_region_ids in region_to_id_map.items(): + send_ops = [ + GatewaySend(target_gateway_id=id, region=next_region, num_connections=self.n_connections) for id in next_region_ids + ] + + # if there is more than one region to forward data to, add MUX_OR + if len(next_region_ids) > 1: + mux_or_op = GatewayMuxOr() + gateway_program.add_operator(op=mux_or_op, parent_handle=mux_op.handle, partition_id=tuple(partition_ids)) + gateway_program.add_operators(ops=send_ops, parent_handle=mux_or_op.handle, partition_id=tuple(partition_ids)) + else: + # otherwise, the parent of send_op is mux_op ("MUX_AND") + assert len(send_ops) == 1 + gateway_program.add_operator(op=send_ops[0], parent_handle=mux_op.handle, partition_id=tuple(partition_ids)) + else: + # only send this partition to a single region + assert len(region_to_id_map) == 1 + next_region = list(region_to_id_map.keys())[0] + ids = [id for next_region_ids in region_to_id_map.values() for id in next_region_ids] + send_ops = [GatewaySend(target_gateway_id=id, region=next_region, num_connections=self.n_connections) for id in ids] + + # if num of gateways > 1, then connect to MUX_OR + if len(ids) > 1: + gateway_program.add_operator(op=mux_op, parent_handle=receive_op.handle, partition_id=tuple(partition_ids)) + gateway_program.add_operators(ops=send_ops, parent_handle=mux_op.handle) + else: + gateway_program.add_operators(ops=send_ops, parent_handle=receive_op.handle, partition_id=tuple(partition_ids)) + + return True + + def add_dst_operator( + self, + solution_graph, + gateway_program: GatewayProgram, + region: str, + partition_ids: List[int], + partition_offset: int, + plan: TopologyPlan, + obj_store: Tuple[str, str] = None, + ): + # operator that receives data + receive_op = GatewayReceive() + gateway_program.add_operator(receive_op, partition_id=tuple(partition_ids)) + + # operator that writes to the object store + write_op = GatewayWriteObjectStore(bucket_name=obj_store[0], bucket_region=obj_store[1], num_connections=self.n_connections) + + g = solution_graph + + # partition_ids are ids that follow the same path from the out edges of the region + any_id = partition_ids[0] - partition_offset + next_regions = set([edge[1] for edge in g.out_edges(region, data=True) if str(any_id) in edge[-1]["partitions"]]) + + # if no regions to forward data to, write to the object store + if len(next_regions) == 0: + gateway_program.add_operator(write_op, parent_handle=receive_op.handle, partition_id=tuple(partition_ids)) + + # otherwise, receive and write to the object store, then forward data to next regions + else: + mux_and_op = GatewayMuxAnd() + # receive and write + gateway_program.add_operator(mux_and_op, parent_handle=receive_op.handle, partition_id=tuple(partition_ids)) + gateway_program.add_operator(write_op, parent_handle=mux_and_op.handle, partition_id=tuple(partition_ids)) + + # forward: destination nodes are also forwarders + self.add_src_or_overlay_operator( + solution_graph, gateway_program, region, partition_ids, partition_offset, plan, dst_op=mux_and_op ) - # update cost per GB - plan.cost_per_gb += compute.CloudProvider.get_transfer_cost(src_region_tag, dst_region_tag) + def logical_plan_to_topology_plan(self, jobs: List[TransferJob], solution_graph: nx.graph) -> TopologyPlan: + """ + Given a logical plan, construct a gateway program for each region in the logical plan for the given jobs. + """ + # get source and destination regions + src_ifaces, dst_ifaces = [job.src_iface for job in jobs], [job.dst_ifaces for job in jobs] + src_region_tag = src_ifaces[0].region_tag() + dst_region_tags = [dst_iface.region_tag() for dst_iface in dst_ifaces[0]] + + # map from the node to the gateway program + region_to_gateway_program = {region: GatewayProgram() for region in solution_graph.nodes} + + # construct TopologyPlan for all the regions in solution_graph + overlay_region_tags = [node for node in solution_graph.nodes if node != src_region_tag and node not in dst_region_tags] + plan = TopologyPlan(src_region_tag=src_region_tag, dest_region_tags=dst_region_tags, overlay_region_tags=overlay_region_tags) + for node in solution_graph.nodes: + plan.add_gateway(node) + + # iterate through all the jobs + for i in range(len(src_ifaces)): + src_bucket = src_ifaces[i].bucket() + dst_buckets = {dst_iface[i].region_tag(): dst_iface[i].bucket() for dst_iface in dst_ifaces} + + # iterate through all the regions in the solution graph + for node in solution_graph.nodes: + node_gateway_program = region_to_gateway_program[node] + partition_to_next_regions = {} + + # give each job a different partition offset i, so we can read/write to different buckets + for j in range(i, i + self.n_partitions): + partition_to_next_regions[j] = set( + [edge[1] for edge in solution_graph.out_edges(node, data=True) if str(j) in edge[-1]["partitions"]] + ) - # set gateway programs - plan.set_gateway_program(src_region_tag, src_program) - plan.set_gateway_program(dst_region_tag, dst_program) + keys_per_set = collections.defaultdict(list) + for key, value in partition_to_next_regions.items(): + keys_per_set[frozenset(value)].append(key) + + list_of_partitions = list(keys_per_set.values()) + + # source node: read from object store or generate random data, then forward data + for partitions in list_of_partitions: + if node == src_region_tag: + self.add_src_or_overlay_operator( + solution_graph, + node_gateway_program, + node, + partitions, + partition_offset=i, + plan=plan, + bucket_info=(src_bucket, node), + ) + + # dst receive data, write to object store, forward data if needed + elif node in dst_region_tags: + dst_bucket = dst_buckets[node] + self.add_dst_operator( + solution_graph, + node_gateway_program, + node, + partitions, + partition_offset=i, + plan=plan, + obj_store=(dst_bucket, node), + ) + + # overlay node only forward data + else: + self.add_src_or_overlay_operator( + solution_graph, node_gateway_program, node, partitions, partition_offset=i, plan=plan, bucket_info=None + ) + region_to_gateway_program[node] = node_gateway_program + assert len(region_to_gateway_program) > 0, f"Empty gateway program {node}" + + for node in solution_graph.nodes: + plan.set_gateway_program(node, region_to_gateway_program[node]) + + for edge in solution_graph.edges.data(): + src_region, dst_region = edge[0], edge[1] + plan.cost_per_gb += compute.CloudProvider.get_transfer_cost(src_region, dst_region) * ( + len(edge[-1]["partitions"]) / self.n_partitions + ) return plan class MulticastDirectPlanner(Planner): - def __init__(self, n_instances: int, n_connections: int, transfer_config: TransferConfig, quota_limits_file: Optional[str] = None): - super().__init__(transfer_config, quota_limits_file) - self.n_instances = n_instances - self.n_connections = n_connections + def __init__(self, n_instances: int, n_connections: int, n_partitions: Optional[int] = 1): + super().__init__(n_instances, n_connections, n_partitions) - def plan(self, jobs: List[TransferJob]) -> TopologyPlan: - src_region_tag = jobs[0].src_iface.region_tag() - dst_region_tags = [iface.region_tag() for iface in jobs[0].dst_ifaces] + def logical_plan(self, src_region: str, dst_regions: List[str]) -> nx.DiGraph: + graph = nx.DiGraph() + graph.add_node(src_region) + for dst_region in dst_regions: + graph.add_node(dst_region) + graph.add_edge(src_region, dst_region, partitions=[str(i) for i in range(self.n_partitions)]) - # jobs must have same sources and destinations - for job in jobs[1:]: - assert job.src_iface.region_tag() == src_region_tag, "All jobs must have same source region" - assert [iface.region_tag() for iface in job.dst_ifaces] == dst_region_tags, "Add jobs must have same destination set" + for node in graph.nodes: + graph.nodes[node]["num_vms"] = self.n_instances + return graph - plan = TopologyPlan(src_region_tag=src_region_tag, dest_region_tags=dst_region_tags) + def plan(self, jobs: List[TransferJob]) -> TopologyPlan: + src_region_tag, dst_region_tags = self.verify_job_src_dsts(jobs) + solution_graph = self.logical_plan(src_region_tag, dst_region_tags) + return self.logical_plan_to_topology_plan(jobs, solution_graph) + + +class MulticastILPPlanner(Planner): + def __init__( + self, + n_instances: int, + n_connections: int, + target_time: Optional[float] = 10000, + n_partitions: Optional[int] = 1, + aws_only: bool = False, + gcp_only: bool = False, + azure_only: bool = False, + ): + super().__init__(n_instances, n_connections, n_partitions) + self.target_time = target_time + self.aws_only = aws_only + self.gcp_only = gcp_only + self.azure_only = azure_only + self.G = super().make_nx_graph() + + def multicast_solution_to_nxgraph(self, solution: BroadcastSolution) -> nx.DiGraph: + """ + Convert ILP solution to logical plan in nx graph + """ + v_result = solution.var_instances_per_region + result = np.array(solution.var_edge_partitions) + result_g = nx.DiGraph() # solution nx graph + for i in range(result.shape[0]): + edge = solution.var_edges[i] + partitions = [str(partition_i) for partition_i in range(result.shape[1]) if result[i][partition_i] > 0.5] + + if len(partitions) == 0: + continue + + src_node, dst_node = edge[0], edge[1] + result_g.add_edge( + src_node, + dst_node, + partitions=partitions, + throughput=self.G[src_node][dst_node]["throughput"], + cost=self.G[src_node][dst_node]["cost"], + ) + + for i in range(len(v_result)): + num_vms = int(v_result[i]) + node = solution.var_nodes[i] + if node in result_g.nodes: + result_g.nodes[node]["num_vms"] = num_vms + + def logical_plan( + self, + src_region: str, + dst_regions: List[str], + gbyte_to_transfer: int = 1, + filter_node: bool = False, + filter_edge: bool = False, + solver_verbose: bool = False, + save_lp_path: Optional[str] = None, + solver: Optional[str] = None, + ) -> nx.DiGraph: + import cvxpy as cp + + if solver is None: + solver = cp.GUROBI + + problem = BroadcastProblem( + src=src_region, + dsts=dst_regions, + gbyte_to_transfer=gbyte_to_transfer, + instance_limit=self.n_instances, + num_partitions=self.n_partitions, + required_time_budget=self.target_time, + ) - # Dynammically calculate n_instances based on quota limits - if src_region_tag.split(":")[0] == "test": - vm_types = None - n_instances = self.n_instances + g = self.G + + # node-approximation + if filter_node: + src_dst_li = [problem.src] + problem.dsts + sampled = [i for i in sample(list(self.G.nodes), 15) if i not in src_dst_li] + g = g.subgraph(src_dst_li + sampled).copy() + print(f"Filter node (only use): {src_dst_li + sampled}") + + cost = np.array([e[2] for e in g.edges(data="cost")]) + tp = np.array([e[2] for e in g.edges(data="throughput")]) + + edges = list(g.edges) + nodes = list(g.nodes) + num_edges, num_nodes = len(edges), len(nodes) + num_dest = len(problem.dsts) + print(f"Num edges: {num_edges}, num nodes: {num_nodes}, num dest: {num_dest}, runtime budget: {problem.required_time_budget}s") + + partition_size_gb = problem.gbyte_to_transfer / problem.num_partitions + partition_size_gbit = partition_size_gb * GBIT_PER_GBYTE + print("Partition size (gbit): ", partition_size_gbit) + + # define variables + p = cp.Variable((num_edges, problem.num_partitions), boolean=True) # whether edge is carrying partition + n = cp.Variable((num_nodes), boolean=True) # whether node transfers partition + f = cp.Variable((num_nodes * problem.num_partitions, num_nodes + 1), integer=True) # enforce flow conservation + v = cp.Variable((num_nodes), integer=True) # number of VMs per region + + # define objective + egress_cost = cp.sum(cost @ p) * partition_size_gb + instance_cost = cp.sum(v) * (problem.cost_per_instance_hr / 3600) * problem.required_time_budget + tot_cost = egress_cost + instance_cost + obj = cp.Minimize(tot_cost) + + # define constants + constraints = [] + + # constraints on VM per region + for i in range(num_nodes): + constraints.append(v[i] <= problem.instance_limit) + constraints.append(v[i] >= 0) + + # constraints to enforce flow between source/dest nodes + for c in range(problem.num_partitions): + for i in range(num_nodes): + for j in range(num_nodes + 1): + if i != j: + if j != num_nodes: + edge = (nodes[i], nodes[j]) + + constraints.append(f[c * num_nodes + i][j] <= p[edges.index(edge)][c] * num_dest) + # p = 0 -> f <= 0 + # p = 1 -> f <= num_dest + constraints.append(f[c * num_nodes + i][j] >= (p[edges.index(edge)][c] - 1) * (num_dest + 1) + 1) + # p = 0 -> f >= -(num_dest) + # p = 1 -> f >= 1 + + constraints.append(f[c * num_nodes + i][j] == -f[c * num_nodes + j][i]) + + # capacity constraint for special node + else: + if nodes[i] in problem.dsts: # only connected to destination nodes + constraints.append(f[c * num_nodes + i][j] <= 1) + else: + constraints.append(f[c * num_nodes + i][j] <= 0) + else: + constraints.append(f[c * num_nodes + i][i] == 0) + + # flow conservation + if nodes[i] != problem.src and i != num_nodes + 1: + constraints.append(cp.sum(f[c * num_nodes + i]) == 0) + + # source must have outgoing flow + constraints.append(cp.sum(f[c * num_nodes + nodes.index(problem.src), :]) == num_dest) + + # special node (connected to all destinations) must recieve all flow + constraints.append(cp.sum(f[c * num_nodes : (c + 1) * num_nodes, -1]) == num_dest) + + # node contained if edge is contained + for edge in edges: + constraints.append(n[nodes.index(edge[0])] >= cp.max(p[edges.index(edge)])) + constraints.append(n[nodes.index(edge[1])] >= cp.max(p[edges.index(edge)])) + + # edge approximation + if filter_edge: + for edge in edges: + if edge[0] != problem.src and edge[1] not in problem.dsts: + # cannot be in graph + constraints.append(p[edges.index(edge)] == 0) + + # throughput constraint + for edge_i in range(num_edges): + node_i = nodes.index(edge[0]) + constraints.append(cp.sum(p[edge_i] * partition_size_gbit) <= problem.required_time_budget * tp[edge_i] * v[node_i]) + + # instance limits + for node in nodes: + region = node.split(":")[0] + if region == "aws": + ingress_limit_gbps, egress_limit_gbps = problem.aws_instance_throughput_limit + elif region == "gcp": + ingress_limit_gbps, egress_limit_gbps = problem.gcp_instance_throughput_limit + elif region == "azure": + ingress_limit_gbps, egress_limit_gbps = problem.azure_instance_throughput_limit + elif region == "cloudflare": # TODO: not supported yet in the tput / cost graph + ingress_limit_gbps, egress_limit_gbps = 1, 1 + + node_i = nodes.index(node) + # egress + i = np.zeros(num_edges) + for e in g.edges: + if e[0] == node: # edge goes to dest + i[edges.index(e)] = 1 + + constraints.append(cp.sum(i @ p) * partition_size_gbit <= problem.required_time_budget * egress_limit_gbps * v[node_i]) + + # ingress + i = np.zeros(num_edges) + for e in g.edges: + # edge goes to dest + if e[1] == node: + i[edges.index(e)] = 1 + constraints.append(cp.sum(i @ p) * partition_size_gbit <= problem.required_time_budget * ingress_limit_gbps * v[node_i]) + + print("Define problem done.") + + # solve + prob = cp.Problem(obj, constraints) + if solver == cp.GUROBI or solver == "gurobi": + solver_options = {} + solver_options["Threads"] = 1 + if save_lp_path: + solver_options["ResultFile"] = str(save_lp_path) + if not solver_verbose: + solver_options["OutputFlag"] = 0 + cost = prob.solve(verbose=solver_verbose, qcp=True, solver=cp.GUROBI, reoptimize=True, **solver_options) + elif solver == cp.CBC or solver == "cbc": + solver_options = {} + solver_options["maximumSeconds"] = 60 + solver_options["numberThreads"] = 1 + cost = prob.solve(verbose=solver_verbose, solver=cp.CBC, **solver_options) else: - vm_types, n_instances = self._get_vm_type_and_instances(src_region_tag=src_region_tag, dst_region_tags=dst_region_tags) - - # TODO: support on-sided transfers but not requiring VMs to be created in source/destination regions - for i in range(n_instances): - plan.add_gateway(src_region_tag, vm_types[src_region_tag] if vm_types else None) - for dst_region_tag in dst_region_tags: - plan.add_gateway(dst_region_tag, vm_types[dst_region_tag] if vm_types else None) - - # initialize gateway programs per region - dst_program = {dst_region: GatewayProgram() for dst_region in dst_region_tags} - src_program = GatewayProgram() - - # iterate through all jobs - for job in jobs: - src_bucket = job.src_iface.bucket() - src_region_tag = job.src_iface.region_tag() - src_provider = src_region_tag.split(":")[0] - - # give each job a different partition id, so we can read/write to different buckets - partition_id = job.uuid - - # source region gateway program - obj_store_read = src_program.add_operator( - GatewayReadObjectStore(src_bucket, src_region_tag, self.n_connections), partition_id=partition_id + cost = prob.solve(solver=solver, verbose=solver_verbose) + + if prob.status == "optimal": + solution = BroadcastSolution( + problem=problem, + is_feasible=True, + var_edges=edges, + var_nodes=nodes, + var_edge_partitions=p.value, + var_node_transfer_partitions=n.value, + var_instances_per_region=v.value, + var_flow=f.value, + cost_egress=egress_cost.value, + cost_instance=instance_cost.value, + cost_total=tot_cost.value, ) + else: + solution = BroadcastSolution(problem=problem, is_feasible=False, extra_data=dict(status=prob.status)) - # send to all destination - mux_and = src_program.add_operator(GatewayMuxAnd(), parent_handle=obj_store_read, partition_id=partition_id) - dst_prefixes = job.dst_prefixes - for i in range(len(job.dst_ifaces)): - dst_iface = job.dst_ifaces[i] - dst_prefix = dst_prefixes[i] - dst_region_tag = dst_iface.region_tag() - dst_bucket = dst_iface.bucket() - dst_gateways = plan.get_region_gateways(dst_region_tag) - - # special case where destination is same region as source - if dst_region_tag == src_region_tag: - src_program.add_operator( - GatewayWriteObjectStore(dst_bucket, dst_region_tag, self.n_connections, key_prefix=dst_prefix), - parent_handle=mux_and, - partition_id=partition_id, - ) - continue - - # can send to any gateway in region - mux_or = src_program.add_operator(GatewayMuxOr(), parent_handle=mux_and, partition_id=partition_id) - for i in range(n_instances): - private_ip = False - if dst_gateways[i].provider == "gcp" and src_provider == "gcp": - # print("Using private IP for GCP to GCP transfer", src_region_tag, dst_region_tag) - private_ip = True - src_program.add_operator( - GatewaySend( - target_gateway_id=dst_gateways[i].gateway_id, - region=dst_region_tag, - num_connections=int(self.n_connections / len(dst_gateways)), - private_ip=private_ip, - compress=self.transfer_config.use_compression, - encrypt=self.transfer_config.use_e2ee, - ), - parent_handle=mux_or, - partition_id=partition_id, - ) + return self.multicast_solution_to_nxgraph(solution) - # each gateway also recieves data from source - recv_op = dst_program[dst_region_tag].add_operator( - GatewayReceive(decompress=self.transfer_config.use_compression, decrypt=self.transfer_config.use_e2ee), - partition_id=partition_id, - ) - dst_program[dst_region_tag].add_operator( - GatewayWriteObjectStore(dst_bucket, dst_region_tag, self.n_connections, key_prefix=dst_prefix), - parent_handle=recv_op, - partition_id=partition_id, - ) + def plan(self, jobs: List[TransferJob]) -> TopologyPlan: + src_region_tag, dst_region_tags = self.verify_job_src_dsts(jobs, multicast=True) + solution_graph = self.logical_plan(src_region_tag, dst_region_tags) + return self.logical_plan_to_topology_plan(jobs, solution_graph) - # update cost per GB - plan.cost_per_gb += compute.CloudProvider.get_transfer_cost(src_region_tag, dst_region_tag) - # set gateway programs - plan.set_gateway_program(src_region_tag, src_program) - for dst_region_tag, program in dst_program.items(): - if dst_region_tag != src_region_tag: # don't overwrite - plan.set_gateway_program(dst_region_tag, program) - return plan +class MulticastMDSTPlanner(Planner): + def __init__(self, n_instances: int, n_connections: int, n_partitions: Optional[int] = 1): + super().__init__(n_instances, n_connections, n_partitions) + self.G = super().make_nx_graph() + def logical_plan(self, src_region: str, dst_regions: List[str]) -> nx.DiGraph: + h = self.G.copy() + h.remove_edges_from(list(h.in_edges(src_region)) + list(nx.selfloop_edges(h))) -class DirectPlannerSourceOneSided(MulticastDirectPlanner): - """Planner that only creates VMs in the source region""" + DST_graph = nx.algorithms.tree.Edmonds(h.subgraph([src_region] + dst_regions)) + opt_DST = DST_graph.find_optimum(attr="cost", kind="min", preserve_attrs=True, style="arborescence") - def plan(self, jobs: List[TransferJob]) -> TopologyPlan: - src_region_tag = jobs[0].src_iface.region_tag() - dst_region_tags = [iface.region_tag() for iface in jobs[0].dst_ifaces] - # jobs must have same sources and destinations - for job in jobs[1:]: - assert job.src_iface.region_tag() == src_region_tag, "All jobs must have same source region" - assert [iface.region_tag() for iface in job.dst_ifaces] == dst_region_tags, "Add jobs must have same destination set" + # Construct MDST graph + MDST_graph = nx.DiGraph() + for edge in list(opt_DST.edges()): + s, d = edge[0], edge[1] + MDST_graph.add_edge(s, d, partitions=[str(i) for i in list(range(self.num_partitions))]) - plan = TopologyPlan(src_region_tag=src_region_tag, dest_region_tags=dst_region_tags) + for node in MDST_graph.nodes: + MDST_graph.nodes[node]["num_vms"] = self.n_instances - # Dynammically calculate n_instances based on quota limits - vm_types, n_instances = self._get_vm_type_and_instances(src_region_tag=src_region_tag) + return MDST_graph - # TODO: support on-sided transfers but not requiring VMs to be created in source/destination regions - for i in range(n_instances): - plan.add_gateway(src_region_tag, vm_types[src_region_tag]) + def plan(self, jobs: List[TransferJob]) -> TopologyPlan: + src_region_tag, dst_region_tags = self.verify_job_src_dsts(jobs, multicast=True) + solution_graph = self.logical_plan(src_region_tag, dst_region_tags) + return self.logical_plan_to_topology_plan(jobs, solution_graph) - # initialize gateway programs per region - src_program = GatewayProgram() - # iterate through all jobs - for job in jobs: - src_bucket = job.src_iface.bucket() - src_region_tag = job.src_iface.region_tag() - src_provider = src_region_tag.split(":")[0] +class UnicastDirectPlanner(Planner): + def __init__(self, n_instances: int, n_connections: int, n_partitions: Optional[int] = 1): + super().__init__(n_instances, n_connections, n_partitions) - # give each job a different partition id, so we can read/write to different buckets - partition_id = job.uuid + def logical_plan(self, src_region: str, dst_regions: List[str]) -> nx.DiGraph: + graph = nx.DiGraph() + graph.add_node(src_region) + for dst_region in dst_regions: + graph.add_node(dst_region) + graph.add_edge(src_region, dst_region, partitions=[str(i) for i in range(self.n_partitions)]) - # source region gateway program - obj_store_read = src_program.add_operator( - GatewayReadObjectStore(src_bucket, src_region_tag, self.n_connections), partition_id=partition_id - ) - # send to all destination - mux_and = src_program.add_operator(GatewayMuxAnd(), parent_handle=obj_store_read, partition_id=partition_id) - dst_prefixes = job.dst_prefixes - for i in range(len(job.dst_ifaces)): - dst_iface = job.dst_ifaces[i] - dst_prefix = dst_prefixes[i] - dst_region_tag = dst_iface.region_tag() - dst_bucket = dst_iface.bucket() - plan.get_region_gateways(dst_region_tag) - - # special case where destination is same region as source - src_program.add_operator( - GatewayWriteObjectStore(dst_bucket, dst_region_tag, self.n_connections, key_prefix=dst_prefix), - parent_handle=mux_and, - partition_id=partition_id, - ) - # update cost per GB - plan.cost_per_gb += compute.CloudProvider.get_transfer_cost(src_region_tag, dst_region_tag) + for node in graph.nodes: + graph.nodes[node]["num_vms"] = self.n_instances - # set gateway programs - plan.set_gateway_program(src_region_tag, src_program) - return plan + return graph + def plan(self, jobs: List[TransferJob]) -> TopologyPlan: + src_region_tag, dst_region_tag = self.verify_job_src_dsts(jobs) + solution_graph = self.logical_plan(src_region_tag, dst_region_tag) + return self.logical_plan_to_topology_plan(jobs, solution_graph) + + +class UnicastILPPlanner(Planner): + def __init__(self, n_instances: int, n_connections: int, required_throughput_gbits: float, n_partitions: Optional[int] = 1): + super().__init__(n_instances, n_connections, n_partitions) + self.solver_required_throughput_gbits = required_throughput_gbits + + def logical_plan(self, src_region: str, dst_regions: List[str]) -> nx.DiGraph: + problem = ThroughputProblem( + src=src_region, + dst=dst_regions, + required_throughput_gbits=self.solver_required_throughput_gbits, + gbyte_to_transfer=1, + instance_limit=self.n_instances, + ) + + with path("skyplane.data", "throughput.csv") as solver_throughput_grid: + tput = ThroughputSolverILP(solver_throughput_grid) + solution = tput.solve_min_cost(problem) -class DirectPlannerDestOneSided(MulticastDirectPlanner): - """Planner that only creates instances in the destination region""" + if not solution.is_feasible: + raise RuntimeError("No feasible solution found") + + return tput.to_replication_topology(solution) def plan(self, jobs: List[TransferJob]) -> TopologyPlan: - # only create in destination region - src_region_tag = jobs[0].src_iface.region_tag() - dst_region_tags = [iface.region_tag() for iface in jobs[0].dst_ifaces] - # jobs must have same sources and destinations - for job in jobs[1:]: - assert job.src_iface.region_tag() == src_region_tag, "All jobs must have same source region" - assert [iface.region_tag() for iface in job.dst_ifaces] == dst_region_tags, "Add jobs must have same destination set" - - plan = TopologyPlan(src_region_tag=src_region_tag, dest_region_tags=dst_region_tags) - - # Dynammically calculate n_instances based on quota limits - vm_types, n_instances = self._get_vm_type_and_instances(dst_region_tags=dst_region_tags) - - # TODO: support on-sided transfers but not requiring VMs to be created in source/destination regions - for i in range(n_instances): - for dst_region_tag in dst_region_tags: - plan.add_gateway(dst_region_tag, vm_types[dst_region_tag]) - - # initialize gateway programs per region - dst_program = {dst_region: GatewayProgram() for dst_region in dst_region_tags} - - # iterate through all jobs - for job in jobs: - src_bucket = job.src_iface.bucket() - src_region_tag = job.src_iface.region_tag() - src_provider = src_region_tag.split(":")[0] - - partition_id = job.uuid - - # send to all destination - dst_prefixes = job.dst_prefixes - for i in range(len(job.dst_ifaces)): - dst_iface = job.dst_ifaces[i] - dst_prefix = dst_prefixes[i] - dst_region_tag = dst_iface.region_tag() - dst_bucket = dst_iface.bucket() - plan.get_region_gateways(dst_region_tag) - - # source region gateway program - obj_store_read = dst_program[dst_region_tag].add_operator( - GatewayReadObjectStore(src_bucket, src_region_tag, self.n_connections), partition_id=partition_id - ) + src_region_tag, dst_region_tag = self.verify_job_src_dsts(jobs) + solution_graph = self.logical_plan(src_region_tag, dst_region_tag) + return self.logical_plan_to_topology_plan(jobs, solution_graph) - dst_program[dst_region_tag].add_operator( - GatewayWriteObjectStore(dst_bucket, dst_region_tag, self.n_connections, key_prefix=dst_prefix), - parent_handle=obj_store_read, - partition_id=partition_id, - ) - # update cost per GB - plan.cost_per_gb += compute.CloudProvider.get_transfer_cost(src_region_tag, dst_region_tag) +class UnicastRONSolverPlanner(Planner): + def __init__(self, n_instances: int, n_connections: int, required_throughput_gbits: float, n_partitions: Optional[int] = 1): + super().__init__(n_instances, n_connections, n_partitions) + self.solver_required_throughput_gbits = required_throughput_gbits - # set gateway programs - for dst_region_tag, program in dst_program.items(): - plan.set_gateway_program(dst_region_tag, program) - return plan + def logical_plan(self, src_region: str, dst_regions: List[str]) -> nx.DiGraph: + raise NotImplementedError("RON solver not implemented yet") + + def plan(self, jobs: List[TransferJob]) -> TopologyPlan: + raise NotImplementedError("RON solver not implemented yet") diff --git a/skyplane/planner/solver.py b/skyplane/planner/solver.py index 5a18f1ed4..65e1a740d 100644 --- a/skyplane/planner/solver.py +++ b/skyplane/planner/solver.py @@ -15,6 +15,86 @@ GBIT_PER_GBYTE = 8 +@dataclass +class BroadcastProblem: + src: str + dsts: List[str] + + gbyte_to_transfer: float + instance_limit: int # max # of vms per region + num_partitions: int + + required_time_budget: float = 10 # ILP specific, default to 10s + + const_throughput_grid_gbits: Optional[np.ndarray] = None # if not set, load from profiles + const_cost_per_gb_grid: Optional[np.ndarray] = None # if not set, load from profiles + + # provider bandwidth limits (ingress, egress) + aws_instance_throughput_limit: Tuple[float, float] = (10, 5) + gcp_instance_throughput_limit: Tuple[float, float] = (16, 7) # limited to 12.5 gbps due to CPU limit + azure_instance_throughput_limit: Tuple[float, float] = (16, 16) # limited to 12.5 gbps due to CPU limit + + # benchmarked_throughput_connections is the number of connections that the iperf3 throughput grid was run at, + # we assume throughput is linear up to this connection limit + benchmarked_throughput_connections = 64 + cost_per_instance_hr = 0.54 # based on m5.8xlarge spot + instance_cost_multiplier = 1.0 + # instance_provision_time_s = 0.0 + + def to_summary_dict(self): + """Simple summary of the problem""" + return { + "src": self.src, + "dsts": self.dsts, + "gbyte_to_transfer": self.gbyte_to_transfer, + "instance_limit": self.instance_limit, + "num_partitions": self.num_partitions, + "required_time_budget": self.required_time_budget, + "aws_instance_throughput_limit": self.aws_instance_throughput_limit, + "gcp_instance_throughput_limit": self.gcp_instance_throughput_limit, + "azure_instance_throughput_limit": self.azure_instance_throughput_limit, + "benchmarked_throughput_connections": self.benchmarked_throughput_connections, + "cost_per_instance_hr": self.cost_per_instance_hr, + "instance_cost_multiplier": self.instance_cost_multiplier + # "instance_provision_time_s": self.instance_provision_time_s, + } + + +@dataclass +class BroadcastSolution: + problem: BroadcastProblem + is_feasible: bool + extra_data: Optional[Dict] = None + + var_edges: Optional[List] = None # need to fix this, just for testing + var_nodes: Optional[List] = None # need to fix this, just for testing + + # solution variables + var_edge_partitions: Optional[np.ndarray] = None # each edge carries each partition or not + var_node_transfer_partitions: Optional[np.ndarray] = None # whether node transfers partition + var_instances_per_region: Optional[np.ndarray] = None # number of VMs per region + var_flow: Optional[np.ndarray] = None # enforce flow conservation, just used for checking + + # solution values + cost_egress: Optional[float] = None + cost_instance: Optional[float] = None + cost_total: Optional[float] = None + transfer_runtime_s: Optional[float] = None # NOTE: might not be able to calculate here + throughput_achieved_gbits: Optional[List[float]] = None # NOTE: might not be able to calculate here + + def to_summary_dict(self): + """Print simple summary of solution.""" + return { + "is_feasible": self.is_feasible, + "solution": { + "cost_egress": self.cost_egress, + "cost_instance": self.cost_instance, + "cost_total": self.cost_total, + "time_budget": self.problem.required_time_budget, + }, + } + + @dataclass class ThroughputProblem: src: str diff --git a/skyplane/planner/solver_ilp.py b/skyplane/planner/solver_ilp.py index e81dbac98..b159572b5 100644 --- a/skyplane/planner/solver_ilp.py +++ b/skyplane/planner/solver_ilp.py @@ -1,18 +1,23 @@ -import cvxpy as cp # type: ignore - from skyplane.planner.solver import ThroughputSolver, ThroughputProblem, GBIT_PER_GBYTE, ThroughputSolution class ThroughputSolverILP(ThroughputSolver): @staticmethod def choose_solver(): + import cvxpy as cp # type: ignore + installed = cp.installed_solvers() order = ["GUROBI", "CBC", "GLPK_MI"] for package in order: if package in installed: return getattr(cp, package) - def solve_min_cost(self, p: ThroughputProblem, solver=cp.GLPK, solver_verbose=False, save_lp_path=None) -> ThroughputSolution: + def solve_min_cost(self, p: ThroughputProblem, solver=None, solver_verbose=False, save_lp_path=None) -> ThroughputSolution: + import cvxpy as cp # type: ignore + + if solver == None: + solver = cp.GLPK # default + regions = self.get_regions() sources = [regions.index(p.src)] sinks = [regions.index(p.dst)] diff --git a/skyplane/planner/topology.py b/skyplane/planner/topology.py index fd5a1898f..e485ce3d3 100644 --- a/skyplane/planner/topology.py +++ b/skyplane/planner/topology.py @@ -63,9 +63,10 @@ class TopologyPlan: The TopologyPlan constains a list of gateway programs corresponding to each gateway in the dataplane. """ - def __init__(self, src_region_tag: str, dest_region_tags: List[str], cost_per_gb: float = 0.0): + def __init__(self, src_region_tag: str, dest_region_tags: List[str], overlay_region_tags: List[str] = [], cost_per_gb: float = 0.0): self.src_region_tag = src_region_tag self.dest_region_tags = dest_region_tags + self.overlay_region_tags = overlay_region_tags self.gateways = {} self.cost_per_gb = cost_per_gb