Cloud Storage#
oxo-flow supports reading and writing workflow inputs and outputs from
cloud object storage, transparently resolving s3:// and gs:// URIs
through its pluggable storage backend system.
Overview#
Workflows can reference remote files using standard URI schemes:
[[rules]]
name = "fetch_data"
input = ["s3://my-bucket/raw/{sample}.fastq.gz"]
output = ["local/{sample}.fastq.gz"]
shell = "cp {input[0]} {output[0]}"
When the pipeline engine encounters an s3:// or gs:// URI, it
detects the remote scheme and logs a warning — the executor does not
yet stage remote files into the workdir or upload outputs back (see
Current Limitations). The storage module is
usable today as a library API: callers can resolve URIs and read,
write, stage, or upload objects programmatically through the
StorageBackend trait.
Prerequisites#
Both backends are feature-gated and are not included by default. Enable them at build time:
The example workflows below illustrate the URI syntax only — remote URIs are not yet staged or uploaded by the executor (see Current Limitations).
AWS S3#
The S3 backend uses the official aws-sdk-s3 Rust SDK with the standard
AWS credential chain. No additional configuration is required beyond
what the AWS SDK normally reads.
Credential Resolution#
The SDK discovers credentials in this order:
- Environment variables (
AWS_ACCESS_KEY_ID,AWS_SECRET_ACCESS_KEY,AWS_SESSION_TOKEN) ~/.aws/credentials(standard AWS config file)- Web identity tokens
- Instance metadata (EC2, ECS)
When using MinIO or LocalStack for testing, set AWS_ENDPOINT_URL to
point to your local S3-compatible service:
export AWS_ENDPOINT_URL=http://localhost:9000
export AWS_ACCESS_KEY_ID=minioadmin
export AWS_SECRET_ACCESS_KEY=minioadmin
Example Workflow#
[workflow]
name = "s3-example"
version = "1.0.0"
[[rules]]
name = "align"
input = ["s3://genomics-bucket/raw/{sample}.fastq.gz"]
output = ["s3://genomics-bucket/aligned/{sample}.bam"]
shell = "bwa mem reference.fa {input[0]} | samtools sort -o {output[0]}"
[rules.resources]
threads = 8
Google Cloud Storage#
The GCS backend uses the GCS XML API with HMAC-SHA1 authentication. HMAC keys can be created in the GCP Console under Cloud Storage → Settings → Interoperability.
Credential Setup#
Set the following environment variables:
For interoperability with tools that use S3-style credentials,
STORAGE_ACCESS_KEY and STORAGE_SECRET_KEY are also accepted.
Example Workflow#
[workflow]
name = "gcs-example"
version = "1.0.0"
[[rules]]
name = "qc"
input = ["gs://my-bucket/raw/{sample}.fastq.gz"]
output = ["gs://my-bucket/qc/{sample}_report.html"]
shell = "fastqc {input[0]} -o {output[0]}"
[rules.resources]
threads = 2
Storage Backend API#
The StorageBackend trait in oxo_flow_core::storage defines the
interface that all backends implement:
| Method | Description |
|---|---|
exists |
Check whether a path exists |
read_to_string |
Read a remote file into a UTF-8 string |
write |
Write bytes to a remote location |
stage |
Download a remote file to a local directory |
upload |
Upload a local file to a remote location |
head |
Object metadata for invalidation: size + content identity (S3 ETag / GCS md5Hash) |
name |
Human-readable backend name for diagnostics |
The StorageResolver maintains a registry of backends keyed by URI
scheme. Custom backends can be registered at runtime:
use oxo_flow_core::storage::{StorageResolver, StorageScheme};
use std::sync::Arc;
let mut resolver = StorageResolver::with_local();
resolver.add_backend(StorageScheme::S3, Arc::new(s3_backend));
Content-addressed invalidation#
Input manifests now record remote inputs alongside local ones. When a rule's
inputs include s3:// / gs:// URIs and a backend is registered for that
scheme, the checkpoint's input manifest stores a remote entry —
(scheme, key, size, etag) — and manifests_match compares them:
- S3: the raw
ETagfromHeadObject. Composite"hash-N"ETags from multipart uploads are recorded verbatim and compared for equality — they are never recomputed locally. Same content re-uploaded with different part boundaries produces a different ETag (a conservative, spurious invalidation). - GCS:
md5Hash(base64) from thex-goog-hashheader — GCS has no native ETag; the md5 hash is the stronger pure content hash.
A same-size remote rewrite with a changed etag invalidates the rule exactly as a local content change does; unchanged etag → skipped. When neither side has an etag, matching degrades to size-only (documented conservative-for-availability fallback). Without a registered backend the entry is skipped with a warning — the run completes and the remaining local entries still invalidate as before. Only exact object references participate: remote globs and directory references are rejected (the same boundary staging enforces).
Remote staging and upload#
When a backend is registered for a scheme, the local executor stages remote inputs before running a rule and uploads remote outputs after it succeeds:
- Inputs download into
.oxo-flow/staged/in/<scheme>/<bucket>/<key>before execution. Downloads are cached against a sidecar metadata file (<file>.meta.jsonholding size + etag): an unchanged object is never re-downloaded, and the cached file keeps its original mtime so the executor's freshness gate keeps working. A changed etag re-downloads atomically (.part→ rename) and the fresh mtime correctly marks the rule stale. - Outputs declared as remote URIs are written locally to
.oxo-flow/staged/out/<scheme>/<bucket>/<key>— reference them in the shell via{output[n]}/{output.name}. After output validation the engine uploads them; an upload failure fails the rule (a declared remote output that did not land is a broken contract). Remote outputs are only "up to date" while the uploaded object still exists (verified with ahead()on every freshness check), so deleting a cloud result re-runs and re-uploads the rule. - Shells see staged paths only through placeholders — the
substitution happens on a copy of the rule, so
{input[n]}renders the staged local path. A raws3://…URI typed directly into the shell text is not rewritten (the engine warns about it). Checkpoint manifests keep recording the original remote URIs, so invalidation stays etag-driven. - Scope — staging is a local-executor feature. Cluster runs
(
BackendDriver) submit scripts unchanged: nodes use their own shared storage.dry-runnever stages or downloads. Remote globs and directory references are rejected (exact object URIs only), remotetemp_outputs are unsupported.
Deleting .oxo-flow/staged/ is safe — a later run re-downloads (and may
re-execute) instead of using the cache.
Current Limitations#
- Feature-gated — Both backends are opt-in at compile time: enable
oxo-flow-cli'ss3-storage/gcs-storagefeatures to register them in the shared run/dry-run storage resolver. The default build includes only the local filesystem backend. Thes3-storagefeature compiles and is tested (unit fixtures + a live MinIO E2E suite, seecrates/oxo-flow-cli/tests/remote_staging.rs). - S3 credentials come from the environment —
AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY/AWS_SESSION_TOKEN(plusAWS_REGION); the SDK's profile-file/IMDS chain lives in aws_config's async loader and is deliberately not loaded — the same env-only contract the GCS backend has. S3-compatible servers (MinIO, LocalStack) additionally needAWS_ENDPOINT_URLandOXO_S3_FORCE_PATH_STYLE=1(path-style addressing). - UTF-8 only —
read_to_stringrequires the content to be valid UTF-8. Binary files should usestageinstead. - Azure Blob Storage — Not yet supported. Contributions are
welcome via the
StorageBackendtrait.