Compare commits
21
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e5c8ff9acd | ||
|
|
706535b41e | ||
|
|
9e88116a2d | ||
|
|
05a306148a | ||
|
|
0dcc796f4c | ||
|
|
3e792b1c86 | ||
|
|
b9c822d49d | ||
|
|
f9b9b83f48 | ||
|
|
aa9e67e7c8 | ||
|
|
16ff02903a | ||
|
|
65d3f2959c | ||
|
|
cb3087ba9b | ||
|
|
5068cfdab8 | ||
|
|
3c81e5d47d | ||
|
|
99e5a3cceb | ||
|
|
7e31043875 | ||
|
|
55d7fbf59a | ||
|
|
c0f4d21c05 | ||
|
|
660c908585 | ||
|
|
377ebc5e85 | ||
|
|
2eca203e74 |
@@ -438,6 +438,7 @@ jobs:
|
||||
release_version: ${{ steps.target.outputs.release_version }}
|
||||
release_commit: ${{ steps.target.outputs.release_commit }}
|
||||
release_tag_object: ${{ steps.target.outputs.release_tag_object }}
|
||||
oss_mirror: ${{ steps.target.outputs.oss_mirror }}
|
||||
is_recovery: ${{ steps.target.outputs.is_recovery }}
|
||||
failed_run_id: ${{ steps.target.outputs.failed_run_id }}
|
||||
channel: ${{ steps.contract.outputs.channel }}
|
||||
@@ -541,8 +542,14 @@ jobs:
|
||||
}
|
||||
tagFields.set(key, line.slice(separator + 2));
|
||||
}
|
||||
const cloudKeys = [
|
||||
"Channel",
|
||||
const ossMirror = tagFields.has("OSS-Mirror")
|
||||
? tagFields.get("OSS-Mirror")
|
||||
: "enabled";
|
||||
if (!["enabled", "deferred"].includes(ossMirror)) {
|
||||
core.setFailed(`${version} contains invalid OSS-Mirror metadata`);
|
||||
return;
|
||||
}
|
||||
const cloudOnlyKeys = [
|
||||
"Release-Run",
|
||||
"Release-Run-Attempt",
|
||||
"Requested-By",
|
||||
@@ -551,7 +558,8 @@ jobs:
|
||||
"Workflow-Commit",
|
||||
"Allocation-Fingerprint",
|
||||
];
|
||||
const hasAnyCloudMetadata = cloudKeys.some((key) => tagFields.has(key));
|
||||
const cloudKeys = ["Channel", ...cloudOnlyKeys];
|
||||
const hasAnyCloudMetadata = cloudOnlyKeys.some((key) => tagFields.has(key));
|
||||
const isCloudSeal = cloudKeys.every((key) => tagFields.has(key));
|
||||
if (hasAnyCloudMetadata && !isCloudSeal) {
|
||||
core.setFailed(`${version} contains incomplete cloud seal metadata`);
|
||||
@@ -655,6 +663,7 @@ jobs:
|
||||
core.setOutput("release_version", version);
|
||||
core.setOutput("release_commit", commit);
|
||||
core.setOutput("release_tag_object", tagObject);
|
||||
core.setOutput("oss_mirror", ossMirror);
|
||||
core.setOutput("is_recovery", isRecovery ? "true" : "false");
|
||||
core.setOutput("failed_run_id", failedRunId);
|
||||
|
||||
@@ -1694,6 +1703,7 @@ jobs:
|
||||
echo "npm $NPM_TAG already supersedes this historical rerun at v$delivered."
|
||||
|
||||
- name: Sync release artifacts to China OSS mirror
|
||||
if: ${{ needs.release-contract.outputs.oss_mirror == 'enabled' }}
|
||||
run: ./scripts/release/sync-to-oss.sh
|
||||
env:
|
||||
VERSION: ${{ needs.release-contract.outputs.release_version }}
|
||||
@@ -1703,7 +1713,7 @@ jobs:
|
||||
OSS_ENDPOINT: ${{ secrets.OSS_ENDPOINT }}
|
||||
OSS_BUCKET: ${{ secrets.OSS_BUCKET }}
|
||||
OSS_PREFIX: ${{ secrets.OSS_PREFIX }}
|
||||
DWS_REQUIRE_OSS: ${{ github.repository_owner == 'DingTalk-Real-AI' && '1' || '0' }}
|
||||
DWS_REQUIRE_OSS: "1"
|
||||
|
||||
mirror-gitee-release:
|
||||
name: Mirror immutable release to Gitee
|
||||
@@ -2218,6 +2228,24 @@ jobs:
|
||||
core.setFailed(`${version} does not resolve to one annotated commit tag`);
|
||||
return;
|
||||
}
|
||||
const tagFields = new Map();
|
||||
for (const line of tag.data.message.split(/\r?\n/)) {
|
||||
const separator = line.indexOf(": ");
|
||||
if (separator < 1) continue;
|
||||
const key = line.slice(0, separator);
|
||||
if (tagFields.has(key)) {
|
||||
core.setFailed(`${version} annotated tag contains duplicate ${key} metadata`);
|
||||
return;
|
||||
}
|
||||
tagFields.set(key, line.slice(separator + 2));
|
||||
}
|
||||
const ossMirror = tagFields.has("OSS-Mirror")
|
||||
? tagFields.get("OSS-Mirror")
|
||||
: "enabled";
|
||||
if (!["enabled", "deferred"].includes(ossMirror)) {
|
||||
core.setFailed(`${version} contains invalid OSS-Mirror metadata`);
|
||||
return;
|
||||
}
|
||||
const commitSha = tag.data.object.sha;
|
||||
const defaultBranch = context.payload.repository.default_branch;
|
||||
const branch = await github.rest.repos.getBranch({
|
||||
@@ -2273,6 +2301,17 @@ jobs:
|
||||
}
|
||||
core.setOutput("commit_sha", commitSha);
|
||||
core.setOutput("tag_object", ref.data.object.sha);
|
||||
core.setOutput("oss_mirror", ossMirror);
|
||||
|
||||
- name: Require sealed OSS policy for repair
|
||||
if: ${{ needs.dispatch-contract.outputs.mode == 'repair_oss' }}
|
||||
env:
|
||||
OSS_MIRROR: ${{ steps.authority.outputs.oss_mirror }}
|
||||
run: |
|
||||
test "$OSS_MIRROR" = enabled || {
|
||||
echo "OSS repair is unavailable because this immutable release deferred the OSS channel." >&2
|
||||
exit 1
|
||||
}
|
||||
|
||||
- name: Check out sealed release source
|
||||
uses: actions/checkout@v4
|
||||
@@ -2375,6 +2414,7 @@ jobs:
|
||||
from_beta: ${{ steps.allocate.outputs.from_beta }}
|
||||
base: ${{ steps.allocate.outputs.base }}
|
||||
refs_fingerprint: ${{ steps.allocate.outputs.refs_fingerprint }}
|
||||
oss_mirror: ${{ steps.allocate.outputs.oss_mirror }}
|
||||
steps:
|
||||
- name: Validate official default-branch release request
|
||||
env:
|
||||
@@ -2425,6 +2465,7 @@ jobs:
|
||||
env:
|
||||
REQUESTED_CHANNEL: ${{ inputs.release_channel }}
|
||||
REQUESTED_BUMP: ${{ inputs.release_bump }}
|
||||
OSS_MIRROR: ${{ vars.ENABLE_OSS_MIRROR == 'true' && 'enabled' || 'deferred' }}
|
||||
run: |
|
||||
set -eu
|
||||
allocation="$RUNNER_TEMP/next-release-version"
|
||||
@@ -2440,15 +2481,20 @@ jobs:
|
||||
value="$(sed -n "s/^$key=//p" "$allocation")"
|
||||
printf '%s=%s\n' "$key" "$value" >> "$GITHUB_OUTPUT"
|
||||
done
|
||||
refs_fingerprint="$(
|
||||
git for-each-ref --format='%(refname)=%(objectname)' \
|
||||
refs/tags/v refs/tags/withdrawn/v \
|
||||
| LC_ALL=C sort \
|
||||
| sha256sum \
|
||||
| awk '{print $1}'
|
||||
)"
|
||||
refs_unsorted="$RUNNER_TEMP/release-refs.unsorted"
|
||||
refs_manifest="$RUNNER_TEMP/release-refs"
|
||||
git for-each-ref --format='%(refname)=%(objectname)' \
|
||||
'refs/tags/v*' 'refs/tags/withdrawn/v*' > "$refs_unsorted"
|
||||
LC_ALL=C sort "$refs_unsorted" > "$refs_manifest"
|
||||
test -s "$refs_manifest" || {
|
||||
echo "release ref manifest is empty after fetching allocated tags" >&2
|
||||
exit 1
|
||||
}
|
||||
refs_fingerprint="$(sha256sum "$refs_manifest" | awk '{print $1}')"
|
||||
echo "release_commit=$GITHUB_SHA" >> "$GITHUB_OUTPUT"
|
||||
echo "refs_fingerprint=$refs_fingerprint" >> "$GITHUB_OUTPUT"
|
||||
case "$OSS_MIRROR" in enabled|deferred) ;; *) exit 1 ;; esac
|
||||
echo "oss_mirror=$OSS_MIRROR" >> "$GITHUB_OUTPUT"
|
||||
|
||||
release_version="$(sed -n 's/^release_version=//p' "$allocation")"
|
||||
from_beta="$(sed -n 's/^from_beta=//p' "$allocation")"
|
||||
@@ -2459,6 +2505,7 @@ jobs:
|
||||
echo "- Version: \`$release_version\`"
|
||||
echo "- Channel: \`$channel\`"
|
||||
echo "- Commit: \`$GITHUB_SHA\`"
|
||||
echo "- OSS mirror: \`$OSS_MIRROR\`"
|
||||
if test -n "$from_beta"; then
|
||||
echo "- Promotes beta: \`$from_beta\`"
|
||||
fi
|
||||
@@ -2474,6 +2521,7 @@ jobs:
|
||||
RELEASE_VERSION: ${{ steps.allocate.outputs.release_version }}
|
||||
RELEASE_CHANNEL: ${{ steps.allocate.outputs.channel }}
|
||||
FROM_BETA: ${{ steps.allocate.outputs.from_beta }}
|
||||
OSS_MIRROR: ${{ steps.allocate.outputs.oss_mirror }}
|
||||
run: |
|
||||
set -eu
|
||||
git config user.name "dws release cloud seal"
|
||||
@@ -2482,6 +2530,7 @@ jobs:
|
||||
printf 'Release %s\n\n' "$RELEASE_VERSION"
|
||||
printf 'Channel: %s\n' "$RELEASE_CHANNEL"
|
||||
test -z "$FROM_BETA" || printf 'From-Beta: %s\n' "$FROM_BETA"
|
||||
printf 'OSS-Mirror: %s\n' "$OSS_MIRROR"
|
||||
} > "$RUNNER_TEMP/candidate-tag-message"
|
||||
git tag -a "$RELEASE_VERSION" HEAD -F "$RUNNER_TEMP/candidate-tag-message"
|
||||
set -- \
|
||||
@@ -2554,6 +2603,7 @@ jobs:
|
||||
RELEASE_CHANNEL: ${{ needs.release-plan.outputs.channel }}
|
||||
FROM_BETA: ${{ needs.release-plan.outputs.from_beta }}
|
||||
REFS_FINGERPRINT: ${{ needs.release-plan.outputs.refs_fingerprint }}
|
||||
OSS_MIRROR: ${{ needs.release-plan.outputs.oss_mirror }}
|
||||
with:
|
||||
script: |
|
||||
const crypto = require("crypto");
|
||||
@@ -2565,6 +2615,7 @@ jobs:
|
||||
const channel = process.env.RELEASE_CHANNEL;
|
||||
const fromBeta = process.env.FROM_BETA;
|
||||
const expectedFingerprint = process.env.REFS_FINGERPRINT;
|
||||
const ossMirror = process.env.OSS_MIRROR;
|
||||
|
||||
if (
|
||||
context.payload.repository.full_name !== expectedRepository ||
|
||||
@@ -2577,7 +2628,8 @@ jobs:
|
||||
if (
|
||||
!/^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)(?:-beta\.[1-9]\d*)?$/.test(version) ||
|
||||
!/^[0-9a-f]{40}$/.test(commit) ||
|
||||
!/^[0-9a-f]{64}$/.test(expectedFingerprint)
|
||||
!/^[0-9a-f]{64}$/.test(expectedFingerprint) ||
|
||||
!["enabled", "deferred"].includes(ossMirror)
|
||||
) {
|
||||
core.setFailed("cloud release plan contains an invalid sealed identity");
|
||||
return;
|
||||
@@ -2646,6 +2698,7 @@ jobs:
|
||||
"",
|
||||
`Channel: ${channel}`,
|
||||
...(fromBeta ? [`From-Beta: ${fromBeta}`] : []),
|
||||
`OSS-Mirror: ${ossMirror}`,
|
||||
`Release-Run: ${context.runId}`,
|
||||
`Release-Run-Attempt: ${process.env.GITHUB_RUN_ATTEMPT}`,
|
||||
`Requested-By: ${context.actor}`,
|
||||
|
||||
@@ -142,7 +142,7 @@ jobs:
|
||||
echo "- Workflow result: ${RESULT}"
|
||||
echo "- Success means every configured channel was verified and the permanent withdrawn/${VERSION} tombstone remains as the version-reuse barrier."
|
||||
echo "- Failure may occur before or after the tombstone/channel mutations; inspect the failed step and rerun the exact same inputs after fixing the cause."
|
||||
echo "- The problem GitHub Release and original tag are removed after npm/OSS/Gitee rollback so GitHub installers stop resolving the bad version while the Homebrew rollback PR is reviewed."
|
||||
echo "- The problem GitHub Release and original tag are removed after npm and every tag-enabled/configured mirror are rolled back, so GitHub installers stop resolving the bad version while the Homebrew rollback PR is reviewed."
|
||||
echo "- npm is deprecated rather than unpublished; already-installed clients cannot be remotely downgraded."
|
||||
echo "- If a Homebrew rollback PR was opened, this run remains failed until that PR is independently reviewed, merged, and the workflow is rerun."
|
||||
} >> "$GITHUB_STEP_SUMMARY"
|
||||
|
||||
@@ -6,6 +6,45 @@ The format is inspired by [Keep a Changelog](https://keepachangelog.com/) and th
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
## [1.0.53] - 2026-07-21
|
||||
|
||||
This release promotes the sealed `v1.0.53-beta.6` contents to stable. It adds enterprise onboarding, declarative shortcuts, Sheet/Aitable writes, multi-account profiles, and broader personal IM events, while hardening authentication and the guarded release path.
|
||||
|
||||
### Added
|
||||
|
||||
- **Enterprise and office command coverage** — adds enterprise creation, employee invitation, and account provisioning commands; 366 declarative service shortcuts; Sheet import commands; and Aitable workflow create/update support with reviewed Schema contracts.
|
||||
- **Multiple accounts in one DingTalk organization** — profiles can distinguish accounts by organization and user, select them explicitly, and log out one account or an entire organization without overwriting another account's credentials.
|
||||
- **Expanded personal IM event subscriptions** (#651) — adds read-receipt, recall, and reaction events for one-to-one and group chats, plus specified-sender subscriptions by staff ID or OpenDingTalk ID.
|
||||
- **Official multi-platform Homebrew channel** — ships separate stable and keg-only beta Formulae for macOS and Linux across amd64 and arm64, with isolated update PRs.
|
||||
|
||||
### Changed
|
||||
|
||||
- **Personal event output contract** (#651) — `event consume` now emits event-specific top-level structured fields; scripts that consumed the former transport envelope must use the flat fields or select `-f raw`, while `--debug-raw-events` retains the diagnostic envelope.
|
||||
- **Guarded release lifecycle** — beta/stable publication now uses explicit promotion, immutable delivery proofs, protected recovery, and tag-bound optional OSS policy; an unprovisioned OSS mirror is sealed as `deferred` so GitHub, npm, and Homebrew are not blocked.
|
||||
|
||||
### Fixed
|
||||
|
||||
- **Authentication and credential reliability** — organization-policy denials stop before mutation or polling, long-running clients reload and refresh access tokens consistently, concurrent credential writes are atomic, and Windows portable-auth commands fail before reading or writing unsupported credential bundles.
|
||||
- **Command validation and compatibility** — invalid Sheet/task targets fail locally, IM shortcuts preserve AI-tag and alias compatibility, and Aitable import uploads require and forward a positive file size.
|
||||
- **Release publication reliability** — GitHub draft publication is bound to one verified release ID and exact assets, preflight uses isolated installer worktrees, guarded local tags remain compatible, and cloud planning fingerprints the actual allocated release refs.
|
||||
|
||||
## [1.0.53-beta.6] - 2026-07-21
|
||||
|
||||
This beta validates guarded local release compatibility and tag-bound OSS deferral so an unprovisioned mirror cannot block the primary release channels.
|
||||
|
||||
### Changed
|
||||
|
||||
- **Tag-bound optional OSS release mirror** — Official cloud Release runs no longer block GitHub, npm, and Homebrew delivery when an OSS bucket has not been provisioned. Cloud tags immutably record `OSS-Mirror: enabled|deferred`; publication, repair, and withdrawal consume that sealed policy instead of the current repository variable. Enabled releases remain fail-closed, while deferred releases skip the nonexistent channel and cannot be backfilled without a future audited repair proof.
|
||||
|
||||
### Fixed
|
||||
|
||||
- **Guarded local release compatibility** — The tag-push Release workflow now accepts the `Channel`-only annotated tags created by the guarded local release entry while continuing to reject any partial cloud-only seal metadata.
|
||||
- **Cloud release tag allocation fingerprint** — Release planning now fingerprints the actual `v*` and `withdrawn/v*` refs fetched from GitHub, matching the seal job's API view instead of hashing an empty non-wildcard ref prefix and rejecting every publish before tag creation.
|
||||
|
||||
## [1.0.53-beta.5] - 2026-07-21
|
||||
|
||||
This beta validates long-running access-token recovery and the faster, recoverable guarded release path introduced after v1.0.53-beta.4.
|
||||
|
||||
### Changed
|
||||
|
||||
- **Fast guarded beta and stable releases** — successful local release checks now leave a six-hour proof bound to the exact version, commit, repository identity, remote `main`, and stable baseline, so the subsequent guarded `--publish` invocation revalidates authority without repeating tests and packaging. A default-branch governance smoke uses the same dedicated immutable-release credential as the tag workflow before any tag is allocated.
|
||||
@@ -13,6 +52,7 @@ The format is inspired by [Keep a Changelog](https://keepachangelog.com/) and th
|
||||
|
||||
### Fixed
|
||||
|
||||
- **Long-running event authentication recovery** — personal and portal event streams resolve the current access token for every ticket request, refresh a server-rejected token with compare-and-refresh semantics, and reconnect with backoff when refresh is temporarily blocked by network failures, rate limits, or 5xx responses.
|
||||
- **Consistent access-token caching and errors** — runtime, recovery, Skill, PAT polling, and personal/portal event clients now resolve user access tokens through one expiry- and publication-aware manager, so long-running processes reload rotated credentials while keychain, refresh, parse, permission, and cancellation failures remain observable instead of being collapsed into “not authenticated.”
|
||||
- **Tag-push GitHub Release publication** — Draft publication now locks one GitHub Release database ID, verifies its exact tag, channel, notes, recovery marker, asset set, and uploaded bytes, then publishes and rechecks that same ID as immutable. Recovery runs use the trusted default-branch release helpers instead of the sealed tag's historical scripts, fixing the Draft-only `GET /releases/tags/{tag}` 404 without allowing the release identity to drift during recovery.
|
||||
- **Release preflight reliability** — source-mode installer tests now use isolated temporary checkouts and HOME directories instead of overwriting and deleting the real repository `dws` binary, release preflight explicitly rebuilds before policy checks, and the full-suite runner gives the growing script package a non-flaky five-minute per-suite budget.
|
||||
|
||||
+27
-15
@@ -108,21 +108,33 @@ commit, and failed tag-push run all match; it then reuses the normal release
|
||||
jobs. Do not put publication secrets in temporary branches or create ad-hoc
|
||||
recovery workflows.
|
||||
|
||||
If the immutable GitHub Release and npm package were delivered but a downstream
|
||||
China mirror failed, dispatch the normal `Release` workflow from the protected
|
||||
default branch with exactly one of `repair_gitee_version` or
|
||||
`repair_oss_version`. Channel repair accepts a failed exact-tag run only when
|
||||
its latest attempt completed the release contract, build, Apple signature,
|
||||
immutable GitHub publication, and npm delivery checks for the exact tagged
|
||||
commit. It then downloads and re-verifies the immutable assets before invoking
|
||||
only the selected mirror. An OSS repair requires the OSS step itself to be the
|
||||
recorded failure. A Gitee repair accepts either a failed Gitee job or a Gitee
|
||||
job that was skipped behind that OSS failure; the latter is an explicit Gitee
|
||||
backfill and does not claim that OSS has been repaired. Gitee repair requires
|
||||
`GITEE_TOKEN`, `GITEE_USER`, and `GITEE_REPO`; OSS repair requires
|
||||
`OSS_ACCESS_KEY_ID`, `OSS_ACCESS_KEY_SECRET`, `OSS_ENDPOINT`, and `OSS_BUCKET`
|
||||
(with optional `OSS_PREFIX`) as Actions secrets. Missing credentials fail the
|
||||
selected repair closed.
|
||||
Cloud-sealed releases mirror to OSS only when the repository variable
|
||||
`ENABLE_OSS_MIRROR` is exactly `true`. Leave the variable unset while no Bucket
|
||||
is provisioned; GitHub, npm, and Homebrew delivery can then complete without
|
||||
running the OSS step. Once enabled, missing credentials, an invalid Bucket, or
|
||||
an upload failure remains fail-closed. The cloud tag immutably records the
|
||||
decision as `OSS-Mirror: enabled|deferred`; publication and withdrawal consume
|
||||
that sealed value instead of the variable's later state. Deferred releases
|
||||
cannot use `repair_oss_version`; enabling OSS applies to later release tags
|
||||
until an audited immutable repair marker is implemented.
|
||||
|
||||
If an immutable GitHub Release and npm package were delivered but an enabled
|
||||
downstream China mirror failed, dispatch the normal `Release` workflow from the
|
||||
protected default branch with exactly one of `repair_gitee_version` or
|
||||
`repair_oss_version`. Channel repair accepts a fully successful exact release,
|
||||
or a failed exact-tag run only when its latest attempt completed the release
|
||||
contract, build, Apple signature, immutable GitHub publication, and npm
|
||||
delivery checks for the exact tagged commit. OSS repair additionally requires
|
||||
the tag's sealed policy to be `enabled`. It then downloads and re-verifies the
|
||||
immutable assets before invoking only the selected mirror. For a failed
|
||||
release, an OSS repair requires the OSS step itself to be the recorded failure.
|
||||
A Gitee repair accepts either a failed Gitee job or a Gitee job that was
|
||||
skipped behind that OSS failure; the latter is an explicit Gitee backfill and
|
||||
does not claim that OSS has been repaired. Gitee repair requires `GITEE_TOKEN`,
|
||||
`GITEE_USER`, and `GITEE_REPO`; OSS repair requires `OSS_ACCESS_KEY_ID`,
|
||||
`OSS_ACCESS_KEY_SECRET`, `OSS_ENDPOINT`, and `OSS_BUCKET` (with optional
|
||||
`OSS_PREFIX`) as Actions secrets. Missing credentials fail the selected repair
|
||||
closed.
|
||||
|
||||
## Handoff Checklist
|
||||
|
||||
|
||||
+11
-9
@@ -13,7 +13,9 @@
|
||||
3. workflow summary 会给出唯一的下一版本。把对应的精确 `CHANGELOG.md` 章节通过 PR 合入 `main`。
|
||||
4. 再次运行,改为 `release_operation=publish`,并输入 `PUBLISH beta` 或 `PUBLISH stable`。
|
||||
|
||||
`plan` 是纯只读操作,不创建 tag、预留版本号或生成包。CHANGELOG 合入期间若另一个发布先占用了该版本,`publish` 会重新分配并因 CHANGELOG 章节不匹配而拒绝,需要重新 plan。`publish` 会先再次确认 dispatch SHA 仍是当前 `main`、Code Admission 和平台治理均通过,再由唯一的 write job 使用 GitHub API 原子创建 annotated tag;同一次 run 随即进入既有的跨平台构建、GitHub/npm/OSS/Gitee 发布和 Homebrew PR DAG。内置 `GITHUB_TOKEN` 创建的 tag 不依赖第二条 workflow 被再次触发。
|
||||
`plan` 是纯只读操作,不创建 tag、预留版本号或生成包。CHANGELOG 合入期间若另一个发布先占用了该版本,`publish` 会重新分配并因 CHANGELOG 章节不匹配而拒绝,需要重新 plan。`publish` 会先再次确认 dispatch SHA 仍是当前 `main`、Code Admission 和平台治理均通过,再由唯一的 write job 使用 GitHub API 原子创建 annotated tag;同一次 run 随即进入既有的跨平台构建、GitHub/npm、可选 OSS/Gitee 发布和 Homebrew PR DAG。内置 `GITHUB_TOKEN` 创建的 tag 不依赖第二条 workflow 被再次触发。
|
||||
|
||||
OSS 镜像默认不参与发布 DAG,适用于尚未创建 Bucket 的仓库。云端封板会把当时的仓库变量 `ENABLE_OSS_MIRROR=true` 记录为不可变 tag 元数据 `OSS-Mirror: enabled`,否则记录为 `deferred`;后续发布和撤回只读取该 sealed policy,不读取变量的当前值。`enabled` 继续对缺失凭据、无效 Bucket、上传、pointer 和撤回失败保持 fail-closed;`deferred` 明确跳过不存在的渠道。为避免补发后撤回遗漏,deferred 版本暂不接受 `repair_oss_version`,启用 OSS 只影响后续新 tag,直到补齐可审计的不可变 repair 证明。
|
||||
|
||||
## 自动版本规则
|
||||
|
||||
@@ -35,8 +37,8 @@
|
||||
|
||||
1. 先创建永久 annotated tag `withdrawn/<version>`,记录原 tag object、commit、原因、申请人和 workflow run。这个墓碑是版本号永久占用记录,永不移动、永不删除。
|
||||
2. 先验证 Homebrew Formula;若它仍指向问题版本,先创建回退 PR,再继续其他渠道撤回。这样 PR 创建失败时只留下可安全续跑的墓碑,不会先造成渠道分裂。若 Formula 尚未指向问题版本或已经处于安全版本,则直接校验。
|
||||
3. GitHub Release 先标记为 withdrawn;npm 精确版本执行 `deprecate`,并把 `latest` / `beta` dist-tag 回退;OSS 先补齐回退版本资产,再移动 `latest.txt` / `beta.txt` 并删除问题版本目录;启用 Gitee 时同样先补齐回退 Release,再删除问题 Release 和 tag。
|
||||
4. npm、OSS 和 Gitee 均已验证安全后,删除 GitHub 上的问题 Release 和原 `v...` tag,并验证 `/releases/latest` 对正式版回到安全版本。若本次创建了 Homebrew PR,run 最后故意保持失败,直到另一名维护者审核合入;合入后,从新的 `main` 使用完全相同的 version、reason 和 confirmation 重跑并完成。永久 `withdrawn/v...` 墓碑始终保留。
|
||||
3. GitHub Release 先标记为 withdrawn;npm 精确版本执行 `deprecate`,并把 `latest` / `beta` dist-tag 回退;只有目标 tag 封存了 `OSS-Mirror: enabled` 时,OSS 才会先补齐回退版本资产,再移动 `latest.txt` / `beta.txt` 并删除问题版本目录;启用 Gitee 时同样先补齐回退 Release,再删除问题 Release 和 tag。
|
||||
4. npm 以及目标 tag 启用或发布时配置的镜像渠道均已验证安全后,删除 GitHub 上的问题 Release 和原 `v...` tag,并验证 `/releases/latest` 对正式版回到安全版本。若本次创建了 Homebrew PR,run 最后故意保持失败,直到另一名维护者审核合入;合入后,从新的 `main` 使用完全相同的 version、reason 和 confirmation 重跑并完成。永久 `withdrawn/v...` 墓碑始终保留。
|
||||
|
||||
GitHub、npm、OSS、Gitee 和 Homebrew 的“回滚”指新的安装、升级和渠道解析不再拿到问题版本。已经装到用户电脑上的二进制无法被服务端强制降级;用户必须重新安装回退版本、安装后续修复版,或使用 CLI 自带的本地 rollback 能力。npm 不执行 `unpublish`:问题版本保留明确的弃用警告,但 `latest` / `beta` 不再指向它;即使 registry 允许删除,已发布过的版本号也不会重新使用。
|
||||
|
||||
@@ -133,15 +135,15 @@ dws-release v1.2.3 --from-beta v1.2.3-beta.1 --publish
|
||||
- 日常 CI 和发布前都会对比“最新已交付正式版”的完整命令树;若长时间预检期间该 baseline 发生变化,会针对新的 baseline 重新比较。
|
||||
- GoReleaser 只构建;Darwin 重签、checksums 重算和 npm 安装验证通过后,才统一上传 GitHub Release 的最终产物。
|
||||
- 六个平台归档会逐个解包并核验二进制内嵌版本;公开资产集合、checksums 集合和 npm tarball integrity 都必须精确一致。npm tarball 固定由 npm `10.9.2` 打包,避免重跑时因 runner 自带 npm 漂移产生不同字节。
|
||||
- stable 发布到 npm `latest`,更新 OSS `latest.txt` 和共享安装脚本;prerelease 发布到 npm `beta`,只更新 OSS `beta.txt`,不会覆盖稳定入口。
|
||||
- stable 发布到 npm `latest`;prerelease 发布到 npm `beta`。启用 `ENABLE_OSS_MIRROR=true` 后,stable 同步 OSS `latest.txt` 和共享安装脚本,prerelease 只同步 OSS `beta.txt`,不会覆盖稳定入口。
|
||||
- Release workflow 使用一个最多容纳 100 个 pending run 的串行 publication queue;版本规划、云端封板、发布、恢复、修复和撤回共享同一发布锁。
|
||||
- 本地 tag push 失败时会删除本次新建的本地 tag。远端 tag 一旦创建,后续发布归 CI 所有;发布中途失败时走受保护恢复,禁止改 tag 指向或复用版本号。只有已经公开版本经过受保护的全渠道撤回并留下永久 `withdrawn/...` 墓碑后,撤回 workflow 才会在最后一步删除原 tag。
|
||||
|
||||
npm 补发只允许从默认分支触发 Release workflow 的 `repair_npm_version`。它只支持启用 immutable releases 后、由本流水线成功产出的公开 immutable release:目标必须是 `main` 历史中的 annotated tag,并且同 commit 的 `Build immutable GitHub Release` job 已成功。即使后续 npm 分发失败,这个独立的产物封存边界仍可作为补发依据。补发会用目标 commit 的 npm 模板重组包,逐平台核验资产和二进制版本,再发布到隔离的 `backfill` dist-tag,不会回滚 `latest` / `beta`。历史 mutable release 不进入自动补发路径,避免把可被替换的资产带入 npm。
|
||||
|
||||
OSS/Gitee 分发失败且 GitHub immutable Release、npm 已交付时,从受保护的默认分支触发
|
||||
已启用的 OSS 或 Gitee 分发失败且 GitHub immutable Release、npm 已交付时,从受保护的默认分支触发
|
||||
Release workflow,并且只填写 `repair_oss_version` 或 `repair_gitee_version` 之一。channel
|
||||
repair 会精确绑定失败 tag run 的最新 attempt;contract、构建、Developer ID 签名、
|
||||
repair 会精确绑定失败 tag run 的最新 attempt,且 OSS repair 要求 tag 的 sealed policy 为 `enabled`;contract、构建、Developer ID 签名、
|
||||
immutable GitHub 发布和 npm delivery 必须全部成功,且只能有一个 OSS/Gitee 下游失败,
|
||||
随后才会下载并重新校验原始资产、修复所选镜像。OSS repair 必须匹配失败的 OSS step;
|
||||
Gitee repair 还允许其 job 因该 OSS 失败而 skipped,此时只代表 Gitee backfill 成功,
|
||||
@@ -163,14 +165,14 @@ dws-release recover v1.2.3-beta.1
|
||||
- 输入精确绑定原 annotated tag object、commit 和失败的 sealed `Release` run;云端 run 还必须与 tag 内的 run ID、attempt、requester 完全一致,commit 必须仍在 `main` 历史中。
|
||||
- 目标只允许不存在 GitHub Release 或仍为 Draft;已经公开的版本不能全量重建:单个下游故障走对应的 channel repair,版本本身有问题则走受保护的全平台 withdrawal。
|
||||
- `release-recovery` environment 必须限制为受保护分支、配置至少一名 required reviewer,并禁止自审;workflow 会通过 API 复核这些设置,未配置时 fail closed。
|
||||
- 恢复复用正常的 contract、构建、Developer ID 签名、资产校验、immutable 发布、Homebrew、npm 和 OSS jobs,不存在 recovery 专用 publisher 或门禁跳过。
|
||||
- 恢复复用正常的 contract、构建、Developer ID 签名、资产校验、immutable 发布、Homebrew、npm,以及已启用的 OSS jobs,不存在 recovery 专用 publisher 或门禁跳过。
|
||||
- 如果 GitHub Release 已在 recovery 中封存、后续 Homebrew/npm 校验发生瞬时失败,只重跑该 run 的 failed jobs;流水线仅在隐藏 run marker、tag object、commit 和 finalized artifact 字节全部精确一致时复用公开 Release。
|
||||
|
||||
成功的默认分支恢复 run 会成为后续 beta → stable 和 stable baseline 验证的可审计交付证据;历史临时分支恢复仍只接受 reviewed manifest 中的固定证据。
|
||||
|
||||
云端 seal 后不要使用 GitHub 的 “Re-run failed jobs” 作为交付修复:annotated tag 永久绑定最初的 run attempt,普通 rerun 不会成为可接受的交付证据。GitHub Release 尚未公开时走上述 protected recovery;已经公开且仅 npm/OSS/Gitee 某一渠道失败时走对应 repair;版本内容本身有问题时走 withdrawal。
|
||||
|
||||
OSS 的 `latest.txt` / `beta.txt` 是镜像频道元数据;当前仓库安装器仍主要从 GitHub/Gitee 解析版本。发布和撤回仍把 OSS 作为受控分发渠道处理,保证一旦外部消费者接入该 pointer,也不会继续解析到已撤回版本。
|
||||
OSS 的 `latest.txt` / `beta.txt` 是镜像频道元数据;当前仓库安装器仍主要从 GitHub/Gitee 解析版本。启用 OSS 后,发布和撤回把它作为受控分发渠道处理,保证一旦外部消费者接入该 pointer,也不会继续解析到已撤回版本;未启用时两条流程都明确跳过不存在的 OSS 渠道。
|
||||
|
||||
Release workflow 会生成 Darwin/Linux 双架构 Formula,并分别为 stable/beta 打开 Homebrew PR;tap 的默认分支仍以独立审核合入为交付边界。撤回 workflow 使用相同模板和回退版本 checksums 打开反向 PR;问题 GitHub Release 会先被移除以阻止新安装,永久墓碑和 workflow 日志承担审计/续跑依据。
|
||||
|
||||
@@ -183,7 +185,7 @@ Release workflow 会生成 Darwin/Linux 双架构 Formula,并分别为 stable/
|
||||
- tag ruleset 还必须覆盖 `withdrawn/v*`:只允许受保护的撤回 workflow 创建墓碑,禁止更新或删除墓碑;同时应允许 Release workflow 创建新的 `v*`,允许撤回 workflow 在全部渠道回退后删除精确的问题 `v*`。若组织级规则阻止这两个 workflow 的预期动作,发布或撤回会 fail closed,不能靠手工移动 tag 绕过。
|
||||
- 配置 `RELEASE_GOVERNANCE_TOKEN` Actions secret,只授予目标仓库 `Administration: read`;内置 `GITHUB_TOKEN` 不具备 immutable-releases API 所需的仓库治理权限。每次本地预检和 tag workflow 都使用这一个身份进行 fail-closed 验证。
|
||||
- 配置 `APPLE_CERTIFICATE_P12_BASE64`、`APPLE_CERTIFICATE_PASSWORD` 和具备发布权限的 `NPM_TOKEN`;撤回还要求该 npm 身份能够执行 `deprecate` 和修改 dist-tag。
|
||||
- 配置 `OSS_ACCESS_KEY_ID`、`OSS_ACCESS_KEY_SECRET`、`OSS_ENDPOINT`、`OSS_BUCKET`,按需配置 `OSS_PREFIX`。撤回身份必须能够补齐安全版本资产、写 `latest.txt` / `beta.txt` 并删除问题版本前缀。
|
||||
- 启用 OSS 镜像时,先创建有效 Bucket,再设置仓库变量 `ENABLE_OSS_MIRROR=true`,并配置 `OSS_ACCESS_KEY_ID`、`OSS_ACCESS_KEY_SECRET`、`OSS_ENDPOINT`、`OSS_BUCKET`,按需配置 `OSS_PREFIX`。启用后发布保持 fail-closed;撤回身份必须能够补齐安全版本资产、写 `latest.txt` / `beta.txt` 并删除问题版本前缀。尚未 provision Bucket 时保持该变量未设置或不等于 `true`,新 tag 会封存 `OSS-Mirror: deferred` 并跳过 OSS;该版本不能通过现有 repair 流程事后改成启用。
|
||||
- 若启用 Gitee fallback,设置 `ENABLE_GITEE_UPLOAD_FALLBACK=true`,并配置 `GITEE_TOKEN`、`GITEE_USER`、`GITEE_REPO`;该身份必须能够创建和删除目标仓库的 Release 与 tag。
|
||||
- 单独配置 `HOMEBREW_PR_TOKEN`,优先使用仅授权本仓库且具备 `Contents: write`、`Pull requests: write` 的 fine-grained PAT;若组织策略不允许该账号使用 fine-grained PAT,则回退到仅带 `public_repo` scope 的专用 classic PAT。治理预检和 tag contract 会验证 token 身份、classic scope,并用 `[skip ci]` 临时分支和 draft PR 完成真实写权限 canary,随后立即关闭 PR、删除分支;任何清理失败都会 fail closed。门禁也会拒绝与治理 token 复用。
|
||||
- 创建 `release-recovery` environment,只允许受保护分支,设置 required reviewer、禁止自审并关闭管理员绕过。workflow 会读取 environment 的 required-reviewer、prevent-self-review 和 protected-branch 规则;规则缺失时紧急恢复会失败,正常 beta/stable tag 发布不受影响。
|
||||
|
||||
@@ -56,6 +56,7 @@ var (
|
||||
eventNewEventSource = newEventSource
|
||||
eventNewDingtalkSource = source.New
|
||||
eventResolveAccessToken = ResolveAuxiliaryAccessToken
|
||||
eventForceRefreshRejected = forceRefreshRejectedAccessToken
|
||||
eventBusRun = bus.Run
|
||||
eventReadyFDFromEnv = busctl.ReadyFDFromEnv
|
||||
eventResolvePersonal = resolvePersonalEventIdentity
|
||||
@@ -436,6 +437,9 @@ func newEventSource(_ context.Context, configDir, clientID, clientSecret string,
|
||||
AccessTokenProvider: func(ctx context.Context) (string, error) {
|
||||
return eventResolveAccessToken(ctx, configDir, "")
|
||||
},
|
||||
ForceRefreshToken: func(ctx context.Context, rejectedToken string) (string, error) {
|
||||
return eventForceRefreshRejected(ctx, configDir, rejectedToken)
|
||||
},
|
||||
SourceID: eventStreamSourceID(streamOpts.SourceID),
|
||||
Mode: streamOpts.Mode,
|
||||
ClientID: portalClientID,
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
package app
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
|
||||
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/source"
|
||||
)
|
||||
|
||||
// TestCrossPlatformCoverageNewEventSourceWiresForceRefreshRejectedToken asserts the portal ticket
|
||||
// source receives a ForceRefreshToken callback that forwards the actual
|
||||
// rejected token into the app-level compare-and-refresh chain.
|
||||
func TestCrossPlatformCoverageNewEventSourceWiresForceRefreshRejectedToken(t *testing.T) {
|
||||
oldNew, oldRefresh := eventNewDingtalkSource, eventForceRefreshRejected
|
||||
t.Cleanup(func() { eventNewDingtalkSource, eventForceRefreshRejected = oldNew, oldRefresh })
|
||||
|
||||
var captured source.Config
|
||||
eventNewDingtalkSource = func(cfg source.Config, _ ...source.SourceOption) (*source.DingtalkSource, error) {
|
||||
captured = cfg
|
||||
return &source.DingtalkSource{}, nil
|
||||
}
|
||||
var gotDir, gotRejected string
|
||||
eventForceRefreshRejected = func(_ context.Context, configDir, rejectedToken string) (string, error) {
|
||||
gotDir, gotRejected = configDir, rejectedToken
|
||||
return "fresh", nil
|
||||
}
|
||||
if _, err := newEventSource(context.Background(), "config-dir", "client", "secret", eventStreamTicketOptions{Mode: "custom"}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if captured.PortalTicket == nil || captured.PortalTicket.ForceRefreshToken == nil {
|
||||
t.Fatal("ForceRefreshToken not wired into portal ticket config")
|
||||
}
|
||||
tok, err := captured.PortalTicket.ForceRefreshToken(context.Background(), "rejected-token")
|
||||
if err != nil || tok != "fresh" {
|
||||
t.Fatalf("force refresh = %q, %v", tok, err)
|
||||
}
|
||||
if gotDir != "config-dir" || gotRejected != "rejected-token" {
|
||||
t.Fatalf("wiring passed dir %q rejected %q", gotDir, gotRejected)
|
||||
}
|
||||
|
||||
fail := errors.New("refresh failed")
|
||||
eventForceRefreshRejected = func(context.Context, string, string) (string, error) { return "", fail }
|
||||
if _, err := captured.PortalTicket.ForceRefreshToken(context.Background(), "x"); !errors.Is(err, fail) {
|
||||
t.Fatalf("refresh error = %v", err)
|
||||
}
|
||||
}
|
||||
@@ -131,6 +131,7 @@ var (
|
||||
personalFindProcess = os.FindProcess
|
||||
personalSignalProcess = (*os.Process).Signal
|
||||
personalResolveAuxiliaryAccessToken = ResolveAuxiliaryAccessToken
|
||||
personalForceRefreshRejectedToken = forceRefreshRejectedAccessToken
|
||||
personalLoadTokenData = authpkg.LoadTokenData
|
||||
personalClientID = authpkg.ClientID
|
||||
personalResolveAppCredentialsStrict = authpkg.ResolveAppCredentialsStrict
|
||||
@@ -832,6 +833,9 @@ func newPersonalStreamSource(ctx context.Context, opts personalStreamSourceOptio
|
||||
AccessTokenProvider: func(ctx context.Context) (string, error) {
|
||||
return personalResolveAuxiliaryAccessToken(ctx, opts.ConfigDir, "")
|
||||
},
|
||||
ForceRefreshToken: func(ctx context.Context, rejectedToken string) (string, error) {
|
||||
return personalForceRefreshRejectedToken(ctx, opts.ConfigDir, rejectedToken)
|
||||
},
|
||||
ClientID: clientID,
|
||||
ClientSecret: clientSecret,
|
||||
SourceID: opts.Identity.SourceID,
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
package app
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
dwsevent "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event"
|
||||
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event/personal"
|
||||
)
|
||||
|
||||
// TestCrossPlatformCoverageNewPersonalStreamSourceWiresForceRefreshRejectedToken asserts the
|
||||
// personal stream source receives a ForceRefreshToken callback that forwards
|
||||
// the rejected token into the app-level compare-and-refresh chain.
|
||||
func TestCrossPlatformCoverageNewPersonalStreamSourceWiresForceRefreshRejectedToken(t *testing.T) {
|
||||
oldAux := personalResolveAuxiliaryAccessToken
|
||||
oldRefresh := personalForceRefreshRejectedToken
|
||||
t.Cleanup(func() {
|
||||
personalResolveAuxiliaryAccessToken = oldAux
|
||||
personalForceRefreshRejectedToken = oldRefresh
|
||||
})
|
||||
personalResolveAuxiliaryAccessToken = func(context.Context, string, string) (string, error) {
|
||||
return "old-token", nil
|
||||
}
|
||||
refreshErr := errors.New("refresh rejected")
|
||||
var gotDir, gotRejected string
|
||||
personalForceRefreshRejectedToken = func(_ context.Context, configDir, rejectedToken string) (string, error) {
|
||||
gotDir, gotRejected = configDir, rejectedToken
|
||||
return "", refreshErr
|
||||
}
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
src, err := newPersonalStreamSource(context.Background(), personalStreamSourceOptions{
|
||||
ConfigDir: "config-dir",
|
||||
Identity: personal.Identity{ClientID: "client", SourceID: "source"},
|
||||
TicketURL: srv.URL,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// The 401 ticket response routes the rejected token through the wired
|
||||
// ForceRefreshToken; the unknown refresh failure stays fatal.
|
||||
if err := src.Start(context.Background(), func(*dwsevent.RawEvent) {}); !errors.Is(err, refreshErr) {
|
||||
t.Fatalf("Start() error = %v, want wrapped refresh error", err)
|
||||
}
|
||||
if gotDir != "config-dir" || gotRejected != "old-token" {
|
||||
t.Fatalf("refresh wiring got dir %q rejected %q", gotDir, gotRejected)
|
||||
}
|
||||
}
|
||||
@@ -722,6 +722,6 @@ func isInvalidGrantError(err error) bool {
|
||||
if err == nil {
|
||||
return false
|
||||
}
|
||||
msg := strings.ToLower(err.Error())
|
||||
msg := strings.ToLower(err.Error() + " " + httpStatusResponseBody(err))
|
||||
return strings.Contains(msg, "invalid_grant") || (strings.Contains(msg, "code") && strings.Contains(msg, "expired"))
|
||||
}
|
||||
|
||||
@@ -237,12 +237,21 @@ func (p *OAuthProvider) postJSON(ctx context.Context, endpoint string, body any)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
data, err := io.ReadAll(io.LimitReader(resp.Body, config.MaxResponseBodySize))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("reading response: %w", err)
|
||||
}
|
||||
data, readErr := io.ReadAll(io.LimitReader(resp.Body, config.MaxResponseBodySize))
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return nil, fmt.Errorf("HTTP %d: %s", resp.StatusCode, truncateBody(data, 200))
|
||||
// Preserve structured HTTP status semantics even when the response
|
||||
// body is truncated. The body is diagnostic-only here, so read it
|
||||
// best-effort and classify retryability from the status code.
|
||||
if readErr != nil {
|
||||
data = nil
|
||||
}
|
||||
return nil, &HTTPStatusError{
|
||||
StatusCode: resp.StatusCode,
|
||||
responseBody: truncateBody(data, 200),
|
||||
}
|
||||
}
|
||||
if readErr != nil {
|
||||
return nil, fmt.Errorf("reading response: %w", readErr)
|
||||
}
|
||||
return data, nil
|
||||
}
|
||||
|
||||
@@ -0,0 +1,69 @@
|
||||
// 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.
|
||||
|
||||
package auth
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"testing"
|
||||
)
|
||||
|
||||
type postJSONRoundTripFunc func(*http.Request) (*http.Response, error)
|
||||
|
||||
func (f postJSONRoundTripFunc) RoundTrip(req *http.Request) (*http.Response, error) {
|
||||
return f(req)
|
||||
}
|
||||
|
||||
type oauthBrokenBody struct{}
|
||||
|
||||
func (oauthBrokenBody) Read([]byte) (int, error) { return 0, io.ErrUnexpectedEOF }
|
||||
func (oauthBrokenBody) Close() error { return nil }
|
||||
|
||||
func TestCrossPlatformCoveragePostJSONTruncatedErrorBodyKeepsHTTPStatus(t *testing.T) {
|
||||
for _, status := range []int{http.StatusTooManyRequests, http.StatusServiceUnavailable} {
|
||||
t.Run(http.StatusText(status), func(t *testing.T) {
|
||||
provider := &OAuthProvider{httpClient: &http.Client{Transport: postJSONRoundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: status, Body: oauthBrokenBody{}, Header: make(http.Header)}, nil
|
||||
})}}
|
||||
|
||||
_, err := provider.postJSON(context.Background(), "https://oauth.test/token", map[string]string{"grantType": "refresh_token"})
|
||||
var statusErr *HTTPStatusError
|
||||
if !errors.As(err, &statusErr) || statusErr.StatusCode != status {
|
||||
t.Fatalf("postJSON() error = %v, want HTTPStatusError %d", err, status)
|
||||
}
|
||||
if errors.Is(err, io.ErrUnexpectedEOF) {
|
||||
t.Fatalf("HTTP status error should not expose diagnostic body read failure: %v", err)
|
||||
}
|
||||
if got := ClassifyRefreshFailure(err); got != RefreshFailureTransient {
|
||||
t.Fatalf("ClassifyRefreshFailure() = %s, want transient", got)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePostJSONOKTruncatedBodyIsTransient(t *testing.T) {
|
||||
provider := &OAuthProvider{httpClient: &http.Client{Transport: postJSONRoundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: http.StatusOK, Body: oauthBrokenBody{}, Header: make(http.Header)}, nil
|
||||
})}}
|
||||
|
||||
_, err := provider.postJSON(context.Background(), "https://oauth.test/token", map[string]string{"grantType": "refresh_token"})
|
||||
if !errors.Is(err, io.ErrUnexpectedEOF) {
|
||||
t.Fatalf("postJSON() error = %v, want io.ErrUnexpectedEOF", err)
|
||||
}
|
||||
if got := ClassifyRefreshFailure(err); got != RefreshFailureTransient {
|
||||
t.Fatalf("ClassifyRefreshFailure() = %s, want transient", got)
|
||||
}
|
||||
}
|
||||
@@ -18,6 +18,7 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"html"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net"
|
||||
@@ -301,7 +302,7 @@ func (p *OAuthProvider) Login(ctx context.Context, force bool) (*TokenData, erro
|
||||
callbackTokenMu.Unlock()
|
||||
|
||||
w.Header().Set("Content-Type", "text/html; charset=utf-8")
|
||||
_, _ = fmt.Fprintf(w, "<html><body><h1>授权失败</h1><p>%s</p></body></html>", exchangeErr.Error())
|
||||
_, _ = fmt.Fprintf(w, "<html><body><h1>授权失败</h1><p>%s</p></body></html>", html.EscapeString(oauthExchangeDisplayError(exchangeErr)))
|
||||
select {
|
||||
case resultCh <- callbackResult{err: exchangeErr}:
|
||||
default:
|
||||
@@ -642,6 +643,14 @@ continueLogin:
|
||||
return tokenData, nil
|
||||
}
|
||||
|
||||
func oauthExchangeDisplayError(err error) string {
|
||||
var statusErr *HTTPStatusError
|
||||
if errors.As(err, &statusErr) && statusErr != nil {
|
||||
return fmt.Sprintf("HTTP %d: token exchange failed", statusErr.StatusCode)
|
||||
}
|
||||
return err.Error()
|
||||
}
|
||||
|
||||
// GetTokenSnapshot returns a valid token together with its expiry metadata.
|
||||
// Storage and refresh failures retain their original cause; only a confirmed
|
||||
// missing credential is reported as ErrTokenDataNotFound.
|
||||
@@ -665,7 +674,12 @@ func (p *OAuthProvider) GetTokenSnapshot(ctx context.Context) (*TokenData, error
|
||||
if rErr == nil {
|
||||
return refreshed, nil
|
||||
}
|
||||
_ = oauthMarkProfile(p.configDir, TokenProfileSelector(data), ProfileStatusExpired)
|
||||
// A network, timeout, rate-limit or 5xx failure does not invalidate the
|
||||
// refresh credential. Keep the profile active so a long-running source
|
||||
// can retry after backoff. Terminal and unknown failures remain fatal.
|
||||
if ClassifyRefreshFailure(rErr) != RefreshFailureTransient {
|
||||
_ = oauthMarkProfile(p.configDir, TokenProfileSelector(data), ProfileStatusExpired)
|
||||
}
|
||||
if p.logger != nil {
|
||||
p.logger.Warn(i18n.T("refresh_token 刷新失败"), "error", rErr)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,89 @@
|
||||
// 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.
|
||||
|
||||
package auth
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
)
|
||||
|
||||
// RefreshFailureClass separates refresh failures that may recover after a
|
||||
// delay from failures that require new credentials or local intervention.
|
||||
type RefreshFailureClass string
|
||||
|
||||
const (
|
||||
RefreshFailureUnknown RefreshFailureClass = "unknown"
|
||||
RefreshFailureTransient RefreshFailureClass = "transient"
|
||||
RefreshFailureTerminal RefreshFailureClass = "terminal"
|
||||
)
|
||||
|
||||
// HTTPStatusError preserves an OAuth endpoint status for structured retry
|
||||
// decisions without copying an untrusted response body into logs.
|
||||
type HTTPStatusError struct {
|
||||
StatusCode int
|
||||
responseBody string
|
||||
}
|
||||
|
||||
func (e *HTTPStatusError) Error() string {
|
||||
if e == nil {
|
||||
return "OAuth endpoint request failed"
|
||||
}
|
||||
return fmt.Sprintf("HTTP %d", e.StatusCode)
|
||||
}
|
||||
|
||||
func httpStatusResponseBody(err error) string {
|
||||
var statusErr *HTTPStatusError
|
||||
if !errors.As(err, &statusErr) || statusErr == nil {
|
||||
return ""
|
||||
}
|
||||
return statusErr.responseBody
|
||||
}
|
||||
|
||||
// ClassifyRefreshFailure uses only structured transport and HTTP signals.
|
||||
// Unknown errors, including parse, keychain and persistence failures, remain
|
||||
// fatal so a long-running source cannot retry an error that needs user action.
|
||||
func ClassifyRefreshFailure(err error) RefreshFailureClass {
|
||||
if err == nil {
|
||||
return RefreshFailureUnknown
|
||||
}
|
||||
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
|
||||
return RefreshFailureTransient
|
||||
}
|
||||
if errors.Is(err, io.ErrUnexpectedEOF) {
|
||||
return RefreshFailureTransient
|
||||
}
|
||||
var netErr net.Error
|
||||
if errors.As(err, &netErr) {
|
||||
return RefreshFailureTransient
|
||||
}
|
||||
var statusErr *HTTPStatusError
|
||||
if !errors.As(err, &statusErr) || statusErr == nil {
|
||||
return RefreshFailureUnknown
|
||||
}
|
||||
if statusErr.StatusCode == http.StatusRequestTimeout ||
|
||||
statusErr.StatusCode == http.StatusTooManyRequests ||
|
||||
statusErr.StatusCode >= http.StatusInternalServerError {
|
||||
return RefreshFailureTransient
|
||||
}
|
||||
if statusErr.StatusCode == http.StatusBadRequest ||
|
||||
statusErr.StatusCode == http.StatusUnauthorized ||
|
||||
statusErr.StatusCode == http.StatusForbidden {
|
||||
return RefreshFailureTerminal
|
||||
}
|
||||
return RefreshFailureUnknown
|
||||
}
|
||||
@@ -0,0 +1,155 @@
|
||||
// Copyright 2026 Alibaba Group
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
|
||||
package auth
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"net/url"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/pkg/edition"
|
||||
)
|
||||
|
||||
func TestCrossPlatformCoverageClassifyRefreshFailureUsesStructuredSignals(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
err error
|
||||
want RefreshFailureClass
|
||||
}{
|
||||
{name: "deadline", err: context.DeadlineExceeded, want: RefreshFailureTransient},
|
||||
{name: "network", err: &url.Error{Op: "Post", URL: "https://oauth.test", Err: context.DeadlineExceeded}, want: RefreshFailureTransient},
|
||||
{name: "request timeout", err: &HTTPStatusError{StatusCode: http.StatusRequestTimeout}, want: RefreshFailureTransient},
|
||||
{name: "rate limited", err: &HTTPStatusError{StatusCode: http.StatusTooManyRequests}, want: RefreshFailureTransient},
|
||||
{name: "server unavailable", err: &HTTPStatusError{StatusCode: http.StatusServiceUnavailable}, want: RefreshFailureTransient},
|
||||
{name: "refresh rejected", err: &HTTPStatusError{StatusCode: http.StatusUnauthorized}, want: RefreshFailureTerminal},
|
||||
{name: "invalid grant", err: &HTTPStatusError{StatusCode: http.StatusBadRequest}, want: RefreshFailureTerminal},
|
||||
{name: "forbidden", err: &HTTPStatusError{StatusCode: http.StatusForbidden}, want: RefreshFailureTerminal},
|
||||
{name: "local persistence", err: errors.New("save refreshed token failed"), want: RefreshFailureUnknown},
|
||||
{name: "nil error", err: nil, want: RefreshFailureUnknown},
|
||||
{name: "dns failure", err: &net.DNSError{Err: "no such host", Name: "oauth.test"}, want: RefreshFailureTransient},
|
||||
{name: "redirect status", err: &HTTPStatusError{StatusCode: http.StatusFound}, want: RefreshFailureUnknown},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
if got := ClassifyRefreshFailure(tt.err); got != tt.want {
|
||||
t.Fatalf("ClassifyRefreshFailure() = %q, want %q", got, tt.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoverageHTTPStatusErrorRetainsStatusThroughWrapping(t *testing.T) {
|
||||
want := &HTTPStatusError{StatusCode: http.StatusTooManyRequests}
|
||||
err := errors.Join(errors.New("refresh failed"), want)
|
||||
if got := ClassifyRefreshFailure(err); got != RefreshFailureTransient {
|
||||
t.Fatalf("ClassifyRefreshFailure() = %q, want transient", got)
|
||||
}
|
||||
var statusErr *HTTPStatusError
|
||||
if !errors.As(err, &statusErr) || statusErr.StatusCode != http.StatusTooManyRequests {
|
||||
t.Fatalf("HTTP status error not retained: %v", err)
|
||||
}
|
||||
if got, want := statusErr.Error(), "HTTP 429"; got != want {
|
||||
t.Fatalf("HTTP status error = %q, want %q", got, want)
|
||||
}
|
||||
var nilStatus *HTTPStatusError
|
||||
if got, want := nilStatus.Error(), "OAuth endpoint request failed"; got != want {
|
||||
t.Fatalf("nil HTTP status error = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoverageOAuthExchangeDisplayErrorFallsBackToPlainError(t *testing.T) {
|
||||
if got, want := oauthExchangeDisplayError(&HTTPStatusError{StatusCode: http.StatusBadGateway}), "HTTP 502: token exchange failed"; got != want {
|
||||
t.Fatalf("status display error = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := oauthExchangeDisplayError(errors.New("exchange failed")), "exchange failed"; got != want {
|
||||
t.Fatalf("plain display error = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePostJSONClassifiesStatusWithoutLoggingResponseBody(t *testing.T) {
|
||||
const secretBody = `{"refreshToken":"must-not-reach-logs"}`
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
w.WriteHeader(http.StatusServiceUnavailable)
|
||||
_, _ = w.Write([]byte(secretBody))
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
provider := &OAuthProvider{httpClient: server.Client()}
|
||||
_, err := provider.postJSON(context.Background(), server.URL, map[string]string{"grantType": "refresh_token"})
|
||||
if got := ClassifyRefreshFailure(err); got != RefreshFailureTransient {
|
||||
t.Fatalf("ClassifyRefreshFailure() = %q, want transient: %v", got, err)
|
||||
}
|
||||
if strings.Contains(err.Error(), "must-not-reach-logs") {
|
||||
t.Fatalf("postJSON error leaked response body: %v", err)
|
||||
}
|
||||
if got := httpStatusResponseBody(err); !strings.Contains(got, "must-not-reach-logs") {
|
||||
t.Fatalf("postJSON did not retain bounded response details for internal classification: %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoverageGetTokenSnapshotOnlyExpiresProfileForNonTransientRefreshFailures(t *testing.T) {
|
||||
oldLoad := oauthLoadToken
|
||||
oldLoadLocked := oauthLoadTokenLocked
|
||||
oldAcquire := oauthAcquireLock
|
||||
oldRefresh := oauthRefreshToken
|
||||
oldMark := oauthMarkProfile
|
||||
oldEdition := edition.Get()
|
||||
t.Cleanup(func() {
|
||||
oauthLoadToken = oldLoad
|
||||
oauthLoadTokenLocked = oldLoadLocked
|
||||
oauthAcquireLock = oldAcquire
|
||||
oauthRefreshToken = oldRefresh
|
||||
oauthMarkProfile = oldMark
|
||||
edition.Override(oldEdition)
|
||||
})
|
||||
edition.Override(&edition.Hooks{})
|
||||
|
||||
expired := &TokenData{
|
||||
AccessToken: "expired-access",
|
||||
ExpiresAt: time.Now().Add(-time.Hour),
|
||||
RefreshToken: "refresh",
|
||||
RefreshExpAt: time.Now().Add(time.Hour),
|
||||
CorpID: "corp",
|
||||
UserID: "user",
|
||||
}
|
||||
oauthLoadToken = func(string) (*TokenData, error) { return expired, nil }
|
||||
oauthLoadTokenLocked = func(string, string) (*TokenData, error) { return expired, nil }
|
||||
oauthAcquireLock = func(context.Context, string) (*DualLock, error) { return &DualLock{}, nil }
|
||||
|
||||
markCalls := 0
|
||||
oauthMarkProfile = func(_, _, status string) error {
|
||||
if status != ProfileStatusExpired {
|
||||
t.Fatalf("profile status = %q, want %q", status, ProfileStatusExpired)
|
||||
}
|
||||
markCalls++
|
||||
return nil
|
||||
}
|
||||
provider := NewOAuthProvider(t.TempDir(), nil)
|
||||
|
||||
oauthRefreshToken = func(*OAuthProvider, context.Context, *TokenData) (*TokenData, error) {
|
||||
return nil, &HTTPStatusError{StatusCode: http.StatusServiceUnavailable}
|
||||
}
|
||||
if _, err := provider.GetTokenSnapshot(context.Background()); ClassifyRefreshFailure(err) != RefreshFailureTransient {
|
||||
t.Fatalf("transient refresh error = %v", err)
|
||||
}
|
||||
if markCalls != 0 {
|
||||
t.Fatalf("transient refresh marked profile expired %d times", markCalls)
|
||||
}
|
||||
|
||||
oauthRefreshToken = func(*OAuthProvider, context.Context, *TokenData) (*TokenData, error) {
|
||||
return nil, &HTTPStatusError{StatusCode: http.StatusUnauthorized}
|
||||
}
|
||||
if _, err := provider.GetTokenSnapshot(context.Background()); ClassifyRefreshFailure(err) != RefreshFailureTerminal {
|
||||
t.Fatalf("terminal refresh error = %v", err)
|
||||
}
|
||||
if markCalls != 1 {
|
||||
t.Fatalf("terminal refresh marked profile expired %d times, want 1", markCalls)
|
||||
}
|
||||
}
|
||||
@@ -414,7 +414,7 @@ func TestCrossPlatformCoveragePortalStartEndToEndAndFailures(t *testing.T) {
|
||||
return s
|
||||
}
|
||||
network := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) { return nil, errSourceInjected })}
|
||||
s := makeSource(&PortalTicketConfig{TicketURL: "https://x", AccessToken: "t", SourceID: "s", HTTPClient: network})
|
||||
s := makeSource(&PortalTicketConfig{TicketURL: "https://x", AccessToken: "t", SourceID: "s", HTTPClient: network, DisableReconnect: true})
|
||||
if err := s.Start(context.Background(), func(*dwsevent.RawEvent) {}); err == nil {
|
||||
t.Fatal("ticket failure expected")
|
||||
}
|
||||
@@ -483,7 +483,7 @@ func TestCrossPlatformCoveragePortalStartHandshakeReadAndAckErrors(t *testing.T)
|
||||
makeSource := func(endpoint string) *DingtalkSource {
|
||||
s, err := New(Config{PortalTicket: &PortalTicketConfig{
|
||||
TicketURL: "https://ticket", AccessToken: "token", SourceID: "source",
|
||||
HTTPClient: staticPersonalTicketClient(endpoint, ""),
|
||||
HTTPClient: staticPersonalTicketClient(endpoint, ""), DisableReconnect: true,
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
@@ -27,6 +27,7 @@ import (
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
authpkg "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/auth"
|
||||
dwsevent "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event"
|
||||
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/pkg/config"
|
||||
"github.com/gorilla/websocket"
|
||||
@@ -43,6 +44,7 @@ const (
|
||||
type PersonalConfig struct {
|
||||
AccessToken string
|
||||
AccessTokenProvider AccessTokenProvider
|
||||
ForceRefreshToken ForceRefreshTokenFn
|
||||
ClientID string
|
||||
ClientSecret string
|
||||
SourceID string
|
||||
@@ -57,6 +59,13 @@ type PersonalConfig struct {
|
||||
|
||||
type AccessTokenProvider func(context.Context) (string, error)
|
||||
|
||||
// ForceRefreshTokenFn rotates an access token that the server has just
|
||||
// rejected (HTTP 401). It receives the exact rejected token so the caller's
|
||||
// compare-and-refresh logic can skip the refresh when another goroutine has
|
||||
// already rotated it, and returns the fresh token to retry with. Optional:
|
||||
// when nil a 401 stays fatal, matching the previous behavior.
|
||||
type ForceRefreshTokenFn func(ctx context.Context, rejectedToken string) (string, error)
|
||||
|
||||
type PersonalSource struct {
|
||||
cfg PersonalConfig
|
||||
machine *Machine
|
||||
@@ -197,8 +206,29 @@ func (s *PersonalSource) runAttempt(ctx context.Context, emit dwsevent.EmitFn) (
|
||||
func (s *PersonalSource) fetchTicket(ctx context.Context) (*ticketResponse, error) {
|
||||
accessToken, err := resolveSourceAccessToken(ctx, s.cfg.AccessTokenProvider, s.cfg.AccessToken, "personal source")
|
||||
if err != nil {
|
||||
// Transient provider failures (network, 429, 5xx) must not kill a
|
||||
// long-running source; the reconnect loop retries after backoff.
|
||||
if authpkg.ClassifyRefreshFailure(err) == authpkg.RefreshFailureTransient {
|
||||
return nil, retryPersonal(err)
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
ticket, status, err := s.fetchTicketAttempt(ctx, accessToken)
|
||||
if status == http.StatusUnauthorized && s.cfg.ForceRefreshToken != nil {
|
||||
refreshed, refreshErr := refreshRejectedSourceToken(ctx, s.cfg.ForceRefreshToken, accessToken, "personal source", err)
|
||||
if refreshErr != nil {
|
||||
if authpkg.ClassifyRefreshFailure(refreshErr) == authpkg.RefreshFailureTransient {
|
||||
return nil, retryPersonal(refreshErr)
|
||||
}
|
||||
return nil, refreshErr
|
||||
}
|
||||
// Retry once with the freshly rotated token; a second 401 stays fatal.
|
||||
ticket, _, err = s.fetchTicketAttempt(ctx, refreshed)
|
||||
}
|
||||
return ticket, err
|
||||
}
|
||||
|
||||
func (s *PersonalSource) fetchTicketAttempt(ctx context.Context, accessToken string) (*ticketResponse, int, error) {
|
||||
body := map[string]any{
|
||||
"sourceId": s.cfg.SourceID,
|
||||
"mode": s.cfg.TicketMode,
|
||||
@@ -210,7 +240,7 @@ func (s *PersonalSource) fetchTicket(ctx context.Context) (*ticketResponse, erro
|
||||
b, _ := json.Marshal(body)
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, s.cfg.TicketURL, bytes.NewReader(b))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("personal source: create ticket request: %w", err)
|
||||
return nil, 0, fmt.Errorf("personal source: create ticket request: %w", err)
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
@@ -221,28 +251,50 @@ func (s *PersonalSource) fetchTicket(ctx context.Context) (*ticketResponse, erro
|
||||
|
||||
resp, err := s.cfg.HTTPClient.Do(req)
|
||||
if err != nil {
|
||||
return nil, retryPersonal(fmt.Errorf("personal source: fetch ticket: %w", err))
|
||||
return nil, 0, retryPersonal(fmt.Errorf("personal source: fetch ticket: %w", err))
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
data, err := io.ReadAll(io.LimitReader(resp.Body, config.MaxResponseBodySize))
|
||||
if err != nil {
|
||||
return nil, retryPersonal(fmt.Errorf("personal source: read ticket response: %w", err))
|
||||
}
|
||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||
// Classify by status before touching the body: a truncated error body
|
||||
// must not upgrade a fatal status (notably 401) into a retryable
|
||||
// error, or the outer reconnect loop would bypass the single
|
||||
// refresh-retry guard. The body is not used here, so drain it only
|
||||
// best-effort for connection reuse.
|
||||
_, _ = io.Copy(io.Discard, io.LimitReader(resp.Body, config.MaxResponseBodySize))
|
||||
err := fmt.Errorf("personal source: ticket HTTP %d", resp.StatusCode)
|
||||
if retryableTicketStatus(resp.StatusCode) {
|
||||
return nil, retryPersonal(err)
|
||||
return nil, resp.StatusCode, retryPersonal(err)
|
||||
}
|
||||
return nil, err
|
||||
return nil, resp.StatusCode, err
|
||||
}
|
||||
data, err := io.ReadAll(io.LimitReader(resp.Body, config.MaxResponseBodySize))
|
||||
if err != nil {
|
||||
return nil, resp.StatusCode, retryPersonal(fmt.Errorf("personal source: read ticket response: %w", err))
|
||||
}
|
||||
ticket, err := decodeTicket(data)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
return nil, resp.StatusCode, err
|
||||
}
|
||||
if ticket.Endpoint == "" || ticket.Ticket == "" {
|
||||
return nil, errors.New("personal source: ticket response missing endpoint or ticket")
|
||||
return nil, resp.StatusCode, errors.New("personal source: ticket response missing endpoint or ticket")
|
||||
}
|
||||
return ticket, nil
|
||||
return ticket, resp.StatusCode, nil
|
||||
}
|
||||
|
||||
// refreshRejectedSourceToken funnels a server-side 401 into the optional
|
||||
// force-refresh callback. It hands the actual rejected token to the caller's
|
||||
// compare-and-refresh logic and returns the rotated token for an immediate
|
||||
// retry. Refresh failures keep the original 401 as context instead of being
|
||||
// dropped.
|
||||
func refreshRejectedSourceToken(ctx context.Context, refresh ForceRefreshTokenFn, rejectedToken, component string, cause error) (string, error) {
|
||||
token, err := refresh(ctx, rejectedToken)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("%s: refresh rejected access token: %w", component, errors.Join(cause, err))
|
||||
}
|
||||
if token = strings.TrimSpace(token); token == "" {
|
||||
return "", fmt.Errorf("%s: refresh rejected access token returned empty token: %w", component, cause)
|
||||
}
|
||||
return token, nil
|
||||
}
|
||||
|
||||
func resolveSourceAccessToken(ctx context.Context, provider AccessTokenProvider, fallback, component string) (string, error) {
|
||||
@@ -381,6 +433,13 @@ func isRetryablePersonalError(err error) bool {
|
||||
func personalRetryLogError(err error) string {
|
||||
message := err.Error()
|
||||
switch {
|
||||
case strings.Contains(message, "resolve access token"), strings.Contains(message, "refresh rejected access token"):
|
||||
// Token resolution/refresh errors may carry provider details; log
|
||||
// only the structured HTTP status.
|
||||
if status := refreshHTTPStatus(err); status != 0 {
|
||||
return fmt.Sprintf("personal source: token refresh HTTP %d", status)
|
||||
}
|
||||
return "personal source: token refresh: temporary network error"
|
||||
case strings.Contains(message, "ticket HTTP"):
|
||||
return message
|
||||
case strings.Contains(message, "fetch ticket"):
|
||||
@@ -398,6 +457,14 @@ func personalRetryLogError(err error) string {
|
||||
}
|
||||
}
|
||||
|
||||
func refreshHTTPStatus(err error) int {
|
||||
var statusErr *authpkg.HTTPStatusError
|
||||
if !errors.As(err, &statusErr) || statusErr == nil {
|
||||
return 0
|
||||
}
|
||||
return statusErr.StatusCode
|
||||
}
|
||||
|
||||
func retryableTicketStatus(status int) bool {
|
||||
return status == http.StatusRequestTimeout ||
|
||||
status == http.StatusTooManyRequests ||
|
||||
|
||||
@@ -20,11 +20,13 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
authpkg "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/auth"
|
||||
dwsevent "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event"
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/open-dingtalk/dingtalk-stream-sdk-go/payload"
|
||||
@@ -33,6 +35,8 @@ import (
|
||||
const (
|
||||
PortalTicketModeNormal = "normal"
|
||||
PortalTicketModeCustom = "custom"
|
||||
portalReconnectMin = time.Second
|
||||
portalReconnectMax = 30 * time.Second
|
||||
)
|
||||
|
||||
// PortalTicketConfig describes the portal-managed user Stream ticket flow.
|
||||
@@ -42,12 +46,45 @@ type PortalTicketConfig struct {
|
||||
TicketURL string
|
||||
AccessToken string
|
||||
AccessTokenProvider AccessTokenProvider
|
||||
ForceRefreshToken ForceRefreshTokenFn
|
||||
SourceID string
|
||||
Mode string
|
||||
ClientID string
|
||||
ClientSecret string
|
||||
UserAgent string
|
||||
HTTPClient *http.Client
|
||||
WebSocketDialer *websocket.Dialer
|
||||
ReconnectMin time.Duration
|
||||
ReconnectMax time.Duration
|
||||
DisableReconnect bool
|
||||
}
|
||||
|
||||
// portalStageError tags a portal stream failure with the stage it happened
|
||||
// in and whether the reconnect loop may retry it. Error() stays free of
|
||||
// untrusted response content so it is safe to log on every reconnect.
|
||||
type portalStageError struct {
|
||||
stage string
|
||||
status int
|
||||
retryable bool
|
||||
cause error
|
||||
}
|
||||
|
||||
func (e *portalStageError) Error() string {
|
||||
if e == nil {
|
||||
return "source: portal stream failed"
|
||||
}
|
||||
message := "source: portal " + strings.ReplaceAll(strings.TrimSpace(e.stage), "_", " ") + " failed"
|
||||
if e.status != 0 {
|
||||
message += fmt.Sprintf(" (HTTP %d)", e.status)
|
||||
}
|
||||
return message
|
||||
}
|
||||
|
||||
func (e *portalStageError) Unwrap() error {
|
||||
if e == nil {
|
||||
return nil
|
||||
}
|
||||
return e.cause
|
||||
}
|
||||
|
||||
var portalWriteMessage = func(conn *websocket.Conn, messageType int, data []byte) error {
|
||||
@@ -97,48 +134,102 @@ func normalizePortalTicketMode(mode string) string {
|
||||
|
||||
func (s *DingtalkSource) startPortalTicket(ctx context.Context, emit dwsevent.EmitFn) error {
|
||||
s.machine.OnConnecting()
|
||||
defer s.machine.OnStopped()
|
||||
|
||||
minBackoff := s.cfg.PortalTicket.ReconnectMin
|
||||
if minBackoff <= 0 {
|
||||
minBackoff = portalReconnectMin
|
||||
}
|
||||
maxBackoff := s.cfg.PortalTicket.ReconnectMax
|
||||
if maxBackoff <= 0 {
|
||||
maxBackoff = portalReconnectMax
|
||||
}
|
||||
if maxBackoff < minBackoff {
|
||||
maxBackoff = minBackoff
|
||||
}
|
||||
backoff := minBackoff
|
||||
for {
|
||||
acked, err := s.runPortalTicketAttempt(ctx, emit)
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
var stageErr *portalStageError
|
||||
if !errors.As(err, &stageErr) || stageErr == nil || !stageErr.retryable || s.cfg.PortalTicket.DisableReconnect {
|
||||
return err
|
||||
}
|
||||
if acked {
|
||||
backoff = minBackoff
|
||||
}
|
||||
s.machine.OnReconnect()
|
||||
slog.Warn("portal source reconnecting",
|
||||
"stage", stageErr.stage,
|
||||
"http_status", stageErr.status,
|
||||
"error_type", fmt.Sprintf("%T", stageErr.cause),
|
||||
"retry_in", backoff,
|
||||
"reconnect_count", s.machine.Snapshot().ReconnectCount,
|
||||
)
|
||||
if err := waitPersonalReconnect(ctx, backoff); err != nil {
|
||||
return err
|
||||
}
|
||||
backoff = nextPersonalBackoff(backoff, maxBackoff)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *DingtalkSource) runPortalTicketAttempt(ctx context.Context, emit dwsevent.EmitFn) (bool, error) {
|
||||
ticket, err := requestPortalTicket(ctx, s.cfg.PortalTicket)
|
||||
if err != nil {
|
||||
s.machine.OnStopped()
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
wsURL, err := websocketURL(ticket)
|
||||
if err != nil {
|
||||
s.machine.OnStopped()
|
||||
return err
|
||||
return false, err
|
||||
}
|
||||
|
||||
userAgent := strings.TrimSpace(s.cfg.PortalTicket.UserAgent)
|
||||
if userAgent == "" {
|
||||
userAgent = "dws-event-consume"
|
||||
}
|
||||
conn, resp, err := (&websocket.Dialer{HandshakeTimeout: 20 * time.Second}).DialContext(ctx, wsURL, http.Header{
|
||||
dialer := s.cfg.PortalTicket.WebSocketDialer
|
||||
if dialer == nil {
|
||||
dialer = &websocket.Dialer{HandshakeTimeout: 20 * time.Second}
|
||||
}
|
||||
conn, resp, err := dialer.DialContext(ctx, wsURL, http.Header{
|
||||
"User-Agent": []string{userAgent},
|
||||
})
|
||||
if err != nil {
|
||||
s.machine.OnStopped()
|
||||
status := 0
|
||||
cause := fmt.Errorf("source: portal stream connect: %w", err)
|
||||
if resp != nil {
|
||||
defer resp.Body.Close()
|
||||
status = resp.StatusCode
|
||||
raw, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
|
||||
return fmt.Errorf("source: portal stream connect HTTP %d: %s: %w",
|
||||
cause = fmt.Errorf("source: portal stream connect HTTP %d: %s: %w",
|
||||
resp.StatusCode, truncatePortalTicketLog(string(raw), 300), err)
|
||||
}
|
||||
return fmt.Errorf("source: portal stream connect: %w", err)
|
||||
return false, &portalStageError{
|
||||
stage: "stream_connect",
|
||||
status: status,
|
||||
retryable: status == 0 || retryableTicketStatus(status),
|
||||
cause: cause,
|
||||
}
|
||||
}
|
||||
defer conn.Close()
|
||||
attemptCtx, cancel := context.WithCancel(ctx)
|
||||
defer func() {
|
||||
cancel()
|
||||
_ = conn.Close()
|
||||
}()
|
||||
closeOnContext(attemptCtx, conn)
|
||||
s.machine.OnConnected()
|
||||
|
||||
closeOnContext(ctx, conn)
|
||||
handler := s.makeHandler(emit)
|
||||
acked := false
|
||||
for {
|
||||
_, message, err := conn.ReadMessage()
|
||||
if err != nil {
|
||||
s.machine.OnStopped()
|
||||
if isContextDone(ctx) {
|
||||
return ctx.Err()
|
||||
return acked, ctx.Err()
|
||||
}
|
||||
return fmt.Errorf("source: portal stream read: %w", err)
|
||||
return acked, &portalStageError{stage: "stream_read", retryable: true, cause: fmt.Errorf("source: portal stream read: %w", err)}
|
||||
}
|
||||
df, err := payload.DecodeDataFrame(message)
|
||||
if err != nil {
|
||||
@@ -147,12 +238,12 @@ func (s *DingtalkSource) startPortalTicket(ctx context.Context, emit dwsevent.Em
|
||||
resp, _ := handler(ctx, df)
|
||||
ensurePortalAckHeaders(resp, df)
|
||||
if err := portalWriteMessage(conn, websocket.TextMessage, resp.Encode()); err != nil {
|
||||
s.machine.OnStopped()
|
||||
if isContextDone(ctx) {
|
||||
return ctx.Err()
|
||||
return acked, ctx.Err()
|
||||
}
|
||||
return fmt.Errorf("source: portal stream ack: %w", err)
|
||||
return acked, &portalStageError{stage: "stream_ack", retryable: true, cause: fmt.Errorf("source: portal stream ack: %w", err)}
|
||||
}
|
||||
acked = true
|
||||
}
|
||||
}
|
||||
|
||||
@@ -164,12 +255,33 @@ type portalStreamTicket struct {
|
||||
func requestPortalTicket(ctx context.Context, cfg *PortalTicketConfig) (portalStreamTicket, error) {
|
||||
accessToken, err := resolveSourceAccessToken(ctx, cfg.AccessTokenProvider, cfg.AccessToken, "source: portal ticket")
|
||||
if err != nil {
|
||||
// Transient provider failures (network, 429, 5xx) must not kill a
|
||||
// long-running source; the reconnect loop retries after backoff.
|
||||
if authpkg.ClassifyRefreshFailure(err) == authpkg.RefreshFailureTransient {
|
||||
return portalStreamTicket{}, &portalStageError{stage: "ticket_auth", status: refreshHTTPStatus(err), retryable: true, cause: err}
|
||||
}
|
||||
return portalStreamTicket{}, err
|
||||
}
|
||||
httpClient := cfg.HTTPClient
|
||||
if httpClient == nil {
|
||||
httpClient = &http.Client{Timeout: 20 * time.Second}
|
||||
}
|
||||
ticket, status, err := requestPortalTicketAttempt(ctx, cfg, httpClient, accessToken)
|
||||
if status == http.StatusUnauthorized && cfg.ForceRefreshToken != nil {
|
||||
refreshed, refreshErr := refreshRejectedSourceToken(ctx, cfg.ForceRefreshToken, accessToken, "source: portal ticket", err)
|
||||
if refreshErr != nil {
|
||||
if authpkg.ClassifyRefreshFailure(refreshErr) == authpkg.RefreshFailureTransient {
|
||||
return portalStreamTicket{}, &portalStageError{stage: "ticket_auth_refresh", status: refreshHTTPStatus(refreshErr), retryable: true, cause: refreshErr}
|
||||
}
|
||||
return portalStreamTicket{}, refreshErr
|
||||
}
|
||||
// Retry once with the freshly rotated token; a second 401 stays fatal.
|
||||
ticket, _, err = requestPortalTicketAttempt(ctx, cfg, httpClient, refreshed)
|
||||
}
|
||||
return ticket, err
|
||||
}
|
||||
|
||||
func requestPortalTicketAttempt(ctx context.Context, cfg *PortalTicketConfig, httpClient *http.Client, accessToken string) (portalStreamTicket, int, error) {
|
||||
body := map[string]string{
|
||||
"sourceId": strings.TrimSpace(cfg.SourceID),
|
||||
"channelType": strings.TrimSpace(cfg.SourceID),
|
||||
@@ -182,7 +294,7 @@ func requestPortalTicket(ctx context.Context, cfg *PortalTicketConfig) (portalSt
|
||||
rawBody, _ := json.Marshal(body)
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, strings.TrimSpace(cfg.TicketURL), bytes.NewReader(rawBody))
|
||||
if err != nil {
|
||||
return portalStreamTicket{}, err
|
||||
return portalStreamTicket{}, 0, err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("Accept", "application/json")
|
||||
@@ -193,18 +305,34 @@ func requestPortalTicket(ctx context.Context, cfg *PortalTicketConfig) (portalSt
|
||||
|
||||
resp, err := httpClient.Do(req)
|
||||
if err != nil {
|
||||
return portalStreamTicket{}, fmt.Errorf("source: portal ticket request: %w", err)
|
||||
return portalStreamTicket{}, 0, &portalStageError{stage: "ticket_request", retryable: true, cause: fmt.Errorf("source: portal ticket request: %w", err)}
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
raw, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
|
||||
if resp.StatusCode >= 400 {
|
||||
return portalStreamTicket{}, fmt.Errorf("source: portal ticket HTTP %d: %s",
|
||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||
// Preserve HTTP status semantics before touching the body. A truncated
|
||||
// 401 body must not become a retryable read error that bypasses the
|
||||
// single token-refresh guard; the body is only best-effort diagnostics.
|
||||
raw, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
|
||||
httpErr := fmt.Errorf("source: portal ticket HTTP %d: %s",
|
||||
resp.StatusCode, truncatePortalTicketLog(string(raw), 300))
|
||||
if retryableTicketStatus(resp.StatusCode) {
|
||||
return portalStreamTicket{}, resp.StatusCode, &portalStageError{stage: "ticket_request", status: resp.StatusCode, retryable: true, cause: httpErr}
|
||||
}
|
||||
return portalStreamTicket{}, resp.StatusCode, httpErr
|
||||
}
|
||||
raw, readErr := io.ReadAll(io.LimitReader(resp.Body, 1<<20))
|
||||
if readErr != nil {
|
||||
return portalStreamTicket{}, resp.StatusCode, &portalStageError{
|
||||
stage: "ticket_request",
|
||||
status: resp.StatusCode,
|
||||
retryable: true,
|
||||
cause: fmt.Errorf("source: portal ticket read: %w", readErr),
|
||||
}
|
||||
}
|
||||
|
||||
var direct portalStreamTicket
|
||||
if err := json.Unmarshal(raw, &direct); err == nil && direct.Endpoint != "" && direct.Ticket != "" {
|
||||
return direct, nil
|
||||
return direct, resp.StatusCode, nil
|
||||
}
|
||||
|
||||
var envelope struct {
|
||||
@@ -214,16 +342,16 @@ func requestPortalTicket(ctx context.Context, cfg *PortalTicketConfig) (portalSt
|
||||
ErrorMsg string `json:"errorMsg"`
|
||||
}
|
||||
if err := json.Unmarshal(raw, &envelope); err != nil {
|
||||
return portalStreamTicket{}, fmt.Errorf("source: portal ticket parse: %w", err)
|
||||
return portalStreamTicket{}, resp.StatusCode, fmt.Errorf("source: portal ticket parse: %w", err)
|
||||
}
|
||||
if !envelope.Success {
|
||||
return portalStreamTicket{}, fmt.Errorf("source: portal ticket failed: %s %s",
|
||||
return portalStreamTicket{}, resp.StatusCode, fmt.Errorf("source: portal ticket failed: %s %s",
|
||||
envelope.ErrorCode, envelope.ErrorMsg)
|
||||
}
|
||||
if envelope.Result.Endpoint == "" || envelope.Result.Ticket == "" {
|
||||
return portalStreamTicket{}, errors.New("source: portal ticket result missing endpoint/ticket")
|
||||
return portalStreamTicket{}, resp.StatusCode, errors.New("source: portal ticket result missing endpoint/ticket")
|
||||
}
|
||||
return envelope.Result, nil
|
||||
return envelope.Result, resp.StatusCode, nil
|
||||
}
|
||||
|
||||
func websocketURL(ticket portalStreamTicket) (string, error) {
|
||||
|
||||
@@ -0,0 +1,485 @@
|
||||
package source
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
dwsevent "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event"
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/open-dingtalk/dingtalk-stream-sdk-go/payload"
|
||||
)
|
||||
|
||||
// TestCrossPlatformCoveragePortalStart401RefreshRetryEndToEnd drives the full production chain
|
||||
// DingtalkSource.Start → startPortalTicket → requestPortalTicket: the first
|
||||
// ticket request is rejected with 401, ForceRefreshToken rotates the token,
|
||||
// the in-chain retry succeeds with the fresh token and a WebSocket event is
|
||||
// delivered to emit.
|
||||
func TestCrossPlatformCoveragePortalStart401RefreshRetryEndToEnd(t *testing.T) {
|
||||
var ticketCalls, refreshCalls atomic.Int64
|
||||
var rejectedSeen atomic.Value
|
||||
|
||||
upgrader := websocket.Upgrader{}
|
||||
var wsURL string
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/ticket", func(w http.ResponseWriter, r *http.Request) {
|
||||
ticketCalls.Add(1)
|
||||
switch r.Header.Get("x-user-access-token") {
|
||||
case "fresh-token":
|
||||
_ = json.NewEncoder(w).Encode(map[string]any{
|
||||
"success": true,
|
||||
"result": map[string]string{"endpoint": wsURL, "ticket": "ticket-1"},
|
||||
})
|
||||
default:
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
_, _ = io.WriteString(w, "token expired")
|
||||
}
|
||||
})
|
||||
mux.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) {
|
||||
conn, err := upgrader.Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
df := payload.DataFrame{Type: "event", Headers: payload.DataFrameHeader{payload.DataFrameHeaderKMessageId: "msg-1"}, Data: `{}`}
|
||||
_ = conn.WriteJSON(df)
|
||||
_, _, _ = conn.ReadMessage()
|
||||
})
|
||||
srv := httptest.NewServer(mux)
|
||||
defer srv.Close()
|
||||
wsURL = "ws" + strings.TrimPrefix(srv.URL, "http") + "/ws"
|
||||
|
||||
s, err := New(Config{PortalTicket: &PortalTicketConfig{
|
||||
TicketURL: srv.URL + "/ticket",
|
||||
AccessTokenProvider: func(context.Context) (string, error) {
|
||||
return "stale-token", nil
|
||||
},
|
||||
ForceRefreshToken: func(_ context.Context, rejectedToken string) (string, error) {
|
||||
refreshCalls.Add(1)
|
||||
rejectedSeen.Store(rejectedToken)
|
||||
return "fresh-token", nil
|
||||
},
|
||||
SourceID: "source",
|
||||
HTTPClient: srv.Client(),
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
emitted := make(chan struct{}, 1)
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- s.Start(ctx, func(*dwsevent.RawEvent) { emitted <- struct{}{} }) }()
|
||||
select {
|
||||
case <-emitted:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("portal event timeout after 401 refresh retry")
|
||||
}
|
||||
cancel()
|
||||
select {
|
||||
case err := <-done:
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("portal stop = %v", err)
|
||||
}
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("portal stop timeout")
|
||||
}
|
||||
|
||||
if got := ticketCalls.Load(); got != 2 {
|
||||
t.Fatalf("ticket calls = %d, want 2", got)
|
||||
}
|
||||
if got := refreshCalls.Load(); got != 1 {
|
||||
t.Fatalf("refresh calls = %d, want 1", got)
|
||||
}
|
||||
if got, _ := rejectedSeen.Load().(string); got != "stale-token" {
|
||||
t.Fatalf("rejected token = %q, want %q", got, "stale-token")
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoverageRequestPortalTicketRetryUsesRotatedTokenDirectly asserts the in-chain
|
||||
// retry sends the token returned by ForceRefreshToken instead of re-invoking
|
||||
// the provider (which could still serve the stale token).
|
||||
func TestCrossPlatformCoverageRequestPortalTicketRetryUsesRotatedTokenDirectly(t *testing.T) {
|
||||
providerCalls := 0
|
||||
var attemptTokens []string
|
||||
client := &http.Client{Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) {
|
||||
token := req.Header.Get("x-user-access-token")
|
||||
attemptTokens = append(attemptTokens, token)
|
||||
if token != "rotated" {
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: io.NopCloser(strings.NewReader("expired")), Header: make(http.Header)}, nil
|
||||
}
|
||||
return &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(`{"endpoint":"wss://x","ticket":"t"}`)), Header: make(http.Header)}, nil
|
||||
})}
|
||||
ticket, err := requestPortalTicket(context.Background(), &PortalTicketConfig{
|
||||
TicketURL: "https://x",
|
||||
AccessTokenProvider: func(context.Context) (string, error) {
|
||||
providerCalls++
|
||||
return "stale", nil
|
||||
},
|
||||
ForceRefreshToken: func(_ context.Context, rejectedToken string) (string, error) {
|
||||
if rejectedToken != "stale" {
|
||||
t.Fatalf("rejected token = %q, want %q", rejectedToken, "stale")
|
||||
}
|
||||
return "rotated", nil
|
||||
},
|
||||
SourceID: "s",
|
||||
HTTPClient: client,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("requestPortalTicket = %v", err)
|
||||
}
|
||||
if ticket.Endpoint != "wss://x" || ticket.Ticket != "t" {
|
||||
t.Fatalf("ticket = %#v", ticket)
|
||||
}
|
||||
if providerCalls != 1 {
|
||||
t.Fatalf("provider calls = %d, want 1", providerCalls)
|
||||
}
|
||||
if len(attemptTokens) != 2 || attemptTokens[0] != "stale" || attemptTokens[1] != "rotated" {
|
||||
t.Fatalf("attempt tokens = %v", attemptTokens)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoverageRequestPortalTicketRefreshFailureKeepsBothErrors asserts a failing
|
||||
// refresh neither retries nor drops the refresh error or the original 401.
|
||||
func TestCrossPlatformCoverageRequestPortalTicketRefreshFailureKeepsBothErrors(t *testing.T) {
|
||||
refreshErr := errors.New("refresh_token exchange failed")
|
||||
attempts := 0
|
||||
client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
attempts++
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: io.NopCloser(strings.NewReader("expired")), Header: make(http.Header)}, nil
|
||||
})}
|
||||
_, err := requestPortalTicket(context.Background(), &PortalTicketConfig{
|
||||
TicketURL: "https://x",
|
||||
AccessToken: "stale",
|
||||
ForceRefreshToken: func(context.Context, string) (string, error) {
|
||||
return "", refreshErr
|
||||
},
|
||||
SourceID: "s",
|
||||
HTTPClient: client,
|
||||
})
|
||||
if !errors.Is(err, refreshErr) {
|
||||
t.Fatalf("error should wrap refresh error, got %v", err)
|
||||
}
|
||||
if err == nil || !strings.Contains(err.Error(), "HTTP 401") {
|
||||
t.Fatalf("error should keep original 401, got %v", err)
|
||||
}
|
||||
if attempts != 1 {
|
||||
t.Fatalf("attempts = %d, want 1 (no retry after failed refresh)", attempts)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoverageRequestPortalTicketWithoutRefreshCallback401StaysFatal covers backward
|
||||
// compatibility: nil ForceRefreshToken keeps the single-attempt fatal 401.
|
||||
func TestCrossPlatformCoverageRequestPortalTicketWithoutRefreshCallback401StaysFatal(t *testing.T) {
|
||||
attempts := 0
|
||||
client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
attempts++
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: io.NopCloser(strings.NewReader("expired")), Header: make(http.Header)}, nil
|
||||
})}
|
||||
_, err := requestPortalTicket(context.Background(), &PortalTicketConfig{
|
||||
TicketURL: "https://x", AccessToken: "stale", SourceID: "s", HTTPClient: client,
|
||||
})
|
||||
if err == nil || !strings.Contains(err.Error(), "HTTP 401") {
|
||||
t.Fatalf("fatal 401 expected, got %v", err)
|
||||
}
|
||||
if attempts != 1 {
|
||||
t.Fatalf("attempts = %d, want 1", attempts)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoverageRequestPortalTicketSecond401IsFatal guards against refresh loops: the
|
||||
// controlled retry happens exactly once even if the rotated token is also
|
||||
// rejected.
|
||||
func TestCrossPlatformCoverageRequestPortalTicketSecond401IsFatal(t *testing.T) {
|
||||
attempts := 0
|
||||
refreshCalls := 0
|
||||
client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
attempts++
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: io.NopCloser(strings.NewReader("expired")), Header: make(http.Header)}, nil
|
||||
})}
|
||||
_, err := requestPortalTicket(context.Background(), &PortalTicketConfig{
|
||||
TicketURL: "https://x",
|
||||
AccessToken: "stale",
|
||||
ForceRefreshToken: func(context.Context, string) (string, error) {
|
||||
refreshCalls++
|
||||
return "rotated-but-still-rejected", nil
|
||||
},
|
||||
SourceID: "s",
|
||||
HTTPClient: client,
|
||||
})
|
||||
if err == nil || !strings.Contains(err.Error(), "HTTP 401") {
|
||||
t.Fatalf("fatal 401 expected after single retry, got %v", err)
|
||||
}
|
||||
if attempts != 2 || refreshCalls != 1 {
|
||||
t.Fatalf("attempts = %d refreshCalls = %d, want 2/1", attempts, refreshCalls)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoverageRequestPortalTicketRefreshEmptyTokenIsFatal asserts an empty rotated
|
||||
// token is rejected instead of being sent to the server.
|
||||
func TestCrossPlatformCoverageRequestPortalTicketRefreshEmptyTokenIsFatal(t *testing.T) {
|
||||
attempts := 0
|
||||
client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
attempts++
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: io.NopCloser(strings.NewReader("expired")), Header: make(http.Header)}, nil
|
||||
})}
|
||||
_, err := requestPortalTicket(context.Background(), &PortalTicketConfig{
|
||||
TicketURL: "https://x",
|
||||
AccessToken: "stale",
|
||||
ForceRefreshToken: func(context.Context, string) (string, error) {
|
||||
return " ", nil
|
||||
},
|
||||
SourceID: "s",
|
||||
HTTPClient: client,
|
||||
})
|
||||
if err == nil || !strings.Contains(err.Error(), "empty token") {
|
||||
t.Fatalf("empty rotated token error expected, got %v", err)
|
||||
}
|
||||
if attempts != 1 {
|
||||
t.Fatalf("attempts = %d, want 1", attempts)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoveragePersonalFetchTicket401RefreshRetry mirrors the portal behavior for the
|
||||
// personal stream ticket path.
|
||||
func TestCrossPlatformCoveragePersonalFetchTicket401RefreshRetry(t *testing.T) {
|
||||
var attemptTokens []string
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessTokenProvider: func(context.Context) (string, error) { return "stale", nil },
|
||||
ForceRefreshToken: func(_ context.Context, rejectedToken string) (string, error) {
|
||||
if rejectedToken != "stale" {
|
||||
t.Fatalf("rejected token = %q, want %q", rejectedToken, "stale")
|
||||
}
|
||||
return "rotated", nil
|
||||
},
|
||||
ClientID: "client",
|
||||
SourceID: "source",
|
||||
TicketURL: "https://ticket.test",
|
||||
HTTPClient: &http.Client{Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) {
|
||||
token := req.Header.Get("x-user-access-token")
|
||||
attemptTokens = append(attemptTokens, token)
|
||||
if token != "rotated" {
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: io.NopCloser(strings.NewReader("expired")), Header: make(http.Header)}, nil
|
||||
}
|
||||
return &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(`{"endpoint":"wss://stream.test","ticket":"ticket"}`)), Header: make(http.Header)}, nil
|
||||
})},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ticket, err := src.fetchTicket(context.Background())
|
||||
if err != nil {
|
||||
t.Fatalf("fetchTicket = %v", err)
|
||||
}
|
||||
if ticket.Endpoint != "wss://stream.test" || ticket.Ticket != "ticket" {
|
||||
t.Fatalf("ticket = %#v", ticket)
|
||||
}
|
||||
if len(attemptTokens) != 2 || attemptTokens[0] != "stale" || attemptTokens[1] != "rotated" {
|
||||
t.Fatalf("attempt tokens = %v", attemptTokens)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoveragePersonalFetchTicket401RefreshFailureStaysFatal asserts a failed refresh
|
||||
// keeps the 401 fatal (not retryable) and wraps the refresh error.
|
||||
func TestCrossPlatformCoveragePersonalFetchTicket401RefreshFailureStaysFatal(t *testing.T) {
|
||||
refreshErr := errors.New("refresh_token exchange failed")
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessToken: "stale",
|
||||
ForceRefreshToken: func(context.Context, string) (string, error) {
|
||||
return "", refreshErr
|
||||
},
|
||||
ClientID: "client",
|
||||
SourceID: "source",
|
||||
TicketURL: "https://ticket.test",
|
||||
HTTPClient: &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: io.NopCloser(strings.NewReader("expired")), Header: make(http.Header)}, nil
|
||||
})},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err = src.fetchTicket(context.Background())
|
||||
if !errors.Is(err, refreshErr) {
|
||||
t.Fatalf("error should wrap refresh error, got %v", err)
|
||||
}
|
||||
if isRetryablePersonalError(err) {
|
||||
t.Fatalf("failed refresh should stay fatal, got retryable %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// brokenBody simulates a response body that fails mid-read, e.g. the server
|
||||
// closing the connection before the error payload is fully written.
|
||||
type brokenBody struct{}
|
||||
|
||||
func (brokenBody) Read([]byte) (int, error) { return 0, io.ErrUnexpectedEOF }
|
||||
func (brokenBody) Close() error { return nil }
|
||||
|
||||
func TestCrossPlatformCoveragePortalTicketNon2xxTruncatedBodyKeepsStatus(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
status int
|
||||
retryable bool
|
||||
}{
|
||||
{status: http.StatusUnauthorized},
|
||||
{status: http.StatusTooManyRequests, retryable: true},
|
||||
{status: http.StatusServiceUnavailable, retryable: true},
|
||||
} {
|
||||
t.Run(http.StatusText(tc.status), func(t *testing.T) {
|
||||
client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: tc.status, Body: brokenBody{}, Header: make(http.Header)}, nil
|
||||
})}
|
||||
|
||||
_, status, err := requestPortalTicketAttempt(context.Background(), &PortalTicketConfig{
|
||||
TicketURL: "https://ticket.test",
|
||||
SourceID: "source",
|
||||
}, client, "token")
|
||||
if status != tc.status || err == nil || !strings.Contains(err.Error(), fmt.Sprintf("HTTP %d", tc.status)) {
|
||||
t.Fatalf("requestPortalTicketAttempt() status=%d err=%v, want HTTP %d", status, err, tc.status)
|
||||
}
|
||||
if errors.Is(err, io.ErrUnexpectedEOF) {
|
||||
t.Fatalf("non-2xx status error should not expose diagnostic body read failure: %v", err)
|
||||
}
|
||||
var stageErr *portalStageError
|
||||
if tc.retryable {
|
||||
if !errors.As(err, &stageErr) || stageErr.status != tc.status || !stageErr.retryable {
|
||||
t.Fatalf("retryable status should return retryable portalStageError, got %T %v", err, err)
|
||||
}
|
||||
} else if errors.As(err, &stageErr) && stageErr.retryable {
|
||||
t.Fatalf("fatal status should not become retryable stage error: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePortalTicket200TruncatedBodyIsRetryable(t *testing.T) {
|
||||
client := &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: http.StatusOK, Body: brokenBody{}, Header: make(http.Header)}, nil
|
||||
})}
|
||||
|
||||
_, status, err := requestPortalTicketAttempt(context.Background(), &PortalTicketConfig{
|
||||
TicketURL: "https://ticket.test",
|
||||
SourceID: "source",
|
||||
}, client, "token")
|
||||
if status != http.StatusOK || !errors.Is(err, io.ErrUnexpectedEOF) {
|
||||
t.Fatalf("requestPortalTicketAttempt() status=%d err=%v, want 200 + io.ErrUnexpectedEOF", status, err)
|
||||
}
|
||||
var stageErr *portalStageError
|
||||
if !errors.As(err, &stageErr) || !stageErr.retryable || stageErr.stage != "ticket_request" {
|
||||
t.Fatalf("2xx truncated body should be retryable ticket_request stage, got %T %v", err, err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoveragePersonalFetchTicket401TruncatedBodyStaysFatal guards the single
|
||||
// refresh-retry protection: a 401 whose body fails with unexpected EOF must
|
||||
// be classified by status (fatal) and never wrapped as retryable, otherwise
|
||||
// the outer reconnect loop would refresh again on every iteration.
|
||||
func TestCrossPlatformCoveragePersonalFetchTicket401TruncatedBodyStaysFatal(t *testing.T) {
|
||||
attempts := 0
|
||||
refreshCalls := 0
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessToken: "stale",
|
||||
ForceRefreshToken: func(context.Context, string) (string, error) {
|
||||
refreshCalls++
|
||||
return "rotated", nil
|
||||
},
|
||||
ClientID: "client",
|
||||
SourceID: "source",
|
||||
TicketURL: "https://ticket.test",
|
||||
HTTPClient: &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
attempts++
|
||||
return &http.Response{StatusCode: http.StatusUnauthorized, Body: brokenBody{}, Header: make(http.Header)}, nil
|
||||
})},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err = src.fetchTicket(context.Background())
|
||||
if err == nil || !strings.Contains(err.Error(), "HTTP 401") {
|
||||
t.Fatalf("fatal 401 error expected, got %v", err)
|
||||
}
|
||||
if isRetryablePersonalError(err) {
|
||||
t.Fatalf("401 with truncated body must stay fatal, got retryable %v", err)
|
||||
}
|
||||
if refreshCalls != 1 {
|
||||
t.Fatalf("refresh calls = %d, want 1", refreshCalls)
|
||||
}
|
||||
if attempts != 2 {
|
||||
t.Fatalf("attempts = %d, want 2", attempts)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoveragePersonalFetchTicket200TruncatedBodyStaysRetryable pins the existing
|
||||
// behavior for success responses: a body read failure on 2xx is a transient
|
||||
// transport problem and remains retryable.
|
||||
func TestCrossPlatformCoveragePersonalFetchTicket200TruncatedBodyStaysRetryable(t *testing.T) {
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessToken: "token",
|
||||
ClientID: "client",
|
||||
SourceID: "source",
|
||||
TicketURL: "https://ticket.test",
|
||||
HTTPClient: &http.Client{Transport: roundTripFunc(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: 200, Body: brokenBody{}, Header: make(http.Header)}, nil
|
||||
})},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
_, err = src.fetchTicket(context.Background())
|
||||
if !errors.Is(err, io.ErrUnexpectedEOF) {
|
||||
t.Fatalf("error should wrap io.ErrUnexpectedEOF, got %v", err)
|
||||
}
|
||||
if !isRetryablePersonalError(err) {
|
||||
t.Fatalf("2xx body read failure should stay retryable, got %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// TestCrossPlatformCoveragePersonalFetchTicketAttemptTransportAndPayloadEdges covers the
|
||||
// remaining fetchTicketAttempt branches: transport failures and retryable
|
||||
// statuses stay retryable, while a well-formed response missing the endpoint
|
||||
// or ticket fields stays fatal.
|
||||
func TestCrossPlatformCoveragePersonalFetchTicketAttemptTransportAndPayloadEdges(t *testing.T) {
|
||||
newSource := func(rt roundTripFunc) *PersonalSource {
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessToken: "token",
|
||||
ClientID: "client",
|
||||
SourceID: "source",
|
||||
TicketURL: "https://ticket.test",
|
||||
HTTPClient: &http.Client{Transport: rt},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return src
|
||||
}
|
||||
|
||||
dialErr := errors.New("dial tcp: connection refused")
|
||||
src := newSource(func(*http.Request) (*http.Response, error) { return nil, dialErr })
|
||||
_, _, err := src.fetchTicketAttempt(context.Background(), "token")
|
||||
if !errors.Is(err, dialErr) || !isRetryablePersonalError(err) {
|
||||
t.Fatalf("transport failure should stay retryable and wrap cause, got %v", err)
|
||||
}
|
||||
|
||||
src = newSource(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: http.StatusServiceUnavailable, Body: io.NopCloser(strings.NewReader("busy")), Header: make(http.Header)}, nil
|
||||
})
|
||||
_, status, err := src.fetchTicketAttempt(context.Background(), "token")
|
||||
if status != http.StatusServiceUnavailable || !isRetryablePersonalError(err) {
|
||||
t.Fatalf("503 should stay retryable, got status %d err %v", status, err)
|
||||
}
|
||||
|
||||
src = newSource(func(*http.Request) (*http.Response, error) {
|
||||
return &http.Response{StatusCode: 200, Body: io.NopCloser(strings.NewReader(`{"endpoint":"","ticket":""}`)), Header: make(http.Header)}, nil
|
||||
})
|
||||
_, _, err = src.fetchTicketAttempt(context.Background(), "token")
|
||||
if err == nil || isRetryablePersonalError(err) || !strings.Contains(err.Error(), "missing endpoint or ticket") {
|
||||
t.Fatalf("missing ticket fields should stay fatal, got %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,295 @@
|
||||
// Copyright 2026 Alibaba Group
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
|
||||
package source
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
authpkg "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/auth"
|
||||
dwsevent "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/event"
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/open-dingtalk/dingtalk-stream-sdk-go/payload"
|
||||
)
|
||||
|
||||
func transientRefreshError() error {
|
||||
return &authpkg.HTTPStatusError{StatusCode: http.StatusServiceUnavailable}
|
||||
}
|
||||
|
||||
func terminalRefreshError() error {
|
||||
return &authpkg.HTTPStatusError{StatusCode: http.StatusUnauthorized}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePersonalRetryLogErrorReportsOnlySafeTransientStatus(t *testing.T) {
|
||||
cause := errors.New("personal source: ticket HTTP 401 secret detail")
|
||||
_, err := refreshRejectedSourceToken(context.Background(), func(context.Context, string) (string, error) {
|
||||
return "", transientRefreshError()
|
||||
}, "rejected", "personal source", cause)
|
||||
if got, want := personalRetryLogError(retryPersonal(err)), "personal source: token refresh HTTP 503"; got != want {
|
||||
t.Fatalf("personalRetryLogError() = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePersonalSourceRetriesTransientTokenResolutionFailure(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
var calls atomic.Int32
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessTokenProvider: func(context.Context) (string, error) {
|
||||
if calls.Add(1) == 2 {
|
||||
cancel()
|
||||
}
|
||||
return "", transientRefreshError()
|
||||
},
|
||||
ClientID: "client",
|
||||
SourceID: "open",
|
||||
TicketURL: "https://ticket.invalid",
|
||||
ReconnectMin: time.Millisecond,
|
||||
ReconnectMax: time.Millisecond,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = src.Start(ctx, func(*dwsevent.RawEvent) {})
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("Start() error = %v, want context canceled after retry", err)
|
||||
}
|
||||
if calls.Load() != 2 || src.State().ReconnectCount != 1 {
|
||||
t.Fatalf("provider calls=%d reconnects=%d, want 2 calls and 1 reconnect", calls.Load(), src.State().ReconnectCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePersonalSourceDoesNotRetryTerminalTokenResolutionFailure(t *testing.T) {
|
||||
var calls atomic.Int32
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessTokenProvider: func(context.Context) (string, error) {
|
||||
calls.Add(1)
|
||||
return "", terminalRefreshError()
|
||||
},
|
||||
ClientID: "client",
|
||||
SourceID: "open",
|
||||
TicketURL: "https://ticket.invalid",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = src.Start(context.Background(), func(*dwsevent.RawEvent) {})
|
||||
if authpkg.ClassifyRefreshFailure(err) != authpkg.RefreshFailureTerminal {
|
||||
t.Fatalf("Start() error = %v, want terminal refresh failure", err)
|
||||
}
|
||||
if calls.Load() != 1 || src.State().ReconnectCount != 0 {
|
||||
t.Fatalf("provider calls=%d reconnects=%d, want 1 call and no reconnect", calls.Load(), src.State().ReconnectCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePersonalSourceRetriesTransientRejectedTokenRefresh(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
var ticketCalls atomic.Int32
|
||||
var refreshCalls atomic.Int32
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
ticketCalls.Add(1)
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
src, err := NewPersonal(PersonalConfig{
|
||||
AccessTokenProvider: func(context.Context) (string, error) { return "old-token", nil },
|
||||
ForceRefreshToken: func(_ context.Context, rejected string) (string, error) {
|
||||
if rejected != "old-token" {
|
||||
t.Fatalf("rejected token = %q, want old-token", rejected)
|
||||
}
|
||||
if refreshCalls.Add(1) == 2 {
|
||||
cancel()
|
||||
}
|
||||
return "", transientRefreshError()
|
||||
},
|
||||
ClientID: "client",
|
||||
SourceID: "open",
|
||||
TicketURL: srv.URL,
|
||||
HTTPClient: srv.Client(),
|
||||
ReconnectMin: time.Millisecond,
|
||||
ReconnectMax: time.Millisecond,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = src.Start(ctx, func(*dwsevent.RawEvent) {})
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("Start() error = %v, want context canceled after retry", err)
|
||||
}
|
||||
if ticketCalls.Load() != 2 || refreshCalls.Load() != 2 || src.State().ReconnectCount != 1 {
|
||||
t.Fatalf("ticket calls=%d refresh calls=%d reconnects=%d", ticketCalls.Load(), refreshCalls.Load(), src.State().ReconnectCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePersonalRetryLogErrorFallsBackToNetworkMessage(t *testing.T) {
|
||||
err := retryPersonal(fmt.Errorf("personal source: resolve access token: %w", errors.New("dial tcp: lookup oauth.invalid")))
|
||||
if got, want := personalRetryLogError(err), "personal source: token refresh: temporary network error"; got != want {
|
||||
t.Fatalf("personalRetryLogError() = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePortalStageErrorNilAndUnwrap(t *testing.T) {
|
||||
var nilErr *portalStageError
|
||||
if got, want := nilErr.Error(), "source: portal stream failed"; got != want {
|
||||
t.Fatalf("nil stage error = %q, want %q", got, want)
|
||||
}
|
||||
if nilErr.Unwrap() != nil {
|
||||
t.Fatal("nil stage error should unwrap to nil")
|
||||
}
|
||||
cause := errors.New("cause")
|
||||
stageErr := &portalStageError{stage: "stream_read", retryable: true, cause: cause}
|
||||
if !errors.Is(stageErr, cause) {
|
||||
t.Fatalf("stage error should unwrap to cause: %v", stageErr)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePortalSourceRetriesTransientTokenResolutionFailure(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
var calls atomic.Int32
|
||||
src, err := New(Config{PortalTicket: &PortalTicketConfig{
|
||||
TicketURL: "https://ticket.invalid",
|
||||
AccessTokenProvider: func(context.Context) (string, error) {
|
||||
if calls.Add(1) == 2 {
|
||||
cancel()
|
||||
}
|
||||
return "", transientRefreshError()
|
||||
},
|
||||
SourceID: "open",
|
||||
// Min above max exercises the reconnect clamp.
|
||||
ReconnectMin: 2 * time.Millisecond,
|
||||
ReconnectMax: time.Millisecond,
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = src.Start(ctx, func(*dwsevent.RawEvent) {})
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("Start() error = %v, want context canceled after retry", err)
|
||||
}
|
||||
if calls.Load() != 2 || src.State().ReconnectCount != 1 {
|
||||
t.Fatalf("provider calls=%d reconnects=%d, want 2 calls and 1 reconnect", calls.Load(), src.State().ReconnectCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePortalSourceResetsBackoffAfterAckedAttempt(t *testing.T) {
|
||||
var ticketCalls atomic.Int32
|
||||
upgrader := websocket.Upgrader{}
|
||||
var wsURL string
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/ticket", func(w http.ResponseWriter, _ *http.Request) {
|
||||
if ticketCalls.Add(1) > 1 {
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
_, _ = io.WriteString(w, `{"endpoint":`+strconvQuote(wsURL)+`,"ticket":"t"}`)
|
||||
})
|
||||
mux.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) {
|
||||
conn, err := upgrader.Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
df := payload.DataFrame{Type: "event", Headers: payload.DataFrameHeader{payload.DataFrameHeaderKMessageId: "m"}, Data: `{}`}
|
||||
_ = conn.WriteJSON(df)
|
||||
// Wait for the ACK, then close so the read fails retryably with an
|
||||
// acked attempt behind it, which resets the reconnect backoff.
|
||||
_, _, _ = conn.ReadMessage()
|
||||
})
|
||||
srv := httptest.NewServer(mux)
|
||||
defer srv.Close()
|
||||
wsURL = "ws" + strings.TrimPrefix(srv.URL, "http") + "/ws"
|
||||
|
||||
src, err := New(Config{PortalTicket: &PortalTicketConfig{
|
||||
TicketURL: srv.URL + "/ticket",
|
||||
AccessToken: "t",
|
||||
SourceID: "open",
|
||||
HTTPClient: srv.Client(),
|
||||
ReconnectMin: time.Millisecond,
|
||||
ReconnectMax: time.Millisecond,
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var events atomic.Int32
|
||||
err = src.Start(context.Background(), func(*dwsevent.RawEvent) { events.Add(1) })
|
||||
if err == nil || !strings.Contains(err.Error(), "portal ticket HTTP 400") {
|
||||
t.Fatalf("Start() error = %v, want fatal ticket HTTP 400 after reconnect", err)
|
||||
}
|
||||
if events.Load() != 1 || ticketCalls.Load() != 2 || src.State().ReconnectCount != 1 {
|
||||
t.Fatalf("events=%d ticket calls=%d reconnects=%d, want 1/2/1", events.Load(), ticketCalls.Load(), src.State().ReconnectCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePortalSourceDoesNotRetryTerminalTokenResolutionFailure(t *testing.T) {
|
||||
var calls atomic.Int32
|
||||
src, err := New(Config{PortalTicket: &PortalTicketConfig{
|
||||
TicketURL: "https://ticket.invalid",
|
||||
AccessTokenProvider: func(context.Context) (string, error) {
|
||||
calls.Add(1)
|
||||
return "", terminalRefreshError()
|
||||
},
|
||||
SourceID: "open",
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = src.Start(context.Background(), func(*dwsevent.RawEvent) {})
|
||||
if authpkg.ClassifyRefreshFailure(err) != authpkg.RefreshFailureTerminal {
|
||||
t.Fatalf("Start() error = %v, want terminal refresh failure", err)
|
||||
}
|
||||
if calls.Load() != 1 || src.State().ReconnectCount != 0 {
|
||||
t.Fatalf("provider calls=%d reconnects=%d, want 1 call and no reconnect", calls.Load(), src.State().ReconnectCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCrossPlatformCoveragePortalSourceRetriesTransientRejectedTokenRefresh(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
var ticketCalls atomic.Int32
|
||||
var refreshCalls atomic.Int32
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
||||
ticketCalls.Add(1)
|
||||
w.WriteHeader(http.StatusUnauthorized)
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
src, err := New(Config{PortalTicket: &PortalTicketConfig{
|
||||
TicketURL: srv.URL,
|
||||
AccessTokenProvider: func(context.Context) (string, error) { return "old-token", nil },
|
||||
ForceRefreshToken: func(_ context.Context, rejected string) (string, error) {
|
||||
if rejected != "old-token" {
|
||||
t.Fatalf("rejected token = %q, want old-token", rejected)
|
||||
}
|
||||
if refreshCalls.Add(1) == 2 {
|
||||
cancel()
|
||||
}
|
||||
return "", transientRefreshError()
|
||||
},
|
||||
SourceID: "open",
|
||||
HTTPClient: srv.Client(),
|
||||
ReconnectMin: time.Millisecond,
|
||||
ReconnectMax: time.Millisecond,
|
||||
}})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
err = src.Start(ctx, func(*dwsevent.RawEvent) {})
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("Start() error = %v, want context canceled after retry", err)
|
||||
}
|
||||
if ticketCalls.Load() != 2 || refreshCalls.Load() != 2 || src.State().ReconnectCount != 1 {
|
||||
t.Fatalf("ticket calls=%d refresh calls=%d reconnects=%d", ticketCalls.Load(), refreshCalls.Load(), src.State().ReconnectCount)
|
||||
}
|
||||
}
|
||||
Executable
+52
@@ -0,0 +1,52 @@
|
||||
#!/bin/sh
|
||||
set -eu
|
||||
|
||||
TAG="${1:-}"
|
||||
[ -n "$TAG" ] || {
|
||||
printf 'usage: release-tag-oss-mode.sh <tag>\n' >&2
|
||||
exit 2
|
||||
}
|
||||
|
||||
tag_ref="refs/tags/$TAG"
|
||||
git rev-parse --verify --quiet "$tag_ref" >/dev/null || {
|
||||
printf 'release tag is not available: %s\n' "$TAG" >&2
|
||||
exit 1
|
||||
}
|
||||
tag_object="$(git rev-parse "$tag_ref")"
|
||||
[ "$(git cat-file -t "$tag_object")" = tag ] || {
|
||||
printf 'release tag must be annotated: %s\n' "$TAG" >&2
|
||||
exit 1
|
||||
}
|
||||
tag_contents="$(git cat-file tag "$tag_object")" || {
|
||||
printf 'release tag could not be read: %s\n' "$TAG" >&2
|
||||
exit 1
|
||||
}
|
||||
|
||||
mode="$(
|
||||
printf '%s\n' "$tag_contents" |
|
||||
awk '
|
||||
{
|
||||
line = $0
|
||||
sub(/\r$/, "", line)
|
||||
}
|
||||
line ~ /^OSS-Mirror: / {
|
||||
count++
|
||||
value = substr(line, length("OSS-Mirror: ") + 1)
|
||||
}
|
||||
END {
|
||||
if (count > 1) exit 2
|
||||
if (count == 0) print "enabled"
|
||||
else print value
|
||||
}
|
||||
'
|
||||
)" || {
|
||||
printf 'release tag contains duplicate OSS-Mirror metadata: %s\n' "$TAG" >&2
|
||||
exit 1
|
||||
}
|
||||
case "$mode" in
|
||||
enabled|deferred) printf '%s\n' "$mode" ;;
|
||||
*)
|
||||
printf 'release tag contains invalid OSS-Mirror metadata: %s (%s)\n' "$TAG" "$mode" >&2
|
||||
exit 1
|
||||
;;
|
||||
esac
|
||||
@@ -24,6 +24,7 @@ TOMBSTONE="withdrawn/${VERSION}"
|
||||
GITEE_API="${GITEE_API:-https://gitee.com/api/v5}"
|
||||
GITEE_REPO="${GITEE_REPO:-DingTalk-Real-AI/dingtalk-workspace-cli}"
|
||||
OSS_PREFIX="${OSS_PREFIX:-dws}"
|
||||
OSS_ENABLED=""
|
||||
DELIVERY_VERIFIER="${DWS_DELIVERY_VERIFIER:-$SCRIPT_DIR/verify-release-workflow-delivery.sh}"
|
||||
STABLE_DELIVERY_VERIFIER="${DWS_STABLE_DELIVERY_VERIFIER:-$SCRIPT_DIR/verify-delivered-stable.sh}"
|
||||
GITHUB_DOWNLOAD_HELPER="${DWS_GITHUB_DOWNLOAD_HELPER:-$SCRIPT_DIR/download-github-release-assets.sh}"
|
||||
@@ -404,6 +405,7 @@ prepare_tombstone_message() {
|
||||
printf 'Original-Commit: %s\n' "$TARGET_COMMIT"
|
||||
printf 'Original-Release-ID: %s\n' "$TARGET_RELEASE_ID"
|
||||
printf 'Channel: %s\n' "$CHANNEL"
|
||||
printf 'OSS-Mirror: %s\n' "$TARGET_OSS_MODE"
|
||||
printf 'Reason-SHA256: %s\n' "$reason_sha"
|
||||
printf 'Reason: %s\n' "$REASON"
|
||||
printf 'Requested-By: %s\n' "$GITHUB_ACTOR"
|
||||
@@ -414,7 +416,7 @@ prepare_tombstone_message() {
|
||||
verify_existing_tombstone() {
|
||||
local reason_sha="$1" ref_sha tombstone_json parsed status
|
||||
local actual_target original_tag_object original_commit original_release_id
|
||||
local tombstone_channel tombstone_version tombstone_reason_sha tombstone_reason
|
||||
local tombstone_channel tombstone_version tombstone_oss_mode tombstone_reason_sha tombstone_reason
|
||||
set +e
|
||||
ref_sha="$(
|
||||
github_optional \
|
||||
@@ -478,6 +480,9 @@ if not re.fullmatch(r"[1-9][0-9]*", fields["Original-Release-ID"]):
|
||||
raise SystemExit(1)
|
||||
if not re.fullmatch(r"[0-9a-f]{64}", fields["Reason-SHA256"]):
|
||||
raise SystemExit(1)
|
||||
oss_mode = fields.get("OSS-Mirror", "enabled")
|
||||
if oss_mode not in {"enabled", "deferred"}:
|
||||
raise SystemExit(1)
|
||||
print("\t".join((
|
||||
obj["sha"],
|
||||
fields["Original-Tag-Object"],
|
||||
@@ -485,6 +490,7 @@ print("\t".join((
|
||||
fields["Original-Release-ID"],
|
||||
fields["Channel"],
|
||||
fields["Version"],
|
||||
oss_mode,
|
||||
fields["Reason-SHA256"],
|
||||
fields["Reason"],
|
||||
)))' "$TOMBSTONE"
|
||||
@@ -495,11 +501,13 @@ print("\t".join((
|
||||
original_release_id="$(printf '%s\n' "$parsed" | cut -f4)"
|
||||
tombstone_channel="$(printf '%s\n' "$parsed" | cut -f5)"
|
||||
tombstone_version="$(printf '%s\n' "$parsed" | cut -f6)"
|
||||
tombstone_reason_sha="$(printf '%s\n' "$parsed" | cut -f7)"
|
||||
tombstone_reason="$(printf '%s\n' "$parsed" | cut -f8-)"
|
||||
tombstone_oss_mode="$(printf '%s\n' "$parsed" | cut -f7)"
|
||||
tombstone_reason_sha="$(printf '%s\n' "$parsed" | cut -f8)"
|
||||
tombstone_reason="$(printf '%s\n' "$parsed" | cut -f9-)"
|
||||
[ "$actual_target" = "$original_commit" ] &&
|
||||
[ "$tombstone_version" = "$VERSION" ] &&
|
||||
[ "$tombstone_channel" = "$CHANNEL" ] &&
|
||||
[ "$tombstone_oss_mode" = "${TARGET_OSS_MODE:-$tombstone_oss_mode}" ] &&
|
||||
[ "$tombstone_reason_sha" = "$reason_sha" ] &&
|
||||
[ "$tombstone_reason" = "$REASON" ] ||
|
||||
err "existing tombstone $TOMBSTONE has different immutable withdrawal metadata"
|
||||
@@ -515,6 +523,7 @@ print("\t".join((
|
||||
TARGET_TAG_OBJECT="$original_tag_object"
|
||||
TARGET_COMMIT="$original_commit"
|
||||
TARGET_RELEASE_ID="$original_release_id"
|
||||
TARGET_OSS_MODE="$tombstone_oss_mode"
|
||||
}
|
||||
|
||||
create_tombstone() {
|
||||
@@ -1000,14 +1009,6 @@ for command in git gh npm curl python3 awk sed grep sort unzip ruby cmp cut; do
|
||||
done
|
||||
need_env GITHUB_TOKEN "${GITHUB_TOKEN:-}"
|
||||
need_env NODE_AUTH_TOKEN "${NODE_AUTH_TOKEN:-}"
|
||||
need_env OSS_ACCESS_KEY_ID "${OSS_ACCESS_KEY_ID:-}"
|
||||
need_env OSS_ACCESS_KEY_SECRET "${OSS_ACCESS_KEY_SECRET:-}"
|
||||
need_env OSS_ENDPOINT "${OSS_ENDPOINT:-}"
|
||||
need_env OSS_BUCKET "${OSS_BUCKET:-}"
|
||||
case "$OSS_PREFIX" in
|
||||
''|/*|*'..'*|*'//'*) err "OSS_PREFIX must be a non-empty safe relative prefix" ;;
|
||||
*[!A-Za-z0-9._/-]*) err "OSS_PREFIX contains unsupported characters" ;;
|
||||
esac
|
||||
|
||||
validate_context
|
||||
git fetch --force --tags origin
|
||||
@@ -1018,6 +1019,7 @@ TARGET_COMMIT=""
|
||||
TARGET_RELEASE_ID=""
|
||||
TARGET_TAG_EXISTS=false
|
||||
TARGET_RELEASE_EXISTS=false
|
||||
TARGET_OSS_MODE=""
|
||||
|
||||
set +e
|
||||
remote_tag_object="$(
|
||||
@@ -1038,6 +1040,8 @@ case "$remote_tag_status" in
|
||||
err "$VERSION must be an annotated GitHub release tag"
|
||||
printf '%s\n' "$TARGET_COMMIT" | grep -Eq '^[0-9a-f]{40}$' ||
|
||||
err "could not resolve the release commit for $VERSION"
|
||||
TARGET_OSS_MODE="$("$SCRIPT_DIR/release-tag-oss-mode.sh" "$VERSION")" ||
|
||||
err "could not resolve immutable OSS policy for $VERSION"
|
||||
;;
|
||||
4) ;;
|
||||
*) err "could not inspect the authoritative GitHub tag $VERSION" ;;
|
||||
@@ -1084,6 +1088,22 @@ if [ "$TARGET_TAG_EXISTS" != true ] || [ "$TARGET_RELEASE_EXISTS" != true ]; the
|
||||
say "Resuming withdrawal for $VERSION from exact permanent tombstone metadata."
|
||||
fi
|
||||
|
||||
case "$TARGET_OSS_MODE" in
|
||||
enabled)
|
||||
OSS_ENABLED=true
|
||||
need_env OSS_ACCESS_KEY_ID "${OSS_ACCESS_KEY_ID:-}"
|
||||
need_env OSS_ACCESS_KEY_SECRET "${OSS_ACCESS_KEY_SECRET:-}"
|
||||
need_env OSS_ENDPOINT "${OSS_ENDPOINT:-}"
|
||||
need_env OSS_BUCKET "${OSS_BUCKET:-}"
|
||||
case "$OSS_PREFIX" in
|
||||
''|/*|*'..'*|*'//'*) err "OSS_PREFIX must be a non-empty safe relative prefix" ;;
|
||||
*[!A-Za-z0-9._/-]*) err "OSS_PREFIX contains unsupported characters" ;;
|
||||
esac
|
||||
;;
|
||||
deferred) OSS_ENABLED=false ;;
|
||||
*) err "could not resolve immutable OSS policy for $VERSION" ;;
|
||||
esac
|
||||
|
||||
[ "$(npm_exact_version "$VERSION" || true)" = "${VERSION#v}" ] ||
|
||||
err "npm exact version ${PACKAGE_NAME}@${VERSION#v} is not published"
|
||||
if [ "$TARGET_TAG_EXISTS" = true ] && [ "$TARGET_RELEASE_EXISTS" = true ]; then
|
||||
@@ -1109,23 +1129,28 @@ NPM_POINTER_BEFORE="$(
|
||||
[ -n "$NPM_POINTER_BEFORE" ] || err "npm $NPM_TAG is empty"
|
||||
require_safe_channel_pointer "v$NPM_POINTER_BEFORE" "$CHANNEL"
|
||||
|
||||
ensure_ossutil
|
||||
if [ -z "${OSS_REGION:-}" ]; then
|
||||
OSS_REGION="$(
|
||||
printf '%s' "$OSS_ENDPOINT" |
|
||||
sed -n 's#^\(https\{0,1\}://\)\{0,1\}oss-\([a-z0-9-]*[a-z0-9]\)\.aliyuncs\.com.*#\2#p' |
|
||||
sed 's#-internal$##'
|
||||
)"
|
||||
fi
|
||||
[ -n "$OSS_REGION" ] || err "could not derive OSS_REGION; set it explicitly"
|
||||
export OSS_REGION
|
||||
if [ "$CHANNEL" = stable ]; then
|
||||
OSS_POINTER_NAME=latest.txt
|
||||
if [ "$OSS_ENABLED" = true ]; then
|
||||
ensure_ossutil
|
||||
if [ -z "${OSS_REGION:-}" ]; then
|
||||
OSS_REGION="$(
|
||||
printf '%s' "$OSS_ENDPOINT" |
|
||||
sed -n 's#^\(https\{0,1\}://\)\{0,1\}oss-\([a-z0-9-]*[a-z0-9]\)\.aliyuncs\.com.*#\2#p' |
|
||||
sed 's#-internal$##'
|
||||
)"
|
||||
fi
|
||||
[ -n "$OSS_REGION" ] || err "could not derive OSS_REGION; set it explicitly"
|
||||
export OSS_REGION
|
||||
if [ "$CHANNEL" = stable ]; then
|
||||
OSS_POINTER_NAME=latest.txt
|
||||
else
|
||||
OSS_POINTER_NAME=beta.txt
|
||||
fi
|
||||
OSS_POINTER_BEFORE="$(read_oss_pointer "$OSS_POINTER_NAME")"
|
||||
require_safe_channel_pointer "$OSS_POINTER_BEFORE" "$CHANNEL"
|
||||
else
|
||||
OSS_POINTER_NAME=beta.txt
|
||||
OSS_POINTER_NAME=""
|
||||
say "OSS mirroring is disabled; no OSS rollback is required."
|
||||
fi
|
||||
OSS_POINTER_BEFORE="$(read_oss_pointer "$OSS_POINTER_NAME")"
|
||||
require_safe_channel_pointer "$OSS_POINTER_BEFORE" "$CHANNEL"
|
||||
|
||||
GITEE_ENABLED="${DWS_GITEE_ENABLED:-false}"
|
||||
case "$GITEE_ENABLED" in
|
||||
@@ -1168,8 +1193,12 @@ create_tombstone
|
||||
update_homebrew
|
||||
update_github_release
|
||||
update_npm_channel
|
||||
ensure_oss_rollback
|
||||
update_oss_channel "$OSS_POINTER_NAME"
|
||||
if [ "$OSS_ENABLED" = true ]; then
|
||||
ensure_oss_rollback
|
||||
update_oss_channel "$OSS_POINTER_NAME"
|
||||
else
|
||||
say "OSS did not contain $VERSION because mirroring was disabled; no mirror mutation was needed."
|
||||
fi
|
||||
if [ "$GITEE_ENABLED" = true ]; then
|
||||
ensure_gitee_rollback
|
||||
withdraw_gitee
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"runtime"
|
||||
"sort"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
@@ -836,7 +837,8 @@ func TestReleaseWorkflowPlansAndSealsCurrentMainInTheCloud(t *testing.T) {
|
||||
"persist-credentials: false",
|
||||
`refs/remotes/origin/main)" = "$GITHUB_SHA`,
|
||||
"next-release-version.sh",
|
||||
"refs/tags/withdrawn/v",
|
||||
`'refs/tags/v*' 'refs/tags/withdrawn/v*'`,
|
||||
"release ref manifest is empty after fetching allocated tags",
|
||||
"refs_fingerprint",
|
||||
"Validate the candidate release contract before sealing",
|
||||
"release-contract.sh",
|
||||
@@ -850,6 +852,9 @@ func TestReleaseWorkflowPlansAndSealsCurrentMainInTheCloud(t *testing.T) {
|
||||
if strings.Contains(plan, "contents: write") {
|
||||
t.Error("cloud release planning must remain read-only")
|
||||
}
|
||||
if strings.Contains(plan, "refs/tags/v refs/tags/withdrawn/v") {
|
||||
t.Error("cloud release planning must use wildcard ref patterns that match the seal API prefixes")
|
||||
}
|
||||
|
||||
for _, required := range []string{
|
||||
"name: Seal cloud release tag",
|
||||
@@ -884,6 +889,136 @@ func TestReleaseWorkflowPlansAndSealsCurrentMainInTheCloud(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestReleaseFingerprintRefPatternsMatchAllAllocatedTags(t *testing.T) {
|
||||
t.Parallel()
|
||||
repo := t.TempDir()
|
||||
mustRun(t, repo, "git", "init", "-b", "main")
|
||||
mustRun(t, repo, "git", "config", "user.name", "Release Fingerprint Test")
|
||||
mustRun(t, repo, "git", "config", "user.email", "release-fingerprint@example.com")
|
||||
mustWriteFile(t, filepath.Join(repo, "tracked"), []byte("fixture\n"), 0o644)
|
||||
mustRun(t, repo, "git", "add", "tracked")
|
||||
mustRun(t, repo, "git", "commit", "-m", "fixture")
|
||||
|
||||
allocatedTags := []string{
|
||||
"v1.0.52",
|
||||
"v1.0.53-beta.5",
|
||||
"withdrawn/v1.0.51",
|
||||
}
|
||||
for _, tag := range allocatedTags {
|
||||
mustRun(t, repo, "git", "tag", "-a", tag, "-m", "Release "+tag)
|
||||
}
|
||||
mustRun(t, repo, "git", "tag", "-a", "release/v1.0.52", "-m", "unrelated namespace")
|
||||
legacy := exec.Command(
|
||||
"git", "for-each-ref", "--format=%(refname)=%(objectname)",
|
||||
"refs/tags/v", "refs/tags/withdrawn/v",
|
||||
)
|
||||
legacy.Dir = repo
|
||||
legacyOutput, err := legacy.CombinedOutput()
|
||||
if err != nil {
|
||||
t.Fatalf("legacy git for-each-ref error = %v\noutput:\n%s", err, legacyOutput)
|
||||
}
|
||||
if strings.TrimSpace(string(legacyOutput)) != "" {
|
||||
t.Fatalf("legacy component patterns unexpectedly matched flat release refs:\n%s", legacyOutput)
|
||||
}
|
||||
|
||||
cmd := exec.Command(
|
||||
"git", "for-each-ref", "--format=%(refname)=%(objectname)",
|
||||
"refs/tags/v*", "refs/tags/withdrawn/v*",
|
||||
)
|
||||
cmd.Dir = repo
|
||||
output, err := cmd.CombinedOutput()
|
||||
if err != nil {
|
||||
t.Fatalf("git for-each-ref error = %v\noutput:\n%s", err, output)
|
||||
}
|
||||
got := strings.Split(strings.TrimSpace(string(output)), "\n")
|
||||
sort.Strings(got)
|
||||
|
||||
want := make([]string, 0, len(allocatedTags))
|
||||
for _, tag := range allocatedTags {
|
||||
object := strings.TrimSpace(mustOutput(t, repo, "git", "rev-parse", "refs/tags/"+tag))
|
||||
want = append(want, "refs/tags/"+tag+"="+object)
|
||||
}
|
||||
sort.Strings(want)
|
||||
if strings.Join(got, "\n") != strings.Join(want, "\n") {
|
||||
t.Fatalf("release ref set differs from the seal API set\ngot:\n%s\nwant:\n%s", strings.Join(got, "\n"), strings.Join(want, "\n"))
|
||||
}
|
||||
workflowDigest := sha256.Sum256([]byte(strings.Join(got, "\n") + "\n"))
|
||||
sealDigest := sha256.Sum256([]byte(strings.Join(want, "\n") + "\n"))
|
||||
if workflowDigest != sealDigest {
|
||||
t.Fatalf("release ref fingerprint differs from seal fingerprint: workflow=%x seal=%x", workflowDigest, sealDigest)
|
||||
}
|
||||
}
|
||||
|
||||
func TestReleaseWorkflowAcceptsGuardedLocalTagMetadata(t *testing.T) {
|
||||
t.Parallel()
|
||||
workflow := readReleaseWorkflow(t)
|
||||
|
||||
for _, required := range []string{
|
||||
`const cloudOnlyKeys = [`,
|
||||
`const cloudKeys = ["Channel", ...cloudOnlyKeys];`,
|
||||
`const hasAnyCloudMetadata = cloudOnlyKeys.some((key) => tagFields.has(key));`,
|
||||
`const isCloudSeal = cloudKeys.every((key) => tagFields.has(key));`,
|
||||
} {
|
||||
if !strings.Contains(workflow, required) {
|
||||
t.Errorf("local tag metadata compatibility is missing %q", required)
|
||||
}
|
||||
}
|
||||
if strings.Contains(workflow, `const hasAnyCloudMetadata = cloudKeys.some((key) => tagFields.has(key));`) {
|
||||
t.Error("Channel-only guarded local tags must not be classified as partial cloud seals")
|
||||
}
|
||||
}
|
||||
|
||||
func TestReleaseWorkflowRequiresOSSOnlyWhenMirrorIsEnabled(t *testing.T) {
|
||||
t.Parallel()
|
||||
workflow := readReleaseWorkflow(t)
|
||||
releaseContract := releaseWorkflowSection(t, workflow, " release-contract:\n", "\n release:\n")
|
||||
targetAuthority := releaseWorkflowSection(
|
||||
t,
|
||||
releaseContract,
|
||||
" - name: Resolve and verify exact release target\n",
|
||||
"\n - name: Check out repository\n",
|
||||
)
|
||||
ossStep := releaseWorkflowSection(
|
||||
t,
|
||||
workflow,
|
||||
" - name: Sync release artifacts to China OSS mirror\n",
|
||||
"\n mirror-gitee-release:\n",
|
||||
)
|
||||
|
||||
for _, required := range []string{
|
||||
`if: ${{ needs.release-contract.outputs.oss_mirror == 'enabled' }}`,
|
||||
`run: ./scripts/release/sync-to-oss.sh`,
|
||||
`DWS_REQUIRE_OSS: "1"`,
|
||||
} {
|
||||
if !strings.Contains(ossStep, required) {
|
||||
t.Errorf("opt-in OSS publication is missing %q", required)
|
||||
}
|
||||
}
|
||||
for _, required := range []string{
|
||||
`OSS_MIRROR: ${{ vars.ENABLE_OSS_MIRROR == 'true' && 'enabled' || 'deferred' }}`,
|
||||
`OSS-Mirror: ${ossMirror}`,
|
||||
`core.setOutput("oss_mirror", ossMirror);`,
|
||||
} {
|
||||
if !strings.Contains(workflow, required) {
|
||||
t.Errorf("immutable OSS release policy is missing %q", required)
|
||||
}
|
||||
}
|
||||
if strings.Contains(ossStep, "vars.ENABLE_OSS_MIRROR") {
|
||||
t.Error("channel publication must use the immutable tag policy, not the current repository variable")
|
||||
}
|
||||
for _, required := range []string{
|
||||
`const ossMirror = tagFields.has("OSS-Mirror")`,
|
||||
`? tagFields.get("OSS-Mirror")`,
|
||||
`: "enabled";`,
|
||||
`!["enabled", "deferred"].includes(ossMirror)`,
|
||||
`core.setOutput("oss_mirror", ossMirror);`,
|
||||
} {
|
||||
if !strings.Contains(targetAuthority, required) {
|
||||
t.Errorf("release target OSS policy authority is missing %q", required)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestReleaseWorkflowChannelRepairUsesSealedReleaseAuthority(t *testing.T) {
|
||||
t.Parallel()
|
||||
workflow := readReleaseWorkflow(t)
|
||||
@@ -897,6 +1032,12 @@ func TestReleaseWorkflowChannelRepairUsesSealedReleaseAuthority(t *testing.T) {
|
||||
t.Fatal("release workflow channel repair job is missing its end marker")
|
||||
}
|
||||
repair := workflow[start : start+end]
|
||||
authority := releaseWorkflowSection(
|
||||
t,
|
||||
repair,
|
||||
" - name: Verify immutable release authority\n",
|
||||
"\n - name: Require sealed OSS policy for repair\n",
|
||||
)
|
||||
tagAuthority := releaseWorkflowSection(
|
||||
t,
|
||||
repair,
|
||||
@@ -945,6 +1086,8 @@ func TestReleaseWorkflowChannelRepairUsesSealedReleaseAuthority(t *testing.T) {
|
||||
"assetNames.length !== expectedAssets.size",
|
||||
"new Set(assetNames).size !== expectedAssets.size",
|
||||
`core.setOutput("tag_object", ref.data.object.sha)`,
|
||||
"Require sealed OSS policy for repair",
|
||||
"OSS repair is unavailable because this immutable release deferred the OSS channel.",
|
||||
`ref: ${{ steps.authority.outputs.commit_sha }}`,
|
||||
"path: release-source",
|
||||
"verify-github-tag-authority.sh",
|
||||
@@ -975,6 +1118,25 @@ func TestReleaseWorkflowChannelRepairUsesSealedReleaseAuthority(t *testing.T) {
|
||||
t.Errorf("channel repair authority is missing %q", required)
|
||||
}
|
||||
}
|
||||
for _, required := range []string{
|
||||
`const tagFields = new Map();`,
|
||||
`const ossMirror = tagFields.has("OSS-Mirror")`,
|
||||
`? tagFields.get("OSS-Mirror")`,
|
||||
`: "enabled";`,
|
||||
`!["enabled", "deferred"].includes(ossMirror)`,
|
||||
`core.setOutput("oss_mirror", ossMirror);`,
|
||||
} {
|
||||
if !strings.Contains(authority, required) {
|
||||
t.Errorf("channel repair tag policy authority is missing %q", required)
|
||||
}
|
||||
}
|
||||
npmRepair := releaseWorkflowSection(t, workflow, " repair-npm:\n", "\n release-delivery-gate:\n")
|
||||
if strings.Contains(npmRepair, `const ossMirror`) || strings.Contains(npmRepair, `core.setOutput("oss_mirror"`) {
|
||||
t.Error("npm repair must not parse or export the channel-only OSS policy")
|
||||
}
|
||||
if strings.Contains(repair, "ENABLE_OSS_MIRROR") {
|
||||
t.Error("OSS repair must use the immutable tag policy, not the current repository variable")
|
||||
}
|
||||
for _, asset := range []string{
|
||||
"dws-darwin-amd64.tar.gz",
|
||||
"dws-darwin-arm64.tar.gz",
|
||||
|
||||
@@ -943,6 +943,11 @@ esac
|
||||
if output, err := run(); err != nil || !strings.Contains(output, "cloud release run 42") {
|
||||
t.Fatalf("tag-bound cloud delivery was rejected: err=%v\noutput:\n%s", err, output)
|
||||
}
|
||||
if output, err := runWithArgs(
|
||||
[]string{"--channel-repair", "oss", tag, commit},
|
||||
); err != nil || !strings.Contains(output, "cloud release run 42") {
|
||||
t.Fatalf("OSS repair could not use a successful release with the mirror enabled: err=%v\noutput:\n%s", err, output)
|
||||
}
|
||||
if output, err := run("RUN_ACTOR=renamed-release-user"); err != nil ||
|
||||
!strings.Contains(output, "cloud release run 42") {
|
||||
t.Fatalf("cloud delivery broke after a harmless login rename: err=%v\noutput:\n%s", err, output)
|
||||
|
||||
@@ -82,6 +82,11 @@ func TestWithdrawReleaseScriptDeletesProblemReleaseLastAndUsesPermanentTombstone
|
||||
`github_api --method PATCH`,
|
||||
`npm deprecate "${PACKAGE_NAME}@${VERSION#v}"`,
|
||||
`npm dist-tag add "${PACKAGE_NAME}@${ROLLBACK_VERSION#v}"`,
|
||||
`TARGET_OSS_MODE="$("$SCRIPT_DIR/release-tag-oss-mode.sh" "$VERSION")"`,
|
||||
`printf 'OSS-Mirror: %s\n' "$TARGET_OSS_MODE"`,
|
||||
`oss_mode = fields.get("OSS-Mirror", "enabled")`,
|
||||
`if [ "$OSS_ENABLED" = true ]; then`,
|
||||
`*) err "could not resolve immutable OSS policy for $VERSION" ;;`,
|
||||
`"$OSSUTIL" rm -rf`,
|
||||
`curl -fsS -X DELETE`,
|
||||
`git push "$GITEE_GIT_REMOTE" ":refs/tags/${VERSION}"`,
|
||||
@@ -114,7 +119,7 @@ func TestWithdrawReleaseScriptDeletesProblemReleaseLastAndUsesPermanentTombstone
|
||||
tombstone := strings.LastIndex(script, "\ncreate_tombstone\n")
|
||||
githubMutation := strings.LastIndex(script, "\nupdate_github_release\n")
|
||||
npmMutation := strings.LastIndex(script, "\nupdate_npm_channel\n")
|
||||
ossMutation := strings.LastIndex(script, "\nupdate_oss_channel ")
|
||||
ossMutation := strings.LastIndex(script, `update_oss_channel "$OSS_POINTER_NAME"`)
|
||||
homebrewGate := strings.LastIndex(script, "\nupdate_homebrew\n")
|
||||
githubDelete := strings.LastIndex(script, "\ndelete_github_release_and_tag\n")
|
||||
if tombstone < 0 || githubMutation < 0 || npmMutation < 0 || ossMutation < 0 ||
|
||||
@@ -128,6 +133,122 @@ func TestWithdrawReleaseScriptDeletesProblemReleaseLastAndUsesPermanentTombstone
|
||||
}
|
||||
}
|
||||
|
||||
func TestReleaseTagOSSModeIsImmutableAndFailClosed(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
scriptPath, err := filepath.Abs(filepath.Join("..", "..", "scripts", "release", "release-tag-oss-mode.sh"))
|
||||
if err != nil {
|
||||
t.Fatalf("Abs(OSS mode script) error = %v", err)
|
||||
}
|
||||
repo := t.TempDir()
|
||||
mustRun(t, repo, "git", "init", "-b", "main")
|
||||
mustRun(t, repo, "git", "config", "user.name", "OSS Mode Test")
|
||||
mustRun(t, repo, "git", "config", "user.email", "oss-mode@example.com")
|
||||
mustWriteFile(t, filepath.Join(repo, "tracked"), []byte("fixture\n"), 0o644)
|
||||
mustRun(t, repo, "git", "add", "tracked")
|
||||
mustRun(t, repo, "git", "commit", "-m", "fixture")
|
||||
|
||||
tag := func(version, message string) {
|
||||
t.Helper()
|
||||
messagePath := filepath.Join(repo, version+".message")
|
||||
mustWriteFile(t, messagePath, []byte(message), 0o644)
|
||||
mustRun(t, repo, "git", "tag", "-a", version, "-F", messagePath)
|
||||
}
|
||||
tag("v1.0.1", "Release v1.0.1\n")
|
||||
tag("v1.0.2", "Release v1.0.2\n\nOSS-Mirror: enabled\n")
|
||||
tag("v1.0.3", "Release v1.0.3\n\nOSS-Mirror: deferred\n")
|
||||
tag("v1.0.4", "Release v1.0.4\n\nOSS-Mirror: invalid\n")
|
||||
tag("v1.0.5", "Release v1.0.5\n\nOSS-Mirror: enabled\nOSS-Mirror: deferred\n")
|
||||
mustRun(t, repo, "git", "tag", "v1.0.6")
|
||||
rawTag := func(version, message string) {
|
||||
t.Helper()
|
||||
commit := strings.TrimSpace(mustOutput(t, repo, "git", "rev-parse", "HEAD"))
|
||||
payload := strings.Join([]string{
|
||||
"object " + commit,
|
||||
"type commit",
|
||||
"tag " + version,
|
||||
"tagger OSS Mode Test <oss-mode@example.com> 1700000000 +0000",
|
||||
"",
|
||||
message,
|
||||
}, "\n")
|
||||
cmd := exec.Command("git", "mktag")
|
||||
cmd.Dir = repo
|
||||
cmd.Stdin = strings.NewReader(payload)
|
||||
output, err := cmd.CombinedOutput()
|
||||
if err != nil {
|
||||
t.Fatalf("git mktag %s error = %v\noutput:\n%s", version, err, output)
|
||||
}
|
||||
mustRun(t, repo, "git", "update-ref", "refs/tags/"+version, strings.TrimSpace(string(output)))
|
||||
}
|
||||
rawTag("v1.0.7", "Release v1.0.7\n\nOSS-Mirror: \n")
|
||||
rawTag("v1.0.8", "Release v1.0.8\r\n\r\nOSS-Mirror: deferred\r\n")
|
||||
|
||||
for _, test := range []struct {
|
||||
version string
|
||||
want string
|
||||
ok bool
|
||||
}{
|
||||
{version: "v1.0.1", want: "enabled", ok: true},
|
||||
{version: "v1.0.2", want: "enabled", ok: true},
|
||||
{version: "v1.0.3", want: "deferred", ok: true},
|
||||
{version: "v1.0.4", want: "invalid OSS-Mirror metadata", ok: false},
|
||||
{version: "v1.0.5", want: "duplicate OSS-Mirror metadata", ok: false},
|
||||
{version: "v1.0.6", want: "must be annotated", ok: false},
|
||||
{version: "v1.0.7", want: "invalid OSS-Mirror metadata", ok: false},
|
||||
{version: "v1.0.8", want: "deferred", ok: true},
|
||||
} {
|
||||
t.Run(test.version, func(t *testing.T) {
|
||||
cmd := exec.Command(scriptPath, test.version)
|
||||
cmd.Dir = repo
|
||||
output, err := cmd.CombinedOutput()
|
||||
if test.ok && err != nil {
|
||||
t.Fatalf("release-tag-oss-mode.sh error = %v\noutput:\n%s", err, output)
|
||||
}
|
||||
if !test.ok && err == nil {
|
||||
t.Fatalf("release-tag-oss-mode.sh unexpectedly accepted %s", test.version)
|
||||
}
|
||||
if !strings.Contains(string(output), test.want) {
|
||||
t.Fatalf("release-tag-oss-mode.sh output missing %q:\n%s", test.want, output)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithdrawReleaseDefersOSSRequirementUntilImmutablePolicyResolution(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
scriptPath, err := filepath.Abs(filepath.Join("..", "..", "scripts", "release", "withdraw-release.sh"))
|
||||
if err != nil {
|
||||
t.Fatalf("Abs(withdraw script) error = %v", err)
|
||||
}
|
||||
cmd := exec.Command("bash", scriptPath,
|
||||
"v1.2.3",
|
||||
"critical startup regression",
|
||||
"WITHDRAW v1.2.3",
|
||||
)
|
||||
cmd.Env = append(os.Environ(),
|
||||
"GITHUB_TOKEN=test-github-token",
|
||||
"NODE_AUTH_TOKEN=test-npm-token",
|
||||
"GITHUB_REPOSITORY=DingTalk-Real-AI/dingtalk-workspace-cli",
|
||||
"GITHUB_REF_NAME=main",
|
||||
"GITHUB_SHA="+strings.Repeat("a", 40),
|
||||
"GITHUB_RUN_ID=1",
|
||||
"GITHUB_ACTOR=test-operator",
|
||||
"GITHUB_EVENT_DEFAULT_BRANCH=main",
|
||||
"GITHUB_ACTIONS=false",
|
||||
)
|
||||
output, err := cmd.CombinedOutput()
|
||||
if err == nil {
|
||||
t.Fatalf("withdraw-release.sh unexpectedly ran outside GitHub Actions:\n%s", output)
|
||||
}
|
||||
if !strings.Contains(string(output), "withdrawal may run only inside GitHub Actions") {
|
||||
t.Fatalf("withdrawal did not reach the protected execution gate before resolving OSS policy:\n%s", output)
|
||||
}
|
||||
if strings.Contains(string(output), "missing required environment variable OSS_") {
|
||||
t.Fatalf("withdrawal required OSS credentials before resolving immutable tag policy:\n%s", output)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWithdrawReleaseRejectsInvalidInputsBeforeMutation(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
@@ -180,6 +301,20 @@ func TestWithdrawReleaseRejectsInvalidInputsBeforeMutation(t *testing.T) {
|
||||
}
|
||||
|
||||
func TestWithdrawReleaseRollsBackConfiguredChannelsAndStopsForHomebrewReview(t *testing.T) {
|
||||
for _, test := range []struct {
|
||||
name string
|
||||
ossDeferred bool
|
||||
}{
|
||||
{name: "legacy tag defaults OSS to enabled"},
|
||||
{name: "deferred OSS survives tombstone retry", ossDeferred: true},
|
||||
} {
|
||||
t.Run(test.name, func(t *testing.T) {
|
||||
testWithdrawReleaseRollsBackConfiguredChannelsAndStopsForHomebrewReview(t, test.ossDeferred)
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func testWithdrawReleaseRollsBackConfiguredChannelsAndStopsForHomebrewReview(t *testing.T, ossDeferred bool) {
|
||||
root := t.TempDir()
|
||||
sourceRoot, err := filepath.Abs(filepath.Join("..", ".."))
|
||||
if err != nil {
|
||||
@@ -188,6 +323,7 @@ func TestWithdrawReleaseRollsBackConfiguredChannelsAndStopsForHomebrewReview(t *
|
||||
|
||||
for _, rel := range []string{
|
||||
"scripts/release/withdraw-release.sh",
|
||||
"scripts/release/release-tag-oss-mode.sh",
|
||||
"scripts/release/release-lib.sh",
|
||||
"build/homebrew-release.rb.tmpl",
|
||||
} {
|
||||
@@ -235,6 +371,10 @@ printf 'gitee-sync %s\n' "$VERSION" >> "$CALL_LOG"
|
||||
cp "$DWS_FORMULA_SOURCE" "$MOCK_STATE/homebrew-formula"
|
||||
printf 'homebrew-pr %s %s\n' "$DWS_TAP_PR_TITLE" "$DWS_TAP_PR_BRANCH" >> "$CALL_LOG"
|
||||
exit 0
|
||||
`), 0o755)
|
||||
mustWriteFile(t, filepath.Join(root, "scripts", "release", "unexpected-ossutil"), []byte(`#!/bin/sh
|
||||
printf 'unexpected-oss-call\n' >> "$CALL_LOG"
|
||||
exit 99
|
||||
`), 0o755)
|
||||
|
||||
stateDir := filepath.Join(root, "state")
|
||||
@@ -262,7 +402,11 @@ end
|
||||
`), 0o644)
|
||||
mustRun(t, root, "git", "add", "Formula/dingtalk-workspace-cli.rb")
|
||||
mustRun(t, root, "git", "commit", "-m", "target")
|
||||
mustRun(t, root, "git", "tag", "-a", "v1.0.52", "-m", "Release v1.0.52")
|
||||
tagArgs := []string{"tag", "-a", "v1.0.52", "-m", "Release v1.0.52"}
|
||||
if ossDeferred {
|
||||
tagArgs = append(tagArgs, "-m", "OSS-Mirror: deferred")
|
||||
}
|
||||
mustRun(t, root, "git", tagArgs...)
|
||||
targetCommit := strings.TrimSpace(mustOutput(t, root, "git", "rev-parse", "HEAD"))
|
||||
targetTagObject := strings.TrimSpace(mustOutput(t, root, "git", "rev-parse", "refs/tags/v1.0.52"))
|
||||
|
||||
@@ -308,11 +452,6 @@ end
|
||||
"GITHUB_EVENT_DEFAULT_BRANCH=main",
|
||||
"GITHUB_TOKEN=github-token",
|
||||
"NODE_AUTH_TOKEN=npm-token",
|
||||
"OSS_ACCESS_KEY_ID=oss-id",
|
||||
"OSS_ACCESS_KEY_SECRET=oss-secret",
|
||||
"OSS_ENDPOINT=https://oss-cn-hangzhou.aliyuncs.com",
|
||||
"OSS_BUCKET=dws-test",
|
||||
"OSSUTIL="+filepath.Join(fakeBin, "ossutil"),
|
||||
"DWS_GITEE_ENABLED=true",
|
||||
"GITEE_TOKEN=gitee-token",
|
||||
"GITEE_USER=gitee-user",
|
||||
@@ -326,6 +465,23 @@ end
|
||||
"DWS_GITEE_SYNC_HELPER="+filepath.Join(root, "scripts", "release", "sync-gitee"),
|
||||
"ORIGIN_GIT="+origin,
|
||||
)
|
||||
if ossDeferred {
|
||||
withdrawEnv = append(withdrawEnv,
|
||||
"OSS_ACCESS_KEY_ID=",
|
||||
"OSS_ACCESS_KEY_SECRET=",
|
||||
"OSS_ENDPOINT=",
|
||||
"OSS_BUCKET=",
|
||||
"OSSUTIL="+filepath.Join(root, "scripts", "release", "unexpected-ossutil"),
|
||||
)
|
||||
} else {
|
||||
withdrawEnv = append(withdrawEnv,
|
||||
"OSS_ACCESS_KEY_ID=oss-id",
|
||||
"OSS_ACCESS_KEY_SECRET=oss-secret",
|
||||
"OSS_ENDPOINT=https://oss-cn-hangzhou.aliyuncs.com",
|
||||
"OSS_BUCKET=dws-test",
|
||||
"OSSUTIL="+filepath.Join(fakeBin, "ossutil"),
|
||||
)
|
||||
}
|
||||
cmd.Env = withdrawEnv
|
||||
output, err := cmd.CombinedOutput()
|
||||
if err == nil {
|
||||
@@ -338,9 +494,13 @@ end
|
||||
assertFileEquals(t, filepath.Join(stateDir, "latest"), "v1.0.51")
|
||||
assertFileEquals(t, filepath.Join(stateDir, "npm-latest"), "1.0.51")
|
||||
assertFileContains(t, filepath.Join(stateDir, "npm-deprecated"), "WITHDRAWN v1.0.52")
|
||||
assertFileEquals(t, filepath.Join(stateDir, "oss-latest"), "v1.0.51")
|
||||
if _, err := os.Stat(filepath.Join(stateDir, "oss-removed")); err != nil {
|
||||
t.Fatalf("OSS withdrawn prefix was not removed: %v", err)
|
||||
if ossDeferred {
|
||||
assertFileEquals(t, filepath.Join(stateDir, "oss-latest"), "v1.0.52")
|
||||
} else {
|
||||
assertFileEquals(t, filepath.Join(stateDir, "oss-latest"), "v1.0.51")
|
||||
if _, err := os.Stat(filepath.Join(stateDir, "oss-removed")); err != nil {
|
||||
t.Fatalf("OSS withdrawn prefix was not removed: %v", err)
|
||||
}
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(stateDir, "gitee-release-v1.0.52")); !os.IsNotExist(err) {
|
||||
t.Fatalf("Gitee release still exists, stat error = %v", err)
|
||||
@@ -357,6 +517,11 @@ end
|
||||
assertFileContains(t, filepath.Join(stateDir, "tombstone-message"), "Original-Commit: "+targetCommit)
|
||||
assertFileContains(t, filepath.Join(stateDir, "tombstone-message"), "Original-Tag-Object: "+targetTagObject)
|
||||
assertFileContains(t, filepath.Join(stateDir, "tombstone-message"), "Original-Release-ID: 52")
|
||||
ossMode := "enabled"
|
||||
if ossDeferred {
|
||||
ossMode = "deferred"
|
||||
}
|
||||
assertFileContains(t, filepath.Join(stateDir, "tombstone-message"), "OSS-Mirror: "+ossMode)
|
||||
assertFileContains(t, filepath.Join(stateDir, "tombstone-message"), "Reason: critical startup regression")
|
||||
assertFileContains(t, callLog, "homebrew-pr revert: withdraw v1.0.52 and restore v1.0.51")
|
||||
|
||||
@@ -369,20 +534,49 @@ end
|
||||
githubIndex := strings.Index(logText, "github-withdraw")
|
||||
npmIndex := strings.Index(logText, "npm-deprecate")
|
||||
ossIndex := strings.Index(logText, "oss-remove")
|
||||
if tombstoneIndex < 0 || githubIndex < 0 || npmIndex < 0 || ossIndex < 0 {
|
||||
if tombstoneIndex < 0 || githubIndex < 0 || npmIndex < 0 || (!ossDeferred && ossIndex < 0) {
|
||||
t.Fatalf("missing mutation audit entries:\n%s", logText)
|
||||
}
|
||||
homebrewIndex := strings.Index(logText, "homebrew-pr")
|
||||
if !(tombstoneIndex < homebrewIndex && homebrewIndex < githubIndex &&
|
||||
githubIndex < npmIndex && npmIndex < ossIndex) {
|
||||
if !(tombstoneIndex < homebrewIndex && homebrewIndex < githubIndex && githubIndex < npmIndex) ||
|
||||
(!ossDeferred && npmIndex >= ossIndex) {
|
||||
t.Fatalf("tombstone was not durable before channel mutations:\n%s", logText)
|
||||
}
|
||||
deleteReleaseIndex := strings.Index(logText, "github-delete-release")
|
||||
deleteTagIndex := strings.Index(logText, "github-delete-tag")
|
||||
channelsCompleteIndex := npmIndex
|
||||
if !ossDeferred {
|
||||
channelsCompleteIndex = ossIndex
|
||||
}
|
||||
if deleteReleaseIndex < 0 || deleteTagIndex < 0 || homebrewIndex < 0 ||
|
||||
!(homebrewIndex < ossIndex && ossIndex < deleteReleaseIndex && deleteReleaseIndex < deleteTagIndex) {
|
||||
!(homebrewIndex < channelsCompleteIndex && channelsCompleteIndex < deleteReleaseIndex && deleteReleaseIndex < deleteTagIndex) {
|
||||
t.Fatalf("GitHub problem release was not removed before the Homebrew review pause:\n%s", logText)
|
||||
}
|
||||
if ossDeferred {
|
||||
if strings.Contains(logText, "unexpected-oss-call") || ossIndex >= 0 {
|
||||
t.Fatalf("deferred withdrawal invoked OSS tooling:\n%s", logText)
|
||||
}
|
||||
for _, want := range []string{
|
||||
"OSS mirroring is disabled; no OSS rollback is required.",
|
||||
"OSS did not contain v1.0.52 because mirroring was disabled",
|
||||
} {
|
||||
if !strings.Contains(string(output), want) {
|
||||
t.Fatalf("deferred withdrawal output missing %q:\n%s", want, output)
|
||||
}
|
||||
}
|
||||
}
|
||||
if !ossDeferred {
|
||||
tombstonePath := filepath.Join(stateDir, "tombstone-message")
|
||||
tombstoneMessage, err := os.ReadFile(tombstonePath)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadFile(legacy tombstone) error = %v", err)
|
||||
}
|
||||
legacyMessage := strings.Replace(string(tombstoneMessage), "OSS-Mirror: enabled\n", "", 1)
|
||||
if legacyMessage == string(tombstoneMessage) {
|
||||
t.Fatal("could not convert tombstone fixture to the legacy format")
|
||||
}
|
||||
mustWriteFile(t, tombstonePath, []byte(legacyMessage), 0o644)
|
||||
}
|
||||
|
||||
mergedFormula, err := os.ReadFile(filepath.Join(stateDir, "homebrew-formula"))
|
||||
if err != nil {
|
||||
@@ -420,6 +614,9 @@ end
|
||||
t.Fatalf("withdrawal retry output missing %q:\n%s", want, string(retryOutput))
|
||||
}
|
||||
}
|
||||
if ossDeferred && !strings.Contains(string(retryOutput), "OSS mirroring is disabled; no OSS rollback is required.") {
|
||||
t.Fatalf("deferred tombstone retry did not restore the sealed OSS policy:\n%s", retryOutput)
|
||||
}
|
||||
assertFileEquals(t, filepath.Join(stateDir, "latest"), "v1.0.51")
|
||||
if _, err := os.Stat(filepath.Join(stateDir, "release-v1.0.52.json")); !os.IsNotExist(err) {
|
||||
t.Fatalf("GitHub problem release still exists, stat error = %v", err)
|
||||
|
||||
Reference in New Issue
Block a user