feat: demo managed functions with harbor-backed operator setup

This commit is contained in:
Naeel
2026-03-14 18:16:53 +03:00
parent cbd2c8c44c
commit 8dd5b676c0
14 changed files with 687 additions and 1 deletions
@@ -0,0 +1,45 @@
# Изменено: 2026-03-14
# Функция event-cleaner: удаляет N самых старых строк из таблицы events.
# Вызывается через HTTP POST из Node-RED (который слушает RabbitMQ).
# Env: POSTGRES_DSN — строка подключения к PostgreSQL.
import os
import json
import psycopg2
def handle(request):
"""Удаляет N старейших строк из таблицы events."""
dsn = os.environ["POSTGRES_DSN"]
body = {}
if request.get_data():
try:
body = json.loads(request.get_data())
except Exception:
pass
# Количество строк для удаления — из тела запроса или дефолт 10
delete_n = int(body.get("delete_n", 10))
# Защита от случайного удаления слишком большого количества строк
delete_n = min(delete_n, 100)
conn = psycopg2.connect(dsn)
try:
with conn.cursor() as cur:
cur.execute("""
DELETE FROM events
WHERE id IN (
SELECT id FROM events ORDER BY created_at ASC LIMIT %s
)
""", (delete_n,))
deleted = cur.rowcount
cur.execute("SELECT COUNT(*) FROM events")
remaining = cur.fetchone()[0]
conn.commit()
return json.dumps({
"ok": True,
"deleted": deleted,
"remaining": remaining
}), 200, {"Content-Type": "application/json"}
finally:
conn.close()
@@ -0,0 +1,47 @@
# Изменено: 2026-03-14
# Функция event-monitor: считает строки в events.
# Если больше 50 — публикует сообщение в RabbitMQ queue "cleanup-needed".
# Запускается по cron (каждую минуту).
# Env:
# POSTGRES_DSN — строка подключения к PostgreSQL
# RABBITMQ_URL — amqp://sless:sless123@rabbitmq.sless.svc.cluster.local:5672/
import os
import json
import psycopg2
import pika
THRESHOLD = 50
def handle(request):
"""Мониторит таблицу events. При переполнении шлёт в RabbitMQ."""
dsn = os.environ["POSTGRES_DSN"]
rabbit_url = os.environ["RABBITMQ_URL"]
conn = psycopg2.connect(dsn)
try:
with conn.cursor() as cur:
# Считаем количество событий
cur.execute("SELECT COUNT(*) FROM events")
count = cur.fetchone()[0]
finally:
conn.close()
result = {"count": count, "threshold": THRESHOLD, "action": "none"}
if count > THRESHOLD:
# Публикуем в очередь — event-cleaner получит и удалит старые строки
params = pika.URLParameters(rabbit_url)
connection = pika.BlockingConnection(params)
channel = connection.channel()
channel.queue_declare(queue="cleanup-needed", durable=True)
channel.basic_publish(
exchange="",
routing_key="cleanup-needed",
body=json.dumps({"count": count, "delete_n": 10}),
properties=pika.BasicProperties(delivery_mode=2) # persistent
)
connection.close()
result["action"] = "cleanup_requested"
return json.dumps(result), 200, {"Content-Type": "application/json"}
@@ -0,0 +1,49 @@
# Изменено: 2026-03-14
# Функция event-writer: принимает HTTP POST, пишет одну строку в таблицу events.
# Таблица создаётся автоматически при первом запуске.
# Env: POSTGRES_DSN — строка подключения к PostgreSQL.
import os
import json
import psycopg2
from datetime import datetime, timezone
def handle(request):
"""Записывает одно событие в таблицу events."""
dsn = os.environ["POSTGRES_DSN"]
body = {}
if request.get_data():
try:
body = json.loads(request.get_data())
except Exception:
pass
source = body.get("source", "node-red")
message = body.get("message", "ping")
conn = psycopg2.connect(dsn)
try:
with conn.cursor() as cur:
# Создаём таблицу если нет — безопасно вызывать при каждом запросе
cur.execute("""
CREATE TABLE IF NOT EXISTS events (
id SERIAL PRIMARY KEY,
source VARCHAR(100) NOT NULL DEFAULT 'unknown',
message TEXT NOT NULL DEFAULT '',
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
)
""")
cur.execute(
"INSERT INTO events (source, message) VALUES (%s, %s) RETURNING id, created_at",
(source, message)
)
row = cur.fetchone()
conn.commit()
return json.dumps({
"ok": True,
"id": row[0],
"created_at": row[1].isoformat()
}), 200, {"Content-Type": "application/json"}
finally:
conn.close()