#!/usr/bin/env python3 """Read and review an ACK-ready Feishu Base view through the official lark-cli. The only mutation is a bounded draft write to configured Coordinator fields. The adapter never reads the active profile and emits one JSON document only on success. """ from __future__ import annotations import argparse import hashlib import json import math import os import pwd import re import resource import signal import stat import subprocess import sys import tempfile import time from pathlib import Path from typing import Any from approval_payload import approval_payload_hash, review_items from validate_tasks import validate_builtin as validate_task_board from yaml_subset import DuplicateKeyError, YamlSubsetError, load_json_unique, load_yaml_subset, make_unique_pyyaml_loader LEGACY_REQUIRED_FIELDS = ("title", "actual", "expected", "stepsToReproduce", "acceptance", "attachments", "updatedAt") LEGACY_OPTIONAL_FIELDS = ("priority", "fixLogic") LEGACY_FIELD_ORDER = ("title", "actual", "expected", "stepsToReproduce", "acceptance", "priority", "attachments", "updatedAt", "fixLogic") CLARIFIED_REQUIRED_FIELDS = ("title", "details", "problemStatement", "expectedOutcome", "acceptance", "intakeStatus", "ackTaskId", "attachments", "updatedAt") CLARIFIED_FIELD_ORDER = CLARIFIED_REQUIRED_FIELDS SOURCE_FACT_FIELDS = ("title", "actual", "expected", "updatedAt") COORDINATOR_FIELDS = ("steps", "acceptance", "priority") BUG_CONTENT_FIELDS = ("title", "actual", "expected", "fixLogic", *COORDINATOR_FIELDS) MAX_PAGES = 100 MAX_RECORDS = 10_000 PAGE_SIZE = 100 MAX_CLI_STDOUT = 1024 * 1024 MAX_CLI_STDERR = 64 * 1024 MAX_ATTACHMENTS_PER_RECORD = 10 MAX_TOTAL_ATTACHMENTS = 100 MAX_ATTACHMENT_BYTES = 20 * 1024 * 1024 MAX_TOTAL_ATTACHMENT_BYTES = 200 * 1024 * 1024 MAX_ATTACHMENT_BATCH_SECONDS = 300 MAX_DRAFT_BYTES = 64 * 1024 SAFE_VALUE = re.compile(r"^[^\s\x00-\x1f]{1,256}$") PROFILE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$") RECORD_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,255}$") SOURCE_REF = re.compile(r"^feishu-base:sha256:[0-9a-f]{64}$") DRAFT_REVISION = re.compile(r"^sha256:[0-9a-f]{64}$") WORKFLOWS = {"read-only-v1", "reviewed-writeback-v1", "clarified-writeback-v1"} INTAKE_STATUSES = ("待整理", "需补充", "待审核", "已确认", "已导入") TARGET_BASE_FIELDS = ( ("标题", "text"), ("详细描述", "text"), ("附件", "attachment"), ("问题说明", "text"), ("期望效果", "text"), ("验收标准", "text"), ("处理状态", "select"), ("ACK任务ID", "text"), ("更新时间", "updated_at"), ) class IntakeError(Exception): pass def account_home() -> Path: """Return the actual account home, never a caller-controlled HOME value.""" try: home = Path(pwd.getpwuid(os.getuid()).pw_dir).resolve(strict=True) except (KeyError, OSError) as exc: raise IntakeError("cannot resolve current account home") from exc if not home.is_dir(): raise IntakeError("current account home is not a directory") return home def trusted_lark_cli_dirs() -> list[Path]: """Fixed account and system locations; intentionally never consult PATH.""" home = account_home() candidates = ( home / ".local" / "bin", home / ".local" / "share" / "mise" / "shims", home / ".cargo" / "bin", Path("/home/linuxbrew/.linuxbrew/bin"), Path("/usr/local/go/bin"), Path("/usr/local/bin"), Path("/usr/bin"), Path("/bin"), ) result: list[Path] = [] for candidate in candidates: try: resolved = candidate.resolve(strict=True) except OSError: continue if resolved.is_dir() and resolved not in result: result.append(resolved) return result def resolve_lark_cli() -> Path: """Resolve a safe lark-cli from the fixed trusted locations only.""" for directory in trusted_lark_cli_dirs(): candidate = directory / "lark-cli" try: candidate_metadata = os.lstat(candidate) resolved = candidate.resolve(strict=True) metadata = resolved.stat() except OSError: continue if not (stat.S_ISREG(candidate_metadata.st_mode) or stat.S_ISLNK(candidate_metadata.st_mode)): continue if candidate_metadata.st_uid not in {0, os.getuid()}: continue if not stat.S_ISREG(metadata.st_mode) or not os.access(resolved, os.X_OK): continue if metadata.st_uid not in {0, os.getuid()} or stat.S_IMODE(metadata.st_mode) & 0o022: continue if resolved.name == "lark-cli": return resolved if not stat.S_ISLNK(candidate_metadata.st_mode): continue official_binary = official_npm_binary(resolved) if official_binary is not None: return official_binary raise IntakeError("lark-cli is not installed in a trusted account or system directory") def official_npm_binary(path: Path) -> Path | None: """Resolve a validated npm wrapper to its downloaded native CLI binary.""" if path.name != "run.js" or path.parent.name != "scripts": return None manifest = path.parent.parent / "package.json" try: metadata = manifest.stat() if not stat.S_ISREG(metadata.st_mode): return None if metadata.st_uid not in {0, os.getuid()} or stat.S_IMODE(metadata.st_mode) & 0o022: return None package = json.loads(manifest.read_text(encoding="utf-8")) except (OSError, UnicodeError, json.JSONDecodeError): return None if not ( isinstance(package, dict) and package.get("name") == "@larksuite/cli" and isinstance(package.get("bin"), dict) and package["bin"].get("lark-cli") == "scripts/run.js" ): return None native = path.parent.parent / "bin" / "lark-cli" try: native_lstat = os.lstat(native) resolved = native.resolve(strict=True) metadata = resolved.stat() except OSError: return None if not stat.S_ISREG(native_lstat.st_mode) or resolved.name != "lark-cli": return None if not stat.S_ISREG(metadata.st_mode) or not os.access(resolved, os.X_OK): return None if metadata.st_uid not in {0, os.getuid()} or stat.S_IMODE(metadata.st_mode) & 0o022: return None return resolved def cli_environment() -> dict[str, str]: """Build a minimal environment so env credentials cannot override `--profile`.""" environment = { "HOME": str(account_home()), "PATH": os.pathsep.join(str(path) for path in trusted_lark_cli_dirs()), } for name in ("LANG", "LC_ALL", "LC_CTYPE"): value = os.environ.get(name) if value and "\x00" not in value and len(value) <= 256: environment[name] = value return environment def load_board(path: Path) -> dict[str, Any]: try: text = path.read_text(encoding="utf-8") if path.suffix.lower() == ".json": value = load_json_unique(text) else: try: import yaml # type: ignore value = yaml.load(text, Loader=make_unique_pyyaml_loader(yaml)) except ImportError: value = load_yaml_subset(text) except (OSError, json.JSONDecodeError, DuplicateKeyError, YamlSubsetError) as exc: raise IntakeError(f"cannot read task board: {exc}") from exc except Exception as exc: # PyYAML errors are intentionally not exposed verbatim. raise IntakeError("cannot parse task board") from exc if not isinstance(value, dict): raise IntakeError("task board must be an object") return value def load_draft(path: Path, workflow: str = "reviewed-writeback-v1") -> dict[str, Any]: """Load one bounded, regular JSON file with the two writable draft fields.""" descriptor: int | None = None try: before = path.lstat() if not stat.S_ISREG(before.st_mode) or path.is_symlink(): raise IntakeError("draft input must be a regular file") flags = os.O_RDONLY | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_NOFOLLOW", 0) descriptor = os.open(path, flags) metadata = os.fstat(descriptor) if ( not stat.S_ISREG(metadata.st_mode) or (before.st_dev, before.st_ino) != (metadata.st_dev, metadata.st_ino) ): raise IntakeError("draft input changed while opening") if metadata.st_size <= 0 or metadata.st_size > MAX_DRAFT_BYTES: raise IntakeError("draft input size is invalid") chunks: list[bytes] = [] total = 0 while total <= MAX_DRAFT_BYTES: chunk = os.read(descriptor, min(64 * 1024, MAX_DRAFT_BYTES + 1 - total)) if not chunk: break chunks.append(chunk) total += len(chunk) content = b"".join(chunks) if len(content) != metadata.st_size: raise IntakeError("draft input changed while reading") value = load_json_unique(content.decode("utf-8")) except IntakeError: raise except (OSError, UnicodeError, json.JSONDecodeError, DuplicateKeyError) as exc: raise IntakeError("cannot read draft input") from exc finally: if descriptor is not None: os.close(descriptor) if workflow == "clarified-writeback-v1": required = {"problemStatement", "expectedOutcome", "acceptance"} if not isinstance(value, dict) or set(value) != required: raise IntakeError("draft input must contain exactly problemStatement, expectedOutcome and acceptance") for field in ("problemStatement", "expectedOutcome"): if not isinstance(value[field], str) or not value[field].strip(): raise IntakeError(f"draft {field} must be a non-empty string") acceptance = value["acceptance"] if not isinstance(acceptance, list) or not acceptance or any( not isinstance(item, str) or not item.strip() for item in acceptance ): raise IntakeError("draft acceptance must be a non-empty string list") return { "problemStatement": value["problemStatement"].strip(), "expectedOutcome": value["expectedOutcome"].strip(), "acceptance": [item.strip() for item in acceptance], } if not isinstance(value, dict) or set(value) != {"fixLogic", "acceptance"}: raise IntakeError("draft input must contain exactly fixLogic and acceptance") fix_logic = value["fixLogic"] acceptance = value["acceptance"] if not isinstance(fix_logic, str) or not fix_logic.strip(): raise IntakeError("draft fixLogic must be a non-empty string") if ( not isinstance(acceptance, list) or not acceptance or any(not isinstance(item, str) or not item.strip() for item in acceptance) ): raise IntakeError("draft acceptance must be a non-empty string list") return { "fixLogic": fix_logic.strip(), "acceptance": [item.strip() for item in acceptance], } def config_from_board(board: dict[str, Any]) -> dict[str, Any]: project = board.get("project") if not isinstance(project, dict) or "bugIntake" not in project: raise IntakeError("project.bugIntake is not configured") config = project["bugIntake"] if not isinstance(config, dict): raise IntakeError("project.bugIntake must be an object") allowed = {"provider", "workflow", "profile", "baseToken", "tableId", "viewId", "fields"} unknown = sorted(set(config) - allowed) if unknown: raise IntakeError("bugIntake has unknown fields") if config.get("provider") != "feishu-base": raise IntakeError("bugIntake.provider must be feishu-base") workflow = config.get("workflow", "read-only-v1") if workflow not in WORKFLOWS: raise IntakeError("bugIntake.workflow is invalid") profile = config.get("profile") if not isinstance(profile, str) or not PROFILE.fullmatch(profile): raise IntakeError("bugIntake.profile is invalid") for key in ("baseToken", "tableId", "viewId"): value = config.get(key) if not isinstance(value, str) or not SAFE_VALUE.fullmatch(value): raise IntakeError(f"bugIntake.{key} is invalid") fields = config.get("fields") if workflow == "clarified-writeback-v1": required = CLARIFIED_REQUIRED_FIELDS supported = set(required) else: required = LEGACY_REQUIRED_FIELDS supported = set(required) | set(LEGACY_OPTIONAL_FIELDS) if ( not isinstance(fields, dict) or not set(required).issubset(fields) or not set(fields).issubset(supported) ): raise IntakeError("bugIntake.fields must map all required and only supported logical fields") if any(not isinstance(value, str) or not SAFE_VALUE.fullmatch(value) for value in fields.values()): raise IntakeError("bugIntake.fields values are invalid") if len(set(fields.values())) != len(fields): raise IntakeError("bugIntake.fields values must be unique") if workflow == "reviewed-writeback-v1": missing_review_fields = {"fixLogic", "priority"} - set(fields) if missing_review_fields: raise IntakeError( "reviewed writeback requires bugIntake.fields.fixLogic and priority" ) return config def limit_child_file_size(limit: int) -> None: """Bound regular-file writes by the CLI and any child spawned by its wrapper.""" _, hard = resource.getrlimit(resource.RLIMIT_FSIZE) bounded = limit if hard == resource.RLIM_INFINITY else min(limit, hard) resource.setrlimit(resource.RLIMIT_FSIZE, (bounded, bounded)) def run_cli( args: list[str], *, allow_profile_list: bool = False, max_file_bytes: int = 0, timeout: float = 60, cwd: Path | None = None, ) -> Any: """Run the official CLI and accept only explicit successful JSON shapes.""" executable = resolve_lark_cli() file_limit = max(MAX_CLI_STDOUT, MAX_CLI_STDERR, max_file_bytes) with tempfile.TemporaryFile() as stdout_file, tempfile.TemporaryFile() as stderr_file: try: process = subprocess.Popen( [str(executable), *args], shell=False, stdin=subprocess.DEVNULL, stdout=stdout_file, stderr=stderr_file, env=cli_environment(), cwd=cwd, start_new_session=True, preexec_fn=lambda: limit_child_file_size(file_limit), ) try: returncode = process.wait(timeout=timeout) except subprocess.TimeoutExpired as exc: try: os.killpg(process.pid, signal.SIGKILL) except ProcessLookupError: pass process.wait() raise IntakeError("lark-cli failed to execute") from exc except (OSError, subprocess.SubprocessError) as exc: raise IntakeError("lark-cli failed to execute") from exc stdout_size = os.fstat(stdout_file.fileno()).st_size stderr_size = os.fstat(stderr_file.fileno()).st_size if stdout_size > MAX_CLI_STDOUT or stderr_size > MAX_CLI_STDERR: raise IntakeError("lark-cli output exceeded the safety limit") if returncode: raise IntakeError("lark-cli command failed") stdout_file.seek(0) try: stdout = stdout_file.read(MAX_CLI_STDOUT + 1).decode("utf-8") except UnicodeDecodeError as exc: raise IntakeError("lark-cli returned malformed JSON") from exc try: value = json.loads( stdout, parse_constant=lambda value: (_ for _ in ()).throw( ValueError(f"non-finite JSON constant: {value}") ), ) except (json.JSONDecodeError, ValueError) as exc: raise IntakeError("lark-cli returned malformed JSON") from exc if isinstance(value, list): if allow_profile_list: return value raise IntakeError("lark-cli returned an unexpected JSON array") if not isinstance(value, dict): raise IntakeError("lark-cli returned an invalid JSON response") if "ok" in value and value["ok"] is not True: raise IntakeError("lark-cli returned an error response") if "code" in value and (not isinstance(value["code"], int) or isinstance(value["code"], bool) or value["code"] != 0): raise IntakeError("lark-cli returned an error response") if "ok" not in value and "code" not in value: raise IntakeError("lark-cli returned an ambiguous JSON response") return value def profile_check(config: dict[str, Any]) -> None: # `profile list` is the official non-mutating profile inspection command. value = run_cli(["profile", "list"], allow_profile_list=True) if isinstance(value, list): profiles = value else: # Compatibility wrapper: only a successful envelope with a direct list # is accepted. Do not loosen this into arbitrary nested objects. profiles = value.get("data") if not isinstance(profiles, list): raise IntakeError("profile check returned an invalid response") matching_profile: dict[str, Any] | None = None for item in profiles: # Match the official profile-list item shape. `user` and # `tokenStatus` are optional and deliberately never propagated. if ( not isinstance(item, dict) or not isinstance(item.get("name"), str) or not PROFILE.fullmatch(item["name"]) or not isinstance(item.get("appId"), str) or not isinstance(item.get("brand"), str) or not isinstance(item.get("active"), bool) ): raise IntakeError("profile check returned an invalid profile entry") if item["name"] == config["profile"]: matching_profile = item if matching_profile is None: raise IntakeError("configured lark-cli profile does not exist") if matching_profile["brand"] != "feishu": raise IntakeError("configured lark-cli profile must use the feishu brand") def text(value: Any) -> str: if value is None: return "" if isinstance(value, str): return " ".join(value.split()) if isinstance(value, float) and not math.isfinite(value): raise IntakeError("text field contains a non-finite number") if isinstance(value, (int, float, bool)): return str(value) if isinstance(value, list): return "\n".join(part for part in (text(item) for item in value) if part) if isinstance(value, dict): for key in ("text", "name", "value"): if key in value: return text(value[key]) raise IntakeError("text field contains an unsupported object") raise IntakeError("text field contains an unsupported value") def attachment_items(value: Any) -> list[tuple[dict[str, Any], str]]: if value in (None, ""): return [] if not isinstance(value, list): raise IntakeError("attachments cell must be a list") if len(value) > MAX_ATTACHMENTS_PER_RECORD: raise IntakeError("record exceeded the attachment count limit") attachments: list[tuple[dict[str, Any], str]] = [] for item in value: if not isinstance(item, dict): raise IntakeError("attachment metadata must be an object") token = item.get("file_token", item.get("token")) if not isinstance(token, str) or not SAFE_VALUE.fullmatch(token): raise IntakeError("attachment token is invalid") metadata = {"name": text(item.get("name")), "type": text(item.get("type", item.get("mime_type"))), "size": item.get("size")} if ( not isinstance(metadata["size"], int) or isinstance(metadata["size"], bool) or metadata["size"] < 0 or metadata["size"] > MAX_ATTACHMENT_BYTES ): raise IntakeError("attachment size is invalid") attachments.append((metadata, token)) return attachments def matrix_from_response(response: dict[str, Any], field_ids: list[str]) -> tuple[list[str], list[list[Any]]]: data = response.get("data", response) if not isinstance(data, dict): raise IntakeError("record list data is invalid") fields = data.get("fields") ids = data.get("record_id_list", data.get("recordIds")) rows = data.get("data", data.get("records", data.get("items", data.get("rows")))) if not isinstance(fields, list) or not all(isinstance(x, str) for x in fields): raise IntakeError("record list fields are invalid") if len(set(fields)) != len(fields) or len(set(field_ids)) != len(field_ids): raise IntakeError("record list field projection contains duplicates") if set(fields) != set(field_ids): raise IntakeError( "record list fields do not match configured projection: " f"expected={field_ids!r}, actual={fields!r}" ) if not isinstance(ids, list) or not all(isinstance(x, str) and RECORD_ID.fullmatch(x) for x in ids): raise IntakeError("record list record_id_list is invalid") if not isinstance(rows, list) or len(rows) != len(ids) or any(not isinstance(row, list) or len(row) != len(fields) for row in rows): raise IntakeError("record list matrix does not match fields and record_id_list") if fields == field_ids: return ids, rows positions = {field: index for index, field in enumerate(fields)} return ids, [ [row[positions[field_id]] for field_id in field_ids] for row in rows ] def fetch_pages(config: dict[str, Any]) -> list[tuple[str, list[Any]]]: field_order = ( CLARIFIED_FIELD_ORDER if config.get("workflow") == "clarified-writeback-v1" else LEGACY_FIELD_ORDER ) logical_fields = [ logical for logical in field_order if logical in config["fields"] ] field_ids = [config["fields"][logical] for logical in logical_fields] all_rows: list[tuple[str, list[Any]]] = [] offset = 0 for _ in range(MAX_PAGES): args = ["base", "+record-list", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--view-id", config["viewId"], "--format", "json", "--offset", str(offset), "--limit", str(PAGE_SIZE)] for field_id in field_ids: args.extend(["--field-id", field_id]) response = run_cli(args) ids, rows = matrix_from_response(response, field_ids) if len(ids) > PAGE_SIZE: raise IntakeError("record list exceeded requested page size") all_rows.extend(zip(ids, rows)) if len(all_rows) > MAX_RECORDS: raise IntakeError("record list exceeded record limit") data = response.get("data", response) has_more = data.get("has_more", data.get("hasMore", False)) if not isinstance(has_more, bool): raise IntakeError("record list pagination marker is invalid") if not has_more: return all_rows if not ids: raise IntakeError("record list pagination made no progress") offset += len(ids) raise IntakeError("record list exceeded page limit") def download( config: dict[str, Any], record_id: str, file_token: str, output_dir: Path, expected_size: int, timeout: float, ) -> str: try: output_dir.mkdir(parents=True, exist_ok=False) except OSError as exc: raise IntakeError("attachment output directory is unsafe") from exc if not output_dir.is_dir() or output_dir.is_symlink(): raise IntakeError("attachment output directory is unsafe") run_cli( ["base", "+record-download-attachment", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--record-id", record_id, "--file-token", file_token, "--output", output_dir.name], max_file_bytes=expected_size, timeout=timeout, cwd=output_dir.parent, ) try: created = list(output_dir.iterdir()) except OSError as exc: raise IntakeError("attachment download output is unreadable") from exc if len(created) != 1 or not created[0].is_file() or created[0].is_symlink(): raise IntakeError("attachment download did not produce one safe file") root = output_dir.resolve() resolved = created[0].resolve() if root not in resolved.parents: raise IntakeError("attachment download escaped output directory") if resolved.stat().st_size != expected_size: raise IntakeError("attachment download size did not match metadata") return str(resolved) def source_ref(config: dict[str, Any], record_id: str) -> str: """Return a stable opaque identity without serializing configured identifiers.""" identity = "\x1f".join(("ack-feishu-base-source-ref-v1", config["profile"], config["baseToken"], config["tableId"], record_id)) return f"feishu-base:sha256:{hashlib.sha256(identity.encode('utf-8')).hexdigest()}" def draft_revision( record: dict[str, Any], attachment_tokens: list[str] | None = None, ) -> str: """Bind approval to the normalized source facts and review-controlled fields.""" tokens = attachment_tokens or [] if len(tokens) != len(record["attachments"]): raise IntakeError("draft revision attachment identity is incomplete") stable = { key: record[key] for key in record if key not in { "attachments", "warnings", "enrichmentRequired", "draftRevision", "recordId", "intakeStatus", "ackTaskId", } } stable["attachments"] = [ { **{key: attachment.get(key) for key in ("name", "type", "size")}, "tokenDigest": f"sha256:{hashlib.sha256(('ack-feishu-attachment-v1\x1f' + token).encode('utf-8')).hexdigest()}", } for attachment, token in zip(record["attachments"], tokens) ] encoded = json.dumps( stable, ensure_ascii=False, sort_keys=True, separators=(",", ":"), ).encode("utf-8") return f"sha256:{hashlib.sha256(encoded).hexdigest()}" def fetch(config: dict[str, Any], output_dir: Path | None) -> dict[str, Any]: profile_check(config) prepared: list[tuple[dict[str, Any], list[tuple[dict[str, Any], str]]]] = [] batch_warnings: list[dict[str, str]] = [] total_attachments = 0 total_attachment_bytes = 0 for record_id, row in fetch_pages(config): workflow = config.get("workflow", "read-only-v1") field_order = CLARIFIED_FIELD_ORDER if workflow == "clarified-writeback-v1" else LEGACY_FIELD_ORDER logical_fields = [logical for logical in field_order if logical in config["fields"]] cells = dict(zip(logical_fields, row)) attachment_data = attachment_items(cells["attachments"]) total_attachments += len(attachment_data) total_attachment_bytes += sum(metadata["size"] for metadata, _ in attachment_data) if total_attachments > MAX_TOTAL_ATTACHMENTS: raise IntakeError("batch exceeded the attachment count limit") if total_attachment_bytes > MAX_TOTAL_ATTACHMENT_BYTES: raise IntakeError("batch exceeded the attachment byte limit") if workflow == "clarified-writeback-v1": record = { "sourceRef": source_ref(config, record_id), "recordId": record_id, "updatedAt": text(cells["updatedAt"]), "title": text(cells["title"]), "details": text(cells["details"]), "problemStatement": text(cells["problemStatement"]), "expectedOutcome": text(cells["expectedOutcome"]), "acceptance": text(cells["acceptance"]), "intakeStatus": text(cells["intakeStatus"]), "ackTaskId": text(cells["ackTaskId"]), "attachments": [metadata for metadata, _ in attachment_data], "warnings": [], } else: record = {"sourceRef": source_ref(config, record_id), "recordId": record_id, "updatedAt": text(cells["updatedAt"]), "title": text(cells["title"]), "actual": text(cells["actual"]), "expected": text(cells["expected"]), "steps": text(cells["stepsToReproduce"]), "fixLogic": text(cells.get("fixLogic")), "acceptance": text(cells["acceptance"]), "priority": text(cells.get("priority")), "attachments": [metadata for metadata, _ in attachment_data], "warnings": []} content_fields = ( ("title", "details", "problemStatement", "expectedOutcome", "acceptance") if workflow == "clarified-writeback-v1" else BUG_CONTENT_FIELDS ) if not attachment_data and not any(record.get(field) for field in content_fields): batch_warnings.append({"recordId": record_id, "code": "blank_record_skipped"}) continue source_fields = ("title", "updatedAt") if workflow == "clarified-writeback-v1" else SOURCE_FACT_FIELDS for field in source_fields: if not record[field]: raise IntakeError(f"record {field} must not be empty") enrichment_fields = (["problemStatement", "expectedOutcome", "acceptance"] if workflow == "clarified-writeback-v1" else list(COORDINATOR_FIELDS)) if workflow != "clarified-writeback-v1" and "fixLogic" in config["fields"]: enrichment_fields.append("fixLogic") record["enrichmentRequired"] = [ field for field in enrichment_fields if not record[field] ] record["draftRevision"] = draft_revision( record, [token for _, token in attachment_data], ) prepared.append((record, attachment_data)) download_root: Path | None = None if output_dir is not None: try: output_dir.mkdir(parents=True, exist_ok=True) if not output_dir.is_dir() or output_dir.is_symlink(): raise OSError("unsafe output directory") download_root = output_dir.resolve(strict=True) except OSError as exc: raise IntakeError("attachment output directory is unsafe") from exc records = [] attachment_deadline = time.monotonic() + MAX_ATTACHMENT_BATCH_SECONDS for record, attachment_data in prepared: if download_root is not None: for index, (attachment, file_token) in enumerate(attachment_data, start=1): remaining = attachment_deadline - time.monotonic() if remaining <= 0: raise IntakeError("attachment batch exceeded the time limit") attachment_dir = download_root / record["recordId"] / f"attachment-{index:02d}" attachment["localPath"] = download( config, record["recordId"], file_token, attachment_dir, attachment["size"], min(60, remaining), ) records.append(record) return {"provider": "feishu-base", "workflow": config.get("workflow", "read-only-v1"), "profile": config["profile"], "tableId": config["tableId"], "viewId": config["viewId"], "records": records, "warnings": batch_warnings} def review_record( config: dict[str, Any], record_id: str, expected_source_ref: str, expected_revision: str, ) -> dict[str, Any]: """Resolve one record inside the configured view and bind its reviewed version.""" if RECORD_ID.fullmatch(record_id) is None: raise IntakeError("review record id is invalid") if SOURCE_REF.fullmatch(expected_source_ref) is None: raise IntakeError("expected source reference is invalid") if DRAFT_REVISION.fullmatch(expected_revision) is None: raise IntakeError("expected draft revision is invalid") matching = [ record for record in fetch(config, None)["records"] if record["recordId"] == record_id ] if len(matching) != 1: raise IntakeError("configured review view did not contain exactly one record") if ( matching[0]["sourceRef"] != expected_source_ref or matching[0]["draftRevision"] != expected_revision ): raise IntakeError("review record changed before the requested operation") return matching[0] def write_draft( config: dict[str, Any], record_id: str, expected_source_ref: str, expected_revision: str, draft_path: Path, ) -> dict[str, Any]: """Overwrite only the review-owned cells for the configured workflow.""" workflow = config.get("workflow", "read-only-v1") if workflow not in {"reviewed-writeback-v1", "clarified-writeback-v1"}: raise IntakeError("draft writeback requires a writeback workflow") if workflow == "reviewed-writeback-v1" and "fixLogic" not in config["fields"]: raise IntakeError("bugIntake.fields.fixLogic is required for draft writeback") review_record(config, record_id, expected_source_ref, expected_revision) draft = load_draft(draft_path, workflow) if workflow == "clarified-writeback-v1": patch = { config["fields"]["problemStatement"]: draft["problemStatement"], config["fields"]["expectedOutcome"]: draft["expectedOutcome"], config["fields"]["acceptance"]: "\n".join( f"{index}. {item}" for index, item in enumerate(draft["acceptance"], start=1) ), config["fields"]["intakeStatus"]: "待审核", } written = ["problemStatement", "expectedOutcome", "acceptance", "intakeStatus"] else: patch = { config["fields"]["fixLogic"]: draft["fixLogic"], config["fields"]["acceptance"]: "\n".join( f"{index}. {item}" for index, item in enumerate(draft["acceptance"], start=1) ), } written = ["fixLogic", "acceptance"] profile_check(config) run_cli([ "base", "+record-upsert", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--record-id", record_id, "--json", json.dumps(patch, ensure_ascii=False, separators=(",", ":")), "--format", "json", ]) matching = [record for record in fetch(config, None)["records"] if record["recordId"] == record_id] if len(matching) != 1 or matching[0]["sourceRef"] != expected_source_ref: raise IntakeError("draft writeback readback did not find exactly one record") expected_acceptance = text(patch[config["fields"]["acceptance"]]) if workflow == "clarified-writeback-v1": matched = ( matching[0]["problemStatement"] == text(draft["problemStatement"]) and matching[0]["expectedOutcome"] == text(draft["expectedOutcome"]) and matching[0]["acceptance"] == expected_acceptance and matching[0]["intakeStatus"] == "待审核" ) else: matched = matching[0]["fixLogic"] == text(draft["fixLogic"]) and matching[0]["acceptance"] == expected_acceptance if not matched: raise IntakeError("draft writeback readback did not match the submitted draft") return { "provider": "feishu-base", "recordId": record_id, "written": written, "draftRevision": matching[0]["draftRevision"], "ok": True, } def import_approved( config: dict[str, Any], record_id: str, expected_source_ref: str, expected_revision: str, ) -> dict[str, Any]: """Emit the canonical task payload for one explicitly approved draft revision.""" workflow = config.get("workflow", "read-only-v1") if workflow not in {"reviewed-writeback-v1", "clarified-writeback-v1"}: raise IntakeError("approved import requires a writeback workflow") record = review_record( config, record_id, expected_source_ref, expected_revision, ) acceptance = review_items(record["acceptance"]) if workflow == "clarified-writeback-v1": if record["intakeStatus"] != "已确认": raise IntakeError("approved record must have intakeStatus 已确认") if not record["problemStatement"] or not record["expectedOutcome"] or not acceptance: raise IntakeError("approved record is missing prepared clarification fields") task_draft = clarified_task_draft(record) return {"provider": "feishu-base", "recordId": record_id, "draftRevision": record["draftRevision"], "taskDraft": task_draft, "ok": True} steps = review_items(record["steps"]) if not record["priority"] or not steps or not record["fixLogic"] or not acceptance: raise IntakeError("approved record is missing prepared review fields") task_draft: dict[str, Any] = { "title": record["title"], "priority": record["priority"], "description": record["title"], "actual": record["actual"], "expected": record["expected"], "stepsToReproduce": steps, "fixLogic": record["fixLogic"], "acceptanceCriteria": acceptance, "source": { "kind": "feishu-base", "workflow": "reviewed-writeback-v1", "ref": record["sourceRef"], "recordId": record["recordId"], "updatedAt": record["updatedAt"], "approvedRevision": record["draftRevision"], }, } task_draft["source"]["approvedPayloadHash"] = approval_payload_hash(task_draft) return { "provider": "feishu-base", "recordId": record_id, "draftRevision": record["draftRevision"], "taskDraft": task_draft, "ok": True, } def clarified_task_draft(record: dict[str, Any]) -> dict[str, Any]: """Map one normalized clarified record to its immutable reviewed task fields.""" task_draft: dict[str, Any] = { "title": record["title"], "description": record["problemStatement"], "actual": record["details"] or record["title"], "expected": record["expectedOutcome"], "acceptanceCriteria": review_items(record["acceptance"]), "source": { "kind": "feishu-base", "workflow": "clarified-writeback-v1", "ref": record["sourceRef"], "recordId": record["recordId"], "updatedAt": record["updatedAt"], "approvedRevision": record["draftRevision"], }, } task_draft["source"]["approvedPayloadHash"] = approval_payload_hash(task_draft) return task_draft def mark_imported( board: dict[str, Any], config: dict[str, Any], record_id: str, expected_source_ref: str, expected_revision: str, task_id: str, ) -> dict[str, Any]: """Bind a confirmed Base record to the validated ACK task created from it.""" if config.get("workflow") != "clarified-writeback-v1": raise IntakeError("mark-imported requires clarified-writeback-v1 workflow") if not isinstance(task_id, str) or not task_id: raise IntakeError("ACK task id is invalid") if validate_task_board(board): raise IntakeError("task board is invalid for import finalization") record = review_record(config, record_id, expected_source_ref, expected_revision) tasks = board.get("tasks") if not isinstance(tasks, list): raise IntakeError("task board tasks must be a list") matches = [ task for task in tasks if isinstance(task, dict) and task.get("id") == task_id ] if len(matches) != 1: raise IntakeError("task board did not contain exactly one imported ACK task") task = matches[0] source = task.get("source") expected_task = clarified_task_draft(record) reviewed_fields = ( "title", "description", "actual", "expected", "acceptanceCriteria", ) if ( not isinstance(source, dict) or source != expected_task["source"] or any(task.get(field) != expected_task[field] for field in reviewed_fields) ): raise IntakeError("ACK task does not match the approved Base record") if record["intakeStatus"] == "已导入" and record["ackTaskId"] == task_id: return { "provider": "feishu-base", "recordId": record_id, "intakeStatus": "已导入", "ackTaskId": task_id, "draftRevision": expected_revision, "ok": True, } if record["intakeStatus"] != "已确认" or record["ackTaskId"]: raise IntakeError("Base record is not ready to mark as imported") profile_check(config) run_cli([ "base", "+record-upsert", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--record-id", record_id, "--json", json.dumps({ config["fields"]["intakeStatus"]: "已导入", config["fields"]["ackTaskId"]: task_id, }, ensure_ascii=False, separators=(",", ":")), "--format", "json", ]) matching = [ item for item in fetch(config, None)["records"] if item["recordId"] == record_id ] if ( len(matching) != 1 or matching[0]["sourceRef"] != expected_source_ref or matching[0]["draftRevision"] != expected_revision or matching[0]["intakeStatus"] != "已导入" or matching[0]["ackTaskId"] != task_id ): raise IntakeError("import marker readback did not match the ACK task") return { "provider": "feishu-base", "recordId": record_id, "intakeStatus": "已导入", "ackTaskId": task_id, "draftRevision": expected_revision, "ok": True, } def plan_actions(board: dict[str, Any], records: list[dict[str, Any]]) -> list[dict[str, Any]]: """Plan idempotent Coordinator actions without mutating the task board.""" tasks = board.get("tasks") if not isinstance(tasks, list): raise IntakeError("task board tasks must be a list") existing: dict[str, dict[str, Any]] = {} for task in tasks: if not isinstance(task, dict): continue source = task.get("source") if not isinstance(source, dict) or source.get("kind") != "feishu-base": continue ref = source.get("ref") task_id = task.get("id") status = task.get("status") updated_at = source.get("updatedAt") if ( not isinstance(ref, str) or SOURCE_REF.fullmatch(ref) is None or not isinstance(task_id, str) or not task_id or not isinstance(status, str) or not status or not isinstance(updated_at, str) or not updated_at ): raise IntakeError("existing Feishu task source is invalid") if ref in existing: raise IntakeError("task board contains duplicate Feishu source references") existing[ref] = task actions: list[dict[str, Any]] = [] seen_records: set[str] = set() for record in records: ref = record.get("sourceRef") record_id = record.get("recordId") updated_at = record.get("updatedAt") revision = record.get("draftRevision") if ( not isinstance(ref, str) or SOURCE_REF.fullmatch(ref) is None or not isinstance(record_id, str) or RECORD_ID.fullmatch(record_id) is None or not isinstance(updated_at, str) or not updated_at or not isinstance(revision, str) or DRAFT_REVISION.fullmatch(revision) is None ): raise IntakeError("normalized Feishu record identity is invalid") if ref in seen_records: raise IntakeError("fetched records contain a duplicate source reference") seen_records.add(ref) task = existing.get(ref) if task is None: planned_action: dict[str, Any] = {"sourceRef": ref, "recordId": record_id, "draftRevision": record["draftRevision"], "action": "create"} enrichment_required = record.get("enrichmentRequired") if enrichment_required: planned_action["enrichmentRequired"] = enrichment_required actions.append(planned_action) continue source = task["source"] approved_revision = source.get("approvedRevision") if isinstance(approved_revision, str): if approved_revision == revision: action_name = "unchanged" elif task["status"] == "open": action_name = "refresh" else: action_name = "drift" elif source["updatedAt"] == updated_at: action_name = "unchanged" elif task["status"] == "open": action_name = "refresh" else: action_name = "drift" planned_action = { "sourceRef": ref, "recordId": record_id, "draftRevision": record["draftRevision"], "taskId": task["id"], "status": task["status"], "action": action_name, } enrichment_required = record.get("enrichmentRequired") if enrichment_required and action_name == "refresh": planned_action["enrichmentRequired"] = enrichment_required actions.append(planned_action) return actions def field_list(config: dict[str, Any]) -> list[dict[str, str]]: """Return the bounded Base field inventory used by schema migration.""" profile_check(config) response = run_cli([ "base", "+field-list", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--format", "json", ]) data = response.get("data", response) items = data.get("fields") if isinstance(data, dict) else None if not isinstance(items, list) or len(items) > 256: raise IntakeError("field list returned an invalid response") result: list[dict[str, str]] = [] for item in items: if not isinstance(item, dict) or not all( isinstance(item.get(key), str) and item[key] for key in ("id", "name", "type") ): raise IntakeError("field list returned an invalid field") result.append({key: item[key] for key in ("id", "name", "type")}) if len({item["name"] for item in result}) != len(result): raise IntakeError("field list contains duplicate names") return result def schema_target(config: dict[str, Any]) -> dict[str, str]: token_digest = hashlib.sha256( ("ack-feishu-schema-target-v1\x1f" + config["baseToken"]).encode("utf-8") ).hexdigest() return { "profile": config["profile"], "baseTokenDigest": f"sha256:{token_digest}", "tableId": config["tableId"], "viewId": config["viewId"], } def schema_fingerprint(config: dict[str, Any], fields: list[dict[str, str]]) -> str: encoded = json.dumps( { "contract": "clarified-writeback-v1", "target": schema_target(config), "fields": sorted(fields, key=lambda item: item["id"]), }, ensure_ascii=False, sort_keys=True, separators=(",", ":"), ).encode("utf-8") return f"sha256:{hashlib.sha256(encoded).hexdigest()}" def schema_plan(config: dict[str, Any]) -> dict[str, Any]: if config.get("workflow") != "clarified-writeback-v1": raise IntakeError("schema migration requires clarified-writeback-v1 workflow") fields = field_list(config) by_name = {item["name"]: item for item in fields} missing = [name for name, _ in TARGET_BASE_FIELDS if name not in by_name] wrong_type = [ {"name": name, "expected": field_type, "actual": by_name[name]["type"]} for name, field_type in TARGET_BASE_FIELDS if name in by_name and by_name[name]["type"] != field_type ] return { "provider": "feishu-base", "target": schema_target(config), "schemaFingerprint": schema_fingerprint(config, fields), "missingFields": missing, "typeConflicts": wrong_type, "visibleFields": [name for name, _ in TARGET_BASE_FIELDS], "legacyFieldsPreserved": [ name for name in ("期望结果", "问题澄清", "复现步骤", "ACK Ready") if name in by_name ], "ok": not wrong_type, } def create_target_field(config: dict[str, Any], name: str) -> None: field_type = dict(TARGET_BASE_FIELDS)[name] if field_type in {"attachment", "updated_at"}: raise IntakeError("schema migration cannot create a missing system/source field") payload: dict[str, Any] = {"name": name, "type": field_type} if name == "处理状态": payload.update({"multiple": False, "options": [{"name": value} for value in INTAKE_STATUSES]}) run_cli([ "base", "+field-create", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--json", json.dumps(payload, ensure_ascii=False, separators=(",", ":")), "--format", "json", ]) def migration_rows(config: dict[str, Any], fields: list[dict[str, str]]) -> list[tuple[str, dict[str, Any]]]: by_name = {item["name"]: item for item in fields} names = ["标题", "详细描述", "期望结果", "处理状态"] present = [name for name in names if name in by_name] # record-list projects cells by configured field name and returns those # names in its matrix, even when the REST field inventory exposes IDs. ids = present args = [ "base", "+record-list", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--view-id", config["viewId"], "--format", "json", "--offset", "0", "--limit", str(PAGE_SIZE), ] for field_id in ids: args.extend(["--field-id", field_id]) response = run_cli(args) record_ids, rows = matrix_from_response(response, ids) data = response.get("data", response) if data.get("has_more", data.get("hasMore", False)): raise IntakeError("schema migration view exceeded one bounded page") return [(record_id, dict(zip(present, row))) for record_id, row in zip(record_ids, rows)] def schema_apply(config: dict[str, Any], expected_fingerprint: str) -> dict[str, Any]: if DRAFT_REVISION.fullmatch(expected_fingerprint) is None: raise IntakeError("expected schema fingerprint is invalid") before = schema_plan(config) if before["schemaFingerprint"] != expected_fingerprint: raise IntakeError("Base schema changed after planning") if before["typeConflicts"]: raise IntakeError("Base schema has incompatible target field types") for name in before["missingFields"]: create_target_field(config, name) fields = field_list(config) for _ in range(4): if all(name in {item["name"] for item in fields} for name, _ in TARGET_BASE_FIELDS): break time.sleep(0.5) fields = field_list(config) by_name = {item["name"]: item for item in fields} if any(name not in by_name for name, _ in TARGET_BASE_FIELDS): raise IntakeError("schema migration did not create all target fields") migrated: list[str] = [] for record_id, cells in migration_rows(config, fields): details = text(cells.get("详细描述")) legacy_expected = text(cells.get("期望结果")) patch: dict[str, Any] = {} marker = f"用户原始期望:{legacy_expected}" if legacy_expected else "" if marker and marker not in details: patch["详细描述"] = f"{details}\n\n{marker}".strip() if not text(cells.get("处理状态")): patch["处理状态"] = "待整理" if patch: run_cli([ "base", "+record-upsert", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--record-id", record_id, "--json", json.dumps(patch, ensure_ascii=False, separators=(",", ":")), "--format", "json", ]) migrated.append(record_id) visible_ids = [by_name[name]["id"] for name, _ in TARGET_BASE_FIELDS] run_cli([ "base", "+view-set-visible-fields", "--profile", config["profile"], "--base-token", config["baseToken"], "--table-id", config["tableId"], "--view-id", config["viewId"], "--json", json.dumps({"visible_fields": visible_ids}, separators=(",", ":")), "--format", "json", ]) return { "provider": "feishu-base", "createdFields": before["missingFields"], "migratedRecordIds": migrated, "visibleFields": [name for name, _ in TARGET_BASE_FIELDS], "schemaFingerprint": schema_fingerprint(config, fields), "ok": True, } def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser(description="Read and review a configured Feishu Base bug intake") sub = parser.add_subparsers(dest="command", required=True) for name in ("check", "fetch", "plan", "write-draft", "import-approved", "mark-imported", "schema-plan", "schema-apply"): command = sub.add_parser(name) command.add_argument("tasks", type=Path, help="ACK tasks.yaml or JSON board") if name in {"fetch", "plan"}: command.add_argument("--output-dir", type=Path, help="explicit directory for downloaded attachments") if name in {"write-draft", "import-approved", "mark-imported"}: command.add_argument("--record-id", required=True, help="existing Feishu Base record id") command.add_argument("--expected-source-ref", required=True, help="sourceRef returned by fetch") command.add_argument("--expected-draft-revision", required=True, help="draftRevision returned by fetch") if name == "write-draft": command.add_argument("--input", type=Path, required=True, help="bounded JSON draft file") if name == "mark-imported": command.add_argument("--task-id", required=True, help="validated ACK task id") if name == "schema-apply": command.add_argument("--expected-schema-fingerprint", required=True) args = parser.parse_args(argv) try: board = load_board(args.tasks) config = config_from_board(board) if args.command == "check": profile_check(config) output = {"provider": "feishu-base", "workflow": config.get("workflow", "read-only-v1"), "profile": config["profile"], "ok": True} elif args.command == "fetch": output = fetch(config, args.output_dir) elif args.command == "plan": output = fetch(config, args.output_dir) output["actions"] = plan_actions(board, output["records"]) elif args.command == "schema-plan": output = schema_plan(config) elif args.command == "schema-apply": output = schema_apply(config, args.expected_schema_fingerprint) elif args.command == "write-draft": output = write_draft( config, args.record_id, args.expected_source_ref, args.expected_draft_revision, args.input, ) elif args.command == "mark-imported": output = mark_imported( board, config, args.record_id, args.expected_source_ref, args.expected_draft_revision, args.task_id, ) else: output = import_approved( config, args.record_id, args.expected_source_ref, args.expected_draft_revision, ) except IntakeError as exc: sys.stderr.write(f"Feishu bug intake failed: {exc}\n") return 1 except (OSError, subprocess.SubprocessError): sys.stderr.write("Feishu bug intake failed: local I/O failed\n") return 1 print(json.dumps(output, ensure_ascii=False, separators=(",", ":"))) return 0 if __name__ == "__main__": raise SystemExit(main())