resume Agent recording and result recovery at startup
This commit is contained in:
@@ -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 {
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user