refactor(token-meter): merge measurement snapshots (round 1)
This commit is contained in:
14 files changed
+119
-115
No files matched your search
@@ -123,13 +123,7 @@ export class BasicCompactService extends CompactService {
|
||||
|
||||
let result: CompactionResult | null = null
|
||||
for (let attempt = 0; attempt <= this.config.compactionRetries; attempt += 1) {
|
||||
const surface = meter.measureSurface(agent.session)
|
||||
if (surface.logRevision !== measurement.logRevision) {
|
||||
throw new Error(
|
||||
`compaction: pressure revision ${measurement.logRevision} does not match surface revision ${surface.logRevision}`,
|
||||
)
|
||||
}
|
||||
const range = selectCompactableRange(agent.session, surface, this.config.retainTokens)
|
||||
const range = selectCompactableRange(agent.session, measurement, this.config.retainTokens)
|
||||
if (range === null) {
|
||||
/* v8 ignore else -- concrete replacement preserves a compactable checkpoint; subclass hooks cannot mutate it. */
|
||||
if (result === null) return null
|
||||
|
||||
@@ -10,7 +10,7 @@ import {
|
||||
toolPairingBalancedBefore,
|
||||
} from '@deepseek-ai/dsh-compact'
|
||||
import type { CompactionResult } from '@deepseek-ai/dsh-compact'
|
||||
import type { TokenMeterService, TokenSurfaceMeasurement } from '@deepseek-ai/dsh-token-meter'
|
||||
import type { TokenMeasurement, TokenMeterService } from '@deepseek-ai/dsh-token-meter'
|
||||
import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import type { Agent } from '@deepseek-ai/dsh-agent'
|
||||
import { frameSummary } from './summarizer.ts'
|
||||
@@ -25,16 +25,16 @@ interface RegionDependencies {
|
||||
* Resolve the next head-anchored range while retaining a priced recent tail
|
||||
* and never splitting an assistant tool-call/result pair.
|
||||
* @param session - session supplying authoritative current surface positions.
|
||||
* @param pricedSurface - same-revision surface measurement from the conversation meter.
|
||||
* @param measurement - unified pressure and surface measurement from the conversation meter.
|
||||
* @param retainTokens - minimum recent tail budget retained verbatim.
|
||||
* @returns the inclusive positional seq range to compact, or `null`.
|
||||
*/
|
||||
export function selectCompactableRange(
|
||||
session: Session,
|
||||
pricedSurface: TokenSurfaceMeasurement,
|
||||
measurement: TokenMeasurement,
|
||||
retainTokens: number,
|
||||
): { start: number; end: number } | null {
|
||||
const pricedNodes = pricedSurface.nodes
|
||||
const pricedNodes = measurement.nodes
|
||||
if (pricedNodes.length === 0) return null
|
||||
|
||||
const surfaceNodes = session.surface.nodes
|
||||
@@ -115,8 +115,8 @@ export async function compactSurfaceRegion(
|
||||
try {
|
||||
// Capture after the lock event so any later durable append, including a
|
||||
// log-only one, invalidates the async selection before replacement.
|
||||
const lockedSurface = dependencies.meter.measureSurface(session)
|
||||
const selected = lockedSurface.nodes.slice(startIdx, endIdx + 1)
|
||||
const lockedMeasurement = dependencies.meter.measure(session)
|
||||
const selected = lockedMeasurement.nodes.slice(startIdx, endIdx + 1)
|
||||
if (selected.length !== shadowedSeqs.length
|
||||
|| selected.some((node, index) => node.seq !== shadowedSeqs[index])) {
|
||||
throw new Error('compaction: selected surface changed before summarization began')
|
||||
@@ -125,8 +125,8 @@ export async function compactSurfaceRegion(
|
||||
const text = renderTranscript(session.events, shadowedSeqs)
|
||||
const { summary, model, maxTokens } = await dependencies.summarize(text, agent, signal)
|
||||
|
||||
const currentSurface = dependencies.meter.measureSurface(session)
|
||||
if (currentSurface.logRevision !== lockedSurface.logRevision) {
|
||||
const currentMeasurement = dependencies.meter.measure(session)
|
||||
if (currentMeasurement.logRevision !== lockedMeasurement.logRevision) {
|
||||
throw new Error('compaction: session log changed during summarization')
|
||||
}
|
||||
const framedSummary = frameSummary(summary)
|
||||
|
||||
@@ -251,17 +251,15 @@ describe('pressure measurement and retention', () => {
|
||||
expect(await compactIfNeeded(compact, retained, MODEL, 'x'.repeat(100_000))).toBeNull()
|
||||
})
|
||||
|
||||
it('detects scalar/surface revision disagreement', async () => {
|
||||
it('uses one unified measurement for each pressure-and-retention decision', async () => {
|
||||
const ctx = createContext()
|
||||
const meter = ctx.tokenMeter
|
||||
const original = meter.measureSurface.bind(meter)
|
||||
vi.spyOn(meter, 'measureSurface').mockImplementation((session) => {
|
||||
const measurement = original(session)
|
||||
return { ...measurement, logRevision: measurement.logRevision - 1 }
|
||||
})
|
||||
const compact = service(compactConfig, ctx)
|
||||
const measure = vi.spyOn(ctx.tokenMeter, 'measure')
|
||||
const stop = new Error('stop after first decision')
|
||||
vi.spyOn(compact, 'compactRegion').mockRejectedValueOnce(stop)
|
||||
|
||||
await expect(compactIfNeeded(compact, conversation(4))).rejects.toThrow(/revision/)
|
||||
await expect(compactIfNeeded(compact, conversation(4))).rejects.toBe(stop)
|
||||
expect(measure).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('bounds retries when a shrinking checkpoint remains above threshold', async () => {
|
||||
@@ -303,7 +301,7 @@ describe('pressure measurement and retention', () => {
|
||||
it('rejects a priced surface that is not the current positional surface', () => {
|
||||
const ctx = createContext()
|
||||
const session = conversation(2)
|
||||
const priced = ctx.tokenMeter.measureSurface(session)
|
||||
const priced = ctx.tokenMeter.measure(session)
|
||||
expect(() => selectCompactableRange(session, {
|
||||
...priced,
|
||||
nodes: priced.nodes.slice(1),
|
||||
@@ -331,7 +329,7 @@ describe('pressure measurement and retention', () => {
|
||||
}, { surfaceOp: 'append' })
|
||||
session.append('step/end', { turn: 1, step: 1 })
|
||||
|
||||
const priced = ctx.tokenMeter.measureSurface(session)
|
||||
const priced = ctx.tokenMeter.measure(session)
|
||||
expect(selectCompactableRange(session, priced, 1)).toBeNull()
|
||||
})
|
||||
})
|
||||
@@ -474,8 +472,8 @@ describe('compaction region transaction', () => {
|
||||
it('rejects a meter snapshot that changed before summarization began', async () => {
|
||||
const ctx = createContext()
|
||||
const meter = ctx.tokenMeter
|
||||
const original = meter.measureSurface.bind(meter)
|
||||
vi.spyOn(meter, 'measureSurface').mockImplementationOnce((session) => {
|
||||
const original = meter.measure.bind(meter)
|
||||
vi.spyOn(meter, 'measure').mockImplementationOnce((session) => {
|
||||
const measurement = original(session)
|
||||
return { ...measurement, nodes: measurement.nodes.slice(1) }
|
||||
})
|
||||
|
||||
Reference in New Issue
Block a user