Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 17 additions & 0 deletions app/core/router.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,16 @@
# SERVICE_MAP entry: this gateway proxies to it, it does not implement
# any ServiceNow logic itself.
"servicenow": "http://servicenow:8097",
# Public API v1 ("M1": GET /v1/usage): the gateway proxies this one
# read-only call into omnibioai-billing's existing
# GET /billing/organizations/{id}/subscription/usage-limits, reusing
# the same build_upstream_headers()-forwarded bearer token every
# other SERVICE_MAP entry relies on -- billing-service independently
# verifies it against the same shared-platform JWT secret (see its
# app/core/iam.py::_verify_caller), exactly like omnibioai-rag does
# for the "rag" entry above. See SERVICE_PERMISSION_MAP below for the
# permission this requires.
"billing": "http://billing-service:8005",
}


Expand Down Expand Up @@ -51,6 +61,13 @@ def resolve_service(service: str) -> str | None:
# finer check" pattern omnibioai-rag's dataset.read/app/api/iam.py
# already established, not a new authorization model.
"servicenow": "servicenow_incident.read",
# Public API v1 (GET /v1/usage): omnibioai-auth's Permission Registry
# already has "usage.read" registered as "reserved -- not yet
# enforced by any route" -- the exact same state dataset.read/
# model.use/workflow.execute were in before this gateway's IAM
# Foundation integration made them real. This is usage.read's first
# real consumer.
"billing": "usage.read",
}


Expand Down
180 changes: 167 additions & 13 deletions app/routes/v1.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,13 +27,20 @@
from app.core.config import Config
from app.core.router import resolve_service
from app.routes.gateway import build_upstream_headers, proxy
from app.services.literature_contract import UnsupportedRequestError, build_public_answer, build_rag_query
from app.services.literature_contract import (
UnsupportedRequestError,
build_public_answer,
build_public_search,
build_rag_query,
build_rag_search_query,
)
from app.services.v1_store import V1Store, request_fingerprint

router = APIRouter(prefix="/v1")
store = V1Store(Config.V1_REDIS_URL, Config.USAGE_REDIS_URL)

ANSWER_RESOURCE = "literature.answer"
SEARCH_RESOURCE = "literature.search"
_IDEMPOTENCY_KEY = re.compile(r"^[\x21-\x7e]{1,255}$")


Expand Down Expand Up @@ -70,13 +77,17 @@ async def _rate_limited(request: Request, subject: str, request_id: str):
return headers, None


async def _forward(request: Request, method: str, path: str, body=None):
url = f"{resolve_service('rag')}/{path}"
async def _forward_to(service: str, request: Request, method: str, path: str, body=None):
url = f"{resolve_service(service)}/{path}"
if request.url.query:
url = f"{url}?{request.url.query}"
return await proxy.forward(url=url, method=method, headers=build_upstream_headers(request), body=body)


async def _forward(request: Request, method: str, path: str, body=None):
return await _forward_to("rag", request, method, path, body)


def _upstream_error(status: int, response, request_id: str, headers: dict):
if status >= 500:
return _error(502 if status == 500 else status, "upstream_error",
Expand All @@ -85,8 +96,18 @@ def _upstream_error(status: int, response, request_id: str, headers: dict):
"The literature service rejected the request.", request_id, headers, detail=response)


@router.post("/literature/answers")
async def literature_answers(request: Request):
async def _handle_billable_literature_call(
request: Request, *, resource: str, build_rag_body, build_public_response, quota_exceeded_message: str,
):
"""Shared lifecycle for every billable /v1/literature/* call
(currently /answers and /search): auth'd-org check, contract
translation, rate limit, idempotency replay, quota check, forward to
RAG, response translation, usage emission, quota consumption. The two
callers differ only in which resource they bill, how they translate
their request/response, and their quota-exceeded wording -- every
other step (in particular the idempotency/quota/usage sequencing)
must stay identical between them, so it lives here once rather than
as two copies that could silently drift apart."""
request_id = getattr(request.state, "trace_id", "")
subject, org_id, user_id = _caller(request)
if not org_id:
Expand All @@ -99,7 +120,7 @@ async def literature_answers(request: Request):
return _error(400, "invalid_request", "Request body must be JSON.", request_id)

try:
rag_body = build_rag_query(body)
rag_body = build_rag_body(body)
except UnsupportedRequestError as exc:
return _error(400, "unsupported_request", exc.message, request_id, detail={"field": exc.field})

Expand Down Expand Up @@ -127,13 +148,11 @@ async def literature_answers(request: Request):
return _error(409, "idempotency_in_progress",
"A request with this Idempotency-Key is still running.", request_id, headers)

remaining = await store.quota_remaining(org_id, ANSWER_RESOURCE)
remaining = await store.quota_remaining(org_id, resource)
if remaining is not None and remaining <= 0:
if idempotency_key is not None:
await store.idempotency_finish(subject, idempotency_key, fingerprint, 402, None)
return _error(402, "quota_exceeded",
"Your organization has used its included answers. Add a payment method or upgrade the plan.",
request_id, headers)
return _error(402, "quota_exceeded", quota_exceeded_message, request_id, headers)

started = time.monotonic()
status, response = await _forward(request, "POST", "v1/query", rag_body)
Expand All @@ -144,7 +163,7 @@ async def literature_answers(request: Request):
await store.idempotency_finish(subject, idempotency_key, fingerprint, status, response)
return _upstream_error(status, response, request_id, headers)

public_response = build_public_answer(
public_response = build_public_response(
response, domain=body.get("domain"), request_id=request_id, latency_ms=latency_ms,
)
if idempotency_key is not None:
Expand All @@ -154,7 +173,7 @@ async def literature_answers(request: Request):
await store.emit_usage(
org_id=org_id,
user_id=user_id,
resource=ANSWER_RESOURCE,
resource=resource,
trace_id=request_id,
dedup_key=f"{subject}:{idempotency_key}" if idempotency_key else request_id,
metadata={
Expand All @@ -164,10 +183,38 @@ async def literature_answers(request: Request):
"idempotency_key_sha256": request_fingerprint(idempotency_key) if idempotency_key else None,
},
)
await store.consume_quota(org_id, ANSWER_RESOURCE)
await store.consume_quota(org_id, resource)
return JSONResponse(public_response, status_code=status, headers={**headers, "X-Request-Id": request_id})


@router.post("/literature/answers")
async def literature_answers(request: Request):
return await _handle_billable_literature_call(
request,
resource=ANSWER_RESOURCE,
build_rag_body=build_rag_query,
build_public_response=build_public_answer,
quota_exceeded_message="Your organization has used its included answers. "
"Add a payment method or upgrade the plan.",
)


@router.post("/literature/search")
async def literature_search(request: Request):
"""Billable unit: 1 search (see the design doc's pricing table --
priced around 1/10th of an answer). Retrieval only: never invokes an
LLM, via RAG's mode="search" (app/services/literature_contract.py's
build_rag_search_query sets it)."""
return await _handle_billable_literature_call(
request,
resource=SEARCH_RESOURCE,
build_rag_body=build_rag_search_query,
build_public_response=build_public_search,
quota_exceeded_message="Your organization has used its included searches. "
"Add a payment method or upgrade the plan.",
)


@router.get("/literature/studies")
async def literature_studies(request: Request):
"""Free: the queryable studies/domains, from omnibioai-rag GET /v1/studies.
Expand All @@ -181,3 +228,110 @@ async def literature_studies(request: Request):
if not 200 <= status < 300:
return _upstream_error(status, response, request_id, headers)
return JSONResponse(response, status_code=status, headers={**headers, "X-Request-Id": request_id})


@router.get("/literature/domains")
async def literature_domains(request: Request):
"""Free: the queryable research domains, under the frozen public
name the design doc uses ("domain", not RAG's internal "study").
Same underlying data as /literature/studies above (kept as-is for
backward compatibility) via the same omnibioai-rag GET /v1/studies
call, reshaped to the public contract. Rate-limited like every /v1
call, never billed."""
request_id = getattr(request.state, "trace_id", "")
subject, _, _ = _caller(request)
headers, limited = await _rate_limited(request, subject, request_id)
if limited:
return limited
status, response = await _forward(request, "GET", "v1/studies")
if not 200 <= status < 300:
return _upstream_error(status, response, request_id, headers)
domains = [
{"name": s.get("name"), "abstract_count": s.get("abstract_count")}
for s in (response.get("studies") or [])
]
return JSONResponse({"domains": domains}, status_code=status, headers={**headers, "X-Request-Id": request_id})


@router.get("/usage")
async def literature_usage(request: Request):
"""Free: the caller's organization's included/used/remaining units
for the current billing period, from omnibioai-billing's existing
GET /billing/organizations/{id}/subscription/usage-limits -- the
gateway's first synchronous call into billing-service (see
app/core/router.py's SERVICE_MAP/SERVICE_PERMISSION_MAP entries,
gated on usage.read). Rate-limited like every /v1 call, never
billed.

Estimated charge in dollars (also mentioned in the design doc) is
deliberately omitted: that needs billing's cost-summary endpoint and
its own start_date/end_date period math, which this does not yet do.
Reporting a wrong number would be worse than omitting it -- the same
principle app/services/literature_contract.py applies to token
counts.
"""
request_id = getattr(request.state, "trace_id", "")
subject, org_id, _ = _caller(request)
if not org_id:
return _error(403, "organization_required",
"This API is billed to an organization; your account has none.", request_id)

headers, limited = await _rate_limited(request, subject, request_id)
if limited:
return limited

status, response = await _forward_to(
"billing", request, "GET", f"billing/organizations/{org_id}/subscription/usage-limits",
)
if status == 404:
return _error(404, "no_active_plan", "Your organization has no active billing plan.", request_id, headers)
if not 200 <= status < 300:
return _error(502 if status >= 500 else status, "upstream_error",
"The billing service failed to report usage.", request_id, headers,
detail=response if status < 500 else None)

usage = [
{
"resource": item.get("resource"),
"unit": item.get("unit"),
"period": item.get("period"),
"included": item.get("included"),
"used": item.get("used"),
"remaining": item.get("remaining"),
}
for item in (response.get("limits") or [])
]
return JSONResponse(
{"plan": response.get("plan_name"), "as_of": response.get("as_of"), "usage": usage},
status_code=200, headers={**headers, "X-Request-Id": request_id},
)


@router.get("/models")
async def literature_models(request: Request):
"""Free: the model catalog -- answer/embedding models currently
served, with source and price. Rate-limited like every /v1 call,
never billed. No upstream call: there is exactly one model path
today, RAG's own GPU-hosted default (see the design's "largest
gaps" #4 -- Claude/OpenAI routing and bring-your-own-key do not
exist yet, and build_rag_query already rejects any request that
would need one).

price is null, not a placeholder dollar figure: the design doc's own
pricing section says to measure real GPU cost per answer first
(milestone M0, 1,000 representative questions) before setting
prices, and that measurement has not been run. The model actually
used for a given answer is already reported per-call in
/v1/literature/answers' response (its `model` field is the real
value RAG used; "default" here is the stable identifier a caller
passes back as this endpoint's own `model` request field).
"""
request_id = getattr(request.state, "trace_id", "")
subject, _, _ = _caller(request)
headers, limited = await _rate_limited(request, subject, request_id)
if limited:
return limited
models = [
{"model": "default", "source": "omnibioai_gpu", "billed_by": "query", "price": None},
]
return JSONResponse({"models": models}, status_code=200, headers={**headers, "X-Request-Id": request_id})
56 changes: 56 additions & 0 deletions app/services/literature_contract.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,62 @@ def build_rag_query(body: dict) -> dict:
return query


def build_rag_search_query(body: dict) -> dict:
"""Translate a public /v1/literature/search request body into RAG's
POST /v1/query body (with mode="search" set by the caller, not here --
that's the one field this function doesn't own, since it's not part of
the public contract). Raises UnsupportedRequestError for a missing/
blank question, the only field this endpoint requires.

Unlike build_rag_query above, model/use_own_key/stream are not
validated here at all: search is retrieval-only regardless of what a
caller sends for them, so there is no "unsupported" case for fields
that were never going to change this call's behavior.
"""
if not isinstance(body, dict):
raise UnsupportedRequestError("body", "Request body must be a JSON object.")

question = body.get("question")
if not isinstance(question, str) or not question.strip():
raise UnsupportedRequestError("question", "\"question\" is required and must be a non-empty string.")

query: dict = {"query": question, "study": body.get("domain") or "default", "mode": "search"}
max_results = body.get("max_results")
if isinstance(max_results, (int, float)) and not isinstance(max_results, bool) and max_results > 0:
query["top_k"] = int(max_results)
return query


def build_public_search(rag_result: dict, *, domain, request_id: str, latency_ms: int) -> dict:
"""Translate RAG's /v1/query (mode="search") response into the frozen
public /v1/literature/search response shape: ranked documents, no
generated answer."""
documents = rag_result.get("documents") or []
results = [
{
"pmid": doc.get("pmid"),
"title": doc.get("title"),
"year": doc.get("year"),
"score": doc.get("citation_confidence", doc.get("similarity_score")),
# The design doc calls this a "snippet"; RAG only has the full
# abstract text to offer, not a separately-generated excerpt,
# so that's what's returned under this name rather than
# fabricating a truncation.
"snippet": doc.get("abstract"),
}
for doc in documents
]
return {
"id": f"srch_{request_id or uuid.uuid4().hex}",
"results": results,
"domain": rag_result.get("study", domain),
"usage": {
"searches": 1,
"latency_ms": latency_ms,
},
}


def build_public_answer(rag_result: dict, *, domain, request_id: str, latency_ms: int) -> dict:
"""Translate RAG's /v1/query response into the frozen public
/v1/literature/answers response shape."""
Expand Down
Loading
Loading