mirror of
https://github.com/microsoft/agent-framework.git
synced 2026-06-16 21:04:09 +08:00
977c3adfb2
* python: replace pre-commit with prek, add PEP 723 script deps, clean up dev dependencies - Replace pre-commit with prek (Rust-native, faster pre-commit alternative) - Move supported hooks to repo: builtin for zero-clone speed - Add new builtin hooks: trailing-whitespace, check-merge-conflict, detect-private-key, check-added-large-files - Update all hook versions to latest (pre-commit-hooks v6, pyupgrade v3.21.2, bandit 1.9.3, uv-pre-commit 0.10.0) - Add PEP 723 inline script metadata to 34 samples with external deps - Remove autogen-agentchat/autogen-ext from dev deps (now declared per-sample) - Remove unused dev deps: pytest-env, tomli-w - Add agent-framework-core>=1.0.0b260130 lower bound to all 21 packages - Update CI workflow to use j178/prek-action - Update docs: DEV_SETUP.md, AGENTS.md, CODING_STANDARD.md, SAMPLE_GUIDELINES.md * updated lock * python: fix prek config paths for local execution and CI workflow Remove global 'files: ^python/' filter and strip python/ prefix from all path patterns in .pre-commit-config.yaml so prek finds files when run from the python/ directory. Update CI workflow to use --cd python instead of --config path. Include trailing whitespace fixes and dev dependency cleanup. * python: move helper scripts to scripts/ folder and exclude from checks * python: exclude AGENTS.md from prek markdown code lint * python: exclude AGENTS.md and azure_ai_search sample from markdown lint * fix m365 sample * python: ignore CPY rule for samples with PEP 723 headers * fix in dev_setup * python: replace aiofiles with regular open in samples * python: suppress reportUnusedImport in markdown code block checker * python: use samples pyright config for markdown code block checker Write a temp pyrightconfig.json matching pyrightconfig.samples.json rules (typeCheckingMode=off, only reportMissingImports and reportAttributeAccessIssue). Filter output to only fail on these rules since syntax-level errors (top-level await, undefined vars) are expected in README documentation snippets. * python: use markdown-code-lint with fixed globs instead of prek file list The prek-markdown-code-lint task received all changed files including non-README markdown and files with pre-existing broken imports. Replace with the standard markdown-code-lint task which uses the correct glob patterns (README.md, packages/**/README.md, samples/**/*.md). * python: exclude READMEs with pre-existing broken imports from markdown lint * python: fix broken README code snippets instead of excluding them - ag-ui: replace TextContent (removed) with content.type == 'text' - durabletask: fix import path to durabletask.worker.TaskHubGrpcWorker - orchestrations: use constructor params instead of .participants() method - observability: mark deprecated code blocks as plain text, filter reportMissingImports to agent_framework modules only - remove README excludes from markdown-code-lint task * add revision to gaia download * feat(python): parallelize checks across packages Run (package × task) cross-product in parallel using ThreadPoolExecutor and subprocesses. Key changes: - Add scripts/task_runner.py with shared parallel execution engine - Update run_tasks_in_packages_if_exists.py to accept multiple tasks - Update run_tasks_in_changed_packages.py with --files flag and parallel support - Add check-packages poe task (fmt+lint+pyright+mypy in parallel) - Add prek-markdown-code-lint and prek-samples-check with change detection - Split CI code quality workflow into parallel prek and mypy jobs - Update DEV_SETUP.md to document new parallel behavior Core package changes still trigger checks on all packages. * feat(ci): split code quality into 4 parallel jobs Split the single prek job into parallel jobs: - pre-commit-hooks: lightweight hooks (SKIP=poe-check) - package-checks: fmt/lint/pyright/mypy via check-packages - samples-markdown: samples-lint, samples-syntax, markdown-code-lint - mypy: change-detected mypy checks All 4 jobs run concurrently (×2 Python versions = 8 runners). * feat(ci): use only Python 3.10 for code quality checks * refactor(python): add future annotations and remove quoted types Add `from __future__ import annotations` to 93 package files that used quoted string annotations, then run pyupgrade --py310-plus to remove the now-unnecessary quotes. Fixes https://github.com/microsoft/agent-framework/issues/3578
200 lines
5.3 KiB
TypeScript
200 lines
5.3 KiB
TypeScript
/**
|
|
* Streaming State Persistence
|
|
*
|
|
* Manages browser storage of streaming response state to enable:
|
|
* - Resume interrupted streams after page refresh
|
|
* - Replay cached events before fetching new ones
|
|
* - Graceful recovery from network disconnections
|
|
*/
|
|
|
|
import type { ExtendedResponseStreamEvent } from "@/types/openai";
|
|
|
|
export interface StreamingState {
|
|
conversationId: string;
|
|
responseId: string;
|
|
lastMessageId?: string;
|
|
lastSequenceNumber: number;
|
|
events: ExtendedResponseStreamEvent[];
|
|
timestamp: number; // When this state was last updated
|
|
completed: boolean; // Whether the stream completed successfully
|
|
accumulatedText?: string; // Accumulated text content for quick restoration
|
|
}
|
|
|
|
const STORAGE_KEY_PREFIX = "devui_streaming_state_";
|
|
const STATE_EXPIRY_MS = 24 * 60 * 60 * 1000; // 24 hours
|
|
|
|
/**
|
|
* Storage key for a specific conversation
|
|
*/
|
|
function getStorageKey(conversationId: string): string {
|
|
return `${STORAGE_KEY_PREFIX}${conversationId}`;
|
|
}
|
|
|
|
/**
|
|
* Extract accumulated text from events (for quick restoration)
|
|
*/
|
|
function extractAccumulatedText(events: ExtendedResponseStreamEvent[]): string {
|
|
let text = "";
|
|
for (const event of events) {
|
|
if (event.type === "response.output_text.delta" && "delta" in event) {
|
|
text += event.delta;
|
|
}
|
|
}
|
|
return text;
|
|
}
|
|
|
|
/**
|
|
* Save streaming state to browser storage
|
|
*/
|
|
export function saveStreamingState(state: StreamingState): void {
|
|
try {
|
|
const key = getStorageKey(state.conversationId);
|
|
const data = JSON.stringify(state);
|
|
localStorage.setItem(key, data);
|
|
} catch (error) {
|
|
console.error("Failed to save streaming state:", error);
|
|
// If storage is full, try to clear old states
|
|
try {
|
|
clearExpiredStreamingStates();
|
|
// Try again
|
|
const key = getStorageKey(state.conversationId);
|
|
const data = JSON.stringify(state);
|
|
localStorage.setItem(key, data);
|
|
} catch {
|
|
console.error("Failed to save streaming state even after cleanup");
|
|
}
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Load streaming state from browser storage
|
|
*/
|
|
export function loadStreamingState(conversationId: string): StreamingState | null {
|
|
try {
|
|
const key = getStorageKey(conversationId);
|
|
const data = localStorage.getItem(key);
|
|
|
|
if (!data) {
|
|
return null;
|
|
}
|
|
|
|
const state: StreamingState = JSON.parse(data);
|
|
|
|
// Check if state has expired
|
|
const age = Date.now() - state.timestamp;
|
|
if (age > STATE_EXPIRY_MS) {
|
|
clearStreamingState(conversationId);
|
|
return null;
|
|
}
|
|
|
|
// If stream was completed, no need to resume
|
|
if (state.completed) {
|
|
return null;
|
|
}
|
|
|
|
return state;
|
|
} catch (error) {
|
|
console.error("Failed to load streaming state:", error);
|
|
return null;
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Update streaming state with a new event
|
|
*/
|
|
export function updateStreamingState(
|
|
conversationId: string,
|
|
event: ExtendedResponseStreamEvent,
|
|
responseId: string,
|
|
lastMessageId?: string
|
|
): void {
|
|
try {
|
|
const existing = loadStreamingState(conversationId);
|
|
const sequenceNumber = "sequence_number" in event ? event.sequence_number : undefined;
|
|
|
|
const newEvents = existing ? [...existing.events, event] : [event];
|
|
|
|
const state: StreamingState = {
|
|
conversationId,
|
|
responseId,
|
|
lastMessageId,
|
|
lastSequenceNumber: sequenceNumber ?? (existing?.lastSequenceNumber ?? -1),
|
|
events: newEvents,
|
|
timestamp: Date.now(),
|
|
completed: event.type === "response.completed" || event.type === "response.failed",
|
|
accumulatedText: extractAccumulatedText(newEvents),
|
|
};
|
|
|
|
saveStreamingState(state);
|
|
} catch (error) {
|
|
console.error("Failed to update streaming state:", error);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Mark streaming state as completed
|
|
*/
|
|
export function markStreamingCompleted(conversationId: string): void {
|
|
try {
|
|
const existing = loadStreamingState(conversationId);
|
|
if (existing) {
|
|
existing.completed = true;
|
|
existing.timestamp = Date.now();
|
|
saveStreamingState(existing);
|
|
}
|
|
} catch (error) {
|
|
console.error("Failed to mark streaming as completed:", error);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Clear streaming state for a conversation
|
|
*/
|
|
export function clearStreamingState(conversationId: string): void {
|
|
try {
|
|
const key = getStorageKey(conversationId);
|
|
localStorage.removeItem(key);
|
|
} catch (error) {
|
|
console.error("Failed to clear streaming state:", error);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Clear all expired streaming states
|
|
*/
|
|
export function clearExpiredStreamingStates(): void {
|
|
try {
|
|
const keys = Object.keys(localStorage);
|
|
const now = Date.now();
|
|
|
|
for (const key of keys) {
|
|
if (key.startsWith(STORAGE_KEY_PREFIX)) {
|
|
try {
|
|
const data = localStorage.getItem(key);
|
|
if (data) {
|
|
const state: StreamingState = JSON.parse(data);
|
|
const age = now - state.timestamp;
|
|
|
|
if (age > STATE_EXPIRY_MS || state.completed) {
|
|
localStorage.removeItem(key);
|
|
}
|
|
}
|
|
} catch {
|
|
// Invalid state, remove it
|
|
localStorage.removeItem(key);
|
|
}
|
|
}
|
|
}
|
|
} catch (error) {
|
|
console.error("Failed to clear expired streaming states:", error);
|
|
}
|
|
}
|
|
|
|
/**
|
|
* Initialize streaming state management (call on app startup)
|
|
*/
|
|
export function initStreamingState(): void {
|
|
// Clear expired states on startup
|
|
clearExpiredStreamingStates();
|
|
}
|