<~I
zh|Qim#hRrfp)@(8RA32SGE;`VT7u3~dc&W#WzD6$u$2HTN<;gELD{aWU=NhnXcS-<
zMz8X=QCP_ZaO0i(xdf4oUQcF24bGiLzxtY}PJP4vCh9Jjo%_s9v@i;6IajF$wo(1f
zuGM6^SUeGM9O-M4?%Kdt=dOH@`4wvT4;9PkF=6eGSuRfiwijR-*g}bEsI&o3z-r_>
z3Fz?wJb}A8xE(GC(XV0ydx5gBFsETngC;g+BTT^Dv%fYRNFjjObm~BMdEJ&Z^9wgy
z3ZD*!v9aU1tQOw!F1QXyP^YeN7ubVp5qdyTr+#o3DCz+y
zI<{44k5*RH_EK-ORc2#D?7R){i6l0`SeZU)OLs#Xv-AvxXL}d_$ay$B)u}@FDLI|$
za=(%jI@;D};?{*YTc1&s4;5kSGh(_DTkAeeaL0~I07H$0Ju?d(4}u&maIo2+e+a>B
zgw+OtLkO^YBicBrz1=;!6ME{tpoEH4hOB()3t0)ZTcY;-zln2^=%aAL2SxxRfDyn5
zU<5D%7y*m`MgSv#5x@v+0suxn4Lh1w^IrXA2Jippzn3Q7|F5DJ3vr2u#g+EM0qt$K@bmxLe5$d&!6APBpZ~tb
b6n_4n|Ko!3{(pG?Km7c^O$ZJ3@GJfgLA^ZT
delta 204
zcmZozz}%29L7J6`fq{W}qJq636N6rxCNGf3#CL>&Uxx322PC~$G+7ixbmR6+&DiLO}fD{L-T2RJbZb1BkNZ{JfIXijv9m`Q;@7f;@d4gCZ5YT_Y7V
aGE+3Uzz$;&06M9be_}%|8yC=>2rdAht~6-?
diff --git a/models/__pycache__/base.cpython-312.pyc b/models/__pycache__/base.cpython-312.pyc
deleted file mode 100644
index bb5d27c68e0bf52b77c8e78be6fd9c8f2c1a0352..0000000000000000000000000000000000000000
GIT binary patch
literal 0
HcmV?d00001
literal 2066
zcmZux&2Q936d&8WUa$8nY(faMK-6vv+k_md3aJ6LDnyMCC4ibstJ^k6IJT$$0lg3rDRk5ds)_@*v`H^L^^F~GmMSCp`Mr5>-pqUN
z_uen1l8Ioo&;Px+l|krFDWogA%WRJ;2(2K3h~gn%u@qmmR9~|+KVxMSxvzRzU$=BW
zXXRp9^9UiV6V*(&=LtD+zk(Qu69J|ZSHSG7cM4T0xO^Syme4v_*Wu9Z?t
zQofIRw++X>D&T8?4y6{>q0WJN^_6s0&j|f|=(YVonCINag@DrrWx|-UIh}QVDzfw@
z4R~D@B|;s~W|Ux?b0!87(dO<=is$W!))iq8IK|-<@R@VTUfHpI>R}$@D6p?naFh#M
z+zo?Bq~;s$<(t#zd+=hX5k5i7IS@@b?|lC5m`4N^tQ8G!T4U8bNOwL#~1RR?k;Mo
zvBhg7t=#A6kupWC7?+@AMn;`q)`64AtKyTfozGzg@Nc62mw?1A%x
z_r3SDGWNt6+lFm4Q_nMLButFKf7^pMiXz50+P0TCmXbXmab`OlV5ADsln6%>k8lP3XSeYC(sXP8bAGL_o68xy>4sNAgXh
zqoZ*~#PDm;V9$-XggZM4Mn~a`60<^H>C+NwkVEGr*(Q~r}Sae3NMM$q;=
ziA<3rbirn5R*k>c_2ay!ds%?FbYSo
z94MlYoR0-|TpI3SPtHUpNe#Z}10X5(sv{3Z9!~vu^5^lN#y6_d-HX5W?>xEs#pcy(
z>sPOBRKL7;F>xr}%+yDO8Qt(~&sn5?XAl1@exa;^II^do0BlV3O~a%^46sUqsQ^+I
z+{OgNZ_7qzrgTtjbJvR|oY42fAek6Nm*Cz=WJT$qBu%Rq$aFd>FYnII>Qs{EJ{5ag
z8V^KNBDK#^+*^eNe-VZbiy;^%?0=aQ#@S>be8eWfmmFBp*Fd&3MNyuk$){-YIr{V^
zI{Ffg{*F#OM`K$$D!j9FV$-axo3)2*!+gJ+duG&@@|#9&-Kae^-tA_eq3k9aUPr^r
zw>Hr69m7z@mu7bm{I*7rI@mq9sZ`gM>hkFaA3r7=${837N(}~tdUi{4JH?EGm1X@O
JB>!=h{sWIS;f4SJ
diff --git a/models/__pycache__/log_alert.cpython-310.pyc b/models/__pycache__/log_alert.cpython-310.pyc
new file mode 100644
index 0000000000000000000000000000000000000000..a77da786a92ea1492acf148ce5168edc889c4588
GIT binary patch
literal 896
zcmYjQ&2H2%5VoB^Z+5r4Ko1D|0y(q?j(||1l?oE7klK^AlI5<=HpNc1j#GBe%eAk8
z#1rry9DU`~2jIeqaZ<5$2KaViuc(B_?HQ7mIXeD5DrDx*$S)ZOcTUHAnpR1Mg^cYglLn`)DspeW(;mRAuFGQo}
zf$!l+0UW|?ZonwSutbbV%pw)5L=t$S+lZxrG8u311~VEWunz9m9M&P=QI`V|nab>z
zurc5rNz~+ODk*S0+T!lr;5KHDWezHbo_?$sZ?&?H<+@&Q5!~jyAbz;Cr4e?`zo@lO
z%1Wx$-3zGnZWoNtma3AV9F%5KeHN{Du!mNy*5U#&>3<8Z1nDqpJSOhI6Z+Iv%#MVP9$oejxiiqz3c%)$z;O(P}k2
zDc$?_;*+hf)Z8^Q=-htSx&21jn{uvZAWG>j(6{tpxwZpf2(t+tq$tKaWP~ZEgyJcl
zlBc-KI4y>KS$T7V?a!Gl|8DSXDq@8Irr#9JA?w->_n
Vb>IFsI20Csw-MTdvQn}e{R5Io>XHBe
literal 0
HcmV?d00001
diff --git a/models/base.py b/models/base.py
deleted file mode 100644
index d35671e..0000000
--- a/models/base.py
+++ /dev/null
@@ -1,40 +0,0 @@
-from sqlalchemy import Column, BigInteger, DateTime, event
-from sqlalchemy.ext.declarative import declared_attr, declarative_base
-
-from datetime import datetime
-
-from config.database import Base
-from utils.common import camel_to_snake
-from id_generator import options, generator
-
-# https://github.com/yitter/IdGenerator/tree/master/Python
-options = options.IdGeneratorOptions(worker_id=23)
-idgen = generator.DefaultIdGenerator()
-idgen.set_id_generator(options)
-
-
-# 第二层基类:包含ID
-class IdBase(Base):
- __abstract__ = True
-
- id = Column(BigInteger, primary_key=True, index=True)
-
- @declared_attr
- def __tablename__(cls):
- # 自动把数据库实体类名驼峰转为数据库表名下划线
- return camel_to_snake(cls.__name__)
-
-
-# 自动填充id
-@event.listens_for(IdBase, 'before_insert', propagate=True)
-def before_insert_listener(mapper, connection, target):
- if target.id is None:
- target.id = idgen.next_id()
-
-
-# 第二层基类:包含ID和审计字段
-class AuditBase(IdBase):
- __abstract__ = True
-
- create_time = Column(DateTime, nullable=True, default=datetime.now)
- update_time = Column(DateTime, nullable=True, default=datetime.now, onupdate=datetime.now)
diff --git a/models/log_alert.py b/models/log_alert.py
index 6d8f97e..791d8ad 100644
--- a/models/log_alert.py
+++ b/models/log_alert.py
@@ -1,7 +1,8 @@
-from sqlalchemy import Column, Integer, String, DateTime, Text
-from sqlalchemy.ext.declarative import declarative_base
from datetime import datetime
+from sqlalchemy import Column, Integer, String, Text, DateTime
+from sqlalchemy.ext.declarative import declarative_base
+
Base = declarative_base()
@@ -9,9 +10,12 @@ class LogAlert(Base):
__tablename__ = "log_alerts"
id = Column(Integer, primary_key=True, index=True)
- timestamp = Column(DateTime)
- log_level = Column(String(20))
+ alter_name = Column(String(50))
+ alter_timestamp = Column(Integer)
+
+ timestamp = Column(Integer)
message = Column(Text)
- source = Column(String(100))
- context = Column(Text) # 存储JSON格式的上下文信息
status = Column(String(20), default="pending") # pending, notified, resolved
+
+ create_time = Column(DateTime, nullable=True, default=datetime.now)
+ update_time = Column(DateTime, nullable=True, default=datetime.now, onupdate=datetime.now)
diff --git a/requirements.txt b/requirements.txt
index 36adcf2..ad92be6 100644
--- a/requirements.txt
+++ b/requirements.txt
@@ -1 +1 @@
-fastapi~=0.119.0
pydantic~=2.12.2
SQLAlchemy~=2.0.44
requests~=2.32.5
\ No newline at end of file
+fastapi~=0.119.0
pydantic~=2.12.2
SQLAlchemy~=2.0.44
requests~=2.32.5
pydantic-settings~=2.11.0
uvicorn~=0.38.0
\ No newline at end of file
diff --git a/schemas/__pycache__/log_alert.cpython-310.pyc b/schemas/__pycache__/log_alert.cpython-310.pyc
new file mode 100644
index 0000000000000000000000000000000000000000..2c2bf8a96370d14b3af5e2627deda178133592ad
GIT binary patch
literal 795
zcmZva&2AGh5XbHPNOrdg%^|=I<+@1a)DwzO3#dT|0tmRQB3W)cO&580Q`%aioP<4AyItFAjNMq#oyt!iQj$KtqwI9_r=t!6_~1!+JO)($R|wg*os+KaT1(Nobw8OI;DLhtnZAC$
zcy)fhnAPsh=J=!OzQNMjh3-~+k`A+H!F;cmu&~QhXk^HUka@P=8p27__D~9z@c|F9
z|4`g=e~X3ppdciO0#X$d^Q)-p1x!7})35r$&d|=Q405b3(B46813&t(hPFGVW5TBm
zn5{36lW~HG`3zbWEQ8&^d~#yC2H&o);6!e;^QqLev_7S{xTfCmiLEc-;-LkNBqmQM
zl$!jjkL{VBmxZDGnm$QJvO_|v4ZXP;lRP4MOp=q_BEPYQ6nlU|NlQt>m+>(o9PyBu
nm#NRzTP0gpFHJ-X!tTT}6x_F`&!#sG@q#WK-etr_@i_epsyMb2
literal 0
HcmV?d00001
diff --git a/service/__pycache__/log_alert_service.cpython-310.pyc b/service/__pycache__/log_alert_service.cpython-310.pyc
new file mode 100644
index 0000000000000000000000000000000000000000..e6c1bb9334ecca2df98fd010af73f37262b71edf
GIT binary patch
literal 2170
zcmY*a&2Jk;6rY*>^x96GkETs4AY1@PY6Lw(6sQ_LBxpdS5@?mmT6-pr-CaAG8MjTf
zKDiR85{Fhrm1t4z3sMhoLY%nqU(6L1krTZjp(0w|+l>=8&`=~
zfQ1@wQz+Hyc0lz%={C>D;0AqqtY4C?W>ZKo&2@K0r(1{pS49nFt05reoDv;gyCxbt
zk*+0>ENTi+*Zd@sJ
zPj(s%-{o*1C@#zkFx4Fpnxv#nIUTSq)@OZ7v!Ru;)EcqAo!Ud2vk~n(+DV-O<<=D@
z&v1K_q;%+NPx}KpENH(SV2gIikbji1gjm<)laDTNuk7@Rsfq?&>ne4d{LHz1{cp
zl%8%^Qg6U0(K9JYoskIVyswn
zj3=&(Mh^;TLx99W8Ry!Lkitx7b)6`-v)z(Pn;^G$qD06NZYrR_&Ik(^&wu{$`lnxP
zoHM0*q#HZ6oIO=vmN94wgb4srnS*G2WC9)%p;N#*IfHph2rIAW$cE+ofrm_aMC8j5
z^K^av?X3&fuWx
zkZK7;g_dYQOO{P-%cWD4J+m$LpYPc8lvM(|#Gve~4GXmKz#2L*(WcIU&_o0JL-rNf
z2F7uA(@*I-831egz&YTZHG=B_GEuZivnXB#0e!9`-BVST
zd(i>y$X4Do7LfB1&4EM&9P`_%Uk#>ps7+G)ae*
z)wtV;Vzma!c=IJV50)*bpnElndqT+*_y`xWCoNsI^JgPNawgeICpu$;d<~6wAbA?Z
z>o8RXgbygfS^lw~2W6kyEc3olUxe)FdiW9+DH#H2njuI-YeW#5KGjYJA3)b<2;t><
z1TeMQF4W#?dy8b~>jG%OfPyRn4sbb8vNxwoDdft5h4e2B1a}jrLWFa!Kp@?K>Z8lp
zSf2a^vx*;=-CvEPq|T$YcX4}8k{d7W^@QBdsXSTZJQ9)~!CWn8rJZPcPXQv4RK|+p
zID7grJ&YUeB5U$?KT2Zk>>FWcV1#bgK`IOyaRlFo?Uvxi$!rrZkWr>Mam;t{_=wa3
zbd^sl6i}kUGvATtz_?L8X^M?*j9(u3;;6hTf!d3sxUnOW{hTt0p_w|iuE1wPVjr4H
z(&ZvnD@VPb(`!cI?r75PVYqjASLC-6c_`Tqewi6LzO
literal 0
HcmV?d00001
diff --git a/service/__pycache__/message_service.cpython-310.pyc b/service/__pycache__/message_service.cpython-310.pyc
new file mode 100644
index 0000000000000000000000000000000000000000..4fabce2867068222fb7d9ed1534df6db042d5119
GIT binary patch
literal 1092
zcmZWoOK%e~5VpPE%_iBTp_LFl@VIbTq!y_sgb*sUrAVj&K_t*B(rQ_|$tL?qZEs&{
zPgLCc4>-V)1HXa);VUP802jo8@n$0`Vael}`8*GMW{OUyg<#$N_G|pzLgn#ZewK7u-hQceERZd6e-m)U7bg3Rb3u_rma?jM5sRgPnH+
z?e0B!{bKv&+rDltxQ(w!bRdLqaJY4RJlsvx^YZ;GQA~KOwUDLmX*04so=dRucBYXdMo?M1U$n#7Sxiv0dW7P+>j{Nx!lk;ijc3LZu>
zPLfl*RQqcWd`maRJYrl(?aU$(Wl~cXsc7JP)gKEp2BK|<2SmGqACz1wxHc=K66SyP
z@)3{6kqVFa`*BfB!&0RG=^zPBt9)!$W-|!K^SEH#x5Se1-UAYJYe#h@|8%HfOnd+a7&xAsmbTCL{^&`1@ArO|>h}|Z^8L%N#helH!!K@2gp22>W*?0pf|exp
zo|Uv>DWg6P%An#YZ)0ADRg^{;htLrbY6s$zNW`N1nWtSr-;zP{6Kj#f0d@U@$>eA+
zGe!QQu?tMZ2COaQ++e)*x(2C^4e0mqrUAWjiLsfsGh0;PVvu4JMH+0QWsdOwOMp>~
zIz~0SXe?P$%a(zl__41{1eUMKnhJi-F2YMvP#fWX2Q~i0M0n0WvNgNtUJ?Wx9iDAv
zbzK&DW(y^sERB)_=Atamp|)qev62>~9Zps?xbU6FI$%1)w%YSI2wZ0YnE-7Z_jeu7
z{J4pWHQ;GuF<&cVwTGQPP&aeW8eM9RS(j8wKsgglRck+K`smIvWJzWf+|AaRF&~r`
z`d*9T29a-UY9$SPBYGP*|LnhfIvt;$PG1%FpgBI$Y6-bD(^AccSqa(>{ck98m|pwa
z-LG<^dgVHEurs{n$u~P_IOUWvM)}_z-lNg)IMEM1nnd@}x&Ac+BTva@virG`C&hep
z^%h!I+_uRVP-SL>#b+&~4KG0 bool:
+ start_time, end_time = get_timestamp_range(log_alter.timestamp, 60)
+ result = openobserve_service.get_log(start_time, end_time)
+ hits = result.get("hits")
- url = 'http://192.168.1.7:5080/api/default/_search'
- username = 'njcxx0822@163.com'
- password = '19940822Cxx'
+ # 判断日志是否为空
+ if len(hits) == 0:
+ return False
- response = requests.post(
- url,
- headers={
- 'accept': 'application/json',
- 'Content-Type': 'application/json'
- },
- json={
- "query": OpenobserveQuery(
- start_time=1760773184953420,
- end_time=1760773189953431
- ).model_dump()
- },
- auth=HTTPBasicAuth(username, password)
+ # 判断日志告警是否已经存在
+ log_db_alter = db.execute(
+ select(LogAlert)
+ .where(LogAlert.alter_timestamp == log_alter.timestamp)
+ ).first()
+
+ if log_db_alter is not None:
+ return False
+
+ # 企业微信通知
+ alter_message = f"🔔 **告警服务**:{log_alter.alter_name} \n 🕒 **告警时间**:{format_timestamp(log_alter.timestamp)}"
+ send_wechat_message(MessageEnum.MARKDOWN2, alter_message)
+
+ # 批量插入
+ return batch_insert(log_alter, hits, db)
+
+
+def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool:
+ stmt = insert(LogAlert).values(
+ [
+ {
+ "alter_name": log_alter.alter_name,
+ "alter_timestamp": log_alter.timestamp,
+ "timestamp": hit["_timestamp"],
+ "message": hit["message"],
+ "status": "pending"
+ }
+ for hit in hits
+ ]
)
- return response.json().get("hits")
+ db.execute(stmt)
+ db.commit()
+
+ return True
+
+
+def query_alert(db: Session):
+ alerts = db.execute(select(LogAlert)).scalars().all()
+
+ # 按照name和timestamp分类
+ classified = defaultdict(lambda: defaultdict(list))
+ for alert in alerts:
+ classified[alert.alter_name][alert.alter_timestamp].append({
+ "timestamp": format_timestamp(alert.timestamp),
+ "message": alert.message
+ })
+
+ return dict(classified)
diff --git a/service/message_service.py b/service/message_service.py
new file mode 100644
index 0000000..c10d057
--- /dev/null
+++ b/service/message_service.py
@@ -0,0 +1,32 @@
+import json
+from enum import Enum
+
+import requests
+
+from config.setting import settings
+
+
+class MessageEnum(Enum):
+ TEXT = 'text'
+ MARKDOWN = 'markdown'
+ MARKDOWN2 = 'markdown2'
+
+
+def send_wechat_message(message_type: MessageEnum, message: str):
+ wechat_message = {}
+ match message_type:
+ case MessageEnum.TEXT:
+ wechat_message = {"msgtype": "text", "text": {"content": message}}
+ case MessageEnum.MARKDOWN:
+ wechat_message = {"msgtype": "markdown", "markdown": {"content": message}}
+ case MessageEnum.MARKDOWN2:
+ wechat_message = {"msgtype": "markdown_v2", "markdown_v2": {"content": message}}
+
+ requests.post(
+ settings.wechat_webhook_url,
+ headers={"Content-Type": "application/json"},
+ params={'key': settings.wechat_webhook_key},
+ data=json.dumps(wechat_message, ensure_ascii=False).encode('utf-8')
+ )
+
+ return True
diff --git a/service/openobserve_service.py b/service/openobserve_service.py
index 961d813..14ef81b 100644
--- a/service/openobserve_service.py
+++ b/service/openobserve_service.py
@@ -1,27 +1,19 @@
import requests
from requests.auth import HTTPBasicAuth
+from config.setting import settings
from schemas.log_alert import OpenobserveQuery
-def get_log():
- url = 'http://192.168.1.7:5080/api/default/_search'
- username = 'njcxx0822@163.com'
- password = '19940822Cxx'
-
+def get_log(start_time: int, end_time: int) -> str:
response = requests.post(
- url,
+ settings.openobserve_url,
headers={
'accept': 'application/json',
'Content-Type': 'application/json'
},
- json={
- "query": OpenobserveQuery(
- start_time=1760773184953420,
- end_time=1760773189953431
- ).model_dump()
- },
- auth=HTTPBasicAuth(username, password)
+ json={"query": OpenobserveQuery(start_time=start_time, end_time=end_time).model_dump()},
+ auth=HTTPBasicAuth(settings.openobserve_username, settings.openobserve_password)
)
- print(response)
+ return response.json()
diff --git a/utils/__pycache__/time.cpython-310.pyc b/utils/__pycache__/time.cpython-310.pyc
new file mode 100644
index 0000000000000000000000000000000000000000..ef6824413006460e849bc98412061764fe35a44a
GIT binary patch
literal 899
zcmY*Xy-yTT5P$FO+r0(u1O*XeL*W%$oSn^S3=tD)B*wGxatT@Yy**d=3y#OALidVdtN)tsLA{RyM}Mc?(ByHZwajZ{E!OW?ob(^&_C<`|o=@
z1fieGI8Fm#OaRBXKrqCxj{;OjSn0&a0jU$v2bs<_W&jsBWhN(Fudl{O-^0PimbFi>jpbfL)K6i@*)doCbK1%Ucf=uZznCEcXc?)ih<{dDCSNSvN$6$
z&Pt5?(w^zLXcBrdGX$532pvtOJ0x{q@g#6u1(8Dk53n#8sUfB@%5-M1!dtq5Qo{62
zGOwjL)ke@fO;Jk6AU>sMkPFe4nbaUb|2R-*=XHDi-QMcw?&{Xw;-~h~a{I-~e?9Gg
zc)7o}(0Tr(y}tG9%VK-^&FjBH}qfi^AFb9{8Ne5X6YzLZR&qrY9F@q-SPEB}78a-B
zJjZ+>hE@2G8;QU^hGY!9vC1@ph0>>CF{-b8%vRx`h#}`%Uf2}ssb!iM%1nbgIllbT
eu_|iDA}?iL4kDKL{EC7atA!L%5uevahyMYebnkrt
literal 0
HcmV?d00001
diff --git a/utils/time.py b/utils/time.py
new file mode 100644
index 0000000..0aac705
--- /dev/null
+++ b/utils/time.py
@@ -0,0 +1,25 @@
+from datetime import datetime
+from typing import Tuple
+
+
+def get_timestamp_range(ts: int, delta_seconds: int = 5, unit: str = "microseconds") -> Tuple[int, int]:
+ """
+ 返回时间戳前后delta_seconds秒的范围(单位与输入一致)
+ """
+ if ts <= 0:
+ return 0, 0
+
+ scale = {
+ "seconds": 1,
+ "milliseconds": 1000,
+ "microseconds": 1_000_000,
+ "nanoseconds": 1_000_000_000
+
+ }.get(unit, 1_000_000) # 默认微秒
+
+ delta = delta_seconds * scale
+ return ts - delta, ts + delta
+
+
+def format_timestamp(timestamp: int) -> str:
+ return datetime.fromtimestamp(timestamp / 1_000_000).strftime("%Y-%m-%d %H:%M:%S")