Compare commits

...
Author SHA1 Message Date
玉澜 8802dcf4dc fix(corecmd): support strongly-typed DTOs in event/auto wait modes
waitResource previously asserted result.Data() to map[string]any, which
caused event mode and auto mode to fail with 'data is not an object'
when commands returned strongly-typed DTOs (struct or struct pointer).

Now normalizes via JSON round-trip for non-map types, preserving the
fast path for map[string]any. Added regression tests covering struct
values, struct pointers, nested dotted queries, and error cases.
2026-08-19 15:14:22 +08:00
玉澜 28bb377e66 fix: remove stale allowlist entry for drive_tree_list.py
The file only exists in skills/mono/scripts/, not in skills/multi/,
so the allowlist entry was causing TestMonoMultiSkillContentG4Drift
to fail. The script is already referenced in both mono and multi
documentation, so no allowlist entry is needed.
2026-08-19 15:02:54 +08:00
玉澜 a0fe8d70c2 Merge remote-tracking branch 'upstream/main' into feat/wait-framework 2026-08-19 14:29:57 +08:00
玉澜 4770e5a8e6 Merge remote-tracking branch 'upstream/main' into feat/wait-framework 2026-08-19 14:22:29 +08:00
玉澜 0bcc2f27c6 fix(corecmd): sync operation.state on terminal wait close and canonicalize wait statuses
Terminal success/failure closes kept the acceptance-phase operation state,
emitting self-contradicting envelopes (outcome=success with state=processing).
WithOperationTerminalState now closes the envelope at the observed terminal
status and clears timed_out.

WaitSpec.Validate trimmed status values only for checking and never wrote them
back, so padded declarations passed validation, published a padded Schema, and
then fail-closed at runtime as unknown statuses. NormalizeWaitSpec is now the
single canonical form shared by declaration (corecmd.New / AttachContract),
ToolSpec, and validation: trimmed values, duplicate/conflict rejection,
defensive copy.

Also covers the waitTimeoutDuration secs<=0 branch flagged by the coverage
gate.
2026-08-19 14:17:11 +08:00
john 0b3abfad4b Merge branch 'main' into feat/wait-framework 2026-08-19 11:21:16 +08:00
玉澜 e625da4c27 fix(corecmd): reject overflowing --wait-timeout seconds
int(secs)*time.Second can wrap a pflag-legal MaxInt64 into a non-positive
duration, which skipped the wait deadline and waited forever. Convert with
an overflow check and return a validation error instead.
2026-08-18 16:59:28 +08:00
玉澜 c68603fea0 fix(corecmd): wait only on pending and honor wait-timeout on leaf I/O
Only a pending ResultInvoke envelope enters the wait phase, so success,
failure, and partial results are returned unchanged. WaitPoll/WaitEvents
now receive the --wait-timeout deadline (and Command().Context() is bound
to the same loop context) so a blocked poll or subscribe cannot hang past
the declared timeout.

Also allowlist the leftover multi drive_tree_list.py orphan that broke CI
after merging main.
2026-08-18 16:32:08 +08:00
john f118a369b0 Merge branch 'main' into feat/wait-framework 2026-08-18 15:34:12 +08:00
玉澜 b7aa6bddf5 feat(corecmd): implement event and auto wait modes
Completes the wait capability per review guidance ("add corresponding
execution hooks and fallback tests for the modes"):

- contract.WaitSpec restores event/auto modes with event_key,
  match_field, and a new resource_query (dotted path into the accepted
  result data yielding the identifier events correlate against);
  per-mode validation of required fields
- internal/wait adds EventStream (leaf-owned transport) and RunEvent:
  correlated-event filtering, the same terminal/pending/unknown mapping
  as polling, timed-out pending on deadline during consumption, and an
  ErrEventStreamEnded sentinel distinguishing stream termination from
  fail-closed status errors
- Spec.WaitEvents hook; validateWaitDecl pairs mode with hooks
  (poll<->WaitPoll, event<->WaitEvents, auto<->both; surplus hooks
  rejected too)
- the wait phase runs event-first in auto mode and falls back to polling
  when the stream ends or the subscription fails, under one deadline
  spanning both phases; strict event mode surfaces stream errors
- output.CommandResult gains Data() (deep copy) so the framework can
  resolve the resource identifier without exposing mutable state

Changed-code coverage re-verified at 100% (CI cross-package recipe).
2026-08-15 19:29:22 +08:00
玉澜 64ad5f22b0 test(cli): cover Wait capability projection (validate/normalize/payload) 2026-08-15 18:53:48 +08:00
玉澜 c68207ad4b fix(corecmd): CR feedback — poll-only wait, deadline-safe loop, ResultInvoke pairing
Addresses the three P1 findings from review 4942891040:

1. event/auto modes were declared but always executed polls. WaitSpec now
   accepts poll only (event/auto fail validation with a not-implemented
   message); event_key/match_field dead fields removed. Event waiting will
   land with its own execution path and mode constant.
2. a deadline reached during the between-poll sleep re-polled with a
   cancelled context, so a context-aware poller surfaced its error as a poll
   failure instead of the contracted timed-out pending. The wait between
   polls now uses a timer + select on ctx.Done(), the deadline is checked
   before each poll, and a poll error on a cancelled context closes as
   timed-out pending with the last observed status.
3. legacy Invoke/Orchestrate/RunE commands declaring Wait observed a failure
   terminal while still exiting 0. validateWaitDecl now requires the
   ResultInvoke dispatcher (the only path whose unified envelope can be
   closed); the wait phase no longer wraps legacy paths.

Also: dropped the unreachable nonPendingTerminal branch, simplified
waitTimeoutSecs to the flag value (registration always seeds the reviewed
default), and raised changed-code coverage to 100% (new wait-engine edge
tests, output With* unit tests, contractfinal deep-copy coverage,
AttachContract invalid-Wait panic path).
2026-08-15 18:16:24 +08:00
玉澜 4f57967c56 Merge remote-tracking branch 'upstream/main' into feat/wait-framework
# Conflicts:
#	internal/corecmd/corecmd.go
2026-08-15 17:59:10 +08:00
玉澜 b49bc0ed14 feat(corecmd): add declarative Wait capability (Contract.Wait + wait phase)
Framework-only: adds the reviewed wait contract mirroring the DryRunSpec
pattern (types declaration -> ContractDecl -> ContractFinal -> ToolSpec ->
Schema wait key). No business command declares it yet.

- contract.WaitSpec (mode poll/event/auto, poll_command, status_query,
  terminal status->success/failure map, pending_values, event_key,
  match_field, default_timeout_secs) with closed-set Validate
- Spec.WaitPoll hook pairs with the declaration at construction time
  (declared without hook / hook without declaration both panic)
- declared leaves register --wait / --wait-timeout natively (never
  FlagSpec, so they cannot enter MCP toolArgs); undeclared leaves reject
  the flags as unknown instead of ignoring them
- internal/wait engine: immediate-first-poll, x1.5 backoff capped 30s,
  dotted status extraction, fail-closed on unknown status, timeout ->
  pending
- ResultInvoke path closes the unified envelope: success terminal ->
  success, failure terminal -> failure with new wire-stable
  error.type "wait" (exit code 8, additive like partial=7), timeout ->
  pending + meta.operation.timed_out with last observed state (exit 0)
- output.WithOutcome / WithErrorInfo / WithOperationTimedOut preserve
  envelope invariants (I2/I3, pending requires meta.operation)

Verified: go test ./... green; check-generated-drift.sh ok (schema
assembly deterministic, wire unchanged); check-schema-catalog.sh ok
(27 products, 1121 tools).
2026-08-15 12:41:46 +08:00
25 changed files with 2754 additions and 3 deletions
+26
View File
@@ -0,0 +1,26 @@
---
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.
@@ -360,6 +360,7 @@ Definition(仅声明;不可编译)
| | `idempotency` | 评审源(或未来 Contract) | reviewed metadata | 今日非框架声明;不得推断 |
| | `effect_source` / provenance | 组装派生物 | resolver 写入 `FieldProvenance` | 派生,不手写 |
| **DryRun** | `preview_kind`, `remote_reads` | 评审源 | `schema_dry_run_capabilities`(正能力声明) | 否;无条目 ≠ 推断「不支持」之外的假能力 |
| **Wait** | `mode`(`poll`/`event`/`auto`), `poll_command`, `status_query`, `terminal`(状态→success/failure), `pending_values`, `event_key`/`match_field`/`resource_query`(event/auto), `default_timeout_secs` | **声明**(`ContractDecl.Wait` 正能力声明,且必须搭配 ResultInvoke dispatcher + 按模式的 hook:poll↔`WaitPoll`、event↔`WaitEvents`、auto↔两者,构造期配对校验,多余 hook 同样拒绝) | 声明后注册 `--wait`/`--wait-timeout`(框架 flag,不进 toolArgs);Schema 投影 `wait` 键;auto = 事件优先、流终止/订阅失败回退轮询,一个 deadline 覆盖两阶段并传入 `WaitPoll`/`WaitEvents`(及 `Command().Context()`);仅 pending 初始结果进入等待,success/failure/partial 原样返回 | 否;未声明命令传 `--wait` = unknown flag。终态失败经统一信封 `error.type: "wait"`(rc=8),超时保持 pending + `meta.operation.timed_out`(rc=0);轮询间/轮询中/事件消费中超时一律按 pending 关闭 |
| **Interface** | `interface_mode`, `interface_ref`, `availability`, `reason` | 评审源 | MCP meta + agent metadata 解析 | 否;与 CLI Identity 分离 |
| **Selection** | `agent_summary`, `use_when`, `avoid_when`, `examples`, `prerequisites`, `tips`, `workflow_refs`, … | 声明(`ContractDecl.Selection` / `ProductDecl`) | `ContractDecl` / `ProductDecl`(`schema_hints/` 已退役) | 可声明;声明载荷**不得携带** `Reviewed`(旧路径专用),携带即组装报错 |
| **FieldProvenance** | 各字段 winner / candidates | 组装派生物 | Schema 组装器 | 派生;须与 delivered value 一致 |
+1 -1
View File
@@ -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,
"parameters": true, "constraints": true, "positionals": true, "dry_run": true, "wait": true,
"result": true, "pagination": true,
"examples": true, "use_when": true, "avoid_when": true,
}
+1
View File
@@ -81,6 +81,7 @@ var schemaCatalogToolOptionalKeys = []string{
"pagination",
"positionals",
"result",
"wait",
}
var schemaCatalogToolEnums = map[string][]string{
+21
View File
@@ -60,6 +60,7 @@ type ToolSpec struct {
Constraints RuntimeSchemaConstraints
Positionals []contract.RuntimeSchemaPositional
DryRun *contract.DryRunSpec
Wait *contract.WaitSpec
Result *contract.ResultSpec
Pagination *contract.PaginationSpec
Safety contract.SafetySpec
@@ -135,6 +136,7 @@ type RuntimeToolSpecInput struct {
Constraints RuntimeSchemaConstraints
Positionals []contract.RuntimeSchemaPositional
DryRun *contract.DryRunSpec
Wait *contract.WaitSpec
Result *contract.ResultSpec
Pagination *contract.PaginationSpec
Safety contract.SafetySpec
@@ -542,6 +544,11 @@ 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
@@ -744,6 +751,16 @@ 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 {
@@ -982,6 +999,10 @@ 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,3 +974,57 @@ func TestFinalProvenanceCoverageDoesNotInventOptionalInterfaceReason(t *testing.
t.Fatalf("optional local interface_reason should not require invented provenance: %v", err)
}
}
func TestToolSpecWaitCapabilityIsPositiveOnly(t *testing.T) {
base := RuntimeToolSpecInput{Identity: contract.ToolIdentitySpec{
ProductID: "sample",
Name: "waitrun",
CLIName: "waitrun",
CLIPath: "sample waitrun",
}}
withoutCapability, err := ToolSpecFromRuntime(base)
if err != nil {
t.Fatalf("ToolSpecFromRuntime() error = %v", err)
}
payload, err := withoutCapability.ToPayload()
if err != nil {
t.Fatalf("ToPayload() error = %v", err)
}
if _, ok := payload["wait"]; ok {
t.Fatalf("nil capability unexpectedly emitted wait: %#v", payload["wait"])
}
base.Wait = &contract.WaitSpec{Mode: "webhook"}
if _, err := ToolSpecFromRuntime(base); err == nil || !strings.Contains(err.Error(), "unknown mode") {
t.Fatalf("invalid mode error = %v", err)
}
base.Wait = &contract.WaitSpec{Mode: contract.WaitModeEvent, StatusQuery: "status", Terminal: map[string]contract.ResultOutcome{"DONE": contract.ResultOutcomeSuccess}}
if _, err := ToolSpecFromRuntime(base); err == nil || !strings.Contains(err.Error(), "requires event_key") {
t.Fatalf("event mode body error = %v", err)
}
base.Wait = &contract.WaitSpec{
Mode: contract.WaitModePoll,
PollCommand: "sample status get",
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{"COMPLETED": contract.ResultOutcomeSuccess},
PendingValues: []string{"NEW"},
DefaultTimeoutSecs: 120,
}
withCapability, err := ToolSpecFromRuntime(base)
if err != nil {
t.Fatalf("ToolSpecFromRuntime(valid wait) error = %v", err)
}
if withCapability.Wait == nil || withCapability.Wait.Mode != contract.WaitModePoll {
t.Fatalf("wait capability lost through normalization: %#v", withCapability.Wait)
}
payload, err = withCapability.ToPayload()
if err != nil {
t.Fatalf("ToPayload(valid wait) error = %v", err)
}
wait := payload["wait"].(map[string]any)
if wait["mode"] != contract.WaitModePoll || wait["poll_command"] != "sample status get" {
t.Fatalf("wait payload=%#v", wait)
}
}
+1
View File
@@ -354,6 +354,7 @@ func runtimeToolSpecFromContractFinal(entry runtimeSchemaEntry, final contract.C
Constraints: constraints,
Positionals: positionals,
DryRun: final.DryRun,
Wait: final.Wait,
Result: result,
Pagination: pagination,
Safety: safety,
+2
View File
@@ -65,6 +65,7 @@ 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"`
@@ -269,6 +270,7 @@ func schemaToolSpecFromWire(wire schemaToolWire) (ToolSpec, error) {
Constraints: wire.Constraints,
Positionals: wire.Positionals,
DryRun: wire.DryRun,
Wait: wire.Wait,
Result: wire.Result,
Pagination: wire.Pagination,
Safety: contract.SafetySpec{
+1
View File
@@ -30,6 +30,7 @@ type ContractFinalPayload struct {
Parameters []ParamDecl
Safety *SafetySpec
DryRun *DryRunSpec
Wait *WaitSpec
Result *ResultSpec
Pagination *PaginationSpec
Interface *InterfaceSpec
+149
View File
@@ -62,6 +62,155 @@ type DryRunSpec struct {
RemoteReads bool `json:"remote_reads,omitempty"`
}
// Wait modes. Poll executes the leaf's WaitPoll hook on a cadence. Event
// consumes the leaf's WaitEvents push stream and correlates events to the
// accepted resource. Auto prefers the event stream and falls back to polling
// when the stream ends before a terminal status.
const (
WaitModePoll = "poll"
WaitModeEvent = "event"
WaitModeAuto = "auto"
)
// WaitSpec is a positive capability declaration for terminal-state waiting
// (approval flows, async exports, batch jobs). A nil ToolSpec.Wait means the
// command has not declared reviewed --wait support; the flag is not
// registered and the Schema does not publish the capability.
//
// Like DryRunSpec, the object is one atomic contract field: Schema only
// projects the reviewed capability; runtime execution stays owned by the
// command runner through the leaf's WaitPoll / WaitEvents hooks. PollCommand
// names the read command that observes status — it is a declared,
// catalog-visible fact (the same command an agent would poll manually), not
// a framework-owned invocation: how one poll or event subscription executes
// is decided by the leaf.
type WaitSpec struct {
Mode string `json:"mode"`
PollCommand string `json:"poll_command,omitempty"`
StatusQuery string `json:"status_query"`
Terminal map[string]ResultOutcome `json:"terminal"`
PendingValues []string `json:"pending_values,omitempty"`
// EventKey is the push channel key the WaitEvents hook subscribes to
// (event/auto modes). Declared for the catalog; the transport stays
// leaf-owned.
EventKey string `json:"event_key,omitempty"`
// MatchField is the event-document path holding the resource identifier
// (event/auto modes); its value must equal the ResourceQuery resolution
// of the accepted result.
MatchField string `json:"match_field,omitempty"`
// ResourceQuery is the dotted path into the accepted result data that
// yields the resource identifier correlated against MatchField
// (event/auto modes).
ResourceQuery string `json:"resource_query,omitempty"`
// DefaultTimeoutSecs is the reviewed default for --wait-timeout. Zero
// means the framework default (300s); the user flag always wins.
DefaultTimeoutSecs int `json:"default_timeout_secs"`
}
// Validate checks mode requirements and the terminal/pending status maps.
// Unknown terminal outcomes, unknown modes, and mode/body mismatches fail at
// declaration so a malformed wait capability cannot reach the wire.
// Validation delegates to NormalizeWaitSpec so the acceptance rules can never
// drift from the normalization the wire and the runtime wait engine share.
func (w WaitSpec) Validate(canonical string) error {
_, err := NormalizeWaitSpec(&w, canonical)
return err
}
// NormalizeWaitSpec returns a validated, canonical, defensively copied wait
// contract. It is shared by declaration (corecmd.New / AttachContract),
// ToolSpec, and snapshot paths, mirroring NormalizeResultSpec. Status values
// are trimmed into their wire form: the wait engine compares backend
// statuses verbatim against these tables, so a padded declaration
// (" processing ") would publish a Schema that its own runtime treats as an
// unknown status. Values collapsing onto one value after trimming (duplicate
// pending values, duplicate terminal keys, terminal/pending conflicts) are
// rejected instead of silently merged.
func NormalizeWaitSpec(in *WaitSpec, canonical string) (*WaitSpec, error) {
if in == nil {
return nil, nil
}
canonical = defaultString(strings.TrimSpace(canonical), "<unknown>")
out := &WaitSpec{
Mode: strings.TrimSpace(in.Mode),
PollCommand: strings.TrimSpace(in.PollCommand),
StatusQuery: strings.TrimSpace(in.StatusQuery),
EventKey: strings.TrimSpace(in.EventKey),
MatchField: strings.TrimSpace(in.MatchField),
ResourceQuery: strings.TrimSpace(in.ResourceQuery),
DefaultTimeoutSecs: in.DefaultTimeoutSecs,
}
if out.Mode == "" {
return nil, fmt.Errorf("schema tool %s wait has no mode", canonical)
}
switch out.Mode {
case WaitModePoll, WaitModeEvent, WaitModeAuto:
default:
return nil, fmt.Errorf("schema tool %s wait has unknown mode %q", canonical, out.Mode)
}
needsPoll := out.Mode == WaitModePoll || out.Mode == WaitModeAuto
if needsPoll && out.PollCommand == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires poll_command", canonical, out.Mode)
}
needsEvent := out.Mode == WaitModeEvent || out.Mode == WaitModeAuto
if needsEvent {
if out.EventKey == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires event_key", canonical, out.Mode)
}
if out.MatchField == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires match_field", canonical, out.Mode)
}
if out.ResourceQuery == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires resource_query", canonical, out.Mode)
}
}
if out.StatusQuery == "" {
return nil, fmt.Errorf("schema tool %s wait mode %s requires status_query", canonical, out.Mode)
}
if len(in.Terminal) == 0 {
return nil, fmt.Errorf("schema tool %s wait has no terminal states", canonical)
}
out.Terminal = make(map[string]ResultOutcome, len(in.Terminal))
for status, outcome := range in.Terminal {
status = strings.TrimSpace(status)
if status == "" {
return nil, fmt.Errorf("schema tool %s wait has a blank terminal status", canonical)
}
if _, dup := out.Terminal[status]; dup {
return nil, fmt.Errorf("schema tool %s wait has duplicate terminal status %q", canonical, status)
}
// Terminal states must close into success or failure. Pending and
// partial are not wait outcomes: pending is expressed through
// timeout, and partial requires the typed multi-status payload only
// the leaf can construct.
if outcome != ResultOutcomeSuccess && outcome != ResultOutcomeFailure {
return nil, fmt.Errorf(
"schema tool %s wait terminal status %q must map to success or failure, got %q",
canonical, status, outcome)
}
out.Terminal[status] = outcome
}
seenPending := make(map[string]bool, len(in.PendingValues))
for _, value := range in.PendingValues {
value = strings.TrimSpace(value)
if value == "" {
return nil, fmt.Errorf("schema tool %s wait has a blank pending value", canonical)
}
if _, conflict := out.Terminal[value]; conflict {
return nil, fmt.Errorf("schema tool %s wait status %q is both terminal and pending", canonical, value)
}
if seenPending[value] {
return nil, fmt.Errorf("schema tool %s wait has duplicate pending value %q", canonical, value)
}
seenPending[value] = true
out.PendingValues = append(out.PendingValues, value)
}
if in.DefaultTimeoutSecs < 0 {
return nil, fmt.Errorf("schema tool %s wait default_timeout_secs must be >= 0", canonical)
}
return out, nil
}
// ResultOutcome is one closed unified-output envelope outcome.
type ResultOutcome string
+227
View File
@@ -0,0 +1,227 @@
// Copyright 2026 Alibaba Group
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package contract
import "testing"
func TestWaitSpecValidateAcceptsReviewedShapes(t *testing.T) {
cases := []WaitSpec{
{
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "result.status",
Terminal: map[string]ResultOutcome{
"COMPLETED": ResultOutcomeSuccess,
"REJECTED": ResultOutcomeFailure,
},
PendingValues: []string{"NEW", "RUNNING"},
DefaultTimeoutSecs: 600,
},
}
for i, spec := range cases {
if err := spec.Validate("sample.tool"); err != nil {
t.Fatalf("case %d: unexpected error: %v", i, err)
}
}
}
func TestWaitSpecValidateRejectsMalformedShapes(t *testing.T) {
terminal := map[string]ResultOutcome{"COMPLETED": ResultOutcomeSuccess}
cases := map[string]WaitSpec{
"no mode": {
Terminal: terminal,
},
"unknown mode": {
Mode: "webhook",
Terminal: terminal,
},
"poll without poll_command": {
Mode: WaitModePoll,
StatusQuery: "status",
Terminal: terminal,
},
"poll without status_query": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
Terminal: terminal,
},
"event without event_key": {
Mode: WaitModeEvent,
MatchField: "process_instance_id",
ResourceQuery: "id",
StatusQuery: "result.status",
Terminal: terminal,
},
"event without match_field": {
Mode: WaitModeEvent,
EventKey: "bpms_instance_change",
ResourceQuery: "id",
StatusQuery: "result.status",
Terminal: terminal,
},
"event without resource_query": {
Mode: WaitModeEvent,
EventKey: "bpms_instance_change",
MatchField: "process_instance_id",
StatusQuery: "result.status",
Terminal: terminal,
},
"auto missing poll_command": {
Mode: WaitModeAuto,
EventKey: "export_finished",
MatchField: "job_id",
ResourceQuery: "job_id",
StatusQuery: "status",
Terminal: terminal,
},
"no terminal states": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
},
"blank terminal status": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{" ": ResultOutcomeSuccess},
},
"terminal outcome pending": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcomePending, "REJECTED": ResultOutcomeFailure},
},
"terminal outcome partial": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcomePartialFailure, "REJECTED": ResultOutcomeFailure},
},
"terminal outcome outside closed set": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcome("explosion")},
},
"only pending terminal outcome": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"NEW": ResultOutcomePending},
},
"status both terminal and pending": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{"COMPLETED"},
},
"blank pending value": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{" "},
},
"negative timeout default": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
DefaultTimeoutSecs: -1,
},
}
for name, spec := range cases {
if err := spec.Validate("sample.tool"); err == nil {
t.Fatalf("%s: expected error, got nil", name)
}
}
}
func TestNormalizeWaitSpecTrimsStatusValuesIntoWireForm(t *testing.T) {
in := &WaitSpec{
Mode: " poll ",
PollCommand: " oa approval-instance get ",
StatusQuery: " result.status ",
Terminal: map[string]ResultOutcome{" COMPLETED ": ResultOutcomeSuccess, "REJECTED": ResultOutcomeFailure},
PendingValues: []string{" NEW ", "RUNNING"},
DefaultTimeoutSecs: 60,
}
out, err := NormalizeWaitSpec(in, "sample.tool")
if err != nil {
t.Fatalf("NormalizeWaitSpec() error = %v", err)
}
if out.Mode != WaitModePoll || out.PollCommand != "oa approval-instance get" || out.StatusQuery != "result.status" {
t.Fatalf("normalized scalars: %#v", out)
}
if len(out.Terminal) != 2 {
t.Fatalf("terminal=%#v, want two trimmed keys", out.Terminal)
}
if got := out.Terminal["COMPLETED"]; got != ResultOutcomeSuccess {
t.Fatalf("terminal[COMPLETED]=%q, want success (key must be trimmed)", got)
}
if _, padded := out.Terminal[" COMPLETED "]; padded {
t.Fatal("padded terminal key survived normalization")
}
for i, want := range []string{"NEW", "RUNNING"} {
if out.PendingValues[i] != want {
t.Fatalf("pending[%d]=%q, want %q", i, out.PendingValues[i], want)
}
}
// The input declaration must stay untouched (defensive copy).
if _, padded := in.Terminal[" COMPLETED "]; !padded {
t.Fatal("NormalizeWaitSpec mutated its input terminal map")
}
if in.PendingValues[0] != " NEW " {
t.Fatal("NormalizeWaitSpec mutated its input pending values")
}
}
func TestNormalizeWaitSpecRejectsDuplicatesAndConflictsAfterTrim(t *testing.T) {
terminal := map[string]ResultOutcome{"COMPLETED": ResultOutcomeSuccess}
cases := map[string]*WaitSpec{
"terminal keys collapsing after trim": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: map[string]ResultOutcome{"COMPLETED": ResultOutcomeSuccess, " COMPLETED ": ResultOutcomeFailure},
},
"pending values collapsing after trim": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{"NEW", " NEW "},
},
"terminal/pending conflict hidden by padding": {
Mode: WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "status",
Terminal: terminal,
PendingValues: []string{" COMPLETED "},
},
}
for name, spec := range cases {
if _, err := NormalizeWaitSpec(spec, "sample.tool"); err == nil {
t.Fatalf("%s: expected error, got nil", name)
}
}
}
func TestNormalizeWaitSpecNilReturnsNil(t *testing.T) {
out, err := NormalizeWaitSpec(nil, "sample.tool")
if err != nil || out != nil {
t.Fatalf("NormalizeWaitSpec(nil) = %#v, %v", out, err)
}
}
+4
View File
@@ -42,6 +42,7 @@ type ContractDecl struct {
Positionals []contract.RuntimeSchemaPositional
Parameters []contract.ParamDecl
DryRun *contract.DryRunSpec
Wait *contract.WaitSpec
Result *contract.ResultSpec
Pagination *contract.PaginationSpec
Interface *contract.InterfaceSpec
@@ -146,6 +147,9 @@ 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,6 +200,13 @@ 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"},
@@ -218,6 +225,8 @@ 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
@@ -226,6 +235,10 @@ func TestFrameworkContractFinalDeepCopyAndSafetyConflicts(t *testing.T) {
if again.Parameters[0].Enum[0] != "a" || !*again.Parameters[0].Required || *again.Selection.ExampleDispositions[0].Index != 1 || !*again.Selection.Reviewed {
t.Fatalf("stored payload aliased input: %#v", again)
}
if again.Wait == payload.Wait || again.Wait.Terminal["COMPLETED"] != contract.ResultOutcomeSuccess || again.Wait.PendingValues[0] != "NEW" {
t.Fatalf("wait spec aliased input: %#v", again.Wait)
t.Fatalf("stored payload aliased input: %#v", again)
}
matching := &cobra.Command{Use: "matching"}
t.Cleanup(func() { ClearRuntimeContractFinalForTest(matching) })
+9
View File
@@ -74,6 +74,15 @@ func cloneContractFinalPayload(in contract.ContractFinalPayload) contract.Contra
value := *in.DryRun
out.DryRun = &value
}
if in.Wait != nil {
value := *in.Wait
value.Terminal = make(map[string]contract.ResultOutcome, len(in.Wait.Terminal))
for status, outcome := range in.Wait.Terminal {
value.Terminal[status] = outcome
}
value.PendingValues = cloneSlice(in.Wait.PendingValues)
out.Wait = &value
}
if in.Result != nil {
value := *in.Result
value.Outcomes = cloneSlice(in.Result.Outcomes)
+330
View File
@@ -49,12 +49,16 @@ package corecmd
import (
"bufio"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"math"
"os"
"strconv"
"strings"
"time"
"github.com/mattn/go-isatty"
"github.com/spf13/cobra"
@@ -64,6 +68,7 @@ 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"
)
@@ -282,6 +287,21 @@ 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.
@@ -355,6 +375,15 @@ 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") }
@@ -374,6 +403,8 @@ 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)
@@ -390,6 +421,7 @@ 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)
@@ -455,6 +487,10 @@ 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)
@@ -506,6 +542,293 @@ 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
@@ -1408,6 +1731,13 @@ func AttachContract(cmd *cobra.Command, safety contract.SafetySpec, decl Contrac
d.PreviewKind = strings.TrimSpace(d.PreviewKind)
payload.DryRun = &d
}
if decl.Wait != nil && strings.TrimSpace(decl.Wait.Mode) != "" {
waitSpec, err := contract.NormalizeWaitSpec(decl.Wait, decl.Identity.CanonicalPath)
if err != nil {
panic(fmt.Sprintf("command %q has invalid Contract.Wait: %v", cmd.Name(), err))
}
payload.Wait = waitSpec
}
if decl.Result != nil {
result, err := contract.NormalizeResultSpec(decl.Result, decl.Identity.CanonicalPath)
if err != nil {
+109
View File
@@ -26,6 +26,7 @@ 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"
)
@@ -2120,3 +2121,111 @@ func TestCrossPlatformCoverageEmbedContractSkipsBlankAndHiddenFlags(t *testing.T
}
}
}
// TestWaitResourceStronglyTypedDTOs verifies waitResource handles struct and
// struct pointer results via JSON normalization (P1 fix for auto-CR).
func TestWaitResourceStronglyTypedDTOs(t *testing.T) {
type TaskDTO struct {
TaskID string `json:"task_id"`
Status string `json:"status"`
}
type NestedDTO struct {
Meta struct {
ResourceID string `json:"resource_id"`
} `json:"meta"`
}
tests := []struct {
name string
data any
query string
wantResource string
wantErrSubstr string
}{
{
name: "map[string]any fast path",
data: map[string]any{"task_id": "abc123"},
query: "task_id",
wantResource: "abc123",
},
{
name: "struct value",
data: TaskDTO{TaskID: "struct-456", Status: "running"},
query: "task_id",
wantResource: "struct-456",
},
{
name: "struct pointer",
data: &TaskDTO{TaskID: "ptr-789", Status: "pending"},
query: "task_id",
wantResource: "ptr-789",
},
{
name: "nested struct dotted query",
data: NestedDTO{},
query: "meta.resource_id",
wantResource: "",
wantErrSubstr: "not found",
},
{
name: "nested struct with value",
data: func() NestedDTO {
var d NestedDTO
d.Meta.ResourceID = "nested-xyz"
return d
}(),
query: "meta.resource_id",
wantResource: "nested-xyz",
},
{
name: "nil data",
data: nil,
query: "task_id",
wantErrSubstr: "nil",
},
{
name: "non-object data (string)",
data: "not-an-object",
query: "task_id",
wantErrSubstr: "not an object",
},
{
name: "non-object data (slice)",
data: []string{"a", "b"},
query: "task_id",
wantErrSubstr: "not an object",
},
{
name: "missing query field in struct",
data: TaskDTO{TaskID: "abc"},
query: "nonexistent",
wantErrSubstr: "not found",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
decl := &contract.WaitSpec{
ResourceQuery: tt.query,
}
result := output.Success(tt.data)
resource, err := waitResource(decl, result)
if tt.wantErrSubstr != "" {
if err == nil {
t.Fatalf("expected error containing %q, got nil", tt.wantErrSubstr)
}
if !strings.Contains(err.Error(), tt.wantErrSubstr) {
t.Errorf("error = %q; want substring %q", err.Error(), tt.wantErrSubstr)
}
return
}
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if resource != tt.wantResource {
t.Errorf("resource = %q; want %q", resource, tt.wantResource)
}
})
}
}
+902
View File
@@ -0,0 +1,902 @@
// Copyright 2026 Alibaba Group
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package corecmd
import (
"bytes"
"context"
"errors"
"io"
"math"
"strconv"
"strings"
"testing"
"time"
"github.com/spf13/cobra"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contract"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contractfinal"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/output"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/wait"
)
func waitTestDecl() ContractDecl {
return ContractDecl{
Title: "Wait Title",
Description: "Wait Desc",
Wait: &contract.WaitSpec{
Mode: contract.WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{"COMPLETED": contract.ResultOutcomeSuccess, "REJECTED": contract.ResultOutcomeFailure},
PendingValues: []string{"NEW", "RUNNING"},
DefaultTimeoutSecs: 60,
},
Interface: &contract.InterfaceSpec{Mode: "local", Availability: "available"},
Selection: contract.SelectionSpec{
AgentSummary: "summary",
UseWhen: []string{"when wait"},
AvoidWhen: []string{"when nowait"},
Examples: []string{"dws wait-sample --wait"},
},
Identity: contract.ToolIdentitySpec{ProductID: "sample", Name: "waitsample", CanonicalPath: "sample.waitsample", CLIPath: "wait-sample", PrimaryCLIPath: "wait-sample"},
}
}
func baseWaitSpec(decl ContractDecl, poll func(context.Context, *Ctx) (wait.PollDoc, error)) Spec {
return Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: decl,
WaitPoll: poll,
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1",
State: "NEW",
NextCommand: "dws wait-sample --id job-1",
}), nil
},
}
}
func TestWaitFlagsOnlyRegisteredWhenDeclared(t *testing.T) {
declared := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
if flag := declared.Flags().Lookup(waitFlagName); flag == nil {
t.Fatal("declared command missing --wait flag")
}
if flag := declared.Flags().Lookup(waitTimeoutFlagName); flag == nil {
t.Fatal("declared command missing --wait-timeout flag")
}
undeclared := New(Spec{
Use: "nowait",
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Invoke: func(*Ctx, map[string]any) error { return nil },
})
if flag := undeclared.Flags().Lookup(waitFlagName); flag != nil {
t.Fatal("undeclared command registered --wait")
}
undeclared.SetArgs([]string{"--wait"})
if err := undeclared.Execute(); err == nil || !strings.Contains(err.Error(), "unknown flag") {
t.Fatalf("err=%v want unknown-flag", err)
}
}
func TestValidateWaitDeclPairsDeclarationWithImplementation(t *testing.T) {
decl := waitTestDecl()
spec := baseWaitSpec(decl, nil)
expectPanic(t, func() { New(spec) }, "WaitPoll")
spec.WaitPoll = func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil }
expectPanic(t, func() {
New(Spec{
Use: "hook-only",
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Invoke: func(*Ctx, map[string]any) error { return nil },
WaitPoll: func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil },
})
}, "Contract.Wait")
}
func expectPanic(t *testing.T, fn func(), want string) {
t.Helper()
defer func() {
recovered := recover()
if recovered == nil {
t.Fatalf("expected panic containing %q", want)
}
if message, ok := recovered.(string); !ok || !strings.Contains(message, want) {
t.Fatalf("panic=%v want containing %q", recovered, want)
}
}()
fn()
}
func TestWaitTimeoutFlagDefaultsComeFromDeclaration(t *testing.T) {
stub := func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}
declared := New(baseWaitSpec(waitTestDecl(), stub))
if value, err := declared.Flags().GetInt(waitTimeoutFlagName); err != nil || value != 60 {
t.Fatalf("declared default=%d/%v, want reviewed 60", value, err)
}
decl := waitTestDecl()
decl.Wait.DefaultTimeoutSecs = 0
fallback := New(baseWaitSpec(decl, stub))
if value, err := fallback.Flags().GetInt(waitTimeoutFlagName); err != nil || value != DefaultWaitTimeoutSecs {
t.Fatalf("fallback default=%d/%v, want framework %d", value, err, DefaultWaitTimeoutSecs)
}
// A non-positive explicit value falls back to the framework default.
fallback.SetArgs([]string{"--wait-timeout", "0", "--wait"})
if err := fallback.Flags().Set(waitTimeoutFlagName, "0"); err != nil {
t.Fatal(err)
}
if got := waitTimeoutSecs(fallback); got != DefaultWaitTimeoutSecs {
t.Fatalf("waitTimeoutSecs=%d, want %d", got, DefaultWaitTimeoutSecs)
}
}
func TestWaitTimeoutDurationRejectsOverflowingSeconds(t *testing.T) {
// math.MaxInt64 (9223372036854775807) is a legal pflag int on 64-bit
// platforms and overflows time.Duration(secs)*time.Second to a negative
// value, which would disable the wait deadline.
if _, err := waitTimeoutDuration(math.MaxInt64); err == nil || !strings.Contains(err.Error(), "超出可表示范围") {
t.Fatalf("err=%v, want overflow validation", err)
}
d, err := waitTimeoutDuration(maxWaitTimeoutSecs)
if err != nil {
t.Fatal(err)
}
if d <= 0 || d != time.Duration(maxWaitTimeoutSecs)*time.Second {
t.Fatalf("duration=%d, want the largest representable timeout", d)
}
// Non-positive second counts fall back to the framework default instead
// of disabling the deadline (waitTimeoutSecs already maps a zero/negative
// flag to the default; this keeps the conversion itself fail-safe).
if got, err := waitTimeoutDuration(0); err != nil || got != time.Duration(DefaultWaitTimeoutSecs)*time.Second {
t.Fatalf("duration/err=%d/%v, want framework default", got, err)
}
if got, err := waitTimeoutDuration(-5); err != nil || got != time.Duration(DefaultWaitTimeoutSecs)*time.Second {
t.Fatalf("duration/err=%d/%v, want framework default", got, err)
}
}
func TestResultInvokeWaitTimeoutOverflowIsValidationError(t *testing.T) {
if int64(math.MaxInt) <= maxWaitTimeoutSecs {
t.Skip("platform int cannot overflow time.Duration")
}
polled := false
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", strconv.Itoa(math.MaxInt)})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !strings.Contains(err.Error(), "超出可表示范围") {
t.Fatalf("err=%v, want overflow validation", err)
}
if polled {
t.Fatal("overflowing --wait-timeout must not start the wait loop")
}
}
func TestResultInvokeWaitPollErrorFailsTheCommand(t *testing.T) {
boom := errors.New("rpc down")
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return nil, boom
}))
cmd.SetArgs([]string{"--wait"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !strings.Contains(err.Error(), "rpc down") {
t.Fatalf("err=%v, want poll error surfaced", err)
}
}
func TestResultInvokeWaitUnknownStatusFailsClosed(t *testing.T) {
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "Mystery"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !wait.IsUnknownStatus(err) {
t.Fatalf("err=%v, want unknown-status", err)
}
}
func TestWaitCtxAccessorsExposeDeclaredCapability(t *testing.T) {
var gotWait bool
var gotTimeout int
cmd := New(baseWaitSpec(waitTestDecl(), func(_ context.Context, c *Ctx) (wait.PollDoc, error) {
gotWait = c.Wait()
gotTimeout = c.WaitTimeoutSecs()
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", "90"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if !gotWait || gotTimeout != 90 {
t.Fatalf("ctx accessors=%v/%d", gotWait, gotTimeout)
}
}
func eventTestDecl(mode string) ContractDecl {
decl := waitTestDecl()
decl.Wait.Mode = mode
decl.Wait.EventKey = "bpms_instance_change"
decl.Wait.MatchField = "process_instance_id"
decl.Wait.ResourceQuery = "id"
return decl
}
type scriptedStream struct {
events []wait.PollDoc
err error
}
func (s *scriptedStream) Recv(context.Context) (wait.PollDoc, error) {
if len(s.events) > 0 {
doc := s.events[0]
s.events = s.events[1:]
return doc, nil
}
if s.err != nil {
return nil, s.err
}
return nil, io.EOF
}
func TestValidateWaitDeclPairsModeWithHooks(t *testing.T) {
poll := func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil }
events := func(context.Context, *Ctx) (wait.EventStream, error) { return nil, nil }
cases := []struct {
name string
mode string
waitPoll bool
waitEvents bool
want string
}{
{"event without WaitEvents", contract.WaitModeEvent, false, false, "WaitEvents"},
{"auto without WaitPoll", contract.WaitModeAuto, false, true, "WaitPoll"},
{"poll with WaitEvents", contract.WaitModePoll, true, true, "WaitEvents"},
{"event with WaitPoll", contract.WaitModeEvent, true, true, "WaitPoll"},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
expectPanic(t, func() {
New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(tc.mode),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{}, nil), nil
},
WaitPoll: hookOrNil(tc.waitPoll, poll),
WaitEvents: eventHookOrNil(tc.waitEvents, events),
})
}, tc.want)
})
}
}
func hookOrNil(set bool, hook func(context.Context, *Ctx) (wait.PollDoc, error)) func(context.Context, *Ctx) (wait.PollDoc, error) {
if !set {
return nil
}
return hook
}
func eventHookOrNil(set bool, hook func(context.Context, *Ctx) (wait.EventStream, error)) func(context.Context, *Ctx) (wait.EventStream, error) {
if !set {
return nil
}
return hook
}
func runWaitModeCommand(t *testing.T, decl ContractDecl, poll func(context.Context, *Ctx) (wait.PollDoc, error), events func(context.Context, *Ctx) (wait.EventStream, error), args ...string) (string, error) {
t.Helper()
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: decl,
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample --id job-1",
}), nil
},
WaitPoll: poll,
WaitEvents: events,
})
cmd.SetArgs(append([]string{"--wait"}, args...))
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
err := cmd.Execute()
return stdout.String(), err
}
func TestEventModeClosesEnvelopeFromCorrelatedEvent(t *testing.T) {
stream := &scriptedStream{events: []wait.PollDoc{
{"process_instance_id": "other", "result": map[string]any{"status": "COMPLETED"}},
{"process_instance_id": "job-1", "result": map[string]any{"status": "REJECTED"}},
}}
stdout, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeEvent), nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return stream, nil
})
if err != nil {
t.Fatal(err)
}
if !strings.Contains(stdout, `"outcome": "failure"`) || !strings.Contains(stdout, `"type": "wait"`) {
t.Fatalf("stdout=%s", stdout)
}
}
func TestEventModeSurfacesStreamEndAsError(t *testing.T) {
_, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeEvent), nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil
})
if err == nil || !errors.Is(err, wait.ErrEventStreamEnded) {
t.Fatalf("err=%v, want stream-ended", err)
}
}
func TestEventModeRejectsUnresolvableResource(t *testing.T) {
decl := eventTestDecl(contract.WaitModeEvent)
decl.Wait.ResourceQuery = "missing"
_, err := runWaitModeCommand(t, decl, nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil
})
if err == nil || !strings.Contains(err.Error(), "resource query") {
t.Fatalf("err=%v", err)
}
}
func TestAutoModeFallsBackToPollOnStreamEnd(t *testing.T) {
polled := false
stdout, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeAuto),
func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
},
func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil // ends immediately
})
if err != nil {
t.Fatal(err)
}
if !polled {
t.Fatal("auto mode did not fall back to polling")
}
if !strings.Contains(stdout, `"outcome": "success"`) {
t.Fatalf("stdout=%s", stdout)
}
}
func TestAutoModeFallsBackToPollOnSubscriptionFailure(t *testing.T) {
polled := false
_, err := runWaitModeCommand(t, eventTestDecl(contract.WaitModeAuto),
func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
},
func(context.Context, *Ctx) (wait.EventStream, error) {
return nil, errors.New("no subscriber credential")
})
if err != nil {
t.Fatal(err)
}
if !polled {
t.Fatal("auto mode did not fall back to polling on subscription failure")
}
}
func TestResultInvokeWaitClosesEnvelopeOutcome(t *testing.T) {
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "REJECTED"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 8 {
t.Fatalf("stored code/emitted=%d/%v, want dedicated wait-terminal-failure code 8", code, emitted)
}
if !strings.Contains(stdout.String(), `"type": "wait"`) {
t.Fatalf("stdout=%s, want error.type wait", stdout.String())
}
if !strings.Contains(stdout.String(), `"outcome": "failure"`) {
t.Fatalf("stdout=%s", stdout.String())
}
// The final emitted envelope must carry the observed terminal status in
// meta.operation.state — the acceptance-phase state ("NEW") must not
// survive the close (P1 regression guard).
if !strings.Contains(stdout.String(), `"state": "REJECTED"`) {
t.Fatalf("stdout=%s, want operation.state synced to the terminal status", stdout.String())
}
if strings.Contains(stdout.String(), `"state": "NEW"`) {
t.Fatalf("stdout=%s, acceptance-phase operation.state leaked into the terminal envelope", stdout.String())
}
if strings.Contains(stdout.String(), `"timed_out": true`) {
t.Fatalf("stdout=%s, terminal close must not claim timed_out", stdout.String())
}
}
func TestResultInvokeWaitSuccessCloseSyncsOperationState(t *testing.T) {
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v, want success exit 0", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "success"`) {
t.Fatalf("stdout=%s", stdout.String())
}
// Success close must publish the terminal status, never the stale
// acceptance-phase state (no outcome=success with state=processing/NEW).
if !strings.Contains(stdout.String(), `"state": "COMPLETED"`) {
t.Fatalf("stdout=%s, want operation.state synced to the terminal status", stdout.String())
}
if strings.Contains(stdout.String(), `"state": "NEW"`) {
t.Fatalf("stdout=%s, acceptance-phase operation.state leaked into the success envelope", stdout.String())
}
// Operation identity (id / next_command) survives the terminal close.
if !strings.Contains(stdout.String(), `"id": "job-1"`) || !strings.Contains(stdout.String(), `"next_command"`) {
t.Fatalf("stdout=%s, want operation id/next_command preserved", stdout.String())
}
}
func TestResultInvokeWaitTimeoutKeepsPending(t *testing.T) {
polls := 0
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
polls++
return wait.PollDoc{"result": map[string]any{"status": "RUNNING"}}, nil
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
if polls == 0 {
t.Fatal("wait phase never polled")
}
}
func TestWaitDeclRequiresResultInvokeDispatcher(t *testing.T) {
// A declared wait on the legacy Invoke path would observe a failure
// terminal while still exiting 0 — construction must reject it.
expectPanic(t, func() {
New(Spec{
Use: "wait-sample",
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: waitTestDecl(),
WaitPoll: func(context.Context, *Ctx) (wait.PollDoc, error) { return nil, nil },
Invoke: func(*Ctx, map[string]any) error { return nil },
})
}, "ResultInvoke")
}
func TestResultInvokeWithoutWaitFlagSkipsPhase(t *testing.T) {
polled := false
cmd := New(baseWaitSpec(waitTestDecl(), func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs(nil)
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if polled {
t.Fatal("wait phase ran without --wait")
}
}
func TestAttachContractPanicsOnInvalidWaitDeclaration(t *testing.T) {
decl := waitTestDecl()
decl.Wait.Mode = "event" // not implemented
defer func() {
recovered := recover()
if recovered == nil {
t.Fatal("expected panic on invalid Contract.Wait")
}
if message, ok := recovered.(string); !ok || !strings.Contains(message, "Contract.Wait") {
t.Fatalf("panic=%v", recovered)
}
}()
AttachContract(&cobra.Command{Use: "wait-sample"}, contract.SafetySpec{
Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent",
}, decl, "", "")
}
func TestWaitDeclPaddedStatusValuesAreNormalized(t *testing.T) {
// A declaration whose status values carry surrounding whitespace must be
// canonicalized at construction so the runtime wait engine and the
// published Schema agree: the backend returns "COMPLETED" verbatim, and
// a padded terminal key would fail closed as an unknown status.
decl := waitTestDecl()
decl.Wait.Terminal = map[string]contract.ResultOutcome{
" COMPLETED ": contract.ResultOutcomeSuccess,
"\tREJECTED": contract.ResultOutcomeFailure,
}
decl.Wait.PendingValues = []string{" NEW ", "RUNNING "}
cmd := New(baseWaitSpec(decl, func(context.Context, *Ctx) (wait.PollDoc, error) {
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
}))
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatalf("padded declaration must still reach the terminal status: %v", err)
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v, want success", code, emitted)
}
final, ok := contractfinal.RuntimeContractFinal(cmd)
if !ok || final.Wait == nil {
t.Fatal("registered ContractFinal lost the wait capability")
}
if _, ok := final.Wait.Terminal["COMPLETED"]; !ok {
t.Fatalf("registered terminal table not trimmed: %#v", final.Wait.Terminal)
}
for _, value := range final.Wait.PendingValues {
if strings.TrimSpace(value) != value {
t.Fatalf("registered pending value %q not trimmed", value)
}
}
}
func TestNewPanicsOnDuplicateOrConflictingWaitStatusesAfterTrim(t *testing.T) {
// Values that collapse onto one status after trimming are programming
// errors: silently merging them would pick one outcome for two authored
// declarations.
dupTerminal := waitTestDecl()
dupTerminal.Wait.Terminal = map[string]contract.ResultOutcome{
"COMPLETED": contract.ResultOutcomeSuccess,
" COMPLETED": contract.ResultOutcomeFailure,
}
expectPanic(t, func() { New(baseWaitSpec(dupTerminal, nil)) }, "Contract.Wait")
conflict := waitTestDecl()
conflict.Wait.PendingValues = []string{" COMPLETED "}
expectPanic(t, func() { New(baseWaitSpec(conflict, nil)) }, "Contract.Wait")
}
func TestContractDeclEmptyTreatsWaitAsAuthored(t *testing.T) {
// Only Wait is authored: empty() must report non-empty through the Wait
// branch (before validateContractDecl then fails on the missing prose).
decl := ContractDecl{Wait: &contract.WaitSpec{
Mode: contract.WaitModePoll,
PollCommand: "oa approval-instance get",
StatusQuery: "result.status",
Terminal: map[string]contract.ResultOutcome{"COMPLETED": contract.ResultOutcomeSuccess},
}}
if decl.Empty() {
t.Fatal("Wait-only declaration must count as authored")
}
defer func() {
if recover() == nil {
t.Fatal("expected validateContractDecl to reject the missing prose")
}
}()
validateContractDecl(Spec{Use: "wait-only", Contract: decl})
}
func TestEventModeRejectsNonObjectResultData(t *testing.T) {
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeEvent),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending([]any{"not", "an", "object"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample",
}), nil
},
WaitEvents: func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil
},
})
cmd.SetArgs([]string{"--wait"})
ctx, _ := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
err := cmd.Execute()
if err == nil || !strings.Contains(err.Error(), "not an object") {
t.Fatalf("err=%v, want non-object data rejection", err)
}
}
func TestEventModeSubscriptionFailureSurfacesInStrictMode(t *testing.T) {
decl := eventTestDecl(contract.WaitModeEvent)
_, err := runWaitModeCommand(t, decl, nil, func(context.Context, *Ctx) (wait.EventStream, error) {
return nil, errors.New("no subscriber credential")
})
if err == nil || !strings.Contains(err.Error(), "subscription failed") {
t.Fatalf("err=%v, want subscription failure surfaced", err)
}
}
func TestResultInvokeNonPendingSkipsWaitPhase(t *testing.T) {
partial, err := output.NewPartialData(2,
[]any{map[string]any{"id": "ok"}},
[]output.PartialFailedEntry{{ID: "bad", Error: &output.ErrorInfo{Type: "api", Message: "item failed"}}},
nil)
if err != nil {
t.Fatal(err)
}
cases := []struct {
name string
result output.CommandResult
want string
}{
{"failure", output.Failure(&output.ErrorInfo{Type: "api", Message: "business failed"}), "failure"},
{"success", output.Success(map[string]any{"id": "job-1"}), "success"},
{"partial", output.Partial(partial), "partial_failure"},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
polled := false
subscribed := false
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeAuto),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return tc.result, nil
},
WaitPoll: func(context.Context, *Ctx) (wait.PollDoc, error) {
polled = true
return wait.PollDoc{"result": map[string]any{"status": "COMPLETED"}}, nil
},
WaitEvents: func(context.Context, *Ctx) (wait.EventStream, error) {
subscribed = true
return &scriptedStream{}, nil
},
})
cmd.SetArgs([]string{"--wait"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
if err := cmd.Execute(); err != nil {
t.Fatal(err)
}
if polled || subscribed {
t.Fatal("wait phase must not call WaitPoll/WaitEvents for a non-pending initial result")
}
if _, emitted := output.StoredExitCode(store); !emitted {
t.Fatal("initial result was not stored")
}
if !strings.Contains(stdout.String(), `"outcome": "`+tc.want+`"`) {
t.Fatalf("stdout=%s, want outcome %s preserved", stdout.String(), tc.want)
}
if strings.Contains(stdout.String(), `"type": "wait"`) {
t.Fatalf("stdout=%s, wait phase overwrote the original envelope", stdout.String())
}
})
}
}
func TestWaitTimeoutCancelsBlockingPoll(t *testing.T) {
started := make(chan struct{})
cmd := New(baseWaitSpec(waitTestDecl(), func(ctx context.Context, c *Ctx) (wait.PollDoc, error) {
close(started)
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-c.Command().Context().Done():
return nil, c.Command().Context().Err()
}
}))
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
done := make(chan error, 1)
go func() { done <- cmd.Execute() }()
select {
case <-started:
case <-time.After(2 * time.Second):
t.Fatal("blocking poll never started")
}
select {
case err := <-done:
if err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
case <-time.After(3 * time.Second):
t.Fatal("blocked poll was not cancelled by --wait-timeout")
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
}
func TestWaitTimeoutCancelsBlockingSubscribe(t *testing.T) {
started := make(chan struct{})
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeEvent),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample --id job-1",
}), nil
},
WaitEvents: func(ctx context.Context, c *Ctx) (wait.EventStream, error) {
close(started)
// Leaf subscribe may wait on either the hook ctx or the cobra
// command context; both must carry the wait-timeout deadline.
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-c.Command().Context().Done():
return nil, c.Command().Context().Err()
}
},
})
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
done := make(chan error, 1)
go func() { done <- cmd.Execute() }()
select {
case <-started:
case <-time.After(2 * time.Second):
t.Fatal("blocking subscribe never started")
}
select {
case err := <-done:
if err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
case <-time.After(3 * time.Second):
t.Fatal("blocked subscribe was not cancelled by --wait-timeout")
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
}
func TestWaitTimeoutCancelsBlockingPollAfterAutoFallback(t *testing.T) {
started := make(chan struct{})
cmd := New(Spec{
Use: "wait-sample",
OutputRollout: output.RolloutUnifiedActive,
Safety: contract.SafetySpec{Effect: "read", Risk: "low", Confirmation: "not_required", Idempotency: "idempotent"},
Contract: eventTestDecl(contract.WaitModeAuto),
ResultInvoke: func(*Ctx, map[string]any) (output.CommandResult, error) {
return output.Pending(map[string]any{"id": "job-1"}, &output.OperationInfo{
ID: "job-1", State: "NEW", NextCommand: "dws wait-sample --id job-1",
}), nil
},
WaitEvents: func(context.Context, *Ctx) (wait.EventStream, error) {
return &scriptedStream{}, nil // ends immediately → poll fallback
},
WaitPoll: func(ctx context.Context, _ *Ctx) (wait.PollDoc, error) {
close(started)
<-ctx.Done()
return nil, ctx.Err()
},
})
cmd.SetArgs([]string{"--wait", "--wait-timeout", "1"})
ctx, store := output.WithResultStore(context.Background())
cmd.SetContext(ctx)
var stdout bytes.Buffer
cmd.SetOut(&stdout)
cmd.PersistentPostRunE = func(executed *cobra.Command, _ []string) error {
_, _, err := output.EmitStoredResult(executed)
return err
}
done := make(chan error, 1)
go func() { done <- cmd.Execute() }()
select {
case <-started:
case <-time.After(2 * time.Second):
t.Fatal("auto-fallback blocking poll never started")
}
select {
case err := <-done:
if err != nil {
t.Fatalf("timeout wait must exit 0 (pending is not failure): %v", err)
}
case <-time.After(3 * time.Second):
t.Fatal("blocked auto-fallback poll was not cancelled by --wait-timeout")
}
if code, emitted := output.StoredExitCode(store); !emitted || code != 0 {
t.Fatalf("stored code/emitted=%d/%v", code, emitted)
}
if !strings.Contains(stdout.String(), `"outcome": "pending"`) {
t.Fatalf("stdout=%s", stdout.String())
}
}
+2
View File
@@ -58,6 +58,8 @@ 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
+3
View File
@@ -435,6 +435,7 @@ const (
exitCodeInternal = 5
exitCodeDiscovery = 6
exitCodePartial = 7 // partial_failure 专用(契约 §4;规划 WS2 第4项;B142 将在 errors 侧补同源常量)
exitCodeWait = 8 // wait 终态失败专用(--wait 观察到失败终态;与 partial 同为"仅新增专用码")
)
// subtypeConfirmationRequired 是门禁拦截的 failure 子类标记(契约规范 §2.4),
@@ -491,6 +492,8 @@ func exitCodeForErrorInfo(info *ErrorInfo) int {
return exitCodePermission
case "discovery":
return exitCodeDiscovery
case "wait":
return exitCodeWait
default:
return exitCodeInternal
}
+4 -2
View File
@@ -283,9 +283,11 @@ 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).
// 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).
switch errorType {
case "api", "auth", "validation", "permission", "discovery", "internal":
case "api", "auth", "validation", "permission", "discovery", "internal", "wait":
default:
return fmt.Errorf("output: unsupported failure error.type %q", e.Type)
}
@@ -19,6 +19,7 @@ type forgedResult struct {
func (r forgedResult) Outcome() Outcome { return r.env.Outcome }
func (r forgedResult) ExitCode() int { return r.exit }
func (r forgedResult) Data() any { return r.env.Data }
func (r forgedResult) envelope() *Envelope { copy := r.env; return &copy }
type cloneNode struct {
+90
View File
@@ -19,6 +19,10 @@ 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
}
@@ -30,6 +34,7 @@ type commandResult struct {
func (r *commandResult) Outcome() Outcome { return r.env.Outcome }
func (r *commandResult) ExitCode() int { return r.exitCode }
func (r *commandResult) Data() any { return cloneResultData(r.env.Data) }
func (r *commandResult) envelope() *Envelope {
copy := cloneEnvelope(r.env)
return &copy
@@ -71,6 +76,91 @@ func Partial(data *PartialData, opts ...ResultOption) CommandResult {
return newCommandResult(OutcomePartialFailure, data, nil, opts...)
}
// WithMeta(meta *Meta) ResultOption was declared above; the two options below
// exist for the wait phase.
// WithErrorInfo replaces the error info of the envelope. Used when a wait
// phase closes an accepted result into failure: envelope invariant I3
// requires an error iff the outcome is failure.
func WithErrorInfo(info *ErrorInfo) ResultOption {
return ResultOption{apply: func(env *Envelope) {
if info != nil {
info = cloneErrorInfo(info)
}
env.Error = info
}}
}
// WithOperationTimedOut marks the envelope's async operation as timed out at
// the last observed state, preserving the declared id / next_command resume
// facts (契约规范 §2.2: 超时必须保持 State 真实值并置 TimedOut:true). A result
// without operation info keeps nil — the pending envelope invariant then fails
// at emission, surfacing the leaf bug instead of synthesizing fake resume
// facts.
func WithOperationTimedOut(state string) ResultOption {
return ResultOption{apply: func(env *Envelope) {
if env.Meta == nil || env.Meta.Operation == nil {
return
}
operation := *env.Meta.Operation
// A subscribe/poll that never observed a status still times out
// against the accepted pending result: keep the original state
// rather than wiping it to empty (pending requires operation.state).
if strings.TrimSpace(state) != "" {
operation.State = state
}
operation.TimedOut = true
env.Meta.Operation = &operation
}}
}
// WithOperationTerminalState closes the envelope's async operation at the observed
// terminal status (契约规范 §2.2: 终态封装必须同步 operation.state — a success or
// failure close that kept the acceptance-phase state would emit a
// self-contradicting envelope such as outcome=success with
// operation.state=processing). The declared id / next_command facts are kept
// as the operation identity, and timed_out is cleared: the §2.2 anti-spoof
// rule forbids a timed-out claim on an operation that reached a terminal
// state. A result without operation info is left untouched.
func WithOperationTerminalState(state string) ResultOption {
return ResultOption{apply: func(env *Envelope) {
if env.Meta == nil || env.Meta.Operation == nil {
return
}
operation := *env.Meta.Operation
if strings.TrimSpace(state) != "" {
operation.State = state
}
operation.TimedOut = false
env.Meta.Operation = &operation
}}
}
// WithOutcome rewraps an existing result with a new outcome, preserving data,
// meta, identity, and any error info (subject to the opts). The corecmd wait
// phase uses it to close an accepted result into its terminal (or timed-out
// pending) outcome; the exit code is re-derived from the new envelope.
func WithOutcome(result CommandResult, outcome Outcome, opts ...ResultOption) CommandResult {
env := *result.envelope()
for _, opt := range opts {
if opt.apply != nil {
opt.apply(&env)
}
}
env.Outcome = outcome
env.OK = outcome == OutcomeSuccess || outcome == OutcomePending
if outcome == OutcomeFailure {
// Invariant I3: data and error are mutually exclusive. Closing into
// failure replaces the accepted data with the failure error.
env.Data = nil
}
exitCode := ExitCodeForEnvelope(&env)
if env.Error != nil {
env.Error.ExitCode = exitCode
}
return &commandResult{env: env, exitCode: exitCode}
}
// Failure constructs an immutable typed failure result.
func Failure(info *ErrorInfo, opts ...ResultOption) CommandResult {
if info != nil {
+165
View File
@@ -0,0 +1,165 @@
// Copyright 2026 Alibaba Group
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package output
import (
"strings"
"testing"
)
func pendingAcceptedResult() CommandResult {
return Pending(map[string]any{"id": "job-1"}, &OperationInfo{
ID: "job-1",
State: "NEW",
NextCommand: "dws wait-sample --id job-1",
})
}
func TestWithOutcomeClosesSuccessPreservingDataAndMeta(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomeSuccess)
if result.Outcome() != OutcomeSuccess || result.ExitCode() != 0 {
t.Fatalf("outcome=%s exit=%d", result.Outcome(), result.ExitCode())
}
env := result.envelope()
if env.Meta == nil || env.Meta.Operation == nil {
t.Fatal("operation info lost on success close")
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOutcomeFailureDropsDataAndCarriesError(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomeFailure, WithErrorInfo(&ErrorInfo{
Type: "wait",
Subtype: "terminal_failure",
Message: "等待到达失败终态:REJECTED",
}))
if result.Outcome() != OutcomeFailure || result.ExitCode() != 8 {
t.Fatalf("outcome=%s exit=%d, want failure/8", result.Outcome(), result.ExitCode())
}
env := result.envelope()
if env.Data != nil {
t.Fatal("failure close must drop data (I3)")
}
if env.Error == nil || env.Error.Type != "wait" {
t.Fatal("failure close must carry error info (I3)")
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTimedOutMarksStateAndKeepsResumeFacts(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomePending, WithOperationTimedOut("RUNNING"))
env := result.envelope()
op := env.Meta.Operation
if op.State != "RUNNING" || !op.TimedOut || op.ID != "job-1" || op.NextCommand != "dws wait-sample --id job-1" {
t.Fatalf("operation=%+v", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTimedOutEmptyStatePreservesExisting(t *testing.T) {
result := WithOutcome(pendingAcceptedResult(), OutcomePending, WithOperationTimedOut(""))
op := result.envelope().Meta.Operation
if op.State != "NEW" || !op.TimedOut {
t.Fatalf("operation=%+v, want original state kept and timed_out set", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTimedOutWithoutOperationInfoLeavesEnvelopeUntouched(t *testing.T) {
// A result without operation info keeps nil — ValidateResult must then
// reject the pending envelope instead of the option synthesizing fake
// resume facts.
result := WithOutcome(Success(map[string]any{"ok": true}), OutcomePending, WithOperationTimedOut("RUNNING"))
env := result.envelope()
if env.Meta != nil && env.Meta.Operation != nil {
t.Fatalf("operation=%+v, want untouched", env.Meta.Operation)
}
err := ValidateResult(result)
if err == nil || !strings.Contains(err.Error(), "meta.operation") {
t.Fatalf("err=%v, want pending-requires-operation rejection", err)
}
}
func TestDataAccessorReturnsDeepCopy(t *testing.T) {
result := pendingAcceptedResult()
data, ok := result.Data().(map[string]any)
if !ok {
t.Fatalf("data=%#v", result.Data())
}
data["id"] = "mutated"
again := result.Data().(map[string]any)
if again["id"] != "job-1" {
t.Fatalf("Data() aliased internal state: %#v", again)
}
}
func TestWithOperationTerminalStateSyncsStateAndClearsTimedOut(t *testing.T) {
for _, tc := range []struct {
name string
outcome Outcome
}{
{"success", OutcomeSuccess},
{"failure", OutcomeFailure},
} {
t.Run(tc.name, func(t *testing.T) {
opts := []ResultOption{WithOperationTerminalState("COMPLETED")}
if tc.outcome == OutcomeFailure {
opts = append(opts, WithErrorInfo(&ErrorInfo{
Type: "wait", Subtype: "terminal_failure", Message: "等待到达失败终态:COMPLETED",
}))
}
result := WithOutcome(pendingAcceptedResult(), tc.outcome, opts...)
op := result.envelope().Meta.Operation
if op.State != "COMPLETED" || op.TimedOut {
t.Fatalf("operation=%+v, want terminal state synced and timed_out cleared", op)
}
if op.ID != "job-1" || op.NextCommand != "dws wait-sample --id job-1" {
t.Fatalf("operation=%+v, want resume identity facts preserved", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
})
}
}
func TestWithOperationTerminalStateBlankKeepsObservedState(t *testing.T) {
// A terminal observation that carries no status must keep the last known
// state rather than wipe it: operation.state may never be emptied by a
// close transition.
result := WithOutcome(pendingAcceptedResult(), OutcomeSuccess, WithOperationTerminalState(" "))
op := result.envelope().Meta.Operation
if op.State != "NEW" || op.TimedOut {
t.Fatalf("operation=%+v, want original state kept and timed_out cleared", op)
}
if err := ValidateResult(result); err != nil {
t.Fatal(err)
}
}
func TestWithOperationTerminalStateWithoutOperationInfoLeavesEnvelopeUntouched(t *testing.T) {
result := WithOutcome(Success(map[string]any{"ok": true}), OutcomeSuccess, WithOperationTerminalState("COMPLETED"))
env := result.envelope()
if env.Meta != nil && env.Meta.Operation != nil {
t.Fatalf("operation=%+v, want untouched", env.Meta.Operation)
}
}
+273
View File
@@ -0,0 +1,273 @@
// Copyright 2026 Alibaba Group
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
// Package wait is the framework terminal-state wait engine behind the
// reviewed contract.WaitSpec capability. It owns polling cadence, status
// extraction, and status→outcome mapping; it knows nothing about Cobra, MCP,
// or any product backend. How one poll executes is supplied by the leaf's
// WaitPoll hook (corecmd), so "poll = an existing read command" stays a leaf
// decision rather than a framework assumption.
package wait
import (
"context"
"errors"
"fmt"
"strconv"
"strings"
"time"
"github.com/DingTalk-Real-AI/dingtalk-workspace-cli/internal/corecmd/contract"
)
// DefaultPollInterval is the cadence between polls when LoopSpec.Interval is
// zero. The first poll runs immediately so an already-terminal resource does
// not pay a sleep tax.
const DefaultPollInterval = 2 * time.Second
// MaxPollInterval caps the exponential backoff growth between polls so a long
// wait cannot degenerate into effectively-blind polling.
const MaxPollInterval = 30 * time.Second
// PollDoc is one decoded poll response document (typically the unified-output
// envelope data of the poll command).
type PollDoc map[string]any
// Poller executes one poll. Returning an error fails the wait phase; the
// engine never retries a poller error because read commands failing is a real
// failure, not a "not yet" signal.
type Poller func(ctx context.Context) (PollDoc, error)
// LoopSpec is the runtime-resolved projection of contract.WaitSpec plus the
// caller-provided timeout.
type LoopSpec struct {
StatusQuery string
Terminal map[string]contract.ResultOutcome
Pending []string
Timeout time.Duration
Interval time.Duration
}
// Outcome is the closed result of a wait loop. TimedOut reports deadline
// exhaustion (Outcome is then pending — an accepted-but-not-terminal state is
// not a process failure per the exit-code contract); Status is the last
// observed status value.
type Outcome struct {
Status string
Outcome contract.ResultOutcome
Attempts int
TimedOut bool
}
// ErrUnknownStatus reports a status value that is neither declared terminal
// nor declared pending. Unknown fails closed: mapping it to pending could
// hide a real state change until timeout, mapping it to success is worse.
type ErrUnknownStatus struct {
Status string
Query string
}
func (e *ErrUnknownStatus) Error() string {
return fmt.Sprintf("wait: status %q (from %q) is neither terminal nor pending", e.Status, e.Query)
}
// Run polls poller until a declared terminal status, deadline exhaustion, or
// a poller error. The first poll is immediate; subsequent polls back off
// exponentially (×1.5) from Interval, capped at MaxPollInterval. Deadline
// exhaustion anywhere — before a poll, during a poll (a context-aware poller
// returns ctx.Err()), or during the wait between polls — always closes as
// timed-out pending with the last observed status, never as a poll failure.
func Run(ctx context.Context, spec LoopSpec, poll Poller) (Outcome, error) {
if spec.Interval <= 0 {
spec.Interval = DefaultPollInterval
}
if spec.Timeout > 0 {
var cancel context.CancelFunc
ctx, cancel = context.WithTimeout(ctx, spec.Timeout)
defer cancel()
}
pending := make(map[string]bool, len(spec.Pending))
for _, value := range spec.Pending {
pending[value] = true
}
timedOut := func(status string, attempts int) Outcome {
return Outcome{Status: status, Outcome: contract.ResultOutcomePending, Attempts: attempts, TimedOut: true}
}
interval := spec.Interval
attempts := 0
lastStatus := ""
for {
if ctx.Err() != nil {
return timedOut(lastStatus, attempts), nil
}
doc, err := poll(ctx)
if err != nil {
if ctx.Err() != nil {
return timedOut(lastStatus, attempts), nil
}
return Outcome{Attempts: attempts}, fmt.Errorf("wait: poll failed: %w", err)
}
attempts++
status, ok := ExtractStatus(doc, spec.StatusQuery)
if !ok {
return Outcome{Attempts: attempts}, fmt.Errorf(
"wait: status query %q not found in poll result", spec.StatusQuery)
}
lastStatus = status
if outcome, ok := spec.Terminal[status]; ok {
return Outcome{Status: status, Outcome: outcome, Attempts: attempts}, nil
}
if !pending[status] {
return Outcome{Status: status, Attempts: attempts}, &ErrUnknownStatus{Status: status, Query: spec.StatusQuery}
}
timer := time.NewTimer(interval)
select {
case <-ctx.Done():
timer.Stop()
return timedOut(status, attempts), nil
case <-timer.C:
}
interval = nextInterval(interval)
}
}
func nextInterval(current time.Duration) time.Duration {
next := current * 3 / 2
if next > MaxPollInterval {
next = MaxPollInterval
}
return next
}
// ExtractStatus resolves a dotted status query against a poll document. Each
// segment walks one map level; array indexes are not supported because wait
// targets a single resource. Numeric segments are stringified, so a document
// decoded with json.Number keys still resolves.
func ExtractStatus(doc PollDoc, query string) (string, bool) {
query = strings.TrimSpace(query)
if query == "" {
return "", false
}
// PollDoc is a defined type, so its dynamic type does not satisfy a
// map[string]any assertion — convert once at the boundary; nested values
// from JSON decoding are plain maps.
var current any = map[string]any(doc)
for _, segment := range strings.Split(query, ".") {
segment = strings.TrimSpace(segment)
if segment == "" {
return "", false
}
node, ok := current.(map[string]any)
if !ok {
return "", false
}
value, ok := node[segment]
if !ok {
return "", false
}
current = value
}
switch value := current.(type) {
case string:
return value, true
case fmt.Stringer:
return value.String(), true
case bool:
return strconv.FormatBool(value), true
case int:
return strconv.Itoa(value), true
case int64:
return strconv.FormatInt(value, 10), true
case float64:
return strconv.FormatFloat(value, 'f', -1, 64), true
default:
return "", false
}
}
// IsUnknownStatus reports whether err is the closed fail-on-unknown error.
func IsUnknownStatus(err error) bool {
var unknown *ErrUnknownStatus
return errors.As(err, &unknown)
}
// EventStream is the leaf-owned push subscription consumed by the event
// phase (the WaitEvents hook in corecmd). Recv delivers the next decoded
// event document; it returns an error or io.EOF-style termination when the
// stream ends — the engine treats non-terminal termination as a stream
// failure the caller (auto mode) may fall back from.
type EventStream interface {
Recv(ctx context.Context) (PollDoc, error)
}
// EventLoopSpec is the event-phase projection of contract.WaitSpec.
type EventLoopSpec struct {
StatusQuery string
MatchField string
Terminal map[string]contract.ResultOutcome
Pending []string
Timeout time.Duration
}
// ErrEventStreamEnded reports a stream that terminated before a terminal
// status. Auto mode uses it to fall back to polling; strict event mode
// surfaces it as a wait failure.
var ErrEventStreamEnded = errors.New("wait: event stream ended before a terminal status")
// RunEvent consumes stream until a correlated event reaches a declared
// terminal status, the deadline exhausts, or the stream ends. Events whose
// MatchField value does not equal resource are ignored (other resources on
// the same channel); a correlated event with an unknown status fails closed
// exactly like a poll would.
func RunEvent(ctx context.Context, spec EventLoopSpec, resource string, stream EventStream) (Outcome, error) {
if spec.Timeout > 0 {
var cancel context.CancelFunc
ctx, cancel = context.WithTimeout(ctx, spec.Timeout)
defer cancel()
}
pending := make(map[string]bool, len(spec.Pending))
for _, value := range spec.Pending {
pending[value] = true
}
attempts := 0
lastStatus := ""
for {
doc, err := stream.Recv(ctx)
if err != nil {
if ctx.Err() != nil {
return Outcome{Status: lastStatus, Outcome: contract.ResultOutcomePending, Attempts: attempts, TimedOut: true}, nil
}
// Wrap with ErrEventStreamEnded so auto mode can fall back to
// polling while correlated-status failures (unknown status,
// missing status query) stay non-recoverable.
return Outcome{Attempts: attempts}, fmt.Errorf("%w: %v", ErrEventStreamEnded, err)
}
attempts++
correlated, ok := ExtractStatus(doc, spec.MatchField)
if !ok || correlated != resource {
continue
}
status, ok := ExtractStatus(doc, spec.StatusQuery)
if !ok {
return Outcome{Attempts: attempts}, fmt.Errorf(
"wait: status query %q not found in event document", spec.StatusQuery)
}
lastStatus = status
if outcome, ok := spec.Terminal[status]; ok {
return Outcome{Status: status, Outcome: outcome, Attempts: attempts}, nil
}
if !pending[status] {
return Outcome{Status: status, Attempts: attempts}, &ErrUnknownStatus{Status: status, Query: spec.StatusQuery}
}
}
}
+365
View File
@@ -0,0 +1,365 @@
// 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)
}
}