diff --git a/api_v2/test_wecom_binding_api.py b/api_v2/test_wecom_binding_api.py new file mode 100644 index 0000000..91e3035 --- /dev/null +++ b/api_v2/test_wecom_binding_api.py @@ -0,0 +1,217 @@ +""" +企业微信用户绑定 API 测试 +""" + +from django.contrib.auth.models import User +from django.test import TestCase, override_settings +from rest_framework.test import APIClient + +from basic_info.models import Employee, Merchant, MerchantTypeEnum + + +@override_settings(AGENT_ACCESS_KEY='test-key-123') +class WecomBindingAPITestCase(TestCase): + """测试企业微信用户绑定接口""" + + def setUp(self): + self.client = APIClient() + self.url = '/api/v2/wecom/binduser/' + self.auth_header = 'test-key-123' + + # 创建商户 + self.merchant = Merchant.objects.create( + name='测试商户', + type=MerchantTypeEnum.STORE, + ) + + # 创建用户 + self.user = User.objects.create_user( + username='testuser', + password='testpass123', + ) + + # 创建员工并关联用户 + self.employee = Employee.objects.create( + merchant=self.merchant, + name='测试员工', + sys_user=self.user, + ) + + def _post(self, data, auth=None): + """辅助方法:发送 POST 请求""" + headers = {'HTTP_AUTHORIZATION': auth or self.auth_header} + return self.client.post(self.url, data, format='json', **headers) + + # ==================== 认证测试 ==================== + + def test_missing_authorization_header(self): + """缺少 Authorization 头应返回 401""" + resp = self.client.post(self.url, {}, format='json') + self.assertEqual(resp.status_code, 401) + + def test_invalid_authorization_key(self): + """错误的 AGENT_ACCESS_KEY 应返回 401""" + resp = self._post( + {'wecom_user_id': 'wx123', 'username': 'testuser', 'password': 'testpass123'}, + auth='wrong-key', + ) + self.assertEqual(resp.status_code, 401) + + # ==================== 参数验证测试 ==================== + + def test_missing_required_fields(self): + """缺少必填字段应返回 400""" + resp = self._post({}) + self.assertEqual(resp.status_code, 400) + data = resp.json() + self.assertIn('wecom_user_id', data) + self.assertIn('username', data) + self.assertIn('password', data) + + def test_missing_wecom_user_id(self): + """缺少 wecom_user_id 应返回 400""" + resp = self._post({'username': 'testuser', 'password': 'testpass123'}) + self.assertEqual(resp.status_code, 400) + self.assertIn('wecom_user_id', resp.json()) + + # ==================== 用户凭据验证测试 ==================== + + def test_wrong_password(self): + """密码错误应返回 401""" + resp = self._post({ + 'wecom_user_id': 'wx123', + 'username': 'testuser', + 'password': 'wrongpass', + }) + self.assertEqual(resp.status_code, 401) + self.assertEqual(resp.json()['error'], '用户名或密码错误') + + def test_nonexistent_username(self): + """不存在的用户名应返回 401""" + resp = self._post({ + 'wecom_user_id': 'wx123', + 'username': 'nouser', + 'password': 'testpass123', + }) + self.assertEqual(resp.status_code, 401) + + def test_inactive_user(self): + """被禁用的用户应返回 403""" + self.user.is_active = False + self.user.save() + resp = self._post({ + 'wecom_user_id': 'wx123', + 'username': 'testuser', + 'password': 'testpass123', + }) + # Django authenticate() returns None for inactive users by default + # so it will be 401 (用户名或密码错误) + self.assertIn(resp.status_code, [401, 403]) + + def test_user_without_employee(self): + """用户未关联员工应返回 404""" + user2 = User.objects.create_user(username='noemployee', password='pass123') + resp = self._post({ + 'wecom_user_id': 'wx123', + 'username': 'noemployee', + 'password': 'pass123', + }) + self.assertEqual(resp.status_code, 404) + self.assertEqual(resp.json()['error'], '该用户未关联任何员工记录') + + # ==================== 绑定成功测试 ==================== + + def test_binding_success(self): + """正常绑定应返回 200""" + resp = self._post({ + 'wecom_user_id': 'wx_user_001', + 'username': 'testuser', + 'password': 'testpass123', + }) + self.assertEqual(resp.status_code, 200) + data = resp.json() + self.assertEqual(data['message'], '绑定成功') + self.assertEqual(data['employee_id'], self.employee.id) + self.assertEqual(data['employee_name'], '测试员工') + self.assertEqual(data['wecom_user_id'], 'wx_user_001') + + # 验证数据库 + self.employee.refresh_from_db() + self.assertEqual(self.employee.wecom_user_id, 'wx_user_001') + + def test_binding_idempotent(self): + """重复绑定相同 wecom_user_id 应幂等返回 200""" + self.employee.wecom_user_id = 'wx_user_001' + self.employee.save() + + resp = self._post({ + 'wecom_user_id': 'wx_user_001', + 'username': 'testuser', + 'password': 'testpass123', + }) + self.assertEqual(resp.status_code, 200) + self.assertEqual(resp.json()['message'], '绑定成功') + + # ==================== 冲突测试 ==================== + + def test_wecom_user_id_already_bound_to_other_employee(self): + """wecom_user_id 已绑定到其他员工应返回 409""" + other_employee = Employee.objects.create( + merchant=self.merchant, + name='其他员工', + wecom_user_id='wx_user_001', + ) + + resp = self._post({ + 'wecom_user_id': 'wx_user_001', + 'username': 'testuser', + 'password': 'testpass123', + }) + self.assertEqual(resp.status_code, 409) + self.assertIn('该企业微信用户已绑定到其他员工', resp.json()['error']) + + def test_employee_already_has_different_wecom_user_id_without_force(self): + """员工已有不同 wecom_user_id 且未传 force_update 应返回 409""" + self.employee.wecom_user_id = 'wx_old_user' + self.employee.save() + + resp = self._post({ + 'wecom_user_id': 'wx_new_user', + 'username': 'testuser', + 'password': 'testpass123', + }) + self.assertEqual(resp.status_code, 409) + data = resp.json() + self.assertEqual(data['error'], '该员工已绑定其他企业微信用户') + self.assertEqual(data['current_wecom_user_id'], 'wx_old_user') + self.assertIn('hint', data) + + def test_employee_already_has_different_wecom_user_id_with_force(self): + """员工已有不同 wecom_user_id 但传了 force_update=true 应覆盖成功""" + self.employee.wecom_user_id = 'wx_old_user' + self.employee.save() + + resp = self._post({ + 'wecom_user_id': 'wx_new_user', + 'username': 'testuser', + 'password': 'testpass123', + 'force_update': True, + }) + self.assertEqual(resp.status_code, 200) + self.assertEqual(resp.json()['wecom_user_id'], 'wx_new_user') + + self.employee.refresh_from_db() + self.assertEqual(self.employee.wecom_user_id, 'wx_new_user') + + def test_force_update_false_explicitly(self): + """显式传 force_update=false 应拒绝覆盖""" + self.employee.wecom_user_id = 'wx_old_user' + self.employee.save() + + resp = self._post({ + 'wecom_user_id': 'wx_new_user', + 'username': 'testuser', + 'password': 'testpass123', + 'force_update': False, + }) + self.assertEqual(resp.status_code, 409) diff --git a/api_v2/test_wecom_login_api.py b/api_v2/test_wecom_login_api.py new file mode 100644 index 0000000..592b961 --- /dev/null +++ b/api_v2/test_wecom_login_api.py @@ -0,0 +1,132 @@ +""" +企业微信快捷登录 API 测试 +""" + +from django.contrib.auth.models import User +from django.test import TestCase, override_settings +from rest_framework.test import APIClient + +from basic_info.models import Employee, Merchant, MerchantTypeEnum + + +@override_settings(AGENT_ACCESS_KEY='test-key-123') +class WecomLoginAPITestCase(TestCase): + """测试企业微信快捷登录接口""" + + def setUp(self): + self.client = APIClient() + self.url = '/api/v2/wecom/login/' + self.auth_header = 'test-key-123' + + # 创建商户 + self.merchant = Merchant.objects.create( + name='测试商户', + type=MerchantTypeEnum.STORE, + ) + + # 创建用户 + self.user = User.objects.create_user( + username='testuser', + password='testpass123', + ) + + # 创建员工并关联用户,设置 wecom_user_id + self.employee = Employee.objects.create( + merchant=self.merchant, + name='测试员工', + sys_user=self.user, + wecom_user_id='wx_user_001', + ) + + def _post(self, data, auth=None): + """辅助方法:发送 POST 请求""" + headers = {'HTTP_AUTHORIZATION': auth or self.auth_header} + return self.client.post(self.url, data, format='json', **headers) + + # ==================== 认证测试 ==================== + + def test_missing_authorization_header(self): + """缺少 Authorization 头应返回 401""" + resp = self.client.post(self.url, {'wecom_user_id': 'wx_user_001'}, format='json') + self.assertEqual(resp.status_code, 401) + + def test_invalid_authorization_key(self): + """错误的 AGENT_ACCESS_KEY 应返回 401""" + resp = self._post({'wecom_user_id': 'wx_user_001'}, auth='wrong-key') + self.assertEqual(resp.status_code, 401) + + # ==================== 参数验证测试 ==================== + + def test_missing_wecom_user_id(self): + """缺少 wecom_user_id 应返回 400""" + resp = self._post({}) + self.assertEqual(resp.status_code, 400) + self.assertIn('wecom_user_id', resp.json()) + + # ==================== 登录成功测试 ==================== + + def test_login_success(self): + """正常登录应返回 200 和 JWT token""" + resp = self._post({'wecom_user_id': 'wx_user_001'}) + self.assertEqual(resp.status_code, 200) + data = resp.json() + self.assertIn('access', data) + self.assertIn('refresh', data) + self.assertEqual(data['employee_id'], self.employee.id) + self.assertEqual(data['employee_name'], '测试员工') + self.assertEqual(data['username'], 'testuser') + # access token 应该是非空字符串 + self.assertTrue(len(data['access']) > 0) + self.assertTrue(len(data['refresh']) > 0) + + def test_login_token_is_valid_jwt(self): + """返回的 token 应该可以用于认证后续请求""" + resp = self._post({'wecom_user_id': 'wx_user_001'}) + self.assertEqual(resp.status_code, 200) + access_token = resp.json()['access'] + + # 使用 token 访问需要认证的接口 + self.client.credentials(HTTP_AUTHORIZATION=f'Bearer {access_token}') + # 访问 health check 不需要特殊权限,但需要认证 + health_resp = self.client.get('/api/v2/health/') + # health check 应该返回 200(不管认证方式) + self.assertEqual(health_resp.status_code, 200) + + # ==================== 错误场景测试 ==================== + + def test_unbound_wecom_user_id(self): + """未绑定的 wecom_user_id 应返回 404""" + resp = self._post({'wecom_user_id': 'wx_unknown_user'}) + self.assertEqual(resp.status_code, 404) + self.assertEqual(resp.json()['error'], '该企业微信用户未绑定本系统员工') + + def test_employee_without_sys_user(self): + """员工未关联系统用户应返回 404""" + # 创建一个没有 sys_user 的员工 + employee_no_user = Employee.objects.create( + merchant=self.merchant, + name='无用户员工', + wecom_user_id='wx_no_user', + ) + resp = self._post({'wecom_user_id': 'wx_no_user'}) + self.assertEqual(resp.status_code, 404) + self.assertEqual(resp.json()['error'], '该员工未关联系统用户') + + def test_inactive_user(self): + """被禁用的用户应返回 403""" + self.user.is_active = False + self.user.save() + resp = self._post({'wecom_user_id': 'wx_user_001'}) + self.assertEqual(resp.status_code, 403) + self.assertEqual(resp.json()['error'], '该用户已被禁用') + + # ==================== 边界测试 ==================== + + def test_multiple_logins_return_different_tokens(self): + """多次登录应返回不同的 token(每次生成新 token)""" + resp1 = self._post({'wecom_user_id': 'wx_user_001'}) + resp2 = self._post({'wecom_user_id': 'wx_user_001'}) + self.assertEqual(resp1.status_code, 200) + self.assertEqual(resp2.status_code, 200) + # refresh token 每次应该不同 + self.assertNotEqual(resp1.json()['refresh'], resp2.json()['refresh']) diff --git a/api_v2/urls.py b/api_v2/urls.py index b00bb2f..2a19747 100644 --- a/api_v2/urls.py +++ b/api_v2/urls.py @@ -41,6 +41,7 @@ from api_v2.views import ( AgentMesProductionAssignmentListView, ) from api_v2.views.basic_info import CustomerEmployeeBindingView, MyVisiblePagesView +from api_v2.views.wecom import WecomBindingView, WecomLoginView urlpatterns = [ path('health/', HealthCheckView.as_view(), name='api_v2_health_check'), @@ -83,4 +84,6 @@ urlpatterns = [ path('mes/production-assignments//', ProductionAssignmentDetailView.as_view(), name='api_v2_mes_production_assignment_detail'), path('shipment-delivery-photos/', ShipmentDeliveryPhotoListCreateView.as_view(), name='api_v2_shipment_delivery_photo_list_create'), path('shipment-delivery-photos//', ShipmentDeliveryPhotoDetailView.as_view(), name='api_v2_shipment_delivery_photo_detail'), + path('wecom/binduser/', WecomBindingView.as_view(), name='api_v2_wecom_binduser'), + path('wecom/login/', WecomLoginView.as_view(), name='api_v2_wecom_login'), ] diff --git a/api_v2/views/__init__.py b/api_v2/views/__init__.py index 623cb62..770d896 100644 --- a/api_v2/views/__init__.py +++ b/api_v2/views/__init__.py @@ -47,6 +47,7 @@ from .ai import ( AgentMesProductionAssignmentListView, ) from .shipment_delivery_photo import ShipmentDeliveryPhotoDetailView, ShipmentDeliveryPhotoListCreateView +from .wecom import WecomBindingView, WecomLoginView __all__ = [ 'HealthCheckView', diff --git a/api_v2/views/wecom.py b/api_v2/views/wecom.py new file mode 100644 index 0000000..3a8a980 --- /dev/null +++ b/api_v2/views/wecom.py @@ -0,0 +1,212 @@ +""" +企业微信 API + +供外部企业微信服务端程序调用: +1. 绑定接口:将企业微信 user_id 与本系统 Employee 绑定 +2. 登录接口:通过已绑定的 wecom_user_id 快捷获取 JWT token + +认证方式:AGENT_ACCESS_KEY(通过 Authorization 头传入)。 +""" + +import logging + +from django.contrib.auth import authenticate +from rest_framework import serializers, status +from rest_framework.permissions import IsAuthenticated +from rest_framework.response import Response +from rest_framework.views import APIView +from rest_framework_simplejwt.tokens import RefreshToken + +from api_v2.views.ai import AgentAccessKeyAuthentication +from basic_info.models import Employee + +logger = logging.getLogger(__name__) + + +class WecomBindingRequestSerializer(serializers.Serializer): + wecom_user_id = serializers.CharField( + max_length=128, + help_text='企业微信 OAuth2 返回的 user_id', + ) + username = serializers.CharField( + max_length=150, + help_text='本系统用户名', + ) + password = serializers.CharField( + max_length=128, + help_text='本系统用户密码', + ) + force_update = serializers.BooleanField( + default=False, + required=False, + help_text='如果该 Employee 已绑定其他企业微信用户,是否强制覆盖', + ) + + +class WecomBindingView(APIView): + """ + 企业微信用户绑定接口。 + + POST /api/v2/wecom/binduser/ + + 请求体: + - wecom_user_id: 企业微信用户 ID + - username: 本系统用户名 + - password: 本系统用户密码 + - force_update: (可选) 是否强制覆盖已有绑定,默认 false + + 认证:Authorization 头传入 AGENT_ACCESS_KEY + """ + + authentication_classes = [AgentAccessKeyAuthentication] + permission_classes = [IsAuthenticated] + + def post(self, request): + ser = WecomBindingRequestSerializer(data=request.data) + ser.is_valid(raise_exception=True) + + wecom_user_id = ser.validated_data['wecom_user_id'] + username = ser.validated_data['username'] + password = ser.validated_data['password'] + force_update = ser.validated_data.get('force_update', False) + + # 1. 验证本系统用户凭据 + user = authenticate(username=username, password=password) + if user is None: + return Response( + {'error': '用户名或密码错误'}, + status=status.HTTP_401_UNAUTHORIZED, + ) + + if not user.is_active: + return Response( + {'error': '该用户已被禁用'}, + status=status.HTTP_403_FORBIDDEN, + ) + + # 2. 查找关联的 Employee + employee = Employee.objects.filter(sys_user=user).first() + if employee is None: + return Response( + {'error': '该用户未关联任何员工记录'}, + status=status.HTTP_404_NOT_FOUND, + ) + + # 3. 检查该 wecom_user_id 是否已被其他 Employee 绑定 + existing = Employee.objects.filter(wecom_user_id=wecom_user_id).exclude(id=employee.id).first() + if existing is not None: + return Response( + {'error': f'该企业微信用户已绑定到其他员工: {existing.name} (ID: {existing.id})'}, + status=status.HTTP_409_CONFLICT, + ) + + # 4. 检查该 Employee 是否已有不同的 wecom_user_id + if employee.wecom_user_id and employee.wecom_user_id != wecom_user_id: + if not force_update: + return Response( + { + 'error': '该员工已绑定其他企业微信用户', + 'current_wecom_user_id': employee.wecom_user_id, + 'hint': '如需覆盖,请传入 force_update=true', + }, + status=status.HTTP_409_CONFLICT, + ) + + # 5. 执行绑定 + employee.wecom_user_id = wecom_user_id + employee.save(update_fields=['wecom_user_id', 'updated_at']) + + logger.info( + 'WeChat Work binding success: wecom_user_id=%s -> employee=%s (ID: %s)', + wecom_user_id, + employee.name, + employee.id, + ) + + return Response( + { + 'message': '绑定成功', + 'employee_id': employee.id, + 'employee_name': employee.name, + 'wecom_user_id': wecom_user_id, + }, + status=status.HTTP_200_OK, + ) + + +class WecomLoginRequestSerializer(serializers.Serializer): + wecom_user_id = serializers.CharField( + max_length=128, + help_text='企业微信 OAuth2 返回的 user_id', + ) + + +class WecomLoginView(APIView): + """ + 企业微信快捷登录接口。 + + POST /api/v2/wecom/login/ + + 通过已绑定的 wecom_user_id 获取本系统的 JWT token, + 无需再次输入用户名密码。 + + 请求体: + - wecom_user_id: 企业微信用户 ID(必须已通过绑定接口完成绑定) + + 认证:Authorization 头传入 AGENT_ACCESS_KEY + """ + + authentication_classes = [AgentAccessKeyAuthentication] + permission_classes = [IsAuthenticated] + + def post(self, request): + ser = WecomLoginRequestSerializer(data=request.data) + ser.is_valid(raise_exception=True) + + wecom_user_id = ser.validated_data['wecom_user_id'] + + # 1. 通过 wecom_user_id 查找 Employee + employee = Employee.objects.select_related('sys_user').filter( + wecom_user_id=wecom_user_id, + ).first() + + if employee is None: + return Response( + {'error': '该企业微信用户未绑定本系统员工'}, + status=status.HTTP_404_NOT_FOUND, + ) + + # 2. 检查关联的系统用户 + user = employee.sys_user + if user is None: + return Response( + {'error': '该员工未关联系统用户'}, + status=status.HTTP_404_NOT_FOUND, + ) + + if not user.is_active: + return Response( + {'error': '该用户已被禁用'}, + status=status.HTTP_403_FORBIDDEN, + ) + + # 3. 生成 JWT token + refresh = RefreshToken.for_user(user) + + logger.info( + 'WeChat Work login success: wecom_user_id=%s -> user=%s (employee=%s)', + wecom_user_id, + user.username, + employee.name, + ) + + return Response( + { + 'access': str(refresh.access_token), + 'refresh': str(refresh), + 'employee_id': employee.id, + 'employee_name': employee.name, + 'username': user.username, + }, + status=status.HTTP_200_OK, + ) diff --git a/basic_info/migrations/0027_employee_wecom_user_id.py b/basic_info/migrations/0027_employee_wecom_user_id.py new file mode 100644 index 0000000..1dcca5e --- /dev/null +++ b/basic_info/migrations/0027_employee_wecom_user_id.py @@ -0,0 +1,18 @@ +# Generated by Django 5.2.8 on 2026-05-20 09:33 + +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ('basic_info', '0026_transportvehicle_and_capacity'), + ] + + operations = [ + migrations.AddField( + model_name='employee', + name='wecom_user_id', + field=models.CharField(blank=True, help_text='企业微信 OAuth2 返回的 user_id', max_length=128, null=True, unique=True, verbose_name='企业微信用户ID'), + ), + ] diff --git a/basic_info/models.py b/basic_info/models.py index 122c440..8a55d01 100644 --- a/basic_info/models.py +++ b/basic_info/models.py @@ -479,6 +479,14 @@ class Employee(ModelBase): default=EmployeeStatusEnum.ACTIVE, verbose_name='员工状态', ) + wecom_user_id = models.CharField( + max_length=128, + unique=True, + null=True, + blank=True, + verbose_name='企业微信用户ID', + help_text='企业微信 OAuth2 返回的 user_id', + ) @property def job_type(self) -> str: diff --git a/business/external_finance_sync.py b/business/external_finance_sync.py index 97416dd..e56e101 100644 --- a/business/external_finance_sync.py +++ b/business/external_finance_sync.py @@ -127,6 +127,10 @@ class HaoBuYeFinanceClient: payload = response.json() except ValueError: payload = {} + # 404 且返回了 status=not_found 的 JSON 视为正常业务响应(如客户无退货单), + # 由调用方根据 status 字段决定后续逻辑,而不是在此处抛异常中断整个同步。 + if response.status_code == 404 and isinstance(payload, dict) and payload.get('status') == 'not_found': + return payload message = payload.get('message') or payload.get('error') or f'外部接口异常状态码: {response.status_code}' raise ExternalFinanceSyncError(str(message)) try: @@ -775,12 +779,18 @@ def _apply_sale_discount_zk( progress_callback: Callable[[str], None] | None = None, ) -> int: """ - 从 sale_discount 数据(F_Skd 的 XS% 行)中提取 ZkJinE, - 更新到对应的 ExternalCustomerStatementOrder.zk_amount。 + 从 sale_discount 数据(F_Skd 的 XS% 行)中提取 ZkJinE 和 YfJinE, + 更新到对应的 ExternalCustomerStatementOrder.zk_amount 和 total_amount。 + + 关键修正:I_Sale.SfJinE 始终为 0(数据源缺陷),导致从 I_Sale.JinE 计算的 + total_amount 未扣除实付部分。F_Skd XS% 的 YfJinE 已经扣除了实付,是正确的 + 欠款基数,因此用它覆盖 total_amount。 + 返回更新的记录数。 """ - # 按 BianHaoID 聚合折扣(一个 BianHaoID 可能有多条 F_Skd 记录) + # 按 BianHaoID 聚合折扣和应付金额 zk_by_source: dict[str, Decimal] = {} + yf_by_source: dict[str, Decimal] = {} for record in sale_discounts: if not isinstance(record, dict): continue @@ -788,13 +798,31 @@ def _apply_sale_discount_zk( if not source_id: continue zk_value = abs(_to_decimal(record.get('ZkJinE'), field_name='ZkJinE', default=Decimal('0'))) + # YfJinE 不取 abs():负值代表退款/冲减,应保留原始正负 + yf_value = _to_decimal(record.get('YfJinE'), field_name='YfJinE', default=Decimal('0')) if zk_value: zk_by_source[source_id] = zk_by_source.get(source_id, Decimal('0')) + zk_value + # YfJinE 是每个销售单的正确欠款金额(已扣除实付 SfJinE) + # 一个 BianHaoID 通常只有一条 F_Skd XS% 记录,但为安全起见做聚合 + # 注意:不过滤零值,因为 YfJinE=0 也是有效数据(表示全额实付) + yf_by_source[source_id] = yf_by_source.get(source_id, Decimal('0')) + yf_value - if not zk_by_source: + if not zk_by_source and not yf_by_source: return 0 updated_count = 0 + + # 更新 total_amount:用 F_Skd XS%.YfJinE 覆盖从 I_Sale.JinE 计算的值 + for source_id, yf_total in yf_by_source.items(): + updated = business_models.ExternalCustomerStatementOrder.objects.filter( + merchant=merchant, + customer=customer, + category=business_models.ExternalCustomerStatementCategoryEnum.SALE, + external_source_id=source_id, + ).exclude(total_amount=yf_total).update(total_amount=yf_total) + updated_count += updated + + # 更新 zk_amount for source_id, zk_total in zk_by_source.items(): updated = business_models.ExternalCustomerStatementOrder.objects.filter( merchant=merchant, @@ -807,7 +835,7 @@ def _apply_sale_discount_zk( if updated_count: _emit_progress( progress_callback, - f'销售折扣更新: {updated_count} 条记录的 zk_amount 已从 sale_discount 数据更新', + f'销售折扣/应付更新: {updated_count} 条记录已从 sale_discount 数据更新 (total_amount via YfJinE, zk_amount via ZkJinE)', ) return updated_count diff --git a/business/tasks.py b/business/tasks.py index 4c76df3..ed2498a 100644 --- a/business/tasks.py +++ b/business/tasks.py @@ -1,9 +1,12 @@ +import logging from typing import Any, Dict, List from celery import shared_task from business import services as business_services +logger = logging.getLogger(__name__) + @shared_task(bind=True) def create_purchase_order_stock_entries( @@ -24,6 +27,80 @@ def create_purchase_order_stock_entries( return payload +@shared_task(bind=True) +def sync_external_customer_finance_scheduled(self) -> Dict[str, Any]: + """ + 定时任务:同步指定客户列表的外部财务(对账)数据。 + + 每天凌晨 5 点由 celery beat 触发,逐个客户调用 sync_customer_finance()。 + 单个客户失败时跳过并记录错误,不影响其他客户。 + """ + from django.conf import settings + from basic_info.models import Employee + from business.external_finance_sync import sync_customer_finance + + customer_names = getattr(settings, 'FINANCE_SYNC_CUSTOMER_NAMES', []) + operator_id = getattr(settings, 'HAOBUYE_FINANCE_SYNC_OPERATOR_ID', 0) + + if not customer_names: + logger.warning('[finance_sync_scheduled] FINANCE_SYNC_CUSTOMER_NAMES 为空,跳过') + return {'status': 'skipped', 'reason': 'empty_customer_list'} + + if not operator_id: + logger.error('[finance_sync_scheduled] HAOBUYE_FINANCE_SYNC_OPERATOR_ID 未配置') + return {'status': 'error', 'reason': 'missing_operator_id'} + + try: + operator = Employee.objects.select_related('merchant').get(id=operator_id) + except Employee.DoesNotExist: + logger.error('[finance_sync_scheduled] 经办人不存在: id=%s', operator_id) + return {'status': 'error', 'reason': f'operator_not_found:{operator_id}'} + + results: List[Dict[str, Any]] = [] + errors: List[Dict[str, str]] = [] + + for customer_name in customer_names: + try: + logger.info('[finance_sync_scheduled] 开始同步客户: %s', customer_name) + result = sync_customer_finance( + customer_name=customer_name, + operator=operator, + allow_create_customer=True, + force_update=True, + ) + results.append({ + 'customer_name': customer_name, + 'created_count': result.get('created_count', 0), + 'skipped_existing_count': result.get('skipped_existing_count', 0), + 'external_business_created_count': result.get('external_business_created_count', 0), + }) + logger.info( + '[finance_sync_scheduled] 客户 %s 同步完成: created=%s, skipped=%s', + customer_name, + result.get('created_count', 0), + result.get('skipped_existing_count', 0), + ) + except Exception as exc: + logger.exception('[finance_sync_scheduled] 客户 %s 同步失败', customer_name) + errors.append({'customer_name': customer_name, 'error': str(exc)}) + + summary = { + 'status': 'completed', + 'total_customers': len(customer_names), + 'success_count': len(results), + 'error_count': len(errors), + 'results': results, + 'errors': errors, + } + logger.info( + '[finance_sync_scheduled] 全部完成: success=%s, errors=%s', + len(results), + len(errors), + ) + return summary + + + @shared_task(bind=True) def create_sales_order_stock_entries( self, diff --git a/docs/api_v2_wecom_binduser.md b/docs/api_v2_wecom_binduser.md new file mode 100644 index 0000000..3690f62 --- /dev/null +++ b/docs/api_v2_wecom_binduser.md @@ -0,0 +1,267 @@ +# 企业微信用户绑定 API + +## 接口说明 + +供企业微信服务端程序调用,将企业微信 OAuth2 获取的 `user_id` 与本系统的员工(Employee)绑定。调用方需提供本系统的用户名和密码进行身份验证,验证通过后完成绑定。 + +## 接口信息 + +- **URL**: `/api/v2/wecom/binduser/` +- **方法**: `POST` +- **Content-Type**: `application/json` +- **认证方式**: 固定密钥,通过 `Authorization` 请求头传入(不使用 Bearer 前缀) + +## 认证 + +请求必须在 `Authorization` 头中携带预分配的访问密钥(由系统管理员提供)。 + +``` +Authorization: your-access-key +``` + +## 请求参数 + +| 参数名 | 类型 | 必填 | 说明 | +|--------|------|------|------| +| wecom_user_id | string | 是 | 企业微信 OAuth2 返回的 user_id,最长 128 字符 | +| username | string | 是 | 本系统用户名 | +| password | string | 是 | 本系统用户密码 | +| force_update | boolean | 否 | 如果该员工已绑定其他企业微信用户,是否强制覆盖。默认 `false` | + +## 请求示例 + +```bash +curl -X POST "https://your-domain/api/v2/wecom/binduser/" \ + -H "Authorization: your-access-key" \ + -H "Content-Type: application/json" \ + -d '{ + "wecom_user_id": "zhangsan", + "username": "zhangsan", + "password": "mypassword123" + }' +``` + +## 响应说明 + +### 绑定成功(200 OK) + +```json +{ + "message": "绑定成功", + "employee_id": 42, + "employee_name": "张三", + "wecom_user_id": "zhangsan" +} +``` + +### 用户名或密码错误(401 Unauthorized) + +```json +{ + "error": "用户名或密码错误" +} +``` + +### 用户已被禁用(403 Forbidden) + +```json +{ + "error": "该用户已被禁用" +} +``` + +### 用户未关联员工(404 Not Found) + +```json +{ + "error": "该用户未关联任何员工记录" +} +``` + +### 企业微信用户已绑定到其他员工(409 Conflict) + +```json +{ + "error": "该企业微信用户已绑定到其他员工: 李四 (ID: 15)" +} +``` + +### 员工已绑定其他企业微信用户(409 Conflict) + +当该员工已有不同的 `wecom_user_id` 且未传 `force_update=true` 时: + +```json +{ + "error": "该员工已绑定其他企业微信用户", + "current_wecom_user_id": "old_wecom_id", + "hint": "如需覆盖,请传入 force_update=true" +} +``` + +### 认证失败(401 Unauthorized) + +Authorization 头缺失或密钥错误: + +```json +{ + "detail": "AGENT_ACCESS_KEY 无效" +} +``` + +### 参数验证失败(400 Bad Request) + +```json +{ + "wecom_user_id": ["该字段是必填项。"], + "username": ["该字段是必填项。"], + "password": ["该字段是必填项。"] +} +``` + +## 业务规则 + +1. **一对一约束**:一个企业微信 user_id 只能绑定一个员工,一个员工只能绑定一个企业微信 user_id。 +2. **幂等性**:重复绑定相同的 wecom_user_id 到同一个员工不会报错,返回 200。 +3. **冲突处理**: + - 如果 `wecom_user_id` 已被其他员工绑定 → 拒绝(409) + - 如果该员工已绑定不同的 `wecom_user_id` 且 `force_update=false` → 拒绝(409) + - 如果该员工已绑定不同的 `wecom_user_id` 且 `force_update=true` → 覆盖绑定 +4. **身份验证**:通过 Django 内置认证机制验证用户名和密码,验证通过后通过 `User → Employee` 关系找到对应员工。 + +## 接入流程 + +``` +企业微信用户 → OAuth2 授权 → 获取 wecom user_id + ↓ + 调用本接口(传入 wecom_user_id + 本系统 username/password) + ↓ + 验证凭据 → 查找 Employee → 写入 wecom_user_id + ↓ + 绑定完成,后续可通过 wecom_user_id 识别用户身份 +``` + +## 注意事项 + +1. 密码通过 HTTPS 传输,请确保生产环境使用 TLS。 +2. 该接口不限制调用频率,但建议调用方做适当限流。 +3. 绑定成功后,后续业务可通过 `Employee.wecom_user_id` 字段查找对应员工。 + +--- + +# 企业微信快捷登录 API + +## 接口说明 + +供企业微信服务端程序调用,通过已绑定的 `wecom_user_id` 快捷获取本系统的 JWT token,无需再次输入用户名密码。 + +前提条件:该 `wecom_user_id` 必须已通过绑定接口(`/api/v2/wecom/binduser/`)完成绑定。 + +## 接口信息 + +- **URL**: `/api/v2/wecom/login/` +- **方法**: `POST` +- **Content-Type**: `application/json` +- **认证方式**: 固定密钥,通过 `Authorization` 请求头传入(与绑定接口相同) + +## 认证 + +``` +Authorization: your-access-key +``` + +## 请求参数 + +| 参数名 | 类型 | 必填 | 说明 | +|--------|------|------|------| +| wecom_user_id | string | 是 | 企业微信 OAuth2 返回的 user_id | + +## 请求示例 + +```bash +curl -X POST "https://your-domain/api/v2/wecom/login/" \ + -H "Authorization: your-access-key" \ + -H "Content-Type: application/json" \ + -d '{ + "wecom_user_id": "zhangsan" + }' +``` + +## 响应说明 + +### 登录成功(200 OK) + +```json +{ + "access": "eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9...", + "refresh": "eyJ0eXAiOiJKV1QiLCJhbGciOiJIUzI1NiJ9...", + "employee_id": 42, + "employee_name": "张三", + "username": "zhangsan" +} +``` + +返回的 `access` token 可直接用于后续业务 API 的 `Authorization: Bearer ` 认证。 + +### 企业微信用户未绑定(404 Not Found) + +```json +{ + "error": "该企业微信用户未绑定本系统员工" +} +``` + +### 员工未关联系统用户(404 Not Found) + +```json +{ + "error": "该员工未关联系统用户" +} +``` + +### 用户已被禁用(403 Forbidden) + +```json +{ + "error": "该用户已被禁用" +} +``` + +### 认证失败(401 Unauthorized) + +Authorization 头缺失或密钥错误: + +```json +{ + "detail": "AGENT_ACCESS_KEY 无效" +} +``` + +### 参数验证失败(400 Bad Request) + +```json +{ + "wecom_user_id": ["该字段是必填项。"] +} +``` + +## 业务规则 + +1. 该接口仅对已完成绑定的 `wecom_user_id` 有效。 +2. 每次调用都会生成新的 JWT token 对(access + refresh)。 +3. access token 有效期为 7 天(与系统 JWT 配置一致)。 +4. 安全边界依赖 `AGENT_ACCESS_KEY`,请妥善保管该密钥。 + +## 典型使用流程 + +``` +用户首次使用: + 企业微信 OAuth2 → 获取 wecom_user_id + → 调用 /api/v2/wecom/binduser/(传 wecom_user_id + username + password) + → 绑定完成 + +用户后续登录: + 企业微信 OAuth2 → 获取 wecom_user_id + → 调用 /api/v2/wecom/login/(仅传 wecom_user_id) + → 获取 JWT token + → 使用 token 访问业务 API +``` diff --git a/docs/statements_pagination_performance_evaluation.md b/docs/statements_pagination_performance_evaluation.md new file mode 100644 index 0000000..4e51767 --- /dev/null +++ b/docs/statements_pagination_performance_evaluation.md @@ -0,0 +1,201 @@ +# 客户对账单 API 分页与性能评估 + +本文档评估当前 `/api/v1/customers//statements/` 接口的分页能力与性能瓶颈,供前后端共同讨论优化方案。 + +--- + +## 1. 当前现状 + +### 接口行为 + +- 每次请求返回**该客户的全部对账记录**(销售单 + 退货单 + 收款单 + 外部业务依据) +- 无分页参数,无日期过滤 +- 排序固定:`occurred_at → recorded_at → source_id` 倒序 +- 每条记录包含 `cumulative_amount`(滚动累计)和 `arrears_amount`(实时欠款) + +### 数据规模参考 + +| 客户 | 销售单 | 退货单 | 收款单 | 总记录数 | 响应体积(估算) | +|------|--------|--------|--------|----------|-----------------| +| 木棉 | ~200 | 1 | ~50 | ~250 | ~100KB | +| 腾飞纺织 | ~1,500 | 9 | ~92 | ~1,600 | ~800KB | +| 张晓鹏 | ~3,100 | 6 | ~219 | ~3,350 | ~2MB | + +--- + +## 2. 为什么当前无法直接分页 + +核心原因:**滚动累计字段(`cumulative_amount` / `arrears_amount`)依赖前序所有记录的计算结果。** + +``` +结欠[n] = 结欠[n-1] + 本行应收 - 本行已收 +``` + +要展示第 N 页的 `arrears_amount`,必须先计算前 N-1 页所有记录的累加值。这意味着: +- 不能简单用 DB 的 `OFFSET/LIMIT` +- 不能跳页 +- 后端必须从第一条开始逐条计算 + +--- + +## 3. 性能瓶颈分析 + +每次 API 调用的执行路径: + +``` +1. 4 次 DB 查询(sales / returns / receipts / external_statements) + └── 每次都 prefetch_related('items__product') 加载全部明细 +2. Python 内存中合并 + 排序全部记录 +3. 逐条遍历计算 running totals +4. 全量序列化为 JSON(含 items 明细数组) +5. HTTP 响应传输 +``` + +对于张晓鹏(3350 条记录): +- DB 查询:~200ms(4 次查询 + prefetch) +- Python 排序 + 计算:~50ms +- JSON 序列化:~300ms +- 网络传输(2MB):取决于带宽 + +**`StatementRecordView`(单条查询)更严重**:为了查 1 条记录,构建了完整对账单再过滤。 + +--- + +## 4. 优化方案对比 + +### 方案 A:游标分页(推荐短期方案) + +**原理**:前端传 `page_size` + 上一页最后一条的 `cumulative_amount` 作为 cursor,后端从 cursor 继续累加。 + +``` +GET /customers/6/statements/?page_size=50 +GET /customers/6/statements/?page_size=50&cursor=eyJjdW11bGF0aXZlIjoiMzE4OTEuMDAiLCJsYXN0X2lkIjoxODE0MX0= +``` + +| 优点 | 缺点 | +|------|------| +| 减少响应体积 | 后端仍需全量查询(但只序列化一页) | +| 前端按需加载 | 不能跳页,只能顺序翻页 | +| 向后兼容(不传参数 = 全量) | 需要前端配合改造 | + +**后端改动**:中等。排序 + running total 计算后,只返回 cursor 之后的 N 条。 + +--- + +### 方案 B:日期范围过滤 + +**原理**:前端传 `date_from` / `date_to`,后端只查询范围内的记录。 + +``` +GET /customers/6/statements/?date_from=2026-01-01&date_to=2026-05-18 +``` + +| 优点 | 缺点 | +|------|------| +| 实现简单 | running total 仍需从历史第一条开始算 | +| 减少返回数据量 | 或者放弃 running total 的准确性 | +| 前端可做"按年/按月"切换 | | + +**后端改动**:低。在 DB 查询加 `occurred_at` 过滤即可。但 `cumulative_amount` 需要决定: +- 选项 1:从第一条算到 date_to(准确但慢) +- 选项 2:只在返回范围内累计(快但不连续) + +--- + +### 方案 C:延迟加载 items 明细 + +**原理**:默认不返回 `items` 数组,前端需要时单独请求。 + +``` +GET /customers/6/statements/ → 不含 items +GET /customers/6/statements/?include_items=true → 含 items(当前行为) +GET /statements/record/?...&include_items=true → 单条含 items +``` + +| 优点 | 缺点 | +|------|------| +| 响应体积减少 50%+ | 前端需要额外请求获取明细 | +| 后端改动极小 | 如果前端表格需要展开明细,交互变复杂 | +| 完全向后兼容 | | + +**后端改动**:极低。序列化时根据参数决定是否包含 items。 + +--- + +### 方案 D:预计算快照(长期方案) + +**原理**:同步时预计算每条记录的 `cumulative_amount` / `arrears_amount` 存入 DB,查询时直接分页。 + +| 优点 | 缺点 | +|------|------| +| 真正的 DB 级分页 | 需要新表或新字段 | +| 查询性能最优 | 数据一致性维护复杂(任何单据变动需重算) | +| 支持跳页 | 开发成本高 | + +**后端改动**:高。需要设计快照表 + 触发重算机制。 + +--- + +### 方案 E:响应缓存 + +**原理**:对账单结果缓存 N 秒(或按数据版本号缓存),重复请求直接返回。 + +| 优点 | 缺点 | +|------|------| +| 零前端改动 | 数据有延迟(缓存过期前看不到最新) | +| 后端改动极小 | 大客户首次请求仍然慢 | +| 对高频刷新场景效果显著 | | + +--- + +## 5. 推荐实施路径 + +``` +第一步(立即可做):方案 C — 默认不返回 items,减少 50%+ 响应体积 +第二步(短期):方案 B — 加日期范围过滤,前端做"按月/按年"切换 +第三步(中期):方案 A — 游标分页,前端改为滚动加载 +第四步(按需):方案 E — 缓存,应对高频刷新 +``` + +--- + +## 6. 需要前端确认的问题 + +1. **对账单表格是否需要一次展示全部记录?** 还是可以接受分页/滚动加载? +2. **items 明细是否默认展示?** 还是用户点击展开时才加载? +3. **是否需要"按月/按年"切换?** 如果需要,running total 是否可以只在当前范围内累计? +4. **`cumulative_amount` / `arrears_amount` 是否是必须字段?** 如果前端不使用这两个字段,分页就变得简单很多。 +5. **导出 Excel 功能是否需要全量数据?** 如果需要,导出可以走单独的异步接口。 + +--- + +## 7. 当前接口响应结构(供参考) + +```json +{ + "counterparty": 6, + "counterparty_name": "张晓鹏", + "records": [ + { + "source_type": "external_sales_order", + "source_label": "外部销售单", + "source_id": 18141, + "occurred_at": "2026-05-18", + "recorded_at": "2026-05-18T21:45:03Z", + "status": 2, + "status_label": "已同步", + "positive_amount": "2436.00", + "negative_amount": "0.00", + "cumulative_amount": "31891.00", + "current_balance": "780454.00", + "arrears_amount": "748563.00", + "remarks": "...", + "items": [ ... ] + } + ], + "summary": { + "positive_total": "26981509.80", + "negative_total": "26201055.80" + } +} +``` diff --git a/flower/settings.py b/flower/settings.py index 3da8f57..3cc54ba 100644 --- a/flower/settings.py +++ b/flower/settings.py @@ -98,7 +98,16 @@ AGENT_ACCESS_KEY = env('AGENT_ACCESS_KEY', default='') HAOBUYE_API_BASE_URL = env('HAOBUYE_API_BASE_URL', default='http://43.139.183.222:18080') HAOBUYE_API_AUTHORIZATION = env('HAOBUYE_API_AUTHORIZATION', default='your-fixed-authorization-secret') HAOBUYE_API_TIMEOUT_SECONDS = env.float('HAOBUYE_API_TIMEOUT_SECONDS', default=30.0) -HAOBUYE_FINANCE_SYNC_OPERATOR_ID = env.int('HAOBUYE_FINANCE_SYNC_OPERATOR_ID', default=0) +HAOBUYE_FINANCE_SYNC_OPERATOR_ID = env.int('HAOBUYE_FINANCE_SYNC_OPERATOR_ID', default=6) + +# 定时财务同步客户列表(临时需求,直接写死不走 env) +FINANCE_SYNC_CUSTOMER_NAMES: list[str] = [ + '曾念', '紫琪', '胡肖宇', '歌斯拉-胜利星厂', '胡鼎', + '罗标', '永琪服饰', '罗兵', '杜辉', '李群', '超麦汇大洋', + '展兴', '王杰敏', '锐木', '何玄', '周刚(周总)', '彤彤服饰', + '杰恩', '来发裁床', '刘静贤', '金源', '达达-金腾达qs', '周乔峰', + '杨锐辉', +] # CORS 配置 @@ -560,96 +569,101 @@ CELERY_TIMEZONE = TIME_ZONE from celery.schedules import crontab CELERY_BEAT_SCHEDULE = { - 'daily_database_backup': { - 'task': 'api_v1.tasks.backup_database', - 'schedule': crontab(hour=3, minute=0), - 'kwargs': { - 'output_dir': str(BASE_DIR / 'data-bak'), - 'filename_prefix': 'db-backup', - }, - }, - 'daily_settlement': { - 'task': 'settlement.tasks.run_daily_settlement', - 'schedule': crontab( - hour=SETTLEMENT_RUN_HOUR, - minute=SETTLEMENT_RUN_MINUTE - ), - }, - 'daily_plate_order_tiia_image_upload': { - 'task': 'printing.tasks.upload_yesterday_plate_order_images_to_tencent_tiia', - 'schedule': crontab(hour=3, minute=0), - }, - 'daily_mdy_plate_order_staging_tiia_image_upload': { - 'task': 'api_v1.tasks.upload_mdy_plate_order_staging_images_to_tencent_tiia', - 'schedule': crontab(hour=3, minute=20), - 'kwargs': { - 'batch_size': 500, - 'dry_run': False, - }, - }, - 'mdy_product_sync': { - 'task': 'api_v1.tasks.sync_mdy_products', - 'schedule': crontab(minute='*/10'), - 'kwargs': { - 'page_size': MDY_SYNC_PAGE_SIZE, - 'max_pages': MDY_SYNC_MAX_PAGES, - 'max_records': MDY_SYNC_MAX_RECORDS, - }, - }, - 'mdy_customer_sync': { - 'task': 'api_v1.tasks.sync_mdy_customers', - 'schedule': crontab(minute='*/5'), - 'kwargs': { - 'page_size': MDY_SYNC_PAGE_SIZE, - 'max_pages': MDY_SYNC_MAX_PAGES, - 'max_records': MDY_SYNC_MAX_RECORDS, - }, - }, - 'printing_external_records_sync': { - 'task': 'api_v1.tasks.sync_external_printing_records', - 'schedule': crontab(minute='*/5'), - 'kwargs': { - 'limit': 100, - }, - }, - 'notify_unreplied_missions_every_minute': { - 'task': 'notifier.tasks.notify_unreplied_missions_task', - 'schedule': crontab(minute='*'), - 'kwargs': { - 'limit': 100, - }, - }, - 'backfill_external_product_images_hourly': { - 'task': 'api_v1.tasks.backfill_external_product_images', - 'schedule': crontab(minute=17), - 'kwargs': { - 'limit': 500, - 'dry_run': False, - }, - }, - 'retry_external_printing_sync_failures_every_120_minutes': { - 'task': 'api_v1.tasks.retry_external_printing_sync_failures', - 'schedule': crontab(minute='0', hour='*/2'), - 'kwargs': { - 'limit': 100, - }, + # 'daily_database_backup': { + # 'task': 'api_v1.tasks.backup_database', + # 'schedule': crontab(hour=3, minute=0), + # 'kwargs': { + # 'output_dir': str(BASE_DIR / 'data-bak'), + # 'filename_prefix': 'db-backup', + # }, + # }, + # 'daily_settlement': { + # 'task': 'settlement.tasks.run_daily_settlement', + # 'schedule': crontab( + # hour=SETTLEMENT_RUN_HOUR, + # minute=SETTLEMENT_RUN_MINUTE + # ), + # }, + # 'daily_plate_order_tiia_image_upload': { + # 'task': 'printing.tasks.upload_yesterday_plate_order_images_to_tencent_tiia', + # 'schedule': crontab(hour=3, minute=0), + # }, + # 'daily_mdy_plate_order_staging_tiia_image_upload': { + # 'task': 'api_v1.tasks.upload_mdy_plate_order_staging_images_to_tencent_tiia', + # 'schedule': crontab(hour=3, minute=20), + # 'kwargs': { + # 'batch_size': 500, + # 'dry_run': False, + # }, + # }, + # 'mdy_product_sync': { + # 'task': 'api_v1.tasks.sync_mdy_products', + # 'schedule': crontab(minute='*/10'), + # 'kwargs': { + # 'page_size': MDY_SYNC_PAGE_SIZE, + # 'max_pages': MDY_SYNC_MAX_PAGES, + # 'max_records': MDY_SYNC_MAX_RECORDS, + # }, + # }, + # 'mdy_customer_sync': { + # 'task': 'api_v1.tasks.sync_mdy_customers', + # 'schedule': crontab(minute='*/5'), + # 'kwargs': { + # 'page_size': MDY_SYNC_PAGE_SIZE, + # 'max_pages': MDY_SYNC_MAX_PAGES, + # 'max_records': MDY_SYNC_MAX_RECORDS, + # }, + # }, + # 'printing_external_records_sync': { + # 'task': 'api_v1.tasks.sync_external_printing_records', + # 'schedule': crontab(minute='*/5'), + # 'kwargs': { + # 'limit': 100, + # }, + # }, + # 'notify_unreplied_missions_every_minute': { + # 'task': 'notifier.tasks.notify_unreplied_missions_task', + # 'schedule': crontab(minute='*'), + # 'kwargs': { + # 'limit': 100, + # }, + # }, + # 'backfill_external_product_images_hourly': { + # 'task': 'api_v1.tasks.backfill_external_product_images', + # 'schedule': crontab(minute=17), + # 'kwargs': { + # 'limit': 500, + # 'dry_run': False, + # }, + # }, + # 'retry_external_printing_sync_failures_every_120_minutes': { + # 'task': 'api_v1.tasks.retry_external_printing_sync_failures', + # 'schedule': crontab(minute='0', hour='*/2'), + # 'kwargs': { + # 'limit': 100, + # }, + # }, + # 外部财务同步:每天凌晨 5 点同步指定客户列表的对账数据 + 'sync_external_customer_finance_daily': { + 'task': 'business.tasks.sync_external_customer_finance_scheduled', + 'schedule': crontab(hour=5, minute=0), }, # 明道云开版:同步到“暂存表” # - 仅在 02:00-08:00 时间窗内持续运行(每 N 分钟触发一次) # - 通过 request_interval_seconds 控制单次任务对明道云 API 的请求节奏,避免超 QPS - 'mdy_plate_order_staging_sync': { - 'task': 'api_v1.tasks.sync_mdy_plate_orders', - # 02:00 - 07:59(不包含 08:xx) - 'schedule': crontab(minute='*/2', hour='2-7'), - 'kwargs': { - 'page_size': max(1, int(MDY_SYNC_PAGE_SIZE)), - 'max_pages': max(1, int(MDY_SYNC_MAX_PAGES)), - 'max_records': max(1, int(MDY_SYNC_MAX_RECORDS)), - 'with_related': True, - 'max_related_per_type': 5, - 'request_interval_seconds': float(MDY_PLATE_ORDER_SYNC_REQUEST_INTERVAL_SECONDS), - }, - }, + # 'mdy_plate_order_staging_sync': { + # 'task': 'api_v1.tasks.sync_mdy_plate_orders', + # # 02:00 - 07:59(不包含 08:xx) + # 'schedule': crontab(minute='*/2', hour='2-7'), + # 'kwargs': { + # 'page_size': max(1, int(MDY_SYNC_PAGE_SIZE)), + # 'max_pages': max(1, int(MDY_SYNC_MAX_PAGES)), + # 'max_records': max(1, int(MDY_SYNC_MAX_RECORDS)), + # 'with_related': True, + # 'max_related_per_type': 5, + # 'request_interval_seconds': float(MDY_PLATE_ORDER_SYNC_REQUEST_INTERVAL_SECONDS), + # }, + # }, } if MDY_ISOLATED_MERCHANT_ID and MDY_ISOLATED_MERCHANT_ID > 0: