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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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.
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
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)
¶
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
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
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
capabilities(context)
async
¶
Return write operations allowed for this session and source.
Source code in src/agora_workbench/data_lake/managed.py
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
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
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
artifact_reference(path, artifact_id=None)
¶
Return the normalized effective reference used for authorization.
Source code in src/agora_workbench/data_lake/managed.py
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
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
1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 | |
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
promote(request, context=RequestContext())
async
¶
Explicitly copy a scratch output into durable managed storage.
Source code in src/agora_workbench/data_lake/managed.py
remove(request, context=RequestContext())
async
¶
Commit a tombstone before collecting any managed bytes.
Source code in src/agora_workbench/data_lake/managed.py
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
2235 2236 2237 2238 2239 2240 2241 2242 2243 2244 2245 2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256 2257 2258 2259 2260 2261 2262 2263 2264 2265 2266 2267 2268 2269 2270 2271 2272 2273 2274 2275 2276 2277 2278 2279 2280 2281 2282 2283 2284 2285 2286 2287 2288 2289 2290 2291 2292 2293 2294 2295 2296 2297 2298 2299 2300 2301 2302 2303 2304 2305 2306 2307 2308 2309 2310 2311 2312 2313 2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324 2325 2326 2327 2328 2329 2330 2331 2332 2333 2334 2335 2336 2337 2338 2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407 2408 2409 | |
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
¶
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
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
search(request, context)
async
¶
Search with authorization enforced before effective pagination.
Source code in src/agora_workbench/data_lake/policy.py
list(request, context)
async
¶
List with authorization enforced before effective pagination.
Source code in src/agora_workbench/data_lake/policy.py
get(reference, context)
async
¶
Get an artifact without distinguishing denied from absent artifacts.
Source code in src/agora_workbench/data_lake/policy.py
resolve(reference, context)
async
¶
Resolve an artifact without distinguishing denied from absent artifacts.
Source code in src/agora_workbench/data_lake/policy.py
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.
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
¶
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
list(provider, request, context, authorizer)
async
¶
List only artifacts authorized for the current caller.
Source code in src/agora_workbench/data_lake/protocols.py
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
resolve(provider, reference, context, authorizer)
async
¶
Resolve only an authorized canonical artifact.
Source code in src/agora_workbench/data_lake/protocols.py
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
¶
search(request, context)
async
¶
Search artifacts using provider-defined ranking and opaque cursors.
list(request, context)
async
¶
List artifacts using provider-defined filters and opaque cursors.
get(reference, context)
async
¶
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.
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
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
PolicyEnforcedCatalog
¶
Bases: Protocol
Caller-aware catalog surface produced by policy composition.
policy_mode
property
¶
Return the configured enforcement granularity.
capabilities(context)
async
¶
search(request, context)
async
¶
Search within the current caller's effective authorization scope.
list(request, context)
async
¶
List within the current caller's effective authorization scope.
get(reference, context)
async
¶
resolve(reference, context)
async
¶
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
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
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
canonicalize_azure_uri(uri)
¶
Canonicalize a supported Azure Blob or DFS URI without credentials.
Source code in src/agora_workbench/data_lake/identity.py
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
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
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
normalize_logical_path(path)
¶
Return a portable, source-relative POSIX path.
Source code in src/agora_workbench/data_lake/identity.py
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
sanitize_uri_for_display(uri)
¶
Remove credentials, query parameters, and fragments from a URI.
Source code in src/agora_workbench/data_lake/identity.py
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
stable_source_id(source_type, root)
¶
Derive a stable fallback source ID.
Source code in src/agora_workbench/data_lake/identity.py
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
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
managed_writer_extension_factory(writer)
¶
Adapt a managed writer to CatalogIntegration.capability_extension_factory.
Source code in src/agora_workbench/data_lake/managed.py
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
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
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
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
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
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
618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 738 739 740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 | |
upsert_artifacts_batch(artifacts)
¶
Atomically upsert a batch of artifacts.
Source code in src/agora_workbench/code_execution/data_access/catalog/db.py
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
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
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
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
has_vector(artifact_id)
¶
Return whether the current artifact has an indexed embedding.
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
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
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
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
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
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
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
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
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
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
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
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
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
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
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
356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 | |
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
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)
¶
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
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
search(request, context)
async
¶
Search with authorization enforced before effective pagination.
Source code in src/agora_workbench/data_lake/policy.py
list(request, context)
async
¶
List with authorization enforced before effective pagination.
Source code in src/agora_workbench/data_lake/policy.py
get(reference, context)
async
¶
Get an artifact without distinguishing denied from absent artifacts.
Source code in src/agora_workbench/data_lake/policy.py
resolve(reference, context)
async
¶
Resolve an artifact without distinguishing denied from absent artifacts.
Source code in src/agora_workbench/data_lake/policy.py
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
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
load()
async
¶
Load or refresh all manifests, retaining the previous valid generation on failure.
readiness()
¶
Return readiness, stale bounds, and per-source refresh state.
Source code in src/agora_workbench/data_lake/providers.py
aclose()
async
¶
Close embedding resources and the private per-reader SQLite cache.
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
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
¶
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
list(provider, request, context, authorizer)
async
¶
List only artifacts authorized for the current caller.
Source code in src/agora_workbench/data_lake/protocols.py
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
resolve(provider, reference, context, authorizer)
async
¶
Resolve only an authorized canonical artifact.
Source code in src/agora_workbench/data_lake/protocols.py
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
¶
search(request, context)
async
¶
Search artifacts using provider-defined ranking and opaque cursors.
list(request, context)
async
¶
List artifacts using provider-defined filters and opaque cursors.
get(reference, context)
async
¶
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.
PolicyEnforcedCatalog
¶
Bases: Protocol
Caller-aware catalog surface produced by policy composition.
policy_mode
property
¶
Return the configured enforcement granularity.
capabilities(context)
async
¶
search(request, context)
async
¶
Search within the current caller's effective authorization scope.
list(request, context)
async
¶
List within the current caller's effective authorization scope.
get(reference, context)
async
¶
resolve(reference, context)
async
¶
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
capabilities(context)
async
¶
Return write operations allowed for this session and source.
Source code in src/agora_workbench/data_lake/managed.py
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
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
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
artifact_reference(path, artifact_id=None)
¶
Return the normalized effective reference used for authorization.
Source code in src/agora_workbench/data_lake/managed.py
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
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
1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 | |
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
promote(request, context=RequestContext())
async
¶
Explicitly copy a scratch output into durable managed storage.
Source code in src/agora_workbench/data_lake/managed.py
remove(request, context=RequestContext())
async
¶
Commit a tombstone before collecting any managed bytes.
Source code in src/agora_workbench/data_lake/managed.py
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
2235 2236 2237 2238 2239 2240 2241 2242 2243 2244 2245 2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256 2257 2258 2259 2260 2261 2262 2263 2264 2265 2266 2267 2268 2269 2270 2271 2272 2273 2274 2275 2276 2277 2278 2279 2280 2281 2282 2283 2284 2285 2286 2287 2288 2289 2290 2291 2292 2293 2294 2295 2296 2297 2298 2299 2300 2301 2302 2303 2304 2305 2306 2307 2308 2309 2310 2311 2312 2313 2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324 2325 2326 2327 2328 2329 2330 2331 2332 2333 2334 2335 2336 2337 2338 2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407 2408 2409 | |
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
¶
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
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
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
canonicalize_azure_uri(uri)
¶
Canonicalize a supported Azure Blob or DFS URI without credentials.
Source code in src/agora_workbench/data_lake/identity.py
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
normalize_logical_path(path)
¶
Return a portable, source-relative POSIX path.
Source code in src/agora_workbench/data_lake/identity.py
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
sanitize_uri_for_display(uri)
¶
Remove credentials, query parameters, and fragments from a URI.
Source code in src/agora_workbench/data_lake/identity.py
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
stable_source_id(source_type, root)
¶
Derive a stable fallback source ID.
Source code in src/agora_workbench/data_lake/identity.py
managed_writer_extension_factory(writer)
¶
Adapt a managed writer to CatalogIntegration.capability_extension_factory.
Source code in src/agora_workbench/data_lake/managed.py
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
¶
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
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
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
artifact_reference(path, artifact_id=None)
¶
Return the normalized effective reference used for authorization.
Source code in src/agora_workbench/data_lake/managed.py
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
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
1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 1105 1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143 1144 1145 1146 1147 1148 1149 1150 1151 1152 1153 1154 1155 1156 1157 1158 1159 1160 1161 1162 1163 1164 1165 1166 1167 1168 1169 1170 1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181 1182 1183 1184 1185 1186 1187 1188 1189 1190 1191 1192 1193 1194 1195 1196 1197 1198 1199 1200 1201 1202 1203 1204 1205 1206 1207 1208 1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 | |
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
promote(request, context=RequestContext())
async
¶
Explicitly copy a scratch output into durable managed storage.
Source code in src/agora_workbench/data_lake/managed.py
remove(request, context=RequestContext())
async
¶
Commit a tombstone before collecting any managed bytes.
Source code in src/agora_workbench/data_lake/managed.py
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
2235 2236 2237 2238 2239 2240 2241 2242 2243 2244 2245 2246 2247 2248 2249 2250 2251 2252 2253 2254 2255 2256 2257 2258 2259 2260 2261 2262 2263 2264 2265 2266 2267 2268 2269 2270 2271 2272 2273 2274 2275 2276 2277 2278 2279 2280 2281 2282 2283 2284 2285 2286 2287 2288 2289 2290 2291 2292 2293 2294 2295 2296 2297 2298 2299 2300 2301 2302 2303 2304 2305 2306 2307 2308 2309 2310 2311 2312 2313 2314 2315 2316 2317 2318 2319 2320 2321 2322 2323 2324 2325 2326 2327 2328 2329 2330 2331 2332 2333 2334 2335 2336 2337 2338 2339 2340 2341 2342 2343 2344 2345 2346 2347 2348 2349 2350 2351 2352 2353 2354 2355 2356 2357 2358 2359 2360 2361 2362 2363 2364 2365 2366 2367 2368 2369 2370 2371 2372 2373 2374 2375 2376 2377 2378 2379 2380 2381 2382 2383 2384 2385 2386 2387 2388 2389 2390 2391 2392 2393 2394 2395 2396 2397 2398 2399 2400 2401 2402 2403 2404 2405 2406 2407 2408 2409 | |
AuthorizedManagedCatalogWriter(writer, authorizer)
¶
Authorize every mutation before metadata or byte side effects.
Source code in src/agora_workbench/data_lake/managed.py
capabilities(context)
async
¶
Return write operations allowed for this session and source.
Source code in src/agora_workbench/data_lake/managed.py
managed_writer_extension_factory(writer)
¶
Adapt a managed writer to CatalogIntegration.capability_extension_factory.
Source code in src/agora_workbench/data_lake/managed.py
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 |
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 |
None
|
Source code in src/agora_workbench/code_execution/data_access/artifact_resolvers.py
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
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. |
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
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
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.
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
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
|
None
|
username
|
str | None
|
Optional username filter. When multiple accounts
exist in the cache, only the matching account is used.
If |
None
|
authority
|
str
|
Azure AD authority URL. Defaults to the
|
'https://login.microsoftonline.com/organizations'
|
Source code in src/agora_workbench/code_execution/data_access/credentials.py
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. |
()
|
Returns:
| Type | Description |
|---|---|
'AccessToken'
|
An |
Raises:
| Type | Description |
|---|---|
CredentialUnavailableError
|
If the cache file is missing,
no accounts are found, or silent acquisition fails.
This allows |
Source code in src/agora_workbench/code_execution/data_access/credentials.py
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 |
None
|
Source code in src/agora_workbench/code_execution/data_access/fetchers.py
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
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
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
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
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
close()
async
¶
can_handle(qualified_name)
¶
Check if this is a blob storage URL.
Source code in src/agora_workbench/code_execution/data_access/fetchers.py
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
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
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
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
|
Source code in src/agora_workbench/code_execution/data_access/fetchers.py
close()
async
¶
__del__()
¶
Defensively release retained descriptors when explicit cleanup is missed.
can_handle(qualified_name)
¶
Check if this is a local filesystem path.
Source code in src/agora_workbench/code_execution/data_access/fetchers.py
fetch(qualified_name)
async
¶
Read a local file into memory.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
qualified_name
|
str
|
Local file path or |
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
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 |
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
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
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
|
extra_fetchers
|
list[AssetFetcher] | None
|
Optional list of additional |
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
|
None
|
artifact_resolver
|
ArtifactResolver | None
|
Optional resolver turning |
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 |
Source code in src/agora_workbench/code_execution/data_access/manager.py
225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 | |
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
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 |
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
388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 | |
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
list_available()
¶
List assets available in this session's cache.
Returns:
| Type | Description |
|---|---|
list[str]
|
List of qualified names currently cached |
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
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
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
__del__()
¶
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 |
None
|
Source code in src/agora_workbench/code_execution/data_access/publishers.py
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. |
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
|
|
str
|
or |
Source code in src/agora_workbench/code_execution/data_access/publishers.py
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
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. |
required |
Returns:
| Type | Description |
|---|---|
bool
|
|
Source code in src/agora_workbench/code_execution/data_access/publishers.py
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
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.
|
required |
container
|
str
|
Container name to upload into. |
required |
credential
|
'AsyncTokenCredential | None'
|
An |
None
|
Source code in src/agora_workbench/code_execution/data_access/publishers.py
close()
async
¶
Close the underlying BlobServiceClient.
Source code in src/agora_workbench/code_execution/data_access/publishers.py
can_handle(destination)
¶
Return True for <blob>…</blob> destinations.
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
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
740 741 742 743 744 745 746 747 748 749 750 751 752 753 754 755 756 757 758 759 760 761 762 763 764 765 766 767 768 769 770 771 772 773 774 775 776 777 778 779 780 781 782 783 784 785 786 787 788 789 790 791 792 793 794 795 796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 | |
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. |
required |
Source code in src/agora_workbench/code_execution/data_access/publishers.py
can_handle(destination)
¶
Return True for <gui>…</gui> destinations.
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
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
close()
async
¶
Close the retained trusted root descriptor.
can_handle(destination)
¶
Return True for <local>…</local> destinations.
Source code in src/agora_workbench/code_execution/data_access/publishers.py
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
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
1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224 1225 1226 1227 1228 1229 1230 1231 1232 1233 1234 1235 1236 1237 1238 1239 1240 1241 1242 1243 1244 1245 1246 1247 1248 1249 1250 1251 1252 1253 1254 1255 1256 1257 1258 1259 1260 1261 1262 1263 1264 1265 1266 1267 1268 1269 | |
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. |
required |
target_url
|
str | None
|
Full base URL of the target MCP server. If |
None
|
timeout
|
float
|
HTTP read/write timeout in seconds for the transfer request. |
60.0
|
trust_http
|
bool
|
When |
False
|
Source code in src/agora_workbench/code_execution/data_access/publishers.py
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
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
1455 1456 1457 1458 1459 1460 1461 1462 1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494 1495 1496 1497 1498 1499 1500 1501 1502 1503 1504 1505 1506 1507 1508 1509 1510 1511 1512 1513 1514 1515 1516 1517 1518 1519 1520 1521 1522 1523 1524 1525 1526 1527 1528 1529 1530 1531 1532 1533 1534 1535 1536 1537 1538 1539 1540 1541 1542 1543 1544 1545 1546 1547 1548 1549 1550 1551 1552 1553 1554 1555 1556 1557 1558 1559 1560 1561 1562 1563 1564 1565 1566 1567 1568 1569 1570 1571 1572 1573 1574 1575 1576 1577 1578 1579 1580 1581 1582 1583 1584 1585 1586 1587 1588 1589 1590 1591 1592 1593 1594 1595 1596 1597 1598 1599 1600 1601 1602 1603 1604 1605 1606 1607 1608 1609 1610 1611 1612 1613 1614 1615 1616 1617 1618 1619 1620 1621 1622 1623 1624 1625 1626 1627 1628 1629 1630 1631 1632 1633 1634 1635 1636 1637 1638 1639 1640 1641 1642 1643 1644 1645 1646 1647 1648 1649 1650 1651 1652 1653 1654 1655 1656 1657 1658 1659 1660 1661 1662 1663 1664 1665 1666 1667 1668 1669 1670 | |
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 |
None
|
Returns:
| Type | Description |
|---|---|
AsyncTokenCredential
|
An |
AsyncTokenCredential
|
|
Source code in src/agora_workbench/code_execution/data_access/credentials.py
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
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
looks_like_qualified_name(value)
¶
Check if a string looks like a DataLake qualified name.
Recognizes type-tagged artifact format:
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
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 |