This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/camel-performance-tests.git
commit 47f3605dcb46c22d64fa7ce550bb4ee23edb479e Author: Claus Ibsen <[email protected]> AuthorDate: Tue Oct 6 17:08:51 2026 +0200 A Kamelet side check, a bare mode, and a failed reload no longer passes a step gen_kamelet.py and run-kamelet-check.sh write the same tasks with Kamelets and with components, plus a custom Kamelet the model writes itself. BENCH_BARE=1 offers the file tools only. A step that ends on a failed reload fails: its WARN quotes the route, so a log check on a value matched the failure itself. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Claude-Session: https://claude.ai/code/session_01STT6whBgK1AqsSsUKrnE8m --- ai-benchmark/agent_mcp_stepwise.py | 26 +++- ai-benchmark/gen_kamelet.py | 308 +++++++++++++++++++++++++++++++++++++ ai-benchmark/run-kamelet-check.sh | 42 +++++ ai-benchmark/summarize_stepwise.py | 3 +- 4 files changed, 375 insertions(+), 4 deletions(-) diff --git a/ai-benchmark/agent_mcp_stepwise.py b/ai-benchmark/agent_mcp_stepwise.py index a821f48..3c794b5 100755 --- a/ai-benchmark/agent_mcp_stepwise.py +++ b/ai-benchmark/agent_mcp_stepwise.py @@ -35,6 +35,9 @@ REFERENCE = os.environ.get("BENCH_REFERENCE") == "1" # CAMEL-24834: before each step, ask camel_runtime_tool_groups what the app has and offer its groups' tools and guidance, # as the AI panel of the camel-jbang views does for a local model (the tool list changes only when the fingerprint does) TOOL_GROUPS = os.environ.get("BENCH_TOOL_GROUPS") == "1" +# the Kamelet side check: bare offers the file tools only, no catalog, validation or runtime feedback (a write of +# invalid YAML is still refused); what the model knows by itself +BARE = os.environ.get("BENCH_BARE") == "1" TOOL_RESULT_CAP = 6000 SHARED = ["camel_catalog_doc", "camel_catalog_find", "camel_catalog_sample", "camel_validate_source", "camel_get_files", "camel_write_file", @@ -44,6 +47,8 @@ EXTRA = ["camel_catalog_docs", "camel_component_properties", "camel_configuratio # CAMEL-24834 tool-group experiment: BENCH_EXTRA_TOOLS=camel_execute_sql,camel_get_datasources adds a group to the set EXTRA += [t for t in os.environ.get("BENCH_EXTRA_TOOLS", "").split(",") if t] +BARE_TOOLS = ["camel_get_files", "camel_write_file", "camel_edit_file"] + SYSTEM = ( "You are an Apache Camel assistant helping a developer edit a running Camel integration through the Camel MCP server.\n\n" "The project directory is {directory}. The integration {name} is already running from it in dev mode: files you write are " @@ -291,7 +296,10 @@ def score(step, project, cfg, mcp, name, before): def lvl(l): return (l.get("level") or "").upper() # a multi-line message (a pretty-printed body) comes as one record with a detail block: match on both def msg(l): return (l.get("message") or l.get("msg") or "") + ("\n" + l["detail"] if l.get("detail") else "") - recent = [l for l in lines if isinstance(l, dict) and lvl(l) in ("INFO", "WARN")] + # a failed reload is a WARN that quotes the route (the endpoint URI with its values), so a log_regex on a value + # matched the failure itself; those records are not what the route logged (Kamelet side check, kb-1) + def reload_failure(l): return any(t in msg(l) for t in ("Error reloading routes", "Reload failed", "Failed to create route")) + recent = [l for l in lines if isinstance(l, dict) and lvl(l) in ("INFO", "WARN") and not reload_failure(l)] msgs = [msg(l) for l in recent] # round 2: log_regex may be a list (all must match), or absent (a step with nothing to see in the log) regexes = chk.get("log_regex") @@ -321,6 +329,13 @@ def score(step, project, cfg, mcp, name, before): result["log_ok"] = False except Exception: pass + # the step ends on a failed reload: the app runs the routes of before (or none), not what the model wrote + fresh_reload_recs = [l for l in lines if isinstance(l, dict) and is_reload(l) and reload_key(l) not in seen_reloads] + if fresh_reload_recs: + last_reload = max(fresh_reload_recs, key=lambda l: tod(l.get("time") or l.get("timestamp")) or 0) + if "Error reloading routes" in msg(last_reload): + result["reload_failed"] = True + result["log_ok"] = False # round 2: errors are counted per step (the log window and the error file keep an earlier step's stack traces) seen = before.get("error_records") or set() result["errors"] = sum(1 for l in lines if isinstance(l, dict) and lvl(l) == "ERROR" and error_key(l) not in seen) @@ -367,7 +382,7 @@ def main(): os.makedirs(OUT, exist_ok=True) log = open(os.path.join(OUT, "run.log"), "a") mcp = McpClient(MCP_URL); mcp.initialize(); all_tools = mcp.list_tools() - tools = to_ollama_tools(all_tools, SHARED + EXTRA) + tools = to_ollama_tools(all_tools, BARE_TOOLS if BARE else SHARED + EXTRA) print(f"tools: {[t['function']['name'] for t in tools]}", file=log, flush=True) proc = None @@ -418,13 +433,18 @@ def main(): time.sleep(6) base_system = SYSTEM.replace("{directory}", project).replace("{name}", name) + if BARE: + # the guidelines about tools it does not have would only confuse it + base_system = "\n".join(l for l in base_system.splitlines() + if not any(t in l for t in ("camel_validate_source", "camel_catalog_sample"))) + base_system = base_system.replace(" (camel_catalog_doc has the option names)", "") messages = [{"role": "system", "content": base_system}] groups_fp = None results = [] try: for step in cfg["steps"]: sid = step["id"] - if TOOL_GROUPS and not REFERENCE: + if TOOL_GROUPS and not REFERENCE and not BARE: data, raw = jcall(mcp, "camel_runtime_tool_groups", {"nameOrPid": name}) if isinstance(data, dict) and data.get("fingerprint") != groups_fp: groups_fp = data.get("fingerprint") diff --git a/ai-benchmark/gen_kamelet.py b/ai-benchmark/gen_kamelet.py new file mode 100644 index 0000000..025377a --- /dev/null +++ b/ai-benchmark/gen_kamelet.py @@ -0,0 +1,308 @@ +#!/usr/bin/env python3 +"""Builds the Kamelet side check: steps-kamelet/<name>.json and stepwise-kamelet/<name>/. + +Do models know Kamelets as well as components? Each task is written twice with the same requests: once with Kamelets +(kamelet-*) and once with plain components (component-*), so a gap between the two is about Kamelets, not the task. +Run it with tools (the Camel MCP server, as the ladder) and bare (BENCH_BARE=1: file tools only, no catalog, no +runtime feedback). Not part of the ladder or its baselines. + +Usage: gen_kamelet.py [<name>] writes every steps file and project, or one example's +""" +import json, os, shutil, sys + +HERE = os.path.dirname(os.path.abspath(__file__)) +STEPS_DIR = os.path.join(HERE, "steps-kamelet") +PROJECTS = os.path.join(HERE, "stepwise-kamelet") + +ROUTE = "orders.camel.yaml" +INITIAL = {ROUTE: """- route: + id: orders + from: + uri: timer + parameters: + timerName: tick + period: 5000 + steps: + - log: + message: "Waiting for orders" +""", "application.properties": ""} + +PAID = '{"orderId":"ORD-1","status":"paid","amount":42}' +PENDING = '{"orderId":"ORD-2","status":"pending","amount":42}' + +# ---------------------------------------------------------------- orders pipeline: source, extract a field, filter +K_SOURCE = """ uri: kamelet:timer-source + parameters: + period: 2000 + message: '%s' + contentType: application/json +""" +K_LOG = """ - to: + uri: kamelet:log-sink +""" +K_EXTRACT = """ - to: + uri: kamelet:extract-field-action + parameters: + field: orderId +""" +K_FILTER = """ - to: + uri: kamelet:predicate-filter-action + parameters: + expression: "@.status == 'paid'" +""" +C_SOURCE = """ uri: timer + parameters: + timerName: orders + period: 2000 +""" +C_BODY = """ - setBody: + constant: '%s' +""" +C_LOG = """ - to: + uri: log:orders +""" +C_EXTRACT = """ - setBody: + jsonpath: "$.orderId" +""" +C_FILTER_OPEN = """ - filter: + jsonpath: "$[?(@.status == 'paid')]" + steps: +""" + + +def orders(src, steps): + return "- route:\n id: orders\n from:\n" + src + " steps:\n" + steps + + +def indent(s, n): + return "".join(" " * n + l + "\n" if l else "\n" for l in s.splitlines()) + + +PIPE_STEPS = [ + {"k": "Replace the route's source with the timer-source Kamelet: every 2 seconds it emits the message %s as JSON " + "(content type application/json). Send each message to the log-sink Kamelet." % PAID, + "c": "Replace the route's source with the timer component: every 2 seconds set the body to the message %s. " + "Send each message to the log component (logger name orders)." % PAID, + "check": {"log_regex": "ORD-1"}}, + {"k": "Before the log-sink, use the extract-field-action Kamelet so only the orderId value is logged, not the whole " + "order.", + "c": "Before the log, set the body to only the orderId value of the JSON message, so only that is logged, not the " + "whole order.", + "check": {"log_regex": "ORD-1", "log_not_regex": "status"}}, + {"k": "Change the message to %s and add the predicate-filter-action Kamelet first, so only orders with status paid " + "go on to the extract and the log." % PENDING, + "c": "Change the message to %s and add a filter first, so only orders with status paid go on to the extract and " + "the log." % PENDING, + "check": {"log_not_regex": "ORD-2"}}, +] +PIPE_K = [ + orders(K_SOURCE % PAID, K_LOG), + orders(K_SOURCE % PAID, K_EXTRACT + K_LOG), + orders(K_SOURCE % PENDING, K_FILTER + K_EXTRACT + K_LOG), +] +PIPE_C = [ + orders(C_SOURCE, C_BODY % PAID + C_LOG), + orders(C_SOURCE, C_BODY % PAID + C_EXTRACT + C_LOG), + orders(C_SOURCE, C_BODY % PENDING + C_FILTER_OPEN + indent(C_EXTRACT + C_LOG, 6)), +] +PIPE_FILE = [{"k": "timer-source", "c": r"uri: timer\b"}, + {"k": "extract-field-action", "c": r'orderId(?!":)'}, + {"k": "predicate-filter-action", "c": "filter"}] + +# ---------------------------------------------------------------- kafka round trip +MSG = "Order ORD-7 placed" +K_KAFKA_OUT = """- route: + id: orders + from: + uri: kamelet:timer-source + parameters: + period: 3000 + message: "%s" + steps: + - to: + uri: kamelet:kafka-sink + parameters: + topic: orders + bootstrapServers: localhost:9092 +""" % MSG +K_KAFKA_IN = """- route: + id: received + from: + uri: kamelet:kafka-source + parameters: + topic: orders + bootstrapServers: localhost:9092 + steps: + - to: + uri: kamelet:log-sink +""" +C_KAFKA_OUT = """- route: + id: orders + from: + uri: timer + parameters: + timerName: orders + period: 3000 + steps: + - setBody: + constant: "%s" + - to: + uri: kafka + parameters: + topic: orders + brokers: localhost:9092 +""" % MSG +C_KAFKA_IN = """- route: + id: received + from: + uri: kafka + parameters: + topic: orders + brokers: localhost:9092 + steps: + - to: + uri: log:received +""" +KAFKA_STEPS = [ + {"k": "Change the route so it sends the message \"%s\" every 3 seconds to the Kafka topic orders on localhost:9092 " + "(the broker has no authentication), using the timer-source and kafka-sink Kamelets." % MSG, + "c": "Change the route so it sends the message \"%s\" every 3 seconds to the Kafka topic orders on localhost:9092 " + "(the broker has no authentication), using the timer and kafka components." % MSG, + "check": {}}, + {"k": "Add a second route that consumes the topic orders with the kafka-source Kamelet and sends each message to the " + "log-sink Kamelet.", + "c": "Add a second route that consumes the topic orders with the kafka component and sends each message to the log " + "component (logger name received).", + "check": {"log_regex": MSG}}, +] +KAFKA_K = [K_KAFKA_OUT, K_KAFKA_OUT + K_KAFKA_IN] +KAFKA_C = [C_KAFKA_OUT, C_KAFKA_OUT + C_KAFKA_IN] +KAFKA_FILE = [{"k": "kafka-sink", "c": r"uri: kafka\b|kafka:"}, + {"k": "kafka-source", "c": r"(?s)kafka.*kafka"}] + + +# ---------------------------------------------------------------- a custom Kamelet: a business building block +# Not a twin: the model writes the Kamelet itself (the definition with its properties, and the template), then uses +# it and changes it. What custom Kamelets are for: a block of the business, not a general purpose one like kafka. +TAG_FILE = "tag-order-action.kamelet.yaml" + + +def tag_kamelet(prefix): + props = """ tag: + title: Tag + description: The tag to append to the order + type: string +""" + expr = "${body} [{{tag}}]" + if prefix: + props += """ prefix: + title: Prefix + description: Goes before the tag + type: string + default: "#" +""" + expr = "${body} [{{prefix}}{{tag}}]" + return """apiVersion: camel.apache.org/v1 +kind: Kamelet +metadata: + name: tag-order-action + labels: + camel.apache.org/kamelet.type: action +spec: + definition: + title: Tag Order + description: Appends a tag to the order in the message body + required: + - tag + type: object + properties: +""" + props + """ template: + from: + uri: kamelet:source + steps: + - setBody: + expression: + simple: + expression: "%s" +""" % expr + + +TAG_ROUTE = """- route: + id: orders + from: + uri: timer + parameters: + timerName: orders + period: 2000 + steps: + - setBody: + expression: + constant: + expression: Order ORD-5 + - to: + uri: kamelet:tag-order-action + parameters: + tag: priority + - log: + message: "Tagged: ${body}" +""" +CUSTOM = {"name": "kamelet-custom", "steps": [ + {"id": 1, + "request": "Write a custom action Kamelet named tag-order-action in the file %s in the project. It has one " + "required property tag and appends \" [<tag>]\" to the message body. Then change the route: every 2 " + "seconds set the body to \"Order ORD-5\", send it through tag-order-action with tag priority, and log " + "\"Tagged: <body>\"." % TAG_FILE, + "check": {"file_regex": "tag-order-action", "files": {TAG_FILE: "kind:\\s*Kamelet"}, + "log_regex": "Tagged: Order ORD-5 \\[priority\\]"}, + "reference": {TAG_FILE: tag_kamelet(False), ROUTE: TAG_ROUTE}}, + {"id": 2, + "request": "Give tag-order-action a second, optional property prefix with the default value #, which goes before " + "the tag, so the route logs \"Tagged: Order ORD-5 [#priority]\" without changing the route.", + "check": {"files": {TAG_FILE: "prefix"}, "log_regex": "Tagged: Order ORD-5 \\[#priority\\]"}, + "reference": {TAG_FILE: tag_kamelet(True)}}, +]} + + +def example(name, kind, steps, refs, files, infra=None): + out = [] + for i, (s, ref, f) in enumerate(zip(steps, refs, files), 1): + chk = dict(s["check"]); chk["file_regex"] = f[kind] + out.append({"id": i, "request": s[kind], "check": chk, "reference": {ROUTE: ref}}) + ex = {"name": name, "steps": out} + if infra: + ex["infra"] = infra + return ex + + +EXAMPLES = [ + example("kamelet-orders", "k", PIPE_STEPS, PIPE_K, PIPE_FILE), + example("component-orders", "c", PIPE_STEPS, PIPE_C, PIPE_FILE), + example("kamelet-kafka", "k", KAFKA_STEPS, KAFKA_K, KAFKA_FILE, ["kafka"]), + example("component-kafka", "c", KAFKA_STEPS, KAFKA_C, KAFKA_FILE, ["kafka"]), + CUSTOM, +] + + +def main(): + os.makedirs(STEPS_DIR, exist_ok=True) + only = sys.argv[1] if len(sys.argv) > 1 else None + for ex in EXAMPLES: + if only and ex["name"] != only: + continue + cfg = {"project": "stepwise-kamelet/" + ex["name"], "route_file": ROUTE, "props_file": "application.properties", + "wait_seconds": 8, "source_dir": True, "initial": INITIAL, "steps": ex["steps"]} + if ex.get("infra"): + cfg["infra"] = ex["infra"] + with open(os.path.join(STEPS_DIR, ex["name"] + ".json"), "w") as f: + json.dump(cfg, f, indent=1) + d = os.path.join(PROJECTS, ex["name"]) + shutil.rmtree(d, ignore_errors=True) + os.makedirs(d) + for fname, content in INITIAL.items(): + with open(os.path.join(d, fname), "w") as f: + f.write(content) + print(f"{ex['name']}: {len(ex['steps'])} steps, project stepwise-kamelet/{ex['name']}") + + +if __name__ == "__main__": + main() diff --git a/ai-benchmark/run-kamelet-check.sh b/ai-benchmark/run-kamelet-check.sh new file mode 100755 index 0000000..bf15339 --- /dev/null +++ b/ai-benchmark/run-kamelet-check.sh @@ -0,0 +1,42 @@ +#!/bin/zsh +# Runs the Kamelet side check: every steps-kamelet/<name>.json, each from a fresh project, k times. +# Usage: run-kamelet-check.sh <tag> [k] Results: stepwise/<tag>-<i>/<name>/ (results.json, run.log, traces) +# BENCH_BARE=1 offers the model the file tools only (see gen_kamelet.py). +# BENCH_REFERENCE=1 applies the reference files instead of asking the model (the reference pass of the steps files). +# BENCH_ONLY="<name> [<name> ...]" runs only those examples. +set -u +TAG="${1:?usage: run-kamelet-check.sh <tag> [k]}" +K="${2:-1}" +cd "$(dirname "$0")" +export MCP_URL="${MCP_URL:-http://127.0.0.1:9090/mcp}" +command -v caffeinate > /dev/null && caffeinate -i -s -w $$ & +for i in $(seq 1 "$K"); do + if (( K > 1 )); then T="$TAG-$i"; else T="$TAG"; fi + : > "stepwise-$T.log" # a fresh log per run, so a later wait on its DONE line cannot see an earlier run's + for f in steps-kamelet/*.json; do + name=$(basename "$f" .json) + if [[ -n "${BENCH_ONLY:-}" && " $BENCH_ONLY " != *" $name "* ]]; then continue; fi + python3 gen_kamelet.py "$name" > /dev/null + echo "[$T] $name start $(date +%T)" | tee -a "stepwise-$T.log" + BENCH_STEPS="$f" BENCH_TAG="$T/$name" python3 agent_mcp_stepwise.py > "stepwise-$T-$name.out" 2>&1 + tail -1 "stepwise/$T/$name/run.log" | tee -a "stepwise-$T.log" + # safety net: the harness stops its integration, but an interrupted run leaves one behind, and an app the model + # started itself with camel_run, or renamed, is not stopped by name (r5: an orphan held port 8080 for six hours + # and failed every later HTTP rung). The machine runs only the benchmark during a series, so stop them all. + camel stop > /dev/null 2>&1 || true + sleep 3 + if lsof -nP -iTCP:8080 -sTCP:LISTEN > /dev/null 2>&1; then + echo "[$T] WARNING port 8080 still in use after $name: $(lsof -nP -iTCP:8080 -sTCP:LISTEN | tail -1)" | tee -a "stepwise-$T.log" + fi + done + # the services an example needed stay up across its passes; stop them once the run is over + for svc in $(python3 -c " +import glob, json +out = set() +for f in glob.glob('steps-kamelet/*.json'): + try: out.update(json.load(open(f)).get('infra') or []) + except Exception: pass +print(' '.join(sorted(out)))"); do camel infra stop "$svc" > /dev/null 2>&1 || true; done + echo "[$T] DONE $(date +%T)" | tee -a "stepwise-$T.log" +done +BENCH_STEPS_DIR=steps-kamelet python3 summarize_stepwise.py "$TAG" "$K" diff --git a/ai-benchmark/summarize_stepwise.py b/ai-benchmark/summarize_stepwise.py index 0ba1a72..7ddde33 100644 --- a/ai-benchmark/summarize_stepwise.py +++ b/ai-benchmark/summarize_stepwise.py @@ -8,7 +8,8 @@ errors were logged after the model's last accepted write; runs before it was rec import json, os, sys tag = sys.argv[1]; k = int(sys.argv[2]) if len(sys.argv) > 2 else 1 tags = [f"{tag}-{i}" for i in range(1, k + 1)] if k > 1 else [tag] -names = sorted(os.path.splitext(f)[0] for f in os.listdir("steps-ladder") if f.endswith(".json")) +steps_dir = os.environ.get("BENCH_STEPS_DIR", "steps-ladder") +names = sorted(os.path.splitext(f)[0] for f in os.listdir(steps_dir) if f.endswith(".json")) print(f"| example | steps | passes per step ({' '.join(tags)}) | all steps passed | last version right | final state right |") print("|---|---|---|---|---|---|") total = 0; possible = 0; clean = 0; final = 0; lastw = 0
