Commit 7aae58aa authored by Hoanganhvu123's avatar Hoanganhvu123

chore: codebase cleanup, healthcheck, graceful shutdown, and learning loop integration

parent 3e89d304
No preview for this file type
# Claude API via Zunef Gateway (OpenAI-compatible endpoint) # Claude API via Zunef Gateway (OpenAI-compatible endpoint)
OPENAI_API_KEY=zf_tITHSHgVdaceZgj6idE214yPxAZewkhE OPENAI_API_KEY=your_zunef_api_key_here
DEFAULT_MODEL=claude-sonnet-4-6 DEFAULT_MODEL=claude-sonnet-4-6
# Alternative: Direct Anthropic API (if not using zunef) # Alternative: Direct Anthropic API (if not using zunef)
...@@ -9,6 +9,8 @@ DEFAULT_MODEL=claude-sonnet-4-6 ...@@ -9,6 +9,8 @@ DEFAULT_MODEL=claude-sonnet-4-6
DS2API_BASE_URL=http://your-ds2-api:port DS2API_BASE_URL=http://your-ds2-api:port
DS2API_API_KEY=your_ds2_key DS2API_API_KEY=your_ds2_key
GORQ_API_KEY=your_gorq_key GORQ_API_KEY=your_gorq_key
# Ultra Descriptions Agent: comma-separated Groq API keys for round-robin
GROQ_API_KEYS=gsk_key1,gsk_key2,gsk_key3
GOOGLE_API_KEY=your_google_key GOOGLE_API_KEY=your_google_key
# Server # Server
...@@ -22,10 +24,11 @@ REDIS_PORT=6379 ...@@ -22,10 +24,11 @@ REDIS_PORT=6379
STARROCKS_HOST=your_starrocks_host STARROCKS_HOST=your_starrocks_host
STARROCKS_PORT=9030 STARROCKS_PORT=9030
# Langfuse Tracing (optional) # Langfuse (observability)
LANGFUSE_TRACING=false LANGFUSE_TRACING=false
LANGFUSE_PUBLIC_KEY= # LANGFUSE_SECRET_KEY=your_langfuse_secret
LANGFUSE_SECRET_KEY= # LANGFUSE_PUBLIC_KEY=your_langfuse_public
# LANGFUSE_HOST=https://your-langfuse-instance
# Disable Auth (for development) # Auth
DISABLE_AUTH=true DISABLE_AUTH=true
DS2API_BASE_URL=http://localhost:5001 DS2API_BASE_URL=http://localhost:5001
DS2API_API_KEY=your-api-key-1 DS2API_API_KEY=your-test-api-key
DEFAULT_MODEL=deepseek-v4-flash DEFAULT_MODEL=deepseek-v4-flash
OPENAI_API_KEY= OPENAI_API_KEY=
GORQ_API_KEY= GORQ_API_KEY=
GOOGLE_API_KEY= GOOGLE_API_KEY=
PORT=5000 PORT=5000
REDIS_HOST=172.16.2.192 REDIS_HOST=127.0.0.1
REDIS_PORT=6379 REDIS_PORT=6379
STARROCKS_HOST= STARROCKS_HOST=
STARROCKS_PORT=9030 STARROCKS_PORT=9030
......
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
"""
ChatKit Protocol — Unified SSE Event Types for Agent Streaming.
Inspired by OpenAI ChatKit SDK (`chatkit-python`), this module provides
a standardised set of Pydantic models for all Server-Sent Events (SSE)
emitted by our agents (report_agent, reaction_agent).
Benefits:
- Frontend handles ONE consistent event vocabulary
- Structured workflow tracking (tasks with status transitions)
- Built-in cancel/retry support
- Type-safe serialization
Usage:
from agent.chatkit_protocol import (
ProgressEvent, WorkflowEvent, StreamDeltaEvent,
ToolTask, DoneEvent, ErrorEvent, sse_event,
)
"""
from __future__ import annotations
import json
import logging
from datetime import UTC, datetime
from enum import Enum
from typing import Any, Literal
from pydantic import BaseModel, Field
logger = logging.getLogger(__name__)
# ═══════════════════════════════════════════════════════════════════════
# ENUMS
# ═══════════════════════════════════════════════════════════════════════
class TaskStatus(str, Enum):
"""Lifecycle states for a workflow task."""
pending = "pending"
running = "running"
done = "done"
error = "error"
class WorkflowAction(str, Enum):
"""Actions that mutate a workflow's state."""
started = "started"
task_added = "task_added"
task_updated = "task_updated"
done = "done"
# ═══════════════════════════════════════════════════════════════════════
# TASK MODELS
# ═══════════════════════════════════════════════════════════════════════
class ToolTask(BaseModel):
"""A single tool execution within a workflow."""
index: int
name: str
purpose: str = ""
params: dict[str, Any] = Field(default_factory=dict)
status: TaskStatus = TaskStatus.pending
preview: dict[str, Any] | None = None
success: bool | None = None
class SubTask(BaseModel):
"""A high-level sub-task in the agent's plan."""
label: str
status: TaskStatus = TaskStatus.pending
# ═══════════════════════════════════════════════════════════════════════
# SSE EVENT MODELS
# ═══════════════════════════════════════════════════════════════════════
class StreamOptionsEvent(BaseModel):
"""Emitted at stream start to declare stream capabilities.
Mirrors ChatKit's StreamOptionsEvent — tells the frontend
whether it can cancel the stream mid-flight.
"""
type: Literal["stream_options"] = "stream_options"
allow_cancel: bool = True
class ProgressEvent(BaseModel):
"""Emitted to show status text / thinking steps to the user.
Replaces the old ``{"type": "thinking", "step": "..."}`` pattern.
Mirrors ChatKit's ProgressUpdateEvent.
"""
type: Literal["progress"] = "progress"
text: str
icon: str | None = None
metadata: dict[str, Any] = Field(default_factory=dict)
class WorkflowEvent(BaseModel):
"""Tracks workflow lifecycle: start → add tasks → update tasks → done.
Mirrors ChatKit's WorkflowTaskAdded / WorkflowTaskUpdated pattern.
The ``action`` field drives a state machine on the frontend:
- ``started``: A new workflow begins (with optional sub_tasks plan)
- ``task_added``: A new tool task is queued
- ``task_updated``: A tool task changed status (running → done / error)
- ``done``: The workflow is complete
"""
type: Literal["workflow"] = "workflow"
action: WorkflowAction
cycle: int = 0
thinking: str = ""
current_task: str = ""
# For 'started'
sub_tasks: list[SubTask] = Field(default_factory=list)
tools_count: int = 0
# For 'task_added' / 'task_updated'
task: ToolTask | None = None
# For 'done' or reflect summary
data_sufficient: bool | None = None
missing: list[str] = Field(default_factory=list)
completed_sub_tasks: list[str] = Field(default_factory=list)
drill_down_opportunities: list[str] = Field(default_factory=list)
next_task: str = ""
class StreamDeltaEvent(BaseModel):
"""Incremental text delta for streaming content.
Mirrors ChatKit's ``AssistantMessageContentPartTextDelta`` and
``WidgetStreamingTextValueDelta``.
``stream_id`` identifies WHICH stream this delta belongs to:
- ``"thinking"`` — agent reasoning tokens
- ``"html"`` — report HTML body tokens
- ``"answer"`` — direct text answer tokens
When ``done=True``, the stream is finalised and ``content``
contains the full accumulated text (if available).
"""
type: Literal["stream_delta"] = "stream_delta"
stream_id: str
delta: str = ""
done: bool = False
content: str | None = None # Full content when done=True
class DirectResponseEvent(BaseModel):
"""Agent provides an immediate answer (no tool execution needed).
Used for greetings, simple questions answered from context, etc.
"""
type: Literal["direct_response"] = "direct_response"
message: str
class ReportCompleteEvent(BaseModel):
"""Final event when a report is fully generated.
Contains the full HTML body and generation metadata.
"""
type: Literal["report_complete"] = "report_complete"
html: str
tools_used: list[str] = Field(default_factory=list)
cycles_count: int = 0
report_id: int | None = None
conversation_id: str | None = None
persisted: bool = False
appendix_label: str | None = None # For follow-up reports
class ErrorEvent(BaseModel):
"""Structured error event with optional retry hint.
Mirrors ChatKit's ErrorEvent — gives the frontend enough info
to show a meaningful error message and optionally offer retry.
"""
type: Literal["error"] = "error"
message: str
code: str = "agent_error"
allow_retry: bool = False
metadata: dict[str, Any] = Field(default_factory=dict)
class DoneEvent(BaseModel):
"""Stream completion marker.
``result`` can carry final summary data (e.g. reaction simulation results).
"""
type: Literal["done"] = "done"
result: dict[str, Any] | None = None
class ReflectEvent(BaseModel):
"""Emitted after the reflect node evaluates data sufficiency.
Provides transparency into the agent's decision-making loop.
"""
type: Literal["reflect"] = "reflect"
cycle: int
thinking: str = ""
data_sufficient: bool = False
missing: list[str] = Field(default_factory=list)
current_task: str = ""
completed_sub_tasks: list[str] = Field(default_factory=list)
drill_down_opportunities: list[str] = Field(default_factory=list)
next_task: str = ""
next_sub_tasks: list[str] = Field(default_factory=list)
class ContextCompressedEvent(BaseModel):
"""Emitted when Hermes ContextManager compresses tool results."""
type: Literal["context_compressed"] = "context_compressed"
items_compressed: int = 0
items_pruned: int = 0
class ErrorRecoveryEvent(BaseModel):
"""Emitted when Hermes error recovery classifies an exception."""
type: Literal["error_recovery"] = "error_recovery"
reason: str
retryable: bool = False
message: str = ""
# ═══════════════════════════════════════════════════════════════════════
# SSE SERIALISATION
# ═══════════════════════════════════════════════════════════════════════
# Union of all event types for type-checking
ChatKitEvent = (
StreamOptionsEvent
| ProgressEvent
| WorkflowEvent
| StreamDeltaEvent
| DirectResponseEvent
| ReportCompleteEvent
| ReflectEvent
| ContextCompressedEvent
| ErrorRecoveryEvent
| ErrorEvent
| DoneEvent
)
def sse_event(data: dict | BaseModel) -> str:
"""Serialize an event (dict or Pydantic model) as an SSE data line.
Handles both legacy dict events and new Pydantic protocol events.
"""
if isinstance(data, BaseModel):
payload = data.model_dump(exclude_none=True)
else:
payload = data
return f"data: {json.dumps(payload, ensure_ascii=False, default=str)}\n\n"
def event_to_dict(event: ChatKitEvent | dict) -> dict:
"""Convert a protocol event to a plain dict for yielding through LangGraph."""
if isinstance(event, BaseModel):
return event.model_dump(exclude_none=True)
return event
...@@ -82,7 +82,7 @@ async def extract_and_save_user_insight(json_content: str, identity_key: str) -> ...@@ -82,7 +82,7 @@ async def extract_and_save_user_insight(json_content: str, identity_key: str) ->
# Save to Redis # Save to Redis
await save_user_insight_to_redis(identity_key, insight_str) await save_user_insight_to_redis(identity_key, insight_str)
elapsed = time.time() - start_time elapsed = time.time() - start_time
logger.warning(f"✅ [user_insight] Extracted + saved in {elapsed:.2f}s | Key: {identity_key}") logger.info(f"✅ [user_insight] Extracted + saved in {elapsed:.2f}s | Key: {identity_key}")
return insight_dict return insight_dict
except json.JSONDecodeError: except json.JSONDecodeError:
continue # Try next pattern continue # Try next pattern
......
...@@ -197,7 +197,7 @@ def extract_product_ids(messages: list) -> list[dict]: ...@@ -197,7 +197,7 @@ def extract_product_ids(messages: list) -> list[dict]:
# Legacy format: {"products": [...]} # Legacy format: {"products": [...]}
product_list = tool_result["products"] product_list = tool_result["products"]
logger.warning(f"🛠️ [EXTRACT] Extracted {len(product_list)} products from tool") logger.info(f"🛠️ [EXTRACT] Extracted {len(product_list)} products from tool")
for product in product_list: for product in product_list:
# ⚡ FLATTEN: Tách variants thành separate products (mỗi màu = 1 product) # ⚡ FLATTEN: Tách variants thành separate products (mỗi màu = 1 product)
......
"""Push check_is_stock + tool_routing to Langfuse."""
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
sys.stdout.reconfigure(encoding="utf-8")
from dotenv import load_dotenv
load_dotenv(os.path.join(os.path.dirname(__file__), "..", "..", ".env"))
from langfuse import Langfuse
lf = Langfuse()
BASE = os.path.dirname(os.path.abspath(__file__))
# 1. Push check_is_stock tool prompt
stock_path = os.path.join(BASE, "..", "tool_prompts", "check_is_stock.txt")
with open(stock_path, encoding="utf-8") as f:
content = f.read()
lf.create_prompt(name="check_is_stock", prompt=content, labels=["production"], tags=["canifa", "tool-prompt"], type="text")
print(f"✅ check_is_stock updated ({len(content):,} chars)")
# 2. Push tool_routing
with open(os.path.join(BASE, "05_tool_routing.txt"), encoding="utf-8") as f:
content2 = f.read()
lf.create_prompt(name="canifa-05-tool-routing", prompt=content2, labels=["production"], tags=["canifa", "system-core"], type="text")
print(f"✅ canifa-05-tool-routing updated ({len(content2):,} chars)")
# 3. Push data_retrieval_tool prompt
drt_path = os.path.join(BASE, "..", "tool_prompts", "data_retrieval_tool.txt")
with open(drt_path, encoding="utf-8") as f:
content3 = f.read()
lf.create_prompt(name="data_retrieval_tool", prompt=content3, labels=["production"], tags=["canifa", "tool-prompt"], type="text")
print(f"✅ data_retrieval_tool updated ({len(content3):,} chars)")
lf.flush()
print("Done!")
Ahh mình hiểu rồi! Bạn muốn tags rõ ràng hơn — nhìn vào tag phải biết ngay đây là phần nào của system prompt. Đây:
Tags cải tiến — Mỗi prompt có 2 tags:
Tag 1: canifa (luôn có, filter toàn bộ project)
Tag 2: Chỉ rõ nó thuộc nhóm nào + là part gì
Name Tags
──────────────────────────────────────────────────────────
canifa-prompt-core canifa system-1-core
canifa-prompt-rules canifa system-2-rules
canifa-prompt-context canifa system-3-context
canifa-prompt-sales-core canifa system-4a-sales-core
canifa-prompt-sales-thaomai canifa system-4b-sales-thaomai
canifa-prompt-sales-upsell canifa system-4c-sales-upsell
canifa-prompt-sales-urgency canifa system-4d-sales-urgency
canifa-prompt-tool-routing canifa system-5-tool-routing
canifa-prompt-user-insight canifa system-6-user-insight
canifa-prompt-output-format canifa system-7-output-format
canifa-prompt-season canifa season
──────────────────────────────────────────────────────────
canifa-tool-brand-knowledge canifa tool-prompt
canifa-tool-check-stock canifa tool-prompt
canifa-tool-data-retrieval canifa tool-prompt
canifa-tool-promotion canifa tool-prompt
canifa-tool-store-search canifa tool-prompt
Nhìn vào là biết ngay:
Tag Ý nghĩa
system-1-core System prompt, phần 1, LÕI (identity, danh mục, liên hệ)
system-2-rules System prompt, phần 2, QUY TẮC (trung thực + ngôn ngữ)
system-3-context System prompt, phần 3, CONTEXT AWARENESS
system-4a-sales-core System prompt, phần 4a, SALES CƠ BẢN
system-4b-sales-thaomai System prompt, phần 4b, THẢO MAI
system-4c-sales-upsell System prompt, phần 4c, UPSELL
system-4d-sales-urgency System prompt, phần 4d, URGENCY
system-5-tool-routing System prompt, phần 5, TOOL ROUTING
system-6-user-insight System prompt, phần 6, USER INSIGHT
system-7-output-format System prompt, phần 7, OUTPUT FORMAT
season Seasonal/Event (marketing)
tool-prompt Prompt cho tool
Ưu điểm:
Có số thứ tự → biết ghép theo order nào (1 → 7)
Có tên → biết nội dung gì
Filter system-4 → thấy tất cả sales-related
Filter system- → thấy toàn bộ system prompt parts
Ai mới vào nhìn tags cũng hiểu ngay cấu trúc
Rõ chưa bro? OK làm đi thì mình cắt file luôn nhé! ✂️
\ No newline at end of file
"""
Push ALL prompts (system + tools) to Langfuse.
This is the master script — run this when setting up a new Langfuse instance.
Usage: python agent/prompt_module/push_all_to_langfuse.py
"""
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
sys.stdout.reconfigure(encoding="utf-8")
from dotenv import load_dotenv
load_dotenv(os.path.join(os.path.dirname(__file__), "..", "..", ".env"))
from langfuse import Langfuse
PROMPT_DIR = os.path.dirname(os.path.abspath(__file__))
TOOL_PROMPT_DIR = os.path.join(PROMPT_DIR, "..", "tool_prompts")
def read_file(path: str) -> str:
with open(path, encoding="utf-8") as f:
return f.read()
def main():
print("🚀 Push ALL Prompts to Langfuse")
print(f" Langfuse URL: {os.getenv('LANGFUSE_BASE_URL')}")
lf = Langfuse()
# ============================================================
# PART 1: System Prompt Sub-Modules (02-07 + season)
# ============================================================
print("\n" + "=" * 60)
print("PART 1: System Prompt Sub-Modules")
print("=" * 60)
SUB_MODULES = [
("02_rules.txt", "canifa-02-rules", ["canifa", "system-core"]),
("03_context.txt", "canifa-03-context", ["canifa", "system-core"]),
("04a_sales_core.txt", "canifa-04a-sales-core", ["canifa", "system-sales"]),
("04b_sales_thaomai.txt", "canifa-04b-sales-thaomai", ["canifa", "system-sales"]),
("04c_sales_upsell.txt", "canifa-04c-sales-upsell", ["canifa", "system-sales"]),
("04d_sales_urgency.txt", "canifa-04d-sales-urgency", ["canifa", "system-sales"]),
("05_tool_routing.txt", "canifa-05-tool-routing", ["canifa", "system-core"]),
("06_user_insight.txt", "canifa-06-user-insight", ["canifa", "system-core"]),
("07_output_format.txt", "canifa-07-output-format", ["canifa", "system-core"]),
]
for filename, langfuse_name, tags in SUB_MODULES:
content = read_file(os.path.join(PROMPT_DIR, filename))
lf.create_prompt(name=langfuse_name, prompt=content, labels=["production"], tags=tags, type="text")
print(f" ✅ {filename:30s} → {langfuse_name} ({len(content):,} chars)")
# Season prompt
SEASON_CONTENT = """## HƯỚNG DẪN TƯ VẤN THEO MÙA / EVENT
**Thời điểm hiện tại:** Tháng 3/2026 — Mùa Xuân, chuyển giao Đông → Hè
**Ưu tiên sản phẩm mùa này:**
- Áo khoác nhẹ, cardigan (trời se lạnh buổi sáng/tối)
- Áo phông, áo thun (ban ngày ấm)
- Sơ mi dài tay (đi làm)
- Quần jeans, quần kaki (đa năng)
**Khi khách hỏi chung chung ("có gì hot?", "gợi ý đi"):**
→ Ưu tiên giới thiệu sản phẩm phù hợp thời tiết hiện tại
→ Nhắc sale/khuyến mãi nếu có
**Event đang diễn ra:**
- (Marketing cập nhật event tại đây)
"""
lf.create_prompt(name="canifa-08-season", prompt=SEASON_CONTENT, labels=["production"], tags=["canifa", "system-addon"], type="text")
print(f" ✅ {'(season template)':30s} → canifa-08-season ({len(SEASON_CONTENT):,} chars)")
# ============================================================
# PART 2: Core System Prompt (01_core + references)
# ============================================================
print("\n" + "=" * 60)
print("PART 2: Core System Prompt (composable)")
print("=" * 60)
core_content = read_file(os.path.join(PROMPT_DIR, "01_core.txt"))
references = "\n".join(
f"@@@langfusePrompt:name={name}|label=production@@@"
for _, name, _ in SUB_MODULES
)
references += "\n@@@langfusePrompt:name=canifa-08-season|label=production@@@"
composed = core_content + "\n" + references + "\n"
lf.create_prompt(name="canifa-stylist-system-prompt", prompt=composed, labels=["production"], tags=["canifa", "system-prompt"], type="text")
print(f" ✅ canifa-stylist-system-prompt ({len(composed):,} chars, {len(SUB_MODULES) + 1} refs)")
# ============================================================
# PART 3: Tool Prompts (with correct Langfuse names!)
# ============================================================
print("\n" + "=" * 60)
print("PART 3: Tool Prompts")
print("=" * 60)
TOOL_PROMPTS = {
"brand_knowledge_tool.txt": "canifa-tool-brand-knowledge",
"check_is_stock.txt": "canifa-tool-check-stock",
"data_retrieval_tool.txt": "canifa-tool-data-retrieval",
"promotion_canifa_tool.txt": "canifa-tool-promotion",
"store_search_tool.txt": "canifa-tool-store-search",
}
for filename, langfuse_name in TOOL_PROMPTS.items():
path = os.path.join(TOOL_PROMPT_DIR, filename)
content = read_file(path)
lf.create_prompt(name=langfuse_name, prompt=content, labels=["production"], tags=["canifa", "tool-prompt"], type="text")
print(f" ✅ {filename:35s} → {langfuse_name} ({len(content):,} chars)")
# ============================================================
# VERIFY
# ============================================================
print("\n" + "=" * 60)
print("VERIFICATION")
print("=" * 60)
prompt = lf.get_prompt("canifa-stylist-system-prompt", label="production", cache_ttl_seconds=0)
assembled = prompt.prompt
print(f" System prompt assembled: {len(assembled):,} chars")
for tool_name, langfuse_name in TOOL_PROMPTS.items():
try:
p = lf.get_prompt(langfuse_name, label="production", cache_ttl_seconds=0)
print(f" ✅ {langfuse_name}: {len(p.prompt):,} chars")
except Exception as e:
print(f" ❌ {langfuse_name}: {e}")
lf.flush()
print("\n🎉 ALL PROMPTS PUSHED SUCCESSFULLY!")
if __name__ == "__main__":
main()
"""
Push prompt modules to Langfuse as composable prompts.
Creates individual sub-module prompts, then updates the core system prompt
to reference them via @@@langfusePrompt:...@@@ tags.
Usage: python push_modules_to_langfuse.py
"""
import os
import sys
sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", ".."))
sys.stdout.reconfigure(encoding="utf-8")
from dotenv import load_dotenv
load_dotenv(os.path.join(os.path.dirname(__file__), "..", "..", ".env"))
from langfuse import Langfuse
PROMPT_DIR = os.path.dirname(os.path.abspath(__file__))
# ---- Module definitions: local file → Langfuse name (NUMBERED) + tags ----
SUB_MODULES = [
("02_rules.txt", "canifa-02-rules", ["canifa", "system-core"]),
("03_context.txt", "canifa-03-context", ["canifa", "system-core"]),
("04a_sales_core.txt", "canifa-04a-sales-core", ["canifa", "system-sales"]),
("04b_sales_thaomai.txt", "canifa-04b-sales-thaomai", ["canifa", "system-sales"]),
("04c_sales_upsell.txt", "canifa-04c-sales-upsell", ["canifa", "system-sales"]),
("04d_sales_urgency.txt", "canifa-04d-sales-urgency", ["canifa", "system-sales"]),
("05_tool_routing.txt", "canifa-05-tool-routing", ["canifa", "system-core"]),
("06_user_insight.txt", "canifa-06-user-insight", ["canifa", "system-core"]),
("07_output_format.txt", "canifa-07-output-format", ["canifa", "system-core"]),
]
CORE_FILE = "01_core.txt"
CORE_PROMPT_NAME = "canifa-stylist-system-prompt"
SEASON_PROMPT_NAME = "canifa-08-season"
# Default season content
DEFAULT_SEASON_CONTENT = """## HƯỚNG DẪN TƯ VẤN THEO MÙA / EVENT
**Thời điểm hiện tại:** Tháng 3/2026 — Mùa Xuân, chuyển giao Đông → Hè
**Ưu tiên sản phẩm mùa này:**
- Áo khoác nhẹ, cardigan (trời se lạnh buổi sáng/tối)
- Áo phông, áo thun (ban ngày ấm)
- Sơ mi dài tay (đi làm)
- Quần jeans, quần kaki (đa năng)
**Khi khách hỏi chung chung ("có gì hot?", "gợi ý đi"):**
→ Ưu tiên giới thiệu sản phẩm phù hợp thời tiết hiện tại
→ Nhắc sale/khuyến mãi nếu có
**Event đang diễn ra:**
- (Marketing cập nhật event tại đây)
"""
def read_file(filename: str) -> str:
path = os.path.join(PROMPT_DIR, filename)
with open(path, encoding="utf-8") as f:
return f.read()
def push_sub_modules(lf: Langfuse):
"""Push each sub-module file as a separate Langfuse text prompt."""
print("\n" + "=" * 60)
print("STEP 1: Pushing sub-module prompts (NUMBERED)")
print("=" * 60)
for filename, langfuse_name, tags in SUB_MODULES:
content = read_file(filename)
lf.create_prompt(
name=langfuse_name,
prompt=content,
labels=["production"],
tags=tags,
type="text",
)
print(f" ✅ {filename:30s} → {langfuse_name} ({len(content):,} chars)")
# Season prompt with default content
lf.create_prompt(
name=SEASON_PROMPT_NAME,
prompt=DEFAULT_SEASON_CONTENT,
labels=["production"],
tags=["canifa", "system-addon"],
type="text",
)
print(f" ✅ {'(season template)':30s} → {SEASON_PROMPT_NAME} ({len(DEFAULT_SEASON_CONTENT):,} chars)")
def push_core_with_references(lf: Langfuse):
"""Push the core prompt with @@@langfusePrompt:...@@@ references."""
print("\n" + "=" * 60)
print("STEP 2: Pushing core prompt with composable references")
print("=" * 60)
core_content = read_file(CORE_FILE)
# Build the composed prompt: core inline + references to sub-modules
references = "\n".join(
f"@@@langfusePrompt:name={name}|label=production@@@"
for _, name, _ in SUB_MODULES
)
# Add season reference at the end
references += f"\n@@@langfusePrompt:name={SEASON_PROMPT_NAME}|label=production@@@"
composed = core_content + "\n" + references + "\n"
lf.create_prompt(
name=CORE_PROMPT_NAME,
prompt=composed,
labels=["production"],
tags=["canifa", "system-prompt"],
type="text",
)
print(f" ✅ Core prompt updated: {CORE_PROMPT_NAME}")
print(f" Inline: {CORE_FILE} ({len(core_content):,} chars)")
print(f" References: {len(SUB_MODULES) + 1} sub-modules")
def verify(lf: Langfuse):
"""Verify the assembled prompt contains all sections."""
print("\n" + "=" * 60)
print("STEP 3: Verification")
print("=" * 60)
# Fetch the composed prompt — Langfuse auto-resolves references
prompt = lf.get_prompt(CORE_PROMPT_NAME, label="production", cache_ttl_seconds=0)
assembled = prompt.prompt
print(f" Assembled prompt length: {len(assembled):,} chars")
# Check key section headings exist in assembled output
checks = [
("01 Identity (core)", "Canifa-AI Stylist"),
("02 Rules", "QUY TẮC TRUNG THỰC"),
("03 Context", "CONTEXT AWARENESS"),
("04a Sales Core", "PHONG CÁCH TƯ VẤN"),
("04b Thảo Mai", "THẢO MAI SALES"),
("04c Upsell", "UPSELL & CROSS-SELL"),
("04d Urgency", "URGENCY TACTICS"),
("05 Tool Routing", "KHI NÀO GỌI TOOL"),
("06 User Insight", "USER INSIGHT 2.0"),
("07 Output Format", "FORMAT ĐẦU RA"),
("08 Season", "HƯỚNG DẪN TƯ VẤN THEO MÙA"),
]
all_ok = True
for label, keyword in checks:
found = keyword in assembled
status = "✅" if found else "❌ MISSING"
print(f" {status} {label}: '{keyword}'")
if not found:
all_ok = False
if all_ok:
print("\n 🎉 ALL CHECKS PASSED!")
else:
print("\n ⚠️ SOME CHECKS FAILED — review above")
return all_ok
def main():
print("🚀 Push Prompt Modules to Langfuse (NUMBERED)")
print(f" Langfuse URL: {os.getenv('LANGFUSE_BASE_URL')}")
lf = Langfuse()
push_sub_modules(lf)
push_core_with_references(lf)
ok = verify(lf)
lf.flush() # ensure all API calls complete
print("\n✅ Done!" if ok else "\n⚠️ Done with warnings")
if __name__ == "__main__":
main()
This diff is collapsed.
"""
Split system_prompt.txt into 10 module files.
Cut at exact heading boundaries. Verify reassembly = original.
"""
import os
import re
import sys
sys.stdout.reconfigure(encoding='utf-8')
PROMPT_DIR = os.path.dirname(os.path.abspath(__file__))
SOURCE = os.path.join(os.path.dirname(PROMPT_DIR), "system_prompt.txt")
# Read file
with open(SOURCE, encoding="utf-8") as f:
full_text = f.read()
lines = full_text.split("\n")
print(f"Total lines: {len(lines)}")
print(f"Total chars: {len(full_text)}")
# --- Find exact line indices for each section boundary ---
# We search for the FIRST occurrence of each heading pattern
boundaries = {}
heading_patterns = [
("rules_start", r"^## 1\. QUY T"),
("context_start", r"^## 3\. CONTEXT"),
("sales_start", r"^## 4\. PHONG C"),
("thaomai_start", r"^### 4\.5\. .* TH"),
("upsell_start", r"^### 4\.6\. .* UPSELL"),
("urgency_start", r"^### 4\.7\. .* KHU"),
("tool_start", r"^## 5\. KHI N"),
("insight_start", r"^## 8\. USER INSIGHT"),
("output_start", r"^## 9\. FORMAT"),
]
for key, pattern in heading_patterns:
for i, line in enumerate(lines):
if re.search(pattern, line):
boundaries[key] = i
print(f" {key} = line {i+1} (0-indexed: {i}): {line[:60].strip()}")
break
else:
print(f" WARNING: {key} NOT FOUND!")
# Now we also need to find where to cut BEFORE each heading.
# Some headings have a "---" separator above them that belongs to the PREVIOUS section.
# We want to cut so the "---" before a heading goes with the previous file.
# Let's look just before each heading to see if there's a "---" separator
def find_cut_point(line_idx):
"""Find the actual cut point: include preceding blank lines and --- in prev section."""
# The heading itself starts the new section
# But we want to include any preceding "---" and blank lines in the PREVIOUS section
# So the new section starts AT the heading line
return line_idx
# Define file mappings: (filename, start_line_idx, end_line_idx_exclusive)
# We'll compute them from boundaries
b = boundaries
# Helper: check if the line before heading is "---" and blank lines
# If so, include them in prev section
def prev_separator_end(heading_idx):
"""Find where the CURRENT section truly starts (at or after heading_idx).
We look backward: if there's a --- and blank lines before the heading,
those belong to the PREVIOUS section. The heading line itself starts the new section."""
# But actually for the "---" that comes BEFORE section X, it's usually the
# separator at the END of section X-1. So we want:
# - Previous section includes up to (and including) the "---" before the heading
# - New section starts at the heading line
# Look backward from heading
idx = heading_idx
# Check if preceding lines are blank or "---"
while idx > 0 and lines[idx-1].strip() in ("", "---"):
idx -= 1
# The "---" and blanks before heading belong to previous section
# So previous section ends at heading_idx - 1 (inclusive)
# New section starts at heading_idx
return heading_idx
modules = [
("01_core.txt", 0, b["rules_start"]),
("02_rules.txt", b["rules_start"], b["context_start"]),
("03_context.txt", b["context_start"], b["sales_start"]),
("04a_sales_core.txt", b["sales_start"], b["thaomai_start"]),
("04b_sales_thaomai.txt", b["thaomai_start"], b["upsell_start"]),
("04c_sales_upsell.txt", b["upsell_start"], b["urgency_start"]),
("04d_sales_urgency.txt", b["urgency_start"], b["tool_start"]),
("05_tool_routing.txt", b["tool_start"], b["insight_start"]),
("06_user_insight.txt", b["insight_start"], b["output_start"]),
("07_output_format.txt", b["output_start"], len(lines)),
]
# Write each module
print("\n--- Writing module files ---")
for filename, start, end in modules:
content = "\n".join(lines[start:end])
filepath = os.path.join(PROMPT_DIR, filename)
with open(filepath, "w", encoding="utf-8") as f:
f.write(content)
line_count = end - start
print(f" {filename}: lines {start+1}-{end} ({line_count} lines, {len(content)} chars)")
# Verify: reassemble and compare
print("\n--- Verification ---")
reassembled_parts = []
for filename, start, end in modules:
filepath = os.path.join(PROMPT_DIR, filename)
with open(filepath, encoding="utf-8") as f:
reassembled_parts.append(f.read())
reassembled = "\n".join(reassembled_parts)
if reassembled == full_text:
print("✅ PASS: Reassembled content MATCHES original exactly!")
print(f" Original: {len(full_text)} chars")
print(f" Reassembled: {len(reassembled)} chars")
else:
print("❌ FAIL: Reassembled content does NOT match original!")
print(f" Original: {len(full_text)} chars")
print(f" Reassembled: {len(reassembled)} chars")
# Find first difference
for i, (a, b_char) in enumerate(zip(full_text, reassembled)):
if a != b_char:
print(f" First diff at char {i}: original='{a!r}' reassembled='{b_char!r}'")
print(f" Context original: ...{full_text[max(0,i-20):i+20]!r}...")
print(f" Context reassembled: ...{reassembled[max(0,i-20):i+20]!r}...")
break
print("\nDone!")
...@@ -7,11 +7,12 @@ Simulate phản ứng cộng đồng khi Canifa launch chiến dịch mới. ...@@ -7,11 +7,12 @@ Simulate phản ứng cộng đồng khi Canifa launch chiến dịch mới.
Dùng LLM generate realistic reactions từ các persona segments. Dùng LLM generate realistic reactions từ các persona segments.
""" """
from .reaction_agent import CAMPAIGN_TYPES, PERSONA_SEGMENTS, run_reaction_simulation, simulate_reaction_for_persona from .reaction_agent import CAMPAIGN_TYPES, PERSONA_SEGMENTS, run_reaction_simulation, run_reaction_simulation_streaming, simulate_reaction_for_persona
__all__ = [ __all__ = [
"CAMPAIGN_TYPES", "CAMPAIGN_TYPES",
"PERSONA_SEGMENTS", "PERSONA_SEGMENTS",
"run_reaction_simulation", "run_reaction_simulation",
"run_reaction_simulation_streaming",
"simulate_reaction_for_persona", "simulate_reaction_for_persona",
] ]
...@@ -184,12 +184,24 @@ def parse_json(raw: str) -> dict: ...@@ -184,12 +184,24 @@ def parse_json(raw: str) -> dict:
} }
class ThinkingStreamer: class ThinkingStreamer:
"""Helper to cleanly stream ONLY text inside <thinking>...</thinking> tags.""" """Streams tokens inside ``<thinking>...</thinking>`` tags via ChatKit protocol.
Emits ``StreamDeltaEvent(stream_id="thinking")`` for each token chunk.
"""
def __init__(self, writer): def __init__(self, writer):
from agent.chatkit_protocol import StreamDeltaEvent, event_to_dict
self.writer = writer self.writer = writer
self.buffer = "" self.buffer = ""
self.is_thinking = False self.is_thinking = False
self.done_thinking = False self.done_thinking = False
self._StreamDeltaEvent = StreamDeltaEvent
self._event_to_dict = event_to_dict
def _emit(self, token: str):
self.writer(self._event_to_dict(self._StreamDeltaEvent(
stream_id="thinking", delta=token,
)))
def feed(self, token: str): def feed(self, token: str):
if self.done_thinking: if self.done_thinking:
...@@ -203,9 +215,7 @@ class ThinkingStreamer: ...@@ -203,9 +215,7 @@ class ThinkingStreamer:
self.is_thinking = True self.is_thinking = True
parts = self.buffer.split(start_tag, 1) parts = self.buffer.split(start_tag, 1)
self.buffer = parts[1] self.buffer = parts[1]
# Fall through to process stream
else: else:
# Keep buffer small in case tag spans across chunks
if len(self.buffer) > 15: if len(self.buffer) > 15:
self.buffer = self.buffer[-15:] self.buffer = self.buffer[-15:]
return return
...@@ -217,12 +227,11 @@ class ThinkingStreamer: ...@@ -217,12 +227,11 @@ class ThinkingStreamer:
self.done_thinking = True self.done_thinking = True
parts = self.buffer.split(end_tag, 1) parts = self.buffer.split(end_tag, 1)
if parts[0]: if parts[0]:
self.writer({"type": "thinking_token", "token": parts[0]}) self._emit(parts[0])
self.buffer = "" self.buffer = ""
# Stream everything except last 15 chars to catch potential end tag splitting
elif len(self.buffer) > 15: elif len(self.buffer) > 15:
to_send = self.buffer[:-15] to_send = self.buffer[:-15]
self.writer({"type": "thinking_token", "token": to_send}) self._emit(to_send)
self.buffer = self.buffer[-15:] self.buffer = self.buffer[-15:]
...@@ -721,7 +730,8 @@ async def load_report_context(report_id: int) -> dict[str, Any] | None: ...@@ -721,7 +730,8 @@ async def load_report_context(report_id: int) -> dict[str, Any] | None:
# ─── SSE Helper ────────────────────────────────────────────────────── # ─── SSE Helper ──────────────────────────────────────────────────────
# Delegates to the unified ChatKit protocol module.
# Kept here as re-export for backward compatibility with existing imports.
from agent.chatkit_protocol import sse_event # noqa: F401, E402
def sse_event(data: dict) -> str:
"""Format dict as SSE event string."""
return f"data: {json.dumps(data, ensure_ascii=False, default=str)}\n\n"
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
import sqlite3
conn = sqlite3.connect('database/canifa_ai_dump.sqlite')
c = conn.cursor()
c.execute("SELECT match_role, COUNT(*) FROM pg__dashboard_canifa__ai_outfit_product_matches GROUP BY match_role")
print(c.fetchall())
Binary files a/backend/temp_query.py and /dev/null differ Binary files a/backend/temp_query.py and /dev/null differ
This diff is collapsed.
This diff is collapsed.
This diff is collapsed.
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment