-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtask_decompose_a2a_server.py
More file actions
1662 lines (1496 loc) · 79.9 KB
/
Copy pathtask_decompose_a2a_server.py
File metadata and controls
1662 lines (1496 loc) · 79.9 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
#!/usr/bin/env python3
# Copyright (c) 2026 hidemi-k
# Licensed under the MIT License.
"""
task_decompose A2A Hub Server — A2A + FastAPI 統合版 (port:8000)
====================================================================
A2A エンドポイント(エージェント間通信):
POST / JSON-RPC message/send
GET /.well-known/agent.json Agent Card
FastAPI エンドポイント(NiceGUI UI 向け):
GET /healthz 死活監視
POST /execute Dry-run(A2A Hub 経由、deploy=False)
POST /deploy/{trace_id} 実機投入(A2A Hub 経由、deploy=True)
GET /diff/{trace_id} 差分・結果取得
WS /ws/updates ログストリーミング
起動:
python task_decompose_a2a_server.py
環境変数:
A2A_PORT : このサーバのポート(デフォルト: 8000)
NETCONF_A2A_URL: NETCONF サーバURL(デフォルト: http://localhost:8001)
EAPI_A2A_URL : eAPI サーバURL(デフォルト: http://localhost:8002)
HTTP_TIMEOUT : 転送タイムアウト秒(デフォルト: 120)
"""
import asyncio, json, logging, os, re, sys, uuid, configparser
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
import httpx
import uvicorn
# ── 多言語対応 ────────────────────────────────────────────────────────────
from i18n import get_msg, locale_from_request, LOCALE
from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect
from fastapi.middleware.cors import CORSMiddleware
from pydantic import BaseModel
from starlette.routing import Route, WebSocketRoute
# A2A SDK
from a2a.server.agent_execution.agent_executor import AgentExecutor
from a2a.server.agent_execution.context import RequestContext
from a2a.server.events.event_queue_v2 import EventQueue
from a2a.server.request_handlers import DefaultRequestHandler
from a2a.server.tasks.inmemory_task_store import InMemoryTaskStore
from a2a.server.routes.fastapi_routes import add_a2a_routes_to_fastapi
from a2a.server.routes.agent_card_routes import create_agent_card_routes
from a2a.server.routes.jsonrpc_routes import create_jsonrpc_routes
from a2a.server.routes.rest_routes import create_rest_routes
from a2a.types import (
AgentCard, AgentCapabilities, AgentSkill,
)
from a2a.utils.errors import UnsupportedOperationError
# LangChain
from langchain_openai import ChatOpenAI
# ── LLM ファクトリ(Groq Primary / Azure Fallback 共通モジュール) ─────────────
from llm_factory import (
build_llm, build_llm_with_fallback,
invoke_with_fallback, log_llm_config,
LLM_PROVIDER_NAME,
)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
handlers=[logging.StreamHandler(sys.stdout)],
)
logger = logging.getLogger("task_decompose_hub_v2")
# ── 設定 ──────────────────────────────────────────────────────────────────────
VERSION = "1.4.0"
BUILD_DATE = "2026-05-16"
BASE_DIR = os.path.dirname(os.path.abspath(__file__))
CONFIG_PATH = os.getenv("SASE_CONFIG",
os.path.join(BASE_DIR, "./config.ini"))
A2A_HOST = os.getenv("A2A_HOST", "0.0.0.0")
A2A_PORT = int(os.getenv("A2A_PORT", "8000"))
A2A_PUBLIC_URL = os.getenv("A2A_PUBLIC_URL", f"http://localhost:{A2A_PORT}")
NETCONF_A2A_URL = os.getenv("NETCONF_A2A_URL", "http://localhost:8001")
EAPI_A2A_URL = os.getenv("EAPI_A2A_URL", "http://localhost:8002")
EAPI_SDIFF_URL = os.getenv("EAPI_SDIFF_URL", "http://localhost:8009") # session-diff REST
XDP_A2A_URL = os.getenv("XDP_A2A_URL", "http://localhost:8003")
ANTA_A2A_URL = os.getenv("ANTA_A2A_URL", "http://localhost:8004") # ANTA Snapshot 検証
EAPI_CONFIG_A2A_URL = os.getenv("EAPI_CONFIG_A2A_URL", "http://localhost:8006") # eAPI configure session
HTTP_TIMEOUT = float(os.getenv("HTTP_TIMEOUT", "120"))
# LLM インスタンス(Groq Primary / Azure Fallback 自動切り替え)
# classify_query() 等で llm.invoke() の代わりに invoke_with_fallback(llm, prompt) を使う
llm = build_llm_with_fallback()
# ── トレース結果ストア ────────────────────────────────────────────────────────
_trace_store: Dict[str, Dict] = {}
_ws_clients: Dict[str, list] = {}
# ═══════════════════════════════════════════════════════════════════════════════
# クエリ分類(ルーター)
# ═══════════════════════════════════════════════════════════════════════════════
READ_KEYWORDS = [
"確認","表示","一覧","状態","参照","見せ","教え","show","get",
"list","check","status","display","バージョン","version",
"見たい","知りたい","調べたい","確かめたい", # ★ 追加(2026-05-20): 参照意図の口語表現
"acl", # ★ 追加(2026-05-20): show ip access-lists は eAPI 参照コマンド
# ★ 削除(2026-05-20 解決策B): "情報" を除去
# 理由: 「セキュリティの統計情報」「XDP情報」等、あらゆる文脈で使われるため
# 単語単体では read/security を決定できない → LLM fallback に委ねる
# 影響: 「バージョン情報を確認して」等は「確認」「バージョン」が残るため影響なし
# 「セキュリティの統計情報を調べて」はLLM fallbackに到達→正しくsecurityへ
]
WRITE_KEYWORDS = [
"設定","変更","追加","削除","適用","投入","作成","修正",
"configure","set","add","delete","remove","update","apply","create",
]
# ── Security キーワード設計方針 ────────────────────────────────────────────────
# 【必須キーワード】これが含まれる場合のみ security に分類する。
# XDP/ファイアウォール操作に固有の用語に限定し、
# NTP・BGP・インターフェース等の一般デバイス参照と混同しない。
#
# 【削除した旧キーワード】
# "統計","フロー","パケット","stats","flow" — 汎用的すぎる(NTP状態参照等と衝突)
# "情報" — READ_KEYWORDS と重複
# "ips" — "tips" 等の誤マッチリスク
#
# 【移動先】
# "脅威","攻撃","flood","syn flood" → SECURITY_CONTEXT_KEYWORDS(後述)
# これらは単独ではなく、SECURITY_REQUIRED との組み合わせで security 判定する。
# ──────────────────────────────────────────────────────────────────────────────
# XDP/FW に明示的に言及しているキーワード(これだけで security 確定)
# ★ 2026-05-20 修正: READ_OVERRIDE 方式導入(方針B)
# 「BGPのACLを確認して」→ acl が SECURITY_REQUIRED にヒットして security 誤判定
# していた問題を修正。acl/security/qos 等は参照目的でも使われる語なので除去し、
# READ_OVERRIDE_WORDS + READ動詞の AND 条件で read に倒す仕組みを追加。
# XDP固有の操作系ワード(ブロック・遮断・xdp・ebpf等)のみここに残す。
SECURITY_REQUIRED = [
# 日本語: XDP/FW操作を示す明示的な操作系単語
"ブロック", "遮断",
# XDP/eBPF 固有(Arista デバイス設定とは別ドメイン)
"xdp", "ebpf", "firewall",
# ブロック操作の英語表現(drop単体は除外 → "drop list" "drop/block" のみ)
"block", "drop list", "drop/block", "drop/unblock",
"ブロックリスト",
# QoS帯域制限(操作系のみ残す)
# ★ "qos" 単体は除去 → "QoSの設定を確認" が eAPI に正しく流れるように
"帯域制限", "rate limit", "rate-limit",
# ★ 除去済みキーワード(READ_OVERRIDE_WORDS に移動):
# "acl" → show ip access-lists は eAPI 参照
# "security" → セキュリティポリシー確認は eAPI 参照
# "セキュリティ" → 同上
# "qos" → QoS設定確認は eAPI 参照
# "ファイアウォール" → FWルール表示は eAPI 参照の場合がある
]
# READ優先上書きワード:
# SECURITY_REQUIRED にヒットしても、このワード + READ動詞がある場合は read に倒す。
# 「BGPのACLを確認して」「QoSの設定を見せて」「セキュリティポリシーを表示」等が対象。
# XDP固有ワード(xdp/block/遮断)は SECURITY_REQUIRED に残してあるので上書きしない。
READ_OVERRIDE_WORDS = [
"acl", # show ip access-lists → eAPI 参照
"security", # セキュリティポリシーを確認 → eAPI 参照
"セキュリティ", # 同上(日本語)
"qos", # QoSの設定を確認 → eAPI 参照
"ファイアウォール", # ファイアウォールのルールを表示 → eAPI 参照
]
# 脅威・攻撃系キーワード(SECURITY_REQUIRED との AND 条件で security 判定)
# 単独では read に倒す(例: "攻撃を調べて" → eAPI で確認、"攻撃をブロック" → security)
SECURITY_CONTEXT_KEYWORDS = [
"脅威", "攻撃", "flood", "syn flood", "syn-flood",
"統計", "フロー", "パケット", "stats", "flow",
]
# ANTA Snapshot 検証キーワード(security/read/write より先にチェック)
VERIFY_KEYWORDS = [
"anta","snapshot","スナップショット","事後","post_check","post-check",
"post check","verify","検証","事後検証","副作用","影響確認",
"anta テスト","anta test","ネットワーク検証","network verify",
]
def classify_query(query: str) -> str:
"""
クエリを eapi_config / verify / security / write / read の5種に分類する。
判定ロジック:
0. VXLAN/EVPN × 設定変更 → eapi_config(NETCONF非対応のため最優先)
1. VERIFY_KEYWORDS → verify (ANTA検証)
2. SECURITY_REQUIRED → security(XDP/FW操作の明示キーワード)
3. READ + SECURITY_CONTEXT の両方 → read を優先
(例: "フローの状態を確認" → read, "フローをブロック" → security)
4. READ / WRITE キーワードで判定
5. いずれも該当しない → LLMフォールバック(改善プロンプト)
"""
q = query.lower()
# 0. VXLAN/EVPN × 設定変更 → eapi_config(最優先)
# NETCONF/OpenConfig では設定不可のため eAPI configure session で処理する
#
# ★ 参照系動詞(READ_KEYWORDS)が含まれる場合は eapi_config へ送らず read に倒す。
# 例: "VXLANの設定を確認して" → "設定"(write_verb) + "確認"(read_verb)
# → 参照系が優先 → eapi_show へ
# 例: "VXLAN VNI 100 に vni 10000 を設定して" → write_verb のみ → eapi_config
_has_vxlan_evpn = any(k in q for k in [
"vxlan", "evpn", "vni", "vtep", "route-target",
"ルートターゲット", "mac flooding", "mac フラッディング",
"address-family evpn", "evpn アドレスファミリー",
])
# ★ BGP network advertise / redistribute も eapi_config へ
# NETCONF の BGP YANG は複数ツリーをまたぐ複雑な操作のため CLI 方式が確実
_has_bgp_cli = (
any(k in q for k in ["bgp", "ルータ bgp", "router bgp"])
and any(k in q for k in [
"network ", "redistribute", "advertise", "アドバタイズ",
"bgp network", "bgp advertise",
])
)
_has_write_verb = any(k in q for k in [
"設定", "変更", "追加", "削除", "適用", "投入", "作成", "修正",
"configure", "set", "add", "delete", "remove", "update", "apply", "create",
])
_has_read_verb = any(k in q for k in READ_KEYWORDS)
if (_has_vxlan_evpn and _has_write_verb and not _has_read_verb):
return "eapi_config"
# BGP network/redistribute/advertise は write_verb なしでも eapi_config へ
# (「redistribute connected して」の「して」は write_verb にヒットしないため)
if _has_bgp_cli and not _has_read_verb:
return "eapi_config"
# VXLAN/EVPN + 参照系動詞 → read(設定変更ではなく状態確認)
# 例: "VXLANの設定を確認して" → read(eapi_show へ)
# ※ mixed に落ちないよう、ここで明示的に read を返す
if _has_vxlan_evpn and _has_read_verb:
return "read"
# 1. ANTA 検証を最優先
if any(k in q for k in VERIFY_KEYWORDS):
return "verify"
# 2. XDP/FW 明示キーワードがあれば security 確定
# ★ READ_OVERRIDE_WORDS + READ動詞がある場合は read を優先(方針B)
# 例: "BGPのACLを確認して" → acl がヒットするが "確認" もある → read
# 例: "10.0.1.30をブロックして" → ブロックがヒット、READ_OVERRIDE にない → security
#
# ★ ただし XDP固有ワード(xdp/ebpf/firewall)が含まれる場合は
# READ_OVERRIDE を無効にする(2026-05-21 修正)
# 例: "XDPでQoSの設定を確認して"
# → xdp(SECURITY_REQUIRED)+ qos(READ_OVERRIDE)+ 確認(READ)
# → xdp が明示されているので override 無効 → security
# 例: "QoSの設定を確認して"
# → qos(READ_OVERRIDE)+ 確認(READ)、xdp なし → override 有効 → read
_XDP_EXPLICIT = ["xdp", "ebpf", "firewall", "ファイアウォール"]
if any(k in q for k in SECURITY_REQUIRED):
has_xdp_explicit = any(k in q for k in _XDP_EXPLICIT)
override = any(k in q for k in READ_OVERRIDE_WORDS)
has_read_verb = any(k in q for k in READ_KEYWORDS)
# XDP固有ワードが明示されている場合は override を無効化
if override and has_read_verb and not has_xdp_explicit:
return "read"
return "security"
# 3. セキュリティ系曖昧語 × XDP文脈語 の AND → security(step3より前に判定)
# ★ 解決策B (2026-05-20): 「セキュリティの統計を見せて」等を正しく security へ
# 「セキュリティ」単体は曖昧(デバイス設定参照でも使われる)だが、
# 「統計/アラート/異常/監視/フロー」と組み合わさると XDP 文脈と判断できる。
# これを READ_KEYWORDS チェックより先に評価することで、
# 「見せ」「確認」等の READ 動詞があっても security に正しく流れる。
_sec_ambiguous = ["セキュリティ", "security"]
_sec_ctx_combine = ["統計", "stats", "アラート", "alert", "異常", "監視", "monitor",
"脅威", "フロー", "flow", "攻撃", "attack"]
if (any(k in q for k in _sec_ambiguous)
and any(k in q for k in _sec_ctx_combine)):
return "security"
has_read = any(k in q for k in READ_KEYWORDS)
has_write = any(k in q for k in WRITE_KEYWORDS)
if has_read and not has_write: return "read"
if has_write and not has_read: return "write"
if has_read and has_write: return "mixed" # 参照+変更の混在 → UI で警告
# 4. 脅威系のみ(READ/WRITE どちらもなし)→ security
if any(k in q for k in SECURITY_CONTEXT_KEYWORDS):
return "security"
# 5. LLM フォールバック(改善プロンプト)
# ★ 解決策B (2026-05-20): 改善プロンプト
# 旧プロンプトの問題: security の定義が「操作系のみ」だったため
# 「セキュリティの統計情報を調べて」等の「参照系 XDP クエリ」を
# LLM が read と誤判定していた。
# 改善: security に「XDP に関するあらゆる参照・分析・統計」も含むと明示。
# また「セキュリティ+統計/監視/アラート → security」の判定例を追加。
result = invoke_with_fallback(
llm,
"あなたは Arista ネットワーク管理システムのルーター AI です。\n"
"以下のクエリを 5 種類のルートのうち 1 つに分類し、一単語のみで回答してください。\n\n"
"【ルート定義】\n"
" eapi_config : VXLAN/EVPN の設定変更(NETCONF/OpenConfig 非対応)\n"
" 例: VXLAN VNI 設定 / EVPN RD/RT 設定 / VTEP 設定 /\n"
" address-family evpn 有効化 / MAC フラッディング追加\n"
" read : Arista cEOS デバイスの状態参照(show コマンド相当)\n"
" 例: BGP 状態確認 / インターフェース確認 / ACL 表示 /\n"
" ルーティングテーブル / NTP 状態 / VLAN 一覧\n"
" write : Arista cEOS デバイスへの設定変更(VXLAN/EVPN 以外)\n"
" 例: VLAN 作成 / インターフェース設定 / BGP neighbor 追加\n"
" security : XDP/eBPF ファイアウォールに関するあらゆる操作・参照・分析\n"
" 例: IP ブロック / 遮断 / QoS 帯域制限 / フロー統計参照 /\n"
" セキュリティ統計情報 / セキュリティアラート /\n"
" セキュリティ脅威分析 / XDP 監視 / 攻撃検知\n"
" verify : ANTA によるネットワーク検証・スナップショット\n"
" 例: 事後検証 / post-check / anta テスト / 副作用確認\n\n"
"【判定ルール(重要)】\n"
" - 「VXLAN」「EVPN」「VNI」「VTEP」「route-target」+「設定/変更/追加」→ eapi_config\n"
" - 「セキュリティ」+「統計/監視/アラート/分析/フロー」→ security\n"
" - 「セキュリティポリシー」「ACL」「FW ルール」等デバイス設定の参照 → read\n"
" - 「ブロック」「遮断」「xdp」「ebpf」単独 → security\n"
" - 「show」で始まる CLI コマンド → read\n"
" - NTP / BGP / OSPF / インターフェース / ルート の参照 → read\n\n"
f"クエリ: {query}\n"
"回答(eapi_config / read / write / security / verify のいずれか一単語のみ):"
).strip().lower()
if "eapi_config" in result: return "eapi_config"
if "write" in result: return "write"
if "security" in result: return "security"
if "verify" in result: return "verify"
return "read"
# ═══════════════════════════════════════════════════════════════════════════════
# A2A 通信ユーティリティ
# ═══════════════════════════════════════════════════════════════════════════════
def _make_a2a_request(payload: dict, msg_id: str = None) -> dict:
# v1.1.0: jsonrpc ラッパーなし、REST 形式で直接送信
mid = msg_id or f"hub-{datetime.now().strftime('%H%M%S%f')}"
return {
"message": {
"role": "ROLE_USER",
"parts": [{"text": json.dumps(payload, ensure_ascii=False)}],
"messageId": mid,
},
"configuration": {"returnImmediately": False},
}
# v1.1.0: A2A-Version ヘッダーが必須
_A2A_HEADERS = {"A2A-Version": "1.0", "Content-Type": "application/json"}
def _extract_text(a2a_resp: dict) -> dict:
"""
A2A v1.1.0 レスポンスからテキストを取り出し JSON として返す。
【#1774 対応】
message.parts(会話・要約)と artifacts(最終成果物)を両方ハンドリングする。
Artifact が存在する場合は _artifact_{name} キーとして parsed にマージし、
"result" という名前の Artifact は parsed["result"] を上書きする(成果物が正)。
既知の Artifact 名:
anta_report → parsed["_artifact_anta_report"] (ANTA テストレポート全体)
xdp_log → parsed["_artifact_xdp_log"] (XDP 統計・分析ログ)
report → parsed["_artifact_report"] (診断レポート全体)
diff → parsed["_artifact_diff"] (差分テキスト)
"""
try:
# ── 1. Artifact を先に収集 ─────────────────────────────────────────────
artifact_map: dict = {}
for art in a2a_resp.get("artifacts", []):
if not isinstance(art, dict):
continue
art_name = art.get("name", "")
for part in art.get("parts", []):
text = part.get("text") if isinstance(part, dict) else None
if text is None:
continue
try:
artifact_map[art_name] = json.loads(text)
except (json.JSONDecodeError, TypeError):
artifact_map[art_name] = {"_raw_text": text}
break # 各 Artifact の parts 先頭のみ使用
# ── 2. message.parts から会話テキスト(要約)を取得 ───────────────────
parts = (
a2a_resp.get("message", {}).get("parts", [])
or a2a_resp.get("result", {}).get("parts", [])
or a2a_resp.get("result", {}).get("message", {}).get("parts", [])
)
parsed: dict = {}
for part in parts:
text = part.get("text") if isinstance(part, dict) else None
if text is None:
continue
try:
parsed = json.loads(text)
except (json.JSONDecodeError, TypeError):
parsed = {"_raw_text": text}
break
# Hub 経由のネスト展開(後方互換)
inner = parsed.get("result")
if isinstance(inner, dict) and "message" in inner:
parsed["result"] = _extract_text(inner)
break
if not parsed and not artifact_map:
return {"_raw_a2a": a2a_resp}
# ── 3. Artifact を parsed にマージ(Artifact が正の成果物)─────────────
for art_name, art_data in artifact_map.items():
parsed[f"_artifact_{art_name}"] = art_data
if "result" in artifact_map:
parsed["result"] = artifact_map["result"]
return parsed
except Exception as e:
return {"_parse_error": str(e)}
async def _forward(target_url: str, payload: dict) -> dict:
"""
ダウンストリームの A2A サーバへリクエストを転送し、レスポンスを返す。
_extract_text() が Artifact を _artifact_* キーとして展開済みの dict を返す。
"""
# v1.1.0: /message:send エンドポイント + A2A-Version ヘッダー必須
a2a_req = _make_a2a_request(payload)
endpoint = target_url.rstrip("/") + "/message:send"
async with httpx.AsyncClient(timeout=HTTP_TIMEOUT) as client:
resp = await client.post(endpoint, json=a2a_req, headers=_A2A_HEADERS)
resp.raise_for_status()
return _extract_text(resp.json())
async def _push_log(trace_id: str, message: str):
for ws in _ws_clients.get(trace_id, []):
try:
await ws.send_text(json.dumps({"trace_id": trace_id, "log": message}))
except Exception:
pass
# ═══════════════════════════════════════════════════════════════════════════════
# A2A AgentExecutor(A2A プロトコル経由のルーティング)
# ═══════════════════════════════════════════════════════════════════════════════
class TaskDecomposeExecutor(AgentExecutor):
def _parse_request(self, text: str) -> dict:
text = text.strip()
try:
p = json.loads(text)
if isinstance(p, dict) and "query" in p:
return p
except json.JSONDecodeError:
pass
return {"query": text}
async def execute(self, context: RequestContext, event_queue: EventQueue) -> None:
from a2a.types.a2a_pb2 import Part as _Part, Message as _Message, Role as _Role
import uuid as _uuid
def _make_message(text):
msg = _Message()
msg.role = _Role.ROLE_AGENT
msg.message_id = str(_uuid.uuid4())
if context.task_id:
msg.task_id = context.task_id
if context.context_id:
msg.context_id = context.context_id
msg.parts.append(_Part(text=text))
return msg
async def _send_text(text):
await event_queue.enqueue_event(_make_message(text))
raw_text = "".join(
part.text for part in context.message.parts
if part.HasField("text")
)
if not raw_text.strip():
await _send_text(get_msg("ws_empty"))
return
params = self._parse_request(raw_text)
query = params.get("query", raw_text)
device_ip = params.get("device_ip")
username = params.get("username")
password = params.get("password")
# port はルート別に使い分け(NETCONF:830 / eAPI:各サーバデフォルト)
port = params.get("port") # None の場合は転送先サーバのデフォルトを使用
deploy = params.get("deploy", False)
logger.info(f"[A2A] 受信: {query[:80]} deploy={deploy}")
# action フィールドが明示されている場合は classify_query より優先
_action = params.get("action", "")
if _action in ("verify", "snapshot", "post_check", "compare"):
route = "verify"
elif _action in ("block", "unblock", "qos_set", "qos_list", "qos_get",
"drop_list", "stats", "top", "info", "analyze"):
route = "security"
else:
route = classify_query(query)
if route == "security":
xdp_payload = {"query": query}
try:
inner = await _forward(XDP_A2A_URL, xdp_payload)
result = {"query": query, "route": "security",
"routed_to": XDP_A2A_URL, "status": "success",
"result": inner}
except httpx.ConnectError as e:
result = {"query": query, "route": "security",
"routed_to": XDP_A2A_URL,
"status": "error", "message": f"XDP A2A 接続エラー: {e}"}
except Exception as e:
result = {"query": query, "route": "security",
"routed_to": XDP_A2A_URL,
"status": "error", "message": str(e)}
elif route == "verify":
anta_payload = {
"query": query,
"action": params.get("action", "verify"),
"snapshot_id": params.get("snapshot_id", ""),
"tests": params.get("tests"),
"device_ip": device_ip or "",
"username": username or "",
"password": password or "",
}
try:
inner = await _forward(ANTA_A2A_URL, anta_payload)
result = {"query": query, "route": "verify",
"routed_to": ANTA_A2A_URL, "status": "success",
"result": inner}
except httpx.ConnectError as e:
result = {"query": query, "route": "verify",
"routed_to": ANTA_A2A_URL,
"status": "error",
"message": f"ANTA A2A 接続エラー: {e} (port:8004 未起動の可能性)"}
except Exception as e:
result = {"query": query, "route": "verify", "routed_to": ANTA_A2A_URL,
"status": "error", "message": str(e)}
else:
target_url = NETCONF_A2A_URL if route == "write" else EAPI_A2A_URL
forward_payload = {
"query": query, "device_ip": device_ip or "",
"username": username or "", "password": password or "",
"deploy": deploy,
}
# NETCONF(write) のみ port を転送。read は eAPI サーバのデフォルト(443)を使用
if route == "write" and port:
forward_payload["port"] = port
elif route == "write":
forward_payload["port"] = "830"
try:
inner = await _forward(target_url, forward_payload)
result = {"query": query, "route": route,
"routed_to": target_url, "status": "success",
"result": inner}
except httpx.ConnectError as e:
result = {"query": query, "route": route, "routed_to": target_url,
"status": "error", "message": f"接続エラー: {e}"}
except Exception as e:
result = {"query": query, "route": route, "routed_to": target_url,
"status": "error", "message": str(e)}
await _send_text(json.dumps(result, ensure_ascii=False, indent=2))
async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
raise UnsupportedOperationError(get_msg("cancel_unsupported"))
# ═══════════════════════════════════════════════════════════════════════════════
# FastAPI アプリ(NiceGUI UI 向け REST + WebSocket)
# ═══════════════════════════════════════════════════════════════════════════════
rest_app = FastAPI(
title="Arista Network Agent Hub API",
version=VERSION,
description="task_decompose A2A Hub + NiceGUI 向け REST API",
)
rest_app.add_middleware(
CORSMiddleware, allow_origins=["*"],
allow_methods=["*"], allow_headers=["*"],
)
# ── スキーマ ─────────────────────────────────────────────────────────────────
class DeviceConfig(BaseModel):
ip: str = "172.20.100.31"
port: str = "830"
username: str = "admin"
password: str = "admin"
class ExecuteRequest(BaseModel):
query: str
device: DeviceConfig = DeviceConfig()
class DeployRequest(BaseModel):
device: DeviceConfig = DeviceConfig()
snapshot_id: str = "" # CNV: Before Snapshot ID(空文字 = Post-Check スキップ)
# ── GET /healthz ──────────────────────────────────────────────────────────────
@rest_app.get("/healthz", tags=["ops"])
async def healthz():
servers = {}
async with httpx.AsyncClient(timeout=5) as client:
for name, url in [("netconf", NETCONF_A2A_URL),
("eapi", EAPI_A2A_URL),
("xdp", XDP_A2A_URL),
("anta", ANTA_A2A_URL),
("eapi_config", EAPI_CONFIG_A2A_URL)]:
try:
r = await client.get(f"{url}/.well-known/agent.json")
servers[name] = {"status": "ok", "name": r.json().get("name")}
except Exception as e:
servers[name] = {"status": "error", "message": str(e)[:60]}
return {
"status": "ok",
"version": VERSION,
"build_date": BUILD_DATE,
"timestamp": datetime.now(timezone.utc).isoformat(),
"hub_port": A2A_PORT,
"downstream": servers,
"routes": {
"write": NETCONF_A2A_URL,
"read": EAPI_A2A_URL,
"security": XDP_A2A_URL,
"verify": ANTA_A2A_URL,
"eapi_config": EAPI_CONFIG_A2A_URL,
},
}
# ── POST /validate ────────────────────────────────────────────────────────────
class ValidateRequest(BaseModel):
xml: str
@rest_app.post("/validate", tags=["ops"])
async def validate_xml(req: ValidateRequest):
"""
XML 構文チェック。ET.fromstring() で parse を試み、
成功なら valid=True、失敗なら valid=False + エラーメッセージを返す。
Hub 側で完結するため、バックエンドサーバへの通信は不要。
"""
import xml.etree.ElementTree as ET
xml_str = req.xml.strip()
if not xml_str:
return {"valid": False, "message": "XML が空です"}
try:
ET.fromstring(xml_str)
return {"valid": True, "message": "✅ XML 構文チェック OK"}
except ET.ParseError as e:
return {"valid": False, "message": f"❌ XML 構文エラー: {e}"}
# ── POST /execute ─────────────────────────────────────────────────────────────
@rest_app.post("/execute", tags=["netconf"])
async def execute(req: ExecuteRequest):
"""
Dry-run 実行。A2A Hub 経由でルーティングし、deploy=False で XML/結果を返す。
変更系: NETCONF XML を生成して返す(実機未投入)
参照系: eAPI show を即時実行して結果を返す
"""
trace_id = str(uuid.uuid4())[:8]
route = classify_query(req.query)
logger.info(f"[{trace_id}] /execute: {req.query!r} route={route}")
# ── VLAN名のみ削除ガード ──────────────────────────────────────────────────
# ネットワーク運用の原則: VLAN操作はVLAN IDで行う(名前は一意でないため危険)。
# 削除系クエリで VLAN IDの数値が含まれず、VLAN名のみ指定されている場合は
# バックエンドに転送せず即エラーを返す。
_delete_verbs = ["削除", "delete", "remove", "消し", "消す", "なくし"]
_q_lower = req.query.lower()
_has_delete = any(k in _q_lower for k in _delete_verbs)
_has_vlan_kw = "vlan" in _q_lower
_has_vlan_id = bool(re.search(r'\b\d+\b', req.query))
if route == "write" and _has_delete and _has_vlan_kw and not _has_vlan_id:
msg = (
"VLAN削除にはVLAN IDの指定が必要です。\n"
"VLAN名からの自動解決は行いません(名前は一意でないため危険)。\n"
"例: 'VLAN 102 を削除して' のようにVLAN IDで指定してください。"
)
logger.warning(f"[{trace_id}] VLAN名のみ削除をブロック: {req.query!r}")
error_response = {
"trace_id": trace_id,
"route": "write",
"is_read": False,
"status": "blocked",
"overall_status": "blocked",
"summary": msg,
"xml": "",
"session_diff": {},
"task_summaries": [{
"task_id": "task_1",
"operation": "delete_vlan",
"target": req.query,
"deploy_status": "blocked",
"audit_message": msg,
}],
}
_trace_store[trace_id] = {
**error_response,
"device": req.device.model_dump(),
"query": req.query,
"is_read": False,
"executed_at": datetime.now(timezone.utc).isoformat(),
}
from fastapi.responses import Response as _VlanErrResp
return _VlanErrResp(
content=json.dumps(error_response, ensure_ascii=True),
media_type="application/json",
)
is_read = (route == "read")
# eAPI と NETCONF でポートを使い分ける
# eAPI: HTTPS/443(実機確認済み。NETCONF port 830 を渡すと SSL エラーになる)
# NETCONF: req.device.port をそのまま使う(デフォルト 830)
EAPI_DEFAULT_PORT = int(os.getenv("EAPI_PORT", "443"))
payload = {
"query": req.query,
"device_ip": req.device.ip,
"username": req.device.username,
"password": req.device.password,
"port": str(EAPI_DEFAULT_PORT) if is_read else req.device.port,
"deploy": is_read, # 参照系のみ即時実行
}
if route == "security":
# analyze クエリは action="analyze" を付与して送る
_analyze_keywords = ["分析", "解析", "analyze", "ai解析", "提案"]
_is_analyze = any(k in req.query.lower() for k in _analyze_keywords)
xdp_payload = {"query": req.query, "deploy": False}
if _is_analyze:
xdp_payload["action"] = "analyze"
try:
xdp_result = await _forward(XDP_A2A_URL, xdp_payload)
except httpx.ConnectError as e:
raise HTTPException(status_code=503,
detail=f"XDP A2A Server ({XDP_A2A_URL}) に接続できません: {e}")
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
# analyze アクション時は exec_tags / analysis をトップレベルに引き上げて
# UI(app_a2a.py)が直接参照しやすくする
# ── #1774: _artifact_* キーをトップレベルに透過伝播 ─────────────────────
_artifact_passthrough = {k: v for k, v in xdp_result.items()
if k.startswith("_artifact_")}
response = {
"trace_id": trace_id,
"route": "security",
"is_read": False,
"status": xdp_result.get("status", "unknown"),
"summary": xdp_result.get("summary", ""),
"routed_to": XDP_A2A_URL,
"result": xdp_result,
# analyze 時の追加フィールド(非 analyze 時は空)
"analysis": xdp_result.get("analysis", ""),
"exec_tags": xdp_result.get("exec_tags", []),
# security 操作では xml/session_diff は不要
"xml": "",
"session_diff": {},
**_artifact_passthrough, # _artifact_xdp_log 等を透過
}
_trace_store[trace_id] = {
**response,
"device": req.device.model_dump(),
"query": req.query,
"executed_at": datetime.now(timezone.utc).isoformat(),
}
from fastapi.responses import Response
return Response(
content=json.dumps(response, ensure_ascii=True),
media_type="application/json",
)
# ── mixed ルート: 参照+変更の混在クエリ → eAPI で参照だけ実行して警告 ───────
# 設定変更は実行せず参照結果のみ返す。UI 側で警告バブルを追加表示する。
if route == "mixed":
_warn_msg = (
"⚠️ 参照と設定変更が混在しています。\n"
"参照結果のみ表示しました。設定変更は別途入力してください。\n"
"例: まず「VLANの状態を確認して」→ 次に「VLAN ID 103 の DEV3_VLAN を作成して」"
)
logger.info(f"[{trace_id}] mixed クエリ検出: 参照のみ実行")
mixed_payload = {
"query": req.query,
"device_ip": req.device.ip,
"username": req.device.username,
"password": req.device.password,
"port": str(EAPI_DEFAULT_PORT),
"deploy": True, # eAPI は即時実行
}
try:
mixed_result = await _forward(EAPI_A2A_URL, mixed_payload)
except Exception as e:
mixed_result = {"status": "error", "message": str(e)}
response = {
"trace_id": trace_id,
"route": "read", # UI は read として処理
"is_read": True,
"status": mixed_result.get("status", "unknown"),
"summary": mixed_result.get("summary", ""),
"result": mixed_result,
"formatted": mixed_result.get("formatted_text",
mixed_result.get("formatted", "")),
"xml": "",
"session_diff": {},
"mixed_warning": _warn_msg, # ★ UI が警告バブルを表示するフラグ
}
_trace_store[trace_id] = {
**response,
"device": req.device.model_dump(),
"query": req.query,
"is_read": True,
"executed_at": datetime.now(timezone.utc).isoformat(),
}
from fastapi.responses import Response as _MixedResp
return _MixedResp(
content=json.dumps(response, ensure_ascii=True),
media_type="application/json",
)
# ── eapi_config ルート: VXLAN/EVPN 設定変更 (8006) dry-run ─────────────────
if route == "eapi_config":
eapi_cfg_payload = {
"query": req.query,
"deploy": False, # dry-run(Phase1)
}
try:
eapi_cfg_result = await _forward(EAPI_CONFIG_A2A_URL, eapi_cfg_payload)
except httpx.ConnectError as e:
raise HTTPException(status_code=503,
detail=f"eAPI Config Server ({EAPI_CONFIG_A2A_URL}) に接続できません: {e}")
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
# eAPI config の diff を diff_history 互換の session_diff 形式に変換する
raw_diff = eapi_cfg_result.get("diff", "")
diff_lines = []
for line in raw_diff.splitlines():
if line.startswith("+") and not line.startswith("+++"):
diff_lines.append({"op": "+", "text": line[1:].lstrip()})
elif line.startswith("-") and not line.startswith("---"):
diff_lines.append({"op": "-", "text": line[1:].lstrip()})
elif line.startswith("@@"):
diff_lines.append({"op": "@@", "text": line})
session_diff = {
"status": "ok" if raw_diff else "no_change",
"diff_lines": diff_lines,
"diff_text": raw_diff,
"message": eapi_cfg_result.get("message", ""),
}
# AI 要約(LLMによるdiff解釈)を生成する
if raw_diff:
try:
ai_summary = invoke_with_fallback(
llm,
"以下は Arista cEOS の設定変更差分(configure session diffs)です。\n"
"この差分を日本語で2〜3文に要約してください。\n"
"変更の意図と影響を具体的に述べてください。\n\n"
f"差分:\n{raw_diff[:1500]}\n\n要約:"
).strip()
session_diff["ai_summary"] = ai_summary
except Exception:
session_diff["ai_summary"] = ""
cmds = eapi_cfg_result.get("cmds", [])
status = eapi_cfg_result.get("status", "unknown")
response = {
"trace_id": trace_id,
"route": "eapi_config",
"is_read": False,
"status": status, # "plan" | "blocked" | "error"
"summary": (
f"eAPI Config dry-run: {len(cmds)}件のコマンドを計画しました"
if status == "plan" else eapi_cfg_result.get("message", "")
),
"xml": "", # NETCONF XML は不要
"session_diff": session_diff,
"routed_to": EAPI_CONFIG_A2A_URL,
"result": eapi_cfg_result,
# eapi_config 専用フィールド(UI が参照)
"eapi_cmds": cmds,
"eapi_diff": raw_diff,
"eapi_session": eapi_cfg_result.get("session", ""),
"eapi_warning": eapi_cfg_result.get("warning", ""),
}
_trace_store[trace_id] = {
**response,
"device": req.device.model_dump(),
"query": req.query,
"executed_at": datetime.now(timezone.utc).isoformat(),
}
from fastapi.responses import Response as _ECResp
return _ECResp(
content=json.dumps(response, ensure_ascii=True),
media_type="application/json",
)
# ── verify ルート: ANTA Snapshot 検証 (8004) ─────────────────────────────
if route == "verify":
anta_payload = {
"query": req.query,
"action": "verify",
"device_ip": req.device.ip,
"username": req.device.username,
"password": req.device.password,
}
try:
anta_result = await _forward(ANTA_A2A_URL, anta_payload)
except httpx.ConnectError as e:
raise HTTPException(status_code=503,
detail=f"ANTA A2A Server ({ANTA_A2A_URL}) に接続できません: {e}")
except Exception as e:
raise HTTPException(status_code=500, detail=str(e))
# ── #1774: _artifact_* キーをトップレベルに透過伝播 ─────────────────────
_artifact_passthrough = {k: v for k, v in anta_result.items()
if k.startswith("_artifact_")}
response = {
"trace_id": trace_id,
"route": "verify",
"is_read": True,
"status": anta_result.get("status", "unknown"),
"summary": anta_result.get("summary", ""),
"routed_to": ANTA_A2A_URL,
"result": anta_result,
"xml": "",
"session_diff": {},
"snapshot_id": anta_result.get("snapshot_id", ""),
"tests_total": anta_result.get("tests_total", 0),
"tests_passed": anta_result.get("tests_passed", 0),
"tests_failed": anta_result.get("tests_failed", 0),
**_artifact_passthrough, # _artifact_anta_report 等を透過
}
_trace_store[trace_id] = {
**response,
"device": req.device.model_dump(),
"query": req.query,
"executed_at": datetime.now(timezone.utc).isoformat(),
}
from fastapi.responses import Response as _Resp
return _Resp(
content=json.dumps(response, ensure_ascii=True),
media_type="application/json",
)
target_url = EAPI_A2A_URL if is_read else NETCONF_A2A_URL
result = None
for attempt in range(2):
try:
result = await _forward(target_url, payload)
break
except Exception as e:
logger.warning(f"[{trace_id}] attempt {attempt+1} failed: {e}")
if attempt < 1:
await asyncio.sleep(2.0)
if result is None:
raise HTTPException(status_code=500, detail=get_msg("hub_conn_error"))
# ── 変更系: NETCONF dry-run 結果 + session diff(事前 diff) ────────────
xml_out = ""
session_diff_result = {}
if not is_read:
# NETCONFサーバの response_payload は _extract_text() で展開されて result に入る。
# final_xml / generated_xml を以下の優先順で取得する:
# 1. result 直下(今回追加)
# 2. result["result"] のネスト内(旧形式互換)
# 3. task_summaries の先頭タスクの final_xml
raw = result.get("result", result)
xml_out = (
result.get("final_xml", "") # ★ NETCONFサーバが直接返す(今回修正)
or result.get("generated_xml", "")
or raw.get("final_xml", "")
or raw.get("generated_xml", "")
)
# task_summaries からも探す(フォールバック)
if not xml_out:
for ts in result.get("task_summaries", []):
_fx = ts.get("final_xml") or ts.get("generated_xml", "")
if _fx:
xml_out = _fx
break