fix(task-020): support s3 assets in publishing dry-run and commit

This commit is contained in:
2026-05-22 02:10:33 +03:00
parent 7667819c46
commit 5ea0d5978a
4 changed files with 161 additions and 6 deletions
+54 -6
View File
@@ -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):
@@ -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,
*,
@@ -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],