Commit 39fac7e4 authored by Hoanganhvu123's avatar Hoanganhvu123

feat: port BrightBean Studio social management suite into Canifa pipeline

Phase 2 - AI Content Pipeline & Social Management:

== Backend Modules ==
- common/notification/: Notification engine (in-app, email, webhook, Slack, HMAC-signed)
- common/social/approval_gate.py: Editorial state machine (draft -> review -> approved -> published)
- common/social/post_queue.py: Post queue + calendar slot scheduling
- common/social/scheduler.py: Background publish engine (parallel, exponential backoff)
- common/media/image_processor.py: Media processor (Pillow resize, ffmpeg thumbnails)
- common/content_templates.py: 30+ fashion content templates + Vietnamese RSS feeds (ported from BrightBean builtin_templates.py 968L)

== API Routes ==
- api/notification_route.py: Notification CRUD + test endpoints
- api/content_approval_route.py: Approval workflow (approve/reject/submit)
- api/queue_route.py: Queue + calendar + posting slots
- api/media_route.py: Media library (upload/resize/delete)
- api/templates_route.py: Content templates + RSS feed catalog
- api/social_inbox_route.py: Unified social inbox (FB/IG/TikTok messages)

== Static Dashboards (zero-build) ==
- static/content-approval/index.html: Approval queue UI
- static/content-calendar/index.html: Weekly calendar + queue sidebar
- static/media-library/index.html: Drag-drop upload + asset grid + platform resize
- static/social-inbox/index.html: 3-column inbox (filters/list/detail) + reply/note composer
- static/content-composer/index.html: Post editor + templates + AI enhancement + schedule

== Server ==
- server.py: Registered all new routers + publish engine autostart + dashboard links in banner
parent 5d075d1e
This diff is collapsed.
...@@ -22,7 +22,9 @@ import { useDark, useToggle } from '@vueuse/core' ...@@ -22,7 +22,9 @@ import { useDark, useToggle } from '@vueuse/core'
const DISABLE_AUTH = import.meta.env.VITE_DISABLE_AUTH === 'true' const DISABLE_AUTH = import.meta.env.VITE_DISABLE_AUTH === 'true'
const route = useRoute() const route = useRoute()
const isDark = useDark()
// Default to LIGHT mode — set dark only if user explicitly toggled it
const isDark = useDark({ initialValue: false })
const toggleDark = useToggle(isDark) const toggleDark = useToggle(isDark)
// ─── SVG Icon definitions ──────────────────────────── // ─── SVG Icon definitions ────────────────────────────
...@@ -214,7 +216,7 @@ const currentBreadcrumb = computed(() => { ...@@ -214,7 +216,7 @@ const currentBreadcrumb = computed(() => {
<!-- Main App Content Area --> <!-- Main App Content Area -->
<SidebarInset> <SidebarInset>
<header class="flex h-14 shrink-0 items-center justify-between gap-2 border-b bg-background px-4 lg:px-6 transition-[width,height] ease-linear group-has-[[data-collapsible=icon]]/sidebar-wrapper:h-12 sticky top-0 z-50"> <header class="flex h-14 shrink-0 items-center justify-between gap-2 border-b border-border/80 bg-white/90 backdrop-blur-sm px-4 lg:px-6 transition-[width,height] ease-linear group-has-[[data-collapsible=icon]]/sidebar-wrapper:h-12 sticky top-0 z-50">
<!-- Left Side: Trigger + Breadcrumb --> <!-- Left Side: Trigger + Breadcrumb -->
<div class="flex items-center gap-1.5 lg:gap-2"> <div class="flex items-center gap-1.5 lg:gap-2">
......
This diff is collapsed.
...@@ -18,6 +18,7 @@ const routes = [ ...@@ -18,6 +18,7 @@ const routes = [
{ path: '/home/product', component: () => import('../pages/product.vue'), meta: { requiresAuth: true } }, { path: '/home/product', component: () => import('../pages/product.vue'), meta: { requiresAuth: true } },
{ path: '/home/product-desc', component: () => import('../pages/product-desc.vue'), meta: { requiresAuth: true } }, { path: '/home/product-desc', component: () => import('../pages/product-desc.vue'), meta: { requiresAuth: true } },
{ path: '/home/fashion-matches',component: () => import('../pages/fashion-matches.vue'), meta: { requiresAuth: true } }, { path: '/home/fashion-matches',component: () => import('../pages/fashion-matches.vue'), meta: { requiresAuth: true } },
{ path: '/home/social-inbox', component: () => import('../pages/social-inbox.vue'), meta: { requiresAuth: true } },
{ path: '/home/ai-report', component: () => import('../pages/ai-report.vue'), meta: { requiresAuth: true } }, { path: '/home/ai-report', component: () => import('../pages/ai-report.vue'), meta: { requiresAuth: true } },
{ path: '/home/ai-sql', component: () => import('../pages/ai-sql.vue'), meta: { requiresAuth: true } }, { path: '/home/ai-sql', component: () => import('../pages/ai-sql.vue'), meta: { requiresAuth: true } },
{ path: '/home/live-monitor', component: () => import('../pages/live-monitor.vue'), meta: { requiresAuth: true } }, { path: '/home/live-monitor', component: () => import('../pages/live-monitor.vue'), meta: { requiresAuth: true } },
......
...@@ -120,17 +120,51 @@ ...@@ -120,17 +120,51 @@
--font-mono: "JetBrains Mono", ui-monospace, SFMono-Regular, Menlo, Monaco, Consolas, --font-mono: "JetBrains Mono", ui-monospace, SFMono-Regular, Menlo, Monaco, Consolas,
"Liberation Mono", "Courier New", monospace; "Liberation Mono", "Courier New", monospace;
/* ── Sidebar (dark, handled by DashboardLayout inline styles) ── */ /* ── Sidebar — clean white, shadcn/ui style ── */
--sidebar: oklch(0.100 0.008 260); --sidebar: oklch(0.992 0.002 240);
--sidebar-foreground: oklch(0.920 0.006 248); --sidebar-foreground: oklch(0.200 0.012 260);
--sidebar-primary: oklch(0.680 0.165 55); --sidebar-primary: oklch(0.680 0.165 55);
--sidebar-primary-foreground: oklch(0.980 0.005 80); --sidebar-primary-foreground: oklch(0.980 0.005 80);
--sidebar-accent: oklch(0.180 0.010 260); --sidebar-accent: oklch(0.960 0.006 248);
--sidebar-accent-foreground: oklch(0.920 0.006 248); --sidebar-accent-foreground: oklch(0.200 0.012 260);
--sidebar-border: oklch(0.200 0.010 260); --sidebar-border: oklch(0.920 0.006 248);
--sidebar-ring: oklch(0.680 0.165 55); --sidebar-ring: oklch(0.680 0.165 55);
} }
/* ══════════════════════════════════════════════════════════════════════════
DARK THEME — activated via .dark class on <html> by VueUse useDark()
══════════════════════════════════════════════════════════════════════════ */
.dark {
color-scheme: dark;
--background: oklch(0.095 0.010 260);
--foreground: oklch(0.930 0.006 248);
--card: oklch(0.130 0.010 260);
--card-foreground: oklch(0.930 0.006 248);
--popover: oklch(0.130 0.010 260);
--popover-foreground: oklch(0.930 0.006 248);
--primary: oklch(0.720 0.160 55);
--primary-foreground: oklch(0.100 0.008 260);
--secondary: oklch(0.190 0.012 260);
--secondary-foreground: oklch(0.930 0.006 248);
--muted: oklch(0.190 0.010 260);
--muted-foreground: oklch(0.580 0.012 258);
--accent: oklch(0.200 0.015 60);
--accent-foreground: oklch(0.930 0.006 248);
--destructive: oklch(0.520 0.220 27);
--destructive-foreground: oklch(0.930 0.006 248);
--border: oklch(0.220 0.012 260);
--input: oklch(0.220 0.012 260);
--ring: oklch(0.720 0.160 55);
--sidebar: oklch(0.110 0.010 260);
--sidebar-foreground: oklch(0.900 0.006 248);
--sidebar-primary: oklch(0.720 0.160 55);
--sidebar-primary-foreground: oklch(0.100 0.008 260);
--sidebar-accent: oklch(0.180 0.012 260);
--sidebar-accent-foreground: oklch(0.900 0.006 248);
--sidebar-border: oklch(0.210 0.012 260);
--sidebar-ring: oklch(0.720 0.160 55);
}
/* ══════════════════════════════════════════════════════════════════════════ /* ══════════════════════════════════════════════════════════════════════════
BASE LAYER BASE LAYER
══════════════════════════════════════════════════════════════════════════ */ ══════════════════════════════════════════════════════════════════════════ */
......
"""
Content Approval REST API Route
POST /api/content — tạo content mới (AI agent)
GET /api/content — list content (filter by status)
GET /api/content/{id} — chi tiết content
POST /api/content/{id}/submit — submit for review
POST /api/content/{id}/approve — marketing approve
POST /api/content/{id}/reject — marketing reject
POST /api/content/{id}/changes — request changes
POST /api/content/{id}/schedule — assign vào queue
POST /api/content/{id}/publish-now — publish ngay lập tức
GET /api/content/stats — approval stats
"""
import logging
from typing import Optional
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from common.social.approval_gate import (
approve_content, create_content, get_approval_stats,
get_content, list_contents, reject_content,
request_changes, schedule_content,
)
from common.social.post_queue import add_to_queue, assign_queue_slots
from common.social.scheduler import manual_publish
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/content", tags=["Content Approval"])
# ─── Request models ───────────────────────────────────────────────────────────
class CreateContentBody(BaseModel):
title: str
caption: str
platforms: list[str]
image_urls: list[str] = []
tags: list[str] = []
created_by: str = "ai_agent"
metadata: dict = {}
auto_submit: bool = False # Nếu True: tạo xong → submit ngay
class ReviewBody(BaseModel):
actor: str = "marketing_team"
comment: str = ""
class ScheduleBody(BaseModel):
scheduled_at: str
priority: bool = False
# ─── Endpoints ───────────────────────────────────────────────────────────────
@router.post("")
async def create_new_content(body: CreateContentBody):
"""Tạo content mới. Thường được AI stylist engine gọi."""
content = create_content(
title=body.title,
caption=body.caption,
platforms=body.platforms,
image_urls=body.image_urls,
created_by=body.created_by,
tags=body.tags,
metadata=body.metadata,
)
if body.auto_submit:
from common.social.approval_gate import submit_for_review
content = submit_for_review(content["id"], submitted_by=body.created_by)
return {"ok": True, "content": content}
@router.get("")
async def list_all_contents(
status: Optional[str] = None,
platform: Optional[str] = None,
limit: int = 50,
offset: int = 0,
):
"""List tất cả content với filter."""
return list_contents(status=status, platform=platform, limit=limit, offset=offset)
@router.get("/stats")
async def approval_stats():
"""Thống kê trạng thái content."""
return get_approval_stats()
@router.get("/{content_id}")
async def get_content_detail(content_id: str):
content = get_content(content_id)
if not content:
raise HTTPException(status_code=404, detail="Content not found")
return content
@router.post("/{content_id}/submit")
async def submit_content(content_id: str, body: ReviewBody):
from common.social.approval_gate import submit_for_review
try:
content = submit_for_review(content_id, submitted_by=body.actor)
return {"ok": True, "content": content}
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/{content_id}/approve")
async def approve(content_id: str, body: ReviewBody):
try:
content = approve_content(content_id, approved_by=body.actor, comment=body.comment)
# Tự động thêm vào queue
add_to_queue(content_id)
return {"ok": True, "content": content, "queued": True}
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/{content_id}/reject")
async def reject(content_id: str, body: ReviewBody):
if not body.comment:
raise HTTPException(status_code=422, detail="Phải có lý do khi từ chối.")
try:
content = reject_content(content_id, rejected_by=body.actor, comment=body.comment)
return {"ok": True, "content": content}
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/{content_id}/changes")
async def request_change(content_id: str, body: ReviewBody):
if not body.comment:
raise HTTPException(status_code=422, detail="Phải mô tả cần sửa gì.")
try:
content = request_changes(content_id, requested_by=body.actor, comment=body.comment)
return {"ok": True, "content": content}
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/{content_id}/schedule")
async def schedule(content_id: str, body: ScheduleBody):
try:
content = schedule_content(content_id, scheduled_at=body.scheduled_at)
add_to_queue(content_id, priority=body.priority)
return {"ok": True, "content": content}
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
@router.post("/{content_id}/publish-now")
async def publish_now(content_id: str):
"""Publish ngay lập tức — bypass queue timing."""
result = manual_publish(content_id)
if not result.get("success"):
raise HTTPException(status_code=500, detail=result.get("error", "Publish failed"))
return {"ok": True, "result": result}
"""
Media Library REST API
POST /api/media/upload — upload ảnh/video
GET /api/media — list assets
GET /api/media/{id} — serve raw file
GET /api/media/{id}/thumb — serve thumbnail
GET /api/media/{id}/info — asset metadata
DELETE /api/media/{id} — xóa asset
GET /api/media/platforms — list platform sizes
POST /api/media/{id}/resize — resize theo platform
"""
import logging
import os
from fastapi import APIRouter, File, HTTPException, Query, UploadFile
from fastapi.responses import FileResponse
from common.media.image_processor import (
PLATFORM_SIZES, delete_asset, get_asset,
list_assets, resize_for_platform, save_uploaded_asset,
)
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/media", tags=["Media Library"])
ALLOWED_CONTENT_TYPES = {
"image/jpeg", "image/jpg", "image/png", "image/webp", "image/gif",
"video/mp4", "video/quicktime", "video/webm",
}
MAX_FILE_SIZE = 50 * 1024 * 1024 # 50MB
@router.post("/upload")
async def upload_asset(
file: UploadFile = File(...),
tags: str = Query("", description="Comma-separated tags"),
uploaded_by: str = Query("user"),
):
"""Upload ảnh hoặc video vào media library."""
if file.content_type not in ALLOWED_CONTENT_TYPES:
raise HTTPException(
status_code=415,
detail=f"Unsupported file type: {file.content_type}. Allowed: {', '.join(ALLOWED_CONTENT_TYPES)}",
)
content = await file.read()
if len(content) > MAX_FILE_SIZE:
raise HTTPException(status_code=413, detail="File quá lớn (max 50MB)")
tag_list = [t.strip() for t in tags.split(",") if t.strip()] if tags else []
asset = save_uploaded_asset(
filename=file.filename or "upload",
content_bytes=content,
content_type=file.content_type,
uploaded_by=uploaded_by,
tags=tag_list,
)
return {"ok": True, "asset": asset}
@router.get("")
async def list_media(
content_type: str = Query(None, description="Filter: 'image' or 'video'"),
tag: str = Query(None),
limit: int = 50,
offset: int = 0,
):
"""List tất cả assets trong media library."""
return list_assets(content_type=content_type, tag=tag, limit=limit, offset=offset)
@router.get("/platforms")
async def platform_sizes():
"""Danh sách platform sizes để resize."""
return {
"platforms": [
{"key": k, "width": v[0], "height": v[1]}
for k, v in PLATFORM_SIZES.items()
]
}
@router.get("/{asset_id}/info")
async def asset_info(asset_id: str):
"""Lấy metadata của asset."""
asset = get_asset(asset_id)
if not asset:
raise HTTPException(status_code=404, detail="Asset not found")
return asset
@router.get("/{asset_id}")
async def serve_asset(asset_id: str):
"""Serve file raw."""
asset = get_asset(asset_id)
if not asset:
raise HTTPException(status_code=404, detail="Asset not found")
path = asset.get("path", "")
if not os.path.exists(path):
raise HTTPException(status_code=404, detail="File not found on disk")
return FileResponse(path, media_type=asset.get("content_type", "application/octet-stream"))
@router.get("/{asset_id}/thumb")
async def serve_thumbnail(asset_id: str):
"""Serve thumbnail image."""
asset = get_asset(asset_id)
if not asset:
raise HTTPException(status_code=404, detail="Asset not found")
thumb_path = asset.get("thumb_path", "")
if not thumb_path or not os.path.exists(thumb_path):
# Fallback to original
return await serve_asset(asset_id)
return FileResponse(thumb_path, media_type="image/jpeg")
@router.post("/{asset_id}/resize")
async def resize_asset(asset_id: str, platform: str):
"""Resize asset theo platform spec. Returns resized image bytes."""
from fastapi.responses import Response
asset = get_asset(asset_id)
if not asset:
raise HTTPException(status_code=404, detail="Asset not found")
if platform not in PLATFORM_SIZES:
raise HTTPException(
status_code=400,
detail=f"Unknown platform: {platform}. Available: {list(PLATFORM_SIZES.keys())}",
)
path = asset.get("path", "")
if not os.path.exists(path):
raise HTTPException(status_code=404, detail="File not found on disk")
with open(path, "rb") as f:
content_bytes = f.read()
resized = resize_for_platform(content_bytes, platform)
size = PLATFORM_SIZES[platform]
return Response(
content=resized,
media_type="image/jpeg",
headers={
"Content-Disposition": f'inline; filename="{asset_id}_{platform}.jpg"',
"X-Dimensions": f"{size[0]}x{size[1]}",
},
)
@router.delete("/{asset_id}")
async def delete_media_asset(asset_id: str):
"""Xóa asset khỏi library."""
ok = delete_asset(asset_id)
if not ok:
raise HTTPException(status_code=404, detail="Asset not found")
return {"ok": True}
"""
Notification REST API Route
GET /api/notifications — lấy danh sách (in-app bell)
PATCH /api/notifications/{id}/read — đánh dấu đã đọc
POST /api/notifications/read-all — đánh dấu tất cả đã đọc
GET /api/notifications/unread-count — số chưa đọc (badge)
POST /api/notifications/test — test fire một notification
"""
import logging
from typing import Optional
from fastapi import APIRouter, HTTPException
from fastapi.responses import JSONResponse
from common.notification.engine import (
get_notifications, get_unread_count,
mark_all_read, mark_read, notify, check_negative_spike,
)
from common.notification.events import EventType
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/notifications", tags=["Notifications"])
@router.get("")
async def list_notifications(
limit: int = 50,
unread_only: bool = False,
):
"""Lấy danh sách notifications cho in-app bell."""
notifs = get_notifications(limit=limit, unread_only=unread_only)
return {
"notifications": notifs,
"total": len(notifs),
"unread_count": get_unread_count(),
}
@router.get("/unread-count")
async def unread_count():
"""Số notification chưa đọc — dùng cho badge."""
return {"unread_count": get_unread_count()}
@router.patch("/{notification_id}/read")
async def mark_notification_read(notification_id: str):
"""Đánh dấu một notification đã đọc."""
success = mark_read(notification_id)
if not success:
raise HTTPException(status_code=404, detail="Notification not found")
return {"ok": True, "unread_count": get_unread_count()}
@router.post("/read-all")
async def mark_all_notifications_read():
"""Đánh dấu tất cả notifications đã đọc."""
count = mark_all_read()
return {"ok": True, "marked_read": count}
@router.post("/test")
async def fire_test_notification(body: dict):
"""
Fire một test notification — dùng khi debug/demo.
Body: { "event_type": "negative_spike", "title": "...", "body": "..." }
"""
event_type = body.get("event_type", EventType.NEW_SOCIAL_MESSAGE)
title = body.get("title", "Test Notification")
msg_body = body.get("body", "Đây là test notification từ Canifa AI Dashboard.")
notif = notify(event_type=event_type, title=title, body=msg_body)
return {"ok": True, "notification": notif}
@router.post("/check-spike")
async def check_spike():
"""Manually trigger negative spike check."""
spiked = check_negative_spike()
return {"spiked": spiked, "unread_count": get_unread_count()}
"""
Content Queue & Calendar REST API
GET /api/queue — xem queue hiện tại
POST /api/queue/{content_id} — thêm vào queue
DELETE /api/queue/{content_id} — xóa khỏi queue
POST /api/queue/reorder — drag-drop reorder
GET /api/queue/calendar — calendar view (theo tuần)
GET /api/queue/slots — xem posting slots
PUT /api/queue/slots — update slots
POST /api/queue/slots/reset — reset về default
GET /api/queue/retry-state — trạng thái retry publisher
POST /api/queue/publish-tick — manual trigger publish engine tick
"""
import logging
from typing import Optional
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from common.social.post_queue import (
add_to_queue, get_calendar_week, get_queue,
get_slots, remove_from_queue, reorder_queue,
reset_slots_to_default, save_slots,
)
from common.social.scheduler import get_retry_state, _engine
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/queue", tags=["Content Queue"])
class ReorderBody(BaseModel):
ordered_content_ids: list[str]
class AddToQueueBody(BaseModel):
priority: bool = False
@router.get("")
async def get_current_queue():
"""Lấy queue hiện tại kèm content data."""
queue = get_queue()
return {"queue": queue, "total": len(queue)}
@router.post("/{content_id}")
async def add_content_to_queue(content_id: str, body: AddToQueueBody = AddToQueueBody()):
"""Thêm content vào queue."""
queue = add_to_queue(content_id, priority=body.priority)
return {"ok": True, "queue": queue}
@router.delete("/{content_id}")
async def remove_content_from_queue(content_id: str):
"""Xóa content khỏi queue."""
queue = remove_from_queue(content_id)
return {"ok": True, "queue": queue}
@router.post("/reorder")
async def reorder(body: ReorderBody):
"""Drag-drop reorder — pass danh sách content_id theo thứ tự mới."""
queue = reorder_queue(body.ordered_content_ids)
return {"ok": True, "queue": queue}
@router.get("/calendar")
async def calendar_week(week_offset: int = 0):
"""Calendar view theo tuần. week_offset=0 → tuần này."""
return get_calendar_week(week_offset=week_offset)
@router.get("/slots")
async def get_posting_slots():
"""Lấy posting slots hiện tại."""
return {"slots": get_slots()}
@router.put("/slots")
async def update_slots(slots: list[dict]):
"""Cập nhật posting slots. Input: [{day, hour, minute, label}]"""
saved = save_slots(slots)
return {"ok": True, "slots": saved}
@router.post("/slots/reset")
async def reset_slots():
"""Reset posting slots về default của BrightBean."""
slots = reset_slots_to_default()
return {"ok": True, "slots": slots}
@router.get("/retry-state")
async def retry_state():
"""Xem trạng thái retry của publisher engine."""
return {"retry_state": get_retry_state()}
@router.post("/publish-tick")
async def manual_publish_tick():
"""Manual trigger một publish engine tick (dùng khi debug/testing)."""
import asyncio
loop = asyncio.get_event_loop()
count = await loop.run_in_executor(None, _engine.poll_and_publish)
return {"ok": True, "published_count": count}
"""
Social Inbox API Route
Nhận webhook từ Facebook/Instagram/TikTok → phân tích sentiment → feed vào Canifa learning loop.
Cũng expose REST API để FE dashboard đọc messages.
"""
import json
import logging
import os
from datetime import datetime, timezone
from typing import Optional
from fastapi import APIRouter, BackgroundTasks, Header, HTTPException, Query, Request
from fastapi.responses import JSONResponse, PlainTextResponse
from common.social.inbox_webhook import (
handle_facebook_verify,
parse_facebook_webhook,
parse_tiktok_webhook,
verify_facebook_signature,
)
from common.social.sentiment_vi import analyze_sentiment_detail
from common.feedback_tracker import save_feedback
logger = logging.getLogger(__name__)
router = APIRouter(prefix="/api/social", tags=["Social Inbox"])
# File lưu social messages (lightweight, tương tự feedback_tracker)
SOCIAL_MESSAGES_FILE = os.path.join(
os.path.dirname(os.path.abspath(__file__)), "..", "data", "social_messages.json"
)
# ─────────────────────────────────────────────
# Storage helpers
# ─────────────────────────────────────────────
def _ensure_dir():
os.makedirs(os.path.dirname(SOCIAL_MESSAGES_FILE), exist_ok=True)
def _load_messages() -> list[dict]:
_ensure_dir()
if not os.path.exists(SOCIAL_MESSAGES_FILE):
return []
try:
with open(SOCIAL_MESSAGES_FILE, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
return []
def _save_messages(messages: list[dict]):
_ensure_dir()
with open(SOCIAL_MESSAGES_FILE, "w", encoding="utf-8") as f:
json.dump(messages, f, ensure_ascii=False, indent=2, default=str)
def _append_messages(new_messages: list[dict]):
"""Append new messages, dedup by platform + body."""
existing = _load_messages()
existing_keys = {(m.get("platform"), m.get("body")) for m in existing}
added = 0
for msg in new_messages:
key = (msg.get("platform"), msg.get("body"))
if key not in existing_keys:
existing.insert(0, msg) # newest first
existing_keys.add(key)
added += 1
# Giữ tối đa 500 messages
_save_messages(existing[:500])
return added
# ─────────────────────────────────────────────
# Background: Feed message vào learning loop
# ─────────────────────────────────────────────
def _feed_to_feedback_tracker(messages: list[dict]):
"""
Feed social messages vào Canifa feedback_tracker.
Negative messages → learning loop để AI cải thiện.
"""
for msg in messages:
if msg.get("sentiment") == "negative":
try:
save_feedback(
trace_id=f"social_{msg.get('platform')}_{datetime.now(timezone.utc).timestamp()}",
rating=0, # 0 = dislike (negative feedback)
comment=msg.get("body", ""),
user_category="social_feedback",
ai_category="social_negative_comment",
ai_severity="low",
ai_summary=f"[{msg.get('platform', '').upper()}] Negative comment từ @{msg.get('sender_name', 'Unknown')}: {msg.get('body', '')[:100]}",
)
logger.info(f"📥 Social negative feedback fed to learning loop: {msg.get('platform')} / {msg.get('body', '')[:50]}")
except Exception as e:
logger.error(f"Failed to feed to feedback tracker: {e}")
# ─────────────────────────────────────────────
# Facebook/Instagram Webhook
# ─────────────────────────────────────────────
@router.get("/webhook/facebook")
async def facebook_webhook_verify(
hub_mode: str = Query(None, alias="hub.mode"),
hub_verify_token: str = Query(None, alias="hub.verify_token"),
hub_challenge: str = Query("", alias="hub.challenge"),
):
"""Facebook webhook verification (GET handshake)."""
challenge = handle_facebook_verify(hub_mode or "", hub_verify_token or "", hub_challenge)
if challenge is None:
raise HTTPException(status_code=403, detail="Webhook verification failed")
return PlainTextResponse(challenge)
@router.post("/webhook/facebook")
async def facebook_webhook_receive(
request: Request,
background_tasks: BackgroundTasks,
x_hub_signature_256: Optional[str] = Header(None, alias="X-Hub-Signature-256"),
):
"""Nhận Facebook/Instagram comment/DM webhook events."""
body = await request.body()
# Verify signature
if not verify_facebook_signature(body, x_hub_signature_256 or ""):
raise HTTPException(status_code=403, detail="Invalid webhook signature")
try:
payload = await request.json()
except Exception:
raise HTTPException(status_code=400, detail="Invalid JSON")
messages = parse_facebook_webhook(payload)
if messages:
added = _append_messages(messages)
background_tasks.add_task(_feed_to_feedback_tracker, messages)
logger.info(f"📬 Facebook/IG webhook: {len(messages)} events, {added} mới")
return JSONResponse({"ok": True, "received": len(messages)})
# ─────────────────────────────────────────────
# TikTok Webhook
# ─────────────────────────────────────────────
@router.post("/webhook/tiktok")
async def tiktok_webhook_receive(
request: Request,
background_tasks: BackgroundTasks,
):
"""Nhận TikTok comment webhook events."""
try:
payload = await request.json()
except Exception:
raise HTTPException(status_code=400, detail="Invalid JSON")
messages = parse_tiktok_webhook(payload)
if messages:
added = _append_messages(messages)
background_tasks.add_task(_feed_to_feedback_tracker, messages)
logger.info(f"📬 TikTok webhook: {len(messages)} comments, {added} mới")
return JSONResponse({"ok": True, "received": len(messages)})
# ─────────────────────────────────────────────
# REST API cho FE Dashboard
# ─────────────────────────────────────────────
@router.get("/messages")
async def get_social_messages(
platform: Optional[str] = Query(None, description="facebook | instagram | tiktok"),
sentiment: Optional[str] = Query(None, description="positive | negative | neutral"),
limit: int = Query(50, ge=1, le=200),
offset: int = Query(0, ge=0),
):
"""
Lấy danh sách social messages cho dashboard.
Filter theo platform và/hoặc sentiment.
"""
messages = _load_messages()
if platform:
messages = [m for m in messages if m.get("platform") == platform]
if sentiment:
messages = [m for m in messages if m.get("sentiment") == sentiment]
total = len(messages)
paginated = messages[offset: offset + limit]
return {
"total": total,
"limit": limit,
"offset": offset,
"messages": paginated,
}
@router.get("/messages/stats")
async def get_social_stats():
"""Thống kê tổng hợp social messages theo platform + sentiment."""
messages = _load_messages()
total = len(messages)
stats: dict = {
"total": total,
"by_platform": {},
"by_sentiment": {"positive": 0, "negative": 0, "neutral": 0},
"negative_rate": 0.0,
}
for msg in messages:
platform = msg.get("platform", "unknown")
sentiment = msg.get("sentiment", "neutral")
if platform not in stats["by_platform"]:
stats["by_platform"][platform] = {"total": 0, "positive": 0, "negative": 0, "neutral": 0}
stats["by_platform"][platform]["total"] += 1
stats["by_platform"][platform][sentiment] = stats["by_platform"][platform].get(sentiment, 0) + 1
stats["by_sentiment"][sentiment] = stats["by_sentiment"].get(sentiment, 0) + 1
if total > 0:
stats["negative_rate"] = round(stats["by_sentiment"]["negative"] / total * 100, 1)
return stats
@router.post("/messages/analyze")
async def analyze_text_sentiment(request: Request):
"""
Test sentiment analysis — nhập text, trả về phân tích.
Dùng cho debug / demo trên dashboard.
"""
body = await request.json()
text = body.get("text", "")
if not text:
raise HTTPException(status_code=400, detail="text is required")
result = analyze_sentiment_detail(text)
return result
@router.delete("/messages/clear")
async def clear_social_messages():
"""Xóa toàn bộ messages (dùng khi debug)."""
_save_messages([])
return {"ok": True, "message": "Đã xóa tất cả social messages"}
"""
API route: Content Templates + RSS Feeds
Ported từ BrightBean apps/composer — cung cấp 30+ content templates và RSS feed catalog
"""
from fastapi import APIRouter, Query
from common.content_templates import (
get_templates, get_template, get_categories, search_templates,
get_feeds, get_feed_categories,
)
router = APIRouter(prefix="/api/templates", tags=["templates"])
@router.get("")
def list_templates(category: str = "all", limit: int = 50):
"""Liệt kê content templates theo category."""
return {
"templates": get_templates(category=category, limit=limit),
"categories": get_categories(),
}
@router.get("/categories")
def list_categories():
return {"categories": get_categories()}
@router.get("/search")
def search(q: str = Query(..., min_length=2)):
return {"templates": search_templates(q)}
@router.get("/{template_id}")
def get_one(template_id: int):
t = get_template(template_id)
if not t:
from fastapi import HTTPException
raise HTTPException(404, "Template not found")
return t
@router.get("/feeds/categories")
def feed_categories():
return {"categories": get_feed_categories()}
@router.get("/feeds/{category}")
def feeds_for_category(category: str):
return {"feeds": get_feeds(category)}
This diff is collapsed.
"""Media package."""
from .image_processor import (
generate_image_thumbnail, resize_for_platform, get_image_dimensions,
generate_video_thumbnail, get_video_metadata,
save_uploaded_asset, list_assets, get_asset, delete_asset,
PLATFORM_SIZES,
)
__all__ = [
"generate_image_thumbnail", "resize_for_platform", "get_image_dimensions",
"generate_video_thumbnail", "get_video_metadata",
"save_uploaded_asset", "list_assets", "get_asset", "delete_asset",
"PLATFORM_SIZES",
]
"""
Image & Video Media Processor
Adapted từ BrightBean Studio apps/media_library/services.py
Xử lý ảnh sản phẩm Canifa cho social media: resize/crop/thumbnail
"""
import io
import json
import logging
import os
import subprocess
import uuid
from datetime import datetime, timezone
from pathlib import Path
from typing import Optional
logger = logging.getLogger(__name__)
MEDIA_DIR = os.path.join(
os.path.dirname(os.path.abspath(__file__)), "..", "..", "data", "media"
)
MEDIA_INDEX_FILE = os.path.join(MEDIA_DIR, "index.json")
# Platform-specific output sizes (px)
PLATFORM_SIZES = {
"instagram_square": (1080, 1080), # Instagram 1:1
"instagram_portrait": (1080, 1350), # Instagram 4:5
"instagram_story": (1080, 1920), # Instagram 9:16 / Story
"tiktok": (1080, 1920), # TikTok 9:16
"facebook": (1200, 630), # Facebook link preview
"facebook_square": (1080, 1080), # Facebook post
"thumbnail": (400, 400), # Dashboard thumbnail
}
def _ensure_dir():
os.makedirs(MEDIA_DIR, exist_ok=True)
thumbs_dir = os.path.join(MEDIA_DIR, "thumbs")
os.makedirs(thumbs_dir, exist_ok=True)
def _load_index() -> list[dict]:
_ensure_dir()
if not os.path.exists(MEDIA_INDEX_FILE):
return []
try:
with open(MEDIA_INDEX_FILE, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
return []
def _save_index(data: list[dict]):
with open(MEDIA_INDEX_FILE, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2, default=str)
def _append_asset(asset: dict):
index = _load_index()
index.insert(0, asset)
_save_index(index[:1000])
# ─── Image Processing (requires Pillow) ──────────────────────────────────────
def generate_image_thumbnail(
image_bytes: bytes,
size: tuple[int, int] = (400, 400),
output_format: str = "JPEG",
) -> bytes:
"""
Generate thumbnail từ image bytes.
Từ BrightBean: LANCZOS resize, RGBA→RGB convert để đảm bảo JPEG compat.
"""
try:
from PIL import Image
except ImportError:
logger.warning("Pillow không có — skip thumbnail generation")
return image_bytes
img = Image.open(io.BytesIO(image_bytes))
# Convert RGBA, P modes sang RGB (JPEG không hỗ trợ)
if img.mode in ("RGBA", "P", "LA"):
background = Image.new("RGB", img.size, (255, 255, 255))
if img.mode == "RGBA":
background.paste(img, mask=img.split()[3])
else:
background.paste(img)
img = background
elif img.mode != "RGB":
img = img.convert("RGB")
# LANCZOS (Antialiased resize — chất lượng tốt nhất)
img.thumbnail(size, Image.LANCZOS)
out = io.BytesIO()
img.save(out, format=output_format, quality=85, optimize=True)
return out.getvalue()
def resize_for_platform(
image_bytes: bytes,
platform: str,
) -> bytes:
"""
Resize và crop ảnh theo đúng kích thước platform.
Canifa use case: ảnh sản phẩm → fit cho từng social platform.
"""
try:
from PIL import Image, ImageOps
except ImportError:
return image_bytes
size = PLATFORM_SIZES.get(platform)
if not size:
return image_bytes
img = Image.open(io.BytesIO(image_bytes))
if img.mode not in ("RGB", "L"):
img = img.convert("RGB")
# ImageOps.fit: crop center + resize giữ nguyên tỷ lệ
img = ImageOps.fit(img, size, method=Image.LANCZOS, centering=(0.5, 0.5))
out = io.BytesIO()
img.save(out, format="JPEG", quality=90, optimize=True)
return out.getvalue()
def get_image_dimensions(image_bytes: bytes) -> dict:
"""Extract width/height từ image bytes."""
try:
from PIL import Image
img = Image.open(io.BytesIO(image_bytes))
return {"width": img.width, "height": img.height, "mode": img.mode}
except Exception:
return {"width": 0, "height": 0, "mode": "unknown"}
# ─── Video Processing (requires ffmpeg) ──────────────────────────────────────
def generate_video_thumbnail(video_path: str, timestamp_seconds: float = 1.0) -> bytes | None:
"""
Extract frame từ video dùng ffmpeg.
Từ BrightBean: dùng ffmpeg -ss {ts} -vframes 1 -f image2pipe.
"""
cmd = [
"ffmpeg",
"-ss", str(timestamp_seconds),
"-i", video_path,
"-vframes", "1",
"-f", "image2",
"-vf", "scale=400:400:force_original_aspect_ratio=decrease",
"pipe:1",
"-loglevel", "quiet",
]
try:
result = subprocess.run(cmd, capture_output=True, timeout=30)
if result.returncode == 0 and result.stdout:
return result.stdout
except (FileNotFoundError, subprocess.TimeoutExpired) as e:
logger.warning("ffmpeg thumbnail failed: %s", e)
return None
def get_video_metadata(video_path: str) -> dict:
"""
Extract video metadata via ffprobe.
Từ BrightBean: parse JSON output từ ffprobe.
"""
cmd = [
"ffprobe",
"-v", "quiet",
"-print_format", "json",
"-show_streams",
"-show_format",
video_path,
]
try:
result = subprocess.run(cmd, capture_output=True, text=True, timeout=15)
if result.returncode != 0:
return {}
data = json.loads(result.stdout)
video_stream = next(
(s for s in data.get("streams", []) if s.get("codec_type") == "video"),
{},
)
fmt = data.get("format", {})
return {
"duration": float(fmt.get("duration", 0)),
"width": video_stream.get("width", 0),
"height": video_stream.get("height", 0),
"codec": video_stream.get("codec_name", ""),
"size_bytes": int(fmt.get("size", 0)),
}
except Exception as e:
logger.warning("ffprobe Error: %s", e)
return {}
# ─── Asset Management ─────────────────────────────────────────────────────────
def save_uploaded_asset(
filename: str,
content_bytes: bytes,
content_type: str,
uploaded_by: str = "user",
tags: list[str] | None = None,
) -> dict:
"""
Lưu asset vào media library, generate thumbnail.
Returns asset metadata dict.
"""
_ensure_dir()
ext = Path(filename).suffix.lower()
asset_id = uuid.uuid4().hex
stored_name = f"{asset_id}{ext}"
asset_path = os.path.join(MEDIA_DIR, stored_name)
with open(asset_path, "wb") as f:
f.write(content_bytes)
is_image = content_type.startswith("image/")
is_video = content_type.startswith("video/")
# Generate thumbnail
thumb_path = None
try:
if is_image:
thumb_bytes = generate_image_thumbnail(content_bytes)
thumb_path = os.path.join(MEDIA_DIR, "thumbs", f"{asset_id}.jpg")
with open(thumb_path, "wb") as f:
f.write(thumb_bytes)
elif is_video:
thumb_bytes = generate_video_thumbnail(asset_path)
if thumb_bytes:
thumb_path = os.path.join(MEDIA_DIR, "thumbs", f"{asset_id}.jpg")
with open(thumb_path, "wb") as f:
f.write(thumb_bytes)
except Exception as e:
logger.warning("Thumbnail generation failed: %s", e)
# Metadata
if is_image:
dims = get_image_dimensions(content_bytes)
elif is_video:
dims = get_video_metadata(asset_path)
else:
dims = {}
asset = {
"id": asset_id,
"filename": filename,
"stored_name": stored_name,
"content_type": content_type,
"size_bytes": len(content_bytes),
"path": asset_path,
"thumb_path": thumb_path,
"url": f"/api/media/{asset_id}",
"thumb_url": f"/api/media/{asset_id}/thumb" if thumb_path else None,
"uploaded_by": uploaded_by,
"tags": tags or [],
"metadata": dims,
"created_at": datetime.now(timezone.utc).isoformat(),
}
_append_asset(asset)
logger.info("📁 Asset saved: %s (%d bytes)", filename, len(content_bytes))
return asset
def list_assets(
content_type: str | None = None,
tag: str | None = None,
limit: int = 50,
offset: int = 0,
) -> dict:
assets = _load_index()
if content_type:
assets = [a for a in assets if a.get("content_type", "").startswith(content_type)]
if tag:
assets = [a for a in assets if tag in a.get("tags", [])]
total = len(assets)
return {"total": total, "assets": assets[offset: offset + limit]}
def get_asset(asset_id: str) -> dict | None:
for a in _load_index():
if a.get("id") == asset_id:
return a
return None
def delete_asset(asset_id: str) -> bool:
index = _load_index()
asset = next((a for a in index if a.get("id") == asset_id), None)
if not asset:
return False
try:
if os.path.exists(asset.get("path", "")):
os.remove(asset["path"])
if asset.get("thumb_path") and os.path.exists(asset["thumb_path"]):
os.remove(asset["thumb_path"])
except OSError as e:
logger.warning("delete_asset: %s", e)
_save_index([a for a in index if a.get("id") != asset_id])
return True
"""Notification package."""
from .engine import notify, get_notifications, mark_read, mark_all_read, get_unread_count, check_negative_spike
from .events import EventType, Channel
__all__ = [
"notify", "get_notifications", "mark_read", "mark_all_read",
"get_unread_count", "check_negative_spike",
"EventType", "Channel",
]
This diff is collapsed.
"""
Notification Events — Constants
Ported & adapted từ BrightBean Studio apps/notifications/models.py
"""
class EventType:
# Social feedback events
NEW_SOCIAL_MESSAGE = "new_social_message"
NEGATIVE_SPIKE = "negative_spike" # Negative rate vượt ngưỡng
SOCIAL_SENTIMENT_REPORT= "social_sentiment_report" # Báo cáo hàng tuần
# AI Content pipeline
CONTENT_SUBMITTED = "content_submitted" # AI tạo content chờ review
CONTENT_APPROVED = "content_approved" # Marketing đã approve
CONTENT_REJECTED = "content_rejected" # Marketing reject
CONTENT_CHANGES = "content_changes_requested"
CONTENT_PUBLISHED = "content_published" # Đã publish lên social
CONTENT_FAILED = "content_failed" # Publish thất bại
# Learning loop events
AI_RULES_UPDATED = "ai_rules_updated" # Feedback loop cập nhật rules
FEEDBACK_MILESTONE = "feedback_milestone" # Đạt N feedbacks
# System
PUBLISH_RETRY = "publish_retry" # Đang retry publish
API_TOKEN_EXPIRED = "api_token_expired" # Token FB/IG hết hạn
ALL = [
NEW_SOCIAL_MESSAGE, NEGATIVE_SPIKE, SOCIAL_SENTIMENT_REPORT,
CONTENT_SUBMITTED, CONTENT_APPROVED, CONTENT_REJECTED,
CONTENT_CHANGES, CONTENT_PUBLISHED, CONTENT_FAILED,
AI_RULES_UPDATED, FEEDBACK_MILESTONE,
PUBLISH_RETRY, API_TOKEN_EXPIRED,
]
class Channel:
IN_APP = "in_app"
EMAIL = "email"
WEBHOOK = "webhook"
# Default channels per event type
DEFAULT_CHANNELS: dict[str, dict[str, bool]] = {
EventType.NEW_SOCIAL_MESSAGE: {Channel.IN_APP: True, Channel.EMAIL: False, Channel.WEBHOOK: False},
EventType.NEGATIVE_SPIKE: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: True},
EventType.SOCIAL_SENTIMENT_REPORT: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: False},
EventType.CONTENT_SUBMITTED: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: False},
EventType.CONTENT_APPROVED: {Channel.IN_APP: True, Channel.EMAIL: False, Channel.WEBHOOK: False},
EventType.CONTENT_REJECTED: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: False},
EventType.CONTENT_CHANGES: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: False},
EventType.CONTENT_PUBLISHED: {Channel.IN_APP: True, Channel.EMAIL: False, Channel.WEBHOOK: True},
EventType.CONTENT_FAILED: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: True},
EventType.AI_RULES_UPDATED: {Channel.IN_APP: True, Channel.EMAIL: False, Channel.WEBHOOK: False},
EventType.FEEDBACK_MILESTONE: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: False},
EventType.PUBLISH_RETRY: {Channel.IN_APP: True, Channel.EMAIL: False, Channel.WEBHOOK: False},
EventType.API_TOKEN_EXPIRED: {Channel.IN_APP: True, Channel.EMAIL: True, Channel.WEBHOOK: False},
}
# Non-critical: suppress khi quiet hours
NON_CRITICAL = {
EventType.NEW_SOCIAL_MESSAGE,
EventType.AI_RULES_UPDATED,
EventType.CONTENT_PUBLISHED,
EventType.SOCIAL_SENTIMENT_REPORT,
}
"""
Social Media Integration Layer
Ported & adapted từ BrightBean Studio (open-source)
Tích hợp: Nhận feedback từ Facebook/Instagram/TikTok comment → Canifa AI Learning Loop
"""
from .sentiment_vi import analyze_sentiment
from .inbox_webhook import parse_facebook_webhook, parse_tiktok_webhook
__all__ = ["analyze_sentiment", "parse_facebook_webhook", "parse_tiktok_webhook"]
This diff is collapsed.
"""
Social Inbox Webhook Parser
Ported & adapted từ BrightBean Studio apps/inbox/webhooks.py
Chuyển từ Django → FastAPI/asyncio, bỏ ORM dependency, thuần Python
Hỗ trợ: Facebook/Instagram webhook, TikTok comment webhook
Output: Chuẩn hóa thành SocialMessage dict → feed vào Canifa feedback_tracker
"""
import hashlib
import hmac
import json
import logging
import os
from datetime import datetime, timezone
from typing import Any
from .sentiment_vi import analyze_sentiment
logger = logging.getLogger(__name__)
# Lấy từ env (set trong .env của Canifa backend)
FACEBOOK_APP_SECRET = os.getenv("FACEBOOK_APP_SECRET", "")
FACEBOOK_WEBHOOK_VERIFY_TOKEN = os.getenv("FACEBOOK_WEBHOOK_VERIFY_TOKEN", "canifa_verify_token")
# ─────────────────────────────────────────────
# Data Model (simple dict — không cần ORM)
# ─────────────────────────────────────────────
def _build_social_message(
platform: str, # "facebook" | "instagram" | "tiktok"
message_type: str, # "comment" | "dm" | "mention"
sender_name: str,
sender_id: str,
body: str,
post_id: str = "",
raw: dict | None = None,
) -> dict:
"""Tạo chuẩn hóa SocialMessage dict để feed vào feedback_tracker."""
sentiment = analyze_sentiment(body)
return {
"source": "social",
"platform": platform,
"message_type": message_type,
"sender_name": sender_name,
"sender_id": sender_id,
"body": body,
"post_id": post_id,
"sentiment": sentiment, # "positive" | "negative" | "neutral"
"received_at": datetime.now(timezone.utc).isoformat(),
"raw": raw or {},
}
# ─────────────────────────────────────────────
# Facebook / Instagram Webhook Handler
# (Adapted từ BrightBean webhooks.py)
# ─────────────────────────────────────────────
def verify_facebook_signature(body: bytes, signature_header: str) -> bool:
"""Verify HMAC-SHA256 signature từ Facebook. Bảo vệ khỏi webhook giả."""
if not FACEBOOK_APP_SECRET:
logger.warning("⚠️ FACEBOOK_APP_SECRET chưa set — bỏ qua verify (dev mode)")
return True # Dev mode: không verify
expected = "sha256=" + hmac.new(
FACEBOOK_APP_SECRET.encode(),
body,
hashlib.sha256,
).hexdigest()
return hmac.compare_digest(expected, signature_header)
def parse_facebook_webhook(payload: dict) -> list[dict]:
"""
Parse Facebook/Instagram webhook payload → list[SocialMessage].
Hỗ trợ: comments, mentions, DM (messaging events).
"""
messages = []
for entry in payload.get("entry", []):
# Facebook Page comments & mentions
for change in entry.get("changes", []):
field = change.get("field", "")
value = change.get("value", {})
if field == "feed":
msg = _parse_fb_feed_change(value)
if msg:
messages.append(msg)
elif field == "mention":
msg = _parse_fb_mention(value)
if msg:
messages.append(msg)
# Instagram / Facebook DM (messaging events)
for messaging in entry.get("messaging", []):
msg = _parse_fb_messaging(messaging)
if msg:
messages.append(msg)
return messages
def _parse_fb_feed_change(value: dict) -> dict | None:
"""Parse một Facebook comment từ feed change event."""
comment_id = value.get("comment_id") or value.get("id")
if not comment_id:
return None
from_data = value.get("from", {})
body = value.get("message", "").strip()
if not body:
return None
post_id = value.get("post_id", "")
platform = "instagram" if value.get("item") == "comment" and "instagram" in str(value).lower() else "facebook"
return _build_social_message(
platform=platform,
message_type="comment",
sender_name=from_data.get("name", "Unknown"),
sender_id=from_data.get("id", ""),
body=body,
post_id=post_id,
raw=value,
)
def _parse_fb_mention(value: dict) -> dict | None:
"""Parse Facebook mention event."""
mention_id = value.get("post_id") or value.get("id")
if not mention_id:
return None
from_data = value.get("from", {})
body = value.get("message", "").strip()
if not body:
return None
return _build_social_message(
platform="facebook",
message_type="mention",
sender_name=from_data.get("name", "Unknown"),
sender_id=from_data.get("id", ""),
body=body,
post_id=mention_id,
raw=value,
)
def _parse_fb_messaging(messaging: dict) -> dict | None:
"""Parse Facebook/Instagram DM message."""
message_data = messaging.get("message", {})
mid = message_data.get("mid")
if not mid:
return None
sender = messaging.get("sender", {})
body = message_data.get("text", "").strip()
if not body:
return None # Bỏ qua message ảnh/sticker không có text
return _build_social_message(
platform="instagram",
message_type="dm",
sender_name=sender.get("name", sender.get("id", "Unknown")),
sender_id=sender.get("id", ""),
body=body,
post_id="",
raw=messaging,
)
# ─────────────────────────────────────────────
# TikTok Webhook Handler
# TikTok gửi comment events qua webhook
# ─────────────────────────────────────────────
def parse_tiktok_webhook(payload: dict) -> list[dict]:
"""
Parse TikTok comment webhook payload → list[SocialMessage].
TikTok webhook format khác Facebook.
"""
messages = []
# TikTok comment notification
comments = payload.get("data", {}).get("comments", [])
for comment in comments:
body = comment.get("text", "").strip()
if not body:
continue
user = comment.get("user", {})
messages.append(_build_social_message(
platform="tiktok",
message_type="comment",
sender_name=user.get("display_name", user.get("unique_id", "Unknown")),
sender_id=user.get("open_id", ""),
body=body,
post_id=comment.get("video_id", ""),
raw=comment,
))
return messages
# ─────────────────────────────────────────────
# Facebook Webhook Verification (for GET requests)
# ─────────────────────────────────────────────
def handle_facebook_verify(mode: str, token: str, challenge: str) -> str | None:
"""
Xử lý Facebook webhook verification handshake (GET request).
Returns challenge string nếu valid, None nếu invalid.
"""
if mode == "subscribe" and token == FACEBOOK_WEBHOOK_VERIFY_TOKEN:
logger.info("✅ Facebook webhook verified!")
return challenge
logger.warning("❌ Facebook webhook verification failed. Token mismatch.")
return None
"""
Content Post Queue + Calendar Scheduler
Ported & adapted từ BrightBean Studio apps/calendar/services.py
Logic: define PostingSlots theo weekday → auto-assign datetime cho queued content
"""
import json
import logging
import os
from datetime import datetime, time, timedelta, timezone
from typing import Optional
logger = logging.getLogger(__name__)
QUEUE_FILE = os.path.join(
os.path.dirname(os.path.abspath(__file__)), "..", "..", "data", "post_queue.json"
)
SLOTS_FILE = os.path.join(
os.path.dirname(os.path.abspath(__file__)), "..", "..", "data", "posting_slots.json"
)
# ─── Default slots (từ BrightBean — non-round times để tránh shadow ban) ────
# Canifa: Đăng T2/T3/T4/T5 vào giờ "vàng" sáng
DEFAULT_SLOTS = [
{"day": 0, "hour": 9, "minute": 24, "label": "Thứ 2 sáng"}, # Monday
{"day": 0, "hour": 11, "minute": 10, "label": "Thứ 2 trưa"},
{"day": 1, "hour": 9, "minute": 55, "label": "Thứ 3 sáng"}, # Tuesday
{"day": 2, "hour": 9, "minute": 30, "label": "Thứ 4 sáng"}, # Wednesday
{"day": 2, "hour": 11, "minute": 32, "label": "Thứ 4 trưa"},
{"day": 3, "hour": 9, "minute": 38, "label": "Thứ 5 sáng"}, # Thursday
{"day": 4, "hour": 10, "minute": 17, "label": "Thứ 6 sáng"}, # Friday
{"day": 4, "hour": 14, "minute": 0, "label": "Thứ 6 chiều"},
{"day": 5, "hour": 10, "minute": 0, "label": "Thứ 7 sáng"}, # Saturday
]
# ─── Storage ─────────────────────────────────────────────────────────────────
def _ensure_dir():
for f in [QUEUE_FILE, SLOTS_FILE]:
os.makedirs(os.path.dirname(f), exist_ok=True)
def _load_queue() -> list[dict]:
_ensure_dir()
if not os.path.exists(QUEUE_FILE):
return []
try:
with open(QUEUE_FILE, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
return []
def _save_queue(data: list[dict]):
_ensure_dir()
with open(QUEUE_FILE, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2, default=str)
def _load_slots() -> list[dict]:
_ensure_dir()
if not os.path.exists(SLOTS_FILE):
_save_slots(DEFAULT_SLOTS)
return DEFAULT_SLOTS
try:
with open(SLOTS_FILE, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
return DEFAULT_SLOTS
def _save_slots(data: list[dict]):
_ensure_dir()
with open(SLOTS_FILE, "w", encoding="utf-8") as f:
json.dump(data, f, ensure_ascii=False, indent=2, default=str)
# ─── Core slot logic (from BrightBean calendar/services.py) ─────────────────
def _next_slot_datetimes(after_dt: datetime, count: int = 20) -> list[datetime]:
"""
Tính toán N thời điểm publish tiếp theo dựa trên PostingSlots.
Walk 60 ngày tới để tìm đủ slots.
"""
slots = _load_slots()
if not slots:
return []
results = []
current_date = after_dt.date()
for day_offset in range(60):
check_date = current_date + timedelta(days=day_offset)
weekday = check_date.weekday() # 0=Monday
for slot in slots:
if slot.get("day") != weekday:
continue
slot_dt = datetime.combine(
check_date,
time(slot.get("hour", 9), slot.get("minute", 0)),
tzinfo=timezone.utc,
)
if slot_dt <= after_dt:
continue
results.append(slot_dt)
if len(results) >= count:
return results
return results
def assign_queue_slots() -> list[dict]:
"""
Recalculate scheduled_at cho tất cả queued content.
Gọi sau mỗi khi add/remove/reorder queue.
"""
queue = _load_queue()
if not queue:
return []
now = datetime.now(timezone.utc)
slot_times = _next_slot_datetimes(now, count=len(queue) + 5)
for idx, entry in enumerate(queue):
if idx < len(slot_times):
entry["scheduled_at"] = slot_times[idx].isoformat()
else:
entry["scheduled_at"] = None
_save_queue(queue)
return queue
def add_to_queue(content_id: str, priority: bool = False) -> list[dict]:
"""
Thêm content vào queue.
priority=True → đẩy lên đầu (cho campaign quan trọng).
"""
queue = _load_queue()
# Remove nếu đã có
queue = [e for e in queue if e.get("content_id") != content_id]
entry = {
"content_id": content_id,
"position": 0,
"added_at": datetime.now(timezone.utc).isoformat(),
"scheduled_at": None,
}
if priority:
queue.insert(0, entry)
else:
queue.append(entry)
# Update positions
for i, e in enumerate(queue):
e["position"] = i
_save_queue(queue)
return assign_queue_slots()
def remove_from_queue(content_id: str) -> list[dict]:
"""Xoá content khỏi queue."""
queue = _load_queue()
queue = [e for e in queue if e.get("content_id") != content_id]
for i, e in enumerate(queue):
e["position"] = i
_save_queue(queue)
return assign_queue_slots()
def reorder_queue(ordered_content_ids: list[str]) -> list[dict]:
"""Reorder queue theo danh sách content_id (drag-drop)."""
queue = _load_queue()
queue_map = {e["content_id"]: e for e in queue}
new_queue = []
for idx, cid in enumerate(ordered_content_ids):
if cid in queue_map:
entry = queue_map[cid]
entry["position"] = idx
new_queue.append(entry)
# Append entries không có trong ordered list
existing_ids = set(ordered_content_ids)
for e in queue:
if e["content_id"] not in existing_ids:
new_queue.append(e)
_save_queue(new_queue)
return assign_queue_slots()
def get_queue() -> list[dict]:
"""Lấy queue hiện tại với thông tin content đính kèm."""
queue = _load_queue()
# Enrich với content data
try:
from common.social.approval_gate import get_content
enriched = []
for entry in queue:
content = get_content(entry["content_id"])
enriched.append({**entry, "content": content})
return enriched
except Exception:
return queue
def get_due_for_publish(now: datetime | None = None) -> list[dict]:
"""
Lấy danh sách queue entries đã đến giờ publish.
Dùng bởi PublishEngine để biết publish cái gì.
"""
if now is None:
now = datetime.now(timezone.utc)
queue = _load_queue()
due = []
for entry in queue:
scheduled = entry.get("scheduled_at")
if not scheduled:
continue
scheduled_dt = datetime.fromisoformat(scheduled)
if scheduled_dt <= now:
due.append(entry)
return due
# ─── Slots management ────────────────────────────────────────────────────────
def get_slots() -> list[dict]:
return _load_slots()
def save_slots(slots: list[dict]) -> list[dict]:
"""Save custom posting slots."""
_save_slots(slots)
assign_queue_slots() # Recalculate queue
return slots
def reset_slots_to_default() -> list[dict]:
_save_slots(DEFAULT_SLOTS)
assign_queue_slots()
return DEFAULT_SLOTS
# ─── Calendar view ────────────────────────────────────────────────────────────
def get_calendar_week(week_offset: int = 0) -> dict:
"""
Lấy data cho calendar view theo tuần.
week_offset=0 → tuần này, 1 → tuần sau, -1 → tuần trước.
"""
today = datetime.now(timezone.utc)
monday = today - timedelta(days=today.weekday()) + timedelta(weeks=week_offset)
days = []
queue = _load_queue()
content_by_date: dict[str, list] = {}
for entry in queue:
scheduled = entry.get("scheduled_at")
if scheduled:
dt = datetime.fromisoformat(scheduled)
date_key = dt.strftime("%Y-%m-%d")
content_by_date.setdefault(date_key, []).append(entry)
day_names = ["Thứ 2", "Thứ 3", "Thứ 4", "Thứ 5", "Thứ 6", "Thứ 7", "CN"]
for i in range(7):
day = monday + timedelta(days=i)
date_key = day.strftime("%Y-%m-%d")
days.append({
"date": date_key,
"day_name": day_names[i],
"is_today": date_key == today.strftime("%Y-%m-%d"),
"entries": content_by_date.get(date_key, []),
})
return {
"week_start": monday.strftime("%Y-%m-%d"),
"week_offset": week_offset,
"days": days,
}
"""Social providers package."""
from .base import SocialProvider, APIError, RateLimitError
from .facebook import FacebookProvider, InstagramProvider
from .tiktok import TikTokProvider
__all__ = [
"SocialProvider",
"APIError",
"RateLimitError",
"FacebookProvider",
"InstagramProvider",
"TikTokProvider",
]
"""
Social Provider Base Class
Ported & simplified từ BrightBean Studio providers/base.py
Bỏ Django dependency, chuyển sang httpx async cho FastAPI
"""
from __future__ import annotations
import logging
from abc import ABC, abstractmethod
from typing import Any
import httpx
logger = logging.getLogger(__name__)
REQUEST_TIMEOUT = 30.0
class APIError(Exception):
def __init__(self, message: str, status_code: int = 0, platform: str = "", raw_response: dict | None = None):
super().__init__(message)
self.status_code = status_code
self.platform = platform
self.raw_response = raw_response or {}
class RateLimitError(APIError):
def __init__(self, message: str, retry_after: int | None = None, **kwargs):
super().__init__(message, **kwargs)
self.retry_after = retry_after
class SocialProvider(ABC):
"""
Abstract base class cho tất cả social platform providers.
Ported từ BrightBean Studio — simplified cho Canifa use case.
Mục đích: publish outfit suggestions + thu thập social feedback.
"""
def __init__(self, credentials: dict | None = None):
self.credentials = credentials or {}
@property
@abstractmethod
def platform_name(self) -> str:
"""Tên platform (vd: 'Facebook', 'Instagram', 'TikTok')."""
@abstractmethod
def get_profile(self, access_token: str) -> dict:
"""Lấy thông tin tài khoản đã kết nối."""
@abstractmethod
def publish_post(self, access_token: str, content: dict) -> dict:
"""
Publish nội dung lên platform.
content = {
"text": str,
"image_url": str | None,
"link": str | None,
}
Returns: {"post_id": str, "url": str, "success": bool}
"""
# ─────────────────────────────────────────
# Optional methods (override khi cần)
# ─────────────────────────────────────────
def get_comments(self, access_token: str, post_id: str) -> list[dict]:
"""Lấy comments của một post. Override theo từng platform."""
raise NotImplementedError(f"{self.platform_name} chưa hỗ trợ get_comments")
def reply_comment(self, access_token: str, comment_id: str, text: str) -> dict:
"""Reply comment. Override theo từng platform."""
raise NotImplementedError(f"{self.platform_name} chưa hỗ trợ reply_comment")
def validate_token(self, access_token: str) -> bool:
"""Kiểm tra token còn hiệu lực không."""
try:
self.get_profile(access_token)
return True
except Exception:
return False
# ─────────────────────────────────────────
# HTTP helper (sync — dùng cho background worker)
# ─────────────────────────────────────────
def _request(
self,
method: str,
url: str,
*,
access_token: str | None = None,
headers: dict | None = None,
params: dict | None = None,
json: dict | None = None,
data: Any = None,
files: dict | None = None,
timeout: float = REQUEST_TIMEOUT,
) -> httpx.Response:
"""HTTP request với error handling chuẩn. Raises APIError / RateLimitError."""
req_headers = {}
if access_token:
req_headers["Authorization"] = f"Bearer {access_token}"
if headers:
req_headers.update(headers)
request_kwargs: dict = {
"headers": req_headers,
"params": params,
"json": json,
"files": files,
}
if isinstance(data, bytes):
request_kwargs["content"] = data
else:
request_kwargs["data"] = data
with httpx.Client(timeout=timeout) as client:
response = client.request(method, url, **request_kwargs)
if response.status_code == 429:
retry_after = response.headers.get("Retry-After")
raise RateLimitError(
f"Rate limit {self.platform_name}: {response.text[:300]}",
retry_after=int(retry_after) if retry_after else None,
platform=self.platform_name,
raw_response=self._safe_json(response),
)
if response.status_code >= 400:
raise APIError(
f"{self.platform_name} API {response.status_code}: {response.text[:300]}",
status_code=response.status_code,
platform=self.platform_name,
raw_response=self._safe_json(response),
)
return response
@staticmethod
def _safe_json(response: httpx.Response) -> dict:
try:
return response.json()
except Exception:
return {}
"""
Facebook & Instagram Provider
Ported & adapted từ BrightBean Studio providers/facebook.py + providers/instagram.py
Simplified: chỉ giữ phần publish post + get comments (Canifa use case)
"""
from __future__ import annotations
import logging
import os
from .base import SocialProvider, APIError
logger = logging.getLogger(__name__)
FACEBOOK_GRAPH_API = "https://graph.facebook.com/v19.0"
# App credentials từ env
FACEBOOK_APP_ID = os.getenv("FACEBOOK_APP_ID", "")
FACEBOOK_APP_SECRET = os.getenv("FACEBOOK_APP_SECRET", "")
class FacebookProvider(SocialProvider):
"""
Facebook Page provider — publish posts + đọc comments.
Dùng cho: Canifa auto-post outfit suggestions lên Facebook Page.
"""
@property
def platform_name(self) -> str:
return "Facebook"
def get_profile(self, access_token: str) -> dict:
"""Lấy thông tin Facebook Page."""
response = self._request(
"GET",
f"{FACEBOOK_GRAPH_API}/me",
access_token=access_token,
params={"fields": "id,name,fan_count,picture"},
)
return response.json()
def publish_post(self, access_token: str, content: dict) -> dict:
"""
Publish một post lên Facebook Page.
content = {"text": str, "link": str | None, "image_url": str | None}
"""
page_id = content.get("page_id", "me")
payload: dict = {"message": content.get("text", "")}
if content.get("link"):
payload["link"] = content["link"]
# Nếu có ảnh → dùng photos endpoint
if content.get("image_url"):
return self._publish_with_photo(access_token, page_id, content)
response = self._request(
"POST",
f"{FACEBOOK_GRAPH_API}/{page_id}/feed",
access_token=access_token,
json=payload,
)
data = response.json()
return {
"post_id": data.get("id", ""),
"success": True,
"platform": "facebook",
}
def _publish_with_photo(self, access_token: str, page_id: str, content: dict) -> dict:
"""Upload ảnh + publish caption lên Facebook."""
payload = {
"url": content["image_url"],
"caption": content.get("text", ""),
"published": True,
}
response = self._request(
"POST",
f"{FACEBOOK_GRAPH_API}/{page_id}/photos",
access_token=access_token,
json=payload,
)
data = response.json()
return {
"post_id": data.get("id", ""),
"success": True,
"platform": "facebook",
}
def get_comments(self, access_token: str, post_id: str) -> list[dict]:
"""Lấy comments của một Facebook post để phân tích feedback."""
response = self._request(
"GET",
f"{FACEBOOK_GRAPH_API}/{post_id}/comments",
access_token=access_token,
params={
"fields": "id,message,from,created_time,like_count",
"order": "chronological",
"limit": 100,
},
)
data = response.json()
comments = []
for item in data.get("data", []):
comments.append({
"platform_comment_id": item.get("id", ""),
"body": item.get("message", ""),
"sender_name": item.get("from", {}).get("name", "Unknown"),
"sender_id": item.get("from", {}).get("id", ""),
"created_at": item.get("created_time", ""),
"like_count": item.get("like_count", 0),
})
return comments
def reply_comment(self, access_token: str, comment_id: str, text: str) -> dict:
"""Reply vào một Facebook comment."""
response = self._request(
"POST",
f"{FACEBOOK_GRAPH_API}/{comment_id}/comments",
access_token=access_token,
json={"message": text},
)
data = response.json()
return {"reply_id": data.get("id", ""), "success": True}
class InstagramProvider(SocialProvider):
"""
Instagram Business Account provider.
Dùng cho: Canifa auto-post outfit / fashion content lên Instagram.
"""
@property
def platform_name(self) -> str:
return "Instagram"
def get_profile(self, access_token: str) -> dict:
"""Lấy thông tin Instagram Business Account."""
response = self._request(
"GET",
f"{FACEBOOK_GRAPH_API}/me",
access_token=access_token,
params={
"fields": "id,name,username,biography,followers_count,media_count,profile_picture_url",
},
)
return response.json()
def publish_post(self, access_token: str, content: dict) -> dict:
"""
Publish Instagram post (ảnh + caption).
Instagram requires media — must have image_url.
content = {"text": str, "image_url": str (required), "ig_user_id": str}
"""
ig_user_id = content.get("ig_user_id", "me")
image_url = content.get("image_url")
if not image_url:
raise ValueError("Instagram publish requires image_url")
# Step 1: Tạo media container
container_response = self._request(
"POST",
f"{FACEBOOK_GRAPH_API}/{ig_user_id}/media",
access_token=access_token,
json={
"image_url": image_url,
"caption": content.get("text", ""),
},
)
container_id = container_response.json().get("id")
if not container_id:
raise APIError("Không tạo được Instagram media container", platform="Instagram")
# Step 2: Publish container
publish_response = self._request(
"POST",
f"{FACEBOOK_GRAPH_API}/{ig_user_id}/media_publish",
access_token=access_token,
json={"creation_id": container_id},
)
data = publish_response.json()
return {
"post_id": data.get("id", ""),
"success": True,
"platform": "instagram",
}
def get_comments(self, access_token: str, post_id: str) -> list[dict]:
"""Lấy comments của một Instagram post."""
response = self._request(
"GET",
f"{FACEBOOK_GRAPH_API}/{post_id}/comments",
access_token=access_token,
params={
"fields": "id,text,username,timestamp,like_count",
"limit": 100,
},
)
data = response.json()
comments = []
for item in data.get("data", []):
comments.append({
"platform_comment_id": item.get("id", ""),
"body": item.get("text", ""),
"sender_name": item.get("username", "Unknown"),
"sender_id": item.get("username", ""),
"created_at": item.get("timestamp", ""),
"like_count": item.get("like_count", 0),
})
return comments
"""
TikTok Provider
Ported & adapted từ BrightBean Studio providers/tiktok.py
Simplified cho Canifa use case: publish fashion video + đọc comments
"""
from __future__ import annotations
import logging
import os
from .base import SocialProvider, APIError
logger = logging.getLogger(__name__)
TIKTOK_API_BASE = "https://open.tiktokapis.com/v2"
TIKTOK_CLIENT_KEY = os.getenv("TIKTOK_CLIENT_KEY", "")
TIKTOK_CLIENT_SECRET = os.getenv("TIKTOK_CLIENT_SECRET", "")
class TikTokProvider(SocialProvider):
"""
TikTok provider — publish video + đọc comments.
Dùng cho: Canifa auto-post fashion reels/videos lên TikTok.
"""
@property
def platform_name(self) -> str:
return "TikTok"
def get_profile(self, access_token: str) -> dict:
"""Lấy thông tin TikTok user."""
response = self._request(
"GET",
f"{TIKTOK_API_BASE}/user/info/",
access_token=access_token,
params={
"fields": "open_id,union_id,display_name,avatar_url,follower_count,following_count,likes_count",
},
)
return response.json().get("data", {}).get("user", {})
def publish_post(self, access_token: str, content: dict) -> dict:
"""
Publish TikTok video.
TikTok cần video_url và dùng Content Posting API.
content = {
"title": str, ← caption/title
"video_url": str, ← direct link tới video file
"privacy": str, ← "PUBLIC" | "FRIENDS" | "PRIVATE"
}
"""
# TikTok direct post via URL
payload = {
"post_info": {
"title": content.get("title", ""),
"privacy_level": content.get("privacy", "PUBLIC_TO_EVERYONE"),
"disable_duet": False,
"disable_comment": False,
"disable_stitch": False,
},
"source_info": {
"source": "PULL_FROM_URL",
"video_url": content.get("video_url", ""),
},
}
response = self._request(
"POST",
f"{TIKTOK_API_BASE}/post/publish/video/init/",
access_token=access_token,
json=payload,
)
data = response.json().get("data", {})
return {
"post_id": data.get("publish_id", ""),
"success": True,
"platform": "tiktok",
}
def get_comments(self, access_token: str, post_id: str) -> list[dict]:
"""Lấy comments của một TikTok video."""
response = self._request(
"POST",
f"{TIKTOK_API_BASE}/video/comment/list/",
access_token=access_token,
json={
"video_id": post_id,
"max_count": 100,
},
)
data = response.json().get("data", {})
comments = []
for item in data.get("comments", []):
comments.append({
"platform_comment_id": item.get("id", ""),
"body": item.get("text", ""),
"sender_name": item.get("display_name", "Unknown"),
"sender_id": item.get("open_id", ""),
"created_at": str(item.get("create_time", "")),
"like_count": item.get("like_count", 0),
})
return comments
This diff is collapsed.
"""
Sentiment Analysis Engine — Tiếng Việt + Tiếng Anh
Ported & extended từ BrightBean Studio apps/inbox/sentiment.py
Thêm: keywords tiếng Việt fashion/thương mại điện tử cho Canifa use case
"""
import re
# ─────────────────────────────────────────────
# ENGLISH Keywords (từ BrightBean Studio)
# ─────────────────────────────────────────────
_POSITIVE_EN = {
"love", "great", "amazing", "thank", "thanks", "excellent", "awesome",
"perfect", "best", "fantastic", "wonderful", "beautiful", "brilliant",
"incredible", "outstanding", "impressive", "helpful", "appreciate",
"happy", "glad", "excited", "recommend", "superb", "delighted", "thrilled",
"good", "nice", "cute", "pretty", "quality", "satisfied", "worth",
}
_NEGATIVE_EN = {
"hate", "terrible", "awful", "worst", "disappointed", "broken", "scam",
"refund", "horrible", "disgusting", "pathetic", "useless", "angry",
"furious", "unacceptable", "frustrating", "annoying", "poor", "waste",
"trash", "spam", "fake", "misleading", "rude", "unprofessional",
"cheap", "bad", "ugly", "tight", "loose", "wrong", "return", "complaint",
}
# ─────────────────────────────────────────────
# VIETNAMESE Keywords — Thêm mới cho Canifa
# ─────────────────────────────────────────────
_POSITIVE_VI = {
# Khen ngợi chung
"tuyệt", "tuyệt vời", "xuất sắc", "hoàn hảo", "đẹp", "đẹp lắm",
"đẹp quá", "xinh", "xinh lắm", "xinh xắn", "dễ thương", "sang",
"sang trọng", "chất", "chất lắm", "hài lòng", "ưng", "ưng lắm",
"thích", "thích lắm", "thích quá", "yêu", "mê", "ok", "ổn",
# Chất lượng sản phẩm
"chất lượng", "mịn", "mềm", "nhẹ", "thoáng", "thoải mái", "vừa vặn",
"đúng size", "chuẩn", "bền", "đúng màu", "đúng mô tả", "y hình",
# Giao hàng / Dịch vụ
"nhanh", "giao nhanh", "đóng gói đẹp", "cẩn thận", "nhiệt tình",
"tư vấn tốt", "phục vụ tốt", "hỗ trợ nhiệt tình",
# Sẽ mua lại
"sẽ mua lại", "ủng hộ tiếp", "mua lần 2", "recommend", "giới thiệu",
"đáng mua", "đáng tiền", "xứng đáng", "worth it",
# Canifa specific
"canifa ngon", "canifa tốt", "thương hiệu uy tín",
"5 sao", "5⭐", "⭐⭐⭐⭐⭐", "👍", "❤️", "🔥",
}
_NEGATIVE_VI = {
# Chê về sản phẩm
"xấu", "kém", "tệ", "không đẹp", "không ưng", "thất vọng",
"chán", "nhạt", "nhàm", "mình không thích", "không như mô tả",
"khác hình", "màu lệch", "phai màu", "nhăn nhúm", "lỗi", "lỗi hàng",
"đường may xấu", "lỏng chỉ", "bung chỉ",
# Size / Fit
"không vừa", "sai size", "nhỏ hơn", "lớn hơn", "chật", "rộng thùng thình",
"ngắn quá", "dài quá", "bé quá",
# Giao hàng / Dịch vụ
"giao chậm", "giao trễ", "sai hàng", "thiếu hàng", "đóng gói ẩu",
"bị bẩn", "phục vụ tệ", "nhân viên thái độ", "không hỗ trợ",
"không phản hồi", "lờ đi", "im hơi lặng tiếng",
# Trả hàng / Hoàn tiền
"trả hàng", "đổi hàng", "hoàn tiền", "khiếu nại", "phàn nàn",
"tố cáo", "scam", "lừa đảo", "hàng giả",
# Canifa specific
"canifa tệ", "không mua nữa", "lần cuối", "thất vọng với canifa",
"1 sao", "1⭐", "👎", "😡", "😤", "😠",
}
# ─────────────────────────────────────────────
# Core analyzer
# ─────────────────────────────────────────────
def _tokenize(text: str) -> set[str]:
"""Lowercase và tokenize text, giữ lại cả nguyên câu (cho multi-word keywords)."""
text_lower = text.lower().strip()
# Word tokens (strip punctuation)
words = set(re.sub(r"[^\w\s]", " ", text_lower).split())
# Thêm nguyên text để match multi-word phrases
words.add(text_lower)
return words
def analyze_sentiment(text: str) -> str:
"""
Phân tích sentiment của text (Tiếng Việt + Tiếng Anh).
Returns:
'positive' | 'negative' | 'neutral'
"""
if not text or not text.strip():
return "neutral"
tokens = _tokenize(text)
text_lower = text.lower()
# Score tiếng Anh
pos_en = sum(1 for kw in _POSITIVE_EN if kw in tokens)
neg_en = sum(1 for kw in _NEGATIVE_EN if kw in tokens)
# Score tiếng Việt (check phrase match trong full text)
pos_vi = sum(1 for kw in _POSITIVE_VI if kw in text_lower)
neg_vi = sum(1 for kw in _NEGATIVE_VI if kw in text_lower)
pos_total = pos_en + pos_vi
neg_total = neg_en + neg_vi
if neg_total > pos_total:
return "negative"
elif pos_total > neg_total:
return "positive"
return "neutral"
def analyze_sentiment_detail(text: str) -> dict:
"""
Phân tích chi tiết — trả về score + label + matched keywords.
Dùng cho debug / dashboard.
"""
if not text or not text.strip():
return {"label": "neutral", "pos_score": 0, "neg_score": 0, "matched": []}
tokens = _tokenize(text)
text_lower = text.lower()
matched_pos = [kw for kw in (_POSITIVE_EN | _POSITIVE_VI) if kw in tokens or kw in text_lower]
matched_neg = [kw for kw in (_NEGATIVE_EN | _NEGATIVE_VI) if kw in tokens or kw in text_lower]
pos_score = len(matched_pos)
neg_score = len(matched_neg)
if neg_score > pos_score:
label = "negative"
elif pos_score > neg_score:
label = "positive"
else:
label = "neutral"
return {
"label": label,
"pos_score": pos_score,
"neg_score": neg_score,
"matched_positive": matched_pos[:5], # Top 5 để không spam
"matched_negative": matched_neg[:5],
}
...@@ -214,6 +214,29 @@ app.include_router(mock_auth_router) # Mock Auth (identity linking test) ...@@ -214,6 +214,29 @@ app.include_router(mock_auth_router) # Mock Auth (identity linking test)
from api.feedback_agent_route import router as feedback_agent_router from api.feedback_agent_route import router as feedback_agent_router
app.include_router(feedback_agent_router) # Lõi Agent Rút Kinh Nghiệm (Langfuse -> Rules) app.include_router(feedback_agent_router) # Lõi Agent Rút Kinh Nghiệm (Langfuse -> Rules)
from api.social_inbox_route import router as social_inbox_router
app.include_router(social_inbox_router) # Social Inbox (Facebook/Instagram/TikTok → Learning Loop)
# ─── Phase 2: AI Content Pipeline ───────────────────────────────────────────
from api.notification_route import router as notification_router
app.include_router(notification_router) # In-app + Email + Webhook + Slack notifications
from api.content_approval_route import router as content_approval_router
app.include_router(content_approval_router) # Content approval gate (draft → review → publish)
from api.queue_route import router as queue_router
app.include_router(queue_router) # Post queue + Calendar scheduling
from api.media_route import router as media_router
app.include_router(media_router) # Media library (upload/resize/serve)
from api.templates_route import router as templates_router
app.include_router(templates_router) # Content templates + RSS feeds (ported from BrightBean)
# ─── Start publish engine background loop ───────────────────────────────────
from common.social.scheduler import start_publish_engine
start_publish_engine(app) # Auto-publish scheduled content every 30s
if __name__ == "__main__": if __name__ == "__main__":
print("=" * 60) print("=" * 60)
...@@ -222,6 +245,11 @@ if __name__ == "__main__": ...@@ -222,6 +245,11 @@ if __name__ == "__main__":
print(f"📡 REST API: http://localhost:{PORT}") print(f"📡 REST API: http://localhost:{PORT}")
print(f"📡 Test Chatbot: http://localhost:{PORT}/static/index.html") print(f"📡 Test Chatbot: http://localhost:{PORT}/static/index.html")
print(f"📚 API Docs: http://localhost:{PORT}/docs") print(f"📚 API Docs: http://localhost:{PORT}/docs")
print(f"📋 Approval: http://localhost:{PORT}/static/content-approval/index.html")
print(f"📅 Calendar: http://localhost:{PORT}/static/content-calendar/index.html")
print(f"🖼️ Media Library: http://localhost:{PORT}/static/media-library/index.html")
print(f"📬 Social Inbox: http://localhost:{PORT}/static/social-inbox/index.html")
print(f"✍️ Composer: http://localhost:{PORT}/static/content-composer/index.html")
print("=" * 60) print("=" * 60)
ENABLE_RELOAD = False ENABLE_RELOAD = False
......
This diff is collapsed.
This diff is collapsed.
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