11import * as Tool from "./tool"
22import DESCRIPTION from "./compact_bulk.txt"
3+ import type { Agent } from "@/agent/agent"
34import { Session } from "@/session/session"
5+ import { LLM } from "@/session/llm"
6+ import { MessageID } from "@/session/schema"
47import { Provider } from "@/provider/provider"
8+ import { errorMessage } from "@/util/error"
59import { SessionV1 } from "@opencode-ai/core/v1/session"
10+ import { LLMEvent } from "@opencode-ai/llm"
611import {
712 COMPACTION_SUMMARY_MAX_CHARS ,
813 compactionInput ,
@@ -13,8 +18,7 @@ import {
1318 setCompaction ,
1419 withCompactionLock ,
1520} from "@/session/compaction-pruning"
16- import { generateObject } from "ai"
17- import { Cause , Effect , Exit , Ref , Schema } from "effect"
21+ import { Cause , Effect , Exit , Option , Ref , Schema , Stream } from "effect"
1822import { Buffer } from "node:buffer"
1923
2024const id = "compact_bulk"
@@ -65,6 +69,8 @@ const Summaries = Schema.Struct({
6569 ) ,
6670} )
6771
72+ const decodeSummaries = Schema . decodeUnknownOption ( Schema . fromJsonString ( Summaries ) )
73+
6874const SYSTEM = [
6975 "You compress segments of an AI agent's own conversation transcript so it can reclaim context while staying effective." ,
7076 "You are given structured data containing a focus, the request being served, and an ordered transcript of messages." ,
@@ -74,8 +80,19 @@ const SYSTEM = [
7480 "For each item with an `id`, write a concise summary that PRESERVES every detail relevant to the focus — file paths, identifiers, decisions, numeric values, error messages, commands, and conclusions — and DROPS boilerplate, repeated verbatim output, and noise." ,
7581 'Write each summary so it stands alone once the original is gone: name what the item was and what came of it, rather than referring to it as "this output".' ,
7682 "Return exactly one entry per item with an `id`, echoing that id verbatim. Never summarize a context-only item, never invent ids, and never merge items." ,
83+ 'Reply with only a JSON object of this shape and no other text: {"summaries":[{"id":"<item id>","summary":"<summary>"}]}' ,
7784] . join ( "\n" )
7885
86+ // Its prompt replaces the provider's default agent prompt, so the summarizer
87+ // sees only its own instructions.
88+ const SUMMARIZER : Agent . Info = {
89+ name : id ,
90+ mode : "primary" ,
91+ permission : [ ] ,
92+ options : { } ,
93+ prompt : SYSTEM ,
94+ }
95+
7996type Outcome =
8097 | { partID : string ; status : "compacted" ; label : string ; freed : number }
8198 | { partID : string ; status : "skipped" ; reason : string }
@@ -97,11 +114,16 @@ type Metadata = {
97114 persistenceError ?: string
98115}
99116
100- export const CompactBulkTool = Tool . define < typeof Parameters , Metadata , Session . Service | Provider . Service > (
117+ export const CompactBulkTool = Tool . define <
118+ typeof Parameters ,
119+ Metadata ,
120+ Session . Service | Provider . Service | LLM . Service
121+ > (
101122 id ,
102123 Effect . gen ( function * ( ) {
103124 const sessions = yield * Session . Service
104125 const provider = yield * Provider . Service
126+ const llm = yield * LLM . Service
105127
106128 return {
107129 description : DESCRIPTION ,
@@ -260,13 +282,17 @@ export const CompactBulkTool = Tool.define<typeof Parameters, Metadata, Session.
260282 const target =
261283 turnModel ( messages ) ??
262284 ( yield * provider . defaultModel ( ) . pipe ( Effect . catch ( ( ) => Effect . succeed ( undefined ) ) ) )
263- const language = target
285+ const model = target
264286 ? yield * provider . getModel ( target . providerID , target . modelID ) . pipe (
265- Effect . flatMap ( ( model ) => provider . getLanguage ( model ) ) ,
287+ // Loaded here only to degrade before any spend. The session
288+ // stream reuses Provider's cached language model.
289+ Effect . tap ( ( model ) => provider . getLanguage ( model ) ) ,
266290 Effect . catch ( ( ) => Effect . succeed ( undefined ) ) ,
267291 )
268292 : undefined
269- if ( ! language )
293+ // A GitLab Duo workflow model runs its own agent loop, and a session
294+ // stream rewires the shared model instance that the running turn uses.
295+ if ( ! model || ( model . providerID === "gitlab" && model . api . id . startsWith ( "duo-workflow-" ) ) )
270296 return {
271297 title : "Bulk compaction unavailable" ,
272298 metadata : {
@@ -300,44 +326,63 @@ export const CompactBulkTool = Tool.define<typeof Parameters, Metadata, Session.
300326 ( group ) =>
301327 Effect . gen ( function * ( ) {
302328 yield * Ref . update ( callsRef , ( calls ) => calls + 1 )
303- const result = yield * Effect . tryPromise ( {
304- try : ( signal ) =>
305- abortable (
306- generateObject ( {
307- model : language ,
308- schema : Object . assign (
309- Schema . toStandardSchemaV1 ( Summaries ) ,
310- Schema . toStandardJSONSchemaV1 ( Summaries ) ,
311- ) ,
312- temperature : 0.2 ,
313- system : SYSTEM ,
314- prompt : renderPrompt ( params . focus , group , surrounding ) ,
315- abortSignal : AbortSignal . any ( [ ctx . abort , signal ] ) ,
316- } ) ,
317- ctx . abort ,
318- ) ,
319- catch : ( error ) => ( error instanceof Error ? error : new Error ( String ( error ) ) ) ,
320- } )
329+ // The session request path applies capability-aware sampling,
330+ // provider options, and headers. Some providers accept only
331+ // streaming calls, so the JSON arrives as text. ctx.abort
332+ // interrupts this fiber through the compaction lock, and the
333+ // interruption aborts the provider request.
334+ const events = yield * llm
335+ . stream ( {
336+ agent : SUMMARIZER ,
337+ user : {
338+ id : MessageID . ascending ( ) ,
339+ sessionID : ctx . sessionID ,
340+ role : "user" ,
341+ time : { created : Date . now ( ) } ,
342+ agent : SUMMARIZER . name ,
343+ model : { providerID : model . providerID , modelID : model . id } ,
344+ } ,
345+ sessionID : ctx . sessionID ,
346+ model,
347+ system : [ ] ,
348+ messages : [ { role : "user" , content : renderPrompt ( params . focus , group , surrounding ) } ] ,
349+ tools : { } ,
350+ retries : 2 ,
351+ } )
352+ . pipe ( Stream . runCollect )
353+ const failure = events . find ( LLMEvent . is . providerError )
354+ if ( failure ) return rejectAll ( group , `secondary summarization failed — ${ failure . message } ` )
355+ const finish = events . findLast ( LLMEvent . is . finish )
321356 yield * Ref . update ( usageRef , ( usage ) => ( {
322- inputTokens : usage . inputTokens + ( result . usage . inputTokens ?? 0 ) ,
323- outputTokens : usage . outputTokens + ( result . usage . outputTokens ?? 0 ) ,
324- totalTokens : usage . totalTokens + ( result . usage . totalTokens ?? 0 ) ,
357+ inputTokens : usage . inputTokens + ( finish ? .usage ? .inputTokens ?? 0 ) ,
358+ outputTokens : usage . outputTokens + ( finish ? .usage ? .outputTokens ?? 0 ) ,
359+ totalTokens : usage . totalTokens + ( finish ? .usage ? .totalTokens ?? 0 ) ,
325360 } ) )
326- return validateSummaries ( group , result . object . summaries )
361+ // A provider safety classifier can refuse a transcript that quotes
362+ // tool instructions. That is not a parse failure, and retrying the
363+ // same data is unlikely to help.
364+ if ( finish ?. reason === "content-filter" )
365+ return rejectAll ( group , "summarizer refused the request (content filter)" )
366+ const parsed = parseSummaries (
367+ events
368+ . filter ( LLMEvent . is . textDelta )
369+ . map ( ( event ) => event . text )
370+ . join ( "" ) ,
371+ )
372+ if ( ! parsed ) return rejectAll ( group , "summarizer returned invalid JSON" )
373+ return validateSummaries ( group , parsed . summaries )
327374 } ) . pipe (
328375 // Contained per chunk: one failed call must not discard the
329376 // summaries the other calls already paid for and got right.
330- Effect . catch ( ( error ) =>
331- ctx . abort . aborted
332- ? Effect . interrupt
333- : Effect . succeed ( {
334- accepted : [ ] as { id : string ; summary : string } [ ] ,
335- rejected : group . map ( ( item ) => ( {
336- id : item . id ,
337- reason : `secondary summarization failed — ${ error . message } ` ,
338- } ) ) ,
339- } ) ,
340- ) ,
377+ // Plugin hooks on the request path fail as defects, so the
378+ // whole cause is contained, not only typed errors.
379+ Effect . catchCause ( ( cause ) => {
380+ if ( ctx . abort . aborted ) return Effect . interrupt
381+ if ( Cause . hasInterrupts ( cause ) ) return Effect . failCause ( cause )
382+ return Effect . succeed (
383+ rejectAll ( group , `secondary summarization failed — ${ errorMessage ( Cause . squash ( cause ) ) } ` ) ,
384+ )
385+ } ) ,
341386 ) ,
342387 { concurrency : 4 } ,
343388 )
@@ -657,6 +702,19 @@ function validateInput(focus: string, partIDs: readonly string[]) {
657702 if ( oversizedID ) return `Part IDs are limited to ${ MAX_PART_ID_BYTES } bytes.`
658703}
659704
705+ // Without a structured-output mode, models often wrap the JSON in a code fence
706+ // or a sentence. The first candidate that decodes to the exact shape wins.
707+ function parseSummaries ( text : string ) {
708+ const start = text . search ( / \{ \s * " s u m m a r i e s " \s * : / )
709+ return [
710+ text ,
711+ ...Array . from ( text . matchAll ( / ` ` ` [ a - z ] * \s * ( [ \s \S ] * ?) ` ` ` / gi) , ( match ) => match [ 1 ] ?? "" ) ,
712+ ...( start === - 1 ? [ ] : [ text . slice ( start , text . lastIndexOf ( "}" ) + 1 ) ] ) ,
713+ ]
714+ . map ( ( candidate ) => decodeSummaries ( candidate ) )
715+ . find ( Option . isSome ) ?. value
716+ }
717+
660718// Judged one part at a time. A summarizer that returns junk for a single item —
661719// or silently drops it — used to cost the caller every other summary in the same
662720// call, which is the opposite of what a batch tool is for: each part either gets
@@ -692,6 +750,13 @@ function validateSummaries(
692750 return { accepted, rejected }
693751}
694752
753+ function rejectAll ( group : readonly { id : string } [ ] , reason : string ) {
754+ return {
755+ accepted : [ ] as { id : string ; summary : string } [ ] ,
756+ rejected : group . map ( ( item ) => ( { id : item . id , reason } ) ) ,
757+ }
758+ }
759+
695760function rejected ( reason : string , skipped : number ) {
696761 return {
697762 title : "Bulk compaction rejected" ,
@@ -700,22 +765,4 @@ function rejected(reason: string, skipped: number) {
700765 }
701766}
702767
703- function abortable < T > ( promise : Promise < T > , signal : AbortSignal ) {
704- if ( signal . aborted ) return Promise . reject ( signal . reason )
705- return new Promise < T > ( ( resolve , reject ) => {
706- const abort = ( ) => reject ( signal . reason )
707- signal . addEventListener ( "abort" , abort , { once : true } )
708- promise . then (
709- ( value ) => {
710- signal . removeEventListener ( "abort" , abort )
711- resolve ( value )
712- } ,
713- ( error ) => {
714- signal . removeEventListener ( "abort" , abort )
715- reject ( error )
716- } ,
717- )
718- } )
719- }
720-
721768export * as CompactBulk from "./compact_bulk"
0 commit comments