Skip to content

Data Access

Records and protocols

agora_workbench.data_lake

Backend-neutral public contracts for data-lake access.

ArtifactNotFoundError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

The requested catalog artifact does not exist.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

BackendUnavailableError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

The catalog backend cannot currently serve the request.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

ConflictError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

A concurrent mutation conflicted with the requested operation.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

DataLakeError(message, *, resource_id=None, operation=None)

Bases: Exception

Base class for errors crossing the public data-lake boundary.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

DataLakeErrorCode

Bases: StrEnum

Stable machine-readable categories for data-lake failures.

InvalidRequestError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

The request is malformed or internally inconsistent.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

PermissionDeniedError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

The current caller is not permitted to perform an operation.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

PreconditionFailedError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

An ownership, generation, or revision precondition was not satisfied.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

ReconciliationError(message, *, failures=None, resource_id=None, operation=None)

Bases: DataLakeError

Recovery could not safely classify or clean an interrupted operation.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    failures: Mapping[str, str] | None = None,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message, resource_id=resource_id, operation=operation)
    self.failures = MappingProxyType(dict(failures or {}))

RetryExhaustedError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

Bounded optimistic-concurrency retries were exhausted.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

TransferCancelledError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

A transfer was cancelled before it committed its destination.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

TransferChecksumError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

Transferred bytes did not match the required checksum.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

TransferLimitError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

A transfer exceeded its object-size or caller-quota bound.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

TransferTimeoutError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

A transfer exceeded its configured end-to-end timeout.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

UnsafePathError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

A local or provider path escaped its configured namespace.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

UnsupportedOperationError(message, *, resource_id=None, operation=None)

Bases: DataLakeError

The provider or effective caller capabilities do not support an operation.

Source code in src/agora_workbench/data_lake/errors.py
def __init__(
    self,
    message: str,
    *,
    resource_id: str | None = None,
    operation: str | None = None,
) -> None:
    super().__init__(message)
    self.message = message
    self.resource_id = resource_id
    self.operation = operation

AzureBlobScope(account, container, prefix='') dataclass

Configured Azure account/container/prefix boundary for transfers.

from_uri(uri) classmethod

Build a validated scope from any supported Azure URI form.

Source code in src/agora_workbench/data_lake/identity.py
@classmethod
def from_uri(cls, uri: str) -> "AzureBlobScope":
    """Build a validated scope from any supported Azure URI form."""
    account, container, prefix = parse_azure_uri(uri)
    return cls(account, container, prefix)

contains(account, container, object_path)

Return whether an object lies on this scope's prefix boundary.

Source code in src/agora_workbench/data_lake/identity.py
def contains(self, account: str, container: str, object_path: str) -> bool:
    """Return whether an object lies on this scope's prefix boundary."""
    if (account, container) != (self.account, self.container):
        return False
    if not self.prefix:
        return True
    if object_path.startswith(f"{self.prefix}/"):
        return True
    return not self._descendants_only and object_path == self.prefix

ArtifactPresentation(name, description=None, media_type=None, size_bytes=None) dataclass

Human-facing catalog metadata.

ArtifactReference(artifact_id, source_id, revision=None) dataclass

Stable logical identity, optionally pinned to a provider-honored revision.

Providers and adapters that accept a non-None revision must resolve that exact retained revision or reject the request explicitly; they must not silently return the current revision.

is_current property

Whether the reference follows the current artifact revision.

CatalogAuthorizationRequest(operation, source_id, reference=None) dataclass

One source- or artifact-scoped authorization check.

CatalogArtifact(reference, presentation, locator=None, download=None, metadata=dict(), revision=None, content_revision=None, metadata_revision=None, checksum_sha256=None, deleted_at=None, score=None) dataclass

Backend-neutral artifact returned by catalog operations.

is_deleted property

Whether the catalog record is a deletion tombstone.

CatalogOperation

Bases: StrEnum

Catalog operations a provider or writer may support.

CatalogPolicyMode

Bases: StrEnum

Granularity at which caller policy is enforced.

DownloadInfo(url, filename=None, expires_at=None) dataclass

Optional presentation-layer download information.

ListRequest(source_ids=(), page=PageRequest(), filters=dict()) dataclass

Catalog listing request.

Page(items, next_cursor=None) dataclass

Bases: Generic[T]

One page of results and an opaque provider cursor.

PageRequest(limit=50, cursor=None) dataclass

Cursor-based pagination request, capped at :data:MAX_PAGE_LIMIT.

RequestContext(request_id=None, caller_id=None, attributes=dict()) dataclass

Request-scoped caller metadata copied for safe propagation.

ResolvedArtifact(reference, locator) dataclass

A logical artifact reference resolved to a physical locator.

ResourceLease(resource, ownership=ResourceOwnership.BORROWED) dataclass

Bases: Generic[T]

A resource paired with its explicit cleanup ownership.

should_close property

Whether the recipient is responsible for closing the resource.

ResourceOwnership

Bases: StrEnum

Whether the recipient owns cleanup of a supplied resource.

SearchRequest(query, source_ids=(), page=PageRequest(), filters=dict()) dataclass

Catalog search request.

SourceCapabilities(source_id, supported_operations) dataclass

Read operations authoritatively supported by one provider source.

supports(operation)

Return whether the source reports support for operation.

Source code in src/agora_workbench/data_lake/models.py
def supports(self, operation: CatalogOperation) -> bool:
    """Return whether the source reports support for *operation*."""
    return operation in self.supported_operations

StorageLocator(uri) dataclass

Physical locator understood by a fetcher or storage provider.

CatalogManifest(version, generation, artifacts, commit_fences=()) dataclass

An authoritative manifest generation.

to_mapping()

Serialize a canonical JSON-compatible manifest mapping.

Source code in src/agora_workbench/data_lake/manifest.py
def to_mapping(self) -> dict[str, object]:
    """Serialize a canonical JSON-compatible manifest mapping."""
    return {
        "version": self.version,
        "generation": self.generation,
        "artifacts": [artifact.to_mapping() for artifact in self.artifacts],
        **({"commit_fences": list(self.commit_fences)} if self.commit_fences else {}),
    }

ManifestArtifact(path, artifact_id=None, name=None, description=None, domain=None, media_type=None, size_bytes=None, content_revision=None, metadata_revision=None, checksum_sha256=None, aliases=(), storage_path=None, revision_id=None, revisions=(), removals=(), ownership=None, deleted_at=None, provenance=None) dataclass

One authoritative artifact registration.

to_mapping()

Serialize the artifact using the stable manifest field names.

Source code in src/agora_workbench/data_lake/manifest.py
def to_mapping(self) -> dict[str, object]:
    """Serialize the artifact using the stable manifest field names."""
    value = asdict(self)
    value["ownership"] = self.ownership.value if self.ownership is not None else None
    value["revisions"] = [
        {**asdict(revision), "ownership": revision.ownership.value} for revision in self.revisions
    ]
    return {key: item for key, item in value.items() if item is not None and item != ()}

ManifestOwnership

Bases: StrEnum

Ownership of bytes referenced by a manifest revision.

ManifestProvenance(operation_id, kind, created_at, caller_id=None, source_uri=None, session_id=None, output_name=None) dataclass

Audit provenance attached to a managed revision.

ManifestRemoval(operation_id, generation, revision_id, storage_path, garbage_collect, revisions=(), provenance=None) dataclass

Durable evidence of one committed tombstone operation.

ManifestRevision(revision_id, storage_path, content_revision, created_at, operation_id, ownership, checksum_sha256=None, size_bytes=None, version_token=None, committed_generation=None, provenance=None) dataclass

One immutable physical revision retained by the manifest.

ArtifactMetadata(name=None, description=None, domain=None, media_type=None, aliases=()) dataclass

Metadata supplied for a durable artifact registration.

AuthorizedManagedCatalogWriter(writer, authorizer)

Authorize every mutation before metadata or byte side effects.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, writer: ManagedCatalogWriter, authorizer: CatalogAuthorizer) -> None:
    self._writer = writer
    self._authorizer = authorizer

capabilities(context) async

Return write operations allowed for this session and source.

Source code in src/agora_workbench/data_lake/managed.py
async def capabilities(self, context: RequestContext) -> tuple[SourceCapabilities, ...]:
    """Return write operations allowed for this session and source."""
    allowed = {
        operation
        for operation in WRITE_OPERATIONS
        if await self._authorizer.authorize(
            CatalogAuthorizationRequest(operation, self._writer.source_id),
            context,
        )
    }
    return (SourceCapabilities(self._writer.source_id, frozenset(allowed)),) if allowed else ()

BlobManagedStorage(container_client, *, prefix='')

Blob-container adapter using create-only objects and ETag manifest CAS.

The caller owns the container client and its prefix. Transfers apply the shared safety options while this adapter retains lifecycle-specific ownership metadata and conditional state operations.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, container_client: object, *, prefix: str = "") -> None:
    self._container_client = container_client
    self._prefix = prefix.strip("/")

DeleteOutcome

Bases: StrEnum

Ownership-aware deletion postcondition.

LocalManagedStorage(root)

Create-exclusive local objects and atomically replaced manifests.

Writers use an advisory flock covering the full mutation. Files and containing directories are fsynced before the lock is released. Readers do not take the lock; atomic replacement means they observe either the old or new complete manifest.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, root: str | Path) -> None:
    self.root = Path(root).resolve()
    self.root.mkdir(parents=True, exist_ok=True)
    self._lock_path = self.root / _WRITER_LOCK_PATH

ManagedCatalogWriter(source_id, backend, *, max_conflict_retries=4, revision_retention=1, operation_lease_seconds=30.0, interruption_hook=None)

Storage-neutral managed write state machine.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(
    self,
    source_id: str,
    backend: _ManagedStorageBackend,
    *,
    max_conflict_retries: int = 4,
    revision_retention: int = 1,
    operation_lease_seconds: float = 30.0,
    interruption_hook: Callable[[str, str], None] | None = None,
) -> None:
    if not source_id:
        raise ValueError("source_id must be non-empty")
    if max_conflict_retries < 0:
        raise ValueError("max_conflict_retries must be non-negative")
    if revision_retention < 0:
        raise ValueError("revision_retention must be non-negative")
    if operation_lease_seconds <= 0:
        raise ValueError("operation_lease_seconds must be positive")
    self.source_id = source_id
    self._backend = backend
    self._manifest_path = RESERVED_MANIFEST_PATH
    self._max_conflict_retries = max_conflict_retries
    self._revision_retention = revision_retention
    self._operation_lease_seconds = operation_lease_seconds
    self._interruption_hook = interruption_hook

artifact_reference(path, artifact_id=None)

Return the normalized effective reference used for authorization.

Source code in src/agora_workbench/data_lake/managed.py
def artifact_reference(self, path: str, artifact_id: str | None = None) -> ArtifactReference:
    """Return the normalized effective reference used for authorization."""
    normalized = _normalize_catalog_path(path)
    return ArtifactReference(
        artifact_id or logical_artifact_id(self.source_id, normalized),
        self.source_id,
    )

read_manifest(minimum_generation=None) async

Read committed state directly from the writer's backing store.

Source code in src/agora_workbench/data_lake/managed.py
async def read_manifest(self, minimum_generation: int | None = None) -> CatalogManifest:
    """Read committed state directly from the writer's backing store."""
    raw, _ = await self._backend.read_json(self._manifest_path)
    manifest = (
        CatalogManifest.from_mapping(raw)
        if raw is not None
        else CatalogManifest(MANIFEST_VERSION, 0, ())  # internal empty baseline
    )
    if minimum_generation is not None and manifest.generation < minimum_generation:
        raise PreconditionFailedError(
            "The requested committed generation is not visible.",
            operation="read_manifest",
        )
    return manifest

register(request, context=RequestContext()) async

Snapshot caller-owned bytes without modifying or owning the source.

Source code in src/agora_workbench/data_lake/managed.py
async def register(
    self,
    request: RegisterArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Snapshot caller-owned bytes without modifying or owning the source."""
    operation_id = _validate_operation_id(request.operation_id)
    path = _normalize_catalog_path(request.path)
    storage_path = normalize_logical_path(request.storage_path)
    if is_reserved_provider_path(storage_path):
        raise InvalidRequestError("External registrations cannot target reserved managed paths.")
    if request.checksum_sha256 is None:
        raise InvalidRequestError(
            "External registration requires checksum_sha256 for an immutable snapshot.",
            operation="register",
        )
    if (
        request.transfer_options.expected_sha256 is not None
        and request.transfer_options.expected_sha256 != request.checksum_sha256
    ):
        raise InvalidRequestError(
            "TransferOptions checksum must match the registered checksum.",
            operation="register",
        )
    artifact_id = request.artifact_id or logical_artifact_id(self.source_id, path)
    revision_path = _revision_path(artifact_id, operation_id, ".data")
    provenance = self._provenance(CatalogOperation.REGISTER, operation_id, context, request)
    intent = self._intent(
        CatalogOperation.REGISTER,
        operation_id,
        path,
        artifact_id,
        request,
        owned_path=revision_path,
        storage_path=revision_path,
        checksum_sha256=request.checksum_sha256,
        metadata=request.metadata,
        provenance=provenance,
    )
    prior, lease_id = await self._begin_operation(operation_id, intent)
    if prior is not None:
        return await self._resume_remove_cleanup(prior)
    try:
        self._interrupt("after_intent", operation_id)
        await emit_transfer_diagnostic(
            request.transfer_options,
            TransferDiagnostic("register", "started", context, revision_path),
        )
        async with self._heartbeat(operation_id, lease_id):
            try:
                version = await self._backend.create_from_storage(
                    storage_path,
                    revision_path,
                    operation_id,
                    request.checksum_sha256,
                    request.transfer_options,
                )
            except FileNotFoundError as exc:
                raise ArtifactNotFoundError("External bytes were not found.", operation="register") from exc
            except _StorageConflict:
                version = await self._backend.exists(revision_path)
                if (
                    version is None
                    or version.operation_id != operation_id
                    or version.checksum_sha256 != request.checksum_sha256
                ):
                    raise ConflictError(
                        "Immutable snapshot path is occupied by bytes not owned by this operation.",
                        operation="register",
                    ) from None
            except BaseException as exc:
                await emit_transfer_diagnostic(
                    request.transfer_options,
                    TransferDiagnostic(
                        "register",
                        "failed",
                        context,
                        revision_path,
                        error_type=type(exc).__name__,
                    ),
                )
                raise
            if request.size_bytes is not None and request.size_bytes != version.size_bytes:
                await self._backend.delete_owned(revision_path, operation_id, version.token)
                error = PreconditionFailedError(
                    "Registered size did not match the immutable snapshot.",
                    operation="register",
                )
                await emit_transfer_diagnostic(
                    request.transfer_options,
                    TransferDiagnostic(
                        "register",
                        "failed",
                        context,
                        revision_path,
                        error_type=type(error).__name__,
                    ),
                )
                raise error
        await emit_transfer_diagnostic(
            request.transfer_options,
            TransferDiagnostic(
                "register",
                "completed",
                context,
                revision_path,
                version.size_bytes,
                version.checksum_sha256,
            ),
        )
        await self._update_operation(
            operation_id,
            lease_id,
            ownership_verified=True,
            owned_version_token=version.token,
        )
        self._interrupt("after_object", operation_id)
        await self._update_operation(operation_id, lease_id, state="commit_ready")
        revision = ManifestRevision(
            revision_id=operation_id,
            storage_path=revision_path,
            content_revision=request.content_revision or request.checksum_sha256,
            created_at=provenance.created_at,
            operation_id=operation_id,
            ownership=ManifestOwnership.MANAGED,
            checksum_sha256=request.checksum_sha256,
            size_bytes=version.size_bytes,
            version_token=version.token,
            provenance=provenance,
        )
        return await self._commit_revision(
            CatalogOperation.REGISTER,
            operation_id,
            path,
            artifact_id,
            request.metadata,
            revision,
            provenance,
            request.expected_generation,
            request.expected_revision_id,
            lease_id,
        )
    except BaseException:
        await self._abandon_operation(operation_id, lease_id)
        raise

upload(request, context=RequestContext()) async

Create an immutable managed revision and commit it to the manifest.

Source code in src/agora_workbench/data_lake/managed.py
async def upload(
    self,
    request: UploadArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Create an immutable managed revision and commit it to the manifest."""
    return await self._upload(request, context, CatalogOperation.UPLOAD)

promote(request, context=RequestContext()) async

Explicitly copy a scratch output into durable managed storage.

Source code in src/agora_workbench/data_lake/managed.py
async def promote(
    self,
    request: PromoteOutputRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Explicitly copy a scratch output into durable managed storage."""
    if not request.session_id or not request.output_name:
        raise InvalidRequestError("Promotion requires session_id and output_name.", operation="promote")
    return await self._upload(request, context, CatalogOperation.PROMOTE)

remove(request, context=RequestContext()) async

Commit a tombstone before collecting any managed bytes.

Source code in src/agora_workbench/data_lake/managed.py
async def remove(
    self,
    request: RemoveArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Commit a tombstone before collecting any managed bytes."""
    operation_id = _validate_operation_id(request.operation_id)
    if request.reference.source_id != self.source_id:
        raise ArtifactNotFoundError("Artifact was not found.", operation="remove")
    if request.reference.revision is not None:
        raise InvalidRequestError(
            "Managed removal uses expected_revision_id rather than a catalog revision number.",
            operation="remove",
        )
    intent = self._intent(
        CatalogOperation.REMOVE,
        operation_id,
        "",
        request.reference.artifact_id,
        request,
        provenance=self._provenance(CatalogOperation.REMOVE, operation_id, context, request),
    )
    prior, lease_id = await self._begin_operation(operation_id, intent)
    if prior is not None:
        return await self._resume_remove_cleanup(prior)
    provenance = self._provenance(CatalogOperation.REMOVE, operation_id, context, request)
    try:
        self._interrupt("after_intent", operation_id)
        await self._update_operation(operation_id, lease_id, state="commit_ready")
        return await self._remove_under_lease(request, provenance, operation_id, lease_id)
    except BaseException:
        await self._abandon_operation(operation_id, lease_id)
        raise

reconcile(*, grace_seconds=300.0) async

Claim abandoned operations before recovery or owned-orphan cleanup.

Source code in src/agora_workbench/data_lake/managed.py
async def reconcile(self, *, grace_seconds: float = 300.0) -> ReconciliationReport:
    """Claim abandoned operations before recovery or owned-orphan cleanup."""
    if grace_seconds < 0:
        raise ValueError("grace_seconds must be non-negative")
    recovered: list[str] = []
    removed: list[str] = []
    deferred: list[str] = []
    active_operation_ids: set[str] = set()
    failures: dict[str, str] = {}
    for intent_path, listed_intent in await self._backend.list_json(RESERVED_OPERATIONS_PREFIX.rstrip("/")):
        operation_id = str(listed_intent.get("operation_id", ""))
        try:
            async with self._backend.serialized():
                intent, token = await self._backend.read_json(intent_path)
                if intent is None:
                    continue
                receipt, _ = await self._backend.read_json(self._receipt_path(operation_id))
                if receipt is not None:
                    result = self._result_from_mapping(receipt)
                    resumed = await self._resume_remove_cleanup(result)
                    if resumed != result:
                        recovered.append(operation_id)
                    continue
                manifest, _ = await self._load()
                artifact = self._find_operation_artifact(manifest, operation_id)
                if artifact is not None:
                    pass
                else:
                    created_at = datetime.fromisoformat(str(intent["created_at"]).replace("Z", "+00:00"))
                    age = (datetime.now(timezone.utc) - created_at).total_seconds()
                    if age < grace_seconds or (
                        intent.get("state") in {"active", "commit_ready"} and not self._lease_expired(intent)
                    ):
                        active_operation_ids.add(operation_id)
                        deferred.append(operation_id)
                        continue
                    if intent.get("state") in {"commit_ready", "fencing"}:
                        artifact, manifest, intent = await self._fence_expired_commit(
                            intent_path,
                            intent,
                            token,
                            operation_id,
                        )
                        if manifest is None or intent is None:
                            deferred.append(operation_id)
                            continue
                    elif intent.get("state") == "reconciling":
                        pass
                    else:
                        claim_id = uuid.uuid4().hex
                        claimed = {**intent, "state": "reconciling", "reconciliation_claim": claim_id}
                        try:
                            await self._backend.replace_json(intent_path, claimed, token)
                        except _StorageConflict:
                            deferred.append(operation_id)
                            continue
                        intent = claimed
                        manifest, _ = await self._load()
                        artifact = self._find_operation_artifact(manifest, operation_id)
            if artifact is not None:
                removal = next(
                    (item for item in artifact.removals if item.operation_id == operation_id),
                    None,
                )
                committed_revision = next(
                    (item for item in artifact.revisions if item.operation_id == operation_id),
                    None,
                )
                result = ManagedWriteResult(
                    operation_id,
                    self.source_id,
                    artifact.artifact_id or "",
                    (
                        removal.generation
                        if removal is not None
                        else (
                            committed_revision.committed_generation
                            if committed_revision is not None
                            and committed_revision.committed_generation is not None
                            else manifest.generation
                        )
                    ),
                    removal.revision_id
                    if removal is not None
                    else (
                        committed_revision.revision_id if committed_revision is not None else artifact.revision_id
                    ),
                    removal.storage_path
                    if removal is not None
                    else (
                        committed_revision.storage_path if committed_revision is not None else artifact.storage_path
                    ),
                    deleted=removal is not None or artifact.deleted_at is not None,
                    cleanup_pending=removal.garbage_collect if removal is not None else False,
                )
                await self._write_receipt(result)
                result = await self._resume_remove_cleanup(result)
                recovered.append(operation_id)
                continue
            owned_path = intent.get("owned_path")
            owned_version_token = intent.get("owned_version_token")
            if isinstance(owned_path, str) and (
                intent.get("ownership_verified") is not True or not isinstance(owned_version_token, str)
            ):
                owned = await self._backend.exists(owned_path)
                if owned is not None and owned.operation_id == operation_id:
                    current, token = await self._backend.read_json(intent_path)
                    if (
                        current is not None
                        and current.get("state") in {"fencing", "reconciling"}
                        and current.get("reconciliation_claim") == intent.get("reconciliation_claim")
                    ):
                        recovered_intent = {
                            **current,
                            "ownership_verified": True,
                            "owned_version_token": owned.token,
                            "checksum_sha256": owned.checksum_sha256,
                        }
                        try:
                            await self._backend.replace_json(intent_path, recovered_intent, token)
                        except _StorageConflict:
                            deferred.append(operation_id)
                            continue
                        intent = recovered_intent
                        owned_version_token = owned.token
            if (
                isinstance(owned_path, str)
                and intent.get("ownership_verified") is True
                and isinstance(owned_version_token, str)
            ):
                outcome = await self._backend.delete_owned(owned_path, operation_id, owned_version_token)
                if outcome in {DeleteOutcome.DELETED, DeleteOutcome.ABSENT}:
                    if intent.get("state") != "fencing" or await self._finish_fenced_operation(
                        intent_path,
                        intent,
                        operation_id,
                    ):
                        if intent.get("state") != "reconciling" or await self._finish_cleanup_operation(
                            intent_path,
                            intent,
                        ):
                            removed.append(operation_id)
                        else:
                            deferred.append(operation_id)
                    else:
                        deferred.append(operation_id)
                else:
                    failures[operation_id] = "Owned orphan could not be verified or removed."
            else:
                if intent.get("state") == "fencing":
                    await self._finish_fenced_operation(intent_path, intent, operation_id)
                elif intent.get("state") == "reconciling":
                    await self._finish_cleanup_operation(intent_path, intent)
                deferred.append(operation_id)
        except Exception as exc:
            failures[operation_id] = type(exc).__name__
    if failures:
        raise ReconciliationError(
            "One or more interrupted operations could not be reconciled.",
            failures=failures,
            operation="reconcile",
        )
    async with self._backend.serialized():
        removed_staging, deferred_staging = await self._backend.reconcile_staging(
            frozenset(active_operation_ids),
            grace_seconds,
        )
    return ReconciliationReport(
        tuple(recovered),
        tuple(removed),
        tuple(deferred),
        removed_staging,
        deferred_staging,
        failures,
    )

ManagedWriteResult(operation_id, source_id, artifact_id, generation, revision_id, storage_path, deleted=False, cleanup_pending=False) dataclass

Committed write result and its read-after-write generation.

PromoteOutputRequest(operation_id, path, local_path, artifact_id=None, metadata=ArtifactMetadata(), expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions(), session_id='', output_name='') dataclass

Bases: UploadArtifactRequest

Explicitly promote a scratch/session output to a durable artifact.

ReconciliationReport(recovered=(), removed_orphans=(), deferred=(), removed_staging=(), deferred_staging=(), failures=dict()) dataclass

Outcome of scanning interrupted operation intents.

RegisterArtifactRequest(operation_id, path, storage_path, artifact_id=None, metadata=ArtifactMetadata(), content_revision=None, checksum_sha256=None, size_bytes=None, expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions()) dataclass

Snapshot caller-owned bytes into a durable immutable managed revision.

RemoveArtifactRequest(operation_id, reference, expected_generation=None, expected_revision_id=None, garbage_collect=True) dataclass

Commit a tombstone and then optionally collect owned revision bytes.

UploadArtifactRequest(operation_id, path, local_path, artifact_id=None, metadata=ArtifactMetadata(), expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions()) dataclass

Upload local bytes as a new immutable managed revision.

AuthorizedCatalogProvider(provider, authorizer, *, mode, per_artifact_enforcer=None)

Compose caller policy around an untrusted-to-authorize provider.

Providers still enforce ordinary query constraints such as source_ids; they never decide caller policy. Per-artifact mode additionally requires a backend-specific policy enforcer because generic post-filtering cannot preserve ranking, pagination, aggregation, or alias confidentiality.

Source code in src/agora_workbench/data_lake/policy.py
def __init__(
    self,
    provider: CatalogProvider,
    authorizer: CatalogAuthorizer,
    *,
    mode: CatalogPolicyMode,
    per_artifact_enforcer: CatalogPolicyEnforcer | None = None,
) -> None:
    if mode is not CatalogPolicyMode.HOMOGENEOUS_SOURCE and mode is not CatalogPolicyMode.PER_ARTIFACT:
        raise ValueError(f"Unknown catalog policy mode: {mode!r}.")
    if mode is CatalogPolicyMode.PER_ARTIFACT and per_artifact_enforcer is None:
        raise ValueError("Per-artifact policy requires a backend-specific CatalogPolicyEnforcer.")
    if mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE and per_artifact_enforcer is not None:
        raise ValueError("A per-artifact enforcer cannot be used with homogeneous-source policy.")
    self._provider = provider
    self._authorizer = authorizer
    self._mode = mode
    self._per_artifact_enforcer = per_artifact_enforcer

policy_mode property

Return the configured enforcement granularity.

authorizer property

Return the caller-scoped authorizer used by this policy wrapper.

capabilities(context) async

Return provider support intersected with current caller policy.

Source code in src/agora_workbench/data_lake/policy.py
async def capabilities(self, context: RequestContext) -> tuple[SourceCapabilities, ...]:
    """Return provider support intersected with current caller policy."""
    effective: list[SourceCapabilities] = []
    for capability in await self._provider_capabilities():
        allowed_operations: set[CatalogOperation] = set()
        for operation in capability.supported_operations:
            if await self._authorize_source(operation, capability.source_id, context):
                allowed_operations.add(operation)
        allowed = frozenset(allowed_operations)
        if allowed:
            effective.append(SourceCapabilities(capability.source_id, allowed))
    return tuple(effective)

search(request, context) async

Search with authorization enforced before effective pagination.

Source code in src/agora_workbench/data_lake/policy.py
async def search(self, request: SearchRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """Search with authorization enforced before effective pagination."""
    constrained, allowed_sources = await self._constrain_request(request, CatalogOperation.SEARCH, context)
    if not allowed_sources:
        return Page(())
    if self._uses_per_artifact_enforcement():
        assert self._per_artifact_enforcer is not None
        page = await self._per_artifact_enforcer.search(
            self._provider,
            constrained,
            context,
            self._authorizer,
        )
    elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
        page = await self._provider.search(constrained, context)
    else:
        self._raise_enforcement_failure(CatalogOperation.SEARCH)
    await self._validate_page(page, allowed_sources, CatalogOperation.SEARCH, context)
    return page

list(request, context) async

List with authorization enforced before effective pagination.

Source code in src/agora_workbench/data_lake/policy.py
async def list(self, request: ListRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """List with authorization enforced before effective pagination."""
    constrained, allowed_sources = await self._constrain_request(request, CatalogOperation.LIST, context)
    if not allowed_sources:
        return Page(())
    if self._uses_per_artifact_enforcement():
        assert self._per_artifact_enforcer is not None
        page = await self._per_artifact_enforcer.list(
            self._provider,
            constrained,
            context,
            self._authorizer,
        )
    elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
        page = await self._provider.list(constrained, context)
    else:
        self._raise_enforcement_failure(CatalogOperation.LIST)
    await self._validate_page(page, allowed_sources, CatalogOperation.LIST, context)
    return page

get(reference, context) async

Get an artifact without distinguishing denied from absent artifacts.

Source code in src/agora_workbench/data_lake/policy.py
async def get(self, reference: ArtifactReference, context: RequestContext) -> CatalogArtifact:
    """Get an artifact without distinguishing denied from absent artifacts."""
    await self._require_source_operation(reference, CatalogOperation.GET, context)
    backend_denied = False
    try:
        if self._uses_per_artifact_enforcement():
            assert self._per_artifact_enforcer is not None
            artifact = await self._per_artifact_enforcer.get(
                self._provider,
                reference,
                context,
                self._authorizer,
            )
        elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
            artifact = await self._provider.get(reference, context)
        else:
            self._raise_enforcement_failure(CatalogOperation.GET)
    except ArtifactNotFoundError:
        pass
    except PermissionDeniedError:
        backend_denied = True
    else:
        self._validate_reference(artifact.reference, reference, CatalogOperation.GET)
        if self._uses_per_artifact_enforcement():
            await self._require_artifact_operation(artifact.reference, CatalogOperation.GET, context)
        return artifact
    if backend_denied:
        self._raise_enforcement_failure(CatalogOperation.GET)
    self._raise_not_found(CatalogOperation.GET)

resolve(reference, context) async

Resolve an artifact without distinguishing denied from absent artifacts.

Source code in src/agora_workbench/data_lake/policy.py
async def resolve(self, reference: ArtifactReference, context: RequestContext) -> ResolvedArtifact:
    """Resolve an artifact without distinguishing denied from absent artifacts."""
    await self._require_source_operation(reference, CatalogOperation.RESOLVE, context)
    backend_denied = False
    try:
        if self._uses_per_artifact_enforcement():
            assert self._per_artifact_enforcer is not None
            resolved = await self._per_artifact_enforcer.resolve(
                self._provider,
                reference,
                context,
                self._authorizer,
            )
        elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
            resolved = await self._provider.resolve(reference, context)
        else:
            self._raise_enforcement_failure(CatalogOperation.RESOLVE)
    except ArtifactNotFoundError:
        pass
    except PermissionDeniedError:
        backend_denied = True
    else:
        self._validate_reference(resolved.reference, reference, CatalogOperation.RESOLVE)
        if self._uses_per_artifact_enforcement():
            await self._require_artifact_operation(resolved.reference, CatalogOperation.RESOLVE, context)
        return resolved
    if backend_denied:
        self._raise_enforcement_failure(CatalogOperation.RESOLVE)
    self._raise_not_found(CatalogOperation.RESOLVE)

DenyAllCatalogAuthorizer

Fail-closed policy suitable as an explicit production baseline.

DevelopmentAllowAllCatalogAuthorizer

Explicit development-only policy that permits every catalog request.

ArtifactResolver

Bases: Protocol

Resolve an opaque artifact ID to a fetchable storage location.

Implementations may define async def aclose(self) -> None for cleanup, but it is intentionally not required by the protocol.

unavailable_reason property

Return an operator-facing unavailability reason, or None when ready.

resolve(artifact_id) async

Resolve an artifact ID to a qualified name or URL.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(self, artifact_id: str) -> str:
    """Resolve an artifact ID to a qualified name or URL."""
    ...

CatalogAuthorizer

Bases: Protocol

Application-supplied caller policy, composed outside a catalog provider.

A request without reference is source-scoped. In per-artifact mode, source-scoped approval advertises an operation as potentially available; the policy enforcer must additionally authorize each artifact before it can affect ranking, pagination, aggregation, lookup results, or errors.

authorize(request, context) async

Return whether the caller may perform the requested operation.

Source code in src/agora_workbench/data_lake/protocols.py
async def authorize(self, request: CatalogAuthorizationRequest, context: RequestContext) -> bool:
    """Return whether the caller may perform the requested operation."""
    ...

CatalogPolicyEnforcer

Bases: Protocol

Backend-specific, trusted enforcement for per-artifact policy.

Implementations must apply policy before ranking, pagination, aggregation, alias resolution, and not-found decisions. Filtering one provider page after retrieval does not satisfy this contract.

search(provider, request, context, authorizer) async

Search only artifacts authorized for the current caller.

Source code in src/agora_workbench/data_lake/protocols.py
async def search(
    self,
    provider: CatalogProvider,
    request: SearchRequest,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> Page[CatalogArtifact]:
    """Search only artifacts authorized for the current caller."""
    ...

list(provider, request, context, authorizer) async

List only artifacts authorized for the current caller.

Source code in src/agora_workbench/data_lake/protocols.py
async def list(
    self,
    provider: CatalogProvider,
    request: ListRequest,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> Page[CatalogArtifact]:
    """List only artifacts authorized for the current caller."""
    ...

get(provider, reference, context, authorizer) async

Authorize lookup semantics before returning the exact requested logical reference.

Source code in src/agora_workbench/data_lake/protocols.py
async def get(
    self,
    provider: CatalogProvider,
    reference: ArtifactReference,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> CatalogArtifact:
    """Authorize lookup semantics before returning the exact requested logical reference."""
    ...

resolve(provider, reference, context, authorizer) async

Resolve only an authorized canonical artifact.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(
    self,
    provider: CatalogProvider,
    reference: ArtifactReference,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> ResolvedArtifact:
    """Resolve only an authorized canonical artifact."""
    ...

CatalogProvider

Bases: Protocol

Backend-neutral asynchronous catalog interface.

Runtime checks only verify that named members exist. They do not validate signatures, async behavior, return types, or capability truthfulness. capabilities values are authoritative for provider support; callers must not infer support from method presence. Authorization policy is composed outside the provider protocol.

capabilities() async

Return authoritative read support for each source exposed by the provider.

Source code in src/agora_workbench/data_lake/protocols.py
async def capabilities(self) -> tuple[SourceCapabilities, ...]:
    """Return authoritative read support for each source exposed by the provider."""
    ...

search(request, context) async

Search artifacts using provider-defined ranking and opaque cursors.

Source code in src/agora_workbench/data_lake/protocols.py
async def search(self, request: SearchRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """Search artifacts using provider-defined ranking and opaque cursors."""
    ...

list(request, context) async

List artifacts using provider-defined filters and opaque cursors.

Source code in src/agora_workbench/data_lake/protocols.py
async def list(self, request: ListRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """List artifacts using provider-defined filters and opaque cursors."""
    ...

get(reference, context) async

Get one artifact, honoring or explicitly rejecting a pinned revision.

Source code in src/agora_workbench/data_lake/protocols.py
async def get(self, reference: ArtifactReference, context: RequestContext) -> CatalogArtifact:
    """Get one artifact, honoring or explicitly rejecting a pinned revision."""
    ...

resolve(reference, context) async

Resolve a logical reference, honoring or explicitly rejecting a pinned revision.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(self, reference: ArtifactReference, context: RequestContext) -> ResolvedArtifact:
    """Resolve a logical reference, honoring or explicitly rejecting a pinned revision."""
    ...

DetailedStreamingArtifactFetcher

Bases: Protocol

Optional detailed-result seam for managed-read and audit integrations.

fetch_to_file_result(qualified_name, dest_path, *, options=None, context=None) async

Stream an artifact and return integrity and caller diagnostics.

Source code in src/agora_workbench/data_lake/protocols.py
async def fetch_to_file_result(
    self,
    qualified_name: str,
    dest_path: Path,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> TransferResult:
    """Stream an artifact and return integrity and caller diagnostics."""
    raise NotImplementedError

DetailedStreamingArtifactPublisher

Bases: Protocol

Optional detailed-result seam for managed-write integrations.

publish_with_result(local_path, name, session_id, *, options=None, context=None) async

Publish a file and return its locator plus transfer result.

Source code in src/agora_workbench/data_lake/protocols.py
async def publish_with_result(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> tuple[str, TransferResult]:
    """Publish a file and return its locator plus transfer result."""
    raise NotImplementedError

PolicyEnforcedCatalog

Bases: Protocol

Caller-aware catalog surface produced by policy composition.

policy_mode property

Return the configured enforcement granularity.

capabilities(context) async

Return provider support intersected with current caller policy.

Source code in src/agora_workbench/data_lake/protocols.py
async def capabilities(self, context: RequestContext) -> tuple[SourceCapabilities, ...]:
    """Return provider support intersected with current caller policy."""
    ...

search(request, context) async

Search within the current caller's effective authorization scope.

Source code in src/agora_workbench/data_lake/protocols.py
async def search(self, request: SearchRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """Search within the current caller's effective authorization scope."""
    ...

list(request, context) async

List within the current caller's effective authorization scope.

Source code in src/agora_workbench/data_lake/protocols.py
async def list(self, request: ListRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """List within the current caller's effective authorization scope."""
    ...

get(reference, context) async

Get an artifact without disclosing unauthorized existence.

Source code in src/agora_workbench/data_lake/protocols.py
async def get(self, reference: ArtifactReference, context: RequestContext) -> CatalogArtifact:
    """Get an artifact without disclosing unauthorized existence."""
    ...

resolve(reference, context) async

Resolve an artifact without disclosing unauthorized existence.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(self, reference: ArtifactReference, context: RequestContext) -> ResolvedArtifact:
    """Resolve an artifact without disclosing unauthorized existence."""
    ...

StreamingArtifactFetcher

Bases: Protocol

Fetcher capability that commits a bounded stream to a local file.

This is deliberately separate from full-memory fetch() conveniences. Implementations leave an existing destination unchanged on failure.

fetch_to_file(qualified_name, dest_path, *, options=None, context=None) async

Stream an artifact to a file and return the committed byte count.

Source code in src/agora_workbench/data_lake/protocols.py
async def fetch_to_file(
    self,
    qualified_name: str,
    dest_path: Path,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> int:
    """Stream an artifact to a file and return the committed byte count."""
    raise NotImplementedError

StreamingArtifactPublisher

Bases: Protocol

Publisher capability that uploads a local file with bounded memory.

publish(local_path, name, session_id, *, options=None, context=None) async

Publish a file after enforcing the supplied transfer guarantees.

Source code in src/agora_workbench/data_lake/protocols.py
async def publish(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> str:
    """Publish a file after enforcing the supplied transfer guarantees."""
    raise NotImplementedError

TransferDiagnostic(operation, state, context, resource=None, bytes_transferred=0, checksum_sha256=None, error_type=None) dataclass

Credential-safe event delivered to an operator-provided audit hook.

TransferOptions(max_bytes=DEFAULT_TRANSFER_MAX_BYTES, quota_bytes=None, timeout_seconds=DEFAULT_TRANSFER_TIMEOUT_SECONDS, chunk_size=DEFAULT_TRANSFER_CHUNK_BYTES, expected_sha256=None, cancellation_event=None, diagnostic_hook=None, create_exclusive=False, object_metadata=dict(), allow_reserved=False) dataclass

Limits and integrity requirements for one streaming transfer.

max_bytes bounds the individual object while quota_bytes represents caller/provider capacity remaining before the operation starts. The smaller non-None value is enforced. fetch() convenience methods still load the complete object in memory; these options apply to streaming methods.

effective_max_bytes property

Return the tightest configured object/quota bound.

TransferResult(bytes_transferred, checksum_sha256, context, resource=None, elapsed_seconds=0.0, created=None, object_metadata=dict()) dataclass

Completed bounded transfer details, including the original caller context.

azure_uri_from_blob_name(account, container, blob_name)

Build a canonical URI from an SDK-decoded blob name, quoting exactly once.

Source code in src/agora_workbench/data_lake/identity.py
def azure_uri_from_blob_name(account: str, container: str, blob_name: str) -> str:
    """Build a canonical URI from an SDK-decoded blob name, quoting exactly once."""
    account, container, _ = parse_azure_uri(f"az://{account}/{container}")
    encoded_path = quote(blob_name, safe="/-._~")
    return f"az://{account}/{container}/{encoded_path}" if encoded_path else f"az://{account}/{container}"

canonicalize_azure_uri(uri)

Canonicalize a supported Azure Blob or DFS URI without credentials.

Source code in src/agora_workbench/data_lake/identity.py
def canonicalize_azure_uri(uri: str) -> str:
    """Canonicalize a supported Azure Blob or DFS URI without credentials."""
    account, container, object_path = parse_azure_uri(uri)
    return azure_uri_from_blob_name(account, container, object_path)

is_reserved_provider_path(path)

Return whether a source-relative path belongs to provider-managed state.

Source code in src/agora_workbench/data_lake/identity.py
def is_reserved_provider_path(path: str) -> bool:
    """Return whether a source-relative path belongs to provider-managed state."""
    normalized = path.replace("\\", "/").lstrip("/")
    return normalized == RESERVED_PROVIDER_PREFIX.rstrip("/") or normalized.startswith(RESERVED_PROVIDER_PREFIX)

is_scan_excluded_path(path)

Return whether scan discovery must prune a hidden or provider-managed path.

Source code in src/agora_workbench/data_lake/identity.py
def is_scan_excluded_path(path: str) -> bool:
    """Return whether scan discovery must prune a hidden or provider-managed path."""
    normalized = path.replace("\\", "/").strip("/")
    return is_reserved_provider_path(normalized) or any(segment.startswith(".") for segment in normalized.split("/"))

logical_artifact_id(source_id, logical_path)

Generate a location-independent ID for a newly discovered logical path.

Source code in src/agora_workbench/data_lake/identity.py
def logical_artifact_id(source_id: str, logical_path: str) -> str:
    """Generate a location-independent ID for a newly discovered logical path."""
    normalized = normalize_logical_path(logical_path)
    return uuid.uuid5(_LOGICAL_ID_NAMESPACE, f"{source_id}\0{normalized}").hex

normalize_logical_path(path)

Return a portable, source-relative POSIX path.

Source code in src/agora_workbench/data_lake/identity.py
def normalize_logical_path(path: str) -> str:
    """Return a portable, source-relative POSIX path."""
    candidate = path.replace("\\", "/")
    if candidate.startswith("/") or re.match(r"^[a-zA-Z]:[\\/]", path):
        raise _invalid_identity("Artifact path must be source-relative.")
    normalized = posixpath.normpath(candidate)
    if normalized in {"", "."}:
        raise _invalid_identity("Artifact path must identify an object.")
    if normalized == ".." or normalized.startswith("../"):
        raise _invalid_identity("Artifact path must stay within its source.")
    return str(PurePosixPath(normalized))

parse_azure_uri(uri)

Parse a supported Azure URI into account, container, and decoded object path.

Source code in src/agora_workbench/data_lake/identity.py
def parse_azure_uri(uri: str) -> tuple[str, str, str]:
    """Parse a supported Azure URI into account, container, and decoded object path."""
    if _MALFORMED_PERCENT_RE.search(uri):
        raise _invalid_identity("Azure storage URI contains malformed percent encoding.")
    parsed = urlsplit(uri)
    scheme = parsed.scheme.lower()
    if scheme != "abfss" and (parsed.username is not None or parsed.password is not None):
        raise _invalid_identity("Azure storage URI must not contain user information.")
    if scheme != "abfss":
        try:
            if parsed.port is not None:
                raise _invalid_identity("Azure storage URI ports are not supported.")
        except ValueError as exc:
            raise _invalid_identity("Azure storage URI contains an invalid port.") from exc

    if scheme == "az":
        account = parsed.hostname or ""
        parts = parsed.path.lstrip("/").split("/", 1)
        container = parts[0] if parts else ""
        encoded_path = parts[1] if len(parts) > 1 else ""
    elif scheme in {"http", "https"}:
        host = (parsed.hostname or "").lower()
        suffix = next(
            (
                candidate
                for candidate in (".blob.core.windows.net", ".dfs.core.windows.net")
                if host.endswith(candidate)
            ),
            None,
        )
        if suffix is None:
            raise _invalid_identity("Unsupported Azure storage URI host.")
        account = host[: -len(suffix)]
        parts = parsed.path.lstrip("/").split("/", 1)
        container = parts[0] if parts else ""
        encoded_path = parts[1] if len(parts) > 1 else ""
    elif scheme == "abfss":
        if parsed.netloc.count("@") != 1:
            raise _invalid_identity("Malformed abfss URI.")
        encoded_container, host = parsed.netloc.split("@", 1)
        if ":" in host:
            raise _invalid_identity("Azure storage URI ports are not supported.")
        suffix = ".dfs.core.windows.net"
        if not host.lower().endswith(suffix):
            raise _invalid_identity("Unsupported Azure storage URI host.")
        account = host[: -len(suffix)]
        container = encoded_container
        encoded_path = parsed.path.lstrip("/")
    else:
        raise _invalid_identity("Unsupported Azure storage URI scheme.")

    if _ENCODED_SEPARATOR_RE.search(encoded_path):
        raise _invalid_identity("Azure object paths must not contain encoded separators.")
    account = account.lower()
    container = unquote(container).lower()
    if not _ACCOUNT_RE.fullmatch(account):
        raise _invalid_identity("Azure storage account name is malformed.")
    if container not in _SYSTEM_CONTAINERS and (not _CONTAINER_RE.fullmatch(container) or "--" in container):
        raise _invalid_identity("Azure storage container name is malformed.")
    object_path = unquote(encoded_path)
    validate_azure_object_path(object_path, allow_empty=True, allow_reserved=True)
    return account, container, object_path

sanitize_uri_for_display(uri)

Remove credentials, query parameters, and fragments from a URI.

Source code in src/agora_workbench/data_lake/identity.py
def sanitize_uri_for_display(uri: str) -> str:
    """Remove credentials, query parameters, and fragments from a URI."""
    parsed = urlsplit(uri)
    if parsed.scheme.lower() == "abfss":
        container = parsed.username or ""
        hostname = parsed.hostname or ""
        netloc = f"{container}@{hostname}" if container else hostname
        return urlunsplit((parsed.scheme, netloc, parsed.path, "", ""))
    hostname = parsed.hostname or ""
    try:
        port = parsed.port
    except ValueError:
        port = None
    if port is not None:
        hostname = f"{hostname}:{port}"
    return urlunsplit((parsed.scheme, hostname, parsed.path, "", ""))

split_alias(value, default_namespace='artifact-id')

Split a namespaced alias while retaining compatibility with opaque IDs.

Source code in src/agora_workbench/data_lake/identity.py
def split_alias(value: str, default_namespace: str = "artifact-id") -> tuple[str, str]:
    """Split a namespaced alias while retaining compatibility with opaque IDs."""
    if ":" not in value:
        return default_namespace, value
    namespace, alias = value.split(":", 1)
    if not namespace or not alias:
        raise _invalid_identity("Artifact aliases require non-empty namespace and value.")
    return namespace, alias

stable_source_id(source_type, root)

Derive a stable fallback source ID.

Source code in src/agora_workbench/data_lake/identity.py
def stable_source_id(source_type: str, root: str) -> str:
    """Derive a stable fallback source ID."""
    identity_root = canonicalize_azure_uri(root) if source_type == "blob" else str(root)
    digest = hashlib.sha256(f"{source_type}\0{identity_root}".encode()).hexdigest()[:20]
    return f"{source_type}-{digest}"

validate_azure_object_path(path, *, allow_empty=False, allow_reserved=False)

Validate one SDK-decoded Blob object path without changing its spelling.

Source code in src/agora_workbench/data_lake/identity.py
def validate_azure_object_path(path: str, *, allow_empty: bool = False, allow_reserved: bool = False) -> str:
    """Validate one SDK-decoded Blob object path without changing its spelling."""
    if not path:
        if allow_empty:
            return ""
        raise _invalid_identity("Azure storage URI must identify an object.")
    if "\\" in path or "\x00" in path:
        raise _invalid_identity("Azure object path contains an ambiguous separator or NUL.")
    segments = path.split("/")
    path_segments = segments[:-1] if segments[-1] == "" else segments
    if any(segment in {"", ".", ".."} for segment in path_segments):
        raise _invalid_identity("Azure object path contains empty or dot segments.")
    if not allow_reserved and is_reserved_provider_path(path):
        raise _invalid_identity("Azure object path is reserved for provider metadata.")
    return path

validate_managed_revision_path(path)

Validate the narrow reserved storage namespace available to managed writes.

Source code in src/agora_workbench/data_lake/identity.py
def validate_managed_revision_path(path: str) -> str:
    """Validate the narrow reserved storage namespace available to managed writes."""
    normalized = validate_azure_object_path(path, allow_reserved=True)
    if not normalized.startswith(RESERVED_REVISIONS_PREFIX):
        raise _invalid_identity("Managed storage path must be within the provider revisions prefix.")
    return normalized

managed_writer_extension_factory(writer)

Adapt a managed writer to CatalogIntegration.capability_extension_factory.

Source code in src/agora_workbench/data_lake/managed.py
def managed_writer_extension_factory(writer: ManagedCatalogWriter):
    """Adapt a managed writer to ``CatalogIntegration.capability_extension_factory``."""

    def create(_session: object, catalog: object, _context: RequestContext) -> AuthorizedManagedCatalogWriter:
        authorizer = getattr(catalog, "authorizer", None)
        if authorizer is None:
            raise TypeError("Managed writer integration requires an authorized catalog.")
        return AuthorizedManagedCatalogWriter(writer, authorizer)

    return create

Catalog

agora_workbench.data_lake.catalog

Public catalog contracts and compatibility exports.

CatalogConfig

Bases: BaseModel

Top-level catalog.yaml configuration.

from_yaml(path) classmethod

Load configuration from a YAML file.

Source code in src/agora_workbench/code_execution/data_access/catalog/config.py
@classmethod
def from_yaml(cls, path: str | Path) -> "CatalogConfig":
    """Load configuration from a YAML file."""
    config_path = Path(path)
    if not config_path.exists():
        raise FileNotFoundError(f"Catalog config not found: {config_path}")
    raw = yaml.safe_load(config_path.read_text(encoding="utf-8")) or {}
    return cls.model_validate(raw)

convert(raw) classmethod

Convert and report the existing public catalog format without storage access.

Source code in src/agora_workbench/code_execution/data_access/catalog/config.py
@classmethod
def convert(cls, raw: dict[str, Any]) -> tuple["CatalogConfig", dict[str, object]]:
    """Convert and report the existing public catalog format without storage access."""
    config = cls.model_validate(raw)
    sources = []
    for source in config.sources:
        sources.append(
            {
                "source_id": source.source_id
                or (
                    "derived at load time from the canonical source root"
                    if source.discovery is DiscoveryMode.SCAN
                    else None
                ),
                "source_type": source.source_type,
                "discovery": source.discovery.value,
                "path": source.path,
                "manifest": source.manifest,
                "max_stale_seconds": source.max_stale_seconds,
                "files_are_metadata_overrides": source.files is not None,
                "configuration_valid": True,
                "manifest_checked": False,
                "manifest_content_valid": None,
            }
        )
    return config, {
        "version": config.version,
        "configuration_valid": True,
        "manifest_checked": False,
        "manifest_content_valid": None,
        "storage_accessed": False,
        "sources": sources,
    }

CatalogConfigConversionReport(source, destination, changed, written, rendered_yaml, summary) dataclass

Result of converting legacy catalog YAML to explicit versioned source configuration.

CatalogDB(db_path=':memory:', vec_dimensions=768)

SQLite catalog with durable identity, history, FTS5, and optional vectors.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def __init__(self, db_path: str | Path = ":memory:", vec_dimensions: int | None = 768):
    if vec_dimensions is not None and vec_dimensions <= 0:
        raise ValueError("vec_dimensions must be greater than zero")
    self._db_path = str(db_path)
    self._vec_dimensions = vec_dimensions
    self._conn: Optional[sqlite3.Connection] = None
    self._vector_loaded = False
    self._write_lock = threading.RLock()

vec_dimensions property

Configured or inferred vector dimensions for this catalog.

open()

Open the database and migrate known older schemas atomically.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def open(self) -> None:
    """Open the database and migrate known older schemas atomically."""
    connection = sqlite3.connect(self._db_path, timeout=5.0, check_same_thread=False)
    try:
        connection.row_factory = sqlite3.Row
        connection.execute("PRAGMA busy_timeout = 5000")
        connection.execute("PRAGMA foreign_keys = ON")
        version = connection.execute("PRAGMA user_version").fetchone()[0]
        if version > SCHEMA_VERSION:
            raise RuntimeError(
                f"Catalog schema version {version} is newer than supported version {SCHEMA_VERSION}; "
                "the database was not modified."
            )
        if self._db_path != ":memory:":
            try:
                journal_mode = connection.execute("PRAGMA journal_mode = WAL").fetchone()[0]
                if journal_mode.lower() != "wal":
                    LOGGER.warning("WAL mode is unavailable for the catalog database")
            except sqlite3.DatabaseError:
                LOGGER.warning("WAL mode is unavailable for the catalog database")

        self._conn = connection
        self._probe_fts5()
        if version < 2 and self._has_legacy_schema():
            self._migrate_legacy_schema()
        else:
            connection.executescript(f"BEGIN IMMEDIATE;\n{_SCHEMA_SQL}")
            try:
                refresh_columns = {
                    row["name"] for row in connection.execute("PRAGMA table_info(catalog_source_refreshes)")
                }
                if "manifest_generation" not in refresh_columns:
                    connection.execute(
                        "ALTER TABLE catalog_source_refreshes ADD COLUMN manifest_generation INTEGER"
                    )
                if "manifest_etag" not in refresh_columns:
                    connection.execute("ALTER TABLE catalog_source_refreshes ADD COLUMN manifest_etag TEXT")
                connection.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
                connection.commit()
            except Exception:
                connection.rollback()
                raise
    except BaseException:
        connection.close()
        self._conn = None
        raise

execute_readonly(sql, max_rows=100)

Execute a SELECT from one committed snapshot.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def execute_readonly(self, sql: str, max_rows: int = 100) -> list[dict]:
    """Execute a SELECT from one committed snapshot."""
    stripped = sql.strip().upper()
    write_keywords = ("INSERT", "UPDATE", "DELETE", "DROP", "ALTER", "CREATE", "REPLACE")
    if any(stripped.startswith(keyword) for keyword in write_keywords):
        raise ValueError(f"Write operations are not permitted. Query starts with: {stripped.split()[0]}")
    with self._read_snapshot(vectors=_references_vector_table(sql)) as connection:
        previous_query_only = connection.execute("PRAGMA query_only").fetchone()[0]
        connection.execute("PRAGMA query_only = ON")
        try:
            cursor = connection.execute(sql)
            return [dict(row) for row in cursor.fetchmany(max_rows)]
        finally:
            connection.execute(f"PRAGMA query_only = {int(previous_query_only)}")

upsert_artifact(artifact_id, name, storage_uri, description=None, domain=None, source_type=None, content_type=None, size_bytes=None, indexed_at=None, embedding=None, *, source_id='legacy', logical_path=None, source_root=None, content_revision=None, metadata_revision=None, checksum_sha256=None, aliases=(), _replace_embedding=False, _allow_move=False, _commit=True)

Insert or update an artifact and return its persisted logical ID.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def upsert_artifact(
    self,
    artifact_id: str | None,
    name: str,
    storage_uri: str,
    description: Optional[str] = None,
    domain: Optional[str] = None,
    source_type: Optional[str] = None,
    content_type: Optional[str] = None,
    size_bytes: Optional[int] = None,
    indexed_at: Optional[str] = None,
    embedding: Optional[list[float]] = None,
    *,
    source_id: str = "legacy",
    logical_path: str | None = None,
    source_root: str | None = None,
    content_revision: str | None = None,
    metadata_revision: str | None = None,
    checksum_sha256: str | None = None,
    aliases: tuple[str, ...] | list[str] = (),
    _replace_embedding: bool = False,
    _allow_move: bool = False,
    _commit: bool = True,
) -> str:
    """Insert or update an artifact and return its persisted logical ID."""
    now = indexed_at or datetime.now(timezone.utc).isoformat()
    source_type = source_type or "legacy"
    if not source_id:
        raise ValueError("Source ID must be non-empty")
    if artifact_id == "":
        raise ValueError("Artifact ID must be non-empty")
    if _commit:
        with self._write_transaction():
            return self.upsert_artifact(
                artifact_id,
                name,
                storage_uri,
                description,
                domain,
                source_type,
                content_type,
                size_bytes,
                indexed_at,
                embedding,
                source_id=source_id,
                logical_path=logical_path,
                source_root=source_root,
                content_revision=content_revision,
                metadata_revision=metadata_revision,
                checksum_sha256=checksum_sha256,
                aliases=aliases,
                _replace_embedding=_replace_embedding,
                _allow_move=_allow_move,
                _commit=False,
            )
    if embedding is not None:
        self._validate_vector_dimensions(embedding, "artifact embedding")
        self._ensure_vector_capability(self.conn)
    existing_by_id = (
        self.conn.execute("SELECT * FROM artifacts WHERE id=?", (artifact_id,)).fetchone() if artifact_id else None
    )
    requested_logical_path = logical_path
    if existing_by_id is not None and requested_logical_path is None:
        source_id = existing_by_id["source_id"]
        requested_logical_path = existing_by_id["logical_path"]
    logical_path = normalize_logical_path(requested_logical_path or name)
    content_revision = content_revision or (
        f"sha256:{checksum_sha256}"
        if checksum_sha256 is not None
        else (f"size:{size_bytes}" if size_bytes is not None else "unknown")
    )
    metadata_revision = metadata_revision or _digest([name, description, domain, source_type, content_type])

    with nullcontext():
        self._ensure_source(source_id, source_type, source_root, now)
        existing_by_path = self.conn.execute(
            "SELECT * FROM artifacts WHERE source_id = ? AND logical_path = ?",
            (source_id, logical_path),
        ).fetchone()
        configured_alias = None
        if existing_by_id is not None:
            chosen_id = existing_by_id["id"]
        elif existing_by_path is not None:
            chosen_id = existing_by_path["id"]
            if artifact_id is not None and artifact_id != chosen_id:
                configured_alias = artifact_id
        else:
            chosen_id = artifact_id or logical_artifact_id(source_id, logical_path)
        adopting_legacy = (
            existing_by_id is not None
            and existing_by_id["id"] == chosen_id
            and existing_by_id["source_id"].startswith("legacy-")
            and existing_by_id["logical_path"].startswith("imported/")
        )
        if existing_by_id is not None and existing_by_id["source_id"] != source_id and not adopting_legacy:
            raise ValueError(
                f"Artifact ID {chosen_id!r} is already assigned to source {existing_by_id['source_id']!r}"
            )
        if (
            existing_by_id is not None
            and existing_by_id["source_id"] == source_id
            and existing_by_id["logical_path"] != logical_path
            and not adopting_legacy
            and not _allow_move
        ):
            raise ValueError(
                f"Artifact ID {chosen_id!r} is already assigned to {source_id}:{existing_by_id['logical_path']}"
            )
        existing = existing_by_id or existing_by_path

        conflicting_alias = self.conn.execute(
            "SELECT artifact_id FROM artifact_aliases WHERE namespace='artifact-id' AND alias=?",
            (chosen_id,),
        ).fetchone()
        if conflicting_alias is not None and conflicting_alias["artifact_id"] != chosen_id:
            raise ValueError(f"Artifact ID collides with an existing alias: {chosen_id!r}")
        if ":" in chosen_id:
            namespace, alias = split_alias(chosen_id)
            qualified_alias = self.conn.execute(
                "SELECT artifact_id FROM artifact_aliases WHERE namespace=? AND alias=?",
                (namespace, alias),
            ).fetchone()
            if qualified_alias is not None and qualified_alias["artifact_id"] != chosen_id:
                raise ValueError(f"Artifact ID collides with an existing alias: {chosen_id!r}")
        if existing_by_path is not None and existing_by_path["id"] != chosen_id:
            raise ValueError(
                f"Logical path {source_id}:{logical_path} is already assigned to artifact {existing_by_path['id']}"
            )

        if existing is None:
            id_owner = self.conn.execute("SELECT source_id FROM artifacts WHERE id = ?", (chosen_id,)).fetchone()
            if id_owner is not None:
                raise ValueError(f"Artifact ID collision: {chosen_id!r}")
            self._insert_artifact(
                artifact_id=chosen_id,
                source_id=source_id,
                logical_path=logical_path,
                name=name,
                storage_uri=storage_uri,
                description=description,
                domain=domain,
                source_type=source_type,
                content_type=content_type,
                size_bytes=size_bytes,
                indexed_at=now,
                content_revision=content_revision,
                metadata_revision=metadata_revision,
                checksum_sha256=checksum_sha256,
            )
        else:
            changed = any(
                (
                    existing["storage_uri"] != storage_uri,
                    existing["source_id"] != source_id,
                    existing["logical_path"] != logical_path,
                    existing["name"] != name,
                    existing["description"] != description,
                    existing["domain"] != domain,
                    existing["source_type"] != source_type,
                    existing["content_type"] != content_type,
                    existing["size_bytes"] != size_bytes,
                    existing["content_revision"] != content_revision,
                    existing["metadata_revision"] != metadata_revision,
                    existing["checksum_sha256"] != checksum_sha256,
                    existing["deleted_at"] is not None,
                )
            )
            if changed:
                revision = existing["current_revision"] + 1
                self.conn.execute(
                    """UPDATE artifacts SET
                           source_id=?, logical_path=?, name=?, storage_uri=?, description=?,
                           domain=?, source_type=?, content_type=?, size_bytes=?, indexed_at=?,
                           current_revision=?, content_revision=?, metadata_revision=?,
                           checksum_sha256=?, deleted_at=NULL
                       WHERE id=?""",
                    (
                        source_id,
                        logical_path,
                        name,
                        storage_uri,
                        description,
                        domain,
                        source_type,
                        content_type,
                        size_bytes,
                        now,
                        revision,
                        content_revision,
                        metadata_revision,
                        checksum_sha256,
                        chosen_id,
                    ),
                )
                if adopting_legacy:
                    self.conn.execute(
                        "UPDATE artifact_aliases SET source_id=? WHERE artifact_id=?",
                        (source_id, chosen_id),
                    )
                self._insert_revision(chosen_id)

        for alias in aliases:
            namespace, value = split_alias(alias)
            self.add_alias(namespace, value, source_id, chosen_id, commit=False)
        if configured_alias is not None:
            namespace, value = split_alias(configured_alias)
            self.add_alias(namespace, value, source_id, chosen_id, commit=False)
        if _replace_embedding or embedding is not None:
            self._ensure_vector_capability(self.conn)
            self.conn.execute("DELETE FROM artifacts_vec WHERE id = ?", (chosen_id,))
        if embedding is not None:
            vector_state = self.conn.execute(
                "SELECT model_id FROM catalog_vector_state WHERE singleton=1"
            ).fetchone()
            if vector_state["model_id"] is None:
                vector_count = self.conn.execute("SELECT COUNT(*) FROM artifacts_vec").fetchone()[0]
                if vector_count:
                    raise ValueError(
                        "Catalog contains vectors with unknown model identity; rebuild vectors before writing"
                    )
                self.conn.execute("UPDATE catalog_vector_state SET model_id='manual' WHERE singleton=1")
            self.conn.execute(
                "INSERT INTO artifacts_vec (id, embedding) VALUES (?, ?)",
                (chosen_id, _serialize_vector(embedding)),
            )
    return chosen_id

upsert_artifacts_batch(artifacts)

Atomically upsert a batch of artifacts.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def upsert_artifacts_batch(self, artifacts: list[dict]) -> None:
    """Atomically upsert a batch of artifacts."""
    with self._write_transaction():
        for artifact in artifacts:
            self.upsert_artifact(**artifact, _commit=False)

apply_refresh_batch(artifacts, stale_artifact_ids, source_results, vector_model_id=None)

Apply one refresh generation atomically across rows, search indexes, and vectors.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def apply_refresh_batch(
    self,
    artifacts: list[dict],
    stale_artifact_ids: list[str],
    source_results: list[dict],
    vector_model_id: str | None = None,
) -> None:
    """Apply one refresh generation atomically across rows, search indexes, and vectors."""
    with self._write_transaction():
        if vector_model_id is not None:
            self._ensure_vector_capability(self.conn)
            dimensions = self._vec_dimensions
            if dimensions is None:
                raise ValueError("Catalog vector dimensions are unknown after vector capability initialization")
            self._validate_vector_state(self.conn, vector_model_id, dimensions)
            self.conn.execute(
                "UPDATE catalog_vector_state SET model_id=? WHERE singleton=1",
                (vector_model_id,),
            )
        for result in source_results:
            self._ensure_source(
                result["source_id"],
                result["source_type"],
                result.get("root_uri"),
                result["attempted_at"],
            )
        refresh_timestamp = source_results[0]["attempted_at"] if source_results else None
        self._delete_artifacts(stale_artifact_ids, deleted_at=refresh_timestamp)
        for artifact in artifacts:
            self.upsert_artifact(**artifact, _commit=False)
        for result in source_results:
            self._record_source_refresh(**result)

get_source_refresh_state(source_id)

Return the latest observable refresh state for a source.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def get_source_refresh_state(self, source_id: str) -> SourceRefreshState | None:
    """Return the latest observable refresh state for a source."""
    with self._read_snapshot() as connection:
        row = connection.execute(
            "SELECT * FROM catalog_source_refreshes WHERE source_id=?", (source_id,)
        ).fetchone()
    return SourceRefreshState(**dict(row)) if row is not None else None

list_source_refresh_states()

Return latest refresh states in deterministic source order.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def list_source_refresh_states(self) -> list[SourceRefreshState]:
    """Return latest refresh states in deterministic source order."""
    with self._read_snapshot() as connection:
        rows = connection.execute("SELECT * FROM catalog_source_refreshes ORDER BY source_id").fetchall()
    return [SourceRefreshState(**dict(row)) for row in rows]

validate_vector_state(model_id, dimensions)

Reject incompatible vector model/dimension reuse instead of mixing embeddings.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def validate_vector_state(self, model_id: str, dimensions: int) -> None:
    """Reject incompatible vector model/dimension reuse instead of mixing embeddings."""
    with self._write_lock:
        if self._vec_dimensions is None:
            self._vec_dimensions = dimensions
    with self._read_snapshot(vectors=True) as connection:
        self._validate_vector_state(connection, model_id, dimensions)

has_vector(artifact_id)

Return whether the current artifact has an indexed embedding.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def has_vector(self, artifact_id: str) -> bool:
    """Return whether the current artifact has an indexed embedding."""
    return artifact_id not in self.missing_vectors([artifact_id])

missing_vectors(artifact_ids)

Return artifact IDs without vectors using one operation-scoped snapshot.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def missing_vectors(self, artifact_ids: list[str]) -> set[str]:
    """Return artifact IDs without vectors using one operation-scoped snapshot."""
    unique_ids = list(dict.fromkeys(artifact_ids))
    if not unique_ids:
        return set()
    present_ids: set[str] = set()
    with self._read_snapshot(vectors=True) as connection:
        for offset in range(0, len(unique_ids), 500):
            batch = unique_ids[offset : offset + 500]
            placeholders = ",".join("?" for _ in batch)
            rows = connection.execute(
                f"SELECT id FROM artifacts_vec WHERE id IN ({placeholders})",
                batch,
            ).fetchall()
            present_ids.update(row["id"] for row in rows)
    return set(unique_ids) - present_ids

add_alias(namespace, alias, source_id, artifact_id, *, commit=True)

Add an unambiguous namespaced alias.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def add_alias(
    self,
    namespace: str,
    alias: str,
    source_id: str,
    artifact_id: str,
    *,
    commit: bool = True,
) -> None:
    """Add an unambiguous namespaced alias."""
    if commit:
        with self._write_transaction():
            self.add_alias(namespace, alias, source_id, artifact_id, commit=False)
        return
    if not namespace or not alias:
        raise ValueError("Alias namespace and value must be non-empty")
    canonical_candidates = (
        (alias, f"{namespace}:{alias}") if namespace == "artifact-id" else (f"{namespace}:{alias}",)
    )
    placeholders = ",".join("?" for _ in canonical_candidates)
    direct = self.conn.execute(
        f"SELECT id FROM artifacts WHERE id IN ({placeholders}) AND id != ?",
        (*canonical_candidates, artifact_id),
    ).fetchone()
    if direct is not None:
        raise ValueError(f"Alias collides with canonical artifact ID: {direct['id']!r}")
    existing = self.conn.execute(
        "SELECT artifact_id FROM artifact_aliases WHERE source_id=? AND namespace=? AND alias=?",
        (source_id, namespace, alias),
    ).fetchone()
    if existing is not None and existing["artifact_id"] != artifact_id:
        raise ValueError(f"Alias collision in source {source_id!r}: {namespace}:{alias}")
    self.conn.execute(
        """INSERT OR IGNORE INTO artifact_aliases(namespace, alias, source_id, artifact_id, created_at)
           VALUES (?, ?, ?, ?, ?)""",
        (namespace, alias, source_id, artifact_id, datetime.now(timezone.utc).isoformat()),
    )

resolve_artifact_id(value, source_id=None)

Resolve a canonical ID or namespaced compatibility alias.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def resolve_artifact_id(self, value: str, source_id: str | None = None) -> str | None:
    """Resolve a canonical ID or namespaced compatibility alias."""
    with self._read_snapshot() as connection:
        return self._resolve_artifact_id(connection, value, source_id)

has_canonical_artifact_id(artifact_id)

Return whether an exact canonical artifact ID is retained.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def has_canonical_artifact_id(self, artifact_id: str) -> bool:
    """Return whether an exact canonical artifact ID is retained."""
    with self._read_snapshot() as connection:
        return (
            connection.execute(
                "SELECT 1 FROM artifacts WHERE id=?",
                (artifact_id,),
            ).fetchone()
            is not None
        )

resolve_scan_alias(value, source_id)

Resolve a scan alias within its source or against one unadopted migration record.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def resolve_scan_alias(self, value: str, source_id: str) -> str | None:
    """Resolve a scan alias within its source or against one unadopted migration record."""
    with self._read_snapshot() as connection:
        resolved = self._resolve_artifact_id(connection, value, source_id)
        if resolved is not None:
            return resolved
        namespace, alias = split_alias(value)
        rows = connection.execute(
            """SELECT DISTINCT a.id FROM artifact_aliases aa
               JOIN artifacts a ON a.id=aa.artifact_id
               WHERE aa.namespace=? AND aa.alias=?
                 AND a.source_id LIKE 'legacy-%'
                 AND a.logical_path LIKE 'imported/%'""",
            (namespace, alias),
        ).fetchall()
        if len(rows) > 1:
            raise ValueError(f"Ambiguous migrated artifact alias: {namespace}:{alias}")
        return rows[0]["id"] if rows else None

export_v0_json(destination=None)

Export current live artifacts as a deterministic v0-compatible JSON array.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def export_v0_json(self, destination: str | Path | None = None) -> str:
    """Export current live artifacts as a deterministic v0-compatible JSON array."""
    with self._read_snapshot() as connection:
        rows = connection.execute(
            """SELECT id, name, storage_uri, description, domain, source_type,
                      content_type, size_bytes, indexed_at
               FROM artifacts WHERE deleted_at IS NULL ORDER BY source_id, logical_path, id"""
        ).fetchall()
        records = []
        for row in rows:
            legacy_id = artifact_id_from_uri(row["storage_uri"])
            legacy_alias = connection.execute(
                """SELECT alias FROM artifact_aliases
                   WHERE artifact_id=? AND namespace='artifact-id'
                   ORDER BY CASE WHEN alias=? THEN 0 ELSE 1 END, alias LIMIT 1""",
                (row["id"], legacy_id),
            ).fetchone()
            record = dict(row)
            record["id"] = legacy_alias["alias"] if legacy_alias is not None else row["id"]
            records.append(record)
    payload = json.dumps(records, ensure_ascii=False, indent=2, sort_keys=True) + "\n"
    if destination is not None:
        Path(destination).write_text(payload, encoding="utf-8")
    return payload

get_artifact(artifact_id, *, source_id=None, revision=None, include_deleted=False)

Retrieve a current or revision-qualified record by ID or alias.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def get_artifact(
    self,
    artifact_id: str,
    *,
    source_id: str | None = None,
    revision: int | None = None,
    include_deleted: bool = False,
) -> Optional[ArtifactRecord]:
    """Retrieve a current or revision-qualified record by ID or alias."""
    with self._read_snapshot() as connection:
        resolved_id = self._resolve_artifact_id(connection, artifact_id, source_id)
        if resolved_id is None:
            return None
        if revision is None:
            row = connection.execute("SELECT * FROM artifacts WHERE id = ?", (resolved_id,)).fetchone()
        else:
            row = connection.execute(
                """SELECT artifact_id AS id, source_id, logical_path, name, storage_uri,
                          description, domain, source_type, content_type, size_bytes,
                          indexed_at, revision, content_revision, metadata_revision,
                          checksum_sha256, deleted_at
                   FROM artifact_revisions WHERE artifact_id=? AND revision=?""",
                (resolved_id, revision),
            ).fetchone()
        if row is None or (row["deleted_at"] is not None and not include_deleted):
            return None
        return _record_from_row(row)

list_revisions(artifact_id, *, source_id=None)

Return all retained revisions, including deletion tombstones.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def list_revisions(self, artifact_id: str, *, source_id: str | None = None) -> list[ArtifactRecord]:
    """Return all retained revisions, including deletion tombstones."""
    with self._read_snapshot() as connection:
        resolved_id = self._resolve_artifact_id(connection, artifact_id, source_id)
        if resolved_id is None:
            return []
        rows = connection.execute(
            """SELECT artifact_id AS id, source_id, logical_path, name, storage_uri,
                      description, domain, source_type, content_type, size_bytes,
                      indexed_at, revision, content_revision, metadata_revision,
                      checksum_sha256, deleted_at
               FROM artifact_revisions WHERE artifact_id=? ORDER BY revision""",
            (resolved_id,),
        ).fetchall()
        return [_record_from_row(row) for row in rows]

records_by_source_path(source_id, *, include_deleted=False)

Return one source's records keyed by normalized logical path.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def records_by_source_path(self, source_id: str, *, include_deleted: bool = False) -> dict[str, ArtifactRecord]:
    """Return one source's records keyed by normalized logical path."""
    deleted_filter = "" if include_deleted else " AND deleted_at IS NULL"
    with self._read_snapshot() as connection:
        rows = connection.execute(
            f"SELECT * FROM artifacts WHERE source_id=?{deleted_filter}",
            (source_id,),
        ).fetchall()
        return {row["logical_path"]: _record_from_row(row) for row in rows}

retained_artifact_ids_by_source_path(source_id)

Return current and revision artifact IDs grouped by retained logical path.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def retained_artifact_ids_by_source_path(self, source_id: str) -> dict[str, set[str]]:
    """Return current and revision artifact IDs grouped by retained logical path."""
    with self._read_snapshot() as connection:
        rows = connection.execute(
            """SELECT logical_path, id AS artifact_id
               FROM artifacts
               WHERE source_id=?
               UNION
               SELECT logical_path, artifact_id
               FROM artifact_revisions
               WHERE source_id=?""",
            (source_id, source_id),
        ).fetchall()
    retained: dict[str, set[str]] = {}
    for row in rows:
        retained.setdefault(row["logical_path"], set()).add(row["artifact_id"])
    return retained

delete_artifacts(artifact_ids, *, deleted_at=None)

Create tombstone revisions; history is retained until explicit purge.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def delete_artifacts(self, artifact_ids: list[str], *, deleted_at: str | None = None) -> None:
    """Create tombstone revisions; history is retained until explicit purge."""
    with self._write_transaction():
        self._delete_artifacts(artifact_ids, deleted_at=deleted_at)

purge_deleted(before)

Permanently remove tombstones older than before and all their revisions.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def purge_deleted(self, before: str) -> int:
    """Permanently remove tombstones older than *before* and all their revisions."""
    with self._write_transaction():
        rows = self.conn.execute(
            """SELECT id FROM artifacts
               WHERE deleted_at IS NOT NULL
                 AND julianday(deleted_at) < julianday(?)""",
            (before,),
        ).fetchall()
        ids = [row["id"] for row in rows]
        if not ids:
            return 0
        placeholders = ",".join("?" for _ in ids)
        has_vector_table = self._vector_loaded or self._has_vector_table(self.conn)
        if has_vector_table:
            self._ensure_vector_capability(self.conn, create_table=False)
            self.conn.execute(f"DELETE FROM artifacts_vec WHERE id IN ({placeholders})", ids)
        self.conn.execute(f"DELETE FROM artifact_aliases WHERE artifact_id IN ({placeholders})", ids)
        self.conn.execute(f"DELETE FROM artifact_revisions WHERE artifact_id IN ({placeholders})", ids)
        self.conn.execute(f"DELETE FROM artifacts WHERE id IN ({placeholders})", ids)
    return len(ids)

list_artifacts(*, source_ids=(), domain=None, source_type=None, limit=50, offset=0)

List current artifacts deterministically for public provider adapters.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def list_artifacts(
    self,
    *,
    source_ids: tuple[str, ...] = (),
    domain: str | None = None,
    source_type: str | None = None,
    limit: int = 50,
    offset: int = 0,
) -> list[ArtifactRecord]:
    """List current artifacts deterministically for public provider adapters."""
    conditions = ["deleted_at IS NULL"]
    params: list[object] = []
    if source_ids:
        placeholders = ",".join("?" for _ in source_ids)
        conditions.append(f"source_id IN ({placeholders})")
        params.extend(source_ids)
    if domain:
        conditions.append("domain = ?")
        params.append(domain)
    if source_type:
        conditions.append("source_type = ?")
        params.append(source_type)
    with self._read_snapshot() as connection:
        rows = connection.execute(
            f"""SELECT * FROM artifacts
                WHERE {" AND ".join(conditions)}
                ORDER BY source_id, logical_path, id LIMIT ? OFFSET ?""",
            (*params, max(0, limit), max(0, offset)),
        ).fetchall()
    return [_record_from_row(row) for row in rows]

CatalogDryRunReport(sources, configuration_valid=True) dataclass

Mutation-free catalog refresh preview.

CatalogIndexer(config, db, credential_provider=None, embedding_provider=None)

Scans configured sources and populates the catalog database.

Parameters:

Name Type Description Default
config CatalogConfig

Catalog configuration.

required
db CatalogDB

The catalog database instance.

required
credential_provider Optional[CredentialProvider]

Optional CredentialProvider from the auth module. Used for authenticating to blob storage. If None, blob sources will fall back to DefaultAzureCredential.

None
embedding_provider Optional[EmbeddingProvider]

Optional pre-configured EmbeddingProvider instance. If provided, bypasses the factory and uses this provider directly. Useful for custom embedding backends not covered by the built-in factory (e.g. Cohere, Ollama, HuggingFace Inference API).

None
Source code in src/agora_workbench/code_execution/data_access/catalog/indexer.py
def __init__(
    self,
    config: CatalogConfig,
    db: CatalogDB,
    credential_provider: Optional[CredentialProvider] = None,
    embedding_provider: Optional[EmbeddingProvider] = None,
):
    """
    Args:
        config: Catalog configuration.
        db: The catalog database instance.
        credential_provider: Optional CredentialProvider from the auth module.
            Used for authenticating to blob storage. If None, blob sources
            will fall back to DefaultAzureCredential.
        embedding_provider: Optional pre-configured EmbeddingProvider instance.
            If provided, bypasses the factory and uses this provider directly.
            Useful for custom embedding backends not covered by the built-in
            factory (e.g. Cohere, Ollama, HuggingFace Inference API).
    """
    self._config = config
    self._db = db
    self._credential_provider = credential_provider
    self._embedding_provider: Optional[EmbeddingProvider] = embedding_provider
    self._embedding_provider_lock = threading.Lock()
    self._manifest_revisions: dict[str, tuple[int, str]] = {}
    if not _USE_POSIX_DIR_FDS and any(source.source_type == "local" for source in config.sources):
        raise ValueError("Local catalog sources require POSIX descriptor-relative path operations.")

embedding_provider property

Resolve the embedding provider from config (None = keyword-only).

Re-resolves while unset; a None result (keyword-only / BM25) is cheap to recompute, and tests may inject _embedding_provider directly.

index() async

Run a full index pass: enumerate all sources (local + blob), diff against existing entries, compute embeddings, upsert.

Returns:

Type Description
int

Number of artifacts indexed (new + updated).

Source code in src/agora_workbench/code_execution/data_access/catalog/indexer.py
async def index(self) -> int:
    """
    Run a full index pass: enumerate all sources (local + blob),
    diff against existing entries, compute embeddings, upsert.

    Returns:
        Number of artifacts indexed (new + updated).
    """
    sources = self._validated_sources()
    self._manifest_revisions = {}
    enumeration = await self._enumerate_all_sources(sources)
    enumeration = self._isolate_manifest_candidate_conflicts(enumeration, sources)
    all_artifacts = enumeration.artifacts
    attempted_at = datetime.now(timezone.utc).isoformat()

    discovered_by_source: dict[str, set[str]] = {
        source_id: set() for source_id in enumeration.successful_source_ids
    }
    for artifact in all_artifacts:
        discovered_by_source.setdefault(artifact["source_id"], set()).add(artifact["logical_path"])

    records_by_source: dict[str, dict[str, ArtifactRecord]] = {}
    stale_ids: list[str] = []
    for source_id, discovered_paths in discovered_by_source.items():
        current = self._db.records_by_source_path(source_id, include_deleted=True)
        records_by_source[source_id] = current
        stale_ids.extend(
            record.id
            for path, record in current.items()
            if record.deleted_at is None and path not in discovered_paths
        )

    provider = self.embedding_provider
    existing_live_ids = [
        record.id
        for records in records_by_source.values()
        for record in records.values()
        if record.deleted_at is None
    ]
    missing_vector_ids = (
        self._db.missing_vectors(existing_live_ids) if provider is not None and existing_live_ids else set()
    )
    to_index: list[tuple[dict, bool]] = []
    for artifact in all_artifacts:
        existing = records_by_source[artifact["source_id"]].get(artifact["logical_path"])
        needs_vector = (
            existing is not None
            and existing.deleted_at is None
            and provider is not None
            and existing.id in missing_vector_ids
        )
        if (
            existing is None
            or existing.deleted_at is not None
            or existing.storage_uri != artifact["storage_uri"]
            or existing.content_revision != artifact["content_revision"]
            or existing.metadata_revision != artifact["metadata_revision"]
            or needs_vector
        ):
            searchable_change = (
                existing is None
                or existing.deleted_at is not None
                or existing.name != artifact["name"]
                or existing.description != artifact["description"]
                or existing.domain != artifact["domain"]
                or needs_vector
            )
            to_index.append((artifact, provider is not None and searchable_change))

    rows = await self._compute_rows(to_index)
    artifact_counts: dict[str, int] = {}
    for artifact in all_artifacts:
        artifact_counts[artifact["source_id"]] = artifact_counts.get(artifact["source_id"], 0) + 1
    source_results = []
    for source, source_id in sources:
        succeeded = source_id in enumeration.successful_source_ids
        manifest_revision = self._manifest_revisions.get(source_id)
        source_results.append(
            {
                "source_id": source_id,
                "source_type": source.source_type,
                "root_uri": (
                    canonicalize_azure_uri(source.path)
                    if source.source_type == "blob"
                    else str(Path(source.path).resolve())
                ),
                "attempted_at": attempted_at,
                "succeeded": succeeded,
                "artifact_count": artifact_counts.get(source_id, 0) if succeeded else None,
                "error": None if succeeded else enumeration.errors.get(source_id, "Source enumeration failed"),
                "manifest_generation": manifest_revision[0] if succeeded and manifest_revision else None,
                "manifest_etag": manifest_revision[1] if succeeded and manifest_revision else None,
            }
        )

    vector_model_id = (
        _embedding_model_id(self._config) if provider is not None and self._db.vec_dimensions is not None else None
    )
    self._db.apply_refresh_batch(rows, stale_ids, source_results, vector_model_id)
    if stale_ids:
        LOGGER.info("Tombstoned %d stale artifacts.", len(stale_ids))

    manifest_errors = {
        source_id: error
        for source_id, error in enumeration.errors.items()
        if next(source.discovery for source, candidate_id in sources if candidate_id == source_id)
        is DiscoveryMode.MANIFEST
    }
    if manifest_errors:
        raise ManifestRefreshError(manifest_errors)

    # Log warning for artifacts without descriptions
    no_desc_count = sum(1 for a in all_artifacts if not a.get("description"))
    if no_desc_count:
        LOGGER.warning(
            "%d artifacts have no description — search quality will be reduced.",
            no_desc_count,
        )

    if not to_index:
        if enumeration.errors:
            LOGGER.warning(
                "Catalog refresh preserved last valid state for %d failed source(s).",
                len(enumeration.errors),
            )
        else:
            LOGGER.info("Catalog up to date (%d artifacts).", len(all_artifacts))
        return 0
    LOGGER.info("Indexed %d new or updated artifacts (%d total).", len(to_index), len(all_artifacts))
    return len(to_index)

dry_run() async

Enumerate and diff sources without writing SQLite or computing embeddings.

Source code in src/agora_workbench/code_execution/data_access/catalog/indexer.py
async def dry_run(self) -> CatalogDryRunReport:
    """Enumerate and diff sources without writing SQLite or computing embeddings."""
    sources = self._validated_sources()
    self._manifest_revisions = {}
    enumeration = await self._enumerate_all_sources(sources)
    enumeration = self._isolate_manifest_candidate_conflicts(enumeration, sources)
    artifacts_by_source: dict[str, list[dict]] = {}
    for artifact in enumeration.artifacts:
        artifacts_by_source.setdefault(artifact["source_id"], []).append(artifact)

    planned: list[SourceDryRun] = []
    for source, source_id in sources:
        error = enumeration.errors.get(source_id)
        revision = self._manifest_revisions.get(source_id)
        if error is not None:
            planned.append(
                SourceDryRun(
                    source_id=source_id,
                    discovery=source.discovery.value,
                    manifest_generation=revision[0] if revision else None,
                    manifest_etag=revision[1] if revision else None,
                    manifest_checked=source.discovery is DiscoveryMode.MANIFEST,
                    manifest_content_valid=(False if source.discovery is DiscoveryMode.MANIFEST else None),
                    error=error,
                )
            )
            continue
        existing = self._db.records_by_source_path(source_id, include_deleted=True)
        discovered = {artifact["logical_path"]: artifact for artifact in artifacts_by_source.get(source_id, [])}
        added = updated = unchanged = 0
        for logical_path, artifact in discovered.items():
            record = existing.get(logical_path)
            if record is None:
                added += 1
            elif (
                record.deleted_at is not None
                or record.storage_uri != artifact["storage_uri"]
                or record.content_revision != artifact["content_revision"]
                or record.metadata_revision != artifact["metadata_revision"]
            ):
                updated += 1
            else:
                unchanged += 1
        deleted = sum(
            record.deleted_at is None and logical_path not in discovered
            for logical_path, record in existing.items()
        )
        planned.append(
            SourceDryRun(
                source_id=source_id,
                discovery=source.discovery.value,
                added=added,
                updated=updated,
                deleted=deleted,
                unchanged=unchanged,
                manifest_generation=revision[0] if revision else None,
                manifest_etag=revision[1] if revision else None,
                manifest_checked=source.discovery is DiscoveryMode.MANIFEST,
                manifest_content_valid=(True if source.discovery is DiscoveryMode.MANIFEST else None),
            )
        )
    return CatalogDryRunReport(tuple(planned))

DiscoveryMode

Bases: StrEnum

How a source exposes artifacts to the catalog.

SearchConfig

Bases: BaseModel

Search/embedding configuration.

SourceConfig

Bases: BaseModel

A single data source (directory or blob prefix).

source_type property

Infer storage type from path prefix.

SourceDryRun(source_id, discovery, added=0, updated=0, deleted=0, unchanged=0, manifest_generation=None, manifest_etag=None, configuration_valid=True, manifest_checked=False, manifest_content_valid=None, error=None) dataclass

Planned changes for one source without mutating SQLite.

SourceRefreshState(source_id, attempt_generation, successful_generation, status, artifact_count, last_attempt_at, last_success_at, error, manifest_generation=None, manifest_etag=None) dataclass

Most recent attempted and successful refresh state for one source.

ArtifactPresentation(name, description=None, media_type=None, size_bytes=None) dataclass

Human-facing catalog metadata.

ArtifactReference(artifact_id, source_id, revision=None) dataclass

Stable logical identity, optionally pinned to a provider-honored revision.

Providers and adapters that accept a non-None revision must resolve that exact retained revision or reject the request explicitly; they must not silently return the current revision.

is_current property

Whether the reference follows the current artifact revision.

CatalogAuthorizationRequest(operation, source_id, reference=None) dataclass

One source- or artifact-scoped authorization check.

CatalogArtifact(reference, presentation, locator=None, download=None, metadata=dict(), revision=None, content_revision=None, metadata_revision=None, checksum_sha256=None, deleted_at=None, score=None) dataclass

Backend-neutral artifact returned by catalog operations.

is_deleted property

Whether the catalog record is a deletion tombstone.

CatalogOperation

Bases: StrEnum

Catalog operations a provider or writer may support.

CatalogPolicyMode

Bases: StrEnum

Granularity at which caller policy is enforced.

DownloadInfo(url, filename=None, expires_at=None) dataclass

Optional presentation-layer download information.

ListRequest(source_ids=(), page=PageRequest(), filters=dict()) dataclass

Catalog listing request.

Page(items, next_cursor=None) dataclass

Bases: Generic[T]

One page of results and an opaque provider cursor.

PageRequest(limit=50, cursor=None) dataclass

Cursor-based pagination request, capped at :data:MAX_PAGE_LIMIT.

RequestContext(request_id=None, caller_id=None, attributes=dict()) dataclass

Request-scoped caller metadata copied for safe propagation.

ResolvedArtifact(reference, locator) dataclass

A logical artifact reference resolved to a physical locator.

SearchRequest(query, source_ids=(), page=PageRequest(), filters=dict()) dataclass

Catalog search request.

SourceCapabilities(source_id, supported_operations) dataclass

Read operations authoritatively supported by one provider source.

supports(operation)

Return whether the source reports support for operation.

Source code in src/agora_workbench/data_lake/models.py
def supports(self, operation: CatalogOperation) -> bool:
    """Return whether the source reports support for *operation*."""
    return operation in self.supported_operations

StorageLocator(uri) dataclass

Physical locator understood by a fetcher or storage provider.

AuthorizedCatalogProvider(provider, authorizer, *, mode, per_artifact_enforcer=None)

Compose caller policy around an untrusted-to-authorize provider.

Providers still enforce ordinary query constraints such as source_ids; they never decide caller policy. Per-artifact mode additionally requires a backend-specific policy enforcer because generic post-filtering cannot preserve ranking, pagination, aggregation, or alias confidentiality.

Source code in src/agora_workbench/data_lake/policy.py
def __init__(
    self,
    provider: CatalogProvider,
    authorizer: CatalogAuthorizer,
    *,
    mode: CatalogPolicyMode,
    per_artifact_enforcer: CatalogPolicyEnforcer | None = None,
) -> None:
    if mode is not CatalogPolicyMode.HOMOGENEOUS_SOURCE and mode is not CatalogPolicyMode.PER_ARTIFACT:
        raise ValueError(f"Unknown catalog policy mode: {mode!r}.")
    if mode is CatalogPolicyMode.PER_ARTIFACT and per_artifact_enforcer is None:
        raise ValueError("Per-artifact policy requires a backend-specific CatalogPolicyEnforcer.")
    if mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE and per_artifact_enforcer is not None:
        raise ValueError("A per-artifact enforcer cannot be used with homogeneous-source policy.")
    self._provider = provider
    self._authorizer = authorizer
    self._mode = mode
    self._per_artifact_enforcer = per_artifact_enforcer

policy_mode property

Return the configured enforcement granularity.

authorizer property

Return the caller-scoped authorizer used by this policy wrapper.

capabilities(context) async

Return provider support intersected with current caller policy.

Source code in src/agora_workbench/data_lake/policy.py
async def capabilities(self, context: RequestContext) -> tuple[SourceCapabilities, ...]:
    """Return provider support intersected with current caller policy."""
    effective: list[SourceCapabilities] = []
    for capability in await self._provider_capabilities():
        allowed_operations: set[CatalogOperation] = set()
        for operation in capability.supported_operations:
            if await self._authorize_source(operation, capability.source_id, context):
                allowed_operations.add(operation)
        allowed = frozenset(allowed_operations)
        if allowed:
            effective.append(SourceCapabilities(capability.source_id, allowed))
    return tuple(effective)

search(request, context) async

Search with authorization enforced before effective pagination.

Source code in src/agora_workbench/data_lake/policy.py
async def search(self, request: SearchRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """Search with authorization enforced before effective pagination."""
    constrained, allowed_sources = await self._constrain_request(request, CatalogOperation.SEARCH, context)
    if not allowed_sources:
        return Page(())
    if self._uses_per_artifact_enforcement():
        assert self._per_artifact_enforcer is not None
        page = await self._per_artifact_enforcer.search(
            self._provider,
            constrained,
            context,
            self._authorizer,
        )
    elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
        page = await self._provider.search(constrained, context)
    else:
        self._raise_enforcement_failure(CatalogOperation.SEARCH)
    await self._validate_page(page, allowed_sources, CatalogOperation.SEARCH, context)
    return page

list(request, context) async

List with authorization enforced before effective pagination.

Source code in src/agora_workbench/data_lake/policy.py
async def list(self, request: ListRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """List with authorization enforced before effective pagination."""
    constrained, allowed_sources = await self._constrain_request(request, CatalogOperation.LIST, context)
    if not allowed_sources:
        return Page(())
    if self._uses_per_artifact_enforcement():
        assert self._per_artifact_enforcer is not None
        page = await self._per_artifact_enforcer.list(
            self._provider,
            constrained,
            context,
            self._authorizer,
        )
    elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
        page = await self._provider.list(constrained, context)
    else:
        self._raise_enforcement_failure(CatalogOperation.LIST)
    await self._validate_page(page, allowed_sources, CatalogOperation.LIST, context)
    return page

get(reference, context) async

Get an artifact without distinguishing denied from absent artifacts.

Source code in src/agora_workbench/data_lake/policy.py
async def get(self, reference: ArtifactReference, context: RequestContext) -> CatalogArtifact:
    """Get an artifact without distinguishing denied from absent artifacts."""
    await self._require_source_operation(reference, CatalogOperation.GET, context)
    backend_denied = False
    try:
        if self._uses_per_artifact_enforcement():
            assert self._per_artifact_enforcer is not None
            artifact = await self._per_artifact_enforcer.get(
                self._provider,
                reference,
                context,
                self._authorizer,
            )
        elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
            artifact = await self._provider.get(reference, context)
        else:
            self._raise_enforcement_failure(CatalogOperation.GET)
    except ArtifactNotFoundError:
        pass
    except PermissionDeniedError:
        backend_denied = True
    else:
        self._validate_reference(artifact.reference, reference, CatalogOperation.GET)
        if self._uses_per_artifact_enforcement():
            await self._require_artifact_operation(artifact.reference, CatalogOperation.GET, context)
        return artifact
    if backend_denied:
        self._raise_enforcement_failure(CatalogOperation.GET)
    self._raise_not_found(CatalogOperation.GET)

resolve(reference, context) async

Resolve an artifact without distinguishing denied from absent artifacts.

Source code in src/agora_workbench/data_lake/policy.py
async def resolve(self, reference: ArtifactReference, context: RequestContext) -> ResolvedArtifact:
    """Resolve an artifact without distinguishing denied from absent artifacts."""
    await self._require_source_operation(reference, CatalogOperation.RESOLVE, context)
    backend_denied = False
    try:
        if self._uses_per_artifact_enforcement():
            assert self._per_artifact_enforcer is not None
            resolved = await self._per_artifact_enforcer.resolve(
                self._provider,
                reference,
                context,
                self._authorizer,
            )
        elif self._mode is CatalogPolicyMode.HOMOGENEOUS_SOURCE:
            resolved = await self._provider.resolve(reference, context)
        else:
            self._raise_enforcement_failure(CatalogOperation.RESOLVE)
    except ArtifactNotFoundError:
        pass
    except PermissionDeniedError:
        backend_denied = True
    else:
        self._validate_reference(resolved.reference, reference, CatalogOperation.RESOLVE)
        if self._uses_per_artifact_enforcement():
            await self._require_artifact_operation(resolved.reference, CatalogOperation.RESOLVE, context)
        return resolved
    if backend_denied:
        self._raise_enforcement_failure(CatalogOperation.RESOLVE)
    self._raise_not_found(CatalogOperation.RESOLVE)

DenyAllCatalogAuthorizer

Fail-closed policy suitable as an explicit production baseline.

DevelopmentAllowAllCatalogAuthorizer

Explicit development-only policy that permits every catalog request.

CatalogArtifactResolver(provider, source_id, context=None)

Legacy ArtifactResolver adapter for one source of a CatalogProvider.

Source code in src/agora_workbench/data_lake/providers.py
def __init__(
    self,
    provider: CatalogProvider,
    source_id: str,
    context: RequestContext | None = None,
):
    self._provider = provider
    self._source_id = source_id
    self._context = context or RequestContext()
    self._closed = False

CatalogReadiness(ready, stale, sources, reason=None) dataclass

Observable readiness and retained-refresh state.

ManifestCatalogProvider(config, *, db_path=':memory:', credential_provider=None, max_stale_seconds=300.0)

Bases: SQLiteCatalogProvider

Authoritative manifest catalog with a private rebuildable SQLite index.

Source code in src/agora_workbench/data_lake/providers.py
def __init__(
    self,
    config: CatalogConfig,
    *,
    db_path: str | Path = ":memory:",
    credential_provider: CredentialProvider | None = None,
    max_stale_seconds: float = 300.0,
):
    if not config.sources:
        raise ValueError("ManifestCatalogProvider requires at least one source")
    if any(source.discovery is not DiscoveryMode.MANIFEST for source in config.sources):
        raise ValueError("ManifestCatalogProvider accepts only discovery='manifest' sources")
    if max_stale_seconds < 0:
        raise ValueError("max_stale_seconds must be non-negative")
    if not isfinite(max_stale_seconds):
        raise ValueError("max_stale_seconds must be finite")
    self._config = config
    self._closed = False
    self._embedding_closed = False
    self._db_closed = False
    self._lifecycle_lock = asyncio.Lock()
    self._active_reads = 0
    self._reads_drained = asyncio.Event()
    self._reads_drained.set()
    self._last_states: tuple[SourceRefreshState, ...] = ()
    self._db_owned = CatalogDB(db_path, vec_dimensions=config.search.embedding_dimensions)
    try:
        self._db_owned.open()
        self._indexer = CatalogIndexer(config, self._db_owned, credential_provider=credential_provider)
        self._source_stale_limits = {
            source.source_id or "": (
                source.max_stale_seconds if source.max_stale_seconds is not None else max_stale_seconds
            )
            for source in config.sources
        }
        self._last_error: str | None = None
        self._load_attempted = False
        super().__init__(
            self._db_owned,
            tuple(source.source_id or "" for source in config.sources),
            query_embedder=self._embed_query,
            hybrid_alpha=config.search.hybrid_alpha,
        )
    except BaseException:
        self._db_owned.close()
        self._db_closed = True
        self._closed = True
        raise

load() async

Load or refresh all manifests, retaining the previous valid generation on failure.

Source code in src/agora_workbench/data_lake/providers.py
async def load(self) -> int:
    """Load or refresh all manifests, retaining the previous valid generation on failure."""
    async with self._lifecycle_lock:
        await self._reads_drained.wait()
        return await self._load_unlocked()

readiness()

Return readiness, stale bounds, and per-source refresh state.

Source code in src/agora_workbench/data_lake/providers.py
def readiness(self) -> CatalogReadiness:
    """Return readiness, stale bounds, and per-source refresh state."""
    states = self._current_source_states()
    if self._closed:
        return CatalogReadiness(
            False,
            self._last_error is not None,
            states,
            self._last_error or "Manifest catalog is closed.",
        )
    if not self._load_attempted:
        return CatalogReadiness(False, False, states, "Manifest catalog load() has not completed.")
    expected = set(self._source_ids)
    successful = {
        state.source_id
        for state in states
        if state.successful_generation > 0 and state.manifest_generation is not None
    }
    if not expected <= successful:
        return CatalogReadiness(False, False, states, self._last_error or "No valid manifest generation loaded.")
    if self._last_error is None:
        return CatalogReadiness(True, False, states)
    now = datetime.now(timezone.utc)
    expired_sources: list[str] = []
    for state in states:
        if state.last_success_at is None:
            expired_sources.append(state.source_id)
            continue
        try:
            success_at = _parse_success_timestamp(state.last_success_at)
        except (TypeError, ValueError):
            expired_sources.append(state.source_id)
            continue
        age = max(0.0, (now - success_at).total_seconds())
        if age > self._source_stale_limits[state.source_id]:
            expired_sources.append(state.source_id)
    if expired_sources:
        return CatalogReadiness(
            False,
            True,
            states,
            "Stale catalog data expired after a refresh failure for source(s): "
            + ", ".join(sorted(expired_sources)),
        )
    return CatalogReadiness(True, True, states, "Serving the last valid generation after a refresh failure.")

aclose() async

Close embedding resources and the private per-reader SQLite cache.

Source code in src/agora_workbench/data_lake/providers.py
async def aclose(self) -> None:
    """Close embedding resources and the private per-reader SQLite cache."""
    async with self._lifecycle_lock:
        await self._reads_drained.wait()
        await self._aclose_unlocked()

SQLiteCatalogProvider(db, source_ids, *, query_embedder=None, hybrid_alpha=0.5)

CatalogProvider adapter over an opened, caller-owned CatalogDB.

Source code in src/agora_workbench/data_lake/providers.py
def __init__(
    self,
    db: CatalogDB,
    source_ids: tuple[str, ...],
    *,
    query_embedder: Callable[[str], Awaitable[list[float] | None]] | None = None,
    hybrid_alpha: float = 0.5,
):
    self._db = db
    self._source_ids = tuple(dict.fromkeys(source_ids))
    self._query_embedder = query_embedder
    self._hybrid_alpha = hybrid_alpha
    if not self._source_ids or any(not source_id for source_id in self._source_ids):
        raise ValueError("SQLiteCatalogProvider requires at least one non-empty source_id")
    if not 0.0 <= hybrid_alpha <= 1.0:
        raise ValueError("hybrid_alpha must be between 0 and 1")

CatalogAuthorizer

Bases: Protocol

Application-supplied caller policy, composed outside a catalog provider.

A request without reference is source-scoped. In per-artifact mode, source-scoped approval advertises an operation as potentially available; the policy enforcer must additionally authorize each artifact before it can affect ranking, pagination, aggregation, lookup results, or errors.

authorize(request, context) async

Return whether the caller may perform the requested operation.

Source code in src/agora_workbench/data_lake/protocols.py
async def authorize(self, request: CatalogAuthorizationRequest, context: RequestContext) -> bool:
    """Return whether the caller may perform the requested operation."""
    ...

CatalogPolicyEnforcer

Bases: Protocol

Backend-specific, trusted enforcement for per-artifact policy.

Implementations must apply policy before ranking, pagination, aggregation, alias resolution, and not-found decisions. Filtering one provider page after retrieval does not satisfy this contract.

search(provider, request, context, authorizer) async

Search only artifacts authorized for the current caller.

Source code in src/agora_workbench/data_lake/protocols.py
async def search(
    self,
    provider: CatalogProvider,
    request: SearchRequest,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> Page[CatalogArtifact]:
    """Search only artifacts authorized for the current caller."""
    ...

list(provider, request, context, authorizer) async

List only artifacts authorized for the current caller.

Source code in src/agora_workbench/data_lake/protocols.py
async def list(
    self,
    provider: CatalogProvider,
    request: ListRequest,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> Page[CatalogArtifact]:
    """List only artifacts authorized for the current caller."""
    ...

get(provider, reference, context, authorizer) async

Authorize lookup semantics before returning the exact requested logical reference.

Source code in src/agora_workbench/data_lake/protocols.py
async def get(
    self,
    provider: CatalogProvider,
    reference: ArtifactReference,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> CatalogArtifact:
    """Authorize lookup semantics before returning the exact requested logical reference."""
    ...

resolve(provider, reference, context, authorizer) async

Resolve only an authorized canonical artifact.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(
    self,
    provider: CatalogProvider,
    reference: ArtifactReference,
    context: RequestContext,
    authorizer: CatalogAuthorizer,
) -> ResolvedArtifact:
    """Resolve only an authorized canonical artifact."""
    ...

CatalogProvider

Bases: Protocol

Backend-neutral asynchronous catalog interface.

Runtime checks only verify that named members exist. They do not validate signatures, async behavior, return types, or capability truthfulness. capabilities values are authoritative for provider support; callers must not infer support from method presence. Authorization policy is composed outside the provider protocol.

capabilities() async

Return authoritative read support for each source exposed by the provider.

Source code in src/agora_workbench/data_lake/protocols.py
async def capabilities(self) -> tuple[SourceCapabilities, ...]:
    """Return authoritative read support for each source exposed by the provider."""
    ...

search(request, context) async

Search artifacts using provider-defined ranking and opaque cursors.

Source code in src/agora_workbench/data_lake/protocols.py
async def search(self, request: SearchRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """Search artifacts using provider-defined ranking and opaque cursors."""
    ...

list(request, context) async

List artifacts using provider-defined filters and opaque cursors.

Source code in src/agora_workbench/data_lake/protocols.py
async def list(self, request: ListRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """List artifacts using provider-defined filters and opaque cursors."""
    ...

get(reference, context) async

Get one artifact, honoring or explicitly rejecting a pinned revision.

Source code in src/agora_workbench/data_lake/protocols.py
async def get(self, reference: ArtifactReference, context: RequestContext) -> CatalogArtifact:
    """Get one artifact, honoring or explicitly rejecting a pinned revision."""
    ...

resolve(reference, context) async

Resolve a logical reference, honoring or explicitly rejecting a pinned revision.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(self, reference: ArtifactReference, context: RequestContext) -> ResolvedArtifact:
    """Resolve a logical reference, honoring or explicitly rejecting a pinned revision."""
    ...

PolicyEnforcedCatalog

Bases: Protocol

Caller-aware catalog surface produced by policy composition.

policy_mode property

Return the configured enforcement granularity.

capabilities(context) async

Return provider support intersected with current caller policy.

Source code in src/agora_workbench/data_lake/protocols.py
async def capabilities(self, context: RequestContext) -> tuple[SourceCapabilities, ...]:
    """Return provider support intersected with current caller policy."""
    ...

search(request, context) async

Search within the current caller's effective authorization scope.

Source code in src/agora_workbench/data_lake/protocols.py
async def search(self, request: SearchRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """Search within the current caller's effective authorization scope."""
    ...

list(request, context) async

List within the current caller's effective authorization scope.

Source code in src/agora_workbench/data_lake/protocols.py
async def list(self, request: ListRequest, context: RequestContext) -> Page[CatalogArtifact]:
    """List within the current caller's effective authorization scope."""
    ...

get(reference, context) async

Get an artifact without disclosing unauthorized existence.

Source code in src/agora_workbench/data_lake/protocols.py
async def get(self, reference: ArtifactReference, context: RequestContext) -> CatalogArtifact:
    """Get an artifact without disclosing unauthorized existence."""
    ...

resolve(reference, context) async

Resolve an artifact without disclosing unauthorized existence.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(self, reference: ArtifactReference, context: RequestContext) -> ResolvedArtifact:
    """Resolve an artifact without disclosing unauthorized existence."""
    ...

ArtifactMetadata(name=None, description=None, domain=None, media_type=None, aliases=()) dataclass

Metadata supplied for a durable artifact registration.

AuthorizedManagedCatalogWriter(writer, authorizer)

Authorize every mutation before metadata or byte side effects.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, writer: ManagedCatalogWriter, authorizer: CatalogAuthorizer) -> None:
    self._writer = writer
    self._authorizer = authorizer

capabilities(context) async

Return write operations allowed for this session and source.

Source code in src/agora_workbench/data_lake/managed.py
async def capabilities(self, context: RequestContext) -> tuple[SourceCapabilities, ...]:
    """Return write operations allowed for this session and source."""
    allowed = {
        operation
        for operation in WRITE_OPERATIONS
        if await self._authorizer.authorize(
            CatalogAuthorizationRequest(operation, self._writer.source_id),
            context,
        )
    }
    return (SourceCapabilities(self._writer.source_id, frozenset(allowed)),) if allowed else ()

BlobManagedStorage(container_client, *, prefix='')

Blob-container adapter using create-only objects and ETag manifest CAS.

The caller owns the container client and its prefix. Transfers apply the shared safety options while this adapter retains lifecycle-specific ownership metadata and conditional state operations.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, container_client: object, *, prefix: str = "") -> None:
    self._container_client = container_client
    self._prefix = prefix.strip("/")

DeleteOutcome

Bases: StrEnum

Ownership-aware deletion postcondition.

LocalManagedStorage(root)

Create-exclusive local objects and atomically replaced manifests.

Writers use an advisory flock covering the full mutation. Files and containing directories are fsynced before the lock is released. Readers do not take the lock; atomic replacement means they observe either the old or new complete manifest.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, root: str | Path) -> None:
    self.root = Path(root).resolve()
    self.root.mkdir(parents=True, exist_ok=True)
    self._lock_path = self.root / _WRITER_LOCK_PATH

ManagedCatalogWriter(source_id, backend, *, max_conflict_retries=4, revision_retention=1, operation_lease_seconds=30.0, interruption_hook=None)

Storage-neutral managed write state machine.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(
    self,
    source_id: str,
    backend: _ManagedStorageBackend,
    *,
    max_conflict_retries: int = 4,
    revision_retention: int = 1,
    operation_lease_seconds: float = 30.0,
    interruption_hook: Callable[[str, str], None] | None = None,
) -> None:
    if not source_id:
        raise ValueError("source_id must be non-empty")
    if max_conflict_retries < 0:
        raise ValueError("max_conflict_retries must be non-negative")
    if revision_retention < 0:
        raise ValueError("revision_retention must be non-negative")
    if operation_lease_seconds <= 0:
        raise ValueError("operation_lease_seconds must be positive")
    self.source_id = source_id
    self._backend = backend
    self._manifest_path = RESERVED_MANIFEST_PATH
    self._max_conflict_retries = max_conflict_retries
    self._revision_retention = revision_retention
    self._operation_lease_seconds = operation_lease_seconds
    self._interruption_hook = interruption_hook

artifact_reference(path, artifact_id=None)

Return the normalized effective reference used for authorization.

Source code in src/agora_workbench/data_lake/managed.py
def artifact_reference(self, path: str, artifact_id: str | None = None) -> ArtifactReference:
    """Return the normalized effective reference used for authorization."""
    normalized = _normalize_catalog_path(path)
    return ArtifactReference(
        artifact_id or logical_artifact_id(self.source_id, normalized),
        self.source_id,
    )

read_manifest(minimum_generation=None) async

Read committed state directly from the writer's backing store.

Source code in src/agora_workbench/data_lake/managed.py
async def read_manifest(self, minimum_generation: int | None = None) -> CatalogManifest:
    """Read committed state directly from the writer's backing store."""
    raw, _ = await self._backend.read_json(self._manifest_path)
    manifest = (
        CatalogManifest.from_mapping(raw)
        if raw is not None
        else CatalogManifest(MANIFEST_VERSION, 0, ())  # internal empty baseline
    )
    if minimum_generation is not None and manifest.generation < minimum_generation:
        raise PreconditionFailedError(
            "The requested committed generation is not visible.",
            operation="read_manifest",
        )
    return manifest

register(request, context=RequestContext()) async

Snapshot caller-owned bytes without modifying or owning the source.

Source code in src/agora_workbench/data_lake/managed.py
async def register(
    self,
    request: RegisterArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Snapshot caller-owned bytes without modifying or owning the source."""
    operation_id = _validate_operation_id(request.operation_id)
    path = _normalize_catalog_path(request.path)
    storage_path = normalize_logical_path(request.storage_path)
    if is_reserved_provider_path(storage_path):
        raise InvalidRequestError("External registrations cannot target reserved managed paths.")
    if request.checksum_sha256 is None:
        raise InvalidRequestError(
            "External registration requires checksum_sha256 for an immutable snapshot.",
            operation="register",
        )
    if (
        request.transfer_options.expected_sha256 is not None
        and request.transfer_options.expected_sha256 != request.checksum_sha256
    ):
        raise InvalidRequestError(
            "TransferOptions checksum must match the registered checksum.",
            operation="register",
        )
    artifact_id = request.artifact_id or logical_artifact_id(self.source_id, path)
    revision_path = _revision_path(artifact_id, operation_id, ".data")
    provenance = self._provenance(CatalogOperation.REGISTER, operation_id, context, request)
    intent = self._intent(
        CatalogOperation.REGISTER,
        operation_id,
        path,
        artifact_id,
        request,
        owned_path=revision_path,
        storage_path=revision_path,
        checksum_sha256=request.checksum_sha256,
        metadata=request.metadata,
        provenance=provenance,
    )
    prior, lease_id = await self._begin_operation(operation_id, intent)
    if prior is not None:
        return await self._resume_remove_cleanup(prior)
    try:
        self._interrupt("after_intent", operation_id)
        await emit_transfer_diagnostic(
            request.transfer_options,
            TransferDiagnostic("register", "started", context, revision_path),
        )
        async with self._heartbeat(operation_id, lease_id):
            try:
                version = await self._backend.create_from_storage(
                    storage_path,
                    revision_path,
                    operation_id,
                    request.checksum_sha256,
                    request.transfer_options,
                )
            except FileNotFoundError as exc:
                raise ArtifactNotFoundError("External bytes were not found.", operation="register") from exc
            except _StorageConflict:
                version = await self._backend.exists(revision_path)
                if (
                    version is None
                    or version.operation_id != operation_id
                    or version.checksum_sha256 != request.checksum_sha256
                ):
                    raise ConflictError(
                        "Immutable snapshot path is occupied by bytes not owned by this operation.",
                        operation="register",
                    ) from None
            except BaseException as exc:
                await emit_transfer_diagnostic(
                    request.transfer_options,
                    TransferDiagnostic(
                        "register",
                        "failed",
                        context,
                        revision_path,
                        error_type=type(exc).__name__,
                    ),
                )
                raise
            if request.size_bytes is not None and request.size_bytes != version.size_bytes:
                await self._backend.delete_owned(revision_path, operation_id, version.token)
                error = PreconditionFailedError(
                    "Registered size did not match the immutable snapshot.",
                    operation="register",
                )
                await emit_transfer_diagnostic(
                    request.transfer_options,
                    TransferDiagnostic(
                        "register",
                        "failed",
                        context,
                        revision_path,
                        error_type=type(error).__name__,
                    ),
                )
                raise error
        await emit_transfer_diagnostic(
            request.transfer_options,
            TransferDiagnostic(
                "register",
                "completed",
                context,
                revision_path,
                version.size_bytes,
                version.checksum_sha256,
            ),
        )
        await self._update_operation(
            operation_id,
            lease_id,
            ownership_verified=True,
            owned_version_token=version.token,
        )
        self._interrupt("after_object", operation_id)
        await self._update_operation(operation_id, lease_id, state="commit_ready")
        revision = ManifestRevision(
            revision_id=operation_id,
            storage_path=revision_path,
            content_revision=request.content_revision or request.checksum_sha256,
            created_at=provenance.created_at,
            operation_id=operation_id,
            ownership=ManifestOwnership.MANAGED,
            checksum_sha256=request.checksum_sha256,
            size_bytes=version.size_bytes,
            version_token=version.token,
            provenance=provenance,
        )
        return await self._commit_revision(
            CatalogOperation.REGISTER,
            operation_id,
            path,
            artifact_id,
            request.metadata,
            revision,
            provenance,
            request.expected_generation,
            request.expected_revision_id,
            lease_id,
        )
    except BaseException:
        await self._abandon_operation(operation_id, lease_id)
        raise

upload(request, context=RequestContext()) async

Create an immutable managed revision and commit it to the manifest.

Source code in src/agora_workbench/data_lake/managed.py
async def upload(
    self,
    request: UploadArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Create an immutable managed revision and commit it to the manifest."""
    return await self._upload(request, context, CatalogOperation.UPLOAD)

promote(request, context=RequestContext()) async

Explicitly copy a scratch output into durable managed storage.

Source code in src/agora_workbench/data_lake/managed.py
async def promote(
    self,
    request: PromoteOutputRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Explicitly copy a scratch output into durable managed storage."""
    if not request.session_id or not request.output_name:
        raise InvalidRequestError("Promotion requires session_id and output_name.", operation="promote")
    return await self._upload(request, context, CatalogOperation.PROMOTE)

remove(request, context=RequestContext()) async

Commit a tombstone before collecting any managed bytes.

Source code in src/agora_workbench/data_lake/managed.py
async def remove(
    self,
    request: RemoveArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Commit a tombstone before collecting any managed bytes."""
    operation_id = _validate_operation_id(request.operation_id)
    if request.reference.source_id != self.source_id:
        raise ArtifactNotFoundError("Artifact was not found.", operation="remove")
    if request.reference.revision is not None:
        raise InvalidRequestError(
            "Managed removal uses expected_revision_id rather than a catalog revision number.",
            operation="remove",
        )
    intent = self._intent(
        CatalogOperation.REMOVE,
        operation_id,
        "",
        request.reference.artifact_id,
        request,
        provenance=self._provenance(CatalogOperation.REMOVE, operation_id, context, request),
    )
    prior, lease_id = await self._begin_operation(operation_id, intent)
    if prior is not None:
        return await self._resume_remove_cleanup(prior)
    provenance = self._provenance(CatalogOperation.REMOVE, operation_id, context, request)
    try:
        self._interrupt("after_intent", operation_id)
        await self._update_operation(operation_id, lease_id, state="commit_ready")
        return await self._remove_under_lease(request, provenance, operation_id, lease_id)
    except BaseException:
        await self._abandon_operation(operation_id, lease_id)
        raise

reconcile(*, grace_seconds=300.0) async

Claim abandoned operations before recovery or owned-orphan cleanup.

Source code in src/agora_workbench/data_lake/managed.py
async def reconcile(self, *, grace_seconds: float = 300.0) -> ReconciliationReport:
    """Claim abandoned operations before recovery or owned-orphan cleanup."""
    if grace_seconds < 0:
        raise ValueError("grace_seconds must be non-negative")
    recovered: list[str] = []
    removed: list[str] = []
    deferred: list[str] = []
    active_operation_ids: set[str] = set()
    failures: dict[str, str] = {}
    for intent_path, listed_intent in await self._backend.list_json(RESERVED_OPERATIONS_PREFIX.rstrip("/")):
        operation_id = str(listed_intent.get("operation_id", ""))
        try:
            async with self._backend.serialized():
                intent, token = await self._backend.read_json(intent_path)
                if intent is None:
                    continue
                receipt, _ = await self._backend.read_json(self._receipt_path(operation_id))
                if receipt is not None:
                    result = self._result_from_mapping(receipt)
                    resumed = await self._resume_remove_cleanup(result)
                    if resumed != result:
                        recovered.append(operation_id)
                    continue
                manifest, _ = await self._load()
                artifact = self._find_operation_artifact(manifest, operation_id)
                if artifact is not None:
                    pass
                else:
                    created_at = datetime.fromisoformat(str(intent["created_at"]).replace("Z", "+00:00"))
                    age = (datetime.now(timezone.utc) - created_at).total_seconds()
                    if age < grace_seconds or (
                        intent.get("state") in {"active", "commit_ready"} and not self._lease_expired(intent)
                    ):
                        active_operation_ids.add(operation_id)
                        deferred.append(operation_id)
                        continue
                    if intent.get("state") in {"commit_ready", "fencing"}:
                        artifact, manifest, intent = await self._fence_expired_commit(
                            intent_path,
                            intent,
                            token,
                            operation_id,
                        )
                        if manifest is None or intent is None:
                            deferred.append(operation_id)
                            continue
                    elif intent.get("state") == "reconciling":
                        pass
                    else:
                        claim_id = uuid.uuid4().hex
                        claimed = {**intent, "state": "reconciling", "reconciliation_claim": claim_id}
                        try:
                            await self._backend.replace_json(intent_path, claimed, token)
                        except _StorageConflict:
                            deferred.append(operation_id)
                            continue
                        intent = claimed
                        manifest, _ = await self._load()
                        artifact = self._find_operation_artifact(manifest, operation_id)
            if artifact is not None:
                removal = next(
                    (item for item in artifact.removals if item.operation_id == operation_id),
                    None,
                )
                committed_revision = next(
                    (item for item in artifact.revisions if item.operation_id == operation_id),
                    None,
                )
                result = ManagedWriteResult(
                    operation_id,
                    self.source_id,
                    artifact.artifact_id or "",
                    (
                        removal.generation
                        if removal is not None
                        else (
                            committed_revision.committed_generation
                            if committed_revision is not None
                            and committed_revision.committed_generation is not None
                            else manifest.generation
                        )
                    ),
                    removal.revision_id
                    if removal is not None
                    else (
                        committed_revision.revision_id if committed_revision is not None else artifact.revision_id
                    ),
                    removal.storage_path
                    if removal is not None
                    else (
                        committed_revision.storage_path if committed_revision is not None else artifact.storage_path
                    ),
                    deleted=removal is not None or artifact.deleted_at is not None,
                    cleanup_pending=removal.garbage_collect if removal is not None else False,
                )
                await self._write_receipt(result)
                result = await self._resume_remove_cleanup(result)
                recovered.append(operation_id)
                continue
            owned_path = intent.get("owned_path")
            owned_version_token = intent.get("owned_version_token")
            if isinstance(owned_path, str) and (
                intent.get("ownership_verified") is not True or not isinstance(owned_version_token, str)
            ):
                owned = await self._backend.exists(owned_path)
                if owned is not None and owned.operation_id == operation_id:
                    current, token = await self._backend.read_json(intent_path)
                    if (
                        current is not None
                        and current.get("state") in {"fencing", "reconciling"}
                        and current.get("reconciliation_claim") == intent.get("reconciliation_claim")
                    ):
                        recovered_intent = {
                            **current,
                            "ownership_verified": True,
                            "owned_version_token": owned.token,
                            "checksum_sha256": owned.checksum_sha256,
                        }
                        try:
                            await self._backend.replace_json(intent_path, recovered_intent, token)
                        except _StorageConflict:
                            deferred.append(operation_id)
                            continue
                        intent = recovered_intent
                        owned_version_token = owned.token
            if (
                isinstance(owned_path, str)
                and intent.get("ownership_verified") is True
                and isinstance(owned_version_token, str)
            ):
                outcome = await self._backend.delete_owned(owned_path, operation_id, owned_version_token)
                if outcome in {DeleteOutcome.DELETED, DeleteOutcome.ABSENT}:
                    if intent.get("state") != "fencing" or await self._finish_fenced_operation(
                        intent_path,
                        intent,
                        operation_id,
                    ):
                        if intent.get("state") != "reconciling" or await self._finish_cleanup_operation(
                            intent_path,
                            intent,
                        ):
                            removed.append(operation_id)
                        else:
                            deferred.append(operation_id)
                    else:
                        deferred.append(operation_id)
                else:
                    failures[operation_id] = "Owned orphan could not be verified or removed."
            else:
                if intent.get("state") == "fencing":
                    await self._finish_fenced_operation(intent_path, intent, operation_id)
                elif intent.get("state") == "reconciling":
                    await self._finish_cleanup_operation(intent_path, intent)
                deferred.append(operation_id)
        except Exception as exc:
            failures[operation_id] = type(exc).__name__
    if failures:
        raise ReconciliationError(
            "One or more interrupted operations could not be reconciled.",
            failures=failures,
            operation="reconcile",
        )
    async with self._backend.serialized():
        removed_staging, deferred_staging = await self._backend.reconcile_staging(
            frozenset(active_operation_ids),
            grace_seconds,
        )
    return ReconciliationReport(
        tuple(recovered),
        tuple(removed),
        tuple(deferred),
        removed_staging,
        deferred_staging,
        failures,
    )

ManagedWriteResult(operation_id, source_id, artifact_id, generation, revision_id, storage_path, deleted=False, cleanup_pending=False) dataclass

Committed write result and its read-after-write generation.

PromoteOutputRequest(operation_id, path, local_path, artifact_id=None, metadata=ArtifactMetadata(), expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions(), session_id='', output_name='') dataclass

Bases: UploadArtifactRequest

Explicitly promote a scratch/session output to a durable artifact.

ReconciliationReport(recovered=(), removed_orphans=(), deferred=(), removed_staging=(), deferred_staging=(), failures=dict()) dataclass

Outcome of scanning interrupted operation intents.

RegisterArtifactRequest(operation_id, path, storage_path, artifact_id=None, metadata=ArtifactMetadata(), content_revision=None, checksum_sha256=None, size_bytes=None, expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions()) dataclass

Snapshot caller-owned bytes into a durable immutable managed revision.

RemoveArtifactRequest(operation_id, reference, expected_generation=None, expected_revision_id=None, garbage_collect=True) dataclass

Commit a tombstone and then optionally collect owned revision bytes.

UploadArtifactRequest(operation_id, path, local_path, artifact_id=None, metadata=ArtifactMetadata(), expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions()) dataclass

Upload local bytes as a new immutable managed revision.

artifact_id_from_uri(uri)

Return the legacy URI-derived artifact ID.

New catalog records use persisted logical IDs. This helper is retained for importing and resolving IDs written by earlier schema versions.

Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
def artifact_id_from_uri(uri: str) -> str:
    """Return the legacy URI-derived artifact ID.

    New catalog records use persisted logical IDs. This helper is retained for
    importing and resolving IDs written by earlier schema versions.
    """
    return hashlib.sha256(uri.encode()).hexdigest()[:16]

convert_catalog_config(source, destination=None, *, dry_run=True)

Validate and render explicit version-1 catalog configuration.

Dry-run is the default and never writes. A non-dry-run conversion requires a separate destination so an existing public configuration is not overwritten implicitly.

Source code in src/agora_workbench/code_execution/data_access/catalog/config.py
def convert_catalog_config(
    source: str | Path,
    destination: str | Path | None = None,
    *,
    dry_run: bool = True,
) -> CatalogConfigConversionReport:
    """Validate and render explicit version-1 catalog configuration.

    Dry-run is the default and never writes. A non-dry-run conversion requires a
    separate destination so an existing public configuration is not overwritten
    implicitly.
    """
    source_path = Path(source)
    if not source_path.exists():
        raise FileNotFoundError(f"Catalog config not found: {source_path}")
    raw = yaml.safe_load(source_path.read_text(encoding="utf-8")) or {}
    if not isinstance(raw, dict):
        raise ValueError("Catalog configuration must contain a YAML object")
    config, summary = CatalogConfig.convert(raw)
    explicit_sources = []
    for source_config in config.sources:
        rendered_source = source_config.model_dump(
            mode="json",
            exclude_none=True,
            exclude_defaults=True,
        )
        rendered_source["path"] = source_config.path
        rendered_source["discovery"] = source_config.discovery.value
        explicit_sources.append(rendered_source)
    explicit: dict[str, object] = {
        "version": config.version,
        "sources": explicit_sources,
    }
    if "search" in raw or config.search != SearchConfig():
        explicit["search"] = config.search.model_dump(
            mode="json",
            exclude_none=True,
            exclude_defaults=True,
        )
    rendered = yaml.safe_dump(explicit, sort_keys=False)
    destination_path = Path(destination) if destination is not None else None
    if not dry_run:
        if destination_path is None:
            raise ValueError("A destination is required when dry_run is False")
        if destination_path.resolve() == source_path.resolve():
            raise ValueError("Catalog conversion destination must differ from the source")
        destination_path.write_text(rendered, encoding="utf-8")
    return CatalogConfigConversionReport(
        source=source_path,
        destination=destination_path,
        changed=explicit != raw,
        written=not dry_run,
        rendered_yaml=rendered,
        summary=summary,
    )

azure_uri_from_blob_name(account, container, blob_name)

Build a canonical URI from an SDK-decoded blob name, quoting exactly once.

Source code in src/agora_workbench/data_lake/identity.py
def azure_uri_from_blob_name(account: str, container: str, blob_name: str) -> str:
    """Build a canonical URI from an SDK-decoded blob name, quoting exactly once."""
    account, container, _ = parse_azure_uri(f"az://{account}/{container}")
    encoded_path = quote(blob_name, safe="/-._~")
    return f"az://{account}/{container}/{encoded_path}" if encoded_path else f"az://{account}/{container}"

canonicalize_azure_uri(uri)

Canonicalize a supported Azure Blob or DFS URI without credentials.

Source code in src/agora_workbench/data_lake/identity.py
def canonicalize_azure_uri(uri: str) -> str:
    """Canonicalize a supported Azure Blob or DFS URI without credentials."""
    account, container, object_path = parse_azure_uri(uri)
    return azure_uri_from_blob_name(account, container, object_path)

logical_artifact_id(source_id, logical_path)

Generate a location-independent ID for a newly discovered logical path.

Source code in src/agora_workbench/data_lake/identity.py
def logical_artifact_id(source_id: str, logical_path: str) -> str:
    """Generate a location-independent ID for a newly discovered logical path."""
    normalized = normalize_logical_path(logical_path)
    return uuid.uuid5(_LOGICAL_ID_NAMESPACE, f"{source_id}\0{normalized}").hex

normalize_logical_path(path)

Return a portable, source-relative POSIX path.

Source code in src/agora_workbench/data_lake/identity.py
def normalize_logical_path(path: str) -> str:
    """Return a portable, source-relative POSIX path."""
    candidate = path.replace("\\", "/")
    if candidate.startswith("/") or re.match(r"^[a-zA-Z]:[\\/]", path):
        raise _invalid_identity("Artifact path must be source-relative.")
    normalized = posixpath.normpath(candidate)
    if normalized in {"", "."}:
        raise _invalid_identity("Artifact path must identify an object.")
    if normalized == ".." or normalized.startswith("../"):
        raise _invalid_identity("Artifact path must stay within its source.")
    return str(PurePosixPath(normalized))

parse_azure_uri(uri)

Parse a supported Azure URI into account, container, and decoded object path.

Source code in src/agora_workbench/data_lake/identity.py
def parse_azure_uri(uri: str) -> tuple[str, str, str]:
    """Parse a supported Azure URI into account, container, and decoded object path."""
    if _MALFORMED_PERCENT_RE.search(uri):
        raise _invalid_identity("Azure storage URI contains malformed percent encoding.")
    parsed = urlsplit(uri)
    scheme = parsed.scheme.lower()
    if scheme != "abfss" and (parsed.username is not None or parsed.password is not None):
        raise _invalid_identity("Azure storage URI must not contain user information.")
    if scheme != "abfss":
        try:
            if parsed.port is not None:
                raise _invalid_identity("Azure storage URI ports are not supported.")
        except ValueError as exc:
            raise _invalid_identity("Azure storage URI contains an invalid port.") from exc

    if scheme == "az":
        account = parsed.hostname or ""
        parts = parsed.path.lstrip("/").split("/", 1)
        container = parts[0] if parts else ""
        encoded_path = parts[1] if len(parts) > 1 else ""
    elif scheme in {"http", "https"}:
        host = (parsed.hostname or "").lower()
        suffix = next(
            (
                candidate
                for candidate in (".blob.core.windows.net", ".dfs.core.windows.net")
                if host.endswith(candidate)
            ),
            None,
        )
        if suffix is None:
            raise _invalid_identity("Unsupported Azure storage URI host.")
        account = host[: -len(suffix)]
        parts = parsed.path.lstrip("/").split("/", 1)
        container = parts[0] if parts else ""
        encoded_path = parts[1] if len(parts) > 1 else ""
    elif scheme == "abfss":
        if parsed.netloc.count("@") != 1:
            raise _invalid_identity("Malformed abfss URI.")
        encoded_container, host = parsed.netloc.split("@", 1)
        if ":" in host:
            raise _invalid_identity("Azure storage URI ports are not supported.")
        suffix = ".dfs.core.windows.net"
        if not host.lower().endswith(suffix):
            raise _invalid_identity("Unsupported Azure storage URI host.")
        account = host[: -len(suffix)]
        container = encoded_container
        encoded_path = parsed.path.lstrip("/")
    else:
        raise _invalid_identity("Unsupported Azure storage URI scheme.")

    if _ENCODED_SEPARATOR_RE.search(encoded_path):
        raise _invalid_identity("Azure object paths must not contain encoded separators.")
    account = account.lower()
    container = unquote(container).lower()
    if not _ACCOUNT_RE.fullmatch(account):
        raise _invalid_identity("Azure storage account name is malformed.")
    if container not in _SYSTEM_CONTAINERS and (not _CONTAINER_RE.fullmatch(container) or "--" in container):
        raise _invalid_identity("Azure storage container name is malformed.")
    object_path = unquote(encoded_path)
    validate_azure_object_path(object_path, allow_empty=True, allow_reserved=True)
    return account, container, object_path

sanitize_uri_for_display(uri)

Remove credentials, query parameters, and fragments from a URI.

Source code in src/agora_workbench/data_lake/identity.py
def sanitize_uri_for_display(uri: str) -> str:
    """Remove credentials, query parameters, and fragments from a URI."""
    parsed = urlsplit(uri)
    if parsed.scheme.lower() == "abfss":
        container = parsed.username or ""
        hostname = parsed.hostname or ""
        netloc = f"{container}@{hostname}" if container else hostname
        return urlunsplit((parsed.scheme, netloc, parsed.path, "", ""))
    hostname = parsed.hostname or ""
    try:
        port = parsed.port
    except ValueError:
        port = None
    if port is not None:
        hostname = f"{hostname}:{port}"
    return urlunsplit((parsed.scheme, hostname, parsed.path, "", ""))

split_alias(value, default_namespace='artifact-id')

Split a namespaced alias while retaining compatibility with opaque IDs.

Source code in src/agora_workbench/data_lake/identity.py
def split_alias(value: str, default_namespace: str = "artifact-id") -> tuple[str, str]:
    """Split a namespaced alias while retaining compatibility with opaque IDs."""
    if ":" not in value:
        return default_namespace, value
    namespace, alias = value.split(":", 1)
    if not namespace or not alias:
        raise _invalid_identity("Artifact aliases require non-empty namespace and value.")
    return namespace, alias

stable_source_id(source_type, root)

Derive a stable fallback source ID.

Source code in src/agora_workbench/data_lake/identity.py
def stable_source_id(source_type: str, root: str) -> str:
    """Derive a stable fallback source ID."""
    identity_root = canonicalize_azure_uri(root) if source_type == "blob" else str(root)
    digest = hashlib.sha256(f"{source_type}\0{identity_root}".encode()).hexdigest()[:20]
    return f"{source_type}-{digest}"

managed_writer_extension_factory(writer)

Adapt a managed writer to CatalogIntegration.capability_extension_factory.

Source code in src/agora_workbench/data_lake/managed.py
def managed_writer_extension_factory(writer: ManagedCatalogWriter):
    """Adapt a managed writer to ``CatalogIntegration.capability_extension_factory``."""

    def create(_session: object, catalog: object, _context: RequestContext) -> AuthorizedManagedCatalogWriter:
        authorizer = getattr(catalog, "authorizer", None)
        if authorizer is None:
            raise TypeError("Managed writer integration requires an authorized catalog.")
        return AuthorizedManagedCatalogWriter(writer, authorizer)

    return create

Managed writes

agora_workbench.data_lake.managed

Recoverable, optimistic-concurrency managed catalog writes.

Managed writes deliberately sit above asset transfer. Revision bytes are created first at immutable paths, then a manifest generation is conditionally committed. Those two steps are recoverable, but are not an atomic transaction.

ArtifactMetadata(name=None, description=None, domain=None, media_type=None, aliases=()) dataclass

Metadata supplied for a durable artifact registration.

RegisterArtifactRequest(operation_id, path, storage_path, artifact_id=None, metadata=ArtifactMetadata(), content_revision=None, checksum_sha256=None, size_bytes=None, expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions()) dataclass

Snapshot caller-owned bytes into a durable immutable managed revision.

UploadArtifactRequest(operation_id, path, local_path, artifact_id=None, metadata=ArtifactMetadata(), expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions()) dataclass

Upload local bytes as a new immutable managed revision.

PromoteOutputRequest(operation_id, path, local_path, artifact_id=None, metadata=ArtifactMetadata(), expected_generation=None, expected_revision_id=None, transfer_options=TransferOptions(), session_id='', output_name='') dataclass

Bases: UploadArtifactRequest

Explicitly promote a scratch/session output to a durable artifact.

RemoveArtifactRequest(operation_id, reference, expected_generation=None, expected_revision_id=None, garbage_collect=True) dataclass

Commit a tombstone and then optionally collect owned revision bytes.

ManagedWriteResult(operation_id, source_id, artifact_id, generation, revision_id, storage_path, deleted=False, cleanup_pending=False) dataclass

Committed write result and its read-after-write generation.

ReconciliationReport(recovered=(), removed_orphans=(), deferred=(), removed_staging=(), deferred_staging=(), failures=dict()) dataclass

Outcome of scanning interrupted operation intents.

DeleteOutcome

Bases: StrEnum

Ownership-aware deletion postcondition.

LocalManagedStorage(root)

Create-exclusive local objects and atomically replaced manifests.

Writers use an advisory flock covering the full mutation. Files and containing directories are fsynced before the lock is released. Readers do not take the lock; atomic replacement means they observe either the old or new complete manifest.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, root: str | Path) -> None:
    self.root = Path(root).resolve()
    self.root.mkdir(parents=True, exist_ok=True)
    self._lock_path = self.root / _WRITER_LOCK_PATH

BlobManagedStorage(container_client, *, prefix='')

Blob-container adapter using create-only objects and ETag manifest CAS.

The caller owns the container client and its prefix. Transfers apply the shared safety options while this adapter retains lifecycle-specific ownership metadata and conditional state operations.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, container_client: object, *, prefix: str = "") -> None:
    self._container_client = container_client
    self._prefix = prefix.strip("/")

ManagedCatalogWriter(source_id, backend, *, max_conflict_retries=4, revision_retention=1, operation_lease_seconds=30.0, interruption_hook=None)

Storage-neutral managed write state machine.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(
    self,
    source_id: str,
    backend: _ManagedStorageBackend,
    *,
    max_conflict_retries: int = 4,
    revision_retention: int = 1,
    operation_lease_seconds: float = 30.0,
    interruption_hook: Callable[[str, str], None] | None = None,
) -> None:
    if not source_id:
        raise ValueError("source_id must be non-empty")
    if max_conflict_retries < 0:
        raise ValueError("max_conflict_retries must be non-negative")
    if revision_retention < 0:
        raise ValueError("revision_retention must be non-negative")
    if operation_lease_seconds <= 0:
        raise ValueError("operation_lease_seconds must be positive")
    self.source_id = source_id
    self._backend = backend
    self._manifest_path = RESERVED_MANIFEST_PATH
    self._max_conflict_retries = max_conflict_retries
    self._revision_retention = revision_retention
    self._operation_lease_seconds = operation_lease_seconds
    self._interruption_hook = interruption_hook

artifact_reference(path, artifact_id=None)

Return the normalized effective reference used for authorization.

Source code in src/agora_workbench/data_lake/managed.py
def artifact_reference(self, path: str, artifact_id: str | None = None) -> ArtifactReference:
    """Return the normalized effective reference used for authorization."""
    normalized = _normalize_catalog_path(path)
    return ArtifactReference(
        artifact_id or logical_artifact_id(self.source_id, normalized),
        self.source_id,
    )

read_manifest(minimum_generation=None) async

Read committed state directly from the writer's backing store.

Source code in src/agora_workbench/data_lake/managed.py
async def read_manifest(self, minimum_generation: int | None = None) -> CatalogManifest:
    """Read committed state directly from the writer's backing store."""
    raw, _ = await self._backend.read_json(self._manifest_path)
    manifest = (
        CatalogManifest.from_mapping(raw)
        if raw is not None
        else CatalogManifest(MANIFEST_VERSION, 0, ())  # internal empty baseline
    )
    if minimum_generation is not None and manifest.generation < minimum_generation:
        raise PreconditionFailedError(
            "The requested committed generation is not visible.",
            operation="read_manifest",
        )
    return manifest

register(request, context=RequestContext()) async

Snapshot caller-owned bytes without modifying or owning the source.

Source code in src/agora_workbench/data_lake/managed.py
async def register(
    self,
    request: RegisterArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Snapshot caller-owned bytes without modifying or owning the source."""
    operation_id = _validate_operation_id(request.operation_id)
    path = _normalize_catalog_path(request.path)
    storage_path = normalize_logical_path(request.storage_path)
    if is_reserved_provider_path(storage_path):
        raise InvalidRequestError("External registrations cannot target reserved managed paths.")
    if request.checksum_sha256 is None:
        raise InvalidRequestError(
            "External registration requires checksum_sha256 for an immutable snapshot.",
            operation="register",
        )
    if (
        request.transfer_options.expected_sha256 is not None
        and request.transfer_options.expected_sha256 != request.checksum_sha256
    ):
        raise InvalidRequestError(
            "TransferOptions checksum must match the registered checksum.",
            operation="register",
        )
    artifact_id = request.artifact_id or logical_artifact_id(self.source_id, path)
    revision_path = _revision_path(artifact_id, operation_id, ".data")
    provenance = self._provenance(CatalogOperation.REGISTER, operation_id, context, request)
    intent = self._intent(
        CatalogOperation.REGISTER,
        operation_id,
        path,
        artifact_id,
        request,
        owned_path=revision_path,
        storage_path=revision_path,
        checksum_sha256=request.checksum_sha256,
        metadata=request.metadata,
        provenance=provenance,
    )
    prior, lease_id = await self._begin_operation(operation_id, intent)
    if prior is not None:
        return await self._resume_remove_cleanup(prior)
    try:
        self._interrupt("after_intent", operation_id)
        await emit_transfer_diagnostic(
            request.transfer_options,
            TransferDiagnostic("register", "started", context, revision_path),
        )
        async with self._heartbeat(operation_id, lease_id):
            try:
                version = await self._backend.create_from_storage(
                    storage_path,
                    revision_path,
                    operation_id,
                    request.checksum_sha256,
                    request.transfer_options,
                )
            except FileNotFoundError as exc:
                raise ArtifactNotFoundError("External bytes were not found.", operation="register") from exc
            except _StorageConflict:
                version = await self._backend.exists(revision_path)
                if (
                    version is None
                    or version.operation_id != operation_id
                    or version.checksum_sha256 != request.checksum_sha256
                ):
                    raise ConflictError(
                        "Immutable snapshot path is occupied by bytes not owned by this operation.",
                        operation="register",
                    ) from None
            except BaseException as exc:
                await emit_transfer_diagnostic(
                    request.transfer_options,
                    TransferDiagnostic(
                        "register",
                        "failed",
                        context,
                        revision_path,
                        error_type=type(exc).__name__,
                    ),
                )
                raise
            if request.size_bytes is not None and request.size_bytes != version.size_bytes:
                await self._backend.delete_owned(revision_path, operation_id, version.token)
                error = PreconditionFailedError(
                    "Registered size did not match the immutable snapshot.",
                    operation="register",
                )
                await emit_transfer_diagnostic(
                    request.transfer_options,
                    TransferDiagnostic(
                        "register",
                        "failed",
                        context,
                        revision_path,
                        error_type=type(error).__name__,
                    ),
                )
                raise error
        await emit_transfer_diagnostic(
            request.transfer_options,
            TransferDiagnostic(
                "register",
                "completed",
                context,
                revision_path,
                version.size_bytes,
                version.checksum_sha256,
            ),
        )
        await self._update_operation(
            operation_id,
            lease_id,
            ownership_verified=True,
            owned_version_token=version.token,
        )
        self._interrupt("after_object", operation_id)
        await self._update_operation(operation_id, lease_id, state="commit_ready")
        revision = ManifestRevision(
            revision_id=operation_id,
            storage_path=revision_path,
            content_revision=request.content_revision or request.checksum_sha256,
            created_at=provenance.created_at,
            operation_id=operation_id,
            ownership=ManifestOwnership.MANAGED,
            checksum_sha256=request.checksum_sha256,
            size_bytes=version.size_bytes,
            version_token=version.token,
            provenance=provenance,
        )
        return await self._commit_revision(
            CatalogOperation.REGISTER,
            operation_id,
            path,
            artifact_id,
            request.metadata,
            revision,
            provenance,
            request.expected_generation,
            request.expected_revision_id,
            lease_id,
        )
    except BaseException:
        await self._abandon_operation(operation_id, lease_id)
        raise

upload(request, context=RequestContext()) async

Create an immutable managed revision and commit it to the manifest.

Source code in src/agora_workbench/data_lake/managed.py
async def upload(
    self,
    request: UploadArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Create an immutable managed revision and commit it to the manifest."""
    return await self._upload(request, context, CatalogOperation.UPLOAD)

promote(request, context=RequestContext()) async

Explicitly copy a scratch output into durable managed storage.

Source code in src/agora_workbench/data_lake/managed.py
async def promote(
    self,
    request: PromoteOutputRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Explicitly copy a scratch output into durable managed storage."""
    if not request.session_id or not request.output_name:
        raise InvalidRequestError("Promotion requires session_id and output_name.", operation="promote")
    return await self._upload(request, context, CatalogOperation.PROMOTE)

remove(request, context=RequestContext()) async

Commit a tombstone before collecting any managed bytes.

Source code in src/agora_workbench/data_lake/managed.py
async def remove(
    self,
    request: RemoveArtifactRequest,
    context: RequestContext = RequestContext(),
) -> ManagedWriteResult:
    """Commit a tombstone before collecting any managed bytes."""
    operation_id = _validate_operation_id(request.operation_id)
    if request.reference.source_id != self.source_id:
        raise ArtifactNotFoundError("Artifact was not found.", operation="remove")
    if request.reference.revision is not None:
        raise InvalidRequestError(
            "Managed removal uses expected_revision_id rather than a catalog revision number.",
            operation="remove",
        )
    intent = self._intent(
        CatalogOperation.REMOVE,
        operation_id,
        "",
        request.reference.artifact_id,
        request,
        provenance=self._provenance(CatalogOperation.REMOVE, operation_id, context, request),
    )
    prior, lease_id = await self._begin_operation(operation_id, intent)
    if prior is not None:
        return await self._resume_remove_cleanup(prior)
    provenance = self._provenance(CatalogOperation.REMOVE, operation_id, context, request)
    try:
        self._interrupt("after_intent", operation_id)
        await self._update_operation(operation_id, lease_id, state="commit_ready")
        return await self._remove_under_lease(request, provenance, operation_id, lease_id)
    except BaseException:
        await self._abandon_operation(operation_id, lease_id)
        raise

reconcile(*, grace_seconds=300.0) async

Claim abandoned operations before recovery or owned-orphan cleanup.

Source code in src/agora_workbench/data_lake/managed.py
async def reconcile(self, *, grace_seconds: float = 300.0) -> ReconciliationReport:
    """Claim abandoned operations before recovery or owned-orphan cleanup."""
    if grace_seconds < 0:
        raise ValueError("grace_seconds must be non-negative")
    recovered: list[str] = []
    removed: list[str] = []
    deferred: list[str] = []
    active_operation_ids: set[str] = set()
    failures: dict[str, str] = {}
    for intent_path, listed_intent in await self._backend.list_json(RESERVED_OPERATIONS_PREFIX.rstrip("/")):
        operation_id = str(listed_intent.get("operation_id", ""))
        try:
            async with self._backend.serialized():
                intent, token = await self._backend.read_json(intent_path)
                if intent is None:
                    continue
                receipt, _ = await self._backend.read_json(self._receipt_path(operation_id))
                if receipt is not None:
                    result = self._result_from_mapping(receipt)
                    resumed = await self._resume_remove_cleanup(result)
                    if resumed != result:
                        recovered.append(operation_id)
                    continue
                manifest, _ = await self._load()
                artifact = self._find_operation_artifact(manifest, operation_id)
                if artifact is not None:
                    pass
                else:
                    created_at = datetime.fromisoformat(str(intent["created_at"]).replace("Z", "+00:00"))
                    age = (datetime.now(timezone.utc) - created_at).total_seconds()
                    if age < grace_seconds or (
                        intent.get("state") in {"active", "commit_ready"} and not self._lease_expired(intent)
                    ):
                        active_operation_ids.add(operation_id)
                        deferred.append(operation_id)
                        continue
                    if intent.get("state") in {"commit_ready", "fencing"}:
                        artifact, manifest, intent = await self._fence_expired_commit(
                            intent_path,
                            intent,
                            token,
                            operation_id,
                        )
                        if manifest is None or intent is None:
                            deferred.append(operation_id)
                            continue
                    elif intent.get("state") == "reconciling":
                        pass
                    else:
                        claim_id = uuid.uuid4().hex
                        claimed = {**intent, "state": "reconciling", "reconciliation_claim": claim_id}
                        try:
                            await self._backend.replace_json(intent_path, claimed, token)
                        except _StorageConflict:
                            deferred.append(operation_id)
                            continue
                        intent = claimed
                        manifest, _ = await self._load()
                        artifact = self._find_operation_artifact(manifest, operation_id)
            if artifact is not None:
                removal = next(
                    (item for item in artifact.removals if item.operation_id == operation_id),
                    None,
                )
                committed_revision = next(
                    (item for item in artifact.revisions if item.operation_id == operation_id),
                    None,
                )
                result = ManagedWriteResult(
                    operation_id,
                    self.source_id,
                    artifact.artifact_id or "",
                    (
                        removal.generation
                        if removal is not None
                        else (
                            committed_revision.committed_generation
                            if committed_revision is not None
                            and committed_revision.committed_generation is not None
                            else manifest.generation
                        )
                    ),
                    removal.revision_id
                    if removal is not None
                    else (
                        committed_revision.revision_id if committed_revision is not None else artifact.revision_id
                    ),
                    removal.storage_path
                    if removal is not None
                    else (
                        committed_revision.storage_path if committed_revision is not None else artifact.storage_path
                    ),
                    deleted=removal is not None or artifact.deleted_at is not None,
                    cleanup_pending=removal.garbage_collect if removal is not None else False,
                )
                await self._write_receipt(result)
                result = await self._resume_remove_cleanup(result)
                recovered.append(operation_id)
                continue
            owned_path = intent.get("owned_path")
            owned_version_token = intent.get("owned_version_token")
            if isinstance(owned_path, str) and (
                intent.get("ownership_verified") is not True or not isinstance(owned_version_token, str)
            ):
                owned = await self._backend.exists(owned_path)
                if owned is not None and owned.operation_id == operation_id:
                    current, token = await self._backend.read_json(intent_path)
                    if (
                        current is not None
                        and current.get("state") in {"fencing", "reconciling"}
                        and current.get("reconciliation_claim") == intent.get("reconciliation_claim")
                    ):
                        recovered_intent = {
                            **current,
                            "ownership_verified": True,
                            "owned_version_token": owned.token,
                            "checksum_sha256": owned.checksum_sha256,
                        }
                        try:
                            await self._backend.replace_json(intent_path, recovered_intent, token)
                        except _StorageConflict:
                            deferred.append(operation_id)
                            continue
                        intent = recovered_intent
                        owned_version_token = owned.token
            if (
                isinstance(owned_path, str)
                and intent.get("ownership_verified") is True
                and isinstance(owned_version_token, str)
            ):
                outcome = await self._backend.delete_owned(owned_path, operation_id, owned_version_token)
                if outcome in {DeleteOutcome.DELETED, DeleteOutcome.ABSENT}:
                    if intent.get("state") != "fencing" or await self._finish_fenced_operation(
                        intent_path,
                        intent,
                        operation_id,
                    ):
                        if intent.get("state") != "reconciling" or await self._finish_cleanup_operation(
                            intent_path,
                            intent,
                        ):
                            removed.append(operation_id)
                        else:
                            deferred.append(operation_id)
                    else:
                        deferred.append(operation_id)
                else:
                    failures[operation_id] = "Owned orphan could not be verified or removed."
            else:
                if intent.get("state") == "fencing":
                    await self._finish_fenced_operation(intent_path, intent, operation_id)
                elif intent.get("state") == "reconciling":
                    await self._finish_cleanup_operation(intent_path, intent)
                deferred.append(operation_id)
        except Exception as exc:
            failures[operation_id] = type(exc).__name__
    if failures:
        raise ReconciliationError(
            "One or more interrupted operations could not be reconciled.",
            failures=failures,
            operation="reconcile",
        )
    async with self._backend.serialized():
        removed_staging, deferred_staging = await self._backend.reconcile_staging(
            frozenset(active_operation_ids),
            grace_seconds,
        )
    return ReconciliationReport(
        tuple(recovered),
        tuple(removed),
        tuple(deferred),
        removed_staging,
        deferred_staging,
        failures,
    )

AuthorizedManagedCatalogWriter(writer, authorizer)

Authorize every mutation before metadata or byte side effects.

Source code in src/agora_workbench/data_lake/managed.py
def __init__(self, writer: ManagedCatalogWriter, authorizer: CatalogAuthorizer) -> None:
    self._writer = writer
    self._authorizer = authorizer

capabilities(context) async

Return write operations allowed for this session and source.

Source code in src/agora_workbench/data_lake/managed.py
async def capabilities(self, context: RequestContext) -> tuple[SourceCapabilities, ...]:
    """Return write operations allowed for this session and source."""
    allowed = {
        operation
        for operation in WRITE_OPERATIONS
        if await self._authorizer.authorize(
            CatalogAuthorizationRequest(operation, self._writer.source_id),
            context,
        )
    }
    return (SourceCapabilities(self._writer.source_id, frozenset(allowed)),) if allowed else ()

managed_writer_extension_factory(writer)

Adapt a managed writer to CatalogIntegration.capability_extension_factory.

Source code in src/agora_workbench/data_lake/managed.py
def managed_writer_extension_factory(writer: ManagedCatalogWriter):
    """Adapt a managed writer to ``CatalogIntegration.capability_extension_factory``."""

    def create(_session: object, catalog: object, _context: RequestContext) -> AuthorizedManagedCatalogWriter:
        authorizer = getattr(catalog, "authorizer", None)
        if authorizer is None:
            raise TypeError("Managed writer integration requires an authorized catalog.")
        return AuthorizedManagedCatalogWriter(writer, authorizer)

    return create

Resolvers

agora_workbench.data_lake.resolvers

Public resolver protocol and compatibility exports.

SearchIndexArtifactResolver(credential, endpoint=None, index_name=None, credential_init_error=None)

Resolves artifact IDs against an Azure AI Search blob-details index.

The artifact ID is used directly as the document key; the blob URL is read from the document's metadata_storage_path field. Resolved URLs are cached so repeated resolutions of the same artifact skip the round-trip.

The credential is borrowed, not owned: :meth:aclose closes the search client but never the credential, which remains the caller's to close.

Initialize the resolver.

Parameters:

Name Type Description Default
credential AsyncTokenCredential | None

Async token credential used to query the search index. When None, resolution is unavailable.

required
endpoint str | None

Azure AI Search endpoint. When falsy, resolution is unavailable and the resolver reports that the endpoint is unset.

None
index_name str | None

Name of the blob-details index.

None
credential_init_error str | None

Optional "Type: message" string describing a credential failure that happened before this resolver was constructed, so the deferred error can be surfaced at resolve time rather than being replaced by a generic message.

None
Source code in src/agora_workbench/code_execution/data_access/artifact_resolvers.py
def __init__(
    self,
    credential: "AsyncTokenCredential | None",
    endpoint: str | None = None,
    index_name: str | None = None,
    credential_init_error: str | None = None,
):
    """
    Initialize the resolver.

    Args:
        credential: Async token credential used to query the search index.
            When ``None``, resolution is unavailable.
        endpoint: Azure AI Search endpoint. When falsy, resolution is
            unavailable and the resolver reports that the endpoint is unset.
        index_name: Name of the blob-details index.
        credential_init_error: Optional ``"Type: message"`` string describing
            a credential failure that happened before this resolver was
            constructed, so the deferred error can be surfaced at resolve
            time rather than being replaced by a generic message.
    """
    self._endpoint = endpoint
    self._index_name = index_name
    self._credential_init_error = credential_init_error
    self._url_cache: dict[str, str] = {}  # Maps artifact_id -> resolved blob URL
    self._search_client: Any = None
    self._closed = False

    if not endpoint:
        LOGGER.info("No Azure Search endpoint configured; blob artifact ID resolution is disabled")
        return

    if credential is None:
        LOGGER.warning("Azure Search endpoint configured, but no Azure credential is available")
        return

    try:
        from azure.search.documents.aio import SearchClient

        self._search_client = SearchClient(
            endpoint=endpoint,
            index_name=index_name or "",
            credential=credential,
        )
        LOGGER.info(f"Initialized blob-details search client: {endpoint}/{index_name}")
    except (ImportError, RuntimeError, TypeError, ValueError) as e:
        self._credential_init_error = f"{type(e).__name__}: {e}"
        LOGGER.warning(f"Failed to initialize Azure data access components: {e}")

unavailable_reason property

Reason blob artifact resolution cannot run, or None when ready.

from_env(credential, credential_init_error=None) classmethod

Build a resolver from DATA_LAKE_SEARCH_ENDPOINT / DATA_LAKE_BLOB_DETAILS_INDEX.

The index name is only read when an endpoint is configured, so an endpoint-less deployment is distinguishable from one whose index name was explicitly set to the empty string.

Source code in src/agora_workbench/code_execution/data_access/artifact_resolvers.py
@classmethod
def from_env(
    cls,
    credential: "AsyncTokenCredential | None",
    credential_init_error: str | None = None,
) -> "SearchIndexArtifactResolver":
    """
    Build a resolver from ``DATA_LAKE_SEARCH_ENDPOINT`` / ``DATA_LAKE_BLOB_DETAILS_INDEX``.

    The index name is only read when an endpoint is configured, so an
    endpoint-less deployment is distinguishable from one whose index name
    was explicitly set to the empty string.
    """
    endpoint = os.getenv("DATA_LAKE_SEARCH_ENDPOINT")
    index_name = os.getenv("DATA_LAKE_BLOB_DETAILS_INDEX", DEFAULT_BLOB_DETAILS_INDEX) if endpoint else None
    return cls(
        credential=credential,
        endpoint=endpoint,
        index_name=index_name,
        credential_init_error=credential_init_error,
    )

resolve(artifact_id) async

Retrieve a blob storage URL for artifact_id from the blob-details index.

Parameters:

Name Type Description Default
artifact_id str

Base64-encoded artifact identifier from the index.

required

Returns:

Type Description
str

The blob storage URL (e.g. https://account.blob.core.windows.net/container/path).

Raises:

Type Description
ValueError

If resolution is unavailable, the artifact is not found, or the stored path is not a valid storage URL.

Source code in src/agora_workbench/code_execution/data_access/artifact_resolvers.py
async def resolve(self, artifact_id: str) -> str:
    """
    Retrieve a blob storage URL for *artifact_id* from the blob-details index.

    Args:
        artifact_id: Base64-encoded artifact identifier from the index.

    Returns:
        The blob storage URL (e.g. ``https://account.blob.core.windows.net/container/path``).

    Raises:
        ValueError: If resolution is unavailable, the artifact is not found,
            or the stored path is not a valid storage URL.
    """
    if artifact_id in self._url_cache:
        LOGGER.debug(f"URL cache hit for artifact {artifact_id[:40]}...")
        return self._url_cache[artifact_id]

    if self._search_client is None:
        raise ValueError(self.unavailable_reason)

    try:
        # Query the blob-details index using artifact_id as the document key
        # The artifact_id field is the unique key in the blob-details index
        result = await self._search_client.get_document(key=artifact_id)

        if not result:
            raise ValueError(f"Artifact not found in blob-details index: {artifact_id}")

        # Extract metadata_storage_path and strip any trailing whitespace
        blob_url = result.get("metadata_storage_path", "").strip()

        if not blob_url:
            raise ValueError(f"Artifact {artifact_id} has no metadata_storage_path in blob-details index")

        if not blob_url.startswith(("https://", "abfss://", "az://")):
            raise ValueError(f"Retrieved storage path is not a valid URL: {blob_url!r}")

        LOGGER.info(f"Retrieved blob URL for artifact {artifact_id[:40]}...")
        self._url_cache[artifact_id] = blob_url
        return blob_url

    except Exception as e:
        LOGGER.error(f"Failed to retrieve blob URL for artifact {artifact_id}: {e}")
        raise ValueError(f"Failed to resolve blob artifact {artifact_id}: {e}") from e

aclose() async

Clear the URL cache and close the search client (never the credential).

Idempotent: the client reference is dropped so a second call is a no-op and :attr:unavailable_reason stops reporting readiness after shutdown.

Source code in src/agora_workbench/code_execution/data_access/artifact_resolvers.py
async def aclose(self) -> None:
    """Clear the URL cache and close the search client (never the credential).

    Idempotent: the client reference is dropped so a second call is a no-op
    and :attr:`unavailable_reason` stops reporting readiness after shutdown.
    """
    self._closed = True
    self._url_cache.clear()

    search_client, self._search_client = self._search_client, None
    if search_client is not None:
        try:
            await search_client.close()
        except Exception as e:
            LOGGER.debug(f"Error closing search client: {e}")

ArtifactResolver

Bases: Protocol

Resolve an opaque artifact ID to a fetchable storage location.

Implementations may define async def aclose(self) -> None for cleanup, but it is intentionally not required by the protocol.

unavailable_reason property

Return an operator-facing unavailability reason, or None when ready.

resolve(artifact_id) async

Resolve an artifact ID to a qualified name or URL.

Source code in src/agora_workbench/data_lake/protocols.py
async def resolve(self, artifact_id: str) -> str:
    """Resolve an artifact ID to a qualified name or URL."""
    ...

CatalogArtifactResolver(provider, source_id, context=None)

Legacy ArtifactResolver adapter for one source of a CatalogProvider.

Source code in src/agora_workbench/data_lake/providers.py
def __init__(
    self,
    provider: CatalogProvider,
    source_id: str,
    context: RequestContext | None = None,
):
    self._provider = provider
    self._source_id = source_id
    self._context = context or RequestContext()
    self._closed = False

Fetching and publishing

agora_workbench.data_lake.execution

Compatibility exports for execution-integrated data-lake implementations.

MsalCacheCredential(cache_path=None, username=None, authority='https://login.microsoftonline.com/organizations')

Async credential that reads Azure CLI's MSAL token cache directly.

Acquires tokens silently using the refresh token stored by az login. Does not require the az binary — only the msal Python library and a readable msal_token_cache.json.

Implements the AsyncTokenCredential protocol from azure.core.credentials_async so it can be used with BlobPublisher, BlobFetcher, BlobServiceClient, etc.

The cache file is re-read on every get_token() call so the credential picks up tokens refreshed by the host's az CLI without requiring a server restart.

Parameters:

Name Type Description Default
cache_path str | Path | None

Path to the MSAL token cache file. Defaults to ~/.azure/msal_token_cache.json.

None
username str | None

Optional username filter. When multiple accounts exist in the cache, only the matching account is used. If None, the first account is selected.

None
authority str

Azure AD authority URL. Defaults to the organizations endpoint (multi-tenant).

'https://login.microsoftonline.com/organizations'
Source code in src/agora_workbench/code_execution/data_access/credentials.py
def __init__(
    self,
    cache_path: str | Path | None = None,
    username: str | None = None,
    authority: str = "https://login.microsoftonline.com/organizations",
):
    self._cache_path = Path(cache_path) if cache_path else _DEFAULT_CACHE_PATH
    self._username = username
    self._authority = authority

get_token(*scopes, **kwargs) async

Acquire a token for the requested scopes from the MSAL cache.

Parameters:

Name Type Description Default
*scopes str

One or more Azure resource scopes (e.g. "https://storage.azure.com/.default").

()

Returns:

Type Description
'AccessToken'

An AccessToken with the token string and expiry.

Raises:

Type Description
CredentialUnavailableError

If the cache file is missing, no accounts are found, or silent acquisition fails. This allows ChainedTokenCredential to fall through to the next credential in the chain.

Source code in src/agora_workbench/code_execution/data_access/credentials.py
async def get_token(self, *scopes: str, **kwargs: Any) -> "AccessToken":
    """Acquire a token for the requested scopes from the MSAL cache.

    Args:
        *scopes: One or more Azure resource scopes
            (e.g. ``"https://storage.azure.com/.default"``).

    Returns:
        An ``AccessToken`` with the token string and expiry.

    Raises:
        CredentialUnavailableError: If the cache file is missing,
            no accounts are found, or silent acquisition fails.
            This allows ``ChainedTokenCredential`` to fall through
            to the next credential in the chain.
    """
    import msal

    try:
        from azure.core.credentials import AccessToken
        from azure.identity import CredentialUnavailableError
    except ImportError as exc:
        raise RuntimeError("Azure storage credentials require the 'agora-workbench[azure]' extra.") from exc

    if not self._cache_path.is_file():
        raise CredentialUnavailableError(message=f"MSAL token cache not found at {self._cache_path}")

    cache = msal.SerializableTokenCache()
    cache.deserialize(self._cache_path.read_text())

    app = msal.PublicClientApplication(
        _AZURE_CLI_CLIENT_ID,
        authority=self._authority,
        token_cache=cache,
    )

    accounts = app.get_accounts(username=self._username)
    if not accounts:
        raise CredentialUnavailableError(message="No accounts found in MSAL token cache")

    account = accounts[0]
    scope_list = list(scopes)

    # Try cached access token first, then force a refresh via RT.
    result = app.acquire_token_silent(scope_list, account=account)
    if not result or "access_token" not in result:
        result = app.acquire_token_silent(scope_list, account=account, force_refresh=True)

    if not result or "access_token" not in result:
        error_desc = (
            result.get("error_description", "unknown error") if result else "acquire_token_silent returned None"
        )
        raise CredentialUnavailableError(message=f"MSAL silent token acquisition failed: {error_desc}")

    LOGGER.debug(
        "MsalCacheCredential: acquired token for %s (account=%s)",
        scopes[0][:50],
        account.get("username", "?"),
    )
    expires_on = int(time.time()) + int(result["expires_in"])
    return AccessToken(result["access_token"], expires_on)

close() async

No-op — no persistent connections to release.

Source code in src/agora_workbench/code_execution/data_access/credentials.py
async def close(self) -> None:
    """No-op — no persistent connections to release."""

AssetFetcher(credential=None)

Bases: ABC

Base class for asset fetchers.

fetch is a full-memory convenience. Use fetch_to_file (or fetch_to_file_result when diagnostics are needed) for bounded-memory transfer.

Initialize fetcher with an optional async token credential.

Parameters:

Name Type Description Default
credential AsyncTokenCredential | None

An AsyncTokenCredential that provides tokens for downstream Azure resources (e.g. ManagedIdentityCredential). May be None for fetchers that don't require credentials (e.g. local filesystem).

None
Source code in src/agora_workbench/code_execution/data_access/fetchers.py
def __init__(self, credential: "AsyncTokenCredential | None" = None):
    """
    Initialize fetcher with an optional async token credential.

    Args:
        credential: An ``AsyncTokenCredential`` that provides tokens for
                   downstream Azure resources (e.g. ManagedIdentityCredential).
                   May be ``None`` for fetchers that don't require credentials
                   (e.g. local filesystem).
    """
    self.credential = credential

fetch(qualified_name) abstractmethod async

Fetch asset data into memory.

Parameters:

Name Type Description Default
qualified_name str

DataLake asset qualified name

required

Returns:

Type Description
Any

Raw data (bytes, DataFrame, etc.)

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
@abstractmethod
async def fetch(self, qualified_name: str) -> Any:
    """
    Fetch asset data into memory.

    Args:
        qualified_name: DataLake asset qualified name

    Returns:
        Raw data (bytes, DataFrame, etc.)
    """
    pass

fetch_to_file(qualified_name, dest_path, *, options=None, context=None) abstractmethod async

Fetch asset data and stream directly to a file.

Streams data to disk to avoid loading large assets into memory.

Parameters:

Name Type Description Default
qualified_name str

DataLake asset qualified name

required
dest_path Any

Destination file path (Path object or string)

required

Returns:

Type Description
int

Number of bytes written

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
@abstractmethod
async def fetch_to_file(
    self,
    qualified_name: str,
    dest_path: Any,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> int:
    """
    Fetch asset data and stream directly to a file.

    Streams data to disk to avoid loading large assets into memory.

    Args:
        qualified_name: DataLake asset qualified name
        dest_path: Destination file path (Path object or string)

    Returns:
        Number of bytes written
    """
    pass

fetch_to_file_result(qualified_name, dest_path, *, options=None, context=None) async

Return detailed transfer diagnostics when the fetcher supports them.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def fetch_to_file_result(
    self,
    qualified_name: str,
    dest_path: Any,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> TransferResult:
    """Return detailed transfer diagnostics when the fetcher supports them."""
    raise UnsupportedOperationError(
        f"{type(self).__name__} does not implement bounded streaming.",
        operation="download",
    )

can_handle(qualified_name) abstractmethod

Check if this fetcher can handle the given qualified name.

Parameters:

Name Type Description Default
qualified_name str

DataLake asset qualified name

required

Returns:

Type Description
bool

True if this fetcher supports the asset type

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
@abstractmethod
def can_handle(self, qualified_name: str) -> bool:
    """
    Check if this fetcher can handle the given qualified name.

    Args:
        qualified_name: DataLake asset qualified name

    Returns:
        True if this fetcher supports the asset type
    """
    pass

BlobFetcher(credential=None, *, allowed_locations=None, account_endpoints=None)

Bases: AssetFetcher

Fetcher for Azure Blob Storage / ADLS Gen2 assets.

Maintains a per-account client cache to amortize TCP/TLS handshake and token acquisition costs across multiple fetches.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
def __init__(
    self,
    credential: "AsyncTokenCredential | None" = None,
    *,
    allowed_locations: list[str | AzureBlobScope] | None = None,
    account_endpoints: Mapping[str, str] | None = None,
):
    super().__init__(credential=credential)
    # Cache of account_url -> BlobServiceClient for connection reuse
    self._clients: dict[str, "AzureBlobServiceClient"] = {}
    self._allowed_scopes = tuple(
        value if isinstance(value, AzureBlobScope) else AzureBlobScope.from_uri(value)
        for value in (allowed_locations or [])
    )
    self._account_endpoints = {
        account.lower(): self._validate_account_endpoint(account, endpoint)
        for account, endpoint in (account_endpoints or {}).items()
    }

close() async

Close all cached blob service clients.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def close(self) -> None:
    """Close all cached blob service clients."""
    for client in self._clients.values():
        await client.close()
    self._clients.clear()

can_handle(qualified_name)

Check if this is a blob storage URL.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
def can_handle(self, qualified_name: str) -> bool:
    """Check if this is a blob storage URL."""
    if qualified_name.startswith("abfss://"):
        return True

    if qualified_name.startswith("az://"):
        # az://account/container/blob — the scheme emitted by the catalog
        # indexer. Structural validation happens in _parse_blob_url.
        return True

    if qualified_name.startswith("https://"):
        # Properly parse URL and check hostname to avoid substring injection
        try:
            parsed = urlsplit(qualified_name)
            hostname = parsed.netloc.lower()
            return hostname.endswith(".blob.core.windows.net") or hostname.endswith(".dfs.core.windows.net")
        except Exception:
            return False

    return False

fetch(qualified_name) async

Fetch data from Azure Blob Storage.

Supports: - abfss://container@storage.dfs.core.windows.net/path/to/file - az://account/container/path/to/file - https://storage.blob.core.windows.net/container/path/to/file

Parameters:

Name Type Description Default
qualified_name str

Blob URL

required

Returns:

Type Description
bytes

Raw bytes of the file

Raises:

Type Description
ClientAuthenticationError

If access is denied

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def fetch(self, qualified_name: str) -> bytes:
    """
    Fetch data from Azure Blob Storage.

    Supports:
    - abfss://container@storage.dfs.core.windows.net/path/to/file
    - az://account/container/path/to/file
    - https://storage.blob.core.windows.net/container/path/to/file

    Args:
        qualified_name: Blob URL

    Returns:
        Raw bytes of the file

    Raises:
        azure.core.exceptions.ClientAuthenticationError: If access is denied
    """
    # Parse the URL first to sanitize for logging (removes query params like SAS tokens)
    storage_account, container, blob_path = self._parse_blob_url(qualified_name)
    self._require_allowed(storage_account, container, blob_path)
    sanitized_url = f"{storage_account}/{container}/{blob_path}"
    LOGGER.info(f"Fetching blob asset: {sanitized_url}")

    # Get or create authenticated client (connection reuse)
    account_url = self._account_url(storage_account)
    client = self._get_client(account_url)
    blob_client = client.get_blob_client(container=container, blob=blob_path)

    # Download blob data with parallel range requests
    stream = await blob_client.download_blob(max_concurrency=_BLOB_MAX_CONCURRENCY)
    data = await stream.readall()

    LOGGER.info(f"Successfully fetched {len(data)} bytes from {sanitized_url}")
    return data

fetch_to_file(qualified_name, dest_path, *, options=None, context=None) async

Fetch blob data and stream directly to a file.

Streams data in chunks to avoid loading large files into memory. Uses parallel range requests for improved throughput on large blobs.

Parameters:

Name Type Description Default
qualified_name str

Blob URL

required
dest_path Any

Destination file path

required

Returns:

Type Description
int

Number of bytes written

Raises:

Type Description
ClientAuthenticationError

If access is denied

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def fetch_to_file(
    self,
    qualified_name: str,
    dest_path: Any,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> int:
    """
    Fetch blob data and stream directly to a file.

    Streams data in chunks to avoid loading large files into memory.
    Uses parallel range requests for improved throughput on large blobs.

    Args:
        qualified_name: Blob URL
        dest_path: Destination file path

    Returns:
        Number of bytes written

    Raises:
        azure.core.exceptions.ClientAuthenticationError: If access is denied
    """
    result = await self.fetch_to_file_result(
        qualified_name,
        dest_path,
        options=options,
        context=context,
    )
    return result.bytes_transferred

fetch_to_file_result(qualified_name, dest_path, *, options=None, context=None) async

Download one Blob object with bounded memory and atomic local commit.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def fetch_to_file_result(
    self,
    qualified_name: str,
    dest_path: Any,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> TransferResult:
    """Download one Blob object with bounded memory and atomic local commit."""
    return await self._fetch_to_file_result(
        qualified_name,
        dest_path,
        options=options,
        context=context,
        allow_reserved=False,
    )

LocalFileFetcher(allowed_roots=None)

Bases: AssetFetcher

Fetcher for local filesystem paths.

Handles absolute paths, relative paths, and file:// URIs. No credentials are required.

Security

An allowed_roots list restricts which directories the fetcher may read from. Every resolved path is checked against these roots before any I/O occurs. If allowed_roots is empty, all paths are permitted (use only inside a sandboxed container).

Initialize the local file fetcher.

Parameters:

Name Type Description Default
allowed_roots list[str] | None

Optional list of directory paths that the fetcher is allowed to read from. Paths are resolved to absolute form. If None or empty, all paths are permitted.

None
Source code in src/agora_workbench/code_execution/data_access/fetchers.py
def __init__(self, allowed_roots: list[str] | None = None):
    """
    Initialize the local file fetcher.

    Args:
        allowed_roots: Optional list of directory paths that the fetcher
            is allowed to read from.  Paths are resolved to absolute form.
            If ``None`` or empty, all paths are permitted.
    """
    super().__init__(credential=None)
    self._allowed_roots: list[Path] = [Path(r).resolve() for r in (allowed_roots or [])]
    self._allowed_root_fds: list[int] = []
    self._allowed_root_identities: list[tuple[int, int]] = []
    if os.name == "posix":
        try:
            for root in self._allowed_roots:
                descriptor = os.open(
                    root,
                    os.O_RDONLY | getattr(os, "O_DIRECTORY", 0) | getattr(os, "O_NOFOLLOW", 0),
                )
                stat_result = os.fstat(descriptor)
                self._allowed_root_fds.append(descriptor)
                self._allowed_root_identities.append((stat_result.st_dev, stat_result.st_ino))
        except BaseException:
            for descriptor in self._allowed_root_fds:
                os.close(descriptor)
            self._allowed_root_fds.clear()
            self._allowed_root_identities.clear()
            raise

close() async

Close retained trusted allowed-root descriptors.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def close(self) -> None:
    """Close retained trusted allowed-root descriptors."""
    self._close_root_descriptors()

__del__()

Defensively release retained descriptors when explicit cleanup is missed.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
def __del__(self) -> None:
    """Defensively release retained descriptors when explicit cleanup is missed."""
    with suppress(Exception):
        self._close_root_descriptors()

can_handle(qualified_name)

Check if this is a local filesystem path.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
def can_handle(self, qualified_name: str) -> bool:
    """Check if this is a local filesystem path."""
    return (
        qualified_name.startswith("/")
        or qualified_name.startswith("./")
        or qualified_name.startswith("../")
        or qualified_name.startswith("file://")
    )

fetch(qualified_name) async

Read a local file into memory.

Parameters:

Name Type Description Default
qualified_name str

Local file path or file:// URI.

required

Returns:

Type Description
bytes

Raw bytes of the file.

Raises:

Type Description
FileNotFoundError

If the file does not exist.

PermissionError

If the resolved path is outside allowed_roots.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def fetch(self, qualified_name: str) -> bytes:
    """
    Read a local file into memory.

    Args:
        qualified_name: Local file path or ``file://`` URI.

    Returns:
        Raw bytes of the file.

    Raises:
        FileNotFoundError: If the file does not exist.
        PermissionError: If the resolved path is outside *allowed_roots*.
    """
    path, descriptor = self._open_checked(qualified_name)
    LOGGER.info(f"Reading local file: {path}")
    with os.fdopen(descriptor, "rb", closefd=True) as input_file:
        data = input_file.read()
    LOGGER.info(f"Read {len(data)} bytes from {path}")
    return data

fetch_to_file(qualified_name, dest_path, *, options=None, context=None) async

Copy a local file to dest_path.

Parameters:

Name Type Description Default
qualified_name str

Local file path or file:// URI.

required
dest_path Any

Destination file path.

required

Returns:

Type Description
int

Number of bytes written.

Raises:

Type Description
FileNotFoundError

If the source file does not exist.

PermissionError

If the resolved path is outside allowed_roots.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def fetch_to_file(
    self,
    qualified_name: str,
    dest_path: Any,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> int:
    """
    Copy a local file to *dest_path*.

    Args:
        qualified_name: Local file path or ``file://`` URI.
        dest_path: Destination file path.

    Returns:
        Number of bytes written.

    Raises:
        FileNotFoundError: If the source file does not exist.
        PermissionError: If the resolved path is outside *allowed_roots*.
    """
    result = await self.fetch_to_file_result(
        qualified_name,
        dest_path,
        options=options,
        context=context,
    )
    return result.bytes_transferred

fetch_to_file_result(qualified_name, dest_path, *, options=None, context=None) async

Copy a local file through a descriptor that cannot follow symlinks.

Source code in src/agora_workbench/code_execution/data_access/fetchers.py
async def fetch_to_file_result(
    self,
    qualified_name: str,
    dest_path: Any,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> TransferResult:
    """Copy a local file through a descriptor that cannot follow symlinks."""
    return await self._fetch_to_file_result(
        qualified_name,
        dest_path,
        options=options,
        context=context,
        allow_reserved=False,
    )

DataLakeDataManager(allowed_local_roots=None, extra_fetchers=None, credential=None, artifact_resolver=None, transfer_options=None, credential_ownership=ResourceOwnership.BORROWED, _catalog_integration_token=None)

Manages data asset retrieval and caching for code execution sessions.

Fetches data assets from DataLake-cataloged sources and caches them to disk in their original format. Tools receive Path objects and handle loading the data as needed.

Authentication uses a credential chain (create_storage_credential): the mounted az login MSAL cache for local development, falling back to managed identity in production, for downstream Azure resources (Storage, AI Search).

Blob artifact IDs are resolved through a pluggable ArtifactResolver, defaulting to Azure AI Search; see artifact_resolvers.py.

Supports: - Azure Blob Storage (abfss://, az://, https://) - Local filesystem (absolute paths, relative paths, file:// URIs)

Initialize the data manager.

Uses a credential chain (az login MSAL cache locally, managed identity in production) for Azure Blob access unless a credential is provided. Blob URL fetching is available whenever a credential can be initialized. Blob artifact ID resolution is delegated to an ArtifactResolver; the default one queries Azure AI Search and so additionally requires DATA_LAKE_SEARCH_ENDPOINT.

Parameters:

Name Type Description Default
allowed_local_roots list[str] | None

Optional list of directory paths the local file fetcher is allowed to read from. If None or empty, all paths are permitted (suitable for sandboxed containers).

None
extra_fetchers list[AssetFetcher] | None

Optional list of additional AssetFetcher instances to register. These are checked before the built-in fetchers (local file, blob), allowing custom fetchers to override default handling for specific URL schemes or patterns.

None
credential AsyncTokenCredential | None

Optional async token credential to use for Azure Blob Storage and Azure AI Search access. When omitted, the manager creates the same storage credential chain as before, resolving AZURE_CLIENT_ID for user-assigned managed identity binding.

None
artifact_resolver ArtifactResolver | None

Optional resolver turning <blob>id</blob> identifiers into fetchable URLs, for deployments whose catalog is not an Azure AI Search index. When omitted, a SearchIndexArtifactResolver is built from the environment, preserving existing behavior. A supplied resolver is used as-is and is not given the manager's credential, so it must arrange its own authentication.

None
transfer_options TransferOptions | None

Default bounded-transfer policy used when a call does not provide an explicit override.

None
credential_ownership ResourceOwnership

Whether the manager closes a supplied credential. Existing callers retain borrowed semantics by default; integrations that create one credential per session can explicitly transfer ownership.

BORROWED

Raises:

Type Description
TypeError

If artifact_resolver does not implement the ArtifactResolver protocol.

Source code in src/agora_workbench/code_execution/data_access/manager.py
def __init__(
    self,
    allowed_local_roots: list[str] | None = None,
    extra_fetchers: list[AssetFetcher] | None = None,
    credential: "AsyncTokenCredential | None" = None,
    artifact_resolver: ArtifactResolver | None = None,
    transfer_options: TransferOptions | None = None,
    credential_ownership: ResourceOwnership = ResourceOwnership.BORROWED,
    _catalog_integration_token: object | None = None,
):
    """
    Initialize the data manager.

    Uses a credential chain (``az login`` MSAL cache locally, managed
    identity in production) for Azure Blob access unless a credential is
    provided. Blob URL fetching is available whenever a credential can be
    initialized. Blob artifact ID resolution is delegated to an
    ``ArtifactResolver``; the default one queries Azure AI Search and so
    additionally requires ``DATA_LAKE_SEARCH_ENDPOINT``.

    Args:
        allowed_local_roots: Optional list of directory paths the local
            file fetcher is allowed to read from. If ``None`` or empty,
            all paths are permitted (suitable for sandboxed containers).
        extra_fetchers: Optional list of additional ``AssetFetcher``
            instances to register. These are checked *before* the
            built-in fetchers (local file, blob), allowing custom
            fetchers to override default handling for specific URL
            schemes or patterns.
        credential: Optional async token credential to use for Azure Blob
            Storage and Azure AI Search access. When omitted, the manager
            creates the same storage credential chain as before, resolving
            ``AZURE_CLIENT_ID`` for user-assigned managed identity binding.
        artifact_resolver: Optional resolver turning ``<blob>id</blob>``
            identifiers into fetchable URLs, for deployments whose catalog
            is not an Azure AI Search index. When omitted, a
            ``SearchIndexArtifactResolver`` is built from the environment,
            preserving existing behavior. A supplied resolver is used as-is
            and is *not* given the manager's credential, so it must arrange
            its own authentication.
        transfer_options: Default bounded-transfer policy used when a call
            does not provide an explicit override.
        credential_ownership: Whether the manager closes a supplied
            credential. Existing callers retain borrowed semantics by
            default; integrations that create one credential per session
            can explicitly transfer ownership.

    Raises:
        TypeError: If ``artifact_resolver`` does not implement the
            ``ArtifactResolver`` protocol.
    """
    if artifact_resolver is not None:
        _validate_artifact_resolver(artifact_resolver)

    self._cache_dir = Path(tempfile.mkdtemp(prefix="data_lake_cache_"))
    self._cache_index = {}  # Maps artifact_id -> cache file path
    self._cache_generation = 0
    self._full_cache_generation = 0
    self._transfer_options = transfer_options or TransferOptions()
    self._catalog_managed_revision_access = _catalog_integration_token is _CATALOG_INTEGRATION_MANAGER_TOKEN

    self._credential_init_error: str | None = None
    self._credential: "AsyncTokenCredential | None" = None
    self._owns_credential = credential is None or credential_ownership is ResourceOwnership.OWNED

    # Initialize fetchers — custom fetchers take priority over built-ins
    self._fetchers: list[AssetFetcher] = list(extra_fetchers or [])
    self._fetchers.append(LocalFileFetcher(allowed_roots=allowed_local_roots))

    try:
        if credential is not None:
            self._credential = credential
        else:
            mi_client_id = (os.getenv("AZURE_CLIENT_ID") or "").strip() or None
            # MSAL cache (mounted az login) locally, managed identity in
            # production. Pass the AZURE_CLIENT_ID-resolved id through so
            # prod keeps binding to the same user-assigned identity.
            self._credential = create_storage_credential(client_id=mi_client_id)
    except ImportError as e:
        self._credential_init_error = f"{type(e).__name__}: {e}"
        LOGGER.info(
            "Azure SDK is not installed; cloud data access is disabled. "
            "Install the 'agora-workbench[azure]' extra to enable it."
        )
    except (RuntimeError, TypeError, ValueError) as e:
        self._credential_init_error = f"{type(e).__name__}: {e}"
        LOGGER.warning(f"Failed to initialize Azure storage credential: {e}")

    if self._credential is not None:
        self._fetchers.append(BlobFetcher(credential=self._credential))

    if artifact_resolver is not None:
        self._artifact_resolver: ArtifactResolver = artifact_resolver
    else:
        # The deferred credential error is handed over so the resolver can
        # explain *why* it is unavailable rather than blaming configuration.
        self._artifact_resolver = SearchIndexArtifactResolver.from_env(
            credential=self._credential,
            credential_init_error=self._credential_init_error,
        )
    self._catalog_artifact_resolver: ArtifactResolver | None = None

bind_catalog_resolver(resolver, *, fetchers=())

Attach a caller-scoped catalog resolver and its backend fetchers.

Catalog references are routed to resolver while legacy blob artifact IDs continue to use the manager's existing resolver. Supplied fetchers become session-owned and take priority over previously configured fetchers. Ownership transfers only after this method returns successfully.

Source code in src/agora_workbench/code_execution/data_access/manager.py
def bind_catalog_resolver(
    self,
    resolver: ArtifactResolver,
    *,
    fetchers: Sequence[AssetFetcher] = (),
) -> None:
    """Attach a caller-scoped catalog resolver and its backend fetchers.

    Catalog references are routed to *resolver* while legacy blob artifact
    IDs continue to use the manager's existing resolver. Supplied fetchers
    become session-owned and take priority over previously configured
    fetchers. Ownership transfers only after this method returns
    successfully.
    """
    _validate_artifact_resolver(resolver)
    if self._catalog_artifact_resolver is not None and self._catalog_artifact_resolver is not resolver:
        raise RuntimeError("A catalog resolver is already bound to this data manager.")
    existing_fetchers = {id(fetcher) for fetcher in self._fetchers}
    if len({id(fetcher) for fetcher in fetchers}) != len(fetchers):
        raise ValueError("Catalog fetchers must be distinct instances.")
    if any(id(fetcher) in existing_fetchers for fetcher in fetchers):
        raise ValueError("Catalog fetchers are already registered with this data manager.")
    self._catalog_artifact_resolver = resolver
    self._fetchers[0:0] = list(fetchers)

get_cache_path(qualified_name, *, context=None, transfer_options=None) async

Get the filesystem path where the asset is cached.

Ensures the asset is fetched and cached to disk in its original format, then returns the path. This allows kernel subprocesses to load the asset directly.

Parameters:

Name Type Description Default
qualified_name AssetId

Type-tagged artifact format base64_id Examples: id, id

required

Returns:

Type Description
Path

Path to the cached file

Raises:

Type Description
ValueError

If not in tagged format, unsupported type, or artifact not found

Source code in src/agora_workbench/code_execution/data_access/manager.py
async def get_cache_path(
    self,
    qualified_name: "AssetId",
    *,
    context: RequestContext | None = None,
    transfer_options: TransferOptions | None = None,
) -> Path:
    """
    Get the filesystem path where the asset is cached.

    Ensures the asset is fetched and cached to disk in its original format,
    then returns the path. This allows kernel subprocesses to load the
    asset directly.

    Args:
        qualified_name: Type-tagged artifact format <type>base64_id</type>
                      Examples: <blob>id</blob>, <sql>id</sql>

    Returns:
        Path to the cached file

    Raises:
        ValueError: If not in tagged format, unsupported type, or artifact not found
    """
    # Extract artifact type and ID from tagged format
    artifact_match = re.match(r"^<(\w+)>([^<>]+)</\1>$", qualified_name.strip())
    if not artifact_match:
        # Fallback: accept unclosed tags like "<blob>id" (LLM sometimes omits closing tag)
        artifact_match = re.match(r"^<(\w+)>([^<>]+)$", qualified_name.strip())
    if not artifact_match:
        raise ValueError(
            "Invalid artifact format - expected <type>id</type>. "
            f"{self._asset_tag_guidance()} {agent_guidance.DISCOVER_DATA}"
        )

    artifact_type = artifact_match.group(1)
    artifact_id = artifact_match.group(2)
    cache_generation = self._cache_generation
    full_cache_generation = self._full_cache_generation
    generation_scoped = artifact_type == "blob" and artifact_id.startswith("catalog-v1:")

    def cache_was_invalidated() -> bool:
        return full_cache_generation != self._full_cache_generation or (
            generation_scoped and cache_generation != self._cache_generation
        )

    display_id = sanitize_uri_for_display(artifact_id) if "://" in artifact_id else artifact_id
    LOGGER.info("Resolving %s artifact: %s", artifact_type, display_id)

    # Check if already cached (use artifact_id as cache key). A concurrent
    # authorization refresh invalidates catalog entries; retry once against
    # the new generation, then fall through to a fresh resolution.
    for _ in range(2):
        cache_path = self._cache_index.get(artifact_id)
        if cache_path is None or not cache_path.exists():
            break
        validated_cache_path = cache_path
        options = transfer_options or self._transfer_options

        if generation_scoped:
            try:
                await self._get_blob_url_from_artifact_id(artifact_id)
            except Exception as exc:
                if cache_was_invalidated():
                    cache_generation = self._cache_generation
                    full_cache_generation = self._full_cache_generation
                    continue
                if not isinstance(exc, (ArtifactNotFoundError, PermissionDeniedError, PermissionError)):
                    raise
                if self._cache_index.get(artifact_id) == validated_cache_path:
                    self._cache_index.pop(artifact_id, None)
                try:
                    validated_cache_path.unlink(missing_ok=True)
                except OSError:
                    LOGGER.debug(
                        "Failed to remove unauthorized catalog cache entry %s",
                        validated_cache_path,
                        exc_info=True,
                    )
                raise

        async def validate_cached_file() -> tuple[int, int]:
            check_transfer_cancelled(options, operation="download", resource=str(validated_cache_path))
            cache_file = await _run_blocking_io(
                lambda: _open_cached_file_no_follow(validated_cache_path),
                options=options,
                operation="download",
                resource=str(validated_cache_path),
            )
            try:
                file_stat = await _run_blocking_io(
                    lambda: os.fstat(cache_file.fileno()),
                    options=options,
                    operation="download",
                    resource=str(validated_cache_path),
                )
                check_transfer_size(
                    file_stat.st_size,
                    options,
                    operation="download",
                    resource=str(validated_cache_path),
                )
                if options.expected_sha256 is not None:
                    await hash_file(
                        cache_file,
                        options=options,
                        context=context or RequestContext(),
                        operation="download",
                        resource=str(validated_cache_path),
                    )
                return file_stat.st_dev, file_stat.st_ino
            finally:
                await _run_blocking_io(cache_file.close)

        try:
            await await_transfer(
                validate_cached_file(),
                options,
                operation="download",
                resource=str(validated_cache_path),
            )
        except Exception:
            if not cache_was_invalidated():
                raise
        else:
            if not cache_was_invalidated():
                LOGGER.debug(f"Asset already cached: {cache_path}")
                return cache_path
        self._remove_validated_cache_entry(artifact_id, validated_cache_path)
        cache_generation = self._cache_generation
        full_cache_generation = self._full_cache_generation

    for fetch_attempt in range(2):
        # Route to appropriate resolver based on artifact type
        if artifact_type == "blob":
            resource_url = await self._get_blob_url_from_artifact_id(artifact_id)
        elif artifact_type == "local":
            # Local artifacts: the artifact_id is the file path itself
            resource_url = artifact_id
        else:
            raise ValueError(f"Unsupported artifact type: {artifact_type}. {self._asset_tag_guidance()}")

        # Fetch and cache the asset
        LOGGER.debug(f"Fetching and caching {artifact_type} asset")
        cache_path = self._get_cache_file_path(
            resource_url,
            cache_salt=(
                f"{full_cache_generation}:{cache_generation if generation_scoped else ''}"
                if full_cache_generation or generation_scoped
                else None
            ),
        )

        # Stream asset directly to file to avoid loading into memory
        bytes_written = await self._fetch_asset_to_file(
            resource_url,
            cache_path,
            context=context,
            transfer_options=transfer_options,
            trusted_catalog_reference=self._catalog_managed_revision_access and generation_scoped,
        )

        if cache_was_invalidated():
            cache_path.unlink(missing_ok=True)
            if fetch_attempt == 1:
                raise BackendUnavailableError(
                    "Catalog authorization changed repeatedly during download.",
                    resource_id=safe_artifact_reference(artifact_id),
                    operation="download",
                )
            cache_generation = self._cache_generation
            full_cache_generation = self._full_cache_generation
            continue

        # Update index (use artifact_id as key)
        self._cache_index[artifact_id] = cache_path

        LOGGER.debug(f"Cached asset to disk ({bytes_written} bytes)")
        return cache_path

    raise AssertionError("bounded cache fetch loop exited unexpectedly")

get_asset_info(qualified_name)

Get information about a cached asset.

Parameters:

Name Type Description Default
qualified_name str

Asset identifier

required

Returns:

Type Description
dict

Dict with asset metadata

Source code in src/agora_workbench/code_execution/data_access/manager.py
def get_asset_info(self, qualified_name: str) -> dict:
    """
    Get information about a cached asset.

    Args:
        qualified_name: Asset identifier

    Returns:
        Dict with asset metadata
    """
    info = {
        "qualified_name": safe_artifact_reference(qualified_name),
        "cached": qualified_name in self._cache_index,
    }

    if qualified_name in self._cache_index:
        cache_path = self._cache_index[qualified_name]
        if cache_path.exists():
            info["size_bytes"] = cache_path.stat().st_size
            info["cache_location"] = str(cache_path)

    return info

list_available()

List assets available in this session's cache.

Returns:

Type Description
list[str]

List of qualified names currently cached

Source code in src/agora_workbench/code_execution/data_access/manager.py
def list_available(self) -> list[str]:
    """
    List assets available in this session's cache.

    Returns:
        List of qualified names currently cached
    """
    return list(self._cache_index.keys())

invalidate_cache_entries(*, artifact_id_prefix=None)

Forget cached artifacts matching a resolver namespace.

Entries are removed from the index so a subsequent lookup must pass through resolution and authorization again. Their files remain until session cleanup because a kernel may still be opening a previously returned path.

Source code in src/agora_workbench/code_execution/data_access/manager.py
def invalidate_cache_entries(self, *, artifact_id_prefix: str | None = None) -> None:
    """Forget cached artifacts matching a resolver namespace.

    Entries are removed from the index so a subsequent lookup must pass
    through resolution and authorization again. Their files remain until
    session cleanup because a kernel may still be opening a previously
    returned path.
    """
    self._cache_generation += 1
    if artifact_id_prefix is None:
        self._full_cache_generation += 1
    keys = [
        artifact_id
        for artifact_id in self._cache_index
        if artifact_id_prefix is None or artifact_id.startswith(artifact_id_prefix)
    ]
    for artifact_id in keys:
        self._cache_index.pop(artifact_id, None)

cleanup()

Clean up cache directory, credentials, and resources.

Removes the temporary cache directory if it was created by this manager. Call this when the session is ending to free up disk space.

The cache directory removal (and index clearing) always completes synchronously before this method returns — it does not depend on an event loop draining a deferred task. Only the inherently-async resource closes (fetchers, resolver, owned credential) are deferred as a task when a loop is already running; that task is returned so a caller (e.g. __del__) can retain or await it, but losing it never leaks the on-disk cache directory.

Source code in src/agora_workbench/code_execution/data_access/manager.py
def cleanup(self) -> asyncio.Task[None] | None:
    """
    Clean up cache directory, credentials, and resources.

    Removes the temporary cache directory if it was created by this manager.
    Call this when the session is ending to free up disk space.

    The cache directory removal (and index clearing) always completes
    synchronously before this method returns — it does not depend on an
    event loop draining a deferred task. Only the inherently-async
    resource closes (fetchers, resolver, owned credential) are deferred
    as a task when a loop is already running; that task is returned so a
    caller (e.g. ``__del__``) can retain or await it, but losing it never
    leaks the on-disk cache directory.
    """
    self._cache_index.clear()

    try:
        loop = asyncio.get_running_loop()
    except RuntimeError:
        loop = None

    task: asyncio.Task[None] | None = None
    try:
        if loop is not None:
            task = loop.create_task(self._aclose_async_resources())
        else:
            asyncio.run(self._aclose_async_resources())
    finally:
        self._remove_cache_dir()
    return task

aclose() async

Async cleanup — preferred over sync cleanup() when inside an event loop.

Source code in src/agora_workbench/code_execution/data_access/manager.py
async def aclose(self) -> None:
    """Async cleanup — preferred over sync cleanup() when inside an event loop."""
    self._cache_index.clear()
    try:
        await self._aclose_async_resources()
    finally:
        self._remove_cache_dir()

__del__()

Cleanup on garbage collection.

Source code in src/agora_workbench/code_execution/data_access/manager.py
def __del__(self):
    """Cleanup on garbage collection."""
    try:
        self.cleanup()
    except Exception as e:
        LOGGER.exception(f"Error during DataLakeDataManager cleanup: {e}")

AssetPublisher(credential=None)

Bases: ABC

Base class for artifact publishers.

Publishers are the symmetric counterpart to :class:~.fetchers.AssetFetcher: they push a local file produced by an agent session to a remote storage destination.

Concrete implementations are configured at server startup and registered with the :class:~code_execution.server.CodeExecutionServer. Operators control which destinations are reachable by which publishers they register — no separate allowlist environment variable is needed.

Initialise the publisher with an optional async token credential.

Parameters:

Name Type Description Default
credential 'AsyncTokenCredential | None'

An AsyncTokenCredential that provides tokens for downstream Azure resources (e.g. ManagedIdentityCredential). May be None for publishers that don't require credentials (e.g. local filesystem).

None
Source code in src/agora_workbench/code_execution/data_access/publishers.py
def __init__(self, credential: "AsyncTokenCredential | None" = None):
    """
    Initialise the publisher with an optional async token credential.

    Args:
        credential: An ``AsyncTokenCredential`` that provides tokens for
            downstream Azure resources (e.g. ``ManagedIdentityCredential``).
            May be ``None`` for publishers that don't require credentials
            (e.g. local filesystem).
    """
    self.credential = credential

destination_name abstractmethod property

Logical name used for routing in the unified send tool.

This name is what agents pass as the to parameter, e.g. "blob", "user", "gis". Must be unique across all publishers registered on a single server.

publish(local_path, name, session_id, *, options=None, context=None) abstractmethod async

Publish a local artifact to this publisher's configured destination.

The publisher owns path placement logic — it combines its configured base path with the session context and the logical name to derive the full destination path.

Parameters:

Name Type Description Default
local_path Path

Absolute path to the file to publish.

required
name str

Logical name (relative path-like value from the tag inner text, e.g. "results.csv" or "subdir/report.pdf").

required
session_id str

Active session ID used to scope the upload path.

required

Returns:

Type Description
str

The remote URI of the published artifact (e.g.

str

"https://account.blob.core.windows.net/container/session/name"

str

or "/mnt/shared/outputs/session/name").

Source code in src/agora_workbench/code_execution/data_access/publishers.py
@abstractmethod
async def publish(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> str:
    """Publish a local artifact to this publisher's configured destination.

    The publisher owns path placement logic — it combines its configured
    base path with the session context and the logical name to derive the
    full destination path.

    Args:
        local_path: Absolute path to the file to publish.
        name: Logical name (relative path-like value from the tag inner
            text, e.g. ``"results.csv"`` or ``"subdir/report.pdf"``).
        session_id: Active session ID used to scope the upload path.

    Returns:
        The remote URI of the published artifact (e.g.
        ``"https://account.blob.core.windows.net/container/session/name"``
        or ``"/mnt/shared/outputs/session/name"``).
    """
    raise NotImplementedError

publish_with_result(local_path, name, session_id, *, options=None, context=None) async

Publish with detailed transfer data when supported by the provider.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish_with_result(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> tuple[str, TransferResult]:
    """Publish with detailed transfer data when supported by the provider."""
    raise UnsupportedOperationError(
        f"{type(self).__name__} does not implement bounded file publishing.",
        operation="upload",
    )

can_handle(destination) abstractmethod

Check whether this publisher handles the given tagged destination.

Parameters:

Name Type Description Default
destination str

Tagged destination string, e.g. "<blob>results.csv</blob>".

required

Returns:

Type Description
bool

True if this publisher accepts the tag type.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
@abstractmethod
def can_handle(self, destination: str) -> bool:
    """Check whether this publisher handles the given tagged destination.

    Args:
        destination: Tagged destination string, e.g. ``"<blob>results.csv</blob>"``.

    Returns:
        ``True`` if this publisher accepts the tag type.
    """
    raise NotImplementedError

close() async

Release any resources held by this publisher.

The default implementation is a no-op; override when the publisher holds pooled connections or clients that need explicit teardown.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def close(self) -> None:
    """Release any resources held by this publisher.

    The default implementation is a no-op; override when the publisher
    holds pooled connections or clients that need explicit teardown.
    """

BlobPublisher(account_url, container, credential=None, *, prefix='', staging_dir=None)

Bases: AssetPublisher

Publisher that uploads artifacts to Azure Blob Storage.

Configured at startup with a storage account URL and container name. Files are placed at {container}/{session_id}/{name} inside the configured account.

Maintains a per-account BlobServiceClient cache to amortise TCP/TLS handshake and token acquisition costs across multiple publishes.

Handles destination tags of the form <blob>name</blob>.

Initialise the BlobPublisher.

Parameters:

Name Type Description Default
account_url str

Azure Storage account URL, e.g. "https://myaccount.blob.core.windows.net".

required
container str

Container name to upload into.

required
credential 'AsyncTokenCredential | None'

An AsyncTokenCredential for blob auth (typically managed identity). Reuse the same credential instance as the server's :class:~.fetchers.BlobFetcher to avoid redundant token refreshes.

None
Source code in src/agora_workbench/code_execution/data_access/publishers.py
def __init__(
    self,
    account_url: str,
    container: str,
    credential: "AsyncTokenCredential | None" = None,
    *,
    prefix: str = "",
    staging_dir: Path | str | None = None,
):
    """
    Initialise the BlobPublisher.

    Args:
        account_url: Azure Storage account URL, e.g.
            ``"https://myaccount.blob.core.windows.net"``.
        container: Container name to upload into.
        credential: An ``AsyncTokenCredential`` for blob auth (typically
            managed identity).  Reuse the same credential instance as the
            server's :class:`~.fetchers.BlobFetcher` to avoid redundant
            token refreshes.
    """
    super().__init__(credential=credential)
    parsed = urlsplit(account_url)
    if (
        parsed.scheme.lower() != "https"
        or parsed.username is not None
        or parsed.password is not None
        or "@" in parsed.netloc
        or parsed.path not in {"", "/"}
        or parsed.query
        or parsed.fragment
    ):
        raise ValueError("BlobPublisher account_url must be a credential-free Azure HTTPS account URL.")
    account, _, _ = parse_azure_uri(f"{account_url.rstrip('/')}/container")
    scope = AzureBlobScope(account, container, prefix)
    self._account_url = f"https://{scope.account}.blob.core.windows.net"
    self._container = scope.container
    self._prefix = scope.prefix
    configured_staging = staging_dir or (
        Path(os.getenv("MCP_ASSET_CACHE_DIR", os.getcwd())) / ".agora-transfer-staging"
    )
    self._staging_dir = Path(os.path.abspath(os.fspath(configured_staging)))
    if not _USE_POSIX_DIR_FDS:
        raise UnsupportedOperationError(
            "BlobPublisher requires POSIX descriptor-relative filesystem support for secure staging.",
            operation="upload",
        )
    self._staging_fd: int | None = None
    self._staging_identity: tuple[int, int] | None = None
    self._staging_fd = _open_or_create_posix_directory(self._staging_dir)
    stat_result = os.fstat(self._staging_fd)
    self._staging_identity = (stat_result.st_dev, stat_result.st_ino)
    self._client = None  # lazily initialised

close() async

Close the underlying BlobServiceClient.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def close(self) -> None:
    """Close the underlying BlobServiceClient."""
    try:
        if self._client is not None:
            await self._client.close()
            self._client = None
    finally:
        if self._staging_fd is not None:
            os.close(self._staging_fd)
            self._staging_fd = None

can_handle(destination)

Return True for <blob>…</blob> destinations.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
def can_handle(self, destination: str) -> bool:
    """Return ``True`` for ``<blob>…</blob>`` destinations."""
    parsed = parse_destination_tag(destination)
    return parsed is not None and parsed[0] == "blob"

publish(local_path, name, session_id, *, options=None, context=None) async

Upload local_path to {container}/{session_id}/{name}.

Parameters:

Name Type Description Default
local_path Path

Absolute path to the file to upload.

required
name str

Logical name / relative path within the session's namespace.

required
session_id str

Session ID used to scope the blob path.

required

Returns:

Type Description
str

The full HTTPS URL of the uploaded blob.

Raises:

Type Description
FileNotFoundError

If local_path does not exist.

ClientAuthenticationError

If the credential is not authorised to write to the container.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> str:
    """Upload *local_path* to ``{container}/{session_id}/{name}``.

    Args:
        local_path: Absolute path to the file to upload.
        name: Logical name / relative path within the session's namespace.
        session_id: Session ID used to scope the blob path.

    Returns:
        The full HTTPS URL of the uploaded blob.

    Raises:
        FileNotFoundError: If *local_path* does not exist.
        azure.core.exceptions.ClientAuthenticationError: If the credential
            is not authorised to write to the container.
    """
    remote_uri, _ = await self.publish_with_result(
        local_path,
        name,
        session_id,
        options=options,
        context=context,
    )
    return remote_uri

publish_with_result(local_path, name, session_id, *, options=None, context=None) async

Upload a regular file with bounded reads, timeout, cancellation, and checksum.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish_with_result(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> tuple[str, TransferResult]:
    """Upload a regular file with bounded reads, timeout, cancellation, and checksum."""
    options = options or TransferOptions()
    context = context or RequestContext()
    _validate_blob_metadata(dict(options.object_metadata))
    if not local_path.is_file():
        raise FileNotFoundError(f"Artifact not found at {local_path}")
    _validate_artifact_name(name, allow_reserved=options.allow_reserved)
    if session_id:
        _validate_artifact_name(session_id)
    relative_path = "/".join(part for part in (self._prefix, session_id, name) if part)
    blob_path = _validate_publish_path(relative_path, allow_reserved=options.allow_reserved)
    remote_uri = azure_uri_from_blob_name(
        parse_azure_uri(f"{self._account_url}/{self._container}")[0],
        self._container,
        blob_path,
    )
    display_uri = safe_transfer_resource(remote_uri)
    LOGGER.info(
        "BlobPublisher: uploading %s → %s/%s/%s",
        local_path,
        self._account_url,
        self._container,
        blob_path,
    )

    client = self._get_client()
    blob_client = client.get_blob_client(container=self._container, blob=blob_path)
    started = time.monotonic()
    await emit_transfer_diagnostic(options, TransferDiagnostic("upload", "started", context, display_uri))
    snapshot_name = f"{secrets.token_hex(16)}.upload"
    staging_fd: int | None = None
    snapshot_fd: int | None = None
    snapshot_created = False
    try:
        staging_fd = self._open_verified_staging_root()
        snapshot_fd = os.open(
            snapshot_name,
            os.O_RDWR | os.O_CREAT | os.O_EXCL | getattr(os, "O_NOFOLLOW", 0),
            0o600,
            dir_fd=staging_fd,
        )
        snapshot_created = True

        async def perform_upload() -> TransferResult:
            snapshot = await _copy_local_descriptors(local_path, snapshot_fd, options, context)
            os.lseek(snapshot_fd, 0, os.SEEK_SET)
            check_transfer_cancelled(options, operation="upload", resource=display_uri)
            metadata = dict(options.object_metadata)
            with _TrackedUploadStream(os.fdopen(os.dup(snapshot_fd), "rb", buffering=0)) as tracked_source:
                if options.create_exclusive and metadata:
                    upload = blob_client.upload_blob(
                        tracked_source,
                        overwrite=False,
                        if_none_match="*",
                        metadata=metadata,
                    )
                elif options.create_exclusive:
                    upload = blob_client.upload_blob(
                        tracked_source,
                        overwrite=False,
                        if_none_match="*",
                    )
                elif metadata:
                    upload = blob_client.upload_blob(
                        tracked_source,
                        overwrite=True,
                        metadata=metadata,
                    )
                else:
                    upload = blob_client.upload_blob(tracked_source, overwrite=True)
                await await_transfer(
                    upload,
                    TransferOptions(
                        max_bytes=options.max_bytes,
                        quota_bytes=options.quota_bytes,
                        timeout_seconds=None,
                        chunk_size=options.chunk_size,
                        cancellation_event=options.cancellation_event,
                    ),
                    operation="upload",
                    resource=display_uri,
                )
                if not tracked_source.fully_consumed(snapshot.bytes_transferred):
                    raise OSError("Blob provider did not consume the complete immutable upload snapshot.")
            return snapshot

        try:
            if options.timeout_seconds is None:
                uploaded = await perform_upload()
            else:
                async with asyncio.timeout(options.timeout_seconds):
                    uploaded = await perform_upload()
        except TimeoutError as exc:
            message = (
                "Provider transfer timed out."
                if options.timeout_seconds is None
                else f"Transfer exceeded the configured {options.timeout_seconds:g}-second timeout."
            )
            raise TransferTimeoutError(
                message,
                resource_id=display_uri,
                operation="upload",
            ) from exc
    except BaseException as exc:
        await emit_transfer_diagnostic(
            options,
            TransferDiagnostic("upload", "failed", context, display_uri, error_type=type(exc).__name__),
        )
        raise
    finally:
        if snapshot_fd is not None:
            try:
                os.close(snapshot_fd)
            except OSError:
                LOGGER.warning("Could not close BlobPublisher snapshot descriptor.", exc_info=True)
        try:
            if snapshot_created and staging_fd is not None:
                try:
                    os.unlink(snapshot_name, dir_fd=staging_fd)
                except FileNotFoundError:
                    LOGGER.debug("BlobPublisher snapshot was already removed: %s", snapshot_name)
                except OSError:
                    LOGGER.warning("Could not remove BlobPublisher snapshot %s.", snapshot_name, exc_info=True)
        finally:
            if staging_fd is not None:
                try:
                    os.close(staging_fd)
                except OSError:
                    LOGGER.warning("Could not close BlobPublisher staging descriptor.", exc_info=True)
    result = TransferResult(
        uploaded.bytes_transferred,
        uploaded.checksum_sha256,
        context,
        display_uri,
        time.monotonic() - started,
        True if options.create_exclusive else None,
        options.object_metadata,
    )
    await emit_transfer_diagnostic(
        options,
        TransferDiagnostic(
            "upload",
            "completed",
            context,
            display_uri,
            result.bytes_transferred,
            result.checksum_sha256,
        ),
    )
    LOGGER.info("BlobPublisher: uploaded %d bytes → %s", result.bytes_transferred, display_uri)
    return f"{self._account_url}{urlsplit(remote_uri).path}", result

GuiPublisher(public_url_fn)

Bases: AssetPublisher

Publisher that exposes artifacts via the server's download endpoint.

Unlike Blob or Local publishers that transfer the file elsewhere, this publisher simply returns the server's /artifacts/ download URL for the already-registered artifact. The activity UI surfaces this URL so the user can download the file from their browser.

This publisher is auto-registered on every CodeExecutionServer so agents can always use <gui>filename</gui> to make outputs downloadable without requiring operator-configured storage backends.

Handles destination tags of the form <gui>name</gui>.

Initialise the GuiPublisher.

Parameters:

Name Type Description Default
public_url_fn 'Callable[[], str]'

A callable returning the server's public base URL (e.g. server.public_url). Deferred so the URL reflects the actual bind address after run_http() is called.

required
Source code in src/agora_workbench/code_execution/data_access/publishers.py
def __init__(self, public_url_fn: "Callable[[], str]"):
    """
    Initialise the GuiPublisher.

    Args:
        public_url_fn: A callable returning the server's public base URL
            (e.g. ``server.public_url``).  Deferred so the URL reflects
            the actual bind address after ``run_http()`` is called.
    """
    super().__init__(credential=None)
    self._public_url_fn = public_url_fn

can_handle(destination)

Return True for <gui>…</gui> destinations.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
def can_handle(self, destination: str) -> bool:
    """Return ``True`` for ``<gui>…</gui>`` destinations."""
    parsed = parse_destination_tag(destination)
    return parsed is not None and parsed[0] == "gui"

publish(local_path, name, session_id, *, options=None, context=None) async

Return the download URL for the artifact.

The artifact must already be registered in the session manager's artifact registry (populated by the snapshot-diff after execution). The caller (the publish tool wrapper) is responsible for passing the download token via the _download_token attribute set on this instance before calling publish().

Parameters:

Name Type Description Default
local_path Path

Absolute path to the artifact file.

required
name str

Logical name / relative path for the URL's filename segment.

required
session_id str

Session ID scoping the artifact.

required

Returns:

Type Description
str

The fully-qualified download URL.

Raises:

Type Description
FileNotFoundError

If local_path does not exist.

RuntimeError

If no download token was provided.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> str:
    """Return the download URL for the artifact.

    The artifact must already be registered in the session manager's
    artifact registry (populated by the snapshot-diff after execution).
    The caller (the publish tool wrapper) is responsible for passing the
    download token via the ``_download_token`` attribute set on this
    instance before calling ``publish()``.

    Args:
        local_path: Absolute path to the artifact file.
        name: Logical name / relative path for the URL's filename segment.
        session_id: Session ID scoping the artifact.

    Returns:
        The fully-qualified download URL.

    Raises:
        FileNotFoundError: If *local_path* does not exist.
        RuntimeError: If no download token was provided.
    """
    del options, context
    if not local_path.is_file():
        raise FileNotFoundError(f"Artifact not found at {local_path}")

    _validate_artifact_name(name)

    token = getattr(self, "_download_token", None)
    if not token:
        raise RuntimeError(
            "GuiPublisher requires a download token set via _download_token before publish() is called."
        )

    public_base = (os.getenv("SERVER_PUBLIC_URL") or self._public_url_fn()).rstrip("/")
    download_url = f"{public_base}/artifacts/{session_id}/{token}/{name}"
    LOGGER.info("GuiPublisher: exposing %s for session %s as %s", local_path, session_id, name)
    return download_url

LocalFilePublisher(base_dir)

Bases: AssetPublisher

Publisher that copies artifacts to a local directory.

Configured at startup with a base directory. Files are placed at {base_dir}/{session_id}/{name}.

Handles destination tags of the form <local>name</local>.

Initialise the LocalFilePublisher.

Parameters:

Name Type Description Default
base_dir Path | str

Base directory under which session sub-directories are created. Must be an absolute path (or will be resolved to one).

required
Source code in src/agora_workbench/code_execution/data_access/publishers.py
def __init__(self, base_dir: Path | str):
    """
    Initialise the LocalFilePublisher.

    Args:
        base_dir: Base directory under which session sub-directories are
            created.  Must be an absolute path (or will be resolved to
            one).
    """
    super().__init__(credential=None)
    self._base_dir = Path(os.path.abspath(os.fspath(base_dir)))
    self._anchor_fd: int | None = None
    self._anchor_path: Path | None = None
    self._anchor_identity: tuple[int, int] | None = None
    self._root_parts: tuple[str, ...] = ()
    self._root_identity: tuple[int, int] | None = None
    self._portable_anchor_path: Path | None = None
    self._portable_anchor_identity: tuple[int, int] | None = None
    if _USE_POSIX_DIR_FDS:
        self._initialize_root_anchor()
    else:
        self._initialize_portable_root()

close() async

Close the retained trusted root descriptor.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def close(self) -> None:
    """Close the retained trusted root descriptor."""
    if self._anchor_fd is not None:
        os.close(self._anchor_fd)
        self._anchor_fd = None

can_handle(destination)

Return True for <local>…</local> destinations.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
def can_handle(self, destination: str) -> bool:
    """Return ``True`` for ``<local>…</local>`` destinations."""
    parsed = parse_destination_tag(destination)
    return parsed is not None and parsed[0] == "local"

publish(local_path, name, session_id, *, options=None, context=None) async

Copy local_path to {base_dir}/{session_id}/{name}.

Parameters:

Name Type Description Default
local_path Path

Absolute path to the file to copy.

required
name str

Logical name / relative path within the session's namespace.

required
session_id str

Session ID used to scope the destination path.

required

Returns:

Type Description
str

The absolute path of the copied file as a string.

Raises:

Type Description
FileNotFoundError

If local_path does not exist.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> str:
    """Copy *local_path* to ``{base_dir}/{session_id}/{name}``.

    Args:
        local_path: Absolute path to the file to copy.
        name: Logical name / relative path within the session's namespace.
        session_id: Session ID used to scope the destination path.

    Returns:
        The absolute path of the copied file as a string.

    Raises:
        FileNotFoundError: If *local_path* does not exist.
    """
    destination, _ = await self.publish_with_result(
        local_path,
        name,
        session_id,
        options=options,
        context=context,
    )
    return destination

publish_with_result(local_path, name, session_id, *, options=None, context=None) async

Copy into the configured root without following destination symlinks.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish_with_result(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> tuple[str, TransferResult]:
    """Copy into the configured root without following destination symlinks."""
    options = options or TransferOptions()
    context = context or RequestContext()
    if not local_path.is_file():
        raise FileNotFoundError(f"Artifact not found at {local_path}")
    _validate_artifact_name(name, allow_reserved=options.allow_reserved)
    if session_id:
        _validate_artifact_name(session_id)
    if options.object_metadata:
        raise UnsupportedOperationError(
            "LocalFilePublisher does not support object metadata.",
            operation="upload",
        )
    relative_text = normalize_logical_path("/".join(part for part in (session_id, name) if part))
    _validate_publish_path(relative_text, allow_reserved=options.allow_reserved)
    relative = Path(relative_text)
    destination = self._base_dir / relative
    started = time.monotonic()
    await emit_transfer_diagnostic(
        options,
        TransferDiagnostic("upload", "started", context, str(destination)),
    )
    try:
        result = await self._copy_secure(local_path, relative, options, context)
    except BaseException as exc:
        await emit_transfer_diagnostic(
            options,
            TransferDiagnostic("upload", "failed", context, str(destination), error_type=type(exc).__name__),
        )
        raise
    result = TransferResult(
        result.bytes_transferred,
        result.checksum_sha256,
        context,
        str(destination),
        time.monotonic() - started,
        result.created,
        result.object_metadata,
    )
    await emit_transfer_diagnostic(
        options,
        TransferDiagnostic(
            "upload",
            "completed",
            context,
            str(destination),
            result.bytes_transferred,
            result.checksum_sha256,
        ),
    )
    LOGGER.info("LocalFilePublisher: copied %d bytes → %s", result.bytes_transferred, destination)
    return str(destination), result

ServerPublisher(server_name, target_url=None, timeout=60.0, trust_http=False)

Bases: AssetPublisher

Publisher that pushes serialized objects to a peer MCP server's kernel.

Each instance represents a single peer server destination. Operators register one ServerPublisher per reachable peer in the publishers list at server construction time.

The publisher reads a dill-serialized file, base64-encodes it, and POSTs it to the target server's /object-transfer/receive endpoint. URL validation (HTTPS enforcement, host allow-lists) is applied before sending credentials.

.. note::

When ``target_url`` is not provided, the URL is auto-expanded from
``server_name`` using the Docker Compose convention
(``http://{name}-server:8000``). By default, plain HTTP URLs are
rejected by ``_validate_target_url`` unless the hostname appears in
the ``OBJECT_TRANSFER_TRUSTED_HTTP_HOSTS`` environment variable.
Operators deploying with Docker Compose internal networking should
set this variable to include the auto-expanded hostnames.

Session resolution on the target: when session_id is empty the receive endpoint selects the first active session owned by the caller (identified by the forwarded bearer token). If no active session exists on the target, the receive endpoint returns a 404 — the agent must ensure a session exists on the target (e.g. via execute_{target}_code) before sending.

Initialise the ServerPublisher.

Parameters:

Name Type Description Default
server_name str

Logical name of the target server (e.g. "gis"). Used as destination_name for routing and for Docker URL expansion if target_url is not provided.

required
target_url str | None

Full base URL of the target MCP server. If None, the URL is auto-expanded from server_name using the Docker Compose convention (http://{name}-server:8000).

None
timeout float

HTTP read/write timeout in seconds for the transfer request.

60.0
trust_http bool

When True, plain-HTTP transfers to this peer are permitted without listing the host in OBJECT_TRANSFER_TRUSTED_HTTP_HOSTS. Set by the unified send tool when the publisher is built on demand from an operator-curated peer registry, where the operator already chose the scheme in the configured URL.

False
Source code in src/agora_workbench/code_execution/data_access/publishers.py
def __init__(
    self,
    server_name: str,
    target_url: str | None = None,
    timeout: float = 60.0,
    trust_http: bool = False,
):
    """
    Initialise the ServerPublisher.

    Args:
        server_name: Logical name of the target server (e.g. ``"gis"``).
            Used as ``destination_name`` for routing and for Docker URL
            expansion if ``target_url`` is not provided.
        target_url: Full base URL of the target MCP server.  If ``None``,
            the URL is auto-expanded from ``server_name`` using the Docker
            Compose convention (``http://{name}-server:8000``).
        timeout: HTTP read/write timeout in seconds for the transfer request.
        trust_http: When ``True``, plain-HTTP transfers to this peer are
            permitted without listing the host in
            ``OBJECT_TRANSFER_TRUSTED_HTTP_HOSTS``. Set by the unified send
            tool when the publisher is built on demand from an
            operator-curated peer registry, where the operator already
            chose the scheme in the configured URL.
    """
    super().__init__(credential=None)
    self._server_name = server_name
    self._timeout = timeout
    self._trust_http = trust_http
    if target_url:
        self._target_url = target_url.rstrip("/")
    else:
        hostname = server_name if server_name.endswith("-server") else f"{server_name}-server"
        self._target_url = f"http://{hostname}:8000"

can_handle(destination)

ServerPublisher does not use tag-based routing.

It is routed via destination_name in the unified send tool. For backwards-compat with the publish_artifact tag-based flow, this always returns False.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
def can_handle(self, destination: str) -> bool:
    """ServerPublisher does not use tag-based routing.

    It is routed via ``destination_name`` in the unified send tool.
    For backwards-compat with the ``publish_artifact`` tag-based flow,
    this always returns ``False``.
    """
    return False

publish(local_path, name, session_id, *, options=None, context=None) async

Serialize and push the file to the target server's kernel.

The file at local_path is expected to be a dill-serialized pickle. It is base64-encoded and sent to the target's /object-transfer/receive endpoint.

Validates the target URL before sending credentials to prevent SSRF and bearer-token leakage (see _validate_target_url).

Parameters:

Name Type Description Default
local_path Path

Path to the serialized (pickle) file to transfer.

required
name str

Variable name to inject on the target kernel.

required
session_id str

Session ID on the target server (empty string to let the target resolve via bearer token).

required

Returns:

Type Description
str

A human-readable confirmation message.

Raises:

Type Description
FileNotFoundError

If local_path does not exist.

RuntimeError

If no user token was set before calling publish.

ValueError

If the target URL fails validation.

ObjectTransferError

On structured non-2xx responses from the target.

HTTPStatusError

On non-2xx responses from the target.

RequestError

On connection / timeout errors.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish(
    self,
    local_path: Path,
    name: str,
    session_id: str,
    *,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> str:
    """Serialize and push the file to the target server's kernel.

    The file at ``local_path`` is expected to be a dill-serialized pickle.
    It is base64-encoded and sent to the target's ``/object-transfer/receive``
    endpoint.

    Validates the target URL before sending credentials to prevent SSRF and
    bearer-token leakage (see ``_validate_target_url``).

    Args:
        local_path: Path to the serialized (pickle) file to transfer.
        name: Variable name to inject on the target kernel.
        session_id: Session ID on the target server (empty string to
            let the target resolve via bearer token).

    Returns:
        A human-readable confirmation message.

    Raises:
        FileNotFoundError: If *local_path* does not exist.
        RuntimeError: If no user token was set before calling publish.
        ValueError: If the target URL fails validation.
        ObjectTransferError: On structured non-2xx responses from the target.
        httpx.HTTPStatusError: On non-2xx responses from the target.
        httpx.RequestError: On connection / timeout errors.
    """
    import base64
    import re as _re

    import httpx

    from ..object_transfer import (
        MAX_TRANSFER_SIZE_BYTES,
        STREAMING_TRANSFER_INFO_HEADER,
        STREAMING_TRANSFER_VERSION,
        STREAMING_TRANSFER_VERSION_HEADER,
        _validate_target_url,
        encode_streaming_transfer_info,
    )

    options = options or TransferOptions()
    context = context or RequestContext()

    if not local_path.is_file():
        raise FileNotFoundError(f"Transfer file not found at {local_path}")

    # user_token is injected by the send tool before calling publish
    user_token = getattr(self, "_user_token", "")
    if not user_token:
        raise RuntimeError("ServerPublisher requires _user_token to be set before publish() is called.")

    _validate_target_url(self._target_url, trust_http=self._trust_http)

    # Strip common MCP path suffixes so the agent can pass the MCP
    # endpoint URL directly (e.g. http://gis-server:8000/mcp).
    base = _re.sub(r"/mcp/?$", "", self._target_url.rstrip("/"))
    receive_url = f"{base}/object-transfer/receive"

    configured_max = options.max_bytes
    bounded_options = replace(
        options,
        max_bytes=MAX_TRANSFER_SIZE_BYTES
        if configured_max is None
        else min(configured_max, MAX_TRANSFER_SIZE_BYTES),
    )
    display_resource = safe_transfer_resource(receive_url)
    started = time.monotonic()
    await emit_transfer_diagnostic(
        bounded_options,
        TransferDiagnostic("upload", "started", context, display_resource),
    )
    snapshot = tempfile.TemporaryFile(mode="w+b")
    try:
        snapshot_result = await _copy_local_descriptors(local_path, snapshot.fileno(), bounded_options, context)
        await _run_blocking_io(
            lambda: snapshot.seek(0),
            options=bounded_options,
            operation="upload",
            resource=display_resource,
        )
    except BaseException as exc:
        await emit_transfer_diagnostic(
            bounded_options,
            TransferDiagnostic("upload", "failed", context, display_resource, error_type=type(exc).__name__),
        )
        await _run_blocking_io(snapshot.close, options=None, operation="upload", resource=display_resource)
        raise

    metadata = {
        "source_server": getattr(self, "_source_server", "unknown"),
        "transfer_id": getattr(self, "_transfer_id", ""),
    }
    transfer_info = encode_streaming_transfer_info(
        variable_name=name,
        session_id=session_id,
        metadata=metadata,
        size_bytes=snapshot_result.bytes_transferred,
        checksum_sha256=snapshot_result.checksum_sha256,
    )
    prefix = ('{"variable_name":' + json.dumps(name) + ',"data":"').encode()
    suffix = (
        '","metadata":'
        + json.dumps(metadata, separators=(",", ":"))
        + (',"session_id":' + json.dumps(session_id) if session_id else "")
        + "}"
    ).encode()

    async def payload_chunks():
        yield prefix
        remainder = b""
        while True:
            check_transfer_cancelled(bounded_options, operation="upload", resource=display_resource)
            chunk = await _run_blocking_io(
                lambda: snapshot.read(bounded_options.chunk_size),
                options=bounded_options,
                operation="upload",
                resource=display_resource,
            )
            if not chunk:
                break
            data = remainder + chunk
            encoded_length = len(data) - (len(data) % 3)
            if encoded_length:
                yield base64.b64encode(data[:encoded_length])
            remainder = data[encoded_length:]
        if remainder:
            yield base64.b64encode(remainder)
        yield suffix

    elapsed = time.monotonic() - started
    remaining_timeout = None if options.timeout_seconds is None else options.timeout_seconds - elapsed
    if remaining_timeout is not None and remaining_timeout <= 0:
        await emit_transfer_diagnostic(
            bounded_options,
            TransferDiagnostic(
                "upload",
                "failed",
                context,
                display_resource,
                error_type=TransferTimeoutError.__name__,
            ),
        )
        await _run_blocking_io(snapshot.close, options=None, operation="upload", resource=display_resource)
        raise TransferTimeoutError(
            f"Transfer exceeded the configured {options.timeout_seconds:g}-second timeout.",
            resource_id=display_resource,
            operation="upload",
        )
    request_options = replace(bounded_options, timeout_seconds=remaining_timeout, expected_sha256=None)

    try:
        async with httpx.AsyncClient(
            timeout=httpx.Timeout(connect=10.0, read=self._timeout, write=self._timeout, pool=10.0),
        ) as client:
            response = await await_transfer(
                client.post(
                    receive_url,
                    content=payload_chunks(),
                    headers={
                        "Authorization": f"Bearer {user_token}",
                        "Content-Type": "application/json",
                        STREAMING_TRANSFER_VERSION_HEADER: STREAMING_TRANSFER_VERSION,
                        STREAMING_TRANSFER_INFO_HEADER: transfer_info,
                    },
                ),
                request_options,
                operation="upload",
                resource=display_resource,
            )
            try:
                result = response.json()
            except ValueError:
                response.raise_for_status()
                raise

            try:
                response.raise_for_status()
            except httpx.HTTPStatusError as exc:
                if isinstance(result, dict):
                    raise ObjectTransferError(
                        server_name=self._server_name,
                        status_code=response.status_code,
                        response_body=result,
                    ) from exc
                raise
    except BaseException as exc:
        await emit_transfer_diagnostic(
            bounded_options,
            TransferDiagnostic("upload", "failed", context, display_resource, error_type=type(exc).__name__),
        )
        raise
    finally:
        await _run_blocking_io(snapshot.close, options=None, operation="upload", resource=display_resource)

    await emit_transfer_diagnostic(
        bounded_options,
        TransferDiagnostic(
            "upload",
            "completed",
            context,
            display_resource,
            snapshot_result.bytes_transferred,
            snapshot_result.checksum_sha256,
        ),
    )

    return f"Injected '{name}' into {self._server_name} kernel (response: {result})"

TransferDiagnostic(operation, state, context, resource=None, bytes_transferred=0, checksum_sha256=None, error_type=None) dataclass

Credential-safe event delivered to an operator-provided audit hook.

TransferOptions(max_bytes=DEFAULT_TRANSFER_MAX_BYTES, quota_bytes=None, timeout_seconds=DEFAULT_TRANSFER_TIMEOUT_SECONDS, chunk_size=DEFAULT_TRANSFER_CHUNK_BYTES, expected_sha256=None, cancellation_event=None, diagnostic_hook=None, create_exclusive=False, object_metadata=dict(), allow_reserved=False) dataclass

Limits and integrity requirements for one streaming transfer.

max_bytes bounds the individual object while quota_bytes represents caller/provider capacity remaining before the operation starts. The smaller non-None value is enforced. fetch() convenience methods still load the complete object in memory; these options apply to streaming methods.

effective_max_bytes property

Return the tightest configured object/quota bound.

TransferResult(bytes_transferred, checksum_sha256, context, resource=None, elapsed_seconds=0.0, created=None, object_metadata=dict()) dataclass

Completed bounded transfer details, including the original caller context.

create_storage_credential(client_id=None)

Create an async credential chain suitable for Azure Storage access.

The returned ChainedTokenCredential first tries MsalCacheCredential (the mounted az login cache for local development), then ManagedIdentityCredential (the Azure-assigned production identity).

The chain raises CredentialUnavailableError from each link until one succeeds, making the credential work transparently in both environments.

Parameters:

Name Type Description Default
client_id str | None

Optional user-assigned managed identity client ID. When provided it takes precedence over the DEFAULT_IDENTITY_CLIENT_ID environment variable. Callers that resolve the MI client id from a different convention (e.g. AZURE_CLIENT_ID) should pass it here so production keeps binding to the same user-assigned identity.

None

Returns:

Type Description
AsyncTokenCredential

An AsyncTokenCredential usable with BlobPublisher,

AsyncTokenCredential

BlobFetcher, BlobServiceClient, etc.

Source code in src/agora_workbench/code_execution/data_access/credentials.py
def create_storage_credential(client_id: str | None = None) -> AsyncTokenCredential:
    """Create an async credential chain suitable for Azure Storage access.

    The returned ``ChainedTokenCredential`` first tries
    ``MsalCacheCredential`` (the mounted ``az login`` cache for local
    development), then ``ManagedIdentityCredential`` (the Azure-assigned
    production identity).

    The chain raises ``CredentialUnavailableError`` from each link until one
    succeeds, making the credential work transparently in both environments.

    Args:
        client_id: Optional user-assigned managed identity client ID.  When
            provided it takes precedence over the ``DEFAULT_IDENTITY_CLIENT_ID``
            environment variable.  Callers that resolve the MI client id from a
            different convention (e.g. ``AZURE_CLIENT_ID``) should pass it here
            so production keeps binding to the same user-assigned identity.

    Returns:
        An ``AsyncTokenCredential`` usable with ``BlobPublisher``,
        ``BlobFetcher``, ``BlobServiceClient``, etc.
    """
    try:
        from azure.identity.aio import (
            ChainedTokenCredential,
            ManagedIdentityCredential,
        )
    except ImportError as exc:
        raise RuntimeError("Azure storage credentials require the 'agora-workbench[azure]' extra.") from exc

    managed_identity_client_id = (client_id.strip() if client_id is not None else None) or os.getenv(
        "DEFAULT_IDENTITY_CLIENT_ID"
    )

    credentials: list[MsalCacheCredential | ManagedIdentityCredential] = [MsalCacheCredential()]

    if managed_identity_client_id:
        credentials.append(ManagedIdentityCredential(client_id=managed_identity_client_id))
    else:
        credentials.append(ManagedIdentityCredential())

    return ChainedTokenCredential(*credentials)

publish_compat(publisher, *, local_path, name, session_id, options=None, context=None) async

Call modern or legacy publishers without masking implementation errors.

Source code in src/agora_workbench/code_execution/data_access/publishers.py
async def publish_compat(
    publisher: Any,
    *,
    local_path: Path,
    name: str,
    session_id: str,
    options: TransferOptions | None = None,
    context: RequestContext | None = None,
) -> str:
    """Call modern or legacy publishers without masking implementation errors."""
    publish = publisher.publish
    implementation = getattr(publish, "__func__", None)
    if implementation is not None and implementation is getattr(type(publisher), "publish", None):
        supports_kwargs, supports_options, supports_context = _publish_capabilities(implementation)
    else:
        supports_kwargs, supports_options, supports_context = _inspect_publish_capabilities(publish)
    kwargs: dict[str, object] = {
        "local_path": local_path,
        "name": name,
        "session_id": session_id,
    }
    if supports_kwargs or supports_options:
        kwargs["options"] = options
    if supports_kwargs or supports_context:
        kwargs["context"] = context
    return await publish(**kwargs)

Asset Resolution

agora_workbench.code_execution.data_access.resolution

Utilities for detecting and resolving DataLake-cataloged asset references.

AssetResolutionMiddleware(server)

Bases: Middleware

FastMCP middleware that resolves DataLake asset references before Pydantic validation.

When the agent sends a tool call with tagged asset references like <blob>base64_id</blob>, this middleware intercepts the arguments, resolves each reference to a local cache path via the session's data manager, and replaces the argument value with the path string. This happens before FastMCP/Pydantic coerces arguments against the function signature, so parameters with their natural types (Path, bool, int, etc.) receive properly typed values instead of being mangled by premature coercion.

Resolution metadata (qualified name, cache path, parameter name) is stored in a ContextVar so that the tool callback can inject the asset into the kernel for the generated execution code.

Source code in src/agora_workbench/code_execution/data_access/resolution.py
def __init__(self, server: "CodeExecutionServer"):
    self.server = server

looks_like_qualified_name(value)

Check if a string looks like a DataLake qualified name.

Recognizes type-tagged artifact format: id Examples: abc123, xyz789

Parameters:

Name Type Description Default
value str

String to check

required

Returns:

Type Description
bool

True if value matches the type-tagged artifact pattern

Source code in src/agora_workbench/code_execution/data_access/resolution.py
def looks_like_qualified_name(value: str) -> bool:
    """
    Check if a string looks like a DataLake qualified name.

    Recognizes type-tagged artifact format: <type>id</type>
    Examples: <blob>abc123</blob>, <sql>xyz789</sql>

    Args:
        value: String to check

    Returns:
        True if value matches the type-tagged artifact pattern
    """
    if not isinstance(value, str):
        return False

    stripped = value.strip()
    if not stripped:
        return False

    # Accept both <type>id</type> and unclosed <type>id (LLM sometimes omits closing tag)
    return bool(re.match(r"^<(\w+)>([^<>]+)</\1>$", stripped)) or bool(re.match(r"^<(\w+)>([^<>]+)$", stripped))

should_resolve_as_asset(value)

Determine if a parameter value should be resolved as a DataLake asset.

Detection is purely value-based: any string matching the type-tagged format <type>id</type> will be resolved.

Parameters:

Name Type Description Default
value Any

Parameter value to check

required

Returns:

Type Description
bool

True if value is a type-tagged asset reference

Source code in src/agora_workbench/code_execution/data_access/resolution.py
def should_resolve_as_asset(value: Any) -> bool:
    """
    Determine if a parameter value should be resolved as a DataLake asset.

    Detection is purely value-based: any string matching the type-tagged
    format ``<type>id</type>`` will be resolved.

    Args:
        value: Parameter value to check

    Returns:
        True if value is a type-tagged asset reference
    """
    if not isinstance(value, str):
        return False

    return looks_like_qualified_name(value)