Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
4987aae380 | ||
|
|
c8c997fff3 |
@@ -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.
|
||||
@@ -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/
|
||||
|
||||
@@ -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
@@ -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 一致 |
|
||||
|
||||
@@ -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 不提供的高级多事件控制",
|
||||
},
|
||||
|
||||
@@ -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] {
|
||||
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
@@ -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 := ¶mAliasCaptureCaller{}
|
||||
_, 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 := ¶mAliasCaptureCaller{}
|
||||
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 := ¶mAliasCaptureCaller{}
|
||||
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 := ¶mAliasCaptureCaller{}
|
||||
_, 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 := ¶mAliasCaptureCaller{}
|
||||
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 {
|
||||
|
||||
@@ -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
@@ -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,
|
||||
}
|
||||
|
||||
@@ -81,7 +81,6 @@ var schemaCatalogToolOptionalKeys = []string{
|
||||
"pagination",
|
||||
"positionals",
|
||||
"result",
|
||||
"wait",
|
||||
}
|
||||
|
||||
var schemaCatalogToolEnums = map[string][]string{
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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{
|
||||
|
||||
@@ -30,7 +30,6 @@ type ContractFinalPayload struct {
|
||||
Parameters []ParamDecl
|
||||
Safety *SafetySpec
|
||||
DryRun *DryRunSpec
|
||||
Wait *WaitSpec
|
||||
Result *ResultSpec
|
||||
Pagination *PaginationSpec
|
||||
Interface *InterfaceSpec
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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) })
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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"}},
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
} {
|
||||
|
||||
@@ -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: "审批单终止",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 © }
|
||||
|
||||
type cloneNode struct {
|
||||
|
||||
@@ -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 ©
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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,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"
|
||||
---
|
||||
|
||||
|
||||
@@ -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 的订阅。
|
||||
|
||||
@@ -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` 命令;直接告知用户请在钉钉客户端操作。
|
||||
|
||||
|
||||
@@ -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`,否则会跳过自动退订。
|
||||
|
||||
@@ -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`
|
||||
|
||||
@@ -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",
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user