from __future__ import annotations import json import re import shutil import subprocess import tempfile from datetime import UTC, datetime from pathlib import Path 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, PublishCommitCreateResponse, PublishCommitListResponse, PublishCommitSummary, PublishingDryRunResponse, PublishingStatus, PublishingStatusResponse, ) _VALIDATION_LABEL = "Best-effort content-shape validation only." _FINAL_APPROVAL_EVENT = "FINAL_APPROVAL_GRANTED" _RUNTIME_DIR = ".pipeline-runtime" _PUBLISHING_YAML_PATH = f"{_RUNTIME_DIR}/publishing.yaml" _TRANSFORM_SCRIPT_PATH = f"{_RUNTIME_DIR}/transform.mjs" _TRANSFORM_RUNNER_PATH = f"{_RUNTIME_DIR}/run-transform.mjs" _TRANSFORM_INPUT_PATH = f"{_RUNTIME_DIR}/transform-input.json" _TRANSFORM_OUTPUT_PATH = f"{_RUNTIME_DIR}/transform-output.json" def run_publishing_dry_run( repository: object, *, article_id: UUID, actor_user_id: UUID, ) -> PublishingDryRunResponse: article = repository.articles.get(article_id) final_approval_event = _require_final_approval(repository, article_id=article_id) if article.status not in { ArticleWorkflowStatus.PUBLISH_DRY_RUN_REQUIRED, ArticleWorkflowStatus.PUBLISH_COMMIT_READY, }: raise PermissionError("Publishing dry run is allowed only after final approval.") bundle_context = _load_bundle_context( repository, article=article, final_approval_event=final_approval_event, ) validation_errors = _validate_content_shape(bundle_context["markdown_body"]) content_shape_valid = len(validation_errors) == 0 now = _now() status = ( PublishingStatus.PUBLISH_COMMIT_READY if content_shape_valid else PublishingStatus.PUBLISH_DRY_RUN_FAILED ) manifest = _build_manifest( bundle_context=bundle_context, status=status, validation_errors=validation_errors, ) publish_commit = repository.publish_commits.create( article_id=article.id, target_site_id=article.target_site_id, repository_url=bundle_context["repository_url"], branch=bundle_context["branch"], commit_sha=None, content_bundle_manifest=manifest, status=status, deployment_status=None, created_at=now, ) updated = _update_article_for_dry_run_result( repository, article=article, content_shape_valid=content_shape_valid, actor_user_id=actor_user_id, validation_errors=validation_errors, created_at=now, ) return PublishingDryRunResponse( article=updated, publish_commit=publish_commit, content_shape_valid=content_shape_valid, validation_label=_VALIDATION_LABEL, errors=validation_errors, ) def create_publish_commit( repository: object, *, article_id: UUID, actor_user_id: UUID, ) -> PublishCommitCreateResponse: article = repository.articles.get(article_id) _require_final_approval(repository, article_id=article_id) latest_dry_run = repository.publish_commits.latest_for_article_with_statuses( article_id, statuses={PublishingStatus.PUBLISH_COMMIT_READY}, ) if latest_dry_run is None: raise PermissionError("Publish commit requires successful dry run.") if article.status != ArticleWorkflowStatus.PUBLISH_COMMIT_READY: raise PermissionError("Publish commit can run only from PUBLISH_COMMIT_READY status.") manifest = dict(latest_dry_run.content_bundle_manifest or {}) git_info = manifest.get("git") if isinstance(manifest.get("git"), dict) else {} expected_base_head_sha = git_info.get("base_head_sha") repository_url = latest_dry_run.repository_url branch = latest_dry_run.branch current_head_sha = _resolve_remote_branch_head_sha(repository_url, branch) if expected_base_head_sha and current_head_sha != expected_base_head_sha: _record_publish_failure( repository, article=article, actor_user_id=actor_user_id, repository_url=repository_url, branch=branch, previous_manifest=manifest, detail=( "Non-fast-forward detected: remote branch advanced since dry run; " "re-run dry run before creating commit." ), ) raise PermissionError( "Non-fast-forward detected: remote branch advanced since dry run." ) final_approval_event = _require_final_approval(repository, article_id=article_id) bundle_context = _load_bundle_context( repository, article=article, final_approval_event=final_approval_event, ) commit_sha = _create_and_push_commit(bundle_context) now = _now() publish_commit = repository.publish_commits.create( article_id=article.id, target_site_id=article.target_site_id, repository_url=repository_url, branch=branch, commit_sha=commit_sha, content_bundle_manifest=_build_manifest( bundle_context=bundle_context, status=PublishingStatus.PUBLISH_COMMIT_CREATED, validation_errors=[], commit_sha=commit_sha, ), status=PublishingStatus.PUBLISH_COMMIT_CREATED, deployment_status="PENDING", created_at=now, ) updated = repository.articles.update_status( article_id=article.id, status=ArticleWorkflowStatus.PUBLISH_COMMIT_CREATED, updated_at=now, ) updated = repository.articles.update_publishing_status( article_id=article.id, publishing_status=PublishingStatus.PUBLISH_COMMIT_CREATED, updated_at=now, ) repository.articles.create_workflow_event( article_id=article.id, event_type="PUBLISH_COMMIT_CREATED", from_status=article.status, to_status=ArticleWorkflowStatus.PUBLISH_COMMIT_CREATED, actor_user_id=actor_user_id, payload={ "publish_commit_id": str(publish_commit.id), "repository_url": repository_url, "branch": branch, "commit_sha": commit_sha, }, created_at=now, ) return PublishCommitCreateResponse(article=updated, publish_commit=publish_commit) def get_publishing_status( repository: object, *, article_id: UUID, ) -> PublishingStatusResponse: article = repository.articles.get(article_id) latest_dry_run = repository.publish_commits.latest_for_article_with_statuses( article_id, statuses={PublishingStatus.PUBLISH_COMMIT_READY, PublishingStatus.PUBLISH_DRY_RUN_FAILED}, ) latest_publish_commit = repository.publish_commits.latest_for_article_with_statuses( article_id, statuses={PublishingStatus.PUBLISH_COMMIT_CREATED}, ) return PublishingStatusResponse( article=article, latest_dry_run=latest_dry_run, latest_publish_commit=latest_publish_commit, validation_label=_VALIDATION_LABEL, ) def list_publish_commits( repository: object, *, article_id: UUID, ) -> PublishCommitListResponse: repository.articles.get(article_id) return PublishCommitListResponse(commits=repository.publish_commits.list_for_article(article_id)) def _load_bundle_context( repository: object, *, article: ArticleSummary, final_approval_event: object, ) -> dict[str, Any]: target_site = repository.target_sites.get_by_id(article.target_site_id) if target_site.active_script_config_version_id is None: raise PermissionError("Active script config version is required for publishing.") script_config_version = repository.script_config_versions.get_by_id( target_site.active_script_config_version_id ) draft = repository.article_drafts.latest_for_article(article.id) if draft is None: raise PermissionError("Latest draft is required for publishing.") settings_payload = final_approval_event.payload.get("publishing_settings", {}) if not isinstance(settings_payload, dict): settings_payload = {} content_path = _string_value(settings_payload.get("content_path")) if not content_path: raise PermissionError("Final approval publishing settings are missing content_path.") author = _string_value(settings_payload.get("author")) if not author: raise PermissionError("Final approval publishing settings are missing author.") frontmatter_input = ( settings_payload.get("frontmatter") if isinstance(settings_payload.get("frontmatter"), dict) else {} ) slug = _slug_from_content_path(content_path) or _slugify(draft.title or article.working_title or "article") content_rel_path = _apply_template( target_site.publishing_rules.content_path_template, { "slug": slug, "article_id": str(article.id), "content_path": content_path.lstrip("/"), "format": target_site.publishing_rules.content_format, "language": article.language, }, ) frontmatter = _build_frontmatter( frontmatter_mapping=target_site.publishing_rules.frontmatter_mapping, frontmatter_input=frontmatter_input, author=author, draft=draft, 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 { "article": article, "draft": draft, "target_site": target_site, "script_config_version": script_config_version, "slug": slug, "content_rel_path": content_rel_path, "frontmatter": frontmatter, "markdown_body": draft.body_markdown, "markdown_with_frontmatter": markdown_with_frontmatter, "assets": assets, "repository_url": target_site.publishing_rules.repository_url, "branch": target_site.publishing_rules.production_branch, "content_format": target_site.publishing_rules.content_format, } def _build_frontmatter( *, frontmatter_mapping: dict[str, Any], frontmatter_input: dict[str, Any], author: str, draft: object, published_at: str, ) -> dict[str, Any]: source = { **frontmatter_input, "title": frontmatter_input.get("title") or getattr(draft, "title", None), "meta_description": getattr(draft, "meta_description", None), "author": author, "published_at": published_at, } mapped: dict[str, Any] = {} for output_key, source_key in frontmatter_mapping.items(): if not isinstance(output_key, str) or not output_key: continue if not isinstance(source_key, str) or not source_key: continue if source_key in source and source[source_key] is not None: mapped[output_key] = source[source_key] for key, value in frontmatter_input.items(): if key not in mapped: mapped[key] = value if "author" not in mapped: mapped["author"] = author return mapped def _collect_assets( repository: object, *, 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]] = [] for asset in assets: if asset.status.value != "APPROVED": continue if not asset.file_url: continue 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_file_name, "article_id": str(article_id), "asset_id": str(asset.id), }, ) output.append( { "asset_id": str(asset.id), "asset_type": asset.asset_type.value, "source_url": asset.file_url, "source_content": source_content, "target_path": target_path, } ) return output def _create_and_push_commit(bundle_context: dict[str, Any]) -> str: repository_url = bundle_context["repository_url"] branch = bundle_context["branch"] with tempfile.TemporaryDirectory(prefix="publish-commit-") as tmp_dir: workspace = Path(tmp_dir) / "site" _run_cmd( ["git", "clone", "--branch", branch, "--single-branch", repository_url, str(workspace)], cwd=None, error_prefix="GIT_CHECKOUT_FAILED", ) runtime_dir = workspace / _RUNTIME_DIR runtime_dir.mkdir(parents=True, exist_ok=True) script_config = bundle_context["script_config_version"] (workspace / _PUBLISHING_YAML_PATH).write_text( str(script_config["publishing_yaml"]), encoding="utf-8", ) (workspace / _TRANSFORM_SCRIPT_PATH).write_text( 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": transform_assets, "content_path": bundle_context["content_rel_path"], } (workspace / _TRANSFORM_INPUT_PATH).write_text( json.dumps(transform_input, ensure_ascii=True, indent=2), encoding="utf-8", ) (workspace / _TRANSFORM_RUNNER_PATH).write_text( _transform_runner_script(), encoding="utf-8", ) _run_cmd( [ "node", _TRANSFORM_RUNNER_PATH, _TRANSFORM_SCRIPT_PATH, _TRANSFORM_INPUT_PATH, _TRANSFORM_OUTPUT_PATH, ], cwd=workspace, error_prefix="PUBLISH_DRY_RUN_FAILED", ) content_file = workspace / bundle_context["content_rel_path"] content_file.parent.mkdir(parents=True, exist_ok=True) content_file.write_text(bundle_context["markdown_with_frontmatter"], encoding="utf-8") for asset in bundle_context["assets"]: target_path = workspace / str(asset["target_path"]) target_path.parent.mkdir(parents=True, exist_ok=True) target_path.write_bytes(bytes(asset["source_content"])) shutil.rmtree(runtime_dir, ignore_errors=True) _run_cmd( ["git", "config", "user.name", "Pipeline Bot"], cwd=workspace, error_prefix="GIT_COMMIT_FAILED", ) _run_cmd( ["git", "config", "user.email", "pipeline-bot@example.com"], cwd=workspace, error_prefix="GIT_COMMIT_FAILED", ) _run_cmd(["git", "add", "."], cwd=workspace, error_prefix="GIT_COMMIT_FAILED") commit_message = ( f"Publish article {bundle_context['article'].id}: {bundle_context['draft'].title}" ) _run_cmd( ["git", "commit", "--allow-empty", "-m", commit_message], cwd=workspace, error_prefix="GIT_COMMIT_FAILED", ) try: _run_cmd( ["git", "push", "origin", branch], cwd=workspace, error_prefix="GIT_PUSH_NON_FAST_FORWARD", ) except RuntimeError as error: if "non-fast-forward" in str(error) or "[rejected]" in str(error): raise PermissionError( "Non-fast-forward push rejected. Re-run dry run and retry commit." ) from error raise commit_sha = _run_cmd( ["git", "rev-parse", "HEAD"], cwd=workspace, error_prefix="GIT_COMMIT_FAILED", ).strip() return commit_sha def _update_article_for_dry_run_result( repository: object, *, article: ArticleSummary, content_shape_valid: bool, actor_user_id: UUID, validation_errors: list[str], created_at: datetime, ) -> ArticleSummary: if content_shape_valid: updated = repository.articles.update_status( article_id=article.id, status=ArticleWorkflowStatus.PUBLISH_COMMIT_READY, updated_at=created_at, ) updated = repository.articles.update_publishing_status( article_id=article.id, publishing_status=PublishingStatus.PUBLISH_COMMIT_READY, updated_at=created_at, ) repository.articles.create_workflow_event( article_id=article.id, event_type="PUBLISH_DRY_RUN_SUCCEEDED", from_status=article.status, to_status=ArticleWorkflowStatus.PUBLISH_COMMIT_READY, actor_user_id=actor_user_id, payload={"validation_label": _VALIDATION_LABEL}, created_at=created_at, ) return updated updated = repository.articles.update_status( article_id=article.id, status=ArticleWorkflowStatus.PUBLISH_DRY_RUN_REQUIRED, updated_at=created_at, ) updated = repository.articles.update_publishing_status( article_id=article.id, publishing_status=PublishingStatus.PUBLISH_DRY_RUN_FAILED, updated_at=created_at, ) repository.articles.create_workflow_event( article_id=article.id, event_type="PUBLISH_DRY_RUN_FAILED", from_status=article.status, to_status=ArticleWorkflowStatus.PUBLISH_DRY_RUN_REQUIRED, actor_user_id=actor_user_id, payload={ "validation_label": _VALIDATION_LABEL, "errors": validation_errors, }, created_at=created_at, ) return updated def _record_publish_failure( repository: object, *, article: ArticleSummary, actor_user_id: UUID, repository_url: str, branch: str, previous_manifest: dict[str, Any], detail: str, ) -> None: now = _now() failed_manifest = dict(previous_manifest) failed_manifest["failure"] = detail repository.publish_commits.create( article_id=article.id, target_site_id=article.target_site_id, repository_url=repository_url, branch=branch, commit_sha=None, content_bundle_manifest=failed_manifest, status=PublishingStatus.PUBLISH_VERIFICATION_FAILED, deployment_status="FAILED", created_at=now, ) repository.articles.update_publishing_status( article_id=article.id, publishing_status=PublishingStatus.PUBLISH_VERIFICATION_FAILED, updated_at=now, ) repository.articles.create_workflow_event( article_id=article.id, event_type="PUBLISH_COMMIT_FAILED", from_status=article.status, to_status=article.status, actor_user_id=actor_user_id, payload={"detail": detail}, created_at=now, ) def _build_manifest( *, bundle_context: dict[str, Any], status: PublishingStatus, validation_errors: list[str], commit_sha: str | None = None, ) -> dict[str, Any]: base_head_sha = _resolve_remote_branch_head_sha( bundle_context["repository_url"], bundle_context["branch"] ) script_config_version = bundle_context["script_config_version"] return { "status": status.value, "generated_at": _now().isoformat(), "validation": { "label": _VALIDATION_LABEL, "content_shape_valid": len(validation_errors) == 0, "errors": validation_errors, }, "content": { "format": bundle_context["content_format"], "path": bundle_context["content_rel_path"], "slug": bundle_context["slug"], }, "frontmatter": bundle_context["frontmatter"], "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"] ], "config_version": { "version_id": str(script_config_version["id"]), "version": int(script_config_version["version"]), "publishing_yaml_hash": str(script_config_version["publishing_yaml_hash"]), "transform_script_hash": str(script_config_version["transform_script_hash"]), }, "git": { "repository_url": bundle_context["repository_url"], "branch": bundle_context["branch"], "base_head_sha": base_head_sha, "commit_sha": commit_sha, }, "draft_version": int(bundle_context["draft"].version), } def _validate_content_shape(markdown: str) -> list[str]: errors: list[str] = [] if not markdown.strip(): errors.append("Markdown body is empty.") if re.search(r"^#{1,6}\s+\S", markdown, re.MULTILINE) is None: errors.append("Markdown body must include at least one heading.") if markdown.count("```") % 2 != 0: errors.append("Markdown body has unbalanced fenced code blocks.") return errors def _compose_markdown(frontmatter: dict[str, Any], markdown_body: str) -> str: lines = ["---"] for key in sorted(frontmatter): lines.append(f"{key}: {json.dumps(frontmatter[key], ensure_ascii=True)}") lines.extend(["---", "", markdown_body.strip(), ""]) return "\n".join(lines) def _resolve_remote_branch_head_sha(repository_url: str, branch: str) -> str | None: output = _run_cmd( ["git", "ls-remote", repository_url, f"refs/heads/{branch}"], cwd=None, error_prefix="GIT_CHECKOUT_FAILED", ) line = output.strip() if not line: return None return line.split("\t", 1)[0] def _run_cmd( command: list[str], *, cwd: Path | None, error_prefix: str, ) -> str: result = subprocess.run( command, cwd=str(cwd) if cwd is not None else None, capture_output=True, text=True, check=False, ) if result.returncode != 0: message = result.stderr.strip() or result.stdout.strip() or "command failed" raise RuntimeError(f"{error_prefix}: {message}") return result.stdout def _apply_template(template: str, values: dict[str, Any]) -> str: class _SafeValues(dict[str, Any]): def __missing__(self, key: str) -> str: return "{" + key + "}" rendered = template.format_map(_SafeValues(values)) return rendered.lstrip("/") def _local_path_from_file_url(file_url: str) -> Path | None: parsed = urlparse(file_url) if parsed.scheme != "file": return 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): if event.event_type == _FINAL_APPROVAL_EVENT: return event raise PermissionError("Final approval is required before publishing.") def _slug_from_content_path(content_path: str) -> str: value = content_path.strip().strip("/") if not value: return "" tail = value.split("/")[-1] if "." in tail: tail = tail.rsplit(".", 1)[0] return _slugify(tail) def _slugify(value: str) -> str: normalized = re.sub(r"[^a-zA-Z0-9]+", "-", value.strip().lower()).strip("-") return normalized or "article" def _string_value(value: Any) -> str: if isinstance(value, str): return value.strip() return "" def _transform_runner_script() -> str: return """ import { readFileSync, writeFileSync } from "node:fs"; import { pathToFileURL } from "node:url"; const [scriptPath, inputPath, outputPath] = process.argv.slice(2); const scriptUrl = pathToFileURL(scriptPath).href; const moduleRef = await import(scriptUrl + `?t=${Date.now()}`); const transform = typeof moduleRef.transformArticle === "function" ? moduleRef.transformArticle : typeof moduleRef.default === "function" ? moduleRef.default : null; if (!transform) { throw new Error("Transform script must export transformArticle(article)."); } const raw = readFileSync(inputPath, "utf8"); const article = JSON.parse(raw); const transformed = await transform(article); const output = transformed ?? article; writeFileSync(outputPath, JSON.stringify(output, null, 2), "utf8"); """.strip() def _now() -> datetime: return datetime.now(UTC)