fix(event): address retry storm review

This commit is contained in:
wxianfeng
2026-07-30 16:08:19 +08:00
parent ec99654854
commit 2808e71cb6
19 changed files with 446 additions and 138 deletions
+44 -13
View File
@@ -92,41 +92,72 @@ Ownership-based cleanup:
- T4c (control): `kill -9` leaves subscribe_id lingering (documented risk;
we only guarantee SIGTERM is clean, we do not fix kill -9 itself).
### 5. Subscription-create retry budget
### 5. Subscription-create retry orchestration and local guard
This policy covers all 16 public personal-event keys and every logical
subscription in a multi-event command. It applies only before the ready
marker; reconnecting an established Stream remains a separate mechanism.
- The `0/2/1` limits below are an **Agent/host orchestration contract**, not
a CLI-enforced persisted total-attempt cap. Each `dws event consume`
process sends at most one subscription-create HTTP request for a logical
subscription and performs no in-process automatic retry. The CLI persists
only the `in_flight`, `cooldown`, and `terminal_hold` guard states; it does
not persist or enforce the Agent/host attempt count across invocations.
- ID resolution, `event consume`, and later `event status/stop` must use the
same `--profile`. A user or conversation ID resolved under another profile
must not be reused for the current subscription.
- A logical subscription is keyed by the current profile/identity, event key,
rule type, target, and filters. A new `subscribe_id`, `trace_id`, or process
does not reset its retry budget.
- `retryable=false` means `max_additional_attempts=0`.
- `retryable=true` means `max_additional_attempts=2`. A caller must honor
`retry_after_seconds` or `next_retry_at` when present and must not retry
early.
- An omitted retryable value (`retryable=unknown`) means
`max_additional_attempts=1`; a second unknown failure stops the operation.
does not create a new logical operation or reset the Agent/host budget.
- For the Agent/host, `retryable=false` means
`max_additional_attempts=0`.
- For the Agent/host, `retryable=true` means
`max_additional_attempts=2`. It must honor `retry_after_seconds` or
`next_retry_at` when present and must not retry early.
- For the Agent/host, an omitted retryable value
(`retryable=unknown`) means `max_additional_attempts=1`; a second unknown
failure stops the operation.
- `in_flight` means the original logical request is still running.
`cooldown` and `terminal_hold` mean a guard is already delaying or blocking
it. These states must not recursively launch `event consume`, start a
parallel equivalent subscription, or reset the budget with a new subId or
trace. The caller waits for the original request/guard or stops.
parallel equivalent subscription, or bypass the guard with a new subId or
trace. The caller waits for the original request/guard or stops, while the
Agent/host keeps its own orchestration count.
- A multi-event command remains one original operation. A caller must not
split out a failed event, reorder events, or restart the command to bypass
a budget. Existing startup rollback cleans subscriptions created before a
later item fails.
#### Local guard state operations
- The default open-edition state file is
`~/.dws/events/open/personal_stream/<identity_hash>/personal_subscription_attempts.json`.
The config root follows `DWS_CONFIG_DIR` when set, and another edition uses
that edition's directory instead of `open`.
- The identity directory is mode `0700`; both
`personal_subscription_attempts.json` and
`personal_subscription_attempts.lock` are mode `0600`.
- A failure streak resets after 24h without another failure. A
`terminal_hold` lasts 1h. Prefer waiting until the reported
`next_retry_at`; do not clear the file as a normal retry mechanism.
- For emergency recovery, first ensure that no subscription-create process is
running for that identity. Delete only
`personal_subscription_attempts.json`, never the lock file. This clears
every protection record for that identity, not just one event.
**Verification**
- T5a: non-retryable, retryable, and unknown failures allow respectively
0, 2, and 1 additional attempts for the same logical subscription.
- T5b: a changed subId/trace or process restart does not increase the budget.
- T5a (policy): skill/docs tests pin the Agent/host 0/2/1 orchestration
contract and explicitly reject describing it as a CLI-persisted hard cap.
- T5b (CLI): one process issues at most one create request per logical
subscription; a changed subId/trace or process restart does not bypass the
persisted fingerprint guard.
- T5c: `in_flight`/`cooldown` does not recursively issue another create.
- T5d: multi-event startup cannot be split or reordered to bypass the guard,
and a partial startup still rolls back earlier subscriptions.
- T5e: state-store tests cover `0700`/`0600` permissions, 24h reset, 1h
`terminal_hold`, and identity-scoped cleanup; skill/docs tests pin the
operational recovery instructions.
## Out of scope (next branch)
+25 -2
View File
@@ -194,6 +194,12 @@ func (r *personalSubscriptionAttemptReservation) completeFailure(
if r == nil {
return cause
}
if r.store == nil || r.claim == nil {
return personalSubscriptionGuardError(errors.Join(
cause,
errors.New("personal event: subscription attempt reservation is incomplete"),
))
}
if failedIndex < 0 || failedIndex >= len(r.items) ||
succeededCount < 0 || succeededCount > failedIndex {
return personalSubscriptionGuardError(errors.Join(
@@ -262,9 +268,18 @@ func classifyPersonalSubscriptionFailure(err error, now time.Time) personalSubsc
apiErr.HTTPStatus >= http.StatusInternalServerError:
classification.retryability = personal.RetryabilityRetryable
classification.reason = "personal_subscription_transient_http"
case apiErr.HTTPStatus == http.StatusUnauthorized ||
apiErr.HTTPStatus == http.StatusForbidden:
classification.retryability = personal.RetryabilityNonRetryable
classification.reason = "personal_subscription_auth"
case personalSubscriptionTerminalBusinessCode(apiErr.Code):
classification.retryability = personal.RetryabilityNonRetryable
classification.reason = "personal_subscription_business_rejected"
case personalSubscriptionErrorHasSubscribeID(apiErr):
// A few legacy/proxy error shapes include an existing subscription
// ID without a stable server contract. Keep the response as an
// error, but do not turn that unverified shape into a one-hour hold.
classification.reason = "personal_subscription_unverified_existing_id"
case apiErr.HTTPStatus >= http.StatusBadRequest:
classification.retryability = personal.RetryabilityNonRetryable
classification.reason = "personal_subscription_http_rejected"
@@ -308,6 +323,14 @@ func classifyPersonalSubscriptionFailure(err error, now time.Time) personalSubsc
return classification
}
func personalSubscriptionErrorHasSubscribeID(apiErr *personal.APIError) bool {
if apiErr == nil {
return false
}
subscribeID, ok := apiErr.Details["subscribe_id"].(string)
return ok && strings.TrimSpace(subscribeID) != ""
}
func personalAPIRetryDelay(apiErr *personal.APIError, now time.Time) time.Duration {
if apiErr == nil {
return 0
@@ -521,9 +544,9 @@ func personalSubscriptionValidationError(cause error) error {
)
}
func personalSubscriptionLocalFailure(cause error) personalSubscriptionFailureClass {
func personalSubscriptionLocalFailure() personalSubscriptionFailureClass {
return personalSubscriptionFailureClass{
retryability: personal.RetryabilityNonRetryable,
retryability: personal.RetryabilityUnknown,
reason: "personal_subscription_local_failure",
}
}
+174 -19
View File
@@ -96,7 +96,7 @@ func (s *personalRecordingAttemptStore) Release(*personal.AttemptClaim) error {
return nil
}
func TestPersonalSubscriptionProtectionCoversAllPublicEvents(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionProtectionCoversAllPublicEvents(t *testing.T) {
oldFactory := personalNewSubscriptionAttemptStore
t.Cleanup(func() { personalNewSubscriptionAttemptStore = oldFactory })
@@ -162,7 +162,7 @@ func TestPersonalSubscriptionProtectionCoversAllPublicEvents(t *testing.T) {
}
}
func TestPersonalSubscriptionFingerprintChangesWithLogicalInputs(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionFingerprintChangesWithLogicalInputs(t *testing.T) {
identity := personal.Identity{
CorpID: "corp", UserID: "self", ClientID: "client", SourceID: "source",
}
@@ -209,7 +209,7 @@ func TestPersonalSubscriptionFingerprintChangesWithLogicalInputs(t *testing.T) {
}
}
func TestPersonalSubscriptionRejectsMalformedEndpointBeforeClaim(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionRejectsMalformedEndpointBeforeClaim(t *testing.T) {
oldFactory := personalNewSubscriptionAttemptStore
t.Cleanup(func() { personalNewSubscriptionAttemptStore = oldFactory })
factoryCalls := 0
@@ -247,7 +247,7 @@ func (fn personalRoundTripFunc) RoundTrip(request *http.Request) (*http.Response
return fn(request)
}
func TestPersonalSubscriptionProtectionConcurrentHundredAllowsOneHTTPRequest(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionProtectionConcurrentHundredAllowsOneHTTPRequest(t *testing.T) {
oldFactory := personalNewSubscriptionAttemptStore
t.Cleanup(func() { personalNewSubscriptionAttemptStore = oldFactory })
personalNewSubscriptionAttemptStore = func(workDir string) personalSubscriptionAttemptStore {
@@ -347,7 +347,7 @@ func TestPersonalSubscriptionProtectionConcurrentHundredAllowsOneHTTPRequest(t *
}
}
func TestPersonalSubscriptionFailureClassification(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionFailureClassification(t *testing.T) {
retryable, nonRetryable := true, false
now := time.Date(2026, 7, 30, 10, 0, 0, 0, time.UTC)
networkErr := &url.Error{
@@ -366,7 +366,9 @@ func TestPersonalSubscriptionFailureClassification(t *testing.T) {
{
name: "explicit false",
err: &personal.APIError{
Code: "GROUP_NOT_BELONG_TO_ORG", Retryable: &nonRetryable,
Code: "GROUP_NOT_BELONG_TO_ORG",
Retryable: &nonRetryable,
Details: map[string]any{"subscribe_id": "sub-existing"},
},
want: personal.RetryabilityNonRetryable,
wantReason: "personal_subscription_server_non_retryable",
@@ -382,7 +384,9 @@ func TestPersonalSubscriptionFailureClassification(t *testing.T) {
{
name: "terminal business code",
err: &personal.APIError{
Code: "USER_NOT_FOUND", HTTPStatus: http.StatusOK,
Code: "USER_NOT_FOUND",
HTTPStatus: http.StatusOK,
Details: map[string]any{"subscribe_id": "sub-existing"},
},
want: personal.RetryabilityNonRetryable,
wantReason: "personal_subscription_business_rejected",
@@ -396,6 +400,17 @@ func TestPersonalSubscriptionFailureClassification(t *testing.T) {
wantReason: "personal_subscription_business_rejected",
wantAuth: true,
},
{
name: "403 with existing id remains an auth rejection",
err: &personal.APIError{
Code: "PROXY_AUTH_FAILURE",
HTTPStatus: http.StatusForbidden,
Details: map[string]any{"subscribe_id": "sub-existing"},
},
want: personal.RetryabilityNonRetryable,
wantReason: "personal_subscription_auth",
wantAuth: true,
},
{
name: "429",
err: &personal.APIError{
@@ -452,6 +467,46 @@ func TestPersonalSubscriptionFailureClassification(t *testing.T) {
want: personal.RetryabilityUnknown,
wantReason: "personal_subscription_unknown",
},
{
name: "legacy DUP with existing id remains an unknown error",
err: &personal.APIError{
Code: "DUP",
HTTPStatus: http.StatusBadRequest,
Details: map[string]any{"subscribe_id": "sub-existing"},
},
want: personal.RetryabilityUnknown,
wantReason: "personal_subscription_unverified_existing_id",
},
{
name: "legacy SUBSCRIPTION_ALREADY_EXIST with existing id remains an unknown error",
err: &personal.APIError{
Code: "SUBSCRIPTION_ALREADY_EXIST",
HTTPStatus: http.StatusBadRequest,
Details: map[string]any{"subscribe_id": "sub-existing"},
},
want: personal.RetryabilityUnknown,
wantReason: "personal_subscription_unverified_existing_id",
},
{
name: "legacy ALREADY_SUBSCRIBED with existing id remains an unknown error",
err: &personal.APIError{
Code: "ALREADY_SUBSCRIBED",
HTTPStatus: http.StatusBadRequest,
Details: map[string]any{"subscribe_id": "sub-existing"},
},
want: personal.RetryabilityUnknown,
wantReason: "personal_subscription_unverified_existing_id",
},
{
name: "legacy DUPLICATE with existing id remains an unknown error",
err: &personal.APIError{
Code: "DUPLICATE",
HTTPStatus: http.StatusBadRequest,
Details: map[string]any{"subscribe_id": "sub-existing"},
},
want: personal.RetryabilityUnknown,
wantReason: "personal_subscription_unverified_existing_id",
},
{
name: "network",
err: networkErr,
@@ -503,7 +558,7 @@ func TestPersonalSubscriptionFailureClassification(t *testing.T) {
}
}
func TestPersonalSubscriptionFailureErrorPreservesTriStateAndDiagnostics(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionFailureErrorPreservesTriStateAndDiagnostics(t *testing.T) {
now := time.Date(2026, 7, 30, 10, 0, 0, 0, time.UTC)
hold := personal.AttemptHold{
State: personal.AttemptStateTerminalHold,
@@ -552,7 +607,7 @@ func TestPersonalSubscriptionFailureErrorPreservesTriStateAndDiagnostics(t *test
}
}
func TestPersonalSubscriptionBatchClaimFailureMakesZeroCreateCalls(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionBatchClaimFailureMakesZeroCreateCalls(t *testing.T) {
restore := installPersonalManySeams(t)
defer restore()
t.Setenv("DWS_CONFIG_DIR", t.TempDir())
@@ -607,7 +662,7 @@ func TestPersonalSubscriptionBatchClaimFailureMakesZeroCreateCalls(t *testing.T)
}
}
func TestPersonalSubscriptionAttemptGuardBypassesNonCreatePaths(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionAttemptGuardBypassesNonCreatePaths(t *testing.T) {
oldFactory := personalNewSubscriptionAttemptStore
oldIdentity := personalResolveEventIdentity
oldEnsure := personalEnsureSubscription
@@ -756,7 +811,7 @@ func TestPersonalSubscriptionAttemptGuardBypassesNonCreatePaths(t *testing.T) {
assertNoAttemptClaim("stop")
}
func TestPersonalSubscriptionLocalValidationRunsBeforeClaimAndCreate(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionLocalValidationRunsBeforeClaimAndCreate(t *testing.T) {
restore := installPersonalManySeams(t)
defer restore()
t.Setenv("DWS_CONFIG_DIR", t.TempDir())
@@ -794,9 +849,52 @@ func TestPersonalSubscriptionLocalValidationRunsBeforeClaimAndCreate(t *testing.
if recording.claimCalls != 0 || createCalls != 0 {
t.Fatalf("claim calls = %d, create calls = %d", recording.claimCalls, createCalls)
}
err = runPersonalEventConsumeSingle(
newPersonalCoverageCommand(),
personalConsumeOptions{
EventKey: personal.EventMention,
Flatten: true,
DebugRawEvents: true,
},
)
if err == nil {
t.Fatal("conflicting output modes succeeded")
}
if recording.claimCalls != 0 || createCalls != 0 {
t.Fatalf("output validation reached claim/create: %d/%d", recording.claimCalls, createCalls)
}
err = runPersonalEventConsumeSingle(
newPersonalCoverageCommand(),
personalConsumeOptions{
EventKey: personal.EventInChat,
Common: commonConsumeOptions{DryRun: true},
},
)
if err == nil {
t.Fatal("invalid dry-run subscription options succeeded")
}
if recording.claimCalls != 0 || createCalls != 0 {
t.Fatalf("dry-run validation reached claim/create: %d/%d", recording.claimCalls, createCalls)
}
claimErr := errors.New("claim failed")
recording.claimErr = claimErr
personalValidateConsumeConfig = func(consume.Config) error { return nil }
err = runPersonalEventConsumeSingle(
newPersonalCoverageCommand(),
personalConsumeOptions{EventKey: personal.EventMention},
)
if !errors.Is(err, claimErr) {
t.Fatalf("claim error = %v, want %v", err, claimErr)
}
if recording.claimCalls != 1 || createCalls != 0 {
t.Fatalf("failed claim calls = %d, create calls = %d", recording.claimCalls, createCalls)
}
}
func TestPersonalSubscriptionSingleAttemptCompletionErrors(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionSingleAttemptCompletionErrors(t *testing.T) {
tests := []struct {
name string
cancelContext bool
@@ -808,6 +906,7 @@ func TestPersonalSubscriptionSingleAttemptCompletionErrors(t *testing.T) {
wantComplete int
wantDelete int
wantCanceledGC bool
wantUnknown bool
}{
{
name: "nil subscription records failure",
@@ -820,6 +919,13 @@ func TestPersonalSubscriptionSingleAttemptCompletionErrors(t *testing.T) {
wantComplete: 1,
wantDelete: 1,
},
{
name: "local state failure records unknown cooldown",
upsertErr: errors.New("save state failed"),
wantFailure: 1,
wantDelete: 1,
wantUnknown: true,
},
{
name: "canceled local failure releases and uses canceled cleanup context",
cancelContext: true,
@@ -905,11 +1011,20 @@ func TestPersonalSubscriptionSingleAttemptCompletionErrors(t *testing.T) {
deleteCalls,
)
}
if test.wantUnknown {
if recording.lastFailure.Retryability != personal.RetryabilityUnknown {
t.Fatalf("local failure retryability = %q", recording.lastFailure.Retryability)
}
var typed *apperrors.Error
if !errors.As(err, &typed) || typed.RetryableSet {
t.Fatalf("local failure error = %#v, %v", typed, err)
}
}
})
}
}
func TestPersonalSubscriptionManyCompletionFailureCleansBatch(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionManyCompletionFailureCleansBatch(t *testing.T) {
restore := installPersonalManySeams(t)
defer restore()
t.Setenv("DWS_CONFIG_DIR", t.TempDir())
@@ -971,7 +1086,7 @@ func TestPersonalSubscriptionManyCompletionFailureCleansBatch(t *testing.T) {
}
}
func TestPersonalSubscriptionRejectsNonPublicEventsOnChangedPaths(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionRejectsNonPublicEventsOnChangedPaths(t *testing.T) {
const nonPublicEvent = personal.EventMention
oldLookup := personalLookupDefinition
oldGet := personalGetSubscription
@@ -1024,7 +1139,7 @@ func TestPersonalSubscriptionRejectsNonPublicEventsOnChangedPaths(t *testing.T)
}
}
func TestPersonalSubscriptionAttemptLeaseBounds(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionAttemptLeaseBounds(t *testing.T) {
client := personal.NewClient("https://mcp.example.test/dws", personal.Identity{})
tests := []struct {
timeout time.Duration
@@ -1045,7 +1160,7 @@ func TestPersonalSubscriptionAttemptLeaseBounds(t *testing.T) {
}
}
func TestPersonalSubscriptionAttemptGuardHelperEdges(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionAttemptGuardHelperEdges(t *testing.T) {
oldFactory := personalNewSubscriptionAttemptStore
t.Cleanup(func() { personalNewSubscriptionAttemptStore = oldFactory })
@@ -1132,6 +1247,18 @@ func TestPersonalSubscriptionAttemptGuardHelperEdges(t *testing.T) {
}
})
t.Run("invalid prepared subscription", func(t *testing.T) {
if _, err := reservePersonalSubscriptionAttempts(
t.TempDir(),
validClient,
identity,
"profile-a",
[]personalConsumeOptions{{EventKey: personal.EventInChat}},
); err == nil {
t.Fatal("invalid prepared subscription was accepted")
}
})
t.Run("incomplete success reservation", func(t *testing.T) {
if err := (&personalSubscriptionAttemptReservation{}).completeSuccess(); err == nil {
t.Fatal("incomplete reservation succeeded")
@@ -1166,6 +1293,8 @@ func TestPersonalSubscriptionAttemptGuardHelperEdges(t *testing.T) {
t.Run("invalid failure indexes", func(t *testing.T) {
wantErr := errors.New("create failed")
reservation := &personalSubscriptionAttemptReservation{
store: &personalRecordingAttemptStore{},
claim: &personal.AttemptClaim{AttemptID: "attempt"},
items: []personalSubscriptionAttemptItem{{fingerprint: "fingerprint"}},
}
if err := reservation.completeFailure(
@@ -1179,6 +1308,24 @@ func TestPersonalSubscriptionAttemptGuardHelperEdges(t *testing.T) {
}
})
t.Run("incomplete cancellation reservation", func(t *testing.T) {
wantErr := context.Canceled
err := (&personalSubscriptionAttemptReservation{}).completeFailure(
context.Background(),
0,
0,
wantErr,
nil,
)
if !errors.Is(err, wantErr) {
t.Fatalf("incomplete-cancellation error = %v, want joined cause %v", err, wantErr)
}
var typed *apperrors.Error
if !errors.As(err, &typed) || typed.Reason != "personal_subscription_guard_failed" {
t.Fatalf("incomplete-cancellation guard error = %#v, %v", typed, err)
}
})
t.Run("failure completion failure", func(t *testing.T) {
createErr := errors.New("create failed")
completeErr := errors.New("complete failure failed")
@@ -1203,7 +1350,7 @@ func TestPersonalSubscriptionAttemptGuardHelperEdges(t *testing.T) {
})
}
func TestPersonalSubscriptionFailureClassificationEdges(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionFailureClassificationEdges(t *testing.T) {
now := time.Date(2026, 7, 30, 10, 0, 0, 0, time.UTC)
tests := []struct {
@@ -1248,6 +1395,14 @@ func TestPersonalSubscriptionFailureClassificationEdges(t *testing.T) {
if delay := personalAPIRetryDelay(nil, now); delay != 0 {
t.Fatalf("nil API retry delay = %s", delay)
}
if personalSubscriptionErrorHasSubscribeID(nil) {
t.Fatal("nil API error reported a subscription ID")
}
if personalSubscriptionErrorHasSubscribeID(&personal.APIError{
Details: map[string]any{"subscribe_id": 42},
}) {
t.Fatal("non-string subscription ID was accepted")
}
httpDate := now.Add(75 * time.Second).Format(http.TimeFormat)
if delay := personalAPIRetryDelay(&personal.APIError{
Details: map[string]any{"retry_after": httpDate},
@@ -1269,7 +1424,7 @@ func TestPersonalSubscriptionFailureClassificationEdges(t *testing.T) {
}
}
func TestPersonalSubscriptionErrorConstructionEdges(t *testing.T) {
func TestCrossPlatformCoveragePersonalSubscriptionErrorConstructionEdges(t *testing.T) {
nonRetryable := personalSubscriptionFailureClass{
retryability: personal.RetryabilityNonRetryable,
reason: "personal_subscription_auth",
@@ -1320,7 +1475,7 @@ func TestPersonalSubscriptionErrorConstructionEdges(t *testing.T) {
}
}
func TestNewPersonalEventControlClientSetsVersionHeaders(t *testing.T) {
func TestCrossPlatformCoverageNewPersonalEventControlClientSetsVersionHeaders(t *testing.T) {
oldVersion := version
t.Cleanup(func() { version = oldVersion })
version = "1.2.3-test"
+2 -2
View File
@@ -393,7 +393,7 @@ func runPersonalEventConsumeSingle(c *cobra.Command, opts personalConsumeOptions
if personalSubscriptionCanceled(ctx, wrapped) {
cleanupCtx = ctx
}
classification := personalSubscriptionLocalFailure(wrapped)
classification := personalSubscriptionLocalFailure()
wrapped = attempt.completeFailure(ctx, 0, 0, wrapped, &classification)
cleanup(cleanupCtx)
}
@@ -595,7 +595,7 @@ func runPersonalEventConsumeMany(c *cobra.Command, opts personalConsumeOptions)
IdentityHash: identityHash,
}); err != nil {
cause := fmt.Errorf("save run state for %s: %w", eventKey, err)
classification := personalSubscriptionLocalFailure(cause)
classification := personalSubscriptionLocalFailure()
cause = failAndCleanup(i, len(created)-1, cause, &classification)
return fmt.Errorf("event consume --as user: %w", cause)
}
+5 -5
View File
@@ -252,7 +252,7 @@ func TestPreparePersonalMultiOptionsRejectsSingleOnlyFlags(t *testing.T) {
}
}
func TestEventConsumeMultiRejectsExplicitSingleOnlyFlagsEvenWhenEmpty(t *testing.T) {
func TestCrossPlatformCoverageEventConsumeMultiRejectsExplicitSingleOnlyFlagsEvenWhenEmpty(t *testing.T) {
oldRun := eventRunPersonalConsume
defer func() { eventRunPersonalConsume = oldRun }()
eventRunPersonalConsume = func(*cobra.Command, personalConsumeOptions) error {
@@ -380,7 +380,7 @@ func TestRunPersonalEventConsumeManyRollsBackPartialCreation(t *testing.T) {
}
}
func TestRunPersonalEventConsumeManyPersistsFailureBeforeRollback(t *testing.T) {
func TestCrossPlatformCoverageRunPersonalEventConsumeManyPersistsFailureBeforeRollback(t *testing.T) {
restore := installPersonalManySeams(t)
defer restore()
t.Setenv("DWS_CONFIG_DIR", t.TempDir())
@@ -429,7 +429,7 @@ func TestRunPersonalEventConsumeManyPersistsFailureBeforeRollback(t *testing.T)
}
}
func TestRunPersonalEventConsumeSinglePersistsLocalFailureBeforeRollback(t *testing.T) {
func TestCrossPlatformCoverageRunPersonalEventConsumeSinglePersistsLocalFailureBeforeRollback(t *testing.T) {
restore := installPersonalManySeams(t)
defer restore()
t.Setenv("DWS_CONFIG_DIR", t.TempDir())
@@ -475,7 +475,7 @@ func TestRunPersonalEventConsumeSinglePersistsLocalFailureBeforeRollback(t *test
}
}
func TestRunPersonalEventConsumeManyCancellationReleasesBeforeCanceledCleanup(t *testing.T) {
func TestCrossPlatformCoverageRunPersonalEventConsumeManyCancellationReleasesBeforeCanceledCleanup(t *testing.T) {
restore := installPersonalManySeams(t)
defer restore()
t.Setenv("DWS_CONFIG_DIR", t.TempDir())
@@ -531,7 +531,7 @@ func TestRunPersonalEventConsumeManyCancellationReleasesBeforeCanceledCleanup(t
}
}
func TestRunPersonalEventConsumeManyRejectsInvalidSubscriptionResults(t *testing.T) {
func TestCrossPlatformCoverageRunPersonalEventConsumeManyRejectsInvalidSubscriptionResults(t *testing.T) {
for _, test := range []struct {
name string
ensure func(int, personalConsumeOptions) *personal.Subscription
@@ -1,6 +1,6 @@
{
"version": 1,
"source_hash": "sha256:a4998e1d96d1e816e6978125958b69002b5c9f82e6dfc44a0971ba8003fc6e9b",
"source_hash": "sha256:7484b82d1ca793a6129acfaa192ff08fc0864d47a6e75ecd16d8e7fa32857e16",
"surface_hash": "sha256:60eee8e2f37d6d9d60689efce85082798eb9ad38b7ba7c0b471c3de676a85a16",
"coverage": {
"surface_products": 26,
@@ -1,6 +1,6 @@
{
"version": 1,
"source_hash": "sha256:a4998e1d96d1e816e6978125958b69002b5c9f82e6dfc44a0971ba8003fc6e9b",
"source_hash": "sha256:7484b82d1ca793a6129acfaa192ff08fc0864d47a6e75ecd16d8e7fa32857e16",
"surface_hash": "sha256:60eee8e2f37d6d9d60689efce85082798eb9ad38b7ba7c0b471c3de676a85a16",
"source_files": 160,
"hint_files": 54,
+2 -2
View File
@@ -1,12 +1,12 @@
{
"version": 1,
"surface_hash": "sha256:60eee8e2f37d6d9d60689efce85082798eb9ad38b7ba7c0b471c3de676a85a16",
"source_hash": "sha256:d1d19f794effe155a1691c02315546f58aee3860fc694fc0b4456868b0dea3af",
"source_hash": "sha256:ba691e70f1c232fc3378a743c3fcd34113960d30b46bffa7d1fc8b524115052c",
"catalog": {
"agent_metadata": {
"products_with_metadata": 26,
"source": "embedded-skill-metadata",
"source_hash": "sha256:a4998e1d96d1e816e6978125958b69002b5c9f82e6dfc44a0971ba8003fc6e9b",
"source_hash": "sha256:7484b82d1ca793a6129acfaa192ff08fc0864d47a6e75ecd16d8e7fa32857e16",
"surface_hash": "sha256:60eee8e2f37d6d9d60689efce85082798eb9ad38b7ba7c0b471c3de676a85a16",
"surface_products": 26,
"surface_tools": 845,
+2 -2
View File
@@ -78,7 +78,7 @@ func TestPrintJSON(t *testing.T) {
}
}
func TestRetryabilityTriStateAndRetryTiming(t *testing.T) {
func TestCrossPlatformCoverageRetryabilityTriStateAndRetryTiming(t *testing.T) {
t.Parallel()
next := time.Date(2026, time.July, 30, 4, 5, 6, 0, time.FixedZone("CST", 8*60*60))
@@ -158,7 +158,7 @@ func TestRetryabilityTriStateAndRetryTiming(t *testing.T) {
}
}
func TestRetryTimingOptionsIgnoreInvalidValues(t *testing.T) {
func TestCrossPlatformCoverageRetryTimingOptionsIgnoreInvalidValues(t *testing.T) {
t.Parallel()
err := NewAPI(
+21 -21
View File
@@ -67,7 +67,7 @@ func newAttemptTestStore(t *testing.T, clock *attemptTestClock) *AttemptStore {
)
}
func TestAttemptFingerprintIsCanonicalScopedAndOpaque(t *testing.T) {
func TestCrossPlatformCoverageAttemptFingerprintIsCanonicalScopedAndOpaque(t *testing.T) {
base := Fingerprint(" HTTPS://MCP.Example.Test/dws/ ", " idem-1 ")
if base == "" || len(base) != 64 {
t.Fatalf("Fingerprint() = %q", base)
@@ -90,7 +90,7 @@ func TestAttemptFingerprintIsCanonicalScopedAndOpaque(t *testing.T) {
}
}
func TestRetryabilityValuePreservesUnknown(t *testing.T) {
func TestCrossPlatformCoverageRetryabilityValuePreservesUnknown(t *testing.T) {
if value, known := RetryabilityRetryable.Value(); !known || !value {
t.Fatalf("retryable Value() = %v, %v", value, known)
}
@@ -117,7 +117,7 @@ func TestRetryabilityValuePreservesUnknown(t *testing.T) {
}
}
func TestAttemptStoreClaimFailureBackoffAndRetryAfter(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreClaimFailureBackoffAndRetryAfter(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("one")
@@ -183,7 +183,7 @@ func TestAttemptStoreClaimFailureBackoffAndRetryAfter(t *testing.T) {
}
}
func TestAttemptStoreBackoffSequenceAndTwentyFourHourReset(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreBackoffSequenceAndTwentyFourHourReset(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("backoff")
@@ -235,7 +235,7 @@ func TestAttemptStoreBackoffSequenceAndTwentyFourHourReset(t *testing.T) {
}
}
func TestAttemptStoreLongRetryAfterSurvivesFailureCountResetWindow(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreLongRetryAfterSurvivesFailureCountResetWindow(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("long-retry-after")
@@ -284,7 +284,7 @@ func TestAttemptStoreLongRetryAfterSurvivesFailureCountResetWindow(t *testing.T)
}
}
func TestAttemptStoreExpiredInFlightResetsOldFailureHistory(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreExpiredInFlightResetsOldFailureHistory(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("expired-inflight-reset")
@@ -328,7 +328,7 @@ func TestAttemptStoreExpiredInFlightResetsOldFailureHistory(t *testing.T) {
}
}
func TestAttemptStoreTerminalHoldAndSuccessClear(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreTerminalHoldAndSuccessClear(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("terminal")
@@ -378,7 +378,7 @@ func TestAttemptStoreTerminalHoldAndSuccessClear(t *testing.T) {
}
}
func TestAttemptStorePartialFailureClearsSuccessAndReleasesUnexecuted(t *testing.T) {
func TestCrossPlatformCoverageAttemptStorePartialFailureClearsSuccessAndReleasesUnexecuted(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprintA := attemptTestFingerprint("batch-a")
@@ -433,7 +433,7 @@ func TestAttemptStorePartialFailureClearsSuccessAndReleasesUnexecuted(t *testing
}
}
func TestAttemptStoreReleaseRestoresPreviousState(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreReleaseRestoresPreviousState(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprintA := attemptTestFingerprint("release-a")
@@ -475,7 +475,7 @@ func TestAttemptStoreReleaseRestoresPreviousState(t *testing.T) {
}
}
func TestAttemptStoreCASRejectsOldCompletion(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreCASRejectsOldCompletion(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("cas")
@@ -497,7 +497,7 @@ func TestAttemptStoreCASRejectsOldCompletion(t *testing.T) {
}
}
func TestAttemptStoreConcurrentClaimsOnlyOneWins(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreConcurrentClaimsOnlyOneWins(t *testing.T) {
workDir := t.TempDir()
fingerprint := attemptTestFingerprint("concurrent")
spec := AttemptSpec{Fingerprint: fingerprint, EventKey: EventMention}
@@ -539,7 +539,7 @@ func TestAttemptStoreConcurrentClaimsOnlyOneWins(t *testing.T) {
}
}
func TestAttemptStoreBatchClaimIsAtomicWhenOneItemBlocked(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreBatchClaimIsAtomicWhenOneItemBlocked(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprintA := attemptTestFingerprint("atomic-a")
@@ -568,7 +568,7 @@ func TestAttemptStoreBatchClaimIsAtomicWhenOneItemBlocked(t *testing.T) {
}
}
func TestAttemptStorePersistsNoRawFingerprintInputs(t *testing.T) {
func TestCrossPlatformCoverageAttemptStorePersistsNoRawFingerprintInputs(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
const (
@@ -598,7 +598,7 @@ func TestAttemptStorePersistsNoRawFingerprintInputs(t *testing.T) {
}
}
func TestAttemptStoreFailsClosedOnCorruptStateAndWriteFailure(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreFailsClosedOnCorruptStateAndWriteFailure(t *testing.T) {
t.Run("corrupt state", func(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
@@ -680,7 +680,7 @@ func TestAttemptStoreFailsClosedOnCorruptStateAndWriteFailure(t *testing.T) {
})
}
func TestAttemptStoreFailsClosedOnInjectedIOErrors(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreFailsClosedOnInjectedIOErrors(t *testing.T) {
testClaim := func(store *AttemptStore, label string) error {
_, err := store.Claim([]AttemptSpec{{
Fingerprint: attemptTestFingerprint(label),
@@ -773,7 +773,7 @@ func TestAttemptStoreFailsClosedOnInjectedIOErrors(t *testing.T) {
})
}
func TestAttemptStoreTightensStateAndLockPermissions(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreTightensStateAndLockPermissions(t *testing.T) {
if runtime.GOOS == "windows" {
t.Skip("Windows does not expose Unix permission bits")
}
@@ -811,7 +811,7 @@ func TestAttemptStoreTightensStateAndLockPermissions(t *testing.T) {
}
}
func TestAttemptStoreRejectsInvalidPersistedRecords(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreRejectsInvalidPersistedRecords(t *testing.T) {
fingerprint := attemptTestFingerprint("invalid-record")
longValue := strings.Repeat("x", attemptMaxFieldLength+1)
tests := []struct {
@@ -882,7 +882,7 @@ func TestAttemptStoreRejectsInvalidPersistedRecords(t *testing.T) {
}
}
func TestAttemptStoreLockTimeoutAndIDFailureAreFailClosed(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreLockTimeoutAndIDFailureAreFailClosed(t *testing.T) {
t.Run("lock timeout", func(t *testing.T) {
var tick atomic.Int64
base := time.Date(2026, 7, 30, 10, 0, 0, 0, time.UTC)
@@ -921,7 +921,7 @@ func TestAttemptStoreLockTimeoutAndIDFailureAreFailClosed(t *testing.T) {
})
}
func TestAttemptStoreValidationAndFailureCAS(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreValidationAndFailureCAS(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("validation")
@@ -987,7 +987,7 @@ func TestAttemptStoreValidationAndFailureCAS(t *testing.T) {
}
}
func TestAttemptStoreBoundsPersistedDiagnostics(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreBoundsPersistedDiagnostics(t *testing.T) {
clock := newAttemptTestClock()
store := newAttemptTestStore(t, clock)
fingerprint := attemptTestFingerprint("bounded")
@@ -1018,7 +1018,7 @@ func TestAttemptStoreBoundsPersistedDiagnostics(t *testing.T) {
}
}
func TestAttemptStoreChangedCoverageEdges(t *testing.T) {
func TestCrossPlatformCoverageAttemptStoreChangedCoverageEdges(t *testing.T) {
t.Run("endpoint fallback", func(t *testing.T) {
const raw = "https://example.test/%zz"
if got := normalizeAttemptEndpoint(raw); got != raw {
-21
View File
@@ -185,18 +185,6 @@ func (c *Client) CreateSubscription(ctx context.Context, req CreateSubscriptionR
}
var sub Subscription
if err := c.do(ctx, http.MethodPost, "/subscription/user", nil, c.buildCreateRequest(req), &sub); err != nil {
var apiErr *APIError
if errors.As(err, &apiErr) && isDuplicateSubscriptionCode(apiErr.Code) {
if subID, ok := apiErr.Details["subscribe_id"].(string); ok && subID != "" {
return &Subscription{
SubscribeID: subID,
EventKey: req.EventKey,
RuleType: req.RuleType,
Status: "active",
SourceID: c.Identity.SourceID,
}, nil
}
}
return nil, err
}
if sub.EventKey == "" {
@@ -809,15 +797,6 @@ func isNotFound(err error) bool {
return errors.As(err, &apiErr) && (apiErr.Code == "PERSONAL_EVENT_NOT_FOUND" || apiErr.Code == "NOT_FOUND")
}
func isDuplicateSubscriptionCode(code string) bool {
switch strings.ToUpper(strings.TrimSpace(code)) {
case "DUP", "DUPLICATE_SUBSCRIPTION", "SUBSCRIPTION_ALREADY_EXISTS":
return true
default:
return false
}
}
func dwsStatusString(raw json.RawMessage) string {
raw = bytes.TrimSpace(raw)
if len(raw) == 0 || string(raw) == "null" {
+40 -28
View File
@@ -297,7 +297,7 @@ func TestClientBusinessErrorHTTP200(t *testing.T) {
}
}
func TestClientBusinessErrorPreservesRetryContractAndClientHeaders(t *testing.T) {
func TestCrossPlatformCoverageClientBusinessErrorPreservesRetryContractAndClientHeaders(t *testing.T) {
nextRetryAt := "2026-07-30T04:05:06Z"
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if got := r.Header.Get("X-Cli-Version"); got != "1.2.3" {
@@ -347,7 +347,7 @@ func TestClientBusinessErrorPreservesRetryContractAndClientHeaders(t *testing.T)
}
}
func TestClientErrorDiagnosticIdentityPrecedence(t *testing.T) {
func TestCrossPlatformCoverageClientErrorDiagnosticIdentityPrecedence(t *testing.T) {
tests := []struct {
name string
body map[string]any
@@ -451,7 +451,7 @@ func TestClientErrorDiagnosticIdentityPrecedence(t *testing.T) {
}
}
func TestClientHTTPErrorPreservesRetryAfterAndHeaderTrace(t *testing.T) {
func TestCrossPlatformCoverageClientHTTPErrorPreservesRetryAfterAndHeaderTrace(t *testing.T) {
tests := []struct {
name string
retryAfter string
@@ -507,7 +507,7 @@ func TestClientHTTPErrorPreservesRetryAfterAndHeaderTrace(t *testing.T) {
}
}
func TestClientErrorDiagnosticHelpersHandleNilAndMalformedInputs(t *testing.T) {
func TestCrossPlatformCoverageClientErrorDiagnosticHelpersHandleNilAndMalformedInputs(t *testing.T) {
if got := withHTTPResponseDetails(nil, http.MethodPost, "/subscription/user", nil, "request-1", "trace-1"); got != nil {
t.Fatalf("withHTTPResponseDetails(nil) = %#v, want nil", got)
}
@@ -529,31 +529,43 @@ func TestClientErrorDiagnosticHelpersHandleNilAndMalformedInputs(t *testing.T) {
}
}
func TestClientCreateSubscriptionDoesNotTreatArbitraryErrorSubscribeIDAsSuccess(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
_ = json.NewEncoder(w).Encode(map[string]any{
"error": map[string]any{
"code": "SYSTEM_ERROR",
"message": "registration failed",
"details": map[string]any{"subscribe_id": "pending-sub"},
},
})
}))
defer srv.Close()
func TestCrossPlatformCoverageClientCreateSubscriptionRejectsErrorsWithSubscribeID(t *testing.T) {
for _, code := range []string{
"SYSTEM_ERROR",
"DUP",
"DUPLICATE_SUBSCRIPTION",
"SUBSCRIPTION_ALREADY_EXISTS",
"SUBSCRIPTION_ALREADY_EXIST",
"ALREADY_SUBSCRIBED",
"DUPLICATE",
} {
t.Run(code, func(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusInternalServerError)
_ = json.NewEncoder(w).Encode(map[string]any{
"error": map[string]any{
"code": code,
"message": "registration failed",
"details": map[string]any{"subscribe_id": "pending-sub"},
},
})
}))
defer srv.Close()
c := NewClient(srv.URL, Identity{AccessToken: "token", ClientID: "client", SourceID: "open"})
sub, err := c.CreateSubscription(t.Context(), CreateSubscriptionRequest{
EventKey: EventMention,
RuleType: "at",
RuleParam: map[string]any{},
})
if err == nil || sub != nil {
t.Fatalf("subscription = %#v, error = %v; arbitrary errors must stay failures", sub, err)
}
var apiErr *APIError
if !errors.As(err, &apiErr) || apiErr.Code != "SYSTEM_ERROR" {
t.Fatalf("error = %#v", err)
c := NewClient(srv.URL, Identity{AccessToken: "token", ClientID: "client", SourceID: "open"})
sub, err := c.CreateSubscription(t.Context(), CreateSubscriptionRequest{
EventKey: EventMention,
RuleType: "at",
RuleParam: map[string]any{},
})
if err == nil || sub != nil {
t.Fatalf("subscription = %#v, error = %v; API errors must stay failures", sub, err)
}
var apiErr *APIError
if !errors.As(err, &apiErr) || apiErr.Code != code {
t.Fatalf("error = %#v", err)
}
})
}
}
@@ -219,8 +219,8 @@ func TestCrossPlatformCoveragePersonalClientHelpersAndOperations(t *testing.T) {
}
base.HTTPClient = personalHTTPClient(400, `{"error":{"code":"DUP","details":{"subscribe_id":"existing"}}}`)
created, err := base.CreateSubscription(t.Context(), CreateSubscriptionRequest{EventKey: "e", RuleType: "r"})
if err != nil || created.SubscribeID != "existing" {
t.Fatalf("duplicate create = %#v, %v", created, err)
if err == nil || created != nil {
t.Fatalf("duplicate error create = %#v, %v", created, err)
}
request := base.buildCreateRequest(CreateSubscriptionRequest{
EventKey: "e", RuleType: "r", Name: "n", RuleParam: map[string]any{"bad": make(chan int)},
+1 -1
View File
@@ -662,7 +662,7 @@ func (r *devAppFailingCountRunner) Run(_ context.Context, invocation executor.In
return executor.Result{Invocation: invocation}, r.err
}
func TestDevAppEventSubscribeRunnerFailureIsNotRetried(t *testing.T) {
func TestCrossPlatformCoverageDevAppEventSubscribeRunnerFailureIsNotRetried(t *testing.T) {
wantErr := stderrors.New("event subscription failed")
runner := &devAppFailingCountRunner{err: wantErr}
root := newDevAppTestRoot(runner)
+1 -1
View File
@@ -292,7 +292,7 @@ func TestCallToolUsesJSONRPCMethod(t *testing.T) {
}
}
func TestCallToolDevAppEventSubscribeRetriesAreBounded(t *testing.T) {
func TestCrossPlatformCoverageCallToolDevAppEventSubscribeRetriesAreBounded(t *testing.T) {
attempts := 0
httpClient := &http.Client{
Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) {
+13 -5
View File
@@ -175,14 +175,22 @@ dws event stop --all --yes
以下约束适用于上表全部 16 个事件以及多事件命令中的每一项,只治理 `[event] ready` 之前的订阅创建;ready 之后的 Stream 断线由长连接重连机制处理。
- `0/2/1` 是 **Agent/host 编排约束**,不是 CLI 持久化硬总次数上限。每次 `dws event consume` 调用对每个逻辑订阅最多发送一次订阅创建 HTTP 请求,进程内不会自动重试。CLI 本地状态只持久化 `in_flight`、`cooldown`、`terminal_hold` 三种保护状态,不持久化或计算跨调用的 Agent/host 尝试次数。
- 解析人名或群名、执行 `event consume` 以及后续 `event status/stop` 必须使用同一个 `--profile`。不得把其它 profile 下解析出的 userId、openDingtalkId 或 openConversationId 直接带入当前 profile 的订阅。
- 同一逻辑订阅由当前 profile / 身份、event key、rule type、目标和过滤条件共同确定。`subscribe_id`、`trace_id` 以及重新启动进程都只是诊断或执行信息,不会生成新预算。
- `retryable=false`:`max_additional_attempts=0`,立即停止,不得自动重跑。
- `retryable=true`:`max_additional_attempts=2`,初次失败后最多再尝试 2 次;错误若给出 `retry_after_seconds` 或 `next_retry_at`,不得提前重试。
- 未返回 `retryable`(即 `retryable=unknown`):`max_additional_attempts=1`,最多补偿尝试 1 次;仍无法确认时停止并上报错误与 trace。
- `in_flight` 表示同一逻辑订阅已有请求执行中;`cooldown` 或 `terminal_hold` 表示当前被退避或终态保护。遇到这些状态不得递归调用 `event consume`、并行启动相同订阅或通过新 subId / trace 绕过;等待原请求或保护时间结束,并继续消耗原预算。
- 同一逻辑订阅由当前 profile / 身份、event key、rule type、目标和过滤条件共同确定。`subscribe_id`、`trace_id` 以及重新启动进程都只是诊断或执行信息,不会生成新的逻辑操作;Agent/host 必须自行延续原编排预算。
- Agent/host 收到 `retryable=false`:`max_additional_attempts=0`,立即停止,不得自动重跑。
- Agent/host 收到 `retryable=true`:`max_additional_attempts=2`,初次失败后最多再尝试 2 次;错误若给出 `retry_after_seconds` 或 `next_retry_at`,不得提前重试。
- 未返回 `retryable`(即 `retryable=unknown`):Agent/host 使用 `max_additional_attempts=1`,最多补偿尝试 1 次;仍无法确认时停止并上报错误与 trace。
- `in_flight` 表示同一逻辑订阅已有请求执行中;`cooldown` 或 `terminal_hold` 表示当前被退避或终态保护。遇到这些状态不得递归调用 `event consume`、并行启动相同订阅或通过新 subId / trace 绕过;等待原请求或保护时间结束,同时由 Agent/host 继续维护自己的编排次数。
- 多事件命令必须作为同一次原始操作治理。不得把失败事件拆成新的单事件命令、调整顺序或反复重启来重置预算;启动中任一项失败时,由 CLI 回滚本次已经创建的订阅。
### 本地保护状态运维
- open 版默认状态文件为 `~/.dws/events/open/personal_stream/<identity_hash>/personal_subscription_attempts.json`。设置 `DWS_CONFIG_DIR` 后配置根目录随之变化;其它 edition 使用对应 edition 目录,不固定为 `open`。
- identity 目录权限为 `0700`;`personal_subscription_attempts.json` 与 `personal_subscription_attempts.lock` 权限均为 `0600`。
- 连续 24h 没有失败后失败计数重置;`terminal_hold` 持续 1h。正常处理优先等待错误中的 `next_retry_at`,不要把删状态文件当成常规重试手段。
- 紧急恢复时,先确认该 identity 没有正在创建订阅的进程;只删除 `personal_subscription_attempts.json`,不要删除 lock 文件。删除 JSON 会清空该 identity 的全部保护记录,不只影响一个事件。
## Output parsing
- 推荐 `--flatten -f ndjson`:顶层业务字段,一行一个事件 JSON,适合 Agent 管道读取。
+13 -5
View File
@@ -78,14 +78,22 @@ description: 钉钉个人 IM 事件长连接监听、订阅与消费,覆盖消
以下约束适用于上表全部 16 个事件以及多事件命令中的每一项,只治理 `[event] ready` 之前的订阅创建;ready 之后的 Stream 断线由长连接重连机制处理。
- `0/2/1` 是 **Agent/host 编排约束**,不是 CLI 持久化硬总次数上限。每次 `dws event consume` 调用对每个逻辑订阅最多发送一次订阅创建 HTTP 请求,进程内不会自动重试。CLI 本地状态只持久化 `in_flight`、`cooldown`、`terminal_hold` 三种保护状态,不持久化或计算跨调用的 Agent/host 尝试次数。
- 解析人名或群名、执行 `event consume` 以及后续 `event status/stop` 必须使用同一个 `--profile`。不得把其它 profile 下解析出的 userId、openDingtalkId 或 openConversationId 直接带入当前 profile 的订阅。
- 同一逻辑订阅由当前 profile / 身份、event key、rule type、目标和过滤条件共同确定。`subscribe_id`、`trace_id` 以及重新启动进程都只是诊断或执行信息,不会生成新预算。
- `retryable=false`:`max_additional_attempts=0`,立即停止,不得自动重跑。
- `retryable=true`:`max_additional_attempts=2`,初次失败后最多再尝试 2 次;错误若给出 `retry_after_seconds` 或 `next_retry_at`,不得提前重试。
- 未返回 `retryable`(即 `retryable=unknown`):`max_additional_attempts=1`,最多补偿尝试 1 次;仍无法确认时停止并上报错误与 trace。
- `in_flight` 表示同一逻辑订阅已有请求执行中;`cooldown` 或 `terminal_hold` 表示当前被退避或终态保护。遇到这些状态不得递归调用 `event consume`、并行启动相同订阅或通过新 subId / trace 绕过;等待原请求或保护时间结束,并继续消耗原预算。
- 同一逻辑订阅由当前 profile / 身份、event key、rule type、目标和过滤条件共同确定。`subscribe_id`、`trace_id` 以及重新启动进程都只是诊断或执行信息,不会生成新的逻辑操作;Agent/host 必须自行延续原编排预算。
- Agent/host 收到 `retryable=false`:`max_additional_attempts=0`,立即停止,不得自动重跑。
- Agent/host 收到 `retryable=true`:`max_additional_attempts=2`,初次失败后最多再尝试 2 次;错误若给出 `retry_after_seconds` 或 `next_retry_at`,不得提前重试。
- 未返回 `retryable`(即 `retryable=unknown`):Agent/host 使用 `max_additional_attempts=1`,最多补偿尝试 1 次;仍无法确认时停止并上报错误与 trace。
- `in_flight` 表示同一逻辑订阅已有请求执行中;`cooldown` 或 `terminal_hold` 表示当前被退避或终态保护。遇到这些状态不得递归调用 `event consume`、并行启动相同订阅或通过新 subId / trace 绕过;等待原请求或保护时间结束,同时由 Agent/host 继续维护自己的编排次数。
- 多事件命令必须作为同一次原始操作治理。不得把失败事件拆成新的单事件命令、调整顺序或反复重启来重置预算;启动中任一项失败时,由 CLI 回滚本次已经创建的订阅。
### 本地保护状态运维
- open 版默认状态文件为 `~/.dws/events/open/personal_stream/<identity_hash>/personal_subscription_attempts.json`。设置 `DWS_CONFIG_DIR` 后配置根目录随之变化;其它 edition 使用对应 edition 目录,不固定为 `open`。
- identity 目录权限为 `0700`;`personal_subscription_attempts.json` 与 `personal_subscription_attempts.lock` 权限均为 `0600`。
- 连续 24h 没有失败后失败计数重置;`terminal_hold` 持续 1h。正常处理优先等待错误中的 `next_retry_at`,不要把删状态文件当成常规重试手段。
- 紧急恢复时,先确认该 identity 没有正在创建订阅的进程;只删除 `personal_subscription_attempts.json`,不要删除 lock 文件。删除 JSON 会清空该 identity 的全部保护记录,不只影响一个事件。
## Call flow
1. 从用户意图选择事件码;人名或群名先解析成必填 ID。
@@ -373,14 +373,22 @@ dws event stop --all --yes
以下约束适用于上表全部 16 个事件以及多事件命令中的每一项,只治理 `[event] ready` 之前的订阅创建;ready 之后的 Stream 断线由长连接重连机制处理。
- `0/2/1` 是 **Agent/host 编排约束**,不是 CLI 持久化硬总次数上限。每次 `dws event consume` 调用对每个逻辑订阅最多发送一次订阅创建 HTTP 请求,进程内不会自动重试。CLI 本地状态只持久化 `in_flight`、`cooldown`、`terminal_hold` 三种保护状态,不持久化或计算跨调用的 Agent/host 尝试次数。
- 解析人名或群名、执行 `event consume` 以及后续 `event status/stop` 必须使用同一个 `--profile`。不得把其它 profile 下解析出的 userId、openDingtalkId 或 openConversationId 直接带入当前 profile 的订阅。
- 同一逻辑订阅由当前 profile / 身份、event key、rule type、目标和过滤条件共同确定。`subscribe_id`、`trace_id` 以及重新启动进程都只是诊断或执行信息,不会生成新预算。
- `retryable=false`:`max_additional_attempts=0`,立即停止,不得自动重跑。
- `retryable=true`:`max_additional_attempts=2`,初次失败后最多再尝试 2 次;错误若给出 `retry_after_seconds` 或 `next_retry_at`,不得提前重试。
- 未返回 `retryable`(即 `retryable=unknown`):`max_additional_attempts=1`,最多补偿尝试 1 次;仍无法确认时停止并上报错误与 trace。
- `in_flight` 表示同一逻辑订阅已有请求执行中;`cooldown` 或 `terminal_hold` 表示当前被退避或终态保护。遇到这些状态不得递归调用 `event consume`、并行启动相同订阅或通过新 subId / trace 绕过;等待原请求或保护时间结束,并继续消耗原预算。
- 同一逻辑订阅由当前 profile / 身份、event key、rule type、目标和过滤条件共同确定。`subscribe_id`、`trace_id` 以及重新启动进程都只是诊断或执行信息,不会生成新的逻辑操作;Agent/host 必须自行延续原编排预算。
- Agent/host 收到 `retryable=false`:`max_additional_attempts=0`,立即停止,不得自动重跑。
- Agent/host 收到 `retryable=true`:`max_additional_attempts=2`,初次失败后最多再尝试 2 次;错误若给出 `retry_after_seconds` 或 `next_retry_at`,不得提前重试。
- 未返回 `retryable`(即 `retryable=unknown`):Agent/host 使用 `max_additional_attempts=1`,最多补偿尝试 1 次;仍无法确认时停止并上报错误与 trace。
- `in_flight` 表示同一逻辑订阅已有请求执行中;`cooldown` 或 `terminal_hold` 表示当前被退避或终态保护。遇到这些状态不得递归调用 `event consume`、并行启动相同订阅或通过新 subId / trace 绕过;等待原请求或保护时间结束,同时由 Agent/host 继续维护自己的编排次数。
- 多事件命令必须作为同一次原始操作治理。不得把失败事件拆成新的单事件命令、调整顺序或反复重启来重置预算;启动中任一项失败时,由 CLI 回滚本次已经创建的订阅。
### 本地保护状态运维
- open 版默认状态文件为 `~/.dws/events/open/personal_stream/<identity_hash>/personal_subscription_attempts.json`。设置 `DWS_CONFIG_DIR` 后配置根目录随之变化;其它 edition 使用对应 edition 目录,不固定为 `open`。
- identity 目录权限为 `0700`;`personal_subscription_attempts.json` 与 `personal_subscription_attempts.lock` 权限均为 `0600`。
- 连续 24h 没有失败后失败计数重置;`terminal_hold` 持续 1h。正常处理优先等待错误中的 `next_retry_at`,不要把删状态文件当成常规重试手段。
- 紧急恢复时,先确认该 identity 没有正在创建订阅的进程;只删除 `personal_subscription_attempts.json`,不要删除 lock 文件。删除 JSON 会清空该 identity 的全部保护记录,不只影响一个事件。
## Troubleshooting
- 没有输出:单事件确认 stderr 已出现 `ready event_key=...`;多事件确认已出现 `ready event_count=...`,不要把某条 `subscription` 行误认为整体就绪。
+86 -2
View File
@@ -130,7 +130,7 @@ func TestEventSkillUsesFlatOutputContract(t *testing.T) {
}
}
func TestEventSkillPinsSubscriptionRetryBudget(t *testing.T) {
func TestCrossPlatformCoverageEventSkillPinsSubscriptionRetryOrchestrationContract(t *testing.T) {
_, filename, _, ok := runtime.Caller(0)
if !ok {
t.Fatal("runtime.Caller(0) failed")
@@ -148,9 +148,12 @@ func TestEventSkillPinsSubscriptionRetryBudget(t *testing.T) {
t.Fatalf("read %s: %v", path, err)
}
text := string(content)
normalizedText := strings.Join(strings.Fields(text), " ")
for _, required := range []string{
"16",
"--profile",
"Agent/host",
"0/2/1",
"retryable=false",
"max_additional_attempts=0",
"retryable=true",
@@ -161,11 +164,92 @@ func TestEventSkillPinsSubscriptionRetryBudget(t *testing.T) {
"next_retry_at",
"in_flight",
"cooldown",
"terminal_hold",
"subscribe_id",
"trace_id",
} {
if !strings.Contains(text, required) {
t.Errorf("%s missing subscription retry contract %q", path, required)
t.Errorf("%s missing subscription retry orchestration contract %q", path, required)
}
}
if path == filepath.Join(root, "docs", "event-subprocess-contract.md") {
for _, required := range []string{
"not a CLI-enforced persisted total-attempt cap",
"performs no in-process automatic retry",
"does not persist or enforce the Agent/host attempt count",
} {
if !strings.Contains(normalizedText, required) {
t.Errorf("%s overstates CLI retry enforcement; missing %q", path, required)
}
}
continue
}
for _, required := range []string{
"不是 CLI 持久化硬总次数上限",
"进程内不会自动重试",
"不持久化或计算跨调用的 Agent/host 尝试次数",
} {
if !strings.Contains(normalizedText, required) {
t.Errorf("%s overstates CLI retry enforcement; missing %q", path, required)
}
}
}
}
func TestCrossPlatformCoverageEventSkillDocumentsSubscriptionGuardOperations(t *testing.T) {
_, filename, _, ok := runtime.Caller(0)
if !ok {
t.Fatal("runtime.Caller(0) failed")
}
root := filepath.Clean(filepath.Join(filepath.Dir(filename), "..", ".."))
paths := []string{
filepath.Join(root, "skills", "multi", "dingtalk-event", "SKILL.md"),
filepath.Join(root, "skills", "multi", "dingtalk-event", "references", "event-im.md"),
filepath.Join(root, "skills", "mono", "references", "products", "event.md"),
filepath.Join(root, "docs", "event-subprocess-contract.md"),
}
for _, path := range paths {
content, err := os.ReadFile(path)
if err != nil {
t.Fatalf("read %s: %v", path, err)
}
text := string(content)
normalizedText := strings.Join(strings.Fields(text), " ")
for _, required := range []string{
"~/.dws/events/open/personal_stream/<identity_hash>/personal_subscription_attempts.json",
"DWS_CONFIG_DIR",
"personal_subscription_attempts.json",
"personal_subscription_attempts.lock",
"0700",
"0600",
"24h",
"1h",
"terminal_hold",
"next_retry_at",
} {
if !strings.Contains(text, required) {
t.Errorf("%s missing subscription guard operations %q", path, required)
}
}
if path == filepath.Join(root, "docs", "event-subprocess-contract.md") {
for _, required := range []string{
"Delete only",
"never the lock file",
"every protection record for that identity",
} {
if !strings.Contains(normalizedText, required) {
t.Errorf("%s missing emergency guard-clear warning %q", path, required)
}
}
continue
}
for _, required := range []string{
"只删除 `personal_subscription_attempts.json`",
"不要删除 lock 文件",
"该 identity 的全部保护记录",
} {
if !strings.Contains(normalizedText, required) {
t.Errorf("%s missing emergency guard-clear warning %q", path, required)
}
}
}