Files
.pouch/skills/ack/scripts/feishu_bug_intake.py
T

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())