"""Lambda handler shim for AlphaAgent code-interpreter sessions.

This module is the ENTRYPOINT for every environment Lambda function.
It receives a JSON event from the LambdaBackend and executes Python code
or shell commands inside the Lambda execution environment.

Every ``exec_*`` action runs in a FRESH working directory under /tmp (always
writable in Lambda; ``mkdtemp``, 0700, removed when the exec ends - see
"Fresh working directory per exec", below) that mirrors the session's
project prefix
(``users/{uid}/projects/conversations/{slug}/``): the requested project
directories are pulled down first (``sync_dirs`` — ETag-diffed against a
per-project cache, so a warm Lambda re-downloads only what changed), the code
runs with the project root as cwd, then only the files the run created or
changed (mtime/size snapshot) are uploaded and reported as ``uploaded`` so
the caller can emit one ``file_added`` per file.

Connector credentials are fetched from Secrets Manager at invocation time
and handed to the child process as environment variables - THIS invocation's
only: the child's environment is the image's own plus this invocation's
variables (`_child_env`), and nothing per-invocation is ever written to the
handler's ``os.environ``, because one warm container serves every agent
pinned to the environment (see the "Per-invocation environment" section).
Shell commands run under bash when the image has it: the variables
are named ``CONN_{connector_id}_{KEY}`` and a connector id is a UUID, so the
names carry hyphens, and Debian's ``/bin/sh`` (dash) drops such entries from
its environment at startup - see `_shell_argv`.

Dependency contract: stdlib + boto3 ONLY — this file ships inside every env
image (see environments/alphaagent-python/Dockerfile) on release.

Runtime delivery: since Studio 1.0.150 the platform does not rely on the
copy baked into the image. agent-management points every
environment function's ``ImageConfig.EntryPoint`` at an inline bootstrap that
downloads this file and ``aa_env.py`` from a release-staged S3 prefix
(``AA_HANDLER_URI``) into ``/tmp/aa_handler``, imports the handler from there
and only then starts ``awslambdaric``; when the download fails it runs the
image's own copy and records why in ``AA_HANDLER_DELIVERY``. This file
reports which copy is live (`_handler_delivery`, in ``probe_capabilities``)
and, because ``/tmp`` is writable by the sandbox uid while ``/opt`` is not,
restores ``aa_env.py`` next to itself before every exec
(`_ensure_overlay_intact`). The bootstrap never edits this file: it is the
same bytes wherever it runs from.

Private HOME per invoking scope: the image's HOME is unusable for the
runtime uid on the live env Lambda, and the fallback used to be ONE
``/tmp/home`` shared by every agent the warm container serves - a place one
agent's script could plant ``~/.local/lib/python3.X/site-packages/
usercustomize.py`` (imported by ``site`` at every interpreter start), a
``.pth`` file, a shadow module, ``~/.config/pip/pip.conf`` or rc files for
the NEXT agent's child to execute with that agent's credentials. Every child
now gets ``/tmp/home-<sha256(scope)[:16]>`` where the scope is the invoking
agent, else the owning user, else the exec itself (`_home_scope`,
`_scope_home`): the same agent's later warm turns find what it installed
(``pip install --user`` keeps working), a different agent never does.

Fresh working directory per exec: the mirror used to live at a
PREDICTABLE path - ``/tmp/users_<uid>_projects_conversations_<slug>`` -
that persisted while the container was warm, was only ever mirrored DOWN (a
file gone from S3 stayed on disk) and whose ``code/`` sits on the child's
``PYTHONPATH`` ahead of the standard library and site-packages: a hostile
child of user A could plant ``code/pandas.py`` into user B's path and B's
next turn on the same container executed it with B's credentials. Now each
exec gets ``/tmp/.aa-work-<random>`` (`_get_session_workdir`), the child runs
there and the directory is removed after the up-mirror (`_release_workdir`).
The warm ETag diff is kept through a per-project CACHE,
``/tmp/ws-<sha256(prefix)[:16]>`` (`_workspace_cache`), whose manifest is
this process's memory (`_WORKSPACE_INDEX`) - the one place a child cannot
write - and every cache hit is re-verified by sha256 before the working copy
is hard-linked from it (`_sync_dirs_down`, `_materialize_workdir`); anything
in the cache the index does not vouch for is removed (`_scrub_cache_dir`).
Bound (the same one the private HOME has): every child is one uid on one
/tmp, so a resident process a child leaves behind is not stopped by any of
this.

Credential isolation: the runtime hands this process its IAM role
credentials as environment variables and the user's code runs as a same-uid
child, so ``/proc/<ppid>/environ`` used to hand a script a role that read
every user's workspace and every agent's connector secret. The handler now
re-executes itself at import with a clean environment (credentials over a
pipe, into memory) FROM AN EXEC-ONLY COPY of the interpreter, so the serving
process is non-dumpable from its first instruction (a child gets
``PermissionError`` on ``/proc/<ppid>/environ``, ``mem`` and ``maps``;
``prctl(PR_SET_DUMPABLE)`` is refused by Lambda's seccomp, so the exec is
what carries that), and does every S3 / Secrets Manager call for the user on
STS session credentials scoped to ``users/<uid>/`` and the agent's own
secret - the "Credential isolation" section below has the four layers and
their bounds.
"""

import glob
import hashlib
import json
import mimetypes
import os
import re
import resource
import shutil
import stat
import subprocess
import sys
import tempfile
import time
import traceback
import errno
from datetime import datetime, timezone
from typing import Any, Dict, List, NamedTuple, Optional, Tuple

import boto3

#: Refused uploads reported back to the caller: entries kept, chars per entry.
_UPLOAD_ERRORS_MAX = 50
_UPLOAD_ERROR_MAX_CHARS = 300
MAX_OUTPUT_BYTES = 1024 * 1024  # 1 MB cap on stdout/stderr returned to caller

# The project mirror: a fresh working directory per exec, fed from a
# per-project cache whose manifest lives in this process's memory.
#: Working directory name prefix under `TMP_ROOT`: ``mkdtemp``, removed per exec.
WORKDIR_PREFIX = ".aa-work-"
#: Per-project mirror cache under `TMP_ROOT`: ``ws-<sha256(prefix)[:16]>``, 0700.
WORKSPACE_CACHE_PREFIX = "ws-"
#: ``{cache_dir: {rel_path: (etag, sha256)}}`` - what `_sync_dirs_down` has
#: verified in each cache. Handler memory, never a file: a sandbox child can
#: write anything under /tmp, but not here. Same lifetime as /tmp (the warm
#: container), so nothing is lost by keeping it in memory.
_WORKSPACE_INDEX: Dict[str, Dict[str, Tuple[str, str]]] = {}
#: Project directories a coder run always mirrors (see coder_worker.run_in_env;
#: ``brand/`` since v1.0.120 so the style resolver sees the project's brand).
DEFAULT_SYNC_DIRS = ("code", "nodes", "brand")
#: Size guards for the mirror-down (env fallbacks; the event may override).
_DEFAULT_SYNC_MAX_BYTES = 500 * 1024 * 1024
_DEFAULT_SYNC_SKIP_OVER_BYTES = 100 * 1024 * 1024
#: Where an inline ``code`` script lands inside the mirrored project, so the
#: one legacy path that still accepts source text never litters the root.
INLINE_SCRIPT_DIR = "code/scratch"
#: Where an ``ephemeral`` inline script lands: a HIDDEN dir, so a service-side
#: helper run (the preview renderer) never uploads its program into the
#: user's workspace.
EPHEMERAL_SCRIPT_DIR = ".aa-scratch"
#: S3 key prefixes that mark a prefetch entry as a FULL key rather than a
#: path relative to the session prefix.
_FULL_KEY_PREFIXES = ("users/", "workspaces/", "conversations/")
#: Version stamp the document toolkit image writes (reported as
#: ``toolkit_version`` by the capability probe); absent on an image built
#: before it existed.
TOOLKIT_LOCK_PATH = "/opt/alphaagent/toolkit.lock"
#: Read to detect QEMU user-mode emulation (an amd64 image on an arm64 host):
#: x86_64 with no x86 ``flags:`` line. Chromium cannot run there, so the
#: lock's ``needs_browser`` checks are reported skipped, never pass/fail.
CPUINFO_PATH = "/proc/cpuinfo"
#: Per-command ceiling for the capability probe (seconds). soffice is the
#: slow one on a cold container; everything else answers in well under 1s.
#: The lock may override it (``check_timeout_s``).
PROBE_TIMEOUT_S = 15
#: Ceiling for the lock's ``needs_browser`` checks: a cold Chromium launch in
#: Lambda (2 GB, no /dev/shm, single-process) takes well over 15 s the first
#: time. The lock may override it (``browser_check_timeout_s``).
BROWSER_CHECK_TIMEOUT_S = 60
#: The toolkit.lock version this handler was cut with (environments/toolkit/
#: toolkit.lock). An installed lock BELOW this is ``toolkit_outdated`` — e.g. a
#: customer image built from the 1.0.0 lock has Chrome launch flags that crash
#: in Lambda and fill /tmp with core dumps (toolkit 1.1.0).
HANDLER_TOOLKIT_VERSION = "1.1.0"
#: Contract version of THIS handler: bumped when the request/response
#: contract between this file and the platform changes, so the platform can
#: tell an outdated copy from a probe. Every handler that predates the field
#: is implicitly "1" (it reports nothing). "2": runtime delivery, the
#: per-invocation child environment, budget/reserve reporting. Additive
#: under "2" (1.0.154): ``credential_isolation`` in the probe.
HANDLER_CONTRACT_VERSION = "2"
#: Set on the FUNCTION by agent-management: the release-staged S3 prefix the
#: bootstrap downloads the handler pair from (``s3://<bucket>/.../<release>``).
HANDLER_URI_ENV = "AA_HANDLER_URI"
#: Set in the PROCESS by the bootstrap before awslambdaric starts: ``overlay``
#: when the downloaded pair is what runs, ``baked:<reason>`` when it fell back
#: to the image's copy. Absent when the function's entry point is the legacy
#: ``python -m awslambdaric`` (no bootstrap ran).
HANDLER_DELIVERY_ENV = "AA_HANDLER_DELIVERY"
#: Lambda's only writable filesystem. Core dumps land here
#: (``core_pattern=/tmp/core.%e.%p`` on the env Lambda), so they are swept
#: before every action. When the image's HOME is unusable for the runtime uid
#: (the live env Lambda runs as uid 993 while /home/agent is drwx------ uid
#: 1000: unreadable, so choreographer/kaleido, fontconfig and Chromium all
#: fail there) a child's HOME/XDG_CACHE_HOME point at a directory under here
#: that is private to the invoking scope - `_scope_home`. The old single
#: shared fallback directory is never used again.
TMP_ROOT = "/tmp"
#: ``/tmp/<HOME_DIR_PREFIX><sha256(scope)[:16]>`` is a scope's HOME.
HOME_DIR_PREFIX = "home-"
#: The scope of the handler's OWN children - the capability probes and lock
#: checks (`_probe_run`, `probe_capabilities`). Never an agent's scope: agent
#: code never runs with this HOME and these children carry no credentials.
TOOLING_HOME_SCOPE = "handler:tooling"
#: Scope prefix of a HOME that belongs to one exec only (no agent, no owner
#: in the event): removed with the exec's helper directory.
EXEC_HOME_SCOPE_PREFIX = "exec:"
#: Never handed to a child: files an earlier invocation could plant and a
#: later child's interpreter or shell would source at startup. ``BASH_ENV``
#: is the ONE file a non-interactive ``bash -c`` reads; ``PYTHONSTARTUP``
#: the interactive interpreter's. NOT ``ENV``: POSIX sh reads ``$ENV`` only
#: when interactive (never the case here), and it is a plausible name for a
#: customer's own variable on the function.
_CHILD_STARTUP_HOOK_KEYS = ("PYTHONSTARTUP", "BASH_ENV")
CORE_DUMP_GLOB = "core.*"

# Checkpoint/continuation contract (see aa_env.py). When a script checkpoints and
# asks to be re-invoked it exits with this code; the coder harness re-invokes a
# fresh Lambda with resume=true. Kept in one place so the three parties (this
# handler, aa_env, coder_worker) agree.
EXIT_CONTINUE = 75
#: coreutils ``timeout`` exit status when the command outlived the wall. The
#: handler also normalises the SIGKILL that ``timeout -k`` delivers to a
#: SIGTERM-trapping script (-9 / 137 seen at the wall) to this value.
EXIT_TIMEOUT = 124
# Seconds reserved AFTER the script exits for the workdir→S3 upload, and seconds
# of grace before the hard kill in which a cooperative script should checkpoint.
# Both are CAPS, not flat subtractions: `_compute_budget` scales them to the
# function's own remaining time (reserve = min(cap, max(5, remaining // 4))), so
# a 60 s environment gives code ~45 s rather than the 5 s floor a flat 60 s
# reserve would leave. Overridable via env (set on the function by the
# backend) or per event (``upload_reserve_s`` / ``yield_grace_s``).
_DEFAULT_UPLOAD_RESERVE_S = 60
_DEFAULT_YIELD_GRACE_S = 30
#: Floors for the scaled reserve / grace and for the script's own budget.
_MIN_UPLOAD_RESERVE_S = 5
_MIN_YIELD_GRACE_S = 5
_MIN_HARD_TIMEOUT_S = 5
#: Upper bound on ``timeout -k``: a script that traps SIGTERM is SIGKILLed this
#: many seconds after the wall (paid for out of the upload reserve).
_MAX_KILL_AFTER_S = 5

# VENDORED, not imported (stdlib + boto3 only): mirrors the service's own
# run-data tagging definition (``RUN_DATA_TAGGING`` / ``RUN_DATA_KEY_PATTERN``
# / ``TEXT_CONTENT_TYPES`` / ``put_args``), which carries the key and value
# agent-runtime puts on its own writes (``workflow_inputs.RUN_DATA_TAG_KEY`` /
# ``RUN_DATA_TAG_VALUE``). The workspaces bucket's lifecycle rule filters on
# this object tag, so every file this handler uploads under a run root
# (``users/{uid}/projects/workflows/{slug}/runs/{execution_id}/`` or a
# conversation's ``runs/{execution_id}/`` scratch) - the deliverables under
# ``outputs/``, ``outputs/result.json``, ``code/scratch/*`` - must carry it
# or it never ages out. Files outside a run root stay untagged. A test pins
# the three copies equal.
RUN_DATA_TAGGING = "aa-data-class=run-data"
_RUN_DATA_KEY_RE = re.compile(r"^users/[^/]+/projects/(?:workflows|conversations)/[^/]+/runs/[^/]+/")
#: Extension -> ``ContentType`` for the text types a run's files are made of;
#: Python's ``mimetypes`` has no ``.md``, so without this a Markdown
#: deliverable reads back as ``application/octet-stream``.
_TEXT_CONTENT_TYPES: Dict[str, str] = {
    ".md": "text/markdown",
    ".markdown": "text/markdown",
    ".txt": "text/plain",
    ".log": "text/plain",
    ".csv": "text/csv",
    ".tsv": "text/tab-separated-values",
    ".json": "application/json",
    ".jsonl": "application/x-ndjson",
    ".ndjson": "application/x-ndjson",
    ".yaml": "text/yaml",
    ".yml": "text/yaml",
    ".toml": "text/plain",
    ".ini": "text/plain",
    ".cfg": "text/plain",
    ".py": "text/x-python",
    ".sh": "text/x-shellscript",
    ".sql": "application/sql",
    ".xml": "application/xml",
    ".html": "text/html",
    ".htm": "text/html",
    ".css": "text/css",
    ".js": "text/javascript",
    ".svg": "image/svg+xml",
}


def _s3_content_type(name: str) -> Optional[str]:
    """``ContentType`` from the extension of ``name``, or ``None`` (omit it)."""
    base = str(name or "").rsplit("/", 1)[-1]
    dot = base.rfind(".")
    if dot <= 0:
        return None
    known = _TEXT_CONTENT_TYPES.get(base[dot:].lower())
    if known:
        return known
    guessed, _ = mimetypes.guess_type(base)
    return guessed or None


def _s3_put_args(key: str, content_type: Optional[str] = None) -> Dict[str, str]:
    """Extra ``put_object`` kwargs for ``key``: ``Tagging`` when it is run
    data, ``ContentType`` = the explicit one, else the extension's when known.
    Never contains a ``None`` value (botocore rejects ``Tagging=None``)."""
    extra: Dict[str, str] = {}
    if key and _RUN_DATA_KEY_RE.match(str(key)):
        extra["Tagging"] = RUN_DATA_TAGGING
    ctype = content_type or _s3_content_type(key)
    if ctype:
        extra["ContentType"] = ctype
    return extra


def _load_connector_credentials(agent_id: str, connector_ids: list) -> Dict[str, str]:
    """Fetch connector credentials from Secrets Manager and return as env
    vars - through this invocation's scoped session (`_secrets_client`),
    whose session policy names this agent's secret and no other."""
    if not connector_ids:
        return {}

    try:
        sm = _secrets_client()
        secret_name = f"alphaagent-connector-creds-{agent_id}"
        print(f"[HANDLER] Fetching connector secret: {secret_name}")
        resp = sm.get_secret_value(SecretId=secret_name)
        creds = json.loads(resp["SecretString"])
        env_vars: Dict[str, str] = {}
        for key, value in creds.items():
            env_vars[key] = str(value) if value is not None else ""
        print(f"[HANDLER] Loaded {len(env_vars)} connector creds")
        return env_vars
    except Exception as exc:
        print(f"[HANDLER] ERROR loading connector creds: {exc}")
        return {}


def _cap_output(data: bytes) -> str:
    """Decode and truncate output to stay within Lambda response limits."""
    text = data.decode("utf-8", errors="replace")
    if len(text) > MAX_OUTPUT_BYTES:
        return text[:MAX_OUTPUT_BYTES] + "\n[OUTPUT TRUNCATED]"
    return text


# ---------------------------------------------------------------------------
# Credential isolation.
#
# The Lambda runtime vends this execution environment's IAM role credentials
# as plain environment variables on the process `handler()` runs in, and the
# user's own script or shell command runs as a CHILD of that process - same
# uid, same /tmp. Two facts made that a tenant-isolation hole rather than an
# untidiness:
#
# * ``/proc/<pid>/environ`` is the environment block a process was STARTED
#   with. ``os.environ.pop`` rewrites the C ``environ`` pointers, not that
#   block, so a child that reads ``/proc/<ppid>/environ`` recovers every
#   variable the runtime set. Stripping the child's own ``env=`` (which this
#   file has done since D.1.a - `_sandboxed_env`) never closed that.
# * The role those credentials belong to (``<project>-<env>-lambda-execution-
#   role``) was granted the WHOLE workspaces bucket and EVERY agent's
#   connector secret, so the recovered credentials reached other users'
#   ``users/<uid>/`` trees and other agents' secrets directly with boto3.
#
# Four layers, each deterministic - no prompt and no model in any of them:
#
# 1. **Re-exec with a clean environment** (`_reexec_with_clean_environment`).
#    At import - the delivery bootstrap imports this module BEFORE it starts
#    ``awslambdaric`` - the handler writes the credential variables into a
#    pipe, drops them from the environment and ``execve``s itself: the same
#    interpreter, the handler directory first on ``sys.path``, then
#    ``awslambdaric`` via ``runpy`` exactly as the bootstrap would have
#    started it. The process that serves invocations was therefore STARTED
#    without credentials: its ``/proc/<pid>/environ`` has none to give. The
#    next image reads them back from the pipe into `_ROLE_CREDENTIALS` -
#    process memory, never a file, never the environment. ``awslambdaric``
#    itself never reads AWS credentials (it reads ``AWS_LAMBDA_RUNTIME_API``
#    and ``_LAMBDA_TELEMETRY_LOG_FD``, and sets ``_X_AMZN_TRACE_ID``), and the
#    runtime vends the role credentials ONCE per execution environment - no
#    per-invocation refresh - so the in-memory copy has exactly the lifetime
#    the variables had. Skipped, and reported as such, when this module is
#    imported by an already-running ``awslambdaric`` (the legacy ``python -m
#    awslambdaric`` entry point, or the bootstrap's baked fallback: its
#    telemetry fd is already consumed and an exec would lose it) - layer 2
#    stands alone there.
# 2. **Non-dumpable process** (`_nondumpable_image`, `_set_non_dumpable`,
#    `_measure_non_dumpable`). For a task whose dumpable flag is not
#    ``SUID_DUMP_USER`` the kernel's ``ptrace_may_access`` refuses a
#    same-uid, non-root process ``/proc/<pid>/environ``, ``mem``, ``maps``
#    and ``ptrace`` attach (only ``CAP_SYS_PTRACE`` gets through), whatever
#    ``kernel.yama.ptrace_scope`` says - and ``task_dump_owner`` hands those
#    /proc files to root, so plain file permissions refuse them too. Two
#    ways to get there, both taken:
#    * the clean image is started from an EXEC-ONLY COPY of the interpreter
#      (mode 0100, in the overlay directory, pre-flighted as a child first).
#      ``execve`` of a file its caller may not read sets
#      ``BINPRM_FLAGS_ENFORCE_NONDUMP`` (``would_dump``) and the new mm gets
#      ``fs.suid_dumpable`` (2 on Lambda's kernel, 0 elsewhere) instead of
#      ``SUID_DUMP_USER``: non-dumpable from its first instruction, with no
#      syscall the process has to be allowed. ``argv[0]`` stays the real
#      interpreter, so ``sys.executable`` and ``sys.prefix`` are unchanged
#      and every child still runs the installation's python.
#    * ``prctl(PR_SET_DUMPABLE, 0)`` in the serving process, where the kernel
#      lets it. Under the Lambda runtime it does NOT: the sandbox runs under
#      seccomp (``Seccomp: 2``, six filters, ``NoNewPrivs: 1``, kernel 5.10
#      amzn2) that answers EVERY prctl - PR_GET_DUMPABLE included - with
#      EPERM (measured live, inside a sandbox). The first release of this
#      layer was this call alone, so it was never in effect there; the
#      refusal is now recorded (``credential_isolation.prctl``), never
#      swallowed.
#    ``execve`` resets the flag to ``SUID_DUMP_USER`` for a readable binary,
#    which is why the copy must be the very file the clean image is exec'd
#    from. ``PR_SET_PTRACER`` adds nothing: it is a Yama hook, Yama is absent
#    from the Lambda kernel (``ptrace_scope: absent``), so EINVAL at best -
#    and EPERM under that seccomp anyway. Measured, not assumed: the owner of
#    ``/proc/self/status`` (root <=> non-dumpable, for a non-root process) at
#    import and again on every invocation (``credential_isolation.
#    non_dumpable``), plus a REAL child of this process trying
#    ``/proc/<ppid>/environ``, ``maps`` and ``mem``
#    (``child_can_read_parent_environ``, ``child_can_open_parent_mem``).
#    Bound: the baked fallback and the legacy entry point (awslambdaric
#    already running, no re-exec possible) have only the prctl half, which
#    Lambda refuses - reported as ``non_dumpable: false`` on the row and a
#    WARNING on every invocation, never silent.
# 3. **Least privilege per invocation** (`AccessScope`, `_scoped_session`).
#    Every S3 and Secrets Manager call this file makes on the user's behalf
#    - mirror down and up, prefetch, ``fetch_object``, ``read_file`` /
#    ``write_file`` / ``list_files``, the connector secret - uses STS session
#    credentials from assuming ``<project>-<env>-sandbox-scoped-role``
#    (core.yaml; trusted by the execution role alone) WITH a session policy
#    narrowed to this invocation: read on ``users/<uid>/`` plus whatever
#    ``read_prefixes`` the invoking service adds (a governed deployment passes
#    ``users/``: reads deployment-wide, writes the writer's own folder), write
#    on ``users/<uid>/`` only, and the ONE secret
#    ``alphaagent-connector-creds-<agent_id>``. The role's own policy is the
#    old bucket-wide grant; the session policy is what makes it per user,
#    and `_scoped_session` never calls AssumeRole without one. Sessions are
#    cached per (bucket, prefixes, agent) for the container's life and
#    renewed 20 minutes before expiry (an invocation can run for 15). The
#    execution role keeps only what starting the function takes: the overlay
#    read, logs, ENI, ECR, and AssumeRole. `_s3_client` and `_secrets_client`
#    are the ONLY places this file builds those clients, and they refuse to
#    build one outside an invocation (`_CURRENT_SCOPE`).
# 4. **Child hardening** (`_exec_env`). The credential variables
#    (`_LAMBDA_CREDENTIAL_ENV_KEYS`) and the variables an SDK would follow to
#    find credentials (`_CREDENTIAL_SOURCE_ENV_KEYS`) never reach a child, and
#    ``AWS_EC2_METADATA_DISABLED=true`` makes a stray ``boto3`` in the sandbox
#    raise ``NoCredentialsError`` at once instead of waiting on a metadata
#    service Lambda does not have.
#
# When the scoped role cannot be assumed - a deployment whose core stack
# predates it (STS answers AccessDenied), or STS unreachable - the handler
# falls back to the execution role's own credentials for THAT invocation
# only, never cached across invocations, logs it, and reports
# ``credential_scope`` in the probe. On a current core stack that fallback
# reaches nothing in S3 (the grant moved to the scoped role), so a failure
# is loud, never silent.
# ---------------------------------------------------------------------------
_HANDLER_DIR = os.path.dirname(os.path.abspath(__file__))

#: The runtime's own role credentials, as it vends them: plain environment
#: variables. Moved to memory at import; never in a child's environment.
_LAMBDA_CREDENTIAL_ENV_KEYS = (
    "AWS_ACCESS_KEY_ID",
    "AWS_SECRET_ACCESS_KEY",
    "AWS_SESSION_TOKEN",
    "AWS_SECURITY_TOKEN",  # legacy alias some SDKs still read
)
#: Not credentials, but where an SDK goes looking for some - never in a
#: child's environment either. Exact names: a connector's
#: ``CONN_<id>_AWS_PROFILE`` is a different name and survives.
_CREDENTIAL_SOURCE_ENV_KEYS = (
    "AWS_CONTAINER_CREDENTIALS_FULL_URI",
    "AWS_CONTAINER_CREDENTIALS_RELATIVE_URI",
    "AWS_CONTAINER_AUTHORIZATION_TOKEN",
    "AWS_CONTAINER_AUTHORIZATION_TOKEN_FILE",
    "AWS_WEB_IDENTITY_TOKEN_FILE",
    "AWS_ROLE_ARN",
    "AWS_ROLE_SESSION_NAME",
    "AWS_CREDENTIAL_EXPIRATION",
    "AWS_SHARED_CREDENTIALS_FILE",
    "AWS_CONFIG_FILE",
    "AWS_PROFILE",
)
#: Set in every child: boto3 fails fast (``NoCredentialsError``) instead of
#: waiting on an instance metadata service Lambda does not have.
_CHILD_NO_METADATA_ENV = {"AWS_EC2_METADATA_DISABLED": "true"}
#: The re-exec'd image reads the role credentials (JSON) from this fd.
CREDENTIALS_FD_ENV = "AA_ROLE_CREDENTIALS_FD"
#: Function-level override for the scoped role's ARN; otherwise derived from
#: the execution role's name (`_scoped_role_arn`) - core.yaml pins both.
SANDBOX_ROLE_ARN_ENV = "AA_SANDBOX_ROLE_ARN"
EXECUTION_ROLE_SUFFIX = "-lambda-execution-role"
SANDBOX_ROLE_SUFFIX = "-sandbox-scoped-role"
#: Where the delivery bootstrap unpacks the pair (``AA_HANDLER_DIR``, default
#: as agent-management's ``OVERLAY_DIR``): the re-exec runs only for the
#: delivered overlay, mirroring the bootstrap's own shadow check.
OVERLAY_DIR_ENV = "AA_HANDLER_DIR"
DEFAULT_OVERLAY_DIR = "/tmp/aa_handler"
#: The exec-only copy of the interpreter the clean image runs from
#: (`_nondumpable_image`), next to this file in the overlay; mode 0100.
NONDUMPABLE_IMAGE_NAME = ".python-nondumpable"
#: How the previous image started this one - ``exec-only:<pre-flight
#: verdict>`` or ``plain:<why not>`` - set for the re-exec, popped at import.
EXEC_IMAGE_ENV = "AA_HANDLER_EXEC_IMAGE"
#: The exec-only copy is run once as a child before this process commits to
#: it; a cold start's init phase has 10 s in total.
IMAGE_PREFLIGHT_TIMEOUT_S = 5
#: Optional event field, trusted because only the platform's services can
#: invoke the function: extra key prefixes this invocation may READ, each
#: rooted at ``users/``. A governed deployment passes ``["users/"]``.
READ_PREFIXES_EVENT_KEY = "read_prefixes"
#: One secret per agent (``connector_creds.py`` creates it under this name).
#: Secrets Manager appends ``-`` and six random characters to the name in
#: the ARN, hence ``??????`` in the session policy.
CONNECTOR_SECRET_PREFIX = "alphaagent-connector-creds-"
#: Role chaining (the execution role is itself an assumed role) caps
#: AssumeRole at one hour; renew well before an invocation could outlive it.
_SCOPED_SESSION_DURATION_S = 3600
_SCOPED_SESSION_REFRESH_MARGIN_S = 20 * 60
_SCOPED_SESSION_CACHE_MAX = 64
_STS_CONNECT_TIMEOUT_S = 3
_STS_READ_TIMEOUT_S = 10
#: What may be embedded in a session policy: no wildcards, no ``$`` (IAM
#: policy variables), no path tricks. A uid or agent id is a single segment.
#: An owner id may carry a colon (``apikey:<key_id>``, ``service:<name>``);
#: the colon is not an IAM metacharacter and S3 keys may contain it.
_POLICY_ID_RE = re.compile(r"^[A-Za-z0-9_.@:-]{1,128}$")
_POLICY_PREFIX_RE = re.compile(r"^[A-Za-z0-9_./@:-]{1,512}$")
_USERS_PREFIX_RE = re.compile(r"^users/([^/]+)/")

#: This execution environment's role credentials - process memory only.
_ROLE_CREDENTIALS: Optional[Dict[str, str]] = None
#: What the isolation layers did at import, for the probe (`_isolation_report`).
_ISOLATION: Dict[str, Any] = {}


def _sandboxed_env(base: Dict[str, str]) -> Dict[str, str]:
    """`base`, minus this Lambda's own IAM credentials (layer 4, first half;
    `_exec_env` adds the rest). Exactly these four names and no others: every
    connector variable is ``CONN_<connector_id>_...`` (the shape the
    service's ``connector_loader.build_connector_env_vars`` produces and this
    file's own secret payload carries), so removing these cannot drop a real
    connector credential."""
    return {k: v for k, v in base.items() if k not in _LAMBDA_CREDENTIAL_ENV_KEYS}


# -- layer 1: re-exec ------------------------------------------------------

def _reexec_argv(handler_dir: str, handler_args: List[str]) -> List[str]:
    """The command the clean image runs: this interpreter with the flags it
    was started with, an inline bootstrap that puts ``handler_dir`` first on
    ``sys.path``, imports ``handler`` and starts ``awslambdaric`` the way the
    delivery bootstrap does (``runpy``), then the handler name
    ``awslambdaric`` reads as ``argv[1]``."""
    inline = (
        f"import sys;sys.path.insert(0,{handler_dir!r});import handler;"
        "import runpy;runpy.run_module('awslambdaric',run_name='__main__',alter_sys=True)"
    )
    orig = list(getattr(sys, "orig_argv", None) or [])
    flags = orig[1:orig.index("-c")] if "-c" in orig[1:] else []
    return [sys.executable, *flags, "-c", inline, *handler_args]


def _reexec_skip_reason(module_name: str) -> Optional[str]:
    """Why NOT to re-exec at this import, or None when it should happen. The
    conditions are those under which the delivery bootstrap imported us:
    real runtime, credentials in the environment, ``awslambdaric`` not yet
    started, imported as ``handler`` from the delivered overlay."""
    env = os.environ
    if env.get(CREDENTIALS_FD_ENV):
        return "done"
    if not (env.get("AWS_ACCESS_KEY_ID") and env.get("AWS_SECRET_ACCESS_KEY")):
        return "not needed: no role credentials in the environment"
    if not env.get("AWS_LAMBDA_RUNTIME_API"):
        return "skipped: not running under the Lambda runtime API"
    if "awslambdaric" in sys.modules:
        return "skipped: awslambdaric is already running (image entry point or baked fallback)"
    if module_name != "handler":
        return f"skipped: imported as {module_name!r}, not as the handler module"
    overlay = os.path.realpath(env.get(OVERLAY_DIR_ENV) or DEFAULT_OVERLAY_DIR)
    here = os.path.realpath(_HANDLER_DIR)
    if here != overlay and not here.startswith(overlay + os.sep):
        return f"skipped: not the delivered overlay ({_HANDLER_DIR})"
    return None


def _reexec_with_clean_environment(handler_args: List[str]) -> str:
    """Layer 1. Hand the role credentials to the next image of this process
    over a pipe and ``execve`` it with an environment that carries none.
    Returns only on failure, with the reason; on success this process is
    replaced and nothing after the call runs."""
    creds = {k: os.environ[k] for k in _LAMBDA_CREDENTIAL_ENV_KEYS if os.environ.get(k)}
    env = {k: v for k, v in os.environ.items() if k not in _LAMBDA_CREDENTIAL_ENV_KEYS}
    # The bootstrap records its verdict AFTER importing us; this import
    # never returns to it, so record the verdict it was about to.
    env.setdefault(HANDLER_DELIVERY_ENV, "overlay")
    # Layer 2 rides on this exec: the clean image runs from an exec-only copy
    # of the interpreter and is non-dumpable from its first instruction - the
    # kernel path Lambda's seccomp cannot refuse (`_nondumpable_image`).
    image, note = _nondumpable_image(_HANDLER_DIR, env)
    env[EXEC_IMAGE_ENV] = f"exec-only:{note}" if image else f"plain:{note}"
    read_fd, write_fd = os.pipe()
    try:
        os.set_inheritable(read_fd, True)
        payload = json.dumps(creds).encode("utf-8")
        written = 0
        while written < len(payload):
            written += os.write(write_fd, payload[written:])
        os.close(write_fd)
        write_fd = -1
        env[CREDENTIALS_FD_ENV] = str(read_fd)
        argv = _reexec_argv(_HANDLER_DIR, handler_args)
        print(f"[HANDLER] credential isolation: re-executing with a clean environment "
              f"({len(creds)} credential variable(s) moved to memory; image={env[EXEC_IMAGE_ENV]})", flush=True)
        sys.stdout.flush()
        sys.stderr.flush()
        if image:
            try:
                os.execve(image, argv, env)
            except OSError as exc:
                # Never a failed cold start: the plain interpreter, reported.
                env[EXEC_IMAGE_ENV] = f"plain:execve of the exec-only copy failed: {type(exc).__name__}: {exc}"
                print(f"[HANDLER] credential isolation: {env[EXEC_IMAGE_ENV]}", flush=True)
        os.execve(sys.executable, argv, env)
    except OSError as exc:
        for fd in (read_fd, write_fd):
            if fd >= 0:
                try:
                    os.close(fd)
                except OSError:
                    pass
        return f"failed: {type(exc).__name__}: {exc}"
    return "failed: execve returned"


def _read_fd_fully(fd: int) -> bytes:
    chunks: List[bytes] = []
    while True:
        chunk = os.read(fd, 65536)
        if not chunk:
            return b"".join(chunks)
        chunks.append(chunk)


def _capture_role_credentials(scrub_environment: bool) -> Tuple[Optional[Dict[str, str]], str, Optional[str]]:
    """Move the role credentials into memory: from the pipe the previous image
    of this process left (`_reexec_with_clean_environment`), else from the
    environment - the paths that did not re-exec - which, under the runtime,
    is then cleared of them so they are not base names for a child
    (`_BASE_ENV_KEYS`). The previous image's other hand-over, how it started
    this one (`EXEC_IMAGE_ENV`: ``exec-only:<verdict>`` | ``plain:<why>``),
    leaves the environment here too. Returns ``(credentials or None, "pipe"
    | "environment" | "none", exec_image or None)``."""
    exec_image = os.environ.pop(EXEC_IMAGE_ENV, None)
    fd_raw = os.environ.pop(CREDENTIALS_FD_ENV, None)
    if fd_raw:
        try:
            fd = int(fd_raw)
            try:
                raw = _read_fd_fully(fd)
            finally:
                os.close(fd)
            creds = {k: str(v) for k, v in json.loads(raw.decode("utf-8")).items()
                     if k in _LAMBDA_CREDENTIAL_ENV_KEYS and v}
            if creds.get("AWS_ACCESS_KEY_ID") and creds.get("AWS_SECRET_ACCESS_KEY"):
                return creds, "pipe", exec_image
            print("[HANDLER] credential isolation: the pipe carried no usable role credentials")
        except (OSError, ValueError) as exc:
            print(f"[HANDLER] credential isolation: could not read the role credentials from the pipe: {exc}")
    creds = {k: os.environ[k] for k in _LAMBDA_CREDENTIAL_ENV_KEYS if os.environ.get(k)}
    if scrub_environment:
        for k in _LAMBDA_CREDENTIAL_ENV_KEYS:
            os.environ.pop(k, None)
    if creds.get("AWS_ACCESS_KEY_ID") and creds.get("AWS_SECRET_ACCESS_KEY"):
        return creds, "environment", exec_image
    return None, "none", exec_image


# -- layer 2: non-dumpable ---------------------------------------------------

#: Runs INSIDE the exec-only copy, as a child of the pre-exec image, before
#: that image commits to it: does the copy start at all, and does this kernel
#: make a process started from it non-dumpable? One word on stdout.
_IMAGE_PREFLIGHT = (
    "import os\n"
    "print('unknown' if os.geteuid()==0 else "
    "('nondumpable' if os.stat('/proc/self/status').st_uid==0 else 'dumpable'))"
)
_PR_GET_DUMPABLE, _PR_SET_DUMPABLE = 3, 4


def _nondumpable_image(handler_dir: str, env: Dict[str, str]) -> Tuple[Optional[str], str]:
    """Layer 2, at the re-exec: an exec-only copy (mode 0100) of this
    interpreter for the clean image to run from. ``execve`` of a file its
    caller may not read makes the new process non-dumpable from its first
    instruction (``would_dump`` -> ``BINPRM_FLAGS_ENFORCE_NONDUMP`` ->
    ``set_dumpable(mm, fs.suid_dumpable)``): the kernel path that needs no
    ``prctl``, which Lambda's seccomp refuses. The copy is pre-flighted as a
    child before this process commits to it, so an interpreter that cannot
    run from a copy (an ``$ORIGIN`` rpath, a noexec /tmp) means a plain
    re-exec, never a failed cold start. Returns ``(path, verdict)`` - the
    pre-flight's ``nondumpable`` | ``dumpable`` | ``unknown`` (root) - or
    ``(None, reason)``. Linux only: a macOS framework build cannot start
    from a copy, and there is no /proc to protect there."""
    if not sys.platform.startswith("linux"):
        return None, "not linux"
    path = os.path.join(handler_dir, NONDUMPABLE_IMAGE_NAME)
    part = path + ".part"
    try:
        shutil.copyfile(os.path.realpath(sys.executable), part)
        os.chmod(part, 0o100)
        os.replace(part, path)
        # argv[0] stays the real interpreter: getpath derives sys.executable
        # and sys.prefix from it, so the clean image - and every child it
        # starts - keeps the installation's paths.
        proc = subprocess.run([sys.executable, "-c", _IMAGE_PREFLIGHT], executable=path, env=env,
                              capture_output=True, text=True, timeout=IMAGE_PREFLIGHT_TIMEOUT_S)
        lines = (proc.stdout or "").strip().splitlines()
        verdict = lines[-1].strip() if lines else ""
        if proc.returncode != 0 or verdict not in ("nondumpable", "dumpable", "unknown"):
            detail = (proc.stderr or proc.stdout or "").strip()[-200:] or "no verdict"
            raise OSError(f"pre-flight exit {proc.returncode}: {detail}")
        return path, verdict
    except (OSError, subprocess.SubprocessError) as exc:
        for stale in (part, path):
            try:
                os.remove(stale)
            except OSError:
                pass
        return None, f"{type(exc).__name__}: {exc}"


def _prctl(option: int, arg2: int = 0) -> Tuple[int, int]:
    """``prctl(option, arg2, 0, 0, 0)`` -> ``(rc, errno)``; ``(-1, ENOSYS)``
    where there is no prctl to call (macOS, no ctypes)."""
    if not sys.platform.startswith("linux"):
        return -1, errno.ENOSYS
    try:
        import ctypes
        fn = ctypes.CDLL(None, use_errno=True).prctl
        fn.argtypes = [ctypes.c_int, ctypes.c_ulong, ctypes.c_ulong, ctypes.c_ulong, ctypes.c_ulong]
        fn.restype = ctypes.c_int
        rc = fn(option, arg2, 0, 0, 0)
        return rc, (ctypes.get_errno() if rc != 0 else 0)
    except (OSError, AttributeError, ValueError):
        return -1, errno.ENOSYS


def _set_non_dumpable() -> Dict[str, Any]:
    """Layer 2 by ``prctl(PR_SET_DUMPABLE, 0)``, where the kernel lets this
    process call prctl at all. Returns ``{"dumpable": <the flag afterwards,
    or None>, "prctl": "ok" | "<ERRNO>: <message>"}`` - a refusal is
    recorded, never swallowed: under Lambda's seccomp every prctl is EPERM
    and the exec-only image (`_nondumpable_image`) carries the layer."""
    rc, err = _prctl(_PR_SET_DUMPABLE, 0)
    if rc != 0:
        return {"dumpable": None, "prctl": f"{errno.errorcode.get(err, err)}: {os.strerror(err)}"}
    rc, _ = _prctl(_PR_GET_DUMPABLE)
    return {"dumpable": rc if rc >= 0 else None, "prctl": "ok"}


def _measure_non_dumpable(dumpable: Optional[int] = None) -> Dict[str, Any]:
    """What the kernel will actually do about a same-uid reader, measured
    without prctl: ``task_dump_owner`` gives a non-dumpable task's
    ``/proc/<pid>/{status,environ,maps,mem}`` to root, so ``/proc/self/
    status`` owned by root while this process is not root means the READ /
    ATTACH gate is closed. Returns ``{"non_dumpable": True | False | None,
    "proc_owner_uid": ...}``; None where there is no /proc, or the process is
    root (root owns its files either way) and prctl reported no flag."""
    out: Dict[str, Any] = {"non_dumpable": None, "proc_owner_uid": None}
    try:
        out["proc_owner_uid"] = os.stat("/proc/self/status").st_uid
    except OSError:
        pass
    if out["proc_owner_uid"] is not None and os.geteuid() != 0:
        out["non_dumpable"] = out["proc_owner_uid"] == 0
    elif dumpable is not None:
        out["non_dumpable"] = dumpable != 1
    return out


def _startup_environ_credential_names(at_import: List[str]) -> Tuple[Optional[List[str]], str]:
    """Credential variable names in THIS process's start-up environment block
    - what a child would find in ``/proc/<ppid>/environ`` if it could open
    it. Read from ``/proc/self/environ`` where the kernel lets this process
    read its own (a dumpable one); a process that STARTED non-dumpable (the
    exec-only image) is refused its own block - root owns it - and reports
    ``os.environ`` as it stood at import instead, which for a process that
    has not yet touched its environment IS that block. Returns ``(names or
    None, source)``."""
    try:
        with open("/proc/self/environ", "rb") as fh:
            block = fh.read()
    except OSError:
        if sys.platform.startswith("linux"):
            return list(at_import), "os.environ at import"
        return None, "no /proc"
    names = {entry.split(b"=", 1)[0].decode("utf-8", "replace") for entry in block.split(b"\0") if entry}
    return sorted(k for k in _LAMBDA_CREDENTIAL_ENV_KEYS if k in names), "/proc/self/environ"


def _reassert_non_dumpable() -> None:
    """Layer 2 on every invocation: re-apply prctl where it works, re-measure,
    and record. Nothing in this process resets the flag (only its own execve
    or a credential change would; it does neither), so this is the check,
    not the mechanism - and a serving process found dumpable is a WARNING in
    the log on every invocation, never silent."""
    if not sys.platform.startswith("linux"):
        return
    state = _set_non_dumpable()
    state.update(_measure_non_dumpable(state["dumpable"]))
    _ISOLATION.update(state)
    if state["non_dumpable"] is False:
        print("[HANDLER] WARNING: credential isolation: the serving process is dumpable - a same-uid child can "
              f"open its /proc/<pid>/environ, maps and mem (exec_image={_ISOLATION.get('exec_image')}, "
              f"prctl={state['prctl']})", flush=True)


def _read_ptrace_scope() -> Optional[str]:
    """``kernel.yama.ptrace_scope`` as the runtime's kernel reports it, or
    None when Yama is not there. Recorded, not relied on (layer 2)."""
    try:
        with open("/proc/sys/kernel/yama/ptrace_scope", "r", encoding="ascii") as fh:
            return fh.read().strip() or None
    except OSError:
        return None


def _isolate_at_import(module_name: str) -> None:
    """Layers 1 and 2, in import order: re-exec if this is the import the
    bootstrap made (never returns then), capture the credentials, snapshot
    what /proc would show a child, go non-dumpable."""
    global _ROLE_CREDENTIALS
    # The environment as this process was started with it, before anything
    # here pops a name (`_startup_environ_credential_names` may need it).
    at_import = sorted(k for k in _LAMBDA_CREDENTIAL_ENV_KEYS if k in os.environ)
    reason = _reexec_skip_reason(module_name)
    if reason is None:
        reason = _reexec_with_clean_environment(sys.argv[1:])
    _ISOLATION["reexec"] = reason
    under_runtime = bool(os.environ.get("AWS_LAMBDA_RUNTIME_API"))
    # ... plus how the previous image started this one: "exec-only:<verdict>"
    # or "plain:<why not>" (popped there - never a base name for a child).
    _ROLE_CREDENTIALS, _ISOLATION["credential_source"], _ISOLATION["exec_image"] = _capture_role_credentials(
        scrub_environment=under_runtime)
    leaked, source = _startup_environ_credential_names(at_import)
    _ISOLATION["proc_environ_credential_names"] = leaked
    _ISOLATION["proc_environ_clean"] = None if leaked is None else not leaked
    _ISOLATION["proc_environ_source"] = source
    state = _set_non_dumpable()
    state.update(_measure_non_dumpable(state["dumpable"]))
    _ISOLATION.update(state)
    _ISOLATION["ptrace_scope"] = _read_ptrace_scope() or "absent"
    _ISOLATION["environment_credential_names"] = sorted(k for k in _LAMBDA_CREDENTIAL_ENV_KEYS if k in os.environ)


_isolate_at_import(__name__)


# -- layer 3: per-invocation scope -------------------------------------------

class AccessScope(NamedTuple):
    """What one invocation may reach, derived from the event by
    `_scope_from_event`: the bucket, the prefix it may WRITE under
    (``users/<uid>/``; a legacy ``workspaces/sessions/...`` prefix is taken
    as itself), the prefixes it may READ (its own plus the caller's
    `READ_PREFIXES_EVENT_KEY`), and the agent whose ONE connector secret it
    may read (empty when the event names no connectors)."""
    bucket: str
    write_prefix: str
    read_prefixes: Tuple[str, ...]
    agent_id: str

    def grants_anything(self) -> bool:
        return bool(self.bucket and self.write_prefix) or bool(self.agent_id)

    def cache_key(self) -> str:
        return json.dumps(self._asdict(), sort_keys=True, separators=(",", ":"))

    def describe(self) -> str:
        s3 = f"s3://{self.bucket}/ write={self.write_prefix} read={list(self.read_prefixes)}" if self.bucket else "no S3"
        return f"{s3}; secret={CONNECTOR_SECRET_PREFIX + self.agent_id if self.agent_id else 'none'}"

    def policy(self, partition: str = "aws") -> Dict[str, Any]:
        """The STS session policy: the intersection with the scoped role's
        own grant is exactly this invocation's reach. ``ListBucket`` is
        conditioned on ``s3:prefix`` (a listing of ``users/`` alone has no
        matching prefix and is refused); objects are named by prefix."""
        statements: List[Dict[str, Any]] = []
        if self.bucket and self.write_prefix:
            bucket_arn = f"arn:{partition}:s3:::{self.bucket}"
            statements.append({
                "Sid": "ListWithinScope",
                "Effect": "Allow",
                "Action": ["s3:ListBucket"],
                "Resource": bucket_arn,
                "Condition": {"StringLike": {"s3:prefix": [f"{p}*" for p in self.read_prefixes]}},
            })
            statements.append({
                "Sid": "ReadWithinScope",
                "Effect": "Allow",
                "Action": ["s3:GetObject", "s3:GetObjectTagging"],
                "Resource": [f"{bucket_arn}/{p}*" for p in self.read_prefixes],
            })
            statements.append({
                "Sid": "WriteOwnPrefix",
                "Effect": "Allow",
                "Action": ["s3:PutObject", "s3:PutObjectTagging"],
                "Resource": [f"{bucket_arn}/{self.write_prefix}*"],
            })
        if self.agent_id:
            statements.append({
                "Sid": "OwnConnectorSecret",
                "Effect": "Allow",
                "Action": ["secretsmanager:GetSecretValue"],
                "Resource": [f"arn:{partition}:secretsmanager:*:*:secret:{CONNECTOR_SECRET_PREFIX}{self.agent_id}-??????"],
            })
        return {"Version": "2012-10-17", "Statement": statements}


def _clean_prefix(raw: Any) -> Optional[str]:
    """A key prefix fit for a policy: no leading slash, a trailing one, no
    ``..`` segment, no wildcard or policy-variable characters."""
    text = str(raw or "").strip().lstrip("/")
    if not text or ".." in text.split("/") or not _POLICY_PREFIX_RE.match(text):
        return None
    return text.rstrip("/") + "/"


def _scope_from_event(event: Dict[str, Any]) -> AccessScope:
    """The invocation's reach, from what the trusted caller put in the event.
    The owner is read off ``s3_prefix`` (``users/<uid>/...``): the service
    already resolved which tree this session acts in (its own
    ``resolve_owned_project_prefix``), and a governed preview of another
    user's file legitimately acts in that user's tree. Anything that does
    not parse grants nothing, and `_s3_client` then refuses."""
    bucket = str(event.get("s3_bucket") or "").strip()
    prefix = _clean_prefix(event.get("s3_prefix"))
    write_prefix = ""
    if prefix:
        m = _USERS_PREFIX_RE.match(prefix)
        if m:
            write_prefix = f"users/{m.group(1)}/" if _POLICY_ID_RE.match(m.group(1)) else ""
        else:
            write_prefix = prefix           # legacy layout: exactly this tree
    if not _POLICY_ID_RE.match(bucket or "-"):
        bucket = ""
    reads: List[str] = [write_prefix] if write_prefix else []
    extra = event.get(READ_PREFIXES_EVENT_KEY)
    for item in (extra if isinstance(extra, list) else []):
        cleaned = _clean_prefix(item)
        if cleaned and cleaned.startswith("users/") and cleaned not in reads:
            reads.append(cleaned)
    agent_id = str(event.get("agent_id") or "").strip()
    if not event.get("connector_ids") or not _POLICY_ID_RE.match(agent_id):
        agent_id = ""
    return AccessScope(
        bucket=bucket if write_prefix else "",
        write_prefix=write_prefix,
        read_prefixes=tuple(reads),
        agent_id=agent_id,
    )


def _aws_region() -> Optional[str]:
    return os.environ.get("AWS_REGION") or os.environ.get("AWS_DEFAULT_REGION") or None


def _error_code(exc: BaseException) -> str:
    response = getattr(exc, "response", None)
    if isinstance(response, dict):
        code = (response.get("Error") or {}).get("Code")
        if code:
            return str(code)
    return type(exc).__name__


def _role_session():
    """A session on the execution role's own credentials (memory), or on the
    default chain when none were captured (a local run). Used for two calls
    - ``sts:GetCallerIdentity`` and ``sts:AssumeRole`` - and for the fallback
    `_scoped_session` documents; never handed to user work directly."""
    if _ROLE_CREDENTIALS:
        return boto3.Session(
            aws_access_key_id=_ROLE_CREDENTIALS.get("AWS_ACCESS_KEY_ID"),
            aws_secret_access_key=_ROLE_CREDENTIALS.get("AWS_SECRET_ACCESS_KEY"),
            aws_session_token=(_ROLE_CREDENTIALS.get("AWS_SESSION_TOKEN")
                               or _ROLE_CREDENTIALS.get("AWS_SECURITY_TOKEN") or None),
            region_name=_aws_region(),
        )
    return boto3.Session(region_name=_aws_region())


def _sts_endpoint(region: str) -> str:
    suffix = "amazonaws.com.cn" if region.startswith("cn-") else "amazonaws.com"
    return f"https://sts.{region}.{suffix}"


def _sts_client():
    """STS on the REGIONAL endpoint - core.yaml gives the VPC an STS interface
    endpoint whose private DNS covers ``sts.<region>.amazonaws.com``, not the
    global name - with short timeouts, so an unreachable STS fails in
    seconds rather than a minute."""
    kwargs: Dict[str, Any] = {}
    config_cls = getattr(boto3.session, "Config", None)
    if config_cls is not None:
        kwargs["config"] = config_cls(connect_timeout=_STS_CONNECT_TIMEOUT_S,
                                      read_timeout=_STS_READ_TIMEOUT_S,
                                      retries={"max_attempts": 2})
    region = _aws_region()
    if region:
        kwargs["endpoint_url"] = _sts_endpoint(region)
    return _role_session().client("sts", **kwargs)


def _derive_sandbox_role_arn(caller_arn: str) -> Optional[str]:
    """``arn:<partition>:sts::<account>:assumed-role/<project>-<env>-lambda-
    execution-role/<function>`` -> ``arn:<partition>:iam::<account>:role/
    <project>-<env>-sandbox-scoped-role``; None for any other shape."""
    parts = caller_arn.split(":")
    if len(parts) != 6:
        return None
    partition, account, resource = parts[1], parts[4], parts[5]
    segments = resource.split("/")
    if len(segments) < 2 or segments[0] != "assumed-role":
        return None
    role_name = segments[1]
    if not role_name.endswith(EXECUTION_ROLE_SUFFIX) or not partition or not account:
        return None
    return f"arn:{partition}:iam::{account}:role/{role_name[:-len(EXECUTION_ROLE_SUFFIX)]}{SANDBOX_ROLE_SUFFIX}"


#: The derived scoped-role ARN once known; "" when this role's name does not
#: follow the convention (a fact about the deployment, cached); None = not yet
#: resolved (an STS failure is NOT cached: the next invocation tries again).
_SCOPED_ROLE_ARN: Optional[str] = None


def _scoped_role_arn() -> Optional[str]:
    """The sandbox role to assume: `SANDBOX_ROLE_ARN_ENV` when the function
    configuration names it, else derived once per container from this
    process's own identity by the name convention core.yaml pins."""
    global _SCOPED_ROLE_ARN
    override = (os.environ.get(SANDBOX_ROLE_ARN_ENV) or "").strip()
    if override:
        return override
    if _SCOPED_ROLE_ARN is not None:
        return _SCOPED_ROLE_ARN or None
    try:
        caller_arn = str(_sts_client().get_caller_identity()["Arn"])
    except Exception as exc:
        print(f"[HANDLER] WARNING: could not resolve this function's own identity ({_error_code(exc)})")
        return None
    derived = _derive_sandbox_role_arn(caller_arn)
    if derived is None:
        print(f"[HANDLER] WARNING: no sandbox role can be derived from {caller_arn}; "
              f"set {SANDBOX_ROLE_ARN_ENV} on the function")
    _SCOPED_ROLE_ARN = derived or ""
    return derived


def _partition_of(arn: str) -> str:
    parts = arn.split(":")
    return parts[1] if len(parts) > 2 and parts[1] else "aws"


def _session_name(scope: AccessScope) -> str:
    return "aa-env-" + hashlib.sha256(scope.cache_key().encode("utf-8")).hexdigest()[:24]


#: ``{scope.cache_key(): (session, expires_epoch)}`` - scoped sessions only;
#: a fallback session is never cached here.
_SCOPED_SESSIONS: Dict[str, Tuple[Any, float]] = {}


def _scoped_session(scope: AccessScope) -> Tuple[Any, str]:
    """A boto3 session whose credentials reach exactly ``scope``: AssumeRole on
    the sandbox role with ``scope.policy()`` as the session policy, cached
    until `_SCOPED_SESSION_REFRESH_MARGIN_S` before expiry. Returns
    ``(session, "scoped")``, or ``(execution-role session,
    "execution-role-fallback:<reason>")`` when the role cannot be assumed."""
    if not scope.grants_anything():
        raise RuntimeError("this invocation's scope grants no S3 or Secrets Manager access")
    now = time.time()
    key = scope.cache_key()
    hit = _SCOPED_SESSIONS.get(key)
    if hit and hit[1] - now > _SCOPED_SESSION_REFRESH_MARGIN_S:
        return hit[0], "scoped"
    role_arn = _scoped_role_arn()
    if not role_arn:
        return _role_session(), "execution-role-fallback:no-scoped-role"
    try:
        response = _sts_client().assume_role(
            RoleArn=role_arn,
            RoleSessionName=_session_name(scope),
            Policy=json.dumps(scope.policy(_partition_of(role_arn)), separators=(",", ":")),
            DurationSeconds=_SCOPED_SESSION_DURATION_S,
        )
    except Exception as exc:
        code = _error_code(exc)
        print(f"[HANDLER] WARNING: could not assume {role_arn} ({code}); this invocation runs on the "
              f"execution role's own credentials")
        return _role_session(), f"execution-role-fallback:{code}"
    granted = response["Credentials"]
    session = boto3.Session(
        aws_access_key_id=granted["AccessKeyId"],
        aws_secret_access_key=granted["SecretAccessKey"],
        aws_session_token=granted["SessionToken"],
        region_name=_aws_region(),
    )
    expiry = granted.get("Expiration")
    expires_at = expiry.timestamp() if hasattr(expiry, "timestamp") else now + _SCOPED_SESSION_DURATION_S
    if len(_SCOPED_SESSIONS) >= _SCOPED_SESSION_CACHE_MAX:
        _SCOPED_SESSIONS.pop(min(_SCOPED_SESSIONS, key=lambda k: _SCOPED_SESSIONS[k][1]), None)
    _SCOPED_SESSIONS[key] = (session, expires_at)
    return session, "scoped"


#: Set by `handler` for the duration of one invocation; None between them.
_CURRENT_SCOPE: Optional[AccessScope] = None
#: The session this invocation resolved to (scoped or fallback), once.
_INVOCATION_SESSION: Optional[Tuple[Any, str]] = None


def _invocation_session():
    """The one session this invocation's S3/Secrets clients are built from.
    Resolved once per invocation (a fallback is therefore decided once, not
    re-tried on every client); refused outside `handler`."""
    global _INVOCATION_SESSION
    if _CURRENT_SCOPE is None:
        raise RuntimeError("S3 and Secrets Manager are reachable only inside handler(): no invocation scope is set")
    if _INVOCATION_SESSION is None:
        _INVOCATION_SESSION = _scoped_session(_CURRENT_SCOPE)
        _ISOLATION["credential_scope"] = _INVOCATION_SESSION[1]
    return _INVOCATION_SESSION[0]


def _s3_client():
    """THE S3 client for user work - built from this invocation's scoped
    session, nothing else."""
    return _invocation_session().client("s3")


def _secrets_client():
    """THE Secrets Manager client for the connector secret - same session."""
    return _invocation_session().client("secretsmanager")


def _no_core_dumps() -> None:
    """``preexec_fn`` for every child: RLIMIT_CORE soft=0, hard untouched.
    Lowering the soft limit is always permitted (proven for the sandbox uid;
    ``ulimit -c 0`` from sh was EPERM there). Without this one crashing Chrome
    wrote ~16 x 82 MB cores per attempt and a warm instance's 512 MB /tmp was
    full after six renders — every later action died with ENOSPC."""
    _soft, hard = resource.getrlimit(resource.RLIMIT_CORE)
    resource.setrlimit(resource.RLIMIT_CORE, (0, hard))


def _home_usable(path: Optional[str]) -> bool:
    """Readable AND writable AND searchable by this uid. ``os.access`` says no
    for a nonexistent path, for another uid's 0700 dir (unreadable — the live
    Lambda case) and for a read-only one alike."""
    return bool(path) and os.access(path, os.R_OK | os.W_OK | os.X_OK)


def _home_scope(event: Dict[str, Any], exec_token: str) -> str:
    """Whose HOME a child gets, from what the invocation already carries:
    the invoking agent (``agent:<agent_id>``), else the owning user read off
    the project prefix ``users/<uid>/...`` (``user:<uid>``), else this exec
    alone (``exec:<token>`` - the exec's private helper directory's random
    name, so two anonymous execs never meet). A legacy
    ``workspaces/sessions/...`` prefix names no owner."""
    agent_id = str(event.get("agent_id") or "").strip()
    if agent_id:
        return f"agent:{agent_id}"
    m = re.match(r"^/?users/([^/]+)/", str(event.get("s3_prefix") or ""))
    if m:
        return f"user:{m.group(1)}"
    return f"{EXEC_HOME_SCOPE_PREFIX}{exec_token}"


def _scope_path(scope: str) -> str:
    """Where `_scope_home` puts ``scope`` - without creating it."""
    digest = hashlib.sha256(scope.encode("utf-8")).hexdigest()[:16]
    return os.path.join(TMP_ROOT, f"{HOME_DIR_PREFIX}{digest}")


def _scope_home(scope: str) -> str:
    """``/tmp/home-<sha256(scope)[:16]>``, created 0700 on demand with its
    ``.cache``. Deterministic per scope, so the same agent's later warm turns
    find the same directory - a ``pip install --user`` from one turn is
    importable in the next; ``chmod`` after ``makedirs`` because the latter
    is subject to the umask and the directory may pre-exist. Never raises: a
    full /tmp is the child's problem to report, not the handler's to die on."""
    home = _scope_path(scope)
    try:
        os.makedirs(home, mode=0o700, exist_ok=True)
        os.chmod(home, 0o700)
        os.makedirs(os.path.join(home, ".cache"), exist_ok=True)
    except OSError:
        pass
    return home


def _exec_env(base: Dict[str, str], scope: str = TOOLING_HOME_SCOPE) -> Dict[str, str]:
    """The environment every child gets: `_sandboxed_env`, minus the startup
    hooks (`_CHILD_STARTUP_HOOK_KEYS`) and the variables an SDK follows to
    find credentials (`_CREDENTIAL_SOURCE_ENV_KEYS`), plus
    ``AWS_EC2_METADATA_DISABLED=true`` (credential isolation, layer 4), with
    ``PYTHONSAFEPATH=1`` when the
    interpreter knows the flag (3.11+: neither the script's directory nor the
    cwd is put on ``sys.path`` - the project's ``code/`` is, explicitly, via
    `_child_pythonpath`; ``PYTHONNOUSERSITE`` is deliberately NOT set, it
    would break ``pip install --user``), plus a HOME the runtime uid can
    actually use. When the configured HOME is not usable, HOME/XDG_CACHE_HOME
    point at ``scope``'s private directory under /tmp (`_scope_home`) and
    TMPDIR at /tmp — Chromium, puppeteer, fontconfig, soffice and
    choreographer (kaleido: its browser lookup stats ``~/.local/share/...``
    first and raises PermissionError on an unreadable HOME before it ever
    reads BROWSER_PATH) all need that."""
    env = _sandboxed_env(base)
    for key in (*_CHILD_STARTUP_HOOK_KEYS, *_CREDENTIAL_SOURCE_ENV_KEYS):
        env.pop(key, None)
    env.update(_CHILD_NO_METADATA_ENV)
    if tuple(sys.version_info[:2]) >= (3, 11):
        env["PYTHONSAFEPATH"] = "1"
    if not _home_usable(env.get("HOME")):
        home = _scope_home(scope)
        env["HOME"] = home
        env["XDG_CACHE_HOME"] = os.path.join(home, ".cache")
        env["TMPDIR"] = TMP_ROOT
    return env


# ---------------------------------------------------------------------------
# Per-invocation environment.
#
# A Lambda container is WARM across invocations and one environment function
# serves every agent - and every principal - pinned to that environment. The
# handler used to assign each injected ``CONN_*`` variable straight into its
# own ``os.environ`` and never removed the previous invocation's keys, so a warm
# container accumulated every connector it had ever been invoked with, and
# both exec paths - which built the child's environment from the process
# environment - handed the union to the next agent's sandbox, other
# principals' connector credentials included.
#
# Two independent layers, each sufficient on its own:
#
# 1. `_purge_previous_invocation_env` runs FIRST in every invocation (every
#    action, not only exec) and removes from the process environment the
#    names the previous invocation published (`_PREVIOUS_INVOCATION_KEYS`,
#    recorded by `_track_invocation_env`) and, defensively, every name that
#    is per-invocation by shape (``CONN_*``, ``AA_DEADLINE_EPOCH``) - whatever
#    wrote it, including an older handler on a container that was not
#    recycled.
# 2. `_child_env` builds every child's environment from the image's own
#    environment - the NAMES present when this module was imported
#    (`_BASE_ENV_KEYS`: the image ENV, the function configuration, the
#    runtime's init variables), values read live - plus this invocation's
#    variables only. The handler never writes a per-invocation variable into
#    ``os.environ``, and a name that reaches the process environment later by
#    any other route is not a base name, so it cannot reach a child even when
#    the purge did not run.
#
# NOT per-invocation, never purged: ``AA_UPLOAD_RESERVE_S`` /
# ``AA_YIELD_GRACE_S`` (function configuration read in-process by `_int_env`)
# and a customer's own environment variables (function configuration too) -
# both are base names. Connector variables come from the per-agent secret
# only; a ``CONN_*`` in the function configuration is not a supported shape.
# ---------------------------------------------------------------------------
#: The names the process environment had when this module was imported.
_BASE_ENV_KEYS = frozenset(os.environ)
#: Set by the Lambda runtime per invocation, after import; allowed through by name.
_RUNTIME_PASSTHROUGH_KEYS = ("_X_AMZN_TRACE_ID",)
#: Name shapes that are only ever published per invocation: the connector
#: variables (`connector_loader.build_connector_env_vars`) and the
#: cooperative deadline `aa_env` reads.
_PER_INVOCATION_PREFIXES = ("CONN_",)
_PER_INVOCATION_KEYS = ("AA_DEADLINE_EPOCH",)
#: The previous invocation's variable names (layer 1's tracked set).
_PREVIOUS_INVOCATION_KEYS: set = set()


def _is_per_invocation_name(name: str) -> bool:
    return name in _PER_INVOCATION_KEYS or name.startswith(tuple(_PER_INVOCATION_PREFIXES))


def _purge_previous_invocation_env() -> int:
    """Layer 1. Remove from ``os.environ`` every variable a previous
    invocation could have left behind: the tracked names and every
    per-invocation-shaped name. Nothing of this invocation is in the process
    environment yet (nothing ever is - see `_child_env`), so there is nothing
    to keep. Returns the number removed, for the log line."""
    stale = [k for k in os.environ if k in _PREVIOUS_INVOCATION_KEYS or _is_per_invocation_name(k)]
    for k in stale:
        os.environ.pop(k, None)
    _PREVIOUS_INVOCATION_KEYS.clear()
    return len(stale)


def _track_invocation_env(invocation_env: Dict[str, str]) -> None:
    """Record this invocation's variable names so the NEXT invocation's purge
    removes them by name, independently of the shape rule."""
    _PREVIOUS_INVOCATION_KEYS.clear()
    _PREVIOUS_INVOCATION_KEYS.update(invocation_env)


def _base_env() -> Dict[str, str]:
    """The image's own environment as it stands now: only names that were
    present at import (plus the runtime's per-invocation passthrough), with
    their current values."""
    return {k: v for k, v in os.environ.items()
            if k in _BASE_ENV_KEYS or k in _RUNTIME_PASSTHROUGH_KEYS}


def _child_env(invocation_env: Dict[str, str], scope: str = TOOLING_HOME_SCOPE) -> Dict[str, str]:
    """Layer 2. The environment a child gets: `_exec_env` over the base plus
    THIS invocation's variables - never anything an earlier invocation
    published, by construction - with ``scope``'s HOME (the exec paths pass
    `_home_scope`, the handler's own probes take the default)."""
    return _exec_env({**_base_env(), **invocation_env}, scope=scope)


def _sweep_core_dumps() -> Tuple[int, int, int]:
    """Remove ``/tmp/core.*`` left by a crashed child and return
    ``(removed, tmp_used_bytes, tmp_total_bytes)``. Never raises: a full or
    missing /tmp is reported, not fatal."""
    removed = 0
    for path in glob.glob(os.path.join(TMP_ROOT, CORE_DUMP_GLOB)):
        try:
            if os.path.isfile(path):
                os.remove(path)
                removed += 1
        except OSError:
            pass
    try:
        usage = shutil.disk_usage(TMP_ROOT)
        return removed, int(usage.used), int(usage.total)
    except OSError:
        return removed, 0, 0


#: Directories / suffixes that never leave the Lambda: hidden entries (the
#: sync manifest, .pytest_cache, editor state) plus Python bytecode caches a
#: `python -m pytest code/tests` run leaves behind.
_SKIP_DIRS = ("__pycache__",)
_SKIP_SUFFIXES = (".pyc", ".pyo")


def _walk_visible_files(workdir: str):
    """Yield ``(local_path, rel_path)`` for every uploadable file under
    ``workdir`` — not hidden, not bytecode (see ``_SKIP_DIRS`` /
    ``_SKIP_SUFFIXES``)."""
    for dirpath, dirnames, filenames in os.walk(workdir):
        dirnames[:] = [d for d in dirnames if not d.startswith(".") and d not in _SKIP_DIRS]
        for filename in filenames:
            if filename.startswith(".") or filename.endswith(_SKIP_SUFFIXES):
                continue
            local_path = os.path.join(dirpath, filename)
            yield local_path, os.path.relpath(local_path, workdir).replace(os.sep, "/")


def _snapshot_workdir(workdir: str) -> Dict[str, Tuple[int, int]]:
    """``{rel_path: (mtime_ns, size)}`` of every visible file. Taken right
    before the user code runs; `_upload_workdir_files(snapshot=...)` uploads
    only what differs afterwards, so a warm Lambda's mirrored project (which
    S3 already holds) is never re-uploaded wholesale on every step."""
    snap: Dict[str, Tuple[int, int]] = {}
    for local_path, rel in _walk_visible_files(workdir):
        try:
            st = os.stat(local_path)
        except OSError:
            continue
        snap[rel] = (st.st_mtime_ns, st.st_size)
    return snap


def _sha256_file(path: str) -> Optional[str]:
    """Hex sha256 of a regular file's bytes; ``None`` when it is missing, not
    a regular file (a symlink counts as not one) or unreadable."""
    try:
        if not stat.S_ISREG(os.lstat(path).st_mode):
            return None
        digest = hashlib.sha256()
        with open(path, "rb") as fh:
            for chunk in iter(lambda: fh.read(1024 * 1024), b""):
                digest.update(chunk)
        return digest.hexdigest()
    except OSError:
        return None


def _unlink_quiet(path: str) -> None:
    """Remove ``path`` if it exists - a file or a symlink, never a tree."""
    try:
        if os.path.lexists(path):
            os.unlink(path)
    except OSError:
        pass


def _workspace_cache_path(prefix: str) -> str:
    """Where `_workspace_cache` keeps ``prefix``'s mirror - without creating
    it. Keyed by the PREFIX (what is mirrored), not by the invoking scope:
    two agents on one project share one verified cache, one agent on two
    projects gets two, and with the index in this process's memory the key
    is a matter of cache hits only - never of trust."""
    digest = hashlib.sha256(prefix.strip("/").encode("utf-8")).hexdigest()[:16]
    return os.path.join(TMP_ROOT, f"{WORKSPACE_CACHE_PREFIX}{digest}")


def _workspace_cache(prefix: str) -> str:
    """``/tmp/ws-<sha256(prefix)[:16]>``, created 0700 on demand (``chmod``
    after ``makedirs``: the umask, and the directory may pre-exist). A
    symlink or a plain file sitting at that path was put there by a child
    (this handler only ever makes a directory) and is removed first, so the
    cache is never another directory in disguise."""
    cache_dir = _workspace_cache_path(prefix)
    if os.path.islink(cache_dir) or (os.path.exists(cache_dir) and not os.path.isdir(cache_dir)):
        _unlink_quiet(cache_dir)
    os.makedirs(cache_dir, mode=0o700, exist_ok=True)
    try:
        os.chmod(cache_dir, 0o700)
    except OSError:
        pass
    return cache_dir


def _scrub_cache_dir(cache_dir: str, dirs: List[str]) -> int:
    """Before a sync: under each requested project directory of the cache,
    remove every symlink (S3 objects are files; a download never makes one)
    and every file the index does not vouch for - whatever put it there, it
    was not `_sync_dirs_down`. Directories left empty go too. Returns the
    number of entries removed."""
    index = _WORKSPACE_INDEX.get(cache_dir) or {}
    removed = 0
    for d in dirs:
        top = os.path.join(cache_dir, *d.split("/"))
        if os.path.islink(top):
            _unlink_quiet(top)
            removed += 1
            continue
        if not os.path.isdir(top):
            continue
        for dirpath, dirnames, filenames in os.walk(top, topdown=False):
            for name in filenames:
                path = os.path.join(dirpath, name)
                rel = os.path.relpath(path, cache_dir).replace(os.sep, "/")
                if os.path.islink(path) or rel not in index:
                    _unlink_quiet(path)
                    removed += 1
            for name in dirnames:
                path = os.path.join(dirpath, name)
                if os.path.islink(path):
                    _unlink_quiet(path)
                    removed += 1
                else:
                    try:
                        os.rmdir(path)          # succeeds only when empty
                    except OSError:
                        pass
    return removed


def _materialize_workdir(cache_dir: str, rels: List[str], workdir: str) -> int:
    """Build this exec's working copy: hard-link every verified cache file in
    ``rels`` to the same relative path under ``workdir`` (no bytes copied, no
    extra /tmp; ``copy2`` when the filesystem refuses a link). Returns the
    number of files placed. A child's in-place write through a link reaches
    the cache inode - the next exec's sha256 check catches exactly that."""
    placed = 0
    for rel in rels:
        src = os.path.join(cache_dir, *rel.split("/"))
        dst = os.path.join(workdir, *rel.split("/"))
        try:
            os.makedirs(os.path.dirname(dst) or workdir, exist_ok=True)
            _unlink_quiet(dst)
            try:
                os.link(src, dst)
            except OSError:
                shutil.copy2(src, dst)
            placed += 1
        except OSError as e:
            print(f"[HANDLER] WARNING: could not place {rel} in the working directory: {e}")
    return placed


def _cache_uploaded(cache_dir: str, rel: str, local_path: str, content: bytes, etag: str) -> None:
    """After a successful upload of a file under a synced directory: make the
    cache hold exactly the bytes S3 now has (link the working copy's inode
    in; nothing to do when the child wrote through the link) and record them,
    so the next exec does not download what this one just wrote."""
    cache_path = os.path.join(cache_dir, *rel.split("/"))
    try:
        if not (os.path.exists(cache_path) and os.path.samefile(cache_path, local_path)):
            os.makedirs(os.path.dirname(cache_path) or cache_dir, exist_ok=True)
            _unlink_quiet(cache_path)
            try:
                os.link(local_path, cache_path)
            except OSError:
                shutil.copy2(local_path, cache_path)
    except OSError as e:
        print(f"[HANDLER] WARNING: could not cache {rel}: {e}")
        _WORKSPACE_INDEX.setdefault(cache_dir, {}).pop(rel, None)
        return
    _WORKSPACE_INDEX.setdefault(cache_dir, {})[rel] = (etag, hashlib.sha256(content).hexdigest())


def _upload_workdir_files(
    bucket: str,
    prefix: str,
    workdir: str,
    snapshot: Optional[Dict[str, Tuple[int, int]]] = None,
    cache_dir: Optional[str] = None,
    cache_dirs: Tuple[str, ...] = (),
    failures: Optional[List[Tuple[str, str]]] = None,
) -> List[str]:
    """Upload files from workdir to S3. Returns the uploaded relative paths.

    With ``snapshot`` (from `_snapshot_workdir`, taken before the run) only
    NEW or CHANGED files (different mtime or size) are uploaded — that list
    is exactly what the caller turns into ``file_added`` events. Without it
    every visible file goes up (legacy full upload; `_fetch_object` uses it).
    With ``cache_dir`` (this project's mirror cache, `_workspace_cache`) a
    successful upload under one of ``cache_dirs`` - the directories this exec
    synced - is placed in the cache and recorded with its new ETag and sha256
    (`_cache_uploaded`), so the next `_sync_dirs_down` does not re-download a
    file this very Lambda just wrote.
    With ``failures`` (a list the caller owns) every object S3 refused is
    appended as ``(rel_path, "<error class>: <message>")`` so the caller can
    report it; each refusal is also logged, and once more as a summary.
    """
    s3 = _s3_client()
    uploaded: List[str] = []
    refused: List[Tuple[str, str]] = []
    cache_top = tuple(d.strip("/") for d in cache_dirs if d.strip("/"))

    for local_path, rel_path in _walk_visible_files(workdir):
        if snapshot is not None:
            try:
                st = os.stat(local_path)
            except OSError:
                continue
            if snapshot.get(rel_path) == (st.st_mtime_ns, st.st_size):
                continue
        s3_key = f"{prefix.rstrip('/')}/{rel_path}"
        try:
            with open(local_path, "rb") as f:
                content = f.read()
            resp = s3.put_object(Bucket=bucket, Key=s3_key, Body=content, **_s3_put_args(s3_key))
            uploaded.append(rel_path)
            if cache_dir and any(rel_path.startswith(d + "/") for d in cache_top):
                etag = resp.get("ETag", "") if isinstance(resp, dict) else ""
                _cache_uploaded(cache_dir, rel_path, local_path, content, etag or "")
        except Exception as e:
            print(f"[HANDLER] WARNING: failed to upload {rel_path}: {e}")
            refused.append((rel_path, f"{type(e).__name__}: {e}"[:_UPLOAD_ERROR_MAX_CHARS]))

    if uploaded:
        print(f"[HANDLER] Uploaded {len(uploaded)} files to s3://{bucket}/{prefix.rstrip('/')}/")
    if refused:
        print(f"[HANDLER] WARNING: {len(refused)} file(s) could not be uploaded to "
              f"s3://{bucket}/{prefix.rstrip('/')}/ ({refused[0][0]}: {refused[0][1]})")
        if failures is not None:
            failures.extend(refused)
    return uploaded


def _sync_dirs_down(
    bucket: str,
    prefix: str,
    cache_dir: str,
    dirs: List[str],
    max_bytes: Optional[int] = None,
    skip_over_bytes: Optional[int] = None,
) -> Dict[str, Any]:
    """Bring the project's mirror CACHE (`_workspace_cache`) up to date with
    the given project directories (``code``, ``nodes``, ``inputs``…) from
    ``s3://{bucket}/{prefix}{dir}/``, preserving relative paths, and say
    which files are now present and verified.

    ETag-diffed against this process's index (`_WORKSPACE_INDEX`): an object
    whose ETag the index recorded AND whose cached bytes still hash to what
    was recorded is ``unchanged`` and not downloaded, so a warm Lambda pays
    only for what changed since its last invocation - while a cached file a
    child tampered with (the cache is writable by the sandbox uid) fails the
    hash and is fetched again. `_scrub_cache_dir` first removes what
    the index does not vouch for; a cached file S3 no longer lists is pruned
    once that directory's listing completed. Size guards: an object larger
    than ``skip_over_bytes`` is listed but not downloaded; once the running
    total would exceed ``max_bytes`` the rest is listed as skipped and
    ``truncated`` is set — never a silent partial mirror.

    Returns ``{"downloaded": [...], "skipped": [...], "unchanged": n,
    "bytes": n, "truncated": bool, "present": [...]}`` - ``present`` is every
    rel path the working copy should carry (`_materialize_workdir`).
    Best-effort: a listing failure for a directory is logged and that
    directory's VERIFIED cache entries stand in (the run proceeds with the
    last good mirror, as before); a download failure is reported as skipped.
    """
    report: Dict[str, Any] = {"downloaded": [], "skipped": [], "unchanged": 0, "bytes": 0,
                              "truncated": False, "present": []}
    if not bucket or not prefix or not dirs:
        return report
    max_bytes = int(max_bytes if max_bytes is not None else _int_env("AA_SYNC_MAX_BYTES", _DEFAULT_SYNC_MAX_BYTES))
    skip_over_bytes = int(
        skip_over_bytes if skip_over_bytes is not None
        else _int_env("AA_SYNC_SKIP_OVER_BYTES", _DEFAULT_SYNC_SKIP_OVER_BYTES)
    )
    clean_dirs = [d for d in (str(d).strip("/") for d in dirs) if d and ".." not in d.split("/")]
    scrubbed = _scrub_cache_dir(cache_dir, clean_dirs)
    index = _WORKSPACE_INDEX.setdefault(cache_dir, {})
    present: List[str] = []
    mismatched = 0
    s3 = _s3_client()
    base = prefix.rstrip("/") + "/"
    paginator = s3.get_paginator("list_objects_v2")
    for d_clean in clean_dirs:
        listed: List[str] = []
        try:
            for page in paginator.paginate(Bucket=bucket, Prefix=f"{base}{d_clean}/"):
                for obj in page.get("Contents", []) or []:
                    key = obj["Key"]
                    rel = key[len(base):]
                    if not rel or rel.endswith("/"):
                        continue
                    listed.append(rel)
                    size = int(obj.get("Size", 0) or 0)
                    etag = obj.get("ETag", "") or ""
                    dst = os.path.join(cache_dir, *rel.split("/"))
                    known = index.get(rel)
                    if etag and known and known[0] == etag:
                        if _sha256_file(dst) == known[1]:
                            report["unchanged"] += 1
                            present.append(rel)
                            continue
                        mismatched += 1
                    if size > skip_over_bytes:
                        report["skipped"].append(rel)
                        continue
                    if report["bytes"] + size > max_bytes:
                        report["skipped"].append(rel)
                        report["truncated"] = True
                        continue
                    index.pop(rel, None)
                    try:
                        os.makedirs(os.path.dirname(dst) or cache_dir, exist_ok=True)
                        _unlink_quiet(dst)          # never write through a link or a symlink
                        s3.download_file(bucket, key, dst)
                        digest = _sha256_file(dst)
                        if digest is None:
                            raise OSError(f"downloaded file is not readable: {dst}")
                    except Exception as e:
                        print(f"[HANDLER] WARNING: sync-down failed for {rel}: {e}")
                        _unlink_quiet(dst)
                        report["skipped"].append(rel)
                        continue
                    report["downloaded"].append(rel)
                    report["bytes"] += size
                    index[rel] = (etag, digest)
                    present.append(rel)
        except Exception as e:
            print(f"[HANDLER] WARNING: sync-down listing failed for {d_clean}/: {e} - the verified cache stands in")
            for rel, (_etag, digest) in list(index.items()):
                if rel.startswith(d_clean + "/") and rel not in present:
                    if _sha256_file(os.path.join(cache_dir, *rel.split("/"))) == digest:
                        present.append(rel)
                    else:
                        index.pop(rel, None)
            continue
        # This directory's listing completed: what S3 no longer holds leaves the cache.
        listed_set = set(listed)
        for rel in [r for r in index if r.startswith(d_clean + "/") and r not in listed_set]:
            index.pop(rel, None)
            _unlink_quiet(os.path.join(cache_dir, *rel.split("/")))
    if mismatched:
        print(f"[HANDLER] sync: {mismatched} cache entr{'y' if mismatched == 1 else 'ies'} did not match "
              f"what was downloaded - re-fetched from S3")
    report["present"] = present
    print(
        f"[HANDLER] sync: {len(report['downloaded'])} downloaded, {report['unchanged']} unchanged, "
        f"{len(report['skipped'])} skipped ({report['bytes']} bytes), {scrubbed} scrubbed from the cache, "
        f"from s3://{bucket}/{base} dirs={clean_dirs}"
    )
    return report


def _event_sync_dirs(event: Dict[str, Any]) -> List[str]:
    """``event["sync_dirs"]`` — a list of project directory names to mirror
    before running. Absent/empty means "no mirror" (legacy callers), so the
    pre-existing exec paths behave exactly as before unless asked."""
    dirs = event.get("sync_dirs")
    if not isinstance(dirs, list):
        return []
    return [str(d).strip("/") for d in dirs if str(d).strip("/")]


def _child_pythonpath(workdir: str, base_env: Dict[str, str], helper_dir: str) -> str:
    """``{helper_dir}:{workdir}/code[:existing]`` — `import aa_env` resolves
    from THIS exec's private helper directory (`_helper_dir`) and the
    project's ``code/`` package root is importable from anywhere (``from
    src.x import f`` in scripts and in ``code/tests/`` alike, with cwd =
    project root).

    Never the handler's own directory: under runtime delivery that
    is ``/tmp/aa_handler`` - shared by every agent the warm container serves
    and writable by the sandbox uid, where ``/opt`` was read-only - so a file
    one invocation's script planted there (``sitecustomize.py``, a shadow
    ``pandas.py``) would be imported by the next agent's child, with that
    agent's credentials in its environment. A fresh, unpredictably named
    directory per exec cannot be pre-planted."""
    parts = [helper_dir, os.path.join(workdir, "code")]
    if base_env.get("PYTHONPATH"):
        parts.append(base_env["PYTHONPATH"])
    return os.pathsep.join(parts)


def _helper_dir() -> str:
    """A fresh private directory under /tmp holding this exec's ``aa_env.py``,
    written from the import-time copy (`_AA_ENV_BYTES`) - the only thing a
    child imports from outside its project and the image. Removed by the
    caller once the child has exited. ``mkdtemp``: mode 0700, random name."""
    directory = tempfile.mkdtemp(prefix=".aa-helper-", dir=TMP_ROOT)
    if _AA_ENV_BYTES is not None:
        with open(os.path.join(directory, "aa_env.py"), "wb") as fh:
            fh.write(_AA_ENV_BYTES)
    return directory


def _release_exec_dirs(helper_dir: str, scope: str, child_env: Dict[str, str]) -> None:
    """After the child has exited (or the wall raised): remove the exec's
    helper directory and, when the HOME was this exec's alone (no agent, no
    owner in the event - `_home_scope`), that HOME too. An agent's or owner's
    HOME is theirs to keep across warm turns; a usable configured HOME is
    never touched (it is not under `_scope_home`)."""
    shutil.rmtree(helper_dir, ignore_errors=True)
    home = child_env.get("HOME", "")
    if scope.startswith(EXEC_HOME_SCOPE_PREFIX) and home == _scope_path(scope):
        shutil.rmtree(home, ignore_errors=True)


# ---------------------------------------------------------------------------
# Runtime delivery: which handler is live, and keeping the pair intact when
# it runs from a writable directory.
# ---------------------------------------------------------------------------
_AA_ENV_PATH = os.path.join(_HANDLER_DIR, "aa_env.py")


def _read_bytes(path: str) -> Optional[bytes]:
    try:
        with open(path, "rb") as fh:
            return fh.read()
    except OSError:
        return None


#: ``aa_env.py`` as it stood next to this file at import - the delivered copy.
_AA_ENV_BYTES = _read_bytes(_AA_ENV_PATH)


def _handler_delivery() -> Dict[str, Any]:
    """The identity of the handler that is running, for ``probe_capabilities``
    (and so the row, the run's ``env_capabilities`` event and the API)::

        {handler_version, handler_delivery: "overlay" | "baked",
         handler_delivery_reason, handler_release, handler_path}

    Read lazily from the two bootstrap variables - never written here."""
    marker = os.environ.get(HANDLER_DELIVERY_ENV, "")
    uri = os.environ.get(HANDLER_URI_ENV, "").rstrip("/")
    if marker == "overlay":
        delivery, reason = "overlay", None
    elif marker.startswith("baked:"):
        delivery, reason = "baked", marker[len("baked:"):] or "the bootstrap fell back to the image's copy"
    else:
        delivery, reason = "baked", "no bootstrap ran - the function's entry point runs the image's own copy"
    return {
        "handler_version": HANDLER_CONTRACT_VERSION,
        "handler_delivery": delivery,
        "handler_delivery_reason": reason,
        "handler_release": (uri.rsplit("/", 1)[-1] or None) if (delivery == "overlay" and uri) else None,
        "handler_path": os.path.abspath(__file__),
    }


def _ensure_overlay_intact() -> bool:
    """Restore ``aa_env.py`` next to this handler from the import-time copy
    when the on-disk bytes differ. Defence in depth: children never import
    from this directory any more (`_child_pythonpath` hands each exec a
    private copy), but under runtime delivery the directory is
    ``/tmp/aa_handler``, writable by the sandbox uid, and the pair that sits
    there should stay the pair that was delivered. ``/opt`` is read-only in
    Lambda, so for a baked pair this is a no-op. Returns True when a restore
    happened."""
    if _AA_ENV_BYTES is None:
        return False
    if _read_bytes(_AA_ENV_PATH) == _AA_ENV_BYTES:
        return False
    tmp = _AA_ENV_PATH + ".restore"
    try:
        with open(tmp, "wb") as fh:
            fh.write(_AA_ENV_BYTES)
        os.replace(tmp, _AA_ENV_PATH)
    except OSError as e:
        print(f"[HANDLER] aa_env.py next to the handler was modified and could not be restored: {e}")
        return False
    print("[HANDLER] aa_env.py next to the handler had been modified since import; restored the delivered copy")
    return True


def _safe_project_relpath(path: str) -> Optional[str]:
    """Normalise a project-relative path; ``None`` if it is absolute or escapes."""
    raw = str(path or "").strip()
    if not raw or raw.startswith("/") or "\\" in raw:
        return None
    norm = os.path.normpath(raw).replace(os.sep, "/")
    if norm.startswith("../") or norm == ".." or norm.startswith("/"):
        return None
    return norm


def _sync_prefix_down(bucket: str, prefix: str, workdir: str) -> int:
    """Download the ENTIRE session prefix into the workdir. Used on checkpoint
    resume so a fresh Lambda's /tmp is reconstituted exactly (manifest + shards +
    prior outputs) — the durable S3 workspace is the source of truth, never /tmp.
    Best-effort: a failure here degrades to recomputing from the last checkpoint.
    """
    if not bucket or not prefix:
        return 0
    s3 = _s3_client()
    count = 0
    try:
        paginator = s3.get_paginator("list_objects_v2")
        for page in paginator.paginate(Bucket=bucket, Prefix=prefix):
            for obj in page.get("Contents", []):
                key = obj["Key"]
                rel = key[len(prefix):].lstrip("/")
                if not rel or rel.endswith("/"):
                    continue
                dst = os.path.join(workdir, rel)
                os.makedirs(os.path.dirname(dst) or workdir, exist_ok=True)
                _unlink_quiet(dst)      # never write through a link into the mirror cache
                s3.download_file(bucket, key, dst)
                count += 1
        print(f"[HANDLER] resume: synced {count} file(s) down from s3://{bucket}/{prefix}")
    except Exception as e:
        print(f"[HANDLER] WARNING: resume sync-down failed: {e}")
    return count


def _int_env(name: str, default: int) -> int:
    try:
        return int(os.environ.get(name) or default)
    except (TypeError, ValueError):
        return default


class Budget(NamedTuple):
    """One exec action's time budget (see `_compute_budget`)."""
    hard_timeout: int              #: coreutils cap on the script, seconds
    deadline: float                #: epoch at which a cooperative script should yield (AA_DEADLINE_EPOCH)
    reserve: int                   #: seconds kept after the kill for the workdir → S3 upload
    grace: int                     #: seconds before the kill in which aa_env checkpoints
    bound: str                     #: "environment" (the Lambda clock decided) | "requested" (event["timeout"] decided)
    remaining: Optional[float]     #: Lambda remaining time at entry, seconds; None off-Lambda
    limit_seconds: Optional[int]   #: event["env_timeout_seconds"]: the environment's configured limit


def _positive_int(value: Any) -> Optional[int]:
    if isinstance(value, bool) or not isinstance(value, (int, float)) or value <= 0:
        return None
    return int(value)


def _compute_budget(event: Dict[str, Any], context: Any) -> Budget:
    """Return the `Budget` for an exec action.

    The hard timeout is the coreutils cap on the script; the deadline is the
    earlier soft point (exported as ``AA_DEADLINE_EPOCH``) at which a cooperative
    checkpointing script should yield, leaving grace to flush + upload.

    Derived from the LIVE Lambda remaining time so a checkpointed chunk is as
    large as safely possible, minus an upload reserve that SCALES with the
    function's own timeout (an environment is provisioned with its
    ``resource_config.timeout_seconds`` as the Lambda ``Timeout``, 1–900 s)::

        reserve      = min(upload_reserve_s [60], max(5, remaining // 4))
        hard_timeout = max(5, min(remaining - reserve, event["timeout"]))
        grace        = min(yield_grace_s [30],  max(5, hard_timeout // 3))
        deadline     = now + (hard_timeout - grace)

    so a 60 s environment gives code 45 s (reserve 15), a 120 s one 90 s
    (reserve 30) and a 900 s one 840 s (reserve 60) — the flat 60 s reserve
    used to leave every environment under ~65 s with the 5 s floor.
    ``event["timeout"]`` (the caller's ceiling — the coder's
    ``CODER_ENV_INVOKE_TIMEOUT_SECONDS``, code-interpreter's request timeout,
    both clamped to the environment's limit) caps it further; ``bound`` records
    which of the two won so a kill can be reported as the environment's limit
    or as the caller's request. Falls back to the event timeout when no Lambda
    context is available (local/tests).
    """
    def _param(key: str, env_name: str, default: int) -> int:
        v = _positive_int(event.get(key))
        return v if v is not None else _int_env(env_name, default)

    reserve_cap = _param("upload_reserve_s", "AA_UPLOAD_RESERVE_S", _DEFAULT_UPLOAD_RESERVE_S)
    grace_cap = _param("yield_grace_s", "AA_YIELD_GRACE_S", _DEFAULT_YIELD_GRACE_S)
    requested = _positive_int(event.get("timeout"))
    limit = _positive_int(event.get("env_timeout_seconds"))
    remaining: Optional[float] = None
    try:
        if context is not None and hasattr(context, "get_remaining_time_in_millis"):
            remaining = context.get_remaining_time_in_millis() / 1000.0
    except Exception:
        remaining = None

    bound = "requested"
    if remaining is not None:
        reserve = min(reserve_cap, max(_MIN_UPLOAD_RESERVE_S, int(remaining // 4)))
        hard = remaining - reserve
        bound = "environment"
        if requested is not None and requested < hard:
            hard = float(requested)
            bound = "requested"
    elif requested is not None:
        reserve = min(reserve_cap, max(_MIN_UPLOAD_RESERVE_S, requested // 4))
        hard = float(requested)
    else:
        reserve = reserve_cap
        hard = 300.0
    hard_int = int(max(_MIN_HARD_TIMEOUT_S, hard))
    grace = min(grace_cap, max(_MIN_YIELD_GRACE_S, hard_int // 3))
    deadline = time.time() + max(1, hard_int - grace)
    return Budget(hard_int, deadline, int(reserve), int(grace), bound, remaining, limit)


def _coerce_budget(budget: Any) -> Budget:
    """`_exec_python`/`_exec_shell` take a `Budget`; a plain int (direct
    callers, tests, local harnesses) is a requested-bound budget of that many
    seconds."""
    if isinstance(budget, Budget):
        return budget
    return _compute_budget({"timeout": int(budget)}, None)


def _wall_hit(returncode: int, elapsed: float, hard_timeout: int) -> bool:
    """True when the wall stopped the child. coreutils ``timeout`` exits 124 on
    its SIGTERM; after ``-k`` it SIGKILLs the whole process group (itself
    included) and the parent sees -9 (137 through a shell) — counted as the
    wall only when the clock says so, so an early SIGKILL (OOM) stays what it
    is."""
    if returncode == EXIT_TIMEOUT:
        return True
    return returncode in (-9, 137) and elapsed >= max(0, hard_timeout - 1)


def _wall_report(budget: Budget, elapsed: float) -> Dict[str, Any]:
    """The structured fields an exec response carries when the wall was hit.
    code-interpreter (`lambda_backend._map_handler_timeout`) and the coder
    (`coder_worker.run_in_env`) turn ``error_code == "environment_timeout"``
    into the user-facing "‹environment› exceeded its N-second limit" copy;
    ``bound == "requested"`` means the caller's own ``timeout`` ceiling, not
    the environment, stopped the code."""
    hard, reserve = budget.hard_timeout, budget.reserve
    if budget.bound == "requested":
        message = f"Execution stopped after {hard} s: the {hard}-second timeout requested for this run."
    elif budget.limit_seconds:
        message = (f"Execution stopped after {hard} s: the environment's {budget.limit_seconds}-second limit "
                   f"leaves {hard} s for code ({reserve} s is reserved for saving results).")
    else:
        message = (f"Execution stopped after {hard} s: the environment's execution-time limit leaves "
                   f"{hard} s for code ({reserve} s is reserved for saving results).")
    return {
        "error_type": "infrastructure_error",
        "error_code": "environment_timeout",
        "error": message,
        "limit_seconds": budget.limit_seconds,
        "budget_seconds": hard,
        "reserve_seconds": reserve,
        "bound": budget.bound,
        "elapsed_seconds": round(elapsed, 1),
    }


def _shell_argv(login_shell: bool = False) -> List[str]:
    """The shell `_exec_shell` runs the coder's command under: ``bash`` when the
    image has it, else ``sh``. bash gets ``--noprofile --norc`` (no rc file -
    ``~/.bashrc`` is under a HOME the child can write to - is ever
    sourced; a non-interactive ``bash -c`` reads none by itself, only
    ``$BASH_ENV``, which `_exec_env` strips) unless the invocation asks for a
    login shell (``login_shell: true``), which reads its own scope's profile.
    The connector variables never depended on an rc file: they are in the
    environment the handler spawns the shell with.

    Every connector variable is named ``CONN_{connector_id}_{KEY}`` and a
    connector id is a UUID, so the name carries hyphens and is not a valid
    shell identifier. Debian's ``/bin/sh`` is dash - the base of every shipped
    image - and dash discards such entries from its environment at startup,
    so a command run under ``sh -c`` and every process it spawned (the coder's
    ``python file.py``) saw NONE of the injected credentials, while an
    ``exec_python`` child - no shell in between - saw them all (the handler
    logged the variables it had injected while ``env | grep -c ^CONN_`` under
    ``sh -c`` printed 0). bash keeps invalid-name entries in the environment
    it passes to its children (they are simply not addressable as
    ``$CONN_...``), so under bash a script run through the shell reads them
    from ``os.environ`` exactly as an ``exec_python`` script does."""
    if not shutil.which("bash"):
        # dash: a non-interactive sh reads no rc file ($ENV is interactive-only).
        return ["sh", "-l", "-c"] if login_shell else ["sh", "-c"]
    if login_shell:
        return ["bash", "-l", "-c"]
    return ["bash", "--noprofile", "--norc", "-c"]


def _run_under_wall(argv: List[str], workdir: str, env: Dict[str, str], budget: Budget,
                    clock=time.monotonic) -> Tuple[Any, Optional[Dict[str, Any]]]:
    """Run ``argv`` under ``timeout -k <kill_after> <hard_timeout>`` with cwd =
    the project root, capturing stdout/stderr. Returns ``(result, wall)`` where
    ``wall`` is `_wall_report` when the wall stopped the child (its
    ``returncode`` normalised to `EXIT_TIMEOUT`) and ``None`` otherwise.

    Output the child flushed before the kill IS in ``result.stdout`` — the
    pipe survives the signal (verified on coreutils 9 / python:3.12-slim); only
    output still sitting in the child's own buffer is lost, which is why
    `_exec_python` runs with ``PYTHONUNBUFFERED=1``.
    """
    kill_after = max(1, min(_MAX_KILL_AFTER_S, budget.reserve // 3))
    t0 = clock()
    result = subprocess.run(
        ["timeout", "-k", str(kill_after), str(budget.hard_timeout), *argv],
        capture_output=True,
        cwd=workdir,
        env=env,
        preexec_fn=_no_core_dumps,
    )
    elapsed = clock() - t0
    wall = None
    if _wall_hit(result.returncode, elapsed, budget.hard_timeout):
        wall = _wall_report(budget, elapsed)
        result.returncode = EXIT_TIMEOUT
    return result, wall


# ---------------------------------------------------------------------------
# Capability probe (v1.0.120): what document tooling this image
# actually has. Cheap version checks with short timeouts; a missing tool is
# ``False``, never an exception, so the coder's prompt can say "missing:
# node, chromium" instead of the coder discovering it mid-deliverable.
# ---------------------------------------------------------------------------

_VERSION_RE = re.compile(r"(\d+(?:\.\d+)+)")


def _probe_run(args: List[str], env: Optional[Dict[str, str]] = None,
               timeout: Optional[int] = None) -> Optional[subprocess.CompletedProcess]:
    """Run one probe command; ``None`` when the tool is absent or hangs.
    ``timeout`` defaults to `PROBE_TIMEOUT_S`; the big binaries pass the
    lock's browser ceiling instead (a cold start pages them in slowly)."""
    try:
        return subprocess.run(args, capture_output=True, text=True,
                              timeout=timeout if timeout else PROBE_TIMEOUT_S,
                              env=env if env is not None else _child_env({}),
                              preexec_fn=_no_core_dumps)
    except (FileNotFoundError, PermissionError, subprocess.TimeoutExpired, OSError):
        return None


def _probe_version(args: List[str], env: Optional[Dict[str, str]] = None, require_rc0: bool = True,
                   timeout: Optional[int] = None):
    """First dotted version number in stdout+stderr, or ``False``."""
    r = _probe_run(args, env, timeout)
    if r is None or (require_rc0 and r.returncode != 0):
        return False
    m = _VERSION_RE.search((r.stdout or "") + "\n" + (r.stderr or ""))
    return m.group(1) if m else False


#: The lock check that vouches for each probed capability when the tool's
#: own version probe came back False. On a COLD container-image
#: Lambda the first touch of soffice/pandoc/weasyprint/chromium pages in a
#: multi-hundred-MB binary and exceeds `PROBE_TIMEOUT_S`; the same probe on a
#: warm instance answers in a second. The checks run first (warming the
#: binaries) and a passing check is proof the tool is there, so the boolean
#: is derived from it rather than reported False for the whole session.
#: Lock checks that touch a big, lazily paged-in binary; they run with the
#: browser ceiling (`_run_toolkit_checks`) for the same cold-start reason.
_SLOW_CHECK_NAMES = ("soffice", "pandoc", "weasyprint")
_CHECK_FOR_CAPABILITY = {
    "libreoffice": "soffice", "poppler": "pdftoppm", "pandoc": "pandoc", "weasyprint": "weasyprint",
    "node": "node", "pptxgenjs": "pptxgenjs", "graphviz": "graphviz", "mermaid": "mmdc",
    # mermaid-render launches the toolkit's Chromium; passing means the browser runs.
    "chromium": "mermaid-render",
}
def _derive_from_checks(caps: Dict[str, Any], checks: List[Dict[str, Any]], outputs: Dict[str, str]) -> None:
    """For every capability still False whose vouching check passed
    (``ok is True`` — never a skipped or failed one), set the version parsed
    from the check's output, or True when the output carries no version."""
    by_name = {c.get("name"): c for c in checks if isinstance(c, dict)}
    for key, check_name in _CHECK_FOR_CAPABILITY.items():
        if caps.get(key) is not False:
            continue
        check = by_name.get(check_name)
        if not check or check.get("ok") is not True:
            continue
        m = _VERSION_RE.search(outputs.get(check_name) or "")
        caps[key] = m.group(1) if m else True


def _read_toolkit_lock() -> Tuple[Optional[str], Dict[str, Any]]:
    """``(toolkit_version, lock_dict)`` from `TOOLKIT_LOCK_PATH`. The lock is
    the JSON the document toolkit image ships (``toolkit_version``,
    ``check_semantics``, ``checks[]``, package lists — unknown keys ignored);
    a bare version line or ``version: x`` text is accepted for the version
    alone. ``(None, {})`` when the file is absent or empty."""
    try:
        with open(TOOLKIT_LOCK_PATH, "r", encoding="utf-8") as f:
            text = f.read().strip()
    except OSError:
        return None, {}
    if not text:
        return None, {}
    try:
        data = json.loads(text)
    except ValueError:
        data = None
    if isinstance(data, dict):
        version = data.get("toolkit_version") or data.get("version")
        return (str(version) if version else None), data
    m = re.search(r"version\s*[:=]\s*\"?([^\s\"]+)", text, re.IGNORECASE)
    if m:
        return m.group(1), {}
    return (text.splitlines()[0].strip() or None), {}


def _version_tuple(version: Optional[str]) -> Optional[Tuple[int, ...]]:
    """``"1.10.2"`` -> ``(1, 10, 2)``; None for anything non-numeric."""
    if not version:
        return None
    m = re.match(r"\s*v?(\d+(?:\.\d+)*)", str(version))
    if not m:
        return None
    return tuple(int(x) for x in m.group(1).split("."))


def _toolkit_outdated(installed: Optional[str]) -> bool:
    """True when the installed lock's version is numerically below
    `HANDLER_TOOLKIT_VERSION`. An absent lock (None) or a non-numeric version
    is not "outdated" — it is already reported as missing/odd elsewhere."""
    have, want = _version_tuple(installed), _version_tuple(HANDLER_TOOLKIT_VERSION)
    if have is None or want is None:
        return False
    return have < want


def _int_or(value: Any, default: int) -> int:
    try:
        v = int(value)
        return v if v > 0 else default
    except (TypeError, ValueError):
        return default


def _machine() -> str:
    try:
        return os.uname().machine
    except Exception:  # noqa: BLE001
        return ""


def _browser_emulated() -> Optional[str]:
    """Reason string when Chromium cannot run here (x86_64 under QEMU
    user-mode emulation: no x86 ``flags:`` line in /proc/cpuinfo), else None."""
    if _machine() != "x86_64":
        return None
    try:
        with open(CPUINFO_PATH, "r", encoding="utf-8", errors="replace") as f:
            for line in f:
                if line.startswith("flags"):
                    return None
    except OSError:
        return None
    return "needs_browser: skipped under QEMU user-mode emulation (x86_64 with no cpuinfo flags line)"


def _run_toolkit_checks(lock: Dict[str, Any], env: Dict[str, str],
                        outputs: Optional[Dict[str, str]] = None) -> List[Dict[str, Any]]:
    """The lock's ``checks[]`` exactly per its ``check_semantics``: each
    ``cmd`` via ``/bin/sh -c`` (offline, runtime user), pass = exit 0 AND
    ``re.search(expect, stdout + stderr)``; ``needs_browser`` checks are
    ``skipped`` (ok=None) with a reason under emulation. Per-check timeout:
    the lock's ``check_timeout_s`` (default `PROBE_TIMEOUT_S`), or for
    ``needs_browser`` checks its ``browser_check_timeout_s`` (default
    `BROWSER_CHECK_TIMEOUT_S` — a cold Chromium in Lambda is slow). Malformed
    entries are reported, never raised. ``outputs``, when given,
    collects each executed check's combined stdout+stderr by name so
    `_derive_from_checks` can read a version out of it; the reported entries
    keep their ``{name, ok, skipped[, reason]}`` shape."""
    results: List[Dict[str, Any]] = []
    checks = lock.get("checks")
    if not isinstance(checks, list):
        return results
    emulated = _browser_emulated()
    default_timeout = _int_or(lock.get("check_timeout_s"), PROBE_TIMEOUT_S)
    browser_timeout = _int_or(lock.get("browser_check_timeout_s"), BROWSER_CHECK_TIMEOUT_S)
    for i, check in enumerate(checks, start=1):
        if not isinstance(check, dict):
            continue
        name = str(check.get("name") or f"check-{i}")
        entry: Dict[str, Any] = {"name": name, "ok": None, "skipped": False}
        if check.get("needs_browser") and emulated:
            entry["skipped"] = True
            entry["reason"] = emulated
            results.append(entry)
            continue
        cmd = check.get("cmd")
        expect = check.get("expect")
        if not isinstance(cmd, str) or not cmd.strip() or not isinstance(expect, str):
            entry["ok"] = False
            entry["reason"] = "malformed check: cmd and expect are required strings"
            results.append(entry)
            continue
        try:
            pattern = re.compile(expect)
        except re.error as e:
            entry["ok"] = False
            entry["reason"] = f"invalid expect pattern: {e}"
            results.append(entry)
            continue
        # With the checks running first, a big binary's check IS its
        # first touch on a cold instance — the very touch that blew the 15 s
        # probe ceiling. Those checks get the browser ceiling too (deliberate
        # deviation from the lock's check_semantics prose: a longer ceiling
        # cannot turn a passing check into a failure).
        slow = check.get("needs_browser") or name in _SLOW_CHECK_NAMES
        timeout = browser_timeout if slow else default_timeout
        try:
            r = subprocess.run(["/bin/sh", "-c", cmd], capture_output=True, text=True,
                               timeout=timeout, env=env, preexec_fn=_no_core_dumps)
        except subprocess.TimeoutExpired:
            entry["ok"] = False
            entry["reason"] = f"timeout after {timeout}s"
            results.append(entry)
            continue
        except OSError as e:
            entry["ok"] = False
            entry["reason"] = f"could not run: {e}"
            results.append(entry)
            continue
        combined = (r.stdout or "") + (r.stderr or "")
        if outputs is not None:
            outputs[name] = combined[:4000]
        matched = pattern.search(combined) is not None
        if r.returncode == 0 and matched:
            entry["ok"] = True
        else:
            entry["ok"] = False
            why = []
            if r.returncode != 0:
                why.append(f"exit {r.returncode}")
            if not matched:
                why.append(f"expect {expect!r} not found")
            tail = combined.strip().splitlines()[-1][-200:] if combined.strip() else ""
            entry["reason"] = "; ".join(why) + (f" (last line: {tail})" if tail else "")
        results.append(entry)
    return results


#: Resolves the pptxgenjs package root from ``require.resolve`` (its main
#: entry, ``dist/pptxgen.cjs.js``) and reads the nearest package.json above it.
#: ``require('pptxgenjs/package.json')`` is NOT an exported subpath of that
#: package (ERR_PACKAGE_PATH_NOT_EXPORTED), which is why the probe said
#: ``pptxgenjs: false`` on our own image.
_PPTXGENJS_VERSION_JS = (
    "const fs=require('fs'),path=require('path');"
    "let d=path.dirname(require.resolve('pptxgenjs'));"
    "for(;;){const pj=path.join(d,'package.json');"
    "if(fs.existsSync(pj)){console.log(JSON.parse(fs.readFileSync(pj,'utf8')).version);break;}"
    "const up=path.dirname(d);if(up===d){process.exit(1);}d=up;}"
)


def _probe_chromium(lock: Dict[str, Any], env: Dict[str, str], timeout: Optional[int] = None):
    """The toolkit's browser is the lock's ``chromium.executable_link``
    (``/opt/alphaagent/chrome``), never a system ``chromium`` binary — probe
    that first, launched with the lock's Lambda-safe ``args`` (``--version``
    still needs a launchable configuration). System names remain a fallback
    for a customer image that installed a distro Chromium instead."""
    chromium = lock.get("chromium") if isinstance(lock.get("chromium"), dict) else {}
    link = chromium.get("executable_link")
    args = chromium.get("args") if isinstance(chromium.get("args"), list) else []
    if isinstance(link, str) and link:
        v = _probe_version([link, *[str(a) for a in args], "--version"], env, timeout=timeout)
        if v:
            return v
    for exe in ("chromium", "chromium-browser", "google-chrome"):
        v = _probe_version([exe, "--version"], env, timeout=timeout)
        if v:
            return v
    return False


#: Run as a real child of the handler process by the probe: does a same-uid
#: sibling of the sandbox still read this process's start-up environment?
_CHILD_ENVIRON_PROBE = (
    "import os,json\n"
    "p=os.getppid();r={'ppid_environ':None,'credential_names':None,'error':None,'maps':None,'mem':None}\n"
    "try:\n"
    "  d=open('/proc/%d/environ' % p,'rb').read()\n"
    "  r['ppid_environ']='readable'\n"
    "  r['credential_names']=sorted(k for k in ('AWS_ACCESS_KEY_ID','AWS_SECRET_ACCESS_KEY','AWS_SESSION_TOKEN','AWS_SECURITY_TOKEN') if any(e.startswith(k.encode()+b'=') for e in d.split(b'\\0')))\n"
    "except Exception as e:\n"
    "  r['ppid_environ']='denied' if isinstance(e,PermissionError) else 'unavailable'; r['error']=type(e).__name__\n"
    "for f in ('maps','mem'):\n"
    "  try:\n"
    "    open('/proc/%d/%s' % (p,f),'rb').close(); r[f]='opened'\n"
    "  except Exception as e:\n"
    "    r[f]='denied' if isinstance(e,PermissionError) else type(e).__name__\n"
    "print(json.dumps(r))"
)
_DENY_ALL_POLICY = json.dumps({"Version": "2012-10-17",
                               "Statement": [{"Effect": "Deny", "Action": "*", "Resource": "*"}]})


def _check_scoped_role() -> Dict[str, Any]:
    """Probe-time check that the sandbox role exists and trusts this role: one
    AssumeRole with a deny-everything session policy, credentials discarded.
    Reported as ``credential_isolation.scoped_role`` on every environment
    row, so the role's rollout is visible rather than assumed. Only under the
    real runtime - a local probe never reaches STS."""
    out: Dict[str, Any] = {"arn": None, "assumable": None, "error": None}
    if not os.environ.get("AWS_LAMBDA_RUNTIME_API"):
        out["error"] = "not checked outside the Lambda runtime"
        return out
    arn = _scoped_role_arn()
    out["arn"] = arn
    if not arn:
        out["error"] = f"no scoped role ARN: not derivable from this role's name and {SANDBOX_ROLE_ARN_ENV} unset"
        return out
    try:
        _sts_client().assume_role(RoleArn=arn, RoleSessionName="aa-env-probe",
                                  Policy=_DENY_ALL_POLICY, DurationSeconds=900)
        out["assumable"] = True
    except Exception as exc:
        out["assumable"] = False
        out["error"] = _error_code(exc)
    return out


def _isolation_report(env: Dict[str, str]) -> Dict[str, Any]:
    """`_ISOLATION` (what import did, re-measured at this invocation) plus two
    live measurements: a real child of this process trying
    ``/proc/<ppid>/environ``, ``maps`` and ``mem`` (Linux) - the READ and
    ATTACH gates a hostile child would have to pass - and whether the sandbox
    role can be assumed (`_check_scoped_role`)."""
    report: Dict[str, Any] = dict(_ISOLATION)
    report["child_can_read_parent_environ"] = None
    report["child_sees_credential_names"] = None
    report["child_can_read_parent_maps"] = None
    report["child_can_open_parent_mem"] = None
    if sys.platform.startswith("linux"):
        r = _probe_run([sys.executable, "-c", _CHILD_ENVIRON_PROBE], env)
        if r is not None and r.returncode == 0:
            try:
                data = json.loads((r.stdout or "").strip().splitlines()[-1])
                report["child_can_read_parent_environ"] = data.get("ppid_environ") == "readable"
                report["child_sees_credential_names"] = data.get("credential_names")
                report["child_probe_error"] = data.get("error")
                report["child_can_read_parent_maps"] = data.get("maps") == "opened"
                report["child_can_open_parent_mem"] = data.get("mem") == "opened"
            except (ValueError, IndexError):
                report["child_probe_error"] = "unparseable child output"
    report["scoped_role"] = _check_scoped_role()
    report.setdefault("credential_scope", None)
    return report


def probe_capabilities() -> Dict[str, Any]:
    """The capability map: ``{toolkit_version?, libreoffice, poppler, pandoc,
    weasyprint, node, pptxgenjs, chromium, graphviz, mermaid, fonts:[...],
    checked_at}`` — a version string (or True) when present, False when not.
    Plus (v1.0.120) the lock's ``checks`` / ``toolkit_ok`` and (v1.0.121)
    ``toolkit_outdated`` / ``toolkit_expected_version`` and the /tmp gauge
    ``tmp_used_bytes`` / ``tmp_total_bytes``; (1.0.154) ``credential_isolation``
    (`_isolation_report`)."""
    env = _child_env({})
    # soffice writes a profile on first run; without a writable HOME + a
    # writable UserInstallation it can exit 0 having printed nothing, or hang.
    lo_profile = os.path.join(TMP_ROOT, "lo_probe")
    try:
        os.makedirs(lo_profile, exist_ok=True)
    except OSError:
        pass
    toolkit_version, lock = _read_toolkit_lock()
    # The lock's checks run FIRST. On a cold container-image instance
    # they are what pages soffice/pandoc/weasyprint/chromium in, so the version
    # probes below meet warm binaries; their outputs are kept (by name) so a
    # probe that still misses can be derived from its passing check.
    check_outputs: Dict[str, str] = {}
    checks = _run_toolkit_checks(lock, env, outputs=check_outputs) if lock else []
    slow_timeout = _int_or(lock.get("browser_check_timeout_s"), BROWSER_CHECK_TIMEOUT_S) if lock else BROWSER_CHECK_TIMEOUT_S
    caps: Dict[str, Any] = {
        "toolkit_version": toolkit_version,
        "toolkit_expected_version": HANDLER_TOOLKIT_VERSION,
        "toolkit_outdated": _toolkit_outdated(toolkit_version),
        "libreoffice": _probe_version(
            ["soffice", "--headless", "--norestore", "--nolockcheck",
             f"-env:UserInstallation=file://{lo_profile}", "--version"], env, timeout=slow_timeout),
        # pdftoppm -v prints to stderr and exits 0 (99 on very old builds).
        "poppler": _probe_version(["pdftoppm", "-v"], env, require_rc0=False),
        "pandoc": _probe_version(["pandoc", "-v"], env, timeout=slow_timeout),
        "weasyprint": _probe_version(
            [sys.executable, "-c", "import weasyprint; print(weasyprint.__version__)"], env, timeout=slow_timeout),
        "node": _probe_version(["node", "-v"], env),
        "pptxgenjs": _probe_version(["node", "-e", _PPTXGENJS_VERSION_JS], env),
        "chromium": _probe_chromium(lock, env, timeout=slow_timeout),
        "graphviz": _probe_version(["dot", "-V"], env, require_rc0=False),
        "mermaid": _probe_version(["mmdc", "-V"], env),
        "fonts": [],
        "checked_at": datetime.now(timezone.utc).isoformat(timespec="seconds"),
    }
    r = _probe_run(["fc-list", ":", "family"], env)
    if r is not None and r.returncode == 0:
        fams = set()
        for line in (r.stdout or "").splitlines():
            fam = line.split(",", 1)[0].strip()
            if fam:
                fams.add(fam)
        caps["fonts"] = sorted(fams)
    # The toolkit's own contract: the lock's checks (run above, per its
    # check_semantics). Additive to the boolean capability map.
    caps["checks"] = checks
    caps["toolkit_ok"] = bool(lock) and all(c["ok"] is True for c in caps["checks"] if not c["skipped"])
    _derive_from_checks(caps, checks, check_outputs)
    # The credential-isolation layers as this process and a real child see
    # them (1.0.154): re-exec, /proc snapshot, dumpable flag, ptrace_scope,
    # whether a child can read /proc/<ppid>/environ, the scoped role.
    caps["credential_isolation"] = _isolation_report(env)
    # Which handler answered: contract version, delivery path, release.
    caps.update(_handler_delivery())
    try:
        usage = shutil.disk_usage(TMP_ROOT)
        caps["tmp_used_bytes"], caps["tmp_total_bytes"] = int(usage.used), int(usage.total)
    except OSError:
        caps["tmp_used_bytes"], caps["tmp_total_bytes"] = 0, 0
    return caps


def handler(event: Dict[str, Any], context: Any) -> Dict[str, Any]:
    """Lambda entrypoint.

    Event schema::

        {
            "action": "exec_python" | "exec_shell" | "write_file" | "read_file" | "list_files" | "fetch_object"
                      | "probe_capabilities",
            "s3_bucket": "alphaagent-workspaces-...",
            "s3_prefix": "users/{uid}/projects/conversations/{slug}/",   # (legacy: workspaces/sessions/{id}/workspace/)
            "timeout": 300,           # seconds: the caller's CEILING on one exec (the live budget can be lower)
            "env_timeout_seconds": 60, # the environment's configured limit (its Lambda Timeout); named in the timeout error
            "script_path": "...",     # project-relative script to run (exec_python; preferred)
            "code": "...",            # inline Python source (exec_python; lands in code/scratch/)
            "script_name": "...",     # filename for inline code (exec_python)
            "ephemeral": bool,        # inline code lands in a hidden dir, never uploaded
            "command": "...",         # Shell command string (exec_shell)
            "login_shell": bool,      # exec_shell: run a login shell (reads the scope's own profile);
                                      # default false: bash --noprofile --norc
            "sync_dirs": ["code", "nodes", "inputs"],   # project dirs to mirror down first
            "sync_max_bytes": N, "sync_skip_over_bytes": N,   # optional size guards
            "prefetch": [...],        # extra objects to hydrate (relative or full keys)
            "preserve_paths": bool,   # prefetch keeps relative paths (default: basename)
            "resume": bool,           # checkpoint continuation: full prefix sync first
            "path": "...",            # Relative path (file operations)
            "content": "...",         # File content (write_file)
            "agent_id": "...",        # For credential lookup
            "connector_ids": [...]    # Connectors to inject
        }

    exec responses carry ``uploaded`` (project-relative paths the run created
    or changed) and, when ``sync_dirs`` was given, a ``sync`` summary. When the
    wall stops the script the response is ``exit_code`` 124 with the captured
    ``stdout``/``stderr`` PLUS ``error_type: "infrastructure_error"``,
    ``error_code: "environment_timeout"``, a readable ``error``,
    ``limit_seconds`` (the environment's limit), ``budget_seconds`` (what the
    code got), ``reserve_seconds`` and ``bound`` (``"environment"`` |
    ``"requested"``) — see `_compute_budget` / `_wall_report`.
    """
    action = event.get("action", "")
    s3_bucket = event.get("s3_bucket", "")
    s3_prefix = event.get("s3_prefix", "")
    timeout = event.get("timeout", 300)

    print(f"[HANDLER] entry: action={action}, bucket={s3_bucket}, prefix={s3_prefix}, "
          f"timeout={timeout}, resume={bool(event.get('resume'))}")
    delivery = _handler_delivery()
    print(f"[HANDLER] uid={os.getuid()}, gid={os.getgid()}, handler_version={delivery['handler_version']}, "
          f"delivery={delivery['handler_delivery']} ({delivery['handler_delivery_reason'] or delivery['handler_release']}) "
          f"from {delivery['handler_path']}")
    # A crashed child (Chrome under the 1.0.0 flags) leaves ~82 MB cores under
    # /tmp; on a warm instance six of those fill the disk and every later
    # action dies with ENOSPC. Sweep first, and make /tmp usage visible in
    # CloudWatch on every invocation.
    removed, tmp_used, tmp_total = _sweep_core_dumps()
    print(f"[HANDLER] tmp sweep: core dumps removed={removed}, tmp_used_bytes={tmp_used}, "
          f"tmp_total_bytes={tmp_total}")

    # FIRST, before anything else looks at the environment, drop
    # whatever a previous invocation of this warm container could have left
    # in the process environment - every action, not only exec.
    purged = _purge_previous_invocation_env()
    print(f"[HANDLER] env: purged {purged} per-invocation variable(s) left by an earlier invocation")

    # Credential isolation, layer 2, every invocation: re-assert and
    # re-measure this process's non-dumpable state (a WARNING if it is not).
    _reassert_non_dumpable()
    print(f"[HANDLER] credential isolation: exec_image={_ISOLATION.get('exec_image')} "
          f"non_dumpable={_ISOLATION.get('non_dumpable')} prctl={_ISOLATION.get('prctl')}")

    # This invocation's reach (credential isolation, layer 3): every S3 and
    # Secrets Manager client built until this returns comes from the session
    # scoped to it - `_s3_client` / `_secrets_client` refuse otherwise.
    global _CURRENT_SCOPE, _INVOCATION_SESSION
    _CURRENT_SCOPE = _scope_from_event(event)
    _INVOCATION_SESSION = None
    print(f"[HANDLER] scope: {_CURRENT_SCOPE.describe()}")
    try:
        return _dispatch(event, context, action, s3_bucket, s3_prefix)
    finally:
        _CURRENT_SCOPE = None
        _INVOCATION_SESSION = None


def _dispatch(event: Dict[str, Any], context: Any, action: str, s3_bucket: str, s3_prefix: str) -> Dict[str, Any]:
    """`handler`'s body once the environment is purged and the scope is set."""
    # This invocation's variables. They reach the child through `_child_env`
    # and are never written to os.environ.
    invocation_env: Dict[str, str] = {}
    budget: Optional[Budget] = None
    if action in ("exec_python", "exec_shell"):
        # Derive the real per-invocation budget from the LIVE Lambda clock and
        # publish the soft deadline so a cooperative checkpointing script (aa_env)
        # can yield before the hard kill. A non-checkpointed script simply runs
        # under the hard timeout as before.
        budget = _compute_budget(event, context)
        invocation_env["AA_DEADLINE_EPOCH"] = str(budget.deadline)
        print(f"[HANDLER] budget: hard_timeout={budget.hard_timeout}s reserve={budget.reserve}s "
              f"grace={budget.grace}s bound={budget.bound} limit={budget.limit_seconds} "
              f"remaining={budget.remaining} deadline_epoch={budget.deadline:.0f}")
        connector_env = _load_connector_credentials(
            event.get("agent_id", ""),
            event.get("connector_ids", []),
        )
        invocation_env.update(connector_env)
        print(f"[HANDLER] Injected {len(connector_env)} connector env vars")
    _track_invocation_env(invocation_env)

    try:
        if action == "exec_python":
            return _exec_python(event, s3_bucket, s3_prefix, budget, invocation_env)
        elif action == "exec_shell":
            return _exec_shell(event, s3_bucket, s3_prefix, budget, invocation_env)
        elif action == "write_file":
            return _write_file(event, s3_bucket, s3_prefix)
        elif action == "read_file":
            return _read_file(event, s3_bucket, s3_prefix)
        elif action == "list_files":
            return _list_files(event, s3_bucket, s3_prefix)
        elif action == "fetch_object":
            return _fetch_object(event, s3_bucket, s3_prefix)
        elif action == "probe_capabilities":
            return {"exit_code": 0, "stdout": "", "stderr": "", "capabilities": probe_capabilities()}
        else:
            return {"exit_code": 1, "stdout": "", "stderr": f"Unknown action: {action}"}
    except Exception as exc:
        print(f"[HANDLER] Unhandled exception: {traceback.format_exc()}")
        return {
            "exit_code": 1,
            "stdout": "",
            "stderr": f"Handler error: {traceback.format_exc()}",
        }


def _prefetch_files(
    prefetch: List[str],
    s3_bucket: str,
    s3_prefix: str,
    workdir: str,
    preserve_paths: bool = False,
) -> None:
    """Download declared inputs into the workdir before running code.

    Each entry is either a full S3 key (``users/``, ``workspaces/`` or
    ``conversations/``) or a path relative to ``s3_prefix``. With
    ``preserve_paths`` a relative entry lands at its own relative path
    (``scripts/x.py`` -> ``{workdir}/scripts/x.py``) — the coder's
    ``run_in_env`` always asks for this. The default still flattens to the
    basename because workflow's Script node executor runs ``python
    <basename>`` after prefetching ``scripts/<name>.py``; a full key always
    lands at its basename (there is no project-relative path to preserve).
    Best-effort: a missing/failed prefetch never aborts the run.
    """
    if not prefetch or not s3_bucket:
        return
    s3 = _s3_client()
    for entry in prefetch:
        try:
            entry = str(entry)
            if entry.startswith(_FULL_KEY_PREFIXES):
                key = entry
                dst = os.path.join(workdir, os.path.basename(key))
            else:
                rel = entry.lstrip("/")
                key = f"{s3_prefix.rstrip('/')}/{rel}"
                safe_rel = _safe_project_relpath(rel) if preserve_paths else None
                dst = (
                    os.path.join(workdir, *safe_rel.split("/"))
                    if safe_rel else os.path.join(workdir, os.path.basename(key))
                )
            os.makedirs(os.path.dirname(dst) or workdir, exist_ok=True)
            _unlink_quiet(dst)          # never write through a link into the mirror cache
            s3.download_file(s3_bucket, key, dst)
            print(f"[HANDLER] prefetched {key} -> {dst}")
        except Exception as e:
            print(f"[HANDLER] WARNING: prefetch failed for {entry}: {e}")


def _fetch_object(event: Dict, s3_bucket: str, s3_prefix: str) -> Dict[str, Any]:
    """load_artifact's hydrate step: download ONE S3 object into this
    session's live workdir. Runs in THIS handler process - the trusted
    parent, same credential context `_prefetch_files` and
    `_load_connector_credentials` already use - never as code handed to
    the sandboxed `_exec_python`/`_exec_shell` child, which no longer
    carries this Lambda's own role credentials and so has nothing
    for a `boto3.client('s3')` call there to authenticate with.
    """
    s3_key = event.get("s3_key", "")
    # basename, not the raw caller-supplied value: `target_filename` is a
    # single path segment for "where in workdir does this land", not a path
    # itself - a `../`-laden value must not be able to write outside workdir.
    target_filename = os.path.basename(event.get("target_filename") or os.path.basename(s3_key))
    if not s3_key or not target_filename:
        return {"exit_code": 1, "stdout": "", "stderr": "fetch_object requires s3_key"}

    workdir = _get_session_workdir(s3_prefix)
    try:
        dst = os.path.join(workdir, target_filename)
        try:
            _s3_client().download_file(s3_bucket, s3_key, dst)
            size = os.path.getsize(dst)
        except Exception as exc:
            print(f"[HANDLER] fetch_object failed: {exc}")
            return {"exit_code": 1, "stdout": "", "stderr": f"fetch_object failed: {exc}"}

        print(f"[HANDLER] fetch_object: {s3_key} -> {dst} ({size} bytes)")
        # Same "upload back to this session's S3 prefix" step _exec_python takes
        # after running code, so a subsequent read_workspace_file / list_files
        # sees the fetched object without requiring a run_in_env first. The
        # directory is this call's own fresh one, so exactly the fetched object
        # goes up.
        if s3_bucket and s3_prefix:
            _upload_workdir_files(s3_bucket, s3_prefix, workdir)
        return {"exit_code": 0, "stdout": f"LOADED {size} {target_filename}", "stderr": ""}
    finally:
        _release_workdir(workdir)


def _get_session_workdir(s3_prefix: str) -> str:
    """A FRESH working directory for this one exec: ``mkdtemp``
    under /tmp, mode 0700, random name (`WORKDIR_PREFIX`), removed by
    `_release_workdir` when the exec ends. It used to be a stable per-prefix
    path derived from ``s3_prefix`` - predictable to every child the warm
    container served, and persistent between execs. ``s3_prefix`` stays in
    the signature for the callers; it names nothing any more."""
    return tempfile.mkdtemp(prefix=WORKDIR_PREFIX, dir=TMP_ROOT)


def _release_workdir(workdir: str) -> None:
    """Remove an exec's working directory once the up-mirror has read it.
    Only a directory `_get_session_workdir` made - directly under `TMP_ROOT`,
    named `WORKDIR_PREFIX`… - is ever removed; anything else is left alone.
    Hard links into the mirror cache are just unlinked (the cache keeps its
    inodes); ``ignore_errors``: a stray file is not worth failing the exec."""
    if os.path.dirname(workdir) == TMP_ROOT and os.path.basename(workdir).startswith(WORKDIR_PREFIX):
        shutil.rmtree(workdir, ignore_errors=True)


def _prepare_workdir(event: Dict, s3_bucket: str, s3_prefix: str, workdir: str) -> Dict[str, Any]:
    """Shared pre-run step for both exec actions, into this exec's FRESH
    ``workdir``: bring the project's mirror cache up to date and link the
    requested project dirs from it, reconstitute on a checkpoint resume,
    hydrate the prefetch list. Returns the sync report (empty when no
    ``sync_dirs`` were requested) - carrying ``cache_dir`` and ``dirs`` so
    `_finish_run` can record what the run uploads back into the cache."""
    sync_dirs = _event_sync_dirs(event)
    sync_report: Dict[str, Any] = {}
    if sync_dirs and s3_bucket and s3_prefix:
        cache_dir = _workspace_cache(s3_prefix)
        sync_report = _sync_dirs_down(
            s3_bucket, s3_prefix, cache_dir, sync_dirs,
            max_bytes=event.get("sync_max_bytes"),
            skip_over_bytes=event.get("sync_skip_over_bytes"),
        )
        placed = _materialize_workdir(cache_dir, sync_report.get("present") or [], workdir)
        sync_report["cache_dir"] = cache_dir
        sync_report["dirs"] = list(sync_dirs)
        print(f"[HANDLER] workdir: {workdir} (fresh), {placed} file(s) linked from the verified cache {cache_dir}")
    # On a checkpoint resume, reconstitute the workdir from the durable S3
    # workspace (manifest + shards + prior outputs) BEFORE prefetch, so aa_env
    # continues from the last committed cursor. /tmp is never relied upon.
    if event.get("resume"):
        _sync_prefix_down(s3_bucket, s3_prefix, workdir)
    _prefetch_files(
        event.get("prefetch", []), s3_bucket, s3_prefix, workdir,
        preserve_paths=bool(event.get("preserve_paths")),
    )
    return sync_report


def _finish_run(
    result: Any,
    s3_bucket: str,
    s3_prefix: str,
    workdir: str,
    snapshot: Dict[str, Tuple[int, int]],
    sync_report: Dict[str, Any],
    wall: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
    """Incremental upload of what the run produced + the response envelope.
    ``uploaded`` is the list `coder_worker.run_in_env` turns into one
    ``file_added`` event per file. ``upload_errors`` names every object S3
    refused (path and error, capped at `_UPLOAD_ERRORS_MAX`); when refusals
    left nothing uploaded, stderr ends with a line saying so, where the run
    and the coder read it. ``wall`` (`_wall_report`) adds the structured
    timeout fields and appends the readable line to stderr; the captured
    stdout/stderr are kept as-is."""
    uploaded: List[str] = []
    refused: List[Tuple[str, str]] = []
    if s3_bucket and s3_prefix:
        uploaded = _upload_workdir_files(
            s3_bucket, s3_prefix, workdir, snapshot=snapshot,
            cache_dir=sync_report.get("cache_dir"), cache_dirs=tuple(sync_report.get("dirs") or ()),
            failures=refused,
        )
    stderr_text = _cap_output(result.stderr)
    if wall:
        stderr_text = (stderr_text.rstrip("\n") + "\n" if stderr_text.strip() else "") + wall["error"]
    if refused and not uploaded:
        note = (f"[HANDLER] {len(refused)} file(s) could not be uploaded to S3 "
                f"({refused[0][0]}: {refused[0][1]}); the run's outputs are incomplete.")
        stderr_text = (stderr_text.rstrip("\n") + "\n" if stderr_text.strip() else "") + note
    out: Dict[str, Any] = {
        "exit_code": result.returncode,
        "stdout": _cap_output(result.stdout),
        "stderr": stderr_text,
        "uploaded": uploaded,
        "upload_errors": [{"path": p, "error": e} for p, e in refused[:_UPLOAD_ERRORS_MAX]],
    }
    if wall:
        out.update(wall)
    if sync_report:
        out["sync"] = {
            "downloaded": len(sync_report.get("downloaded") or []),
            "unchanged": sync_report.get("unchanged", 0),
            "skipped": list(sync_report.get("skipped") or [])[:50],
            "truncated": bool(sync_report.get("truncated")),
        }
    return out


def _exec_python(event: Dict, s3_bucket: str, s3_prefix: str, budget: Any,
                 invocation_env: Optional[Dict[str, str]] = None) -> Dict[str, Any]:
    """Run Python inside the mirrored project under ``budget`` (a `Budget`, or
    an int number of seconds for direct callers). ``invocation_env`` is THIS
    invocation's connector variables and deadline (`handler` passes them; a
    direct caller that passes none gets a child with no connector variables -
    never the process environment's).

    Two ways to say WHAT to run, in order of preference:

    * ``script_path`` — a project-relative file (``code/src/run.py``) that the
      coder wrote with ``Write``/``Edit`` (its PostToolUse sync hook uploaded
      it) and that ``sync_dirs`` just mirrored down. Runs with cwd = project
      root and ``code/`` on PYTHONPATH.

    The project root is this exec's own fresh directory: mirrored
    in, run, mirrored up, removed - nothing on disk carries over to the next
    exec; only the workspace does.
    * ``code`` + ``script_name`` — inline source (legacy callers and quick
      probes). Written to ``code/scratch/{script_name}`` inside the mirror so
      it is uploaded with everything else instead of littering the root.
    """
    code = event.get("code", "")
    script_name = os.path.basename(event.get("script_name") or "_lambda_exec.py")
    script_path = event.get("script_path")

    workdir = _get_session_workdir(s3_prefix)
    try:
        sync_report = _prepare_workdir(event, s3_bucket, s3_prefix, workdir)

        # Snapshot BEFORE the inline script is written so it counts as "new" and
        # is uploaded (under code/scratch/) along with whatever the run produces.
        snapshot = _snapshot_workdir(workdir)

        if script_path:
            rel = _safe_project_relpath(script_path)
            if rel is None:
                return {"exit_code": 1, "stdout": "", "stderr": (
                    f"script_path must be a project-relative path (got {script_path!r}); "
                    "absolute paths and '..' are refused."
                ), "uploaded": []}
            tmp_script = os.path.join(workdir, *rel.split("/"))
            if not os.path.isfile(tmp_script):
                return {"exit_code": 1, "stdout": "", "stderr": (
                    f"script_path {rel!r} not found in the project mirror after syncing "
                    f"{_event_sync_dirs(event) or list(DEFAULT_SYNC_DIRS)}. Write the file "
                    "first (Write/Edit under code/ — the sync hook uploads it), then run it."
                ), "uploaded": []}
            print(f"[HANDLER] _exec_python: script_path={rel}, workdir={workdir}")
        else:
            script_dir = EPHEMERAL_SCRIPT_DIR if event.get("ephemeral") else INLINE_SCRIPT_DIR
            rel = f"{script_dir}/{script_name}"
            tmp_script = os.path.join(workdir, *rel.split("/"))
            if code:
                os.makedirs(os.path.dirname(tmp_script), exist_ok=True)
                with open(tmp_script, "w") as f:
                    f.write(code)
                print(f"[HANDLER] Script written to {tmp_script}")
            elif not os.path.exists(tmp_script):
                return {"exit_code": 1, "stdout": "", "stderr": "No code provided and script not found",
                        "uploaded": []}
            print(f"[HANDLER] _exec_python: script={rel}, code_len={len(code)}, workdir={workdir}")

        # The user script runs with cwd = project root. PYTHONPATH carries the
        # aa_env checkpoint helper (a private per-exec copy) and the
        # project's code/ root — see _child_pythonpath. No bytecode is written
        # anywhere the child can reach later. HOME is the invoking scope's own
        # - see _home_scope.
        helper_dir = _helper_dir()
        scope = _home_scope(event, os.path.basename(helper_dir))
        child_env = _child_env({**(invocation_env or {}), "PYTHONUNBUFFERED": "1", "PYTHONDONTWRITEBYTECODE": "1"},
                               scope=scope)
        child_env["PYTHONPATH"] = _child_pythonpath(workdir, child_env, helper_dir)
        _ensure_overlay_intact()

        try:
            result, wall = _run_under_wall([sys.executable, tmp_script], workdir, child_env, _coerce_budget(budget))
        finally:
            _release_exec_dirs(helper_dir, scope, child_env)

        print(f"[HANDLER] _exec_python done: exit_code={result.returncode}, stdout_len={len(result.stdout)}, stderr_len={len(result.stderr)}")
        if wall:
            print(f"[HANDLER] wall: {wall['error']} (elapsed {wall['elapsed_seconds']}s, bound={wall['bound']})")
        if result.returncode != 0:
            print(f"[HANDLER] stderr: {result.stderr[:500]}")

        return _finish_run(result, s3_bucket, s3_prefix, workdir, snapshot, sync_report, wall=wall)
    finally:
        _release_workdir(workdir)       # after the up-mirror has read it


def _exec_shell(event: Dict, s3_bucket: str, s3_prefix: str, budget: Any,
                invocation_env: Optional[Dict[str, str]] = None) -> Dict[str, Any]:
    """Run a shell command with cwd = the mirrored project root (e.g. ``python
    -m pytest code/tests -q``). Same sync / snapshot / incremental-upload /
    wall envelope - and the same ``invocation_env`` contract - as
    `_exec_python`. The command runs under bash when the
    image has it (`_shell_argv`): dash, Debian's ``sh``, drops the hyphenated
    ``CONN_<connector id>_*`` variables, so under ``sh -c`` no connector
    credential reached the command or anything it ran."""
    command = event.get("command", "")
    if not command:
        return {"exit_code": 1, "stdout": "", "stderr": "No command provided"}

    workdir = _get_session_workdir(s3_prefix)
    try:
        sync_report = _prepare_workdir(event, s3_bucket, s3_prefix, workdir)
        snapshot = _snapshot_workdir(workdir)

        helper_dir = _helper_dir()
        scope = _home_scope(event, os.path.basename(helper_dir))
        child_env = _child_env({**(invocation_env or {}), "PYTHONDONTWRITEBYTECODE": "1"}, scope=scope)
        child_env["PYTHONPATH"] = _child_pythonpath(workdir, child_env, helper_dir)
        _ensure_overlay_intact()

        shell = _shell_argv(login_shell=bool(event.get("login_shell")))
        if shell[0] == "sh":
            dropped = [k for k in child_env if k.startswith("CONN_") and not k.isidentifier()]
            if dropped:
                print(f"[HANDLER] exec_shell: bash not found in this image; sh drops {len(dropped)} connector "
                      f"variable(s) whose names are not valid shell identifiers - the command and the programs it "
                      f"runs will not see them. Install bash in the image (Debian/Ubuntu bases carry it).")

        try:
            result, wall = _run_under_wall([*shell, command], workdir, child_env, _coerce_budget(budget))
        finally:
            _release_exec_dirs(helper_dir, scope, child_env)
        if wall:
            print(f"[HANDLER] wall: {wall['error']} (elapsed {wall['elapsed_seconds']}s, bound={wall['bound']})")

        return _finish_run(result, s3_bucket, s3_prefix, workdir, snapshot, sync_report, wall=wall)
    finally:
        _release_workdir(workdir)       # after the up-mirror has read it


def _write_file(event: Dict, s3_bucket: str, s3_prefix: str) -> Dict[str, Any]:
    rel_path = event.get("path", "")
    content = event.get("content", "")
    if not rel_path:
        return {"exit_code": 1, "stdout": "", "stderr": "No path provided"}

    s3_key = f"{s3_prefix.rstrip('/')}/{rel_path.lstrip('/')}"
    try:
        _s3_client().put_object(
            Bucket=s3_bucket,
            Key=s3_key,
            Body=content.encode("utf-8"),
            **_s3_put_args(s3_key),
        )
        return {"exit_code": 0, "stdout": f"Written {len(content)} bytes to {rel_path}", "stderr": ""}
    except Exception as e:
        return {"exit_code": 1, "stdout": "", "stderr": f"Write failed: {e}"}


def _read_file(event: Dict, s3_bucket: str, s3_prefix: str) -> Dict[str, Any]:
    rel_path = event.get("path", "")
    if not rel_path:
        return {"exit_code": 1, "stdout": "", "stderr": "No path provided"}

    s3_key = f"{s3_prefix.rstrip('/')}/{rel_path.lstrip('/')}"
    try:
        resp = _s3_client().get_object(Bucket=s3_bucket, Key=s3_key)
        content = resp["Body"].read().decode("utf-8", errors="replace")
        return {"exit_code": 0, "stdout": content, "stderr": ""}
    except Exception as e:
        if "NoSuchKey" in str(e):
            return {"exit_code": 1, "stdout": "", "stderr": f"File not found: {rel_path}"}
        return {"exit_code": 1, "stdout": "", "stderr": f"Read failed: {e}"}


def _list_files(event: Dict, s3_bucket: str, s3_prefix: str) -> Dict[str, Any]:
    try:
        s3 = _s3_client()
        paginator = s3.get_paginator("list_objects_v2")
        pages = paginator.paginate(Bucket=s3_bucket, Prefix=s3_prefix)

        lines = []
        for page in pages:
            for obj in page.get("Contents", []):
                key = obj["Key"]
                rel = key[len(s3_prefix):]
                if not rel or rel.endswith("/"):
                    continue
                size = obj.get("Size", 0)
                lines.append(f"f {size} {rel}")

        return {"exit_code": 0, "stdout": "\n".join(lines), "stderr": ""}
    except Exception as e:
        return {"exit_code": 1, "stdout": "", "stderr": f"List failed: {e}"}
