feat(accounting): 独立异步记账进程——模型用量出账(客户付/商户营收/供应商成本),同步记账留web,进程由start/stop.sh管理

This commit is contained in:
yumoqing 2026-09-04 19:32:43 +08:00
parent 4e3521dfd5
commit e0dbac6f9d
3 changed files with 113 additions and 0 deletions

View File

@ -0,0 +1,83 @@
#!/usr/bin/env python3
"""独立异步记账进程 — 产线模型用量出账(三方账)。
职责分工2026-09 定夺
- web 服务器同步记账账号/存储实时购买当场落账
- 本进程异步记账模型调用时只记成本侧 + pending 流水
本进程扫 llm_usage 出三方账客户付/商户营收/供应商成本
只加载记账链路所需模块 HTTP RBAC 路由
用法
py3/bin/python app/pipeline_accounting_worker.py -w <workdir>
部署 start.sh nohup 启动pid 文件 pipeline-accounting.pid
"""
import os, sys, asyncio, argparse, logging
logger = logging.getLogger("pipeline.accounting_worker")
def init_worker():
"""初始化:配置 → DB → ServerEnv → 记账链路模块(与主应用同一套加载,裁剪到必需)。"""
app_dir = os.path.dirname(os.path.abspath(__file__))
root_dir = os.path.dirname(app_dir)
sys.path.insert(0, root_dir)
sys.path.insert(0, app_dir)
from appPublic.jsonConfig import getConfig
from appPublic.folderUtils import ProgramPath
from sqlor.dbpools import DBPools
from ahserver.serverenv import ServerEnv
from appPublic.event_dispatcher import EventDispatcher
config = getConfig(root_dir, NS={'workdir': root_dir, 'ProgramPath': ProgramPath()})
DBPools(config.databases)
env = ServerEnv()
env.event_dispatcher = EventDispatcher()
def get_module_dbname(mname):
"""All modules use the pipeline database."""
return 'pipeline'
env.get_module_dbname = get_module_dbname
# 记账链路必需模块顺序同主应用pricing 先于资源/产品模块)
from pricing.init import load_pricing
from accounting.init import load_accounting
from discount.init import load_discount
from supplychain.init import load_supplychain
from product_management import load_product_management
from pipeline_llm.init import load_llm
load_pricing() # buffered_charging定价引擎
load_accounting() # consume_accounting / 余额 / 业务日期
load_discount() # get_min_product_discount客户折扣→售价
load_supplychain() # calculate_sale_amounts / get_distribution_chain供应商折扣→成本
load_product_management() # product_accounting_generic三方账落账
load_llm() # pipeline_llm流水扫描 + 结算分流)
logger.info("[accounting_worker] 记账链路模块加载完成")
async def _main():
from pipeline_llm.accounting import llm_usage_accounting_loop
await llm_usage_accounting_loop()
def main():
parser = argparse.ArgumentParser()
parser.add_argument('-w', '--workdir', default=os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
args = parser.parse_args()
os.chdir(args.workdir)
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s[%(levelname)s][%(name)s]%(message)s')
logger.info("[accounting_worker] 启动workdir=%s", args.workdir)
init_worker()
try:
asyncio.run(_main())
except KeyboardInterrupt:
logger.info("[accounting_worker] 收到中断信号,退出")
if __name__ == '__main__':
main()

View File

@ -35,3 +35,22 @@ export PATH="$WORKDIR/bin:$PATH"
$WORKDIR/py3/bin/python $WORKDIR/app/pipeline_app.py -p $PORT -w $WORKDIR >> $WORKDIR/logs/pipeline.log 2>&1 &
echo $! > "$PIDFILE"
echo "Started PID $(cat $PIDFILE)"
# 独立异步记账进程(模型用量出账:客户付/商户营收/供应商成本)
# 只随主进程/worker 启动web 单独模式不启动(同步记账在 web 内)
if [ "$MODE" != "web" ]; then
ACC_PIDFILE="pipeline-accounting.pid"
if [ -f "$ACC_PIDFILE" ]; then
acc_pid=$(cat "$ACC_PIDFILE")
if kill -0 "$acc_pid" 2>/dev/null; then
echo "Accounting worker already running (PID $acc_pid)"
else
rm -f "$ACC_PIDFILE"
fi
fi
if [ ! -f "$ACC_PIDFILE" ]; then
$WORKDIR/py3/bin/python $WORKDIR/app/pipeline_accounting_worker.py -w $WORKDIR >> $WORKDIR/logs/pipeline-accounting.log 2>&1 &
echo $! > "$ACC_PIDFILE"
echo "Started accounting worker PID $(cat $ACC_PIDFILE)"
fi
fi

11
stop.sh
View File

@ -14,3 +14,14 @@ if [ -f pipeline.pid ]; then
else
echo "No pid file found"
fi
if [ -f pipeline-accounting.pid ]; then
pid=$(cat pipeline-accounting.pid)
if kill -0 "$pid" 2>/dev/null; then
kill "$pid"
echo "Stopped accounting worker PID $pid"
else
echo "Accounting worker $pid not running"
fi
rm -f pipeline-accounting.pid
fi