Move Agent RPC schema to unversioned internal package
This commit is contained in:
@@ -5,16 +5,16 @@ import (
|
||||
"errors"
|
||||
"log/slog"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
)
|
||||
|
||||
// mockAuthorizedOriginator is deliberately isolated from SIP and Asterisk.
|
||||
// Mixed/real modes have no adapter until separately authorized and verified.
|
||||
func mockAuthorizedOriginator(mode string) func(context.Context, *agentv1.ExecuteAuthorizedRequest) error {
|
||||
func mockAuthorizedOriginator(mode string) func(context.Context, *agentpb.ExecuteAuthorizedRequest) error {
|
||||
if mode != "mock" {
|
||||
return nil
|
||||
}
|
||||
return func(ctx context.Context, request *agentv1.ExecuteAuthorizedRequest) error {
|
||||
return func(ctx context.Context, request *agentpb.ExecuteAuthorizedRequest) error {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -4,7 +4,7 @@ import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
)
|
||||
|
||||
func TestAuthorizedOriginatorIsExplicitlyMockOnly(t *testing.T) {
|
||||
@@ -15,7 +15,7 @@ func TestAuthorizedOriginatorIsExplicitlyMockOnly(t *testing.T) {
|
||||
if originator == nil {
|
||||
t.Fatal("mock mode has no isolated originator")
|
||||
}
|
||||
request := &agentv1.ExecuteAuthorizedRequest{Binding: &agentv1.ExecutionBinding{ExecutionId: "execution-1"}, SelectedTrunkId: "trunk-1"}
|
||||
request := &agentpb.ExecuteAuthorizedRequest{Binding: &agentpb.ExecutionBinding{ExecutionId: "execution-1"}, SelectedTrunkId: "trunk-1"}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
if err := originator(ctx, request); err == nil {
|
||||
|
||||
+10
-10
@@ -15,7 +15,7 @@ import (
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/ai"
|
||||
"git.ipao.vip/rogee/go-sip/internal/callflow"
|
||||
@@ -363,18 +363,18 @@ func serveAgentRPC(cfg config.Config, spool *agent.Spool, report agent.RecoveryR
|
||||
grpcServer := grpc.NewServer(grpc.Creds(credentials.NewTLS(tlsConfig)))
|
||||
handler := rpc.NewServer(rpc.ServerOptions{
|
||||
Mode: cfg.Mode,
|
||||
Status: &agentv1.AgentStatus{
|
||||
Status: &agentpb.AgentStatus{
|
||||
AgentId: cfg.AgentID,
|
||||
CellId: cfg.CellID,
|
||||
BootId: bootID,
|
||||
SoftwareVersion: cfg.Version,
|
||||
ProtocolVersion: "agent.v1",
|
||||
AdmissionState: agentv1.AdmissionState_ADMISSION_STATE_CLOSED,
|
||||
AdmissionState: agentpb.AdmissionState_ADMISSION_STATE_CLOSED,
|
||||
StatusReason: fmt.Sprintf("recovered_unknown=%d quarantined=%d spool=%s", len(report.Unknown), len(report.Quarantined), spool.Root()),
|
||||
Capabilities: []*agentv1.Capability{{Name: "mode", Value: cfg.Mode}, {Name: "grpc_transport", Value: "unary-mtls"}, {Name: "resource_sample", Value: "partial-unknown"}},
|
||||
Capabilities: []*agentpb.Capability{{Name: "mode", Value: cfg.Mode}, {Name: "grpc_transport", Value: "unary-mtls"}, {Name: "resource_sample", Value: "partial-unknown"}},
|
||||
Resources: health.Sampler{}.Sample(context.Background(), spool.Root()),
|
||||
},
|
||||
UploadPolicy: &agentv1.UploadPolicy{Enabled: cfg.Mode == "mock", MaxAssetBytes: 16 << 20},
|
||||
UploadPolicy: &agentpb.UploadPolicy{Enabled: cfg.Mode == "mock", MaxAssetBytes: 16 << 20},
|
||||
StaticArtifactRaw: staticArtifactRaw,
|
||||
StaticArtifactExpected: contract.StaticArtifactExpectation{CellID: cfg.CellID, Mode: cfg.Mode},
|
||||
RequirePeerCertificate: true,
|
||||
@@ -383,7 +383,7 @@ func serveAgentRPC(cfg config.Config, spool *agent.Spool, report agent.RecoveryR
|
||||
CallLogger: callLogger,
|
||||
MockAuthorizedOriginate: mockAuthorizedOriginator(cfg.Mode),
|
||||
})
|
||||
agentv1.RegisterAgentControlServiceServer(grpcServer, handler)
|
||||
agentpb.RegisterAgentControlServiceServer(grpcServer, handler)
|
||||
serveCtx, cancel := signalContext()
|
||||
defer cancel()
|
||||
stopUploadRecovery, err := startUploadNotificationRecovery(serveCtx, cfg, spool, handler.ActiveSessionMeta)
|
||||
@@ -557,7 +557,7 @@ func newDispatcherCommand() *cobra.Command {
|
||||
return errors.New("MTLS_PEER_CERT_FINGERPRINTS is required when Dispatcher gRPC is enabled")
|
||||
}
|
||||
allowedAgentIDs := parseCSVSet(cfg.DispatcherGRPCAllowedAgentIDs)
|
||||
authorizeAgentSession := func(meta *agentv1.RequestMeta) error {
|
||||
authorizeAgentSession := func(meta *agentpb.RequestMeta) error {
|
||||
if connectedAgents == nil || connectedAgents.coordinator == nil {
|
||||
return fmt.Errorf("no activated Agent coordinator: %w", store.ErrCommandConflict)
|
||||
}
|
||||
@@ -565,13 +565,13 @@ func newDispatcherCommand() *cobra.Command {
|
||||
}
|
||||
uploadHandler, handlerErr := rpc.NewDispatcherUploadServerWithOptions(st, uploadClient, time.Now, rpc.DispatcherUploadOptions{
|
||||
RequirePeer: true, PeerCertificateFingerprints: peerFingerprints, AllowedAgentIDs: allowedAgentIDs,
|
||||
LocalV3Authorize: func(ctx context.Context, req *agentv1.RequestUploadRequest) error {
|
||||
LocalV3Authorize: func(ctx context.Context, req *agentpb.RequestUploadRequest) error {
|
||||
if err := authorizeAgentSession(req.Meta); err != nil {
|
||||
return err
|
||||
}
|
||||
return d.AuthorizeLocalMockUpload(ctx, req)
|
||||
},
|
||||
LocalV3Complete: func(ctx context.Context, req *agentv1.CompleteUploadRequest, record store.UploadRecord) (bool, error) {
|
||||
LocalV3Complete: func(ctx context.Context, req *agentpb.CompleteUploadRequest, record store.UploadRecord) (bool, error) {
|
||||
if err := authorizeAgentSession(req.Meta); err != nil {
|
||||
return false, err
|
||||
}
|
||||
@@ -605,7 +605,7 @@ func newDispatcherCommand() *cobra.Command {
|
||||
}
|
||||
return handler(ctx, req)
|
||||
}))
|
||||
agentv1.RegisterAgentControlServiceServer(dispatcherGRPC, dispatcherHandler)
|
||||
agentpb.RegisterAgentControlServiceServer(dispatcherGRPC, dispatcherHandler)
|
||||
go func() {
|
||||
if serveErr := dispatcherGRPC.Serve(dispatcherListener); serveErr != nil && !errors.Is(serveErr, grpc.ErrServerStopped) {
|
||||
slog.Error("Dispatcher gRPC stopped", "error", serveErr)
|
||||
|
||||
@@ -10,7 +10,7 @@ import (
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/callruntime"
|
||||
"git.ipao.vip/rogee/go-sip/internal/config"
|
||||
@@ -38,7 +38,7 @@ func uploadCallRecordings(ctx context.Context, cfg config.Config, result callrun
|
||||
if executionID == "" {
|
||||
executionID = result.ChannelID
|
||||
}
|
||||
binding := &agentv1.ExecutionBinding{
|
||||
binding := &agentpb.ExecutionBinding{
|
||||
TenantId: cfg.CallTenantID,
|
||||
TenantKey: cfg.CallTenantKey,
|
||||
ExecutionId: executionID,
|
||||
@@ -69,8 +69,8 @@ func uploadCallRecordings(ctx context.Context, cfg config.Config, result callrun
|
||||
return nil, fmt.Errorf("recording %d is missing path, size or checksum", index+1)
|
||||
}
|
||||
assetID := recordingAssetID(recording.Segment, recording.Path, index+1)
|
||||
asset := &agentv1.AssetDescriptor{
|
||||
Kind: agentv1.AssetKind_ASSET_KIND_RECORDING,
|
||||
asset := &agentpb.AssetDescriptor{
|
||||
Kind: agentpb.AssetKind_ASSET_KIND_RECORDING,
|
||||
AssetId: assetID,
|
||||
CallId: result.ChannelID,
|
||||
ExecutionId: executionID,
|
||||
@@ -119,15 +119,15 @@ func recordingAssetID(segment, path string, index int) string {
|
||||
return assetID
|
||||
}
|
||||
|
||||
func stableUploadID(binding *agentv1.ExecutionBinding, asset *agentv1.AssetDescriptor) string {
|
||||
func stableUploadID(binding *agentpb.ExecutionBinding, asset *agentpb.AssetDescriptor) string {
|
||||
value := binding.ExecutionId + "\x00" + asset.AssetId + "\x00" + asset.ChecksumSha256
|
||||
digest := sha256.Sum256([]byte(value))
|
||||
return "upload-" + hex.EncodeToString(digest[:16])
|
||||
}
|
||||
|
||||
func uploadMeta(cfg config.Config, phase, uploadID string) *agentv1.RequestMeta {
|
||||
func uploadMeta(cfg config.Config, phase, uploadID string) *agentpb.RequestMeta {
|
||||
operationID := "recording-" + phase + "-" + uploadID
|
||||
return &agentv1.RequestMeta{
|
||||
return &agentpb.RequestMeta{
|
||||
ProtocolVersion: "agent.v1",
|
||||
RequestId: operationID,
|
||||
TraceId: operationID,
|
||||
|
||||
@@ -10,18 +10,18 @@ import (
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
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"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
type recordingUploadRPC interface {
|
||||
RequestUpload(context.Context, *agentv1.RequestUploadRequest) (*agentv1.RequestUploadResponse, error)
|
||||
CompleteUpload(context.Context, *agentv1.CompleteUploadRequest) (*agentv1.CompleteUploadResponse, error)
|
||||
RequestUpload(context.Context, *agentpb.RequestUploadRequest) (*agentpb.RequestUploadResponse, error)
|
||||
CompleteUpload(context.Context, *agentpb.CompleteUploadRequest) (*agentpb.CompleteUploadResponse, error)
|
||||
}
|
||||
|
||||
func uploadRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, uploader agent.UploadClient, spool *agent.Spool, binding *agentv1.ExecutionBinding, asset *agentv1.AssetDescriptor, path string) (returned agent.UploadAttempt, returnErr error) {
|
||||
func uploadRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, uploader agent.UploadClient, spool *agent.Spool, binding *agentpb.ExecutionBinding, asset *agentpb.AssetDescriptor, path string) (returned agent.UploadAttempt, returnErr error) {
|
||||
if binding == nil || asset == nil {
|
||||
return returned, errors.New("upload binding and asset are required")
|
||||
}
|
||||
@@ -31,7 +31,7 @@ func uploadRecording(ctx context.Context, cfg config.Config, client recordingUpl
|
||||
return returned, err
|
||||
}
|
||||
defer func() { returnErr = errors.Join(returnErr, lock.Close()) }()
|
||||
identityBytes, err := (proto.MarshalOptions{Deterministic: true}).Marshal(&agentv1.RequestUploadRequest{Binding: binding, Asset: asset, UploadId: id})
|
||||
identityBytes, err := (proto.MarshalOptions{Deterministic: true}).Marshal(&agentpb.RequestUploadRequest{Binding: binding, Asset: asset, UploadId: id})
|
||||
if err != nil {
|
||||
return agent.UploadAttempt{}, err
|
||||
}
|
||||
@@ -39,11 +39,11 @@ func uploadRecording(ctx context.Context, cfg config.Config, client recordingUpl
|
||||
identity := hex.EncodeToString(digest[:])
|
||||
record, err := spool.LoadUploadAttempt(id)
|
||||
if errors.Is(err, os.ErrNotExist) {
|
||||
response, err := client.RequestUpload(ctx, &agentv1.RequestUploadRequest{Meta: uploadMeta(cfg, "request", id), Binding: binding, Asset: asset, UploadId: id})
|
||||
response, err := client.RequestUpload(ctx, &agentpb.RequestUploadRequest{Meta: uploadMeta(cfg, "request", id), Binding: binding, Asset: asset, UploadId: id})
|
||||
if err != nil {
|
||||
return record, fmt.Errorf("request upload %s: %w", id, err)
|
||||
}
|
||||
if response.GetReceipt().GetResult() != agentv1.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil {
|
||||
if response.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil {
|
||||
return record, fmt.Errorf("request upload %s has no accepted grant", id)
|
||||
}
|
||||
grant := response.Grant
|
||||
@@ -106,11 +106,11 @@ func notifyUploadedRecordingLocked(ctx context.Context, cfg config.Config, clien
|
||||
return record, errors.New("successful upload and original notification metadata are required")
|
||||
}
|
||||
id := record.UploadID
|
||||
response, err := client.CompleteUpload(ctx, &agentv1.CompleteUploadRequest{Meta: uploadMeta(cfg, "complete", id), Binding: record.Binding, Asset: record.Asset, UploadId: id, UploadedSizeBytes: record.Result.SizeBytes, UploadedChecksumSha256: record.Result.SHA256})
|
||||
response, err := client.CompleteUpload(ctx, &agentpb.CompleteUploadRequest{Meta: uploadMeta(cfg, "complete", id), Binding: record.Binding, Asset: record.Asset, UploadId: id, UploadedSizeBytes: record.Result.SizeBytes, UploadedChecksumSha256: record.Result.SHA256})
|
||||
if err != nil {
|
||||
return record, fmt.Errorf("notify upload %s: %w", id, err)
|
||||
}
|
||||
if response.GetReceipt().GetResult() != agentv1.ResultCode_RESULT_CODE_ACCEPTED || response.GetState() != agentv1.UploadState_UPLOAD_STATE_COMPLETED {
|
||||
if response.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_ACCEPTED || response.GetState() != agentpb.UploadState_UPLOAD_STATE_COMPLETED {
|
||||
return record, errors.New("upload notification has not completed MQ delivery")
|
||||
}
|
||||
if err := spool.CompleteUploadNotification(id); err != nil {
|
||||
|
||||
@@ -12,28 +12,28 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
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"
|
||||
)
|
||||
|
||||
type uploadRPCStub struct {
|
||||
grant *agentv1.UploadGrant
|
||||
grant *agentpb.UploadGrant
|
||||
requests, notifications int
|
||||
pending bool
|
||||
}
|
||||
|
||||
func (s *uploadRPCStub) RequestUpload(_ context.Context, r *agentv1.RequestUploadRequest) (*agentv1.RequestUploadResponse, error) {
|
||||
func (s *uploadRPCStub) RequestUpload(_ context.Context, r *agentpb.RequestUploadRequest) (*agentpb.RequestUploadResponse, error) {
|
||||
s.requests++
|
||||
s.grant.UploadId = r.UploadId
|
||||
return &agentv1.RequestUploadResponse{Grant: s.grant, Receipt: &agentv1.OperationReceipt{Result: agentv1.ResultCode_RESULT_CODE_ACCEPTED}}, nil
|
||||
return &agentpb.RequestUploadResponse{Grant: s.grant, Receipt: &agentpb.OperationReceipt{Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED}}, nil
|
||||
}
|
||||
func (s *uploadRPCStub) CompleteUpload(_ context.Context, _ *agentv1.CompleteUploadRequest) (*agentv1.CompleteUploadResponse, error) {
|
||||
func (s *uploadRPCStub) CompleteUpload(_ context.Context, _ *agentpb.CompleteUploadRequest) (*agentpb.CompleteUploadResponse, error) {
|
||||
s.notifications++
|
||||
if s.pending {
|
||||
return nil, errors.New("notification pending")
|
||||
}
|
||||
return &agentv1.CompleteUploadResponse{State: agentv1.UploadState_UPLOAD_STATE_COMPLETED, Receipt: &agentv1.OperationReceipt{Result: agentv1.ResultCode_RESULT_CODE_ACCEPTED}}, nil
|
||||
return &agentpb.CompleteUploadResponse{State: agentpb.UploadState_UPLOAD_STATE_COMPLETED, Receipt: &agentpb.OperationReceipt{Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED}}, nil
|
||||
}
|
||||
|
||||
func TestRecordingNotificationRecoveryDoesNotPUTAgain(t *testing.T) {
|
||||
@@ -46,9 +46,9 @@ func TestRecordingNotificationRecoveryDoesNotPUTAgain(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sum := sha256.Sum256([]byte("audio"))
|
||||
binding := &agentv1.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a", CallId: "call-a"}
|
||||
asset := &agentv1.AssetDescriptor{AssetId: "recording-a", SizeBytes: 5, ChecksumSha256: hex.EncodeToString(sum[:])}
|
||||
remote := &uploadRPCStub{pending: true, grant: &agentv1.UploadGrant{ObjectKey: "recording-a", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}}
|
||||
binding := &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a", CallId: "call-a"}
|
||||
asset := &agentpb.AssetDescriptor{AssetId: "recording-a", SizeBytes: 5, ChecksumSha256: hex.EncodeToString(sum[:])}
|
||||
remote := &uploadRPCStub{pending: true, grant: &agentpb.UploadGrant{ObjectKey: "recording-a", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}}
|
||||
spool, err := agent.NewSpool(root, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
@@ -6,7 +6,7 @@ import (
|
||||
"net/http"
|
||||
"testing"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
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"
|
||||
"github.com/google/uuid"
|
||||
@@ -17,13 +17,13 @@ type mockFailureRecoveryClient struct {
|
||||
fail bool
|
||||
}
|
||||
|
||||
func (c *mockFailureRecoveryClient) ReportExecutionEvent(_ context.Context, req *agentv1.ReportExecutionEventRequest) (*agentv1.ReportExecutionEventResponse, error) {
|
||||
func (c *mockFailureRecoveryClient) ReportExecutionEvent(_ context.Context, req *agentpb.ReportExecutionEventRequest) (*agentpb.ReportExecutionEventResponse, error) {
|
||||
c.reports++
|
||||
if c.fail {
|
||||
return nil, errors.New("Dispatcher unavailable")
|
||||
}
|
||||
return &agentv1.ReportExecutionEventResponse{Receipt: &agentv1.OperationReceipt{
|
||||
Result: agentv1.ResultCode_RESULT_CODE_ACCEPTED, FactId: req.Fact.FactId, ContentSha256: req.Fact.ContentSha256,
|
||||
return &agentpb.ReportExecutionEventResponse{Receipt: &agentpb.OperationReceipt{
|
||||
Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED, FactId: req.Fact.FactId, ContentSha256: req.Fact.ContentSha256,
|
||||
}}, nil
|
||||
}
|
||||
|
||||
@@ -32,7 +32,7 @@ func TestMockFailureRecoveryWaitsForActivationAndRejectsRealMode(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
active := &agentv1.RequestMeta{
|
||||
active := &agentpb.RequestMeta{
|
||||
ProtocolVersion: "agent.v1", AgentId: "agent-a", CellId: "cell-a", BootId: "boot-a",
|
||||
DispatcherEpoch: "epoch-a", SessionGeneration: 1,
|
||||
}
|
||||
@@ -41,8 +41,8 @@ func TestMockFailureRecoveryWaitsForActivationAndRejectsRealMode(t *testing.T) {
|
||||
t.Helper()
|
||||
if err := spool.ClaimUpload(agent.UploadAttempt{
|
||||
UploadID: id, Identity: "identity-" + id, State: "attempted", RequestID: uuid.NewString(),
|
||||
Binding: &agentv1.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-" + id},
|
||||
Asset: &agentv1.AssetDescriptor{AssetId: "recording-" + id},
|
||||
Binding: &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-" + id},
|
||||
Asset: &agentpb.AssetDescriptor{AssetId: "recording-" + id},
|
||||
}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
@@ -53,12 +53,12 @@ func TestMockFailureRecoveryWaitsForActivationAndRejectsRealMode(t *testing.T) {
|
||||
}
|
||||
claimAndFail("upload-a")
|
||||
cfg := config.Config{Mode: "mock"}
|
||||
missingSession := func() (*agentv1.RequestMeta, error) { return nil, errors.New("not activated") }
|
||||
missingSession := func() (*agentpb.RequestMeta, error) { return nil, errors.New("not activated") }
|
||||
if err := recoverMockFailureNotifications(context.Background(), cfg, remote, spool, missingSession); err == nil || remote.reports != 1 {
|
||||
t.Fatalf("unactivated Agent reported a durable fact: reports=%d err=%v", remote.reports, err)
|
||||
}
|
||||
remote.fail = false
|
||||
currentSession := func() (*agentv1.RequestMeta, error) { return active, nil }
|
||||
currentSession := func() (*agentpb.RequestMeta, error) { return active, nil }
|
||||
if err := recoverMockFailureNotifications(context.Background(), cfg, remote, spool, currentSession); err != nil || remote.reports != 2 {
|
||||
t.Fatalf("Mock failure was not recovered through the active session: reports=%d err=%v", remote.reports, err)
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ import (
|
||||
"log/slog"
|
||||
"time"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
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"
|
||||
@@ -30,7 +30,7 @@ func recoverUploadNotifications(ctx context.Context, cfg config.Config, client r
|
||||
// recoverMockFailureNotifications never treats a prior boot's persisted fact
|
||||
// as a current session. The report uses metadata from the newly activated
|
||||
// session, while preserving the original immutable fact identity.
|
||||
func recoverMockFailureNotifications(ctx context.Context, cfg config.Config, client agent.MockFailureFactClient, spool *agent.Spool, activeMeta func() (*agentv1.RequestMeta, error)) error {
|
||||
func recoverMockFailureNotifications(ctx context.Context, cfg config.Config, client agent.MockFailureFactClient, spool *agent.Spool, activeMeta func() (*agentpb.RequestMeta, error)) error {
|
||||
pending, err := spool.PendingUploadFailures()
|
||||
if err != nil || len(pending) == 0 {
|
||||
return err
|
||||
@@ -50,7 +50,7 @@ func recoverMockFailureNotifications(ctx context.Context, cfg config.Config, cli
|
||||
|
||||
// startUploadNotificationRecovery never requests a token or opens a source file.
|
||||
// The worker owns only durable metadata-to-Dispatcher notifications.
|
||||
func startUploadNotificationRecovery(ctx context.Context, cfg config.Config, spool *agent.Spool, activeMeta func() (*agentv1.RequestMeta, error)) (func(), error) {
|
||||
func startUploadNotificationRecovery(ctx context.Context, cfg config.Config, spool *agent.Spool, activeMeta func() (*agentpb.RequestMeta, error)) (func(), error) {
|
||||
pending, err := spool.PendingUploadNotifications()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read upload recovery journal: %w", err)
|
||||
|
||||
@@ -8,7 +8,7 @@ import (
|
||||
"net/url"
|
||||
"strings"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
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"
|
||||
"github.com/google/uuid"
|
||||
@@ -34,7 +34,7 @@ func retryUploadRecording(ctx context.Context, cfg config.Config, client recordi
|
||||
if record.State != "attempted" || record.Binding == nil || record.Asset == nil || record.RequestID == "" || record.RequestID == requestID {
|
||||
return record, errors.New("only an unsuccessful upload may explicitly request a new grant")
|
||||
}
|
||||
raw, err := (proto.MarshalOptions{Deterministic: true}).Marshal(&agentv1.RequestUploadRequest{Binding: record.Binding, Asset: record.Asset, UploadId: uploadID})
|
||||
raw, err := (proto.MarshalOptions{Deterministic: true}).Marshal(&agentpb.RequestUploadRequest{Binding: record.Binding, Asset: record.Asset, UploadId: uploadID})
|
||||
if err != nil {
|
||||
return record, err
|
||||
}
|
||||
@@ -47,11 +47,11 @@ func retryUploadRecording(ctx context.Context, cfg config.Config, client recordi
|
||||
meta.OperationId = requestID
|
||||
meta.IdempotencyKey = requestID
|
||||
meta.TraceId = requestID
|
||||
response, err := client.RequestUpload(ctx, &agentv1.RequestUploadRequest{Meta: meta, Binding: record.Binding, Asset: record.Asset, UploadId: uploadID})
|
||||
response, err := client.RequestUpload(ctx, &agentpb.RequestUploadRequest{Meta: meta, Binding: record.Binding, Asset: record.Asset, UploadId: uploadID})
|
||||
if err != nil {
|
||||
return record, err
|
||||
}
|
||||
if response.GetReceipt().GetResult() != agentv1.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil {
|
||||
if response.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil {
|
||||
return record, errors.New("explicit upload request was not granted")
|
||||
}
|
||||
grant := response.Grant
|
||||
|
||||
@@ -12,7 +12,7 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
|
||||
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"
|
||||
)
|
||||
@@ -33,9 +33,9 @@ func TestExplicitUploadRetryIsOneNewRequestAndOnePUT(t *testing.T) {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sum := sha256.Sum256([]byte("audio"))
|
||||
binding := &agentv1.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a", CallId: "call-a"}
|
||||
asset := &agentv1.AssetDescriptor{AssetId: "recording-a", SizeBytes: 5, ChecksumSha256: hex.EncodeToString(sum[:])}
|
||||
remote := &uploadRPCStub{grant: &agentv1.UploadGrant{ObjectKey: "recording-a", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(-time.Minute).UnixMilli()}}
|
||||
binding := &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a", CallId: "call-a"}
|
||||
asset := &agentpb.AssetDescriptor{AssetId: "recording-a", SizeBytes: 5, ChecksumSha256: hex.EncodeToString(sum[:])}
|
||||
remote := &uploadRPCStub{grant: &agentpb.UploadGrant{ObjectKey: "recording-a", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(-time.Minute).UnixMilli()}}
|
||||
spool, err := agent.NewSpool(root, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
Reference in New Issue
Block a user