Compare commits
2 Commits
7227d2d2bd
...
58bfa07385
| Author | SHA1 | Date | |
|---|---|---|---|
| 58bfa07385 | |||
| be8103c1d2 |
16
.env.example
Normal file
16
.env.example
Normal file
@@ -0,0 +1,16 @@
|
|||||||
|
GO_API_BASE_URL=
|
||||||
|
NATS_URL=
|
||||||
|
REDIS_HOST=
|
||||||
|
REDIS_PORT=
|
||||||
|
REDIS_PASSWORD=
|
||||||
|
INTERNAL_API_KEY=
|
||||||
|
DB_HOST=
|
||||||
|
DB_PORT=
|
||||||
|
DB_NAME=
|
||||||
|
DB_USER=
|
||||||
|
DB_PASSWORD=
|
||||||
|
NATS_USER=
|
||||||
|
NATS_PASSWORD=
|
||||||
|
ANTHROPIC_API_KEY=
|
||||||
|
LLM_MODEL=claude-opus-4-8
|
||||||
|
LOG_LEVEL=INFO
|
||||||
5
.gitignore
vendored
5
.gitignore
vendored
@@ -7,3 +7,8 @@ venv/
|
|||||||
|
|
||||||
.idea/
|
.idea/
|
||||||
.vscode/
|
.vscode/
|
||||||
|
|
||||||
|
# secrets
|
||||||
|
.env
|
||||||
|
.env.*
|
||||||
|
!.env.example
|
||||||
|
|||||||
@@ -1,27 +1,28 @@
|
|||||||
"""System configuration for LogiFlow AI."""
|
"""System configuration for LogiFlow AI."""
|
||||||
import os
|
import os
|
||||||
|
|
||||||
# Infrastructure Connections — values come from environment / .env file.
|
# Infrastructure Connections — values come from environment / .env file
|
||||||
# Hardcoded strings are last-resort defaults; set the matching env var in production.
|
# (see .env.example). Defaults point at localhost with no credentials so the
|
||||||
NATS_URL = os.getenv("NATS_URL", "nats://doormile:Package@321#@66.116.226.161:4223")
|
# system runs in local-fallback mode; real hosts and secrets are never in code.
|
||||||
NATS_HOST = os.getenv("NATS_HOST", "66.116.226.161")
|
NATS_URL = os.getenv("NATS_URL", "nats://localhost:4222")
|
||||||
NATS_PORT = int(os.getenv("NATS_PORT", "4223"))
|
NATS_HOST = os.getenv("NATS_HOST", "localhost")
|
||||||
NATS_USER = os.getenv("NATS_USER", "doormile")
|
NATS_PORT = int(os.getenv("NATS_PORT", "4222"))
|
||||||
NATS_PASSWORD = os.getenv("NATS_PASSWORD", "Package@321#")
|
NATS_USER = os.getenv("NATS_USER", "")
|
||||||
|
NATS_PASSWORD = os.getenv("NATS_PASSWORD", "")
|
||||||
|
|
||||||
REDIS_HOST = os.getenv("REDIS_HOST", "66.116.226.255")
|
REDIS_HOST = os.getenv("REDIS_HOST", "localhost")
|
||||||
REDIS_PORT = int(os.getenv("REDIS_PORT", "6380"))
|
REDIS_PORT = int(os.getenv("REDIS_PORT", "6379"))
|
||||||
REDIS_PASSWORD = os.getenv("REDIS_PASSWORD", "Package@321#")
|
REDIS_PASSWORD = os.getenv("REDIS_PASSWORD", "")
|
||||||
|
|
||||||
# Postgres — individual params to avoid @ in password breaking DSN parsing
|
# Postgres — individual params to avoid @ in password breaking DSN parsing
|
||||||
DB_HOST = os.getenv("DB_HOST", "31.97.228.132")
|
DB_HOST = os.getenv("DB_HOST", "localhost")
|
||||||
DB_PORT = int(os.getenv("DB_PORT", "5433"))
|
DB_PORT = int(os.getenv("DB_PORT", "5432"))
|
||||||
DB_NAME = os.getenv("DB_NAME", "logistics")
|
DB_NAME = os.getenv("DB_NAME", "logistics")
|
||||||
DB_USER = os.getenv("DB_USER", "admin")
|
DB_USER = os.getenv("DB_USER", "postgres")
|
||||||
DB_PASSWORD = os.getenv("DB_PASSWORD", "Package@321#")
|
DB_PASSWORD = os.getenv("DB_PASSWORD", "")
|
||||||
|
|
||||||
GO_API_BASE_URL = os.getenv("GO_API_BASE_URL", "http://localhost:8080")
|
GO_API_BASE_URL = os.getenv("GO_API_BASE_URL", "http://localhost:8080")
|
||||||
INTERNAL_API_KEY = os.getenv("INTERNAL_API_KEY", "doormile-internal-2024")
|
INTERNAL_API_KEY = os.getenv("INTERNAL_API_KEY", "")
|
||||||
|
|
||||||
# Route Optimization API (Valhalla-backed road routing). The ExpressDispatchAgent
|
# Route Optimization API (Valhalla-backed road routing). The ExpressDispatchAgent
|
||||||
# calls its Doormile endpoint (/api/v1/optimization/doormile/sequence) to order a
|
# calls its Doormile endpoint (/api/v1/optimization/doormile/sequence) to order a
|
||||||
|
|||||||
@@ -110,24 +110,33 @@ class QuestionManager:
|
|||||||
raise TimeoutError(f"Question {question.question_id} timed out waiting for human input")
|
raise TimeoutError(f"Question {question.question_id} timed out waiting for human input")
|
||||||
finally:
|
finally:
|
||||||
self._futures.pop(question.question_id, None)
|
self._futures.pop(question.question_id, None)
|
||||||
|
self._pending.pop(question.question_id, None)
|
||||||
|
|
||||||
def answer_question(self, question_id: str, answer: Any) -> bool:
|
def answer_question(self, question_id: str, answer: Any) -> bool:
|
||||||
"""Called when human operator submits an answer from the UI."""
|
"""Called when human operator submits an answer from the UI.
|
||||||
|
|
||||||
|
Returns False if the question is unknown or was already answered."""
|
||||||
question = self._pending.get(question_id)
|
question = self._pending.get(question_id)
|
||||||
if not question:
|
if not question:
|
||||||
logger.warning(f"No pending question found for id {question_id}")
|
logger.warning(f"No pending question found for id {question_id}")
|
||||||
return False
|
return False
|
||||||
|
if question.answered:
|
||||||
|
logger.warning(f"Question {question_id} already answered; ignoring")
|
||||||
|
return False
|
||||||
|
|
||||||
question.answered = True
|
question.answered = True
|
||||||
question.answer = answer
|
question.answer = answer
|
||||||
|
|
||||||
fut = self._futures.get(question_id)
|
fut = self._futures.get(question_id)
|
||||||
if fut and not fut.done():
|
if fut and not fut.done():
|
||||||
fut.set_result(answer)
|
fut.set_result(answer) # wait_for_answer() removes it from _pending
|
||||||
return True
|
else:
|
||||||
|
self._pending.pop(question_id, None) # nobody waiting; don't leak
|
||||||
return True
|
return True
|
||||||
|
|
||||||
|
def list_pending(self) -> List[Question]:
|
||||||
|
return [q for q in self._pending.values() if not q.answered]
|
||||||
|
|
||||||
def get_pending(self, question_id: str) -> Optional[Question]:
|
def get_pending(self, question_id: str) -> Optional[Question]:
|
||||||
return self._pending.get(question_id)
|
return self._pending.get(question_id)
|
||||||
|
|
||||||
|
|||||||
@@ -1,33 +1,27 @@
|
|||||||
"""Structured logging for LogiFlow AI — loguru setup with standard logging fallback."""
|
"""Structured logging for LogiFlow AI — single loguru setup imported by all modules."""
|
||||||
import sys
|
import sys
|
||||||
import os
|
import os
|
||||||
import logging
|
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
from loguru import logger
|
||||||
|
|
||||||
try:
|
Path("logs").mkdir(exist_ok=True)
|
||||||
from loguru import logger
|
|
||||||
|
|
||||||
Path("logs").mkdir(exist_ok=True)
|
logger.remove()
|
||||||
|
|
||||||
logger.remove()
|
logger.add(
|
||||||
|
sys.stderr,
|
||||||
|
format="{time:YYYY-MM-DD HH:mm:ss} | {level: <8} | {name}:{function}:{line} - {message}",
|
||||||
|
level=os.getenv("LOG_LEVEL", "INFO"),
|
||||||
|
colorize=False,
|
||||||
|
)
|
||||||
|
|
||||||
logger.add(
|
logger.add(
|
||||||
sys.stderr,
|
"logs/logiflow_{time:YYYY-MM-DD}.log",
|
||||||
format="{time:YYYY-MM-DD HH:mm:ss} | {level: <8} | {name}:{function}:{line} - {message}",
|
rotation="100 MB",
|
||||||
level=os.getenv("LOG_LEVEL", "INFO"),
|
retention="30 days",
|
||||||
colorize=False,
|
compression="gz",
|
||||||
)
|
level="DEBUG",
|
||||||
|
format="{time:YYYY-MM-DD HH:mm:ss} | {level: <8} | {name}:{function}:{line} - {message}",
|
||||||
logger.add(
|
)
|
||||||
"logs/logiflow_{time:YYYY-MM-DD}.log",
|
|
||||||
rotation="100 MB",
|
|
||||||
retention="30 days",
|
|
||||||
compression="gz",
|
|
||||||
level="DEBUG",
|
|
||||||
format="{time:YYYY-MM-DD HH:mm:ss} | {level: <8} | {name}:{function}:{line} - {message}",
|
|
||||||
)
|
|
||||||
except ImportError:
|
|
||||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s | %(levelname)-8s | %(message)s")
|
|
||||||
logger = logging.getLogger("logiflow")
|
|
||||||
|
|
||||||
__all__ = ["logger"]
|
__all__ = ["logger"]
|
||||||
|
|||||||
@@ -8,11 +8,8 @@ from collections import defaultdict
|
|||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from enum import Enum
|
from enum import Enum
|
||||||
|
|
||||||
try:
|
import nats
|
||||||
import nats
|
import nats.js.errors
|
||||||
import nats.js.errors
|
|
||||||
except ImportError:
|
|
||||||
nats = None
|
|
||||||
|
|
||||||
from core.types import AgentMessage, MessageType
|
from core.types import AgentMessage, MessageType
|
||||||
from core.logger import logger
|
from core.logger import logger
|
||||||
@@ -59,9 +56,6 @@ class MessageBus:
|
|||||||
|
|
||||||
async def connect(self):
|
async def connect(self):
|
||||||
"""Connect to NATS and create the logistics JetStream stream."""
|
"""Connect to NATS and create the logistics JetStream stream."""
|
||||||
if nats is None:
|
|
||||||
logger.warning("nats-py not installed; message bus running in local in-memory mode")
|
|
||||||
return
|
|
||||||
self._nc = await nats.connect(
|
self._nc = await nats.connect(
|
||||||
servers=[f"nats://{NATS_HOST}:{NATS_PORT}"],
|
servers=[f"nats://{NATS_HOST}:{NATS_PORT}"],
|
||||||
user=NATS_USER,
|
user=NATS_USER,
|
||||||
|
|||||||
0
core/skills/__init__.py
Normal file
0
core/skills/__init__.py
Normal file
@@ -33,7 +33,7 @@ class Tool:
|
|||||||
metadata: Dict[str, Any] = field(default_factory=dict)
|
metadata: Dict[str, Any] = field(default_factory=dict)
|
||||||
|
|
||||||
def to_schema(self) -> Dict[str, Any]:
|
def to_schema(self) -> Dict[str, Any]:
|
||||||
"""Generate OpenAI/Claude compatible JSON Schema for tool calling."""
|
"""Generate an Anthropic Messages API tool definition (name / description / input_schema)."""
|
||||||
properties = {}
|
properties = {}
|
||||||
required = []
|
required = []
|
||||||
|
|
||||||
@@ -56,7 +56,7 @@ class Tool:
|
|||||||
return {
|
return {
|
||||||
"name": self.name,
|
"name": self.name,
|
||||||
"description": self.description,
|
"description": self.description,
|
||||||
"parameters": {
|
"input_schema": {
|
||||||
"type": "object",
|
"type": "object",
|
||||||
"properties": properties,
|
"properties": properties,
|
||||||
"required": required,
|
"required": required,
|
||||||
|
|||||||
@@ -3,6 +3,7 @@
|
|||||||
Doormile Full System Test Suite v2
|
Doormile Full System Test Suite v2
|
||||||
Fixed: correct routes, 2-step login, pricing payload, health endpoint
|
Fixed: correct routes, 2-step login, pricing payload, health endpoint
|
||||||
"""
|
"""
|
||||||
|
import os
|
||||||
import asyncio
|
import asyncio
|
||||||
import aiohttp
|
import aiohttp
|
||||||
import json
|
import json
|
||||||
@@ -13,14 +14,17 @@ import redis.asyncio as aioredis
|
|||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
|
|
||||||
# ── Config ────────────────────────────────────────────────────────────────────
|
# ── Config ────────────────────────────────────────────────────────────────────
|
||||||
API_BASE = "https://api.doormile.com"
|
from dotenv import load_dotenv
|
||||||
NATS_URL = "nats://66.116.226.161:4223"
|
load_dotenv()
|
||||||
NATS_USER = "doormile"
|
|
||||||
NATS_PASSWORD = "Package@321#"
|
API_BASE = os.getenv("GO_API_BASE_URL", "https://api.doormile.com")
|
||||||
REDIS_HOST = "66.116.226.255"
|
NATS_URL = os.environ["NATS_URL"]
|
||||||
REDIS_PORT = 6380
|
NATS_USER = os.getenv("NATS_USER", "")
|
||||||
REDIS_PASSWORD = "Package@321#"
|
NATS_PASSWORD = os.getenv("NATS_PASSWORD", "")
|
||||||
INTERNAL_KEY = "doormile-internal-2024"
|
REDIS_HOST = os.environ["REDIS_HOST"]
|
||||||
|
REDIS_PORT = int(os.getenv("REDIS_PORT", "6379"))
|
||||||
|
REDIS_PASSWORD = os.getenv("REDIS_PASSWORD", "")
|
||||||
|
INTERNAL_KEY = os.environ["INTERNAL_API_KEY"]
|
||||||
|
|
||||||
# Test customer — uses the one created by Window 1 test
|
# Test customer — uses the one created by Window 1 test
|
||||||
TEST_PHONE = "9900000001"
|
TEST_PHONE = "9900000001"
|
||||||
|
|||||||
9
main.py
9
main.py
@@ -204,11 +204,12 @@ Usage:
|
|||||||
python main.py --portal Launch Streamlit customer portal
|
python main.py --portal Launch Streamlit customer portal
|
||||||
python main.py --help Show this help
|
python main.py --help Show this help
|
||||||
|
|
||||||
Infrastructure (set in .env):
|
Infrastructure (set in .env — see .env.example):
|
||||||
NATS nats://doormile@66.116.226.161:4223
|
NATS NATS_URL / NATS_USER / NATS_PASSWORD
|
||||||
Redis 66.116.226.255:6380
|
Redis REDIS_HOST / REDIS_PORT / REDIS_PASSWORD
|
||||||
PG DB_HOST / DB_PORT / DB_NAME / DB_USER / DB_PASSWORD
|
PG DB_HOST / DB_PORT / DB_NAME / DB_USER / DB_PASSWORD
|
||||||
API GO_API_BASE_URL
|
API GO_API_BASE_URL / INTERNAL_API_KEY
|
||||||
|
LLM ANTHROPIC_API_KEY / LLM_MODEL
|
||||||
""")
|
""")
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
3
pytest.ini
Normal file
3
pytest.ini
Normal file
@@ -0,0 +1,3 @@
|
|||||||
|
[pytest]
|
||||||
|
testpaths = tests
|
||||||
|
asyncio_mode = auto
|
||||||
@@ -27,7 +27,8 @@ async def test_tool_registry_registration_and_execution():
|
|||||||
schemas = registry.get_schemas()
|
schemas = registry.get_schemas()
|
||||||
assert len(schemas) == 1
|
assert len(schemas) == 1
|
||||||
assert schemas[0]["name"] == "test_tool"
|
assert schemas[0]["name"] == "test_tool"
|
||||||
assert "order_id" in schemas[0]["parameters"]["properties"]
|
assert "order_id" in schemas[0]["input_schema"]["properties"]
|
||||||
|
assert "parameters" not in schemas[0] # Anthropic shape, not OpenAI
|
||||||
|
|
||||||
result = await registry.execute_tool("test_tool", order_id="ORD-100")
|
result = await registry.execute_tool("test_tool", order_id="ORD-100")
|
||||||
assert result == {"order": "ORD-100", "items": 1}
|
assert result == {"order": "ORD-100", "items": 1}
|
||||||
@@ -50,6 +51,32 @@ async def test_question_manager_flow():
|
|||||||
assert q.options[0].badge == "12 AM–9 AM"
|
assert q.options[0].badge == "12 AM–9 AM"
|
||||||
assert not q.answered
|
assert not q.answered
|
||||||
|
|
||||||
qm.answer_question(q.question_id, "morning")
|
assert qm.list_pending() == [q]
|
||||||
|
assert qm.answer_question(q.question_id, "morning") is True
|
||||||
assert q.answered is True
|
assert q.answered is True
|
||||||
assert q.answer == "morning"
|
assert q.answer == "morning"
|
||||||
|
assert qm.get_pending(q.question_id) is None # cleaned up, not leaked
|
||||||
|
assert qm.answer_question(q.question_id, "afternoon") is False # double-answer rejected
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_question_manager_wait_and_answer():
|
||||||
|
qm = QuestionManager()
|
||||||
|
q = qm.create_question(prompt="Reassign?", options=["yes", "no"], allow_custom_input=False)
|
||||||
|
|
||||||
|
async def operator():
|
||||||
|
await asyncio.sleep(0.01)
|
||||||
|
qm.answer_question(q.question_id, "yes")
|
||||||
|
|
||||||
|
asyncio.create_task(operator())
|
||||||
|
assert await qm.wait_for_answer(q, timeout_s=1.0) == "yes"
|
||||||
|
assert qm.get_pending(q.question_id) is None
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_question_manager_timeout_cleans_up():
|
||||||
|
qm = QuestionManager()
|
||||||
|
q = qm.create_question(prompt="Reassign?", options=["yes", "no"])
|
||||||
|
with pytest.raises(TimeoutError):
|
||||||
|
await qm.wait_for_answer(q, timeout_s=0.01)
|
||||||
|
assert qm.get_pending(q.question_id) is None
|
||||||
|
|||||||
Reference in New Issue
Block a user