Distributed Concurrent Editor

document.go at [6cd08b113e]
Login

document.go at [6cd08b113e]

File document/document.go artifact d27057e3a3 part of check-in 6cd08b113e


     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
package document

import (
	"errors"
	"fmt"
	"math"
	"sync/atomic"

	"github.com/rs/zerolog"
	bolt "go.etcd.io/bbolt"
	"wellquite.org/actors"
	"wellquite.org/actors/mailbox"
	protocol "wellquite.org/edist/protocol/go"
)

// --- Client ---

type DocumentClientFactory struct {
	factory *actors.BackPressureClientBaseFactory
}

func (self *DocumentClientFactory) NewClient() *DocumentClient {
	return &DocumentClient{BackPressureClientBase: self.factory.NewClient()}
}

type DocumentClient struct {
	*actors.BackPressureClientBase
}

var _ actors.Client = (*DocumentClient)(nil)

type DocumentSubscription struct {
	fun     func(updates []byte)
	Client  *DocumentClient
	expired uint32
}

func (self *DocumentSubscription) Cancel() {
	if atomic.CompareAndSwapUint32(&self.expired, 0, 1) {
		self.Client.Send((*documentUnsubscribeMsg)(self))
	}
}

type documentSubscribeMsg struct {
	actors.MsgSyncBase
	subscription *DocumentSubscription
}
type documentUnsubscribeMsg DocumentSubscription

func (self *DocumentClient) SubscribeToDocumentUpdates(fun func(updates []byte)) *DocumentSubscription {
	subscription := &DocumentSubscription{
		fun:    fun,
		Client: self,
	}
	msg := &documentSubscribeMsg{
		subscription: subscription,
	}
	if self.SendSync(msg, true) {
		return subscription
	} else {
		return nil
	}
}

type documentProtocolMessage struct {
	protoMsg *protocol.Message
	client   *DocumentClient
}

func (self *DocumentClient) Message(msg *protocol.Message) bool {
	return self.Send(&documentProtocolMessage{
		protoMsg: msg,
		client:   self,
	})
}

// --- Server ---

type documentServer struct {
	actors.BackPressureServerBase
	documentName string
	db           *DB
	clientId     uint64

	subscribers map[*DocumentSubscription]struct{}

	wordManager *WordManager

	currentFrame *frame
	loadedEvents frames
}

var _ actors.Server = (*documentServer)(nil)

func SpawnDocument(db *bolt.DB, manager actors.ManagerClient, clientId uint64, documentName string) (*DocumentClientFactory, error) {
	server := &documentServer{
		documentName: documentName,
		db:           NewDB(db, documentName),
		clientId:     clientId,
	}
	clientBase, err := manager.Spawn(server, fmt.Sprintf(`document %q`, documentName))
	if err != nil {
		return nil, err
	}
	return &DocumentClientFactory{
		factory: actors.NewBackPressureClientBaseFactory(clientBase),
	}, nil
}

func (self *documentServer) Init(log zerolog.Logger, mailboxReader *mailbox.MailboxReader, selfClient *actors.ClientBase) (err error) {
	if err := self.BackPressureServerBase.Init(log, mailboxReader, selfClient); err != nil {
		return err
	}

	self.subscribers = make(map[*DocumentSubscription]struct{})
	self.loadEvents(math.MaxUint64)
	self.buildWordsFromEvents()

	if self.currentFrame == nil {
		self.Log.Info().Msg("Creating new document")
		self.maybeCreateCheckpoint()
	} else {
		self.Log.Info().Msg("Loaded existing document")
	}
	return nil
}

func (self *documentServer) HandleMsg(msg mailbox.Msg) (err error) {
	switch msgT := msg.(type) {
	case *documentProtocolMessage:
		protoMsg := msgT.protoMsg
		switch {
		case protoMsg.Updates != nil:
			return self.updates(msgT.client, protoMsg)

		case protoMsg.Undo != nil:
			return self.undo(protoMsg)

		case protoMsg.Redo != nil:
			return self.redo(protoMsg)

		default:
			return errors.New(`Illegal message: updates, undo, redo all nil`)
		}

	case *documentSubscribeMsg:
		msgT.MarkProcessed()
		self.subscribers[msgT.subscription] = struct{}{}
		self.Log.Debug().Int("subscription count", len(self.subscribers)).Msg("subscriber added")
		msgT.subscription.fun(self.wordManager.CreateCheckpoint().MarshalBebop())
		return nil

	case *documentUnsubscribeMsg:
		delete(self.subscribers, (*DocumentSubscription)(msgT))
		self.Log.Debug().Int("subscription count", len(self.subscribers)).Msg("subscriber removed")
		if len(self.subscribers) == 0 {
			return actors.ErrNormalActorTermination
		} else {
			return nil
		}

	default:
		return self.BackPressureServerBase.HandleMsg(msg)
	}
}

func (self *documentServer) updates(sender *DocumentClient, msg *protocol.Message) error {
	self.Log.Trace().Int("Updated words", len(*msg.Updates)).Msg("Received Updates")

	self.maybeCreateCheckpoint()

	filteredUpdates := self.wordManager.FilterUpdatesForNewer(*msg.Updates)
	self.wordManager.ApplyUpdates(filteredUpdates)
	self.wordManager.GC()
	updates := self.wordManager.FilterUpdatesForKnown(filteredUpdates)

	if len(updates) == 0 {
		return nil
	}

	msg.Updates = &updates
	eventNumber := self.db.WriteMessage(msg)

	msgBites := msg.MarshalBebop()
	for subscriber := range self.subscribers {
		if subscriber.Client != sender {
			subscriber.fun(msgBites)
		}
	}

	currentFrame := &frame{
		eventNumber: eventNumber,
		message:     msg,
		previous:    self.currentFrame,
	}
	if self.currentFrame != nil {
		self.currentFrame.next = currentFrame
	}
	self.currentFrame = currentFrame
	self.loadedEvents = append(self.loadedEvents, currentFrame)

	return nil
}

func (self *documentServer) undo(msg *protocol.Message) error {
	if self.currentFrame == nil {
		// illegal undo. Ignore it.
		self.Log.Trace().Msg("Received illegal Undo.")

		return nil
	}

	self.loadEvents(self.currentFrame.eventNumber)
	self.buildWordsFromEvents()
	if self.currentFrame.previous == nil {
		// illegal undo. Ignore it.
		self.Log.Trace().Msg("Received illegal Undo.")
		return nil
	}
	self.Log.Trace().Msg("Received Undo.")

	oldWordManager := self.wordManager

	currentFrame := self.currentFrame
	currentFrame.undoCount++
	self.loadedEvents = append(self.loadedEvents, currentFrame)
	self.buildWordsFromEvents()

	self.db.WriteMessage(msg)

	deltaMsgBites := self.wordManager.Delta(oldWordManager).MarshalBebop()
	for subscriber := range self.subscribers {
		subscriber.fun(deltaMsgBites)
	}

	return nil
}

func (self *documentServer) redo(msg *protocol.Message) error {
	if self.currentFrame == nil {
		// illegal redo. Ignore it.
		self.Log.Trace().Msg("Received illegal Redo")
		return nil
	}

	self.loadEvents(self.currentFrame.eventNumber)
	self.buildWordsFromEvents()

	currentFrame := self.currentFrame.next
	if currentFrame == nil {
		// illegal redo. Ignore it.
		self.Log.Trace().Msg("Received illegal Redo")
		return nil
	}
	self.Log.Trace().Msg("Received Redo")

	oldWordManager := self.wordManager

	currentFrame.redoCount++
	self.loadedEvents = append(self.loadedEvents, currentFrame)
	self.buildWordsFromEvents()

	self.db.WriteMessage(msg)

	deltaMsgBites := self.wordManager.Delta(oldWordManager).MarshalBebop()
	for subscriber := range self.subscribers {
		subscriber.fun(deltaMsgBites)
	}

	return nil
}

func (self *documentServer) loadEvents(maxEventNumber uint64) {
	if len(self.loadedEvents) > 0 {
		frame0 := self.loadedEvents[0]
		if frame0.eventNumber == 1 {
			// can't go back any further anyway
			return
		} else if frame0.isCheckpoint && frame0.eventNumber < maxEventNumber {
			// already loaded enough
			return
		}
	}
	self.loadedEvents = self.db.LoadEvents(maxEventNumber)
}

type frameApplyState struct {
	updates   []*UpdatePair
	undoCount uint64
	redoCount uint64
	undoNext  bool
}

func (self *documentServer) buildWordsFromEvents() {
	wordMgr := NewWordManager()
	states := make(map[*frame]*frameApplyState, len(self.loadedEvents))

	var currentFrame *frame
	for _, frame := range self.loadedEvents {
		fas, found := states[frame]
		if !found {
			if currentFrame != nil {
				currentFrame.next = frame
				frame.previous = currentFrame
			}
			currentFrame = frame

			// first time we've seen it, so this is the initial update
			fas = &frameApplyState{
				updates:   wordMgr.FilterUpdatesForNewer(*frame.message.Updates),
				undoCount: frame.undoCount,
				redoCount: frame.redoCount,
				undoNext:  true,
			}
			states[frame] = fas
			// even if we don't end up applying this update, we still
			// need to keep the version applied correctly
			for _, update := range fas.updates {
				if update.updated.Version.GreaterThan(update.original.Version) {
					update.original.Version = update.updated.Version
				}
			}

		} else {
			if fas.undoNext {
				if currentFrame != frame {
					panic(fmt.Sprintf("undo. currentFrame %v, frame %v", currentFrame, frame))
				}
				currentFrame = currentFrame.previous
				fas.undoCount--

			} else { // redoNext
				currentFrame = currentFrame.next
				if currentFrame != frame {
					panic(fmt.Sprintf("redo. currentFrame %v, frame %v", currentFrame, frame))
				}
				fas.redoCount--
			}
			fas.undoNext = !fas.undoNext

			for _, update := range fas.updates {
				// only bump it if the update was a modification of an
				// existing word. If the update was adding this word, then
				// we don't bump it.
				if !update.isNewWord {
					update.original.Version.Bump(self.clientId)
				}
			}
		}

		if fas.undoCount == 0 && fas.redoCount == 0 && fas.undoNext {
			// simplest case is undoCount=0, redoCount=0. So undoNext will
			// be true.
			wordMgr.ApplyUpdates(fas.updates)
		}
	}

	wordMgr.GC()
	self.wordManager = wordMgr
	self.currentFrame = currentFrame
}

const checkpointEvery = 16

func (self *documentServer) maybeCreateCheckpoint() {
	if self.db.largestUsedEventNumber%checkpointEvery == 0 {
		self.Log.Trace().Msg("Creating checkpoint")
		checkpoint := self.wordManager.CreateCheckpoint()
		eventNumber := self.db.WriteCheckpoint(checkpoint)

		currentFrame := &frame{
			eventNumber:  eventNumber,
			message:      checkpoint,
			isCheckpoint: true,
		}
		self.loadedEvents = frames{currentFrame}
		self.currentFrame = currentFrame
	}
}