Skip to content
Merged
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
135 changes: 41 additions & 94 deletions CHANGELOG.md

Large diffs are not rendered by default.

66 changes: 56 additions & 10 deletions lithops/executors.py
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,15 @@ def _missing_plotting_extra(method_name: str) -> ModuleNotFoundError:
)


def _output_settled(future) -> bool:
"""
Whether nothing is left to read from storage for this call: its result
was downloaded or it produced none, it failed, or it handed back
futures of its own
"""
return future.done or future.futures


def _group_futures_by_job(
futures: List[Any]
) -> List[Tuple[str, str, List[Any], List[Any]]]:
Expand Down Expand Up @@ -449,7 +458,21 @@ def _cleanup_jobs(self, futures, exception=None, force=False):
self.compute_handler.clear(present_jobs)
else:
self.compute_handler.clear(present_jobs, exception=exception)
self.clean(clean_cloudobjects=False, force=force)
self.clean(fs=futures, clean_cloudobjects=False, force=force)

def _release_finished_from_monitor(self, futures):
"""
Drops futures that have already reported back, so the monitor does
not keep listing their prefixes. The thread stays up: the next
map() of this executor adds to it instead of joining a stopped
one and spawning another
"""
finished = [
f for f in futures
if getattr(f, 'ready', False) or f.success or f.done
]
if finished:
self.job_monitor.remove(finished)

def _stop_monitor_if_idle(self, extra_fs=None):
"""
Expand Down Expand Up @@ -764,17 +787,25 @@ def wait(
futures_from_executor_wait=not fs,
)

self._stop_monitor_if_idle(futures)
self._release_finished_from_monitor(futures)
if do_clean and return_when == ALL_COMPLETED:
self._cleanup_jobs(futures)

except (KeyboardInterrupt, Exception) as e:
self.invoker.stop(wait=True)
if isinstance(e, KeyboardInterrupt):
self.invoker.stop(wait=True)
else:
# Only the jobs waited on end here. Another job of this
# executor may still have calls queued for a free worker
self.invoker.discard_pending(
{f.job_key for f in futures},
{f.job_id for f in futures},
)
self.job_monitor.remove(futures)
for future in futures:
future._set_exception()
self._stop_monitor_if_idle(futures)
if self.data_cleaner:
if do_clean:
self._cleanup_jobs(futures, exception=e, force=True)
raise

Expand Down Expand Up @@ -808,7 +839,7 @@ def get_result(
:return: The result of the future/s
"""
pending_to_read = (
len(fs) if fs
len(self._as_future_list(fs)) if fs
else sum(1 for f in self.futures if not f._read and not f.futures)
)

Expand Down Expand Up @@ -961,11 +992,26 @@ def clean(
})

futures = self._as_future_list(fs or self.futures)
present_jobs = {
create_job_key(f.executor_id, f.job_id)
for f in futures
if (f.executor_id.count('-') == 1 and f.done) or force
}
if force or on_exit:
# On exit nothing will read the leftover results, so a job
# that still has one unread call would otherwise stay forever
present_jobs = {
create_job_key(f.executor_id, f.job_id) for f in futures
}
else:
# A job's data is one prefix, so it goes only once no call of
# the job, including the ones not passed here, still has a
# result to read from it
unread_jobs = {
create_job_key(f.executor_id, f.job_id)
for f in list(self.futures) + list(futures)
if not _output_settled(f)
}
present_jobs = {
create_job_key(f.executor_id, f.job_id)
for f in futures
if f.executor_id.count('-') == 1 and f.done
} - unread_jobs
jobs_to_clean = present_jobs - self.cleaned_jobs

if jobs_to_clean:
Expand Down
10 changes: 9 additions & 1 deletion lithops/future.py
Original file line number Diff line number Diff line change
Expand Up @@ -277,6 +277,14 @@ def _raise_call_exception(self, throw_except):
except Exception:
pass
fn_exc.args = (fn_exc.args[1],)
# tblib rebuilds a SystemExit from its args and leaves code at
# None, which would end the client with status 0 whatever the
# function passed to sys.exit()
if isinstance(fn_exc, SystemExit) and fn_exc.code is None \
and fn_exc.args:
fn_exc.code = (
fn_exc.args[0] if len(fn_exc.args) == 1 else fn_exc.args
)
else:
fn_exctype = Exception
fn_exc = Exception(self._exception['exc_value'])
Expand Down Expand Up @@ -452,7 +460,7 @@ def status(

if 'new_futures' in self._call_status and not self._new_futures:
self._resolve_new_futures()
elif self._call_status['func_result_size'] == 0:
elif self._call_status.get('func_result_size', 0) == 0:
self._produce_output = False

if 'result' in self._call_status:
Expand Down
50 changes: 50 additions & 0 deletions lithops/invokers.py
Original file line number Diff line number Diff line change
Expand Up @@ -309,6 +309,13 @@ def stop(self, wait: bool = False):
"""
pass

def discard_pending(self, job_keys, job_ids=None):
"""
Drops the calls of the given jobs not invoked yet. Only an invoker
that queues calls has any
"""
pass


class BatchInvoker(Invoker):
"""
Expand Down Expand Up @@ -540,6 +547,45 @@ def _invoke_job_remote(self, job):
return
raise Exception('Unable to spawn remote invoker')

def _empty_token_bucket(self):
while True:
try:
self.job_monitor.token_bucket_q.get(block=False)
except queue.Empty:
return

def discard_pending(self, job_keys, job_ids=None):
"""
Drops the calls of the given jobs that are still waiting for a
worker, leaving the queued calls of every other job in place.

The monitor also stops tracking those jobs, so their workers will
not hand tokens back. When no other job is still queued, the count
of running workers is forgotten; otherwise a later map() of the
same executor stays capped by workers that never free
"""
kept = []
while True:
try:
item = self.pending_calls_q.get(block=False)
except queue.Empty:
break
job, _ = item
if job is None or job.job_key not in job_keys:
kept.append(item)
other_jobs = False
for item in kept:
job, _ = item
if job is not None:
other_jobs = True
self.pending_calls_q.put(item)
if other_jobs:
return
self._empty_token_bucket()
self.running_workers = 0
if job_ids and self.job_monitor is not None:
self.job_monitor.close_jobs(job_ids)

def _drain_token_bucket(self):
"""
Takes back the tokens left over by previous jobs, one per worker that
Expand Down Expand Up @@ -595,6 +641,10 @@ def _invoke_job(self, job):
prefix = log_prefix(job.executor_id, job.job_id)

if not self.should_run:
# Tokens a monitor handed back after the stop belong to workers
# this restart no longer counts, and would each invoke one more
# worker than max_workers allows
self._empty_token_bucket()
self.running_workers = 0
self.should_run = True
self._start_async_invokers()
Expand Down
15 changes: 8 additions & 7 deletions lithops/localhost/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,12 +30,10 @@
LOCALHOST_EXECUTION_TIMEOUT = 3600

_WINDOWS_PATH = re.compile(r'^(?:[A-Za-z]:[\\/]|\\\\)')
# Interpreters like python, python3, python3.12, python.exe — not docker tags
# such as python:3.12.
_PYTHON_INTERPRETER = re.compile(
r'^python(\d+(\.\d+)*)?(\.exe)?$',
re.IGNORECASE,
)
# Interpreters like python3.12, python3.13t, python3-intel64 or pythonw.exe.
# Docker tags and repositories such as python:3.12 or registry/python always
# carry a ':' or a '/', which an interpreter name never does.
_PYTHON_INTERPRETER = re.compile(r'^python[^:/\\]*$', re.IGNORECASE)


class LocalhostEnvironment(Enum):
Expand All @@ -53,7 +51,10 @@ def get_environment(runtime_name: str) -> LocalhostEnvironment:
if (
runtime_name.startswith('/')
or _WINDOWS_PATH.match(runtime_name) is not None
or _PYTHON_INTERPRETER.match(basename) is not None
or (
'/' not in runtime_name
and _PYTHON_INTERPRETER.match(basename) is not None
)
):
return LocalhostEnvironment.DEFAULT
return LocalhostEnvironment.CONTAINER
Expand Down
33 changes: 10 additions & 23 deletions lithops/monitoring/backends/aws_sqs/status.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,15 +12,11 @@
# limitations under the License.
#

import logging
from functools import cached_property

from lithops.monitoring.backends.aws_sqs import aws_sqs as sqs_backend
from lithops.monitoring.monitor import is_named_error
from lithops.monitoring.status import MessageCallStatus

logger = logging.getLogger(__name__)


class SqsCallStatus(MessageCallStatus):
"""
Expand All @@ -29,6 +25,7 @@ class SqsCallStatus(MessageCallStatus):
"""

service_name = 'SQS'
MAX_MESSAGE_SIZE = 256 * 1024

def __init__(self, job, internal_storage):
super().__init__(job, internal_storage)
Expand All @@ -51,33 +48,23 @@ def _queue_url(self, name):
"""
The URL of a queue by name, looked up once per call.

The monitor of the executor created the queue before any worker was
invoked, so this normally just resolves it; it is created here only
when it is really not there, which is what a status published to an
executor further up the chain can run into
The monitor that reads the queue created it before any worker was
invoked. One that is not there belongs to a reader that is gone, and
creating it again would leave a queue nobody deletes
"""
url = self._urls.get(name)
if url:
return url
try:
url = self.client.get_queue_url(QueueName=name)['QueueUrl']
except Exception as exc:
if not is_named_error(
exc, 'QueueDoesNotExist', 'NonExistentQueue'
):
raise
logger.debug(f'The SQS queue {name} is not there; creating it')
url = self.client.create_queue(QueueName=name)['QueueUrl']
url = self.client.get_queue_url(QueueName=name)['QueueUrl']
self._urls[name] = url
return url

def close(self) -> None:
self._urls.clear()
super().close()

def _publish(self, payload: str) -> None:
for name in self._targets():
self.client.send_message(
QueueUrl=self._queue_url(name),
MessageBody=payload,
)
def _publish_to(self, target: str, payload: str) -> None:
self.client.send_message(
QueueUrl=self._queue_url(target),
MessageBody=payload,
)
13 changes: 7 additions & 6 deletions lithops/monitoring/backends/azure_queue/status.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,11 @@ class AzureQueueCallStatus(MessageCallStatus):
"""

service_name = 'Azure Queue'
#: The service takes 64 KiB per message, measured once the SDK has
#: encoded it: XML-escaped text by default, base64 when the queue client
#: is configured so, which is 4/3 of the text. Three quarters of the
#: limit fits either way
MAX_MESSAGE_SIZE = 48 * 1024

def __init__(self, job, internal_storage):
super().__init__(job, internal_storage)
Expand All @@ -47,9 +52,6 @@ def service(self):
),
)

def _targets(self):
return [azure_queue_name(name) for name in super()._targets()]

def _queue(self, name):
name = azure_queue_name(name)
client = self._queues.get(name)
Expand All @@ -71,6 +73,5 @@ def close(self) -> None:
self._queues.clear()
super().close()

def _publish(self, payload: str) -> None:
for name in self._targets():
self._queue(name).send_message(payload)
def _publish_to(self, target: str, payload: str) -> None:
self._queue(target).send_message(payload)
Loading
Loading