Skip to content
Open
Show file tree
Hide file tree
Changes from all 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
21 changes: 14 additions & 7 deletions src/taskgraph/optimize/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -301,6 +301,18 @@ def replace_tasks(
target_task_graph.graph.links_and_reverse_links_dict()
)

# Many tasks share dependents (e.g. docker images or toolchains), so
# resolve the deadline of each dependent only once.
now = datetime.datetime.now(datetime.timezone.utc)
deadlines = {}

def get_deadline(label):
if label not in deadlines:
deadlines[label] = resolve_timestamps(
now, target_task_graph.tasks[label].task["deadline"]
)
return deadlines[label]

for label in target_task_graph.graph.visit_postorder():
logger.debug(f"replace_tasks: {label}")
# if we're not allowed to optimize, that's easy..
Expand Down Expand Up @@ -331,14 +343,9 @@ def replace_tasks(
opt_by, opt, arg = optimizations(label)

# compute latest deadline of dependents (if any)
dependents = [target_task_graph.tasks[l] for l in dependents_of[label]]
deadline = None
if dependents:
now = datetime.datetime.now(datetime.timezone.utc)
deadline = max(
resolve_timestamps(now, task.task["deadline"])
for task in dependents # type: ignore
)
if dependents_of[label]:
deadline = max(get_deadline(l) for l in dependents_of[label])

if isinstance(opt, IndexSearch):
arg = arg, index_to_taskid, taskid_to_status
Expand Down
12 changes: 10 additions & 2 deletions src/taskgraph/optimize/strategies.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import functools
import logging
from datetime import datetime

Expand All @@ -10,6 +11,13 @@
logger = logging.getLogger("optimization")


@functools.cache
def _parse_time(timestamp, fmt):
# Many tasks share the same deadline or replacement task, so avoid parsing
# the same timestamps over and over.
return datetime.strptime(timestamp, fmt)


@register_strategy("index-search")
class IndexSearch(OptimizationStrategy):
# A task with no dependencies remaining after optimization will be replaced
Expand Down Expand Up @@ -56,10 +64,10 @@ def should_replace_task(self, task, params, deadline, arg):
)
continue

if deadline and datetime.strptime(
if deadline and _parse_time(
status["expires"], # type: ignore
self.fmt,
) < datetime.strptime(deadline, self.fmt):
) < _parse_time(deadline, self.fmt):
logger.debug(
f"not replacing {task.label} with {task_id} because it expires before {deadline}"
)
Expand Down
Loading