340 lines
15 KiB
Python
340 lines
15 KiB
Python
from sqlalchemy import create_engine, text
|
|
from sqlalchemy.orm import sessionmaker
|
|
|
|
from app.config import settings
|
|
|
|
|
|
engine = create_engine(settings.database_url, pool_pre_ping=True)
|
|
|
|
SessionLocal = sessionmaker(
|
|
autocommit=False,
|
|
autoflush=False,
|
|
bind=engine,
|
|
)
|
|
|
|
|
|
_CORE_SCHEMA_READY = False
|
|
|
|
|
|
def ensure_core_schema() -> None:
|
|
"""Cria o schema base do ClientFlow numa base PostgreSQL vazia.
|
|
|
|
As versões anteriores assumiam que `tasks`, `messages`, `raw_events` e
|
|
`action_runs` já existiam. Numa instalação limpa isso fazia o arranque
|
|
falhar quando a camada de oportunidades tentava executar:
|
|
|
|
ALTER TABLE tasks ADD COLUMN ...
|
|
|
|
Esta função é aditiva e segura para bases existentes: cria tabelas e
|
|
índices apenas se não existirem e adiciona colunas opcionais em falta.
|
|
"""
|
|
global _CORE_SCHEMA_READY
|
|
if _CORE_SCHEMA_READY:
|
|
return
|
|
|
|
with engine.begin() as conn:
|
|
# Necessário para DEFAULT gen_random_uuid(). Em PostgreSQL moderno,
|
|
# pgcrypto é a forma mais simples de gerar UUIDs no próprio servidor.
|
|
conn.execute(text("CREATE EXTENSION IF NOT EXISTS pgcrypto"))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS schema_migrations (
|
|
version TEXT PRIMARY KEY,
|
|
name TEXT NOT NULL,
|
|
applied_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
)
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS raw_events (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
source_system TEXT NOT NULL DEFAULT 'clientflow',
|
|
event_type TEXT,
|
|
source_event_id TEXT,
|
|
conversation_id TEXT,
|
|
contact_id TEXT,
|
|
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
processed BOOLEAN NOT NULL DEFAULT FALSE,
|
|
ignored BOOLEAN NOT NULL DEFAULT FALSE,
|
|
processing_error TEXT,
|
|
message_id UUID,
|
|
action_run_id UUID,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
processed_at TIMESTAMPTZ
|
|
)
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS ux_raw_events_source_event
|
|
ON raw_events(source_system, source_event_id)
|
|
WHERE source_event_id IS NOT NULL
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS messages (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
raw_event_id UUID,
|
|
source_system TEXT NOT NULL DEFAULT 'clientflow',
|
|
source_event_id TEXT,
|
|
conversation_id TEXT,
|
|
contact_id TEXT,
|
|
direction TEXT NOT NULL DEFAULT 'inbound',
|
|
raw_body TEXT,
|
|
clean_body TEXT,
|
|
previous_context TEXT,
|
|
metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
)
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS action_runs (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
message_id UUID,
|
|
raw_event_id UUID,
|
|
conversation_id TEXT,
|
|
contact_id TEXT,
|
|
source_system TEXT NOT NULL DEFAULT 'clientflow',
|
|
model TEXT,
|
|
provider TEXT,
|
|
openrouter_generation_id TEXT,
|
|
decision_source TEXT,
|
|
action_decision JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
action_result JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
prompt_tokens INTEGER,
|
|
completion_tokens INTEGER,
|
|
total_tokens INTEGER,
|
|
cost NUMERIC(12,6),
|
|
usage JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
needs_review BOOLEAN NOT NULL DEFAULT FALSE,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
)
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS tasks (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
action_run_id UUID,
|
|
message_id UUID,
|
|
raw_event_id UUID,
|
|
opportunity_id UUID,
|
|
conversation_id TEXT,
|
|
contact_id TEXT,
|
|
customer_id TEXT,
|
|
action_code TEXT NOT NULL DEFAULT 'REVIEW_MANUALLY',
|
|
route TEXT NOT NULL DEFAULT 'rever',
|
|
action TEXT NOT NULL DEFAULT 'Rever manualmente',
|
|
note TEXT,
|
|
action_required BOOLEAN NOT NULL DEFAULT FALSE,
|
|
safe_to_post BOOLEAN NOT NULL DEFAULT FALSE,
|
|
status TEXT NOT NULL DEFAULT 'pending',
|
|
source_system TEXT NOT NULL DEFAULT 'clientflow',
|
|
source_event_id TEXT,
|
|
idempotency_key TEXT,
|
|
due_at TIMESTAMPTZ,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
done_at TIMESTAMPTZ,
|
|
done_by TEXT,
|
|
metadata JSONB NOT NULL DEFAULT '{}'::jsonb
|
|
)
|
|
"""))
|
|
|
|
# Colunas adicionadas em upgrades anteriores. Mantidas aqui para
|
|
# compatibilidade quando a tabela já existe numa base antiga.
|
|
for statement in [
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS opportunity_id UUID",
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS due_at TIMESTAMPTZ",
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS done_at TIMESTAMPTZ",
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS done_by TEXT",
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS metadata JSONB NOT NULL DEFAULT '{}'::jsonb",
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS idempotency_key TEXT",
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS priority TEXT NOT NULL DEFAULT 'normal'",
|
|
"ALTER TABLE tasks ADD COLUMN IF NOT EXISTS assigned_to TEXT",
|
|
]:
|
|
conn.execute(text(statement))
|
|
|
|
conn.execute(text("""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS ux_tasks_idempotency_key
|
|
ON tasks(idempotency_key)
|
|
WHERE idempotency_key IS NOT NULL
|
|
"""))
|
|
conn.execute(text("CREATE UNIQUE INDEX IF NOT EXISTS idx_tasks_idempotency_key ON tasks(idempotency_key)"))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS task_preparations (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
task_id UUID NOT NULL,
|
|
conversation_id TEXT,
|
|
contact_id TEXT,
|
|
prep_type TEXT NOT NULL DEFAULT 'generic',
|
|
status TEXT NOT NULL DEFAULT 'draft',
|
|
extracted_data JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
missing_fields JSONB NOT NULL DEFAULT '[]'::jsonb,
|
|
suggested_reply TEXT,
|
|
confidence NUMERIC(4,3),
|
|
model TEXT,
|
|
provider TEXT,
|
|
total_tokens INTEGER NOT NULL DEFAULT 0,
|
|
cost NUMERIC(12,6) NOT NULL DEFAULT 0,
|
|
raw_response JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
)
|
|
"""))
|
|
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_task_preparations_task ON task_preparations(task_id, created_at DESC)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_task_preparations_conversation ON task_preparations(conversation_id)"))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS task_events (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
task_id UUID NOT NULL,
|
|
event_type TEXT NOT NULL,
|
|
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
created_by TEXT NOT NULL DEFAULT 'system',
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
)
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS business_events (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
event_type TEXT NOT NULL,
|
|
task_id UUID,
|
|
action_run_id UUID,
|
|
message_id UUID,
|
|
raw_event_id UUID,
|
|
customer_id TEXT,
|
|
conversation_id TEXT,
|
|
contact_id TEXT,
|
|
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
idempotency_key TEXT,
|
|
created_by TEXT NOT NULL DEFAULT 'system',
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
)
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS ux_business_events_idempotency_key
|
|
ON business_events(idempotency_key)
|
|
WHERE idempotency_key IS NOT NULL
|
|
"""))
|
|
conn.execute(text("CREATE UNIQUE INDEX IF NOT EXISTS idx_business_events_idempotency_key ON business_events(idempotency_key)"))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS integration_outbox (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
business_event_id UUID,
|
|
target_system TEXT NOT NULL,
|
|
action_type TEXT NOT NULL,
|
|
payload JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
status TEXT NOT NULL DEFAULT 'pending',
|
|
retry_count INTEGER NOT NULL DEFAULT 0,
|
|
idempotency_key TEXT,
|
|
last_error TEXT,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
sent_at TIMESTAMPTZ
|
|
)
|
|
"""))
|
|
|
|
# Older ClientFlow installs may already have integration_outbox created
|
|
# without newer operational columns. The outbox worker selects/updates
|
|
# these columns directly, so guarantee them on every schema init.
|
|
for statement in [
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS business_event_id UUID",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS retry_count INTEGER NOT NULL DEFAULT 0",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS idempotency_key TEXT",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS last_error TEXT",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS updated_at TIMESTAMPTZ NOT NULL DEFAULT now()",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS sent_at TIMESTAMPTZ",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS locked_at TIMESTAMPTZ",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS lock_owner TEXT",
|
|
"ALTER TABLE integration_outbox ADD COLUMN IF NOT EXISTS ignored_at TIMESTAMPTZ",
|
|
]:
|
|
conn.execute(text(statement))
|
|
|
|
conn.execute(text("""
|
|
CREATE UNIQUE INDEX IF NOT EXISTS ux_integration_outbox_idempotency_key
|
|
ON integration_outbox(idempotency_key)
|
|
WHERE idempotency_key IS NOT NULL
|
|
"""))
|
|
conn.execute(text("CREATE UNIQUE INDEX IF NOT EXISTS idx_integration_outbox_idempotency_key ON integration_outbox(idempotency_key)"))
|
|
conn.execute(text("""
|
|
CREATE INDEX IF NOT EXISTS idx_integration_outbox_pending_dispatch
|
|
ON integration_outbox(target_system, created_at)
|
|
WHERE status = 'pending'
|
|
"""))
|
|
|
|
conn.execute(text("""
|
|
CREATE TABLE IF NOT EXISTS external_mappings (
|
|
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
|
|
local_system TEXT NOT NULL DEFAULT 'clientflow',
|
|
local_entity_type TEXT NOT NULL,
|
|
local_entity_id TEXT NOT NULL,
|
|
external_system TEXT NOT NULL,
|
|
external_entity_type TEXT NOT NULL,
|
|
external_entity_id TEXT,
|
|
external_url TEXT,
|
|
match_key TEXT,
|
|
match_value TEXT,
|
|
confidence NUMERIC(4,3),
|
|
metadata JSONB NOT NULL DEFAULT '{}'::jsonb,
|
|
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
|
|
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
)
|
|
"""))
|
|
|
|
# Índices de consulta mais usados no admin e webhooks.
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_status ON tasks(status)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_route ON tasks(route)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_conversation ON tasks(conversation_id)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_contact ON tasks(contact_id)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_tasks_created_at ON tasks(created_at DESC)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_raw_events_created_at ON raw_events(created_at DESC)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_raw_events_conversation ON raw_events(conversation_id)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_messages_conversation ON messages(conversation_id)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_action_runs_created_at ON action_runs(created_at DESC)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_task_events_task ON task_events(task_id, created_at DESC)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_business_events_created_at ON business_events(created_at DESC)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_outbox_status ON integration_outbox(status)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_outbox_target ON integration_outbox(target_system)"))
|
|
conn.execute(text("CREATE INDEX IF NOT EXISTS idx_external_mappings_local ON external_mappings(local_system, local_entity_type, local_entity_id)"))
|
|
|
|
_CORE_SCHEMA_READY = True
|
|
|
|
|
|
def init_db() -> None:
|
|
"""Inicialização do schema ClientFlow.
|
|
|
|
Ordem importante para bases novas:
|
|
1. core: raw_events/messages/action_runs/tasks/outbox;
|
|
2. oportunidades: depende de tasks;
|
|
3. produtos: depende de opportunities.
|
|
"""
|
|
try:
|
|
ensure_core_schema()
|
|
|
|
from app.opportunity_service import ensure_opportunity_schema
|
|
from app.product_service import ensure_product_schema
|
|
from app.operation_service import ensure_operation_schema
|
|
from app.commercial_service import ensure_commercial_schema
|
|
from app.communication_service import ensure_communication_schema
|
|
from app.reply_assistant_service import ensure_reply_assistant_schema
|
|
from app.reconciliation_service import ensure_reconciliation_schema
|
|
from app.revenue_forecast_service import ensure_revenue_forecast_schema
|
|
|
|
ensure_opportunity_schema()
|
|
ensure_product_schema()
|
|
ensure_operation_schema()
|
|
ensure_commercial_schema()
|
|
ensure_communication_schema()
|
|
ensure_reply_assistant_schema()
|
|
ensure_reconciliation_schema()
|
|
ensure_revenue_forecast_schema()
|
|
except Exception as exc:
|
|
print(f"ClientFlow schema init failed: {exc}", flush=True)
|
|
raise
|
|
return None
|