summaryrefslogtreecommitdiff
path: root/src/services/api
diff options
context:
space:
mode:
Diffstat (limited to 'src/services/api')
-rw-r--r--src/services/api/background.py179
-rw-r--r--src/services/api/graphql/README.graphql2
-rw-r--r--src/services/api/graphql/bindings.py2
-rwxr-xr-xsrc/services/api/graphql/generate/generate_schema.py2
-rwxr-xr-xsrc/services/api/graphql/generate/schema_from_composite.py4
-rwxr-xr-xsrc/services/api/graphql/generate/schema_from_config_session.py4
-rwxr-xr-xsrc/services/api/graphql/generate/schema_from_op_mode.py4
-rw-r--r--src/services/api/graphql/graphql/auth_token_mutation.py2
-rw-r--r--src/services/api/graphql/graphql/directives.py2
-rw-r--r--src/services/api/graphql/graphql/mutations.py2
-rw-r--r--src/services/api/graphql/graphql/queries.py2
-rw-r--r--src/services/api/graphql/libs/key_auth.py2
-rw-r--r--src/services/api/graphql/libs/op_mode.py2
-rw-r--r--src/services/api/graphql/libs/token_auth.py2
-rw-r--r--src/services/api/graphql/routers.py4
-rwxr-xr-xsrc/services/api/graphql/session/composite/system_status.py2
-rw-r--r--src/services/api/graphql/session/override/remove_firewall_address_group_members.py2
-rw-r--r--src/services/api/graphql/session/session.py2
-rw-r--r--src/services/api/rest/models.py31
-rw-r--r--src/services/api/rest/routers.py343
-rw-r--r--src/services/api/session.py2
21 files changed, 508 insertions, 89 deletions
diff --git a/src/services/api/background.py b/src/services/api/background.py
new file mode 100644
index 000000000..2b25af307
--- /dev/null
+++ b/src/services/api/background.py
@@ -0,0 +1,179 @@
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
+#
+# This library is free software; you can redistribute it and/or
+# modify it under the terms of the GNU Lesser General Public
+# License as published by the Free Software Foundation; either
+# version 2.1 of the License, or (at your option) any later version.
+#
+# This library is distributed in the hope that it will be useful,
+# but WITHOUT ANY WARRANTY; without even the implied warranty of
+# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
+# Lesser General Public License for more details.
+#
+# You should have received a copy of the GNU Lesser General Public License
+# along with this library. If not, see <http://www.gnu.org/licenses/>.
+
+import time
+import functools
+from collections import deque
+from enum import Enum
+from threading import Lock
+from typing import Any
+from typing import Callable
+from typing import Optional
+from uuid import uuid4
+
+from fastapi import BackgroundTasks
+from pydantic import BaseModel
+from pydantic import StrictStr
+from pydantic import StrictInt
+
+
+def _ts():
+ """Return current Unix timestamp (seconds since epoch)"""
+ return int(time.time())
+
+
+class BackgroundOpStatus(str, Enum):
+ queued = 'queued'
+ running = 'running'
+ succeeded = 'succeeded'
+ failed = 'failed'
+
+ @property
+ def is_completed(self):
+ """True if the operation is in a terminal state (succeeded/failed)"""
+ return self in (BackgroundOpStatus.succeeded, BackgroundOpStatus.failed)
+
+
+class BackgroundOpRecord(BaseModel):
+ """Metadata and outcome for a single background operation"""
+
+ op_id: StrictStr
+ created_at: StrictInt
+ started_at: Optional[StrictInt] = None
+ finished_at: Optional[StrictInt] = None
+ status: BackgroundOpStatus = BackgroundOpStatus.queued
+ result: Optional[Any] = None
+ error: Optional[StrictStr] = None
+
+
+class BackgroundOpError(Exception):
+ """Raised when a background operation cannot be enqueued/executed"""
+
+ pass
+
+
+class BackgroundOpManager:
+ """
+ In-memory FIFO operation queue.
+
+ Uses BackgroundTasks to schedule a `drain()` call after the response,
+ so `enqueue()` is fast and non-blocking for the client.
+ """
+
+ DEFAULT_MAX_QUEUE_SIZE = 128
+
+ def __init__(self, max_queue_size: int = DEFAULT_MAX_QUEUE_SIZE):
+ # max number of queued (pending) operations allowed at a time
+ self._max_queue_size = max_queue_size
+
+ # FIFO queue of operation IDs waiting to be executed
+ self._queue = deque()
+ self._jobs = {}
+ self._workers = {}
+
+ # protects _queue/_jobs/_workers/_drain_scheduled from concurrent access
+ self._mx = Lock()
+
+ # whether a drain task has already been scheduled via BackgroundTasks
+ self._drain_scheduled = False
+
+ def enqueue(
+ self,
+ background_tasks: BackgroundTasks,
+ func: Callable,
+ *args,
+ **kwargs,
+ ) -> BackgroundOpRecord:
+ """Enqueue a function for background execution and return its record"""
+
+ assert isinstance(background_tasks, BackgroundTasks)
+ assert callable(func), '`func` argument should be function or lambda'
+
+ with self._mx:
+ if len(self._queue) >= self._max_queue_size:
+ raise BackgroundOpError(
+ f'Background operation queue is full ({self._max_queue_size})'
+ )
+
+ op_id = str(uuid4())
+ record = BackgroundOpRecord(op_id=op_id, created_at=_ts())
+
+ self._jobs[op_id] = record
+ # store the callable for later execution (outside the lock)
+ self._workers[op_id] = functools.partial(func, *args, **kwargs)
+ self._queue.append(op_id)
+
+ if not self._drain_scheduled:
+ # schedule a single drain() call after the current response
+ background_tasks.add_task(self.drain)
+ self._drain_scheduled = True
+
+ # Best-effort pruning: keep history bounded by dropping oldest completed records
+ if len(self._jobs) > self._max_queue_size:
+ oldest = min(self._jobs.values(), key=lambda record: record.created_at)
+ if oldest.status.is_completed:
+ del self._jobs[oldest.op_id]
+
+ return record
+
+ def drain(self):
+ """Run queued operations sequentially until the queue is empty"""
+
+ while True:
+ with self._mx:
+ if not self._queue:
+ # allow future enqueue() calls to schedule the next drain()
+ self._drain_scheduled = False
+ return
+
+ op_id = self._queue.popleft()
+ record = self._jobs[op_id]
+ func = self._workers.pop(op_id)
+
+ record.status = BackgroundOpStatus.running
+ record.started_at = _ts()
+
+ # execute outside the lock to avoid blocking enqueues/status reads
+ result = error = status = None
+ try:
+ result = func()
+ except Exception as e: # noqa: BLE001
+ status = BackgroundOpStatus.failed
+ error = str(e)
+ else:
+ status = BackgroundOpStatus.succeeded
+
+ with self._mx:
+ record.result = result
+ record.error = error
+ record.status = status
+ record.finished_at = _ts()
+
+ def get_record(self, op_id: str) -> BackgroundOpRecord | None:
+ """Return a deep copy of a single record"""
+
+ with self._mx:
+ record = self._jobs.get(op_id)
+ return record.copy(deep=True) if record else None
+
+ def get_records(self) -> list:
+ """Return deep copies of all records, sorted oldest-first by created_at"""
+
+ with self._mx:
+ records = [record.copy(deep=True) for record in self._jobs.values()]
+
+ # stable-ish ordering (oldest first)
+ records.sort(key=lambda record: record.created_at)
+ return records
diff --git a/src/services/api/graphql/README.graphql b/src/services/api/graphql/README.graphql
index 1133d79ed..0f43ac356 100644
--- a/src/services/api/graphql/README.graphql
+++ b/src/services/api/graphql/README.graphql
@@ -64,7 +64,7 @@ save to /config/config.boot; to save to an alternative path, specify
fileName.
Similarly, using an analogous 'endpoint' (meaning the form of the request
-and resolver; the actual enpoint for all GraphQL requests is
+and resolver; the actual endpoint for all GraphQL requests is
https://hostname/graphql), one can load an arbitrary config file from a
path.
diff --git a/src/services/api/graphql/bindings.py b/src/services/api/graphql/bindings.py
index ebf745f32..7380dbb5f 100644
--- a/src/services/api/graphql/bindings.py
+++ b/src/services/api/graphql/bindings.py
@@ -1,4 +1,4 @@
-# Copyright 2021-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/generate/generate_schema.py b/src/services/api/graphql/generate/generate_schema.py
index dd5e7ea56..bb36a4c04 100755
--- a/src/services/api/graphql/generate/generate_schema.py
+++ b/src/services/api/graphql/generate/generate_schema.py
@@ -1,6 +1,6 @@
#!/usr/bin/env python3
#
-# Copyright (C) 2023 VyOS maintainers and contributors
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License version 2 or later as
diff --git a/src/services/api/graphql/generate/schema_from_composite.py b/src/services/api/graphql/generate/schema_from_composite.py
index 06e74032d..9a07f88fe 100755
--- a/src/services/api/graphql/generate/schema_from_composite.py
+++ b/src/services/api/graphql/generate/schema_from_composite.py
@@ -1,6 +1,6 @@
#!/usr/bin/env python3
#
-# Copyright (C) 2022-2023 VyOS maintainers and contributors
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License version 2 or later as
@@ -15,7 +15,7 @@
# along with this program. If not, see <http://www.gnu.org/licenses/>.
#
#
-# A utility to generate GraphQL schema defintions from typing information of
+# A utility to generate GraphQL schema definitions from typing information of
# composite functions comprising several requests.
import os
diff --git a/src/services/api/graphql/generate/schema_from_config_session.py b/src/services/api/graphql/generate/schema_from_config_session.py
index 1d5ff1e53..bfa4bc006 100755
--- a/src/services/api/graphql/generate/schema_from_config_session.py
+++ b/src/services/api/graphql/generate/schema_from_config_session.py
@@ -1,6 +1,6 @@
#!/usr/bin/env python3
#
-# Copyright (C) 2022-2023 VyOS maintainers and contributors
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License version 2 or later as
@@ -15,7 +15,7 @@
# along with this program. If not, see <http://www.gnu.org/licenses/>.
#
#
-# A utility to generate GraphQL schema defintions from typing information of
+# A utility to generate GraphQL schema definitions from typing information of
# (wrappers of) native configsession functions.
import os
diff --git a/src/services/api/graphql/generate/schema_from_op_mode.py b/src/services/api/graphql/generate/schema_from_op_mode.py
index ab7cb691f..618ea2e61 100755
--- a/src/services/api/graphql/generate/schema_from_op_mode.py
+++ b/src/services/api/graphql/generate/schema_from_op_mode.py
@@ -1,6 +1,6 @@
#!/usr/bin/env python3
#
-# Copyright (C) 2022-2023 VyOS maintainers and contributors
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License version 2 or later as
@@ -15,7 +15,7 @@
# along with this program. If not, see <http://www.gnu.org/licenses/>.
#
#
-# A utility to generate GraphQL schema defintions from standardized op-mode
+# A utility to generate GraphQL schema definitions from standardized op-mode
# scripts.
import os
diff --git a/src/services/api/graphql/graphql/auth_token_mutation.py b/src/services/api/graphql/graphql/auth_token_mutation.py
index c74364603..a8020d149 100644
--- a/src/services/api/graphql/graphql/auth_token_mutation.py
+++ b/src/services/api/graphql/graphql/auth_token_mutation.py
@@ -1,4 +1,4 @@
-# Copyright 2022-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/graphql/directives.py b/src/services/api/graphql/graphql/directives.py
index 3927aee58..037f09204 100644
--- a/src/services/api/graphql/graphql/directives.py
+++ b/src/services/api/graphql/graphql/directives.py
@@ -1,4 +1,4 @@
-# Copyright 2021-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/graphql/mutations.py b/src/services/api/graphql/graphql/mutations.py
index 0b391c070..c979d06e8 100644
--- a/src/services/api/graphql/graphql/mutations.py
+++ b/src/services/api/graphql/graphql/mutations.py
@@ -1,4 +1,4 @@
-# Copyright 2021-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/graphql/queries.py b/src/services/api/graphql/graphql/queries.py
index 9303fe909..3a8d12344 100644
--- a/src/services/api/graphql/graphql/queries.py
+++ b/src/services/api/graphql/graphql/queries.py
@@ -1,4 +1,4 @@
-# Copyright 2021-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/libs/key_auth.py b/src/services/api/graphql/libs/key_auth.py
index ffd7f32b2..dc3322fea 100644
--- a/src/services/api/graphql/libs/key_auth.py
+++ b/src/services/api/graphql/libs/key_auth.py
@@ -1,4 +1,4 @@
-# Copyright 2021-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/libs/op_mode.py b/src/services/api/graphql/libs/op_mode.py
index 86e38eae6..fa726264c 100644
--- a/src/services/api/graphql/libs/op_mode.py
+++ b/src/services/api/graphql/libs/op_mode.py
@@ -1,4 +1,4 @@
-# Copyright 2022-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/libs/token_auth.py b/src/services/api/graphql/libs/token_auth.py
index 4f743a096..73c52bdf0 100644
--- a/src/services/api/graphql/libs/token_auth.py
+++ b/src/services/api/graphql/libs/token_auth.py
@@ -1,4 +1,4 @@
-# Copyright 2021-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/routers.py b/src/services/api/graphql/routers.py
index ed3ee1e8c..c6886ba1c 100644
--- a/src/services/api/graphql/routers.py
+++ b/src/services/api/graphql/routers.py
@@ -1,4 +1,4 @@
-# Copyright 2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
@@ -32,7 +32,7 @@ def graphql_init(app: 'FastAPI'):
state = SessionState()
- # import after initializaion of state
+ # import after initialization of state
from .bindings import generate_schema
schema = generate_schema()
diff --git a/src/services/api/graphql/session/composite/system_status.py b/src/services/api/graphql/session/composite/system_status.py
index 516a4eff6..1674b2c2b 100755
--- a/src/services/api/graphql/session/composite/system_status.py
+++ b/src/services/api/graphql/session/composite/system_status.py
@@ -1,6 +1,6 @@
#!/usr/bin/env python3
#
-# Copyright (C) 2022-2024 VyOS maintainers and contributors
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This program is free software; you can redistribute it and/or modify
# it under the terms of the GNU General Public License version 2 or later as
diff --git a/src/services/api/graphql/session/override/remove_firewall_address_group_members.py b/src/services/api/graphql/session/override/remove_firewall_address_group_members.py
index b91932e14..9f39465a1 100644
--- a/src/services/api/graphql/session/override/remove_firewall_address_group_members.py
+++ b/src/services/api/graphql/session/override/remove_firewall_address_group_members.py
@@ -1,4 +1,4 @@
-# Copyright 2021 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/graphql/session/session.py b/src/services/api/graphql/session/session.py
index 619534f43..e4725e752 100644
--- a/src/services/api/graphql/session/session.py
+++ b/src/services/api/graphql/session/session.py
@@ -1,4 +1,4 @@
-# Copyright 2021-2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
diff --git a/src/services/api/rest/models.py b/src/services/api/rest/models.py
index dda50010f..bfea17344 100644
--- a/src/services/api/rest/models.py
+++ b/src/services/api/rest/models.py
@@ -1,4 +1,4 @@
-# Copyright 2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
@@ -26,6 +26,7 @@ from typing import Self
from pydantic import BaseModel
from pydantic import StrictStr
+from pydantic import StrictInt
from pydantic import field_validator
from pydantic import model_validator
from fastapi.responses import HTMLResponse
@@ -47,7 +48,7 @@ def success(data):
# Pydantic models for validation
# Pydantic will cast when possible, so use StrictStr validators added as
# needed for additional constraints
-# json_schema_extra adds anotations to OpenAPI to add examples
+# json_schema_extra adds annotations to OpenAPI to add examples
class ApiModel(BaseModel):
@@ -71,6 +72,8 @@ class BaseConfigureModel(BasePathModel):
class ConfigureModel(ApiModel, BaseConfigureModel):
+ confirm_time: StrictInt = 0
+
class Config:
json_schema_extra = {
'example': {
@@ -81,8 +84,12 @@ class ConfigureModel(ApiModel, BaseConfigureModel):
}
+class ConfirmModel(ApiModel):
+ op: StrictStr
+
class ConfigureListModel(ApiModel):
commands: List[BaseConfigureModel]
+ confirm_time: StrictInt = 0
class Config:
json_schema_extra = {
@@ -134,13 +141,17 @@ class RetrieveModel(ApiModel):
class ConfigFileModel(ApiModel):
op: StrictStr
file: StrictStr = None
+ string: StrictStr = None
+ confirm_time: StrictInt = 0
+ destructive: bool = False
class Config:
json_schema_extra = {
'example': {
'key': 'id_key',
- 'op': 'save | load',
+ 'op': 'save | load | merge | confirm',
'file': 'filename',
+ 'string': 'config_string'
}
}
@@ -251,6 +262,20 @@ class RebootModel(ApiModel):
}
+class RenewModel(ApiModel):
+ op: StrictStr
+ path: List[StrictStr]
+
+ class Config:
+ json_schema_extra = {
+ 'example': {
+ 'key': 'id_key',
+ 'op': 'renew',
+ 'path': ['op', 'mode', 'path'],
+ }
+ }
+
+
class ResetModel(ApiModel):
op: StrictStr
path: List[StrictStr]
diff --git a/src/services/api/rest/routers.py b/src/services/api/rest/routers.py
index e52c77fda..fe67d4612 100644
--- a/src/services/api/rest/routers.py
+++ b/src/services/api/rest/routers.py
@@ -1,4 +1,4 @@
-# Copyright 2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public
@@ -18,6 +18,7 @@
# pylint: disable=wildcard-import,unused-wildcard-import
# pylint: disable=broad-exception-caught
+import asyncio
import json
import copy
import logging
@@ -28,12 +29,14 @@ from typing import Callable
from typing import TYPE_CHECKING
from fastapi import Depends
+from fastapi import Query
from fastapi import Request
from fastapi import Response
from fastapi import HTTPException
from fastapi import APIRouter
from fastapi import BackgroundTasks
from fastapi.routing import APIRoute
+from fastapi.concurrency import run_in_threadpool
from starlette.datastructures import FormData
from starlette.formparsers import FormParser
from starlette.formparsers import MultiPartParser
@@ -45,12 +48,15 @@ from vyos.configtree import ConfigTree
from vyos.configdiff import get_config_diff
from vyos.configsession import ConfigSessionError
+from ..background import BackgroundOpManager
+from ..background import BackgroundOpError
from ..session import SessionState
from .models import success
from .models import error
from .models import responses
from .models import ApiModel
from .models import ConfigureModel
+from .models import ConfirmModel
from .models import ConfigureListModel
from .models import ConfigSectionModel
from .models import ConfigSectionListModel
@@ -66,6 +72,7 @@ from .models import GenerateModel
from .models import ShowModel
from .models import RebootModel
from .models import ResetModel
+from .models import RenewModel
from .models import ImportPkiModel
from .models import PoweroffModel
from .models import TracerouteModel
@@ -79,6 +86,7 @@ LOG = logging.getLogger('http_api.routers')
lock = Lock()
+asynclock = asyncio.Lock()
def check_auth(key_list, key):
key_id = None
@@ -99,7 +107,7 @@ def auth_required(data: ApiModel):
# override Request and APIRoute classes in order to convert form request to json;
-# do all explicit validation here, for backwards compatability of error messages;
+# do all explicit validation here, for backwards compatibility of error messages;
# the explicit validation may be dropped, if desired, in favor of native
# validation by FastAPI/Pydantic, as is used for application/json requests
class MultipartRequest(Request):
@@ -227,7 +235,7 @@ class MultipartRequest(Request):
400,
f"Malformed command '{0}': 'path' field must be a list of strings",
)
- if endpoint in ('/configure'):
+ if endpoint in ('/configure',):
if not c['path']:
self.form_err = (
400,
@@ -238,7 +246,7 @@ class MultipartRequest(Request):
400,
f"Malformed command '{c}': 'value' field must be a string",
)
- if endpoint in ('/configure-section'):
+ if endpoint in ('/configure-section',):
if 'section' not in c and 'config' not in c:
self.form_err = (
400,
@@ -290,6 +298,10 @@ router = APIRouter(
self_ref_msg = 'Requested HTTP API server configuration change; commit will be called in the background'
+# Global background-op manager used by the REST API to run long config commits after the response
+background_op_manager = BackgroundOpManager()
+
+
def call_commit(s: SessionState):
try:
s.session.commit()
@@ -301,24 +313,71 @@ def call_commit(s: SessionState):
LOG.warning(f'ConfigSessionError: {e}')
-def _configure_op(
+def call_commit_confirm(s: SessionState):
+ env = s.session.get_session_env()
+ env['IN_COMMIT_CONFIRM'] = 't'
+ try:
+ s.session.commit()
+ s.session.commit_confirm(minutes=s.confirm_time)
+ except ConfigSessionError as e:
+ s.session.discard()
+ if s.debug:
+ LOG.warning(f'ConfigSessionError:\n {traceback.format_exc()}')
+ else:
+ LOG.warning(f'ConfigSessionError: {e}')
+ finally:
+ del env['IN_COMMIT_CONFIRM']
+
+
+def run_commit(s: SessionState):
+ try:
+ out = s.session.commit()
+ return out, None
+ except Exception as e:
+ return None, e
+
+
+def run_commit_confirm(s: SessionState):
+ env = s.session.get_session_env()
+ env['IN_COMMIT_CONFIRM'] = 't'
+ try:
+ out_c = s.session.commit()
+ out_cc = s.session.commit_confirm(minutes=s.confirm_time)
+ out = out_c + '\n' + out_cc
+ return out, None
+ except Exception as e:
+ return None, e
+ finally:
+ del env['IN_COMMIT_CONFIRM']
+
+
+def _execute_configure_op(
data: Union[
+ ConfirmModel,
ConfigureModel,
ConfigureListModel,
ConfigSectionModel,
ConfigSectionListModel,
ConfigSectionTreeModel,
],
- _request: Request,
- background_tasks: BackgroundTasks,
+ background_tasks: BackgroundTasks | None = None,
):
# pylint: disable=too-many-branches,too-many-locals,too-many-nested-blocks,too-many-statements
# pylint: disable=consider-using-with
+ # True when invoked by the background operation
+ # runner (no FastAPI BackgroundTasks context passed in)
+ is_background_job = background_tasks is None
+
state = SessionState()
session = state.session
env = session.get_session_env()
+ # A non-zero confirm_time will start commit-confirm timer on commit
+ confirm_time = 0
+ if isinstance(data, (ConfigureModel, ConfigureListModel, ConfigFileModel)):
+ confirm_time = data.confirm_time
+
# Allow users to pass just one command
if not isinstance(data, (ConfigureListModel, ConfigSectionListModel)):
data = [data]
@@ -338,10 +397,18 @@ def _configure_op(
try:
for c in data:
op = c.op
- if not isinstance(c, BaseConfigSectionTreeModel):
+ op_error = ConfigSessionError(f"'{op}' is not a valid operation")
+
+ if not isinstance(c, (ConfirmModel, BaseConfigSectionTreeModel)):
path = c.path
- if isinstance(c, BaseConfigureModel):
+ if isinstance(c, ConfirmModel):
+ if op == 'confirm':
+ msg = session.confirm()
+ else:
+ raise op_error
+
+ elif isinstance(c, BaseConfigureModel):
if c.value:
value = c.value
else:
@@ -354,8 +421,8 @@ def _configure_op(
section = c.section
elif isinstance(c, BaseConfigSectionTreeModel):
- mask = c.mask
- config = c.config
+ mask_dict = c.mask
+ config_dict = c.config
if isinstance(c, BaseConfigureModel):
if op == 'set':
@@ -369,7 +436,7 @@ def _configure_op(
elif op == 'comment':
session.comment(path, value=value)
else:
- raise ConfigSessionError(f"'{op}' is not a valid operation")
+ raise op_error
elif isinstance(c, BaseConfigSectionModel):
if op == 'set':
@@ -377,26 +444,50 @@ def _configure_op(
elif op == 'load':
session.load_section(path, section)
else:
- raise ConfigSessionError(f"'{op}' is not a valid operation")
+ raise op_error
elif isinstance(c, BaseConfigSectionTreeModel):
if op == 'set':
- session.set_section_tree(config)
+ session.set_section_tree(config_dict)
elif op == 'load':
- session.load_section_tree(mask, config)
+ config_tree = config.get_config_tree()
+ session.load_section_tree(config_tree, mask_dict, config_dict)
else:
- raise ConfigSessionError(f"'{op}' is not a valid operation")
+ raise op_error
# end for
+
config = Config(session_env=env)
d = get_config_diff(config)
- if d.is_node_changed(['service', 'https']):
- background_tasks.add_task(call_commit, state)
- msg = self_ref_msg
+ state.confirm_time = confirm_time if confirm_time else 0
+
+ if not d.is_node_changed(['service', 'https']):
+ if confirm_time:
+ out, err = run_commit_confirm(state)
+ if err:
+ raise err
+ msg = msg + out if msg else out
+ else:
+ out, err = run_commit(state)
+ if err:
+ raise err
+ msg = msg + out if msg else out
else:
- # capture non-fatal warnings
- out = session.commit()
- msg = out if out else msg
+ if is_background_job:
+ # If already running as a background job, commit synchronously here
+ if confirm_time:
+ call_commit_confirm(state)
+ else:
+ call_commit(state)
+ else:
+ # Otherwise schedule the commit to run after the HTTP response
+ if confirm_time:
+ background_tasks.add_task(call_commit_confirm, state)
+ else:
+ background_tasks.add_task(call_commit, state)
+
+ out = self_ref_msg
+ msg = msg + out if msg else out
LOG.info(f"Configuration modified via HTTP API using key '{state.id}'")
except ConfigSessionError as e:
@@ -411,16 +502,68 @@ def _configure_op(
status = 500
# Don't give the details away to the outer world
- error_msg = 'An internal error occured. Check the logs for details.'
+ error_msg = 'An internal error occurred. Check the logs for details.'
finally:
+ if 'IN_COMMIT_CONFIRM' in env:
+ del env['IN_COMMIT_CONFIRM']
lock.release()
+ # Background jobs return raw success text or raise on failure;
+ # the API wrapper formats HTTP responses and returns it
+ if is_background_job:
+ if status == 200:
+ return msg
+ else:
+ raise RuntimeError(error_msg)
+
if status != 200:
return error(status, error_msg)
return success(msg)
+async def _configure_op(
+ data: Union[
+ ConfirmModel,
+ ConfigureModel,
+ ConfigureListModel,
+ ConfigSectionModel,
+ ConfigSectionListModel,
+ ConfigSectionTreeModel,
+ ],
+ background_tasks: BackgroundTasks,
+ in_background: bool = False,
+):
+ """
+ API wrapper for configure operations.
+
+ If `in_background=True`: enqueue the whole configure
+ workflow and return an operation record immediately.
+ Otherwise: run the configure workflow in a threadpool
+ and return the normal API response.
+ """
+
+ if in_background:
+ try:
+ # Enqueue and return an operation handle that
+ # can be polled via `/retrieve/background-operations`
+ record = background_op_manager.enqueue(
+ background_tasks,
+ _execute_configure_op,
+ data,
+ )
+ except BackgroundOpError as e:
+ return error(500, str(e))
+
+ return success({'operation': record.model_dump()})
+
+ return await run_in_threadpool(
+ _execute_configure_op,
+ data,
+ background_tasks=background_tasks,
+ )
+
+
def create_path_import_pki_no_prompt(path):
correct_paths = ['ca', 'certificate', 'key-pair']
if path[1] not in correct_paths:
@@ -431,21 +574,23 @@ def create_path_import_pki_no_prompt(path):
@router.post('/configure')
-def configure_op(
- data: Union[ConfigureModel, ConfigureListModel],
+async def configure_op(
+ data: Union[ConfigureModel, ConfigureListModel, ConfirmModel],
request: Request,
background_tasks: BackgroundTasks,
+ in_background: bool = Query(False),
):
- return _configure_op(data, request, background_tasks)
+ return await _configure_op(data, background_tasks, in_background)
@router.post('/configure-section')
-def configure_section_op(
+async def configure_section_op(
data: Union[ConfigSectionModel, ConfigSectionListModel, ConfigSectionTreeModel],
request: Request,
background_tasks: BackgroundTasks,
+ in_background: bool = Query(False),
):
- return _configure_op(data, request, background_tasks)
+ return await _configure_op(data, background_tasks, in_background)
@router.post('/retrieve')
@@ -487,49 +632,98 @@ async def retrieve_op(data: RetrieveModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
+@router.post('/retrieve/background-operations')
+async def retrieve_background_operations(
+ op_id: str = Query(None),
+):
+ if op_id:
+ # Return only that record
+ record = background_op_manager.get_record(op_id)
+ records = [record] if record else []
+ else:
+ # Return the full in-memory operation history (oldest first)
+ records = background_op_manager.get_records()
+
+ result = {
+ 'operations': [record.model_dump() for record in records],
+ }
+
+ return success(result)
+
+
@router.post('/config-file')
-def config_file_op(data: ConfigFileModel, background_tasks: BackgroundTasks):
+async def config_file_op(data: ConfigFileModel, background_tasks: BackgroundTasks):
state = SessionState()
session = state.session
env = session.get_session_env()
op = data.op
msg = None
- try:
- if op == 'save':
- if data.file:
- path = data.file
- else:
- path = '/config/config.boot'
- msg = session.save_config(path)
- elif op == 'load':
- if data.file:
- path = data.file
- else:
- return error(400, 'Missing required field "file"')
+ # A non-zero confirm_time will start commit-confirm timer on commit
+ confirm_time = data.confirm_time
+
+ # Serialize config operations without blocking the event loop
+ async with asynclock:
+ try:
+ if op == 'save':
+ path = data.file or '/config/config.boot'
+ msg = session.save_config(path)
+
+ elif op in ('load', 'merge'):
+ if data.file:
+ path = data.file
+ elif data.string:
+ path = '/tmp/config.file'
+ with open(path, 'w') as f:
+ f.write(data.string)
+ else:
+ return error(400, 'Missing required field "file | string"')
+
+ match op:
+ case 'load':
+ session.migrate_and_load_config(path)
+ case 'merge':
+ session.merge_config(path, destructive=data.destructive)
- session.migrate_and_load_config(path)
+ config = Config(session_env=env)
+ d = get_config_diff(config)
- config = Config(session_env=env)
- d = get_config_diff(config)
+ state.confirm_time = confirm_time if confirm_time else 0
- if d.is_node_changed(['service', 'https']):
- background_tasks.add_task(call_commit, state)
- msg = self_ref_msg
+ if not d.is_node_changed(['service', 'https']):
+ if confirm_time:
+ out, err = await run_in_threadpool(run_commit_confirm, state)
+ else:
+ out, err = await run_in_threadpool(run_commit, state)
+
+ if err:
+ raise err
+ msg = (msg or '') + (out or '')
+ else:
+ if confirm_time:
+ background_tasks.add_task(call_commit_confirm, state)
+ else:
+ background_tasks.add_task(call_commit, state)
+ out = self_ref_msg
+ msg = (msg or '') + (out or '')
+ elif op == 'confirm':
+ msg = session.confirm()
else:
- session.commit()
- else:
- return error(400, f"'{op}' is not a valid operation")
- except ConfigSessionError as e:
- return error(400, str(e))
- except Exception:
- LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(400, f"'{op}' is not a valid operation")
+
+ except ConfigSessionError as e:
+ return error(400, str(e))
+ except Exception:
+ LOG.critical(traceback.format_exc())
+ return error(500, 'An internal error occurred. Check the logs for details.')
+ finally:
+ if 'IN_COMMIT_CONFIRM' in env:
+ del env['IN_COMMIT_CONFIRM']
return success(msg)
@@ -554,7 +748,7 @@ def image_op(data: ImageModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
@@ -587,7 +781,7 @@ def container_image_op(data: ContainerImageModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
@@ -603,13 +797,14 @@ def generate_op(data: GenerateModel):
try:
if op == 'generate':
res = session.generate(path)
+ session.commit()
else:
return error(400, f"'{op}' is not a valid operation")
except ConfigSessionError as e:
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
@@ -631,7 +826,7 @@ def show_op(data: ShowModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
@@ -653,10 +848,30 @@ def reboot_op(data: RebootModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
+@router.post('/renew')
+def renew_op(data: RenewModel):
+ state = SessionState()
+ session = state.session
+
+ op = data.op
+ path = data.path
+
+ try:
+ if op == 'renew':
+ res = session.renew(path)
+ else:
+ return error(400, f"'{op}' is not a valid operation")
+ except ConfigSessionError as e:
+ return error(400, str(e))
+ except Exception:
+ LOG.critical(traceback.format_exc())
+ return error(500, 'An internal error occurred. Check the logs for details.')
+
+ return success(res)
@router.post('/reset')
def reset_op(data: ResetModel):
@@ -675,7 +890,7 @@ def reset_op(data: ResetModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
@@ -715,7 +930,7 @@ def import_pki(data: ImportPkiModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
finally:
lock.release()
@@ -739,7 +954,7 @@ def poweroff_op(data: PoweroffModel):
return error(400, str(e))
except Exception:
LOG.critical(traceback.format_exc())
- return error(500, 'An internal error occured. Check the logs for details.')
+ return error(500, 'An internal error occurred. Check the logs for details.')
return success(res)
diff --git a/src/services/api/session.py b/src/services/api/session.py
index ad3ef660c..c25a444e9 100644
--- a/src/services/api/session.py
+++ b/src/services/api/session.py
@@ -1,4 +1,4 @@
-# Copyright 2024 VyOS maintainers and contributors <maintainers@vyos.io>
+# Copyright VyOS maintainers and contributors <maintainers@vyos.io>
#
# This library is free software; you can redistribute it and/or
# modify it under the terms of the GNU Lesser General Public