feat: implement local P1 contracts and Mock call flow

This commit is contained in:
2026-09-25 16:38:22 +08:00
parent ef84a0663d
commit 0b4b901d85
174 changed files with 18044 additions and 1740 deletions
+28
View File
@@ -0,0 +1,28 @@
package main
import (
"context"
"errors"
"log/slog"
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
)
// 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 {
if mode != "mock" {
return nil
}
return func(ctx context.Context, request *agentv1.ExecuteAuthorizedRequest) error {
if err := ctx.Err(); err != nil {
return err
}
if request == nil || request.Binding == nil || request.Binding.ExecutionId == "" {
return errors.New("mock authorized origination has no execution identity")
}
slog.Info("isolated Mock Agent simulated authorized originate; no SIP sent",
"execution_id", request.Binding.ExecutionId, "trunk_id", request.SelectedTrunkId)
return nil
}
}
@@ -0,0 +1,27 @@
package main
import (
"context"
"testing"
agentv1 "git.ipao.vip/rogee/go-sip/gen/agent/v1"
)
func TestAuthorizedOriginatorIsExplicitlyMockOnly(t *testing.T) {
if mockAuthorizedOriginator("mixed") != nil || mockAuthorizedOriginator("real") != nil || mockAuthorizedOriginator("") != nil {
t.Fatal("non-mock mode exposed an authorized mock originator")
}
originator := mockAuthorizedOriginator("mock")
if originator == nil {
t.Fatal("mock mode has no isolated originator")
}
request := &agentv1.ExecuteAuthorizedRequest{Binding: &agentv1.ExecutionBinding{ExecutionId: "execution-1"}, SelectedTrunkId: "trunk-1"}
ctx, cancel := context.WithCancel(context.Background())
cancel()
if err := originator(ctx, request); err == nil {
t.Fatal("cancelled mock invocation appeared to complete")
}
if err := originator(context.Background(), request); err != nil {
t.Fatalf("isolated mock invocation: %v", err)
}
}
-27
View File
@@ -1,27 +0,0 @@
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:
}
}
}
+67 -81
View File
@@ -7,6 +7,7 @@ import (
"fmt"
"log/slog"
"net"
"net/http"
"os"
"os/signal"
"path/filepath"
@@ -22,6 +23,7 @@ import (
"git.ipao.vip/rogee/go-sip/internal/callruntime"
"git.ipao.vip/rogee/go-sip/internal/callwindow"
"git.ipao.vip/rogee/go-sip/internal/config"
"git.ipao.vip/rogee/go-sip/internal/configread"
"git.ipao.vip/rogee/go-sip/internal/contract"
"git.ipao.vip/rogee/go-sip/internal/dispatcher"
"git.ipao.vip/rogee/go-sip/internal/health"
@@ -376,11 +378,12 @@ func serveAgentRPC(cfg config.Config, spool *agent.Spool, report agent.RecoveryR
PeerCertificateFingerprints: peerCertificateFingerprints,
StatePath: filepath.Join(cfg.SpoolRoot, "rpc-session.json"),
CallLogger: callLogger,
MockAuthorizedOriginate: mockAuthorizedOriginator(cfg.Mode),
})
agentv1.RegisterAgentControlServiceServer(grpcServer, handler)
serveCtx, cancel := signalContext()
defer cancel()
stopUploadRecovery, err := startUploadNotificationRecovery(serveCtx, cfg, spool)
stopUploadRecovery, err := startUploadNotificationRecovery(serveCtx, cfg, spool, handler.ActiveSessionMeta)
if err != nil {
return err
}
@@ -451,34 +454,6 @@ func (r *connectedAgentRuntime) Close() error {
return firstErr
}
const dispatcherOutboxFlushInterval = 250 * time.Millisecond
func flushDispatcherOutbox(ctx context.Context, d *dispatcher.Dispatcher, batch int, interval time.Duration) {
flush := func() {
published, err := d.FlushOutbox(ctx, batch)
if err != nil {
if ctx.Err() == nil {
slog.Error("flush Dispatcher outbox", "error", err)
}
return
}
if published > 0 {
slog.Debug("flushed Dispatcher outbox", "published", published)
}
}
flush()
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
flush()
}
}
}
func openDispatcherStore(cfg config.Config) (*store.Store, error) {
st, err := store.Open(cfg.DBPath)
if err != nil {
@@ -497,8 +472,6 @@ func newDispatcherCommand() *cobra.Command {
cfg := config.FromEnv()
var configFile string
var once bool
var consume bool
var tenantKey string
cmd := &cobra.Command{
Use: "dispatcher",
Short: "run the single-active Dispatcher process",
@@ -509,21 +482,18 @@ func newDispatcherCommand() *cobra.Command {
if err := cfg.Validate("dispatcher"); err != nil {
return err
}
if cfg.Mode != "mock" && !once {
return errors.New("Dispatcher V3 authorized execution and upload are isolated Mock only; mixed/real requires separate approval and implementation")
}
var publisher mq.Publisher
var broker *mq.Broker
var broker *mq.V3Broker
if cfg.RabbitURL != "" {
var err error
broker, err = mq.Open(cfg.RabbitURL, cfg.DispatcherID)
broker, err = mq.OpenV3(cfg.RabbitURL, cfg.DispatcherID, 1)
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)
@@ -558,7 +528,7 @@ func newDispatcherCommand() *cobra.Command {
}
}()
}
d, err := dispatcher.New(st, publisher, nil)
d, err := dispatcher.NewV3(cfg.DispatcherID, st, publisher, time.Now)
if err != nil {
return err
}
@@ -581,8 +551,26 @@ 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 {
if connectedAgents == nil || connectedAgents.coordinator == nil {
return fmt.Errorf("no activated Agent coordinator: %w", store.ErrCommandConflict)
}
return connectedAgents.coordinator.AuthorizeInboundMeta(meta)
}
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 {
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) {
if err := authorizeAgentSession(req.Meta); err != nil {
return false, err
}
return d.CompleteLocalMockRecording(ctx, req, record)
},
})
if handlerErr != nil {
return handlerErr
@@ -590,6 +578,8 @@ func newDispatcherCommand() *cobra.Command {
eventHandler, eventErr := rpc.NewDispatcherEventServer(st, rpc.DispatcherEventServerOptions{
RequirePeer: true, PeerCertificateFingerprints: peerFingerprints,
AllowedAgentIDs: allowedAgentIDs, Now: time.Now,
LocalV3RecordingFailure: d.RecordLocalMockRecordingFailure,
LocalV3SessionCheck: authorizeAgentSession,
})
if eventErr != nil {
return eventErr
@@ -634,55 +624,45 @@ func newDispatcherCommand() *cobra.Command {
}
result["agent_sessions"] = sessions
}
if once && consume {
return errors.New("--once cannot be combined with --consume")
}
if publisher != nil && !once && (consume || dispatcherGRPC != nil) {
flushCtx, cancelFlush := context.WithCancel(leaseCtx)
flushDone := make(chan struct{})
go func() {
defer close(flushDone)
flushDispatcherOutbox(flushCtx, d, cfg.OutboxBatch, dispatcherOutboxFlushInterval)
}()
defer func() {
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 {
if publisher == nil {
return errors.New("--once requires RABBITMQ_URL or a configured publisher")
return errors.New("--once requires RABBITMQ_URL")
}
count, err := d.FlushOutbox(cmd.Context(), cfg.OutboxBatch)
if err != nil {
return err
}
result["published"] = count
}
if dispatcherGRPC == nil {
return writeResult(result)
}
if broker == nil {
return errors.New("RABBITMQ_URL is required for the Dispatcher V3 runtime")
}
if connectedAgents == nil || len(connectedAgents.sessions) != 1 {
return errors.New("the Dispatcher V3 runtime requires exactly one configured Agent")
}
configClient, err := configread.NewClient(os.Getenv("DISPATCHER_CONFIG_READ_BASE_URL"), cfg.DispatcherID,
os.Getenv("DISPATCHER_SECRET_KEY"), &http.Client{Timeout: 15 * time.Second})
if err != nil {
return fmt.Errorf("configure Dispatcher config-read client: %w", err)
}
session := connectedAgents.sessions[0]
verifier := &dispatcher.AgentSIPConfigVerifier{
Probe: connectedAgents.coordinator, AgentID: session.AgentID, CellID: session.CellID,
}
runtime, err := dispatcher.NewLocalV01Runtime(d, broker, configClient, verifier, connectedAgents.coordinator)
if err != nil {
return err
}
if err := runtime.EnableMockAuthorizedOrigination(session.AgentID, connectedAgents.coordinator); err != nil {
return fmt.Errorf("enable isolated Mock Agent execution: %w", err)
}
result["queue_protocol"] = "v3"
if err := writeResult(result); err != nil {
return err
}
runtimeDone := make(chan error, 1)
go func() { runtimeDone <- runtime.Run(leaseCtx) }()
select {
case <-leaseCtx.Done():
if errors.Is(context.Cause(leaseCtx), errDispatcherIdentityLost) {
@@ -691,6 +671,14 @@ func newDispatcherCommand() *cobra.Command {
return nil
case err := <-lease.Lost():
return fmt.Errorf("dispatcher lease lost: %w", err)
case err := <-runtimeDone:
if err != nil {
return err
}
if leaseCtx.Err() != nil {
return nil
}
return errors.New("Dispatcher V3 runtime stopped unexpectedly")
}
},
}
@@ -700,8 +688,6 @@ func newDispatcherCommand() *cobra.Command {
cmd.Flags().StringVar(&cfg.RabbitURL, "rabbit-url", cfg.RabbitURL, "RabbitMQ URL")
cmd.Flags().IntVar(&cfg.OutboxBatch, "outbox-batch", cfg.OutboxBatch, "maximum outbox messages per run")
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")
cmd.Flags().BoolVar(&once, "once", false, "flush one outbox batch and exit")
return cmd
}
-46
View File
@@ -2,15 +2,10 @@ package main
import (
"bytes"
"context"
"testing"
"time"
"git.ipao.vip/rogee/go-sip/contracts"
"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) {
@@ -76,44 +71,3 @@ func TestBuildCallEndpointUsesArtifactRoute(t *testing.T) {
t.Fatal("disallowed target was accepted")
}
}
type recordingPublisher struct{}
func (recordingPublisher) Publish(context.Context, string, string, []byte) error { return nil }
func TestFlushDispatcherOutboxPublishesPending(t *testing.T) {
st, err := store.Open(":memory:")
if err != nil {
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 := testfixture.Execute()
if err != nil {
t.Fatal(err)
}
if _, err := d.AcceptCommand(raw, testfixture.InboundKey("tenant-demo-key")); err != nil {
t.Fatal(err)
}
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
go flushDispatcherOutbox(ctx, d, 1, time.Millisecond)
deadline := time.Now().Add(time.Second)
for time.Now().Before(deadline) {
var status string
if err := st.DB().QueryRow(`SELECT status FROM outbox LIMIT 1`).Scan(&status); err == nil && status == "published" {
return
}
time.Sleep(time.Millisecond)
}
t.Fatal("pending outbox was not published")
}
@@ -0,0 +1,75 @@
package main
import (
"context"
"errors"
"net/http"
"testing"
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"
)
type mockFailureRecoveryClient struct {
reports int
fail bool
}
func (c *mockFailureRecoveryClient) ReportExecutionEvent(_ context.Context, req *agentv1.ReportExecutionEventRequest) (*agentv1.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,
}}, nil
}
func TestMockFailureRecoveryWaitsForActivationAndRejectsRealMode(t *testing.T) {
spool, err := agent.NewSpool(t.TempDir(), nil)
if err != nil {
t.Fatal(err)
}
active := &agentv1.RequestMeta{
ProtocolVersion: "agent.v1", AgentId: "agent-a", CellId: "cell-a", BootId: "boot-a",
DispatcherEpoch: "epoch-a", SessionGeneration: 1,
}
remote := &mockFailureRecoveryClient{fail: true}
claimAndFail := func(id string) {
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},
}); err != nil {
t.Fatal(err)
}
known, err := spool.ReportMockUploadFailure(context.Background(), remote, active, id, &agent.UploadHTTPError{StatusCode: http.StatusForbidden})
if !known || err == nil {
t.Fatalf("failed Mock PUT was not retained for recovery: known=%t err=%v", known, err)
}
}
claimAndFail("upload-a")
cfg := config.Config{Mode: "mock"}
missingSession := func() (*agentv1.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 }
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)
}
pending, err := spool.PendingUploadFailures()
if err != nil || len(pending) != 0 {
t.Fatalf("acknowledged failure was still pending: count=%d err=%v", len(pending), err)
}
remote.fail = true
claimAndFail("upload-b")
cfg.Mode = "real"
if err := recoverMockFailureNotifications(context.Background(), cfg, remote, spool, currentSession); err == nil || remote.reports != 3 {
t.Fatalf("Mock failure escaped into real mode: reports=%d err=%v", remote.reports, err)
}
}
+37 -5
View File
@@ -7,6 +7,7 @@ import (
"log/slog"
"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"
"git.ipao.vip/rogee/go-sip/internal/rpc"
@@ -26,16 +27,44 @@ func recoverUploadNotifications(ctx context.Context, cfg config.Config, client r
return errors.Join(failures...)
}
// 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 {
pending, err := spool.PendingUploadFailures()
if err != nil || len(pending) == 0 {
return err
}
if cfg.Mode != "mock" {
return fmt.Errorf("pending Mock recording failure cannot be reported in %q mode", cfg.Mode)
}
if activeMeta == nil {
return errors.New("pending Mock recording failure requires an active Agent session")
}
meta, err := activeMeta()
if err != nil {
return fmt.Errorf("load active Agent session for Mock recording failure: %w", err)
}
return spool.RecoverMockUploadFailures(ctx, client, meta)
}
// 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) {
// 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) {
pending, err := spool.PendingUploadNotifications()
if err != nil {
return nil, fmt.Errorf("read upload recovery journal: %w", err)
}
failureFacts, err := spool.PendingUploadFailures()
if err != nil {
return nil, fmt.Errorf("read Mock failure recovery journal: %w", err)
}
if len(failureFacts) > 0 && cfg.Mode != "mock" {
return nil, errors.New("pending Mock recording failure cannot start in mixed/real mode")
}
if cfg.DispatcherGRPCEndpoint == "" {
if len(pending) > 0 {
return nil, errors.New("pending upload notifications require Dispatcher gRPC endpoint")
if len(pending) > 0 || len(failureFacts) > 0 {
return nil, errors.New("pending upload facts require Dispatcher gRPC endpoint")
}
return func() {}, nil
}
@@ -51,7 +80,10 @@ func startUploadNotificationRecovery(ctx context.Context, cfg config.Config, spo
defer ticker.Stop()
for {
attemptCtx, finish := context.WithTimeout(workerCtx, 10*time.Second)
err := recoverUploadNotifications(attemptCtx, cfg, client, spool)
err := errors.Join(
recoverUploadNotifications(attemptCtx, cfg, client, spool),
recoverMockFailureNotifications(attemptCtx, cfg, client, spool, activeMeta),
)
finish()
if err != nil && workerCtx.Err() == nil {
slog.Error("upload notification recovery failed; original facts retained", "error", err)