344 lines
14 KiB
JavaScript
344 lines
14 KiB
JavaScript
// 直播 SFU 分流 Demo 前端逻辑。控制面走 JSON 网关(protojson),媒体面走反向代理直连 SFU。
|
||
// 已集成登录与 Casbin 权限:自动附加 Authorization header,未登录重定向到 /login
|
||
window.LiveSFU = (function () {
|
||
const ENUM = { cloudflare: "BACKEND_KIND_CLOUDFLARE", srs: "BACKEND_KIND_SRS" };
|
||
const NAME = { BACKEND_KIND_CLOUDFLARE: "Cloudflare Realtime", BACKEND_KIND_SRS: "SRS" };
|
||
|
||
function getRoom() {
|
||
const p = new URLSearchParams(location.search).get("room");
|
||
return (p && p.trim()) || "demo";
|
||
}
|
||
function shortName(enumStr) { return NAME[enumStr] || enumStr; }
|
||
|
||
// ----- Auth helpers -----
|
||
function getToken() {
|
||
return localStorage.getItem("token") || "";
|
||
}
|
||
function setToken(t) {
|
||
if (t) localStorage.setItem("token", t);
|
||
}
|
||
function clearToken() {
|
||
localStorage.removeItem("token");
|
||
}
|
||
function authHeaders() {
|
||
const h = {};
|
||
const tok = getToken();
|
||
if (tok) h["Authorization"] = "Bearer " + tok;
|
||
return h;
|
||
}
|
||
function handleAuthError(status) {
|
||
if (status === 401) {
|
||
const next = encodeURIComponent(location.pathname + location.search);
|
||
if (!location.pathname.startsWith("/login")) {
|
||
location.href = "/login?next=" + next;
|
||
}
|
||
}
|
||
}
|
||
|
||
async function apiGet(path) {
|
||
const r = await fetch(path, { headers: { ...authHeaders() } });
|
||
if (!r.ok) {
|
||
handleAuthError(r.status);
|
||
throw new Error(path + " -> " + r.status + " " + (await r.text()));
|
||
}
|
||
return r.json();
|
||
}
|
||
async function apiGetWithAuth(path) {
|
||
const r = await fetch(path, { headers: { ...authHeaders() } });
|
||
const text = await r.text();
|
||
if (!r.ok) {
|
||
handleAuthError(r.status);
|
||
throw new Error(text || (path + " -> " + r.status));
|
||
}
|
||
return text ? JSON.parse(text) : {};
|
||
}
|
||
async function apiPost(path, body) {
|
||
const r = await fetch(path, {
|
||
method: "POST",
|
||
headers: { "Content-Type": "application/json", ...authHeaders() },
|
||
body: JSON.stringify(body),
|
||
});
|
||
const text = await r.text();
|
||
if (!r.ok) {
|
||
handleAuthError(r.status);
|
||
throw new Error((text || r.status));
|
||
}
|
||
return text ? JSON.parse(text) : {};
|
||
}
|
||
async function apiPostWithAuth(path, body) {
|
||
return apiPost(path, body);
|
||
}
|
||
|
||
async function login(username, password) {
|
||
const r = await fetch("/api/auth/login", {
|
||
method: "POST",
|
||
headers: { "Content-Type": "application/json" },
|
||
body: JSON.stringify({ username, password }),
|
||
});
|
||
const text = await r.text();
|
||
if (!r.ok) throw new Error(text || ("login -> " + r.status));
|
||
const data = text ? JSON.parse(text) : {};
|
||
if (data.token) setToken(data.token);
|
||
return data;
|
||
}
|
||
async function register(username, password, role) {
|
||
const r = await fetch("/api/auth/register", {
|
||
method: "POST",
|
||
headers: { "Content-Type": "application/json", ...authHeaders() },
|
||
body: JSON.stringify({ username, password, role }),
|
||
});
|
||
const text = await r.text();
|
||
if (!r.ok) throw new Error(text || ("register -> " + r.status));
|
||
const data = text ? JSON.parse(text) : {};
|
||
if (data.token) setToken(data.token);
|
||
return data;
|
||
}
|
||
async function getMe() {
|
||
return apiGetWithAuth("/api/auth/me");
|
||
}
|
||
async function authCheck() {
|
||
try {
|
||
const r = await fetch("/api/auth/check", { headers: { ...authHeaders() } });
|
||
return await r.json();
|
||
} catch(e) {
|
||
return { enabled: false };
|
||
}
|
||
}
|
||
|
||
async function loadConfig(el) {
|
||
try {
|
||
const cfg = await apiGet("/api/config");
|
||
el.innerHTML = "";
|
||
(cfg.backends || []).forEach((b) => {
|
||
const span = document.createElement("span");
|
||
span.className = "badge" + (b.primary ? " primary" : "");
|
||
span.textContent = shortName(b.kind) + (b.primary ? " · 主" : "") + (b.configured ? " · 就绪" : " · 未配置");
|
||
el.appendChild(span);
|
||
});
|
||
const chk = await authCheck();
|
||
const authBadge = document.createElement("span");
|
||
authBadge.className = "badge " + (chk.authenticated ? "ok" : "bad");
|
||
authBadge.textContent = chk.authenticated ? ("已登录:" + chk.username + " / " + chk.role) : "未登录";
|
||
authBadge.style.cursor="pointer";
|
||
authBadge.onclick=()=>location.href="/login";
|
||
authBadge.title="点击去登录";
|
||
el.appendChild(authBadge);
|
||
} catch (e) {
|
||
el.innerHTML = '<span class="badge bad">控制面不可达</span>';
|
||
}
|
||
}
|
||
|
||
function log(el, msg) {
|
||
const t = new Date().toLocaleTimeString();
|
||
el.textContent += "[" + t + "] " + msg + "\\n";
|
||
el.scrollTop = el.scrollHeight;
|
||
}
|
||
|
||
async function publishCF(sessionId, localStream, iceServers) {
|
||
const pc = new RTCPeerConnection({ iceServers });
|
||
localStream.getTracks().forEach((t) => pc.addTrack(t, localStream));
|
||
const offer = await pc.createOffer();
|
||
await pc.setLocalDescription(offer);
|
||
const resp = await fetch("/api/cf/sessions/" + sessionId + "/tracks/new", {
|
||
method: "POST",
|
||
headers: { "Content-Type": "application/json", ...authHeaders() },
|
||
body: JSON.stringify({
|
||
sessionDescription: { type: "offer", sdp: offer.sdp },
|
||
tracks: localStream.getTracks().map((t) => ({ location: "local", trackName: t.kind, kind: t.kind })),
|
||
autoDiscover: true,
|
||
}),
|
||
});
|
||
if (!resp.ok) {
|
||
handleAuthError(resp.status);
|
||
throw new Error("cf publish failed: " + resp.status);
|
||
}
|
||
const data = await resp.json();
|
||
await pc.setRemoteDescription({ type: data.sessionDescription.type, sdp: data.sessionDescription.sdp });
|
||
return pc;
|
||
}
|
||
|
||
async function watchCF(viewerSessionId, publisherSessionId, iceServers, video) {
|
||
const pc = new RTCPeerConnection({ iceServers });
|
||
pc.ontrack = (e) => { video.srcObject = e.streams[0]; };
|
||
const offer = await pc.createOffer();
|
||
await pc.setLocalDescription(offer);
|
||
const resp = await fetch("/api/cf/sessions/" + viewerSessionId + "/tracks/new", {
|
||
method: "POST",
|
||
headers: { "Content-Type": "application/json", ...authHeaders() },
|
||
body: JSON.stringify({
|
||
sessionDescription: { type: "offer", sdp: offer.sdp },
|
||
tracks: [
|
||
{ location: "remote", sessionId: publisherSessionId, trackName: "audio" },
|
||
{ location: "remote", sessionId: publisherSessionId, trackName: "video" },
|
||
],
|
||
autoDiscover: true,
|
||
}),
|
||
});
|
||
if (!resp.ok) {
|
||
handleAuthError(resp.status);
|
||
throw new Error("cf watch failed: " + resp.status);
|
||
}
|
||
const data = await resp.json();
|
||
await pc.setRemoteDescription({ type: data.sessionDescription.type, sdp: data.sessionDescription.sdp });
|
||
return pc;
|
||
}
|
||
|
||
async function publishSRS(room, token, stream, iceServers, localStream) {
|
||
const pc = new RTCPeerConnection({ iceServers });
|
||
localStream.getTracks().forEach((t) => pc.addTrack(t, localStream));
|
||
const offer = await pc.createOffer();
|
||
await pc.setLocalDescription(offer);
|
||
const url = "/rtc/v1/whip/?app=live&stream=" + encodeURIComponent(stream) + "&token=" + encodeURIComponent(token);
|
||
const resp = await fetch(url, { method: "POST", headers: { "Content-Type": "application/sdp", ...authHeaders() }, body: offer.sdp });
|
||
if (!resp.ok) {
|
||
handleAuthError(resp.status);
|
||
throw new Error("SRS publish failed: " + resp.status + " " + (await resp.text()));
|
||
}
|
||
const answer = await resp.text();
|
||
await pc.setRemoteDescription({ type: "answer", sdp: answer });
|
||
return pc;
|
||
}
|
||
|
||
async function watchSRS(room, stream, iceServers, video) {
|
||
const pc = new RTCPeerConnection({ iceServers });
|
||
pc.ontrack = (e) => { video.srcObject = e.streams[0]; };
|
||
const offer = await pc.createOffer();
|
||
await pc.setLocalDescription(offer);
|
||
const url = "/rtc/v1/whep/?app=live&stream=" + encodeURIComponent(stream);
|
||
const resp = await fetch(url, { method: "POST", headers: { "Content-Type": "application/sdp", ...authHeaders() }, body: offer.sdp });
|
||
if (!resp.ok) {
|
||
handleAuthError(resp.status);
|
||
throw new Error("SRS watch failed: " + resp.status + " " + (await resp.text()));
|
||
}
|
||
const answer = await resp.text();
|
||
await pc.setRemoteDescription({ type: "answer", sdp: answer });
|
||
return pc;
|
||
}
|
||
|
||
function findTarget(targets, enumStr) {
|
||
return (targets || []).find((t) => t.backend === enumStr);
|
||
}
|
||
|
||
function initPublish() {
|
||
const room = getRoom();
|
||
document.getElementById("roomName").textContent = room;
|
||
const localVideo = document.getElementById("local");
|
||
const logEl = document.getElementById("log");
|
||
const statusEl = document.getElementById("status");
|
||
const pcs = {};
|
||
authCheck().then(chk=>{
|
||
if (!chk.authenticated) {
|
||
log(logEl, "提示:未登录,发布需要 publisher 权限,将自动跳转登录页");
|
||
} else if (chk.permissions && !chk.permissions["room:publish"]) {
|
||
log(logEl, "提示:当前角色 "+chk.role+" 无发布权限,需要 publisher 或 admin");
|
||
}
|
||
});
|
||
document.getElementById("start").onclick = async () => {
|
||
try {
|
||
const backends = Array.from(document.querySelectorAll('input[name="backend"]:checked')).map((i) => i.value);
|
||
if (!backends.length) return log(logEl, "请至少选择一个分发后端");
|
||
const stream = await navigator.mediaDevices.getUserMedia({ video: true, audio: true });
|
||
localVideo.srcObject = stream;
|
||
for (const b of backends) {
|
||
try {
|
||
const resp = await apiPost("/api/publish", { room, backend: ENUM[b] });
|
||
if (b === "cloudflare") {
|
||
pcs[b] = await publishCF(resp.sessionId, stream, resp.iceServers);
|
||
} else {
|
||
pcs[b] = await publishSRS(room, resp.publishToken, resp.stream, resp.iceServers, stream);
|
||
}
|
||
log(logEl, "已在 " + shortName(resp.target.backend) + " 发布(" + (resp.stream || resp.sessionId) + ")");
|
||
} catch (e) {
|
||
log(logEl, "发布到 " + b + " 失败:" + e.message);
|
||
}
|
||
}
|
||
renderStatus(statusEl, pcs);
|
||
} catch (e) {
|
||
log(logEl, "获取摄像头失败:" + e.message);
|
||
}
|
||
};
|
||
document.getElementById("stop").onclick = async () => {
|
||
for (const b of Object.keys(pcs)) {
|
||
try { await apiPost("/api/stop", { room, backend: ENUM[b] }); } catch (e) {}
|
||
if (pcs[b]) pcs[b].close();
|
||
delete pcs[b];
|
||
}
|
||
if (localVideo.srcObject) localVideo.srcObject.getTracks().forEach((t) => t.stop());
|
||
log(logEl, "已停止推流");
|
||
renderStatus(statusEl, pcs);
|
||
};
|
||
}
|
||
|
||
function renderStatus(el, pcs) {
|
||
const keys = Object.keys(pcs);
|
||
el.innerHTML = keys.length
|
||
? keys.map((k) => '<span class="target ' + k + '">分发中:' + shortName(ENUM[k]) + "</span>").join("")
|
||
: '<span class="muted">未发布</span>';
|
||
}
|
||
|
||
function initWatch() {
|
||
const room = getRoom();
|
||
document.getElementById("roomName").textContent = room;
|
||
const remoteVideo = document.getElementById("remote");
|
||
const logEl = document.getElementById("log");
|
||
const statusEl = document.getElementById("status");
|
||
let pc = null;
|
||
let es = null;
|
||
let watching = false;
|
||
function backend() { return document.querySelector('input[name="watchBackend"]:checked').value; }
|
||
authCheck().then(chk=>{
|
||
if (!chk.authenticated) {
|
||
log(logEl, "提示:未登录,观看需要登录,将自动跳转");
|
||
} else if (chk.permissions && !chk.permissions["room:subscribe"]) {
|
||
log(logEl, "提示:当前角色 "+chk.role+" 无订阅权限");
|
||
}
|
||
});
|
||
async function subscribe() {
|
||
if (watching) return;
|
||
const b = backend();
|
||
try {
|
||
const resp = await apiPost("/api/subscribe", { room, backend: ENUM[b] });
|
||
if (b === "cloudflare") {
|
||
pc = await watchCF(resp.sessionId, resp.publisherSessionId, resp.iceServers, remoteVideo);
|
||
} else {
|
||
pc = await watchSRS(room, resp.stream, resp.iceServers, remoteVideo);
|
||
}
|
||
watching = true;
|
||
log(logEl, "已从 " + shortName(ENUM[b]) + " 拉流");
|
||
statusEl.innerHTML = '<span class="target ' + b + '">观看中:' + shortName(ENUM[b]) + "</span>";
|
||
} catch (e) {
|
||
log(logEl, "订阅失败:" + e.message);
|
||
}
|
||
}
|
||
document.getElementById("start").onclick = () => {
|
||
const token = getToken();
|
||
const esUrl = "/api/room/" + encodeURIComponent(room) + "/events" + (token ? "?token="+encodeURIComponent(token) : "");
|
||
es = new EventSource(esUrl);
|
||
es.addEventListener("room", (ev) => {
|
||
const data = JSON.parse(ev.data);
|
||
const t = findTarget(data.targets, ENUM[backend()]);
|
||
statusEl.innerHTML = t
|
||
? '<span class="target ' + backend() + '">房间在 ' + shortName(ENUM[backend()]) + " 已直播</span>"
|
||
: '<span class="muted">等待 ' + shortName(ENUM[backend()]) + " 推流…</span>";
|
||
if (t && !watching) subscribe();
|
||
});
|
||
es.onerror = (e)=>{
|
||
log(logEl, "房间事件流错误,可能是未登录或权限不足");
|
||
};
|
||
log(logEl, "已连接房间事件流");
|
||
};
|
||
document.getElementById("stop").onclick = () => {
|
||
if (pc) pc.close();
|
||
pc = null;
|
||
if (es) es.close();
|
||
es = null;
|
||
watching = false;
|
||
remoteVideo.srcObject = null;
|
||
statusEl.innerHTML = '<span class="muted">未观看</span>';
|
||
log(logEl, "已停止观看");
|
||
};
|
||
}
|
||
|
||
return { ENUM, getRoom, loadConfig, log, initPublish, initWatch, apiGet, apiPost, apiGetWithAuth, apiPostWithAuth, login, register, getMe, authCheck, getToken, setToken, clearToken, authHeaders };
|
||
})();
|