Skip to content

Sessions

Session Container

agora_workbench.code_execution.sessions.session.Session(session_id, data, session_type, user_identity, user_token, token_claims, metadata=None, data_manager=None, extensions=None)

Bases: Generic[T]

Generic session container managing stateful data across MCP server interactions.

Sessions provide lifecycle management, ownership tracking, and metadata storage for persistent resources like code execution environments, database connections, or computation state. Each session is owned by a specific user (identified via JWT token claims) and tracks access patterns, status transitions, and cleanup requirements.

Attributes:

Name Type Description
session_id str

Unique identifier for the session.

data T

The session's payload data, type-parameterized for type safety.

session_type str

Categorizes session type (e.g., "python", "database").

user_identity str

Owner's composite identifier from JWT token (oid@tid).

user_token str

User's bearer token for authentication.

metadata Dict

Optional key-value metadata for session configuration.

token_claims Dict

Optional cached JWT token claims for session authorization. These claims are used to restore authentication context without re-validating the JWT token. Intentionally excluded from get_info() for security.

created_at datetime

Timestamp when session was created.

last_accessed datetime

Timestamp of most recent session access.

status str

Current session state (e.g., "created", "active", "error").

data_manager DataLakeDataManager

Manager for DataLake asset access. Owned by the session — :meth:cleanup tears it down, so an injected manager must not be shared between sessions.

Example

from .session import Session session = Session( ... session_id="sess_123", ... data={"counter": 0}, ... session_type="demo", ... user_identity="user-oid-xyz", ... user_token="eyJ...", ... metadata={"version": "1.0"}, ... token_claims={"oid": "user-oid-xyz", "exp": 1234567890}, ... ) session.touch() # Update last accessed time session.update_status("active") info = session.get_info()

Initialize a session.

Parameters:

Name Type Description Default
session_id str

Unique identifier for the session.

required
data T

The session's payload data.

required
session_type str

Categorizes session type (e.g. "python").

required
user_identity str

Owner's composite identifier from JWT token (oid@tid).

required
user_token str

User's bearer token for authentication.

required
token_claims Dict

Cached JWT claims for the user token.

required
metadata Optional[Dict]

Optional key-value metadata for session configuration.

None
data_manager Optional[DataLakeDataManager]

Optional pre-built data manager for DataLake asset access. When omitted, a default :class:DataLakeDataManager is constructed. The session takes ownership of whichever manager it ends up with — :meth:cleanup calls cleanup() on it — so a caller supplying one must pass a fresh instance per session rather than a shared singleton.

None
Source code in src/agora_workbench/code_execution/sessions/session.py
def __init__(
    self,
    session_id: str,
    data: T,
    session_type: str,
    user_identity: str,
    user_token: str,
    token_claims: Dict,
    metadata: Optional[Dict] = None,
    data_manager: Optional["DataLakeDataManager"] = None,
    extensions: Optional[Dict[str, Any]] = None,
):
    """
    Initialize a session.

    Args:
        session_id: Unique identifier for the session.
        data: The session's payload data.
        session_type: Categorizes session type (e.g. ``"python"``).
        user_identity: Owner's composite identifier from JWT token (``oid@tid``).
        user_token: User's bearer token for authentication.
        token_claims: Cached JWT claims for the user token.
        metadata: Optional key-value metadata for session configuration.
        data_manager: Optional pre-built data manager for DataLake asset
            access. When omitted, a default :class:`DataLakeDataManager` is
            constructed. The session takes ownership of whichever manager it
            ends up with — :meth:`cleanup` calls ``cleanup()`` on it — so a
            caller supplying one must pass a fresh instance per session
            rather than a shared singleton.
    """
    self.session_id = session_id
    self.data = data
    self.created_at = datetime.now()
    self.last_accessed = datetime.now()
    self.metadata = metadata or {}
    self.user_identity = user_identity
    self.user_token = user_token
    self.token_claims = token_claims
    self.status = "created"
    self.session_type = session_type
    self.extensions = extensions or {}
    self._asset_counter: int = 0
    self._status_history = [("created", datetime.now())]
    self._scheduled_cleanup_tasks: set[asyncio.Task[None]] = set()
    self._claimed_cleanup_errors: list[Exception] = []
    self._session_file_cleanup_claimed = False
    self._session_file_cleanup_attempts = 0

    # Initialize data manager for DataLake asset access. Constructing the
    # default lazily matters: DataLakeDataManager.__init__ eagerly allocates
    # a temp cache dir, so building one only to discard it would leak it.
    if data_manager is None:
        from ..data_access.manager import DataLakeDataManager

        data_manager = DataLakeDataManager()

    self.data_manager = data_manager

    # Initialize object store for asset objects
    self.object_store = ObjectStore()

touch()

Update last accessed timestamp.

Source code in src/agora_workbench/code_execution/sessions/session.py
def touch(self):
    """Update last accessed timestamp."""
    self.last_accessed = datetime.now()

update_status(new_status)

Update session status with history tracking.

Source code in src/agora_workbench/code_execution/sessions/session.py
def update_status(self, new_status: str):
    """Update session status with history tracking."""
    self.status = new_status
    self._status_history.append((new_status, datetime.now()))
    self.touch()

get_info()

Return session information.

Source code in src/agora_workbench/code_execution/sessions/session.py
def get_info(self) -> Dict[str, Any]:
    """Return session information."""
    age_seconds = (datetime.now() - self.created_at).total_seconds()
    idle_seconds = (datetime.now() - self.last_accessed).total_seconds()

    return {
        "session_id": self.session_id,
        "session_type": self.session_type,
        "status": self.status,
        "created_at": self.created_at.isoformat(),
        "last_accessed": self.last_accessed.isoformat(),
        "age_seconds": age_seconds,
        "idle_seconds": idle_seconds,
        "metadata": self.metadata,
        "user_identity": self.user_identity,
        "status_history": [{"status": s, "timestamp": t.isoformat()} for s, t in self._status_history],
    }

cleanup()

Start cleanup of every session resource, including async-only clients.

Async cleanup tasks are retained until the SessionManager claims them; callers that need completion should use :meth:aclose.

Source code in src/agora_workbench/code_execution/sessions/session.py
def cleanup(self) -> None:
    """
    Start cleanup of every session resource, including async-only clients.

    Async cleanup tasks are retained until the SessionManager claims them;
    callers that need completion should use :meth:`aclose`.
    """
    claimed_cleanup_errors = self._take_claimed_cleanup_errors()
    errors: list[Exception] = []
    cancellations: list[asyncio.CancelledError] = []
    self._cleanup_resource_sync(self.data_manager, "data manager", errors, cancellations)
    for resource in self.extensions.values():
        self._cleanup_resource_sync(resource, f"extension {type(resource).__name__}", errors, cancellations)
    self._cleanup_resource_sync(self.data, "session payload", errors, cancellations)
    try:
        self._cleanup_session_file()
    except Exception as exc:
        errors.extend(claimed_cleanup_errors)
        errors.append(exc)
    else:
        if self._session_file_cleanup_attempts >= _MAX_SESSION_FILE_CLEANUP_ATTEMPTS:
            errors.extend(claimed_cleanup_errors)
    if cancellations:
        if errors:
            cancellations[0].add_note(str(ExceptionGroup("Additional session cleanup failures.", errors)))
        raise cancellations[0]
    if errors:
        raise ExceptionGroup("Session cleanup failed.", errors)

aclose() async

Attempt all asynchronous cleanup steps, then report aggregated failures.

Source code in src/agora_workbench/code_execution/sessions/session.py
async def aclose(self) -> None:
    """Attempt all asynchronous cleanup steps, then report aggregated failures."""
    claimed_cleanup_errors = self._take_claimed_cleanup_errors()
    errors: list[Exception] = []
    cancellations: list[asyncio.CancelledError] = []
    await self._cleanup_resource_async(self.data_manager, "data manager", errors, cancellations)
    for resource in self.extensions.values():
        await self._cleanup_resource_async(
            resource,
            f"extension {type(resource).__name__}",
            errors,
            cancellations,
        )
    await self._cleanup_resource_async(self.data, "session payload", errors, cancellations)
    try:
        self._cleanup_session_file()
    except Exception as exc:
        errors.extend(claimed_cleanup_errors)
        errors.append(exc)
    else:
        if self._session_file_cleanup_attempts >= _MAX_SESSION_FILE_CLEANUP_ATTEMPTS:
            errors.extend(claimed_cleanup_errors)
    retry_tasks = tuple(self._scheduled_cleanup_tasks)
    if retry_tasks:
        results = await asyncio.gather(
            *(asyncio.shield(task) for task in retry_tasks),
            return_exceptions=True,
        )
        for task, result in zip(retry_tasks, results):
            if task.done():
                self._scheduled_cleanup_tasks.discard(task)
            if isinstance(result, asyncio.CancelledError):
                cancellations.append(result)
            elif isinstance(result, Exception):
                errors.append(result)
    if cancellations:
        if errors:
            cancellations[0].add_note(str(ExceptionGroup("Additional session cleanup failures.", errors)))
        raise cancellations[0]
    if errors:
        raise ExceptionGroup("Session cleanup failed.", errors)

take_cleanup_tasks()

Transfer ownership of sync-scheduled cleanup tasks to the manager.

Source code in src/agora_workbench/code_execution/sessions/session.py
def take_cleanup_tasks(self) -> tuple[asyncio.Task[None], ...]:
    """Transfer ownership of sync-scheduled cleanup tasks to the manager."""
    tasks = tuple(self._scheduled_cleanup_tasks)
    self._scheduled_cleanup_tasks.clear()
    return tasks

claim_session_file_cleanup()

Remove the owned session file before its session ID can be reused.

Source code in src/agora_workbench/code_execution/sessions/session.py
def claim_session_file_cleanup(self) -> None:
    """Remove the owned session file before its session ID can be reused."""
    if self._session_file_cleanup_claimed:
        return
    try:
        self._cleanup_session_file()
    except Exception as exc:
        self._claimed_cleanup_errors.append(exc)
        self._session_file_cleanup_attempts += 1
        if self._session_file_cleanup_attempts >= _MAX_SESSION_FILE_CLEANUP_ATTEMPTS:
            # Give up rather than pin the session ID forever. The claim is
            # terminal, so this session never deletes the path again and a
            # replacement session reusing the ID keeps its own file.
            self._session_file_cleanup_claimed = True
            LOGGER.warning(
                "Giving up on removing the session file for session %s after %d attempts: %s",
                self.session_id,
                self._session_file_cleanup_attempts,
                exc,
            )

Session Manager

agora_workbench.code_execution.sessions.manager.SessionManager(config=None, kernel_name='tools-py')

Manages the lifecycle of code-execution sessions and their Jupyter kernels.

Responsibilities include session creation, retrieval, timeout-based cleanup, kernel provisioning (one AsyncKernelManager per session), background-job tracking, and artifact registration for the /artifacts download endpoint.

Thread-safety is provided by an internal RLock for session-lifecycle mutations and per-session asyncio.Lock instances for kernel access.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def __init__(self, config: Optional[SessionConfig] = None, kernel_name: str = "tools-py"):
    self.config = config or SessionConfig()
    self.kernel_name = kernel_name
    self.storage = self.config.storage_backend
    self._last_cleanup = datetime.now()
    self._session_lifecycle_lock = RLock()
    self._session_lifecycle_condition = Condition(self._session_lifecycle_lock)
    self._closing_session_ids: set[str] = set()
    self._closing_sessions: dict[str, Session] = {}
    timeout_seconds = self.config.timeout.total_seconds()
    self.execution_session_keepalive_seconds = max(0.5, min(timeout_seconds / 10.0, 60.0))

    # session_id -> (KernelManager, KernelClient)
    self._kernels: dict[str, Tuple[AsyncKernelManager, "AsyncKernelClient"]] = {}
    self._kernel_last_used: dict[str, float] = {}  # session_id -> timestamp
    self._kernel_tokens: dict[str, Optional[str]] = {}  # session_id -> last injected user token
    self._kernel_session_generations: dict[str, int | None] = {}
    self._kernel_network_isolation_validated = False
    self._kernel_network_isolation_error: str | None = None
    self._kernel_network_isolation_lock = asyncio.Lock()
    self._kernel_ipc_dirs: dict[AsyncKernelManager, Path] = {}
    # Per-session lock that serializes execute_code_for_session calls so the
    # shared Jupyter kernel client (single iopub queue) cannot be raced by
    # concurrent callers (e.g. four parallel push_object MCP calls).
    self._kernel_execute_locks: dict[str, asyncio.Lock] = {}
    self._kernel_execute_lock_users: dict[str, int] = {}
    self._kernel_execute_drained: dict[str, asyncio.Event] = {}
    self._retired_kernel_execute_locks: set[str] = set()
    self._session_resource_users: dict[str, int] = {}
    self._session_resources_drained: dict[str, asyncio.Event] = {}
    self._pending_session_resource_cleanup: dict[str, Session] = {}
    self._session_resource_cleanup_tasks: dict[str, asyncio.Task[None]] = {}
    self._session_owned_cleanup_tasks: dict[str, set[asyncio.Task[None]]] = {}
    self._background_jobs: dict[str, _BackgroundJob] = {}
    self._session_running_jobs: dict[str, str] = {}
    # Artifact pipeline state: the token -> record map used by the HTTP
    # download endpoint.
    self._session_artifacts: dict[str, dict[str, _ArtifactRecord]] = {}

    # Kernel bootstrap tracking.  Every kernel process gets a globally
    # unique, monotonically increasing *generation* id when it is started.
    # One-time bootstrap work (the AGORA_OUTPUT_DIR preamble, tool proxy
    # injection, ...) is recorded against that generation rather than
    # against the session id, so a session whose kernel is rebuilt is
    # re-bootstrapped automatically.  Because generations are never
    # reused, correctness does not depend on any teardown path
    # remembering to clear this state -- the entries below are dropped in
    # :meth:`_shutdown_kernel` purely to bound memory.
    self._kernel_generation_seq: int = 0
    self._kernel_generations: dict[str, int] = {}  # session_id -> generation
    self._kernel_bootstrap_state: dict[int, set[str]] = {}  # generation -> completed keys

    # In-flight kernel teardowns, so concurrent closes coalesce onto one
    # task instead of racing, callers can await teardown, and the task is
    # strongly referenced for its lifetime (an unreferenced task may be
    # garbage-collected mid-flight).
    self._kernel_shutdown_tasks: dict[str, "asyncio.Task[None]"] = {}
    self._kernel_start_tasks: dict[str, dict[asyncio.Task[Any], int]] = {}
    self._resource_cleanup_tasks: set[asyncio.Task[None]] = set()
    self._background_lease_release_tasks: set[asyncio.Task[None]] = set()
    self._resource_cleanup_errors: list[Exception] = []
    self._resource_cleanup_cancellations: list[asyncio.CancelledError] = []
    self._session_generation_seq = 0
    self._session_generations: dict[str, int] = {}

    LOGGER.info(
        f"Initialized SessionManager: max_sessions={self.config.max_sessions}, "
        f"timeout={self.config.timeout.total_seconds() / 60}min, "
        f"kernel_network_mode={self.config.kernel_network_mode}"
    )

set_kernel_network_mode(kernel_network_mode)

Change the launch policy only when no kernel can retain the old policy.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def set_kernel_network_mode(self, kernel_network_mode: Literal["inherit", "isolated"]) -> None:
    """Change the launch policy only when no kernel can retain the old policy."""
    if kernel_network_mode not in ("inherit", "isolated"):
        raise ValueError("kernel_network_mode must be 'inherit' or 'isolated'")

    with self._session_lifecycle_lock:
        if self.config.kernel_network_mode == kernel_network_mode:
            return
        active_session_ids = sorted(
            set(self._kernels) | {session_id for session_id, starts in self._kernel_start_tasks.items() if starts}
        )
        if active_session_ids:
            raise RuntimeError(
                "cannot change kernel_network_mode while kernels are running or starting: "
                f"{', '.join(active_session_ids)}"
            )
        self.config.kernel_network_mode = kernel_network_mode
        self._kernel_network_isolation_validated = False
        self._kernel_network_isolation_error = None

create_session(data, user_identity, user_token, token_claims, metadata=None, session_id=None)

Create a new session and return its ID.

Parameters:

Name Type Description Default
data Any

The data/object to store in the session

required
user_identity str

User identity (Entra ID, email, etc.)

required
user_token str

User's bearer token for authentication

required
metadata Optional[dict]

Optional metadata dict

None
session_id Optional[str]

Optional custom session ID (default: UUID)

None
token_claims dict

Cached JWT claims for the user token. Stored on the session so that background tasks can restore the full auth context without re-validating the token.

required

Returns:

Type Description
str

Session ID string

Note

When SessionConfig.data_manager_factory is configured it is invoked here, once per session, and the resulting manager is passed to the Session. The session owns it and cleans it up.

Raises:

Type Description
MaxSessionsReachedError

If the session limit has been reached.

TypeError

If a configured data_manager_factory returns something that is not a usable data manager.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def create_session(
    self,
    data: Any,
    user_identity: str,
    user_token: str,
    token_claims: dict,
    metadata: Optional[dict] = None,
    session_id: Optional[str] = None,
) -> str:
    """
    Create a new session and return its ID.

    Args:
        data: The data/object to store in the session
        user_identity: User identity (Entra ID, email, etc.)
        user_token: User's bearer token for authentication
        metadata: Optional metadata dict
        session_id: Optional custom session ID (default: UUID)
        token_claims: Cached JWT claims for the user token.
            Stored on the session so that background tasks can restore
            the full auth context without re-validating the token.

    Returns:
        Session ID string

    Note:
        When ``SessionConfig.data_manager_factory`` is configured it is
        invoked here, once per session, and the resulting manager is passed
        to the ``Session``. The session owns it and cleans it up.

    Raises:
        MaxSessionsReachedError: If the session limit has been reached.
        TypeError: If a configured ``data_manager_factory`` returns
            something that is not a usable data manager.
    """
    # Run periodic cleanup
    self._maybe_cleanup()

    with self._session_lifecycle_lock:
        # Generate session ID
        if session_id is None:
            session_id = str(uuid.uuid4())
        else:
            while session_id in self._closing_session_ids:
                try:
                    asyncio.get_running_loop()
                except RuntimeError:
                    self._session_lifecycle_condition.wait()
                else:
                    raise ValueError(f"Session {session_id} is still closing; retry after cleanup completes.")
            if self.storage.retrieve(session_id) is not None:
                raise ValueError(f"Session {session_id} already exists.")

        # Enforce max sessions limit
        self._enforce_max_sessions()

        # Build a customized data manager when a factory is configured, so
        # the Session never constructs (and immediately discards) a default
        # one — DataLakeDataManager allocates a temp cache dir eagerly.
        data_manager = None
        extensions: dict[str, Any] = {}
        if self.config.data_manager_factory is not None:
            factory_result = self.config.data_manager_factory(
                SessionContext(
                    session_id=session_id,
                    user_identity=user_identity,
                    user_token=user_token,
                    token_claims=token_claims,
                    session_type="default",
                    metadata=metadata or {},
                )
            )
            if isinstance(factory_result, SessionResources):
                data_manager = factory_result.data_manager
                extensions = dict(factory_result.extensions)
            else:
                data_manager = factory_result
            # Validate eagerly. A factory returning None is the dangerous
            # case: Session would fall back to building a default manager,
            # silently discarding the customization, so the mistake would
            # surface later as unexplained default behaviour instead of an
            # error at the point of the bug.
            if data_manager is None or not callable(getattr(data_manager, "cleanup", None)):
                error = TypeError(
                    "SessionConfig.data_manager_factory must return a data manager instance with a "
                    f"cleanup() method, but it returned {type(data_manager).__name__}. Returning None "
                    "would silently fall back to a default DataLakeDataManager and discard the "
                    "customization the factory exists to provide."
                )
                cleanup_error = self._cleanup_unclaimed_session_resources(data_manager, extensions)
                if cleanup_error is not None:
                    error.add_note(f"Factory resource rollback also failed: {cleanup_error!r}")
                raise error

        # Create session
        session = None
        try:
            session = Session(
                session_id=session_id,
                data=data,
                session_type="default",
                user_identity=user_identity,
                user_token=user_token,
                token_claims=token_claims,
                metadata=metadata,
                data_manager=data_manager,
                extensions=extensions if self.config.data_manager_factory is not None else None,
            )

            # Store
            self.storage.store(session_id, session)
        except BaseException as creation_error:
            if session is None:
                cleanup_error = self._cleanup_unclaimed_session_resources(data_manager, extensions)
            else:
                try:
                    self._start_session_cleanup(session)
                except BaseException as exc:
                    cleanup_error = exc
                else:
                    cleanup_error = None
            if cleanup_error is not None:
                creation_error.add_note(f"Factory resource rollback also failed: {cleanup_error!r}")
            raise
        self._session_generation_seq += 1
        self._session_generations[session_id] = self._session_generation_seq
        self._retired_kernel_execute_locks.discard(session_id)

        # Create the per-session outputs directory.  Done eagerly so the
        # kernel can write to it on the very first execute.  Failure to
        # create is non-fatal: artifact discovery just becomes a no-op for
        # this session.
        try:
            self._get_outputs_dir(session_id).mkdir(parents=True, exist_ok=True)
        except OSError:
            LOGGER.warning(
                "Could not create outputs dir for session %s under %s; "
                "artifact download will be disabled for this session.",
                session_id,
                _OUTPUTS_BASE_DIR,
                exc_info=True,
            )

    LOGGER.info(f"Created session {session_id} (total={self.storage.count()})")

    return session_id

get_session(session_id)

Get a session by ID, updating its access time.

Returns:

Type Description
Session

Session object

Raises:

Type Description
ValueError

If session not found or expired

Source code in src/agora_workbench/code_execution/sessions/manager.py
def get_session(self, session_id: str) -> Session:
    """
    Get a session by ID, updating its access time.

    Returns:
        Session object

    Raises:
        ValueError: If session not found or expired
    """
    self._maybe_cleanup()

    with self._session_lifecycle_lock:
        session = self.storage.retrieve(session_id)

        if session is None:
            raise ValueError(
                f"Session {session_id} not found. It may have expired or been "
                f"cleaned up. Active sessions: {self.storage.count()}"
            )

        # Update access time
        session.touch()
        self.storage.store(session_id, session)

    return session

update_session(session_id, session)

Update an existing session.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def update_session(self, session_id: str, session: Session) -> None:
    """Update an existing session."""
    with self._session_lifecycle_lock:
        if self.storage.retrieve(session_id) is None:
            raise ValueError(f"Session {session_id} not found")

        session.touch()
        self.storage.store(session_id, session)

update_status(session_id, status)

Update the status of a session.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def update_status(self, session_id: str, status: str) -> None:
    """Update the status of a session."""
    with self._session_lifecycle_lock:
        session = self.get_session(session_id)
        session.update_status(status)
        self.storage.store(session_id, session)

close_session(session_id)

Explicitly close a session.

Attempts to clean up resources first. If cleanup fails, the session is still deleted to prevent session accumulation, but the error is logged.

Kernel teardown is asynchronous. This method schedules it and returns the task, so the session is removed from storage immediately but its kernel process may still be running when the call returns. Callers that need the kernel's resources actually released — GPU memory, for instance — should use :meth:aclose_session instead, or await the returned task.

Parameters:

Name Type Description Default
session_id str

ID of the session to close

required

Returns:

Type Description
Optional[Task[None]]

The teardown task, or None when there was no kernel to tear

Optional[Task[None]]

down or no running event loop to schedule it on.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def close_session(self, session_id: str) -> "Optional[asyncio.Task[None]]":
    """
    Explicitly close a session.

    Attempts to clean up resources first. If cleanup fails, the session
    is still deleted to prevent session accumulation, but the error is logged.

    Kernel teardown is asynchronous. This method schedules it and returns
    the task, so the session is removed from storage immediately but its
    kernel process may still be running when the call returns. Callers that
    need the kernel's resources actually released — GPU memory, for
    instance — should use :meth:`aclose_session` instead, or await the
    returned task.

    Args:
        session_id: ID of the session to close

    Returns:
        The teardown task, or ``None`` when there was no kernel to tear
        down or no running event loop to schedule it on.
    """
    return self._close_session_sync(session_id, caller="close_session()")

rollback_factory_resources(data_manager, extensions)

Roll back resources from a factory result that cannot be retained.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def rollback_factory_resources(
    self,
    data_manager: object | None,
    extensions: dict[str, Any],
) -> BaseException | None:
    """Roll back resources from a factory result that cannot be retained."""
    return self._cleanup_unclaimed_session_resources(data_manager, extensions)

await_resource_cleanup() async

Wait for all async cleanup started by synchronous session closure.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def await_resource_cleanup(self) -> None:
    """Wait for all async cleanup started by synchronous session closure."""
    errors: list[Exception] = []
    cancelled: asyncio.CancelledError | None = None
    while self._resource_cleanup_tasks:
        tasks = tuple(self._resource_cleanup_tasks)
        try:
            await asyncio.gather(
                *(asyncio.shield(task) for task in tasks),
                return_exceptions=True,
            )
        except asyncio.CancelledError:
            # Cleanup continues under shield. Do not drop strong references:
            # the next drain must still observe and await these tasks.
            raise
        for task in tasks:
            self._on_resource_cleanup_done(task)
    for session_id, tasks in tuple(self._session_owned_cleanup_tasks.items()):
        tasks.difference_update(task for task in tuple(tasks) if task.done())
        if not tasks:
            self._session_owned_cleanup_tasks.pop(session_id, None)
            self._finalize_closed_session(session_id)
    errors.extend(self._resource_cleanup_errors)
    self._resource_cleanup_errors.clear()
    if self._resource_cleanup_cancellations:
        cancelled = self._resource_cleanup_cancellations[0]
        self._resource_cleanup_cancellations.clear()
    if cancelled is not None:
        if errors:
            cancelled.add_note(str(ExceptionGroup("Additional resource cleanup failures.", errors)))
        raise cancelled
    if errors:
        raise ExceptionGroup("Session resource cleanup failed.", errors)

aclose_session(session_id, *, expected_generation=None, expected_closing_session=None) async

Close a session and wait for its kernel to actually shut down.

The awaitable counterpart to :meth:close_session. Prefer this wherever the kernel's resources must be released before proceeding — freeing GPU memory, tearing down a batch's child sessions, or reclaiming capacity before starting new work.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def aclose_session(
    self,
    session_id: str,
    *,
    expected_generation: int | None = None,
    expected_closing_session: Session | None = None,
) -> None:
    """Close a session and wait for its kernel to actually shut down.

    The awaitable counterpart to :meth:`close_session`. Prefer this
    wherever the kernel's resources must be released before proceeding —
    freeing GPU memory, tearing down a batch's child sessions, or
    reclaiming capacity before starting new work.
    """
    shutdown_task, session = self._claim_session_close(
        session_id,
        caller="aclose_session()",
        expected_generation=expected_generation,
        expected_closing_session=expected_closing_session,
    )
    cleanup_task: asyncio.Task[None] | None = None
    if session is not None:
        cleanup_task = self._schedule_session_cleanup_after_operations(session)
    elif expected_generation is None or expected_closing_session is not None:
        with self._session_lifecycle_lock:
            cleanup_task = self._session_resource_cleanup_tasks.get(session_id)
            pending_session = self._pending_session_resource_cleanup.get(session_id)
            closing_session = self._closing_sessions.get(session_id)
        if cleanup_task is None and pending_session is not None:
            session = pending_session
            cleanup_task = self._schedule_session_cleanup_after_operations(session)
        elif cleanup_task is not None or self._session_owned_cleanup_tasks.get(session_id):
            session = closing_session

    async def finish_cleanup() -> None:
        cleanup_error: BaseException | None = None
        try:
            if cleanup_task is not None:
                _ = await cleanup_task
            await self._await_session_owned_cleanup(session_id)
        except asyncio.CancelledError as exc:
            cleanup_error = exc
        except Exception as exc:
            cleanup_error = exc
        finally:
            if session is not None:
                self._track_session_cleanup_tasks(session)
                self._finalize_closed_session(session_id)
        if shutdown_task is not None:
            _ = await shutdown_task
        else:
            await self.await_kernel_shutdown(session_id)
        if cleanup_error is not None:
            raise cleanup_error

    completion = asyncio.create_task(finish_cleanup())
    cancelled: asyncio.CancelledError | None = None
    try:
        while True:
            try:
                await asyncio.shield(completion)
                break
            except asyncio.CancelledError as exc:
                cancelled = cancelled or exc
                if completion.done():
                    break
    finally:
        if cancelled is not None:
            if completion.done() and not completion.cancelled():
                try:
                    completion.result()
                except BaseException as cleanup_error:
                    cancelled.add_note(f"Session cleanup also failed: {cleanup_error!r}")
            raise cancelled

aclose_all_sessions() async

Close every active session and wait for all kernel/resource teardown.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def aclose_all_sessions(self) -> None:
    """Close every active session and wait for all kernel/resource teardown."""
    with self._session_lifecycle_lock:
        sessions: list[tuple[str, int | None, Session | None]] = [
            (session_id, self._adopt_session_generation_locked(session_id), None)
            for session_id in self.storage.list_all()
        ]
        active_session_ids = {session_id for session_id, _, _ in sessions}
        sessions.extend(
            (session_id, None, session)
            for session_id, session in self._closing_sessions.items()
            if session_id not in active_session_ids
        )
    tasks = []
    for session_id, generation, closing_session in sessions:
        kwargs: dict[str, Any] = {"expected_generation": generation}
        if closing_session is not None:
            kwargs["expected_closing_session"] = closing_session
        tasks.append(asyncio.create_task(self.aclose_session(session_id, **kwargs)))
    outer_cancellation: asyncio.CancelledError | None = None
    results: list[BaseException | None] = []
    if tasks:
        completion = asyncio.gather(*tasks, return_exceptions=True)
        try:
            while True:
                try:
                    results = list(await asyncio.shield(completion))
                    break
                except asyncio.CancelledError as exc:
                    outer_cancellation = outer_cancellation or exc
                    if completion.done():
                        results = list(completion.result())
                        break
        finally:
            for (session_id, _, _), task in zip(sessions, tasks):
                if task.done() and not task.cancelled() and task.exception() is not None:
                    LOGGER.error("Failed to close session %s: %s", session_id, task.exception())
    errors: list[Exception] = [
        RuntimeError(f"Failed to close session {session_id}: {result}")
        for (session_id, _, _), result in zip(sessions, results)
        if isinstance(result, Exception)
    ]
    failed_session_ids = {
        session_id for (session_id, _, _), result in zip(sessions, results) if isinstance(result, Exception)
    }
    cancelled = next(
        (result for result in results if isinstance(result, asyncio.CancelledError)),
        outer_cancellation,
    )
    while self._kernel_shutdown_tasks:
        shutdowns = tuple(self._kernel_shutdown_tasks.values())
        try:
            await asyncio.gather(*(asyncio.shield(task) for task in shutdowns), return_exceptions=True)
        except asyncio.CancelledError as exc:
            cancelled = cancelled or exc
    resource_cleanup = asyncio.create_task(self.await_resource_cleanup())
    while True:
        try:
            await asyncio.shield(resource_cleanup)
            break
        except asyncio.CancelledError as exc:
            cancelled = cancelled or exc
            if resource_cleanup.done():
                break
        except Exception as exc:
            errors.append(exc)
            break
    with self._session_lifecycle_lock:
        remaining_session_ids = sorted(self._closing_sessions)
    for _ in range(3):
        if not remaining_session_ids:
            break
        retry_sessions = []
        with self._session_lifecycle_lock:
            retry_sessions = [
                (session_id, self._closing_sessions[session_id])
                for session_id in remaining_session_ids
                if session_id in self._closing_sessions
            ]
        retry_tasks = [
            asyncio.create_task(self.aclose_session(session_id, expected_closing_session=closing_session))
            for session_id, closing_session in retry_sessions
        ]
        retry_results: list[BaseException | None] = []
        if retry_tasks:
            completion = asyncio.gather(*retry_tasks, return_exceptions=True)
            while True:
                try:
                    retry_results = list(await asyncio.shield(completion))
                    break
                except asyncio.CancelledError as exc:
                    cancelled = cancelled or exc
                    if completion.done():
                        retry_results = list(completion.result())
                        break
        for (session_id, _), result in zip(retry_sessions, retry_results):
            if isinstance(result, asyncio.CancelledError):
                cancelled = cancelled or result
            elif isinstance(result, Exception) and session_id not in failed_session_ids:
                errors.append(RuntimeError(f"Failed to close session {session_id}: {result}"))
                failed_session_ids.add(session_id)
        with self._session_lifecycle_lock:
            remaining_session_ids = sorted(self._closing_sessions)
    unreported_session_ids = [
        session_id for session_id in remaining_session_ids if session_id not in failed_session_ids
    ]
    if unreported_session_ids:
        errors.append(RuntimeError(f"Session teardown did not complete for: {', '.join(unreported_session_ids)}"))
    if cancelled is not None:
        if errors:
            cancelled.add_note(str(ExceptionGroup("Additional session cleanup failures.", errors)))
        raise cancelled
    if errors:
        raise ExceptionGroup("One or more sessions failed to close.", errors)

get_kernel_generation(session_id)

Return the generation id of the session's current kernel.

Generation ids are globally unique and monotonically increasing across the lifetime of this manager: restarting a session's kernel always yields a strictly greater id, and an id is never reused. Callers can therefore cache per-kernel state keyed on this value and detect a rebuilt kernel by comparing against the value they captured.

Returns None when the session has no live kernel.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def get_kernel_generation(self, session_id: str) -> Optional[int]:
    """Return the generation id of the session's current kernel.

    Generation ids are globally unique and monotonically increasing across
    the lifetime of this manager: restarting a session's kernel always
    yields a strictly greater id, and an id is never reused. Callers can
    therefore cache per-kernel state keyed on this value and detect a
    rebuilt kernel by comparing against the value they captured.

    Returns ``None`` when the session has no live kernel.
    """
    return self._kernel_generations.get(session_id)

is_kernel_bootstrapped(session_id, key)

Whether key has been completed against the session's current kernel.

Returns False when the session has no kernel, or when its kernel has been rebuilt since the bootstrap step ran.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def is_kernel_bootstrapped(self, session_id: str, key: str) -> bool:
    """Whether ``key`` has been completed against the session's *current* kernel.

    Returns ``False`` when the session has no kernel, or when its kernel
    has been rebuilt since the bootstrap step ran.
    """
    generation = self._kernel_generations.get(session_id)
    if generation is None:
        return False
    return key in self._kernel_bootstrap_state.get(generation, frozenset())

mark_kernel_bootstrapped(session_id, key)

Record that key has been completed against the current kernel.

Returns False (and records nothing) when the session has no live kernel, so the caller's work will be retried against the next one.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def mark_kernel_bootstrapped(self, session_id: str, key: str) -> bool:
    """Record that ``key`` has been completed against the current kernel.

    Returns ``False`` (and records nothing) when the session has no live
    kernel, so the caller's work will be retried against the next one.
    """
    generation = self._kernel_generations.get(session_id)
    if generation is None:
        return False
    self._kernel_bootstrap_state.setdefault(generation, set()).add(key)
    return True

get_artifact_record(session_id, token)

Look up an artifact for the download endpoint.

Returns None for unknown sessions/tokens and for tokens whose backing file has been removed — both surface as 404 at the HTTP layer.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def get_artifact_record(self, session_id: str, token: str) -> Optional[_ArtifactRecord]:
    """Look up an artifact for the download endpoint.

    Returns ``None`` for unknown sessions/tokens *and* for tokens whose
    backing file has been removed — both surface as 404 at the HTTP
    layer.
    """
    record = self._session_artifacts.get(session_id, {}).get(token)
    if record is None:
        return None
    if not record.path.is_file():
        return None
    return record

find_artifact_by_name(session_id, artifact_name)

Find a registered artifact by its relative name within a session.

Looks up a previously registered artifact (i.e. one that was discovered during a snapshot-diff after an execute) by its relative filename. This is used by the publish pipeline to resolve the local path of an artifact before pushing it to remote storage.

Parameters:

Name Type Description Default
session_id str

The session that owns the artifact.

required
artifact_name str

Relative filename as registered (e.g. "results.csv" or "subdir/report.pdf").

required

Returns:

Name Type Description
The Optional[_ArtifactRecord]

class:_ArtifactRecord if found and the backing file still

Optional[_ArtifactRecord]

exists, otherwise None.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def find_artifact_by_name(self, session_id: str, artifact_name: str) -> Optional[_ArtifactRecord]:
    """Find a registered artifact by its relative name within a session.

    Looks up a previously registered artifact (i.e. one that was discovered
    during a snapshot-diff after an execute) by its relative filename.  This
    is used by the publish pipeline to resolve the local path of an artifact
    before pushing it to remote storage.

    Args:
        session_id: The session that owns the artifact.
        artifact_name: Relative filename as registered (e.g. ``"results.csv"``
            or ``"subdir/report.pdf"``).

    Returns:
        The :class:`_ArtifactRecord` if found and the backing file still
        exists, otherwise ``None``.
    """
    registry = self._session_artifacts.get(session_id, {})
    for record in registry.values():
        if record.name == artifact_name:
            if record.path.is_file():
                return record
            return None
    return None

start_background_execution_for_session(session_id, code, timeout, working_dir=None) async

Start execution and return immediately with a background job id.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def start_background_execution_for_session(
    self, session_id: str, code: str, timeout: float, working_dir: Optional[str] = None
) -> dict[str, str]:
    """Start execution and return immediately with a background job id."""
    async with self.session_resource_operation(session_id):
        return await self._start_background_execution_for_session(session_id, code, timeout, working_dir)

start_promoted_execution_for_session(session_id, code, timeout, promotion_threshold_s, working_dir=None) async

Execute code synchronously, promoting to a background job if it exceeds the threshold.

Starts executing like the normal foreground path. If the kernel reaches idle within promotion_threshold_s seconds the result is returned as a (stdout, stderr, success, displays, artifacts) tuple — identical to :meth:execute_code_for_session. If the threshold expires while the kernel is still busy, the in-flight execution is registered as a :class:_BackgroundJob and a job-handle dict is returned instead.

Parameters:

Name Type Description Default
session_id str

Session identifier.

required
code str

Python code to execute. This method applies the outputs and token preambles internally (same as the foreground path), so callers should pass the raw user code — not pre-preambled.

required
timeout float

Total execution timeout in seconds.

required
promotion_threshold_s float

Seconds to wait before promoting.

required
working_dir Optional[str]

Optional working directory for the kernel.

None

Returns:

Type Description
Tuple[str, str, bool, list[dict], list[dict]] | dict[str, Any]

Either the 5-tuple (stdout, stderr, success, displays, artifacts)

Tuple[str, str, bool, list[dict], list[dict]] | dict[str, Any]

when the execution completes within the threshold, or a dict with

Tuple[str, str, bool, list[dict], list[dict]] | dict[str, Any]

job_id, status, session_id, and promoted=True.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def start_promoted_execution_for_session(
    self,
    session_id: str,
    code: str,
    timeout: float,
    promotion_threshold_s: float,
    working_dir: Optional[str] = None,
) -> "Tuple[str, str, bool, list[dict], list[dict]] | dict[str, Any]":
    """Execute code synchronously, promoting to a background job if it exceeds the threshold.

    Starts executing like the normal foreground path. If the kernel reaches
    ``idle`` within *promotion_threshold_s* seconds the result is returned
    as a ``(stdout, stderr, success, displays, artifacts)`` tuple — identical
    to :meth:`execute_code_for_session`.  If the threshold expires while the
    kernel is still busy, the in-flight execution is registered as a
    :class:`_BackgroundJob` and a job-handle dict is returned instead.

    Args:
        session_id: Session identifier.
        code: Python code to execute.  This method applies the outputs and
            token preambles internally (same as the foreground path), so
            callers should pass the raw user code — **not** pre-preambled.
        timeout: Total execution timeout in seconds.
        promotion_threshold_s: Seconds to wait before promoting.
        working_dir: Optional working directory for the kernel.

    Returns:
        Either the 5-tuple ``(stdout, stderr, success, displays, artifacts)``
        when the execution completes within the threshold, or a dict with
        ``job_id``, ``status``, ``session_id``, and ``promoted=True``.
    """
    if timeout <= promotion_threshold_s:
        raise ValueError(
            f"timeout ({timeout}s) must be greater than promotion_threshold_s ({promotion_threshold_s}s)"
        )
    async with self.session_resource_operation(session_id):
        return await self._start_promoted_execution_for_session(
            session_id,
            code,
            timeout,
            promotion_threshold_s,
            working_dir,
        )

check_background_job(job_id, caller_identity=None)

Return current status/output for a background job.

Parameters:

Name Type Description Default
job_id str

Background job identifier.

required
caller_identity Optional[str]

When provided, the caller's user identity is compared to the identity that submitted the job. Both a missing job and an identity mismatch raise ValueError("Job {job_id} not found") so that callers cannot use differing error messages to probe whether a job ID exists.

None
Source code in src/agora_workbench/code_execution/sessions/manager.py
def check_background_job(self, job_id: str, caller_identity: Optional[str] = None) -> dict[str, Any]:
    """Return current status/output for a background job.

    Args:
        job_id: Background job identifier.
        caller_identity: When provided, the caller's user identity is compared to
            the identity that submitted the job. Both a missing job and an identity
            mismatch raise ``ValueError("Job {job_id} not found")`` so that callers
            cannot use differing error messages to probe whether a job ID exists.
    """
    job = self._background_jobs.get(job_id)
    if not job:
        raise ValueError(f"Job {job_id} not found")

    if caller_identity and job.user_identity and caller_identity != job.user_identity:
        raise ValueError(f"Job {job_id} not found")

    now = time.monotonic()
    elapsed_seconds = (job.completed_at or now) - job.start_time
    result: dict[str, Any] = {
        "job_id": job.job_id,
        "session_id": job.session_id,
        "status": job.status,
        "elapsed_seconds": round(elapsed_seconds, 3),
    }

    if job.status == "running":
        result["stdout"] = "".join(job.stdout_parts)
        result["stderr"] = "".join(job.stderr_parts)
        return result

    result["stdout"] = "".join(job.stdout_parts)
    result["stderr"] = "".join(job.stderr_parts)
    result["success"] = job.success
    if job.error:
        result["error"] = job.error
    # Artifacts are populated by _finalize_background_artifacts when the
    # job reaches a terminal state; carry them here without download URLs
    # — URL composition happens in the MCP server layer where the public
    # base URL is known (matches the foreground path's contract).
    result["artifacts"] = list(job.artifacts)
    return result

await_background_job(job_id) async

Wait for a background job's task to reach a terminal state, then return its status.

Returns None if the job was never registered or has already been purged. Used by the activity publisher to emit job_finished events; callers must treat exceptions on the underlying task as terminal (a failed task still leaves _BackgroundJob.status set by _collect_background_job).

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def await_background_job(self, job_id: str) -> Optional[dict[str, Any]]:
    """Wait for a background job's task to reach a terminal state, then return its status.

    Returns ``None`` if the job was never registered or has already been purged.
    Used by the activity publisher to emit ``job_finished`` events; callers must
    treat exceptions on the underlying task as terminal (a failed task still
    leaves ``_BackgroundJob.status`` set by ``_collect_background_job``).
    """
    job = self._background_jobs.get(job_id)
    if job is None or job.task is None:
        return None
    try:
        await job.task
    except Exception:
        LOGGER.debug("Background job task raised for job_id=%s", job_id, exc_info=True)
    try:
        return self.check_background_job(job_id)
    except ValueError:
        return None

session_resource_operation(session_id)

Keep session-owned data resources alive for one admitted operation.

Source code in src/agora_workbench/code_execution/sessions/manager.py
def session_resource_operation(self, session_id: str) -> AsyncContextManager[None]:
    """Keep session-owned data resources alive for one admitted operation."""
    with self._session_lifecycle_lock:
        session_generation = self._session_generations.get(session_id)
        if (
            session_generation is None
            and session_id not in self._closing_session_ids
            and self.storage.retrieve(session_id) is not None
        ):
            session_generation = self._adopt_session_generation_locked(session_id)

    @asynccontextmanager
    async def operation() -> AsyncIterator[None]:
        with self._session_lifecycle_lock:
            if (
                session_generation is None
                or self._session_generations.get(session_id) != session_generation
                or session_id in self._closing_session_ids
                or self.storage.retrieve(session_id) is None
            ):
                raise ValueError(f"Session {session_id} not found or is closing")
            event = self._session_resources_drained.get(session_id)
            if event is None:
                event = asyncio.Event()
                self._session_resources_drained[session_id] = event
            event.clear()
            self._session_resource_users[session_id] = self._session_resource_users.get(session_id, 0) + 1
        try:
            yield
        finally:
            pending_cleanup: Session | None = None
            with self._session_lifecycle_lock:
                remaining = self._session_resource_users[session_id] - 1
                if remaining:
                    self._session_resource_users[session_id] = remaining
                else:
                    self._session_resource_users.pop(session_id, None)
                    drained = self._session_resources_drained.pop(session_id, None)
                    if drained is not None:
                        drained.set()
                    pending_cleanup = self._pending_session_resource_cleanup.get(session_id)
            if pending_cleanup is not None:
                self._schedule_session_cleanup_after_operations(pending_cleanup)
            self._finalize_closed_session(session_id)

    return operation()

execute_code_for_session(session_id, code, timeout, working_dir=None) async

Execute code in the session's Jupyter kernel.

Concurrent calls for the same session_id are serialized by a per-session asyncio.Lock so they cannot race on the shared Jupyter KernelClient iopub stream.

Parameters:

Name Type Description Default
session_id str

Session identifier

required
code str

Python code to execute

required
timeout float

Execution timeout in seconds

required
working_dir Optional[str]

Optional working directory for kernel

None

Returns:

Type Description
str

Tuple of (stdout, stderr, success, displays, artifacts).

str

displays contains rich-output payloads emitted by the kernel,

bool

with MIME type, data, and metadata fields. artifacts contains

list[dict]

metadata for each new or modified file under the session output

list[dict]

directory, including its name, size, MIME type, modification time,

Tuple[str, str, bool, list[dict], list[dict]]

and download token.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def execute_code_for_session(
    self, session_id: str, code: str, timeout: float, working_dir: Optional[str] = None
) -> Tuple[str, str, bool, list[dict], list[dict]]:
    """
    Execute code in the session's Jupyter kernel.

    Concurrent calls for the same ``session_id`` are serialized by a
    per-session asyncio.Lock so they cannot race on the shared
    Jupyter ``KernelClient`` iopub stream.

    Args:
        session_id: Session identifier
        code: Python code to execute
        timeout: Execution timeout in seconds
        working_dir: Optional working directory for kernel

    Returns:
        Tuple of ``(stdout, stderr, success, displays, artifacts)``.
        ``displays`` contains rich-output payloads emitted by the kernel,
        with MIME type, data, and metadata fields. ``artifacts`` contains
        metadata for each new or modified file under the session output
        directory, including its name, size, MIME type, modification time,
        and download token.
    """
    resource_operation = self.session_resource_operation(session_id)
    try:
        await resource_operation.__aenter__()
    except ValueError:
        return (
            "",
            (
                f"Session {session_id} is no longer available (expired or cleaned up). "
                "Please create a new session and retry."
            ),
            False,
            [],
            [],
        )
    try:
        return await self._execute_code_for_admitted_session(session_id, code, timeout, working_dir)
    finally:
        await resource_operation.__aexit__(None, None, None)

await_kernel_shutdown(session_id) async

Wait for any in-flight teardown of this session's kernel to finish.

No-op when none is running. Shielded, so a cancelled caller does not cancel the teardown itself; failures are already reported by the done-callback, so awaiting is purely for sequencing.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def await_kernel_shutdown(self, session_id: str) -> None:
    """Wait for any in-flight teardown of this session's kernel to finish.

    No-op when none is running. Shielded, so a cancelled caller does not
    cancel the teardown itself; failures are already reported by the
    done-callback, so awaiting is purely for sequencing.
    """
    task = self._kernel_shutdown_tasks.get(session_id)
    if task is None or task.done():
        return
    try:
        await asyncio.shield(task)
    except asyncio.CancelledError:
        raise
    except Exception:
        LOGGER.debug("Awaited kernel shutdown for session %s failed", session_id, exc_info=True)

cleanup_idle_kernels(max_idle_time=3600.0) async

Cleanup kernels that have been idle for too long.

Source code in src/agora_workbench/code_execution/sessions/manager.py
async def cleanup_idle_kernels(self, max_idle_time: float = 3600.0):
    """Cleanup kernels that have been idle for too long."""
    now = time.time()
    with self._session_lifecycle_lock:
        idle_sessions = [
            (
                sid,
                self._kernel_session_generations.get(sid),
                self._kernel_generations.get(sid),
            )
            for sid, last_used in self._kernel_last_used.items()
            if now - last_used > max_idle_time
        ]

    for session_id, session_generation, kernel_generation in idle_sessions:
        LOGGER.info(f"Cleaning up idle kernel for session {session_id}")
        joined_shutdown: "Optional[asyncio.Task[None]]" = None
        async with self._kernel_execution_lock(session_id):
            with self._session_lifecycle_lock:
                last_used = self._kernel_last_used.get(session_id)
                if (
                    last_used is None
                    or time.time() - last_used <= max_idle_time
                    or self._kernel_session_generations.get(session_id) != session_generation
                    or self._kernel_generations.get(session_id) != kernel_generation
                ):
                    continue
                existing_shutdown = self._kernel_shutdown_tasks.get(session_id)
                joined_existing = existing_shutdown is not None and not existing_shutdown.done()
                shutdown_task = self._schedule_kernel_shutdown(
                    session_id,
                    caller="cleanup_idle_kernels()",
                )
            if shutdown_task is None:
                continue
            if joined_existing:
                # A teardown started elsewhere may be waiting for executions to
                # drain, and this sweep's own execute-lock lease keeps that drain
                # from completing. Join it only after releasing the lease.
                joined_shutdown = shutdown_task
            else:
                _ = await asyncio.shield(shutdown_task)
        if joined_shutdown is not None:
            _ = await asyncio.shield(joined_shutdown)

list_sessions()

List all active sessions with metadata.

Returns:

Type Description
list[dict[str, Any]]

List of session info dicts

Source code in src/agora_workbench/code_execution/sessions/manager.py
def list_sessions(self) -> list[dict[str, Any]]:
    """
    List all active sessions with metadata.

    Returns:
        List of session info dicts
    """
    self._maybe_cleanup()

    sessions = self.storage.list_all()

    result = []
    for session in sessions.values():
        result.append(session.get_info())

    return result

Session Configuration

agora_workbench.code_execution.sessions.manager.SessionConfig(max_sessions=100, timeout_minutes=30, cleanup_interval_seconds=300, storage_backend=None, data_manager_factory=None, kernel_network_mode='inherit')

Configuration for session manager.

Initialize session manager configuration.

Parameters:

Name Type Description Default
max_sessions int

Maximum number of concurrent sessions.

100
timeout_minutes int

Idle time after which a session is cleaned up.

30
cleanup_interval_seconds int

Minimum interval between cleanup sweeps.

300
storage_backend Optional[SessionStorageBackend]

Optional session storage backend.

None
data_manager_factory Optional[Callable[[SessionContext], DataLakeDataManager | SessionResources]]

Optional callable invoked once per session to build its :class:DataLakeDataManager, receiving a :class:SessionContext describing the session being created. Use it to supply a customized manager — a different credential, a custom artifact resolver, extra fetchers, or configuration derived from user_identity / user_token. When the server mounts a catalog, DataLakeDataManager instances are automatically bound to its caller-scoped resolver through the CatalogAwareDataManager protocol.

The factory must return a fresh data manager or :class:SessionResources bundle per call. The session takes ownership of the manager and extensions and cleans them up when the session ends, so returning shared resources would let the first session torn down destroy resources still in use by others.

When omitted, each session builds a default DataLakeDataManager(), matching previous behavior.

A factory that returns None — or any object without a cleanup() method — raises TypeError from create_session, rather than silently falling back to the default manager.

None
kernel_network_mode Literal['inherit', 'isolated']

"inherit" lets kernels use the host network. "isolated" launches each kernel in an empty Linux network namespace and uses Unix IPC sockets for Jupyter communication.

'inherit'
Source code in src/agora_workbench/code_execution/sessions/manager.py
def __init__(
    self,
    max_sessions: int = 100,
    timeout_minutes: int = 30,
    cleanup_interval_seconds: int = 300,  # 5 minutes
    storage_backend: Optional[SessionStorageBackend] = None,
    data_manager_factory: Optional[Callable[[SessionContext], "DataLakeDataManager | SessionResources"]] = None,
    kernel_network_mode: Literal["inherit", "isolated"] = "inherit",
):
    """
    Initialize session manager configuration.

    Args:
        max_sessions: Maximum number of concurrent sessions.
        timeout_minutes: Idle time after which a session is cleaned up.
        cleanup_interval_seconds: Minimum interval between cleanup sweeps.
        storage_backend: Optional session storage backend.
        data_manager_factory: Optional callable invoked once per session to
            build its :class:`DataLakeDataManager`, receiving a
            :class:`SessionContext` describing the session being created.
            Use it to supply a customized manager — a different credential,
            a custom artifact resolver, extra fetchers, or configuration
            derived from ``user_identity`` / ``user_token``.
            When the server mounts a catalog, ``DataLakeDataManager``
            instances are automatically bound to its caller-scoped
            resolver through the ``CatalogAwareDataManager`` protocol.

            The factory **must return a fresh data manager or
            :class:`SessionResources` bundle per call**. The session takes
            ownership of the manager and extensions and cleans them up when
            the session ends, so returning shared resources would let the
            first session torn down destroy resources still in use by
            others.

            When omitted, each session builds a default
            ``DataLakeDataManager()``, matching previous behavior.

            A factory that returns ``None`` — or any object without a
            ``cleanup()`` method — raises ``TypeError`` from
            ``create_session``, rather than silently falling back to the
            default manager.
        kernel_network_mode: ``"inherit"`` lets kernels use the host network.
            ``"isolated"`` launches each kernel in an empty Linux network
            namespace and uses Unix IPC sockets for Jupyter communication.
    """
    self.max_sessions = max_sessions
    self.timeout = timedelta(minutes=timeout_minutes)
    self.cleanup_interval = timedelta(seconds=cleanup_interval_seconds)
    self.storage_backend = storage_backend or InMemoryStorage()
    self.data_manager_factory = data_manager_factory
    if kernel_network_mode not in ("inherit", "isolated"):
        raise ValueError("kernel_network_mode must be 'inherit' or 'isolated'")
    self.kernel_network_mode = kernel_network_mode