Compare commits

...
Author SHA1 Message Date
Dennis a82d945f54 fix(doc): reject explicit revert failure states 2026-08-17 14:48:13 +08:00
Dennis 6846326445 fix(doc): reject revert request echo evidence 2026-08-17 14:48:11 +08:00
Dennis 54c2054a5c fix(doc): ignore generated JSONML defaults 2026-08-17 14:48:09 +08:00
Dennis e5bf332b05 fix(doc): address readback review findings 2026-08-17 14:48:06 +08:00
Dennis 9ed55978d9 fix(doc): cancel readback retry waits 2026-08-17 14:48:04 +08:00
Dennis 7ffbbc4a51 test(doc): cover stable pagination identities 2026-08-17 14:48:02 +08:00
Dennis 2db73a8185 fix(doc): distinguish identical pagination pages 2026-08-17 14:48:00 +08:00
Dennis a62332be93 fix(doc): trust only explicit inserted block IDs 2026-08-17 14:47:58 +08:00
Dennis 1083093cbc test(doc): complete readback coverage evidence 2026-08-17 14:47:56 +08:00
Dennis c6ebe307cd fix(doc): verify inline media from jsonml readback 2026-08-17 14:47:54 +08:00
Dennis 3ee66d4373 fix(doc): harden mutation readback verification 2026-08-17 14:47:51 +08:00
github-actions[bot] a0be395ccc Merge pull request #1006 from DingTalk-Real-AI/codex/fix-aitable-pagination-minutes-unshare
fix(shortcut): harden Aitable pagination and Minutes unshare
2026-08-17 06:37:16 +00:00
Dennis c1a549cd64 fix: close delete readback continuations 2026-08-17 14:13:51 +08:00
Dennis 5a414999ef fix: validate record query previews 2026-08-17 14:13:49 +08:00
Dennis 7aa8240629 fix: preserve record query preview contract 2026-08-17 14:13:47 +08:00
Dennis 37b9a1dc31 fix: bound exact aitable record queries 2026-08-17 14:13:45 +08:00
Dennis 2ab8748c4d test: use native minutes path separators 2026-08-17 14:13:43 +08:00
Dennis f041275811 fix: make minutes polling portable 2026-08-17 14:13:41 +08:00
Dennis f486105836 fix: bound empty aitable pagination 2026-08-17 14:13:38 +08:00
Dennis e14de2b4c2 test: close shortcut fix review gates 2026-08-17 14:13:36 +08:00
Dennis fe2f3ca92f fix: harden aitable pagination and minutes unshare 2026-08-17 14:13:33 +08:00
github-actions[bot] 8e4519cacd Merge pull request #1014 from FloralTide/codex/fix-windows-event-bus
fix(event): support Windows bus lifecycle
2026-08-17 14:12:46 +08:00
炳昱 16abb481e8 Merge remote-tracking branch 'upstream/main' into codex/fix-windows-event-bus 2026-08-17 13:00:07 +08:00
炳昱 7ad82bbf0a fix(event): accept bus exit at stop timeout boundary 2026-08-17 13:00:07 +08:00
chichuan 104eb715c4 Merge pull request #989 from maoqxxmm/codex/sheet-dropdown-source-range
feat(sheet): support SourceRange dropdowns and read completion
2026-08-17 12:19:00 +08:00
chichuan 97ea887ea5 Merge branch 'main' into codex/sheet-dropdown-source-range 2026-08-17 11:49:42 +08:00
github-actions[bot] bfeb9f6af0 chore: update beta formula for v1.0.59-beta.2 [skip ci] 2026-08-17 03:35:41 +00:00
毛球 e26f278112 Merge branch 'main' into codex/sheet-dropdown-source-range 2026-08-17 11:17:45 +08:00
chichuan e6b5938bd8 Merge pull request #1025 from DingTalk-Real-AI/codex/changelog-v1.0.59-beta.2
docs: seal changelog for v1.0.59-beta.2
2026-08-17 11:03:27 +08:00
chichuan 4f95373420 docs: seal changelog for v1.0.59-beta.2 2026-08-17 10:59:11 +08:00
炳昱 afb90009f6 Merge remote-tracking branch 'upstream/main' into codex/fix-windows-event-bus 2026-08-17 10:49:25 +08:00
github-actions[bot] 6411d26a95 Merge pull request #1023 from DingTalk-Real-AI/fix/app-partition-parallel-jobs
fix(ci): parallelize app test partitions and drop race from the schema partition
2026-08-17 02:45:57 +00:00
毛球 31117d1b89 Merge branch 'main' into codex/sheet-dropdown-source-range 2026-08-17 10:43:01 +08:00
炳昱 e3553fe7a5 test(event): cover bus ownership validation failures 2026-08-17 10:40:21 +08:00
炳昱 067aff179f Merge remote-tracking branch 'upstream/main' into codex/fix-windows-event-bus 2026-08-17 10:32:08 +08:00
炳昱 353454abb2 fix(event): verify bus owner before fallback stop 2026-08-17 10:32:04 +08:00
chichuan 96b9cbce02 Merge branch 'main' into fix/app-partition-parallel-jobs 2026-08-17 10:22:20 +08:00
chichuan 36877d00dc Merge pull request #1024 from DingTalk-Real-AI/perf/schema-json-projection
perf: skip redundant JSON validation when projecting typed Schema values
2026-08-17 10:21:48 +08:00
xiatian a9a97c2746 Merge remote-tracking branch 'upstream/main' into codex/sheet-dropdown-source-range 2026-08-17 09:43:54 +08:00
chichuan 55d94d3b58 perf: skip redundant JSON validation when projecting typed Schema values
typedJSONValue marshaled a typed value and then routed the result through
rawJSONValue, which runs json.Valid before decoding. On that path the input is
whatever json.Marshal has just produced, so the validation scan can only ever
succeed: it re-read every marshaled document for nothing.

The decode step is now shared by both entry points. rawJSONValue keeps its
json.Valid check, because it still accepts untrusted input, while typedJSONValue
decodes what it marshaled directly. Across the 1121-tool set this removes about a
third of the Schema Catalog projection work: the internal/app schema suite goes
from 26.0s to 17.2s uninstrumented, and from 291.1s to 241.0s under -race.

The delivered Catalog is byte-for-byte unchanged. check-generated-drift,
check-schema-catalog and check-schema-binary each regenerate the same
source_hash sha256:93b8d44eb163bd2898c78397d22af92d378e3dc4e20f56b33277b51e4342e2e6,
and the two error contracts are preserved: typedJSONValue still rejects a value
json.Marshal cannot encode, and rawJSONValue still rejects invalid JSON.
2026-08-16 22:21:16 +08:00
chichuan bfd0976b31 fix(ci): run the app test partitions as parallel shards
The five internal/app partitions ran end to end inside one job, so the app
shard's wall clock was the sum of all five: 780s in CI, of which the schema
partition owned 357s. Each partition is now its own matrix shard, so they run
concurrently and the shard's wall clock is set by its slowest partition rather
than by their total. Every partition shard still selects the same single
internal/app package, so the impacted-package query maps the shard name back to
app and the partition only chooses which tests run.

The helper gains a partition argument and a list-partitions mode. APP_PARTITIONS
is the single source of truth for the set, and the discovery pass still runs in
every job, so each one independently verifies that the partition patterns cover
every top-level test exactly once before running the one it was asked for.

Two fail-closed checks guard the split, because the helper's own coverage check
can no longer prove the whole package ran once the partitions are separate jobs:

- The helper cross-checks APP_PARTITIONS against the coverage counters in both
  directions, so a counted partition that nothing dispatches and a dispatchable
  partition with no counter both fail instead of silently skipping tests.
- TestCIAppRacePartitionMatrixMatchesHelper pins the workflow's app-<partition>
  shards to list-partitions output in both directions, so a partition cannot
  lose its job while every job stays green.

The discovery loop variable is renamed from partition to spec: it would
otherwise shadow the partition requested on the command line, which run mode
reads after the discovery pass completes.
2026-08-16 22:18:04 +08:00
chichuan 4a33e7e893 fix(ci): drop race instrumentation from the app schema partition
The schema partition's 52 tests assert structural Schema-to-Cobra contracts over
a single goroutine: none of them call t.Parallel or start a goroutine, so the
race detector has no concurrent access to observe there. The process-global lazy
metadata that does need race coverage (schema_source_root's atomic.Value, the
parameter-binding lazy loaders) is exercised by internal/cli's concurrent tests,
which stay instrumented.

The instrumentation was not free here. The partition shares a single sync.Once
Catalog build whose work is allocation-heavy, and -race made it roughly 11x
slower: 26s -> 291s locally, and 357s of the app shard's 780s in CI. Within that
partition TestFinalSchemaToolsHaveExecutableBaseCommands alone accounted for
262s, not because the test is expensive but because it is the first caller to pay
for the shared snapshot; its 1121 subtests together measure 0.00s.

run_partition now takes the instrumentation mode explicitly and fails closed on
an unrecognized value, so a typo cannot silently drop -race from a partition that
is supposed to carry it.
2026-08-16 22:17:14 +08:00
github-actions[bot] ee74765383 Merge pull request #1019 from DingTalk-Real-AI/feat/help-feedback-entry
feat: add feedback survey entry to root help
2026-08-16 08:09:01 +08:00
xiatian 9fbd8addbe Merge remote-tracking branch 'upstream/main' into codex/sheet-dropdown-source-range 2026-08-15 13:40:55 +08:00
xiatian 92195a58a3 fix(sheet): align SourceRange review contract 2026-08-15 13:40:47 +08:00
炳昱 e742a6c269 Merge remote-tracking branch 'upstream/main' into codex/fix-windows-event-bus 2026-08-14 19:40:58 +08:00
chichuan 05868610f0 Merge branch 'main' into codex/sheet-dropdown-source-range 2026-08-14 18:33:24 +08:00
炳昱 22649e96ef test(event): cover Unix spawn validation on Windows 2026-08-14 18:11:42 +08:00
炳昱 f68a11f11d test(event): cover Windows lifecycle edges 2026-08-14 18:05:34 +08:00
炳昱 abe5129306 fix(event): support Windows bus lifecycle 2026-08-14 17:56:51 +08:00
xiatian 8cd2b0259d Merge remote-tracking branch 'upstream/main' into codex/sheet-dropdown-source-range 2026-08-13 18:04:51 +08:00
xiatian 4b8d94c24e ci: retry interrupted app race shard 2026-08-13 17:15:22 +08:00
xiatian 76e5a8c4d9 Merge remote-tracking branch 'upstream/main' into codex/sheet-dropdown-source-range 2026-08-13 16:37:12 +08:00
xiatian 6abffce4e5 fix(sheet): preserve dropdown schema compatibility 2026-08-13 16:13:30 +08:00
xiatian 86b78e45d7 Merge remote-tracking branch 'upstream/main' into codex/sheet-dropdown-source-range 2026-08-13 16:08:28 +08:00
xiatian 5065e4bfb6 Merge remote-tracking branch 'upstream/main' into codex/sheet-dropdown-source-range 2026-08-13 13:49:52 +08:00
xiatian 2778bef5bd feat(sheet): support source range dropdowns and read completion 2026-08-13 11:46:58 +08:00
72 changed files with 4394 additions and 304 deletions
@@ -0,0 +1,11 @@
---
category: Fixed
---
- **Aitable pagination and Minutes unshare verification** (#1006) — keeps
record queries on the service's 20-record page boundary so multi-page reads
and mutation readbacks no longer report false retryable failures, preserves
`totalCount` when supplied, validates `--dry-run` plans before transport,
follows active deletion readback continuations before proving absence, and
rejects Minutes unshare success until the listening note exists and the
service acknowledges the exact task and member targets.
@@ -0,0 +1,5 @@
---
category: Fixed
---
- **Document write verification** (#960) — avoids false partial-success results when normalized Markdown, paginated blocks, inline images, or version reverts are confirmed by server readback. Document reverts and media inserts now require explicit readback evidence and report partial success when the server cannot prove the requested result.
@@ -0,0 +1,8 @@
---
category: Changed
---
- **Faster Schema Catalog assembly** — projects typed values into payload JSON
without re-running a validation scan over documents `json.Marshal` has just
produced, cutting roughly a third of the projection work across the full tool
set. Untrusted JSON input keeps its existing validation.
+6
View File
@@ -0,0 +1,6 @@
---
category: Added
---
- **Sheet SourceRange dropdowns** — supports range-backed dropdowns across direct, cell, and batch write paths, with structured readback for valid and invalid references. Batch `set-dropdown` now rejects unsupported top-level `colors` / `source-colors`; Inline colors belong in `options[].color`, while SourceRange color writes remain unsupported.
- **Sheet read completion metadata** — documents and preserves returned ranges, truncation reasons, and partial-read status for large range and CSV reads.
+5
View File
@@ -0,0 +1,5 @@
---
category: Fixed
---
- **Windows event bus lifecycle** — start event consumers without unsupported inherited file descriptors, stop buses through local IPC with a termination fallback, and preserve subscription cleanup when startup fails.
+63 -22
View File
@@ -523,12 +523,22 @@ jobs:
# shards at full-suite scope; release-scripts is included because its
# dedicated job only runs at full-suite or release-sensitive scope, and
# dropping it here would stop testing test/scripts changes entirely.
# internal/app is carried by one shard per bounded partition rather than a
# single app shard: the partitions used to run end to end inside one job,
# where the Schema partition alone owned most of the wall clock. The
# app-<partition> names are pinned to the helper's partition set by
# TestCIAppRacePartitionMatrixMatchesHelper, so a partition can never lose
# its job silently.
timeout-minutes: 20
strategy:
fail-fast: false
matrix:
shard:
- app
- app-schema
- app-a-b
- app-c
- app-d-r
- app-s-z-example-fuzz
- generators
- helpers
- cli
@@ -570,9 +580,16 @@ jobs:
TEST_SHARD: ${{ matrix.shard }}
run: |
set -euo pipefail
# Every app partition shard tests the same single internal/app
# package, so the impacted-package query uses the base shard name and
# the partition only selects which tests run.
package_shard="$TEST_SHARD"
case "$TEST_SHARD" in
app-*) package_shard=app ;;
esac
package_output="$(
./scripts/ci/changed-test-packages.sh \
list-shard "$TEST_SHARD" "$TEST_BASE_REF" "$TEST_HEAD_REF"
list-shard "$package_shard" "$TEST_BASE_REF" "$TEST_HEAD_REF"
)"
if [ -z "$package_output" ]; then
echo "No buildable Go package in shard $TEST_SHARD is affected by this revision." \
@@ -612,14 +629,19 @@ jobs:
exit 1
}
done
if [ "$TEST_SHARD" = "app" ]; then
# A single long-lived app test process retains every constructed
# command tree in framework registries. Isolate Schema assembly and
# bounded name ranges so each process releases that state on exit.
test "${#packages[@]}" -eq 1
./scripts/ci/run-app-race-tests.sh run "${packages[0]}"
exit 0
fi
case "$TEST_SHARD" in
app-*)
# A single long-lived app test process retains every constructed
# command tree in framework registries. Each partition is its own
# job, so that state is released when the process exits and the
# partitions run concurrently instead of end to end. The helper
# still verifies that the partition patterns cover every top-level
# test exactly once before running the one it was asked for.
test "${#packages[@]}" -eq 1
./scripts/ci/run-app-race-tests.sh run "${packages[0]}" "${TEST_SHARD#app-}"
exit 0
;;
esac
if [ "$TEST_SHARD" = "release-scripts" ]; then
# Mirror the dedicated release-contract job: these suites shell out
# to archive tooling and are not race-instrumented there.
@@ -640,14 +662,21 @@ jobs:
needs: lint
if: ${{ needs.lint.outputs.changelog_only != 'true' && needs.lint.outputs.docs_only != 'true' && needs.lint.outputs.full_suite == 'true' }}
runs-on: ubuntu-latest
# app runs several independently bounded processes; cli/smoke need headroom
# beyond go test -timeout for setup + assembly.
# internal/app is split across one shard per bounded partition so the
# partitions run concurrently and each releases its framework registries
# when the process exits; cli/smoke need headroom beyond go test -timeout for
# setup + assembly. The app-<partition> names are pinned to the helper's
# partition set by TestCIAppRacePartitionMatrixMatchesHelper.
timeout-minutes: 20
strategy:
fail-fast: false
matrix:
shard:
- app
- app-schema
- app-a-b
- app-c
- app-d-r
- app-s-z-example-fuzz
- generators
- helpers
- cli
@@ -673,18 +702,30 @@ jobs:
TEST_SHARD: ${{ matrix.shard }}
run: |
set -euo pipefail
package_output="$(./scripts/ci/test-packages.sh list "$TEST_SHARD")"
# Every app partition shard tests the same single internal/app
# package, so the package query uses the base shard name and the
# partition only selects which tests run.
package_shard="$TEST_SHARD"
case "$TEST_SHARD" in
app-*) package_shard=app ;;
esac
package_output="$(./scripts/ci/test-packages.sh list "$package_shard")"
test -n "$package_output"
mapfile -t packages <<< "$package_output"
test "${#packages[@]}" -gt 0
if [ "$TEST_SHARD" = "app" ]; then
# A single long-lived app test process retains every constructed
# command tree in framework registries. Isolate Schema assembly and
# bounded name ranges so each process releases that state on exit.
test "${#packages[@]}" -eq 1
./scripts/ci/run-app-race-tests.sh run "${packages[0]}"
exit 0
fi
case "$TEST_SHARD" in
app-*)
# A single long-lived app test process retains every constructed
# command tree in framework registries. Each partition is its own
# job, so that state is released when the process exits and the
# partitions run concurrently instead of end to end. The helper
# still verifies that the partition patterns cover every top-level
# test exactly once before running the one it was asked for.
test "${#packages[@]}" -eq 1
./scripts/ci/run-app-race-tests.sh run "${packages[0]}" "${TEST_SHARD#app-}"
exit 0
;;
esac
# cli/smoke own heavy NewRootCommand / Schema assembly under -race;
# give them a dedicated package timeout on slower hosted runners.
timeout_budget=12m
+31
View File
@@ -6,6 +6,37 @@ The format is inspired by [Keep a Changelog](https://keepachangelog.com/) and th
## [Unreleased]
## [1.0.59-beta.2] - 2026-08-17
### Added
- **Privacy-safe CLI telemetry** (#1009) — reports reviewed command outcomes and profile identity dimensions while excluding command arguments, output, paths, device fingerprints, and automatic system dimensions; `DO_NOT_TRACK=1` disables reporting.
- **Feedback survey entry in root help** (#1019) — `dws --help` now closes with a Feedback section linking the user-experience survey form.
- **Wiki Shortcut workflows** — publishes 20 reviewed space, member, node, and
activity shortcuts with strict collection validation, cursor handling,
write-terminal evidence, safe read-backs where the backend supports them,
task-oriented routing, and documented backend
boundaries.
### Changed
- **Chat IM ID flags** (#954) — standardizes chat command entry points on `--conversation-id` for conversation IDs and `--message-id` for message IDs, so help, Schema, and Agent recommendations use the same canonical flags.
- **Legacy chat flag compatibility** (#954) — keeps older chat IM ID flags such as `--group`, `--id`, `--chat`, `--open-conversation-id`, `--msg-id`, and `--open-message-id` working as compatibility aliases where applicable, while hiding migrated aliases from recommended help and Schema surfaces.
- **Chat group bots target flag** (#954) — keeps `dws chat group bots` on the visible `--group` flag; this command does not register `--group-name`, and `--group` accepts either an openConversationId or a uniquely resolved group name.
- **Faster Schema Catalog assembly** — projects typed values into payload JSON
without re-running a validation scan over documents `json.Marshal` has just
produced, cutting roughly a third of the projection work across the full tool
set. Untrusted JSON input keeps its existing validation.
### Fixed
- **Chat card update evidence** — distinguishes an accepted update request from an independently verified visible update, preserving the real `bizId` and warning callers not to repeat an unverified write.
- **Chat command guidance** — splits message and group references by task and explains that `--from` is ambiguous between sender and time-range intent.
## [1.0.59-beta.1] - 2026-08-14
### Added
+11 -11
View File
@@ -1,33 +1,33 @@
class DingtalkWorkspaceCliBeta < Formula
desc "Automate DingTalk workspace tasks from the terminal (beta channel)"
homepage "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli"
version "1.0.59-beta.1"
version "1.0.59-beta.2"
license "Apache-2.0"
keg_only "it is the beta channel and conflicts with dingtalk-workspace-cli"
on_macos do
if Hardware::CPU.arm?
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.1/dws-darwin-arm64.tar.gz"
sha256 "36a30f3496e0f759c15c0b09f67dbd23b8ecdfff2eebe572f88125b26485830f"
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.2/dws-darwin-arm64.tar.gz"
sha256 "7f11218d3222f3e93c3b1e94b3a004c061eb0a297b206447ea95fe6b2b1ec674"
else
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.1/dws-darwin-amd64.tar.gz"
sha256 "e7a04906380efd8da88cd112e6a512bb6470a3956dc370150037ed6e314db445"
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.2/dws-darwin-amd64.tar.gz"
sha256 "72f06a334cf29d23123639fabe13bf2951765066136e47ccc6453667c382f8c2"
end
end
on_linux do
if Hardware::CPU.arm?
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.1/dws-linux-arm64.tar.gz"
sha256 "f59ab055f3e841e4cebc964ae3ef969668475548abaaf8bede44afdca9a3e28d"
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.2/dws-linux-arm64.tar.gz"
sha256 "3edbfabb7718b53914a2d9efe7b1152e9ebd70765ff0b6f19ad3125e4a2458ed"
else
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.1/dws-linux-amd64.tar.gz"
sha256 "2c8f919489d958c7d49262615e81faac70a9fbcae2d589ab54a0bb3c5700a057"
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.2/dws-linux-amd64.tar.gz"
sha256 "e1c610070a9c1b3763656cbd53c818290f199097adc07cb9fe629ed765c5e1e0"
end
end
resource "skills" do
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.1/dws-skills.zip"
sha256 "25f4a7e1d01fa4d771d79201b34b11ee8a24182bdcdc94bfb98d2bd5845bed3b"
url "https://github.com/DingTalk-Real-AI/dingtalk-workspace-cli/releases/download/v1.0.59-beta.2/dws-skills.zip"
sha256 "aa2854651eaa2c857b526aaec73fd1358d31b11999fe7fd3e99e5060b2522a1f"
end
def install
+1
View File
@@ -18,6 +18,7 @@ require (
github.com/muesli/termenv v0.16.0
github.com/open-dingtalk/dingtalk-stream-sdk-go v0.9.2-beta.1
github.com/spf13/cobra v1.10.2
github.com/yuin/goldmark v1.8.5
github.com/zalando/go-keyring v0.2.8
gitlab.alibaba-inc.com/aes/aem-go-sdk v0.3.0
golang.org/x/crypto v0.49.0
+2
View File
@@ -105,6 +105,8 @@ github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e h1:JVG44RsyaB9T2KIHavMF/ppJZNG9ZpyihvCd0w101no=
github.com/xo/terminfo v0.0.0-20220910002029-abceb7e1c41e/go.mod h1:RbqR21r5mrJuqunuUZ/Dhy/avygyECGrLceyNeo4LiM=
github.com/yuin/goldmark v1.8.5 h1:r6N5afV5qj/5S4UTch8agZHJ8UxNCMwX7WjkkJam2NA=
github.com/yuin/goldmark v1.8.5/go.mod h1:ip/1k0VRfGynBgxOz0yCqHrbZXhcjxyuS66Brc7iBKg=
github.com/zalando/go-keyring v0.2.8 h1:6sD/Ucpl7jNq10rM2pgqTs0sZ9V3qMrqfIIy5YPccHs=
github.com/zalando/go-keyring v0.2.8/go.mod h1:tsMo+VpRq5NGyKfxoBVjCuMrG47yj8cmakZDO5QGii0=
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
+2 -1
View File
@@ -1220,7 +1220,8 @@ func newEventStopCommandWithFlags(globalFlags ...*GlobalFlags) *cobra.Command {
editionName := editionNameOrDefault()
clientIDHash := dwsevent.ClientIDHash(clientID)
workDir := eventWorkDir(configDir, editionName, dwsevent.SourceKindAppStream, clientIDHash)
if err := eventStopBus(busctl.StopConfig{WorkDir: workDir}); err != nil {
ipcEndpoint := defaultIPCEndpoint(workDir, editionName, dwsevent.SourceKindAppStream, clientIDHash)
if err := eventStopBus(busctl.StopConfig{WorkDir: workDir, IPCEndpoint: ipcEndpoint}); err != nil {
if errors.Is(err, busctl.ErrNotRunning) {
fmt.Fprintln(c.OutOrStdout(), "bus is not running")
return nil
+1 -1
View File
@@ -1220,7 +1220,7 @@ func runPersonalEventStop(c *cobra.Command, opts personalStopOptions) error {
}
busState := "personal bus stopped"
if err := personalStopBus(busctl.StopConfig{WorkDir: workDir}); err != nil {
if err := personalStopBus(busctl.StopConfig{WorkDir: workDir, IPCEndpoint: ipcEndpoint}); err != nil {
if errors.Is(err, busctl.ErrNotRunning) {
busState = "personal bus is not running"
} else {
+3 -1
View File
@@ -52,6 +52,8 @@ func (c *paramAliasCaptureCaller) paramAliasResponseForTool(tool string) string
switch tool {
case "list_calendar_events":
return `{"result":{"events":[]}}`
case "query_records":
return `{"success":true,"status":"success","error":{},"data":{}}`
case "search_mail_users":
return `{"users":[{"name":"Fixture User","email":"fixture@example.com","id":"fixture-user"}]}`
case "search_dept_by_keyword":
@@ -63,7 +65,7 @@ func (c *paramAliasCaptureCaller) paramAliasResponseForTool(tool string) string
case "list_doc_versions":
return `{"result":{"items":[{"version":3}]}}`
case "revert_doc_version":
return `{"version":3}`
return `{"revertedToVersion":3}`
case "search_doc_templates":
return `{"result":[{"templateId":"fixture-template-id"}]}`
case "create_document":
+18 -2
View File
@@ -1127,21 +1127,37 @@ func putRawJSON(payload map[string]any, key string, raw json.RawMessage) error {
return nil
}
// rawJSONValue decodes a JSON value that may come from an untrusted source, so
// it validates before decoding. Callers holding output that json.Marshal just
// produced should use typedJSONValue instead of paying the validation scan.
func rawJSONValue(raw json.RawMessage) (any, error) {
if !json.Valid(raw) {
return nil, fmt.Errorf("invalid JSON value")
}
return decodeValidJSONValue(raw), nil
}
// decodeValidJSONValue decodes JSON whose validity the caller has already
// established, either by json.Valid or by having just marshaled it. Decode
// errors are unreachable under that precondition and are therefore discarded,
// exactly as this path behaved when the decode was inlined into rawJSONValue.
func decodeValidJSONValue(raw json.RawMessage) any {
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.UseNumber()
var value any
_ = decoder.Decode(&value)
return value, nil
return value
}
// typedJSONValue projects a typed value into the generic JSON shape the payload
// renderers consume. json.Marshal output is valid by construction, so this path
// decodes it directly: routing through rawJSONValue re-scanned every marshaled
// document with json.Valid, which measured ~34% of Schema Catalog assembly time
// across the 1121-tool set (26.0s -> 17.2s for the internal/app schema suite).
func typedJSONValue(value any) (any, error) {
data, err := json.Marshal(value)
if err != nil {
return nil, err
}
return rawJSONValue(data)
return decodeValidJSONValue(data), nil
}
+8 -3
View File
@@ -361,13 +361,18 @@ func TestCrossPlatformCoverageDaemonMethodEdges(t *testing.T) {
}
return nil, net.ErrClosed
}}
d := &daemon{listener: l, log: logger, hub: NewHub(1), idleStop: make(chan struct{})}
d := &daemon{listener: l, log: logger, hub: NewHub(1), idleStop: make(chan struct{}), stopReq: make(chan string, 1)}
d.acceptLoop(context.Background())
d.shuttingDown.Store(true)
d.acceptLoop(context.Background())
d.triggerShutdown("test")
if !l.closed {
t.Fatal("trigger did not close listener")
select {
case reason := <-d.stopReq:
if reason != "test" {
t.Fatalf("stop reason = %q", reason)
}
default:
t.Fatal("trigger did not notify lifecycle loop")
}
d = &daemon{listener: &scriptedListener{accept: func() (net.Conn, error) { return nil, net.ErrClosed }}, log: logger, hub: NewHub(1), idleStop: make(chan struct{})}
+14 -7
View File
@@ -186,6 +186,7 @@ func Run(ctx context.Context, cfg Config) error {
dedup: dd,
started: time.Now().UTC(),
idleStop: make(chan struct{}),
stopReq: make(chan string, 1),
}
// 4. Signal ready BEFORE accepting consumers (avoids a slow-fork
@@ -263,6 +264,9 @@ func Run(ctx context.Context, cfg Config) error {
}
case <-d.idleStop:
log.Info("bus: idle timeout reached, shutting down")
case reason := <-d.stopReq:
shutdownReason = reason
log.Info("bus: shutdown requested via IPC", "reason", reason)
}
// 7. Graceful shutdown — cancel runCtx first so all background
@@ -292,6 +296,7 @@ type daemon struct {
shutdownMu sync.Mutex
shuttingDown atomic.Bool
idleStop chan struct{}
stopReq chan string
credentialHandoffMu sync.Mutex
terminalMu sync.RWMutex
@@ -749,15 +754,17 @@ func (d *daemon) idleWatch(ctx context.Context) {
}
}
// triggerShutdown is called from RPC handlers that want to end the bus.
// It works by closing the listener (which unblocks Run's select via the
// source error path, indirectly). For v1 a full ctx-cancellation hook is
// out of scope; busctl/stop also sends SIGTERM which is the authoritative
// shutdown path.
// triggerShutdown is called from RPC handlers that want to end the bus. The
// buffered request wakes Run's lifecycle select, which then cancels the cloud
// source and performs the same graceful shutdown used for signals and idle
// expiry. It is intentionally non-blocking so repeated stop requests cannot
// strand connection handlers.
func (d *daemon) triggerShutdown(reason string) {
d.log.Info("bus: shutdown triggered via IPC", "reason", reason)
_ = d.listener.Close() // unblocks Accept(), but doesn't kill Source
// Best-effort: a future version wires a context.CancelFunc here.
select {
case d.stopReq <- reason:
default:
}
}
// shutdown performs the graceful tear-down sequence:
+61
View File
@@ -141,6 +141,67 @@ func TestDaemon_RunStartsAndShutsDownCleanly(t *testing.T) {
}
}
func TestCrossPlatformCoverageDaemonIPCStopShutsDownCleanly(t *testing.T) {
workDir := shortTempDir(t)
identityHash := dwsevent.IdentityHash(workDir)
endpoint := dwsevent.IPCEndpoint(workDir, "open", dwsevent.SourceKindAppStream, identityHash)
runDone := make(chan error, 1)
go func() {
runDone <- Run(context.Background(), Config{
WorkDir: workDir,
IPCEndpoint: endpoint,
ClientID: "ding_ipc_stop_test",
Edition: "open",
SourceKind: dwsevent.SourceKindAppStream,
IdentityHash: identityHash,
Source: &fakeSource{},
})
}()
var conn net.Conn
deadline := time.Now().Add(3 * time.Second)
for time.Now().Before(deadline) {
var err error
conn, err = transport.Dial(endpoint)
if err == nil {
break
}
time.Sleep(10 * time.Millisecond)
}
if conn == nil {
t.Fatal("bus IPC endpoint did not become ready")
}
w := transport.NewWriter(conn)
r := transport.NewReader(conn)
if err := w.WriteJSON(transport.Hello{
Type: transport.FrameTypeHello,
ConsumerPID: os.Getpid(),
Role: transport.HelloRoleStop,
}); err != nil {
t.Fatalf("write stop hello: %v", err)
}
var bye transport.Bye
if err := r.ReadJSON(&bye); err != nil {
t.Fatalf("read stop response: %v", err)
}
_ = conn.Close()
if bye.Type != transport.FrameTypeBye || bye.Reason != "stop_request" {
t.Fatalf("stop response = %#v", bye)
}
select {
case err := <-runDone:
if err != nil {
t.Fatalf("Run returned after IPC stop: %v", err)
}
case <-time.After(3 * time.Second):
t.Fatal("Run did not return after IPC stop")
}
if pid := ReadHolderPID(filepath.Join(workDir, LockFileName)); pid != 0 {
t.Fatalf("bus lock retained pid %d after IPC stop", pid)
}
}
func TestDaemon_ConsumerReceivesEvents(t *testing.T) {
skipOnWindows(t, "uses Unix socket dial")
workDir := shortTempDir(t)
+41
View File
@@ -114,6 +114,47 @@ func ReadHolderPID(path string) int {
return parsePID(b)
}
// ValidateHolderOwnership reports whether pid still owns the exclusive lock at
// path. If the lock can be acquired, no daemon owns it: the PID payload is
// stale and is cleared while the probe holds the lock. The PID is read again
// after a busy result so a daemon replacement cannot be mistaken for the
// original owner.
//
// Callers use this immediately before sending a process-level termination
// fallback. A live PID by itself is not proof of ownership because operating
// systems may reuse PIDs after an unclean daemon exit.
func ValidateHolderOwnership(path string, pid int) (bool, error) {
if pid <= 0 || ReadHolderPID(path) != pid {
return false, nil
}
probe, err := busTryAcquire(path)
if err != nil {
if errors.Is(err, ErrBusy) {
return ReadHolderPID(path) == pid, nil
}
return false, fmt.Errorf("bus: verify lock owner: %w", err)
}
defer probe.Close()
// We own the lock, so the recorded PID cannot own it. Clear only the PID
// we validated above; a changed payload is left untouched defensively.
f := probe.File()
if _, err := busSeek(f, 0, io.SeekStart); err != nil {
return false, fmt.Errorf("bus: seek stale lock: %w", err)
}
current, err := busReadAll(f)
if err != nil {
return false, fmt.Errorf("bus: read stale lock: %w", err)
}
if parsePID(current) == pid {
if err := truncateAndWritePID(f, 0); err != nil {
return false, fmt.Errorf("bus: clear stale PID: %w", err)
}
}
return false, nil
}
// Close releases the flock and best-effort blanks the PID body so a stale
// reader (e.g. `event status` racing our shutdown) does not see our
// long-dead PID and try to signal it. The lock file itself is NOT removed
+107
View File
@@ -16,9 +16,13 @@ package bus
import (
"errors"
"fmt"
"io"
"os"
"path/filepath"
"testing"
eventlock "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/lock"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/testseam"
)
func TestAcquire_WritesOurPID(t *testing.T) {
@@ -112,6 +116,109 @@ func TestReadHolderPID_MalformedReturnsZero(t *testing.T) {
}
}
func TestCrossPlatformCoverageValidateHolderOwnershipHeldLock(t *testing.T) {
path := filepath.Join(t.TempDir(), LockFileName)
held, err := busTryAcquire(path)
if err != nil {
t.Fatalf("hold lock: %v", err)
}
defer held.Close()
if err := truncateAndWritePID(held.File(), os.Getpid()); err != nil {
t.Fatalf("write holder PID: %v", err)
}
owned, err := ValidateHolderOwnership(path, os.Getpid())
if err != nil {
t.Fatalf("ValidateHolderOwnership: %v", err)
}
if !owned {
t.Fatal("held lock was not attributed to its recorded PID")
}
}
func TestCrossPlatformCoverageValidateHolderOwnershipClearsReusedPID(t *testing.T) {
path := filepath.Join(t.TempDir(), LockFileName)
// The current PID is alive but deliberately does not hold this lock. This
// models a stale daemon PID that the OS has reassigned to another process.
if err := os.WriteFile(path, []byte(fmt.Sprintf("%d\n", os.Getpid())), 0o600); err != nil {
t.Fatalf("write stale PID: %v", err)
}
owned, err := ValidateHolderOwnership(path, os.Getpid())
if err != nil {
t.Fatalf("ValidateHolderOwnership: %v", err)
}
if owned {
t.Fatal("live reused PID without the bus lock was accepted as owner")
}
if got := ReadHolderPID(path); got != 0 {
t.Fatalf("stale PID after validation = %d, want cleared", got)
}
}
func TestCrossPlatformCoverageValidateHolderOwnershipRejectsMismatch(t *testing.T) {
path := filepath.Join(t.TempDir(), LockFileName)
if err := os.WriteFile(path, []byte("123\n"), 0o600); err != nil {
t.Fatal(err)
}
owned, err := ValidateHolderOwnership(path, 456)
if err != nil || owned {
t.Fatalf("ValidateHolderOwnership mismatch = %v, %v", owned, err)
}
}
func TestCrossPlatformCoverageValidateHolderOwnershipErrors(t *testing.T) {
errInjected := errors.New("injected ownership validation failure")
writeLivePID := func(t *testing.T) string {
t.Helper()
path := filepath.Join(t.TempDir(), LockFileName)
if err := os.WriteFile(path, []byte(fmt.Sprintf("%d\n", os.Getpid())), 0o600); err != nil {
t.Fatal(err)
}
return path
}
t.Run("probe", func(t *testing.T) {
path := writeLivePID(t)
testseam.Swap(t, &busTryAcquire, func(string) (*eventlock.File, error) {
return nil, errInjected
})
if _, err := ValidateHolderOwnership(path, os.Getpid()); !errors.Is(err, errInjected) {
t.Fatalf("probe error = %v", err)
}
})
t.Run("seek", func(t *testing.T) {
path := writeLivePID(t)
testseam.Swap(t, &busSeek, func(*os.File, int64, int) (int64, error) {
return 0, errInjected
})
if _, err := ValidateHolderOwnership(path, os.Getpid()); !errors.Is(err, errInjected) {
t.Fatalf("seek error = %v", err)
}
})
t.Run("read", func(t *testing.T) {
path := writeLivePID(t)
testseam.Swap(t, &busReadAll, func(io.Reader) ([]byte, error) {
return nil, errInjected
})
if _, err := ValidateHolderOwnership(path, os.Getpid()); !errors.Is(err, errInjected) {
t.Fatalf("read error = %v", err)
}
})
t.Run("clear", func(t *testing.T) {
path := writeLivePID(t)
testseam.Swap(t, &busTruncate, func(*os.File, int64) error {
return errInjected
})
if _, err := ValidateHolderOwnership(path, os.Getpid()); !errors.Is(err, errInjected) {
t.Fatalf("clear error = %v", err)
}
})
}
func TestAcquire_AfterReleaseReclaimable(t *testing.T) {
path := filepath.Join(t.TempDir(), LockFileName)
l1, err := Acquire(path)
@@ -284,11 +284,13 @@ func TestCrossPlatformCoverageQueryStatusAndEntryEdges(t *testing.T) {
func TestCrossPlatformCoverageStopInjectedEdges(t *testing.T) {
origRead := stopReadHolderPID
origValidateOwner := stopValidateHolderOwner
origAlive := stopAlive
origFind := stopFindProcess
origSignal := stopSignalProcess
t.Cleanup(func() {
stopReadHolderPID = origRead
stopValidateHolderOwner = origValidateOwner
stopAlive = origAlive
stopFindProcess = origFind
stopSignalProcess = origSignal
@@ -307,6 +309,7 @@ func TestCrossPlatformCoverageStopInjectedEdges(t *testing.T) {
t.Fatal(err)
}
stopFindProcess = func(int) (*os.Process, error) { return proc, nil }
stopValidateHolderOwner = func(string, int) (bool, error) { return true, nil }
stopSignalProcess = func(*os.Process, os.Signal) error { return os.ErrProcessDone }
if err := Stop(StopConfig{WorkDir: "x"}); err != nil {
t.Fatalf("done signal = %v", err)
+3 -2
View File
@@ -112,8 +112,9 @@ func Discover(cfg DiscoverConfig) (net.Conn, error) {
return nil, fmt.Errorf("busctl: mkdir workdir: %w", err)
}
_, spawnErr := cfg.Spawn(SpawnConfig{
ClientID: cfg.ClientID,
ExtraArgs: cfg.SpawnExtraArgs,
ClientID: cfg.ClientID,
IPCEndpoint: cfg.IPCEndpoint,
ExtraArgs: cfg.SpawnExtraArgs,
})
if spawnErr != nil && !errors.Is(spawnErr, ErrSpawnFailed) {
// Hard error (couldn't even exec the child). Stop here — no bus
+4 -1
View File
@@ -122,7 +122,10 @@ func TestDiscover_NoBus_SpawnSucceeds(t *testing.T) {
closer()
}
})
fakeSpawn := func(SpawnConfig) (int, error) {
fakeSpawn := func(cfg SpawnConfig) (int, error) {
if cfg.IPCEndpoint != sock {
t.Fatalf("spawn IPC endpoint = %q, want %q", cfg.IPCEndpoint, sock)
}
closer = startStubBus(t, sock)
return 12345, nil
}
+15 -14
View File
@@ -43,18 +43,14 @@ var (
spawnPipe = os.Pipe
)
// ErrSpawnFailed is returned when the child reports startup failure via
// the ready pipe ('E' byte). The child's exit error / log file holds the
// actual cause; this sentinel just lets the caller distinguish "ready
// pipe said no" from "ready pipe timed out / closed early".
var ErrSpawnFailed = errors.New("busctl: bus child reported startup failure on ready pipe")
// ErrSpawnFailed is returned when the child reports startup failure through
// the Unix ready pipe or exits before binding the Windows named pipe.
var ErrSpawnFailed = errors.New("busctl: bus child reported startup failure")
// ErrSpawnTimeout is returned when ReadyTimeout elapses without any signal.
var ErrSpawnTimeout = errors.New("busctl: bus child did not signal readiness within deadline")
// SpawnConfig describes one spawn attempt. ClientID is the only field
// inspected by the child; the rest govern process attributes the parent
// applies before exec.
// SpawnConfig describes one spawn attempt.
type SpawnConfig struct {
// ExecPath is the dws binary to exec. Default os.Executable().
ExecPath string
@@ -62,18 +58,23 @@ type SpawnConfig struct {
// ClientID is passed as `--client-id` to `dws event _bus`. Required.
ClientID string
// IPCEndpoint is the bus endpoint the child binds. Windows uses the
// existing named pipe as its readiness handshake because os/exec does not
// support ExtraFiles there. Unix keeps the inherited ready-pipe protocol.
IPCEndpoint string
// ExtraArgs are appended after `--client-id`. Empty for normal use; tests
// pass `--extra-flag-for-test` etc.
ExtraArgs []string
// Env to pass to the child. Defaults to os.Environ(). The ReadyFDEnv
// entry is appended automatically.
// Env to pass to the child. Defaults to os.Environ(). Unix appends the
// ReadyFDEnv entry automatically; Windows removes any inherited copy.
Env []string
}
// Spawn forks a detached `dws event _bus --client-id <id>` child process and
// waits for it to signal readiness via the ready pipe. Returns the child's
// PID on success — the caller can then dial the bus IPC endpoint.
// spawnWithReadyPipe is the Unix implementation behind Spawn. It forks a
// detached `dws event _bus --client-id <id>` child and waits for the inherited
// ready pipe. Windows has a named-pipe implementation in spawn_windows.go.
//
// stdio detach (plan invariant #7):
// - cmd.Stdout / cmd.Stderr set to nil so the child's own writes don't
@@ -89,7 +90,7 @@ type SpawnConfig struct {
// Parent (this function):
// - Holds the read end open until either 1 byte is read or ReadyTimeout
// - Returns ErrSpawnFailed for 'E', ErrSpawnTimeout otherwise
func Spawn(cfg SpawnConfig) (pid int, err error) {
func spawnWithReadyPipe(cfg SpawnConfig) (pid int, err error) {
if cfg.ClientID == "" {
return 0, errors.New("busctl: SpawnConfig.ClientID is required")
}
+24 -1
View File
@@ -22,6 +22,8 @@ import (
"strings"
"testing"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/transport"
)
// Test child mode: when DWS_BUSCTL_TEST_CHILD is set, this test binary
@@ -35,7 +37,8 @@ import (
// env-marker pattern is what Go's own os/exec tests use and stays
// confined to this file.
const (
childEnvMarker = "DWS_BUSCTL_TEST_CHILD"
childEnvMarker = "DWS_BUSCTL_TEST_CHILD"
childEndpointEnv = "DWS_BUSCTL_TEST_ENDPOINT"
// values:
// "ready" — write 'R' then sleep 30s (parent should see ready)
// "fail" — write 'E' then exit (parent should see ErrSpawnFailed)
@@ -70,6 +73,26 @@ func TestMain(m *testing.M) {
writeReady('R')
time.Sleep(30 * time.Second)
os.Exit(0)
case "windows-ready":
listener, err := transport.Listen(os.Getenv(childEndpointEnv))
if err != nil {
os.Exit(3)
}
defer listener.Close()
conn, err := listener.Accept()
if err != nil {
os.Exit(4)
}
_ = conn.Close()
time.Sleep(30 * time.Second)
os.Exit(0)
case "windows-fail":
os.Exit(5)
case "windows-exit":
os.Exit(0)
case "windows-stall":
time.Sleep(30 * time.Second)
os.Exit(0)
}
os.Exit(m.Run())
}
+7
View File
@@ -20,6 +20,13 @@ import (
"syscall"
)
// Spawn starts the detached bus and waits for its inherited ready pipe.
// ExtraFiles is intentionally confined to Unix: Go does not support it on
// Windows.
func Spawn(cfg SpawnConfig) (pid int, err error) {
return spawnWithReadyPipe(cfg)
}
// applyDetach configures the child to live past parent death and not share
// the parent's controlling terminal. Setsid puts the child in a new
// session, so SIGHUP on the parent's controlling tty (e.g. SSH disconnect)
+87
View File
@@ -16,14 +16,101 @@
package busctl
import (
"fmt"
"os"
"os/exec"
"strings"
"syscall"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/transport"
)
// CREATE_NEW_PROCESS_GROUP (0x00000200) prevents the child from receiving
// the parent's Ctrl+C signal, similar in spirit to Setsid on Unix.
const createNewProcessGroup = 0x00000200
var (
spawnWindowsDial = transport.Dial
spawnWindowsPollInterval = 25 * time.Millisecond
)
// Spawn starts a detached Windows bus without cmd.ExtraFiles (unsupported by
// Go on Windows). Readiness is confirmed by dialing the bus's existing named
// pipe. The child is reaped in the background after readiness, while an early
// process exit is surfaced as ErrSpawnFailed.
func Spawn(cfg SpawnConfig) (pid int, err error) {
if strings.TrimSpace(cfg.ClientID) == "" {
return 0, fmt.Errorf("busctl: SpawnConfig.ClientID is required")
}
if strings.TrimSpace(cfg.IPCEndpoint) == "" {
return 0, fmt.Errorf("busctl: SpawnConfig.IPCEndpoint is required on Windows")
}
if cfg.ExecPath == "" {
execPath, err := spawnExecutable()
if err != nil {
return 0, fmt.Errorf("busctl: locate executable: %w", err)
}
cfg.ExecPath = execPath
}
if cfg.Env == nil {
cfg.Env = os.Environ()
}
args := append([]string{"event", "_bus", "--client-id", cfg.ClientID}, cfg.ExtraArgs...)
cmd := exec.Command(cfg.ExecPath, args...)
cmd.Env = withoutReadyFDEnv(cfg.Env)
cmd.Stdin = nil
cmd.Stdout = nil
cmd.Stderr = nil
applyDetach(cmd)
if err := cmd.Start(); err != nil {
return 0, fmt.Errorf("busctl: start %s: %w", cfg.ExecPath, err)
}
pid = cmd.Process.Pid
waitDone := make(chan error, 1)
go func() { waitDone <- cmd.Wait() }()
timer := time.NewTimer(ReadyTimeout)
defer timer.Stop()
ticker := time.NewTicker(spawnWindowsPollInterval)
defer ticker.Stop()
for {
if conn, dialErr := spawnWindowsDial(cfg.IPCEndpoint); dialErr == nil {
_ = conn.Close()
return pid, nil
}
select {
case waitErr := <-waitDone:
if waitErr == nil {
return pid, fmt.Errorf("%w: child exited before binding named pipe", ErrSpawnFailed)
}
return pid, fmt.Errorf("%w: %v", ErrSpawnFailed, waitErr)
case <-timer.C:
_ = cmd.Process.Kill()
select {
case <-waitDone:
case <-time.After(time.Second):
}
return pid, ErrSpawnTimeout
case <-ticker.C:
}
}
}
func withoutReadyFDEnv(env []string) []string {
prefix := strings.ToUpper(ReadyFDEnv) + "="
out := make([]string, 0, len(env))
for _, entry := range env {
if strings.HasPrefix(strings.ToUpper(entry), prefix) {
continue
}
out = append(out, entry)
}
return out
}
func applyDetach(cmd *exec.Cmd) {
cmd.SysProcAttr = &syscall.SysProcAttr{
CreationFlags: createNewProcessGroup,
+132
View File
@@ -0,0 +1,132 @@
// Copyright 2026 Alibaba Group
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//go:build windows
package busctl
import (
"errors"
"fmt"
"os"
"strings"
"testing"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/testseam"
)
func TestCrossPlatformCoverageWindowsSpawnUsesNamedPipeReadiness(t *testing.T) {
endpoint := fmt.Sprintf(`\\.\pipe\dws-event-spawn-test-%d-%d`, os.Getpid(), time.Now().UnixNano())
pid, err := spawnWithMarker(t, "windows-ready", func(cfg *SpawnConfig) {
cfg.IPCEndpoint = endpoint
cfg.Env = append(cfg.Env, childEndpointEnv+"="+endpoint)
})
if err != nil {
t.Fatalf("Spawn() = %v", err)
}
if pid <= 0 {
t.Fatalf("Spawn() pid = %d", pid)
}
proc, err := os.FindProcess(pid)
if err != nil {
t.Fatal(err)
}
_ = proc.Kill()
}
func TestCrossPlatformCoverageWindowsSpawnReportsEarlyExit(t *testing.T) {
endpoint := fmt.Sprintf(`\\.\pipe\dws-event-spawn-fail-%d-%d`, os.Getpid(), time.Now().UnixNano())
_, err := spawnWithMarker(t, "windows-fail", func(cfg *SpawnConfig) {
cfg.IPCEndpoint = endpoint
})
if !errors.Is(err, ErrSpawnFailed) {
t.Fatalf("Spawn() error = %v, want ErrSpawnFailed", err)
}
}
func TestCrossPlatformCoverageWindowsSpawnReportsCleanEarlyExit(t *testing.T) {
endpoint := fmt.Sprintf(`\\.\pipe\dws-event-spawn-exit-%d-%d`, os.Getpid(), time.Now().UnixNano())
_, err := spawnWithMarker(t, "windows-exit", func(cfg *SpawnConfig) {
cfg.IPCEndpoint = endpoint
})
if !errors.Is(err, ErrSpawnFailed) {
t.Fatalf("Spawn() error = %v, want ErrSpawnFailed", err)
}
}
func TestCrossPlatformCoverageWindowsSpawnResolvesExecutable(t *testing.T) {
endpoint := fmt.Sprintf(`\\.\pipe\dws-event-spawn-resolve-%d-%d`, os.Getpid(), time.Now().UnixNano())
testseam.Swap(t, &spawnExecutable, func() (string, error) { return os.Args[0], nil })
pid, err := spawnWithMarker(t, "windows-ready", func(cfg *SpawnConfig) {
cfg.ExecPath = ""
cfg.IPCEndpoint = endpoint
cfg.Env = append(cfg.Env, childEndpointEnv+"="+endpoint)
})
if err != nil {
t.Fatalf("Spawn() = %v", err)
}
proc, err := os.FindProcess(pid)
if err != nil {
t.Fatal(err)
}
_ = proc.Kill()
}
func TestCrossPlatformCoverageWindowsSpawnValidationAndStartErrors(t *testing.T) {
endpoint := fmt.Sprintf(`\\.\pipe\dws-event-spawn-errors-%d-%d`, os.Getpid(), time.Now().UnixNano())
if _, err := Spawn(SpawnConfig{ClientID: "client"}); err == nil {
t.Fatal("missing endpoint unexpectedly succeeded")
}
testseam.Swap(t, &spawnExecutable, func() (string, error) {
return "", errors.New("executable unavailable")
})
if _, err := Spawn(SpawnConfig{ClientID: "client", IPCEndpoint: endpoint}); err == nil {
t.Fatal("executable lookup unexpectedly succeeded")
}
if _, err := Spawn(SpawnConfig{
ClientID: "client",
IPCEndpoint: endpoint,
ExecPath: `C:\\definitely-missing-dws.exe`,
}); err == nil {
t.Fatal("missing executable unexpectedly started")
}
}
func TestCrossPlatformCoverageWindowsSpawnReadyTimeoutTerminatesChild(t *testing.T) {
endpoint := fmt.Sprintf(`\\.\pipe\dws-event-spawn-stall-%d-%d`, os.Getpid(), time.Now().UnixNano())
testseam.Swap(t, &ReadyTimeout, 100*time.Millisecond)
_, err := spawnWithMarker(t, "windows-stall", func(cfg *SpawnConfig) {
cfg.IPCEndpoint = endpoint
})
if !errors.Is(err, ErrSpawnTimeout) {
t.Fatalf("Spawn() error = %v, want ErrSpawnTimeout", err)
}
}
func TestCrossPlatformCoverageWindowsSpawnRemovesInheritedReadyFD(t *testing.T) {
env := withoutReadyFDEnv([]string{
"PATH=C:\\Windows",
ReadyFDEnv + "=3",
strings.ToLower(ReadyFDEnv) + "=4",
})
if len(env) != 1 || !strings.HasPrefix(env[0], "PATH=") {
t.Fatalf("withoutReadyFDEnv() = %#v", env)
}
}
func TestCrossPlatformCoverageWindowsUnixReadyPipeValidation(t *testing.T) {
if _, err := spawnWithReadyPipe(SpawnConfig{}); err == nil {
t.Fatal("spawnWithReadyPipe() unexpectedly accepted an empty client ID")
}
}
+95 -24
View File
@@ -17,15 +17,17 @@ import (
"errors"
"fmt"
"os"
"strings"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/bus"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/process"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/transport"
)
// DefaultStopTimeout is the wall-clock budget Stop waits for the bus to
// exit after the signal is sent. 5s covers the bus's own graceful tear-down
// (broadcast Bye → consumer goroutines drain → cleanup) with margin.
// DefaultStopTimeout is the wall-clock budget Stop waits for each graceful or
// fallback termination phase. 5s covers the bus's own tear-down (broadcast
// Bye → consumer goroutines drain → cleanup) with margin.
const DefaultStopTimeout = 5 * time.Second
// ErrNotRunning indicates bus.lock either does not exist or its recorded
@@ -33,20 +35,32 @@ const DefaultStopTimeout = 5 * time.Second
// distinguish "nothing to stop" from "failed to stop".
var ErrNotRunning = errors.New("busctl: bus is not running")
// ErrOwnerUnverified indicates the PID from bus.lock is alive but no longer
// owns that lock. Stop refuses to send a process-level signal in this state
// because the operating system may have reused the stale PID.
var ErrOwnerUnverified = errors.New("busctl: bus process ownership could not be verified")
var (
stopReadHolderPID = bus.ReadHolderPID
stopAlive = process.Alive
stopFindProcess = os.FindProcess
stopSignalProcess = func(proc *os.Process, signal os.Signal) error { return proc.Signal(signal) }
stopReadHolderPID = bus.ReadHolderPID
stopValidateHolderOwner = bus.ValidateHolderOwnership
stopAlive = process.Alive
stopFindProcess = os.FindProcess
stopSignalProcess = func(proc *os.Process, signal os.Signal) error { return proc.Signal(signal) }
stopDial = transport.Dial
stopRequest = requestBusStop
stopWaitForBusExit = waitForBusExit
)
// StopConfig identifies the target bus and tunes timing.
type StopConfig struct {
// WorkDir holds bus.lock; Stop reads the PID from there.
WorkDir string
// Timeout is the total wall-clock budget for graceful exit. After this,
// Stop returns an error; it does NOT escalate to SIGKILL — leave that
// to the operator.
// IPCEndpoint enables the cross-platform graceful stop RPC. Callers that
// know the bus identity should always provide it. Empty preserves the
// legacy signal-only fallback for compatibility and focused tests.
IPCEndpoint string
// Timeout is the wall-clock budget for graceful exit and, if required,
// the subsequent platform termination fallback.
Timeout time.Duration
}
@@ -54,13 +68,12 @@ type StopConfig struct {
// for the process to actually die. Returns ErrNotRunning if no bus is
// running for that work dir.
//
// Implementation note: on Unix we send SIGTERM. The bus daemon's Run loop
// watches its parent ctx for cancellation; the cobra `event _bus` command
// wires signal.NotifyContext so SIGTERM triggers ctx.Done() → graceful
// shutdown path. On Windows we use os.Process.Signal(os.Interrupt) which
// the Go runtime maps to TerminateProcess for processes outside our
// console group; for v1 that's acceptable (Windows graceful shutdown is
// future work — plan §16 v2).
// Stop first asks the bus to shut down through its owner-only IPC endpoint.
// This is the normal path on every platform and lets the daemon cancel its
// cloud source, notify consumers, and release its lock. If the endpoint is
// unavailable (for example an older bus) or graceful shutdown times out, the
// platform signal is used as a compatibility fallback: SIGTERM on Unix and
// TerminateProcess via os.Kill on Windows.
func Stop(cfg StopConfig) error {
if cfg.WorkDir == "" {
return errors.New("busctl: StopConfig.WorkDir is required")
@@ -81,6 +94,32 @@ func Stop(cfg StopConfig) error {
if err != nil {
return fmt.Errorf("busctl: find process %d: %w", pid, err)
}
if strings.TrimSpace(cfg.IPCEndpoint) != "" {
if err := stopRequest(cfg.IPCEndpoint); err == nil && stopWaitForBusExit(pid, cfg.Timeout) {
return nil
}
}
// The bus may finish shutting down immediately after the graceful wait
// reaches its deadline (or after a failed IPC request). Recheck before
// inspecting ownership so a completed stop is not reported as stale.
if !stopAlive(pid) {
return nil
}
owner, err := stopValidateHolderOwner(LockPath(cfg.WorkDir), pid)
if err != nil {
return fmt.Errorf("busctl: verify bus pid=%d ownership: %w", pid, err)
}
if !owner {
// ValidateHolderOwnership observes a released lock both for a stale,
// reused PID and for a bus that exits during the ownership probe. Only
// the still-alive case is unsafe to signal.
if !stopAlive(pid) {
return nil
}
return fmt.Errorf("%w: pid=%d does not own %s; stale PID was not signalled", ErrOwnerUnverified, pid, LockPath(cfg.WorkDir))
}
if err := stopSignalProcess(proc, stopSignal()); err != nil {
// On many Unix platforms Signal returns "process already finished"
// when the bus has just exited on its own — treat that as success.
@@ -90,13 +129,45 @@ func Stop(cfg StopConfig) error {
return fmt.Errorf("busctl: signal bus pid=%d: %w", pid, err)
}
// Poll for actual exit.
deadline := time.Now().Add(cfg.Timeout)
for time.Now().Before(deadline) {
if !stopAlive(pid) {
return nil
}
time.Sleep(50 * time.Millisecond)
if stopWaitForBusExit(pid, cfg.Timeout) {
return nil
}
return fmt.Errorf("busctl: bus pid=%d did not exit within %s", pid, cfg.Timeout)
}
func requestBusStop(endpoint string) error {
conn, err := stopDial(endpoint)
if err != nil {
return fmt.Errorf("busctl: dial bus for stop: %w", err)
}
defer conn.Close()
_ = conn.SetDeadline(time.Now().Add(DefaultStatusRPCTimeout))
w := transport.NewWriter(conn)
r := transport.NewReader(conn)
if err := w.WriteJSON(transport.Hello{
Type: transport.FrameTypeHello,
ConsumerPID: os.Getpid(),
Role: transport.HelloRoleStop,
}); err != nil {
return fmt.Errorf("busctl: write stop hello: %w", err)
}
var bye transport.Bye
if err := r.ReadJSON(&bye); err != nil {
return fmt.Errorf("busctl: read stop response: %w", err)
}
if bye.Type != transport.FrameTypeBye || bye.Reason != "stop_request" {
return fmt.Errorf("busctl: unexpected stop response type=%q reason=%q", bye.Type, bye.Reason)
}
return nil
}
func waitForBusExit(pid int, timeout time.Duration) bool {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
if !stopAlive(pid) {
return true
}
time.Sleep(50 * time.Millisecond)
}
return !stopAlive(pid)
}
+270 -1
View File
@@ -16,6 +16,7 @@ package busctl
import (
"context"
"errors"
"net"
"os"
"os/exec"
"path/filepath"
@@ -23,7 +24,10 @@ import (
"testing"
"time"
dwsevent "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/bus"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/transport"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/testseam"
)
func TestStop_NotRunningWhenLockMissing(t *testing.T) {
@@ -34,7 +38,7 @@ func TestStop_NotRunningWhenLockMissing(t *testing.T) {
}
}
func TestStop_NotRunningWhenPIDDead(t *testing.T) {
func TestCrossPlatformCoverageStopNotRunningWhenPIDDead(t *testing.T) {
dir := shortTempDir(t)
// Write a definitely-dead PID into bus.lock.
if err := os.WriteFile(LockPath(dir), []byte("2147483646\n"), 0o600); err != nil {
@@ -46,6 +50,269 @@ func TestStop_NotRunningWhenPIDDead(t *testing.T) {
}
}
func TestCrossPlatformCoverageStopPrefersGracefulIPC(t *testing.T) {
const pid = 4242
stopped := false
signals := 0
testseam.Swap(t, &stopReadHolderPID, func(string) int { return pid })
testseam.Swap(t, &stopAlive, func(int) bool { return !stopped })
proc, err := os.FindProcess(os.Getpid())
if err != nil {
t.Fatal(err)
}
testseam.Swap(t, &stopFindProcess, func(int) (*os.Process, error) { return proc, nil })
testseam.Swap(t, &stopRequest, func(endpoint string) error {
if endpoint != "test-endpoint" {
t.Fatalf("stop endpoint = %q", endpoint)
}
stopped = true
return nil
})
testseam.Swap(t, &stopSignalProcess, func(*os.Process, os.Signal) error {
signals++
return nil
})
if err := Stop(StopConfig{
WorkDir: "test-workdir",
IPCEndpoint: "test-endpoint",
Timeout: time.Second,
}); err != nil {
t.Fatalf("Stop() = %v", err)
}
if signals != 0 {
t.Fatalf("graceful IPC stop used %d process signals", signals)
}
}
func TestCrossPlatformCoverageStopDoesNotSignalUnverifiedReusedPID(t *testing.T) {
const pid = 4242
testseam.Swap(t, &stopReadHolderPID, func(string) int { return pid })
testseam.Swap(t, &stopAlive, func(int) bool { return true })
proc, err := os.FindProcess(os.Getpid())
if err != nil {
t.Fatal(err)
}
testseam.Swap(t, &stopFindProcess, func(int) (*os.Process, error) { return proc, nil })
testseam.Swap(t, &stopRequest, func(string) error { return errors.New("stale endpoint") })
testseam.Swap(t, &stopValidateHolderOwner, func(path string, gotPID int) (bool, error) {
if path != LockPath("test-workdir") || gotPID != pid {
t.Fatalf("ownership check path=%q pid=%d", path, gotPID)
}
return false, nil
})
signals := 0
testseam.Swap(t, &stopSignalProcess, func(*os.Process, os.Signal) error {
signals++
return nil
})
err = Stop(StopConfig{
WorkDir: "test-workdir",
IPCEndpoint: "stale-endpoint",
Timeout: time.Millisecond,
})
if !errors.Is(err, ErrOwnerUnverified) {
t.Fatalf("Stop() error = %v, want ErrOwnerUnverified", err)
}
if signals != 0 {
t.Fatalf("unverified reused PID received %d signals", signals)
}
}
func TestCrossPlatformCoverageStopAcceptsExitAtGracefulTimeoutBoundary(t *testing.T) {
const pid = 4242
t.Run("before ownership validation", func(t *testing.T) {
exited := false
validated := false
testseam.Swap(t, &stopReadHolderPID, func(string) int { return pid })
testseam.Swap(t, &stopAlive, func(int) bool { return !exited })
proc, err := os.FindProcess(os.Getpid())
if err != nil {
t.Fatal(err)
}
testseam.Swap(t, &stopFindProcess, func(int) (*os.Process, error) { return proc, nil })
testseam.Swap(t, &stopRequest, func(string) error { return nil })
testseam.Swap(t, &stopWaitForBusExit, func(int, time.Duration) bool {
exited = true
return false
})
testseam.Swap(t, &stopValidateHolderOwner, func(string, int) (bool, error) {
validated = true
return false, nil
})
if err := Stop(StopConfig{WorkDir: "test-workdir", IPCEndpoint: "test-endpoint"}); err != nil {
t.Fatalf("Stop() = %v, want success after bus exit", err)
}
if validated {
t.Fatal("ownership was validated after the bus had already exited")
}
})
t.Run("during ownership validation", func(t *testing.T) {
exited := false
signals := 0
testseam.Swap(t, &stopReadHolderPID, func(string) int { return pid })
testseam.Swap(t, &stopAlive, func(int) bool { return !exited })
proc, err := os.FindProcess(os.Getpid())
if err != nil {
t.Fatal(err)
}
testseam.Swap(t, &stopFindProcess, func(int) (*os.Process, error) { return proc, nil })
testseam.Swap(t, &stopRequest, func(string) error { return nil })
testseam.Swap(t, &stopWaitForBusExit, func(int, time.Duration) bool { return false })
testseam.Swap(t, &stopValidateHolderOwner, func(string, int) (bool, error) {
exited = true
return false, nil
})
testseam.Swap(t, &stopSignalProcess, func(*os.Process, os.Signal) error {
signals++
return nil
})
if err := Stop(StopConfig{WorkDir: "test-workdir", IPCEndpoint: "test-endpoint"}); err != nil {
t.Fatalf("Stop() = %v, want success after exit during ownership validation", err)
}
if signals != 0 {
t.Fatalf("exited bus received %d fallback signals", signals)
}
})
}
func TestCrossPlatformCoverageStopReportsOwnershipValidationError(t *testing.T) {
const pid = 4242
errInjected := errors.New("ownership validation failed")
testseam.Swap(t, &stopReadHolderPID, func(string) int { return pid })
testseam.Swap(t, &stopAlive, func(int) bool { return true })
proc, err := os.FindProcess(os.Getpid())
if err != nil {
t.Fatal(err)
}
testseam.Swap(t, &stopFindProcess, func(int) (*os.Process, error) { return proc, nil })
testseam.Swap(t, &stopValidateHolderOwner, func(string, int) (bool, error) {
return false, errInjected
})
signals := 0
testseam.Swap(t, &stopSignalProcess, func(*os.Process, os.Signal) error {
signals++
return nil
})
err = Stop(StopConfig{WorkDir: "test-workdir"})
if !errors.Is(err, errInjected) {
t.Fatalf("Stop() error = %v, want injected ownership error", err)
}
if signals != 0 {
t.Fatalf("ownership validation error sent %d signals", signals)
}
}
func TestCrossPlatformCoverageRequestBusStopProtocol(t *testing.T) {
dir := shortTempDir(t)
endpoint := dwsevent.IPCEndpoint(
dir,
"open",
dwsevent.SourceKindAppStream,
dwsevent.IdentityHash(dir),
)
listener, err := transport.Listen(endpoint)
if err != nil {
t.Fatalf("Listen() = %v", err)
}
defer listener.Close()
serverDone := make(chan error, 1)
go func() {
conn, err := listener.Accept()
if err != nil {
serverDone <- err
return
}
defer conn.Close()
r := transport.NewReader(conn)
w := transport.NewWriter(conn)
var hello transport.Hello
if err := r.ReadJSON(&hello); err != nil {
serverDone <- err
return
}
if hello.Type != transport.FrameTypeHello || hello.Role != transport.HelloRoleStop {
serverDone <- errors.New("unexpected stop hello")
return
}
serverDone <- w.WriteJSON(transport.Bye{
Type: transport.FrameTypeBye,
Reason: "stop_request",
})
}()
if err := requestBusStop(endpoint); err != nil {
t.Fatalf("requestBusStop() = %v", err)
}
if err := <-serverDone; err != nil {
t.Fatalf("stop protocol server = %v", err)
}
}
func TestCrossPlatformCoverageRequestBusStopErrors(t *testing.T) {
t.Run("dial", func(t *testing.T) {
testseam.Swap(t, &stopDial, func(string) (net.Conn, error) {
return nil, errors.New("dial failed")
})
if err := requestBusStop("test-endpoint"); err == nil {
t.Fatal("requestBusStop() unexpectedly succeeded")
}
})
t.Run("write", func(t *testing.T) {
testseam.Swap(t, &stopDial, func(string) (net.Conn, error) {
return &queryErrorConn{failAt: 1}, nil
})
if err := requestBusStop("test-endpoint"); err == nil {
t.Fatal("requestBusStop() unexpectedly succeeded")
}
})
t.Run("read", func(t *testing.T) {
testseam.Swap(t, &stopDial, func(string) (net.Conn, error) {
return &queryErrorConn{}, nil
})
if err := requestBusStop("test-endpoint"); err == nil {
t.Fatal("requestBusStop() unexpectedly succeeded")
}
})
t.Run("unexpected response", func(t *testing.T) {
client, server := net.Pipe()
t.Cleanup(func() {
_ = client.Close()
_ = server.Close()
})
testseam.Swap(t, &stopDial, func(string) (net.Conn, error) {
return client, nil
})
serverDone := make(chan error, 1)
go func() {
var hello transport.Hello
if err := transport.NewReader(server).ReadJSON(&hello); err != nil {
serverDone <- err
return
}
serverDone <- transport.NewWriter(server).WriteJSON(transport.Bye{
Type: transport.FrameTypeBye,
Reason: "unexpected",
})
}()
if err := requestBusStop("test-endpoint"); err == nil {
t.Fatal("requestBusStop() unexpectedly succeeded")
}
if err := <-serverDone; err != nil {
t.Fatalf("stop protocol server = %v", err)
}
})
}
func TestStop_SignalsLiveProcess(t *testing.T) {
skipOnWindows(t)
dir := shortTempDir(t)
@@ -63,6 +330,7 @@ func TestStop_SignalsLiveProcess(t *testing.T) {
}
}()
pid := cmd.Process.Pid
testseam.Swap(t, &stopValidateHolderOwner, func(string, int) (bool, error) { return true, nil })
// Reap the child in background so Wait doesn't leave a zombie.
waited := make(chan error, 1)
@@ -111,6 +379,7 @@ func TestStop_TimeoutWhenChildIgnoresSignal(t *testing.T) {
_, _ = cmd.Process.Wait()
}()
pid := cmd.Process.Pid
testseam.Swap(t, &stopValidateHolderOwner, func(string, int) (bool, error) { return true, nil })
if err := os.WriteFile(LockPath(dir), []byte(strconv.Itoa(pid)+"\n"), 0o600); err != nil {
t.Fatal(err)
+4 -5
View File
@@ -17,8 +17,7 @@ package busctl
import "os"
// stopSignal returns the graceful-shutdown signal for Windows. The Go
// runtime maps os.Interrupt to TerminateProcess for non-console-group
// processes — not truly graceful, but acceptable for v1 (Windows graceful
// shutdown via Ctrl+Break is in the v2 backlog, plan §16).
func stopSignal() os.Signal { return os.Interrupt }
// stopSignal is only the fallback after the graceful IPC stop path fails or
// times out. Go maps os.Kill to TerminateProcess on Windows; os.Interrupt is
// unsupported and returns syscall.EWINDOWS.
func stopSignal() os.Signal { return os.Kill }
@@ -0,0 +1,27 @@
// Copyright 2026 Alibaba Group
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//go:build windows
package busctl
import (
"os"
"testing"
)
func TestCrossPlatformCoverageWindowsStopFallbackUsesKill(t *testing.T) {
if stopSignal() != os.Kill {
t.Fatalf("stopSignal() = %v, want os.Kill", stopSignal())
}
}
+281 -6
View File
@@ -9,6 +9,7 @@ import (
"net/url"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"time"
@@ -492,10 +493,11 @@ func defaultHTTPGetFile(ctx context.Context, url string, headers map[string]stri
return nil
}
// runMediaInsert implements the three-step flow for inserting an attachment into a document:
// runMediaInsert implements the four-step flow for inserting an attachment into a document:
// 1. get_doc_attachment_upload_info → obtain uploadUrl + resourceId
// 2. HTTP PUT file content to OSS
// 3. insert_document_block with attachment element
// 4. list_document_blocks → prove the uploaded resource is visible in the document
func runMediaInsert(cmd *cobra.Command, _ []string) error {
nodeID, err := mustFlagOrFallback(cmd, "node", "url", "id", "node-id", "doc-id", "file-id")
if err != nil {
@@ -551,7 +553,7 @@ func runMediaInsert(cmd *cobra.Command, _ []string) error {
ctx := cmd.Context()
// Step 1: get upload credentials (uploadUrl + resourceId)
deps.Out.PrintInfo(fmt.Sprintf("[1/3] 获取附件上传凭证 (%s, %d bytes)...", fileName, fileSize))
deps.Out.PrintInfo(fmt.Sprintf("[1/4] 获取附件上传凭证 (%s, %d bytes)...", fileName, fileSize))
credText, err := callMCPToolReturnText(ctx, "get_doc_attachment_upload_info", map[string]any{
"nodeId": nodeID,
@@ -569,7 +571,7 @@ func runMediaInsert(cmd *cobra.Command, _ []string) error {
}
// Step 2: HTTP PUT file to OSS
deps.Out.PrintInfo("[2/3] 上传文件到 OSS...")
deps.Out.PrintInfo("[2/4] 上传文件到 OSS...")
ossHeaders := map[string]string{
"Content-Type": mimeType,
@@ -592,7 +594,7 @@ func runMediaInsert(cmd *cobra.Command, _ []string) error {
}
// Step 3: insert block into document
deps.Out.PrintInfo("[3/3] 插入块到文档...")
deps.Out.PrintInfo("[3/4] 插入块到文档...")
const maxInlineImageSize = 20 * 1024 * 1024 // 20MB
@@ -644,7 +646,8 @@ func runMediaInsert(cmd *cobra.Command, _ []string) error {
insertArgs["referenceBlockId"] = v
}
if err := callMCPTool("insert_document_block", insertArgs); err != nil {
insertText, err := callMCPToolReturnText(ctx, "insert_document_block", insertArgs)
if err != nil {
return apperrors.NewAPI(
"附件已上传,但正文 block 插入结果未知;请先检查媒体列表,不要重复上传或插入",
apperrors.WithOperation("doc.media_insert"),
@@ -665,6 +668,23 @@ func runMediaInsert(cmd *cobra.Command, _ []string) error {
apperrors.WithCause(err),
)
}
insertResult := map[string]any{}
if strings.TrimSpace(insertText) != "" {
if err := json.Unmarshal([]byte(insertText), &insertResult); err != nil {
return docMediaInsertVerificationError(nodeID, resourceID, resourceURL, fileName,
fmt.Errorf("解析 insert_document_block 响应失败: %w", err))
}
}
insertedBlockID := insertedDocBlockID(insertResult)
deps.Out.PrintInfo("[4/4] 回读验证媒体块...")
verifiedBlockID, verifyErr := verifyInsertedDocMedia(ctx, nodeID, insertedBlockID, resourceID, resourceURL)
if verifyErr != nil {
return docMediaInsertVerificationError(nodeID, resourceID, resourceURL, fileName, verifyErr)
}
if insertedBlockID == "" {
insertedBlockID = verifiedBlockID
}
return deps.Out.PrintJSON(map[string]any{
"contractVersion": "doc.operation.v1",
@@ -674,16 +694,271 @@ func runMediaInsert(cmd *cobra.Command, _ []string) error {
"operation": "doc.media_insert",
"data": map[string]any{
"nodeId": nodeID, "resourceId": resourceID, "resourceUrl": resourceURL,
"fileName": fileName, "mimeType": mimeType, "sizeBytes": fileSize, "inserted": true,
"blockId": insertedBlockID, "fileName": fileName, "mimeType": mimeType, "sizeBytes": fileSize,
"inserted": true, "verified": true,
},
"steps": []map[string]any{
{"name": "resolve_upload", "status": "success"},
{"name": "upload_oss", "status": "success"},
{"name": "insert_block", "status": "success"},
{"name": "verify", "status": "success"},
},
})
}
func docMediaInsertVerificationError(nodeID, resourceID, resourceURL, fileName string, cause error) error {
return apperrors.NewAPI(
"附件已上传且插块请求已执行,但回读未能证明媒体块落库;不要直接重试上传或插入",
apperrors.WithOperation("doc.media_insert"),
apperrors.WithReason("doc_media_insert_verification_failed"),
apperrors.WithFailureStage("verify"),
apperrors.WithExecutionStarted(true),
apperrors.WithRetryable(false),
apperrors.WithActions("运行 dws doc +media-list 检查 resourceId", "确认媒体不存在后再决定是否重新执行"),
apperrors.WithDetails(map[string]any{
"contractVersion": "doc.operation.v1", "status": "partial_success", "nodeId": nodeID,
"resourceId": resourceID, "resourceUrl": resourceURL, "fileName": fileName, "verified": false,
"steps": []map[string]any{
{"name": "resolve_upload", "status": "success"},
{"name": "upload_oss", "status": "success"},
{"name": "insert_block", "status": "success"},
{"name": "verify", "status": "failed"},
},
}),
apperrors.WithCause(cause),
)
}
var docMediaVerifyWait = waitForDocVerification
func verifyInsertedDocMedia(ctx context.Context, nodeID, blockID, resourceID, resourceURL string) (string, error) {
delays := []time.Duration{250 * time.Millisecond, 500 * time.Millisecond, time.Second, 2 * time.Second, 4 * time.Second, 8 * time.Second}
var lastErr error
for attempt := 0; attempt <= len(delays); attempt++ {
blocks, err := readAllDocBlocksForVerification(ctx, nodeID)
if err != nil {
lastErr = err
} else {
lastErr = nil
if found := findVerifiedMediaBlock(blocks, blockID, resourceID, resourceURL); found != "" {
return found, nil
}
}
if attempt < len(delays) {
if err := docMediaVerifyWait(ctx, delays[attempt]); err != nil {
return "", err
}
}
}
if lastErr != nil {
return "", fmt.Errorf("媒体资源在有界回读窗口内仍无法读取: %w", lastErr)
}
return "", fmt.Errorf("媒体资源在有界回读窗口内仍不可见")
}
func waitForDocVerification(ctx context.Context, delay time.Duration) error {
if ctx == nil {
ctx = context.Background()
}
timer := time.NewTimer(delay)
defer timer.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-timer.C:
return nil
}
}
func readAllDocBlocksForVerification(ctx context.Context, nodeID string) ([]any, error) {
const pageSize = 50
const maxItems = 5000
all := make([]any, 0, pageSize)
seenPageIdentities := map[string]bool{}
for start := 0; start < maxItems; start += pageSize {
text, err := callMCPToolReturnTextOnServer(ctx, "doc", "list_document_blocks", map[string]any{
"nodeId": nodeID, "format": "jsonml", "startIndex": start, "endIndex": start + pageSize - 1,
})
if err != nil {
return nil, err
}
var payload map[string]any
if err := json.Unmarshal([]byte(text), &payload); err != nil {
return nil, fmt.Errorf("解析 list_document_blocks 回读失败: %w", err)
}
payload = nestedDocMap(payload)
blocks, ok := payload["blocks"].([]any)
if !ok {
return nil, fmt.Errorf("list_document_blocks 回读缺少 blocks 数组")
}
pageIdentity := docBlockPageIdentity(blocks)
if pageIdentity != "" && seenPageIdentities[pageIdentity] {
return nil, fmt.Errorf("list_document_blocks 分页停滞")
}
if pageIdentity != "" {
seenPageIdentities[pageIdentity] = true
}
all = append(all, blocks...)
hasMore, hasMoreKnown := payload["hasMore"].(bool)
if hasMoreKnown && !hasMore {
return all, nil
}
if !hasMoreKnown {
if total, ok := docNumberAsInt(payload["totalCount"]); ok && len(all) >= total {
return all, nil
}
}
if !hasMoreKnown && len(blocks) < pageSize {
return all, nil
}
if len(blocks) == 0 {
return nil, fmt.Errorf("list_document_blocks 声明仍有下一页但当前页为空")
}
}
return nil, fmt.Errorf("文档块超过安全回读上限")
}
func docBlockPageIdentity(blocks []any) string {
if len(blocks) == 0 {
return ""
}
ids := make([]string, 0, len(blocks))
for _, value := range blocks {
id := ""
switch block := value.(type) {
case map[string]any:
id = directDocBlockIdentity(block)
if id == "" {
if element, ok := block["element"].(map[string]any); ok {
id = directDocBlockIdentity(element)
}
}
case []any:
if len(block) > 1 {
if attributes, ok := block[1].(map[string]any); ok {
id = directDocBlockIdentity(attributes)
}
}
}
if id == "" {
return ""
}
ids = append(ids, id)
}
encoded, _ := json.Marshal(ids)
return string(encoded)
}
func directDocBlockIdentity(block map[string]any) string {
for _, key := range []string{"blockId", "id", "uuid", "elementId"} {
if text, ok := block[key].(string); ok && strings.TrimSpace(text) != "" {
return strings.TrimSpace(text)
}
}
return ""
}
func nestedDocMap(data map[string]any) map[string]any {
for _, key := range []string{"result", "data"} {
if nested, ok := data[key].(map[string]any); ok {
return nestedDocMap(nested)
}
}
return data
}
func nestedDocString(value any, keys ...string) string {
switch typed := value.(type) {
case map[string]any:
for _, key := range keys {
if text, ok := typed[key].(string); ok && strings.TrimSpace(text) != "" {
return strings.TrimSpace(text)
}
}
orderedKeys := make([]string, 0, len(typed))
for key := range typed {
orderedKeys = append(orderedKeys, key)
}
sort.Strings(orderedKeys)
for _, key := range orderedKeys {
if text := nestedDocString(typed[key], keys...); text != "" {
return text
}
}
case []any:
for _, child := range typed {
if text := nestedDocString(child, keys...); text != "" {
return text
}
}
}
return ""
}
// insertedDocBlockID only accepts explicit block IDs from the insert result or
// known response wrappers. Arbitrary IDs may belong to the document, operator,
// or request and must not become a hard constraint for the media readback.
func insertedDocBlockID(data map[string]any) string {
for _, key := range []string{"blockId", "elementId"} {
if text, ok := data[key].(string); ok && strings.TrimSpace(text) != "" {
return strings.TrimSpace(text)
}
}
for _, wrapper := range []string{"result", "data", "content"} {
if inner, ok := data[wrapper].(map[string]any); ok {
if text := insertedDocBlockID(inner); text != "" {
return text
}
}
}
return ""
}
func findVerifiedMediaBlock(blocks []any, blockID, resourceID, resourceURL string) string {
for _, value := range blocks {
candidateID := nestedDocString(value, "blockId", "id", "uuid")
if candidateID == "" || (blockID != "" && candidateID != blockID) {
continue
}
mediaValue := docMediaReadbackValue(value)
if resourceID != "" && nestedDocString(mediaValue, "resourceId") == resourceID {
return candidateID
}
if resourceURL != "" && nestedDocString(mediaValue, "resourceUrl", "src") == resourceURL {
return candidateID
}
}
return ""
}
func docMediaReadbackValue(value any) any {
block, ok := value.(map[string]any)
if !ok {
return value
}
encoded, ok := block["jsonml"].(string)
if !ok || strings.TrimSpace(encoded) == "" {
return value
}
var decoded any
if json.Unmarshal([]byte(encoded), &decoded) != nil {
return value
}
return decoded
}
func docNumberAsInt(value any) (int, bool) {
switch typed := value.(type) {
case float64:
if typed >= 0 && typed == float64(int(typed)) {
return int(typed), true
}
case int:
return typed, typed >= 0
}
return 0, false
}
// parseAttachmentUploadInfo extracts uploadUrl, resourceId and resourceUrl from the MCP tool response.
func parseAttachmentUploadInfo(text string) (uploadURL, resourceID, resourceURL string, err error) {
var data map[string]any
+226 -1
View File
@@ -3,7 +3,9 @@ package helpers
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/http/httptest"
@@ -13,6 +15,7 @@ import (
"testing"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/testseam"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/pkg/edition"
"github.com/spf13/cobra"
)
@@ -74,6 +77,8 @@ func TestCrossPlatformCoverageDocUploadAndMediaErrorEdges(t *testing.T) {
t.Cleanup(func() { os.Args = oldArgs })
oldPut, oldGet := httpPutFile, httpGetFile
t.Cleanup(func() { httpPutFile, httpGetFile = oldPut, oldGet })
testseam.Swap(t, &helperSleep, func(time.Duration) {})
testseam.Swap(t, &docMediaVerifyWait, func(context.Context, time.Duration) error { return nil })
file := filepath.Join(t.TempDir(), "file.txt")
if err := os.WriteFile(file, []byte("content"), 0o600); err != nil {
t.Fatal(err)
@@ -193,6 +198,47 @@ func TestCrossPlatformCoverageDocUploadAndMediaErrorEdges(t *testing.T) {
t.Fatal("media insert failure returned nil")
}
})
t.Run("media insert response parse failure", func(t *testing.T) {
httpPutFile = func(context.Context, string, map[string]string, string, int64) error { return nil }
caller := &scriptedToolCaller{steps: []scriptedToolStep{
{text: `{"uploadUrl":"https://upload","resourceId":"resource"}`},
{text: `{`},
}}
if err := mediaCommand(t, caller, file, "text/plain"); err == nil {
t.Fatal("invalid media insert response returned nil")
}
})
t.Run("small image verifies src readback when upload also returns resource ID", func(t *testing.T) {
httpPutFile = func(context.Context, string, map[string]string, string, int64) error { return nil }
caller := &scriptedToolCaller{steps: []scriptedToolStep{
{text: `{"uploadUrl":"https://upload","resourceId":"resource","resourceUrl":"https://image"}`},
{text: `{"blockId":"media-block"}`},
{text: `{"blocks":[{"blockId":"media-block","jsonml":"[\"p\",{\"uuid\":\"media-block\"},[\"img\",{\"src\":\"https://image\"}]]"}],"hasMore":false}`},
}}
if err := mediaCommand(t, caller, file, "image/png"); err != nil {
t.Fatal(err)
}
if caller.calls != 3 {
t.Fatalf("media insert calls = %d, want credential, insert, and one readback", caller.calls)
}
if caller.args["format"] != "jsonml" {
t.Fatalf("media readback format = %#v, want jsonml", caller.args["format"])
}
})
t.Run("unrelated insert response IDs do not constrain media readback", func(t *testing.T) {
httpPutFile = func(context.Context, string, map[string]string, string, int64) error { return nil }
caller := &scriptedToolCaller{steps: []scriptedToolStep{
{text: `{"uploadUrl":"https://upload","resourceId":"resource","resourceUrl":"https://image"}`},
{text: `{"id":"document-id","operator":{"id":"operator-id"},"data":{"request":{"id":"request-id"}}}`},
{text: `{"blocks":[{"blockId":"media-block","jsonml":"[\"p\",{\"uuid\":\"media-block\"},[\"img\",{\"src\":\"https://image\"}]]"}],"hasMore":false}`},
}}
if err := mediaCommand(t, caller, file, "image/png"); err != nil {
t.Fatal(err)
}
if caller.calls != 3 {
t.Fatalf("media insert calls = %d, want credential, insert, and one readback", caller.calls)
}
})
t.Run("large image becomes attachment", func(t *testing.T) {
httpPutFile = func(context.Context, string, map[string]string, string, int64) error { return nil }
large := filepath.Join(t.TempDir(), "large.png")
@@ -202,7 +248,11 @@ func TestCrossPlatformCoverageDocUploadAndMediaErrorEdges(t *testing.T) {
if err := os.Truncate(large, 21*1024*1024); err != nil {
t.Fatal(err)
}
caller := &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"uploadUrl":"https://upload","resourceId":"resource","resourceUrl":"https://image"}`}, {text: `{}`}}}
caller := &scriptedToolCaller{steps: []scriptedToolStep{
{text: `{"uploadUrl":"https://upload","resourceId":"resource","resourceUrl":"https://image"}`},
{text: `{}`},
{text: `{"blocks":[{"blockId":"media-block","element":{"attachment":{"resourceId":"resource"}}}],"hasMore":false}`},
}}
if err := mediaCommand(t, caller, large, "image/png"); err != nil {
t.Fatal(err)
}
@@ -240,6 +290,181 @@ func TestCrossPlatformCoverageDefaultDocHTTPTransportEdges(t *testing.T) {
}
}
func TestCrossPlatformCoverageDocMediaReadbackDefensiveEdges(t *testing.T) {
testseam.Swap(t, &helperSleep, func(time.Duration) {})
testseam.Swap(t, &docMediaVerifyWait, func(context.Context, time.Duration) error { return nil })
ctx := context.Background()
for _, tc := range []struct {
name string
steps []scriptedToolStep
}{
{"call failure", []scriptedToolStep{{err: errors.New("read")}}},
{"invalid json", []scriptedToolStep{{text: `{`}}},
{"missing blocks", []scriptedToolStep{{text: `{}`}}},
{"stalled page", []scriptedToolStep{{text: `{"blocks":[{"id":"a"}],"hasMore":true}`}, {text: `{"blocks":[{"id":"a"}],"hasMore":true}`}}},
{"empty continued page", []scriptedToolStep{{text: `{"blocks":[],"hasMore":true}`}}},
} {
t.Run(tc.name, func(t *testing.T) {
installScriptedCaller(t, &scriptedToolCaller{steps: tc.steps})
if _, err := readAllDocBlocksForVerification(ctx, "node"); err == nil {
t.Fatal("defensive readback returned nil")
}
})
}
t.Run("total count and nested payload", func(t *testing.T) {
installScriptedCaller(t, &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"data":{"blocks":[{"id":"a"}],"totalCount":1}}`}}})
blocks, err := readAllDocBlocksForVerification(ctx, "node")
if err != nil || len(blocks) != 1 {
t.Fatalf("blocks=%#v err=%v", blocks, err)
}
})
t.Run("identical adjacent pages advance by requested indexes", func(t *testing.T) {
blocks := make([]any, 50)
for index := range blocks {
blocks[index] = map[string]any{"blockType": "paragraph"}
}
first, err := json.Marshal(map[string]any{"blocks": blocks, "hasMore": true, "totalCount": 100})
if err != nil {
t.Fatal(err)
}
second, err := json.Marshal(map[string]any{"blocks": blocks, "hasMore": false, "totalCount": 100})
if err != nil {
t.Fatal(err)
}
caller := &scriptedToolCaller{steps: []scriptedToolStep{{text: string(first)}, {text: string(second)}}}
installScriptedCaller(t, caller)
got, err := readAllDocBlocksForVerification(ctx, "node")
if err != nil || len(got) != 100 || caller.calls != 2 {
t.Fatalf("blocks=%d calls=%d err=%v, want 100 blocks from two pages", len(got), caller.calls, err)
}
})
t.Run("explicit has more overrides inconsistent total count", func(t *testing.T) {
caller := &scriptedToolCaller{steps: []scriptedToolStep{
{text: `{"blocks":[{"id":"first"}],"hasMore":true,"totalCount":1}`},
{text: `{"blocks":[{"id":"second"}],"hasMore":false,"totalCount":2}`},
}}
installScriptedCaller(t, caller)
got, err := readAllDocBlocksForVerification(ctx, "node")
if err != nil || len(got) != 2 || caller.calls != 2 {
t.Fatalf("blocks=%d calls=%d err=%v, want both explicitly advertised pages", len(got), caller.calls, err)
}
})
t.Run("page identity accepts nested and JSONML block IDs", func(t *testing.T) {
if got := docBlockPageIdentity([]any{map[string]any{"element": map[string]any{"blockId": "nested-block"}}}); got != `["nested-block"]` {
t.Fatalf("nested block page identity = %q", got)
}
if got := docBlockPageIdentity([]any{[]any{"p", map[string]any{"uuid": "jsonml-block"}}}); got != `["jsonml-block"]` {
t.Fatalf("JSONML block page identity = %q", got)
}
})
t.Run("unknown pagination short page", func(t *testing.T) {
installScriptedCaller(t, &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"blocks":[{"id":"a"}]}`}}})
if blocks, err := readAllDocBlocksForVerification(ctx, "node"); err != nil || len(blocks) != 1 {
t.Fatalf("blocks=%#v err=%v", blocks, err)
}
})
t.Run("bounded retry failure", func(t *testing.T) {
installScriptedCaller(t, &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"blocks":[],"hasMore":false}`}}})
if _, err := verifyInsertedDocMedia(ctx, "node", "", "missing", ""); err == nil {
t.Fatal("missing media unexpectedly verified")
}
})
t.Run("transient read failure retries", func(t *testing.T) {
caller := &scriptedToolCaller{steps: []scriptedToolStep{
{err: errors.New("temporary read failure")},
{text: `{"blocks":[{"blockId":"media-block","element":{"attachment":{"resourceId":"resource"}}}],"hasMore":false}`},
}}
installScriptedCaller(t, caller)
blockID, err := verifyInsertedDocMedia(ctx, "node", "media-block", "resource", "")
if err != nil || blockID != "media-block" || caller.calls != 2 {
t.Fatalf("blockID=%q calls=%d err=%v", blockID, caller.calls, err)
}
})
t.Run("block identity alone does not verify media", func(t *testing.T) {
installScriptedCaller(t, &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"blocks":[{"blockId":"media-block","element":{"attachment":{"resourceId":"other-resource"}}}],"hasMore":false}`}}})
if _, err := verifyInsertedDocMedia(ctx, "node", "media-block", "resource", ""); err == nil {
t.Fatal("matching block ID with a different resource unexpectedly verified")
}
})
t.Run("block read safety limit", func(t *testing.T) {
steps := make([]scriptedToolStep, 100)
for index := range steps {
steps[index] = scriptedToolStep{text: fmt.Sprintf(`{"blocks":[{"id":"block-%d"}],"hasMore":true}`, index)}
}
installScriptedCaller(t, &scriptedToolCaller{steps: steps})
if _, err := readAllDocBlocksForVerification(ctx, "node"); err == nil {
t.Fatal("oversized block read returned nil")
}
})
if got := nestedDocMap(map[string]any{"result": map[string]any{"data": map[string]any{"ok": true}}}); got["ok"] != true {
t.Fatalf("nested map=%#v", got)
}
if nestedDocString(map[string]any{"x": []any{map[string]any{"id": " nested "}}}, "id") != "nested" || nestedDocString(3, "id") != "" {
t.Fatal("nested string traversal failed")
}
if got := nestedDocString(map[string]any{"z": map[string]any{"id": "last"}, "a": map[string]any{"id": "first"}}, "id"); got != "first" {
t.Fatalf("nested string traversal = %q, want deterministic key order", got)
}
if got := insertedDocBlockID(map[string]any{"id": "document", "result": map[string]any{"data": map[string]any{"blockId": " block "}}}); got != "block" {
t.Fatalf("trusted inserted block ID = %q, want block", got)
}
if got := insertedDocBlockID(map[string]any{"id": "document", "operator": map[string]any{"id": "operator"}, "data": map[string]any{"request": map[string]any{"id": "request"}}}); got != "" {
t.Fatalf("untrusted inserted block ID = %q, want empty", got)
}
blocks := []any{map[string]any{"id": "block", "resourceId": "rid"}, map[string]any{"id": "url-block", "resourceUrl": "https://media"}, map[string]any{"id": "src-block", "jsonml": `["p",{},["img",{"src":"https://image"}]]`}}
if findVerifiedMediaBlock(blocks, "block", "rid", "") != "block" || findVerifiedMediaBlock(blocks, "", "rid", "") != "block" || findVerifiedMediaBlock(blocks, "", "", "https://media") != "url-block" || findVerifiedMediaBlock(blocks, "src-block", "upload-resource", "https://image") != "src-block" || findVerifiedMediaBlock(blocks, "other-block", "upload-resource", "https://image") != "" || findVerifiedMediaBlock(blocks, "block", "wrong", "") != "" || findVerifiedMediaBlock(blocks, "", "missing", "") != "" {
t.Fatal("media block matching failed")
}
if findVerifiedMediaBlock([]any{map[string]any{"id": "bad-jsonml", "jsonml": `[`}}, "bad-jsonml", "", "https://image") != "" {
t.Fatal("invalid JSONML unexpectedly verified")
}
if got := docMediaReadbackValue("plain"); got != "plain" {
t.Fatalf("non-object media readback = %#v, want unchanged value", got)
}
for _, tc := range []struct {
value any
want bool
}{
{float64(3), true}, {float64(-1), false}, {1.5, false}, {3, true}, {-1, false}, {"3", false},
} {
_, ok := docNumberAsInt(tc.value)
if ok != tc.want {
t.Fatalf("docNumberAsInt(%#v) ok=%v want=%v", tc.value, ok, tc.want)
}
}
}
func TestCrossPlatformCoverageDocMediaReadbackStopsOnCancellation(t *testing.T) {
if err := waitForDocVerification(nil, time.Nanosecond); err != nil {
t.Fatalf("completed verification wait = %v", err)
}
cancelled, cancel := context.WithCancel(context.Background())
cancel()
if err := waitForDocVerification(cancelled, time.Hour); !errors.Is(err, context.Canceled) {
t.Fatalf("cancelled verification wait = %v, want context.Canceled", err)
}
caller := &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"blocks":[],"hasMore":false}`}}}
installScriptedCaller(t, caller)
if _, err := verifyInsertedDocMedia(cancelled, "node", "", "missing", ""); !errors.Is(err, context.Canceled) {
t.Fatalf("cancelled media verification = %v, want context.Canceled", err)
}
if caller.calls != 1 {
t.Fatalf("cancelled media verification calls = %d, want 1", caller.calls)
}
}
func TestCrossPlatformCoverageDocCreateUpdateAndBlockCommandEdges(t *testing.T) {
oldDeps, oldArgs, oldPut, oldGet := deps, os.Args, httpPutFile, httpGetFile
t.Cleanup(func() {
@@ -9,6 +9,7 @@ import (
"path/filepath"
"strings"
"testing"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/testseam"
"github.com/spf13/cobra"
@@ -194,6 +195,7 @@ func TestCrossPlatformCoverageDocDeprecationWrappersCoverage(t *testing.T) {
}
func TestCrossPlatformCoverageRunDocUploadDownloadAndMediaCoverage(t *testing.T) {
testseam.Swap(t, &docMediaVerifyWait, func(context.Context, time.Duration) error { return nil })
oldArgs := os.Args
os.Args = []string{"dws", "doc"}
t.Cleanup(func() { os.Args = oldArgs })
+45 -1
View File
@@ -60,6 +60,11 @@ func translateBatchOp(op map[string]any) (map[string]any, error) {
if !ok {
return nil, fmt.Errorf("unsupported toolName %q: must be a CLI command name (e.g. \"range clear\", \"range update\", \"merge-cells\"). Run 'dws sheet batch-update --help' for the full list", toolName)
}
if toolName == "set-dropdown" {
if err := validateBatchSetDropdownInput(input); err != nil {
return nil, err
}
}
return map[string]any{
"toolName": mapping.mcpTool,
@@ -67,6 +72,33 @@ func translateBatchOp(op map[string]any) (map[string]any, error) {
}, nil
}
func validateBatchSetDropdownInput(input map[string]any) error {
if _, exists := input["source-colors"]; exists {
return fmt.Errorf("set-dropdown: 顶层 source-colors 不受支持;Inline 颜色请写入 options[].color,SourceRange 颜色写入暂不支持")
}
if _, exists := input["colors"]; exists {
return fmt.Errorf("set-dropdown: 顶层 colors 不受支持;Inline 颜色请写入 options[].color,SourceRange 颜色写入暂不支持")
}
options, hasOptions := input["options"]
hasOptions = hasOptions && options != nil
sourceRange := batchStr(input, "source-range")
sourceSheetID := batchStr(input, "source-sheet-id")
hasSourceRange := sourceRange != ""
if hasOptions == hasSourceRange {
return fmt.Errorf("set-dropdown: options 与 source-range 必须且只能指定一个")
}
if hasSourceRange != (sourceSheetID != "") {
return fmt.Errorf("set-dropdown: source-range 与 source-sheet-id 必须同时指定")
}
if hasSourceRange {
if err := validateDropdownSourceRangeInput(sourceSheetID, sourceRange); err != nil {
return fmt.Errorf("set-dropdown: %w", err)
}
}
return nil
}
// ── BuildXxxArgs: CLI flag → MCP param 转换函数 ──────────────────────────────────
// 每个函数接收 CLI flag 名(kebab-case)的 map,输出 MCP 参数名(camelCase)的 map。
// 目前集中放在此文件;后续拆分命令文件时可移到各命令所在文件。
@@ -215,7 +247,15 @@ func BuildSetDropdownArgs(input map[string]any) map[string]any {
args := map[string]any{
"sheetId": batchStr(input, "sheet-id"),
"range": batchStr(input, "range"),
"options": input["options"],
}
if options, ok := input["options"]; ok && options != nil {
args["options"] = options
}
if sourceRange := batchStr(input, "source-range"); sourceRange != "" {
args["sourceRange"] = map[string]any{
"sheetId": batchStr(input, "source-sheet-id"),
"a1Notation": sourceRange,
}
}
if v, ok := input["multi-select"]; ok {
args["enableMultiSelect"] = v
@@ -417,6 +457,8 @@ func newBatchUpdateCmd() *cobra.Command {
toolName 使用 CLI 命令名(与原子命令一致),input 的键用 CLI flag 名去掉 --。
CLI 层自动翻译为 MCP toolName + 参数名,无需记忆 MCP 参数名。
source-range 的语义按 toolName 隔离:set-dropdown 中表示下拉候选项来源,
range fill/copy-to/move-to 中表示待填充、复制或移动的数据源区域。
支持的 CLI 命令名:
range clear / range update / merge-cells / unmerge-cells / update-dimension
@@ -425,6 +467,8 @@ CLI 层自动翻译为 MCP toolName + 参数名,无需记忆 MCP 参数名。
set-dropdown / delete-dropdown / csv-put / delete-float-image
其中 csv-put 与独立命令语义一致:CSV 字段值以 = 开头时按公式解析,
前加单引号(例如 "'=1+1")时写入以 = 开头的字面文本。
set-dropdown 不接受顶层 colors/source-colors;Inline 颜色应写在 options[].color,
SourceRange 颜色写入暂不支持。
注意:batch-update 中 group-dimension 适合默认展开分组;需要 --group-state fold 时请使用独立
dws sheet group-dimension 命令。
+9 -1
View File
@@ -23,7 +23,15 @@ func TestCrossPlatformCoverageSheetBatchOperationTranslationCoversEveryMapping(t
"group-state": "fold",
}
for name, mapping := range batchOpDispatch {
got, err := translateBatchOp(map[string]any{"toolName": name, "input": input})
opInput := input
if name == "set-dropdown" {
opInput = make(map[string]any, len(input))
for key, value := range input {
opInput[key] = value
}
delete(opInput, "source-range")
}
got, err := translateBatchOp(map[string]any{"toolName": name, "input": opInput})
if err != nil {
t.Errorf("translateBatchOp(%q): %v", name, err)
continue
+10 -2
View File
@@ -388,7 +388,10 @@ range update 与合并区域冲突时返回 MERGED_CELLS_CONFLICT 的行为。
csv 带 [row=N] 行号前缀的 CSV 文本
colIndices 列字母映射数组(定位列用 colIndices[j],禁止手数逗号)
rowIndices 行号映射数组
hasMore 是否因 maxChars 截断
hasMore 目标范围是否还有未返回数据
truncationReasons 部分返回原因:max_cells / max_chars
resolvedRange 未传 --range 时底层解析出的完整目标范围
returnedRange 本次实际完整返回的范围
取值模式(--value-render-option):
formatted_value 格式化后的展示值(默认),如 ¥1,000.00、2025-06-01
@@ -398,9 +401,14 @@ range update 与合并区域冲突时返回 MERGED_CELLS_CONFLICT 的行为。
与 range read 的区别:
- CSV 格式 token 消耗约为 JSON 的 1/3
- 支持选择取值模式
- 自动防爆(max-chars 截断 + has_more 标志)
- 自动防爆(30,000 单元格 / max-chars 上限 + hasMore 标志)
- [row=N] 前缀防止行号计算错误
当 hasMore=true 时本命令不会自动续读;请结合目标范围和
returnedRange 从下一行显式传 --range。max_cells 不能通过增大
--max-chars 解决。收到 forbidden.document.sizeOverLimit 表示工作簿整体
无法装载,应创建更小副本或拆分工作簿,缩小 --range 不能解决。
注意:csv-get 不返回合并单元格结构。查看合并范围请使用
dws sheet info --node NODE_ID --sheet-id SHEET_ID --format json,并读取 mergedRanges。`,
Example: ` dws sheet csv-get --node NODE_ID
+93 -22
View File
@@ -6,6 +6,7 @@ import (
"strconv"
"strings"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/cli"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contract"
"github.com/spf13/cobra"
)
@@ -18,6 +19,23 @@ func sizeTypeEnumHint(dimension string) string {
return "pixel / standard"
}
func validateDropdownSourceRangeInput(sourceSheetID, sourceRange string) error {
if strings.TrimSpace(sourceSheetID) == "" {
return fmt.Errorf("--source-sheet-id 不能为空")
}
sourceRange = strings.TrimSpace(sourceRange)
if sourceRange == "" {
return fmt.Errorf("--source-range 不能为空")
}
if strings.Contains(sourceRange, "!") {
return fmt.Errorf("--source-range 不能包含工作表前缀;请通过 --source-sheet-id 指定来源工作表")
}
if strings.HasPrefix(sourceRange, "=") || strings.Contains(sourceRange, ",") {
return fmt.Errorf("--source-range 必须是单一连续区域,不能是公式、表达式或多区域引用")
}
return nil
}
// newDimensionCmds creates dimension-related commands: insert/delete/update/move/add-dimension,
// merge-cells, unmerge-cells, and dropdown commands (set/get/delete-dropdown).
func newDimensionCmds() []*cobra.Command {
@@ -775,15 +793,24 @@ sheetId 支持传入工作表 ID 或工作表名称,可通过 sheet list 获
Short: "设置下拉列表",
Long: `在钉钉表格的指定单元格范围内设置下拉列表。
设置后,用户可以在这些单元格中从预定义的选项列表中选择值。
支持自定义每个选项的颜色和是否允许多选。
Inline 模式通过 --options 提供静态选项,支持选项颜色。
SourceRange 模式通过 --source-sheet-id + --source-range 引用同一工作簿内的来源区域,
支持跨工作表、整行和整列引用;读取配置时不会展开来源区域的当前值或颜色。
两种模式都支持 --multi-select,并且 --options 与 --source-range 必须且只能指定一个。
如果目标范围已存在下拉列表,会被新的配置覆盖。
nodeId 支持传入文档链接 URL 或文档 ID(dentryUuid),系统自动识别。
sheetId 支持传入工作表 ID 或工作表名称,可通过 sheet list 获取。
--options 为 JSON 数组,每个元素包含 value(必填)和 color(可选)。
选项值不能包含英文逗号。`,
选项值不能包含英文逗号。
--source-range 只接受不含工作表前缀的 A1 区域,例如 T1:T3、T:T 或 1:3;
来源工作表由 --source-sheet-id 单独指定。SourceRange 颜色写入暂不支持。
已验证行为:工作表重命名、在引用前插入行/列、删除引用前的行会自动调整引用并保持 valid;
已验证的 move-dimension 场景会使引用变为 invalid。列删除、删除整个来源区域或来源工作表等未覆盖场景不能预设结果;
结构操作后应先回读 sourceRangeStatus,只有 invalid 时才重新选择来源并写入。`,
Example: ` # 设置单选下拉列表
dws sheet set-dropdown --node NODE_ID --sheet-id SHEET_ID --range "A2:A100" \
--options '[{"value":"选项1"},{"value":"选项2"},{"value":"选项3"}]'
@@ -793,32 +820,63 @@ sheetId 支持传入工作表 ID 或工作表名称,可通过 sheet list 获
--options '[{"value":"高","color":"#ff0000"},{"value":"中","color":"#ffaa00"},{"value":"低","color":"#00ff00"}]' \
--multi-select
# 引用同一工作簿内另一工作表的区域作为候选项来源
dws sheet set-dropdown --node NODE_ID --sheet-id TARGET_SHEET_ID --range "C2:C100" \
--source-sheet-id SOURCE_SHEET_ID --source-range "T1:T3"
# 使用文档链接 URL
dws sheet set-dropdown --node "https://alidocs.dingtalk.com/i/nodes/<DOC_UUID>" \
--sheet-id SHEET_ID --range "C1:C10" --options '[{"value":"是"},{"value":"否"}]'`,
RunE: func(cmd *cobra.Command, args []string) error {
if err := validateRequiredFlags(cmd, "node", "sheet-id", "range"); err != nil {
return err
}
optionsStr := mustGetFlag(cmd, "options")
var options []map[string]any
if err := json.Unmarshal([]byte(optionsStr), &options); err != nil {
return fmt.Errorf("--options JSON 解析失败: %w", err)
sourceSheetID := mustGetFlag(cmd, "source-sheet-id")
sourceRange := mustGetFlag(cmd, "source-range")
hasOptions := optionsStr != ""
hasSourceRange := sourceRange != ""
if hasOptions == hasSourceRange {
return fmt.Errorf("--options 与 --source-range 必须且只能指定一个")
}
if len(options) == 0 {
return fmt.Errorf("--options 至少包含 1 个选项")
if hasSourceRange && sourceSheetID == "" {
return fmt.Errorf("使用 --source-range 时必须同时指定 --source-sheet-id")
}
for i, opt := range options {
val, ok := opt["value"].(string)
if !ok || val == "" {
return fmt.Errorf("--options[%d] 缺少必填的 value 字段或 value 为空", i)
}
if strings.Contains(val, ",") {
return fmt.Errorf("--options[%d].value 不能包含英文逗号: %q", i, val)
}
if !hasSourceRange && sourceSheetID != "" {
return fmt.Errorf("使用 --source-sheet-id 时必须同时指定 --source-range")
}
toolArgs := map[string]any{
"nodeId": mustGetFlag(cmd, "node"),
"sheetId": mustGetFlag(cmd, "sheet-id"),
"range": mustGetFlag(cmd, "range"),
"options": options,
}
if hasOptions {
var options []map[string]any
if err := json.Unmarshal([]byte(optionsStr), &options); err != nil {
return fmt.Errorf("--options JSON 解析失败: %w", err)
}
if len(options) == 0 {
return fmt.Errorf("--options 至少包含 1 个选项")
}
for i, opt := range options {
val, ok := opt["value"].(string)
if !ok || val == "" {
return fmt.Errorf("--options[%d] 缺少必填的 value 字段或 value 为空", i)
}
if strings.Contains(val, ",") {
return fmt.Errorf("--options[%d].value 不能包含英文逗号: %q", i, val)
}
}
toolArgs["options"] = options
} else {
if err := validateDropdownSourceRangeInput(sourceSheetID, sourceRange); err != nil {
return err
}
toolArgs["sourceRange"] = map[string]any{
"sheetId": sourceSheetID,
"a1Notation": sourceRange,
}
}
if multiSelect, _ := cmd.Flags().GetBool("multi-select"); multiSelect {
toolArgs["enableMultiSelect"] = true
@@ -839,36 +897,49 @@ sheetId 支持传入工作表 ID 或工作表名称,可通过 sheet list 获
CLIPath: "sheet set-dropdown",
PrimaryCLIPath: "sheet set-dropdown",
},
Description: "为指定范围设置下拉列表(可多选、可带颜色)。",
Description: "为指定范围设置 Inline 或 SourceRange 下拉列表(可多选;颜色仅 Inline 支持)。",
Interface: &contract.InterfaceSpec{
Mode: "mcp",
Availability: "available",
Ref: &contract.InterfaceRefSpec{ProductID: "sheet", RPCName: "set_dropdown_lists"},
},
Selection: contract.SelectionSpec{
AgentSummary: "为指定范围设置下拉列表(可多选、可带颜色)。",
AgentSummary: "为指定范围设置 Inline 或 SourceRange 下拉列表;SourceRange 可跨同一工作簿内的工作表。",
UseWhen: []string{"需要给单元格配置可选值下拉约束时"},
AvoidWhen: []string{"查看已有下拉用 get-dropdown;移除下拉用 delete-dropdown"},
Examples: []string{"dws sheet set-dropdown --node <NODE_ID> --sheet-id <SHEET_ID> --range \"A2:A100\" --options '[{\"value\":\"选项1\"},{\"value\":\"选项2\"}]'"},
Examples: []string{
"dws sheet set-dropdown --node <NODE_ID> --sheet-id <SHEET_ID> --range \"A2:A100\" --options '[{\"value\":\"选项1\"},{\"value\":\"选项2\"}]'",
"dws sheet set-dropdown --node <NODE_ID> --sheet-id <TARGET_SHEET_ID> --range \"B2:B100\" --source-sheet-id <SOURCE_SHEET_ID> --source-range \"T1:T3\"",
},
},
Parameters: []contract.ParamDecl{
{Name: "multi-select", Property: "enableMultiSelect"},
{Name: "node", Property: "nodeId"},
{Name: "source-sheet-id", Property: "sourceRange.sheetId", Required: boolPtr(false), RequiredWhen: "--source-range is provided"},
{Name: "source-range", Property: "sourceRange.a1Notation", Required: boolPtr(false), RequiredWhen: "exactly one of --options or --source-range must be provided"},
},
},
})
setDropdownCmd.Flags().String("node", "", "表格文档 ID 或 URL (必填)")
setDropdownCmd.Flags().String("sheet-id", "", "工作表 ID 或名称 (必填)")
setDropdownCmd.Flags().String("range", "", "目标单元格范围,A1 表示法,如 A2:A100 (必填)")
setDropdownCmd.Flags().String("options", "", `下拉选项 JSON 数组 (必填),如 '[{"value":"选项1","color":"#ff0000"}]'`)
setDropdownCmd.Flags().String("options", "", `Inline 下拉选项 JSON 数组,与 --source-range 二选一,如 '[{"value":"选项1","color":"#ff0000"}]'`)
setDropdownCmd.Flags().String("source-sheet-id", "", "SourceRange 来源工作表 ID;使用 --source-range 时必填")
setDropdownCmd.Flags().String("source-range", "", "SourceRange 来源区域,与 --options 二选一;不含工作表前缀,如 T1:T3、T:T 或 1:3")
setDropdownCmd.Flags().Bool("multi-select", false, "是否允许多选(默认单选)")
cli.AnnotateRuntimeFlagFormat(setDropdownCmd, "source-range", "a1-range")
setDropdownCmd.MarkFlagsOneRequired("options", "source-range")
setDropdownCmd.MarkFlagsMutuallyExclusive("options", "source-range")
setDropdownCmd.MarkFlagsRequiredTogether("source-sheet-id", "source-range")
getDropdownCmd := &cobra.Command{
Use: "get-dropdown",
Short: "获取下拉列表配置",
Long: `查询钉钉表格指定范围内的下拉列表配置。
返回范围内所有单元格的下拉列表配置信息,包括选项值和颜色。
Inline 配置返回 conditionValues 和 options;SourceRange 配置返回
sourceType/sourceRangeStatus/enableMultiSelect,且仅在 sourceRangeStatus="valid" 时返回
sourceRange。invalid 结果仍保留配置组,但省略 sourceRange;两种状态都不展开来源区域的候选值或颜色。
如果范围内存在多个不同的下拉列表配置,会分别返回每组配置及其覆盖的单元格列表。
如果范围内没有设置下拉列表,返回空。
@@ -0,0 +1,248 @@
package helpers
import (
"strings"
"testing"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/runtimeannotate"
)
func TestCrossPlatformCoverageSetDropdownSourceRangeCommand(t *testing.T) {
caller := &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"success":true}`}}}
installScriptedCaller(t, caller)
err := executeDimensionCoverage(t, "set-dropdown",
"--node", "node-1",
"--sheet-id", "target-sheet",
"--range", "B2:B100",
"--source-sheet-id", "source-sheet",
"--source-range", "$t$1:$t$3",
"--multi-select",
)
if err != nil {
t.Fatalf("set SourceRange dropdown: %v", err)
}
if caller.tool != "set_dropdown_lists" {
t.Fatalf("tool = %q, want set_dropdown_lists", caller.tool)
}
if _, exists := caller.args["options"]; exists {
t.Fatalf("SourceRange request must omit options: %#v", caller.args)
}
source, ok := caller.args["sourceRange"].(map[string]any)
if !ok {
t.Fatalf("sourceRange = %#v", caller.args["sourceRange"])
}
if source["sheetId"] != "source-sheet" || source["a1Notation"] != "$t$1:$t$3" {
t.Fatalf("sourceRange = %#v", source)
}
if caller.args["enableMultiSelect"] != true {
t.Fatalf("enableMultiSelect = %#v", caller.args["enableMultiSelect"])
}
}
func TestCrossPlatformCoverageSetDropdownInlineCommandStillOmitsSourceRange(t *testing.T) {
caller := &scriptedToolCaller{steps: []scriptedToolStep{{text: `{"success":true}`}}}
installScriptedCaller(t, caller)
err := executeDimensionCoverage(t, "set-dropdown",
"--node", "node-1",
"--sheet-id", "sheet-1",
"--range", "A1:A3",
"--options", `[{"value":"one","color":"#ff0000"}]`,
)
if err != nil {
t.Fatalf("set inline dropdown: %v", err)
}
if _, exists := caller.args["sourceRange"]; exists {
t.Fatalf("inline request must omit sourceRange: %#v", caller.args)
}
options, ok := caller.args["options"].([]map[string]any)
if !ok || len(options) != 1 || options[0]["value"] != "one" {
t.Fatalf("options = %#v", caller.args["options"])
}
}
func TestCrossPlatformCoverageSetDropdownModeConstraints(t *testing.T) {
base := []string{"--node", "node", "--sheet-id", "sheet", "--range", "A1"}
tests := []struct {
name string
args []string
want string
}{
{name: "neither", want: "one of the flags"},
{name: "both", args: []string{"--options", `[{"value":"one"}]`, "--source-sheet-id", "source", "--source-range", "A1:A3"}, want: "none of the others"},
{name: "source sheet without range", args: []string{"--options", `[{"value":"one"}]`, "--source-sheet-id", "source"}, want: "must all be set"},
{name: "range without source sheet", args: []string{"--source-range", "A1:A3"}, want: "must all be set"},
{name: "sheet prefix", args: []string{"--source-sheet-id", "source", "--source-range", "Sheet2!A1:A3"}, want: "不能包含工作表前缀"},
{name: "formula", args: []string{"--source-sheet-id", "source", "--source-range", "=A1:A3"}, want: "必须是单一连续区域"},
{name: "multi region", args: []string{"--source-sheet-id", "source", "--source-range", "A1:A3,C1:C3"}, want: "必须是单一连续区域"},
{name: "blank source range", args: []string{"--source-sheet-id", "source", "--source-range", " \t "}, want: "--source-range 不能为空"},
{name: "blank source sheet", args: []string{"--source-sheet-id", " \t ", "--source-range", "A1:A3"}, want: "--source-sheet-id 不能为空"},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
caller := &scriptedToolCaller{}
installScriptedCaller(t, caller)
args := append(append([]string{}, base...), tc.args...)
err := executeDimensionCoverage(t, "set-dropdown", args...)
if err == nil || !strings.Contains(err.Error(), tc.want) {
t.Fatalf("error = %v, want contains %q", err, tc.want)
}
if caller.calls != 0 {
t.Fatalf("calls = %d, validation must fail before request", caller.calls)
}
})
}
}
func TestCrossPlatformCoverageSetDropdownRunEValidation(t *testing.T) {
tests := []struct {
name string
options string
source string
rangeID string
want string
}{
{name: "neither mode", want: "必须且只能指定一个"},
{name: "source range without sheet", rangeID: "A1:A3", want: "必须同时指定 --source-sheet-id"},
{name: "source sheet without range", options: `[{"value":"one"}]`, source: "source", want: "必须同时指定 --source-range"},
{name: "invalid options JSON", options: "{", want: "JSON 解析失败"},
{name: "missing option value", options: `[{}]`, want: "value 为空"},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
cmd := dimensionCoverageCommand(t, "set-dropdown")
for name, value := range map[string]string{
"node": "node", "sheet-id": "sheet", "range": "A1",
"options": tc.options, "source-sheet-id": tc.source, "source-range": tc.rangeID,
} {
if err := cmd.Flags().Set(name, value); err != nil {
t.Fatalf("set --%s: %v", name, err)
}
}
err := cmd.RunE(cmd, nil)
if err == nil || !strings.Contains(err.Error(), tc.want) {
t.Fatalf("error = %v, want contains %q", err, tc.want)
}
})
}
}
func TestCrossPlatformCoverageSetDropdownPreservesSchemaConstraints(t *testing.T) {
cmd := dimensionCoverageCommand(t, "set-dropdown")
if constraints := runtimeannotate.CommandConstraints(cmd); !runtimeannotate.ConstraintsEmpty(constraints) {
t.Fatalf("constraints = %#v, want no new schema-level constraints", constraints)
}
}
func TestCrossPlatformCoverageBatchSetDropdownSourceRange(t *testing.T) {
got, err := translateBatchOp(map[string]any{
"toolName": "set-dropdown",
"input": map[string]any{
"sheet-id": "target-sheet",
"range": "B2:B100",
"source-sheet-id": "source-sheet",
"source-range": "T:T",
"multi-select": true,
},
})
if err != nil {
t.Fatalf("translate SourceRange batch op: %v", err)
}
input := got["input"].(map[string]any)
if _, exists := input["options"]; exists {
t.Fatalf("SourceRange batch input must omit options: %#v", input)
}
source := input["sourceRange"].(map[string]any)
if source["sheetId"] != "source-sheet" || source["a1Notation"] != "T:T" {
t.Fatalf("sourceRange = %#v", source)
}
invalid := []map[string]any{
{},
{"options": []any{map[string]any{"value": "one"}}, "source-sheet-id": "source", "source-range": "A1:A3"},
{"source-range": "A1:A3"},
{"source-sheet-id": "source", "source-range": "A1:A3", "colors": []any{"#fff"}},
{"source-sheet-id": "source", "source-range": "A1:A3", "source-colors": []any{"#fff"}},
{"source-sheet-id": "source", "source-range": "Sheet2!A1:A3"},
}
for _, value := range invalid {
if _, err := translateBatchOp(map[string]any{"toolName": "set-dropdown", "input": value}); err == nil {
t.Errorf("invalid batch SourceRange input %#v returned nil", value)
}
}
}
func TestCrossPlatformCoverageBatchSetDropdownValidationGuidance(t *testing.T) {
tests := []struct {
name string
input map[string]any
want string
}{
{
name: "inline top-level colors",
input: map[string]any{"options": []any{map[string]any{"value": "one"}}, "colors": []any{"#fff"}},
want: "Inline 颜色请写入 options[].color",
},
{
name: "source top-level colors",
input: map[string]any{"source-sheet-id": "source", "source-range": "A1:A3", "source-colors": []any{"#fff"}},
want: "SourceRange 颜色写入暂不支持",
},
{
name: "blank source range",
input: map[string]any{"source-sheet-id": "source", "source-range": " \t "},
want: "--source-range 不能为空",
},
{
name: "blank source sheet",
input: map[string]any{"source-sheet-id": " \t ", "source-range": "A1:A3"},
want: "--source-sheet-id 不能为空",
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
_, err := translateBatchOp(map[string]any{"toolName": "set-dropdown", "input": tc.input})
if err == nil || !strings.Contains(err.Error(), tc.want) {
t.Fatalf("error = %v, want contains %q", err, tc.want)
}
})
}
}
func TestCrossPlatformCoverageDropdownHelpDescribesDynamicAndInvalidContracts(t *testing.T) {
setDropdown := dimensionCoverageCommand(t, "set-dropdown")
if !strings.Contains(setDropdown.Long, "会自动调整引用并保持 valid") ||
!strings.Contains(setDropdown.Long, "只有 invalid 时才重新选择来源并写入") {
t.Fatalf("set-dropdown help does not describe verified structural behavior:\n%s", setDropdown.Long)
}
getDropdown := dimensionCoverageCommand(t, "get-dropdown")
if !strings.Contains(getDropdown.Long, `仅在 sourceRangeStatus="valid" 时返回`) ||
!strings.Contains(getDropdown.Long, "invalid 结果仍保留配置组,但省略 sourceRange") {
t.Fatalf("get-dropdown help does not describe conditional sourceRange readback:\n%s", getDropdown.Long)
}
}
func TestCrossPlatformCoverageValidateSourceRangeDataValidation(t *testing.T) {
valid := []any{
map[string]any{"type": "dropdown", "sourceRange": map[string]any{"sheetId": "source", "a1Notation": "T1:T3"}},
map[string]any{"type": "dropdown", "sourceRange": map[string]any{"sheetId": "missing-sheet", "a1Notation": "not-validated-locally"}, "enableMultiSelect": true},
}
for _, value := range valid {
if err := validateDataValidation(value, "dv"); err != nil {
t.Errorf("valid SourceRange data validation %#v: %v", value, err)
}
}
invalid := []any{
map[string]any{"type": "dropdown", "options": []any{map[string]any{"value": "one"}}, "sourceRange": map[string]any{"sheetId": "source", "a1Notation": "T1:T3"}},
map[string]any{"type": "dropdown", "sourceRange": "T1:T3"},
map[string]any{"type": "dropdown", "sourceRange": map[string]any{"a1Notation": "T1:T3"}},
map[string]any{"type": "dropdown", "sourceRange": map[string]any{"sheetId": "source"}},
map[string]any{"type": "dropdown", "sourceRange": map[string]any{"sheetId": "source", "a1Notation": "T1:T3", "colors": []any{"#fff"}}},
map[string]any{"type": "dropdown", "sourceRange": map[string]any{"sheetId": "source", "a1Notation": "T1:T3"}, "enableMultiSelect": "yes"},
}
for _, value := range invalid {
if err := validateDataValidation(value, "dv"); err == nil {
t.Errorf("invalid SourceRange data validation %#v returned nil", value)
}
}
}
+18 -1
View File
@@ -26,6 +26,16 @@ func newRangeCmd() *cobra.Command {
value 单元格值(内容由 --value-render-option 决定)
dataValidation 数据验证配置(下拉列表/复选框),无则为 null
hyperlink 单元格级超链接(path/sheet/range),无则省略
顶层还会返回 rowIndices / colIndices 与完成度字段:
hasMore 目标范围是否还有未返回数据
truncationReasons 部分返回原因(max_cells)
resolvedRange 未传 --range 时底层解析出的完整目标范围
returnedRange 本次实际完整返回的范围
单次最多返回 30,000 个单元格。hasMore=true 是部分成功,本命令
不会自动续读;请结合目标范围和 returnedRange 从下一行显式传
--range。收到 forbidden.document.sizeOverLimit 则表示工作簿整体无法装载,
应创建更小副本或拆分工作簿,缩小 --range 不能解决。
注意:range read/get 不返回合并单元格结构。查看合并范围请使用
dws sheet info --node NODE_ID --sheet-id SHEET_ID --format json,并读取 mergedRanges。
@@ -157,7 +167,12 @@ dws sheet info --node NODE_ID --sheet-id SHEET_ID --format json,并读取 merg
2) {"dataValidation":{"type":"none"}} → 显式清除该单元格 DV
3) {"dataValidation":{"type":"dropdown",...}} → 写新 dropdown(覆盖)
{"dataValidation":{"type":"checkbox",...}} → 写新 checkbox(覆盖)
dropdown: {"dataValidation":{"type":"dropdown","options":[{"value":"选项1"}],"enableMultiSelect":false}}
Inline dropdown:
{"dataValidation":{"type":"dropdown","options":[{"value":"选项1"}],"enableMultiSelect":false}}
SourceRange dropdown(同一工作簿内可跨工作表;不展开来源 values,不支持 colors):
{"dataValidation":{"type":"dropdown","sourceRange":{"sheetId":"SOURCE_SHEET_ID","a1Notation":"T1:T3"},"enableMultiSelect":false}}
dropdown 的 options 与 sourceRange 必须且只能传一个。SourceRange 支持普通区域、整行和整列;
a1Notation 不带工作表前缀,来源工作表单独写在 sheetId。
checkbox: {"dataValidation":{"type":"checkbox","checked":true}}
可与 text/richText 共存,也可单独使用(如 {dataValidation:{type:"none"}} 仅清除 DV 不写值)
@@ -169,6 +184,8 @@ dws sheet info --node NODE_ID --sheet-id SHEET_ID --format json,并读取 merg
- 只设样式或批量刷整片区域样式请用 dws sheet range set-style;写值同时设置少量 cell 样式可用 cellStyles
- 目标范围与已有合并区域冲突时,range update 会返回 MERGED_CELLS_CONFLICT;先用 sheet info 查看 mergedRanges,取消合并后写入,必要时再重新合并
- csv-put 的合并处理不同:目标区域含合并单元格时会打散合并并写入 CSV 值或公式
- SourceRange 在已验证的重命名、引用前插入行/列、删除引用前行的场景会自动调整;move-dimension 及其他未覆盖的删除/移动场景后先回读 sourceRangeStatus,仅 invalid 时重新选源写入
- 同一 cell 的 value/style 已写入但 SourceRange 校验失败时,服务端可返回 success=true,并通过 message 明确下拉未创建;必须检查 message,必要时重新读取确认
- 清空整片区域请用 dws sheet range clear`,
Example: ` # 写入文本
dws sheet range update --node NODE_ID --sheet-id SHEET_ID --range "A1:B2" \
@@ -0,0 +1,202 @@
package helpers
import (
"bytes"
"encoding/json"
"errors"
"strings"
"testing"
)
func executeSheetReadCompletionCommand(t *testing.T, caller *scriptedToolCaller, args ...string) (*bytes.Buffer, error) {
t.Helper()
installScriptedCaller(t, caller)
installSheetProductArgs(t)
stdout := &bytes.Buffer{}
deps.Out.w = stdout
cmd := newSheetCommand()
cmd.SilenceErrors = true
cmd.SilenceUsage = true
cmd.SetArgs(args)
return stdout, cmd.Execute()
}
func decodeSheetReadOutput(t *testing.T, stdout *bytes.Buffer) map[string]any {
t.Helper()
var got map[string]any
if err := json.Unmarshal(stdout.Bytes(), &got); err != nil {
t.Fatalf("decode sheet read output %q: %v", stdout.String(), err)
}
return got
}
func assertSheetReadCompletionMetadata(t *testing.T, got map[string]any) {
t.Helper()
if hasMore, ok := got["hasMore"].(bool); !ok || !hasMore {
t.Fatalf("hasMore = %#v, want true", got["hasMore"])
}
if got["resolvedRange"] != "A1:A30001" {
t.Fatalf("resolvedRange = %#v", got["resolvedRange"])
}
if got["returnedRange"] != "A1:A30000" {
t.Fatalf("returnedRange = %#v", got["returnedRange"])
}
reasons, ok := got["truncationReasons"].([]any)
if !ok || len(reasons) == 0 || reasons[0] != "max_cells" {
t.Fatalf("truncationReasons = %#v", got["truncationReasons"])
}
}
// csv-get intentionally projects the MCP response rather than inventing a
// client-side page model. Keep all completion metadata visible and make one
// request only, even when the server reports a partial success.
func TestSheetCSVGetPreservesCompletionMetadataWithoutAutoPagination(t *testing.T) {
const response = `{"csv":"[row=1]1\n","colIndices":["A"],"rowIndices":[1],"hasMore":true,"truncationReasons":["max_cells","max_chars"],"resolvedRange":"A1:A30001","returnedRange":"A1:A30000","message":"partial"}`
caller := &scriptedToolCaller{
format: "raw",
steps: []scriptedToolStep{{text: response}},
}
stdout, err := executeSheetReadCompletionCommand(t, caller,
"csv-get", "--node", "NODE")
if err != nil {
t.Fatalf("csv-get: %v", err)
}
if caller.calls != 1 || caller.tool != "get_range_as_csv" {
t.Fatalf("calls/tool = %d/%q, want one get_range_as_csv call", caller.calls, caller.tool)
}
if _, exists := caller.args["range"]; exists {
t.Fatalf("omitted range was unexpectedly synthesized by the CLI: %#v", caller.args["range"])
}
if got := strings.TrimSpace(stdout.String()); got != response {
t.Fatalf("raw csv-get response changed:\n got: %s\nwant: %s", got, response)
}
got := decodeSheetReadOutput(t, stdout)
assertSheetReadCompletionMetadata(t, got)
reasons := got["truncationReasons"].([]any)
if len(reasons) != 2 || reasons[1] != "max_chars" {
t.Fatalf("truncationReasons = %#v, want max_cells and max_chars", reasons)
}
}
// range read parses the get_cell_infos envelope to remove empty per-cell
// schema shells. That cleanup must not discard new top-level completion fields.
func TestSheetRangeReadPreservesCompletionMetadataWithoutAutoPagination(t *testing.T) {
const response = `{"cells":[[{"value":"1","dataValidation":{"type":null},"hyperlink":{"url":null}}]],"colIndices":["A"],"rowIndices":[1],"hasMore":true,"truncationReasons":["max_cells"],"resolvedRange":"A1:A30001","returnedRange":"A1:A30000","message":"partial"}`
caller := &scriptedToolCaller{
format: "json",
steps: []scriptedToolStep{{text: response}},
}
stdout, err := executeSheetReadCompletionCommand(t, caller,
"range", "read", "--node", "NODE")
if err != nil {
t.Fatalf("range read: %v", err)
}
if caller.calls != 1 || caller.tool != "get_cell_infos" {
t.Fatalf("calls/tool = %d/%q, want one get_cell_infos call", caller.calls, caller.tool)
}
if _, exists := caller.args["range"]; exists {
t.Fatalf("omitted range was unexpectedly synthesized by the CLI: %#v", caller.args["range"])
}
got := decodeSheetReadOutput(t, stdout)
assertSheetReadCompletionMetadata(t, got)
cells := got["cells"].([]any)
row := cells[0].([]any)
cell := row[0].(map[string]any)
if _, exists := cell["dataValidation"]; exists {
t.Fatalf("empty dataValidation shell was not removed: %#v", cell)
}
if _, exists := cell["hyperlink"]; exists {
t.Fatalf("empty hyperlink shell was not removed: %#v", cell)
}
}
// sheet +read calls get_cell_infos through this direct passthrough instead of
// callMCPToolCellInfos. Lock that third read surface to the same raw envelope
// contract and to one request only.
func TestSheetCellInfosDirectPassthroughPreservesCompletionMetadata(t *testing.T) {
const response = `{"cells":[[{"value":"1"}]],"colIndices":["A"],"rowIndices":[1],"hasMore":true,"truncationReasons":["max_cells"],"resolvedRange":"A1:A30001","returnedRange":"A1:A30000","message":"partial"}`
caller := &scriptedToolCaller{
format: "raw",
steps: []scriptedToolStep{{text: response}},
}
installScriptedCaller(t, caller)
stdout := &bytes.Buffer{}
deps.Out.w = stdout
if err := CallMCPToolOnServer("sheet", "get_cell_infos", map[string]any{"nodeId": "NODE"}); err != nil {
t.Fatalf("direct get_cell_infos: %v", err)
}
if caller.calls != 1 || caller.server != "sheet" || caller.tool != "get_cell_infos" {
t.Fatalf("calls/server/tool = %d/%q/%q", caller.calls, caller.server, caller.tool)
}
if got := strings.TrimSpace(stdout.String()); got != response {
t.Fatalf("raw get_cell_infos response changed:\n got: %s\nwant: %s", got, response)
}
assertSheetReadCompletionMetadata(t, decodeSheetReadOutput(t, stdout))
}
// Large-CP failure is not a pageable success. Both the raw csv-get path and
// both get_cell_infos paths must retain the backend code and safe user message
// exactly, matching the generic get_all_sheets error-classification path.
func TestSheetReadPreservesWorkbookSizeOverLimitError(t *testing.T) {
const response = `{"success":false,"errorCode":"forbidden.document.sizeOverLimit","errorMessage":"The workbook data is too large to process. Use a smaller copy or split the workbook, then try again."}`
tests := []struct {
name string
args []string
tool string
direct bool
}{
{name: "csv-get", args: []string{"csv-get", "--node", "NODE"}, tool: "get_range_as_csv"},
{name: "range-read", args: []string{"range", "read", "--node", "NODE"}, tool: "get_cell_infos"},
{name: "cell-infos-direct", tool: "get_cell_infos", direct: true},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
caller := &scriptedToolCaller{
format: "json",
steps: []scriptedToolStep{{text: response}},
}
var stdout *bytes.Buffer
var err error
if test.direct {
installScriptedCaller(t, caller)
stdout = &bytes.Buffer{}
deps.Out.w = stdout
err = CallMCPToolOnServer("sheet", test.tool, map[string]any{"nodeId": "NODE"})
} else {
stdout, err = executeSheetReadCompletionCommand(t, caller, test.args...)
}
if err == nil {
t.Fatal("large-CP response was accepted as success")
}
var cliErr *CLIError
if !errors.As(err, &cliErr) || cliErr.Code != CodeMCPToolError {
t.Fatalf("error = %T %v, want MCP tool error", err, err)
}
if cliErr.Message != response {
t.Fatalf("backend error payload changed:\n got: %s\nwant: %s", cliErr.Message, response)
}
for _, want := range []string{
"forbidden.document.sizeOverLimit",
"The workbook data is too large to process.",
} {
if !strings.Contains(err.Error(), want) {
t.Errorf("error %q does not preserve %q", err, want)
}
}
if stdout.Len() != 0 {
t.Fatalf("failure wrote a success body: %q", stdout.String())
}
if caller.calls != 1 || caller.tool != test.tool {
t.Fatalf("calls/tool = %d/%q, want one %s call", caller.calls, caller.tool, test.tool)
}
})
}
}
+40 -14
View File
@@ -239,22 +239,48 @@ func validateDataValidation(dvRaw any, path string) error {
}
switch dvType {
case "dropdown":
optionsRaw, ok := dv["options"]
if !ok {
return fmt.Errorf("%s: type=dropdown 必须包含 options 数组", path)
optionsRaw, hasOptions := dv["options"]
sourceRangeRaw, hasSourceRange := dv["sourceRange"]
hasOptions = hasOptions && optionsRaw != nil
hasSourceRange = hasSourceRange && sourceRangeRaw != nil
if hasOptions == hasSourceRange {
return fmt.Errorf("%s: type=dropdown 的 options 与 sourceRange 必须且只能包含一个", path)
}
options, ok := optionsRaw.([]any)
if !ok || len(options) == 0 {
return fmt.Errorf("%s.options 必须为非空数组", path)
}
for i, opt := range options {
optMap, ok := opt.(map[string]any)
if !ok {
return fmt.Errorf("%s.options[%d] 必须为 object(如 {\"value\":\"选项\"})", path, i)
if hasOptions {
options, ok := optionsRaw.([]any)
if !ok || len(options) == 0 {
return fmt.Errorf("%s.options 必须为非空数组", path)
}
val, ok := optMap["value"].(string)
if !ok || val == "" {
return fmt.Errorf("%s.options[%d].value 必须为非空字符串", path, i)
for i, opt := range options {
optMap, ok := opt.(map[string]any)
if !ok {
return fmt.Errorf("%s.options[%d] 必须为 object(如 {\"value\":\"选项\"})", path, i)
}
val, ok := optMap["value"].(string)
if !ok || val == "" {
return fmt.Errorf("%s.options[%d].value 必须为非空字符串", path, i)
}
}
} else {
sourceRange, ok := sourceRangeRaw.(map[string]any)
if !ok {
return fmt.Errorf("%s.sourceRange 必须为 object", path)
}
if _, exists := sourceRange["colors"]; exists {
return fmt.Errorf("%s.sourceRange.colors 暂不支持", path)
}
sheetID, ok := sourceRange["sheetId"].(string)
if !ok || strings.TrimSpace(sheetID) == "" {
return fmt.Errorf("%s.sourceRange.sheetId 必须为非空字符串", path)
}
a1Notation, ok := sourceRange["a1Notation"].(string)
if !ok || strings.TrimSpace(a1Notation) == "" {
return fmt.Errorf("%s.sourceRange.a1Notation 必须为非空字符串", path)
}
}
if multi, exists := dv["enableMultiSelect"]; exists {
if _, ok := multi.(bool); !ok {
return fmt.Errorf("%s.enableMultiSelect 必须为 boolean", path)
}
}
case "checkbox":
+1 -1
View File
@@ -733,7 +733,7 @@ var RecordQuery = shortcut.Shortcut{
if rt.Changed("cursor") {
params["cursor"] = rt.Str("cursor")
}
return rt.CallMCP("query_records", params)
return executeRecordQuery(rt, params)
},
}
+24 -11
View File
@@ -98,7 +98,7 @@ func executeRecordDeleteBatches(rt *shortcut.RuntimeContext) error {
step.Status = "unknown"
step.Error = writeErr.Error()
}
remaining, verifyErr := queryRecordsByIDs(rt, baseID, tableID, batch)
remaining, verifyErr := queryDeletedRecordsByIDs(rt, baseID, tableID, batch)
if verifyErr == nil && len(remaining) > 0 {
verifyErr = fmt.Errorf("read-back still contains deleted record IDs: %s", strings.Join(recordIDs(remaining), ","))
}
@@ -303,20 +303,33 @@ func verifyUpsertBatch(rt *shortcut.RuntimeContext, baseID, tableID string, batc
}
func queryRecordsByIDs(rt *shortcut.RuntimeContext, baseID, tableID string, ids []string) ([]map[string]any, error) {
data, err := rt.CallMCPData(serverMain, "query_records", map[string]any{
"baseId": baseID, "tableId": tableID, "recordIds": ids, "limit": len(ids),
})
window, err := queryRecordWindow(rt, map[string]any{
"baseId": baseID, "tableId": tableID, "recordIds": ids,
}, len(ids))
if err != nil {
return nil, err
}
records, found := findRecords(data)
if !found {
return nil, fmt.Errorf("query_records read-back is missing records")
// Exact-ID verification below compares every requested ID with the returned
// records. The service can publish a continuation even after all requested
// IDs are present, so hasMore is not evidence that this bounded read-back is
// incomplete.
return window.Records, nil
}
func queryDeletedRecordsByIDs(rt *shortcut.RuntimeContext, baseID, tableID string, ids []string) ([]map[string]any, error) {
remaining := make([]map[string]any, 0)
for offset := 0; offset < len(ids); offset += recordQueryServicePageSize {
end := minInt(offset+recordQueryServicePageSize, len(ids))
chunk := ids[offset:end]
window, err := queryRecordWindow(rt, map[string]any{
"baseId": baseID, "tableId": tableID, "recordIds": chunk,
}, len(chunk))
if err != nil {
return nil, err
}
remaining = append(remaining, window.Records...)
}
if responseHasMore(data) {
return nil, fmt.Errorf("query_records read-back is incomplete")
}
return records, nil
return remaining, nil
}
func recordIDs(records []map[string]any) []string {
+233 -1
View File
@@ -10,13 +10,245 @@ import (
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/shortcut"
)
// query_records is backed by a service pager whose stable typed-envelope
// boundary is 20 records. Asking that layer to aggregate more than one page can
// drop the envelope fields used by the record projector, turning a successful
// read into an empty collection. Keep every remote request on one service page
// and aggregate only after each page has passed the CLI's shape validation.
const (
recordQueryServicePageSize = 20
recordQueryMaxConsecutiveEmptyPages = 3
)
type recordQueryWindow struct {
Records []map[string]any
HasMore bool
NextCursor string
Pages int
TotalCount any
HasTotalCount bool
}
func queryRecordWindow(rt *shortcut.RuntimeContext, params map[string]any, limit int) (recordQueryWindow, error) {
if limit <= 0 {
return recordQueryWindow{}, fmt.Errorf("record query limit must be positive, got %d", limit)
}
request := cloneAnyMap(params)
cursor, _ := request["cursor"].(string)
cursor = strings.TrimSpace(cursor)
seen := map[string]bool{}
if cursor != "" {
seen[cursor] = true
}
window := recordQueryWindow{Records: make([]map[string]any, 0, limit)}
consecutiveEmptyPages := 0
for len(window.Records) < limit {
pageSize := minInt(recordQueryServicePageSize, limit-len(window.Records))
request["limit"] = pageSize
if cursor == "" {
delete(request, "cursor")
} else {
request["cursor"] = cursor
}
data, err := rt.CallMCPData(serverMain, "query_records", request)
if err != nil {
return recordQueryWindow{}, err
}
if !window.HasTotalCount {
window.TotalCount, window.HasTotalCount = responseTotalCount(data)
}
records, found := findRecords(data)
if !found {
if explicitEmptyRecordQuery(data) {
window.Pages++
window.HasMore = false
window.NextCursor = ""
return window, nil
}
return recordQueryWindow{}, fmt.Errorf("query_records page %d is missing the records collection", window.Pages+1)
}
window.Records = append(window.Records, records...)
window.Pages++
if !responseHasMore(data) {
window.HasMore = false
window.NextCursor = ""
return window, nil
}
if len(records) == 0 {
consecutiveEmptyPages++
if consecutiveEmptyPages >= recordQueryMaxConsecutiveEmptyPages {
return recordQueryWindow{}, fmt.Errorf(
"query_records made no progress for %d consecutive pages",
consecutiveEmptyPages,
)
}
} else {
consecutiveEmptyPages = 0
}
next := responseCursor(data)
if next == "" {
return recordQueryWindow{}, fmt.Errorf("query_records page %d reports more data but no next cursor", window.Pages)
}
if seen[next] {
return recordQueryWindow{}, fmt.Errorf("query_records cursor cycle detected at %q", next)
}
seen[next] = true
cursor = next
window.HasMore = true
window.NextCursor = next
}
return window, nil
}
// explicitEmptyRecordQuery recognizes the service's reviewed zero-match wire
// shape. It is deliberately stricter than "records is absent": only a complete
// success envelope with an empty data object is accepted as an empty set.
func explicitEmptyRecordQuery(data map[string]any) bool {
if data == nil || data["success"] != true || data["status"] != "success" {
return false
}
payload, ok := data["data"].(map[string]any)
if !ok || len(payload) != 0 {
return false
}
if rawError, exists := data["error"]; exists && rawError != nil {
errorObject, ok := rawError.(map[string]any)
if !ok || len(errorObject) != 0 {
return false
}
}
return true
}
func executeRecordQuery(rt *shortcut.RuntimeContext, params map[string]any) error {
limit := 100
if rt.Changed("limit") {
limit = rt.Int("limit")
}
if limit < 1 || limit > recordBatchSize {
return fmt.Errorf("--limit must be in [1,%d], got %d", recordBatchSize, limit)
}
requestedIDs, exactIDQuery, err := recordQueryRequestedIDs(params)
if err != nil {
return err
}
if exactIDQuery {
if len(requestedIDs) > recordBatchSize {
return fmt.Errorf("--record-ids accepts at most %d unique IDs, got %d", recordBatchSize, len(requestedIDs))
}
params = cloneAnyMap(params)
params["recordIds"] = requestedIDs
limit = minInt(limit, len(requestedIDs))
}
if rt.DryRun() {
return rt.Output(map[string]any{
"dry_run": true,
"executed": false,
"tool": "query_records",
"arguments": params,
})
}
window, err := queryRecordWindow(rt, params, limit)
if err != nil {
return err
}
if exactIDQuery {
complete, err := validateExactRecordQuery(window.Records, requestedIDs)
if err != nil {
return err
}
if complete {
// The service may advertise unrelated continuation after every
// requested ID is already present. Exact-ID completion is stronger
// evidence than that residual pager state.
window.HasMore = false
window.NextCursor = ""
}
}
records := make([]any, 0, len(window.Records))
for _, record := range window.Records {
records = append(records, record)
}
data := map[string]any{
"records": records,
"hasMore": window.HasMore,
"page": window.Pages,
"size": len(window.Records),
}
if window.NextCursor != "" {
data["nextCursor"] = window.NextCursor
}
if window.HasTotalCount {
data["totalCount"] = window.TotalCount
}
return rt.Output(map[string]any{
"success": true,
"status": "success",
"data": data,
})
}
func responseTotalCount(data map[string]any) (any, bool) {
for {
if totalCount, exists := data["totalCount"]; exists {
return totalCount, true
}
nested, ok := data["data"].(map[string]any)
if !ok {
return nil, false
}
data = nested
}
}
func recordQueryRequestedIDs(params map[string]any) ([]string, bool, error) {
raw, exists := params["recordIds"]
if !exists {
return nil, false, nil
}
values, ok := raw.([]string)
if !ok {
return nil, false, fmt.Errorf("recordIds must be a string list, got %T", raw)
}
ids, err := parseRecordIDs(values)
if err != nil {
return nil, false, err
}
return ids, true, nil
}
func validateExactRecordQuery(records []map[string]any, requestedIDs []string) (bool, error) {
wanted := make(map[string]bool, len(requestedIDs))
for _, id := range requestedIDs {
wanted[id] = true
}
seen := make(map[string]bool, len(records))
for index, record := range records {
id := recordID(record)
if id == "" {
return false, fmt.Errorf("query_records exact-ID result %d is missing recordId", index)
}
if !wanted[id] {
return false, fmt.Errorf("query_records exact-ID result contains unexpected recordId %q", id)
}
if seen[id] {
return false, fmt.Errorf("query_records exact-ID result contains duplicate recordId %q", id)
}
seen[id] = true
}
return len(seen) == len(wanted), nil
}
func queryAllRecords(rt *shortcut.RuntimeContext, params map[string]any, maxRecords int) ([]map[string]any, error) {
all := make([]map[string]any, 0)
cursor := ""
seen := map[string]bool{}
for page := 0; ; page++ {
request := cloneAnyMap(params)
request["limit"] = recordBatchSize
request["limit"] = recordQueryServicePageSize
if cursor != "" {
request["cursor"] = cursor
}
@@ -9,6 +9,7 @@ import (
"encoding/json"
"errors"
"fmt"
"slices"
"strings"
"testing"
@@ -265,6 +266,446 @@ func recordListJSON(t *testing.T, records []map[string]any) string {
return string(raw)
}
func pagedRecordQueryResponse(t *testing.T, records []map[string]any, args map[string]any) string {
t.Helper()
limit, ok := args["limit"].(int)
if !ok || limit < 1 || limit > recordQueryServicePageSize {
t.Fatalf("query_records limit = %#v, want 1..%d", args["limit"], recordQueryServicePageSize)
}
offset := 0
if cursor, _ := args["cursor"].(string); cursor != "" {
if _, err := fmt.Sscanf(cursor, "offset-%d", &offset); err != nil {
t.Fatalf("invalid test cursor %q: %v", cursor, err)
}
}
end := minInt(offset+limit, len(records))
pageRecords := records[offset:end]
data := map[string]any{"records": pageRecords, "hasMore": end < len(records), "totalCount": len(records)}
if end < len(records) {
data["nextCursor"] = fmt.Sprintf("offset-%d", end)
}
raw, err := json.Marshal(map[string]any{
"success": true,
"hasMore": end < len(records),
"page": offset/recordQueryServicePageSize + 1,
"size": limit,
"data": data,
})
if err != nil {
t.Fatal(err)
}
return string(raw)
}
func runRecordQueryShortcutCLI(t *testing.T, caller *upsertByKeyCaller, limit int, extra ...string) (map[string]any, error) {
t.Helper()
helpers.InitDepsForTest(t, caller)
root := &cobra.Command{Use: "dws", SilenceErrors: true, SilenceUsage: true}
root.PersistentFlags().Bool("yes", false, "")
root.PersistentFlags().Bool("dry-run", false, "")
root.PersistentFlags().String("format", "json", "")
root.AddCommand(shortcut.Commands()...)
stdout := &bytes.Buffer{}
root.SetOut(stdout)
root.SetErr(&bytes.Buffer{})
args := []string{"aitable", "+record-query", "--base-id", "base", "--table-id", "table", "--limit", fmt.Sprint(limit)}
root.SetArgs(append(args, extra...))
err := root.Execute()
if stdout.Len() == 0 {
return nil, err
}
var payload map[string]any
if decodeErr := json.Unmarshal(stdout.Bytes(), &payload); decodeErr != nil {
t.Fatalf("decode record query output %q: %v", stdout.String(), decodeErr)
}
return payload, err
}
func TestCrossPlatformCoverageRecordQueryServicePageBoundariesE2E(t *testing.T) {
for _, size := range []int{20, 21, 22, 100} {
t.Run(fmt.Sprintf("size_%d", size), func(t *testing.T) {
records := updateFixtureRecords(0, size, "可见")
caller := &upsertByKeyCaller{}
caller.callFn = func(_ int, _, tool string, args map[string]any) (string, error) {
if tool != "query_records" {
return "", fmt.Errorf("unexpected tool %s", tool)
}
return pagedRecordQueryResponse(t, records, args), nil
}
payload, err := runRecordQueryShortcutCLI(t, caller, size)
if err != nil {
t.Fatalf("record query size %d: %v", size, err)
}
data, ok := payload["data"].(map[string]any)
if !ok || payload["success"] != true || data["hasMore"] != false || data["totalCount"] != float64(size) || len(data["records"].([]any)) != size {
t.Fatalf("record query size %d payload = %#v", size, payload)
}
wantCalls := (size + recordQueryServicePageSize - 1) / recordQueryServicePageSize
if len(caller.calls) != wantCalls || data["page"] != float64(wantCalls) || data["size"] != float64(size) {
t.Fatalf("record query size %d calls=%d payload=%#v, want calls=%d", size, len(caller.calls), payload, wantCalls)
}
})
}
}
func TestCrossPlatformCoverageRecordQueryDryRunStopsBeforeTransportE2E(t *testing.T) {
caller := &upsertByKeyCaller{}
payload, err := runRecordQueryShortcutCLI(t, caller, 5, "--query", "needle", "--dry-run")
if err != nil {
t.Fatalf("record query dry-run error = %v", err)
}
if len(caller.calls) != 0 {
t.Fatalf("record query dry-run crossed transport: %#v", caller.calls)
}
arguments, _ := payload["arguments"].(map[string]any)
if payload["dry_run"] != true || payload["executed"] != false || payload["tool"] != "query_records" ||
arguments["baseId"] != "base" || arguments["tableId"] != "table" || arguments["keyword"] != "needle" || arguments["limit"] != float64(5) {
t.Fatalf("record query dry-run payload = %#v", payload)
}
}
func TestCrossPlatformCoverageRecordQueryDryRunRejectsInvalidLocalPlanE2E(t *testing.T) {
for _, testCase := range []struct {
name string
limit int
extra []string
want string
}{
{name: "zero limit", limit: 0, extra: []string{"--dry-run"}, want: "--limit must be in [1,100]"},
{name: "excessive limit", limit: recordBatchSize + 1, extra: []string{"--dry-run"}, want: "--limit must be in [1,100]"},
{name: "empty record IDs", limit: 100, extra: []string{"--record-ids", " ", "--dry-run"}, want: "至少包含一个非空 recordId"},
{
name: "excessive record IDs",
limit: 100,
extra: []string{"--record-ids", strings.Join(recordIDFixtures(recordBatchSize+1), ","), "--dry-run"},
want: "at most 100 unique IDs",
},
} {
t.Run(testCase.name, func(t *testing.T) {
caller := &upsertByKeyCaller{}
payload, err := runRecordQueryShortcutCLI(t, caller, testCase.limit, testCase.extra...)
if err == nil || !strings.Contains(err.Error(), testCase.want) {
t.Fatalf("record query dry-run error = %v, want %q (payload=%#v)", err, testCase.want, payload)
}
if len(caller.calls) != 0 {
t.Fatalf("invalid record query dry-run crossed transport: %#v", caller.calls)
}
})
}
}
func TestCrossPlatformCoverageRecordQueryExactIDsIgnoreResidualContinuationE2E(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(call int, _, tool string, args map[string]any) (string, error) {
if tool != "query_records" {
return "", fmt.Errorf("unexpected tool %s", tool)
}
if call != 0 {
return "", errors.New("exact-ID query followed residual continuation")
}
if args["limit"] != 2 {
t.Fatalf("exact-ID service limit = %#v, want 2", args["limit"])
}
ids, ok := args["recordIds"].([]string)
if !ok || !slices.Equal(ids, []string{"r1", "r2"}) {
t.Fatalf("exact-ID request IDs = %#v", args["recordIds"])
}
return `{"success":true,"data":{"records":[{"recordId":"r1"},{"recordId":"r2"}],"hasMore":true,"nextCursor":"residual"}}`, nil
}}
payload, err := runRecordQueryShortcutCLI(t, caller, 100, "--record-ids", "r2,r1,r2")
data, _ := payload["data"].(map[string]any)
if err != nil || len(caller.calls) != 1 || data["hasMore"] != false || data["size"] != float64(2) {
t.Fatalf("exact-ID payload=%#v calls=%d err=%v", payload, len(caller.calls), err)
}
if _, exists := data["nextCursor"]; exists {
t.Fatalf("exact-ID payload retained residual cursor: %#v", payload)
}
}
func TestCrossPlatformCoverageRecordQueryFailureAndContinuationBranches(t *testing.T) {
t.Run("non-positive window limit", func(t *testing.T) {
caller := &upsertByKeyCaller{}
helpers.InitDepsForTest(t, caller)
rt := shortcut.RuntimeContextForTest(&cobra.Command{Use: "query"}, RecordQuery)
if _, err := queryRecordWindow(rt, nil, 0); err == nil {
t.Fatal("non-positive record query window limit accepted")
}
})
t.Run("initial cursor cycle", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(_ int, _, tool string, _ map[string]any) (string, error) {
if tool != "query_records" {
return "", fmt.Errorf("unexpected tool %s", tool)
}
return `{"success":true,"data":{"records":[{"recordId":"r1"}],"hasMore":true,"nextCursor":"seed"}}`, nil
}}
helpers.InitDepsForTest(t, caller)
rt := shortcut.RuntimeContextForTest(&cobra.Command{Use: "query"}, RecordQuery)
if _, err := queryRecordWindow(rt, map[string]any{"cursor": " seed "}, 2); err == nil || !strings.Contains(err.Error(), "cursor cycle") {
t.Fatalf("initial cursor cycle error = %v", err)
}
})
t.Run("strict explicit empty shape", func(t *testing.T) {
for _, data := range []map[string]any{
{"success": true, "status": "success", "data": map[string]any{"unexpected": true}},
{"success": true, "status": "success", "data": map[string]any{}, "error": "bad"},
} {
if explicitEmptyRecordQuery(data) {
t.Fatalf("invalid empty-query envelope accepted: %#v", data)
}
}
})
t.Run("explicit empty page", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(_ int, _, _ string, _ map[string]any) (string, error) {
return `{"success":true,"status":"success","error":{},"data":{}}`, nil
}}
helpers.InitDepsForTest(t, caller)
rt := shortcut.RuntimeContextForTest(&cobra.Command{Use: "query"}, RecordQuery)
window, err := queryRecordWindow(rt, map[string]any{}, 1)
if err != nil || len(window.Records) != 0 || window.Pages != 1 || window.HasMore || window.NextCursor != "" {
t.Fatalf("explicit empty window=%#v err=%v", window, err)
}
})
t.Run("missing records collection", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(_ int, _, _ string, _ map[string]any) (string, error) {
return `{}`, nil
}}
helpers.InitDepsForTest(t, caller)
rt := shortcut.RuntimeContextForTest(&cobra.Command{Use: "query"}, RecordQuery)
if _, err := queryRecordWindow(rt, map[string]any{}, 1); err == nil || !strings.Contains(err.Error(), "missing the records collection") {
t.Fatalf("missing records error = %v", err)
}
})
t.Run("changing cursors cannot hide sustained empty pages", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(call int, _, tool string, _ map[string]any) (string, error) {
if tool != "query_records" {
return "", fmt.Errorf("unexpected tool %s", tool)
}
return fmt.Sprintf(
`{"success":true,"data":{"records":[],"hasMore":true,"nextCursor":"empty-%d"}}`,
call,
), nil
}}
helpers.InitDepsForTest(t, caller)
rt := shortcut.RuntimeContextForTest(&cobra.Command{Use: "query"}, RecordQuery)
if _, err := queryRecordWindow(rt, map[string]any{}, 1); err == nil || !strings.Contains(err.Error(), "no progress") {
t.Fatalf("sustained empty pages error = %v", err)
}
if len(caller.calls) != recordQueryMaxConsecutiveEmptyPages {
t.Fatalf("query calls = %d, want %d", len(caller.calls), recordQueryMaxConsecutiveEmptyPages)
}
})
t.Run("invalid shortcut limit", func(t *testing.T) {
caller := &upsertByKeyCaller{}
helpers.InitDepsForTest(t, caller)
cmd := &cobra.Command{Use: "query"}
cmd.Flags().Int("limit", 0, "")
if err := cmd.Flags().Set("limit", "0"); err != nil {
t.Fatal(err)
}
rt := shortcut.RuntimeContextForTest(cmd, RecordQuery)
if err := executeRecordQuery(rt, map[string]any{}); err == nil {
t.Fatal("invalid shortcut limit accepted")
}
})
t.Run("invalid exact-ID request", func(t *testing.T) {
for _, testCase := range []struct {
name string
params map[string]any
want string
}{
{name: "wrong type", params: map[string]any{"recordIds": 1}, want: "string list"},
{name: "empty", params: map[string]any{"recordIds": []string{" "}}, want: "至少包含一个非空 recordId"},
{name: "over service limit", params: map[string]any{"recordIds": recordIDFixtures(recordBatchSize + 1)}, want: "at most 100 unique IDs"},
} {
t.Run(testCase.name, func(t *testing.T) {
caller := &upsertByKeyCaller{}
helpers.InitDepsForTest(t, caller)
rt := shortcut.RuntimeContextForTest(&cobra.Command{Use: "query"}, RecordQuery)
err := executeRecordQuery(rt, testCase.params)
if err == nil || !strings.Contains(err.Error(), testCase.want) {
t.Fatalf("exact-ID request error = %v, want %q", err, testCase.want)
}
if len(caller.calls) != 0 {
t.Fatalf("invalid exact-ID request made %d service calls", len(caller.calls))
}
})
}
})
t.Run("invalid exact-ID response", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(_ int, _, tool string, _ map[string]any) (string, error) {
if tool != "query_records" {
return "", fmt.Errorf("unexpected tool %s", tool)
}
return `{"success":true,"data":{"records":[{"recordId":"other"}]}}`, nil
}}
if _, err := runRecordQueryShortcutCLI(t, caller, 100, "--record-ids", "r1"); err == nil || !strings.Contains(err.Error(), "unexpected recordId") {
t.Fatalf("unexpected exact-ID response error = %v", err)
}
if _, err := validateExactRecordQuery([]map[string]any{{}}, []string{"r1"}); err == nil || !strings.Contains(err.Error(), "missing recordId") {
t.Fatalf("missing exact-ID response error = %v", err)
}
duplicate := []map[string]any{{"recordId": "r1"}, {"recordId": "r1"}}
if _, err := validateExactRecordQuery(duplicate, []string{"r1", "r2"}); err == nil || !strings.Contains(err.Error(), "duplicate recordId") {
t.Fatalf("duplicate exact-ID response error = %v", err)
}
})
t.Run("shortcut query failure", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(_ int, _, _ string, _ map[string]any) (string, error) {
return "", errors.New("query unavailable")
}}
helpers.InitDepsForTest(t, caller)
rt := shortcut.RuntimeContextForTest(&cobra.Command{Use: "query"}, RecordQuery)
if err := executeRecordQuery(rt, map[string]any{}); err == nil || !strings.Contains(err.Error(), "query unavailable") {
t.Fatalf("shortcut query failure = %v", err)
}
})
t.Run("bounded continuation is published", func(t *testing.T) {
records := updateFixtureRecords(0, 21, "visible")
caller := &upsertByKeyCaller{callFn: func(_ int, _, tool string, args map[string]any) (string, error) {
if tool != "query_records" {
return "", fmt.Errorf("unexpected tool %s", tool)
}
return pagedRecordQueryResponse(t, records, args), nil
}}
payload, err := runRecordQueryShortcutCLI(t, caller, 20)
data, _ := payload["data"].(map[string]any)
if err != nil || data["hasMore"] != true || data["nextCursor"] != "offset-20" {
t.Fatalf("bounded continuation payload=%#v err=%v", payload, err)
}
})
t.Run("delete readback follows active continuation to absence", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(call int, _, tool string, args map[string]any) (string, error) {
if tool != "query_records" {
return `{"deletedCount":1}`, nil
}
if call == 1 {
if _, exists := args["cursor"]; exists {
t.Fatalf("first deletion readback unexpectedly has cursor: %#v", args)
}
return `{"success":true,"data":{"records":[],"hasMore":true,"nextCursor":"c1"}}`, nil
}
if args["cursor"] != "c1" {
t.Fatalf("continued deletion readback cursor = %#v", args["cursor"])
}
return `{"success":true,"data":{"records":[],"hasMore":false}}`, nil
}}
out, err := runRecordDeleteCLI(t, caller, []string{"r1"})
if err != nil || out == "" || len(caller.calls) != 3 {
t.Fatalf("delete continued absence output=%q error=%v calls=%#v", out, err, caller.calls)
}
})
t.Run("delete readback follows active continuation to remaining record", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(call int, _, tool string, _ map[string]any) (string, error) {
if tool != "query_records" {
return `{"deletedCount":1}`, nil
}
if call == 1 {
return `{"success":true,"data":{"records":[],"hasMore":true,"nextCursor":"c1"}}`, nil
}
return `{"success":true,"data":{"records":[{"recordId":"r1"}],"hasMore":false}}`, nil
}}
out, err := runRecordDeleteCLI(t, caller, []string{"r1"})
if err == nil || out != "" || len(caller.calls) != 3 {
t.Fatalf("delete continued remaining output=%q error=%v calls=%#v", out, err, caller.calls)
}
var typed *apperrors.Error
if !errors.As(err, &typed) || typed.Reason != "aitable_composite_unknown" {
t.Fatalf("delete continued remaining error = %#v", err)
}
})
t.Run("delete readback fails closed on empty continuation stall", func(t *testing.T) {
caller := &upsertByKeyCaller{callFn: func(call int, _, tool string, _ map[string]any) (string, error) {
if tool != "query_records" {
return `{"deletedCount":1}`, nil
}
return fmt.Sprintf(`{"success":true,"data":{"records":[],"hasMore":true,"nextCursor":"c%d"}}`, call), nil
}}
out, err := runRecordDeleteCLI(t, caller, []string{"r1"})
if err == nil || out != "" || len(caller.calls) != recordQueryMaxConsecutiveEmptyPages+1 {
t.Fatalf("delete stalled continuation output=%q error=%v calls=%#v", out, err, caller.calls)
}
var typed *apperrors.Error
if !errors.As(err, &typed) || typed.Reason != "aitable_composite_unknown" {
t.Fatalf("delete stalled continuation error = %#v", err)
}
})
}
func TestCrossPlatformCoverageRecordWriteReadbackUsesStableServicePagesE2E(t *testing.T) {
for _, operation := range []string{"update", "upsert", "delete"} {
for _, size := range []int{21, 22, 100} {
t.Run(fmt.Sprintf("%s_%d", operation, size), func(t *testing.T) {
records := updateFixtureRecords(0, size, "完成")
caller := &upsertByKeyCaller{}
queryCalls := 0
caller.callFn = func(_ int, _, tool string, args map[string]any) (string, error) {
if tool != "query_records" {
return `{"success":true}`, nil
}
queryCalls++
if operation == "delete" {
ids := args["recordIds"].([]string)
empty := make([]map[string]any, len(ids))
response := pagedRecordQueryResponse(t, empty, args)
var payload map[string]any
if err := json.Unmarshal([]byte(response), &payload); err != nil {
t.Fatal(err)
}
page := payload["data"].(map[string]any)
page["records"] = []any{}
raw, _ := json.Marshal(payload)
return string(raw), nil
}
response := pagedRecordQueryResponse(t, records, args)
var payload map[string]any
if err := json.Unmarshal([]byte(response), &payload); err != nil {
t.Fatal(err)
}
page := payload["data"].(map[string]any)
if len(page["records"].([]any))+recordQueryServicePageSize*(queryCalls-1) >= size {
// The live service can retain a continuation after it has
// returned every requested record ID. Exact-ID verification
// must use ID coverage, not this advisory bit.
page["hasMore"] = true
page["nextCursor"] = "ignored-after-complete-id-set"
payload["hasMore"] = true
}
raw, _ := json.Marshal(payload)
return string(raw), nil
}
var out string
var err error
if operation == "delete" {
out, err = runRecordDeleteCLI(t, caller, recordIDs(records))
} else {
out, err = runRecordBatchCLI(t, caller, "+record-"+operation, records)
}
if err != nil || out == "" {
t.Fatalf("%s size %d false negative: output=%q err=%v", operation, size, out, err)
}
wantQueries := (size + recordQueryServicePageSize - 1) / recordQueryServicePageSize
if queryCalls != wantQueries {
t.Fatalf("%s size %d query calls=%d, want %d", operation, size, queryCalls, wantQueries)
}
})
}
}
}
func TestCrossPlatformCoverageRecordUpdateAutoChunksAndVerifiesE2E(t *testing.T) {
records := updateFixtureRecords(0, 101, "完成")
caller := &upsertByKeyCaller{steps: []upsertByKeyStep{
@@ -395,6 +836,10 @@ func TestCrossPlatformCoverageRecordDeleteAutoChunksAndProvesAbsenceE2E(t *testi
caller := &upsertByKeyCaller{steps: []upsertByKeyStep{
{text: `{"deletedCount":100}`},
{text: `{"data":{"records":[]}}`},
{text: `{"data":{"records":[]}}`},
{text: `{"data":{"records":[]}}`},
{text: `{"data":{"records":[]}}`},
{text: `{"data":{"records":[]}}`},
{text: `{"deletedCount":1}`},
{text: `{"data":{"records":[]}}`},
}}
@@ -407,7 +852,7 @@ func TestCrossPlatformCoverageRecordDeleteAutoChunksAndProvesAbsenceE2E(t *testi
t.Fatalf("delete output missing %s: %s", want, out)
}
}
if len(caller.calls) != 4 || caller.calls[0].tool != "delete_records" || caller.calls[1].tool != "query_records" {
if len(caller.calls) != 8 || caller.calls[0].tool != "delete_records" || caller.calls[1].tool != "query_records" || caller.calls[6].tool != "delete_records" {
t.Fatalf("delete call sequence = %#v", caller.calls)
}
}
@@ -421,6 +866,17 @@ func TestCrossPlatformCoverageRecordDeleteEmptyReplyRecoveredOnlyByAbsenceE2E(t
}
})
t.Run("explicit service success with empty data proves deletion", func(t *testing.T) {
caller := &upsertByKeyCaller{steps: []upsertByKeyStep{
{text: `{"deletedCount":1}`},
{text: `{"success":true,"status":"success","error":{},"data":{}}`},
}}
out, err := runRecordDeleteCLI(t, caller, []string{"r1"})
if err != nil || !strings.Contains(out, `"status": "verified_absent"`) {
t.Fatalf("delete explicit empty success = output:%q err:%v", out, err)
}
})
t.Run("missing records contract is unknown", func(t *testing.T) {
caller := &upsertByKeyCaller{steps: []upsertByKeyStep{{text: `{"deletedCount":1}`}, {text: `{"data":{}}`}}}
out, err := runRecordDeleteCLI(t, caller, []string{"r1"})
@@ -450,6 +906,10 @@ func TestCrossPlatformCoverageRecordDeletePartialCheckpointE2E(t *testing.T) {
caller := &upsertByKeyCaller{steps: []upsertByKeyStep{
{text: `{"deletedCount":100}`},
{text: `{"records":[]}`},
{text: `{"records":[]}`},
{text: `{"records":[]}`},
{text: `{"records":[]}`},
{text: `{"records":[]}`},
{text: `{"deletedCount":0}`},
{text: `{"records":[{"recordId":"r100","cells":{}}]}`},
}}
+382 -33
View File
@@ -4,10 +4,14 @@
package doc
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"os"
"path/filepath"
"reflect"
"strings"
"time"
"unicode"
@@ -18,6 +22,9 @@ import (
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/localio"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/shortcut"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/shortcut/docresolver"
"github.com/yuin/goldmark"
"github.com/yuin/goldmark/extension"
"github.com/yuin/goldmark/renderer/html"
)
var (
@@ -28,6 +35,21 @@ var (
docMkdirTemp = os.MkdirTemp
docRemoveAll = os.RemoveAll
docDownload = localio.Download
docVerifyWait = waitForDocVerification
docVerifyDelays = []time.Duration{250 * time.Millisecond, 500 * time.Millisecond, time.Second, 2 * time.Second, 4 * time.Second, 8 * time.Second}
docMarkdown = goldmark.New(
goldmark.WithExtensions(extension.Table),
goldmark.WithRendererOptions(html.WithUnsafe()),
)
docMarkdownConvert = func(source []byte, writer io.Writer) error {
return docMarkdown.Convert(source, writer)
}
)
const (
docBlockReadPageSize = 50
docBlockReadMaxItems = 5000
docMarkdownVerifyMax = 2 * 1024 * 1024
)
var Create = shortcut.Shortcut{
@@ -133,7 +155,9 @@ var Create = shortcut.Shortcut{
verifyTool = "get_document_content"
verifyParams["format"] = format
}
verification, err := rt.CallMCPData(productDoc, verifyTool, verifyParams)
verification, err := readDocVerification(rt, verifyTool, verifyParams, func(data map[string]any) bool {
return content == "" || verifyUpdatedDocumentContent(data, content, "overwrite", format)
})
if err != nil {
return docVerificationError("doc.create", "verify", nodeID, err, append(steps, map[string]any{"name": "verify", "status": "failed"}))
}
@@ -422,7 +446,9 @@ var CheckpointUpdate = shortcut.Shortcut{
append(steps, map[string]any{"name": "update", "status": "failed"}, map[string]any{"name": "verify", "status": "not_started"}))
}
steps = append(steps, map[string]any{"name": "update", "status": "success"})
verification, err := rt.CallMCPData(productDoc, "get_document_content", map[string]any{"nodeId": rt.Str("node"), "format": "markdown"})
verification, err := readDocVerification(rt, "get_document_content", map[string]any{"nodeId": rt.Str("node"), "format": "markdown"}, func(data map[string]any) bool {
return verifyUpdatedDocumentContent(data, content, rt.Str("mode"), "markdown")
})
if err != nil {
return checkpointPartialWriteError(rt.Str("node"), checkpoint, "verify", "doc_checkpoint_verification_failed", err,
append(steps, map[string]any{"name": "verify", "status": "failed"}))
@@ -539,7 +565,7 @@ func executeUpdate(rt *shortcut.RuntimeContext) error {
}
referenceBlockID := rt.Str("after-block-id")
return executeVerifiedDocMutation(rt, "doc.update", "insert_document_block", params, node,
"list_document_blocks", map[string]any{"nodeId": node, "format": verificationFormat},
"list_document_blocks", map[string]any{"nodeId": node, "format": verificationFormat, "__allBlocks": true},
func(result, data map[string]any) bool {
return verifyInsertedBlock(result, data, referenceBlockID, content, rt.Str("doc-format"))
})
@@ -553,14 +579,14 @@ func executeUpdate(rt *shortcut.RuntimeContext) error {
params["element"] = map[string]any{"blockType": "paragraph", "paragraph": map[string]any{"text": content}}
}
return executeVerifiedDocMutation(rt, "doc.update", "update_document_block", params, node,
"list_document_blocks", map[string]any{"nodeId": node, "blockId": blockID, "format": verificationFormat},
"list_document_blocks", map[string]any{"nodeId": node, "format": verificationFormat, "__allBlocks": true},
func(_, data map[string]any) bool {
return blockContentEquals(data, blockID, content, rt.Str("doc-format"))
})
case "block_delete":
blockID := rt.Str("block-id")
return executeVerifiedDocMutation(rt, "doc.update", "delete_document_block", map[string]any{"nodeId": node, "blockId": blockID}, node,
"list_document_blocks", map[string]any{"nodeId": node, "format": "element"},
"list_document_blocks", map[string]any{"nodeId": node, "format": "element", "__allBlocks": true},
func(_, data map[string]any) bool { return findBlock(data, blockID) == nil })
case "str_replace":
return executePlainTextReplace(rt, node)
@@ -609,7 +635,7 @@ func nestedRevision(value any) (int, bool) {
}
func executePlainTextReplace(rt *shortcut.RuntimeContext, nodeID string) error {
data, err := rt.CallMCPData(productDoc, "list_document_blocks", map[string]any{"nodeId": nodeID, "format": "element"})
data, err := readAllDocumentBlocks(rt, map[string]any{"nodeId": nodeID, "format": "element"})
if err != nil {
return err
}
@@ -643,12 +669,12 @@ func executePlainTextReplace(rt *shortcut.RuntimeContext, nodeID string) error {
blockID := matches[0].blockID
return executeVerifiedDocMutation(rt, "doc.update", "update_document_block",
map[string]any{"nodeId": nodeID, "blockId": blockID, "element": map[string]any{"blockType": "paragraph", "paragraph": map[string]any{"text": updated}}}, nodeID,
"list_document_blocks", map[string]any{"nodeId": nodeID, "blockId": blockID, "format": "element"},
"list_document_blocks", map[string]any{"nodeId": nodeID, "format": "element", "__allBlocks": true},
func(_, data map[string]any) bool { return blockContentEquals(data, blockID, updated, "markdown") })
}
func executeBlockCopy(rt *shortcut.RuntimeContext, nodeID string) error {
data, err := rt.CallMCPData(productDoc, "list_document_blocks", map[string]any{"nodeId": nodeID, "blockId": rt.Str("block-id"), "format": "element"})
data, err := readAllDocumentBlocks(rt, map[string]any{"nodeId": nodeID, "format": "element"})
if err != nil {
return err
}
@@ -664,7 +690,7 @@ func executeBlockCopy(rt *shortcut.RuntimeContext, nodeID string) error {
referenceBlockID := rt.Str("after-block-id")
return executeVerifiedDocMutation(rt, "doc.update", "insert_document_block",
map[string]any{"nodeId": nodeID, "referenceBlockId": referenceBlockID, "where": "after", "element": block}, nodeID,
"list_document_blocks", map[string]any{"nodeId": nodeID, "format": "element"},
"list_document_blocks", map[string]any{"nodeId": nodeID, "format": "element", "__allBlocks": true},
func(result, data map[string]any) bool {
return verifyInsertedCanonicalBlock(result, data, referenceBlockID, expectedContent, "markdown")
})
@@ -684,7 +710,9 @@ func executeVerifiedDocMutation(
return docUnknownWriteError(operation, tool, nodeID, err)
}
steps[0]["status"] = "success"
verification, err := rt.CallMCPData(productDoc, verifyTool, verifyParams)
verification, err := readDocVerification(rt, verifyTool, verifyParams, func(data map[string]any) bool {
return verify == nil || verify(result, data)
})
if err != nil {
return docVerificationError(operation, "verify", nodeID, err, append(steps, map[string]any{"name": "verify", "status": "failed"}))
}
@@ -732,7 +760,9 @@ func executeVerifiedDocContentMutation(rt *shortcut.RuntimeContext, firstParams
}
steps = append(steps, map[string]any{"name": stepName, "status": "success"})
}
verification, err := rt.CallMCPData(productDoc, "get_document_content", map[string]any{"nodeId": nodeID, "format": format})
verification, err := readDocVerification(rt, "get_document_content", map[string]any{"nodeId": nodeID, "format": format}, func(data map[string]any) bool {
return verifyUpdatedDocumentContent(data, content, mode, format)
})
if err != nil {
return docVerificationError("doc.update", "verify", nodeID, err, append(steps, map[string]any{"name": "verify", "status": "failed"}))
}
@@ -745,6 +775,180 @@ func executeVerifiedDocContentMutation(rt *shortcut.RuntimeContext, firstParams
}, steps...))
}
func readDocVerification(rt *shortcut.RuntimeContext, tool string, rawParams map[string]any, verify func(map[string]any) bool) (map[string]any, error) {
params := cloneMap(rawParams)
allBlocks, _ := params["__allBlocks"].(bool)
delete(params, "__allBlocks")
var last map[string]any
var lastErr error
for attempt := 0; attempt <= len(docVerifyDelays); attempt++ {
var data map[string]any
var err error
if allBlocks && tool == "list_document_blocks" {
data, err = readAllDocumentBlocks(rt, params)
} else {
data, err = rt.CallMCPData(productDoc, tool, params)
}
if err != nil {
lastErr = err
} else {
last = data
lastErr = nil
if verify == nil || verify(data) {
return data, nil
}
}
if attempt < len(docVerifyDelays) {
if err := docVerifyWait(rt.Command().Context(), docVerifyDelays[attempt]); err != nil {
return nil, err
}
}
}
if lastErr != nil {
return nil, lastErr
}
return last, nil
}
func waitForDocVerification(ctx context.Context, delay time.Duration) error {
if ctx == nil {
ctx = context.Background()
}
timer := time.NewTimer(delay)
defer timer.Stop()
select {
case <-ctx.Done():
return ctx.Err()
case <-timer.C:
return nil
}
}
func readAllDocumentBlocks(rt *shortcut.RuntimeContext, base map[string]any) (map[string]any, error) {
all := make([]any, 0, docBlockReadPageSize)
seenPageIdentities := map[string]bool{}
for start := 0; start < docBlockReadMaxItems; start += docBlockReadPageSize {
params := cloneMap(base)
params["startIndex"] = start
params["endIndex"] = start + docBlockReadPageSize - 1
page, err := rt.CallMCPData(productDoc, "list_document_blocks", params)
if err != nil {
return nil, err
}
blocks, ok := documentBlockEntries(page)
if !ok {
return nil, fmt.Errorf("list_document_blocks 回读缺少 blocks 数组")
}
pageIdentity := documentBlockPageIdentity(blocks)
if pageIdentity != "" && seenPageIdentities[pageIdentity] {
return nil, fmt.Errorf("list_document_blocks 分页停滞,无法证明回读完整")
}
if pageIdentity != "" {
seenPageIdentities[pageIdentity] = true
}
all = append(all, blocks...)
hasMore, known, _ := docPageState(page)
if known && !hasMore {
return map[string]any{"blocks": all, "hasMore": false, "totalCount": len(all)}, nil
}
if !known {
if total, ok := nestedNonNegativeInt(page, "totalCount", "total_count"); ok && len(all) >= total {
return map[string]any{"blocks": all, "hasMore": false, "totalCount": total}, nil
}
}
if !known && len(blocks) < docBlockReadPageSize {
return map[string]any{"blocks": all, "hasMore": false, "totalCount": len(all)}, nil
}
if len(blocks) == 0 {
return nil, fmt.Errorf("list_document_blocks 声明仍有下一页但当前页为空,无法证明回读完整")
}
}
return nil, fmt.Errorf("list_document_blocks 超过 %d 个块,无法在安全上限内完成回读", docBlockReadMaxItems)
}
func documentBlockPageIdentity(blocks []any) string {
if len(blocks) == 0 {
return ""
}
ids := make([]string, 0, len(blocks))
for _, value := range blocks {
id := ""
switch block := value.(type) {
case map[string]any:
id = blockIdentity(block, "")
if id == "" {
if element, ok := block["element"].(map[string]any); ok {
id = blockIdentity(element, "")
}
}
case []any:
id = jsonMLBlockIdentity(block)
}
if id == "" {
return ""
}
ids = append(ids, id)
}
encoded, _ := json.Marshal(ids)
return string(encoded)
}
func documentBlockEntries(value any) ([]any, bool) {
switch typed := value.(type) {
case map[string]any:
for _, key := range []string{"blocks", "items"} {
if blocks, ok := typed[key].([]any); ok {
return blocks, true
}
}
if encoded, ok := typed["jsonml"].(string); ok {
var decoded any
if json.Unmarshal([]byte(encoded), &decoded) == nil {
blocks := orderedJSONMLBlocks(decoded)
values := make([]any, len(blocks))
for index := range blocks {
values[index] = blocks[index]
}
return values, true
}
}
for _, key := range []string{"result", "data"} {
if nested, ok := typed[key]; ok {
if blocks, found := documentBlockEntries(nested); found {
return blocks, true
}
}
}
}
return nil, false
}
func nestedNonNegativeInt(value any, keys ...string) (int, bool) {
switch typed := value.(type) {
case map[string]any:
for _, key := range keys {
if raw, ok := typed[key]; ok {
switch number := raw.(type) {
case float64:
if number >= 0 && number == float64(int(number)) {
return int(number), true
}
case int:
if number >= 0 {
return number, true
}
}
}
}
for _, key := range []string{"result", "data"} {
if result, ok := nestedNonNegativeInt(typed[key], keys...); ok {
return result, true
}
}
}
return 0, false
}
func splitDocMarkdown(content string, maxRunes int) []string {
if maxRunes <= 0 {
return []string{content}
@@ -796,11 +1000,17 @@ func containsText(value any, needle string) bool {
}
func verifyUpdatedDocumentContent(value any, expected, mode, format string) bool {
expected = normalizeDocumentContentForVerification(expected, format)
expectedRaw := expected
expected = normalizeDocumentContentForVerification(expectedRaw, format)
for _, candidate := range documentContentCandidates(value, format) {
actual := normalizeDocumentContentForVerification(candidate, format)
actualRaw := candidate
actual := normalizeDocumentContentForVerification(actualRaw, format)
if mode == "overwrite" {
if actual == expected {
if actual == expected || (format == "markdown" && stripReadbackDocumentTitle(actual) == expected) {
return true
}
if format == "markdown" && (markdownSemanticallyEquivalent(actualRaw, expectedRaw) ||
markdownSemanticallyEquivalent(stripReadbackDocumentTitle(actualRaw), expectedRaw)) {
return true
}
continue
@@ -808,10 +1018,48 @@ func verifyUpdatedDocumentContent(value any, expected, mode, format string) bool
if actual == expected || strings.HasSuffix(actual, "\n"+expected) {
return true
}
if format == "markdown" && markdownSemanticallyEndsWith(actualRaw, expectedRaw) {
return true
}
}
return false
}
func markdownSemanticallyEquivalent(left, right string) bool {
leftFingerprint, leftOK := markdownSemanticFingerprint(left)
rightFingerprint, rightOK := markdownSemanticFingerprint(right)
return leftOK && rightOK && leftFingerprint == rightFingerprint
}
func markdownSemanticallyEndsWith(content, suffix string) bool {
contentFingerprint, contentOK := markdownSemanticFingerprint(content)
suffixFingerprint, suffixOK := markdownSemanticFingerprint(suffix)
return contentOK && suffixOK && strings.HasSuffix(contentFingerprint, suffixFingerprint)
}
func markdownSemanticFingerprint(source string) (string, bool) {
if len(source) > docMarkdownVerifyMax {
return "", false
}
var rendered bytes.Buffer
if err := docMarkdownConvert([]byte(source), &rendered); err != nil {
return "", false
}
return rendered.String(), true
}
func stripReadbackDocumentTitle(content string) string {
lines := strings.Split(content, "\n")
if len(lines) == 0 || !strings.HasPrefix(strings.TrimSpace(lines[0]), "# ") {
return content
}
lines = lines[1:]
for len(lines) > 0 && strings.TrimSpace(lines[0]) == "" {
lines = lines[1:]
}
return strings.Join(lines, "\n")
}
func verifyInsertedBlock(result, data map[string]any, referenceBlockID, expected, format string) bool {
return verifyInsertedCanonicalBlock(result, data, referenceBlockID, normalizeDocumentContentForVerification(expected, format), format)
}
@@ -848,6 +1096,11 @@ func blockContentEquals(data map[string]any, blockID, expected, format string) b
}
func canonicalBlockContent(value any, format string) string {
if values, ok := value.(map[string]any); ok {
if element, ok := values["element"].(map[string]any); ok {
value = element
}
}
if format == "jsonml" {
if values, ok := value.(map[string]any); ok {
if encoded, ok := values["jsonml"].(string); ok {
@@ -871,7 +1124,7 @@ func canonicalBlockContent(value any, format string) string {
return
}
for key, child := range typed {
if key == "id" || key == "blockId" || key == "uuid" {
if key == "id" || key == "blockId" || key == "uuid" || key == "blockType" {
continue
}
walk(child)
@@ -1005,6 +1258,10 @@ func orderedDocumentBlocks(value any) []map[string]any {
walk = func(current any) {
switch typed := current.(type) {
case map[string]any:
if element, ok := typed["element"].(map[string]any); ok && blockIdentity(element, "") != "" {
blocks = append(blocks, element)
return
}
if blockIdentity(typed, "") != "" {
blocks = append(blocks, typed)
return
@@ -1065,9 +1322,23 @@ func normalizeDocumentContentForVerification(raw, format string) string {
func normalizeMarkdownForVerification(raw string) string {
raw = strings.ReplaceAll(strings.ReplaceAll(raw, "\r\n", "\n"), "\r", "\n")
lines := make([]string, 0, strings.Count(raw, "\n")+1)
inFence := false
for _, line := range strings.Split(raw, "\n") {
line = strings.TrimSpace(line)
trimmed := strings.TrimSpace(line)
if strings.HasPrefix(trimmed, "```") || strings.HasPrefix(trimmed, "~~~") {
inFence = !inFence
lines = append(lines, trimmed)
continue
}
if inFence {
lines = append(lines, strings.TrimRight(line, " \t"))
continue
}
line = trimmed
if line == "" {
if len(lines) > 0 && lines[len(lines)-1] != "" {
lines = append(lines, "")
}
continue
}
if strings.Contains(line, "|") {
@@ -1081,6 +1352,9 @@ func normalizeMarkdownForVerification(raw string) string {
}
lines = append(lines, line)
}
for len(lines) > 0 && lines[len(lines)-1] == "" {
lines = lines[:len(lines)-1]
}
return strings.Join(lines, "\n")
}
@@ -1089,33 +1363,108 @@ func normalizeJSONMLForVerification(raw string) string {
if err := json.Unmarshal([]byte(raw), &value); err != nil {
return normalizeMarkdownForVerification(raw)
}
var tokens []string
var walk func(any)
walk = func(current any) {
var normalize func(any) any
normalize = func(current any) any {
switch typed := current.(type) {
case []any:
start := 0
if len(typed) > 0 {
if tag, ok := typed[0].(string); ok {
tokens = append(tokens, "<"+strings.ToLower(strings.TrimSpace(tag))+">")
start = 1
if len(typed) == 0 {
return []any{}
}
tag, isElement := typed[0].(string)
if !isElement {
out := make([]any, 0, len(typed))
for _, child := range typed {
out = append(out, normalize(child))
}
return out
}
start := 1
attrs := map[string]any{}
if len(typed) > 1 {
if declared, ok := typed[1].(map[string]any); ok {
attrs, _ = normalize(declared).(map[string]any)
attrs = removeGeneratedJSONMLDefaults(tag, attrs)
start = 2
}
}
children := make([]any, 0, len(typed)-start)
for _, child := range typed[start:] {
walk(child)
normalized := normalize(child)
if normalized != nil {
children = append(children, normalized)
}
}
if strings.EqualFold(tag, "span") && isGeneratedTextSpan(attrs) {
if len(children) == 1 {
return children[0]
}
return children
}
out := []any{strings.ToLower(tag), attrs}
out = append(out, children...)
return out
case map[string]any:
// JSONML maps contain element attributes. Server-generated UUIDs and
// default attributes do not change the authored document content.
return
case string:
if text := normalizeMarkdownForVerification(typed); text != "" {
tokens = append(tokens, text)
out := make(map[string]any, len(typed))
for key, child := range typed {
normalizedKey := strings.ToLower(strings.NewReplacer("_", "", "-", "").Replace(key))
if normalizedKey == "uuid" || normalizedKey == "blockid" || normalizedKey == "elementid" || normalizedKey == "index" {
continue
}
out[normalizedKey] = normalize(child)
}
return out
case string:
return strings.ReplaceAll(strings.ReplaceAll(typed, "\r\n", "\n"), "\r", "\n")
}
return current
}
walk(value)
return strings.Join(tokens, "\n")
// normalize only receives values decoded by encoding/json, so the resulting
// tree is always JSON-marshalable.
encoded, _ := json.Marshal(normalize(value))
return string(encoded)
}
var generatedJSONMLAttributeDefaults = map[string]map[string]any{
"hr": {
"sz": float64(1),
},
"tc": {
"colspan": float64(1), "rowspan": float64(1), "valign": "middle",
},
"code": {
"code": "", "syntax": "plaintext", "theme": "default",
"wrap": true, "showlinenumber": true, "fold": false,
},
}
// removeGeneratedJSONMLDefaults drops only defaults declared by the reviewed
// JSONML schema, plus empty server style objects. Other attributes remain part
// of the semantic fingerprint so links, formatting, and table layout stay
// strict.
func removeGeneratedJSONMLDefaults(tag string, attrs map[string]any) map[string]any {
defaults := generatedJSONMLAttributeDefaults[strings.ToLower(tag)]
out := make(map[string]any, len(attrs))
for key, value := range attrs {
if object, ok := value.(map[string]any); ok && len(object) == 0 {
continue
}
if defaultValue, ok := defaults[key]; ok && reflect.DeepEqual(value, defaultValue) {
continue
}
out[key] = value
}
return out
}
func isGeneratedTextSpan(attrs map[string]any) bool {
if len(attrs) == 0 {
return true
}
if len(attrs) != 1 {
return false
}
value, ok := attrs["datatype"].(string)
return ok && (value == "text" || value == "leaf")
}
func executeExport(rt *shortcut.RuntimeContext) error {
@@ -0,0 +1,401 @@
// Copyright 2026 Alibaba Group
// Licensed under the Apache License, Version 2.0
package doc
import (
"context"
"encoding/json"
"errors"
"fmt"
"io"
"strings"
"testing"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/helpers"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/shortcut"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/testseam"
"github.com/spf13/cobra"
)
func TestCrossPlatformCoverageDocReadbackRetriesStaleContent(t *testing.T) {
testseam.Swap(t, &docVerifyWait, func(context.Context, time.Duration) error { return nil })
testseam.Swap(t, &docVerifyDelays, []time.Duration{time.Millisecond})
caller := &docCoverageCaller{responses: map[string][]map[string]any{
"get_document_content": {{"markdown": "old"}, {"markdown": "old\nnew"}},
}}
if err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "append", "--content", "new", "--yes"); err != nil {
t.Fatal(err)
}
reads := 0
for _, call := range caller.history {
if call.tool == "get_document_content" {
reads++
}
}
if reads != 2 {
t.Fatalf("readback calls = %d, want 2; history=%#v", reads, caller.history)
}
}
func TestCrossPlatformCoverageDocReadbackStopsOnCancellation(t *testing.T) {
if err := waitForDocVerification(nil, time.Nanosecond); err != nil {
t.Fatalf("completed verification wait = %v", err)
}
cancelled, cancel := context.WithCancel(context.Background())
cancel()
if err := waitForDocVerification(cancelled, time.Hour); !errors.Is(err, context.Canceled) {
t.Fatalf("cancelled verification wait = %v, want context.Canceled", err)
}
cmd := &cobra.Command{Use: "verify"}
cmd.SetContext(cancelled)
caller := &docCoverageCaller{responses: map[string][]map[string]any{
"get_document_content": {{"markdown": "stale"}},
}}
helpers.InitDeps(caller)
rt := shortcut.RuntimeContextForTest(cmd, Update)
_, err := readDocVerification(rt, "get_document_content", map[string]any{"nodeId": "n"}, func(map[string]any) bool { return false })
if !errors.Is(err, context.Canceled) {
t.Fatalf("cancelled document verification = %v, want context.Canceled", err)
}
if len(caller.history) != 1 {
t.Fatalf("cancelled document verification calls = %d, want 1", len(caller.history))
}
}
func TestCrossPlatformCoverageDocDeleteReadbackConsumesEveryPage(t *testing.T) {
testseam.Swap(t, &docVerifyWait, func(context.Context, time.Duration) error { return nil })
testseam.Swap(t, &docVerifyDelays, []time.Duration{time.Millisecond})
firstPage := make([]any, 50)
for index := range firstPage {
firstPage[index] = map[string]any{"id": fmt.Sprintf("block-%d", index), "text": "body"}
}
caller := &docCoverageCaller{responses: map[string][]map[string]any{
"list_document_blocks": {
{"blocks": firstPage, "hasMore": true, "totalCount": 51},
{"blocks": []any{map[string]any{"id": "target", "text": "stale"}}, "hasMore": false, "totalCount": 51},
{"blocks": firstPage, "hasMore": false, "totalCount": 50},
},
}}
if err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "block_delete", "--block-id", "target", "--yes"); err != nil {
t.Fatal(err)
}
starts := []int{}
for _, call := range caller.history {
if call.tool == "list_document_blocks" {
starts = append(starts, call.params["startIndex"].(int))
}
}
if fmt.Sprint(starts) != "[0 50 0]" {
t.Fatalf("pagination starts = %v, want [0 50 0]", starts)
}
}
func TestCrossPlatformCoverageDocReplacePreflightIsGloballyUnique(t *testing.T) {
firstPage := make([]any, 50)
for index := range firstPage {
text := "body"
if index == 0 {
text = "unique needle"
}
firstPage[index] = map[string]any{"id": fmt.Sprintf("block-%d", index), "text": text}
}
caller := &docCoverageCaller{responses: map[string][]map[string]any{
"list_document_blocks": {
{"blocks": firstPage, "hasMore": true, "totalCount": 51},
{"blocks": []any{map[string]any{"id": "block-50", "text": "another needle"}}, "hasMore": false, "totalCount": 51},
},
}}
err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "str_replace", "--old", "needle", "--new", "changed", "--yes")
if err == nil {
t.Fatal("str_replace accepted a second match on a later page")
}
for _, call := range caller.history {
if call.tool == "update_document_block" {
t.Fatalf("ambiguous replace executed a write: %#v", caller.history)
}
}
}
func TestCrossPlatformCoverageDocVerificationPreservesMeaning(t *testing.T) {
expected := "[\"p\",{},\"text\"]"
serverExpanded := "[\"p\",{\"uuid\":\"generated\",\"style\":{}},[\"span\",{\"data-type\":\"text\"},[\"span\",{\"data-type\":\"leaf\"},\"text\"]]]"
if normalizeJSONMLForVerification(expected) != normalizeJSONMLForVerification(serverExpanded) {
t.Fatal("generated JSONML text wrappers should not change document meaning")
}
defaultFree := `[["hr",{}],["code",{}],["table",{},["tr",{},["tc",{},["p",{},"cell"]]]]]`
serverDefaulted := `[["hr",{"sz":1}],["code",{"code":"","syntax":"plaintext","theme":"default","wrap":true,"showLineNumber":true,"fold":false}],["table",{},["tr",{},["tc",{"colSpan":1,"rowSpan":1,"vAlign":"middle"},["p",{},"cell"]]]]]`
if !verifyUpdatedDocumentContent(map[string]any{"jsonml": serverDefaulted}, defaultFree, "overwrite", "jsonml") {
t.Fatal("server-generated JSONML schema defaults should not fail readback verification")
}
linkA := "[\"a\",{\"href\":\"https://example.com/a\"},\"text\"]"
linkB := "[\"a\",{\"href\":\"https://example.com/b\"},\"text\"]"
if normalizeJSONMLForVerification(linkA) == normalizeJSONMLForVerification(linkB) {
t.Fatal("semantic JSONML attributes were ignored")
}
tableA := `[["table",{"jc":"center"},["tr",{},["tc",{},["p",{},"cell"]]]]]`
tableB := `[["table",{"jc":"right"},["tr",{},["tc",{},["p",{},"cell"]]]]]`
if normalizeJSONMLForVerification(tableA) == normalizeJSONMLForVerification(tableB) {
t.Fatal("semantic JSONML table alignment was ignored")
}
codeA := "~~~go\n return nil\n~~~"
codeB := "~~~go\nreturn nil\n~~~"
if normalizeMarkdownForVerification(codeA) == normalizeMarkdownForVerification(codeB) {
t.Fatal("fenced code indentation was ignored")
}
if !verifyUpdatedDocumentContent(map[string]any{"markdown": "# Server title\n\nbody"}, "body", "overwrite", "markdown") {
t.Fatal("server-generated document title prevented body verification")
}
}
func TestCrossPlatformCoverageMarkdownSemanticRoundTrip(t *testing.T) {
input := strings.Join([]string{
"sales_data.xlsx",
"### 1.",
"+10.22%",
"| name | value |",
"| -------- | -------- |",
"| sales_data.xlsx | +10.22% |",
}, "\n")
server := strings.Join([]string{
`sales\_data.xlsx`,
`### 1\.`,
`\+10.22%`,
"|name|value|",
"|---|---|",
`|sales\_data.xlsx|\+10.22%|`,
}, "\n")
if !markdownSemanticallyEquivalent(input, server) {
t.Fatal("server Markdown escaping and table delimiter normalization changed the semantic fingerprint")
}
if !verifyUpdatedDocumentContent(map[string]any{"markdown": server}, input, "overwrite", "markdown") {
t.Fatal("equivalent server Markdown failed overwrite verification")
}
if !verifyUpdatedDocumentContent(map[string]any{"markdown": "existing\n\n" + server}, input, "append", "markdown") {
t.Fatal("equivalent server Markdown failed append verification")
}
}
func TestCrossPlatformCoverageMarkdownSemanticDifferencesRemainStrict(t *testing.T) {
for _, test := range []struct {
name string
left string
right string
}{
{name: "emphasis", left: `*important*`, right: `\*important\*`},
{name: "inline code", left: "`sales_data`", right: "`sales\\_data`"},
{name: "fenced code", left: "```\nsales_data\n```", right: "```\nsales\\_data\n```"},
{name: "table alignment", left: "|a|\n|---|\n|x|", right: "|a|\n|:---|\n|x|"},
{name: "table columns", left: "|a|b|\n|---|---|\n|x|y|", right: "|a|\n|---|\n|x|"},
{name: "table content", left: "|a|\n|---|\n|x|", right: "|a|\n|---|\n|y|"},
} {
t.Run(test.name, func(t *testing.T) {
if markdownSemanticallyEquivalent(test.left, test.right) {
t.Fatal("meaningful Markdown difference was ignored")
}
})
}
oversized := strings.Repeat("x", docMarkdownVerifyMax+1)
if _, ok := markdownSemanticFingerprint(oversized); ok {
t.Fatal("oversized Markdown entered semantic verification")
}
testseam.Swap(t, &docMarkdownConvert, func([]byte, io.Writer) error { return errors.New("render") })
if _, ok := markdownSemanticFingerprint("body"); ok {
t.Fatal("failed Markdown render produced a semantic fingerprint")
}
}
func TestCrossPlatformCoverageDocElementReadbackUsesNestedElement(t *testing.T) {
wrapper := map[string]any{
"blockType": "paragraph",
"element": map[string]any{"id": "inserted", "blockType": "paragraph", "paragraph": map[string]any{"text": "body"}},
}
if got := canonicalBlockContent(wrapper, "markdown"); got != "body" {
t.Fatalf("nested element content = %q, want body", got)
}
blocks := orderedDocumentBlocks(map[string]any{"blocks": []any{wrapper}})
if len(blocks) != 1 || blockIdentity(blocks[0], "") != "inserted" {
t.Fatalf("nested element blocks = %#v", blocks)
}
}
func TestCrossPlatformCoverageVersionRevertRequiresTargetEvidence(t *testing.T) {
if revertResultMatchesVersion(map[string]any{"ok": true}, 3) || currentDocumentMatchesRestoredVersion(map[string]any{"version": 99}, 3) {
t.Fatal("readability or an unrelated current version must not prove a revert")
}
if revertResultMatchesVersion(map[string]any{"version": 3}, 3) {
t.Fatal("the request version parameter must not prove its own revert")
}
if revertResultMatchesVersion(map[string]any{"data": map[string]any{"request": map[string]any{"targetVersion": 3}}}, 3) {
t.Fatal("request echo containers must not provide version evidence")
}
for _, failed := range []map[string]any{
{"data": map[string]any{"success": "false", "revertedToVersion": 3}},
{"data": map[string]any{"ok": false, "revertedToVersion": 3}},
{"data": map[string]any{"status": "FAILED", "revertedToVersion": 3}},
{"data": map[string]any{"state": "failure", "revertedToVersion": 3}},
{"data": []any{map[string]any{"error_code": "REVERT_FAILED", "revertedToVersion": 3}}},
{"data": map[string]any{"errorCode": 500.0, "revertedToVersion": 3}},
{"data": map[string]any{"code": json.Number("500"), "revertedToVersion": 3}},
} {
if revertResultMatchesVersion(failed, 3) {
t.Fatalf("explicit failure %#v must override target-version evidence", failed)
}
}
for _, succeeded := range []map[string]any{
{"revertedToVersion": 3},
{"status": "SUCCESS", "revertedToVersion": 3},
{"state": "succeeded", "revertedToVersion": 3},
{"errorCode": "0", "revertedToVersion": 3},
{"code": 200, "revertedToVersion": 3},
{"code": "204", "revertedToVersion": 3},
} {
if !revertResultMatchesVersion(succeeded, 3) {
t.Fatalf("explicit success %#v suppressed target-version evidence", succeeded)
}
}
for _, absentOrSuccess := range []any{nil, "", "OK", "SUCCESS", json.Number("0"), json.Number("200"), 0.0, 201.0, 0, 202} {
if revertErrorCodeIsFailure(absentOrSuccess) {
t.Fatalf("success code %#v was treated as an explicit failure", absentOrSuccess)
}
}
for _, failedCode := range []any{json.Number("bad"), json.Number("500"), 1.0, 500.0, 1, 500} {
if !revertErrorCodeIsFailure(failedCode) {
t.Fatalf("failure code %#v was not treated as an explicit failure", failedCode)
}
}
if revertStatusIsFailure(500) || revertStatusIsFailure("PROCESSING") {
t.Fatal("non-failure status was treated as an explicit failure")
}
}
func TestCrossPlatformCoverageDocReadbackDefensiveEdges(t *testing.T) {
testseam.Swap(t, &docVerifyWait, func(context.Context, time.Duration) error { return nil })
for _, tc := range []struct {
name string
responses []map[string]any
failAt int
}{
{"call failure", nil, 1},
{"missing blocks", []map[string]any{{"ok": true}}, 0},
{"stalled page", []map[string]any{{"blocks": []any{map[string]any{"id": "a"}}, "hasMore": true}, {"blocks": []any{map[string]any{"id": "a"}}, "hasMore": true}}, 0},
{"empty continued page", []map[string]any{{"blocks": []any{}, "hasMore": true}}, 0},
} {
t.Run(tc.name, func(t *testing.T) {
testseam.Swap(t, &docVerifyDelays, []time.Duration{})
caller := &docCoverageCaller{failAt: tc.failAt, responses: map[string][]map[string]any{"list_document_blocks": tc.responses}}
err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "block_delete", "--block-id", "target", "--yes")
if err == nil {
t.Fatal("defensive readback unexpectedly succeeded")
}
})
}
t.Run("identical adjacent pages advance by requested indexes", func(t *testing.T) {
testseam.Swap(t, &docVerifyDelays, []time.Duration{})
blocks := make([]any, docBlockReadPageSize)
for index := range blocks {
blocks[index] = map[string]any{"blockType": "paragraph"}
}
caller := &docCoverageCaller{responses: map[string][]map[string]any{"list_document_blocks": {
{"blocks": blocks, "hasMore": true, "totalCount": 2 * docBlockReadPageSize},
{"blocks": blocks, "hasMore": false, "totalCount": 2 * docBlockReadPageSize},
}}}
if err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "block_delete", "--block-id", "target", "--yes"); err != nil {
t.Fatal(err)
}
starts := []int{}
for _, call := range caller.history {
if call.tool == "list_document_blocks" {
starts = append(starts, call.params["startIndex"].(int))
}
}
if fmt.Sprint(starts) != "[0 50]" {
t.Fatalf("pagination starts = %v, want [0 50]", starts)
}
})
t.Run("total count terminates pagination", func(t *testing.T) {
testseam.Swap(t, &docVerifyDelays, []time.Duration{})
caller := &docCoverageCaller{responses: map[string][]map[string]any{"list_document_blocks": {{"blocks": []any{map[string]any{"id": "other"}}, "totalCount": 1}}}}
if err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "block_delete", "--block-id", "target", "--yes"); err != nil {
t.Fatal(err)
}
})
t.Run("explicit has more overrides inconsistent total count", func(t *testing.T) {
testseam.Swap(t, &docVerifyDelays, []time.Duration{})
caller := &docCoverageCaller{responses: map[string][]map[string]any{"list_document_blocks": {
{"blocks": []any{map[string]any{"id": "first"}}, "hasMore": true, "totalCount": 1},
{"blocks": []any{map[string]any{"id": "second"}}, "hasMore": false, "totalCount": 2},
}}}
if err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "block_delete", "--block-id", "target", "--yes"); err != nil {
t.Fatal(err)
}
reads := 0
for _, call := range caller.history {
if call.tool == "list_document_blocks" {
reads++
}
}
if reads != 2 {
t.Fatalf("read calls = %d, want both explicitly advertised pages", reads)
}
})
t.Run("block read safety limit", func(t *testing.T) {
testseam.Swap(t, &docVerifyDelays, []time.Duration{})
pages := make([]map[string]any, docBlockReadMaxItems/docBlockReadPageSize)
for index := range pages {
pages[index] = map[string]any{"blocks": []any{map[string]any{"id": fmt.Sprintf("block-%d", index)}}, "hasMore": true}
}
caller := &docCoverageCaller{responses: map[string][]map[string]any{"list_document_blocks": pages}}
if err := runDocCoverage(t, Update, caller, "--node", "n", "--command", "block_delete", "--block-id", "target", "--yes"); err == nil {
t.Fatal("oversized block read returned nil")
}
})
if blocks, ok := documentBlockEntries(map[string]any{"jsonml": `["root",{},["p",{"uuid":"a"},"x"]]`}); !ok || len(blocks) == 0 {
t.Fatalf("jsonml blocks=%#v ok=%v", blocks, ok)
}
if _, ok := documentBlockEntries(map[string]any{"jsonml": `{`}); ok {
t.Fatal("invalid jsonml produced blocks")
}
if blocks, ok := documentBlockEntries(map[string]any{"data": map[string]any{"items": []any{"x"}}}); !ok || len(blocks) != 1 {
t.Fatalf("nested items=%#v ok=%v", blocks, ok)
}
if _, ok := documentBlockEntries(nil); ok {
t.Fatal("nil produced blocks")
}
for _, tc := range []struct {
value any
want bool
}{
{map[string]any{"totalCount": float64(2)}, true},
{map[string]any{"totalCount": float64(-1)}, false},
{map[string]any{"totalCount": 2.5}, false},
{map[string]any{"data": map[string]any{"total_count": 2}}, true},
{map[string]any{"totalCount": -1}, false},
{nil, false},
} {
_, ok := nestedNonNegativeInt(tc.value, "totalCount", "total_count")
if ok != tc.want {
t.Fatalf("nestedNonNegativeInt(%#v) ok=%v want=%v", tc.value, ok, tc.want)
}
}
for _, raw := range []string{
`[]`, `[1,["p",{},"x"]]`, `["span",{},"a","b"]`,
`["p",{"block_id":"x","custom":true},"x"]`, `true`,
} {
if normalizeJSONMLForVerification(raw) == "" {
t.Fatalf("empty normalized JSONML for %s", raw)
}
}
if !isGeneratedTextSpan(nil) || isGeneratedTextSpan(map[string]any{"a": 1, "b": 2}) || isGeneratedTextSpan(map[string]any{"data-type": 3}) || isGeneratedTextSpan(map[string]any{"data-type": "other"}) {
t.Fatal("generated span classification failed")
}
}
+55 -3
View File
@@ -15,6 +15,7 @@ import (
"strings"
"sync"
"testing"
"time"
"unicode/utf8"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd"
@@ -57,7 +58,7 @@ func (f *docCoverageCaller) CallTool(_ context.Context, _, tool string, params m
}
value := docCoveragePayload(tool)
if tool == "revert_doc_version" {
value = map[string]any{"version": params["version"]}
value = map[string]any{"revertedToVersion": params["version"]}
}
if queue := f.responses[tool]; len(queue) > 0 {
value = queue[0]
@@ -112,6 +113,7 @@ func runDocCoverageInput(t *testing.T, declaration shortcut.Shortcut, caller *do
func runDocCoveragePath(t *testing.T, declaration shortcut.Shortcut, caller *docCoverageCaller, input io.Reader, commandPath string, args ...string) error {
t.Helper()
testseam.Swap(t, &docVerifyWait, func(context.Context, time.Duration) error { return nil })
helpers.InitDeps(caller)
root := &cobra.Command{Use: "dws", SilenceErrors: true, SilenceUsage: true}
root.PersistentFlags().Bool("yes", false, "")
@@ -491,6 +493,7 @@ func TestCrossPlatformCoverageDocUpdateAliasReachesNestedBranches(t *testing.T)
}
func TestCrossPlatformCoverageDocWritesStopOnUnknownCommitAndRequireVerification(t *testing.T) {
testseam.Swap(t, &docVerifyDelays, []time.Duration{})
unknown := &docCoverageCaller{failAt: 1, responses: map[string][]map[string]any{}}
err := runDocCoverage(t, Update, unknown, "--node", "n", "--command", "append", "--content", "x", "--yes")
var typed *apperrors.Error
@@ -512,6 +515,7 @@ func TestCrossPlatformCoverageDocWritesStopOnUnknownCommitAndRequireVerification
}
func TestCrossPlatformCoverageDocCreateRejectsSuccessfulMismatchedReadback(t *testing.T) {
testseam.Swap(t, &docVerifyDelays, []time.Duration{})
caller := &docCoverageCaller{responses: map[string][]map[string]any{
"get_document_content": {{"markdown": "truncated"}},
}}
@@ -577,11 +581,59 @@ func TestCrossPlatformCoverageDocVersionRevertPaginationAndVerification(t *testi
"revert_doc_version": {{}},
"get_document_info": {{"nodeId": "n", "revision": 99.0}},
}}
if err := runDocCoverage(t, VersionRevert, caller, "--node", "n", "--version", "3", "--yes"); err != nil {
t.Fatal(err)
err := runDocCoverage(t, VersionRevert, caller, "--node", "n", "--version", "3", "--yes")
var typed *apperrors.Error
if !errors.As(err, &typed) || typed.Reason != "doc_history_revert_target_unproven" || typed.FailureStage != "verify" || typed.Details["status"] != "partial_success" {
t.Fatalf("unproven revert error = %#v", err)
}
data, _ := typed.Details["data"].(map[string]any)
steps, _ := typed.Details["steps"].([]map[string]any)
if data["verified"] != false || len(steps) != 3 || steps[2]["status"] != "failed" {
t.Fatalf("unproven revert details = %#v", typed.Details)
}
})
for _, test := range []struct {
name string
response map[string]any
current map[string]any
wantOK bool
}{
{name: "explicit server evidence", response: map[string]any{"data": map[string]any{"revertResult": map[string]any{"revertedToVersion": 3}}}, wantOK: true},
{name: "bare request parameter echo", response: map[string]any{"version": 3}},
{name: "nested request parameter echo", response: map[string]any{"data": map[string]any{"request": map[string]any{"nodeId": "n", "version": 3}}}},
{name: "accepted field inside request echo", response: map[string]any{"data": map[string]any{"request": map[string]any{"targetVersion": 3}}}},
{name: "nested business failure with request echo", response: map[string]any{"data": map[string]any{"success": false, "errorCode": "REVERT_FAILED", "request": map[string]any{"version": 3}}}},
{name: "failed status overrides target evidence", response: map[string]any{"data": map[string]any{"status": "FAILED", "revertResult": map[string]any{"revertedToVersion": 3}}}},
{name: "failed code overrides target evidence", response: map[string]any{"data": map[string]any{"code": "REVERT_FAILED", "revertResult": map[string]any{"revertedToVersion": 3}}}},
{
name: "nested business failure overrides all evidence",
response: map[string]any{"data": map[string]any{"state": "FAILURE", "revertResult": map[string]any{"revertedToVersion": 3}}},
current: map[string]any{"targetVersion": 3},
},
} {
t.Run(test.name, func(t *testing.T) {
responses := map[string][]map[string]any{
"revert_doc_version": {test.response},
}
if test.current != nil {
responses["get_document_info"] = []map[string]any{test.current}
}
caller := &docCoverageCaller{responses: responses}
err := runDocCoverage(t, VersionRevert, caller, "--node", "n", "--version", "3", "--yes")
if test.wantOK {
if err != nil {
t.Fatal(err)
}
return
}
var typed *apperrors.Error
if !errors.As(err, &typed) || typed.Reason != "doc_history_revert_target_unproven" || typed.Details["status"] != "partial_success" {
t.Fatalf("request echo response %#v produced %#v", test.response, err)
}
})
}
for _, test := range []struct {
name string
responses []map[string]any
@@ -178,15 +178,51 @@ func executeHistoryRevert(rt *shortcut.RuntimeContext) error {
map[string]any{"available": false, "reason": "the requested revert completed; verify the current document before any further write"},
)
}
verified := !revertResponseHasExplicitFailure(reverted) &&
(revertResultMatchesVersion(reverted, target) || currentDocumentMatchesRestoredVersion(current, target))
if !verified {
return docPartialWriteError(
"doc.history_revert", "doc_history_revert_target_unproven", "verify",
fmt.Sprintf("版本 %d 的回滚请求已执行且文档可读,但响应没有提供目标版本证据;不要直接重试回滚", target),
fmt.Errorf("回读缺少目标版本 %d 的明确证据", target),
map[string]any{
"nodeId": nodeID, "version": target, "reverted": true, "verified": false,
"revertResult": reverted, "current": current,
},
[]map[string]any{
{"name": "preflight", "status": "success"},
{"name": "revert", "status": "success"},
{"name": "verify", "status": "failed"},
},
map[string]any{"available": false, "reason": "the revert may have completed; inspect version history before any further revert"},
)
}
return rt.Output(docEnvelope("doc.history_revert", map[string]any{
"version": target, "revertResult": reverted, "current": current, "verified": true,
"verification": "revert_acknowledged_and_document_readable",
"verification": "target_version_proven",
},
map[string]any{"name": "preflight", "status": "success"},
map[string]any{"name": "revert", "status": "success"},
map[string]any{"name": "verify", "status": "success"}))
}
func revertResultMatchesVersion(value map[string]any, target int) bool {
if revertResponseHasExplicitFailure(value) {
return false
}
return versionEvidenceMatches(value, target, map[string]bool{
"targetversion": true, "appliedversion": true,
"restoredversion": true, "revertedversion": true, "revertedtoversion": true, "sourceversion": true,
})
}
func currentDocumentMatchesRestoredVersion(value map[string]any, target int) bool {
return versionEvidenceMatches(value, target, map[string]bool{
"restoredfromversion": true, "revertedfromversion": true,
"sourceversion": true, "appliedversion": true, "targetversion": true,
})
}
func findHistoryVersion(rt *shortcut.RuntimeContext, nodeID string, target int) (bool, error) {
const maxPages = 20
cursor := ""
@@ -240,6 +276,9 @@ func versionEvidenceMatches(value any, target int, acceptedKeys map[string]bool)
case map[string]any:
for key, child := range typed {
normalized := strings.ToLower(strings.NewReplacer("_", "", "-", "").Replace(key))
if versionEvidenceRequestEchoKeys[normalized] {
continue
}
if acceptedKeys[normalized] {
if versionNumberMatches(child, target) {
return true
@@ -259,6 +298,87 @@ func versionEvidenceMatches(value any, target int, acceptedKeys map[string]bool)
return false
}
var versionEvidenceRequestEchoKeys = map[string]bool{
"args": true, "arguments": true, "input": true, "inputs": true,
"params": true, "parameters": true, "request": true, "requestbody": true,
"requestparams": true, "toolargs": true, "toolarguments": true,
}
func revertResponseHasExplicitFailure(value any) bool {
switch typed := value.(type) {
case map[string]any:
for key, child := range typed {
normalized := strings.ToLower(strings.NewReplacer("_", "", "-", "").Replace(key))
if versionEvidenceRequestEchoKeys[normalized] {
continue
}
if normalized == "success" || normalized == "ok" {
if success, ok := child.(bool); ok && !success {
return true
}
if success, ok := child.(string); ok && strings.EqualFold(strings.TrimSpace(success), "false") {
return true
}
}
if (normalized == "status" || normalized == "state") && revertStatusIsFailure(child) {
return true
}
if normalized == "errorcode" || normalized == "code" {
if revertErrorCodeIsFailure(child) {
return true
}
}
if revertResponseHasExplicitFailure(child) {
return true
}
}
case []any:
for _, child := range typed {
if revertResponseHasExplicitFailure(child) {
return true
}
}
}
return false
}
func revertStatusIsFailure(value any) bool {
status, ok := value.(string)
if !ok {
return false
}
switch strings.ToLower(strings.TrimSpace(status)) {
case "fail", "failed", "failure", "error", "errored", "reject", "rejected", "deny", "denied", "cancel", "cancelled", "canceled", "abort", "aborted":
return true
default:
return false
}
}
func revertErrorCodeIsFailure(value any) bool {
switch typed := value.(type) {
case string:
normalized := strings.ToLower(strings.TrimSpace(typed))
switch normalized {
case "", "ok", "success", "succeed", "succeeded":
return false
}
if parsed, err := strconv.Atoi(normalized); err == nil {
return parsed != 0 && (parsed < 200 || parsed >= 300)
}
return true
case float64:
return typed != 0 && (typed < 200 || typed >= 300)
case json.Number:
parsed, err := typed.Int64()
return err != nil || (parsed != 0 && (parsed < 200 || parsed >= 300))
case int:
return typed != 0 && (typed < 200 || typed >= 300)
default:
return false
}
}
func versionNumberMatches(value any, target int) bool {
switch number := value.(type) {
case float64:
+3 -7
View File
@@ -618,15 +618,11 @@ func completeMinutesUpload(rt *shortcut.RuntimeContext, sessionID string, timeou
}
message := strings.ToLower(err.Error())
retryable := strings.Contains(message, "still uploading") || strings.Contains(message, "wait and retry") || strings.Contains(message, "正在上传")
if !retryable || time.Now().Add(interval).After(deadline) {
if !retryable || minutesPollDeadlineReached(deadline, interval) {
return nil, attempts, err
}
timer := time.NewTimer(interval)
select {
case <-rt.Command().Context().Done():
timer.Stop()
return nil, attempts, rt.Command().Context().Err()
case <-timer.C:
if err := waitMinutesInterval(rt, interval); err != nil {
return nil, attempts, err
}
}
}
+30 -6
View File
@@ -4,6 +4,7 @@
package minutes
import (
"context"
"encoding/json"
"fmt"
"os"
@@ -372,7 +373,7 @@ func runMinutesMindmap(rt *shortcut.RuntimeContext, id string, timeout, interval
payload := map[string]any{"operation": "minutes.mindmap", "complete": false, "taskUuid": id, "taskStatus": status, "attempts": attempts, "result": result, "recovery": map[string]any{"taskUuid": id, "nextAction": "inspect source transcript; do not assume an empty mind map"}}
return payload, minutesCompositeError("minutes_mindmap_failed", "poll", payload)
}
if time.Now().Add(interval).After(deadline) {
if minutesPollDeadlineReached(deadline, interval) {
payload := map[string]any{"operation": "minutes.mindmap", "complete": false, "taskUuid": id, "taskStatus": status, "attempts": attempts, "recovery": map[string]any{"taskUuid": id, "nextAction": "dws minutes mind-graph status --id <taskUuid>"}}
return payload, minutesCompositeError("minutes_mindmap_timeout", "poll", payload)
}
@@ -422,7 +423,7 @@ func runMinutesSpeakerInsights(rt *shortcut.RuntimeContext, id string, timeout,
payload := map[string]any{"operation": "minutes.speaker_insights", "complete": false, "taskUuid": id, "taskId": taskID, "attempts": attempts, "stage": "poll", "recovery": map[string]any{"taskUuid": id, "taskId": taskID, "nextAction": "dws minutes speaker summary get --ids <taskUuid>"}}
return payload, callErr
}
if time.Now().Add(interval).After(deadline) {
if minutesPollDeadlineReached(deadline, interval) {
payload := map[string]any{"operation": "minutes.speaker_insights", "complete": false, "taskUuid": id, "taskId": taskID, "attempts": attempts, "stage": "poll", "recovery": map[string]any{"taskUuid": id, "taskId": taskID, "nextAction": "dws minutes speaker summary get --ids <taskUuid>"}}
return payload, minutesCompositeError("minutes_speaker_insights_timeout", "poll", payload)
}
@@ -614,6 +615,17 @@ func executeMinutesShare(rt *shortcut.RuntimeContext) error {
}
func executeMinutesUnshare(rt *shortcut.RuntimeContext) error {
if !rt.DryRun() {
for _, id := range minutesIDs(rt) {
data, err := rt.CallMCPData("minutes", "get_minutes_basic_info", map[string]any{"taskUuid": id})
if err != nil {
return fmt.Errorf("minutes unshare preflight for %s: %w", id, err)
}
if _, err := minutesdata.Basic(id, data); err != nil {
return fmt.Errorf("minutes unshare preflight for %s: %w", id, err)
}
}
}
return executeMinutesPermissionLedger(rt, "unshare", "remove_member_permission", func(member string) map[string]any {
return map[string]any{"uuids": minutesIDs(rt), "memberUids": []string{member}}
})
@@ -631,7 +643,11 @@ func executeMinutesPermissionLedger(rt *shortcut.RuntimeContext, operation, tool
for index, member := range members {
data, err := rt.CallMCPWriteDataStrict("minutes", tool, params(member))
if err == nil {
err = minutesdata.RequireWriteAcknowledgement(operation, data)
if operation == "unshare" {
err = minutesdata.RequirePermissionMutationAcknowledgement(operation, minutesIDs(rt), []string{member}, data)
} else {
err = minutesdata.RequireWriteAcknowledgement(operation, data)
}
}
if err != nil {
failures = append(failures, map[string]any{"memberUid": member, "error": err.Error()})
@@ -708,7 +724,7 @@ func waitMinutesArtifacts(rt *shortcut.RuntimeContext, id string, artifacts []st
for {
attempts++
bundle, failures := collectMinutesArtifactsOnce(rt, id, artifacts, pageLimit)
if len(failures) == 0 || time.Now().Add(interval).After(deadline) {
if len(failures) == 0 || minutesPollDeadlineReached(deadline, interval) {
return bundle, failures, attempts
}
if err := waitMinutesInterval(rt, interval); err != nil {
@@ -720,14 +736,22 @@ func waitMinutesArtifacts(rt *shortcut.RuntimeContext, id string, artifacts []st
func waitMinutesInterval(rt *shortcut.RuntimeContext, interval time.Duration) error {
timer := time.NewTimer(interval)
defer timer.Stop()
commandContext := rt.Command().Context()
if commandContext == nil {
commandContext = context.Background()
}
select {
case <-rt.Command().Context().Done():
return rt.Command().Context().Err()
case <-commandContext.Done():
return commandContext.Err()
case <-timer.C:
return nil
}
}
func minutesPollDeadlineReached(deadline time.Time, interval time.Duration) bool {
return !time.Now().Add(interval).Before(deadline)
}
func outputWorkflowResult(rt *shortcut.RuntimeContext, payload map[string]any, failed bool, reason, stage string) error {
if err := rt.Output(payload); err != nil {
return err
@@ -323,7 +323,10 @@ func TestCrossPlatformCoverageMinutesPermissionLedgerBranchesE2E(t *testing.T) {
if err := runMinutesAlignmentCLIWithWriter(t, &minutesE2ECaller{}, minutesFailWriter{}, "minutes", "+unshare", "--id", "u1", "--member-uids", "m1", "--dry-run"); err == nil {
t.Fatal("unshare dry-run output failure accepted")
}
unshare := &minutesE2ECaller{responses: map[string][]string{"minutes/remove_member_permission": {`{"success":true,"result":{}}`}}}
unshare := &minutesE2ECaller{responses: map[string][]string{
"minutes/get_minutes_basic_info": {`{"success":true,"result":{"taskUuid":"u1"}}`},
"minutes/remove_member_permission": {`{"success":true,"result":{"resultMap":{"u1":["m1"]}}}`},
}}
payload, _, err := runMinutesAlignmentCLI(t, unshare, "minutes", "+unshare", "--id", "u1", "--member-uids", "m1", "--yes")
if err != nil || payload["complete"] != true || payload["succeeded"] != float64(1) {
t.Fatalf("unshare payload=%#v err=%v", payload, err)
@@ -337,11 +340,31 @@ func TestCrossPlatformCoverageMinutesPermissionLedgerBranchesE2E(t *testing.T) {
if args["coverPermission"] != "true" || len(args["roleSubResourceIds"].([]string)) != 2 {
t.Fatalf("share args=%#v", args)
}
continueFailure := &minutesE2ECaller{responses: map[string][]string{"minutes/remove_member_permission": {`{"result":{}}`, `{"success":true,"result":{}}`}}}
continueFailure := &minutesE2ECaller{responses: map[string][]string{
"minutes/get_minutes_basic_info": {`{"success":true,"result":{"taskUuid":"u1"}}`},
"minutes/remove_member_permission": {
`{"result":{}}`,
`{"success":true,"result":{"resultMap":{"u1":["m2"]}}}`,
},
}}
payload, output, err := runMinutesAlignmentCLI(t, continueFailure, "minutes", "+unshare", "--id", "u1", "--member-uids", "m1,m2", "--failure-policy", "continue", "--yes")
if err == nil || output == "" || payload["failed"] != float64(1) || payload["succeeded"] != float64(1) {
t.Fatalf("continue ledger payload=%#v err=%v", payload, err)
}
missing := &minutesE2ECaller{responses: map[string][]string{
"minutes/get_minutes_basic_info": {`{"success":true,"result":{}}`},
}}
payload, output, err = runMinutesAlignmentCLI(t, missing, "minutes", "+unshare", "--id", "missing", "--member-uids", "m1", "--yes")
if err == nil || payload != nil || output != "" || missing.counts["minutes/remove_member_permission"] != 0 {
t.Fatalf("missing minutes reached unshare: payload=%#v output=%q err=%v calls=%#v", payload, output, err, missing.counts)
}
preflightFailure := &minutesE2ECaller{failAt: map[string]int{"minutes/get_minutes_basic_info": 1}}
payload, output, err = runMinutesAlignmentCLI(t, preflightFailure, "minutes", "+unshare", "--id", "unavailable", "--member-uids", "m1", "--yes")
if err == nil || payload != nil || output != "" || preflightFailure.counts["minutes/remove_member_permission"] != 0 {
t.Fatalf("failed preflight reached unshare: payload=%#v output=%q err=%v calls=%#v", payload, output, err, preflightFailure.counts)
}
}
func TestCrossPlatformCoverageMinutesArtifactCollectorBranches(t *testing.T) {
@@ -419,10 +442,9 @@ func TestCrossPlatformCoverageMinutesArtifactWaitAndOutput(t *testing.T) {
t.Fatalf("cancel failures=%#v", failures)
}
cmd = &cobra.Command{Use: "wait"}
cmd.SetContext(context.Background())
rt = shortcut.RuntimeContextForTest(cmd, ExportPack)
if err := waitMinutesInterval(rt, 0); err != nil {
t.Fatalf("timer wait: %v", err)
t.Fatalf("timer wait with unset command context: %v", err)
}
cmd.SetOut(minutesFailWriter{})
@@ -564,7 +586,7 @@ func TestCrossPlatformCoverageMinutesExportPackBranchesE2E(t *testing.T) {
}
func TestCrossPlatformCoverageMinutesExportPathAndJSONFaults(t *testing.T) {
if _, _, _, err := prepareExportTarget("../escape"); err == nil {
if _, _, _, err := prepareExportTarget(filepath.Join("..", "escape")); err == nil {
t.Fatal("unsafe target accepted")
}
for _, test := range []struct {
@@ -590,7 +612,7 @@ func TestCrossPlatformCoverageMinutesExportPathAndJSONFaults(t *testing.T) {
}
}},
{name: "escape", run: func(t *testing.T) {
testseam.Swap(t, &minutesRel, func(string, string) (string, error) { return "../escape", nil })
testseam.Swap(t, &minutesRel, func(string, string) (string, error) { return filepath.Join("..", "escape"), nil })
if _, _, _, err := prepareExportTarget("pack"); err == nil {
t.Fatal("rel escape accepted")
}
@@ -602,7 +624,7 @@ func TestCrossPlatformCoverageMinutesExportPathAndJSONFaults(t *testing.T) {
}
}},
{name: "parent", run: func(t *testing.T) {
testseam.Swap(t, &minutesRel, func(string, string) (string, error) { return "parent/pack", nil })
testseam.Swap(t, &minutesRel, func(string, string) (string, error) { return filepath.Join("parent", "pack"), nil })
testseam.Swap(t, &minutesLstat, func(string) (os.FileInfo, error) { return nil, errors.New("parent") })
if _, _, _, err := prepareExportTarget("pack"); err == nil {
t.Fatal("unsafe parent accepted")
@@ -224,6 +224,33 @@ func TestCrossPlatformCoverageMinutesWorkflowCompletion(t *testing.T) {
if err := RequireWriteAcknowledgement("write", map[string]any{"result": map[string]any{"updated": true}}); err != nil {
t.Fatal(err)
}
if err := RequirePermissionMutationAcknowledgement("unshare", []string{"u1"}, []string{"202397"}, map[string]any{
"success": true,
"result": map[string]any{"resultMap": map[string]any{"u1": []any{float64(202397)}}},
}); err != nil {
t.Fatalf("valid permission acknowledgement rejected: %v", err)
}
for _, data := range []map[string]any{
{"success": false, "errorMsg": "denied"},
{"success": "yes"},
{"success": true, "result": "bad"},
{"success": true, "result": map[string]any{}},
{"success": true, "result": map[string]any{"resultMap": map[string]any{}}},
{"success": true, "result": map[string]any{"resultMap": map[string]any{"other": []any{"202397"}}}},
{"success": true, "result": map[string]any{"resultMap": map[string]any{"u1": "bad"}}},
{"success": true, "result": map[string]any{"resultMap": map[string]any{"u1": []any{nil}}}},
{"success": true, "result": map[string]any{"resultMap": map[string]any{"u1": []any{"other"}}}},
} {
if err := RequirePermissionMutationAcknowledgement("unshare", []string{"u1"}, []string{"202397"}, data); err == nil {
t.Fatalf("invalid permission acknowledgement accepted: %#v", data)
}
}
if err := RequirePermissionMutationAcknowledgement("unshare", []string{"u1"}, []string{"202397", "other"}, map[string]any{
"success": true,
"result": map[string]any{"resultMap": map[string]any{"u1": []any{"202397"}}},
}); err == nil {
t.Fatal("incomplete member coverage accepted")
}
for _, tc := range []struct {
cmd, id string
data map[string]any
+55
View File
@@ -30,6 +30,61 @@ func RequireWriteAcknowledgement(operation string, data map[string]any) error {
return fmt.Errorf("minutes %s response has no explicit successful acknowledgement", operation)
}
// RequirePermissionMutationAcknowledgement validates the target-level payload
// returned by add/remove_member_permission. The backend's top-level
// success=true is not sufficient: historically it also accompanied a
// resultMap that merely echoed nonexistent task UUIDs.
func RequirePermissionMutationAcknowledgement(operation string, taskUUIDs, memberUIDs []string, data map[string]any) error {
if err := validateEnvelope(data); err != nil {
return err
}
if success, ok := data["success"].(bool); !ok || !success {
return fmt.Errorf("minutes %s response has no explicit success=true", operation)
}
result, ok := data["result"].(map[string]any)
if !ok {
return fmt.Errorf("minutes %s response has no result object", operation)
}
resultMap, ok := result["resultMap"].(map[string]any)
if !ok {
return fmt.Errorf("minutes %s response has no result.resultMap object", operation)
}
if len(resultMap) != len(taskUUIDs) {
return fmt.Errorf("minutes %s resultMap covers %d task UUIDs, want %d", operation, len(resultMap), len(taskUUIDs))
}
wantMembers := make(map[string]bool, len(memberUIDs))
for _, member := range memberUIDs {
wantMembers[strings.TrimSpace(member)] = true
}
for _, taskUUID := range taskUUIDs {
rawMembers, exists := resultMap[taskUUID]
if !exists {
return fmt.Errorf("minutes %s resultMap is missing task UUID %s", operation, taskUUID)
}
members, ok := rawMembers.([]any)
if !ok {
return fmt.Errorf("minutes %s resultMap[%s] has type %T, want array", operation, taskUUID, rawMembers)
}
gotMembers := make(map[string]bool, len(members))
for index, member := range members {
value := stringField(map[string]any{"member": member}, "member")
if value == "" {
return fmt.Errorf("minutes %s resultMap[%s][%d] is not a member UID", operation, taskUUID, index)
}
gotMembers[value] = true
}
if len(gotMembers) != len(wantMembers) {
return fmt.Errorf("minutes %s resultMap[%s] member coverage mismatch", operation, taskUUID)
}
for member := range wantMembers {
if !gotMembers[member] {
return fmt.Errorf("minutes %s resultMap[%s] is missing member UID %s", operation, taskUUID, member)
}
}
}
return nil
}
// RecordResult validates the gateway's observed listening-note command result.
func RecordResult(expectedCmd, taskUUID string, data map[string]any) (map[string]any, error) {
if err := RequireWriteAcknowledgement("record "+expectedCmd, data); err != nil {
+115 -21
View File
@@ -6,16 +6,35 @@ ROOT="$(CDPATH= cd -- "$(dirname -- "$0")/../.." && pwd)"
usage() {
printf '%s\n' \
"usage: $0 verify <app-package>" \
" $0 run <app-package>" >&2
" $0 run <app-package> [partition]" \
" $0 list-partitions" >&2
exit 2
}
[ "$#" -eq 2 ] || usage
mode="$1"
app_package="$2"
# Single source of truth for the partition set. CI runs one job per partition and
# pins its shard names to this list, so a name that appears here without a
# dispatch entry below fails closed rather than silently skipping tests.
APP_PARTITIONS='schema a-b c d-r s-z-example-fuzz'
mode="${1:-}"
partition=""
case "$mode" in
verify|run) ;;
list-partitions)
[ "$#" -eq 1 ] || usage
for name in $APP_PARTITIONS; do
printf '%s\n' "$name"
done
exit 0
;;
verify)
[ "$#" -eq 2 ] || usage
app_package="$2"
;;
run)
[ "$#" -eq 2 ] || [ "$#" -eq 3 ] || usage
app_package="$2"
partition="${3:-}"
;;
*) usage ;;
esac
@@ -70,21 +89,49 @@ if [ "$unmatched_count" -ne 0 ]; then
exit 1
fi
for partition in \
# The loop variable is deliberately not named "partition": that name holds the
# partition requested on the command line, and run mode still executes this
# discovery pass before dispatching.
classified=''
for spec in \
"schema:$schema_count" \
"a-b:$ab_count" \
"c:$c_count" \
"d-r:$dr_count" \
"s-z-example-fuzz:$sz_count"
do
name="${partition%%:*}"
count="${partition#*:}"
name="${spec%%:*}"
count="${spec#*:}"
classified="$classified $name"
if [ "$count" -eq 0 ]; then
printf 'app race partition %s is empty\n' "$name" >&2
exit 1
fi
done
# The counters above and APP_PARTITIONS must describe the same set in both
# directions. A counted partition missing from APP_PARTITIONS would never be
# dispatched by any CI job, and a dispatchable partition with no counter would
# escape the exact-coverage check above.
for name in $APP_PARTITIONS; do
case " $classified " in
*" $name "*) ;;
*)
printf 'app partition %s has no coverage counter\n' "$name" >&2
exit 1
;;
esac
done
for name in $classified; do
case " $APP_PARTITIONS " in
*" $name "*) ;;
*)
printf 'counted partition %s is not a dispatchable app partition\n' "$name" >&2
exit 1
;;
esac
done
total_count="$(wc -l < "$tests" | tr -d ' ')"
assigned_count=$((schema_count + ab_count + c_count + dr_count + sz_count))
if [ "$assigned_count" -ne "$total_count" ]; then
@@ -101,26 +148,73 @@ fi
run_partition() {
name="$1"
run_pattern="$2"
skip_pattern="${3:-}"
instrumentation="$2"
run_pattern="$3"
skip_pattern="${4:-}"
printf 'running internal/app race partition %s\n' "$name"
# Fail closed on an unrecognized mode: a typo must not silently drop race
# instrumentation from a partition that is supposed to carry it.
case "$instrumentation" in
race|no-race) ;;
*)
printf 'unknown instrumentation %s for app partition %s\n' \
"$instrumentation" "$name" >&2
exit 1
;;
esac
printf 'running internal/app %s partition %s\n' "$instrumentation" "$name"
set -- -v -count=1 -timeout=15m -run "$run_pattern"
if [ -n "$skip_pattern" ]; then
go test -v -race -count=1 -timeout=15m \
-run "$run_pattern" -skip "$skip_pattern" "$app_package"
else
go test -v -race -count=1 -timeout=15m \
-run "$run_pattern" "$app_package"
set -- "$@" -skip "$skip_pattern"
fi
if [ "$instrumentation" = race ]; then
set -- -race "$@"
fi
go test "$@" "$app_package"
}
# Schema assembly has the largest transient memory footprint. Run those tests
# in a fresh process, then keep each remaining name range in its own process so
# command trees retained by process-global registries are released between
# partitions. The complementary run/skip patterns preserve the full test set.
#
# The schema partition runs uninstrumented. Its tests assert structural
# Schema-to-Cobra contracts over a single goroutine: none of them call
# t.Parallel or start a goroutine, so the race detector has no concurrent access
# to observe here. The process-global lazy metadata that does need race coverage
# (schema_source_root's atomic.Value, the parameter-binding lazy loaders) is
# exercised by internal/cli's concurrent tests, which stay instrumented. The
# instrumentation is not free on this partition: its shared sync.Once Catalog
# build is allocation-heavy, and -race made the partition roughly 11x slower
# (26s -> 291s locally, 357s in CI) without being able to report anything.
schema_pattern='^Test.*Schema'
run_partition schema "$schema_pattern"
run_partition a-b '^Test[A-B]' "$schema_pattern"
run_partition c '^TestC' "$schema_pattern"
run_partition d-r '^Test[D-R]' "$schema_pattern"
run_partition s-z-example-fuzz '^(Test[S-Z]|Example|Fuzz)' "$schema_pattern"
# Dispatch table for the partition set declared in APP_PARTITIONS. CI passes one
# partition per job so they run concurrently; running without a partition keeps
# the original end-to-end behaviour for local use and for any caller that wants
# the whole package in one invocation.
run_named_partition() {
case "$1" in
schema) run_partition schema no-race "$schema_pattern" ;;
a-b) run_partition a-b race '^Test[A-B]' "$schema_pattern" ;;
c) run_partition c race '^TestC' "$schema_pattern" ;;
d-r) run_partition d-r race '^Test[D-R]' "$schema_pattern" ;;
s-z-example-fuzz)
run_partition s-z-example-fuzz race '^(Test[S-Z]|Example|Fuzz)' "$schema_pattern"
;;
*)
printf 'unknown app partition: %s\n' "$1" >&2
exit 1
;;
esac
}
if [ -n "$partition" ]; then
run_named_partition "$partition"
exit 0
fi
for name in $APP_PARTITIONS; do
run_named_partition "$name"
done
@@ -14,7 +14,7 @@
| delete-old-doc | 1. `doc search --query "<关键词>"` 或 `doc list --folder <FOLDER>` 取 `nodeId`<br>2. **必须先向用户展示要删除的文档标题/路径**,等用户确认<br>3. `doc delete --node <nodeId> --yes`(注意:是 `doc delete`,不是 `doc block delete`;前者删整篇,后者删块)<br>4. 验证:`doc info --node <nodeId>` 应返回 not-found |
| export-doc-as-docx | 1. `doc search --query "<关键词>"` → 取 `nodeId`<br>2. `doc export --node <nodeId> --output /tmp/<name>.docx --timeout-sec 600`(CLI 内置渐进式退避轮询,**不要自己拼 GET downloadUrl**)<br>3. 落盘失败兜底:`doc export get --job-id <ID> --output /tmp/<name>.docx` 续等<br>4. (可选)`doc-to-message` 把 docx 路径作为附件发给用户 |
| grant-doc-access | 1. `doc search --query "<关键词>"` → 取 `nodeId`<br>2. `contact user search --query "<姓名>"` → 取 `userId`(**注意**: doc permission 用 `userId`,不是 `openDingTalkId`)<br>3. **节点级**授权:`doc permission add --node <nodeId> --user <UID1,UID2> --role EDITOR`<br>  role 取值: MANAGER / EDITOR / DOWNLOADER / READER(**不要传 OWNER**)<br>  单次最多 30 个 userId<br>4. 查权限:`doc permission list --node <nodeId>` 确认<br>5. **不要混淆**: 给"知识库"加成员用 `wiki member add --workspace`(容器级);给"单个文档"加权限用 `doc permission add --node`(节点级)|
| insert-image-to-doc | 1. `doc search --query "<关键词>"` → 取 `nodeId`<br>2. `doc media insert --node <nodeId> --file ./image.png`(3 步一体化:获取凭证 → PUT OSS → 插入块;图片 ≤20MB 自动作为内联图片,其他作为附件块)<br>3. (可选)`doc read --node <nodeId>` 回读确认<br>**禁止**: 自己拿 uploadUrl 写 PUT —— 90%+ 会因 Content-Type 未清空触发 SignatureDoesNotMatch |
| insert-image-to-doc | 1. `doc search --query "<关键词>"` → 取 `nodeId`<br>2. `doc media insert --node <nodeId> --file ./image.png`(4 步一体化:获取凭证 → PUT OSS → 插入块 → 自动回读验证;图片 ≤20MB 自动作为内联图片,其他作为附件块)<br>3. 检查结果中的 `verified=true`;若服务端无法提供明确回读证据,命令会返回 partial success<br>**禁止**: 自己拿 uploadUrl 写 PUT —— 90%+ 会因 Content-Type 未清空触发 SignatureDoesNotMatch |
| download-doc-attachment | 1. `doc block list --node <nodeId>` → 找到 `blockType=attachment` 的块取 `resourceId`<br>2. `doc media download --node <nodeId> --resource-id <RID>` → 返回 OSS 临时 URL(**注意 expirationSeconds**)<br>3. 调用方自行 GET 该 URL 落盘 |
---
@@ -83,7 +83,8 @@ Example:
{"toolName":"range update","input":{"sheet-id":"Sheet1","range":"A1","values":[[{"type":"text","text":"hello"}]]}},
{"toolName":"merge-cells","input":{"sheet-id":"Sheet1","range":"A1:B1","merge-type":"mergeAll"}},
{"toolName":"update-dimension","input":{"sheet-id":"Sheet1","dimension":"ROWS","start-index":"1","length":1,"pixel-size":40}},
{"toolName":"group-dimension","input":{"sheet-id":"Sheet1","range":"3:7","group-state":"expand"}}
{"toolName":"group-dimension","input":{"sheet-id":"Sheet1","range":"3:7","group-state":"expand"}},
{"toolName":"set-dropdown","input":{"sheet-id":"Sheet1","range":"C2:C100","source-sheet-id":"SourceSheet","source-range":"T1:T3"}}
]'
dws sheet batch-update --node <NODE_ID> --continue-on-error --operations '[...]'
Flags:
@@ -101,6 +102,9 @@ Notes:
- operations 最多 20 条
- 当需要对多个区域执行相同清除时,优先使用 `range batch-clear`(更简洁)
- `csv-put` 子操作与独立命令语义一致:CSV 字段值以 `=` 开头时按公式解析;前加单引号时写入以 `=` 开头的字面文本
- `set-dropdown` 的 `input` 中,Inline 使用 `options`;SourceRange 使用 `source-sheet-id` + `source-range`,两种模式必须且只能选一个。顶层 `colors` / `source-colors` 会被拒绝;Inline 颜色写在 `options[].color`,SourceRange 颜色写入暂不支持
- `set-dropdown` SourceRange 在已验证的重命名、引用前插入行/列、删除引用前行的场景会自动调整;已验证的 `move-dimension` 会使其变为 `invalid`。其他未覆盖删除/移动场景后先回读,仅 `invalid` 时重新选源写入
- `source-range` 按 `toolName` 解释:`set-dropdown` 中是下拉候选项来源;`range fill` / `range copy-to` / `range move-to` 中是待填充、复制或移动的数据源区域
- 典型场景:先插入行列再写入数据、先清除再写入、批量合并+调整行高列宽
- `group-dimension` 在 batch 中只适合默认展开分组;需要 `--group-state fold` 时请使用独立 `dws sheet group-dimension`
- `table-put` 不支持放进 batch-update;结构化 table 请用独立 `dws sheet table-put`
@@ -7,6 +7,7 @@
用户说"设置下拉列表/下拉选项/下拉菜单/添加下拉/配置下拉":
- 设置下拉列表 → `set-dropdown`
- 设置多选下拉 → `set-dropdown --multi-select`
- 引用单元格区域作为候选项 → `set-dropdown --source-sheet-id ... --source-range ...`
用户说"查看下拉列表/获取下拉配置/下拉列表有哪些选项":
- 获取下拉列表配置 → `get-dropdown`
@@ -29,18 +30,25 @@ Example:
dws sheet set-dropdown --node <NODE_ID> --sheet-id <SHEET_ID> --range "B2:B50" \
--options '[{"value":"高","color":"#ff0000"},{"value":"中","color":"#ffaa00"},{"value":"低","color":"#00ff00"}]' \
--multi-select
# 引用同一工作簿内另一工作表的区域作为候选项来源
dws sheet set-dropdown --node <NODE_ID> --sheet-id <TARGET_SHEET_ID> --range "C2:C100" \
--source-sheet-id <SOURCE_SHEET_ID> --source-range "T1:T3"
Flags:
--node string 表格文档 ID 或 URL (必填)
--sheet-id string 工作表 ID 或名称 (必填)
--range string 目标单元格范围,A1 表示法,如 A2:A100 (必填)
--options string 下拉选项 JSON 数组 (必填),如 '[{"value":"选项1","color":"#ff0000"}]'
--multi-select 是否允许多选(默认单选)
--node string 表格文档 ID 或 URL (必填)
--sheet-id string 工作表 ID 或名称 (必填)
--range string 目标单元格范围,A1 表示法,如 A2:A100 (必填)
--options string Inline 下拉选项 JSON 数组,与 --source-range 二选一
--source-sheet-id string SourceRange 来源工作表 ID,与 --source-range 同时指定
--source-range string SourceRange 来源区域,与 --options 二选一;不带工作表前缀
--multi-select 是否允许多选(默认单选)
```
在指定单元格范围内设置下拉列表。设置后用户可从预定义选项中选择值。
- **用途**:为单元格配置下拉列表,支持自定义选项颜色和多选。
在指定单元格范围内设置下拉列表。Inline 模式直接存储选项;SourceRange 模式引用同一工作簿内的来源区域,可跨工作表,并支持普通区域、整行和整列。
- **用途**:为单元格配置静态选项或区域来源下拉,两种模式都支持多选;颜色仅 Inline 支持。
- **场景**:规范数据输入,如状态选择(完成/进行中/待处理)、优先级(高/中/低)等。
- **注意**:选项值不能包含英文逗号;如果目标范围已存在下拉列表,会被新配置覆盖。
- **注意**:`--options` 与 `--source-range` 必须且只能指定一个。`--source-range` 只写 `T1:T3`、`T:T`、`1:3` 这类 A1 区域,来源工作表通过 `--source-sheet-id` 单独指定;不接受工作表前缀、公式或多区域。SourceRange 颜色写入暂不支持。
- **结构操作行为**:已验证的工作表重命名、在引用前插入行/列、删除引用前的行会自动调整引用并保持 `valid`;已验证的 `move-dimension` 场景会使其变为 `invalid`。列删除、删除整个来源区域或来源工作表等场景未覆盖,不能预设结果;结构操作后先回读 `sourceRangeStatus`,仅在 `invalid` 时重新选择来源并写入。
### 获取下拉列表配置
```
@@ -55,10 +63,10 @@ Flags:
--range string 查询范围,A1 表示法,如 A1:A100 (必填)
```
查询指定范围内的下拉列表配置信息,包括选项值、颜色和是否多选。
查询指定范围内的下拉列表配置信息。
- **用途**:查看单元格已设置的下拉列表选项和配置。
- **场景**:在修改下拉列表前先查询现有配置;确认下拉列表是否设置成功。
- **返回**:`dataValidations` 数组,相同选项的单元格聚合为一组,每组包含 `conditionValues`(选项值)、`ranges`(覆盖范围)、`options`(含 `enableMultiSelect` 和 `colorValueMap`)。范围内无下拉列表时 `hasDropdown` 为 false。
- **返回**:`dataValidations` 数组由底层按配置聚合。Inline 组返回 `sourceType:"inline"`、`conditionValues`、`ranges` 和 `options`;SourceRange 组始终返回 `sourceType:"sourceRange"`、`sourceRangeStatus:"valid"/"invalid"`、`enableMultiSelect` 和 `ranges`,仅在 `sourceRangeStatus:"valid"` 时返回 `sourceRange:{sheetId,a1Notation}`。`invalid` 时仍保留配置组,但省略 `sourceRange`,不得依赖旧坐标修复。SourceRange 不会展开候选值,因此不返回 `conditionValues`、`options` 或颜色。范围内无下拉列表时 `hasDropdown` 为 false。
### 删除下拉列表
```
@@ -81,14 +89,15 @@ Flags:
| 操作 | 从返回中提取 | 用于 |
|------|-------------|------|
| `set-dropdown` | `range` 实际设置范围、`optionCount` 选项数量、`enableMultiSelect` 是否多选 | 确认下拉列表设置成功 |
| `get-dropdown` | `hasDropdown` 是否存在下拉、`dataValidations` 下拉配置列表(含 `conditionValues`、`ranges`、`options`) | 查看已有下拉配置 |
| `set-dropdown` | `range` 实际设置范围、`enableMultiSelect` 是否多选;仅 Inline 模式返回 `optionCount` | 确认下拉列表设置成功 |
| `get-dropdown` | `hasDropdown`、`dataValidations`;按 `sourceType` 区分 Inline 与 SourceRange | 查看已有下拉配置 |
| `delete-dropdown` | `range` 实际删除范围 | 确认下拉列表删除完成 |
| `list` | 工作表的 `sheetId` | info / range read / range update / find 的 --sheet-id |
## 注意事项
- ★ **`--sheet-id` 获取规范(强制)**:`sheetId` 未知时必须先通过 `dws sheet list --node <NODE_ID> --format json` 查询,禁止凭空编造(如臆测为 `Sheet1`、`sheet1`、`0`、`default` 等)
- `set-dropdown` 在指定范围内设置下拉列表,`--options` 为 JSON 数组,每个元素包含 `value`(必填)和 `color`(可选,`#RRGGBB` 格式)。选项值不能包含英文逗号。`--multi-select` 启用多选模式。如果目标范围已存在下拉列表,会被新配置覆盖
- `get-dropdown` 查询指定范围内的下拉列表配置,返回 `dataValidations` 数组,相同选项的单元格聚合为一组。无下拉列表时 `hasDropdown` 为 false
- `set-dropdown` 的 Inline 模式使用 `--options`,每个元素包含 `value`(必填)和 `color`(可选,`#RRGGBB`);SourceRange 模式使用 `--source-sheet-id` + `--source-range`。两种模式均可用 `--multi-select`,并会覆盖目标范围已有下拉
- SourceRange 在已验证的重命名、引用前插入行/列、删除引用前行的场景会自动调整;已验证的 `move-dimension` 会使其变为 `invalid`。其他未覆盖删除/移动场景后先回读,仅 `invalid` 时重新选源写入;颜色写入暂不支持
- `get-dropdown` 查询指定范围内的下拉配置,聚合发生在底层服务。SourceRange 即使无效也保留一组并以 `sourceRangeStatus:"invalid"` 表示,但省略 `sourceRange`,不回退展开候选值
- `delete-dropdown` 删除指定范围内的下拉列表配置,单元格恢复为普通文本格式。已填写的值不会被清除。目标范围不存在下拉列表时操作仍返回成功
@@ -13,7 +13,7 @@
| 读取目的 | 推荐命令 | 说明 |
|---------|---------|------|
| 快速查看纯值、数据分析、大表分批读取 | `csv-get` | CSV 格式,token 消耗约为 JSON 的 1/3,内置 maxChars 防爆 |
| 快速查看纯值、数据分析、大表分批读取 | `csv-get` | CSV 格式,token 消耗约为 JSON 的 1/3,内置 30,000 单元格和 maxChars 防爆 |
| 按 table/dataframe 协议读取 | `table-get` | 返回 `columns` / `data` / `dtypes` / `formats`;默认首行为表头 |
| 查看数据验证配置(下拉/复选框) | `range read` | 返回 per-cell 结构,含 dataValidation |
| 查看单元格样式(背景色/字体/对齐等) | `range read` | 返回 per-cell 结构,含 cellStyles(仅显式设置的样式) |
@@ -52,7 +52,10 @@ Flags:
- `csv` — CSV 文本,每逻辑行前加 `[row=N]` 前缀标注真实表格行号。行号一律从此前缀读取,禁止手算
- `colIndices` — 列字母映射数组(如 `["A","B","C"]`)。定位列字母用 `colIndices[j]`,禁止手数逗号
- `rowIndices` — 行号映射数组(如 `[1,2,3]`)
- `hasMore` — 是否因 maxChars 截断。为 true 时需要调整 `--range` 继续分页读取
- `hasMore` — 完整目标范围是否还有未返回数据。`true` 是部分成功,不能宣称已读完
- `truncationReasons` — 部分返回原因。`max_cells` 表示命中单次 30,000 单元格上限;`max_chars` 表示命中 CSV 字符上限,两者可同时出现
- `resolvedRange` — 仅在未传 `--range` 时返回,表示底层解析出的完整目标范围
- `returnedRange` — 本次实际完整返回的精确范围;它不是服务端续读游标
`csv-get` 不返回合并单元格结构。若 CSV 中出现合并区域的非左上角单元格为空,不能据此判断该区域"无内容";需要先用 `dws sheet info --node <NODE_ID> --sheet-id <SHEET_ID> --format json` 读取 `mergedRanges`,再结合左上角单元格理解合并区域语义。
@@ -63,10 +66,15 @@ Flags:
| `raw_value` | 原始值(如 1000、45808) | 数据处理、计算 |
| `formula` | 公式文本(如 =SUM(A1:A10)),无公式时回退原始值 | 查看/复制公式 |
**大表分批读取**:当 `hasMore=true` 或数据量很大时,按行窗口分批:
- 先通过 `info` 获取 `nonEmptyRange.range`,或用 `nonEmptyRange.lastRow` / `nonEmptyRange.lastColumn` 确定 A1 边界
- 分批读取:`--range "A1:J500"`、`--range "A501:J1000"` ……
- 单次建议 ≤5000 单元格
**大表分批读取**:`csv-get` 不会自动发起后续请求。当 `hasMore=true` 时:
- 显式传了 `--range` 时,该范围是完整目标;未传时,以 `resolvedRange` 为完整目标
- 根据 `returnedRange` 的结束行,从下一行开始显式构造新的 `--range`;例如目标是 `A1:A30001`、本次返回 `A1:A30000`,则继续读 `A30001:A30001`
- `max_cells` 不能通过增大 `--max-chars` 解决;`max_chars` 且没有 `returnedRange` 时,增大 `--max-chars` 或缩小列范围
- 上述完成度协议覆盖未传 `--range` 和显式有限矩形范围;`A:A`、`1:1` 等非有限范围可能不返回完成度字段,建议先用 `sheet info` 的 `nonEmptyRange` 换算成有限范围再读
- 即使未命中硬上限,单次仍建议控制在 5000 单元格以内,减少超时和 token 消耗
**工作簿过大失败与部分成功不同**:收到顶层
`errorCode=forbidden.document.sizeOverLimit` 和对应 `errorMessage` 时,表示工作簿整体无法装载。应停止读取并提示创建更小副本或拆分工作簿;缩小 `--range` 不能解决这类错误。不要把它当成 `hasMore=true` 的可续读结果。
### 以 table/dataframe 协议读取结构化数据
```
@@ -130,6 +138,8 @@ Flags:
```
**返回字段说明**:
- `rowIndices` / `colIndices` — 本次返回的真实行号与列字母映射
- `hasMore` / `truncationReasons` / `resolvedRange` / `returnedRange` — 语义与 `csv-get` 相同;`range read` 的 `truncationReasons` 主要为 `max_cells`
- `cells` — 二维数组,第一维为行,第二维为列。每个元素为 per-cell 对象,字段如下:
| 字段 | 类型 | 是否必有 | 说明 |
@@ -146,10 +156,16 @@ Flags:
| type | 字段 | 说明 |
|------|------|------|
| `dropdown` | `options: [{value: string, color?: string}]` | 下拉选项列表 |
| `dropdown` | `enableMultiSelect: boolean` | 是否允许多选 |
| `dropdown` Inline | `sourceType: "inline"` | 静态选项模式 |
| `dropdown` Inline | `options: [{value: string, color?: string}]` | 下拉选项列表 |
| `dropdown` Inline | `enableMultiSelect: boolean` | 是否允许多选 |
| `dropdown` SourceRange | `sourceType: "sourceRange"` | 区域来源模式 |
| `dropdown` SourceRange | `sourceRange: {sheetId, a1Notation}` | 仅 `sourceRangeStatus:"valid"` 时返回的来源工作表和规范化 A1 区域;`invalid` 时省略 |
| `dropdown` SourceRange | `sourceRangeStatus: "valid" / "invalid"` | 底层引用是否可解析;无效引用仍返回配置,但不返回 `sourceRange` |
| `checkbox` | `checked: boolean` | 当前勾选状态 |
SourceRange 模式不动态拉取或展开来源区域的候选值,因此不返回 `options` 或颜色。`invalid` 时只能依赖 `sourceRangeStatus` 判定需要重新选源,不得假定仍能从 `sourceRange` 取回旧坐标。
**hyperlink 结构**:
| type | 字段 | 说明 |
@@ -216,7 +232,7 @@ Flags:
**公式校验建议**:写公式后不要只看写入返回结果。先用 `formula` 模式确认公式文本已落表,再运行 `formula-verify` 聚合扫描错误;关键业务数值继续用 `raw_value` 抽样对账。详见 [sheet-formula](./sheet-formula.md)。
**超时处理建议**:读取大范围数据时若出现超时或响应过慢,请主动缩小 `--range` 查询范围,**建议单次读取的单元格数量控制在 5000 个以内**(例如 50 行 × 100 列、100 行 × 50 列)。对于大表可采用分页读取策略:
**超时处理建议**:读取大范围数据时若出现超时或响应过慢,请主动缩小 `--range` 查询范围,**建议单次读取的单元格数量控制在 5000 个以内**(例如 50 行 × 100 列、100 行 × 50 列)。对于大表可采用分批读取策略,每次都必须检查 `hasMore`:
- 先通过 `info` 获取 `nonEmptyRange.range`,或用 `nonEmptyRange.lastRow` / `nonEmptyRange.lastColumn` 确定 A1 边界
- 按行分批读取,如 `A1:J500`、`A501:J1000`、`A1001:J1500` ……
- 避免不传 `--range` 直接读取整个大工作表
@@ -232,7 +248,7 @@ dws sheet list --node <NODE_ID> --format json
# 2. 查看工作表详情(行列数、最后非空位置、mergedRanges 等)
dws sheet info --node <NODE_ID> --sheet-id <SHEET_ID> --format json
# 3. 读取全部数据
# 3. 读取数据(只有 hasMore=false 才代表目标范围已完整返回)
dws sheet range read --node <NODE_ID> --sheet-id <SHEET_ID> --format json
# 4. 读取指定区域
@@ -251,6 +267,8 @@ dws sheet range read --node <NODE_ID> --sheet-id <SHEET_ID> --range "A1:D10" --f
- ★ **`--sheet-id` 获取规范(强制)**:`sheetId` 未知时必须先通过 `dws sheet list --node <NODE_ID> --format json` 查询真实的 `sheetId` / 工作表名称后再调用,禁止凭空编造(如臆测为 `Sheet1`、`sheet1`、`0`、`default` 等);用户仅给出工作表名称时,也应通过 `list` 校验该名称是否存在,避免名称大小写或拼写不一致导致失败
- `range read` 不传 `--range` 时默认读取整个工作表的全部非空数据
- `range read` 的 `--range` 支持 `Sheet1!A1:D10` 格式直接指定工作表(此时忽略 `--sheet-id`)
- ★ `csv-get` 和 `range read` 都必须检查 `hasMore`;`true` 时结合目标范围与 `returnedRange` 显式分批,CLI 不会自动续读
- ★ `forbidden.document.sizeOverLimit` 是工作簿装载失败,不是 30,000 单元格部分成功;缩小 `--range` 不能修复
- ★ `csv-get` / `range read` / `range get` 不返回合并单元格结构;查看合并范围必须用 `sheet info` 的 `mergedRanges`
- `table-get` 返回 table/dataframe 结构,适合需要 `columns` / `data` / `dtypes` / `formats` 的场景;它不返回 per-cell 元数据,也不替代 `range read`
- ★ 大整数和长数字标识符回读校验:精确 ID 应是字符串 / `object`;若看到 `float64`,说明它已经按数值路径处理,超过 `9007199254740991` 时可能已丢精度
@@ -386,9 +386,14 @@ Flags:
**dropdown(下拉列表)**:
```json
{ "type": "text", "text": "High", "dataValidation": { "type": "dropdown", "options": [{"value":"High","color":"#00ff00"},{"value":"Low","color":"#ff0000"}], "enableMultiSelect": false } }
{ "dataValidation": { "type": "dropdown", "sourceRange": {"sheetId":"SOURCE_SHEET_ID","a1Notation":"T1:T3"}, "enableMultiSelect": false } }
```
- `options`:必填,`[{value, color?}]` 数组
- `options` 与 `sourceRange` 必须且只能提供一个
- Inline:`options` 为非空 `[{value, color?}]` 数组
- SourceRange:`sourceRange` 为 `{sheetId,a1Notation}`;可引用同一工作簿内其他工作表,支持普通区域、整行和整列。`a1Notation` 不带工作表前缀,不接受公式或多区域;颜色写入暂不支持
- `enableMultiSelect`:可选,是否多选,默认 false
- SourceRange 在已验证的工作表重命名、引用前插入行/列、删除引用前行的场景会自动调整;已验证的 `move-dimension` 会使其变为 `invalid`。其他未覆盖删除/移动场景后先回读 `sourceRangeStatus`,仅 `invalid` 时重新选源写入
- 复合 cell 写入中,如果值或样式已成功落盘但 SourceRange 校验失败,服务端可返回 `success:true` 并通过 `message` 说明下拉未创建;必须检查 `message` 并按需回读
**checkbox(复选框)**:
```json
@@ -83,7 +83,8 @@ Example:
{"toolName":"range update","input":{"sheet-id":"Sheet1","range":"A1","values":[[{"type":"text","text":"hello"}]]}},
{"toolName":"merge-cells","input":{"sheet-id":"Sheet1","range":"A1:B1","merge-type":"mergeAll"}},
{"toolName":"update-dimension","input":{"sheet-id":"Sheet1","dimension":"ROWS","start-index":"1","length":1,"pixel-size":40}},
{"toolName":"group-dimension","input":{"sheet-id":"Sheet1","range":"3:7","group-state":"expand"}}
{"toolName":"group-dimension","input":{"sheet-id":"Sheet1","range":"3:7","group-state":"expand"}},
{"toolName":"set-dropdown","input":{"sheet-id":"Sheet1","range":"C2:C100","source-sheet-id":"SourceSheet","source-range":"T1:T3"}}
]'
dws sheet batch-update --node <NODE_ID> --continue-on-error --operations '[...]'
Flags:
@@ -101,6 +102,9 @@ Notes:
- operations 最多 20 条
- 当需要对多个区域执行相同清除时,优先使用 `range batch-clear`(更简洁)
- `csv-put` 子操作与独立命令语义一致:CSV 字段值以 `=` 开头时按公式解析;前加单引号时写入以 `=` 开头的字面文本
- `set-dropdown` 的 `input` 中,Inline 使用 `options`;SourceRange 使用 `source-sheet-id` + `source-range`,两种模式必须且只能选一个。顶层 `colors` / `source-colors` 会被拒绝;Inline 颜色写在 `options[].color`,SourceRange 颜色写入暂不支持
- `set-dropdown` SourceRange 在已验证的重命名、引用前插入行/列、删除引用前行的场景会自动调整;已验证的 `move-dimension` 会使其变为 `invalid`。其他未覆盖删除/移动场景后先回读,仅 `invalid` 时重新选源写入
- `source-range` 按 `toolName` 解释:`set-dropdown` 中是下拉候选项来源;`range fill` / `range copy-to` / `range move-to` 中是待填充、复制或移动的数据源区域
- 典型场景:先插入行列再写入数据、先清除再写入、批量合并+调整行高列宽
- `group-dimension` 在 batch 中只适合默认展开分组;需要 `--group-state fold` 时请使用独立 `dws sheet group-dimension`
- `table-put` 不支持放进 batch-update;结构化 table 请用独立 `dws sheet table-put`
@@ -7,6 +7,7 @@
用户说"设置下拉列表/下拉选项/下拉菜单/添加下拉/配置下拉":
- 设置下拉列表 → `set-dropdown`
- 设置多选下拉 → `set-dropdown --multi-select`
- 引用单元格区域作为候选项 → `set-dropdown --source-sheet-id ... --source-range ...`
用户说"查看下拉列表/获取下拉配置/下拉列表有哪些选项":
- 获取下拉列表配置 → `get-dropdown`
@@ -29,18 +30,25 @@ Example:
dws sheet set-dropdown --node <NODE_ID> --sheet-id <SHEET_ID> --range "B2:B50" \
--options '[{"value":"高","color":"#ff0000"},{"value":"中","color":"#ffaa00"},{"value":"低","color":"#00ff00"}]' \
--multi-select
# 引用同一工作簿内另一工作表的区域作为候选项来源
dws sheet set-dropdown --node <NODE_ID> --sheet-id <TARGET_SHEET_ID> --range "C2:C100" \
--source-sheet-id <SOURCE_SHEET_ID> --source-range "T1:T3"
Flags:
--node string 表格文档 ID 或 URL (必填)
--sheet-id string 工作表 ID 或名称 (必填)
--range string 目标单元格范围,A1 表示法,如 A2:A100 (必填)
--options string 下拉选项 JSON 数组 (必填),如 '[{"value":"选项1","color":"#ff0000"}]'
--multi-select 是否允许多选(默认单选)
--node string 表格文档 ID 或 URL (必填)
--sheet-id string 工作表 ID 或名称 (必填)
--range string 目标单元格范围,A1 表示法,如 A2:A100 (必填)
--options string Inline 下拉选项 JSON 数组,与 --source-range 二选一
--source-sheet-id string SourceRange 来源工作表 ID,与 --source-range 同时指定
--source-range string SourceRange 来源区域,与 --options 二选一;不带工作表前缀
--multi-select 是否允许多选(默认单选)
```
在指定单元格范围内设置下拉列表。设置后用户可从预定义选项中选择值。
- **用途**:为单元格配置下拉列表,支持自定义选项颜色和多选。
在指定单元格范围内设置下拉列表。Inline 模式直接存储选项;SourceRange 模式引用同一工作簿内的来源区域,可跨工作表,并支持普通区域、整行和整列。
- **用途**:为单元格配置静态选项或区域来源下拉,两种模式都支持多选;颜色仅 Inline 支持。
- **场景**:规范数据输入,如状态选择(完成/进行中/待处理)、优先级(高/中/低)等。
- **注意**:选项值不能包含英文逗号;如果目标范围已存在下拉列表,会被新配置覆盖。
- **注意**:`--options` 与 `--source-range` 必须且只能指定一个。`--source-range` 只写 `T1:T3`、`T:T`、`1:3` 这类 A1 区域,来源工作表通过 `--source-sheet-id` 单独指定;不接受工作表前缀、公式或多区域。SourceRange 颜色写入暂不支持。
- **结构操作行为**:已验证的工作表重命名、在引用前插入行/列、删除引用前的行会自动调整引用并保持 `valid`;已验证的 `move-dimension` 场景会使其变为 `invalid`。列删除、删除整个来源区域或来源工作表等场景未覆盖,不能预设结果;结构操作后先回读 `sourceRangeStatus`,仅在 `invalid` 时重新选择来源并写入。
### 获取下拉列表配置
```
@@ -55,10 +63,10 @@ Flags:
--range string 查询范围,A1 表示法,如 A1:A100 (必填)
```
查询指定范围内的下拉列表配置信息,包括选项值、颜色和是否多选。
查询指定范围内的下拉列表配置信息。
- **用途**:查看单元格已设置的下拉列表选项和配置。
- **场景**:在修改下拉列表前先查询现有配置;确认下拉列表是否设置成功。
- **返回**:`dataValidations` 数组,相同选项的单元格聚合为一组,每组包含 `conditionValues`(选项值)、`ranges`(覆盖范围)、`options`(含 `enableMultiSelect` 和 `colorValueMap`)。范围内无下拉列表时 `hasDropdown` 为 false。
- **返回**:`dataValidations` 数组由底层按配置聚合。Inline 组返回 `sourceType:"inline"`、`conditionValues`、`ranges` 和 `options`;SourceRange 组始终返回 `sourceType:"sourceRange"`、`sourceRangeStatus:"valid"/"invalid"`、`enableMultiSelect` 和 `ranges`,仅在 `sourceRangeStatus:"valid"` 时返回 `sourceRange:{sheetId,a1Notation}`。`invalid` 时仍保留配置组,但省略 `sourceRange`,不得依赖旧坐标修复。SourceRange 不会展开候选值,因此不返回 `conditionValues`、`options` 或颜色。范围内无下拉列表时 `hasDropdown` 为 false。
### 删除下拉列表
```
@@ -81,14 +89,15 @@ Flags:
| 操作 | 从返回中提取 | 用于 |
|------|-------------|------|
| `set-dropdown` | `range` 实际设置范围、`optionCount` 选项数量、`enableMultiSelect` 是否多选 | 确认下拉列表设置成功 |
| `get-dropdown` | `hasDropdown` 是否存在下拉、`dataValidations` 下拉配置列表(含 `conditionValues`、`ranges`、`options`) | 查看已有下拉配置 |
| `set-dropdown` | `range` 实际设置范围、`enableMultiSelect` 是否多选;仅 Inline 模式返回 `optionCount` | 确认下拉列表设置成功 |
| `get-dropdown` | `hasDropdown`、`dataValidations`;按 `sourceType` 区分 Inline 与 SourceRange | 查看已有下拉配置 |
| `delete-dropdown` | `range` 实际删除范围 | 确认下拉列表删除完成 |
| `list` | 工作表的 `sheetId` | info / range read / range update / find 的 --sheet-id |
## 注意事项
- ★ **`--sheet-id` 获取规范(强制)**:`sheetId` 未知时必须先通过 `dws sheet list --node <NODE_ID> --format json` 查询,禁止凭空编造(如臆测为 `Sheet1`、`sheet1`、`0`、`default` 等)
- `set-dropdown` 在指定范围内设置下拉列表,`--options` 为 JSON 数组,每个元素包含 `value`(必填)和 `color`(可选,`#RRGGBB` 格式)。选项值不能包含英文逗号。`--multi-select` 启用多选模式。如果目标范围已存在下拉列表,会被新配置覆盖
- `get-dropdown` 查询指定范围内的下拉列表配置,返回 `dataValidations` 数组,相同选项的单元格聚合为一组。无下拉列表时 `hasDropdown` 为 false
- `set-dropdown` 的 Inline 模式使用 `--options`,每个元素包含 `value`(必填)和 `color`(可选,`#RRGGBB`);SourceRange 模式使用 `--source-sheet-id` + `--source-range`。两种模式均可用 `--multi-select`,并会覆盖目标范围已有下拉
- SourceRange 在已验证的重命名、引用前插入行/列、删除引用前行的场景会自动调整;已验证的 `move-dimension` 会使其变为 `invalid`。其他未覆盖删除/移动场景后先回读,仅 `invalid` 时重新选源写入;颜色写入暂不支持
- `get-dropdown` 查询指定范围内的下拉配置,聚合发生在底层服务。SourceRange 即使无效也保留一组并以 `sourceRangeStatus:"invalid"` 表示,但省略 `sourceRange`,不回退展开候选值
- `delete-dropdown` 删除指定范围内的下拉列表配置,单元格恢复为普通文本格式。已填写的值不会被清除。目标范围不存在下拉列表时操作仍返回成功
@@ -13,7 +13,7 @@
| 读取目的 | 推荐命令 | 说明 |
|---------|---------|------|
| 快速查看纯值、数据分析、大表分批读取 | `csv-get` | CSV 格式,token 消耗约为 JSON 的 1/3,内置 maxChars 防爆 |
| 快速查看纯值、数据分析、大表分批读取 | `csv-get` | CSV 格式,token 消耗约为 JSON 的 1/3,内置 30,000 单元格和 maxChars 防爆 |
| 按 table/dataframe 协议读取 | `table-get` | 返回 `columns` / `data` / `dtypes` / `formats`;默认首行为表头 |
| 查看数据验证配置(下拉/复选框) | `range read` | 返回 per-cell 结构,含 dataValidation |
| 查看单元格样式(背景色/字体/对齐等) | `range read` | 返回 per-cell 结构,含 cellStyles(仅显式设置的样式) |
@@ -52,7 +52,10 @@ Flags:
- `csv` — CSV 文本,每逻辑行前加 `[row=N]` 前缀标注真实表格行号。行号一律从此前缀读取,禁止手算
- `colIndices` — 列字母映射数组(如 `["A","B","C"]`)。定位列字母用 `colIndices[j]`,禁止手数逗号
- `rowIndices` — 行号映射数组(如 `[1,2,3]`)
- `hasMore` — 是否因 maxChars 截断。为 true 时需要调整 `--range` 继续分页读取
- `hasMore` — 完整目标范围是否还有未返回数据。`true` 是部分成功,不能宣称已读完
- `truncationReasons` — 部分返回原因。`max_cells` 表示命中单次 30,000 单元格上限;`max_chars` 表示命中 CSV 字符上限,两者可同时出现
- `resolvedRange` — 仅在未传 `--range` 时返回,表示底层解析出的完整目标范围
- `returnedRange` — 本次实际完整返回的精确范围;它不是服务端续读游标
`csv-get` 不返回合并单元格结构。若 CSV 中出现合并区域的非左上角单元格为空,不能据此判断该区域"无内容";需要先用 `dws sheet info --node <NODE_ID> --sheet-id <SHEET_ID> --format json` 读取 `mergedRanges`,再结合左上角单元格理解合并区域语义。
@@ -63,10 +66,15 @@ Flags:
| `raw_value` | 原始值(如 1000、45808) | 数据处理、计算 |
| `formula` | 公式文本(如 =SUM(A1:A10)),无公式时回退原始值 | 查看/复制公式 |
**大表分批读取**:当 `hasMore=true` 或数据量很大时,按行窗口分批:
- 先通过 `info` 获取 `nonEmptyRange.range`,或用 `nonEmptyRange.lastRow` / `nonEmptyRange.lastColumn` 确定 A1 边界
- 分批读取:`--range "A1:J500"`、`--range "A501:J1000"` ……
- 单次建议 ≤5000 单元格
**大表分批读取**:`csv-get` 不会自动发起后续请求。当 `hasMore=true` 时:
- 显式传了 `--range` 时,该范围是完整目标;未传时,以 `resolvedRange` 为完整目标
- 根据 `returnedRange` 的结束行,从下一行开始显式构造新的 `--range`;例如目标是 `A1:A30001`、本次返回 `A1:A30000`,则继续读 `A30001:A30001`
- `max_cells` 不能通过增大 `--max-chars` 解决;`max_chars` 且没有 `returnedRange` 时,增大 `--max-chars` 或缩小列范围
- 上述完成度协议覆盖未传 `--range` 和显式有限矩形范围;`A:A`、`1:1` 等非有限范围可能不返回完成度字段,建议先用 `sheet info` 的 `nonEmptyRange` 换算成有限范围再读
- 即使未命中硬上限,单次仍建议控制在 5000 单元格以内,减少超时和 token 消耗
**工作簿过大失败与部分成功不同**:收到顶层
`errorCode=forbidden.document.sizeOverLimit` 和对应 `errorMessage` 时,表示工作簿整体无法装载。应停止读取并提示创建更小副本或拆分工作簿;缩小 `--range` 不能解决这类错误。不要把它当成 `hasMore=true` 的可续读结果。
### 以 table/dataframe 协议读取结构化数据
```
@@ -130,6 +138,8 @@ Flags:
```
**返回字段说明**:
- `rowIndices` / `colIndices` — 本次返回的真实行号与列字母映射
- `hasMore` / `truncationReasons` / `resolvedRange` / `returnedRange` — 语义与 `csv-get` 相同;`range read` 的 `truncationReasons` 主要为 `max_cells`
- `cells` — 二维数组,第一维为行,第二维为列。每个元素为 per-cell 对象,字段如下:
| 字段 | 类型 | 是否必有 | 说明 |
@@ -146,10 +156,16 @@ Flags:
| type | 字段 | 说明 |
|------|------|------|
| `dropdown` | `options: [{value: string, color?: string}]` | 下拉选项列表 |
| `dropdown` | `enableMultiSelect: boolean` | 是否允许多选 |
| `dropdown` Inline | `sourceType: "inline"` | 静态选项模式 |
| `dropdown` Inline | `options: [{value: string, color?: string}]` | 下拉选项列表 |
| `dropdown` Inline | `enableMultiSelect: boolean` | 是否允许多选 |
| `dropdown` SourceRange | `sourceType: "sourceRange"` | 区域来源模式 |
| `dropdown` SourceRange | `sourceRange: {sheetId, a1Notation}` | 仅 `sourceRangeStatus:"valid"` 时返回的来源工作表和规范化 A1 区域;`invalid` 时省略 |
| `dropdown` SourceRange | `sourceRangeStatus: "valid" / "invalid"` | 底层引用是否可解析;无效引用仍返回配置,但不返回 `sourceRange` |
| `checkbox` | `checked: boolean` | 当前勾选状态 |
SourceRange 模式不动态拉取或展开来源区域的候选值,因此不返回 `options` 或颜色。`invalid` 时只能依赖 `sourceRangeStatus` 判定需要重新选源,不得假定仍能从 `sourceRange` 取回旧坐标。
**hyperlink 结构**:
| type | 字段 | 说明 |
@@ -216,7 +232,7 @@ Flags:
**公式校验建议**:写公式后不要只看写入返回结果。先用 `formula` 模式确认公式文本已落表,再运行 `formula-verify` 聚合扫描错误;关键业务数值继续用 `raw_value` 抽样对账。详见 [sheet-formula](./sheet-formula.md)。
**超时处理建议**:读取大范围数据时若出现超时或响应过慢,请主动缩小 `--range` 查询范围,**建议单次读取的单元格数量控制在 5000 个以内**(例如 50 行 × 100 列、100 行 × 50 列)。对于大表可采用分页读取策略:
**超时处理建议**:读取大范围数据时若出现超时或响应过慢,请主动缩小 `--range` 查询范围,**建议单次读取的单元格数量控制在 5000 个以内**(例如 50 行 × 100 列、100 行 × 50 列)。对于大表可采用分批读取策略,每次都必须检查 `hasMore`:
- 先通过 `info` 获取 `nonEmptyRange.range`,或用 `nonEmptyRange.lastRow` / `nonEmptyRange.lastColumn` 确定 A1 边界
- 按行分批读取,如 `A1:J500`、`A501:J1000`、`A1001:J1500` ……
- 避免不传 `--range` 直接读取整个大工作表
@@ -232,7 +248,7 @@ dws sheet list --node <NODE_ID> --format json
# 2. 查看工作表详情(行列数、最后非空位置、mergedRanges 等)
dws sheet info --node <NODE_ID> --sheet-id <SHEET_ID> --format json
# 3. 读取全部数据
# 3. 读取数据(只有 hasMore=false 才代表目标范围已完整返回)
dws sheet range read --node <NODE_ID> --sheet-id <SHEET_ID> --format json
# 4. 读取指定区域
@@ -251,6 +267,8 @@ dws sheet range read --node <NODE_ID> --sheet-id <SHEET_ID> --range "A1:D10" --f
- ★ **`--sheet-id` 获取规范(强制)**:`sheetId` 未知时必须先通过 `dws sheet list --node <NODE_ID> --format json` 查询真实的 `sheetId` / 工作表名称后再调用,禁止凭空编造(如臆测为 `Sheet1`、`sheet1`、`0`、`default` 等);用户仅给出工作表名称时,也应通过 `list` 校验该名称是否存在,避免名称大小写或拼写不一致导致失败
- `range read` 不传 `--range` 时默认读取整个工作表的全部非空数据
- `range read` 的 `--range` 支持 `Sheet1!A1:D10` 格式直接指定工作表(此时忽略 `--sheet-id`)
- ★ `csv-get` 和 `range read` 都必须检查 `hasMore`;`true` 时结合目标范围与 `returnedRange` 显式分批,CLI 不会自动续读
- ★ `forbidden.document.sizeOverLimit` 是工作簿装载失败,不是 30,000 单元格部分成功;缩小 `--range` 不能修复
- ★ `csv-get` / `range read` / `range get` 不返回合并单元格结构;查看合并范围必须用 `sheet info` 的 `mergedRanges`
- `table-get` 返回 table/dataframe 结构,适合需要 `columns` / `data` / `dtypes` / `formats` 的场景;它不返回 per-cell 元数据,也不替代 `range read`
- ★ 大整数和长数字标识符回读校验:精确 ID 应是字符串 / `object`;若看到 `float64`,说明它已经按数值路径处理,超过 `9007199254740991` 时可能已丢精度
@@ -386,9 +386,14 @@ Flags:
**dropdown(下拉列表)**:
```json
{ "type": "text", "text": "High", "dataValidation": { "type": "dropdown", "options": [{"value":"High","color":"#00ff00"},{"value":"Low","color":"#ff0000"}], "enableMultiSelect": false } }
{ "dataValidation": { "type": "dropdown", "sourceRange": {"sheetId":"SOURCE_SHEET_ID","a1Notation":"T1:T3"}, "enableMultiSelect": false } }
```
- `options`:必填,`[{value, color?}]` 数组
- `options` 与 `sourceRange` 必须且只能提供一个
- Inline:`options` 为非空 `[{value, color?}]` 数组
- SourceRange:`sourceRange` 为 `{sheetId,a1Notation}`;可引用同一工作簿内其他工作表,支持普通区域、整行和整列。`a1Notation` 不带工作表前缀,不接受公式或多区域;颜色写入暂不支持
- `enableMultiSelect`:可选,是否多选,默认 false
- SourceRange 在已验证的工作表重命名、引用前插入行/列、删除引用前行的场景会自动调整;已验证的 `move-dimension` 会使其变为 `invalid`。其他未覆盖删除/移动场景后先回读 `sourceRangeStatus`,仅 `invalid` 时重新选源写入
- 复合 cell 写入中,如果值或样式已成功落盘但 SourceRange 校验失败,服务端可返回 `success:true` 并通过 `message` 说明下拉未创建;必须检查 `message` 并按需回读
**checkbox(复选框)**:
```json
+21 -15
View File
@@ -1157,7 +1157,7 @@ func TestChangelogPRFastPathWorkflowContract(t *testing.T) {
`test "$(git rev-parse HEAD^1)" = "$PR_BASE_SHA"`,
`test "$(git rev-parse HEAD^2)" = "$PR_HEAD_SHA"`,
`echo "TEST_HEAD_REF=$(git rev-parse HEAD)"`,
`list-shard "$TEST_SHARD" "$TEST_BASE_REF" "$TEST_HEAD_REF"`,
`list-shard "$package_shard" "$TEST_BASE_REF" "$TEST_HEAD_REF"`,
"needs.lint.outputs.full_suite != 'true'",
`name: "Test (race: ${{ matrix.shard }})"`,
"name: Test (workflow and release contracts)",
@@ -1219,16 +1219,20 @@ func TestChangelogPRFastPathWorkflowContract(t *testing.T) {
}
// The focused path fans the impacted set across the same shards as test-race
// and runs each shard the way test-race runs it, so no single job carries
// internal/app together with its reverse dependencies. internal/app keeps its
// package-level headroom through the process-isolating helper instead of one
// long -timeout, which is strictly stronger: every process releases the
// framework registries it populated. release-scripts is asserted because its
// dedicated job only runs at full-suite or release-sensitive scope, so losing
// it here would silently stop testing test/scripts changes.
// internal/app together with its reverse dependencies. internal/app is split
// further into one shard per bounded partition, which keeps its package-level
// headroom through the process-isolating helper instead of one long -timeout
// and is strictly stronger than a single app job: every partition process
// releases the framework registries it populated, and the partitions run
// concurrently rather than end to end. Each partition shard still selects the
// same single internal/app package, so the impacted-package query maps the
// shard name back to app. release-scripts is asserted because its dedicated
// job only runs at full-suite or release-sensitive scope, so losing it here
// would silently stop testing test/scripts changes.
for _, want := range []string{
`if [ "$TEST_SHARD" = "app" ]; then`,
`app-*) package_shard=app ;;`,
`test "${#packages[@]}" -eq 1`,
`./scripts/ci/run-app-race-tests.sh run "${packages[0]}"`,
`./scripts/ci/run-app-race-tests.sh run "${packages[0]}" "${TEST_SHARD#app-}"`,
`if [ "$TEST_SHARD" = "release-scripts" ]; then`,
`go test -v -count=1 -timeout=10m "${packages[@]}"`,
"timeout_budget=12m",
@@ -1249,14 +1253,16 @@ func TestChangelogPRFastPathWorkflowContract(t *testing.T) {
t.Fatal("Code Admission workflow missing race test job boundaries")
}
raceJob := admission[raceStart:raceEnd]
// The app shard uses independently bounded test processes so process-global
// command registries are released before the Schema assembly peak. Other full
// race shards retain the dynamic package timeout: default/floor 12m, with
// cli/smoke raised to 15m on slower hosted runners.
// internal/app is carried by one shard per bounded partition, so process-global
// command registries are released with each partition process and the Schema
// assembly peak no longer sits in front of the other partitions. Each
// partition shard resolves back to the same single internal/app package.
// Other full race shards retain the dynamic package timeout: default/floor
// 12m, with cli/smoke raised to 15m on slower hosted runners.
for _, want := range []string{
`if [ "$TEST_SHARD" = "app" ]; then`,
`app-*) package_shard=app ;;`,
`test "${#packages[@]}" -eq 1`,
`./scripts/ci/run-app-race-tests.sh run "${packages[0]}"`,
`./scripts/ci/run-app-race-tests.sh run "${packages[0]}" "${TEST_SHARD#app-}"`,
"timeout_budget=12m",
`if [ "$TEST_SHARD" = "cli" ] ||`,
`[ "$TEST_SHARD" = "smoke" ]; then`,
+68
View File
@@ -99,6 +99,74 @@ func TestCIAppRacePartitionsCoverTopLevelTestsExactlyOnce(t *testing.T) {
}
}
// TestCIAppRacePartitionMatrixMatchesHelper pins the workflow's app partition
// shards to the partition set the helper actually runs. The partitions are
// separate CI jobs now, so the helper's own "covered exactly once" check can no
// longer prove the whole package ran: a partition the helper knows about but no
// matrix shard dispatches would silently stop running while every job stays
// green. Both directions are asserted so a stale matrix shard fails too.
func TestCIAppRacePartitionMatrixMatchesHelper(t *testing.T) {
root := testPackagePlanRoot(t)
script := filepath.Join(root, "scripts", "ci", "run-app-race-tests.sh")
cmd := exec.Command("sh", script, "list-partitions")
cmd.Dir = root
output, err := cmd.CombinedOutput()
if err != nil {
t.Fatalf("%s list-partitions failed: %v\n%s", script, err, output)
}
partitions := strings.Fields(string(output))
if len(partitions) == 0 {
t.Fatalf("list-partitions returned no partitions: %q", output)
}
workflow, err := os.ReadFile(filepath.Join(root, ".github", "workflows", "ci.yml"))
if err != nil {
t.Fatalf("read ci.yml: %v", err)
}
admission := string(workflow)
for _, job := range []struct {
name string
startMark string
endMark string
}{
{"test-focused", "\n test-focused:\n", "\n test-race:\n"},
{"test-race", "\n test-race:\n", "\n test-release-scripts:\n"},
} {
start := strings.Index(admission, job.startMark)
end := strings.Index(admission, job.endMark)
if start < 0 || end <= start {
t.Fatalf("ci.yml is missing %s job boundaries", job.name)
}
body := admission[start:end]
for _, partition := range partitions {
want := "- app-" + partition
if !strings.Contains(body, want) {
t.Errorf("%s matrix is missing shard %q for a partition the helper runs", job.name, want)
}
}
for _, line := range strings.Split(body, "\n") {
shard := strings.TrimSpace(line)
if !strings.HasPrefix(shard, "- app-") {
continue
}
name := strings.TrimPrefix(shard, "- app-")
matched := false
for _, partition := range partitions {
if partition == name {
matched = true
break
}
}
if !matched {
t.Errorf("%s matrix shard %q has no matching helper partition", job.name, shard)
}
}
}
}
func TestCITestPackagePlanFailsClosedWhenGoListFails(t *testing.T) {
root := testPackagePlanRoot(t)
fakeBin := t.TempDir()