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
8 changes: 8 additions & 0 deletions .env.exemple
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,14 @@ STORAGE_ACCESS_KEY=minioadmin
STORAGE_SECRET_KEY=minioadmin
STORAGE_REGION=us-east-1

# ── Realtime (Redis pub/sub) ──
# Fire-and-forget SSE hints consumed by crm-backend's realtime hub. Optional
# everywhere, including production: an empty URL makes core.realtime.emit() a
# silent no-op, so a broker outage degrades the UI to polling and never affects
# a business write. db 1 on purpose — db 0 is crm-backend's cache.
REALTIME_ENABLED=false
REALTIME_REDIS_URL=redis://:changeme_redis_password@localhost:6379/1

# ── CORS ──
# Comma-separated list of allowed origins. "*" is allowed only in local.
ALLOW_ORIGIN=*
Expand Down
6 changes: 5 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,11 @@ pyrightconfig.json
[Ll]ib
[Ll]ib64
[Ll]ocal
[Ss]cripts
# [Ss]cripts — REMOVED. This unanchored toptal venv entry matched this
# project's own top-level scripts/ directory, which silently kept every
# file under scripts/sql/migrations/ out of the repository. The virtualenv
# case it was meant to cover is already handled by the `.venv` and `venv/`
# rules above, so the entry was pure downside.
pyvenv.cfg
pip-selfcheck.json

Expand Down
43 changes: 43 additions & 0 deletions api/generation/mappers.py
Original file line number Diff line number Diff line change
@@ -1,10 +1,15 @@
from api.generation.schemas import (
AllocationKeyGenerated,
ConsumerGenerated,
CrmDataPreview,
Generation,
IncompleteMeter,
IterationGenerated,
PartialAllocationKeyGenerated,
PreviewBlocker,
)
from core.i18n import translate
from shared.crm_preflight import Preflight
from shared.models.crm_models import AllocationKeyModel, ConsumerModel, IterationModel
from shared.models.local_models import (
AllocationKeyGeneratedModel,
Expand Down Expand Up @@ -91,3 +96,41 @@ def to_allocation_key_crm(
id_community=allocation_key.id_community,
iterations=[to_iteration_crm(i) for i in allocation_key.iterations],
)


def to_crm_data_preview(preflight: Preflight, locale: str) -> CrmDataPreview:
"""Render a pre-flight verdict for the manager's screen.

Blocker messages are localised here rather than in the service so that
translation stays at the API edge — the worker reaches the same verdict via
``Preflight`` and needs no locale at all.
"""
summary = preflight.summary
return CrmDataPreview(
can_generate=preflight.ok,
# The participant count, not the raw meter count: injection-only sites
# contribute production but never receive a share.
meter_count=len(preflight.consumer_eans),
reading_count=summary.total_rows,
first_timestamp=summary.first_timestamp,
last_timestamp=summary.last_timestamp,
total_consumption_kwh=summary.total_consumption_kwh,
total_injection_kwh=summary.total_injection_kwh,
incomplete_meters=[
IncompleteMeter(
ean=e.ean,
readings=e.distinct_ts,
expected=summary.grid_size,
missing=summary.grid_size - e.distinct_ts,
)
for e in summary.incomplete
],
blockers=[
PreviewBlocker(
error_code=b.error.code,
message=translate(b.error.key, locale=locale),
detail=b.detail,
)
for b in preflight.blockers
],
)
51 changes: 51 additions & 0 deletions api/generation/routes.py
Original file line number Diff line number Diff line change
@@ -1,12 +1,16 @@
import json
from datetime import date
from typing import Annotated, Any

from fastapi import APIRouter, Body, Depends, File, Form, Query, UploadFile
from sqlalchemy.ext.asyncio import AsyncSession

from algorithms.registry import AlgorithmMetadata, registry
from api.generation.mappers import to_crm_data_preview
from api.generation.schemas import (
AllocationKeyGenerated,
CrmDataPreview,
GenerateFromCrmRequest,
GenerateRequest,
GenerateResponse,
Generation,
Expand Down Expand Up @@ -95,6 +99,34 @@ async def get_algorithm_inputs(algorithm_name: str):
return ApiResponse[LocalizedAlgorithmMetadata](data=data)


# GET (/crm-data-preview) : What the CRM holds for a sharing operation + period
#
# Declared BEFORE `GET /{id}`: FastAPI matches in declaration order, so putting
# this after the integer catch-all would make "/crm-data-preview" try to parse
# as an id and 422.
@generation_routes.get("/crm-data-preview", response_model=ApiResponse[CrmDataPreview])
@with_default_error(default_error=errors.generation.GET_CRM_PREVIEW)
async def get_crm_data_preview(
local_session: Annotated[AsyncSession, Depends(get_local_session)],
crm_session: Annotated[AsyncSession, Depends(get_crm_session)],
id_sharing_operation: Annotated[int, Query(description="CRM sharing operation id.")],
period_start: Annotated[date, Query(description="First day of the period (inclusive).")],
period_end: Annotated[date, Query(description="Last day of the period (inclusive).")],
):
internal_community_id = current_internal_community_id.get()
if internal_community_id is None:
raise ErrorException(error=errors.auth.UNAUTHORIZED, status_code=401)
service = GenerationService(local_session, crm_session)
preflight = await service.preview_crm_data(
id_sharing_operation=id_sharing_operation,
period_start=period_start,
period_end=period_end,
community_id=internal_community_id,
)
locale = current_locale.get().split("_")[0]
return ApiResponse[CrmDataPreview](data=to_crm_data_preview(preflight, locale))


# GET (/key/{id}) : Key generated
@generation_routes.get("/key/{id_key}", response_model=ApiResponse[AllocationKeyGenerated])
@with_default_error(default_error=errors.generation.GET_ALLOCATION_KEY)
Expand Down Expand Up @@ -172,6 +204,25 @@ async def start_generation(
return ApiResponse[GenerateResponse](data=data)


# POST (/from-crm) Generate from CRM meter data
#
# Plain JSON, unlike POST / — with no file part nothing forces multipart. It
# also correctly keeps the default 2 MB body cap rather than the upload one.
@generation_routes.post("/from-crm", response_model=ApiResponse[GenerateResponse])
@with_default_error(default_error=errors.generation.START_GENERATION)
async def start_generation_from_crm(
body: Annotated[GenerateFromCrmRequest, Body()],
local_session: Annotated[AsyncSession, Depends(get_local_session)],
crm_session: Annotated[AsyncSession, Depends(get_crm_session)],
):
internal_community_id = current_internal_community_id.get()
if internal_community_id is None:
raise ErrorException(error=errors.auth.UNAUTHORIZED, status_code=401)
service = GenerationService(local_session, crm_session)
data = await service.start_generation_from_crm(body, internal_community_id)
return ApiResponse[GenerateResponse](data=data)


# POST (/save): Save a key
@generation_routes.post("/save", response_model=ApiResponse[str])
@with_default_error(default_error=errors.generation.SAVE_KEY)
Expand Down
64 changes: 64 additions & 0 deletions api/generation/schemas.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import datetime
from typing import Any

from pydantic import BaseModel, Field
Expand Down Expand Up @@ -102,3 +103,66 @@ class LocalizedAlgorithmMetadata(BaseModel):
input_schema: dict[str, Any]
tags: list[str]
timeout_seconds: int | None


# ---------------------------------------------------------------------------
# CRM-sourced generation (source = DataSource.CRM)
# ---------------------------------------------------------------------------


class GenerateFromCrmRequest(BaseModel):
"""Body of ``POST /from-crm``.

Unlike ``GenerateRequest`` this *is* the FastAPI body model: with no file
part there is nothing forcing multipart, so the request is plain JSON.

``injection_name`` is deliberately absent — the production profile is summed
from the meters themselves, which is the whole reason this path is simpler
for the user than uploading a file.
"""

name: str = Field(..., min_length=1, description="User-facing label for the generation.")
algorithm_name: str = Field(..., description="Algorithm registry key, e.g. 'olagsa'.")
inputs: dict[str, Any] = Field(
...,
description=(
"Algorithm-specific input parameters; validated against the algorithm's input schema."
),
)
id_sharing_operation: int = Field(..., description="CRM sharing operation to read meters from.")
period_start: datetime.date = Field(..., description="First day of the period (inclusive).")
period_end: datetime.date = Field(..., description="Last day of the period (inclusive).")


class IncompleteMeter(BaseModel):
"""A meter missing part of the period. Zero-filled, not fatal."""

ean: str
readings: int = Field(..., description="Distinct timestamps this meter actually has.")
expected: int = Field(..., description="Distinct timestamps across the whole operation.")
missing: int = Field(..., description="expected - readings.")


class PreviewBlocker(BaseModel):
"""A reason the period cannot be used, already localised."""

error_code: int = Field(..., description="Matches the error_code of the eventual 4xx.")
message: str = Field(..., description="Localised, manager-facing explanation.")
detail: str = Field(..., description="Which meters/values triggered it.")


class CrmDataPreview(BaseModel):
"""What ``GET /crm-data-preview`` shows before the manager commits to a run.

``can_generate`` is the single flag the UI binds its submit button to.
"""

can_generate: bool
meter_count: int = Field(..., description="Meters that drew energy and will be participants.")
reading_count: int
first_timestamp: datetime.datetime | None
last_timestamp: datetime.datetime | None
total_consumption_kwh: float
total_injection_kwh: float
incomplete_meters: list[IncompleteMeter]
blockers: list[PreviewBlocker]
129 changes: 128 additions & 1 deletion api/generation/service.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import datetime
import logging
from uuid import uuid4

Expand All @@ -14,6 +15,7 @@
from api.generation.repository import GenerationRepository
from api.generation.schemas import (
AllocationKeyGenerated,
GenerateFromCrmRequest,
GenerateRequest,
GenerateResponse,
Generation,
Expand All @@ -29,7 +31,10 @@
from core.middleware.request_limits import UPLOAD_MAX_BODY_BYTES
from core.queue.helper import Event, send_event
from core.queue.init import get_jetstream
from shared.const import GenerationStatus
from shared import crm_preflight
from shared.const import DataSource, GenerationStatus
from shared.crm_meter_repository import CrmMeterRepository
from shared.crm_preflight import Preflight
from shared.crm_repository import CRMRepository
from shared.custom_errors import errors
from shared.models.local_models import GenerationModel
Expand Down Expand Up @@ -255,6 +260,128 @@ async def _mark_failed_to_queue(generation_id: int, reason: str) -> None:
)
await crm_session.commit()

# ------------------------------------------------------------------
# CRM-sourced generation
# ------------------------------------------------------------------

async def preview_crm_data(
self,
*,
id_sharing_operation: int,
period_start: datetime.date,
period_end: datetime.date,
community_id: int,
) -> Preflight:
"""Aggregate the requested period and classify it, without running anything.

Also the pre-flight for ``start_generation_from_crm`` — one code path, so
the answer the manager saw and the answer that gates the run cannot drift.
"""
if period_start > period_end:
raise ErrorException(error=errors.generation.INVALID_PERIOD, status_code=422)

crm_meters = CrmMeterRepository(self.crm_session)
# Explicit tenant check: without it a foreign operation id is
# indistinguishable from an empty period, which is a confusing 422 for a
# legitimate user and a soft information leak for everyone else.
if not await crm_meters.sharing_operation_exists(
id_community=community_id, id_sharing_operation=id_sharing_operation
):
raise ErrorException(
error=errors.generation.SHARING_OPERATION_NOT_FOUND, status_code=404
)

summary = await crm_meters.summarize(
id_community=community_id,
id_sharing_operation=id_sharing_operation,
period_start=period_start,
period_end=period_end,
)
return crm_preflight.evaluate(summary)

async def start_generation_from_crm(
self, req: GenerateFromCrmRequest, community_id: int
) -> GenerateResponse:
"""Queue a generation that reads its input from the CRM.

Same ordering as the file path minus the upload: validate, commit, then
publish. There is no object to roll back, so the ``_best_effort_delete``
branches have no counterpart here.
"""
if req.algorithm_name not in registry:
raise ErrorException(error=errors.generation.ALGORITHM_NOT_FOUND, status_code=404)
meta = registry.metadata(req.algorithm_name)

try:
validated_inputs = meta.input_schema.model_validate(req.inputs)
except ValidationError as e:
logger.info("Invalid inputs for algorithm '%s': %s", req.algorithm_name, e)
raise ErrorException(
error=errors.generation.INVALID_ALGORITHM_INPUTS, status_code=422
) from e

# Re-run the pre-flight rather than trusting whatever the client saw:
# the preview may be minutes old, and an import can have landed since.
preflight = await self.preview_crm_data(
id_sharing_operation=req.id_sharing_operation,
period_start=req.period_start,
period_end=req.period_end,
community_id=community_id,
)
if preflight.blockers:
first = preflight.blockers[0]
logger.info(
"CRM generation refused for community %d op %d: %s",
community_id,
req.id_sharing_operation,
first.detail,
)
raise ErrorException(error=first.error, status_code=422)

model = GenerationModel(
name=req.name,
id_community=community_id,
source=DataSource.CRM,
id_sharing_operation=req.id_sharing_operation,
period_start=req.period_start,
period_end=req.period_end,
algorithm_name=meta.name,
algorithm_version=meta.version,
inputs=validated_inputs.model_dump(mode="json"),
status=GenerationStatus.PENDING,
data_warnings=preflight.warnings,
)
await self.repository.create_generation(model)
await self.local_session.commit()
generation_id = model.id
app_metrics.generations_created.add(1, {"algorithm": meta.name})
await self.audit_log_service.log(
AuditLogInput(
action=AuditActions.GENERATION_CREATED,
entity_type="generation",
entity_id=str(generation_id),
payload={
"name": req.name,
"algorithm_name": meta.name,
"algorithm_version": meta.version,
"source": DataSource.CRM.name,
"id_sharing_operation": req.id_sharing_operation,
"period_start": req.period_start.isoformat(),
"period_end": req.period_end.isoformat(),
},
)
)

event = Event(type="generation.requested", data={"generation_id": generation_id})
try:
await send_event(get_jetstream(), meta.queue, event)
except Exception as exc:
logger.exception("Failed to publish generation %d to %s", generation_id, meta.queue)
await self._mark_failed_to_queue(generation_id, str(exc))
raise ErrorException(error=errors.generation.START_GENERATION, status_code=500) from exc

return GenerateResponse(id=generation_id, status=GenerationStatus.PENDING)

async def save_key(self, saved_key: SaveKey):
# Retrieve it in this database
key = await self.repository.get_allocation_key(saved_key.id_key)
Expand Down
Loading
Loading