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
1 change: 1 addition & 0 deletions .github/workflows/lif_semantic_search_mcp_server.yml
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ on:
- cloudformation/lif-semantic-search-taskdef-includes.yml
- components/lif/logging/**
- components/lif/datatypes/**
- components/lif/graphql_client/**
- components/lif/openapi_schema_parser/**
- components/lif/semantic_search_service/**
- components/lif/string_utils/**
Expand Down
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
completeness" toggle, because the preview is an intentionally partial document. `format` is not
validated, matching the runtime translator, which calls `jsonschema.validate` without a
`format_checker`
- Query Planner query statistics record the caller as `client`, from an optional `X-LIF-Client`
request header: `unknown` when absent (never a rejected query), `invalid` when not a short
lowercase name. Learner Data Export sends `learner-data-export`, the MCP server
`semantic-search-mcp`, and GraphQL forwards its caller's name or sends `graphql`

### Changed

Expand Down
4 changes: 4 additions & 0 deletions bases/lif/query_planner_restapi/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,10 @@ The planner reads YAML at startup that describes available information sources (

A **malformed** value for either timeout — `"120s"`, or `0` — stops the service at startup with a message naming the variable, per the convention decided in #1179. Unset or empty falls back to the code default.

## Caller identity

`POST /query` and `POST /query_async` read an optional `X-LIF-Client` header naming the caller, and every query statistics event (`LIF_QUERY_STATISTICS` log lines, #341) records it as `client` (#1272). It is optional by design: a missing header is recorded as `unknown`, never a rejected query, so the planner works standalone. A value that is not a short lowercase name (`[a-z0-9][a-z0-9._-]{0,63}`) is recorded as `invalid` rather than as itself, which keeps free text and any person data out of the logs. The in-repo callers send `learner-data-export` (`query_planner_client`) and `graphql`; GraphQL forwards its own caller's name instead when it has one, so MCP traffic arrives as `semantic-search-mcp`.

## Composes
- `datatypes` — `LIFQuery`, `LIFRecord`, `LIFUpdate`, planner-side types
- `exceptions`
Expand Down
19 changes: 12 additions & 7 deletions bases/lif/query_planner_restapi/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,10 +2,10 @@
import yaml
from datetime import datetime
from pathlib import Path
from typing import List
from typing import Annotated, List

from asyncio import sleep
from fastapi import FastAPI, HTTPException, Response, status
from fastapi import FastAPI, Header, HTTPException, Response, status

from lif.datatypes import (
OrchestratorJobResults,
Expand All @@ -17,6 +17,7 @@
)
from lif.exceptions.core import LIFException
from lif.logging.core import get_logger
from lif.query_planner_service import statistics
from lif.query_planner_service.core import LIFQueryPlannerService
from lif.query_planner_service.datatypes import LIFQueryPlannerConfig, LIFQueryPlannerInfoSourceConfig

Expand Down Expand Up @@ -144,14 +145,16 @@ def root() -> dict:
# temporary, and will be removed soon.
# -------------------------------------------------------------------------
@app.post("/query", status_code=status.HTTP_200_OK, response_model=List[LIFRecord])
async def do_run_query_sync(query: LIFQuery, response: Response) -> List[LIFRecord]:
async def do_run_query_sync(
query: LIFQuery, response: Response, client: Annotated[str | None, Header(alias=statistics.CLIENT_HEADER)] = None
) -> List[LIFRecord]:
logger.info("CALL RECEIVED TO /query (sync) API")
try:
# Counted from before the first run_query: that call makes the cache read and the
# orchestrator submission, so starting the clock after it left those round trips
# outside the budget entirely (#571).
start_time = datetime.now()
result = await service.run_query(query, first_run=True)
result = await service.run_query(query, first_run=True, client=client)
if isinstance(result, LIFQueryStatusResponse):
logger.info("Query is still processing, entering polling loop")
delay_in_seconds: int = MIN_POLLING_DELAY_SECONDS
Expand All @@ -173,7 +176,7 @@ async def do_run_query_sync(query: LIFQuery, response: Response) -> List[LIFReco
result = await service.get_query_status(result.query_id)
if result.status == "COMPLETED":
logger.info("Query completed successfully, retrieving results")
result = await service.run_query(query, first_run=False)
result = await service.run_query(query, first_run=False, client=client)
if isinstance(result, list):
logger.info(f"Query completed successfully, returning {len(result)} record(s)")
return result
Expand Down Expand Up @@ -203,10 +206,12 @@ async def do_run_query_sync(query: LIFQuery, response: Response) -> List[LIFReco
# to /query in the future.
# -------------------------------------------------------------------------
@app.post("/query_async", response_model=List[LIFRecord] | LIFQueryStatusResponse)
async def do_run_query(query: LIFQuery, response: Response) -> List[LIFRecord] | LIFQueryStatusResponse:
async def do_run_query(
query: LIFQuery, response: Response, client: Annotated[str | None, Header(alias=statistics.CLIENT_HEADER)] = None
) -> List[LIFRecord] | LIFQueryStatusResponse:
logger.info("CALL RECEIVED TO /query_async API")
try:
result = await service.run_query(query, first_run=True)
result = await service.run_query(query, first_run=True, client=client)
if isinstance(result, LIFQueryStatusResponse):
response.status_code = status.HTTP_202_ACCEPTED
response.headers["Location"] = f"/query/{result.query_id}/status"
Expand Down
2 changes: 1 addition & 1 deletion components/lif/graphql_client/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ Authenticated HTTP client for calling the LIF GraphQL API. Wraps the boilerplate
from lif.graphql_client import graphql_query, graphql_mutation, GraphQLClientException
```

Both functions send `X-API-Key` from `LIF_GRAPHQL_API_KEY` (when set) as the auth header — see CLAUDE.md § "GraphQL API Key Authentication" for the server-side configuration.
Both functions send `X-API-Key` from `LIF_GRAPHQL_API_KEY` (when set) as the auth header — see CLAUDE.md § "GraphQL API Key Authentication" for the server-side configuration. They also always send `X-LIF-Client: semantic-search-mcp`, which GraphQL forwards so the Query Planner's statistics can tell this traffic apart (#1272).

| Function | Purpose |
|---|---|
Expand Down
11 changes: 8 additions & 3 deletions components/lif/graphql_client/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,13 +12,18 @@
GRAPHQL_TIMEOUT_READ = float(os.getenv("SEMANTIC_SEARCH_SERVICE__GRAPHQL_TIMEOUT__READ", "300"))


# Names this caller in the Query Planner's query statistics; GraphQL forwards it (#1272).
LIF_CLIENT_NAME = "semantic-search-mcp"


def _build_headers(api_key: str = "") -> dict:
"""Build request headers, including X-API-Key if provided."""
"""Build request headers: X-LIF-Client always, X-API-Key if provided."""
headers = {"X-LIF-Client": LIF_CLIENT_NAME}
if not api_key:
api_key = LIF_GRAPHQL_API_KEY
if api_key:
return {"X-API-Key": api_key}
return {}
headers["X-API-Key"] = api_key
return headers


async def _post(url: str, json: dict, headers: dict, timeout: httpx.Timeout | None = None) -> httpx.Response:
Expand Down
1 change: 1 addition & 0 deletions components/lif/openapi_to_graphql/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ from lif.openapi_to_graphql import schema_tools
- **Strawberry `info` typing:** dynamic resolvers must annotate the `info` parameter as `strawberry.types.Info` (not `object` / `Any`). Strawberry 0.297+ identifies the parameter by type, not by name.
- **Backend failures raise:** the root query resolver raises on a non-200 from the Query Planner so Strawberry emits an `errors` entry, matching the update mutation. Returning `[]` made every failure indistinguishable from a learner with no data (#1264).
- **Field name preservation:** uses `strawberry.field(name=field_name)` so the wire shape preserves PascalCase entity / camelCase scalar conventions ([`docs/specs/data-model-rules.md`](../../../docs/specs/data-model-rules.md)).
- **Caller forwarding:** the root query resolver forwards the incoming `X-LIF-Client` header to the Query Planner, or sends `graphql` when there is none, so the planner's statistics see the original caller (#1272). It reads the request from `info.context["request"]`, which Strawberry's FastAPI router provides; it does not validate the value, since the planner does.

## Used by
- `bases/lif/api_graphql` — single consumer; the GraphQL service's whole reason for existing is this component.
22 changes: 21 additions & 1 deletion components/lif/openapi_to_graphql/type_factory.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,13 @@
os.getenv("LIF_GRAPHQL_CLIENT_TIMEOUT_SECONDS") or os.getenv("LIF_QUERY_TIMEOUT_SECONDS") or "20"
)

# Callers name themselves to the Query Planner's statistics in this header (#1272). GraphQL
# forwards the name it received, so a query that arrives through here -- the MCP server's --
# keeps its origin; it names itself only when its own caller did not. Validating the value is
# the planner's job, not this relay's.
LIF_CLIENT_HEADER = "X-LIF-Client"
LIF_CLIENT_NAME = "graphql"


# === Constants ===

Expand Down Expand Up @@ -768,6 +775,19 @@ def create_input_type(
return create_nested_input_type(type_name, schema, openapi, created_types, input_type_cache)


def lif_client_headers(info: Any) -> Dict[str, str]:
"""
The X-LIF-Client header to send the Query Planner: the incoming one, else this service's name.

Strawberry's FastAPI router puts the request in `info.context["request"]`; a schema
executed without one (tests, scripts) simply names itself.
"""
context = info.context if isinstance(info.context, dict) else {}
request = context.get("request")
incoming = request.headers.get(LIF_CLIENT_HEADER) if request is not None else None
return {LIF_CLIENT_HEADER: incoming or LIF_CLIENT_NAME}


# === Root Query Type Construction ===


Expand Down Expand Up @@ -822,7 +842,7 @@ async def resolver(self: Any, info: Any, filter: Optional[Any] = None) -> List[A
logger.info(f"Query: {query}")
# Make the backend API call
async with httpx.AsyncClient(timeout=httpx.Timeout(LIF_GRAPHQL_CLIENT_TIMEOUT_SECONDS)) as client:
response = await client.post(query_planner_query_url, json=query)
response = await client.post(query_planner_query_url, json=query, headers=lif_client_headers(info))

if response.status_code == 200:
response_json = response.json()
Expand Down
5 changes: 4 additions & 1 deletion components/lif/query_planner_client/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,9 @@
# Default timeout for Query Planner client API calls (in seconds)
DEFAULT_QUERY_PLANNER_CLIENT_TIMEOUT_SECONDS = 30

# Names this caller in the Query Planner's query statistics (#1272).
LIF_CLIENT_HEADERS = {"X-LIF-Client": "learner-data-export"}


def _get_query_planner_timeout_seconds() -> int:
return int(os.getenv("QUERY_PLANNER_CLIENT_TIMEOUT_SECONDS", str(DEFAULT_QUERY_PLANNER_CLIENT_TIMEOUT_SECONDS)))
Expand Down Expand Up @@ -39,7 +42,7 @@ async def fetch_query_from_query_planner(base_url: str, query: dict) -> list[dic

try:
async for client in _get_query_planner_client():
response = await client.post(url, json=query)
response = await client.post(url, json=query, headers=LIF_CLIENT_HEADERS)
except httpx.TimeoutException as e:
msg = f"Query Planner request timed out due to: {e}"
logger.error(msg)
Expand Down
30 changes: 24 additions & 6 deletions components/lif/query_planner_service/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,7 @@ def _emit_query_planned(
paths_not_in_cache: List[str],
lif_query_plan: LIFQueryPlan | None = None,
correlation_id: str | None = None,
client: str = statistics.CLIENT_UNKNOWN,
) -> None:
"""
Emit the planning-phase statistics event. Never raises -- statistics must not fail a query.
Expand All @@ -72,7 +73,7 @@ def _emit_query_planned(
logger.info(
statistics.format_event(
statistics.build_query_planned_event(
outcome, requested_paths, paths_not_in_cache, lif_query_plan, correlation_id
outcome, requested_paths, paths_not_in_cache, lif_query_plan, correlation_id, client
)
)
)
Expand All @@ -81,12 +82,15 @@ def _emit_query_planned(

# Main function to run a query
# -------------------------------------------------------------------------
async def run_query(self, query: LIFQuery, first_run: bool) -> List[LIFRecord] | LIFQueryStatusResponse:
async def run_query(
self, query: LIFQuery, first_run: bool, client: str | None = None
) -> List[LIFRecord] | LIFQueryStatusResponse:
"""
Execute a LIF query.

Args:
query (LIFQuery): Input query with filter and selected fields.
client (str | None): The raw `X-LIF-Client` header value, or None when absent.

Returns:
List[LIFRecord]: List of matching LIF records (persons) from the database, with only
Expand All @@ -95,6 +99,8 @@ async def run_query(self, query: LIFQuery, first_run: bool) -> List[LIFRecord] |
Raises:
LIFException: If the query fails.
"""
# Reduced here, at the service boundary, so no raw header value reaches a log line.
client = statistics.normalize_client(client)
try:
# Send the query to the LIF Cache service
lif_records: List[LIFRecord] = await query_lif_cache(
Expand All @@ -120,7 +126,9 @@ async def run_query(self, query: LIFQuery, first_run: bool) -> List[LIFRecord] |
# Guarded on first_run: the sync /query path calls run_query twice, and the
# second call lands here once the orchestrator has filled the cache (#341 gap 3).
if first_run:
self._emit_query_planned(statistics.OUTCOME_SERVED_FROM_CACHE, lif_fragment_paths, [])
self._emit_query_planned(
statistics.OUTCOME_SERVED_FROM_CACHE, lif_fragment_paths, [], client=client
)
return lif_records
logger.info(f"LIF Record does not contain all requested fields, missing: {lif_fragment_paths_not_found}")

Expand All @@ -140,7 +148,10 @@ async def run_query(self, query: LIFQuery, first_run: bool) -> List[LIFRecord] |
f"No information sources found for the requested LIF fragment paths: {lif_fragment_paths}. Returning LIF records found in cache."
)
self._emit_query_planned(
statistics.OUTCOME_NO_SOURCES_AVAILABLE, lif_fragment_paths, lif_fragment_paths_not_found
statistics.OUTCOME_NO_SOURCES_AVAILABLE,
lif_fragment_paths,
lif_fragment_paths_not_found,
client=client,
)
return lif_records

Expand All @@ -162,11 +173,12 @@ async def run_query(self, query: LIFQuery, first_run: bool) -> List[LIFRecord] |
lif_fragment_paths,
lif_fragment_paths_not_found,
lif_query_plan,
client=client,
)
return lif_records

lif_query_planner_job = LIFQueryPlannerJob(
job_id=orchestrator_job_request_response.run_id, query=query, status="PENDING"
job_id=orchestrator_job_request_response.run_id, query=query, status="PENDING", client=client
)

prune_job_store()
Expand All @@ -179,6 +191,7 @@ async def run_query(self, query: LIFQuery, first_run: bool) -> List[LIFRecord] |
lif_fragment_paths_not_found,
lif_query_plan,
lif_query_planner_job.job_id,
client,
)
query_status_response = LIFQueryStatusResponse(query_id=lif_query_planner_job.job_id, status="PENDING")
return query_status_response
Expand Down Expand Up @@ -304,7 +317,9 @@ async def run_post_orchestration_results(self, results: OrchestratorJobResults)

try:
logger.info(
statistics.format_event(statistics.build_query_completed_event(results, lif_fragment_paths))
statistics.format_event(
statistics.build_query_completed_event(results, lif_fragment_paths, job.client)
)
)
except Exception:
logger.exception("Failed to emit query statistics")
Expand Down Expand Up @@ -408,6 +423,7 @@ class LIFQueryPlannerJob(BaseModel):
status (str): Status of the job (e.g., 'pending', 'running', 'completed').
created_timestamp (str): Timestamp of when the job was created.
updated_timestamp (str): Timestamp of when the job was last updated.
client (str): The caller that submitted the query, for the completion statistics event.
"""

job_id: str = Field(..., description="Unique identifier for the job")
Expand All @@ -423,6 +439,8 @@ class LIFQueryPlannerJob(BaseModel):
description="Timestamp of when the job was last updated",
default_factory=lambda: datetime.now(timezone.utc).isoformat(),
)
# The orchestrator's results callback carries no caller, so the job remembers it (#1272).
client: str = Field(statistics.CLIENT_UNKNOWN, description="The caller that submitted the query")


# -------------------------------------------------------------------------
Expand Down
Loading
Loading