update: bfcl cookbook to support using ReMe service

This commit is contained in:
caozouying.czy 2025-08-29 00:59:37 +08:00
parent 5ce1697738
commit 15795221dd
6 changed files with 340 additions and 115 deletions

View file

@ -63,14 +63,14 @@ class BFCLAgent:
max_response_size: int = 2000,
num_runs: int = 1,
enable_thinking: bool = False,
use_experience: bool = False,
use_fixed_experience: bool = True,
use_experience_deletion: bool = False,
use_memory: bool = False,
use_memory_addition: bool = False,
use_memory_deletion: bool = False,
delete_freq: int = 10,
freq_threshold: int = 5,
utility_threshold: float = 0.5,
experience_base_url: str = "http://0.0.0.0:8001/",
experience_workspace_id: str = "bfcl_8b_0725"):
memory_base_url: str = "http://0.0.0.0:8001/",
memory_workspace_id: str = "bfcl_8b_0725"):
self.index: int = index
self.task_ids: List[str] = task_ids
@ -84,17 +84,17 @@ class BFCLAgent:
self.max_response_size: int = max_response_size
self.num_runs: int = num_runs
self.enable_thinking: bool = enable_thinking
self.use_experience: bool = use_experience
self.use_fixed_experience: bool = use_fixed_experience if use_experience else True
self.use_experience_deletion: bool = use_experience_deletion if not use_fixed_experience else False
self.use_memory: bool = use_memory
self.use_memory_addition: bool = use_memory_addition if use_memory else False
self.use_memory_deletion: bool = use_memory_deletion if use_memory else False
self.delete_freq: int = delete_freq
self.freq_threshold: int = freq_threshold
self.utility_threshold: float = utility_threshold
self.experience_base_url: str = experience_base_url
self.experience_workspace_id: str = experience_workspace_id
self.memory_base_url: str = memory_base_url
self.memory_workspace_id: str = memory_workspace_id
self.history: List[List[List[dict]]] = [[] for _ in range(num_runs)]
self.retrieved_experience_ids: List[List[List[str]]] = [[] for _ in range(num_runs)]
self.retrieved_memory_list: List[List[List[Any]]] = [[] for _ in range(num_runs)]
self.test_entry: List[List[Dict[str, Any]]] = [[] for _ in range(num_runs)]
self.original_test_entry: List[List[Dict[str, Any]]] = [[] for _ in range(num_runs)]
self.tool_schema: List[List[List[dict]]] = [[] for _ in range(num_runs)]
@ -112,74 +112,71 @@ class BFCLAgent:
# 初始历史
msg = self.test_entry[run_id][i].get("messages", [])[0]
if self.use_experience:
if self.use_memory:
query = msg["content"]
response = self.get_experience(query)
if len(response["experience_list"]):
self.retrieved_experience_ids[run_id].append([e["experience_id"] for e in response["experience_list"]])
exp: str = response["experience_merged"]
print(f"experience_merged={exp}")
self.history[run_id].append([self.get_query_with_experience(query, exp)])
self.update_experience_freq(self.retrieved_experience_ids[run_id][i])
response = self.get_memory(query)
if len(response["metadata"]["memory_list"]):
self.retrieved_memory_list[run_id].append(response["metadata"]["memory_list"])
exp: str = response["answer"]
# print(f"memory_merged={exp}")
self.history[run_id].append([self.get_query_with_memory(query, exp)])
else:
self.retrieved_experience_ids[run_id].append([])
self.retrieved_memory_list[run_id].append([])
self.history[run_id].append([msg])
else:
self.history[run_id].append([msg])
self.current_turn[run_id][i] = 1
def get_query_with_experience(self, query: str, experience: str):
def get_query_with_memory(self, query: str, memory: str):
return {
"role": "user",
"content": "Task:\n" + query + "\n\nSome Related Experience to help you to complete the task:\n" + experience
"content": "Task:\n" + query + "\n\nSome Related Experience to help you to complete the task:\n" + memory
}
def get_experience(self, query: str):
response = requests.post(url=self.experience_base_url + "retriever", json={
"workspace_id": self.experience_workspace_id,
def get_traj_from_task_history(self, task_id: str, task_history: list, reward: float):
return {
"task_id": task_id,
"messages": task_history,
"score": reward
}
def get_memory(self, query: str):
response = requests.post(url=self.memory_base_url + "retrieve_task_memory", json={
"workspace_id": self.memory_workspace_id,
"query": query,
"top_k": 5
})
logger.info(f"query:{query}")
if response.status_code != 200:
logger.info(response.text)
return ""
response = response.json()
logger.info(response)
logger.info(f"query: {query}, response: {response}")
return response
def add_experience(self, trajectories):
response = requests.post(url=self.experience_base_url + "summarizer", json={
"workspace_id": self.experience_workspace_id,
"traj_list": trajectories,
def add_memory(self, trajectories):
response = requests.post(url=self.memory_base_url + "summary_task_memory", json={
"workspace_id": self.memory_workspace_id,
"trajectories": trajectories,
})
response.raise_for_status()
response = response.json()
logger.info(f"add new experiences: {response["experience_list"]}")
logger.info(f"add new memorys: {response["metadata"]["memory_list"]}")
def update_experience_freq(self, experience_ids):
response = requests.post(url=self.experience_base_url + "vector_store", json={
"workspace_id": self.experience_workspace_id,
"action": "update_freq",
"experience_ids": experience_ids,
def update_memory_information(self, memory_list, update_utility: bool=False):
response = requests.post(url=self.memory_base_url + "record_task_memory", json={
"workspace_id": self.memory_workspace_id,
"memory_dicts": memory_list,
"update_utility": update_utility
})
response.raise_for_status()
logger.info(response.json())
def update_experience_utility(self, experience_ids):
response = requests.post(url=self.experience_base_url + "vector_store", json={
"workspace_id": self.experience_workspace_id,
"action": "update_utility",
"experience_ids": experience_ids,
})
response.raise_for_status()
def delete_experience(self):
response = requests.post(url=self.experience_base_url + "vector_store", json={
"workspace_id": self.experience_workspace_id,
"action": "utility_based_delete",
def delete_memory(self):
response = requests.post(url=self.memory_base_url + "delete_task_memory", json={
"workspace_id": self.memory_workspace_id,
"freq_threshold": self.freq_threshold,
"utility_threshold": self.utility_threshold
})
@ -563,21 +560,18 @@ class BFCLAgent:
break
reward = self.get_reward(run_id, task_index)
if reward == 1 and self.use_experience and not self.use_fixed_experience:
# selectively add experiences when succeed
self.add_experience([{
"task_id":task_id,
"messages":self.history[run_id][task_index],
"score":reward
}])
if self.use_memory:
if reward == 1 and self.use_memory_addition: # selectively add memories when succeed
new_traj_list = [self.get_traj_from_task_history(task_id, self.history[run_id][task_index], reward)]
self.add_memory(new_traj_list)
if len(self.retrieved_experience_ids[run_id][task_index]):
# update the utility-related attributes of retrieved experiences
self.update_experience_utility(self.retrieved_experience_ids[run_id][task_index])
# update the freq & utility attributes of retrieved memories
update_utility: bool = (reward == 1)
self.update_memory_information(self.retrieved_memory_list[run_id][task_index], update_utility)
counter += 1
if self.use_experience_deletion and counter % self.delete_freq == 0:
self.delete_experience()
if self.use_memory_deletion and counter % self.delete_freq == 0:
self.delete_memory()
t_result = {
"run_id": run_id,

View file

@ -0,0 +1,223 @@
import json
import requests
import argparse
from pathlib import Path
from typing import List, Dict, Any
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor, as_completed
def load_task_case(data_path: str, task_id: str | None) -> Dict[str, Any]:
"""
load training cases by id
"""
if not Path(data_path).exists():
raise FileNotFoundError(f"BFCL data file '{data_path}' not found")
if task_id is None:
raise ValueError("task_id is required")
with open(data_path, "r", encoding="utf-8") as f:
if str(task_id).isdigit():
idx = int(task_id)
for line_no, line in enumerate(f):
if line_no == idx:
return json.loads(line)
raise ValueError(f"Task case index {idx} not found in {data_path}")
else:
for line in f:
data = json.loads(line)
if data.get("id") == task_id:
return data
raise ValueError(f"Task case id '{task_id}' not found in {data_path}")
def get_tool_prompt(tools):
tool_prompt = "\n\n# Tools\n\nYou may call one or more functions to assist with the user query.\n\nYou are provided with function signatures within <tools></tools> XML tags:\n<tools>"
for tool in tools:
tool_prompt += "\n" + json.dumps(tool)
tool_prompt += "\n</tools>\n\nFor each function call, return a json object with function name and arguments within <tool_call></tool_call> XML tags:\n<tool_call>\n{\"name\": <function-name>, \"arguments\": <args-json-object>}\n</tool_call>"
return tool_prompt
def group_trajectories_by_task_id(jsonl_entries: List[Dict[str, Any]]) -> List[List[Any]]:
"""
group trajectories by task_id
Args:
jsonl_entries: JSONL entry list
Returns:
List[List[Any]]: trajectory list grouped by task_id
"""
grouped = defaultdict(list)
for entry in jsonl_entries:
task_id = entry.get("task_id", "")
taks_case = load_task_case("data/multiturn_data_base.jsonl", task_id)
tools = taks_case.get("tools", [{}])
from bfcl_utils import extract_tool_schema
tool_schema = extract_tool_schema(tools)
entry["task_history"][0]["content"] += get_tool_prompt(tool_schema)
grouped[task_id].append(entry)
# retain only the two with the highest and lowest rewards
filtered_groups = []
for key, trajectories in grouped.items():
if len(trajectories) == 1:
# when only one trajectory, retain it
filtered_groups.append(trajectories)
elif len(trajectories) == 2:
# when there are two trajectories, retain them
filtered_groups.append(trajectories)
else:
# when there are more than two trajectories, choose the two with the highest and lowest rewards
trajectories.sort(key=lambda t: t["reward"])
min_reward_traj = trajectories[0] # highest reward
max_reward_traj = trajectories[-1] # lowest reward
filtered_groups.append([min_reward_traj, max_reward_traj])
return filtered_groups
def post_to_summarizer(trajectories: List[Any], service_url: str, workspace_id: str) -> Dict[str, Any]:
trajectory_dicts = [{
"task_id": traj["task_id"],
"messages": traj["task_history"],
"score": traj["reward"]
} for traj in trajectories]
request_data = {
"trajectories": trajectory_dicts,
"workspace_id": workspace_id
}
try:
response = requests.post(f"{service_url}/summary_task_memory", json=request_data)
response.raise_for_status()
return response.json()
except Exception as e:
return {"error": str(e), "trajectories_count": len(trajectories)}
def process_trajectories_with_threads(grouped_trajectories: List[List[Any]],
service_url: str,
workspace_id: str,
n_threads: int = 4) -> List[Dict[str, Any]]:
"""
use threads to process trajectories
Args:
grouped_trajectories: group trajectory list by task_id
service_url: memory summarizer service URL
workspace_id: workspace ID
n_threads: number of threads
Returns:
all results
"""
results = []
with ThreadPoolExecutor(max_workers=n_threads) as executor:
future_to_group = {
executor.submit(post_to_summarizer, group, service_url, workspace_id): i
for i, group in enumerate(grouped_trajectories)
}
for future in as_completed(future_to_group):
group_index = future_to_group[future]
try:
result = future.result()
result["group_index"] = group_index
result["group_size"] = len(grouped_trajectories[group_index])
results.append(result)
print(f"✅ Group {group_index} processed: {result["metadata"].get('memory_list', 0) if 'memory_list' in result["metadata"] else 'error'}")
except Exception as e:
error_result = {
"group_index": group_index,
"group_size": len(grouped_trajectories[group_index]),
"error": str(e)
}
results.append(error_result)
print(f"❌ Group {group_index} failed: {e}")
return results
def main():
parser = argparse.ArgumentParser(description='Convert JSONL to memories using ReMe service')
parser.add_argument('--jsonl_file', type=str, required=True, help='Path to the JSONL file')
parser.add_argument('--service_url', type=str, default='http://localhost:8001', help='Reme service URL')
parser.add_argument('--workspace_id', type=str, required=True, help='Workspace ID for the task memory pool')
parser.add_argument('--output_file', type=str, help='Output file to save results (optional)')
parser.add_argument('--n_threads', type=int, default=4, help='Number of threads for processing')
args = parser.parse_args()
print(f"Processing JSONL file: {args.jsonl_file}")
print(f"Service URL: {args.service_url}")
print(f"Workspace ID: {args.workspace_id}")
print(f"Threads: {args.n_threads}")
with open(args.jsonl_file, "r") as f:
data = [json.loads(line) for line in f]
print(f"Loaded {len(data)} entries from JSONL file")
grouped_trajectories = group_trajectories_by_task_id(data)
print(f"Total groups: {len(grouped_trajectories)}")
results = process_trajectories_with_threads(
grouped_trajectories,
args.service_url,
args.workspace_id,
n_threads=args.n_threads
)
print(f"Processed {len(results)} groups")
success_count = sum(1 for r in results if 'error' not in r)
error_count = len(results) - success_count
total_memories = sum(len(r["metadata"].get('memory_list', [])) for r in results if 'memory_list' in r["metadata"])
print(f"✅ Success: {success_count}")
print(f"❌ Errors: {error_count}")
print(f"📊 Total task memories created: {total_memories}")
if args.output_file:
try:
summary = {
"workspace_id": args.workspace_id,
"jsonl_file": args.jsonl_file,
"total_groups": len(grouped_trajectories),
"success_count": success_count,
"error_count": error_count,
"total_task_memories": total_memories,
"results": results
}
with open(args.output_file, 'w') as f:
json.dump(summary, f, indent=2)
print(f"Results saved to: {args.output_file}")
except Exception as e:
print(f"Error saving results: {e}")
if __name__ == "__main__":
import sys
if len(sys.argv) > 1:
main()
else:
print("Running in compatibility mode...")
with open("exp_result/qwen3-14b/no_think/bfcl-multi-turn-base_wo-exp.jsonl", "r") as f:
data = [json.loads(line) for line in f]
grouped_trajectories = group_trajectories_by_task_id(data)
print(f"Total groups: {len(grouped_trajectories)}")
results = process_trajectories_with_threads(
grouped_trajectories,
"http://localhost:8001",
"bfcl_test",
n_threads=4
)
print(f"Processed {len(results)} groups")

View file

@ -1,21 +1,27 @@
import json
with open("../../file_vector_store/bfcl_train50_extract_compare_validate.jsonl", 'r') as f:
appworld = [json.loads(line) for line in f]
with open("../../file_vector_store/bfcl_test.jsonl", 'r') as f:
bfcl = [json.loads(line) for line in f]
new_appworld = []
for exp in appworld:
new_bfcl = []
for exp in bfcl:
new_exp = {}
new_exp["workspace_id"] = exp["workspace_id"]
new_exp["experience_id"] = exp["unique_id"]
new_exp["experience_type"] = "text"
new_exp["memory_id"] = exp["unique_id"]
new_exp["memory_type"] = exp["metadata"]["memory_type"]
new_exp["when_to_use"] = exp["content"]
new_exp["content"] = exp["metadata"]["experience_content"]
new_exp["content"] = exp["metadata"]["content"]
new_exp["score"] = exp["metadata"]["score"]
new_exp["time_created"] = exp["metadata"]["time_created"]
new_exp["time_modified"] = exp["metadata"]["time_modified"]
new_exp["author"] = exp["metadata"]["author"]
new_exp["metadata"]= exp["metadata"]["metadata"]
new_appworld.append(new_exp)
new_bfcl.append(new_exp)
with open('../../library/bfcl_train50_extract_compare_validate.jsonl', 'w', encoding='utf-8') as f:
f.writelines(json.dumps(item, ensure_ascii=False) + '\n' for item in new_appworld)
with open('../../library/bfcl_test.jsonl', 'w', encoding='utf-8') as f:
f.writelines(json.dumps(item, ensure_ascii=False) + '\n' for item in new_bfcl)

View file

@ -1,6 +1,6 @@
# BFCL Experiment Quick Start Guide
This guide helps you quickly set up and run BFCL experiments with ExperienceMaker integration.
This guide helps you quickly set up and run BFCL experiments with ReMe integration.
## Env Setup
@ -32,7 +32,7 @@ cp -r bfcl_eval/data {/path/to/bfcl/data}
### 2. Collect agent trajectories on training data set
Run the main experiment script to collect agent trajectories on training data set without experience(`use_experience=False`):
Run the main experiment script to collect agent trajectories on training data set without task memory(`use_memory=False`):
```bash
python run_bfcl.py
@ -40,29 +40,30 @@ python run_bfcl.py
**Note**:
- `max_workers`: Number of parallel workers (default: `4`)
- `num_runs`: Number of times each task is repeated (default: `4`)
- `num_runs`: Number of times each task is repeated (default: `1`)
- `model_name`: LLM model name (default: `qwen3-8b`)
- `enable_thinking`: Control the model's thinking mode (default: `False`)
- `data_path`: Path to the training dataset (default: `./data/multiturn_data_base_train.jsonl`)
- `answer_path`: Path to the possible answer, which are used to evaluate the model's output function (default: `./data/possible_answer`)
- Results are automatically saved to `./exp_result/{model_name}/{no_think/with_think}` directory
### 3. Start ExperienceMaker Service and Init the Experience Pool
### 3. Start ReMe Service and Init the task memory pool
After collecting trajectories, Launch the ExperienceMaker service to enable experience library functionality:
After collecting trajectories, Launch the ReMe service to enable memory library functionality:
```bash
# Go back to the project root
cd ../..
experiencemaker \
http_service.port=8001 \
reme \
backend=http \
http.port=8001 \
llm.default.model_name=qwen-max-2025-01-25 \
embedding_model.default.model_name=text-embedding-v4 \
vector_store.default.backend=local_file
vector_store.default.backend=local
```
and then init the experience pool:
and then init the task memory pool:
```bash
python init_exp_pool.py
@ -70,29 +71,29 @@ python init_exp_pool.py
**Configuration options in `init_exp_pool.py`:**
- `jsonl_file`: Path to the collloaded trajectories
- `service_url`: Experience maker service URL (default: `http://localhost:8001`)
- `workspace_id`: Workspace ID for the experience (default: `bfcl_v1`)
- `service_url`: ReMe service URL (default: `http://localhost:8001`)
- `workspace_id`: Workspace ID for the task memory pool (default: `bfcl_test`)
- `n_threads`: Number of threads for processing (default: `4`)
- `output_file`: Output file to save results (optional)
Now you have inited the experience pool using `local_file` backend (start on `http://localhost:8001`). The `local_file_to_library.py` script or use the following `curl` command:
Now you have inited the task memory pool using `local` backend (start on `http://localhost:8001`). The `local_file_to_library.py` script or use the following `curl` command:
```bash
curl -X POST "http://0.0.0.0:8001/vector_store" \
-H "Content-Type: application/json" \
-d '{
"workspace_id": "bfcl_v1",
"workspace_id": "bfcl_test",
"action": "dump",
"path": "./library"
}'
```
can convert the local file to the experience library (default in `./library/bfcl_v1.jsonl`).
can convert the local file to the memory library (default in `./library/bfcl_test.jsonl`).
Next time, you can import this previously exported experience data to populate the new started workspace with existing knowledge:
Next time, you can import this previously exported task memory data to populate the new started workspace with existing knowledge:
```bash
curl -X POST "http://0.0.0.0:8001/vector_store" \
-H "Content-Type: application/json" \
-d '{
"workspace_id": "bfcl_v1",
"workspace_id": "bfcl_test",
"action": "load",
"path": "./library"
}'
@ -101,7 +102,7 @@ curl -X POST "http://0.0.0.0:8001/vector_store" \
### 4. Run Experiments on Validation Set
Run you can compare agent performance on the validation set with experience (`use_experience=True`) and without experience:
Run you can compare agent performance on the validation set with task memory (`use_memory=True`) and without task memory:
```bash
# remember to change the configuration options, e.g., `data_path=./data/multiturn_data_base_val.jsonl`

View file

@ -1,7 +1,8 @@
import os
import time
import ray
from ray import logger
# from ray import logger
from loguru import logger
from dotenv import load_dotenv
@ -21,15 +22,15 @@ def run_agent(dataset_name: str,
model_name: str = "qwen3-8b",
data_path: str = "data/multiturn_data_base_val.jsonl",
answer_path: Path = Path("data/possible_answer"),
use_experience: bool = False,
use_fixed_experience: bool = True,
use_experience_deletion: bool = False,
use_memory: bool = False,
use_memory_addition: bool = True,
use_memory_deletion: bool = False,
delete_freq: int = 10,
freq_threshold: int = 5,
utility_threshold: float = 0.5,
enable_thinking: bool = False,
experience_base_url: str = "http://0.0.0.0:8001/",
experience_workspace_id: str = "bfcl_8b_0725"):
memory_base_url: str = "http://0.0.0.0:8001/",
memory_workspace_id: str = "bfcl_test"):
experiment_name = dataset_name + "_" + experiment_suffix
path: Path = Path(f"./exp_result/{model_name}/with_think" if enable_thinking else f"./exp_result/{model_name}/no_think")
path.mkdir(parents=True, exist_ok=True)
@ -54,15 +55,15 @@ def run_agent(dataset_name: str,
answer_path=answer_path,
model_name=model_name,
num_runs=num_runs,
use_experience=use_experience,
use_fixed_experience=use_fixed_experience,
use_experience_deletion=use_experience_deletion,
use_memory=use_memory,
use_memory_addition=use_memory_addition,
use_memory_deletion=use_memory_deletion,
delete_freq=delete_freq,
freq_threshold=freq_threshold,
utility_threshold=utility_threshold,
enable_thinking=enable_thinking,
experience_base_url=experience_base_url,
experience_workspace_id=experience_workspace_id
memory_base_url=memory_base_url,
memory_workspace_id=memory_workspace_id
)
future = actor.execute.remote()
future_list.append(future)
@ -83,32 +84,32 @@ def run_agent(dataset_name: str,
def main():
max_workers = 4
num_runs = 8
use_experience = False
use_fixed_experience = False
use_experience_deletion = True
experience_base_url = "http://0.0.0.0:8003/"
experience_workspace_id = "bfcl_train50_extract_compare_validate_add_delete_2"
num_runs = 1
use_memory = False
use_memory_addition = False
use_memory_deletion = False
memory_base_url = "http://0.0.0.0:8001/"
memory_workspace_id = "bfcl_test"
if max_workers > 1:
ray.init(num_cpus=max_workers)
for run_id in range(num_runs):
run_agent(
dataset_name="bfcl-multi-turn-base",
experiment_suffix=f"wo-exp",
model_name="qwen3-14b",
model_name="qwen3-8b",
max_workers=max_workers,
num_runs=1,
data_path="data/multiturn_data_base_train.jsonl",
data_path="data/multiturn_data_base_val.jsonl",
answer_path=Path("data/possible_answer"),
enable_thinking=True,
use_experience=use_experience,
use_fixed_experience=use_fixed_experience,
use_experience_deletion=use_experience_deletion,
enable_thinking=False,
use_memory=use_memory,
use_memory_addition=use_memory_addition,
use_memory_deletion=use_memory_deletion,
delete_freq=5,
freq_threshold=5,
utility_threshold=0.5,
experience_base_url=experience_base_url,
experience_workspace_id=experience_workspace_id,
memory_base_url=memory_base_url,
memory_workspace_id=memory_workspace_id,
)

View file

@ -61,7 +61,7 @@ def get_possible_k_values(total_runs: int) -> list:
def run_exp_statistic():
path: Path = Path(f"./exp_result/qwen-max-latest/no_think")
path: Path = Path(f"./exp_result/qwen3-8b/no_think")
# Store results for all experiments
all_results = {}