1222 lines
55 KiB
Python
1222 lines
55 KiB
Python
#!/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())
|