kaxil commented on code in PR #73559:
URL: https://github.com/apache/airflow/pull/73559#discussion_r4123171475


##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/base.py:
##########
@@ -109,12 +137,89 @@ class SandboxSpec:
         than a hostname list gives, so it needs no opt-in. Both lists may be 
set
         together; how a backend combines them, and what that costs, is the
         backend's to document.
+    :param owner: Who the sandbox is for, when a task provisions it for an 
agent
+        task to attach to later. A 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`
+        attaching to the sandbox has to present the same value, and by default 
it
+        presents the Dag run it is part of, so the provisioning task in the 
same
+        run writes ``owner=dag_run_owner(context)``. Unset for a sandbox nobody
+        will attach to. A backend that cannot record it must refuse it.
     """
 
     env: Mapping[str, str] | None = None
     block_network: bool = True
     allow_egress_to: Sequence[str] | None = None
     allow_egress_to_cidrs: Sequence[str] | None = None
+    owner: str | None = None
+
+
+def dag_run_owner(context: Mapping[str, Any]) -> str:

Review Comment:
   No convention I could find in the provider or the SDK, so I kept the noun 
form: it reads as a value at the call site, 
`SandboxSpec(owner=dag_run_owner(context))`. Happy to rename if you prefer.



##########
providers/common/ai/src/airflow/providers/common/ai/sandbox/base.py:
##########
@@ -109,12 +137,89 @@ class SandboxSpec:
         than a hostname list gives, so it needs no opt-in. Both lists may be 
set
         together; how a backend combines them, and what that costs, is the
         backend's to document.
+    :param owner: Who the sandbox is for, when a task provisions it for an 
agent
+        task to attach to later. A 
:class:`~airflow.providers.common.ai.toolsets.sandbox.SandboxToolset`
+        attaching to the sandbox has to present the same value, and by default 
it
+        presents the Dag run it is part of, so the provisioning task in the 
same
+        run writes ``owner=dag_run_owner(context)``. Unset for a sandbox nobody
+        will attach to. A backend that cannot record it must refuse it.
     """
 
     env: Mapping[str, str] | None = None
     block_network: bool = True
     allow_egress_to: Sequence[str] | None = None
     allow_egress_to_cidrs: Sequence[str] | None = None
+    owner: str | None = None
+
+
+def dag_run_owner(context: Mapping[str, Any]) -> str:
+    """
+    Return the owner token naming the Dag run a task is part of: 
``"<dag_id>/<run_id>"``.
+
+    This is what a ``SandboxToolset`` presents when it attaches to a sandbox
+    without an explicit ``owner``, so a task provisioning a sandbox for an 
agent
+    task in the same Dag run stamps it with 
``SandboxSpec(owner=dag_run_owner(context))``.
+    The pair is unique across the deployment where a bare ``run_id`` is not: 
two
+    Dags on the same schedule share their run ids. ``context`` is the task 
context,
+    as a ``@task`` receives it in ``**context`` or ``get_current_context`` 
returns it.
+    """
+    ti = context["ti"]
+    return f"{ti.dag_id}/{ti.run_id}"
+
+
+def encode_network_policy(spec: SandboxSpec) -> str:
+    """
+    Serialize a spec's network policy for a sandbox tag, so an attaching 
toolset can read it back.
+
+    The toolset tells the model what the sandbox can reach, because a model 
that has to
+    discover a denied network by failing wastes a turn, or a whole command 
budget. An
+    attached sandbox was provisioned under a spec the toolset never sees, so 
the backend
+    records the policy on the sandbox at create and 
:func:`decode_network_policy` turns
+    it back into a spec. Compact JSON with sorted keys, so the same policy 
always encodes
+    the same way.
+    """
+    return json.dumps(
+        {
+            "block_network": spec.block_network,
+            "allow_egress_to": list(spec.allow_egress_to or ()),
+            "allow_egress_to_cidrs": list(spec.allow_egress_to_cidrs or ()),
+        },
+        separators=(",", ":"),
+        sort_keys=True,
+    )
+
+
+def decode_network_policy(value: str | None) -> SandboxSpec | None:
+    """Read a :data:`NETWORK_TAG` value back into a spec carrying only the 
network fields, or ``None``."""
+    if not value:
+        return None

Review Comment:
   Done: it raises `ValueError` now. `attach` catches it, logs a warning naming 
the sandbox and the stamp, and reports no network policy rather than describing 
one nobody asked for.



##########
providers/common/ai/src/airflow/providers/common/ai/toolsets/sandbox.py:
##########
@@ -230,26 +311,146 @@ async def for_run(self, ctx: RunContext[Any]) -> 
AbstractToolset[Any]:
         # not silently degrade to this class on every run.
         return type(self)(
             self._backend,
-            spec=self._spec,
+            # Attach mode refuses a spec, and the default one filled in above 
is
+            # not the author's, so it is not handed back.
+            spec=None if self._attach_mode else self._spec,
             default_command_timeout=self._default_command_timeout,
             max_command_timeout=self._max_command_timeout,
             max_output_lines=self._max_output_lines,
             max_output_bytes=self._max_output_bytes,
             max_read_bytes=self._max_read_bytes,
             tool_prefix=self._tool_prefix,
+            attach_to=self._attached_handle() if self._attach_mode else None,
+            owner=self._owner,
         )
 
+    def _attached_handle(self) -> str:
+        """
+        Return the rendered handle, checked, because the templater bypasses 
the constructor.
+
+        ``attach_to`` is templated so the handle can come from an upstream 
XCom, and
+        an XCom that was never pushed renders to ``None`` under native 
rendering or to
+        the string ``"None"`` otherwise. Neither is a sandbox, and silently 
falling back
+        to provisioning one would run the agent in an empty workspace with no 
error.
+        """
+        handle = self.attach_to
+        if not is_sandbox_handle(handle):
+            raise SandboxTerminalError(
+                f"attach_to rendered to {handle!r}, which is not a sandbox 
handle. The task that "
+                "provisions the sandbox pushed nothing, or this task does not 
depend on it and ran "
+                "first; check the upstream task and the XCom it returns."
+            )
+        return handle
+
     async def __aenter__(self) -> Self:
-        # The sandbox is provisioned lazily on first use, not here: a durable
+        # An owned sandbox is provisioned lazily on first use, not here: a 
durable
         # replay that only serves cached tool results must not provision one, 
and
-        # nothing leaks if the run fails before any tool executes.
+        # nothing leaks if the run fails before any tool executes. An attached 
one is
+        # claimed now, so a wrong handle or a held sandbox fails the run 
before the
+        # model has spent anything, and the tool descriptions can state the 
lifetime.
+        if self._attach_mode:
+            await self._attach(self._attached_handle())
         return self
 
+    async def _attach(self, handle: str) -> None:
+        owner, holder = self._identity()
+        backend = self._attachable_backend()
+        try:
+            attached = await asyncio.to_thread(backend.attach, handle, 
owner=owner, holder=holder)
+        except SandboxTerminalError:
+            raise
+        except SandboxError as e:
+            # Nothing the model does can change whether this sandbox can be 
attached
+            # to, so a recoverable label here is one it could not act on.
+            raise SandboxTerminalError(
+                f"Could not attach to sandbox {handle!r} on backend 
{backend.name!r}: {e}"
+            ) from e
+        self._sandbox = handle
+        self._holder = holder
+        remaining = attached.remaining_lifetime
+        self._expires_at = None if remaining is None else time.monotonic() + 
remaining
+        self._attach_note = self._describe_attached(attached.network, 
remaining)
+        log.info(
+            "Attached to sandbox %s on backend %s as %s; %s of its lifetime 
remain",
+            handle,
+            backend.name,
+            holder,
+            "an unknown number of seconds" if remaining is None else 
f"{remaining:.0f}s",
+        )
+
+    @classmethod
+    def _describe_attached(cls, network: SandboxSpec | None, 
remaining_lifetime: float | None) -> str:
+        whose = (
+            "This sandbox was set up by an earlier task, and your files stay 
in it after this run "
+            "for a later task to collect."
+        )
+        policy = (
+            cls._describe_network(network)
+            if network is not None
+            else "Its network access is whatever the task that set it up 
allowed; test before relying on it."
+        )
+        if remaining_lifetime is None:
+            clock = "How long it has left is not known."
+        elif remaining_lifetime < 60:
+            clock = "Under a minute of its lifetime remained when this run 
began, so finish up."
+        else:
+            minutes = round(remaining_lifetime / 60)
+            unit = "minute" if minutes == 1 else "minutes"
+            clock = f"About {minutes} {unit} of its lifetime remained when 
this run began."
+        return f"{whose} {policy} {clock}"
+
+    def _attachable_backend(self) -> AttachableSandboxBackend:

Review Comment:
   Done, and the same for the other two no-argument helpers on the class 
(`_attached_handle`, `_network_note`).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to