Compare commits

..
Author SHA1 Message Date
南润 4987aae380 Merge remote-tracking branch 'origin/main' into calendar_skill_opt1
# Conflicts:
#	skills/multi/dingtalk-calendar/SKILL.md
#	skills/multi/dingtalk-calendar/references/03-meeting.md
#	skills/multi/dingtalk-calendar/references/lite-recipes.md
2026-08-19 10:51:54 +08:00
南润 c8c997fff3 refactor(calendar): align skill with golden routes 2026-08-19 10:43:17 +08:00
49 changed files with 116 additions and 3906 deletions
-26
View File
@@ -1,26 +0,0 @@
---
category: Added
---
- **Wait framework capability** — adds the reviewed `Contract.Wait`
declaration (`contract.WaitSpec`) with three execution modes: `poll`
(cadence-poll the leaf's `WaitPoll` hook), `event` (consume the leaf's
`WaitEvents` push stream, correlate events to the accepted resource via
`match_field`/`resource_query`, apply the same terminal map), and `auto`
(event first, fall back to polling when the stream ends or the
subscription fails — one deadline spans both phases). Declared commands
must use the `ResultInvoke` dispatcher; mode and hooks are paired at
construction (poll↔WaitPoll, event↔WaitEvents, auto↔both; surplus hooks
are rejected too). Declared commands register `--wait` /
`--wait-timeout` (framework-owned flags that never enter MCP toolArgs);
undeclared commands reject the flags as unknown. The wait phase closes
the unified envelope exactly once: terminal success → `success`,
terminal failure → `failure` with new wire-stable `error.type: "wait"`
(exit code 8), timeout → `pending` with `meta.operation.timed_out: true`
and the last observed state (exit 0). Deadline exhaustion during a poll,
during event consumption, or between polls always closes as timed-out
pending, never as a poll/stream failure; a correlated event with an
unknown status fails closed exactly like a poll. The capability is
projected into the Schema catalog (`wait` key) alongside `dry_run`. No
business command declares it yet; approval/export/batch adoption lands
separately.
+3
View File
@@ -73,3 +73,6 @@ coverage-*.txt
# stray compiled generator binary (source lives in internal/generator/cmd_param_aliases/)
/cmd_param_aliases
# Local product-skill design materials (not for repository pushes)
/design-dws-product-skills/
+2 -3
View File
@@ -482,7 +482,7 @@ Env vars: `DWS_SKILL_MODE=mono|multi` (also honored by `install.sh` / `install.p
<details>
<summary><strong>Personal Event Subscription</strong> — real-time DingTalk messages for event-driven agents</summary>
`dws event consume` subscribes as the currently logged-in user over a managed Stream WebSocket and emits each event as one NDJSON line on stdout. The public catalog covers scoped and all one-to-one/group messages, specified senders, read/recall/reaction events, group lifecycle events, and seven OA approval task/instance events.
`dws event consume` subscribes as the currently logged-in user over a managed Stream WebSocket and emits each event as one NDJSON line on stdout. The public catalog covers scoped and all one-to-one/group messages, specified senders, read/recall/reaction events, group lifecycle events, and six OA approval task/instance events.
The default `ndjson`, `json`, and `pretty` output preserves the transport envelope (`type`, `event_type`, string `data`, and `headers`) for existing scripts; `compact` retains its existing processor. Add `--flatten` to emit the stable top-level business fields used by Agent workflows. `--format` controls JSON serialization; `--flatten` controls the data structure and cannot be combined with `-f raw` or `--debug-raw-events`.
@@ -530,13 +530,12 @@ dws event consume user_im_group_disbanded --group <openConversationId> --flatten
dws event +listen-im --kind sender --user <userId> \
--events message,read,recall -f ndjson
# Listen for all seven public OA approval events in one process
# Listen for all six public OA approval events in one process
dws event consume \
user_oa_approval_task_created \
user_oa_approval_task_finished \
user_oa_approval_task_redirected \
user_oa_approval_instance_started \
user_oa_approval_instance_cc \
user_oa_approval_instance_terminated \
user_oa_approval_instance_finished \
--flatten -f ndjson
+2 -3
View File
@@ -476,7 +476,7 @@ multi setup 或 upgrade 后,DWS 会把官方 bundle 快照和统一所有权
<details>
<summary><strong>个人事件订阅</strong> — 实时接收钉钉消息,驱动事件触发的 Agent</summary>
`dws event consume` 使用当前 OAuth 登录用户建立托管的 Stream WebSocket 长连接,并把每条事件以 NDJSON 一行输出到 stdout。当前公开目录覆盖指定范围和全量单聊/群消息、指定发送人、已读/撤回/表情回应、群生命周期,以及七个 OA 审批任务/实例事件。
`dws event consume` 使用当前 OAuth 登录用户建立托管的 Stream WebSocket 长连接,并把每条事件以 NDJSON 一行输出到 stdout。当前公开目录覆盖指定范围和全量单聊/群消息、指定发送人、已读/撤回/表情回应、群生命周期,以及六个 OA 审批任务/实例事件。
默认 `ndjson`、`json`、`pretty` 输出保留兼容 transport envelope(`type`、`event_type`、字符串 `data`、`headers`),`compact` 继续沿用原 processor。Agent 或新脚本显式加 `--flatten` 后,输出稳定的顶层业务字段。`--format` 控制 JSON 序列化,`--flatten` 控制数据结构,且不能与 `-f raw` 或 `--debug-raw-events` 同时使用。
@@ -524,13 +524,12 @@ dws event consume user_im_group_disbanded --group <openConversationId> --flatten
dws event +listen-im --kind sender --user <userId> \
--events message,read,recall -f ndjson
# 一个进程监听全部七个公开 OA 审批事件
# 一个进程监听全部六个公开 OA 审批事件
dws event consume \
user_oa_approval_task_created \
user_oa_approval_task_finished \
user_oa_approval_task_redirected \
user_oa_approval_instance_started \
user_oa_approval_instance_cc \
user_oa_approval_instance_terminated \
user_oa_approval_instance_finished \
--flatten -f ndjson
@@ -360,7 +360,6 @@ Definition(仅声明;不可编译)
| | `idempotency` | 评审源(或未来 Contract) | reviewed metadata | 今日非框架声明;不得推断 |
| | `effect_source` / provenance | 组装派生物 | resolver 写入 `FieldProvenance` | 派生,不手写 |
| **DryRun** | `preview_kind`, `remote_reads` | 评审源 | `schema_dry_run_capabilities`(正能力声明) | 否;无条目 ≠ 推断「不支持」之外的假能力 |
| **Wait** | `mode`(`poll`/`event`/`auto`), `poll_command`, `status_query`, `terminal`(状态→success/failure), `pending_values`, `event_key`/`match_field`/`resource_query`(event/auto), `default_timeout_secs` | **声明**(`ContractDecl.Wait` 正能力声明,且必须搭配 ResultInvoke dispatcher + 按模式的 hook:poll↔`WaitPoll`、event↔`WaitEvents`、auto↔两者,构造期配对校验,多余 hook 同样拒绝) | 声明后注册 `--wait`/`--wait-timeout`(框架 flag,不进 toolArgs);Schema 投影 `wait` 键;auto = 事件优先、流终止/订阅失败回退轮询,一个 deadline 覆盖两阶段并传入 `WaitPoll`/`WaitEvents`(及 `Command().Context()`);仅 pending 初始结果进入等待,success/failure/partial 原样返回 | 否;未声明命令传 `--wait` = unknown flag。终态失败经统一信封 `error.type: "wait"`(rc=8),超时保持 pending + `meta.operation.timed_out`(rc=0);轮询间/轮询中/事件消费中超时一律按 pending 关闭 |
| **Interface** | `interface_mode`, `interface_ref`, `availability`, `reason` | 评审源 | MCP meta + agent metadata 解析 | 否;与 CLI Identity 分离 |
| **Selection** | `agent_summary`, `use_when`, `avoid_when`, `examples`, `prerequisites`, `tips`, `workflow_refs`, … | 声明(`ContractDecl.Selection` / `ProductDecl`) | `ContractDecl` / `ProductDecl`(`schema_hints/` 已退役) | 可声明;声明载荷**不得携带** `Reviewed`(旧路径专用),携带即组装报错 |
| **FieldProvenance** | 各字段 winner / candidates | 组装派生物 | Schema 组装器 | 派生;须与 delivered value 一致 |
+1 -1
View File
@@ -403,7 +403,7 @@ SIGTERM、关 stdin,或先用 dws event stop <subscribe_id> --dry-run 预览
Selection: contract.SelectionSpec{
AgentSummary: "消费 OA、群生命周期或需要底层控制的个人事件流;Agent 通常使用 --flatten 输出 NDJSON",
UseWhen: []string{
"需要监听七个公开 OA 审批任务/实例 EventKey 中的一个或多个事件",
"需要监听六个公开 OA 审批任务/实例 EventKey 中的一个或多个事件",
"需要监听指定群的标题变更、成员进退群或群解散事件",
"用户显式给出原始 EventKey、Filter DSL、subscribe_id,要求原始 transport envelope,或需要普通 IM facade 不提供的高级多事件控制",
},
+2 -2
View File
@@ -152,8 +152,8 @@ func TestCrossPlatformCoveragePersonalSubscriptionProtectionCoversAllPublicEvent
}
}
if publicCount != 23 {
t.Fatalf("public personal events = %d, want 23 (16 IM + 7 OA)", publicCount)
if publicCount != 22 {
t.Fatalf("public personal events = %d, want 22 (16 IM + 6 OA)", publicCount)
}
for _, ruleType := range []string{"at", "all", "singleChat", "sender", "group"} {
if !ruleTypes[ruleType] {
-9
View File
@@ -61,13 +61,6 @@ func TestPersonalOAEventListAndSchemaCommands(t *testing.T) {
"process_code", "title", "status", "create_time", "event_time",
},
},
{
eventKey: personal.EventOAApprovalInstanceCC,
properties: []string{
"type", "event_id", "timestamp", "subscribe_id", "process_instance_id",
"process_code", "title", "status", "create_time", "event_time",
},
},
{
eventKey: personal.EventOAApprovalInstanceTerminated,
properties: []string{
@@ -168,7 +161,6 @@ func TestPersonalOAEventConsumeDryRunAndValidation(t *testing.T) {
personal.EventOAApprovalTaskFinished,
personal.EventOAApprovalTaskRedirected,
personal.EventOAApprovalInstanceStarted,
personal.EventOAApprovalInstanceCC,
personal.EventOAApprovalInstanceTerminated,
personal.EventOAApprovalInstanceFinished,
}
@@ -422,7 +414,6 @@ func TestPersonalOAMultiConsumeCreatesIndependentAllSubscriptionsOnSharedBus(t *
personal.EventOAApprovalTaskFinished,
personal.EventOAApprovalTaskRedirected,
personal.EventOAApprovalInstanceStarted,
personal.EventOAApprovalInstanceCC,
personal.EventOAApprovalInstanceTerminated,
personal.EventOAApprovalInstanceFinished,
}
+1 -67
View File
@@ -51,38 +51,7 @@ func (c *paramAliasCaptureCaller) CallTool(_ context.Context, server, tool strin
func (c *paramAliasCaptureCaller) paramAliasResponseForTool(tool string) string {
switch tool {
case "list_calendar_events":
return `{"success":true,"result":{"events":[],"hasMore":false,"nextCursor":""}}`
case "get_calendar_detail":
return c.paramAliasCalendarDetailResponse()
case "get_calendar_participants":
return `{"success":true,"result":{"participants":[{"userId":"fixture-user","displayName":"Fixture User"},{"userId":"user-2","displayName":"User Two"}]}}`
case "search_calendar":
return `{"success":true,"result":{"calendars":[]}}`
case "search_rooms":
return `{"success":true,"result":{"rooms":[]}}`
case "query_available_meeting_room":
return `{"success":true,"result":{"rooms":[],"hasMore":false}}`
case "list_meeting_room_groups":
return `{"success":true,"result":{"groups":[]}}`
case "query_busy_status":
return `{"success":true,"result":[]}`
case "list_suggested_event_times":
return `{"success":true,"result":{"recommendEventTimes":[]}}`
case "create_calendar_event":
return `{"success":true,"result":{"eventId":"event-1"}}`
case "update_calendar_event", "delete_calendar_event", "add_calendar_participant", "remove_calendar_participant":
return `{"success":true}`
case "respond":
status := "accepted"
if call := c.lastParamAliasCall(); call != nil {
if value, ok := call.args["responseStatus"].(string); ok && value != "" {
status = value
}
}
encoded, _ := json.Marshal(map[string]any{"success": true, "result": map[string]any{"responseStatus": status}})
return string(encoded)
case "get_current_user_profile":
return `{"success":true,"result":{"userId":"user-1","name":"Fixture Current User"}}`
return `{"result":{"events":[]}}`
case "query_records":
return `{"success":true,"status":"success","error":{},"data":{}}`
case "search_mail_users":
@@ -158,41 +127,6 @@ func (c *paramAliasCaptureCaller) paramAliasResponseForTool(tool string) string
}
}
func (c *paramAliasCaptureCaller) lastParamAliasCall() *paramAliasToolCall {
if len(c.calls) == 0 {
return nil
}
return &c.calls[len(c.calls)-1]
}
func (c *paramAliasCaptureCaller) paramAliasCalendarDetailResponse() string {
event := map[string]any{
"eventId": "event-1",
"summary": "Fixture Meeting",
"description": "fixture description",
"startDateTime": "2026-03-10T09:00:00+08:00",
"endDateTime": "2026-03-10T10:00:00+08:00",
}
for _, call := range c.calls {
switch call.tool {
case "create_calendar_event", "update_calendar_event":
for _, key := range []string{"eventId", "summary", "description", "startDateTime", "endDateTime", "timeZone", "location", "freeBusy"} {
if value, ok := call.args[key]; ok {
event[key] = value
}
}
case "respond":
if value, ok := call.args["responseStatus"]; ok {
event["responseStatus"] = value
}
case "delete_calendar_event":
event["status"] = "cancelled"
}
}
encoded, _ := json.Marshal(map[string]any{"success": true, "result": event})
return string(encoded)
}
func (*paramAliasCaptureCaller) Format() string { return "json" }
func (*paramAliasCaptureCaller) DryRun() bool { return false }
func (*paramAliasCaptureCaller) Fields() string { return "" }
@@ -5,8 +5,6 @@ package app
import (
"errors"
"os"
"os/exec"
"reflect"
"strings"
"testing"
@@ -19,8 +17,6 @@ import (
const (
appFixtureCurrentDOpenID = "DAAAAAAAAAAAiE"
appFixtureCurrentDOpenID2 = "DAQEBAQEBAQEiE"
paramAliasCalendarPayloadChildEnv = "DWS_TEST_CALENDAR_PARAM_ALIAS_PAYLOAD_CHILD"
)
// paramAliasCompleteCommands is deliberately keyed by the exact reviewed
@@ -50,37 +46,7 @@ var paramAliasCompleteCommands = map[string][]string{
"aitable workflow run": {"aitable", "workflow", "run", "--base-id", "base-1", "--workflow-id", "workflow-1", "--table-id", "table-1", "--record-ids", "record-1", "--yes"},
"attendance check result": {"attendance", "check", "result", "--users", "user-1,user-2", "--start", "2026-03-01", "--end", "2026-03-02"},
"attendance +check-result": {"attendance", "+check-result", "--users", "user-1,user-2", "--start", "2026-03-01", "--end", "2026-03-02"},
"calendar +agenda": {"calendar", "+agenda", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00", "--calendar-id", "primary", "--cursor", "cursor-1", "--limit", "7"},
"calendar +attendee-list": {"calendar", "+attendee-list", "--event", "event-1", "--calendar-id", "primary"},
"calendar +book": {"calendar", "+book", "--title", "Fixture Meeting", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T10:00:00+08:00", "--with", "Fixture User", "--yes"},
"calendar +book-search": {"calendar", "+book-search", "--query", "fixture"},
"calendar +cancel-event": {"calendar", "+cancel-event", "--event", "event-1", "--yes"},
"calendar +conflicts": {"calendar", "+conflicts", "--in-days", "1"},
"calendar +create": {"calendar", "+create", "--title", "Fixture Meeting", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T10:00:00+08:00", "--desc", "fixture description", "--attendees", "user-1,user-2", "--rooms", "room-1,room-2", "--calendar-id", "primary", "--yes"},
"calendar +free": {"calendar", "+free", "--who", "Fixture User", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00"},
"calendar +free-slots": {"calendar", "+free-slots", "--from", "9", "--to", "18", "--in-days", "1"},
"calendar +freebusy": {"calendar", "+freebusy", "--users", "user-1,user-2", "--rooms", "room-1,room-2", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00"},
"calendar +get": {"calendar", "+get", "--event", "event-1", "--calendar-id", "primary"},
"calendar +invite": {"calendar", "+invite", "--event", "event-1", "--with", "Fixture User", "--yes"},
"calendar +my-free": {"calendar", "+my-free", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00"},
"calendar +reschedule": {"calendar", "+reschedule", "--event", "event-1", "--start", "2026-03-10T10:00:00+08:00", "--end", "2026-03-10T11:00:00+08:00", "--yes"},
"calendar +room-find": {"calendar", "+room-find", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T10:00:00+08:00", "--room-name", "Fixture Room", "--group-id", "group-1", "--page", "1", "--limit", "7"},
"calendar +room-groups": {"calendar", "+room-groups", "--page", "1", "--limit", "7"},
"calendar +room-search": {"calendar", "+room-search", "--room-name", "Fixture Room"},
"calendar +rsvp": {"calendar", "+rsvp", "--event", "event-1", "--status", "accept", "--calendar-id", "primary", "--yes"},
"calendar +search-event": {"calendar", "+search-event", "--query", "fixture", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00", "--calendar-id", "primary", "--cursor", "cursor-1", "--limit", "7"},
"calendar +suggest-time": {"calendar", "+suggest-time", "--with", "Fixture User", "--duration", "30", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00"},
"calendar +suggestion": {"calendar", "+suggestion", "--users", "user-1,user-2", "--duration", "30", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00", "--timezone", "Asia/Shanghai"},
"calendar +update": {"calendar", "+update", "--event", "event-1", "--title", "Fixture Updated Meeting", "--desc", "fixture updated description", "--start", "2026-03-10T10:00:00+08:00", "--end", "2026-03-10T11:00:00+08:00", "--add-attendees", "user-2", "--remove-attendees", "user-1", "--yes"},
"calendar busy search": {"calendar", "busy", "search", "--users", "user-1,user-2", "--rooms", "room-1,room-2", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00"},
"calendar event create": {"calendar", "event", "create", "--title", "Fixture Meeting", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T10:00:00+08:00", "--remind-minutes", "15", "--timezone", "Asia/Shanghai", "--rooms", "room-1,room-2"},
"calendar event list": {"calendar", "event", "list", "--start", "2026-03-10T14:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00", "--calendar-id", "primary", "--cursor", "cursor-1", "--limit", "7"},
"calendar event respond": {"calendar", "event", "respond", "--id", "event-1", "--status", "accepted"},
"calendar event suggest": {"calendar", "event", "suggest", "--users", "user-1,user-2", "--duration", "30", "--start", "2026-03-10T09:00:00+08:00", "--end", "2026-03-10T18:00:00+08:00", "--timezone", "Asia/Shanghai"},
"calendar event update": {"calendar", "event", "update", "--id", "event-1", "--timezone", "Asia/Shanghai"},
"calendar room add": {"calendar", "room", "add", "--event", "event-1", "--rooms", "room-1,room-2"},
"calendar room delete": {"calendar", "room", "delete", "--event", "event-1", "--rooms", "room-1,room-2"},
"calendar room search": {"calendar", "room", "search", "--room-name", "Fixture Room", "--group-id", "group-1", "--start", "2027-03-10T09:00:00+08:00", "--end", "2027-03-10T10:00:00+08:00", "--page", "1", "--limit", "7"},
"chat +chat-messages": {"chat", "+chat-messages", "--group", "fixture-conversation"},
"chat +chat-add-bot": {"chat", "+chat-add-bot", "--id", "fixture-conversation", "--robot-code", "robot-1", "--yes"},
"chat +chat-audit-join": {"chat", "+chat-audit-join", "--group", "fixture-conversation", "--record-id", "7", "--applicant", "user-1", "--inviter", "user-2", "--status", "AuditApprove", "--yes"},
@@ -600,108 +566,6 @@ var paramAliasRepresentativePayloadCases = map[string]bool{
paramAliasPayloadCaseKey("report list", "from-date"): true, // date-range concept alias
}
// paramAliasCalendarPayloadCases keeps the full reviewed Calendar expansion
// separate from the long-lived app-c race process. Each case still executes
// both canonical and alias argv through the real PreParse/Cobra path and
// compares the final captured transport calls; the owning top-level test runs
// these allocations in a short-lived race-instrumented subprocess so all Root
// registrations are released together when that process exits.
var paramAliasCalendarPayloadCases = map[string]bool{
paramAliasPayloadCaseKey("calendar +agenda", "from"): true,
paramAliasPayloadCaseKey("calendar +agenda", "to"): true,
paramAliasPayloadCaseKey("calendar +agenda", "max-results"): true,
paramAliasPayloadCaseKey("calendar +agenda", "next-cursor"): true,
paramAliasPayloadCaseKey("calendar +agenda", "calendar-book-id"): true,
paramAliasPayloadCaseKey("calendar +attendee-list", "event-id"): true,
paramAliasPayloadCaseKey("calendar +attendee-list", "calendar-book-id"): true,
paramAliasPayloadCaseKey("calendar +book", "summary"): true,
paramAliasPayloadCaseKey("calendar +book", "attendee-names"): true,
paramAliasPayloadCaseKey("calendar +book-search", "keyword"): true,
paramAliasPayloadCaseKey("calendar +book-search", "search"): true,
paramAliasPayloadCaseKey("calendar +book-search", "name"): true,
paramAliasPayloadCaseKey("calendar +cancel-event", "event-id"): true,
paramAliasPayloadCaseKey("calendar +cancel-event", "id"): true,
paramAliasPayloadCaseKey("calendar +free", "name"): true,
paramAliasPayloadCaseKey("calendar +free-slots", "start-hour"): true,
paramAliasPayloadCaseKey("calendar +free-slots", "end-hour"): true,
paramAliasPayloadCaseKey("calendar +free-slots", "day-offset"): true,
paramAliasPayloadCaseKey("calendar +freebusy", "user-ids"): true,
paramAliasPayloadCaseKey("calendar +freebusy", "room-ids"): true,
paramAliasPayloadCaseKey("calendar +freebusy", "room-id"): true,
paramAliasPayloadCaseKey("calendar +my-free", "from"): true,
paramAliasPayloadCaseKey("calendar +my-free", "to"): true,
paramAliasPayloadCaseKey("calendar +invite", "id"): true,
paramAliasPayloadCaseKey("calendar +invite", "participant-names"): true,
paramAliasPayloadCaseKey("calendar +reschedule", "id"): true,
paramAliasPayloadCaseKey("calendar +reschedule", "from"): true,
paramAliasPayloadCaseKey("calendar +reschedule", "to"): true,
paramAliasPayloadCaseKey("calendar +room-groups", "page-size"): true,
paramAliasPayloadCaseKey("calendar +room-groups", "page-index"): true,
paramAliasPayloadCaseKey("calendar +room-search", "query"): true,
paramAliasPayloadCaseKey("calendar +suggest-time", "duration-minutes"): true,
paramAliasPayloadCaseKey("calendar +suggest-time", "attendee-names"): true,
paramAliasPayloadCaseKey("calendar +conflicts", "day-offset"): true,
paramAliasPayloadCaseKey("calendar busy search", "room-id"): true,
paramAliasPayloadCaseKey("calendar event create", "reminder-minutes"): true,
paramAliasPayloadCaseKey("calendar event create", "tz"): true,
paramAliasPayloadCaseKey("calendar event create", "room-id"): true,
paramAliasPayloadCaseKey("calendar event respond", "response-status"): true,
paramAliasPayloadCaseKey("calendar event suggest", "duration-minutes"): true,
paramAliasPayloadCaseKey("calendar event update", "tz"): true,
paramAliasPayloadCaseKey("calendar room add", "room-id"): true,
paramAliasPayloadCaseKey("calendar room delete", "room-id"): true,
paramAliasPayloadCaseKey("calendar room search", "room-group-id"): true,
paramAliasPayloadCaseKey("calendar +create", "summary"): true,
paramAliasPayloadCaseKey("calendar +create", "description"): true,
paramAliasPayloadCaseKey("calendar +create", "user-ids"): true,
paramAliasPayloadCaseKey("calendar +create", "room-ids"): true,
paramAliasPayloadCaseKey("calendar +create", "room-id"): true,
paramAliasPayloadCaseKey("calendar +create", "calendar-book-id"): true,
paramAliasPayloadCaseKey("calendar +create", "to"): true,
paramAliasPayloadCaseKey("calendar +create", "from"): true,
paramAliasPayloadCaseKey("calendar +get", "event-id"): true,
paramAliasPayloadCaseKey("calendar +get", "calendar-book-id"): true,
paramAliasPayloadCaseKey("calendar +room-find", "from"): true,
paramAliasPayloadCaseKey("calendar +room-find", "to"): true,
paramAliasPayloadCaseKey("calendar +room-find", "page-size"): true,
paramAliasPayloadCaseKey("calendar +room-find", "page-index"): true,
paramAliasPayloadCaseKey("calendar +room-find", "room-group-id"): true,
paramAliasPayloadCaseKey("calendar +room-find", "query"): true,
paramAliasPayloadCaseKey("calendar +rsvp", "event-id"): true,
paramAliasPayloadCaseKey("calendar +rsvp", "response-status"): true,
paramAliasPayloadCaseKey("calendar +search-event", "keyword"): true,
paramAliasPayloadCaseKey("calendar +search-event", "from"): true,
paramAliasPayloadCaseKey("calendar +search-event", "to"): true,
paramAliasPayloadCaseKey("calendar +search-event", "next-cursor"): true,
paramAliasPayloadCaseKey("calendar +search-event", "max-results"): true,
paramAliasPayloadCaseKey("calendar +suggestion", "user-ids"): true,
paramAliasPayloadCaseKey("calendar +suggestion", "duration-minutes"): true,
paramAliasPayloadCaseKey("calendar +suggestion", "from"): true,
paramAliasPayloadCaseKey("calendar +suggestion", "to"): true,
paramAliasPayloadCaseKey("calendar +suggestion", "tz"): true,
paramAliasPayloadCaseKey("calendar +update", "event-id"): true,
paramAliasPayloadCaseKey("calendar +update", "from"): true,
paramAliasPayloadCaseKey("calendar +update", "summary"): true,
paramAliasPayloadCaseKey("calendar +update", "description"): true,
paramAliasPayloadCaseKey("calendar +update", "add-user-ids"): true,
paramAliasPayloadCaseKey("calendar +update", "remove-user-ids"): true,
}
// paramAliasCalendarConfirmationCases selects one newly reviewed alias for
// every Calendar Shortcut whose runtime contract requires user confirmation.
// The complete Calendar matrix proves confirmed canonical/alias payload
// equality; these representatives additionally prove semantic normalization
// cannot cross the confirmation boundary before the first transport call.
var paramAliasCalendarConfirmationCases = map[string]bool{
paramAliasPayloadCaseKey("calendar +book", "summary"): true,
paramAliasPayloadCaseKey("calendar +cancel-event", "event-id"): true,
paramAliasPayloadCaseKey("calendar +create", "summary"): true,
paramAliasPayloadCaseKey("calendar +invite", "id"): true,
paramAliasPayloadCaseKey("calendar +reschedule", "from"): true,
paramAliasPayloadCaseKey("calendar +rsvp", "response-status"): true,
paramAliasPayloadCaseKey("calendar +update", "event-id"): true,
}
func TestCrossPlatformCoverageReviewedParamAliasesHaveCompleteTemplatesAndRepresentativeFinalPayloads(t *testing.T) {
concepts, err := cli.LoadParamConcepts()
if err != nil {
@@ -735,7 +599,28 @@ func TestCrossPlatformCoverageReviewedParamAliasesHaveCompleteTemplatesAndRepres
}
executedRepresentatives[caseKey] = true
t.Run(fixture.Command+"/"+fixture.Emitted, func(t *testing.T) {
assertParamAliasFinalPayloadEquivalent(t, fixture.Command, canonicalArgs, aliasArgs)
canonicalCaller := &paramAliasCaptureCaller{}
_, canonicalErr := executeParamAliasPayloadE2E(t, canonicalCaller, canonicalArgs...)
if canonicalErr != nil {
t.Fatalf("complete canonical command failed: %v\nargs=%v\ncalls=%#v", canonicalErr, canonicalArgs, canonicalCaller.calls)
}
if len(canonicalCaller.calls) == 0 {
t.Fatalf("complete canonical command reached no final transport payload: args=%v", canonicalArgs)
}
aliasCaller := &paramAliasCaptureCaller{}
ctx, aliasErr := executeParamAliasPayloadE2E(t, aliasCaller, aliasArgs...)
if aliasErr != nil {
t.Fatalf("complete alias command failed: %v\nargs=%v\ncalls=%#v", aliasErr, aliasArgs, aliasCaller.calls)
}
if ctx == nil {
t.Fatal("complete alias command skipped PreParse")
}
normalizeParamAliasVolatileDefaults(fixture.Command, canonicalCaller, aliasCaller)
if !reflect.DeepEqual(aliasCaller.calls, canonicalCaller.calls) {
t.Fatalf("final transport calls differ\ncanonical args: %v\nalias args: %v\ncanonical calls: %#v\nalias calls: %#v", canonicalArgs, aliasArgs, canonicalCaller.calls, aliasCaller.calls)
}
})
}
@@ -765,118 +650,6 @@ func TestCrossPlatformCoverageReviewedParamAliasesHaveCompleteTemplatesAndRepres
}
}
func TestCrossPlatformCoverageReviewedCalendarParamAliasesReachCanonicalEquivalentFinalPayloads(t *testing.T) {
if os.Getenv(paramAliasCalendarPayloadChildEnv) != "1" {
command := exec.Command(
os.Args[0],
"-test.run=^TestCrossPlatformCoverageReviewedCalendarParamAliasesReachCanonicalEquivalentFinalPayloads$",
"-test.count=1",
"-test.timeout=5m",
)
command.Env = append(os.Environ(), paramAliasCalendarPayloadChildEnv+"=1")
output, err := command.CombinedOutput()
if err != nil {
t.Fatalf("Calendar param-alias payload subprocess failed: %v\n%s", err, strings.TrimSpace(string(output)))
}
return
}
concepts, err := cli.LoadParamConcepts()
if err != nil {
t.Fatalf("LoadParamConcepts() error = %v", err)
}
executed := make(map[string]bool)
executedConfirmation := make(map[string]bool)
for _, fixture := range concepts.Fixture {
caseKey := paramAliasPayloadCaseKey(fixture.Command, fixture.Emitted)
if !paramAliasCalendarPayloadCases[caseKey] {
continue
}
executed[caseKey] = true
fixture := fixture
t.Run(fixture.Command+"/"+fixture.Emitted, func(t *testing.T) {
complete, ok := paramAliasCompleteCommand(fixture.Command, fixture.Expect)
if !ok {
t.Fatal("reviewed Calendar alias has no complete-command E2E template")
}
canonicalArgs := append([]string(nil), complete...)
aliasArgs, replacements := replaceLongFlag(canonicalArgs, fixture.Expect, fixture.Emitted)
if replacements != 1 {
t.Fatalf("complete Calendar command must contain canonical --%s exactly once; replacements=%d args=%v", fixture.Expect, replacements, canonicalArgs)
}
assertParamAliasFinalPayloadEquivalent(t, fixture.Command, canonicalArgs, aliasArgs)
if paramAliasCalendarConfirmationCases[caseKey] {
executedConfirmation[caseKey] = true
assertParamAliasCannotBypassConfirmation(t, aliasArgs)
}
})
}
for caseKey := range paramAliasCalendarPayloadCases {
if !executed[caseKey] {
t.Errorf("Calendar final-payload case %q has no active reviewed fixture", caseKey)
}
}
if len(executed) != len(paramAliasCalendarPayloadCases) {
t.Fatalf("Calendar final-payload coverage = %d, want %d", len(executed), len(paramAliasCalendarPayloadCases))
}
for caseKey := range paramAliasCalendarConfirmationCases {
if !executedConfirmation[caseKey] {
t.Errorf("Calendar confirmation case %q has no active reviewed fixture", caseKey)
}
}
if len(executedConfirmation) != len(paramAliasCalendarConfirmationCases) {
t.Fatalf("Calendar confirmation coverage = %d, want %d", len(executedConfirmation), len(paramAliasCalendarConfirmationCases))
}
}
func assertParamAliasCannotBypassConfirmation(t *testing.T, aliasArgs []string) {
t.Helper()
unconfirmedArgs, removals := removeExactArg(aliasArgs, "--yes")
if removals != 1 {
t.Fatalf("confirmation template must contain --yes exactly once; removals=%d args=%v", removals, aliasArgs)
}
caller := &paramAliasCaptureCaller{}
ctx, err := executeParamAliasPayloadE2E(t, caller, unconfirmedArgs...)
if ctx == nil {
t.Fatal("unconfirmed Calendar alias command skipped PreParse")
}
var appErr *apperrors.Error
if !errors.As(err, &appErr) || appErr.Reason != "confirmation_required" {
t.Fatalf("unconfirmed Calendar alias command error = %#v, want confirmation_required\nargs=%v", err, unconfirmedArgs)
}
if len(caller.calls) != 0 {
t.Fatalf("unconfirmed Calendar alias crossed the transport boundary: args=%v calls=%#v", unconfirmedArgs, caller.calls)
}
}
func assertParamAliasFinalPayloadEquivalent(t *testing.T, command string, canonicalArgs, aliasArgs []string) {
t.Helper()
canonicalCaller := &paramAliasCaptureCaller{}
_, canonicalErr := executeParamAliasPayloadE2E(t, canonicalCaller, canonicalArgs...)
if canonicalErr != nil {
t.Fatalf("complete canonical command failed: %v\nargs=%v\ncalls=%#v", canonicalErr, canonicalArgs, canonicalCaller.calls)
}
if len(canonicalCaller.calls) == 0 {
t.Fatalf("complete canonical command reached no final transport payload: args=%v", canonicalArgs)
}
aliasCaller := &paramAliasCaptureCaller{}
ctx, aliasErr := executeParamAliasPayloadE2E(t, aliasCaller, aliasArgs...)
if aliasErr != nil {
t.Fatalf("complete alias command failed: %v\nargs=%v\ncalls=%#v", aliasErr, aliasArgs, aliasCaller.calls)
}
if ctx == nil {
t.Fatal("complete alias command skipped PreParse")
}
normalizeParamAliasVolatileDefaults(command, canonicalCaller, aliasCaller)
if !reflect.DeepEqual(aliasCaller.calls, canonicalCaller.calls) {
t.Fatalf("final transport calls differ\ncanonical args: %v\nalias args: %v\ncanonical calls: %#v\nalias calls: %#v", canonicalArgs, aliasArgs, canonicalCaller.calls, aliasCaller.calls)
}
}
func TestCrossPlatformCoverageNewIMParamAliasesReachCanonicalEquivalentFinalPayloads(t *testing.T) {
activeAliases := 0
for _, test := range paramAliasNewIMCases {
-481
View File
@@ -1773,439 +1773,6 @@ var generatedParamAliases = []ParamAliasEntry{
},
Blocked: []string{"at-user-ids", "staff-id", "uid", "user", "user-id", "userid"},
},
{
CLIPath: "calendar +agenda",
Aliases: map[string]string{
"begin": "start",
"calendar": "calendar-id",
"calendar-book-id": "calendar-id",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"max-result": "limit",
"max-results": "limit",
"max-time": "end",
"min-time": "start",
"next-cursor": "cursor",
"next-page-token": "cursor",
"next-token": "cursor",
"page-size": "limit",
"page-token": "cursor",
"per-page": "limit",
"since": "start",
"size": "limit",
"start-date": "start",
"start-time": "start",
"take": "limit",
"time-max": "end",
"time-min": "start",
"to": "end",
"top": "limit",
},
Blocked: []string{"acl-id", "count", "date", "event", "event-id", "id", "offset", "page", "room-id", "time"},
},
{
CLIPath: "calendar +attendee-list",
Aliases: map[string]string{
"calendar": "calendar-id",
"calendar-book-id": "calendar-id",
"calendar-event-id": "event",
"event-id": "event",
},
Blocked: []string{"acl-id", "room-id"},
Ambiguous: []string{"id"},
},
{
CLIPath: "calendar +book",
Aliases: map[string]string{
"attendee-names": "with",
"begin": "start",
"end-time": "end",
"from": "start",
"names": "with",
"participant-names": "with",
"start-time": "start",
"subject": "title",
"summary": "title",
"to": "end",
},
Blocked: []string{"attendees", "calendar-id", "calendar-name", "date", "end-date", "name", "open-dingtalk-ids", "participants", "room-id", "room-ids", "room-name", "rooms", "start-date", "time", "time-max", "time-min", "user", "user-id", "user-ids", "users"},
},
{
CLIPath: "calendar +book-search",
Aliases: map[string]string{
"keyword": "query",
"keywords": "query",
"name": "query",
"q": "query",
"search": "query",
"search-word": "query",
},
Blocked: []string{"subject", "text", "title"},
},
{
CLIPath: "calendar +cancel-event",
Aliases: map[string]string{
"calendar-event-id": "event",
"event-id": "event",
"id": "event",
},
Blocked: []string{"acl-id", "calendar-book-id", "calendar-id", "room-id"},
},
{
CLIPath: "calendar +conflicts",
Aliases: map[string]string{
"day-offset": "in-days",
"days-from-today": "in-days",
},
Blocked: []string{"days", "duration", "end", "from", "start", "to"},
},
{
CLIPath: "calendar +create",
Aliases: map[string]string{
"attendee-ids": "attendees",
"begin": "start",
"calendar": "calendar-id",
"calendar-book-id": "calendar-id",
"description": "desc",
"end-time": "end",
"freebusy": "free-busy",
"from": "start",
"room-id": "rooms",
"room-ids": "rooms",
"start-time": "start",
"subject": "title",
"summary": "title",
"time-zone": "timezone",
"to": "end",
"tz": "timezone",
"user-ids": "attendees",
"users": "attendees",
},
Blocked: []string{"acl-id", "attendee-name", "attendee-names", "calendar-name", "config", "date", "end-date", "event", "event-id", "field-description", "group-id", "id", "locale", "name", "offset", "open-dingtalk-ids", "participant-name", "participant-names", "rich-text-desc", "room", "room-name", "start-date", "time", "time-max", "time-min", "utc-offset", "who", "with"},
},
{
CLIPath: "calendar +free",
Aliases: map[string]string{
"begin": "start",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"max-time": "end",
"min-time": "start",
"name": "who",
"person": "who",
"person-name": "who",
"since": "start",
"start-date": "start",
"start-time": "start",
"time-max": "end",
"time-min": "start",
"to": "end",
},
Blocked: []string{"date", "time", "user", "user-id", "user-ids", "users", "with"},
},
{
CLIPath: "calendar +free-slots",
Aliases: map[string]string{
"day-offset": "in-days",
"days-from-today": "in-days",
"end-hour": "to",
"start-hour": "from",
},
Blocked: []string{"days", "duration", "end", "end-time", "start", "start-time", "time-max", "time-min"},
},
{
CLIPath: "calendar +freebusy",
Aliases: map[string]string{
"begin": "start",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"max-time": "end",
"min-time": "start",
"room-id": "rooms",
"room-ids": "rooms",
"since": "start",
"start-date": "start",
"start-time": "start",
"time-max": "end",
"time-min": "start",
"to": "end",
"user-ids": "users",
},
Blocked: []string{"at-user-ids", "attendee-names", "date", "group-id", "location", "name", "names", "participant-names", "room", "room-name", "staff-id", "time", "uid", "user", "user-id", "userid", "who", "with"},
},
{
CLIPath: "calendar +get",
Aliases: map[string]string{
"calendar": "calendar-id",
"calendar-book-id": "calendar-id",
"calendar-event-id": "event",
"event-id": "event",
},
Blocked: []string{"acl-id", "room-id"},
Ambiguous: []string{"id"},
},
{
CLIPath: "calendar +invite",
Aliases: map[string]string{
"attendee-names": "with",
"calendar-event-id": "event",
"event-id": "event",
"id": "event",
"names": "with",
"participant-names": "with",
},
Blocked: []string{"acl-id", "attendees", "calendar-book-id", "calendar-id", "open-dingtalk-ids", "participants", "room-id", "user", "user-id", "user-ids", "users"},
},
{
CLIPath: "calendar +my-free",
Aliases: map[string]string{
"begin": "start",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"max-time": "end",
"min-time": "start",
"since": "start",
"start-date": "start",
"start-time": "start",
"time-max": "end",
"time-min": "start",
"to": "end",
},
Blocked: []string{"date", "time"},
},
{
CLIPath: "calendar +reschedule",
Aliases: map[string]string{
"begin": "start",
"calendar-event-id": "event",
"end-time": "end",
"event-id": "event",
"from": "start",
"id": "event",
"start-time": "start",
"to": "end",
},
Blocked: []string{"acl-id", "calendar-book-id", "calendar-id", "date", "end-date", "room-id", "start-date", "time", "time-max", "time-min"},
},
{
CLIPath: "calendar +room-find",
Aliases: map[string]string{
"begin": "start",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"group": "group-id",
"max-result": "limit",
"max-results": "limit",
"max-time": "end",
"min-time": "start",
"name": "room-name",
"page-index": "page",
"page-size": "limit",
"per-page": "limit",
"query": "room-name",
"room-group-id": "group-id",
"since": "start",
"size": "limit",
"start-date": "start",
"start-time": "start",
"take": "limit",
"time-max": "end",
"time-min": "start",
"to": "end",
"top": "limit",
},
Blocked: []string{"count", "cursor", "date", "location", "next-cursor", "page-token", "room", "room-id", "room-ids", "rooms", "time"},
},
{
CLIPath: "calendar +room-groups",
Aliases: map[string]string{
"max-result": "limit",
"max-results": "limit",
"page-index": "page",
"page-size": "limit",
"per-page": "limit",
"size": "limit",
"take": "limit",
"top": "limit",
},
Blocked: []string{"count", "cursor", "page-token"},
},
{
CLIPath: "calendar +room-search",
Aliases: map[string]string{
"name": "room-name",
"query": "room-name",
},
Blocked: []string{"group-id", "location", "room", "room-id", "room-ids", "rooms"},
},
{
CLIPath: "calendar +rsvp",
Aliases: map[string]string{
"calendar": "calendar-id",
"calendar-book-id": "calendar-id",
"calendar-event-id": "event",
"event-id": "event",
"response": "status",
"response-status": "status",
},
Blocked: []string{"acl-id", "availability", "done", "free-busy", "room-id", "state"},
Ambiguous: []string{"id"},
},
{
CLIPath: "calendar +search-event",
Aliases: map[string]string{
"begin": "start",
"calendar": "calendar-id",
"calendar-book-id": "calendar-id",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"keyword": "query",
"keywords": "query",
"max-result": "limit",
"max-results": "limit",
"max-time": "end",
"min-time": "start",
"next-cursor": "cursor",
"next-page-token": "cursor",
"next-token": "cursor",
"page-size": "limit",
"page-token": "cursor",
"per-page": "limit",
"q": "query",
"search": "query",
"search-word": "query",
"since": "start",
"size": "limit",
"start-date": "start",
"start-time": "start",
"take": "limit",
"time-max": "end",
"time-min": "start",
"to": "end",
"top": "limit",
},
Blocked: []string{"acl-id", "count", "date", "event", "event-id", "id", "name", "offset", "page", "page-index", "room-id", "subject", "text", "time", "title"},
},
{
CLIPath: "calendar +suggest-time",
Aliases: map[string]string{
"attendee-names": "with",
"begin": "start",
"duration-minutes": "duration",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"max-time": "end",
"meeting-duration-minutes": "duration",
"min-time": "start",
"names": "with",
"participant-names": "with",
"since": "start",
"start-date": "start",
"start-time": "start",
"time-max": "end",
"time-min": "start",
"to": "end",
},
Blocked: []string{"attendees", "date", "open-dingtalk-ids", "participants", "remind-minutes", "time", "user", "user-id", "user-ids", "users"},
},
{
CLIPath: "calendar +suggestion",
Aliases: map[string]string{
"begin": "start",
"duration-minutes": "duration",
"end-date": "end",
"end-time": "end",
"from": "start",
"from-date": "start",
"max-time": "end",
"meeting-duration-minutes": "duration",
"min-time": "start",
"since": "start",
"start-date": "start",
"start-time": "start",
"time-max": "end",
"time-min": "start",
"time-zone": "timezone",
"to": "end",
"tz": "timezone",
"user-ids": "users",
},
Blocked: []string{"at-user-ids", "attendee-name", "attendee-names", "date", "locale", "name", "names", "offset", "open-dingtalk-ids", "participant-name", "participant-names", "remind-minutes", "room-id", "room-ids", "room-name", "rooms", "staff-id", "time", "uid", "user", "user-id", "userid", "utc-offset", "who", "with"},
},
{
CLIPath: "calendar +update",
Aliases: map[string]string{
"add-user-ids": "add-attendees",
"add-users": "add-attendees",
"begin": "start",
"calendar": "calendar-id",
"calendar-book-id": "calendar-id",
"calendar-event-id": "event",
"description": "desc",
"end-time": "end",
"event-id": "event",
"freebusy": "free-busy",
"from": "start",
"remove-user-ids": "remove-attendees",
"remove-users": "remove-attendees",
"start-time": "start",
"subject": "title",
"summary": "title",
"time-zone": "timezone",
"to": "end",
"tz": "timezone",
},
Blocked: []string{"acl-id", "attendee-name", "attendee-names", "calendar-name", "config", "date", "end-date", "field-description", "group-id", "locale", "name", "offset", "open-dingtalk-ids", "participant-name", "participant-names", "remind-minutes", "reminder-minutes", "rich-text-desc", "room-id", "room-ids", "room-name", "rooms", "start-date", "time", "time-max", "time-min", "utc-offset", "who", "with"},
Ambiguous: []string{"attendees", "id", "user-ids", "users"},
},
{
CLIPath: "calendar acl delete",
Blocked: []string{"calendar-book-id", "calendar-id", "event", "event-id", "room-id", "user-id"},
},
{
CLIPath: "calendar attachment add",
Blocked: []string{"attachments", "file", "file-id", "file-ids"},
},
{
CLIPath: "calendar attendee add",
Blocked: []string{"attendee-name", "attendee-names", "open-dingtalk-ids", "participant-name", "participant-names", "who", "with"},
},
{
CLIPath: "calendar attendee delete",
Blocked: []string{"attendee-name", "attendee-names", "open-dingtalk-ids", "participant-name", "participant-names", "who", "with"},
},
{
CLIPath: "calendar busy search",
Aliases: map[string]string{
"room-id": "rooms",
},
Blocked: []string{"attendee-name", "attendee-names", "group-id", "location", "name", "names", "participant-names", "room", "room-name", "who", "with"},
},
{
CLIPath: "calendar event create",
Aliases: map[string]string{
"reminder-minutes": "remind-minutes",
"reminder-offset-minutes": "remind-minutes",
"room-id": "rooms",
"time-zone": "timezone",
"tz": "timezone",
},
Blocked: []string{"at", "attendee-name", "attendee-names", "due", "duration", "group-id", "locale", "offset", "participant-name", "participant-names", "remind-at", "reminder-time", "room", "room-name", "utc-offset", "who", "with"},
},
{
CLIPath: "calendar event list",
Aliases: map[string]string{
@@ -2222,54 +1789,6 @@ var generatedParamAliases = []ParamAliasEntry{
},
Blocked: []string{"offset", "page", "time"},
},
{
CLIPath: "calendar event respond",
Aliases: map[string]string{
"response": "status",
"response-status": "status",
},
Blocked: []string{"availability", "done", "free-busy", "state"},
},
{
CLIPath: "calendar event suggest",
Aliases: map[string]string{
"duration-minutes": "duration",
"meeting-duration-minutes": "duration",
"time-zone": "timezone",
"tz": "timezone",
},
Blocked: []string{"attendee-name", "attendee-names", "from", "locale", "name", "names", "offset", "open-dingtalk-ids", "participant-names", "remind-minutes", "to", "utc-offset", "who", "with"},
},
{
CLIPath: "calendar event update",
Aliases: map[string]string{
"time-zone": "timezone",
"tz": "timezone",
},
Blocked: []string{"attendees", "group-id", "locale", "offset", "participants", "remind-minutes", "reminder-minutes", "room-id", "room-ids", "room-name", "rooms", "utc-offset"},
},
{
CLIPath: "calendar room add",
Aliases: map[string]string{
"room-id": "rooms",
},
Blocked: []string{"group-id", "location", "room", "room-name"},
},
{
CLIPath: "calendar room delete",
Aliases: map[string]string{
"room-id": "rooms",
},
Blocked: []string{"group-id", "location", "room", "room-name"},
},
{
CLIPath: "calendar room search",
Aliases: map[string]string{
"group": "group-id",
"room-group-id": "group-id",
},
Blocked: []string{"location", "room", "room-id", "room-ids", "rooms"},
},
{
CLIPath: "chat +bot-find",
Aliases: map[string]string{
File diff suppressed because one or more lines are too long
+1 -1
View File
@@ -1058,7 +1058,7 @@ var schemaCompactPayloadKeys = map[string]bool{
"agent_summary": true, "description": true,
"effect": true, "risk": true, "confirmation": true, "idempotency": true,
"interface_mode": true, "availability": true, "interface_reason": true,
"parameters": true, "constraints": true, "positionals": true, "dry_run": true, "wait": true,
"parameters": true, "constraints": true, "positionals": true, "dry_run": true,
"result": true, "pagination": true,
"examples": true, "use_when": true, "avoid_when": true,
}
-1
View File
@@ -81,7 +81,6 @@ var schemaCatalogToolOptionalKeys = []string{
"pagination",
"positionals",
"result",
"wait",
}
var schemaCatalogToolEnums = map[string][]string{
-21
View File
@@ -60,7 +60,6 @@ type ToolSpec struct {
Constraints RuntimeSchemaConstraints
Positionals []contract.RuntimeSchemaPositional
DryRun *contract.DryRunSpec
Wait *contract.WaitSpec
Result *contract.ResultSpec
Pagination *contract.PaginationSpec
Safety contract.SafetySpec
@@ -136,7 +135,6 @@ type RuntimeToolSpecInput struct {
Constraints RuntimeSchemaConstraints
Positionals []contract.RuntimeSchemaPositional
DryRun *contract.DryRunSpec
Wait *contract.WaitSpec
Result *contract.ResultSpec
Pagination *contract.PaginationSpec
Safety contract.SafetySpec
@@ -544,11 +542,6 @@ func (t ToolSpec) Validate() error {
return err
}
}
if t.Wait != nil {
if err := t.Wait.Validate(id.CanonicalPath); err != nil {
return err
}
}
if t.Result != nil {
if _, err := contract.NormalizeResultSpec(t.Result, id.CanonicalPath); err != nil {
return err
@@ -751,16 +744,6 @@ func (t ToolSpec) normalized() ToolSpec {
dryRun.PreviewKind = strings.TrimSpace(dryRun.PreviewKind)
out.DryRun = &dryRun
}
if t.Wait != nil {
// NormalizeWaitSpec is the single canonical form shared with the
// declaration path: trimmed status values, duplicate/conflict
// rejection, defensive copy. Invalid declarations are rejected by
// ToolSpec.Validate below, which runs the same normalization
// through WaitSpec.Validate.
if wait, err := contract.NormalizeWaitSpec(t.Wait, id.CanonicalPath); err == nil {
out.Wait = wait
}
}
if t.Result != nil {
result, err := contract.NormalizeResultSpec(t.Result, id.CanonicalPath)
if err == nil {
@@ -999,10 +982,6 @@ func (t ToolSpec) ToPayload() (map[string]any, error) {
value, _ := typedJSONValue(t.DryRun)
payload["dry_run"] = value
}
if t.Wait != nil {
value, _ := typedJSONValue(t.Wait)
payload["wait"] = value
}
if t.Result != nil {
value, _ := typedJSONValue(t.Result)
payload["result"] = value
@@ -974,57 +974,3 @@ func TestFinalProvenanceCoverageDoesNotInventOptionalInterfaceReason(t *testing.
t.Fatalf("optional local interface_reason should not require invented provenance: %v", err)
}
}
func TestToolSpecWaitCapabilityIsPositiveOnly(t *testing.T) {
base := RuntimeToolSpecInput{Identity: contract.ToolIdentitySpec{
ProductID: "sample",
Name: "waitrun",
CLIName: "waitrun",
CLIPath: "sample waitrun",
}}
withoutCapability, err := ToolSpecFromRuntime(base)
if err != nil {
t.Fatalf("ToolSpecFromRuntime() error = %v", err)
}
payload, err := withoutCapability.ToPayload()
if err != nil {
t.Fatalf("ToPayload() error = %v", err)
}
if _, ok := payload["wait"]; ok {
t.Fatalf("nil capability unexpectedly emitted wait: %#v", payload["wait"])
}
base.Wait = &contract.WaitSpec{Mode: "webhook"}
if _, err := ToolSpecFromRuntime(base); err == nil || !strings.Contains(err.Error(), "unknown mode") {
t.Fatalf("invalid mode error = %v", err)
}
base.Wait = &contract.WaitSpec{Mode: contract.WaitModeEvent, StatusQuery: "status", Terminal: map[string]contract.ResultOutcome{"DONE": contract.ResultOutcomeSuccess}}
if _, err := ToolSpecFromRuntime(base); err == nil || !strings.Contains(err.Error(), "requires event_key") {
t.Fatalf("event mode body error = %v", err)
}
base.Wait = &contract.WaitSpec{
Mode: contract.WaitModePoll,
PollCommand: "sample status get",
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{"COMPLETED": contract.ResultOutcomeSuccess},
PendingValues: []string{"NEW"},
DefaultTimeoutSecs: 120,
}
withCapability, err := ToolSpecFromRuntime(base)
if err != nil {
t.Fatalf("ToolSpecFromRuntime(valid wait) error = %v", err)
}
if withCapability.Wait == nil || withCapability.Wait.Mode != contract.WaitModePoll {
t.Fatalf("wait capability lost through normalization: %#v", withCapability.Wait)
}
payload, err = withCapability.ToPayload()
if err != nil {
t.Fatalf("ToPayload(valid wait) error = %v", err)
}
wait := payload["wait"].(map[string]any)
if wait["mode"] != contract.WaitModePoll || wait["poll_command"] != "sample status get" {
t.Fatalf("wait payload=%#v", wait)
}
}
-1
View File
@@ -354,7 +354,6 @@ func runtimeToolSpecFromContractFinal(entry runtimeSchemaEntry, final contract.C
Constraints: constraints,
Positionals: positionals,
DryRun: final.DryRun,
Wait: final.Wait,
Result: result,
Pagination: pagination,
Safety: safety,
-2
View File
@@ -65,7 +65,6 @@ type schemaToolWire struct {
Constraints RuntimeSchemaConstraints `json:"constraints"`
Positionals []contract.RuntimeSchemaPositional `json:"positionals"`
DryRun *contract.DryRunSpec `json:"dry_run"`
Wait *contract.WaitSpec `json:"wait"`
Result *contract.ResultSpec `json:"result"`
Pagination *contract.PaginationSpec `json:"pagination"`
Effect string `json:"effect"`
@@ -270,7 +269,6 @@ func schemaToolSpecFromWire(wire schemaToolWire) (ToolSpec, error) {
Constraints: wire.Constraints,
Positionals: wire.Positionals,
DryRun: wire.DryRun,
Wait: wire.Wait,
Result: wire.Result,
Pagination: wire.Pagination,
Safety: contract.SafetySpec{
-1
View File
@@ -30,7 +30,6 @@ type ContractFinalPayload struct {
Parameters []ParamDecl
Safety *SafetySpec
DryRun *DryRunSpec
Wait *WaitSpec
Result *ResultSpec
Pagination *PaginationSpec
Interface *InterfaceSpec
-149
View File
@@ -62,155 +62,6 @@ type DryRunSpec struct {
RemoteReads bool `json:"remote_reads,omitempty"`
}
// Wait modes. Poll executes the leaf's WaitPoll hook on a cadence. Event
// consumes the leaf's WaitEvents push stream and correlates events to the
// accepted resource. Auto prefers the event stream and falls back to polling
// when the stream ends before a terminal status.
const (
WaitModePoll = "poll"
WaitModeEvent = "event"
WaitModeAuto = "auto"
)
// WaitSpec is a positive capability declaration for terminal-state waiting
// (approval flows, async exports, batch jobs). A nil ToolSpec.Wait means the
// command has not declared reviewed --wait support; the flag is not
// registered and the Schema does not publish the capability.
//
// Like DryRunSpec, the object is one atomic contract field: Schema only
// projects the reviewed capability; runtime execution stays owned by the
// command runner through the leaf's WaitPoll / WaitEvents hooks. PollCommand
// names the read command that observes status — it is a declared,
// catalog-visible fact (the same command an agent would poll manually), not
// a framework-owned invocation: how one poll or event subscription executes
// is decided by the leaf.
type WaitSpec struct {
Mode string `json:"mode"`
PollCommand string `json:"poll_command,omitempty"`
StatusQuery string `json:"status_query"`
Terminal map[string]ResultOutcome `json:"terminal"`
PendingValues []string `json:"pending_values,omitempty"`
// EventKey is the push channel key the WaitEvents hook subscribes to
// (event/auto modes). Declared for the catalog; the transport stays
// leaf-owned.
EventKey string `json:"event_key,omitempty"`
// MatchField is the event-document path holding the resource identifier
// (event/auto modes); its value must equal the ResourceQuery resolution
// of the accepted result.
MatchField string `json:"match_field,omitempty"`
// ResourceQuery is the dotted path into the accepted result data that
// yields the resource identifier correlated against MatchField
// (event/auto modes).
ResourceQuery string `json:"resource_query,omitempty"`
// DefaultTimeoutSecs is the reviewed default for --wait-timeout. Zero
// means the framework default (300s); the user flag always wins.
DefaultTimeoutSecs int `json:"default_timeout_secs"`
}
// Validate checks mode requirements and the terminal/pending status maps.
// Unknown terminal outcomes, unknown modes, and mode/body mismatches fail at
// declaration so a malformed wait capability cannot reach the wire.
// Validation delegates to NormalizeWaitSpec so the acceptance rules can never
// drift from the normalization the wire and the runtime wait engine share.
func (w WaitSpec) Validate(canonical string) error {
_, err := NormalizeWaitSpec(&w, canonical)
return err
}
// NormalizeWaitSpec returns a validated, canonical, defensively copied wait
// contract. It is shared by declaration (corecmd.New / AttachContract),
// ToolSpec, and snapshot paths, mirroring NormalizeResultSpec. Status values
// are trimmed into their wire form: the wait engine compares backend
// statuses verbatim against these tables, so a padded declaration
// (" processing ") would publish a Schema that its own runtime treats as an
// unknown status. Values collapsing onto one value after trimming (duplicate
// pending values, duplicate terminal keys, terminal/pending conflicts) are
// rejected instead of silently merged.
func NormalizeWaitSpec(in *WaitSpec, canonical string) (*WaitSpec, error) {
if in == nil {
return nil, nil
}
canonical = defaultString(strings.TrimSpace(canonical), "<unknown>")
out := &WaitSpec{
Mode: strings.TrimSpace(in.Mode),
PollCommand: strings.TrimSpace(in.PollCommand),
StatusQuery: strings.TrimSpace(in.StatusQuery),
EventKey: strings.TrimSpace(in.EventKey),
MatchField: strings.TrimSpace(in.MatchField),
ResourceQuery: strings.TrimSpace(in.ResourceQuery),
DefaultTimeoutSecs: in.DefaultTimeoutSecs,
}
if out.Mode == "" {
return nil, fmt.Errorf("schema tool %s wait has no mode", canonical)
}
switch out.Mode {
case WaitModePoll, WaitModeEvent, WaitModeAuto:
default:
return nil, fmt.Errorf("schema tool %s wait has unknown mode %q", canonical, out.Mode)
}
needsPoll := out.Mode == WaitModePoll || out.Mode == WaitModeAuto
if needsPoll && out.PollCommand == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires poll_command", canonical, out.Mode)
}
needsEvent := out.Mode == WaitModeEvent || out.Mode == WaitModeAuto
if needsEvent {
if out.EventKey == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires event_key", canonical, out.Mode)
}
if out.MatchField == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires match_field", canonical, out.Mode)
}
if out.ResourceQuery == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires resource_query", canonical, out.Mode)
}
}
if out.StatusQuery == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires status_query", canonical, out.Mode)
}
if len(in.Terminal) == 0 {
return nil, fmt.Errorf("schema tool %s wait has no terminal states", canonical)
}
out.Terminal = make(map[string]ResultOutcome, len(in.Terminal))
for status, outcome := range in.Terminal {
status = strings.TrimSpace(status)
if status == "" {
return nil, fmt.Errorf("schema tool %s wait has a blank terminal status", canonical)
}
if _, dup := out.Terminal[status]; dup {
return nil, fmt.Errorf("schema tool %s wait has duplicate terminal status %q", canonical, status)
}
// Terminal states must close into success or failure. Pending and
// partial are not wait outcomes: pending is expressed through
// timeout, and partial requires the typed multi-status payload only
// the leaf can construct.
if outcome != ResultOutcomeSuccess && outcome != ResultOutcomeFailure {
return nil, fmt.Errorf(
"schema tool %s wait terminal status %q must map to success or failure, got %q",
canonical, status, outcome)
}
out.Terminal[status] = outcome
}
seenPending := make(map[string]bool, len(in.PendingValues))
for _, value := range in.PendingValues {
value = strings.TrimSpace(value)
if value == "" {
return nil, fmt.Errorf("schema tool %s wait has a blank pending value", canonical)
}
if _, conflict := out.Terminal[value]; conflict {
return nil, fmt.Errorf("schema tool %s wait status %q is both terminal and pending", canonical, value)
}
if seenPending[value] {
return nil, fmt.Errorf("schema tool %s wait has duplicate pending value %q", canonical, value)
}
seenPending[value] = true
out.PendingValues = append(out.PendingValues, value)
}
if in.DefaultTimeoutSecs < 0 {
return nil, fmt.Errorf("schema tool %s wait default_timeout_secs must be >= 0", canonical)
}
return out, nil
}
// ResultOutcome is one closed unified-output envelope outcome.
type ResultOutcome string
-227
View File
@@ -1,227 +0,0 @@
// 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 contract
import "testing"
func TestWaitSpecValidateAcceptsReviewedShapes(t *testing.T) {
cases := []WaitSpec{
{
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "result.status",
Terminal: map[string]ResultOutcome{
"COMPLETED": ResultOutcomeSuccess,
"REJECTED": ResultOutcomeFailure,
},
PendingValues: []string{"NEW", "RUNNING"},
DefaultTimeoutSecs: 600,
},
}
for i, spec := range cases {
if err := spec.Validate("sample.tool"); err != nil {
t.Fatalf("case %d: unexpected error: %v", i, err)
}
}
}
func TestWaitSpecValidateRejectsMalformedShapes(t *testing.T) {
terminal := map[string]ResultOutcome{"COMPLETED": ResultOutcomeSuccess}
cases := map[string]WaitSpec{
"no mode": {
Terminal: terminal,
},
"unknown mode": {
Mode: "webhook",
Terminal: terminal,
},
"poll without poll_command": {
Mode: WaitModePoll,
StatusQuery: "status",
Terminal: terminal,
},
"poll without status_query": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
Terminal: terminal,
},
"event without event_key": {
Mode: WaitModeEvent,
MatchField: "process_instance_id",
ResourceQuery: "id",
StatusQuery: "result.status",
Terminal: terminal,
},
"event without match_field": {
Mode: WaitModeEvent,
EventKey: "bpms_instance_change",
ResourceQuery: "id",
StatusQuery: "result.status",
Terminal: terminal,
},
"event without resource_query": {
Mode: WaitModeEvent,
EventKey: "bpms_instance_change",
MatchField: "process_instance_id",
StatusQuery: "result.status",
Terminal: terminal,
},
"auto missing poll_command": {
Mode: WaitModeAuto,
EventKey: "export_finished",
MatchField: "job_id",
ResourceQuery: "job_id",
StatusQuery: "status",
Terminal: terminal,
},
"no terminal states": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
},
"blank terminal status": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{" ": ResultOutcomeSuccess},
},
"terminal outcome pending": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcomePending, "REJECTED": ResultOutcomeFailure},
},
"terminal outcome partial": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcomePartialFailure, "REJECTED": ResultOutcomeFailure},
},
"terminal outcome outside closed set": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcome("explosion")},
},
"only pending terminal outcome": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"NEW": ResultOutcomePending},
},
"status both terminal and pending": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{"COMPLETED"},
},
"blank pending value": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{" "},
},
"negative timeout default": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
DefaultTimeoutSecs: -1,
},
}
for name, spec := range cases {
if err := spec.Validate("sample.tool"); err == nil {
t.Fatalf("%s: expected error, got nil", name)
}
}
}
func TestNormalizeWaitSpecTrimsStatusValuesIntoWireForm(t *testing.T) {
in := &WaitSpec{
Mode: " poll ",
PollCommand: " oa approval-instance get ",
StatusQuery: " result.status ",
Terminal: map[string]ResultOutcome{" COMPLETED ": ResultOutcomeSuccess, "REJECTED": ResultOutcomeFailure},
PendingValues: []string{" NEW ", "RUNNING"},
DefaultTimeoutSecs: 60,
}
out, err := NormalizeWaitSpec(in, "sample.tool")
if err != nil {
t.Fatalf("NormalizeWaitSpec() error = %v", err)
}
if out.Mode != WaitModePoll || out.PollCommand != "oa approval-instance get" || out.StatusQuery != "result.status" {
t.Fatalf("normalized scalars: %#v", out)
}
if len(out.Terminal) != 2 {
t.Fatalf("terminal=%#v, want two trimmed keys", out.Terminal)
}
if got := out.Terminal["COMPLETED"]; got != ResultOutcomeSuccess {
t.Fatalf("terminal[COMPLETED]=%q, want success (key must be trimmed)", got)
}
if _, padded := out.Terminal[" COMPLETED "]; padded {
t.Fatal("padded terminal key survived normalization")
}
for i, want := range []string{"NEW", "RUNNING"} {
if out.PendingValues[i] != want {
t.Fatalf("pending[%d]=%q, want %q", i, out.PendingValues[i], want)
}
}
// The input declaration must stay untouched (defensive copy).
if _, padded := in.Terminal[" COMPLETED "]; !padded {
t.Fatal("NormalizeWaitSpec mutated its input terminal map")
}
if in.PendingValues[0] != " NEW " {
t.Fatal("NormalizeWaitSpec mutated its input pending values")
}
}
func TestNormalizeWaitSpecRejectsDuplicatesAndConflictsAfterTrim(t *testing.T) {
terminal := map[string]ResultOutcome{"COMPLETED": ResultOutcomeSuccess}
cases := map[string]*WaitSpec{
"terminal keys collapsing after trim": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcomeSuccess, " COMPLETED ": ResultOutcomeFailure},
},
"pending values collapsing after trim": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{"NEW", " NEW "},
},
"terminal/pending conflict hidden by padding": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{" COMPLETED "},
},
}
for name, spec := range cases {
if _, err := NormalizeWaitSpec(spec, "sample.tool"); err == nil {
t.Fatalf("%s: expected error, got nil", name)
}
}
}
func TestNormalizeWaitSpecNilReturnsNil(t *testing.T) {
out, err := NormalizeWaitSpec(nil, "sample.tool")
if err != nil || out != nil {
t.Fatalf("NormalizeWaitSpec(nil) = %#v, %v", out, err)
}
}
-4
View File
@@ -42,7 +42,6 @@ type ContractDecl struct {
Positionals []contract.RuntimeSchemaPositional
Parameters []contract.ParamDecl
DryRun *contract.DryRunSpec
Wait *contract.WaitSpec
Result *contract.ResultSpec
Pagination *contract.PaginationSpec
Interface *contract.InterfaceSpec
@@ -147,9 +146,6 @@ func (s ContractDecl) empty() bool {
if s.DryRun != nil && strings.TrimSpace(s.DryRun.PreviewKind) != "" {
return false
}
if s.Wait != nil && strings.TrimSpace(s.Wait.Mode) != "" {
return false
}
if s.Result != nil {
return false
}
@@ -200,13 +200,6 @@ func TestFrameworkContractFinalDeepCopyAndSafetyConflicts(t *testing.T) {
Parameters: []contract.ParamDecl{{Name: "mode", Enum: []string{"a"}, Required: boolPointer(true)}},
Safety: &contract.SafetySpec{Effect: " read ", EffectSource: " source ", Risk: " low ", Confirmation: " not_required ", Idempotency: " idempotent "},
DryRun: &contract.DryRunSpec{PreviewKind: "plan"},
Wait: &contract.WaitSpec{
Mode: contract.WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{"COMPLETED": contract.ResultOutcomeSuccess, "REJECTED": contract.ResultOutcomeFailure},
PendingValues: []string{"NEW"},
},
Result: &contract.ResultSpec{
Outcomes: []contract.ResultOutcome{contract.ResultOutcomeSuccess},
DataSchema: []byte(`{"type":"object"}`), SensitivePaths: []string{"token"},
@@ -225,8 +218,6 @@ func TestFrameworkContractFinalDeepCopyAndSafetyConflicts(t *testing.T) {
if !ok || got.Result == payload.Result || got.Pagination == payload.Pagination || got.Interface == payload.Interface || got.Selection == payload.Selection || got.Identity == payload.Identity {
t.Fatalf("payload not deeply cloned: %#v", got)
}
payload.Wait.Terminal["COMPLETED"] = contract.ResultOutcomeFailure
payload.Wait.PendingValues[0] = "mutated"
payload.Parameters[0].Enum[0] = "changed"
*payload.Parameters[0].Required = false
*payload.Selection.ExampleDispositions[0].Index = 9
@@ -235,10 +226,6 @@ func TestFrameworkContractFinalDeepCopyAndSafetyConflicts(t *testing.T) {
if again.Parameters[0].Enum[0] != "a" || !*again.Parameters[0].Required || *again.Selection.ExampleDispositions[0].Index != 1 || !*again.Selection.Reviewed {
t.Fatalf("stored payload aliased input: %#v", again)
}
if again.Wait == payload.Wait || again.Wait.Terminal["COMPLETED"] != contract.ResultOutcomeSuccess || again.Wait.PendingValues[0] != "NEW" {
t.Fatalf("wait spec aliased input: %#v", again.Wait)
t.Fatalf("stored payload aliased input: %#v", again)
}
matching := &cobra.Command{Use: "matching"}
t.Cleanup(func() { ClearRuntimeContractFinalForTest(matching) })
-9
View File
@@ -74,15 +74,6 @@ func cloneContractFinalPayload(in contract.ContractFinalPayload) contract.Contra
value := *in.DryRun
out.DryRun = &value
}
if in.Wait != nil {
value := *in.Wait
value.Terminal = make(map[string]contract.ResultOutcome, len(in.Wait.Terminal))
for status, outcome := range in.Wait.Terminal {
value.Terminal[status] = outcome
}
value.PendingValues = cloneSlice(in.Wait.PendingValues)
out.Wait = &value
}
if in.Result != nil {
value := *in.Result
value.Outcomes = cloneSlice(in.Result.Outcomes)
-330
View File
@@ -49,16 +49,12 @@ package corecmd
import (
"bufio"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"math"
"os"
"strconv"
"strings"
"time"
"github.com/mattn/go-isatty"
"github.com/spf13/cobra"
@@ -68,7 +64,6 @@ import (
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/runtimeannotate"
apperrors "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/errors"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/output"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/wait"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/pkg/cmdutil"
)
@@ -287,21 +282,6 @@ type Spec struct {
// Orchestrate executes a multi-step command; it assembles whatever payloads
// it needs from the Ctx.
Orchestrate func(c *Ctx) error
// WaitPoll executes one poll of the declared Contract.Wait capability.
// Exactly one poll is one call; cadence, status extraction, and outcome
// mapping belong to the framework wait phase. Required for poll/auto
// declarations — a declared capability without a runtime implementation
// can never honor --wait, so New rejects the pairing at construction.
// ctx is the wait-phase deadline (--wait-timeout); leaf I/O must honor
// it so a blocked poll cannot outlive the declared timeout.
WaitPoll func(ctx context.Context, c *Ctx) (wait.PollDoc, error)
// WaitEvents opens the push subscription of the declared Contract.Wait
// capability (event/auto modes). The framework owns correlation and
// status mapping; the leaf owns the transport. Auto mode falls back to
// WaitPoll when the stream ends before a terminal status. ctx is the
// same wait-phase deadline as WaitPoll; subscription setup must honor
// it so --wait-timeout can cancel a blocked subscribe.
WaitEvents func(ctx context.Context, c *Ctx) (wait.EventStream, error)
}
// Ctx is the framework-neutral execution context handed to Invoke/Orchestrate.
@@ -375,15 +355,6 @@ func (c *Ctx) Changed(name string) bool { return c.cmd.Flags().Changed(name) }
// DryRun reports the effective global --dry-run.
func (c *Ctx) DryRun() bool { return BoolFlag(c.cmd, "dry-run") }
// Wait reports the effective --wait flag. It is false on commands that did
// not declare the capability: the flag is not registered there, so passing it
// is an unknown-flag error rather than a silently ignored value.
func (c *Ctx) Wait() bool { return BoolFlag(c.cmd, waitFlagName) }
// WaitTimeoutSecs reports the effective --wait-timeout in seconds (flag
// value, then the declared default, then the framework default).
func (c *Ctx) WaitTimeoutSecs() int { return waitTimeoutSecs(c.cmd) }
// Yes reports the effective global --yes.
func (c *Ctx) Yes() bool { return BoolFlag(c.cmd, "yes") }
@@ -403,8 +374,6 @@ func New(spec Spec) *cobra.Command {
validateDispatchDecl(spec)
validateSafetySpec(spec)
validateContractDecl(spec)
normalizeWaitDecl(&spec)
validateWaitDecl(spec)
validateInputSpecs(spec.Use, spec.Flags)
// Help prose inherits the declaration when not authored separately:
// Selection.Examples (already contract-validated against the real flags)
@@ -421,7 +390,6 @@ func New(spec Spec) *cobra.Command {
Hidden: spec.Hidden,
}
RegisterFlags(cmd, spec.Flags)
registerWaitFlags(cmd, spec)
ValidateConstraintDecls(spec.Use, spec.Flags, spec.Constraints)
embedContractIntoSchema(cmd, spec)
AnnotateConstraints(cmd, spec.Constraints)
@@ -487,10 +455,6 @@ func New(spec Spec) *cobra.Command {
if err != nil {
return err
}
result, err = runDeclaredWaitPhase(cmd, args, spec, result)
if err != nil {
return err
}
return output.StoreResult(cmd.Context(), result)
}
return spec.Invoke(ctx, toolArgs)
@@ -542,293 +506,6 @@ func runDeclaredPreflight(cmd *cobra.Command, args []string, spec Spec) error {
return nil
}
// Wait-phase framework flags. They are registered natively on the leaf (never
// as FlagSpec) so they cannot leak into toolArgs / MCP payloads: --wait is a
// client-side execution modifier, not a backend parameter. On commands that
// did not declare Contract.Wait the flags do not exist, so passing --wait
// fails as an unknown flag instead of being silently ignored.
const (
waitFlagName = "wait"
waitTimeoutFlagName = "wait-timeout"
defaultWaitTimeoutS = 300
)
// DefaultWaitTimeoutSecs is the framework default for --wait-timeout when the
// declaration carries no reviewed default.
const DefaultWaitTimeoutSecs = defaultWaitTimeoutS
// normalizeWaitDecl rewrites the spec's declared Wait in place with its
// canonical form (contract.NormalizeWaitSpec): trimmed status values,
// duplicate/conflict rejection, defensive copy. The runtime wait phase and
// AttachContract both read spec.Contract.Wait, so normalizing once at
// construction guarantees the wait engine, the Schema wire, and the
// registered ContractFinal payload all see identical status tables — a
// padded declaration can no longer publish a Schema its own runtime treats
// as unknown statuses. An invalid declaration panics here, next to the
// authoring mistake, with the same message Validate reports.
func normalizeWaitDecl(spec *Spec) {
decl := spec.Contract.Wait
if decl == nil || strings.TrimSpace(decl.Mode) == "" {
return
}
normalized, err := contract.NormalizeWaitSpec(decl, spec.Contract.Identity.CanonicalPath)
if err != nil {
panic(fmt.Sprintf("command %q has invalid Contract.Wait: %v", spec.Use, err))
}
spec.Contract.Wait = normalized
}
// validateWaitDecl enforces the declaration ⇄ implementation pairing at build
// time: a declared Contract.Wait without a WaitPoll hook is a capability the
// command can never honor, and a WaitPoll hook without the declaration has no
// flags or Schema capability to serve. The declaration also requires the
// ResultInvoke dispatcher: only the unified-result envelope can be closed
// into the terminal outcome (error.type "wait", exit code 8) and the timed-out
// pending form — legacy Invoke/Orchestrate/RunE paths emit their own output
// and would observe a failure terminal while still exiting 0. All three
// mismatches are programming errors.
func validateWaitDecl(spec Spec) {
decl := spec.Contract.Wait
declared := decl != nil && strings.TrimSpace(decl.Mode) != ""
if !declared {
if spec.WaitPoll != nil || spec.WaitEvents != nil {
panic(fmt.Sprintf(
"command %q sets a wait hook without declaring Contract.Wait: the wait flags and Schema capability come from the declaration",
spec.Use))
}
return
}
if spec.ResultInvoke == nil {
panic(fmt.Sprintf(
"command %q declares Contract.Wait without ResultInvoke: wait closes the unified-result envelope, which legacy Invoke/Orchestrate/RunE paths cannot rewrite",
spec.Use))
}
mode := strings.TrimSpace(decl.Mode)
needsPoll := mode == contract.WaitModePoll || mode == contract.WaitModeAuto
needsEvent := mode == contract.WaitModeEvent || mode == contract.WaitModeAuto
if needsPoll && spec.WaitPoll == nil {
panic(fmt.Sprintf(
"command %q declares wait mode %s but sets no WaitPoll: a declared wait capability must carry its runtime poll implementation",
spec.Use, mode))
}
if needsEvent && spec.WaitEvents == nil {
panic(fmt.Sprintf(
"command %q declares wait mode %s but sets no WaitEvents: a declared event wait must carry its runtime subscription",
spec.Use, mode))
}
if !needsPoll && spec.WaitPoll != nil {
panic(fmt.Sprintf(
"command %q declares wait mode %s but sets WaitPoll: the declaration decides which hooks run",
spec.Use, mode))
}
if !needsEvent && spec.WaitEvents != nil {
panic(fmt.Sprintf(
"command %q declares wait mode %s but sets WaitEvents: the declaration decides which hooks run",
spec.Use, mode))
}
}
// registerWaitFlags adds --wait / --wait-timeout to a leaf that declared
// Contract.Wait. The timeout default is the reviewed declaration, falling
// back to DefaultWaitTimeoutSecs.
func registerWaitFlags(cmd *cobra.Command, spec Spec) {
decl := spec.Contract.Wait
if decl == nil || strings.TrimSpace(decl.Mode) == "" {
return
}
cmd.Flags().Bool(waitFlagName, false,
"等待到达命令声明的终态(如审批完成、导出结束)后再返回;未声明该能力的命令不接受此 flag")
timeoutDefault := decl.DefaultTimeoutSecs
if timeoutDefault <= 0 {
timeoutDefault = DefaultWaitTimeoutSecs
}
cmd.Flags().Int(waitTimeoutFlagName, timeoutDefault,
"等待超时秒数;超时以 pending 结束(异步受理不是失败)")
}
// runDeclaredWaitPhase runs the declared wait loop after a successful
// ResultInvoke dispatch and closes the accepted unified envelope into the
// wait outcome (validateWaitDecl guarantees the ResultInvoke pairing).
// Only a pending accepted result is waitable: success / failure / partial
// are already terminal and must be returned unchanged. Waiting on a
// business failure would let WithOutcome(..., success) overwrite it into
// an illegal success-with-error envelope.
func runDeclaredWaitPhase(cmd *cobra.Command, args []string, spec Spec, result output.CommandResult) (output.CommandResult, error) {
if !BoolFlag(cmd, waitFlagName) {
return result, nil
}
if result == nil || result.Outcome() != output.OutcomePending {
return result, nil
}
decl := spec.Contract.Wait
timeout, err := waitTimeoutDuration(int64(waitTimeoutSecs(cmd)))
if err != nil {
return result, err
}
ctx := newCtx(cmd, args, spec.Flags)
outcome, err := runWaitLoop(cmd.Context(), decl, timeout, spec, ctx, result)
if err != nil {
return result, err
}
if outcome.TimedOut {
cmd.PrintErrf("等待超时(%s):当前状态 %q,未到达终态,以 pending 结束\n", timeout, outcome.Status)
return output.WithOutcome(result, output.OutcomePending,
output.WithOperationTimedOut(outcome.Status)), nil
}
if outcome.Outcome == contract.ResultOutcomeFailure {
return output.WithOutcome(result, output.OutcomeFailure,
output.WithOperationTerminalState(outcome.Status),
output.WithErrorInfo(&output.ErrorInfo{
Type: "wait",
Subtype: "terminal_failure",
Message: fmt.Sprintf("等待到达失败终态:%s", outcome.Status),
})), nil
}
return output.WithOutcome(result, output.OutcomeSuccess,
output.WithOperationTerminalState(outcome.Status)), nil
}
// runWaitLoop executes the declared wait mode. One deadline spans the event
// phase and an auto-mode poll fallback (the inner loops run without their
// own timeouts and inherit this context's deadline). The deadline is
// forwarded to WaitPoll / WaitEvents and bound onto the cobra command so
// leaf I/O that reads either the hook ctx or Command().Context() is
// cancelled when --wait-timeout expires.
func runWaitLoop(parent context.Context, decl *contract.WaitSpec, timeout time.Duration, spec Spec, ctx *Ctx, result output.CommandResult) (wait.Outcome, error) {
loopCtx := parent
if timeout > 0 {
var cancel context.CancelFunc
loopCtx, cancel = context.WithTimeout(parent, timeout)
defer cancel()
}
if ctx != nil && ctx.cmd != nil {
prev := ctx.cmd.Context()
ctx.cmd.SetContext(loopCtx)
defer ctx.cmd.SetContext(prev)
}
mode := strings.TrimSpace(decl.Mode)
if mode == contract.WaitModePoll {
return wait.Run(loopCtx, wait.LoopSpec{
StatusQuery: decl.StatusQuery,
Terminal: decl.Terminal,
Pending: decl.PendingValues,
}, func(pollCtx context.Context) (wait.PollDoc, error) {
return spec.WaitPoll(pollCtx, ctx)
})
}
resource, err := waitResource(decl, result)
if err != nil {
return wait.Outcome{}, err
}
stream, err := spec.WaitEvents(loopCtx, ctx)
if err != nil {
if loopCtx.Err() != nil {
// Subscribe blocked until the wait deadline: same contract as a
// cancelled poll — close as timed-out pending, do not surface
// ctx.Err() as a subscription failure (and do not poll-fallback
// in auto mode; the shared deadline is already exhausted).
return wait.Outcome{Outcome: contract.ResultOutcomePending, TimedOut: true}, nil
}
if mode == contract.WaitModeAuto {
// Subscription failed before any event: fall back to polling.
return pollWithSpec(loopCtx, decl, spec, ctx)
}
return wait.Outcome{}, fmt.Errorf("wait: event subscription failed: %w", err)
}
eventSpec := wait.EventLoopSpec{
StatusQuery: decl.StatusQuery,
MatchField: decl.MatchField,
Terminal: decl.Terminal,
Pending: decl.PendingValues,
}
outcome, err := wait.RunEvent(loopCtx, eventSpec, resource, stream)
if err == nil {
return outcome, nil
}
if mode == contract.WaitModeAuto && errors.Is(err, wait.ErrEventStreamEnded) {
// Stream ended before a terminal status: fall back to polling under
// the same deadline.
return pollWithSpec(loopCtx, decl, spec, ctx)
}
return outcome, err
}
// pollWithSpec runs the poll loop for an auto-mode fallback.
func pollWithSpec(loopCtx context.Context, decl *contract.WaitSpec, spec Spec, ctx *Ctx) (wait.Outcome, error) {
return wait.Run(loopCtx, wait.LoopSpec{
StatusQuery: decl.StatusQuery,
Terminal: decl.Terminal,
Pending: decl.PendingValues,
}, func(pollCtx context.Context) (wait.PollDoc, error) {
return spec.WaitPoll(pollCtx, ctx)
})
}
// waitResource resolves the resource identifier an event stream correlates
// against, from the accepted result data via the declared ResourceQuery.
// result.Data() returns any deep-copied business payload, which may be a
// map[string]any, struct, or struct pointer. We normalize via JSON round-trip
// to support all valid result types uniformly.
func waitResource(decl *contract.WaitSpec, result output.CommandResult) (string, error) {
raw := result.Data()
if raw == nil {
return "", fmt.Errorf("wait: accepted result data is nil; cannot resolve resource %q", decl.ResourceQuery)
}
// Fast path: already a map.
if data, ok := raw.(map[string]any); ok {
resource, ok := wait.ExtractStatus(wait.PollDoc(data), decl.ResourceQuery)
if !ok || strings.TrimSpace(resource) == "" {
return "", fmt.Errorf("wait: resource query %q not found in accepted result data", decl.ResourceQuery)
}
return resource, nil
}
// Slow path: struct or struct pointer. Normalize via JSON round-trip.
jsonBytes, err := json.Marshal(raw)
if err != nil {
return "", fmt.Errorf("wait: accepted result data cannot be serialized to JSON: %w", err)
}
var data map[string]any
if err := json.Unmarshal(jsonBytes, &data); err != nil {
return "", fmt.Errorf("wait: accepted result data is not an object; cannot resolve resource %q", decl.ResourceQuery)
}
resource, ok := wait.ExtractStatus(wait.PollDoc(data), decl.ResourceQuery)
if !ok || strings.TrimSpace(resource) == "" {
return "", fmt.Errorf("wait: resource query %q not found in accepted result data", decl.ResourceQuery)
}
return resource, nil
}
// waitTimeoutSecs resolves the effective timeout. The flag is registered
// with the reviewed declaration default (or the framework default), so the
// flag value is authoritative; a non-positive explicit value falls back to
// the framework default.
func waitTimeoutSecs(cmd *cobra.Command) int {
if value, err := cmd.Flags().GetInt(waitTimeoutFlagName); err == nil && value > 0 {
return value
}
return DefaultWaitTimeoutSecs
}
// maxWaitTimeoutSecs is the largest second count that still fits in a
// time.Duration. Multiplying a larger int by time.Second overflows to a
// non-positive duration, which would skip the deadline and wait forever.
const maxWaitTimeoutSecs = math.MaxInt64 / int64(time.Second)
// waitTimeoutDuration converts a resolved second count into the wait-phase
// deadline. Values that cannot be represented as a positive time.Duration
// are rejected as validation errors instead of silently disabling timeout.
func waitTimeoutDuration(secs int64) (time.Duration, error) {
if secs <= 0 {
secs = DefaultWaitTimeoutSecs
}
if secs > maxWaitTimeoutSecs {
return 0, apperrors.NewValidation(fmt.Sprintf(
"参数 --%s 取值 %d 超出可表示范围(最大 %d 秒)",
waitTimeoutFlagName, secs, maxWaitTimeoutSecs))
}
return time.Duration(secs) * time.Second, nil
}
// validateDispatchDecl enforces "exactly one dispatcher" at build time. Like
// ValidateConstraintDecls this panics: a spec with no runnable body (or with two
// competing ones) is a programming error that every test and startup path should
@@ -1731,13 +1408,6 @@ func AttachContract(cmd *cobra.Command, safety contract.SafetySpec, decl Contrac
d.PreviewKind = strings.TrimSpace(d.PreviewKind)
payload.DryRun = &d
}
if decl.Wait != nil && strings.TrimSpace(decl.Wait.Mode) != "" {
waitSpec, err := contract.NormalizeWaitSpec(decl.Wait, decl.Identity.CanonicalPath)
if err != nil {
panic(fmt.Sprintf("command %q has invalid Contract.Wait: %v", cmd.Name(), err))
}
payload.Wait = waitSpec
}
if decl.Result != nil {
result, err := contract.NormalizeResultSpec(decl.Result, decl.Identity.CanonicalPath)
if err != nil {
-109
View File
@@ -26,7 +26,6 @@ import (
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contractfinal"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/runtimeannotate"
apperrors "github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/errors"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/output"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/testseam"
"github.com/spf13/cobra"
)
@@ -2121,111 +2120,3 @@ func TestCrossPlatformCoverageEmbedContractSkipsBlankAndHiddenFlags(t *testing.T
}
}
}
// TestWaitResourceStronglyTypedDTOs verifies waitResource handles struct and
// struct pointer results via JSON normalization (P1 fix for auto-CR).
func TestWaitResourceStronglyTypedDTOs(t *testing.T) {
type TaskDTO struct {
TaskID string `json:"task_id"`
Status string `json:"status"`
}
type NestedDTO struct {
Meta struct {
ResourceID string `json:"resource_id"`
} `json:"meta"`
}
tests := []struct {
name string
data any
query string
wantResource string
wantErrSubstr string
}{
{
name: "map[string]any fast path",
data: map[string]any{"task_id": "abc123"},
query: "task_id",
wantResource: "abc123",
},
{
name: "struct value",
data: TaskDTO{TaskID: "struct-456", Status: "running"},
query: "task_id",
wantResource: "struct-456",
},
{
name: "struct pointer",
data: &TaskDTO{TaskID: "ptr-789", Status: "pending"},
query: "task_id",
wantResource: "ptr-789",
},
{
name: "nested struct dotted query",
data: NestedDTO{},
query: "meta.resource_id",
wantResource: "",
wantErrSubstr: "not found",
},
{
name: "nested struct with value",
data: func() NestedDTO {
var d NestedDTO
d.Meta.ResourceID = "nested-xyz"
return d
}(),
query: "meta.resource_id",
wantResource: "nested-xyz",
},
{
name: "nil data",
data: nil,
query: "task_id",
wantErrSubstr: "nil",
},
{
name: "non-object data (string)",
data: "not-an-object",
query: "task_id",
wantErrSubstr: "not an object",
},
{
name: "non-object data (slice)",
data: []string{"a", "b"},
query: "task_id",
wantErrSubstr: "not an object",
},
{
name: "missing query field in struct",
data: TaskDTO{TaskID: "abc"},
query: "nonexistent",
wantErrSubstr: "not found",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
decl := &contract.WaitSpec{
ResourceQuery: tt.query,
}
result := output.Success(tt.data)
resource, err := waitResource(decl, result)
if tt.wantErrSubstr != "" {
if err == nil {
t.Fatalf("expected error containing %q, got nil", tt.wantErrSubstr)
}
if !strings.Contains(err.Error(), tt.wantErrSubstr) {
t.Errorf("error = %q; want substring %q", err.Error(), tt.wantErrSubstr)
}
return
}
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if resource != tt.wantResource {
t.Errorf("resource = %q; want %q", resource, tt.wantResource)
}
})
}
}
-902
View File
@@ -1,902 +0,0 @@
// 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 corecmd
import (
"bytes"
"context"
"errors"
"io"
"math"
"strconv"
"strings"
"testing"
"time"
"github.com/spf13/cobra"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contract"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contractfinal"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/output"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/wait"
)
func waitTestDecl() ContractDecl {
return ContractDecl{
Title: "Wait Title",
Description: "Wait Desc",
Wait: &contract.WaitSpec{
Mode: contract.WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{"COMPLETED": contract.ResultOutcomeSuccess, "REJECTED": contract.ResultOutcomeFailure},
PendingValues: []string{"NEW", "RUNNING"},
DefaultTimeoutSecs: 60,
},
Interface: &contract.InterfaceSpec{Mode: "local", Availability: "available"},
Selection: contract.SelectionSpec{
AgentSummary: "summary",
UseWhen: []string{"when wait"},
AvoidWhen: []string{"when nowait"},
Examples: []string{"dws wait-sample --wait"},
},
Identity: contract.ToolIdentitySpec{ProductID: "sample", Name: "waitsample", CanonicalPath: "sample.waitsample", CLIPath: "wait-sample", PrimaryCLIPath: "wait-sample"},
}
}
func baseWaitSpec(decl ContractDecl, poll func(context.Context, *Ctx) (wait.PollDoc, error)) Spec {
return Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: decl,
WaitPoll: poll,
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1",
State: "NEW",
NextCommand: "dws wait-sample --id job-1",
}), nil
},
}
}
func TestWaitFlagsOnlyRegisteredWhenDeclared(t *testing.T) {
declared := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
if flag := declared.Flags().Lookup(waitFlagName); flag == nil {
t.Fatal("declared command missing --wait flag")
}
if flag := declared.Flags().Lookup(waitTimeoutFlagName); flag == nil {
t.Fatal("declared command missing --wait-timeout flag")
}
undeclared := New(Spec{
Use: "nowait",
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Invoke: func(*Ctx, map[string]any) error { return nil },
})
if flag := undeclared.Flags().Lookup(waitFlagName); flag != nil {
t.Fatal("undeclared command registered --wait")
}
undeclared.SetArgs([]string{"--wait"})
if err := undeclared.Execute(); err == nil || !strings.Contains(err.Error(), "unknown flag") {
t.Fatalf("err=%v want unknown-flag", err)
}
}
func TestValidateWaitDeclPairsDeclarationWithImplementation(t *testing.T) {
decl := waitTestDecl()
spec := baseWaitSpec(decl, nil)
expectPanic(t, func() { New(spec) }, "WaitPoll")
spec.WaitPoll = func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil }
expectPanic(t, func() {
New(Spec{
Use: "hook-only",
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Invoke: func(*Ctx, map[string]any) error { return nil },
WaitPoll: func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil },
})
}, "Contract.Wait")
}
func expectPanic(t *testing.T, fn func(), want string) {
t.Helper()
defer func() {
recovered := recover()
if recovered == nil {
t.Fatalf("expected panic containing %q", want)
}
if message, ok := recovered.(string); !ok || !strings.Contains(message, want) {
t.Fatalf("panic=%v want containing %q", recovered, want)
}
}()
fn()
}
func TestWaitTimeoutFlagDefaultsComeFromDeclaration(t *testing.T) {
stub := func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}
declared := New(baseWaitSpec(waitTestDecl(), stub))
if value, err := declared.Flags().GetInt(waitTimeoutFlagName); err != nil || value != 60 {
t.Fatalf("declared default=%d/%v, want reviewed 60", value, err)
}
decl := waitTestDecl()
decl.Wait.DefaultTimeoutSecs = 0
fallback := New(baseWaitSpec(decl, stub))
if value, err := fallback.Flags().GetInt(waitTimeoutFlagName); err != nil || value != DefaultWaitTimeoutSecs {
t.Fatalf("fallback default=%d/%v, want framework %d", value, err, DefaultWaitTimeoutSecs)
}
// A non-positive explicit value falls back to the framework default.
fallback.SetArgs([]string{"--wait-timeout", "0", "--wait"})
if err := fallback.Flags().Set(waitTimeoutFlagName, "0"); err != nil {
t.Fatal(err)
}
if got := waitTimeoutSecs(fallback); got != DefaultWaitTimeoutSecs {
t.Fatalf("waitTimeoutSecs=%d, want %d", got, DefaultWaitTimeoutSecs)
}
}
func TestWaitTimeoutDurationRejectsOverflowingSeconds(t *testing.T) {
// math.MaxInt64 (9223372036854775807) is a legal pflag int on 64-bit
// platforms and overflows time.Duration(secs)*time.Second to a negative
// value, which would disable the wait deadline.
if _, err := waitTimeoutDuration(math.MaxInt64); err == nil || !strings.Contains(err.Error(), "超出可表示范围") {
t.Fatalf("err=%v, want overflow validation", err)
}
d, err := waitTimeoutDuration(maxWaitTimeoutSecs)
if err != nil {
t.Fatal(err)
}
if d <= 0 || d != time.Duration(maxWaitTimeoutSecs)*time.Second {
t.Fatalf("duration=%d, want the largest representable timeout", d)
}
// Non-positive second counts fall back to the framework default instead
// of disabling the deadline (waitTimeoutSecs already maps a zero/negative
// flag to the default; this keeps the conversion itself fail-safe).
if got, err := waitTimeoutDuration(0); err != nil || got != time.Duration(DefaultWaitTimeoutSecs)*time.Second {
t.Fatalf("duration/err=%d/%v, want framework default", got, err)
}
if got, err := waitTimeoutDuration(-5); err != nil || got != time.Duration(DefaultWaitTimeoutSecs)*time.Second {
t.Fatalf("duration/err=%d/%v, want framework default", got, err)
}
}
func TestResultInvokeWaitTimeoutOverflowIsValidationError(t *testing.T) {
if int64(math.MaxInt) <= maxWaitTimeoutSecs {
t.Skip("platform int cannot overflow time.Duration")
}
polled := false
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", strconv.Itoa(math.MaxInt)})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !strings.Contains(err.Error(), "超出可表示范围") {
t.Fatalf("err=%v, want overflow validation", err)
}
if polled {
t.Fatal("overflowing --wait-timeout must not start the wait loop")
}
}
func TestResultInvokeWaitPollErrorFailsTheCommand(t *testing.T) {
boom := errors.New("rpc down")
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return nil, boom
}))
cmd.SetArgs([]string{"--wait"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !strings.Contains(err.Error(), "rpc down") {
t.Fatalf("err=%v, want poll error surfaced", err)
}
}
func TestResultInvokeWaitUnknownStatusFailsClosed(t *testing.T) {
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "Mystery"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !wait.IsUnknownStatus(err) {
t.Fatalf("err=%v, want unknown-status", err)
}
}
func TestWaitCtxAccessorsExposeDeclaredCapability(t *testing.T) {
var gotWait bool
var gotTimeout int
cmd := New(baseWaitSpec(waitTestDecl(), func(_ context.Context, c *Ctx) (wait.PollDoc, error) {
gotWait = c.Wait()
gotTimeout = c.WaitTimeoutSecs()
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", "90"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if !gotWait || gotTimeout != 90 {
t.Fatalf("ctx accessors=%v/%d", gotWait, gotTimeout)
}
}
func eventTestDecl(mode string) ContractDecl {
decl := waitTestDecl()
decl.Wait.Mode = mode
decl.Wait.EventKey = "bpms_instance_change"
decl.Wait.MatchField = "process_instance_id"
decl.Wait.ResourceQuery = "id"
return decl
}
type scriptedStream struct {
events []wait.PollDoc
err error
}
func (s *scriptedStream) Recv(context.Context) (wait.PollDoc, error) {
if len(s.events) > 0 {
doc := s.events[0]
s.events = s.events[1:]
return doc, nil
}
if s.err != nil {
return nil, s.err
}
return nil, io.EOF
}
func TestValidateWaitDeclPairsModeWithHooks(t *testing.T) {
poll := func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil }
events := func(context.Context, *Ctx) (wait.EventStream, error) { return nil, nil }
cases := []struct {
name string
mode string
waitPoll bool
waitEvents bool
want string
}{
{"event without WaitEvents", contract.WaitModeEvent, false, false, "WaitEvents"},
{"auto without WaitPoll", contract.WaitModeAuto, false, true, "WaitPoll"},
{"poll with WaitEvents", contract.WaitModePoll, true, true, "WaitEvents"},
{"event with WaitPoll", contract.WaitModeEvent, true, true, "WaitPoll"},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
expectPanic(t, func() {
New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(tc.mode),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{}, nil), nil
},
WaitPoll: hookOrNil(tc.waitPoll, poll),
WaitEvents: eventHookOrNil(tc.waitEvents, events),
})
}, tc.want)
})
}
}
func hookOrNil(set bool, hook func(context.Context, *Ctx) (wait.PollDoc, error)) func(context.Context, *Ctx) (wait.PollDoc, error) {
if !set {
return nil
}
return hook
}
func eventHookOrNil(set bool, hook func(context.Context, *Ctx) (wait.EventStream, error)) func(context.Context, *Ctx) (wait.EventStream, error) {
if !set {
return nil
}
return hook
}
func runWaitModeCommand(t *testing.T, decl ContractDecl, poll func(context.Context, *Ctx) (wait.PollDoc, error), events func(context.Context, *Ctx) (wait.EventStream, error), args ...string) (string, error) {
t.Helper()
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: decl,
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample --id job-1",
}), nil
},
WaitPoll: poll,
WaitEvents: events,
})
cmd.SetArgs(append([]string{"--wait"}, args...))
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
err := cmd.Execute()
return stdout.String(), err
}
func TestEventModeClosesEnvelopeFromCorrelatedEvent(t *testing.T) {
stream := &scriptedStream{events: []wait.PollDoc{
{"process_instance_id": "other", "result": map[string]any{"status": "COMPLETED"}},
{"process_instance_id": "job-1", "result": map[string]any{"status": "REJECTED"}},
}}
stdout, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeEvent), nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return stream, nil
})
if err != nil {
t.Fatal(err)
}
if !strings.Contains(stdout, `"outcome": "failure"`) || !strings.Contains(stdout, `"type": "wait"`) {
t.Fatalf("stdout=%s", stdout)
}
}
func TestEventModeSurfacesStreamEndAsError(t *testing.T) {
_, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeEvent), nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil
})
if err == nil || !errors.Is(err, wait.ErrEventStreamEnded) {
t.Fatalf("err=%v, want stream-ended", err)
}
}
func TestEventModeRejectsUnresolvableResource(t *testing.T) {
decl := eventTestDecl(contract.WaitModeEvent)
decl.Wait.ResourceQuery = "missing"
_, err := runWaitModeCommand(t, decl, nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil
})
if err == nil || !strings.Contains(err.Error(), "resource query") {
t.Fatalf("err=%v", err)
}
}
func TestAutoModeFallsBackToPollOnStreamEnd(t *testing.T) {
polled := false
stdout, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeAuto),
func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
},
func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil // ends immediately
})
if err != nil {
t.Fatal(err)
}
if !polled {
t.Fatal("auto mode did not fall back to polling")
}
if !strings.Contains(stdout, `"outcome": "success"`) {
t.Fatalf("stdout=%s", stdout)
}
}
func TestAutoModeFallsBackToPollOnSubscriptionFailure(t *testing.T) {
polled := false
_, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeAuto),
func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
},
func(context.Context, *Ctx) (wait.EventStream, error) {
return nil, errors.New("no subscriber credential")
})
if err != nil {
t.Fatal(err)
}
if !polled {
t.Fatal("auto mode did not fall back to polling on subscription failure")
}
}
func TestResultInvokeWaitClosesEnvelopeOutcome(t *testing.T) {
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "REJECTED"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 8 {
t.Fatalf("stored code/emitted=%d/%v, want dedicated wait-terminal-failure code 8", code, emitted)
}
if !strings.Contains(stdout.String(), `"type": "wait"`) {
t.Fatalf("stdout=%s, want error.type wait", stdout.String())
}
if !strings.Contains(stdout.String(), `"outcome": "failure"`) {
t.Fatalf("stdout=%s", stdout.String())
}
// The final emitted envelope must carry the observed terminal status in
// meta.operation.state — the acceptance-phase state ("NEW") must not
// survive the close (P1 regression guard).
if !strings.Contains(stdout.String(), `"state": "REJECTED"`) {
t.Fatalf("stdout=%s, want operation.state synced to the terminal status", stdout.String())
}
if strings.Contains(stdout.String(), `"state": "NEW"`) {
t.Fatalf("stdout=%s, acceptance-phase operation.state leaked into the terminal envelope", stdout.String())
}
if strings.Contains(stdout.String(), `"timed_out": true`) {
t.Fatalf("stdout=%s, terminal close must not claim timed_out", stdout.String())
}
}
func TestResultInvokeWaitSuccessCloseSyncsOperationState(t *testing.T) {
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v, want success exit 0", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "success"`) {
t.Fatalf("stdout=%s", stdout.String())
}
// Success close must publish the terminal status, never the stale
// acceptance-phase state (no outcome=success with state=processing/NEW).
if !strings.Contains(stdout.String(), `"state": "COMPLETED"`) {
t.Fatalf("stdout=%s, want operation.state synced to the terminal status", stdout.String())
}
if strings.Contains(stdout.String(), `"state": "NEW"`) {
t.Fatalf("stdout=%s, acceptance-phase operation.state leaked into the success envelope", stdout.String())
}
// Operation identity (id / next_command) survives the terminal close.
if !strings.Contains(stdout.String(), `"id": "job-1"`) || !strings.Contains(stdout.String(), `"next_command"`) {
t.Fatalf("stdout=%s, want operation id/next_command preserved", stdout.String())
}
}
func TestResultInvokeWaitTimeoutKeepsPending(t *testing.T) {
polls := 0
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
polls++
return wait.PollDoc{"result": map[string]any{"status": "RUNNING"}}, nil
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
if polls == 0 {
t.Fatal("wait phase never polled")
}
}
func TestWaitDeclRequiresResultInvokeDispatcher(t *testing.T) {
// A declared wait on the legacy Invoke path would observe a failure
// terminal while still exiting 0 — construction must reject it.
expectPanic(t, func() {
New(Spec{
Use: "wait-sample",
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: waitTestDecl(),
WaitPoll: func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil },
Invoke: func(*Ctx, map[string]any) error { return nil },
})
}, "ResultInvoke")
}
func TestResultInvokeWithoutWaitFlagSkipsPhase(t *testing.T) {
polled := false
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs(nil)
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if polled {
t.Fatal("wait phase ran without --wait")
}
}
func TestAttachContractPanicsOnInvalidWaitDeclaration(t *testing.T) {
decl := waitTestDecl()
decl.Wait.Mode = "event" // not implemented
defer func() {
recovered := recover()
if recovered == nil {
t.Fatal("expected panic on invalid Contract.Wait")
}
if message, ok := recovered.(string); !ok || !strings.Contains(message, "Contract.Wait") {
t.Fatalf("panic=%v", recovered)
}
}()
AttachContract(&cobra.Command{Use: "wait-sample"}, contract.SafetySpec{
Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent",
}, decl, "", "")
}
func TestWaitDeclPaddedStatusValuesAreNormalized(t *testing.T) {
// A declaration whose status values carry surrounding whitespace must be
// canonicalized at construction so the runtime wait engine and the
// published Schema agree: the backend returns "COMPLETED" verbatim, and
// a padded terminal key would fail closed as an unknown status.
decl := waitTestDecl()
decl.Wait.Terminal = map[string]contract.ResultOutcome{
" COMPLETED ": contract.ResultOutcomeSuccess,
"\tREJECTED": contract.ResultOutcomeFailure,
}
decl.Wait.PendingValues = []string{" NEW ", "RUNNING "}
cmd := New(baseWaitSpec(decl, func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatalf("padded declaration must still reach the terminal status: %v", err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v, want success", code, emitted)
}
final, ok := contractfinal.RuntimeContractFinal(cmd)
if !ok || final.Wait == nil {
t.Fatal("registered ContractFinal lost the wait capability")
}
if _, ok := final.Wait.Terminal["COMPLETED"]; !ok {
t.Fatalf("registered terminal table not trimmed: %#v", final.Wait.Terminal)
}
for _, value := range final.Wait.PendingValues {
if strings.TrimSpace(value) != value {
t.Fatalf("registered pending value %q not trimmed", value)
}
}
}
func TestNewPanicsOnDuplicateOrConflictingWaitStatusesAfterTrim(t *testing.T) {
// Values that collapse onto one status after trimming are programming
// errors: silently merging them would pick one outcome for two authored
// declarations.
dupTerminal := waitTestDecl()
dupTerminal.Wait.Terminal = map[string]contract.ResultOutcome{
"COMPLETED": contract.ResultOutcomeSuccess,
" COMPLETED": contract.ResultOutcomeFailure,
}
expectPanic(t, func() { New(baseWaitSpec(dupTerminal, nil)) }, "Contract.Wait")
conflict := waitTestDecl()
conflict.Wait.PendingValues = []string{" COMPLETED "}
expectPanic(t, func() { New(baseWaitSpec(conflict, nil)) }, "Contract.Wait")
}
func TestContractDeclEmptyTreatsWaitAsAuthored(t *testing.T) {
// Only Wait is authored: empty() must report non-empty through the Wait
// branch (before validateContractDecl then fails on the missing prose).
decl := ContractDecl{Wait: &contract.WaitSpec{
Mode: contract.WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{"COMPLETED": contract.ResultOutcomeSuccess},
}}
if decl.Empty() {
t.Fatal("Wait-only declaration must count as authored")
}
defer func() {
if recover() == nil {
t.Fatal("expected validateContractDecl to reject the missing prose")
}
}()
validateContractDecl(Spec{Use: "wait-only", Contract: decl})
}
func TestEventModeRejectsNonObjectResultData(t *testing.T) {
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeEvent),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending([]any{"not", "an", "object"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample",
}), nil
},
WaitEvents: func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil
},
})
cmd.SetArgs([]string{"--wait"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !strings.Contains(err.Error(), "not an object") {
t.Fatalf("err=%v, want non-object data rejection", err)
}
}
func TestEventModeSubscriptionFailureSurfacesInStrictMode(t *testing.T) {
decl := eventTestDecl(contract.WaitModeEvent)
_, err := runWaitModeCommand(t, decl, nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return nil, errors.New("no subscriber credential")
})
if err == nil || !strings.Contains(err.Error(), "subscription failed") {
t.Fatalf("err=%v, want subscription failure surfaced", err)
}
}
func TestResultInvokeNonPendingSkipsWaitPhase(t *testing.T) {
partial, err := output.NewPartialData(2,
[]any{map[string]any{"id": "ok"}},
[]output.PartialFailedEntry{{ID: "bad", Error: &output.ErrorInfo{Type: "api", Message: "item failed"}}},
nil)
if err != nil {
t.Fatal(err)
}
cases := []struct {
name string
result output.CommandResult
want string
}{
{"failure", output.Failure(&output.ErrorInfo{Type: "api", Message: "business failed"}), "failure"},
{"success", output.Success(map[string]any{"id": "job-1"}), "success"},
{"partial", output.Partial(partial), "partial_failure"},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
polled := false
subscribed := false
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeAuto),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return tc.result, nil
},
WaitPoll: func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
},
WaitEvents: func(context.Context, *Ctx) (wait.EventStream, error) {
subscribed = true
return &scriptedStream{}, nil
},
})
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if polled || subscribed {
t.Fatal("wait phase must not call WaitPoll/WaitEvents for a non-pending initial result")
}
if _, emitted := output.StoredExitCode(store); !emitted {
t.Fatal("initial result was not stored")
}
if !strings.Contains(stdout.String(), `"outcome": "`+tc.want+`"`) {
t.Fatalf("stdout=%s, want outcome %s preserved", stdout.String(), tc.want)
}
if strings.Contains(stdout.String(), `"type": "wait"`) {
t.Fatalf("stdout=%s, wait phase overwrote the original envelope", stdout.String())
}
})
}
}
func TestWaitTimeoutCancelsBlockingPoll(t *testing.T) {
started := make(chan struct{})
cmd := New(baseWaitSpec(waitTestDecl(), func(ctx context.Context, c *Ctx) (wait.PollDoc, error) {
close(started)
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-c.Command().Context().Done():
return nil, c.Command().Context().Err()
}
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
done := make(chan error, 1)
go func() { done <- cmd.Execute() }()
select {
case <-started:
case <-time.After(2 * time.Second):
t.Fatal("blocking poll never started")
}
select {
case err := <-done:
if err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
case <-time.After(3 * time.Second):
t.Fatal("blocked poll was not cancelled by --wait-timeout")
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
}
func TestWaitTimeoutCancelsBlockingSubscribe(t *testing.T) {
started := make(chan struct{})
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeEvent),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample --id job-1",
}), nil
},
WaitEvents: func(ctx context.Context, c *Ctx) (wait.EventStream, error) {
close(started)
// Leaf subscribe may wait on either the hook ctx or the cobra
// command context; both must carry the wait-timeout deadline.
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-c.Command().Context().Done():
return nil, c.Command().Context().Err()
}
},
})
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
done := make(chan error, 1)
go func() { done <- cmd.Execute() }()
select {
case <-started:
case <-time.After(2 * time.Second):
t.Fatal("blocking subscribe never started")
}
select {
case err := <-done:
if err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
case <-time.After(3 * time.Second):
t.Fatal("blocked subscribe was not cancelled by --wait-timeout")
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
}
func TestWaitTimeoutCancelsBlockingPollAfterAutoFallback(t *testing.T) {
started := make(chan struct{})
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeAuto),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample --id job-1",
}), nil
},
WaitEvents: func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil // ends immediately → poll fallback
},
WaitPoll: func(ctx context.Context, _ *Ctx) (wait.PollDoc, error) {
close(started)
<-ctx.Done()
return nil, ctx.Err()
},
})
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
done := make(chan error, 1)
go func() { done <- cmd.Execute() }()
select {
case <-started:
case <-time.After(2 * time.Second):
t.Fatal("auto-fallback blocking poll never started")
}
select {
case err := <-done:
if err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
case <-time.After(3 * time.Second):
t.Fatal("blocked auto-fallback poll was not cancelled by --wait-timeout")
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
}
-2
View File
@@ -58,8 +58,6 @@ const (
// 5 internal (CategoryInternal 与兜底:非结构化错误、panic 收敛均归 5)
// 6 discovery (CategoryDiscovery)
// 7 partial_failure(部分成功专用码,见 ExitCodePartial)
// 8 wait (--wait 观察到失败终态的专用码,见 internal/output
// 的 exitCodeWait;不设 Category,仅经统一信封产出)
//
// ExitCodePartial is the partial-result exit code shared with internal/output.
// It is not returned for CategoryPartial errors because they lack the typed
-1
View File
@@ -132,7 +132,6 @@ func TestClientCreateRuleBasedSubscriptionsUsesDocumentedRuleParam(t *testing.T)
{"oa_approval_task_finished", EventOAApprovalTaskFinished, RuleOptions{}, map[string]any{}},
{"oa_approval_task_redirected", EventOAApprovalTaskRedirected, RuleOptions{}, map[string]any{}},
{"oa_approval_instance_started", EventOAApprovalInstanceStarted, RuleOptions{}, map[string]any{}},
{"oa_approval_instance_cc", EventOAApprovalInstanceCC, RuleOptions{}, map[string]any{}},
{"oa_approval_instance_terminated", EventOAApprovalInstanceTerminated, RuleOptions{}, map[string]any{}},
{"oa_approval_instance_finished", EventOAApprovalInstanceFinished, RuleOptions{}, map[string]any{}},
{"read_group", EventReadGroup, RuleOptions{GroupID: "cid-1"}, map[string]any{"openConversationId": "cid-1"}},
-29
View File
@@ -180,19 +180,6 @@ type OAApprovalInstanceStartedOutput struct {
EventTime int64 `json:"event_time" description:"审批实例事件业务时间" format:"timestamp_ms"`
}
type OAApprovalInstanceCCOutput struct {
Type string `json:"type" description:"事件类型,固定为当前 event_key"`
EventID string `json:"event_id" description:"事件 ID,可用于去重"`
Timestamp int64 `json:"timestamp" description:"事件发生时间戳" format:"timestamp_ms"`
SubscribeID string `json:"subscribe_id" description:"订阅 ID"`
ProcessInstanceID string `json:"process_instance_id" description:"审批实例 ID"`
ProcessCode string `json:"process_code" description:"审批流程模板编码"`
Title string `json:"title" description:"审批标题"`
Status string `json:"status" description:"审批实例到达抄送节点时的状态"`
CreateTime int64 `json:"create_time" description:"审批实例创建时间" format:"timestamp_ms"`
EventTime int64 `json:"event_time" description:"审批抄送事件业务时间" format:"timestamp_ms"`
}
type OAApprovalInstanceTerminatedOutput struct {
Type string `json:"type" description:"事件类型,固定为当前 event_key"`
EventID string `json:"event_id" description:"事件 ID,可用于去重"`
@@ -680,19 +667,6 @@ func projectOAApprovalEvent(ev transport.Event, base baseEventOutput, raw json.R
CreateTime: payload.Body.CreateTime,
EventTime: payload.EventTime,
}, nil
case EventOAApprovalInstanceCC:
return OAApprovalInstanceCCOutput{
Type: base.Type,
EventID: base.EventID,
Timestamp: base.Timestamp,
SubscribeID: base.SubscribeID,
ProcessInstanceID: payload.Body.ProcessInstanceID,
ProcessCode: payload.Body.ProcessCode,
Title: payload.Body.Title,
Status: payload.Body.Status,
CreateTime: payload.Body.CreateTime,
EventTime: payload.EventTime,
}, nil
case EventOAApprovalInstanceTerminated:
return OAApprovalInstanceTerminatedOutput{
Type: base.Type,
@@ -897,8 +871,6 @@ func outputTypeForEvent(eventKey string) reflect.Type {
return reflect.TypeOf(OAApprovalTaskRedirectedOutput{})
case eventKey == EventOAApprovalInstanceStarted:
return reflect.TypeOf(OAApprovalInstanceStartedOutput{})
case eventKey == EventOAApprovalInstanceCC:
return reflect.TypeOf(OAApprovalInstanceCCOutput{})
case eventKey == EventOAApprovalInstanceTerminated:
return reflect.TypeOf(OAApprovalInstanceTerminatedOutput{})
case eventKey == EventOAApprovalInstanceFinished:
@@ -934,7 +906,6 @@ func isOAEvent(eventKey string) bool {
eventKey == EventOAApprovalTaskFinished ||
eventKey == EventOAApprovalTaskRedirected ||
eventKey == EventOAApprovalInstanceStarted ||
eventKey == EventOAApprovalInstanceCC ||
eventKey == EventOAApprovalInstanceTerminated ||
eventKey == EventOAApprovalInstanceFinished
}
-18
View File
@@ -173,8 +173,6 @@ func personalOAData(eventKey string) string {
body["finishTime"] = int64(1785229199000)
case EventOAApprovalInstanceStarted:
body["status"] = "RUNNING"
case EventOAApprovalInstanceCC:
body["status"] = "RUNNING"
case EventOAApprovalInstanceTerminated:
body["status"] = "TERMINATED"
body["finishTime"] = int64(1785229199000)
@@ -523,21 +521,6 @@ func TestCrossPlatformCoverageProjectOutputOAEvents(t *testing.T) {
EventTime: 1785229199000,
},
},
{
eventKey: EventOAApprovalInstanceCC,
want: OAApprovalInstanceCCOutput{
Type: EventOAApprovalInstanceCC,
EventID: "oa-event",
Timestamp: 1785229200123,
SubscribeID: "outer-sub",
ProcessInstanceID: "process-instance-1",
ProcessCode: "PROC-TEST-1",
Title: "测试审批",
Status: "RUNNING",
CreateTime: 1785229100000,
EventTime: 1785229199000,
},
},
{
eventKey: EventOAApprovalInstanceTerminated,
want: OAApprovalInstanceTerminatedOutput{
@@ -802,7 +785,6 @@ func TestCrossPlatformCoverageProjectOutputRejectsInvalidOAPayloads(t *testing.T
EventOAApprovalTaskFinished,
EventOAApprovalTaskRedirected,
EventOAApprovalInstanceStarted,
EventOAApprovalInstanceCC,
EventOAApprovalInstanceTerminated,
EventOAApprovalInstanceFinished,
} {
-12
View File
@@ -43,7 +43,6 @@ const (
EventOAApprovalTaskFinished = "user_oa_approval_task_finished"
EventOAApprovalTaskRedirected = "user_oa_approval_task_redirected"
EventOAApprovalInstanceStarted = "user_oa_approval_instance_started"
EventOAApprovalInstanceCC = "user_oa_approval_instance_cc"
EventOAApprovalInstanceTerminated = "user_oa_approval_instance_terminated"
EventOAApprovalInstanceFinished = "user_oa_approval_instance_finished"
)
@@ -324,17 +323,6 @@ var definitions = []Definition{
Auth: map[string]any{"identity": "user"},
Public: true,
},
{
EventKey: EventOAApprovalInstanceCC,
DisplayName: "审批单抄送",
Description: "审批实例到达抄送节点,发送给被抄送人",
Category: "oa",
RuleType: "all",
Status: StatusEnabled,
RequiredParams: nil,
Auth: map[string]any{"identity": "user"},
Public: true,
},
{
EventKey: EventOAApprovalInstanceTerminated,
DisplayName: "审批单终止",
-12
View File
@@ -50,7 +50,6 @@ func TestCatalogEnabledEvents(t *testing.T) {
EventOAApprovalTaskFinished,
EventOAApprovalTaskRedirected,
EventOAApprovalInstanceStarted,
EventOAApprovalInstanceCC,
EventOAApprovalInstanceTerminated,
EventOAApprovalInstanceFinished,
}
@@ -66,7 +65,6 @@ func TestOAEventCatalogDefinitions(t *testing.T) {
EventOAApprovalTaskFinished,
EventOAApprovalTaskRedirected,
EventOAApprovalInstanceStarted,
EventOAApprovalInstanceCC,
EventOAApprovalInstanceTerminated,
EventOAApprovalInstanceFinished,
}
@@ -160,7 +158,6 @@ func TestSchemaDocumentsDefaultToTransportEnvelope(t *testing.T) {
EventOAApprovalTaskFinished,
EventOAApprovalTaskRedirected,
EventOAApprovalInstanceStarted,
EventOAApprovalInstanceCC,
EventOAApprovalInstanceTerminated,
EventOAApprovalInstanceFinished,
} {
@@ -512,13 +509,6 @@ func TestOAEventSchemaDocumentsMatchOutputDTO(t *testing.T) {
"process_code", "title", "status", "create_time", "event_time",
},
},
{
eventKey: EventOAApprovalInstanceCC,
properties: []string{
"type", "event_id", "timestamp", "subscribe_id", "process_instance_id",
"process_code", "title", "status", "create_time", "event_time",
},
},
{
eventKey: EventOAApprovalInstanceTerminated,
properties: []string{
@@ -643,7 +633,6 @@ func TestBuildRuleParamAllEvents(t *testing.T) {
EventOAApprovalTaskFinished,
EventOAApprovalTaskRedirected,
EventOAApprovalInstanceStarted,
EventOAApprovalInstanceCC,
EventOAApprovalInstanceTerminated,
EventOAApprovalInstanceFinished,
} {
@@ -856,7 +845,6 @@ func TestSupportsMessageFilter(t *testing.T) {
EventOAApprovalTaskFinished,
EventOAApprovalTaskRedirected,
EventOAApprovalInstanceStarted,
EventOAApprovalInstanceCC,
EventOAApprovalInstanceTerminated,
EventOAApprovalInstanceFinished,
"unknown_event",
-3
View File
@@ -435,7 +435,6 @@ const (
exitCodeInternal = 5
exitCodeDiscovery = 6
exitCodePartial = 7 // partial_failure 专用(契约 §4;规划 WS2 第4项;B142 将在 errors 侧补同源常量)
exitCodeWait = 8 // wait 终态失败专用(--wait 观察到失败终态;与 partial 同为"仅新增专用码")
)
// subtypeConfirmationRequired 是门禁拦截的 failure 子类标记(契约规范 §2.4),
@@ -492,8 +491,6 @@ func exitCodeForErrorInfo(info *ErrorInfo) int {
return exitCodePermission
case "discovery":
return exitCodeDiscovery
case "wait":
return exitCodeWait
default:
return exitCodeInternal
}
+2 -4
View File
@@ -283,11 +283,9 @@ func (e *ErrorInfo) Validate() error {
}
// error.type is a wire-stable Agent branch key, not an open-ended label.
// Keep this set aligned with exitCodeForErrorInfo. "permission" is the
// compatibility projection for PAT failures (rc=4); "wait" is the wait
// phase's terminal-failure projection (rc=8, e.g. an approval observed
// REJECTED after --wait).
// compatibility projection for PAT failures (rc=4).
switch errorType {
case "api", "auth", "validation", "permission", "discovery", "internal", "wait":
case "api", "auth", "validation", "permission", "discovery", "internal":
default:
return fmt.Errorf("output: unsupported failure error.type %q", e.Type)
}
@@ -19,7 +19,6 @@ type forgedResult struct {
func (r forgedResult) Outcome() Outcome { return r.env.Outcome }
func (r forgedResult) ExitCode() int { return r.exit }
func (r forgedResult) Data() any { return r.env.Data }
func (r forgedResult) envelope() *Envelope { copy := r.env; return &copy }
type cloneNode struct {
-90
View File
@@ -19,10 +19,6 @@ import (
type CommandResult interface {
Outcome() Outcome
ExitCode() int
// Data returns the accepted payload (already deep-copied). The wait
// phase reads it to resolve the resource identifier an event stream
// correlates against.
Data() any
envelope() *Envelope
}
@@ -34,7 +30,6 @@ type commandResult struct {
func (r *commandResult) Outcome() Outcome { return r.env.Outcome }
func (r *commandResult) ExitCode() int { return r.exitCode }
func (r *commandResult) Data() any { return cloneResultData(r.env.Data) }
func (r *commandResult) envelope() *Envelope {
copy := cloneEnvelope(r.env)
return &copy
@@ -76,91 +71,6 @@ func Partial(data *PartialData, opts ...ResultOption) CommandResult {
return newCommandResult(OutcomePartialFailure, data, nil, opts...)
}
// WithMeta(meta *Meta) ResultOption was declared above; the two options below
// exist for the wait phase.
// WithErrorInfo replaces the error info of the envelope. Used when a wait
// phase closes an accepted result into failure: envelope invariant I3
// requires an error iff the outcome is failure.
func WithErrorInfo(info *ErrorInfo) ResultOption {
return ResultOption{apply: func(env *Envelope) {
if info != nil {
info = cloneErrorInfo(info)
}
env.Error = info
}}
}
// WithOperationTimedOut marks the envelope's async operation as timed out at
// the last observed state, preserving the declared id / next_command resume
// facts (契约规范 §2.2: 超时必须保持 State 真实值并置 TimedOut:true). A result
// without operation info keeps nil — the pending envelope invariant then fails
// at emission, surfacing the leaf bug instead of synthesizing fake resume
// facts.
func WithOperationTimedOut(state string) ResultOption {
return ResultOption{apply: func(env *Envelope) {
if env.Meta == nil || env.Meta.Operation == nil {
return
}
operation := *env.Meta.Operation
// A subscribe/poll that never observed a status still times out
// against the accepted pending result: keep the original state
// rather than wiping it to empty (pending requires operation.state).
if strings.TrimSpace(state) != "" {
operation.State = state
}
operation.TimedOut = true
env.Meta.Operation = &operation
}}
}
// WithOperationTerminalState closes the envelope's async operation at the observed
// terminal status (契约规范 §2.2: 终态封装必须同步 operation.state — a success or
// failure close that kept the acceptance-phase state would emit a
// self-contradicting envelope such as outcome=success with
// operation.state=processing). The declared id / next_command facts are kept
// as the operation identity, and timed_out is cleared: the §2.2 anti-spoof
// rule forbids a timed-out claim on an operation that reached a terminal
// state. A result without operation info is left untouched.
func WithOperationTerminalState(state string) ResultOption {
return ResultOption{apply: func(env *Envelope) {
if env.Meta == nil || env.Meta.Operation == nil {
return
}
operation := *env.Meta.Operation
if strings.TrimSpace(state) != "" {
operation.State = state
}
operation.TimedOut = false
env.Meta.Operation = &operation
}}
}
// WithOutcome rewraps an existing result with a new outcome, preserving data,
// meta, identity, and any error info (subject to the opts). The corecmd wait
// phase uses it to close an accepted result into its terminal (or timed-out
// pending) outcome; the exit code is re-derived from the new envelope.
func WithOutcome(result CommandResult, outcome Outcome, opts ...ResultOption) CommandResult {
env := *result.envelope()
for _, opt := range opts {
if opt.apply != nil {
opt.apply(&env)
}
}
env.Outcome = outcome
env.OK = outcome == OutcomeSuccess || outcome == OutcomePending
if outcome == OutcomeFailure {
// Invariant I3: data and error are mutually exclusive. Closing into
// failure replaces the accepted data with the failure error.
env.Data = nil
}
exitCode := ExitCodeForEnvelope(&env)
if env.Error != nil {
env.Error.ExitCode = exitCode
}
return &commandResult{env: env, exitCode: exitCode}
}
// Failure constructs an immutable typed failure result.
func Failure(info *ErrorInfo, opts ...ResultOption) CommandResult {
if info != nil {
-165
View File
@@ -1,165 +0,0 @@
// 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 output
import (
"strings"
"testing"
)
func pendingAcceptedResult() CommandResult {
return Pending(map[string]any{"id": "job-1"}, &OperationInfo{
ID: "job-1",
State: "NEW",
NextCommand: "dws wait-sample --id job-1",
})
}
func TestWithOutcomeClosesSuccessPreservingDataAndMeta(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomeSuccess)
if result.Outcome() != OutcomeSuccess || result.ExitCode() != 0 {
t.Fatalf("outcome=%s exit=%d", result.Outcome(), result.ExitCode())
}
env := result.envelope()
if env.Meta == nil || env.Meta.Operation == nil {
t.Fatal("operation info lost on success close")
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOutcomeFailureDropsDataAndCarriesError(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomeFailure, WithErrorInfo(&ErrorInfo{
Type: "wait",
Subtype: "terminal_failure",
Message: "等待到达失败终态:REJECTED",
}))
if result.Outcome() != OutcomeFailure || result.ExitCode() != 8 {
t.Fatalf("outcome=%s exit=%d, want failure/8", result.Outcome(), result.ExitCode())
}
env := result.envelope()
if env.Data != nil {
t.Fatal("failure close must drop data (I3)")
}
if env.Error == nil || env.Error.Type != "wait" {
t.Fatal("failure close must carry error info (I3)")
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTimedOutMarksStateAndKeepsResumeFacts(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomePending, WithOperationTimedOut("RUNNING"))
env := result.envelope()
op := env.Meta.Operation
if op.State != "RUNNING" || !op.TimedOut || op.ID != "job-1" || op.NextCommand != "dws wait-sample --id job-1" {
t.Fatalf("operation=%+v", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTimedOutEmptyStatePreservesExisting(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomePending, WithOperationTimedOut(""))
op := result.envelope().Meta.Operation
if op.State != "NEW" || !op.TimedOut {
t.Fatalf("operation=%+v, want original state kept and timed_out set", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTimedOutWithoutOperationInfoLeavesEnvelopeUntouched(t *testing.T) {
// A result without operation info keeps nil — ValidateResult must then
// reject the pending envelope instead of the option synthesizing fake
// resume facts.
result := WithOutcome(Success(map[string]any{"ok": true}), OutcomePending, WithOperationTimedOut("RUNNING"))
env := result.envelope()
if env.Meta != nil && env.Meta.Operation != nil {
t.Fatalf("operation=%+v, want untouched", env.Meta.Operation)
}
err := ValidateResult(result)
if err == nil || !strings.Contains(err.Error(), "meta.operation") {
t.Fatalf("err=%v, want pending-requires-operation rejection", err)
}
}
func TestDataAccessorReturnsDeepCopy(t *testing.T) {
result := pendingAcceptedResult()
data, ok := result.Data().(map[string]any)
if !ok {
t.Fatalf("data=%#v", result.Data())
}
data["id"] = "mutated"
again := result.Data().(map[string]any)
if again["id"] != "job-1" {
t.Fatalf("Data() aliased internal state: %#v", again)
}
}
func TestWithOperationTerminalStateSyncsStateAndClearsTimedOut(t *testing.T) {
for _, tc := range []struct {
name string
outcome Outcome
}{
{"success", OutcomeSuccess},
{"failure", OutcomeFailure},
} {
t.Run(tc.name, func(t *testing.T) {
opts := []ResultOption{WithOperationTerminalState("COMPLETED")}
if tc.outcome == OutcomeFailure {
opts = append(opts, WithErrorInfo(&ErrorInfo{
Type: "wait", Subtype: "terminal_failure", Message: "等待到达失败终态:COMPLETED",
}))
}
result := WithOutcome(pendingAcceptedResult(), tc.outcome, opts...)
op := result.envelope().Meta.Operation
if op.State != "COMPLETED" || op.TimedOut {
t.Fatalf("operation=%+v, want terminal state synced and timed_out cleared", op)
}
if op.ID != "job-1" || op.NextCommand != "dws wait-sample --id job-1" {
t.Fatalf("operation=%+v, want resume identity facts preserved", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
})
}
}
func TestWithOperationTerminalStateBlankKeepsObservedState(t *testing.T) {
// A terminal observation that carries no status must keep the last known
// state rather than wipe it: operation.state may never be emptied by a
// close transition.
result := WithOutcome(pendingAcceptedResult(), OutcomeSuccess, WithOperationTerminalState(" "))
op := result.envelope().Meta.Operation
if op.State != "NEW" || op.TimedOut {
t.Fatalf("operation=%+v, want original state kept and timed_out cleared", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTerminalStateWithoutOperationInfoLeavesEnvelopeUntouched(t *testing.T) {
result := WithOutcome(Success(map[string]any{"ok": true}), OutcomeSuccess, WithOperationTerminalState("COMPLETED"))
env := result.envelope()
if env.Meta != nil && env.Meta.Operation != nil {
t.Fatalf("operation=%+v, want untouched", env.Meta.Operation)
}
}
-273
View File
@@ -1,273 +0,0 @@
// 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 wait is the framework terminal-state wait engine behind the
// reviewed contract.WaitSpec capability. It owns polling cadence, status
// extraction, and status→outcome mapping; it knows nothing about Cobra, MCP,
// or any product backend. How one poll executes is supplied by the leaf's
// WaitPoll hook (corecmd), so "poll = an existing read command" stays a leaf
// decision rather than a framework assumption.
package wait
import (
"context"
"errors"
"fmt"
"strconv"
"strings"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contract"
)
// DefaultPollInterval is the cadence between polls when LoopSpec.Interval is
// zero. The first poll runs immediately so an already-terminal resource does
// not pay a sleep tax.
const DefaultPollInterval = 2 * time.Second
// MaxPollInterval caps the exponential backoff growth between polls so a long
// wait cannot degenerate into effectively-blind polling.
const MaxPollInterval = 30 * time.Second
// PollDoc is one decoded poll response document (typically the unified-output
// envelope data of the poll command).
type PollDoc map[string]any
// Poller executes one poll. Returning an error fails the wait phase; the
// engine never retries a poller error because read commands failing is a real
// failure, not a "not yet" signal.
type Poller func(ctx context.Context) (PollDoc, error)
// LoopSpec is the runtime-resolved projection of contract.WaitSpec plus the
// caller-provided timeout.
type LoopSpec struct {
StatusQuery string
Terminal map[string]contract.ResultOutcome
Pending []string
Timeout time.Duration
Interval time.Duration
}
// Outcome is the closed result of a wait loop. TimedOut reports deadline
// exhaustion (Outcome is then pending — an accepted-but-not-terminal state is
// not a process failure per the exit-code contract); Status is the last
// observed status value.
type Outcome struct {
Status string
Outcome contract.ResultOutcome
Attempts int
TimedOut bool
}
// ErrUnknownStatus reports a status value that is neither declared terminal
// nor declared pending. Unknown fails closed: mapping it to pending could
// hide a real state change until timeout, mapping it to success is worse.
type ErrUnknownStatus struct {
Status string
Query string
}
func (e *ErrUnknownStatus) Error() string {
return fmt.Sprintf("wait: status %q (from %q) is neither terminal nor pending", e.Status, e.Query)
}
// Run polls poller until a declared terminal status, deadline exhaustion, or
// a poller error. The first poll is immediate; subsequent polls back off
// exponentially (×1.5) from Interval, capped at MaxPollInterval. Deadline
// exhaustion anywhere — before a poll, during a poll (a context-aware poller
// returns ctx.Err()), or during the wait between polls — always closes as
// timed-out pending with the last observed status, never as a poll failure.
func Run(ctx context.Context, spec LoopSpec, poll Poller) (Outcome, error) {
if spec.Interval <= 0 {
spec.Interval = DefaultPollInterval
}
if spec.Timeout > 0 {
var cancel context.CancelFunc
ctx, cancel = context.WithTimeout(ctx, spec.Timeout)
defer cancel()
}
pending := make(map[string]bool, len(spec.Pending))
for _, value := range spec.Pending {
pending[value] = true
}
timedOut := func(status string, attempts int) Outcome {
return Outcome{Status: status, Outcome: contract.ResultOutcomePending, Attempts: attempts, TimedOut: true}
}
interval := spec.Interval
attempts := 0
lastStatus := ""
for {
if ctx.Err() != nil {
return timedOut(lastStatus, attempts), nil
}
doc, err := poll(ctx)
if err != nil {
if ctx.Err() != nil {
return timedOut(lastStatus, attempts), nil
}
return Outcome{Attempts: attempts}, fmt.Errorf("wait: poll failed: %w", err)
}
attempts++
status, ok := ExtractStatus(doc, spec.StatusQuery)
if !ok {
return Outcome{Attempts: attempts}, fmt.Errorf(
"wait: status query %q not found in poll result", spec.StatusQuery)
}
lastStatus = status
if outcome, ok := spec.Terminal[status]; ok {
return Outcome{Status: status, Outcome: outcome, Attempts: attempts}, nil
}
if !pending[status] {
return Outcome{Status: status, Attempts: attempts}, &ErrUnknownStatus{Status: status, Query: spec.StatusQuery}
}
timer := time.NewTimer(interval)
select {
case <-ctx.Done():
timer.Stop()
return timedOut(status, attempts), nil
case <-timer.C:
}
interval = nextInterval(interval)
}
}
func nextInterval(current time.Duration) time.Duration {
next := current * 3 / 2
if next > MaxPollInterval {
next = MaxPollInterval
}
return next
}
// ExtractStatus resolves a dotted status query against a poll document. Each
// segment walks one map level; array indexes are not supported because wait
// targets a single resource. Numeric segments are stringified, so a document
// decoded with json.Number keys still resolves.
func ExtractStatus(doc PollDoc, query string) (string, bool) {
query = strings.TrimSpace(query)
if query == "" {
return "", false
}
// PollDoc is a defined type, so its dynamic type does not satisfy a
// map[string]any assertion — convert once at the boundary; nested values
// from JSON decoding are plain maps.
var current any = map[string]any(doc)
for _, segment := range strings.Split(query, ".") {
segment = strings.TrimSpace(segment)
if segment == "" {
return "", false
}
node, ok := current.(map[string]any)
if !ok {
return "", false
}
value, ok := node[segment]
if !ok {
return "", false
}
current = value
}
switch value := current.(type) {
case string:
return value, true
case fmt.Stringer:
return value.String(), true
case bool:
return strconv.FormatBool(value), true
case int:
return strconv.Itoa(value), true
case int64:
return strconv.FormatInt(value, 10), true
case float64:
return strconv.FormatFloat(value, 'f', -1, 64), true
default:
return "", false
}
}
// IsUnknownStatus reports whether err is the closed fail-on-unknown error.
func IsUnknownStatus(err error) bool {
var unknown *ErrUnknownStatus
return errors.As(err, &unknown)
}
// EventStream is the leaf-owned push subscription consumed by the event
// phase (the WaitEvents hook in corecmd). Recv delivers the next decoded
// event document; it returns an error or io.EOF-style termination when the
// stream ends — the engine treats non-terminal termination as a stream
// failure the caller (auto mode) may fall back from.
type EventStream interface {
Recv(ctx context.Context) (PollDoc, error)
}
// EventLoopSpec is the event-phase projection of contract.WaitSpec.
type EventLoopSpec struct {
StatusQuery string
MatchField string
Terminal map[string]contract.ResultOutcome
Pending []string
Timeout time.Duration
}
// ErrEventStreamEnded reports a stream that terminated before a terminal
// status. Auto mode uses it to fall back to polling; strict event mode
// surfaces it as a wait failure.
var ErrEventStreamEnded = errors.New("wait: event stream ended before a terminal status")
// RunEvent consumes stream until a correlated event reaches a declared
// terminal status, the deadline exhausts, or the stream ends. Events whose
// MatchField value does not equal resource are ignored (other resources on
// the same channel); a correlated event with an unknown status fails closed
// exactly like a poll would.
func RunEvent(ctx context.Context, spec EventLoopSpec, resource string, stream EventStream) (Outcome, error) {
if spec.Timeout > 0 {
var cancel context.CancelFunc
ctx, cancel = context.WithTimeout(ctx, spec.Timeout)
defer cancel()
}
pending := make(map[string]bool, len(spec.Pending))
for _, value := range spec.Pending {
pending[value] = true
}
attempts := 0
lastStatus := ""
for {
doc, err := stream.Recv(ctx)
if err != nil {
if ctx.Err() != nil {
return Outcome{Status: lastStatus, Outcome: contract.ResultOutcomePending, Attempts: attempts, TimedOut: true}, nil
}
// Wrap with ErrEventStreamEnded so auto mode can fall back to
// polling while correlated-status failures (unknown status,
// missing status query) stay non-recoverable.
return Outcome{Attempts: attempts}, fmt.Errorf("%w: %v", ErrEventStreamEnded, err)
}
attempts++
correlated, ok := ExtractStatus(doc, spec.MatchField)
if !ok || correlated != resource {
continue
}
status, ok := ExtractStatus(doc, spec.StatusQuery)
if !ok {
return Outcome{Attempts: attempts}, fmt.Errorf(
"wait: status query %q not found in event document", spec.StatusQuery)
}
lastStatus = status
if outcome, ok := spec.Terminal[status]; ok {
return Outcome{Status: status, Outcome: outcome, Attempts: attempts}, nil
}
if !pending[status] {
return Outcome{Status: status, Attempts: attempts}, &ErrUnknownStatus{Status: status, Query: spec.StatusQuery}
}
}
}
-365
View File
@@ -1,365 +0,0 @@
// 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 wait
import (
"context"
"errors"
"io"
"strings"
"testing"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contract"
)
func loopSpec() LoopSpec {
return LoopSpec{
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{
"COMPLETED": contract.ResultOutcomeSuccess,
"REJECTED": contract.ResultOutcomeFailure,
},
Pending: []string{"NEW", "RUNNING"},
Interval: time.Millisecond,
}
}
func TestExtractStatusResolvesDottedPaths(t *testing.T) {
doc := PollDoc{
"result": map[string]any{
"instance": map[string]any{"status": "RUNNING"},
"count": float64(3),
},
}
if status, ok := ExtractStatus(doc, "result.instance.status"); !ok || status != "RUNNING" {
t.Fatalf("status=%q ok=%v", status, ok)
}
if status, ok := ExtractStatus(doc, "result.count"); !ok || status != "3" {
t.Fatalf("numeric status=%q ok=%v", status, ok)
}
if _, ok := ExtractStatus(doc, "result.missing"); ok {
t.Fatal("missing path resolved")
}
if _, ok := ExtractStatus(doc, "result.instance.status.deep"); ok {
t.Fatal("descending into a scalar resolved")
}
if _, ok := ExtractStatus(doc, ""); ok {
t.Fatal("empty query resolved")
}
}
func TestRunReturnsTerminalOnFirstPoll(t *testing.T) {
polls := 0
outcome, err := Run(context.Background(), loopSpec(), func(context.Context) (PollDoc, error) {
polls++
return PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
})
if err != nil {
t.Fatal(err)
}
if polls != 1 || outcome.Attempts != 1 {
t.Fatalf("polls=%d attempts=%d", polls, outcome.Attempts)
}
if outcome.Outcome != contract.ResultOutcomeSuccess || outcome.Status != "COMPLETED" {
t.Fatalf("outcome=%s status=%s", outcome.Outcome, outcome.Status)
}
}
func TestRunPollsUntilTerminal(t *testing.T) {
seen := []string{"NEW", "RUNNING", "RUNNING", "COMPLETED"}
index := 0
outcome, err := Run(context.Background(), loopSpec(), func(context.Context) (PollDoc, error) {
status := seen[index]
index++
return PollDoc{"result": map[string]any{"status": status}}, nil
})
if err != nil {
t.Fatal(err)
}
if outcome.Attempts != len(seen) || outcome.Outcome != contract.ResultOutcomeSuccess {
t.Fatalf("attempts=%d outcome=%s", outcome.Attempts, outcome.Outcome)
}
}
func TestRunTimesOutAsPendingDuringWait(t *testing.T) {
spec := loopSpec()
spec.Timeout = 5 * time.Millisecond
polls := 0
outcome, err := Run(context.Background(), spec, func(context.Context) (PollDoc, error) {
polls++
return PollDoc{"result": map[string]any{"status": "RUNNING"}}, nil
})
if err != nil {
t.Fatal(err)
}
if !outcome.TimedOut || outcome.Outcome != contract.ResultOutcomePending {
t.Fatalf("timedOut=%v outcome=%s", outcome.TimedOut, outcome.Outcome)
}
if outcome.Status != "RUNNING" {
t.Fatalf("status=%q, want last observed", outcome.Status)
}
if polls == 0 {
t.Fatal("timeout during wait must still have polled at least once")
}
}
func TestRunTimesOutAsPendingWhenPollerRespectsDeadline(t *testing.T) {
spec := loopSpec()
spec.Timeout = 5 * time.Millisecond
polls := 0
// A context-aware poller: blocks until the deadline, then reports the
// cancellation as an error — the loop must close it as timed-out pending,
// never as a poll failure.
outcome, err := Run(context.Background(), spec, func(ctx context.Context) (PollDoc, error) {
polls++
<-ctx.Done()
return nil, ctx.Err()
})
if err != nil {
t.Fatalf("deadline during poll closed as error: %v", err)
}
if !outcome.TimedOut || outcome.Outcome != contract.ResultOutcomePending {
t.Fatalf("timedOut=%v outcome=%s", outcome.TimedOut, outcome.Outcome)
}
}
func TestRunTimesOutBeforeFirstPoll(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel()
outcome, err := Run(ctx, loopSpec(), func(context.Context) (PollDoc, error) {
t.Fatal("poller ran on a pre-cancelled context")
return nil, nil
})
if err != nil {
t.Fatal(err)
}
if !outcome.TimedOut || outcome.Attempts != 0 || outcome.Outcome != contract.ResultOutcomePending {
t.Fatalf("timedOut=%v attempts=%d outcome=%s", outcome.TimedOut, outcome.Attempts, outcome.Outcome)
}
}
func TestRunFailsClosedOnUnknownStatus(t *testing.T) {
_, err := Run(context.Background(), loopSpec(), func(context.Context) (PollDoc, error) {
return PollDoc{"result": map[string]any{"status": "Mystery"}}, nil
})
if !IsUnknownStatus(err) {
t.Fatalf("err=%v want unknown-status", err)
}
}
func TestRunFailsOnMissingStatusQuery(t *testing.T) {
_, err := Run(context.Background(), loopSpec(), func(context.Context) (PollDoc, error) {
return PollDoc{"unexpected": true}, nil
})
if err == nil || !errors.Is(err, err) {
t.Fatalf("err=%v", err)
}
}
func TestRunPropagatesPollerError(t *testing.T) {
boom := errors.New("rpc down")
_, err := Run(context.Background(), loopSpec(), func(context.Context) (PollDoc, error) {
return nil, boom
})
if !errors.Is(err, boom) {
t.Fatalf("err=%v", err)
}
}
func TestUnknownStatusErrorCarriesStatusAndQuery(t *testing.T) {
err := &ErrUnknownStatus{Status: "Mystery", Query: "result.status"}
message := err.Error()
if !strings.Contains(message, "Mystery") || !strings.Contains(message, "result.status") {
t.Fatalf("message=%q", message)
}
}
func TestRunAppliesDefaultIntervalWhenUnset(t *testing.T) {
spec := loopSpec()
spec.Interval = 0
polls := 0
// Terminal on the second poll forces one interval wait; with Interval=0
// the loop must still work using DefaultPollInterval (not spin/panic).
_, err := Run(context.Background(), spec, func(context.Context) (PollDoc, error) {
polls++
if polls == 1 {
return PollDoc{"result": map[string]any{"status": "NEW"}}, nil
}
return PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
})
if err != nil {
t.Fatal(err)
}
if polls != 2 {
t.Fatalf("polls=%d", polls)
}
}
func TestRunTimesOutDuringWaitBetweenPolls(t *testing.T) {
spec := loopSpec()
spec.Interval = time.Hour // the deadline wins long before the next poll
spec.Timeout = 5 * time.Millisecond
outcome, err := Run(context.Background(), spec, func(context.Context) (PollDoc, error) {
return PollDoc{"result": map[string]any{"status": "RUNNING"}}, nil
})
if err != nil {
t.Fatal(err)
}
if !outcome.TimedOut || outcome.Outcome != contract.ResultOutcomePending || outcome.Status != "RUNNING" {
t.Fatalf("outcome=%+v", outcome)
}
}
type stringStatus string
func (s stringStatus) String() string { return string(s) }
func TestExtractStatusCoversScalarShapes(t *testing.T) {
doc := PollDoc{
"result": map[string]any{
"flag": true,
"small": 7,
"big": int64(9007199254740993),
"fraction": 1.5,
"custom": stringStatus("CUSTOM"),
"nested": map[string]any{"deep": "x"},
},
}
cases := map[string]string{
"result.flag": "true",
"result.small": "7",
"result.big": "9007199254740993",
"result.fraction": "1.5",
"result.custom": "CUSTOM",
}
for query, want := range cases {
if got, ok := ExtractStatus(doc, query); !ok || got != want {
t.Fatalf("query=%s got=%q ok=%v want=%q", query, got, ok, want)
}
}
if _, ok := ExtractStatus(doc, "result.nested"); ok {
t.Fatal("non-scalar nested map must not resolve")
}
if _, ok := ExtractStatus(doc, "result..flag"); ok {
t.Fatal("empty segment must not resolve")
}
}
func TestNextIntervalCapsAtMax(t *testing.T) {
if got := nextInterval(MaxPollInterval); got != MaxPollInterval {
t.Fatalf("nextInterval(max)=%s", got)
}
if got := nextInterval(10 * time.Millisecond); got != 15*time.Millisecond {
t.Fatalf("nextInterval(10ms)=%s", got)
}
}
type fakeEventStream struct {
events []PollDoc
err error // returned after events are exhausted (nil = clean end)
block bool // hold until the context deadline
}
func (f *fakeEventStream) Recv(ctx context.Context) (PollDoc, error) {
if f.block {
<-ctx.Done()
return nil, ctx.Err()
}
if len(f.events) > 0 {
doc := f.events[0]
f.events = f.events[1:]
return doc, nil
}
if f.err != nil {
return nil, f.err
}
return nil, io.EOF
}
func eventLoopSpec() EventLoopSpec {
return EventLoopSpec{
StatusQuery: "result.status",
MatchField: "process_instance_id",
Terminal: map[string]contract.ResultOutcome{
"COMPLETED": contract.ResultOutcomeSuccess,
"REJECTED": contract.ResultOutcomeFailure,
},
Pending: []string{"RUNNING"},
}
}
func approvalEvent(instance, status string) PollDoc {
return PollDoc{"process_instance_id": instance, "result": map[string]any{"status": status}}
}
func TestRunEventReturnsCorrelatedTerminal(t *testing.T) {
stream := &fakeEventStream{events: []PollDoc{
approvalEvent("other-instance", "COMPLETED"), // other resource: ignored
approvalEvent("job-1", "RUNNING"), // correlated pending: kept waiting
approvalEvent("job-1", "COMPLETED"),
}}
outcome, err := RunEvent(context.Background(), eventLoopSpec(), "job-1", stream)
if err != nil {
t.Fatal(err)
}
if outcome.Outcome != contract.ResultOutcomeSuccess || outcome.Status != "COMPLETED" {
t.Fatalf("outcome=%+v", outcome)
}
}
func TestRunEventFailsClosedOnUnknownCorrelatedStatus(t *testing.T) {
stream := &fakeEventStream{events: []PollDoc{approvalEvent("job-1", "Mystery")}}
_, err := RunEvent(context.Background(), eventLoopSpec(), "job-1", stream)
if !IsUnknownStatus(err) {
t.Fatalf("err=%v", err)
}
}
func TestRunEventRejectsCorrelatedEventWithoutStatus(t *testing.T) {
stream := &fakeEventStream{events: []PollDoc{
{"process_instance_id": "job-1"}, // correlated but no status document
}}
_, err := RunEvent(context.Background(), eventLoopSpec(), "job-1", stream)
if err == nil || !strings.Contains(err.Error(), "status query") {
t.Fatalf("err=%v", err)
}
}
func TestRunEventStreamEndSurfacesFallbackSentinel(t *testing.T) {
stream := &fakeEventStream{events: []PollDoc{approvalEvent("job-1", "RUNNING")}}
_, err := RunEvent(context.Background(), eventLoopSpec(), "job-1", stream)
if !errors.Is(err, ErrEventStreamEnded) {
t.Fatalf("err=%v, want ErrEventStreamEnded", err)
}
failing := &fakeEventStream{err: errors.New("transport reset")}
_, err = RunEvent(context.Background(), eventLoopSpec(), "job-1", failing)
if !errors.Is(err, ErrEventStreamEnded) {
t.Fatalf("err=%v, want ErrEventStreamEnded wrapping the transport error", err)
}
}
func TestRunEventTimesOutAsPendingWhileBlocked(t *testing.T) {
spec := eventLoopSpec()
spec.Timeout = 5 * time.Millisecond
stream := &fakeEventStream{block: true}
outcome, err := RunEvent(context.Background(), spec, "job-1", stream)
if err != nil {
t.Fatal(err)
}
if !outcome.TimedOut || outcome.Outcome != contract.ResultOutcomePending {
t.Fatalf("outcome=%+v", outcome)
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
---
name: dws
description: 管理钉钉产品能力(AI表格/AI搜问/日历/通讯录/群聊与机器人/待办/审批/考勤/日志/DING消息/开放平台文档/钉钉文档/钉钉云盘/原生Markdown文件/AI听记/邮箱/在线电子表格/知识库等)。当用户需要操作表格数据、管理日程会议、模糊找人/查谁负责某事项、查询通讯录、管理群聊、机器人发消息、创建待办、提交审批、查看考勤、提交日报周报(钉钉日志模版)、读写钉钉文档、上传下载云盘文件、读取或修改原生.md文件、查询听记纪要、收发邮件、读写在线电子表格(axls)、管理钉钉知识库,或订阅个人 IM 事件或 OA 审批事件、实时监听群成员加入、群成员退出、群改名和群解散、审批实例发起/抄送/终止/完成,以及审批任务创建/完成/转交时使用。
description: 管理钉钉产品能力(AI表格/AI搜问/日历/通讯录/群聊与机器人/待办/审批/考勤/日志/DING消息/开放平台文档/钉钉文档/钉钉云盘/原生Markdown文件/AI听记/邮箱/在线电子表格/知识库等)。当用户需要操作表格数据、管理日程会议、模糊找人/查谁负责某事项、查询通讯录、管理群聊、机器人发消息、创建待办、提交审批、查看考勤、提交日报周报(钉钉日志模版)、读写钉钉文档、上传下载云盘文件、读取或修改原生.md文件、查询听记纪要、收发邮件、读写在线电子表格(axls)、管理钉钉知识库,或订阅个人 IM 事件或 OA 审批事件、实时监听群成员加入、群成员退出、群改名和群解散、审批实例发起/终止/完成,以及审批任务创建/完成/转交时使用。
cli_version: ">=1.0.15"
---
+4 -9
View File
@@ -47,11 +47,10 @@
| `user_oa_approval_task_finished` | 审批任务已完成 | 无 |
| `user_oa_approval_task_redirected` | 审批任务已转交 | 无 |
| `user_oa_approval_instance_started` | 审批实例已发起 | 无 |
| `user_oa_approval_instance_cc` | 审批实例到达抄送节点,发送给被抄送人 | 无 |
| `user_oa_approval_instance_terminated` | 审批实例已终止 | 无 |
| `user_oa_approval_instance_finished` | 审批实例完成,发送给审批单发起人 | 无 |
只承认上表 23 个事件码。默认身份就是当前用户,使用当前用户 OAuth 登录态,不要额外加身份切换 flag。七个 OA 事件订阅当前用户相关的全部审批事件,规则均为 `all`、空 `filterRule`,不需要目标参数。
只承认上表 22 个事件码。默认身份就是当前用户,使用当前用户 OAuth 登录态,不要额外加身份切换 flag。六个 OA 事件订阅当前用户相关的全部审批事件,规则均为 `all`、空 `filterRule`,不需要目标参数。
## Intent mapping
@@ -77,10 +76,9 @@
| "审批任务完成时通知我" | `event consume`,事件码 `user_oa_approval_task_finished`,参数 `--flatten -f ndjson` |
| "审批任务被转交时通知我" | `event consume`,事件码 `user_oa_approval_task_redirected`,参数 `--flatten -f ndjson` |
| "有审批单发起时通知我" | `event consume`,事件码 `user_oa_approval_instance_started`,参数 `--flatten -f ndjson` |
| "有审批抄送给我时通知我" | `event consume`,事件码 `user_oa_approval_instance_cc`,参数 `--flatten -f ndjson` |
| "有审批单终止时通知我" | `event consume`,事件码 `user_oa_approval_instance_terminated`,参数 `--flatten -f ndjson` |
| "监听我发起的审批何时完成" / "审批实例完成时通知我" | `event consume`,事件码 `user_oa_approval_instance_finished`,参数 `--flatten -f ndjson` |
| "同时监听全部已公开 OA 事件" | 一个 consume 放入七个 OA event key,不加目标或消息过滤参数 |
| "同时监听全部已公开 OA 事件" | 一个 consume 放入六个 OA event key,不加目标或消息过滤参数 |
| "查看个人事件 schema" | `dws event schema <event_key> --flatten` |
| "看个人事件订阅状态" | `dws event status --event <event_key>` |
| "停止这个个人事件订阅" | `dws event stop <subscribe_id> --dry-run`,确认后改用 `--yes` |
@@ -132,7 +130,6 @@ dws event schema user_oa_approval_task_created --flatten
dws event schema user_oa_approval_task_finished --flatten
dws event schema user_oa_approval_task_redirected --flatten
dws event schema user_oa_approval_instance_started --flatten
dws event schema user_oa_approval_instance_cc --flatten
dws event schema user_oa_approval_instance_terminated --flatten
dws event schema user_oa_approval_instance_finished --flatten
```
@@ -160,7 +157,6 @@ dws event consume user_oa_approval_task_created --flatten -f ndjson
dws event consume user_oa_approval_task_finished --flatten -f ndjson
dws event consume user_oa_approval_task_redirected --flatten -f ndjson
dws event consume user_oa_approval_instance_started --flatten -f ndjson
dws event consume user_oa_approval_instance_cc --flatten -f ndjson
dws event consume user_oa_approval_instance_terminated --flatten -f ndjson
dws event consume user_oa_approval_instance_finished --flatten -f ndjson
```
@@ -189,14 +185,13 @@ dws event consume \
user_oa_approval_task_finished \
user_oa_approval_task_redirected \
user_oa_approval_instance_started \
user_oa_approval_instance_cc \
user_oa_approval_instance_terminated \
user_oa_approval_instance_finished \
--flatten \
-f ndjson
```
用户类事件共享 `--user` 或 `--open-dingtalk-id`,群类事件共享 `--group`,无目标 IM 事件可加入任一组合。用户类与群类、不同目标或不同过滤条件要拆成多个进程。七个 OA 事件可以同进程消费并共享 personal bus,但各自建立独立订阅。多事件共享 `--query` / `--filter-json` 时,所选事件必须全部是 IM 消息接收事件;OA 事件单独或组合消费都禁止使用这两个消息过滤参数。
用户类事件共享 `--user` 或 `--open-dingtalk-id`,群类事件共享 `--group`,无目标 IM 事件可加入任一组合。用户类与群类、不同目标或不同过滤条件要拆成多个进程。六个 OA 事件可以同进程消费并共享 personal bus,但各自建立独立订阅。多事件共享 `--query` / `--filter-json` 时,所选事件必须全部是 IM 消息接收事件;OA 事件单独或组合消费都禁止使用这两个消息过滤参数。
上述所有 `*_o2o` 命令和 `user_im_message_receive_user` 都可将 `--user <userId>` 替换为 `--open-dingtalk-id <openDingtalkId>`,但两个参数不能同时使用。
@@ -221,7 +216,7 @@ dws event stop --all --yes
## 订阅创建失败与重试预算
以下约束适用于上表全部 23 个公开个人事件(16 个 IM + 7 个 OA)以及多事件命令中的每一项,只治理 `[event] ready` 之前的订阅创建;ready 之后的 Stream 断线由长连接重连机制处理。
以下约束适用于上表全部 22 个公开个人事件(16 个 IM + 6 个 OA)以及多事件命令中的每一项,只治理 `[event] ready` 之前的订阅创建;ready 之后的 Stream 断线由长连接重连机制处理。
- `0/2/1` 是 **Agent/host 编排约束**,不是 CLI 持久化硬总次数上限。每次 `dws event consume` 调用对每个逻辑订阅最多发送一次订阅创建 HTTP 请求,进程内不会自动重试。CLI 本地状态只持久化 `in_flight`、`cooldown`、`terminal_hold` 三种保护状态,不持久化或计算跨调用的 Agent/host 尝试次数。
- 解析人名或群名、执行 `event consume` 以及后续 `event status/stop` 必须使用同一个 `--profile`。不得把其它 profile 下解析出的 userId、openDingtalkId 或 openConversationId 直接带入当前 profile 的订阅。
+42 -62
View File
@@ -11,11 +11,13 @@ metadata:
# 钉钉日历 Skill
## 前置条件 — 执行操作前必读
## 执行契约
> **CRITICAL — 执行任何 `dws` 操作前,MUST 先用 Read 工具完整读取 [`dingtalk-shared`](../dingtalk-shared/SKILL.md)。**该轻量文件包含全局执行契约、安全底线及 shared references 的按需加载导航;不要预加载其全部 references。
> 命令参考:[calendar.md](references/calendar.md);剧本:[03-meeting.md](references/03-meeting.md)。
- 明确的 Calendar 请求直接按本 Skill 执行;仅在跨产品、profile、确认或错误恢复需要时读取 [`dingtalk-shared`](../dingtalk-shared/SKILL.md) 的对应章节。
- 已知意图直接走下方 Golden Route。只有 leaf 参数或安全语义不确定时读取单个 compact Schema,只有 Cobra flag 不确定时读取精确 leaf Help。
- 所有目标解析、读取、写入和验证使用同一 profile。`eventId`、`roomId`、`calendarId`、`userId` 只取真实返回;零命中或多候选时停止消歧。
- 写操作按 Runtime confirmation gate 执行:先解析并展示目标;需要确认时获得确认后才加 `--yes`。退出码为 0 不等于业务完成,必须检查结构化 outcome 和读回证据。
- 默认只加载一个操作 Reference:通用原子命令读 [calendar.md](references/calendar.md);涉及会议室预订读 [03-meeting.md](references/03-meeting.md)。
<!-- VISIBLE_SHORTCUTS_START -->
## Shortcuts(无专用脚本/recipe 时优先)
@@ -47,70 +49,48 @@ metadata:
| `dws calendar +week` | read | 列出我本周的日程(自动按周一为周首计算本周起止时间,无需手动填时间范围) |
<!-- VISIBLE_SHORTCUTS_END -->
## 意图表
## Golden Routes
| 用户说 | 命令 |
|--------|------|
| "今天 / 明天 / 本周日程" | `python scripts/calendar_today_agenda.py [today\|tomorrow\|week]` |
| "约会议(含参会人 + 会议室)" | `python scripts/calendar_schedule_meeting.py --title "<主题>" --start "<起>" --end "<止>" [--users <ids>] [--book-room]` |
| "多人共同空闲" | `python scripts/calendar_free_slot_finder.py --users <ids> --date <yyyy-MM-dd>` |
| "查闲忙" | `dws calendar busy search --users <userIds> --start "<ISO>" --end "<ISO>"` |
| "加参会人" / "订房" / "取消" | `dws calendar attendee add` / `room add` / `event delete` |
| 用户意图 | 首选入口 | 身份与边界 |
|---|---|---|
| 今天 / 明天 / 本周日程 | `dws calendar +today|+tomorrow|+week` | 当前 profile 的主日历;无需手算时间窗 |
| 任意时段日程 | `dws calendar +agenda --start "<ISO>" --end "<ISO>"` | 保留 `eventId` 和分页证据;`hasMore` 时继续翻页 |
| 按标题、描述或地点找日程 | `dws calendar +search-event --query "<关键词>"` | 单页零命中且 `hasMore=true` 不是全局零命中 |
| 创建个人日程或按姓名约人 | `dws calendar +book --title "<主题>" --start "<ISO>" --end "<ISO>" [--with "张三,李四"]` | 姓名必须唯一解析;写后读回;不含会议室预订 |
| 查看最近一场 | `dws calendar +next-event` | 默认未来 7 天;不要用它代替完整列表 |
| 查某人 / 会议室闲忙 | 姓名用 `+free`;ID 用 `+freebusy` | 必须有明确时段;至少指定 users/rooms 一类 |
| 推荐多人共同时间 | `dws calendar +suggest-time --with "张三,李四" --start "<ISO>" --end "<ISO>"` | 只推荐,不创建日程 |
| 找可用会议室 | `dws calendar +room-find --start "<ISO>" --end "<ISO>"` | 名称定位但不查可用性时才用 `+room-search` |
| 邀请参会人 | `dws calendar +invite --event <EVENT_ID> --with "张三,李四"` | 只修改指定已有日程 |
| 改期 | `dws calendar +reschedule --event <EVENT_ID> --start "<ISO>" --end "<ISO>"` | 只改起止时间;其他字段保持不变 |
| 取消日程 | `dws calendar +cancel-event --event <EVENT_ID>` | 高风险删除;先读目标,确认后执行并验证不存在 |
## 标准 SOP(必遵流程)
当用户需要 shortcut 未公开的字段、共享日历、循环规则、附件、ACL 或会议室绑定时,才降级到 [calendar.md](references/calendar.md) 的单个原子 leaf。
> 命中以下意图**必须**按对应 SOP 顺序执行;**禁止**跳步、替换命令、编造 userId/eventId。每条命令必须带 `--format json`,时间参数**必须**是 ISO-8601(如 `2026-07-03T14:00:00+08:00`)。
## 资源与安全约束
### SOP-1 查日程(list-events)
- 时间必须是带时区的 ISO-8601,且 `end > start`。预约意图缺少起止时间时先追问;不要自设全天窗口。`+today/+tomorrow/+week` 和自身声明默认窗口的只读 shortcut 除外。
- 多轮任务持续复用同一 `eventId`;更新、邀请、订房和取消不得通过再次创建来替代。
- 用户点名日程但未给 `eventId` 时,用标题和合理时间窗检索;零命中停止,多候选列出标题、时间和 `eventId` 让用户选择,禁止选第一条。
- 用户姓名必须唯一解析到当前 profile 的身份;零命中、多候选或跨 profile 不复用 ID。
- `roomId` 只能来自同一 profile、同一目标时段下的 `+room-find` / `room search` 返回。会议室名、楼层编号和地点文本都不是 `roomId`;`--location` 也不等于预订。
- 用户指定时间或会议室时不得擅自换时间、换房或扩大地点范围。允许范围无空房时停止并说明,若要继续必须让用户明确放宽条件。
- 分页必须跟随 `nextCursor` 或 page 语义,检测 cursor 丢失、停滞和循环;不得把第一页当全量。
- 非事务多步写入发生 partial、pending 或 commit-unknown 时如实报告已完成与未完成步骤,先读回协调,禁止盲目重试非幂等创建。
**触发**:今天/明天/本周日程/我有什么会/某时段日程。
## 写后验证
1. **首选脚本(必须)**:`python scripts/calendar_today_agenda.py today|tomorrow|week`(聚合今日议程)。
2. **降级 CLI(必须)**:脚本不可用时 `dws calendar event list --start "<起始ISO>" --end "<结束ISO>" --format json`;不传 `--start/--end` 默认查今天(00:00:00~23:59:59)。`hasMore=true` 用 `--limit`/翻页。
3. **解析(必须)**:取真实 `eventId`、`attendees[]`、`start/end`;按需抽取,**禁止**把整段 JSON 原样贴出。
- 创建:以结构化结果中的 `eventId` 为身份,读回核对标题、起止时间和预期参会人。
- 邀请 / 移除参会人:读取参会人列表核对目标;底层不提供稳定 userId 时不得伪造身份字段。
- 改期 / 更新:读回同一 `eventId`,核对实际变更字段并确认未意外覆盖其他字段。
- 订房 / 换房:读回同一日程,并用对应时段的会议室闲忙或日程详情确认绑定;仅有空响应不能证明成功。
- 取消:读回明确为不存在;权限错误、超时或缺少不存在证据时不得报告成功。
**禁止**:用 `event list` 替代闲忙查询(查闲忙走 SOP-3)、编造时间窗口、用非 ISO 时间格式。
## 产品边界
### SOP-2 建日程(create-event)
- 视频会议发起、入会链接、会中控制:当前 Calendar CLI 不支持,不能臆造 conference 命令。
- AI 听记、会后摘要与转写:转 `dingtalk-minutes`;待办任务与独立截止提醒:转 `dingtalk-todo`。
- 按姓名解析人员由 Calendar 的 `+book/+invite/+free/+suggest-time` 优先完成;只有原子命令路径才先用 `dingtalk-aisearch`。
- 日历事件是 Calendar 资源;“给自己留时间块”不是 Todo。
**触发**:建日程/约会议/加日程。
1. **解析与会人(必须)**:对每个姓名 `dws aisearch person --query "<姓名>" --dimension name --format json` 取 `userId`,多人逗号拼接。
2. **执行(必须)**:`dws calendar event create --title "<主题>" --start "<ISO>" --end "<ISO>" --attendees <userId1,userId2> --format json`(按需加 `--location`/`--desc`/`--rooms`)。
3. **验证(必须)**:从返回 `result.id` 取日程 ID(下游参数语义称 `eventId`),再执行 `dws calendar event list --start "<ISO>" --end "<ISO>" --format json` 复核标题、描述和时段。
**禁止**:跳过与会人 userId 解析直接传姓名、编造会议室 roomId。
### SOP-3 查闲忙(check-busy)
**触发**:某人/会议室是否有空/找空闲时段/避免冲突。
1. **解析对象(必须)**:姓名 → `dws aisearch person --query "<姓名>" --dimension name --format json` 取 `userId`;会议室用 `roomId`。
2. **收敛时段(必须)**:`--start`/`--end` **必须**由用户给出或明确收敛;时段不明确**必须先追问**,**禁止**默认全天窗口。
3. **执行(必须)**:`dws calendar busy search --users <userId1,userId2> --start "<ISO>" --end "<ISO>" --format json`(查会议室换 `--rooms <roomId...>`,可同时传)。**禁止**用 `event list` 扫日程替代闲忙查询。
4. **空闲时段(必须)**:找共同空闲用 `python scripts/calendar_free_slot_finder.py`。
**禁止**:用 `event list` 冒充 `busy search`、未确认时段就默认全天查询。
## 执行硬约束
- 多轮日程任务必须保留 `eventId`,后续加人、移人、订房、换房、改描述、删除都基于同一个 `eventId` 执行;不要重新创建重复日程。
- 用户明确说"帮我订一个空闲会议室"时,`room search` 返回可用会议室后直接选择第一个可预订且不需要自定义审批的 `roomId` 执行 `room add`;不要把选择权抛回用户导致任务停住。
- 已有日程订房:`dws calendar room search --start ... --end ... --format json` → `dws calendar room add --event <EVENT_ID> --rooms <ROOM_ID> --format json` → `event get` 或 `room/busy` 验证。
- 换会议室:先 `room delete --event <EVENT_ID> --rooms <OLD_ROOM_ID>`,再 `room add --event <EVENT_ID> --rooms <NEW_ROOM_ID>`,最后回查;不要只更新 `--location`。
- 参会人变化用 `attendee add/delete`,日程描述变化用 `event update --desc`,删除日程用 `event delete --id`。用户当前消息已明确要求删除/取消时可直接执行;否则先确认。
- 脚本失败或参数不完整时,立即降级到明确的 `dws calendar event/attendee/room` 命令,不要停在"我要查看用法"。
- 所有 dws 命令带 `--format json`;查询时间必须显式 `--start` / `--end`。
## 跨产品协作
- 视频会议发起 / 入会链接 / 邀请入会 / 会中控制 → 当前 CLI **不支持**;请在钉钉客户端完成
- 会后摘要 / 待办 → 切到 `dingtalk-minutes`
- 参会人按人名 → 先用 `dingtalk-aisearch` 解析
## 注意
`schedule-meeting` 必须读 [03-meeting.md](references/03-meeting.md) 中的「两准则」「搜房失败硬门禁」,禁止假设 `roomId`。
## 局部意图与短流程
- [局部意图消歧](references/intent-guide.md);[短流程](references/lite-recipes.md)。
涉及“创建日程并订会议室”的组合流程时读取 [03-meeting.md](references/03-meeting.md),严格执行时段、地点、`roomId` 来源与失败收束规则。
@@ -14,11 +14,11 @@
**`room search --available`**(与传入的 `--start` / `--end` 配对):返回的是在**该整段时段内**可被预订的空闲会议室(不是「有一段空就算」);脚本与用户手工选房均应沿用同一时间窗,避免误以为分段凑满即等价于整段可用。
**`dws calendar room search` 合法参数**(与 [calendar.md](./calendar.md) 一致):仅 `--start`、`--end`、`--group-id`(可选)、`--available`(可选)、`--format json` 等;**禁止使用 `--query`**,否则会报 `unknown flag: --query`。
**会议室名称参数**:shortcut 使用 `dws calendar +room-find --room-name "<核心专名>"`;原子命令使用 `dws calendar room search --room-name "<核心专名>"`。两者都不支持 `--query`。只知道名称、无需检查时段可用性时用 `+room-search --room-name`。
**地点归组早停**:若用户给的是同一地点范围(如“西溪园区 C6 楼 3-5 层”或具体楼层/楼栋),先用 `room list-groups` 找到**最相关的承载 group**(通常是该楼层;若楼层下无会议室则为直接挂会议室的上一级)。在这个最相关 group 下查不到有效 `rooms[].roomId` 或空房时,**不得**再跳去别的同级/异地 group 继续搜;同一地点的会议室不会散落在别的 group 里。只有用户明确放宽到别的楼层、楼栋或园区,才能重新解析新的 group 并继续。
**用户点名具体会议室(如「C6-4-06-N / 贡嘎山」)**:**不要**尝试 `room search --query "<名称>"`;**禁止**把用户原文(含「C6-4-06-N 贡嘎山」整句)或展示名当作 `room add --rooms` 的 `roomId`。用户输入**几乎从不会是**有效 `roomId`。须先 `dws calendar room list-groups` 定位所在楼层/分组的 `group-id`,再 `dws calendar room search --start "<ISO>" --end "<ISO>" --group-id <GROUP_ID> [--available] --format json`,在返回 `rooms[]` 中对 `roomName`、`name` 等与用户表述匹配,**仅**取 JSON 里的 `roomId`(典型为小写十六进制串,长度以返回为准),最后 `dws calendar room add --event <eventId> --rooms <roomId>`。该时段无匹配或房间忙 → 如实告知;**禁止**为通过校验而编造、拼接或猜测 `roomId`。
**用户点名具体会议室(如「C6-4-06-N / 贡嘎山」)**:**不要**尝试 `--query`;**禁止**把用户原文或展示名当作 `room add --rooms` 的 `roomId`。用 `dws calendar +room-find --room-name "<核心专名>" --start "<ISO>" --end "<ISO>" [--group-id <GROUP_ID>] --format json`,在返回会议室中按名称匹配,**仅**取真实 `roomId`,最后 `dws calendar room add --event <eventId> --rooms <roomId>`。该时段无匹配或房间忙则如实告知;禁止编造、拼接或猜测 `roomId`。
### 搜房失败硬门禁(园区/范围搜尽仍无 roomId)
@@ -52,5 +52,5 @@
| Recipe | 行动指南(固定路线) |
| ------------------ | ------------------- |
| schedule-meeting | **见上文「两准则」**、**「搜房失败硬门禁」**。**未给时段且仅说「发起/开个会」**→ 不走本 recipe;当前 CLI 不支持实时视频会议,告知用户请在钉钉客户端操作。**未给时段但有预约意图**("安排""约""定"等词):追问具体开始/结束时间。**已有时段后**,按固定顺序执行:1. `dws calendar event create` 建日程;2. 有参会人则 `dws calendar participant add`;3. 再处理会议室。**无明确会议室范围**:可直接 `python scripts/calendar_schedule_meeting.py --title "<主题>" --start "<起始>" --end "<结束>" [--users <userIds>] [--book-room] [--dry-run]`。**有明确范围(某楼/层)**:先 `dws calendar room list-groups`,锁定该地点**最相关的承载 group**;若只有一个地点,`--room-group-id` 应只传这个最相关 group,**不要**把同楼内多个楼层 group 打包传入碰运气。只有用户明确给出多个允许地点时,才把这些 `group-id` 一并传给 `python scripts/calendar_schedule_meeting.py ... --book-room --room-group-id "<id1,id2,...>"`。**用户点名具体会议室**:须手工 `dws calendar room search --start "<ISO>" --end "<ISO>" --group-id <GROUP_ID> [--available] --format json`(**无** `--query`),在 JSON 中匹配名称取 **`rooms[].roomId` 唯一真值** → `dws calendar room add --event <eventId> --rooms <roomId>`;**不得**把用户输入的会议室名当 `roomId`。**一旦连续 2 次空结果 / 任意一次 `roomId invalid`**:**必须回读本节并立即收束判断**;若整园/限定范围内搜尽仍无 roomId 或无空房 → **下一条消息必须直接向用户汇报失败结论**;否则只能向用户确认是否放宽范围/改时间。**禁止**假设 roomId、禁止无 ID 调用 `room add`、禁止用日程详情绕路、禁止继续猜测 Mock/测试环境。细则见「会议室搜索早停」。 |
| schedule-meeting | **见上文「两准则」**、**「搜房失败硬门禁」**。**未给时段且仅说「发起/开个会」**→ 当前 CLI 不支持实时视频会议,告知用户请在钉钉客户端操作。**未给时段但有预约意图**:追问具体开始/结束时间。**已有时段后**:1. 先用 `+room-find` 在用户允许范围内取得真实 `roomId`(用户没要求会议室则跳过);2. 用 `+book` 创建日程并按姓名邀请参会人;3. 有 `roomId` 时用 `room add --event <eventId> --rooms <roomId>` 绑定;4. 读回同一 `eventId` 验证。用户点名会议室时给 `+room-find` 传 `--room-name`,不得使用 `--query`。连续 2 次空结果或任意一次 `roomId invalid` 时立即收束;若限定范围内无 roomId 或无空房,直接汇报失败,除非用户明确放宽范围或改时间。 |
| reschedule-meeting | 1. `calendar event list --start "<起始ISO>" --end "<结束ISO>"` → 取 `eventId` 2. `calendar event update --id <eventId> --start "<新起始ISO>" --end "<新结束ISO>"` 更新时间 3. `chat search --query "<群名>"` → 取 `openConversationId` → `chat message send --conversation-id <openConversationId> --content "<变更通知>"` 通知变更 |
@@ -6,20 +6,20 @@
### list-today-meetings
**优先**:`python scripts/calendar_today_agenda.py [today|tomorrow|week]`
备选:`dws calendar event list --start "<今日起始ISO>" --end "<今日结束ISO>"`(须加 `--format json`)
**优先**:`dws calendar +today|+tomorrow|+week --format json`
任意时段:`dws calendar +agenda --start "<起始ISO>" --end "<结束ISO>" --format json`
### check-users-busy
查询多人在某时段内的闲忙(**busy**,不是用 `event list` 扫日程):
1. 解析用户:对每个姓名执行 `aisearch person --query "<姓名>" --dimension name` → `userId`;多人将 `userId` 用英文逗号拼接(无空格或按 [calendar.md](./calendar.md) `busy search` 要求)。
2. 确认时段:用户须给出或可收敛为明确的 `--start` / `--end`(ISO-8601);若未给出,**先追问**起止时间,禁止用任意默认全天窗口代替用户意图。
3. 执行:`dws calendar busy search --users <userId1,userId2,...> --start "<ISO>" --end "<ISO>" --format json`
1. 确认时段:用户须给出或可收敛为明确的 `--start` / `--end`(ISO-8601);若未给出,先追问起止时间。
2. 按姓名查一人:`dws calendar +free --who "<姓名>" --start "<ISO>" --end "<ISO>" --format json`。
3. 多人共同时间:`dws calendar +suggest-time --with "张三,李四" --start "<ISO>" --end "<ISO>" --format json`。
4. 已有 userId/roomId:`dws calendar +freebusy --users <userIds> --rooms <roomIds> --start "<ISO>" --end "<ISO>" --format json`;users/rooms 至少一类。
详见 [calendar.md](./calendar.md) 中「查询用户闲忙状态」。
### start-conference
> 当前 CLI 不提供视频会议(conference)发起/入会/会中控制能力。触发「发起会议」「开个会」「创建会议」且**没有给出具体时间**时,不要构造 `conference` 命令;直接告知用户请在钉钉客户端操作。
+5 -5
View File
@@ -1,6 +1,6 @@
---
name: dingtalk-event
description: 钉钉个人 IM 与 OA 审批事件长连接监听。Use when 用户说监听消息/@我/某人/某群/全部消息、已读/撤回/reaction、群成员加入/群成员退出/群状态变化,或监听审批任务创建/完成/转交、审批实例发起/抄送/终止/完成。命令前缀:dws event。
description: 钉钉个人 IM 与 OA 审批事件长连接监听。Use when 用户说监听消息/@我/某人/某群/全部消息、已读/撤回/reaction、群成员加入/群成员退出/群状态变化,或监听审批任务创建/完成/转交、审批实例发起/终止/完成。命令前缀:dws event。
metadata:
cli_version: ">=0.2.14"
category: product
@@ -45,7 +45,7 @@ metadata:
- `group` 必须且只能传 `--chat-id` 或 `--chat-query` 之一。
- `--query` 只用于纯 `message` 监听;混入 reaction/read/recall 时不得使用。
OA 事件不进入 `+listen-im`。七个公开 OA EventKey 都订阅当前 OAuth 用户相关的全部审批事件,使用 `ruleType=all`、`filterRule={}`;不接受 `--user`、`--open-dingtalk-id`、`--group`、`--query` 或 `--filter-json`。七项可放入同一个 consume,每项建立独立订阅并共享 bus。
OA 事件不进入 `+listen-im`。六个公开 OA EventKey 都订阅当前 OAuth 用户相关的全部审批事件,使用 `ruleType=all`、`filterRule={}`;不接受 `--user`、`--open-dingtalk-id`、`--group`、`--query` 或 `--filter-json`。六项可放入同一个 consume,每项建立独立订阅并共享 bus。
自然姓名和群名由 CLI 内部唯一解析:零命中或多候选返回结构化失败,在创建任何订阅前停止。`--dry-run` 走同一解析链。解析、监听、状态和停止必须使用同一个 `--profile`,不得跨组织搬运 ID。
@@ -65,7 +65,7 @@ user_im_group_member_added user_im_group_member_exited
user_im_group_disbanded
```
七个 OA EventKey 及其输出字段见 [OA 事件参考](references/event-oa.md)。
六个 OA EventKey 及其输出字段见 [OA 事件参考](references/event-oa.md)。
用户类事件传 `--user` 或 `--open-dingtalk-id`,群类事件传 `--group`。群生命周期输出可含 `operator_open_dingtalk_id` 和 `members`;成员项使用 `open_dingtalk_id`。精确组合、兼容性和 Filter 规则见 reference。
@@ -98,7 +98,7 @@ kind + events + target
- `event stop` 会取消订阅并影响本地 consumer:先 `--dry-run`,用户确认后再加 `--yes`。
- 多事件属于一次原始操作;任一订阅启动失败时 Runtime 回滚本次已创建项,不拆成新命令绕过重试预算。
- 这套 `0/2/1` 是 **Agent/host** 编排预算,适用于全部 23 个公开个人 EventKey(16 个 IM + 7 个 OA):`retryable=false` 对应 `max_additional_attempts=0`;`retryable=true` 对应 `max_additional_attempts=2`;`retryable=unknown` 对应 `max_additional_attempts=1`。它不是 CLI 持久化硬总次数上限;每次调用最多创建一次,进程内不会自动重试,CLI 也不持久化或计算跨调用的 Agent/host 尝试次数。
- 这套 `0/2/1` 是 **Agent/host** 编排预算,适用于全部 22 个公开个人 EventKey(16 个 IM + 6 个 OA):`retryable=false` 对应 `max_additional_attempts=0`;`retryable=true` 对应 `max_additional_attempts=2`;`retryable=unknown` 对应 `max_additional_attempts=1`。它不是 CLI 持久化硬总次数上限;每次调用最多创建一次,进程内不会自动重试,CLI 也不持久化或计算跨调用的 Agent/host 尝试次数。
- 重试必须遵守 `retry_after_seconds` / `next_retry_at`。遇到 `in_flight`、`cooldown`、`terminal_hold` 不并发或递归重启同一逻辑订阅,也不换 `subscribe_id` / `trace_id` 绕过保护。
- 认证、profile、订阅保护状态和 bus 排障按失败类型读取 [订阅运维](references/event-im-operations.md),不要在正常路径预加载完整运维手册。
@@ -125,4 +125,4 @@ kind + events + target
| ready、bounded consume 与退出清理 | [event-im-lifecycle.md](references/event-im-lifecycle.md) | 启动/托管/关闭 consumer |
| 扁平字段与事件到 Chat 交接 | [event-im-output.md](references/event-im-output.md) | 解析事件或自动回复 |
| Filter、status/stop、重试与排障 | [event-im-operations.md](references/event-im-operations.md) | 订阅控制或失败恢复 |
| OA 审批事件 | [event-oa.md](references/event-oa.md) | 选择七个 OA EventKey、组合消费或解析审批字段 |
| OA 审批事件 | [event-oa.md](references/event-oa.md) | 选择六个 OA EventKey、组合消费或解析审批字段 |
@@ -1,6 +1,6 @@
# OA 个人审批事件
先读事件产品入口 [SKILL.md](../SKILL.md) 的命令规则、调用流和子进程契约。本参考覆盖当前公开的七个 OA 个人事件:审批实例发起、抄送、终止和完成,以及审批任务创建、完成和转交。
先读事件产品入口 [SKILL.md](../SKILL.md) 的命令规则、调用流和子进程契约。本参考覆盖当前公开的六个 OA 个人事件:审批实例发起、终止和完成,以及审批任务创建、完成和转交。
<!-- dws-intent: event.listen.oa -->实时监听审批事件必须使用 `dws event consume` 长连接,不要轮询 OA 待办或审批实例列表来模拟事件。
@@ -22,11 +22,10 @@ dws auth login
| `user_oa_approval_task_finished` | `all` | 审批任务已完成 | 无 |
| `user_oa_approval_task_redirected` | `all` | 审批任务已转交 | 无 |
| `user_oa_approval_instance_started` | `all` | 审批实例已发起 | 无 |
| `user_oa_approval_instance_cc` | `all` | 审批实例到达抄送节点,发送给被抄送人 | 无 |
| `user_oa_approval_instance_terminated` | `all` | 审批实例已终止 | 无 |
| `user_oa_approval_instance_finished` | `all` | 审批实例完成,发送给审批单发起人 | 无 |
只承认上表 7 个 OA 事件码。CLI 为每个事件发送 `ruleType=all`、`filterRule={}` 的独立订阅请求;不要添加 `--user`、`--open-dingtalk-id`、`--group`、`--query` 或 `--filter-json`。
只承认上表 6 个 OA 事件码。CLI 为每个事件发送 `ruleType=all`、`filterRule={}` 的独立订阅请求;不要添加 `--user`、`--open-dingtalk-id`、`--group`、`--query` 或 `--filter-json`。
## Intent mapping
@@ -36,14 +35,13 @@ dws auth login
| “审批任务完成时通知我” | `dws event consume user_oa_approval_task_finished --flatten -f ndjson` |
| “审批任务被转交时通知我” | `dws event consume user_oa_approval_task_redirected --flatten -f ndjson` |
| “有审批单发起时通知我” | `dws event consume user_oa_approval_instance_started --flatten -f ndjson` |
| “有审批抄送给我时通知我” | `dws event consume user_oa_approval_instance_cc --flatten -f ndjson` |
| “有审批单终止时通知我” | `dws event consume user_oa_approval_instance_terminated --flatten -f ndjson` |
| “监听我发起的审批何时完成” / “审批实例完成时通知我” | `dws event consume user_oa_approval_instance_finished --flatten -f ndjson` |
| “同时监听全部已公开 OA 事件” | 一个 consume 放入七个 OA event key,不加目标或过滤参数 |
| “同时监听全部已公开 OA 事件” | 一个 consume 放入六个 OA event key,不加目标或过滤参数 |
| “查看 OA 事件目录” | `dws event list --category oa` |
| “查看 OA 事件输出字段” | 对对应事件运行 `dws event schema <event_key> --flatten` |
三个审批任务事件分别表达任务已创建、已完成和已转交;四个审批实例事件分别表达实例已发起、到达抄送节点、已终止和已完成。`status` 和 `result` 保留服务端原值,不把当前样本值推断为完整枚举。
三个审批任务事件分别表达任务已创建、已完成和已转交;三个审批实例事件分别表达实例已发起、已终止和已完成。扁平字段来自六类事件的预发联调样本;`status` 和 `result` 保留服务端原值,不把当前样本值推断为完整枚举。
## Commands
@@ -54,7 +52,6 @@ dws event schema user_oa_approval_task_created --flatten
dws event schema user_oa_approval_task_finished --flatten
dws event schema user_oa_approval_task_redirected --flatten
dws event schema user_oa_approval_instance_started --flatten
dws event schema user_oa_approval_instance_cc --flatten
dws event schema user_oa_approval_instance_terminated --flatten
dws event schema user_oa_approval_instance_finished --flatten
```
@@ -66,12 +63,11 @@ dws event consume user_oa_approval_task_created --flatten -f ndjson
dws event consume user_oa_approval_task_finished --flatten -f ndjson
dws event consume user_oa_approval_task_redirected --flatten -f ndjson
dws event consume user_oa_approval_instance_started --flatten -f ndjson
dws event consume user_oa_approval_instance_cc --flatten -f ndjson
dws event consume user_oa_approval_instance_terminated --flatten -f ndjson
dws event consume user_oa_approval_instance_finished --flatten -f ndjson
```
同时监听七种事件:
同时监听六种事件:
```bash
dws event consume \
@@ -79,14 +75,13 @@ dws event consume \
user_oa_approval_task_finished \
user_oa_approval_task_redirected \
user_oa_approval_instance_started \
user_oa_approval_instance_cc \
user_oa_approval_instance_terminated \
user_oa_approval_instance_finished \
--flatten \
-f ndjson
```
多事件 consume 会为七个 event key 分别创建订阅和逻辑 consumer,并共享当前组织的 personal bus、远程连接、stdout 和生命周期。不要给 OA 命令加 `--query` 或 `--filter-json`;这两个 flag 只用于兼容的 IM 消息接收事件。
多事件 consume 会为六个 event key 分别创建订阅和逻辑 consumer,并共享当前组织的 personal bus、远程连接、stdout 和生命周期。不要给 OA 命令加 `--query` 或 `--filter-json`;这两个 flag 只用于兼容的 IM 消息接收事件。
## Output contract
@@ -110,7 +105,7 @@ dws event consume \
- `type` 是当前 event key;`event_id` 可用于去重;`timestamp` 是 transport 事件发生时间;`subscribe_id` 标识对应的独立订阅。
- `process_instance_id` 是审批实例 ID,可传给 OA 审批命令的 `--instance-id`;`process_code` 是审批流程模板编码。
- `create_time`、`finish_time` 和 `event_time` 都是毫秒时间戳。`event_time` 是审批业务事件时间,`timestamp` 是 transport 事件时间。
- 七类事件的额外字段如下;具体事件始终以 `dws event schema <event_key> --flatten` 为准。
- 六类事件的额外字段如下;具体事件始终以 `dws event schema <event_key> --flatten` 为准。
| 事件 | 额外顶层字段 |
|---|---|
@@ -118,7 +113,6 @@ dws event consume \
| `user_oa_approval_task_finished` | `task_id`、`result`、`finish_time` |
| `user_oa_approval_task_redirected` | `task_id`、`result`、`finish_time` |
| `user_oa_approval_instance_started` | 无 |
| `user_oa_approval_instance_cc` | 无 |
| `user_oa_approval_instance_terminated` | `finish_time` |
| `user_oa_approval_instance_finished` | `result`、`finish_time` |
@@ -150,6 +144,6 @@ dws event consume \
## Lifecycle
- 单事件等待 `[event] ready event_key=<key> bus_pid=<pid> subscribe_id=<id>`。
- 七事件先保存七条 `[event] subscription event_key=<key> subscribe_id=<id>`,再等待 `[event] ready event_count=7 bus_pid=<pid>`。
- 六事件先保存六条 `[event] subscription event_key=<key> subscribe_id=<id>`,再等待 `[event] ready event_count=6 bus_pid=<pid>`。
- 临时验证使用 `--max-events 1` 或 `--duration 10m`;任务完成后优雅结束 consume,本次新建的订阅会自动取消。
- 外部停止已有订阅时先运行 `dws event stop <subscribe_id> --dry-run`,确认后再加 `--yes`。不要 `kill -9`,否则会跳过自动退订。
-4
View File
@@ -1190,10 +1190,6 @@ Agent 安装 dws skill 后,仅依据 skill 提供的参考文档,将自然
- Prompt: 有和我相关的审批实例发起时实时通知我
- Expected: `dws event consume user_oa_approval_instance_started --flatten -f ndjson`
**event_event_consume_oa_instance_cc_001**
- Prompt: 有审批实例抄送给我时实时通知我
- Expected: `dws event consume user_oa_approval_instance_cc --flatten -f ndjson`
**event_event_consume_oa_instance_terminated_001**
- Prompt: 和我相关的审批实例终止时实时通知我
- Expected: `dws event consume user_oa_approval_instance_terminated --flatten -f ndjson`
+2 -3
View File
@@ -280,7 +280,7 @@ func TestEventSkillFrontmatterAdvertisesGroupMemberLifecycle(t *testing.T) {
"群成员加入",
"群成员退出",
"审批任务创建/完成/转交",
"审批实例发起/抄送/终止/完成",
"审批实例发起/终止/完成",
} {
if !strings.Contains(frontmatter, required) {
t.Errorf("%s frontmatter missing event discovery trigger %q", path, required)
@@ -306,7 +306,7 @@ func TestStandaloneEventSkillOwnsAllPersonalEventContracts(t *testing.T) {
"<!-- dws-intent: event.listen.im -->",
"<!-- dws-intent: event.listen.oa -->",
"16 个 EventKey",
"23 个公开个人 EventKey",
"22 个公开个人 EventKey",
} {
if !strings.Contains(string(skillContent), required) {
t.Errorf("%s missing standalone event contract %q", skillPath, required)
@@ -360,7 +360,6 @@ func TestStandaloneEventSkillOwnsAllPersonalEventContracts(t *testing.T) {
"user_oa_approval_task_finished",
"user_oa_approval_task_redirected",
"user_oa_approval_instance_started",
"user_oa_approval_instance_cc",
"user_oa_approval_instance_terminated",
"user_oa_approval_instance_finished",
}