fix(session-projection): count registrants sharing one projection key
One unit definition already serves every session — its cells are keyed by `Session` — but registrants became per-session when agent presets started mounting tool packages per agent. N sessions on one preset register the same key N times. The first registration won and owned the only disposer, so ending one session stripped `goal`, `todos`, `plan`, `tokenUsage` and `contextPressure` from every other live session's snapshot. Measured against the shipped `standard` preset: two concurrent sessions each had eight projection keys, and disposing the first left the second with three — the ones host rows register. Count the registrants instead and remove the key when the last one goes. A differing `stateVersion` still refuses to share: it is the one incompatibility a runtime comparison can name, since everything else about a definition is functions.
This commit is contained in:
@@ -134,10 +134,22 @@ interface UnitCell {
|
|||||||
observedSeq: number
|
observedSeq: number
|
||||||
}
|
}
|
||||||
|
|
||||||
/** One live registration: the unit plus its per-session cells (dropped whole on disposal). */
|
/**
|
||||||
|
* One live registration: the unit plus its per-session cells (dropped whole
|
||||||
|
* once the last registrant releases it).
|
||||||
|
*
|
||||||
|
* `refs` exists because one unit definition already serves every session — the
|
||||||
|
* cells are keyed by `Session` — while the registrants are now per-session:
|
||||||
|
* an agent preset mounts the same tool package once per agent, so N sessions
|
||||||
|
* on one preset register the same key N times. Without a count the first
|
||||||
|
* registrant would own the disposer, and its session ending would strip the
|
||||||
|
* projection from every other live session.
|
||||||
|
*/
|
||||||
interface Registration {
|
interface Registration {
|
||||||
readonly def: ErasedDefinition
|
readonly def: ErasedDefinition
|
||||||
readonly cells: WeakMap<Session, UnitCell>
|
readonly cells: WeakMap<Session, UnitCell>
|
||||||
|
/** Live registrants sharing this unit; the last one out removes the key. */
|
||||||
|
refs: number
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -149,9 +161,12 @@ interface Registration {
|
|||||||
* older than the registry, folds `init` over the in-memory log on first
|
* older than the registry, folds `init` over the in-memory log on first
|
||||||
* touch (event or read). Registration is an effect (disposer rides the
|
* touch (event or read). Registration is an effect (disposer rides the
|
||||||
* calling fiber): an unloaded domain plugin's key disappears from snapshots
|
* calling fiber): an unloaded domain plugin's key disappears from snapshots
|
||||||
* and clients read it as capability absence. Duplicate keys throw. Domain
|
* and clients read it as capability absence. Domain
|
||||||
* plugins register under `ctx.inject(['sessionProjections'], …)` so headless
|
* plugins register under `ctx.inject(['sessionProjections'], …)` so headless
|
||||||
* assemblies without the registry stay unaffected.
|
* assemblies without the registry stay unaffected. Registrants sharing a key
|
||||||
|
* share one unit and are counted: the same tool package mounted in N agent
|
||||||
|
* presets registers N times, and the key survives until the last one
|
||||||
|
* unloads.
|
||||||
*/
|
*/
|
||||||
export class SessionProjectionRegistry extends Service {
|
export class SessionProjectionRegistry extends Service {
|
||||||
private readonly registrations = new Map<string, Registration>()
|
private readonly registrations = new Map<string, Registration>()
|
||||||
@@ -182,12 +197,25 @@ export class SessionProjectionRegistry extends Service {
|
|||||||
}
|
}
|
||||||
const dispose = this.ctx.effect(function* (this: SessionProjectionRegistry) {
|
const dispose = this.ctx.effect(function* (this: SessionProjectionRegistry) {
|
||||||
const key = definition.key as string
|
const key = definition.key as string
|
||||||
if (this.registrations.has(key)) {
|
const existing = this.registrations.get(key)
|
||||||
throw new Error(`session projection key ${JSON.stringify(key)} is already registered`)
|
if (existing === undefined) {
|
||||||
|
this.registrations.set(key, { def: definition, cells: new WeakMap(), refs: 1 })
|
||||||
|
} else {
|
||||||
|
// A differing `stateVersion` is the one incompatibility this can name:
|
||||||
|
// the versioned contract says the cached state shape differs, so the
|
||||||
|
// two registrants cannot share cells. Anything else about a definition
|
||||||
|
// is functions, which no runtime comparison can tell apart.
|
||||||
|
if (existing.def.stateVersion !== definition.stateVersion) {
|
||||||
|
throw new Error(`session projection key ${JSON.stringify(key)} is already registered at stateVersion ${String(existing.def.stateVersion)}; refusing to share it with stateVersion ${String(definition.stateVersion)}`)
|
||||||
|
}
|
||||||
|
existing.refs += 1
|
||||||
}
|
}
|
||||||
this.registrations.set(key, { def: definition, cells: new WeakMap() })
|
|
||||||
yield () => {
|
yield () => {
|
||||||
this.registrations.delete(key)
|
const live = this.registrations.get(key)
|
||||||
|
/* v8 ignore next -- the disposer runs once per successful registration, so the entry it counted is still here */
|
||||||
|
if (live === undefined) return
|
||||||
|
live.refs -= 1
|
||||||
|
if (live.refs === 0) this.registrations.delete(key)
|
||||||
}
|
}
|
||||||
}.bind(this), 'sessionProjections.register()')
|
}.bind(this), 'sessionProjections.register()')
|
||||||
return () => void dispose()
|
return () => void dispose()
|
||||||
|
|||||||
@@ -127,14 +127,45 @@ describe('SessionProjectionRegistry drive', () => {
|
|||||||
expect(snapshot.values['test/marks']).toEqual({ marks: [] })
|
expect(snapshot.values['test/marks']).toEqual({ marks: [] })
|
||||||
})
|
})
|
||||||
|
|
||||||
it('rejects duplicate keys loud and keeps the first unit', async () => {
|
it('shares one unit between registrants of the same key', async () => {
|
||||||
const { ctx, session } = await harness()
|
const { ctx, session } = await harness()
|
||||||
ctx.sessionProjections.register(marksUnit())
|
ctx.sessionProjections.register(marksUnit())
|
||||||
expect(() => ctx.sessionProjections.register(marksUnit())).toThrow(/"test\/marks" is already registered/)
|
|
||||||
|
// One definition already serves every session (cells are keyed by
|
||||||
|
// Session), and registrants are per-session now: an agent preset mounts
|
||||||
|
// the same tool package once per agent.
|
||||||
|
expect(() => ctx.sessionProjections.register(marksUnit())).not.toThrow()
|
||||||
mark(session, ['kept'])
|
mark(session, ['kept'])
|
||||||
expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['kept'] })
|
expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['kept'] })
|
||||||
})
|
})
|
||||||
|
|
||||||
|
it('keeps the unit until the last registrant releases it', async () => {
|
||||||
|
const { ctx, session } = await harness()
|
||||||
|
const first = ctx.sessionProjections.register(marksUnit())
|
||||||
|
const second = ctx.sessionProjections.register(marksUnit())
|
||||||
|
mark(session, ['kept'])
|
||||||
|
|
||||||
|
first()
|
||||||
|
|
||||||
|
// The regression this counts against: one session ending used to strip
|
||||||
|
// the projection from every other live session, because the first
|
||||||
|
// registrant owned the only disposer.
|
||||||
|
expect(ctx.sessionProjections.snapshot(session).values['test/marks']).toEqual({ marks: ['kept'] })
|
||||||
|
second()
|
||||||
|
expect(ctx.sessionProjections.snapshot(session).values).toEqual({})
|
||||||
|
})
|
||||||
|
|
||||||
|
it('refuses to share a key across a stateVersion change', async () => {
|
||||||
|
const { ctx } = await harness()
|
||||||
|
ctx.sessionProjections.register(marksUnit())
|
||||||
|
|
||||||
|
// The one incompatibility a runtime comparison can name: the versioned
|
||||||
|
// contract says the cached state shape differs, so the two cannot share
|
||||||
|
// cells. Everything else about a definition is functions.
|
||||||
|
expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: 9 }))
|
||||||
|
.toThrow(/already registered at stateVersion 1; refusing to share it with stateVersion 9/)
|
||||||
|
})
|
||||||
|
|
||||||
it('rejects a non-integer or negative stateVersion at register time', async () => {
|
it('rejects a non-integer or negative stateVersion at register time', async () => {
|
||||||
const { ctx } = await harness()
|
const { ctx } = await harness()
|
||||||
expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: -1 })).toThrow(/stateVersion/)
|
expect(() => ctx.sessionProjections.register({ ...marksUnit(), stateVersion: -1 })).toThrow(/stateVersion/)
|
||||||
|
|||||||
Reference in New Issue
Block a user