| from __future__ import annotations |
|
|
| import json |
| import threading |
| import time |
| from collections.abc import Callable |
| from datetime import datetime |
| from pathlib import Path |
| from typing import Any |
|
|
| from services.config import DATA_DIR, config |
| from services.content_filter import request_text |
| from services.log_service import LOG_TYPE_CALL, log_service |
| from services.protocol import openai_v1_image_edit, openai_v1_image_generations |
| from services.time_utils import china_from_timestamp_string, china_now_string |
|
|
| TASK_STATUS_QUEUED = "queued" |
| TASK_STATUS_RUNNING = "running" |
| TASK_STATUS_SUCCESS = "success" |
| TASK_STATUS_ERROR = "error" |
| TERMINAL_STATUSES = {TASK_STATUS_SUCCESS, TASK_STATUS_ERROR} |
| UNFINISHED_STATUSES = {TASK_STATUS_QUEUED, TASK_STATUS_RUNNING} |
| TASK_CLEANUP_INTERVAL_SECONDS = 60 * 60 |
|
|
|
|
| def _now_iso() -> str: |
| return china_now_string() |
|
|
|
|
| def _timestamp(value: object) -> float: |
| if not isinstance(value, str) or not value.strip(): |
| return 0.0 |
| for fmt in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%dT%H:%M:%S.%f", "%Y-%m-%dT%H:%M:%S"): |
| try: |
| return datetime.strptime(value[:26], fmt).timestamp() |
| except ValueError: |
| continue |
| try: |
| return datetime.fromisoformat(value.replace("Z", "+00:00")).timestamp() |
| except Exception: |
| return 0.0 |
|
|
|
|
| def _clean(value: object, default: str = "") -> str: |
| return str(value or default).strip() |
|
|
|
|
| def _owner_id(identity: dict[str, object]) -> str: |
| return _clean(identity.get("id")) or "anonymous" |
|
|
|
|
| def _memory_conversation_id(value: object) -> str: |
| return _clean(value) or "default" |
|
|
|
|
| def _account_pool_strategy(value: object) -> str: |
| normalized = _clean(value).lower().replace("-", "_") |
| return normalized if normalized in {"own_first", "admin_first", "shared_first", "own_only", "admin_only", "shared_only"} else "own_first" |
|
|
|
|
| def _task_key(owner_id: str, task_id: str) -> str: |
| return f"{owner_id}:{task_id}" |
|
|
|
|
| def _collect_image_urls(data: list[Any]) -> list[str]: |
| urls: list[str] = [] |
| for item in data: |
| if isinstance(item, dict): |
| url = item.get("url") |
| if isinstance(url, str) and url: |
| urls.append(url) |
| return urls |
|
|
|
|
| def _count_image_items(data: object) -> int: |
| if not isinstance(data, list): |
| return 0 |
| count = 0 |
| for item in data: |
| if isinstance(item, dict) and (_clean(item.get("url")) or _clean(item.get("b64_json"))): |
| count += 1 |
| return count |
|
|
|
|
| def _collect_third_party_image_api(result: object) -> dict[str, Any]: |
| if isinstance(result, dict) and isinstance(result.get("_third_party_image_api"), dict): |
| return result["_third_party_image_api"] |
| return {} |
|
|
|
|
| def _public_task(task: dict[str, Any]) -> dict[str, Any]: |
| item = { |
| "id": task.get("id"), |
| "status": task.get("status"), |
| "mode": task.get("mode"), |
| "model": task.get("model"), |
| "size": task.get("size"), |
| "created_at": task.get("created_at"), |
| "started_at": task.get("started_at"), |
| "updated_at": task.get("updated_at"), |
| } |
| if task.get("progress"): |
| item["progress"] = task.get("progress") |
| if isinstance(task.get("progress_events"), list): |
| item["progress_events"] = [ |
| event |
| for event in task.get("progress_events", []) |
| if isinstance(event, dict) and _clean(event.get("step")) |
| ] |
| if task.get("duration_ms") is not None: |
| item["duration_ms"] = task.get("duration_ms") |
| if task.get("status") in UNFINISHED_STATUSES: |
| base_ts = task.get("started_ts") if task.get("status") == TASK_STATUS_RUNNING else task.get("created_ts") |
| if not base_ts: |
| base_ts = _timestamp(task.get("started_at") if task.get("status") == TASK_STATUS_RUNNING else task.get("created_at")) |
| if base_ts: |
| item["elapsed_secs"] = round(time.time() - float(base_ts), 1) |
| if task.get("data") is not None: |
| item["data"] = task.get("data") |
| if task.get("error"): |
| item["error"] = task.get("error") |
| return item |
|
|
|
|
| class ImageTaskService: |
| def __init__( |
| self, |
| path: Path, |
| *, |
| generation_handler: Callable[[dict[str, Any]], dict[str, Any]] = openai_v1_image_generations.handle, |
| edit_handler: Callable[[dict[str, Any]], dict[str, Any]] = openai_v1_image_edit.handle, |
| retention_days_getter: Callable[[], int] | None = None, |
| ): |
| self.path = path |
| self.generation_handler = generation_handler |
| self.edit_handler = edit_handler |
| self.retention_days_getter = retention_days_getter or (lambda: config.image_retention_days) |
| self._lock = threading.RLock() |
| self._tasks: dict[str, dict[str, Any]] = {} |
| self._last_cleanup_at = time.monotonic() |
| self.path.parent.mkdir(parents=True, exist_ok=True) |
| with self._lock: |
| self._tasks = self._load_locked() |
| changed = self._recover_unfinished_locked() |
| changed = self._cleanup_locked() or changed |
| if changed: |
| self._save_locked() |
|
|
| def submit_generation( |
| self, |
| identity: dict[str, object], |
| *, |
| client_task_id: str, |
| prompt: str, |
| model: str, |
| size: str | None, |
| base_url: str, |
| memory_enabled: bool = False, |
| memory_conversation_id: str = "", |
| memory_reset: bool = False, |
| account_pool_strategy: str = "own_first", |
| credit_user_id: str = "", |
| credit_amount: int = 0, |
| quality: str = "", |
| resolution: str = "", |
| third_party_image_group_id: str = "", |
| third_party_wallet_user_id: str = "", |
| third_party_wallet_group_id: str = "", |
| third_party_wallet_model_id: str = "", |
| third_party_wallet_amount: int = 0, |
| third_party_wallet_cost_per_image: int = 1, |
| third_party_wallet_fallback_credit_user_id: str = "", |
| third_party_wallet_fallback_credit_amount: int = 0, |
| ) -> dict[str, Any]: |
| payload = { |
| "prompt": prompt, |
| "model": model, |
| "n": 1, |
| "size": size, |
| "quality": quality, |
| "resolution": resolution, |
| "third_party_image_group_id": third_party_image_group_id, |
| "response_format": "url", |
| "base_url": base_url, |
| "account_owner_id": _owner_id(identity), |
| "account_pool_strategy": _account_pool_strategy(account_pool_strategy), |
| "account_allow_shared": not memory_enabled, |
| } |
| if memory_enabled: |
| payload["memory"] = True |
| payload["memory_owner_id"] = _owner_id(identity) |
| payload["memory_conversation_id"] = _memory_conversation_id(memory_conversation_id) |
| if memory_reset: |
| payload["memory_reset"] = True |
| return self._submit( |
| identity, |
| client_task_id=client_task_id, |
| mode="generate", |
| payload=payload, |
| credit_user_id=credit_user_id, |
| credit_amount=credit_amount, |
| third_party_wallet_user_id=third_party_wallet_user_id, |
| third_party_wallet_group_id=third_party_wallet_group_id, |
| third_party_wallet_model_id=third_party_wallet_model_id, |
| third_party_wallet_amount=third_party_wallet_amount, |
| third_party_wallet_cost_per_image=third_party_wallet_cost_per_image, |
| third_party_wallet_fallback_credit_user_id=third_party_wallet_fallback_credit_user_id, |
| third_party_wallet_fallback_credit_amount=third_party_wallet_fallback_credit_amount, |
| ) |
|
|
| def submit_edit( |
| self, |
| identity: dict[str, object], |
| *, |
| client_task_id: str, |
| prompt: str, |
| model: str, |
| size: str | None, |
| base_url: str, |
| images: list[tuple[bytes, str, str]], |
| memory_enabled: bool = False, |
| memory_conversation_id: str = "", |
| memory_reset: bool = False, |
| account_pool_strategy: str = "own_first", |
| credit_user_id: str = "", |
| credit_amount: int = 0, |
| quality: str = "", |
| resolution: str = "", |
| third_party_image_group_id: str = "", |
| third_party_wallet_user_id: str = "", |
| third_party_wallet_group_id: str = "", |
| third_party_wallet_model_id: str = "", |
| third_party_wallet_amount: int = 0, |
| third_party_wallet_cost_per_image: int = 1, |
| third_party_wallet_fallback_credit_user_id: str = "", |
| third_party_wallet_fallback_credit_amount: int = 0, |
| ) -> dict[str, Any]: |
| payload = { |
| "prompt": prompt, |
| "images": images, |
| "model": model, |
| "n": 1, |
| "size": size, |
| "quality": quality, |
| "resolution": resolution, |
| "third_party_image_group_id": third_party_image_group_id, |
| "response_format": "url", |
| "base_url": base_url, |
| "account_owner_id": _owner_id(identity), |
| "account_pool_strategy": _account_pool_strategy(account_pool_strategy), |
| "account_allow_shared": not memory_enabled, |
| } |
| if memory_enabled: |
| payload["memory"] = True |
| payload["memory_owner_id"] = _owner_id(identity) |
| payload["memory_conversation_id"] = _memory_conversation_id(memory_conversation_id) |
| if memory_reset: |
| payload["memory_reset"] = True |
| return self._submit( |
| identity, |
| client_task_id=client_task_id, |
| mode="edit", |
| payload=payload, |
| credit_user_id=credit_user_id, |
| credit_amount=credit_amount, |
| third_party_wallet_user_id=third_party_wallet_user_id, |
| third_party_wallet_group_id=third_party_wallet_group_id, |
| third_party_wallet_model_id=third_party_wallet_model_id, |
| third_party_wallet_amount=third_party_wallet_amount, |
| third_party_wallet_cost_per_image=third_party_wallet_cost_per_image, |
| third_party_wallet_fallback_credit_user_id=third_party_wallet_fallback_credit_user_id, |
| third_party_wallet_fallback_credit_amount=third_party_wallet_fallback_credit_amount, |
| ) |
|
|
| def list_tasks(self, identity: dict[str, object], task_ids: list[str]) -> dict[str, Any]: |
| owner = _owner_id(identity) |
| requested_ids = [_clean(task_id) for task_id in task_ids if _clean(task_id)] |
| with self._lock: |
| if self._cleanup_if_due_locked(): |
| self._save_locked() |
| items = [] |
| missing_ids = [] |
| for task_id in requested_ids: |
| task = self._tasks.get(_task_key(owner, task_id)) |
| if task is None: |
| missing_ids.append(task_id) |
| else: |
| items.append(_public_task(task)) |
| if not requested_ids: |
| items = [ |
| _public_task(task) |
| for task in self._tasks.values() |
| if task.get("owner_id") == owner |
| ] |
| items.sort(key=lambda item: str(item.get("updated_at") or ""), reverse=True) |
| missing_ids = [] |
| return {"items": items, "missing_ids": missing_ids} |
|
|
| def _submit( |
| self, |
| identity: dict[str, object], |
| *, |
| client_task_id: str, |
| mode: str, |
| payload: dict[str, Any], |
| credit_user_id: str = "", |
| credit_amount: int = 0, |
| third_party_wallet_user_id: str = "", |
| third_party_wallet_group_id: str = "", |
| third_party_wallet_model_id: str = "", |
| third_party_wallet_amount: int = 0, |
| third_party_wallet_cost_per_image: int = 1, |
| third_party_wallet_fallback_credit_user_id: str = "", |
| third_party_wallet_fallback_credit_amount: int = 0, |
| ) -> dict[str, Any]: |
| task_id = _clean(client_task_id) |
| if not task_id: |
| raise ValueError("client_task_id is required") |
| owner = _owner_id(identity) |
| key = _task_key(owner, task_id) |
| now = _now_iso() |
| should_start = False |
| with self._lock: |
| cleaned = self._cleanup_if_due_locked() |
| task = self._tasks.get(key) |
| if task is not None: |
| if cleaned: |
| self._save_locked() |
| return _public_task(task) |
| normalized_credit_user_id = _clean(credit_user_id) |
| normalized_credit_amount = max(0, int(credit_amount or 0)) |
| normalized_fallback_credit_user_id = _clean(third_party_wallet_fallback_credit_user_id) |
| normalized_fallback_credit_amount = max(0, int(third_party_wallet_fallback_credit_amount or 0)) |
| daily_reserved = 0 |
| paid_reserved = 0 |
| package_reservations: list[dict[str, object]] = [] |
| grant_reservations: list[dict[str, object]] = [] |
| normalized_wallet_user_id = _clean(third_party_wallet_user_id) |
| normalized_wallet_group_id = _clean(third_party_wallet_group_id) |
| normalized_wallet_model_id = _clean(third_party_wallet_model_id) |
| normalized_wallet_amount = max(0, int(third_party_wallet_amount or 0)) |
| normalized_wallet_cost_per_image = max(0, int(third_party_wallet_cost_per_image or 0)) |
| if normalized_wallet_user_id and normalized_wallet_group_id and normalized_wallet_model_id and normalized_wallet_amount > 0: |
| from services.auth_service import auth_service |
|
|
| try: |
| auth_service.reserve_third_party_image_wallet( |
| normalized_wallet_user_id, |
| group_id=normalized_wallet_group_id, |
| model_id=normalized_wallet_model_id, |
| amount=normalized_wallet_amount, |
| ) |
| except ValueError: |
| if not normalized_fallback_credit_user_id or normalized_fallback_credit_amount <= 0: |
| raise |
| normalized_wallet_user_id = "" |
| normalized_wallet_group_id = "" |
| normalized_wallet_model_id = "" |
| normalized_wallet_amount = 0 |
| normalized_wallet_cost_per_image = 0 |
| normalized_credit_user_id = normalized_fallback_credit_user_id |
| normalized_credit_amount = normalized_fallback_credit_amount |
| if normalized_credit_user_id and normalized_credit_amount > 0: |
| from services.auth_service import auth_service |
|
|
| reservation = auth_service.reserve_credits(normalized_credit_user_id, normalized_credit_amount).get("_credit_reservation", {}) |
| if not isinstance(reservation, dict): |
| reservation = {} |
| daily_reserved = max(0, int(reservation.get("daily") or 0)) |
| paid_reserved = max(0, int(reservation.get("permanent_paid") if "permanent_paid" in reservation else reservation.get("paid") or 0)) |
| package_reservations = reservation.get("packages") if isinstance(reservation.get("packages"), list) else [] |
| grant_reservations = reservation.get("grants") if isinstance(reservation.get("grants"), list) else [] |
| task = { |
| "id": task_id, |
| "owner_id": owner, |
| "status": TASK_STATUS_QUEUED, |
| "mode": mode, |
| "model": _clean(payload.get("model"), "gpt-image-2"), |
| "size": _clean(payload.get("size")), |
| "created_at": now, |
| "updated_at": now, |
| "created_ts": time.time(), |
| "updated_ts": time.time(), |
| "progress_events": [], |
| } |
| if normalized_credit_user_id and normalized_credit_amount > 0: |
| task["credit_user_id"] = normalized_credit_user_id |
| task["credit_reserved"] = normalized_credit_amount |
| task["credit_daily_reserved"] = daily_reserved |
| task["credit_paid_reserved"] = paid_reserved |
| task["credit_package_reservations"] = package_reservations |
| task["credit_grant_reservations"] = grant_reservations |
| task["credit_settled"] = False |
| if normalized_wallet_user_id and normalized_wallet_amount > 0: |
| task["third_party_wallet_user_id"] = normalized_wallet_user_id |
| task["third_party_wallet_group_id"] = normalized_wallet_group_id |
| task["third_party_wallet_model_id"] = normalized_wallet_model_id |
| task["third_party_wallet_reserved"] = normalized_wallet_amount |
| task["third_party_wallet_cost_per_image"] = normalized_wallet_cost_per_image |
| task["third_party_wallet_settled"] = False |
| self._tasks[key] = task |
| self._save_locked() |
| should_start = True |
|
|
| if should_start: |
| thread = threading.Thread( |
| target=self._run_task, |
| args=(key, mode, payload, dict(identity), _clean(payload.get("model"), "gpt-image-2")), |
| name=f"image-task-{task_id[:16]}", |
| daemon=True, |
| ) |
| thread.start() |
| return _public_task(task) |
|
|
| def _run_task( |
| self, |
| key: str, |
| mode: str, |
| payload: dict[str, Any], |
| identity: dict[str, object], |
| model: str, |
| ) -> None: |
| started = time.time() |
| self._update_task(key, status=TASK_STATUS_RUNNING, started_at=_now_iso(), error="") |
| account_email = "" |
|
|
| def progress_callback(step: str) -> None: |
| cleaned_step = _clean(step) or "generating" |
| self._record_progress(key, cleaned_step, set_started_ts=cleaned_step == "image_stream_resolve_start") |
|
|
| payload_with_progress = {**payload, "progress_callback": progress_callback} |
| try: |
| handler = self.edit_handler if mode == "edit" else self.generation_handler |
| result = handler(payload_with_progress) |
| if not isinstance(result, dict): |
| raise RuntimeError("image task returned streaming result unexpectedly") |
| data = result.get("data") |
| account_email = _clean(result.get("_account_email") or result.get("account_email")) |
| if not isinstance(data, list) or not data: |
| upstream = _clean(result.get("message")) |
| if upstream: |
| message = upstream |
| else: |
| message = "号池中没有可用账号或所有账号均被限流,请检查号池状态(账号额度、是否被封禁、是否到达生图上限)" |
| raise RuntimeError(message) |
| duration_ms = int((time.time() - started) * 1000) |
| self._update_task(key, status=TASK_STATUS_SUCCESS, data=data, error="", duration_ms=duration_ms) |
| self._settle_credit(key, success=True) |
| self._settle_third_party_wallet(key, success_count=_count_image_items(data)) |
| self._log_call( |
| identity, |
| mode, |
| model, |
| started, |
| "调用完成", |
| request_preview=request_text(payload.get("prompt")), |
| urls=_collect_image_urls(data), |
| account_email=account_email, |
| image_api=_collect_third_party_image_api(result), |
| ) |
| except Exception as exc: |
| error_message = str(exc) or "image task failed" |
| account_email = _clean(getattr(exc, "account_email", account_email)) |
| duration_ms = int((time.time() - started) * 1000) |
| self._update_task(key, status=TASK_STATUS_ERROR, error=error_message, data=[], duration_ms=duration_ms) |
| self._settle_credit(key, success=False) |
| self._settle_third_party_wallet(key, success_count=0) |
| self._log_call( |
| identity, |
| mode, |
| model, |
| started, |
| "调用失败", |
| request_preview=request_text(payload.get("prompt")), |
| status="failed", |
| error=error_message, |
| account_email=account_email, |
| ) |
|
|
| def _log_call( |
| self, |
| identity: dict[str, object], |
| mode: str, |
| model: str, |
| started: float, |
| suffix: str, |
| *, |
| request_preview: str = "", |
| status: str = "success", |
| error: str = "", |
| urls: list[str] | None = None, |
| account_email: str = "", |
| image_api: dict[str, Any] | None = None, |
| ) -> None: |
| endpoint = "/v1/images/edits" if mode == "edit" else "/v1/images/generations" |
| summary_prefix = "图生图" if mode == "edit" else "文生图" |
| detail = { |
| "key_id": identity.get("id"), |
| "key_name": identity.get("name"), |
| "role": identity.get("role"), |
| "endpoint": endpoint, |
| "model": model, |
| "started_at": china_from_timestamp_string(started), |
| "ended_at": _now_iso(), |
| "duration_ms": int((time.time() - started) * 1000), |
| "status": status, |
| } |
| if request_preview: |
| detail["request_text"] = request_preview |
| if error: |
| detail["error"] = error |
| if account_email: |
| detail["account_email"] = account_email |
| if image_api: |
| detail["image_api_name"] = " / ".join( |
| item |
| for item in [ |
| _clean(image_api.get("group_name") or image_api.get("group_id")), |
| _clean(image_api.get("connection_name") or image_api.get("connection_id")), |
| ] |
| if item |
| ) |
| detail["image_api_group"] = image_api.get("group_name") or image_api.get("group_id") |
| detail["image_api_connection"] = image_api.get("connection_name") or image_api.get("connection_id") |
| detail["image_api_model"] = image_api.get("model_id") |
| detail["image_api_upstream_model"] = image_api.get("upstream_model") |
| detail["image_api_resolution"] = image_api.get("resolution") |
| if image_api.get("response_format"): |
| detail["image_api_response_format"] = image_api.get("response_format") |
| if image_api.get("size"): |
| detail["image_api_size"] = image_api.get("size") |
| if urls: |
| detail["urls"] = list(dict.fromkeys(urls)) |
| try: |
| log_service.add(LOG_TYPE_CALL, f"{summary_prefix}{suffix}", detail) |
| except Exception: |
| pass |
|
|
| def _update_task(self, key: str, **updates: Any) -> None: |
| with self._lock: |
| task = self._tasks.get(key) |
| if task is None: |
| return |
| task.update(updates) |
| task["updated_at"] = _now_iso() |
| task["updated_ts"] = time.time() |
| self._save_locked() |
|
|
| def _record_progress(self, key: str, step: str, *, set_started_ts: bool = False) -> None: |
| cleaned_step = _clean(step) or "generating" |
| now_ts = time.time() |
| with self._lock: |
| task = self._tasks.get(key) |
| if task is None: |
| return |
| events = task.get("progress_events") if isinstance(task.get("progress_events"), list) else [] |
| normalized_events = [ |
| event for event in events |
| if isinstance(event, dict) and _clean(event.get("step")) |
| ] |
| if not normalized_events or _clean(normalized_events[-1].get("step")) != cleaned_step: |
| normalized_events.append({ |
| "step": cleaned_step, |
| "at": _now_iso(), |
| "ts": now_ts, |
| }) |
| task["progress_events"] = normalized_events[-20:] |
| task["progress"] = cleaned_step |
| if set_started_ts: |
| task["started_ts"] = now_ts |
| task["updated_at"] = _now_iso() |
| task["updated_ts"] = now_ts |
| self._save_locked() |
|
|
| def _settle_credit(self, key: str, *, success: bool) -> None: |
| with self._lock: |
| task = self._tasks.get(key) |
| if task is None or bool(task.get("credit_settled")): |
| return |
| user_id = _clean(task.get("credit_user_id")) |
| amount = max(0, int(task.get("credit_reserved") or 0)) |
| daily_amount = max(0, int(task.get("credit_daily_reserved") or 0)) |
| paid_amount = max(0, int(task.get("credit_paid_reserved") or 0)) |
| package_reservations = task.get("credit_package_reservations") if isinstance(task.get("credit_package_reservations"), list) else [] |
| grant_reservations = task.get("credit_grant_reservations") if isinstance(task.get("credit_grant_reservations"), list) else [] |
| package_amount = sum(max(0, int(item.get("amount") or 0)) for item in package_reservations if isinstance(item, dict)) |
| grant_amount = sum(max(0, int(item.get("amount") or 0)) for item in grant_reservations if isinstance(item, dict)) |
| if amount > 0 and daily_amount + paid_amount + package_amount + grant_amount == 0: |
| paid_amount = amount |
| if not user_id or amount <= 0: |
| task["credit_settled"] = True |
| self._save_locked() |
| return |
| try: |
| from services.auth_service import auth_service |
|
|
| if success: |
| auth_service.consume_reserved_credits( |
| user_id, |
| amount, |
| daily_amount=daily_amount, |
| paid_amount=paid_amount, |
| package_reservations=package_reservations, |
| grant_reservations=grant_reservations, |
| ) |
| else: |
| auth_service.refund_reserved_credits( |
| user_id, |
| amount, |
| daily_amount=daily_amount, |
| paid_amount=paid_amount, |
| package_reservations=package_reservations, |
| grant_reservations=grant_reservations, |
| ) |
| finally: |
| task["credit_settled"] = True |
| task["updated_at"] = _now_iso() |
| self._save_locked() |
|
|
| def _settle_third_party_wallet(self, key: str, *, success_count: int) -> None: |
| with self._lock: |
| task = self._tasks.get(key) |
| if task is None or bool(task.get("third_party_wallet_settled")): |
| return |
| user_id = _clean(task.get("third_party_wallet_user_id")) |
| group_id = _clean(task.get("third_party_wallet_group_id")) |
| model_id = _clean(task.get("third_party_wallet_model_id")) |
| amount = max(0, int(task.get("third_party_wallet_reserved") or 0)) |
| cost_per_image = max(0, int(task.get("third_party_wallet_cost_per_image") or 0)) |
| if not user_id or not group_id or not model_id or amount <= 0: |
| task["third_party_wallet_settled"] = True |
| self._save_locked() |
| return |
| consume_amount = min(amount, max(0, int(success_count or 0)) * cost_per_image) |
| refund_amount = amount - consume_amount |
| try: |
| from services.auth_service import auth_service |
|
|
| if consume_amount > 0: |
| auth_service.consume_third_party_image_wallet( |
| user_id, |
| group_id=group_id, |
| model_id=model_id, |
| amount=consume_amount, |
| ) |
| if refund_amount > 0: |
| auth_service.refund_third_party_image_wallet( |
| user_id, |
| group_id=group_id, |
| model_id=model_id, |
| amount=refund_amount, |
| ) |
| finally: |
| task["third_party_wallet_settled"] = True |
| task["updated_at"] = _now_iso() |
| self._save_locked() |
|
|
| def _load_locked(self) -> dict[str, dict[str, Any]]: |
| if not self.path.exists(): |
| return {} |
| try: |
| raw = json.loads(self.path.read_text(encoding="utf-8")) |
| except Exception: |
| return {} |
| raw_items = raw.get("tasks") if isinstance(raw, dict) else raw |
| if not isinstance(raw_items, list): |
| return {} |
| tasks: dict[str, dict[str, Any]] = {} |
| for item in raw_items: |
| if not isinstance(item, dict): |
| continue |
| task_id = _clean(item.get("id")) |
| owner = _clean(item.get("owner_id")) |
| if not task_id or not owner: |
| continue |
| status = _clean(item.get("status")) |
| if status not in {TASK_STATUS_QUEUED, TASK_STATUS_RUNNING, TASK_STATUS_SUCCESS, TASK_STATUS_ERROR}: |
| status = TASK_STATUS_ERROR |
| task = { |
| "id": task_id, |
| "owner_id": owner, |
| "status": status, |
| "mode": "edit" if item.get("mode") == "edit" else "generate", |
| "model": _clean(item.get("model"), "gpt-image-2"), |
| "size": _clean(item.get("size")), |
| "created_at": _clean(item.get("created_at"), _now_iso()), |
| "started_at": _clean(item.get("started_at")), |
| "updated_at": _clean(item.get("updated_at"), _clean(item.get("created_at"), _now_iso())), |
| "created_ts": item.get("created_ts"), |
| "updated_ts": item.get("updated_ts"), |
| "started_ts": item.get("started_ts"), |
| "duration_ms": item.get("duration_ms"), |
| } |
| progress = _clean(item.get("progress")) |
| if progress: |
| task["progress"] = progress |
| if isinstance(item.get("progress_events"), list): |
| task["progress_events"] = [ |
| event |
| for event in item.get("progress_events", [])[-20:] |
| if isinstance(event, dict) and _clean(event.get("step")) |
| ] |
| credit_user_id = _clean(item.get("credit_user_id")) |
| credit_reserved = max(0, int(item.get("credit_reserved") or 0)) |
| if credit_user_id and credit_reserved > 0: |
| task["credit_user_id"] = credit_user_id |
| task["credit_reserved"] = credit_reserved |
| task["credit_daily_reserved"] = max(0, int(item.get("credit_daily_reserved") or 0)) |
| task["credit_paid_reserved"] = max(0, int(item.get("credit_paid_reserved") or 0)) |
| task["credit_package_reservations"] = item.get("credit_package_reservations") if isinstance(item.get("credit_package_reservations"), list) else [] |
| task["credit_grant_reservations"] = item.get("credit_grant_reservations") if isinstance(item.get("credit_grant_reservations"), list) else [] |
| task["credit_settled"] = bool(item.get("credit_settled")) |
| wallet_user_id = _clean(item.get("third_party_wallet_user_id")) |
| wallet_reserved = max(0, int(item.get("third_party_wallet_reserved") or 0)) |
| if wallet_user_id and wallet_reserved > 0: |
| task["third_party_wallet_user_id"] = wallet_user_id |
| task["third_party_wallet_group_id"] = _clean(item.get("third_party_wallet_group_id")) |
| task["third_party_wallet_model_id"] = _clean(item.get("third_party_wallet_model_id")) |
| task["third_party_wallet_reserved"] = wallet_reserved |
| task["third_party_wallet_cost_per_image"] = max(0, int(item.get("third_party_wallet_cost_per_image") or 0)) |
| task["third_party_wallet_settled"] = bool(item.get("third_party_wallet_settled")) |
| data = item.get("data") |
| if isinstance(data, list): |
| task["data"] = data |
| error = _clean(item.get("error")) |
| if error: |
| task["error"] = error |
| tasks[_task_key(owner, task_id)] = task |
| return tasks |
|
|
| def _save_locked(self) -> None: |
| items = sorted(self._tasks.values(), key=lambda item: str(item.get("updated_at") or ""), reverse=True) |
| tmp_path = self.path.with_suffix(self.path.suffix + ".tmp") |
| tmp_path.write_text(json.dumps({"tasks": items}, ensure_ascii=False, indent=2) + "\n", encoding="utf-8") |
| tmp_path.replace(self.path) |
|
|
| def _recover_unfinished_locked(self) -> bool: |
| changed = False |
| for task in self._tasks.values(): |
| if task.get("status") in UNFINISHED_STATUSES: |
| self._refund_task_credit_locked(task) |
| self._refund_task_third_party_wallet_locked(task) |
| task["status"] = TASK_STATUS_ERROR |
| task["error"] = "服务已重启,未完成的图片任务已中断" |
| task["updated_at"] = _now_iso() |
| changed = True |
| return changed |
|
|
| def _refund_task_credit_locked(self, task: dict[str, Any]) -> None: |
| if bool(task.get("credit_settled")): |
| return |
| user_id = _clean(task.get("credit_user_id")) |
| amount = max(0, int(task.get("credit_reserved") or 0)) |
| daily_amount = max(0, int(task.get("credit_daily_reserved") or 0)) |
| paid_amount = max(0, int(task.get("credit_paid_reserved") or 0)) |
| package_reservations = task.get("credit_package_reservations") if isinstance(task.get("credit_package_reservations"), list) else [] |
| grant_reservations = task.get("credit_grant_reservations") if isinstance(task.get("credit_grant_reservations"), list) else [] |
| package_amount = sum(max(0, int(item.get("amount") or 0)) for item in package_reservations if isinstance(item, dict)) |
| grant_amount = sum(max(0, int(item.get("amount") or 0)) for item in grant_reservations if isinstance(item, dict)) |
| if amount > 0 and daily_amount + paid_amount + package_amount + grant_amount == 0: |
| paid_amount = amount |
| if not user_id or amount <= 0: |
| task["credit_settled"] = True |
| return |
| try: |
| from services.auth_service import auth_service |
|
|
| auth_service.refund_reserved_credits( |
| user_id, |
| amount, |
| daily_amount=daily_amount, |
| paid_amount=paid_amount, |
| package_reservations=package_reservations, |
| grant_reservations=grant_reservations, |
| ) |
| finally: |
| task["credit_settled"] = True |
|
|
| def _refund_task_third_party_wallet_locked(self, task: dict[str, Any]) -> None: |
| if bool(task.get("third_party_wallet_settled")): |
| return |
| user_id = _clean(task.get("third_party_wallet_user_id")) |
| group_id = _clean(task.get("third_party_wallet_group_id")) |
| model_id = _clean(task.get("third_party_wallet_model_id")) |
| amount = max(0, int(task.get("third_party_wallet_reserved") or 0)) |
| if not user_id or not group_id or not model_id or amount <= 0: |
| task["third_party_wallet_settled"] = True |
| return |
| try: |
| from services.auth_service import auth_service |
|
|
| auth_service.refund_third_party_image_wallet( |
| user_id, |
| group_id=group_id, |
| model_id=model_id, |
| amount=amount, |
| ) |
| finally: |
| task["third_party_wallet_settled"] = True |
|
|
| def _cleanup_locked(self) -> bool: |
| try: |
| retention_days = max(1, int(self.retention_days_getter())) |
| except Exception: |
| retention_days = 30 |
| cutoff = time.time() - retention_days * 86400 |
| removed_keys = [ |
| key |
| for key, task in self._tasks.items() |
| if task.get("status") in TERMINAL_STATUSES and _timestamp(task.get("updated_at")) < cutoff |
| ] |
| for key in removed_keys: |
| self._tasks.pop(key, None) |
| return bool(removed_keys) |
|
|
| def _cleanup_if_due_locked(self) -> bool: |
| now = time.monotonic() |
| if now - self._last_cleanup_at < TASK_CLEANUP_INTERVAL_SECONDS: |
| return False |
| self._last_cleanup_at = now |
| return self._cleanup_locked() |
|
|
|
|
| image_task_service = ImageTaskService(DATA_DIR / "image_tasks.json") |
|
|