Compare commits

..
14 Commits
Author SHA256 Message Date
lixen 80e0017140 v0.17.0 mtls added to network chain, sinks, and sources 2026-08-29 18:44:44 -04:00
lixen b2e36be53f v0.16.1 doc update 2026-08-29 16:37:20 -04:00
lixen 325b51d840 v0.16.0 tls added to tcp/http sources and sinks 2026-07-18 04:10:22 -04:00
lixen f0aca019f3 v0.15.0 deprecated http/tcp plugins (fasthttp/gnet2), migrated stdtcp/stdhttp to tcp/http 2026-07-17 18:29:26 -04:00
lixen dc3e9910a3 v0.14.0 tcp/http sinks with standard library added, chain aggregate test script added 2026-07-17 17:21:50 -04:00
lixen e96ede32ae v0.13.1 folder restructure, test script added, format adapter async fix 2026-07-17 09:09:00 -04:00
lixen a48514a2eb v0.13.0 doc update, refactor, http/tcp chain source and sink added 2026-07-17 05:58:44 -04:00
lixen fa7f41c059 v0.12.2 config validation updated, distributed to their own pacakges and plugins 2026-01-08 12:34:40 -05:00
lixen 36e7f8a6a9 v0.12.1 config cleanup 2026-01-05 02:20:38 -05:00
lixen 15f484f98b v0.12.0 tcp and http server sinks added antested, no tls or access control 2026-01-04 17:21:08 -05:00
lixen 70bf6a8060 v0.11.0 external formatter and sanitizer integrated, refactored 2025-12-10 08:37:26 -05:00
lixen 98ace914f7 v0.10.0 flow and plugin structure, networking and commands removed, dirty 2025-11-11 16:42:09 -05:00
lixen 22652f9e53 v0.9.0 restructure for flow architecture, dirty 2025-11-09 15:08:50 -05:00
lixen dcf803bac1 v0.8.0 decoupled session management and auth, auth deprecated except mtls, session management, tls and mtls flows fixed, docs and config outdated 2025-11-06 16:43:46 -05:00
12 changed files with 105 additions and 306 deletions
-40
View File
@@ -1,40 +0,0 @@
# Builder pin and go.mod directive are the same patch release deliberately;
# an older builder reports the mismatch only after downloading the module graph.
ARG GO_VERSION=1.27.1
FROM docker.io/library/golang:${GO_VERSION}-alpine AS build
# Git supplies Go's VCS build information; only /out/logwisp crosses stages.
RUN apk add --no-cache git
WORKDIR /src
COPY go.mod go.sum ./
RUN go mod download
COPY . .
ARG TARGETOS=linux
ARG TARGETARCH=amd64
ARG VERSION=dev
ARG REVISION=unknown
RUN CGO_ENABLED=0 GOOS=${TARGETOS} GOARCH=${TARGETARCH} \
go build -trimpath \
-ldflags="-s -w -X logwisp/internal/version.Version=${VERSION} -X logwisp/internal/version.GitCommit=${REVISION}" \
-o /out/logwisp ./cmd/logwisp
FROM scratch
ARG VERSION=dev
ARG REVISION=unknown
LABEL org.opencontainers.image.title="logwisp" \
org.opencontainers.image.description="Log transport: sources, flow, sinks" \
org.opencontainers.image.source="https://github.com/lixenwraith/logwisp" \
org.opencontainers.image.revision="${REVISION}" \
org.opencontainers.image.version="${VERSION}" \
org.opencontainers.image.licenses="BSD-3-Clause"
COPY --from=build /out/logwisp /logwisp
# Numeric identity is required in scratch and satisfies a restricted pod spec.
USER 65532:65532
ENTRYPOINT ["/logwisp"]
-3
View File
@@ -152,9 +152,6 @@ directory = "./" # Directory to monitor (required, not re
pattern = "*.log" # Glob pattern (* and ? only) pattern = "*.log" # Glob pattern (* and ? only)
## Tailing an already-open file polls at a fixed 100ms, regardless of this value. ## Tailing an already-open file polls at a fixed 100ms, regardless of this value.
check_interval_ms = 100 # Directory rescan interval (min 10) check_interval_ms = 100 # Directory rescan interval (min 10)
## raw = true never parses a line; with format type "raw" the file is relayed byte for byte.
raw = false # Keep the whole line as the message
from = "end" # "end" or "start" of a newly discovered file
## Console source (stdin, single instance per pipeline) ## Console source (stdin, single instance per pipeline)
# [[pipelines.plugin_sources]] # [[pipelines.plugin_sources]]
+8 -16
View File
@@ -30,10 +30,8 @@ Omitting `[pipelines.flow.format]` entirely selects `raw`.
### raw ### raw
Passthrough. `FlagRaw` bypasses formatting and sanitization: the message reaches Passthrough. `FlagRaw` bypasses both formatting and sanitization, so the
the sink exactly as the source produced it, with no timestamp, level or source message reaches the sink exactly as the source produced it.
prefix added. An entry that also carries `fields` gets the fields JSON appended
verbatim after a single space — `raw` never drops data and never re-encodes it.
```toml ```toml
[pipelines.flow.format] [pipelines.flow.format]
@@ -44,12 +42,6 @@ Fastest option, and the right one when you are relaying text that is already in
its final form. Note that it also bypasses sanitization, so control characters its final form. Note that it also bypasses sanitization, so control characters
in the source data reach your sinks intact. in the source data reach your sinks intact.
Byte-exact transport needs a source that does not split the line: the `console`
source, or the `file` source with `raw = true`. Both put the whole line —
newline included — in the message and leave `fields` empty. The `file` source's
JSON branch splits a line into message and fields, so `raw` reassembles it as
`<msg> <fields>` rather than reproducing the original object.
### txt ### txt
Human-readable line output with a timestamp and level. Human-readable line output with a timestamp and level.
@@ -97,7 +89,7 @@ you need to override the defaults.
With `flags = 0` the formatter selects `1` for `type = "raw"` and `6` With `flags = 0` the formatter selects `1` for `type = "raw"` and `6`
(timestamp + level) for every other type. `8` is added automatically whenever an (timestamp + level) for every other type. `8` is added automatically whenever an
entry carries parseable `fields` and `1` is not set; `1` always wins. entry carries parseable `fields`.
Examples: `flags = 4` for level only, no timestamp; `flags = 2` for timestamp Examples: `flags = 4` for level only, no timestamp; `flags = 2` for timestamp
only, no level. only, no level.
@@ -141,11 +133,11 @@ when writing downstream parsers or grep patterns.
## Structured Fields ## Structured Fields
When an entry carries `Fields` (raw JSON) and `FlagRaw` is not set, the When an entry carries `Fields` (raw JSON), the formatter parses it and switches
formatter parses it and switches to structured rendering by adding the to structured rendering by adding the `StructuredJSON` flag automatically.
`StructuredJSON` flag automatically. Fields reach a pipeline in two ways: from Fields reach a pipeline in two ways: from the `file` source when a tailed line
the `file` source when a tailed line parses as JSON with a `fields` key, and parses as JSON with a `fields` key, and from the heartbeat generator when
from the heartbeat generator when `include_stats = true`. `include_stats = true`.
## Choosing a Configuration ## Choosing a Configuration
+15 -26
View File
@@ -4,7 +4,7 @@
- **Operating systems**: Linux (kernel 6.10+), FreeBSD (14.0+) - **Operating systems**: Linux (kernel 6.10+), FreeBSD (14.0+)
- **Architecture**: amd64 - **Architecture**: amd64
- **Go**: 1.27.1 or newer, to build from source - **Go**: 1.26 or newer, to build from source
## Building from Source ## Building from Source
@@ -37,24 +37,6 @@ go build -o bin/logwisp ./cmd/logwisp
`go install github.com/lixenwraith/logwisp/cmd/logwisp@latest` also works, with `go install github.com/lixenwraith/logwisp/cmd/logwisp@latest` also works, with
the same loss of version metadata. the same loss of version metadata.
## Container Image
The root `Dockerfile` builds the same package into `scratch` under UID 65532,
static and stripped. There is no shell and no config in the image: mount one and
name it, as the binary has no daemon mode and no built-in defaults worth running.
```bash
REV=$(git rev-parse HEAD)
docker build -t "logwisp:$(git rev-parse --short HEAD)" \
--build-arg VERSION="$(git describe --tags --always)" \
--build-arg REVISION="$REV" .
docker run --rm -v /etc/logwisp:/etc/logwisp:ro logwisp:... -c /etc/logwisp/logwisp.toml
```
Sinks that listen (`http`, `tcp`) need their ports published; the read-only
root filesystem and dropped capabilities a restricted runtime imposes are all
compatible with it, provided a `file` sink's directory is writable by 65532.
## Configuration ## Configuration
Copy the annotated reference configuration and edit it: Copy the annotated reference configuration and edit it:
@@ -191,20 +173,27 @@ mode; see [Operations](operations.md#checking-a-configuration).
## Test Scripts ## Test Scripts
End-to-end scripts under `test/` run against a local build: Two end-to-end scripts under `test/` build multi-node chain topologies against a
local build:
```bash ```bash
make make
./test/chain-test.sh --auto # two independent relay pipelines ./test/chain-test.sh --auto # two independent relay pipelines
./test/chain-aggregate-test.sh --auto # fan-in: both edges into one pipeline ./test/chain-aggregate-test.sh --auto # fan-in: both edges into one pipeline
./test/mtls-chain-test.sh --auto # the same fan-in under mTLS
./test/passthrough-test.sh # file source relays a wide envelope intact
``` ```
Without `--auto` the chain scripts run the relay in the foreground for Without `--auto` they run the relay in the foreground for interactive
interactive inspection. They need bash 5+, coreutils, and curl, and they bind inspection. They need bash 5+, coreutils, and curl, and they bind ports
ports 1580115804. The pass-through test binds nothing. Generated configuration 1580115804. Generated configuration and logs land in `test/run/`.
and logs land in `test/run/`.
> Two of the three `--auto` assertions currently report `FAIL` against a
> working build. They grep the sink output for `"source":"edge-tcp/` and
> `"node":"edge-http"`, but the JSON formatter emits the `node/source` label
> under the key `trace`. The transport itself is healthy — the
> `total_processed` assertion passes and the streamed entries carry
> `"trace":"edge-tcp/random_rand"` as expected. Until the assertions are
> updated, verify the streams by eye with `nc 127.0.0.1 15803` and
> `curl -sN http://127.0.0.1:15804/stream`.
## Uninstall ## Uninstall
+1 -2
View File
@@ -231,8 +231,7 @@ pipeline `total_dropped_by_sink` (sink backed up?), sink `total_processed`.
- The watcher seeks to end-of-file on start; only content appended afterwards is - The watcher seeks to end-of-file on start; only content appended afterwards is
read. Positions are in memory, so a restart re-seeks to end and anything read. Positions are in memory, so a restart re-seeks to end and anything
written during the downtime is lost. `from = "start"` reads each file whole written during the downtime is lost.
instead, and replays it on every restart.
- `pattern` is a filename glob with `*` and `?` only, and matching is not - `pattern` is a filename glob with `*` and `?` only, and matching is not
recursive. recursive.
- `check_interval_ms` governs how quickly a *new file* is noticed; tailing an - `check_interval_ms` governs how quickly a *new file* is noticed; tailing an
+5 -17
View File
@@ -32,8 +32,6 @@ type = "file"
directory = "/var/log/myapp" directory = "/var/log/myapp"
pattern = "*.log" pattern = "*.log"
check_interval_ms = 100 check_interval_ms = 100
raw = false
from = "end"
``` ```
| Option | Type | Default | Description | | Option | Type | Default | Description |
@@ -41,8 +39,6 @@ from = "end"
| `directory` | string | **required** | Directory to scan; not recursive | | `directory` | string | **required** | Directory to scan; not recursive |
| `pattern` | string | `*` | Glob over filenames; `*` and `?` only | | `pattern` | string | `*` | Glob over filenames; `*` and `?` only |
| `check_interval_ms` | int | `100` | Directory rescan interval; minimum `10` | | `check_interval_ms` | int | `100` | Directory rescan interval; minimum `10` |
| `raw` | bool | `false` | Never parse a line: the whole line is the message |
| `from` | string | `end` | Where a new watcher starts: `end` or `start` of the file |
**Behaviour** **Behaviour**
@@ -53,23 +49,15 @@ from = "end"
stopped and removed on the next scan. stopped and removed on the next scan.
- A new watcher seeks to end-of-file. Positions live in memory only, so a - A new watcher seeks to end-of-file. Positions live in memory only, so a
restart resumes from the current end of each file and content written while restart resumes from the current end of each file and content written while
LogWisp was down is not read. `from = "start"` reads each file whole when its LogWisp was down is not read.
watcher is created instead — what a process writing beside LogWisp needs, at
the cost of replaying a file already on disk at every restart.
- Rotation is detected from size decrease, modification-time reset, a position - Rotation is detected from size decrease, modification-time reset, a position
beyond end-of-file, or an inode change. An inode change where the new file is beyond end-of-file, or an inode change. An inode change where the new file is
already larger than the recorded position is treated as an atomic save, not a already larger than the recorded position is treated as an atomic save, not a
rotation, and the position is preserved. rotation, and the position is preserved.
- A line is parsed as JSON only when it is an object whose top-level keys are - Lines are parsed as JSON when they contain `time`, `level`, `msg`, and
all drawn from `time`, `level`, `msg` and `fields` — the four an entry can `fields` keys; `time` is read as RFC3339Nano. Anything else is kept as plain
carry. `time` is read as RFC3339Nano. Any other key, and any non-object line, text with the level inferred from common markers (`[ERROR]`, `WARN:`, and so
is kept whole as text with the level inferred from common markers on).
(`[ERROR]`, `WARN:`, and so on), because parsing it would drop the rest.
- `raw = true` skips the JSON branch entirely. The line, plus its newline,
becomes the message; `fields` stays empty, the time is the read time, and the
level is inferred from the text as for any unparsed line. Paired with
`format.type = "raw"` this is byte-exact transport for records LogWisp's
envelope cannot hold — see [Formatters](formatters.md#raw).
- `Source` is set to the file's base name. - `Source` is set to the file's base name.
**Statistics**: per-watcher size, position, entries read, rotation count, and **Statistics**: per-watcher size, position, entries read, rotation count, and
+1 -1
View File
@@ -1,6 +1,6 @@
module logwisp module logwisp
go 1.27.1 go 1.26.5
require ( require (
github.com/lixenwraith/config v0.1.1-0.20260712172228-ccd280ba6a98 github.com/lixenwraith/config v0.1.1-0.20260712172228-ccd280ba6a98
-2
View File
@@ -198,8 +198,6 @@ type FileSourceOptions struct {
Directory string `toml:"directory"` Directory string `toml:"directory"`
Pattern string `toml:"pattern"` // glob pattern Pattern string `toml:"pattern"` // glob pattern
CheckIntervalMS int64 `toml:"check_interval_ms"` CheckIntervalMS int64 `toml:"check_interval_ms"`
Raw bool `toml:"raw"` // keep the whole line as the message, never parse it
From string `toml:"from"` // "end" (default) or "start" of a newly discovered file
} }
// ConsoleSourceOptions defines settings for a stdin-based source // ConsoleSourceOptions defines settings for a stdin-based source
+43 -30
View File
@@ -90,44 +90,57 @@ func NewFormatterAdapter(cfg *config.FormatConfig) (*FormatterAdapter, error) {
// Format implements Formatter interface // Format implements Formatter interface
func (a *FormatterAdapter) Format(entry core.LogEntry) ([]byte, error) { func (a *FormatterAdapter) Format(entry core.LogEntry) ([]byte, error) {
return a.serialize(entry, a.flags), nil // Map logwisp LogEntry to formatter args
level := mapLevel(entry.Level)
// syslog-style origin prefix for chained entries
src := sourceLabel(entry)
// Build args based on whether we have structured fields
var args []any
effectiveFlags := a.flags
if len(entry.Fields) > 0 {
// Parse fields JSON
var fields map[string]any
if err := json.Unmarshal(entry.Fields, &fields); err == nil && len(fields) > 0 {
// Use structured JSON format for fields
args = []any{entry.Message, fields}
// Add structured flag to properly format fields as JSON object
effectiveFlags |= formatter.FlagStructuredJSON
return a.formatter.Format(effectiveFlags, entry.Time, level, src, args), nil
}
}
if args == nil {
args = []any{entry.Message}
}
a.mu.Lock()
out := bytes.Clone(a.formatter.Format(effectiveFlags, entry.Time, level, src, args))
a.mu.Unlock()
return out, nil
} }
// FormatWithFlags allows custom flags for specific formatting needs // FormatWithFlags allows custom flags for specific formatting needs
func (a *FormatterAdapter) FormatWithFlags(entry core.LogEntry, customFlags int64) ([]byte, error) { func (a *FormatterAdapter) FormatWithFlags(entry core.LogEntry, customFlags int64) ([]byte, error) {
return a.serialize(entry, customFlags), nil level := mapLevel(entry.Level)
} src := sourceLabel(entry)
// serialize renders an entry under the given flags. The returned slice is a var args []any
// copy: the underlying formatter reuses one buffer and sinks retain payloads. if len(entry.Fields) > 0 {
func (a *FormatterAdapter) serialize(entry core.LogEntry, flags int64) []byte { var fields map[string]any
args, flags := formatArgs(entry, flags) if err := json.Unmarshal(entry.Fields, &fields); err == nil && len(fields) > 0 {
args = []any{entry.Message, fields}
customFlags |= formatter.FlagStructuredJSON
}
}
if args == nil {
args = []any{entry.Message}
}
a.mu.Lock() a.mu.Lock()
out := bytes.Clone(a.formatter.Format(flags, entry.Time, mapLevel(entry.Level), sourceLabel(entry), args)) out := bytes.Clone(a.formatter.Format(customFlags, entry.Time, level, src, args))
a.mu.Unlock() a.mu.Unlock()
return out return out, nil
}
// formatArgs pairs the entry with its flags. FlagRaw keeps the fields JSON
// verbatim beside the message rather than silently overriding the caller's
// choice of passthrough; every other mode renders it as a JSON object.
func formatArgs(entry core.LogEntry, flags int64) ([]any, int64) {
if len(entry.Fields) == 0 {
return []any{entry.Message}, flags
}
if flags&formatter.FlagRaw != 0 {
if entry.Message == "" {
return []any{[]byte(entry.Fields)}, flags
}
return []any{entry.Message, []byte(entry.Fields)}, flags
}
var fields map[string]any
if err := json.Unmarshal(entry.Fields, &fields); err != nil || len(fields) == 0 {
return []any{entry.Message}, flags
}
return []any{entry.Message, fields}, flags | formatter.FlagStructuredJSON
} }
// Name returns formatter type // Name returns formatter type
+2 -10
View File
@@ -61,7 +61,6 @@ const (
DefaultFileSourcePattern = "*" DefaultFileSourcePattern = "*"
DefaultFileSourceCheckIntervalMS = 100 DefaultFileSourceCheckIntervalMS = 100
MinFileSourceCheckIntervalMS = 10 MinFileSourceCheckIntervalMS = 10
DefaultFileSourceFrom = "end"
) )
// NewFileSourcePlugin creates a file source through plugin factory // NewFileSourcePlugin creates a file source through plugin factory
@@ -91,11 +90,6 @@ func NewFileSourcePlugin(
} else if opts.CheckIntervalMS < MinFileSourceCheckIntervalMS { } else if opts.CheckIntervalMS < MinFileSourceCheckIntervalMS {
return nil, fmt.Errorf("check_interval_ms: must be >= %d", MinFileSourceCheckIntervalMS) return nil, fmt.Errorf("check_interval_ms: must be >= %d", MinFileSourceCheckIntervalMS)
} }
if opts.From == "" {
opts.From = DefaultFileSourceFrom
} else if err := lconfig.OneOf("start", "end")(opts.From); err != nil {
return nil, fmt.Errorf("from: %w", err)
}
// Create and return plugin instance // Create and return plugin instance
fs := &FileSource{ fs := &FileSource{
@@ -122,9 +116,7 @@ func NewFileSourcePlugin(
"component", "file_source", "component", "file_source",
"instance_id", id, "instance_id", id,
"directory", opts.Directory, "directory", opts.Directory,
"pattern", opts.Pattern, "pattern", opts.Pattern)
"raw", opts.Raw,
"from", opts.From)
return fs, nil return fs, nil
} }
@@ -269,7 +261,7 @@ func (fs *FileSource) ensureWatcher(path string) {
return return
} }
w := newFileWatcher(path, fs.config.Raw, fs.config.From == "start", fs.publish, fs.logger) w := newFileWatcher(path, fs.publish, fs.logger)
fs.watchers[path] = w fs.watchers[path] = w
fs.logger.Debug("msg", "Created file watcher", fs.logger.Debug("msg", "Created file watcher",
+30 -67
View File
@@ -34,7 +34,6 @@ type WatcherInfo struct {
type fileWatcher struct { type fileWatcher struct {
directory string directory string
callback func(core.LogEntry) callback func(core.LogEntry)
raw bool
position int64 position int64
size int64 size int64
inode uint64 inode uint64
@@ -47,18 +46,12 @@ type fileWatcher struct {
logger *log.Logger logger *log.Logger
} }
// newFileWatcher creates a new watcher for a specific file path. // newFileWatcher creates a new watcher for a specific file path
// A start position of 0 reads an existing file whole; -1 seeks to its end. func newFileWatcher(directory string, callback func(core.LogEntry), logger *log.Logger) *fileWatcher {
func newFileWatcher(directory string, raw, fromStart bool, callback func(core.LogEntry), logger *log.Logger) *fileWatcher {
position := int64(-1)
if fromStart {
position = 0
}
w := &fileWatcher{ w := &fileWatcher{
directory: directory, directory: directory,
callback: callback, callback: callback,
raw: raw, position: -1,
position: position,
logger: logger, logger: logger,
} }
w.lastReadTime.Store(time.Time{}) w.lastReadTime.Store(time.Time{})
@@ -67,8 +60,8 @@ func newFileWatcher(directory string, raw, fromStart bool, callback func(core.Lo
// watch starts the main monitoring loop for the file // watch starts the main monitoring loop for the file
func (w *fileWatcher) watch(ctx context.Context) error { func (w *fileWatcher) watch(ctx context.Context) error {
if err := w.initPosition(); err != nil { if err := w.seekToEnd(); err != nil {
return fmt.Errorf("initPosition failed: %w", err) return fmt.Errorf("seekToEnd failed: %w", err)
} }
ticker := time.NewTicker(core.FileWatcherPollInterval) ticker := time.NewTicker(core.FileWatcherPollInterval)
@@ -304,9 +297,8 @@ func (w *fileWatcher) checkFile() error {
return nil return nil
} }
// initPosition records the file's metadata and, unless the watcher was created // seekToEnd sets the initial read position to the end of the file
// to read from the start, sets the initial read position to the end func (w *fileWatcher) seekToEnd() error {
func (w *fileWatcher) initPosition() error {
file, err := os.Open(w.directory) file, err := os.Open(w.directory)
if err != nil { if err != nil {
if os.IsNotExist(err) { if os.IsNotExist(err) {
@@ -330,6 +322,8 @@ func (w *fileWatcher) initPosition() error {
w.mu.Lock() w.mu.Lock()
defer w.mu.Unlock() defer w.mu.Unlock()
// Keep existing position (including 0)
// First time initialization seeks to the end of the file
if w.position == -1 { if w.position == -1 {
pos, err := file.Seek(0, io.SeekEnd) pos, err := file.Seek(0, io.SeekEnd)
if err != nil { if err != nil {
@@ -354,67 +348,36 @@ func (w *fileWatcher) isStopped() bool {
return w.stopped return w.stopped
} }
// parseLine converts a line into an entry, as JSON when nothing would be lost // parseLine attempts to parse a line as JSON, falling back to plain text
func (w *fileWatcher) parseLine(line string) core.LogEntry { func (w *fileWatcher) parseLine(line string) core.LogEntry {
if w.raw { var jsonLog struct {
// Newline restored: sinks write the payload as it stands Time string `json:"time"`
Level string `json:"level"`
Message string `json:"msg"`
Fields json.RawMessage `json:"fields"`
}
if err := json.Unmarshal([]byte(line), &jsonLog); err == nil {
timestamp, err := time.Parse(time.RFC3339Nano, jsonLog.Time)
if err != nil {
timestamp = time.Now()
}
return core.LogEntry{ return core.LogEntry{
Time: time.Now(), Time: timestamp,
Source: filepath.Base(w.directory), Source: filepath.Base(w.directory),
Level: source.ExtractLogLevel(line), Level: jsonLog.Level,
Message: line + "\n", Message: jsonLog.Message,
Fields: jsonLog.Fields,
} }
} }
if entry, ok := w.parseJSON(line); ok { level := source.ExtractLogLevel(line)
return entry
}
return core.LogEntry{ return core.LogEntry{
Time: time.Now(), Time: time.Now(),
Source: filepath.Base(w.directory), Source: filepath.Base(w.directory),
Level: source.ExtractLogLevel(line), Level: level,
Message: line, Message: line,
} }
} }
// parseJSON decodes a line into the entry envelope. A top-level key LogEntry
// cannot carry refuses the whole line, so a richer record reaches the pipeline
// as text rather than silently reduced to the four keys kept here.
func (w *fileWatcher) parseJSON(line string) (core.LogEntry, bool) {
if len(line) == 0 || line[0] != '{' {
return core.LogEntry{}, false
}
var obj map[string]json.RawMessage
if err := json.Unmarshal([]byte(line), &obj); err != nil || len(obj) == 0 {
return core.LogEntry{}, false
}
entry := core.LogEntry{Time: time.Now(), Source: filepath.Base(w.directory)}
for key, val := range obj {
var err error
switch key {
case "time":
var ts string
if json.Unmarshal(val, &ts) == nil {
if t, terr := time.Parse(time.RFC3339Nano, ts); terr == nil {
entry.Time = t
}
}
case "level":
err = json.Unmarshal(val, &entry.Level)
case "msg":
err = json.Unmarshal(val, &entry.Message)
case "fields":
entry.Fields = val
default:
return core.LogEntry{}, false
}
if err != nil {
return core.LogEntry{}, false
}
}
return entry, true
}
-92
View File
@@ -1,92 +0,0 @@
#!/usr/bin/env bash
# logwisp file source pass-through test
#
# file src (raw, from=start) --> raw format --> file sink byte-exact relay
# file src (defaults) --> raw format --> file sink no key dropped
#
# The fixture is a record whose envelope is wider than time/level/msg/fields,
# which is what the narrow JSON branch used to reduce to an empty message.
#
# Usage: ./passthrough-test.sh
# Requires: bash 5+, coreutils (timeout). Linux dev host only.
set -u
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
BIN="${LOGWISP_BIN:-$SCRIPT_DIR/../bin/logwisp}"
RUN="$SCRIPT_DIR/run/passthrough"
[[ -x $BIN ]] || { echo "no binary at $BIN; run make build" >&2; exit 1; }
rm -rf "$RUN"
mkdir -p "$RUN/in" "$RUN/out-raw" "$RUN/out-parsed"
cat > "$RUN/in/wide.jsonl" <<'EOF'
{"time":"2026-09-08T21:35:44.178372768-04:00","level":"INFO","sub":"app","run":0,"tick":0,"frame":0,"fields":{"msg":"init begin","mode":"play"}}
{"time":"2026-09-08T21:35:44.178581366-04:00","level":"PROC","run":0,"tick":0,"frame":0,"fields":{"seq":1,"seed":1788917744178374836}}
plain text line, not JSON at all
EOF
conf() { # sink_dir raw
cat <<EOF
quiet = true
status_reporter = false
[logging]
output = "stderr"
level = "error"
[[pipelines]]
name = "passthrough"
[pipelines.flow.format]
type = "raw"
sanitizer_policy = "raw"
[[pipelines.plugin_sources]]
id = "wide"
type = "file"
[pipelines.plugin_sources.config]
directory = "$RUN/in"
pattern = "*.jsonl"
raw = $2
from = "start"
[[pipelines.plugin_sinks]]
id = "out"
type = "file"
[pipelines.plugin_sinks.config]
directory = "$1"
name = "relay"
flush_interval_ms = 100
EOF
}
conf "$RUN/out-raw" true > "$RUN/raw.toml"
conf "$RUN/out-parsed" false > "$RUN/parsed.toml"
for c in raw parsed; do
timeout 5 "$BIN" -c "$RUN/$c.toml" > "$RUN/$c.out" 2>&1
done
fail=0
check() { # label condition_result
if (( $2 )); then echo "PASS: $1"; else echo "FAIL: $1"; fail=1; fi
}
# 1. raw = true relays the file byte for byte
if diff -q "$RUN/in/wide.jsonl" "$RUN/out-raw/relay.log" > /dev/null; then
check "raw = true: output identical to input" 1
else
check "raw = true: output identical to input" 0
diff "$RUN/in/wide.jsonl" "$RUN/out-raw/relay.log" | head -6
fi
# 2. the default parse keeps every key, JSON branch refused on the wide envelope
n=$(grep -c '"sub":"app"' "$RUN/out-parsed/relay.log")
check "defaults: wide envelope reaches the sink whole ($n line(s) carry sub)" $(( n == 1 ))
n=$(grep -c '1788917744178374836' "$RUN/out-parsed/relay.log")
check "defaults: large integers are not re-encoded through float64 ($n)" $(( n == 1 ))
echo "================================================================"
if (( fail == 0 )); then
echo "RESULT: ALL PASS"
else
echo "RESULT: FAILURES — inspect $RUN/"
fi
exit "$fail"