File size: 19,108 Bytes
b8370f5
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
import Foundation

/// Read-only offline cache seam for chat sessions and transcripts.
///
/// The cache only pre-paints cold opens and covers offline browsing; connected
/// reads always come from the gateway and replace cached content wholesale.
/// Implementations must scope every row by gateway and agent identity so one
/// shared installation database can safely serve all paired gateways.
public protocol OpenClawChatTranscriptCache: Sendable {
    func loadSessions() async -> [OpenClawChatSessionEntry]
    func loadSessions(agentID: String?) async -> [OpenClawChatSessionEntry]
    func loadTranscript(sessionKey: String) async -> [OpenClawChatMessage]
    func loadTranscript(sessionKey: String, agentID: String?) async -> [OpenClawChatMessage]
    func storeSessions(_ sessions: [OpenClawChatSessionEntry]) async
    func storeSessions(_ sessions: [OpenClawChatSessionEntry], agentID: String?) async
    /// Canonical gateway rows can prove that an ambiguously delivered local
    /// command landed after cancellation and must override local suppression.
    func storeCanonicalTranscript(
        sessionKey: String,
        agentID: String?,
        messages: [OpenClawChatMessage],
        canonicalMessageIdempotencyKeys: Set<String>) async
    /// Synchronous observation closes the session.message -> cancellation
    /// race before asynchronous SQLite confirmation starts.
    func observeCanonicalMessageIdempotencyKeys(_ keys: Set<String>)
}

extension OpenClawChatTranscriptCache {
    public func loadSessions(agentID: String?) async -> [OpenClawChatSessionEntry] {
        // Legacy conformers have no agent partition. Scoped access must fail
        // closed or an ownerless roster can cross an agent switch.
        guard agentID == nil else { return [] }
        return await self.loadSessions()
    }

    public func storeSessions(_ sessions: [OpenClawChatSessionEntry], agentID: String?) async {
        guard agentID == nil else { return }
        await self.storeSessions(sessions)
    }

    public func loadTranscript(sessionKey: String, agentID: String?) async -> [OpenClawChatMessage] {
        guard agentID == nil else { return [] }
        return await self.loadTranscript(sessionKey: sessionKey)
    }

    public func observeCanonicalMessageIdempotencyKeys(_: Set<String>) {}
}

/// Optional atomic merge seam for cache owners that also provide a durable
/// outbox. Keeping this separate preserves source compatibility for read-only
/// transcript-cache conformers.
protocol OpenClawChatCanonicalTranscriptMerging: OpenClawChatTranscriptCache {
    func mergeCanonicalTranscriptMessage(
        sessionKey: String,
        agentID: String?,
        message: OpenClawChatMessage,
        canonicalMessageIdempotencyKey: String) async
}

/// Durable branch ownership is scoped exactly like outbox delivery routing.
public struct OpenClawChatOutboxScope: Hashable, Sendable {
    public let sessionKey: String
    public let agentID: String?

    public init(sessionKey: String, agentID: String?) {
        self.sessionKey = sessionKey
        let normalizedAgentID = agentID?.trimmingCharacters(in: .whitespacesAndNewlines).lowercased()
        self.agentID = normalizedAgentID?.isEmpty == false ? normalizedAgentID : nil
    }
}

/// Persisted branch ownership captured before bootstrap can advance the transcript tip.
public struct OpenClawChatOutboxBranchState: Equatable, Sendable {
    public let epoch: Int
    public let lastActiveLeafEntryID: String?
    public let hadPendingCommands: Bool
    public let switchPendingSince: TimeInterval?
    public let needsReconciliation: Bool
    public let revision: Int

    public init(
        epoch: Int,
        lastActiveLeafEntryID: String?,
        hadPendingCommands: Bool = false,
        switchPendingSince: TimeInterval? = nil,
        needsReconciliation: Bool = false,
        revision: Int = 0)
    {
        self.epoch = epoch
        self.lastActiveLeafEntryID = lastActiveLeafEntryID
        self.hadPendingCommands = hadPendingCommands
        self.switchPendingSince = switchPendingSince
        self.needsReconciliation = needsReconciliation
        self.revision = revision
    }
}

public struct OpenClawChatOutboxRetryExpectation: Equatable, Sendable {
    public let attemptVersion: Int
    public let retryCount: Int
    public let lastError: String?

    public init(attemptVersion: Int, retryCount: Int, lastError: String?) {
        self.attemptVersion = attemptVersion
        self.retryCount = retryCount
        self.lastError = lastError
    }
}

/// One attachment captured with a durable chat command.
public struct OpenClawChatOutboxAttachment: Codable, Hashable, Sendable {
    public let type: String
    public let mimeType: String
    public let fileName: String
    public let data: Data
    public let durationSeconds: Double?

    public init(
        type: String,
        mimeType: String,
        fileName: String,
        data: Data,
        durationSeconds: Double? = nil)
    {
        self.type = type
        self.mimeType = mimeType
        self.fileName = fileName
        self.data = data
        self.durationSeconds = durationSeconds
    }
}

/// One durable queued chat command. `id` is the client UUID
/// that becomes the transport idempotency key on flush, so at-least-once
/// delivery stays safe across retries and app restarts.
///
/// Naming mirrors the watch-side `QueuedCommand` shape (WatchChatCoordinator)
/// so the two queues can merge into one owner later.
public struct OpenClawChatOutboxCommand: Hashable, Sendable, Identifiable {
    static let legacyUnboundRoutingContract = "legacy-unbound"

    public enum Status: String, Sendable {
        case queued
        case sending
        case awaitingConfirmation = "awaiting_confirmation"
        case failed
    }

    public let id: String
    /// Presentation/cache key captured when the user queued the command.
    public let sessionKey: String
    /// Canonical transport key captured at enqueue time. This must never be
    /// re-resolved from a mutable main/default alias during reconnect.
    public let deliverySessionKey: String
    /// Gateway main-routing contract (scope, main key, default agent) captured
    /// with the command. A changed contract must fail closed before replay.
    public let routingContract: String?
    /// Durable routing owner, required for the literal `global` session and
    /// retained for ownership checks on canonical agent-scoped keys.
    public let agentID: String?
    /// Local branch generation captured when this delivery attempt was queued.
    public let branchEpoch: Int
    /// Scope epoch observed alongside this row snapshot.
    public let scopeBranchEpoch: Int?
    public let text: String
    /// Attachment bytes remain owned by SQLite until canonical history proves
    /// delivery or the user explicitly deletes the command.
    public let attachments: [OpenClawChatOutboxAttachment]
    /// Thinking level captured when the command was queued, so a later flush
    /// never borrows the setting of whichever session is visible then.
    public let thinking: String
    /// Permission and tool state captured with this command. Durable replay
    /// must use this command-owned fence, never the currently visible session.
    public let expectedSessionSettings: OpenClawChatSessionSettingsExpectation?
    /// Seconds since 1970; flush order is strictly ascending `createdAt`.
    public let createdAt: Double
    public var status: Status
    /// Immutable ownership token for one delivery lifecycle. Every automatic
    /// or user-initiated retry increments it before another send can start.
    public let attemptVersion: Int
    public var retryCount: Int
    public var lastError: String?

    public init(
        id: String,
        sessionKey: String,
        deliverySessionKey: String? = nil,
        routingContract: String? = nil,
        agentID: String? = nil,
        branchEpoch: Int = 0,
        scopeBranchEpoch: Int? = nil,
        text: String,
        attachments: [OpenClawChatOutboxAttachment] = [],
        thinking: String,
        expectedSessionSettings: OpenClawChatSessionSettingsExpectation? = nil,
        createdAt: Double,
        status: Status,
        attemptVersion: Int = 1,
        retryCount: Int,
        lastError: String?)
    {
        self.id = id
        self.sessionKey = sessionKey
        if let deliverySessionKey {
            self.deliverySessionKey = deliverySessionKey.trimmingCharacters(in: .whitespacesAndNewlines)
        } else {
            self.deliverySessionKey = sessionKey
        }
        let normalizedRoutingContract = routingContract?.trimmingCharacters(in: .whitespacesAndNewlines)
        self.routingContract = normalizedRoutingContract?.isEmpty == false ? normalizedRoutingContract : nil
        let normalizedAgentID = agentID?.trimmingCharacters(in: .whitespacesAndNewlines).lowercased()
        self.agentID = normalizedAgentID?.isEmpty == false ? normalizedAgentID : nil
        self.branchEpoch = branchEpoch
        self.scopeBranchEpoch = scopeBranchEpoch ?? branchEpoch
        self.text = text
        self.attachments = attachments
        self.thinking = thinking
        self.expectedSessionSettings = expectedSessionSettings
        self.createdAt = createdAt
        self.status = status
        self.attemptVersion = attemptVersion
        self.retryCount = retryCount
        self.lastError = lastError
    }
}

public enum OpenClawChatOutboxUpdateResult: Equatable, Sendable {
    case updated
    case confirmed
    case missing
    case superseded
    case unavailable
}

public enum OpenClawChatOutboxChange: Equatable, Sendable {
    case canceled(gatewayID: String, id: String)
    case confirmed(gatewayID: String, id: String)
    case invalidated(gatewayID: String, scope: OpenClawChatOutboxScope)

    var gatewayID: String {
        switch self {
        case let .canceled(gatewayID, _), let .confirmed(gatewayID, _), let .invalidated(gatewayID, _):
            gatewayID
        }
    }
}

/// Durable offline outbox for chat commands. Implementations expose one
/// gateway-scoped facade over installation-wide client state so queued sends
/// survive app restarts and flush on reconnect.
public protocol OpenClawChatCommandOutbox: Sendable {
    /// Returns false when the row or attachment-byte budget is full, or
    /// storage is unavailable; callers surface that instead of dropping text.
    func enqueueCommand(_ command: OpenClawChatOutboxCommand) async -> Bool
    /// Gateway-scoped rows in `createdAt` order. Applies the staleness gate:
    /// old queued or unconfirmed rows become failed so reconnect never sends
    /// stale or ambiguously delivered commands silently.
    func loadCommands() async -> [OpenClawChatOutboxCommand]
    /// Availability-aware read used by the FIFO restoration gate. Nil means
    /// storage was not readable, not that the queue was empty.
    func loadCommandsIfAvailable() async -> [OpenClawChatOutboxCommand]?
    /// Crash safety: rows stuck in 'sending' from a previous process become
    /// failed once per store lifetime. Delivery is ambiguous after a crash,
    /// so only explicit user retry may replay them; acknowledged rows stay
    /// awaiting canonical history confirmation.
    /// Returns false while storage is unavailable so callers can retry later.
    @discardableResult
    func recoverInterruptedSends() async -> Bool
    /// Atomically claims the oldest queued row when no other row is sending.
    /// Nil means another flusher owns the queue or no deliverable row remains.
    func claimNextCommand() async -> OpenClawChatOutboxCommand?
    /// Safe automatic retry: only the completing attempt may requeue the row,
    /// and a successful requeue mints the next attempt version atomically.
    func markCommandQueued(
        id: String,
        attemptVersion: Int,
        retryCount: Int,
        lastError: String?) async -> OpenClawChatOutboxUpdateResult
    func markCommandAwaitingConfirmation(
        id: String,
        attemptVersion: Int) async -> OpenClawChatOutboxUpdateResult
    /// Result-bearing terminal transition for callers that must stop their
    /// FIFO when durable storage is unavailable.
    func markCommandFailedIfPresent(
        id: String,
        attemptVersion: Int,
        retryCount: Int,
        lastError: String?) async -> OpenClawChatOutboxUpdateResult
    /// Captures the persisted scope state before bootstrap history can advance its tip.
    func branchState(for scope: OpenClawChatOutboxScope) async -> OpenClawChatOutboxBranchState?
    /// Installs the cross-view-model transcript-mutation barrier only when no
    /// delivery is already unresolved for the scope.
    func beginBranchSwitch(_ scope: OpenClawChatOutboxScope) async -> Bool
    /// Rolls back a barrier when the server rejected the switch.
    func cancelBranchSwitch(_ scope: OpenClawChatOutboxScope) async -> Bool
    /// The server changed the branch but local refresh failed; block replay
    /// until reconciliation establishes the active leaf.
    func demoteBranchSwitchToReconcile(_ scope: OpenClawChatOutboxScope) async -> Bool
    /// Reconciles a bootstrap branch snapshot before automatic replay is enabled.
    /// A nil active leaf represents a successfully listed empty transcript.
    func reconcileBranchScope(
        _ scope: OpenClawChatOutboxScope,
        previousState: OpenClawChatOutboxBranchState,
        activeLeafEntryID: String?,
        branchLeafEntryIDs: Set<String>,
        activeTranscriptEntryIDs: Set<String>,
        lastError: String) async -> [OpenClawChatOutboxCommand]?
    /// Atomically records a confirmed server-side branch change and parks rows
    /// stamped with the superseded generation.
    func confirmBranchChange(
        _ scope: OpenClawChatOutboxScope,
        activeLeafEntryID: String,
        lastError: String) async -> [OpenClawChatOutboxCommand]?
    /// Advances the observed transcript tip only while branch ownership still
    /// matches the epoch captured by the caller.
    func updateLastActiveLeafEntryID(
        _ leafEntryID: String,
        expectedEpoch: Int,
        for scope: OpenClawChatOutboxScope) async -> Bool
    /// Retry only if the failed row still matches the version shown to the user.
    /// The default fails closed so a store cannot bypass branch-change parking.
    func markCommandRetriedIfPresent(
        id: String,
        expectation: OpenClawChatOutboxRetryExpectation,
        agentID: String?,
        deliverySessionKey: String,
        routingContract: String,
        expectedSessionSettings: OpenClawChatSessionSettingsExpectation,
        replacementID: String?) async -> OpenClawChatOutboxUpdateResult
    /// Persistently parks automatic replay after a failed settings mutation.
    func parkQueuedCommands(
        in scope: OpenClawChatOutboxScope,
        lastError: String) async -> Bool
    /// User cancellation succeeds only before a sender claims the row. The
    /// status predicate is the cross-view-model cancellation boundary.
    func cancelCommand(id: String) async -> OpenClawChatOutboxUpdateResult
    /// Canonical gateway history may complete the matching attempt, including
    /// a sending row whose request ACK was lost.
    func confirmCommand(id: String, attemptVersion: Int) async -> OpenClawChatOutboxUpdateResult
    /// Cross-view-model invalidation.
    func changes() -> AsyncStream<OpenClawChatOutboxChange>
}

extension OpenClawChatCommandOutbox {
    public func parkQueuedCommands(
        in _: OpenClawChatOutboxScope,
        lastError _: String) async -> Bool
    {
        false
    }

    public func markCommandQueued(
        id _: String,
        attemptVersion _: Int,
        retryCount _: Int,
        lastError _: String?) async -> OpenClawChatOutboxUpdateResult
    {
        .unavailable
    }

    public func markCommandAwaitingConfirmation(
        id _: String,
        attemptVersion _: Int) async -> OpenClawChatOutboxUpdateResult
    {
        .unavailable
    }

    public func markCommandFailedIfPresent(
        id _: String,
        attemptVersion _: Int,
        retryCount _: Int,
        lastError _: String?) async -> OpenClawChatOutboxUpdateResult
    {
        .unavailable
    }

    public func confirmCommand(
        id _: String,
        attemptVersion _: Int) async -> OpenClawChatOutboxUpdateResult
    {
        .unavailable
    }

    public func branchState(for _: OpenClawChatOutboxScope) async -> OpenClawChatOutboxBranchState? {
        nil
    }

    public func beginBranchSwitch(_: OpenClawChatOutboxScope) async -> Bool {
        false
    }

    public func cancelBranchSwitch(_: OpenClawChatOutboxScope) async -> Bool {
        false
    }

    public func demoteBranchSwitchToReconcile(_: OpenClawChatOutboxScope) async -> Bool {
        false
    }

    public func reconcileBranchScope(
        _: OpenClawChatOutboxScope,
        previousState _: OpenClawChatOutboxBranchState,
        activeLeafEntryID _: String?,
        branchLeafEntryIDs _: Set<String>,
        activeTranscriptEntryIDs _: Set<String>,
        lastError _: String) async -> [OpenClawChatOutboxCommand]?
    {
        nil
    }

    public func confirmBranchChange(
        _: OpenClawChatOutboxScope,
        activeLeafEntryID _: String,
        lastError _: String) async -> [OpenClawChatOutboxCommand]?
    {
        nil
    }

    public func updateLastActiveLeafEntryID(
        _: String,
        expectedEpoch _: Int,
        for _: OpenClawChatOutboxScope) async -> Bool
    {
        false
    }

    // periphery:ignore - protocol-typed callers require this forwarding convenience overload.
    public func markCommandRetriedIfPresent(
        id _: String,
        expectation _: OpenClawChatOutboxRetryExpectation,
        agentID _: String?,
        deliverySessionKey _: String,
        routingContract _: String,
        expectedSessionSettings _: OpenClawChatSessionSettingsExpectation,
        replacementID _: String? = nil) async -> OpenClawChatOutboxUpdateResult
    {
        .unavailable
    }
}

public struct OpenClawChatSessionRoutingIdentity: Equatable, Sendable {
    public let scope: String
    public let mainSessionKey: String
    public let defaultAgentID: String
    public let contract: String

    public init?(contract: String?) {
        guard let components = OpenClawChatSessionRoutingContract.parse(contract) else { return nil }
        self.scope = components.scope
        self.mainSessionKey = components.mainKey
        self.defaultAgentID = components.defaultAgentID
        self.contract = "\(components.scope)|\(components.mainKey)|\(components.defaultAgentID)"
    }

    public init?(scope: String?, mainSessionKey: String?, defaultAgentID: String?) {
        guard let contract = OpenClawChatSessionRoutingContract.make(
            scope: scope,
            mainKey: mainSessionKey,
            defaultAgentID: defaultAgentID)
        else { return nil }
        self.init(contract: contract)
    }
}