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