mirror of
https://github.com/bitechdev/ResolveSpec.git
synced 2026-10-01 19:20:31 +00:00
feat(websocketspec): fire BeforeDisconnect/AfterDisconnect hooks on close
Hooks fire from Connection.Close() once per connection. ConnectionManager Shutdown now closes connections outside its lock so hooks can call back into the manager. Update the single-transaction audit plan status.
This commit is contained in:
+16
-6
@@ -10,7 +10,7 @@
|
||||
- Each un-transacted call takes its own pool connection → bursts with a small pool (see `dbtrace`).
|
||||
- Already fixed: read/create hooks in `resolvespec` + `restheadspec` (commit `47708fc`, tag >= v1.1.28). Consumers on older tags still show the bug.
|
||||
|
||||
## Current state (verified by reading code; not yet by `dbtrace`)
|
||||
## Current state — BEFORE this work (historical baseline; everything below is now fixed, see Progress and Status)
|
||||
| Spec | Read | Create | Update | Delete |
|
||||
|---|---|---|---|---|
|
||||
| restheadspec | tx; `AfterRead` post-commit on pool | tx; `AfterCreate` post-commit on pool | tx; re-fetch + `BeforeScan` post-commit on pool (`:1667-1674`) | **single: no tx, hook + select + delete on pool (`:1945-1994`)**; batch: tx, per-item `BeforeDelete` inside |
|
||||
@@ -65,8 +65,18 @@
|
||||
- Update re-fetch is a plain SELECT in that second tx. No `RETURNING`.
|
||||
- `OnTxBegin` failure aborts the whole request, rolls back, returns an error with no detail leaked to the client.
|
||||
|
||||
## Open
|
||||
- Consumer's ResolveSpec version: confirm it is >= v1.1.28 (read/create already in tx). Not blocking.
|
||||
## Status summary
|
||||
**Done**
|
||||
- P0-P7 all DONE (baseline, delete in one tx, `OnTxBegin` + `runInTx`, second short tx, websocketspec/mqttspec, resolvemcp, funcspec, security stamping).
|
||||
- Regression tests in all six specs, plus source guard `pkg/common/tx_guard_test.go`.
|
||||
- `AfterRead` decided and done (restheadspec: second short tx; websocketspec/mqttspec: inside the read tx).
|
||||
- websocketspec `BeforeDisconnect`/`AfterDisconnect` wired (see Progress); `unwiredHooks` allowlist is now empty.
|
||||
|
||||
**Not done**
|
||||
- Real-Postgres `dbtrace` measurement for websocketspec, mqttspec, resolvemcp, restheadspec, funcspec (only resolvespec measured: `pooled=0` on every op). "Done when" bullet 1 is proven for resolvespec only.
|
||||
- resolvespec batch delete: per-item `BeforeDelete` (one hook per request today). Deferred on purpose: behavior change.
|
||||
- Confirm the consumer's ResolveSpec version is >= v1.1.28 (read/create already in tx). Not blocking; needs the consumer.
|
||||
- Known, pre-existing, not ours: `pkg/security` `TestDatabaseAuthenticator` fails with `-count=2` (use `-count=1`); mqttspec integration tests need a DB.
|
||||
|
||||
## Phases
|
||||
| # | Status | Change | Files | Notes |
|
||||
@@ -88,7 +98,7 @@
|
||||
- DONE P2 (resolvespec + restheadspec): `common.TxHookName`, `common.TxContext` (`SetTx` only; no abort/context accessors needed since `Execute` already returns an error on abort), `common.RunRequestTx`; per-spec `OnTxBegin`, `HookContext.SetTx`, `Handler.runInTx`. Every `RunInTransaction` in both handlers now goes through it. Tests: `pkg/*/on_tx_begin_test.go` (once, first, on tx, failure rolls back). Not yet: the post-commit second tx (P3) and the security stamping registration (P7).
|
||||
- DONE P3: restheadspec update re-fetch + `BeforeScan` + `AfterUpdate` and `AfterCreate` run in a second short `runInTx`; resolvespec update re-fetches (single, both batch paths) run in a second short `runInTx`. Fixed the pool reads inside the first tx (resolvespec single/batch update existing-record select, restheadspec update existence select) to use `tx`. Tests: `pkg/*/update_tx_test.go` (restheadspec uses the bun adapter; the pgsql adapter cannot build model-based updates).
|
||||
- NOTE: resolvespec fires no `AfterCreate`/`AfterRead`/`AfterUpdate`-post-commit hooks other than `AfterUpdate` inside the tx; nothing more to move there.
|
||||
- OPEN: restheadspec `AfterRead` still runs post-commit with `Tx = h.db` (`:1004`); decision says read has no second tx. Needs a call: run it inside the read tx, or in a short second tx.
|
||||
- RESOLVED: restheadspec `AfterRead` question (see DONE AfterRead below).
|
||||
- DONE P4: websocketspec + mqttspec. `OnTxBegin` (mqttspec re-exports the websocketspec constant), `HookContext.SetTx`, per-handler `runInTx`/`sendTxError`. Per message: read = 1 tx (Before/After hooks + queries); delete = 1 tx (Before, delete, After); create/update = tx 1 (Before + write) then tx 2 (re-fetch + `BeforeScan` + After). `create()`/`update()` no longer re-fetch; `read*`/`create`/`update`/`delete` use `hookCtx.Tx`. websocketspec `FetchRowNumber` keeps its public signature and delegates to a new tx-aware `fetchRowNumber`. A failure in begin/`OnTxBegin`/commit answers `transaction_error` with no detail. Tests: `pkg/websocketspec/tx_test.go` (sqlmock), `pkg/mqttspec/tx_test.go` (sqlite); mqttspec `update` tests now pass `Tx`.
|
||||
- DECIDED in P4 (follow `AfterRead` question above): websocketspec/mqttspec run `AfterRead` inside the read tx (keeps "read has no second tx").
|
||||
- DONE P5: resolvemcp. `OnTxBegin`, `HookContext.SetTx`, `Handler.runInTx`. Read = 1 tx (`BeforeRead`, count, scan, `AfterRead`; `readInTx`). Delete = 1 tx (`BeforeDelete` moved inside, after `OnTxBegin`). Create (single and batch, unified) = tx 1 (`BeforeCreate` + inserts) then tx 2 (re-fetch + `AfterCreate`); the old single-record pool insert/re-fetch is gone. Update = tx 1 (select, `BeforeUpdate`, update, `AfterUpdate`) then tx 2 (re-fetch). `BeforeHandle` still runs before any tx with `Tx = h.db`. Tests: `pkg/resolvemcp/tx_test.go` (sqlmock).
|
||||
@@ -103,7 +113,7 @@
|
||||
|
||||
## Tests
|
||||
- Existing: per-spec `handler_test.go`, `hooks_test.go`, `integration_test.go`; models in `pkg/testmodels/business.go`; `dbtrace` unit tests.
|
||||
- Done: delete tx tests (`pkg/*/delete_tx_test.go`, sqlmock, 1-conn pool detects pool use). Missing: same for read/create/update, `OnTxBegin`, other specs.
|
||||
- Done: delete tx tests (`pkg/*/delete_tx_test.go`, sqlmock, 1-conn pool detects pool use), and read/create/update/`OnTxBegin`/other-spec tests (see "DONE regression tests" in Progress). Still missing: `dbtrace` `pooled == 0` on real Postgres for all specs except resolvespec.
|
||||
- Add per spec/op: hook `Tx` is not the pool; `OnTxBegin` fires once per tx, before other hooks; single-ID delete = 1 tx; `dbtrace` `pooled == 0` on the request path.
|
||||
- Test data: reuse `pkg/testmodels`; **ask before generating new data** (per project rule).
|
||||
- Regression: full `go test -race` for security, dbmanager, common, restheadspec, resolvespec, websocketspec, mqttspec, resolvemcp, funcspec. Known pre-existing failures: mqttspec integration (no DB).
|
||||
@@ -118,5 +128,5 @@
|
||||
- `dbtrace` shows `pooled=0` for every handler op on a hooked model.
|
||||
- RLS GUC set in `OnTxBegin` is visible to read, create, update, delete queries and hooks.
|
||||
- No `Tx: h.db` / `hookCtx.Tx = h.db` left in spec handlers.
|
||||
- OPEN: websocketspec `BeforeDisconnect`/`AfterDisconnect` are defined but never executed (connection lifecycle, not DB). Allowlisted in `TestEveryDefinedHookHasACallSite`; wire them to remove the entry.
|
||||
- DONE: websocketspec `BeforeDisconnect`/`AfterDisconnect` fire from `Connection.Close()` (single close path, `closedOnce`, so exactly once per registered connection however it closes: read error, write error, slow-consumer eviction, shutdown). Connection lifecycle, not DB: `Tx` is not set. The hook context is detached from the connection cancel (`context.WithoutCancel`) so `AfterDisconnect` still has a live context. Errors are logged and never block the close. A connection rejected by `BeforeConnect` gets no disconnect hooks. `ConnectionManager.Shutdown` now closes connections outside its lock (a hook calling `Count()` would have deadlocked). Allowlist in `TestEveryDefinedHookHasACallSite` is now empty. Tests: `pkg/websocketspec/connection_test.go`. mqttspec already fired these in `Handler.Shutdown`; its per-client disconnect is unchanged.
|
||||
- DONE: column-level hide/mask columns are dropped from create/update payloads (`security.ApplyWriteColumnSecurity`); rules preloaded in `BeforeHandle` for create/update. resolvemcp update now runs `BeforeHandle`.
|
||||
|
||||
@@ -102,10 +102,7 @@ func TestSpecHandlersDoNotQueryThePoolDirectly(t *testing.T) {
|
||||
// Anything else defined in a spec's hooks.go must have an Execute call site: an
|
||||
// unwired hook silently disables whatever is registered on it (resolvespec's
|
||||
// AfterRead skipped column-level security masking until it was wired).
|
||||
var unwiredHooks = map[string]string{
|
||||
"websocketspec/BeforeDisconnect": "connection close is not hooked yet",
|
||||
"websocketspec/AfterDisconnect": "connection close is not hooked yet",
|
||||
}
|
||||
var unwiredHooks = map[string]string{}
|
||||
|
||||
var hookConstRE = regexp.MustCompile(`(?m)^\s*([A-Z][A-Za-z0-9]*)\s+HookType\s*=`)
|
||||
|
||||
|
||||
@@ -193,12 +193,17 @@ func (cm *ConnectionManager) GetConnection(id string) (*Connection, bool) {
|
||||
func (cm *ConnectionManager) Shutdown() {
|
||||
cm.cancel()
|
||||
|
||||
// Close all connections
|
||||
cm.mu.Lock()
|
||||
// Close all connections outside the lock: disconnect hooks may call back
|
||||
// into the manager (Count, GetConnection).
|
||||
cm.mu.RLock()
|
||||
conns := make([]*Connection, 0, len(cm.connections))
|
||||
for _, conn := range cm.connections {
|
||||
conns = append(conns, conn)
|
||||
}
|
||||
cm.mu.RUnlock()
|
||||
for _, conn := range conns {
|
||||
conn.Close()
|
||||
}
|
||||
cm.mu.Unlock()
|
||||
}
|
||||
|
||||
// ReadPump reads messages from the WebSocket connection
|
||||
@@ -300,6 +305,14 @@ func (c *Connection) SendJSON(v interface{}) error {
|
||||
// Close closes the connection
|
||||
func (c *Connection) Close() {
|
||||
c.closedOnce.Do(func() {
|
||||
// Disconnect hooks fire once per connection, however it closes. They
|
||||
// are best-effort: a failure is logged and never blocks the close.
|
||||
hookCtx := c.disconnectHookContext()
|
||||
c.logDisconnectHook(BeforeDisconnect, c.executeHook(func(r *HookRegistry) error { return r.Execute(BeforeDisconnect, hookCtx) }))
|
||||
defer func() {
|
||||
c.logDisconnectHook(AfterDisconnect, c.executeHook(func(r *HookRegistry) error { return r.Execute(AfterDisconnect, hookCtx) }))
|
||||
}()
|
||||
|
||||
if c.cancel != nil {
|
||||
c.cancel()
|
||||
}
|
||||
@@ -321,6 +334,34 @@ func (c *Connection) Close() {
|
||||
})
|
||||
}
|
||||
|
||||
// disconnectHookContext builds the hook context for the disconnect hooks. The
|
||||
// context is detached from the connection's cancellation so hooks that touch
|
||||
// the database still run after the connection context is cancelled.
|
||||
func (c *Connection) disconnectHookContext() *HookContext {
|
||||
parent := c.ctx
|
||||
if parent == nil {
|
||||
parent = context.Background()
|
||||
}
|
||||
return &HookContext{
|
||||
Context: context.WithoutCancel(parent),
|
||||
Handler: c.handler,
|
||||
Connection: c,
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Connection) executeHook(run func(*HookRegistry) error) error {
|
||||
if c.handler == nil || c.handler.hooks == nil {
|
||||
return nil
|
||||
}
|
||||
return run(c.handler.hooks)
|
||||
}
|
||||
|
||||
func (c *Connection) logDisconnectHook(hook HookType, err error) {
|
||||
if err != nil {
|
||||
logger.Warn("[WebSocketSpec] %s hook failed for connection %s: %v", hook, c.ID, err)
|
||||
}
|
||||
}
|
||||
|
||||
// AddSubscription adds a subscription to this connection
|
||||
func (c *Connection) AddSubscription(sub *Subscription) {
|
||||
c.mu.Lock()
|
||||
|
||||
@@ -594,3 +594,57 @@ func TestConnectionManager_CompleteLifecycle(t *testing.T) {
|
||||
_, exists = cm.GetConnection("conn-1")
|
||||
assert.False(t, exists)
|
||||
}
|
||||
|
||||
func TestConnection_Close_FiresDisconnectHooksOnce(t *testing.T) {
|
||||
h := &Handler{hooks: NewHookRegistry()}
|
||||
var order []string
|
||||
var ctxErrAfter error
|
||||
h.hooks.Register(BeforeDisconnect, func(hc *HookContext) error {
|
||||
order = append(order, "before")
|
||||
assert.NoError(t, hc.Context.Err())
|
||||
return nil
|
||||
})
|
||||
h.hooks.Register(AfterDisconnect, func(hc *HookContext) error {
|
||||
order = append(order, "after")
|
||||
ctxErrAfter = hc.Context.Err()
|
||||
return nil
|
||||
})
|
||||
|
||||
conn := createTestConnection("c1")
|
||||
conn.handler = h
|
||||
|
||||
conn.Close()
|
||||
conn.Close() // second close must not re-fire
|
||||
|
||||
assert.Equal(t, []string{"before", "after"}, order)
|
||||
assert.NoError(t, ctxErrAfter, "hook context must survive the connection cancel")
|
||||
}
|
||||
|
||||
func TestConnection_Close_HookErrorDoesNotBlockClose(t *testing.T) {
|
||||
h := &Handler{hooks: NewHookRegistry()}
|
||||
h.hooks.Register(BeforeDisconnect, func(*HookContext) error { return assert.AnError })
|
||||
|
||||
conn := createTestConnection("c2")
|
||||
conn.handler = h
|
||||
conn.Close()
|
||||
|
||||
assert.Error(t, conn.ctx.Err(), "connection must still be cancelled")
|
||||
}
|
||||
|
||||
func TestConnectionManager_Shutdown_HookMayCallManager(t *testing.T) {
|
||||
cm := NewConnectionManager(context.Background())
|
||||
h := &Handler{hooks: NewHookRegistry(), connManager: cm}
|
||||
h.hooks.Register(BeforeDisconnect, func(*HookContext) error { cm.Count(); return nil })
|
||||
|
||||
conn := createTestConnection("c3")
|
||||
conn.handler = h
|
||||
cm.connections[conn.ID] = conn
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() { cm.Shutdown(); close(done) }()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("Shutdown deadlocked on a hook calling the manager")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user