A multi-agent AI orchestration platform where a central LangGraph orchestrator reasons, remembers, and delegates work to dynamically registered specialized agents. Features dual-LLM routing (Groq Cloud + local Ollama), SSE token streaming, and reliable task delivery with agent callbacks + a background Task Watcher.
Krastix is a domain-agnostic agentic engine for building AI-powered workflows. The orchestrator:
- Reasons — Plans responses using LLMs (Ollama local + Groq cloud) with RAG context injection
- Remembers — Stores and retrieves semantic memories scoped per domain via pgvector
- Delegates — Routes sub-tasks to specialized agents via a registry-driven dispatcher
- Streams — Delivers LLM tokens in real-time to the frontend via Server-Sent Events
- Validates — Enforces JSON Schema on all entity writes with optimistic concurrency
- Watches — Background Task Watcher catches stale/lost tasks as a safety net
Current agents: CRM (universal entity management), Form Builder (Tally.so), Research (Firecrawl + LinkedIn scraping), Doc Agent (PDF/image extraction via Groq Vision), Communication Agent (approved Gmail sends via Celery).
- Docker & Docker Compose
- Ollama running on a reachable host (default: Tailscale LAN node)
- Supabase project (PostgreSQL + pgvector)
- Groq API key (for vision + fast text models)
- Google OAuth credentials + integration encryption key (for Gmail connect/send)
- GROK-compatible API key (for email summarization in Communications view)
- API keys: Firecrawl, ScrapeCreators (for research agent)
cp .env.example .envEdit .env with your values:
# Database (Supabase PostgreSQL)
DATABASE_URL=postgresql://postgres.xxx:password@host:5432/postgres
# LLM — Local (Ollama via Tailscale)
OLLAMA_BASE_URL=http://100.115.107.20:11434
OLLAMA_MODEL=qwen2.5:14b-instruct-q5_K_M
# LLM — Cloud (Groq)
GROQ_API_KEY=gsk_... # Vision + fast text inference
GROQ_VISION_MODEL=llama-3.2-90b-vision-preview
GROQ_TEXT_MODEL=llama-3.3-70b-versatile
# Communications summarization key
# Supports Grok/xAI key or Groq-compatible key (provider auto-detected by orchestrator)
GROK_API_KEY=gsk_or_xai_...
# Message Queue
REDIS_URL=redis://redis:6379/0
# External APIs
FIRECRAWL_API_KEY=fc_...
SCRAPECREATORS_API_KEY=...
# Google OAuth (Communication Agent)
GOOGLE_CLIENT_ID=...
GOOGLE_CLIENT_SECRET=...
GOOGLE_REDIRECT_URI=http://localhost:8000/api/v1/integrations/google/oauth/callback
# Integration secret encryption (Fernet key, generate once and keep private)
INTEGRATIONS_ENCRYPTION_KEY=...
# Callback signing secret for agent -> orchestrator task callbacks
CALLBACK_SIGNING_SECRET=...
# Optional keepalive token for external cron pings
KEEPALIVE_TOKEN=...
# Supabase Auth (optional)
SUPABASE_URL=https://xxx.supabase.co
SUPABASE_ANON_KEY=...
SUPABASE_SERVICE_ROLE_KEY=...Run the full schema against your Supabase PostgreSQL:
psql "$DATABASE_URL" -f init.sqlOr for existing databases, apply the incremental migration:
psql "$DATABASE_URL" -f migrations/002_universal_engine.sql
psql "$DATABASE_URL" -f migrations/003_doc_agent.sql
psql "$DATABASE_URL" -f migrations/004_communication_agent.sql
psql "$DATABASE_URL" -f migrations/005_phase0_stabilization.sql
psql "$DATABASE_URL" -f migrations/006_phase1_planner_scheduler.sqldocker compose up --build -dThis starts 7 containers:
| Service | Container | Port | Role |
|---|---|---|---|
| Redis | krastix-redis | 6379 | Celery broker |
| Orchestrator | krastix-orchestrator | 8000 | LangGraph brain + SSE streaming + Task Watcher |
| Research Agent | krastix-research-agent | 8001 | Web research (FastAPI) |
| Doc Agent | krastix-doc-agent | 8002 | Document extraction (Groq Vision + LangGraph) |
| CRM Agent | krastix-crm-agent | — | Entity management (Celery) |
| Form Agent | krastix-form-agent | — | Tally.so forms (Celery) |
| Communication Agent | krastix-communication-agent | — | Gmail send after draft approval (Celery) |
# Health checks
curl http://localhost:8000/health
curl http://localhost:8001/health
# Send a chat message
curl -X POST http://localhost:8000/api/v1/chat \
-H "Content-Type: application/json" \
-d '{
"user_id": "YOUR_UUID",
"domain": "recruitment",
"message": "Hello, what can you do?",
"session_id": "test-session-001"
}'flowchart LR
U[User]
FE[React Frontend]
O[Orchestrator\nFastAPI + LangGraph]
MEM[(Memory\npgvector)]
REG[(Agent Registry)]
CRM[CRM Agent\nCelery]
FORM[Form Agent\nCelery]
RESEARCH[Research Agent\nHTTP]
DOC[Doc Agent\nHTTP]
COM[Communication Agent\nCelery]
CB[/callbacks/task-completed/]
GMAIL[Gmail API]
DB[(PostgreSQL\nSupabase)]
U --> FE
FE -->|SSE chat| O
O --> MEM
O --> REG
O --> CRM
O --> FORM
O --> RESEARCH
O --> DOC
O --> COM
CRM --> CB
FORM --> CB
RESEARCH --> CB
DOC --> CB
COM --> CB
CB --> O
O --> GMAIL
O --> DB
Key patterns:
- Dual-LLM Routing — Groq (cloud) for vision/speed-critical tasks; Ollama (local) for text reasoning at zero API cost.
- SSE Token Streaming —
/chat/streamdelivers LLM tokens to the frontend in real-time via Server-Sent Events. - Agent Registry — Agents register in
agent_registrywith capabilities + dispatch method. Orchestrator queries at planning time. - Callback + Task Watcher — Agents call back on completion; background watcher catches anything missed after 10 minutes.
- Schema-on-Demand — Entity types defined in
entity_definitionswith JSON Schema. New types via SQL. - Optimistic Concurrency —
versioncolumn on entities. Concurrent writes detected and retried automatically. - Namespace Isolation — Memory searches scoped by
domain_keyto prevent cross-domain RAG leakage. - Communication Workflow — Gmail inbox list is fetched via orchestrator API, then summarize and reply draft are generated on-demand in full-email view before human approval and send.
See ARCHITECTURE.md for the full system diagram, database schema, and design decisions.
| Method | Endpoint | Description |
|---|---|---|
POST |
/api/v1/chat |
Send a message (non-streaming response) |
POST |
/api/v1/chat/stream |
Send a message (SSE token streaming) |
POST |
/callbacks/task-completed |
Agent callback on task completion |
POST |
/memory/ingest |
Ingest text into semantic memory |
POST |
/api/v1/batch/process |
Process pending batch jobs |
GET |
/api/v1/integrations/{user_id} |
List integration statuses |
DELETE |
/api/v1/integrations/{user_id}/{provider} |
Disconnect an integration |
GET |
/api/v1/communications/gmail/primary |
Fetch today's inbox emails (after_ts baseline optional) |
POST |
/api/v1/communications/gmail/summarize |
Generate on-demand email summary |
POST |
/api/v1/communications/gmail/reply/draft |
Generate editable reply draft |
POST |
/api/v1/communications/gmail/reply/send |
Send approved Gmail reply |
GET |
/health |
Health check |
GET |
/health/keepalive |
Lightweight DB keepalive ping (optional token via X-Keepalive-Token) |
{
"user_id": "uuid",
"domain": "recruitment",
"message": "Find me senior React developers",
"session_id": "optional-thread-id"
}Response:
{
"response": "I'll research that for you...",
"task_id": "uuid-of-delegated-task (if any)",
"session_id": "thread-id"
}Same request body as /api/v1/chat. Returns text/event-stream:
data: {"event": "token", "data": "I'll"}
data: {"event": "token", "data": " research"}
data: {"event": "tool_start", "data": {"tool": "DelegateTask", "input": {...}}}
data: {"event": "tool_result", "data": {"tool": "DelegateTask", "output": "..."}}
data: {"event": "done", "data": {"response": "...", "task_id": "uuid"}}
The frontend connects via fetch() + ReadableStream and falls back to /api/v1/chat on failure.
{
"user_id": "uuid",
"domain": "recruitment",
"content": "Research results about React developers...",
"metadata": { "source": "research_agent", "task_type": "GENERAL_SEARCH" }
}GET /api/v1/communications/gmail/primary?user_id=<uuid>&limit=25&after_ts=<unix_seconds>
- Returns today's inbox messages (
in:inbox -in:spam -in:trash) with full body + snippet. - If tokens cannot be decrypted or are missing, response includes
statusvalues likereauth_requiredornot_connected.
POST /api/v1/communications/gmail/summarize
{
"user_id": "uuid",
"subject": "Partnership Opportunity",
"sender": "Founder <founder@company.com>",
"body": "Full email body text...",
"snippet": "Optional snippet"
}POST /api/v1/communications/gmail/reply/draft
{
"user_id": "uuid",
"subject": "Partnership Opportunity",
"sender": "Founder <founder@company.com>",
"body": "Full email body text...",
"thread_id": "gmail-thread-id",
"message_id_header": "<message-id@domain.com>"
}POST /api/v1/communications/gmail/reply/send
{
"user_id": "uuid",
"to": "founder@company.com",
"subject": "Re: Partnership Opportunity",
"body": "Thanks for reaching out...",
"thread_id": "gmail-thread-id",
"in_reply_to": "<message-id@domain.com>",
"references": "<message-id@domain.com>"
}| Method | Endpoint | Description |
|---|---|---|
POST |
/research/run |
Execute a research task |
GET |
/health |
Health check |
{
"user_id": "uuid",
"task_type": "GENERAL_SEARCH",
"query_or_url": "latest AI recruitment tools 2025",
"context_metadata": { "task_id": "uuid", "session_id": "thread-id" }
}Supported task_type values: GENERAL_SEARCH, QUICK_SCRAPE, SITE_MAP, LINKEDIN_PROFILE
This repo includes a scheduler at .github/workflows/db-keepalive.yml that pings your PostgreSQL database directly every 12 hours.
This does not require the orchestrator app to be running.
Go to GitHub -> Settings -> Secrets and variables -> Actions -> New repository secret:
DATABASE_URL: full PostgreSQL connection string used by your deployment
Important: store DATABASE_URL as a single line (no wrapping quotes, no trailing newline).
Run the workflow manually once from Actions tab (DB Keepalive) and ensure it passes.
- Register in database:
INSERT INTO agent_registry (agent_key, display_name, queue_or_url, dispatch_method, capabilities, supported_domains)
VALUES (
'my_agent_v1',
'My Custom Agent',
'my_queue',
'celery',
'{"actions": ["do_thing"], "description": "Does the thing", "celery_task_name": "agents.my_worker.execute_task"}',
'["recruitment", "sales"]'
);- Add queue to domain config:
UPDATE domain_configs
SET allowed_agent_queues = allowed_agent_queues || '"my_queue"'
WHERE domain_key = 'recruitment';-
Create the agent in
agents/my_agent/:Dockerfile— copy from crm_agentrequirements.txt— celery, redis, asyncpg, your depssrc/worker.py— Celery task matching thecelery_task_nameabove
-
Add to
docker-compose.ymlfollowing the crm_agent pattern.
The orchestrator will automatically discover your agent via the registry and include its capabilities in LLM prompts for matching domains.
No code changes required. Just add a schema definition:
INSERT INTO entity_definitions (entity_type, display_name, validation_schema) VALUES (
'invoice',
'Invoice',
'{
"type": "object",
"required": ["amount", "currency", "client_name"],
"properties": {
"amount": { "type": "number", "minimum": 0 },
"currency": { "type": "string", "enum": ["USD", "EUR", "GBP"] },
"client_name": { "type": "string" },
"due_date": { "type": "string", "format": "date" }
}
}'
);The CRM agent will validate all upsert_entity calls for type invoice against this schema automatically.
krastix/
├── docker-compose.yml # 7 services: redis, orchestrator, crm, form, research, doc, communication
├── init.sql # Full PostgreSQL schema (entities, registry, pgvector)
├── migrations/ # Incremental SQL migrations
├── ARCHITECTURE.md # Detailed architecture docs
│
├── orchestrator/src/
│ ├── main.py # FastAPI app (/chat, /chat/stream SSE, Task Watcher)
│ ├── graph.py # LangGraph: Planner + Dispatcher + stream_message()
│ ├── schemas.py # DelegateTask/QueueBatch tools + registry helpers
│ └── services/memory.py # Namespace-isolated RAG (pgvector + nomic-embed-text)
│
├── agents/
│ ├── crm_agent/src/worker.py # Universal entity CRUD + OCC + callback
│ ├── form_agent/src/worker.py # Tally.so form management + callback
│ ├── communication_agent/src/worker.py # Gmail send execution + callback
│ ├── research_agent/src/ # Firecrawl + LinkedIn (FastAPI + LangGraph)
│ └── doc_agent/src/ # PDF/image VDU pipeline (Groq Vision + Ollama)
│ ├── llm_router.py # Dual-LLM factory (vision→Groq, text→Ollama)
│ └── pipeline/ # preprocess → extract → ground → audit
│
├── shared/
│ ├── database.py # Async PostgreSQL pool + stale task queries
│ ├── callbacks.py # Agent → Orchestrator callback (httpx)
│ └── mq.py # Celery configuration
│
└── frontend/src/ # React + Vite chat interface (SSE consumer)
| Component | Technology |
|---|---|
| LLM (Text) | Ollama (qwen2.5:14b-instruct-q5_K_M) — local, free |
| LLM (Vision) | Groq (llama-3.2-90b-vision-preview) — cloud, fast |
| Embeddings | nomic-embed-text (768d) |
| Orchestration | LangGraph + LangChain |
| API | FastAPI (async) + SSE streaming |
| Queue | Celery + Redis 7 |
| Database | PostgreSQL (Supabase) + pgvector |
| Communications | Gmail API (OAuth2) + on-demand summarize/reply flow |
| Validation | JSON Schema (jsonschema) |
| Concurrency | Optimistic Concurrency Control |
| Reliability | Agent callbacks + Task Watcher |
| Containers | Docker Compose (7 services) |
| Frontend | React + Vite (SSE consumer) |
# Orchestrator (with hot reload)
cd orchestrator && uvicorn src.main:app --reload --port 8000
# Research Agent
cd agents/research_agent && uvicorn src.main:app --reload --port 8001
# Doc Agent
cd agents/doc_agent && uvicorn src.main:app --reload --port 8002
# CRM Worker
celery -A shared.mq:celery_app worker -Q crm_queue --loglevel=info
# Form Worker
celery -A shared.mq:celery_app worker -Q form_queue --loglevel=info
# Communication Worker
celery -A shared.mq:celery_app worker -Q communication_queue --loglevel=info# All containers
docker compose logs -f
# Specific service
docker compose logs -f orchestrator
docker compose logs -f crm_agent# Send a chat message from inside the orchestrator container
docker exec krastix-orchestrator python -c "
import httpx, asyncio
async def test():
async with httpx.AsyncClient(timeout=300) as c:
r = await c.post('http://localhost:8000/api/v1/chat', json={
'user_id': 'YOUR_UUID',
'domain': 'recruitment',
'message': 'Research the latest AI tools',
'session_id': 'test-001'
})
print(r.status_code, r.json())
asyncio.run(test())
"Private / Proprietary