Add one-click remote node upgrade via WebSocket (v1.3.3).
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
parent
d2046904a0
commit
6437e42c30
11 changed files with 488 additions and 31 deletions
|
|
@ -64,11 +64,16 @@ export default function NodesPage() {
|
||||||
const [editScanAt, setEditScanAt] = useState("03:30");
|
const [editScanAt, setEditScanAt] = useState("03:30");
|
||||||
const [editRoots, setEditRoots] = useState<ScanRootRow[]>([]);
|
const [editRoots, setEditRoots] = useState<ScanRootRow[]>([]);
|
||||||
const [editBusy, setEditBusy] = useState(false);
|
const [editBusy, setEditBusy] = useState(false);
|
||||||
|
const [upgradingId, setUpgradingId] = useState<string | null>(null);
|
||||||
|
const [targetNodeVersion, setTargetNodeVersion] = useState("1.3.3");
|
||||||
|
|
||||||
function loadNodes() {
|
function loadNodes() {
|
||||||
fetch("/api/v1/admin/nodes", { credentials: "include" })
|
fetch("/api/v1/admin/nodes", { credentials: "include" })
|
||||||
.then((r) => r.json())
|
.then((r) => r.json())
|
||||||
.then((d) => setNodes(d.nodes ?? []));
|
.then((d) => {
|
||||||
|
setNodes(d.nodes ?? []);
|
||||||
|
if (d.targetNodeVersion) setTargetNodeVersion(d.targetNodeVersion);
|
||||||
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
function loadShelves() {
|
function loadShelves() {
|
||||||
|
|
@ -164,24 +169,30 @@ export default function NodesPage() {
|
||||||
}
|
}
|
||||||
|
|
||||||
async function upgradeNode(id: string, name: string) {
|
async function upgradeNode(id: string, name: string) {
|
||||||
const res = await fetch(`/api/v1/admin/nodes/${id}/upgrade-command`, {
|
if (!confirm(`Node "${name}" upgraden naar ${targetNodeVersion}?`)) return;
|
||||||
credentials: "include",
|
setUpgradingId(id);
|
||||||
});
|
try {
|
||||||
if (!res.ok) {
|
const res = await fetch(`/api/v1/admin/nodes/${id}/upgrade`, {
|
||||||
|
method: "POST",
|
||||||
|
credentials: "include",
|
||||||
|
});
|
||||||
const data = await res.json();
|
const data = await res.json();
|
||||||
alert(data.error?.message ?? "Upgrade-commando ophalen mislukt");
|
if (data.manualRequired) {
|
||||||
return;
|
await navigator.clipboard.writeText(data.command);
|
||||||
|
alert(
|
||||||
|
`${data.message ?? "Handmatige upgrade nodig."}\n\nCommando gekopieerd — plak in SSH op de ${data.platform === "synology" ? "Synology" : "Unraid"}-host:\n\n${data.command}`
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (!res.ok) {
|
||||||
|
alert(data.error?.message ?? "Upgrade mislukt");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
alert(data.message ?? `Upgrade naar ${targetNodeVersion} gestart.`);
|
||||||
|
loadNodes();
|
||||||
|
} finally {
|
||||||
|
setUpgradingId(null);
|
||||||
}
|
}
|
||||||
const data = (await res.json()) as {
|
|
||||||
command: string;
|
|
||||||
currentVersion: string | null;
|
|
||||||
targetVersion: string;
|
|
||||||
platform: string;
|
|
||||||
};
|
|
||||||
await navigator.clipboard.writeText(data.command);
|
|
||||||
alert(
|
|
||||||
`Upgrade-commando gekopieerd voor "${name}".\n\nPlak in SSH/terminal op de ${data.platform === "synology" ? "Synology" : "Unraid"}-host:\n\n${data.command}\n\nHuidig: ${data.currentVersion ?? "?"} → ${data.targetVersion}`
|
|
||||||
);
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async function revokeNode(id: string) {
|
async function revokeNode(id: string) {
|
||||||
|
|
@ -536,16 +547,20 @@ export default function NodesPage() {
|
||||||
<td>{node.activeStreams}</td>
|
<td>{node.activeStreams}</td>
|
||||||
<td>
|
<td>
|
||||||
{node.version ?? "-"}
|
{node.version ?? "-"}
|
||||||
{node.version && node.version !== "1.3.2" && (
|
{node.version && node.version !== targetNodeVersion && (
|
||||||
<span style={{ color: "#f59e0b", fontSize: "0.7rem", marginLeft: 4 }}>update</span>
|
<span style={{ color: "#f59e0b", fontSize: "0.7rem", marginLeft: 4 }}>update</span>
|
||||||
)}
|
)}
|
||||||
</td>
|
</td>
|
||||||
<td style={{ display: "flex", gap: "0.5rem", flexWrap: "wrap" }}>
|
<td style={{ display: "flex", gap: "0.5rem", flexWrap: "wrap" }}>
|
||||||
<button onClick={() => openEdit(node)}>Edit</button>
|
<button onClick={() => openEdit(node)}>Edit</button>
|
||||||
<button onClick={() => rescanNode(node.id)}>Rescan</button>
|
<button onClick={() => rescanNode(node.id)}>Rescan</button>
|
||||||
{node.version !== "1.3.2" && (
|
{node.version !== targetNodeVersion && (
|
||||||
<button onClick={() => upgradeNode(node.id, node.name)} style={{ background: "#7c3aed" }}>
|
<button
|
||||||
Upgrade
|
onClick={() => upgradeNode(node.id, node.name)}
|
||||||
|
disabled={upgradingId === node.id}
|
||||||
|
style={{ background: "#7c3aed" }}
|
||||||
|
>
|
||||||
|
{upgradingId === node.id ? "Upgraden…" : "Upgrade"}
|
||||||
</button>
|
</button>
|
||||||
)}
|
)}
|
||||||
<button onClick={() => restartNode(node.id)} style={{ background: "#0e7490" }}>
|
<button onClick={() => restartNode(node.id)} style={{ background: "#0e7490" }}>
|
||||||
|
|
|
||||||
|
|
@ -8,6 +8,13 @@ import {
|
||||||
} from "../security/crypto";
|
} from "../security/crypto";
|
||||||
import { AppError } from "../security/errors";
|
import { AppError } from "../security/errors";
|
||||||
import { requireAdmin } from "../auth/routes";
|
import { requireAdmin } from "../auth/routes";
|
||||||
|
import {
|
||||||
|
MEDIA_NODE_TARGET_VERSION,
|
||||||
|
buildUpgradeCommand,
|
||||||
|
mediaNodeBinaryForArch,
|
||||||
|
nodeSupportsRemoteUpgrade,
|
||||||
|
} from "./version";
|
||||||
|
import { randomUUID } from "crypto";
|
||||||
|
|
||||||
export async function registerNodeRoutes(app: FastifyInstance) {
|
export async function registerNodeRoutes(app: FastifyInstance) {
|
||||||
app.post("/api/v1/nodes/enroll", async (request) => {
|
app.post("/api/v1/nodes/enroll", async (request) => {
|
||||||
|
|
@ -137,6 +144,7 @@ export async function registerNodeRoutes(app: FastifyInstance) {
|
||||||
...serializeNode(n),
|
...serializeNode(n),
|
||||||
scanStatus: scanStatuses[n.id] ?? (n.status === "ONLINE" ? "idle" : "offline"),
|
scanStatus: scanStatuses[n.id] ?? (n.status === "ONLINE" ? "idle" : "offline"),
|
||||||
})),
|
})),
|
||||||
|
targetNodeVersion: MEDIA_NODE_TARGET_VERSION,
|
||||||
};
|
};
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
@ -382,18 +390,70 @@ export async function registerNodeRoutes(app: FastifyInstance) {
|
||||||
...node.moviesPaths,
|
...node.moviesPaths,
|
||||||
...node.seriesPaths,
|
...node.seriesPaths,
|
||||||
];
|
];
|
||||||
const isSynology = paths.some((p) => p.startsWith("/volume1"));
|
|
||||||
const appdata = isSynology ? "/volume1/docker/media-node" : "/mnt/user/appdata/media-node";
|
|
||||||
const masterUrl = loadConfig().PUBLIC_URL.replace(/\/$/, "");
|
const masterUrl = loadConfig().PUBLIC_URL.replace(/\/$/, "");
|
||||||
const command = `curl -fsSL ${masterUrl}/install/upgrade-node.sh | bash -s -- --appdata=${appdata}`;
|
const { command, appdata, platform } = buildUpgradeCommand(masterUrl, paths);
|
||||||
|
|
||||||
return {
|
return {
|
||||||
command,
|
command,
|
||||||
appdata,
|
appdata,
|
||||||
platform: isSynology ? "synology" : "unraid",
|
platform,
|
||||||
currentVersion: node.version,
|
currentVersion: node.version,
|
||||||
targetVersion: "1.3.2",
|
targetVersion: MEDIA_NODE_TARGET_VERSION,
|
||||||
needsUpgrade: node.version !== "1.3.2",
|
needsUpgrade: node.version !== MEDIA_NODE_TARGET_VERSION,
|
||||||
|
};
|
||||||
|
});
|
||||||
|
|
||||||
|
app.post("/api/v1/admin/nodes/:id/upgrade", { preHandler: requireAdmin }, async (request) => {
|
||||||
|
const { id } = request.params as { id: string };
|
||||||
|
const node = await prisma.node.findUnique({
|
||||||
|
where: { id },
|
||||||
|
include: { scanRoots: { select: { path: true } } },
|
||||||
|
});
|
||||||
|
if (!node) throw new AppError("NOT_FOUND", "Node not found", 404);
|
||||||
|
|
||||||
|
const paths = [
|
||||||
|
...node.scanRoots.map((r) => r.path),
|
||||||
|
...node.moviesPaths,
|
||||||
|
...node.seriesPaths,
|
||||||
|
];
|
||||||
|
const masterUrl = loadConfig().PUBLIC_URL.replace(/\/$/, "");
|
||||||
|
const { command, platform } = buildUpgradeCommand(masterUrl, paths);
|
||||||
|
|
||||||
|
if (!nodeSupportsRemoteUpgrade(node.version)) {
|
||||||
|
return {
|
||||||
|
ok: false,
|
||||||
|
manualRequired: true,
|
||||||
|
command,
|
||||||
|
platform,
|
||||||
|
currentVersion: node.version,
|
||||||
|
targetVersion: MEDIA_NODE_TARGET_VERSION,
|
||||||
|
message:
|
||||||
|
"Deze node ondersteunt nog geen remote upgrade. Voer het commando éénmalig handmatig uit; daarna werkt Upgrade vanuit admin.",
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
const { nodeConnectionManager } = await import("../websocket/manager");
|
||||||
|
if (!nodeConnectionManager.isOnline(id)) {
|
||||||
|
throw new AppError("NODE_OFFLINE", "Node is niet verbonden", 503);
|
||||||
|
}
|
||||||
|
|
||||||
|
const binary = mediaNodeBinaryForArch(node.architecture);
|
||||||
|
const upgradeId = randomUUID();
|
||||||
|
const result = await nodeConnectionManager.upgradeAcked(id, {
|
||||||
|
upgradeId,
|
||||||
|
url: `${masterUrl}/install/${binary}`,
|
||||||
|
});
|
||||||
|
|
||||||
|
if (!result.ok) {
|
||||||
|
throw new AppError("UPGRADE_FAILED", result.error ?? "Upgrade mislukt", 500);
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
ok: true,
|
||||||
|
upgradeId,
|
||||||
|
currentVersion: node.version,
|
||||||
|
targetVersion: MEDIA_NODE_TARGET_VERSION,
|
||||||
|
message: "Upgrade voltooid; node herstart.",
|
||||||
};
|
};
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
|
||||||
56
apps/master-api/src/nodes/version.ts
Normal file
56
apps/master-api/src/nodes/version.ts
Normal file
|
|
@ -0,0 +1,56 @@
|
||||||
|
/** Target media-node version served from /install/media-node-linux-* */
|
||||||
|
export const MEDIA_NODE_TARGET_VERSION =
|
||||||
|
process.env.MEDIA_NODE_VERSION?.trim() || "1.3.3";
|
||||||
|
|
||||||
|
/** Minimum node version that handles UPGRADE over WebSocket */
|
||||||
|
const REMOTE_UPGRADE_MIN_VERSION = "1.3.3";
|
||||||
|
|
||||||
|
function parseVersion(version: string): [number, number, number] {
|
||||||
|
const parts = version.trim().split(".");
|
||||||
|
return [
|
||||||
|
Number.parseInt(parts[0] ?? "0", 10) || 0,
|
||||||
|
Number.parseInt(parts[1] ?? "0", 10) || 0,
|
||||||
|
Number.parseInt(parts[2] ?? "0", 10) || 0,
|
||||||
|
];
|
||||||
|
}
|
||||||
|
|
||||||
|
export function versionGte(version: string | null | undefined, min: string): boolean {
|
||||||
|
if (!version?.trim()) return false;
|
||||||
|
const [a0, a1, a2] = parseVersion(version);
|
||||||
|
const [b0, b1, b2] = parseVersion(min);
|
||||||
|
if (a0 !== b0) return a0 > b0;
|
||||||
|
if (a1 !== b1) return a1 > b1;
|
||||||
|
return a2 >= b2;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function nodeSupportsRemoteUpgrade(version: string | null | undefined): boolean {
|
||||||
|
return versionGte(version, REMOTE_UPGRADE_MIN_VERSION);
|
||||||
|
}
|
||||||
|
|
||||||
|
export function mediaNodeBinaryForArch(architecture: string | null | undefined): string {
|
||||||
|
const arch = (architecture ?? "amd64").toLowerCase();
|
||||||
|
if (arch === "arm64" || arch === "aarch64") {
|
||||||
|
return "media-node-linux-arm64";
|
||||||
|
}
|
||||||
|
return "media-node-linux-amd64";
|
||||||
|
}
|
||||||
|
|
||||||
|
export function inferNodeAppdata(paths: string[]): {
|
||||||
|
appdata: string;
|
||||||
|
platform: "synology" | "unraid";
|
||||||
|
} {
|
||||||
|
const isSynology = paths.some((p) => p.startsWith("/volume1"));
|
||||||
|
return {
|
||||||
|
appdata: isSynology ? "/volume1/docker/media-node" : "/mnt/user/appdata/media-node",
|
||||||
|
platform: isSynology ? "synology" : "unraid",
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
export function buildUpgradeCommand(
|
||||||
|
masterUrl: string,
|
||||||
|
paths: string[]
|
||||||
|
): { command: string; appdata: string; platform: "synology" | "unraid" } {
|
||||||
|
const { appdata, platform } = inferNodeAppdata(paths);
|
||||||
|
const command = `curl -fsSL ${masterUrl}/install/upgrade-node.sh | bash -s -- --appdata=${appdata}`;
|
||||||
|
return { command, appdata, platform };
|
||||||
|
}
|
||||||
|
|
@ -19,6 +19,8 @@ import type {
|
||||||
ImportMediaPayload,
|
ImportMediaPayload,
|
||||||
WriteSubtitleAckPayload,
|
WriteSubtitleAckPayload,
|
||||||
WriteSubtitlePayload,
|
WriteSubtitlePayload,
|
||||||
|
UpgradeAckPayload,
|
||||||
|
UpgradePayload,
|
||||||
} from "@media-cluster/protocol";
|
} from "@media-cluster/protocol";
|
||||||
|
|
||||||
interface NodeConnection {
|
interface NodeConnection {
|
||||||
|
|
@ -43,6 +45,11 @@ interface PendingWriteSubtitleAck {
|
||||||
timer: ReturnType<typeof setTimeout>;
|
timer: ReturnType<typeof setTimeout>;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
interface PendingUpgradeAck {
|
||||||
|
resolve: (result: UpgradeAckPayload) => void;
|
||||||
|
timer: ReturnType<typeof setTimeout>;
|
||||||
|
}
|
||||||
|
|
||||||
interface SyncState {
|
interface SyncState {
|
||||||
syncId: string;
|
syncId: string;
|
||||||
fileIds: Set<string>;
|
fileIds: Set<string>;
|
||||||
|
|
@ -56,6 +63,7 @@ class NodeConnectionManager {
|
||||||
private pendingSessionAcks = new Map<string, PendingAck>();
|
private pendingSessionAcks = new Map<string, PendingAck>();
|
||||||
private pendingImportAcks = new Map<string, PendingImportAck>();
|
private pendingImportAcks = new Map<string, PendingImportAck>();
|
||||||
private pendingWriteSubtitleAcks = new Map<string, PendingWriteSubtitleAck>();
|
private pendingWriteSubtitleAcks = new Map<string, PendingWriteSubtitleAck>();
|
||||||
|
private pendingUpgradeAcks = new Map<string, PendingUpgradeAck>();
|
||||||
private syncState = new Map<string, SyncState>();
|
private syncState = new Map<string, SyncState>();
|
||||||
private syncQueues = new Map<string, Promise<void>>();
|
private syncQueues = new Map<string, Promise<void>>();
|
||||||
|
|
||||||
|
|
@ -157,6 +165,9 @@ class NodeConnectionManager {
|
||||||
case "WRITE_SUBTITLE_ACK":
|
case "WRITE_SUBTITLE_ACK":
|
||||||
this.handleWriteSubtitleAck(message.payload as WriteSubtitleAckPayload);
|
this.handleWriteSubtitleAck(message.payload as WriteSubtitleAckPayload);
|
||||||
break;
|
break;
|
||||||
|
case "UPGRADE_ACK":
|
||||||
|
this.handleUpgradeAck(message.payload as UpgradeAckPayload);
|
||||||
|
break;
|
||||||
default:
|
default:
|
||||||
this.send(socket, createMessage("ERROR", { message: `Unknown type: ${message.type}` }));
|
this.send(socket, createMessage("ERROR", { message: `Unknown type: ${message.type}` }));
|
||||||
}
|
}
|
||||||
|
|
@ -451,6 +462,14 @@ class NodeConnectionManager {
|
||||||
pending.resolve(payload);
|
pending.resolve(payload);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private handleUpgradeAck(payload: UpgradeAckPayload): void {
|
||||||
|
const pending = this.pendingUpgradeAcks.get(payload.upgradeId);
|
||||||
|
if (!pending) return;
|
||||||
|
clearTimeout(pending.timer);
|
||||||
|
this.pendingUpgradeAcks.delete(payload.upgradeId);
|
||||||
|
pending.resolve(payload);
|
||||||
|
}
|
||||||
|
|
||||||
private async handlePlaybackEnded(payload: PlaybackSessionEndedPayload): Promise<void> {
|
private async handlePlaybackEnded(payload: PlaybackSessionEndedPayload): Promise<void> {
|
||||||
await prisma.playbackSession.updateMany({
|
await prisma.playbackSession.updateMany({
|
||||||
where: { id: payload.sessionId },
|
where: { id: payload.sessionId },
|
||||||
|
|
@ -514,6 +533,36 @@ class NodeConnectionManager {
|
||||||
return this.sendToNode(nodeId, "RESTART", { reason });
|
return this.sendToNode(nodeId, "RESTART", { reason });
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async upgradeAcked(
|
||||||
|
nodeId: string,
|
||||||
|
payload: UpgradePayload,
|
||||||
|
timeoutMs = 180_000
|
||||||
|
): Promise<UpgradeAckPayload> {
|
||||||
|
if (!this.isOnline(nodeId)) {
|
||||||
|
return { upgradeId: payload.upgradeId, ok: false, error: "Node offline" };
|
||||||
|
}
|
||||||
|
|
||||||
|
const acked = new Promise<UpgradeAckPayload>((resolve) => {
|
||||||
|
const timer = setTimeout(() => {
|
||||||
|
this.pendingUpgradeAcks.delete(payload.upgradeId);
|
||||||
|
resolve({ upgradeId: payload.upgradeId, ok: false, error: "Upgrade timeout" });
|
||||||
|
}, timeoutMs);
|
||||||
|
this.pendingUpgradeAcks.set(payload.upgradeId, { resolve, timer });
|
||||||
|
});
|
||||||
|
|
||||||
|
const pushed = this.sendToNode(nodeId, "UPGRADE", payload);
|
||||||
|
if (!pushed) {
|
||||||
|
const pending = this.pendingUpgradeAcks.get(payload.upgradeId);
|
||||||
|
if (pending) {
|
||||||
|
clearTimeout(pending.timer);
|
||||||
|
this.pendingUpgradeAcks.delete(payload.upgradeId);
|
||||||
|
}
|
||||||
|
return { upgradeId: payload.upgradeId, ok: false, error: "Kan niet naar node sturen" };
|
||||||
|
}
|
||||||
|
|
||||||
|
return acked;
|
||||||
|
}
|
||||||
|
|
||||||
isOnline(nodeId: string): boolean {
|
isOnline(nodeId: string): boolean {
|
||||||
const conn = this.connections.get(nodeId);
|
const conn = this.connections.get(nodeId);
|
||||||
return !!conn && conn.socket.readyState === 1;
|
return !!conn && conn.socket.readyState === 1;
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,7 @@ WORKDIR /app
|
||||||
|
|
||||||
FROM golang:1.22-alpine AS node-build
|
FROM golang:1.22-alpine AS node-build
|
||||||
WORKDIR /src
|
WORKDIR /src
|
||||||
ARG MEDIA_NODE_VERSION=1.3.2
|
ARG MEDIA_NODE_VERSION=1.3.3
|
||||||
RUN apk add --no-cache git ca-certificates
|
RUN apk add --no-cache git ca-certificates
|
||||||
COPY node/media-node/ ./
|
COPY node/media-node/ ./
|
||||||
RUN go mod tidy \
|
RUN go mod tidy \
|
||||||
|
|
|
||||||
|
|
@ -61,7 +61,7 @@ mv "$BIN_PATH.new" "$BIN_PATH"
|
||||||
if docker ps -a --format '{{.Names}}' | grep -qx "$CONTAINER"; then
|
if docker ps -a --format '{{.Names}}' | grep -qx "$CONTAINER"; then
|
||||||
echo "Container herstarten..."
|
echo "Container herstarten..."
|
||||||
docker restart "$CONTAINER" >/dev/null
|
docker restart "$CONTAINER" >/dev/null
|
||||||
echo "Klaar. Controleer versie in Admin → Nodes (verwacht 1.3.2+)."
|
echo "Klaar. Controleer versie in Admin → Nodes (verwacht ${MEDIA_NODE_VERSION:-1.3.3}+)."
|
||||||
else
|
else
|
||||||
echo "Waarschuwing: container '$CONTAINER' niet gevonden — binary is wel bijgewerkt."
|
echo "Waarschuwing: container '$CONTAINER' niet gevonden — binary is wel bijgewerkt."
|
||||||
fi
|
fi
|
||||||
|
|
|
||||||
|
|
@ -20,4 +20,4 @@ Na een Master-redeploy staat de nieuwste binary op de Master. Op de Unraid-host
|
||||||
curl -fsSL https://master.vonas.nl/install/upgrade-node.sh | bash -s -- --appdata=/mnt/user/appdata/media-node
|
curl -fsSL https://master.vonas.nl/install/upgrade-node.sh | bash -s -- --appdata=/mnt/user/appdata/media-node
|
||||||
```
|
```
|
||||||
|
|
||||||
Of via Admin → Nodes → **Upgrade** (kopieert het commando). Controleer daarna of de versie **1.3.2** is.
|
Of via Admin → Nodes → **Upgrade** (één klik; nodes op 1.3.3+). Oudere nodes krijgen éénmalig een SSH-commando. Controleer daarna of de versie **1.3.3** is.
|
||||||
|
|
|
||||||
|
|
@ -24,7 +24,7 @@ import (
|
||||||
"github.com/sthmedia/media-node/internal/streaming"
|
"github.com/sthmedia/media-node/internal/streaming"
|
||||||
)
|
)
|
||||||
|
|
||||||
var version = "1.3.2"
|
var version = "1.3.3"
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
if len(os.Args) < 2 {
|
if len(os.Args) < 2 {
|
||||||
|
|
|
||||||
|
|
@ -390,6 +390,8 @@ func (c *Client) handleMessage(msgType string, payload json.RawMessage) {
|
||||||
}
|
}
|
||||||
case "RESTART":
|
case "RESTART":
|
||||||
go c.doRestart()
|
go c.doRestart()
|
||||||
|
case "UPGRADE":
|
||||||
|
go c.doUpgrade(payload)
|
||||||
case "IMPORT_MEDIA":
|
case "IMPORT_MEDIA":
|
||||||
go c.handleImportMedia(payload)
|
go c.handleImportMedia(payload)
|
||||||
case "WRITE_SUBTITLE":
|
case "WRITE_SUBTITLE":
|
||||||
|
|
|
||||||
261
node/media-node/internal/control/upgrade.go
Normal file
261
node/media-node/internal/control/upgrade.go
Normal file
|
|
@ -0,0 +1,261 @@
|
||||||
|
package control
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
"fmt"
|
||||||
|
"io"
|
||||||
|
"log"
|
||||||
|
"net"
|
||||||
|
"net/http"
|
||||||
|
"os"
|
||||||
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
|
"time"
|
||||||
|
)
|
||||||
|
|
||||||
|
type upgradePayload struct {
|
||||||
|
UpgradeID string `json:"upgradeId"`
|
||||||
|
URL string `json:"url"`
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *Client) doUpgrade(payload json.RawMessage) {
|
||||||
|
var p upgradePayload
|
||||||
|
if err := json.Unmarshal(payload, &p); err != nil || p.URL == "" || p.UpgradeID == "" {
|
||||||
|
c.sendUpgradeAck("", false, "ongeldig upgrade-verzoek")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := c.runUpgrade(p.URL); err != nil {
|
||||||
|
log.Printf("UPGRADE failed: %v", err)
|
||||||
|
c.sendUpgradeAck(p.UpgradeID, false, err.Error())
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Printf("UPGRADE ok, herstarten (upgradeId=%s)", p.UpgradeID)
|
||||||
|
c.sendUpgradeAck(p.UpgradeID, true, "")
|
||||||
|
time.Sleep(500 * time.Millisecond)
|
||||||
|
c.doRestart()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *Client) sendUpgradeAck(upgradeID string, ok bool, errMsg string) {
|
||||||
|
ack := controlMessage("UPGRADE_ACK", map[string]interface{}{
|
||||||
|
"upgradeId": upgradeID,
|
||||||
|
"ok": ok,
|
||||||
|
"version": c.version,
|
||||||
|
"error": errMsg,
|
||||||
|
})
|
||||||
|
if err := c.write(ack); err != nil {
|
||||||
|
log.Printf("failed to send UPGRADE_ACK: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *Client) runUpgrade(downloadURL string) error {
|
||||||
|
dataDir := strings.TrimSpace(os.Getenv("MEDIA_NODE_DATA"))
|
||||||
|
if dataDir == "" {
|
||||||
|
dataDir = "/var/lib/media-node"
|
||||||
|
}
|
||||||
|
stagingDir := filepath.Join(dataDir, "upgrade")
|
||||||
|
stagingFile := filepath.Join(stagingDir, "media-node")
|
||||||
|
if err := os.MkdirAll(stagingDir, 0o755); err != nil {
|
||||||
|
return fmt.Errorf("staging map: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := downloadFile(downloadURL, stagingFile); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := os.Chmod(stagingFile, 0o755); err != nil {
|
||||||
|
return fmt.Errorf("chmod staging: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if _, err := os.Stat("/var/run/docker.sock"); err != nil {
|
||||||
|
return fmt.Errorf("docker.sock niet beschikbaar — gebruik handmatig upgrade-script")
|
||||||
|
}
|
||||||
|
|
||||||
|
containerName := strings.TrimSpace(os.Getenv("MEDIA_NODE_CONTAINER_NAME"))
|
||||||
|
if containerName == "" {
|
||||||
|
containerName = "media-node"
|
||||||
|
}
|
||||||
|
|
||||||
|
binHostDir, dataHostDir, err := dockerContainerBindDirs(containerName)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
stagingHost := filepath.Join(dataHostDir, "upgrade", "media-node")
|
||||||
|
shellCmd := "cp /staging/media-node /target/media-node && chmod +x /target/media-node"
|
||||||
|
binds := []string{
|
||||||
|
binHostDir + ":/target:rw",
|
||||||
|
filepath.Dir(stagingHost) + ":/staging:ro",
|
||||||
|
}
|
||||||
|
if err := dockerRunOnce("alpine:3.20", shellCmd, binds); err != nil {
|
||||||
|
return fmt.Errorf("binary installeren: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func downloadFile(url, dest string) error {
|
||||||
|
log.Printf("UPGRADE: download %s", url)
|
||||||
|
tmp := dest + ".download"
|
||||||
|
client := &http.Client{Timeout: 5 * time.Minute}
|
||||||
|
res, err := client.Get(url)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("download: %w", err)
|
||||||
|
}
|
||||||
|
defer res.Body.Close()
|
||||||
|
if res.StatusCode >= 300 {
|
||||||
|
return fmt.Errorf("download HTTP %d", res.StatusCode)
|
||||||
|
}
|
||||||
|
|
||||||
|
out, err := os.Create(tmp)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("schrijven staging: %w", err)
|
||||||
|
}
|
||||||
|
if _, err := io.Copy(out, res.Body); err != nil {
|
||||||
|
out.Close()
|
||||||
|
_ = os.Remove(tmp)
|
||||||
|
return fmt.Errorf("download schrijven: %w", err)
|
||||||
|
}
|
||||||
|
if err := out.Close(); err != nil {
|
||||||
|
_ = os.Remove(tmp)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if err := os.Rename(tmp, dest); err != nil {
|
||||||
|
_ = os.Remove(tmp)
|
||||||
|
return fmt.Errorf("staging afronden: %w", err)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
type dockerMount struct {
|
||||||
|
Source string `json:"Source"`
|
||||||
|
Destination string `json:"Destination"`
|
||||||
|
}
|
||||||
|
|
||||||
|
type dockerInspect struct {
|
||||||
|
Mounts []dockerMount `json:"Mounts"`
|
||||||
|
}
|
||||||
|
|
||||||
|
func dockerContainerBindDirs(containerName string) (binHostDir, dataHostDir string, err error) {
|
||||||
|
body, err := dockerAPIRequest(http.MethodGet, "/containers/"+containerName+"/json", nil)
|
||||||
|
if err != nil {
|
||||||
|
return "", "", fmt.Errorf("container inspect: %w", err)
|
||||||
|
}
|
||||||
|
var info dockerInspect
|
||||||
|
if err := json.Unmarshal(body, &info); err != nil {
|
||||||
|
return "", "", fmt.Errorf("inspect parse: %w", err)
|
||||||
|
}
|
||||||
|
for _, m := range info.Mounts {
|
||||||
|
switch m.Destination {
|
||||||
|
case "/usr/local/bin/media-node":
|
||||||
|
binHostDir = filepath.Dir(m.Source)
|
||||||
|
case "/var/lib/media-node":
|
||||||
|
dataHostDir = m.Source
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if binHostDir == "" {
|
||||||
|
return "", "", fmt.Errorf("binary-mount niet gevonden op container %s", containerName)
|
||||||
|
}
|
||||||
|
if dataHostDir == "" {
|
||||||
|
return "", "", fmt.Errorf("data-mount niet gevonden op container %s", containerName)
|
||||||
|
}
|
||||||
|
return binHostDir, dataHostDir, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func dockerRunOnce(image, shellCmd string, binds []string) error {
|
||||||
|
createBody := map[string]interface{}{
|
||||||
|
"Image": image,
|
||||||
|
"Cmd": []string{"sh", "-c", shellCmd},
|
||||||
|
"HostConfig": map[string]interface{}{
|
||||||
|
"Binds": binds,
|
||||||
|
"AutoRemove": true,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
raw, err := json.Marshal(createBody)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
created, err := dockerAPIRequest(http.MethodPost, "/containers/create", raw)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("container create: %w", err)
|
||||||
|
}
|
||||||
|
var createdResp struct {
|
||||||
|
ID string `json:"Id"`
|
||||||
|
}
|
||||||
|
if err := json.Unmarshal(created, &createdResp); err != nil || createdResp.ID == "" {
|
||||||
|
return fmt.Errorf("container create parse: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if _, err := dockerAPIRequest(http.MethodPost, "/containers/"+createdResp.ID+"/start", nil); err != nil {
|
||||||
|
return fmt.Errorf("container start: %w", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
waitCtx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
|
||||||
|
defer cancel()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-waitCtx.Done():
|
||||||
|
return fmt.Errorf("container wait timeout")
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
body, err := dockerAPIRequest(http.MethodGet, "/containers/"+createdResp.ID+"/json", nil)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
var state struct {
|
||||||
|
State struct {
|
||||||
|
Status string `json:"Status"`
|
||||||
|
ExitCode int `json:"ExitCode"`
|
||||||
|
Running bool `json:"Running"`
|
||||||
|
FinishedAt string `json:"FinishedAt"`
|
||||||
|
} `json:"State"`
|
||||||
|
}
|
||||||
|
if err := json.Unmarshal(body, &state); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
if state.State.Running {
|
||||||
|
time.Sleep(200 * time.Millisecond)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if state.State.ExitCode != 0 {
|
||||||
|
return fmt.Errorf("install-container exit %d", state.State.ExitCode)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func dockerAPIRequest(method, path string, body []byte) ([]byte, error) {
|
||||||
|
httpc := http.Client{
|
||||||
|
Transport: &http.Transport{
|
||||||
|
DialContext: func(_ context.Context, _, _ string) (net.Conn, error) {
|
||||||
|
return net.Dial("unix", "/var/run/docker.sock")
|
||||||
|
},
|
||||||
|
},
|
||||||
|
Timeout: 2 * time.Minute,
|
||||||
|
}
|
||||||
|
var reqBody io.Reader
|
||||||
|
if body != nil {
|
||||||
|
reqBody = bytes.NewReader(body)
|
||||||
|
}
|
||||||
|
req, err := http.NewRequest(method, "http://localhost"+path, reqBody)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if body != nil {
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
}
|
||||||
|
res, err := httpc.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
defer res.Body.Close()
|
||||||
|
out, err := io.ReadAll(res.Body)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
if res.StatusCode >= 300 {
|
||||||
|
return nil, fmt.Errorf("docker API %s %s: %d %s", method, path, res.StatusCode, strings.TrimSpace(string(out)))
|
||||||
|
}
|
||||||
|
return out, nil
|
||||||
|
}
|
||||||
|
|
@ -23,6 +23,8 @@ export type ControlMessageType =
|
||||||
| "RESCAN"
|
| "RESCAN"
|
||||||
| "CONFIG_UPDATE"
|
| "CONFIG_UPDATE"
|
||||||
| "RESTART"
|
| "RESTART"
|
||||||
|
| "UPGRADE"
|
||||||
|
| "UPGRADE_ACK"
|
||||||
| "IMPORT_MEDIA"
|
| "IMPORT_MEDIA"
|
||||||
| "IMPORT_MEDIA_ACK"
|
| "IMPORT_MEDIA_ACK"
|
||||||
| "WRITE_SUBTITLE"
|
| "WRITE_SUBTITLE"
|
||||||
|
|
@ -114,6 +116,18 @@ export interface RestartPayload {
|
||||||
reason?: string;
|
reason?: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
export interface UpgradePayload {
|
||||||
|
upgradeId: string;
|
||||||
|
url: string;
|
||||||
|
}
|
||||||
|
|
||||||
|
export interface UpgradeAckPayload {
|
||||||
|
upgradeId: string;
|
||||||
|
ok: boolean;
|
||||||
|
version?: string;
|
||||||
|
error?: string;
|
||||||
|
}
|
||||||
|
|
||||||
/** Master → node: move finished Download Station output into 4K library folder. */
|
/** Master → node: move finished Download Station output into 4K library folder. */
|
||||||
export interface ImportMediaPayload {
|
export interface ImportMediaPayload {
|
||||||
jobId: string;
|
jobId: string;
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue