pr0nstar

XIII

......@@ -6,7 +6,7 @@ from starlette_caches.utils import cache as scache
from aiocache import Cache
REDIS_CONN = open('./docs/postgres').read().strip()
REDIS_CONN = open('./docs/redis').read().strip()
REDIS_CONN = REDIS_CONN.split(':')
# cache = Cache(
# Cache.REDIS,
......
import logging
import threading
import time
import traceback
WINDOW_SECONDS = 1
MAX_KEYS = 1024
class DuplicateErrorFilter(logging.Filter):
def __init__(self, window_seconds=WINDOW_SECONDS, max_keys=MAX_KEYS):
super().__init__()
self.window_seconds = window_seconds
self.max_keys = max_keys
self._lock = threading.Lock()
self._entries = {}
def filter(self, record):
if record.levelno < logging.ERROR:
return True
now = time.monotonic()
key = self.build_key(record)
with self._lock:
self.prune(now)
opened_until = self._entries.get(key, 0.)
if opened_until > now:
return False
self._entries[key] = now + self.window_seconds
return True
def build_key(self, record):
exc_type = ""
frame = ("", 0, "")
exc_info = record.exc_info
if isinstance(exc_info, tuple) and len(exc_info) == 3 and exc_info[0] is not None:
exc_type = exc_info[0].__name__
extracted = traceback.extract_tb(exc_info[2])
if extracted:
last = extracted[-1]
frame = (last.filename, last.lineno, last.name)
return (record.name, exc_type, frame, record.msg)
def prune(self, now):
expired = [
key
for key, opened_until in self._entries.items()
if opened_until <= now
]
for key in expired:
self._entries.pop(key, None)
if len(self._entries) <= self.max_keys:
return
oldest = sorted(
self._entries.items(),
key=lambda item: item[1],
)[:len(self._entries) - self.max_keys]
for key, _ in oldest:
self._entries.pop(key, None)
......@@ -15,7 +15,7 @@ from starlette.middleware import Middleware
from starlette_caches.middleware import CacheMiddleware
from starlette_caches.rules import Rule
from presupuestov1 import cache
from presupuestov1 import cache, metrics
from presupuestov1.routers import (
project,
entity,
......@@ -23,17 +23,10 @@ from presupuestov1.routers import (
classifier_income,
daily,
search,
private_metrics,
)
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s %(message)s',
filename='log.log',
)
logging.getLogger('api_proy2').setLevel(logging.DEBUG)
app = FastAPI(
title='MEFP presupuestosv1 API',
middleware=[
......@@ -47,13 +40,16 @@ app = FastAPI(
CacheMiddleware,
cache=cache.cache,
rules=[
Rule(match=re.compile(r"^/api/_metrics/.*"), ttl=0),
Rule(match=re.compile(r"^/api/search/.*"), ttl=0),
Rule(match=re.compile(r"^/api/programa_proyecto/.*"), ttl=0),
Rule(match=re.compile(r'^/api/.+'), ttl=86400, status=200),
],
),
],
root_path='/api'
)
app.include_router(entity.router_entidad)
app.include_router(entity.router_ubigeo)
app.include_router(entity.router_da)
......@@ -65,6 +61,9 @@ app.include_router(classifier_income.router_organismo)
app.include_router(project.router)
app.include_router(daily.router)
app.include_router(search.router)
app.include_router(private_metrics.router)
# metrics.install(app)
# @app.on_event('startup')
......
import logging
import time
from datetime import datetime, timedelta, timezone
import redis.asyncio as redis
from starlette.routing import Match
PREFIX = "metrics:v1"
TTL_SECONDS = 7 * 24 * 60 * 60
STREAM_MAXLEN = 10000
HTTP_BUCKETS_MS = (100, 300, 600, 1000, 2000, 3000)
PRIVATE_API_PREFIXES = ("/_metrics", "/api/_metrics")
redis_client = None
logger = logging.getLogger("api_proy2")
###############################################################################
# middleware
###############################################################################
class MetricsMiddleware:
def __init__(self, app, routes):
self.app = app
self.routes = routes
async def __call__(self, scope, receive, send):
if scope["type"] != "http":
await self.app(scope, receive, send)
return
path = scope.get("path", "")
if path.startswith(PRIVATE_API_PREFIXES):
await self.app(scope, receive, send)
return
started_at = time.perf_counter()
status_code = 500
response_headers = []
async def send_wrapper(message):
nonlocal status_code, response_headers
if message["type"] == "http.response.start":
status_code = message["status"]
response_headers = message.get("headers", [])
await send(message)
try:
await self.app(scope, receive, send_wrapper)
finally:
try:
await record_http(
method=scope.get("method", "GET"),
route=find_route_template(self.routes, scope),
path=scope.get("path", ""),
client_ip=get_client_ip(scope.get("headers", [])),
status_code=status_code,
duration_ms=(time.perf_counter() - started_at) * 1000,
cache_status=get_header_value(response_headers, b"x-cache"),
)
except Exception:
logger.exception("metrics write failed")
def install(app):
app.add_middleware(MetricsMiddleware, routes=app.router.routes)
###############################################################################
# redis lifecycle
###############################################################################
async def get_redis():
global redis_client
if redis_client is not None:
return redis_client
with open("./docs/redis") as f:
redis_url = f.read().strip()
if "://" not in redis_url:
redis_url = f"redis://{redis_url}"
redis_client = redis.from_url(redis_url, decode_responses=True)
return redis_client
async def close_redis():
global redis_client
if redis_client is None:
return
await redis_client.aclose()
redis_client = None
###############################################################################
# write api
###############################################################################
async def record_http(
method,
route,
path,
client_ip,
status_code,
duration_ms,
cache_status=None
):
redis_conn = await get_redis()
minute = get_current_minute()
http_minute_key = get_http_minute_key(minute)
ip_minute_key = get_ip_minute_key(minute)
route_label = join_label(method, route, int(status_code))
route_key = join_label(method, route)
pipe = redis_conn.pipeline(transaction=False)
pipe.hincrby(http_minute_key, f"count|{route_label}", 1)
pipe.hincrbyfloat(http_minute_key, f"sum_ms|{route_label}", float(duration_ms))
pipe.hincrby(http_minute_key, f"status|{int(status_code)}", 1)
if cache_status in {"hit", "miss"}:
pipe.hincrby(http_minute_key, f"cache_{cache_status}|{route_key}", 1)
for bucket_ms in HTTP_BUCKETS_MS:
if duration_ms <= bucket_ms:
pipe.hincrby(http_minute_key, f"le_{bucket_ms}|{route_label}", 1)
pipe.hincrby(ip_minute_key, normalize_label_part(client_ip), 1)
pipe.expire(http_minute_key, TTL_SECONDS)
pipe.expire(ip_minute_key, TTL_SECONDS)
# pipe.xadd(
# get_requests_stream_key(),
# {
# "ts": get_now_utc().isoformat(),
# "method": str(method),
# "route": str(route),
# "path": str(path),
# "client_ip": str(client_ip or ""),
# "status_code": str(int(status_code)),
# "duration_ms": str(round(float(duration_ms), 3)),
# "cache_status": str(cache_status or ""),
# },
# maxlen=STREAM_MAXLEN,
# approximate=True,
# )
await pipe.execute()
###############################################################################
# read api
###############################################################################
async def read_summary(minutes=60, include_ts=False):
if include_ts:
return await read_summary_series(minutes)
return await read_summary_window(minutes)
async def read_routes(minutes=60, include_ts=False):
if include_ts:
return await read_route_series(minutes)
return await read_route_window(minutes)
async def read_ips(minutes=60, include_ts=False):
if include_ts:
return await read_ip_series(minutes)
return await read_ip_window(minutes)
# async def read_requests(limit=100):
# redis_conn = await get_redis()
# rows = await redis_conn.xrevrange(get_requests_stream_key(), count=limit)
# return [{"id": row_id, "fields": payload} for row_id, payload in rows]
async def read_http(minutes):
return merge_minute_maps(await load_http_minutes(minutes))
###############################################################################
# summary readers
###############################################################################
async def read_summary_window(minutes):
http_minutes = await load_http_minutes(minutes)
http_fields = merge_minute_maps(http_minutes)
total_count = get_total_count(http_fields)
total_sum_ms = get_total_sum_ms(http_fields)
return {
"window_minutes": minutes,
"generated_at": get_now_utc().isoformat(),
"http": {
"count": total_count,
"sum_ms": total_sum_ms,
"avg_ms": round(total_sum_ms / total_count, 3) if total_count else 0.0,
"status": get_status_totals(http_fields),
"cache": get_cache_totals(http_fields),
"buckets_ms": get_bucket_totals(http_fields),
},
}
async def read_summary_series(minutes):
http_minutes = await load_http_minutes(minutes)
http_fields = merge_minute_maps(http_minutes)
total_count = get_total_count(http_fields)
total_sum_ms = get_total_sum_ms(http_fields)
series = []
for minute, minute_fields in http_minutes:
minute_count = get_total_count(minute_fields)
minute_sum_ms = get_total_sum_ms(minute_fields)
series.append({
"minute": minute,
"http": {
"count": minute_count,
"sum_ms": minute_sum_ms,
"avg_ms": round(minute_sum_ms / minute_count, 3) if minute_count else 0.0,
"status": get_status_totals(minute_fields),
"cache": get_cache_totals(minute_fields),
"buckets_ms": get_bucket_totals(minute_fields),
},
})
return {
"window_minutes": minutes,
"generated_at": get_now_utc().isoformat(),
"http": {
"count": total_count,
"sum_ms": total_sum_ms,
"avg_ms": round(total_sum_ms / total_count, 3) if total_count else 0.0,
"status": get_status_totals(http_fields),
"cache": get_cache_totals(http_fields),
"buckets_ms": get_bucket_totals(http_fields),
},
"series": series,
}
###############################################################################
# route readers
###############################################################################
async def read_route_window(minutes):
return build_route_rows(await read_http(minutes))
async def read_route_series(minutes):
rows = []
for minute, minute_fields in await load_http_minutes(minutes):
rows.extend(build_route_rows(minute_fields, minute=minute))
rows.sort(key=lambda row: (
row["minute"], -row["count"], -row["sum_ms"], row["route"]
))
return rows
###############################################################################
# ip readers
###############################################################################
async def read_ip_window(minutes):
ip_fields = merge_minute_maps(await load_ip_minutes(minutes))
rows = []
for ip, count in ip_fields.items():
rows.append({
"ip": ip,
"count": int(count),
})
rows.sort(key=lambda row: (-row["count"], row["ip"]))
return rows
async def read_ip_series(minutes):
rows = []
for minute, minute_fields in await load_ip_minutes(minutes):
for ip, count in minute_fields.items():
rows.append({
"minute": minute,
"ip": ip,
"count": int(count),
})
rows.sort(key=lambda row: (row["minute"], -row["count"], row["ip"]))
return rows
###############################################################################
# minute loaders
###############################################################################
async def load_http_minutes(minutes):
minute_keys = get_minute_buckets(minutes)
if not minute_keys:
return []
redis_conn = await get_redis()
pipe = redis_conn.pipeline(transaction=False)
for minute in minute_keys:
pipe.hgetall(get_http_minute_key(minute))
rows = await pipe.execute()
return build_minute_maps(minute_keys, rows)
async def load_ip_minutes(minutes):
minute_keys = get_minute_buckets(minutes)
if not minute_keys:
return []
redis_conn = await get_redis()
pipe = redis_conn.pipeline(transaction=False)
for minute in minute_keys:
pipe.hgetall(get_ip_minute_key(minute))
rows = await pipe.execute()
return build_minute_maps(minute_keys, rows)
###############################################################################
# route matching
###############################################################################
def find_route_template(routes, scope):
route_template = match_route_template(routes, scope)
if route_template is not None:
return route_template
root_path = scope.get("root_path", "")
path = scope.get("path", "")
if root_path and path.startswith(root_path):
stripped_scope = dict(scope)
stripped_scope["path"] = path[len(root_path):] or "/"
route_template = match_route_template(routes, stripped_scope)
if route_template is not None:
return route_template
return path
def match_route_template(routes, scope):
for route in routes:
match, _ = route.matches(scope)
if match == Match.FULL:
return getattr(
route, "path_format", getattr(route, "path", scope.get("path", ""))
)
return None
###############################################################################
# request parsing
###############################################################################
def get_header_value(headers, key):
for header_key, header_value in headers:
if header_key.lower() == key:
return header_value.decode()
return None
def get_client_ip(headers):
forwarded_for = get_header_value(headers, b"x-forwarded-for")
if forwarded_for:
return forwarded_for.split(",", 1)[0].strip()
real_ip = get_header_value(headers, b"x-real-ip")
if real_ip:
return real_ip.strip()
return None
###############################################################################
# row builders
###############################################################################
def build_route_rows(fields, minute=None):
routes = {}
for field, value in fields.items():
parts = field.split("|")
if len(parts) < 4:
continue
metric_name = parts[0]
method = parts[1]
route = parts[2]
status_code = parts[3]
if metric_name not in {"count", "sum_ms"} and not metric_name.startswith("le_"):
continue
route_key = (method, route, status_code)
if route_key not in routes:
routes[route_key] = {
"method": method,
"route": route,
"status_code": int(status_code),
"count": 0,
"sum_ms": 0.0,
"avg_ms": 0.0,
"buckets_ms": {},
}
if minute is not None:
routes[route_key]["minute"] = minute
if metric_name == "count":
routes[route_key]["count"] += int(value)
elif metric_name == "sum_ms":
routes[route_key]["sum_ms"] += float(value)
else:
bucket_ms = metric_name[3:]
routes[route_key]["buckets_ms"][bucket_ms] = (
routes[route_key]["buckets_ms"].get(bucket_ms, 0) + int(value)
)
rows = []
for row in routes.values():
if row["count"]:
row["avg_ms"] = round(row["sum_ms"] / row["count"], 3)
row["sum_ms"] = round(row["sum_ms"], 3)
rows.append(row)
rows.sort(key=lambda row: (-row["count"], -row["sum_ms"], row["route"]))
return rows
def build_minute_maps(minute_keys, rows):
minute_maps = []
for minute, row in zip(minute_keys, rows):
minute_fields = {}
for field, raw_value in row.items():
minute_fields[field] = float(raw_value)
minute_maps.append((minute, minute_fields))
return minute_maps
###############################################################################
# aggregation helpers
###############################################################################
def merge_minute_maps(minute_maps):
merged = {}
for _, minute_fields in minute_maps:
for field, value in minute_fields.items():
merged[field] = merged.get(field, 0.0) + float(value)
return merged
def get_total_count(fields):
count = 0
for field, value in fields.items():
if field.startswith("count|"):
count += int(value)
return count
def get_total_sum_ms(fields):
total = 0.0
for field, value in fields.items():
if field.startswith("sum_ms|"):
total += float(value)
return round(total, 3)
def get_status_totals(fields):
totals = {}
for field, value in fields.items():
if field.startswith("status|"):
status_code = field.split("|", 1)[1]
totals[status_code] = totals.get(status_code, 0) + int(value)
return totals
def get_cache_totals(fields):
totals = {"hit": 0, "miss": 0}
for field, value in fields.items():
if field.startswith("cache_hit|"):
totals["hit"] += int(value)
elif field.startswith("cache_miss|"):
totals["miss"] += int(value)
return totals
def get_bucket_totals(fields):
totals = {}
for field, value in fields.items():
if field.startswith("le_"):
bucket_ms = field.split("|", 1)[0][3:]
totals[bucket_ms] = totals.get(bucket_ms, 0) + int(value)
return totals
###############################################################################
# time and key helpers
###############################################################################
def get_now_utc():
return datetime.now(timezone.utc)
def get_current_minute(now=None):
now = now or get_now_utc()
return now.strftime("%Y%m%d%H%M")
def get_minute_buckets(minutes, now=None):
if minutes <= 0:
return []
now = now or get_now_utc()
base = now.replace(second=0, microsecond=0)
return [
(base - timedelta(minutes=offset)).strftime("%Y%m%d%H%M")
for offset in range(minutes)
]
def get_http_minute_key(minute):
return f"{PREFIX}:http:{minute}"
def get_ip_minute_key(minute):
return f"{PREFIX}:ip:{minute}"
def get_requests_stream_key():
return f"{PREFIX}:requests"
###############################################################################
# label helpers
###############################################################################
def join_label(*parts):
return "|".join(normalize_label_part(part) for part in parts)
def normalize_label_part(value):
if value is None:
return "__none__"
text = str(value).strip()
if not text:
return "__empty__"
return text.replace("|", "_")
......@@ -17,7 +17,6 @@ SessionLocal = sessionmaker(
autocommit=False,
autoflush=False,
bind=engine,
# pool_size=1,
)
Base = declarative_base()
......
......@@ -4,6 +4,7 @@ from . import (
finfun,
objeto,
treemap,
project,
entity,
rubro,
organismo,
......
......@@ -57,6 +57,7 @@ class ActecoEntidad(Base):
Index("idx_acteco_entidad_gestion", "gestion"),
Index("idx_acteco_entidad_acteco", "acteco"),
Index("idx_acteco_entidad_acteco_entidad", "acteco", "entidad"),
Index("idx_acteco_entidad_acteco_gestion", "acteco", "gestion"),
{"schema": "web"},
)
......@@ -81,6 +82,7 @@ class ActecoUbigeo(Base):
Index("idx_acteco_ubigeo_gestion", "gestion"),
Index("idx_acteco_ubigeo_acteco", "acteco"),
Index("idx_acteco_ubigeo_acteco_entidad", "acteco", "ubigeo"),
Index("idx_acteco_ubigeo_acteco_gestion", "acteco", "gestion"),
{"schema": "web"},
)
......
......@@ -19,6 +19,7 @@ class EntidadDistribuciones(Base):
Index("idx_entidad_distribuciones_codigo", "codigo"),
Index("idx_entidad_distribuciones_dimension", "dimension"),
Index("idx_entidad_distribuciones_gestion", "gestion"),
Index("idx_entidad_distribuciones_compos1", "tipo_codigo", "codigo", "dimension", "gestion"),
{"schema": "web"},
)
......@@ -44,6 +45,7 @@ class EntidadResumenes(Base):
Index("idx_entidad_resumenes_tipo_codigo", "tipo_codigo"),
Index("idx_entidad_resumenes_codigo", "codigo"),
Index("idx_entidad_resumenes_gestion", "gestion"),
Index("idx_entidad_resumenes_compos1", "tipo_codigo", "codigo"),
{"schema": "web"},
)
......
......@@ -71,6 +71,7 @@ class FinfunEntidad(Base):
Index("idx_finfun_entidad_gestion", "gestion"),
Index("idx_finfun_entidad_finfun", "finfun"),
Index("idx_finfun_entidad_finfun_entidad", "finfun", "entidad"),
Index("idx_finfun_entidad_finfun_gestion", "finfun", "gestion"),
{"schema": "web"},
)
......@@ -95,6 +96,7 @@ class FinfunUbigeo(Base):
Index("idx_finfun_ubigeo_gestion", "gestion"),
Index("idx_finfun_ubigeo_finfun", "finfun"),
Index("idx_finfun_ubigeo_finfun_entidad", "finfun", "ubigeo"),
Index("idx_finfun_ubigeo_finfun_gestion", "finfun", "gestion"),
{"schema": "web"},
)
......
......@@ -72,6 +72,7 @@ class ObjetoEntidad(Base):
Index("idx_objeto_entidad_gestion", "gestion"),
Index("idx_objeto_entidad_objeto", "objeto"),
Index("idx_objeto_entidad_objeto_entidad", "objeto", "entidad"),
Index("idx_objeto_entidad_objeto_gestion", "objeto", "gestion"),
{"schema": "web"},
)
......@@ -96,6 +97,7 @@ class ObjetoUbigeo(Base):
Index("idx_objeto_ubigeo_gestion", "gestion"),
Index("idx_objeto_ubigeo_objeto", "objeto"),
Index("idx_objeto_ubigeo_objeto_entidad", "objeto", "ubigeo"),
Index("idx_objeto_ubigeo_objeto_gestion", "objeto", "gestion"),
{"schema": "web"},
)
......
......@@ -57,6 +57,7 @@ class OrganismoEntidad(Base):
Index("idx_organismo_entidad_gestion", "gestion"),
Index("idx_organismo_entidad_objeto", "organismo"),
Index("idx_organismo_entidad_rubro_entidad", "organismo", "entidad"),
Index("idx_organismo_entidad_rubro_gestion", "organismo", "gestion"),
{"schema": "web"},
)
......
from sqlalchemy.orm import Mapped, mapped_column
from presupuestov1.model import Base
from sqlalchemy import (
Integer,
String,
Text,
Numeric,
Index,
MetaData,
)
class ProyectoResumenes(Base):
__tablename__ = "proyecto_resumenes"
__table_args__ = (
Index("idx_proyecto_resumenes_gestion", "gestion"),
Index("idx_proyecto_resumenes_nivel", "nivel"),
Index("idx_proyecto_resumenes_codigo", "codigo"),
Index("idx_proyecto_resumenes_codigo_padre", "codigo_padre"),
{"schema": "web"},
)
id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True)
gestion: Mapped[int] = mapped_column(Integer)
nivel: Mapped[str] = mapped_column(Text)
codigo: Mapped[str] = mapped_column(String(32))
codigo_padre: Mapped[str | None] = mapped_column(String(32))
devengado: Mapped[float] = mapped_column(Numeric(30, 15))
n_entidad: Mapped[int] = mapped_column(Integer)
n_objeto: Mapped[int] = mapped_column(Integer)
n_finfun: Mapped[int] = mapped_column(Integer)
n_acteco: Mapped[int] = mapped_column(Integer)
desc: Mapped[str] = mapped_column(Text)
desc_padre: Mapped[str | None] = mapped_column(Text)
entidad: Mapped[str] = mapped_column(Text)
objeto: Mapped[str] = mapped_column(Text)
finfun: Mapped[str] = mapped_column(Text)
acteco: Mapped[str] = mapped_column(Text)
......@@ -57,6 +57,7 @@ class RubroEntidad(Base):
Index("idx_rubro_entidad_gestion", "gestion"),
Index("idx_rubro_entidad_objeto", "rubro"),
Index("idx_rubro_entidad_rubro_entidad", "rubro", "entidad"),
Index("idx_rubro_entidad_rubro_gestion", "rubro", "gestion"),
{"schema": "web"},
)
......
......@@ -15,6 +15,8 @@ from sqlalchemy import and_, or_, select
class TreemapObjeto(Base):
__tablename__ = "treemap_objeto"
__table_args__ = (
Index("idx_treemap_objeto_objeto", "objeto"),
Index("idx_treemap_objeto_objeto_gestion", "objeto", "gestion"),
Index("idx_treemap_objeto_gestion_entidad", "gestion", "entidad"),
Index("idx_treemap_objeto_nivel", "nivel"),
{"schema": "web"},
......@@ -39,6 +41,8 @@ class TreemapObjeto(Base):
class TreemapFinFun(Base):
__tablename__ = "treemap_finfun"
__table_args__ = (
Index("idx_treemap_finfun_finfun", "finfun"),
Index("idx_treemap_finfun_finfun_gestion", "finfun", "gestion"),
Index("idx_treemap_finfun_gestion_entidad", "gestion", "entidad"),
Index("idx_treemap_finfun_nivel", "nivel"),
{"schema": "web"},
......@@ -63,6 +67,8 @@ class TreemapFinFun(Base):
class TreemapActeco(Base):
__tablename__ = "treemap_acteco"
__table_args__ = (
Index("idx_treemap_acteco_acteco", "acteco"),
Index("idx_treemap_acteco_acteco_gestion", "acteco", "gestion"),
Index("idx_treemap_acteco_gestion_entidad", "gestion", "entidad"),
Index("idx_treemap_acteco_nivel", "nivel"),
{"schema": "web"},
......@@ -87,6 +93,8 @@ class TreemapActeco(Base):
class TreemapRubro(Base):
__tablename__ = "treemap_rubro"
__table_args__ = (
Index("idx_treemap_rubro_rubro", "rubro"),
Index("idx_treemap_rubro_rubro_gestion", "rubro", "gestion"),
Index("idx_treemap_rubro_gestion_entidad", "gestion", "entidad"),
Index("idx_treemap_rubro_nivel", "nivel"),
{"schema": "web"},
......
......@@ -194,7 +194,7 @@ def fetch_acteco_ubigeo(
query: schema.QueryParams = Depends(),
):
gestion = query.gestion or 0
ClasObjetos = models.objeto.ActecoUbigeo
ClasObjetos = models.acteco.ActecoUbigeo
return db.query(ClasObjetos).filter(
ClasObjetos.acteco == acteco_id,
......
......@@ -81,42 +81,66 @@ def fetch_timeline(
page_size = 10
if query.entidad is not None:
return db.query(Daily).filter(
Daily.entidad == query.entidad
).order_by(
Daily.devengado.desc()
).offset(
(page - 1) * page_size
).limit(page_size).all()
# top 10 entidades
top_entities = db.query(
Daily.entidad
).group_by(
Daily.entidad
).order_by(
func.sum(Daily.devengado).desc()
).limit(10).subquery()
# ordena las filas por devengado desc por cada una de las top 10 entidades
ranked = db.query(
Daily,
stmt = (
select(Daily)
.where(Daily.entidad == query.entidad)
.order_by(Daily.devengado.desc(), Daily.id)
.offset((page - 1) * page_size)
.limit(page_size)
)
return db.execute(stmt).scalars().all()
total_devengado = func.sum(Daily.devengado).label("total_devengado")
top_entities = (
select(
Daily.entidad.label("entidad"),
total_devengado,
)
.group_by(Daily.entidad)
.order_by(total_devengado.desc(), Daily.entidad)
.limit(10)
.subquery()
)
ranked = (
select(
Daily.entidad.label("entidad"),
Daily.devengado.label("devengado"),
Daily.entidad_desc_area.label("entidad_desc_area"),
Daily.entidad_desc_entidad.label("entidad_desc_entidad"),
Daily.desc_programa.label("desc_programa"),
Daily.desc_actividad_proyecto.label("desc_actividad_proyecto"),
Daily.id.label("id"),
top_entities.c.total_devengado,
func.row_number().over(
partition_by=Daily.entidad,
order_by=Daily.devengado.desc()
).label('row_num')
).filter(
Daily.entidad.in_(top_entities)
).subquery()
# top 10 filas por cada una de las top 10 entidades
RankedTimeline = aliased(Daily, ranked)
return db.query(RankedTimeline).filter(
ranked.c.row_num <= page_size
).order_by(
RankedTimeline.entidad,
RankedTimeline.devengado.desc()
).all()
order_by=(Daily.devengado.desc(), Daily.id),
).label("row_num"),
)
.join(top_entities, top_entities.c.entidad == Daily.entidad)
.subquery()
)
stmt = (
select(
ranked.c.entidad,
ranked.c.devengado,
ranked.c.entidad_desc_area,
ranked.c.entidad_desc_entidad,
ranked.c.desc_programa,
ranked.c.desc_actividad_proyecto,
)
.where(ranked.c.row_num <= page_size)
.order_by(
ranked.c.total_devengado.desc(),
ranked.c.entidad,
ranked.c.devengado.desc(),
ranked.c.id,
)
)
return db.execute(stmt).mappings().all()
@router.get('/timeseries', response_model=List[schemas.daily.TimeSeries])
......
......@@ -46,7 +46,6 @@ def fetch_ubigeo_classifier(
return db.query(models.classifier.ClasGeografico).all()
# base
@router_entidad.get('/{entity_id:str}', response_model=schemas.entity.ClasInstitucionalIngresosGastos)
......
from fastapi import APIRouter, Query
from presupuestov1 import metrics
WINDOW_MINUTES_MAX = 7 * 24 * 60
router = APIRouter(
prefix='/_metrics',
include_in_schema=False,
)
###############################################################################
# GET
###############################################################################
@router.get('/summary')
async def fetch_metrics_summary(
minutes: int = Query(default=60, ge=1, le=WINDOW_MINUTES_MAX),
include_ts: bool = Query(default=False),
):
return await metrics.read_summary(minutes, include_ts)
@router.get('/routes')
async def fetch_metrics_routes(
minutes: int = Query(default=60, ge=1, le=WINDOW_MINUTES_MAX),
include_ts: bool = Query(default=False),
):
return await metrics.read_routes(minutes, include_ts)
@router.get('/ips')
async def fetch_metrics_ips(
minutes: int = Query(default=60, ge=1, le=WINDOW_MINUTES_MAX),
include_ts: bool = Query(default=False),
):
return await metrics.read_ips(minutes, include_ts)
import time
import pydantic_core
from typing import List, Any, Optional, Union
from fastapi import APIRouter, Depends, Request, HTTPException, Body
from sqlalchemy.orm import Session
from sqlalchemy import and_, or_, select, func, desc
from presupuestov1 import model, schema, constant, schemas, models, search
from presupuestov1.routers._decorators import not_found
......@@ -19,23 +20,55 @@ router = APIRouter(
# GET
###############################################################################
@router.get('/{project_id}')
@router.get('/{project_id}', response_model=List[schemas.project.ProyectoResumen])
@not_found
def fetch_programa_proyecto(
project_id: int,
project_id: str,
db: Session = Depends(model.get_db),
):
return []
ProyectoResumenes = models.project.ProyectoResumenes
project_items = db.query(ProyectoResumenes).filter(
ProyectoResumenes.codigo == project_id
).all()
# TODO if devengado > mucho_dinero: cache
for _ in ['entidad', 'objeto', 'finfun', 'acteco']:
for __ in project_items:
obj_ = getattr(__, _)
setattr(__, _, pydantic_core.from_json(obj_))
return project_items
@router.get('/{project_id}/proyectos')
@not_found
def fetch_programa_proyecto_proyectos(
project_id: int,
project_id: str,
db: Session = Depends(model.get_db),
query: schema.QueryParams = Depends(),
query: schema.SearchQueryParams = Depends(),
):
return []
page = query.page or 1
page_size = 20
ProyectoResumenes = models.project.ProyectoResumenes
return db.execute(
select(
ProyectoResumenes.nivel,
ProyectoResumenes.codigo,
ProyectoResumenes.desc,
func.sum(ProyectoResumenes.devengado).label('devengado'),
).where(
ProyectoResumenes.codigo_padre == project_id,
ProyectoResumenes.nivel != 'programa',
).group_by(
ProyectoResumenes.nivel,
ProyectoResumenes.codigo,
ProyectoResumenes.desc,
).order_by(desc('devengado')).offset(
(page - 1) * page_size
).limit(page_size)
).mappings().all()
@router.get('/{project_id}/objetos')
......
from typing import List, Any, Optional, Union
from fastapi import APIRouter, Depends, Request, HTTPException, Body
from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy.orm import Session
from presupuestov1 import model, schema, constant, schemas, models, search
......@@ -28,6 +28,13 @@ def fetch_search(
'hits': [],
}
try:
return search.do_search(
query.q, query.page, query.is_class, query.class_, query.order
)
except search.SearchUnavailable as exc:
raise HTTPException(
status_code=503,
detail='Search temporarily unavailable',
headers={'Retry-After': str(exc.retry_after)},
) from exc
......
......@@ -19,14 +19,14 @@ class ActecoEstado(OrmBaseModel):
top1_entidad_desc: str | None
top1_monto: float | None
top1_pct: float | None
top2_entidad: int | None | float
top2_entidad: int | None | float | str
top2_entidad_desc: str | None
top2_monto: float | None
top2_pct: float | None
top3_entidad: int | None | float
top2_monto: float | None | str
top2_pct: float | None | str
top3_entidad: int | None | float | str
top3_entidad_desc: str | None
top3_monto: float | None
top3_pct: float | None
top3_monto: float | None | str
top3_pct: float | None | str
class ActecoEntidad(OrmBaseModel):
......
......@@ -26,14 +26,14 @@ class FinfunEstado(OrmBaseModel):
top1_entidad_desc: str | None
top1_monto: float | None
top1_pct: float | None
top2_entidad: int | None | float
top2_entidad: int | None | float | str
top2_entidad_desc: str | None
top2_monto: float | None
top2_pct: float | None
top3_entidad: int | None | float
top2_monto: float | None | str
top2_pct: float | None | str
top3_entidad: int | None | float | str
top3_entidad_desc: str | None
top3_monto: float | None
top3_pct: float | None
top3_monto: float | None | str
top3_pct: float | None | str
class FinfunEntidad(OrmBaseModel):
......
......@@ -26,14 +26,14 @@ class ObjetoEstado(OrmBaseModel):
top1_entidad_desc: str | None
top1_monto: float | None
top1_pct: float | None
top2_entidad: int | None
top2_entidad: int | None | str
top2_entidad_desc: str | None
top2_monto: float | None
top2_pct: float | None
top3_entidad: int | None
top2_monto: float | None | str
top2_pct: float | None | str
top3_entidad: int | None | str
top3_entidad_desc: str | None
top3_monto: float | None
top3_pct: float | None
top3_monto: float | None | str
top3_pct: float | None | str
class ObjetoEntidad(OrmBaseModel):
......
......@@ -20,14 +20,14 @@ class OrganismoEstado(OrmBaseModel):
top1_entidad_desc: str | None
top1_monto: float | None
top1_pct: float | None
top2_entidad: int | None | float
top2_entidad: int | None | float | str
top2_entidad_desc: str | None
top2_monto: float | None
top2_pct: float | None
top3_entidad: int | None | float
top2_monto: float | None | str
top2_pct: float | None | str
top3_entidad: int | None | float | str
top3_entidad_desc: str | None
top3_monto: float | None
top3_pct: float | None
top3_monto: float | None | str
top3_pct: float | None | str
class OrganismoEntidad(OrmBaseModel):
......
from typing import List, Any, Dict
from typing import List
from presupuestov1.schema import BaseModel
from presupuestov1.schema import BaseModel, OrmBaseModel
class ProyectoResumenObj(BaseModel):
codigo: int | str
desc: str | None
desc_padre: str | None
devengado: float
class ProyectoResumen(OrmBaseModel):
gestion: int
codigo: str
codigo_padre: str | None
devengado: float
n_entidad: int
n_objeto: int
n_finfun: int
n_acteco: int
desc: str
desc_padre: str | None
entidad: str | List[ProyectoResumenObj]
objeto: str | List[ProyectoResumenObj]
finfun: str | List[ProyectoResumenObj]
acteco: str | List[ProyectoResumenObj]
class ProyectoProyectos(OrmBaseModel):
nivel: str
codigo: str
desc: str
class ProjectSearch(BaseModel):
......
......@@ -20,14 +20,14 @@ class RubroEstado(OrmBaseModel):
top1_entidad_desc: str | None
top1_monto: float | None
top1_pct: float | None
top2_entidad: int | None | float
top2_entidad: int | None | float | str
top2_entidad_desc: str | None
top2_monto: float | None
top2_pct: float | None
top3_entidad: int | None | float
top2_monto: float | None | str
top2_pct: float | None | str
top3_entidad: int | None | float | str
top3_entidad_desc: str | None
top3_monto: float | None
top3_pct: float | None
top3_monto: float | None | str
top3_pct: float | None | str
class RubroEntidad(OrmBaseModel):
......
import math
import threading
import time
import typesense
from typesense import exceptions as typesense_exceptions
TYPESENSE_API = open('./docs/typesense').read().strip()
client = typesense.Client({
......@@ -8,7 +13,7 @@ client = typesense.Client({
'port': '8108',
'protocol': 'http'
}],
'connection_timeout_seconds': 2
'connection_timeout_seconds': 0.5
})
COLLECTION = 'presup_actpro_fts'
......@@ -34,9 +39,70 @@ search_parameters = {
AVAILABLE_CLASSES = [
'programa_proyecto', 'entidad', 'objeto', 'finfun', 'acteco'
'programa_proyecto',
'entidad',
'objeto',
'finfun',
'acteco',
'ubigeo',
'rubro',
'organismo',
]
BACKEND_ERRORS = (
typesense_exceptions.Timeout,
typesense_exceptions.ServiceUnavailable,
typesense_exceptions.ServerError,
typesense_exceptions.HTTPStatus0Error,
OSError,
)
FAILURE_THRESHOLD = 2
OPEN_INTERVAL_SECONDS = 15.0
class SearchUnavailable(Exception):
def __init__(self, retry_after: int):
self.retry_after = retry_after
super().__init__('search backend unavailable')
class SearchGuard:
def __init__(self):
self._lock = threading.Lock()
self._failures = 0
self._opened_until = 0.0
def enter(self):
now = time.monotonic()
with self._lock:
if self._opened_until > now:
raise SearchUnavailable(max(1, math.ceil(self._opened_until - now)))
def record_success(self):
with self._lock:
self._failures = 0
self._opened_until = 0.0
def record_failure(self):
now = time.monotonic()
with self._lock:
self._failures += 1
if self._failures >= FAILURE_THRESHOLD:
self._opened_until = now + OPEN_INTERVAL_SECONDS
return int(OPEN_INTERVAL_SECONDS)
return 1
search_guard = SearchGuard()
def do_search(query, page, is_class, class_=None, order=None):
search_guard.enter()
search_parameters_ = search_parameters.copy()
search_parameters_['q'] = query
......@@ -60,6 +126,12 @@ def do_search(query, page, is_class, class_=None, order=None):
if order == 'asc':
search_parameters_['sort_by'] = 'devengado:asc'
return client.collections[COLLECTION].documents.search(
try:
response = client.collections[COLLECTION].documents.search(
search_parameters_
)
except BACKEND_ERRORS:
raise SearchUnavailable(search_guard.record_failure())
search_guard.record_success()
return response
......