diff --git a/src/app/guards/admin-guard.ts b/src/app/guards/admin-guard.ts
index 6cf1199..eb8da11 100644
--- a/src/app/guards/admin-guard.ts
+++ b/src/app/guards/admin-guard.ts
@@ -1,20 +1,9 @@
-import { inject } from '@angular/core';
-import { CanActivateFn, Router } from '@angular/router';
-import { Authorization } from '@services/authorization/authorization';
-import { map } from 'rxjs';
+import { CanActivateFn } from '@angular/router';
+import { createSessionGuard } from './session-guard';
-export const adminGuard: CanActivateFn = (_route, state) => {
- const authorization = inject(Authorization);
- const router = inject(Router);
-
- return authorization.validateSession().pipe(
- map((user) => {
- if (user?.role === 'ADMIN') {
- return true;
- }
-
- console.warn(`[adminGuard] Access denied to ${state.url} — user is not admin, redirecting to /`);
- return router.parseUrl('/');
- })
- );
-};
+export const adminGuard: CanActivateFn = createSessionGuard(
+ (user) => user?.role === 'ADMIN',
+ '/',
+ 'adminGuard',
+ 'user is not admin'
+);
diff --git a/src/app/guards/auth-guard.ts b/src/app/guards/auth-guard.ts
index aa25a10..2e884fc 100644
--- a/src/app/guards/auth-guard.ts
+++ b/src/app/guards/auth-guard.ts
@@ -1,20 +1,9 @@
-import { inject } from '@angular/core';
-import { CanActivateFn, Router } from '@angular/router';
-import { Authorization } from '@services/authorization/authorization';
-import { map } from 'rxjs';
+import { CanActivateFn } from '@angular/router';
+import { createSessionGuard } from './session-guard';
-export const authGuard: CanActivateFn = (_route, state) => {
- const authService = inject(Authorization);
- const router = inject(Router);
-
- return authService.validateSession().pipe(
- map((user) => {
- if (user) {
- return true;
- }
-
- console.warn(`[authGuard] Access denied to ${state.url} — user not logged in, redirecting to /login`);
- return router.parseUrl('/login');
- })
- );
-};
+export const authGuard: CanActivateFn = createSessionGuard(
+ (user) => !!user,
+ '/login',
+ 'authGuard',
+ 'user not logged in'
+);
diff --git a/src/app/guards/session-guard.ts b/src/app/guards/session-guard.ts
new file mode 100644
index 0000000..8414908
--- /dev/null
+++ b/src/app/guards/session-guard.ts
@@ -0,0 +1,33 @@
+import { inject } from '@angular/core';
+import { CanActivateFn, Router } from '@angular/router';
+import { Authorization } from '@services/authorization/authorization';
+import { User } from '@models/user';
+import { map } from 'rxjs';
+
+/**
+ * Builds a route guard that validates the session and redirects when `predicate`
+ * rejects the logged-in user (or there is none). Shared by authGuard/adminGuard
+ * so the two only differ in who is allowed through and where they land otherwise.
+ */
+export function createSessionGuard(
+ predicate: (user: User | null) => boolean,
+ redirectUrl: string,
+ logLabel: string,
+ denialReason: string
+): CanActivateFn {
+ return (_route, state) => {
+ const authorization = inject(Authorization);
+ const router = inject(Router);
+
+ return authorization.validateSession().pipe(
+ map((user) => {
+ if (predicate(user)) {
+ return true;
+ }
+
+ console.warn(`[${logLabel}] Access denied to ${state.url} — ${denialReason}, redirecting to ${redirectUrl}`);
+ return router.parseUrl(redirectUrl);
+ })
+ );
+ };
+}
diff --git a/src/app/layouts/app-layout/app-layout.ts b/src/app/layouts/app-layout/app-layout.ts
index 2a056e7..31be555 100644
--- a/src/app/layouts/app-layout/app-layout.ts
+++ b/src/app/layouts/app-layout/app-layout.ts
@@ -10,6 +10,7 @@ import { Authorization } from '@services/authorization/authorization';
import { BlocksService } from '@services/blocks/blocks';
import { ContainersService } from '@services/containers/containers';
import { ChangePasswordDialogComponent } from '@shared/change-password-dialog/change-password-dialog';
+import { scheduleSignalClear } from '@utilities/temporary-signal';
@Component({
selector: 'app-app-layout',
@@ -79,9 +80,7 @@ logout() {
this.changePasswordOpen.set(false);
this.changePasswordError.set(null);
this.changePasswordSuccess.set('Password changed successfully.');
- setTimeout(() => {
- this.changePasswordSuccess.set(null);
- }, 3000);
+ scheduleSignalClear(this.changePasswordSuccess);
},
error: (error) => {
this.changePasswordSaving.set(false);
diff --git a/src/app/layouts/tasks-executor/tasks-executor.ts b/src/app/layouts/tasks-executor/tasks-executor.ts
index 6f9ce55..4d5ff99 100644
--- a/src/app/layouts/tasks-executor/tasks-executor.ts
+++ b/src/app/layouts/tasks-executor/tasks-executor.ts
@@ -9,6 +9,7 @@ import {
TasksExecutionsListComponent
} from '@shared/tasks-executions-list/tasks-executions-list';
import { TaskExecutionViewerComponent } from '@shared/task-execution-viewer/task-execution-viewer';
+import { formatDuration } from '@shared/task-execution-viewer/execution-viewer.utils';
import { BlocksService } from '@services/blocks/blocks';
import { ContainersService } from '@services/containers/containers';
import { ConfirmDialogService } from '@services/dialogs/confirm-dialog';
@@ -151,7 +152,7 @@ export class TasksExecutor {
creationTime: execution.creationTime,
runNumber: typeof execution.runNumber === 'number' ? execution.runNumber : fallbackRunNumber,
rerunOfExecutionId: execution.rerunOfExecutionId ?? null,
- duration: this.formatDuration(execution.context.startTime ?? null, execution.context.endTime ?? null),
+ duration: this.formatExecutionDuration(execution.context.startTime ?? null, execution.context.endTime ?? null),
simulated: execution.interactionSimulationEnabled === true
};
}
@@ -166,29 +167,8 @@ export class TasksExecutor {
return `${yyyy}-${mm}-${dd} ${hh}:${mi}`;
}
- private formatDuration(startTime: number | null, endTime: number | null): string {
+ private formatExecutionDuration(startTime: number | null, endTime: number | null): string {
if (!startTime || !endTime) return '0 sec';
- const diffMs = Math.max(0, endTime - startTime);
- const totalSeconds = Math.floor(diffMs / 1000);
- const totalMinutes = Math.floor(totalSeconds / 60);
- const totalHours = Math.floor(totalMinutes / 60);
- const totalDays = Math.floor(totalHours / 24);
-
- if (totalSeconds < 60) {
- return `${totalSeconds} sec`;
- }
-
- if (totalMinutes < 60) {
- const seconds = totalSeconds % 60;
- return seconds > 0 ? `${totalMinutes} min ${seconds} sec` : `${totalMinutes} min`;
- }
-
- if (totalHours < 24) {
- const minutes = totalMinutes % 60;
- return minutes > 0 ? `${totalHours} h ${minutes} min` : `${totalHours} h`;
- }
-
- const hours = totalHours % 24;
- return hours > 0 ? `${totalDays} gg ${hours} h` : `${totalDays} gg`;
+ return formatDuration(startTime, endTime);
}
}
diff --git a/src/app/pages/admin/admin-access.util.spec.ts b/src/app/pages/admin/admin-access.util.spec.ts
new file mode 100644
index 0000000..5144c4a
--- /dev/null
+++ b/src/app/pages/admin/admin-access.util.spec.ts
@@ -0,0 +1,27 @@
+import { Router } from '@angular/router';
+import { vi } from 'vitest';
+import { redirectOnAdminAccessDenied } from './admin-access.util';
+
+describe('redirectOnAdminAccessDenied', () => {
+ it('resets busy state and redirects to /editor on "Admin access required."', () => {
+ const router = { navigateByUrl: vi.fn() } as unknown as Router;
+ const onRedirect = vi.fn();
+
+ const handled = redirectOnAdminAccessDenied(new Error('Admin access required.'), router, onRedirect);
+
+ expect(handled).toBe(true);
+ expect(onRedirect).toHaveBeenCalled();
+ expect(router.navigateByUrl).toHaveBeenCalledWith('/editor');
+ });
+
+ it('leaves other errors untouched', () => {
+ const router = { navigateByUrl: vi.fn() } as unknown as Router;
+ const onRedirect = vi.fn();
+
+ const handled = redirectOnAdminAccessDenied(new Error('Unable to create user.'), router, onRedirect);
+
+ expect(handled).toBe(false);
+ expect(onRedirect).not.toHaveBeenCalled();
+ expect(router.navigateByUrl).not.toHaveBeenCalled();
+ });
+});
diff --git a/src/app/pages/admin/admin-access.util.ts b/src/app/pages/admin/admin-access.util.ts
new file mode 100644
index 0000000..f5074bd
--- /dev/null
+++ b/src/app/pages/admin/admin-access.util.ts
@@ -0,0 +1,22 @@
+import { Router } from '@angular/router';
+
+/**
+ * Admin API calls fail with this exact message when the caller's session
+ * lost admin privileges mid-flow (e.g. role changed in another tab).
+ */
+const ADMIN_ACCESS_REQUIRED_MESSAGE = 'Admin access required.';
+
+/**
+ * Detects the "admin access required" failure from an admin API call and, if
+ * matched, resets whatever busy-state `onRedirect` clears and navigates away.
+ * Returns whether the error was handled, so callers can bail out of their own
+ * error handling with `if (redirectOnAdminAccessDenied(...)) return;`.
+ */
+export function redirectOnAdminAccessDenied(error: unknown, router: Router, onRedirect: () => void): boolean {
+ const message = error instanceof Error ? error.message : '';
+ if (message !== ADMIN_ACCESS_REQUIRED_MESSAGE) return false;
+
+ onRedirect();
+ router.navigateByUrl('/editor');
+ return true;
+}
diff --git a/src/app/pages/admin/admin-create-user/admin-create-user.ts b/src/app/pages/admin/admin-create-user/admin-create-user.ts
index facdc42..60d5b0a 100644
--- a/src/app/pages/admin/admin-create-user/admin-create-user.ts
+++ b/src/app/pages/admin/admin-create-user/admin-create-user.ts
@@ -11,7 +11,9 @@ import { MatSelectModule } from '@angular/material/select';
import { AdminCreateUserRequest, UserRole } from '@models/user';
import { Router } from '@angular/router';
import { AdminService } from '@services/admin/admin';
+import { redirectOnAdminAccessDenied } from '@pages/admin/admin-access.util';
import { FormUtility } from '@utilities/form-utility';
+import { scheduleSignalClear } from '@utilities/temporary-signal';
import { hasValidPasswordComplexity, evaluatePasswordChecks, initialPasswordChecks, PASSWORD_MIN_LENGTH } from '@utilities/password-validation';
@Component({
@@ -102,10 +104,10 @@ export class AdminCreateUserPage extends FormUtility {
password: '',
role: 'USER'
});
- setTimeout(() => this.successMessage.set(null), 3000);
+ scheduleSignalClear(this.successMessage);
},
error: (error) => {
- if (this.redirectOnAdminAccessDenied(error)) return;
+ if (redirectOnAdminAccessDenied(error, this.router, () => this.createSaving.set(false))) return;
this.createSaving.set(false);
const message = error instanceof Error ? error.message : 'Unable to create user.';
if (message === 'INVALID_EMAIL') {
@@ -120,13 +122,4 @@ export class AdminCreateUserPage extends FormUtility {
}
});
}
-
- private redirectOnAdminAccessDenied(error: unknown): boolean {
- const message = error instanceof Error ? error.message : '';
- if (message !== 'Admin access required.') return false;
-
- this.createSaving.set(false);
- this.router.navigateByUrl('/editor');
- return true;
- }
}
diff --git a/src/app/pages/admin/admin-users-list/admin-users-list.ts b/src/app/pages/admin/admin-users-list/admin-users-list.ts
index 0db20f1..088228b 100644
--- a/src/app/pages/admin/admin-users-list/admin-users-list.ts
+++ b/src/app/pages/admin/admin-users-list/admin-users-list.ts
@@ -11,6 +11,8 @@ import { Router } from '@angular/router';
import { AdminService } from '@services/admin/admin';
import { ConfirmDialogService } from '@services/dialogs/confirm-dialog';
import { AdminResetPasswordDialogComponent } from '@shared/admin-reset-password-dialog/admin-reset-password-dialog';
+import { redirectOnAdminAccessDenied } from '@pages/admin/admin-access.util';
+import { scheduleSignalClear } from '@utilities/temporary-signal';
@Component({
selector: 'app-admin-users-list-page',
@@ -64,7 +66,7 @@ export class AdminUsersListPage {
this.loading.set(false);
},
error: (error) => {
- if (this.redirectOnAdminAccessDenied(error)) return;
+ if (this.redirectOnAccessDenied(error)) return;
this.pageError.set(error instanceof Error ? error.message : 'Unable to load users.');
this.loading.set(false);
}
@@ -91,10 +93,10 @@ export class AdminUsersListPage {
this.roleSavingByUser.update((current) => ({ ...current, [user.username]: false }));
this.successMessage.set(`Role updated for ${user.username}.`);
this.loadUsers();
- setTimeout(() => this.successMessage.set(null), 3000);
+ scheduleSignalClear(this.successMessage);
},
error: (error) => {
- if (this.redirectOnAdminAccessDenied(error)) return;
+ if (this.redirectOnAccessDenied(error)) return;
this.roleSavingByUser.update((current) => ({ ...current, [user.username]: false }));
this.roleErrorByUser.update((current) => ({
...current,
@@ -123,10 +125,10 @@ export class AdminUsersListPage {
this.resetPasswordSaving.set(false);
this.resetPasswordDialogUser.set(null);
this.successMessage.set(`Password updated for ${event.username}.`);
- setTimeout(() => this.successMessage.set(null), 3000);
+ scheduleSignalClear(this.successMessage);
},
error: (error) => {
- if (this.redirectOnAdminAccessDenied(error)) return;
+ if (this.redirectOnAccessDenied(error)) return;
this.resetPasswordSaving.set(false);
this.resetPasswordError.set(error instanceof Error ? error.message : 'Unable to update password.');
}
@@ -143,10 +145,10 @@ export class AdminUsersListPage {
this.deleteBusyByUser.update((current) => ({ ...current, [user.username]: false }));
this.successMessage.set(`User ${user.username} deleted.`);
this.loadUsers();
- setTimeout(() => this.successMessage.set(null), 3000);
+ scheduleSignalClear(this.successMessage);
},
error: (error) => {
- if (this.redirectOnAdminAccessDenied(error)) return;
+ if (this.redirectOnAccessDenied(error)) return;
this.deleteBusyByUser.update((current) => ({ ...current, [user.username]: false }));
this.roleErrorByUser.update((current) => ({
...current,
@@ -156,13 +158,10 @@ export class AdminUsersListPage {
});
}
- private redirectOnAdminAccessDenied(error: unknown): boolean {
- const message = error instanceof Error ? error.message : '';
- if (message !== 'Admin access required.') return false;
-
- this.loading.set(false);
- this.resetPasswordSaving.set(false);
- this.router.navigateByUrl('/editor');
- return true;
+ private redirectOnAccessDenied(error: unknown): boolean {
+ return redirectOnAdminAccessDenied(error, this.router, () => {
+ this.loading.set(false);
+ this.resetPasswordSaving.set(false);
+ });
}
}
diff --git a/src/app/pages/admin/admin-users/admin-users.css b/src/app/pages/admin/admin-users/admin-users.css
deleted file mode 100644
index 01c0a0f..0000000
--- a/src/app/pages/admin/admin-users/admin-users.css
+++ /dev/null
@@ -1,187 +0,0 @@
-.admin-users-page {
- display: flex;
- flex-direction: column;
- gap: 16px;
- height: 100%;
- padding: 8px 4px 16px;
- overflow: auto;
- background: linear-gradient(180deg, #f8fafc 0%, #eef6ff 100%);
-}
-
-.admin-users-page__header {
- display: flex;
- align-items: flex-start;
- justify-content: space-between;
- gap: 16px;
-}
-
-.admin-users-page__eyebrow {
- color: #2563eb;
- font-size: 12px;
- font-weight: 700;
- letter-spacing: 0.08em;
- text-transform: uppercase;
-}
-
-.admin-users-page__title {
- margin: 4px 0 0;
- color: #0f172a;
- font-size: 28px;
- font-weight: 800;
-}
-
-.admin-users-page__subtitle {
- margin: 6px 0 0;
- color: #64748b;
- font-size: 14px;
-}
-
-.admin-users-page__success,
-.admin-users-card__error,
-.admin-users-row__error {
- border-radius: 14px;
- padding: 12px 14px;
- font-size: 13px;
- font-weight: 600;
-}
-
-.admin-users-page__success {
- border: 1px solid #86efac;
- background: #f0fdf4;
- color: #166534;
-}
-
-.admin-users-card {
- border-radius: 22px;
- padding: 20px;
-}
-
-.admin-users-card__title {
- margin-bottom: 14px;
- color: #0f172a;
- font-size: 18px;
- font-weight: 700;
-}
-
-.admin-users-card__error,
-.admin-users-row__error {
- border: 1px solid #fecaca;
- background: #fff1f2;
- color: #b91c1c;
-}
-
-.admin-users-create {
- display: flex;
- flex-direction: column;
- gap: 14px;
-}
-
-.admin-users-create__grid {
- display: grid;
- grid-template-columns: repeat(3, minmax(0, 1fr));
- gap: 12px;
-}
-
-.admin-users-create__password {
- width: 100%;
-}
-
-.admin-users-checklist {
- display: grid;
- grid-template-columns: repeat(2, minmax(0, 1fr));
- gap: 8px 12px;
-}
-
-.admin-users-check {
- display: flex;
- align-items: center;
- gap: 8px;
- color: #b91c1c;
- font-size: 12px;
-}
-
-.admin-users-check--ok {
- color: #15803d;
-}
-
-.admin-users-create__actions {
- display: flex;
- justify-content: flex-end;
-}
-
-.admin-users-state {
- color: #64748b;
- font-size: 14px;
- padding: 12px 4px;
-}
-
-.admin-users-list {
- display: flex;
- flex-direction: column;
- gap: 12px;
-}
-
-.admin-users-row {
- display: grid;
- grid-template-columns: minmax(180px, 1.2fr) minmax(220px, 1fr) auto;
- gap: 16px;
- align-items: center;
- border: 1px solid #e2e8f0;
- border-radius: 18px;
- background: #f8fbff;
- padding: 16px;
-}
-
-.admin-users-row__identity {
- min-width: 0;
-}
-
-.admin-users-row__username {
- color: #0f172a;
- font-size: 15px;
- font-weight: 700;
-}
-
-.admin-users-row__email {
- color: #64748b;
- font-size: 13px;
- margin-top: 4px;
- word-break: break-word;
-}
-
-.admin-users-row__role {
- display: flex;
- align-items: center;
- gap: 10px;
-}
-
-.admin-users-row__role-field {
- min-width: 150px;
-}
-
-.admin-users-row__actions {
- display: flex;
- align-items: center;
- justify-content: flex-end;
- gap: 8px;
- flex-wrap: wrap;
-}
-
-.admin-users-row__error {
- grid-column: 1 / -1;
-}
-
-@media (max-width: 1100px) {
- .admin-users-create__grid,
- .admin-users-row {
- grid-template-columns: 1fr;
- }
-
- .admin-users-checklist {
- grid-template-columns: 1fr;
- }
-
- .admin-users-row__actions {
- justify-content: flex-start;
- }
-}
diff --git a/src/app/pages/admin/admin-users/admin-users.html b/src/app/pages/admin/admin-users/admin-users.html
deleted file mode 100644
index 9322098..0000000
--- a/src/app/pages/admin/admin-users/admin-users.html
+++ /dev/null
@@ -1,140 +0,0 @@
-
-
-
- @if (successMessage()) {
-
{{ successMessage() }}
- }
-
-
- Create User
-
-
-
-
- Users
-
- @if (loading()) {
- Loading users...
- } @else if (pageError()) {
- {{ pageError() }}
- } @else if (!users().length) {
- No users found.
- } @else {
-
- @for (user of users(); track user.username) {
-
-
-
{{ user.username }}
-
{{ user.email || 'No email' }}
-
-
-
-
- Role
-
- USER
- ADMIN
-
-
-
-
-
-
-
-
-
-
- @if (roleErrorByUser()[user.username]) {
-
{{ roleErrorByUser()[user.username] }}
- }
-
- }
-
- }
-
-
- @if (resetPasswordDialogUser(); as user) {
-
-
- }
-
diff --git a/src/app/pages/admin/admin-users/admin-users.ts b/src/app/pages/admin/admin-users/admin-users.ts
deleted file mode 100644
index cd17f8e..0000000
--- a/src/app/pages/admin/admin-users/admin-users.ts
+++ /dev/null
@@ -1,258 +0,0 @@
-import { CommonModule } from '@angular/common';
-import { ChangeDetectionStrategy, Component, effect, inject, signal } from '@angular/core';
-import { FormsModule } from '@angular/forms';
-import { Field, form, minLength, required, validate } from '@angular/forms/signals';
-import { MatButtonModule } from '@angular/material/button';
-import { MatCardModule } from '@angular/material/card';
-import { MatFormFieldModule } from '@angular/material/form-field';
-import { MatIconModule } from '@angular/material/icon';
-import { MatInputModule } from '@angular/material/input';
-import { MatSelectModule } from '@angular/material/select';
-import {
- AdminCreateUserRequest,
- AdminUser,
- UserRole
-} from '@models/user';
-import { AdminService } from '@services/admin/admin';
-import { ConfirmDialogService } from '@services/dialogs/confirm-dialog';
-import { AdminResetPasswordDialogComponent } from '@shared/admin-reset-password-dialog/admin-reset-password-dialog';
-import { FormUtility } from '@utilities/form-utility';
-import { hasValidPasswordComplexity, evaluatePasswordChecks, initialPasswordChecks, PASSWORD_MIN_LENGTH } from '@utilities/password-validation';
-
-@Component({
- selector: 'app-admin-users',
- imports: [
- CommonModule,
- FormsModule,
- Field,
- MatButtonModule,
- MatCardModule,
- MatFormFieldModule,
- MatIconModule,
- MatInputModule,
- MatSelectModule,
- AdminResetPasswordDialogComponent
- ],
- templateUrl: './admin-users.html',
- styleUrl: './admin-users.css',
- changeDetection: ChangeDetectionStrategy.OnPush
-})
-export class AdminUsersPage extends FormUtility {
- private adminService = inject(AdminService);
- private confirmDialog = inject(ConfirmDialogService);
-
- readonly users = signal([]);
- readonly loading = signal(true);
- readonly pageError = signal(null);
- readonly createError = signal(null);
- readonly createEmailError = signal(null);
- readonly createPasswordError = signal(null);
- readonly createSaving = signal(false);
- readonly roleSavingByUser = signal>({});
- readonly roleErrorByUser = signal>({});
- readonly deleteBusyByUser = signal>({});
- readonly resetPasswordDialogUser = signal(null);
- readonly resetPasswordSaving = signal(false);
- readonly resetPasswordError = signal(null);
- readonly successMessage = signal(null);
- readonly roleDraftByUser = signal>({});
-
- readonly createModel = signal({
- username: '',
- email: '',
- password: '',
- role: 'USER'
- });
-
- readonly createForm = form(this.createModel, (model) => {
- required(model.username, { message: 'Username is required' });
- minLength(model.username, 3, { message: 'Username must be at least 3 characters' });
- required(model.email, { message: 'Email is required' });
- validate(model.email, ({ value }) => {
- const email = value();
- if (!email || /^[^\s@]+@[^\s@]+\.[^\s@]+$/.test(email)) return null;
- return {
- kind: 'invalidEmail',
- message: 'Invalid email address'
- };
- });
- required(model.password, { message: 'Password is required' });
- minLength(model.password, PASSWORD_MIN_LENGTH, { message: `Password must be at least ${PASSWORD_MIN_LENGTH} characters long` });
- validate(model.password, ({ value }) => {
- const password = value();
- if (!password || hasValidPasswordComplexity(password)) return null;
- return {
- kind: 'passwordComplexity',
- message: 'Password must include uppercase, lowercase, number and special character, with no spaces.'
- };
- });
- });
-
- readonly createPasswordChecks = signal(initialPasswordChecks());
-
- constructor() {
- super();
-
- effect(() => {
- const password = this.createModel().password;
- this.createPasswordChecks.set(evaluatePasswordChecks(password));
- });
- }
-
- ngOnInit() {
- this.loadUsers();
- }
-
- loadUsers() {
- this.loading.set(true);
- this.pageError.set(null);
- this.adminService.listAdminUsers().subscribe({
- next: (users) => {
- this.users.set(users);
- this.roleDraftByUser.set(
- users.reduce>((acc, user) => {
- acc[user.username] = user.role;
- return acc;
- }, {})
- );
- this.loading.set(false);
- },
- error: (error) => {
- this.pageError.set(error instanceof Error ? error.message : 'Unable to load users.');
- this.loading.set(false);
- }
- });
- }
-
- onCreateUser() {
- if (this.createForm().invalid()) return;
-
- this.createSaving.set(true);
- this.createError.set(null);
- this.createEmailError.set(null);
- this.createPasswordError.set(null);
-
- this.adminService.createAdminUser(this.createModel()).subscribe({
- next: () => {
- this.createSaving.set(false);
- this.successMessage.set('User created successfully.');
- this.createModel.set({
- username: '',
- email: '',
- password: '',
- role: 'USER'
- });
- this.loadUsers();
- setTimeout(() => this.successMessage.set(null), 3000);
- },
- error: (error) => {
- this.createSaving.set(false);
- const message = error instanceof Error ? error.message : 'Unable to create user.';
- if (message === 'INVALID_EMAIL') {
- this.createEmailError.set('Invalid email address');
- return;
- }
- if (message === 'INVALID_PASSWORD') {
- this.createPasswordError.set('Password does not satisfy the required policy.');
- return;
- }
- this.createError.set(message);
- }
- });
- }
-
- roleDraft(username: string): UserRole {
- return this.roleDraftByUser()[username] ?? 'USER';
- }
-
- setCreateRole(role: UserRole) {
- this.createModel.update((current) => ({
- ...current,
- role
- }));
- }
-
- setRoleDraft(username: string, role: UserRole) {
- this.roleDraftByUser.update((current) => ({
- ...current,
- [username]: role
- }));
- this.roleErrorByUser.update((current) => ({
- ...current,
- [username]: null
- }));
- }
-
- saveRole(user: AdminUser) {
- const nextRole = this.roleDraft(user.username);
- if (nextRole === user.role) return;
-
- this.roleSavingByUser.update((current) => ({ ...current, [user.username]: true }));
- this.roleErrorByUser.update((current) => ({ ...current, [user.username]: null }));
- this.adminService.changeAdminUserRole(user.username, { role: nextRole }).subscribe({
- next: () => {
- this.roleSavingByUser.update((current) => ({ ...current, [user.username]: false }));
- this.successMessage.set(`Role updated for ${user.username}.`);
- this.loadUsers();
- setTimeout(() => this.successMessage.set(null), 3000);
- },
- error: (error) => {
- this.roleSavingByUser.update((current) => ({ ...current, [user.username]: false }));
- this.roleErrorByUser.update((current) => ({
- ...current,
- [user.username]: error instanceof Error ? error.message : 'Unable to update role.'
- }));
- }
- });
- }
-
- openResetPasswordDialog(user: AdminUser) {
- this.resetPasswordError.set(null);
- this.resetPasswordDialogUser.set(user);
- }
-
- closeResetPasswordDialog() {
- if (this.resetPasswordSaving()) return;
- this.resetPasswordDialogUser.set(null);
- this.resetPasswordError.set(null);
- }
-
- submitResetPassword(event: { username: string; newPassword: string }) {
- this.resetPasswordSaving.set(true);
- this.resetPasswordError.set(null);
- this.adminService.changeAdminUserPassword(event.username, { newPassword: event.newPassword }).subscribe({
- next: () => {
- this.resetPasswordSaving.set(false);
- this.resetPasswordDialogUser.set(null);
- this.successMessage.set(`Password updated for ${event.username}.`);
- setTimeout(() => this.successMessage.set(null), 3000);
- },
- error: (error) => {
- this.resetPasswordSaving.set(false);
- this.resetPasswordError.set(error instanceof Error ? error.message : 'Unable to update password.');
- }
- });
- }
-
- async deleteUser(user: AdminUser) {
- const confirmed = await this.confirmDialog.open(`Delete user ${user.username}?`);
- if (!confirmed) return;
-
- this.deleteBusyByUser.update((current) => ({ ...current, [user.username]: true }));
- this.adminService.deleteAdminUser(user.username).subscribe({
- next: () => {
- this.deleteBusyByUser.update((current) => ({ ...current, [user.username]: false }));
- this.successMessage.set(`User ${user.username} deleted.`);
- this.loadUsers();
- setTimeout(() => this.successMessage.set(null), 3000);
- },
- error: (error) => {
- this.deleteBusyByUser.update((current) => ({ ...current, [user.username]: false }));
- this.roleErrorByUser.update((current) => ({
- ...current,
- [user.username]: error instanceof Error ? error.message : 'Unable to delete user.'
- }));
- }
- });
- }
-}
diff --git a/src/app/pages/auth/login/login.ts b/src/app/pages/auth/login/login.ts
index 1179b87..aa8f113 100644
--- a/src/app/pages/auth/login/login.ts
+++ b/src/app/pages/auth/login/login.ts
@@ -7,6 +7,7 @@ import { MatIconModule } from '@angular/material/icon';
import { MatInputModule } from '@angular/material/input';
import { Router, RouterLink } from "@angular/router";
import { FormUtility } from '@utilities/form-utility';
+import { scheduleSignalClear } from '@utilities/temporary-signal';
import { Authorization } from '@services/authorization/authorization';
import { Field, form, required } from '@angular/forms/signals';
@@ -40,12 +41,12 @@ export class Login extends FormUtility {
super();
effect(() => {
if (this.error() != null) {
- setTimeout(() => this.error.set(null), 3000);
+ scheduleSignalClear(this.error, 3000);
}
});
effect(() => {
if (this.registeredUser() != null) {
- setTimeout(() => this.registeredUser.set(null), 5000);
+ scheduleSignalClear(this.registeredUser, 5000);
}
});
}
diff --git a/src/app/pages/main/editor-sidebar/editor-sidebar.ts b/src/app/pages/main/editor-sidebar/editor-sidebar.ts
index c273354..441835f 100644
--- a/src/app/pages/main/editor-sidebar/editor-sidebar.ts
+++ b/src/app/pages/main/editor-sidebar/editor-sidebar.ts
@@ -72,12 +72,10 @@ export class EditorSidebar {
}
this.creatingFlow.set(true);
- console.log('Creating new flow...');
this.flowService.createNewFlow().pipe(
finalize(() => this.creatingFlow.set(false))
).subscribe({
next: flow => {
- console.log('New flow created:', flow);
this.flowState.openDocument(flow);
},
error: err => {
@@ -90,8 +88,4 @@ export class EditorSidebar {
this.createWithAiRequested.emit();
}
- createNewBlock() {
- console.log('Creating new block...');
- }
-
}
diff --git a/src/app/services/admin/admin-call.ts b/src/app/services/admin/admin-call.ts
index e96291a..0ab383e 100644
--- a/src/app/services/admin/admin-call.ts
+++ b/src/app/services/admin/admin-call.ts
@@ -7,11 +7,12 @@ import {
UserStatistics,
UserRole
} from "@models/user";
-import { HttpClient, HttpErrorResponse } from "@angular/common/http";
+import { HttpClient } from "@angular/common/http";
import { inject } from "@angular/core";
import { environment } from "@environment";
-import { catchError, map, Observable, throwError } from "rxjs";
+import { catchError, map, Observable } from "rxjs";
import { AdminCallServiceBase } from "./admin-call.base";
+import { toHttpError } from "@services/shared/http-error.util";
export class AdminCallService extends AdminCallServiceBase {
private readonly http = inject(HttpClient);
@@ -21,7 +22,7 @@ export class AdminCallService extends AdminCallServiceBase {
.get(`${environment.apiUrl}/auth/admin/users`)
.pipe(
map((raw) => Array.isArray(raw) ? raw.map((item) => this.adminUserFromApi(item)) : []),
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
403: 'Admin access required.'
}))
);
@@ -31,7 +32,7 @@ export class AdminCallService extends AdminCallServiceBase {
return this.http
.post(`${environment.apiUrl}/auth/admin/users`, request)
.pipe(
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
400: 'Unable to create user.',
403: 'Admin access required.'
}))
@@ -42,7 +43,7 @@ export class AdminCallService extends AdminCallServiceBase {
return this.http
.put(`${environment.apiUrl}/auth/admin/users/${encodeURIComponent(username)}/password`, request)
.pipe(
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
400: 'Unable to update password.',
403: 'Admin access required.',
404: `User ${username} not found`
@@ -54,7 +55,7 @@ export class AdminCallService extends AdminCallServiceBase {
return this.http
.put(`${environment.apiUrl}/auth/admin/users/${encodeURIComponent(username)}/role`, request)
.pipe(
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
400: 'Unable to update role.',
403: 'Admin access required.',
404: `User ${username} not found`,
@@ -67,7 +68,7 @@ export class AdminCallService extends AdminCallServiceBase {
return this.http
.delete(`${environment.apiUrl}/auth/admin/users/${encodeURIComponent(username)}`)
.pipe(
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
403: 'Admin access required.',
404: `User ${username} not found`,
409: 'LAST_ADMIN'
@@ -80,7 +81,7 @@ export class AdminCallService extends AdminCallServiceBase {
.get(`${environment.apiUrl}/stats`)
.pipe(
map((raw) => this.operationsStatisticsFromApi(raw)),
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
401: 'Unauthenticated',
403: 'You are not allowed to view user statistics'
}))
@@ -94,7 +95,7 @@ export class AdminCallService extends AdminCallServiceBase {
map((raw) => Array.isArray(raw)
? raw.filter((item): item is string => typeof item === 'string').map((item) => item.trim()).filter((item) => item.length > 0)
: []),
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
401: 'Unauthenticated',
403: 'You are not allowed to view user statistics'
}))
@@ -106,7 +107,7 @@ export class AdminCallService extends AdminCallServiceBase {
.get(`${environment.apiUrl}/stats/users/${encodeURIComponent(username)}`)
.pipe(
map((raw) => this.userStatisticsFromApi(raw, username)),
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
401: 'Unauthenticated',
403: 'You are not allowed to view user statistics',
404: 'User not found'
@@ -114,29 +115,6 @@ export class AdminCallService extends AdminCallServiceBase {
);
}
- private extractHttpErrorMessage(error: HttpErrorResponse): string | null {
- const payload = error.error;
- if (typeof payload === 'string' && payload.trim().length > 0) {
- return payload.trim();
- }
- if (payload && typeof payload === 'object') {
- const record = payload as Record;
- const directMessage = record['message'];
- if (typeof directMessage === 'string' && directMessage.trim().length > 0) {
- return directMessage.trim();
- }
- const errorMessage = record['error'];
- if (typeof errorMessage === 'string' && errorMessage.trim().length > 0) {
- return errorMessage.trim();
- }
- const details = record['details'];
- if (typeof details === 'string' && details.trim().length > 0) {
- return details.trim();
- }
- }
- return null;
- }
-
private normalizeRole(value: unknown): UserRole {
return String(value ?? '').toUpperCase() === 'ADMIN' ? 'ADMIN' : 'USER';
}
@@ -187,16 +165,4 @@ export class AdminCallService extends AdminCallServiceBase {
};
}
- private toHttpError(error: unknown, fallbackByStatus: Record): Observable {
- if (error instanceof HttpErrorResponse) {
- const message = this.extractHttpErrorMessage(error)
- ?? fallbackByStatus[error.status]
- ?? 'Request failed.';
- return throwError(() => new Error(message));
- }
- if (error instanceof Error) {
- return throwError(() => error);
- }
- return throwError(() => new Error('Request failed.'));
- }
}
diff --git a/src/app/services/assistant/assistant-call.fake.spec.ts b/src/app/services/assistant/assistant-call.fake.spec.ts
new file mode 100644
index 0000000..834f4f3
--- /dev/null
+++ b/src/app/services/assistant/assistant-call.fake.spec.ts
@@ -0,0 +1,44 @@
+import { catchError, of } from 'rxjs';
+import { AssistantCallServiceFake } from './assistant-call.fake';
+
+describe('AssistantCallServiceFake', () => {
+ let service: AssistantCallServiceFake;
+
+ beforeEach(() => {
+ service = new AssistantCallServiceFake();
+ });
+
+ it('surfaces "call not found" as an observable error catchError can intercept, not a synchronous throw', async () => {
+ expect(() => service.getCall('missing-call')).not.toThrow();
+
+ let caught: unknown = null;
+ await new Promise((resolve) => {
+ service.getCall('missing-call').pipe(
+ catchError((error) => {
+ caught = error;
+ return of(null);
+ })
+ ).subscribe(() => resolve());
+ });
+
+ expect(caught).toBeInstanceOf(Error);
+ expect((caught as Error).message).toContain('missing-call');
+ });
+
+ it('surfaces "session not found" as an observable error catchError can intercept, not a synchronous throw', async () => {
+ expect(() => service.getSession('missing-session')).not.toThrow();
+
+ let caught: unknown = null;
+ await new Promise((resolve) => {
+ service.getSession('missing-session').pipe(
+ catchError((error) => {
+ caught = error;
+ return of(null);
+ })
+ ).subscribe(() => resolve());
+ });
+
+ expect(caught).toBeInstanceOf(Error);
+ expect((caught as Error).message).toContain('missing-session');
+ });
+});
diff --git a/src/app/services/assistant/assistant-call.fake.ts b/src/app/services/assistant/assistant-call.fake.ts
index 1b5511c..1fcf7b4 100644
--- a/src/app/services/assistant/assistant-call.fake.ts
+++ b/src/app/services/assistant/assistant-call.fake.ts
@@ -9,7 +9,7 @@ import {
AssistantValidationIssue
} from '@models/assistant';
import { FlowData } from '@models/flow';
-import { Observable, of } from 'rxjs';
+import { defer, Observable, of } from 'rxjs';
import { AssistantCallServiceBase } from './assistant-call.base';
type FakeCallRecord = {
@@ -62,99 +62,109 @@ export class AssistantCallServiceFake extends AssistantCallServiceBase {
}
override sendMessage(sessionId: string, request: AssistantSendMessageRequest): Observable<{ callId: string }> {
- const session = this.sessions.get(sessionId);
- if (!session) {
- throw new Error(`Assistant session ${sessionId} not found`);
- }
- session.messages = [
- ...session.messages,
- {
- id: crypto.randomUUID(),
- role: 'user',
- content: request.message
+ return defer(() => {
+ const session = this.sessions.get(sessionId);
+ if (!session) {
+ throw new Error(`Assistant session ${sessionId} not found`);
}
- ];
+ session.messages = [
+ ...session.messages,
+ {
+ id: crypto.randomUUID(),
+ role: 'user',
+ content: request.message
+ }
+ ];
- const callId = crypto.randomUUID();
- const phases = this.buildPhases(request.message);
- this.calls.set(callId, {
- id: callId,
- sessionId,
- content: request.message,
- phaseIndex: 0,
- phases,
- completed: false,
- failed: false,
- cancelled: false
+ const callId = crypto.randomUUID();
+ const phases = this.buildPhases(request.message);
+ this.calls.set(callId, {
+ id: callId,
+ sessionId,
+ content: request.message,
+ phaseIndex: 0,
+ phases,
+ completed: false,
+ failed: false,
+ cancelled: false
+ });
+ session.lastCallId = callId;
+
+ return of({ callId });
});
- session.lastCallId = callId;
-
- return of({ callId });
}
override getCall(callId: string): Observable {
- const call = this.calls.get(callId);
- if (!call) {
- throw new Error(`Assistant call ${callId} not found`);
- }
-
- if (!call.completed && !call.failed && !call.cancelled) {
- if (call.phaseIndex < call.phases.length - 1) {
- call.phaseIndex += 1;
- } else {
- call.completed = true;
- this.applyCallResult(call);
+ return defer(() => {
+ const call = this.calls.get(callId);
+ if (!call) {
+ throw new Error(`Assistant call ${callId} not found`);
}
- }
- const phase = call.completed
- ? 'completed'
- : call.failed
- ? 'failed'
- : call.cancelled
- ? 'cancelled'
- : call.phases[call.phaseIndex];
+ if (!call.completed && !call.failed && !call.cancelled) {
+ if (call.phaseIndex < call.phases.length - 1) {
+ call.phaseIndex += 1;
+ } else {
+ call.completed = true;
+ this.applyCallResult(call);
+ }
+ }
- return of({
- id: call.id,
- sessionId: call.sessionId,
- status: call.failed
- ? 'FAILED'
- : call.cancelled
- ? 'CANCELLED'
- : call.completed
- ? 'COMPLETED'
- : call.phaseIndex === 0
- ? 'QUEUED'
- : 'RUNNING',
- phase,
- progressMessage: call.cancelled ? 'Assistant request cancelled' : undefined,
- errorMessage: call.failed ? 'Fake assistant call failed.' : undefined
+ const phase = call.completed
+ ? 'completed'
+ : call.failed
+ ? 'failed'
+ : call.cancelled
+ ? 'cancelled'
+ : call.phases[call.phaseIndex];
+
+ const result: AssistantCallState = {
+ id: call.id,
+ sessionId: call.sessionId,
+ status: call.failed
+ ? 'FAILED'
+ : call.cancelled
+ ? 'CANCELLED'
+ : call.completed
+ ? 'COMPLETED'
+ : call.phaseIndex === 0
+ ? 'QUEUED'
+ : 'RUNNING',
+ phase,
+ progressMessage: call.cancelled ? 'Assistant request cancelled' : undefined,
+ errorMessage: call.failed ? 'Fake assistant call failed.' : undefined
+ };
+ return of(result);
});
}
override cancelCall(callId: string): Observable {
- const call = this.calls.get(callId);
- if (!call) {
- throw new Error(`Assistant call ${callId} not found`);
- }
+ return defer(() => {
+ const call = this.calls.get(callId);
+ if (!call) {
+ throw new Error(`Assistant call ${callId} not found`);
+ }
- call.cancelled = true;
- return of({
- id: call.id,
- sessionId: call.sessionId,
- status: 'CANCELLED',
- phase: 'cancelled',
- progressMessage: 'Assistant request cancelled'
+ call.cancelled = true;
+ const result: AssistantCallState = {
+ id: call.id,
+ sessionId: call.sessionId,
+ status: 'CANCELLED',
+ phase: 'cancelled',
+ progressMessage: 'Assistant request cancelled'
+ };
+ return of(result);
});
}
override getSession(sessionId: string): Observable {
- const session = this.sessions.get(sessionId);
- if (!session) {
- throw new Error(`Assistant session ${sessionId} not found`);
- }
- return of(structuredClone(session));
+ return defer(() => {
+ const session = this.sessions.get(sessionId);
+ if (!session) {
+ throw new Error(`Assistant session ${sessionId} not found`);
+ }
+ return of(structuredClone(session));
+ });
}
private buildPhases(content: string): AssistantCallPhase[] {
diff --git a/src/app/services/authorization/authorization-call.ts b/src/app/services/authorization/authorization-call.ts
index e542cac..a220375 100644
--- a/src/app/services/authorization/authorization-call.ts
+++ b/src/app/services/authorization/authorization-call.ts
@@ -9,6 +9,7 @@ import { catchError, map, Observable, of, throwError } from "rxjs";
import { HttpClient, HttpErrorResponse } from "@angular/common/http";
import { inject } from "@angular/core";
import { environment } from "@environment";
+import { extractHttpErrorMessage, toHttpError } from "@services/shared/http-error.util";
export class AuthorizationCallService extends AuthorizationCallServiceBase {
private readonly http = inject(HttpClient);
@@ -35,7 +36,7 @@ export class AuthorizationCallService extends AuthorizationCallServiceBase {
.get(`${environment.apiUrl}/auth/me`)
.pipe(
map((raw) => this.userFromApi(raw)),
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
401: 'Unauthenticated'
}))
);
@@ -44,7 +45,7 @@ export class AuthorizationCallService extends AuthorizationCallServiceBase {
override register(userRegistration: UserRegistration): Observable {
return this.http.post(`${environment.apiUrl}/auth/register`, userRegistration)
.pipe(
- catchError((error: unknown) => this.toHttpError(error, {
+ catchError((error: unknown) => toHttpError(error, {
400: 'Unable to register user.'
}))
);
@@ -57,7 +58,7 @@ export class AuthorizationCallService extends AuthorizationCallServiceBase {
map(() => undefined),
catchError((error: unknown) => {
if (error instanceof HttpErrorResponse) {
- const message = this.extractHttpErrorMessage(error)
+ const message = extractHttpErrorMessage(error)
?? (error.status === 400 ? 'Missing required fields.' : null)
?? (error.status === 401 ? 'Current password is invalid' : null)
?? (error.status === 404 ? 'User not found' : null)
@@ -69,29 +70,6 @@ export class AuthorizationCallService extends AuthorizationCallServiceBase {
);
}
- private extractHttpErrorMessage(error: HttpErrorResponse): string | null {
- const payload = error.error;
- if (typeof payload === 'string' && payload.trim().length > 0) {
- return payload.trim();
- }
- if (payload && typeof payload === 'object') {
- const record = payload as Record;
- const directMessage = record['message'];
- if (typeof directMessage === 'string' && directMessage.trim().length > 0) {
- return directMessage.trim();
- }
- const errorMessage = record['error'];
- if (typeof errorMessage === 'string' && errorMessage.trim().length > 0) {
- return errorMessage.trim();
- }
- const details = record['details'];
- if (typeof details === 'string' && details.trim().length > 0) {
- return details.trim();
- }
- }
- return null;
- }
-
private userFromApi(raw: unknown, fallbackUsername?: string): User {
const payload = (raw ?? {}) as Record;
const userSource =
@@ -118,14 +96,4 @@ export class AuthorizationCallService extends AuthorizationCallServiceBase {
catchError(() => of(undefined))
);
}
-
- private toHttpError(error: unknown, fallbackByStatus: Record): Observable {
- if (error instanceof HttpErrorResponse) {
- const message = this.extractHttpErrorMessage(error)
- ?? fallbackByStatus[error.status]
- ?? 'Request failed.';
- return throwError(() => new Error(message));
- }
- return throwError(() => error);
- }
}
diff --git a/src/app/services/bias/bias-error.util.ts b/src/app/services/bias/bias-error.util.ts
new file mode 100644
index 0000000..c023d4f
--- /dev/null
+++ b/src/app/services/bias/bias-error.util.ts
@@ -0,0 +1,28 @@
+import { BiasSideEffectError } from '@models/bias-impact';
+
+/**
+ * Extracts a human-readable message for a bias-related request failure.
+ * Handles, in order: the typed side-effect conflict shape produced by
+ * `TaskExecutionsService.toBiasOperationError` (409s), the backend's
+ * `application/problem+json` body (`errors[].message`/`detail`, used by the
+ * bias endpoints for 400s), a plain `Error`, then `fallback`.
+ */
+export function extractBiasErrorMessage(error: unknown, fallback: string): string {
+ const sideEffectError = error as Partial;
+ if (sideEffectError.reason === 'SIDE_EFFECT_BLOCKED' || sideEffectError.reason === 'CONFIRMATION_REQUIRED') {
+ return sideEffectError.message ?? fallback;
+ }
+
+ const body = (error as { error?: unknown })?.error;
+ if (body && typeof body === 'object') {
+ const record = body as Record;
+ const errors = Array.isArray(record['errors']) ? record['errors'] : [];
+ const first = errors[0];
+ if (first && typeof first === 'object' && typeof (first as Record)['message'] === 'string') {
+ return (first as Record)['message'] as string;
+ }
+ if (typeof record['detail'] === 'string' && record['detail']) return record['detail'];
+ }
+
+ return error instanceof Error ? error.message : fallback;
+}
diff --git a/src/app/services/blocks/blocks-call.ts b/src/app/services/blocks/blocks-call.ts
index f467e1b..1676ea2 100644
--- a/src/app/services/blocks/blocks-call.ts
+++ b/src/app/services/blocks/blocks-call.ts
@@ -5,6 +5,7 @@ import { inject } from "@angular/core";
import { environment } from "@environment";
import { catchError, map, Observable, of, switchMap, take, throwError } from "rxjs";
import { BlockDraftContext, BlocksCallServiceBase } from "./block-call.base";
+import { attachSharedDefinitions, toApiPath, toNullableString, toPorts, toPosition, toRecord, toSchema, toValueKinds } from "@services/shared/flow-node-mapping";
export class BlocksCallService extends BlocksCallServiceBase {
private readonly http = inject(HttpClient);
@@ -88,7 +89,7 @@ export class BlocksCallService extends BlocksCallServiceBase {
const descriptor = types.find((type) => type.type === blockType);
const payload = this.buildBlockConfigurationPayload(
blockType,
- this.toRecord(configuration?.specificConfiguration ?? configuration),
+ toRecord(configuration?.specificConfiguration ?? configuration),
descriptor?.schema ?? null
);
@@ -106,8 +107,8 @@ export class BlocksCallService extends BlocksCallServiceBase {
map((raw) =>
this.flowBlockFromApi(
{
- ...(this.toRecord(raw)),
- id: this.toRecord(raw)["id"] ?? blockId
+ ...(toRecord(raw)),
+ id: toRecord(raw)["id"] ?? blockId
},
blockType,
payload
@@ -127,20 +128,20 @@ export class BlocksCallService extends BlocksCallServiceBase {
}
private parseCatalogResponse(raw: unknown): BlockType[] {
- const value = this.toRecord(raw);
+ const value = toRecord(raw);
const descriptors = value["descriptors"];
if (!Array.isArray(descriptors)) {
throw new Error('Invalid block catalog response: expected reduced catalog format with a descriptors array');
}
- const sharedDefinitions = this.toSchema(value["sharedDefinitions"]);
+ const sharedDefinitions = toSchema(value["sharedDefinitions"]);
return descriptors.map((descriptor) => this.blockTypeFromApi(descriptor, sharedDefinitions));
}
private blockTypeFromApi(raw: unknown, sharedDefinitions?: Record | null): BlockType {
- const value = this.toRecord(raw);
- const schema = this.attachSharedDefinitions(
- this.toSchema(value["schema"] ?? value["configurationSchema"] ?? null),
+ const value = toRecord(raw);
+ const schema = attachSharedDefinitions(
+ toSchema(value["schema"] ?? value["configurationSchema"] ?? null),
sharedDefinitions ?? null
);
@@ -151,27 +152,27 @@ export class BlocksCallService extends BlocksCallServiceBase {
userInteractive: Boolean(value["userInteractive"] ?? value["interactive"] ?? false),
interactionContract: this.toInteractionContract(value["interactionContract"]),
hasExampleBlock: Boolean(value["hasExampleBlock"] ?? false),
- exampleBlockEndpoint: this.toApiPath(value["exampleBlockEndpoint"]),
- configurationType: this.toNullableString(value["configurationType"]),
- configurationClass: this.toNullableString(value["configurationClass"]),
+ exampleBlockEndpoint: toApiPath(value["exampleBlockEndpoint"]),
+ configurationType: toNullableString(value["configurationType"]),
+ configurationClass: toNullableString(value["configurationClass"]),
schema
};
}
private flowBlockFromApi(raw: unknown, fallbackTypeName = "LLMBlock", fallbackConfig?: Record): FlowBlock {
- const root = this.toRecord(raw);
- const value = this.toRecord(root["block"] ?? root["node"] ?? root["data"] ?? root);
+ const root = toRecord(raw);
+ const value = toRecord(root["block"] ?? root["node"] ?? root["data"] ?? root);
const specificConfigurationRaw = value["specificConfiguration"] ?? value["configuration"] ?? value["blockConfiguration"] ?? fallbackConfig ?? {};
- const specificConfiguration = this.toRecord(specificConfigurationRaw);
+ const specificConfiguration = toRecord(specificConfigurationRaw);
const typeName = String(value["typeName"] ?? value["blockType"] ?? specificConfiguration["typeName"] ?? fallbackTypeName);
const io = this.defaultIOForBlockType(typeName);
return {
id: String(value["id"] ?? crypto.randomUUID()),
name: String(value["name"] ?? specificConfiguration["name"] ?? typeName),
- position: this.toPosition(value["position"]),
- inputs: this.toPorts(value["inputs"], io.inputs),
- outputs: this.toPorts(value["outputs"], io.outputs),
+ position: toPosition(value["position"]),
+ inputs: toPorts(value["inputs"], io.inputs),
+ outputs: toPorts(value["outputs"], io.outputs),
specificConfiguration,
typeName,
nodeFamily: 'block',
@@ -182,13 +183,13 @@ export class BlocksCallService extends BlocksCallServiceBase {
}
private biasAnnotationsDescriptorFromApi(raw: unknown): BiasAnnotationsDescriptor {
- const value = this.toRecord(raw);
- const rawOptions = this.toRecord(value["options"]);
+ const value = toRecord(raw);
+ const rawOptions = toRecord(value["options"]);
const options: Record = {};
for (const [field, entries] of Object.entries(rawOptions)) {
if (!Array.isArray(entries)) continue;
options[field] = entries
- .map((entry) => this.toRecord(entry))
+ .map((entry) => toRecord(entry))
.filter((entry) => typeof entry["value"] === "string")
.map((entry) => ({
value: String(entry["value"]),
@@ -203,9 +204,9 @@ export class BlocksCallService extends BlocksCallServiceBase {
blockProperty: String(value["blockProperty"] ?? "biasAnnotations"),
multiple: value["multiple"] !== false,
maxItems: Number.isFinite(maxItems) && maxItems >= 0 ? maxItems : null,
- schema: this.toRecord(value["schema"]),
+ schema: toRecord(value["schema"]),
options,
- defaults: this.toRecord(value["defaults"]),
+ defaults: toRecord(value["defaults"]),
serverGeneratedFields: Array.isArray(value["serverGeneratedFields"])
? value["serverGeneratedFields"].map(String)
: []
@@ -213,7 +214,7 @@ export class BlocksCallService extends BlocksCallServiceBase {
}
private biasCapabilitiesFromApi(raw: unknown, fallbackBlockType: string): BiasCapabilities {
- const value = this.toRecord(raw);
+ const value = toRecord(raw);
const activationModes = Array.isArray(value['activationModes'])
? value['activationModes']
.filter((mode): mode is string => typeof mode === 'string')
@@ -231,69 +232,6 @@ export class BlocksCallService extends BlocksCallServiceBase {
};
}
- private toPorts(raw: unknown, fallback: Array<{ name: string; type: string; multiple: boolean }>) {
- if (!Array.isArray(raw)) return fallback;
- return raw
- .map((port) => this.toRecord(port))
- .filter((port) => typeof port["name"] === "string" && (port["name"] as string).length > 0)
- .map((port) => {
- const type = String(port["type"] ?? "TEXT");
- const multiple = Boolean(port["multiple"] ?? false);
- return {
- ...port,
- name: String(port["name"]),
- type,
- multiple,
- valueKinds: this.toValueKinds(port["valueKinds"], { type, multiple })
- };
- });
- }
-
- private toValueKinds(raw: unknown, fallback: { type: string; multiple: boolean }) {
- if (!Array.isArray(raw)) {
- return [{ type: fallback.type, multiple: fallback.multiple }];
- }
-
- const kinds = raw
- .map((item) => this.toRecord(item))
- .filter((item) => typeof item["type"] === "string")
- .map((item) => ({
- type: String(item["type"] ?? fallback.type),
- multiple: Boolean(item["multiple"] ?? false)
- }));
-
- return kinds.length ? kinds : [{ type: fallback.type, multiple: fallback.multiple }];
- }
-
- private toPosition(raw: unknown): { x: number; y: number } | undefined {
- const value = this.toRecord(raw);
- const x = value["x"];
- const y = value["y"];
- if (typeof x !== "number" || typeof y !== "number") return undefined;
- return { x, y };
- }
-
- private toSchema(raw: unknown): Record | null {
- if (!raw || typeof raw !== "object" || Array.isArray(raw)) return null;
- return raw as Record;
- }
-
- private attachSharedDefinitions(
- schema: Record | null,
- sharedDefinitions: Record | null
- ): Record | null {
- if (!schema) return null;
- if (!sharedDefinitions || !Object.keys(sharedDefinitions).length) return schema;
-
- return {
- ...schema,
- sharedDefinitions: {
- ...sharedDefinitions,
- ...this.toRecord(schema["sharedDefinitions"])
- }
- };
- }
-
private toInteractionContract(raw: unknown): BlockType["interactionContract"] {
if (!raw || typeof raw !== "object" || Array.isArray(raw)) return null;
const value = raw as Record;
@@ -315,21 +253,6 @@ export class BlocksCallService extends BlocksCallServiceBase {
};
}
- private toRecord(value: unknown): Record {
- if (!value || typeof value !== "object" || Array.isArray(value)) return {};
- return value as Record;
- }
-
- private toNullableString(value: unknown): string | null {
- return typeof value === "string" && value.length > 0 ? value : null;
- }
-
- private toApiPath(value: unknown): string | null {
- if (typeof value !== "string" || value.length === 0) return null;
- if (/^https?:\/\//.test(value)) return value;
- return `${environment.apiUrl}${value.startsWith("/") ? value : `/${value}`}`;
- }
-
private buildBlockConfigurationPayload(
blockType: string,
configuration: Record,
@@ -394,8 +317,8 @@ export class BlocksCallService extends BlocksCallServiceBase {
private buildObjectFromSchema(node: unknown, root: unknown): Record {
const resolved = this.resolveRef(node, root);
- const resolvedRecord = this.toRecord(resolved);
- const properties = this.toRecord(resolvedRecord["properties"]);
+ const resolvedRecord = toRecord(resolved);
+ const properties = toRecord(resolvedRecord["properties"]);
const result: Record = {};
for (const [key, propSchema] of Object.entries(properties)) {
@@ -406,7 +329,7 @@ export class BlocksCallService extends BlocksCallServiceBase {
private buildValueFromSchema(node: unknown, root: unknown): unknown {
const resolved = this.resolveRef(node, root);
- const value = this.toRecord(resolved);
+ const value = toRecord(resolved);
if (Object.prototype.hasOwnProperty.call(value, "default")) {
return value["default"];
@@ -437,8 +360,8 @@ export class BlocksCallService extends BlocksCallServiceBase {
if (!schemaNode || !schemaRoot) return { ...configuration };
const resolved = this.resolveRef(schemaNode, schemaRoot);
- const schemaRecord = this.toRecord(resolved);
- const properties = this.toRecord(schemaRecord["properties"]);
+ const schemaRecord = toRecord(resolved);
+ const properties = toRecord(schemaRecord["properties"]);
if (!Object.keys(properties).length) {
return { ...configuration };
}
@@ -446,7 +369,7 @@ export class BlocksCallService extends BlocksCallServiceBase {
const sanitized: Record = {};
for (const [key, value] of Object.entries(configuration)) {
if (!Object.prototype.hasOwnProperty.call(properties, key)) continue;
- const propertySchema = this.toRecord(properties[key]);
+ const propertySchema = toRecord(properties[key]);
sanitized[key] = this.sanitizeSchemaValue(value, propertySchema, schemaRoot);
}
@@ -461,15 +384,15 @@ export class BlocksCallService extends BlocksCallServiceBase {
if (!schemaNode || !schemaRoot || value == null) return value;
const resolved = this.resolveRef(schemaNode, schemaRoot);
- const schemaRecord = this.toRecord(resolved);
+ const schemaRecord = toRecord(resolved);
const type = schemaRecord["type"];
if ((type === "object" || schemaRecord["properties"]) && value && typeof value === "object" && !Array.isArray(value)) {
- return this.sanitizeConfigurationBySchema(this.toRecord(value), schemaRecord, schemaRoot);
+ return this.sanitizeConfigurationBySchema(toRecord(value), schemaRecord, schemaRoot);
}
if (type === "array" && Array.isArray(value)) {
- const itemSchema = this.toRecord(schemaRecord["items"]);
+ const itemSchema = toRecord(schemaRecord["items"]);
return value.map((item) => this.sanitizeSchemaValue(item, itemSchema, schemaRoot));
}
@@ -477,7 +400,7 @@ export class BlocksCallService extends BlocksCallServiceBase {
}
private resolveRef(node: unknown, root: unknown): unknown {
- const value = this.toRecord(node);
+ const value = toRecord(node);
const ref = value["$ref"];
if (typeof ref !== "string" || !ref.startsWith("#/")) return node;
diff --git a/src/app/services/blocks/blocks.spec.ts b/src/app/services/blocks/blocks.spec.ts
index b6c9524..5094ebd 100644
--- a/src/app/services/blocks/blocks.spec.ts
+++ b/src/app/services/blocks/blocks.spec.ts
@@ -1,4 +1,6 @@
import { TestBed } from '@angular/core/testing';
+import { throwError, of } from 'rxjs';
+import { vi } from 'vitest';
import { BlocksService } from './blocks';
@@ -13,4 +15,17 @@ describe('BlocksService', () => {
it('should be created', () => {
expect(service).toBeTruthy();
});
+
+ it('retries loading the block catalog after a failed initial fetch', async () => {
+ const retrieveAllBlocksTypes = vi.fn()
+ .mockReturnValueOnce(throwError(() => new Error('network down')))
+ .mockReturnValueOnce(of([]));
+ service.blocksCallService = { retrieveAllBlocksTypes } as unknown as typeof service.blocksCallService;
+
+ await expect(service.getAllBlocksTypes()).rejects.toThrow('network down');
+ expect(retrieveAllBlocksTypes).toHaveBeenCalledTimes(1);
+
+ await service.getAllBlocksTypes();
+ expect(retrieveAllBlocksTypes).toHaveBeenCalledTimes(2);
+ });
});
diff --git a/src/app/services/blocks/blocks.ts b/src/app/services/blocks/blocks.ts
index 90f0c59..de9422d 100644
--- a/src/app/services/blocks/blocks.ts
+++ b/src/app/services/blocks/blocks.ts
@@ -1,30 +1,31 @@
-import { computed, Injectable, signal } from '@angular/core';
+import { Injectable, Signal, signal } from '@angular/core';
import { environment } from '@environment';
import { BiasAnnotationsDescriptor, BlockType, BlockTypeName, FlowBlock } from '@models/flow';
import { BiasCapabilities } from '@models/bias-impact';
import { BlockDraftContext, BlocksCallServiceBase } from './block-call.base';
-import { catchError, finalize, firstValueFrom, map, Observable, of, shareReplay, tap, throwError } from 'rxjs';
+import { CatalogStore } from '@services/shared/catalog-store';
+import { EmptyNodeCache } from '@services/shared/empty-node-cache';
+import { PendingSyncCounter } from '@services/shared/pending-sync-counter';
+import { catchError, firstValueFrom, Observable, of, tap, throwError } from 'rxjs';
@Injectable({
providedIn: 'root',
})
-export class BlocksService {
+export class BlocksService extends CatalogStore {
blocksCallService: BlocksCallServiceBase = new environment.blocksCallService();
- toInit: boolean = true;
- private loadingPromise: Promise | null = null;
- private readonly _catalogLoading = signal(false);
- private readonly emptyBlockCache = new Map();
- private readonly pendingEmptyBlockRequests = new Map>();
- private readonly pendingServerSyncCount = signal(0);
+ protected readonly loadErrorLabel = 'Retrieve blocks types failed';
+
+ private readonly emptyBlockCache = new EmptyNodeCache();
+ private readonly serverSync = new PendingSyncCounter();
- private _blockTypes = signal([]);
private readonly _biasAnnotationsDescriptor = signal(null);
private readonly _biasCapabilities = signal>({});
private biasDescriptorPromise: Promise | null = null;
- readonly hasPendingServerSync = computed(() => this.pendingServerSyncCount() > 0);
- readonly blockTypes = this._blockTypes.asReadonly();
- readonly catalogLoading = this._catalogLoading.asReadonly();
+
+ readonly hasPendingServerSync = this.serverSync.active;
+ readonly blockTypes = this.types;
+ readonly catalogLoading = this.loading;
readonly biasAnnotationsDescriptor = this._biasAnnotationsDescriptor.asReadonly();
readonly biasCapabilities = this._biasCapabilities.asReadonly();
@@ -58,97 +59,26 @@ export class BlocksService {
}
hasLoadedBlockTypes() {
- return this._blockTypes().length > 0 || (!this.toInit && !this.loadingPromise);
+ return this.hasLoadedTypes();
}
- async getAllBlocksTypes() {
- if (this.toInit) {
- this.toInit = false;
- await this.refresh();
- } else if (this.loadingPromise) {
- await this.loadingPromise;
- }
-
- return this._blockTypes.asReadonly();
+ getAllBlocksTypes(): Promise> {
+ return this.getAllTypes();
}
- async refresh(force = false): Promise {
- if (this.loadingPromise && !force) {
- return this.loadingPromise;
- }
-
- this.loadingPromise = firstValueFrom(this.blocksCallService.retrieveAllBlocksTypes())
- .finally(() => {
- this._catalogLoading.set(false);
- })
- .then((blockTypes) => {
- this._blockTypes.set(blockTypes);
- this.clearEmptyBlockCache();
- })
- .catch((err) => {
- console.error('Retrieve blocks types failed', err);
- throw err;
- })
- .finally(() => {
- this.loadingPromise = null;
- });
-
- this._catalogLoading.set(true);
-
- return this.loadingPromise;
+ async getBlockType(typeName: BlockTypeName): Promise {
+ return this.getTypeOrFetch((blockType) => blockType.type === typeName);
}
- async getBlockType(typeName: BlockTypeName) {
- const current = this._blockTypes().find((blockType) => blockType.type === typeName);
- if (current) return current;
-
- if (this.loadingPromise) {
- await this.loadingPromise;
- return this._blockTypes().find((blockType) => blockType.type === typeName);
- }
-
- this._catalogLoading.set(true);
- const blockTypes = await firstValueFrom(this.blocksCallService.retrieveAllBlocksTypes())
- .finally(() => {
- this._catalogLoading.set(false);
- });
- this._blockTypes.set(blockTypes);
- this.clearEmptyBlockCache();
- return blockTypes.find((blockType) => blockType.type === typeName);
- }
-
- peekBlockType(typeName: BlockTypeName) {
- return this._blockTypes().find((blockType) => blockType.type === typeName) ?? null;
+ peekBlockType(typeName: BlockTypeName): BlockType | null {
+ return this.peekType((blockType) => blockType.type === typeName);
}
createEmptyBlock(blockType: BlockTypeName, context?: BlockDraftContext) {
const flowId = typeof context?.flowId === 'string' && context.flowId.trim().length > 0 ? context.flowId.trim() : '';
const cacheKey = `${String(blockType)}::${flowId}`;
- const cached = this.emptyBlockCache.get(cacheKey);
- if (cached) {
- return of(this.cloneEmptyBlock(cached));
- }
- const pending = this.pendingEmptyBlockRequests.get(cacheKey);
- if (pending) {
- return pending.pipe(map((block) => this.cloneEmptyBlock(block)));
- }
-
- const request = this.blocksCallService.createEmptyBlock(blockType, context).pipe(
- map((block) => {
- this.emptyBlockCache.set(cacheKey, this.cloneEmptyBlock(block));
- return block;
- }),
- finalize(() => {
- this.pendingEmptyBlockRequests.delete(cacheKey);
- }),
- shareReplay(1)
- );
-
- this.pendingEmptyBlockRequests.set(cacheKey, request);
-
- return request.pipe(
- map((block) => this.cloneEmptyBlock(block)),
+ return this.emptyBlockCache.getOrCreate(cacheKey, () => this.blocksCallService.createEmptyBlock(blockType, context)).pipe(
catchError((err) => {
console.error('Create empty block failed', err);
return throwError(() => err);
@@ -157,11 +87,7 @@ export class BlocksService {
}
updateBlock(blockId: string, configuration: any, context?: BlockDraftContext) {
- this.pendingServerSyncCount.update((count) => count + 1);
- return this.blocksCallService.updateBlock(blockId, configuration, context).pipe(
- finalize(() => {
- this.pendingServerSyncCount.update((count) => Math.max(0, count - 1));
- }),
+ return this.serverSync.track(this.blocksCallService.updateBlock(blockId, configuration, context)).pipe(
catchError((err) => {
console.error('Update block failed', err);
return throwError(() => err);
@@ -169,28 +95,11 @@ export class BlocksService {
);
}
- private clearEmptyBlockCache() {
+ protected fetchAll(): Observable {
+ return this.blocksCallService.retrieveAllBlocksTypes();
+ }
+
+ protected override onLoaded(): void {
this.emptyBlockCache.clear();
- this.pendingEmptyBlockRequests.clear();
- }
-
- private cloneEmptyBlock(block: FlowBlock): FlowBlock {
- const clone = this.deepClone(block);
- return {
- ...clone,
- id: globalThis.crypto?.randomUUID?.() ?? `${Date.now()}`,
- position: undefined
- };
- }
-
- private deepClone(value: T): T {
- if (typeof globalThis.structuredClone === 'function') {
- try {
- return globalThis.structuredClone(value);
- } catch {
- // Some cached payloads may carry non-cloneable runtime fields.
- }
- }
- return JSON.parse(JSON.stringify(value)) as T;
}
}
diff --git a/src/app/services/containers/containers-call.ts b/src/app/services/containers/containers-call.ts
index 2027558..c2858b3 100644
--- a/src/app/services/containers/containers-call.ts
+++ b/src/app/services/containers/containers-call.ts
@@ -9,6 +9,7 @@ import {
} from "@models/flow";
import { map, Observable } from "rxjs";
import { ContainersCallServiceBase } from "./container-call.base";
+import { attachSharedDefinitions, toApiPath, toNullableString, toPorts, toPosition, toRecord, toSchema, toValueKinds } from "@services/shared/flow-node-mapping";
export class ContainersCallService extends ContainersCallServiceBase {
private readonly http = inject(HttpClient);
@@ -39,7 +40,7 @@ export class ContainersCallService extends ContainersCallServiceBase {
const containerType = String(configuration?.typeName ?? configuration?.type ?? "GenericContainer");
const payload = this.buildContainerConfigurationPayload(
containerType,
- this.toRecord(configuration?.specificConfiguration ?? configuration)
+ toRecord(configuration?.specificConfiguration ?? configuration)
);
return this.http
@@ -54,8 +55,8 @@ export class ContainersCallService extends ContainersCallServiceBase {
map((raw) =>
this.flowContainerFromApi(
{
- ...(this.toRecord(raw)),
- id: this.toRecord(raw)["id"] ?? containerId
+ ...(toRecord(raw)),
+ id: toRecord(raw)["id"] ?? containerId
},
containerType
)
@@ -64,7 +65,7 @@ export class ContainersCallService extends ContainersCallServiceBase {
}
override validateContainerSubflow(subFlow: FlowData, validationUrl?: string | null): Observable {
- const resolvedValidationUrl = this.toApiPath(validationUrl) ?? `${environment.apiUrl}/containers/validate-subflow`;
+ const resolvedValidationUrl = toApiPath(validationUrl) ?? `${environment.apiUrl}/containers/validate-subflow`;
return this.http
.post(
resolvedValidationUrl,
@@ -74,20 +75,20 @@ export class ContainersCallService extends ContainersCallServiceBase {
}
private parseCatalogResponse(raw: unknown): BlockType[] {
- const value = this.toRecord(raw);
+ const value = toRecord(raw);
const descriptors = value["descriptors"];
if (!Array.isArray(descriptors)) {
throw new Error('Invalid container catalog response: expected reduced catalog format with a descriptors array');
}
- const sharedDefinitions = this.toSchema(value["sharedDefinitions"]);
+ const sharedDefinitions = toSchema(value["sharedDefinitions"]);
return descriptors.map((descriptor) => this.containerTypeFromApi(descriptor, sharedDefinitions));
}
private containerTypeFromApi(raw: unknown, sharedDefinitions?: Record | null): BlockType {
- const value = this.toRecord(raw);
- const schema = this.attachSharedDefinitions(
- this.toSchema(value["schema"] ?? value["configurationSchema"] ?? null),
+ const value = toRecord(raw);
+ const schema = attachSharedDefinitions(
+ toSchema(value["schema"] ?? value["configurationSchema"] ?? null),
sharedDefinitions ?? null
);
@@ -97,26 +98,26 @@ export class ContainersCallService extends ContainersCallServiceBase {
description: String(value["description"] ?? ""),
userInteractive: Boolean(value["userInteractive"] ?? value["interactive"] ?? false),
hasExampleBlock: Boolean(value["hasExampleBlock"] ?? value["hasExampleContainer"] ?? false),
- exampleBlockEndpoint: this.toApiPath(value["exampleBlockEndpoint"] ?? value["exampleContainerEndpoint"]),
- configurationType: this.toNullableString(value["configurationType"]),
- configurationClass: this.toNullableString(value["configurationClass"]),
+ exampleBlockEndpoint: toApiPath(value["exampleBlockEndpoint"] ?? value["exampleContainerEndpoint"]),
+ configurationType: toNullableString(value["configurationType"]),
+ configurationClass: toNullableString(value["configurationClass"]),
schema
};
}
private flowContainerFromApi(raw: unknown, fallbackTypeName = "GenericContainer"): FlowContainer {
- const root = this.toRecord(raw);
- const value = this.toRecord(root["container"] ?? root["node"] ?? root["data"] ?? root);
+ const root = toRecord(raw);
+ const value = toRecord(root["container"] ?? root["node"] ?? root["data"] ?? root);
const specificConfigurationRaw = value["specificConfiguration"] ?? value["configuration"] ?? value["containerConfiguration"] ?? {};
- const specificConfiguration = this.toRecord(specificConfigurationRaw);
+ const specificConfiguration = toRecord(specificConfigurationRaw);
const typeName = String(value["typeName"] ?? value["containerType"] ?? specificConfiguration["typeName"] ?? fallbackTypeName);
return {
id: String(value["id"] ?? crypto.randomUUID()),
name: String(value["name"] ?? specificConfiguration["name"] ?? typeName),
- position: this.toPosition(value["position"]),
- inputs: this.toPorts(value["inputs"]),
- outputs: this.toPorts(value["outputs"]),
+ position: toPosition(value["position"]),
+ inputs: toPorts(value["inputs"]),
+ outputs: toPorts(value["outputs"]),
specificConfiguration,
typeName,
nodeFamily: 'container'
@@ -124,17 +125,17 @@ export class ContainersCallService extends ContainersCallServiceBase {
}
private subflowValidationFromApi(raw: unknown): FlowSubflowValidationResult {
- const value = this.toRecord(raw);
+ const value = toRecord(raw);
const rawErrors = Array.isArray(value['errors']) ? value['errors'] : [];
return {
valid: Boolean(value['valid'] ?? false),
errors: rawErrors
- .map((item) => this.toRecord(item))
+ .map((item) => toRecord(item))
.map((item) => ({
- entity: this.toNullableString(item['entity']) ?? undefined,
- id: this.toNullableString(item['id']) ?? undefined,
- field: this.toNullableString(item['field']) ?? undefined,
+ entity: toNullableString(item['entity']) ?? undefined,
+ id: toNullableString(item['id']) ?? undefined,
+ field: toNullableString(item['field']) ?? undefined,
message: String(item['message'] ?? 'Invalid subflow')
})),
openInputs: this.toOpenInputs(value['openInputs']),
@@ -145,17 +146,17 @@ export class ContainersCallService extends ContainersCallServiceBase {
private toOpenInputs(raw: unknown) {
if (!Array.isArray(raw)) return [];
return raw
- .map((item) => this.toRecord(item))
+ .map((item) => toRecord(item))
.map((item) => {
- const io = this.toRecord(item['io']);
- const port = this.toPorts([Object.keys(io).length ? io : item])[0];
+ const io = toRecord(item['io']);
+ const port = toPorts([Object.keys(io).length ? io : item])[0];
if (!port) return null;
return {
...port,
- targetBlockId: this.toNullableString(item['targetBlockId'] ?? item['blockId'] ?? item['nodeId']) ?? undefined,
- targetInputName: this.toNullableString(item['targetInputName'] ?? item['inputName'] ?? io['name']) ?? undefined,
- blockId: this.toNullableString(item['blockId'] ?? item['nodeId']) ?? undefined,
- inputName: this.toNullableString(item['inputName'] ?? io['name']) ?? undefined
+ targetBlockId: toNullableString(item['targetBlockId'] ?? item['blockId'] ?? item['nodeId']) ?? undefined,
+ targetInputName: toNullableString(item['targetInputName'] ?? item['inputName'] ?? io['name']) ?? undefined,
+ blockId: toNullableString(item['blockId'] ?? item['nodeId']) ?? undefined,
+ inputName: toNullableString(item['inputName'] ?? io['name']) ?? undefined
};
})
.filter((item): item is NonNullable => !!item);
@@ -164,100 +165,22 @@ export class ContainersCallService extends ContainersCallServiceBase {
private toOpenOutputs(raw: unknown) {
if (!Array.isArray(raw)) return [];
return raw
- .map((item) => this.toRecord(item))
+ .map((item) => toRecord(item))
.map((item) => {
- const io = this.toRecord(item['io']);
- const port = this.toPorts([Object.keys(io).length ? io : item])[0];
+ const io = toRecord(item['io']);
+ const port = toPorts([Object.keys(io).length ? io : item])[0];
if (!port) return null;
return {
...port,
- sourceBlockId: this.toNullableString(item['sourceBlockId'] ?? item['blockId'] ?? item['nodeId']) ?? undefined,
- sourceOutputName: this.toNullableString(item['sourceOutputName'] ?? item['outputName'] ?? io['name']) ?? undefined,
- blockId: this.toNullableString(item['blockId'] ?? item['nodeId']) ?? undefined,
- outputName: this.toNullableString(item['outputName'] ?? io['name']) ?? undefined
+ sourceBlockId: toNullableString(item['sourceBlockId'] ?? item['blockId'] ?? item['nodeId']) ?? undefined,
+ sourceOutputName: toNullableString(item['sourceOutputName'] ?? item['outputName'] ?? io['name']) ?? undefined,
+ blockId: toNullableString(item['blockId'] ?? item['nodeId']) ?? undefined,
+ outputName: toNullableString(item['outputName'] ?? io['name']) ?? undefined
};
})
.filter((item): item is NonNullable => !!item);
}
- private toPorts(raw: unknown) {
- if (!Array.isArray(raw)) return [];
- return raw
- .map((port) => this.toRecord(port))
- .filter((port) => typeof port["name"] === "string" && (port["name"] as string).length > 0)
- .map((port) => {
- const type = String(port["type"] ?? "TEXT");
- const multiple = Boolean(port["multiple"] ?? false);
- return {
- ...port,
- name: String(port["name"]),
- type,
- multiple,
- valueKinds: this.toValueKinds(port["valueKinds"], { type, multiple })
- };
- });
- }
-
- private toValueKinds(raw: unknown, fallback: { type: string; multiple: boolean }) {
- if (!Array.isArray(raw)) {
- return [{ type: fallback.type, multiple: fallback.multiple }];
- }
-
- const kinds = raw
- .map((item) => this.toRecord(item))
- .filter((item) => typeof item["type"] === "string")
- .map((item) => ({
- type: String(item["type"] ?? fallback.type),
- multiple: Boolean(item["multiple"] ?? false)
- }));
-
- return kinds.length ? kinds : [{ type: fallback.type, multiple: fallback.multiple }];
- }
-
- private toPosition(raw: unknown): { x: number; y: number } | undefined {
- const value = this.toRecord(raw);
- const x = value["x"];
- const y = value["y"];
- if (typeof x !== "number" || typeof y !== "number") return undefined;
- return { x, y };
- }
-
- private toSchema(raw: unknown): Record | null {
- if (!raw || typeof raw !== "object" || Array.isArray(raw)) return null;
- return raw as Record;
- }
-
- private attachSharedDefinitions(
- schema: Record | null,
- sharedDefinitions: Record | null
- ): Record | null {
- if (!schema) return null;
- if (!sharedDefinitions || !Object.keys(sharedDefinitions).length) return schema;
-
- return {
- ...schema,
- sharedDefinitions: {
- ...sharedDefinitions,
- ...this.toRecord(schema["sharedDefinitions"])
- }
- };
- }
-
- private toRecord(value: unknown): Record {
- if (!value || typeof value !== "object" || Array.isArray(value)) return {};
- return value as Record;
- }
-
- private toNullableString(value: unknown): string | null {
- return typeof value === "string" && value.length > 0 ? value : null;
- }
-
- private toApiPath(value: unknown): string | null {
- if (typeof value !== "string" || value.length === 0) return null;
- if (/^https?:\/\//.test(value)) return value;
- return `${environment.apiUrl}${value.startsWith("/") ? value : `/${value}`}`;
- }
-
private resolveExampleEndpoint(typeName: string, descriptor?: BlockType): string {
if (descriptor?.hasExampleBlock && descriptor.exampleBlockEndpoint) {
return descriptor.exampleBlockEndpoint;
@@ -292,10 +215,10 @@ export class ContainersCallService extends ContainersCallServiceBase {
: [];
if (!required.length) return normalized;
- const properties = this.toRecord(schema["properties"]);
+ const properties = toRecord(schema["properties"]);
for (const key of required) {
if (normalized[key] !== undefined) continue;
- const propertySchema = this.toRecord(properties[key]);
+ const propertySchema = toRecord(properties[key]);
const defaultValue = propertySchema["default"];
if (typeof defaultValue === "boolean" || typeof defaultValue === "number" || typeof defaultValue === "string") {
normalized[key] = defaultValue;
@@ -310,7 +233,7 @@ export class ContainersCallService extends ContainersCallServiceBase {
}
private resolveConfigurationType(containerType: string, configuration: Record) {
- const explicitType = this.toNullableString(configuration["type"]);
+ const explicitType = toNullableString(configuration["type"]);
if (explicitType) return explicitType;
const descriptor = this.containerTypesCache
@@ -331,7 +254,7 @@ export class ContainersCallService extends ContainersCallServiceBase {
if (!typeProperty || typeof typeProperty !== "object" || Array.isArray(typeProperty)) return null;
const typeSchema = typeProperty as Record;
- const defaultValue = this.toNullableString(typeSchema["default"]);
+ const defaultValue = toNullableString(typeSchema["default"]);
if (defaultValue) return defaultValue;
const enumValues = Array.isArray(typeSchema["enum"])
diff --git a/src/app/services/containers/containers.spec.ts b/src/app/services/containers/containers.spec.ts
new file mode 100644
index 0000000..301388f
--- /dev/null
+++ b/src/app/services/containers/containers.spec.ts
@@ -0,0 +1,31 @@
+import { TestBed } from '@angular/core/testing';
+import { throwError, of } from 'rxjs';
+import { vi } from 'vitest';
+
+import { ContainersService } from './containers';
+
+describe('ContainersService', () => {
+ let service: ContainersService;
+
+ beforeEach(() => {
+ TestBed.configureTestingModule({});
+ service = TestBed.inject(ContainersService);
+ });
+
+ it('should be created', () => {
+ expect(service).toBeTruthy();
+ });
+
+ it('retries loading the container catalog after a failed initial fetch', async () => {
+ const retrieveAllContainerTypes = vi.fn()
+ .mockReturnValueOnce(throwError(() => new Error('network down')))
+ .mockReturnValueOnce(of([]));
+ service.containersCallService = { retrieveAllContainerTypes } as unknown as typeof service.containersCallService;
+
+ await expect(service.getAllContainerTypes()).rejects.toThrow('network down');
+ expect(retrieveAllContainerTypes).toHaveBeenCalledTimes(1);
+
+ await service.getAllContainerTypes();
+ expect(retrieveAllContainerTypes).toHaveBeenCalledTimes(2);
+ });
+});
diff --git a/src/app/services/containers/containers.ts b/src/app/services/containers/containers.ts
index 79d844f..2ff52c2 100644
--- a/src/app/services/containers/containers.ts
+++ b/src/app/services/containers/containers.ts
@@ -1,118 +1,48 @@
-import { computed, Injectable, signal } from '@angular/core';
+import { Injectable, Signal } from '@angular/core';
import { environment } from '@environment';
import { BlockType, BlockTypeName, FlowData, FlowNode } from '@models/flow';
-import { catchError, finalize, firstValueFrom, map, Observable, of, shareReplay, throwError } from 'rxjs';
+import { catchError, Observable, throwError } from 'rxjs';
import { ContainersCallServiceBase } from './container-call.base';
+import { CatalogStore } from '@services/shared/catalog-store';
+import { EmptyNodeCache } from '@services/shared/empty-node-cache';
+import { PendingSyncCounter } from '@services/shared/pending-sync-counter';
+import { deepClone } from '@services/shared/deep-clone';
@Injectable({
providedIn: 'root',
})
-export class ContainersService {
+export class ContainersService extends CatalogStore {
containersCallService: ContainersCallServiceBase = new environment.containersCallService();
- toInit = true;
- private loadingPromise: Promise | null = null;
- private readonly _catalogLoading = signal(false);
- private readonly emptyContainerCache = new Map();
- private readonly pendingEmptyContainerRequests = new Map>();
- private readonly pendingServerSyncCount = signal(0);
+ protected readonly loadErrorLabel = 'Retrieve container types failed';
- private _containerTypes = signal([]);
- readonly hasPendingServerSync = computed(() => this.pendingServerSyncCount() > 0);
- readonly containerTypes = this._containerTypes.asReadonly();
- readonly catalogLoading = this._catalogLoading.asReadonly();
+ private readonly emptyContainerCache = new EmptyNodeCache();
+ private readonly serverSync = new PendingSyncCounter();
+
+ readonly hasPendingServerSync = this.serverSync.active;
+ readonly containerTypes = this.types;
+ readonly catalogLoading = this.loading;
hasLoadedContainerTypes() {
- return this._containerTypes().length > 0 || (!this.toInit && !this.loadingPromise);
+ return this.hasLoadedTypes();
}
- async getAllContainerTypes() {
- if (this.toInit) {
- this.toInit = false;
- await this.refresh();
- } else if (this.loadingPromise) {
- await this.loadingPromise;
- }
-
- return this._containerTypes.asReadonly();
+ getAllContainerTypes(): Promise> {
+ return this.getAllTypes();
}
- async refresh(force = false): Promise {
- if (this.loadingPromise && !force) {
- return this.loadingPromise;
- }
-
- this.loadingPromise = firstValueFrom(this.containersCallService.retrieveAllContainerTypes())
- .finally(() => {
- this._catalogLoading.set(false);
- })
- .then((containerTypes) => {
- this._containerTypes.set(containerTypes);
- this.clearEmptyContainerCache();
- })
- .catch((err) => {
- console.error('Retrieve container types failed', err);
- throw err;
- })
- .finally(() => {
- this.loadingPromise = null;
- });
-
- this._catalogLoading.set(true);
-
- return this.loadingPromise;
+ async getContainerType(typeName: BlockTypeName): Promise {
+ return this.getTypeOrFetch((containerType) => containerType.type === typeName);
}
- async getContainerType(typeName: BlockTypeName) {
- const current = this._containerTypes().find((containerType) => containerType.type === typeName);
- if (current) return current;
-
- if (this.loadingPromise) {
- await this.loadingPromise;
- return this._containerTypes().find((containerType) => containerType.type === typeName);
- }
-
- this._catalogLoading.set(true);
- const containerTypes = await firstValueFrom(this.containersCallService.retrieveAllContainerTypes())
- .finally(() => {
- this._catalogLoading.set(false);
- });
- this._containerTypes.set(containerTypes);
- this.clearEmptyContainerCache();
- return containerTypes.find((containerType) => containerType.type === typeName);
- }
-
- peekContainerType(typeName: BlockTypeName) {
- return this._containerTypes().find((containerType) => containerType.type === typeName) ?? null;
+ peekContainerType(typeName: BlockTypeName): BlockType | null {
+ return this.peekType((containerType) => containerType.type === typeName);
}
createEmptyContainer(containerType: BlockTypeName) {
const cacheKey = String(containerType);
- const cached = this.emptyContainerCache.get(cacheKey);
- if (cached) {
- return of(this.cloneEmptyNode(cached));
- }
- const pending = this.pendingEmptyContainerRequests.get(cacheKey);
- if (pending) {
- return pending.pipe(map((container) => this.cloneEmptyNode(container)));
- }
-
- const request = this.containersCallService.createEmptyContainer(containerType).pipe(
- map((container) => {
- this.emptyContainerCache.set(cacheKey, this.cloneEmptyNode(container));
- return container;
- }),
- finalize(() => {
- this.pendingEmptyContainerRequests.delete(cacheKey);
- }),
- shareReplay(1)
- );
-
- this.pendingEmptyContainerRequests.set(cacheKey, request);
-
- return request.pipe(
- map((container) => this.cloneEmptyNode(container)),
+ return this.emptyContainerCache.getOrCreate(cacheKey, () => this.containersCallService.createEmptyContainer(containerType)).pipe(
catchError((err) => {
console.error('Create empty container failed', err);
return throwError(() => err);
@@ -121,11 +51,7 @@ export class ContainersService {
}
createContainer(containerId: string, configuration: any) {
- this.pendingServerSyncCount.update((count) => count + 1);
- return this.containersCallService.createContainer(containerId, configuration).pipe(
- finalize(() => {
- this.pendingServerSyncCount.update((count) => Math.max(0, count - 1));
- }),
+ return this.serverSync.track(this.containersCallService.createContainer(containerId, configuration)).pipe(
catchError((err) => {
console.error('Create container failed', err);
return throwError(() => err);
@@ -134,7 +60,7 @@ export class ContainersService {
}
validateContainerSubflow(subFlow: FlowData, validationUrl?: string | null) {
- return this.containersCallService.validateContainerSubflow(this.deepClone(subFlow), validationUrl).pipe(
+ return this.containersCallService.validateContainerSubflow(deepClone(subFlow), validationUrl).pipe(
catchError((err) => {
console.error('Validate container subflow failed', err);
return throwError(() => err);
@@ -142,28 +68,11 @@ export class ContainersService {
);
}
- private clearEmptyContainerCache() {
+ protected fetchAll(): Observable {
+ return this.containersCallService.retrieveAllContainerTypes();
+ }
+
+ protected override onLoaded(): void {
this.emptyContainerCache.clear();
- this.pendingEmptyContainerRequests.clear();
- }
-
- private cloneEmptyNode(node: FlowNode): FlowNode {
- const clone = this.deepClone(node);
- return {
- ...clone,
- id: globalThis.crypto?.randomUUID?.() ?? `${Date.now()}`,
- position: undefined
- };
- }
-
- private deepClone(value: T): T {
- if (typeof globalThis.structuredClone === 'function') {
- try {
- return globalThis.structuredClone(value);
- } catch {
- // Some cached payloads may carry non-cloneable runtime fields.
- }
- }
- return JSON.parse(JSON.stringify(value)) as T;
}
}
diff --git a/src/app/services/flows/flows-call.fake.spec.ts b/src/app/services/flows/flows-call.fake.spec.ts
new file mode 100644
index 0000000..da2329d
--- /dev/null
+++ b/src/app/services/flows/flows-call.fake.spec.ts
@@ -0,0 +1,56 @@
+import { TestBed } from '@angular/core/testing';
+import { catchError, of } from 'rxjs';
+import { vi } from 'vitest';
+import { Authorization } from '@services/authorization/authorization';
+import { FlowsCallServiceFake } from './flows-call.fake';
+
+describe('FlowsCallServiceFake', () => {
+ let service: FlowsCallServiceFake;
+
+ beforeEach(() => {
+ TestBed.configureTestingModule({
+ providers: [
+ { provide: Authorization, useValue: { loggedInUser: vi.fn().mockReturnValue({ username: 'Alice', email: null, role: 'USER' }) } }
+ ]
+ });
+ service = TestBed.runInInjectionContext(() => new FlowsCallServiceFake());
+ });
+
+ it('surfaces "flow not found" as an observable error catchError can intercept, not a synchronous throw', async () => {
+ let caught: unknown = null;
+ await new Promise((resolve) => {
+ service.getFlowById('missing-flow').pipe(
+ catchError((error) => {
+ caught = error;
+ return of(null);
+ })
+ ).subscribe(() => resolve());
+ });
+
+ expect(caught).toBeInstanceOf(Error);
+ expect((caught as Error).message).toContain('missing-flow');
+ });
+
+ it('surfaces "flow is finalized" on updateFlow as an observable error, not a synchronous throw', async () => {
+ const finalizedFlow = {
+ id: '1', name: 'A Flow', data: { blocks: [], containers: [], connections: [], dependencies: [] },
+ visibility: 'PUBLIC' as const, author: 'Alice', createdAt: new Date(), status: 'EXECUTABLE' as const,
+ updatedAt: new Date(), finalized: true
+ };
+
+ expect(() => service.updateFlow(finalizedFlow)).not.toThrow();
+
+ let caught: unknown = null;
+ await new Promise((resolve) => {
+ service.updateFlow(finalizedFlow).pipe(
+ catchError((error) => {
+ caught = error;
+ return of(null);
+ })
+ ).subscribe(() => resolve());
+ });
+
+ expect(caught).toBeInstanceOf(Error);
+ expect((caught as Error).message).toBe('Flow is finalized');
+ });
+});
diff --git a/src/app/services/flows/flows-call.fake.ts b/src/app/services/flows/flows-call.fake.ts
index 7956be6..446ac31 100644
--- a/src/app/services/flows/flows-call.fake.ts
+++ b/src/app/services/flows/flows-call.fake.ts
@@ -1,6 +1,6 @@
import { Flow, FlowValidationError } from "@models/flow";
import { FlowsCallServiceBase } from "./flows-call.base";
-import { Observable, of } from "rxjs";
+import { defer, Observable, of } from "rxjs";
import { Authorization } from "@services/authorization/authorization";
import { inject } from "@angular/core";
import { flowFromApi } from "./flow-mapper";
@@ -25,7 +25,7 @@ export class FlowsCallServiceFake extends FlowsCallServiceBase {
}
override getFlowById(flowId: string): Observable {
- return of(this.requireFlow(flowId));
+ return defer(() => of(this.requireFlow(flowId)));
}
authorizationService = inject(Authorization);
@@ -41,11 +41,13 @@ export class FlowsCallServiceFake extends FlowsCallServiceBase {
}
override updateFlow(flow: Flow) {
- if (flow.finalized) {
- throw new Error('Flow is finalized');
- }
- this.data[flow.id] = flow;
- return of(flow);
+ return defer(() => {
+ if (flow.finalized) {
+ throw new Error('Flow is finalized');
+ }
+ this.data[flow.id] = flow;
+ return of(flow);
+ });
}
override createFlow(flow: Pick): Observable {
@@ -75,44 +77,50 @@ export class FlowsCallServiceFake extends FlowsCallServiceBase {
}
override deleteFlow(flowId: string): Observable {
- const flow = this.requireFlow(flowId);
- if (flow.finalized) {
- throw new Error('Flow is finalized');
- }
- delete this.data[flowId];
- return of(void 0);
+ return defer(() => {
+ const flow = this.requireFlow(flowId);
+ if (flow.finalized) {
+ throw new Error('Flow is finalized');
+ }
+ delete this.data[flowId];
+ return of(void 0);
+ });
}
override updatePublished(flowId: string, value: boolean): Observable {
- const flow = this.requireFlow(flowId);
- this.requireOwner(flow);
- const updated = {
- ...flow,
- published: value,
- visibility: value ? 'PUBLIC' : 'PRIVATE',
- updatedAt: new Date()
- } satisfies Flow;
- this.data[flowId] = updated;
- return of(updated);
+ return defer(() => {
+ const flow = this.requireFlow(flowId);
+ this.requireOwner(flow);
+ const updated = {
+ ...flow,
+ published: value,
+ visibility: value ? 'PUBLIC' : 'PRIVATE',
+ updatedAt: new Date()
+ } satisfies Flow;
+ this.data[flowId] = updated;
+ return of(updated);
+ });
}
override finalizeFlow(flowId: string): Observable {
- const flow = this.requireFlow(flowId);
- this.requireOwner(flow);
- if (flow.finalized) {
- return of(flow);
- }
- const updated = {
- ...flow,
- finalized: true,
- updatedAt: new Date()
- } satisfies Flow;
- this.data[flowId] = updated;
- return of(updated);
+ return defer(() => {
+ const flow = this.requireFlow(flowId);
+ this.requireOwner(flow);
+ if (flow.finalized) {
+ return of(flow);
+ }
+ const updated = {
+ ...flow,
+ finalized: true,
+ updatedAt: new Date()
+ } satisfies Flow;
+ this.data[flowId] = updated;
+ return of(updated);
+ });
}
override getFlowValidation(flowId: string): Observable {
- return of(this.requireFlow(flowId).validationErrors ?? []);
+ return defer(() => of(this.requireFlow(flowId).validationErrors ?? []));
}
}
const testDataFlow ={
diff --git a/src/app/services/shared/catalog-store.ts b/src/app/services/shared/catalog-store.ts
new file mode 100644
index 0000000..57c6433
--- /dev/null
+++ b/src/app/services/shared/catalog-store.ts
@@ -0,0 +1,96 @@
+import { Signal, signal } from '@angular/core';
+import { firstValueFrom, Observable } from 'rxjs';
+
+/**
+ * Shared loading/caching state machine for a "type catalog" (block types,
+ * container types): a signal holding the last-loaded list, a single in-flight
+ * load shared across concurrent callers, and automatic retry after a failed
+ * initial load (a first failed fetch no longer leaves the catalog permanently
+ * empty — the next call retries instead of silently returning nothing).
+ *
+ * Subclasses provide `fetchAll()` (the HTTP call) and the domain-specific
+ * public method names (`getAllBlocksTypes`, `getAllContainerTypes`, ...) that
+ * delegate to the protected methods here.
+ */
+export abstract class CatalogStore {
+ private toInit = true;
+ private loadingPromise: Promise | null = null;
+ private readonly _loading = signal(false);
+ private readonly _types = signal([]);
+
+ protected readonly loading = this._loading.asReadonly();
+ protected readonly types = this._types.asReadonly();
+
+ /** Fetches the full catalog from the backend. */
+ protected abstract fetchAll(): Observable;
+ /** Label used in the `console.error` logged when a fetch fails. */
+ protected abstract readonly loadErrorLabel: string;
+ /** Called whenever a fresh catalog is stored, e.g. to invalidate derived caches. */
+ protected onLoaded(): void {}
+
+ protected hasLoadedTypes(): boolean {
+ return this._types().length > 0 || (!this.toInit && !this.loadingPromise);
+ }
+
+ protected async getAllTypes(): Promise> {
+ if (this.toInit) {
+ this.toInit = false;
+ try {
+ await this.refresh();
+ } catch (err) {
+ this.toInit = true;
+ throw err;
+ }
+ } else if (this.loadingPromise) {
+ await this.loadingPromise;
+ }
+
+ return this.types;
+ }
+
+ protected async refresh(force = false): Promise {
+ if (this.loadingPromise && !force) {
+ return this.loadingPromise;
+ }
+
+ this.loadingPromise = firstValueFrom(this.fetchAll())
+ .finally(() => {
+ this._loading.set(false);
+ })
+ .then((types) => {
+ this._types.set(types);
+ this.onLoaded();
+ })
+ .catch((err) => {
+ console.error(this.loadErrorLabel, err);
+ throw err;
+ })
+ .finally(() => {
+ this.loadingPromise = null;
+ });
+
+ this._loading.set(true);
+
+ return this.loadingPromise;
+ }
+
+ protected async getTypeOrFetch(predicate: (type: T) => boolean): Promise {
+ const current = this._types().find(predicate);
+ if (current) return current;
+
+ if (this.loadingPromise) {
+ await this.loadingPromise;
+ return this._types().find(predicate);
+ }
+
+ this._loading.set(true);
+ const types = await firstValueFrom(this.fetchAll()).finally(() => this._loading.set(false));
+ this._types.set(types);
+ this.onLoaded();
+ return types.find(predicate);
+ }
+
+ protected peekType(predicate: (type: T) => boolean): T | null {
+ return this._types().find(predicate) ?? null;
+ }
+}
diff --git a/src/app/services/shared/deep-clone.ts b/src/app/services/shared/deep-clone.ts
new file mode 100644
index 0000000..400d1ba
--- /dev/null
+++ b/src/app/services/shared/deep-clone.ts
@@ -0,0 +1,11 @@
+/** `structuredClone` with a JSON round-trip fallback for non-cloneable runtime fields. */
+export function deepClone(value: T): T {
+ if (typeof globalThis.structuredClone === 'function') {
+ try {
+ return globalThis.structuredClone(value);
+ } catch {
+ // Some cached payloads may carry non-cloneable runtime fields.
+ }
+ }
+ return JSON.parse(JSON.stringify(value)) as T;
+}
diff --git a/src/app/services/shared/empty-node-cache.ts b/src/app/services/shared/empty-node-cache.ts
new file mode 100644
index 0000000..59c8b45
--- /dev/null
+++ b/src/app/services/shared/empty-node-cache.ts
@@ -0,0 +1,54 @@
+import { finalize, map, Observable, of, shareReplay } from 'rxjs';
+import { deepClone } from './deep-clone';
+
+/**
+ * Caches "empty" node templates (an empty block/container fresh from the
+ * backend) keyed by a cache key (usually the type name), de-duplicating
+ * concurrent requests for the same key and handing every caller its own
+ * clone with a fresh id so mutating one instance never leaks into another.
+ */
+export class EmptyNodeCache {
+ private readonly cache = new Map();
+ private readonly pendingRequests = new Map>();
+
+ getOrCreate(cacheKey: string, request: () => Observable): Observable {
+ const cached = this.cache.get(cacheKey);
+ if (cached) {
+ return of(this.cloneWithNewId(cached));
+ }
+
+ const pending = this.pendingRequests.get(cacheKey);
+ if (pending) {
+ return pending.pipe(map((node) => this.cloneWithNewId(node)));
+ }
+
+ const shared = request().pipe(
+ map((node) => {
+ this.cache.set(cacheKey, this.cloneWithNewId(node));
+ return node;
+ }),
+ finalize(() => {
+ this.pendingRequests.delete(cacheKey);
+ }),
+ shareReplay(1)
+ );
+
+ this.pendingRequests.set(cacheKey, shared);
+
+ return shared.pipe(map((node) => this.cloneWithNewId(node)));
+ }
+
+ clear(): void {
+ this.cache.clear();
+ this.pendingRequests.clear();
+ }
+
+ private cloneWithNewId(node: T): T {
+ const clone = deepClone(node);
+ return {
+ ...clone,
+ id: globalThis.crypto?.randomUUID?.() ?? `${Date.now()}`,
+ position: undefined
+ };
+ }
+}
diff --git a/src/app/services/shared/flow-node-mapping.spec.ts b/src/app/services/shared/flow-node-mapping.spec.ts
new file mode 100644
index 0000000..c2cf9a6
--- /dev/null
+++ b/src/app/services/shared/flow-node-mapping.spec.ts
@@ -0,0 +1,113 @@
+import { attachSharedDefinitions, toApiPath, toNullableString, toPorts, toPosition, toRecord, toSchema, toValueKinds } from './flow-node-mapping';
+
+describe('toRecord', () => {
+ it('returns the object as-is', () => {
+ expect(toRecord({ a: 1 })).toEqual({ a: 1 });
+ });
+
+ it('returns an empty object for arrays, null, primitives', () => {
+ expect(toRecord([1, 2])).toEqual({});
+ expect(toRecord(null)).toEqual({});
+ expect(toRecord('x')).toEqual({});
+ });
+});
+
+describe('toNullableString', () => {
+ it('passes through non-empty strings and nulls everything else', () => {
+ expect(toNullableString('hi')).toBe('hi');
+ expect(toNullableString('')).toBeNull();
+ expect(toNullableString(42)).toBeNull();
+ expect(toNullableString(null)).toBeNull();
+ });
+});
+
+describe('toApiPath', () => {
+ it('passes absolute http(s) URLs through unchanged', () => {
+ expect(toApiPath('https://example.com/x')).toBe('https://example.com/x');
+ });
+
+ it('prefixes a relative path with the API base URL', () => {
+ expect(toApiPath('/retriever/x')).toMatch(/\/retriever\/x$/);
+ expect(toApiPath('retriever/x')).toMatch(/\/retriever\/x$/);
+ });
+
+ it('returns null for anything that is not a non-empty string', () => {
+ expect(toApiPath('')).toBeNull();
+ expect(toApiPath(null)).toBeNull();
+ });
+});
+
+describe('toPosition', () => {
+ it('reads a valid {x,y} pair', () => {
+ expect(toPosition({ x: 1, y: 2 })).toEqual({ x: 1, y: 2 });
+ });
+
+ it('returns undefined when x/y are missing or not numeric', () => {
+ expect(toPosition({ x: 1 })).toBeUndefined();
+ expect(toPosition({ x: '1', y: 2 })).toBeUndefined();
+ expect(toPosition(null)).toBeUndefined();
+ });
+});
+
+describe('toSchema', () => {
+ it('accepts a plain object and rejects arrays/primitives/null', () => {
+ expect(toSchema({ type: 'object' })).toEqual({ type: 'object' });
+ expect(toSchema([1])).toBeNull();
+ expect(toSchema('x')).toBeNull();
+ expect(toSchema(null)).toBeNull();
+ });
+});
+
+describe('attachSharedDefinitions', () => {
+ it('returns null when there is no schema', () => {
+ expect(attachSharedDefinitions(null, { a: 1 })).toBeNull();
+ });
+
+ it('returns the schema unchanged when there are no shared definitions', () => {
+ const schema = { type: 'object' };
+ expect(attachSharedDefinitions(schema, null)).toBe(schema);
+ expect(attachSharedDefinitions(schema, {})).toBe(schema);
+ });
+
+ it('merges shared definitions into the schema, letting schema-local ones win', () => {
+ const schema = { type: 'object', sharedDefinitions: { a: 'local' } };
+ expect(attachSharedDefinitions(schema, { a: 'shared', b: 'shared' })).toEqual({
+ type: 'object',
+ sharedDefinitions: { a: 'local', b: 'shared' }
+ });
+ });
+});
+
+describe('toValueKinds', () => {
+ it('falls back to a single kind when raw is not an array', () => {
+ expect(toValueKinds(null, { type: 'TEXT', multiple: false })).toEqual([{ type: 'TEXT', multiple: false }]);
+ });
+
+ it('maps well-formed entries and drops entries without a string type', () => {
+ expect(toValueKinds([{ type: 'FILE', multiple: true }, { multiple: true }], { type: 'TEXT', multiple: false }))
+ .toEqual([{ type: 'FILE', multiple: true }]);
+ });
+
+ it('falls back when the array yields no usable entries', () => {
+ expect(toValueKinds([{ multiple: true }], { type: 'TEXT', multiple: false })).toEqual([{ type: 'TEXT', multiple: false }]);
+ });
+});
+
+describe('toPorts', () => {
+ it('maps named ports and derives valueKinds from type/multiple', () => {
+ expect(toPorts([{ name: 'input', type: 'FILE', multiple: true }])).toEqual([
+ { name: 'input', type: 'FILE', multiple: true, valueKinds: [{ type: 'FILE', multiple: true }] }
+ ]);
+ });
+
+ it('drops ports without a usable name', () => {
+ expect(toPorts([{ type: 'TEXT' }])).toEqual([]);
+ });
+
+ it('returns the fallback (default []) when raw is not an array', () => {
+ expect(toPorts(null)).toEqual([]);
+ expect(toPorts(null, [{ name: 'fallback', type: 'TEXT', multiple: false }])).toEqual([
+ { name: 'fallback', type: 'TEXT', multiple: false }
+ ]);
+ });
+});
diff --git a/src/app/services/shared/flow-node-mapping.ts b/src/app/services/shared/flow-node-mapping.ts
new file mode 100644
index 0000000..adf5867
--- /dev/null
+++ b/src/app/services/shared/flow-node-mapping.ts
@@ -0,0 +1,89 @@
+import { environment } from '@environment';
+
+/**
+ * Response-mapping helpers shared by `blocks-call.ts` and `containers-call.ts`
+ * (and their `.fake.ts` counterparts): both map the same wire shape — ports,
+ * value kinds, position, JSON-schema — into the app's `FlowBlock`/`FlowNode`
+ * domain types.
+ */
+
+export function toRecord(value: unknown): Record {
+ if (!value || typeof value !== 'object' || Array.isArray(value)) return {};
+ return value as Record;
+}
+
+export function toNullableString(value: unknown): string | null {
+ return typeof value === 'string' && value.length > 0 ? value : null;
+}
+
+export function toApiPath(value: unknown): string | null {
+ if (typeof value !== 'string' || value.length === 0) return null;
+ if (/^https?:\/\//.test(value)) return value;
+ return `${environment.apiUrl}${value.startsWith('/') ? value : `/${value}`}`;
+}
+
+export function toPosition(raw: unknown): { x: number; y: number } | undefined {
+ const value = toRecord(raw);
+ const x = value['x'];
+ const y = value['y'];
+ if (typeof x !== 'number' || typeof y !== 'number') return undefined;
+ return { x, y };
+}
+
+export function toSchema(raw: unknown): Record | null {
+ if (!raw || typeof raw !== 'object' || Array.isArray(raw)) return null;
+ return raw as Record;
+}
+
+export function attachSharedDefinitions(
+ schema: Record | null,
+ sharedDefinitions: Record | null
+): Record | null {
+ if (!schema) return null;
+ if (!sharedDefinitions || !Object.keys(sharedDefinitions).length) return schema;
+
+ return {
+ ...schema,
+ sharedDefinitions: {
+ ...sharedDefinitions,
+ ...toRecord(schema['sharedDefinitions'])
+ }
+ };
+}
+
+export function toValueKinds(raw: unknown, fallback: { type: string; multiple: boolean }): Array<{ type: string; multiple: boolean }> {
+ if (!Array.isArray(raw)) {
+ return [{ type: fallback.type, multiple: fallback.multiple }];
+ }
+
+ const kinds = raw
+ .map((item) => toRecord(item))
+ .filter((item) => typeof item['type'] === 'string')
+ .map((item) => ({
+ type: String(item['type'] ?? fallback.type),
+ multiple: Boolean(item['multiple'] ?? false)
+ }));
+
+ return kinds.length ? kinds : [{ type: fallback.type, multiple: fallback.multiple }];
+}
+
+export function toPorts(
+ raw: unknown,
+ fallback: Array<{ name: string; type: string; multiple: boolean }> = []
+) {
+ if (!Array.isArray(raw)) return fallback;
+ return raw
+ .map((port) => toRecord(port))
+ .filter((port) => typeof port['name'] === 'string' && (port['name'] as string).length > 0)
+ .map((port) => {
+ const type = String(port['type'] ?? 'TEXT');
+ const multiple = Boolean(port['multiple'] ?? false);
+ return {
+ ...port,
+ name: String(port['name']),
+ type,
+ multiple,
+ valueKinds: toValueKinds(port['valueKinds'], { type, multiple })
+ };
+ });
+}
diff --git a/src/app/services/shared/http-error.util.spec.ts b/src/app/services/shared/http-error.util.spec.ts
new file mode 100644
index 0000000..048888b
--- /dev/null
+++ b/src/app/services/shared/http-error.util.spec.ts
@@ -0,0 +1,46 @@
+import { HttpErrorResponse } from '@angular/common/http';
+import { lastValueFrom } from 'rxjs';
+import { extractHttpErrorMessage, toHttpError } from './http-error.util';
+
+describe('extractHttpErrorMessage', () => {
+ it('reads a plain string error body', () => {
+ const error = new HttpErrorResponse({ error: 'Something broke', status: 400 });
+ expect(extractHttpErrorMessage(error)).toBe('Something broke');
+ });
+
+ it('prefers message, then error, then details from an object body', () => {
+ expect(extractHttpErrorMessage(new HttpErrorResponse({ error: { message: 'msg' }, status: 400 }))).toBe('msg');
+ expect(extractHttpErrorMessage(new HttpErrorResponse({ error: { error: 'err' }, status: 400 }))).toBe('err');
+ expect(extractHttpErrorMessage(new HttpErrorResponse({ error: { details: 'det' }, status: 400 }))).toBe('det');
+ });
+
+ it('returns null when nothing usable is present', () => {
+ expect(extractHttpErrorMessage(new HttpErrorResponse({ error: {}, status: 500 }))).toBeNull();
+ expect(extractHttpErrorMessage(new HttpErrorResponse({ error: null, status: 500 }))).toBeNull();
+ });
+});
+
+describe('toHttpError', () => {
+ it('prefers the backend message over the status fallback', async () => {
+ const error = new HttpErrorResponse({ error: { message: 'Backend said no' }, status: 400 });
+ await expect(lastValueFrom(toHttpError(error, { 400: 'Fallback message' }))).rejects.toThrow('Backend said no');
+ });
+
+ it('falls back to the status-keyed message when the backend gives nothing usable', async () => {
+ const error = new HttpErrorResponse({ error: {}, status: 403 });
+ await expect(lastValueFrom(toHttpError(error, { 403: 'Admin access required.' }))).rejects.toThrow('Admin access required.');
+ });
+
+ it('falls back to a generic message when the status is not mapped', async () => {
+ const error = new HttpErrorResponse({ error: {}, status: 418 });
+ await expect(lastValueFrom(toHttpError(error, {}))).rejects.toThrow('Request failed.');
+ });
+
+ it('passes an existing Error straight through', async () => {
+ await expect(lastValueFrom(toHttpError(new Error('boom'), {}))).rejects.toThrow('boom');
+ });
+
+ it('wraps a non-Error, non-HttpErrorResponse value in a generic error', async () => {
+ await expect(lastValueFrom(toHttpError('a raw string', {}))).rejects.toThrow('Request failed.');
+ });
+});
diff --git a/src/app/services/shared/http-error.util.ts b/src/app/services/shared/http-error.util.ts
new file mode 100644
index 0000000..0030ea2
--- /dev/null
+++ b/src/app/services/shared/http-error.util.ts
@@ -0,0 +1,47 @@
+import { HttpErrorResponse } from '@angular/common/http';
+import { Observable, throwError } from 'rxjs';
+
+/**
+ * Reads a human-readable message out of a backend error body, trying the
+ * conventional `message`/`error`/`details` string fields in that order.
+ */
+export function extractHttpErrorMessage(error: HttpErrorResponse): string | null {
+ const payload = error.error;
+ if (typeof payload === 'string' && payload.trim().length > 0) {
+ return payload.trim();
+ }
+ if (payload && typeof payload === 'object') {
+ const record = payload as Record;
+ const directMessage = record['message'];
+ if (typeof directMessage === 'string' && directMessage.trim().length > 0) {
+ return directMessage.trim();
+ }
+ const errorMessage = record['error'];
+ if (typeof errorMessage === 'string' && errorMessage.trim().length > 0) {
+ return errorMessage.trim();
+ }
+ const details = record['details'];
+ if (typeof details === 'string' && details.trim().length > 0) {
+ return details.trim();
+ }
+ }
+ return null;
+}
+
+/**
+ * Converts any thrown/caught value into an `Observable` error carrying a
+ * human-readable `Error`: prefers the backend's own message, falls back to
+ * a status-keyed message, then a generic one. Existing `Error`s pass through.
+ */
+export function toHttpError(error: unknown, fallbackByStatus: Record): Observable {
+ if (error instanceof HttpErrorResponse) {
+ const message = extractHttpErrorMessage(error)
+ ?? fallbackByStatus[error.status]
+ ?? 'Request failed.';
+ return throwError(() => new Error(message));
+ }
+ if (error instanceof Error) {
+ return throwError(() => error);
+ }
+ return throwError(() => new Error('Request failed.'));
+}
diff --git a/src/app/services/shared/pending-sync-counter.ts b/src/app/services/shared/pending-sync-counter.ts
new file mode 100644
index 0000000..9a261fd
--- /dev/null
+++ b/src/app/services/shared/pending-sync-counter.ts
@@ -0,0 +1,15 @@
+import { computed, signal } from '@angular/core';
+import { finalize, Observable } from 'rxjs';
+
+/** Tracks how many "sync this to the server" requests are currently in flight. */
+export class PendingSyncCounter {
+ private readonly count = signal(0);
+ readonly active = computed(() => this.count() > 0);
+
+ track(source: Observable): Observable {
+ this.count.update((current) => current + 1);
+ return source.pipe(
+ finalize(() => this.count.update((current) => Math.max(0, current - 1)))
+ );
+ }
+}
diff --git a/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.css b/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.css
index 009fca1..525f718 100644
--- a/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.css
+++ b/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.css
@@ -1,94 +1 @@
-.admin-reset-password-backdrop {
- position: fixed;
- inset: 0;
- z-index: 9999;
- display: flex;
- align-items: center;
- justify-content: center;
- padding: 16px;
- background: rgba(15, 23, 42, 0.5);
-}
-
-.admin-reset-password-modal {
- width: min(520px, calc(100vw - 32px));
- max-height: min(80vh, 720px);
- overflow: auto;
- border: 1px solid #e2e8f0;
- border-radius: 20px;
- background: #ffffff;
- box-shadow: 0 24px 64px rgba(15, 23, 42, 0.28);
- padding: 24px;
-}
-
-.admin-reset-password-header {
- display: flex;
- align-items: flex-start;
- justify-content: space-between;
- gap: 16px;
-}
-
-.admin-reset-password-title {
- margin: 0;
- color: #0f172a;
- font-size: 20px;
- font-weight: 700;
-}
-
-.admin-reset-password-subtitle {
- margin: 4px 0 0;
- color: #64748b;
- font-size: 13px;
-}
-
-.admin-reset-password-form {
- display: flex;
- flex-direction: column;
- gap: 14px;
- margin-top: 20px;
-}
-
-.admin-reset-password-field {
- width: 100%;
-}
-
-.admin-reset-password-checklist {
- display: flex;
- flex-direction: column;
- gap: 6px;
- margin-top: -6px;
- padding: 0 2px;
-}
-
-.admin-reset-password-check {
- display: flex;
- align-items: center;
- gap: 8px;
- color: #b91c1c;
- font-size: 12px;
-}
-
-.admin-reset-password-check--ok {
- color: #15803d;
-}
-
-.admin-reset-password-check-icon {
- width: 16px;
- height: 16px;
- font-size: 16px;
-}
-
-.admin-reset-password-error {
- border: 1px solid #fecaca;
- border-radius: 12px;
- background: #fff1f2;
- color: #b91c1c;
- font-size: 13px;
- padding: 10px 12px;
-}
-
-.admin-reset-password-actions {
- display: flex;
- justify-content: flex-end;
- gap: 10px;
- margin-top: 4px;
-}
+@import '../password-dialog-chrome.css';
diff --git a/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.html b/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.html
index eff6eb7..72288c0 100644
--- a/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.html
+++ b/src/app/shared/admin-reset-password-dialog/admin-reset-password-dialog.html
@@ -1,15 +1,15 @@
-