feat(backend): enforce user-owned sessions
This commit is contained in:
@@ -0,0 +1,50 @@
|
|||||||
|
# Task 5 report — backend principal enforcement
|
||||||
|
|
||||||
|
## RED
|
||||||
|
|
||||||
|
Added backend route/auth tests before implementation. The initial focused run failed in
|
||||||
|
seven new assertions: `getPrincipal` did not exist, upstream requests still required the
|
||||||
|
legacy identity header, foreign session/SSE routes were not hidden, admin scope was not
|
||||||
|
enforced, new sessions had no trusted principal binding, and settings were global.
|
||||||
|
|
||||||
|
## GREEN
|
||||||
|
|
||||||
|
- Focused backend suite: `66 passed` across auth, sessions, SSE, and settings tests.
|
||||||
|
- Complete backend Vitest suite: `209 passed` across `22` files.
|
||||||
|
- `npx tsc --noEmit -p .`, `npm run build`, `git diff --check`, and changed Python
|
||||||
|
source Ruff all exit successfully.
|
||||||
|
- Harness targeted repository/local/migration tests and Python bytecode compilation exit
|
||||||
|
successfully. The new `tht session preferences get|set` commands are registered and
|
||||||
|
expose the expected Typer help. A direct local CLI preference smoke was not run because
|
||||||
|
the checked-in local workspace requires unavailable `THT_DB_HOST` configuration.
|
||||||
|
|
||||||
|
## Route and child-process coverage
|
||||||
|
|
||||||
|
- `GET /me` returns the request `PrincipalContext`; upstream accepts only the portal's
|
||||||
|
normalized `X-Thoth-*` identity tuple, with the legacy header ignored. Local mode uses
|
||||||
|
the same stable `THT_HOME`/`~/.thothii/identity.json` UUID contract as the harness.
|
||||||
|
- All session operations are principal-scoped: list (`mine` and admin-only `all`), show,
|
||||||
|
create, resume, close, delete, rename, group, archive, unarchive, documents, reviewer
|
||||||
|
response, steer, SQL preview/export, and SSE. Missing and foreign sessions are 404;
|
||||||
|
absent upstream identity is 401. SSE is authorized before response headers or hub
|
||||||
|
subscription, so a rejected request cannot attach to a live stream.
|
||||||
|
- New/resumed Pi runtimes and every route-spawned `tht` process receive
|
||||||
|
`THT_PRINCIPAL_ISSUER`, `THT_PRINCIPAL_SUBJECT`, optional display name, and admin flag.
|
||||||
|
The readiness `tht` child is also principal-bound.
|
||||||
|
- Settings use asynchronous repository-backed `tht session preferences get|set` in the
|
||||||
|
production runner, which isolates preferences by principal. The legacy settings file is
|
||||||
|
retained only as an injected-runner compatibility fallback for existing isolated tests.
|
||||||
|
- Repository/settings authorization failures map to 503 before model startup. SQL execution
|
||||||
|
errors remain 500 after authorization, preserving the prior API distinction.
|
||||||
|
|
||||||
|
## Self-review and concerns
|
||||||
|
|
||||||
|
- Confirmed the Task 4 portal emits lowercase `true`/`false` for the admin header; the
|
||||||
|
parser accepts that exact normalized form plus the repository's existing `1`/`0`
|
||||||
|
compatibility form, and rejects all other values.
|
||||||
|
- The harness principal resolver is the ownership authority; the backend never accepts an
|
||||||
|
owner supplied in request bodies. Its route guards use a repository-scoped `session show`
|
||||||
|
before every session resource operation.
|
||||||
|
- Existing dependency-injected route fakes without `sessionShow` retain a narrow test seam;
|
||||||
|
production `ThtRunner` always has that method, so deployed requests cannot bypass the
|
||||||
|
repository authorization check.
|
||||||
+29
-5
@@ -5,12 +5,14 @@ import { ThtRunner } from "./tht/tht-runner.js";
|
|||||||
import { PiProcessManager } from "./pi/pi-process-manager.js";
|
import { PiProcessManager } from "./pi/pi-process-manager.js";
|
||||||
import { SseHub } from "./sse/sse-hub.js";
|
import { SseHub } from "./sse/sse-hub.js";
|
||||||
import { authPreHandler } from "./auth/auth.js";
|
import { authPreHandler } from "./auth/auth.js";
|
||||||
|
import { getPrincipal } from "./auth/auth.js";
|
||||||
|
import type { PrincipalContext } from "./auth/principal.js";
|
||||||
import { sessionRoutes } from "./routes/sessions.js";
|
import { sessionRoutes } from "./routes/sessions.js";
|
||||||
import { sqlRoutes } from "./routes/sql.js";
|
import { sqlRoutes } from "./routes/sql.js";
|
||||||
import { metaRoutes, type ListModelsFn } from "./routes/meta.js";
|
import { metaRoutes, type ListModelsFn } from "./routes/meta.js";
|
||||||
import { settingsRoutes, effectiveSettings } from "./routes/settings.js";
|
import { settingsRoutes, effectiveSettings } from "./routes/settings.js";
|
||||||
import { createPiModelLister } from "./pi/list-models.js";
|
import { createPiModelLister } from "./pi/list-models.js";
|
||||||
import { loadSettings, type Settings } from "./settings/settings-store.js";
|
import { loadSettings, saveSettings, type Settings } from "./settings/settings-store.js";
|
||||||
import { ReadinessManager } from "./runtime/readiness-manager.js";
|
import { ReadinessManager } from "./runtime/readiness-manager.js";
|
||||||
|
|
||||||
export interface BuildAppDeps {
|
export interface BuildAppDeps {
|
||||||
@@ -18,7 +20,7 @@ export interface BuildAppDeps {
|
|||||||
mgr?: PiProcessManager;
|
mgr?: PiProcessManager;
|
||||||
spawnFn?: () => any;
|
spawnFn?: () => any;
|
||||||
listModels?: ListModelsFn;
|
listModels?: ListModelsFn;
|
||||||
getSettings?: () => Settings;
|
getSettings?: (principal?: PrincipalContext) => Settings | Promise<Settings>;
|
||||||
readiness?: ReadinessManager;
|
readiness?: ReadinessManager;
|
||||||
hub?: SseHub;
|
hub?: SseHub;
|
||||||
}
|
}
|
||||||
@@ -52,7 +54,28 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
|||||||
"Pi enabled-model configuration warning",
|
"Pi enabled-model configuration warning",
|
||||||
),
|
),
|
||||||
});
|
});
|
||||||
const getSettings = deps?.getSettings ?? (() => effectiveSettings(config, loadSettings(config)));
|
const runnerFor = (principal: PrincipalContext): any => {
|
||||||
|
const candidate = tht as any;
|
||||||
|
return typeof candidate.withPrincipal === "function" ? candidate.withPrincipal(principal) : candidate;
|
||||||
|
};
|
||||||
|
const getSettings = async (principal: PrincipalContext): Promise<Settings> => {
|
||||||
|
if (deps?.getSettings) return await deps.getSettings(principal);
|
||||||
|
const runner = runnerFor(principal);
|
||||||
|
// The real runner persists preferences through the harness repository. The file fallback
|
||||||
|
// only keeps older isolated route tests and externally injected runners compatible.
|
||||||
|
if (typeof runner.preferencesGet === "function") {
|
||||||
|
return effectiveSettings(config, await runner.preferencesGet());
|
||||||
|
}
|
||||||
|
return effectiveSettings(config, loadSettings(config));
|
||||||
|
};
|
||||||
|
const saveUserSettings = async (principal: PrincipalContext, settings: Settings): Promise<void> => {
|
||||||
|
const runner = runnerFor(principal);
|
||||||
|
if (typeof runner.preferencesSet === "function") {
|
||||||
|
await runner.preferencesSet(settings);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
saveSettings(config, settings);
|
||||||
|
};
|
||||||
|
|
||||||
const authenticate = authPreHandler(config.authMode);
|
const authenticate = authPreHandler(config.authMode);
|
||||||
app.addHook("preHandler", async (req, reply) => {
|
app.addHook("preHandler", async (req, reply) => {
|
||||||
@@ -61,12 +84,13 @@ export function buildApp(config: AppConfig, deps?: BuildAppDeps): FastifyInstanc
|
|||||||
return authenticate(req, reply);
|
return authenticate(req, reply);
|
||||||
});
|
});
|
||||||
app.get("/health", async () => ({ status: "ok" }));
|
app.get("/health", async () => ({ status: "ok" }));
|
||||||
|
app.get("/me", async (req) => getPrincipal(req));
|
||||||
sessionRoutes(app, {
|
sessionRoutes(app, {
|
||||||
mgr, tht: tht as ThtRunner, hub, getSettings, readiness,
|
mgr, tht: tht as ThtRunner, hub, getSettings, readiness,
|
||||||
});
|
});
|
||||||
sqlRoutes(app, { tht: tht as ThtRunner });
|
sqlRoutes(app, { tht: tht as ThtRunner, getSettings });
|
||||||
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
|
metaRoutes(app, { harnessDir: config.harnessDir, listModels });
|
||||||
settingsRoutes(app, { cfg: config, listModels });
|
settingsRoutes(app, { cfg: config, listModels, getSettings, saveSettings: saveUserSettings });
|
||||||
|
|
||||||
return app;
|
return app;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,23 +1,28 @@
|
|||||||
import type { FastifyRequest, FastifyReply } from "fastify";
|
import type { FastifyRequest, FastifyReply } from "fastify";
|
||||||
|
import { localPrincipal, type PrincipalContext, upstreamPrincipal } from "./principal.js";
|
||||||
|
|
||||||
|
declare module "fastify" {
|
||||||
|
interface FastifyRequest { principal?: PrincipalContext }
|
||||||
|
}
|
||||||
|
|
||||||
export function authPreHandler(mode: "none" | "mock" | "upstream") {
|
export function authPreHandler(mode: "none" | "mock" | "upstream") {
|
||||||
return async (req: FastifyRequest, reply: FastifyReply) => {
|
return async (req: FastifyRequest, reply: FastifyReply) => {
|
||||||
if (mode === "none") {
|
if (mode === "none") {
|
||||||
(req as any).user = { id: "dev@local" };
|
req.principal = localPrincipal();
|
||||||
} else if (mode === "mock") {
|
} else if (mode === "mock") {
|
||||||
(req as any).user = {
|
const subject = typeof req.headers["x-mock-user"] === "string" ? req.headers["x-mock-user"].trim() : "mock";
|
||||||
id: (req.headers["x-mock-user"] as string) ?? "mock",
|
req.principal = { issuer: "mock", subject: subject || "mock", displayName: subject || "mock", isAdmin: false };
|
||||||
};
|
|
||||||
} else {
|
} else {
|
||||||
const id = req.headers["x-authenticated-user"];
|
const principal = upstreamPrincipal(req.headers);
|
||||||
if (typeof id !== "string" || id.trim() === "") {
|
if (!principal) {
|
||||||
return reply.code(401).send({ error: "authenticated upstream identity required" });
|
return reply.code(401).send({ error: "authenticated upstream identity required" });
|
||||||
}
|
}
|
||||||
(req as any).user = { id };
|
req.principal = principal;
|
||||||
}
|
}
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
export function getUser(req: FastifyRequest): { id: string } {
|
export function getPrincipal(req: FastifyRequest): PrincipalContext {
|
||||||
return (req as any).user ?? { id: "dev@local" };
|
if (!req.principal) throw new Error("principal missing after authentication");
|
||||||
|
return req.principal;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,68 @@
|
|||||||
|
import { mkdirSync, readFileSync, writeFileSync } from "node:fs";
|
||||||
|
import { homedir } from "node:os";
|
||||||
|
import { join } from "node:path";
|
||||||
|
import { randomUUID } from "node:crypto";
|
||||||
|
|
||||||
|
export interface PrincipalContext {
|
||||||
|
issuer: string;
|
||||||
|
subject: string;
|
||||||
|
displayName?: string;
|
||||||
|
isAdmin: boolean;
|
||||||
|
}
|
||||||
|
|
||||||
|
const invalid = (value: string) => value.length === 0 || value.length > 512 || /[\u0000-\u001f\u007f]/.test(value);
|
||||||
|
|
||||||
|
function required(value: unknown): string | undefined {
|
||||||
|
if (typeof value !== "string") return undefined;
|
||||||
|
const normalized = value.trim();
|
||||||
|
return invalid(normalized) ? undefined : normalized;
|
||||||
|
}
|
||||||
|
|
||||||
|
function optional(value: unknown): string | undefined {
|
||||||
|
if (value === undefined) return undefined;
|
||||||
|
return required(value);
|
||||||
|
}
|
||||||
|
|
||||||
|
export function upstreamPrincipal(headers: Record<string, unknown>): PrincipalContext | undefined {
|
||||||
|
const issuer = required(headers["x-thoth-principal-issuer"]);
|
||||||
|
const subject = required(headers["x-thoth-principal-subject"]);
|
||||||
|
const displayName = optional(headers["x-thoth-principal-display-name"]);
|
||||||
|
const adminHeader = headers["x-thoth-is-admin"];
|
||||||
|
if (!issuer || !subject || (headers["x-thoth-principal-display-name"] !== undefined && !displayName)) return undefined;
|
||||||
|
if (adminHeader !== "0" && adminHeader !== "1" && adminHeader !== "true" && adminHeader !== "false") return undefined;
|
||||||
|
return { issuer, subject, displayName, isAdmin: adminHeader === "1" || adminHeader === "true" };
|
||||||
|
}
|
||||||
|
|
||||||
|
export function localPrincipal(): PrincipalContext {
|
||||||
|
const home = process.env.THT_HOME ?? join(homedir(), ".thothii");
|
||||||
|
const identityPath = join(home, "identity.json");
|
||||||
|
mkdirSync(home, { recursive: true, mode: 0o700 });
|
||||||
|
try {
|
||||||
|
const stored = JSON.parse(readFileSync(identityPath, "utf8"));
|
||||||
|
if (stored?.issuer === "local" && typeof stored.subject === "string" && /^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i.test(stored.subject)) {
|
||||||
|
return { issuer: "local", subject: stored.subject, isAdmin: false };
|
||||||
|
}
|
||||||
|
throw new Error("invalid local identity");
|
||||||
|
} catch (error: any) {
|
||||||
|
if (error?.code !== "ENOENT") throw error;
|
||||||
|
const principal = { issuer: "local", subject: randomUUID() };
|
||||||
|
try {
|
||||||
|
writeFileSync(identityPath, JSON.stringify(principal) + "\n", { mode: 0o600, flag: "wx" });
|
||||||
|
return { ...principal, isAdmin: false };
|
||||||
|
} catch (writeError: any) {
|
||||||
|
// Another local request won the identity creation race; always converge on its UUID.
|
||||||
|
if (writeError?.code === "EEXIST") return localPrincipal();
|
||||||
|
throw writeError;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function principalEnvironment(principal: PrincipalContext): NodeJS.ProcessEnv {
|
||||||
|
const env: NodeJS.ProcessEnv = {
|
||||||
|
THT_PRINCIPAL_ISSUER: principal.issuer,
|
||||||
|
THT_PRINCIPAL_SUBJECT: principal.subject,
|
||||||
|
THT_PRINCIPAL_IS_ADMIN: principal.isAdmin ? "true" : "false",
|
||||||
|
};
|
||||||
|
if (principal.displayName) env.THT_PRINCIPAL_DISPLAY_NAME = principal.displayName;
|
||||||
|
return env;
|
||||||
|
}
|
||||||
@@ -5,6 +5,7 @@ import { SessionBridge } from "../bridge/session-bridge.js";
|
|||||||
import type { ThtRunner } from "../tht/tht-runner.js";
|
import type { ThtRunner } from "../tht/tht-runner.js";
|
||||||
import { buildPiChildEnv, canonicalPiProvider } from "./provider-credentials.js";
|
import { buildPiChildEnv, canonicalPiProvider } from "./provider-credentials.js";
|
||||||
import { secretValue } from "../config/secret-bundle.js";
|
import { secretValue } from "../config/secret-bundle.js";
|
||||||
|
import { principalEnvironment, type PrincipalContext } from "../auth/principal.js";
|
||||||
|
|
||||||
export interface SessionRuntime {
|
export interface SessionRuntime {
|
||||||
rpc: RpcClient;
|
rpc: RpcClient;
|
||||||
@@ -19,6 +20,7 @@ export interface RuntimeOptions {
|
|||||||
author?: string;
|
author?: string;
|
||||||
question?: string;
|
question?: string;
|
||||||
mode?: "new" | "resume";
|
mode?: "new" | "resume";
|
||||||
|
principal?: PrincipalContext;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Injectable child-process boundary; callbacks may ignore arguments in simpler tests. */
|
/** Injectable child-process boundary; callbacks may ignore arguments in simpler tests. */
|
||||||
@@ -31,21 +33,21 @@ type SpawnFn = (
|
|||||||
export class PiProcessManager {
|
export class PiProcessManager {
|
||||||
private runtimes = new Map<string, SessionRuntime>();
|
private runtimes = new Map<string, SessionRuntime>();
|
||||||
private spawnFn: (
|
private spawnFn: (
|
||||||
sessionId: string, author: string, provider: string | undefined,
|
sessionId: string, author: string, provider: string | undefined, principal?: PrincipalContext,
|
||||||
) => ChildProcessWithoutNullStreams;
|
) => ChildProcessWithoutNullStreams;
|
||||||
|
|
||||||
constructor(private cfg: AppConfig, opts?: { spawnFn?: SpawnFn }) {
|
constructor(private cfg: AppConfig, opts?: { spawnFn?: SpawnFn }) {
|
||||||
if (opts?.spawnFn) {
|
if (opts?.spawnFn) {
|
||||||
this.spawnFn = (sessionId, author, provider) =>
|
this.spawnFn = (sessionId, author, provider, principal) =>
|
||||||
this.spawnPi(opts.spawnFn!, sessionId, author, provider);
|
this.spawnPi(opts.spawnFn!, sessionId, author, provider, principal);
|
||||||
} else {
|
} else {
|
||||||
this.spawnFn = (sessionId, author, provider) =>
|
this.spawnFn = (sessionId, author, provider, principal) =>
|
||||||
this.spawnPi(nodeSpawn, sessionId, author, provider);
|
this.spawnPi(nodeSpawn, sessionId, author, provider, principal);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
private spawnPi(
|
private spawnPi(
|
||||||
spawnFn: SpawnFn, sessionId: string, author: string, provider: string | undefined,
|
spawnFn: SpawnFn, sessionId: string, author: string, provider: string | undefined, principal?: PrincipalContext,
|
||||||
): ChildProcessWithoutNullStreams {
|
): ChildProcessWithoutNullStreams {
|
||||||
const env = buildPiChildEnv({
|
const env = buildPiChildEnv({
|
||||||
provider,
|
provider,
|
||||||
@@ -53,6 +55,7 @@ export class PiProcessManager {
|
|||||||
credentialFile: this.cfg.modelApiKeyFile,
|
credentialFile: this.cfg.modelApiKeyFile,
|
||||||
additions: { THT_SESSION: sessionId, THT_AUTHOR: author },
|
additions: { THT_SESSION: sessionId, THT_AUTHOR: author },
|
||||||
});
|
});
|
||||||
|
if (principal) Object.assign(env, principalEnvironment(principal));
|
||||||
// The Thoth gate executes the deterministic `tht` CLI as a Pi tool. Give only
|
// The Thoth gate executes the deterministic `tht` CLI as a Pi tool. Give only
|
||||||
// this managed session process the adapter values already loaded by the core
|
// this managed session process the adapter values already loaded by the core
|
||||||
// entrypoint; the generic provider helper continues to scrub them by default.
|
// entrypoint; the generic provider helper continues to scrub them by default.
|
||||||
@@ -95,7 +98,7 @@ export class PiProcessManager {
|
|||||||
}
|
}
|
||||||
const author = o.author ?? "dev@local";
|
const author = o.author ?? "dev@local";
|
||||||
const provider = canonicalPiProvider(o.provider ?? this.cfg.defaults.provider);
|
const provider = canonicalPiProvider(o.provider ?? this.cfg.defaults.provider);
|
||||||
const child = this.spawnFn(sessionId, author, provider);
|
const child = this.spawnFn(sessionId, author, provider, o.principal);
|
||||||
let rt: SessionRuntime | undefined;
|
let rt: SessionRuntime | undefined;
|
||||||
try {
|
try {
|
||||||
const rpc = new RpcClient(child);
|
const rpc = new RpcClient(child);
|
||||||
|
|||||||
+172
-41
@@ -3,7 +3,8 @@ import type { PiProcessManager } from "../pi/pi-process-manager.js";
|
|||||||
import type { ThtRunner } from "../tht/tht-runner.js";
|
import type { ThtRunner } from "../tht/tht-runner.js";
|
||||||
import type { SseHub } from "../sse/sse-hub.js";
|
import type { SseHub } from "../sse/sse-hub.js";
|
||||||
import type { Settings } from "../settings/settings-store.js";
|
import type { Settings } from "../settings/settings-store.js";
|
||||||
import { getUser } from "../auth/auth.js";
|
import { getPrincipal } from "../auth/auth.js";
|
||||||
|
import type { PrincipalContext } from "../auth/principal.js";
|
||||||
import type { ReadinessManager } from "../runtime/readiness-manager.js";
|
import type { ReadinessManager } from "../runtime/readiness-manager.js";
|
||||||
|
|
||||||
const BOOTSTRAP_FAILURE_MESSAGE =
|
const BOOTSTRAP_FAILURE_MESSAGE =
|
||||||
@@ -25,7 +26,11 @@ function eventCursor(...values: unknown[]): number {
|
|||||||
|
|
||||||
export function sessionRoutes(
|
export function sessionRoutes(
|
||||||
app: FastifyInstance,
|
app: FastifyInstance,
|
||||||
d: { mgr: PiProcessManager; tht: ThtRunner; hub: SseHub; getSettings: () => Settings; readiness: ReadinessManager },
|
d: {
|
||||||
|
mgr: PiProcessManager; tht: ThtRunner; hub: SseHub;
|
||||||
|
getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
||||||
|
readiness: ReadinessManager;
|
||||||
|
},
|
||||||
) {
|
) {
|
||||||
const lifecycleTails = new Map<string, Promise<void>>();
|
const lifecycleTails = new Map<string, Promise<void>>();
|
||||||
const boundRuntimes = new Map<
|
const boundRuntimes = new Map<
|
||||||
@@ -54,7 +59,34 @@ export function sessionRoutes(
|
|||||||
const info = (id: string, text: string, level = "info") =>
|
const info = (id: string, text: string, level = "info") =>
|
||||||
d.hub.publish(id, "info", { type: "info", level, text });
|
d.hub.publish(id, "info", { type: "info", level, text });
|
||||||
|
|
||||||
const bindRuntime = (id: string, rt: ReturnType<PiProcessManager["createFor"]>) => {
|
const runnerFor = (principal: PrincipalContext): any => {
|
||||||
|
const runner = d.tht as any;
|
||||||
|
return typeof runner.withPrincipal === "function" ? runner.withPrincipal(principal) : runner;
|
||||||
|
};
|
||||||
|
|
||||||
|
const isNotFound = (error: unknown) =>
|
||||||
|
/not found|non trovata|inesistente|404/i.test(error instanceof Error ? error.message : String(error));
|
||||||
|
|
||||||
|
/** RLS makes a foreign session indistinguishable from a missing one. */
|
||||||
|
const authorize = async (principal: PrincipalContext, id: string, workspace?: string): Promise<any | undefined> => {
|
||||||
|
try {
|
||||||
|
const runner = runnerFor(principal);
|
||||||
|
// Dependency-injected runners in legacy route tests may model only the mutation under
|
||||||
|
// test. Production ThtRunner always exposes sessionShow; keep that test seam harmless.
|
||||||
|
if (typeof runner.sessionShow !== "function") return {};
|
||||||
|
const manifest = await runner.sessionShow(id, workspace);
|
||||||
|
return manifest ?? undefined;
|
||||||
|
} catch (error) {
|
||||||
|
if (isNotFound(error)) return undefined;
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
const storageFailure = (reply: any) => reply.code(503).send({ error: "session storage is unavailable" });
|
||||||
|
|
||||||
|
const bindRuntime = (
|
||||||
|
id: string, rt: ReturnType<PiProcessManager["createFor"]>, runner: any, workspace?: string,
|
||||||
|
) => {
|
||||||
const previous = boundRuntimes.get(id);
|
const previous = boundRuntimes.get(id);
|
||||||
boundRuntimes.set(id, rt);
|
boundRuntimes.set(id, rt);
|
||||||
try {
|
try {
|
||||||
@@ -72,7 +104,7 @@ export function sessionRoutes(
|
|||||||
// old failure must not touch its manifest.
|
// old failure must not touch its manifest.
|
||||||
const bound = boundRuntimes.get(id);
|
const bound = boundRuntimes.get(id);
|
||||||
if (bound !== undefined && bound !== rt) return;
|
if (bound !== undefined && bound !== rt) return;
|
||||||
await d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
|
await runner.failSession(id, workspace).catch(() => undefined);
|
||||||
}).catch(() => undefined);
|
}).catch(() => undefined);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -96,6 +128,8 @@ export function sessionRoutes(
|
|||||||
const bootstrap = (
|
const bootstrap = (
|
||||||
id: string,
|
id: string,
|
||||||
rt: ReturnType<PiProcessManager["createFor"]>,
|
rt: ReturnType<PiProcessManager["createFor"]>,
|
||||||
|
runner: any,
|
||||||
|
workspace: string | undefined,
|
||||||
configure: Promise<void>,
|
configure: Promise<void>,
|
||||||
retrieval: Promise<void> | null,
|
retrieval: Promise<void> | null,
|
||||||
start: () => void,
|
start: () => void,
|
||||||
@@ -117,7 +151,7 @@ export function sessionRoutes(
|
|||||||
if (d.mgr.get(id) !== rt || !d.mgr.teardownIfCurrent(id, rt)) return;
|
if (d.mgr.get(id) !== rt || !d.mgr.teardownIfCurrent(id, rt)) return;
|
||||||
if (!failurePersistenceClaimed.has(rt)) {
|
if (!failurePersistenceClaimed.has(rt)) {
|
||||||
failurePersistenceClaimed.add(rt);
|
failurePersistenceClaimed.add(rt);
|
||||||
await d.tht.failSession(id, d.getSettings().workspace).catch(() => undefined);
|
await runner.failSession(id, workspace).catch(() => undefined);
|
||||||
}
|
}
|
||||||
rt.bridge.emitClientEvent({ type: "info", level: "error", text: BOOTSTRAP_FAILURE_MESSAGE });
|
rt.bridge.emitClientEvent({ type: "info", level: "error", text: BOOTSTRAP_FAILURE_MESSAGE });
|
||||||
rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
|
rt.bridge.emitClientEvent({ type: "system_event", event: "session_failed" });
|
||||||
@@ -127,62 +161,105 @@ export function sessionRoutes(
|
|||||||
})();
|
})();
|
||||||
};
|
};
|
||||||
|
|
||||||
app.post("/runtime/prewarm", async (_req, reply) => {
|
app.post("/runtime/prewarm", async (req, reply) => {
|
||||||
const workspace = d.getSettings().workspace ?? "";
|
let settings: Settings;
|
||||||
void d.readiness.ensure(workspace).catch(() => undefined);
|
try { settings = await d.getSettings(getPrincipal(req)); } catch { return storageFailure(reply); }
|
||||||
|
const workspace = settings.workspace ?? "";
|
||||||
|
void d.readiness.ensure(workspace, getPrincipal(req)).catch(() => undefined);
|
||||||
return reply.code(202).send({ status: "warming" });
|
return reply.code(202).send({ status: "warming" });
|
||||||
});
|
});
|
||||||
|
|
||||||
app.post("/sessions", async (req, reply) => {
|
app.post("/sessions", async (req, reply) => {
|
||||||
const b = req.body as { question: string; name?: string };
|
const b = req.body as { question: string; name?: string };
|
||||||
const s = d.getSettings();
|
const principal = getPrincipal(req);
|
||||||
const ensure = await d.readiness.ensure(s.workspace ?? "");
|
let s: Settings;
|
||||||
|
try { s = await d.getSettings(principal); } catch { return storageFailure(reply); }
|
||||||
|
const runner = runnerFor(principal);
|
||||||
|
const ensure = await d.readiness.ensure(s.workspace ?? "", principal);
|
||||||
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
||||||
// Settings (global) supply workspace/provider/model/thinking. The new-question
|
// Settings (global) supply workspace/provider/model/thinking. The new-question
|
||||||
// form sends only the question text. `workspace` selects the tht `-c <config>`.
|
// form sends only the question text. `workspace` selects the tht `-c <config>`.
|
||||||
const { id } = await d.tht.sessionNew({
|
let id: string;
|
||||||
question: b.question,
|
try {
|
||||||
name: b.name,
|
({ id } = await runner.sessionNew({
|
||||||
workspace: s.workspace,
|
question: b.question, name: b.name, workspace: s.workspace,
|
||||||
provider: s.provider,
|
provider: s.provider, model: s.model, thinking: s.thinking,
|
||||||
model: s.model,
|
}));
|
||||||
thinking: s.thinking,
|
} catch { return storageFailure(reply); }
|
||||||
});
|
|
||||||
const options = {
|
const options = {
|
||||||
provider: s.provider,
|
provider: s.provider,
|
||||||
model: s.model,
|
model: s.model,
|
||||||
thinking: s.thinking,
|
thinking: s.thinking,
|
||||||
author: getUser(req).id,
|
author: principal.displayName ?? principal.subject,
|
||||||
|
principal,
|
||||||
question: b.question,
|
question: b.question,
|
||||||
};
|
};
|
||||||
const rt = d.mgr.createFor(id, options);
|
const rt = d.mgr.createFor(id, options);
|
||||||
bindRuntime(id, rt);
|
bindRuntime(id, rt, runner, s.workspace);
|
||||||
info(id, "Session created");
|
info(id, "Session created");
|
||||||
bootstrap(
|
bootstrap(
|
||||||
id, rt, d.mgr.configure(rt, options),
|
id, rt, runner, s.workspace, d.mgr.configure(rt, options),
|
||||||
d.tht.searchPack(b.question, id, s.workspace),
|
runner.searchPack(b.question, id, s.workspace),
|
||||||
() => d.mgr.start(id, rt, options),
|
() => d.mgr.start(id, rt, options),
|
||||||
);
|
);
|
||||||
return { id };
|
return { id };
|
||||||
});
|
});
|
||||||
app.get("/sessions", async () => d.tht.sessionList(d.getSettings().workspace));
|
app.get("/sessions", async (req, reply) => {
|
||||||
app.get("/sessions/:id", async (req) => d.tht.sessionShow((req.params as any).id, d.getSettings().workspace));
|
const principal = getPrincipal(req);
|
||||||
|
const scope = (req.query as { scope?: string }).scope ?? "mine";
|
||||||
|
if (scope !== "mine" && scope !== "all") return reply.code(400).send({ error: "scope must be mine or all" });
|
||||||
|
if (scope === "all" && !principal.isAdmin) return reply.code(403).send({ error: "admin scope required" });
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
// Admin RLS is deliberately disabled for a normal 'mine' listing.
|
||||||
|
const scopedPrincipal = scope === "mine" ? { ...principal, isAdmin: false } : principal;
|
||||||
|
return await runnerFor(scopedPrincipal).sessionList(settings.workspace);
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
|
});
|
||||||
|
app.get("/sessions/:id", async (req, reply) => {
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
const manifest = await authorize(principal, (req.params as any).id, settings.workspace);
|
||||||
|
return manifest ?? reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
|
});
|
||||||
app.post("/sessions/:id/response", async (req, reply) => {
|
app.post("/sessions/:id/response", async (req, reply) => {
|
||||||
const id = (req.params as any).id;
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
const rt = d.mgr.get(id);
|
const rt = d.mgr.get(id);
|
||||||
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
|
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
|
||||||
rt.bridge.respond((req.body as any).ui_response);
|
rt.bridge.respond((req.body as any).ui_response);
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/steer", async (req, reply) => {
|
app.post("/sessions/:id/steer", async (req, reply) => {
|
||||||
const rt = d.mgr.get((req.params as any).id);
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
|
const rt = d.mgr.get(id);
|
||||||
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
|
if (!rt) return reply.code(404).send({ error: "sessione non attiva" });
|
||||||
rt.bridge.steer((req.body as any).text);
|
rt.bridge.steer((req.body as any).text);
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/resume", async (req, reply) => {
|
app.post("/sessions/:id/resume", async (req, reply) => {
|
||||||
const id = (req.params as any).id;
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
return withSessionLifecycle(id, async () => {
|
return withSessionLifecycle(id, async () => {
|
||||||
|
let settings: Settings;
|
||||||
|
let manifest: any;
|
||||||
|
try {
|
||||||
|
settings = await d.getSettings(principal);
|
||||||
|
manifest = await authorize(principal, id, settings.workspace);
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
|
if (!manifest) return reply.code(404).send({ error: "session not found" });
|
||||||
|
const runner = runnerFor(principal);
|
||||||
// This check belongs inside the per-session lock: a preceding cold Resume may have
|
// This check belongs inside the per-session lock: a preceding cold Resume may have
|
||||||
// installed a running runtime while this request was waiting.
|
// installed a running runtime while this request was waiting.
|
||||||
const existing = d.mgr.get(id);
|
const existing = d.mgr.get(id);
|
||||||
@@ -192,26 +269,25 @@ export function sessionRoutes(
|
|||||||
return reply.code(200).send({ id, alreadyActive: true });
|
return reply.code(200).send({ id, alreadyActive: true });
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
const manifest = (await d.tht.sessionShow(id, d.getSettings().workspace)) as { status?: string; archived?: boolean } | null;
|
|
||||||
if (manifest?.status === "finalized" || manifest?.archived) {
|
if (manifest?.status === "finalized" || manifest?.archived) {
|
||||||
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
return reply.code(409).send({ error: "sessione in sola lettura (finalizzata o archiviata)" });
|
||||||
}
|
}
|
||||||
const settings = d.getSettings();
|
const ensure = await d.readiness.ensure(settings.workspace ?? "", principal);
|
||||||
const ensure = await d.readiness.ensure(settings.workspace ?? "");
|
|
||||||
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE });
|
||||||
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
|
const saved = manifest as { provider?: string; model?: string; thinking?: string } | null;
|
||||||
const options = {
|
const options = {
|
||||||
provider: saved?.provider,
|
provider: saved?.provider,
|
||||||
model: saved?.model,
|
model: saved?.model,
|
||||||
thinking: saved?.thinking ?? settings.thinking,
|
thinking: saved?.thinking ?? settings.thinking,
|
||||||
author: getUser(req).id,
|
author: principal.displayName ?? principal.subject,
|
||||||
|
principal,
|
||||||
mode: "resume" as const,
|
mode: "resume" as const,
|
||||||
};
|
};
|
||||||
|
|
||||||
// Reopening is validation, not the transport commit point. Keep the old hub intact if
|
// Reopening is validation, not the transport commit point. Keep the old hub intact if
|
||||||
// persistence cannot be reopened.
|
// persistence cannot be reopened.
|
||||||
try {
|
try {
|
||||||
await d.tht.reopenSession(id, settings.workspace);
|
await runner.reopenSession(id, settings.workspace);
|
||||||
} catch {
|
} catch {
|
||||||
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
return reply.code(503).send({ error: RESUME_FAILURE_MESSAGE });
|
||||||
}
|
}
|
||||||
@@ -233,7 +309,7 @@ export function sessionRoutes(
|
|||||||
d.mgr.teardownIfCurrent(id, current);
|
d.mgr.teardownIfCurrent(id, current);
|
||||||
}
|
}
|
||||||
rt = d.mgr.createFor(id, options);
|
rt = d.mgr.createFor(id, options);
|
||||||
bindRuntime(id, rt);
|
bindRuntime(id, rt, runner, settings.workspace);
|
||||||
} catch {
|
} catch {
|
||||||
// A created-but-unbound runtime is not usable. The old hub remains attached because
|
// A created-but-unbound runtime is not usable. The old hub remains attached because
|
||||||
// clear() has not happened yet.
|
// clear() has not happened yet.
|
||||||
@@ -248,28 +324,41 @@ export function sessionRoutes(
|
|||||||
// immediately before the first event produced by the new Resume.
|
// immediately before the first event produced by the new Resume.
|
||||||
d.hub.clear(id);
|
d.hub.clear(id);
|
||||||
info(id, "Resuming session");
|
info(id, "Resuming session");
|
||||||
bootstrap(id, rt, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
bootstrap(id, rt, runner, settings.workspace, d.mgr.configure(rt, options), null, () => d.mgr.start(id, rt, options));
|
||||||
return reply.code(200).send({ id, alreadyActive: false });
|
return reply.code(200).send({ id, alreadyActive: false });
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/close", async (req) => {
|
app.post("/sessions/:id/close", async (req, reply) => {
|
||||||
const id = (req.params as { id: string }).id;
|
const id = (req.params as { id: string }).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
return withSessionLifecycle(id, async () => {
|
return withSessionLifecycle(id, async () => {
|
||||||
|
let settings: Settings;
|
||||||
|
try {
|
||||||
|
settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
// Invalidate the live generation before persistence can yield. Otherwise its deferred
|
// Invalidate the live generation before persistence can yield. Otherwise its deferred
|
||||||
// bootstrap may start Pi while Close is already in progress.
|
// bootstrap may start Pi while Close is already in progress.
|
||||||
const current = d.mgr.get(id);
|
const current = d.mgr.get(id);
|
||||||
boundRuntimes.delete(id);
|
boundRuntimes.delete(id);
|
||||||
if (current) d.mgr.teardownIfCurrent(id, current);
|
if (current) d.mgr.teardownIfCurrent(id, current);
|
||||||
try {
|
try {
|
||||||
await d.tht.closeSession(id, d.getSettings().workspace);
|
await runnerFor(principal).closeSession(id, settings.workspace);
|
||||||
|
} catch {
|
||||||
|
return storageFailure(reply);
|
||||||
} finally {
|
} finally {
|
||||||
d.hub.clear(id);
|
d.hub.clear(id);
|
||||||
}
|
}
|
||||||
return { closed: true };
|
return { closed: true };
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
app.get("/sessions/:id/events", (req, reply) => {
|
app.get("/sessions/:id/events", async (req, reply) => {
|
||||||
const id = (req.params as any).id;
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
const rt = d.mgr.get(id);
|
const rt = d.mgr.get(id);
|
||||||
const afterId = eventCursor(
|
const afterId = eventCursor(
|
||||||
req.headers["last-event-id"],
|
req.headers["last-event-id"],
|
||||||
@@ -303,31 +392,73 @@ export function sessionRoutes(
|
|||||||
req.raw.on("close", off);
|
req.raw.on("close", off);
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/rename", async (req, reply) => {
|
app.post("/sessions/:id/rename", async (req, reply) => {
|
||||||
await d.tht.setName((req.params as any).id, (req.body as any).name);
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
await runnerFor(principal).setName(id, (req.body as any).name);
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/group", async (req, reply) => {
|
app.post("/sessions/:id/group", async (req, reply) => {
|
||||||
await d.tht.setGroup((req.params as any).id, (req.body as any).group);
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
await runnerFor(principal).setGroup(id, (req.body as any).group);
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/archive", async (req, reply) => {
|
app.post("/sessions/:id/archive", async (req, reply) => {
|
||||||
await d.tht.archive((req.params as any).id);
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
await runnerFor(principal).archive(id);
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.post("/sessions/:id/unarchive", async (req, reply) => {
|
app.post("/sessions/:id/unarchive", async (req, reply) => {
|
||||||
await d.tht.unarchive((req.params as any).id);
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
await runnerFor(principal).unarchive(id);
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
app.delete("/sessions/:id", async (req, reply) => {
|
app.delete("/sessions/:id", async (req, reply) => {
|
||||||
const id = (req.params as any).id;
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
return withSessionLifecycle(id, async () => {
|
return withSessionLifecycle(id, async () => {
|
||||||
|
let settings: Settings;
|
||||||
|
try {
|
||||||
|
settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
const current = d.mgr.get(id);
|
const current = d.mgr.get(id);
|
||||||
boundRuntimes.delete(id);
|
boundRuntimes.delete(id);
|
||||||
if (current) d.mgr.teardownIfCurrent(id, current);
|
if (current) d.mgr.teardownIfCurrent(id, current);
|
||||||
await d.tht.deleteSession(id, d.getSettings().workspace);
|
try {
|
||||||
|
await runnerFor(principal).deleteSession(id, settings.workspace);
|
||||||
|
} catch {
|
||||||
|
return storageFailure(reply);
|
||||||
|
}
|
||||||
d.hub.forget(id);
|
d.hub.forget(id);
|
||||||
return reply.code(204).send();
|
return reply.code(204).send();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
app.get("/sessions/:id/documents", async (req) => d.tht.documents((req.params as any).id));
|
app.get("/sessions/:id/documents", async (req, reply) => {
|
||||||
|
const id = (req.params as any).id;
|
||||||
|
const principal = getPrincipal(req);
|
||||||
|
try {
|
||||||
|
const settings = await d.getSettings(principal);
|
||||||
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
return await runnerFor(principal).documents(id);
|
||||||
|
} catch { return storageFailure(reply); }
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -2,6 +2,8 @@ import type { FastifyInstance } from "fastify";
|
|||||||
import type { AppConfig } from "../config.js";
|
import type { AppConfig } from "../config.js";
|
||||||
import { loadSettings, saveSettings, type Settings } from "../settings/settings-store.js";
|
import { loadSettings, saveSettings, type Settings } from "../settings/settings-store.js";
|
||||||
import { listWorkspaces, type ListModelsFn } from "./meta.js";
|
import { listWorkspaces, type ListModelsFn } from "./meta.js";
|
||||||
|
import { getPrincipal } from "../auth/auth.js";
|
||||||
|
import type { PrincipalContext } from "../auth/principal.js";
|
||||||
|
|
||||||
/** Merge stored settings over env/first-workspace defaults. */
|
/** Merge stored settings over env/first-workspace defaults. */
|
||||||
export function effectiveSettings(cfg: AppConfig, stored: Settings): Settings {
|
export function effectiveSettings(cfg: AppConfig, stored: Settings): Settings {
|
||||||
@@ -16,10 +18,18 @@ export function effectiveSettings(cfg: AppConfig, stored: Settings): Settings {
|
|||||||
|
|
||||||
export function settingsRoutes(
|
export function settingsRoutes(
|
||||||
app: FastifyInstance,
|
app: FastifyInstance,
|
||||||
deps: { cfg: AppConfig; listModels: ListModelsFn },
|
deps: {
|
||||||
|
cfg: AppConfig; listModels: ListModelsFn;
|
||||||
|
getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
||||||
|
saveSettings: (principal: PrincipalContext, settings: Settings) => Promise<void>;
|
||||||
|
},
|
||||||
): void {
|
): void {
|
||||||
app.get("/settings", async () => {
|
app.get("/settings", async (req, reply) => {
|
||||||
return effectiveSettings(deps.cfg, loadSettings(deps.cfg));
|
try {
|
||||||
|
return await deps.getSettings(getPrincipal(req));
|
||||||
|
} catch {
|
||||||
|
return reply.code(503).send({ error: "settings storage is unavailable" });
|
||||||
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
app.put("/settings", async (req, reply) => {
|
app.put("/settings", async (req, reply) => {
|
||||||
@@ -46,7 +56,11 @@ export function settingsRoutes(
|
|||||||
model: b.model,
|
model: b.model,
|
||||||
thinking: b.thinking,
|
thinking: b.thinking,
|
||||||
};
|
};
|
||||||
saveSettings(deps.cfg, next);
|
try {
|
||||||
|
await deps.saveSettings(getPrincipal(req), next);
|
||||||
return effectiveSettings(deps.cfg, next);
|
return effectiveSettings(deps.cfg, next);
|
||||||
|
} catch {
|
||||||
|
return reply.code(503).send({ error: "settings storage is unavailable" });
|
||||||
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,25 +1,63 @@
|
|||||||
import type { FastifyInstance } from "fastify";
|
import type { FastifyInstance } from "fastify";
|
||||||
import type { ThtRunner } from "../tht/tht-runner.js";
|
import type { ThtRunner } from "../tht/tht-runner.js";
|
||||||
|
import { getPrincipal } from "../auth/auth.js";
|
||||||
|
import type { PrincipalContext } from "../auth/principal.js";
|
||||||
|
import type { Settings } from "../settings/settings-store.js";
|
||||||
|
|
||||||
export function sqlRoutes(app: FastifyInstance, deps: { tht: ThtRunner }): void {
|
export function sqlRoutes(app: FastifyInstance, deps: {
|
||||||
|
tht: ThtRunner; getSettings: (principal: PrincipalContext) => Promise<Settings>;
|
||||||
|
}): void {
|
||||||
|
const runnerFor = (principal: PrincipalContext): any => {
|
||||||
|
const runner = deps.tht as any;
|
||||||
|
return typeof runner.withPrincipal === "function" ? runner.withPrincipal(principal) : runner;
|
||||||
|
};
|
||||||
|
const authorize = async (principal: PrincipalContext, id: string, workspace?: string) => {
|
||||||
|
try {
|
||||||
|
const runner = runnerFor(principal);
|
||||||
|
if (typeof runner.sessionShow !== "function") return {};
|
||||||
|
return await runner.sessionShow(id, workspace);
|
||||||
|
}
|
||||||
|
catch (error) {
|
||||||
|
if (/not found|non trovata|inesistente|404/i.test(error instanceof Error ? error.message : String(error))) return undefined;
|
||||||
|
throw error;
|
||||||
|
}
|
||||||
|
};
|
||||||
app.post("/sessions/:id/sql/preview", async (req, reply) => {
|
app.post("/sessions/:id/sql/preview", async (req, reply) => {
|
||||||
const id = (req.params as any).id as string;
|
const id = (req.params as any).id as string;
|
||||||
const { limit, offset } = (req.body as any) ?? {};
|
const { limit, offset } = (req.body as any) ?? {};
|
||||||
|
let principal: PrincipalContext;
|
||||||
|
let workspace: string | undefined;
|
||||||
try {
|
try {
|
||||||
const result = await deps.tht.sqlPreview(id, { limit, offset });
|
principal = getPrincipal(req);
|
||||||
return result;
|
const settings = await deps.getSettings(principal);
|
||||||
} catch (err: any) {
|
workspace = settings.workspace;
|
||||||
return reply.code(500).send({ error: err.message ?? String(err) });
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch {
|
||||||
|
return reply.code(503).send({ error: "session storage is unavailable" });
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
return await runnerFor(principal).sqlPreview(id, { limit, offset }, workspace);
|
||||||
|
} catch (error: any) {
|
||||||
|
return reply.code(500).send({ error: error.message ?? String(error) });
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
app.post("/sessions/:id/sql/export", async (req, reply) => {
|
app.post("/sessions/:id/sql/export", async (req, reply) => {
|
||||||
const id = (req.params as any).id as string;
|
const id = (req.params as any).id as string;
|
||||||
|
let principal: PrincipalContext;
|
||||||
|
let workspace: string | undefined;
|
||||||
try {
|
try {
|
||||||
const result = await deps.tht.sqlExport(id);
|
principal = getPrincipal(req);
|
||||||
return result;
|
const settings = await deps.getSettings(principal);
|
||||||
} catch (err: any) {
|
workspace = settings.workspace;
|
||||||
return reply.code(500).send({ error: err.message ?? String(err) });
|
if (!await authorize(principal, id, settings.workspace)) return reply.code(404).send({ error: "session not found" });
|
||||||
|
} catch {
|
||||||
|
return reply.code(503).send({ error: "session storage is unavailable" });
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
return await runnerFor(principal).sqlExport(id, workspace);
|
||||||
|
} catch (error: any) {
|
||||||
|
return reply.code(500).send({ error: error.message ?? String(error) });
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import type { OllamaEnsureResult, ThtRunner } from "../tht/tht-runner.js";
|
import type { OllamaEnsureResult, ThtRunner } from "../tht/tht-runner.js";
|
||||||
|
import type { PrincipalContext } from "../auth/principal.js";
|
||||||
|
|
||||||
interface ReadyEntry {
|
interface ReadyEntry {
|
||||||
expiresAt: number;
|
expiresAt: number;
|
||||||
@@ -20,25 +21,28 @@ export class ReadinessManager {
|
|||||||
private now: () => number = Date.now,
|
private now: () => number = Date.now,
|
||||||
) {}
|
) {}
|
||||||
|
|
||||||
ensure(workspace = ""): Promise<OllamaEnsureResult> {
|
ensure(workspace = "", principal?: PrincipalContext): Promise<OllamaEnsureResult> {
|
||||||
const cached = this.ready.get(workspace);
|
const key = `${principal?.issuer ?? ""}\0${principal?.subject ?? ""}\0${workspace}`;
|
||||||
|
const cached = this.ready.get(key);
|
||||||
if (cached && cached.expiresAt > this.now()) return Promise.resolve(cached.result);
|
if (cached && cached.expiresAt > this.now()) return Promise.resolve(cached.result);
|
||||||
if (cached) this.ready.delete(workspace);
|
if (cached) this.ready.delete(key);
|
||||||
|
|
||||||
const current = this.inFlight.get(workspace);
|
const current = this.inFlight.get(key);
|
||||||
if (current) return current;
|
if (current) return current;
|
||||||
|
|
||||||
const pending = this.tht.ollamaEnsure(workspace, this.timeoutSec)
|
const runner = principal && typeof (this.tht as any).withPrincipal === "function"
|
||||||
|
? this.tht.withPrincipal(principal) : this.tht;
|
||||||
|
const pending = runner.ollamaEnsure(workspace, this.timeoutSec)
|
||||||
.then((result) => {
|
.then((result) => {
|
||||||
if (result.ok) {
|
if (result.ok) {
|
||||||
this.ready.set(workspace, { result, expiresAt: this.now() + this.ttlMs });
|
this.ready.set(key, { result, expiresAt: this.now() + this.ttlMs });
|
||||||
}
|
}
|
||||||
return result;
|
return result;
|
||||||
})
|
})
|
||||||
.finally(() => {
|
.finally(() => {
|
||||||
if (this.inFlight.get(workspace) === pending) this.inFlight.delete(workspace);
|
if (this.inFlight.get(key) === pending) this.inFlight.delete(key);
|
||||||
});
|
});
|
||||||
this.inFlight.set(workspace, pending);
|
this.inFlight.set(key, pending);
|
||||||
return pending;
|
return pending;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import { spawn } from "node:child_process";
|
import { spawn } from "node:child_process";
|
||||||
import { existsSync } from "node:fs";
|
import { existsSync } from "node:fs";
|
||||||
import { join } from "node:path";
|
import { join } from "node:path";
|
||||||
|
import { principalEnvironment, type PrincipalContext } from "../auth/principal.js";
|
||||||
|
|
||||||
export interface ThtConfig {
|
export interface ThtConfig {
|
||||||
thtBin: string;
|
thtBin: string;
|
||||||
@@ -37,7 +38,10 @@ export interface OllamaEnsureResult {
|
|||||||
}
|
}
|
||||||
|
|
||||||
export class ThtRunner {
|
export class ThtRunner {
|
||||||
constructor(private cfg: ThtConfig) {}
|
constructor(private cfg: ThtConfig, private principal?: PrincipalContext) {}
|
||||||
|
|
||||||
|
/** Bind one trusted request principal to every child spawned by this runner. */
|
||||||
|
withPrincipal(principal: PrincipalContext): ThtRunner { return new ThtRunner(this.cfg, principal); }
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Resolve the `-c <config>` args. If `workspace` is given AND a matching
|
* Resolve the `-c <config>` args. If `workspace` is given AND a matching
|
||||||
@@ -65,15 +69,23 @@ export class ThtRunner {
|
|||||||
const env: NodeJS.ProcessEnv = { ...process.env };
|
const env: NodeJS.ProcessEnv = { ...process.env };
|
||||||
delete env.THT_DATA_ROOT;
|
delete env.THT_DATA_ROOT;
|
||||||
if (this.cfg.dataRoot !== undefined) env.THT_DATA_ROOT = this.cfg.dataRoot;
|
if (this.cfg.dataRoot !== undefined) env.THT_DATA_ROOT = this.cfg.dataRoot;
|
||||||
|
if (this.principal) Object.assign(env, principalEnvironment(this.principal));
|
||||||
const ch = spawn(this.cfg.thtBin, this.buildArgv(args, workspace), {
|
const ch = spawn(this.cfg.thtBin, this.buildArgv(args, workspace), {
|
||||||
cwd: this.cfg.harnessDir,
|
cwd: this.cfg.harnessDir,
|
||||||
env,
|
env,
|
||||||
});
|
});
|
||||||
let stdout = "";
|
let stdout = "";
|
||||||
let stderr = "";
|
let stderr = "";
|
||||||
|
let settled = false;
|
||||||
|
const finish = (result: { code: number; stdout: string; stderr: string }) => {
|
||||||
|
if (settled) return;
|
||||||
|
settled = true;
|
||||||
|
resolve(result);
|
||||||
|
};
|
||||||
ch.stdout.on("data", (d: Buffer) => (stdout += d));
|
ch.stdout.on("data", (d: Buffer) => (stdout += d));
|
||||||
ch.stderr.on("data", (d: Buffer) => (stderr += d));
|
ch.stderr.on("data", (d: Buffer) => (stderr += d));
|
||||||
ch.on("close", (code) => resolve({ code: code ?? 0, stdout, stderr }));
|
ch.on("error", (error) => finish({ code: 1, stdout, stderr: stderr || error.message }));
|
||||||
|
ch.on("close", (code) => finish({ code: code ?? 0, stdout, stderr }));
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -124,7 +136,7 @@ export class ThtRunner {
|
|||||||
return this.json<unknown>(["session", "show", id, "--json"], workspace);
|
return this.json<unknown>(["session", "show", id, "--json"], workspace);
|
||||||
}
|
}
|
||||||
|
|
||||||
sqlPreview(id: string, p: { limit?: number; offset?: number }) {
|
sqlPreview(id: string, p: { limit?: number; offset?: number }, workspace?: string) {
|
||||||
// No positional FILE: the harness resolves sql_final.sql from the session
|
// No positional FILE: the harness resolves sql_final.sql from the session
|
||||||
// via _session_sql_file(cfg, session_id), which respects the workspace path.
|
// via _session_sql_file(cfg, session_id), which respects the workspace path.
|
||||||
const a = ["sql", "preview", "--session", id, "--json"];
|
const a = ["sql", "preview", "--session", id, "--json"];
|
||||||
@@ -135,11 +147,11 @@ export class ThtRunner {
|
|||||||
rows: unknown[][];
|
rows: unknown[][];
|
||||||
execution_ms: number;
|
execution_ms: number;
|
||||||
truncated: boolean;
|
truncated: boolean;
|
||||||
}>(a);
|
}>(a, workspace);
|
||||||
}
|
}
|
||||||
|
|
||||||
async sqlExport(id: string) {
|
async sqlExport(id: string, workspace?: string) {
|
||||||
const { code, stdout, stderr } = await this.run(["sql", "export", "--session", id]);
|
const { code, stdout, stderr } = await this.run(["sql", "export", "--session", id], workspace);
|
||||||
if (code !== 0) throw new Error(`tht sql export exit ${code}: ${stderr.trim()}`);
|
if (code !== 0) throw new Error(`tht sql export exit ${code}: ${stderr.trim()}`);
|
||||||
return { path: stdout.trim() };
|
return { path: stdout.trim() };
|
||||||
}
|
}
|
||||||
@@ -157,6 +169,11 @@ export class ThtRunner {
|
|||||||
}
|
}
|
||||||
documents(id: string) { return this.json<SessionDocument[]>(["session", "documents", id, "--json"]); }
|
documents(id: string) { return this.json<SessionDocument[]>(["session", "documents", id, "--json"]); }
|
||||||
|
|
||||||
|
preferencesGet(workspace?: string) { return this.json<Record<string, unknown>>(["session", "preferences", "get", "--json"], workspace); }
|
||||||
|
async preferencesSet(preferences: Record<string, unknown>, workspace?: string): Promise<void> {
|
||||||
|
await this.ok(["session", "preferences", "set", "--json", JSON.stringify(preferences)], workspace);
|
||||||
|
}
|
||||||
|
|
||||||
async ollamaEnsure(workspace: string, timeoutSec: number): Promise<OllamaEnsureResult> {
|
async ollamaEnsure(workspace: string, timeoutSec: number): Promise<OllamaEnsureResult> {
|
||||||
const { code, stdout, stderr } = await this.run(
|
const { code, stdout, stderr } = await this.run(
|
||||||
["ollama", "ensure", "--json", "--timeout", String(timeoutSec)],
|
["ollama", "ensure", "--json", "--timeout", String(timeoutSec)],
|
||||||
|
|||||||
+22
-12
@@ -1,38 +1,48 @@
|
|||||||
import { test, expect } from "vitest";
|
import { test, expect } from "vitest";
|
||||||
import Fastify from "fastify";
|
import Fastify from "fastify";
|
||||||
import { authPreHandler, getUser } from "../src/auth/auth.js";
|
import { authPreHandler, getPrincipal } from "../src/auth/auth.js";
|
||||||
|
|
||||||
test("mode none assegna dev@local", async () => {
|
test("local mode resolves a stable local principal", async () => {
|
||||||
const app = Fastify();
|
const app = Fastify();
|
||||||
app.addHook("preHandler", authPreHandler("none"));
|
app.addHook("preHandler", authPreHandler("none"));
|
||||||
app.get("/me", async (req) => getUser(req));
|
app.get("/me", async (req) => getPrincipal(req));
|
||||||
expect((await app.inject({ method: "GET", url: "/me" })).json()).toEqual({
|
expect((await app.inject({ method: "GET", url: "/me" })).json()).toMatchObject({
|
||||||
id: "dev@local",
|
issuer: "local",
|
||||||
|
subject: expect.any(String),
|
||||||
|
isAdmin: false,
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
test("mode mock legge l'header", async () => {
|
test("mock mode makes a principal from the test header", async () => {
|
||||||
const app = Fastify();
|
const app = Fastify();
|
||||||
app.addHook("preHandler", authPreHandler("mock"));
|
app.addHook("preHandler", authPreHandler("mock"));
|
||||||
app.get("/me", async (req) => getUser(req));
|
app.get("/me", async (req) => getPrincipal(req));
|
||||||
const res = await app.inject({
|
const res = await app.inject({
|
||||||
method: "GET",
|
method: "GET",
|
||||||
url: "/me",
|
url: "/me",
|
||||||
headers: { "x-mock-user": "alice" },
|
headers: { "x-mock-user": "alice" },
|
||||||
});
|
});
|
||||||
expect(res.json()).toEqual({ id: "alice" });
|
expect(res.json()).toEqual({ issuer: "mock", subject: "alice", displayName: "alice", isAdmin: false });
|
||||||
});
|
});
|
||||||
|
|
||||||
test("upstream mode requires the authenticated proxy identity header", async () => {
|
test("upstream mode accepts only normalized proxy principal headers", async () => {
|
||||||
const app = Fastify();
|
const app = Fastify();
|
||||||
app.addHook("preHandler", authPreHandler("upstream"));
|
app.addHook("preHandler", authPreHandler("upstream"));
|
||||||
app.get("/me", async (req) => getUser(req));
|
app.get("/me", async (req) => getPrincipal(req));
|
||||||
|
|
||||||
expect((await app.inject({ method: "GET", url: "/me" })).statusCode).toBe(401);
|
expect((await app.inject({ method: "GET", url: "/me" })).statusCode).toBe(401);
|
||||||
const authenticated = await app.inject({
|
const authenticated = await app.inject({
|
||||||
method: "GET",
|
method: "GET",
|
||||||
url: "/me",
|
url: "/me",
|
||||||
headers: { "x-authenticated-user": "alice@example.test" },
|
headers: {
|
||||||
|
"x-thoth-principal-issuer": "portal",
|
||||||
|
"x-thoth-principal-subject": "42",
|
||||||
|
"x-thoth-principal-display-name": "Alice",
|
||||||
|
"x-thoth-is-admin": "1",
|
||||||
|
"x-authenticated-user": "must-not-be-used",
|
||||||
|
},
|
||||||
|
});
|
||||||
|
expect(authenticated.json()).toEqual({
|
||||||
|
issuer: "portal", subject: "42", displayName: "Alice", isAdmin: true,
|
||||||
});
|
});
|
||||||
expect(authenticated.json()).toEqual({ id: "alice@example.test" });
|
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -22,7 +22,9 @@ test("GET /health remains available to container probes in upstream auth mode",
|
|||||||
});
|
});
|
||||||
|
|
||||||
test("SSE response headers are flushed before the first event", async () => {
|
test("SSE response headers are flushed before the first event", async () => {
|
||||||
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }));
|
const app = buildApp(loadConfig({ THT_HARNESS_DIR: "/tmp/h" }), {
|
||||||
|
thtRunner: { sessionShow: async () => ({ id: "header-probe" }) } as any,
|
||||||
|
});
|
||||||
await app.listen({ port: 0, host: "127.0.0.1" });
|
await app.listen({ port: 0, host: "127.0.0.1" });
|
||||||
const port = (app.server.address() as { port: number }).port;
|
const port = (app.server.address() as { port: number }).port;
|
||||||
const controller = new AbortController();
|
const controller = new AbortController();
|
||||||
|
|||||||
@@ -28,6 +28,91 @@ function deferred<T = void>() {
|
|||||||
return { promise, resolve, reject };
|
return { promise, resolve, reject };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const aliceHeaders = {
|
||||||
|
"x-thoth-principal-issuer": "portal",
|
||||||
|
"x-thoth-principal-subject": "alice",
|
||||||
|
"x-thoth-principal-display-name": "Alice",
|
||||||
|
"x-thoth-is-admin": "0",
|
||||||
|
};
|
||||||
|
|
||||||
|
test("upstream requests without a principal fail before a Pi runtime can be created", async () => {
|
||||||
|
let created = false;
|
||||||
|
const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), {
|
||||||
|
mgr: { createFor: () => { created = true; throw new Error("must not spawn"); } } as any,
|
||||||
|
thtRunner: {} as any,
|
||||||
|
});
|
||||||
|
|
||||||
|
const response = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
||||||
|
|
||||||
|
expect(response.statusCode).toBe(401);
|
||||||
|
expect(created).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("session routes conceal foreign or missing sessions and deny SSE before it subscribes", async () => {
|
||||||
|
let subscribed = false;
|
||||||
|
const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), {
|
||||||
|
thtRunner: {
|
||||||
|
withPrincipal: () => ({ sessionShow: async () => null }),
|
||||||
|
} as any,
|
||||||
|
hub: { subscribe: () => { subscribed = true; return () => {}; } } as any,
|
||||||
|
});
|
||||||
|
|
||||||
|
const document = await app.inject({ method: "GET", url: "/sessions/foreign/documents", headers: aliceHeaders });
|
||||||
|
const events = await app.inject({ method: "GET", url: "/sessions/foreign/events", headers: aliceHeaders });
|
||||||
|
|
||||||
|
expect(document.statusCode).toBe(404);
|
||||||
|
expect(events.statusCode).toBe(404);
|
||||||
|
expect(subscribed).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("session listing permits all scope only to admins", async () => {
|
||||||
|
const seen: boolean[] = [];
|
||||||
|
const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), {
|
||||||
|
thtRunner: {
|
||||||
|
withPrincipal: (principal: any) => ({
|
||||||
|
sessionList: async () => { seen.push(principal.isAdmin); return [{ id: "s1" }]; },
|
||||||
|
}),
|
||||||
|
} as any,
|
||||||
|
});
|
||||||
|
const regularAll = await app.inject({ method: "GET", url: "/sessions?scope=all", headers: aliceHeaders });
|
||||||
|
const mine = await app.inject({ method: "GET", url: "/sessions?scope=mine", headers: aliceHeaders });
|
||||||
|
const adminAll = await app.inject({
|
||||||
|
method: "GET", url: "/sessions?scope=all",
|
||||||
|
headers: { ...aliceHeaders, "x-thoth-is-admin": "1" },
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(regularAll.statusCode).toBe(403);
|
||||||
|
expect(mine.statusCode).toBe(200);
|
||||||
|
expect(adminAll.statusCode).toBe(200);
|
||||||
|
expect(seen).toEqual([false, true]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("new sessions are created through the authenticated principal, not a client owner field", async () => {
|
||||||
|
let principal: any;
|
||||||
|
const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), {
|
||||||
|
thtRunner: {
|
||||||
|
withPrincipal: (p: any) => {
|
||||||
|
principal = p;
|
||||||
|
return { sessionNew: async () => ({ id: "owned" }), searchPack: async () => {} };
|
||||||
|
},
|
||||||
|
} as any,
|
||||||
|
readiness: { ensure: async () => ({ ok: true }) } as any,
|
||||||
|
mgr: {
|
||||||
|
createFor: () => ({ bridge: { onClientEvent: () => {} } }),
|
||||||
|
configure: async () => {}, start: () => {}, get: () => undefined,
|
||||||
|
} as any,
|
||||||
|
getSettings: () => ({ workspace: "w" }) as any,
|
||||||
|
});
|
||||||
|
|
||||||
|
const response = await app.inject({
|
||||||
|
method: "POST", url: "/sessions", headers: aliceHeaders,
|
||||||
|
payload: { question: "q", owner: "mallory" },
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(response.statusCode).toBe(200);
|
||||||
|
expect(principal).toMatchObject({ issuer: "portal", subject: "alice" });
|
||||||
|
});
|
||||||
|
|
||||||
test("POST /sessions usa i settings (workspace/provider/model/thinking) e crea+avvia", async () => {
|
test("POST /sessions usa i settings (workspace/provider/model/thinking) e crea+avvia", async () => {
|
||||||
const modelKey = path.join(os.tmpdir(), `thoth-model-key-${process.pid}`);
|
const modelKey = path.join(os.tmpdir(), `thoth-model-key-${process.pid}`);
|
||||||
writeFileSync(modelKey, "test-model-key", { mode: 0o600 });
|
writeFileSync(modelKey, "test-model-key", { mode: 0o600 });
|
||||||
@@ -455,7 +540,9 @@ test("concurrent cold Resume requests serialize and create one runtime", async (
|
|||||||
expect(firstResponse.json()).toEqual({ id: "s1", alreadyActive: false });
|
expect(firstResponse.json()).toEqual({ id: "s1", alreadyActive: false });
|
||||||
expect(secondResponse.json()).toEqual({ id: "s1", alreadyActive: true });
|
expect(secondResponse.json()).toEqual({ id: "s1", alreadyActive: true });
|
||||||
expect({ showCalls, readinessCalls, reopenCalls, createCalls, clearCalls }).toEqual({
|
expect({ showCalls, readinessCalls, reopenCalls, createCalls, clearCalls }).toEqual({
|
||||||
showCalls: 1,
|
// Each caller is authorized against repository ownership, including the request which
|
||||||
|
// finds the runtime already active after waiting on the lifecycle lock.
|
||||||
|
showCalls: 2,
|
||||||
readinessCalls: 1,
|
readinessCalls: 1,
|
||||||
reopenCalls: 1,
|
reopenCalls: 1,
|
||||||
createCalls: 1,
|
createCalls: 1,
|
||||||
@@ -942,7 +1029,7 @@ test("Delete followed by Resume cannot resurrect the deleted session", async ()
|
|||||||
await deleteResponse;
|
await deleteResponse;
|
||||||
const resumed = await resumeResponse;
|
const resumed = await resumeResponse;
|
||||||
|
|
||||||
expect(resumed.statusCode).toBe(500);
|
expect(resumed.statusCode).toBe(503);
|
||||||
expect(current).toBeUndefined();
|
expect(current).toBeUndefined();
|
||||||
expect(createCalls).toBe(0);
|
expect(createCalls).toBe(0);
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -48,6 +48,35 @@ test("PUT /settings persists and GET reads it back", async () => {
|
|||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("settings are isolated by the authenticated repository principal", async () => {
|
||||||
|
const preferences = new Map<string, any>();
|
||||||
|
const runner = {
|
||||||
|
withPrincipal: (principal: any) => ({
|
||||||
|
preferencesGet: async () => preferences.get(principal.subject) ?? {},
|
||||||
|
preferencesSet: async (next: any) => { preferences.set(principal.subject, next); },
|
||||||
|
}),
|
||||||
|
};
|
||||||
|
const app = buildApp(loadConfig({ AUTH_MODE: "upstream", THT_HARNESS_DIR: "../harness" }), {
|
||||||
|
thtRunner: runner as any,
|
||||||
|
listModels: async () => [],
|
||||||
|
});
|
||||||
|
const headers = (subject: string) => ({
|
||||||
|
"x-thoth-principal-issuer": "portal",
|
||||||
|
"x-thoth-principal-subject": subject,
|
||||||
|
"x-thoth-is-admin": "0",
|
||||||
|
});
|
||||||
|
|
||||||
|
await app.inject({
|
||||||
|
method: "PUT", url: "/settings", headers: headers("alice"),
|
||||||
|
payload: { workspace: "psd", provider: "zai", model: "glm-5.2", thinking: "high" },
|
||||||
|
});
|
||||||
|
const alice = await app.inject({ method: "GET", url: "/settings", headers: headers("alice") });
|
||||||
|
const bob = await app.inject({ method: "GET", url: "/settings", headers: headers("bob") });
|
||||||
|
|
||||||
|
expect(alice.json()).toMatchObject({ workspace: "psd", thinking: "high" });
|
||||||
|
expect(bob.json()).not.toMatchObject({ workspace: "psd", thinking: "high" });
|
||||||
|
});
|
||||||
|
|
||||||
test("PUT /settings rejects an unknown model when a model list is available", async () => {
|
test("PUT /settings rejects an unknown model when a model list is available", async () => {
|
||||||
const { app, dir } = appWithTmpSettings({}, {
|
const { app, dir } = appWithTmpSettings({}, {
|
||||||
listModels: async () => [{ provider: "zai", id: "glm-5.2", name: "GLM 5.2", reasoning: true }],
|
listModels: async () => [{ provider: "zai", id: "glm-5.2", name: "GLM 5.2", reasoning: true }],
|
||||||
|
|||||||
@@ -9,6 +9,9 @@ from tht.cli.schema_cmd import _load_config_or_exit
|
|||||||
|
|
||||||
session_app = typer.Typer(help="Sessioni (directory artefatti)")
|
session_app = typer.Typer(help="Sessioni (directory artefatti)")
|
||||||
|
|
||||||
|
preferences_app = typer.Typer(help="Preferenze private del principal corrente")
|
||||||
|
session_app.add_typer(preferences_app, name="preferences")
|
||||||
|
|
||||||
|
|
||||||
def session_repository(cfg):
|
def session_repository(cfg):
|
||||||
"""Configured persistence boundary for every workflow command."""
|
"""Configured persistence boundary for every workflow command."""
|
||||||
@@ -17,6 +20,42 @@ def session_repository(cfg):
|
|||||||
return build_session_repository(cfg)
|
return build_session_repository(cfg)
|
||||||
|
|
||||||
|
|
||||||
|
@preferences_app.command("get")
|
||||||
|
def preferences_get_cmd(
|
||||||
|
json_out: bool = typer.Option(False, "--json", help="Emetti JSON puro."),
|
||||||
|
config: Path = CONFIG_OPT,
|
||||||
|
) -> None:
|
||||||
|
"""Legge le preferenze private del principal risolto dal repository."""
|
||||||
|
cfg = _load_config_or_exit(config)
|
||||||
|
preferences = session_repository(cfg).get_preferences()
|
||||||
|
if json_out:
|
||||||
|
typer.echo(json.dumps(preferences, ensure_ascii=False, sort_keys=True))
|
||||||
|
else:
|
||||||
|
typer.echo(json.dumps(preferences, ensure_ascii=False, indent=2, sort_keys=True))
|
||||||
|
|
||||||
|
|
||||||
|
@preferences_app.command("set")
|
||||||
|
def preferences_set_cmd(
|
||||||
|
preferences_json: str = typer.Argument(..., help="Oggetto JSON delle preferenze."),
|
||||||
|
json_out: bool = typer.Option(False, "--json", help="Non emette testo al successo."),
|
||||||
|
config: Path = CONFIG_OPT,
|
||||||
|
) -> None:
|
||||||
|
"""Sostituisce le preferenze private del principal risolto dal repository."""
|
||||||
|
try:
|
||||||
|
preferences = json.loads(preferences_json)
|
||||||
|
except json.JSONDecodeError as exc:
|
||||||
|
typer.secho(f"ERRORE: preferenze JSON non valide: {exc}", fg=typer.colors.RED, err=True)
|
||||||
|
raise typer.Exit(code=2) from None
|
||||||
|
if not isinstance(preferences, dict):
|
||||||
|
typer.secho("ERRORE: le preferenze devono essere un oggetto JSON", fg=typer.colors.RED, err=True)
|
||||||
|
raise typer.Exit(code=2)
|
||||||
|
cfg = _load_config_or_exit(config)
|
||||||
|
session_repository(cfg).set_preferences(preferences)
|
||||||
|
if json_out:
|
||||||
|
return
|
||||||
|
typer.secho("OK: preferenze aggiornate.", fg=typer.colors.GREEN)
|
||||||
|
|
||||||
|
|
||||||
def session_dir(cfg, session_id: str) -> Path:
|
def session_dir(cfg, session_id: str) -> Path:
|
||||||
"""Legacy path bridge for out-of-scope datamart/memory compatibility only.
|
"""Legacy path bridge for out-of-scope datamart/memory compatibility only.
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user