|
| 1 | +import datetime |
| 2 | +from typing import Callable, Dict, Iterable, List, Literal, Optional, Tuple, Union, TYPE_CHECKING |
| 3 | + |
| 4 | + |
| 5 | +import flyte |
| 6 | +from flyte import Image, Resources, TaskEnvironment |
| 7 | +from flyte._doc import Documentation |
| 8 | +from flyte._task import AsyncFunctionTaskTemplate, P, R |
| 9 | + |
| 10 | +import flytekit |
| 11 | + |
| 12 | +if TYPE_CHECKING: |
| 13 | + from flytekit import Cache, Resources, Secret, ImageSpec, Documentation, PodTemplate |
| 14 | + from flytekit.core.base_task import T, TaskResolverMixin |
| 15 | + from flytekit.core.python_function_task import PythonFunctionTask |
| 16 | + from flytekit.core.task import FuncOut |
| 17 | + from flytekit.deck import DeckField |
| 18 | + from flytekit.extras.accelerators import BaseAccelerator |
| 19 | + |
| 20 | + |
| 21 | +def task_shim( |
| 22 | + _task_function: Optional[Callable[P, "FuncOut"]] = None, |
| 23 | + task_config: Optional["T"] = None, |
| 24 | + cache: Union[bool, "Cache"] = False, |
| 25 | + retries: int = 0, |
| 26 | + interruptible: Optional[bool] = None, |
| 27 | + deprecated: str = "", |
| 28 | + timeout: Union[datetime.timedelta, int] = 0, |
| 29 | + container_image: Optional[Union[str, "ImageSpec"]] = None, |
| 30 | + environment: Optional[Dict[str, str]] = None, |
| 31 | + requests: Optional[Resources] = None, |
| 32 | + limits: Optional[Resources] = None, |
| 33 | + secret_requests: Optional[List["Secret"]] = None, |
| 34 | + docs: Optional["Documentation"] = None, |
| 35 | + disable_deck: Optional[bool] = None, |
| 36 | + enable_deck: Optional[bool] = None, |
| 37 | + pod_template: Optional["PodTemplate"] = None, |
| 38 | + pod_template_name: Optional[str] = None, |
| 39 | + accelerator: Optional["BaseAccelerator"] = None, |
| 40 | + pickle_untyped: bool = False, |
| 41 | + shared_memory: Optional[Union[Literal[True], str]] = None, |
| 42 | + resources: Optional[Resources] = None, |
| 43 | + labels: Optional[dict[str, str]] = None, |
| 44 | + annotations: Optional[dict[str, str]] = None, |
| 45 | + **kwargs, |
| 46 | +) -> Union[AsyncFunctionTaskTemplate, Callable[P, R]]: |
| 47 | + plugin_config = task_config |
| 48 | + pod_template = ( |
| 49 | + flyte.PodTemplate( |
| 50 | + pod_spec=pod_template.pod_spec, |
| 51 | + primary_container_name=pod_template.primary_container_name, |
| 52 | + labels=pod_template.labels, |
| 53 | + annotations=pod_template.annotations, |
| 54 | + ) |
| 55 | + if pod_template |
| 56 | + else None |
| 57 | + ) |
| 58 | + |
| 59 | + if isinstance(container_image, flytekit.ImageSpec): |
| 60 | + image = Image.from_debian_base() |
| 61 | + if container_image.apt_packages: |
| 62 | + image = image.with_apt_packages(*container_image.apt_packages) |
| 63 | + pip_packages = ["flytekit"] |
| 64 | + if container_image.packages: |
| 65 | + pip_packages.extend(container_image.packages) |
| 66 | + image = image.with_pip_packages(*pip_packages) |
| 67 | + elif isinstance(container_image, str): |
| 68 | + image = Image.from_base(container_image).with_pip_packages("flyte") |
| 69 | + else: |
| 70 | + image = Image.from_debian_base().with_pip_packages("flytekit") |
| 71 | + |
| 72 | + docs = Documentation(description=docs.short_description) if docs else None |
| 73 | + |
| 74 | + env = TaskEnvironment( |
| 75 | + name="flytekit", |
| 76 | + resources=Resources(cpu=0.8, memory="800Mi"), |
| 77 | + image=image, |
| 78 | + cache="enabled" if cache else "disable", |
| 79 | + plugin_config=plugin_config, |
| 80 | + ) |
| 81 | + return env.task(retries=retries, pod_template=pod_template_name or pod_template, docs=docs) |
| 82 | + |
| 83 | + |
| 84 | +flytekit.task = task_shim |
0 commit comments