clusternode.java
来自「jsr170接口的java实现。是个apache的开源项目。」· Java 代码 · 共 1,278 行 · 第 1/3 页
JAVA
1,278 行
* @param register <code>true</code>, if this is a register operation; * <code>false</code> otherwise */ private void process(Collection c, boolean register) { if (nodeTypeListener == null) { String msg = "NodeType listener unavailable."; log.error(msg); return; } try { if (register) { nodeTypeListener.externalRegistered(c); } else { nodeTypeListener.externalUnregistered(c); } } catch (InvalidNodeTypeDefException e) { String msg = "Unable to deliver node type operation: " + e.getMessage(); log.error(msg); } catch (RepositoryException e) { String msg = "Unable to deliver node type operation: " + e.getMessage(); log.error(msg); } } /** * Process a node type re-registration. * * @param ntDef node type definition */ private void process(NodeTypeDef ntDef) { if (nodeTypeListener == null) { String msg = "NodeType listener unavailable."; log.error(msg); return; } try { nodeTypeListener.externalReregistered(ntDef); } catch (InvalidNodeTypeDefException e) { String msg = "Unable to deliver node type operation: " + e.getMessage(); log.error(msg); } catch (RepositoryException e) { String msg = "Unable to deliver node type operation: " + e.getMessage(); log.error(msg); } } /** * Invoked when a record ends. */ private void end() { UpdateEventListener listener = null; if (workspace != null) { listener = (UpdateEventListener) wspUpdateListeners.get(workspace); if (listener == null) { try { clusterContext.updateEventsReady(workspace); } catch (RepositoryException e) { String msg = "Error making update listener for workspace " + workspace + " online: " + e.getMessage(); log.warn(msg); } listener = (UpdateEventListener) wspUpdateListeners.get(workspace); if (listener == null) { String msg = "Update listener unavailable for workspace: " + workspace; log.error(msg); return; } } } else { if (versionUpdateListener != null) { listener = versionUpdateListener; } else { String msg = "Version update listener unavailable."; log.error(msg); return; } } try { listener.externalUpdate(changeLog, events); } catch (RepositoryException e) { String msg = "Unable to deliver update events: " + e.getMessage(); log.error(msg); } } //-------------------------------------------------------< RecordConsumer > /** * {@inheritDoc} */ public String getId() { return PRODUCER_ID; } /** * {@inheritDoc} */ public long getRevision() { try { return instanceRevision.get(); } catch (JournalException e) { log.warn("Unable to return current revision.", e); return Long.MAX_VALUE; } } /** * {@inheritDoc} */ public void consume(Record record) { log.info("Processing revision: " + record.getRevision()); String workspace = null; try { workspace = record.readString(); start(workspace); for (;;) { char c = record.readChar(); if (c == '\0') { break; } if (c == 'N') { NodeOperation operation = NodeOperation.create(record.readByte()); operation.setId(record.readNodeId()); process(operation); } else if (c == 'P') { PropertyOperation operation = PropertyOperation.create(record.readByte()); operation.setId(record.readPropertyId()); process(operation); } else if (c == 'E') { int type = record.readByte(); NodeId parentId = record.readNodeId(); Path parentPath = record.readPath(); NodeId childId = record.readNodeId(); Path.PathElement childRelPath = record.readPathElement(); QName ntName = record.readQName(); Set mixins = new HashSet(); int mixinCount = record.readInt(); for (int i = 0; i < mixinCount; i++) { mixins.add(record.readQName()); } String userId = record.readString(); process(createEventState(type, parentId, parentPath, childId, childRelPath, ntName, mixins, userId)); } else if (c == 'L') { NodeId nodeId = record.readNodeId(); boolean isLock = record.readBoolean(); if (isLock) { boolean isDeep = record.readBoolean(); String owner = record.readString(); process(nodeId, isDeep, owner); } else { process(nodeId); } } else if (c == 'S') { String oldPrefix = record.readString(); String newPrefix = record.readString(); String uri = record.readString(); process(oldPrefix, newPrefix, uri); } else if (c == 'T') { int size = record.readInt(); int opcode = size & NTREG_MASK; size &= ~NTREG_MASK; switch (opcode) { case NTREG_REGISTER: HashSet ntDefs = new HashSet(); for (int i = 0; i < size; i++) { ntDefs.add(record.readNodeTypeDef()); } process(ntDefs, true); break; case NTREG_REREGISTER: process(record.readNodeTypeDef()); break; case NTREG_UNREGISTER: HashSet ntNames = new HashSet(); for (int i = 0; i < size; i++) { ntNames.add(record.readQName()); } process(ntNames, false); break; default: throw new IllegalArgumentException("Unknown opcode: " + opcode); } } else { throw new IllegalArgumentException("Unknown entry type: " + c); } } end(); } catch (JournalException e) { String msg = "Unable to read revision '" + record.getRevision() + "'."; log.error(msg, e); } catch (IllegalArgumentException e) { String msg = "Error while processing revision " + record.getRevision() + ": " + e.getMessage(); log.error(msg); } } /** * {@inheritDoc} */ public void setRevision(long revision) { try { instanceRevision.set(revision); } catch (JournalException e) { log.warn("Unable to set current revision to " + revision + ".", e); } } /** * Create an event state. * * @param type event type * @param parentId parent id * @param parentPath parent path * @param childId child id * @param childRelPath child relative path * @param ntName ndoe type name * @param userId user id * @return event */ private EventState createEventState(int type, NodeId parentId, Path parentPath, NodeId childId, Path.PathElement childRelPath, QName ntName, Set mixins, String userId) { switch (type) { case Event.NODE_ADDED: return EventState.childNodeAdded(parentId, parentPath, childId, childRelPath, ntName, mixins, getOrCreateSession(userId), true); case Event.NODE_REMOVED: return EventState.childNodeRemoved(parentId, parentPath, childId, childRelPath, ntName, mixins, getOrCreateSession(userId), true); case Event.PROPERTY_ADDED: return EventState.propertyAdded(parentId, parentPath, childRelPath, ntName, mixins, getOrCreateSession(userId), true); case Event.PROPERTY_CHANGED: return EventState.propertyChanged(parentId, parentPath, childRelPath, ntName, mixins, getOrCreateSession(userId), true); case Event.PROPERTY_REMOVED: return EventState.propertyRemoved(parentId, parentPath, childRelPath, ntName, mixins, getOrCreateSession(userId), true); default: String msg = "Unexpected event type: " + type; throw new IllegalArgumentException(msg); } } /** * Return a session matching a certain user id. * * @param userId user id * @return session */ private Session getOrCreateSession(String userId) { if (lastSession == null || !lastSession.getUserID().equals(userId)) { lastSession = new ClusterSession(userId); } return lastSession; } //-----------------------------------------------< Record writing methods > private static void write(Record record, ChangeLog changeLog, EventStateCollection esc) throws JournalException { Iterator deletedStates = changeLog.deletedStates(); while (deletedStates.hasNext()) { ItemState state = (ItemState) deletedStates.next(); if (state.isNode()) { write(record, NodeDeletedOperation.create((NodeState) state)); } else { write(record, PropertyDeletedOperation.create((PropertyState) state)); } } Iterator modifiedStates = changeLog.modifiedStates(); while (modifiedStates.hasNext()) { ItemState state = (ItemState) modifiedStates.next(); if (state.isNode()) { write(record, NodeModifiedOperation.create((NodeState) state)); } else { write(record, PropertyModifiedOperation.create((PropertyState) state)); } } Iterator addedStates = changeLog.addedStates(); while (addedStates.hasNext()) { ItemState state = (ItemState) addedStates.next(); if (state.isNode()) { write(record, NodeAddedOperation.create((NodeState) state)); } else { write(record, PropertyAddedOperation.create((PropertyState) state)); } } Iterator events = esc.getEvents().iterator(); while (events.hasNext()) { EventState event = (EventState) events.next(); write(record, event); } } private static void write(Record record, String oldPrefix, String newPrefix, String uri) throws JournalException { record.writeChar('S'); record.writeString(oldPrefix); record.writeString(newPrefix); record.writeString(uri); } private static void write(Record record, Collection c, boolean register) throws JournalException { record.writeChar('T'); int size = c.size(); if (!register) { size |= NTREG_UNREGISTER; } record.writeInt(size); Iterator iter = c.iterator(); while (iter.hasNext()) { if (register) { record.writeNodeTypeDef((NodeTypeDef) iter.next()); } else { record.writeQName((QName) iter.next()); } } } private static void write(Record record, NodeTypeDef ntDef) throws JournalException { record.writeChar('T'); int size = 1; size |= NTREG_REREGISTER; record.writeInt(size); record.writeNodeTypeDef(ntDef); } private static void write(Record record, PropertyOperation operation) throws JournalException { record.writeChar('P'); record.writeByte(operation.getOperationType()); record.writePropertyId(operation.getId()); } private static void write(Record record, NodeOperation operation) throws JournalException { record.writeChar('N'); record.writeByte(operation.getOperationType()); record.writeNodeId(operation.getId()); } /** * Log an event. Subclass responsibility. * * @param event event to log */ private static void write(Record record, EventState event) throws JournalException { record.writeChar('E'); record.writeByte(event.getType()); record.writeNodeId(event.getParentId()); record.writePath(event.getParentPath()); record.writeNodeId(event.getChildId()); record.writePathElement(event.getChildRelPath()); record.writeQName(event.getNodeType()); Set mixins = event.getMixinNames(); record.writeInt(mixins.size()); Iterator iter = mixins.iterator(); while (iter.hasNext()) { record.writeQName((QName) iter.next()); } record.writeString(event.getUserId()); } /** * Invoked when a cluster operation has ended. If <code>successful</code>, * attempts to fill the journal record and update it, otherwise cancels * the update. * * @param operation cluster operation * @param successful <code>true</code> if the operation was successful and * the journal record should be updated; * <code>false</code> to revoke changes */ public void ended(AbstractClusterOperation operation, boolean successful) { Record record = operation.getRecord(); boolean succeeded = false; try { if (successful) { record = operation.getRecord(); record.writeString(operation.getWorkspace()); operation.write(); record.writeChar('\0'); record.update(); setRevision(record.getRevision()); succeeded = true; } } catch (JournalException e) { String msg = "Unable to create log entry: " + e.getMessage(); log.error(msg); } catch (Throwable e) { String msg = "Unexpected error while creating log entry."; log.error(msg, e); } finally { if (!succeeded) { record.cancelUpdate(); } } }}
⌨️ 快捷键说明
复制代码Ctrl + C
搜索代码Ctrl + F
全屏模式F11
增大字号Ctrl + =
减小字号Ctrl + -
显示快捷键?