diff --git a/cmd/sip-go-agent/current_dispatcher_command.go b/cmd/sip-go-agent/current_dispatcher_command.go index a376e92..ee33281 100644 --- a/cmd/sip-go-agent/current_dispatcher_command.go +++ b/cmd/sip-go-agent/current_dispatcher_command.go @@ -82,7 +82,7 @@ func runCurrentDispatcher(ctx context.Context, mode string) (result error) { if err != nil { return fmt.Errorf("initialize bounded OSS grant signer: %w", err) } - database, err := store.OpenCurrent(settings.SQLitePath) + database, err := store.Open(settings.SQLitePath) if err != nil { return err } diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index d777ad7..349e410 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -105,6 +105,7 @@ - SQLite 代码归属:当前 Store 所需的 SQLite 驱动和 Dispatcher 身份/命令冲突错误已迁到当前入口;随后删除旧 Store 的源码、旧测试和旧迁移 SQL 文件。当前 Store 的独立建表、拒绝旧布局及外来 Dispatcher 抢占测试继续通过;这些操作仅移除仓库内过时源码,不读取、改写或清理现存 SQLite、spool 或 outbox 业务数据。 - MQ 旧任务队列分支:旧 Store 的任务队列引用删除后,移除 `V3Broker` 和三份只验证旧拓扑的代码/测试;共享结果队列和无配置权限下的投递仍由当前 MQ 隔离测试验证。真实 SaaS/MQ 接收与应用收讫仍未验证。 - MQ 旧声明拓扑分支:移除旧 `Broker`、旧租户 Topic 路由/队列声明及其测试;独立保留当前消费端的消息处理器、永久错误分类与 Dispatcher UUID v4 校验,补错误解包和规范身份的回归测试。当前 Broker 仅被动核对 SaaS 预建拓扑;隔离 MQ 测试通过不代表真实 SaaS 应用收讫。 +- Store 名称收敛:移除现行 Store 的自有 `Current*` 类型、错误和 `OpenCurrent` 入口,保留单一 `Store`/`Open`;20 份 Go 源码与测试文件改为不带代次的路径,原现行任务/录音/结果测试仍执行。只重命名源码与调用,不修改现有 SQLite 表、记录或恢复数据。 ## 验收台账 diff --git a/internal/dispatcher/current_ai_test.go b/internal/dispatcher/current_ai_test.go index 5665def..d21975a 100644 --- a/internal/dispatcher/current_ai_test.go +++ b/internal/dispatcher/current_ai_test.go @@ -55,7 +55,7 @@ func TestCurrentBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(t *test if err != nil { t.Fatal(err) } - db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + db, err := store.Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } @@ -79,7 +79,7 @@ func TestCurrentBootstrapRejectsSDKUnsupportedTaskBeforeOpeningAdmission(t *test } func TestCurrentExecuteCorruptAISnapshotFailsClosedBeforeReservation(t *testing.T) { - db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + db, err := store.Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_config.go b/internal/dispatcher/current_config.go index 93fcec6..ca4b147 100644 --- a/internal/dispatcher/current_config.go +++ b/internal/dispatcher/current_config.go @@ -16,7 +16,7 @@ import ( // opening task admission fails. type CurrentBootstrap struct { Client *configread.Client - Store *store.CurrentStore + Store *store.Store DispatcherID string VerifySIP func(context.Context, configread.CurrentSIP) error DrainControls func(context.Context) error diff --git a/internal/dispatcher/current_config_test.go b/internal/dispatcher/current_config_test.go index 21b025c..54152fc 100644 --- a/internal/dispatcher/current_config_test.go +++ b/internal/dispatcher/current_config_test.go @@ -56,7 +56,7 @@ func TestCurrentBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { if err != nil { t.Fatal(err) } - db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "current.db")) + db, err := store.Open(filepath.Join(t.TempDir(), "current.db")) if err != nil { t.Fatal(err) } @@ -91,7 +91,7 @@ func TestCurrentBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { if err := db.NoteSIPChange(id, 9); err != nil { t.Fatal(err) } - if err := bootstrap.Run(context.Background()); !errors.Is(err, store.ErrCurrentSIPPending) { + if err := bootstrap.Run(context.Background()); !errors.Is(err, store.ErrSIPPending) { t.Fatalf("startup lost durable newer SIP notification: %v", err) } if cursor != "opaque-end-token" || approvedSIP.Revision != 8 { @@ -134,7 +134,7 @@ func TestCurrentBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { if err != nil { t.Fatal(err) } - db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "current.db")) + db, err := store.Open(filepath.Join(t.TempDir(), "current.db")) if err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_control.go b/internal/dispatcher/current_control.go index e652bd8..a00a96a 100644 --- a/internal/dispatcher/current_control.go +++ b/internal/dispatcher/current_control.go @@ -29,7 +29,7 @@ type CurrentControlAgent interface { type CurrentControlController struct { DispatcherID string - Store *store.CurrentStore + Store *store.Store Client *configread.Client Agent CurrentControlAgent VerifySIP func(context.Context, configread.CurrentSIP) error @@ -108,7 +108,7 @@ func (c *CurrentControlController) ProcessControl(ctx context.Context, body []by } } if err := c.Store.PrepareControl(event.DispatcherID, event.TenantID, event.Payload.TaskID, event.Payload.Action); err != nil { - if errors.Is(err, store.ErrCurrentControlRejected) { + if errors.Is(err, store.ErrControlRejected) { if rejectErr := c.Store.RejectControl(event.DispatcherID, event.TenantID, event.EventID); rejectErr != nil { return fmt.Errorf("persist rejected control %q: %w", event.EventID, rejectErr) } diff --git a/internal/dispatcher/current_control_test.go b/internal/dispatcher/current_control_test.go index 2f1063e..bd99ef5 100644 --- a/internal/dispatcher/current_control_test.go +++ b/internal/dispatcher/current_control_test.go @@ -31,7 +31,7 @@ func (a *currentFakeControlAgent) SendControl(_ context.Context, spec CurrentCon return a.err } -func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *currentFakeControlAgent, *store.CurrentStore) { +func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *currentFakeControlAgent, *store.Store) { t.Helper() id := "c046b893-8628-4589-ae50-619d049248a6" snapshot := currentPolicySnapshot(t) @@ -63,7 +63,7 @@ func newCurrentControlFixture(t *testing.T) (*CurrentControlController, *current if err != nil { t.Fatal(err) } - s, err := store.OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + s, err := store.Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } @@ -133,7 +133,7 @@ func TestCurrentControlPauseDefaultHangupAckAfterAgentDispatch(t *testing.T) { func TestCurrentControlStopCannotResumeOrDispatchPending(t *testing.T) { controller, agent, s := newCurrentControlFixture(t) - command := store.CurrentExecuteCommand{DispatcherID: controller.DispatcherID, EventID: "queued-1", TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:00:00Z"} + command := store.ExecuteCommand{DispatcherID: controller.DispatcherID, EventID: "queued-1", TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:00:00Z"} if _, _, err := s.RecordExecute(command); err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_discovery.go b/internal/dispatcher/current_discovery.go index e29c5ca..870733b 100644 --- a/internal/dispatcher/current_discovery.go +++ b/internal/dispatcher/current_discovery.go @@ -16,7 +16,7 @@ import ( type CurrentDiscoveryFollower struct { DispatcherID string Client *configread.Client - Store *store.CurrentStore + Store *store.Store ApprovedSIP configread.CurrentSIP VerifySIP func(context.Context, configread.CurrentSIP) error Cursor string diff --git a/internal/dispatcher/current_execute.go b/internal/dispatcher/current_execute.go index eb271a5..ab04032 100644 --- a/internal/dispatcher/current_execute.go +++ b/internal/dispatcher/current_execute.go @@ -41,7 +41,7 @@ type CurrentPublisher interface { type CurrentExecuteController struct { DispatcherID string - Store *store.CurrentStore + Store *store.Store Originator CurrentOriginator Publisher CurrentPublisher Now func() time.Time @@ -82,13 +82,13 @@ func (c *CurrentExecuteController) ProcessExecute(ctx context.Context, body []by if event.EventType != "call.execute" || event.Payload.TaskID == "" || event.Payload.Callee == "" || event.DispatcherID != c.DispatcherID { return errors.New("call instruction event type, task, or Dispatcher owner mismatch") } - cmd := store.CurrentExecuteCommand{ + cmd := store.ExecuteCommand{ DispatcherID: event.DispatcherID, EventID: event.EventID, TenantID: event.TenantID, TaskID: event.Payload.TaskID, Callee: event.Payload.Callee, IssuedAt: event.IssuedAt, } stored, _, err := c.Store.RecordExecute(cmd) - if errors.Is(err, store.ErrCurrentStopped) { + if errors.Is(err, store.ErrStopped) { // The stopped task never accepted this old queued command. RabbitMQ // may ACK it without creating an execution or a per-call result. return nil @@ -127,7 +127,7 @@ func (c *CurrentExecuteController) ProcessPending(ctx context.Context) error { return errors.Join(failures...) } -func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd store.CurrentExecuteCommand) error { +func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd store.ExecuteCommand) error { at := c.Now() issued, err := time.Parse(time.RFC3339Nano, cmd.IssuedAt) if err != nil { @@ -158,13 +158,13 @@ func (c *CurrentExecuteController) dispatchPending(ctx context.Context, cmd stor if err != nil { return fmt.Errorf("task %q rule validation: %w", cmd.TaskID, err) } - reservation := store.CurrentCallReservation{ + reservation := store.CallReservation{ TrunkID: choice.TrunkID, SIPRevision: snapshot.SIP.Revision, CallerID: choice.CallerID, DialedCallee: choice.DialedCallee, Deadline: choice.Deadline, } if err := c.Store.ReserveExecute(cmd.DispatcherID, cmd.EventID, reservation, at); err != nil { - if errors.Is(err, store.ErrCurrentCapacity) || errors.Is(err, store.ErrCurrentNotReady) || errors.Is(err, store.ErrCurrentAlreadyStarted) { + if errors.Is(err, store.ErrCapacity) || errors.Is(err, store.ErrNotReady) || errors.Is(err, store.ErrAlreadyStarted) { return nil } return fmt.Errorf("reserve call %q: %w", cmd.EventID, err) diff --git a/internal/dispatcher/current_execute_test.go b/internal/dispatcher/current_execute_test.go index b15122f..d007b86 100644 --- a/internal/dispatcher/current_execute_test.go +++ b/internal/dispatcher/current_execute_test.go @@ -44,9 +44,9 @@ func (f *currentFakePublisher) Publish(_ context.Context, exchange, key string, return f.err } -func newCurrentExecuteFixture(t *testing.T) (*CurrentExecuteController, *currentFakeOriginator, *currentFakePublisher, *store.CurrentStore) { +func newCurrentExecuteFixture(t *testing.T) (*CurrentExecuteController, *currentFakeOriginator, *currentFakePublisher, *store.Store) { t.Helper() - s, err := store.OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + s, err := store.Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_gate_test.go b/internal/dispatcher/current_gate_test.go index 16cd9c7..4bb884e 100644 --- a/internal/dispatcher/current_gate_test.go +++ b/internal/dispatcher/current_gate_test.go @@ -58,7 +58,7 @@ func TestCurrentFutureIssuedAtNeverDialsEarly(t *testing.T) { if err := controller.ProcessPending(context.Background()); err != nil || len(originator.calls) != 1 { t.Fatalf("future command did not dispatch after issued_at: %d %v", len(originator.calls), err) } - if _, _, err := s.RecordExecute(store.CurrentExecuteCommand{DispatcherID: controller.DispatcherID, EventID: "future-1", TenantID: 1001, TaskID: "task-asr", Callee: "15830461047", IssuedAt: "2026-09-21T01:00:00Z"}); err == nil { + if _, _, err := s.RecordExecute(store.ExecuteCommand{DispatcherID: controller.DispatcherID, EventID: "future-1", TenantID: 1001, TaskID: "task-asr", Callee: "15830461047", IssuedAt: "2026-09-21T01:00:00Z"}); err == nil { t.Fatal("conflicting identity silently accepted") } } diff --git a/internal/dispatcher/current_runtime.go b/internal/dispatcher/current_runtime.go index 8d5dee1..c541e1f 100644 --- a/internal/dispatcher/current_runtime.go +++ b/internal/dispatcher/current_runtime.go @@ -88,7 +88,7 @@ func (r *CurrentRuntime) Serve(ctx context.Context) (result error) { return nil } if err := r.Bootstrap.Run(ctx); err != nil { - if !errors.Is(err, store.ErrCurrentSIPPending) { + if !errors.Is(err, store.ErrSIPPending) { return fmt.Errorf("bootstrap current Dispatcher: %w", err) } r.Logger.Warn("control and outbox stay active while newer SIP revision waits; task admission remains closed", "dispatcher_id", r.Bootstrap.DispatcherID, "error", err) @@ -248,7 +248,7 @@ func (r *CurrentRuntime) syncTaskConsumers(ctx context.Context, consumers map[st if err != nil { return err } - wanted := make(map[string]store.CurrentAssignedTask, len(assigned)) + wanted := make(map[string]store.AssignedTask, len(assigned)) for _, task := range assigned { route, err := tenant.CurrentTaskRoute(r.Bootstrap.DispatcherID, task.TaskID) if err != nil { diff --git a/internal/dispatcher/current_runtime_integration_test.go b/internal/dispatcher/current_runtime_integration_test.go index 06c514f..1f2035c 100644 --- a/internal/dispatcher/current_runtime_integration_test.go +++ b/internal/dispatcher/current_runtime_integration_test.go @@ -132,7 +132,7 @@ func TestCurrentRuntimeIsolatedControlBacklogExecuteAndSharedResult(t *testing.T if err != nil { t.Fatal(err) } - db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "dispatcher.db")) + db, err := store.Open(filepath.Join(t.TempDir(), "dispatcher.db")) if err != nil { t.Fatal(err) } diff --git a/internal/dispatcher/current_sip_conflict_test.go b/internal/dispatcher/current_sip_conflict_test.go index 965f100..2bd6a40 100644 --- a/internal/dispatcher/current_sip_conflict_test.go +++ b/internal/dispatcher/current_sip_conflict_test.go @@ -47,7 +47,7 @@ func TestCurrentBootstrapRejectsChangedSIPContentUnderSameRevision(t *testing.T) if err != nil { t.Fatal(err) } - db, err := store.OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + db, err := store.Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } diff --git a/internal/rpc/recording_server.go b/internal/rpc/recording_server.go index a2da5b2..bdc0a52 100644 --- a/internal/rpc/recording_server.go +++ b/internal/rpc/recording_server.go @@ -23,7 +23,7 @@ import ( // when the active Agent session, mTLS peer or original execution is unknown. type RecordingServer struct { agentpb.UnimplementedAgentControlServiceServer - Store *store.CurrentStore + Store *store.Store OSS *oss.Client DispatcherID string TrustedFingerprints map[string]struct{} @@ -89,7 +89,7 @@ func (s *RecordingServer) RequestRecordingUpload(ctx context.Context, req *agent log.Printf("recording grant failed dispatcher=%s event=%s stage=signed_target_mismatch", req.GetDispatcherId(), req.GetSourceEventId()) return nil, status.Error(codes.Internal, "OSS grant differs from approved original asset") } - binding := store.CurrentRecordingGrant{ + binding := store.RecordingGrant{ DispatcherID: req.GetDispatcherId(), SourceEventID: req.GetSourceEventId(), UploadID: req.GetUploadId(), RecordingID: asset.GetAssetId(), Bucket: grant.GetBucket(), ObjectKey: grant.GetObjectKey(), ChecksumSHA256: asset.GetChecksumSha256(), SizeBytes: asset.GetSizeBytes(), Format: asset.GetFormat(), Channels: int(asset.GetChannels()), SampleRateHz: int(asset.GetSampleRateHz()), DurationMS: asset.GetDurationMs(), @@ -98,9 +98,9 @@ func (s *RecordingServer) RequestRecordingUpload(ctx context.Context, req *agent if err != nil { log.Printf("recording grant failed dispatcher=%s event=%s stage=bind error_class=%T", req.GetDispatcherId(), req.GetSourceEventId(), err) switch { - case errors.Is(err, store.ErrCurrentUploadConflict), errors.Is(err, store.ErrCurrentResultConflict): + case errors.Is(err, store.ErrUploadConflict), errors.Is(err, store.ErrResultConflict): return nil, status.Error(codes.AlreadyExists, "original recording target cannot be changed") - case errors.Is(err, store.ErrCurrentUploadAlreadyConfirmed): + case errors.Is(err, store.ErrUploadAlreadyConfirmed): return nil, status.Error(codes.FailedPrecondition, "recording is already confirmed; another PUT is forbidden") default: return nil, status.Error(codes.Internal, "original recording target could not be persisted") @@ -129,11 +129,11 @@ func (s *RecordingServer) ReportCallResult(ctx context.Context, req *agentpb.Rep if len(req.GetResultPayloadJson()) == 0 { return nil, status.Error(codes.InvalidArgument, "final result payload is required") } - var event store.CurrentOutboxEvent + var event store.OutboxEvent var created bool var err error if observation := req.GetUpload(); observation != nil { - event, created, err = s.Store.RecordUploadedCallResult(req.GetDispatcherId(), req.GetSourceEventId(), req.GetResultPayloadJson(), store.CurrentUploadProof{ + event, created, err = s.Store.RecordUploadedCallResult(req.GetDispatcherId(), req.GetSourceEventId(), req.GetResultPayloadJson(), store.UploadProof{ UploadID: observation.GetUploadId(), RecordingID: observation.GetRecordingId(), StatusCode: int(observation.GetPutStatusCode()), SizeBytes: observation.GetSizeBytes(), SHA256: observation.GetChecksumSha256(), }) @@ -143,11 +143,11 @@ func (s *RecordingServer) ReportCallResult(ctx context.Context, req *agentpb.Rep if err != nil { log.Printf("final result not committed dispatcher=%s event=%s stage=outbox error_class=%T", req.GetDispatcherId(), req.GetSourceEventId(), err) switch { - case errors.Is(err, store.ErrCurrentResultInvalid): + case errors.Is(err, store.ErrResultInvalid): return nil, status.Error(codes.InvalidArgument, "final result violates the approved MQ contract") - case errors.Is(err, store.ErrCurrentUploadUnverified), errors.Is(err, store.ErrCurrentEndUnconfirmed): + case errors.Is(err, store.ErrUploadUnverified), errors.Is(err, store.ErrEndUnconfirmed): return nil, status.Error(codes.FailedPrecondition, "final result requires confirmed call end and original upload outcome") - case errors.Is(err, store.ErrCurrentResultConflict): + case errors.Is(err, store.ErrResultConflict): return nil, status.Error(codes.AlreadyExists, "call already has a different final result") default: return nil, status.Error(codes.Internal, "final result could not be persisted") diff --git a/internal/rpc/recording_server_flow_test.go b/internal/rpc/recording_server_flow_test.go index 463c453..6ec4f58 100644 --- a/internal/rpc/recording_server_flow_test.go +++ b/internal/rpc/recording_server_flow_test.go @@ -26,9 +26,9 @@ import ( const recordingDispatcherID = "c046b893-8628-4589-ae50-619d049248a6" -func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.CurrentStore, context.Context, *agentpb.RequestRecordingUploadRequest, configread.CurrentSnapshot) { +func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.Store, context.Context, *agentpb.RequestRecordingUploadRequest, configread.CurrentSnapshot) { t.Helper() - database, err := store.OpenCurrent(filepath.Join(t.TempDir(), "recordings.db")) + database, err := store.Open(filepath.Join(t.TempDir(), "recordings.db")) if err != nil { t.Fatal(err) } @@ -65,12 +65,12 @@ func recordingRPCFixture(t *testing.T) (*RecordingServer, *store.CurrentStore, c if err := database.MarkReadyForSIP(recordingDispatcherID, 8); err != nil { t.Fatal(err) } - command := store.CurrentExecuteCommand{DispatcherID: recordingDispatcherID, EventID: "call-recording-1", TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:30:00Z"} + command := store.ExecuteCommand{DispatcherID: recordingDispatcherID, EventID: "call-recording-1", TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:30:00Z"} if _, _, err := database.RecordExecute(command); err != nil { t.Fatal(err) } now := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) - if err := database.ReserveExecute(recordingDispatcherID, command.EventID, store.CurrentCallReservation{ + if err := database.ReserveExecute(recordingDispatcherID, command.EventID, store.CallReservation{ TrunkID: "trunk-mock", SIPRevision: 8, CallerID: "BD00000000", DialedCallee: command.Callee, Deadline: now.Add(2 * time.Minute), }, now); err != nil { t.Fatal(err) @@ -135,7 +135,7 @@ func TestRecordingServerAuthenticatesOriginalGrantAndExplicitReissue(t *testing. if _, err := server.RequestRecordingUpload(context.Background(), request); status.Code(err) != codes.Unauthenticated { t.Fatalf("upload token issued without verified mTLS peer: %v", err) } - if _, err := database.LoadRecordingUpload(request.DispatcherId, request.SourceEventId); !errors.Is(err, store.ErrCurrentUploadNotFound) { + if _, err := database.LoadRecordingUpload(request.DispatcherId, request.SourceEventId); !errors.Is(err, store.ErrUploadNotFound) { t.Fatalf("unauthenticated Agent persisted a grant: %v", err) } wrongDispatcher := proto.Clone(request).(*agentpb.RequestRecordingUploadRequest) diff --git a/internal/store/current_calls.go b/internal/store/calls.go similarity index 84% rename from internal/store/current_calls.go rename to internal/store/calls.go index 091b2f0..cb356fb 100644 --- a/internal/store/current_calls.go +++ b/internal/store/calls.go @@ -13,13 +13,13 @@ import ( ) var ( - ErrCurrentCapacity = errors.New("current task, tenant, or trunk capacity exhausted") - ErrCurrentNotReady = errors.New("current task admission is not ready") - ErrCurrentAlreadyStarted = errors.New("call instruction already began dispatch") - ErrCurrentStopped = errors.New("stopped task does not accept old call instructions") + ErrCapacity = errors.New("current task, tenant, or trunk capacity exhausted") + ErrNotReady = errors.New("current task admission is not ready") + ErrAlreadyStarted = errors.New("call instruction already began dispatch") + ErrStopped = errors.New("stopped task does not accept old call instructions") ) -type CurrentExecuteCommand struct { +type ExecuteCommand struct { DispatcherID string EventID string TenantID int64 @@ -28,13 +28,13 @@ type CurrentExecuteCommand struct { IssuedAt string } -type CurrentExecution struct { - CurrentExecuteCommand +type Execution struct { + ExecuteCommand Status string SelectedTrunkID string } -type CurrentCallReservation struct { +type CallReservation struct { TrunkID string SIPRevision int64 CallerID string @@ -42,7 +42,7 @@ type CurrentCallReservation struct { Deadline time.Time } -type CurrentOutboxEvent struct { +type OutboxEvent struct { EventID string EventType string RoutingKey string @@ -51,50 +51,50 @@ type CurrentOutboxEvent struct { // RecordExecute durably accepts the transport message before MQ ACK. A // redelivery of the same message identity never creates another execution. -func (s *CurrentStore) RecordExecute(cmd CurrentExecuteCommand) (CurrentExecution, bool, error) { +func (s *Store) RecordExecute(cmd ExecuteCommand) (Execution, bool, error) { if cmd.DispatcherID == "" || cmd.EventID == "" || len(cmd.EventID) > 255 || cmd.TenantID <= 0 || cmd.TaskID == "" || cmd.Callee == "" { - return CurrentExecution{}, false, errors.New("invalid call command identity or callee") + return Execution{}, false, errors.New("invalid call command identity or callee") } if _, err := time.Parse(time.RFC3339Nano, cmd.IssuedAt); err != nil { - return CurrentExecution{}, false, fmt.Errorf("invalid call issued_at: %w", err) + return Execution{}, false, fmt.Errorf("invalid call issued_at: %w", err) } tx, err := s.db.Begin() if err != nil { - return CurrentExecution{}, false, err + return Execution{}, false, err } defer tx.Rollback() var controlState string var present int err = tx.QueryRow(`SELECT control_state,present FROM dispatcher_tasks WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, cmd.DispatcherID, cmd.TenantID, cmd.TaskID).Scan(&controlState, &present) if errors.Is(err, sql.ErrNoRows) || present != 1 { - return CurrentExecution{}, false, ErrCurrentNotReady + return Execution{}, false, ErrNotReady } if err != nil { - return CurrentExecution{}, false, fmt.Errorf("read call task ownership: %w", err) + return Execution{}, false, fmt.Errorf("read call task ownership: %w", err) } if controlState == "stopped" || controlState == "stopping" { - return CurrentExecution{}, false, ErrCurrentStopped + return Execution{}, false, ErrStopped } result, err := tx.Exec(`INSERT INTO dispatcher_inbox(dispatcher_id,event_id,tenant_id,task_id,callee,issued_at,status) VALUES(?,?,?,?,?,?,'pending') ON CONFLICT(dispatcher_id,event_id) DO NOTHING`, cmd.DispatcherID, cmd.EventID, cmd.TenantID, cmd.TaskID, cmd.Callee, cmd.IssuedAt) if err != nil { - return CurrentExecution{}, false, fmt.Errorf("persist call inbox: %w", err) + return Execution{}, false, fmt.Errorf("persist call inbox: %w", err) } affected, err := result.RowsAffected() if err != nil { - return CurrentExecution{}, false, err + return Execution{}, false, err } - var found CurrentExecution + var found Execution err = tx.QueryRow(`SELECT dispatcher_id,event_id,tenant_id,task_id,callee,issued_at,status,COALESCE(selected_trunk_id,'') FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, cmd.DispatcherID, cmd.EventID).Scan(&found.DispatcherID, &found.EventID, &found.TenantID, &found.TaskID, &found.Callee, &found.IssuedAt, &found.Status, &found.SelectedTrunkID) if err != nil { - return CurrentExecution{}, false, fmt.Errorf("read call inbox: %w", err) + return Execution{}, false, fmt.Errorf("read call inbox: %w", err) } - if found.CurrentExecuteCommand != cmd { - return CurrentExecution{}, false, fmt.Errorf("conflicting call message identity %q", cmd.EventID) + if found.ExecuteCommand != cmd { + return Execution{}, false, fmt.Errorf("conflicting call message identity %q", cmd.EventID) } if err := tx.Commit(); err != nil { - return CurrentExecution{}, false, fmt.Errorf("commit call inbox: %w", err) + return Execution{}, false, fmt.Errorf("commit call inbox: %w", err) } return found, affected == 1, nil } @@ -102,7 +102,7 @@ func (s *CurrentStore) RecordExecute(cmd CurrentExecuteCommand) (CurrentExecutio // ReserveExecute is the final transactional quota/fencing boundary before // sending a single originate. A crash or timeout after this point leaves an // occupied dispatching/unknown execution, never an automatic redial. -func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected CurrentCallReservation, at time.Time) error { +func (s *Store) ReserveExecute(dispatcherID, eventID string, selected CallReservation, at time.Time) error { if dispatcherID == "" || eventID == "" || selected.TrunkID == "" || selected.SIPRevision <= 0 || selected.CallerID == "" || selected.DialedCallee == "" || !selected.Deadline.After(at) { return errors.New("invalid originate reservation or expired deadline") } @@ -111,14 +111,14 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur return err } defer tx.Rollback() - var cmd CurrentExecuteCommand + var cmd ExecuteCommand var status string err = tx.QueryRow(`SELECT tenant_id,task_id,callee,issued_at,status FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, eventID).Scan(&cmd.TenantID, &cmd.TaskID, &cmd.Callee, &cmd.IssuedAt, &status) if err != nil { return fmt.Errorf("load durable call instruction: %w", err) } if status != "pending" { - return ErrCurrentAlreadyStarted + return ErrAlreadyStarted } cmd.DispatcherID, cmd.EventID = dispatcherID, eventID issued, err := time.Parse(time.RFC3339Nano, cmd.IssuedAt) @@ -126,7 +126,7 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur return fmt.Errorf("decode persisted call issued_at: %w", err) } if issued.After(at) { - return fmt.Errorf("%w: call issued_at is in the future", ErrCurrentNotReady) + return fmt.Errorf("%w: call issued_at is in the future", ErrNotReady) } var body []byte var revision int64 @@ -135,13 +135,13 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur JOIN dispatcher_configs c ON c.dispatcher_id=t.dispatcher_id AND c.tenant_id=t.tenant_id AND c.task_id=t.task_id AND c.task_revision=t.task_revision WHERE t.dispatcher_id=? AND t.tenant_id=? AND t.task_id=? AND t.present=1 AND t.status='running' AND t.control_state='' AND ds.discovery_ready=1`, dispatcherID, cmd.TenantID, cmd.TaskID).Scan(&body, &revision) if errors.Is(err, sql.ErrNoRows) { - return ErrCurrentNotReady + return ErrNotReady } if err != nil { return fmt.Errorf("load current task admission: %w", err) } if revision != selected.SIPRevision { - return fmt.Errorf("%w: SIP revision changed", ErrCurrentNotReady) + return fmt.Errorf("%w: SIP revision changed", ErrNotReady) } var snapshot struct { Task configread.CurrentTask `json:"task"` @@ -155,7 +155,7 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur return errors.New("admission snapshot identity mismatch") } if snapshot.Task.MaxConcurrentCalls <= 0 || snapshot.Quota.MaxConcurrentCalls <= 0 { - return fmt.Errorf("%w: missing quota", ErrCurrentNotReady) + return fmt.Errorf("%w: missing quota", ErrNotReady) } var trunks []struct { TrunkID string `json:"trunk_id"` @@ -173,7 +173,7 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur } } if trunkLimit <= 0 { - return fmt.Errorf("%w: selected trunk is disabled or has unknown quota", ErrCurrentNotReady) + return fmt.Errorf("%w: selected trunk is disabled or has unknown quota", ErrNotReady) } count := func(query string, args ...any) (int64, error) { var n int64 @@ -194,7 +194,7 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur return fmt.Errorf("count trunk occupancy: %w", err) } if tenantUsed >= snapshot.Quota.MaxConcurrentCalls || taskUsed >= snapshot.Task.MaxConcurrentCalls || trunkUsed >= trunkLimit { - return ErrCurrentCapacity + return ErrCapacity } result, err := tx.Exec(`UPDATE dispatcher_inbox SET status='dispatching',selected_trunk_id=?,caller_id=?,dialed_callee=?,deadline=?,snapshot_json=? WHERE dispatcher_id=? AND event_id=? AND status='pending'`, selected.TrunkID, selected.CallerID, selected.DialedCallee, selected.Deadline.UTC().Format(time.RFC3339Nano), body, dispatcherID, eventID) @@ -206,7 +206,7 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur return err } if n != 1 { - return ErrCurrentAlreadyStarted + return ErrAlreadyStarted } if err := tx.Commit(); err != nil { return fmt.Errorf("commit pre-originate fence: %w", err) @@ -217,7 +217,7 @@ func (s *CurrentStore) ReserveExecute(dispatcherID, eventID string, selected Cur // CancelReservationBeforeOrigin is valid only before invoking the Agent RPC. // It releases a reservation if the clock, SIP load, or schedule changed during // the final pre-dial check; it must never be used after an unknown RPC outcome. -func (s *CurrentStore) CancelReservationBeforeOrigin(dispatcherID, eventID string) error { +func (s *Store) CancelReservationBeforeOrigin(dispatcherID, eventID string) error { result, err := s.db.Exec(`UPDATE dispatcher_inbox SET status=CASE WHEN EXISTS( SELECT 1 FROM dispatcher_tasks t WHERE t.dispatcher_id=dispatcher_inbox.dispatcher_id AND t.tenant_id=dispatcher_inbox.tenant_id AND t.task_id=dispatcher_inbox.task_id @@ -230,7 +230,7 @@ func (s *CurrentStore) CancelReservationBeforeOrigin(dispatcherID, eventID strin return requireOneRow(result, "cancel unused originate reservation") } -func (s *CurrentStore) MarkExecuteUnknown(dispatcherID, eventID string) error { +func (s *Store) MarkExecuteUnknown(dispatcherID, eventID string) error { result, err := s.db.Exec(`UPDATE dispatcher_inbox SET status='unknown' WHERE dispatcher_id=? AND event_id=? AND status='dispatching'`, dispatcherID, eventID) if err != nil { return fmt.Errorf("persist unknown execution: %w", err) @@ -258,7 +258,7 @@ func (s *CurrentStore) MarkExecuteUnknown(dispatcherID, eventID string) error { // end can race the originate RPC acknowledgment, or resolve a previously // unknown RPC outcome. In both cases the acknowledgment and occupancy release // commit together; a missing outbox never frees capacity. -func (s *CurrentStore) FinishExecute(dispatcherID, eventID string) error { +func (s *Store) FinishExecute(dispatcherID, eventID string) error { tx, err := s.db.Begin() if err != nil { return err @@ -297,18 +297,18 @@ func (s *CurrentStore) FinishExecute(dispatcherID, eventID string) error { return nil } -func (s *CurrentStore) RejectExecute(dispatcherID, eventID, reason string) error { +func (s *Store) RejectExecute(dispatcherID, eventID, reason string) error { if reason == "" { return errors.New("rejection reason is required") } return s.closeExecuteWithAck(dispatcherID, eventID, "pending", "rejected", map[string]any{"status": "rejected", "reason_code": nil, "reason_message": reason}) } -func (s *CurrentStore) MarkExecuteDispatched(dispatcherID, eventID string) error { +func (s *Store) MarkExecuteDispatched(dispatcherID, eventID string) error { return s.closeExecuteWithAck(dispatcherID, eventID, "dispatching", "dispatched", map[string]any{"status": "dispatched"}) } -func (s *CurrentStore) closeExecuteWithAck(dispatcherID, eventID, expectedStatus, nextStatus string, payload any) error { +func (s *Store) closeExecuteWithAck(dispatcherID, eventID, expectedStatus, nextStatus string, payload any) error { tx, err := s.db.Begin() if err != nil { return err @@ -380,7 +380,7 @@ func requireExecuteAck(tx *sql.Tx, dispatcherID, eventID string) error { // PendingExecuteCount is scoped to one task so a rule wait pauses only its // SaaS-owned task queue; other tasks and the control queue keep consuming. -func (s *CurrentStore) PendingExecuteCount(dispatcherID string, tenantID int64, taskID string) (int64, error) { +func (s *Store) PendingExecuteCount(dispatcherID string, tenantID int64, taskID string) (int64, error) { if dispatcherID == "" || tenantID <= 0 || taskID == "" { return 0, errors.New("pending task count requires durable task identity") } @@ -393,7 +393,7 @@ func (s *CurrentStore) PendingExecuteCount(dispatcherID string, tenantID int64, // TrunkOccupancy includes dispatching and unknown executions. No timeout or // lease expiry releases a call whose real end has not been confirmed. -func (s *CurrentStore) TrunkOccupancy(dispatcherID string) (map[string]int64, error) { +func (s *Store) TrunkOccupancy(dispatcherID string) (map[string]int64, error) { rows, err := s.db.Query(`SELECT selected_trunk_id, COUNT(*) FROM dispatcher_inbox WHERE dispatcher_id=? AND status IN ('dispatching','dispatched','unknown') AND selected_trunk_id IS NOT NULL GROUP BY selected_trunk_id`, dispatcherID) @@ -416,15 +416,15 @@ func (s *CurrentStore) TrunkOccupancy(dispatcherID string) (map[string]int64, er return occupied, nil } -func (s *CurrentStore) ListPendingExecute(dispatcherID string) ([]CurrentExecuteCommand, error) { +func (s *Store) ListPendingExecute(dispatcherID string) ([]ExecuteCommand, error) { rows, err := s.db.Query(`SELECT dispatcher_id,event_id,tenant_id,task_id,callee,issued_at FROM dispatcher_inbox WHERE dispatcher_id=? AND status='pending' ORDER BY rowid`, dispatcherID) if err != nil { return nil, fmt.Errorf("list pending call commands: %w", err) } defer rows.Close() - var commands []CurrentExecuteCommand + var commands []ExecuteCommand for rows.Next() { - var cmd CurrentExecuteCommand + var cmd ExecuteCommand if err := rows.Scan(&cmd.DispatcherID, &cmd.EventID, &cmd.TenantID, &cmd.TaskID, &cmd.Callee, &cmd.IssuedAt); err != nil { return nil, err } @@ -436,15 +436,15 @@ func (s *CurrentStore) ListPendingExecute(dispatcherID string) ([]CurrentExecute return commands, nil } -func (s *CurrentStore) ListPendingOutbox(dispatcherID string) ([]CurrentOutboxEvent, error) { +func (s *Store) ListPendingOutbox(dispatcherID string) ([]OutboxEvent, error) { rows, err := s.db.Query(`SELECT event_id,event_type,routing_key,body FROM dispatcher_outbox WHERE dispatcher_id=? AND confirmed=0 ORDER BY rowid`, dispatcherID) if err != nil { return nil, fmt.Errorf("list pending SaaS events: %w", err) } defer rows.Close() - var events []CurrentOutboxEvent + var events []OutboxEvent for rows.Next() { - var event CurrentOutboxEvent + var event OutboxEvent if err := rows.Scan(&event.EventID, &event.EventType, &event.RoutingKey, &event.Body); err != nil { return nil, err } @@ -458,7 +458,7 @@ func (s *CurrentStore) ListPendingOutbox(dispatcherID string) ([]CurrentOutboxEv // Confirm means delivery to a bound queue, not application receipt. The row // remains durable and is never deleted without separate explicit evidence. -func (s *CurrentStore) MarkOutboxConfirmed(dispatcherID, eventID string) error { +func (s *Store) MarkOutboxConfirmed(dispatcherID, eventID string) error { result, err := s.db.Exec(`UPDATE dispatcher_outbox SET confirmed=1,confirmed_at=? WHERE dispatcher_id=? AND event_id=? AND confirmed=0`, time.Now().UTC().Format(time.RFC3339Nano), dispatcherID, eventID) if err != nil { return fmt.Errorf("persist MQ queue confirm: %w", err) diff --git a/internal/store/current_calls_test.go b/internal/store/calls_test.go similarity index 88% rename from internal/store/current_calls_test.go rename to internal/store/calls_test.go index d93effe..1242850 100644 --- a/internal/store/current_calls_test.go +++ b/internal/store/calls_test.go @@ -10,9 +10,9 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -func preparedCurrentCallStore(t *testing.T) *CurrentStore { +func preparedCurrentCallStore(t *testing.T) *Store { t.Helper() - s, err := OpenCurrent(filepath.Join(t.TempDir(), "state.db")) + s, err := Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } @@ -31,14 +31,14 @@ func preparedCurrentCallStore(t *testing.T) *CurrentStore { return s } -func currentCall(id string) CurrentExecuteCommand { - return CurrentExecuteCommand{DispatcherID: currentDispatcherID, EventID: id, TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:30:00Z"} +func currentCall(id string) ExecuteCommand { + return ExecuteCommand{DispatcherID: currentDispatcherID, EventID: id, TenantID: 1001, TaskID: "task-asr", Callee: "15003164745", IssuedAt: "2026-09-21T01:30:00Z"} } -func currentReservation() CurrentCallReservation { - return CurrentCallReservation{TrunkID: "trunk-mock", SIPRevision: 8, CallerID: "BD00000000", DialedCallee: "15003164745", Deadline: time.Date(2026, 9, 21, 1, 32, 0, 0, time.UTC)} +func currentReservation() CallReservation { + return CallReservation{TrunkID: "trunk-mock", SIPRevision: 8, CallerID: "BD00000000", DialedCallee: "15003164745", Deadline: time.Date(2026, 9, 21, 1, 32, 0, 0, time.UTC)} } -func TestCurrentExecuteInboxPreventsDuplicateOriginationAfterRestart(t *testing.T) { +func TestExecuteInboxPreventsDuplicateOriginationAfterRestart(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("call-1") first, created, err := s.RecordExecute(cmd) @@ -70,7 +70,7 @@ func TestCurrentExecuteInboxPreventsDuplicateOriginationAfterRestart(t *testing. } } -func TestCurrentExecuteRejectsOnlyOnceWithoutFinalResult(t *testing.T) { +func TestExecuteRejectsOnlyOnceWithoutFinalResult(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("call-invalid") if _, _, err := s.RecordExecute(cmd); err != nil { @@ -101,7 +101,7 @@ func TestCurrentExecuteRejectsOnlyOnceWithoutFinalResult(t *testing.T) { } } -func TestCurrentExecuteQuotaCountsUnknownAndReleasesOnlyConfirmedEnd(t *testing.T) { +func TestExecuteQuotaCountsUnknownAndReleasesOnlyConfirmedEnd(t *testing.T) { s := preparedCurrentCallStore(t) at := time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC) for _, id := range []string{"call-1", "call-2"} { @@ -117,13 +117,13 @@ func TestCurrentExecuteQuotaCountsUnknownAndReleasesOnlyConfirmedEnd(t *testing. if _, _, err := s.RecordExecute(third); err != nil { t.Fatal(err) } - if err := s.ReserveExecute(third.DispatcherID, third.EventID, currentReservation(), at); !errors.Is(err, ErrCurrentCapacity) { + if err := s.ReserveExecute(third.DispatcherID, third.EventID, currentReservation(), at); !errors.Is(err, ErrCapacity) { t.Fatalf("third call exceeded task/trunk capacity: %v", err) } if err := s.MarkExecuteUnknown(currentDispatcherID, "call-1"); err != nil { t.Fatal(err) } - if err := s.ReserveExecute(third.DispatcherID, third.EventID, currentReservation(), at); !errors.Is(err, ErrCurrentCapacity) { + if err := s.ReserveExecute(third.DispatcherID, third.EventID, currentReservation(), at); !errors.Is(err, ErrCapacity) { t.Fatalf("unknown call released occupancy: %v", err) } if err := s.MarkExecuteDispatched(currentDispatcherID, "call-2"); err != nil { @@ -137,7 +137,7 @@ func TestCurrentExecuteQuotaCountsUnknownAndReleasesOnlyConfirmedEnd(t *testing. } } -func TestCurrentConfirmedEndBeforeOriginateAckIsDurableAndIdempotent(t *testing.T) { +func TestConfirmedEndBeforeOriginateAckIsDurableAndIdempotent(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("call-ended-before-ack") if _, _, err := s.RecordExecute(cmd); err != nil { @@ -171,7 +171,7 @@ func TestCurrentConfirmedEndBeforeOriginateAckIsDurableAndIdempotent(t *testing. } } -func TestCurrentConfirmedUnknownCallEndReleasesOnlyAfterDurableAck(t *testing.T) { +func TestConfirmedUnknownCallEndReleasesOnlyAfterDurableAck(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("call-ended-after-unknown") if _, _, err := s.RecordExecute(cmd); err != nil { @@ -196,7 +196,7 @@ func TestCurrentConfirmedUnknownCallEndReleasesOnlyAfterDurableAck(t *testing.T) } } -func TestCurrentEarlyEndOutboxFailureRetainsUnknownOccupancy(t *testing.T) { +func TestEarlyEndOutboxFailureRetainsUnknownOccupancy(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("call-early-end-fault") if _, _, err := s.RecordExecute(cmd); err != nil { @@ -223,7 +223,7 @@ func TestCurrentEarlyEndOutboxFailureRetainsUnknownOccupancy(t *testing.T) { } } -func TestCurrentEndAndOriginateAckRaceNeverDuplicate(t *testing.T) { +func TestEndAndOriginateAckRaceNeverDuplicate(t *testing.T) { for attempt := 0; attempt < 30; attempt++ { s := preparedCurrentCallStore(t) cmd := currentCall("call-concurrent-end") @@ -252,7 +252,7 @@ func TestCurrentEndAndOriginateAckRaceNeverDuplicate(t *testing.T) { } } -func TestCurrentExecuteOutboxFailureDoesNotMarkOriginatedCallDelivered(t *testing.T) { +func TestExecuteOutboxFailureDoesNotMarkOriginatedCallDelivered(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("call-fault") if _, _, err := s.RecordExecute(cmd); err != nil { diff --git a/internal/store/current_canonical_test.go b/internal/store/canonical_test.go similarity index 82% rename from internal/store/current_canonical_test.go rename to internal/store/canonical_test.go index 2c91b18..4409746 100644 --- a/internal/store/current_canonical_test.go +++ b/internal/store/canonical_test.go @@ -6,8 +6,8 @@ import ( "testing" ) -func TestCurrentStoreSameRevisionCanonicalJSON(t *testing.T) { - s, err := OpenCurrent(filepath.Join(t.TempDir(), "state.db")) +func TestStoreSameRevisionCanonicalJSON(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } diff --git a/internal/store/current_control_ack_test.go b/internal/store/control_ack_test.go similarity index 94% rename from internal/store/current_control_ack_test.go rename to internal/store/control_ack_test.go index 6d2bdb0..0df7ad7 100644 --- a/internal/store/current_control_ack_test.go +++ b/internal/store/control_ack_test.go @@ -5,7 +5,7 @@ import ( "testing" ) -func TestCurrentControlAckOnlyAfterDispatchAndInSameTransaction(t *testing.T) { +func TestControlAckOnlyAfterDispatchAndInSameTransaction(t *testing.T) { s := preparedCurrentCallStore(t) if err := s.PrepareControl(currentDispatcherID, 1001, "task-asr", "pause"); err != nil { t.Fatal(err) @@ -39,7 +39,7 @@ func TestCurrentControlAckOnlyAfterDispatchAndInSameTransaction(t *testing.T) { } } -func TestCurrentResumeDoesNotOpenAdmissionWhenAckOutboxFails(t *testing.T) { +func TestResumeDoesNotOpenAdmissionWhenAckOutboxFails(t *testing.T) { s := preparedCurrentCallStore(t) if err := s.PrepareControl(currentDispatcherID, 1001, "task-asr", "pause"); err != nil { t.Fatal(err) @@ -70,7 +70,7 @@ func TestCurrentResumeDoesNotOpenAdmissionWhenAckOutboxFails(t *testing.T) { } } -func TestCurrentStopRejectsResumeAndSuppressesOldCommands(t *testing.T) { +func TestStopRejectsResumeAndSuppressesOldCommands(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("queued-stop") if _, _, err := s.RecordExecute(cmd); err != nil { diff --git a/internal/store/current_control_flow.go b/internal/store/control_flow.go similarity index 90% rename from internal/store/current_control_flow.go rename to internal/store/control_flow.go index e32ebfc..4c14b85 100644 --- a/internal/store/current_control_flow.go +++ b/internal/store/control_flow.go @@ -10,13 +10,13 @@ import ( "git.ipao.vip/rogee/go-sip/internal/tenant" ) -var ErrCurrentControlRejected = errors.New("task control rejected by durable task state") +var ErrControlRejected = errors.New("task control rejected by durable task state") // PrepareControl closes only this task's admission before requesting the // Agent action. No successful control acknowledgment exists at this point. // Repeated commands are prepared and dispatched again; this is not control // deduplication or an expected-revision/CAS API. -func (s *CurrentStore) PrepareControl(dispatcherID string, tenantID int64, taskID, action string) error { +func (s *Store) PrepareControl(dispatcherID string, tenantID int64, taskID, action string) error { if dispatcherID == "" || tenantID <= 0 || taskID == "" { return errors.New("invalid task control identity") } @@ -34,23 +34,23 @@ func (s *CurrentStore) PrepareControl(dispatcherID string, tenantID int64, taskI return fmt.Errorf("load task control state: %w", err) } if present != 1 { - return fmt.Errorf("%w: task is not in the assigned discovery list", ErrCurrentControlRejected) + return fmt.Errorf("%w: task is not in the assigned discovery list", ErrControlRejected) } var prepared string switch action { case "pause": if state == "stopped" || state == "stopping" { - return fmt.Errorf("%w: stopped task cannot be paused", ErrCurrentControlRejected) + return fmt.Errorf("%w: stopped task cannot be paused", ErrControlRejected) } prepared = "pausing" case "stop": prepared = "stopping" case "resume": if state != "paused" && state != "resuming" { - return fmt.Errorf("%w: resume requires a persistently paused, non-stopped task", ErrCurrentControlRejected) + return fmt.Errorf("%w: resume requires a persistently paused, non-stopped task", ErrControlRejected) } if status != "running" { - return fmt.Errorf("%w: fresh task is not running", ErrCurrentControlRejected) + return fmt.Errorf("%w: fresh task is not running", ErrControlRejected) } var configuredRevision int64 err = tx.QueryRow(`SELECT task_revision FROM dispatcher_configs WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, dispatcherID, tenantID, taskID).Scan(&configuredRevision) @@ -78,7 +78,7 @@ func (s *CurrentStore) PrepareControl(dispatcherID string, tenantID int64, taskI // CompleteControl atomically records the applied state and durable outbox // only after an Agent control RPC accepted the action. It does not claim that // active-call drain or hangup has already finished. -func (s *CurrentStore) CompleteControl(dispatcherID string, tenantID int64, taskID, action, eventID string) error { +func (s *Store) CompleteControl(dispatcherID string, tenantID int64, taskID, action, eventID string) error { var expected, final string switch action { case "pause": @@ -115,7 +115,7 @@ func (s *CurrentStore) CompleteControl(dispatcherID string, tenantID int64, task return nil } -func (s *CurrentStore) RejectControl(dispatcherID string, tenantID int64, eventID string) error { +func (s *Store) RejectControl(dispatcherID string, tenantID int64, eventID string) error { tx, err := s.db.Begin() if err != nil { return err diff --git a/internal/store/current_control_test.go b/internal/store/control_test.go similarity index 95% rename from internal/store/current_control_test.go rename to internal/store/control_test.go index 9803886..49b501f 100644 --- a/internal/store/current_control_test.go +++ b/internal/store/control_test.go @@ -5,7 +5,7 @@ import ( "time" ) -func TestCurrentStopSuppressesUnstartedCallsButPreservesUnknown(t *testing.T) { +func TestStopSuppressesUnstartedCallsButPreservesUnknown(t *testing.T) { s := preparedCurrentCallStore(t) pending := currentCall("queued-before-stop") if _, _, err := s.RecordExecute(pending); err != nil { diff --git a/internal/store/current_discovery_test.go b/internal/store/discovery_test.go similarity index 92% rename from internal/store/current_discovery_test.go rename to internal/store/discovery_test.go index 3f1273e..8833bd0 100644 --- a/internal/store/current_discovery_test.go +++ b/internal/store/discovery_test.go @@ -11,8 +11,8 @@ func currentDiscoveredTask(status string, revision int64) configread.CurrentDisc return configread.CurrentDiscoveredTask{TaskID: "task-asr", TenantID: 1001, TaskRevision: revision, Status: status} } -func TestCurrentHTTPStopCannotBeUndoneByStaleRunningAndSuppressesInbox(t *testing.T) { - s, err := OpenCurrent(filepath.Join(t.TempDir(), "state.db")) +func TestHTTPStopCannotBeUndoneByStaleRunningAndSuppressesInbox(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } @@ -50,8 +50,8 @@ func TestCurrentHTTPStopCannotBeUndoneByStaleRunningAndSuppressesInbox(t *testin } } -func TestCurrentHTTPPausedRequiresExplicitMQResume(t *testing.T) { - s, err := OpenCurrent(filepath.Join(t.TempDir(), "state.db")) +func TestHTTPPausedRequiresExplicitMQResume(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } @@ -86,8 +86,8 @@ func TestCurrentHTTPPausedRequiresExplicitMQResume(t *testing.T) { } } -func TestCurrentDiscoveryPageCommitsAtomicallyAndPreservesControl(t *testing.T) { - s, err := OpenCurrent(filepath.Join(t.TempDir(), "state.db")) +func TestDiscoveryPageCommitsAtomicallyAndPreservesControl(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } diff --git a/internal/store/current_global_revision_test.go b/internal/store/global_revision_test.go similarity index 90% rename from internal/store/current_global_revision_test.go rename to internal/store/global_revision_test.go index e7ef1a1..2d08a84 100644 --- a/internal/store/current_global_revision_test.go +++ b/internal/store/global_revision_test.go @@ -10,8 +10,8 @@ import ( "git.ipao.vip/rogee/go-sip/internal/configread" ) -func TestCurrentGlobalSIPAndQuotaRevisionCannotChangeContent(t *testing.T) { - s, err := OpenCurrent(filepath.Join(t.TempDir(), "state.db")) +func TestGlobalSIPAndQuotaRevisionCannotChangeContent(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } @@ -60,8 +60,8 @@ func TestCurrentGlobalSIPAndQuotaRevisionCannotChangeContent(t *testing.T) { } } -func TestCurrentSharedGlobalRevisionConflictAcrossTasks(t *testing.T) { - s, err := OpenCurrent(filepath.Join(t.TempDir(), "state.db")) +func TestSharedGlobalRevisionConflictAcrossTasks(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "state.db")) if err != nil { t.Fatal(err) } diff --git a/internal/store/current_identity_test.go b/internal/store/identity_test.go similarity index 91% rename from internal/store/current_identity_test.go rename to internal/store/identity_test.go index 65744e1..55e3e9c 100644 --- a/internal/store/current_identity_test.go +++ b/internal/store/identity_test.go @@ -6,12 +6,12 @@ import ( "testing" ) -func TestCurrentStoreRejectsForeignDispatcherWithoutTouchingOutbox(t *testing.T) { +func TestStoreRejectsForeignDispatcherWithoutTouchingOutbox(t *testing.T) { const foreignID = "862c8e9b-26b8-4642-9fac-c7d85f487a30" const eventID = "retained-outbox-event" const originalBody = `{"proof":"original"}` path := filepath.Join(t.TempDir(), "dispatcher.sqlite") - owner, err := OpenCurrent(path) + owner, err := Open(path) if err != nil { t.Fatal(err) } @@ -26,7 +26,7 @@ func TestCurrentStoreRejectsForeignDispatcherWithoutTouchingOutbox(t *testing.T) t.Fatal(err) } - foreign, err := OpenCurrent(path) + foreign, err := Open(path) if err != nil { t.Fatal(err) } diff --git a/internal/store/current_pending_test.go b/internal/store/pending_test.go similarity index 92% rename from internal/store/current_pending_test.go rename to internal/store/pending_test.go index f67c7ae..df991e3 100644 --- a/internal/store/current_pending_test.go +++ b/internal/store/pending_test.go @@ -5,7 +5,7 @@ import ( "time" ) -func TestCurrentPendingCountPausesOnlyAffectedTaskQueue(t *testing.T) { +func TestPendingCountPausesOnlyAffectedTaskQueue(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("waiting-for-task-window") if _, _, err := s.RecordExecute(cmd); err != nil { diff --git a/internal/store/current_result.go b/internal/store/result.go similarity index 65% rename from internal/store/current_result.go rename to internal/store/result.go index 76c12c3..1b41464 100644 --- a/internal/store/current_result.go +++ b/internal/store/result.go @@ -14,10 +14,10 @@ import ( "git.ipao.vip/rogee/go-sip/internal/tenant" ) -var ErrCurrentUploadUnverified = errors.New("recording outcome is not verified against its Dispatcher-approved upload") -var ErrCurrentResultConflict = errors.New("call already has a different final result") -var ErrCurrentEndUnconfirmed = errors.New("final result requires confirmed call end and its frozen snapshot") -var ErrCurrentResultInvalid = errors.New("final result input violates the approved contract") +var ErrUploadUnverified = errors.New("recording outcome is not verified against its Dispatcher-approved upload") +var ErrResultConflict = errors.New("call already has a different final result") +var ErrEndUnconfirmed = errors.New("final result requires confirmed call end and its frozen snapshot") +var ErrResultInvalid = errors.New("final result input violates the approved contract") func currentResultEventID(dispatcherID, sourceEventID string) string { identity := sha256.Sum256([]byte("call.execute.result\x00" + dispatcherID + "\x00" + sourceEventID)) @@ -26,20 +26,20 @@ func currentResultEventID(dispatcherID, sourceEventID string) string { // RecordCallResult stores one no-recording final result only after confirmed // call end. A durable upload grant cannot be bypassed with an empty recording. -func (s *CurrentStore) RecordCallResult(dispatcherID, sourceEventID string, payload []byte) (CurrentOutboxEvent, bool, error) { +func (s *Store) RecordCallResult(dispatcherID, sourceEventID string, payload []byte) (OutboxEvent, bool, error) { return s.recordCallResult(dispatcherID, sourceEventID, payload, nil) } // RecordUploadedCallResult accepts the authenticated Agent's successful PUT // observation only for the precise, previously bound OSS asset. Both the // upload fact and sole SaaS result outbox entry commit in the same transaction. -func (s *CurrentStore) RecordUploadedCallResult(dispatcherID, sourceEventID string, payload []byte, proof CurrentUploadProof) (CurrentOutboxEvent, bool, error) { +func (s *Store) RecordUploadedCallResult(dispatcherID, sourceEventID string, payload []byte, proof UploadProof) (OutboxEvent, bool, error) { return s.recordCallResult(dispatcherID, sourceEventID, payload, &proof) } -func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payload []byte, proof *CurrentUploadProof) (CurrentOutboxEvent, bool, error) { +func (s *Store) recordCallResult(dispatcherID, sourceEventID string, payload []byte, proof *UploadProof) (OutboxEvent, bool, error) { if dispatcherID == "" || sourceEventID == "" || len(sourceEventID) > 255 || len(payload) == 0 { - return CurrentOutboxEvent{}, false, fmt.Errorf("%w: durable call identity and payload are required", ErrCurrentResultInvalid) + return OutboxEvent{}, false, fmt.Errorf("%w: durable call identity and payload are required", ErrResultInvalid) } var result struct { TaskID string `json:"task_id"` @@ -51,19 +51,19 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl Recording json.RawMessage `json:"recording"` } if err := json.Unmarshal(payload, &result); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("%w: decode JSON: %v", ErrCurrentResultInvalid, err) + return OutboxEvent{}, false, fmt.Errorf("%w: decode JSON: %v", ErrResultInvalid, err) } start, err := time.Parse(time.RFC3339Nano, result.StartedAt) if err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("%w: invalid start time", ErrCurrentResultInvalid) + return OutboxEvent{}, false, fmt.Errorf("%w: invalid start time", ErrResultInvalid) } end, err := time.Parse(time.RFC3339Nano, result.EndedAt) if err != nil || end.Before(start) { - return CurrentOutboxEvent{}, false, fmt.Errorf("%w: end precedes start or is invalid", ErrCurrentResultInvalid) + return OutboxEvent{}, false, fmt.Errorf("%w: end precedes start or is invalid", ErrResultInvalid) } tx, err := s.db.Begin() if err != nil { - return CurrentOutboxEvent{}, false, err + return OutboxEvent{}, false, err } defer tx.Rollback() var tenantID int64 @@ -71,22 +71,22 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl var snapshotJSON []byte err = tx.QueryRow(`SELECT tenant_id,task_id,callee,COALESCE(selected_trunk_id,''),status,snapshot_json FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, sourceEventID).Scan(&tenantID, &taskID, &callee, &trunkID, &status, &snapshotJSON) if errors.Is(err, sql.ErrNoRows) { - return CurrentOutboxEvent{}, false, errors.New("no durable approved call matches final result") + return OutboxEvent{}, false, errors.New("no durable approved call matches final result") } if err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("load completed call identity: %w", err) + return OutboxEvent{}, false, fmt.Errorf("load completed call identity: %w", err) } if status != "finished" || trunkID == "" || len(snapshotJSON) == 0 { - return CurrentOutboxEvent{}, false, ErrCurrentEndUnconfirmed + return OutboxEvent{}, false, ErrEndUnconfirmed } var snapshot struct { Task configread.CurrentTask `json:"task"` } if err := json.Unmarshal(snapshotJSON, &snapshot); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("decode frozen call task: %w", err) + return OutboxEvent{}, false, fmt.Errorf("decode frozen call task: %w", err) } if snapshot.Task.DispatcherID != dispatcherID || snapshot.Task.TenantID != tenantID || snapshot.Task.TaskID != taskID || result.TaskID != taskID || result.Callee != callee || result.TrunkID != trunkID || result.CallerProfileID != snapshot.Task.CallerProfileID { - return CurrentOutboxEvent{}, false, errors.New("final result identity differs from the approved call") + return OutboxEvent{}, false, errors.New("final result identity differs from the approved call") } var recording struct { Status string `json:"status"` @@ -100,24 +100,24 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl ChecksumSHA256 string `json:"checksum_sha256"` } if err := json.Unmarshal(result.Recording, &recording); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("%w: recording fact is invalid", ErrCurrentResultInvalid) + return OutboxEvent{}, false, fmt.Errorf("%w: recording fact is invalid", ErrResultInvalid) } grant, grantErr := loadRecordingUpload(tx.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, dispatcherID, sourceEventID)) - if grantErr != nil && !errors.Is(grantErr, ErrCurrentUploadNotFound) { - return CurrentOutboxEvent{}, false, fmt.Errorf("inspect original recording target: %w", grantErr) + if grantErr != nil && !errors.Is(grantErr, ErrUploadNotFound) { + return OutboxEvent{}, false, fmt.Errorf("inspect original recording target: %w", grantErr) } if recording.Status == "uploaded" { if proof == nil || grantErr != nil || proof.StatusCode < 200 || proof.StatusCode >= 300 || proof.UploadID != grant.UploadID || proof.RecordingID != grant.RecordingID || proof.SizeBytes != grant.SizeBytes || proof.SHA256 != grant.ChecksumSHA256 || recording.Bucket != grant.Bucket || recording.ObjectKey != grant.ObjectKey || recording.Format != grant.Format || recording.Channels != grant.Channels || recording.SampleRateHz != grant.SampleRateHz || recording.DurationMS != grant.DurationMS || recording.SizeBytes != grant.SizeBytes || recording.ChecksumSHA256 != grant.ChecksumSHA256 { - return CurrentOutboxEvent{}, false, ErrCurrentUploadUnverified + return OutboxEvent{}, false, ErrUploadUnverified } } else if proof != nil || grantErr == nil { - return CurrentOutboxEvent{}, false, ErrCurrentUploadUnverified + return OutboxEvent{}, false, ErrUploadUnverified } route, err := tenant.CurrentResultRoute(dispatcherID) if err != nil { - return CurrentOutboxEvent{}, false, err + return OutboxEvent{}, false, err } eventID := currentResultEventID(dispatcherID, sourceEventID) body, err := json.Marshal(struct { @@ -129,45 +129,45 @@ func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payl Payload json.RawMessage `json:"payload"` }{eventID, "call.execute.result", dispatcherID, tenantID, result.EndedAt, payload}) if err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("encode final result: %w", err) + return OutboxEvent{}, false, fmt.Errorf("encode final result: %w", err) } if err := contract.ValidateCurrent("mq", body); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("%w: MQ contract: %v", ErrCurrentResultInvalid, err) + return OutboxEvent{}, false, fmt.Errorf("%w: MQ contract: %v", ErrResultInvalid, err) } inserted, err := tx.Exec(`INSERT INTO dispatcher_outbox(dispatcher_id,event_id,event_type,routing_key,body) VALUES(?,?,?,?,?) ON CONFLICT(dispatcher_id,event_id) DO NOTHING`, dispatcherID, eventID, "call.execute.result", route.BindingKey, body) if err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("persist final result outbox: %w", err) + return OutboxEvent{}, false, fmt.Errorf("persist final result outbox: %w", err) } count, err := inserted.RowsAffected() if err != nil { - return CurrentOutboxEvent{}, false, err + return OutboxEvent{}, false, err } - var stored CurrentOutboxEvent + var stored OutboxEvent stored.EventID = eventID if err := tx.QueryRow(`SELECT event_type,routing_key,body FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, eventID).Scan(&stored.EventType, &stored.RoutingKey, &stored.Body); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("inspect existing final result: %w", err) + return OutboxEvent{}, false, fmt.Errorf("inspect existing final result: %w", err) } if stored.EventType != "call.execute.result" || stored.RoutingKey != route.BindingKey || !bytes.Equal(stored.Body, body) { - return CurrentOutboxEvent{}, false, ErrCurrentResultConflict + return OutboxEvent{}, false, ErrResultConflict } if proof != nil { if grant.ConfirmedAt != "" && count == 1 { - return CurrentOutboxEvent{}, false, ErrCurrentResultConflict + return OutboxEvent{}, false, ErrResultConflict } confirmed, err := tx.Exec(`UPDATE dispatcher_recordings SET confirmed_at=COALESCE(confirmed_at,?) WHERE dispatcher_id=? AND source_event_id=?`, time.Now().UTC().Format(time.RFC3339Nano), dispatcherID, sourceEventID) if err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("persist confirmed recording with result: %w", err) + return OutboxEvent{}, false, fmt.Errorf("persist confirmed recording with result: %w", err) } affected, err := confirmed.RowsAffected() if err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("inspect recording confirmation: %w", err) + return OutboxEvent{}, false, fmt.Errorf("inspect recording confirmation: %w", err) } if affected != 1 { - return CurrentOutboxEvent{}, false, errors.New("recording confirmation lost its original target") + return OutboxEvent{}, false, errors.New("recording confirmation lost its original target") } } if err := tx.Commit(); err != nil { - return CurrentOutboxEvent{}, false, fmt.Errorf("commit final result outbox: %w", err) + return OutboxEvent{}, false, fmt.Errorf("commit final result outbox: %w", err) } return stored, count == 1, nil } diff --git a/internal/store/current_result_test.go b/internal/store/result_test.go similarity index 93% rename from internal/store/current_result_test.go rename to internal/store/result_test.go index 3f10765..aea1d5c 100644 --- a/internal/store/current_result_test.go +++ b/internal/store/result_test.go @@ -28,7 +28,7 @@ func currentResultPayload(t *testing.T) []byte { return raw } -func currentResultCall(t *testing.T, s *CurrentStore, id string, finish bool) CurrentExecuteCommand { +func currentResultCall(t *testing.T, s *Store, id string, finish bool) ExecuteCommand { t.Helper() cmd := currentCall(id) if _, _, err := s.RecordExecute(cmd); err != nil { @@ -48,7 +48,7 @@ func currentResultCall(t *testing.T, s *CurrentStore, id string, finish bool) Cu return cmd } -func TestCurrentFinalResultNeedsConfirmedEndAndRemainsExactlyOne(t *testing.T) { +func TestFinalResultNeedsConfirmedEndAndRemainsExactlyOne(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "result-call-1", false) payload := currentResultPayload(t) @@ -103,7 +103,7 @@ func TestCurrentFinalResultNeedsConfirmedEndAndRemainsExactlyOne(t *testing.T) { } } -func TestCurrentFinalResultAfterEndRacedOriginateAck(t *testing.T) { +func TestFinalResultAfterEndRacedOriginateAck(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("result-ended-before-ack") if _, _, err := s.RecordExecute(cmd); err != nil { @@ -128,7 +128,7 @@ func TestCurrentFinalResultAfterEndRacedOriginateAck(t *testing.T) { } } -func TestCurrentFinalResultRejectsDifferentBoundTaskCallOrTrunk(t *testing.T) { +func TestFinalResultRejectsDifferentBoundTaskCallOrTrunk(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "result-bound-1", true) for _, change := range []struct{ field, value string }{ @@ -149,7 +149,7 @@ func TestCurrentFinalResultRejectsDifferentBoundTaskCallOrTrunk(t *testing.T) { } } -func TestCurrentFinalResultRejectsUnverifiedOSSAssetAndInventedState(t *testing.T) { +func TestFinalResultRejectsUnverifiedOSSAssetAndInventedState(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "result-unverified-1", true) var claim map[string]any @@ -158,7 +158,7 @@ func TestCurrentFinalResultRejectsUnverifiedOSSAssetAndInventedState(t *testing. } claim["recording"] = map[string]any{"status": "uploaded", "bucket": "mock-bucket", "object_key": "tenant/rec.wav", "format": "wav", "channels": 1, "sample_rate_hz": 16000, "duration_ms": 3000, "size_bytes": 100, "checksum_sha256": strings.Repeat("a", 64)} uploaded, _ := json.Marshal(claim) - if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, uploaded); !errors.Is(err, ErrCurrentUploadUnverified) { + if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, uploaded); !errors.Is(err, ErrUploadUnverified) { t.Fatalf("Agent invented an ungranted uploaded asset: %v", err) } claim["recording"] = map[string]any{"status": "unavailable"} @@ -171,7 +171,7 @@ func TestCurrentFinalResultRejectsUnverifiedOSSAssetAndInventedState(t *testing. } } -func TestCurrentFinalResultOutboxFailureRollsBackAndCanResume(t *testing.T) { +func TestFinalResultOutboxFailureRollsBackAndCanResume(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "result-fault-1", true) if _, err := s.db.Exec(`CREATE TRIGGER fail_final BEFORE INSERT ON dispatcher_outbox WHEN NEW.event_type='call.execute.result' BEGIN SELECT RAISE(ABORT,'injected final outbox failure'); END`); err != nil { @@ -191,7 +191,7 @@ func TestCurrentFinalResultOutboxFailureRollsBackAndCanResume(t *testing.T) { } } -func TestCurrentFinalResultRestartResumesSameOutboxBody(t *testing.T) { +func TestFinalResultRestartResumesSameOutboxBody(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "result-restart-1", true) original, created, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)) @@ -206,7 +206,7 @@ func TestCurrentFinalResultRestartResumesSameOutboxBody(t *testing.T) { if err := s.Close(); err != nil { t.Fatal(err) } - restarted, err := OpenCurrent(databasePath) + restarted, err := Open(databasePath) if err != nil { t.Fatal(err) } diff --git a/internal/store/current_sip.go b/internal/store/sip.go similarity index 89% rename from internal/store/current_sip.go rename to internal/store/sip.go index 5247dd1..b21fbda 100644 --- a/internal/store/current_sip.go +++ b/internal/store/sip.go @@ -6,11 +6,11 @@ import ( "fmt" ) -var ErrCurrentSIPPending = errors.New("new SIP revision is pending drain or verified load") +var ErrSIPPending = errors.New("new SIP revision is pending drain or verified load") // NoteSIPChange durably closes admission before the control queue is ACKed. // Superseded notifications cannot roll back an already applied revision. -func (s *CurrentStore) NoteSIPChange(dispatcherID string, revision int64) error { +func (s *Store) NoteSIPChange(dispatcherID string, revision int64) error { if dispatcherID == "" || revision <= 0 { return errors.New("SIP notification requires Dispatcher identity and positive revision") } @@ -41,7 +41,7 @@ func (s *CurrentStore) NoteSIPChange(dispatcherID string, revision int64) error return nil } -func (s *CurrentStore) SIPState(dispatcherID string) (applied, pending int64, err error) { +func (s *Store) SIPState(dispatcherID string) (applied, pending int64, err error) { err = s.db.QueryRow(`SELECT applied_sip_revision,pending_sip_revision FROM dispatcher_state WHERE dispatcher_id=?`, dispatcherID).Scan(&applied, &pending) if err != nil { return 0, 0, fmt.Errorf("read durable SIP revision state: %w", err) @@ -51,7 +51,7 @@ func (s *CurrentStore) SIPState(dispatcherID string) (applied, pending int64, er // OccupiedCalls includes dispatching, dispatched, and unknown executions. // No lease expiry, SIP reload, or broker restart releases an uncertain call. -func (s *CurrentStore) OccupiedCalls(dispatcherID string) (int64, error) { +func (s *Store) OccupiedCalls(dispatcherID string) (int64, error) { var count int64 if err := s.db.QueryRow(`SELECT COUNT(*) FROM dispatcher_inbox WHERE dispatcher_id=? AND status IN ('dispatching','dispatched','unknown')`, dispatcherID).Scan(&count); err != nil { return 0, fmt.Errorf("count calls before SIP reload: %w", err) @@ -62,7 +62,7 @@ func (s *CurrentStore) OccupiedCalls(dispatcherID string) (int64, error) { // MarkReadyForSIP opens admission only after all present task bindings match // this exact verified SIP revision and any prior SIP calls have drained. // The loaded Agent/Asterisk state is independently checked by the caller. -func (s *CurrentStore) MarkReadyForSIP(dispatcherID string, revision int64) error { +func (s *Store) MarkReadyForSIP(dispatcherID string, revision int64) error { if dispatcherID == "" || revision <= 0 { return errors.New("admission requires a verified SIP revision") } @@ -83,7 +83,7 @@ func (s *CurrentStore) MarkReadyForSIP(dispatcherID string, revision int64) erro return fmt.Errorf("SIP revision %d regresses from applied %d", revision, applied) } if revision < pending { - return fmt.Errorf("%w: verified revision %d is behind pending %d", ErrCurrentSIPPending, revision, pending) + return fmt.Errorf("%w: verified revision %d is behind pending %d", ErrSIPPending, revision, pending) } var missing int64 if err := tx.QueryRow(`SELECT COUNT(*) FROM dispatcher_tasks t LEFT JOIN dispatcher_configs c @@ -102,7 +102,7 @@ func (s *CurrentStore) MarkReadyForSIP(dispatcherID string, revision int64) erro return fmt.Errorf("check old calls before SIP admission: %w", err) } if occupied != 0 { - return fmt.Errorf("%w: %d old or unknown calls have not drained for revision %d", ErrCurrentSIPPending, occupied, revision) + return fmt.Errorf("%w: %d old or unknown calls have not drained for revision %d", ErrSIPPending, occupied, revision) } } result, err := tx.Exec(`UPDATE dispatcher_state SET discovery_ready=1,applied_sip_revision=?,pending_sip_revision=0 WHERE dispatcher_id=?`, revision, dispatcherID) diff --git a/internal/store/current_sip_change_test.go b/internal/store/sip_change_test.go similarity index 93% rename from internal/store/current_sip_change_test.go rename to internal/store/sip_change_test.go index 293521a..d227162 100644 --- a/internal/store/current_sip_change_test.go +++ b/internal/store/sip_change_test.go @@ -5,7 +5,7 @@ import ( "testing" ) -func TestCurrentSIPNotificationRequiresFullDrainAndExactLoadedRevision(t *testing.T) { +func TestSIPNotificationRequiresFullDrainAndExactLoadedRevision(t *testing.T) { s := preparedCurrentCallStore(t) if err := s.MarkReadyForSIP(currentDispatcherID, 8); err != nil { t.Fatal(err) @@ -68,9 +68,9 @@ func TestCurrentSIPNotificationRequiresFullDrainAndExactLoadedRevision(t *testin } } -func TestCurrentPendingSIPChangeSurvivesRestartWithoutOldDataDeletion(t *testing.T) { +func TestPendingSIPChangeSurvivesRestartWithoutOldDataDeletion(t *testing.T) { path := filepath.Join(t.TempDir(), "state.db") - s, err := OpenCurrent(path) + s, err := Open(path) if err != nil { t.Fatal(err) } @@ -80,7 +80,7 @@ func TestCurrentPendingSIPChangeSurvivesRestartWithoutOldDataDeletion(t *testing if err := s.Close(); err != nil { t.Fatal(err) } - reopened, err := OpenCurrent(path) + reopened, err := Open(path) if err != nil { t.Fatal(err) } diff --git a/internal/store/current_snapshot_test.go b/internal/store/snapshot_test.go similarity index 91% rename from internal/store/current_snapshot_test.go rename to internal/store/snapshot_test.go index 015b7a5..6462d9e 100644 --- a/internal/store/current_snapshot_test.go +++ b/internal/store/snapshot_test.go @@ -6,9 +6,9 @@ import ( "testing" ) -func TestCurrentStoreSnapshotPreservesAgentAndAllowsNewSIPRevision(t *testing.T) { +func TestStoreSnapshotPreservesAgentAndAllowsNewSIPRevision(t *testing.T) { path := filepath.Join(t.TempDir(), "state.db") - s, err := OpenCurrent(path) + s, err := Open(path) if err != nil { t.Fatal(err) } @@ -19,7 +19,7 @@ func TestCurrentStoreSnapshotPreservesAgentAndAllowsNewSIPRevision(t *testing.T) if err := s.Close(); err != nil { t.Fatal(err) } - s, err = OpenCurrent(path) + s, err = Open(path) if err != nil { t.Fatal(err) } diff --git a/internal/store/current.go b/internal/store/store.go similarity index 93% rename from internal/store/current.go rename to internal/store/store.go index 73856f7..ea5f432 100644 --- a/internal/store/current.go +++ b/internal/store/store.go @@ -14,16 +14,16 @@ import ( _ "modernc.org/sqlite" ) -// CurrentStore refuses an incompatible database instead of clearing or +// Store refuses an incompatible database instead of clearing or // replaying data whose owner and message identity cannot be proved safe. -type CurrentStore struct{ db *sql.DB } +type Store struct{ db *sql.DB } var ( ErrCommandConflict = errors.New("command id reused with different body") ErrDispatcherIdentityMismatch = errors.New("database belongs to another dispatcher") ) -const currentSchema = ` +const schema = ` CREATE TABLE IF NOT EXISTS dispatcher_state ( dispatcher_id TEXT PRIMARY KEY, discovery_ready INTEGER NOT NULL DEFAULT 0 CHECK(discovery_ready IN (0,1)), @@ -97,7 +97,7 @@ CREATE TABLE IF NOT EXISTS dispatcher_recordings ( FOREIGN KEY(dispatcher_id,source_event_id) REFERENCES dispatcher_inbox(dispatcher_id,event_id) );` -func OpenCurrent(path string) (_ *CurrentStore, err error) { +func Open(path string) (_ *Store, err error) { if path == "" { return nil, errors.New("SQLite path is required") } @@ -176,7 +176,7 @@ func OpenCurrent(path string) (_ *CurrentStore, err error) { return nil, err } defer tx.Rollback() - if _, err := tx.Exec(currentSchema); err != nil { + if _, err := tx.Exec(schema); err != nil { return nil, fmt.Errorf("create current SQLite schema: %w", err) } if _, err := tx.Exec(`PRAGMA user_version=2`); err != nil { @@ -185,16 +185,16 @@ func OpenCurrent(path string) (_ *CurrentStore, err error) { if err := tx.Commit(); err != nil { return nil, fmt.Errorf("commit current SQLite schema: %w", err) } - } else if _, err := db.Exec(currentSchema); err != nil { + } else if _, err := db.Exec(schema); err != nil { return nil, fmt.Errorf("verify current SQLite schema: %w", err) } - return &CurrentStore{db: db}, nil + return &Store{db: db}, nil } -func (s *CurrentStore) Close() error { return s.db.Close() } +func (s *Store) Close() error { return s.db.Close() } // CloseAdmission is durable, including when startup fails before SIP loading. -func (s *CurrentStore) CloseAdmission(dispatcherID string) error { +func (s *Store) CloseAdmission(dispatcherID string) error { if dispatcherID == "" { return errors.New("dispatcher ID is required") } @@ -223,7 +223,7 @@ func (s *CurrentStore) CloseAdmission(dispatcherID string) error { // ApplyDiscoverySnapshot commits the entire cold-start list in one transaction. // It closes admission until the caller has drained the control queue. -func (s *CurrentStore) ApplyDiscoverySnapshot(dispatcherID string, tasks []configread.CurrentDiscoveredTask) (err error) { +func (s *Store) ApplyDiscoverySnapshot(dispatcherID string, tasks []configread.CurrentDiscoveredTask) (err error) { if dispatcherID == "" { return errors.New("dispatcher ID is required") } @@ -247,7 +247,7 @@ func (s *CurrentStore) ApplyDiscoverySnapshot(dispatcherID string, tasks []confi // ApplyDiscoveryPage commits one delta page before its in-memory cursor may // advance. It never treats tasks absent from a delta page as retired. -func (s *CurrentStore) ApplyDiscoveryPage(dispatcherID string, tasks []configread.CurrentDiscoveredTask) error { +func (s *Store) ApplyDiscoveryPage(dispatcherID string, tasks []configread.CurrentDiscoveredTask) error { if dispatcherID == "" || len(tasks) == 0 { return errors.New("discovery delta requires a Dispatcher and nonempty page") } @@ -320,7 +320,7 @@ func validTaskStatus(status string) bool { return status == "running" || status == "paused" || status == "stopped" } -type CurrentAssignedTask struct { +type AssignedTask struct { TenantID int64 TaskID string TaskRevision int64 @@ -330,16 +330,16 @@ type CurrentAssignedTask struct { // ListAssignedTasks is used to attach only the task queues assigned to this // Dispatcher. Stopped task queues can still drain old, unaccepted commands. -func (s *CurrentStore) ListAssignedTasks(dispatcherID string) ([]CurrentAssignedTask, error) { +func (s *Store) ListAssignedTasks(dispatcherID string) ([]AssignedTask, error) { rows, err := s.db.Query(`SELECT tenant_id,task_id,task_revision,status,control_state FROM dispatcher_tasks WHERE dispatcher_id=? AND present=1 ORDER BY task_id`, dispatcherID) if err != nil { return nil, fmt.Errorf("list assigned task queues: %w", err) } defer rows.Close() - var assigned []CurrentAssignedTask + var assigned []AssignedTask for rows.Next() { - var task CurrentAssignedTask + var task AssignedTask if err := rows.Scan(&task.TenantID, &task.TaskID, &task.TaskRevision, &task.Status, &task.ControlState); err != nil { return nil, fmt.Errorf("read assigned task queue: %w", err) } @@ -353,7 +353,7 @@ func (s *CurrentStore) ListAssignedTasks(dispatcherID string) ([]CurrentAssigned // CanAdmit checks persisted discovery and human controls; it does not replace // call-time whitelist, schedule, SIP load, quota or authorization checks. -func (s *CurrentStore) CanAdmit(dispatcherID string, tenantID int64, taskID string) (bool, error) { +func (s *Store) CanAdmit(dispatcherID string, tenantID int64, taskID string) (bool, error) { var count int err := s.db.QueryRow(`SELECT COUNT(*) FROM dispatcher_tasks t JOIN dispatcher_state ds ON ds.dispatcher_id=t.dispatcher_id @@ -368,7 +368,7 @@ func (s *CurrentStore) CanAdmit(dispatcherID string, tenantID int64, taskID stri // SaveSnapshot refuses an immutable task revision with different content. Only // the providers actually referenced by this task are included in its digest. -func (s *CurrentStore) SaveSnapshot(snapshot configread.CurrentSnapshot) error { +func (s *Store) SaveSnapshot(snapshot configread.CurrentSnapshot) error { task := snapshot.Task if task.DispatcherID == "" || task.TenantID <= 0 || task.TaskID == "" || task.TaskRevision <= 0 || snapshot.SIP.DispatcherID != task.DispatcherID || snapshot.Quota.DispatcherID != task.DispatcherID || snapshot.Quota.TenantID != task.TenantID || snapshot.SIP.Revision <= 0 || snapshot.Quota.QuotaRevision <= 0 || !json.Valid(task.Raw) { return errors.New("invalid current task snapshot identity, revision, or body") @@ -413,7 +413,7 @@ func (s *CurrentStore) SaveSnapshot(snapshot configread.CurrentSnapshot) error { return err } defer tx.Rollback() - if err := validateCurrentGlobalRevisions(tx, snapshot); err != nil { + if err := validateGlobalRevisions(tx, snapshot); err != nil { return err } var oldRevision int64 @@ -445,10 +445,10 @@ func (s *CurrentStore) SaveSnapshot(snapshot configread.CurrentSnapshot) error { return tx.Commit() } -// validateCurrentGlobalRevisions prevents two task bindings from silently +// validateGlobalRevisions prevents two task bindings from silently // disagreeing about the same approved Dispatcher SIP or tenant quota revision. // Older revisions also cannot replace a newer binding across tasks. -func validateCurrentGlobalRevisions(tx *sql.Tx, snapshot configread.CurrentSnapshot) error { +func validateGlobalRevisions(tx *sql.Tx, snapshot configread.CurrentSnapshot) error { newSIP, err := json.Marshal(snapshot.SIP) if err != nil { return fmt.Errorf("encode approved SIP revision: %w", err) @@ -515,7 +515,7 @@ func validateCurrentGlobalRevisions(tx *sql.Tx, snapshot configread.CurrentSnaps // ReadSnapshot returns only the current durable binding. Corrupt or mismatched // data is an error, never a signal to fetch an old contract instead. -func (s *CurrentStore) ReadSnapshot(dispatcherID string, tenantID int64, taskID string) (configread.CurrentSnapshot, error) { +func (s *Store) ReadSnapshot(dispatcherID string, tenantID int64, taskID string) (configread.CurrentSnapshot, error) { var body []byte var taskRevision, sipRevision, quotaRevision int64 err := s.db.QueryRow(`SELECT task_revision,sip_revision,quota_revision,snapshot_json FROM dispatcher_configs WHERE dispatcher_id=? AND tenant_id=? AND task_id=?`, dispatcherID, tenantID, taskID).Scan(&taskRevision, &sipRevision, "aRevision, &body) @@ -541,7 +541,7 @@ func (s *CurrentStore) ReadSnapshot(dispatcherID string, tenantID int64, taskID return configread.CurrentSnapshot{Task: task, Providers: decoded.Providers, SIP: decoded.SIP, Quota: decoded.Quota}, nil } -func (s *CurrentStore) ApplyControl(dispatcherID string, tenantID int64, taskID, action string) error { +func (s *Store) ApplyControl(dispatcherID string, tenantID int64, taskID, action string) error { if dispatcherID == "" || tenantID <= 0 || taskID == "" { return errors.New("invalid task control identity") } diff --git a/internal/store/current_test.go b/internal/store/store_test.go similarity index 90% rename from internal/store/current_test.go rename to internal/store/store_test.go index fdac7ea..8845e15 100644 --- a/internal/store/current_test.go +++ b/internal/store/store_test.go @@ -41,9 +41,9 @@ func currentStoreSnapshot(t *testing.T) configread.CurrentSnapshot { return s } -func TestCurrentStoreFreshSnapshotAndControlSurviveRestart(t *testing.T) { +func TestStoreFreshSnapshotAndControlSurviveRestart(t *testing.T) { path := filepath.Join(t.TempDir(), "current.db") - s, err := OpenCurrent(path) + s, err := Open(path) if err != nil { t.Fatal(err) } @@ -69,7 +69,7 @@ func TestCurrentStoreFreshSnapshotAndControlSurviveRestart(t *testing.T) { if err := s.Close(); err != nil { t.Fatal(err) } - s, err = OpenCurrent(path) + s, err = Open(path) if err != nil { t.Fatal(err) } @@ -94,8 +94,8 @@ func TestCurrentStoreFreshSnapshotAndControlSurviveRestart(t *testing.T) { } } -func TestCurrentStoreRejectsSameRevisionDifferentContent(t *testing.T) { - s, err := OpenCurrent(filepath.Join(t.TempDir(), "current.db")) +func TestStoreRejectsSameRevisionDifferentContent(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "current.db")) if err != nil { t.Fatal(err) } @@ -111,7 +111,7 @@ func TestCurrentStoreRejectsSameRevisionDifferentContent(t *testing.T) { } } -func TestCurrentStoreRefusesOldSchemaWithoutDeletingRows(t *testing.T) { +func TestStoreRefusesOldSchemaWithoutDeletingRows(t *testing.T) { path := filepath.Join(t.TempDir(), "legacy.db") legacy, err := sql.Open("sqlite", path) if err != nil { @@ -121,7 +121,7 @@ func TestCurrentStoreRefusesOldSchemaWithoutDeletingRows(t *testing.T) { t.Fatal(err) } legacy.Close() - if s, err := OpenCurrent(path); err == nil { + if s, err := Open(path); err == nil { s.Close() t.Fatal("opened an old database without an explicit data decision") } @@ -136,7 +136,7 @@ func TestCurrentStoreRefusesOldSchemaWithoutDeletingRows(t *testing.T) { } } -func TestCurrentStoreRefusesPreviousCurrentLayoutBeforeModifyingDatabase(t *testing.T) { +func TestStoreRefusesPreviousCurrentLayoutBeforeModifyingDatabase(t *testing.T) { path := filepath.Join(t.TempDir(), "previous-current.db") old, err := sql.Open("sqlite", path) if err != nil { @@ -148,7 +148,7 @@ func TestCurrentStoreRefusesPreviousCurrentLayoutBeforeModifyingDatabase(t *test if err := old.Close(); err != nil { t.Fatal(err) } - if opened, err := OpenCurrent(path); err == nil { + if opened, err := Open(path); err == nil { opened.Close() t.Fatal("accepted prior current SQLite layout without explicit data decision") } diff --git a/internal/store/current_upload.go b/internal/store/upload.go similarity index 67% rename from internal/store/current_upload.go rename to internal/store/upload.go index d95130d..e6035ea 100644 --- a/internal/store/current_upload.go +++ b/internal/store/upload.go @@ -8,14 +8,14 @@ import ( "strings" ) -var ErrCurrentUploadConflict = errors.New("recording upload differs from its original approved target") -var ErrCurrentUploadNotFound = errors.New("no durable original recording target") -var ErrCurrentUploadAlreadyConfirmed = errors.New("recording upload already confirmed; another PUT is forbidden") +var ErrUploadConflict = errors.New("recording upload differs from its original approved target") +var ErrUploadNotFound = errors.New("no durable original recording target") +var ErrUploadAlreadyConfirmed = errors.New("recording upload already confirmed; another PUT is forbidden") -// CurrentUploadProof is the authenticated Agent's successful PUT observation; +// UploadProof is the authenticated Agent's successful PUT observation; // the Dispatcher also verifies every asset field against its durable grant. // It does not claim an independent OSS HEAD or SaaS application receipt. -type CurrentUploadProof struct { +type UploadProof struct { UploadID string RecordingID string StatusCode int @@ -23,9 +23,9 @@ type CurrentUploadProof struct { SHA256 string } -// CurrentRecordingGrant holds the original D-approved OSS destination. Signed +// RecordingGrant holds the original D-approved OSS destination. Signed // URLs, headers and short-lived tokens are deliberately never stored here. -type CurrentRecordingGrant struct { +type RecordingGrant struct { DispatcherID string SourceEventID string UploadID string @@ -44,7 +44,7 @@ type CurrentRecordingGrant struct { // RequireReservedCall ties Agent-reported facts to the numeric tenant and the // durable, already reserved execution. Callers must also authenticate the // active Agent session before accepting any report or signing an OSS token. -func (s *CurrentStore) RequireReservedCall(dispatcherID, sourceEventID string, tenantID int64) error { +func (s *Store) RequireReservedCall(dispatcherID, sourceEventID string, tenantID int64) error { if dispatcherID == "" || sourceEventID == "" || tenantID <= 0 { return errors.New("recording source lacks a Dispatcher execution identity") } @@ -60,78 +60,78 @@ func (s *CurrentStore) RequireReservedCall(dispatcherID, sourceEventID string, t return nil } -func (s *CurrentStore) BindRecordingUpload(proposed CurrentRecordingGrant) (CurrentRecordingGrant, bool, error) { +func (s *Store) BindRecordingUpload(proposed RecordingGrant) (RecordingGrant, bool, error) { if proposed.DispatcherID == "" || proposed.SourceEventID == "" || proposed.UploadID == "" || proposed.RecordingID == "" || proposed.Bucket == "" || proposed.ObjectKey == "" || proposed.SizeBytes <= 0 || proposed.Format != "wav" || (proposed.Channels != 1 && proposed.Channels != 2) || proposed.SampleRateHz != 16000 || proposed.DurationMS < 0 || proposed.ConfirmedAt != "" { - return CurrentRecordingGrant{}, false, errors.New("recording grant lacks a bounded original asset and execution") + return RecordingGrant{}, false, errors.New("recording grant lacks a bounded original asset and execution") } if len(proposed.ChecksumSHA256) != 64 || strings.ToLower(proposed.ChecksumSHA256) != proposed.ChecksumSHA256 { - return CurrentRecordingGrant{}, false, errors.New("recording grant requires a canonical SHA-256 digest") + return RecordingGrant{}, false, errors.New("recording grant requires a canonical SHA-256 digest") } if _, err := hex.DecodeString(proposed.ChecksumSHA256); err != nil { - return CurrentRecordingGrant{}, false, errors.New("recording grant has an invalid SHA-256 digest") + return RecordingGrant{}, false, errors.New("recording grant has an invalid SHA-256 digest") } tx, err := s.db.Begin() if err != nil { - return CurrentRecordingGrant{}, false, err + return RecordingGrant{}, false, err } defer tx.Rollback() var status, trunkID string var snapshot []byte err = tx.QueryRow(`SELECT status,COALESCE(selected_trunk_id,''),snapshot_json FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, proposed.DispatcherID, proposed.SourceEventID).Scan(&status, &trunkID, &snapshot) if errors.Is(err, sql.ErrNoRows) { - return CurrentRecordingGrant{}, false, ErrCurrentUploadNotFound + return RecordingGrant{}, false, ErrUploadNotFound } if err != nil { - return CurrentRecordingGrant{}, false, fmt.Errorf("load approved call for recording: %w", err) + return RecordingGrant{}, false, fmt.Errorf("load approved call for recording: %w", err) } if status != "dispatching" && status != "dispatched" && status != "unknown" && status != "finished" || trunkID == "" || len(snapshot) == 0 { - return CurrentRecordingGrant{}, false, errors.New("recording grant requires an already reserved approved call") + return RecordingGrant{}, false, errors.New("recording grant requires an already reserved approved call") } inserted, err := tx.Exec(`INSERT INTO dispatcher_recordings(dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms) VALUES(?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(dispatcher_id,source_event_id) DO NOTHING`, proposed.DispatcherID, proposed.SourceEventID, proposed.UploadID, proposed.RecordingID, proposed.Bucket, proposed.ObjectKey, proposed.ChecksumSHA256, proposed.SizeBytes, proposed.Format, proposed.Channels, proposed.SampleRateHz, proposed.DurationMS) if err != nil { - return CurrentRecordingGrant{}, false, fmt.Errorf("persist immutable recording target: %w", err) + return RecordingGrant{}, false, fmt.Errorf("persist immutable recording target: %w", err) } count, err := inserted.RowsAffected() if err != nil { - return CurrentRecordingGrant{}, false, err + return RecordingGrant{}, false, err } stored, err := loadRecordingUpload(tx.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, proposed.DispatcherID, proposed.SourceEventID)) if err != nil { - return CurrentRecordingGrant{}, false, fmt.Errorf("inspect persisted recording target: %w", err) + return RecordingGrant{}, false, fmt.Errorf("inspect persisted recording target: %w", err) } if stored.ConfirmedAt != "" { - return CurrentRecordingGrant{}, false, ErrCurrentUploadAlreadyConfirmed + return RecordingGrant{}, false, ErrUploadAlreadyConfirmed } if stored != proposed { - return CurrentRecordingGrant{}, false, ErrCurrentUploadConflict + return RecordingGrant{}, false, ErrUploadConflict } var finalized int err = tx.QueryRow(`SELECT 1 FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, proposed.DispatcherID, currentResultEventID(proposed.DispatcherID, proposed.SourceEventID)).Scan(&finalized) if err == nil { - return CurrentRecordingGrant{}, false, ErrCurrentResultConflict + return RecordingGrant{}, false, ErrResultConflict } if !errors.Is(err, sql.ErrNoRows) { - return CurrentRecordingGrant{}, false, fmt.Errorf("check existing call result before granting an upload: %w", err) + return RecordingGrant{}, false, fmt.Errorf("check existing call result before granting an upload: %w", err) } if err := tx.Commit(); err != nil { - return CurrentRecordingGrant{}, false, fmt.Errorf("commit original recording target: %w", err) + return RecordingGrant{}, false, fmt.Errorf("commit original recording target: %w", err) } return stored, count == 1, nil } -func (s *CurrentStore) LoadRecordingUpload(dispatcherID, sourceEventID string) (CurrentRecordingGrant, error) { +func (s *Store) LoadRecordingUpload(dispatcherID, sourceEventID string) (RecordingGrant, error) { if dispatcherID == "" || sourceEventID == "" { - return CurrentRecordingGrant{}, ErrCurrentUploadNotFound + return RecordingGrant{}, ErrUploadNotFound } return loadRecordingUpload(s.db.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, dispatcherID, sourceEventID)) } -func loadRecordingUpload(row *sql.Row) (CurrentRecordingGrant, error) { - var grant CurrentRecordingGrant +func loadRecordingUpload(row *sql.Row) (RecordingGrant, error) { + var grant RecordingGrant err := row.Scan(&grant.DispatcherID, &grant.SourceEventID, &grant.UploadID, &grant.RecordingID, &grant.Bucket, &grant.ObjectKey, &grant.ChecksumSHA256, &grant.SizeBytes, &grant.Format, &grant.Channels, &grant.SampleRateHz, &grant.DurationMS, &grant.ConfirmedAt) if errors.Is(err, sql.ErrNoRows) { - return CurrentRecordingGrant{}, ErrCurrentUploadNotFound + return RecordingGrant{}, ErrUploadNotFound } return grant, err } diff --git a/internal/store/current_upload_result_test.go b/internal/store/upload_result_test.go similarity index 81% rename from internal/store/current_upload_result_test.go rename to internal/store/upload_result_test.go index e12b91e..cb8363f 100644 --- a/internal/store/current_upload_result_test.go +++ b/internal/store/upload_result_test.go @@ -7,7 +7,7 @@ import ( "testing" ) -func currentUploadedResultPayload(t *testing.T, binding CurrentRecordingGrant) []byte { +func currentUploadedResultPayload(t *testing.T, binding RecordingGrant) []byte { t.Helper() var payload map[string]any if err := json.Unmarshal(currentResultPayload(t), &payload); err != nil { @@ -29,26 +29,26 @@ func currentUploadedResultPayload(t *testing.T, binding CurrentRecordingGrant) [ return raw } -func currentUploadProof(binding CurrentRecordingGrant) CurrentUploadProof { - return CurrentUploadProof{UploadID: binding.UploadID, RecordingID: binding.RecordingID, StatusCode: 200, SizeBytes: binding.SizeBytes, SHA256: binding.ChecksumSHA256} +func currentUploadProof(binding RecordingGrant) UploadProof { + return UploadProof{UploadID: binding.UploadID, RecordingID: binding.RecordingID, StatusCode: 200, SizeBytes: binding.SizeBytes, SHA256: binding.ChecksumSHA256} } -func TestCurrentUploadedResultRequiresOriginalGrantAndConfirmedPUT(t *testing.T) { +func TestUploadedResultRequiresOriginalGrantAndConfirmedPUT(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "uploaded-result-1", true) binding := currentUploadBinding(cmd) payload := currentUploadedResultPayload(t, binding) proof := currentUploadProof(binding) - if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof); !errors.Is(err, ErrCurrentUploadUnverified) { + if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof); !errors.Is(err, ErrUploadUnverified) { t.Fatalf("ungranted asset was reported as uploaded: %v", err) } if _, _, err := s.BindRecordingUpload(binding); err != nil { t.Fatal(err) } - if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)); !errors.Is(err, ErrCurrentUploadUnverified) { + if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)); !errors.Is(err, ErrUploadUnverified) { t.Fatalf("failed OSS path bypassed itself with an empty recording: %v", err) } - if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, CurrentUploadProof{}); !errors.Is(err, ErrCurrentUploadUnverified) { + if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, UploadProof{}); !errors.Is(err, ErrUploadUnverified) { t.Fatalf("missing PUT confirmation was accepted: %v", err) } original, created, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof) @@ -63,12 +63,12 @@ func TestCurrentUploadedResultRequiresOriginalGrantAndConfirmedPUT(t *testing.T) if err != nil || created || repeated.EventID != original.EventID { t.Fatalf("same PUT fact created a second MQ result: created=%t err=%v", created, err) } - if _, _, err := s.BindRecordingUpload(binding); !errors.Is(err, ErrCurrentUploadAlreadyConfirmed) { + if _, _, err := s.BindRecordingUpload(binding); !errors.Is(err, ErrUploadAlreadyConfirmed) { t.Fatalf("confirmed upload received another grant: %v", err) } } -func TestCurrentUploadedResultRejectsChangedOSSMetadataOrProof(t *testing.T) { +func TestUploadedResultRejectsChangedOSSMetadataOrProof(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "uploaded-changed-1", true) binding := currentUploadBinding(cmd) @@ -96,17 +96,17 @@ func TestCurrentUploadedResultRejectsChangedOSSMetadataOrProof(t *testing.T) { } for _, change := range []struct { name string - apply func(*CurrentUploadProof) + apply func(*UploadProof) }{ - {"status", func(p *CurrentUploadProof) { p.StatusCode = 503 }}, - {"upload ID", func(p *CurrentUploadProof) { p.UploadID = "new-upload" }}, - {"recording ID", func(p *CurrentUploadProof) { p.RecordingID = "new-recording" }}, - {"checksum", func(p *CurrentUploadProof) { p.SHA256 = strings.Repeat("b", 64) }}, - {"size", func(p *CurrentUploadProof) { p.SizeBytes++ }}, + {"status", func(p *UploadProof) { p.StatusCode = 503 }}, + {"upload ID", func(p *UploadProof) { p.UploadID = "new-upload" }}, + {"recording ID", func(p *UploadProof) { p.RecordingID = "new-recording" }}, + {"checksum", func(p *UploadProof) { p.SHA256 = strings.Repeat("b", 64) }}, + {"size", func(p *UploadProof) { p.SizeBytes++ }}, } { altered := proof change.apply(&altered) - if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, originalPayload, altered); !errors.Is(err, ErrCurrentUploadUnverified) { + if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, originalPayload, altered); !errors.Is(err, ErrUploadUnverified) { t.Fatalf("changed %s proof was accepted: %v", change.name, err) } } @@ -115,7 +115,7 @@ func TestCurrentUploadedResultRejectsChangedOSSMetadataOrProof(t *testing.T) { } } -func TestCurrentUploadedResultOutboxFailureDoesNotConfirmGrant(t *testing.T) { +func TestUploadedResultOutboxFailureDoesNotConfirmGrant(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "uploaded-outbox-fault-1", true) binding := currentUploadBinding(cmd) @@ -140,21 +140,21 @@ func TestCurrentUploadedResultOutboxFailureDoesNotConfirmGrant(t *testing.T) { } } -func TestCurrentUploadCannotBeginAfterEmptyFinalResult(t *testing.T) { +func TestUploadCannotBeginAfterEmptyFinalResult(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "empty-result-before-grant", true) if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)); err != nil { t.Fatal(err) } - if _, _, err := s.BindRecordingUpload(currentUploadBinding(cmd)); !errors.Is(err, ErrCurrentResultConflict) { + if _, _, err := s.BindRecordingUpload(currentUploadBinding(cmd)); !errors.Is(err, ErrResultConflict) { t.Fatalf("already final no-recording call gained a new upload target: %v", err) } - if _, err := s.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID); !errors.Is(err, ErrCurrentUploadNotFound) { + if _, err := s.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID); !errors.Is(err, ErrUploadNotFound) { t.Fatalf("unexpected recording for finalized empty result: %v", err) } } -func TestCurrentUploadedResultRestartReusesGrantAndPendingOutbox(t *testing.T) { +func TestUploadedResultRestartReusesGrantAndPendingOutbox(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "uploaded-restart-1", true) binding := currentUploadBinding(cmd) @@ -174,7 +174,7 @@ func TestCurrentUploadedResultRestartReusesGrantAndPendingOutbox(t *testing.T) { if err := s.Close(); err != nil { t.Fatal(err) } - restarted, err := OpenCurrent(databasePath) + restarted, err := Open(databasePath) if err != nil { t.Fatal(err) } @@ -190,7 +190,7 @@ func TestCurrentUploadedResultRestartReusesGrantAndPendingOutbox(t *testing.T) { if repeated, created, err := restarted.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof); err != nil || created || repeated.EventID != original.EventID { t.Fatalf("restarted D created a second result: created=%t err=%v", created, err) } - if _, _, err := restarted.BindRecordingUpload(binding); !errors.Is(err, ErrCurrentUploadAlreadyConfirmed) { + if _, _, err := restarted.BindRecordingUpload(binding); !errors.Is(err, ErrUploadAlreadyConfirmed) { t.Fatalf("restarted D issued another PUT grant: %v", err) } } diff --git a/internal/store/current_upload_test.go b/internal/store/upload_test.go similarity index 84% rename from internal/store/current_upload_test.go rename to internal/store/upload_test.go index 7285542..923cadc 100644 --- a/internal/store/current_upload_test.go +++ b/internal/store/upload_test.go @@ -10,8 +10,8 @@ import ( "time" ) -func currentUploadBinding(cmd CurrentExecuteCommand) CurrentRecordingGrant { - return CurrentRecordingGrant{ +func currentUploadBinding(cmd ExecuteCommand) RecordingGrant { + return RecordingGrant{ DispatcherID: cmd.DispatcherID, SourceEventID: cmd.EventID, UploadID: "upload-" + cmd.EventID, RecordingID: "recording-" + cmd.EventID, Bucket: "mock-bucket", ObjectKey: "tenant/1001/" + cmd.EventID + "/rec.wav", @@ -20,7 +20,7 @@ func currentUploadBinding(cmd CurrentExecuteCommand) CurrentRecordingGrant { } } -func TestCurrentUploadBoundWhileOriginateAcknowledgmentIsInFlight(t *testing.T) { +func TestUploadBoundWhileOriginateAcknowledgmentIsInFlight(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentCall("grant-during-dispatch") if _, _, err := s.RecordExecute(cmd); err != nil { @@ -41,7 +41,7 @@ func TestCurrentUploadBoundWhileOriginateAcknowledgmentIsInFlight(t *testing.T) } } -func TestCurrentUploadBindingPersistsOnlyOneOriginalAssetAcrossRestart(t *testing.T) { +func TestUploadBindingPersistsOnlyOneOriginalAssetAcrossRestart(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "grant-stable-1", true) original := currentUploadBinding(cmd) @@ -64,7 +64,7 @@ func TestCurrentUploadBindingPersistsOnlyOneOriginalAssetAcrossRestart(t *testin if err := s.Close(); err != nil { t.Fatal(err) } - restarted, err := OpenCurrent(databasePath) + restarted, err := Open(databasePath) if err != nil { t.Fatal(err) } @@ -79,7 +79,7 @@ func TestCurrentUploadBindingPersistsOnlyOneOriginalAssetAcrossRestart(t *testin } } -func TestCurrentUploadBindingRejectsChangedTargetAndUnreservedCalls(t *testing.T) { +func TestUploadBindingRejectsChangedTargetAndUnreservedCalls(t *testing.T) { s := preparedCurrentCallStore(t) pending := currentCall("grant-not-reserved") if _, _, err := s.RecordExecute(pending); err != nil { @@ -95,19 +95,19 @@ func TestCurrentUploadBindingRejectsChangedTargetAndUnreservedCalls(t *testing.T } for _, change := range []struct { name string - apply func(*CurrentRecordingGrant) + apply func(*RecordingGrant) }{ - {"upload ID", func(x *CurrentRecordingGrant) { x.UploadID = "new-upload" }}, - {"recording ID", func(x *CurrentRecordingGrant) { x.RecordingID = "other-recording" }}, - {"bucket", func(x *CurrentRecordingGrant) { x.Bucket = "changed-bucket" }}, - {"object key", func(x *CurrentRecordingGrant) { x.ObjectKey = "another/key.wav" }}, - {"checksum", func(x *CurrentRecordingGrant) { x.ChecksumSHA256 = strings.Repeat("b", 64) }}, - {"size", func(x *CurrentRecordingGrant) { x.SizeBytes++ }}, - {"channels", func(x *CurrentRecordingGrant) { x.Channels = 2 }}, + {"upload ID", func(x *RecordingGrant) { x.UploadID = "new-upload" }}, + {"recording ID", func(x *RecordingGrant) { x.RecordingID = "other-recording" }}, + {"bucket", func(x *RecordingGrant) { x.Bucket = "changed-bucket" }}, + {"object key", func(x *RecordingGrant) { x.ObjectKey = "another/key.wav" }}, + {"checksum", func(x *RecordingGrant) { x.ChecksumSHA256 = strings.Repeat("b", 64) }}, + {"size", func(x *RecordingGrant) { x.SizeBytes++ }}, + {"channels", func(x *RecordingGrant) { x.Channels = 2 }}, } { changed := original change.apply(&changed) - if _, _, err := s.BindRecordingUpload(changed); !errors.Is(err, ErrCurrentUploadConflict) { + if _, _, err := s.BindRecordingUpload(changed); !errors.Is(err, ErrUploadConflict) { t.Fatalf("%s replaced the immutable original: %v", change.name, err) } } @@ -116,7 +116,7 @@ func TestCurrentUploadBindingRejectsChangedTargetAndUnreservedCalls(t *testing.T } } -func TestCurrentRecordingSourceMustMatchReservedCallAndNumericTenant(t *testing.T) { +func TestRecordingSourceMustMatchReservedCallAndNumericTenant(t *testing.T) { s := preparedCurrentCallStore(t) cmd := currentResultCall(t, s, "grant-auth-source", true) if err := s.RequireReservedCall(cmd.DispatcherID, cmd.EventID, cmd.TenantID); err != nil { @@ -145,7 +145,7 @@ func TestCurrentRecordingSourceMustMatchReservedCallAndNumericTenant(t *testing. } } -func TestCurrentUploadBindingRejectsReusedUploadIDOrOSSObject(t *testing.T) { +func TestUploadBindingRejectsReusedUploadIDOrOSSObject(t *testing.T) { s := preparedCurrentCallStore(t) first := currentUploadBinding(currentResultCall(t, s, "grant-unique-1", true)) second := currentUploadBinding(currentResultCall(t, s, "grant-unique-2", true)) @@ -161,14 +161,14 @@ func TestCurrentUploadBindingRejectsReusedUploadIDOrOSSObject(t *testing.T) { if _, _, err := s.BindRecordingUpload(second); err == nil { t.Fatal("two calls shared one original OSS object") } - if _, err := s.LoadRecordingUpload(second.DispatcherID, second.SourceEventID); !errors.Is(err, ErrCurrentUploadNotFound) { + if _, err := s.LoadRecordingUpload(second.DispatcherID, second.SourceEventID); !errors.Is(err, ErrUploadNotFound) { t.Fatalf("failed uniqueness check left a partial grant: %v", err) } } -func TestCurrentSchemaRejectsPreviousCompleteLayoutWithoutErasingOutbox(t *testing.T) { +func TestSchemaRejectsPreviousCompleteLayoutWithoutErasingOutbox(t *testing.T) { path := filepath.Join(t.TempDir(), "previous-current.db") - s, err := OpenCurrent(path) + s, err := Open(path) if err != nil { t.Fatal(err) } @@ -184,7 +184,7 @@ func TestCurrentSchemaRejectsPreviousCompleteLayoutWithoutErasingOutbox(t *testi if err := s.Close(); err != nil { t.Fatal(err) } - if reopened, err := OpenCurrent(path); err == nil { + if reopened, err := Open(path); err == nil { reopened.Close() t.Fatal("incomplete previous SQLite layout was silently modified or admitted") }