From 5ea0d5978a33523144683381d9d85abdc330b426 Mon Sep 17 00:00:00 2001 From: "E.Gavrilov" Date: Fri, 22 May 2026 02:10:33 +0300 Subject: [PATCH] fix(task-020): support s3 assets in publishing dry-run and commit --- apps/backend/src/application/publishing.py | 60 +++++++++++++++++-- .../src/infrastructure/object_storage.py | 17 ++++++ .../test_publishing_git_flow_public_api.py | 35 +++++++++++ ...demo-publishing-s3-assets-compatibility.md | 55 +++++++++++++++++ 4 files changed, 161 insertions(+), 6 deletions(-) create mode 100644 tasks/020-demo-publishing-s3-assets-compatibility.md diff --git a/apps/backend/src/application/publishing.py b/apps/backend/src/application/publishing.py index 4b52d45..8adbc09 100644 --- a/apps/backend/src/application/publishing.py +++ b/apps/backend/src/application/publishing.py @@ -11,6 +11,7 @@ from typing import Any from urllib.parse import urlparse from uuid import UUID +from src.infrastructure.object_storage import open_object_storage_client from src.domain.contracts import ( ArticleSummary, ArticleWorkflowStatus, @@ -266,11 +267,13 @@ def _load_bundle_context( published_at=_now().isoformat(), ) markdown_with_frontmatter = _compose_markdown(frontmatter, draft.body_markdown) + object_storage = open_object_storage_client() assets = _collect_assets( repository, article_id=article.id, slug=slug, asset_path_template=target_site.publishing_rules.asset_path_template, + object_storage=object_storage, ) return { @@ -328,6 +331,7 @@ def _collect_assets( article_id: UUID, slug: str, asset_path_template: str, + object_storage: object, ) -> list[dict[str, Any]]: assets = repository.assets.list_for_article(article_id) output: list[dict[str, Any]] = [] @@ -336,14 +340,18 @@ def _collect_assets( continue if not asset.file_url: continue - source_path = _local_path_from_file_url(asset.file_url) - if source_path is None or not source_path.exists(): + source_file_name = _asset_filename(asset) + source_content = _load_asset_bytes( + asset=asset, + object_storage=object_storage, + ) + if source_content is None: raise PermissionError(f"Approved asset file is unavailable: {asset.file_url}") target_path = _apply_template( asset_path_template, { "slug": slug, - "filename": source_path.name, + "filename": source_file_name, "article_id": str(article_id), "asset_id": str(asset.id), }, @@ -353,7 +361,7 @@ def _collect_assets( "asset_id": str(asset.id), "asset_type": asset.asset_type.value, "source_url": asset.file_url, - "source_path": str(source_path), + "source_content": source_content, "target_path": target_path, } ) @@ -381,10 +389,19 @@ def _create_and_push_commit(bundle_context: dict[str, Any]) -> str: str(script_config["transform_script"]), encoding="utf-8", ) + transform_assets = [ + { + "asset_id": asset["asset_id"], + "asset_type": asset["asset_type"], + "source_url": asset["source_url"], + "target_path": asset["target_path"], + } + for asset in bundle_context["assets"] + ] transform_input = { "frontmatter": bundle_context["frontmatter"], "body": bundle_context["markdown_body"], - "assets": bundle_context["assets"], + "assets": transform_assets, "content_path": bundle_context["content_rel_path"], } (workspace / _TRANSFORM_INPUT_PATH).write_text( @@ -413,7 +430,7 @@ def _create_and_push_commit(bundle_context: dict[str, Any]) -> str: for asset in bundle_context["assets"]: target_path = workspace / str(asset["target_path"]) target_path.parent.mkdir(parents=True, exist_ok=True) - shutil.copyfile(str(asset["source_path"]), target_path) + target_path.write_bytes(bytes(asset["source_content"])) shutil.rmtree(runtime_dir, ignore_errors=True) _run_cmd( @@ -668,6 +685,37 @@ def _local_path_from_file_url(file_url: str) -> Path | None: return Path(parsed.path) +def _asset_filename(asset: object) -> str: + source_path = _local_path_from_file_url(_string_value(getattr(asset, "file_url", None))) + if source_path is not None: + return source_path.name + object_key = _string_value(getattr(asset, "object_key", None)) + if object_key: + return Path(object_key).name + return "asset.bin" + + +def _load_asset_bytes(*, asset: object, object_storage: object) -> bytes | None: + file_url = _string_value(getattr(asset, "file_url", None)) + local_path = _local_path_from_file_url(file_url) + if local_path is not None: + if not local_path.exists(): + return None + return local_path.read_bytes() + + object_key = _string_value(getattr(asset, "object_key", None)) + parsed = urlparse(file_url) + if parsed.scheme == "s3" and not object_key: + object_key = parsed.path.lstrip("/") + if not object_key: + return None + + try: + return object_storage.get_bytes(object_key=object_key) + except Exception: + return None + + def _require_final_approval(repository: object, *, article_id: UUID) -> object: events = repository.articles.list_workflow_events(article_id) for event in reversed(events): diff --git a/apps/backend/src/infrastructure/object_storage.py b/apps/backend/src/infrastructure/object_storage.py index b7330fc..3ac1270 100644 --- a/apps/backend/src/infrastructure/object_storage.py +++ b/apps/backend/src/infrastructure/object_storage.py @@ -8,6 +8,9 @@ class ObjectStorageClient: def put_text(self, *, object_key: str, content: str) -> str: raise NotImplementedError + def get_bytes(self, *, object_key: str) -> bytes: + raise NotImplementedError + def put_bytes( self, *, @@ -41,6 +44,10 @@ class LocalObjectStorageClient(ObjectStorageClient): path.write_bytes(content) return f"file://{path}" + def get_bytes(self, *, object_key: str) -> bytes: + path = self.root / object_key + return path.read_bytes() + class S3ObjectStorageClient(ObjectStorageClient): def __init__(self) -> None: @@ -64,6 +71,16 @@ class S3ObjectStorageClient(ObjectStorageClient): content_type="application/json; charset=utf-8", ) + def get_bytes(self, *, object_key: str) -> bytes: + response = self.client.get_object( + Bucket=self.bucket, + Key=object_key, + ) + body = response.get("Body") + if body is None: + raise RuntimeError(f"Object storage response missing Body for key: {object_key}") + return body.read() + def put_bytes( self, *, diff --git a/apps/backend/tests/integration/test_publishing_git_flow_public_api.py b/apps/backend/tests/integration/test_publishing_git_flow_public_api.py index d110ac2..10e5a52 100644 --- a/apps/backend/tests/integration/test_publishing_git_flow_public_api.py +++ b/apps/backend/tests/integration/test_publishing_git_flow_public_api.py @@ -135,6 +135,20 @@ class PublishingGitFlowPublicApiTest(unittest.TestCase): self.assertEqual(409, commit_response.status_code, commit_response.text) self.assertIn("non-fast-forward", commit_response.text.lower()) + def test_publishing_dry_run_supports_s3_asset_urls_with_object_keys(self) -> None: + article_id, _ = self._prepare_final_approved_article() + self._rewrite_asset_urls_to_s3(article_id) + + dry_run_response = self.client.post( + f"/api/articles/{article_id}/publishing/dry-run", + headers={DEMO_USER_EMAIL_HEADER: DEMO_EDITOR_EMAIL}, + ) + self.assertEqual(201, dry_run_response.status_code, dry_run_response.text) + body = dry_run_response.json() + self.assertTrue(body["content_shape_valid"]) + self.assertEqual("PUBLISH_COMMIT_READY", body["article"]["status"]) + self.assertEqual("PUBLISH_COMMIT_READY", body["article"]["publishing_status"]) + def test_status_and_commits_endpoints_expose_required_metadata(self) -> None: article_id, _ = self._prepare_final_approved_article() @@ -202,6 +216,27 @@ class PublishingGitFlowPublicApiTest(unittest.TestCase): self._run_git(["remote", "add", "origin", str(self.bare_repo_path)], cwd=self.seed_repo_path) self._run_git(["push", "origin", "main"], cwd=self.seed_repo_path) + def _rewrite_asset_urls_to_s3(self, article_id: str) -> None: + assets = self.repository.assets.list_for_article(UUID(article_id)) + now = self.repository.articles.get(UUID(article_id)).updated_at + for asset in assets: + if asset.status.value != "APPROVED": + continue + self.assertTrue(asset.object_key) + self.repository.assets.update( + article_id=UUID(article_id), + asset_id=asset.id, + section_id=asset.section_id, + title=asset.title, + prompt=asset.prompt, + object_key=asset.object_key, + file_url=f"s3://pipeline-local/{asset.object_key}", + alt_text=asset.alt_text, + caption=asset.caption, + status=asset.status, + updated_at=now, + ) + def _run_git(self, args: list[str], *, cwd: Path | None = None) -> None: subprocess.run( ["git", *args], diff --git a/tasks/020-demo-publishing-s3-assets-compatibility.md b/tasks/020-demo-publishing-s3-assets-compatibility.md new file mode 100644 index 0000000..5d3ae0b --- /dev/null +++ b/tasks/020-demo-publishing-s3-assets-compatibility.md @@ -0,0 +1,55 @@ +# Task 020: Demo Publishing S3 Assets Compatibility + +Development description: Fix publish dry-run/commit flow so approved assets stored in object storage (`s3://...`) are correctly bundled for Git publishing in demo Compose mode. + +## Implementation Details + +- Root cause: + - Publishing bundle loader accepted only `file://` asset URLs. + - In Compose demo mode, approved assets are stored in MinIO and exposed as `s3://bucket/key`. +- Required backend changes: + - Extend object storage client with read capability (`get_bytes`). + - Update publishing asset collector to support both `file://` and `s3://` sources. + - Use `asset.object_key` (or URL-derived key fallback) for object-storage fetch. + - Preserve existing manifest structure and publish workflow statuses. +- Validation: + - Add regression test proving dry-run works when approved assets use `s3://` URLs with valid object keys. + +## Public Interface + +- No API contract changes. +- Existing endpoints must behave identically, except they no longer fail on `s3://` approved assets: + - `POST /api/articles/{article_id}/publishing/dry-run` + - `POST /api/articles/{article_id}/publishing/create-commit` + +## Acceptance Criteria + +- [x] Dry-run no longer fails with `Approved asset file is unavailable: s3://...` when asset object exists. +- [x] Publish bundle copies approved assets from object storage into target repo workspace. +- [x] Existing file-based publishing tests continue to pass. +- [x] New regression test covers `s3://` asset URL compatibility. + +## Verification + +- Run publishing integration suite: + - `PYTHONPATH=/private/tmp/pupline-backend-deps python3 -m unittest apps/backend/tests/integration/test_publishing_git_flow_public_api.py` +- Run task 019 demo smoke: + - `PYTHONPATH=/private/tmp/pupline-backend-deps python3 -m unittest apps/backend/tests/integration/test_end_to_end_demo_stack_smoke_public_api.py` +- Run clean compose checklist manually (`down -v` -> `up --build`) and verify publish dry-run/commit steps. + +## Result + +- Status: Implemented. +- Files changed: + - `apps/backend/src/infrastructure/object_storage.py` + - `apps/backend/src/application/publishing.py` + - `apps/backend/tests/integration/test_publishing_git_flow_public_api.py` +- Green evidence: + - `PYTHONPATH=/private/tmp/pupline-backend-deps python3 -m unittest apps/backend/tests/integration/test_publishing_git_flow_public_api.py` -> `Ran 7 tests ... OK` + - `PYTHONPATH=/private/tmp/pupline-backend-deps python3 -m unittest apps/backend/tests/integration/test_end_to_end_demo_stack_smoke_public_api.py` -> `Ran 1 test ... OK` + - Clean compose check (`docker compose down -v` -> `docker compose up --build -d`) + manual public API flow: + - final approval: `200 PUBLISH_DRY_RUN_REQUIRED` + - publishing dry-run: `201 PUBLISH_COMMIT_READY` + - create-commit: `201 PUBLISH_COMMIT_CREATED` +- Refactor notes: + - Added storage read path to reuse existing object storage abstractions instead of introducing publish-specific S3 calls.