release reservations on explicit Agent refusal before dial
This commit is contained in:
@@ -5,11 +5,14 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
"git.ipao.vip/rogee/go-sip/internal/contract"
|
||||
"git.ipao.vip/rogee/go-sip/internal/store"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
// CallSpec is a frozen instruction handed to one Agent. Its Snapshot
|
||||
@@ -205,8 +208,17 @@ func (c *ExecuteController) dispatchPending(ctx context.Context, cmd store.Execu
|
||||
Deadline: choice.Deadline, Snapshot: snapshot,
|
||||
}
|
||||
if err := c.Originator.Originate(ctx, spec); err != nil {
|
||||
unknownErr := c.Store.MarkExecuteUnknown(cmd.DispatcherID, cmd.EventID)
|
||||
return errors.Join(fmt.Errorf("originator outcome unknown for call %q: %w", cmd.EventID, err), unknownErr)
|
||||
switch status.Code(err) {
|
||||
case codes.FailedPrecondition, codes.InvalidArgument, codes.PermissionDenied:
|
||||
rejectErr := c.Store.RejectUnissuedExecute(cmd.DispatcherID, cmd.EventID, "Agent refused before outbound call")
|
||||
if rejectErr == nil {
|
||||
log.Printf("Dispatcher call refused before dial: event_id=%q agent_code=%s", cmd.EventID, status.Code(err))
|
||||
}
|
||||
return errors.Join(fmt.Errorf("Agent rejected call %q before dialing: %w", cmd.EventID, err), rejectErr)
|
||||
default:
|
||||
unknownErr := c.Store.MarkExecuteUnknown(cmd.DispatcherID, cmd.EventID)
|
||||
return errors.Join(fmt.Errorf("originator outcome unknown for call %q: %w", cmd.EventID, err), unknownErr)
|
||||
}
|
||||
}
|
||||
if err := c.Store.MarkExecuteDispatched(cmd.DispatcherID, cmd.EventID); err != nil {
|
||||
return fmt.Errorf("originated call %q has no durable acknowledgment: %w", cmd.EventID, err)
|
||||
|
||||
@@ -12,6 +12,8 @@ import (
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
"git.ipao.vip/rogee/go-sip/internal/contract"
|
||||
"git.ipao.vip/rogee/go-sip/internal/store"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
type fakeOriginator struct {
|
||||
@@ -136,6 +138,26 @@ func TestExecuteWaitsForRulesThenDispatchesOriginalIdentity(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentExplicitRefusalBeforeDialReleasesReservationAndKeepsUnknownBlocked(t *testing.T) {
|
||||
controller, originator, _, s := newExecuteFixture(t)
|
||||
originator.err = status.Error(codes.FailedPrecondition, "Agent refused before issue")
|
||||
body := executeBody(t, "agent-refused-1", "15003164745")
|
||||
if err := controller.ProcessExecute(context.Background(), body); err == nil {
|
||||
t.Fatal("Agent refusal hidden")
|
||||
}
|
||||
outbox, err := s.ListPendingOutbox(controller.DispatcherID)
|
||||
if err != nil || len(outbox) != 1 || !strings.Contains(string(outbox[0].Body), `"status":"rejected"`) {
|
||||
t.Fatalf("definite refusal lacks durable rejection receipt: %+v %v", outbox, err)
|
||||
}
|
||||
if err := controller.ProcessExecute(context.Background(), body); err != nil || len(originator.calls) != 1 {
|
||||
t.Fatalf("refused command retried: %d %v", len(originator.calls), err)
|
||||
}
|
||||
originator.err = nil
|
||||
if err := controller.ProcessExecute(context.Background(), executeBody(t, "after-refusal-1", "15003164745")); err != nil || len(originator.calls) != 2 {
|
||||
t.Fatalf("definite pre-dial failure held capacity: %d %v", len(originator.calls), err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExecuteTimeoutAndMQFailureDoNotRedial(t *testing.T) {
|
||||
controller, originator, publisher, s := newExecuteFixture(t)
|
||||
originator.err = context.DeadlineExceeded
|
||||
|
||||
@@ -304,6 +304,15 @@ func (s *Store) RejectExecute(dispatcherID, eventID, reason string) error {
|
||||
return s.closeExecuteWithAck(dispatcherID, eventID, "pending", "rejected", map[string]any{"status": "rejected", "reason_code": nil, "reason_message": reason})
|
||||
}
|
||||
|
||||
// RejectUnissuedExecute is only for an explicit Agent refusal before it
|
||||
// accepted the call; transport/ARI-unknown failures must remain occupied.
|
||||
func (s *Store) RejectUnissuedExecute(dispatcherID, eventID, reason string) error {
|
||||
if reason == "" {
|
||||
return errors.New("pre-dial rejection reason is required")
|
||||
}
|
||||
return s.closeExecuteWithAck(dispatcherID, eventID, "dispatching", "rejected", map[string]any{"status": "rejected", "reason_code": nil, "reason_message": reason})
|
||||
}
|
||||
|
||||
func (s *Store) MarkExecuteDispatched(dispatcherID, eventID string) error {
|
||||
return s.closeExecuteWithAck(dispatcherID, eventID, "dispatching", "dispatched", map[string]any{"status": "dispatched"})
|
||||
}
|
||||
|
||||
@@ -101,6 +101,35 @@ func TestExecuteRejectsOnlyOnceWithoutFinalResult(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAgentExplicitPreDialRejectionReleasesReservedCapacity(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
at := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC)
|
||||
cmd := currentCall("call-agent-rejected")
|
||||
if _, _, err := s.RecordExecute(cmd); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.ReserveExecute(cmd.DispatcherID, cmd.EventID, currentReservation(), at); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.RejectUnissuedExecute(cmd.DispatcherID, cmd.EventID, "Agent refused before outbound call"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
outbox, err := s.ListPendingOutbox(cmd.DispatcherID)
|
||||
if err != nil || len(outbox) != 1 || !strings.Contains(string(outbox[0].Body), `"status":"rejected"`) {
|
||||
t.Fatalf("pre-dial failure must produce one durable honest receipt: %+v %v", outbox, err)
|
||||
}
|
||||
if err := s.RejectUnissuedExecute(cmd.DispatcherID, cmd.EventID, "Agent refused before outbound call"); err == nil {
|
||||
t.Fatal("duplicate rejection cannot create a second receipt")
|
||||
}
|
||||
if _, created, err := s.RecordExecute(cmd); err != nil || created {
|
||||
t.Fatalf("rejected call was reaccepted: %v %v", created, err)
|
||||
}
|
||||
occupied, err := s.TrunkOccupancy(cmd.DispatcherID)
|
||||
if err != nil || occupied["trunk-mock"] != 0 {
|
||||
t.Fatalf("definite pre-dial failure still occupies capacity: %+v %v", occupied, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExecuteQuotaCountsUnknownAndReleasesOnlyConfirmedEnd(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
at := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC)
|
||||
|
||||
Reference in New Issue
Block a user