feat: complete MQ-only dispatcher and OSS upload flow

This commit is contained in:
2026-09-22 21:09:06 +08:00
parent 704652bd0d
commit a2fcc4aabb
265 changed files with 18689 additions and 1469 deletions
+27
View File
@@ -0,0 +1,27 @@
package main
import (
"context"
"log/slog"
"time"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
)
func runDispatcherControls(ctx context.Context, d *dispatcher.Dispatcher, controller dispatcher.TaskController, batch int) {
ticker := time.NewTicker(dispatcherOutboxFlushInterval)
defer ticker.Stop()
for {
if ctx.Err() != nil {
return
}
if _, err := d.ProcessTaskControls(ctx, controller, batch); err != nil && ctx.Err() == nil {
slog.Error("process Dispatcher task controls", "error", err)
}
select {
case <-ctx.Done():
return
case <-ticker.C:
}
}
}
@@ -0,0 +1,69 @@
package main
import (
"errors"
"os"
"path/filepath"
"strings"
"testing"
"git.ipao.vip/rogee/go-sip/internal/config"
"git.ipao.vip/rogee/go-sip/internal/store"
)
func TestDispatcherHasNoBusinessHTTPFlags(t *testing.T) {
command := newDispatcherCommand()
for _, name := range []string{"control-listen", "control-token"} {
if command.Flags().Lookup(name) != nil {
t.Fatalf("obsolete HTTP flag remains: --%s", name)
}
}
}
func TestDispatcherRequiresFileBeforeOpeningDatabase(t *testing.T) {
path := filepath.Join(t.TempDir(), "untouched.db")
root := newRootCommand()
root.SetArgs([]string{"dispatcher", "--mode", "mock", "--db", path})
if err := root.Execute(); err == nil || !strings.Contains(err.Error(), "--config") {
t.Fatalf("missing file did not fail early: %v", err)
}
if _, err := os.Stat(path); !errors.Is(err, os.ErrNotExist) {
t.Fatalf("database was touched before configuration validation: %v", err)
}
}
func TestDispatcherStoreRejectsChangedIdentityBeforeRecovery(t *testing.T) {
path := filepath.Join(t.TempDir(), "identity.db")
cfg := config.Config{DBPath: path, DispatcherID: "c046b893-8628-4589-ae50-619d049248a6"}
st, err := openDispatcherStore(cfg)
if err != nil {
t.Fatal(err)
}
// A claimed notice must not be reset by a process presenting another ID.
_, err = st.DB().Exec(`INSERT INTO outbox (event_id, tenant_key, exchange, routing_key, body, status, created_at) VALUES ('notice','key','exchange','route','{}','dispatching','2026-09-21T00:00:00Z')`)
if err != nil {
t.Fatal(err)
}
if err := st.Close(); err != nil {
t.Fatal(err)
}
cfg.DispatcherID = "a50b1569-aa17-4503-bfc1-cf55c10a1c24"
if other, err := openDispatcherStore(cfg); !errors.Is(err, store.ErrDispatcherIdentityMismatch) {
if other != nil {
other.Close()
}
t.Fatalf("wrong identity accepted: %v", err)
}
check, err := store.Open(path)
if err != nil {
t.Fatal(err)
}
defer check.Close()
var status string
if err := check.DB().QueryRow(`SELECT status FROM outbox WHERE event_id='notice'`).Scan(&status); err != nil {
t.Fatal(err)
}
if status != "dispatching" {
t.Fatalf("foreign identity changed recovery state: %s", status)
}
}
+27
View File
@@ -0,0 +1,27 @@
package main
import (
"context"
"errors"
"fmt"
amqp "github.com/rabbitmq/amqp091-go"
)
var errDispatcherIdentityLost = errors.New("Dispatcher MQ identity ownership lost")
func dispatcherIdentityContext(parent context.Context, done <-chan *amqp.Error) (context.Context, context.CancelFunc) {
ctx, cancel := context.WithCancelCause(parent)
go func() {
select {
case <-ctx.Done():
case err := <-done:
if err != nil {
cancel(fmt.Errorf("%w: %v", errDispatcherIdentityLost, err))
} else {
cancel(errDispatcherIdentityLost)
}
}
}()
return ctx, func() { cancel(context.Canceled) }
}
+42
View File
@@ -0,0 +1,42 @@
package main
import (
"context"
"errors"
"testing"
"time"
amqp "github.com/rabbitmq/amqp091-go"
)
func TestDispatcherIdentityLossCancelsAdmissionContext(t *testing.T) {
for _, closed := range []bool{false, true} {
done := make(chan *amqp.Error, 1)
ctx, cancel := dispatcherIdentityContext(context.Background(), done)
if closed {
close(done)
} else {
done <- &amqp.Error{Code: 320, Reason: "local test disconnect"}
}
select {
case <-ctx.Done():
case <-time.After(time.Second):
t.Fatal("identity loss did not stop admission")
}
if !errors.Is(context.Cause(ctx), errDispatcherIdentityLost) {
t.Fatalf("lost cause: %v", context.Cause(ctx))
}
cancel()
}
}
func TestDispatcherIdentityWatchStopsWithParent(t *testing.T) {
parent, cancelParent := context.WithCancel(context.Background())
ctx, cancel := dispatcherIdentityContext(parent, make(chan *amqp.Error))
defer cancel()
cancelParent()
<-ctx.Done()
if !errors.Is(context.Cause(ctx), context.Canceled) {
t.Fatal("normal cancellation reported identity loss")
}
}
+84 -55
View File
@@ -7,7 +7,6 @@ import (
"fmt"
"log/slog"
"net"
"net/http"
"os"
"os/signal"
"path/filepath"
@@ -24,7 +23,6 @@ import (
"git.ipao.vip/rogee/go-sip/internal/callwindow"
"git.ipao.vip/rogee/go-sip/internal/config"
"git.ipao.vip/rogee/go-sip/internal/contract"
"git.ipao.vip/rogee/go-sip/internal/control"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
"git.ipao.vip/rogee/go-sip/internal/health"
"git.ipao.vip/rogee/go-sip/internal/mq"
@@ -33,7 +31,9 @@ import (
"git.ipao.vip/rogee/go-sip/internal/store"
"github.com/spf13/cobra"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials"
"google.golang.org/grpc/status"
)
func main() {
@@ -98,6 +98,7 @@ func newAgentCommand() *cobra.Command {
cmd.Flags().StringVar(&cfg.CallAISnapshotPath, "call-ai-snapshot", cfg.CallAISnapshotPath, "immutable AI snapshot path")
cmd.Flags().IntVar(&cfg.CallMediaPort, "call-media-port", cfg.CallMediaPort, "ExternalMedia UDP port")
cmd.Flags().StringVar(&cfg.CallRecordingDirectory, "call-recording-dir", cfg.CallRecordingDirectory, "local recording directory")
cmd.AddCommand(newUploadRetryCommand())
return cmd
}
@@ -379,6 +380,11 @@ func serveAgentRPC(cfg config.Config, spool *agent.Spool, report agent.RecoveryR
agentv1.RegisterAgentControlServiceServer(grpcServer, handler)
serveCtx, cancel := signalContext()
defer cancel()
stopUploadRecovery, err := startUploadNotificationRecovery(serveCtx, cfg, spool)
if err != nil {
return err
}
defer stopUploadRecovery()
go func() {
<-serveCtx.Done()
grpcServer.GracefulStop()
@@ -473,8 +479,23 @@ func flushDispatcherOutbox(ctx context.Context, d *dispatcher.Dispatcher, batch
}
}
func openDispatcherStore(cfg config.Config) (*store.Store, error) {
st, err := store.Open(cfg.DBPath)
if err != nil {
return nil, err
}
if err := st.BindDispatcherID(cfg.DispatcherID); err != nil {
return nil, errors.Join(err, st.Close())
}
if err := st.RecoverOutbox(); err != nil {
return nil, errors.Join(err, st.Close())
}
return st, nil
}
func newDispatcherCommand() *cobra.Command {
cfg := config.FromEnv()
var configFile string
var once bool
var consume bool
var tenantKey string
@@ -482,16 +503,42 @@ func newDispatcherCommand() *cobra.Command {
Use: "dispatcher",
Short: "run the single-active Dispatcher process",
RunE: func(cmd *cobra.Command, _ []string) error {
if err := cfg.LoadDispatcherFile(configFile); err != nil {
return err
}
if err := cfg.Validate("dispatcher"); err != nil {
return err
}
st, err := store.Open(cfg.DBPath)
var publisher mq.Publisher
var broker *mq.Broker
if cfg.RabbitURL != "" {
var err error
broker, err = mq.Open(cfg.RabbitURL, cfg.DispatcherID)
if err != nil {
return err
}
defer broker.Close()
if tenantKey == "" {
return errors.New("--tenant-key is required with RabbitMQ")
}
if _, err := broker.DeclareTenantQueue(tenantKey); err != nil {
return err
}
publisher = broker
}
st, err := openDispatcherStore(cfg)
if err != nil {
return err
}
defer st.Close()
leaseCtx, cancelLease := signalContext()
defer cancelLease()
signalCtx, cancelSignal := signalContext()
defer cancelSignal()
leaseCtx := signalCtx
if broker != nil {
var cancelIdentity context.CancelFunc
leaseCtx, cancelIdentity = dispatcherIdentityContext(signalCtx, broker.Done())
defer cancelIdentity()
}
lease, err := dispatcher.StartLease(leaseCtx, st, "dispatcher-active-"+cfg.DispatcherID, "dispatcher", cfg.DispatcherID, 30*time.Second)
if err != nil {
return err
@@ -511,16 +558,6 @@ func newDispatcherCommand() *cobra.Command {
}
}()
}
var publisher mq.Publisher
var broker *mq.Broker
if cfg.RabbitURL != "" {
broker, err = mq.Open(cfg.RabbitURL, cfg.Exchange)
if err != nil {
return err
}
defer broker.Close()
publisher = broker
}
d, err := dispatcher.New(st, publisher, nil)
if err != nil {
return err
@@ -566,7 +603,12 @@ func newDispatcherCommand() *cobra.Command {
if err != nil {
return fmt.Errorf("listen Dispatcher gRPC: %w", err)
}
dispatcherGRPC = grpc.NewServer(grpc.Creds(credentials.NewTLS(tlsConfig)))
dispatcherGRPC = grpc.NewServer(grpc.Creds(credentials.NewTLS(tlsConfig)), grpc.UnaryInterceptor(func(ctx context.Context, req any, _ *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (any, error) {
if leaseCtx.Err() != nil {
return nil, status.Error(codes.Unavailable, "Dispatcher ownership is no longer active")
}
return handler(ctx, req)
}))
agentv1.RegisterAgentControlServiceServer(dispatcherGRPC, dispatcherHandler)
go func() {
if serveErr := dispatcherGRPC.Serve(dispatcherListener); serveErr != nil && !errors.Is(serveErr, grpc.ErrServerStopped) {
@@ -595,13 +637,7 @@ func newDispatcherCommand() *cobra.Command {
if once && consume {
return errors.New("--once cannot be combined with --consume")
}
if consume {
if broker == nil || tenantKey == "" {
return errors.New("--consume requires --rabbit-url and --tenant-key")
}
if cfg.ControlListen != "" {
return errors.New("--consume cannot be combined with --control-listen")
}
if publisher != nil && !once && (consume || dispatcherGRPC != nil) {
flushCtx, cancelFlush := context.WithCancel(leaseCtx)
flushDone := make(chan struct{})
go func() {
@@ -612,9 +648,26 @@ func newDispatcherCommand() *cobra.Command {
cancelFlush()
<-flushDone
}()
}
if connectedAgents != nil && !once && (consume || dispatcherGRPC != nil) {
controlCtx, cancelControl := context.WithCancel(leaseCtx)
controlDone := make(chan struct{})
go func() {
defer close(controlDone)
runDispatcherControls(controlCtx, d, connectedAgents.coordinator, cfg.OutboxBatch)
}()
defer func() { cancelControl(); <-controlDone }()
}
if consume {
if broker == nil || tenantKey == "" {
return errors.New("--consume requires --rabbit-url and --tenant-key")
}
if err := d.ConsumeTenant(leaseCtx, broker, tenantKey); err != nil && !errors.Is(err, context.Canceled) {
return err
}
if errors.Is(context.Cause(leaseCtx), errDispatcherIdentityLost) {
return context.Cause(leaseCtx)
}
return nil
}
if once {
@@ -627,49 +680,25 @@ func newDispatcherCommand() *cobra.Command {
}
result["published"] = count
}
if cfg.ControlListen == "" {
if dispatcherGRPC == nil {
return writeResult(result)
}
select {
case <-leaseCtx.Done():
return nil
case err := <-lease.Lost():
return fmt.Errorf("dispatcher lease lost: %w", err)
}
}
if once {
return errors.New("--once cannot be combined with --control-listen")
}
server := &http.Server{Addr: cfg.ControlListen, Handler: control.Handler{Store: st, BearerToken: cfg.ControlToken}}
go func() {
select {
case <-leaseCtx.Done():
case <-lease.Lost():
}
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 5*time.Second)
defer shutdownCancel()
_ = server.Shutdown(shutdownCtx)
}()
if err := server.ListenAndServe(); !errors.Is(err, http.ErrServerClosed) {
return err
if dispatcherGRPC == nil {
return writeResult(result)
}
select {
case <-leaseCtx.Done():
if errors.Is(context.Cause(leaseCtx), errDispatcherIdentityLost) {
return context.Cause(leaseCtx)
}
return nil
case err := <-lease.Lost():
return fmt.Errorf("dispatcher lease lost: %w", err)
default:
return nil
}
},
}
cmd.Flags().StringVar(&cfg.Mode, "mode", cfg.Mode, "mock, mixed, or real")
cmd.Flags().StringVar(&cfg.DBPath, "db", cfg.DBPath, "Dispatcher SQLite path")
cmd.Flags().StringVar(&cfg.DispatcherID, "dispatcher-id", cfg.DispatcherID, "single-active Dispatcher holder identity")
cmd.Flags().StringVar(&configFile, "config", "", "required strict Dispatcher JSON configuration file")
cmd.Flags().StringVar(&cfg.RabbitURL, "rabbit-url", cfg.RabbitURL, "RabbitMQ URL")
cmd.Flags().StringVar(&cfg.Exchange, "exchange", cfg.Exchange, "durable command exchange")
cmd.Flags().IntVar(&cfg.OutboxBatch, "outbox-batch", cfg.OutboxBatch, "maximum outbox messages per run")
cmd.Flags().StringVar(&cfg.ControlListen, "control-listen", cfg.ControlListen, "internal control HTTP listen address; empty disables server")
cmd.Flags().StringVar(&cfg.ControlToken, "control-token", cfg.ControlToken, "bearer token for internal control HTTP")
cmd.Flags().StringVar(&cfg.AgentEndpointsFile, "agent-endpoints-file", cfg.AgentEndpointsFile, "strict JSON file of Dispatcher-owned Agent endpoints")
cmd.Flags().StringVar(&tenantKey, "tenant-key", "", "tenant key to consume from its command queue")
cmd.Flags().BoolVar(&consume, "consume", false, "consume one tenant command queue")
+6 -2
View File
@@ -10,6 +10,7 @@ import (
"git.ipao.vip/rogee/go-sip/internal/contract"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
"git.ipao.vip/rogee/go-sip/internal/store"
"git.ipao.vip/rogee/go-sip/internal/testfixture"
)
func TestRootHasExplicitRoles(t *testing.T) {
@@ -86,16 +87,19 @@ func TestFlushDispatcherOutboxPublishesPending(t *testing.T) {
t.Fatal(err)
}
defer st.Close()
if err := st.BindDispatcherID(testfixture.DispatcherID); err != nil {
t.Fatal(err)
}
d, err := dispatcher.New(st, recordingPublisher{}, time.Now)
if err != nil {
t.Fatal(err)
}
raw, err := contracts.Read("examples/call.execute.json")
raw, err := testfixture.Execute()
if err != nil {
t.Fatal(err)
}
if _, err := d.AcceptCommand(raw, "agent-call.tenant.tenant-demo-key.call.execute"); err != nil {
if _, err := d.AcceptCommand(raw, testfixture.InboundKey("tenant-demo-key")); err != nil {
t.Fatal(err)
}
+10 -37
View File
@@ -6,7 +6,6 @@ import (
"encoding/hex"
"errors"
"fmt"
"net/url"
"path/filepath"
"strings"
"time"
@@ -55,6 +54,10 @@ func uploadCallRecordings(ctx context.Context, cfg config.Config, result callrun
return nil, fmt.Errorf("dial Dispatcher gRPC service: %w", err)
}
defer client.Close()
spool, err := agent.NewSpool(cfg.SpoolRoot, time.Now)
if err != nil {
return nil, err
}
uploader := agent.UploadClient{Now: time.Now}
facts := append(append([]callruntime.RecordingFact(nil), result.InboundRecordings...), result.OutboundRecordings...)
if len(facts) == 0 {
@@ -76,46 +79,16 @@ func uploadCallRecordings(ctx context.Context, cfg config.Config, result callrun
ChecksumSha256: recording.SHA256,
Channels: 1,
SampleRateHz: 16000,
DurationMs: recording.DurationMS,
}
uploadID := stableUploadID(binding, asset)
meta := uploadMeta(cfg, "request", uploadID)
grantResponse, err := client.RequestUpload(ctx, &agentv1.RequestUploadRequest{Meta: meta, Binding: binding, Asset: asset, UploadId: uploadID})
record, err := uploadRecording(ctx, cfg, client, uploader, spool, binding, asset, recording.Path)
if err != nil {
return nil, fmt.Errorf("request upload %s: %w", uploadID, err)
}
if grantResponse == nil || grantResponse.Grant == nil || grantResponse.Receipt == nil || grantResponse.Receipt.Result != agentv1.ResultCode_RESULT_CODE_ACCEPTED {
return nil, fmt.Errorf("request upload %s was rejected", uploadID)
}
parsed, err := url.Parse(grantResponse.Grant.TargetUrl)
if err != nil || parsed.Host == "" {
return nil, fmt.Errorf("upload %s returned invalid target URL", uploadID)
}
uploader.AllowedHosts = map[string]struct{}{strings.ToLower(parsed.Host): {}}
uploadResult, err := uploader.UploadFile(ctx, grantResponse.Grant, recording.Path)
if err != nil {
return nil, fmt.Errorf("upload %s data plane: %w", uploadID, err)
}
if uploadResult.SizeBytes != asset.SizeBytes || !strings.EqualFold(uploadResult.SHA256, asset.ChecksumSha256) {
return nil, fmt.Errorf("upload %s local result does not match asset", uploadID)
}
completeResponse, err := client.CompleteUpload(ctx, &agentv1.CompleteUploadRequest{
Meta: uploadMeta(cfg, "complete", uploadID),
Binding: binding,
Asset: asset,
UploadId: uploadID,
UploadedSizeBytes: uploadResult.SizeBytes,
UploadedChecksumSha256: uploadResult.SHA256,
})
if err != nil {
return nil, fmt.Errorf("complete upload %s: %w", uploadID, err)
}
if completeResponse == nil || completeResponse.Receipt == nil || completeResponse.Receipt.Result != agentv1.ResultCode_RESULT_CODE_ACCEPTED || completeResponse.OssId == "" {
return nil, fmt.Errorf("complete upload %s was not verified", uploadID)
return nil, err
}
uploaded = append(uploaded, map[string]any{
"upload_id": uploadID, "asset_id": asset.AssetId, "object_key": grantResponse.Grant.ObjectKey,
"oss_id": completeResponse.OssId, "bytes": uploadResult.SizeBytes, "sha256": uploadResult.SHA256,
"status_code": uploadResult.StatusCode,
"upload_id": record.UploadID, "asset_id": asset.AssetId, "object_key": record.ObjectKey,
"bytes": record.Result.SizeBytes, "sha256": record.Result.SHA256,
"status_code": record.Result.StatusCode, "notification_state": record.State,
})
}
return uploaded, nil
+121
View File
@@ -0,0 +1,121 @@
package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"net/url"
"os"
"strings"
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
"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)
}
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) {
if binding == nil || asset == nil {
return returned, errors.New("upload binding and asset are required")
}
id := stableUploadID(binding, asset)
lock, err := spool.LockUpload(id)
if err != nil {
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})
if err != nil {
return agent.UploadAttempt{}, err
}
digest := sha256.Sum256(identityBytes)
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})
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 {
return record, fmt.Errorf("request upload %s has no accepted grant", id)
}
grant := response.Grant
if grant.UploadId != id {
return record, errors.New("upload grant identity mismatch")
}
parsed, err := url.Parse(grant.TargetUrl)
if err != nil || parsed.Host == "" {
return record, errors.New("upload grant has invalid target URL")
}
record = agent.UploadAttempt{UploadID: id, Identity: identity, State: "attempted", ObjectKey: grant.ObjectKey, Binding: binding, Asset: asset, RequestID: uploadMeta(cfg, "request", id).OperationId}
// Durable exclusive claim must precede any network PUT, including an attempt
// that ends in an ambiguous transport failure.
if err := spool.ClaimUpload(record); err != nil {
return record, err
}
uploader.AllowedHosts = map[string]struct{}{strings.ToLower(parsed.Host): {}}
result, err := uploader.UploadFile(ctx, grant, path)
if err != nil {
return record, fmt.Errorf("upload %s data plane: %w", id, err)
}
if result.SizeBytes != asset.SizeBytes || !strings.EqualFold(result.SHA256, asset.ChecksumSha256) {
return record, errors.New("upload result does not match asset")
}
if err := spool.RecordUploadResult(id, result); err != nil {
return record, err
}
record.State, record.Result = "uploaded", result
} else if err != nil {
return record, err
} else if record.Identity != identity {
return record, errors.New("persisted upload binding or asset mismatch")
} else if record.State == "attempted" {
return record, errors.New("prior PUT outcome unknown or failed; automatic re-upload is forbidden")
}
if record.State == "completed" {
return record, nil
}
return notifyUploadedRecordingLocked(ctx, cfg, client, spool, record)
}
func notifyUploadedRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool, record agent.UploadAttempt) (returned agent.UploadAttempt, returnErr error) {
lock, err := spool.LockUpload(record.UploadID)
if err != nil {
return returned, err
}
defer func() { returnErr = errors.Join(returnErr, lock.Close()) }()
current, err := spool.LoadUploadAttempt(record.UploadID)
if err != nil {
return returned, err
}
if current.State == "completed" {
return current, nil
}
return notifyUploadedRecordingLocked(ctx, cfg, client, spool, current)
}
func notifyUploadedRecordingLocked(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool, record agent.UploadAttempt) (agent.UploadAttempt, error) {
if record.State != "uploaded" || record.Binding == nil || record.Asset == nil {
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})
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 {
return record, errors.New("upload notification has not completed MQ delivery")
}
if err := spool.CompleteUploadNotification(id); err != nil {
return record, err
}
record.State = "completed"
return record, nil
}
+84
View File
@@ -0,0 +1,84 @@
package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"testing"
"time"
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
"git.ipao.vip/rogee/go-sip/internal/agent"
"git.ipao.vip/rogee/go-sip/internal/config"
)
type uploadRPCStub struct {
grant *agentv1.UploadGrant
requests, notifications int
pending bool
}
func (s *uploadRPCStub) RequestUpload(_ context.Context, r *agentv1.RequestUploadRequest) (*agentv1.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
}
func (s *uploadRPCStub) CompleteUpload(_ context.Context, _ *agentv1.CompleteUploadRequest) (*agentv1.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
}
func TestRecordingNotificationRecoveryDoesNotPUTAgain(t *testing.T) {
puts := 0
server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { puts++; w.WriteHeader(http.StatusOK) }))
defer server.Close()
root := t.TempDir()
path := filepath.Join(root, "audio.wav")
if err := os.WriteFile(path, []byte("audio"), 0600); err != nil {
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()}}
spool, err := agent.NewSpool(root, nil)
if err != nil {
t.Fatal(err)
}
uploader := agent.UploadClient{HTTPClient: server.Client()}
if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, spool, binding, asset, path); err == nil {
t.Fatal("pending MQ notification reported complete")
}
remote.pending = false
restarted, err := agent.NewSpool(root, nil)
if err != nil {
t.Fatal(err)
}
if err := recoverUploadNotifications(context.Background(), config.Config{}, remote, restarted); err != nil {
t.Fatal(err)
}
record, err := restarted.LoadUploadAttempt(stableUploadID(binding, asset))
if err != nil {
t.Fatal(err)
}
if record.State != "completed" || puts != 1 || remote.requests != 1 || remote.notifications != 2 {
t.Fatalf("state=%s PUT=%d grant=%d notification=%d", record.State, puts, remote.requests, remote.notifications)
}
if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, restarted, binding, asset, path); err != nil {
t.Fatal(err)
}
if puts != 1 || remote.notifications != 2 {
t.Fatal("completed upload repeated side effects")
}
if _, err := os.Stat(path); err != nil {
t.Fatal("source recording removed")
}
}
+73
View File
@@ -0,0 +1,73 @@
package main
import (
"context"
"errors"
"fmt"
"log/slog"
"time"
"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 recoverUploadNotifications(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool) error {
pending, err := spool.PendingUploadNotifications()
if err != nil {
return err
}
var failures []error
for _, record := range pending {
if _, err := notifyUploadedRecording(ctx, cfg, client, spool, record); err != nil {
failures = append(failures, err)
}
}
return errors.Join(failures...)
}
// startUploadNotificationRecovery never requests a token or opens a source file.
// The worker owns only the durable metadata-to-Dispatcher notification path.
func startUploadNotificationRecovery(ctx context.Context, cfg config.Config, spool *agent.Spool) (func(), error) {
pending, err := spool.PendingUploadNotifications()
if err != nil {
return nil, fmt.Errorf("read upload recovery journal: %w", err)
}
if cfg.DispatcherGRPCEndpoint == "" {
if len(pending) > 0 {
return nil, errors.New("pending upload notifications require Dispatcher gRPC endpoint")
}
return func() {}, nil
}
client, err := rpc.DialFromFiles(cfg.DispatcherGRPCEndpoint, cfg.MTLSCAFile, cfg.MTLSCertFile, cfg.MTLSKeyFile, cfg.DispatcherGRPCServerName)
if err != nil {
return nil, err
}
workerCtx, cancel := context.WithCancel(ctx)
done := make(chan struct{})
go func() {
defer close(done)
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
for {
attemptCtx, finish := context.WithTimeout(workerCtx, 10*time.Second)
err := recoverUploadNotifications(attemptCtx, cfg, client, spool)
finish()
if err != nil && workerCtx.Err() == nil {
slog.Error("upload notification recovery failed; original facts retained", "error", err)
}
select {
case <-workerCtx.Done():
return
case <-ticker.C:
}
}
}()
return func() {
cancel()
<-done
if err := client.Close(); err != nil {
slog.Error("close upload notification client", "error", err)
}
}, nil
}
+83
View File
@@ -0,0 +1,83 @@
package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"net/url"
"strings"
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
"git.ipao.vip/rogee/go-sip/internal/agent"
"git.ipao.vip/rogee/go-sip/internal/config"
"github.com/google/uuid"
"google.golang.org/protobuf/proto"
)
// retryUploadRecording is reachable only through an explicit operator command.
// A supplied request ID is stable across redelivery and can authorize one PUT.
func retryUploadRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, uploader agent.UploadClient, spool *agent.Spool, uploadID, requestID, path string) (returned agent.UploadAttempt, returnErr error) {
parsedID, err := uuid.Parse(requestID)
if err != nil || parsedID.Version() != 4 || parsedID.Variant() != uuid.RFC4122 || parsedID.String() != requestID {
return returned, errors.New("explicit request ID must be a canonical UUID v4")
}
lock, err := spool.LockUpload(uploadID)
if err != nil {
return returned, err
}
defer func() { returnErr = errors.Join(returnErr, lock.Close()) }()
record, err := spool.LoadUploadAttempt(uploadID)
if err != nil {
return record, err
}
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})
if err != nil {
return record, err
}
sum := sha256.Sum256(raw)
if record.Identity != hex.EncodeToString(sum[:]) || stableUploadID(record.Binding, record.Asset) != uploadID {
return record, errors.New("persisted upload identity mismatch")
}
meta := uploadMeta(cfg, "request", uploadID)
meta.RequestId = requestID
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})
if err != nil {
return record, err
}
if response.GetReceipt().GetResult() != agentv1.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil {
return record, errors.New("explicit upload request was not granted")
}
grant := response.Grant
if grant.UploadId != uploadID || grant.ObjectKey != record.ObjectKey || grant.RequiredChecksumSha256 != record.Asset.ChecksumSha256 || grant.MaxBytes < record.Asset.SizeBytes {
return record, errors.New("replacement grant does not match the original asset")
}
target, err := url.Parse(grant.TargetUrl)
if err != nil || target.Host == "" {
return record, errors.New("replacement grant has invalid target")
}
if err := spool.ReserveUploadRetry(uploadID, requestID); err != nil {
return record, err
}
uploader.AllowedHosts = map[string]struct{}{strings.ToLower(target.Host): {}}
result, err := uploader.UploadFile(ctx, grant, path)
if err != nil {
return record, err
}
if result.SizeBytes != record.Asset.SizeBytes || result.SHA256 != record.Asset.ChecksumSha256 {
return record, errors.New("replacement upload result does not match original asset")
}
if err := spool.RecordUploadResult(uploadID, result); err != nil {
return record, err
}
record.RequestID = requestID
record.State = "uploaded"
record.Result = result
return notifyUploadedRecordingLocked(ctx, cfg, client, spool, record)
}
+50
View File
@@ -0,0 +1,50 @@
package main
import (
"errors"
"time"
"git.ipao.vip/rogee/go-sip/internal/agent"
"git.ipao.vip/rogee/go-sip/internal/config"
"git.ipao.vip/rogee/go-sip/internal/rpc"
"github.com/spf13/cobra"
)
func newUploadRetryCommand() *cobra.Command {
cfg := config.FromEnv()
var uploadID, requestID, path string
cmd := &cobra.Command{
Use: "upload-retry", Short: "explicitly request one new grant for an unsuccessful upload; never dial",
Args: cobra.NoArgs,
RunE: func(cmd *cobra.Command, _ []string) (returnErr error) {
if uploadID == "" || requestID == "" || path == "" {
return errors.New("--upload-id, --request-id and --file are required")
}
if err := cfg.Validate("agent"); err != nil {
return err
}
if cfg.DispatcherGRPCEndpoint == "" || cfg.DispatcherGRPCServerName == "" {
return errors.New("Dispatcher gRPC endpoint and server name are required")
}
spool, err := agent.NewSpool(cfg.SpoolRoot, time.Now)
if err != nil {
return err
}
client, err := rpc.DialFromFiles(cfg.DispatcherGRPCEndpoint, cfg.MTLSCAFile, cfg.MTLSCertFile, cfg.MTLSKeyFile, cfg.DispatcherGRPCServerName)
if err != nil {
return err
}
defer func() { returnErr = errors.Join(returnErr, client.Close()) }()
record, err := retryUploadRecording(cmd.Context(), cfg, client, agent.UploadClient{Now: time.Now}, spool, uploadID, requestID, path)
if err != nil {
return err
}
return writeResult(map[string]any{"upload_id": record.UploadID, "notification_state": record.State, "bytes": record.Result.SizeBytes, "sha256": record.Result.SHA256})
},
}
cmd.Flags().StringVar(&cfg.SpoolRoot, "spool", cfg.SpoolRoot, "existing Agent spool root")
cmd.Flags().StringVar(&uploadID, "upload-id", "", "existing failed or uncertain upload identity")
cmd.Flags().StringVar(&requestID, "request-id", "", "explicit new canonical UUID v4; never reuse a consumed request")
cmd.Flags().StringVar(&path, "file", "", "retained original recording; content must match its persisted identity")
return cmd
}
+70
View File
@@ -0,0 +1,70 @@
package main
import (
"context"
"crypto/sha256"
"encoding/hex"
"io"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"testing"
"time"
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
"git.ipao.vip/rogee/go-sip/internal/agent"
"git.ipao.vip/rogee/go-sip/internal/config"
)
func TestExplicitUploadRetryIsOneNewRequestAndOnePUT(t *testing.T) {
puts := 0
server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
puts++
if _, err := io.Copy(io.Discard, r.Body); err != nil {
t.Error(err)
}
w.WriteHeader(http.StatusOK)
}))
defer server.Close()
root := t.TempDir()
path := filepath.Join(root, "audio.wav")
if err := os.WriteFile(path, []byte("audio"), 0600); err != nil {
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()}}
spool, err := agent.NewSpool(root, nil)
if err != nil {
t.Fatal(err)
}
uploader := agent.UploadClient{HTTPClient: server.Client()}
if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, spool, binding, asset, path); err == nil {
t.Fatal("expired token accepted")
}
if puts != 0 || remote.requests != 1 {
t.Fatal("expired grant caused PUT or auto renewal")
}
remote.grant.ExpiresAtUnixMs = time.Now().Add(15 * time.Minute).UnixMilli()
id := stableUploadID(binding, asset)
if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, spool, binding, asset, path); err == nil {
t.Fatal("ordinary recovery retried a failed attempt")
}
if _, err := retryUploadRecording(context.Background(), config.Config{}, remote, uploader, spool, id, "11111111-1111-4111-8111-111111111111", path); err != nil {
t.Fatal(err)
}
if puts != 1 || remote.requests != 2 || remote.notifications != 1 {
t.Fatalf("PUT=%d grants=%d notifications=%d", puts, remote.requests, remote.notifications)
}
if _, err := retryUploadRecording(context.Background(), config.Config{}, remote, uploader, spool, id, "22222222-2222-4222-8222-222222222222", path); err == nil {
t.Fatal("completed upload retried")
}
if puts != 1 || remote.requests != 2 {
t.Fatal("duplicate retry repeated effects")
}
if _, err := os.Stat(path); err != nil {
t.Fatal("retry removed source file")
}
}