mirror of
https://github.com/Routstr/routstr-core.git
synced 2026-10-05 20:28:23 +00:00
add routstr logic
This commit is contained in:
@@ -6,7 +6,6 @@ Create Date: 2026-02-13 22:36:53.608737
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
import sqlalchemy as sa
|
import sqlalchemy as sa
|
||||||
import sqlmodel
|
|
||||||
from alembic import op
|
from alembic import op
|
||||||
|
|
||||||
# revision identifiers, used by Alembic.
|
# revision identifiers, used by Alembic.
|
||||||
|
|||||||
+187
-1
@@ -809,6 +809,44 @@ class TopupRequest(BaseModel):
|
|||||||
amount: int
|
amount: int
|
||||||
|
|
||||||
|
|
||||||
|
class TopupTokenRequest(BaseModel):
|
||||||
|
token: str
|
||||||
|
|
||||||
|
|
||||||
|
@admin_router.post(
|
||||||
|
"/api/upstream-providers/{provider_id}/topup-token",
|
||||||
|
dependencies=[Depends(require_admin_api)],
|
||||||
|
)
|
||||||
|
async def topup_provider_with_token(
|
||||||
|
provider_id: int, payload: TopupTokenRequest
|
||||||
|
) -> dict:
|
||||||
|
"""Redeem a Cashu token for an upstream provider."""
|
||||||
|
async with create_session() as session:
|
||||||
|
provider = await session.get(UpstreamProviderRow, provider_id)
|
||||||
|
if not provider:
|
||||||
|
raise HTTPException(status_code=404, detail="Provider not found")
|
||||||
|
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
async with httpx.AsyncClient() as client:
|
||||||
|
clean_url = provider.base_url.rstrip("/")
|
||||||
|
resp = await client.post(
|
||||||
|
f"{clean_url}/v1/balance/topup",
|
||||||
|
json={"cashu_token": payload.token},
|
||||||
|
headers={"Authorization": f"Bearer {provider.api_key}"},
|
||||||
|
)
|
||||||
|
|
||||||
|
if resp.status_code == 200:
|
||||||
|
return {"ok": True, "message": "Token redeemed successfully"}
|
||||||
|
else:
|
||||||
|
logger.error(f"Upstream token topup failed: {resp.text}")
|
||||||
|
try:
|
||||||
|
error_detail = resp.json()
|
||||||
|
except Exception:
|
||||||
|
error_detail = resp.text
|
||||||
|
return {"ok": False, "message": f"Upstream error: {error_detail}"}
|
||||||
|
|
||||||
|
|
||||||
@admin_router.post(
|
@admin_router.post(
|
||||||
"/api/upstream-providers/{provider_id}/topup",
|
"/api/upstream-providers/{provider_id}/topup",
|
||||||
dependencies=[Depends(require_admin_api)],
|
dependencies=[Depends(require_admin_api)],
|
||||||
@@ -835,7 +873,49 @@ async def initiate_provider_topup(
|
|||||||
f"Initiating top-up for provider {provider_id}",
|
f"Initiating top-up for provider {provider_id}",
|
||||||
extra={"amount": payload.amount},
|
extra={"amount": payload.amount},
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# For Routstr providers, we might be doing a Lightning top-up or a direct token transfer
|
||||||
|
if provider.provider_type == "routstr":
|
||||||
|
# UI sends sats for Routstr topup
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
async with httpx.AsyncClient() as client:
|
||||||
|
clean_url = provider.base_url.rstrip("/")
|
||||||
|
# Proxy the request to upstream Routstr
|
||||||
|
# Use the actual API key from the database
|
||||||
|
resp = await client.post(
|
||||||
|
f"{clean_url}/v1/balance/lightning/invoice",
|
||||||
|
json={
|
||||||
|
"amount_sats": int(payload.amount),
|
||||||
|
"purpose": "topup",
|
||||||
|
"api_key": provider.api_key,
|
||||||
|
},
|
||||||
|
headers={"Authorization": f"Bearer {provider.api_key}"},
|
||||||
|
)
|
||||||
|
|
||||||
|
if resp.status_code == 200:
|
||||||
|
data = resp.json()
|
||||||
|
return {
|
||||||
|
"ok": True,
|
||||||
|
"topup_data": {
|
||||||
|
"payment_request": data.get("bolt11"),
|
||||||
|
"invoice_id": data.get("invoice_id"),
|
||||||
|
"status": "pending",
|
||||||
|
},
|
||||||
|
}
|
||||||
|
else:
|
||||||
|
logger.error(f"Upstream topup request failed: {resp.text}")
|
||||||
|
# Check if it's JSON error
|
||||||
|
try:
|
||||||
|
error_detail = resp.json()
|
||||||
|
except Exception:
|
||||||
|
error_detail = resp.text
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=resp.status_code, detail=error_detail
|
||||||
|
)
|
||||||
|
|
||||||
topup_data = await upstream_instance.initiate_topup(payload.amount)
|
topup_data = await upstream_instance.initiate_topup(payload.amount)
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"Top-up initiated successfully",
|
"Top-up initiated successfully",
|
||||||
extra={
|
extra={
|
||||||
@@ -886,6 +966,23 @@ async def check_topup_status(provider_id: int, invoice_id: str) -> dict[str, obj
|
|||||||
if not provider:
|
if not provider:
|
||||||
raise HTTPException(status_code=404, detail="Provider not found")
|
raise HTTPException(status_code=404, detail="Provider not found")
|
||||||
|
|
||||||
|
# For Routstr providers, proxy the status check
|
||||||
|
if provider.provider_type == "routstr":
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
async with httpx.AsyncClient() as client:
|
||||||
|
clean_url = provider.base_url.rstrip("/")
|
||||||
|
resp = await client.get(
|
||||||
|
f"{clean_url}/v1/balance/lightning/invoice/{invoice_id}/status",
|
||||||
|
headers={"Authorization": f"Bearer {provider.api_key}"},
|
||||||
|
)
|
||||||
|
if resp.status_code == 200:
|
||||||
|
status_data = resp.json()
|
||||||
|
return {"ok": True, "paid": status_data.get("status") == "paid"}
|
||||||
|
else:
|
||||||
|
logger.error(f"Upstream status check failed: {resp.text}")
|
||||||
|
return {"ok": False, "paid": False}
|
||||||
|
|
||||||
upstream_instance = _instantiate_provider(provider)
|
upstream_instance = _instantiate_provider(provider)
|
||||||
if not upstream_instance:
|
if not upstream_instance:
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
@@ -913,7 +1010,7 @@ async def check_topup_status(provider_id: int, invoice_id: str) -> dict[str, obj
|
|||||||
dependencies=[Depends(require_admin_api)],
|
dependencies=[Depends(require_admin_api)],
|
||||||
)
|
)
|
||||||
async def get_provider_balance(provider_id: int) -> dict[str, object]:
|
async def get_provider_balance(provider_id: int) -> dict[str, object]:
|
||||||
"""Get the current account balance for the upstream provider."""
|
"""Get the current balance for an upstream provider account."""
|
||||||
from ..upstream.helpers import _instantiate_provider
|
from ..upstream.helpers import _instantiate_provider
|
||||||
|
|
||||||
async with create_session() as session:
|
async with create_session() as session:
|
||||||
@@ -921,6 +1018,27 @@ async def get_provider_balance(provider_id: int) -> dict[str, object]:
|
|||||||
if not provider:
|
if not provider:
|
||||||
raise HTTPException(status_code=404, detail="Provider not found")
|
raise HTTPException(status_code=404, detail="Provider not found")
|
||||||
|
|
||||||
|
# For Routstr providers, proxy the balance check
|
||||||
|
if provider.provider_type == "routstr":
|
||||||
|
import httpx
|
||||||
|
|
||||||
|
async with httpx.AsyncClient() as client:
|
||||||
|
clean_url = provider.base_url.rstrip("/")
|
||||||
|
resp = await client.get(
|
||||||
|
f"{clean_url}/v1/balance/info",
|
||||||
|
headers={"Authorization": f"Bearer {provider.api_key}"},
|
||||||
|
)
|
||||||
|
if resp.status_code == 200:
|
||||||
|
data = resp.json()
|
||||||
|
# Return balance in sats
|
||||||
|
balance = data.get("balance", 0)
|
||||||
|
if isinstance(balance, (int, float)):
|
||||||
|
return {"ok": True, "balance_data": balance // 1000}
|
||||||
|
return {"ok": True, "balance_data": balance}
|
||||||
|
else:
|
||||||
|
logger.error(f"Failed to fetch Routstr balance: {resp.text}")
|
||||||
|
return {"ok": False, "balance_data": None}
|
||||||
|
|
||||||
upstream_instance = _instantiate_provider(provider)
|
upstream_instance = _instantiate_provider(provider)
|
||||||
if not upstream_instance:
|
if not upstream_instance:
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
@@ -1090,3 +1208,71 @@ async def get_log_dates_api(request: Request) -> dict[str, object]:
|
|||||||
continue
|
continue
|
||||||
|
|
||||||
return {"dates": dates}
|
return {"dates": dates}
|
||||||
|
|
||||||
|
|
||||||
|
@admin_router.post(
|
||||||
|
"/api/upstream-providers/{provider_id}/routstr/refund",
|
||||||
|
dependencies=[Depends(require_admin_api)],
|
||||||
|
)
|
||||||
|
async def refund_routstr_provider_balance(provider_id: int) -> dict[str, object]:
|
||||||
|
"""Refund balance from an upstream Routstr provider back to the local wallet."""
|
||||||
|
from ..upstream.helpers import _instantiate_provider
|
||||||
|
from ..upstream.routstr import RoutstrUpstreamProvider
|
||||||
|
|
||||||
|
async with create_session() as session:
|
||||||
|
provider_row = await session.get(UpstreamProviderRow, provider_id)
|
||||||
|
if not provider_row:
|
||||||
|
raise HTTPException(status_code=404, detail="Provider not found")
|
||||||
|
|
||||||
|
if provider_row.provider_type != "routstr":
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=400, detail="Refund only supported for Routstr providers"
|
||||||
|
)
|
||||||
|
|
||||||
|
provider = _instantiate_provider(provider_row)
|
||||||
|
if not isinstance(provider, RoutstrUpstreamProvider):
|
||||||
|
raise HTTPException(status_code=400, detail="Invalid provider instance")
|
||||||
|
|
||||||
|
try:
|
||||||
|
# Request refund from upstream
|
||||||
|
data = await provider.refund_balance()
|
||||||
|
if "error" in data:
|
||||||
|
# If the upstream returned an OpenAI-style error (like the model unknown error)
|
||||||
|
# it means the request likely didn't even reach the refund endpoint handler
|
||||||
|
# but was intercepted by the proxy layer.
|
||||||
|
error_info = data.get("error", {})
|
||||||
|
message = (
|
||||||
|
error_info.get("message")
|
||||||
|
if isinstance(error_info, dict)
|
||||||
|
else str(error_info)
|
||||||
|
)
|
||||||
|
return {
|
||||||
|
"ok": False,
|
||||||
|
"message": f"Upstream refund failed: {message}",
|
||||||
|
}
|
||||||
|
|
||||||
|
token = data.get("token")
|
||||||
|
if not token:
|
||||||
|
return {"ok": False, "message": "Upstream did not return a token"}
|
||||||
|
|
||||||
|
# Receive token into local wallet
|
||||||
|
from ..wallet import recieve_token
|
||||||
|
|
||||||
|
try:
|
||||||
|
# Use current wallet to receive
|
||||||
|
await recieve_token(token)
|
||||||
|
return {
|
||||||
|
"ok": True,
|
||||||
|
"message": "Successfully received refund from upstream provider",
|
||||||
|
}
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(f"Failed to receive refund token: {e}")
|
||||||
|
return {
|
||||||
|
"ok": False,
|
||||||
|
"message": f"Failed to receive refund token: {str(e)}",
|
||||||
|
"token": token,
|
||||||
|
}
|
||||||
|
|
||||||
|
except Exception as e:
|
||||||
|
logger.exception(f"Refund failed for provider {provider_id}")
|
||||||
|
raise HTTPException(status_code=500, detail=str(e))
|
||||||
|
|||||||
+1
-5
@@ -138,11 +138,6 @@ async def proxy(
|
|||||||
) -> Response | StreamingResponse:
|
) -> Response | StreamingResponse:
|
||||||
headers = dict(request.headers)
|
headers = dict(request.headers)
|
||||||
|
|
||||||
if "x-cashu" not in headers and "authorization" not in headers.keys():
|
|
||||||
return create_error_response(
|
|
||||||
"unauthorized", "Unauthorized", 401, request=request
|
|
||||||
)
|
|
||||||
|
|
||||||
is_responses_api = path.startswith("v1/responses") or path.startswith("responses")
|
is_responses_api = path.startswith("v1/responses") or path.startswith("responses")
|
||||||
request_body = await request.body()
|
request_body = await request.body()
|
||||||
request_body_dict = parse_request_body_json(request_body, path)
|
request_body_dict = parse_request_body_json(request_body, path)
|
||||||
@@ -153,6 +148,7 @@ async def proxy(
|
|||||||
model_id = request_body_dict.get("model", "unknown")
|
model_id = request_body_dict.get("model", "unknown")
|
||||||
|
|
||||||
model_obj = get_model_instance(model_id)
|
model_obj = get_model_instance(model_id)
|
||||||
|
|
||||||
if not model_obj:
|
if not model_obj:
|
||||||
return create_error_response(
|
return create_error_response(
|
||||||
"invalid_model", f"Model '{model_id}' not found", 400, request=request
|
"invalid_model", f"Model '{model_id}' not found", 400, request=request
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
from typing import TYPE_CHECKING, Any, Mapping
|
from typing import TYPE_CHECKING, Any
|
||||||
|
|
||||||
import httpx
|
import httpx
|
||||||
|
|
||||||
@@ -147,3 +147,24 @@ class RoutstrUpstreamProvider(BaseUpstreamProvider):
|
|||||||
extra={"url": url, "error": str(e)},
|
extra={"url": url, "error": str(e)},
|
||||||
)
|
)
|
||||||
return []
|
return []
|
||||||
|
|
||||||
|
async def refund_balance(self) -> dict[str, Any]:
|
||||||
|
"""Request a refund from the upstream Routstr node.
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Dict containing refund result and token
|
||||||
|
"""
|
||||||
|
url = f"{self.base_url}/v1/balance/refund"
|
||||||
|
headers = {"Authorization": f"Bearer {self.api_key}"}
|
||||||
|
|
||||||
|
async with httpx.AsyncClient() as client:
|
||||||
|
try:
|
||||||
|
response = await client.post(url, headers=headers, timeout=30.0)
|
||||||
|
response.raise_for_status()
|
||||||
|
return response.json()
|
||||||
|
except Exception as e:
|
||||||
|
logger.error(
|
||||||
|
"Failed to request refund from upstream Routstr",
|
||||||
|
extra={"url": url, "error": str(e)},
|
||||||
|
)
|
||||||
|
return {"error": str(e)}
|
||||||
|
|||||||
@@ -83,12 +83,6 @@ export function RoutstrProviderCard({
|
|||||||
>
|
>
|
||||||
{provider.enabled ? 'Enabled' : 'Disabled'}
|
{provider.enabled ? 'Enabled' : 'Disabled'}
|
||||||
</Badge>
|
</Badge>
|
||||||
<Badge
|
|
||||||
variant='outline'
|
|
||||||
className='bg-blue-50 text-blue-700 dark:bg-blue-900/20 dark:text-blue-400'
|
|
||||||
>
|
|
||||||
NIP-91
|
|
||||||
</Badge>
|
|
||||||
{!hasMint && (
|
{!hasMint && (
|
||||||
<Badge
|
<Badge
|
||||||
variant='outline'
|
variant='outline'
|
||||||
|
|||||||
Reference in New Issue
Block a user