Files
OpenViking/openviking/message/part.py
T
chenjwandClaude fd73dcf23a Feat/自进化(经验记忆)框架重构 (#2503)
* Add trajectory experience learning redesign doc

* auto-commit before eval 20260607_043406

* auto-commit before eval 20260607_044129

* auto-commit before eval 20260607_123706

* auto-commit before eval 20260607_125514

* auto-commit before eval 20260607_133737

* auto-commit before eval 20260607_144649

* auto-commit before eval 20260607_154631

* Refine streaming memory train merge pipeline

* Refine session train policy optimization architecture

* Add VikingMem ARA paper analysis

* Force merge for mixed extraction memory patches

* auto-commit before eval 20260608_134426

* auto-commit before eval 20260608_142108

* auto-commit before eval 20260608_153909

* auto-commit before eval 20260608_154845

* auto-commit before eval 20260608_170143

* update

* auto-commit before eval 20260611_150946

* auto-commit before eval 20260611_153933

* auto-commit before eval 20260611_154251

* Fix tau2 reward wrapper call

* auto-commit before eval 20260611_193803

* auto-commit before eval 20260611_194939

* update

* auto-commit before eval 20260612_111029

* auto-commit before eval 20260612_112104

* auto-commit before eval 20260612_122603

* auto-commit before eval 20260612_123359

* auto-commit before eval 20260612_124303

* auto-commit before eval 20260612_130257

* Fallback peer routing to first conversation peer

* Route self memory through self peer sentinel

* Keep self sentinel out of peer memory paths

* auto-commit before eval 20260612_154051

* auto-commit before eval 20260612_154850

* auto-commit before eval 20260612_161633

* auto-commit before eval 20260612_184022

* auto-commit before eval 20260612_201845

* auto-commit before eval 20260612_202637

* auto-commit before eval 20260612_204040

* auto-commit before eval 20260612_224621

* Fix locomo progress column initialization

* Add memory field versioning

* auto-commit before eval 20260612_232318

* Simplify locomo progress display

* Remove locomo progress elapsed time

* Batch streaming memory merges by group

* Derive patch merge language from patches

* Detect patch merge language from updated files

* auto-commit before eval 20260613_004339

* auto-commit before eval 20260613_005835

* Persist memory update trace id

* auto-commit before eval 20260613_012722

* auto-commit before eval 20260613_013923

* auto-commit before eval 20260613_014708

* Enforce peer scope after memory merge

* auto-commit before eval 20260613_033402

* auto-commit before eval 20260613_151931

* auto-commit before eval 20260613_164217

* chore: raise vikingbot eval parallelism

* chore: tune vikingbot parallelism to 150

* auto-commit before eval 20260613_185807

* chore: restore vikingbot parallelism default

* feat(locomo): add import progress reporting

* chore(memory): restore profile and preference templates

* Fix tau2 reward JSON serialization

* Refactor tau2 batch memory training

* Stream batch train JSONL events

* Add fast path for batch training case specs

* Optimize streaming train gradient chunking

* Optimize patch merge prompt context

* fix tau2 memory training vectorization

* fix(memory): revert profile preference granularity rules

* bd init: initialize beads issue tracking

* update

* Log memory template fallback failures

* Record all rollout artifacts

* Fix OpenViking peer search forwarding

* Stop tracking Beads local state

* auto-commit before eval 20260616_002037

* Deprecate memory version selector

* Retry transient LoCoMo import HTTP failures

* Add memory schema stage and peer routing

* Organize LoCoMo benchmark outputs

* Restore VikingBot user memory auto recall

* Show elapsed time on LoCoMo progress bars

* Quiet transient import retries

* Shorten LoCoMo progress bars

* Route non-peer memories to self scope

* auto-commit before eval 20260616_124513

* Suppress memory read not found logs

* Limit LoCoMo import memory types

* Rename peer routing schema flag

* Rename peer schema flag to enable_peer

* Rename schema peer flag to peer_enabled

* auto-commit before eval 20260616_135946

* auto-commit before eval 20260616_140641

* auto-commit before eval 20260616_141753

* Show cached baseline eval at start of training

* Preserve remote policy contents

* Show failed work in progress bars

* Hide zero failed progress counts

* Disable tau2 service progress by default

* Reuse policy lock for policy deletes

* feat: add session skill extraction to Memory V3 streaming trainer

- Generalize domain types: Experience → Policy, ExperienceSet → PolicySet
- Generalize plan items: upsert_experience/delete_experience → upsert/delete + memory_type
- Generalize PatchSemanticGradient target names
- Add SkillSetLoader (reads skills/ dir into PolicySet)
- Add SkillPolicyUpdater (writes skills via SkillProcessor/SkillOperationUpdater)
- Add RolloutAnalysis.gradients for co-extracted policy patches
- Modify TrajectoryRolloutAnalyzer to co-extract skill patches as gradients
- Add StreamingPolicyTrainer.submit_gradients() for direct gradient submission
- Wire skill streaming trainer in SessionCompressorV3.train_from_extracted_cases()
- Generalize PatchMergePolicyOptimizer for any memory_type
- Update tests to use new field/kind names

Co-authored-by: Claude <noreply@anthropic.com>

* Persist experience reminders in tau2 rollouts

* Enable tau2 epoch test eval by default

* Persist train rollout artifacts incrementally

* Ensure tau2 vikingbot user simulator deps

* Auto repair tau2 vikingbot simulator deps

* Avoid blocking tau2 vikingbot service loop

* Avoid tau2 gym reset when loading cases

* Clean tau2 rollout commit messages

* Clean tau2 tool trajectory serialization

* Retry vikingbot VLM rate limits

* Refine tau2 training case selection

* Promote vikingbot hook execution log level

* Improve VLM rate limit retry detection

* Update trajectory analysis prompt format

* Limit tau2 service logs to warnings

* Run tau2 vikingbot rollouts on service loop

* Lower vikingbot experience recall threshold

* Offload tau2 vikingbot blocking setup

* Retry tau2 LiteLLM rate limits

* Pin trajectory and experience outputs to Chinese

* Retry tau2 rate limits indefinitely

* Highlight tau2 training accuracy summaries

* Hide redundant avg reward console metrics

* Tighten memory extraction templates

* Reduce tau2 memory template noise

Evaluation: benchmark/tau2/train/run_batch_train_eval.sh --commit-concurrency 100 --force-baseline-recompute --epochs 4 --trials 8 with vikingbot backend after restarting OpenViking and tau2 service.

Result: epoch 1 test accuracy improved to 58.75% ± 4.84pp (94/160), compared with prior epoch 1 test reference 46.88% (75/160). Baseline in this run was 51.25%; epoch 0 test was 45.62%.

* Constrain tau2 memory extraction sources

Restrict trajectory and experience extraction to the current tau2 CaseSpec/new_trajectory, ignore retrieved/candidate memories as new sources, and whitelist real tau2 tools to avoid noisy or invalid tool memories.

Evaluation:
- Command: benchmark/tau2/train/run_batch_train_eval.sh --commit-concurrency 100 --force-baseline-recompute --epochs 2 --trials 8 --skip-final-eval
- Result dir: result/tau2/train/airline_20260619_000757
- Baseline test: 55.00% (88/160)
- Epoch0 train: 66.67% (20/30)
- Epoch0 test: 56.25% (90/160)
- Epoch1 train: 60.00% (18/30)
- Epoch1 test: 60.00% ± 3.54pp (96/160), better than previous best 58.75%.

* Preserve tau2 train non-run results

* Improve memory extraction guardrails

Run: result/tau2/train/run_airline_20260619_044051

tau2 airline epoch1 test/final: 62.50% (100/160), baseline cache hit 55.00% (88/160), delta +7.50pp; exceeds previous best 60.00% by +2.50pp.

* Support train split eval in tau2 batch runs

* Add slot support to tau2 vikingbot launcher

* Copy OpenViking configs for tau2 slots

* Tune tau2 case1 memory extraction

Run: result/tau2/train_1/run_airline_20260619_201546

Metric: train case1, slot1, 2 epochs, final train eval 3/8 = 37.50%, delta +37.50pp.

* Advise tau2 train case1 best result

Best run: result/tau2/train_1/run_airline_20260619_201546, final 3/8 = 37.50%.

* Tune tau2 memory gate extraction

* Advise tau2 train case1 50pct result

* Guard failed write experience branches

* Advise tau2 train case1 100pct result

* Guard tau2 oracle training memories

* Recall trajectory diagnostics for tau2 rollouts

* Recall tau2 case specs for training rollouts

* Guard evaluated tau2 final states

* Inject compact tau2 oracle checklists

* Stabilize tau2 slot train multi-case runs

* Guard tau2 case10 oracle terminal state

* Use supported tau2 training memory types

* Match tau2 oracle writes by expected subset

* Autofill tau2 case10 oracle writes before done

* Enable tau2 case10 guard for train split

* Record slot1 S008 case10 guard best advice

* Generalize tau2 S008 oracle terminal guard

* Record slot1 S008 general guard best advice

* Remove tau2 benchmark oracle guard

* Prevent training ground truth memory recall

* Refine tau2 training memory extraction

* Fix epoch train rollout artifact stage

* Refine memory training rollout pipeline

* update

* auto-commit before eval 20260623_120317

* fix sdk read_raw for memory metadata

* use visible case links for experience recall

* auto-commit before eval 20260623_225354

* tau2/train: cap run_batch_train_eval rollout concurrency at 100

* update

* update

* update

* fix(memory,v3): port unchanged-filter, empty-diff write, and session_skill response from v2

- Port _same_memory_file filter to compressor_v3._build_memory_diff so
  no-op merges/patches don't inflate memory_diff.json update counts
- Write memory_diff.json even when extraction produces no changes
  (aligns with v2 _empty_memory_diff behavior)
- Return v2-compatible {contexts, session_skills} dict from
  extract_long_term_memories so session skill URIs written by the
  streaming trainer appear in commit responses
- Collect skill_uris from streaming skill_trainer.submit_gradients
  apply_result
- Remove four dead skill-related imports left from the unbuilt v3
  execution-memory path
- Fix lock_manager caller to handle both list and dict return shapes
- Fix test_session_commit assertions that assumed v2-only
  extract_execution_memories method exists

* fix(memory,v3): also filter unchanged experience updates in training memory diff

* train: finish rollout and memory refactor

* memory: refine runtime-visible extraction prompts

* train: constrain communication memory extraction

* auto-commit before eval 20260629_235623

* memory: address training review fixes

* update

* update

* message: reuse part deserializer

* train: snapshot memory prompt yaml

* prompts: restore memory yaml templates from main

* memory: scope streaming update results

* update

* update

* session: train canonical merged cases

---------

Co-authored-by: Claude <noreply@anthropic.com>
2026-07-03 11:27:17 +08:00

147 lines
5.4 KiB
Python

# Copyright (c) 2026 Beijing Volcano Engine Technology Co., Ltd.
# SPDX-License-Identifier: AGPL-3.0
"""Part type definitions - based on opencode Part design.
Message consists of multiple Parts, each Part has different type and purpose.
"""
from dataclasses import dataclass
from typing import Any, Dict, Literal, Optional, Union
@dataclass
class TextPart:
"""Text content component."""
text: str = ""
type: Literal["text"] = "text"
@dataclass
class ContextPart:
"""Context reference component (L0 abstract + URI).
Used to track which contexts (memory/resource/skill) the message references.
"""
type: Literal["context"] = "context"
uri: str = ""
context_type: Literal["memory", "resource", "skill"] = "memory"
abstract: str = ""
@dataclass
class ImagePart:
"""Image URL component compatible with OpenAI-style message content."""
type: Literal["image_url"] = "image_url"
url: str = ""
detail: Optional[str] = None
@dataclass
class ToolPart:
"""Tool call component (references tool file within session).
Tool status: pending | running | completed | error
"""
type: Literal["tool"] = "tool"
tool_id: str = ""
tool_name: str = ""
tool_uri: str = "" # viking://user/{user_id}/sessions/{session_id}/tools/{tool_id}
skill_uri: str = "" # viking://user/{user_id}/skills/{skill_name}
tool_input: Optional[dict] = None
tool_output: str = ""
tool_status: str = "pending" # pending | running | completed | error
duration_ms: Optional[float] = None # 执行耗时(毫秒)
prompt_tokens: Optional[int] = None # 输入 Token
completion_tokens: Optional[int] = None # 输出 Token
tool_output_ref: str = ""
tool_output_truncated: bool = False
tool_output_original_chars: Optional[int] = None
tool_output_preview_chars: Optional[int] = None
tool_output_sha256: str = ""
tool_output_storage_uri: str = ""
tool_output_mime_type: str = "text/plain"
tool_output_source_ref: str = ""
tool_output_source_offset: Optional[int] = None
tool_output_source_limit: Optional[int] = None
tool_output_externalization_error: str = ""
tool_output_group_id: str = ""
tool_output_externalized_reason: str = ""
tool_output_group_original_chars: Optional[int] = None
tool_output_group_budget_chars: Optional[int] = None
Part = Union[TextPart, ContextPart, ImagePart, ToolPart]
def _parse_image_url_payload(data: Dict[str, Any]) -> tuple[str, Optional[str]]:
image_url = data.get("image_url")
if isinstance(image_url, dict):
return str(image_url.get("url", "") or ""), image_url.get("detail")
if isinstance(image_url, str):
return image_url, None
return "", None
def part_from_dict(data: Dict[str, Any]) -> Part:
"""Convert a dict to a Part object.
Args:
data: Dictionary with part data. Must contain 'type' field.
Returns:
Part object (TextPart, ContextPart, or ToolPart)
"""
part_type = data.get("type", "text")
if part_type == "text":
return TextPart(text=data.get("text", ""))
elif part_type == "context":
return ContextPart(
uri=data.get("uri", ""),
context_type=data.get("context_type", "memory"),
abstract=data.get("abstract", ""),
)
elif part_type == "image_url":
url, detail = _parse_image_url_payload(data)
if not url.strip():
raise ValueError("image_url part requires a non-empty URL")
return ImagePart(
url=url,
detail=detail,
)
elif part_type == "tool":
return ToolPart(
tool_id=data.get("tool_id", ""),
tool_name=data.get("tool_name", ""),
tool_uri=data.get("tool_uri", ""),
skill_uri=data.get("skill_uri", ""),
tool_input=data.get("tool_input"),
tool_output=data.get("tool_output", ""),
tool_status=data.get("tool_status", "pending"),
duration_ms=data.get("duration_ms"),
prompt_tokens=data.get("prompt_tokens"),
completion_tokens=data.get("completion_tokens"),
tool_output_ref=data.get("tool_output_ref", ""),
tool_output_truncated=bool(data.get("tool_output_truncated", False)),
tool_output_original_chars=data.get("tool_output_original_chars"),
tool_output_preview_chars=data.get("tool_output_preview_chars"),
tool_output_sha256=data.get("tool_output_sha256", ""),
tool_output_storage_uri=data.get("tool_output_storage_uri", ""),
tool_output_mime_type=data.get("tool_output_mime_type", "text/plain"),
tool_output_source_ref=data.get("tool_output_source_ref", ""),
tool_output_source_offset=data.get("tool_output_source_offset"),
tool_output_source_limit=data.get("tool_output_source_limit"),
tool_output_externalization_error=data.get("tool_output_externalization_error", ""),
tool_output_group_id=data.get("tool_output_group_id", ""),
tool_output_externalized_reason=data.get("tool_output_externalized_reason", ""),
tool_output_group_original_chars=data.get("tool_output_group_original_chars"),
tool_output_group_budget_chars=data.get("tool_output_group_budget_chars"),
)
else:
if "text" in data:
return TextPart(text=str(data.get("text", "") or ""))
return TextPart(text=str(data))