diff --git a/AGENTS.md b/AGENTS.md index 46c7870..eeb8aa0 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -95,7 +95,7 @@ - **私有配置位置(本机路径相对本仓库根目录,只读,绝不提交)**:`.local/provider-ai.env` 是 `0600` 的 `KEY=VALUE` 文件;字段名为 `BAILIAN_API_KEY`、`BAILIAN_BASE_URL`、`BAILIAN_WSS_BASE_URL`、`BAILIAN_TTS_VOICE`、`VOLC_ASR_APP_NAME`、`VOLC_ASR_APP_KEY`、`VOLCENGINE_ACCESS_KEY`、`VOLCENGINE_SECRET_KEY`、`VOLCENGINE_REGION`、`VOLCENGINE_DISABLE_SSL`。根目录 `aliyun-oss.env` 也是 `0600`,**不是 shell env 文件**;它以冒号分隔,字段名准确为 `bucket`、`Endpoint`、`Region`,以及 `RAM` 下的 `username`、`accessKeyId`、`accessKeySecret`(大小写须保持原样)。测试机 `rogee` 用户的现行 ARI 文件位于 `~/.config/go-sip-asterisk/{ari.conf,http.conf,ari-secret}`,不是旧 `.local/asterisk-*/ari.conf`;访问测试机前先核对已登记的 SSH 主机指纹,不展示 `ari-secret`。 - **下次安全读取步骤**:先确认工作目录是本仓库,用 `stat` 仅检查本机两份文件是否存在、所有者与权限 `0600`;不满足即停止。按各自格式在受限本机进程中解析所需字段到内存,不执行 `source`、不打印全文/字段值、不写临时明文副本,不把密钥、签名 URL、音频或完整对话带入聊天、日志、提交及长期证据。AI 的历史文件只可作为**获准凭据来源**,模型/voice/速度等仍由当前获批的 task/providers 快照固定,不能用环境变量覆盖。OSS 历史文件也不能直接传给 `DISPATCHER_OSS_CONFIG_FILE`:该运行配置要求私有 JSON、`dispatcher_id` 和 `oss` 字段,并以环境变量**名称引用**密钥;需按现行合同构造并核验授权后才能使用。普通构建和测试不读取这些私有文件;真实服务测试必须显式启用对应 opt-in 并受现行门禁约束。 - **授权边界**:使用者将本轮验收目标改为三条已登记线路、两个既有白名单号码中任意一组真实接通并有 LLM 正常应答;上述 AI 与 OSS 私有配置仍仅用于当前获准的非生产测试,下次任务须重新确认范围和真实服务调用授权,不能沿用本轮或历史一次性授权。不得把历史配置直接当 SaaS 快照、任务授权或真实呼叫准入,不覆盖/清理旧 OSS 对象。 -- Agent 录音经受控双向 TLS 向 D 领取短期 OSS 上传授权,每次尝试只作**一次 HTTPS PUT**;正常上传成功不写录音/结果业务文件。首次明确失败须先完整保存录音与结果两份恢复文件,才从该时刻启动 48 小时重试;按 1、2、4、8、16、32、60 分钟及其后每 60 分钟的固定节奏显式重新申请授权,同一 OSS 目标、同一消息身份。PUT 结果未知不得盲目重传;48 小时届满仍失败时保留文件待人工,**不伪造最终结果或自动清理**。D 不转发文件,已确认结束的通话及时释放执行占用;未知执行仍占用。只有真实终结后才通过唯一 `call.execute.result` 回报录音路径、最终转写和拒联事实;无录音或录音生成失败以空 `recording={}` 和真实结果收口,生成失败须说明原因。不能恢复的录音不声称零丢失,也不伪造 OSS/SaaS 应用回执。凭据/TOKEN/签名 URL 不写入样例、日志、源码或证据。 +- Agent 录音经受控双向 TLS 向 D 领取短期 OSS 上传授权,每次尝试只作**一次 HTTPS PUT**;正常上传不写录音文件,最终结果在 Dispatcher 确认前允许写入 Agent 私有临时结果文件,确认后删除;已确认挂断但结束回报未确认时须保留原结果并重报原结束事实;PUT 前预存的结果在成功未被确认时不得自行报告。首次明确失败须先完整保存录音与结果两份恢复文件,才从该时刻启动 48 小时重试;按 1、2、4、8、16、32、60 分钟及其后每 60 分钟的固定节奏显式重新申请授权,同一 OSS 目标、同一消息身份。PUT 结果未知不得盲目重传;48 小时届满仍失败时保留文件待人工,**不伪造最终结果或自动清理**。D 不转发文件,已确认结束的通话及时释放执行占用;未知执行仍占用。只有真实终结后才通过唯一 `call.execute.result` 回报录音路径、最终转写和拒联事实;无录音或录音生成失败以空 `recording={}` 和真实结果收口,生成失败须说明原因。不能恢复的录音不声称零丢失,也不伪造 OSS/SaaS 应用回执。凭据/TOKEN/签名 URL 不写入样例、日志、源码或证据。 ## SIP 与真实呼叫限制 diff --git a/cmd/sip-go-agent/agent_command.go b/cmd/sip-go-agent/agent_command.go index 0d01c03..5b3d1ab 100644 --- a/cmd/sip-go-agent/agent_command.go +++ b/cmd/sip-go-agent/agent_command.go @@ -46,6 +46,7 @@ func newAgentCommand() *cobra.Command { return errors.New("Agent mTLS listener certificate is invalid") } var handler *rpc.Server + var recoveryClient agentpb.AgentControlServiceClient if mode == "sip-only" { handler, err = newSIPOnlyAgentServer(settings) } else { @@ -76,6 +77,7 @@ func newAgentCommand() *cobra.Command { handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, client) } else { handler, err = newRealAgentServer(cmd.Context(), settings, client) + recoveryClient = client } } if err != nil { @@ -88,6 +90,9 @@ func newAgentCommand() *cobra.Command { defer listener.Close() server := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS))) agentpb.RegisterAgentControlServiceServer(server, handler) + if mode == "nonprod-real" { + startRealAgentRecovery(cmd.Context(), settings, recoveryClient, handler) + } stopped := make(chan struct{}) go func() { select { diff --git a/cmd/sip-go-agent/agent_real.go b/cmd/sip-go-agent/agent_real.go index 3ad85e9..f4fc50f 100644 --- a/cmd/sip-go-agent/agent_real.go +++ b/cmd/sip-go-agent/agent_real.go @@ -82,6 +82,7 @@ func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, d Root: settings.RecoveryRoot, Upload: agent.UploadClient{AllowedHosts: map[string]struct{}{strings.ToLower(settings.OSSAllowedHost): {}}}, }, + Journal: &agent.ResultJournal{Root: settings.RecoveryRoot}, } return (&rpc.ApprovedRecordedRealCall{ Loader: loader, MediaPayloadType: 118, EvidenceRoot: settings.EvidenceRoot, diff --git a/cmd/sip-go-agent/agent_recovery.go b/cmd/sip-go-agent/agent_recovery.go new file mode 100644 index 0000000..8b5449c --- /dev/null +++ b/cmd/sip-go-agent/agent_recovery.go @@ -0,0 +1,163 @@ +package main + +import ( + "context" + "crypto/sha256" + "encoding/json" + "errors" + "fmt" + "log" + "strings" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/agent" + "git.ipao.vip/rogee/go-sip/internal/config" + "git.ipao.vip/rogee/go-sip/internal/rpc" + "google.golang.org/grpc/status" +) + +// Runtime recovery never repeats an OSS PUT with an unknown outcome. The +// recording state machine alone decides whether another attempt is due. +func startRealAgentRecovery(ctx context.Context, settings config.AgentEnvironment, dispatcher agentpb.AgentControlServiceClient, handler *rpc.Server) { + go func() { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + for { + if err := recoverRealAgentOnce(ctx, settings, dispatcher, handler); err != nil && ctx.Err() == nil { + log.Printf("Agent private recovery needs inspection: error_type=%T", err) + } + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + } + }() +} + +func recoverRealAgentOnce(ctx context.Context, settings config.AgentEnvironment, dispatcher agentpb.AgentControlServiceClient, handler *rpc.Server) error { + if handler == nil || dispatcher == nil || strings.TrimSpace(settings.RecoveryRoot) == "" { + return errors.New("real Agent recovery needs a Dispatcher and private directory") + } + journal := agent.ResultJournal{Root: settings.RecoveryRoot} + call := func(id, dispatcherID string, tenantID int64) (agent.RecordingClient, error) { + if dispatcherID != settings.DispatcherID || tenantID <= 0 || id == "" { + return agent.RecordingClient{}, errors.New("original result has no approved Dispatcher or tenant binding") + } + return agent.RecordingClient{Client: dispatcher, DispatcherID: dispatcherID, TenantID: tenantID, SourceEventID: id, + Session: func(context.Context) (*agentpb.RequestMeta, error) { return handler.ActiveSessionMeta() }}, nil + } + resultErr := journal.RecoverWithEnd(ctx, func(ctx context.Context, entry agent.ResultEntry) error { + client, err := call(entry.SourceEventID, entry.DispatcherID, entry.TenantID) + if err == nil { + err = client.ReportEnded(ctx) + } + if err != nil { + logRecoveryError("call_end", entry.SourceEventID, err) + } + return err + }, func(ctx context.Context, entry agent.ResultEntry) error { + client, err := call(entry.SourceEventID, entry.DispatcherID, entry.TenantID) + if err == nil { + err = validateRestoredResult(entry.Payload, entry.Upload) + } + if err == nil { + _, err = client.ReportFinal(ctx, entry.Payload, entry.Upload) + } + if err != nil { + logRecoveryError("result", entry.SourceEventID, err) + } + return err + }) + recovery := agent.RecordingRecovery{Root: settings.RecoveryRoot, Upload: agent.UploadClient{AllowedHosts: map[string]struct{}{strings.ToLower(settings.OSSAllowedHost): {}}}} + _, uploadErr := recovery.RecoverDue(ctx, func(ctx context.Context, entry agent.RecordingRecoveryEntry) (target agent.RecoveryTarget, targetErr error) { + defer func() { + if targetErr != nil { + logRecoveryError("upload_grant", entry.SourceEventID, targetErr) + } + }() + client, err := call(entry.SourceEventID, entry.DispatcherID, entry.TenantID) + if err != nil { + return agent.RecoveryTarget{}, err + } + asset, err := restoredAsset(entry) + if err != nil { + return agent.RecoveryTarget{}, err + } + grant, err := client.RequestUpload(ctx, asset, entry.UploadID) + if err != nil { + return agent.RecoveryTarget{}, err + } + return agent.RecoveryTarget{Bucket: grant.Bucket, Grant: grant}, nil + }, func(ctx context.Context, entry agent.RecordingRecoveryEntry) (reportErr error) { + defer func() { + if reportErr != nil { + logRecoveryError("recording_result", entry.SourceEventID, reportErr) + } + }() + client, err := call(entry.SourceEventID, entry.DispatcherID, entry.TenantID) + if err != nil { + return err + } + if entry.Uploaded == nil { + return errors.New("original recording has no observed PUT success") + } + observation := &agentpb.UploadObservation{UploadId: entry.UploadID, RecordingId: entry.RecordingID, + PutStatusCode: int32(entry.Uploaded.StatusCode), SizeBytes: entry.Uploaded.SizeBytes, ChecksumSha256: entry.Uploaded.SHA256} + if err := validateRestoredResult(entry.ResultPayload, observation); err != nil { + return err + } + _, err = client.ReportFinal(ctx, entry.ResultPayload, observation) + if err != nil { + return err + } + return journal.DiscardPrepared(entry.SourceEventID) + }) + return errors.Join(resultErr, uploadErr) +} + +func logRecoveryError(stage, id string, err error) { + digest := sha256.Sum256([]byte(id)) + log.Printf("Agent recovery blocked stage=%s event_digest=%x grpc_code=%s error_class=%T", stage, digest[:6], status.Code(err), err) +} + +func restoredAsset(entry agent.RecordingRecoveryEntry) (*agentpb.AssetDescriptor, error) { + var payload struct { + Recording struct { + Status string `json:"status"` + Bucket string `json:"bucket"` + ObjectKey string `json:"object_key"` + DurationMS int64 `json:"duration_ms"` + SizeBytes int64 `json:"size_bytes"` + ChecksumSHA256 string `json:"checksum_sha256"` + } `json:"recording"` + } + if err := json.Unmarshal(entry.ResultPayload, &payload); err != nil || payload.Recording.Status != "uploaded" || payload.Recording.Bucket != entry.Bucket || payload.Recording.ObjectKey != entry.ObjectKey || payload.Recording.SizeBytes != entry.SizeBytes || payload.Recording.ChecksumSHA256 != entry.SHA256 || payload.Recording.DurationMS <= 0 { + return nil, errors.New("original recording asset does not match its persisted result") + } + return &agentpb.AssetDescriptor{Kind: agentpb.AssetKind_ASSET_KIND_RECORDING, AssetId: entry.RecordingID, + CallId: entry.SourceEventID, ExecutionId: entry.SourceEventID, Format: "wav", Channels: 1, SampleRateHz: 16000, + DurationMs: payload.Recording.DurationMS, SizeBytes: entry.SizeBytes, ChecksumSha256: entry.SHA256}, nil +} + +func validateRestoredResult(payload []byte, upload *agentpb.UploadObservation) error { + var result struct { + Recording struct { + Status string `json:"status"` + SizeBytes int64 `json:"size_bytes"` + ChecksumSHA256 string `json:"checksum_sha256"` + } `json:"recording"` + } + if err := json.Unmarshal(payload, &result); err != nil { + return errors.New("original result is invalid") + } + if result.Recording.Status == "uploaded" { + if upload == nil || upload.GetPutStatusCode() < 200 || upload.GetPutStatusCode() >= 300 || result.Recording.SizeBytes != upload.GetSizeBytes() || result.Recording.ChecksumSHA256 != upload.GetChecksumSha256() || upload.GetUploadId() == "" || upload.GetRecordingId() == "" { + return errors.New("original recording has no matching confirmed upload") + } + } else if result.Recording.Status != "" || upload != nil { + return fmt.Errorf("original result has inconsistent recording status") + } + return nil +} diff --git a/cmd/sip-go-agent/agent_recovery_test.go b/cmd/sip-go-agent/agent_recovery_test.go new file mode 100644 index 0000000..f9ad97f --- /dev/null +++ b/cmd/sip-go-agent/agent_recovery_test.go @@ -0,0 +1,54 @@ +package main + +import ( + "context" + "os" + "path/filepath" + "testing" + + "git.ipao.vip/rogee/go-sip/internal/agent" + "git.ipao.vip/rogee/go-sip/internal/config" + "git.ipao.vip/rogee/go-sip/internal/rpc" +) + +func TestRealAgentRecoveryKeepsResultUntilSessionAndDispatcherConfirm(t *testing.T) { + root := t.TempDir() + if err := os.Chmod(root, 0700); err != nil { + t.Fatal(err) + } + const dispatcherID = "11111111-1111-4111-8111-111111111111" + settings := config.AgentEnvironment{DispatcherID: dispatcherID, RecoveryRoot: root, OSSAllowedHost: "oss.example.invalid"} + journal := agent.ResultJournal{Root: root} + if err := journal.SaveReady(agent.ResultEntry{SourceEventID: "call-original", DispatcherID: dispatcherID, TenantID: 1, Payload: []byte(`{"recording":{}}`)}); err != nil { + t.Fatal(err) + } + if err := recoverRealAgentOnce(context.Background(), settings, &isolatedAgentRecordingClient{}, rpc.NewServer(rpc.ServerOptions{})); err == nil { + t.Fatal("result without an active approved session was incorrectly confirmed") + } + files, err := os.ReadDir(filepath.Join(root, ".results")) + if err != nil || len(files) != 1 { + t.Fatalf("unconfirmed original result was lost: files=%v err=%v", files, err) + } +} + +func TestRestoredAssetUsesOnlyOriginalRecordingMetadata(t *testing.T) { + entry := agent.RecordingRecoveryEntry{SourceEventID: "call-original", RecordingID: "recording-original", Bucket: "bucket", ObjectKey: "tenant/original.wav", SizeBytes: 123, SHA256: "digest", ResultPayload: []byte(`{"recording":{"status":"uploaded","bucket":"bucket","object_key":"tenant/original.wav","duration_ms":500,"size_bytes":123,"checksum_sha256":"digest"}}`)} + asset, err := restoredAsset(entry) + if err != nil || asset.GetAssetId() != entry.RecordingID || asset.GetDurationMs() != 500 || asset.GetChecksumSha256() != entry.SHA256 { + t.Fatalf("original recording metadata was changed: asset=%+v err=%v", asset, err) + } + entry.ObjectKey = "tenant/different.wav" + if _, err := restoredAsset(entry); err == nil { + t.Fatal("a regrant for a different OSS object was accepted") + } +} + +func TestRestoredResultRejectsUnconfirmedOrMismatchedUploads(t *testing.T) { + payload := []byte(`{"recording":{"status":"uploaded","size_bytes":123,"checksum_sha256":"digest"}}`) + if validateRestoredResult(payload, nil) == nil { + t.Fatal("upload without an observed success was accepted") + } + if validateRestoredResult([]byte(`{"recording":{}}`), nil) != nil { + t.Fatal("a genuine no-recording result was rejected") + } +}