/
/
1"""Setup flow engine: registry and API commands for interactive provider/player setup."""
2
3from __future__ import annotations
4
5import asyncio
6import logging
7import time
8from dataclasses import dataclass
9from functools import partial
10from typing import TYPE_CHECKING, Any
11from uuid import uuid4
12
13from music_assistant_models.auth import Scope
14from music_assistant_models.config_entries import ConfigEntry, ConfigValueOption
15from music_assistant_models.enums import ConfigEntryType, EventType, FlowStepType, MediaType
16from music_assistant_models.errors import (
17 ActionUnavailable,
18 InsufficientPermissions,
19 SetupFailedError,
20)
21from music_assistant_models.setup_flow import SetupFlowStep
22
23from music_assistant.constants import CONF_PLAYERS, CONF_PROVIDERS
24from music_assistant.controllers.config.helpers import _AUTH_ERROR_CODES
25from music_assistant.helpers.api import api_command
26from music_assistant.helpers.util import import_module_in_thread, load_provider_module
27from music_assistant.models.player import Player
28from music_assistant.models.setup_flow import (
29 AbortFlow,
30 FlowReason,
31 SetupFlowContext,
32 SetupFlowError,
33 SetupSession,
34 StepExpiredError,
35)
36
37if TYPE_CHECKING:
38 from collections.abc import Awaitable, Callable
39
40 from music_assistant_models.config_entries import ConfigValueType, ProviderConfig
41 from music_assistant_models.provider import ProviderManifest
42
43 from music_assistant import MusicAssistant
44
45LOGGER = logging.getLogger(__name__)
46
47# a flow with no client interaction or step changes for this long is garbage collected
48IDLE_FLOW_TTL = 15 * 60
49FLOW_SWEEP_INTERVAL = 60
50# how long start/submit wait for the flow coroutine to produce the (next) step;
51# generous because finish() may install requirements and load the provider
52NEXT_STEP_TIMEOUT = 120
53# bound on waiting for a cancelled flow's cleanup (author finally-blocks)
54FLOW_ABORT_CLEANUP_TIMEOUT = 10
55
56
57@dataclass
58class ActiveSetupFlow:
59 """Registry record for a running setup flow."""
60
61 session: SetupSession
62 # target_key: identifies what the flow is (re)configuring; one flow per target
63 target_key: str
64 # required_scope: the scope the starting command required; re-checked on
65 # every flows/* continuation command
66 required_scope: Scope
67 task: asyncio.Task[None] | None = None
68
69
70class SetupFlowMixin:
71 """Mixin providing the setup flow engine for the ConfigController."""
72
73 # registry of running flows, keyed by flow_id (lazily created per instance)
74 _flows: dict[str, ActiveSetupFlow] | None = None
75 _flow_sweep_handle: asyncio.TimerHandle | None = None
76 # required scopes of recently finished flows: terminal steps can publish after
77 # the registry pop (cancel-driven aborts), and the event-scope filter still
78 # needs to resolve them; bounded FIFO
79 _finished_flow_scopes: dict[str, Scope] | None = None
80
81 # Type hints for attributes/methods provided by the class this mixin is used with
82 if TYPE_CHECKING:
83 mass: MusicAssistant
84
85 @property
86 def onboard_done(self) -> bool: ... # noqa: D102
87
88 def get(self, key: str, default: Any = None) -> Any: ... # noqa: D102
89
90 def set(self, key: str, value: Any, immediate: bool = False) -> None: ... # noqa: D102
91
92 def encrypt_string(self, str_value: str) -> str: ... # noqa: D102
93
94 def decrypt_string(self, encrypted_str: str) -> str: ... # noqa: D102
95
96 async def get_provider_configs( # noqa: D102
97 self, provider_type: Any = None, provider_domain: str | None = None
98 ) -> list[ProviderConfig]: ...
99
100 async def get_provider_config(self, instance_id: str) -> ProviderConfig: ... # noqa: D102
101
102 def update_provider_last_error(self, instance_id: str, error: Any) -> None: ... # noqa: D102
103
104 async def get_player_config(self, player_id: str) -> Any: ... # noqa: D102
105
106 async def _create_provider_instance(
107 self,
108 provider_domain: str,
109 values: dict[str, ConfigValueType],
110 setup_data: dict[str, Any] | None = None,
111 ) -> ProviderConfig: ...
112
113 @api_command("config/providers/setup", required_scope=Scope.CONFIG_PROVIDERS_WRITE)
114 async def setup_provider(self, provider_domain: str) -> SetupFlowStep:
115 """
116 Start the setup flow to add a new instance of the given provider.
117
118 For a provider without a setup flow (no user input needed) the instance is
119 created right away and a FINISH step is returned, so the add-provider path
120 is uniform for all providers.
121
122 :param provider_domain: Domain of the provider to add an instance of.
123 """
124 for manifest in self.mass.get_provider_manifests():
125 if manifest.domain == provider_domain:
126 break
127 else:
128 msg = f"Unknown provider domain: {provider_domain}"
129 raise KeyError(msg)
130 owner = f"provider.{provider_domain}"
131 # fail fast on conditions that would otherwise only surface at save
132 existing = await self.get_provider_configs(provider_domain=provider_domain)
133 if existing and not manifest.multi_instance:
134 return self._synthesized_step(FlowStepType.ABORT, owner, reason="already_configured")
135 if manifest.depends_on:
136 dep_configs = await self.get_provider_configs(provider_domain=manifest.depends_on)
137 if not any(dep_conf.enabled for dep_conf in dep_configs):
138 return self._synthesized_step(
139 FlowStepType.ABORT, owner, reason="missing_dependency"
140 )
141 flow_module = await self._get_setup_flow_module(manifest)
142 if flow_module is None:
143 # zero-input provider: create the instance immediately and report it
144 # as an (already) finished flow
145 config = await self._create_provider_instance(provider_domain, {})
146 return self._synthesized_step(
147 FlowStepType.FINISH,
148 owner,
149 step_id=self._provider_finish_step_id(config.instance_id),
150 result={"instance_id": config.instance_id},
151 )
152 context = SetupFlowContext(kind="setup", reason="user", domain=provider_domain)
153 return await self._start_flow(
154 flow_coro=flow_module.run_setup,
155 context=context,
156 target_key=f"provider_setup:{provider_domain}",
157 required_scope=Scope.CONFIG_PROVIDERS_WRITE,
158 finish_handler=self._finish_provider_setup,
159 )
160
161 @api_command("config/providers/reconfigure", required_scope=Scope.CONFIG_PROVIDERS_WRITE)
162 async def reconfigure_provider(self, instance_id: str) -> SetupFlowStep:
163 """
164 Start the reconfigure flow on an existing provider instance (covers reauth).
165
166 :param instance_id: The provider instance to reconfigure.
167 """
168 raw_conf = self.get(f"{CONF_PROVIDERS}/{instance_id}")
169 if not raw_conf:
170 msg = f"No config found for provider id {instance_id}"
171 raise KeyError(msg)
172 domain: str = raw_conf["domain"]
173 manifest = self.mass.get_provider_manifest(domain)
174 owner = f"provider.{domain}"
175 flow_module = await self._get_setup_flow_module(manifest)
176 if flow_module is None:
177 # flow-less providers have nothing to reconfigure;
178 # their failures are environmental (reload/retry covers them)
179 return self._synthesized_step(FlowStepType.ABORT, owner, reason="nothing_to_configure")
180 context = SetupFlowContext(
181 kind="reconfigure",
182 reason=self._reconfigure_reason(raw_conf.get("last_error")),
183 domain=domain,
184 instance_id=instance_id,
185 setup_data=self._decrypt_values(raw_conf.get("setup_data") or {}),
186 values=self._decrypt_values(raw_conf.get("values") or {}),
187 )
188 return await self._start_flow(
189 flow_coro=flow_module.run_setup,
190 context=context,
191 target_key=f"provider_reconfigure:{instance_id}",
192 required_scope=Scope.CONFIG_PROVIDERS_WRITE,
193 finish_handler=self._finish_provider_reconfigure,
194 )
195
196 @api_command("config/players/setup", required_scope=Scope.CONFIG_PLAYERS_WRITE)
197 async def setup_player(self, player_id: str) -> SetupFlowStep:
198 """
199 Start the setup flow for a player (e.g. pairing).
200
201 A player that itself needs no setup but wraps protocol child player(s) that do
202 (universal players, or native players wrapping protocol children) delegates to
203 the child's setup flow: to the single child that needs setup directly, or - when
204 more than one does - via a form that lets the user pick which child to set up.
205 The child's flow persists to the child's own config.
206
207 Also serves on-demand re-runs: when nothing needs setup (anymore), delegation
208 falls back to any child that merely has a flow, so a step the user skipped
209 earlier - an optional pairing, say - remains reachable.
210
211 :param player_id: The player to set up.
212 """
213 # deliberately no raise_unavailable: a player that needs setup is serialized
214 # as unavailable, and that is exactly the player this command targets
215 player = self.mass.players.get_player(player_id)
216 if player is None:
217 msg = f"Player {player_id} not found"
218 raise KeyError(msg)
219 owner = f"provider.{player.provider.domain}"
220 target_key = f"player_setup:{player_id}"
221 if player.implements_setup_flow:
222 # the player implements its own setup flow: run it directly
223 return await self._start_flow(
224 flow_coro=player.run_setup_flow,
225 context=self._player_flow_context(player),
226 target_key=target_key,
227 required_scope=Scope.CONFIG_PLAYERS_WRITE,
228 finish_handler=self._finish_player_setup,
229 )
230 # no direct setup: delegate to protocol child player(s), preferring the ones
231 # that actually need setup and falling back to any that can re-run their flow
232 children = self._protocol_children_with_setup_flow(player, needing_only=True)
233 if not children:
234 children = self._protocol_children_with_setup_flow(player, needing_only=False)
235 if len(children) == 1:
236 child = children[0]
237 return await self._start_flow(
238 flow_coro=child.run_setup_flow,
239 context=self._player_flow_context(child),
240 # key on the child: a direct setup of the child must replace this flow
241 target_key=f"player_setup:{child.player_id}",
242 required_scope=Scope.CONFIG_PLAYERS_WRITE,
243 finish_handler=self._finish_player_setup,
244 )
245 if children:
246 return await self._start_flow(
247 flow_coro=partial(self._run_child_selection_flow, children),
248 context=self._player_flow_context(player),
249 target_key=target_key,
250 required_scope=Scope.CONFIG_PLAYERS_WRITE,
251 finish_handler=self._finish_player_setup,
252 )
253 # nothing on this player (or its children) to configure
254 return self._synthesized_step(FlowStepType.ABORT, owner, reason="nothing_to_configure")
255
256 @api_command("config/flows/submit")
257 async def submit_setup_flow(
258 self, flow_id: str, values: dict[str, ConfigValueType]
259 ) -> SetupFlowStep:
260 """
261 Submit the user's values for the flow's pending FORM step.
262
263 Returns the flow's next step, or the same FORM step (with per-field errors
264 set) when validation failed.
265
266 :param flow_id: The id of the running flow.
267 :param values: The raw values for the form's config entries.
268 """
269 flow = self._get_flow(flow_id)
270 self._check_flow_permission(flow)
271 if (error_step := flow.session.handle_submit(values)) is not None:
272 return error_step
273 # wait (bounded) for the coroutine to produce the next step
274 submitted_step = flow.session.current_step
275 await flow.session.wait_for_step_change(NEXT_STEP_TIMEOUT)
276 step = flow.session.current_step
277 assert step is not None # an accepted submit implies a published FORM step
278 if step is submitted_step:
279 # rare: the coroutine is still working on the next step. The submitted
280 # form's input future is already consumed, so re-serving the form would
281 # invite a doomed resubmit - publish a progress step (so flows/get agrees)
282 # and let the coroutine's next publish deliver the real step
283 flow.session.progress("working")
284 step = flow.session.current_step
285 assert step is not None
286 return step
287
288 @api_command("config/flows/get")
289 async def get_setup_flow(self, flow_id: str) -> SetupFlowStep:
290 """
291 Return the current step of a running flow (idempotent re-render, never advances).
292
293 :param flow_id: The id of the running flow.
294 """
295 flow = self._get_flow(flow_id)
296 self._check_flow_permission(flow)
297 flow.session.last_activity = time.monotonic()
298 if (step := flow.session.current_step) is None:
299 raise ActionUnavailable("The setup flow has not produced a step yet")
300 return step
301
302 @api_command("config/flows/abort")
303 async def abort_setup_flow(self, flow_id: str) -> None:
304 """
305 Abort a running flow (user cancelled).
306
307 :param flow_id: The id of the running flow.
308 """
309 flow = self._get_flow(flow_id)
310 self._check_flow_permission(flow)
311 await self._abort_flow(flow, reason="aborted")
312
313 def get_setup_flow_required_scope(self, flow_id: str) -> Scope | None:
314 """
315 Return the scope required to receive/interact with the given setup flow.
316
317 Also resolves recently finished flows (their terminal step can publish
318 just after the registry pop). Returns None when the flow is unknown.
319
320 :param flow_id: The id of the flow.
321 """
322 if flow := self._setup_flows.get(flow_id):
323 return flow.required_scope
324 if self._finished_flow_scopes:
325 return self._finished_flow_scopes.get(flow_id)
326 return None
327
328 def _pop_flow(self, flow: ActiveSetupFlow) -> None:
329 """Remove a flow from the registry, retaining its scope for late events."""
330 self._setup_flows.pop(flow.session.flow_id, None)
331 if self._finished_flow_scopes is None:
332 self._finished_flow_scopes = {}
333 finished = self._finished_flow_scopes
334 finished[flow.session.flow_id] = flow.required_scope
335 while len(finished) > 64:
336 finished.pop(next(iter(finished)))
337
338 async def _start_flow(
339 self,
340 *,
341 flow_coro: Callable[[SetupSession], Awaitable[Any]],
342 context: SetupFlowContext,
343 target_key: str,
344 required_scope: Scope,
345 finish_handler: Callable[
346 [SetupSession, dict[str, ConfigValueType]], Awaitable[dict[str, str]]
347 ],
348 ) -> SetupFlowStep:
349 """Register and start a new flow, returning its first published step."""
350 # one flow per target: starting anew replaces (aborts) a lingering previous flow.
351 # re-scan after every await: the abort yields, so a concurrent start for the same
352 # target may have registered a new flow in the meantime
353 while existing_flow := next(
354 (f for f in self._setup_flows.values() if f.target_key == target_key), None
355 ):
356 await self._abort_flow(existing_flow, reason="replaced")
357 flow_id = uuid4().hex
358 session = SetupSession(self.mass, flow_id, context, finish_handler)
359 flow = ActiveSetupFlow(
360 session=session, target_key=target_key, required_scope=required_scope
361 )
362 self._setup_flows[flow_id] = flow
363 LOGGER.debug("Starting setup flow %s for %s", flow_id, target_key)
364 flow.task = self.mass.create_task(self._run_flow(flow, flow_coro))
365 self._schedule_flow_sweep()
366 if session.current_step is None:
367 await session.wait_for_step_change(NEXT_STEP_TIMEOUT)
368 if (step := session.current_step) is None:
369 await self._abort_flow(flow, reason="internal_error")
370 raise SetupFailedError("Setup flow did not produce a first step")
371 return step
372
373 async def _run_flow(
374 self, flow: ActiveSetupFlow, flow_coro: Callable[[SetupSession], Awaitable[Any]]
375 ) -> None:
376 """Drive the flow coroutine and convert its outcome into a terminal step."""
377 session = flow.session
378 try:
379 await flow_coro(session)
380 except AbortFlow as err:
381 session.publish_abort(err.reason)
382 except StepExpiredError:
383 session.publish_abort("timed_out")
384 except SetupFlowError as err:
385 # the author did not catch a finish failure: end with the failure message
386 session.publish_abort(str(err) or "internal_error")
387 except asyncio.CancelledError:
388 # abort/replace/shutdown: the author's cleanup (finally blocks) has run;
389 # the canceller publishes the ABORT step. Never swallow the cancellation.
390 raise
391 except Exception:
392 LOGGER.exception("Unhandled error in setup flow for %s", session.context.domain)
393 session.publish_abort("internal_error")
394 else:
395 if not session.finished:
396 LOGGER.error(
397 "Setup flow for %s returned without calling finish()", session.context.domain
398 )
399 session.publish_abort("internal_error")
400 finally:
401 LOGGER.debug("Setup flow %s ended", session.flow_id)
402 session.close()
403 self._pop_flow(flow)
404
405 async def _abort_flow(self, flow: ActiveSetupFlow, reason: str) -> None:
406 """
407 Abort the given flow with the given reason and clean it up.
408
409 Cancelling raises CancelledError inside the flow coroutine so the author's
410 cleanup (try/finally around pairing sessions etc.) runs before the terminal
411 ABORT step goes out.
412 """
413 if flow.task is not None and not flow.task.done():
414 flow.task.cancel()
415 # wait() shields us from the task's CancelledError without
416 # masking a cancellation of the caller itself; the timeout keeps a
417 # wedged author cleanup (e.g. a hanging pairing teardown) from
418 # stalling the abort and any replacement flow indefinitely
419 _, pending = await asyncio.wait([flow.task], timeout=FLOW_ABORT_CLEANUP_TIMEOUT)
420 if pending:
421 LOGGER.warning(
422 "Setup flow for %s did not clean up within %ss after cancellation",
423 flow.session.context.domain,
424 FLOW_ABORT_CLEANUP_TIMEOUT,
425 )
426 # the wedged task never reaches _run_flow's finally: close the
427 # session here so the unauthenticated callback route is dropped
428 flow.session.close()
429 self._pop_flow(flow)
430 current_step = flow.session.current_step
431 if current_step is None or current_step.type not in (
432 FlowStepType.FINISH,
433 FlowStepType.ABORT,
434 ):
435 flow.session.publish_abort(reason)
436
437 async def _finish_provider_setup(
438 self, session: SetupSession, values: dict[str, ConfigValueType]
439 ) -> dict[str, str]:
440 """Finish handler for provider setup flows: create, persist and load the instance."""
441 try:
442 config = await self._create_provider_instance(
443 session.context.domain, {}, setup_data=self._encrypt_values(values)
444 )
445 except Exception as err:
446 raise SetupFlowError(
447 str(err) or err.__class__.__name__,
448 translation_key=getattr(err, "translation_key", None),
449 ) from err
450 session.finish_step_id = self._provider_finish_step_id(config.instance_id)
451 return {"instance_id": config.instance_id}
452
453 async def _finish_provider_reconfigure(
454 self, session: SetupSession, values: dict[str, ConfigValueType]
455 ) -> dict[str, str]:
456 """Finish handler for provider reconfigure flows: merge setup_data and reload."""
457 instance_id = session.context.instance_id
458 assert instance_id is not None # always set for reconfigure flows
459 conf_key = f"{CONF_PROVIDERS}/{instance_id}"
460 raw_conf = self.get(conf_key)
461 if not raw_conf:
462 raise SetupFlowError(f"Provider {instance_id} no longer exists")
463 snapshot = dict(raw_conf.get("setup_data") or {})
464 self.set(f"{conf_key}/setup_data", {**snapshot, **self._encrypt_values(values)})
465 try:
466 config = await self.get_provider_config(instance_id)
467 await self.mass.load_provider_config(config)
468 except asyncio.CancelledError:
469 self.set(f"{conf_key}/setup_data", snapshot)
470 raise
471 except Exception as err:
472 # reloading with the new values failed: restore the previous setup_data
473 self.set(f"{conf_key}/setup_data", snapshot)
474 raise SetupFlowError(
475 str(err) or err.__class__.__name__,
476 translation_key=getattr(err, "translation_key", None),
477 ) from err
478 self.update_provider_last_error(instance_id, None)
479 return {"instance_id": instance_id}
480
481 async def _finish_player_setup(
482 self, session: SetupSession, values: dict[str, ConfigValueType]
483 ) -> dict[str, str]:
484 """Finish handler for player setup flows: persist and apply the collected setup data."""
485 player_id = session.context.player_id
486 assert player_id is not None # always set for player flows
487 conf_key = f"{CONF_PLAYERS}/{player_id}"
488 raw_conf = self.get(conf_key)
489 if not raw_conf:
490 raise SetupFlowError(f"No config found for player {player_id}")
491 snapshot = dict(raw_conf.get("setup_data") or {})
492 self.set(f"{conf_key}/setup_data", {**snapshot, **self._encrypt_values(values)})
493 try:
494 config = await self.get_player_config(player_id)
495 changed_keys = {f"setup_data/{key}" for key in values}
496 await self.mass.players.on_player_config_change(config, changed_keys)
497 except asyncio.CancelledError:
498 self.set(f"{conf_key}/setup_data", snapshot)
499 raise
500 except Exception as err:
501 # reading back or applying the updated config failed: restore the previous setup_data
502 self.set(f"{conf_key}/setup_data", snapshot)
503 raise SetupFlowError(
504 str(err) or err.__class__.__name__,
505 translation_key=getattr(err, "translation_key", None),
506 ) from err
507 self.mass.signal_event(EventType.PLAYER_CONFIG_UPDATED, object_id=player_id, data=config)
508 return {"player_id": player_id}
509
510 def _player_flow_context(self, player: Player) -> SetupFlowContext:
511 """Build the setup flow context (with decrypted prefill) for the given player."""
512 raw_conf = self.get(f"{CONF_PLAYERS}/{player.player_id}") or {}
513 return SetupFlowContext(
514 kind="setup",
515 reason="user",
516 domain=player.provider.domain,
517 instance_id=player.provider.instance_id,
518 player_id=player.player_id,
519 setup_data=self._decrypt_values(raw_conf.get("setup_data") or {}),
520 values=self._decrypt_values(raw_conf.get("values") or {}),
521 )
522
523 def _protocol_children_with_setup_flow(
524 self, player: Player, *, needing_only: bool
525 ) -> list[Player]:
526 """
527 Return the player's protocol child players that implement a setup flow.
528
529 Covers the wrapper case: a universal player, or a native player wrapping
530 protocol children, whose own setup is a no-op but whose linked protocol
531 outputs still require pairing/credentials.
532
533 :param player: The (wrapper) player whose protocol children to inspect.
534 :param needing_only: Only return children that currently need setup.
535 """
536 children: list[Player] = []
537 seen: set[str] = set()
538 for output_protocol in player.output_protocols:
539 child_id = output_protocol.output_protocol_id
540 if output_protocol.is_native or child_id in seen:
541 continue
542 seen.add(child_id)
543 child = self.mass.players.get_player(child_id)
544 if child is None or not child.implements_setup_flow:
545 continue
546 if needing_only and not child.needs_setup:
547 continue
548 children.append(child)
549 return children
550
551 async def _run_child_selection_flow(
552 self, children: list[Player], session: SetupSession
553 ) -> None:
554 """
555 Run the wrapper flow that lets the user pick which protocol child to set up.
556
557 The selection form is owned by the parent; once a child is picked the session is
558 re-pointed at that child so its flow's steps localize and its ``finish()`` persists
559 under the child's own config.
560 """
561 options = [
562 ConfigValueOption(
563 value=child.player_id, title=f"{child.display_name} ({child.provider.name})"
564 )
565 for child in children
566 ]
567 values = await session.form(
568 [
569 ConfigEntry(
570 key="child",
571 type=ConfigEntryType.STRING,
572 required=True,
573 default_value=children[0].player_id,
574 options=options,
575 )
576 ],
577 step_id="select_child",
578 )
579 child_id = str(values["child"])
580 child = next((candidate for candidate in children if candidate.player_id == child_id), None)
581 if child is None:
582 raise AbortFlow("nothing_to_configure")
583 child_raw = self.get(f"{CONF_PLAYERS}/{child_id}") or {}
584 session.retarget(
585 domain=child.provider.domain,
586 instance_id=child.provider.instance_id,
587 player_id=child_id,
588 setup_data=self._decrypt_values(child_raw.get("setup_data") or {}),
589 values=self._decrypt_values(child_raw.get("values") or {}),
590 )
591 await child.run_setup_flow(session)
592
593 async def _get_setup_flow_module(self, manifest: ProviderManifest) -> Any | None:
594 """
595 Import (lazily) the provider's setup_flow module, or None when it has none.
596
597 A provider without a setup_flow module needs no setup input at all.
598 """
599 # ensure the provider's requirements are installed and its package imports
600 # cleanly first: the setup_flow submodule may rely on those requirements
601 await load_provider_module(manifest.domain, manifest.requirements)
602 module_path = f"music_assistant.providers.{manifest.domain}.setup_flow"
603 try:
604 return await import_module_in_thread(module_path)
605 except ModuleNotFoundError as err:
606 if err.name == module_path:
607 # the provider ships no setup_flow module: it needs no setup input
608 return None
609 # an import *inside* setup_flow.py failed: an actual bug, surface it
610 raise
611
612 @property
613 def _setup_flows(self) -> dict[str, ActiveSetupFlow]:
614 """Return the registry of running flows (created lazily)."""
615 if self._flows is None:
616 self._flows = {}
617 return self._flows
618
619 def _get_flow(self, flow_id: str) -> ActiveSetupFlow:
620 """Return the running flow for the given id."""
621 if flow := self._setup_flows.get(flow_id):
622 return flow
623 msg = f"Unknown (or finished) setup flow: {flow_id}"
624 raise KeyError(msg)
625
626 def _check_flow_permission(self, flow: ActiveSetupFlow) -> None:
627 """Verify the calling user holds the scope the flow's start command required."""
628 # imported here: the webserver helpers pull in the full auth stack,
629 # which must not be imported with the config controller at startup
630 from music_assistant.controllers.webserver.helpers.auth_middleware import ( # noqa: PLC0415
631 get_current_user,
632 has_scope,
633 )
634
635 user = get_current_user()
636 # no user context means an internal (server-side) caller, which is trusted
637 if user is not None and not has_scope(user, flow.required_scope):
638 raise InsufficientPermissions(
639 f"This action requires the {flow.required_scope.value} scope"
640 )
641
642 def _synthesized_step(
643 self,
644 step_type: FlowStepType,
645 translation_owner: str,
646 *,
647 step_id: str | None = None,
648 result: dict[str, str] | None = None,
649 reason: str | None = None,
650 ) -> SetupFlowStep:
651 """
652 Return a terminal step for a flow that ended before a session was needed.
653
654 :param step_type: The terminal step type (FINISH or ABORT).
655 :param translation_owner: The namespace the step's strings resolve under.
656 :param step_id: Slug to serve the step's strings under; defaults to the one
657 implied by the step type.
658 :param result: Reference to the created/updated object (FINISH).
659 :param reason: Slug describing why the flow was aborted (ABORT).
660 """
661 default_step_id = "finish" if step_type == FlowStepType.FINISH else "abort"
662 return SetupFlowStep(
663 flow_id=uuid4().hex,
664 step_id=step_id or default_step_id,
665 type=step_type,
666 result=result,
667 reason=reason,
668 translation_owner=translation_owner,
669 )
670
671 def _provider_finish_step_id(self, instance_id: str) -> str:
672 """
673 Return the i18n slug of the FINISH step for a newly set up provider instance.
674
675 :param instance_id: The provider instance the flow just created.
676 """
677 # a provider that imports a library gets the variant explaining that the first import
678 # runs in the background, so a library that still looks empty right after setup does
679 # not read as a failed setup
680 provider = self.mass.get_provider(instance_id, return_unavailable=True)
681 if provider is not None and any(
682 self.mass.music.library_supported(provider, media_type) for media_type in MediaType.ALL
683 ):
684 return "finish_library_sync"
685 return "finish"
686
687 def _reconfigure_reason(self, last_error: Any) -> FlowReason:
688 """Derive the reconfigure flow reason from the provider's stored last_error."""
689 if not last_error:
690 return "user"
691 # legacy settings may still hold a plain string last_error
692 error_code = last_error.get("error_code") if isinstance(last_error, dict) else None
693 if error_code in _AUTH_ERROR_CODES:
694 return "auth"
695 return "error"
696
697 def _encrypt_values(self, values: dict[str, ConfigValueType]) -> dict[str, Any]:
698 """Return a copy of the values with all string values encrypted (at-rest form)."""
699 return {
700 key: self.encrypt_string(value) if isinstance(value, str) else value
701 for key, value in values.items()
702 }
703
704 def _decrypt_values(self, values: dict[str, Any]) -> dict[str, Any]:
705 """Return a copy of the values with all (encrypted) string values decrypted."""
706 return {
707 key: self.decrypt_string(value) if isinstance(value, str) else value
708 for key, value in values.items()
709 }
710
711 def _schedule_flow_sweep(self) -> None:
712 """Arm the periodic idle-flow sweeper (idempotent)."""
713 if self._flow_sweep_handle is not None:
714 return
715 self._flow_sweep_handle = self.mass.loop.call_later(
716 FLOW_SWEEP_INTERVAL, self._sweep_idle_flows
717 )
718
719 def _sweep_idle_flows(self) -> None:
720 """Abort flows that have been idle for longer than the TTL."""
721 self._flow_sweep_handle = None
722 if self.mass.closing:
723 return
724 now = time.monotonic()
725 for flow in list(self._setup_flows.values()):
726 current_step = flow.session.current_step
727 if (
728 current_step is not None
729 and current_step.expires_at is not None
730 and current_step.expires_at > time.time()
731 ):
732 # the step advertises a (longer) countdown to the user; the step
733 # deadline machinery guarantees the flow terminates on its own
734 continue
735 if now - flow.session.last_activity >= IDLE_FLOW_TTL:
736 self.mass.create_task(self._abort_flow(flow, "timed_out"))
737 if self._setup_flows:
738 self._schedule_flow_sweep()
739