diff --git a/CHANGELOG.md b/CHANGELOG.md index b8929db..2489a78 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Added + +- Runtime: optional shared invoke concurrency budgets, cooperative cancellation on exit and shutdown, per-universe accumulator and tracking limits, and strict observer registration checks. Existing constructors preserve zero-value policies. +- Runtime: optional context-aware snapshot capture/restoration and invoke lifecycle interfaces; reject synchronous callback reentry when callers preserve the callback context. +- Debugger bot: cancellable event processing, checked snapshot errors, and an optional history retention limit. + ### Security - Studio: metadata pointer writes and object merges use own data properties, preserving JSON keys without traversing inherited properties. @@ -14,6 +20,12 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- Runtime: honor cancellation while waiting for execution locks and between synchronous callbacks; preserve bounded tracking during failed entry rollback. +- CLI: avoid reflection panics when copying contexts with private fields or typed nil pointers; handle empty lists and history independently of saved checkpoints. +- Studio: reject malformed JSON without throwing and resolve definition references only from own properties, while allowing explicitly declared names such as `constructor`. +- Web Component: support registering multiple tag names and cover attributes, properties, lifecycle, and DOM events with integration tests. +- Studio: fit header, machine panel, and search controls within mobile viewports. Keep Babel 7 overrides scoped to Babel 7 so the Stryker 10 instrumenter can use Babel 8. +- Studio tests: exclude Stryker sandboxes from normal Vitest discovery to prevent duplicate or instrumented suites from running as application tests. - Runtime: validate all included universe snapshots before applying them; restore metadata and tracking instead of retaining newer entries, synchronize metadata restoration with invokes, and reconstruct empty superposition accumulators and final-state flags. - Runtime: synchronize custom executor registration and lookup; accept custom implementations of the public `Event` interface in accumulators. - Runtime: return errors for nil machine/universe models and nil events; reject already-canceled `SendEvent` calls before event admission, including cancellation while waiting for the machine lock. diff --git a/debugger/bot/bot.go b/debugger/bot/bot.go index e822d39..ef92c19 100644 --- a/debugger/bot/bot.go +++ b/debugger/bot/bot.go @@ -50,9 +50,13 @@ func NewBot( return nil, fmt.Errorf("event provider cannot be nil") } + initial, err := captureSnapshot(context.Background(), qm) + if err != nil { + return nil, fmt.Errorf("error capturing initial snapshot: %w", err) + } b := &bot{ qm: qm, - initialSnapshot: qm.GetSnapshot(), + initialSnapshot: initial, eventProvider: eventProvider, initQuantumMachine: initQuantumMachine, } @@ -60,7 +64,9 @@ func NewBot( for _, opt := range opts { opt(b) } - + if b.historyLimit < 0 { + return nil, fmt.Errorf("history limit must not be negative") + } return b, nil } @@ -71,6 +77,7 @@ type bot struct { initQuantumMachine bool history []*EventHistory ignoreUnhandledEvents bool + historyLimit int } type BotOption func(*bot) @@ -81,10 +88,46 @@ func WithIgnoreUnhandledEvents(ignore bool) BotOption { } } +// WithHistoryLimit retains the latest entries; zero preserves unlimited history. +func WithHistoryLimit(limit int) BotOption { + return func(b *bot) { b.historyLimit = limit } +} + +func captureSnapshot(ctx context.Context, qm instrumentation.QuantumMachine) (*instrumentation.MachineSnapshot, error) { + var snapshot *instrumentation.MachineSnapshot + var err error + if contextual, ok := qm.(instrumentation.ContextSnapshotProvider); ok { + snapshot, err = contextual.GetSnapshotContext(ctx) + } else if checked, ok := qm.(instrumentation.SnapshotProvider); ok { + snapshot, err = checked.GetSnapshotWithError() + } else { + snapshot = qm.GetSnapshot() + } + if err != nil { + return nil, err + } + if snapshot == nil { + return nil, fmt.Errorf("snapshot capture returned nil") + } + return snapshot, nil +} + func (b *bot) Run(ctx context.Context, machineContext any) error { + if ctx == nil { + return fmt.Errorf("context must not be nil") + } + if err := ctx.Err(); err != nil { + return err + } b.history = nil - if err := b.qm.LoadSnapshot(b.initialSnapshot, machineContext); err != nil { - return fmt.Errorf("error loading initial snapshot: %w", err) + var restoreErr error + if contextual, ok := b.qm.(instrumentation.ContextSnapshotProvider); ok { + restoreErr = contextual.LoadSnapshotContext(ctx, b.initialSnapshot, machineContext) + } else { + restoreErr = b.qm.LoadSnapshot(b.initialSnapshot, machineContext) + } + if restoreErr != nil { + return fmt.Errorf("error loading initial snapshot: %w", restoreErr) } if b.initQuantumMachine { @@ -94,10 +137,20 @@ func (b *bot) Run(ctx context.Context, machineContext any) error { } for { - event, err := b.eventProvider(b.qm.GetSnapshot()) + if err := ctx.Err(); err != nil { + return err + } + snapshot, err := captureSnapshot(ctx, b.qm) + if err != nil { + return fmt.Errorf("error capturing snapshot: %w", err) + } + event, err := b.eventProvider(snapshot) if err != nil { return err } + if err := ctx.Err(); err != nil { + return err + } if event == nil { break } @@ -114,10 +167,19 @@ func (b *bot) Run(ctx context.Context, machineContext any) error { return fmt.Errorf("event '%s' was not handled", event.GetEventName()) } + snapshot, err = captureSnapshot(ctx, b.qm) + if err != nil { + return fmt.Errorf("error capturing event snapshot: %w", err) + } b.history = append(b.history, &EventHistory{ Event: event, - Snapshot: b.qm.GetSnapshot(), + Snapshot: snapshot, }) + if b.historyLimit > 0 && len(b.history) > b.historyLimit { + copy(b.history, b.history[len(b.history)-b.historyLimit:]) + clear(b.history[b.historyLimit:]) + b.history = b.history[:b.historyLimit] + } } return nil diff --git a/debugger/bot/runtime_controls_test.go b/debugger/bot/runtime_controls_test.go new file mode 100644 index 0000000..762ae13 --- /dev/null +++ b/debugger/bot/runtime_controls_test.go @@ -0,0 +1,130 @@ +package bot_test + +import ( + "context" + "errors" + "testing" + "time" + + statepro "github.com/rendis/statepro/v3" + "github.com/rendis/statepro/v3/builtin" + "github.com/rendis/statepro/v3/debugger/bot" + "github.com/rendis/statepro/v3/instrumentation" + "github.com/rendis/statepro/v3/theoretical" +) + +func TestBot_CancelaRestauracionMientrasMaquinaRealEjecutaCallback(t *testing.T) { + entered, release := make(chan struct{}), make(chan struct{}) + if err := builtin.RegisterAction("test:bot-runtime-block", func(context.Context, instrumentation.ActionExecutorArgs) error { + close(entered) + <-release + return nil + }); err != nil { + t.Fatal(err) + } + defer builtin.RegisterAction("test:bot-runtime-block", nil) + initial := "idle" + qm, err := statepro.NewQuantumMachine(&theoretical.QuantumMachineModel{ + ID: "machine", Initials: []string{"U:main"}, + Universes: map[string]*theoretical.UniverseModel{"main": { + ID: "main", Initial: &initial, Realities: map[string]*theoretical.RealityModel{"idle": { + ID: "idle", Type: "transition", EntryActions: []*theoretical.ActionModel{{Src: "test:bot-runtime-block"}}, + }}, + }}, + }) + if err != nil { + t.Fatal(err) + } + b, err := bot.NewBot(qm, func(*instrumentation.MachineSnapshot) (instrumentation.Event, error) { + t.Error("proveedor ejecutado tras cancelacion") + return nil, nil + }, false) + if err != nil { + t.Fatal(err) + } + done := make(chan error, 1) + go func() { done <- qm.Init(context.Background(), nil) }() + defer func() { + close(release) + if err := <-done; err != nil { + t.Error(err) + } + }() + select { + case <-entered: + case <-time.After(time.Second): + t.Fatal("callback no iniciado") + } + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + defer cancel() + if err := b.Run(ctx, nil); !errors.Is(err, context.DeadlineExceeded) { + t.Fatal(err) + } +} + +func TestBot_CancelacionAntesYDuranteProveedor(t *testing.T) { + for _, preCanceled := range []bool{false, true} { + qm := &hostileQM{MockQuantumMachine: MockQuantumMachine{snapshot: &instrumentation.MachineSnapshot{}}} + ctx, cancel := context.WithCancel(context.Background()) + calls := 0 + provider := func(*instrumentation.MachineSnapshot) (instrumentation.Event, error) { + calls++ + cancel() + return &MockEvent{name: "handled"}, nil + } + b, err := bot.NewBot(qm, provider, false) + if err != nil { + t.Fatal(err) + } + if preCanceled { + cancel() + } + if err := b.Run(ctx, nil); !errors.Is(err, context.Canceled) { + t.Fatal(err) + } + if qm.sendCalls != 0 || (preCanceled && calls != 0) { + t.Fatal("cancelacion ejecuta trabajo adicional") + } + } +} + +func TestBot_LimiteRetieneUltimosEventosYRechazaNegativos(t *testing.T) { + qm := &MockQuantumMachine{snapshot: &instrumentation.MachineSnapshot{}} + count := 0 + provider := func(*instrumentation.MachineSnapshot) (instrumentation.Event, error) { + count++ + if count > 5 { + return nil, nil + } + return &MockEvent{name: "handled"}, nil + } + b, err := bot.NewBot(qm, provider, false, bot.WithHistoryLimit(2)) + if err != nil { + t.Fatal(err) + } + if err := b.Run(context.Background(), nil); err != nil { + t.Fatal(err) + } + if len(b.GetHistory()) != 2 { + t.Fatalf("historial sin limite: %d", len(b.GetHistory())) + } + if _, err := bot.NewBot(qm, provider, false, bot.WithHistoryLimit(-1)); err == nil { + t.Fatal("limite negativo aceptado") + } +} + +type failingSnapshotMachine struct { + MockQuantumMachine + err error +} + +func (m *failingSnapshotMachine) GetSnapshotWithError() (*instrumentation.MachineSnapshot, error) { + return nil, m.err +} +func TestBot_NoOcultaErroresDeCaptura(t *testing.T) { + failure := errors.New("snapshot failure") + qm := &failingSnapshotMachine{err: failure} + if _, err := bot.NewBot(qm, func(*instrumentation.MachineSnapshot) (instrumentation.Event, error) { return nil, nil }, false); !errors.Is(err, failure) { + t.Fatal(err) + } +} diff --git a/debugger/cli/helper.go b/debugger/cli/helper.go index d22d7d0..f77af8a 100644 --- a/debugger/cli/helper.go +++ b/debugger/cli/helper.go @@ -151,11 +151,13 @@ func copyStructPointer(v any) any { } val := reflect.ValueOf(v) + if val.IsNil() { + return nil + } structValue := val.Elem() structCopy := reflect.New(structValue.Type()).Elem() - for i := 0; i < structValue.NumField(); i++ { - structCopy.Field(i).Set(structValue.Field(i)) - } + // Copy the complete value; setting unexported fields individually panics. + structCopy.Set(structValue) return structCopy.Addr().Interface() } diff --git a/debugger/cli/history_viewer_model.go b/debugger/cli/history_viewer_model.go index eca9fb1..eb685a7 100644 --- a/debugger/cli/history_viewer_model.go +++ b/debugger/cli/history_viewer_model.go @@ -39,7 +39,7 @@ var historySnapshotKeys = map[string]snapshotExtractor{ } func buildHistoryViewerModel(prevModel *model, container *smContainer) (tea.Model, tea.Cmd) { - if len(container.snapshots) == 0 { + if len(container.history) == 0 { return prevModel, nil } @@ -79,6 +79,9 @@ func historyViewerModelView(m *model) string { hm := m.helperModel.(*list.Model) item, _ := hm.SelectedItem().(*choice) v1 := appStyle.Render(hm.View()) + if item == nil { + return v1 + } h := item.obj.(*containerHistory) part := buildSnapshotPartFromHistory(h, "m") @@ -102,6 +105,9 @@ func historyViewerModelUpdate(m *model, teaMsg tea.Msg) (tea.Model, tea.Cmd) { return m.prevModel, nil case "r": item, _ := hm.SelectedItem().(*choice) + if item == nil { + return m, nil + } h := item.obj.(*containerHistory) if err := m.container.qm.LoadSnapshot(h.snapshot, m.container.smContext); err != nil { m.err = err @@ -115,6 +121,9 @@ func historyViewerModelUpdate(m *model, teaMsg tea.Msg) (tea.Model, tea.Cmd) { break } item, _ := hm.SelectedItem().(*choice) + if item == nil { + return m, nil + } h := item.obj.(*containerHistory) v := buildSnapshotPartFromHistory(h, msg.String()) return buildJsonViewerModel(m, v.title, v.content) diff --git a/debugger/cli/load_snapshot_model.go b/debugger/cli/load_snapshot_model.go index 4d3042f..0a45df7 100644 --- a/debugger/cli/load_snapshot_model.go +++ b/debugger/cli/load_snapshot_model.go @@ -63,6 +63,9 @@ func loadSnapshotModelView(m *model) string { v1 := appStyle.Render(hm.View()) item, _ := hm.SelectedItem().(*choice) + if item == nil { + return v1 + } h := item.obj.(*debuggerSnapshot) part := buildSnapshotPart(h, "m") @@ -85,11 +88,17 @@ func loadSnapshotModelUpdate(m *model, teaMsg tea.Msg) (tea.Model, tea.Cmd) { break } item, _ := hm.SelectedItem().(*choice) + if item == nil { + return m, nil + } dn := item.obj.(*debuggerSnapshot) v := buildSnapshotPart(dn, msg.String()) return buildJsonViewerModel(m, v.title, v.content) case "l": item, _ := hm.SelectedItem().(*choice) + if item == nil { + return m, nil + } dn := item.obj.(*debuggerSnapshot) if err := m.container.qm.LoadSnapshot(dn.Snapshot, m.container.smContext); err != nil { return m, nil diff --git a/debugger/cli/send_event_model.go b/debugger/cli/send_event_model.go index c90246b..e21f022 100644 --- a/debugger/cli/send_event_model.go +++ b/debugger/cli/send_event_model.go @@ -104,6 +104,9 @@ func sendEventModelView(m *model) string { }) v1 := appStyle.Render(hm.View()) + if len(m.container.history) == 0 { + return v1 + } h := m.container.history[len(m.container.history)-1] part := buildSnapshotPartFromHistory(h, "m") @@ -128,11 +131,17 @@ func sendEventModelUpdate(m *model, teaMsg tea.Msg) (tea.Model, tea.Cmd) { if isFiltering(hm) { break } + if len(m.container.history) == 0 { + return m, nil + } h := m.container.history[len(m.container.history)-1] v := buildSnapshotPartFromHistory(h, msg.String()) return buildJsonViewerModel(m, v.title, v.content) case "enter": item, _ := hm.SelectedItem().(*choice) + if item == nil { + return m, nil + } event := item.obj.(*debuggerEvent) evt := statepro.NewEventBuilder(event.Name). SetData(event.Params). diff --git a/debugger/cli/workflows_test.go b/debugger/cli/workflows_test.go new file mode 100644 index 0000000..e486787 --- /dev/null +++ b/debugger/cli/workflows_test.go @@ -0,0 +1,128 @@ +package cli + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "reflect" + "strings" + "testing" + + "github.com/charmbracelet/bubbles/list" + tea "github.com/charmbracelet/bubbletea" + "github.com/google/uuid" + "github.com/rendis/statepro/v3" + "github.com/rendis/statepro/v3/theoretical" +) + +func TestCLI_EnviarEventoYRestaurarHistorial(t *testing.T) { + initial := "idle" + qm, err := statepro.NewQuantumMachine(&theoretical.QuantumMachineModel{ + ID: "machine", Initials: []string{"U:main"}, Universes: map[string]*theoretical.UniverseModel{ + "main": {ID: "main", CanonicalName: "main", Initial: &initial, Realities: map[string]*theoretical.RealityModel{ + "idle": {ID: "idle", Type: theoretical.RealityTypeTransition, On: map[string][]*theoretical.TransitionModel{"GO": {{Targets: []string{"done"}}}}}, + "done": {ID: "done", Type: theoretical.RealityTypeFinal}, + }}, + }, + }) + if err != nil { + t.Fatal(err) + } + if err := qm.Init(context.Background(), nil); err != nil { + t.Fatal(err) + } + snapshot := qm.GetSnapshot() + container := &smContainer{qm: qm, events: []*debuggerEvent{{Name: "GO", Title: "Continue", uid: uuid.New()}}, history: []*containerHistory{{snapshot: snapshot, event: getSnapshotEvent(), pos: 0}}} + prev := &model{} + built, _ := buildSendEventModel(prev, container) + send := built.(*model) + if !strings.Contains(send.View(), "Send an event") { + t.Fatal("vista ausente") + } + send.Update(tea.KeyMsg{Type: tea.KeyEnter}) + if !container.events[0].Sent || len(container.history) != 2 || qm.GetSnapshot().Resume.FinalizedUniverses["main"] != "done" { + t.Fatal("evento no registrado") + } + // History is independent of the saved-snapshot picker. + historyModel, _ := buildHistoryViewerModel(prev, container) + history := historyModel.(*model) + if history == prev || !strings.Contains(history.View(), "View history") { + t.Fatal("historial no disponible") + } + history.Update(tea.KeyMsg{Type: tea.KeyRunes, Runes: []rune{'r'}}) + if len(container.history) != 1 || container.events[0].Sent || qm.GetSnapshot().Resume.ActiveUniverses["main"] != "idle" { + t.Fatal("rollback no restaura estado y marcas") + } + container.snapshots = []*debuggerSnapshot{{Title: "Checkpoint", Snapshot: snapshot}} + loadModel, _ := buildLoadSnapshotModel(prev, container) + load := loadModel.(*model) + if !strings.Contains(load.View(), "Choose a snapshot") { + t.Fatal("vista de snapshots ausente") + } + load.Update(tea.KeyMsg{Type: tea.KeyRunes, Runes: []rune{'l'}}) + if len(container.history) != 1 || container.history[0].pos != 0 { + t.Fatal("carga no reinicia historial") + } +} + +func TestCLI_ListasSinSeleccionNoProvocanPanic(t *testing.T) { + for _, kind := range []string{"events", "snapshots", "history"} { + t.Run(kind, func(t *testing.T) { + l := list.New(nil, list.NewDefaultDelegate(), 80, 24) + m := &model{helperModel: &l, container: &smContainer{}} + switch kind { + case "events": + m.view, m.update = sendEventModelView, sendEventModelUpdate + case "snapshots": + m.view, m.update = loadSnapshotModelView, loadSnapshotModelUpdate + case "history": + m.view, m.update = historyViewerModelView, historyViewerModelUpdate + } + _ = m.View() + for _, msg := range []tea.KeyMsg{{Type: tea.KeyEnter}, {Type: tea.KeyRunes, Runes: []rune{'v'}}, {Type: tea.KeyRunes, Runes: []rune{'l'}}, {Type: tea.KeyRunes, Runes: []rune{'r'}}} { + m.Update(msg) + } + }) + } +} + +func TestCLI_CopiaContextoConCamposPrivadosYJSONConErrores(t *testing.T) { + type applicationContext struct { + Name string + private int + } + original := &applicationContext{Name: "example", private: 7} + copy := copyStructPointer(original).(*applicationContext) + if copy == original || !reflect.DeepEqual(copy, original) { + t.Fatal("copia de contexto incorrecta") + } + copy.Name = "changed" + if original.Name != "example" { + t.Fatal("copia comparte valor de contexto") + } + var missing *applicationContext + if copyStructPointer(missing) != nil { + t.Fatal("nil tipado no preservado") + } + encoded, err := formatToJson(map[string]any{"count": json.Number("9007199254740993")}, false) + if err != nil || !strings.Contains(string(encoded), "9007199254740993") { + t.Fatal("formato pierde precision") + } + if _, err := formatToJson(make(chan int), false); err == nil { + t.Fatal("error JSON oculto") + } + path := filepath.Join(t.TempDir(), "invalid.json") + if err := os.WriteFile(path, []byte("{"), 0600); err != nil { + t.Fatal(err) + } + if err := loadJSON(&map[string]any{}, path, false); err == nil { + t.Fatal("JSON invalido aceptado") + } + if err := loadJSON(&map[string]any{}, path+"missing", true); err != nil { + t.Fatal(err) + } + if err := loadJSON(&map[string]any{}, path+"missing", false); err == nil { + t.Fatal("archivo requerido ausente aceptado") + } +} diff --git a/docs/mutation-fuzzing.md b/docs/mutation-fuzzing.md index 5e27652..c2708f7 100644 --- a/docs/mutation-fuzzing.md +++ b/docs/mutation-fuzzing.md @@ -55,10 +55,23 @@ Mutates `validateStatePro`, `identifiers`, and `transitionRules` by default. Exp Aim to kill **observable** survivors (wrong branch, wrong sentinel, off-by-one on a documented boundary). Do **not** chase arithmetic on buffer capacity, loop-control on uniquely named maps, or sort comparator `<` vs `<=` when IDs are unique — those are usually equivalent mutants. -Approximate baselines after survivor hardening: - -| Package | Tool | Score (ballpark) | -|---------|------|------------------| -| `builtin/` | Gremlins | ~92% (remaining LIVED: capacity arithmetic) | -| `experimental/` | Gremlins | ~89%+ | -| editor-core validators | Stryker | ~45% (narrow suite; raise by expanding mutate + tests) | +Complete campaign measurements on 2026-10-04: + +| Target | Tool | Efficacy / mutation score | Mutant coverage / covered score | +| --- | --- | --- | --- | +| `builtin/` | Gremlins | 92.00% | 100.00% | +| `debugger/bot/` | Gremlins | 96.00% | 100.00% | +| root package and covered dependencies | Gremlins | 83.11% | 85.43% | +| `experimental/` and covered dependencies | Gremlins | 87.31% | 96.50% | +| editor-core configured three source files | Stryker | 58.11% | 65.12% | + +Stryker ran the expanded 59-test suite against 845 mutants: 489 killed, 2 timed out, +263 survived, and 91 uncovered, with no runner errors. The prior 16-test suite scored 48.88% +on 804 mutants; both test scope and guarded source changed, so the totals differ. +Gremlins efficacy excludes timeout results; Stryker includes timeouts in detected mutants. +Covered scores exclude uncovered mutants and are not comparable to statement coverage. + +The normal Vitest configuration excludes `.stryker-tmp` to avoid discovering duplicated or +instrumented test files. Babel overrides retain the major version required by the Stryker +instrumenter. See the [follow-up report](reviews/2026-10-04-runtime-controls.md) for counts, +reproduction methods, performance measurements, and remaining limits. diff --git a/docs/reviews/2026-10-04-runtime-controls.md b/docs/reviews/2026-10-04-runtime-controls.md new file mode 100644 index 0000000..48e3095 --- /dev/null +++ b/docs/reviews/2026-10-04-runtime-controls.md @@ -0,0 +1,122 @@ +# Runtime controls and follow-up validation + +Reviewed against merged `main` at `4d4959e6020ee32033cb69c270a3a74742ec1cba`. +This review covers the runtime, debugger, Studio validation, Web Component, mobile controls, +and the remaining layout/history performance costs. + +## Confirmed problems and changes + +| Area | Reproduction / finding | Result | +| --- | --- | --- | +| Invokes | Machine constants and reality invokes used separate unbounded goroutine launch paths. Leaving a reality did not cancel its work. | A shared manager covers machine constants, universe constants, and reality/transition invokes. Optional capacity rejects admission before spawning; exit cancellation and shutdown are cooperative. Completion and panic release capacity. | +| Accumulation and tracking | Persistent superposition and repeated transitions could grow retained state indefinitely. | Optional entry budgets reserve event fan-out before appending, validate restored accumulators, and retain recent tracking. Failed entry rollback retains the previous bounded history. Bot history has an independent optional limit. | +| Callbacks | Reentering the owning execution lock from a synchronous callback could deadlock; waiting callers could not cancel promptly. | Context-aware lock waiting and callback scope detection reject synchronous reentry, including nested calls across machines. Async invoke contexts retain cancellation and user values without inheriting the synchronous scope. | +| Unknown observers | An unregistered observer approved by default. | `StrictObservers` rejects missing registrations before event accumulation. Zero-value options preserve the existing permissive contract. | +| Bot | Cancellation was ignored between provider iterations and snapshot capture errors were hidden. | Check cancellation around provider/event processing, use optional contextual capture/restoration, propagate errors, and reject nil captures. A real-machine integration test cancels restoration while another callback holds its lock. | +| CLI | Copying a context with private fields panicked through reflection; empty list selection could panic; history navigation depended on saved checkpoints. | Copy the whole struct value, handle typed nil pointers and empty selections, and use actual history availability. Integration tests exercise a real machine, event history, rollback, checkpoint loading, and precise JSON display. | +| Studio validation | Twelve new cases failed before the fix: malformed JSON threw exceptions, and inherited `constructor` / `toString` values resolved as definition entries. | Type guards return structured errors; own-property lookup requires explicitly declared universes and realities. Valid own entries named `constructor` remain supported. | +| Web Component | Registering the same constructor under a second tag caused `NotSupportedError`. | Alternate tags receive subclasses. Tests cover native registration, React mounting, attributes/properties, disconnection/reconnection, and event/callback delivery. | +| Mobile UI | Header controls crowded each other; the machine panel exceeded the viewport; search overlapped it and expanded beyond the screen. | Responsive spacing, compact labeled controls, bounded panels, mobile search placement, and wrapping/stacked bottom toolbars. Browser assertions cover 320, 390, 768, and 1440 px. | +| Mutation tooling | A global Babel 7 override replaced Babel 8 required by Stryker 10, crashing instrumentation of TypeScript generics. | Scope the Babel 7 override to Babel 7 requests. Keep Stryker's Babel 8 dependency tree intact; frozen installation and dependency audit validate the resulting lockfile. | +| Test discovery | Vitest picked up Stryker sandbox copies, doubling the normal suite and treating instrumented files as application tests. | Exclude mutation sandboxes while retaining Vitest's default exclusions. Repeat the normal suite with the configured scope. | + +All resource policies are opt-in. The mandatory machine/executor interfaces remain unchanged. +See [runtime policies](../runtime.md#runtime-policies-and-resource-limits) for budgets, error behavior, +callback context requirements, and shutdown semantics. + +## Validation + +- Go 1.25: build, vet, race-enabled tests with coverage across all 11 packages. +- CLI statement coverage increased from 1.3% to 51.0%; experimental runtime coverage is 90.1%. +- Three native fuzz smoke targets, five seconds each: references, binary definition validation, + and builtin observer arguments. +- Node 24.19 / pnpm 10.11: frozen install, lint, typecheck, 283 tests (280 editor-core and 3 wrapper tests), builds, and full + dependency audit (including development dependencies): zero advisories. +- Chromium: canvas, library, JSON export, bounded mobile header and search, production layout, + and real Web Component locale/event flow. Wrapper unit tests mock only the editor boundary; + browser tests render the actual editor. +- A regression test confirms that consumer mutations cannot change captured checkpoints, + stored undo snapshots, or the restored graph's source snapshot. + +### Complete configured mutation campaigns + +These campaigns executed mutants, rather than stopping after an instrumented dry-run. Scope is +the configured targets below, not every file in the repository. Gremlins also mutates covered +module dependencies; target results overlap and must not be summed as unique mutants. + +| Target | Killed | Survived | Uncovered | Timed out | Efficacy / total score | Mutant coverage / covered score | +| --- | ---: | ---: | ---: | ---: | ---: | ---: | +| Gremlins: builtin | 23 | 2 | 0 | 0 | 92.00% | 100.00% | +| Gremlins: bot | 24 | 1 | 0 | 2 | 96.00% | 100.00% | +| Gremlins: root and covered dependencies | 492 | 100 | 101 | 8 | 83.11% | 85.43% | +| Gremlins: experimental and covered dependencies | 289 | 42 | 12 | 49 | 87.31% | 96.50% | +| Stryker: configured Studio sources | 489 | 263 | 91 | 2 | 58.11% | 65.12% | + +Gremlins efficacy excludes timeouts; Stryker counts them as detected. The root campaign preceded +one final private async-context adjustment; the completed experimental campaign and final Go +race tests cover that resulting runtime. All four Go campaigns exceeded their configured floors. +The final Studio campaign executed 845 mutants with 59 tests and no runner errors. The prior +16-test suite scored 48.88% on 804 mutants; test scope and guarded source changed together. +`transitionRules` improved from 45.22% to 84.35%; `validateStatePro` remains 45.65%, with much +of its issue-path mapping and diagnostic differentiation still weakly tested. + +Stryker's expanded suite includes serialization, condition ordering, issue mapping, regression, +malformed-input, and transition-reference tests. Neither tool's score is a proof that no bugs remain. +Surviving and uncovered mutants remain a coverage backlog; timeouts and equivalent mutants must +be reviewed before treating them as meaningful defects. + +## ELK size and checkpoint cost + +The production build retains a separate ELK chunk of approximately 1,458 kB (445 kB gzip), with +an approximately 803 kB main JavaScript entry (225 kB gzip). Dynamic import removes ELK from the +main chunk, but the application's automatic initial layout still requests it during startup. +Manual layout reuses that loaded module. This does not eliminate the initial graph's layout cost. + +`elkjs` distributes the full generated worker implementation. Selecting only `layered` in the +constructor controls registered algorithms; it does not remove the other algorithms' shipped +bytes. A smaller generated engine or a separately hosted worker would require a packaging and +hosting design change, not a configuration-only size fix. This change retains the existing +layered routing behavior and deployment contract. + +History comparisons and coalescing already avoid unnecessary full-graph serialization. New +checkpoints and undo/redo still clone graph state so externally mutable data cannot corrupt history. +Removing those copies would require immutable ownership guarantees or a patch-based history +design. The regression test covers the isolation that must be preserved. + +Reproduce the local benchmark after building editor-core: + +```bash +cd studio +pnpm --filter @rendis/statepro-studio-react build +node scripts/benchmark-history.mjs +``` + +The fixture has one universe and 100, 1,000, or 5,000 realities, each with a 256-byte metadata payload +and no edges. For each operation, discard one warm-up, take five samples, and report median elapsed +milliseconds for 20 operations; undo/redo measures 20 pairs. These are Node timings on a local +machine, not browser frame timings or a universal performance threshold. + +| Realities | Capture 20 checkpoints (ms) | Record 20 edits (ms) | Coalesce 20 edits (ms) | 20 undo/redo pairs (ms) | +| --- | ---: | ---: | ---: | ---: | +| 100 | 6.05 | 4.89 | 0.43 | 26.62 | +| 1000 | 62.33 | 98.34 | 4.67 | 239.83 | +| 5000 | 313.87 | 286.76 | 11.87 | 1059.54 | + +At 5,000 realities, coalescing makes twenty edits about 24 times cheaper than recording each edit +in this fixture. Undo/redo still averages about 53 ms per pair, exceeding a 16 ms frame budget. +These measurements justify preserving coalescing and identifying clone-heavy large histories as +an architectural optimization opportunity, rather than silently weakening snapshot isolation. + +## Practical limits + +- Callbacks and invokes must cooperate with cancellation. Existing context-free snapshot methods + can still deadlock when called on their owner inside a synchronous callback; use callback args + or the new contextual capability. Replacing the supplied context bypasses scope detection. +- Limit/cancellation errors do not roll back prior side effects, other universes, or already admitted + invokes. `Close()` requests shutdown; `WaitInvokes` needs a deadline for uncooperative user code. +- Accumulator budgets count entries, not payload bytes; history bounds do not bound individual + snapshot size. The bot's legacy provider cannot be forcibly interrupted and runs are not concurrent. +- Strict observers and finite budgets require application configuration. Existing defaults remain + compatible. ELK download size and full checkpoint copy costs remain architectural tradeoffs. +- CLI coverage is broader but not complete. Browser validation used Chromium; no cross-browser + or distributed-service load campaign is claimed. diff --git a/docs/runtime.md b/docs/runtime.md index c2da5f6..431ea6b 100644 --- a/docs/runtime.md +++ b/docs/runtime.md @@ -57,8 +57,9 @@ superposition indefinitely. - **Actions** run synchronously. Any error stops the transition and the machine remains in the previous state. -- **Invokes** run asynchronously on separate goroutines. They are "fire-and-forget" and do not affect - control flow. +- **Invokes** run asynchronously on separate goroutines. Their function does not return an error + to the transition. With configured runtime limits, an admission failure does return an error + before the rejected invoke is started. - Both receive `instrumentation` executor arguments including the machine context, universe metadata, event payload, and snapshot accessors. @@ -114,7 +115,7 @@ Zero changes to the JSON definition. The existing `on.create-form` transition wi ## Conditions & Observers -- Observers run in parallel; the first success wins. Errors are propagated unless another observer has +- Observers run sequentially; the first success wins. Errors are propagated unless another observer has already authorized the transition. - `TransitionModel.condition` and `conditions` arrays are evaluated sequentially. All must return `true` for the transition to proceed. @@ -192,3 +193,71 @@ The experimental runtime implements all instrumentation interfaces. You can buil 3. Re-registering actions/observers/invokes via the `builtin` package or custom registries. Consult [instrumentation.md](instrumentation.md) for the list of contracts you must satisfy. + +## Runtime policies and resource limits + +Existing constructors preserve unlimited numeric budgets, permissive unknown observers, and invokes +that continue after their reality exits. Configure policies explicitly when those defaults do not +fit an application's workload: + +```go +qm, err := statepro.NewQuantumMachineWithOptions(model, instrumentation.RuntimeOptions{ + MaxConcurrentInvokes: 8, + CancelInvokesOnExit: true, + MaxAccumulatedEvents: 1_000, + MaxTrackingEntries: 100, + StrictObservers: true, +}) +if err != nil { + return err +} +``` + +These are example budgets, not recommended values for every application. Options are copied at +construction; negative numeric limits are rejected. The experimental constructor also exposes +`NewExQuantumMachineWithOptions`. The mandatory machine and executor interfaces are unchanged. + +| Policy | Behavior | +| --- | --- | +| `MaxConcurrentInvokes` | One pool for the entire machine, including machine constants, universe constants, and reality/transition invokes. At capacity, reject immediately with `*instrumentation.ResourceLimitError`; do not enqueue or spawn the rejected task. Slots release on completion or panic, not merely on a cancellation request. | +| `CancelInvokesOnExit` | Request cancellation of existing invokes associated with a successfully exited reality, before launching exit invokes. Also cancel affected universes on valid snapshot replacement or static positioning. Rejected snapshot restoration does not cancel existing tasks. | +| `MaxAccumulatedEvents` | Per-universe entry budget, counting separate copies for different realities. Reserve worst-case fan-out before callbacks or appending entries, even if an early observer might approve. Reject over-budget snapshots before applying them. This is an entry limit, not a byte limit. | +| `MaxTrackingEntries` | Retain the most recent entries per universe, including restored tracking. Failed entry rollback preserves the previous bounded history. | +| `StrictObservers` | Reject empty or unregistered observer sources with `ErrUnknownObserver` before accumulating the event. Every configured observer in a fan-out must be registered. The default still permits unknown sources. | + +The debugger bot independently supports `bot.WithHistoryLimit(n)`, retaining the latest event +snapshots. Zero means unlimited; negative values are rejected. Bot runs check cancellation before +restoration, between events, and after the event provider returns, and propagate checked snapshot +errors. The provider's existing signature has no context argument; a provider that blocks must +arrange its own cancellation. Serialize `Run` and history access at the application level. + +### Cancellation, reentry, and shutdown + +Operations with a context can cancel while waiting for the machine lock. Cancellation is checked +between synchronous callbacks and after each returns. A running callback must cooperate with its +context: Go cannot safely interrupt arbitrary user code or undo its external side effects. +Earlier actions, events in other universes, or admitted invokes may already have executed when an +operation returns an error. Resource admission and cancellation do not make execution transactional. + +Calling the owning machine from a synchronous callback using the callback context (or a derived +context) returns `ErrReentrantCall`. This includes nested synchronous callbacks across machines. +Asynchronous invokes receive a context without that synchronous scope marker, preserving caller +values and cancellation; they can enqueue work on the machine normally. Reusing a callback context +after its callback has completed is supported. + +For capture inside an action, use `instrumentation.GetSnapshotWithError(args)` or `args.GetSnapshot()`. +The optional `instrumentation.ContextSnapshotProvider` adds `GetSnapshotContext(ctx)` and +`LoadSnapshotContext(ctx, snapshot, machineContext)`, with cancellable lock waiting and reentry checks. +The legacy context-free machine snapshot methods cannot identify a callback's caller and can still +deadlock if called directly inside its owning synchronous callback. Replacing the callback context +with `context.Background()` also bypasses reentry detection. Preserve the supplied context. + +The optional `instrumentation.RuntimeLifecycle` provides: + +- `Close()`: idempotently request cancellation of all invokes and reject new execution with + `ErrMachineClosed`. It returns without waiting for user code to stop. Snapshot reads remain available. +- `WaitInvokes(ctx)`: wait for the pool to become idle or return the context error. Call `Close()` first + when a stable shutdown barrier is required; concurrent new admissions can otherwise start later. + +Always give shutdown waiting a deadline. An invoke that ignores cancellation occupies its slot +until it exits and may outlive `Close()`. diff --git a/experimental/machine.go b/experimental/machine.go index fb183f1..95a7065 100644 --- a/experimental/machine.go +++ b/experimental/machine.go @@ -5,7 +5,6 @@ import ( "fmt" "log/slog" "reflect" - "sync" "github.com/rendis/statepro/v3/builtin" "github.com/rendis/statepro/v3/instrumentation" @@ -33,12 +32,21 @@ var qmInitFunctions = map[refType]initFunc{ } func NewExQuantumMachine(qmm *theoretical.QuantumMachineModel, universes []*ExUniverse) (instrumentation.QuantumMachine, error) { + return NewExQuantumMachineWithOptions(qmm, universes, instrumentation.RuntimeOptions{}) +} + +func NewExQuantumMachineWithOptions(qmm *theoretical.QuantumMachineModel, universes []*ExUniverse, options instrumentation.RuntimeOptions) (instrumentation.QuantumMachine, error) { + if err := options.Validate(); err != nil { + return nil, err + } if qmm == nil { return nil, fmt.Errorf("quantum machine model must not be nil") } qm := &ExQuantumMachine{ model: qmm, + options: options, + invokes: newInvokeManager(options.MaxConcurrentInvokes), universes: map[string]*ExUniverse{}, } @@ -55,6 +63,9 @@ func NewExQuantumMachine(qmm *theoretical.QuantumMachineModel, universes []*ExUn return nil, fmt.Errorf("universe '%s' already exists", u.model.ID) } + u.options = options + u.invokes = qm.invokes + u.owner = qm u.constantsLawsExecutor = qm u.getSnapshotFn = qm.snapshotUnlocked u.getSnapshotWithErrorFn = qm.snapshotUnlockedWithError @@ -76,7 +87,9 @@ type ExQuantumMachine struct { universes map[string]*ExUniverse // quantumMachineMtx is the mutex for the quantum machine - quantumMachineMtx sync.Mutex + quantumMachineMtx contextMutex + options instrumentation.RuntimeOptions + invokes *invokeManager } //--------- QuantumMachine interface implementation --------- @@ -93,14 +106,10 @@ func (qm *ExQuantumMachine) SendEvent(ctx context.Context, event instrumentation if event == nil || isNilEvent(event) { return false, fmt.Errorf("event must not be nil") } - if err := ctx.Err(); err != nil { + if err := qm.lockContext(ctx); err != nil { return false, err } - qm.quantumMachineMtx.Lock() defer qm.quantumMachineMtx.Unlock() - if err := ctx.Err(); err != nil { - return false, err - } var pairs []util.Pair[instrumentation.Event, []string] @@ -111,6 +120,9 @@ func (qm *ExQuantumMachine) SendEvent(ctx context.Context, event instrumentation } for _, u := range activeUniverses { + if err := ctx.Err(); err != nil { + return true, err + } externalTargets, err := u.handleEvent(ctx, nil, event, qm.machineContext) if err != nil { return true, err @@ -140,6 +152,21 @@ func isNilEvent(event instrumentation.Event) bool { func (qm *ExQuantumMachine) LoadSnapshot(snapshot *instrumentation.MachineSnapshot, machineContext any) error { qm.quantumMachineMtx.Lock() defer qm.quantumMachineMtx.Unlock() + if qm.invokes != nil && qm.invokes.isClosed() { + return instrumentation.ErrMachineClosed + } + return qm.loadSnapshotUnlocked(context.Background(), snapshot, machineContext) +} + +func (qm *ExQuantumMachine) LoadSnapshotContext(ctx context.Context, snapshot *instrumentation.MachineSnapshot, machineContext any) error { + if err := qm.lockContext(ctx); err != nil { + return err + } + defer qm.quantumMachineMtx.Unlock() + return qm.loadSnapshotUnlocked(ctx, snapshot, machineContext) +} + +func (qm *ExQuantumMachine) loadSnapshotUnlocked(ctx context.Context, snapshot *instrumentation.MachineSnapshot, machineContext any) error { if snapshot == nil { return nil @@ -148,6 +175,9 @@ func (qm *ExQuantumMachine) LoadSnapshot(snapshot *instrumentation.MachineSnapsh // Decode and validate every included universe before changing any live state. prepared := make(map[*ExUniverse]*UniverseInfoSnapshot) for _, u := range qm.universes { + if err := ctx.Err(); err != nil { + return err + } universeSnapshot, ok := snapshot.Snapshots[u.model.ID] if !ok { @@ -160,9 +190,15 @@ func (qm *ExQuantumMachine) LoadSnapshot(snapshot *instrumentation.MachineSnapsh } prepared[u] = decoded } + if err := ctx.Err(); err != nil { + return err + } for u, decoded := range prepared { + if qm.options.CancelInvokesOnExit { + qm.invokes.cancel(u.model.ID, "") + } u.applySnapshot(decoded) - u.tracking = cloneStringSlice(snapshot.Tracking[u.model.ID]) + u.tracking = u.retainTracking(cloneStringSlice(snapshot.Tracking[u.model.ID])) } qm.machineContext = machineContext @@ -233,7 +269,9 @@ func (qm *ExQuantumMachine) snapshotUnlockedWithError() (*instrumentation.Machin } func (qm *ExQuantumMachine) ReplayOnEntry(ctx context.Context) error { - qm.quantumMachineMtx.Lock() + if err := qm.lockContext(ctx); err != nil { + return err + } defer qm.quantumMachineMtx.Unlock() var evt = NewEventBuilder("replayOnEntry"). @@ -254,7 +292,9 @@ func (qm *ExQuantumMachine) ReplayOnEntry(ctx context.Context) error { } func (qm *ExQuantumMachine) PositionMachine(ctx context.Context, machineContext any, universeID string, realityID string, executeFlow bool) error { - qm.quantumMachineMtx.Lock() + if err := qm.lockContext(ctx); err != nil { + return err + } defer qm.quantumMachineMtx.Unlock() // Validate parameters @@ -319,7 +359,9 @@ func (qm *ExQuantumMachine) PositionMachineOnInitial(ctx context.Context, machin } // Get target universe (without lock - PositionMachine will handle locking) - qm.quantumMachineMtx.Lock() + if err := qm.lockContext(ctx); err != nil { + return err + } universe, ok := qm.universes[universeID] qm.quantumMachineMtx.Unlock() @@ -344,7 +386,9 @@ func (qm *ExQuantumMachine) PositionMachineByCanonicalName(ctx context.Context, } // Find universe by canonical name - qm.quantumMachineMtx.Lock() + if err := qm.lockContext(ctx); err != nil { + return err + } var universeID string for id, universe := range qm.universes { if universe.model.CanonicalName == universeCanonicalName { @@ -369,7 +413,9 @@ func (qm *ExQuantumMachine) PositionMachineOnInitialByCanonicalName(ctx context. } // Find universe by canonical name - qm.quantumMachineMtx.Lock() + if err := qm.lockContext(ctx); err != nil { + return err + } var universeID string for id, universe := range qm.universes { if universe.model.CanonicalName == universeCanonicalName { @@ -461,7 +507,9 @@ func (qm *ExQuantumMachine) ExecuteTransitionAction(ctx context.Context, args *i //----------------------------------------------------------- func (qm *ExQuantumMachine) init(ctx context.Context, machineContext any, event instrumentation.Event) error { - qm.quantumMachineMtx.Lock() + if err := qm.lockContext(ctx); err != nil { + return err + } defer qm.quantumMachineMtx.Unlock() // guard: prevent double initialization @@ -476,6 +524,9 @@ func (qm *ExQuantumMachine) init(ctx context.Context, machineContext any, event var pairs []util.Pair[instrumentation.Event, []string] for _, ref := range qm.model.Initials { + if err := ctx.Err(); err != nil { + return err + } // get reference type and parts refT, parts, err := processReference(ref) if err != nil { @@ -528,20 +579,7 @@ func (qm *ExQuantumMachine) executeInvoke(ctx context.Context, invoke theoretica invoke: invoke, } - if fn := builtin.GetInvoke(invoke.Src); fn != nil { - src := invoke.Src - go func() { - defer func() { - if r := recover(); r != nil { - slog.ErrorContext(ctx, "invoke panicked", "src", src, "panic", r) - } - }() - fn(ctx, a) - }() - return - } - - slog.WarnContext(ctx, "invoke not found", "src", invoke.Src) + u.runInvokeExecutor(ctx, a) } func (qm *ExQuantumMachine) executeAction(ctx context.Context, model *theoretical.ActionModel, args *instrumentation.QuantumMachineExecutorArgs, actionType instrumentation.ActionType) error { @@ -570,7 +608,16 @@ func (qm *ExQuantumMachine) executeAction(ctx context.Context, model *theoretica } if fn := builtin.GetAction(model.Src); fn != nil { - return fn(ctx, a) + if err := ctx.Err(); err != nil { + return err + } + callbackCtx, done := qm.callbackContext(ctx) + defer done() + err := fn(callbackCtx, a) + if err != nil { + return err + } + return ctx.Err() } slog.WarnContext(ctx, "action not found", "src", model.Src) @@ -615,6 +662,9 @@ func (qm *ExQuantumMachine) executeExternalTargetPairs(ctx context.Context, pair var jobs []cascadeJob for _, pair := range pairs { + if err := ctx.Err(); err != nil { + return err + } evt, targets := pair.GetAll() if len(targets) == 0 { continue @@ -653,6 +703,9 @@ func (qm *ExQuantumMachine) executeTransitions(ctx context.Context, event instru var newTargets []string for _, target := range targets { + if err := ctx.Err(); err != nil { + return nil, err + } refT, parts, err := processReference(target) if err != nil { return nil, err diff --git a/experimental/runtime_control.go b/experimental/runtime_control.go new file mode 100644 index 0000000..ac54a0b --- /dev/null +++ b/experimental/runtime_control.go @@ -0,0 +1,199 @@ +package experimental + +import ( + "context" + "fmt" + "log/slog" + "sync" + "sync/atomic" + + "github.com/rendis/statepro/v3/instrumentation" +) + +// contextMutex preserves synchronous locking for context-free legacy methods, +// while allowing execution callers to cancel their wait without spawning goroutines. +type contextMutex struct { + once sync.Once + token chan struct{} +} + +func (m *contextMutex) init() { + m.once.Do(func() { m.token = make(chan struct{}, 1); m.token <- struct{}{} }) +} +func (m *contextMutex) Lock() { m.init(); <-m.token } +func (m *contextMutex) Unlock() { m.token <- struct{}{} } +func (m *contextMutex) LockContext(ctx context.Context) error { + if ctx == nil { + return fmt.Errorf("context must not be nil") + } + if err := ctx.Err(); err != nil { + return err + } + m.init() + select { + case <-ctx.Done(): + return ctx.Err() + case <-m.token: + if err := ctx.Err(); err != nil { + m.Unlock() + return err + } + return nil + } +} + +type callbackKey struct{} +type callbackScope struct { + owner *ExQuantumMachine + active atomic.Bool + parent *callbackScope +} + +func (qm *ExQuantumMachine) callbackContext(ctx context.Context) (context.Context, func()) { + parent, _ := ctx.Value(callbackKey{}).(*callbackScope) + scope := &callbackScope{owner: qm, parent: parent} + scope.active.Store(true) + return context.WithValue(ctx, callbackKey{}, scope), func() { scope.active.Store(false) } +} + +func (qm *ExQuantumMachine) lockContext(ctx context.Context) error { + if err := qm.checkCallbackContext(ctx); err != nil { + return err + } + if err := qm.quantumMachineMtx.LockContext(ctx); err != nil { + return err + } + if qm.invokes != nil && qm.invokes.isClosed() { + qm.quantumMachineMtx.Unlock() + return instrumentation.ErrMachineClosed + } + return nil +} + +func (qm *ExQuantumMachine) checkCallbackContext(ctx context.Context) error { + if ctx == nil { + return fmt.Errorf("context must not be nil") + } + for scope, _ := ctx.Value(callbackKey{}).(*callbackScope); scope != nil; scope = scope.parent { + if scope.owner == qm && scope.active.Load() { + return instrumentation.ErrReentrantCall + } + } + return nil +} + +type invokeTask struct { + universe, reality string + cancel context.CancelFunc +} +type invokeManager struct { + mu sync.Mutex + limit int + closed bool + next uint64 + tasks map[uint64]invokeTask + idle chan struct{} +} + +func newInvokeManager(limit int) *invokeManager { + idle := make(chan struct{}) + close(idle) + return &invokeManager{limit: limit, tasks: make(map[uint64]invokeTask), idle: idle} +} + +func (m *invokeManager) start(ctx context.Context, universe, reality, src string, run func(context.Context)) error { + if err := ctx.Err(); err != nil { + return err + } + m.mu.Lock() + if m.closed { + m.mu.Unlock() + return instrumentation.ErrMachineClosed + } + if m.limit > 0 && len(m.tasks) >= m.limit { + m.mu.Unlock() + return &instrumentation.ResourceLimitError{Resource: "concurrent invokes", Limit: m.limit} + } + if len(m.tasks) == 0 { + m.idle = make(chan struct{}) + } + m.next++ + id := m.next + // Async invokes run outside the synchronous callback's lock scope. Preserve + // cancellation and caller values without inheriting its reentry marker. + taskCtx, cancel := context.WithCancel(context.WithValue(ctx, callbackKey{}, (*callbackScope)(nil))) + m.tasks[id] = invokeTask{universe: universe, reality: reality, cancel: cancel} + m.mu.Unlock() + go func() { + defer func() { + if r := recover(); r != nil { + slog.ErrorContext(taskCtx, "invoke panicked", "src", src, "panic", r) + } + cancel() + m.mu.Lock() + delete(m.tasks, id) + if len(m.tasks) == 0 { + close(m.idle) + } + m.mu.Unlock() + }() + if taskCtx.Err() == nil { + run(taskCtx) + } + }() + return nil +} + +func (m *invokeManager) cancel(universe, reality string) { + m.mu.Lock() + defer m.mu.Unlock() + for _, task := range m.tasks { + if task.universe == universe && (reality == "" || task.reality == reality) { + task.cancel() + } + } +} +func (m *invokeManager) isClosed() bool { m.mu.Lock(); defer m.mu.Unlock(); return m.closed } +func (m *invokeManager) close() { + m.mu.Lock() + defer m.mu.Unlock() + m.closed = true + for _, task := range m.tasks { + task.cancel() + } +} +func (m *invokeManager) wait(ctx context.Context) error { + if ctx == nil { + return fmt.Errorf("context must not be nil") + } + if err := ctx.Err(); err != nil { + return err + } + m.mu.Lock() + idle := m.idle + m.mu.Unlock() + select { + case <-idle: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + +func (qm *ExQuantumMachine) Close() error { qm.invokes.close(); return nil } +func (qm *ExQuantumMachine) WaitInvokes(ctx context.Context) error { + if err := qm.checkCallbackContext(ctx); err != nil { + return err + } + return qm.invokes.wait(ctx) +} +func (qm *ExQuantumMachine) GetSnapshotContext(ctx context.Context) (*instrumentation.MachineSnapshot, error) { + if err := qm.checkCallbackContext(ctx); err != nil { + return nil, err + } + if err := qm.quantumMachineMtx.LockContext(ctx); err != nil { + return nil, err + } + defer qm.quantumMachineMtx.Unlock() + return qm.snapshotUnlockedWithError() +} diff --git a/experimental/runtime_control_test.go b/experimental/runtime_control_test.go new file mode 100644 index 0000000..480d3ca --- /dev/null +++ b/experimental/runtime_control_test.go @@ -0,0 +1,399 @@ +package experimental + +import ( + "context" + "errors" + "reflect" + "sync/atomic" + "testing" + "time" + + "github.com/rendis/statepro/v3/instrumentation" + "github.com/rendis/statepro/v3/theoretical" +) + +func buildControlledQM(t *testing.T, options instrumentation.RuntimeOptions, realities map[string]*theoretical.RealityModel) (*ExQuantumMachine, *ExUniverse) { + t.Helper() + base, u := buildQM(t, "stateA", realities) + machine, err := NewExQuantumMachineWithOptions(base.model, []*ExUniverse{u}, options) + if err != nil { + t.Fatal(err) + } + qm := machine.(*ExQuantumMachine) + t.Cleanup(func() { + _ = qm.Close() + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + if err := qm.WaitInvokes(ctx); err != nil { + t.Errorf("invokes pendientes: %v", err) + } + }) + return qm, u +} + +func awaitRuntimeSignal(t *testing.T, signal <-chan struct{}) { + t.Helper() + select { + case <-signal: + case <-time.After(time.Second): + t.Fatal("operacion bloqueada") + } +} + +func TestRuntime_CancelaEsperaDelLockSinEsperarAlCallback(t *testing.T) { + entered, release := make(chan struct{}), make(chan struct{}) + registerTestAction(t, "test:runtime-lock", func(context.Context, instrumentation.ActionExecutorArgs) error { close(entered); <-release; return nil }) + qm, _ := buildControlledQM(t, instrumentation.RuntimeOptions{}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryAction("test:runtime-lock"))}) + done := make(chan error, 1) + go func() { done <- qm.Init(context.Background(), nil) }() + awaitRuntimeSignal(t, entered) + ctx, cancel := context.WithCancel(context.Background()) + waiting := make(chan error, 1) + go func() { _, err := qm.SendEvent(ctx, NewEventBuilder("GO").Build()); waiting <- err }() + cancel() + select { + case err := <-waiting: + if !errors.Is(err, context.Canceled) { + t.Fatal(err) + } + case <-time.After(time.Second): + close(release) + t.Fatal("cancelacion espera el lock") + } + close(release) + if err := <-done; err != nil { + t.Fatal(err) + } + if _, err := qm.GetSnapshotContext(context.Background()); err != nil { + t.Fatal(err) + } +} + +func TestRuntime_ReentradaDetectadaYContextoReutilizableTrasCallback(t *testing.T) { + var qm *ExQuantumMachine + var saved context.Context + registerTestAction(t, "test:runtime-reentry", func(ctx context.Context, args instrumentation.ActionExecutorArgs) error { + saved = ctx + if _, err := qm.SendEvent(ctx, NewEventBuilder("GO").Build()); !errors.Is(err, instrumentation.ErrReentrantCall) { + return errors.New("reentrada no rechazada") + } + if _, err := qm.GetSnapshotContext(ctx); !errors.Is(err, instrumentation.ErrReentrantCall) { + return errors.New("snapshot reentrante no rechazado") + } + if err := qm.LoadSnapshotContext(ctx, nil, nil); !errors.Is(err, instrumentation.ErrReentrantCall) { + return errors.New("restauracion reentrante no rechazada") + } + if err := qm.WaitInvokes(ctx); !errors.Is(err, instrumentation.ErrReentrantCall) { + return errors.New("espera reentrante no rechazada") + } + if _, err := instrumentation.GetSnapshotWithError(args); err != nil { + return err + } + return nil + }) + qm, _ = buildControlledQM(t, instrumentation.RuntimeOptions{}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryAction("test:runtime-reentry"))}) + done := make(chan error, 1) + go func() { done <- qm.Init(context.Background(), nil) }() + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(time.Second): + t.Fatal("reentrada bloqueada") + } + if _, err := qm.GetSnapshotContext(saved); err != nil { + t.Fatal(err) + } +} + +func TestRuntime_CancelacionDuranteAccionEvitaSiguienteCallback(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + var second atomic.Bool + registerTestAction(t, "test:runtime-cancel-action", func(context.Context, instrumentation.ActionExecutorArgs) error { cancel(); return nil }) + registerTestAction(t, "test:runtime-second-action", func(context.Context, instrumentation.ActionExecutorArgs) error { second.Store(true); return nil }) + qm, u := buildControlledQM(t, instrumentation.RuntimeOptions{}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryAction("test:runtime-cancel-action"), withEntryAction("test:runtime-second-action"))}) + if err := qm.Init(ctx, nil); !errors.Is(err, context.Canceled) { + t.Fatal(err) + } + if second.Load() || u.initialized { + t.Fatal("cancelacion deja ejecutar otro callback") + } +} + +func TestRuntime_ReentradaEntreCallbacksDeDosMaquinas(t *testing.T) { + var first, second *ExQuantumMachine + registerTestAction(t, "test:runtime-nested-first", func(ctx context.Context, _ instrumentation.ActionExecutorArgs) error { return second.Init(ctx, nil) }) + registerTestAction(t, "test:runtime-nested-second", func(ctx context.Context, _ instrumentation.ActionExecutorArgs) error { + _, err := first.SendEvent(ctx, NewEventBuilder("GO").Build()) + if !errors.Is(err, instrumentation.ErrReentrantCall) { + return errors.New("reentrada anidada no rechazada") + } + return nil + }) + first, _ = buildControlledQM(t, instrumentation.RuntimeOptions{}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryAction("test:runtime-nested-first"))}) + second, _ = buildControlledQM(t, instrumentation.RuntimeOptions{}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryAction("test:runtime-nested-second"))}) + done := make(chan error, 1) + go func() { done <- first.Init(context.Background(), nil) }() + select { + case err := <-done: + if err != nil { + t.Fatal(err) + } + case <-time.After(time.Second): + t.Fatal("reentrada anidada bloqueada") + } +} + +func TestRuntime_InvokeAsincronoNoHeredaReentradaDelCallback(t *testing.T) { + var first, second *ExQuantumMachine + started := make(chan struct{}) + result := make(chan error, 1) + registerTestAction(t, "test:runtime-async-parent", func(ctx context.Context, _ instrumentation.ActionExecutorArgs) error { + if err := second.Init(ctx, nil); err != nil { + return err + } + awaitRuntimeSignal(t, started) + return nil + }) + registerTestInvoke(t, "test:runtime-async-child", func(ctx context.Context, _ instrumentation.InvokeExecutorArgs) { + err := first.checkCallbackContext(ctx) + close(started) + if err == nil { + _, err = first.SendEvent(ctx, NewEventBuilder("GO").Build()) + } + result <- err + }) + first, _ = buildControlledQM(t, instrumentation.RuntimeOptions{}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryAction("test:runtime-async-parent"))}) + second, _ = buildControlledQM(t, instrumentation.RuntimeOptions{}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryInvoke("test:runtime-async-child"))}) + if err := first.Init(context.Background(), nil); err != nil { + t.Fatal(err) + } + select { + case err := <-result: + if err != nil { + t.Fatal(err) + } + case <-time.After(time.Second): + t.Fatal("invoke asincrono bloqueado") + } +} + +func TestRuntime_LimiteCompartidoEntreUniversos(t *testing.T) { + started := make(chan struct{}) + var rejected atomic.Bool + registerTestInvoke(t, "test:runtime-universe-one", func(ctx context.Context, _ instrumentation.InvokeExecutorArgs) { close(started); <-ctx.Done() }) + registerTestInvoke(t, "test:runtime-universe-two", func(context.Context, instrumentation.InvokeExecutorArgs) { rejected.Store(true) }) + base, first, second := buildMultiUniverseQM(t, + "stateA", map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryInvoke("test:runtime-universe-one"))}, + "stateA", map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryInvoke("test:runtime-universe-two"))}, + ) + configured, err := NewExQuantumMachineWithOptions(base.model, []*ExUniverse{first, second}, instrumentation.RuntimeOptions{MaxConcurrentInvokes: 1}) + if err != nil { + t.Fatal(err) + } + qm := configured.(*ExQuantumMachine) + defer qm.Close() + var limit *instrumentation.ResourceLimitError + if err := qm.Init(context.Background(), nil); !errors.As(err, &limit) { + t.Fatalf("limite por maquina omitido: %v", err) + } + awaitRuntimeSignal(t, started) + if rejected.Load() { + t.Fatal("segundo universo excede limite") + } +} + +func TestRuntime_PanicLiberaSlotYPosicionCancelaInvokes(t *testing.T) { + registerTestInvoke(t, "test:runtime-pool-panic", func(context.Context, instrumentation.InvokeExecutorArgs) { panic("intentional pool panic") }) + qm, _ := buildControlledQM(t, instrumentation.RuntimeOptions{MaxConcurrentInvokes: 1}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryInvoke("test:runtime-pool-panic"))}) + if err := qm.Init(context.Background(), nil); err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + if err := qm.WaitInvokes(ctx); err != nil { + t.Fatal(err) + } + if err := qm.ReplayOnEntry(ctx); err != nil { + t.Fatal(err) + } + if err := qm.WaitInvokes(ctx); err != nil { + t.Fatal(err) + } + started, ended := make(chan struct{}), make(chan struct{}) + registerTestInvoke(t, "test:runtime-position-cancel", func(ctx context.Context, _ instrumentation.InvokeExecutorArgs) { + close(started) + <-ctx.Done() + close(ended) + }) + positioned, _ := buildControlledQM(t, instrumentation.RuntimeOptions{CancelInvokesOnExit: true}, map[string]*theoretical.RealityModel{ + "stateA": newTransitionReality("stateA", withEntryInvoke("test:runtime-position-cancel")), "stateB": newFinalReality("stateB"), + }) + if err := positioned.Init(ctx, nil); err != nil { + t.Fatal(err) + } + awaitRuntimeSignal(t, started) + if err := positioned.PositionMachine(ctx, nil, "u1", "stateB", false); err != nil { + t.Fatal(err) + } + awaitRuntimeSignal(t, ended) +} + +func TestRuntime_InvokesRespetanSalidaOptInYCierre(t *testing.T) { + for _, cancelOnExit := range []bool{false, true} { + t.Run(map[bool]string{false: "contrato_legacy", true: "cancelacion_opt_in"}[cancelOnExit], func(t *testing.T) { + started, ended := make(chan struct{}), make(chan struct{}) + name := "test:runtime-exit-legacy" + if cancelOnExit { + name = "test:runtime-exit-controlled" + } + registerTestInvoke(t, name, func(ctx context.Context, _ instrumentation.InvokeExecutorArgs) { + close(started) + <-ctx.Done() + close(ended) + }) + qm, _ := buildControlledQM(t, instrumentation.RuntimeOptions{CancelInvokesOnExit: cancelOnExit}, map[string]*theoretical.RealityModel{ + "stateA": newTransitionReality("stateA", withEntryInvoke(name), withOnTransition("GO", []string{"stateB"}, nil)), "stateB": newFinalReality("stateB"), + }) + if err := qm.Init(context.Background(), nil); err != nil { + t.Fatal(err) + } + awaitRuntimeSignal(t, started) + if _, err := qm.SendEvent(context.Background(), NewEventBuilder("GO").Build()); err != nil { + t.Fatal(err) + } + if cancelOnExit { + awaitRuntimeSignal(t, ended) + } else { + select { + case <-ended: + t.Fatal("default cancela invoke al salir") + default: + } + } + _ = qm.Close() + awaitRuntimeSignal(t, ended) + if err := qm.ReplayOnEntry(context.Background()); !errors.Is(err, instrumentation.ErrMachineClosed) { + t.Fatal(err) + } + }) + } +} + +func TestRuntime_LimiteGlobalCubreInvokesDeConstantesYRealidad(t *testing.T) { + started := make(chan struct{}) + var rejected atomic.Bool + registerTestInvoke(t, "test:runtime-limited-constant", func(ctx context.Context, _ instrumentation.InvokeExecutorArgs) { close(started); <-ctx.Done() }) + registerTestInvoke(t, "test:runtime-limited-reality", func(context.Context, instrumentation.InvokeExecutorArgs) { rejected.Store(true) }) + qm, _ := buildControlledQM(t, instrumentation.RuntimeOptions{MaxConcurrentInvokes: 1}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryInvoke("test:runtime-limited-reality"))}) + qm.model.UniversalConstants = &theoretical.UniversalConstantsModel{EntryInvokes: []*theoretical.InvokeModel{{Src: "test:runtime-limited-constant"}}} + var limit *instrumentation.ResourceLimitError + if err := qm.Init(context.Background(), nil); !errors.As(err, &limit) || limit.Resource != "concurrent invokes" { + t.Fatalf("limite omitido: %v", err) + } + awaitRuntimeSignal(t, started) + if rejected.Load() { + t.Fatal("invoke rechazado fue lanzado") + } + ctx, cancel := context.WithTimeout(context.Background(), time.Millisecond) + defer cancel() + if err := qm.WaitInvokes(ctx); !errors.Is(err, context.DeadlineExceeded) { + t.Fatal(err) + } +} + +func TestRuntime_RestauracionValidaCancelaInvokesPeroInvalidaNo(t *testing.T) { + started, ended := make(chan struct{}), make(chan struct{}) + registerTestInvoke(t, "test:runtime-restore", func(ctx context.Context, _ instrumentation.InvokeExecutorArgs) { + close(started) + <-ctx.Done() + close(ended) + }) + qm, _ := buildControlledQM(t, instrumentation.RuntimeOptions{CancelInvokesOnExit: true}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withEntryInvoke("test:runtime-restore"))}) + if err := qm.Init(context.Background(), nil); err != nil { + t.Fatal(err) + } + awaitRuntimeSignal(t, started) + snapshot, err := qm.GetSnapshotWithError() + if err != nil { + t.Fatal(err) + } + if err := qm.LoadSnapshot(&instrumentation.MachineSnapshot{Snapshots: map[string]instrumentation.SerializedUniverseSnapshot{"u1": nil}}, nil); err == nil { + t.Fatal("snapshot invalido aceptado") + } + select { + case <-ended: + t.Fatal("rechazar snapshot cancela invokes") + default: + } + if err := qm.LoadSnapshot(snapshot, nil); err != nil { + t.Fatal(err) + } + awaitRuntimeSignal(t, ended) +} + +func TestRuntime_ObserverEstrictoYLimiteDeAcumulacionSinMutacionParcial(t *testing.T) { + registerTestObserver(t, "test:runtime-decline", func(context.Context, instrumentation.ObserverExecutorArgs) (bool, error) { return false, nil }) + for _, strict := range []bool{false, true} { + qm, u := buildControlledQM(t, instrumentation.RuntimeOptions{StrictObservers: strict}, map[string]*theoretical.RealityModel{"stateA": newTransitionReality("stateA", withObserver("test:runtime-missing", nil))}) + u.initOnSuperposition() + _, err := qm.SendEvent(context.Background(), NewEventBuilder("GO").Build()) + if strict && (!errors.Is(err, instrumentation.ErrUnknownObserver) || !u.inSuperposition || u.eventAccumulator.GetStatistics().CountAllEvents() != 0) { + t.Fatalf("modo estricto no rechazo sin acumular: %v", err) + } + if !strict && (err != nil || u.inSuperposition) { + t.Fatalf("default alterado: %v", err) + } + } + qm, u := buildControlledQM(t, instrumentation.RuntimeOptions{MaxAccumulatedEvents: 3}, map[string]*theoretical.RealityModel{ + "stateA": newTransitionReality("stateA", withObserver("test:runtime-decline", nil)), "stateB": newTransitionReality("stateB", withObserver("test:runtime-decline", nil)), + }) + u.initOnSuperposition() + if _, err := qm.SendEvent(context.Background(), NewEventBuilder("GO").Build()); err != nil { + t.Fatal(err) + } + var limit *instrumentation.ResourceLimitError + if _, err := qm.SendEvent(context.Background(), NewEventBuilder("GO").Build()); !errors.As(err, &limit) { + t.Fatal(err) + } + if u.eventAccumulator.GetStatistics().CountAllEvents() != 2 { + t.Fatal("fanout rechazado muta acumulador") + } + snapshot, err := qm.GetSnapshotWithError() + if err != nil { + t.Fatal(err) + } + qm.options.MaxAccumulatedEvents = 1 + u.options.MaxAccumulatedEvents = 1 + if err := qm.LoadSnapshot(snapshot, nil); !errors.As(err, &limit) { + t.Fatalf("restauracion evade limite: %v", err) + } +} + +func TestRuntime_TrackingLimitadoPreservaRollback(t *testing.T) { + failure := errors.New("entry failed") + registerTestAction(t, "test:runtime-tracking-fail", func(context.Context, instrumentation.ActionExecutorArgs) error { return failure }) + qm, u := buildControlledQM(t, instrumentation.RuntimeOptions{MaxTrackingEntries: 2}, map[string]*theoretical.RealityModel{ + "stateA": newTransitionReality("stateA", withOnTransition("GO", []string{"stateB"}, nil)), "stateB": newTransitionReality("stateB", withEntryAction("test:runtime-tracking-fail")), + }) + if err := qm.Init(context.Background(), nil); err != nil { + t.Fatal(err) + } + u.tracking = []string{"older", "stateA"} + if _, err := qm.SendEvent(context.Background(), NewEventBuilder("GO").Build()); !errors.Is(err, failure) { + t.Fatal(err) + } + if !reflect.DeepEqual(u.tracking, []string{"older", "stateA"}) { + t.Fatalf("rollback pierde historial: %v", u.tracking) + } + for i := 0; i < 5; i++ { + if err := qm.PositionMachine(context.Background(), nil, "u1", "stateA", false); err != nil { + t.Fatal(err) + } + } + if !reflect.DeepEqual(u.tracking, []string{"stateA", "stateA"}) { + t.Fatalf("historial sin limite: %v", u.tracking) + } +} diff --git a/experimental/universe.go b/experimental/universe.go index 61532af..7f50314 100644 --- a/experimental/universe.go +++ b/experimental/universe.go @@ -106,6 +106,10 @@ type ExUniverse struct { // getSnapshotFn returns a snapshot without taking the machine mutex. // Used from actions that already run under quantumMachineMtx. + options instrumentation.RuntimeOptions + invokes *invokeManager + owner *ExQuantumMachine + invokeError error // Written only by synchronous admission under the machine lock. getSnapshotFn func() *instrumentation.MachineSnapshot getSnapshotWithErrorFn func() (*instrumentation.MachineSnapshot, error) } @@ -264,6 +268,9 @@ func (u *ExUniverse) decodeSnapshot(universeSnapshot instrumentation.SerializedU } } if snapshot.Accumulator != nil { + if limit := u.options.MaxAccumulatedEvents; limit > 0 && snapshot.Accumulator.CountAllEvents() > limit { + return nil, &instrumentation.ResourceLimitError{Resource: "accumulated events", Limit: limit} + } for reality, events := range snapshot.Accumulator.RealitiesEvents { model, err := u.getRealityModel(reality) if err != nil { @@ -359,6 +366,9 @@ func (u *ExUniverse) positionStatic(realityID string, universeContext any) error u.initialized = true // Set current reality directly + if u.options.CancelInvokesOnExit && u.invokes != nil { + u.invokes.cancel(u.model.ID, "") + } u.currentReality = &realityID u.addStateToTracking(u.currentReality) @@ -389,7 +399,15 @@ func (u *ExUniverse) addStateToTracking(state *string) { if state == nil { return } - u.tracking = append(u.tracking, *state) + u.tracking = u.retainTracking(append(u.tracking, *state)) +} + +func (u *ExUniverse) retainTracking(entries []string) []string { + limit := u.options.MaxTrackingEntries + if limit > 0 && len(entries) > limit { + return append([]string(nil), entries[len(entries)-limit:]...) + } + return entries } func (u *ExUniverse) popTrackingIfLast(state string) { @@ -554,6 +572,7 @@ func (u *ExUniverse) initializeUniverseOn(ctx context.Context, realityName strin } func (u *ExUniverse) establishNewReality(ctx context.Context, reality string, event instrumentation.Event) error { + previousTracking := u.tracking previousReality := u.currentReality previousFinal := u.isFinalReality previousRealityInitialized := u.realityInitialized @@ -568,7 +587,7 @@ func (u *ExUniverse) establishNewReality(ctx context.Context, reality string, ev u.realityInitialized = previousRealityInitialized u.inSuperposition = previousSuperposition u.realityBeforeSuperposition = previousBeforeSuperposition - u.popTrackingIfLast(reality) + u.tracking = previousTracking return errors.Join(fmt.Errorf(errorExecutingOnEntryProcessMsgTemplate, u.model.ID, reality), err) } u.realityInitialized = true @@ -634,6 +653,9 @@ func (u *ExUniverse) doCyclicTransition( visitedTargets := map[string]int{} for { + if err := ctx.Err(); err != nil { + return err + } if approvedTransition == nil || len(approvedTransition.Targets) == 0 { return nil } @@ -658,9 +680,9 @@ func (u *ExUniverse) doCyclicTransition( return errors.Join(fmt.Errorf("error executing transition actions for reality '%s'", *u.currentReality), err) } - u.constantsLawsExecutor.ExecuteTransitionInvokes(ctx, &args) - u.executeUniverseConstantInvokes(ctx, "transition", event) - u.executeInvokes(ctx, approvedTransition.Invokes, event) + if err := u.executeInvokeGroup(ctx, "transition", approvedTransition.Invokes, event, &args); err != nil { + return err + } if approvedTransition.IsNotification() { u.externalTargets = approvedTransition.Targets @@ -692,6 +714,7 @@ func (u *ExUniverse) doCyclicTransition( ) } + previousTracking := u.tracking previousReality := u.currentReality previousFinal := u.isFinalReality @@ -702,7 +725,7 @@ func (u *ExUniverse) doCyclicTransition( u.currentReality = previousReality u.isFinalReality = previousFinal u.realityInitialized = previousReality != nil - u.popTrackingIfLast(next) + u.tracking = previousTracking return errors.Join(fmt.Errorf(errorExecutingOnEntryProcessMsgTemplate, u.model.ID, next), err) } @@ -730,6 +753,9 @@ func (u *ExUniverse) getApprovedTransition( } for _, transition := range transitionModels { + if err := ctx.Err(); err != nil { + return nil, err + } var conditions []*theoretical.ConditionModel if transition.Condition != nil { @@ -866,12 +892,9 @@ func (u *ExUniverse) executeOnEntryProcess(ctx context.Context, event instrument } // execute on entry constants invokes, invokes are executed asynchronously - u.constantsLawsExecutor.ExecuteEntryInvokes(ctx, args) - - u.executeUniverseConstantInvokes(ctx, "entry", event) - - // execute on entry reality invokes, invokes are executed asynchronously - u.executeInvokes(ctx, realityModel.EntryInvokes, event) + if err := u.executeInvokeGroup(ctx, "entry", realityModel.EntryInvokes, event, args); err != nil { + return err + } u.realityInitialized = true @@ -924,6 +947,9 @@ func (u *ExUniverse) processEmittedEvents( } for _, emitted := range emittedEvents { + if err := ctx.Err(); err != nil { + return err + } transitions, ok := realityModel.On[emitted.Name] if !ok { continue @@ -994,9 +1020,12 @@ func (u *ExUniverse) executeOnExitProcess(ctx context.Context, event instrumenta ) } - u.constantsLawsExecutor.ExecuteExitInvokes(ctx, args) - u.executeUniverseConstantInvokes(ctx, "exit", event) - u.executeInvokes(ctx, realityModel.ExitInvokes, event) + if u.options.CancelInvokesOnExit && u.invokes != nil { + u.invokes.cancel(u.model.ID, realityModel.ID) + } + if err := u.executeInvokeGroup(ctx, "exit", realityModel.ExitInvokes, event, args); err != nil { + return err + } u.realityInitialized = false return nil @@ -1015,6 +1044,9 @@ func (u *ExUniverse) executeActions( // execute actions for _, action := range actionModels { + if err := ctx.Err(); err != nil { + return err + } args := &actionExecutorArgs{ context: u.universeContext, realityName: *u.currentReality, @@ -1076,6 +1108,12 @@ func (u *ExUniverse) accumulateEventForReality( return false, nil } + if err := u.validateObservers(realityModel); err != nil { + return false, err + } + if err := u.checkAccumulatorBudget(1); err != nil { + return false, err + } // accumulate Event u.eventAccumulator.Accumulate(realityName, event) @@ -1089,6 +1127,19 @@ func (u *ExUniverse) accumulateEventForReality( } func (u *ExUniverse) accumulateEventForAllRealities(ctx context.Context, event instrumentation.Event) (bool, string, error) { + // Reserve the worst-case fan-out before invoking observers or adding any copies. + copies := 0 + for _, model := range u.model.Realities { + if err := u.validateObservers(model); err != nil { + return false, "", err + } + if len(model.Observers) > 0 { + copies++ + } + } + if err := u.checkAccumulatorBudget(copies); err != nil { + return false, "", err + } for _, reality := range sortedMapKeys(u.model.Realities) { isNewReality, err := u.accumulateEventForReality(ctx, reality, event, false) if err != nil { @@ -1109,6 +1160,9 @@ func (u *ExUniverse) executeObservers( var firstErr error for _, observer := range realityModel.Observers { + if err := ctx.Err(); err != nil { + return false, err + } args := &observerExecutorArgs{ context: u.universeContext, realityName: realityModel.ID, @@ -1154,12 +1208,27 @@ func (u *ExUniverse) canRealityHandleEvent(realityName string, evt instrumentati //------------- executors ------------- func (u *ExUniverse) runObserverExecutor(ctx context.Context, src string, args *observerExecutorArgs) (bool, error) { + if u.options.StrictObservers && (src == "" || builtin.GetObserver(src) == nil) { + return false, fmt.Errorf("%w: %q", instrumentation.ErrUnknownObserver, src) + } if src == "" { return true, nil } + if err := ctx.Err(); err != nil { + return false, err + } if fn := builtin.GetObserver(src); fn != nil { - return fn(ctx, args) + callbackCtx, done := u.owner.callbackContext(ctx) + defer done() + approved, err := fn(callbackCtx, args) + if err != nil { + return false, err + } + if err := ctx.Err(); err != nil { + return false, err + } + return approved, nil } slog.WarnContext(ctx, "observer not found (default: return true)", "src", src) @@ -1171,8 +1240,17 @@ func (u *ExUniverse) runActionExecutor(ctx context.Context, src string, args *ac return nil } + if err := ctx.Err(); err != nil { + return err + } if fn := builtin.GetAction(src); fn != nil { - return fn(ctx, args) + callbackCtx, done := u.owner.callbackContext(ctx) + defer done() + err := fn(callbackCtx, args) + if err != nil { + return err + } + return ctx.Err() } slog.WarnContext(ctx, "action not found", "src", src) @@ -1180,33 +1258,82 @@ func (u *ExUniverse) runActionExecutor(ctx context.Context, src string, args *ac } func (u *ExUniverse) runInvokeExecutor(ctx context.Context, args *invokeExecutorArgs) { - if args.invoke.Src == "" { + if u.invokeError != nil || args.invoke.Src == "" { return } - if fn := builtin.GetInvoke(args.invoke.Src); fn != nil { - src := args.invoke.Src - go func() { - defer func() { - if r := recover(); r != nil { - slog.ErrorContext(ctx, "invoke panicked", "src", src, "panic", r) - } - }() - fn(ctx, args) - }() + if u.invokes == nil { + u.invokes = newInvokeManager(0) + } + u.invokeError = u.invokes.start(ctx, u.model.ID, args.realityName, args.invoke.Src, func(ctx context.Context) { fn(ctx, args) }) + if u.invokeError != nil { + slog.ErrorContext(ctx, "invoke admission failed", "src", args.invoke.Src, "error", u.invokeError) + } return } - slog.WarnContext(ctx, "invoke not found", "src", args.invoke.Src) } +func (u *ExUniverse) executeInvokeGroup(ctx context.Context, phase string, invokes []*theoretical.InvokeModel, event instrumentation.Event, args *instrumentation.QuantumMachineExecutorArgs) error { + if err := ctx.Err(); err != nil { + return err + } + u.invokeError = nil + switch phase { + case "entry": + u.constantsLawsExecutor.ExecuteEntryInvokes(ctx, args) + case "exit": + u.constantsLawsExecutor.ExecuteExitInvokes(ctx, args) + case "transition": + u.constantsLawsExecutor.ExecuteTransitionInvokes(ctx, args) + } + u.executeUniverseConstantInvokes(ctx, phase, event) + u.executeInvokes(ctx, invokes, event) + return u.invokeError +} + +func (u *ExUniverse) validateObservers(model *theoretical.RealityModel) error { + if !u.options.StrictObservers { + return nil + } + for _, observer := range model.Observers { + if observer == nil || observer.Src == "" || builtin.GetObserver(observer.Src) == nil { + return fmt.Errorf("%w in universe %q, reality %q", instrumentation.ErrUnknownObserver, u.model.ID, model.ID) + } + } + return nil +} + +func (u *ExUniverse) checkAccumulatorBudget(copies int) error { + limit := u.options.MaxAccumulatedEvents + if limit == 0 || u.eventAccumulator == nil { + return nil + } + if copies > limit-u.eventAccumulator.GetStatistics().CountAllEvents() { + return &instrumentation.ResourceLimitError{Resource: "accumulated events", Limit: limit} + } + return nil +} + func (u *ExUniverse) runConditionExecutor(ctx context.Context, args *conditionExecutorArgs) (bool, error) { if args.condition.Src == "" { return true, nil } + if err := ctx.Err(); err != nil { + return false, err + } if fn := builtin.GetCondition(args.condition.Src); fn != nil { - return fn(ctx, args) + callbackCtx, done := u.owner.callbackContext(ctx) + defer done() + approved, err := fn(callbackCtx, args) + if err != nil { + return false, err + } + if err := ctx.Err(); err != nil { + return false, err + } + return approved, nil } slog.WarnContext(ctx, "condition not found (default: return false)", "src", args.condition.Src) diff --git a/instrumentation/runtime.go b/instrumentation/runtime.go new file mode 100644 index 0000000..e3e61f2 --- /dev/null +++ b/instrumentation/runtime.go @@ -0,0 +1,54 @@ +package instrumentation + +import ( + "context" + "errors" + "fmt" +) + +var ( + ErrReentrantCall = errors.New("synchronous callback cannot reenter its owning machine; use callback arguments") + ErrMachineClosed = errors.New("quantum machine is closed") + ErrUnknownObserver = errors.New("observer is not registered") +) + +// RuntimeOptions configures opt-in execution policies. Zero limits mean unlimited. +// The zero value preserves the existing observer and invoke lifecycle contracts. +type RuntimeOptions struct { + MaxConcurrentInvokes int + CancelInvokesOnExit bool + MaxAccumulatedEvents int // Per universe, including copies accumulated for different realities. + MaxTrackingEntries int // Retain the most recent entries per universe. + StrictObservers bool +} + +func (o RuntimeOptions) Validate() error { + if o.MaxConcurrentInvokes < 0 || o.MaxAccumulatedEvents < 0 || o.MaxTrackingEntries < 0 { + return fmt.Errorf("runtime limits must not be negative") + } + return nil +} + +// ResourceLimitError reports rejected admission without starting another task or event copy. +type ResourceLimitError struct { + Resource string + Limit int +} + +func (e *ResourceLimitError) Error() string { + return fmt.Sprintf("%s limit reached (%d)", e.Resource, e.Limit) +} + +// RuntimeLifecycle is an optional capability. Close requests cooperative invoke +// cancellation and rejects new execution; WaitInvokes waits until the pool is idle. +// An invoke that ignores its context cannot be forcibly terminated. +type RuntimeLifecycle interface { + Close() error + WaitInvokes(context.Context) error +} + +// ContextSnapshotProvider adds cancellable waiting and callback reentry detection. +type ContextSnapshotProvider interface { + GetSnapshotContext(context.Context) (*MachineSnapshot, error) + LoadSnapshotContext(context.Context, *MachineSnapshot, any) error +} diff --git a/instrumentation/runtime_test.go b/instrumentation/runtime_test.go new file mode 100644 index 0000000..576fd2d --- /dev/null +++ b/instrumentation/runtime_test.go @@ -0,0 +1,16 @@ +package instrumentation + +import "testing" + +func TestRuntimeOptions_RechazaLimitesNegativos(t *testing.T) { + for _, options := range []RuntimeOptions{{MaxConcurrentInvokes: -1}, {MaxAccumulatedEvents: -1}, {MaxTrackingEntries: -1}} { + if err := options.Validate(); err == nil { + t.Fatalf("opcion invalida aceptada: %+v", options) + } + } + for _, options := range []RuntimeOptions{{}, {MaxConcurrentInvokes: 2, MaxAccumulatedEvents: 100, MaxTrackingEntries: 20, StrictObservers: true}} { + if err := options.Validate(); err != nil { + t.Fatal(err) + } + } +} diff --git a/statepro.go b/statepro.go index e77ab7e..d60237b 100644 --- a/statepro.go +++ b/statepro.go @@ -8,6 +8,14 @@ import ( ) func NewQuantumMachine(qmModel *theoretical.QuantumMachineModel) (instrumentation.QuantumMachine, error) { + return NewQuantumMachineWithOptions(qmModel, instrumentation.RuntimeOptions{}) +} + +// NewQuantumMachineWithOptions creates a machine with explicit runtime policies. +func NewQuantumMachineWithOptions(qmModel *theoretical.QuantumMachineModel, options instrumentation.RuntimeOptions) (instrumentation.QuantumMachine, error) { + if err := options.Validate(); err != nil { + return nil, err + } if qmModel == nil { return nil, fmt.Errorf("quantum machine model must not be nil") } @@ -18,7 +26,7 @@ func NewQuantumMachine(qmModel *theoretical.QuantumMachineModel) (instrumentatio } universes = append(universes, experimental.NewExUniverse(model)) } - return experimental.NewExQuantumMachine(qmModel, universes) + return experimental.NewExQuantumMachineWithOptions(qmModel, universes, options) } func NewEventBuilder(eventName string) instrumentation.EventBuilder { diff --git a/studio/package.json b/studio/package.json index 5eb4737..6d79a9e 100644 --- a/studio/package.json +++ b/studio/package.json @@ -34,7 +34,7 @@ "baseline-browser-mapping": "^2.11.0", "brace-expansion@^5.0.0": "^5.0.12", "postcss-selector-parser@^6.0.0": "^6.1.3", - "@babel/core": "^7.29.6", + "@babel/core@^7.0.0": "^7.29.6", "esbuild@^0.27.0": "0.28.2" } }, diff --git a/studio/packages/editor-core/src/StateProEditor.tsx b/studio/packages/editor-core/src/StateProEditor.tsx index 7bfd5a1..5cd55de 100644 --- a/studio/packages/editor-core/src/StateProEditor.tsx +++ b/studio/packages/editor-core/src/StateProEditor.tsx @@ -4321,25 +4321,25 @@ function StateProEditorInner({ return (