From 9b84bc9c702f58a8467dba6af616cce030d5a614 Mon Sep 17 00:00:00 2001 From: Zhiqiang ZHOU Date: Sun, 26 Jul 2026 21:03:38 -0700 Subject: [PATCH] feat(web): add GCP log import flow --- CLAUDE.md | 1 + cmd/lapp/workspace_import.go | 7 +- docs/web-app-plan.md | 4 + frontend/src/import-panel.tsx | 326 ++++++++++++++++++++++++++++++ frontend/src/main.tsx | 42 +++- frontend/src/store.ts | 145 +++++++++++-- frontend/src/styles.css | 225 +++++++++++++++++++++ go.mod | 5 +- go.sum | 24 --- pkg/gcplog/fetch.go | 122 ++++++++--- pkg/webapp/import_service_test.go | 23 ++- pkg/webapp/server.go | 7 +- pkg/workspace/import.go | 190 +++++++++++++---- pkg/workspace/import_test.go | 51 +++-- 14 files changed, 1013 insertions(+), 159 deletions(-) create mode 100644 frontend/src/import-panel.tsx diff --git a/CLAUDE.md b/CLAUDE.md index e0d9c40..f20a0e0 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -54,6 +54,7 @@ integration_test/ Integration tests against Loghub-2.0 datasets `add-log` is a pure copy into `logs/` and never triggers discovery. `workspace import gcp` pulls a snapshot from GCP Cloud Logging through ADC credentials, lands it as enveloped NDJSON in `logs/`, and records provenance under `import-runs//record.json`; it never triggers discovery either. Each `workspace discover` starts a DiscoveryRun: reads ALL files in `logs/`, runs fresh Drain + semantic labeling, and writes run-scoped `patterns/` and `notes/`. When `lapp web` starts, it marks any previous `QUEUED` or `RUNNING` DiscoveryRuns and ImportRuns as failed because those local workers no longer exist. DiscoveryRun records persist structured `progress` and `error` fields; frontend code renders those facts into user-facing text. +The web Imports tab starts GCP imports through standard Operations, shows ImportRun history and final facts, and fills project and filter fields from recent queries across all workspaces. ``` workspace.Discover(ctx, cfg) diff --git a/cmd/lapp/workspace_import.go b/cmd/lapp/workspace_import.go index 1b7794b..1f7b6d2 100644 --- a/cmd/lapp/workspace_import.go +++ b/cmd/lapp/workspace_import.go @@ -139,19 +139,18 @@ func runWorkspaceImportGCP(cmd *cobra.Command, _ []string) error { From: from, To: to, Limit: importGCPLimit, - Fetcher: func(fetchCtx context.Context, req workspace.ImportRequest) (workspace.ImportFetchResult, error) { - fetched, err := gcplog.FetchLines(fetchCtx, gcplog.FetchRequest{ + Fetcher: func(fetchCtx context.Context, req workspace.ImportRequest, writeLine workspace.ImportLineWriter) (workspace.ImportFetchResult, error) { + fetched, err := gcplog.Fetch(fetchCtx, gcplog.FetchRequest{ Project: req.Project, Filter: req.Filter, From: req.From, To: req.To, Limit: req.Limit, - }) + }, writeLine) if err != nil { return workspace.ImportFetchResult{}, err } return workspace.ImportFetchResult{ - Lines: fetched.Lines, Truncated: fetched.Truncated, }, nil }, diff --git a/docs/web-app-plan.md b/docs/web-app-plan.md index f588477..da3f331 100644 --- a/docs/web-app-plan.md +++ b/docs/web-app-plan.md @@ -42,6 +42,10 @@ LAPP Web is a local web app for working with local LAPP workspaces. It starts fr - Show latest discovery status and progress. - Disable log upload, log deletion, and starting another discovery while a discovery run is running. +### GCP Imports + +The Imports tab starts GCP imports with project, filter, time range, and limit fields. Quick time ranges cover the common one hour, six hour, and twenty four hour cases. The frontend polls the standard Operation while work is active. ImportRun records remain the source for workspace history, provenance, final counts, truncation, and structured errors. Recent project and filter pairs come from ImportRun records across every workspace. Choosing one fills both fields without creating another stored resource. + ### Patterns - List patterns from the selected or latest successful discovery run. diff --git a/frontend/src/import-panel.tsx b/frontend/src/import-panel.tsx new file mode 100644 index 0000000..dd73e6a --- /dev/null +++ b/frontend/src/import-panel.tsx @@ -0,0 +1,326 @@ +import { timestampDate } from "@bufbuild/protobuf/wkt"; +import { AlertTriangle, CloudDownload, History } from "lucide-react"; +import { useState } from "react"; +import { ImportRun, ImportRunState } from "./gen/lapp/web/v1/web_pb"; +import { useAppStore } from "./store"; + +export function ImportPanel({ hasActiveRun }: { hasActiveRun: boolean }) { + const importRuns = useAppStore((state) => state.importRuns); + const recentImportQueries = useAppStore((state) => state.recentImportQueries); + const busy = useAppStore((state) => state.busy); + const runAction = useAppStore((state) => state.runAction); + const startImport = useAppStore((state) => state.startImport); + const [project, setProject] = useState(""); + const [filter, setFilter] = useState(""); + const [range, setRange] = useState(() => rangeForHours(1)); + const [selectedHours, setSelectedHours] = useState(1); + const [limit, setLimit] = useState("100000"); + const latestRun = importRuns[0]; + const disabled = busy || hasActiveRun; + + const useQuickRange = (hours: number) => { + setSelectedHours(hours); + setRange(rangeForHours(hours)); + }; + + const submit = async () => { + const from = new Date(range.from); + const to = new Date(range.to); + const parsedLimit = Number(limit); + if (!project.trim()) { + throw new Error("Project is required."); + } + if (Number.isNaN(from.getTime()) || Number.isNaN(to.getTime())) { + throw new Error("From and to are required."); + } + if (from > to) { + throw new Error("From must not be after to."); + } + if (!Number.isInteger(parsedLimit) || parsedLimit < 1) { + throw new Error("Limit must be a positive whole number."); + } + + await startImport({ + project: project.trim(), + filter, + from, + to, + limit: parsedLimit + }); + }; + + return ( +
+
+
{ + event.preventDefault(); + void runAction(submit); + }} + > +
+
+

Import GCP logs

+

The server uses local Application Default Credentials.

+
+ +
+ + + +
+ + +
+ + + +
+ Time range +
+ {[1, 6, 24].map((hours) => ( + + ))} +
+
+ + +
+
+ +
+ +
+
+ +
+
+
+

Latest import

+

Current state and final result.

+
+ +
+ {latestRun ? ( + <> +
+ + {formatTimestamp(latestRun.startedAt)} +
+

{importRunMessage(latestRun)}

+
+
Project
+
{latestRun.project}
+
Range
+
{formatRange(latestRun)}
+
Limit
+
{latestRun.limit.toLocaleString()}
+
+ {latestRun.truncated && ( +
+ + Results reached the limit and were truncated. +
+ )} + + ) : ( +

No imports for this workspace.

+ )} +
+
+ +
+
+
+

Import history

+

Newest runs appear first.

+
+
+ {importRuns.length > 0 ? ( +
+ + + + + + + + + + + + + {importRuns.map((run) => ( + + + + + + + + + ))} + +
StartedProject and filterRangeStateEntriesFile
{formatTimestamp(run.startedAt)} + {run.project} + {run.filter || "All logs"} + {formatRange(run)} + + {run.error && ( + + {run.error.code}: {run.error.message} + + )} + + {run.state === ImportRunState.SUCCEEDED ? run.entryCount.toLocaleString() : ""} + {run.truncated && Truncated} + {run.logFileName}
+
+ ) : ( +

No import history.

+ )} +
+
+ ); +} + +function StateBadge({ state }: { state: ImportRunState }) { + return {importStateLabel(state)}; +} + +function importRunMessage(run: ImportRun) { + switch (run.state) { + case ImportRunState.QUEUED: + return "Waiting to start."; + case ImportRunState.RUNNING: + return "Reading matching log entries from GCP."; + case ImportRunState.SUCCEEDED: + return `${run.entryCount.toLocaleString()} entries written to ${run.logFileName}. Start discovery when ready.`; + case ImportRunState.FAILED: + return run.error ? `${run.error.code}: ${run.error.message}` : "Import failed without an error message."; + default: + return "Waiting for state."; + } +} + +function importStateLabel(state: ImportRunState) { + return ImportRunState[state]?.toLowerCase() ?? "unknown"; +} + +function importStateClass(state: ImportRunState) { + switch (state) { + case ImportRunState.SUCCEEDED: + return "succeeded"; + case ImportRunState.FAILED: + return "failed"; + case ImportRunState.QUEUED: + case ImportRunState.RUNNING: + return "running"; + default: + return ""; + } +} + +function recentQueryLabel(project: string, filter: string) { + return filter ? `${project} · ${filter}` : `${project} · All logs`; +} + +function formatRange(run: ImportRun) { + return `${formatTimestamp(run.from)} to ${formatTimestamp(run.to)}`; +} + +function formatTimestamp(value: ImportRun["from"]) { + return value ? timestampDate(value).toLocaleString() : ""; +} + +function rangeForHours(hours: number) { + const to = new Date(); + const from = new Date(to.getTime() - hours * 60 * 60 * 1000); + return { + from: datetimeLocalValue(from), + to: datetimeLocalValue(to) + }; +} + +function datetimeLocalValue(value: Date) { + const local = new Date(value.getTime() - value.getTimezoneOffset() * 60 * 1000); + return local.toISOString().slice(0, 16); +} diff --git a/frontend/src/main.tsx b/frontend/src/main.tsx index b670f45..6ed661f 100644 --- a/frontend/src/main.tsx +++ b/frontend/src/main.tsx @@ -2,6 +2,7 @@ import React, { useEffect, useMemo, useRef } from "react"; import { createRoot } from "react-dom/client"; import { AlertTriangle, + CloudDownload, FileText, FolderOpen, Play, @@ -14,9 +15,11 @@ import { DiscoveryRun, DiscoveryRunState, DiscoveryStep, + ImportRunState, Pattern, WorkspaceStatus } from "./gen/lapp/web/v1/web_pb"; +import { ImportPanel } from "./import-panel"; import { useAppStore } from "./store"; import "./styles.css"; @@ -24,7 +27,9 @@ function App() { const workspaces = useAppStore((state) => state.workspaces); const selectedWorkspaceName = useAppStore((state) => state.selectedWorkspaceName); const logFiles = useAppStore((state) => state.logFiles); + const importRuns = useAppStore((state) => state.importRuns); const discoveryRuns = useAppStore((state) => state.discoveryRuns); + const activeOperationName = useAppStore((state) => state.activeOperationName); const selectedRunName = useAppStore((state) => state.selectedRunName); const selectedPatternName = useAppStore((state) => state.selectedPatternName); const patterns = useAppStore((state) => state.patterns); @@ -42,8 +47,9 @@ function App() { const runAction = useAppStore((state) => state.runAction); const refreshWorkspaces = useAppStore((state) => state.refreshWorkspaces); const refreshWorkspaceDetails = useAppStore((state) => state.refreshWorkspaceDetails); + const refreshRecentImportQueries = useAppStore((state) => state.refreshRecentImportQueries); + const refreshActiveOperation = useAppStore((state) => state.refreshActiveOperation); const refreshResults = useAppStore((state) => state.refreshResults); - const refreshSelectedWorkspace = useAppStore((state) => state.refreshSelectedWorkspace); const createWorkspace = useAppStore((state) => state.createWorkspace); const deleteWorkspace = useAppStore((state) => state.deleteWorkspace); const uploadFiles = useAppStore((state) => state.uploadFiles); @@ -64,10 +70,16 @@ function App() { [patterns, selectedPatternName] ); const isDiscovering = selectedWorkspace?.status === WorkspaceStatus.DISCOVERING; + const isImporting = importRuns.some( + (run) => run.state === ImportRunState.QUEUED || run.state === ImportRunState.RUNNING + ); + const hasActiveRun = Boolean(activeOperationName) || isDiscovering || isImporting; useEffect(() => { - void runAction(refreshWorkspaces); - }, [refreshWorkspaces, runAction]); + void runAction(async () => { + await Promise.all([refreshWorkspaces(), refreshRecentImportQueries()]); + }); + }, [refreshRecentImportQueries, refreshWorkspaces, runAction]); useEffect(() => { void refreshWorkspaceDetails().catch((error: unknown) => { @@ -86,18 +98,18 @@ function App() { }, [refreshResults, selectedRunName, runAction]); useEffect(() => { - if (!selectedRun || selectedRun.state !== DiscoveryRunState.RUNNING) { + if (!activeOperationName) { return; } const timer = window.setInterval(() => { - void refreshSelectedWorkspace().catch((error: unknown) => { + void refreshActiveOperation().catch((error: unknown) => { void runAction(async () => { throw error; }); }); }, 1500); return () => window.clearInterval(timer); - }, [refreshSelectedWorkspace, runAction, selectedRun?.name, selectedRun?.state]); + }, [activeOperationName, refreshActiveOperation, runAction]); return (
@@ -153,14 +165,14 @@ function App() {
- + @@ -232,7 +248,7 @@ function App() { key={file.name} className="icon-button danger" onClick={() => void runAction(() => deleteLogFile(file))} - disabled={busy || isDiscovering} + disabled={busy || hasActiveRun} title="Delete log file" > @@ -242,6 +258,10 @@ function App() { )} + {activeTab === "imports" && ( + + )} + {activeTab === "patterns" && (
diff --git a/frontend/src/store.ts b/frontend/src/store.ts index 7404862..5c05eaa 100644 --- a/frontend/src/store.ts +++ b/frontend/src/store.ts @@ -1,24 +1,39 @@ import { ConnectError } from "@connectrpc/connect"; -import { anyUnpack } from "@bufbuild/protobuf/wkt"; +import { anyUnpack, timestampFromDate } from "@bufbuild/protobuf/wkt"; import { create } from "zustand"; -import { workspaceClient } from "./api"; +import { operationsClient, workspaceClient } from "./api"; import { DiscoveryRun, DiscoveryRunMetadataSchema, DiscoveryRunState, + ImportRun, + ImportRunMetadataSchema, + ImportRunState, LogFile, Pattern, + RecentImportQuery, UnmatchedErrorLine, Workspace } from "./gen/lapp/web/v1/web_pb"; -export type Tab = "logs" | "patterns" | "errors"; +export type Tab = "logs" | "imports" | "patterns" | "errors"; + +export type ImportInput = { + project: string; + filter: string; + from: Date; + to: Date; + limit: number; +}; type AppState = { workspaces: Workspace[]; selectedWorkspaceName: string; logFiles: LogFile[]; + importRuns: ImportRun[]; + recentImportQueries: RecentImportQuery[]; discoveryRuns: DiscoveryRun[]; + activeOperationName: string; selectedRunName: string; selectedPatternName: string; patterns: Pattern[]; @@ -39,12 +54,14 @@ type AppActions = { runAction: (action: () => Promise) => Promise; refreshWorkspaces: () => Promise; refreshWorkspaceDetails: (workspaceName?: string) => Promise; + refreshRecentImportQueries: () => Promise; + refreshActiveOperation: () => Promise; refreshResults: (runName?: string) => Promise; - refreshSelectedWorkspace: () => Promise; createWorkspace: () => Promise; deleteWorkspace: (workspace: Workspace) => Promise; uploadFiles: (files: FileList | null) => Promise; deleteLogFile: (logFile: LogFile) => Promise; + startImport: (input: ImportInput) => Promise; startDiscovery: () => Promise; }; @@ -58,7 +75,10 @@ export const useAppStore = create()((set, get) => ({ workspaces: [], selectedWorkspaceName: "", logFiles: [], + importRuns: [], + recentImportQueries: [], discoveryRuns: [], + activeOperationName: "", selectedRunName: "", selectedPatternName: "", patterns: [], @@ -77,7 +97,9 @@ export const useAppStore = create()((set, get) => ({ selectedRunName: "", selectedPatternName: "", logFiles: [], + importRuns: [], discoveryRuns: [], + activeOperationName: "", ...emptyResults }), selectRun: (selectedRunName) => set({ selectedRunName, selectedPatternName: "" }), @@ -107,32 +129,75 @@ export const useAppStore = create()((set, get) => ({ if (!workspaceName) { set({ logFiles: [], + importRuns: [], discoveryRuns: [], + activeOperationName: "", selectedRunName: "", ...emptyResults }); return; } - const [filesResponse, runsResponse] = await Promise.all([ + const [filesResponse, importRunsResponse, discoveryRunsResponse] = await Promise.all([ workspaceClient.listLogFiles({ parent: workspaceName }), + workspaceClient.listImportRuns({ parent: workspaceName }), workspaceClient.listDiscoveryRuns({ parent: workspaceName }) ]); const currentRun = get().selectedRunName; const latestSuccessful = - runsResponse.discoveryRuns.find((run) => run.state === DiscoveryRunState.SUCCEEDED)?.name || ""; + discoveryRunsResponse.discoveryRuns.find((run) => run.state === DiscoveryRunState.SUCCEEDED)?.name || ""; const selectedRunName = - currentRun && runsResponse.discoveryRuns.some((run) => run.name === currentRun) + currentRun && discoveryRunsResponse.discoveryRuns.some((run) => run.name === currentRun) ? currentRun : latestSuccessful; + const activeOperationName = + findActiveOperationName(importRunsResponse.importRuns, discoveryRunsResponse.discoveryRuns) || ""; set({ logFiles: filesResponse.logFiles, - discoveryRuns: runsResponse.discoveryRuns, + importRuns: importRunsResponse.importRuns, + discoveryRuns: discoveryRunsResponse.discoveryRuns, + activeOperationName, selectedRunName }); }, + refreshRecentImportQueries: async () => { + const response = await workspaceClient.listRecentImportQueries({}); + set({ recentImportQueries: response.recentQueries }); + }, + + refreshActiveOperation: async () => { + const activeOperationName = get().activeOperationName; + if (!activeOperationName) return; + + const operation = await operationsClient.getOperation({ name: activeOperationName }); + const importMetadata = operation.metadata + ? anyUnpack(operation.metadata, ImportRunMetadataSchema) + : undefined; + const discoveryMetadata = operation.metadata + ? anyUnpack(operation.metadata, DiscoveryRunMetadataSchema) + : undefined; + + set((state) => ({ + importRuns: importMetadata?.importRun + ? replaceRun(state.importRuns, importMetadata.importRun) + : state.importRuns, + discoveryRuns: discoveryMetadata?.discoveryRun + ? replaceRun(state.discoveryRuns, discoveryMetadata.discoveryRun) + : state.discoveryRuns, + selectedRunName: discoveryMetadata?.discoveryRun?.name || state.selectedRunName, + activeOperationName: operation.done ? "" : activeOperationName + })); + + if (operation.done) { + await get().refreshWorkspaceDetails(); + await get().refreshWorkspaces(); + await get().refreshRecentImportQueries(); + await get().refreshResults(); + } + }, + refreshResults: async (runName = get().selectedRunName) => { if (!runName) { set(emptyResults); @@ -154,12 +219,6 @@ export const useAppStore = create()((set, get) => ({ }); }, - refreshSelectedWorkspace: async () => { - await get().refreshWorkspaceDetails(); - await get().refreshResults(); - await get().refreshWorkspaces(); - }, - createWorkspace: async () => { const workspaceId = get().newWorkspaceID.trim(); if (!workspaceId) return; @@ -171,6 +230,8 @@ export const useAppStore = create()((set, get) => ({ selectedWorkspaceName, selectedRunName: "", selectedPatternName: "", + importRuns: [], + activeOperationName: "", ...emptyResults }); await get().refreshWorkspaces(); @@ -186,7 +247,9 @@ export const useAppStore = create()((set, get) => ({ selectedRunName: "", selectedPatternName: "", logFiles: [], + importRuns: [], discoveryRuns: [], + activeOperationName: "", ...emptyResults }); await get().refreshWorkspaces(); @@ -216,13 +279,43 @@ export const useAppStore = create()((set, get) => ({ await get().refreshWorkspaces(); }, + startImport: async (input) => { + const selectedWorkspaceName = get().selectedWorkspaceName; + if (!selectedWorkspaceName) return; + + const operation = await workspaceClient.createImportRun({ + parent: selectedWorkspaceName, + project: input.project, + filter: input.filter, + from: timestampFromDate(input.from), + to: timestampFromDate(input.to), + limit: input.limit + }); + const metadata = operation.metadata + ? anyUnpack(operation.metadata, ImportRunMetadataSchema) + : undefined; + set((state) => ({ + activeTab: "imports", + activeOperationName: operation.name, + importRuns: metadata?.importRun + ? replaceRun(state.importRuns, metadata.importRun) + : state.importRuns + })); + await get().refreshWorkspaceDetails(selectedWorkspaceName); + await get().refreshWorkspaces(); + await get().refreshRecentImportQueries(); + }, + startDiscovery: async () => { const selectedWorkspaceName = get().selectedWorkspaceName; if (!selectedWorkspaceName) return; const response = await workspaceClient.createDiscoveryRun({ parent: selectedWorkspaceName }); const metadata = response.metadata ? anyUnpack(response.metadata, DiscoveryRunMetadataSchema) : undefined; - set({ selectedRunName: metadata?.discoveryRun?.name || "" }); + set({ + activeOperationName: response.name, + selectedRunName: metadata?.discoveryRun?.name || "" + }); await get().refreshWorkspaceDetails(selectedWorkspaceName); await get().refreshWorkspaces(); } @@ -241,3 +334,25 @@ function errorMessage(error: unknown) { function isPattern(value: Pattern | undefined): value is Pattern { return value !== undefined; } + +function replaceRun(runs: T[], replacement: T) { + const index = runs.findIndex((run) => run.name === replacement.name); + if (index < 0) { + return [replacement, ...runs]; + } + return runs.map((run, runIndex) => runIndex === index ? replacement : run); +} + +function findActiveOperationName(importRuns: ImportRun[], discoveryRuns: DiscoveryRun[]) { + const activeImport = importRuns.find( + (run) => run.state === ImportRunState.QUEUED || run.state === ImportRunState.RUNNING + ); + if (activeImport) { + return activeImport.name.replace("/importRuns/", "/operations/"); + } + + const activeDiscovery = discoveryRuns.find( + (run) => run.state === DiscoveryRunState.QUEUED || run.state === DiscoveryRunState.RUNNING + ); + return activeDiscovery?.name.replace("/discoveryRuns/", "/operations/"); +} diff --git a/frontend/src/styles.css b/frontend/src/styles.css index 95a3619..c001880 100644 --- a/frontend/src/styles.css +++ b/frontend/src/styles.css @@ -244,12 +244,232 @@ select { padding: 14px; } +.panel-heading { + display: flex; + align-items: flex-start; + justify-content: space-between; + gap: 12px; +} + +.panel-heading h2, +.panel-heading p { + margin: 0; +} + +.panel-heading h2 { + font-size: 18px; +} + +.panel-heading p { + margin-top: 4px; + color: #66768a; + font-size: 13px; +} + .panel-toolbar { display: flex; justify-content: flex-end; margin-bottom: 12px; } +.import-view { + display: grid; + gap: 12px; +} + +.import-top { + display: grid; + grid-template-columns: minmax(420px, 1.3fr) minmax(300px, 0.7fr); + gap: 12px; + align-items: start; +} + +.import-form { + display: grid; + gap: 14px; +} + +.import-form label { + display: grid; + gap: 6px; + min-width: 0; +} + +.import-form label > span, +.import-form legend { + color: #526174; + font-size: 13px; + font-weight: 600; +} + +.import-form fieldset { + display: grid; + gap: 10px; + min-width: 0; + margin: 0; + border: 1px solid #e4e9ee; + border-radius: 7px; + padding: 10px; +} + +.import-form input, +.import-form select { + width: 100%; +} + +.form-grid { + display: grid; + grid-template-columns: minmax(0, 1fr) minmax(180px, 0.35fr); + gap: 10px; +} + +.quick-ranges { + display: flex; + flex-wrap: wrap; + gap: 7px; +} + +.quick-ranges button { + min-height: 32px; + padding: 0 10px; +} + +.quick-ranges button.selected { + border-color: #3167b1; + background: #edf4ff; + color: #224f89; +} + +.form-actions { + display: flex; + justify-content: flex-end; +} + +.import-latest { + display: grid; + gap: 14px; +} + +.import-latest.failed { + border-color: #e7a7a7; + background: #fffafa; +} + +.import-state-line { + display: flex; + align-items: center; + justify-content: space-between; + gap: 10px; + color: #66768a; + font-size: 13px; +} + +.state-badge { + display: inline-flex; + width: fit-content; + border-radius: 999px; + background: #edf1f5; + color: #526174; + padding: 3px 8px; + font-size: 12px; + font-weight: 700; +} + +.state-badge.running { + background: #e8f2ff; + color: #245b9e; +} + +.state-badge.succeeded { + background: #e6f6eb; + color: #236d3b; +} + +.state-badge.failed { + background: #fde8e8; + color: #9f2424; +} + +.import-message, +.empty-copy { + margin: 0; + color: #526174; +} + +.import-latest dl { + display: grid; + grid-template-columns: auto minmax(0, 1fr); + gap: 7px 12px; + margin: 0; + font-size: 13px; +} + +.import-latest dt { + color: #66768a; +} + +.import-latest dd { + min-width: 0; + margin: 0; + overflow-wrap: anywhere; +} + +.truncation-warning { + display: flex; + align-items: flex-start; + gap: 8px; + border: 1px solid #f0c36c; + border-radius: 6px; + background: #fff7e1; + color: #6f4d00; + padding: 9px; + font-size: 13px; +} + +.import-history { + display: grid; + gap: 12px; +} + +.table-scroll { + overflow-x: auto; +} + +.import-history table { + min-width: 900px; +} + +.import-history td { + max-width: 290px; + font-size: 13px; +} + +.import-history td strong, +.filter-value, +.run-error, +.small-warning { + display: block; +} + +.filter-value { + margin-top: 5px; + overflow-wrap: anywhere; + color: #526174; + white-space: normal; +} + +.run-error { + margin-top: 6px; + color: #9f2424; +} + +.small-warning { + margin-top: 5px; + color: #8a5d00; + font-size: 12px; + font-weight: 700; +} + .file-input { display: none; } @@ -453,6 +673,11 @@ th { justify-content: flex-start; } + .import-top, + .form-grid { + grid-template-columns: 1fr; + } + .pattern-row { grid-template-columns: 1fr; } diff --git a/go.mod b/go.mod index b097338..df9f667 100644 --- a/go.mod +++ b/go.mod @@ -13,6 +13,7 @@ require ( github.com/duckdb/duckdb-go/v2 v2.5.5 github.com/go-errors/errors v1.5.1 github.com/google/uuid v1.6.0 + github.com/googleapis/gax-go/v2 v2.23.0 github.com/jaeyo/go-drain3 v0.1.2 github.com/joho/godotenv v1.5.1 github.com/spf13/cobra v1.10.2 @@ -23,6 +24,7 @@ require ( go.opentelemetry.io/otel/sdk v1.44.0 google.golang.org/api v0.290.0 google.golang.org/genproto/googleapis/rpc v0.0.0-20260706201446-f0a921348800 + google.golang.org/grpc v1.82.0 google.golang.org/protobuf v1.36.11 ) @@ -31,7 +33,6 @@ require ( cloud.google.com/go/auth v0.20.0 // indirect cloud.google.com/go/auth/oauth2adapt v0.2.8 // indirect cloud.google.com/go/compute/metadata v0.9.0 // indirect - cloud.google.com/go/iam v1.11.0 // indirect github.com/apache/arrow-go/v18 v18.5.1 // indirect github.com/bahlo/generic-list-go v0.2.0 // indirect github.com/buger/jsonparser v1.1.1 // indirect @@ -61,7 +62,6 @@ require ( github.com/google/flatbuffers v25.12.19+incompatible // indirect github.com/google/s2a-go v0.1.9 // indirect github.com/googleapis/enterprise-certificate-proxy v0.3.18 // indirect - github.com/googleapis/gax-go/v2 v2.23.0 // indirect github.com/goph/emperror v0.17.2 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.27.7 // indirect github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect @@ -109,6 +109,5 @@ require ( golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da // indirect google.golang.org/genproto v0.0.0-20260319201613-d00831a3d3e7 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260630182238-925bb5da69e7 // indirect - google.golang.org/grpc v1.82.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index 7a16eb7..7785d05 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,3 @@ -cel.dev/expr v0.25.1 h1:1KrZg61W6TWSxuNZ37Xy49ps13NUovb66QLprthtwi4= -cel.dev/expr v0.25.1/go.mod h1:hrXvqGP6G6gyx8UAHSHJ5RGk//1Oj5nXQ2NI02Nrsg4= cloud.google.com/go v0.123.0 h1:2NAUJwPR47q+E35uaJeYoNhuNEM9kM8SjgRgdeOJUSE= cloud.google.com/go v0.123.0/go.mod h1:xBoMV08QcqUGuPW65Qfm1o9Y4zKZBpGS+7bImXLTAZU= cloud.google.com/go/auth v0.20.0 h1:kXTssoVb4azsVDoUiF8KvxAqrsQcQtB53DcSgta74CA= @@ -8,24 +6,12 @@ cloud.google.com/go/auth/oauth2adapt v0.2.8 h1:keo8NaayQZ6wimpNSmW5OPc283g65QNIi cloud.google.com/go/auth/oauth2adapt v0.2.8/go.mod h1:XQ9y31RkqZCcwJWNSx2Xvric3RrU88hAYYbjDWYDL+c= cloud.google.com/go/compute/metadata v0.9.0 h1:pDUj4QMoPejqq20dK0Pg2N4yG9zIkYGdBtwLoEkH9Zs= cloud.google.com/go/compute/metadata v0.9.0/go.mod h1:E0bWwX5wTnLPedCKqk3pJmVgCBSM6qQI1yTBdEb3C10= -cloud.google.com/go/iam v1.11.0 h1:KieQ9Pb+LLPak1O3Rv3GgCxhnmkYf7Xyh0P5HfF1jFM= -cloud.google.com/go/iam v1.11.0/go.mod h1:KP+nKGugNJW4LcLx1uEZcq1ok5sQHFaQehQNl4QDgV4= cloud.google.com/go/logging v1.19.0 h1:NCqhdVUg3wQ8Cobdf16FDSuTGi3+6+hdSBHrY5TsR6Q= cloud.google.com/go/logging v1.19.0/go.mod h1:i40NZCHC9Gqvod4yE+yQfDWwlgwW/SrshkkGibCHxcA= cloud.google.com/go/longrunning v1.2.0 h1:WjYH3YHBGCxGJP9M4dWGHBfXr/cFIjMkNgWcJj7/iMM= cloud.google.com/go/longrunning v1.2.0/go.mod h1:5KMQALFGOCtFoi2xSOA1u3H7WKlhmckgiyFw7+LGQp0= -cloud.google.com/go/monitoring v1.24.3 h1:dde+gMNc0UhPZD1Azu6at2e79bfdztVDS5lvhOdsgaE= -cloud.google.com/go/monitoring v1.24.3/go.mod h1:nYP6W0tm3N9H/bOw8am7t62YTzZY+zUeQ+Bi6+2eonI= -cloud.google.com/go/storage v1.62.3 h1:SZq1t23NCI+e96dH77Dg3PEfsNNEjqO8zE5AnD8gVD0= -cloud.google.com/go/storage v1.62.3/go.mod h1:cpYz/kRVZ+UQAF1uHeea10/9ewcRbxGoGNKsS9daSXA= connectrpc.com/connect v1.20.0 h1:6TNDAB+WeNd2uolWNlYczB5E0KNNaVMNUEx8JEUsPmQ= connectrpc.com/connect v1.20.0/go.mod h1:A2ygJrukXwWy32vkCAAHNVguZrqZ+jeZ9rGRnGR4dN4= -github.com/GoogleCloudPlatform/opentelemetry-operations-go/detectors/gcp v1.32.0 h1:rIkQfkCOVKc1OiRCNcSDD8ml5RJlZbH/Xsq7lbpynwc= -github.com/GoogleCloudPlatform/opentelemetry-operations-go/detectors/gcp v1.32.0/go.mod h1:RD2SsorTmYhF6HkTmDw7KmPYQk8OBYwTkuasChwv7R4= -github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric v0.55.0 h1:UnDZ/zFfG1JhH/DqxIZYU/1CUAlTUScoXD/LcM2Ykk8= -github.com/GoogleCloudPlatform/opentelemetry-operations-go/exporter/metric v0.55.0/go.mod h1:IA1C1U7jO/ENqm/vhi7V9YYpBsp+IMyqNrEN94N7tVc= -github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.55.0 h1:0s6TxfCu2KHkkZPnBfsQ2y5qia0jl3MMrmBhu3nCOYk= -github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.55.0/go.mod h1:Mf6O40IAyB9zR/1J8nGDDPirZQQPbYJni8Yisy7NTMc= github.com/airbrake/gobrake v3.6.1+incompatible/go.mod h1:wM4gu3Cn0W0K7GUuVWnlXZU11AGBXMILnrdOU8Kn00o= github.com/andybalholm/brotli v1.2.0 h1:ukwgCxwYrmACq68yiUqwIWnGY0cTPox/M94sVwToPjQ= github.com/andybalholm/brotli v1.2.0/go.mod h1:rzTDkvFWvIrjDXZHkuS16NPggd91W3kUSvPlQ1pLaKY= @@ -110,8 +96,6 @@ github.com/go-check/check v0.0.0-20180628173108-788fd7840127 h1:0gkP6mzaMqkmpcJY github.com/go-check/check v0.0.0-20180628173108-788fd7840127/go.mod h1:9ES+weclKsC9YodN5RgxqK/VD9HM9JsCSh7rNhMZE98= github.com/go-errors/errors v1.5.1 h1:ZwEMSLRCapFLflTpT7NKaAc7ukJ8ZPEjzlxt8rPN8bk= github.com/go-errors/errors v1.5.1/go.mod h1:sIVyrIiJhuEF+Pj9Ebtd6P/rEYROXFi3BopGUQ5a5Og= -github.com/go-jose/go-jose/v4 v4.1.4 h1:moDMcTHmvE6Groj34emNPLs/qtYXRVcd6S7NHbHz3kA= -github.com/go-jose/go-jose/v4 v4.1.4/go.mod h1:x4oUasVrzR7071A4TnHLGSPpNOm2a21K9Kf04k1rs08= github.com/go-logr/logr v1.2.2/go.mod h1:jdQByPbusPIv2/zmleS9BjJVeZ6kBagPoEUsqbVz/1A= github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= @@ -122,8 +106,6 @@ github.com/go-viper/mapstructure/v2 v2.5.0/go.mod h1:oJDH3BJKyqBA2TXFhDsKDGDTlnd github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4= github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= github.com/gofrs/uuid v3.2.0+incompatible/go.mod h1:b2aQJv3Z4Fp6yNu3cdSllBxTCLRxnplIgP/c0N/04lM= -github.com/golang/groupcache v0.0.0-20210331224755-41bb18bfe9da h1:oI5xCqsCo564l8iNU+DwB5epxmsaqB+rhGL0m5jtYqE= -github.com/golang/groupcache v0.0.0-20210331224755-41bb18bfe9da/go.mod h1:cIg4eruTrX1D+g88fzRXU5OdNfaM+9IcxsU14FzY7Hc= github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc= github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs= github.com/golang/protobuf v1.2.0/go.mod h1:6lQm79b+lXiMfvg/cZm0SGofjICqVBUtrP5yJMmIC1U= @@ -233,8 +215,6 @@ github.com/spf13/cobra v1.10.2 h1:DMTTonx5m65Ic0GOoRY2c16WCbHxOOw6xxezuLaBpcU= github.com/spf13/cobra v1.10.2/go.mod h1:7C1pvHqHw5A4vrJfjNwvOdzYu0Gml16OCs2GRiTUUS4= github.com/spf13/pflag v1.0.9 h1:9exaQaMOCwffKiiiYk6/BndUBv+iRViNW+4lEMi0PvY= github.com/spf13/pflag v1.0.9/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= -github.com/spiffe/go-spiffe/v2 v2.6.0 h1:l+DolpxNWYgruGQVV0xsfeya3CsC7m8iBzDnMpsbLuo= -github.com/spiffe/go-spiffe/v2 v2.6.0/go.mod h1:gm2SeUoMZEtpnzPNs2Csc0D/gX33k1xIx7lEzqblHEs= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.1.1/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= @@ -271,12 +251,8 @@ github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0= github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= -go.opencensus.io v0.24.0 h1:y73uSU6J157QMP2kn2r30vwW1A2W2WFwSCGnAVxeaD0= -go.opencensus.io v0.24.0/go.mod h1:vNK8G9p7aAivkbmorf4v+7Hgx+Zs0yY+0fOtgBfjQKo= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= -go.opentelemetry.io/contrib/detectors/gcp v1.43.0 h1:62yY3dT7/ShwOxzA0RsKRgshBmfElKI4d/Myu2OxDFU= -go.opentelemetry.io/contrib/detectors/gcp v1.43.0/go.mod h1:RyaZMFY7yi1kAs45S6mbFGz8O8rqB0dTY14uzvG4LCs= go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.67.0 h1:yI1/OhfEPy7J9eoa6Sj051C7n5dvpj0QX8g4sRchg04= go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc v0.67.0/go.mod h1:NoUCKYWK+3ecatC4HjkRktREheMeEtrXoQxrqYFeHSc= go.opentelemetry.io/contrib/instrumentation/net/http/otelhttp v0.67.0 h1:OyrsyzuttWTSur2qN/Lm0m2a8yqyIjUVBZcxFPuXq2o= diff --git a/pkg/gcplog/fetch.go b/pkg/gcplog/fetch.go index 92274c4..5986715 100644 --- a/pkg/gcplog/fetch.go +++ b/pkg/gcplog/fetch.go @@ -2,23 +2,31 @@ package gcplog import ( "context" + "encoding/base64" "encoding/json" "fmt" "strings" "time" - "cloud.google.com/go/logging" - "cloud.google.com/go/logging/logadmin" + loggingv2 "cloud.google.com/go/logging/apiv2" + loggingpb "cloud.google.com/go/logging/apiv2/loggingpb" "github.com/go-errors/errors" + "github.com/googleapis/gax-go/v2" "google.golang.org/api/iterator" + "google.golang.org/grpc/codes" "google.golang.org/protobuf/encoding/protojson" - "google.golang.org/protobuf/proto" ) // adcHint is appended to fetch failures so credential problems point people // at Application Default Credentials setup. LAPP never manages provider auth. const adcHint = "if this is a credential problem, run `gcloud auth application-default login`" +const ( + maxPageSize = 10_000 + pageTimeout = 5 * time.Minute + maxRetryDelay = time.Minute +) + // FetchRequest describes one pull from GCP Cloud Logging. Filter is the user // supplied Cloud Logging filter; the time range is combined with it when // querying. Limit caps the number of returned entries. @@ -30,24 +38,31 @@ type FetchRequest struct { Limit int } -// FetchResult carries converted envelope NDJSON lines. Truncated is true -// when the limit cut the result short. +// FetchResult reports facts about one completed fetch. type FetchResult struct { - Lines []string Truncated bool } -// FetchLines pulls matching entries from GCP Cloud Logging using Application -// Default Credentials and converts each one to an envelope NDJSON line. -func FetchLines(ctx context.Context, req FetchRequest) (FetchResult, error) { - client, err := logadmin.NewClient(ctx, req.Project) +// Fetch pulls matching entries from GCP Cloud Logging using Application +// Default Credentials and sends each converted NDJSON line to writeLine. +func Fetch(ctx context.Context, req FetchRequest, writeLine func(string) error) (FetchResult, error) { + if writeLine == nil { + return FetchResult{}, errors.New("log line writer is required") + } + client, err := loggingv2.NewClient(ctx) if err != nil { return FetchResult{}, errors.Errorf("create cloud logging client (%s): %w", adcHint, err) } defer func() { _ = client.Close() }() - it := client.Entries(ctx, logadmin.Filter(CombinedFilter(req.Filter, req.From, req.To))) + it := client.ListLogEntries(ctx, &loggingpb.ListLogEntriesRequest{ + ResourceNames: []string{"projects/" + req.Project}, + Filter: CombinedFilter(req.Filter, req.From, req.To), + PageSize: pageSize(req.Limit), + }, gax.WithTimeout(pageTimeout), gax.WithRetry(newListRetryer)) + var result FetchResult + count := 0 for { entry, err := it.Next() if errors.Is(err, iterator.Done) { @@ -56,19 +71,46 @@ func FetchLines(ctx context.Context, req FetchRequest) (FetchResult, error) { if err != nil { return FetchResult{}, errors.Errorf("list cloud logging entries (%s): %w", adcHint, err) } - if req.Limit > 0 && len(result.Lines) >= req.Limit { + if req.Limit > 0 && count >= req.Limit { result.Truncated = true break } - line, err := ConvertEntry(entryFromLogging(entry)) + converted, err := entryFromProto(entry) if err != nil { return FetchResult{}, err } - result.Lines = append(result.Lines, line) + line, err := ConvertEntry(converted) + if err != nil { + return FetchResult{}, err + } + if err := writeLine(line); err != nil { + return FetchResult{}, errors.Errorf("write cloud logging entry: %w", err) + } + count++ } return result, nil } +func pageSize(limit int) int32 { + if limit > 0 && limit < maxPageSize { + return int32(limit + 1) + } + return maxPageSize +} + +func newListRetryer() gax.Retryer { + return gax.OnCodes([]codes.Code{ + codes.ResourceExhausted, + codes.DeadlineExceeded, + codes.Internal, + codes.Unavailable, + }, gax.Backoff{ + Initial: time.Second, + Max: maxRetryDelay, + Multiplier: 2, + }) +} + // CombinedFilter joins the resolved time range with the user supplied filter. // The time range is a separate concept from the filter; they are only // combined here at query time. @@ -83,31 +125,57 @@ func CombinedFilter(filter string, from, to time.Time) string { return strings.Join(parts, " AND ") } -// entryFromLogging maps a client library entry onto the neutral Entry shape -// consumed by ConvertEntry. -func entryFromLogging(e *logging.Entry) Entry { - out := Entry{Timestamp: e.Timestamp} - if e.Severity != logging.Default { - out.Severity = strings.ToUpper(e.Severity.String()) +// entryFromProto maps a provider entry onto the neutral Entry shape consumed +// by ConvertEntry. +func entryFromProto(e *loggingpb.LogEntry) (Entry, error) { + out := Entry{} + if timestamp := e.GetTimestamp(); timestamp != nil { + if err := timestamp.CheckValid(); err != nil { + return Entry{}, errors.Errorf("invalid cloud logging timestamp: %w", err) + } + out.Timestamp = timestamp.AsTime() + } + if severity := strings.ToUpper(e.GetSeverity().String()); severity != DefaultSeverity { + out.Severity = severity } switch payload := e.Payload.(type) { case nil: - case string: - out.TextPayload = &payload - case proto.Message: - if data, err := protojson.Marshal(payload); err == nil { + case *loggingpb.LogEntry_TextPayload: + out.TextPayload = &payload.TextPayload + case *loggingpb.LogEntry_JsonPayload: + if payload.JsonPayload == nil { + break + } + if data, err := protojson.Marshal(payload.JsonPayload); err == nil { out.JSONPayload = data } else { text := fmt.Sprint(payload) out.TextPayload = &text } - default: - if data, err := json.Marshal(payload); err == nil { + case *loggingpb.LogEntry_ProtoPayload: + if payload.ProtoPayload == nil { + break + } + message, err := payload.ProtoPayload.UnmarshalNew() + if err == nil { + data, marshalErr := protojson.Marshal(message) + if marshalErr == nil { + out.JSONPayload = data + break + } + } + data, marshalErr := json.Marshal(map[string]string{ + "@type": payload.ProtoPayload.GetTypeUrl(), + "value": base64.StdEncoding.EncodeToString(payload.ProtoPayload.GetValue()), + }) + if marshalErr == nil { out.JSONPayload = data } else { text := fmt.Sprint(payload) out.TextPayload = &text } + default: + return Entry{}, errors.Errorf("unknown cloud logging payload type %T", e.Payload) } - return out + return out, nil } diff --git a/pkg/webapp/import_service_test.go b/pkg/webapp/import_service_test.go index 13dd2ef..a25ee5d 100644 --- a/pkg/webapp/import_service_test.go +++ b/pkg/webapp/import_service_test.go @@ -21,7 +21,7 @@ func TestCreateImportRunRejectsInvalidRequestWithoutRecord(t *testing.T) { root := t.TempDir() service, err := NewWorkspaceService(ServiceConfig{ Root: root, - ImportFetcher: func(context.Context, workspace.ImportRequest) (workspace.ImportFetchResult, error) { + ImportFetcher: func(context.Context, workspace.ImportRequest, workspace.ImportLineWriter) (workspace.ImportFetchResult, error) { t.Fatal("fetcher must not run") return workspace.ImportFetchResult{}, nil }, @@ -135,13 +135,16 @@ func TestWorkspaceServiceImportRunAsyncSuccess(t *testing.T) { fetchStarted := make(chan struct{}) releaseFetch := make(chan struct{}) line := `{"ts":"2026-07-25T00:00:01Z","severity":"ERROR","payload":{"message":"failed request"}}` - service, workspaceName := newImportTestService(t, func(_ context.Context, req workspace.ImportRequest) (workspace.ImportFetchResult, error) { + service, workspaceName := newImportTestService(t, func(_ context.Context, req workspace.ImportRequest, writeLine workspace.ImportLineWriter) (workspace.ImportFetchResult, error) { if req.Limit != defaultImportLimit { t.Errorf("fetch limit = %d, want %d", req.Limit, defaultImportLimit) } close(fetchStarted) <-releaseFetch - return workspace.ImportFetchResult{Lines: []string{line}}, nil + if err := writeLine(line); err != nil { + return workspace.ImportFetchResult{}, err + } + return workspace.ImportFetchResult{}, nil }) response, err := service.CreateImportRun(context.Background(), connect.NewRequest(validImportRequest(workspaceName))) @@ -213,7 +216,7 @@ func TestWorkspaceServiceImportRunAsyncSuccess(t *testing.T) { } func TestWorkspaceServiceImportRunFailure(t *testing.T) { - service, workspaceName := newImportTestService(t, func(context.Context, workspace.ImportRequest) (workspace.ImportFetchResult, error) { + service, workspaceName := newImportTestService(t, func(context.Context, workspace.ImportRequest, workspace.ImportLineWriter) (workspace.ImportFetchResult, error) { return workspace.ImportFetchResult{}, stderrors.New("cloud logging unavailable") }) @@ -238,11 +241,11 @@ func TestWorkspaceServiceImportRunFailure(t *testing.T) { } func TestWorkspaceServiceImportRunTruncation(t *testing.T) { - service, workspaceName := newImportTestService(t, func(context.Context, workspace.ImportRequest) (workspace.ImportFetchResult, error) { - return workspace.ImportFetchResult{ - Lines: []string{`{"payload":{"message":"one"}}`}, - Truncated: true, - }, nil + service, workspaceName := newImportTestService(t, func(_ context.Context, _ workspace.ImportRequest, writeLine workspace.ImportLineWriter) (workspace.ImportFetchResult, error) { + if err := writeLine(`{"payload":{"message":"one"}}`); err != nil { + return workspace.ImportFetchResult{}, err + } + return workspace.ImportFetchResult{Truncated: true}, nil }) response, err := service.CreateImportRun(context.Background(), connect.NewRequest(validImportRequest(workspaceName))) @@ -292,7 +295,7 @@ func TestWorkspaceServiceRunMutualExclusion(t *testing.T) { releaseDiscovery := make(chan struct{}) service, err := NewWorkspaceService(ServiceConfig{ Root: t.TempDir(), - ImportFetcher: func(context.Context, workspace.ImportRequest) (workspace.ImportFetchResult, error) { + ImportFetcher: func(context.Context, workspace.ImportRequest, workspace.ImportLineWriter) (workspace.ImportFetchResult, error) { close(importStarted) <-releaseImport return workspace.ImportFetchResult{}, nil diff --git a/pkg/webapp/server.go b/pkg/webapp/server.go index 7b907f6..d317885 100644 --- a/pkg/webapp/server.go +++ b/pkg/webapp/server.go @@ -35,19 +35,18 @@ func NewHandler(config ServerConfig) (http.Handler, error) { Root: config.Root, APIKey: config.APIKey, Model: config.Model, - ImportFetcher: func(ctx context.Context, req workspace.ImportRequest) (workspace.ImportFetchResult, error) { - fetched, err := gcplog.FetchLines(ctx, gcplog.FetchRequest{ + ImportFetcher: func(ctx context.Context, req workspace.ImportRequest, writeLine workspace.ImportLineWriter) (workspace.ImportFetchResult, error) { + fetched, err := gcplog.Fetch(ctx, gcplog.FetchRequest{ Project: req.Project, Filter: req.Filter, From: req.From, To: req.To, Limit: req.Limit, - }) + }, writeLine) if err != nil { return workspace.ImportFetchResult{}, err } return workspace.ImportFetchResult{ - Lines: fetched.Lines, Truncated: fetched.Truncated, }, nil }, diff --git a/pkg/workspace/import.go b/pkg/workspace/import.go index cd690a7..f77d1af 100644 --- a/pkg/workspace/import.go +++ b/pkg/workspace/import.go @@ -1,17 +1,20 @@ package workspace import ( + "bufio" "context" + stderrors "errors" "log/slog" "os" "path/filepath" - "strings" "time" - "github.com/go-errors/errors" + goerrors "github.com/go-errors/errors" "github.com/google/uuid" ) +var errImportLimitReached = stderrors.New("import limit reached") + // ImportRequest is what the fetcher needs to pull entries from a provider. type ImportRequest struct { Project string @@ -21,17 +24,22 @@ type ImportRequest struct { Limit int } -// ImportFetchResult carries provider entries already converted to envelope -// NDJSON lines. Truncated is true when the limit cut the result short. +// ImportFetchResult reports facts about one completed provider fetch. type ImportFetchResult struct { - Lines []string Truncated bool } +// ImportLineWriter writes one envelope NDJSON line without a trailing newline. +type ImportLineWriter func(line string) error + // ImportFetcher pulls matching entries from a provider. It is the narrow // seam between run orchestration and provider network access, so tests can // inject a fake without touching the network (precedent: DiscoveryConfig.Labeler). -type ImportFetcher func(ctx context.Context, req ImportRequest) (ImportFetchResult, error) +type ImportFetcher func( + ctx context.Context, + req ImportRequest, + writeLine ImportLineWriter, +) (ImportFetchResult, error) type ImportConfig struct { Dir string @@ -59,13 +67,13 @@ type ImportResult struct { // Zero matching entries is a success with no log file written. func RunImport(ctx context.Context, config ImportConfig) (ImportResult, error) { if config.Dir == "" { - return ImportResult{}, errors.New("workspace dir is required") + return ImportResult{}, goerrors.New("workspace dir is required") } if config.Provider == "" { - return ImportResult{}, errors.New("import provider is required") + return ImportResult{}, goerrors.New("import provider is required") } if config.Fetcher == nil { - return ImportResult{}, errors.New("import fetcher is required") + return ImportResult{}, goerrors.New("import fetcher is required") } runID, err := importRunID(config.RunID) if err != nil { @@ -80,41 +88,18 @@ func RunImport(ctx context.Context, config ImportConfig) (ImportResult, error) { return ImportResult{}, err } - fetched, err := config.Fetcher(ctx, ImportRequest{ - Project: config.Project, - Filter: config.Filter, - From: record.From, - To: record.To, - Limit: config.Limit, - }) + fetched, errorCode, err := fetchImport(ctx, config, record) if err != nil { - failImportRun(config.Dir, &record, "FETCH_FAILED", err) + failImportRun(config.Dir, &record, errorCode, err) return ImportResult{}, err } - // The limit is enforced here as well as in the fetcher so the cap holds - // regardless of provider behavior. - if config.Limit > 0 && len(fetched.Lines) > config.Limit { - fetched.Lines = fetched.Lines[:config.Limit] - fetched.Truncated = true - } - - if len(fetched.Lines) > 0 { - fileName := config.Provider + "-" + runID + ".ndjson" - logPath := filepath.Join(config.Dir, "logs", fileName) - content := strings.Join(fetched.Lines, "\n") + "\n" - if err := os.WriteFile(logPath, []byte(content), 0o644); err != nil { - err = errors.Errorf("write imported log file: %w", err) - failImportRun(config.Dir, &record, "WRITE_LOG_FAILED", err) - return ImportResult{}, err - } - record.LogFileName = fileName - } finishedAt := timeNow() record.State = ImportRunStateSucceeded record.FinishedAt = &finishedAt - record.EntryCount = len(fetched.Lines) + record.EntryCount = fetched.EntryCount record.Truncated = fetched.Truncated + record.LogFileName = fetched.LogFileName if err := WriteImportRunRecord(config.Dir, record); err != nil { // Best effort: do not leave the run stuck in RUNNING when the final // record write fails. A secondary failure here is ignored. @@ -137,6 +122,131 @@ func RunImport(ctx context.Context, config ImportConfig) (ImportResult, error) { }, nil } +type importFetchOutcome struct { + EntryCount int + Truncated bool + LogFileName string +} + +type importLogOutput struct { + file *os.File + writer *bufio.Writer + tempPath string + logPath string + fileName string + limit int + entryCount int + limitReached bool + writeErr error + closed bool + committed bool +} + +func fetchImport( + ctx context.Context, + config ImportConfig, + record ImportRunRecord, +) (importFetchOutcome, string, error) { + output, err := newImportLogOutput(config, record.ID) + if err != nil { + return importFetchOutcome{}, "WRITE_LOG_FAILED", err + } + defer output.cleanup() + + fetched, err := config.Fetcher(ctx, ImportRequest{ + Project: config.Project, + Filter: config.Filter, + From: record.From, + To: record.To, + Limit: config.Limit, + }, output.writeLine) + if err != nil && !stderrors.Is(err, errImportLimitReached) { + if output.writeErr != nil { + return importFetchOutcome{}, "WRITE_LOG_FAILED", output.writeErr + } + return importFetchOutcome{}, "FETCH_FAILED", err + } + if err := output.commit(); err != nil { + return importFetchOutcome{}, "WRITE_LOG_FAILED", err + } + return importFetchOutcome{ + EntryCount: output.entryCount, + Truncated: fetched.Truncated || output.limitReached, + LogFileName: output.committedFileName(), + }, "", nil +} + +func newImportLogOutput(config ImportConfig, runID string) (*importLogOutput, error) { + fileName := config.Provider + "-" + runID + ".ndjson" + logPath := filepath.Join(config.Dir, "logs", fileName) + file, err := os.CreateTemp(filepath.Dir(logPath), "."+fileName+".*.tmp") + if err != nil { + return nil, goerrors.Errorf("create imported log file: %w", err) + } + return &importLogOutput{ + file: file, + writer: bufio.NewWriter(file), + tempPath: file.Name(), + logPath: logPath, + fileName: fileName, + limit: config.Limit, + }, nil +} + +func (output *importLogOutput) writeLine(line string) error { + if output.limit > 0 && output.entryCount >= output.limit { + output.limitReached = true + return errImportLimitReached + } + if _, err := output.writer.WriteString(line); err != nil { + output.writeErr = goerrors.Errorf("write imported log file: %w", err) + return output.writeErr + } + if err := output.writer.WriteByte('\n'); err != nil { + output.writeErr = goerrors.Errorf("write imported log file: %w", err) + return output.writeErr + } + output.entryCount++ + return nil +} + +func (output *importLogOutput) commit() error { + if err := output.writer.Flush(); err != nil { + return goerrors.Errorf("write imported log file: %w", err) + } + if err := output.file.Chmod(0o644); err != nil { + return goerrors.Errorf("set imported log file permissions: %w", err) + } + if err := output.file.Close(); err != nil { + return goerrors.Errorf("close imported log file: %w", err) + } + output.closed = true + if output.entryCount == 0 { + return nil + } + if err := os.Rename(output.tempPath, output.logPath); err != nil { + return goerrors.Errorf("commit imported log file: %w", err) + } + output.committed = true + return nil +} + +func (output *importLogOutput) committedFileName() string { + if !output.committed { + return "" + } + return output.fileName +} + +func (output *importLogOutput) cleanup() { + if !output.closed { + _ = output.file.Close() + } + if !output.committed { + _ = os.Remove(output.tempPath) + } +} + func claimImportRun(config ImportConfig, runID string) (ImportRunRecord, error) { record := ImportRunRecord{ ID: runID, @@ -151,13 +261,13 @@ func claimImportRun(config ImportConfig, runID string) (ImportRunRecord, error) } existing, err := ReadImportRunRecord(config.Dir, runID) if err != nil { - if errors.Is(err, os.ErrNotExist) { + if goerrors.Is(err, os.ErrNotExist) { return record, nil } - return ImportRunRecord{}, errors.Errorf("check import run %q: %w", runID, err) + return ImportRunRecord{}, goerrors.Errorf("check import run %q: %w", runID, err) } if existing.State != ImportRunStateQueued || !sameImportRequest(existing, record) { - return ImportRunRecord{}, errors.Errorf("import run %q already exists", runID) + return ImportRunRecord{}, goerrors.Errorf("import run %q already exists", runID) } existing.State = ImportRunStateRunning return existing, nil @@ -179,7 +289,7 @@ func importRunID(value string) (string, error) { } id, err := uuid.NewV7() if err != nil { - return "", errors.Errorf("create import run id: %w", err) + return "", goerrors.Errorf("create import run id: %w", err) } return id.String(), nil } diff --git a/pkg/workspace/import_test.go b/pkg/workspace/import_test.go index a7c9a15..0e4016d 100644 --- a/pkg/workspace/import_test.go +++ b/pkg/workspace/import_test.go @@ -34,9 +34,14 @@ func TestRunImportSuccessWritesFileAndRecord(t *testing.T) { `{"ts":"2026-07-25T00:00:02Z","severity":"ERROR","payload":{"message":"b"}}`, } var gotReq ImportRequest - fetcher := func(_ context.Context, req ImportRequest) (ImportFetchResult, error) { + fetcher := func(_ context.Context, req ImportRequest, writeLine ImportLineWriter) (ImportFetchResult, error) { gotReq = req - return ImportFetchResult{Lines: lines}, nil + for _, line := range lines { + if err := writeLine(line); err != nil { + return ImportFetchResult{}, err + } + } + return ImportFetchResult{}, nil } result, err := RunImport(context.Background(), importTestConfig(dir, fetcher)) @@ -80,7 +85,7 @@ func TestRunImportSuccessWritesFileAndRecord(t *testing.T) { func TestRunImportZeroEntriesSucceedsWithoutFile(t *testing.T) { dir := t.TempDir() mustMkdir(t, filepath.Join(dir, "logs")) - fetcher := func(_ context.Context, _ ImportRequest) (ImportFetchResult, error) { + fetcher := func(_ context.Context, _ ImportRequest, _ ImportLineWriter) (ImportFetchResult, error) { return ImportFetchResult{}, nil } @@ -112,11 +117,11 @@ func TestRunImportZeroEntriesSucceedsWithoutFile(t *testing.T) { func TestRunImportTruncationIsRecorded(t *testing.T) { dir := t.TempDir() mustMkdir(t, filepath.Join(dir, "logs")) - fetcher := func(_ context.Context, _ ImportRequest) (ImportFetchResult, error) { - return ImportFetchResult{ - Lines: []string{`{"ts":"2026-07-25T00:00:01Z","severity":"ERROR","payload":{"message":"a"}}`}, - Truncated: true, - }, nil + fetcher := func(_ context.Context, _ ImportRequest, writeLine ImportLineWriter) (ImportFetchResult, error) { + if err := writeLine(`{"ts":"2026-07-25T00:00:01Z","severity":"ERROR","payload":{"message":"a"}}`); err != nil { + return ImportFetchResult{}, err + } + return ImportFetchResult{Truncated: true}, nil } result, err := RunImport(context.Background(), importTestConfig(dir, fetcher)) @@ -139,7 +144,7 @@ func TestRunImportTruncationIsRecorded(t *testing.T) { func TestRunImportFetchFailureWritesFailedRecord(t *testing.T) { dir := t.TempDir() mustMkdir(t, filepath.Join(dir, "logs")) - fetcher := func(_ context.Context, _ ImportRequest) (ImportFetchResult, error) { + fetcher := func(_ context.Context, _ ImportRequest, _ ImportLineWriter) (ImportFetchResult, error) { return ImportFetchResult{}, errors.New("could not find default credentials") } @@ -174,14 +179,17 @@ func TestRunImportFetchFailureWritesFailedRecord(t *testing.T) { func TestRunImportCapsLinesAtLimit(t *testing.T) { dir := t.TempDir() mustMkdir(t, filepath.Join(dir, "logs")) - fetcher := func(_ context.Context, _ ImportRequest) (ImportFetchResult, error) { - return ImportFetchResult{ - Lines: []string{ - `{"ts":"2026-07-25T00:00:01Z","severity":"ERROR","payload":{"message":"a"}}`, - `{"ts":"2026-07-25T00:00:02Z","severity":"ERROR","payload":{"message":"b"}}`, - `{"ts":"2026-07-25T00:00:03Z","severity":"ERROR","payload":{"message":"c"}}`, - }, - }, nil + fetcher := func(_ context.Context, _ ImportRequest, writeLine ImportLineWriter) (ImportFetchResult, error) { + for _, line := range []string{ + `{"ts":"2026-07-25T00:00:01Z","severity":"ERROR","payload":{"message":"a"}}`, + `{"ts":"2026-07-25T00:00:02Z","severity":"ERROR","payload":{"message":"b"}}`, + `{"ts":"2026-07-25T00:00:03Z","severity":"ERROR","payload":{"message":"c"}}`, + } { + if err := writeLine(line); err != nil { + return ImportFetchResult{}, err + } + } + return ImportFetchResult{}, nil } config := importTestConfig(dir, fetcher) @@ -204,10 +212,11 @@ func TestRunImportCapsLinesAtLimit(t *testing.T) { func TestRunImportRejectsDuplicateRunID(t *testing.T) { dir := t.TempDir() mustMkdir(t, filepath.Join(dir, "logs")) - fetcher := func(_ context.Context, _ ImportRequest) (ImportFetchResult, error) { - return ImportFetchResult{ - Lines: []string{`{"ts":"2026-07-25T00:00:01Z","severity":"ERROR","payload":{"message":"a"}}`}, - }, nil + fetcher := func(_ context.Context, _ ImportRequest, writeLine ImportLineWriter) (ImportFetchResult, error) { + if err := writeLine(`{"ts":"2026-07-25T00:00:01Z","severity":"ERROR","payload":{"message":"a"}}`); err != nil { + return ImportFetchResult{}, err + } + return ImportFetchResult{}, nil } config := importTestConfig(dir, fetcher)