Skip to content
Roboto
Esc
↑↓navigate↵open⌘Jpreview
On this page

roboto.storage

Remote file storage I/O.

Whole-file transfer (upload transactions, download sessions, credentials, and the object-store abstraction) for moving files in and out of Roboto storage, plus the range-reader, local cache, and sparse-buffer primitives for streaming byte-range reads that the format decoders in roboto.formats build on.

Submodules

Package Contents

AbortTransactionsRequest

class roboto.storage.AbortTransactionsRequest(/, **data)#View Source

Bases: pydantic.BaseModel

Request payload for aborting file upload transactions.

Used to cancel ongoing file upload transactions, typically when uploads fail or are no longer needed. This cleans up any reserved resources and marks associated files as no longer pending.

Parameters

data Any

Attributes

AbortTransactionsRequest.transaction_ids

transaction_ids list[str] #

List of transaction IDs to abort.

BeginSignedUrlUploadRequest

class roboto.storage.BeginSignedUrlUploadRequest(/, **data)#View Source

Bases: pydantic.BaseModel

Request payload to begin a single file upload with a signed URL.

Used for simpler upload scenarios where a pre-signed URL is preferred over temporary credentials. The returned URL can be used directly for uploading the file content.

Parameters

data Any

Attributes

BeginSignedUrlUploadRequest.association

The entity this file will be associated with (e.g., dataset, topic).

BeginSignedUrlUploadRequest.file_path

file_path str #

Destination path for the file within the association.

BeginSignedUrlUploadRequest.file_size

file_size int #

Size of the file in bytes.

BeginSignedUrlUploadRequest.origination

origination str | None = None #

Optional description of the upload source.

BeginSignedUrlUploadResponse

class roboto.storage.BeginSignedUrlUploadResponse(/, **data)#View Source

Bases: pydantic.BaseModel

Response from beginning a single file upload.

Contains the upload ID for completing the transaction and a pre-signed URL that can be used to upload the file content directly.

Parameters

data Any

Attributes

BeginSignedUrlUploadResponse.upload_id

upload_id str #

Unique identifier for this upload transaction.

BeginSignedUrlUploadResponse.upload_url

upload_url str #

Pre-signed URL for uploading the file content.

BeginUploadRequest

class roboto.storage.BeginUploadRequest(/, **data)#View Source

Bases: pydantic.BaseModel

Request payload to begin a batch file upload transaction.

Used to initiate a multi-file upload transaction for any association type (dataset, topic, etc.). Returns a transaction ID and upload mappings that specify where each file should be uploaded.

Parameters

data Any

Attributes

BeginUploadRequest.association

The entity these files will be associated with (e.g., dataset, topic).

BeginUploadRequest.device_id

device_id str | None = None #

Optional identifier of the device that generated this data.

BeginUploadRequest.origination

origination str #

Description of the upload source (e.g., ‘roboto-sdk v1.0.0’).

BeginUploadRequest.resource_manifest

resource_manifest dict[str, int] #

Dictionary mapping destination file paths to file sizes in bytes.

BeginUploadResponse

class roboto.storage.BeginUploadResponse(/, **data)#View Source

Bases: pydantic.BaseModel

Response from beginning a batch upload transaction.

Contains the transaction ID needed for subsequent progress reporting and completion calls, plus mappings from file paths to their upload URIs.

Parameters

data Any

Attributes

BeginUploadResponse.transaction_id

transaction_id str #

Unique identifier for this upload transaction.

BeginUploadResponse.upload_mappings

upload_mappings dict[str, str] #

Dictionary mapping file paths to their S3 upload URIs.

CachePolicy

class roboto.storage.CachePolicy#View Source

Bases: str, enum.Enum

Governs whether a fetched data file is cached to local disk before reading.

The policy applies to formats with a disk-cache path (Parquet today); a format that always streams (MCAP) ignores it.

Attributes

CachePolicy.ADAPTIVE

ADAPTIVE = 'adaptive' #

Reuse an already-cached file; otherwise download when the read projects enough columns (COLUMN_COUNT_LOCAL_CACHE_THRESHOLD) to justify it, and stream over HTTP when it does not.

CachePolicy.ALWAYS

ALWAYS = 'always' #

Download the file to the local cache before reading, regardless of how much of it the read projects.

CachePolicy.NEVER

NEVER = 'never' #

Always stream over HTTP; never write to local disk.

DownloadableFile

class roboto.storage.DownloadableFile#View Source

Bases: TypedDict

A file to be downloaded from the Roboto Platform.

Attributes

DownloadableFile.bucket_name

bucket_name str #

Name of the bucket where the file is stored.

DownloadableFile.destination_path

destination_path pathlib.Path #

Local path where the file should be saved.

DownloadableFile.source_uri

source_uri str #

Full URI of the file in cloud storage (e.g., ‘s3://bucket/key’).

FileService

class roboto.storage.FileService(roboto_client=None, object_store_registry=None)#View Source

Application service for performing upload and download to the Roboto Platform.

Agnostic to object store provider.

Parameters

roboto_client Optional[roboto.http.RobotoClient]
object_store_registry Optional[roboto.storage.object_store.StoreRegistry]

FileService.download()

download(files, association, caller_org_id=None, on_progress=None)#View Source

Download files from the Roboto Platform.

Parameters

files collections.abc.Sequence[roboto.storage.download_session.DownloadableFile]

Sequence of files to download, each with source_uri and destination_path.

Association of the files to download.

caller_org_id Optional[str]

Optional organization ID for cross-org access.

Optional callback to be periodically called with the number of bytes downloaded.

Return type

None

FileService.upload()

upload(files, association, destination_paths={}, batch_size=_DEFAULT_UPLOAD_BATCH_SIZE, device_id=None, caller_org_id=None, on_progress=None)#View Source

Upload the given files and return which file record each one created.

Parameters

files collections.abc.Iterable[pathlib.Path]
destination_paths collections.abc.Mapping[pathlib.Path, str]
batch_size int
device_id Optional[str]
caller_org_id Optional[str]

Returns

dict[pathlib.Path, str]

Mapping from each uploaded local path to the ID of the file record it created.

Raises

ValueError

If two of the given files resolve to the same destination path: their uploads would overwrite each other and only one could appear in the returned mapping. Files without a destination_paths entry are destined for their own basename, so two like-named files from different directories collide unless given distinct destinations.

OSError

If a given file cannot be read.

HttpRangeReader

class roboto.storage.HttpRangeReader(url, read_ahead_size=_READ_AHEAD_SIZE)#View Source

A seekable, buffered byte-range reader backed by an HTTP URL.

Uses HTTP range requests so only the requested byte ranges are fetched, allowing efficient partial access to remote files (e.g., reading just the MCAP summary/index section at the end of a file without downloading the full data payload).

Reads are satisfied from an in-memory sparse cache. HTTP requests are only issued on a cache miss, fetching READ_AHEAD_SIZE bytes at a time. Unlike a simple single-buffer approach, this cache retains all fetched regions, so seeking back to previously-read data doesn’t trigger re-fetches.

This class implements the IO[bytes] protocol methods needed by mcap.reader.

Uses urllib3 connection pooling to reuse HTTP connections across requests, reducing TCP handshake and TLS negotiation overhead.

Parameters

url str
read_ahead_size int

HttpRangeReader.close()

close()#View Source

Close the reader and release resources.

Return type

None

HttpRangeReader.prefetch_range()

prefetch_range(start, end)#View Source

Prefetch a byte range using parallel HTTP requests.

Byte spans already in the cache (e.g., placed there by the footer read-behind at open, which covers the whole file when it is small) are not re-fetched; only the uncovered gaps are requested.

Parameters

start int

Start byte offset (inclusive)

end int

End byte offset (inclusive)

Return type

None

HttpRangeReader.read()

read(size=-1)#View Source

Parameters

size int

Return type

bytes

HttpRangeReader.readable()

readable()#View Source

Return type

bool

HttpRangeReader.seek()

seek(offset, whence=0)#View Source

Parameters

offset int
whence int

Return type

int

HttpRangeReader.seekable()

seekable()#View Source

Return type

bool

Properties

HttpRangeReader.size

size int #

Get the total size of the remote file in bytes.

Return type: int

HttpRangeReader.tell()

tell()#View Source

Return type

int

HttpRangeReader.writable()

writable()#View Source

Return type

bool

ReportUploadProgressRequest

class roboto.storage.ReportUploadProgressRequest(/, **data)#View Source

Bases: pydantic.BaseModel

Request payload for reporting file upload progress.

Used to notify the platform about the completion status of individual files within a batch upload transaction. This enables progress tracking and partial completion handling for large file uploads.

Parameters

data Any

Attributes

ReportUploadProgressRequest.manifest_items

manifest_items list[str] #

List of file URIs that have completed upload.

ReportUploadProgressResponseItem

class roboto.storage.ReportUploadProgressResponseItem(/, **data)#View Source

Bases: pydantic.BaseModel

One file marked available by a progress report.

The service returns one item per reported URI that matched a file of the transaction, and omits URIs that matched none; the omission does not fail the request on the server side. UploadTransaction is stricter with the response it receives: every URI it reports comes from the transaction’s own upload mappings, so a missing pair means client and service disagree about the transaction’s contents, and it raises RobotoInternalException.

Parameters

data Any

Attributes

ReportUploadProgressResponseItem.file_id

file_id str #

Identifier of the file record marked available.

ReportUploadProgressResponseItem.uri

uri str #

Upload URI of the file, matching an entry of the request’s manifest_items.

RobotoCredentials

class roboto.storage.RobotoCredentials(/, **data)#View Source

Bases: pydantic.BaseModel

Credentials returned from the Roboto Platform

Parameters

data Any

Attributes

RobotoCredentials.access_key_id

access_key_id str #

RobotoCredentials.bucket

bucket str #

RobotoCredentials.expiration

expiration datetime.datetime #

RobotoCredentials.is_expired()

is_expired()#View Source

Return type

bool

Attributes

RobotoCredentials.region

region str #

RobotoCredentials.required_prefix

required_prefix str #

RobotoCredentials.secret_access_key

secret_access_key str #

RobotoCredentials.session_token

session_token str #

RobotoCredentials.to_dict()

to_dict()#View Source

Return type

dict[str, Any]

RobotoCredentials.to_object_store_credentials()

to_object_store_credentials()#View Source

SparseBuffer

class roboto.storage.SparseBuffer(file_size)#View Source

A seekable, read-only file-like object backed by sparse in-memory byte regions.

Stores fetched byte regions and provides a standard IO[bytes] interface for reading from them. Regions are automatically merged when they overlap or are adjacent, keeping the internal representation compact.

This is intended to be used as: 1. The cache backend for HttpRangeReader (sparse storage with smart fetching) 2. The stream for mcap.reader.SeekingReader after bulk-fetching byte ranges

Usage

buf = SparseBuffer(file_size=1000)
buf.add_region(0, b"MCAP_MAGIC")  # header
buf.add_region(900, b"footer_data")  # footer
buf.seek(0)
# 0
buf.read(10)
# b'MCAP_MAGIC'

Parameters

file_size int

SparseBuffer.add_region()

add_region(offset, data)#View Source

Store a byte region at the given file offset.

Merges with any overlapping or adjacent existing regions.

Parameters

offset int

Byte offset within the virtual file.

data bytes

Raw bytes to store at that offset.

Return type

None

SparseBuffer.clear()

clear()#View Source

Remove all cached regions.

Return type

None

SparseBuffer.find_region()

find_region(start, size)#View Source

Check if [start, start+size) is fully contained in a cached region.

Parameters

start int

Start byte offset.

size int

Number of bytes.

Returns

bytes | None

The requested bytes if fully cached, None otherwise.

SparseBuffer.read()

read(size=-1)#View Source

Read up to size bytes from the current position.

If the current position is within a cached region, returns available bytes (may be fewer than requested if the region ends before size bytes). If the current position is not in any cached region (a gap), returns b”“.

This allows callers to detect partial hits and fetch missing data: - len(result) == size: fully satisfied - 0 < len(result) < size: partial hit, more data may be needed - len(result) == 0: gap at current position, caller should fetch

Parameters

size int

Maximum number of bytes to read. -1 means read to end of file.

Returns

bytes

Bytes read from cached regions, or b”” if at a gap or past EOF.

SparseBuffer.readable()

readable()#View Source

Return True - this buffer supports reading.

Return type

bool

Properties

SparseBuffer.regions

regions list[tuple[int, int]] #

List of (start, end) byte ranges currently cached.

End is exclusive. Useful for fetch planning and debugging.

Return type: list[tuple[int, int]]

SparseBuffer.seek()

seek(offset, whence=0)#View Source

Move the read position.

Parameters

offset int

Byte offset relative to the position indicated by whence.

whence int

0=SEEK_SET (start), 1=SEEK_CUR (current), 2=SEEK_END (end).

Returns

int

The new absolute position.

SparseBuffer.seekable()

seekable()#View Source

Return True - this buffer supports seeking.

Return type

bool

Properties

SparseBuffer.size

size int #

Total size of the virtual file.

Return type: int

SparseBuffer.tell()

tell()#View Source

Return the current read position.

Return type

int

SparseBuffer.writable()

writable()#View Source

Return False - this buffer is read-only.

Return type

bool

as_io_bytes()

roboto.storage.as_io_bytes(reader)#View Source

Cast an HttpRangeReader to typing.IO[bytes] for type-checking purposes.

Parameters

Return type

IO[bytes]

Was this page helpful?