diff options
Diffstat (limited to 'src/services/api')
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 |
