Skip to content

Commit dfd6aef

Browse files
committed
log stream client
1 parent 1fc4981 commit dfd6aef

4 files changed

Lines changed: 78 additions & 18 deletions

File tree

src/daemon.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -154,7 +154,7 @@ export default class Daemon {
154154
return this.server;
155155
}
156156

157-
async handleStreamMessage(msg: DaemonMessage, streamController: ReadableStreamController<any>, signal: AbortSignal) {
157+
async handleStreamMessage(msg: DaemonMessage, streamController: ReadableStreamDefaultController, signal: AbortSignal) {
158158

159159
if (!this.initialized) {
160160
await this.initialize();

src/index.ts

Lines changed: 57 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -166,6 +166,52 @@ class BM2CLI {
166166
return { type: "error", error: "Fetch Error", success: false };
167167
}
168168
}
169+
170+
async getDaemonStream<T>(data: DaemonMessage, callback: (data: T | null) => void) {
171+
try {
172+
173+
await this.startDaemon();
174+
175+
const res = await fetch("http://localhost/command", {
176+
unix: DAEMON_SOCKET,
177+
method: "POST",
178+
headers: { "Content-Type": "application/json" },
179+
body: JSON.stringify(data)
180+
})
181+
182+
if (!res.body) {
183+
console.error("No stream received");
184+
process.exit(1);
185+
}
186+
187+
const reader = res.body.getReader();
188+
const decoder = new TextDecoder();
189+
190+
191+
let buffer = "";
192+
193+
while (true) {
194+
195+
const { value, done } = await reader.read();
196+
197+
if (done) break;
198+
199+
buffer += decoder.decode(value, { stream: true });
200+
201+
// split SSE messages
202+
const parts = buffer.split("\n\n");
203+
buffer = parts.pop()!;
204+
205+
for (const part of parts) {
206+
console.log("part===>", part)
207+
}
208+
}
209+
210+
} catch (e: any) {
211+
console.log("callDaemonCmd:", e, e.stack);
212+
return { type: "error", error: "Fetch Error", success: false };
213+
}
214+
}
169215

170216
// -------------------------------------------------------------------------
171217
// Ecosystem config loader
@@ -570,6 +616,7 @@ class BM2CLI {
570616
}
571617

572618
async cmdLogs(args: string[]) {
619+
573620
let target: string | number = "all";
574621
let lines = 20;
575622
let follow = false;
@@ -593,10 +640,6 @@ class BM2CLI {
593640
i++;
594641
}
595642

596-
console.log("target===>", target);
597-
console.log("lines===>", lines);
598-
console.log("follow===>", follow);
599-
600643
const renderLogs = (logs: LogItem[]) => {
601644
for (let log of logs) {
602645
let line;
@@ -612,10 +655,18 @@ class BM2CLI {
612655

613656
if (follow) {
614657

615-
await this.sendToDaemon({
658+
const opts: DaemonMessage = {
616659
type: "streamLogs",
617660
data: { target },
618-
});
661+
mode: "stream"
662+
}
663+
664+
const callback = (data: any) => {
665+
console.log("data===>", data)
666+
}
667+
668+
const res = await this.getDaemonStream(opts, callback);
669+
619670

620671
} else {
621672

src/log-manager.ts

Lines changed: 16 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -179,11 +179,16 @@ export class LogManager {
179179
return sortedLogs
180180
}
181181

182-
async tailLog(filePath: string, streamController: ReadableStreamController<any>, signal: any): Promise<void> {
182+
async tailLog(
183+
name: string,
184+
id: number,
185+
streamController: ReadableStreamDefaultController,
186+
signal: any
187+
): Promise<void> {
183188

184189
let lastSize = (await Bun.file(filePath).exists()) ? Bun.file(filePath).size : 0;
185190

186-
const watcher = watch(filePath, async () => {
191+
const watcher = setInterval(async () => {
187192

188193
const f = Bun.file(filePath);
189194

@@ -192,12 +197,18 @@ export class LogManager {
192197
const chunk = await f.slice(lastSize, f.size).text();
193198
lastSize = f.size;
194199

195-
chunk.split("\n").filter(Boolean).forEach(streamController.enqueue);
200+
chunk
201+
.split("\n")
202+
.filter(Boolean)
203+
.forEach(line => {
204+
const data = `data: ${JSON.stringify(line)}\n\n`
205+
streamController.enqueue(data)
206+
});
196207

197-
});
208+
}, 2000);
198209

199210
signal?.addEventListener("abort", () => {
200-
watcher.close();
211+
clearInterval(watcher)
201212
});
202213
}
203214

src/process-manager.ts

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -320,7 +320,6 @@ import type { ReadableStreamController } from "bun";
320320
return logs.map((log) => ({ name: c.name, id: c.id, ...log }))
321321
}))).flat();
322322

323-
console.log("results===>", results)
324323

325324
let sortedResults = results
326325
.sort((a, b) => (a.ts || "").localeCompare(b.ts || ""))
@@ -329,15 +328,14 @@ import type { ReadableStreamController } from "bun";
329328
return sortedResults;
330329
}
331330

332-
async streamLogs(target: string | number, streamController: ReadableStreamController<any>, signal: AbortSignal) {
331+
async streamLogs(target: string | number, streamController: ReadableStreamDefaultController, signal: AbortSignal) {
333332

334333
const containers = this.resolveTarget(target);
335334
const lm = this.logManager;
336335

337-
await Promise.all(containers.map(async (c) => {
338-
const logPaths = lm.getLogPaths(c.name, c.id);
339-
return lm.tailLog(logPaths.outFile, streamController, signal);
340-
}))
336+
await Promise.all(containers.map(async (c) => (
337+
lm.tailLog(c.name, c.id, streamController, signal)
338+
)))
341339

342340
}
343341

0 commit comments

Comments
 (0)