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)
}
}
|