-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathmain.py
More file actions
229 lines (184 loc) · 7.53 KB
/
Copy pathmain.py
File metadata and controls
229 lines (184 loc) · 7.53 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
"""
ContextCortex Main Application Server.
Provides FastAPI entrypoint, ASGI Authentication Middleware (MCP 2026-07-28 RFC 9728),
Static Admin UI mounting, FastMCP SSE / Streamable HTTP routing, and container cold-start resilience.
"""
import asyncio
from contextlib import asynccontextmanager
import json
import logging
import os
import threading
from fastapi import FastAPI
from fastapi.responses import JSONResponse, RedirectResponse
from fastapi.staticfiles import StaticFiles
import uvicorn
from app.api.routes import router as admin_router
from app.mcp.mcp_server import mcp_server
from app.services.auth import (
AuthenticationError,
ForbiddenError,
get_auth_service,
set_current_auth_context,
)
from app.services.database import init_db
import app.services.indexing as indexer
from app.services.indexing import run_full_indexing, VAULT_PATH
from app.services.poller import start_poller_daemon, stop_poller_daemon
from app.services.vector_store.manager import VectorStoreManager
# Configure logging
logging.basicConfig(level=logging.INFO, format="%(asctime)s [%(levelname)s] %(message)s")
logger = logging.getLogger("contextcortex")
class McpNotificationFilter(logging.Filter):
"""Filters out premature MCP notification warnings caused by client race conditions."""
def filter(self, record: logging.LogRecord) -> bool:
msg = record.getMessage()
if "Failed to validate notification" in msg and "notifications/roots/list_changed" in msg:
return False
return True
logging.getLogger().addFilter(McpNotificationFilter())
logging.getLogger("mcp").addFilter(McpNotificationFilter())
def init_application_database() -> None:
"""
Initializes database schema, seeds default configurations,
and bootstraps initial admin API credentials if configured.
"""
try:
init_db(VAULT_PATH)
admin_initial_key = os.getenv("ADMIN_INITIAL_KEY")
if admin_initial_key:
auth_service = get_auth_service()
key = auth_service.bootstrap_admin_key(admin_initial_key)
if key and key.secret_key:
logger.info(f"Admin API key initialized successfully (prefix: {key.key_prefix})")
except Exception as e:
logger.error(f"Failed to initialize database: {e}")
# Initialize database on module load
init_application_database()
@asynccontextmanager
async def lifespan(app: FastAPI):
"""Application lifespan manager handling startup indexing, polling daemons, and clean shutdown."""
indexer.main_event_loop = asyncio.get_running_loop()
logger.info("ContextCortex Server starting up...")
# Ensure database connection & bootstrapping is verified on startup
init_application_database()
try:
threading.Thread(target=run_full_indexing, daemon=True).start()
except Exception as e:
logger.error(f"Startup indexing error: {e}")
try:
start_poller_daemon()
except Exception as e:
logger.error(f"Startup poller daemon error: {e}")
if hasattr(mcp_server.session_manager, "_has_started"):
mcp_server.session_manager._has_started = False
async with mcp_server.session_manager.run():
yield
if hasattr(mcp_server.session_manager, "_has_started"):
mcp_server.session_manager._has_started = False
logger.info("ContextCortex Server shutting down...")
try:
stop_poller_daemon()
except Exception as e:
logger.error(f"Shutdown poller daemon error: {e}")
try:
VectorStoreManager.reset_instance()
except Exception:
pass
class AuthMiddleware:
"""
ASGI Authentication Middleware enforcing Bearer token / API key security
when AUTH_ENABLED=true, while allowing public access to metadata, health,
webhooks, and static frontend assets.
"""
def __init__(self, app):
self.app = app
async def __call__(self, scope, receive, send):
if scope["type"] != "http":
await self.app(scope, receive, send)
return
path = scope.get("path", "")
# Public bypass routes
if (
path == "/.well-known/oauth-protected-resource"
or path == "/health"
or path == "/healthz"
or path == "/"
or path.startswith("/assets")
or (path.startswith("/admin") and not path.startswith("/admin/api"))
or path.startswith("/api/webhooks")
):
await self.app(scope, receive, send)
return
auth_service = get_auth_service()
if not auth_service.is_auth_enabled():
bypass_ctx = auth_service.authenticate_token(None)
set_current_auth_context(bypass_ctx)
try:
await self.app(scope, receive, send)
finally:
set_current_auth_context(None)
return
headers = dict(scope.get("headers", []))
auth_header = headers.get(b"authorization", b"").decode("latin1")
try:
auth_ctx = auth_service.authenticate_token(auth_header)
set_current_auth_context(auth_ctx)
try:
await self.app(scope, receive, send)
finally:
set_current_auth_context(None)
except AuthenticationError as e:
resource_indicator = auth_service.resource_indicator
res_meta = f"{resource_indicator}/.well-known/oauth-protected-resource"
www_auth = f'Bearer error="invalid_token", error_description="{str(e)}", resource_metadata="{res_meta}"'
body = json.dumps({"detail": str(e), "error": "Unauthorized"}).encode("utf-8")
await send({
"type": "http.response.start",
"status": 401,
"headers": [
(b"content-type", b"application/json"),
(b"www-authenticate", www_auth.encode("latin1")),
(b"content-length", str(len(body)).encode("latin1")),
],
})
await send({
"type": "http.response.body",
"body": body,
})
except ForbiddenError as e:
body = json.dumps({"detail": str(e), "error": "Forbidden"}).encode("utf-8")
await send({
"type": "http.response.start",
"status": 403,
"headers": [
(b"content-type", b"application/json"),
(b"content-length", str(len(body)).encode("latin1")),
],
})
await send({
"type": "http.response.body",
"body": body,
})
app = FastAPI(title="ContextCortex", version="2.17.0", lifespan=lifespan)
app.add_middleware(AuthMiddleware)
# Include API routes
app.include_router(admin_router)
@app.get("/")
async def root_redirect():
return RedirectResponse(url="/admin/")
www_dir = "frontend/dist" if os.path.exists("frontend/dist/index.html") else "www"
assets_dir = os.path.join(www_dir, "assets") if os.path.exists(os.path.join(www_dir, "assets")) else "www/assets"
app.mount("/admin", StaticFiles(directory=www_dir, html=True), name="admin")
app.mount("/assets", StaticFiles(directory=assets_dir), name="assets")
# Mount FastMCP SSE and Streamable HTTP endpoints
for route in mcp_server.sse_app().routes:
app.routes.append(route)
for route in mcp_server.streamable_http_app().routes:
app.routes.append(route)
@app.get("/health")
@app.get("/healthz")
async def health():
return JSONResponse(content={"status": "healthy"})
if __name__ == "__main__":
uvicorn.run(app, host="0.0.0.0", port=3000)