feat(ack): add Feishu bug intake
This commit is contained in:
@@ -0,0 +1,583 @@
|
||||
#!/usr/bin/env python3
|
||||
"""Read an ACK-ready Feishu Base view through the official lark-cli.
|
||||
|
||||
This is deliberately a small, non-mutating adapter. It 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 yaml_subset import DuplicateKeyError, YamlSubsetError, load_json_unique, load_yaml_subset, make_unique_pyyaml_loader
|
||||
|
||||
REQUIRED_FIELDS = ("title", "actual", "expected", "stepsToReproduce", "acceptance", "priority", "attachments", "updatedAt")
|
||||
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
|
||||
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}$")
|
||||
|
||||
|
||||
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 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", "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")
|
||||
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 not isinstance(fields, dict) or set(fields) != set(REQUIRED_FIELDS):
|
||||
raise IntakeError("bugIntake.fields must map exactly the required 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")
|
||||
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 fields != field_ids:
|
||||
raise IntakeError("record list fields do not match configured projection")
|
||||
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")
|
||||
return ids, rows
|
||||
|
||||
|
||||
def fetch_pages(config: dict[str, Any]) -> list[tuple[str, list[Any]]]:
|
||||
field_ids = [config["fields"][logical] for logical in REQUIRED_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 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]]]] = []
|
||||
total_attachments = 0
|
||||
total_attachment_bytes = 0
|
||||
for record_id, row in fetch_pages(config):
|
||||
cells = dict(zip(REQUIRED_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")
|
||||
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"]), "acceptance": text(cells["acceptance"]), "priority": text(cells["priority"]), "attachments": [metadata for metadata, _ in attachment_data], "warnings": []}
|
||||
for field in ("title", "actual", "expected", "steps", "acceptance", "priority", "updatedAt"):
|
||||
if not record[field]:
|
||||
raise IntakeError(f"record {field} must not be empty")
|
||||
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", "profile": config["profile"], "tableId": config["tableId"], "viewId": config["viewId"], "records": records}
|
||||
|
||||
|
||||
def plan_actions(board: dict[str, Any], records: list[dict[str, Any]]) -> list[dict[str, str]]:
|
||||
"""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, str]] = []
|
||||
seen_records: set[str] = set()
|
||||
for record in records:
|
||||
ref = record.get("sourceRef")
|
||||
record_id = record.get("recordId")
|
||||
updated_at = record.get("updatedAt")
|
||||
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
|
||||
):
|
||||
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:
|
||||
actions.append({"sourceRef": ref, "recordId": record_id, "action": "create"})
|
||||
continue
|
||||
source = task["source"]
|
||||
if source["updatedAt"] == updated_at:
|
||||
action = "unchanged"
|
||||
elif task["status"] == "open":
|
||||
action = "refresh"
|
||||
else:
|
||||
action = "drift"
|
||||
actions.append({
|
||||
"sourceRef": ref,
|
||||
"recordId": record_id,
|
||||
"taskId": task["id"],
|
||||
"status": task["status"],
|
||||
"action": action,
|
||||
})
|
||||
return actions
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
parser = argparse.ArgumentParser(description="Read a configured Feishu Base bug intake")
|
||||
sub = parser.add_subparsers(dest="command", required=True)
|
||||
for name in ("check", "fetch", "plan"):
|
||||
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")
|
||||
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", "profile": config["profile"], "ok": True}
|
||||
elif args.command == "fetch":
|
||||
output = fetch(config, args.output_dir)
|
||||
else:
|
||||
output = fetch(config, args.output_dir)
|
||||
output["actions"] = plan_actions(board, output["records"])
|
||||
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())
|
||||
@@ -113,6 +113,15 @@ DELIVERY_STATUSES = {
|
||||
}
|
||||
DELIVERY_ARTIFACT_FIELDS = {"id", "type", "reference", "digest"}
|
||||
DELIVERY_DEPLOYMENT_FIELDS = {"environment", "result", "evidence"}
|
||||
FEISHU_REQUIRED_FIELDS = {
|
||||
"title", "actual", "expected", "stepsToReproduce", "acceptance", "priority",
|
||||
"attachments", "updatedAt",
|
||||
}
|
||||
FEISHU_CONFIG_FIELDS = {"provider", "profile", "baseToken", "tableId", "viewId", "fields"}
|
||||
FEISHU_SOURCE_FIELDS = {"kind", "ref", "recordId", "updatedAt"}
|
||||
FEISHU_PROFILE_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$")
|
||||
FEISHU_SOURCE_REF_RE = re.compile(r"^feishu-base:sha256:[0-9a-f]{64}$")
|
||||
FEISHU_RECORD_ID_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,255}$")
|
||||
DISPATCH_FIELDS = {
|
||||
"taskId",
|
||||
"dispatchId",
|
||||
@@ -654,6 +663,28 @@ def validate_builtin(data: dict) -> list[str]:
|
||||
{"repoPath", "baseUrl", "devWorktree", "overlayFile", "deliveryFile"},
|
||||
"project",
|
||||
)
|
||||
if "bugIntake" in project:
|
||||
intake = project["bugIntake"]
|
||||
if not isinstance(intake, dict):
|
||||
errors.append("project.bugIntake 必须是对象")
|
||||
else:
|
||||
reject_unknown_fields(intake, FEISHU_CONFIG_FIELDS, "project.bugIntake", errors)
|
||||
if intake.get("provider") != "feishu-base":
|
||||
errors.append("project.bugIntake.provider 必须是 feishu-base")
|
||||
profile = intake.get("profile")
|
||||
if not isinstance(profile, str) or FEISHU_PROFILE_RE.fullmatch(profile) is None:
|
||||
errors.append("project.bugIntake.profile 非法")
|
||||
for key in ("baseToken", "tableId", "viewId"):
|
||||
value = intake.get(key)
|
||||
if not isinstance(value, str) or not value.strip() or any(char.isspace() for char in value):
|
||||
errors.append(f"project.bugIntake.{key} 必须是无空白非空字符串")
|
||||
fields = intake.get("fields")
|
||||
if not isinstance(fields, dict) or set(fields) != FEISHU_REQUIRED_FIELDS:
|
||||
errors.append("project.bugIntake.fields 必须且只能映射所需逻辑字段")
|
||||
elif any(not isinstance(v, str) or not v.strip() or any(c.isspace() for c in v) for v in fields.values()):
|
||||
errors.append("project.bugIntake.fields 字段值必须是无空白非空字符串")
|
||||
elif len(set(fields.values())) != len(fields):
|
||||
errors.append("project.bugIntake.fields 字段值不能重复")
|
||||
if (
|
||||
"knowledgeFile" in project
|
||||
and project.get("knowledgeFile") != "docs/ack/knowledge.yaml"
|
||||
@@ -710,6 +741,7 @@ def validate_builtin(data: dict) -> list[str]:
|
||||
return errors
|
||||
|
||||
seen_ids: set[str] = set()
|
||||
seen_source_refs: set[str] = set()
|
||||
for i, task in enumerate(tasks):
|
||||
where = f"tasks[{i}]"
|
||||
if not isinstance(task, dict):
|
||||
@@ -756,6 +788,25 @@ def validate_builtin(data: dict) -> list[str]:
|
||||
)
|
||||
validate_object_fields(task, {"evidence", "verification"}, where)
|
||||
|
||||
if "source" in task:
|
||||
source = task["source"]
|
||||
# `source` was historically an open extension point. Preserve
|
||||
# non-Feishu strings/objects and tighten only the namespaced shape.
|
||||
if isinstance(source, dict) and source.get("kind") == "feishu-base":
|
||||
reject_unknown_fields(source, FEISHU_SOURCE_FIELDS, f"{where}.source", errors)
|
||||
ref = source.get("ref")
|
||||
if not isinstance(ref, str) or FEISHU_SOURCE_REF_RE.fullmatch(ref) is None:
|
||||
errors.append(f"{where}.source.ref: 必须是不透明 feishu-base SHA-256 引用")
|
||||
else:
|
||||
if ref in seen_source_refs:
|
||||
errors.append(f"{where}.source.ref: 来源引用重复")
|
||||
seen_source_refs.add(ref)
|
||||
record_id = source.get("recordId")
|
||||
if not isinstance(record_id, str) or FEISHU_RECORD_ID_RE.fullmatch(record_id) is None:
|
||||
errors.append(f"{where}.source.recordId: 必须是合法飞书记录 ID")
|
||||
if not _nonempty_string(source.get("updatedAt")):
|
||||
errors.append(f"{where}.source.updatedAt: 必须是非空字符串")
|
||||
|
||||
validate_knowledge_fields(task, where, status, errors)
|
||||
|
||||
if "dispatch" not in task:
|
||||
|
||||
Reference in New Issue
Block a user