Skip to content
20 changes: 20 additions & 0 deletions deployments/CI_DEPLOYMENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -43,11 +43,31 @@ It then handles deployments for:

- the database, assuming an existing Azure PostgreSQL flexible server
- the API Container App, assuming an existing Container Apps environment
(see "Log routing" below)
- the entity-linkage Container App in the same environment
- the stitch-llm Container App in the same environment
- the ETL Container App (`etl`) in the same environment, on non-`development`
lanes only (see below)

### Log routing

Container App logs (structured JSON on stdout) are forwarded to a Log Analytics
workspace by the **Container Apps environment** (`appLogsConfiguration`), not by
this pipeline or the app — so every app in an environment, including per-PR
preview apps, shares one workspace. Current topology: `development` and
`staging` lanes both feed the non-prod workspace **`stitch-staging`**
(`STITCH-DEV-RG`); `production` is isolated in its own workspace.

To change where a lane's logs go, edit that lane's environment (named by its
`AZURE_CONTAINER_APP_ENVIRONMENT` variable) — Portal: **Settings → Logging**, or
CLI: `az containerapp env update --logs-destination log-analytics
--logs-workspace-id <customerId> --logs-workspace-key <key>`. Keep "Parse JSON
logs into columns" **off** on every environment sharing a workspace, so the
table schema stays uniform (the KQL in [`PERFORMANCE.md`](./PERFORMANCE.md) and
`tools/analyze_logs.py` assume the raw-`Log_s` form). This binding lives outside
the repo, so recreating an environment reverts it to `None` (a greyed-out Logs
blade) — reconfigure it when that happens.

### ETL pipelines (temporary POC wiring)

The single `etl` Container App is deployed from a pre-built image published by
Expand Down
13 changes: 12 additions & 1 deletion deployments/PERFORMANCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,9 @@ the cloud, and readable straight from the terminal locally — so you can find
slow/frequent queries from real data instead of guessing.

This doc covers the basic loop: **enable capture → drive traffic → analyze**.
For *where* deployed logs are routed (which Log Analytics workspace each lane
feeds, and how to change it), see "Log routing" in
[`CI_DEPLOYMENTS.md`](./CI_DEPLOYMENTS.md).

> The instrumentation lives in the app code
> ([`deployments/api/src/stitch/api/observability/`](api/src/stitch/api/observability/)),
Expand All @@ -23,7 +26,7 @@ Two structured log streams, distinguished by the `logger` field:
| Logger | Emitted | Key fields |
|---|---|---|
| `stitch.observability.request` | once per HTTP request (always) | `route`, `method`, `status_code`, `duration_ms`, `db_query_count`, `db_time_ms`, `request_id` |
| `stitch.api.observability.query` | once per query above the slow threshold | `statement` (parameterized SQL, **no bound values**), `duration_ms`, `rowcount`, `route`, `request_id` |
| `stitch.api.observability.query` | once per query above the slow threshold | `statement` (parameterized SQL, **no bound values**), `duration_ms`, `rowcount`, `route`, `request_id`, `query_name` (when the query runs in a labeled scope) |

> The request summary is emitted by the shared `stitch.observability`
> middleware, so it logs under `stitch.observability.request` (the API's
Expand All @@ -34,6 +37,14 @@ Two structured log streams, distinguished by the `logger` field:
`db_query_count` on a request is the N+1 detector; the `query` stream tells you
*which* statement is expensive.

`query_name` is a stable label (e.g. `resources.list_ids`, `resources.count`,
`resources.filter_options`) attached to the queries a request handler runs, so you
can pick out a specific query without matching on SQL text — useful when two
statements share a near-identical prefix (both list queries open with
`WITH resource_universe AS …`). It is present only for queries executed inside a
labeled scope; unlabeled queries (e.g. ORM-internal statements outside a handler)
omit the field entirely.

---

## Step 1 — Enable capture (the knobs)
Expand Down
82 changes: 47 additions & 35 deletions deployments/api/src/stitch/api/db/merge_candidate_actions.py
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
MergeCandidateStatus,
MergeCandidateView,
)
from stitch.api.observability.context import named_query
from stitch.ogsi.model import OGFieldSource
from stitch.ogsi.model.og_field import OilGasFieldBase
from stitch.ogsi.model.types import OGSISrcKey
Expand Down Expand Up @@ -204,7 +205,8 @@ async def list_merge_candidates(session: AsyncSession) -> list[MergeCandidateVie
.options(selectinload(MergeCandidateModel.items))
.order_by(MergeCandidateModel.created.desc())
)
candidates = (await session.scalars(stmt)).all()
with named_query("merge_candidates.list"):
candidates = (await session.scalars(stmt)).all()
return [_candidate_to_view(candidate) for candidate in candidates]


Expand All @@ -213,7 +215,8 @@ async def get_merge_candidate(
candidate_id: int,
licensed_sources: Collection[OGSISrcKey] | None = None,
) -> MergeCandidateDetailView:
candidate = await _load_candidate_model(session, candidate_id)
with named_query("merge_candidates.detail.load"):
candidate = await _load_candidate_model(session, candidate_id)

resource_ids = [
item.resource_id for item in sorted(candidate.items, key=lambda i: i.position)
Expand All @@ -228,14 +231,16 @@ async def get_merge_candidate(
# never delete). So a null-shell view always means "emptied by a merge",
# never "missing"; no existence check is needed here. Revisit if a resource
# hard-delete path is ever added.
by_id = await coalesce_resources_with_sources(
session, resource_ids, licensed_sources
)
with named_query("merge_candidates.detail.coalesce"):
by_id = await coalesce_resources_with_sources(
session, resource_ids, licensed_sources
)

# `status` compares the resources' coalesced values; `values` lists every
# contributing source tagged with the resource it's attached to, ranked by
# the default source order (winner-first).
default_priority = await _default_source_priority(session)
with named_query("merge_candidates.detail.default_priority"):
default_priority = await _default_source_priority(session)
fallback_priority = max(default_priority.values(), default=0) + 1
sources_with_priority = [
(rid, source, default_priority.get(source.source, fallback_priority))
Expand All @@ -254,14 +259,16 @@ async def create_merge_candidate(
request: MergeCandidateCreateRequest,
) -> MergeCandidateView:
resource_ids = _normalize_resource_ids(request.resource_ids)
await _load_mergeable_resources(session, resource_ids)
with named_query("merge_candidates.create.load_resources"):
await _load_mergeable_resources(session, resource_ids)

fingerprint = _fingerprint(resource_ids)
existing = await session.scalar(
select(MergeCandidateModel)
.options(selectinload(MergeCandidateModel.items))
.where(MergeCandidateModel.fingerprint == fingerprint)
)
with named_query("merge_candidates.create.check_existing"):
existing = await session.scalar(
select(MergeCandidateModel)
.options(selectinload(MergeCandidateModel.items))
.where(MergeCandidateModel.fingerprint == fingerprint)
)
if existing is not None:
if existing.status == MergeCandidateStatus.PENDING:
raise InvalidActionError(
Expand All @@ -275,22 +282,23 @@ async def create_merge_candidate(
f"An approved merge candidate already exists for resources {resource_ids}."
)

candidate = MergeCandidateModel.create(created_by=user, fingerprint=fingerprint)
session.add(candidate)
await session.flush()
with named_query("merge_candidates.create.persist"):
candidate = MergeCandidateModel.create(created_by=user, fingerprint=fingerprint)
session.add(candidate)
await session.flush()

session.add_all(
[
MergeCandidateItemModel(
merge_candidate_id=candidate.id,
resource_id=resource_id,
position=position,
)
for position, resource_id in enumerate(resource_ids)
]
)
await session.flush()
await session.refresh(candidate, ["items"])
session.add_all(
[
MergeCandidateItemModel(
merge_candidate_id=candidate.id,
resource_id=resource_id,
position=position,
)
for position, resource_id in enumerate(resource_ids)
]
)
await session.flush()
await session.refresh(candidate, ["items"])
return _candidate_to_view(candidate)


Expand All @@ -300,7 +308,8 @@ async def approve_merge_candidate(
candidate_id: int,
request: MergeCandidateReviewRequest | None = None,
) -> MergeCandidateView:
candidate = await _load_candidate_model(session, candidate_id)
with named_query("merge_candidates.approve.load"):
candidate = await _load_candidate_model(session, candidate_id)
if candidate.status != MergeCandidateStatus.PENDING:
raise InvalidActionError(
f"Merge candidate {candidate_id} is not pending; current status={candidate.status}."
Expand All @@ -309,7 +318,8 @@ async def approve_merge_candidate(
resource_ids = [
item.resource_id for item in sorted(candidate.items, key=lambda i: i.position)
]
await _load_mergeable_resources(session, resource_ids)
with named_query("merge_candidates.approve.load_resources"):
await _load_mergeable_resources(session, resource_ids)
merged_resource = await apply_resource_merge(
session=session,
user=user,
Expand All @@ -322,9 +332,9 @@ async def approve_merge_candidate(
candidate.reviewed_by_id = user.id
candidate.last_updated_by_id = user.id
candidate.merged_resource_id = merged_resource.id
await session.flush()

candidate = await _load_candidate_model(session, candidate_id)
with named_query("merge_candidates.approve.persist"):
await session.flush()
candidate = await _load_candidate_model(session, candidate_id)
return _candidate_to_view(candidate)


Expand All @@ -334,7 +344,8 @@ async def deny_merge_candidate(
candidate_id: int,
request: MergeCandidateReviewRequest | None = None,
) -> MergeCandidateView:
candidate = await _load_candidate_model(session, candidate_id)
with named_query("merge_candidates.deny.load"):
candidate = await _load_candidate_model(session, candidate_id)
if candidate.status != MergeCandidateStatus.PENDING:
raise InvalidActionError(
f"Merge candidate {candidate_id} is not pending; current status={candidate.status}."
Expand All @@ -345,6 +356,7 @@ async def deny_merge_candidate(
candidate.reviewed_at = datetime.now(timezone.utc)
candidate.reviewed_by_id = user.id
candidate.last_updated_by_id = user.id
await session.flush()
candidate = await _load_candidate_model(session, candidate_id)
with named_query("merge_candidates.deny.persist"):
await session.flush()
candidate = await _load_candidate_model(session, candidate_id)
return _candidate_to_view(candidate)
Loading
Loading