Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions docs/nav/reference/storage-gcs.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
# StorageGCS

::: dotflow.providers.storage_gcs.StorageGCS
25 changes: 25 additions & 0 deletions docs/nav/tutorial/storage-gcs.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
# Storage GCS

Persists task output to Google Cloud Storage. Good for serverless and cloud-native workflows on GCP.

/// note
Requires `pip install dotflow[gcp]`
///

## Example

{* ./docs_src/storage/storage_gcs.py ln[1:34] hl[2,16:22] *}

## Authentication

`StorageGCS` uses Application Default Credentials (ADC):

1. Environment variable: `GOOGLE_APPLICATION_CREDENTIALS` pointing to a service account JSON
2. `gcloud auth application-default login` for local development
3. Service account: automatic on Cloud Run, Cloud Functions, GKE

No credentials are needed in code — the GCP client handles it transparently.

## References

- [StorageGCS](https://dotflow-io.github.io/dotflow/nav/reference/storage-gcs/)
36 changes: 36 additions & 0 deletions docs_src/storage/storage_gcs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
from dotflow import Config, DotFlow, action
from dotflow.providers import StorageGCS


@action
def step_one():
return {"message": "hello from GCS"}


@action
def step_two(previous_context):
print(previous_context.storage)
return "ok"


config = Config(
storage=StorageGCS(
bucket="dotflow-io-bucket",
prefix="workflows/",
project="etl-test",
)
)


def main():
workflow = DotFlow(config=config)

workflow.task.add(step=step_one)
workflow.task.add(step=step_two)
workflow.start()

return workflow


if __name__ == "__main__":
main()
7 changes: 7 additions & 0 deletions dotflow/providers/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
"StorageDefault",
"StorageFile",
"StorageS3",
"StorageGCS",
]


Expand All @@ -21,4 +22,10 @@ def __getattr__(name):
from dotflow.providers.storage_s3 import StorageS3

return StorageS3

if name == "StorageGCS":
from dotflow.providers.storage_gcs import StorageGCS

return StorageGCS

raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
119 changes: 119 additions & 0 deletions dotflow/providers/storage_gcs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,119 @@
"""Storage GCS"""

from collections.abc import Callable
from json import dumps, loads
from typing import Any

from dotflow.abc.storage import Storage
from dotflow.core.context import Context
from dotflow.core.exception import ModuleNotFound


class StorageGCS(Storage):
"""
Import:
You can import the **StorageGCS** class directly from dotflow providers:

from dotflow.providers import StorageGCS

Example:
`class` dotflow.providers.storage_gcs.StorageGCS

from dotflow import Config
from dotflow.providers import StorageGCS

config = Config(
storage=StorageGCS(
bucket="my-dotflow-bucket",
prefix="workflows/",
project="my-gcp-project"
)
)

Args:
bucket (str): GCS bucket name.

prefix (str): Key prefix for all stored objects.

project (str): GCP project ID. Defaults to ADC project.
"""

def __init__(
self,
*args,
bucket: str,
prefix: str = "dotflow/",
project: str = None,
**kwargs,
):
try:
from google.api_core.exceptions import NotFound
from google.cloud import storage as gcs
except ImportError:
raise ModuleNotFound(
module="google-cloud-storage",
library="dotflow[gcp]",
) from None

self._not_found = NotFound
self.client = gcs.Client(project=project)
self.bucket_obj = self.client.bucket(bucket)
self.bucket_obj.reload()
self.prefix = prefix

def post(self, key: str, context: Context) -> None:
task_context = []

if isinstance(context.storage, list):
for item in context.storage:
if isinstance(item, Context):
task_context.append(self._dumps(storage=item.storage))
else:
task_context.append(self._dumps(storage=context.storage))

self._write(key=key, data=task_context)

def get(self, key: str) -> Context:
task_context = self._read(key)

if len(task_context) == 0:
return Context()

if len(task_context) == 1:
return self._loads(storage=task_context[0])

contexts = Context(storage=[])
for context in task_context:
contexts.storage.append(self._loads(storage=context))

return contexts

def key(self, task: Callable):
return f"{task.workflow_id}-{task.task_id}"

def _read(self, key: str) -> list:
blob = self.bucket_obj.blob(f"{self.prefix}{key}")
try:
data = blob.download_as_text()
return loads(data)
except self._not_found:
return []

def _write(self, key: str, data: list) -> None:
blob = self.bucket_obj.blob(f"{self.prefix}{key}")
blob.upload_from_string(
dumps(data),
content_type="application/json",
)

def _loads(self, storage: Any) -> Context:
try:
return Context(storage=loads(storage))
except Exception:
return Context(storage=storage)

def _dumps(self, storage: Any) -> str:
try:
return dumps(storage)
except TypeError:
return str(storage)
2 changes: 2 additions & 0 deletions mkdocs.yml
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,7 @@ nav:
- nav/tutorial/storage-default.md
- nav/tutorial/storage-file.md
- nav/tutorial/storage-s3.md
- nav/tutorial/storage-gcs.md
- nav/tutorial/provider-notify.md
- nav/tutorial/provider-log.md
- Tutorial - User Guide:
Expand Down Expand Up @@ -187,6 +188,7 @@ nav:
- nav/reference/storage-init.md
- nav/reference/storage-file.md
- nav/reference/storage-s3.md
- nav/reference/storage-gcs.md
- Instance:
- nav/reference/task-instance.md
- nav/reference/context-instance.md
Expand Down
Loading
Loading