Skip to content

[Bug]: Missing invocation deadlines and unbounded branch cleanup can pin LMI workers #741

Description

@zhongkechen

Expected Behavior

The SDK should observe the Lambda invocation deadline and have a bounded shutdown strategy for invocation-owned work. Once the deadline or a scope's cancellation condition is reached, it should stop scheduling new user work, allow in-flight checkpoint waiters to settle, and reclaim the affected worker without silently abandoning running work.

This is particularly important for Lambda Managed Instances (LMI), where an invocation timeout does not forcibly terminate the function's code. The SDK problems below are related and need a coordinated fix.

Actual Behavior

1. Invocation deadlines are not observed.

In a local reproduction, an already-expired Lambda context still allows a step side effect and checkpoint calls, and the wrapper returns SUCCEEDED against an in-memory service stub. The remaining-time method is never called. This demonstrates the missing SDK deadline handling; it does not establish that the real service accepts expired checkpoint tokens.

With LMI's documented timeout behavior, a blocked or long-running step can continue consuming its runtime worker after the invocation has been reported as timed out. Its side effects can overlap a retry. A later rejected checkpoint cannot interrupt or undo an external request already in progress.

2. Early completion can block indefinitely in invocation cleanup.

For parallel/map, the coordinator calls pool.shutdown(wait=False, cancel_futures=True). Already-running branches continue until they finish or encounter an orphan check at a subsequent SDK boundary.

At invocation exit, ExecutionState.close() joins every registered pool with shutdown(wait=True), with no deadline. Only after the branches finish does it stop checkpointing. The outer handler executor also waits for its workers during context-manager exit.

Consequently, CompletionConfig.first_successful() can select a winner and let the user handler compute its final result, while the invocation wrapper remains stuck joining a losing branch blocked in user I/O. That branch can still perform side effects after the winner was selected. Under LMI, reaching the function timeout will not forcibly terminate this work or release the worker slot.

The existing cleanup order serves a real purpose: a branch waiting for a synchronous checkpoint needs the checkpoint worker alive to settle its wait. Simply stopping the batcher first can deadlock it. Normal completion and replay tests did not leave SDK threads behind; the defect is the absence of deadline-aware cancellation and a bounded fallback for work that does not finish.

3. New SDK defect discovered by the LMI tests: new operations can bypass orphan rejection after parent completion.

This is a new problem discovered by the regression tests added in #743, after the original two reproductions in this issue. It has been reproduced with the repository's local runner and on real Python 3.14 LMI with PerExecutionEnvironmentMaxConcurrency=2.

Reproduction:

  1. Run parallel with CompletionConfig.first_successful(), a fast winner, and a losing branch held in controlled I/O inside a step.
  2. Let the winner complete and the parent parallel context checkpoint SUCCEEDED.
  3. Release the losing branch's I/O. In that branch's finally block, attempt a new child.step(..., name="late-work").
  4. Expected: reject the new operation with OrphanedChildException before its body or checkpoint runs. Actual: the new step body executes, and the service records both its start and success after the ancestor context has completed.

The current orphan guard checks whether the operation's own ID is in _parent_done. Parent completion marks the descendants known at that time, but later operation registration does not make a new ID inherit its already-completed/orphaned ancestor's state. Consequently, the new operation bypasses the existing orphan check. Relevant source: _reject_if_parent_done and child registration / parent-completion handling.

Real-cloud evidence from run 35904960112, concurrency-2 job, for case marker late-operation-858b9b6e47bb (all times UTC, 2026-09-23):

Service history event Time
race-parallel: ContextSucceeded 19:32:23.361
late-work: StepStarted 19:32:24.837
late-work: StepSucceeded 19:32:24.837

The external lifecycle ledger independently records LATE_ATTEMPT, BODY with operation late-work, and LATE_ACCEPTED in the original invocation. The artifact lmi-3.14-c2-35904960112-1 contains the corresponding histories/late-operation-858b9b6e47bb.json and events.json. This establishes admission of new durable work under a completed ancestor, extending the original finding about already-running I/O continuing after early completion.

Regression coverage: test_abandoned_child_rejects_late_durable_operation and the local late-operation regression. From the #743 checkout, run:

hatch run lmi:regressions -k abandoned_child_cannot_start_another_step

The regression asserts rejection before any late body or service operation; it fails on the current SDK. The test PR does not implement the SDK fix.

Steps to Reproduce

The following self-contained tests use the real decorator, step, and parallel APIs with an in-memory checkpoint service. They make no AWS calls and require no OTel or other instrumentation plugin.

From a checkout of the commit below, save the code as test_lmi_core_issue.py, then run it in a Python 3.13 environment containing boto3 and pytest:

PYTHONPATH="$PWD/packages/aws-durable-execution-sdk-python/src" \
python -m pytest -q -s test_lmi_core_issue.py

Observed:

expired deadline: step side effect + checkpoints still executed; result SUCCEEDED
first_successful: handler result ready, wrapper blocked, batcher still alive
2 passed

These tests assert the current problematic behavior, so passing confirms the reproduction. The second test uses a context with 10 seconds remaining to isolate cleanup from the already-expired-context case. Test-controlled events release blocked work in finally; the emergency waits prevent the reproduction from hanging indefinitely. Unbounded waiting in the implementation is established by the source paths above.

Self-contained reproduction
"""Behavioral reproductions: passing asserts the current bugs, not a fix."""

import json


import threading


import time


from concurrent.futures import ThreadPoolExecutor


from unittest.mock import Mock


from aws_durable_execution_sdk_python import durable_execution


from aws_durable_execution_sdk_python.config import CompletionConfig, ParallelConfig


from aws_durable_execution_sdk_python.execution import (
    DurableExecutionInvocationInputWithClient, InitialExecutionState,
)


from aws_durable_execution_sdk_python.lambda_service import (
    CheckpointOutput, CheckpointUpdatedExecutionState, ContextDetails,
    DurableServiceClient, ExecutionDetails, Operation, OperationStatus,
    OperationType, StepDetails,
)


class MemoryClient:
    def __init__(self):
        self.operations = {}
        self.calls = []

    def checkpoint(self, *, updates, **kwargs):
        self.calls.append(updates)
        changed = []
        for update in updates:
            status = {
                "START": OperationStatus.STARTED,
                "SUCCEED": OperationStatus.SUCCEEDED,
                "FAIL": OperationStatus.FAILED,
            }.get(update.action.value, OperationStatus.STARTED)
            op = Operation(
                operation_id=update.operation_id,
                operation_type=update.operation_type,
                status=status,
                parent_id=update.parent_id,
                name=update.name,
                sub_type=update.sub_type,
                step_details=(StepDetails(result=update.payload)
                              if update.operation_type == OperationType.STEP else None),
                context_details=(ContextDetails(result=update.payload)
                                 if update.operation_type == OperationType.CONTEXT else None),
            )
            self.operations[op.operation_id] = op
            changed.append(op)
        return CheckpointOutput(
            checkpoint_token="test-token-next",
            new_execution_state=CheckpointUpdatedExecutionState(operations=changed),
        )


def event(client, key="A", replay=False):
    initial = Operation(
        operation_id=key, operation_type=OperationType.EXECUTION,
        status=OperationStatus.STARTED,
        execution_details=ExecutionDetails(input_payload=json.dumps({"key": key})),
    )
    return DurableExecutionInvocationInputWithClient(
        durable_execution_arn="arn:test:execution/" + key,
        checkpoint_token="test-token",
        initial_execution_state=InitialExecutionState(
            operations=[initial, *(client.operations.values() if replay else [])], next_marker="",
        ),
        service_client=client,
    )


def context(key="A", remaining=0):
    ctx = Mock()
    ctx.aws_request_id = key
    ctx.client_context = None
    ctx.identity = None
    ctx.invoked_function_arn = "test-function"
    ctx.tenant_id = None
    ctx._epoch_deadline_time_in_ms = int(time.time() * 1000) + remaining
    ctx.get_remaining_time_in_millis.return_value = remaining
    return ctx


def test_expired_invocation_still_executes_step_and_checkpoints():
    effects = []
    client = MemoryClient()
    expired = context(remaining=-1000)

    @durable_execution
    def handler(_, ctx):
        return ctx.step(lambda _: effects.append("external write") or "done", name="write")

    result = handler(event(client), expired)
    assert result["Status"] == "SUCCEEDED"
    assert effects == ["external write"]
    assert client.calls
    expired.get_remaining_time_in_millis.assert_not_called()
    print("expired deadline: step side effect + checkpoints still executed; result SUCCEEDED")


def test_first_successful_parallel_blocks_invocation_cleanup():
    slow_started = threading.Event()
    release_slow = threading.Event()
    result_computed = threading.Event()
    state = []
    effects = []
    client = MemoryClient()

    def slow(ctx):
        def io(_):
            slow_started.set()
            assert release_slow.wait(5)
            effects.append("losing branch side effect")
            return "slow"
        return ctx.step(io, name="slow-io")

    def fast(ctx):
        def io(_):
            assert slow_started.wait(5)
            return "fast"
        return ctx.step(io, name="fast-io")

    @durable_execution
    def handler(_, ctx):
        state.append(ctx.state)
        ctx.parallel([slow, fast], name="race", config=ParallelConfig(
            max_concurrency=2, completion_config=CompletionConfig.first_successful(),
        ))
        result_computed.set()
        return "winner selected"

    with ThreadPoolExecutor(max_workers=1) as caller:
        future = caller.submit(handler, event(client), context(remaining=10000))
        try:
            assert result_computed.wait(3), future.exception() if future.done() else "no result"
            assert not future.done()
            assert not state[0]._checkpointing_stopped.is_set()
            print("first_successful: handler result ready, wrapper blocked, batcher still alive")
            print("threads:", [t.name for t in threading.enumerate()])
        finally:
            release_slow.set()
        assert future.result(timeout=3)["Status"] == "SUCCEEDED"
    assert effects == ["losing branch side effect"]

SDK Version

2.0.1; inspected and reproduced against commit 9157022256ce2dcff075308b71ed6054bf8b8569.

Python Version

Python 3.13.15, Linux.

Is this a regression?

Unknown; no previous working version has been established.

Additional Context

LMI's execution-environment documentation states that the environment does not freeze between invocations and that timing out an invocation does not forcibly terminate its code. Unfinished code continues occupying a worker slot.

The official Python runtime uses separate worker processes, each handling one invocation at a time. This issue is about invocation lifetime and owned threads, not shared Python state between concurrent invocations. PerExecutionEnvironmentMaxConcurrency > 1 amplifies resource consumption and possible overlap with retries, but is not necessary to trigger either defect.

Each active invocation adds two SDK workers (handler and checkpoint batcher), and each active map/parallel creates an additional branch pool. SDK branch concurrency is separate from LMI process concurrency. Stuck branches retain those resources, including the live checkpoint batcher waiting behind the join.

Original validation comprised local source-based reproduction plus 33 existing focused tests covering cleanup, checkpoint error propagation, completion events, and invocation lifecycle. Subsequent tests in #743 added real LMI validation. The newly discovered late-operation defect above is supported by both a local reproduction and the linked Python 3.14 concurrency-2 cloud execution; it was not known when this issue was initially filed.

Possible Fix

  1. Create an invocation-scoped deadline and cancellation controller. Initialize a monotonic deadline from get_remaining_time_in_millis() and reserve a cleanup margin. Pass it explicitly through ExecutionState, user execution, checkpoint waits, and concurrency coordinators. Check it before starting new user operations or branches and during SDK-owned blocking waits. Invocation expiry should propagate as an interruption/retry condition, not become a permanent user step failure.

  2. Give abandoned scopes a separate cancellation signal. When a parallel/map completion policy is satisfied, cancel queued branches and signal already-running descendants cooperatively. Link nested scope cancellation to the invocation controller, but keep it scoped so selecting a parallel winner does not cancel the rest of the handler. Expose a suitable cancellation/deadline surface to user I/O and require bounded downstream timeouts; internal orphan checks alone cannot stop a running step body. The cancellation design must also reject newly registered descendants of an already completed or orphaned scope before user code or checkpoint submission starts.

  3. Bound checkpoint I/O and settle waiters before shutting workers down. Account for connection, read, and retry time in the remaining budget. The cached boto client should not be mutated unsafely to implement per-call timeouts; evaluate a supported transport/client strategy or conservative fixed limits combined with cancellation. Keep the batcher alive while required in-flight checkpoints settle, or explicitly release all pending waiters with an invocation-interruption error under the state locks. Preserve operation identity and persisted outcomes. Do not synthesize SUCCEEDED/PENDING for unacknowledged critical checkpoints.

  4. Make cleanup a bounded protocol with an explicit fallback. Track branch futures/completion so cleanup can wait only for the available grace period, including pools registered by nested work during shutdown. Audit the outer ThreadPoolExecutor context-manager exit as well; otherwise a timeout-aware future.result() can still be followed by an unbounded executor join. Python cannot safely terminate an arbitrary running thread. For non-cooperative work, evaluate a supervised invocation subprocess or runtime-supported retirement/replacement of the affected Python worker, scoped to that worker and coordinated with the runtime. Merely switching to shutdown(wait=False) or Future.cancel() does not stop running work and would allow it to outlive the invocation. This fallback needs design/runtime validation before choosing an implementation.

Suggested regression coverage:

  • An already-expired invocation starts no new user work; expiry during an invocation stops further scheduling and propagates as an invocation-level interruption.
  • A fast winner plus a blocked losing branch cannot pin a cooperative invocation indefinitely; non-cooperative work follows the explicitly chosen bounded fallback.
  • Cancellation races with checkpoint responses settle every waiter, without deadlocks or accepting late checkpoints after shutdown.
  • Nested map/parallel registration during cleanup is covered, and scope cancellation leaves unrelated branches/invocations unaffected.
  • New test-discovered regression: after parent completion, a residual branch attempts a new step from finally. Both previously known and newly created operation IDs must be rejected under the completed/orphaned ancestor, before the late body, START hook, or service checkpoint can execute.
  • Normal return, failure, suspension, and replay preserve checkpointed results and errors and leave no invocation-owned workers running. Completed side effects are not repeated on replay; interrupted side effects remain subject to documented retry/idempotency semantics.
  • An eventual real LMI test verifies worker capacity recovery after deadline expiry while other worker processes continue handling requests.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions