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
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
"""9-5: Map fail-fast (tolerated-failure-count=0) stops after first failure."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, MapConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def map_fn(_ctx: DurableContext, item: str, _index: int, _items: Any) -> str:
if item == "fail":
raise RuntimeError("item failed")
return item


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
result = context.map(
["ok", "fail", "never"],
map_fn,
name="failfast",
config=MapConfig(
max_concurrency=1,
completion_config=CompletionConfig(tolerated_failure_count=0),
),
)
return {
"completionReason": result.completion_reason.value,
"status": result.status.value,
"successCount": result.success_count,
"failureCount": result.failure_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
"""9-7: Map min-successful early completion."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, MapConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def map_fn(_ctx: DurableContext, item: str, _index: int, _items: Any) -> str:
return item


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
result = context.map(
["s0", "s1", "s2", "s3"],
map_fn,
name="min-successful",
config=MapConfig(
max_concurrency=1,
completion_config=CompletionConfig(min_successful=2),
),
)
return {
"completionReason": result.completion_reason.value,
"successCount": result.success_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
"""9-6: Map throw-if-error propagates an item failure to the execution."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, MapConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def map_fn(_ctx: DurableContext, item: str, _index: int, _items: Any) -> str:
if item == "fail":
raise RuntimeError("item failed")
return item


@durable_execution
def handler(event: Any, context: DurableContext) -> list:
result = context.map(
["fail", "never"],
map_fn,
name="throwing",
config=MapConfig(
max_concurrency=1,
completion_config=CompletionConfig(tolerated_failure_count=0),
),
)
# Rethrow the first item failure; uncaught -> the execution fails.
result.throw_if_error()
return result.get_results()
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
"""9-9: Map tolerated-failure-count exceeded (stops early)."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, MapConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def map_fn(_ctx: DurableContext, item: str, _index: int, _items: Any) -> str:
if item != "never":
raise RuntimeError("item failed")
return item


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
result = context.map(
["f0", "f1", "never"],
map_fn,
name="tolerated-exceeded",
config=MapConfig(
max_concurrency=1,
completion_config=CompletionConfig(tolerated_failure_count=1),
),
)
return {
"completionReason": result.completion_reason.value,
"successCount": result.success_count,
"failureCount": result.failure_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
"""9-10: Map tolerated-failure-percentage exceeded (stops early)."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, MapConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def map_fn(_ctx: DurableContext, item: str, _index: int, _items: Any) -> str:
if item != "never":
raise RuntimeError("item failed")
return item


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
result = context.map(
["f0", "f1", "never", "never"],
map_fn,
name="tolerated-pct",
config=MapConfig(
max_concurrency=1,
completion_config=CompletionConfig(tolerated_failure_percentage=25),
),
)
return {
"completionReason": result.completion_reason.value,
"successCount": result.success_count,
"failureCount": result.failure_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
"""8-19: Parallel with invalid max-concurrency raises a validation error."""

from typing import Any

from aws_durable_execution_sdk_python.config import ParallelConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def branch_a(_ctx: DurableContext) -> str:
return "a"


def branch_b(_ctx: DurableContext) -> str:
return "b"


@durable_execution
def handler(event: Any, context: DurableContext) -> list:
# max_concurrency=0 is invalid; the SDK should reject it before any branch runs.
result = context.parallel(
[branch_a, branch_b],
name="bad-concurrency",
config=ParallelConfig(max_concurrency=0),
)
return result.get_results()
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
"""8-18: Parallel with combined completion config (min-successful + tolerated-failure-count)."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, ParallelConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def branch_fail(_ctx: DurableContext) -> str:
raise RuntimeError("branch failed")


def ok2(_ctx: DurableContext) -> str:
return "ok2"


def ok3(_ctx: DurableContext) -> str:
return "ok3"


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
result = context.parallel(
[branch_fail, branch_fail, ok2, ok3],
name="combined",
config=ParallelConfig(
max_concurrency=1,
completion_config=CompletionConfig(
min_successful=3, tolerated_failure_count=1
),
),
)
# totalCount = started branches (succeeded + failed); matches Java's succeeded()+failed().
return {
"completionReason": result.completion_reason.value,
"successCount": result.success_count,
"failureCount": result.failure_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
"""8-6: Parallel fail-fast (tolerated-failure-count=0) stops after first failure."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, ParallelConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def branch_ok(_ctx: DurableContext) -> str:
return "ok"


def branch_fail(_ctx: DurableContext) -> str:
raise RuntimeError("branch failed")


def branch_never(_ctx: DurableContext) -> str:
return "never"


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
# Fail-fast: tolerated_failure_count=0 stops on the first failure (portable across SDKs).
result = context.parallel(
[branch_ok, branch_fail, branch_never],
name="failfast",
config=ParallelConfig(
max_concurrency=1,
completion_config=CompletionConfig(tolerated_failure_count=0),
),
)
# totalCount = started branches (succeeded + failed); early-stopped
# branches are not counted, matching Java's succeeded()+failed().
return {
"completionReason": result.completion_reason.value,
"status": result.status.value,
"successCount": result.success_count,
"failureCount": result.failure_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
"""8-17: Parallel min-successful not reached (all branches run)."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, ParallelConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def ok0(_ctx: DurableContext) -> str:
return "ok0"


def branch_fail(_ctx: DurableContext) -> str:
raise RuntimeError("branch failed")


def ok2(_ctx: DurableContext) -> str:
return "ok2"


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
result = context.parallel(
[ok0, branch_fail, ok2],
name="min-not-reached",
config=ParallelConfig(
max_concurrency=1,
completion_config=CompletionConfig(min_successful=3),
),
)
# totalCount = started branches (succeeded + failed); matches Java's succeeded()+failed().
return {
"completionReason": result.completion_reason.value,
"status": result.status.value,
"successCount": result.success_count,
"failureCount": result.failure_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
"""8-8: Parallel min-successful early completion."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, ParallelConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def branch0(_ctx: DurableContext) -> str:
return "s0"


def branch1(_ctx: DurableContext) -> str:
return "s1"


def branch2(_ctx: DurableContext) -> str:
return "s2"


def branch3(_ctx: DurableContext) -> str:
return "s3"


@durable_execution
def handler(event: Any, context: DurableContext) -> dict:
result = context.parallel(
[branch0, branch1, branch2, branch3],
name="min-successful",
config=ParallelConfig(
max_concurrency=1,
completion_config=CompletionConfig(min_successful=2),
),
)
# totalCount = started branches (succeeded + failed); early-stopped
# branches are not counted, matching Java's succeeded()+failed().
return {
"completionReason": result.completion_reason.value,
"successCount": result.success_count,
"totalCount": result.total_count,
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
"""8-7: Parallel throw-if-error propagates a branch failure to the execution (fail-fast config)."""

from typing import Any

from aws_durable_execution_sdk_python.config import CompletionConfig, ParallelConfig
from aws_durable_execution_sdk_python.context import DurableContext
from aws_durable_execution_sdk_python.execution import durable_execution


def branch_fail(_ctx: DurableContext) -> str:
raise RuntimeError("branch failed")


def branch_never(_ctx: DurableContext) -> str:
return "never"


@durable_execution
def handler(event: Any, context: DurableContext) -> list:
# Fail-fast: tolerated_failure_count=0 stops on the first failure (portable across SDKs).
result = context.parallel(
[branch_fail, branch_never],
name="throwing",
config=ParallelConfig(
max_concurrency=1,
completion_config=CompletionConfig(tolerated_failure_count=0),
),
)
# Rethrow the first branch failure; uncaught -> the execution fails.
result.throw_if_error()
return result.get_results()
Loading
Loading