170 lines
7.1 KiB
Go
170 lines
7.1 KiB
Go
package agent
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"log"
|
|
"time"
|
|
|
|
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
|
"git.ipao.vip/rogee/go-sip/internal/session"
|
|
|
|
"google.golang.org/grpc/codes"
|
|
"google.golang.org/grpc/status"
|
|
"google.golang.org/protobuf/proto"
|
|
)
|
|
|
|
var ErrRecordingClientUnavailable = errors.New("Agent recording delivery requires an active Dispatcher session and Unary client")
|
|
|
|
// RecordingClient belongs to one approved call. It requests upload tokens only
|
|
// when its caller explicitly asks, and it never retries an OSS PUT or invents a
|
|
// successful call end or result after an RPC error.
|
|
type RecordingClient struct {
|
|
Client agentpb.AgentControlServiceClient
|
|
Session func(context.Context) (*agentpb.RequestMeta, error)
|
|
DispatcherID string
|
|
TenantID int64
|
|
SourceEventID string
|
|
}
|
|
|
|
func (c RecordingClient) requestMeta(ctx context.Context, action string) (*agentpb.RequestMeta, error) {
|
|
if c.Client == nil || c.Session == nil || c.DispatcherID == "" || c.TenantID <= 0 || c.SourceEventID == "" {
|
|
return nil, ErrRecordingClientUnavailable
|
|
}
|
|
if err := ctx.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
meta, err := c.Session(ctx)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("obtain current Agent session: %w", err)
|
|
}
|
|
if meta == nil || meta.GetAgentId() == "" || meta.GetCellId() == "" || meta.GetBootId() == "" || meta.GetDispatcherEpoch() == "" || meta.GetSessionGeneration() == 0 {
|
|
return nil, ErrRecordingClientUnavailable
|
|
}
|
|
copy := proto.Clone(meta).(*agentpb.RequestMeta)
|
|
copy.OperationId = c.SourceEventID + "/" + action
|
|
copy.IdempotencyKey = copy.OperationId
|
|
return copy, nil
|
|
}
|
|
|
|
func (c RecordingClient) RequestUpload(ctx context.Context, asset *agentpb.AssetDescriptor, uploadID string) (*agentpb.UploadGrant, error) {
|
|
meta, err := c.requestMeta(ctx, "upload/"+uploadID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if asset == nil || uploadID == "" || asset.GetKind() != agentpb.AssetKind_ASSET_KIND_RECORDING ||
|
|
asset.GetExecutionId() != c.SourceEventID || asset.GetCallId() != c.SourceEventID || asset.GetAssetId() == "" || asset.GetSizeBytes() <= 0 || asset.GetChecksumSha256() == "" {
|
|
return nil, errors.New("recording token request does not describe the approved call and original asset")
|
|
}
|
|
originalAsset := proto.Clone(asset).(*agentpb.AssetDescriptor)
|
|
response, err := reReportVerifiedSessionCutover(ctx, c, meta, "upload/"+uploadID, func(callCtx context.Context, current *agentpb.RequestMeta) (*agentpb.RequestRecordingUploadResponse, error) {
|
|
return c.Client.RequestRecordingUpload(callCtx, &agentpb.RequestRecordingUploadRequest{
|
|
Meta: current, DispatcherId: c.DispatcherID, TenantId: c.TenantID, SourceEventId: c.SourceEventID,
|
|
UploadId: uploadID, Asset: proto.Clone(originalAsset).(*agentpb.AssetDescriptor),
|
|
})
|
|
})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("request original recording upload token: %w", err)
|
|
}
|
|
grant := response.GetGrant()
|
|
if grant == nil || grant.GetUploadId() != uploadID || grant.GetBucket() == "" || grant.GetObjectKey() == "" ||
|
|
grant.GetMaxBytes() != originalAsset.GetSizeBytes() || grant.GetRequiredChecksumSha256() != originalAsset.GetChecksumSha256() {
|
|
return nil, fmt.Errorf("%w: Dispatcher returned a different recording asset", ErrUploadGrantInvalid)
|
|
}
|
|
return grant, nil
|
|
}
|
|
|
|
func (c RecordingClient) ReportEnded(ctx context.Context) error {
|
|
meta, err := c.requestMeta(ctx, "ended")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
response, err := reReportVerifiedSessionCutover(ctx, c, meta, "ended", func(callCtx context.Context, current *agentpb.RequestMeta) (*agentpb.ReportCallEndedResponse, error) {
|
|
return c.Client.ReportCallEnded(callCtx, &agentpb.ReportCallEndedRequest{
|
|
Meta: current, DispatcherId: c.DispatcherID, TenantId: c.TenantID, SourceEventId: c.SourceEventID,
|
|
})
|
|
})
|
|
if err != nil {
|
|
return fmt.Errorf("report confirmed call end: %w", err)
|
|
}
|
|
if receipt := response.GetReceipt(); receipt.GetResult() != agentpb.ResultCode_RESULT_CODE_APPLIED || receipt.GetFactId() != c.SourceEventID {
|
|
return errors.New("Dispatcher did not persist the confirmed call end")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (c RecordingClient) ReportFinal(ctx context.Context, payload []byte, upload *agentpb.UploadObservation) (string, error) {
|
|
meta, err := c.requestMeta(ctx, "result")
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
if len(payload) == 0 {
|
|
return "", errors.New("final call result payload is required")
|
|
}
|
|
originalPayload := bytes.Clone(payload)
|
|
var originalUpload *agentpb.UploadObservation
|
|
if upload != nil {
|
|
originalUpload = proto.Clone(upload).(*agentpb.UploadObservation)
|
|
}
|
|
response, err := reReportVerifiedSessionCutover(ctx, c, meta, "result", func(callCtx context.Context, current *agentpb.RequestMeta) (*agentpb.ReportCallResultResponse, error) {
|
|
request := &agentpb.ReportCallResultRequest{
|
|
Meta: current, DispatcherId: c.DispatcherID, TenantId: c.TenantID, SourceEventId: c.SourceEventID,
|
|
ResultPayloadJson: bytes.Clone(originalPayload),
|
|
}
|
|
if originalUpload != nil {
|
|
request.Upload = proto.Clone(originalUpload).(*agentpb.UploadObservation)
|
|
}
|
|
return c.Client.ReportCallResult(callCtx, request)
|
|
})
|
|
if err != nil {
|
|
return "", fmt.Errorf("persist unique call result: %w", err)
|
|
}
|
|
if receipt := response.GetReceipt(); receipt.GetResult() == agentpb.ResultCode_RESULT_CODE_ACCEPTED && receipt.GetFactId() != "" {
|
|
return receipt.GetFactId(), nil
|
|
}
|
|
return "", errors.New("Dispatcher did not persist the unique call result")
|
|
}
|
|
|
|
func reReportVerifiedSessionCutover[T any](ctx context.Context, c RecordingClient, original *agentpb.RequestMeta, action string, send func(context.Context, *agentpb.RequestMeta) (T, error)) (T, error) {
|
|
response, err := send(ctx, original)
|
|
if !isVerifiedSessionTransition(err) {
|
|
return response, err
|
|
}
|
|
log.Printf("Agent fact fenced during session transition event=%s agent=%s generation=%d action=%s", c.SourceEventID, original.AgentId, original.SessionGeneration, action)
|
|
retryCtx, cancel := context.WithTimeout(ctx, 6*time.Second)
|
|
defer cancel()
|
|
attempts := 1
|
|
for {
|
|
timer := time.NewTimer(100 * time.Millisecond)
|
|
select {
|
|
case <-retryCtx.Done():
|
|
timer.Stop()
|
|
log.Printf("Agent fact re-report deadline event=%s agent=%s action=%s attempts=%d", c.SourceEventID, original.AgentId, action, attempts)
|
|
var zero T
|
|
return zero, fmt.Errorf("Agent fact re-report not confirmed: %w", errors.Join(err, retryCtx.Err()))
|
|
case <-timer.C:
|
|
}
|
|
current, sessionErr := c.requestMeta(retryCtx, action)
|
|
if sessionErr != nil {
|
|
var zero T
|
|
return zero, fmt.Errorf("obtain renewed Agent session for original fact: %w", sessionErr)
|
|
}
|
|
if current.AgentId != original.AgentId || current.CellId != original.CellId || current.BootId != original.BootId || current.DispatcherEpoch != original.DispatcherEpoch {
|
|
var zero T
|
|
return zero, fmt.Errorf("%w: Agent identity changed during session transition", ErrRecordingClientUnavailable)
|
|
}
|
|
response, err = send(retryCtx, current)
|
|
attempts++
|
|
if !isVerifiedSessionTransition(err) {
|
|
return response, err
|
|
}
|
|
}
|
|
}
|
|
|
|
func isVerifiedSessionTransition(err error) bool {
|
|
st, ok := status.FromError(err)
|
|
return ok && st.Code() == codes.Unauthenticated && st.Message() == session.GenerationTransitionMessage
|
|
}
|