137 lines
3.7 KiB
TypeScript
137 lines
3.7 KiB
TypeScript
import { Injectable, inject } from '@angular/core';
|
|
import { AuthService } from './auth.service';
|
|
import { EventSourceLike } from './event-status.pure';
|
|
|
|
interface SseMessage {
|
|
event: string;
|
|
data: string;
|
|
}
|
|
|
|
@Injectable({ providedIn: 'root' })
|
|
export class AuthenticatedEventSourceService {
|
|
private readonly auth = inject(AuthService);
|
|
|
|
open(url: string): EventSourceLike {
|
|
const token = this.auth.getAccessToken();
|
|
const headers: Record<string, string> = { Accept: 'text/event-stream' };
|
|
if (token) headers['Authorization'] = `Bearer ${token}`;
|
|
|
|
const listeners = new Map<string, Array<(ev: MessageEvent | Event) => void>>();
|
|
let controller: AbortController | null = new AbortController();
|
|
let closed = false;
|
|
let buffered: SseMessage[] = [];
|
|
|
|
const dispatch = (type: string, ev: MessageEvent | Event) => {
|
|
const arr = listeners.get(type);
|
|
if (!arr) return;
|
|
for (const fn of arr) {
|
|
try {
|
|
fn(ev);
|
|
} catch {
|
|
// ignore listener errors
|
|
}
|
|
}
|
|
};
|
|
|
|
const flushBuffered = () => {
|
|
if (closed) return;
|
|
const items = buffered;
|
|
buffered = [];
|
|
for (const m of items) {
|
|
const me = new MessageEvent(m.event, { data: m.data });
|
|
dispatch('message', me);
|
|
}
|
|
};
|
|
|
|
const fireError = () => {
|
|
if (closed) return;
|
|
dispatch('error', new Event('error'));
|
|
};
|
|
|
|
const fireUnauthorized = () => {
|
|
if (closed) return;
|
|
dispatch('unauthorized', new Event('unauthorized'));
|
|
};
|
|
|
|
const start = async () => {
|
|
try {
|
|
const res = await fetch(url, {
|
|
method: 'GET',
|
|
headers,
|
|
credentials: 'include',
|
|
signal: controller!.signal,
|
|
});
|
|
if (res.status === 401 || res.status === 403) {
|
|
if (!closed) fireUnauthorized();
|
|
return;
|
|
}
|
|
if (!res.ok || !res.body) {
|
|
fireError();
|
|
return;
|
|
}
|
|
const reader = res.body.getReader();
|
|
const decoder = new TextDecoder();
|
|
let buf = '';
|
|
for (;;) {
|
|
const { value, done } = await reader.read();
|
|
if (done) break;
|
|
buf += decoder.decode(value, { stream: true });
|
|
let idx: number;
|
|
while ((idx = buf.indexOf('\n\n')) >= 0) {
|
|
const raw = buf.slice(0, idx);
|
|
buf = buf.slice(idx + 2);
|
|
parseAndDispatch(raw);
|
|
}
|
|
}
|
|
if (buf.trim().length > 0) parseAndDispatch(buf);
|
|
} catch {
|
|
if (!closed) fireError();
|
|
}
|
|
};
|
|
|
|
const parseAndDispatch = (raw: string) => {
|
|
let event = 'message';
|
|
const dataLines: string[] = [];
|
|
for (const line of raw.split('\n')) {
|
|
if (line.startsWith(':')) continue;
|
|
if (line.startsWith('event:')) event = line.slice(6).trim();
|
|
else if (line.startsWith('data:')) dataLines.push(line.slice(5).trim());
|
|
}
|
|
if (dataLines.length === 0) return;
|
|
const payload = dataLines.join('\n');
|
|
if (listeners.has('message')) {
|
|
dispatch('message', new MessageEvent(event, { data: payload }));
|
|
} else {
|
|
buffered.push({ event, data: payload });
|
|
}
|
|
};
|
|
|
|
void start();
|
|
|
|
return {
|
|
addEventListener(type, listener) {
|
|
const arr = listeners.get(type) ?? [];
|
|
arr.push(listener);
|
|
listeners.set(type, arr);
|
|
if (type === 'message' && buffered.length > 0) {
|
|
queueMicrotask(flushBuffered);
|
|
}
|
|
},
|
|
close() {
|
|
if (closed) return;
|
|
closed = true;
|
|
if (controller) {
|
|
try {
|
|
controller.abort();
|
|
} catch {
|
|
// ignore
|
|
}
|
|
controller = null;
|
|
}
|
|
listeners.clear();
|
|
buffered = [];
|
|
},
|
|
};
|
|
}
|
|
}
|