Retire obsolete Agent failure-event reporting
This commit is contained in:
@@ -115,6 +115,7 @@
|
||||
- Agent 会话持久边界:将会话代际使用的私有原子落盘函数从旧执行日志模块移入独立会话日志文件;原有目录同步、写失败报错和重启代际栅栏测试继续通过。旧执行日志及业务 RPC 尚未删除;此次只移动源码,不修改现存 Agent 恢复文件。
|
||||
- Agent 旧执行状态屏障:批准的 Mock Agent 启动时若发现旧执行日志,立即拒绝启动并明确要求人工处置;隔离测试先复现了原先允许启动的问题,再证明旧日志字节不变、新会话状态未写入。此屏障不迁移、不清理存量文件;没有检查或处置真实环境中的旧记录。
|
||||
- Agent Proto 唯一服务面:新增服务方法集合测试先确认旧 RPC 仍可被服务描述符发现,随后将 `AgentControlService` 收敛为状态/激活、获批执行与控制、加载版本及录音/结果八个方法;重新生成 Go 类型、核验七项来源哈希,并同步改写当前 Proto 错误与恢复说明。旧消息定义、旧服务端实现及专属测试尚待清理,不能把本批当作 P07 完成或真实 Agent 联调。
|
||||
- Agent 旧失败事实分支:旧 Mock 上传失败事实通过已退役的 `ReportExecutionEvent` RPC 回报,现删除该代码及专属测试。现行录音失败、未知 PUT、重启恢复和最终结果仍由 `recording_delivery*` 隔离测试覆盖;不复活额外通话事件或把未知上传当作成功。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -1,141 +0,0 @@
|
||||
package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"os"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/contract"
|
||||
"github.com/google/uuid"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
type MockFailureFactClient interface {
|
||||
ReportExecutionEvent(context.Context, *agentpb.ReportExecutionEventRequest) (*agentpb.ReportExecutionEventResponse, error)
|
||||
}
|
||||
|
||||
// MockUploadFailureCode returns empty for uncertain PUT outcomes. In particular
|
||||
// a transport error, HTTP 408/429 or 5xx never claims that OSS rejected the
|
||||
// object; the Dispatcher retains the unknown outcome until its deadline.
|
||||
func MockUploadFailureCode(cause error) string {
|
||||
switch {
|
||||
case errors.Is(cause, ErrUploadGrantExpired):
|
||||
return "upload_authorization_expired"
|
||||
case errors.Is(cause, ErrUploadGrantInvalid):
|
||||
return "upload_authorization_failed"
|
||||
case errors.Is(cause, ErrUploadChecksumMismatch):
|
||||
return "checksum_mismatch"
|
||||
case errors.Is(cause, os.ErrNotExist):
|
||||
return "upload_failed"
|
||||
}
|
||||
var response *UploadHTTPError
|
||||
if errors.As(cause, &response) && response.StatusCode >= 400 && response.StatusCode < 500 &&
|
||||
response.StatusCode != http.StatusRequestTimeout && response.StatusCode != http.StatusTooManyRequests {
|
||||
if response.StatusCode == http.StatusUnauthorized || response.StatusCode == http.StatusForbidden {
|
||||
return "upload_authorization_failed"
|
||||
}
|
||||
return "upload_failed"
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func validateMockFailureActiveMeta(meta *agentpb.RequestMeta) error {
|
||||
if meta == nil || meta.ProtocolVersion != "agent.v1" || meta.AgentId == "" || meta.CellId == "" ||
|
||||
meta.BootId == "" || meta.DispatcherEpoch == "" || meta.SessionGeneration == 0 {
|
||||
return errors.New("Mock failure requires the current activated Agent session")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ReportMockUploadFailure persists a known, explicit failure before trying to
|
||||
// notify the Dispatcher. It never retries a PUT or claims an uncertain result.
|
||||
// The caller supplies metadata from its activated Agent session, not a newly
|
||||
// invented boot identity. A pending fact is retried with the same ID after
|
||||
// restart, while the request metadata uses the new active session.
|
||||
func (s *Spool) ReportMockUploadFailure(ctx context.Context, client MockFailureFactClient, active *agentpb.RequestMeta, uploadID string, cause error) (bool, error) {
|
||||
code := MockUploadFailureCode(cause)
|
||||
if code == "" {
|
||||
return false, nil
|
||||
}
|
||||
if err := validateMockFailureActiveMeta(active); err != nil {
|
||||
return true, err
|
||||
}
|
||||
record, err := s.LoadUploadAttempt(uploadID)
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
if record.FailureFact == nil {
|
||||
if record.Binding == nil || record.Asset == nil {
|
||||
return true, errors.New("Mock failure has no bound recording attempt")
|
||||
}
|
||||
payload, err := json.Marshal(contract.LocalMockRecordingFailure{
|
||||
SchemaVersion: contract.LocalMockRecordingFailureVersion,
|
||||
UploadID: record.UploadID, RecordingID: record.Asset.AssetId, ErrorCode: code,
|
||||
})
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
sum := sha256.Sum256(payload)
|
||||
fact := &agentpb.ExecutionFact{
|
||||
FactId: uuid.NewString(), Binding: proto.Clone(record.Binding).(*agentpb.ExecutionBinding),
|
||||
Kind: agentpb.FactKind_FACT_KIND_RECORDING_PROGRESS, PayloadJson: payload,
|
||||
ContentSha256: hex.EncodeToString(sum[:]), ObservedAtUnixMs: s.now().UTC().UnixMilli(),
|
||||
SourceBootId: active.BootId, SourceSequence: 1,
|
||||
}
|
||||
if err := s.RecordUploadFailureFact(uploadID, fact); err != nil {
|
||||
return true, fmt.Errorf("persist Mock failure before RPC: %w", err)
|
||||
}
|
||||
record, err = s.LoadUploadAttempt(uploadID)
|
||||
if err != nil {
|
||||
return true, err
|
||||
}
|
||||
}
|
||||
return true, s.deliverMockFailure(ctx, client, active, record)
|
||||
}
|
||||
|
||||
func (s *Spool) RecoverMockUploadFailures(ctx context.Context, client MockFailureFactClient, active *agentpb.RequestMeta) error {
|
||||
pending, err := s.PendingUploadFailures()
|
||||
if err != nil || len(pending) == 0 {
|
||||
return err
|
||||
}
|
||||
if err := validateMockFailureActiveMeta(active); err != nil {
|
||||
return err
|
||||
}
|
||||
for _, record := range pending {
|
||||
if err := s.deliverMockFailure(ctx, client, active, record); err != nil {
|
||||
return fmt.Errorf("recover Mock recording failure %s: %w", record.UploadID, err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Spool) deliverMockFailure(ctx context.Context, client MockFailureFactClient, active *agentpb.RequestMeta, record UploadAttempt) error {
|
||||
if client == nil || record.FailureFact == nil {
|
||||
return errors.New("Mock failure reporter or durable fact is missing")
|
||||
}
|
||||
if err := validateMockFailureActiveMeta(active); err != nil {
|
||||
return err
|
||||
}
|
||||
meta := proto.Clone(active).(*agentpb.RequestMeta)
|
||||
meta.RequestId = record.FailureFact.FactId
|
||||
meta.OperationId = record.FailureFact.FactId
|
||||
meta.IdempotencyKey = "mock-upload-failure:" + record.FailureFact.FactId
|
||||
response, err := client.ReportExecutionEvent(ctx, &agentpb.ReportExecutionEventRequest{
|
||||
Meta: meta, Fact: proto.Clone(record.FailureFact).(*agentpb.ExecutionFact),
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
receipt := response.GetReceipt()
|
||||
if receipt.GetResult() != agentpb.ResultCode_RESULT_CODE_ACCEPTED ||
|
||||
receipt.GetFactId() != record.FailureFact.FactId || receipt.GetContentSha256() != record.FailureFact.ContentSha256 {
|
||||
return errors.New("Dispatcher did not acknowledge the original Mock failure fact")
|
||||
}
|
||||
return s.CompleteUploadFailureReport(record.UploadID, record.FailureFact.FactId)
|
||||
}
|
||||
@@ -1,154 +0,0 @@
|
||||
package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net/http"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"github.com/google/uuid"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
type mockFailureRPC struct {
|
||||
requests []*agentpb.ReportExecutionEventRequest
|
||||
unavailable bool
|
||||
wrongReceipt bool
|
||||
}
|
||||
|
||||
func (m *mockFailureRPC) ReportExecutionEvent(_ context.Context, request *agentpb.ReportExecutionEventRequest) (*agentpb.ReportExecutionEventResponse, error) {
|
||||
m.requests = append(m.requests, proto.Clone(request).(*agentpb.ReportExecutionEventRequest))
|
||||
if m.unavailable {
|
||||
return nil, errors.New("Dispatcher unavailable before receipt")
|
||||
}
|
||||
factID := request.Fact.FactId
|
||||
if m.wrongReceipt {
|
||||
factID = uuid.NewString()
|
||||
}
|
||||
return &agentpb.ReportExecutionEventResponse{Receipt: &agentpb.OperationReceipt{
|
||||
Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED, FactId: factID, ContentSha256: request.Fact.ContentSha256,
|
||||
}}, nil
|
||||
}
|
||||
|
||||
func TestMockUploadFailureReportReusesDurableFactAcrossActiveBoots(t *testing.T) {
|
||||
root := t.TempDir()
|
||||
observed := time.Date(2026, 9, 21, 2, 0, 0, 0, time.UTC)
|
||||
spool, err := NewSpool(root, func() time.Time { return observed })
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
binding := &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a"}
|
||||
if err := spool.ClaimUpload(UploadAttempt{
|
||||
UploadID: "upload-a", Identity: "identity-a", State: "attempted", RequestID: uuid.NewString(),
|
||||
Binding: binding, Asset: &agentpb.AssetDescriptor{AssetId: "recording-a"},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
active := &agentpb.RequestMeta{
|
||||
ProtocolVersion: "agent.v1", AgentId: "agent-a", CellId: "cell-a", BootId: "boot-old",
|
||||
DispatcherEpoch: "epoch-a", SessionGeneration: 1,
|
||||
}
|
||||
remote := &mockFailureRPC{unavailable: true}
|
||||
known, err := spool.ReportMockUploadFailure(context.Background(), remote, active, "upload-a", &UploadHTTPError{StatusCode: http.StatusForbidden})
|
||||
if !known || err == nil || len(remote.requests) != 1 {
|
||||
t.Fatalf("explicit 403 failure was not durably sent: known=%t requests=%d err=%v", known, len(remote.requests), err)
|
||||
}
|
||||
record, err := spool.LoadUploadAttempt("upload-a")
|
||||
if err != nil || record.FailureFact == nil || record.FailureDelivered {
|
||||
t.Fatalf("fact lost before receipt: record=%+v err=%v", record, err)
|
||||
}
|
||||
factID, digest, sourceBoot := record.FailureFact.FactId, record.FailureFact.ContentSha256, record.FailureFact.SourceBootId
|
||||
if sourceBoot != active.BootId || record.FailureFact.ObservedAtUnixMs != observed.UnixMilli() {
|
||||
t.Fatalf("fact lost source boot or observation clock: boot=%q observed=%d", sourceBoot, record.FailureFact.ObservedAtUnixMs)
|
||||
}
|
||||
restarted, err := NewSpool(root, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
activeNew := proto.Clone(active).(*agentpb.RequestMeta)
|
||||
activeNew.BootId = "boot-new"
|
||||
activeNew.SessionGeneration = 2
|
||||
remote.unavailable, remote.wrongReceipt = false, true
|
||||
if err := restarted.RecoverMockUploadFailures(context.Background(), remote, activeNew); err == nil {
|
||||
t.Fatal("unrelated accepted receipt closed a durable failure fact")
|
||||
}
|
||||
record, err = restarted.LoadUploadAttempt("upload-a")
|
||||
if err != nil || record.FailureDelivered {
|
||||
t.Fatalf("false receipt advanced failure state: record=%+v err=%v", record, err)
|
||||
}
|
||||
remote.wrongReceipt = false
|
||||
if err := restarted.RecoverMockUploadFailures(context.Background(), remote, activeNew); err != nil {
|
||||
t.Fatalf("failure fact not recovered after boot change: %v", err)
|
||||
}
|
||||
record, err = restarted.LoadUploadAttempt("upload-a")
|
||||
if err != nil || !record.FailureDelivered || record.State != "attempted" {
|
||||
t.Fatalf("accepted receipt did not close only the fact: record=%+v err=%v", record, err)
|
||||
}
|
||||
if len(remote.requests) != 3 {
|
||||
t.Fatalf("expected original, false receipt and recovery; got %d reports", len(remote.requests))
|
||||
}
|
||||
for _, request := range remote.requests {
|
||||
if request.Fact.FactId != factID || request.Fact.ContentSha256 != digest || request.Fact.SourceBootId != sourceBoot {
|
||||
t.Fatalf("recovery generated a new failure identity: %+v", request.Fact)
|
||||
}
|
||||
}
|
||||
if remote.requests[2].Meta.BootId != activeNew.BootId || remote.requests[2].Meta.SessionGeneration != activeNew.SessionGeneration {
|
||||
t.Fatalf("recovery did not use the current active session: %+v", remote.requests[2].Meta)
|
||||
}
|
||||
if err := restarted.RecoverMockUploadFailures(context.Background(), remote, activeNew); err != nil || len(remote.requests) != 3 {
|
||||
t.Fatalf("acknowledged failure re-reported: requests=%d err=%v", len(remote.requests), err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMockUploadFailureClassificationKeepsUncertainOutcomesUnknown(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
err error
|
||||
code string
|
||||
}{
|
||||
{"expired grant", ErrUploadGrantExpired, "upload_authorization_expired"},
|
||||
{"invalid grant", ErrUploadGrantInvalid, "upload_authorization_failed"},
|
||||
{"checksum mismatch", ErrUploadChecksumMismatch, "checksum_mismatch"},
|
||||
{"explicit 400", &UploadHTTPError{StatusCode: http.StatusBadRequest}, "upload_failed"},
|
||||
{"explicit 403", &UploadHTTPError{StatusCode: http.StatusForbidden}, "upload_authorization_failed"},
|
||||
{"missing recording", os.ErrNotExist, "upload_failed"},
|
||||
{"request timeout", &UploadHTTPError{StatusCode: http.StatusRequestTimeout}, ""},
|
||||
{"throttled", &UploadHTTPError{StatusCode: http.StatusTooManyRequests}, ""},
|
||||
{"server uncertain", &UploadHTTPError{StatusCode: http.StatusInternalServerError}, ""},
|
||||
{"transport uncertain", errors.New("connection reset"), ""},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
if got := MockUploadFailureCode(tc.err); got != tc.code {
|
||||
t.Fatalf("classification %q, want %q", got, tc.code)
|
||||
}
|
||||
})
|
||||
}
|
||||
spool, err := NewSpool(t.TempDir(), nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := spool.ClaimUpload(UploadAttempt{
|
||||
UploadID: "unknown-a", Identity: "identity-a", State: "attempted", RequestID: uuid.NewString(),
|
||||
Binding: &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a"},
|
||||
Asset: &agentpb.AssetDescriptor{AssetId: "recording-a"},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
active := &agentpb.RequestMeta{ProtocolVersion: "agent.v1", AgentId: "agent-a", CellId: "cell-a", BootId: "boot-a", DispatcherEpoch: "epoch-a", SessionGeneration: 1}
|
||||
remote := &mockFailureRPC{}
|
||||
known, err := spool.ReportMockUploadFailure(context.Background(), remote, active, "unknown-a", &UploadHTTPError{StatusCode: 500})
|
||||
if err != nil || known || len(remote.requests) != 0 {
|
||||
t.Fatalf("unknown PUT was reported as failed: known=%t reports=%d err=%v", known, len(remote.requests), err)
|
||||
}
|
||||
pending, err := spool.PendingUploadFailures()
|
||||
if err != nil || len(pending) != 0 {
|
||||
t.Fatalf("unknown PUT produced a terminal fact: pending=%d err=%v", len(pending), err)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(spool.root, ".uploads", "unknown-a", "state.json")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user