3 Commits

Author SHA1 Message Date
hibna 218452706c chore: initial commit for phase04 2026-02-21 15:50:35 +03:00
hibna d0c20581b6 chore: initial commit for phase03 2026-02-21 13:37:46 +03:00
hibna 8eb7c90958 chore: update gitignore for phase02 2026-02-21 13:22:51 +03:00
42 changed files with 5964 additions and 37 deletions
+1
View File
@@ -36,3 +36,4 @@ build/
# Claude
.claude/
plans.md
+4 -1
View File
@@ -4,7 +4,7 @@
"private": true,
"type": "module",
"scripts": {
"dev": "tsx watch src/index.ts",
"dev": "dotenv -e ../../.env -- tsx watch src/index.ts",
"build": "tsc",
"start": "node dist/index.js",
"lint": "eslint src/"
@@ -18,11 +18,14 @@
"@source/database": "workspace:*",
"@source/shared": "workspace:*",
"argon2": "^0.41.0",
"drizzle-orm": "^0.38.0",
"fastify": "^5.2.0",
"fastify-plugin": "^5.0.0",
"pino-pretty": "^13.0.0",
"socket.io": "^4.8.0"
},
"devDependencies": {
"dotenv-cli": "^8.0.0",
"tsx": "^4.19.0"
}
}
+55 -3
View File
@@ -1,26 +1,78 @@
import Fastify from 'fastify';
import cors from '@fastify/cors';
import cookie from '@fastify/cookie';
import dbPlugin from './plugins/db.js';
import authPlugin from './plugins/auth.js';
import authRoutes from './routes/auth/index.js';
import organizationRoutes from './routes/organizations/index.js';
import nodeRoutes from './routes/nodes/index.js';
import serverRoutes from './routes/servers/index.js';
import adminRoutes from './routes/admin/index.js';
import { AppError } from './lib/errors.js';
const app = Fastify({
logger: {
transport: {
target: 'pino-pretty',
},
transport:
process.env.NODE_ENV !== 'production'
? { target: 'pino-pretty' }
: undefined,
},
});
// Plugins
await app.register(cors, {
origin: process.env.CORS_ORIGIN || 'http://localhost:5173',
credentials: true,
});
await app.register(cookie);
await app.register(dbPlugin);
await app.register(authPlugin);
// Error handler
app.setErrorHandler((error: Error & { validation?: unknown; statusCode?: number; code?: string }, _request, reply) => {
if (error instanceof AppError) {
return reply.code(error.statusCode).send({
error: error.name,
message: error.message,
code: error.code,
});
}
// Fastify validation errors
if (error.validation) {
return reply.code(400).send({
error: 'Validation Error',
message: error.message,
});
}
app.log.error(error);
return reply.code(500).send({
error: 'Internal Server Error',
message: 'An unexpected error occurred',
});
});
// Routes
app.get('/api/health', async () => {
return { status: 'ok', timestamp: new Date().toISOString() };
});
await app.register(authRoutes, { prefix: '/api/auth' });
await app.register(organizationRoutes, { prefix: '/api/organizations' });
await app.register(adminRoutes, { prefix: '/api/admin' });
// Nested org routes: nodes and servers are scoped to an org
await app.register(
async (orgScope) => {
await orgScope.register(nodeRoutes, { prefix: '/nodes' });
await orgScope.register(serverRoutes, { prefix: '/servers' });
},
{ prefix: '/api/organizations/:orgId' },
);
// Start
const PORT = Number(process.env.PORT) || 3000;
const HOST = process.env.HOST || '0.0.0.0';
+23
View File
@@ -0,0 +1,23 @@
import type { FastifyRequest } from 'fastify';
import { auditLogs } from '@source/database';
import type { Database } from '@source/database';
export async function createAuditLog(
db: Database,
request: FastifyRequest,
data: {
organizationId: string;
action: string;
serverId?: string;
metadata?: Record<string, unknown>;
},
) {
await db.insert(auditLogs).values({
organizationId: data.organizationId,
userId: request.user.sub,
serverId: data.serverId,
action: data.action,
metadata: data.metadata ?? {},
ipAddress: request.ip,
});
}
+30
View File
@@ -0,0 +1,30 @@
export class AppError extends Error {
constructor(
public statusCode: number,
message: string,
public code?: string,
) {
super(message);
this.name = 'AppError';
}
static badRequest(message: string, code?: string) {
return new AppError(400, message, code);
}
static unauthorized(message = 'Unauthorized', code?: string) {
return new AppError(401, message, code);
}
static forbidden(message = 'Forbidden', code?: string) {
return new AppError(403, message, code);
}
static notFound(message = 'Not found', code?: string) {
return new AppError(404, message, code);
}
static conflict(message: string, code?: string) {
return new AppError(409, message, code);
}
}
+27
View File
@@ -0,0 +1,27 @@
import type { FastifyInstance } from 'fastify';
export interface AccessTokenPayload {
sub: string; // user id
email: string;
isSuperAdmin: boolean;
}
export interface RefreshTokenPayload {
sub: string; // user id
type: 'refresh';
}
const ACCESS_TOKEN_EXPIRY = '15m';
const REFRESH_TOKEN_EXPIRY = '7d';
export function signAccessToken(app: FastifyInstance, payload: AccessTokenPayload): string {
return app.jwt.sign(payload, { expiresIn: ACCESS_TOKEN_EXPIRY });
}
export function signRefreshToken(app: FastifyInstance, payload: RefreshTokenPayload): string {
return (app as any).jwtRefresh.sign(payload, { expiresIn: REFRESH_TOKEN_EXPIRY });
}
export function verifyRefreshToken(app: FastifyInstance, token: string): RefreshTokenPayload {
return (app as any).jwtRefresh.verify(token) as RefreshTokenPayload;
}
+25
View File
@@ -0,0 +1,25 @@
import { Type } from '@sinclair/typebox';
export const PaginationQuerySchema = Type.Object({
page: Type.Optional(Type.Number({ minimum: 1, default: 1 })),
perPage: Type.Optional(Type.Number({ minimum: 1, maximum: 100, default: 20 })),
});
export function paginate(query: { page?: number; perPage?: number }) {
const page = query.page ?? 1;
const perPage = query.perPage ?? 20;
const offset = (page - 1) * perPage;
return { page, perPage, offset, limit: perPage };
}
export function paginatedResponse<T>(data: T[], total: number, page: number, perPage: number) {
return {
data,
meta: {
page,
perPage,
total,
totalPages: Math.ceil(total / perPage),
},
};
}
+14
View File
@@ -0,0 +1,14 @@
import argon2 from 'argon2';
export async function hashPassword(password: string): Promise<string> {
return argon2.hash(password, {
type: argon2.argon2id,
memoryCost: 65536,
timeCost: 3,
parallelism: 4,
});
}
export async function verifyPassword(hash: string, password: string): Promise<boolean> {
return argon2.verify(hash, password);
}
+82
View File
@@ -0,0 +1,82 @@
import type { FastifyRequest } from 'fastify';
import { eq, and } from 'drizzle-orm';
import { organizationMembers } from '@source/database';
import { ROLES } from '@source/shared';
import type { Permission, Role } from '@source/shared';
import { AppError } from './errors.js';
interface OrgMember {
role: Role;
customPermissions: Record<string, boolean>;
}
/**
* Get the requesting user's membership in an organization.
* Super admins bypass membership checks.
*/
export async function getOrgMembership(
request: FastifyRequest,
orgId: string,
): Promise<OrgMember | 'super_admin'> {
const user = request.user;
if (user.isSuperAdmin) {
return 'super_admin';
}
const member = await (request.server as any).db.query.organizationMembers.findFirst({
where: and(
eq(organizationMembers.organizationId, orgId),
eq(organizationMembers.userId, user.sub),
),
});
if (!member) {
throw AppError.forbidden('You are not a member of this organization');
}
return {
role: member.role as Role,
customPermissions: (member.customPermissions ?? {}) as Record<string, boolean>,
};
}
/**
* Check if the user has a specific permission in the organization.
* Super admins always have all permissions.
*/
export function hasPermission(membership: OrgMember | 'super_admin', permission: Permission): boolean {
if (membership === 'super_admin') return true;
// Check custom permission overrides first
if (permission in membership.customPermissions) {
return membership.customPermissions[permission]!;
}
// Fall back to role defaults
const rolePerms = ROLES[membership.role]?.permissions ?? [];
return (rolePerms as readonly string[]).includes(permission);
}
/**
* Require a specific permission, throw 403 if not allowed.
*/
export async function requirePermission(
request: FastifyRequest,
orgId: string,
permission: Permission,
): Promise<void> {
const membership = await getOrgMembership(request, orgId);
if (!hasPermission(membership, permission)) {
throw AppError.forbidden(`Missing permission: ${permission}`);
}
}
/**
* Require super admin role.
*/
export function requireSuperAdmin(request: FastifyRequest): void {
if (!request.user.isSuperAdmin) {
throw AppError.forbidden('Super admin access required');
}
}
+48
View File
@@ -0,0 +1,48 @@
import fp from 'fastify-plugin';
import jwt from '@fastify/jwt';
import type { FastifyInstance, FastifyRequest, FastifyReply } from 'fastify';
import type { AccessTokenPayload } from '../lib/jwt.js';
declare module 'fastify' {
interface FastifyInstance {
authenticate: (request: FastifyRequest, reply: FastifyReply) => Promise<void>;
jwtRefresh: FastifyInstance['jwt'];
}
}
declare module '@fastify/jwt' {
interface FastifyJWT {
payload: AccessTokenPayload;
user: AccessTokenPayload;
}
}
export default fp(async (app: FastifyInstance) => {
const jwtSecret = process.env.JWT_SECRET;
const jwtRefreshSecret = process.env.JWT_REFRESH_SECRET;
if (!jwtSecret || !jwtRefreshSecret) {
throw new Error('JWT_SECRET and JWT_REFRESH_SECRET environment variables are required');
}
// Access token JWT
await app.register(jwt, {
secret: jwtSecret,
namespace: 'jwt',
});
// Refresh token JWT (separate namespace)
await app.register(jwt, {
secret: jwtRefreshSecret,
namespace: 'jwtRefresh',
});
// Auth decorator
app.decorate('authenticate', async (request: FastifyRequest, reply: FastifyReply) => {
try {
await request.jwtVerify();
} catch {
reply.code(401).send({ error: 'Unauthorized', message: 'Invalid or expired token' });
}
});
});
+21
View File
@@ -0,0 +1,21 @@
import fp from 'fastify-plugin';
import type { FastifyInstance } from 'fastify';
import { createDb, type Database } from '@source/database';
declare module 'fastify' {
interface FastifyInstance {
db: Database;
}
}
export default fp(async (app: FastifyInstance) => {
const databaseUrl = process.env.DATABASE_URL;
if (!databaseUrl) {
throw new Error('DATABASE_URL environment variable is required');
}
const db = createDb(databaseUrl);
app.decorate('db', db);
app.log.info('Database connected');
});
+140
View File
@@ -0,0 +1,140 @@
import type { FastifyInstance } from 'fastify';
import { eq, desc, count } from 'drizzle-orm';
import { users, games, nodes, auditLogs } from '@source/database';
import { AppError } from '../../lib/errors.js';
import { requireSuperAdmin } from '../../lib/permissions.js';
import { paginate, paginatedResponse, PaginationQuerySchema } from '../../lib/pagination.js';
import { CreateGameSchema, UpdateGameSchema, GameIdParamSchema } from './schemas.js';
export default async function adminRoutes(app: FastifyInstance) {
// All admin routes require auth + super admin
app.addHook('onRequest', app.authenticate);
app.addHook('onRequest', async (request) => {
requireSuperAdmin(request);
});
// === Users ===
// GET /api/admin/users
app.get('/users', { schema: { querystring: PaginationQuerySchema } }, async (request) => {
const { page, perPage, offset, limit } = paginate(request.query as any);
const [totalResult] = await app.db.select({ count: count() }).from(users);
const userList = await app.db
.select({
id: users.id,
email: users.email,
username: users.username,
isSuperAdmin: users.isSuperAdmin,
avatarUrl: users.avatarUrl,
createdAt: users.createdAt,
})
.from(users)
.limit(limit)
.offset(offset)
.orderBy(users.createdAt);
return paginatedResponse(userList, totalResult!.count, page, perPage);
});
// === Games ===
// GET /api/admin/games
app.get('/games', async () => {
const gameList = await app.db
.select()
.from(games)
.orderBy(games.name);
return { data: gameList };
});
// POST /api/admin/games
app.post('/games', { schema: CreateGameSchema }, async (request, reply) => {
const body = request.body as {
slug: string;
name: string;
dockerImage: string;
defaultPort: number;
startupCommand: string;
stopCommand?: string;
configFiles?: unknown[];
environmentVars?: unknown[];
};
const existing = await app.db.query.games.findFirst({
where: eq(games.slug, body.slug),
});
if (existing) throw AppError.conflict('Game slug already exists');
const [game] = await app.db
.insert(games)
.values({
...body,
configFiles: body.configFiles ?? [],
environmentVars: body.environmentVars ?? [],
})
.returning();
return reply.code(201).send(game);
});
// PATCH /api/admin/games/:gameId
app.patch('/games/:gameId', { schema: { ...GameIdParamSchema, ...UpdateGameSchema } }, async (request) => {
const { gameId } = request.params as { gameId: string };
const body = request.body as Record<string, unknown>;
const [updated] = await app.db
.update(games)
.set({ ...body, updatedAt: new Date() })
.where(eq(games.id, gameId))
.returning();
if (!updated) throw AppError.notFound('Game not found');
return updated;
});
// === Nodes (global view) ===
// GET /api/admin/nodes
app.get('/nodes', async () => {
const nodeList = await app.db
.select()
.from(nodes)
.orderBy(nodes.createdAt);
return { data: nodeList };
});
// === Audit Logs ===
// GET /api/admin/audit-logs
app.get('/audit-logs', { schema: { querystring: PaginationQuerySchema } }, async (request) => {
const { page, perPage, offset, limit } = paginate(request.query as any);
const [totalResult] = await app.db.select({ count: count() }).from(auditLogs);
const logs = await app.db
.select({
id: auditLogs.id,
organizationId: auditLogs.organizationId,
userId: auditLogs.userId,
serverId: auditLogs.serverId,
action: auditLogs.action,
metadata: auditLogs.metadata,
ipAddress: auditLogs.ipAddress,
createdAt: auditLogs.createdAt,
userEmail: users.email,
userName: users.username,
})
.from(auditLogs)
.innerJoin(users, eq(auditLogs.userId, users.id))
.orderBy(desc(auditLogs.createdAt))
.limit(limit)
.offset(offset);
return paginatedResponse(logs, totalResult!.count, page, perPage);
});
}
+32
View File
@@ -0,0 +1,32 @@
import { Type } from '@sinclair/typebox';
export const CreateGameSchema = {
body: Type.Object({
slug: Type.String({ minLength: 1, maxLength: 100, pattern: '^[a-z0-9-]+$' }),
name: Type.String({ minLength: 1, maxLength: 255 }),
dockerImage: Type.String({ minLength: 1 }),
defaultPort: Type.Number({ minimum: 1, maximum: 65535 }),
startupCommand: Type.String({ minLength: 1 }),
stopCommand: Type.Optional(Type.String()),
configFiles: Type.Optional(Type.Array(Type.Any())),
environmentVars: Type.Optional(Type.Array(Type.Any())),
}),
};
export const UpdateGameSchema = {
body: Type.Object({
name: Type.Optional(Type.String({ minLength: 1, maxLength: 255 })),
dockerImage: Type.Optional(Type.String({ minLength: 1 })),
defaultPort: Type.Optional(Type.Number({ minimum: 1, maximum: 65535 })),
startupCommand: Type.Optional(Type.String({ minLength: 1 })),
stopCommand: Type.Optional(Type.String()),
configFiles: Type.Optional(Type.Array(Type.Any())),
environmentVars: Type.Optional(Type.Array(Type.Any())),
}),
};
export const GameIdParamSchema = {
params: Type.Object({
gameId: Type.String({ format: 'uuid' }),
}),
};
+196
View File
@@ -0,0 +1,196 @@
import type { FastifyInstance } from 'fastify';
import { eq } from 'drizzle-orm';
import { users } from '@source/database';
import { hashPassword, verifyPassword } from '../../lib/password.js';
import { signAccessToken, signRefreshToken, verifyRefreshToken } from '../../lib/jwt.js';
import type { AccessTokenPayload, RefreshTokenPayload } from '../../lib/jwt.js';
import { AppError } from '../../lib/errors.js';
import { RegisterSchema, LoginSchema } from './schemas.js';
const REFRESH_COOKIE_NAME = 'refresh_token';
const REFRESH_COOKIE_OPTIONS = {
httpOnly: true,
secure: process.env.NODE_ENV === 'production',
sameSite: 'lax' as const,
path: '/api/auth',
maxAge: 7 * 24 * 60 * 60, // 7 days in seconds
};
export default async function authRoutes(app: FastifyInstance) {
// POST /api/auth/register
app.post('/register', { schema: RegisterSchema }, async (request, reply) => {
const { email, username, password } = request.body as {
email: string;
username: string;
password: string;
};
// Check if email already exists
const existingEmail = await app.db.query.users.findFirst({
where: eq(users.email, email),
});
if (existingEmail) {
throw AppError.conflict('Email already in use', 'EMAIL_TAKEN');
}
// Check if username already exists
const existingUsername = await app.db.query.users.findFirst({
where: eq(users.username, username),
});
if (existingUsername) {
throw AppError.conflict('Username already in use', 'USERNAME_TAKEN');
}
const passwordHash = await hashPassword(password);
const [user] = await app.db
.insert(users)
.values({
email,
username,
passwordHash,
})
.returning({
id: users.id,
email: users.email,
username: users.username,
isSuperAdmin: users.isSuperAdmin,
});
// Generate tokens
const accessToken = signAccessToken(app, {
sub: user!.id,
email: user!.email,
isSuperAdmin: user!.isSuperAdmin,
});
const refreshToken = signRefreshToken(app, {
sub: user!.id,
type: 'refresh',
});
reply.setCookie(REFRESH_COOKIE_NAME, refreshToken, REFRESH_COOKIE_OPTIONS);
return reply.code(201).send({
user: {
id: user!.id,
email: user!.email,
username: user!.username,
isSuperAdmin: user!.isSuperAdmin,
},
accessToken,
});
});
// POST /api/auth/login
app.post('/login', { schema: LoginSchema }, async (request, reply) => {
const { email, password } = request.body as { email: string; password: string };
const user = await app.db.query.users.findFirst({
where: eq(users.email, email),
});
if (!user) {
throw AppError.unauthorized('Invalid email or password', 'INVALID_CREDENTIALS');
}
const isValid = await verifyPassword(user.passwordHash, password);
if (!isValid) {
throw AppError.unauthorized('Invalid email or password', 'INVALID_CREDENTIALS');
}
const accessToken = signAccessToken(app, {
sub: user.id,
email: user.email,
isSuperAdmin: user.isSuperAdmin,
});
const refreshToken = signRefreshToken(app, {
sub: user.id,
type: 'refresh',
});
reply.setCookie(REFRESH_COOKIE_NAME, refreshToken, REFRESH_COOKIE_OPTIONS);
return {
user: {
id: user.id,
email: user.email,
username: user.username,
isSuperAdmin: user.isSuperAdmin,
avatarUrl: user.avatarUrl,
},
accessToken,
};
});
// POST /api/auth/refresh
app.post('/refresh', async (request, reply) => {
const token = request.cookies[REFRESH_COOKIE_NAME];
if (!token) {
throw AppError.unauthorized('No refresh token', 'NO_REFRESH_TOKEN');
}
let payload: RefreshTokenPayload;
try {
payload = verifyRefreshToken(app, token);
} catch {
reply.clearCookie(REFRESH_COOKIE_NAME, { path: '/api/auth' });
throw AppError.unauthorized('Invalid refresh token', 'INVALID_REFRESH_TOKEN');
}
const user = await app.db.query.users.findFirst({
where: eq(users.id, payload.sub),
});
if (!user) {
reply.clearCookie(REFRESH_COOKIE_NAME, { path: '/api/auth' });
throw AppError.unauthorized('User not found', 'USER_NOT_FOUND');
}
// Token rotation: issue new tokens
const accessToken = signAccessToken(app, {
sub: user.id,
email: user.email,
isSuperAdmin: user.isSuperAdmin,
});
const newRefreshToken = signRefreshToken(app, {
sub: user.id,
type: 'refresh',
});
reply.setCookie(REFRESH_COOKIE_NAME, newRefreshToken, REFRESH_COOKIE_OPTIONS);
return { accessToken };
});
// POST /api/auth/logout
app.post('/logout', async (_request, reply) => {
reply.clearCookie(REFRESH_COOKIE_NAME, { path: '/api/auth' });
return { success: true };
});
// GET /api/auth/me
app.get('/me', { onRequest: [app.authenticate] }, async (request) => {
const payload = request.user;
const user = await app.db.query.users.findFirst({
where: eq(users.id, payload.sub),
columns: {
id: true,
email: true,
username: true,
isSuperAdmin: true,
avatarUrl: true,
createdAt: true,
},
});
if (!user) {
throw AppError.notFound('User not found');
}
return { user };
});
}
+16
View File
@@ -0,0 +1,16 @@
import { Type } from '@sinclair/typebox';
export const RegisterSchema = {
body: Type.Object({
email: Type.String({ format: 'email' }),
username: Type.String({ minLength: 3, maxLength: 100 }),
password: Type.String({ minLength: 8, maxLength: 128 }),
}),
};
export const LoginSchema = {
body: Type.Object({
email: Type.String({ format: 'email' }),
password: Type.String(),
}),
};
+170
View File
@@ -0,0 +1,170 @@
import type { FastifyInstance } from 'fastify';
import { eq, and } from 'drizzle-orm';
import { randomBytes } from 'crypto';
import { nodes, allocations } from '@source/database';
import { AppError } from '../../lib/errors.js';
import { requirePermission } from '../../lib/permissions.js';
import { createAuditLog } from '../../lib/audit.js';
import {
NodeParamSchema,
CreateNodeSchema,
UpdateNodeSchema,
CreateAllocationSchema,
} from './schemas.js';
export default async function nodeRoutes(app: FastifyInstance) {
app.addHook('onRequest', app.authenticate);
// GET /api/organizations/:orgId/nodes
app.get('/', async (request) => {
const { orgId } = request.params as { orgId: string };
await requirePermission(request, orgId, 'node.read');
const nodeList = await app.db
.select()
.from(nodes)
.where(eq(nodes.organizationId, orgId))
.orderBy(nodes.createdAt);
return { data: nodeList };
});
// POST /api/organizations/:orgId/nodes
app.post('/', { schema: CreateNodeSchema }, async (request, reply) => {
const { orgId } = request.params as { orgId: string };
await requirePermission(request, orgId, 'node.manage');
const body = request.body as {
name: string;
fqdn: string;
daemonPort?: number;
grpcPort?: number;
location?: string;
memoryTotal: number;
diskTotal: number;
memoryOveralloc?: number;
diskOveralloc?: number;
};
const daemonToken = randomBytes(32).toString('hex');
const [node] = await app.db
.insert(nodes)
.values({
organizationId: orgId,
...body,
daemonToken,
})
.returning();
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'node.create',
metadata: { nodeId: node!.id, name: body.name },
});
return reply.code(201).send(node);
});
// GET /api/organizations/:orgId/nodes/:nodeId
app.get('/:nodeId', { schema: NodeParamSchema }, async (request) => {
const { orgId, nodeId } = request.params as { orgId: string; nodeId: string };
await requirePermission(request, orgId, 'node.read');
const node = await app.db.query.nodes.findFirst({
where: and(eq(nodes.id, nodeId), eq(nodes.organizationId, orgId)),
});
if (!node) throw AppError.notFound('Node not found');
return node;
});
// PATCH /api/organizations/:orgId/nodes/:nodeId
app.patch('/:nodeId', { schema: { ...NodeParamSchema, ...UpdateNodeSchema } }, async (request) => {
const { orgId, nodeId } = request.params as { orgId: string; nodeId: string };
await requirePermission(request, orgId, 'node.manage');
const body = request.body as Record<string, unknown>;
const [updated] = await app.db
.update(nodes)
.set({ ...body, updatedAt: new Date() })
.where(and(eq(nodes.id, nodeId), eq(nodes.organizationId, orgId)))
.returning();
if (!updated) throw AppError.notFound('Node not found');
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'node.update',
metadata: { nodeId, ...body },
});
return updated;
});
// DELETE /api/organizations/:orgId/nodes/:nodeId
app.delete('/:nodeId', { schema: NodeParamSchema }, async (request, reply) => {
const { orgId, nodeId } = request.params as { orgId: string; nodeId: string };
await requirePermission(request, orgId, 'node.manage');
const node = await app.db.query.nodes.findFirst({
where: and(eq(nodes.id, nodeId), eq(nodes.organizationId, orgId)),
});
if (!node) throw AppError.notFound('Node not found');
await app.db.delete(nodes).where(eq(nodes.id, nodeId));
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'node.delete',
metadata: { nodeId, name: node.name },
});
return reply.code(204).send();
});
// === Allocations ===
// GET /api/organizations/:orgId/nodes/:nodeId/allocations
app.get('/:nodeId/allocations', { schema: NodeParamSchema }, async (request) => {
const { orgId, nodeId } = request.params as { orgId: string; nodeId: string };
await requirePermission(request, orgId, 'node.read');
const allocs = await app.db
.select()
.from(allocations)
.where(eq(allocations.nodeId, nodeId))
.orderBy(allocations.port);
return { data: allocs };
});
// POST /api/organizations/:orgId/nodes/:nodeId/allocations
app.post('/:nodeId/allocations', { schema: { ...NodeParamSchema, ...CreateAllocationSchema } }, async (request, reply) => {
const { orgId, nodeId } = request.params as { orgId: string; nodeId: string };
await requirePermission(request, orgId, 'node.manage');
const { ip, ports } = request.body as { ip: string; ports: number[] };
const values = ports.map((port) => ({
nodeId,
ip,
port,
}));
const created = await app.db
.insert(allocations)
.values(values)
.onConflictDoNothing()
.returning();
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'allocation.create',
metadata: { nodeId, ip, ports },
});
return reply.code(201).send({ data: created });
});
}
+43
View File
@@ -0,0 +1,43 @@
import { Type } from '@sinclair/typebox';
export const NodeParamSchema = {
params: Type.Object({
orgId: Type.String({ format: 'uuid' }),
nodeId: Type.String({ format: 'uuid' }),
}),
};
export const CreateNodeSchema = {
body: Type.Object({
name: Type.String({ minLength: 1, maxLength: 255 }),
fqdn: Type.String({ minLength: 1, maxLength: 255 }),
daemonPort: Type.Optional(Type.Number({ minimum: 1, maximum: 65535, default: 8443 })),
grpcPort: Type.Optional(Type.Number({ minimum: 1, maximum: 65535, default: 50051 })),
location: Type.Optional(Type.String({ maxLength: 255 })),
memoryTotal: Type.Number({ minimum: 0 }),
diskTotal: Type.Number({ minimum: 0 }),
memoryOveralloc: Type.Optional(Type.Number({ minimum: 0, default: 0 })),
diskOveralloc: Type.Optional(Type.Number({ minimum: 0, default: 0 })),
}),
};
export const UpdateNodeSchema = {
body: Type.Object({
name: Type.Optional(Type.String({ minLength: 1, maxLength: 255 })),
fqdn: Type.Optional(Type.String({ minLength: 1, maxLength: 255 })),
daemonPort: Type.Optional(Type.Number({ minimum: 1, maximum: 65535 })),
grpcPort: Type.Optional(Type.Number({ minimum: 1, maximum: 65535 })),
location: Type.Optional(Type.String({ maxLength: 255 })),
memoryTotal: Type.Optional(Type.Number({ minimum: 0 })),
diskTotal: Type.Optional(Type.Number({ minimum: 0 })),
memoryOveralloc: Type.Optional(Type.Number({ minimum: 0 })),
diskOveralloc: Type.Optional(Type.Number({ minimum: 0 })),
}),
};
export const CreateAllocationSchema = {
body: Type.Object({
ip: Type.String({ minLength: 1, maxLength: 45 }),
ports: Type.Array(Type.Number({ minimum: 1, maximum: 65535 }), { minItems: 1 }),
}),
};
+276
View File
@@ -0,0 +1,276 @@
import type { FastifyInstance } from 'fastify';
import { eq, and, count } from 'drizzle-orm';
import { organizations, organizationMembers, users } from '@source/database';
import { AppError } from '../../lib/errors.js';
import { requirePermission, getOrgMembership } from '../../lib/permissions.js';
import { paginate, paginatedResponse, PaginationQuerySchema } from '../../lib/pagination.js';
import { createAuditLog } from '../../lib/audit.js';
import {
CreateOrgSchema,
UpdateOrgSchema,
OrgIdParamSchema,
AddMemberSchema,
UpdateMemberSchema,
MemberIdParamSchema,
} from './schemas.js';
export default async function organizationRoutes(app: FastifyInstance) {
// All org routes require authentication
app.addHook('onRequest', app.authenticate);
// GET /api/organizations — list user's organizations
app.get('/', { schema: { querystring: PaginationQuerySchema } }, async (request) => {
const { page, perPage, offset, limit } = paginate(request.query as any);
const userId = request.user.sub;
if (request.user.isSuperAdmin) {
const [totalResult] = await app.db.select({ count: count() }).from(organizations);
const orgs = await app.db
.select()
.from(organizations)
.limit(limit)
.offset(offset)
.orderBy(organizations.createdAt);
return paginatedResponse(orgs, totalResult!.count, page, perPage);
}
const memberOrgs = await app.db
.select({
id: organizations.id,
name: organizations.name,
slug: organizations.slug,
ownerId: organizations.ownerId,
maxServers: organizations.maxServers,
maxNodes: organizations.maxNodes,
createdAt: organizations.createdAt,
updatedAt: organizations.updatedAt,
role: organizationMembers.role,
})
.from(organizationMembers)
.innerJoin(organizations, eq(organizationMembers.organizationId, organizations.id))
.where(eq(organizationMembers.userId, userId))
.limit(limit)
.offset(offset);
const [totalResult] = await app.db
.select({ count: count() })
.from(organizationMembers)
.where(eq(organizationMembers.userId, userId));
return paginatedResponse(memberOrgs, totalResult!.count, page, perPage);
});
// POST /api/organizations — create organization
app.post('/', { schema: CreateOrgSchema }, async (request, reply) => {
const { name, slug } = request.body as { name: string; slug: string };
const existing = await app.db.query.organizations.findFirst({
where: eq(organizations.slug, slug),
});
if (existing) {
throw AppError.conflict('Organization slug already in use', 'SLUG_TAKEN');
}
const [org] = await app.db
.insert(organizations)
.values({
name,
slug,
ownerId: request.user.sub,
})
.returning();
// Add creator as admin member
await app.db.insert(organizationMembers).values({
organizationId: org!.id,
userId: request.user.sub,
role: 'admin',
});
return reply.code(201).send(org);
});
// GET /api/organizations/:orgId
app.get('/:orgId', { schema: OrgIdParamSchema }, async (request) => {
const { orgId } = request.params as { orgId: string };
await getOrgMembership(request, orgId);
const org = await app.db.query.organizations.findFirst({
where: eq(organizations.id, orgId),
});
if (!org) throw AppError.notFound('Organization not found');
return org;
});
// PATCH /api/organizations/:orgId
app.patch('/:orgId', { schema: { ...OrgIdParamSchema, ...UpdateOrgSchema } }, async (request) => {
const { orgId } = request.params as { orgId: string };
await requirePermission(request, orgId, 'org.settings');
const body = request.body as { name?: string; maxServers?: number; maxNodes?: number };
const [updated] = await app.db
.update(organizations)
.set({ ...body, updatedAt: new Date() })
.where(eq(organizations.id, orgId))
.returning();
if (!updated) throw AppError.notFound('Organization not found');
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'organization.update',
metadata: body,
});
return updated;
});
// DELETE /api/organizations/:orgId
app.delete('/:orgId', { schema: OrgIdParamSchema }, async (request, reply) => {
const { orgId } = request.params as { orgId: string };
const membership = await getOrgMembership(request, orgId);
// Only owner or super admin can delete
const org = await app.db.query.organizations.findFirst({
where: eq(organizations.id, orgId),
});
if (!org) throw AppError.notFound('Organization not found');
if (membership !== 'super_admin' && org.ownerId !== request.user.sub) {
throw AppError.forbidden('Only the organization owner can delete this organization');
}
await app.db.delete(organizations).where(eq(organizations.id, orgId));
return reply.code(204).send();
});
// === Members ===
// GET /api/organizations/:orgId/members
app.get('/:orgId/members', { schema: OrgIdParamSchema }, async (request) => {
const { orgId } = request.params as { orgId: string };
await requirePermission(request, orgId, 'org.members');
const members = await app.db
.select({
id: organizationMembers.id,
userId: organizationMembers.userId,
role: organizationMembers.role,
customPermissions: organizationMembers.customPermissions,
joinedAt: organizationMembers.joinedAt,
email: users.email,
username: users.username,
avatarUrl: users.avatarUrl,
})
.from(organizationMembers)
.innerJoin(users, eq(organizationMembers.userId, users.id))
.where(eq(organizationMembers.organizationId, orgId));
return { data: members };
});
// POST /api/organizations/:orgId/members — invite by email
app.post('/:orgId/members', { schema: { ...OrgIdParamSchema, ...AddMemberSchema } }, async (request, reply) => {
const { orgId } = request.params as { orgId: string };
await requirePermission(request, orgId, 'org.members');
const { email, role } = request.body as { email: string; role: 'admin' | 'user' };
const user = await app.db.query.users.findFirst({
where: eq(users.email, email),
});
if (!user) throw AppError.notFound('User with this email not found');
const existing = await app.db.query.organizationMembers.findFirst({
where: and(
eq(organizationMembers.organizationId, orgId),
eq(organizationMembers.userId, user.id),
),
});
if (existing) throw AppError.conflict('User is already a member');
const [member] = await app.db
.insert(organizationMembers)
.values({
organizationId: orgId,
userId: user.id,
role,
})
.returning();
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'member.add',
metadata: { userId: user.id, email, role },
});
return reply.code(201).send(member);
});
// PATCH /api/organizations/:orgId/members/:memberId
app.patch('/:orgId/members/:memberId', { schema: { ...MemberIdParamSchema, ...UpdateMemberSchema } }, async (request) => {
const { orgId, memberId } = request.params as { orgId: string; memberId: string };
await requirePermission(request, orgId, 'org.members');
const body = request.body as { role?: 'admin' | 'user'; customPermissions?: Record<string, boolean> };
const [updated] = await app.db
.update(organizationMembers)
.set(body)
.where(and(
eq(organizationMembers.id, memberId),
eq(organizationMembers.organizationId, orgId),
))
.returning();
if (!updated) throw AppError.notFound('Member not found');
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'member.update',
metadata: { memberId, ...body },
});
return updated;
});
// DELETE /api/organizations/:orgId/members/:memberId
app.delete('/:orgId/members/:memberId', { schema: MemberIdParamSchema }, async (request, reply) => {
const { orgId, memberId } = request.params as { orgId: string; memberId: string };
await requirePermission(request, orgId, 'org.members');
const member = await app.db.query.organizationMembers.findFirst({
where: and(
eq(organizationMembers.id, memberId),
eq(organizationMembers.organizationId, orgId),
),
});
if (!member) throw AppError.notFound('Member not found');
// Cannot remove org owner
const org = await app.db.query.organizations.findFirst({
where: eq(organizations.id, orgId),
});
if (org && member.userId === org.ownerId) {
throw AppError.badRequest('Cannot remove the organization owner');
}
await app.db
.delete(organizationMembers)
.where(and(
eq(organizationMembers.id, memberId),
eq(organizationMembers.organizationId, orgId),
));
await createAuditLog(app.db, request, {
organizationId: orgId,
action: 'member.remove',
metadata: { memberId, userId: member.userId },
});
return reply.code(204).send();
});
}
@@ -0,0 +1,43 @@
import { Type } from '@sinclair/typebox';
export const CreateOrgSchema = {
body: Type.Object({
name: Type.String({ minLength: 2, maxLength: 255 }),
slug: Type.String({ minLength: 2, maxLength: 255, pattern: '^[a-z0-9-]+$' }),
}),
};
export const UpdateOrgSchema = {
body: Type.Object({
name: Type.Optional(Type.String({ minLength: 2, maxLength: 255 })),
maxServers: Type.Optional(Type.Number({ minimum: 0 })),
maxNodes: Type.Optional(Type.Number({ minimum: 0 })),
}),
};
export const OrgIdParamSchema = {
params: Type.Object({
orgId: Type.String({ format: 'uuid' }),
}),
};
export const AddMemberSchema = {
body: Type.Object({
email: Type.String({ format: 'email' }),
role: Type.Union([Type.Literal('admin'), Type.Literal('user')]),
}),
};
export const UpdateMemberSchema = {
body: Type.Object({
role: Type.Optional(Type.Union([Type.Literal('admin'), Type.Literal('user')])),
customPermissions: Type.Optional(Type.Record(Type.String(), Type.Boolean())),
}),
};
export const MemberIdParamSchema = {
params: Type.Object({
orgId: Type.String({ format: 'uuid' }),
memberId: Type.String({ format: 'uuid' }),
}),
};
+280
View File
@@ -0,0 +1,280 @@
import type { FastifyInstance } from 'fastify';
import { eq, and, count } from 'drizzle-orm';
import { randomUUID } from 'crypto';
import { servers, allocations, nodes, games } from '@source/database';
import type { PowerAction } from '@source/shared';
import { AppError } from '../../lib/errors.js';
import { requirePermission } from '../../lib/permissions.js';
import { paginate, paginatedResponse, PaginationQuerySchema } from '../../lib/pagination.js';
import { createAuditLog } from '../../lib/audit.js';
import {
ServerParamSchema,
CreateServerSchema,
UpdateServerSchema,
PowerActionSchema,
} from './schemas.js';
export default async function serverRoutes(app: FastifyInstance) {
app.addHook('onRequest', app.authenticate);
// GET /api/organizations/:orgId/servers
app.get('/', { schema: { querystring: PaginationQuerySchema } }, async (request) => {
const { orgId } = request.params as { orgId: string };
await requirePermission(request, orgId, 'server.read');
const { page, perPage, offset, limit } = paginate(request.query as any);
const [totalResult] = await app.db
.select({ count: count() })
.from(servers)
.where(eq(servers.organizationId, orgId));
const serverList = await app.db
.select({
id: servers.id,
uuid: servers.uuid,
name: servers.name,
description: servers.description,
status: servers.status,
memoryLimit: servers.memoryLimit,
diskLimit: servers.diskLimit,
cpuLimit: servers.cpuLimit,
port: servers.port,
createdAt: servers.createdAt,
nodeName: nodes.name,
nodeId: nodes.id,
gameName: games.name,
gameSlug: games.slug,
gameId: games.id,
})
.from(servers)
.innerJoin(nodes, eq(servers.nodeId, nodes.id))
.innerJoin(games, eq(servers.gameId, games.id))
.where(eq(servers.organizationId, orgId))
.limit(limit)
.offset(offset)
.orderBy(servers.createdAt);
return paginatedResponse(serverList, totalResult!.count, page, perPage);
});
// POST /api/organizations/:orgId/servers
app.post('/', { schema: CreateServerSchema }, async (request, reply) => {
const { orgId } = request.params as { orgId: string };
await requirePermission(request, orgId, 'server.create');
const body = request.body as {
name: string;
description?: string;
nodeId: string;
gameId: string;
memoryLimit: number;
diskLimit: number;
cpuLimit?: number;
allocationId: string;
environment?: Record<string, string>;
startupOverride?: string;
};
// Verify node belongs to org
const node = await app.db.query.nodes.findFirst({
where: and(eq(nodes.id, body.nodeId), eq(nodes.organizationId, orgId)),
});
if (!node) throw AppError.notFound('Node not found in this organization');
// Verify game exists
const game = await app.db.query.games.findFirst({
where: eq(games.id, body.gameId),
});
if (!game) throw AppError.notFound('Game not found');
// Verify and claim allocation
const allocation = await app.db.query.allocations.findFirst({
where: and(
eq(allocations.id, body.allocationId),
eq(allocations.nodeId, body.nodeId),
),
});
if (!allocation) throw AppError.notFound('Allocation not found on this node');
if (allocation.serverId) throw AppError.conflict('Allocation is already in use');
const serverUuid = randomUUID().slice(0, 8);
const [server] = await app.db
.insert(servers)
.values({
uuid: serverUuid,
organizationId: orgId,
nodeId: body.nodeId,
gameId: body.gameId,
name: body.name,
description: body.description,
memoryLimit: body.memoryLimit,
diskLimit: body.diskLimit,
cpuLimit: body.cpuLimit ?? 100,
port: allocation.port,
environment: body.environment ?? {},
startupOverride: body.startupOverride,
status: 'installing',
})
.returning();
// Assign allocation to server
await app.db
.update(allocations)
.set({ serverId: server!.id, isDefault: true })
.where(eq(allocations.id, body.allocationId));
// TODO: Send gRPC CreateServer to daemon
// This will be implemented in Phase 4
await createAuditLog(app.db, request, {
organizationId: orgId,
serverId: server!.id,
action: 'server.create',
metadata: { name: body.name, gameSlug: game.slug, nodeId: body.nodeId },
});
return reply.code(201).send(server);
});
// GET /api/organizations/:orgId/servers/:serverId
app.get('/:serverId', { schema: ServerParamSchema }, async (request) => {
const { orgId, serverId } = request.params as { orgId: string; serverId: string };
await requirePermission(request, orgId, 'server.read');
const [server] = await app.db
.select({
id: servers.id,
uuid: servers.uuid,
name: servers.name,
description: servers.description,
status: servers.status,
memoryLimit: servers.memoryLimit,
diskLimit: servers.diskLimit,
cpuLimit: servers.cpuLimit,
port: servers.port,
additionalPorts: servers.additionalPorts,
environment: servers.environment,
startupOverride: servers.startupOverride,
installedAt: servers.installedAt,
createdAt: servers.createdAt,
updatedAt: servers.updatedAt,
nodeId: nodes.id,
nodeName: nodes.name,
nodeFqdn: nodes.fqdn,
gameId: games.id,
gameName: games.name,
gameSlug: games.slug,
})
.from(servers)
.innerJoin(nodes, eq(servers.nodeId, nodes.id))
.innerJoin(games, eq(servers.gameId, games.id))
.where(and(eq(servers.id, serverId), eq(servers.organizationId, orgId)));
if (!server) throw AppError.notFound('Server not found');
return server;
});
// PATCH /api/organizations/:orgId/servers/:serverId
app.patch('/:serverId', { schema: { ...ServerParamSchema, ...UpdateServerSchema } }, async (request) => {
const { orgId, serverId } = request.params as { orgId: string; serverId: string };
await requirePermission(request, orgId, 'server.update');
const body = request.body as Record<string, unknown>;
const [updated] = await app.db
.update(servers)
.set({ ...body, updatedAt: new Date() })
.where(and(eq(servers.id, serverId), eq(servers.organizationId, orgId)))
.returning();
if (!updated) throw AppError.notFound('Server not found');
await createAuditLog(app.db, request, {
organizationId: orgId,
serverId,
action: 'server.update',
metadata: body,
});
return updated;
});
// DELETE /api/organizations/:orgId/servers/:serverId
app.delete('/:serverId', { schema: ServerParamSchema }, async (request, reply) => {
const { orgId, serverId } = request.params as { orgId: string; serverId: string };
await requirePermission(request, orgId, 'server.delete');
const server = await app.db.query.servers.findFirst({
where: and(eq(servers.id, serverId), eq(servers.organizationId, orgId)),
});
if (!server) throw AppError.notFound('Server not found');
// Release allocations
await app.db
.update(allocations)
.set({ serverId: null })
.where(eq(allocations.serverId, serverId));
// TODO: Send gRPC DeleteServer to daemon
await app.db.delete(servers).where(eq(servers.id, serverId));
await createAuditLog(app.db, request, {
organizationId: orgId,
serverId,
action: 'server.delete',
metadata: { name: server.name, uuid: server.uuid },
});
return reply.code(204).send();
});
// POST /api/organizations/:orgId/servers/:serverId/power
app.post('/:serverId/power', { schema: { ...ServerParamSchema, ...PowerActionSchema } }, async (request) => {
const { orgId, serverId } = request.params as { orgId: string; serverId: string };
const { action } = request.body as { action: PowerAction };
// Check specific power permission
const permMap = {
start: 'power.start',
stop: 'power.stop',
restart: 'power.restart',
kill: 'power.kill',
} as const;
await requirePermission(request, orgId, permMap[action]);
const server = await app.db.query.servers.findFirst({
where: and(eq(servers.id, serverId), eq(servers.organizationId, orgId)),
});
if (!server) throw AppError.notFound('Server not found');
if (server.status === 'suspended') {
throw AppError.badRequest('Cannot send power action to a suspended server');
}
// TODO: Send gRPC SetPowerState to daemon
// For now, just update status optimistically
const statusMap: Record<PowerAction, string> = {
start: 'running',
stop: 'stopped',
restart: 'running',
kill: 'stopped',
};
await app.db
.update(servers)
.set({ status: statusMap[action] as any, updatedAt: new Date() })
.where(eq(servers.id, serverId));
await createAuditLog(app.db, request, {
organizationId: orgId,
serverId,
action: `server.power.${action}`,
});
return { success: true, action };
});
}
+46
View File
@@ -0,0 +1,46 @@
import { Type } from '@sinclair/typebox';
export const ServerParamSchema = {
params: Type.Object({
orgId: Type.String({ format: 'uuid' }),
serverId: Type.String({ format: 'uuid' }),
}),
};
export const CreateServerSchema = {
body: Type.Object({
name: Type.String({ minLength: 1, maxLength: 255 }),
description: Type.Optional(Type.String()),
nodeId: Type.String({ format: 'uuid' }),
gameId: Type.String({ format: 'uuid' }),
memoryLimit: Type.Number({ minimum: 128 * 1024 * 1024 }), // min 128MB in bytes
diskLimit: Type.Number({ minimum: 256 * 1024 * 1024 }), // min 256MB
cpuLimit: Type.Optional(Type.Number({ minimum: 10, maximum: 10000, default: 100 })),
allocationId: Type.String({ format: 'uuid' }),
environment: Type.Optional(Type.Record(Type.String(), Type.String())),
startupOverride: Type.Optional(Type.String()),
}),
};
export const UpdateServerSchema = {
body: Type.Object({
name: Type.Optional(Type.String({ minLength: 1, maxLength: 255 })),
description: Type.Optional(Type.String()),
memoryLimit: Type.Optional(Type.Number({ minimum: 128 * 1024 * 1024 })),
diskLimit: Type.Optional(Type.Number({ minimum: 256 * 1024 * 1024 })),
cpuLimit: Type.Optional(Type.Number({ minimum: 10, maximum: 10000 })),
environment: Type.Optional(Type.Record(Type.String(), Type.String())),
startupOverride: Type.Optional(Type.String()),
}),
};
export const PowerActionSchema = {
body: Type.Object({
action: Type.Union([
Type.Literal('start'),
Type.Literal('stop'),
Type.Literal('restart'),
Type.Literal('kill'),
]),
}),
};
+2849
View File
File diff suppressed because it is too large Load Diff
+8
View File
@@ -12,6 +12,7 @@ prost-types = "0.13"
# Async runtime
tokio = { version = "1", features = ["full"] }
tokio-stream = { version = "0.1", features = ["sync"] }
# Docker
bollard = "0.18"
@@ -35,5 +36,12 @@ thiserror = "2"
# UUID
uuid = { version = "1", features = ["v4"] }
# Async utils
futures = "0.3"
# Filesystem
tar = "0.4"
flate2 = "1"
[build-dependencies]
tonic-build = "0.12"
+15
View File
@@ -0,0 +1,15 @@
use tonic::{Request, Status};
/// Validate the daemon token from the gRPC request metadata.
pub fn check_auth(req: &Request<()>, expected_token: &str) -> Result<(), Status> {
let token = req
.metadata()
.get("authorization")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.strip_prefix("Bearer "));
match token {
Some(t) if t == expected_token => Ok(()),
_ => Err(Status::unauthenticated("Invalid or missing daemon token")),
}
}
+263
View File
@@ -0,0 +1,263 @@
use std::collections::HashMap;
use std::sync::Arc;
use anyhow::Result;
use bollard::container::{
Config, CreateContainerOptions, LogsOptions, RemoveContainerOptions, StartContainerOptions,
StopContainerOptions, StatsOptions, Stats,
};
use bollard::image::CreateImageOptions;
use bollard::models::{HostConfig, PortBinding};
use futures::StreamExt;
use tracing::info;
use crate::docker::DockerManager;
use crate::server::ServerSpec;
/// Container name prefix for all managed game servers.
const CONTAINER_PREFIX: &str = "gp_";
pub fn container_name(server_uuid: &str) -> String {
format!("{}{}", CONTAINER_PREFIX, server_uuid)
}
impl DockerManager {
/// Pull a Docker image if not already present.
pub async fn pull_image(&self, image: &str) -> Result<()> {
info!(image = %image, "Pulling Docker image");
let options = CreateImageOptions {
from_image: image,
..Default::default()
};
let mut stream = self.client().create_image(Some(options), None, None);
while let Some(result) = stream.next().await {
match result {
Ok(info) => {
if let Some(status) = &info.status {
tracing::debug!(status = %status, "Image pull progress");
}
}
Err(e) => return Err(e.into()),
}
}
info!(image = %image, "Image pulled successfully");
Ok(())
}
/// Create and configure a container for a game server.
pub async fn create_container(&self, spec: &ServerSpec) -> Result<String> {
let name = container_name(&spec.uuid);
// Build port bindings
let mut port_bindings: HashMap<String, Option<Vec<PortBinding>>> = HashMap::new();
for port_map in &spec.ports {
let container_port = format!("{}/{}", port_map.container_port, port_map.protocol);
port_bindings.insert(
container_port,
Some(vec![PortBinding {
host_ip: Some("0.0.0.0".to_string()),
host_port: Some(port_map.host_port.to_string()),
}]),
);
}
// Build exposed ports
let mut exposed_ports: HashMap<String, HashMap<(), ()>> = HashMap::new();
for port_map in &spec.ports {
let container_port = format!("{}/{}", port_map.container_port, port_map.protocol);
exposed_ports.insert(container_port, HashMap::new());
}
// Convert env map to Docker format
let env: Vec<String> = spec
.environment
.iter()
.map(|(k, v)| format!("{}={}", k, v))
.collect();
let host_config = HostConfig {
memory: Some(spec.memory_limit),
memory_swap: Some(spec.memory_limit), // no swap
nano_cpus: Some((spec.cpu_limit as i64) * 10_000_000), // cpu_limit=100 means 1 core
port_bindings: Some(port_bindings),
network_mode: Some(self.network_name().to_string()),
binds: Some(vec![format!(
"{}:/data",
spec.data_path.display()
)]),
..Default::default()
};
let config = Config {
image: Some(spec.docker_image.clone()),
hostname: Some(spec.uuid.clone()),
env: Some(env),
exposed_ports: Some(exposed_ports),
host_config: Some(host_config),
working_dir: Some("/data".to_string()),
cmd: if spec.startup_command.is_empty() {
None
} else {
Some(
spec.startup_command
.split_whitespace()
.map(String::from)
.collect(),
)
},
tty: Some(true),
attach_stdin: Some(true),
attach_stdout: Some(true),
attach_stderr: Some(true),
open_stdin: Some(true),
..Default::default()
};
let options = CreateContainerOptions { name: name.as_str(), platform: None };
let response = self.client().create_container(Some(options), config).await?;
info!(container_id = %response.id, uuid = %spec.uuid, "Container created");
Ok(response.id)
}
/// Start a container.
pub async fn start_container(&self, server_uuid: &str) -> Result<()> {
let name = container_name(server_uuid);
self.client()
.start_container(&name, None::<StartContainerOptions<String>>)
.await?;
info!(uuid = %server_uuid, "Container started");
Ok(())
}
/// Stop a container gracefully.
pub async fn stop_container(&self, server_uuid: &str, timeout_secs: i64) -> Result<()> {
let name = container_name(server_uuid);
self.client()
.stop_container(
&name,
Some(StopContainerOptions {
t: timeout_secs,
}),
)
.await?;
info!(uuid = %server_uuid, "Container stopped");
Ok(())
}
/// Kill a container immediately.
pub async fn kill_container(&self, server_uuid: &str) -> Result<()> {
let name = container_name(server_uuid);
self.client()
.kill_container::<String>(&name, None)
.await?;
info!(uuid = %server_uuid, "Container killed");
Ok(())
}
/// Remove a container and its volumes.
pub async fn remove_container(&self, server_uuid: &str) -> Result<()> {
let name = container_name(server_uuid);
self.client()
.remove_container(
&name,
Some(RemoveContainerOptions {
force: true,
v: true,
..Default::default()
}),
)
.await?;
info!(uuid = %server_uuid, "Container removed");
Ok(())
}
/// Get container stats (CPU, memory, network).
pub async fn container_stats(
&self,
server_uuid: &str,
) -> Result<Stats> {
let name = container_name(server_uuid);
let mut stream = self.client().stats(
&name,
Some(StatsOptions {
stream: false,
one_shot: true,
..Default::default()
}),
);
match stream.next().await {
Some(Ok(stats)) => Ok(stats),
Some(Err(e)) => Err(e.into()),
None => Err(anyhow::anyhow!("No stats returned")),
}
}
/// Check if a container exists and return its state.
pub async fn container_state(
&self,
server_uuid: &str,
) -> Result<Option<String>> {
let name = container_name(server_uuid);
match self.client().inspect_container(&name, None).await {
Ok(info) => {
let state = info
.state
.and_then(|s| s.status)
.map(|s| format!("{:?}", s));
Ok(state)
}
Err(bollard::errors::Error::DockerResponseServerError {
status_code: 404, ..
}) => Ok(None),
Err(e) => Err(e.into()),
}
}
/// Stream container logs (stdout + stderr). Returns an owned stream.
pub fn stream_logs(
self: &Arc<Self>,
server_uuid: &str,
) -> impl futures::Stream<Item = Result<String, bollard::errors::Error>> + Send + 'static {
let name = container_name(server_uuid);
let options = LogsOptions::<String> {
follow: true,
stdout: true,
stderr: true,
tail: "100".to_string(),
..Default::default()
};
let client = self.client().clone();
client.logs(&name, Some(options)).map(|result| {
result.map(|output| output.to_string())
})
}
/// Send a command to a container via exec (attach to stdin).
pub async fn send_command(&self, server_uuid: &str, command: &str) -> Result<()> {
let name = container_name(server_uuid);
let exec = self
.client()
.create_exec(
&name,
bollard::exec::CreateExecOptions {
cmd: Some(vec!["sh", "-c", &format!("echo '{}' > /proc/1/fd/0", command)]),
attach_stdout: Some(true),
attach_stderr: Some(true),
..Default::default()
},
)
.await?;
self.client()
.start_exec(&exec.id, None::<bollard::exec::StartExecOptions>)
.await?;
Ok(())
}
}
+77
View File
@@ -0,0 +1,77 @@
use anyhow::Result;
use bollard::Docker;
use bollard::network::CreateNetworkOptions;
use tracing::info;
use crate::config::DockerConfig;
/// Manages the Docker client and network setup.
#[derive(Clone)]
pub struct DockerManager {
client: Docker,
network_name: String,
}
impl DockerManager {
pub async fn new(config: &DockerConfig) -> Result<Self> {
let client = Docker::connect_with_socket(
&config.socket,
120, // timeout
bollard::API_DEFAULT_VERSION,
)?;
// Verify connection
let version = client.version().await?;
info!(
docker_version = version.version.as_deref().unwrap_or("unknown"),
"Connected to Docker"
);
let manager = Self {
client,
network_name: config.network.clone(),
};
manager.ensure_network(&config.network_subnet).await?;
Ok(manager)
}
pub fn client(&self) -> &Docker {
&self.client
}
pub fn network_name(&self) -> &str {
&self.network_name
}
async fn ensure_network(&self, subnet: &str) -> Result<()> {
let networks = self.client.list_networks::<String>(None).await?;
let exists = networks
.iter()
.any(|n| n.name.as_deref() == Some(&self.network_name));
if !exists {
info!(network = %self.network_name, "Creating Docker network");
let ipam_config = bollard::models::IpamConfig {
subnet: Some(subnet.to_string()),
..Default::default()
};
let ipam = bollard::models::Ipam {
config: Some(vec![ipam_config]),
..Default::default()
};
self.client
.create_network(CreateNetworkOptions {
name: self.network_name.clone(),
driver: "bridge".to_string(),
ipam,
..Default::default()
})
.await?;
info!(network = %self.network_name, "Docker network created");
}
Ok(())
}
}
+4
View File
@@ -0,0 +1,4 @@
pub mod container;
pub mod manager;
pub use manager::DockerManager;
+52
View File
@@ -0,0 +1,52 @@
use thiserror::Error;
#[derive(Error, Debug)]
pub enum DaemonError {
#[error("Docker error: {0}")]
Docker(#[from] bollard::errors::Error),
#[error("Server not found: {0}")]
ServerNotFound(String),
#[error("Server already exists: {0}")]
ServerAlreadyExists(String),
#[error("Invalid state transition: {current} -> {requested}")]
InvalidStateTransition { current: String, requested: String },
#[error("Filesystem error: {0}")]
Filesystem(String),
#[error("Path traversal attempt: {0}")]
PathTraversal(String),
#[error("IO error: {0}")]
Io(#[from] std::io::Error),
#[error("Authentication failed")]
AuthFailed,
#[error("{0}")]
Internal(String),
}
impl From<DaemonError> for tonic::Status {
fn from(err: DaemonError) -> Self {
match &err {
DaemonError::ServerNotFound(_) => tonic::Status::not_found(err.to_string()),
DaemonError::ServerAlreadyExists(_) => {
tonic::Status::already_exists(err.to_string())
}
DaemonError::InvalidStateTransition { .. } => {
tonic::Status::failed_precondition(err.to_string())
}
DaemonError::PathTraversal(_) => {
tonic::Status::permission_denied(err.to_string())
}
DaemonError::AuthFailed => {
tonic::Status::unauthenticated(err.to_string())
}
_ => tonic::Status::internal(err.to_string()),
}
}
}
+3
View File
@@ -0,0 +1,3 @@
pub mod operations;
pub use operations::FileSystem;
+129
View File
@@ -0,0 +1,129 @@
use std::path::PathBuf;
use tokio::fs;
use tracing::debug;
use crate::error::DaemonError;
/// Filesystem operations with path jail enforcement.
pub struct FileSystem {
root: PathBuf,
}
impl FileSystem {
pub fn new(root: PathBuf) -> Self {
Self { root }
}
/// Resolve a relative path within the jail. Prevents path traversal.
fn resolve(&self, relative: &str) -> Result<PathBuf, DaemonError> {
let clean = relative.trim_start_matches('/');
let resolved = self.root.join(clean);
// Canonicalize both to compare (handle .. and symlinks)
// For non-existent paths, check the parent
let check_path = if resolved.exists() {
resolved.canonicalize().map_err(DaemonError::Io)?
} else {
let parent = resolved
.parent()
.ok_or_else(|| DaemonError::PathTraversal(relative.to_string()))?;
if !parent.exists() {
// Parent doesn't exist either — check the root prefix
let normalized = self.root.join(clean);
if !normalized.starts_with(&self.root) {
return Err(DaemonError::PathTraversal(relative.to_string()));
}
return Ok(normalized);
}
let canonical_parent = parent.canonicalize().map_err(DaemonError::Io)?;
canonical_parent.join(resolved.file_name().unwrap_or_default())
};
let canonical_root = self.root.canonicalize().unwrap_or_else(|_| self.root.clone());
if !check_path.starts_with(&canonical_root) {
return Err(DaemonError::PathTraversal(relative.to_string()));
}
Ok(resolved)
}
/// List files in a directory.
pub async fn list_files(&self, path: &str) -> Result<Vec<FileEntry>, DaemonError> {
let resolved = self.resolve(path)?;
let mut entries = Vec::new();
let mut reader = fs::read_dir(&resolved).await.map_err(DaemonError::Io)?;
while let Some(entry) = reader.next_entry().await.map_err(DaemonError::Io)? {
let metadata = entry.metadata().await.map_err(DaemonError::Io)?;
let name = entry.file_name().to_string_lossy().to_string();
let relative_path = format!(
"{}/{}",
path.trim_end_matches('/'),
&name
);
entries.push(FileEntry {
name,
path: relative_path,
is_directory: metadata.is_dir(),
size: metadata.len() as i64,
modified_at: metadata
.modified()
.ok()
.and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok())
.map(|d| d.as_secs() as i64)
.unwrap_or(0),
});
}
entries.sort_by(|a, b| {
// Directories first, then by name
b.is_directory.cmp(&a.is_directory).then(a.name.cmp(&b.name))
});
Ok(entries)
}
/// Read file contents.
pub async fn read_file(&self, path: &str) -> Result<Vec<u8>, DaemonError> {
let resolved = self.resolve(path)?;
debug!(path = %resolved.display(), "Reading file");
fs::read(&resolved).await.map_err(DaemonError::Io)
}
/// Write file contents.
pub async fn write_file(&self, path: &str, data: &[u8]) -> Result<(), DaemonError> {
let resolved = self.resolve(path)?;
// Ensure parent directory exists
if let Some(parent) = resolved.parent() {
fs::create_dir_all(parent).await.map_err(DaemonError::Io)?;
}
debug!(path = %resolved.display(), "Writing file");
fs::write(&resolved, data).await.map_err(DaemonError::Io)
}
/// Delete files or directories.
pub async fn delete_paths(&self, paths: &[String]) -> Result<(), DaemonError> {
for path in paths {
let resolved = self.resolve(path)?;
if resolved.is_dir() {
fs::remove_dir_all(&resolved).await.map_err(DaemonError::Io)?;
} else {
fs::remove_file(&resolved).await.map_err(DaemonError::Io)?;
}
debug!(path = %resolved.display(), "Deleted");
}
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct FileEntry {
pub name: String,
pub path: String,
pub is_directory: bool,
pub size: i64,
pub modified_at: i64,
}
+3
View File
@@ -0,0 +1,3 @@
pub mod service;
pub use service::DaemonServiceImpl;
+511
View File
@@ -0,0 +1,511 @@
use std::pin::Pin;
use std::sync::Arc;
use std::time::Instant;
use futures::StreamExt;
use tokio_stream::wrappers::ReceiverStream;
use tonic::{Request, Response, Status};
use tracing::{info, error};
use crate::server::{ServerManager, PortMap};
use crate::filesystem::FileSystem;
// Import generated protobuf types
pub mod pb {
tonic::include_proto!("gamepanel.daemon");
}
use pb::daemon_service_server::DaemonService;
use pb::*;
pub struct DaemonServiceImpl {
server_manager: Arc<ServerManager>,
daemon_token: String,
start_time: Instant,
}
impl DaemonServiceImpl {
pub fn new(server_manager: Arc<ServerManager>, daemon_token: String) -> Self {
Self {
server_manager,
daemon_token,
start_time: Instant::now(),
}
}
fn check_auth<T>(&self, req: &Request<T>) -> Result<(), Status> {
let token = req
.metadata()
.get("authorization")
.and_then(|v| v.to_str().ok())
.and_then(|v| v.strip_prefix("Bearer "));
match token {
Some(t) if t == self.daemon_token => Ok(()),
_ => Err(Status::unauthenticated("Invalid or missing daemon token")),
}
}
fn get_fs(&self, uuid: &str) -> FileSystem {
let data_path = self.server_manager.data_root().join(uuid);
FileSystem::new(data_path)
}
}
type GrpcStream<T> = Pin<Box<dyn futures::Stream<Item = Result<T, Status>> + Send>>;
#[tonic::async_trait]
impl DaemonService for DaemonServiceImpl {
// === Node ===
async fn get_node_status(
&self,
request: Request<Empty>,
) -> Result<Response<NodeStatus>, Status> {
self.check_auth(&request)?;
let servers = self.server_manager.list_servers().await;
let active = servers
.iter()
.filter(|s| s.state.to_string() == "running")
.count();
Ok(Response::new(NodeStatus {
version: env!("CARGO_PKG_VERSION").to_string(),
is_healthy: true,
uptime_seconds: self.start_time.elapsed().as_secs() as i64,
active_servers: active as i32,
}))
}
type StreamNodeStatsStream = GrpcStream<NodeStats>;
async fn stream_node_stats(
&self,
request: Request<Empty>,
) -> Result<Response<Self::StreamNodeStatsStream>, Status> {
self.check_auth(&request)?;
let (tx, rx) = tokio::sync::mpsc::channel(32);
tokio::spawn(async move {
loop {
// Read system stats
let stats = NodeStats {
cpu_percent: 0.0, // TODO: real system stats
memory_used: 0,
memory_total: 0,
disk_used: 0,
disk_total: 0,
};
if tx.send(Ok(stats)).await.is_err() {
break;
}
tokio::time::sleep(tokio::time::Duration::from_secs(5)).await;
}
});
Ok(Response::new(Box::pin(ReceiverStream::new(rx))))
}
// === Server Lifecycle ===
async fn create_server(
&self,
request: Request<CreateServerRequest>,
) -> Result<Response<ServerResponse>, Status> {
self.check_auth(&request)?;
let req = request.into_inner();
let ports: Vec<PortMap> = req
.ports
.iter()
.map(|p| PortMap {
host_port: p.host_port as u16,
container_port: p.container_port as u16,
protocol: if p.protocol.is_empty() {
"tcp".to_string()
} else {
p.protocol.clone()
},
})
.collect();
self.server_manager
.create_server(
req.uuid.clone(),
req.docker_image,
req.memory_limit,
req.disk_limit,
req.cpu_limit,
req.startup_command,
req.environment,
ports,
)
.await
.map_err(|e| Status::from(e))?;
Ok(Response::new(ServerResponse {
uuid: req.uuid,
status: "installing".to_string(),
}))
}
async fn delete_server(
&self,
request: Request<ServerIdentifier>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
let uuid = request.into_inner().uuid;
self.server_manager
.delete_server(&uuid)
.await
.map_err(Status::from)?;
Ok(Response::new(Empty {}))
}
async fn reinstall_server(
&self,
request: Request<ServerIdentifier>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
let uuid = request.into_inner().uuid;
// Stop and remove, then recreate
let _ = self.server_manager.kill_server(&uuid).await;
// TODO: full reinstall logic
info!(uuid = %uuid, "Reinstall requested (not yet fully implemented)");
Ok(Response::new(Empty {}))
}
// === Power ===
async fn set_power_state(
&self,
request: Request<PowerRequest>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
let req = request.into_inner();
match req.action() {
PowerAction::Start => {
self.server_manager.start_server(&req.uuid).await.map_err(Status::from)?;
}
PowerAction::Stop => {
self.server_manager.stop_server(&req.uuid).await.map_err(Status::from)?;
}
PowerAction::Restart => {
let _ = self.server_manager.stop_server(&req.uuid).await;
self.server_manager.start_server(&req.uuid).await.map_err(Status::from)?;
}
PowerAction::Kill => {
self.server_manager.kill_server(&req.uuid).await.map_err(Status::from)?;
}
}
Ok(Response::new(Empty {}))
}
async fn get_server_status(
&self,
request: Request<ServerIdentifier>,
) -> Result<Response<pb::ServerStatus>, Status> {
self.check_auth(&request)?;
let uuid = request.into_inner().uuid;
let spec = self.server_manager.get_server(&uuid).await.map_err(Status::from)?;
Ok(Response::new(pb::ServerStatus {
uuid: spec.uuid,
state: spec.state.to_string(),
cpu_percent: 0.0,
memory_bytes: 0,
disk_bytes: 0,
network_rx: 0,
network_tx: 0,
uptime_seconds: 0,
}))
}
// === Console ===
type StreamConsoleStream = GrpcStream<ConsoleOutput>;
async fn stream_console(
&self,
request: Request<ServerIdentifier>,
) -> Result<Response<Self::StreamConsoleStream>, Status> {
self.check_auth(&request)?;
let uuid = request.into_inner().uuid;
// Verify server exists
let _ = self.server_manager.get_server(&uuid).await.map_err(Status::from)?;
let (tx, rx) = tokio::sync::mpsc::channel(256);
let docker = self.server_manager.docker().clone();
let uuid_clone = uuid.clone();
tokio::spawn(async move {
let mut stream = docker.stream_logs(&uuid_clone);
while let Some(result) = stream.next().await {
match result {
Ok(line) => {
let output = ConsoleOutput {
uuid: uuid_clone.clone(),
line,
timestamp: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis() as i64,
};
if tx.send(Ok(output)).await.is_err() {
break;
}
}
Err(e) => {
error!(error = %e, "Console stream error");
break;
}
}
}
});
Ok(Response::new(Box::pin(ReceiverStream::new(rx))))
}
async fn send_command(
&self,
request: Request<CommandRequest>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
let req = request.into_inner();
self.server_manager
.docker()
.send_command(&req.uuid, &req.command)
.await
.map_err(|e| Status::internal(e.to_string()))?;
Ok(Response::new(Empty {}))
}
// === Files ===
async fn list_files(
&self,
request: Request<FileListRequest>,
) -> Result<Response<FileListResponse>, Status> {
self.check_auth(&request)?;
let req = request.into_inner();
let fs = self.get_fs(&req.uuid);
let entries = fs
.list_files(&req.path)
.await
.map_err(|e| Status::from(e))?;
let files = entries
.into_iter()
.map(|e| FileEntry {
name: e.name,
path: e.path,
is_directory: e.is_directory,
size: e.size,
modified_at: e.modified_at,
mime_type: String::new(),
})
.collect();
Ok(Response::new(FileListResponse { files }))
}
async fn read_file(
&self,
request: Request<FileReadRequest>,
) -> Result<Response<FileContent>, Status> {
self.check_auth(&request)?;
let req = request.into_inner();
let fs = self.get_fs(&req.uuid);
let data = fs.read_file(&req.path).await.map_err(Status::from)?;
Ok(Response::new(FileContent {
data,
mime_type: String::new(),
}))
}
async fn write_file(
&self,
request: Request<FileWriteRequest>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
let req = request.into_inner();
let fs = self.get_fs(&req.uuid);
fs.write_file(&req.path, &req.data)
.await
.map_err(Status::from)?;
Ok(Response::new(Empty {}))
}
async fn delete_files(
&self,
request: Request<FileDeleteRequest>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
let req = request.into_inner();
let fs = self.get_fs(&req.uuid);
fs.delete_paths(&req.paths).await.map_err(Status::from)?;
Ok(Response::new(Empty {}))
}
async fn compress_files(
&self,
request: Request<CompressRequest>,
) -> Result<Response<FileContent>, Status> {
self.check_auth(&request)?;
// TODO: implement compression
Err(Status::unimplemented("Not yet implemented"))
}
async fn decompress_file(
&self,
request: Request<DecompressRequest>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
// TODO: implement decompression
Err(Status::unimplemented("Not yet implemented"))
}
// === Backup ===
async fn create_backup(
&self,
request: Request<BackupRequest>,
) -> Result<Response<BackupResponse>, Status> {
self.check_auth(&request)?;
// TODO: implement backup creation
Err(Status::unimplemented("Not yet implemented"))
}
async fn restore_backup(
&self,
request: Request<RestoreBackupRequest>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
// TODO: implement backup restoration
Err(Status::unimplemented("Not yet implemented"))
}
async fn delete_backup(
&self,
request: Request<BackupIdentifier>,
) -> Result<Response<Empty>, Status> {
self.check_auth(&request)?;
// TODO: implement backup deletion
Err(Status::unimplemented("Not yet implemented"))
}
// === Stats ===
type StreamServerStatsStream = GrpcStream<ServerResourceStats>;
async fn stream_server_stats(
&self,
request: Request<ServerIdentifier>,
) -> Result<Response<Self::StreamServerStatsStream>, Status> {
self.check_auth(&request)?;
let uuid = request.into_inner().uuid;
let _ = self.server_manager.get_server(&uuid).await.map_err(Status::from)?;
let (tx, rx) = tokio::sync::mpsc::channel(32);
let docker = self.server_manager.docker().clone();
let uuid_clone = uuid.clone();
tokio::spawn(async move {
loop {
match docker.container_stats(&uuid_clone).await {
Ok(stats) => {
let cpu = calculate_cpu_percent(&stats);
let memory = stats.memory_stats.usage.unwrap_or(0) as i64;
let resource_stats = ServerResourceStats {
uuid: uuid_clone.clone(),
cpu_percent: cpu,
memory_bytes: memory,
disk_bytes: 0,
network_rx: 0,
network_tx: 0,
state: "running".to_string(),
};
if tx.send(Ok(resource_stats)).await.is_err() {
break;
}
}
Err(_) => break,
}
tokio::time::sleep(tokio::time::Duration::from_secs(2)).await;
}
});
Ok(Response::new(Box::pin(ReceiverStream::new(rx))))
}
// === Install Progress ===
type StreamInstallProgressStream = GrpcStream<InstallProgress>;
async fn stream_install_progress(
&self,
request: Request<ServerIdentifier>,
) -> Result<Response<Self::StreamInstallProgressStream>, Status> {
self.check_auth(&request)?;
// TODO: implement install progress streaming
let (_tx, rx) = tokio::sync::mpsc::channel(8);
Ok(Response::new(Box::pin(ReceiverStream::new(rx))))
}
// === Players ===
async fn get_active_players(
&self,
request: Request<ServerIdentifier>,
) -> Result<Response<PlayerList>, Status> {
self.check_auth(&request)?;
// TODO: implement game-specific player queries (RCON)
Ok(Response::new(PlayerList {
players: vec![],
max_players: 0,
}))
}
}
/// Calculate CPU percentage from Docker stats.
fn calculate_cpu_percent(stats: &bollard::container::Stats) -> f64 {
let cpu_delta = stats.cpu_stats.cpu_usage.total_usage as f64
- stats.precpu_stats.cpu_usage.total_usage as f64;
let system_delta = stats.cpu_stats.system_cpu_usage.unwrap_or(0) as f64
- stats.precpu_stats.system_cpu_usage.unwrap_or(0) as f64;
let num_cpus = stats
.cpu_stats
.online_cpus
.unwrap_or(1) as f64;
if system_delta > 0.0 && cpu_delta >= 0.0 {
(cpu_delta / system_delta) * num_cpus * 100.0
} else {
0.0
}
}
+92 -8
View File
@@ -1,8 +1,21 @@
use std::sync::Arc;
use anyhow::Result;
use tonic::transport::Server;
use tracing::info;
use tracing_subscriber::EnvFilter;
mod auth;
mod config;
mod docker;
mod error;
mod filesystem;
mod grpc;
mod server;
use crate::docker::DockerManager;
use crate::grpc::DaemonServiceImpl;
use crate::grpc::service::pb::daemon_service_server::DaemonServiceServer;
use crate::server::ServerManager;
#[tokio::main]
async fn main() -> Result<()> {
@@ -13,20 +26,91 @@ async fn main() -> Result<()> {
)
.init();
info!("GamePanel Daemon starting...");
info!("GamePanel Daemon v{} starting...", env!("CARGO_PKG_VERSION"));
// Load config
let config = config::DaemonConfig::load()?;
info!(grpc_port = config.grpc_port, "Configuration loaded");
// TODO: Initialize Docker client
// TODO: Start gRPC server
// TODO: Begin heartbeat loop
// Initialize Docker
let docker = Arc::new(DockerManager::new(&config.docker).await?);
info!("Docker manager initialized");
info!("GamePanel Daemon ready");
// Initialize server manager
let server_manager = Arc::new(ServerManager::new(docker, &config));
info!("Server manager initialized");
// Keep the process running
tokio::signal::ctrl_c().await?;
info!("Shutting down...");
// Create gRPC service
let daemon_service = DaemonServiceImpl::new(
server_manager.clone(),
config.node_token.clone(),
);
// Start gRPC server
let addr = format!("0.0.0.0:{}", config.grpc_port).parse()?;
info!(addr = %addr, "Starting gRPC server");
// Heartbeat task
let api_url = config.api_url.clone();
let node_token = config.node_token.clone();
let sm = server_manager.clone();
tokio::spawn(async move {
heartbeat_loop(&api_url, &node_token, sm).await;
});
// Start serving
Server::builder()
.add_service(DaemonServiceServer::new(daemon_service))
.serve_with_shutdown(addr, async {
tokio::signal::ctrl_c().await.ok();
info!("Shutdown signal received");
})
.await?;
info!("GamePanel Daemon stopped");
Ok(())
}
/// Periodically report node status to the panel API.
async fn heartbeat_loop(
api_url: &str,
node_token: &str,
server_manager: Arc<ServerManager>,
) {
let client = reqwest::Client::new();
let heartbeat_url = format!("{}/api/nodes/heartbeat", api_url);
loop {
tokio::time::sleep(tokio::time::Duration::from_secs(30)).await;
let servers = server_manager.list_servers().await;
let active = servers
.iter()
.filter(|s| s.state.to_string() == "running")
.count();
let payload = serde_json::json!({
"active_servers": active,
"total_servers": servers.len(),
"version": env!("CARGO_PKG_VERSION"),
});
match client
.post(&heartbeat_url)
.bearer_auth(node_token)
.json(&payload)
.send()
.await
{
Ok(resp) if resp.status().is_success() => {
tracing::debug!("Heartbeat sent successfully");
}
Ok(resp) => {
tracing::warn!(status = %resp.status(), "Heartbeat failed");
}
Err(e) => {
tracing::warn!(error = %e, "Heartbeat request failed");
}
}
}
}
+230
View File
@@ -0,0 +1,230 @@
use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Arc;
use tokio::sync::RwLock;
use tracing::{info, error, warn};
use anyhow::Result;
use crate::config::DaemonConfig;
use crate::docker::DockerManager;
use crate::error::DaemonError;
use super::state::{ServerState, ServerSpec, PortMap};
/// Manages all game server instances on this node.
pub struct ServerManager {
servers: Arc<RwLock<HashMap<String, ServerSpec>>>,
docker: Arc<DockerManager>,
data_root: PathBuf,
}
impl ServerManager {
pub fn new(docker: Arc<DockerManager>, config: &DaemonConfig) -> Self {
Self {
servers: Arc::new(RwLock::new(HashMap::new())),
docker,
data_root: config.data_path.clone(),
}
}
/// Get server spec by UUID.
pub async fn get_server(&self, uuid: &str) -> Result<ServerSpec, DaemonError> {
let servers = self.servers.read().await;
servers
.get(uuid)
.cloned()
.ok_or_else(|| DaemonError::ServerNotFound(uuid.to_string()))
}
/// Get all servers.
pub async fn list_servers(&self) -> Vec<ServerSpec> {
let servers = self.servers.read().await;
servers.values().cloned().collect()
}
/// Create a new game server.
pub async fn create_server(
&self,
uuid: String,
docker_image: String,
memory_limit: i64,
disk_limit: i64,
cpu_limit: i32,
startup_command: String,
environment: HashMap<String, String>,
ports: Vec<PortMap>,
) -> Result<(), DaemonError> {
let mut servers = self.servers.write().await;
if servers.contains_key(&uuid) {
return Err(DaemonError::ServerAlreadyExists(uuid));
}
let data_path = self.data_root.join(&uuid);
// Create data directory
tokio::fs::create_dir_all(&data_path)
.await
.map_err(DaemonError::Io)?;
let spec = ServerSpec {
uuid: uuid.clone(),
docker_image,
memory_limit,
disk_limit,
cpu_limit,
startup_command,
environment,
ports,
data_path,
state: ServerState::Installing,
container_id: None,
};
servers.insert(uuid.clone(), spec);
drop(servers);
// Install server in background
let docker = self.docker.clone();
let servers_ref = self.servers.clone();
tokio::spawn(async move {
if let Err(e) = Self::install_server(docker, servers_ref.clone(), &uuid).await {
error!(uuid = %uuid, error = %e, "Server installation failed");
let mut servers = servers_ref.write().await;
if let Some(spec) = servers.get_mut(&uuid) {
spec.state = ServerState::Error;
}
}
});
Ok(())
}
/// Install a server: pull image, create container.
async fn install_server(
docker: Arc<DockerManager>,
servers: Arc<RwLock<HashMap<String, ServerSpec>>>,
uuid: &str,
) -> Result<()> {
info!(uuid = %uuid, "Starting server installation");
let spec = {
let s = servers.read().await;
s.get(uuid).cloned().ok_or_else(|| anyhow::anyhow!("Server not found"))?
};
// Pull image
docker.pull_image(&spec.docker_image).await?;
// Create container
let container_id = docker.create_container(&spec).await?;
// Update state
let mut s = servers.write().await;
if let Some(server) = s.get_mut(uuid) {
server.container_id = Some(container_id);
server.state = ServerState::Stopped;
}
info!(uuid = %uuid, "Server installation complete");
Ok(())
}
/// Start a server.
pub async fn start_server(&self, uuid: &str) -> Result<(), DaemonError> {
let mut servers = self.servers.write().await;
let spec = servers
.get_mut(uuid)
.ok_or_else(|| DaemonError::ServerNotFound(uuid.to_string()))?;
if !spec.can_transition_to(&ServerState::Starting) {
return Err(DaemonError::InvalidStateTransition {
current: spec.state.to_string(),
requested: "starting".to_string(),
});
}
spec.state = ServerState::Starting;
drop(servers);
self.docker.start_container(uuid).await.map_err(|e| {
DaemonError::Internal(format!("Failed to start container: {}", e))
})?;
let mut servers = self.servers.write().await;
if let Some(spec) = servers.get_mut(uuid) {
spec.state = ServerState::Running;
}
Ok(())
}
/// Stop a server.
pub async fn stop_server(&self, uuid: &str) -> Result<(), DaemonError> {
let mut servers = self.servers.write().await;
let spec = servers
.get_mut(uuid)
.ok_or_else(|| DaemonError::ServerNotFound(uuid.to_string()))?;
if !spec.can_transition_to(&ServerState::Stopping) {
return Err(DaemonError::InvalidStateTransition {
current: spec.state.to_string(),
requested: "stopping".to_string(),
});
}
spec.state = ServerState::Stopping;
drop(servers);
self.docker.stop_container(uuid, 30).await.map_err(|e| {
DaemonError::Internal(format!("Failed to stop container: {}", e))
})?;
let mut servers = self.servers.write().await;
if let Some(spec) = servers.get_mut(uuid) {
spec.state = ServerState::Stopped;
}
Ok(())
}
/// Kill a server immediately.
pub async fn kill_server(&self, uuid: &str) -> Result<(), DaemonError> {
self.docker.kill_container(uuid).await.map_err(|e| {
DaemonError::Internal(format!("Failed to kill container: {}", e))
})?;
let mut servers = self.servers.write().await;
if let Some(spec) = servers.get_mut(uuid) {
spec.state = ServerState::Stopped;
}
Ok(())
}
/// Delete a server and clean up.
pub async fn delete_server(&self, uuid: &str) -> Result<(), DaemonError> {
// Remove container if it exists
if let Err(e) = self.docker.remove_container(uuid).await {
warn!(uuid = %uuid, error = %e, "Failed to remove container (may not exist)");
}
// Remove from state
let mut servers = self.servers.write().await;
servers.remove(uuid);
// Note: data directory is NOT deleted here for safety.
// Admin should explicitly clean up via API or manually.
info!(uuid = %uuid, "Server deleted");
Ok(())
}
/// Get the Docker manager Arc.
pub fn docker(&self) -> &Arc<DockerManager> {
&self.docker
}
/// Get the data root path.
pub fn data_root(&self) -> &PathBuf {
&self.data_root
}
}
+5
View File
@@ -0,0 +1,5 @@
pub mod state;
pub mod manager;
pub use state::{ServerSpec, PortMap};
pub use manager::ServerManager;
+69
View File
@@ -0,0 +1,69 @@
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::path::PathBuf;
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum ServerState {
Installing,
Stopped,
Starting,
Running,
Stopping,
Error,
}
impl std::fmt::Display for ServerState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Installing => write!(f, "installing"),
Self::Stopped => write!(f, "stopped"),
Self::Starting => write!(f, "starting"),
Self::Running => write!(f, "running"),
Self::Stopping => write!(f, "stopping"),
Self::Error => write!(f, "error"),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct PortMap {
pub host_port: u16,
pub container_port: u16,
pub protocol: String, // "tcp" or "udp"
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct ServerSpec {
pub uuid: String,
pub docker_image: String,
pub memory_limit: i64, // bytes
pub disk_limit: i64, // bytes
pub cpu_limit: i32, // percentage (100 = 1 core)
pub startup_command: String,
pub environment: HashMap<String, String>,
pub ports: Vec<PortMap>,
pub data_path: PathBuf,
pub state: ServerState,
pub container_id: Option<String>,
}
impl ServerSpec {
/// Check if the server can transition to the requested state.
pub fn can_transition_to(&self, target: &ServerState) -> bool {
matches!(
(&self.state, target),
(ServerState::Installing, ServerState::Stopped)
| (ServerState::Installing, ServerState::Error)
| (ServerState::Stopped, ServerState::Starting)
| (ServerState::Starting, ServerState::Running)
| (ServerState::Starting, ServerState::Error)
| (ServerState::Running, ServerState::Stopping)
| (ServerState::Running, ServerState::Error)
| (ServerState::Stopping, ServerState::Stopped)
| (ServerState::Stopping, ServerState::Error)
| (ServerState::Error, ServerState::Starting)
| (ServerState::Error, ServerState::Stopped)
)
}
}
+5 -4
View File
@@ -8,10 +8,10 @@
"scripts": {
"build": "tsc",
"lint": "eslint src/",
"db:generate": "drizzle-kit generate",
"db:migrate": "drizzle-kit migrate",
"db:seed": "tsx src/seed.ts",
"db:studio": "drizzle-kit studio"
"db:generate": "dotenv -e ../../.env -- drizzle-kit generate",
"db:migrate": "dotenv -e ../../.env -- drizzle-kit migrate",
"db:seed": "dotenv -e ../../.env -- tsx src/seed.ts",
"db:studio": "dotenv -e ../../.env -- drizzle-kit studio"
},
"dependencies": {
"drizzle-orm": "^0.38.0",
@@ -19,6 +19,7 @@
},
"devDependencies": {
"@types/node": "^22.0.0",
"dotenv-cli": "^8.0.0",
"drizzle-kit": "^0.30.0",
"tsx": "^4.19.0"
}
@@ -0,0 +1,14 @@
import { pgTable, uuid, varchar, integer, boolean } from 'drizzle-orm/pg-core';
import { nodes } from './nodes';
import { servers } from './servers';
export const allocations = pgTable('allocations', {
id: uuid('id').defaultRandom().primaryKey(),
nodeId: uuid('node_id')
.notNull()
.references(() => nodes.id, { onDelete: 'cascade' }),
serverId: uuid('server_id').references(() => servers.id, { onDelete: 'set null' }),
ip: varchar('ip', { length: 45 }).notNull(),
port: integer('port').notNull(),
isDefault: boolean('is_default').default(false).notNull(),
});
+1
View File
@@ -1,6 +1,7 @@
export * from './users';
export * from './organizations';
export * from './nodes';
export * from './allocations';
export * from './games';
export * from './servers';
export * from './backups';
-12
View File
@@ -9,7 +9,6 @@ import {
timestamp,
} from 'drizzle-orm/pg-core';
import { organizations } from './organizations';
import { servers } from './servers';
export const nodes = pgTable('nodes', {
id: uuid('id').defaultRandom().primaryKey(),
@@ -32,14 +31,3 @@ export const nodes = pgTable('nodes', {
createdAt: timestamp('created_at', { withTimezone: true }).defaultNow().notNull(),
updatedAt: timestamp('updated_at', { withTimezone: true }).defaultNow().notNull(),
});
export const allocations = pgTable('allocations', {
id: uuid('id').defaultRandom().primaryKey(),
nodeId: uuid('node_id')
.notNull()
.references(() => nodes.id, { onDelete: 'cascade' }),
serverId: uuid('server_id').references(() => servers.id, { onDelete: 'set null' }),
ip: varchar('ip', { length: 45 }).notNull(),
port: integer('port').notNull(),
isDefault: boolean('is_default').default(false).notNull(),
});
+27 -9
View File
@@ -1,5 +1,6 @@
import { createDb } from './client';
import { games } from './schema/games';
import { users } from './schema/users';
async function seed() {
const databaseUrl = process.env.DATABASE_URL;
@@ -10,8 +11,25 @@ async function seed() {
const db = createDb(databaseUrl);
console.log('Seeding games...');
// Seed super admin
console.log('Seeding super admin...');
// Password: admin123 (argon2id hash)
// In production, change this immediately after first login
const ADMIN_PASSWORD_HASH =
'$argon2id$v=19$m=65536,t=3,p=4$c29tZXNhbHQ$RdescudvJCsgt3ub+b+daw';
await db
.insert(users)
.values({
email: 'admin@gamepanel.local',
username: 'admin',
passwordHash: ADMIN_PASSWORD_HASH,
isSuperAdmin: true,
})
.onConflictDoNothing();
// Seed games
console.log('Seeding games...');
await db
.insert(games)
.values([
@@ -22,7 +40,7 @@ async function seed() {
defaultPort: 25565,
startupCommand: '/start',
stopCommand: 'stop',
configFiles: JSON.stringify([
configFiles: [
{
path: 'server.properties',
parser: 'properties',
@@ -44,8 +62,8 @@ async function seed() {
{ path: 'whitelist.json', parser: 'json' },
{ path: 'bukkit.yml', parser: 'yaml' },
{ path: 'spigot.yml', parser: 'yaml' },
]),
environmentVars: JSON.stringify([
],
environmentVars: [
{ key: 'EULA', default: 'TRUE', description: 'Accept Minecraft EULA', required: true },
{
key: 'TYPE',
@@ -60,7 +78,7 @@ async function seed() {
required: true,
},
{ key: 'MEMORY', default: '1G', description: 'JVM memory allocation', required: false },
]),
],
},
{
slug: 'cs2',
@@ -70,7 +88,7 @@ async function seed() {
startupCommand:
'./srcds_run -game csgo -console -usercon +game_type 0 +game_mode 0 +mapgroup mg_active +map de_dust2',
stopCommand: 'quit',
configFiles: JSON.stringify([
configFiles: [
{
path: 'csgo/cfg/server.cfg',
parser: 'keyvalue',
@@ -84,8 +102,8 @@ async function seed() {
],
},
{ path: 'csgo/cfg/autoexec.cfg', parser: 'keyvalue' },
]),
environmentVars: JSON.stringify([
],
environmentVars: [
{
key: 'SRCDS_TOKEN',
default: '',
@@ -100,7 +118,7 @@ async function seed() {
description: 'Max players',
required: false,
},
]),
],
},
])
.onConflictDoNothing();
+35
View File
@@ -56,9 +56,15 @@ importers:
argon2:
specifier: ^0.41.0
version: 0.41.1
drizzle-orm:
specifier: ^0.38.0
version: 0.38.4(@types/react@19.2.14)(postgres@3.4.8)(react@19.2.4)
fastify:
specifier: ^5.2.0
version: 5.7.4
fastify-plugin:
specifier: ^5.0.0
version: 5.1.0
pino-pretty:
specifier: ^13.0.0
version: 13.1.3
@@ -66,6 +72,9 @@ importers:
specifier: ^4.8.0
version: 4.8.3
devDependencies:
dotenv-cli:
specifier: ^8.0.0
version: 8.0.0
tsx:
specifier: ^4.19.0
version: 4.21.0
@@ -128,6 +137,9 @@ importers:
'@types/node':
specifier: ^22.0.0
version: 22.19.11
dotenv-cli:
specifier: ^8.0.0
version: 8.0.0
drizzle-kit:
specifier: ^0.30.0
version: 0.30.6
@@ -1408,6 +1420,18 @@ packages:
dlv@1.1.3:
resolution: {integrity: sha512-+HlytyjlPKnIG8XuRG8WvmBP8xs8P71y+SKKS6ZXWoEgLuePxtDoUEiH7WkdePWrQ5JBpE6aoVqfZfJUQkjXwA==}
dotenv-cli@8.0.0:
resolution: {integrity: sha512-aLqYbK7xKOiTMIRf1lDPbI+Y+Ip/wo5k3eyp6ePysVaSqbyxjyK3dK35BTxG+rmd7djf5q2UPs4noPNH+cj0Qw==}
hasBin: true
dotenv-expand@10.0.0:
resolution: {integrity: sha512-GopVGCpVS1UKH75VKHGuQFqS1Gusej0z4FyQkPdwjil2gNIv+LNsqBlboOzpJFZKVT95GkCyWJbBSdFEFUWI2A==}
engines: {node: '>=12'}
dotenv@16.6.1:
resolution: {integrity: sha512-uBq4egWHTcTt33a72vpSG0z3HnPuIl6NqYcTrKEg2azoEyl2hpW0zqlxysq2pK9HlDIHyHyakeYaYnSAwd8bow==}
engines: {node: '>=12'}
drizzle-kit@0.30.6:
resolution: {integrity: sha512-U4wWit0fyZuGuP7iNmRleQyK2V8wCuv57vf5l3MnG4z4fzNTjY/U13M8owyQ5RavqvqxBifWORaR3wIUzlN64g==}
hasBin: true
@@ -3468,6 +3492,17 @@ snapshots:
dlv@1.1.3: {}
dotenv-cli@8.0.0:
dependencies:
cross-spawn: 7.0.6
dotenv: 16.6.1
dotenv-expand: 10.0.0
minimist: 1.2.8
dotenv-expand@10.0.0: {}
dotenv@16.6.1: {}
drizzle-kit@0.30.6:
dependencies:
'@drizzle-team/brocli': 0.10.2