Files
openclaw/src/agents/embedded-agent-runner/run/message-transform-stream-wrapper.ts
T
Peter Steinberger 7eed2c3f21 feat(google): add current-turn native video input (#122074)
* feat(agents): add current-turn Gemini video handoff

* test(google): add live native video regression

* build(ai): emit provider types entrypoint

* fix(google): preserve video shedding on retry
2026-08-11 12:58:32 -07:00

61 lines
2.1 KiB
TypeScript

import type { StreamFn } from "openclaw/plugin-sdk/agent-core";
/**
* Wraps stream functions with pre-call message transforms.
*/
import {
PROVIDER_CONTEXT_HANDOFF,
type ProviderContext,
type ProviderStreamOptions,
} from "../../../../packages/ai/src/provider-types.js";
import type { AgentMessage } from "../../runtime/index.js";
/**
* Stream wrapper for applying message transforms immediately before provider dispatch.
*/
type MessageTransform = (messages: AgentMessage[], model: unknown) => AgentMessage[];
type ProviderContextMaterializer = (input: {
context: Parameters<StreamFn>[1];
signal?: AbortSignal;
}) => Promise<ProviderContext>;
/** Wraps a stream function with a conditional message-list transform. */
export function wrapStreamFnWithMessageTransform(
streamFn: StreamFn,
transform: MessageTransform,
materializeProviderContext?: ProviderContextMaterializer,
): StreamFn {
return (model, context, options) => {
const messages = (context as unknown as { messages?: unknown })?.messages;
const nextMessages = Array.isArray(messages)
? transform(messages as AgentMessage[], model)
: messages;
const nextContext =
Array.isArray(messages) && nextMessages !== messages
? ({
...(context as unknown as Record<string, unknown>),
messages: nextMessages,
} as typeof context)
: context;
if (!materializeProviderContext) {
return streamFn(model, nextContext, options);
}
let availableContext: Parameters<StreamFn>[1] | undefined = nextContext;
const handoff = async (): Promise<ProviderContext> => {
const captured = availableContext;
availableContext = undefined;
if (!captured) {
throw new Error("provider context handoff already consumed");
}
options?.signal?.throwIfAborted();
return await materializeProviderContext({
context: captured,
signal: options?.signal,
});
};
return streamFn(model, nextContext, {
...options,
[PROVIDER_CONTEXT_HANDOFF]: handoff,
} as ProviderStreamOptions & NonNullable<Parameters<StreamFn>[2]>);
};
}