diff --git a/README.md b/README.md index 26d8f1e..4dc9e11 100644 --- a/README.md +++ b/README.md @@ -150,6 +150,7 @@ release memory sooner than the enclosing query would. | `notmuch.TagsIterator`, `notmuch.PropertiesIterator`, `notmuch.PairsIterator`, `notmuch.ValuesIterator` | iteration over tags, message properties, and configuration data | | `notmuch.Error` | the full set of `notmuch_status_t` codes as Zig errors, plus `error.Unexpected` for status codes a function is not documented to return | | `notmuch.compact`, `notmuch.builtWith`, `notmuch.tag_max` | module-level helpers | +| `notmuch.helpers.MessageWriter` | streaming an email through a `std.Io.Writer` into the mail store and indexing it (a convenience unique to these bindings) | Full API documentation is generated from the source and published at . To read it locally instead: diff --git a/src/compile_check.zig b/src/compile_check.zig index 5706203..0d3cad9 100644 --- a/src/compile_check.zig +++ b/src/compile_check.zig @@ -124,6 +124,13 @@ test "every wrapper body compiles" { _ = indexopts.getDecryptPolicy(); indexopts.deinit(); + var message_writer: notmuch.helpers.MessageWriter = undefined; + message_writer = try notmuch.helpers.MessageWriter.init(undefined, db, &.{}, .{}); + _ = message_writer.writer(); + _ = message_writer.filename(); + _ = try message_writer.finish(null); + message_writer.cancel(); + const messages: notmuch.MessagesIterator = undefined; _ = messages.collectTags(); _ = try messages.next(); diff --git a/src/helpers.zig b/src/helpers.zig new file mode 100644 index 0000000..b427b96 --- /dev/null +++ b/src/helpers.zig @@ -0,0 +1,133 @@ +// SPDX-FileCopyrightText: © 2024 Jeffrey C. Ollie +// SPDX-License-Identifier: GPL-3.0-or-later + +//! Convenience helpers layered on top of the bindings. +//! +//! Unlike the rest of the library, these do not wrap libnotmuch APIs +//! directly; they compose the bindings with the Zig standard library. + +const std = @import("std"); + +const Database = @import("Database.zig"); +const IndexOpts = @import("IndexOpts.zig"); +const Message = @import("Message.zig"); + +/// The maildir subdirectory a message is delivered into. +pub const Subdir = enum { + /// Mail that has not yet been seen by the user. + new, + /// Mail the user has already seen. Files delivered here carry an empty + /// `:2,` maildir info suffix, ready for flag synchronization. + cur, +}; + +pub const Options = struct { + /// Which maildir subdirectory the message is delivered into. + subdir: Subdir = .new, +}; + +/// Streams an email into a notmuch database: the message is composed through +/// a `std.Io.Writer` into the mail store, then indexed. +/// +/// A unique maildir filename is generated for the message — prefixed with +/// the current time so that filenames sort in arrival order, like those +/// written by notmuch insert — and delivered into the mail root's `new` or +/// `cur` subdirectory according to `Options.subdir`. The message is written through `std.Io.File.Atomic`, so +/// no partially written file is ever visible in the mail store: the email +/// accumulates in an unnamed ephemeral file and is atomically materialized +/// by `finish` before indexing. `cancel` discards it without a trace. +/// +/// `buffer` is caller-owned and must remain valid until `finish` or `cancel` +/// is called. +pub const MessageWriter = struct { + database: Database, + dir: std.Io.Dir, + atomic: std.Io.File.Atomic, + file_writer: std.Io.File.Writer, + subdir: Subdir, + name_buf: [name_capacity]u8, + name_len: usize, + + const name_capacity = 64; + + /// Prepare to stream an email into the database's mail store. Writes are + /// buffered through `buffer`. + pub fn init(io: std.Io, database: Database, buffer: []u8, options: Options) !MessageWriter { + const root = database.getPath() orelse return error.NoDatabasePath; + + var mw: MessageWriter = undefined; + mw.database = database; + mw.subdir = options.subdir; + + const suffix: []const u8 = switch (options.subdir) { + .new => "", + .cur => ":2,", + }; + var random_bytes: [8]u8 = undefined; + io.random(&random_bytes); + const micros = std.Io.Timestamp.now(io, .real).toMicroseconds(); + // Lead with the delivery time, like notmuch insert's + // .MP. names, so that sorting + // new/ by filename approximates arrival order. + const name = std.fmt.bufPrint(&mw.name_buf, "{d}.M{d}R{x:0>16}.notmuch-zig{s}", .{ + @divTrunc(micros, std.time.us_per_s), + @mod(micros, std.time.us_per_s), + std.mem.readInt(u64, &random_bytes, .big), + suffix, + }) catch unreachable; + mw.name_len = name.len; + + var path_buf: [std.fs.max_path_bytes]u8 = undefined; + const dir_path = std.fmt.bufPrint(&path_buf, "{s}/{t}", .{ root, options.subdir }) catch + return error.NameTooLong; + mw.dir = try std.Io.Dir.openDirAbsolute(io, dir_path, .{}); + errdefer mw.dir.close(io); + mw.atomic = try mw.dir.createFileAtomic(io, name, .{}); + mw.file_writer = mw.atomic.file.writer(io, buffer); + return mw; + } + + /// The `std.Io.Writer` to compose the email into. + pub fn writer(self: *MessageWriter) *std.Io.Writer { + return &self.file_writer.interface; + } + + /// The generated maildir filename (the basename only) the message will + /// be — or after `finish`, has been — delivered under. + pub fn filename(self: *const MessageWriter) []const u8 { + return self.name_buf[0..self.name_len]; + } + + /// Flush the email, atomically materialize it in the mail store, index + /// it into the database, and return the indexed message. The caller owns + /// the returned message. The database must be open in read-write mode. + pub fn finish(self: *MessageWriter, indexopts: ?IndexOpts) !Message { + const io = self.file_writer.io; + defer { + self.atomic.deinit(io); + self.dir.close(io); + self.* = undefined; + } + try self.file_writer.interface.flush(); + + // The struct (and the name buffer inside it) may have moved since + // init, so re-point the atomic file's destination before linking. + self.atomic.dest_sub_path = self.filename(); + try self.atomic.link(io); + + const root = self.database.getPath() orelse return error.NoDatabasePath; + var path_buf: [std.fs.max_path_bytes]u8 = undefined; + const path = std.fmt.bufPrintZ(&path_buf, "{s}/{t}/{s}", .{ root, self.subdir, self.filename() }) catch + return error.NameTooLong; + return self.database.indexFileGetMessage(path, indexopts); + } + + /// Discard the partially written email without indexing it. Nothing ever + /// appears in the mail store. + pub fn cancel(self: *MessageWriter) void { + const io = self.file_writer.io; + self.atomic.deinit(io); + self.dir.close(io); + self.* = undefined; + } +}; diff --git a/src/notmuch.zig b/src/notmuch.zig index cf12c89..ba0c8ce 100644 --- a/src/notmuch.zig +++ b/src/notmuch.zig @@ -24,6 +24,7 @@ pub const Database = @import("Database.zig"); pub const Directory = @import("Directory.zig"); pub const Error = @import("error.zig").Error; pub const FilenamesIterator = @import("FilenamesIterator.zig"); +pub const helpers = @import("helpers.zig"); pub const IndexOpts = @import("IndexOpts.zig"); pub const Message = @import("Message.zig"); pub const MessagesIterator = @import("MessagesIterator.zig"); diff --git a/src/testing.zig b/src/testing.zig index 124fb44..a0f3fc5 100644 --- a/src/testing.zig +++ b/src/testing.zig @@ -9,6 +9,8 @@ const std = @import("std"); const Database = @import("Database.zig"); +const helpers = @import("helpers.zig"); +const Message = @import("Message.zig"); pub const CorpusMessage = struct { filename: []const u8, @@ -31,6 +33,47 @@ pub const TestDatabase = struct { self.database.deinit() catch {}; self.tmp.cleanup(); } + + /// Begin streaming a new message into the maildir using + /// `helpers.MessageWriter`, with the write buffer managed by the test + /// allocator. Compose the email through `TestMessageWriter.writer`, then + /// call `TestMessageWriter.finish` to index it. + pub fn newMessage(self: *TestDatabase, options: helpers.Options) !TestMessageWriter { + const alloc = std.testing.allocator; + + const buffer = try alloc.alloc(u8, 4096); + errdefer alloc.free(buffer); + + return .{ + .inner = try helpers.MessageWriter.init(std.testing.io, self.database, buffer, options), + .buffer = buffer, + }; + } +}; + +/// An email being streamed into the test database. Obtained from +/// `TestDatabase.newMessage`; wraps `helpers.MessageWriter` and owns its +/// write buffer. +pub const TestMessageWriter = struct { + inner: helpers.MessageWriter, + buffer: []u8, + + /// The `std.Io.Writer` to compose the email into. + pub fn writer(self: *TestMessageWriter) *std.Io.Writer { + return self.inner.writer(); + } + + /// Flush and close the message file, index it into the database, and + /// return the indexed message. Frees the helper's resources; the caller + /// owns the returned message. + pub fn finish(self: *TestMessageWriter) !Message { + const alloc = std.testing.allocator; + defer { + alloc.free(self.buffer); + self.* = undefined; + } + return self.inner.finish(null); + } }; /// Create a notmuch database in a temporary directory and index the embedded @@ -107,6 +150,32 @@ pub fn corpusDatabase() !TestDatabase { return .{ .tmp = tmp, .database = database }; } +test TestMessageWriter { + var test_db = try corpusDatabase(); + defer test_db.deinit(); + + var message_writer = try test_db.newMessage(.{}); + const w = message_writer.writer(); + try w.print( + \\From: Frank + \\To: Alice + \\Subject: {s} + \\Date: Fri, 05 Apr 2024 12:00:00 +0000 + \\Message-ID: + \\ + \\Written through a std.Io.Writer. + \\ + , .{"Streamed message"}); + + const message = try message_writer.finish(); + defer message.deinit(); + try std.testing.expectEqualStrings("six@example.org", message.getMessageID() orelse ""); + + const query = try test_db.database.queryCreate("subject:streamed"); + defer query.deinit(); + try std.testing.expectEqual(@as(u32, 1), try query.countMessages()); +} + test corpusDatabase { var test_db = try corpusDatabase(); defer test_db.deinit(); diff --git a/src/tests.zig b/src/tests.zig index 3491519..7039d17 100644 --- a/src/tests.zig +++ b/src/tests.zig @@ -1087,6 +1087,60 @@ test "writes to a read-only database return errors" { try std.testing.expectError(error.ReadOnlyDatabase, message.removeAllProperties(null)); } +test "helpers.MessageWriter atomicity and delivery" { + const io = std.testing.io; + var test_db = try fixture.corpusDatabase(); + defer test_db.deinit(); + const db = test_db.database; + var buffer: [256]u8 = undefined; + + // A cancelled message leaves nothing behind at its destination. + var cancelled = try notmuch.helpers.MessageWriter.init(io, db, &buffer, .{ .subdir = .new }); + try cancelled.writer().writeAll("From: partial \n"); + var cancelled_path_buf: [128]u8 = undefined; + const cancelled_path = try std.fmt.bufPrint(&cancelled_path_buf, "mail/new/{s}", .{cancelled.filename()}); + cancelled.cancel(); + try std.testing.expectError(error.FileNotFound, test_db.tmp.dir.access(io, cancelled_path, .{})); + + // A finished message is materialized in new/ and indexed. The generated + // name leads with the delivery timestamp so filenames sort by arrival. + var finished = try notmuch.helpers.MessageWriter.init(io, db, &buffer, .{ .subdir = .new }); + try std.testing.expect(std.ascii.isDigit(finished.filename()[0])); + try std.testing.expect(std.mem.indexOf(u8, finished.filename(), ".M") != null); + try finished.writer().writeAll( + \\From: Grace + \\To: Alice + \\Subject: Atomic delivery + \\Date: Sat, 06 Apr 2024 09:00:00 +0000 + \\Message-ID: + \\ + \\Materialized atomically. + \\ + ); + const message = try finished.finish(null); + defer message.deinit(); + try std.testing.expectEqualStrings("seven@example.org", message.getMessageID() orelse ""); + try std.testing.expect(std.mem.indexOf(u8, try message.getFilename(), "/new/") != null); + + // Delivery into cur/ carries an empty maildir info suffix. + var seen = try notmuch.helpers.MessageWriter.init(io, db, &buffer, .{ .subdir = .cur }); + try seen.writer().writeAll( + \\From: Heidi + \\To: Alice + \\Subject: Already seen + \\Date: Sat, 06 Apr 2024 10:00:00 +0000 + \\Message-ID: + \\ + \\Delivered as seen mail. + \\ + ); + const seen_message = try seen.finish(null); + defer seen_message.deinit(); + const seen_filename = try seen_message.getFilename(); + try std.testing.expect(std.mem.indexOf(u8, seen_filename, "/cur/") != null); + try std.testing.expect(std.mem.endsWith(u8, seen_filename, ":2,")); +} + test "module helpers" { try std.testing.expect(notmuch.builtWith("compact")); try std.testing.expect(!notmuch.builtWith("time-travel"));