diff --git a/docs/01-basic/report.md b/docs/01-basic/report.md new file mode 100644 index 0000000..60d2765 --- /dev/null +++ b/docs/01-basic/report.md @@ -0,0 +1,113 @@ +# (Q1.1) + +## 1. 按逗号分割日志的语句,以及如何指定每个字段的意义 + +本框架借助第三方 CSV 解析库 **CsvHelper** 来完成按逗号分列的工作: + +```csharp +using var csv = new CsvReader(logFile, config); +csv.Context.RegisterClassMap(); +foreach (var logRecord in csv.GetRecords()) { ... } +``` + +其中真正「按逗号把一行切成多个字段」的工作由 `CsvReader` / `csv.GetRecords()` 在库内部完成(它还能正确处理 `message` 字段两端的双引号以及 JSON 内部出现的逗号)。 + +「每一行的第几个字段代表何种意义」是通过一个继承自 `ClassMap` 的映射类 `LogRecordMap` 来指定的,使用 `Map(...).Index(n)` 把 CSV 的第 `n` 列绑定到 `LogRecord` 的对应属性上: + +```csharp +internal class LogRecordMap : ClassMap +{ + public LogRecordMap() + { + Map(m => m.LineNo).Index(0); // 第 0 列 -> LineNo + Map(m => m.Timestamp).Index(1); // 第 1 列 -> Timestamp + Map(m => m.PodName).Index(2); // 第 2 列 -> PodName + Map(m => m.Message).Index(3); // 第 3 列 -> Message + } +} +``` + +即:`Index(0)` 对应 `lineno`、`Index(1)` 对应 `timestamp`、`Index(2)` 对应 `pod-name`、`Index(3)` 对应 `message`。随后通过 `csv.Context.RegisterClassMap()` 让 CsvHelper 读取时按这个映射把每列填入 `LogRecord`。 + +## 2. 在哪个方法内、用哪几条语句判断日志种类 + +在 `Parser/LineParser.cs` 的 `ParseLine(LogRecord logRecord)` 方法内判断。先用 `JsonDocument` 把 `message` 当作 JSON 解析,再读取其中的 `event` 字段,用 `switch` 表达式根据其取值分流到不同的工厂方法: + +```csharp +using (var doc = JsonDocument.Parse(logRecord.Message)) +{ + var root = doc.RootElement; + if (root.TryGetProperty("event", out var eventElement)) + { + return eventElement.GetString() switch + { + "call" => LineParser.CreateCall(logRecord), + "request" => LineParser.CreateRequest(logRecord), + "internal" => LineParser.CreateInternal(logRecord), + _ => throw new FormatException(...) + }; + } + ... +} +``` + +也就是说,判断种类的语句是 `root.TryGetProperty("event", out var eventElement)` 配合 `eventElement.GetString() switch { "call" => ..., "request" => ..., "internal" => ... }`。 + +## 3. 确定种类后调用哪个库方法解析 JSON + +确定种类后,在对应的工厂方法(如 `CreateCall` / `CreateRequest` / `CreateInternal`)中调用 `System.Text.Json` 提供的: + +```csharp +JsonSerializer.Deserialize(logRecord.Message, options) +``` + +把 JSON 字符串反序列化成一个强类型的 `record`(如 `CallMessage`)。 + +### 3.1 如何防止日志中字段缺失 + +通过在反序列化目标 `record` 的每个属性上标注 `[property: JsonRequired]` 特性,例如: + +```csharp +private record CallMessage( + [property: JsonRequired] string Severity, + [property: JsonRequired] string RequestId, + [property: JsonRequired] string TargetService, + [property: JsonRequired] int DurationMs +); +``` + +`[JsonRequired]` 告诉 JSON 序列化器这些属性是必需的:当 JSON 中缺少对应键时,`JsonSerializer.Deserialize` 会抛出 `JsonException`,从而把「字段缺失」这一异常情况暴露出来。此外,反序列化结果后还跟了一个 `?? throw new FormatException(...)`,用于在结果为 `null` 时也抛出异常,进一步兜底: + +```csharp +var callMessage = JsonSerializer.Deserialize(logRecord.Message, options) + ?? throw new FormatException(...); +``` + +### 3.2 如何让 JSON 解析器完成「烤串命名法 → 大驼峰命名法」的转换 + +通过配置 `JsonSerializerOptions` 的 `PropertyNamingPolicy`: + +```csharp +private static JsonSerializerOptions options = new JsonSerializerOptions +{ + PropertyNamingPolicy = JsonNamingPolicy.KebabCaseLower, +}; +``` + +`JsonNamingPolicy.KebabCaseLower` 作为命名策略,会在反序列化时把 C# 属性名(大驼峰,如 `RequestId`、`TargetService`、`DurationMs`)转换成小写烤串形式(`request-id`、`target-service`、`duration-ms`)再去和 JSON 中的键匹配。这样就在不修改 C# 属性名的前提下,完成了 `abc-def` 与 `AbcDef` 两种命名法之间的映射。 + +# (Q1.2) + +以一个 Call 事件为例,调用 `KeyValueVisitor.Dump(entry)`(其中 `entry` 的静态类型是 `LogEntry`,实际运行时类型是 `CallLogEntry`)后,方法调用链如下(.NET 内置库方法略): + ++ `Dictionary KeyValueVisitor.Dump(LogEntry entry)` + - 内部执行 `return entry.Accept(this);`,由于 `entry` 的运行时类型是 `CallLogEntry`,发生多态分派,调用 `CallLogEntry` 中被 override 的 `Accept` ++ `TResult CallLogEntry.Accept(ILogEntryVisitor visitor)` + - 内部执行 `return visitor.Visit(this);`,此处 `this` 的编译时类型是 `CallLogEntry`,于是通过重载分派选中 `KeyValueVisitor.Visit(CallLogEntry entry)` ++ `Dictionary KeyValueVisitor.Visit(CallLogEntry entry)` + - 构造并返回保存了 `LineNo`、`Timestamp`、`PodName`、`Severity`、`EventType`、`RequestId`、`TargetService`、`DurationMs` 的 `Dictionary` + +这里正是访问者模式的「双重分派(double dispatch)」:第一重由 `entry.Accept(this)` 按 `entry` 的**运行时类型**分派到 `CallLogEntry.Accept`;第二重由 `visitor.Visit(this)` 按 `this` 的**编译时类型**(`CallLogEntry`)分派到 `KeyValueVisitor.Visit(CallLogEntry)` 的重载,从而对外部屏蔽了具体子类,却仍能对每种日志执行不同的行为。 + +# (Q1.3.b) +根据TODO框架和guidance.md完成任务××.AI能给出达成任务要求的代码并自行测试验证。有时候AI会有过度、无效兜底的问题,在这次作业中基本没有出现 \ No newline at end of file diff --git a/docs/02-multithreading/0b487ce620ea8067dee251d1968508c4.png b/docs/02-multithreading/0b487ce620ea8067dee251d1968508c4.png new file mode 100644 index 0000000..a6e8c52 Binary files /dev/null and b/docs/02-multithreading/0b487ce620ea8067dee251d1968508c4.png differ diff --git a/docs/02-multithreading/QQ_1785484373806.png b/docs/02-multithreading/QQ_1785484373806.png new file mode 100644 index 0000000..e48e4d8 Binary files /dev/null and b/docs/02-multithreading/QQ_1785484373806.png differ diff --git a/docs/02-multithreading/QQ_1785484396152-1.png b/docs/02-multithreading/QQ_1785484396152-1.png new file mode 100644 index 0000000..2fdd35c Binary files /dev/null and b/docs/02-multithreading/QQ_1785484396152-1.png differ diff --git a/docs/02-multithreading/QQ_1785484396152-2.png b/docs/02-multithreading/QQ_1785484396152-2.png new file mode 100644 index 0000000..2fdd35c Binary files /dev/null and b/docs/02-multithreading/QQ_1785484396152-2.png differ diff --git a/docs/02-multithreading/QQ_1785484396152-3.png b/docs/02-multithreading/QQ_1785484396152-3.png new file mode 100644 index 0000000..2fdd35c Binary files /dev/null and b/docs/02-multithreading/QQ_1785484396152-3.png differ diff --git a/docs/02-multithreading/QQ_1785484396152.png b/docs/02-multithreading/QQ_1785484396152.png new file mode 100644 index 0000000..2fdd35c Binary files /dev/null and b/docs/02-multithreading/QQ_1785484396152.png differ diff --git a/docs/02-multithreading/QQ_1785484923024.png b/docs/02-multithreading/QQ_1785484923024.png new file mode 100644 index 0000000..16e9353 Binary files /dev/null and b/docs/02-multithreading/QQ_1785484923024.png differ diff --git a/docs/02-multithreading/QQ_1785485348097.png b/docs/02-multithreading/QQ_1785485348097.png new file mode 100644 index 0000000..7e5730b Binary files /dev/null and b/docs/02-multithreading/QQ_1785485348097.png differ diff --git a/docs/02-multithreading/QQ_1785485430494.png b/docs/02-multithreading/QQ_1785485430494.png new file mode 100644 index 0000000..513a656 Binary files /dev/null and b/docs/02-multithreading/QQ_1785485430494.png differ diff --git a/docs/02-multithreading/QQ_1785485472946.png b/docs/02-multithreading/QQ_1785485472946.png new file mode 100644 index 0000000..aa12f07 Binary files /dev/null and b/docs/02-multithreading/QQ_1785485472946.png differ diff --git a/docs/02-multithreading/QQ_1785485519059.png b/docs/02-multithreading/QQ_1785485519059.png new file mode 100644 index 0000000..c67ddb1 Binary files /dev/null and b/docs/02-multithreading/QQ_1785485519059.png differ diff --git a/docs/02-multithreading/e87b65691415cb6e03e2e985cded5bfb.png b/docs/02-multithreading/e87b65691415cb6e03e2e985cded5bfb.png new file mode 100644 index 0000000..80d3df7 Binary files /dev/null and b/docs/02-multithreading/e87b65691415cb6e03e2e985cded5bfb.png differ diff --git a/docs/02-multithreading/report.md b/docs/02-multithreading/report.md new file mode 100644 index 0000000..00bc2ac --- /dev/null +++ b/docs/02-multithreading/report.md @@ -0,0 +1,188 @@ +# 02-multithreading 实验报告 + +## 一、功能介绍 + +本节在 `01-basic` 的单文件日志解析基础上,实现了一个**目录级别的并行日志分析器**,并配有一个简易的交互式控制台界面。整体由三部分组成: + +| 文件 | 任务 | 作用 | +| :--- | :--- | :--- | +| `LogAnalyzer/WorkQueue.cs` | T2.1 | 基于非线程安全 `Queue` 自造的**线程安全阻塞队列** | +| `LogAnalyzer/LogFileAnalyzer.cs` | T2.2 | 扫描目录、调度多线程并行解析、保存结果 | +| `LocalCli/Program.cs` | T2.3 | 与用户交互的控制台菜单,串联上述能力 | + +### 1. 线程安全队列 `WorkQueue`(T2.1) + +共享变量为内部的 `Queue _items` 与「是否结束放入」标记 `_isCompleted`,两者统一用 `lock(_items)` 这同一个互斥量保护。这是一个带「结束放入」语义的无限容量生产者—消费者问题: + +- `Enqueue`:加锁后入队,并 `Monitor.Pulse`(signal)唤醒一个等待中的消费者;若已 `CompleteAdding` 则抛 `InvalidOperationException`。 +- `CompleteAdding`:置位 `_isCompleted`,并 `Monitor.PulseAll`(broadcast)唤醒**全部**正在等待的消费者,使其能够正常退出而不是永远阻塞。 +- `TryDequeue`:加锁后用 **`while`** 循环判断「队列空 且 未结束」才 `Monitor.Wait`;被唤醒后重新检查条件。队列非空则取出返回 `true`,否则(空且已结束)返回 `false` 并把 `item` 置为 `default`。 + +### 2. 并行日志分析 `LogFileAnalyzer`(T2.2) + +- **目录扫描**:`ChangeDirectory` 中 `Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.TopDirectoryOnly)` 扫描当前目录下所有 `.log` 文件,并把每个文件以 `NotAnalyzed` 状态登记进 `_analysisResults`。 +- **状态位 `_isAnalyzing`**:`AnalyzeFiles` 进入分析前置 `true`,在 `try/finally` 的 `finally` 里**加锁**复位为 `false`,保证异常时也能复位。该标志保证同一时刻只允许一个分析任务进行,其余并发请求抛 `InvalidOperationException`。 +- **`RunWorkers`**:先把 `State == NotAnalyzed` 的文件筛选出来,跳过已 `Succeeded`/`Failed` 的文件以节省计算资源,未知文件抛 `InvalidOperationException`,用主线程作为生产者把待解析文件 `Enqueue` 进 `WorkQueue` 后 `CompleteAdding`;再开启 `degreeOfParallelism` 个 worker 线程(入口方法 `WorkerMain`),最后 `Join` 等待全部 worker 结束。 + - `degreeOfParallelism == 0` 表示取 `Environment.ProcessorCount`;并按 `[1, 文件数]` 夹取,避免开多余空转线程。 +- **`WorkerMain`**:每个消费者循环 `TryDequeue` 取文件,用各自的 `LogFileParser` + `StreamReader` 解析;解析成功 → `Succeeded`,解析抛异常 → `Failed`(写 `ErrorMessage`、空 `Entries`)。最后**加 `_syncRoot` 锁**把结果写回共享的 `_analysisResults` 字典。 + - `parser.Parse(reader).ToList()` 中的 `ToList()` 用于强制立即求值:`Parse` 是 `yield return` 的惰性迭代器,若不立刻物化,异常会推迟到 `try` 之外才发生而无法被捕获。 +- **`TryGetAnalysisResult`**:加锁查 `_analysisResults`,存在则返回 `true` 并输出结果,否则返回 `false`。 + +### 3. 控制台交互 `LocalCli/Program.cs`(T2.3) + +菜单提供 6 个功能,`LogFileAnalyzer` 对错误输入,CLI 层把这些异常兜住并提示用户重新输入,保证非法输入不会让程序崩溃。 + +| 选项 | 功能 | 实现 | +| :--: | :--- | :--- | +| —— | `InputDirectory` | 输入目录构造 `analyzer`;目录不存在→提示重输(`ChangeDirectory` 返回 `false`),路径非法→捕获 `ArgumentException` 提示重输 | +| 1 | `ShowLogFiles` | 调用 `GetLogFiles()` 列出目录中全部 `.log` 文件名 | +| 2 | `AnalyzeFiles` | 输入逗号分隔的文件名,`Split` 时去空白,调用 `AnalyzeFiles(0, ...)`;捕获 `ArgumentException`/`InvalidOperationException` | +| 3 | `AnalyzeAll` | 调用 `AnalyzeAll(0)` 分析全部;捕获 `InvalidOperationException` | +| 4 | `GetAnalysisResult` | 输入文件名查结果:不存在→提示;`NotAnalyzed`→提示先分析;`Succeeded`→用 `KeyValueVisitor.Dump` 逐条输出;`Failed`→输出 `ErrorMessage` | +| 5 | ChangeDirectory | 重新输入目录(复用 `InputDirectory`) | +| 6 | Exit | 退出 | + +非法的菜单输入(非数字 / 超出范围)会被 `int.Parse` 的异常捕获或 `default` 分支拦截,提示重输,不会崩溃。 + +--- + +## 二、功能演示 + +### 启动 + 查看日志文件列表 + +![alt text](./0b487ce620ea8067dee251d1968508c4.png) + +### 分析指定文件 + +![alt text](./QQ_1785484396152-3.png) + +### 查看分析成功的结果,查询尚未分析的文件 + +![alt text](./e87b65691415cb6e03e2e985cded5bfb.png) + +### 分析全部 + 查看失败文件的错误信息 + +![alt text](./QQ_1785484923024.png) + + +## 三、鲁棒性测试 + +### 不存在的目录 + +![alt text](./QQ_1785485348097.png) + +### 非法的菜单输入 + +![alt text](./QQ_1785485472946.png) + +### 分析不存在的文件 + +![alt text](./QQ_1785485519059.png) + +### 查询不存在 / 未分析的文件 + +![alt text](./QQ_1785485430494.png) + +--- + +## 四、问答题 + +### (Q2.1) + +共享变量有两个: + +- `Queue _items`:真正存放元素的内部队列; +- `bool _isCompleted`:标记是否已结束放入(`CompleteAdding` 是否被调用过)。 + +两者都通过**以 `_items` 这个引用对象本身作为互斥量**来保护——所有对它们的读写都放在 `lock(_items)` 临界区内: + +```csharp +public void Enqueue(T item) +{ + lock (_items) + { + if (_isCompleted) throw new InvalidOperationException(...); + _items.Enqueue(item); + Monitor.Pulse(_items); + } +} +``` +同步关系(消费者等待 / 生产者唤醒)也建立在同一个 `_items` 上:消费者 `Monitor.Wait(_items)` 释放锁并休眠,生产者用 `Monitor.Pulse`(signal)/ `Monitor.PulseAll`(broadcast)唤醒。这是 C# `Monitor` 实现的 MESA 模型条件变量。 + + `LogFileAnalyzer` 中的共享变量有: + +- `string? _currentDirectory`:当前日志目录; +- `bool _isAnalyzing`:是否正在分析; +- `Dictionary _logFiles`:文件名到 `FileInfo` 的映射; +- `Dictionary _analysisResults`:文件名到分析结果的映射。 + +它们统一由一个专用的互斥量对象 `private readonly object _syncRoot = new();` 保护,所有访问都放在 `lock(_syncRoot)` 内(`ChangeDirectory`、`GetLogFiles`、`TryGetAnalysisResult`,以及 worker 写回结果时): + +```csharp +lock (_syncRoot) +{ + _analysisResults[file.Name] = result; +} +``` + +`RunWorkers` 中还有一个局部构造的 `WorkQueue` 实例,被主线程(生产者)和各 worker(消费者)共享,但它由 `WorkQueue` **内部自己的 `_items` 锁**保护,属于另一套独立的互斥机制,与 `_syncRoot` 无关。`IsAnalyzing`、`IsCompleted` 等属性的 getter 也都通过加锁读取,避免读到未同步的值。 + +用 `if` 而非 `while` 在虚假唤醒下的后果: + +以无限容量生产者—消费者为例,若消费者写成: + +```csharp +lock (mtx) +{ + if (buffer == 0) // 用 if + { + Monitor.Wait(mtx); + } + buffer -= 1; // 直接取用商品 +} +``` + +当发生虚假唤醒(`Wait` 在没有人 `Pulse` 的情况下自行返回)时: + +- 线程被唤醒后**不再重新检查** `buffer == 0`,直接执行 `buffer -= 1`; +- 但此时 `buffer` 仍为 0(根本没生产出商品),于是「取走了一个不存在的商品」,`buffer` 变成 −1,状态不变量被破坏,出现逻辑错误。 + +更严重的是多消费者下的竞争:MESA 模型中,线程被 `Pulse` 唤醒后并不会立即拿到锁,而要重新去抢锁;在它重新拿到锁之前,另一个消费者可能已经把唯一的商品取走了(也可能是纯粹的虚假唤醒)。若用 `if`,这个被唤醒的消费者不会再检查条件,照样去取商品,于是出现「一个商品被消费两次」或「消费了空仓库」的错误。 + +而用 `while`: + +```csharp +while (buffer == 0) { Monitor.Wait(mtx); } +``` + +被唤醒后会**再次判断条件**,若仓库仍为空(无论是虚假唤醒还是被别的消费者抢先),就继续 `Wait`,只有确实非空时才取用——这才是正确的。 + + +### (Q2.2) + +在 `ChangeDirectory` 方法中,这段代码完成了扫描: + +```csharp +var logFiles = Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.TopDirectoryOnly) + .Select(filePath => Path.GetFileName(filePath)) + .OrderBy(fileName => fileName); +foreach (var fileName in logFiles) +{ + _logFiles.Add(fileName, new FileInfo(Path.Join(_currentDirectory, fileName))); + _analysisResults.Add(fileName, new AnalysisResult(...)); +} +``` + +其中真正「扫描目录里所有 `.log` 文件」的是 `Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.TopDirectoryOnly)`:第一个参数是目录、第二个是匹配模式 `*.log`、第三个 `SearchOption.TopDirectoryOnly` 表示只扫描当前目录(不进入子目录);后面的 `Select`/`OrderBy` 只是把得到的路径取出文件名并排序。 + +若要递归扫描所有子目录,把第三个参数改为 `SearchOption.AllDirectories`,即可让 `EnumerateFiles` 递归遍历所有子目录、子子目录……: + +```csharp +Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.AllDirectories) +``` + +一个需要注意的细节:当前代码用 `Path.GetFileName(filePath)`(仅文件名)作为 `_logFiles` / `_analysisResults` 的键。非递归时文件名不会重复,没有问题;但改为递归后,不同子目录下可能存在同名文件(例如两个子目录里都有 `20260701.log`),会造成字典键冲突、后扫到的覆盖先扫到的。因此递归版本更适合改用「相对路径」作为键,例如 `Path.GetRelativePath(directoryPath, filePath)`,以避免重名冲突。 + +### (Q2.3) + +根据TODO框架和guidance.md完成任务××.AI能给出达成任务要求的代码并自行测试验证。有时候AI会有过度、无效兜底的问题,在这次作业中基本没有出现 \ No newline at end of file diff --git a/docs/03-async-grpc/QQ_1785571614496.png b/docs/03-async-grpc/QQ_1785571614496.png new file mode 100644 index 0000000..0af46bb Binary files /dev/null and b/docs/03-async-grpc/QQ_1785571614496.png differ diff --git a/docs/03-async-grpc/QQ_1785571631493.png b/docs/03-async-grpc/QQ_1785571631493.png new file mode 100644 index 0000000..b11feea Binary files /dev/null and b/docs/03-async-grpc/QQ_1785571631493.png differ diff --git a/docs/03-async-grpc/QQ_1785571665044.png b/docs/03-async-grpc/QQ_1785571665044.png new file mode 100644 index 0000000..05ce81d Binary files /dev/null and b/docs/03-async-grpc/QQ_1785571665044.png differ diff --git a/docs/03-async-grpc/QQ_1785571809037.png b/docs/03-async-grpc/QQ_1785571809037.png new file mode 100644 index 0000000..2e36b15 Binary files /dev/null and b/docs/03-async-grpc/QQ_1785571809037.png differ diff --git a/docs/03-async-grpc/QQ_1785572331363.png b/docs/03-async-grpc/QQ_1785572331363.png new file mode 100644 index 0000000..9a450c7 Binary files /dev/null and b/docs/03-async-grpc/QQ_1785572331363.png differ diff --git a/docs/03-async-grpc/QQ_1785572349080.png b/docs/03-async-grpc/QQ_1785572349080.png new file mode 100644 index 0000000..09585e3 Binary files /dev/null and b/docs/03-async-grpc/QQ_1785572349080.png differ diff --git a/docs/03-async-grpc/QQ_1785572368644.png b/docs/03-async-grpc/QQ_1785572368644.png new file mode 100644 index 0000000..dad6c44 Binary files /dev/null and b/docs/03-async-grpc/QQ_1785572368644.png differ diff --git a/docs/03-async-grpc/QQ_1785572407446.png b/docs/03-async-grpc/QQ_1785572407446.png new file mode 100644 index 0000000..d10de54 Binary files /dev/null and b/docs/03-async-grpc/QQ_1785572407446.png differ diff --git a/docs/03-async-grpc/QQ_1785572418459.png b/docs/03-async-grpc/QQ_1785572418459.png new file mode 100644 index 0000000..5935949 Binary files /dev/null and b/docs/03-async-grpc/QQ_1785572418459.png differ diff --git a/docs/03-async-grpc/QQ_1785762240394.png b/docs/03-async-grpc/QQ_1785762240394.png new file mode 100644 index 0000000..15347e0 Binary files /dev/null and b/docs/03-async-grpc/QQ_1785762240394.png differ diff --git a/docs/03-async-grpc/report.md b/docs/03-async-grpc/report.md new file mode 100644 index 0000000..641f73e --- /dev/null +++ b/docs/03-async-grpc/report.md @@ -0,0 +1,127 @@ +# 03-async-grpc 实验报告 + +## 一、T3.2 功能说明 + +`RemoteCli` 是一个对接 Agent gRPC 服务的远程控制台客户端,其交互逻辑与上一节的 `LocalCli` 基本一致,区别在于:所有对 `LogFileAnalyzer` 的本地函数调用都被替换为对应的 **异步 gRPC 调用**,即调用方法名带 `Async` 后缀的版本。 + +### 1. 实现的功能 + +| 菜单选项 | 功能 | 对应的 gRPC 异步调用 | +| :------: | :--- | :--- | +| 启动时 | 输入并切换日志目录 | `ChangeDirectoryAsync` | +| 1 | 列出当前目录中的日志文件 | `GetLogFilesAsync` | +| 2 | 解析指定的日志文件(可指定并行度) | `AnalyzeFilesAsync` | +| 3 | 解析当前目录中的全部日志文件 | `AnalyzeAllAsync` | +| 4 | 查询指定文件的分析结果(**流式**返回) | `GetAnalysisResult`(`ResponseStream.ReadAllAsync`) | +| 5 | 切换日志目录 | `ChangeDirectoryAsync` | +| 6 | 退出 | —— | + +### 2. 关键实现要点 + +1. **全异步调用**:按照本节要求,所有 gRPC 调用均使用 `Async` 版本。对于流式接口 `GetAnalysisResult`,使用 `await foreach (var response in call.ResponseStream.ReadAllAsync())` 逐条读取服务端流式返回的 `GetAnalysisResultResponse`。 + +2. **流式结果的解析**:根据响应中的 `PayloadCase` 分别处理: + - `Header`:根据 `State`(`NotAnalyzed` / `Failed` / `Succeeded`)给出对应的提示;只有 `Succeeded` 时才继续接收后续的日志条目。 + - `LogEntry`:通过 `GrpcTypeConverter.ConvertFromGrpc` 转回内部的 `LogEntry` 类型,再用 `KeyValueVisitor` 以键值对形式打印。 + +3. **切换目录后即时反馈**:利用 `ChangeDirectoryResponse` 中额外返回的 `current_directory` 与 `file_names` 字段,在切换目录成功后立刻打印 Agent 的完整路径与目录中的全部日志文件,方便确认。 + +4. **可配置并行度**:`AnalyzeFiles` / `AnalyzeAll` 在执行前会通过 `ReadDegreeOfParallelism` 让用户输入并行度(`0` 表示使用 `ProcessorCount`),通过 `ReadFileNames` 读取待解析的文件名列表。 + +5. **健壮性(重点)**:Agent 作为常驻服务,绝不应因用户非法输入或内部错误而崩溃。客户端同样做了充分的容错: + - 所有 gRPC 调用均捕获 `RpcException`(网络层错误),输出友好提示而非崩溃。 + - 服务端返回的 `OperationStatusMessage` 中 `success == false` 时,将错误码与错误信息展示给用户,并允许其重新输入。 + - 菜单输入、并行度输入、文件名输入均做了非法值校验与重试。 + +### 3. 运行方式 + +需要**同时**运行 Agent(服务端)与 RemoteCli(客户端)两个程序: + +1. **启动 Agent**:将 `LogAnalyzerAgent` 设为启动项目并运行(或在 `LogAnalyzerAgent` 目录下执行 `dotnet run`),它会在 `http://localhost:5000` 上监听 gRPC 服务。 +2. **启动 RemoteCli**:运行 `RemoteCli`,默认连接 `http://localhost:5000`;也可通过命令行参数或环境变量 `LOG_ANALYZER_AGENT_ADDRESS` 指定 Agent 地址: + + ```bash + dotnet run --project RemoteCli -- http://localhost:5000 + ``` + +--- + +## 二、功能演示截图 + +### 1. 连接 Agent 并切换目录 + +![alt text](./QQ_1785571614496.png) + +### 2. 列出日志文件(选项 1) + +![alt text](./QQ_1785571631493.png) + +### 3. 解析指定文件并查看结果(选项 2 + 选项 4) + +![alt text](./QQ_1785571665044.png) + +### 4. 解析全部文件并查看多日志文件结果(选项 3 + 选项 4) + +![alt text](./QQ_1785571809037.png) + +--- + +## 三、鲁棒性测试截图 + +### 1. 非法目录名 + +![alt text](./QQ_1785572331363.png) + +### 2. 非法菜单选项 + +![alt text](./QQ_1785572349080.png) +90c'x +### 3. 查询不存在的文件 + +![alt text](./QQ_1785572368644.png) + +### 4. 解析不存在的文件 + +### 5. 非法并行度输入 +![alt text](./QQ_1785572407446.png) + + +### 6. 空文件名输入 + +![alt text](./QQ_1785572418459.png) + +--- + +## 四、问答题 + +### (Q3.1) + +**区别:** + +1. **调用方式的本质变化**:非网络程序中,函数调用都是同进程内的本地调用,传参、返回都直接在内存中进行;而网络应用中,跨机器/跨进程的交互变成了远程过程调用(RPC)。表面上 `client.GetLogFilesAsync()` 看起来像普通方法,但背后实际上经历了一次完整的网络往返。 +2. **数据需要序列化**:本地调用直接传递对象引用;网络调用则必须把数据序列化为字节流(本节中是 Protobuf),到达对端再反序列化。这就要求两端有一套共同的接口描述(IDL,即 `.proto` 文件),并且在内部 C# 类型与 Protobuf 消息类型之间编写转换层(本节的 `GrpcTypeConverter` / `GrpcLogEntryVisitor`)。 +3. **必须采用异步编程模型**:网络 I/O 是典型的 I/O 密集场景,若用同步调用,线程会在等待网络响应时被白白阻塞。本节的 `RemoteCli` 因此全部使用 `async` / `await`,在等待响应时让出线程,这正是异步编程相比多线程的优势所在。 +4. **需要同时启动、联合调试两个程序**:以往调试单个可执行文件即可;网络应用必须同时跑起服务端(Agent)和客户端(RemoteCli),且 Visual Studio 一次只能调试一个,另一个要手动启动,调试方式发生了变化。 +5. **状态分布在多个进程**:Agent 是有状态的单例服务(保存当前目录、分析结果),客户端只是远程地读取/修改这些状态,状态不再集中在一个进程内。 + +**额外的难点:** + +1. **网络不可靠**:连接可能失败、请求可能超时、对端可能宕机。必须捕获 `RpcException` 并做容错处理,而非像本地调用那样假定“调用了一定会返回”。本节的鲁棒性测试就体现了这一点。 +2. **错误定位困难**:一次调用失败,可能是客户端参数错、可能是网络层、可能是序列化/反序列化、也可能是服务端逻辑。错误来源横跨两端,排查时需要分别查看客户端与服务端的输出,调试成本显著上升。 +3. **类型系统的割裂与一致性维护**:两端可能用不同语言、不同类型表示同一概念。一旦 `.proto` 改动,两端的生成代码与转换逻辑都要同步更新,否则会出现字段对不上的隐蔽 bug。 +4. **并发与共享状态的同步**:Agent 作为常驻服务,可能同时收到多个客户端请求,对共享状态(当前目录、分析结果)的访问需要加锁(本节的 `LogFileAnalyzer` 用 `_syncRoot` 保护),还要处理“分析进行中再次请求分析”这类并发冲突。 +5. **部署与环境配置**:要关心端口、监听地址(`localhost` 仅本机、`0.0.0.0` 对外)、HTTP/2 协议、CORS、防火墙等,这些在非网络程序里几乎不存在。 +6. **安全性**:网络服务暴露在外,需要考虑鉴权、传输加密、防止恶意请求与 DDoS 等,而本地程序一般无需考虑。 + +**额外的复杂之处:** + +1. 需要额外学习并理解一整套协议栈知识:Protobuf 的消息定义与 `oneof`、gRPC 的四种调用模式(本节用到了服务端流式)、HTTP/2 等。 +2. 需要理解依赖注入、服务注册等框架级概念(本节用 ASP.NET 的 `AddSingleton` 注册有状态服务),这对初学者是不小的认知负担。 +3. 流式 RPC 的处理比一次性返回更复杂:要逐条读取、区分 `header` 与 `log_entry`、处理“文件不存在/未分析/失败/成功”等多种情形。 +4. 调试反馈链路变长:改一处接口往往要重新生成代码、重启服务端、再重启客户端,迭代效率比单机程序低。 + +总的来说,网络应用的核心复杂度来自于**“分布”**二字——计算与状态被分散到了通过网络连接的不同节点上,由此衍生出序列化、异步、容错、并发、安全等一系列非网络程序所没有的问题。 + +### (Q2.2) + +根据TODO框架和guidance.md完成任务××.AI能给出达成任务要求的代码并自行测试验证。有时候AI会有过度、无效兜底的问题,在这次作业中基本没有出现 diff --git a/docs/04-avalonia/QQ_1785813088042.png b/docs/04-avalonia/QQ_1785813088042.png new file mode 100644 index 0000000..6b40279 Binary files /dev/null and b/docs/04-avalonia/QQ_1785813088042.png differ diff --git a/docs/04-avalonia/QQ_1785817655622.png b/docs/04-avalonia/QQ_1785817655622.png new file mode 100644 index 0000000..e22cfce Binary files /dev/null and b/docs/04-avalonia/QQ_1785817655622.png differ diff --git a/docs/04-avalonia/QQ_1785817741740.png b/docs/04-avalonia/QQ_1785817741740.png new file mode 100644 index 0000000..c1b6506 Binary files /dev/null and b/docs/04-avalonia/QQ_1785817741740.png differ diff --git a/docs/04-avalonia/QQ_1785817782805.png b/docs/04-avalonia/QQ_1785817782805.png new file mode 100644 index 0000000..0ffee3c Binary files /dev/null and b/docs/04-avalonia/QQ_1785817782805.png differ diff --git a/docs/04-avalonia/QQ_1785817816941.png b/docs/04-avalonia/QQ_1785817816941.png new file mode 100644 index 0000000..268f941 Binary files /dev/null and b/docs/04-avalonia/QQ_1785817816941.png differ diff --git a/docs/04-avalonia/QQ_1785818118027.png b/docs/04-avalonia/QQ_1785818118027.png new file mode 100644 index 0000000..decf909 Binary files /dev/null and b/docs/04-avalonia/QQ_1785818118027.png differ diff --git a/docs/04-avalonia/QQ_1785818173900.png b/docs/04-avalonia/QQ_1785818173900.png new file mode 100644 index 0000000..489f0c8 Binary files /dev/null and b/docs/04-avalonia/QQ_1785818173900.png differ diff --git a/docs/04-avalonia/QQ_1785819031379.png b/docs/04-avalonia/QQ_1785819031379.png new file mode 100644 index 0000000..b242ea5 Binary files /dev/null and b/docs/04-avalonia/QQ_1785819031379.png differ diff --git a/docs/04-avalonia/QQ_1785827003500.png b/docs/04-avalonia/QQ_1785827003500.png new file mode 100644 index 0000000..4ab6dd5 Binary files /dev/null and b/docs/04-avalonia/QQ_1785827003500.png differ diff --git a/docs/04-avalonia/QQ_1785827109255.png b/docs/04-avalonia/QQ_1785827109255.png new file mode 100644 index 0000000..2409aea Binary files /dev/null and b/docs/04-avalonia/QQ_1785827109255.png differ diff --git a/docs/04-avalonia/QQ_1785827490152-1.png b/docs/04-avalonia/QQ_1785827490152-1.png new file mode 100644 index 0000000..0b06b4b Binary files /dev/null and b/docs/04-avalonia/QQ_1785827490152-1.png differ diff --git a/docs/04-avalonia/QQ_1785827490152.png b/docs/04-avalonia/QQ_1785827490152.png new file mode 100644 index 0000000..0b06b4b Binary files /dev/null and b/docs/04-avalonia/QQ_1785827490152.png differ diff --git a/docs/04-avalonia/QQ_1785827569651.png b/docs/04-avalonia/QQ_1785827569651.png new file mode 100644 index 0000000..1e025bd Binary files /dev/null and b/docs/04-avalonia/QQ_1785827569651.png differ diff --git a/docs/04-avalonia/QQ_1785827606708.png b/docs/04-avalonia/QQ_1785827606708.png new file mode 100644 index 0000000..97ffadf Binary files /dev/null and b/docs/04-avalonia/QQ_1785827606708.png differ diff --git a/docs/04-avalonia/report.md b/docs/04-avalonia/report.md new file mode 100644 index 0000000..1963f28 --- /dev/null +++ b/docs/04-avalonia/report.md @@ -0,0 +1,159 @@ +# 04-avalonia 实验报告 + +## 一、T4.1 功能说明 + +`LogAnalyzerClient` 是一个基于 [Avalonia UI]的跨平台图形界面客户端,用于替代上一节的控制台客户端 `RemoteCli`。它包含 `RemoteCli` 的全部功能:连接 Agent、切换日志目录、刷新文件列表、(按并行度)分析选中 / 全部 / 右键单个文件、以及流式查看分析结果。整个客户端采用 MVVM 模式编写,借助 `CommunityToolkit.Mvvm` 的 `[ObservableProperty]` 与 `[RelayCommand]` 源生成器大幅减少了样板代码。 + +### 1. 实现的功能 + +| 入口 | 功能 | 对应的 gRPC 异步调用 | 说明 | +| :--- | :--- | :--- | :--- | +| `File → Connect...` | 连接到 Agent | `PingAsync` | 通过工厂 `AppService.ClientFactory.CreateClient` 创建 gRPC Client,再 `Ping` 验证连通性 | +| `Change Directory` 按钮 | 切换 Agent 的日志目录 | `ChangeDirectoryAsync` | 框架已实现,切换成功后自动刷新文件列表 | +| `Refresh` 按钮 / `File → Refresh` | 刷新日志目录文件列表 | `GetLogFilesAsync` | 用返回的 `file_names` 重建 `LogFiles` | +| `Analyze → Selected` 按钮 | 分析多选的若干文件 | `AnalyzeFilesAsync` |参数为 `SelectedFiles`(由 `LogFileListBox_SelectionChanged` 维护) | +| `Analyze → All` 按钮 | 分析当前目录全部文件 | `AnalyzeAllAsync` | 在 XAML 中新增 `All` 按钮并绑定 `AnalyzeAllCommand` | +| 右键菜单 `Analyze File` | 分析右键选中的单个文件 | `AnalyzeFilesAsync` | 参数为当前 `SelectedLogFile` | +| 右键菜单 `View Analysis Results` | 查看所选文件分析结果 | `GetAnalysisResult`(`ResponseStream.ReadAllAsync`) | 逐条接收并填充 `ResultEntries` | +| Analysis Result 列表 | 展示分析结果 | —— | `LogFields.Summary` 负责每行的文本格式 | + +### 2. 关键实现要点 + +1. **全异步调用**:图形界面程序只有单一的 UI 线程负责渲染与响应,任何阻塞都会让程序看起来「卡死」。因此所有的 gRPC 调用都使用 `Async` 版本并 `await`;对于服务端流式接口 `GetAnalysisResult`,使用 `await foreach (var response in call.ResponseStream.ReadAllAsync())` 逐条读取。 + +2. **统一的异常兜底 `WithClientNotNull`**:所有需要客户端的命令都套在 `WithClientNotNull` 中。它先检查是否已连接(`_client is null` 时弹出提示),再用 `try/catch (Exception)` 兜住一切 gRPC 网络异常与内部错误,转为消息框,保证「GUI 程序绝不应因用户非法输入或内部错误而崩溃」。 + +3. **输入校验前置**: + - 并行度通过 `TryGetDegreeOfParallelism` 校验,必须是**非负整数**(`0` 表示使用 `ProcessorCount`),非法时弹框提示并不发起请求。 + - 「分析选中文件」会检查 `SelectedFiles.Count == 0`;「右键分析 / 查看结果」会检查 `SelectedLogFile is null`,避免空引用。 + +4. **服务端返回状态二次检查**:除了捕获网络层的 `RpcException`,还对每个响应的 `OperationStatusMessage.Success` 做检查,失败时将 `Code: Message` 通过消息框展示给用户。 + +5. **流式结果解析**(`GetAnalysisResultAsync`):根据 `PayloadCase` 分别处理: + - `Header`:按 `State` 分三种情况——`NotAnalyzed` 显示「尚未分析」提示;`Failed` 显示 `Analysis failed: {ErrorMessage}`;`Succeeded` 继续接收后续条目。 + - `LogEntry`:经 `GrpcTypeConverter.ConvertFromGrpc` 转回内部 `LogEntry`,再用 `KeyValueVisitor` 转成键值对,包装成 `LogFields` 加入 `ResultEntries`。每次查看前先 `ResultEntries.Clear()`,避免与上次结果混淆。 + +6. **`LogFields.Summary` 的显示格式**: + - 普通日志条目:`序号 | Key: Value, Key: Value, ...`(与示例截图一致,序号取日志的 `LineNo`)。 + - 错误 / 未分析:直接展示 `ErrorMessage` 文本(如 `File 'xxx' has not been analyzed yet.`)。 + +7. **新增 `All` 按钮**:在 `MainView.axaml` 的分析操作 `Grid` 中,把列定义从 `Auto,*,Auto,Auto` 扩展为 `Auto,*,Auto,Auto,Auto`,在 `Selected` 之后追加 `All` 按钮并绑定 `AnalyzeAllCommand`。 + +### 3. 运行方式 + +GUI 客户端需要**同时**运行 Agent(服务端)与客户端两个程序。下面所有截图均按以下命令启动: + +```bash +# 终端 1:启动 Agent(gRPC 服务端,监听 http://localhost:5000) +dotnet run --project src/LogAnalyzerAgent + +# 终端 2:启动 Avalonia 桌面客户端 +dotnet run --project src/LogAnalyzerClient/LogAnalyzerClient.Desktop +``` + +> 客户端启动后,点击菜单 `File → Connect...`,在弹窗中输入 `http://localhost:5000` 即可连接。 +> 测试所用日志目录(任选其一填入 *Directory Path* 输入框): +> - `src/dataset`(含 `basic.log`、`basic-fail.log`、`basic-multiple.log`) +> - `src/dataset/multiple-logs`(含 30 个 `20260701.log` ~ `20260730.log`) + +--- + +## 二、功能演示截图(命令 / 操作标注) + +### 1. 启动并连接到 Agent + +![alt text](./QQ_1785817655622.png) + +### 2. 切换目录并刷新出文件列表 + +在 *Directory Path* 输入框输入绝对路径 → 点击 `Change Directory` +![alt text](./QQ_1785813088042.png) + +### 3. 多选文件并分析(Analyze → Selected) + +![alt text](./QQ_1785817741740.png) + +### 4. 分析全部文件(Analyze → All) + +![alt text](./QQ_1785817782805.png) + +### 5. 右键分析单个文件(右键菜单 Analyze File) + +![alt text](./QQ_1785817816941.png) + +### 6. 查看分析结果(右键菜单 View Analysis Results)—— 成功 + +选中已分析成功的文件(如 `basic-multiple.log`)→ 右键 → 点击 `View Analysis Results(V)`(调用流式 `GetAnalysisResult`) + +![alt text](./QQ_1785818118027.png) + +### 7. 查看分析结果 —— 失败 / 尚未分析 + +失败:选中 `basic-fail.log`,用 `All` 或 `Selected` 分析它(会解析失败)→ 右键 → `View Analysis Results(V)` +未分析:连接并切换目录后,**不**进行分析,直接选中某文件 → 右键 → `View Analysis Results(V)` + +![alt text](./QQ_1785818173900.png) +![alt text](./QQ_1785819031379.png) +--- + +## 三、鲁棒性测试截图(命令 / 操作标注) + +### 1. 未连接 Agent 就执行操作 + +**不**点击 `Connect...`,直接点击 `Refresh` / `Change Directory` / `Selected` 等任意按钮 +![alt text](./QQ_1785827003500.png) + +### 2. 连接到不存在的 Agent 地址 + +![alt text](./QQ_1785827109255.png) + +### 3. 非法的目录路径 + +![alt text](./QQ_1785827490152-1.png) + +### 4. 非法的并行度输入 + +![alt text](./QQ_1785827569651.png) + +### 5. 未选中文件就点击 Selected + +![alt text](./QQ_1785827606708.png) + +--- + +## 四、问答题 + +### (Q4.1) + +**GUI 应用与控制台应用的区别:** + +1. **交互范式不同**:控制台应用是「线性的一问一答」——程序主动 `Console.ReadLine` 等待输入,流程是预先确定的;GUI 应用是「事件驱动」——用户可以在任意时刻点击任意按钮、输入任意内容,程序必须随时响应,控制流不再线性。本次实现中,每个按钮 / 菜单项都被绑定到一个独立的 `ICommand`,由用户决定何时触发、以何种顺序触发。 +2. **关注点分离的要求不同**:控制台应用里输入、逻辑、输出往往混在一个 `Main` 里;GUI 应用要求把**界面(View)**、**状态与逻辑(ViewModel)**、**数据(Model)**分层,即 MVVM。本次中 `MainView.axaml` 只管展示与绑定,`MainViewModel` 持有所有状态与命令逻辑,`Models` 定义纯数据结构,三者通过数据绑定协作。 +3. **状态展示方式不同**:控制台靠 `Console.WriteLine` 顺序打印;GUI 靠**数据绑定**——只要 `ObservableProperty` 的值变化,界面自动刷新(如状态栏的 `ConnectStatus`、文件列表 `LogFiles`、结果列表 `ResultEntries`),无需手动「重绘」。 +4. **用户输入的不可控性**:控制台输入基本是字符串;GUI 中用户可能在不该空的输入框留空、输入非法字符、在不该点击时点击、未连接就操作等。必须处处做输入校验与异常兜底。 + +**额外的难点 / 复杂之处:** + +1. **必须理解并正确使用 MVVM 与数据绑定**:要搞清 `OneWay` / `OneWayToSource` / `TwoWay` 等绑定模式。例如文件列表用 `SelectedItem="{Binding SelectedLogFile, Mode=OneWayToSource}"`,而多选时 ListBox 无法把「全部选中项」直接绑定到 ViewModel,必须借助 `MainView.axaml.cs` 里的 `LogFileListBox_SelectionChanged` 回调手动维护 `SelectedFiles`——这是 View 与 ViewModel 边界上一个很别扭的地方。 +2. **UI 线程模型与线程安全**:UI 控件只能由 UI 线程访问,而异步 gRPC 调用的延续可能在别的线程上。`ObservableCollection` 的变更必须回到 UI 线程,否则会抛异常。Avalonia 的绑定机制帮我们处理了大部分,但理解其原理是额外的心智负担。 +3. **异常处理策略完全不同**:控制台里一个未捕获异常最多让程序退出;GUI 里一个未捕获异常会让整个窗口崩溃,体验极差。因此必须用 `WithClientNotNull` 这类统一的兜底,把所有异常转成**消息框**而非崩溃。 +4. **调试与反馈链更长**:除了要同时启动 Agent 与 Client 两个程序(上一节的痛点依然存在),GUI 的状态分布在绑定、ViewModel、控件回调等多处,定位「为什么这一项没更新」往往要检查绑定路径、`Mode`、`x:DataType`、属性通知等多个环节。 + +**对异步 `async` / `await` 的进一步理解:** + +通过编写 GUI 客户端,我对异步的理解确实更深了一层。在控制台里,`await` 更多是「写法上的要求」;但在 GUI 里,`await` 有了**肉眼可见的意义**——如果没有 `await` 而是同步阻塞,UI 线程会被网络 I/O 占住,整个窗口会「卡死」(拖不动、按钮无响应)。`await` 让 UI 线程在等待网络响应时返回消息循环去处理用户的其他操作(比如拖动窗口、点击别的按钮),等结果回来再继续。这让我真正体会到「异步是为了不阻塞调用线程」这句话的含义。 + +**异步带来的额外困扰:** + +1. **「异步传染」**:一旦底层是异步的(gRPC 调用),上层调用链就得一路 `async`/`await` 到底,方法签名都要带 `Async` 后缀和 `Task`,这是无法回避的传播。 +2. **异常捕获的位置变了**:异步方法的异常不会在调用处直接抛出,而是藏在返回的 `Task` 里,必须 `await` 才能观察到,漏 `await` 会导致异常被「吞掉」,排查很困难。 +3. **多选回调与异步命令的时序**:`LogFileListBox_SelectionChanged` 在 UI 线程更新 `SelectedFiles`,而分析命令异步读取它,两者之间没有显式同步——靠的是 UI 单线程模型保证的串行性,理解这一点需要额外的思考。 + +总的来说,GUI 开发相对控制台,本质上是从「**线性流程**」转向「**事件驱动 + 数据绑定 + 多线程协作**」,复杂度显著上升;而异步编程既是 GUI 不卡死的必需品,也确实带来了一些新的心智负担。 + +### (Q4.2) + +根据TODO框架和guidance.md完成任务××.AI能给出达成任务要求的代码并自行测试验证。有时候AI会有过度、无效兜底的问题,在这次作业中基本没有出现 +--- + + diff --git a/docs/05-advanced/report.md b/docs/05-advanced/report.md new file mode 100644 index 0000000..4ddcf98 --- /dev/null +++ b/docs/05-advanced/report.md @@ -0,0 +1,226 @@ +# 05-advanced 实验报告 + +本章在前四章(解析、多线程、异步 gRPC、Avalonia 基础客户端)之上,把一个「能用的 CLI」打磨成一个「可用的 GUI 产品」。围绕「功能性」「美观性」「自由功能」三类要求,共实现六个功能: + +| # | 任务编号 | 功能 | 类别 | +| :-: | :------ | :--- | :--- | +| 1 | T5.1.a.a | Parquet 列式读写 | 功能性 | +| 2 | T5.1.a.b | Token 鉴权 + 多用户隔离 | 功能性 | +| 3 | T5.1.a.c | 多条件查询(Query) | 功能性 | +| 4 | T5.1.a.d | 调用拓扑推断 | 功能性 | +| 5 | T5.1.b.a | 结果表格 + Severity 高亮 | 美观性 | +| 6 | T5.2 | Request ID 链路追踪瀑布图 | 自由功能 | + +## 一、程序编译与运行 + +GUI 客户端采用「Agent(gRPC 服务端)+ 桌面客户端」双进程架构,需要分别启动。 + +**1. 启动 Agent(服务端)** + +```bash +# 终端 1:编译并启动 Agent,监听 http://localhost:5000 +dotnet run --project src/LogAnalyzerAgent +``` + +Agent 启动时会在控制台输出一个**管理员 token**(这是登录的凭据,务必复制保存),形如: + +``` +info: Bootstrap[0] + Admin token generated. Use it to log in the client (File -> Connect...) and manage other tokens: <32 位 token> +``` + +token 由 `RandomNumberGenerator` 生成 24 字节随机数再 base64url 编码(约 32 字符),高熵且 URL 安全。 + +**2. 启动桌面客户端** + +```bash +# 终端 2:编译并启动桌面客户端 +dotnet run --project src/LogAnalyzerClient/LogAnalyzerClient.Desktop +``` + +**3. 连接 Agent** + +客户端启动后界面为空,点击菜单 `File → Connect...`,在弹出的 `Connect` 对话框中: +- **Address**:填 `http://localhost:5000`(占位符已给出示例)。 +- **Token**:粘贴上一步控制台打印的管理员 token。 + +点击 `Connect`,客户端先发一个 `Ping` 校验 token;成功后底部状态栏显示 `Connected`、`[Admin]`、Agent 地址与当前目录,左侧文件列表刷新。若 token 错误会提示 `Authentication failed: the token was rejected by the Agent.` + +> 连接成功后,所有后续操作(列文件、分析、查询、导出……)都会自动在 gRPC 头里携带该 token,无需重复输入。 + +## 二、实现的功能 + +### 1.(T5.1.a.a,功能性)Parquet 列式读写 + +#### 1.1 功能概述 + +把日志分析结果持久化为 **Parquet 列式存储**文件,并支持把导出的 `.parquet` 当作普通日志文件读回、再次分析,形成「`.log → 分析 → 导出 Parquet → 再分析」的闭环。Parquet 的列式 + 游程编码对这类「三种事件类型拼成的稀疏宽表」特别友好——大量 `null` 列几乎不占空间。 + +#### 1.2 实现要点 + +- **依赖**:在 `LogParser` 项目引入 `Parquet.Net 6.0.3`(`LogParser.csproj:11`),上层通过 `LogParser` 的包装 API 间接使用,避免污染客户端依赖。 +- **Schema 设计**:用 POCO `ParquetLogRow`(`LogParser/Parquet/LogParquetSchema.cs`)描述一张 13 列的「宽而稀疏」表——公共列 `LineNo/Timestamp/PodName/Severity/EventType` 全填,类型专属列(Call 的 `TargetService/DurationMs`、Request 的 `Method/Path/StatusCode`、Internal 的 `ExceptionName/ExceptionMessage`)按行类型填、其余留 `null`。`Timestamp` 用 ISO-8601 round-trip(`"O"`)格式存,枚举统一小写存字符串。 +- **写**:`ParquetLogWriter.WriteAsync`(`LogParser/Parquet/ParquetLogWriter.cs`)把每条 `LogEntry` 经 `ToRow` 映射成 `ParquetLogRow`,再用高层 API `ParquetSerializer.SerializeAsync` 一次落盘,返回行数。 +- **读**:`ParquetLogReader.ReadAsync`(`LogParser/Parquet/ParquetLogReader.cs`)用 `ParquetSerializer.DeserializeAsync` 读回,再按 `EventType` 重建出 `CallLogEntry/RequestLogEntry/InternalLogEntry` 三种具体记录;遇到未知 `event_type`/`severity` 抛 `FormatException`,会被 `WorkerMain` 捕获、标记分析失败。 +- **新增 RPC**(`log_analyzer.proto`): + ```proto + rpc ExportAnalysisResult(ExportAnalysisResultRequest) returns (ExportAnalysisResultResponse); + // request: file_name / output_path / overwrite + // response: status / written_path / entry_count + ``` +- **Agent 端**(`AgentSession.ExportAnalysisResultAsync`):校验文件名、输出路径非空 → 文件存在且已分析成功 → 解析路径(相对路径基于 Agent 当前目录,缺 `.parquet` 后缀自动补齐)→ `overwrite=false` 且文件已存在则报错 → 创建父目录 → 调 `ParquetLogWriter.WriteAsync` → 回填 `written_path` 与 `entry_count`。 +- **文件枚举一视同仁**:`LogFileAnalyzer.EnumerateLogFiles` 同时收集 `*.log` 和 `*.parquet`;分析时 `WorkerMain` 按扩展名分发(`.parquet` 走 `ParquetLogReader`,其余走文本解析器),所以导出的 `.parquet` 会直接出现在文件列表里、可被再次分析。 + +#### 1.3 操作方法 + +1. 连接 Agent 后,在 `Directory Path` 输入框填入数据目录(如 `dataset`),点 `Change Directory`,再点 `Refresh`,左侧列表出现日志文件。 +2. 选中一个文件(左键单击高亮),点右侧 `Analyze → Selected`(或右键 `Analyze File`)完成分析。 +3. **右键**该已分析文件 → `Export to Parquet`,弹出 `Export to Parquet` 对话框: + - 上方 `Source` 显示源文件名; + - `Output Path` 填绝对路径,省略 `.parquet` 后缀会自动补齐; + - 勾选 `Overwrite if the file already exists` 可覆盖同名文件。 +4. 点 `Export`,成功后弹出 `Exported {N} entries to: {written_path}`。 +5. 点工具栏 `Refresh`,新生成的 `.parquet` 出现在文件列表中。 +6. 选中该 `.parquet` → `Analyze → Selected` → `View Analysis Results`,结果应与原 `.log` 完全一致,验证读写闭环。 +7. (Browser 端)WASM 无自定义窗口,退化为单行 `prompt`:输入路径,行首加 `!` 表示覆盖。 + +### 2.(T5.1.a.b,功能性)Token 鉴权 + 多用户隔离 + +#### 2.1 功能概述 + +为 Agent 加上 **Bearer Token 应用层鉴权**:每个 RPC 都必须携带合法 token,否则返回 `Unauthenticated`;管理类 RPC 还要求 `Admin` 角色。同时实现**多用户隔离**——每个 token 拥有独立的 `LogFileAnalyzer`,目录与分析结果互不可见。Admin 可在 GUI 里增删 token、改权限,且系统拒绝删除/降级最后一个 admin 以防锁死。 + +#### 2.2 实现要点 + +- **token 生成**:`TokenStore.GenerateToken`(`Auth/TokenStore.cs`)用 `RandomNumberGenerator.Fill` 取 24 字节密码学随机数,base64url 编码。`TokenInfo` 含不可变的 `Token` 与可变的 `Role/Note`,全部读写经同一把锁。 +- **启动引导**:`Program.cs` 启动时调 `tokenStore.CreateAdminToken()`,并把该 admin token **完整**打印一次(其余日志里都用 `Mask()` 只留前 6 位)。 +- **服务端鉴权**:`AgentService.Authorize`(`Services/AgentService.cs`)从 `authorization: Bearer ` 头取 token,查 `TokenStore.TryGet`,失败抛 `RpcException(Unauthenticated)`;`RequireAdmin` 在此基础上校验 `Admin` 角色,否则 `PermissionDenied`。17 个 RPC 全部先 `Authorize`,其中 4 个管理 RPC 额外 `RequireAdmin`。 +- **多用户隔离**:`SessionManager`(`Auth/SessionManager.cs`)是一个 `ConcurrentDictionary`,以 token 为键、懒加载每个用户独有的 `LogFileAnalyzer`。`AgentSession` 里每个业务方法开头都是 `var analyzer = _sessions.GetOrCreate(caller.Token);`,天然隔离目录与结果。 +- **客户端注入**:`TokenInterceptor`(`Services/TokenInterceptor.cs`)实现 gRPC `Interceptor`,覆盖 unary/server-stream/client-stream/duplex 五种入口,在 `WithToken` 里给每次调用注入 `authorization` 头(并幂等地剔除旧头防重复)。Desktop 与 Browser 的 `ClientFactory` 都用 `channel.Intercept(new TokenInterceptor(token))` 包裹通道。 +- **防锁死**:`TokenStore.TryDelete`/`TrySetRole` 在只剩一个 admin 时拒绝删除/降级。 +- **相关 proto**: + ```proto + rpc CreateToken(CreateTokenRequest) returns (CreateTokenResponse); // 需 Admin + rpc DeleteToken(DeleteTokenRequest) returns (OperationStatusMessage); // 需 Admin + rpc ListTokens(Empty) returns (ListTokensResponse); // 需 Admin + rpc SetTokenRole(SetTokenRoleRequest) returns (OperationStatusMessage); // 需 Admin + ``` + `ListTokensResponse.caller_token` 让客户端高亮「自己」那一行。 + +#### 2.3 操作方法 + +1. 用 admin token 连接(见「一」)。连接成功后 `DetectAdminAsync` 自动探测权限,菜单出现 `File → Manage Tokens...`(非 admin 该项隐藏)。 +2. 点 `File → Manage Tokens...` 打开 `Token Management` 窗口,列出全部 token(自己的行带浅蓝底、`YOU` 徽标)。 +3. **创建 token**:`Role` 选 `Normal`/`Admin`,`Note` 填备注(如 `for-alice`),点 `Create`。新 token **仅显示一次**在状态栏,需立即复制发给用户。 +4. **改权限**:点对应该行的 `Promote`/`Demote` 按钮切换 Normal ↔ Admin。 +5. **删除**:点该行 `Delete`(最后一个 admin 删不掉,会报错)。 +6. **验证隔离**:另开一个客户端实例,用不同 token 连接,`Change Directory` 到不同目录、分析不同文件;两个窗口的文件列表与结果互不影响。 +7. **验证鉴权**:用 `File → Connect...` 填一个随便编的 token,连接会被拒,提示认证失败。 + +### 3.(T5.1.a.c,功能性)多条件查询(Query) + +#### 3.1 功能概述 + +对已分析的结果按 **事件类型 / 严重等级 / 服务 / Request ID / 时间范围** 五个维度做组合过滤,结果以表格呈现并支持重排。全部条件「留空即不过滤该维度」,组合关系为 AND。 + +#### 3.2 实现要点 + +- **服务端流式 RPC**(复用 `GetAnalysisResultResponse` 载体): + ```proto + rpc QueryAnalysisResult(QueryAnalysisResultRequest) returns (stream GetAnalysisResultResponse); + // request: file_name / repeated event_types / repeated severities / + // request_id_pattern / service_pattern / optional start_time / optional end_time + ``` + pattern 用子串、大小写不敏感;时间两端均**闭区间**、按 UTC 解释。Internal 日志没有 request-id,设了 request-id 过滤时自动排除。 +- **客户端数据模型**:`QueryFilter`(`Models/QueryFilter.cs`)用 `HashSet<枚举>` 表达类型/等级集合、两个 `string` pattern、两个 `DateTimeOffset?` 时间界。`ToRequest` 序列化为 gRPC 请求。 +- **桌面对话框** `QueryDialog`:用 CheckBox 组表达「类型」「等级」多选(全不勾=任意),TextBox 表达服务/Request ID/起止时间。校验时间可解析且 `start <= end`,否则行内红字报错、窗口不关。 +- **Browser 退化**:`QueryFilterParser`(`Helpers/QueryFilterParser.cs`)把一行文本解析成 `QueryFilter`,语法如 `type=Call,Request severity=Warning,Error service=gateway from=2026-06-05 to=2026-06-05T17:00:00Z`;未知键静默忽略、保证鲁棒。 +- **结果与排序解耦**:流式结果先进 `_loadedEntries`(原始顺序),再套当前排序键生成 `ResultEntries`。`SortKeys` 支持 12 个键(LineNo/Timestamp/Severity/EventType/PodName/RequestId/TargetService/Method/Path/StatusCode/DurationMs/ExceptionName),缺失字段按类型默认值兜底(如非 Call 行的 `Method` 为 `""`)。查询后 `IsResultFiltered=true`,`Show All` 按钮出现,点它调 `GetAnalysisResult` 还原全集。 + +#### 3.3 操作方法 + +1. 选中一个已分析成功的文件(左键高亮,使其成为当前结果来源)。 +2. 点结果面板右上角 `Query...` 按钮,打开 `Query Log Entries` 对话框。 +3. 按需勾选/填写(全部留空 = 等价于 Show All): + - **Event Type**:勾 `Call`/`Request`/`Internal` 中的一个或多个; + - **Severity**:勾 `Info`/`Warning`/`Error`; + - **Service**:填服务名(如 `gateway`,会匹配 `gateway-0`、`gateway-1`); + - **Request ID contains**:填 request-id 子串; + - **Start time / End time**:填日期(`2026-06-05`)或完整 ISO 8601(`2026-06-05T17:00:00Z`),闭区间、按 UTC。 +4. 点 `Query`。表格刷新为命中行,状态栏显示 `Filtered: showing N entr(y/ies).`。 +5. 用 `Sort by` 下拉 + `Descending` 复选框对查询结果重排(实时生效)。 +6. 点 `Show All` 退出过滤、恢复全集。 + +### 4.(T5.1.a.d,功能性)调用拓扑推断 + +#### 4.1 功能概述 + +从已分析文件的 **Call 类型日志**里推断出服务间的有向调用图(节点=服务,边=调用关系,边权=调用次数),在独立窗口里画成节点-箭头图;点某条边即把该调用关系对应的全部 Call 日志加载进结果面板。这对应可观测性「拓扑」支柱。 + +#### 4.2 实现要点 + +- **两个 RPC**(`log_analyzer.proto`): + ```proto + rpc GetCallTopology(GetCallTopologyRequest) returns (GetCallTopologyResponse); // 一元,返回整图 + rpc GetEdgeCallLogs(GetEdgeCallLogsRequest) returns (stream GetAnalysisResultResponse); // 流式,返回该边 Call 日志 + ``` +- **图推断**(`AgentSession.GetCallTopology`):遍历 `result.Entries`,只取 `CallLogEntry`;用 `ExtractService`(正则 `-\d+$` 去掉 pod 副本后缀,如 `gateway-0 → gateway`)把 `PodName` 归一为源服务、`TargetService` 为目标服务;`SortedSet` 收集节点(确定性顺序保证客户端布局可复现),`Dictionary<(src,tgt),int>` 累计边权。`Top` 头部信息含服务数、边数。 +- **点边取日志**:`GetEdgeCallLogs` 先推一个 header(让客户端读文件状态),再把源、目标都匹配的 Call 日志逐条 stream 回去。 +- **Avalonia 绘图**(`TopologyWindow`):固定 720×460 `Canvas` 套 `Viewbox` 等比缩放。节点按圆周等分布置(胶囊形 `Border`,蓝色填充);普通边为灰色 `Line` + 三角箭头,并把两端各回缩 36px 让箭头落在节点边框上;自环画成节点上方的小椭圆。关键技巧:在细线之上叠一条 `StrokeThickness=16`、alpha=1 的近透明 `Line` 作**点击热区**,让细边也能轻松点中。 +- **交互回流**:点边 → `Window.Close(edge)` → `ShowTopologyAsync` 拿到该边 → 发 `GetEdgeCallLogs` → `ConsumeResultStreamAsync` 把 Call 日志灌进结果表格,状态栏显示 `Edge gateway -> userservice (12 calls).`。 + +#### 4.3 操作方法 + +1. 选中并分析一个含 Call 日志的文件(如 `basic.log`)。 +2. **右键**该文件 → `Show Call Topology`,弹出 `Call Topology` 窗口,顶部显示 `Call topology of '' - N service(s), M edge(s)`。 +3. 鼠标悬停节点看服务名 Tooltip;点任意**箭头**(或自环节点)选中该调用关系,窗口自动关闭。 +4. 结果表格被替换为该边的全部 Call 日志,状态栏显示边信息;可继续 `Query...` 或 `Sort by` 二次筛选。 +5. 若文件没有 Call 日志,会提示无法推断拓扑;点 `Close` 或不选边直接关窗则什么都不做。 + +### 5.(T5.1.b.a,美观性)结果表格 + Severity 高亮 + +#### 5.1 功能概述 + +把分析结果从「一行字符串的 ListBox」升级为 **DataGrid 表格**:每条日志一行,按 `#`、`Time`、`Severity`、`Service`、`Type`、`Detail` 分列;不同事件类型的字段摘要统一收进 Detail 列。Severity 列用圆角胶囊按等级着色——Info 蓝、Warning 橙、Error 红,一眼定位错误。 + +#### 5.2 实现要点 + +- **强类型行模型** `LogEntryRowVm`(`Models/LogEntryRowVm.cs`):由 `LogEntry` 直接构造。`Time` 格式化为 `HH:mm:ss.fff`;`Service` 用同样的 `-\d+$` 正则去掉 pod 后缀;`Detail` 由 `BuildDetail` 按类型生成摘要: + - Call → `-> {target} ({dur} ms)` + - Request → `{method} {path} -> {code}` + - Internal → `{ExceptionName}: {ExceptionMessage}` +- **着色**:`SeverityToBrushConverter`(`Converters/SeverityToBrushConverter.cs`,`IValueConverter`)把 `LogSeverity` 枚举映射成预分配的 `SolidColorBrush`(`#2B6CB0`/`#DD6B20`/`#C53030`,static 只分配一次)。XAML 里 Severity 列用 `DataGridTemplateColumn`,胶囊 `Border` 的 `Background="{Binding Severity, Converter={StaticResource SeverityToBrush}}"`,比给 `Classes` 绑定动态值更可靠。 +- **状态反馈**:未分析 / 分析失败 / 命中数为 0 等状态,用表格上方加粗的状态栏文字显示(`Showing all N entries.` / `Filtered: showing N entries.` / `No entries match the query.` 等)。 +- **排序与展示解耦**:`_loadedEntries: List` 保留原始顺序,展示用的 `ResultEntries` 由 `ApplySort` 重排生成;切换 `Sort by` / `Descending` 实时重排。 + +#### 5.3 操作方法 + +1. 分析某文件后右键 `View Analysis Results`(或点过滤后的 `Show All`),结果以六列表格呈现,Severity 列自动着色。 +2. 建议用 `basic.log` 验证 Info/Warning/Error 三色胶囊;用混合日志验证三种 Detail 摘要格式。 +3. 用 `Sort by` 下拉选排序键、勾 `Descending` 调整顺序;列宽自适应,Detail 列占剩余宽度。 + +### 6.(T5.2,自由功能)Request ID 链路追踪瀑布图 + +#### 6.1 功能概述 + +云服务可观测性三大支柱是「指标 / 拓扑 / 追踪」。第 4 节展示了「拓扑」,本功能展示**追踪**——即分布式追踪(Distributed Tracing)的瀑布图:刻画**某一次**请求依次经过了哪些服务、每跳多久、在哪一段出错。在结果表格里选中一条 Call/Request 日志,客户端按其 `request-id` 向 Agent 拉取整条调用链,在弹窗里画成竖向瀑布:每个横条是一次服务调用,从左到右按时间排列,条长代表耗时,出错段标红。 + +#### 6.2 实现要点 + +- **新增 RPC**(`log_analyzer.proto`),**响应复用 `stream GetAnalysisResultResponse`**(与 `GetAnalysisResult`/`GetEdgeCallLogs` 同载体),省去新建 message,客户端复用 `GrpcTypeConverter` 与 header 状态判断: + ```proto + rpc GetTrace(GetTraceRequest) returns (stream GetAnalysisResultResponse); + // request: file_name / request_id + ``` +- **Agent 端**(`AgentSession.GetTrace`):按 `request-id` 精确匹配(Call/Request 取其 `RequestId`,Internal 无 id 自动排除),**按时间升序**返回,便于客户端直接落笔画瀑布。 +- **客户端** `ConsumeTraceStreamAsync`:只取其中的 Call 日志组装成 `TraceSpan`(源服务、目标服务、起始时刻、耗时、是否出错——`Severity==Error` 即红色),不污染主结果表格。 +- **瀑布图**(`TraceWindow`,Canvas + Viewbox,与拓扑窗口同套路):先求整条链路的 `[minStart, maxEnd]` 时间跨度并归一化到画布宽度,每个 span 的左偏移由其起始时刻算出、宽度由耗时按比例算出(设 `MinBarWidth` 防极短条看不见),失败段用红色填充/边框、正常段蓝色;画布高度随 span 数动态增高。两端标注 `+0 ms` 与 `+{total} ms total`。 +- **数据坑(关键)**:原 `gen.py` 给每条日志独立生成 `uuid`,导致每个 `request-id` 只出现一次,画不出多跳链路。改造 `gen.py` 新增 `make_trace`:模拟一次请求沿调用拓扑随机游走多跳(最多 3 跳)、共享同一 `request-id`、时间递增,约 1/4 的链路在末跳失败(用于演示红色错误段)。生成的 `dataset/trace.log` 经统计确认存在出现 ≥2 次的 `request-id`(最多 5 次,即 2 跳链路)。 + +#### 6.3 操作方法 + +1. `Change Directory` 到 `dataset`,分析 `trace.log`(已含多跳链路)。 +2. `View Analysis Results` 查看结果表格。 +3. 在表格里**左键**选中一条 Call 行(使其成为 `SelectedResultEntry`),再**右键** → `View Trace for this Request`。 +4. 弹出 `Request Trace Waterfall` 窗口,顶部显示 `Trace of request - N span(s) - ''`,下方为瀑布图:每条蓝色横条一跳、末跳失败时为红色,鼠标可读 `源 -> 目标 (耗时 ms)` 的标签。 +5. 选中 Internal 类型日志再右键追踪,会提示「无 Request ID」;选中的行没有先左键点过,会提示先左键选中再右键。 diff --git a/src/DataGen/gen.py b/src/DataGen/gen.py index 50def39..78a5374 100644 --- a/src/DataGen/gen.py +++ b/src/DataGen/gen.py @@ -197,6 +197,74 @@ def make_internal_event() -> Tuple[str, Dict[str, object]]: return choose_pod(service), message +# 一次完整请求贯穿的最多跳数与出错概率(用于 T5.2 链路追踪数据)。 +TRACE_MAX_HOPS = 3 +TRACE_ERROR_PROBABILITY = 0.25 + + +def make_trace(start_time: datetime) -> Tuple[List[CsvLogLine], datetime]: + """生成一条完整调用链(T5.2 链路追踪用)。 + + 一次外部请求从入口服务进入,沿 CALL_TOPOLOGY 随机游走若干跳,所有事件共享同一个 + request-id,时间戳严格递增。约 1/4 的链路会在最后一跳失败(call 为 ERROR + 一个 + internal 错误),用于演示瀑布图里的「出错段标红」。 + + 返回 (该链路产生的全部日志行, 链路结束后的时间)。 + """ + request_id = str(uuid.uuid4()) + + # 沿调用拓扑随机游走一条有向路径,至少入口 + 1 跳。 + path: List[str] = [choose_service_with_outgoing_edges()] + for _ in range(TRACE_MAX_HOPS): + targets = CALL_TOPOLOGY.get(path[-1], []) + if not targets: + break + path.append(random.choice(targets)) + + is_error_trace = random.random() < TRACE_ERROR_PROBABILITY + last_hop_index = len(path) - 2 # 最后一跳:src=path[i], dst=path[i+1] + + lines: List[CsvLogLine] = [] + t = start_time + + # 入口服务收到外部请求。 + entry_pod, entry_msg = make_request_event(path[0], request_id) + lines.append(CsvLogLine(0, t, entry_pod, entry_msg)) + + for i in range(len(path) - 1): + src, dst = path[i], path[i + 1] + duration = random.randint(5, 250) + t = t + timedelta(milliseconds=random.randint(MIN_TIME_STEP_MS, 200)) + + fail_here = is_error_trace and i == last_hop_index + call_msg = { + "severity": "ERROR" if fail_here else "INFO", + "event": "call", + "request-id": request_id, + "target-service": dst, + "duration-ms": duration, + } + lines.append(CsvLogLine(0, t, choose_pod(src), call_msg)) + + if fail_here: + # 出错跳:补一条 internal 错误,便于在瀑布图中看到红色错误段。 + t = t + timedelta(milliseconds=random.randint(MIN_TIME_STEP_MS, 100)) + exc_name, exc_message = random.choice(ERROR_EXCEPTIONS) + internal_msg = { + "severity": "ERROR", + "event": "internal", + "exception": f"{exc_name}: {exc_message}", + } + lines.append(CsvLogLine(0, t, choose_pod(dst), internal_msg)) + else: + # 正常跳:下游收到请求。 + t = t + timedelta(milliseconds=random.randint(MIN_TIME_STEP_MS, 100)) + req_pod, req_msg = make_request_event(dst, request_id) + lines.append(CsvLogLine(0, t, req_pod, req_msg)) + + return lines, t + + def next_timestamp(current: datetime) -> datetime: step_ms = random.randint(MIN_TIME_STEP_MS, MAX_TIME_STEP_MS) return current + timedelta(milliseconds=step_ms) @@ -211,10 +279,18 @@ def generate_log_lines(count: int) -> List[CsvLogLine]: lines: List[CsvLogLine] = [] current_time = START_TIME + produced = 0 - for lineno in range(count): - current_time = next_timestamp(current_time) + while produced < count: + remaining = count - produced + # 约 35% 的概率生成一条完整调用链(至少 4 行),其余生成单条独立日志。 + if remaining >= 4 and random.random() < 0.35: + trace_lines, current_time = make_trace(current_time) + lines.extend(trace_lines) + produced += len(trace_lines) + continue + current_time = next_timestamp(current_time) if random.random() < INTERNAL_EVENT_PROBABILITY: pod_name, message = make_internal_event() else: @@ -224,8 +300,14 @@ def generate_log_lines(count: int) -> List[CsvLogLine]: pod_name, message = make_call_event(source_service) else: pod_name, message = make_request_event() - - lines.append(CsvLogLine(lineno, current_time, pod_name, message)) + lines.append(CsvLogLine(0, current_time, pod_name, message)) + produced += 1 + + # trace 与独立行交错生成,按时间戳排序保证 lineno 与 timestamp 全局单调, + # 再统一编号(CsvLogLine 为不可变 dataclass,原地重建)。 + lines.sort(key=lambda line: line.timestamp) + for index, line in enumerate(lines): + lines[index] = CsvLogLine(index, line.timestamp, line.pod_name, line.message) return lines diff --git a/src/LocalCli/Program.cs b/src/LocalCli/Program.cs index 17b30db..a0afc33 100644 --- a/src/LocalCli/Program.cs +++ b/src/LocalCli/Program.cs @@ -1,4 +1,5 @@ using LogAnalyzer; +using LogParser.Models; using LogParser.Visitors; namespace LocalCli @@ -112,22 +113,106 @@ 6. Exit. private static void ShowLogFiles(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + var files = analyzer.GetLogFiles(); + if (files.Count == 0) + { + Console.WriteLine("No log files found in the current directory."); + return; + } + + Console.WriteLine($"Log files in '{analyzer.CurrentDirectory}' ({files.Count}):"); + foreach (var file in files) + { + Console.WriteLine($" {file}"); + } } private static void AnalyzeFiles(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + Console.WriteLine("Please input file names separated by ',' to analyze:"); + var line = Console.ReadLine(); + if (line is null) + { + return; + } + + var fileNames = line.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries); + if (fileNames.Length == 0) + { + Console.WriteLine("No file names provided."); + return; + } + + try + { + Console.WriteLine($"Analyzing {fileNames.Length} file(s) with parallelism 0 (= ProcessorCount = {Environment.ProcessorCount})..."); + analyzer.AnalyzeFiles(0, fileNames); + Console.WriteLine("Analysis completed."); + } + catch (ArgumentException ex) + { + Console.WriteLine($"Error: {ex.Message}"); + } + catch (InvalidOperationException ex) + { + Console.WriteLine($"Error: {ex.Message}"); + } } private static void AnalyzeAll(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + try + { + Console.WriteLine($"Analyzing all log files with parallelism 0 (= ProcessorCount = {Environment.ProcessorCount})..."); + analyzer.AnalyzeAll(0); + Console.WriteLine("Analysis completed."); + } + catch (InvalidOperationException ex) + { + Console.WriteLine($"Error: {ex.Message}"); + } } private static void GetAnalysisResult(LogFileAnalyzer analyzer) { - throw new NotImplementedException("T2.3"); + Console.WriteLine("Please input the file name to get analysis result:"); + var fileName = Console.ReadLine(); + if (fileName is null) + { + return; + } + + fileName = fileName.Trim(); + if (string.IsNullOrEmpty(fileName)) + { + Console.WriteLine("No file name provided."); + return; + } + + if (!analyzer.TryGetAnalysisResult(fileName, out var result)) + { + Console.WriteLine($"File '{fileName}' does not exist in the current directory."); + return; + } + + switch (result!.State) + { + case AnalysisState.NotAnalyzed: + Console.WriteLine($"File '{fileName}' has not been analyzed yet. Please analyze it first."); + break; + case AnalysisState.Succeeded: + Console.WriteLine($"Analysis result for '{fileName}' (parsed by worker {result.WorkerId}, {result.Entries.Count} entries):"); + var visitor = new KeyValueVisitor(); + foreach (var entry in result.Entries) + { + var dump = visitor.Dump(entry); + Console.WriteLine($" [{entry.EventType}] {string.Join(", ", dump.Select(kv => $"{kv.Key}={kv.Value}"))}"); + } + break; + case AnalysisState.Failed: + Console.WriteLine($"Failed to analyze '{fileName}': {result.ErrorMessage}"); + break; + } } } } diff --git a/src/LogAnalyzer/LogFileAnalyzer.cs b/src/LogAnalyzer/LogFileAnalyzer.cs index c3e7691..56a7181 100644 --- a/src/LogAnalyzer/LogFileAnalyzer.cs +++ b/src/LogAnalyzer/LogFileAnalyzer.cs @@ -1,7 +1,7 @@ using LogParser.Models; using LogParser.Parser; +using LogParser.Parquet; using System.Diagnostics.CodeAnalysis; -using System.Security.Cryptography.X509Certificates; namespace LogAnalyzer { @@ -62,9 +62,8 @@ public bool ChangeDirectory(string? directoryPath) _analysisResults.Clear(); if (directoryPath is not null) { - var logFiles = Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.TopDirectoryOnly) - .Select(filePath => Path.GetFileName(filePath)) - .OrderBy(fileName => fileName); + // 同时收集 .log 与 .parquet 两类日志文件,统一按文件名排序。 + var logFiles = EnumerateLogFiles(directoryPath); foreach (var fileName in logFiles) { _logFiles.Add(fileName, new FileInfo(Path.Join(_currentDirectory, fileName))); @@ -98,6 +97,18 @@ public bool TryGetAnalysisResult(string fileName, out AnalysisResult? result) } } + /// + /// 枚举目录中受支持的日志文件(.log 与 .parquet),按文件名排序返回。 + /// + private static IEnumerable EnumerateLogFiles(string directoryPath) + { + var logFiles = Directory.EnumerateFiles(directoryPath, "*.log", SearchOption.TopDirectoryOnly) + .Select(Path.GetFileName)!; + var parquetFiles = Directory.EnumerateFiles(directoryPath, "*.parquet", SearchOption.TopDirectoryOnly) + .Select(Path.GetFileName)!; + return logFiles.Concat(parquetFiles).OrderBy(fileName => fileName)!; + } + public void AnalyzeAll(int degreeOfParallelism) { List fileNames; @@ -141,7 +152,7 @@ public void AnalyzeFiles(int degreeOfParallelism, IEnumerable fileNames) /* * Set _isAnalyzing */ - // TODO: T2.2 + _isAnalyzing = true; } try @@ -154,7 +165,10 @@ public void AnalyzeFiles(int degreeOfParallelism, IEnumerable fileNames) * Unset _isAnalyzing * Remember to lock _syncRoot to prevent data race */ - // TODO: T2.2 + lock (_syncRoot) + { + _isAnalyzing = false; + } } } @@ -169,7 +183,15 @@ private void RunWorkers(int degreeOfParallelism, IReadOnlyList fileLis * Filter unparsed files. * If there is an unknown file, throw System.InvalidOperationException. */ - throw new NotImplementedException("TODO: T2.2"); + if (!_analysisResults.TryGetValue(file.Name, out var existing)) + { + throw new InvalidOperationException($"Unknown file '{file.Name}'."); + } + // 跳过已经分析过且保存了分析结果的文件(Succeeded / Failed),只保留未分析的文件 + if (existing.State == AnalysisState.NotAnalyzed) + { + logFilesToParse.Add(file); + } } } @@ -183,7 +205,11 @@ private void RunWorkers(int degreeOfParallelism, IReadOnlyList fileLis /* * Enqueue log files */ - // TODO: T2.2 + foreach (var file in logFilesToParse) + { + queue.Enqueue(file); + } + queue.CompleteAdding(); degreeOfParallelism = Math.Max(Math.Min(degreeOfParallelism, logFilesToParse.Count), 1); var workers = new Thread[degreeOfParallelism]; @@ -194,13 +220,22 @@ private void RunWorkers(int degreeOfParallelism, IReadOnlyList fileLis /* * Create and start threads to run `WorkerMain` */ - // TODO: T2.2 + var thread = new Thread(() => WorkerMain(workerId, queue)) + { + IsBackground = true, + Name = threadName, + }; + workers[i] = thread; + thread.Start(); } /* * Wait for (join) all threads to end */ - // TODO: T2.2 + foreach (var thread in workers) + { + thread.Join(); + } } private void WorkerMain(int workerId, WorkQueue queue) @@ -212,21 +247,49 @@ private void WorkerMain(int workerId, WorkQueue queue) AnalysisResult result; try { - // Parse file - throw new NotImplementedException("TODO: T2.2"); + // 按扩展名分派解析器:.log 走文本 CSV 解析;.parquet 走 Parquet 读取。 + IReadOnlyList entries = Path.GetExtension(file.Name).Equals(".parquet", StringComparison.OrdinalIgnoreCase) + ? ParquetLogReader.ReadAsync(file.FullName).GetAwaiter().GetResult() + : ParseLogFile(parser, file.FullName); + result = new AnalysisResult( + FileName: file.Name, + FullName: file.FullName, + State: AnalysisState.Succeeded, + Entries: entries, + ErrorMessage: null, + WorkerId: workerId); } catch (Exception ex) { // Save exception message to result - throw new NotImplementedException("TODO: T2.2"); + result = new AnalysisResult( + FileName: file.Name, + FullName: file.FullName, + State: AnalysisState.Failed, + Entries: Array.Empty(), + ErrorMessage: ex.Message, + WorkerId: workerId); } /* * Save parse result. * [!Important] Remember to lock _syncRoot to prevent data race. */ - throw new NotImplementedException("TODO: T2.2"); + lock (_syncRoot) + { + _analysisResults[file.Name] = result; + } } } + + /// + /// 用文本解析器解析 .log 文件,立即求值以确保解析异常在本处被捕获。 + /// + private static IReadOnlyList ParseLogFile(LogFileParser parser, string fullName) + { + using var reader = new StreamReader(fullName); + // 使用 ToList() 强制立即求值,确保解析过程中的异常在本 try 块内被捕获 + return parser.Parse(reader).ToList(); + } } } diff --git a/src/LogAnalyzer/WorkQueue.cs b/src/LogAnalyzer/WorkQueue.cs index 23055a5..e6b42a3 100644 --- a/src/LogAnalyzer/WorkQueue.cs +++ b/src/LogAnalyzer/WorkQueue.cs @@ -20,17 +20,50 @@ public bool IsCompleted public void Enqueue(T item) { - throw new NotImplementedException("TODO: T2.1"); + lock (_items) + { + if (_isCompleted) + { + throw new InvalidOperationException( + "Cannot enqueue after CompleteAdding has been called."); + } + _items.Enqueue(item); + // 唤醒一个正在等待的消费者(signal 操作) + Monitor.Pulse(_items); + } } public bool TryDequeue([NotNullWhen(true)] out T? item) { - throw new NotImplementedException("TODO: T2.1"); + lock (_items) + { + // 用 while 而非 if:避免虚假唤醒(spurious wakeup)带来的错误判断 + while (_items.Count == 0 && !_isCompleted) + { + // 队列为空且尚未结束放入:解锁互斥量并等待(wait 操作) + Monitor.Wait(_items); + } + + if (_items.Count > 0) + { + item = _items.Dequeue(); + return true; + } + + // 队列为空且已结束放入:返回 false + item = default; + return false; + } } public void CompleteAdding() { - throw new NotImplementedException("TODO: T2.1"); + lock (_items) + { + _isCompleted = true; + // 唤醒全部正在等待的消费者(broadcast 操作),使其能正常退出 + Monitor.PulseAll(_items); + } } } } diff --git a/src/LogAnalyzerAgent/Applications/AgentSession.cs b/src/LogAnalyzerAgent/Applications/AgentSession.cs index 2531f22..8e9f630 100644 --- a/src/LogAnalyzerAgent/Applications/AgentSession.cs +++ b/src/LogAnalyzerAgent/Applications/AgentSession.cs @@ -1,30 +1,42 @@ -using Google.Protobuf.WellKnownTypes; +using Google.Protobuf.WellKnownTypes; using Grpc.Core; using LogAnalyzer; +using LogAnalyzerAgent.Auth; using LogAnalyzerRpc.Protos; using LogAnalyzerRpc; +using LogParser.Models; +using LogParser.Parquet; using LogParser.Visitors; +using System.Text.RegularExpressions; namespace LogAnalyzerAgent.Applications { + /// + /// Agent 端的业务逻辑入口(T5.1.a.b 起,所有方法均以 标识调用者)。 + /// + /// 每个 token 代表一个独立用户,由 映射到其专属的 + /// ,从而实现不同用户间目录与分析结果的完全隔离。 + /// public class AgentSession { - private readonly LogFileAnalyzer _analyzer; + private readonly SessionManager _sessions; + private readonly TokenStore _tokens; private readonly ILogger _logger; - public AgentSession(LogFileAnalyzer analyzer, ILoggerFactory loggerFactory) + public AgentSession(SessionManager sessions, TokenStore tokens, ILoggerFactory loggerFactory) { - _analyzer = analyzer; + _sessions = sessions; + _tokens = tokens; _logger = loggerFactory.CreateLogger(); } - private static OperationStatusMessage CreateInternalErrorOperationStatus(Exception ex) + private static OperationStatusMessage CreateInternalErrorOperationStatus(Exception ex, string operation) { return new OperationStatusMessage() { Success = false, Code = AgentErrorCode.InternalError, - Message = $"An error occurred while retrieving agent status: {ex.Message}", + Message = $"An error occurred while {operation}: {ex.Message}", }; } @@ -38,63 +50,775 @@ private static OperationStatusMessage CreateNoErrorOperationStatus() }; } - public Task Ping(Empty empty, CancellationToken cancellationToken) + private static OperationStatusMessage CreateErrorOperationStatus(AgentErrorCode code, string message) { + return new OperationStatusMessage() + { + Success = false, + Code = code, + Message = message, + }; + } + + public Task Ping(Empty empty, TokenInfo caller, CancellationToken cancellationToken) + { + // Ping 仅用于校验连接与 token 是否合法;真正的校验已在 AgentService.Authorize 完成。 + _ = caller; return Task.FromResult(new Empty()); } - public Task GetAgentStatus(Empty empty, CancellationToken cancellationToken) + public Task GetAgentStatus(Empty empty, TokenInfo caller, CancellationToken cancellationToken) { + var analyzer = _sessions.GetOrCreate(caller.Token); var response = new GetAgentStatusResponse(); try { - response.HasDirectory = _analyzer.HasDirectory; - response.CurrentDirectory = _analyzer.CurrentDirectory ?? ""; - response.IsAnalyzing = _analyzer.IsAnalyzing; + response.HasDirectory = analyzer.HasDirectory; + response.CurrentDirectory = analyzer.CurrentDirectory ?? ""; + response.IsAnalyzing = analyzer.IsAnalyzing; response.Status = CreateNoErrorOperationStatus(); } catch (Exception ex) { - response.Status = CreateInternalErrorOperationStatus(ex); + response.Status = CreateInternalErrorOperationStatus(ex, "retrieving agent status"); _logger.LogError(ex, "An error occurred while retrieving agent status."); } return Task.FromResult(response); } - public Task GetLogFiles(Empty empty, CancellationToken cancellationToken) + public Task GetLogFiles(Empty empty, TokenInfo caller, CancellationToken cancellationToken) { + var analyzer = _sessions.GetOrCreate(caller.Token); var response = new GetLogFilesResponse(); try { - response.FileNames.AddRange(_analyzer.GetLogFiles()); + response.FileNames.AddRange(analyzer.GetLogFiles()); response.Status = CreateNoErrorOperationStatus(); } catch (Exception ex) { - response.Status = CreateInternalErrorOperationStatus(ex); + response.Status = CreateInternalErrorOperationStatus(ex, "retrieving log files"); _logger.LogError(ex, "An error occurred while retrieving log files."); } return Task.FromResult(response); } - public Task ChangeDirectory(ChangeDirectoryRequest request, CancellationToken cancellationToken) + public Task ChangeDirectory(ChangeDirectoryRequest request, TokenInfo caller, CancellationToken cancellationToken) + { + var analyzer = _sessions.GetOrCreate(caller.Token); + var response = new ChangeDirectoryResponse(); + try + { + if (analyzer.IsAnalyzing) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidOperation, + "Cannot change directory while analysis is in progress."); + return Task.FromResult(response); + } + + bool changed; + try + { + changed = analyzer.ChangeDirectory(request.DirectoryPath); + } + catch (ArgumentException) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + $"Invalid directory path: '{request.DirectoryPath}'."); + return Task.FromResult(response); + } + + if (!changed) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.DirectoryNotFound, + $"Directory '{request.DirectoryPath}' does not exist."); + return Task.FromResult(response); + } + + response.CurrentDirectory = analyzer.CurrentDirectory ?? ""; + response.FileNames.AddRange(analyzer.GetLogFiles()); + response.Status = CreateNoErrorOperationStatus(); + } + catch (Exception ex) + { + response.Status = CreateInternalErrorOperationStatus(ex, "changing directory"); + _logger.LogError(ex, "An error occurred while changing directory."); + } + return Task.FromResult(response); + } + + public Task AnalyzeAll(AnalyzeAllRequest request, TokenInfo caller, CancellationToken cancellationToken) + { + var analyzer = _sessions.GetOrCreate(caller.Token); + var response = new AnalyzeAllResponse(); + try + { + if (!analyzer.HasDirectory) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidOperation, + "No log directory has been set. Please change directory first."); + return Task.FromResult(response); + } + + if (request.DegreeOfParallelism < 0) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "Degree of parallelism must be non-negative."); + return Task.FromResult(response); + } + + analyzer.AnalyzeAll(request.DegreeOfParallelism); + response.Status = CreateNoErrorOperationStatus(); + } + catch (InvalidOperationException ex) + { + response.Status = CreateErrorOperationStatus(AgentErrorCode.InvalidOperation, ex.Message); + } + catch (Exception ex) + { + response.Status = CreateInternalErrorOperationStatus(ex, "analyzing all log files"); + _logger.LogError(ex, "An error occurred while analyzing all log files."); + } + return Task.FromResult(response); + } + + public Task AnalyzeFiles(AnalyzeFilesRequest request, TokenInfo caller, CancellationToken cancellationToken) + { + var analyzer = _sessions.GetOrCreate(caller.Token); + var response = new AnalyzeFilesResponse(); + try + { + if (!analyzer.HasDirectory) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidOperation, + "No log directory has been set. Please change directory first."); + return Task.FromResult(response); + } + + if (request.DegreeOfParallelism < 0) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "Degree of parallelism must be non-negative."); + return Task.FromResult(response); + } + + if (request.FileNames.Count == 0) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "No file names provided."); + return Task.FromResult(response); + } + + analyzer.AnalyzeFiles(request.DegreeOfParallelism, request.FileNames); + response.Status = CreateNoErrorOperationStatus(); + } + catch (InvalidOperationException ex) + { + response.Status = CreateErrorOperationStatus(AgentErrorCode.InvalidOperation, ex.Message); + } + catch (ArgumentException ex) + { + // LogFileAnalyzer throws ArgumentException when a requested file is not in the directory. + response.Status = CreateErrorOperationStatus(AgentErrorCode.FileNotFound, ex.Message); + } + catch (Exception ex) + { + response.Status = CreateInternalErrorOperationStatus(ex, "analyzing log files"); + _logger.LogError(ex, "An error occurred while analyzing log files."); + } + return Task.FromResult(response); + } + + public IReadOnlyList GetAnalysisResult(GetAnalysisResultRequest request, TokenInfo caller, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var analyzer = _sessions.GetOrCreate(caller.Token); + var responses = new List(); + try + { + if (!analyzer.TryGetAnalysisResult(request.FileName, out var result) || result is null) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.FileNotFound, + $"File '{request.FileName}' does not exist in the current directory."), + }); + return responses; + } + + // First, return the header describing the analysis state of the file. + responses.Add(BuildHeaderResponse(result)); + + // Only stream log entries when the analysis succeeded. + if (result.State == AnalysisState.Succeeded) + { + foreach (var entry in result.Entries) + { + responses.Add(BuildEntryResponse(entry)); + } + } + } + catch (Exception ex) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateInternalErrorOperationStatus(ex, "retrieving analysis result"), + }); + _logger.LogError(ex, "An error occurred while retrieving analysis result."); + } + return responses; } - public Task AnalyzeAll(AnalyzeAllRequest request, CancellationToken cancellationToken) + /// + /// 按条件查询某个日志文件的分析结果。逻辑与 一致, + /// 但在流式返回逐条日志前,先用 中给出的过滤条件筛选。 + /// 任一维度未填写即表示不按该维度过滤。 + /// + public IReadOnlyList QueryAnalysisResult(QueryAnalysisResultRequest request, TokenInfo caller, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var analyzer = _sessions.GetOrCreate(caller.Token); + var responses = new List(); + try + { + if (string.IsNullOrEmpty(request.FileName)) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "file_name must not be empty."), + }); + return responses; + } + + if (!analyzer.TryGetAnalysisResult(request.FileName, out var result) || result is null) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.FileNotFound, + $"File '{request.FileName}' does not exist in the current directory."), + }); + return responses; + } + + // 头部与 GetAnalysisResult 完全一致,客户端可据此判断文件是否已分析 / 是否失败。 + responses.Add(BuildHeaderResponse(result)); + + // 仅在分析成功时才进行过滤并返回日志条目。 + if (result.State == AnalysisState.Succeeded) + { + var predicate = BuildQueryPredicate(request); + foreach (var entry in result.Entries) + { + if (predicate(entry)) + { + responses.Add(BuildEntryResponse(entry)); + } + } + } + } + catch (Exception ex) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateInternalErrorOperationStatus(ex, "querying analysis result"), + }); + _logger.LogError(ex, "An error occurred while querying analysis result."); + } + return responses; } - public Task AnalyzeFiles(AnalyzeFilesRequest request, CancellationToken cancellationToken) + /// + /// 根据某个日志文件中的 Call 类型日志,推断云服务的调用拓扑。 + /// 结点 = 服务(由 pod 名去除末尾的 -<索引> 得到,如 gateway-0 → gateway); + /// 有向边 = 调用关系,由「发出 Call 日志的服务」指向 target-service,并统计该边对应的 Call 日志条数。 + /// + public GetCallTopologyResponse GetCallTopology(GetCallTopologyRequest request, TokenInfo caller, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var analyzer = _sessions.GetOrCreate(caller.Token); + var response = new GetCallTopologyResponse + { + FileName = request.FileName, + }; + try + { + if (string.IsNullOrEmpty(request.FileName)) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "file_name must not be empty."); + return response; + } + + if (!analyzer.TryGetAnalysisResult(request.FileName, out var result) || result is null) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.FileNotFound, + $"File '{request.FileName}' does not exist in the current directory."); + return response; + } + + if (result.State != AnalysisState.Succeeded) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidOperation, + $"File '{request.FileName}' has not been successfully analyzed (state: {result.State}). Please analyze it first."); + return response; + } + + // 用 SortedSet 保证结点顺序确定(前端布局可复现)。 + var nodes = new SortedSet(); + var edgeCounts = new Dictionary<(string Source, string Target), int>(); + + foreach (var entry in result.Entries) + { + if (entry is CallLogEntry call) + { + string source = ExtractService(call.PodName); + string target = ExtractService(call.TargetService); + if (source.Length == 0 && target.Length == 0) + { + continue; + } + nodes.Add(source); + nodes.Add(target); + + var key = (source, target); + edgeCounts[key] = edgeCounts.TryGetValue(key, out int c) ? c + 1 : 1; + } + } + + foreach (var service in nodes) + { + response.Nodes.Add(new TopologyNodeMessage { Service = service }); + } + foreach (var kv in edgeCounts) + { + response.Edges.Add(new TopologyEdgeMessage + { + SourceService = kv.Key.Source, + TargetService = kv.Key.Target, + CallCount = kv.Value, + }); + } + response.Status = CreateNoErrorOperationStatus(); + } + catch (Exception ex) + { + response.Status = CreateInternalErrorOperationStatus(ex, "building call topology"); + _logger.LogError(ex, "An error occurred while building call topology."); + } + return response; } - public IReadOnlyList GetAnalysisResult(GetAnalysisResultRequest request, CancellationToken cancellationToken) + /// + /// 流式返回某条有向边(source_service → target_service)对应的所有 Call 日志。 + /// 复用与 GetAnalysisResult 相同的 header + 逐条日志格式,便于客户端复用结果展示逻辑。 + /// + public IReadOnlyList GetEdgeCallLogs(GetEdgeCallLogsRequest request, TokenInfo caller, CancellationToken cancellationToken) { - throw new NotImplementedException("TODO: T3.1"); + var analyzer = _sessions.GetOrCreate(caller.Token); + var responses = new List(); + try + { + if (string.IsNullOrEmpty(request.FileName) || + string.IsNullOrEmpty(request.SourceService) || + string.IsNullOrEmpty(request.TargetService)) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "file_name, source_service and target_service must not be empty."), + }); + return responses; + } + + if (!analyzer.TryGetAnalysisResult(request.FileName, out var result) || result is null) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.FileNotFound, + $"File '{request.FileName}' does not exist in the current directory."), + }); + return responses; + } + + // 与 GetAnalysisResult 一致:先返回 header,客户端据此判断文件状态。 + responses.Add(BuildHeaderResponse(result)); + + if (result.State == AnalysisState.Succeeded) + { + foreach (var entry in result.Entries) + { + if (entry is CallLogEntry call && + ExtractService(call.PodName) == request.SourceService && + ExtractService(call.TargetService) == request.TargetService) + { + responses.Add(BuildEntryResponse(entry)); + } + } + } + } + catch (Exception ex) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateInternalErrorOperationStatus(ex, "retrieving edge call logs"), + }); + _logger.LogError(ex, "An error occurred while retrieving edge call logs."); + } + return responses; + } + + /// + /// 按 Request ID 追踪某次请求的完整调用链(T5.2):流式返回该 request-id 对应的所有 + /// Call / Request 日志,并按时间升序排列,供客户端绘制瀑布图。Internal 日志无 RequestId,不参与。 + /// 响应复用与 GetAnalysisResult 相同的 header + 逐条日志格式。 + /// + public IReadOnlyList GetTrace(GetTraceRequest request, TokenInfo caller, CancellationToken cancellationToken) + { + var analyzer = _sessions.GetOrCreate(caller.Token); + var responses = new List(); + try + { + if (string.IsNullOrEmpty(request.FileName) || string.IsNullOrEmpty(request.RequestId)) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "file_name and request_id must not be empty."), + }); + return responses; + } + + if (!analyzer.TryGetAnalysisResult(request.FileName, out var result) || result is null) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateErrorOperationStatus( + AgentErrorCode.FileNotFound, + $"File '{request.FileName}' does not exist in the current directory."), + }); + return responses; + } + + // 与 GetAnalysisResult 一致:先返回 header,客户端据此判断文件状态。 + responses.Add(BuildHeaderResponse(result)); + + if (result.State == AnalysisState.Succeeded) + { + // 收集该 request-id 的全部日志,再按时间升序输出,便于客户端直接绘制瀑布图。 + var trace = new List(); + foreach (var entry in result.Entries) + { + string? rid = entry switch + { + CallLogEntry c => c.RequestId, + RequestLogEntry r => r.RequestId, + _ => null, + }; + if (rid == request.RequestId) + { + trace.Add(entry); + } + } + trace.Sort((a, b) => a.Timestamp.CompareTo(b.Timestamp)); + foreach (var entry in trace) + { + responses.Add(BuildEntryResponse(entry)); + } + } + } + catch (Exception ex) + { + responses.Add(new GetAnalysisResultResponse() + { + Status = CreateInternalErrorOperationStatus(ex, "retrieving trace"), + }); + _logger.LogError(ex, "An error occurred while retrieving trace."); + } + return responses; + } + + /// + /// 从 pod 名 / 服务名中提取服务名:去除末尾的 -<数字> 索引后缀。 + /// 例如 gateway-0 → gateway,my-svc-12 → my-svc,authservice → authservice(无后缀则原样返回)。 + /// + private static string ExtractService(string podOrServiceName) + { + if (string.IsNullOrEmpty(podOrServiceName)) + { + return ""; + } + return PodIndexSuffixRegex.Replace(podOrServiceName, ""); + } + + private static readonly Regex PodIndexSuffixRegex = new(@"-\d+$", RegexOptions.Compiled); + + /// + /// 根据查询请求构造一个日志条目过滤谓词。所有维度均为「未填写则不过滤」。 + /// + private static Func BuildQueryPredicate(QueryAnalysisResultRequest request) + { + // 把 gRPC 枚举转换为领域枚举,装入集合便于 O(1) 包含判断。 + var eventTypes = request.EventTypes.Select(GrpcTypeConverter.ConvertFromGrpc).ToHashSet(); + var severities = request.Severities.Select(GrpcTypeConverter.ConvertFromGrpc).ToHashSet(); + + string requestIdPattern = request.RequestIdPattern ?? ""; + string servicePattern = request.ServicePattern ?? ""; + + // message 类型字段默认可空:客户端未设置时为 null,即表示不限定该侧时间边界。 + DateTimeOffset? startTime = request.StartTime is not null + ? request.StartTime.ToDateTimeOffset() + : null; + DateTimeOffset? endTime = request.EndTime is not null + ? request.EndTime.ToDateTimeOffset() + : null; + + return entry => + { + if (eventTypes.Count > 0 && !eventTypes.Contains(entry.EventType)) + { + return false; + } + if (severities.Count > 0 && !severities.Contains(entry.Severity)) + { + return false; + } + if (servicePattern.Length > 0 && + !entry.PodName.Contains(servicePattern, StringComparison.OrdinalIgnoreCase)) + { + return false; + } + if (requestIdPattern.Length > 0) + { + // 仅 Call / Request 日志拥有 Request ID;Internal 日志在按 Request ID 过滤时一律排除。 + string? requestId = entry switch + { + CallLogEntry c => c.RequestId, + RequestLogEntry r => r.RequestId, + _ => null, + }; + if (requestId is null || + !requestId.Contains(requestIdPattern, StringComparison.OrdinalIgnoreCase)) + { + return false; + } + } + if (startTime is not null && entry.Timestamp < startTime.Value) + { + return false; + } + if (endTime is not null && entry.Timestamp > endTime.Value) + { + return false; + } + return true; + }; + } + + private static GetAnalysisResultResponse BuildHeaderResponse(AnalysisResult result) + { + return new GetAnalysisResultResponse() + { + Header = new AnalysisResultHeaderMessage() + { + FileName = result.FileName, + FullName = result.FullName, + State = GrpcTypeConverter.ConvertToGrpc(result.State), + ErrorMessage = result.ErrorMessage ?? "", + WorkerId = result.WorkerId, + }, + Status = CreateNoErrorOperationStatus(), + }; + } + + private static GetAnalysisResultResponse BuildEntryResponse(LogEntry entry) + { + return new GetAnalysisResultResponse() + { + LogEntry = GrpcTypeConverter.ConvertToGrpc(entry), + Status = CreateNoErrorOperationStatus(), + }; + } + + /// + /// 将某个已分析日志文件的结果导出为 Parquet 文件(T5.1.a.a)。 + /// 输出路径可为绝对路径,或相对于当前日志目录的相对路径 / 文件名。 + /// + public async Task ExportAnalysisResultAsync(ExportAnalysisResultRequest request, TokenInfo caller, CancellationToken cancellationToken) + { + var analyzer = _sessions.GetOrCreate(caller.Token); + var response = new ExportAnalysisResultResponse(); + try + { + if (string.IsNullOrEmpty(request.FileName)) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "file_name must not be empty."); + return response; + } + + if (string.IsNullOrWhiteSpace(request.OutputPath)) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + "output_path must not be empty."); + return response; + } + + if (!analyzer.TryGetAnalysisResult(request.FileName, out var result) || result is null) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.FileNotFound, + $"File '{request.FileName}' does not exist in the current directory."); + return response; + } + + if (result.State != AnalysisState.Succeeded) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidOperation, + $"File '{request.FileName}' has not been successfully analyzed (state: {result.State}). Please analyze it first."); + return response; + } + + // 解析输出路径:相对路径以当前日志目录为基准。 + string outputDirectory = analyzer.CurrentDirectory ?? ""; + string outputPath = Path.IsPathRooted(request.OutputPath) + ? request.OutputPath + : Path.GetFullPath(Path.Combine( + string.IsNullOrEmpty(outputDirectory) ? Directory.GetCurrentDirectory() : outputDirectory, + request.OutputPath)); + + if (!outputPath.EndsWith(".parquet", StringComparison.OrdinalIgnoreCase)) + { + outputPath += ".parquet"; + } + + if (File.Exists(outputPath) && !request.Overwrite) + { + response.Status = CreateErrorOperationStatus( + AgentErrorCode.InvalidArgument, + $"Output file already exists: '{outputPath}'. Set overwrite=true to replace it."); + return response; + } + + string? dir = Path.GetDirectoryName(outputPath); + if (!string.IsNullOrEmpty(dir) && !Directory.Exists(dir)) + { + Directory.CreateDirectory(dir); + } + + int count = await ParquetLogWriter.WriteAsync(outputPath, result.Entries, cancellationToken); + response.WrittenPath = outputPath; + response.EntryCount = count; + response.Status = CreateNoErrorOperationStatus(); + } + catch (Exception ex) + { + response.Status = CreateInternalErrorOperationStatus(ex, "exporting analysis result to parquet"); + _logger.LogError(ex, "An error occurred while exporting analysis result to parquet."); + } + return response; + } + + // —— Token 管理(T5.1.a.b)。调用者已被 AgentService 校验为管理员,这里只做业务逻辑。 —— + + public CreateTokenResponse CreateToken(CreateTokenRequest request, TokenInfo caller) + { + var role = ConvertRole(request.Role); + var info = _tokens.CreateToken(role, request.Note ?? ""); + _logger.LogInformation("Admin {Caller} created a new {Role} token.", Mask(caller.Token), role); + return new CreateTokenResponse + { + Status = CreateNoErrorOperationStatus(), + Token = ToMessage(info), + }; + } + + public OperationStatusMessage DeleteToken(DeleteTokenRequest request, TokenInfo caller) + { + var (ok, error) = _tokens.TryDelete(request.Token ?? ""); + if (!ok) + { + return CreateErrorOperationStatus(AgentErrorCode.InvalidArgument, error ?? "Failed to delete token."); + } + _logger.LogInformation("Admin {Caller} deleted token {Token}.", Mask(caller.Token), Mask(request.Token ?? "")); + return CreateNoErrorOperationStatus(); + } + + public ListTokensResponse ListTokens(Empty empty, TokenInfo caller) + { + var all = _tokens.List(); + var response = new ListTokensResponse + { + Status = CreateNoErrorOperationStatus(), + CallerToken = caller.Token, + }; + foreach (var info in all) + { + response.Tokens.Add(ToMessage(info)); + } + return response; + } + + public OperationStatusMessage SetTokenRole(SetTokenRoleRequest request, TokenInfo caller) + { + var role = ConvertRole(request.Role); + var (ok, error) = _tokens.TrySetRole(request.Token ?? "", role); + if (!ok) + { + return CreateErrorOperationStatus(AgentErrorCode.InvalidArgument, error ?? "Failed to set token role."); + } + _logger.LogInformation("Admin {Caller} set token {Token} role to {Role}.", + Mask(caller.Token), Mask(request.Token ?? ""), role); + return CreateNoErrorOperationStatus(); + } + + private static TokenRole ConvertRole(TokenRoleEnum role) => role switch + { + TokenRoleEnum.TokenAdmin => TokenRole.Admin, + _ => TokenRole.Normal, + }; + + private static TokenRoleEnum ConvertRole(TokenRole role) => role switch + { + TokenRole.Admin => TokenRoleEnum.TokenAdmin, + _ => TokenRoleEnum.TokenNormal, + }; + + private static TokenInfoMessage ToMessage(TokenInfo info) => new() + { + Token = info.Token, + Role = ConvertRole(info.Role), + Note = info.Note, + }; + + /// + /// 在日志中对 token 做脱敏:只保留前 6 个字符,避免完整 token 落入日志后被泄露。 + /// 启动时输出的 bootstrap admin token 是例外(由 Program.cs 显式输出,方便使用者取用)。 + /// + private static string Mask(string token) + { + if (string.IsNullOrEmpty(token)) + { + return ""; + } + return token.Length <= 6 ? token + "…" : token[..6] + "…"; } } } diff --git a/src/LogAnalyzerAgent/Auth/SessionManager.cs b/src/LogAnalyzerAgent/Auth/SessionManager.cs new file mode 100644 index 0000000..2be0c58 --- /dev/null +++ b/src/LogAnalyzerAgent/Auth/SessionManager.cs @@ -0,0 +1,31 @@ +using System.Collections.Concurrent; +using LogAnalyzer; + +namespace LogAnalyzerAgent.Auth +{ + /// + /// 每个合法 token 对应一个独立的用户;不同用户的日志分析操作完全互不干扰(T5.1.a.b)。 + /// + /// 实现方式:以 token 为键,为每个用户懒加载一个独立的 实例。 + /// 由于 内部以私有字段保存「当前目录 / 文件列表 / 分析结果」, + /// 给每个用户各分配一个实例,即可天然地实现: + /// + /// 目录隔离:A 改的目录不影响 B; + /// 分析结果隔离:A 分析过的文件、缓存的结果,B 看不到也改不了。 + /// + /// + public sealed class SessionManager + { + private readonly ConcurrentDictionary _analyzers = new(); + + /// + /// 取得(或懒创建)某个 token 对应用户的 Analyzer。token 由 校验通过后传入。 + /// + public LogFileAnalyzer GetOrCreate(string token) + { + // 传入 null 目录:构造一个尚未设置目录的空 Analyzer(HasDirectory == false), + // 用户随后通过 ChangeDirectory 选择自己的日志目录。 + return _analyzers.GetOrAdd(token, _ => new LogFileAnalyzer(null)); + } + } +} diff --git a/src/LogAnalyzerAgent/Auth/TokenStore.cs b/src/LogAnalyzerAgent/Auth/TokenStore.cs new file mode 100644 index 0000000..a551f72 --- /dev/null +++ b/src/LogAnalyzerAgent/Auth/TokenStore.cs @@ -0,0 +1,175 @@ +using System.Security.Cryptography; + +namespace LogAnalyzerAgent.Auth +{ + /// + /// Token 权限等级(T5.1.a.b)。 + /// 与 gRPC 的 TokenRoleEnum 一一对应,但属于 Agent 的领域模型,避免在业务逻辑中直接依赖生成代码。 + /// + public enum TokenRole + { + /// 普通权限:可进行各项日志分析操作。 + Normal, + + /// 管理员权限:在普通权限基础上,可对其他 token 进行增删与权限调整。 + Admin, + } + + /// + /// 一个已签发的 token 的信息。Role / Note 在运行期可被管理员修改,故为可变类, + /// 但所有读写都经由 在同一把锁下进行,保证线程安全。 + /// + public sealed class TokenInfo + { + public string Token { get; } + public TokenRole Role { get; set; } + public string Note { get; set; } + + public TokenInfo(string token, TokenRole role, string note) + { + Token = token; + Role = role; + Note = note ?? ""; + } + + /// + /// 复制一份,供 返回不可变快照使用,避免调用方持有可变引用。 + /// + public TokenInfo Clone() => new(Token, Role, Note); + } + + /// + /// 维护所有已签发 token 的存储(T5.1.a.b)。 + /// + /// 职责: + /// + /// 启动时生成一个管理员 token()。 + /// 校验请求携带的 token()。 + /// 管理员的增 / 删 / 列 / 改权限操作。 + /// + /// + /// 线程安全:所有公开方法都在同一把 _lock 下完成「检查 + 修改」的复合操作, + /// 因此「不能删除 / 降级最后一个管理员」的防锁死策略是原子的。 + /// + public sealed class TokenStore + { + private readonly object _lock = new(); + private readonly Dictionary _tokens = new(); + + /// + /// 启动时签发一个管理员 token 并返回,供 Program.cs 通过 logger 输出。 + /// 该方法在 Agent 生命周期内只应调用一次。 + /// + public string CreateAdminToken() + { + var info = new TokenInfo(GenerateToken(), TokenRole.Admin, "bootstrap admin token"); + lock (_lock) + { + _tokens[info.Token] = info; + } + return info.Token; + } + + /// + /// 校验 token 合法性。返回 null 表示不存在 / 非法。 + /// + public TokenInfo? TryGet(string? token) + { + if (string.IsNullOrEmpty(token)) + { + return null; + } + lock (_lock) + { + return _tokens.TryGetValue(token!, out var info) ? info : null; + } + } + + /// + /// 签发一个新 token。 + /// + public TokenInfo CreateToken(TokenRole role, string note) + { + var info = new TokenInfo(GenerateToken(), role, note); + lock (_lock) + { + _tokens[info.Token] = info; + } + return info; + } + + /// + /// 删除一个 token。带防锁死策略:不允许删除最后一个管理员 token。 + /// 返回 (success, errorMessage)。 + /// + public (bool success, string? error) TryDelete(string token) + { + lock (_lock) + { + if (!_tokens.TryGetValue(token, out var info)) + { + return (false, "Token does not exist."); + } + if (info.Role == TokenRole.Admin && AdminCountNoLock() <= 1) + { + return (false, "Cannot delete the last remaining admin token (would lock out management)."); + } + _tokens.Remove(token); + return (true, null); + } + } + + /// + /// 调整一个 token 的权限。带防锁死策略:不允许把最后一个管理员降级为普通权限。 + /// + public (bool success, string? error) TrySetRole(string token, TokenRole role) + { + lock (_lock) + { + if (!_tokens.TryGetValue(token, out var info)) + { + return (false, "Token does not exist."); + } + if (info.Role == role) + { + return (true, null); + } + if (info.Role == TokenRole.Admin && role == TokenRole.Normal && AdminCountNoLock() <= 1) + { + return (false, "Cannot demote the last remaining admin token (would lock out management)."); + } + info.Role = role; + return (true, null); + } + } + + /// + /// 列出所有 token(返回快照副本)。同时返回当前管理员数量供上层提示。 + /// + public IReadOnlyList List() + { + lock (_lock) + { + return _tokens.Values.Select(t => t.Clone()).ToList(); + } + } + + private int AdminCountNoLock() + { + return _tokens.Values.Count(t => t.Role == TokenRole.Admin); + } + + /// + /// 生成一个高熵、URL 安全的随机 token:24 字节密码学随机数 → base64url。 + /// + private static string GenerateToken() + { + Span bytes = stackalloc byte[24]; + RandomNumberGenerator.Fill(bytes); + return Convert.ToBase64String(bytes) + .Replace('+', '-') + .Replace('/', '_') + .TrimEnd('='); + } + } +} diff --git a/src/LogAnalyzerAgent/Program.cs b/src/LogAnalyzerAgent/Program.cs index 4e247ad..1319b25 100644 --- a/src/LogAnalyzerAgent/Program.cs +++ b/src/LogAnalyzerAgent/Program.cs @@ -1,5 +1,5 @@ -using LogAnalyzer; -using LogAnalyzerAgent.Applications; +using LogAnalyzerAgent.Applications; +using LogAnalyzerAgent.Auth; using LogAnalyzerAgent.Services; using Microsoft.AspNetCore.Server.Kestrel.Core; @@ -17,13 +17,23 @@ builder.Services.AddGrpc(); builder.Services.AddCors(); -// 依赖注入,且有状态服务需要单例 -builder.Services.AddSingleton(); +// 依赖注入,且有状态服务需要单例。 +// T5.1.a.b:以 TokenStore(鉴权)+ SessionManager(按 token 隔离的 Analyzer)替代原先的单个共享 Analyzer。 +builder.Services.AddSingleton(); +builder.Services.AddSingleton(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); var app = builder.Build(); +// 启动时生成一个管理员 token,并以 log 形式输出,供使用者首次登录客户端时取用(T5.1.a.b)。 +var tokenStore = app.Services.GetRequiredService(); +var bootstrapLogger = app.Services.GetRequiredService().CreateLogger("Bootstrap"); +string adminToken = tokenStore.CreateAdminToken(); +bootstrapLogger.LogInformation( + "Admin token generated. Use it to log in the client (File -> Connect...) and manage other tokens: {Token}", + adminToken); + var whiteList = new HashSet() { "http://localhost:5235", diff --git a/src/LogAnalyzerAgent/Services/AgentService.cs b/src/LogAnalyzerAgent/Services/AgentService.cs index 591dcad..ccbe857 100644 --- a/src/LogAnalyzerAgent/Services/AgentService.cs +++ b/src/LogAnalyzerAgent/Services/AgentService.cs @@ -1,55 +1,195 @@ -using Google.Protobuf.WellKnownTypes; +using Google.Protobuf.WellKnownTypes; using Grpc.Core; -using LogAnalyzer; -using LogAnalyzerRpc.Protos; -using LogAnalyzerRpc; -using LogParser.Visitors; using LogAnalyzerAgent.Applications; +using LogAnalyzerAgent.Auth; +using LogAnalyzerRpc; +using LogAnalyzerRpc.Protos; namespace LogAnalyzerAgent.Services { + /// + /// gRPC 服务实现(T5.1.a.b 起承担鉴权职责)。 + /// + /// 鉴权方式:客户端在每次调用的 metadata 中携带 authorization: Bearer <token>, + /// 本服务通过 取出并校验: + /// + /// 缺失 / 非法 → 抛 RpcException(Unauthenticated),拒绝该次调用; + /// 合法 → 得到 ,作为调用者身份传入业务层,并由 + /// 路由到该用户专属的 Analyzer。 + /// + /// 管理员专用的 token 管理 RPC 还会额外经过 校验。 + /// public class AgentService : LogAnalyzerAgentService.LogAnalyzerAgentServiceBase { private readonly AgentSession _session; + private readonly TokenStore _tokens; + + // 标准的 Bearer token 认证头前缀。 + private const string BearerPrefix = "Bearer "; - public AgentService(AgentSession session) + public AgentService(AgentSession session, TokenStore tokens) { _session = session; + _tokens = tokens; + } + + /// + /// 从请求 metadata 中解析并校验 token。失败时抛出 , + /// 由 gRPC 框架转换为对应的状态码返回给客户端。 + /// + private TokenInfo Authorize(ServerCallContext context) + { + string? headerValue = null; + foreach (var entry in context.RequestHeaders) + { + if (entry.Key == "authorization") + { + headerValue = entry.Value; + break; + } + } + + string? token = null; + if (!string.IsNullOrEmpty(headerValue) && headerValue.StartsWith(BearerPrefix, StringComparison.OrdinalIgnoreCase)) + { + token = headerValue[BearerPrefix.Length..].Trim(); + } + + TokenInfo? info = token is null ? null : _tokens.TryGet(token); + if (info is null) + { + throw new RpcException(new Status(StatusCode.Unauthenticated, "Missing or invalid token.")); + } + return info; + } + + /// + /// 要求调用者为管理员,否则抛 。 + /// + private static void RequireAdmin(TokenInfo caller) + { + if (caller.Role != TokenRole.Admin) + { + throw new RpcException(new Status(StatusCode.PermissionDenied, "Admin privilege required for this operation.")); + } } public override Task Ping(Empty empty, ServerCallContext context) { - return _session.Ping(empty, context.CancellationToken); + var caller = Authorize(context); + return _session.Ping(empty, caller, context.CancellationToken); } public override Task GetAgentStatus(Empty empty, ServerCallContext context) { - return _session.GetAgentStatus(empty, context.CancellationToken); + var caller = Authorize(context); + return _session.GetAgentStatus(empty, caller, context.CancellationToken); } public override Task ChangeDirectory(ChangeDirectoryRequest request, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + var caller = Authorize(context); + return _session.ChangeDirectory(request, caller, context.CancellationToken); } public override Task GetLogFiles(Empty empty, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + var caller = Authorize(context); + return _session.GetLogFiles(empty, caller, context.CancellationToken); } public override Task AnalyzeAll(AnalyzeAllRequest request, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + var caller = Authorize(context); + return _session.AnalyzeAll(request, caller, context.CancellationToken); } public override Task AnalyzeFiles(AnalyzeFilesRequest request, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + var caller = Authorize(context); + return _session.AnalyzeFiles(request, caller, context.CancellationToken); } public override async Task GetAnalysisResult(GetAnalysisResultRequest request, IServerStreamWriter responseStream, ServerCallContext context) { - throw new NotImplementedException("TODO: T3.1"); + var caller = Authorize(context); + var responses = _session.GetAnalysisResult(request, caller, context.CancellationToken); + foreach (var response in responses) + { + await responseStream.WriteAsync(response); + } + } + + public override async Task QueryAnalysisResult(QueryAnalysisResultRequest request, IServerStreamWriter responseStream, ServerCallContext context) + { + var caller = Authorize(context); + var responses = _session.QueryAnalysisResult(request, caller, context.CancellationToken); + foreach (var response in responses) + { + await responseStream.WriteAsync(response); + } + } + + public override Task GetCallTopology(GetCallTopologyRequest request, ServerCallContext context) + { + var caller = Authorize(context); + return Task.FromResult(_session.GetCallTopology(request, caller, context.CancellationToken)); + } + + public override async Task GetEdgeCallLogs(GetEdgeCallLogsRequest request, IServerStreamWriter responseStream, ServerCallContext context) + { + var caller = Authorize(context); + var responses = _session.GetEdgeCallLogs(request, caller, context.CancellationToken); + foreach (var response in responses) + { + await responseStream.WriteAsync(response); + } + } + + public override async Task GetTrace(GetTraceRequest request, IServerStreamWriter responseStream, ServerCallContext context) + { + var caller = Authorize(context); + var responses = _session.GetTrace(request, caller, context.CancellationToken); + foreach (var response in responses) + { + await responseStream.WriteAsync(response); + } + } + + public override Task ExportAnalysisResult(ExportAnalysisResultRequest request, ServerCallContext context) + { + var caller = Authorize(context); + return _session.ExportAnalysisResultAsync(request, caller, context.CancellationToken); + } + + // —— Token 管理 RPC(需管理员权限)—— + + public override Task CreateToken(CreateTokenRequest request, ServerCallContext context) + { + var caller = Authorize(context); + RequireAdmin(caller); + return Task.FromResult(_session.CreateToken(request, caller)); + } + + public override Task DeleteToken(DeleteTokenRequest request, ServerCallContext context) + { + var caller = Authorize(context); + RequireAdmin(caller); + return Task.FromResult(_session.DeleteToken(request, caller)); + } + + public override Task ListTokens(Empty request, ServerCallContext context) + { + var caller = Authorize(context); + RequireAdmin(caller); + return Task.FromResult(_session.ListTokens(request, caller)); + } + + public override Task SetTokenRole(SetTokenRoleRequest request, ServerCallContext context) + { + var caller = Authorize(context); + RequireAdmin(caller); + return Task.FromResult(_session.SetTokenRole(request, caller)); } } } diff --git a/src/LogAnalyzerClient/Directory.Packages.props b/src/LogAnalyzerClient/Directory.Packages.props index 8c9efe7..0856df1 100644 --- a/src/LogAnalyzerClient/Directory.Packages.props +++ b/src/LogAnalyzerClient/Directory.Packages.props @@ -6,12 +6,13 @@ - - - + + + + - - + + diff --git a/src/LogAnalyzerClient/LogAnalyzerClient.Browser/Program.cs b/src/LogAnalyzerClient/LogAnalyzerClient.Browser/Program.cs index 3955402..0caf5da 100644 --- a/src/LogAnalyzerClient/LogAnalyzerClient.Browser/Program.cs +++ b/src/LogAnalyzerClient/LogAnalyzerClient.Browser/Program.cs @@ -1,5 +1,6 @@ using Avalonia; using Avalonia.Browser; +using Grpc.Core.Interceptors; using Grpc.Net.Client; using Grpc.Net.Client.Web; using LogAnalyzerClient; @@ -13,15 +14,16 @@ internal sealed partial class Program { internal class ClientFactory : IClientFactory { - public LogAnalyzerAgentServiceClient CreateClient(string address) + public AgentClientHandle CreateClient(string address, string token) { var handler = new GrpcWebHandler(GrpcWebMode.GrpcWeb, new HttpClientHandler()); var channel = GrpcChannel.ForAddress(address, new GrpcChannelOptions() { HttpHandler = handler }); - var client = new LogAnalyzerAgentServiceClient(channel); - return client; + // 用拦截器把 token 附加到每一次 gRPC 调用,满足 Agent 端的鉴权要求(T5.1.a.b)。 + var client = new LogAnalyzerAgentServiceClient(channel.Intercept(new TokenInterceptor(token))); + return new AgentClientHandle(client, channel); } } diff --git a/src/LogAnalyzerClient/LogAnalyzerClient.Desktop/Program.cs b/src/LogAnalyzerClient/LogAnalyzerClient.Desktop/Program.cs index 11c3468..8843a3f 100644 --- a/src/LogAnalyzerClient/LogAnalyzerClient.Desktop/Program.cs +++ b/src/LogAnalyzerClient/LogAnalyzerClient.Desktop/Program.cs @@ -1,4 +1,5 @@ using Avalonia; +using Grpc.Core.Interceptors; using Grpc.Net.Client; using LogAnalyzerClient.Services; using LogAnalyzerRpc.Protos; @@ -9,11 +10,12 @@ namespace LogAnalyzerClient.Desktop { internal class ClientFactory : IClientFactory { - public LogAnalyzerAgentServiceClient CreateClient(string address) + public AgentClientHandle CreateClient(string address, string token) { var channel = GrpcChannel.ForAddress(address); - var client = new LogAnalyzerAgentServiceClient(channel); - return client; + // 用拦截器把 token 附加到每一次 gRPC 调用,满足 Agent 端的鉴权要求(T5.1.a.b)。 + var client = new LogAnalyzerAgentServiceClient(channel.Intercept(new TokenInterceptor(token))); + return new AgentClientHandle(client, channel); } } diff --git a/src/LogAnalyzerClient/LogAnalyzerClient/App.axaml b/src/LogAnalyzerClient/LogAnalyzerClient/App.axaml index 1e08725..76b1847 100644 --- a/src/LogAnalyzerClient/LogAnalyzerClient/App.axaml +++ b/src/LogAnalyzerClient/LogAnalyzerClient/App.axaml @@ -11,6 +11,7 @@ + \ No newline at end of file diff --git a/src/LogAnalyzerClient/LogAnalyzerClient/Converters/SeverityToBrushConverter.cs b/src/LogAnalyzerClient/LogAnalyzerClient/Converters/SeverityToBrushConverter.cs new file mode 100644 index 0000000..4fecd90 --- /dev/null +++ b/src/LogAnalyzerClient/LogAnalyzerClient/Converters/SeverityToBrushConverter.cs @@ -0,0 +1,38 @@ +using System; +using System.Globalization; +using Avalonia.Data.Converters; +using Avalonia.Media; +using LogParser.Models; + +namespace LogAnalyzerClient.Converters +{ + /// + /// 把日志等级映射为高亮背景色(T5.1.b.a):Info=蓝、Warning=橙、Error=红。 + /// 用于结果表格 Severity 列圆角胶囊的 Background 绑定。 + /// + public sealed class SeverityToBrushConverter : IValueConverter + { + private static readonly IBrush SevInfo = new SolidColorBrush(Color.FromRgb(0x2B, 0x6C, 0xB0)); + private static readonly IBrush SevWarning = new SolidColorBrush(Color.FromRgb(0xDD, 0x6B, 0x20)); + private static readonly IBrush SevError = new SolidColorBrush(Color.FromRgb(0xC5, 0x30, 0x30)); + private static readonly IBrush Fallback = new SolidColorBrush(Colors.Gray); + + public object? Convert(object? value, Type targetType, object? parameter, CultureInfo culture) + { + if (value is LogSeverity severity) + { + return severity switch + { + LogSeverity.Info => SevInfo, + LogSeverity.Warning => SevWarning, + LogSeverity.Error => SevError, + _ => Fallback, + }; + } + return Fallback; + } + + public object? ConvertBack(object? value, Type targetType, object? parameter, CultureInfo culture) + => throw new NotSupportedException(); + } +} diff --git a/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml b/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml index 4f44da3..ed473e3 100644 --- a/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml +++ b/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml @@ -4,14 +4,25 @@ xmlns:mc="http://schemas.openxmlformats.org/markup-compatibility/2006" mc:Ignorable="d" x:Class="LogAnalyzerClient.ConnectDialog" - d:DesignWidth="440" d:DesignHeight="180" - Width="440" Height="180" + d:DesignWidth="440" d:DesignHeight="240" + Width="440" Height="240" CanResize="False" WindowStartupLocation="CenterOwner" Title="Connect..."> - - + + + + + + + + + diff --git a/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml.cs b/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml.cs index 7d9ddb1..b8adc3a 100644 --- a/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml.cs +++ b/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ConnectDialog.axaml.cs @@ -1,10 +1,13 @@ -using Avalonia; using Avalonia.Controls; using Avalonia.Interactivity; -using Avalonia.Markup.Xaml; namespace LogAnalyzerClient; +/// +/// 连接对话框的返回值:Agent 地址 + 鉴权 token(T5.1.a.b)。 +/// +internal sealed record ConnectResult(string Address, string Token); + public partial class ConnectDialog : Window { public ConnectDialog() @@ -12,18 +15,14 @@ public ConnectDialog() InitializeComponent(); } - public ConnectDialog(string currentAddress) : this() - { - AddressTextBox.Text = currentAddress; - } - private void ConnectButton_Click(object? sender, RoutedEventArgs e) { - Close(AddressTextBox.Text); + // 返回地址与 token;校验是否为空交给调用方(便于给出统一的错误提示)。 + Close(new ConnectResult(AddressTextBox.Text ?? "", TokenTextBox.Text ?? "")); } private void CancelButton_Click(object? sender, RoutedEventArgs e) { Close(null); } -} \ No newline at end of file +} diff --git a/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ExportDialog.axaml b/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ExportDialog.axaml new file mode 100644 index 0000000..dc7dec8 --- /dev/null +++ b/src/LogAnalyzerClient/LogAnalyzerClient/Dialogs/ExportDialog.axaml @@ -0,0 +1,49 @@ + + + + + + + + + + + + + + + + + + + +